"""Gateway-owned, authenticated receipt ledger for nutrition-onboarding cards."""

from __future__ import annotations

import fcntl
import hashlib
import hmac
import json
import os
import re
import secrets
import stat
import tempfile
from collections.abc import Iterator, Mapping
from contextlib import contextmanager
from dataclasses import dataclass
from pathlib import Path
from typing import Final

_SCHEMA: Final = "telegram-nutrition-onboarding-publication-outbox-v2"
_LEGACY_SCHEMA: Final = "telegram-nutrition-onboarding-publication-outbox-v1"
_CALLBACK_SCHEMA: Final = "telegram-nutrition-onboarding-owner-callback-v1"
_DIGEST: Final = re.compile(r"^[a-f0-9]{64}$")
_STATES: Final = frozenset({"DISPATCHING", "RECEIPTED", "COMMITTED"})
_EMERGENCY_STATES: Final = frozenset({"RECEIPTED", "COMMITTED"})
_RELATIVE_PATH: Final = Path("data/onboarding/telegram-publication-outbox-v1")
_PRIVATE_MODE: Final = 0o600
_DIRECTORY_MODE: Final = 0o700
_NOATIME: Final = getattr(os, "O_NOATIME", 0)


@dataclass(frozen=True, slots=True)
class GatewayPublicationReceipt:
    """One exact publication authority and its optional provider receipt."""

    session_id: str
    generation: int
    payload: dict[str, object]
    route: tuple[str, str]
    role: str
    render_identity: str
    payload_digest: str
    dispatch_identity: str
    state: str
    message_id: int | None = None
    receipt_integrity: str | None = None


@dataclass(frozen=True, slots=True)
class GatewayOwnerCallbackReceipt:
    """One immutable authenticated owner callback event."""

    session_id: str
    customer_key: str
    action: str
    actor_user_id: int
    authority: tuple[int, int, int]
    route: tuple[str, str]
    message_id: int
    update_id: int
    consumed_updates_before: tuple[int, ...]
    callback_data: str
    publication_generation: int
    publication_dispatch_identity: str
    publication_receipt_integrity: str
    event_integrity: str


class GatewayOnboardingPublicationOutbox:
    """Persist authenticated Telegram receipts before profile publication commit."""

    def __init__(self, profile_root: Path | str, *, initialize: bool = True) -> None:
        self.state_dir = Path(profile_root) / _RELATIVE_PATH
        self.ledger_path = self.state_dir / "ledger.json"
        self.legacy_ledger_path = self.state_dir / "ledger-v1-untrusted.json"
        self.lock_path = self.state_dir / ".lock"
        self.emergency_path = self.state_dir / "emergency.json"
        self.emergency_lock_path = self.state_dir / ".emergency.lock"
        self.callback_path = self.state_dir / "owner-callbacks.json"
        self.callback_lock_path = self.state_dir / ".owner-callbacks.lock"
        self.key_path = self.state_dir / ".receipt-key"
        self._callback_authority_absent = False
        if initialize:
            self._prepare_paths()
        else:
            self._inspect_paths()

    def claim(
        self,
        *,
        session_id: str,
        generation: int,
        payload: Mapping[str, object],
        route: tuple[str, str],
        role: str,
        render_identity: str,
    ) -> tuple[GatewayPublicationReceipt, bool]:
        """Claim one exact publication before the provider call."""
        candidate = self._candidate(
            session_id=session_id,
            generation=generation,
            payload=payload,
            route=route,
            role=role,
            render_identity=render_identity,
        )
        with self._locked(self.lock_path):
            records = list(self._read_primary_unlocked())
            match = self._matching_record(records, candidate.session_id, candidate.generation)
            if match is None:
                records.append(candidate)
                self._write_primary_unlocked(records)
                return candidate, True
            if not self._same_authority(match, candidate):
                raise ValueError("onboarding publication authority does not match")
            return match, False

    def get(
        self,
        *,
        session_id: str,
        generation: int,
        payload: Mapping[str, object],
        route: tuple[str, str],
        role: str,
        render_identity: str,
    ) -> GatewayPublicationReceipt | None:
        """Resolve only an exact current publication authority."""
        candidate = self._candidate(
            session_id=session_id,
            generation=generation,
            payload=payload,
            route=route,
            role=role,
            render_identity=render_identity,
        )
        with self._locked(self.lock_path):
            record = self._matching_record(
                self._read_primary_unlocked(), candidate.session_id, candidate.generation
            )
        if record is None:
            return None
        if not self._same_authority(record, candidate):
            raise ValueError("onboarding publication authority does not match")
        return record

    def record_receipt(
        self,
        *,
        session_id: str,
        generation: int,
        chat_id: str,
        topic_id: str,
        message_id: int,
    ) -> GatewayPublicationReceipt:
        """Attach an authenticated Telegram receipt to one dispatch authority."""
        return self._record_receipt(
            session_id=session_id,
            generation=generation,
            chat_id=chat_id,
            topic_id=topic_id,
            message_id=message_id,
            emergency=False,
        )

    def record_emergency_receipt(
        self,
        *,
        session_id: str,
        generation: int,
        chat_id: str,
        topic_id: str,
        message_id: int,
    ) -> GatewayPublicationReceipt:
        """Persist a separate receipt when the primary ledger write fails."""
        return self._record_receipt(
            session_id=session_id,
            generation=generation,
            chat_id=chat_id,
            topic_id=topic_id,
            message_id=message_id,
            emergency=True,
        )

    def reconcile_emergency_receipts(self) -> int:
        """Merge exact emergency evidence into the primary receipt ledger."""
        with self._locked(self.lock_path):
            records = list(self._read_primary_unlocked())
            with self._locked(self.emergency_lock_path):
                emergency_records = self._read_emergency_unlocked()
            changed = 0
            for emergency in emergency_records:
                if emergency.state == "COMMITTED":
                    continue
                index, primary = self._find(records, emergency.session_id, emergency.generation)
                if not self._same_authority(primary, emergency):
                    raise ValueError("emergency publication authority does not match")
                if not self._valid_receipt(emergency):
                    raise ValueError("emergency publication receipt integrity is invalid")
                if primary.state == "DISPATCHING":
                    records[index] = emergency
                    changed += 1
                    continue
                if (
                    primary.state != "RECEIPTED"
                    or primary.message_id != emergency.message_id
                    or not hmac.compare_digest(
                        primary.receipt_integrity or "", emergency.receipt_integrity or ""
                    )
                ):
                    raise ValueError("emergency publication receipt conflicts with primary evidence")
            if changed:
                self._write_primary_unlocked(records)
            return changed

    def mark_committed(
        self,
        *,
        session_id: str,
        generation: int,
        payload: Mapping[str, object],
        route: tuple[str, str],
        role: str,
        render_identity: str,
        message_id: int,
    ) -> GatewayPublicationReceipt:
        """Retire an exact receipt only after the profile commit succeeds."""
        candidate = self._candidate(
            session_id=session_id,
            generation=generation,
            payload=payload,
            route=route,
            role=role,
            render_identity=render_identity,
        )
        with self._locked(self.lock_path):
            records = list(self._read_primary_unlocked())
            index, current = self._find(records, session_id, generation)
            if not self._same_authority(current, candidate):
                raise ValueError("onboarding publication authority does not match")
            if current.message_id != message_id or not self._valid_receipt(current):
                raise ValueError("onboarding publication receipt integrity is invalid")
            if current.state == "COMMITTED":
                return current
            if current.state != "RECEIPTED":
                raise ValueError("onboarding publication has no committed receipt")
            updated = self._with_state(current, "COMMITTED")
            records[index] = updated
            self._write_primary_unlocked(records)
        self._mark_emergency_committed(updated)
        return updated

    def receipted(self) -> tuple[GatewayPublicationReceipt, ...]:
        with self._locked(self.lock_path):
            return tuple(
                record
                for record in self._read_primary_unlocked()
                if record.state == "RECEIPTED"
            )

    def records(self) -> tuple[GatewayPublicationReceipt, ...]:
        with self._locked(self.lock_path):
            return self._read_primary_unlocked()

    def emergency_records(self) -> tuple[GatewayPublicationReceipt, ...]:
        with self._locked(self.emergency_lock_path):
            return self._read_emergency_unlocked()

    def record_owner_callback(
        self,
        *,
        session_id: str,
        customer_key: str,
        action: str,
        actor_user_id: int,
        authority: tuple[int, int, int],
        route: tuple[str, str],
        message_id: int,
        update_id: int,
        consumed_updates_before: tuple[int, ...],
        callback_data: str,
        publication_generation: int,
    ) -> GatewayOwnerCallbackReceipt:
        """Persist the full authenticated event before owner state finalization."""
        with self._locked(self.lock_path):
            publication = self._matching_record(
                self._read_primary_unlocked(), session_id, publication_generation
            )
        if (
            publication is None
            or publication.state != "COMMITTED"
            or publication.role != "owner"
            or publication.route != self._route(*route)
            or publication.message_id != self._message_id(message_id)
            or not self._valid_receipt(publication)
        ):
            raise ValueError("owner callback publication authority does not match")
        receipt = self._owner_callback_candidate(
            session_id=session_id,
            customer_key=customer_key,
            action=action,
            actor_user_id=actor_user_id,
            authority=authority,
            route=route,
            message_id=message_id,
            update_id=update_id,
            consumed_updates_before=consumed_updates_before,
            callback_data=callback_data,
            publication=publication,
        )
        with self._locked(self.callback_lock_path):
            records = list(self._read_owner_callbacks_unlocked())
            matches = [item for item in records if item.session_id == receipt.session_id]
            if matches:
                if len(matches) != 1 or matches[0] != receipt:
                    raise ValueError("owner callback receipt conflicts with immutable evidence")
                return matches[0]
            records.append(receipt)
            self._write_owner_callbacks_unlocked(records)
        return receipt

    def owner_callback(self, session_id: str) -> GatewayOwnerCallbackReceipt | None:
        """Return the sole integrity-checked owner callback for one session."""
        session = self._session_id(session_id)
        if self._callback_authority_absent:
            if self._entry_exists(self.callback_path) or self._entry_exists(self.callback_lock_path):
                raise ValueError("owner callback receipt authority changed during inspection")
            return None
        with self._locked(self.callback_lock_path):
            matches = [
                item
                for item in self._read_owner_callbacks_unlocked()
                if item.session_id == session
            ]
        if not matches:
            return None
        if len(matches) != 1:
            raise ValueError("owner callback receipt is ambiguous")
        return matches[0]

    def _record_receipt(
        self,
        *,
        session_id: str,
        generation: int,
        chat_id: str,
        topic_id: str,
        message_id: int,
        emergency: bool,
    ) -> GatewayPublicationReceipt:
        route = self._route(chat_id, topic_id)
        message = self._message_id(message_id)
        with self._locked(self.lock_path):
            primary = self._matching_record(
                self._read_primary_unlocked(), session_id, generation
            )
        if primary is None or primary.state not in {"DISPATCHING", "RECEIPTED"}:
            raise ValueError("onboarding publication is not awaiting a receipt")
        if primary.route != route:
            raise ValueError("Telegram receipt route does not match dispatch authority")
        received = self._with_receipt(primary, message)
        if primary.state == "RECEIPTED":
            if (
                primary.message_id != received.message_id
                or not hmac.compare_digest(
                    primary.receipt_integrity or "", received.receipt_integrity or ""
                )
            ):
                raise ValueError("Telegram receipt conflicts with dispatch authority")
            return primary
        if emergency:
            with self._locked(self.emergency_lock_path):
                records = list(self._read_emergency_unlocked())
                existing = self._matching_record(records, session_id, generation)
                if existing is not None:
                    if (
                        not self._same_authority(existing, received)
                        or existing.message_id != received.message_id
                        or not hmac.compare_digest(
                            existing.receipt_integrity or "",
                            received.receipt_integrity or "",
                        )
                    ):
                        raise ValueError("emergency publication receipt conflicts with dispatch authority")
                    return existing
                records.append(received)
                self._write_emergency_unlocked(records)
                return received
        with self._locked(self.lock_path):
            records = list(self._read_primary_unlocked())
            index, current = self._find(records, session_id, generation)
            if not self._same_authority(current, primary) or current.state != "DISPATCHING":
                raise ValueError("onboarding publication dispatch authority is stale")
            records[index] = received
            self._write_primary_unlocked(records)
            return received

    def _mark_emergency_committed(self, primary: GatewayPublicationReceipt) -> None:
        with self._locked(self.emergency_lock_path):
            records = list(self._read_emergency_unlocked())
            existing = self._matching_record(records, primary.session_id, primary.generation)
            if existing is None or existing.state == "COMMITTED":
                return
            if (
                not self._same_authority(existing, primary)
                or existing.message_id != primary.message_id
                or not hmac.compare_digest(
                    existing.receipt_integrity or "", primary.receipt_integrity or ""
                )
            ):
                raise ValueError("emergency publication receipt conflicts with primary evidence")
            index, _ = self._find(records, primary.session_id, primary.generation)
            records[index] = self._with_state(existing, "COMMITTED")
            self._write_emergency_unlocked(records)

    def _candidate(
        self,
        *,
        session_id: str,
        generation: int,
        payload: Mapping[str, object],
        route: tuple[str, str],
        role: str,
        render_identity: str,
    ) -> GatewayPublicationReceipt:
        session = self._session_id(session_id)
        current_generation = self._generation(generation)
        current_route = self._route(*route)
        current_role = self._role(role)
        if not isinstance(render_identity, str) or _DIGEST.fullmatch(render_identity) is None:
            raise ValueError("onboarding publication render identity is invalid")
        normalized_payload = self._payload(payload)
        payload_digest = self._payload_digest(
            session,
            current_generation,
            current_route,
            current_role,
            render_identity,
            normalized_payload,
        )
        dispatch_identity = self._dispatch_identity(
            session,
            current_generation,
            current_route,
            current_role,
            render_identity,
            payload_digest,
        )
        return GatewayPublicationReceipt(
            session,
            current_generation,
            normalized_payload,
            current_route,
            current_role,
            render_identity,
            payload_digest,
            dispatch_identity,
            "DISPATCHING",
        )

    def _with_receipt(
        self,
        record: GatewayPublicationReceipt,
        message_id: int,
    ) -> GatewayPublicationReceipt:
        integrity = self._receipt_integrity(record, message_id)
        return GatewayPublicationReceipt(
            record.session_id,
            record.generation,
            record.payload,
            record.route,
            record.role,
            record.render_identity,
            record.payload_digest,
            record.dispatch_identity,
            "RECEIPTED",
            message_id,
            integrity,
        )

    @staticmethod
    def _with_state(
        record: GatewayPublicationReceipt,
        state: str,
    ) -> GatewayPublicationReceipt:
        return GatewayPublicationReceipt(
            record.session_id,
            record.generation,
            record.payload,
            record.route,
            record.role,
            record.render_identity,
            record.payload_digest,
            record.dispatch_identity,
            state,
            record.message_id,
            record.receipt_integrity,
        )

    def _same_authority(
        self,
        left: GatewayPublicationReceipt,
        right: GatewayPublicationReceipt,
    ) -> bool:
        return (
            left.session_id == right.session_id
            and left.generation == right.generation
            and left.payload == right.payload
            and left.route == right.route
            and left.role == right.role
            and hmac.compare_digest(left.render_identity, right.render_identity)
            and hmac.compare_digest(left.payload_digest, right.payload_digest)
            and hmac.compare_digest(left.dispatch_identity, right.dispatch_identity)
        )

    def _valid_receipt(self, record: GatewayPublicationReceipt) -> bool:
        return bool(
            record.state in {"RECEIPTED", "COMMITTED"}
            and record.message_id is not None
            and record.receipt_integrity is not None
            and hmac.compare_digest(
                record.receipt_integrity,
                self._receipt_integrity(record, record.message_id),
            )
        )

    @staticmethod
    def _session_id(value: str) -> str:
        if not isinstance(value, str) or not value or len(value) > 256:
            raise ValueError("onboarding session ID is invalid")
        return value

    @staticmethod
    def _generation(value: int) -> int:
        if type(value) is not int or value < 0:
            raise ValueError("onboarding publication generation is invalid")
        return value

    @staticmethod
    def _route(chat_id: str, topic_id: str) -> tuple[str, str]:
        if (
            not isinstance(chat_id, str)
            or not chat_id
            or not isinstance(topic_id, str)
            or not topic_id
        ):
            raise ValueError("onboarding publication route is invalid")
        return chat_id, topic_id

    @staticmethod
    def _role(value: str) -> str:
        if value not in {"customer", "owner"}:
            raise ValueError("onboarding publication role is invalid")
        return value

    @staticmethod
    def _message_id(value: int) -> int:
        if type(value) is not int or isinstance(value, bool) or value <= 0:
            raise ValueError("Telegram publication message ID is invalid")
        return value

    @staticmethod
    def _payload(value: Mapping[str, object]) -> dict[str, object]:
        try:
            canonical = json.dumps(
                dict(value),
                ensure_ascii=False,
                sort_keys=True,
                separators=(",", ":"),
                allow_nan=False,
            )
            decoded = json.loads(canonical)
        except (TypeError, ValueError) as exc:
            raise ValueError("onboarding publication payload is invalid") from exc
        if not isinstance(decoded, dict) or any(not isinstance(key, str) for key in decoded):
            raise ValueError("onboarding publication payload is invalid")
        return decoded

    @staticmethod
    def _canonical(value: Mapping[str, object]) -> bytes:
        return json.dumps(
            dict(value),
            ensure_ascii=False,
            sort_keys=True,
            separators=(",", ":"),
            allow_nan=False,
        ).encode("utf-8")

    @classmethod
    def _payload_digest(
        cls,
        session_id: str,
        generation: int,
        route: tuple[str, str],
        role: str,
        render_identity: str,
        payload: Mapping[str, object],
    ) -> str:
        return hashlib.sha256(cls._canonical({
            "session_id": session_id,
            "generation": generation,
            "route": list(route),
            "role": role,
            "render_identity": render_identity,
            "payload": dict(payload),
        })).hexdigest()

    @classmethod
    def _dispatch_identity(
        cls,
        session_id: str,
        generation: int,
        route: tuple[str, str],
        role: str,
        render_identity: str,
        payload_digest: str,
    ) -> str:
        return hashlib.sha256(cls._canonical({
            "domain": "telegram-nutrition-onboarding-dispatch-v1",
            "session_id": session_id,
            "generation": generation,
            "route": list(route),
            "role": role,
            "render_identity": render_identity,
            "payload_digest": payload_digest,
        })).hexdigest()

    def _receipt_integrity(
        self,
        record: GatewayPublicationReceipt,
        message_id: int,
    ) -> str:
        payload = self._canonical({
            "domain": "telegram-nutrition-onboarding-receipt-v1",
            "dispatch_identity": record.dispatch_identity,
            "payload_digest": record.payload_digest,
            "route": list(record.route),
            "message_id": message_id,
        })
        return hmac.new(self._read_key(), payload, hashlib.sha256).hexdigest()

    def _prepare_paths(self) -> None:
        self._validate_path_ancestry(require_state_dir=False)
        self.state_dir.mkdir(parents=True, exist_ok=True, mode=_DIRECTORY_MODE)
        directory = os.stat(self.state_dir, follow_symlinks=False)
        if not stat.S_ISDIR(directory.st_mode) or stat.S_ISLNK(directory.st_mode):
            raise ValueError("onboarding publication outbox directory is invalid")
        os.chmod(self.state_dir, _DIRECTORY_MODE)
        self._ensure_private_file(self.lock_path)
        self._ensure_private_file(self.emergency_lock_path)
        self._ensure_private_file(self.callback_lock_path)
        self._ensure_key()
        self._ensure_ledger(self.ledger_path, emergency=False)
        self._ensure_ledger(self.emergency_path, emergency=True)
        self._ensure_callback_ledger()

    def _inspect_paths(self) -> None:
        """Validate existing authorities without creating or repairing any path."""
        self._validate_path_ancestry(require_state_dir=True)
        directory = os.stat(self.state_dir, follow_symlinks=False)
        if (
            not stat.S_ISDIR(directory.st_mode)
            or stat.S_ISLNK(directory.st_mode)
            or stat.S_IMODE(directory.st_mode) != _DIRECTORY_MODE
        ):
            raise ValueError("onboarding publication outbox directory is invalid")
        for path in (self.lock_path, self.emergency_lock_path):
            self._inspect_private_file(path, "onboarding publication outbox authority")
        self._inspect_private_file(self.key_path, "onboarding publication receipt key")
        self._read_key()
        self._inspect_private_file(self.ledger_path, "onboarding publication outbox ledger")
        self._inspect_private_file(self.emergency_path, "onboarding publication outbox ledger")
        if self._ledger_schema(self.ledger_path) == _LEGACY_SCHEMA:
            raise ValueError("onboarding publication outbox ledger requires explicit migration")
        self._read_primary_unlocked()
        self._read_emergency_unlocked()
        callback_exists = self._entry_exists(self.callback_path)
        callback_lock_exists = self._entry_exists(self.callback_lock_path)
        if callback_exists != callback_lock_exists:
            raise ValueError("owner callback receipt authority is incomplete")
        self._callback_authority_absent = not callback_exists
        if callback_exists:
            self._inspect_private_file(
                self.callback_lock_path, "onboarding publication outbox authority"
            )
            self._inspect_private_file(self.callback_path, "owner callback receipt ledger")
            self._read_owner_callbacks_unlocked()

    def _validate_path_ancestry(self, *, require_state_dir: bool) -> None:
        for path in (*reversed(self.state_dir.parents), self.state_dir):
            try:
                info = os.lstat(path)
            except FileNotFoundError:
                if require_state_dir and path == self.state_dir:
                    raise ValueError("onboarding publication outbox directory is unavailable") from None
                continue
            if stat.S_ISLNK(info.st_mode):
                raise ValueError("onboarding publication outbox path must not contain symlinks")

    @staticmethod
    def _entry_exists(path: Path) -> bool:
        try:
            os.lstat(path)
        except FileNotFoundError:
            return False
        return True

    def _inspect_private_file(self, path: Path, label: str) -> None:
        try:
            descriptor = os.open(
                path,
                os.O_RDONLY | os.O_CLOEXEC | getattr(os, "O_NOFOLLOW", 0) | _NOATIME,
            )
        except OSError as exc:
            raise ValueError(f"{label} is unavailable") from exc
        try:
            self._validate_private_descriptor(descriptor, label)
        finally:
            os.close(descriptor)

    def _ensure_callback_ledger(self) -> None:
        try:
            descriptor = os.open(
                self.callback_path,
                os.O_RDONLY | os.O_CLOEXEC | getattr(os, "O_NOFOLLOW", 0) | _NOATIME,
            )
        except FileNotFoundError:
            self._write_owner_callbacks_unlocked(())
            return
        except OSError as exc:
            raise ValueError("owner callback receipt ledger is unavailable") from exc
        try:
            self._validate_private_descriptor(descriptor, "owner callback receipt ledger")
        finally:
            os.close(descriptor)
        self._read_owner_callbacks_unlocked()

    def _ensure_private_file(self, path: Path) -> None:
        try:
            descriptor = os.open(
                path,
                os.O_RDWR | os.O_CREAT | os.O_EXCL | os.O_CLOEXEC | getattr(os, "O_NOFOLLOW", 0),
                _PRIVATE_MODE,
            )
        except FileExistsError:
            descriptor = os.open(
                path,
                os.O_RDWR | os.O_CLOEXEC | getattr(os, "O_NOFOLLOW", 0),
            )
        except OSError as exc:
            raise ValueError("onboarding publication outbox authority is unavailable") from exc
        try:
            self._validate_private_descriptor(descriptor, "onboarding publication outbox authority")
        finally:
            os.close(descriptor)

    def _ensure_key(self) -> None:
        try:
            descriptor = os.open(
                self.key_path,
                os.O_WRONLY | os.O_CREAT | os.O_EXCL | os.O_CLOEXEC | getattr(os, "O_NOFOLLOW", 0),
                _PRIVATE_MODE,
            )
        except FileExistsError:
            self._read_key()
            return
        except OSError as exc:
            raise ValueError("onboarding publication receipt key is unavailable") from exc
        try:
            os.write(descriptor, secrets.token_bytes(32))
            os.fsync(descriptor)
        finally:
            os.close(descriptor)

    def _read_key(self) -> bytes:
        try:
            descriptor = os.open(
                self.key_path,
                os.O_RDONLY | os.O_CLOEXEC | getattr(os, "O_NOFOLLOW", 0) | _NOATIME,
            )
        except OSError as exc:
            raise ValueError("onboarding publication receipt key is unavailable") from exc
        try:
            self._validate_private_descriptor(descriptor, "onboarding publication receipt key")
            key = os.read(descriptor, 33)
        finally:
            os.close(descriptor)
        if len(key) != 32:
            raise ValueError("onboarding publication receipt key is invalid")
        return key

    def _ensure_ledger(self, path: Path, *, emergency: bool) -> None:
        try:
            descriptor = os.open(
                path,
                os.O_RDONLY | os.O_CLOEXEC | getattr(os, "O_NOFOLLOW", 0) | _NOATIME,
            )
        except FileNotFoundError:
            if emergency:
                self._write_emergency_unlocked(())
            else:
                self._write_primary_unlocked(())
            return
        except OSError as exc:
            raise ValueError("onboarding publication outbox ledger is unavailable") from exc
        try:
            self._validate_private_descriptor(descriptor, "onboarding publication outbox ledger")
        finally:
            os.close(descriptor)
        if emergency:
            self._read_emergency_unlocked()
        elif self._ledger_schema(path) == _LEGACY_SCHEMA:
            self._quarantine_legacy_primary_ledger()
        else:
            self._read_primary_unlocked()

    def _ledger_schema(self, path: Path) -> str | None:
        try:
            descriptor = os.open(
                path,
                os.O_RDONLY | os.O_CLOEXEC | getattr(os, "O_NOFOLLOW", 0) | _NOATIME,
            )
        except OSError as exc:
            raise ValueError("onboarding publication outbox ledger is unavailable") from exc
        try:
            self._validate_private_descriptor(
                descriptor,
                "onboarding publication outbox ledger",
            )
            with os.fdopen(
                descriptor,
                "r",
                encoding="utf-8",
                closefd=False,
            ) as handle:
                value = json.load(handle)
        finally:
            os.close(descriptor)
        if not isinstance(value, Mapping):
            return None
        schema = value.get("schema")
        return schema if isinstance(schema, str) else None

    def _quarantine_legacy_primary_ledger(self) -> None:
        """Preserve unauthenticated v1 evidence without letting it authorize sends."""
        with self._locked(self.lock_path):
            schema = self._ledger_schema(self.ledger_path)
            if schema == _SCHEMA:
                return
            if schema != _LEGACY_SCHEMA:
                raise ValueError("onboarding publication outbox ledger integrity is invalid")
            try:
                info = os.lstat(self.legacy_ledger_path)
            except FileNotFoundError:
                pass
            else:
                if stat.S_ISLNK(info.st_mode):
                    raise ValueError("onboarding publication outbox legacy path is invalid")
                raise ValueError("onboarding publication outbox legacy ledger already exists")
            os.replace(self.ledger_path, self.legacy_ledger_path)
            directory = os.open(
                self.state_dir,
                os.O_RDONLY | os.O_DIRECTORY | os.O_CLOEXEC,
            )
            try:
                os.fsync(directory)
            finally:
                os.close(directory)
            self._write_primary_unlocked(())

    @contextmanager
    def _locked(self, path: Path) -> Iterator[None]:
        try:
            descriptor = os.open(
                path,
                os.O_RDWR | os.O_CLOEXEC | getattr(os, "O_NOFOLLOW", 0),
            )
        except OSError as exc:
            raise ValueError("onboarding publication outbox lock is unavailable") from exc
        try:
            self._validate_private_descriptor(descriptor, "onboarding publication outbox lock")
            fcntl.flock(descriptor, fcntl.LOCK_EX)
            yield
        finally:
            os.close(descriptor)

    @staticmethod
    def _validate_private_descriptor(descriptor: int, label: str) -> None:
        info = os.fstat(descriptor)
        if (
            not stat.S_ISREG(info.st_mode)
            or info.st_nlink != 1
            or stat.S_IMODE(info.st_mode) != _PRIVATE_MODE
        ):
            raise ValueError(f"{label} is invalid")

    def _read_primary_unlocked(self) -> tuple[GatewayPublicationReceipt, ...]:
        return self._read_records_unlocked(self.ledger_path, emergency=False)

    def _read_emergency_unlocked(self) -> tuple[GatewayPublicationReceipt, ...]:
        return self._read_records_unlocked(self.emergency_path, emergency=True)

    def _read_records_unlocked(
        self,
        path: Path,
        *,
        emergency: bool,
    ) -> tuple[GatewayPublicationReceipt, ...]:
        try:
            descriptor = os.open(
                path,
                os.O_RDONLY | os.O_CLOEXEC | getattr(os, "O_NOFOLLOW", 0) | _NOATIME,
            )
        except OSError as exc:
            raise ValueError("onboarding publication outbox ledger is unavailable") from exc
        try:
            self._validate_private_descriptor(descriptor, "onboarding publication outbox ledger")
            with os.fdopen(descriptor, "r", encoding="utf-8", closefd=False) as handle:
                value = json.load(handle)
        finally:
            os.close(descriptor)
        if not isinstance(value, Mapping) or value.get("schema") != _SCHEMA:
            raise ValueError("onboarding publication outbox ledger integrity is invalid")
        records = value.get("records")
        if not isinstance(records, list):
            raise ValueError("onboarding publication outbox ledger integrity is invalid")
        parsed = tuple(self._record_from_dict(item, emergency=emergency) for item in records)
        keys = {(record.session_id, record.generation) for record in parsed}
        if len(keys) != len(parsed):
            raise ValueError("onboarding publication outbox ledger is ambiguous")
        return parsed

    def _write_primary_unlocked(
        self,
        records: tuple[GatewayPublicationReceipt, ...] | list[GatewayPublicationReceipt],
    ) -> None:
        self._write_records_unlocked(self.ledger_path, records, emergency=False)

    def _write_emergency_unlocked(
        self,
        records: tuple[GatewayPublicationReceipt, ...] | list[GatewayPublicationReceipt],
    ) -> None:
        self._write_records_unlocked(self.emergency_path, records, emergency=True)

    def _write_records_unlocked(
        self,
        path: Path,
        records: tuple[GatewayPublicationReceipt, ...] | list[GatewayPublicationReceipt],
        *,
        emergency: bool,
    ) -> None:
        document = {
            "schema": _SCHEMA,
            "records": [self._record_to_dict(record) for record in records],
        }
        descriptor, temporary = tempfile.mkstemp(prefix=".ledger-", suffix=".tmp", dir=self.state_dir)
        try:
            os.fchmod(descriptor, _PRIVATE_MODE)
            with os.fdopen(descriptor, "w", encoding="utf-8") as handle:
                json.dump(document, handle, ensure_ascii=False, sort_keys=True, separators=(",", ":"))
                handle.write("\n")
                handle.flush()
                os.fsync(handle.fileno())
            os.replace(temporary, path)
            directory = os.open(self.state_dir, os.O_RDONLY | os.O_DIRECTORY | os.O_CLOEXEC)
            try:
                os.fsync(directory)
            finally:
                os.close(directory)
        except BaseException:
            try:
                os.unlink(temporary)
            except FileNotFoundError:
                pass
            raise

    def _owner_callback_candidate(
        self,
        *,
        session_id: str,
        customer_key: str,
        action: str,
        actor_user_id: int,
        authority: tuple[int, int, int],
        route: tuple[str, str],
        message_id: int,
        update_id: int,
        consumed_updates_before: tuple[int, ...],
        callback_data: str,
        publication: GatewayPublicationReceipt,
    ) -> GatewayOwnerCallbackReceipt:
        session = self._session_id(session_id)
        if not isinstance(customer_key, str) or not customer_key or len(customer_key) > 256:
            raise ValueError("owner callback customer key is invalid")
        if not isinstance(action, str) or not action or len(action) > 64:
            raise ValueError("owner callback action is invalid")
        if (
            type(actor_user_id) is not int
            or actor_user_id < 0
            or len(authority) != 3
            or type(authority[0]) is not int
            or authority[0] < 0
            or type(authority[1]) is not int
            or type(authority[2]) is not int
            or authority[2] < 0
            or type(update_id) is not int
            or update_id < 0
        ):
            raise ValueError("owner callback authority is invalid")
        current_route = self._route(*route)
        current_message_id = self._message_id(message_id)
        if (
            not isinstance(callback_data, str)
            or not callback_data
            or len(callback_data) > 256
            or any(type(item) is not int or item < 0 for item in consumed_updates_before)
            or len(set(consumed_updates_before)) != len(consumed_updates_before)
            or update_id in consumed_updates_before
            or not isinstance(publication.receipt_integrity, str)
        ):
            raise ValueError("owner callback event is invalid")
        values: dict[str, object] = {
            "session_id": session,
            "customer_key": customer_key,
            "action": action,
            "actor_user_id": actor_user_id,
            "authority": list(authority),
            "route": list(current_route),
            "message_id": current_message_id,
            "update_id": update_id,
            "consumed_updates_before": list(consumed_updates_before),
            "callback_data": callback_data,
            "publication_generation": publication.generation,
            "publication_dispatch_identity": publication.dispatch_identity,
            "publication_receipt_integrity": publication.receipt_integrity,
        }
        integrity = hmac.new(
            self._read_key(),
            self._canonical({"domain": "telegram-owner-approve-event-v1", **values}),
            hashlib.sha256,
        ).hexdigest()
        return GatewayOwnerCallbackReceipt(
            session,
            customer_key,
            action,
            actor_user_id,
            authority,
            current_route,
            current_message_id,
            update_id,
            consumed_updates_before,
            callback_data,
            publication.generation,
            publication.dispatch_identity,
            publication.receipt_integrity,
            integrity,
        )

    def _read_owner_callbacks_unlocked(self) -> tuple[GatewayOwnerCallbackReceipt, ...]:
        try:
            descriptor = os.open(
                self.callback_path,
                os.O_RDONLY | os.O_CLOEXEC | getattr(os, "O_NOFOLLOW", 0) | _NOATIME,
            )
        except OSError as exc:
            raise ValueError("owner callback receipt ledger is unavailable") from exc
        try:
            self._validate_private_descriptor(descriptor, "owner callback receipt ledger")
            with os.fdopen(descriptor, "r", encoding="utf-8", closefd=False) as handle:
                value = json.load(handle)
        finally:
            os.close(descriptor)
        if not isinstance(value, Mapping) or value.get("schema") != _CALLBACK_SCHEMA:
            raise ValueError("owner callback receipt ledger integrity is invalid")
        records = value.get("records")
        if not isinstance(records, list):
            raise ValueError("owner callback receipt ledger integrity is invalid")
        parsed = tuple(self._owner_callback_from_dict(item) for item in records)
        if len({item.session_id for item in parsed}) != len(parsed):
            raise ValueError("owner callback receipt is ambiguous")
        return parsed

    def _write_owner_callbacks_unlocked(
        self,
        records: tuple[GatewayOwnerCallbackReceipt, ...] | list[GatewayOwnerCallbackReceipt],
    ) -> None:
        document = {
            "schema": _CALLBACK_SCHEMA,
            "records": [self._owner_callback_to_dict(item) for item in records],
        }
        descriptor, temporary = tempfile.mkstemp(
            prefix=".owner-callbacks-", suffix=".tmp", dir=self.state_dir
        )
        try:
            os.fchmod(descriptor, _PRIVATE_MODE)
            with os.fdopen(descriptor, "w", encoding="utf-8") as handle:
                json.dump(document, handle, sort_keys=True, separators=(",", ":"))
                handle.write("\n")
                handle.flush()
                os.fsync(handle.fileno())
            os.replace(temporary, self.callback_path)
            directory = os.open(self.state_dir, os.O_RDONLY | os.O_DIRECTORY | os.O_CLOEXEC)
            try:
                os.fsync(directory)
            finally:
                os.close(directory)
        except BaseException:
            try:
                os.unlink(temporary)
            except FileNotFoundError:
                pass
            raise

    @staticmethod
    def _owner_callback_to_dict(record: GatewayOwnerCallbackReceipt) -> dict[str, object]:
        return {
            "session_id": record.session_id,
            "customer_key": record.customer_key,
            "action": record.action,
            "actor_user_id": record.actor_user_id,
            "authority": list(record.authority),
            "route": list(record.route),
            "message_id": record.message_id,
            "update_id": record.update_id,
            "consumed_updates_before": list(record.consumed_updates_before),
            "callback_data": record.callback_data,
            "publication_generation": record.publication_generation,
            "publication_dispatch_identity": record.publication_dispatch_identity,
            "publication_receipt_integrity": record.publication_receipt_integrity,
            "event_integrity": record.event_integrity,
        }

    def _owner_callback_from_dict(self, value: object) -> GatewayOwnerCallbackReceipt:
        if not isinstance(value, Mapping):
            raise ValueError("owner callback receipt is invalid")
        record: dict[str, object] = {}
        for key, item in value.items():
            if not isinstance(key, str):
                raise ValueError("owner callback receipt is invalid")
            record[key] = item
        if set(record) != {
            "session_id", "customer_key", "action", "actor_user_id", "authority", "route",
            "message_id", "update_id", "consumed_updates_before", "callback_data",
            "publication_generation", "publication_dispatch_identity",
            "publication_receipt_integrity", "event_integrity",
        }:
            raise ValueError("owner callback receipt is invalid")
        session_id = record["session_id"]
        customer_key = record["customer_key"]
        action = record["action"]
        actor_user_id = record["actor_user_id"]
        authority = record["authority"]
        route = record["route"]
        message_id = record["message_id"]
        update_id = record["update_id"]
        consumed = record["consumed_updates_before"]
        callback_data = record["callback_data"]
        generation = record["publication_generation"]
        dispatch_identity = record["publication_dispatch_identity"]
        receipt_integrity = record["publication_receipt_integrity"]
        event_integrity = record["event_integrity"]
        if (
            not isinstance(session_id, str)
            or not isinstance(customer_key, str)
            or not isinstance(action, str)
            or type(actor_user_id) is not int
            or not isinstance(authority, list)
            or len(authority) != 3
            or not isinstance(route, list)
            or len(route) != 2
            or type(message_id) is not int
            or type(update_id) is not int
            or not isinstance(consumed, list)
            or not isinstance(callback_data, str)
            or type(generation) is not int
            or not isinstance(dispatch_identity, str)
            or not isinstance(receipt_integrity, str)
            or not isinstance(event_integrity, str)
        ):
            raise ValueError("owner callback receipt is invalid")
        owner_user_id, owner_chat_id, owner_topic_id = authority
        route_chat_id, route_topic_id = route
        if (
            type(owner_user_id) is not int
            or type(owner_chat_id) is not int
            or type(owner_topic_id) is not int
            or not isinstance(route_chat_id, str)
            or not isinstance(route_topic_id, str)
        ):
            raise ValueError("owner callback receipt is invalid")
        consumed_updates: list[int] = []
        for item in consumed:
            if type(item) is not int:
                raise ValueError("owner callback receipt is invalid")
            consumed_updates.append(item)
        publication = GatewayPublicationReceipt(
            "placeholder",
            generation,
            {},
            (route_chat_id, route_topic_id),
            "owner",
            "0" * 64,
            "0" * 64,
            dispatch_identity,
            "COMMITTED",
            message_id,
            receipt_integrity,
        )
        candidate = self._owner_callback_candidate(
            session_id=session_id,
            customer_key=customer_key,
            action=action,
            actor_user_id=actor_user_id,
            authority=(owner_user_id, owner_chat_id, owner_topic_id),
            route=(route_chat_id, route_topic_id),
            message_id=message_id,
            update_id=update_id,
            consumed_updates_before=tuple(consumed_updates),
            callback_data=callback_data,
            publication=publication,
        )
        if not hmac.compare_digest(event_integrity, candidate.event_integrity):
            raise ValueError("owner callback receipt integrity is invalid")
        return candidate

    @staticmethod
    def _record_to_dict(record: GatewayPublicationReceipt) -> dict[str, object]:
        return {
            "session_id": record.session_id,
            "generation": record.generation,
            "payload": record.payload,
            "route": list(record.route),
            "role": record.role,
            "render_identity": record.render_identity,
            "payload_digest": record.payload_digest,
            "dispatch_identity": record.dispatch_identity,
            "state": record.state,
            "message_id": record.message_id,
            "receipt_integrity": record.receipt_integrity,
        }

    def _record_from_dict(
        self,
        value: object,
        *,
        emergency: bool,
    ) -> GatewayPublicationReceipt:
        if not isinstance(value, Mapping):
            raise ValueError("onboarding publication outbox record is invalid")
        record: dict[str, object] = {}
        for key, item in value.items():
            if not isinstance(key, str):
                raise ValueError("onboarding publication outbox record is invalid")
            record[key] = item
        if set(record) != {
            "session_id", "generation", "payload", "route", "role", "render_identity",
            "payload_digest", "dispatch_identity", "state", "message_id", "receipt_integrity",
        }:
            raise ValueError("onboarding publication outbox record is invalid")
        session_id = record["session_id"]
        generation = record["generation"]
        payload = record["payload"]
        route = record["route"]
        role = record["role"]
        render_identity = record["render_identity"]
        payload_digest = record["payload_digest"]
        dispatch_identity = record["dispatch_identity"]
        state = record["state"]
        message_id = record["message_id"]
        receipt_integrity = record["receipt_integrity"]
        if (
            not isinstance(session_id, str)
            or not isinstance(generation, int)
            or isinstance(generation, bool)
            or not isinstance(payload, Mapping)
            or not isinstance(route, list)
            or len(route) != 2
            or not isinstance(role, str)
            or not isinstance(render_identity, str)
            or not isinstance(payload_digest, str)
            or not isinstance(dispatch_identity, str)
            or not isinstance(state, str)
        ):
            raise ValueError("onboarding publication outbox record is invalid")
        route_first, route_second = route
        if not isinstance(route_first, str) or not isinstance(route_second, str):
            raise ValueError("onboarding publication outbox record is invalid")
        candidate = self._candidate(
            session_id=session_id,
            generation=generation,
            payload={str(key): item for key, item in payload.items() if isinstance(key, str)},
            route=(route_first, route_second),
            role=role,
            render_identity=render_identity,
        )
        if (
            not hmac.compare_digest(payload_digest, candidate.payload_digest)
            or not hmac.compare_digest(dispatch_identity, candidate.dispatch_identity)
        ):
            raise ValueError("onboarding publication authority integrity is invalid")
        allowed_states = _EMERGENCY_STATES if emergency else _STATES
        if state not in allowed_states:
            raise ValueError("onboarding publication outbox record is invalid")
        if state == "DISPATCHING":
            if message_id is not None or receipt_integrity is not None:
                raise ValueError("dispatching publication cannot have a receipt")
            return candidate
        if (
            type(message_id) is not int
            or isinstance(message_id, bool)
            or message_id <= 0
            or not isinstance(receipt_integrity, str)
            or _DIGEST.fullmatch(receipt_integrity) is None
        ):
            raise ValueError("receipted publication requires valid receipt integrity")
        received = self._with_receipt(candidate, message_id)
        if not hmac.compare_digest(receipt_integrity, received.receipt_integrity or ""):
            raise ValueError("onboarding publication receipt integrity is invalid")
        return self._with_state(received, state)

    @staticmethod
    def _matching_record(
        records: tuple[GatewayPublicationReceipt, ...] | list[GatewayPublicationReceipt],
        session_id: str,
        generation: int,
    ) -> GatewayPublicationReceipt | None:
        matches = [
            record
            for record in records
            if record.session_id == session_id and record.generation == generation
        ]
        if not matches:
            return None
        if len(matches) != 1:
            raise ValueError("onboarding publication receipt is ambiguous")
        return matches[0]

    @classmethod
    def _find(
        cls,
        records: list[GatewayPublicationReceipt],
        session_id: str,
        generation: int,
    ) -> tuple[int, GatewayPublicationReceipt]:
        for index, record in enumerate(records):
            if record.session_id == session_id and record.generation == generation:
                return index, record
        raise ValueError("onboarding publication receipt is unavailable")
