from __future__ import annotations

import asyncio
import logging
from typing import Any

from gateway.platforms.telegram_nutrition_onboarding import (
    OPTIONAL_QUESTION_DEFAULTS,
    log_ingress_stage,
)
from gateway.platforms.telegram_nutrition_onboarding_copy import (
    answer_handoff,
    answer_error,
    parse_answer,
)
from gateway.platforms.telegram_nutrition_onboarding_customer_v2 import (
    CustomerV2RequiredAnswerUnknown,
    customer_v2_unknown_help,
)
from gateway.platforms.telegram_customer_bootstrap import BootstrapState


logger = logging.getLogger(__name__)


def is_consent_withdrawal(text: str | None) -> bool:
    return (
        isinstance(text, str)
        and " ".join(text.split()) == "서비스 중단 및 동의 철회 요청"
    )


class TelegramNutritionOnboardingRuntimeCollectionMixin:
    def _service(self, customer_key: str) -> Any:
        raise NotImplementedError

    def _session_for_customer_message(self, message: Any) -> Any | None:
        raise NotImplementedError

    def _current_authority(self, session: Any) -> Any:
        raise NotImplementedError

    def _evidence(
        self,
        *,
        actor_id: int | str | None,
        message: Any,
        update_id: int | str | None,
    ) -> Any:
        raise NotImplementedError

    async def _publish(
        self,
        session: Any,
        service: Any,
        status: Any,
        *,
        reply_anchor_message_id: int | None = None,
    ) -> None:
        raise NotImplementedError

    async def _record_answer_failure(
        self,
        *,
        session: Any,
        service: Any,
        field: str,
        authority: Any,
        evidence: Any,
        message: Any,
        questionnaire_version: str | None,
        force_handoff: bool,
    ) -> None:
        status = service.record_answer_failure(
            field=field,
            authority=authority,
            evidence=evidence,
            force_handoff=force_handoff,
        )
        if status.state.value == "input_help":
            text = (
                customer_v2_unknown_help(field)
                if force_handoff
                else answer_handoff(
                    field,
                    questionnaire_version=questionnaire_version,
                )
            )
            await message.reply_text(text)
            await self._publish(session, service, status)
            return
        await message.reply_text(
            answer_error(
                field,
                questionnaire_version=questionnaire_version,
            )
        )

    async def handle_text(self, update: Any, message: Any) -> bool:
        lock = getattr(self, "_onboarding_mutation_lock", None)
        if lock is None:
            lock = asyncio.Lock()
            self._onboarding_mutation_lock = lock
        async with lock:
            handler = getattr(self, "_handle_text_locked", None)
            if callable(handler):
                return await handler(update, message)
            return await TelegramNutritionOnboardingRuntimeCollectionMixin._handle_text_locked(
                self,
                update,
                message,
            )

    async def _handle_text_locked(self, update: Any, message: Any) -> bool:
        if is_consent_withdrawal(getattr(message, "text", None)):
            return False
        session = self._session_for_customer_message(message)
        if session is None:
            return False
        update_id = getattr(update, "update_id", None)
        log_ingress_stage("ingress", update_id, logger=logger)
        if session.state is not BootstrapState.AWAITING_ACTIVATION:
            log_ingress_stage(
                "authority",
                update_id,
                reason_code="onboarding_state_unavailable",
                logger=logger,
            )
            return False
        service = self._service(session.customer_key)
        try:
            authority = self._current_authority(session)
        except (OSError, ValueError):
            log_ingress_stage(
                "authority",
                update_id,
                reason_code="authority_unavailable",
                logger=logger,
            )
            await message.reply_text(
                "온보딩 권한 또는 동의가 변경됐습니다."
            )
            return True
        log_ingress_stage("authority", update_id, logger=logger)
        log_ingress_stage("route", update_id, logger=logger)
        status = service.status()
        try:
            publication = service.store.load_session(session.session_id)
        except ValueError:
            log_ingress_stage(
                "validation",
                update_id,
                reason_code="onboarding_session_unavailable",
                logger=logger,
            )
            return True
        reply = getattr(message, "reply_to_message", None)
        actor_id = getattr(
            getattr(message, "from_user", None),
            "id",
            None,
        )
        reply_message_id = getattr(reply, "message_id", None)
        if reply is not None and str(
            getattr(reply, "message_id", "")
        ) != str(publication.message_id):
            log_ingress_stage(
                "validation",
                update_id,
                reason_code="stale_reply",
                logger=logger,
            )
            return True
        if reply is None and status.next_field is None:
            log_ingress_stage(
                "validation",
                update_id,
                reason_code="no_active_question",
                logger=logger,
            )
            return True
        evidence = self._evidence(
            actor_id=getattr(
                getattr(message, "from_user", None),
                "id",
                None,
            ),
            message=message,
            update_id=getattr(update, "update_id", None),
        )
        replayed_status_reader = getattr(
            self,
            "replayed_status_if_consumed",
            None,
        )
        replayed_status = (
            replayed_status_reader(service, update_id)
            if callable(replayed_status_reader)
            else None
        )
        if replayed_status is not None:
            log_ingress_stage(
                "persistence",
                update_id,
                reason_code="duplicate_business_replay",
                logger=logger,
            )
            log_ingress_stage("publication", update_id, logger=logger)
            try:
                await self._publish(
                    session,
                    service,
                    replayed_status,
                    reply_anchor_message_id=int(message.message_id),
                )
            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 True
        log_ingress_stage("validation", update_id, logger=logger)
        if status.next_field is None:
            payload = publication.payload
            field = payload.get("clarification_field")
            answers_digest = payload.get("answers_digest")
            binding = payload.get("revision_binding")
            if (
                status.state.value != "customer_attestation"
                or not isinstance(field, str)
                or not isinstance(answers_digest, str)
                or not isinstance(binding, dict)
                or binding.get("field") != field
                or binding.get("answers_digest") != answers_digest
                or not isinstance(binding.get("issue_id"), str)
                or not isinstance(binding.get("reconciliation_digest"), str)
            ):
                log_ingress_stage(
                    "validation",
                    update_id,
                    reason_code="reconciliation_state_invalid",
                    logger=logger,
                )
                return True
            try:
                value = parse_answer(
                    field,
                    str(message.text),
                    questionnaire_version=status.questionnaire_version,
                )
            except CustomerV2RequiredAnswerUnknown:
                log_ingress_stage(
                    "validation",
                    update_id,
                    reason_code="answer_unknown_required",
                    logger=logger,
                )
                await self._record_answer_failure(
                    session=session,
                    service=service,
                    field=field,
                    authority=authority,
                    evidence=evidence,
                    message=message,
                    questionnaire_version=status.questionnaire_version,
                    force_handoff=True,
                )
                return True
            except ValueError:
                log_ingress_stage(
                    "validation",
                    update_id,
                    reason_code="answer_invalid",
                    logger=logger,
                )
                await self._record_answer_failure(
                    session=session,
                    service=service,
                    field=field,
                    authority=authority,
                    evidence=evidence,
                    message=message,
                    questionnaire_version=status.questionnaire_version,
                    force_handoff=False,
                )
                return True
            try:
                next_status = service.revise_reconciliation_answer(
                    field=field,
                    value=value,
                    expected_issue_id=binding["issue_id"],
                    expected_answers_digest=answers_digest,
                    expected_reconciliation_digest=binding["reconciliation_digest"],
                    authority=authority,
                    evidence=evidence,
                )
            except (KeyError, ValueError):
                log_ingress_stage(
                    "persistence",
                    update_id,
                    reason_code="transition_rejected",
                    logger=logger,
                )
                await message.reply_text(
                    "온보딩 권한 또는 상태가 변경됐습니다."
                )
                return True
            log_ingress_stage("persistence", update_id, logger=logger)
            if next_status.state.value == "input_help":
                await message.reply_text(
                    answer_handoff(
                        field,
                        questionnaire_version=status.questionnaire_version,
                    )
                )
            log_ingress_stage("publication", update_id, logger=logger)
            try:
                await self._publish(
                    session,
                    service,
                    next_status,
                    reply_anchor_message_id=int(message.message_id),
                )
            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 True
        try:
            raw_value = str(message.text).strip()
            if (
                raw_value == "건너뛰기"
                and status.next_field in OPTIONAL_QUESTION_DEFAULTS
            ):
                value = OPTIONAL_QUESTION_DEFAULTS[status.next_field]
            else:
                value = parse_answer(
                    status.next_field,
                    raw_value,
                    questionnaire_version=status.questionnaire_version,
                )
        except CustomerV2RequiredAnswerUnknown:
            log_ingress_stage(
                "validation",
                update_id,
                reason_code="answer_unknown_required",
                logger=logger,
            )
            await self._record_answer_failure(
                session=session,
                service=service,
                field=status.next_field,
                authority=authority,
                evidence=evidence,
                message=message,
                questionnaire_version=status.questionnaire_version,
                force_handoff=True,
            )
            return True
        except ValueError:
            log_ingress_stage(
                "validation",
                update_id,
                reason_code="answer_invalid",
                logger=logger,
            )
            await self._record_answer_failure(
                session=session,
                service=service,
                field=status.next_field,
                authority=authority,
                evidence=evidence,
                message=message,
                questionnaire_version=status.questionnaire_version,
                force_handoff=False,
            )
            return True
        try:
            next_status = service.submit_answer(
                field=status.next_field,
                value=value,
                authority=authority,
                evidence=evidence,
            )
        except ValueError:
            log_ingress_stage(
                "persistence",
                update_id,
                reason_code="transition_rejected",
                logger=logger,
            )
            await message.reply_text("온보딩 권한 또는 상태가 변경됐습니다.")
            return True
        log_ingress_stage("persistence", update_id, logger=logger)
        log_ingress_stage("publication", update_id, logger=logger)
        try:
            await self._publish(
                session,
                service,
                next_status,
                reply_anchor_message_id=int(message.message_id),
            )
        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 True
