"""The entry point.

    python -m app             run the daemon
    python -m app check       verify the deployment and exit
    python -m app migrate     apply pending schema changes and exit

`check` is what a deployment script runs before switching traffic, and what a
person runs when something is wrong. It makes exactly the same assertions the
daemon makes at boot, so its answer means something.
"""

from __future__ import annotations

import argparse
import logging
import os
import re
import socket
import sys
import uuid
from pathlib import Path
from datetime import UTC, datetime, timedelta

from app.config.settings import ConfigurationError
from app.documents import classifier, router
from app.documents.registry import ProcessingStrategy
from app.documents.signals import Signals
from app.runtime import preflight
from app.runtime.container import Container
from app.runtime.daemon import Daemon
from app.runtime import status_report
from app.runtime.server import HttpServer, health_routes, metrics_route

logger = logging.getLogger("taxpilot")

EXIT_OK = 0
EXIT_FAILED = 1
EXIT_MISCONFIGURED = 78  # EX_CONFIG: "you configured this wrong", not "it broke"

#: How often to re-ask the CMS what it speaks.
#:
#: Fifteen minutes. A contract change arrives with a major CMS release, which
#: is not a frequent event — and the readiness probe refreshes the verdict
#: anyway wherever one is configured. This is the floor for a deployment that
#: nothing else is watching.
COMPATIBILITY_INTERVAL = 15 * 60

#: How often to report this deployment's own health to the CMS.
#:
#: Five minutes. The CMS treats a report older than thirty minutes as stale
#: and shows nothing rather than a comfortable lie, so this leaves room for a
#: few failed attempts before an operator sees the state go unknown.
STATUS_INTERVAL = 5 * 60

#: How often to collect queued commands.
#:
#: Five seconds. Each poll is one request that usually answers 204, and the
#: thing on the other end is a person watching a page for a QR code — a slower
#: loop would be felt directly. Everything else in this daemon runs on minutes.
COMMAND_INTERVAL = 5

#: How often to take delivery of whatever the webhook queued.
#:
#: Two seconds. The webhook only queues — OCR inside a request would exhaust the
#: provider's patience and earn a retry of the document already being read — so
#: this is the other half of that split, and the gap between the two is latency
#: somebody waiting on their document feels directly.
INBOX_INTERVAL = 2


def configure_logging(level: str) -> None:
    """Logs to stdout, which is where a container's logs are collected from.

    No file handler on purpose. A process writing its own log files inside a
    container writes them into a layer nobody reads and that vanishes on
    restart.
    """
    logging.basicConfig(
        level=getattr(logging, level.upper(), logging.INFO),
        format="%(asctime)s %(levelname)-8s %(name)s: %(message)s",
        stream=sys.stdout,
    )

    # These are chatty at INFO and say nothing an operator needs.
    for noisy in ("urllib3", "PIL", "paddle", "paddlex"):
        logging.getLogger(noisy).setLevel(logging.WARNING)


def run_checks(container: Container) -> bool:
    checks = [
        # Order matters: check_cms performs the handshake that
        # check_compatibility then judges.
        preflight.check_cms(container.cms),
        preflight.check_compatibility(container.cms),
        preflight.check_database(container.connect),
        preflight.check_ocr(container.ocr),
        # After check_cms, which establishes that the CMS is reachable at all —
        # otherwise this reports "types unverified" for a connection problem
        # already named on the line above.
        preflight.check_document_types(container.cms),
        # Reported at boot so a deployment that configured a public endpoint
        # learns its model stage is off, rather than running for weeks in the
        # belief that a stage which has never once fired is working.
        preflight.check_model(container),
    ]

    return preflight.run(checks)


def command_check(container: Container) -> int:
    return EXIT_OK if run_checks(container) else EXIT_FAILED


def command_migrate(container: Container) -> int:
    if container.connect is None:
        logger.error("No TAXPILOT_DATABASE_URL is set; there is nothing to migrate.")

        return EXIT_MISCONFIGURED

    from app.database.migrator import Migrator

    applied = Migrator(container.connect).apply()

    if applied:
        logger.info("Applied: %s", ", ".join(applied))
    else:
        logger.info("Schema is already up to date.")

    return EXIT_OK


def command_run(container: Container) -> int:
    """Start the daemon.

    Migrations run first, deliberately. A deployment that starts against a schema
    it does not match fails on the first write — halfway through handling a real
    document, rather than at boot where it belongs.
    """
    if container.connect is not None:
        command_migrate(container)

    if not run_checks(container):
        logger.error("Refusing to start. Fix the errors above.")

        return EXIT_MISCONFIGURED

    daemon = Daemon()
    daemon.install_signal_handlers()

    daemon.add(
        "decisions",
        lambda: _poll_decisions(container),
        container.poll_seconds,
    )

    # The CMS updates itself over the air, so the installation this is attached
    # to can change contract while this process is running. A readiness probe
    # would notice — but a deployment without an orchestrator has nothing
    # probing it, and would go on submitting into an API it no longer
    # understands until somebody restarted it.
    daemon.add("compatibility", lambda: _recheck_compatibility(container), COMPATIBILITY_INTERVAL)

    # The CMS has no other way to know any of this — see
    # app/runtime/status_report.py. Reported early and often enough that the
    # operations screen is showing a live picture rather than a memory.
    daemon.add("status", lambda: status_report.send(container, daemon), STATUS_INTERVAL)

    # Work the customer asked for from their own dashboard — connect WhatsApp,
    # disconnect, re-check. Polled far more often than anything else here
    # because a QR code answers somebody stood at a screen: at the decision
    # interval they would wait half a minute to see it.
    daemon.add("commands", lambda: _process_commands(container), COMMAND_INTERVAL)

    # Documents the customer put in their own self-chat. The webhook queued
    # them and returned; this is where the reading actually happens.
    daemon.add("inbox", lambda: _process_inbox(container), INBOX_INTERVAL)

    # Arrivals the CMS has not been told about yet, because it was unreachable
    # when they came in. Retried often: every document in here is one somebody
    # sent and cannot yet see, so the interval is the length of time the firm is
    # kept in the dark rather than a background tidy-up.
    daemon.add("intake_outbox", lambda: _drain_outbox(container), INBOX_INTERVAL)

    # Before any new work: documents that were mid-read when this agent last
    # stopped are still shown as Processing by the CMS, and nothing is processing
    # them. Reconciled once, at boot, so the queue starts honest.
    _recover_in_flight(container)

    # Hourly, not continuously. This exists to catch drift between two databases
    # that have no transaction between them — drift accumulates slowly, and a
    # check that runs constantly is a check whose own failures become the noise.
    daemon.add("reconcile", lambda: _reconcile(container), 60 * 60)

    # Runs waiting on a proposal the CMS no longer has. Hourly and bounded: it
    # costs one CMS call per run checked, and a queue of orphans is a symptom of
    # something that already happened rather than an emergency.
    daemon.add("orphan_runs", lambda: _retire_orphan_runs(container), 60 * 60)

    if container.connect is not None:
        # Once a day is ample: this only moves finished runs out of the working
        # set, and nothing waits on it.
        daemon.add("archive", lambda: _archive(container), 24 * 60 * 60)

    server = HttpServer(container.http_host, container.http_port)
    health_routes(server, container, daemon)
    metrics_route(server, container)

    if container.meta_config is not None:
        # Only Meta has a webhook. Evolution posts to whatever URL its instance
        # was configured with, which is the same path — but the signature and
        # verification handshake are Meta's, so the route is registered only
        # when there are Meta credentials to verify against.
        from app.whatsapp.webhook import meta_webhook_routes

        meta_webhook_routes(server, container.whatsapp, container.receiver, container.inbound_queue)
        logger.info("WhatsApp webhook listening at /webhook/whatsapp")
    else:
        # Evolution. It signs nothing, so the route authenticates with a shared
        # secret and registers only when one is configured — see
        # app/whatsapp/webhook.py.
        from app.whatsapp.webhook import evolution_webhook_routes

        evolution_webhook_routes(
            server, container.receiver, container.inbound_queue, container.evolution_webhook_secret
        )

        if container.evolution_webhook_secret:
            logger.info("Evolution webhook listening at /webhook/evolution")

    try:
        server.start()
    except OSError as exc:
        # A taken port is worth failing on rather than running blind: without
        # health endpoints an orchestrator cannot tell a wedged process from a
        # busy one, and will leave a dead container in rotation.
        logger.error("Could not bind %s:%d — %s", container.http_host, container.http_port, exc)

        return EXIT_MISCONFIGURED

    try:
        daemon.run()
    finally:
        # In a finally so a crash in the loop still releases the port; otherwise
        # a restart finds it taken and fails for a reason unrelated to the fault.
        server.stop()

    return EXIT_OK


def _poll_decisions(container: Container) -> None:
    outcomes = container.poller.poll()

    for outcome in outcomes:
        logger.info(
            "Proposal %s (%s): %s", outcome.proposal_id, outcome.status, outcome.detail
        )


def _process_commands(container: Container) -> None:
    """Collect and run whatever the CMS has queued.

    The processor is built once and kept: registering handlers on every pass
    would rebind the provider each time, and the provider holds the Evolution
    connection.
    """
    processor = getattr(container, "_command_processor", None)

    if processor is None:
        from app.whatsapp import handlers
        from app.workflow.commands import CommandProcessor

        processor = CommandProcessor(container.cms)
        handlers.register(processor, container.whatsapp)
        container._command_processor = processor  # noqa: SLF001

    handled = processor.poll()

    if handled:
        logger.info("Handled %d command(s).", len(handled))


def _process_inbox(container: Container) -> None:
    """Read what the webhook queued.

    The webhook queues and returns; this is the half that takes seconds. Split
    that way because OCR inside a request outlasts the provider's patience and
    earns a retry of the very document already being read.

    Every message here has already passed the self-chat rule — the receiver
    applies it before anything is queued — so this does not re-decide who may
    send documents. It decides only what to do with one.

    PDFs WAIT; PICTURES DO NOT

    A PDF is held, not processed. WhatsApp offers no caption box when a document
    is forwarded — the message arrives with the field absent rather than empty —
    so for the case this feature exists to serve there is no caption to read,
    and reading one would only ever pick up a caption that came from somewhere
    else. The identifier arrives as the next plain text message instead, and the
    PDF is downloaded and queued until it does.

    A photograph is processed at once, exactly as before. An image can carry a
    caption, its content is what identifies the client, and making somebody send
    a second message for every snapshot of a CNIC would be a worse workflow than
    the one already working.

    TEXT IS NOW WORK

    Text used to be the owner talking to themselves and was ignored. While a PDF
    is waiting it is the answer to "whose is this", so it consumes the oldest
    waiting document. With nothing waiting it is ignored exactly as before.
    """
    from app.documents import identifiers
    from app.workflow import catalogue
    from app.workflow.bank_statement_intake import BANK_STATEMENT_INTAKE

    # First, because a document nobody ever named should not wait behind
    # whatever arrived tonight.
    _file_pdfs_nobody_named(container)

    messages = container.inbound_queue.drain()

    if not messages:
        return

    # Every arrival is announced before any is read (ADR-0010: a document is
    # never in a state nobody can see).
    #
    # Announcement used to happen inside the processing loop, at the moment each
    # document's turn came — so four documents sent together appeared in the
    # Intake Center one at a time, minutes apart, while the other three sat in
    # this process's memory. Observed on 2026-08-04: a person sent four, saw
    # one, and asked where the rest were. "In the agent's RAM" is exactly the
    # invisibility this module was built to end.
    #
    # Queued is what the waiting ones say, and it is the honest state: the CMS
    # holds the document and its file, nothing has started, and the tab that
    # used to read zero forever now answers "did it arrive?" the moment the
    # answer is yes.
    for message in messages:
        _announce(container, message)

    for message in messages:
        try:
            if not message.has_document:
                _metadata_arrived(container, message)

                continue

            if _TAKEN.get(message.provider_message_id) == "":
                # The gate already settled this one during the announcement
                # pass — refused with its reason, failed to download, or held
                # in the outbox for a CMS that was unreachable. All of it is
                # recorded; classifying it now would only refuse it a second
                # time and double every counter on the operations page.
                continue

            # Classified without the caption first, because the caption may only
            # be read once the strategy is known and the strategy comes from the
            # classification. The first pass costs nothing — filename and media
            # type, no file opened — and settles the one question that has to be
            # answered before anything else: may this document be read at all?
            classification = _classify_message(container, message, use_caption=False)

            if _wrong_format_for_its_type(container, message, classification):
                continue

            strategy = _strategy_for(container, message, classification)

            if strategy is ProcessingStrategy.USER_METADATA_REQUIRED:
                _hold_for_metadata(container, message, classification)

                continue

            # The strategy reads, so the caption is allowed to speak — about what
            # the document is, and about whose it is. It is deliberately NOT
            # allowed to change the answer above.
            #
            # Letting it would be incoherent: a caption reading "Bank Statement
            # for file 1420" would classify the document into a type whose rule
            # is that captions are ignored, and the identifier in that same
            # caption would then be thrown away and a second message demanded
            # for something already answered. Whether a document may be opened
            # is settled without the caption; what it is, and whose, may be
            # refined by it. The never-read guarantee is unaffected either way,
            # because it belongs to the type and is enforced by the workflow.
            #
            # A PDF is the exception, and the rule follows the FORMAT rather
            # than the branch. WhatsApp offers no caption box when forwarding a
            # document, so a caption on one came from whatever chat it was
            # forwarded out of — another client's, most dangerously. That was
            # true when unknown PDFs were held and is still true now they are
            # read; tying it to the branch is what let it lapse the moment the
            # branch changed.
            classification = _classify_message(container, message, use_caption=not _is_pdf(message))
            destination = router.route(classification)

            logger.info(
                "Message %s classified as %s (%s, %.0f%%) — %s",
                message.provider_message_id,
                classification.label,
                classification.method.value,
                classification.confidence * 100,
                destination.workflow,
            )

            if destination.workflow == BANK_STATEMENT_INTAKE.name:
                # The one workflow with its own download rules — it validates
                # the envelope and never opens the file.
                _file_identified(
                    container,
                    message,
                    identifiers.read(message.text),
                    BANK_STATEMENT_INTAKE,
                    classification,
                )

                continue

            # The same gate the forwarding paths use. It was on those only, so a
            # document arriving here was fetched and read with none of the
            # checks applied — no type list, no size cap, and nothing asking
            # whether a file claiming to be a PDF actually was one.
            path = _take_in(container, message)

            if path is None:
                continue

            workflow = catalogue.require(destination.workflow)
            found = identifiers.read(message.text)

            logger.info("Reading %s from the self-chat for %s.", path.name, workflow.name)

            # Announced before the reading starts, not after. OCR takes minutes,
            # and a record sitting at Received while a workflow is halfway
            # through it is the CMS showing something that stopped being true.
            _push(
                container,
                _reference(message),
                "processing",
                reason=f"Reading the document for {workflow.name}.",
                workflow_name=workflow.name,
                current_step="read",
            )

            run = container.engine.start(
                workflow,
                {
                    "path": str(path),
                    "sender": message.sender,
                    "workflow_name": workflow.name,
                    # So the proposal can adopt the file the CMS already
                    # stored at registration rather than uploading it again.
                    "intake_reference": _reference(message),
                    # What the caption named, so an extraction workflow prefers
                    # a client somebody typed over one printed on the document.
                    "file_number": found.file_number,
                    "cnic": found.cnic,
                    "document_type": classification.filing_type,
                    # Deliberately NOT set, so a document that names nobody is
                    # proposed rather than held.
                    #
                    # It was set here for one day. The rule read well — ask the
                    # person who sent it rather than hand a reviewer an empty
                    # client field — and in use it meant documents vanished: 31
                    # on the live installation in a single working morning, and
                    # three more on staging within minutes, each invisible until
                    # somebody happened to send an identifier or a day elapsed.
                    #
                    # A proposal a reviewer must attach a client to is visible
                    # work. A document waiting silently is work nobody knows
                    # exists. The queue is where an unidentified document
                    # belongs.
                    "may_wait_for_metadata": False,
                    # Level 1 throughout: a method, a score and a sentence. It
                    # reaches the reviewer's screen and answers the first
                    # question asked when a type looks wrong.
                    "classification": {
                        "method": classification.method.value,
                        "confidence": classification.confidence,
                        "reason": classification.reason,
                        "refinement": classification.refinement,
                        "label": classification.label,
                        # Whether the classifier opened the file. The
                        # workflow cannot work this out afterwards, and
                        # a proposal claiming "never read" about a
                        # document that was read is a false claim on the
                        # screen where trust is decided.
                        "was_read": classification.was_read,
                    },
                    "source": {
                        "channel": "whatsapp",
                        "message_id": message.provider_message_id,
                        "number": message.sender,
                        "received_at": message.received_at.isoformat(),
                        "file_name": (message.media.filename if message.media else "") or path.name,
                        "file_size": path.stat().st_size,
                        "mime_type": ((message.media.mime_type if message.media else "") or "").lower(),
                    },
                },
            )

            _push_outcome(container, _reference(message), run, workflow.name,
                          forget=message.provider_message_id)

            if _stopped_to_ask(run):
                # Read, and it named nobody. The file is already downloaded, so
                # this only queues what is on disk — the next plain text message
                # will name it and the run starts again with the identifier.
                _queue_after_reading(container, message, classification, path)

                continue

            logger.info("%s %s finished: %s", workflow.name, run.id, run.state.value)
        except Exception:  # noqa: BLE001
            # One document failing must not lose the rest of the batch, and must
            # not take the daemon down. The run itself is persisted with its
            # failure, so this is not the only record of it.
            logger.exception(
                "Could not process message %s.", message.provider_message_id
            )
        finally:
            # The take-in cache is per-batch bookkeeping: it exists so the
            # announcement pass and the processing pass agree about one
            # message without downloading or registering twice. Once this
            # message's turn is over, nothing asks again.
            _TAKEN.pop(message.provider_message_id, None)


#: How long a PDF waits for the message that names its client before it is sent
#: for manual review instead.
#:
#: The rule says "wait for the next text", which on its own is unbounded: a PDF
#: forwarded and then forgotten would sit in the queue for ever, and the next
#: text — a week later, about something else — would be read as its identifier.
#: A day is long enough to cover somebody being interrupted mid-thought and
#: short enough that the document reaches a person while anybody still
#: remembers sending it.
_DEFAULT_PDF_METADATA_WAIT_MINUTES = 24 * 60


def _pdf_metadata_wait() -> timedelta:
    raw = os.getenv("TAXPILOT_PDF_METADATA_WAIT_MINUTES", "")

    try:
        minutes = int(raw)
    except ValueError:
        minutes = _DEFAULT_PDF_METADATA_WAIT_MINUTES

    return timedelta(minutes=max(minutes, 1))


def _is_pdf(message) -> bool:
    """Whether this is the kind of document that waits.

    The declared type decides it. The filename is consulted only when the
    provider stated no type at all — a message that says `image/jpeg` while
    being called `.pdf` is a JPEG, and treating it as a PDF would refuse it at
    the magic-bytes check rather than reading it as the picture it is.
    """
    media = message.media

    if media is None:
        return False

    mime = (media.mime_type or "").lower().split(";")[0].strip()

    if mime:
        return mime == "application/pdf"

    return (media.filename or "").lower().endswith(".pdf")


def _wrong_format_for_its_type(container, message, classification) -> bool:
    """Refuse a document in a format its own type does not accept.

    The registry's `allowed_file_types` made real. It is checked here, once,
    before anything is downloaded or opened — a type that says it is only ever a
    PDF should not have a photograph fetched, read and filed against it.

    Counted and logged rather than answered, like every other refusal on this
    path: there is no reply to the sender in this version, so the count on the
    operations console is the whole account of it.
    """
    from app.documents import registry

    name = (message.media.filename if message.media else "") or ""
    suffix = _FORWARDABLE_MIMES.get((message.media.mime_type or "").lower()) or name

    if registry.permits_format(classification.filing_type, suffix):
        return False

    container.inbox.forwards_rejected += 1
    logger.info(
        "Refused a %s: %s is not a format that type accepts.",
        classification.label,
        (suffix or "the file").lstrip("."),
    )

    # Settled on the record, not only in a counter.
    #
    # Before arrivals were announced up front, this refusal happened before any
    # registration existed, so the document was invisible — a counter moved and
    # nothing else. Now the record exists (and says Queued), and a refusal that
    # leaves it there has invented a document that waits forever for a turn
    # that already came.
    reason = (
        f"A {classification.label} cannot arrive as "
        f"{(suffix or 'a file of this format').lstrip('.')}."
    )

    if _reference(message) is not None:
        _push(container, _reference(message), "failed", reason=reason)
    else:
        # The announcement failed or never ran; register the refusal the way
        # the gate always has, so the arrival is recorded either way.
        _register(container, message, refused=reason)

    return True


def _strategy_for(container, message, classification):
    """How this document's client will be worked out.

    **Unknown is not a strategy.** A document the classifier could not name has
    not earned a way of being handled, and giving "unknown" its own entry in the
    routing table would freeze a guess into the design. It is a temporary state,
    resolved here from the one thing that is known — the format — and re-decided
    against the real type the moment there is one.
    """
    route = router.route(classification)

    if route.strategy is not None:
        return route.strategy

    # Format no longer decides this.
    #
    # An unrecognised PDF was held unread, reasoning that it is more often a
    # statement than anything else and that opening one on speculation was the
    # risk worth avoiding. In use it cost more than it saved: a scanner's
    # filename — "CamScanner 07-06-2026 09.33.pdf" — names nothing, so ordinary
    # paperwork stopped and waited for a message nobody knew to send, while OCR
    # read comparable documents two thousand characters at a time.
    #
    # So an unrecognised document of either kind is read, and reaches a person
    # either way: filed against a client when the reading names one, and in the
    # approval queue with its candidates when it does not.
    #
    # The residual risk, stated rather than buried: a bank statement whose
    # filename nothing recognises will now be opened once. The filename patterns
    # keep that rare, and a statement that IS recognised still goes to the
    # workflow that never opens it at all.
    return ProcessingStrategy.OCR_WITH_METADATA_FALLBACK


def _stopped_to_ask(run) -> bool:
    """Whether the run read the document, named nobody, and stopped short.

    Read from the filing step's own record rather than from the run state,
    because "no proposal was submitted" is the fact that matters and the step is
    where it is recorded. A run that failed, or that proposed, both leave this
    False.
    """
    from app.workflow.state import StepState

    return any(
        step.name == "file" and step.state is StepState.SKIPPED
        for step in getattr(run, "steps", []) or []
    )


def _queue_after_reading(container, message, classification, path) -> None:
    """Queue a document that has already been read and downloaded.

    The other entry to the same queue. `_hold_for_metadata` puts a document in
    before anything opens it; this puts one in after reading failed to name a
    client, which is the fallback half of the hybrid rule. Both leave the same
    thing waiting for the same message.
    """
    from app.whatsapp.pending import PendingPdf

    destination = router.route(classification)

    container.pending_pdfs.add(
        PendingPdf(
            message_id=message.provider_message_id,
            path=str(path),
            sender=message.sender,
            received_at=message.received_at,
            file_name=(message.media.filename if message.media else "") or path.name,
            mime_type=((message.media.mime_type if message.media else "") or "").lower(),
            size_bytes=path.stat().st_size if path.exists() else 0,
            # Without this the hourly reconcile cannot tell the agent is holding
            # this document, and marks it Failed — while it is being held and
            # waiting exactly as designed. This is the main route into Waiting for
            # Information, so the effect was to destroy that feature's state once
            # an hour, on documents nothing was wrong with.
            intake_reference=_reference(message),
            classification={
                "workflow": destination.workflow,
                "filing_type": classification.filing_type,
                "method": classification.method.value,
                "confidence": classification.confidence,
                "reason": classification.reason,
                "refinement": classification.refinement,
                "label": classification.label,
                "was_read": classification.was_read,
            },
        )
    )

    logger.info(
        "Read %s and it named nobody; waiting for a message to say whose it is.",
        path.name,
    )


#: What _take_in already decided about a message, for the length of one batch.
#:
#: The announcement pass and the processing pass both go through _take_in for
#: the same message, and the second caller must get the first caller's answer —
#: not a second download, a second refusal counter, or a second registration.
#: A path string records success; "" records a refusal (terminal, already
#: registered with its reason). Entries are popped when the message's turn in
#: the processing loop ends, so the dict never outlives a batch.
_TAKEN: dict[str, str] = {}


def _announce(container, message) -> None:
    """Register an arrival with the CMS before its turn to be read comes.

    Best-effort by definition: failing here costs earliness, never the
    document. If this succeeds, the processing loop's own _take_in finds the
    cached answer and the record is simply visible sooner. If it fails, the
    cache stays empty and every pre-announcement path runs exactly as it always
    did — registration at processing time, through the same single gate.
    """
    if not message.has_document:
        return

    try:
        if _take_in(container, message) is None:
            # Refused at the gate. The refusal registered itself with its
            # reason, which is terminal — there is no queue to join.
            return

        _push(
            container,
            _reference(message),
            "queued",
            reason="Waiting its turn to be read.",
        )
    except Exception:  # noqa: BLE001
        if _TAKEN.get(message.provider_message_id) == "":
            # The gate settled it — counted, registered with its reason — and
            # then re-raised so the traceback survives. Logged here exactly the
            # way the processing loop always logged it, because the traceback
            # carries the provider's HTTP status and that is the diagnosis; an
            # operator greps for one line, not one line per pass.
            logger.exception(
                "Could not process message %s.", message.provider_message_id
            )
        else:
            # Failed before the gate could settle anything. Nothing is recorded
            # yet, so the processing pass gets its ordinary turn.
            logger.warning(
                "Could not announce %s on arrival; it will be registered when its turn comes.",
                message.provider_message_id,
                exc_info=True,
            )


def _take_in(container, message):
    """Fetch a document onto disk, or refuse it and say why.

    The single gate every route goes through. The checks are the same ones the
    forwarding path has always made — the declared type, the size the CMS will
    accept, and whether the bytes begin the way that type must — and they were
    on that path only, so a document arriving any other way was fetched and read
    with none of them applied.

    Each of them protects something different. The type list keeps a `.webp` out
    of a queue where a reviewer has to look at the thing. The cap stops a 40 MB
    file being pulled into memory, expanded to base64 and refused by the CMS at
    the far end. The magic bytes are the one that matters most: a file claiming
    to be a PDF and not being one is what turns "read this document" into
    running an unknown binary through a parser.

    Returns the path, or None when the document was refused. Counts either way,
    because a refusal that moves no number is indistinguishable from nothing
    having arrived.
    """
    cached = _TAKEN.get(message.provider_message_id)

    if cached is not None:
        # Already taken in this batch — by the announcement pass, usually. The
        # counters moved and the registration happened exactly once, whichever
        # pass got here first; "" means it was refused, and a refusal does not
        # become less refused by being asked about again.
        return Path(cached) if cached else None

    inbox = container.inbox
    suffix = _FORWARDABLE_MIMES.get((message.media.mime_type or "").lower())

    if suffix is None:
        inbox.forwards_rejected += 1
        reason = (
            f"Only PDF, JPG and PNG can be filed this way; this was "
            f"{message.media.mime_type or 'a file of unstated type'}."
        )
        logger.info("Refused a forwarded %s: only PDF, JPG and PNG can be filed this way.",
                    message.media.mime_type or "file of unstated type")
        _register(container, message, refused=reason)
        _TAKEN[message.provider_message_id] = ""

        return None

    if (message.media.size_bytes or 0) > _MAX_ATTACHMENT_BYTES:
        inbox.forwards_rejected += 1
        reason = f"Larger than the CMS accepts: {message.media.size_bytes} bytes."
        logger.info(
            "Refused a forwarded document of %d bytes: larger than the CMS accepts.",
            message.media.size_bytes,
        )
        _register(container, message, refused=reason)
        _TAKEN[message.provider_message_id] = ""

        return None

    path = container.incoming_dir / f"{_safe_stem(message.provider_message_id)}{suffix}"

    try:
        container.whatsapp.download_media(message.media, path)
    except Exception:
        # Counted, then re-raised so the traceback still carries the provider's
        # HTTP status — which is the actual diagnosis.
        #
        # Cached as settled before raising. One attempt per batch has always
        # been the contract; without this the announcement pass and the
        # processing pass would each take a swing, and one dead download would
        # move two counters and register twice.
        inbox.forwards_failed += 1
        _TAKEN[message.provider_message_id] = ""
        _register(container, message, refused="The document could not be downloaded.")

        raise

    if not _looks_like(path, suffix):
        inbox.forwards_rejected += 1
        reason = f"The file does not begin like the {suffix} it claims to be."
        logger.info("Refused a forwarded %s: the file does not begin like one.", suffix)
        path.unlink(missing_ok=True)
        _register(container, message, refused=reason)
        _TAKEN[message.provider_message_id] = ""

        return None

    inbox.forwards_accepted += 1

    # Registered before anything reads it, and the document does not proceed
    # until the CMS has a record of it. A held document waits in a place with a
    # depth an operator can see; it does not get quietly processed on the
    # assumption somebody will be told about it later.
    reference = _register(container, message, path=path)

    if reference is None:
        # The CMS is unreachable; the outbox is holding the registration with
        # the file's path. Cached as settled so the processing pass does not
        # add the same held registration a second time.
        _TAKEN[message.provider_message_id] = ""

        return None

    # Carried on the message so every later stage can address the record without
    # a lookup, and without _take_in growing a second return value that three
    # call sites would have to unpack.
    _REFERENCES[message.provider_message_id] = reference
    _TAKEN[message.provider_message_id] = str(path)

    return path


#: Which CMS record each in-flight message belongs to.
#:
#: Bounded and disposable: an entry is written when a document is registered and
#: read when its workflow reports where it got to. Losing one costs a status
#: update, never the document — the record itself is in the CMS, which is the
#: whole point of registering before reading.
_REFERENCES: dict[str, str] = {}


def _reference(message) -> str | None:
    return _REFERENCES.get(getattr(message, "provider_message_id", "") or "")


def _push(container, reference: str | None, status: str | None = None, **fields) -> bool:
    """Tell the CMS where a document has got to.

    Best-effort by design, and the asymmetry with registration is deliberate.
    Registration failing means the document is invisible, so it blocks. A status
    push failing means the record exists but reads as stale — worse than
    accurate, far better than absent — so it is logged and the work continues.

    The one caller that must not treat it that way is the hand-off out of
    pending_pdfs, which checks the return value: a document may only leave the
    durable queue once the CMS has agreed it is being processed.
    """
    if reference is None:
        return False

    try:
        container.cms.update_intake(reference, status=status, **fields)

        return True
    except Exception as error:  # noqa: BLE001
        logger.warning("Could not tell the CMS that %s is %s: %s", reference, status, error)

        return False


#: The longest failure reason worth putting on a record.
#:
#: The column holds 255. A reader needs the first sentence — "this image could
#: not be decoded" — and everything after it belongs in the agent's log, where a
#: full traceback is useful rather than in the way.
_REASON_LIMIT = 220


def _why_it_failed(run) -> str:
    """What to tell a person about a run that failed.

    "The workflow failed" was true and useless. A reviewer opening the Intake
    Center could not tell "this image is corrupt, ask the client to resend it"
    from "the agent had a bad afternoon, press retry" — the first needs a phone
    call and the second needs a click, and the record could not distinguish them.

    ## The first thing that broke, not the last

    Using the run's own error was an improvement and still the wrong sentence. A
    workflow keeps going after a step fails, so the error left on the run is the
    *final* symptom rather than the cause. On a corrupt photograph the reader
    failed to decode it, the run carried on, and the filing step failed for want
    of a client — so the record said "nothing identifies a client", and a
    reviewer would have gone looking for the client by hand, never learning the
    image was unreadable.

    The first step to record an error is almost always the cause, and the ones
    after it are consequences of it. That is the sentence somebody can act on.

    Only the first line of it: a traceback's later lines describe this system's
    internals to somebody who wants to know about their client's document.
    """
    steps = getattr(run, "steps", None) or []

    causes = []

    for step in steps:
        # A step's recorded error is the obvious cause.
        if error := (getattr(step, "error", None) or "").strip():
            causes.append(error)

            continue

        # And a step that *succeeded* while reporting a failure is the subtler
        # one. The reader is deliberately tolerant — a document nobody could read
        # must still reach a person, so failing the step would end the workflow
        # instead of queueing it — but it records why, and that record is the
        # sentence somebody can act on. Without this the run named the filing
        # step's complaint about a missing client, and a reviewer would go
        # looking for the client rather than for a readable copy.
        output = getattr(step, "output", None) or {}

        if isinstance(output, dict) and (failure := (output.get("failure") or "")):
            causes.append(str(failure).strip())

    # The run's own error as a fallback, for a run that failed before any step
    # recorded anything — a workflow that could not be built, say.
    error = causes[0] if causes else (getattr(run, "error", None) or "").strip()

    if not error:
        # A run can fail with nothing recorded against it anywhere. Saying so
        # plainly beats inventing a cause.
        return "The workflow failed without recording a reason."

    first = error.splitlines()[0].strip()

    if len(first) > _REASON_LIMIT:
        first = first[: _REASON_LIMIT - 1].rstrip() + "…"

    return first


def _proposal_reference(run) -> str | None:
    """The DOC- handle of the proposal this run produced, if it produced one.

    Read off the step that submitted it rather than tracked separately: the step
    already records what the CMS handed back, and a second copy kept elsewhere is
    a second thing to keep in step.

    Most runs have none — a document that named nobody, or one that failed before
    proposing — and None is the right answer for those rather than an error.
    """
    for step in getattr(run, "steps", None) or []:
        output = getattr(step, "output", None) or {}

        if isinstance(output, dict) and (reference := output.get("proposal_reference")):
            return str(reference)

    return None


def _push_outcome(container, reference: str | None, run, workflow_name: str, forget: str | None = None) -> None:
    """Report a finished run as one of the six states.

    A run that stopped to ask is Waiting, not Approval — it read the document,
    named nobody, and submitted nothing, so there is no proposal for anybody to
    approve. Calling that Awaiting approval would put it in a queue where a
    reviewer would open it and find nothing to decide.
    """
    from app.workflow.state import RunState

    if reference is None:
        return

    if _stopped_to_ask(run):
        status, reason = "waiting", "Read, but the document names nobody."
    elif run.state is RunState.AWAITING_APPROVAL:
        status, reason = "approval", "Waiting for a reviewer."
    elif run.state is RunState.COMPLETED:
        status, reason = "completed", None
    elif run.state is RunState.FAILED:
        status, reason = "failed", _why_it_failed(run)
    else:
        # Still running, which is not an outcome. Nothing to report.
        return

    _push(
        container,
        reference,
        status,
        reason=reason,
        workflow_id=getattr(run, "id", None),
        workflow_name=workflow_name,
        # Links the two records to each other. Without it the intake record and
        # its proposal never knew about one another, so the self-healer had
        # nothing to match and every document that filed successfully stayed at
        # "awaiting a reviewer" for ever — with its file already on the client's
        # record. The endpoint has accepted this since contract 1.1; nothing sent
        # it until now.
        proposal_reference=_proposal_reference(run),
    )

    if forget and status in {"completed", "failed"}:
        # Finished with, as far as this process is concerned. Waiting is kept:
        # an identifier may still arrive, and the run that follows has to be
        # able to report against the same record.
        _REFERENCES.pop(forget, None)


def _recover_in_flight(container) -> None:
    """Find documents the CMS thinks are being processed, and nothing is.

    Run once at startup. A document that was mid-read when the process died left
    the durable queue and never reached a conclusion, so the CMS holds it at
    Processing with nothing behind it — which reads exactly like a document being
    worked on, and stays that way forever.

    They are marked Failed rather than silently retried. A reviewer looking at
    the Failed tab can retry one deliberately; a document that quietly re-reads
    itself on every restart is how a loop nobody ordered gets built. Saying "this
    stopped and nobody noticed" out loud is the more useful answer.
    """
    try:
        # Asked for by state, and for as many as the CMS will give at once.
        # Asking for everything open returned the oldest page of it, so a
        # document stranded mid-read behind a backlog was never seen.
        stranded, truncated = container.cms.pending_intake(status="processing", per_page=100)
    except Exception as error:  # noqa: BLE001
        logger.warning("Could not ask the CMS what it thinks is unfinished: %s", error)

        return

    if truncated:
        # More than one page of stranded documents is not a normal condition,
        # and quietly recovering a fraction of them would hide it.
        logger.warning(
            "More documents are stranded than this pass can reconcile; "
            "the rest will be picked up on the next restart."
        )

    if not stranded:
        return

    logger.warning(
        "%d document(s) were being processed when this agent last stopped.",
        len(stranded),
    )

    for record in stranded:
        _push(
            container,
            record.get("reference"),
            "failed",
            reason="The agent stopped while this document was being read. It was not lost — retry it.",
        )


#: How many runs to check for a vanished proposal in one pass.
#:
#: One CMS call each, so this is a rate-limit budget rather than a performance
#: one. Forty orphans clear in two passes; a pathological number clears over a
#: morning without crowding out the work that matters.
_ORPHAN_CHECK_LIMIT = 25


def _retire_orphan_runs(container) -> None:
    """End runs waiting on a decision that can never come (ADR-0011).

    A run paused at `awaiting_approval` holds the id of the proposal it is
    waiting for. When that proposal no longer exists — a restored database, a
    rolled-back migration, a proposal deleted by somebody with the grant — the
    run waits forever. The decision poller skips a *decision* it cannot match to
    a run; it has no answer for a *run* it cannot match to a proposal, so each
    one is re-read from the store on every poll and the agent does steadily more
    work to reach the same conclusion about documents nobody can act on.

    Staging held forty of them, which is what made the gap visible.

    ## It acts only on a definite answer

    A run is failed only when the CMS says 404 — the proposal is not there. Any
    other outcome, including the CMS being unreachable or refusing, leaves the
    run exactly as it is. The rule is the one the compatibility gate follows:
    refuse on knowledge, never on ignorance. Failing a run because the network
    blinked would destroy work that was merely waiting.

    The document itself is not affected. Its intake record is in the CMS, which
    is the copy that matters; this only stops the agent waiting on a conversation
    the other side has forgotten.
    """
    from app.api.client import NotFound
    from app.workflow.state import RunState, StepState

    try:
        runs = container.engine.store.resumable()
    except Exception:  # noqa: BLE001
        logger.exception("Could not read the workflow store to look for orphaned runs.")

        return

    checked = retired = 0

    for run in runs:
        if checked >= _ORPHAN_CHECK_LIMIT:
            break

        if run.state is not RunState.AWAITING_APPROVAL:
            continue

        paused = next((s for s in run.steps if s.state is StepState.AWAITING_APPROVAL), None)
        proposal_id = paused.output.get("proposal_id") if paused else None

        if not proposal_id:
            continue

        checked += 1

        try:
            container.cms.get_proposal(int(proposal_id))

            continue
        except NotFound:
            pass
        except Exception:  # noqa: BLE001
            # Unreachable, refused, rate limited — all mean "ask again later".
            continue

        run.state = RunState.FAILED
        run.error = (
            f"Proposal {proposal_id} no longer exists in the CMS, so this run "
            "could never be decided. The document's own record is unaffected."
        )

        try:
            container.engine.store.save(run)
            retired += 1
        except Exception:  # noqa: BLE001
            logger.exception("Could not retire orphaned run %s.", run.id)

    if retired:
        logger.warning(
            "Retired %d run(s) waiting on a proposal the CMS no longer has.", retired
        )


def _reconcile(container) -> None:
    """Check the two databases still agree about what exists (ADR-0010).

    The CMS holds the record; the agent holds the work. They are different
    databases on different machines with no transaction between them, so
    consistency is something to be checked rather than assumed — and the whole
    module rests on the claim that nothing is in a state nobody can see. A claim
    that is never verified is a hope.

    Two disagreements matter, and they fail in opposite directions:

      * The CMS thinks a document is Waiting and the agent is not holding it.
        The record is a promise nothing will keep — nobody supplying a file
        number would make anything happen. Marked Failed, because "this stopped"
        is true and "still waiting" is not.

      * The agent is holding a document the CMS has no live record of. That is
        the original bug in miniature, so it is registered again rather than
        reported: the point is that the document becomes visible, not that
        somebody is told it was not.
    """
    try:
        unfinished, _ = container.cms.pending_intake(status="waiting", per_page=100)
    except Exception as error:  # noqa: BLE001
        logger.warning("Could not reconcile with the CMS: %s", error)

        return

    try:
        held = container.pending_pdfs.take_older_than(datetime.now(UTC) + timedelta(days=365))
    except Exception:  # noqa: BLE001
        logger.exception("Could not read the held documents to reconcile them.")

        return

    # take_older_than removes what it returns, so anything read here has to go
    # back. Reconciling must not empty the queue it is checking.
    for pdf in held:
        container.pending_pdfs.add(pdf)

    holding = {pdf.intake_reference for pdf in held if pdf.intake_reference}

    stranded = [
        record for record in unfinished
        if record.get("reference") not in holding
    ]

    for record in stranded:
        logger.warning(
            "%s is shown as Waiting but no document is held for it; marking it failed.",
            record.get("reference"),
        )
        _push(
            container,
            record.get("reference"),
            "failed",
            reason="The agent is no longer holding this document. Nothing was waiting for an identifier.",
        )

    if stranded:
        logger.warning("Reconciled %d document(s) the CMS and the agent disagreed about.", len(stranded))


#: Which process is holding a claim. Recorded so an operator looking at a stuck
#: row can tell "a worker has this" from "a worker had this and died".
_WORKER = f"{socket.gethostname()}:{os.getpid()}"


def _drain_outbox(container) -> None:
    """Send the registrations that were held while the CMS was unreachable.

    Oldest first, and only one is retried per pass. If the CMS is still down the
    first attempt fails and the rest of the queue is left alone rather than
    generating a burst of failing requests against a server that is already
    struggling — and if it is back up, the queue empties over the next few ticks
    in the order the documents arrived.

    Deleted only after the CMS acknowledges. A failure leaves the row where it
    is, with its attempt count raised, which is what makes "held" different from
    "lost".
    """
    from app.api.client import NotFound

    outbox = container.intake_outbox

    # Claimed, not merely read. Two workers reading the oldest row both believed
    # they owned it and both registered the document — one claim, one owner
    # (ADR-0011). The claim is a lease, so a worker that dies holding one does
    # not strand the document: it ages out and the next pass takes it.
    held = outbox.claim(_WORKER)

    if held is None:
        return

    path = Path(held.path) if held.path else None

    if path is not None and not path.exists():
        # The bytes are gone — cleaned up, or the deployment was rebuilt. The
        # arrival is still worth recording, so it goes without the file rather
        # than being dropped: a record saying a document arrived and could not be
        # kept is far better than no record at all.
        logger.warning(
            "%s was held for registration but its file is gone; registering the "
            "arrival without it.",
            held.original_filename or held.provider_message_id,
        )
        path = None

    try:
        record = container.cms.register_intake(held.submission(), path=path)
    except Exception as error:
        outbox.failed(held, str(error))

        # A 404 is not "the CMS is down", and the difference is the whole
        # diagnosis. The outbox exists for a CMS that is temporarily
        # unreachable, where holding on is exactly right and the queue drains
        # when it returns. A CMS with no intake endpoint will never accept this
        # call: the queue fills instead, and once it is full the agent turns new
        # documents away. Intake goes quiet, and every line in the log says
        # "still cannot register", which reads as a network problem.
        #
        # The document is held either way — that part was already right. What
        # was missing was anything telling an operator which of the two they
        # are looking at.
        if isinstance(error, NotFound):
            logger.error(
                "The CMS has no intake endpoint, so %s cannot be registered and "
                "retrying will not help: this CMS predates the Intake Center and "
                "is too old for this release. Update the CMS. (%s)",
                held.original_filename or held.provider_message_id,
                error,
            )
        else:
            logger.warning(
                "Still cannot register %s with the CMS (attempt %d): %s",
                held.original_filename or held.provider_message_id,
                held.attempts + 1,
                error,
            )

        return

    outbox.done(held)

    logger.info(
        "Registered %s as %s after holding it for %.0f seconds.",
        held.original_filename or held.provider_message_id,
        record.get("reference"),
        held.waited(),
    )


def register_arrival(
    container,
    *,
    source: str,
    path=None,
    filename: str | None = None,
    mime_type: str | None = None,
    size_bytes: int | None = None,
    provider_message_id: str | None = None,
    source_account: str | None = None,
    refused: str | None = None,
) -> str | None:
    """Tell the CMS a document arrived, or hold the fact until it can be told.

    Source-agnostic on purpose (ADR-0010). Every way a document can reach this
    firm goes through here — WhatsApp today, a manual submission beside it, and
    email or a portal later — because the guarantee is about documents, not about
    WhatsApp, and a second registration path is a second place for one to slip
    through unrecorded. That is not hypothetical: the manual submission script
    started workflows directly for months, so anything put through it was
    invisible to the Intake Center, which is the exact condition this record
    exists to forbid.

    Called on every path out of the gate, refusals included. A client who forwards
    a .docx would otherwise get silence, because the only trace is a counter
    nobody is watching; registering the refusal turns it into a row somebody can
    read, with a reason attached.

    Returns the INT- reference, or None when the CMS could not be reached and the
    registration was queued instead. **None means "do not proceed"** — not "carry
    on without it" — because a document the CMS has never heard of is exactly the
    invisible state this was built to abolish.
    """
    from app.intake import OutboxFull, PendingRegistration

    label = filename or provider_message_id or (Path(path).name if path else "a document")

    # Every source gets an idempotency key, not just the one that happens to
    # carry a message id.
    #
    # The CMS deduplicates on (source, provider_message_id), which protected
    # WhatsApp and nothing else: a manual document arrived with NULL, so two
    # deliveries of one file became two records in front of a reviewer. Found by
    # the stress harness — the outbox drain has no lock, two drainers took the
    # same row, and only the missing key turned that race into a duplicate.
    #
    # Minted once here and carried on the outbox row, so a retry after an outage
    # sends the same key rather than a new one.
    if provider_message_id is None:
        provider_message_id = f"{source}:{uuid.uuid4().hex}"

    submission = {
        "source": source,
        "source_account": source_account,
        "provider_message_id": provider_message_id,
        "original_filename": filename,
        "mime": mime_type,
        "size_bytes": size_bytes,
    }
    submission = {key: value for key, value in submission.items() if value is not None}

    if refused:
        submission["failure_reason"] = refused

    try:
        record = container.cms.register_intake(submission, path=path)
        reference = record.get("reference")

        logger.info(
            "Registered %s from %s as %s%s.",
            label, source, reference, " (refused)" if refused else "",
        )

        return reference
    except Exception as error:  # noqa: BLE001
        # Not a failure of the document — a failure to reach the CMS. Holding it
        # is the approved trade (ADR-0010 §7): refusing a client's document
        # because our own server blinked is the worse outcome.
        try:
            container.intake_outbox.add(
                PendingRegistration(
                    source=source,
                    source_account=source_account,
                    provider_message_id=provider_message_id,
                    original_filename=filename,
                    mime_type=mime_type,
                    size_bytes=size_bytes,
                    path=str(path) if path else None,
                    failure_reason=refused,
                    received_at=datetime.now(UTC),
                )
            )

            logger.warning(
                "Could not register %s with the CMS (%s); it is held and will be "
                "retried. Nothing will be read until the CMS has the record.",
                label, error,
            )
        except OutboxFull:
            # The cap has been reached, which means the CMS has been unreachable
            # for a long time. Refusing loudly beats hoarding quietly.
            logger.error(
                "The intake outbox is full and the CMS is unreachable. %s was NOT "
                "accepted. Documents are being refused until the CMS returns.",
                label,
            )

        return None


def _register(container, message, path=None, refused: str | None = None) -> str | None:
    """The WhatsApp adapter. Everything it knows comes off the message."""
    return register_arrival(
        container,
        source="whatsapp",
        path=path,
        filename=message.media.filename,
        mime_type=message.media.mime_type,
        size_bytes=message.media.size_bytes,
        provider_message_id=message.provider_message_id,
        source_account=message.sender,
        refused=refused,
    )


def _hold_for_metadata(container, message, classification) -> None:
    """Download the document now; wait for the message that says whose it is.

    Downloaded immediately rather than when the identifier arrives, because the
    media handle is the provider's and does not keep — waiting first and
    fetching later would mean the file expires exactly when somebody finally
    says who it belongs to.

    Nothing here opens it, and the caption is neither read nor stored.
    """
    from app.whatsapp.pending import PendingPdf

    # Counted inside, at the moment the document is taken in rather than when it
    # is finally filed: a count that waited for the naming message would read as
    # nothing having arrived during exactly the window an operator asks about.
    path = _take_in(container, message)

    if path is None:
        return

    destination = router.route(classification)

    container.pending_pdfs.add(
        PendingPdf(
            message_id=message.provider_message_id,
            path=str(path),
            sender=message.sender,
            received_at=message.received_at,
            file_name=(message.media.filename or path.name),
            intake_reference=_reference(message),
            mime_type=(message.media.mime_type or "application/pdf"),
            size_bytes=path.stat().st_size,
            classification={
                "workflow": destination.workflow,
                "filing_type": classification.filing_type,
                "method": classification.method.value,
                "confidence": classification.confidence,
                "reason": classification.reason,
                "refinement": classification.refinement,
                "label": classification.label,
                "was_read": classification.was_read,
            },
        )
    )

    # Held here means Waiting there. This is the state that used to have no name
    # and no screen: the document was in a table only this process could read,
    # and asking the CMS about it produced nothing at all.
    #
    # PROCESSING FIRST, BECAUSE RECEIVED CANNOT REACH WAITING
    #
    # A registered document is Received, and AiIntakeDocument::TRANSITIONS only
    # allows Received → Processing or Failed. Pushing Waiting straight from
    # Received is refused with a 409, and the CMS is right to refuse it: a
    # document that reached Waiting was worked on first, and skipping the step
    # would claim it arrived already stuck.
    #
    # It was refused in production on 2026-08-04, and the failure was quiet in
    # the worst way — the push is best-effort, so the hold still happened and the
    # file was safe, while the Intake Center went on showing Received for a
    # document nothing would ever move again.
    #
    # Processing is not a formality here. This document *was* worked on: it was
    # downloaded, its type classified, and the classification is in the event
    # recorded just above.
    _push(
        container,
        _reference(message),
        "processing",
        reason="Classifying the document and looking for a client.",
        workflow_name=destination.workflow,
    )

    _push(
        container,
        _reference(message),
        "waiting",
        reason="Waiting for a message naming the client by file number or CNIC.",
        workflow_name=destination.workflow,
    )

    logger.info(
        "Holding %s as %s (%s, %.0f%%) until a message names the client.",
        path.name,
        classification.label,
        classification.method.value,
        classification.confidence * 100,
    )


def _metadata_arrived(container, message) -> None:
    """A plain text message: the identifier for the oldest waiting PDF.

    With nothing waiting this is the owner talking to themselves, which is an
    ordinary thing to do and not an instruction — ignored, exactly as before.
    """
    from app.documents import identifiers, replies

    text = (message.text or "").strip()

    if not text:
        logger.debug("Nothing to read in message %s.", message.provider_message_id)

        return

    # A reply naming a document wins, and is tried first.
    #
    # It has to be. The path below takes the OLDEST waiting PDF, which is a
    # sensible default for a message that names nothing and precisely the wrong
    # answer for one that does — a reply saying DOC-1004 would otherwise file
    # the identifier onto DOC-1001 and report success.
    #
    # The two also address different populations: pending_pdfs are held here and
    # have never reached the CMS, while a reference names a proposal already in
    # the queue. Answering an explicit question with a queue-order guess would
    # be wrong even if the populations overlapped.
    reply = replies.read(text)

    if reply is not None:
        _identify_by_reply(container, reply)

        return

    pdf = container.pending_pdfs.take_oldest()

    if pdf is None:
        logger.debug(
            "Message %s names nobody's document; nothing is waiting.",
            message.provider_message_id,
        )

        return

    _file_waiting_pdf(container, pdf, identifiers.read(text))


def _identify_by_reply(container, reply) -> None:
    """Hand the sender's answer to the CMS, which decides whether it is usable.

    Nothing is resolved here. The agent does not hold the client list and must
    not start guessing at one — the CMS owns who a file number belongs to, and
    it is the only side that can refuse an identifier belonging to two people or
    to nobody.

    Every refusal is logged and swallowed. This runs while draining the inbox,
    and a raised exception would stop the remaining messages being read over an
    answer that was simply wrong. The document stays on Waiting for Information
    either way, which is the recoverable outcome.
    """
    try:
        result = container.cms.identify_proposal(reply.reference, reply.identifier)
    except Exception as exc:  # noqa: BLE001
        logger.warning("Could not identify %s from the reply: %s", reply.reference, exc)

        return

    client = (result or {}).get("client") or {}

    logger.info(
        "%s identified as %s from a WhatsApp reply.",
        reply.reference,
        client.get("name", "a client"),
    )


def _file_pdfs_nobody_named(container) -> None:
    """Send on the PDFs whose identifier never came.

    They go to the same workflow they would have gone to anyway, with no client
    attached, which lands them in the approval queue for a person to place. That
    is the same destination as a text that named nobody — the difference is only
    how long we waited first.
    """
    from app.documents import identifiers

    try:
        stale = container.pending_pdfs.take_older_than(datetime.now(UTC) - _pdf_metadata_wait())
    except Exception:  # noqa: BLE001
        # A queue that cannot be read must not stop the inbox being drained.
        logger.exception("Could not check for PDFs still waiting to be named.")

        return

    for pdf in stale:
        try:
            logger.info(
                "No identifier arrived for %s in %.0f minutes; sending it for review.",
                pdf.file_name or pdf.path,
                pdf.waited() / 60,
            )

            # Announced before the work starts. take_older_than() has already
            # removed the row, so from here until the run finishes the document
            # exists in this process and nowhere else — the window in which a
            # crash used to lose a client's file with no trace that it had ever
            # arrived. The CMS record is what closes it: whatever happens next,
            # something outside this process knows the document is being worked
            # on and can be found again.
            _push(
                container,
                pdf.intake_reference or _REFERENCES.get(pdf.message_id),
                "processing",
                reason="No identifier arrived; reading it for review.",
            )

            _file_waiting_pdf(container, pdf, identifiers.Identifiers())
        except Exception:  # noqa: BLE001
            logger.exception("Could not file %s after its wait expired.", pdf.message_id)

            # The row is already gone and the run did not finish. Say so, rather
            # than leaving the record reading Processing forever — a document
            # stuck at Processing with nothing running is indistinguishable from
            # one being worked on, which is the confusion this module exists to
            # remove.
            _push(
                container,
                pdf.intake_reference or _REFERENCES.get(pdf.message_id),
                "failed",
                reason="Reading the document failed after its wait expired.",
            )


def _file_waiting_pdf(container, pdf, found) -> None:
    """Start the run for a PDF now that we know — or know we do not know — whose it is."""
    from pathlib import PurePath

    from app.workflow import catalogue

    stored = dict(pdf.classification or {})
    workflow = catalogue.require(stored.get("workflow") or "document_intake")

    logger.info(
        "Filing %s from the self-chat, identified by %s.",
        PurePath(pdf.path).name,
        found.describe(),
    )

    run = container.engine.start(
        workflow,
        {
            "path": pdf.path,
            "sender": pdf.sender,
            "workflow_name": workflow.name,
            "intake_reference": pdf.intake_reference or _REFERENCES.get(pdf.message_id),
            "file_number": found.file_number,
            "cnic": found.cnic,
            "document_type": stored.get("filing_type"),
            "classification": {
                "method": stored.get("method"),
                "confidence": stored.get("confidence"),
                "reason": stored.get("reason"),
                "refinement": stored.get("refinement"),
                "label": stored.get("label"),
                "was_read": stored.get("was_read", False),
            },
            "source": {
                "channel": "whatsapp",
                "message_id": pdf.message_id,
                "number": pdf.sender,
                "received_at": pdf.received_at.isoformat(),
                "file_name": pdf.file_name or PurePath(pdf.path).name,
                "file_size": pdf.size_bytes or 0,
                "mime_type": (pdf.mime_type or "application/pdf").lower(),
            },
        },
    )

    # The stored reference first: this document may have been held across a
    # restart, and the in-memory map does not survive one.
    if pdf.intake_reference:
        _REFERENCES.setdefault(pdf.message_id, pdf.intake_reference)

    _push_outcome(container, pdf.intake_reference or _REFERENCES.get(pdf.message_id),
                  run, workflow.name, forget=pdf.message_id)

    logger.info("%s %s finished: %s", workflow.name, run.id, run.state.value)


#: What may be forwarded with an identifier, and the suffix each gets.
#:
#: A narrower list than the reader's below, for a different reason. That one is
#: about what PaddleOCR will open; this one is about what a reviewer can look at
#: in a browser before approving it. A .webp or a .tif is perfectly readable by
#: machine and shows a person nothing.
_FORWARDABLE_MIMES = {
    "application/pdf": ".pdf",
    "image/jpeg": ".jpg",
    "image/jpg": ".jpg",
    "image/png": ".png",
}

#: The bytes each accepted type must actually begin with.
#:
#: Checked because the media type is what the sending phone claimed, not what
#: the file is. The common real cause is not deception but a download that
#: stopped early — which has no header at all, and which would otherwise be
#: discovered by a reviewer staring at a preview that will not render.
_MAGIC = {
    ".pdf": (b"%PDF",),
    ".jpg": (b"\xff\xd8\xff",),
    ".png": (b"\x89PNG\r\n\x1a\n",),
}

#: The largest attachment worth downloading, matching the CMS's own limit.
#:
#: Enforced here as well as there, because the alternative is reading a 40 MB
#: file into memory, expanding it to 53 MB of base64, and posting it to be
#: refused.
_MAX_ATTACHMENT_BYTES = 20 * 1024 * 1024


def _file_identified(container, message, found, workflow, classification=None) -> None:
    """File a document whose client the sender named in the caption.

    Nothing here opens the document. Every check is on the envelope — the media
    type, the size, the first few bytes — because the contents are Level 3 and
    the premise of this path is that they are never read (ADR-0002).

    A refusal returns quietly. There is no reply to the sender in this version,
    so the counter and the log line are the whole account of it; both exist
    because the alternative is a document that silently never arrives.
    """
    # The filter itself, not the receiver's view of it — they are the same
    # object, and reaching through the receiver would break the moment it stopped
    # being built from this one.
    inbox = container.inbox
    suffix = _FORWARDABLE_MIMES.get((message.media.mime_type or "").lower())

    if suffix is None:
        inbox.forwards_rejected += 1
        logger.info(
            "Refused a forwarded %s: only PDF, JPG and PNG can be filed this way.",
            message.media.mime_type or "file of unstated type",
        )

        return

    if (message.media.size_bytes or 0) > _MAX_ATTACHMENT_BYTES:
        inbox.forwards_rejected += 1
        logger.info(
            "Refused a forwarded document of %d bytes: larger than the CMS accepts.",
            message.media.size_bytes,
        )

        return

    path = container.incoming_dir / f"{_safe_stem(message.provider_message_id)}{suffix}"

    try:
        container.whatsapp.download_media(message.media, path)
    except Exception:
        # Counted, then re-raised.
        #
        # A document that never arrived is as invisible to the queue as one
        # that was refused — and worse, until this counter existed: a refusal
        # at least moved a number, while a failed fetch moved nothing and lived
        # only in a traceback nobody reads until they are already looking.
        #
        # Separate from forwards_rejected because the two send an operator to
        # different places. Refused means the file is wrong — go and ask whoever
        # sent it. Failed means the file never came — go and look at Evolution.
        #
        # Re-raised so the outer handler still writes the traceback, which
        # carries the provider's HTTP status and is the actual diagnosis.
        inbox.forwards_failed += 1

        raise

    if not _looks_like(path, suffix):
        inbox.forwards_rejected += 1
        logger.info("Refused a forwarded %s: the file does not begin like one.", suffix)
        path.unlink(missing_ok=True)

        return

    inbox.forwards_accepted += 1

    logger.info("Filing %s from the self-chat, identified by %s.", path.name, found.describe())

    # This path never told the Intake Center anything — bank statements were
    # only ever visible once their proposal appeared. Now that every arrival is
    # announced, silence here would strand the record at Queued, so the branch
    # reports its progress the way every other workflow does. The file is still
    # never opened; the status says what is being done with the envelope.
    _push(
        container,
        _reference(message),
        "processing",
        reason="Filing the statement by its caption; the file is never opened.",
        workflow_name=workflow.name,
        current_step="file",
    )

    run = container.engine.start(
        workflow,
        {
            "path": str(path),
            "file_number": found.file_number,
            "cnic": found.cnic,
            "sender": message.sender,
            # How the type was decided, carried through so this path reports it
            # the same way every other workflow does.
            "classification": {
                "method": classification.method.value,
                "confidence": classification.confidence,
                "reason": classification.reason,
                "refinement": classification.refinement,
                "label": classification.label,
                "was_read": classification.was_read,
            } if classification is not None else {},
            # Everything the reviewer is shown about where this came from, and
            # the two fields the CMS deduplicates on. Assembled here because this
            # is the only place that holds the message.
            "source": {
                "channel": "whatsapp",
                "message_id": message.provider_message_id,
                "number": message.sender,
                "received_at": message.received_at.isoformat(),
                "file_name": message.media.filename or path.name,
                "file_size": path.stat().st_size,
                "mime_type": (message.media.mime_type or "").lower(),
            },
        },
    )

    logger.info("Bank statement intake %s finished: %s", run.id, run.state.value)

    _push_outcome(container, _reference(message), run, workflow.name,
                  forget=message.provider_message_id)

    _ask_for_missing_information(container, run, message.sender)


def _ask_for_missing_information(container, run, recipient: str) -> None:
    """Tell the sender their document needs an identifier — at most once, ever.

    The decision is not taken here. The CMS returns a `notify` instruction on
    the first response for a document nobody could be matched to, and stamps the
    row as it answers, so a retried step or a restarted daemon sees nothing left
    to send. This side only carries it out.

    That ordering makes the whole thing at-most-once rather than at-least-once,
    which is the trade worth making: an unsent notification is recoverable — the
    document is sitting on Waiting for Information where somebody will see it —
    and a message storm under a firm's own WhatsApp number is not. This is the
    only message this software sends, and the number it sends from can be banned
    without warning or appeal.

    A failure to send is logged and swallowed. The document is already filed as
    a proposal; raising here would fail a workflow whose real work succeeded.
    """
    notify = _first_notify(run)

    if notify is None:
        return

    from app.whatsapp.notifications import missing_information

    try:
        result = container.message_sender.send(missing_information(
            recipient=recipient,
            reference=str(notify.get("reference", "")),
            required=str(notify.get("required", "a file number or CNIC")),
        ))
    except Exception:
        logger.exception("Could not ask for the missing information on %s.", notify.get("reference"))

        return

    if result.ok:
        logger.info("Asked the sender for the identifier on %s.", notify.get("reference"))
    else:
        logger.warning("The provider refused the request for %s: %s", notify.get("reference"), result.error)


def _first_notify(run) -> dict | None:
    """The one notify instruction in a run's step outputs, if the CMS sent one.

    Deliberately defensive about shape: this reads a field absent from almost
    every response, and an AttributeError here would fail a workflow whose real
    work has already succeeded.
    """
    for record in getattr(run, "steps", []) or []:
        output = getattr(record, "output", None)

        if isinstance(output, dict) and isinstance(output.get("notify"), dict):
            return output["notify"]

    return None


def _classify_message(container, message, *, use_caption: bool = True):
    """What this document is, decided as cheaply as it can be.

    **Nothing is downloaded here.** The caption and the filename are already in
    the message, and between them they settle the overwhelming majority — which
    is the point of the ordering. `read_first_page` is left unset, so the rule
    stage abstains rather than reading a file that is not on disk yet.

    The document is downloaded afterwards, by whichever workflow the router
    chose, under that workflow's own rules — and for a bank statement those
    rules say it is never opened at all (ADR-0002).

    A document neither the caption nor the filename recognises therefore reaches
    the reading workflow as `unknown`, and that workflow OCRs it as it always
    has. Classifying from text before choosing a workflow would mean reading
    every document to discover the ones that must not be read.

    `use_caption=False` for PDFs, whose identifier arrives as a separate message
    rather than a caption. The filename still decides the type — that is what
    routes a bank statement to the workflow that never opens it — but nothing a
    caption says is consulted.
    """
    return classifier.classify(
        Signals(
            path=container.incoming_dir / _incoming_name(message),
            source="whatsapp",
            caption=message.text if use_caption else "",
            filename=(message.media.filename if message.media else "") or "",
            mime_type=(message.media.mime_type if message.media else "") or "",
            size_bytes=(message.media.size_bytes if message.media else 0) or 0,
            # None on every deployment that has not configured one, and
            # refused unless it runs on this host (ADR-0009).
            model=container.model,
        )
    )


def _looks_like(path, suffix: str) -> bool:
    """Whether a file begins the way its type must.

    Not a validation of the document — nothing here parses a PDF. It answers one
    question: is this the kind of file it was announced as?
    """
    expected = _MAGIC.get(suffix)

    if not expected:  # pragma: no cover - every accepted suffix has an entry
        return True

    try:
        with path.open("rb") as handle:
            head = handle.read(8)
    except OSError:
        return False

    return any(head.startswith(prefix) for prefix in expected)


def _safe_stem(message_id: str) -> str:
    """A filename body from a provider id, with nothing a path can act on."""
    return re.sub(r"[^A-Za-z0-9_-]", "_", message_id)[:120] or "document"


#: Suffixes the reader will actually open.
#:
#: PaddleOCR dispatches on the extension and refuses anything else outright —
#: and, because the engine never raises, refusal arrives as an empty read rather
#: than an error. So the extension is not cosmetic: it decides whether a
#: document is read at all.
_READABLE_SUFFIXES = frozenset(
    {".jpg", ".jpeg", ".png", ".webp", ".bmp", ".tif", ".tiff", ".pdf"}
)

#: What to call a file the provider described only by media type.
#:
#: An explicit table rather than mimetypes.guess_extension(), which answers
#: ".jpe" for image/jpeg on some builds — a suffix PaddleOCR does not accept.
_SUFFIX_BY_MIME = {
    "image/jpeg": ".jpg",
    "image/jpg": ".jpg",
    "image/png": ".png",
    "image/webp": ".webp",
    "image/bmp": ".bmp",
    "image/tiff": ".tif",
    "application/pdf": ".pdf",
}


def _incoming_name(message) -> str:
    """A filename that cannot collide or escape the directory, and can be read.

    The body is the provider's message id: unique, already the deduplication
    key, and made of characters the provider generated rather than a client.

    THE EXTENSION DECIDES WHETHER THE DOCUMENT IS READ

    Found on the first real photograph sent through this. WhatsApp images carry
    no filename — only documents do — so every photo was saved as `.bin`,
    PaddleOCR refused the type outright, and because that engine never raises,
    the refusal surfaced as a successful read of zero characters. The classifier
    then said "unknown", the workflow correctly declined to guess, and nothing
    anywhere said the file had simply never been opened.

    So the media type is preferred over the provider's filename: it is a
    constrained vocabulary rather than a path, and it is the field WhatsApp
    actually fills in for an image. The filename is the fallback, and only when
    its suffix is one the reader accepts.
    """
    from pathlib import PurePosixPath

    media = message.media
    suffix = _SUFFIX_BY_MIME.get((media.mime_type or "").split(";")[0].strip().lower(), "")

    if not suffix:
        candidate = PurePosixPath(media.filename or "").suffix.lower()
        suffix = candidate if candidate in _READABLE_SUFFIXES else ""

    # Deliberately still possible to end up with nothing. A file this cannot
    # name is one the reader could not have opened either, and inventing ".jpg"
    # for it would turn "we do not handle this" into a page of noise in front of
    # a reviewer.
    safe = "".join(c for c in message.provider_message_id if c.isalnum() or c in "-_")

    return f"{safe or 'message'}{suffix or '.bin'}"


def _recheck_compatibility(container: Container) -> None:
    """Re-run the handshake, and say something only when the answer changes.

    The CMS updates itself over the air, so the installation this is attached to
    can change contract while this process is running. A verdict reached once at
    start-up goes stale at exactly the moment it matters.
    """
    from app.api.compatibility import Level

    before = container.cms.compatibility

    try:
        verdict = container.cms.handshake()
    except Exception as exc:  # noqa: BLE001 - a network blip is not an incompatibility
        logger.warning("Could not re-check CMS compatibility: %s", exc)

        return

    if before is not None and before.level is verdict.level:
        return

    if verdict.level is Level.INCOMPATIBLE:
        # Loud, because every CMS call is now being refused and the process
        # would otherwise simply look idle.
        logger.error("The CMS is no longer compatible: %s", verdict.detail)
    else:
        logger.info("CMS compatibility: %s", verdict.detail)


def _archive(container: Container) -> None:
    cutoff = datetime.now(UTC) - timedelta(days=container.archive_after_days)
    moved = container.store.archive(older_than=cutoff)

    if moved:
        logger.info("Archived %d finished run(s).", moved)


COMMANDS = {"run": command_run, "check": command_check, "migrate": command_migrate}


def main(argv: list[str] | None = None) -> int:
    parser = argparse.ArgumentParser(prog="taxpilot-ai", description=__doc__)
    parser.add_argument("command", nargs="?", default="run", choices=sorted(COMMANDS))
    parser.add_argument("--log-level", default=os.environ.get("TAXPILOT_LOG_LEVEL", "INFO"))

    arguments = parser.parse_args(argv)
    configure_logging(arguments.log_level)

    try:
        container = Container()
        # Touched here so a missing credential is reported as configuration,
        # with the exit code that says so, rather than as a crash inside a job.
        container.settings
    except ConfigurationError as exc:
        logger.error("%s", exc)

        return EXIT_MISCONFIGURED

    try:
        return COMMANDS[arguments.command](container)
    except KeyboardInterrupt:
        # Reached only if a signal handler was not installed. Not an error.
        logger.info("Interrupted.")

        return EXIT_OK
    except Exception as exc:  # noqa: BLE001 - the last line before a traceback
        logger.exception("Fatal: %s", exc)

        return EXIT_FAILED


if __name__ == "__main__":
    raise SystemExit(main())
