"""The process that actually runs.

Most of this is about stopping and about failing. A daemon that cannot be stopped
within Docker's grace period is killed mid-write; one that exits on the first
transient error turns a thirty-second network blip into an outage lasting until
somebody notices.
"""

from __future__ import annotations

import tempfile
import threading
import time
from pathlib import Path

import pytest

from app.runtime import preflight
from app.runtime.container import Container
from app.runtime.daemon import Daemon
from app.runtime.preflight import Check, Outcome

def fast_daemon() -> Daemon:
    """A daemon that does not idle for five seconds between passes.

    max_wait_seconds is a production setting — a loop that wakes every five
    seconds costs nothing and stays responsive. In a test it is dead time, so
    these ask for a short one rather than changing the default.
    """
    return Daemon(max_wait_seconds=0.01)


BASE_ENV = {
    "TAXPILOT_CMS_URL": "https://cms.test",
    "TAXPILOT_AGENT_KEY": "tpa_test",
    "TAXPILOT_AGENT_SECRET": "s" * 64,
}


class TestJobScheduling:
    def test_a_job_runs_when_it_is_due(self):
        daemon = fast_daemon()
        calls: list[int] = []
        daemon.add("work", lambda: calls.append(1), interval_seconds=0)

        daemon.run(max_iterations=3)

        assert len(calls) == 3

    def test_a_job_not_yet_due_does_not_run_again(self):
        daemon = fast_daemon()
        calls: list[int] = []
        daemon.add("work", lambda: calls.append(1), interval_seconds=3600)

        daemon.run(max_iterations=3)

        # Due immediately the first time, then not for an hour.
        assert len(calls) == 1

    def test_jobs_keep_their_own_intervals(self):
        daemon = fast_daemon()
        often, rarely = [], []
        daemon.add("often", lambda: often.append(1), interval_seconds=0)
        daemon.add("rarely", lambda: rarely.append(1), interval_seconds=3600)

        daemon.run(max_iterations=4)

        assert len(often) == 4
        assert len(rarely) == 1

    def test_a_daemon_with_no_jobs_still_stops(self):
        fast_daemon().run(max_iterations=2)


class TestFailureIsolation:
    def test_a_failing_job_does_not_stop_the_loop(self):
        daemon = fast_daemon()
        survivor: list[int] = []

        def explode():
            raise RuntimeError("the CMS restarted")

        daemon.add("broken", explode, interval_seconds=0)
        daemon.add("fine", lambda: survivor.append(1), interval_seconds=0)

        daemon.run(max_iterations=3)

        # A network blip is an ordinary Tuesday. Exiting on it turns thirty
        # seconds of trouble into an outage lasting until somebody notices.
        assert len(survivor) == 3

    def test_failures_are_counted_and_the_reason_kept(self):
        daemon = fast_daemon()

        def explode():
            raise RuntimeError("connection reset")

        job = daemon.add("broken", explode, interval_seconds=0)
        daemon.run(max_iterations=2)

        assert job.runs == 2
        assert job.failures == 2
        assert job.last_error == "connection reset"

    def test_a_job_that_recovers_clears_its_error(self):
        daemon = fast_daemon()
        attempts = {"n": 0}

        def flaky():
            attempts["n"] += 1
            if attempts["n"] == 1:
                raise RuntimeError("first attempt")

        job = daemon.add("flaky", flaky, interval_seconds=0)
        daemon.run(max_iterations=2)

        assert job.failures == 1
        assert job.last_error is None


class TestStopping:
    def test_stop_ends_the_loop(self):
        daemon = fast_daemon()
        daemon.add("work", daemon.stop, interval_seconds=0)

        daemon.run()   # no max_iterations: this must end on its own

        assert daemon.stopping

    def test_a_stop_during_a_pass_skips_the_remaining_jobs(self):
        daemon = fast_daemon()
        later: list[int] = []

        daemon.add("stopper", daemon.stop, interval_seconds=0)
        daemon.add("later", lambda: later.append(1), interval_seconds=0)

        daemon.run()

        # Checked between jobs as well as between iterations, so a stop arriving
        # during a slow pass is not held up by everything after it.
        assert later == []

    def test_waiting_is_interrupted_by_a_stop(self):
        """The difference between a clean shutdown and SIGKILL."""
        daemon = fast_daemon()
        daemon.add("slow", lambda: None, interval_seconds=3600)

        def stop_shortly():
            time.sleep(0.1)
            daemon.stop()

        threading.Thread(target=stop_shortly, daemon=True).start()

        started = time.monotonic()
        daemon.run()
        elapsed = time.monotonic() - started

        # time.sleep() here would hold the process for up to max_wait_seconds
        # past the request. Docker waits ten seconds and then sends SIGKILL,
        # which arrives mid-write.
        assert elapsed < 2.0

    def test_a_clock_jump_does_not_disturb_scheduling(self):
        # Intervals are monotonic, so an NTP correction or a daylight-saving
        # change cannot make a job run twice or stall for an hour.
        daemon = fast_daemon()
        job = daemon.add("work", lambda: None, interval_seconds=60)
        daemon.run(max_iterations=1)

        assert job.next_due > 0


class TestPreflight:
    class FakeCms:
        """CmsClient's surface as preflight uses it.

        `check_cms` performs the handshake rather than a bare whoami, so one
        call answers both "can we reach it?" and "can we work with it?" and
        leaves the verdict where the compatibility gate reads it.
        """

        def __init__(self, identity=None, error=None):
            self._payload = identity or {}
            self._error = error
            self.identity = None
            self.compatibility = None

        def whoami(self):
            if self._error:
                raise RuntimeError(self._error)

            return self._payload

        def handshake(self):
            from app.api.compatibility import evaluate

            self.identity = self.whoami()
            self.compatibility = evaluate(self.identity)

            return self.compatibility

    def test_a_working_cms_passes(self):
        cms = self.FakeCms({"name": "Intake", "permissions": ["clients.read", "proposals.submit"]})

        check = preflight.check_cms(cms)

        assert check.outcome is Outcome.OK
        assert "Intake" in check.detail

    def test_an_unreachable_cms_is_fatal(self):
        check = preflight.check_cms(self.FakeCms(error="connection refused"))

        assert check.is_fatal
        assert "connection refused" in check.detail

    def test_missing_grants_are_fatal_and_named(self):
        cms = self.FakeCms({"name": "Intake", "permissions": ["clients.read"]})

        # Without proposals.submit every document it reads has nowhere to go:
        # the process would run, read, and silently discard.
        check = preflight.check_cms(cms)

        assert check.is_fatal
        assert "proposals.submit" in check.detail

    def test_no_database_is_fatal(self):
        check = preflight.check_database(None)

        assert check.is_fatal
        assert "TAXPILOT_DATABASE_URL" in check.detail

    def test_an_unreachable_database_is_fatal(self):
        def connect():
            raise RuntimeError("could not connect to server")

        assert preflight.check_database(connect).is_fatal

    def test_missing_ocr_is_degraded_not_fatal(self):
        class Null:
            name = "null"

        check = preflight.check_ocr(Null())

        # Refusing to start would turn a degraded service into no service, and
        # the documents would still arrive.
        assert check.outcome is Outcome.DEGRADED
        assert not check.is_fatal

    def test_a_loaded_engine_passes(self):
        class Paddle:
            name = "paddleocr"

        assert preflight.check_ocr(Paddle()).outcome is Outcome.OK

    def test_every_check_runs_even_after_one_fails(self, caplog):
        checks = [
            Check("first", Outcome.FATAL, "broken"),
            Check("second", Outcome.FATAL, "also broken"),
        ]

        # Reporting only the first means an operator fixes it, restarts, and
        # discovers the second — several round trips one log could have saved.
        assert preflight.run(checks) is False
        assert "broken" in caplog.text and "also broken" in caplog.text

    def test_degraded_alone_still_allows_starting(self):
        assert preflight.run([Check("ocr", Outcome.DEGRADED, "no engine")]) is True


class TestContainer:
    def test_it_refuses_to_build_without_credentials(self):
        from app.config.settings import ConfigurationError

        with pytest.raises(ConfigurationError):
            Container(env={}).settings

    def test_the_registry_holds_the_tools_the_workflow_needs(self):
        tools = Container(env=BASE_ENV).tools

        for required in ("ocr_document", "classify_document", "extract_fields",
                         "search_client", "submit_filing_proposal", "remember"):
            assert tools.has(required), f"{required} is not registered"

    def test_the_autonomous_upload_tool_is_not_registered(self):
        # ADR-0004's promoted path files without a human. A deployment gets it
        # only by deciding to, never by default.
        assert not Container(env=BASE_ENV).tools.has("upload_document")

    def test_it_falls_back_to_in_memory_storage_loudly(self, caplog):
        import logging

        with caplog.at_level(logging.WARNING):
            store = Container(env=BASE_ENV).store

        assert type(store).__name__ == "InMemoryRunStore"
        assert "will not survive a restart" in caplog.text

    def test_no_database_url_means_no_connection_factory(self):
        assert Container(env=BASE_ENV).connect is None

    def test_the_poll_interval_is_configurable(self):
        container = Container(env={**BASE_ENV, "TAXPILOT_POLL_SECONDS": "5"})

        assert container.poll_seconds == 5.0


class TestEntryPoint:
    def test_missing_configuration_exits_with_the_config_code(self, monkeypatch):
        from app.__main__ import EXIT_MISCONFIGURED, main

        for name in ("TAXPILOT_CMS_URL", "TAXPILOT_AGENT_KEY", "TAXPILOT_AGENT_SECRET"):
            monkeypatch.delenv(name, raising=False)

        # 78 rather than 1: "you configured this wrong" is a different thing
        # from "it broke", and a deployment script should be able to tell them
        # apart without reading logs.
        assert main(["check"]) == EXIT_MISCONFIGURED

    def test_migrate_without_a_database_says_so(self, monkeypatch):
        from app.__main__ import EXIT_MISCONFIGURED, command_migrate

        assert command_migrate(Container(env=BASE_ENV)) == EXIT_MISCONFIGURED


class TestTheDaemonsJobsAreActuallyWired:
    """Every name the run loop reaches for has to exist.

    Found by starting the daemon, not by any test: `command_run` registered a
    job as `lambda: _recheck_compatibility(container)` while the function itself
    had never landed in the module. Nothing imported it, so nothing failed at
    import; the job simply raised NameError every fifteen minutes, and the
    daemon — which isolates job failures on purpose — carried on looking healthy.

    Unit tests never touch this wiring because `command_run` needs a real
    container, a database and a CMS. So this reads the function's own syntax
    tree instead and checks that every name it uses resolves.
    """

    def _free_names(self, function) -> set[str]:
        """Names the function USES without binding them itself.

        Everything it binds — parameters, assignments, `with ... as`,
        `except ... as`, and its own local imports — is subtracted, because those
        are supposed to be absent from the module. What remains has to come from
        somewhere else, and that somewhere is the module namespace.
        """
        import ast
        import inspect
        import textwrap

        tree = ast.parse(textwrap.dedent(inspect.getsource(function)))
        used, bound = set(), set()

        for node in ast.walk(tree):
            if isinstance(node, ast.Name):
                (used if isinstance(node.ctx, ast.Load) else bound).add(node.id)
            elif isinstance(node, ast.arg):
                bound.add(node.arg)
            elif isinstance(node, ast.ExceptHandler) and node.name:
                bound.add(node.name)
            elif isinstance(node, ast.alias):
                bound.add((node.asname or node.name).split(".")[0])

        return used - bound

    def test_every_name_command_run_uses_resolves(self):
        import builtins

        from app import __main__ as entry

        unresolved = sorted(
            name
            for name in self._free_names(entry.command_run)
            if not hasattr(entry, name) and not hasattr(builtins, name)
        )

        assert unresolved == [], f"command_run references names that do not exist: {unresolved}"

    def test_the_scheduled_jobs_are_callable(self):
        from app import __main__ as entry

        for name in ("_poll_decisions", "_recheck_compatibility", "_archive"):
            assert callable(getattr(entry, name, None)), f"{name} is missing or not callable"


    def test_every_expected_job_is_registered(self):
        """A job that is missing entirely has no name to fail on.

        The command processor was written, tested and wired to nothing: the
        whole queue was unreachable from the running daemon for a full
        milestone, and every unit test passed because they called the processor
        directly. `command_run` simply never mentioned it.

        The name-resolution test above cannot catch that — there is no
        undefined name when the code was never written. This reads the
        registrations instead and requires the set to be complete.
        """
        import ast
        import inspect
        import textwrap

        from app import __main__ as entry

        tree = ast.parse(textwrap.dedent(inspect.getsource(entry.command_run)))

        registered = {
            node.args[0].value
            for node in ast.walk(tree)
            if isinstance(node, ast.Call)
            and isinstance(node.func, ast.Attribute)
            and node.func.attr == "add"
            and node.args
            and isinstance(node.args[0], ast.Constant)
        }

        # "inbox" is the one this test was written for, twice over: the webhook
        # queued documents and nothing drained the queue, so a message arrived,
        # was accepted, and sat in memory for the life of the process.
        expected = {"decisions", "compatibility", "status", "commands", "archive", "inbox"}
        missing = expected - registered

        assert missing == set(), f"the daemon never registers: {sorted(missing)}"

    def test_commands_are_polled_far_more_often_than_anything_else(self):
        """A QR code answers somebody stood at a screen.

        At the decision interval they would wait half a minute to see it, which
        reads as the button not having worked.
        """
        from app import __main__ as entry

        assert entry.COMMAND_INTERVAL <= 10
        assert entry.COMMAND_INTERVAL < entry.STATUS_INTERVAL

    def test_the_intervals_are_sane(self):
        from app import __main__ as entry

        # Minutes, not seconds: a handshake every second would be a denial of
        # service against the installation it is checking on.
        assert entry.COMPATIBILITY_INTERVAL >= 60
        assert entry.STATUS_INTERVAL >= 60


class TestTheInboxJob:
    """Turning a queued message into a document on disk and a workflow run.

    Everything reaching this has already passed the self-chat rule — the
    receiver applies it before anything is queued — so these are about what
    happens to a message, not about who may send one.
    """

    def _message(
        self,
        filename: str | None = "statement.pdf",
        message_id: str = "evo-1",
        mime_type: str | None = "application/pdf",
    ):
        from app.whatsapp.messages import InboundMessage, MediaReference, MessageType

        return InboundMessage(
            provider_message_id=message_id,
            sender="923001234567",
            type=MessageType.DOCUMENT,
            media=MediaReference(handle="h", mime_type=mime_type, filename=filename),
        )

    def test_a_photograph_is_named_from_its_media_type(self):
        """The first real document sent through this was lost here.

        WhatsApp images carry no filename — only documents do — so every photo
        was saved as `.bin`. PaddleOCR refuses an extension it does not know,
        and because that engine never raises, the refusal arrived as a
        successful read of zero characters. The classifier said "unknown", the
        workflow correctly declined to guess, and nothing said the file had
        never been opened at all.

        Measured on the real photograph: 248 characters as .jpg, 0 as .bin.
        """
        from app import __main__ as entry

        name = entry._incoming_name(
            self._message(filename=None, mime_type="image/jpeg")
        )

        assert name == "evo-1.jpg"

    def test_the_media_type_is_preferred_over_the_providers_filename(self):
        # It is a constrained vocabulary rather than a path, and it is the field
        # WhatsApp actually fills in for an image.
        from app import __main__ as entry

        name = entry._incoming_name(
            self._message(filename="scan.pdf", mime_type="image/png")
        )

        assert name == "evo-1.png"

    def test_a_media_type_with_parameters_is_still_understood(self):
        from app import __main__ as entry

        name = entry._incoming_name(
            self._message(filename=None, mime_type="image/jpeg; charset=binary")
        )

        assert name == "evo-1.jpg"

    def test_an_unreadable_type_is_not_given_a_readable_extension(self):
        """Better a file the reader refuses than a lie about what it holds.

        Naming an unknown type `.jpg` turns "we do not handle this" into a page
        of noise in front of a reviewer.
        """
        from app import __main__ as entry

        name = entry._incoming_name(
            self._message(filename="notes.docx", mime_type="application/msword")
        )

        assert name == "evo-1.bin"

    def test_a_provider_filename_cannot_escape_the_incoming_directory(self):
        """`../../etc/passwd` is a perfectly valid WhatsApp filename.

        The name comes from a client, over an unofficial API, and is used to
        write a file. Only the extension survives, and the body is the
        provider's own message id.
        """
        from app import __main__ as entry

        name = entry._incoming_name(self._message(filename="../../../etc/passwd.pdf"))

        assert "/" not in name and "\\" not in name and ".." not in name
        assert name == "evo-1.pdf"

    def test_an_absurd_extension_is_replaced_rather_than_trusted(self):
        # No usable media type, so the filename is all there is — and it is a
        # client's, arriving over an unofficial API.
        from app import __main__ as entry

        for filename in ("x.thisisnotanextension", "x.<script>", "x."):
            name = entry._incoming_name(self._message(filename=filename, mime_type=None))

            assert name.endswith(".bin"), filename

    def test_a_message_describing_itself_as_nothing_gets_a_safe_name(self):
        from app import __main__ as entry

        assert entry._incoming_name(
            self._message(filename=None, mime_type=None)
        ) == "evo-1.bin"

    def test_a_message_id_full_of_punctuation_is_stripped(self):
        from app import __main__ as entry

        name = entry._incoming_name(self._message(filename="a.pdf", message_id="../../evil"))

        assert name == "evil.pdf"

    def test_text_in_the_self_chat_is_not_work(self):
        """Talking to yourself is an ordinary thing to do.

        Only a document is work. Treating a note as one would start a workflow
        with nothing to read and put a failure in front of a reviewer.
        """
        from app import __main__ as entry
        from app.whatsapp.messages import InboundMessage, MessageType
        from app.whatsapp.pending import InMemoryPendingPdfs
        from app.whatsapp.webhook import InboundQueue

        started: list = []

        class FakeContainer:
            inbound_queue = InboundQueue()
            whatsapp = None
            # Empty, which is the point: a text message with no document
            # waiting on it is still somebody talking to themselves.
            pending_pdfs = InMemoryPendingPdfs()

            class engine:
                @staticmethod
                def start(*args, **kwargs):
                    started.append(args)

        FakeContainer.inbound_queue.put(
            InboundMessage(provider_message_id="t1", sender="923001234567",
                           type=MessageType.TEXT, text="remember to file this")
        )

        entry._process_inbox(FakeContainer())

        assert started == []

    def test_one_failing_document_does_not_lose_the_rest(self):
        from app import __main__ as entry
        from app.intake import InMemoryOutbox as _InMemoryOutbox
        from app.whatsapp.pending import InMemoryPendingPdfs
        from app.whatsapp.webhook import InboundQueue
        from tests.intake_support import RecordingCms as _RecordingCms

        started: list = []

        from app.whatsapp.inbox import InboxFilter
        from app.whatsapp.messages import InboundMessage, MessageType

        class FakeContainer:
            inbound_queue = InboundQueue()
            incoming_dir = Path(tempfile.mkdtemp())
            pending_pdfs = InMemoryPendingPdfs()
            inbox = InboxFilter()
            # The classification stage is off unless a model is configured,
            # which is the state of every deployment today.
            model = None

            cms = _RecordingCms()
            intake_outbox = _InMemoryOutbox()

            class whatsapp:
                @staticmethod
                def download_media(media, destination):
                    if "bad" in str(destination):
                        raise RuntimeError("evolution refused")

                    destination.write_bytes(b"%PDF-")

                    return destination

            class engine:
                @staticmethod
                def start(workflow, context):
                    started.append(context["path"])

                    class Run:
                        id = "r1"

                        class state:
                            value = "succeeded"

                    return Run()

        for message_id in ("bad", "good"):
            FakeContainer.inbound_queue.put(
                TestTheInboxJob()._message(message_id=message_id)
            )

        entry._process_inbox(FakeContainer())

        # The batch continues. A provider that refuses one download must not
        # cost the document behind it — the good one is queued and waiting.
        FakeContainer.inbound_queue.put(
            InboundMessage(provider_message_id="t1", sender="923001234567",
                           type=MessageType.TEXT, text="file 1420")
        )

        entry._process_inbox(FakeContainer())

        assert len(started) == 1
        assert "good" in started[0]


class TestForwardingAnIdentifiedDocument:
    """Which workflow a self-chat attachment goes to, and what never reaches one.

    The routing rule is one line — does the caption name a client? — and every
    test here is about the consequences of getting it wrong in either direction:
    a document read when it should not have been, or a document silently dropped
    with nothing anywhere saying so.
    """

    def _container(self, tmp_path):
        from app.intake import InMemoryOutbox
        from app.whatsapp.inbox import InboxFilter
        from app.whatsapp.pending import InMemoryPendingPdfs
        from app.whatsapp.webhook import InboundQueue
        from tests.intake_support import RecordingCms

        started: list = []

        class FakeContainer:
            inbound_queue = InboundQueue()
            incoming_dir = tmp_path
            inbox = InboxFilter()
            pending_pdfs = InMemoryPendingPdfs()
            # No classification model. The stage is off unless one is
            # configured, which is the state of every deployment today.
            model = None

            # Since ADR-0010 a document does not proceed until the CMS has a
            # record of it, so a container without these processes nothing.
            cms = RecordingCms()
            intake_outbox = InMemoryOutbox()

            class whatsapp:
                @staticmethod
                def download_media(media, destination):
                    destination.write_bytes(
                        (media.raw or {}).get("content") or b"%PDF-1.4 statement"
                    )

                    return destination

            class engine:
                @staticmethod
                def start(workflow, context):
                    started.append((workflow.name, context))

                    class Run:
                        id = "r1"

                        class state:
                            value = "awaiting_approval"

                    return Run()

        return FakeContainer, started

    def _message(
        self,
        *,
        text: str = "",
        mime_type: str | None = "application/pdf",
        size_bytes: int | None = None,
        content: bytes | None = None,
        message_id: str = "evo-9",
        filename: str = "statement.pdf",
    ):
        from app.whatsapp.messages import InboundMessage, MediaReference, MessageType

        media = MediaReference(
            handle="h",
            mime_type=mime_type,
            filename=filename,
            size_bytes=size_bytes,
            # The provider's own fragment, which is where the fake downloader
            # finds the bytes it should write. MediaReference has slots, so
            # there is nowhere else to put them.
            raw={"content": content} if content is not None else {},
        )

        return InboundMessage(
            provider_message_id=message_id,
            sender="923049637232",
            type=MessageType.DOCUMENT,
            text=text,
            media=media,
        )

    def _name_it(self, container, text, message_id="txt-1"):
        """The plain text message that says whose the waiting PDF is."""
        from app.whatsapp.messages import InboundMessage, MessageType

        container.inbound_queue.put(
            InboundMessage(
                provider_message_id=message_id,
                sender="923049637232",
                type=MessageType.TEXT,
                text=text,
            )
        )

    def test_the_next_message_naming_a_client_skips_the_reader(self, tmp_path):
        """The type comes from the filename now, not from a caption.

        Under V1 an identifier in the caption was itself the signal that this
        was a forwarded statement. A forwarded PDF has no caption, so the
        filename is what routes it to the workflow that never opens it, and the
        text that follows only says whose it is.
        """
        from app import __main__ as entry

        container, started = self._container(tmp_path)
        container.inbound_queue.put(self._message(filename="Bank Statement July.pdf"))

        entry._process_inbox(container)

        # Held, not filed: nobody has said whose it is yet.
        assert started == []

        self._name_it(container, "file 1420")
        entry._process_inbox(container)

        assert [name for name, _ in started] == ["bank_statement_intake"]
        assert started[0][1]["file_number"] == "1420"

    def test_a_message_naming_nobody_still_goes_to_the_reader(self, tmp_path):
        """A named-nobody document is not an error; it is a queue item.

        The rule says never guess. A text that carries neither a file number nor
        a CNIC therefore files the document with no client attached, which puts
        it in front of a person instead of against the wrong record.
        """
        from app import __main__ as entry

        container, started = self._container(tmp_path)
        container.inbound_queue.put(self._message())
        entry._process_inbox(container)

        self._name_it(container, "here you go")
        entry._process_inbox(container)

        assert [name for name, _ in started] == ["document_intake"]
        assert started[0][1]["file_number"] is None
        assert started[0][1]["cnic"] is None

    def test_a_cnic_in_the_next_message_is_enough(self, tmp_path):
        from app import __main__ as entry

        container, started = self._container(tmp_path)
        container.inbound_queue.put(self._message(filename="Bank Statement July.pdf"))
        entry._process_inbox(container)

        self._name_it(container, "35202-1234567-1")
        entry._process_inbox(container)

        assert started[0][0] == "bank_statement_intake"
        assert started[0][1]["cnic"] == "3520212345671"

    def test_it_carries_the_message_the_cms_deduplicates_on(self, tmp_path):
        """The document's own id, not the id of the text that named it.

        The CMS deduplicates on this. Recording the text message would let the
        same statement be filed twice by naming it twice.
        """
        from app import __main__ as entry

        container, started = self._container(tmp_path)
        container.inbound_queue.put(self._message(message_id="WA-77"))
        entry._process_inbox(container)

        self._name_it(container, "file 1420", message_id="WA-78")
        entry._process_inbox(container)

        source = started[0][1]["source"]

        assert source["message_id"] == "WA-77"
        assert source["channel"] == "whatsapp"
        assert source["number"] == "923049637232"

    def test_an_unsupported_type_is_refused_and_counted(self, tmp_path):
        """A refused attachment leaves no proposal — so it must leave a count.

        From the Approval Queue's side nothing happened, and the person who
        forwarded it is waiting for a document that is not coming. The counter
        on the operations console is the only sign it ever existed.
        """
        from app import __main__ as entry

        container, started = self._container(tmp_path)
        container.inbound_queue.put(
            self._message(text="file 1420", mime_type="application/msword")
        )

        entry._process_inbox(container)

        assert started == []
        assert container.inbox.forwards_rejected == 1

    def test_a_file_larger_than_the_cms_accepts_is_not_downloaded(self, tmp_path):
        from app import __main__ as entry

        container, started = self._container(tmp_path)
        container.inbound_queue.put(
            self._message(text="file 1420", size_bytes=21 * 1024 * 1024)
        )

        entry._process_inbox(container)

        # Refused before the download. The alternative is reading 40 MB into
        # memory, expanding it to 53 MB of base64, and posting it to be refused.
        assert started == []
        assert container.inbox.forwards_rejected == 1
        assert list(tmp_path.iterdir()) == []

    def test_a_file_that_is_not_what_it_claims_is_refused(self, tmp_path):
        from app import __main__ as entry

        container, started = self._container(tmp_path)
        container.inbound_queue.put(
            self._message(text="file 1420", content=b"this is not a pdf")
        )

        entry._process_inbox(container)

        assert started == []
        assert container.inbox.forwards_rejected == 1

    def test_a_refused_file_is_not_left_on_disk(self, tmp_path):
        from app import __main__ as entry

        container, _ = self._container(tmp_path)
        container.inbound_queue.put(
            self._message(text="file 1420", content=b"truncated")
        )

        entry._process_inbox(container)

        assert list(tmp_path.iterdir()) == []

    def test_a_document_that_never_arrived_is_counted_separately(self, tmp_path):
        """Refused and never-fetched are different problems.

        Found on staging: a forward whose download failed moved neither counter,
        so the only trace was a traceback in the daemon log. Folding it into
        `forwards_rejected` would be worse than nothing — that number tells an
        operator to go and ask whoever sent the file, and the file is not the
        problem.
        """
        from app import __main__ as entry

        container, started = self._container(tmp_path)

        class Broken:
            @staticmethod
            def download_media(media, destination):
                raise RuntimeError("Evolution refused the media request (HTTP 400).")

        container.whatsapp = Broken
        container.inbound_queue.put(self._message(text="file 1420"))

        entry._process_inbox(container)

        assert started == []
        assert container.inbox.forwards_failed == 1
        # Not this one. The file was never seen, so nothing about it was refused.
        assert container.inbox.forwards_rejected == 0
        assert container.inbox.forwards_accepted == 0

    def test_a_failed_download_still_leaves_a_traceback(self, tmp_path, caplog):
        """The counter is additional to the log, not a replacement for it.

        The count says something went wrong; the traceback carries the
        provider's HTTP status, which is the actual diagnosis. Swallowing the
        exception to count it would trade the second for the first.
        """
        from app import __main__ as entry

        container, _ = self._container(tmp_path)

        class Broken:
            @staticmethod
            def download_media(media, destination):
                raise RuntimeError("Evolution refused the media request (HTTP 400).")

        container.whatsapp = Broken
        container.inbound_queue.put(self._message(text="file 1420"))

        with caplog.at_level("ERROR"):
            entry._process_inbox(container)

        assert "Could not process message" in caplog.text
        assert "HTTP 400" in caplog.text

    def test_one_failed_download_does_not_lose_the_forward_behind_it(self, tmp_path):
        from app import __main__ as entry

        container, started = self._container(tmp_path)

        class Flaky:
            @staticmethod
            def download_media(media, destination):
                if "bad" in str(destination):
                    raise RuntimeError("Evolution refused the media request (HTTP 400).")

                destination.write_bytes(b"%PDF-1.4 statement")

                return destination

        container.whatsapp = Flaky

        for message_id in ("bad", "good"):
            container.inbound_queue.put(
                self._message(message_id=message_id, filename="Bank Statement July.pdf")
            )

        entry._process_inbox(container)

        # The failed one never reached the queue, so the naming message belongs
        # to the one that did.
        assert container.inbox.forwards_failed == 1
        assert container.inbox.forwards_accepted == 1

        self._name_it(container, "file 1420")
        entry._process_inbox(container)

        assert len(started) == 1
        assert started[0][1]["source"]["message_id"] == "good"

    def test_an_accepted_forward_is_counted(self, tmp_path):
        from app import __main__ as entry

        container, _ = self._container(tmp_path)
        container.inbound_queue.put(self._message(text="file 1420"))

        entry._process_inbox(container)

        # All three. "It refused everything", "it received nothing" and "it
        # could not fetch anything" look identical with fewer.
        assert container.inbox.forwards_accepted == 1
        assert container.inbox.forwards_rejected == 0
        assert container.inbox.forwards_failed == 0

    def test_a_jpeg_forward_is_accepted(self, tmp_path):
        from app import __main__ as entry

        container, started = self._container(tmp_path)
        container.inbound_queue.put(
            self._message(text="file 1420", mime_type="image/jpeg",
                          content=b"\xff\xd8\xff\xe0 jpeg")
        )

        entry._process_inbox(container)

        assert started[0][0] == "bank_statement_intake"

    def test_a_message_id_full_of_punctuation_cannot_escape_the_directory(self, tmp_path):
        from app import __main__ as entry

        container, started = self._container(tmp_path)
        container.inbound_queue.put(
            self._message(message_id="../../evil", filename="Bank Statement July.pdf")
        )

        entry._process_inbox(container)

        written = [p.name for p in tmp_path.iterdir()]

        assert written == ["______evil.pdf"]

        self._name_it(container, "file 1420")
        entry._process_inbox(container)

        assert started[0][0] == "bank_statement_intake"
