from __future__ import annotations

import hashlib
import json
from datetime import date
from pathlib import Path
from types import SimpleNamespace
from unittest.mock import AsyncMock, Mock

import anyio
import pytest

from checkin_cli import nutrition_onboarding as domain
from checkin_cli.nutrition_onboarding_calculations import (
    UnsafeGoalTrajectoryError,
)
from gateway.platforms.telegram import TelegramAdapter
from gateway.platforms.telegram_nutrition_onboarding import (
    encode_callback_hash,
)
from gateway.platforms.telegram_nutrition_onboarding_operator_ack import (
    GatewayOperatorAttentionStore,
)
from gateway.platforms.telegram_nutrition_onboarding_operator_recovery import (
    load_operator_notification_recovery,
)
from gateway.platforms.telegram_nutrition_onboarding_owner_risk import (
    load_owner_risk_acceptance,
)
from gateway.platforms.telegram_nutrition_onboarding_publication_outbox import (
    GatewayOnboardingPublicationOutbox,
)
from gateway.platforms.telegram_nutrition_onboarding_runtime_publication import (
    TelegramNutritionOnboardingRuntimePublicationMixin,
)
from gateway.platforms.telegram_nutrition_onboarding_runtime import (
    TelegramNutritionOnboardingRuntime,
)
from gateway.platforms.telegram_nutrition_onboarding_runtime_callback import (
    TelegramNutritionOnboardingRuntimeCallbackMixin,
)
from gateway.platforms.telegram_nutrition_onboarding_runtime_publication_transport import (
    TelegramNutritionOnboardingRuntimePublicationTransportMixin,
    operator_attention_payloads_equivalent,
    publication_generation,
)
from gateway.platforms.telegram_customer_bootstrap import BootstrapState


def _canonical(value: object) -> bytes:
    return json.dumps(
        value,
        sort_keys=True,
        separators=(",", ":"),
    ).encode()


def _committed_owner_publication(
    profile_root: Path,
):
    outbox = GatewayOnboardingPublicationOutbox(profile_root)
    receipt, claimed = outbox.claim(
        session_id="session-1",
        generation=26,
        payload={
            "body_digest": "a" * 64,
            "operator_attention_identity": "b" * 64,
            "operator_delivery_epoch": "c" * 64,
            "state": "safety_hold",
        },
        route=("12", "0"),
        role="owner",
        render_identity="d" * 64,
    )
    assert claimed is True
    receipt = outbox.record_receipt(
        session_id=receipt.session_id,
        generation=receipt.generation,
        chat_id=receipt.route[0],
        topic_id=receipt.route[1],
        message_id=901,
    )
    return outbox.mark_committed(
        session_id=receipt.session_id,
        generation=receipt.generation,
        payload=receipt.payload,
        route=receipt.route,
        role=receipt.role,
        render_identity=receipt.render_identity,
        message_id=901,
    )


def test_safety_hold_acknowledgement_is_append_only_and_idempotent(
    tmp_path: Path,
) -> None:
    publication = _committed_owner_publication(tmp_path)
    store = GatewayOperatorAttentionStore(tmp_path)

    first, created = store.acknowledge(
        session_id="session-1",
        customer_key="pilot-1",
        attention_identity="b" * 64,
        candidate_epoch="c" * 64,
        actor_user_id=12,
        authority=(12, 12, 0),
        route=("12", "0"),
        message_id=901,
        update_id=7001,
        callback_data="non2:hold_ack:token",
        publication=publication,
    )
    replay, replay_created = store.acknowledge(
        session_id="session-1",
        customer_key="pilot-1",
        attention_identity="b" * 64,
        candidate_epoch="c" * 64,
        actor_user_id=12,
        authority=(12, 12, 0),
        route=("12", "0"),
        message_id=901,
        update_id=7001,
        callback_data="non2:hold_ack:token",
        publication=publication,
    )

    assert created is True
    assert replay_created is False
    assert replay == first
    assert GatewayOperatorAttentionStore(tmp_path).acknowledged(
        session_id="session-1",
        attention_identity="b" * 64,
    )

    later, later_created = store.acknowledge(
        session_id="session-1",
        customer_key="pilot-1",
        attention_identity="b" * 64,
        candidate_epoch="c" * 64,
        actor_user_id=12,
        authority=(12, 12, 0),
        route=("12", "0"),
        message_id=901,
        update_id=7002,
        callback_data="non2:hold_ack:other",
        publication=publication,
    )
    assert later_created is False
    assert later == first


def test_runtime_acknowledges_hold_without_domain_transition(
    tmp_path: Path,
) -> None:
    publication = _committed_owner_publication(tmp_path)
    attention_store = GatewayOperatorAttentionStore(tmp_path)
    runtime = SimpleNamespace(
        _operator_notification_recovery=SimpleNamespace(
            candidate_epoch="c" * 64,
        ),
        _operator_attention_store=attention_store,
        _publication_outbox=GatewayOnboardingPublicationOutbox(tmp_path),
        _route=lambda _session, _role: ("12", "0"),
    )
    session = SimpleNamespace(session_id="session-1", customer_key="pilot-1")
    current = SimpleNamespace(
        generation=publication.generation,
        message_id=publication.message_id,
        payload=publication.payload,
    )
    authority = SimpleNamespace(
        owner_user_id=12,
        owner_chat_id=12,
        owner_topic_id=0,
    )
    evidence = SimpleNamespace(
        actor_user_id=12,
        chat_id=12,
        topic_id=0,
        message_id=901,
        update_id=7001,
    )

    created = TelegramNutritionOnboardingRuntime._acknowledge_operator_attention(
        runtime,
        session=session,
        publication=current,
        authority=authority,
        evidence=evidence,
        callback_data="non2:hold_ack:token",
    )

    assert created is True
    assert attention_store.acknowledged(
        session_id="session-1",
        attention_identity="b" * 64,
    )


def test_acknowledged_hold_does_not_republish_on_successor_startup(
    tmp_path: Path,
) -> None:
    publication = _committed_owner_publication(tmp_path)
    attention_store = GatewayOperatorAttentionStore(tmp_path)
    attention_store.acknowledge(
        session_id="session-1",
        customer_key="pilot-1",
        attention_identity="b" * 64,
        candidate_epoch="c" * 64,
        actor_user_id=12,
        authority=(12, 12, 0),
        route=("12", "0"),
        message_id=901,
        update_id=7001,
        callback_data="non2:hold_ack:token",
        publication=publication,
    )
    send = AsyncMock()
    runtime = TelegramNutritionOnboardingRuntimePublicationMixin()
    runtime._operator_notification_recovery = SimpleNamespace(
        candidate_epoch="f" * 64,
    )
    runtime._operator_attention_store = attention_store
    runtime._current_authority = lambda _session: SimpleNamespace()
    runtime._operator_attention_payload = lambda **_kwargs: {
        **publication.payload,
        "operator_delivery_epoch": "f" * 64,
    }
    runtime._send_publication = send
    service = SimpleNamespace(
        reconciliation_record=lambda **_kwargs: {"digest": "e" * 64},
    )
    session = SimpleNamespace(session_id="session-1", sid_hash="a" * 64)
    status = SimpleNamespace(
        answer_count=22,
        state=SimpleNamespace(value="safety_hold"),
    )

    anyio.run(runtime._publish, session, service, status)

    send.assert_not_awaited()


def test_safety_hold_card_exposes_owner_risk_action_only_when_authorized(
    tmp_path: Path,
) -> None:
    runtime = TelegramNutritionOnboardingRuntimePublicationMixin()
    runtime._operator_notification_recovery = SimpleNamespace(
        candidate_epoch="c" * 64,
    )
    runtime._owner_risk_acceptance = SimpleNamespace(
        candidate_epoch="c" * 64,
    )
    runtime.domain = domain
    GatewayOnboardingPublicationOutbox(tmp_path)
    runtime._operator_attention_store = GatewayOperatorAttentionStore(
        tmp_path,
    )
    runtime._current_authority = lambda _session: SimpleNamespace()
    runtime._send_publication = AsyncMock()
    service = SimpleNamespace(
        reconciliation_record=lambda **_kwargs: {"digest": "e" * 64},
    )
    session = SimpleNamespace(
        session_id="session-1",
        customer_key="pilot-1",
        sid_hash="a" * 64,
    )
    status = SimpleNamespace(
        answer_count=22,
        baseline_digest="f" * 64,
        state=SimpleNamespace(value="safety_hold"),
    )

    anyio.run(runtime._publish, session, service, status)

    actions = runtime._send_publication.await_args.kwargs["actions"]
    assert [action for action, _label in actions] == [
        "hold_ack",
        "owner_risk",
    ]

    runtime._owner_risk_acceptance = None
    anyio.run(runtime._publish, session, service, status)
    disabled_actions = runtime._send_publication.await_args.kwargs["actions"]
    assert [action for action, _label in disabled_actions] == ["hold_ack"]

    runtime._operator_notification_recovery = None
    runtime._owner_risk_acceptance = SimpleNamespace(
        candidate_epoch="c" * 64,
    )
    anyio.run(runtime._publish, session, service, status)
    owner_risk_only_actions = (
        runtime._send_publication.await_args.kwargs["actions"]
    )
    assert [
        action for action, _label in owner_risk_only_actions
    ] == ["owner_risk"]


def test_enabling_owner_risk_republishes_acknowledged_hold_card(
    tmp_path: Path,
) -> None:
    publication = _committed_owner_publication(tmp_path)
    attention_store = GatewayOperatorAttentionStore(tmp_path)
    attention_store.acknowledge(
        session_id="session-1",
        customer_key="pilot-1",
        attention_identity="b" * 64,
        candidate_epoch="c" * 64,
        actor_user_id=12,
        authority=(12, 12, 0),
        route=("12", "0"),
        message_id=901,
        update_id=7001,
        callback_data="non2:hold_ack:token",
        publication=publication,
    )
    runtime = TelegramNutritionOnboardingRuntimePublicationMixin()
    runtime._operator_notification_recovery = SimpleNamespace(
        candidate_epoch="c" * 64,
    )
    runtime._owner_risk_acceptance = SimpleNamespace(
        candidate_epoch="c" * 64,
    )
    runtime._operator_attention_store = attention_store
    runtime.domain = domain
    runtime._current_authority = lambda _session: SimpleNamespace()
    runtime._send_publication = AsyncMock()
    service = SimpleNamespace(
        reconciliation_record=lambda **_kwargs: {"digest": "e" * 64},
    )
    session = SimpleNamespace(
        session_id="session-1",
        customer_key="pilot-1",
        sid_hash="a" * 64,
    )
    status = SimpleNamespace(
        answer_count=22,
        baseline_digest="f" * 64,
        state=SimpleNamespace(value="safety_hold"),
    )

    anyio.run(runtime._publish, session, service, status)

    runtime._send_publication.assert_awaited_once()
    actions = runtime._send_publication.await_args.kwargs["actions"]
    assert [action for action, _label in actions] == [
        "hold_ack",
        "owner_risk",
    ]


def test_owner_risk_capability_requires_exact_candidate_owner_and_seal(
    tmp_path: Path,
    monkeypatch: pytest.MonkeyPatch,
) -> None:
    candidate = "c" * 64
    core = {
        "capability": "owner_risk_acceptance",
        "candidate_digest": candidate,
        "owner": {"chat_id": 12, "topic_id": 0, "user_id": 12},
        "schema": "owner-risk-acceptance-authority-v1",
        "status": "AUTHORIZED",
    }
    document = {
        **core,
        "authorization_seal_sha256": hashlib.sha256(
            _canonical(core),
        ).hexdigest(),
    }
    authority_root = tmp_path / "authority"
    authority_root.mkdir()
    receipt = authority_root / "owner-risk-acceptance.json"
    receipt.write_bytes(_canonical(document) + b"\n")
    receipt.chmod(0o600)
    pin = authority_root / "task26-authority-pin.json"
    pin.write_text("{}\n", encoding="utf-8")
    pin.chmod(0o600)
    monkeypatch.setenv("TASK26_AUTHORITY_PIN", str(pin))
    config = {
        "room_bootstrap": {
            "owner_risk_acceptance": {
                "enabled": True,
                "receipt_path": str(receipt),
                "receipt_sha256": hashlib.sha256(
                    receipt.read_bytes(),
                ).hexdigest(),
            },
        },
    }

    capability = load_owner_risk_acceptance(
        config,
        candidate_digest=candidate,
        owner_user_id=12,
        owner_chat_id=12,
        owner_topic_id=0,
    )

    assert capability is not None
    assert capability.candidate_epoch == candidate
    assert load_owner_risk_acceptance(
        {"room_bootstrap": {}},
        candidate_digest=candidate,
        owner_user_id=12,
        owner_chat_id=12,
        owner_topic_id=0,
    ) is None
    with pytest.raises(ValueError, match="owner"):
        load_owner_risk_acceptance(
            config,
            candidate_digest=candidate,
            owner_user_id=99,
            owner_chat_id=12,
            owner_topic_id=0,
        )


def test_owner_risk_callback_transitions_and_publishes_owner_review() -> None:
    sid_hash = "a" * 64
    session = SimpleNamespace(
        session_id="session-1",
        customer_key="pilot-1",
        sid_hash=sid_hash,
        state=BootstrapState.AWAITING_ACTIVATION,
    )
    binding = domain.OwnerRiskAcceptanceBinding(
        candidate_epoch="c" * 64,
        bootstrap_session_id="session-1",
        customer_key="pilot-1",
        baseline_digest="f" * 64,
        reconciliation_digest="e" * 64,
    )
    publication = SimpleNamespace(
        generation=27,
        message_id=75,
        payload={
            "operator_delivery_epoch": "c" * 64,
            "owner_risk_acceptance_identity": (
                domain.owner_risk_acceptance_identity(binding)
            ),
            "owner_risk_baseline_digest": "f" * 64,
            "owner_risk_reconciliation_digest": "e" * 64,
            "state": "safety_hold",
        },
    )
    held = SimpleNamespace(state=SimpleNamespace(value="safety_hold"))
    owner_review = SimpleNamespace(
        state=SimpleNamespace(value="owner_review"),
    )
    service = SimpleNamespace(
        store=SimpleNamespace(
            load_session=lambda _session_id: publication,
        ),
        status=lambda: held,
        record_owner_risk_acceptance=Mock(return_value=owner_review),
    )
    authority = SimpleNamespace(
        owner_user_id=12,
        owner_chat_id=12,
        owner_topic_id=0,
    )
    evidence = SimpleNamespace(
        actor_user_id=12,
        chat_id=12,
        topic_id=0,
        message_id=75,
        update_id=7001,
    )
    publish = AsyncMock()
    retire = AsyncMock()
    runtime = SimpleNamespace(
        bootstrap=SimpleNamespace(
            store=SimpleNamespace(
                get_by_sid_hash=lambda _sid_hash: session,
            ),
        ),
        _service=lambda _customer_key: service,
        _publication_callback_is_current=lambda **_kwargs: True,
        _route=lambda _session, _role: ("12", "0"),
        _member_present=AsyncMock(return_value=True),
        _evidence=lambda **_kwargs: evidence,
        _current_authority=lambda _session: authority,
        domain=domain,
        replayed_status_if_consumed=lambda *_args: None,
        _owner_risk_acceptance=SimpleNamespace(
            candidate_epoch="c" * 64,
        ),
        _retire_keyboard=retire,
        _publish=publish,
    )
    query = SimpleNamespace(
        id="callback-1",
        from_user=SimpleNamespace(id=12),
        answer=AsyncMock(),
    )
    message = SimpleNamespace(message_id=75)
    callback = encode_callback_hash(
        action="owner_risk",
        generation=27,
        sid_hash=sid_hash,
    )

    async def exercise() -> None:
        await TelegramNutritionOnboardingRuntimeCallbackMixin.handle_callback(
            runtime,
            query,
            callback,
            message,
            update_id=7001,
        )

    anyio.run(exercise)

    service.record_owner_risk_acceptance.assert_called_once_with(
        request=domain.OwnerRiskAcceptanceRequest(
            binding=binding,
            acceptance_identity=domain.owner_risk_acceptance_identity(
                binding,
            ),
            authority=authority,
            evidence=evidence,
        ),
    )
    retire.assert_awaited_once_with(query)
    publish.assert_awaited_once_with(session, service, owner_review)


def test_owner_approval_preflights_plan_before_committing_review() -> None:
    sid_hash = "a" * 64
    session = SimpleNamespace(
        session_id="session-1",
        customer_key="pilot-1",
        customer_draft=SimpleNamespace(starts_on="2026-08-21"),
        sid_hash=sid_hash,
        state=BootstrapState.AWAITING_ACTIVATION,
    )
    publication = SimpleNamespace(
        dispatch_identity="d" * 64,
        generation=29,
        message_id=77,
        payload={"state": "owner_review"},
        receipt_integrity="e" * 64,
        route=("12", "0"),
    )
    owner_review = SimpleNamespace(
        state=SimpleNamespace(value="owner_review"),
    )
    preflight = Mock(
        side_effect=UnsafeGoalTrajectoryError(
            "unsafe goal trajectory",
        ),
    )
    review = Mock(
        return_value=SimpleNamespace(
            state=SimpleNamespace(value="finalizing"),
        ),
    )
    service = SimpleNamespace(
        preflight_owner_approval=preflight,
        review_as_owner=review,
        status=lambda: owner_review,
        store=SimpleNamespace(
            load_session=lambda _session_id: publication,
        ),
    )
    authority = SimpleNamespace(
        owner_user_id=12,
        owner_chat_id=12,
        owner_topic_id=0,
    )
    evidence = SimpleNamespace(
        actor_user_id=12,
        chat_id=12,
        topic_id=0,
        message_id=77,
        update_id=7002,
    )
    preserve = Mock()
    finalize = Mock(
        side_effect=UnsafeGoalTrajectoryError(
            "unsafe goal trajectory",
        ),
    )
    runtime = SimpleNamespace(
        bootstrap=SimpleNamespace(
            store=SimpleNamespace(
                get_by_sid_hash=lambda _sid_hash: session,
            ),
        ),
        _current_authority=lambda _session: authority,
        _evidence=lambda **_kwargs: evidence,
        _finalize_owner_approval=finalize,
        _member_present=AsyncMock(return_value=True),
        _preserve_owner_callback_receipt=preserve,
        _publication_callback_is_current=lambda **_kwargs: True,
        _publish=AsyncMock(),
        _retire_keyboard=AsyncMock(),
        _route=lambda _session, _role: ("12", "0"),
        _service=lambda _customer_key: service,
        domain=domain,
        replayed_status_if_consumed=lambda *_args: None,
    )
    query = SimpleNamespace(
        id="callback-2",
        from_user=SimpleNamespace(id=12),
        answer=AsyncMock(),
    )
    callback = encode_callback_hash(
        action="owner_ok",
        generation=29,
        sid_hash=sid_hash,
    )

    async def exercise() -> None:
        await TelegramNutritionOnboardingRuntimeCallbackMixin.handle_callback(
            runtime,
            query,
            callback,
            SimpleNamespace(message_id=77),
            update_id=7002,
        )

    anyio.run(exercise)

    preflight.assert_called_once_with(
        starts_on=date(2026, 8, 21),
    )
    preserve.assert_not_called()
    review.assert_not_called()
    finalize.assert_not_called()


def test_recovery_capability_requires_exact_candidate_owner_and_seal(
    tmp_path: Path,
    monkeypatch: pytest.MonkeyPatch,
) -> None:
    candidate = "c" * 64
    core = {
        "capability": "operator_notification_recovery",
        "candidate_digest": candidate,
        "owner": {"chat_id": 12, "topic_id": 0, "user_id": 12},
        "schema": "operator-notification-recovery-authority-v1",
        "status": "AUTHORIZED",
    }
    document = {
        **core,
        "authorization_seal_sha256": hashlib.sha256(_canonical(core)).hexdigest(),
    }
    authority_root = tmp_path / "authority"
    authority_root.mkdir()
    receipt = authority_root / "operator-notification-recovery.json"
    receipt.write_bytes(_canonical(document) + b"\n")
    receipt.chmod(0o600)
    pin = authority_root / "task26-authority-pin.json"
    pin.write_text("{}\n", encoding="utf-8")
    pin.chmod(0o600)
    monkeypatch.setenv("TASK26_AUTHORITY_PIN", str(pin))
    config = {
        "room_bootstrap": {
            "operator_notification_recovery": {
                "enabled": True,
                "receipt_path": str(receipt),
                "receipt_sha256": hashlib.sha256(receipt.read_bytes()).hexdigest(),
            },
        },
    }

    loaded = load_operator_notification_recovery(
        config,
        candidate_digest=candidate,
        owner_user_id=12,
        owner_chat_id=12,
        owner_topic_id=0,
    )

    assert loaded is not None
    assert loaded.candidate_epoch == candidate
    assert load_operator_notification_recovery(
        {"room_bootstrap": {}},
        candidate_digest=candidate,
        owner_user_id=12,
        owner_chat_id=12,
        owner_topic_id=0,
    ) is None
    relative_config = {
        "room_bootstrap": {
            "operator_notification_recovery": {
                **config["room_bootstrap"]["operator_notification_recovery"],
                "receipt_path": receipt.name,
            },
        },
    }
    monkeypatch.delenv("TASK26_AUTHORITY_PIN")
    monkeypatch.setenv("CREDENTIALS_DIRECTORY", str(authority_root))
    assert load_operator_notification_recovery(
        relative_config,
        candidate_digest=candidate,
        owner_user_id=12,
        owner_chat_id=12,
        owner_topic_id=0,
    ) is not None
    monkeypatch.setenv("TASK26_AUTHORITY_PIN", str(pin))
    monkeypatch.delenv("CREDENTIALS_DIRECTORY")
    with pytest.raises(ValueError, match="owner"):
        load_operator_notification_recovery(
            config,
            candidate_digest=candidate,
            owner_user_id=99,
            owner_chat_id=12,
            owner_topic_id=0,
        )


def test_unacknowledged_hold_republishes_once_per_candidate_epoch() -> None:
    session = SimpleNamespace(session_id="session-1", sid_hash="a" * 64)
    service = SimpleNamespace(
        reconciliation_record=lambda **_kwargs: {"digest": "e" * 64},
    )
    status = SimpleNamespace(
        answer_count=22,
        state=SimpleNamespace(value="safety_hold"),
    )
    base = {"body_digest": "a" * 64, "state": "safety_hold"}

    first = TelegramNutritionOnboardingRuntimePublicationMixin()
    first._operator_notification_recovery = SimpleNamespace(
        candidate_epoch="c" * 64,
    )
    first._current_authority = lambda _session: SimpleNamespace()
    first_payload = first._operator_attention_payload(
        session=session,
        service=service,
        status=status,
        payload=base,
    )
    same_epoch_payload = first._operator_attention_payload(
        session=session,
        service=service,
        status=status,
        payload=base,
    )

    successor = TelegramNutritionOnboardingRuntimePublicationMixin()
    successor._operator_notification_recovery = SimpleNamespace(
        candidate_epoch="f" * 64,
    )
    successor._current_authority = lambda _session: SimpleNamespace()
    successor_payload = successor._operator_attention_payload(
        session=session,
        service=service,
        status=status,
        payload=base,
    )

    assert first_payload == same_epoch_payload
    assert first_payload["operator_attention_identity"] == successor_payload[
        "operator_attention_identity"
    ]
    assert first_payload["operator_delivery_epoch"] == "c" * 64
    assert successor_payload["operator_delivery_epoch"] == "f" * 64
    current = SimpleNamespace(generation=26, payload=first_payload)
    assert publication_generation(
        status,
        current=current,
        payload=same_epoch_payload,
    ) == 26
    assert publication_generation(
        status,
        current=current,
        payload=successor_payload,
    ) == 27


def test_startup_recovers_hold_when_bootstrap_consent_handoff_is_missing() -> None:
    session = SimpleNamespace(
        session_id="session-1",
        customer_key="pilot-1",
        state=BootstrapState.AWAITING_ACTIVATION,
        consent_handoff=None,
    )
    status = SimpleNamespace(
        answer_count=22,
        state=SimpleNamespace(value="safety_hold"),
    )
    service = SimpleNamespace(status=lambda: status)
    publish = AsyncMock()
    recover_receipts = AsyncMock()
    runtime = SimpleNamespace(
        _consent_recovery_lock=anyio.Lock(),
        _operator_notification_recovery=SimpleNamespace(
            candidate_epoch="c" * 64,
        ),
        bootstrap=SimpleNamespace(
            store=SimpleNamespace(get=lambda _session_id: session),
        ),
        _current_authority=lambda _session: SimpleNamespace(),
        _service=lambda _customer_key: service,
        _recover_publication_receipts=recover_receipts,
        _publish=publish,
    )

    result = anyio.run(
        TelegramNutritionOnboardingRuntime.recover_waiting_session,
        runtime,
        session,
    )

    assert result is True
    recover_receipts.assert_awaited_once_with(service, session=session)
    publish.assert_awaited_once_with(session, service, status)


def test_status_recovery_then_same_candidate_restart_is_deduplicated() -> None:
    attention = {
        "body_digest": "a" * 64,
        "operator_attention_identity": "b" * 64,
        "operator_delivery_epoch": "c" * 64,
        "state": "safety_hold",
    }
    status_payload = {
        **attention,
        "operator_recovery_request": "d" * 64,
    }
    successor_payload = {
        **attention,
        "operator_delivery_epoch": "f" * 64,
    }

    assert operator_attention_payloads_equivalent(
        status_payload,
        attention,
    )
    assert not operator_attention_payloads_equivalent(
        status_payload,
        successor_payload,
    )

    transport = TelegramNutritionOnboardingRuntimePublicationTransportMixin()
    transport._current_publication = lambda _service, _session_id: SimpleNamespace(
        payload=status_payload,
    )
    transport._route = lambda _session, _role: ("12", "0")
    store = SimpleNamespace(mark_prepared=Mock())
    service = SimpleNamespace(store=store)
    session = SimpleNamespace(session_id="session-1", sid_hash="a" * 64)
    status = SimpleNamespace(
        answer_count=22,
        state=SimpleNamespace(value="safety_hold"),
    )

    async def restart() -> None:
        await transport._send_publication(
            session=session,
            service=service,
            status=status,
            text="hold",
            role="owner",
            payload=attention,
            actions=[("hold_ack", "확인")],
            force_reply=False,
        )

    anyio.run(restart)

    store.mark_prepared.assert_not_called()

    recovery_transport = (
        TelegramNutritionOnboardingRuntimePublicationTransportMixin()
    )
    recovery_transport._current_publication = (
        lambda _service, _session_id: SimpleNamespace(
            generation=26,
            payload=attention,
        )
    )
    recovery_transport._route = lambda _session, _role: ("12", "0")
    prepared = SimpleNamespace(
        generation=27,
        payload=status_payload,
        state="COMMITTED",
    )
    recovery_store = SimpleNamespace(mark_prepared=Mock(return_value=prepared))
    recovery_service = SimpleNamespace(store=recovery_store)

    async def recover() -> None:
        await recovery_transport._send_publication(
            session=session,
            service=recovery_service,
            status=status,
            text="hold",
            role="owner",
            payload=status_payload,
            actions=[("hold_ack", "확인")],
            force_reply=False,
        )

    anyio.run(recover)

    recovery_store.mark_prepared.assert_called_once()


def test_operator_status_recovers_current_card_only_for_exact_owner() -> None:
    session = SimpleNamespace(
        customer_key="pilot-1",
        session_id="session-1",
        state=BootstrapState.AWAITING_ACTIVATION,
    )
    status = SimpleNamespace(
        answer_count=22,
        state=SimpleNamespace(value="safety_hold"),
    )
    service = SimpleNamespace(status=lambda: status)
    publish = AsyncMock()
    runtime = SimpleNamespace(
        _operator_notification_recovery=SimpleNamespace(
            candidate_epoch="c" * 64,
            owner_user_id=12,
            owner_chat_id=12,
            owner_topic_id=0,
        ),
        bootstrap=SimpleNamespace(
            store=SimpleNamespace(list_sessions=lambda: (session,)),
        ),
        _service=lambda _customer_key: service,
        _publish=publish,
    )

    async def exercise() -> tuple[int | None, int | None]:
        handled = await TelegramNutritionOnboardingRuntimePublicationMixin.recover_operator_status(
            runtime,
            actor_user_id=12,
            chat_id=12,
            topic_id=0,
            request_id=7001,
        )
        rejected = await TelegramNutritionOnboardingRuntimePublicationMixin.recover_operator_status(
            runtime,
            actor_user_id=99,
            chat_id=12,
            topic_id=0,
            request_id=7002,
        )
        return handled, rejected

    handled, rejected = anyio.run(exercise)

    assert handled == 1
    assert rejected is None
    publish.assert_awaited_once_with(
        session,
        service,
        status,
        operator_recovery_id="7001",
    )


def test_status_command_uses_deterministic_onboarding_recovery() -> None:
    recover = AsyncMock(return_value=1)
    runtime = SimpleNamespace(recover_operator_status=recover)
    message = SimpleNamespace(
        text="/status",
        message_id=501,
        message_thread_id=None,
        from_user=SimpleNamespace(id=12),
        chat=SimpleNamespace(id=12),
        reply_text=AsyncMock(),
    )
    update = SimpleNamespace(
        update_id=7001,
        effective_message=message,
    )
    adapter = object.__new__(TelegramAdapter)
    adapter._get_nutrition_onboarding_runtime = lambda: runtime

    anyio.run(adapter._handle_command, update, None)

    recover.assert_awaited_once_with(
        actor_user_id=12,
        chat_id=12,
        topic_id=0,
        request_id=7001,
    )
