from __future__ import annotations

import asyncio
import json
import logging
from dataclasses import dataclass
from types import SimpleNamespace
from unittest.mock import AsyncMock, Mock, call

import pytest
from telegram.error import BadRequest

from gateway.platforms.telegram_nutrition_onboarding import (
    encode_callback_hash,
)
from gateway.platforms.telegram_nutrition_onboarding_authority import (
    feature_epoch_digest,
)
from gateway.platforms.telegram_nutrition_onboarding_runtime import (
    TelegramNutritionOnboardingRuntime,
)
from gateway.platforms.telegram_nutrition_onboarding_runtime_collection import (
    TelegramNutritionOnboardingRuntimeCollectionMixin,
)
from gateway.platforms.telegram_nutrition_onboarding_publication_outbox import (
    GatewayOnboardingPublicationOutbox,
)
from gateway.platforms.telegram_customer_bootstrap import BootstrapState


def test_missing_feature_epoch_binds_to_disabled_epoch_without_creating_it(
    tmp_path,
) -> None:
    expected = (
        "2092cad6375cbbf617b88668a3103ad7fb07d854c20bdefa1b3b074b0ee135ac"
    )

    assert (
        feature_epoch_digest(
            profile_root=tmp_path,
            customer_key="client_001",
        )
        == expected
    )
    assert not (
        tmp_path
        / "data"
        / "customers"
        / "client_001"
        / "nutrition-plans"
        / "feature-epoch.json"
    ).exists()


def test_owner_dm_evidence_normalizes_absent_thread_to_zero() -> None:
    runtime = object.__new__(TelegramNutritionOnboardingRuntime)
    runtime.domain = SimpleNamespace(
        MessageEvidence=lambda **values: SimpleNamespace(**values)
    )
    message = SimpleNamespace(
        chat=SimpleNamespace(id=9000000001),
        message_thread_id=None,
        message_id=59,
    )

    evidence = runtime._evidence(
        actor_id=9000000001,
        message=message,
        update_id=123,
    )

    assert evidence.actor_user_id == 9000000001
    assert evidence.chat_id == 9000000001
    assert evidence.topic_id == 0
    assert evidence.message_id == 59
    assert evidence.update_id == 123


@dataclass
class _CallbackActor:
    id: int


@dataclass
class _CallbackQuery:
    id: str
    from_user: _CallbackActor
    answer_mock: AsyncMock
    edit_message_reply_markup_mock: AsyncMock

    async def answer(self, text: str | None = None) -> object:
        if text is None:
            return await self.answer_mock()
        return await self.answer_mock(text=text)

    async def edit_message_reply_markup(
        self,
        *,
        reply_markup: object | None = None,
    ) -> object:
        return await self.edit_message_reply_markup_mock(
            reply_markup=reply_markup
        )


@dataclass
class _CallbackMessage:
    message_id: int


def _runtime() -> tuple[
    TelegramNutritionOnboardingRuntime,
    SimpleNamespace,
    _CallbackQuery,
    _CallbackMessage,
]:
    runtime = object.__new__(TelegramNutritionOnboardingRuntime)
    session = SimpleNamespace(
        state=BootstrapState.AWAITING_ACTIVATION,
        customer_key="mother_pilot",
        session_id="session-1",
        sid_hash="a" * 64,
        chat_id="-1009000000002",
    )
    runtime.bootstrap = SimpleNamespace(
        store=SimpleNamespace(get_by_sid_hash=Mock(return_value=session))
    )
    result = SimpleNamespace(
        state=SimpleNamespace(value="trainer_review"),
        answer_count=22,
        next_field=None,
    )
    service = SimpleNamespace(
        store=SimpleNamespace(
            load_session=Mock(
                return_value=SimpleNamespace(
                    generation=29,
                    message_id=80,
                    payload={"state": "customer_attestation"},
                )
            )
        ),
        status=Mock(
            return_value=SimpleNamespace(
                state=SimpleNamespace(value="customer_attestation")
            )
        ),
        attest_baseline=Mock(return_value=result),
    )
    runtime._service = Mock(return_value=service)
    service.member_present = AsyncMock(return_value=True)
    service.retire_keyboard = AsyncMock()
    service.send_customer_confirmation = AsyncMock()
    service.publish = AsyncMock()
    setattr(runtime, "_member_present", service.member_present)
    runtime._route = Mock(
        side_effect=lambda _session, role: {
            "customer": ("-1009000000002", "0"),
            "trainer": ("-1009000000002", "71"),
            "owner": ("-1009000000002", "90"),
        }[role]
    )
    runtime._current_authority = Mock(return_value="authority")
    runtime._evidence = Mock(return_value="evidence")
    setattr(runtime, "_retire_keyboard", service.retire_keyboard)
    setattr(
        runtime,
        "_send_customer_confirmation",
        service.send_customer_confirmation,
    )
    setattr(runtime, "_publish", service.publish)
    runtime.domain = SimpleNamespace(validate_route=Mock())
    query = _CallbackQuery(
        id="callback-attest",
        from_user=_CallbackActor(id=8000000001),
        answer_mock=AsyncMock(),
        edit_message_reply_markup_mock=AsyncMock(),
    )
    message = _CallbackMessage(message_id=80)
    return runtime, service, query, message


@pytest.mark.asyncio
async def test_restart_recovers_committed_owner_approval_without_second_review(
    tmp_path,
) -> None:
    runtime = object.__new__(TelegramNutritionOnboardingRuntime)
    session = SimpleNamespace(
        state=BootstrapState.AWAITING_ACTIVATION,
        customer_key="client_001",
        session_id="session-1",
        customer_draft=SimpleNamespace(starts_on="2026-08-10"),
    )
    runtime.bootstrap = SimpleNamespace(
        store=SimpleNamespace(get=Mock(return_value=session))
    )
    finalizing = SimpleNamespace(state=SimpleNamespace(value="finalizing"))
    ready = SimpleNamespace(state=SimpleNamespace(value="ready"))
    service = SimpleNamespace(
        status=Mock(return_value=finalizing),
        finalize=Mock(return_value=ready),
    )
    runtime._service = Mock(return_value=service)
    runtime._current_authority = Mock(return_value="authority")
    runtime._recover_publication_receipts = AsyncMock()
    runtime._publish = AsyncMock()
    runtime.profile_root = tmp_path

    assert await runtime.recover_waiting_session(session) is True

    service.finalize.assert_called_once()
    kwargs = service.finalize.call_args.kwargs
    assert kwargs["starts_on"].isoformat() == "2026-08-10"
    assert kwargs["privacy_consent_digest"] is None
    assert kwargs["feature_epoch_digest"] == (
        "2092cad6375cbbf617b88668a3103ad7fb07d854c20bdefa1b3b074b0ee135ac"
    )
    assert kwargs["authority"] == "authority"
    runtime._publish.assert_awaited_once_with(session, service, ready)


@pytest.mark.asyncio
async def test_owner_callback_marks_business_commit_before_finalization_failure(
    caplog: pytest.LogCaptureFixture,
) -> None:
    runtime, service, query, message = _runtime()
    service.status.return_value = SimpleNamespace(
        state=SimpleNamespace(value="owner_review")
    )
    service.store.load_session.return_value.payload = {"state": "owner_review"}
    service.review_as_owner = Mock(
        return_value=SimpleNamespace(state=SimpleNamespace(value="finalizing"))
    )
    setattr(runtime, "_preserve_owner_callback_receipt", Mock())
    setattr(
        runtime,
        "_finalize_owner_approval",
        Mock(side_effect=OSError("post-commit finalization failure")),
    )
    setattr(
        runtime,
        "_route",
        Mock(return_value=("-1009000000002", "90")),
    )
    data = encode_callback_hash(
        action="owner_ok",
        generation=29,
        sid_hash="a" * 64,
    )

    with caplog.at_level(
        logging.INFO,
        logger=(
            "gateway.platforms."
            "telegram_nutrition_onboarding_runtime_callback"
        ),
    ):
        with pytest.raises(OSError, match="post-commit finalization failure"):
            await runtime.handle_callback(query, data, message, update_id=405)

    service.review_as_owner.assert_called_once()
    runtime._preserve_owner_callback_receipt.assert_called_once()
    assert (
        "nutrition_ingress stage=persistence update_id=405 "
        "reason_code=business_state_committed"
        in caplog.text
    )
    service.publish.assert_not_awaited()


@pytest.mark.asyncio
async def test_restart_authenticates_original_owner_update_and_records_terminal_receipt(
    tmp_path,
) -> None:
    session_id = "session-1"
    update_id = 629525050
    owner_message_id = 120
    ready_message_id = 121
    baseline_digest = "b" * 64
    authority = SimpleNamespace(
        owner_user_id=12,
        owner_chat_id=-100,
        owner_topic_id=22,
    )
    domain = SimpleNamespace(
        canonical_digest=lambda value: __import__("hashlib").sha256(
            json.dumps(
                value,
                sort_keys=True,
                separators=(",", ":"),
            ).encode()
        ).hexdigest(),
        read_private_json=lambda path: json.loads(path.read_text(encoding="utf-8")),
    )
    expected_owner_receipt = domain.canonical_digest(
        {
            "schema_version": "1.0",
            "customer_key": "client_001",
            "role": "owner",
            "decision": "approved",
            "baseline_digest": baseline_digest,
            "actor_user_id": 12,
            "chat_id": -100,
            "topic_id": 22,
            "message_id": owner_message_id,
            "update_id": update_id,
        }
    )
    ready_path = tmp_path / "ready.json"
    baseline_path = tmp_path / "baseline-v1.json"
    ready_path.write_text(
        json.dumps(
            {
                "baseline_digest": baseline_digest,
                "consumed_updates": [901, update_id],
                "state": "ready",
            }
        ),
        encoding="utf-8",
    )
    baseline_path.write_text(
        json.dumps(
            {
                "digest": baseline_digest,
                "owner_review_receipt": expected_owner_receipt,
            }
        ),
        encoding="utf-8",
    )
    outbox = GatewayOnboardingPublicationOutbox(tmp_path)
    for generation, state, message_id in (
        (25, "owner_review", owner_message_id),
        (26, "ready", ready_message_id),
    ):
        payload = {"body_digest": str(generation) * 64, "state": state}
        outbox.claim(
            session_id=session_id,
            generation=generation,
            payload=payload,
            route=("-100", "22"),
            role="owner",
            render_identity=str(generation)[-1] * 64,
        )
        outbox.record_receipt(
            session_id=session_id,
            generation=generation,
            chat_id="-100",
            topic_id="22",
            message_id=message_id,
        )
        outbox.mark_committed(
            session_id=session_id,
            generation=generation,
            payload=payload,
            route=("-100", "22"),
            role="owner",
            render_identity=str(generation)[-1] * 64,
            message_id=message_id,
        )

    runtime = object.__new__(TelegramNutritionOnboardingRuntime)
    runtime.domain = domain
    recorded_callback = outbox.record_owner_callback(
        session_id=session_id,
        customer_key="client_001",
        action="Approve",
        actor_user_id=12,
        authority=(12, -100, 22),
        route=("-100", "22"),
        message_id=owner_message_id,
        update_id=update_id,
        consumed_updates_before=(901,),
        callback_data=encode_callback_hash(
            action="owner_ok",
            generation=25,
            sid_hash="a" * 64,
        ),
        publication_generation=25,
    )
    assert (
        outbox.record_owner_callback(
            session_id=session_id,
            customer_key="client_001",
            action="Approve",
            actor_user_id=12,
            authority=(12, -100, 22),
            route=("-100", "22"),
            message_id=owner_message_id,
            update_id=update_id,
            consumed_updates_before=(901,),
            callback_data=encode_callback_hash(
                action="owner_ok",
                generation=25,
                sid_hash="a" * 64,
            ),
            publication_generation=25,
        )
        == recorded_callback
    )

    with pytest.raises(
        ValueError,
        match="owner callback receipt conflicts with immutable evidence",
    ):
        outbox.record_owner_callback(
            session_id=session_id,
            customer_key="client_001",
            action="Approve",
            actor_user_id=13,
            authority=(12, -100, 22),
            route=("-100", "22"),
            message_id=owner_message_id,
            update_id=update_id,
            consumed_updates_before=(901,),
            callback_data=encode_callback_hash(
                action="owner_ok",
                generation=25,
                sid_hash="a" * 64,
            ),
            publication_generation=25,
        )
    assert outbox.owner_callback(session_id).actor_user_id == 12

    runtime._publication_outbox = outbox
    runtime._current_authority = Mock(return_value=authority)
    runtime._route = Mock(return_value=("-100", "22"))
    runtime.adapter = SimpleNamespace(
        _reconcile_telegram_business_commit=Mock()
    )
    session = SimpleNamespace(
        session_id=session_id,
        sid_hash="a" * 64,
        customer_key="client_001",
    )
    service = SimpleNamespace(
        baseline_path=baseline_path,
        ready_path=ready_path,
        store=SimpleNamespace(
            load_session=Mock(
                return_value=SimpleNamespace(
                    generation=26,
                    message_id=ready_message_id,
                    payload={"body_digest": "26" * 64, "state": "ready"},
                    state="COMMITTED",
                )
            )
        ),
    )

    assert runtime._reconcile_owner_callback_receipt(session, service) is True

    runtime.adapter._reconcile_telegram_business_commit.assert_called_once()
    call_kwargs = runtime.adapter._reconcile_telegram_business_commit.call_args.kwargs
    assert call_kwargs["update_id"] == update_id
    assert call_kwargs["actor_id"] == 12
    assert call_kwargs["chat_id"] == -100
    assert call_kwargs["topic_id"] == 22
    assert call_kwargs["message_id"] == owner_message_id
    assert len(call_kwargs["provenance_digest"]) == 64

    callback_path = outbox.callback_path
    original_callback = json.loads(callback_path.read_text(encoding="utf-8"))
    conflicts = {
        "actor": ("actor_user_id", 13),
        "route": ("route", ["-100", "23"]),
        "message": ("message_id", owner_message_id + 1),
        "action": ("action", "Reject"),
        "update": ("update_id", update_id + 1),
        "publication": ("publication_generation", 24),
        "authority": ("authority", [13, -100, 22]),
    }
    for label, (field, conflicting_value) in conflicts.items():
        conflicting = json.loads(json.dumps(original_callback))
        conflicting["records"][0][field] = conflicting_value
        callback_path.write_text(json.dumps(conflicting), encoding="utf-8")
        runtime.adapter._reconcile_telegram_business_commit.reset_mock()
        assert runtime._reconcile_owner_callback_receipt(session, service) is False, label
        runtime.adapter._reconcile_telegram_business_commit.assert_not_called()
    missing_callback = json.loads(json.dumps(original_callback))
    missing_callback["records"] = []
    callback_path.write_text(json.dumps(missing_callback), encoding="utf-8")
    assert runtime._reconcile_owner_callback_receipt(session, service) is False
    runtime.adapter._reconcile_telegram_business_commit.assert_not_called()
    callback_path.write_text(json.dumps(original_callback), encoding="utf-8")

    ready_path.write_text(
        json.dumps(
            {
                "baseline_digest": baseline_digest,
                "consumed_updates": [update_id, 901],
                "state": "ready",
            }
        ),
        encoding="utf-8",
    )
    assert runtime._reconcile_owner_callback_receipt(session, service) is False
    runtime.adapter._reconcile_telegram_business_commit.assert_not_called()


@pytest.mark.asyncio
async def test_attestation_callback_logs_and_advances(
    caplog: pytest.LogCaptureFixture,
) -> None:
    runtime, service, query, message = _runtime()
    data = encode_callback_hash(
        action="attest",
        generation=29,
        sid_hash="a" * 64,
    )

    with caplog.at_level(
        logging.INFO,
        logger=(
            "gateway.platforms."
            "telegram_nutrition_onboarding_runtime_callback"
        ),
    ):
        await runtime.handle_callback(
            query,
            data,
            message,
            update_id=197,
        )

    service.attest_baseline.assert_called_once_with(
        authority="authority",
        evidence="evidence",
    )
    runtime.domain.validate_route.assert_called_once_with(
        "authority",
        "evidence",
        role="customer",
    )
    service.publish.assert_awaited_once()
    assert "actor_id=8000000001" not in caplog.text
    assert data not in caplog.text
    assert query.id not in caplog.text
    assert "nutrition_onboarding_callback_applied" in caplog.text
    assert "action=attest" not in caplog.text
    stage_messages = [
        record.getMessage()
        for record in caplog.records
        if record.getMessage().startswith(
            "nutrition_onboarding_callback_stage"
        )
    ]
    assert stage_messages == [
        "nutrition_onboarding_callback_stage stage=member_check_start",
        "nutrition_onboarding_callback_stage stage=member_check_pass",
        "nutrition_onboarding_callback_stage stage=authority_pass",
        "nutrition_onboarding_callback_stage stage=route_pass",
        "nutrition_onboarding_callback_stage stage=transition_pass",
    ]
    ingress_stages = [
        record.getMessage()
        for record in caplog.records
        if record.getMessage().startswith("nutrition_ingress ")
    ]
    assert ingress_stages == [
        "nutrition_ingress stage=ingress update_id=197",
        "nutrition_ingress stage=authority update_id=197",
        "nutrition_ingress stage=route update_id=197",
        "nutrition_ingress stage=validation update_id=197",
        "nutrition_ingress stage=persistence update_id=197",
        "nutrition_ingress stage=publication update_id=197",
        "nutrition_ingress stage=receipt update_id=197",
    ]


@pytest.mark.asyncio
async def test_restart_callback_requires_the_signed_current_card_route(
    tmp_path,
) -> None:
    outbox = GatewayOnboardingPublicationOutbox(tmp_path)
    payload = {"state": "customer_attestation"}
    outbox.claim(
        session_id="session-1",
        generation=29,
        payload=payload,
        route=("-1009000000002", "0"),
        role="customer",
        render_identity="a" * 64,
    )
    outbox.record_receipt(
        session_id="session-1",
        generation=29,
        chat_id="-1009000000002",
        topic_id="0",
        message_id=80,
    )
    outbox.mark_committed(
        session_id="session-1",
        generation=29,
        payload=payload,
        route=("-1009000000002", "0"),
        role="customer",
        render_identity="a" * 64,
        message_id=80,
    )
    runtime, service, query, message = _runtime()
    runtime._publication_outbox = GatewayOnboardingPublicationOutbox(tmp_path)
    runtime._route = Mock(return_value=("-1009000000002", "0"))
    message.chat = SimpleNamespace(id=-1009000000002)
    message.message_thread_id = 0
    data = encode_callback_hash(
        action="attest",
        generation=29,
        sid_hash="a" * 64,
    )

    await runtime.handle_callback(query, data, message, update_id=197)

    service.attest_baseline.assert_called_once()
    service.publish.assert_awaited_once()

    restarted, restarted_service, stale_query, stale_message = _runtime()
    restarted._publication_outbox = GatewayOnboardingPublicationOutbox(tmp_path)
    restarted._route = Mock(return_value=("-1009000000002", "0"))
    stale_message.chat = SimpleNamespace(id=-1004325423948)
    stale_message.message_thread_id = 0

    await restarted.handle_callback(stale_query, data, stale_message, update_id=198)

    restarted_service.attest_baseline.assert_not_called()
    restarted_service.publish.assert_not_awaited()
    assert stale_query.answer_mock.await_args_list == [
        call(),
        call(text="가장 최근 온보딩 카드를 사용해 주세요."),
    ]


@pytest.mark.asyncio
async def test_stale_attestation_callback_logs_rejection(
    caplog: pytest.LogCaptureFixture,
) -> None:
    runtime, service, query, message = _runtime()
    data = encode_callback_hash(
        action="attest",
        generation=28,
        sid_hash="a" * 64,
    )

    with caplog.at_level(
        logging.INFO,
        logger=(
            "gateway.platforms."
            "telegram_nutrition_onboarding_runtime_callback"
        ),
    ):
        await runtime.handle_callback(
            query,
            data,
            message,
            update_id=198,
        )

    assert query.answer_mock.await_args_list == [
        call(),
        call(text="가장 최근 온보딩 카드를 사용해 주세요."),
    ]
    service.attest_baseline.assert_not_called()
    assert "actor_id=8000000001" not in caplog.text
    assert data not in caplog.text
    assert query.id not in caplog.text
    assert (
        "nutrition_ingress stage=validation update_id=198 "
        "reason_code=stale_publication"
        in caplog.text
    )


@pytest.mark.asyncio
async def test_stale_owner_review_callback_is_rejected_before_mutation() -> None:
    runtime, service, query, message = _runtime()
    service.status.return_value = SimpleNamespace(
        state=SimpleNamespace(value="owner_review")
    )
    service.store.load_session.return_value.payload = {"state": "owner_review"}
    service.review_as_owner = Mock()
    data = encode_callback_hash(
        action="owner_ok",
        generation=28,
        sid_hash="a" * 64,
    )

    await runtime.handle_callback(query, data, message, update_id=199)

    assert query.answer_mock.await_args_list == [
        call(),
        call(text="가장 최근 온보딩 카드를 사용해 주세요."),
    ]
    service.review_as_owner.assert_not_called()
    service.publish.assert_not_awaited()


class _TextIngressRuntime(TelegramNutritionOnboardingRuntimeCollectionMixin):
    def __init__(self) -> None:
        self.session = SimpleNamespace(
            state=BootstrapState.AWAITING_ACTIVATION,
            customer_key="test_customer",
            session_id="session-1",
        )
        self.status = SimpleNamespace(
            state=SimpleNamespace(value="collecting"),
            next_field="activity_rationale",
        )
        self.service = SimpleNamespace(
            store=SimpleNamespace(
                load_session=Mock(
                    return_value=SimpleNamespace(
                        message_id=81,
                        payload={},
                    )
                )
            ),
            status=Mock(return_value=self.status),
            submit_answer=Mock(return_value=self.status),
        )
        self.published = AsyncMock()

    def _session_for_customer_message(self, message: object) -> object | None:
        return self.session

    def _current_authority(self, session: object) -> str:
        assert session is self.session
        return "authority"

    def _service(self, customer_key: str) -> object:
        assert customer_key == "test_customer"
        return self.service

    def _evidence(
        self,
        *,
        actor_id: int | str | None,
        message: object,
        update_id: int | str | None,
    ) -> str:
        assert actor_id == 12
        assert update_id == 313
        return "evidence"

    async def _publish(
        self,
        session: object,
        service: object,
        status: object,
        *,
        reply_anchor_message_id: int | None = None,
    ) -> None:
        assert session is self.session
        assert service is self.service
        assert status is self.status
        assert reply_anchor_message_id == 91
        await self.published()


def _text_message(*, reply_message_id: int | None = 81) -> SimpleNamespace:
    reply = (
        None
        if reply_message_id is None
        else SimpleNamespace(message_id=reply_message_id)
    )
    return SimpleNamespace(
        text="PRIVATE_ANSWER_SHOULD_NOT_APPEAR",
        message_id=91,
        reply_to_message=reply,
        from_user=SimpleNamespace(id=12),
        chat=SimpleNamespace(id=34),
        message_thread_id=0,
        reply_text=AsyncMock(),
    )


@pytest.mark.asyncio
async def test_text_ingress_records_ordered_stages_without_answer_text(
    caplog: pytest.LogCaptureFixture,
    monkeypatch: pytest.MonkeyPatch,
) -> None:
    runtime = _TextIngressRuntime()
    monkeypatch.setattr(
        "gateway.platforms.telegram_nutrition_onboarding_runtime_collection.parse_answer",
        lambda _field, _value: "parsed",
    )

    with caplog.at_level(
        logging.INFO,
        logger=(
            "gateway.platforms."
            "telegram_nutrition_onboarding_runtime_collection"
        ),
    ):
        assert await runtime.handle_text(
            SimpleNamespace(update_id=313),
            _text_message(),
        ) is True

    runtime.service.submit_answer.assert_called_once_with(
        field="activity_rationale",
        value="parsed",
        authority="authority",
        evidence="evidence",
    )
    runtime.published.assert_awaited_once()
    assert [
        record.getMessage()
        for record in caplog.records
        if record.getMessage().startswith("nutrition_ingress ")
    ] == [
        "nutrition_ingress stage=ingress update_id=313",
        "nutrition_ingress stage=authority update_id=313",
        "nutrition_ingress stage=route update_id=313",
        "nutrition_ingress stage=validation update_id=313",
        "nutrition_ingress stage=persistence update_id=313",
        "nutrition_ingress stage=publication update_id=313",
        "nutrition_ingress stage=receipt update_id=313",
    ]
    assert "PRIVATE_ANSWER_SHOULD_NOT_APPEAR" not in caplog.text


@pytest.mark.asyncio
async def test_text_ingress_rejection_has_update_keyed_reason_code(
    caplog: pytest.LogCaptureFixture,
) -> None:
    runtime = _TextIngressRuntime()

    with caplog.at_level(
        logging.INFO,
        logger=(
            "gateway.platforms."
            "telegram_nutrition_onboarding_runtime_collection"
        ),
    ):
        assert await runtime.handle_text(
            SimpleNamespace(update_id=313),
            _text_message(reply_message_id=80),
        ) is True

    runtime.service.submit_answer.assert_not_called()
    runtime.published.assert_not_awaited()
    assert (
        "nutrition_ingress stage=validation update_id=313 "
        "reason_code=stale_reply"
        in caplog.text
    )


@pytest.mark.asyncio
async def test_callback_authority_route_and_persistence_rejections_are_coded(
    caplog: pytest.LogCaptureFixture,
) -> None:
    runtime, service, query, message = _runtime()
    data = encode_callback_hash(
        action="attest",
        generation=29,
        sid_hash="a" * 64,
    )
    runtime._current_authority = Mock(side_effect=OSError("unavailable"))

    with caplog.at_level(
        logging.INFO,
        logger=(
            "gateway.platforms."
            "telegram_nutrition_onboarding_runtime_callback"
        ),
    ):
        await runtime.handle_callback(query, data, message, update_id=401)

    service.attest_baseline.assert_not_called()
    assert (
        "nutrition_ingress stage=authority update_id=401 "
        "reason_code=authority_unavailable"
        in caplog.text
    )

    runtime, service, query, message = _runtime()
    runtime.domain.validate_route = Mock(side_effect=ValueError("rejected"))
    with caplog.at_level(
        logging.INFO,
        logger=(
            "gateway.platforms."
            "telegram_nutrition_onboarding_runtime_callback"
        ),
    ):
        await runtime.handle_callback(query, data, message, update_id=402)

    service.attest_baseline.assert_not_called()
    assert (
        "nutrition_ingress stage=route update_id=402 "
        "reason_code=route_rejected"
        in caplog.text
    )

    runtime, service, query, message = _runtime()
    service.attest_baseline = Mock(side_effect=ValueError("rejected"))
    with caplog.at_level(
        logging.INFO,
        logger=(
            "gateway.platforms."
            "telegram_nutrition_onboarding_runtime_callback"
        ),
    ):
        await runtime.handle_callback(query, data, message, update_id=403)

    service.publish.assert_not_awaited()
    assert (
        "nutrition_ingress stage=persistence update_id=403 "
        "reason_code=transition_rejected"
        in caplog.text
    )


@pytest.mark.asyncio
async def test_callback_publication_failure_has_a_reason_code(
    caplog: pytest.LogCaptureFixture,
) -> None:
    runtime, _service, query, message = _runtime()
    runtime._publish = AsyncMock(side_effect=OSError("unavailable"))
    data = encode_callback_hash(
        action="attest",
        generation=29,
        sid_hash="a" * 64,
    )

    with caplog.at_level(
        logging.INFO,
        logger=(
            "gateway.platforms."
            "telegram_nutrition_onboarding_runtime_callback"
        ),
    ):
        with pytest.raises(OSError, match="unavailable"):
            await runtime.handle_callback(query, data, message, update_id=404)

    assert (
        "nutrition_ingress stage=publication update_id=404 "
        "reason_code=publication_failed"
        in caplog.text
    )


@pytest.mark.asyncio
async def test_expired_onboarding_ack_is_attempted_before_commit_and_publish(
    caplog: pytest.LogCaptureFixture,
) -> None:
    runtime, service, query, message = _runtime()
    events: list[str] = []
    advanced = SimpleNamespace(
        state=SimpleNamespace(value="trainer_review"),
        answer_count=22,
        next_field=None,
    )

    def attest_baseline(**_kwargs: object) -> SimpleNamespace:
        events.append("commit")
        return advanced

    async def expired_answer(**_kwargs: object) -> None:
        events.append("ack")
        raise BadRequest("Query is too old")

    service.attest_baseline = Mock(side_effect=attest_baseline)
    query.answer_mock.side_effect = expired_answer
    query.id = "onboarding-callback-19"
    data = encode_callback_hash(
        action="attest",
        generation=29,
        sid_hash="a" * 64,
    )

    with caplog.at_level(
        logging.WARNING,
        logger=(
            "gateway.platforms."
            "telegram_nutrition_onboarding_runtime_callback"
        ),
    ):
        await runtime.handle_callback(query, data, message, update_id=199)

    assert events[:2] == ["ack", "commit"]
    service.publish.assert_awaited_once_with(
        runtime.bootstrap.store.get_by_sid_hash.return_value,
        service,
        advanced,
    )
    assert (
        "nutrition_onboarding_callback_ack_failed "
        "gate=onboarding_callback update_id=199 error=BadRequest"
        in caplog.text
    )


@pytest.mark.asyncio
async def test_expired_trainer_callback_is_acknowledged_without_domain_or_publication(
    caplog: pytest.LogCaptureFixture,
) -> None:
    runtime, service, query, message = _runtime()
    session = runtime.bootstrap.store.get_by_sid_hash.return_value
    service.store.load_session.return_value = SimpleNamespace(
        generation=30,
        message_id=82,
        payload={"state": "trainer_review"},
    )
    service.status.return_value = SimpleNamespace(
        state=SimpleNamespace(value="trainer_review")
    )
    service.review_as_trainer = Mock()
    query.from_user.id = 6949941558
    query.answer_mock.side_effect = BadRequest(
        "Query is too old and response timeout expired or query id is invalid"
    )
    message.message_id = 82
    data = encode_callback_hash(
        action="train_ok",
        generation=30,
        sid_hash="a" * 64,
    )

    with caplog.at_level(
        logging.WARNING,
        logger=(
            "gateway.platforms."
            "telegram_nutrition_onboarding_runtime_callback"
        ),
    ):
        await runtime.handle_callback(query, data, message, update_id=200)

    service.review_as_trainer.assert_not_called()
    service.publish.assert_not_awaited()
    assert (
        "nutrition_onboarding_callback_ack_failed "
        "gate=onboarding_callback update_id=200 error=BadRequest"
        in caplog.text
    )
