from __future__ import annotations

import hashlib
import json
import logging
from datetime import date
from typing import Any

from checkin_cli.nutrition_onboarding_clarification_policy import (
    compile_clarification_policy,
)
from checkin_cli.nutrition_onboarding_contract import QUESTION_FIELDS

from gateway.platforms.nutrition_onboarding_reconciliation import (
    ReconciliationError,
    ReconciliationUnavailable,
    NutritionOnboardingClarification,
    reconciliation_payload,
    render_authoritative_customer_summary,
)
from gateway.platforms.nutrition_onboarding_reconciliation_copy import (
    render_clarification_text,
)
from gateway.platforms.telegram_nutrition_onboarding import (
    OPTIONAL_QUESTION_DEFAULTS,
)
from gateway.platforms.telegram_nutrition_onboarding_copy import card
from gateway.platforms.telegram_nutrition_onboarding_runtime_publication_transport import (
    TelegramNutritionOnboardingRuntimePublicationTransportMixin,
    publication_generation,
)


logger = logging.getLogger(__name__)


class TelegramNutritionOnboardingRuntimePublicationMixin(
    TelegramNutritionOnboardingRuntimePublicationTransportMixin
):
    domain: Any
    reconciler: Any = None
    adapter: Any
    bootstrap: Any
    _operator_notification_recovery: Any = None
    _owner_risk_acceptance: Any = None
    _operator_attention_store: Any = None

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

    def _route(self, session: Any, role: str) -> tuple[str, str]:
        raise NotImplementedError

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

    def _operator_attention_payload(
        self,
        *,
        session: Any,
        service: Any,
        status: Any,
        payload: dict[str, Any],
        recovery_request_id: str | None = None,
    ) -> dict[str, Any]:
        notification = getattr(
            self,
            "_operator_notification_recovery",
            None,
        )
        owner_risk = getattr(self, "_owner_risk_acceptance", None)
        notification_epoch = getattr(
            notification,
            "candidate_epoch",
            None,
        )
        owner_risk_epoch = getattr(
            owner_risk,
            "candidate_epoch",
            None,
        )
        candidate_epochs = {
            value
            for value in (notification_epoch, owner_risk_epoch)
            if isinstance(value, str)
            and len(value) == 64
            and not set(value) - set("0123456789abcdef")
        }
        if (
            status.state.value not in {"owner_review", "safety_hold"}
            or not candidate_epochs
        ):
            return payload
        if len(candidate_epochs) != 1:
            raise ValueError("operator capability candidate mismatch")
        candidate_epoch = next(iter(candidate_epochs))
        authority = self._current_authority(session)
        record = service.reconciliation_record(authority=authority)
        reconciliation_digest = (
            record.get("digest") if isinstance(record, dict) else None
        )
        if (
            not isinstance(reconciliation_digest, str)
            or len(reconciliation_digest) != 64
            or set(reconciliation_digest) - set("0123456789abcdef")
        ):
            raise ValueError("operator attention authority is unavailable")
        decorated = dict(payload)
        if notification_epoch == candidate_epoch:
            identity = hashlib.sha256(
                json.dumps(
                    {
                        "body_digest": payload.get("body_digest"),
                        "reconciliation_digest": (
                            reconciliation_digest
                        ),
                        "session_id": str(session.session_id),
                        "state": status.state.value,
                    },
                    sort_keys=True,
                    separators=(",", ":"),
                ).encode()
            ).hexdigest()
            decorated.update(
                operator_attention_identity=identity,
                operator_delivery_epoch=candidate_epoch,
            )
        if (
            status.state.value == "safety_hold"
            and owner_risk_epoch == candidate_epoch
            and isinstance(status.baseline_digest, str)
        ):
            binding = self.domain.OwnerRiskAcceptanceBinding(
                candidate_epoch=candidate_epoch,
                bootstrap_session_id=str(session.session_id),
                customer_key=str(session.customer_key),
                baseline_digest=status.baseline_digest,
                reconciliation_digest=reconciliation_digest,
            )
            decorated["owner_risk_acceptance_identity"] = (
                self.domain.owner_risk_acceptance_identity(binding)
            )
            decorated["owner_risk_baseline_digest"] = (
                status.baseline_digest
            )
            decorated["owner_risk_reconciliation_digest"] = (
                reconciliation_digest
            )
            decorated["operator_delivery_epoch"] = candidate_epoch
        if recovery_request_id is not None:
            decorated["operator_recovery_request"] = hashlib.sha256(
                (
                    f"{session.session_id}\0{recovery_request_id}"
                ).encode()
            ).hexdigest()
        return decorated

    async def recover_operator_status(
        self,
        *,
        actor_user_id: int,
        chat_id: int,
        topic_id: int,
        request_id: int,
    ) -> int | None:
        capability = getattr(self, "_operator_notification_recovery", None)
        if (
            type(actor_user_id) is not int
            or type(chat_id) is not int
            or type(topic_id) is not int
            or type(request_id) is not int
            or request_id < 0
            or actor_user_id
            != getattr(capability, "owner_user_id", None)
            or chat_id != getattr(capability, "owner_chat_id", None)
            or topic_id != getattr(capability, "owner_topic_id", None)
        ):
            return None
        recovered = 0
        for session in self.bootstrap.store.list_sessions():
            state = getattr(getattr(session, "state", None), "value", None)
            if state != "AWAITING_ACTIVATION":
                continue
            service = self._service(session.customer_key)
            status = service.status()
            if status.state.value not in {"owner_review", "safety_hold"}:
                continue
            await self._publish(
                session,
                service,
                status,
                operator_recovery_id=str(request_id),
            )
            recovered += 1
        return recovered

    async def _publish(
        self,
        session: Any,
        service: Any,
        status: Any,
        *,
        reply_anchor_message_id: int | None = None,
        operator_recovery_id: str | None = None,
    ) -> None:
        text, action, role = card(status)
        if status.state.value == "ready":
            plan = self.domain.read_private_json(
                service.root / "initial-plan-v1.json",
            )
            if (
                plan.get(
                    "requested_trajectory_within_guardrail"
                )
                is False
                and isinstance(
                    plan.get("recommended_target_date"),
                    str,
                )
            ):
                text += (
                    "\n\n고객이 요청한 목표와 기한은 그대로 보존했습니다."
                    "\n코치 권장 기한: "
                    f"{plan['recommended_target_date']}"
                    "\n권장 기한은 더 안정적인 진행을 위한 비교안이며, "
                    "고객 목표를 대체하지 않습니다."
                )
        if status.state.value == "collecting":
            text = (
                f"[{status.answer_count + 1}/{len(QUESTION_FIELDS)}] "
                f"{text}\n\n"
                "입력창이 자동으로 연결됩니다. "
                "그대로 답을 쓰고 전송하세요."
            )
            if status.next_field in OPTIONAL_QUESTION_DEFAULTS:
                text += "\n답변이 없으면 '건너뛰기'라고 입력하세요."
        clarification = None
        revision_binding = None
        answers_digest = None
        if status.state.value == "customer_attestation":
            try:
                authority = self._current_authority(session)
                answers = service.reconciliation_answers(
                    authority=authority
                )
                policy = compile_clarification_policy(
                    answers,
                    reference_date=date.today(),
                )
                answers_digest = policy.answers_digest
                record = service.reconciliation_record(authority=authority)
                stale_answers_digest = None
                if (
                    record is not None
                    and record.get("answers_digest") != answers_digest
                ):
                    stale_answers_digest = str(record["answers_digest"])
                    record = None
                if record is None:
                    try:
                        advisory = await self.reconciler.reconcile(
                            answers,
                            consent_granted=bool(authority.consent_granted),
                        )
                        advisory_payload = reconciliation_payload(advisory)
                    except (ReconciliationError, ReconciliationUnavailable):
                        advisory_payload = {"status": "unavailable"}
                    record_payload = {
                        "answers_digest": answers_digest,
                        "advisory": advisory_payload,
                        "clarifications": [],
                        "authority": authority,
                        "reference_date": policy.reference_date,
                    }
                    if stale_answers_digest is None:
                        service.record_reconciliation(**record_payload)
                    else:
                        service.replace_stale_reconciliation(
                            expected_stale_answers_digest=stale_answers_digest,
                            **record_payload,
                        )
                    record = service.reconciliation_record(authority=authority)
                if record is not None:
                    raw_issues = record.get("issues", [])
                    issues = raw_issues if isinstance(raw_issues, list) else []
                    if record.get("state") == "clarifying" and issues:
                        issue = issues[0]
                        if not isinstance(issue, dict):
                            raise ValueError("persisted reconciliation issue is invalid")
                        clarification = NutritionOnboardingClarification(
                            field=str(issue["field"]),
                            kind="ambiguity",
                            question_ko=str(issue["question_ko"]),
                        )
                        revision_binding = {
                            "issue_id": str(issue["issue_id"]),
                            "field": clarification.field,
                            "answers_digest": answers_digest,
                            "reconciliation_digest": str(record["digest"]),
                        }
                        text = render_clarification_text(
                            clarification,
                            index=1,
                            total=min(len(issues), 3),
                        )
                        action = None
                    elif record.get("state") == "resolved":
                        text = (
                            "입력 내용을 정리했어요\n\n"
                            + render_authoritative_customer_summary(answers)
                            + "\n\n아래 버튼을 누르면 입력 확인이 저장되고 "
                            "운영자 검토로 넘어갑니다.\n"
                            "이 단계에서 영양 코칭이 시작되거나 "
                            "활성화되지는 않습니다."
                        )
            except (
                OSError,
                ValueError,
                ReconciliationError,
                ReconciliationUnavailable,
            ) as exc:
                logger.warning(
                    "Nutrition onboarding advisory summary unavailable: %s",
                    type(exc).__name__,
                )
        digest = hashlib.sha256(text.encode()).hexdigest()
        payload = {"body_digest": digest, "state": status.state.value}
        if status.state.value == "collecting":
            payload["force_reply_mode"] = (
                "anchored"
                if reply_anchor_message_id is not None
                else "broadcast"
            )
        if answers_digest is not None:
            payload["answers_digest"] = answers_digest
        if clarification is not None:
            payload["clarification_field"] = clarification.field
            payload["revision_binding"] = revision_binding
        payload = self._operator_attention_payload(
            session=session,
            service=service,
            status=status,
            payload=payload,
            recovery_request_id=operator_recovery_id,
        )
        attention_identity = payload.get("operator_attention_identity")
        attention_store = getattr(self, "_operator_attention_store", None)
        if (
            operator_recovery_id is None
            and isinstance(attention_identity, str)
            and attention_store is not None
            and attention_store.acknowledged(
                session_id=str(session.session_id),
                attention_identity=attention_identity,
            )
            and "owner_risk_acceptance_identity" not in payload
        ):
            return
        actions: list[tuple[str, str]] = []
        match status.state.value:
            case "collecting":
                pass
            case "owner_review":
                actions.extend(
                    (
                        ("owner_ok", "Approve"),
                        ("op_rej", "Reject"),
                        ("op_hold", "Safety Hold"),
                    )
                )
            case _ if action == "attest":
                actions.extend(
                    (
                        ("attest", "입력 내용이 맞습니다"),
                        ("revise", "수정"),
                    )
                )
            case "safety_hold":
                if "operator_attention_identity" in payload:
                    actions.append(
                        ("hold_ack", "확인했습니다 - 임상 승인 아님")
                    )
                if (
                    "owner_risk_acceptance_identity" in payload
                ):
                    actions.append(
                        (
                            "owner_risk",
                            "내 책임으로 비의료 코칭 계속",
                        )
                    )
            case (
                "customer_attestation"
                | "rejected"
                | "finalizing"
                | "ready"
            ):
                pass
            case unexpected:
                raise RuntimeError(f"unsupported onboarding publication state: {unexpected}")
        await self._send_publication(
            session=session,
            service=service,
            status=status,
            text=text,
            role=role,
            payload=payload,
            actions=actions,
            force_reply=(
                status.state.value == "collecting"
                or clarification is not None
            ),
            reply_anchor_message_id=reply_anchor_message_id,
        )

    async def _publish_paused(
        self,
        session: Any,
        service: Any,
        status: Any,
    ) -> None:
        text = (
            "온보딩 진행 위치를 저장했습니다.\n"
            f"현재 진행: {status.answer_count}/{len(QUESTION_FIELDS)}\n"
            "준비되면 아래 버튼을 눌러 같은 온보딩을 이어서 진행해 주세요."
        )
        await self._send_publication(
            session=session,
            service=service,
            status=status,
            text=text,
            role="customer",
            payload={
                "body_digest": hashlib.sha256(text.encode()).hexdigest(),
                "state": status.state.value,
                "mode": "paused",
            },
            actions=[("resume", "계속하기")],
            force_reply=False,
        )
