from __future__ import annotations

import asyncio
import base64
from datetime import date
import hashlib
import logging
from typing import Any, Protocol
from gateway.platforms.telegram_nutrition_onboarding import (
    OPTIONAL_QUESTION_DEFAULTS,
    decode_callback,
    log_ingress_stage,
)
from gateway.platforms.telegram_nutrition_onboarding_runtime_notice import (
    TelegramNutritionOnboardingRuntimeNoticeMixin,
)
from gateway.platforms.telegram_nutrition_onboarding_runtime_errors import (
    NutritionOnboardingRuntimeError,
)
from gateway.platforms.telegram_customer_bootstrap import BootstrapState


logger = logging.getLogger(__name__)
_CALLBACK_ACK_TIMEOUT_SECONDS = 2.0


class _CallbackAnswerQuery(Protocol):
    async def answer(self, text: str | None = None) -> object: ...


class _CallbackActor(Protocol):
    @property
    def id(self) -> int | None: ...


class _OnboardingCallbackQuery(_CallbackAnswerQuery, Protocol):
    @property
    def id(self) -> str | None: ...

    @property
    def from_user(self) -> _CallbackActor | None: ...


class _OnboardingCallbackMessage(Protocol):
    @property
    def message_id(self) -> int | str | None: ...


class _BestEffortCallbackQuery:
    def __init__(
        self,
        query: _OnboardingCallbackQuery,
        *,
        update_id: int | None,
    ) -> None:
        self._query = query
        self._update_id = (
            str(update_id)
            if type(update_id) is int and update_id >= 0
            else "unknown"
        )

    @property
    def id(self) -> str | None:
        return self._query.id

    @property
    def from_user(self) -> _CallbackActor | None:
        return self._query.from_user

    async def answer(self, text: str | None = None) -> None:
        try:
            acknowledgement = (
                self._query.answer()
                if text is None
                else self._query.answer(text=text)
            )
            await asyncio.wait_for(
                acknowledgement,
                timeout=_CALLBACK_ACK_TIMEOUT_SECONDS,
            )
        except asyncio.CancelledError as exc:
            task = asyncio.current_task()
            if task is not None and task.cancelling():
                raise
            logger.warning(
                "nutrition_onboarding_callback_ack_failed "
                "gate=onboarding_callback update_id=%s error=%s",
                self._update_id,
                type(exc).__name__,
            )
        except Exception as exc:
            logger.warning(
                "nutrition_onboarding_callback_ack_failed "
                "gate=onboarding_callback update_id=%s error=%s",
                self._update_id,
                type(exc).__name__,
            )


def callback_update_id(query: _OnboardingCallbackQuery) -> int:
    callback_id = str(getattr(query, "id", "") or "")
    if not callback_id:
        raise NutritionOnboardingRuntimeError(
            "callback query id is required"
        )
    digest = hashlib.sha256(callback_id.encode()).digest()
    return int.from_bytes(digest[:8], "big") & ((1 << 63) - 1)


class TelegramNutritionOnboardingRuntimeCallbackMixin(
    TelegramNutritionOnboardingRuntimeNoticeMixin
):
    async def handle_callback(
        self: Any,
        query: _OnboardingCallbackQuery,
        data: str,
        message: _OnboardingCallbackMessage,
        *,
        update_id: int | None = None,
    ) -> None:
        raw_query = query
        query = _BestEffortCallbackQuery(query, update_id=update_id)
        log_ingress_stage("ingress", update_id, logger=logger)
        await query.answer()
        try:
            token = decode_callback(data)
            sid_hash = base64.urlsafe_b64decode(
                token.session_hash + "="
            ).hex()
            session = self.bootstrap.store.get_by_sid_hash(sid_hash)
        except (ValueError, TypeError):
            log_ingress_stage(
                "validation",
                update_id,
                reason_code="callback_invalid",
                logger=logger,
            )
            await query.answer(
                text="만료되거나 잘못된 온보딩 버튼입니다."
            )
            return
        if (
            session is None
            or session.state is not BootstrapState.AWAITING_ACTIVATION
        ):
            log_ingress_stage(
                "validation",
                update_id,
                reason_code="onboarding_state_unavailable",
                logger=logger,
            )
            await query.answer(text="현재 사용할 수 없는 온보딩 버튼입니다.")
            return
        service = self._service(session.customer_key)
        try:
            publication = service.store.load_session(session.session_id)
        except ValueError:
            log_ingress_stage(
                "validation",
                update_id,
                reason_code="onboarding_session_unavailable",
                logger=logger,
            )
            await query.answer(text="현재 온보딩 질문을 확인하지 못했습니다.")
            return
        actor_id = getattr(getattr(query, "from_user", None), "id", None)
        route_roles = {
            "prev": "customer",
            "later": "customer",
            "skip": "customer",
            "resume": "customer",
            "revise": "customer",
            "attest": "customer",
            "owner_ok": "owner",
            "op_rev": "owner",
            "op_rej": "owner",
            "op_hold": "owner",
            "hold_ack": "owner",
            "owner_risk": "owner",
        }
        role = route_roles.get(token.action)
        if (
            token.generation != publication.generation
            or str(getattr(message, "message_id", ""))
            != str(publication.message_id)
            or role is None
            or not self._publication_callback_is_current(
                session=session,
                publication=publication,
                role=role,
                message=message,
            )
        ):
            log_ingress_stage(
                "validation",
                update_id,
                reason_code="stale_publication",
                logger=logger,
            )
            await query.answer(
                text="가장 최근 온보딩 카드를 사용해 주세요."
            )
            return
        try:
            member_chat_id, _ = self._route(session, role)
        except (OSError, TypeError, ValueError):
            log_ingress_stage(
                "authority",
                update_id,
                reason_code="route_unavailable",
                logger=logger,
            )
            await query.answer(
                text="현재 참여자 권한을 확인하지 못했습니다."
            )
            return
        logger.info(
            "nutrition_onboarding_callback_stage stage=member_check_start"
        )
        if not await self._member_present(member_chat_id, actor_id):
            log_ingress_stage(
                "authority",
                update_id,
                reason_code="member_unavailable",
                logger=logger,
            )
            await query.answer(
                text="현재 참여자 권한을 확인하지 못했습니다."
            )
            return
        logger.info(
            "nutrition_onboarding_callback_stage stage=member_check_pass"
        )
        try:
            evidence_update_id = (
                update_id
                if type(update_id) is int and update_id >= 0
                else callback_update_id(query)
            )
        except ValueError:
            log_ingress_stage(
                "validation",
                update_id,
                reason_code="callback_identifier_unavailable",
                logger=logger,
            )
            await query.answer(
                text="만료되었거나 사용할 수 없는 버튼입니다."
            )
            return
        evidence = self._evidence(
            actor_id=actor_id,
            message=message,
            update_id=evidence_update_id,
        )
        try:
            authority = self._current_authority(session)
        except (OSError, ValueError):
            log_ingress_stage(
                "authority",
                update_id,
                reason_code="authority_unavailable",
                logger=logger,
            )
            await query.answer(
                text="온보딩 권한 또는 동의가 변경됐습니다."
            )
            return
        log_ingress_stage("authority", update_id, logger=logger)
        logger.info(
            "nutrition_onboarding_callback_stage stage=authority_pass"
        )
        status = service.status()
        try:
            if role is not None:
                try:
                    self.domain.validate_route(
                        authority,
                        evidence,
                        role=role,
                    )
                except ValueError:
                    log_ingress_stage(
                        "route",
                        update_id,
                        reason_code="route_rejected",
                        logger=logger,
                    )
                    await query.answer(
                        text="온보딩 권한 또는 상태가 변경됐습니다."
                    )
                    return
                logger.info(
                    "nutrition_onboarding_callback_stage stage=route_pass"
                )
            log_ingress_stage("route", update_id, logger=logger)
            replayed_status = self.replayed_status_if_consumed(
                service,
                evidence_update_id,
            )
            if replayed_status is not None:
                log_ingress_stage(
                    "persistence",
                    update_id,
                    reason_code="duplicate_business_replay",
                    logger=logger,
                )
                await self._retire_keyboard(raw_query)
                log_ingress_stage("publication", update_id, logger=logger)
                try:
                    await self._publish(session, service, replayed_status)
                except (OSError, RuntimeError, TypeError, ValueError):
                    log_ingress_stage(
                        "publication",
                        update_id,
                        reason_code="publication_failed",
                        logger=logger,
                    )
                    raise
                log_ingress_stage("receipt", update_id, logger=logger)
                return
            log_ingress_stage("validation", update_id, logger=logger)
            if (
                token.action == "prev"
                and status.state.value == "collecting"
                and status.answer_count > 0
            ):
                status = self._rewind_collection(
                    service=service,
                    authority=authority,
                    evidence=evidence,
                    target_cursor=status.answer_count - 1,
                )
            elif (
                token.action == "later"
                and status.state.value == "collecting"
            ):
                log_ingress_stage("persistence", update_id, logger=logger)
                await self._retire_keyboard(raw_query)
                await query.answer(text="현재 위치를 저장했습니다.")
                log_ingress_stage("publication", update_id, logger=logger)
                await self._continue_later(session, service, status)
                log_ingress_stage("receipt", update_id, logger=logger)
                return
            elif (
                token.action == "resume"
                and status.state.value == "collecting"
                and publication.payload.get("mode") == "paused"
            ):
                log_ingress_stage("persistence", update_id, logger=logger)
                await self._retire_keyboard(raw_query)
                await query.answer(text="온보딩을 이어서 진행합니다.")
                log_ingress_stage("publication", update_id, logger=logger)
                await self._resume(session, service)
                log_ingress_stage("receipt", update_id, logger=logger)
                return
            elif (
                token.action == "skip"
                and status.state.value == "collecting"
                and status.next_field in OPTIONAL_QUESTION_DEFAULTS
            ):
                status = service.submit_answer(
                    field=status.next_field,
                    value=OPTIONAL_QUESTION_DEFAULTS[status.next_field],
                    authority=authority,
                    evidence=evidence,
                )
            elif (
                token.action == "revise"
                and status.state.value == "customer_attestation"
                and status.answer_count > 0
            ):
                status = self._rewind_collection(
                    service=service,
                    authority=authority,
                    evidence=evidence,
                    target_cursor=status.answer_count - 1,
                )
            elif (
                token.action == "attest"
                and status.state.value == "customer_attestation"
            ):
                status = service.attest_baseline(
                    authority=authority,
                    evidence=evidence,
                )
            elif (
                status.state.value == "safety_hold"
                and token.action == "hold_ack"
            ):
                self._acknowledge_operator_attention(
                    session=session,
                    publication=publication,
                    authority=authority,
                    evidence=evidence,
                    callback_data=data,
                )
                log_ingress_stage(
                    "persistence",
                    update_id,
                    reason_code="operator_attention_acknowledged",
                    logger=logger,
                )
                await self._retire_keyboard(raw_query)
                await query.answer(
                    text="확인했습니다. 안전 보류 상태는 그대로 유지됩니다."
                )
                log_ingress_stage("receipt", update_id, logger=logger)
                return
            elif (
                status.state.value == "safety_hold"
                and token.action == "owner_risk"
            ):
                capability = getattr(
                    self,
                    "_owner_risk_acceptance",
                    None,
                )
                payload = getattr(publication, "payload", None)
                acceptance_identity = (
                    payload.get("owner_risk_acceptance_identity")
                    if isinstance(payload, dict)
                    else None
                )
                if (
                    capability is None
                    or not isinstance(payload, dict)
                    or payload.get("operator_delivery_epoch")
                    != capability.candidate_epoch
                    or not isinstance(acceptance_identity, str)
                ):
                    raise ValueError(
                        "owner risk acceptance is unavailable"
                    )
                binding = self.domain.OwnerRiskAcceptanceBinding(
                    candidate_epoch=capability.candidate_epoch,
                    bootstrap_session_id=str(session.session_id),
                    customer_key=str(session.customer_key),
                    baseline_digest=str(
                        payload.get("owner_risk_baseline_digest")
                    ),
                    reconciliation_digest=str(
                        payload.get(
                            "owner_risk_reconciliation_digest"
                        )
                    ),
                )
                status = service.record_owner_risk_acceptance(
                    request=self.domain.OwnerRiskAcceptanceRequest(
                        binding=binding,
                        acceptance_identity=acceptance_identity,
                        authority=authority,
                        evidence=evidence,
                    ),
                )
            elif status.state.value == "owner_review" and token.action in {
                "owner_ok",
                "op_rev",
                "op_rej",
                "op_hold",
            }:
                decision = {
                    "owner_ok": "approved",
                    "op_rev": "revise",
                    "op_rej": "rejected",
                    "op_hold": "safety_hold",
                }[token.action]
                kwargs = {
                    "decision": decision,
                    "authority": authority,
                    "evidence": evidence,
                }
                if decision == "revise":
                    payload = getattr(publication, "payload", None)
                    binding = (
                        payload.get("revision_binding")
                        if isinstance(payload, dict)
                        else None
                    )
                    if not isinstance(binding, dict) or set(binding) != {
                        "issue_id",
                        "field",
                        "answers_digest",
                        "reconciliation_digest",
                    }:
                        raise ValueError("revision binding is unavailable")
                    field = binding.get("field")
                    issue_id = binding.get("issue_id")
                    answers_digest = binding.get("answers_digest")
                    reconciliation_digest = binding.get("reconciliation_digest")
                    if (
                        field not in self.domain.QUESTION_FIELDS
                        or not all(
                            isinstance(value, str)
                            and len(value) == 64
                            and not set(value).difference("0123456789abcdef")
                            for value in (
                                issue_id,
                                answers_digest,
                                reconciliation_digest,
                            )
                        )
                    ):
                        raise ValueError("revision binding is invalid")
                    kwargs.update(
                        revision_issue_id=issue_id,
                        revision_field=field,
                        expected_answers_digest=answers_digest,
                        expected_reconciliation_digest=reconciliation_digest,
                    )
                if decision == "approved":
                    service.preflight_owner_approval(
                        starts_on=date.fromisoformat(
                            session.customer_draft.starts_on,
                        ),
                    )
                    self._preserve_owner_callback_receipt(
                        session=session,
                        service=service,
                        publication=publication,
                        authority=authority,
                        evidence=evidence,
                        callback_data=data,
                    )
                status = service.review_as_owner(**kwargs)
                if decision == "approved":
                    log_ingress_stage(
                        "persistence",
                        update_id,
                        reason_code="business_state_committed",
                        logger=logger,
                    )
                    status = self._finalize_owner_approval(
                        session=session,
                        service=service,
                        authority=authority,
                    )
            else:
                log_ingress_stage(
                    "validation",
                    update_id,
                    reason_code="action_state_invalid",
                    logger=logger,
                )
                await query.answer(
                    text="현재 단계에서 사용할 수 없는 버튼입니다."
                )
                return
        except self.domain.UnsafeGoalTrajectoryError:
            log_ingress_stage(
                "persistence",
                update_id,
                reason_code="unsafe_goal_trajectory",
                logger=logger,
            )
            await query.answer(
                text="목표 체중 또는 목표 날짜를 조정해야 합니다.",
            )
            return
        except ValueError:
            log_ingress_stage(
                "persistence",
                update_id,
                reason_code="transition_rejected",
                logger=logger,
            )
            await query.answer(
                text="온보딩 권한 또는 상태가 변경됐습니다."
            )
            return
        log_ingress_stage("persistence", update_id, logger=logger)
        logger.info(
            "nutrition_onboarding_callback_stage stage=transition_pass"
        )
        logger.info("nutrition_onboarding_callback_applied")
        await self._retire_keyboard(raw_query)
        await query.answer(text="확인했습니다.")
        log_ingress_stage("publication", update_id, logger=logger)
        try:
            if token.action == "attest":
                await self._send_customer_confirmation(session, status)
            lifecycle = {
                "op_rev": "revision",
                "op_rej": "rejection",
            }.get(token.action)
            if lifecycle is not None:
                await self._send_customer_lifecycle_notice(session, lifecycle)
            await self._publish(session, service, status)
        except (OSError, RuntimeError, TypeError, ValueError):
            log_ingress_stage(
                "publication",
                update_id,
                reason_code="publication_failed",
                logger=logger,
            )
            raise
        log_ingress_stage("receipt", update_id, logger=logger)
