"""Registrations the CMS has not acknowledged yet (ADR-0010).

The rule the Intake Center exists to enforce is that a document is never
somewhere nobody can see it. Across a network that cannot be made absolute: the
CMS is a different database on a different machine, so no transaction spans
"bytes arrived" and "the CMS knows about them". Something has to give when the
CMS is unreachable, and the three honest options are to refuse the document, to
hold it, or to drop it.

Holding is the approved choice, and this is the part that keeps it honest:

  * it is **bounded** — past the cap the agent refuses new documents rather than
    accumulating silently, because an outbox that grows without limit is the
    hidden queue again wearing a different name;
  * its depth is a **health signal**, so a non-empty outbox is visible in
    /health/ready and in metrics rather than discovered later;
  * nothing **advances** on a held document — no reading, no proposal, no reply
    to the sender claiming it was received — until the CMS has the record.

A document therefore waits, and waiting is a state somebody can be shown. That
is the difference from what this replaces.
"""

from __future__ import annotations

import logging
from dataclasses import dataclass, replace
from datetime import UTC, datetime
from typing import Any

logger = logging.getLogger(__name__)

#: How long a claim is honoured before the work is considered abandoned.
#:
#: A lease, not a lock. It has to outlast the slowest legitimate attempt — a
#: registration carrying a 20 MB file to a struggling CMS — and expire soon
#: enough that a crash does not strand a client's document for an afternoon.
CLAIM_LEASE_SECONDS = 300

#: Past this the agent stops accepting rather than queueing. Sized for a CMS
#: restart or a short outage, not for a day of downtime: if this fills, the
#: problem is the CMS being down, and quietly hoarding a thousand documents
#: would hide that instead of surfacing it.
DEFAULT_CAPACITY = 200


class OutboxFull(RuntimeError):
    """The CMS has been unreachable long enough that we stop taking documents."""


@dataclass(frozen=True, slots=True)
class PendingRegistration:
    """One arrival the CMS has not been told about."""

    source: str
    received_at: datetime
    provider_message_id: str | None = None
    source_account: str | None = None
    original_filename: str | None = None
    mime_type: str | None = None
    size_bytes: int | None = None
    path: str | None = None
    failure_reason: str | None = None
    workflow_id: str | None = None
    attempts: int = 0
    last_error: str | None = None
    id: int | None = None
    claimed_by: str | None = None

    def waited(self, now: datetime | None = None) -> float:
        """Seconds since the document arrived — not since we last tried."""
        return ((now or datetime.now(UTC)) - self.received_at).total_seconds()

    def submission(self) -> dict[str, Any]:
        """The body to send, minus the file, which is read from disk at send time."""
        body = {
            "source": self.source,
            "source_account": self.source_account,
            "provider_message_id": self.provider_message_id,
            "original_filename": self.original_filename,
            "mime": self.mime_type,
            "size_bytes": self.size_bytes,
            "failure_reason": self.failure_reason,
            "workflow_id": self.workflow_id,
            "received_at": self.received_at.isoformat(),
        }

        return {key: value for key, value in body.items() if value is not None}


class InMemoryOutbox:
    """The fallback, for a deployment with no database.

    Honest rather than convenient: it works, and it loses everything held on
    restart — which for this table means losing exactly the documents it exists
    to protect. The warning that selects it says so.
    """

    def __init__(self, capacity: int = DEFAULT_CAPACITY) -> None:
        self._items: list[PendingRegistration] = []
        self._claims: dict[int, tuple[str, datetime]] = {}
        self._capacity = capacity

    def add(self, registration: PendingRegistration) -> None:
        if registration.provider_message_id is not None and any(
            item.source == registration.source
            and item.provider_message_id == registration.provider_message_id
            for item in self._items
        ):
            return

        if len(self._items) >= self._capacity:
            raise OutboxFull(f"The intake outbox is full ({self._capacity}).")

        self._items.append(registration)

    def oldest(self) -> PendingRegistration | None:
        return next(iter(self.all()), None)

    def claim(self, worker: str, lease: float = CLAIM_LEASE_SECONDS) -> PendingRegistration | None:
        """Take the oldest unclaimed row, or one whose claim has expired.

        The in-memory store has one worker by definition, so this cannot race —
        but it implements the same contract as the durable one, because a
        deployment without a database must not behave differently in a way that
        only shows up when somebody moves between them.
        """
        now = datetime.now(UTC)

        for index, item in enumerate(sorted(self._items, key=lambda i: i.received_at)):
            held = self._claims.get(id(item))

            if held is not None and (now - held[1]).total_seconds() < lease:
                continue

            self._claims[id(item)] = (worker, now)

            return replace(item, claimed_by=worker)

        return None

    def release(self, registration: PendingRegistration) -> None:
        """Give the work back without counting it done."""
        for item in self._items:
            if item == registration or (
                item.source == registration.source
                and item.provider_message_id == registration.provider_message_id
            ):
                self._claims.pop(id(item), None)

    def all(self) -> list[PendingRegistration]:
        # Sorted by arrival, not by insertion, so this store answers in the same
        # order as the Postgres one. Two implementations of one interface that
        # disagree about ordering produce a bug that only appears on whichever
        # deployment has no database.
        return sorted(self._items, key=lambda item: item.received_at)

    def done(self, registration: PendingRegistration) -> None:
        kept = []

        for item in self._items:
            if (
                item.source == registration.source
                and item.provider_message_id == registration.provider_message_id
                and item.received_at == registration.received_at
            ):
                self._claims.pop(id(item), None)

                continue

            kept.append(item)

        self._items = kept

    def failed(self, registration: PendingRegistration, error: str) -> None:
        # Matched on identity rather than equality: claim() hands back a copy
        # carrying claimed_by, so a full comparison finds nothing and the attempt
        # goes uncounted — which reads as a queue that never retries.
        def same(item: PendingRegistration) -> bool:
            return (
                item.source == registration.source
                and item.provider_message_id == registration.provider_message_id
                and item.received_at == registration.received_at
            )

        updated = []

        for item in self._items:
            if same(item):
                # Released as well as counted. The lease exists for workers that
                # died, not for ones that failed and said so.
                self._claims.pop(id(item), None)
                item = replace(item, attempts=item.attempts + 1, last_error=error)

            updated.append(item)

        self._items = updated

    def __len__(self) -> int:
        return len(self._items)

    @property
    def capacity(self) -> int:
        return self._capacity


class PostgresOutbox:
    """The durable one.

    A restart is the case this exists for. An in-flight document that lived only
    in a local variable was lost when the process died — silently, with no record
    anywhere that it had ever arrived. Written down first, it survives.
    """

    def __init__(self, connect, capacity: int = DEFAULT_CAPACITY) -> None:
        # A callable, matching PostgresRunStore and PostgresPendingPdfs: built at
        # startup, possibly before the database is reachable.
        self._connect = connect
        self._capacity = capacity

    def add(self, registration: PendingRegistration) -> None:
        if len(self) >= self._capacity:
            raise OutboxFull(f"The intake outbox is full ({self._capacity}).")

        with self._connect() as connection, connection.cursor() as cursor:
            cursor.execute(
                """
                INSERT INTO intake_outbox
                    (source, source_account, provider_message_id, original_filename,
                     mime_type, size_bytes, path, failure_reason, workflow_id, received_at)
                VALUES (%s, %s, %s, %s, %s, %s, %s, %s, %s, %s)
                ON CONFLICT (source, provider_message_id)
                    WHERE provider_message_id IS NOT NULL
                    DO NOTHING
                """,
                (
                    registration.source,
                    registration.source_account,
                    registration.provider_message_id,
                    registration.original_filename,
                    registration.mime_type,
                    registration.size_bytes,
                    registration.path,
                    registration.failure_reason,
                    registration.workflow_id,
                    registration.received_at,
                ),
            )
            connection.commit()

    def oldest(self) -> PendingRegistration | None:
        with self._connect() as connection, connection.cursor() as cursor:
            cursor.execute(
                """
                SELECT id, source, source_account, provider_message_id, original_filename,
                       mime_type, size_bytes, path, failure_reason, workflow_id,
                       received_at, attempts, last_error
                FROM intake_outbox
                ORDER BY received_at, id
                LIMIT 1
                """
            )

            return _from_row(cursor.fetchone())

    def claim(self, worker: str, lease: float = CLAIM_LEASE_SECONDS) -> PendingRegistration | None:
        """Take exactly one row, and make sure nobody else has it.

        One statement, so there is no window between choosing a row and marking
        it taken — which is where the duplicate came from. FOR UPDATE SKIP
        LOCKED means a second worker running the same statement at the same
        instant selects a *different* row rather than waiting for this one, so
        concurrency scales instead of serialising.

        The lease is what makes a claim safe to take. A worker that dies holding
        one does not strand the document: the claim ages out and the next worker
        picks it up, which is what preserves at-least-once delivery. Combined
        with the idempotency key the row already carries, a redelivery after an
        expired claim is harmless — the CMS returns the record it already made.
        """
        with self._connect() as connection, connection.cursor() as cursor:
            cursor.execute(
                """
                UPDATE intake_outbox
                SET claimed_at = now(), claimed_by = %s
                WHERE id = (
                    SELECT id FROM intake_outbox
                    WHERE claimed_at IS NULL
                       OR claimed_at < now() - (%s * INTERVAL '1 second')
                    ORDER BY received_at, id
                    LIMIT 1
                    FOR UPDATE SKIP LOCKED
                )
                RETURNING id, source, source_account, provider_message_id, original_filename,
                          mime_type, size_bytes, path, failure_reason, workflow_id,
                          received_at, attempts, last_error, claimed_by
                """,
                (worker, lease),
            )
            row = cursor.fetchone()
            connection.commit()

            return _from_row(row)

    def release(self, registration: PendingRegistration) -> None:
        """Hand the work back without counting an attempt against it.

        For the caller that decides not to proceed after claiming — a shutdown,
        say. A failed attempt uses failed() instead, which keeps the claim's
        expiry doing the waiting.
        """
        if registration.id is None:
            return

        with self._connect() as connection, connection.cursor() as cursor:
            cursor.execute(
                "UPDATE intake_outbox SET claimed_at = NULL, claimed_by = NULL WHERE id = %s",
                (registration.id,),
            )
            connection.commit()

    def all(self) -> list[PendingRegistration]:
        with self._connect() as connection, connection.cursor() as cursor:
            cursor.execute(
                """
                SELECT id, source, source_account, provider_message_id, original_filename,
                       mime_type, size_bytes, path, failure_reason, workflow_id,
                       received_at, attempts, last_error
                FROM intake_outbox
                ORDER BY received_at, id
                """
            )

            return [row for row in (_from_row(r) for r in cursor.fetchall()) if row]

    def done(self, registration: PendingRegistration) -> None:
        """Delete only after the CMS has acknowledged it.

        The order matters and is the whole guarantee: acknowledged, then deleted.
        Deleting first would reopen the window this table closes.
        """
        if registration.id is None:
            return

        with self._connect() as connection, connection.cursor() as cursor:
            cursor.execute("DELETE FROM intake_outbox WHERE id = %s", (registration.id,))
            connection.commit()

    def failed(self, registration: PendingRegistration, error: str) -> None:
        if registration.id is None:
            return

        with self._connect() as connection, connection.cursor() as cursor:
            cursor.execute(
                """
                UPDATE intake_outbox
                SET attempts = attempts + 1,
                    last_error = %s,
                    last_tried_at = now(),
                    first_tried_at = COALESCE(first_tried_at, now()),
                    -- Released, not held. The next pass should be free to try
                    -- again immediately; the lease exists for workers that died,
                    -- not for ones that failed and said so.
                    claimed_at = NULL,
                    claimed_by = NULL
                WHERE id = %s
                """,
                (error[:1000], registration.id),
            )
            connection.commit()

    def __len__(self) -> int:
        with self._connect() as connection, connection.cursor() as cursor:
            cursor.execute("SELECT count(*) FROM intake_outbox")
            row = cursor.fetchone()

            return int(row[0]) if row else 0

    @property
    def capacity(self) -> int:
        return self._capacity


def _from_row(row) -> PendingRegistration | None:
    if row is None:
        return None

    return PendingRegistration(
        claimed_by=row[13] if len(row) > 13 else None,
        id=row[0],
        source=row[1],
        source_account=row[2],
        provider_message_id=row[3],
        original_filename=row[4],
        mime_type=row[5],
        size_bytes=row[6],
        path=row[7],
        failure_reason=row[8],
        workflow_id=row[9],
        received_at=row[10],
        attempts=row[11] or 0,
        last_error=row[12],
    )
