"""One claim, one owner (ADR-0011).

The drain used to read the oldest row and delete it once the CMS acknowledged,
with nothing in between. Two workers read the same row, both registered the
document, and the stress harness produced INT-2808 and INT-2809 for one file.

An idempotency key made that harmless. It did not make it correct: two workers
believing they own the same unit of work is a correctness problem the moment
their action is not idempotent, and every future action here is one nobody has
checked yet.

A claim is a lease rather than a lock — held in the row so it outlives the
process, and expiring so a worker that dies does not strand a client's document.
That combination is what keeps at-least-once delivery while making duplicate
*work* rare rather than routine.
"""

from __future__ import annotations

import time
from datetime import UTC, datetime, timedelta

import pytest

from app.intake import InMemoryOutbox, PendingRegistration
from app.intake.outbox import CLAIM_LEASE_SECONDS


def arrival(message_id: str, received_at: datetime | None = None) -> PendingRegistration:
    return PendingRegistration(
        source="whatsapp",
        provider_message_id=message_id,
        original_filename=f"{message_id}.pdf",
        path=f"incoming/{message_id}.pdf",
        received_at=received_at or datetime.now(UTC),
    )


class TestOnlyOneWorkerGetsAnItem:
    def test_two_workers_never_take_the_same_row(self):
        outbox = InMemoryOutbox()
        outbox.add(arrival("A"))
        outbox.add(arrival("B"))

        first = outbox.claim("worker-1")
        second = outbox.claim("worker-2")

        # Different documents, not the same one twice. This is the assertion the
        # whole mechanism exists for.
        assert first.provider_message_id != second.provider_message_id
        assert {first.provider_message_id, second.provider_message_id} == {"A", "B"}

    def test_a_second_worker_finds_nothing_when_everything_is_claimed(self):
        outbox = InMemoryOutbox()
        outbox.add(arrival("only"))

        assert outbox.claim("worker-1") is not None

        # Nothing rather than the same row again. A worker with no work waits,
        # which is correct; a worker handed somebody else's work is the bug.
        assert outbox.claim("worker-2") is None

    def test_the_claim_records_who_holds_it(self):
        outbox = InMemoryOutbox()
        outbox.add(arrival("A"))

        # So an operator looking at a stuck row can tell "a worker has this" from
        # "a worker had this and died".
        assert outbox.claim("worker-7").claimed_by == "worker-7"

    def test_the_oldest_is_still_taken_first(self):
        outbox = InMemoryOutbox()
        now = datetime.now(UTC)

        outbox.add(arrival("NEW", now))
        outbox.add(arrival("OLD", now - timedelta(hours=3)))

        # Claiming must not reorder the queue. The document somebody is asking
        # about is the one that has been waiting longest.
        assert outbox.claim("worker-1").provider_message_id == "OLD"


class TestAbandonedWorkComesBack:
    def test_a_crashed_worker_does_not_strand_the_document(self):
        outbox = InMemoryOutbox()
        outbox.add(arrival("A"))

        outbox.claim("worker-that-dies")

        # Nothing released it; the process is simply gone. With a lock that would
        # be the end of the document.
        assert outbox.claim("worker-2") is None

        # With a lease it is a delay. Expiry is checked against the lease rather
        # than slept through, so the test measures the rule and not the clock.
        recovered = outbox.claim("worker-2", lease=0)

        assert recovered is not None
        assert recovered.provider_message_id == "A"

    def test_the_document_is_still_there_after_the_claim_expires(self):
        outbox = InMemoryOutbox()
        outbox.add(arrival("A"))

        outbox.claim("worker-that-dies")

        # Claiming is not consuming. At-least-once means the row survives every
        # attempt that did not end in the CMS acknowledging it.
        assert len(outbox) == 1

    def test_a_failure_hands_the_work_straight_back(self):
        outbox = InMemoryOutbox()
        outbox.add(arrival("A"))

        held = outbox.claim("worker-1")
        outbox.failed(held, "connection refused")

        # No waiting for the lease. The lease is for workers that vanished, not
        # for ones that failed and said so.
        again = outbox.claim("worker-2")

        assert again is not None
        assert again.attempts == 1
        assert again.last_error == "connection refused"

    def test_releasing_does_not_count_an_attempt(self):
        outbox = InMemoryOutbox()
        outbox.add(arrival("A"))

        held = outbox.claim("worker-1")
        outbox.release(held)

        # For a worker shutting down cleanly: it never tried, so the retry count
        # must not suggest the document is troublesome.
        assert outbox.claim("worker-2").attempts == 0


class TestDeliveryGuarantees:
    def test_the_row_only_goes_when_the_cms_has_it(self):
        outbox = InMemoryOutbox()
        outbox.add(arrival("A"))

        held = outbox.claim("worker-1")
        assert len(outbox) == 1

        outbox.done(held)
        assert len(outbox) == 0

    def test_a_retry_reuses_the_same_idempotency_key(self):
        outbox = InMemoryOutbox()
        outbox.add(arrival("STABLE-KEY"))

        first = outbox.claim("worker-1")
        outbox.failed(first, "boom")
        second = outbox.claim("worker-2")

        # This is what makes at-least-once safe. The key travels on the row, so a
        # redelivery after an expired claim is the same document to the CMS —
        # which returns the record it already made instead of making another.
        assert first.provider_message_id == second.provider_message_id == "STABLE-KEY"

    def test_a_document_survives_being_claimed_and_abandoned_repeatedly(self):
        outbox = InMemoryOutbox()
        outbox.add(arrival("PERSISTENT"))

        for _ in range(5):
            outbox.claim("worker-that-dies", lease=0)

        # Five crashes, one document, still queued. Zero lost is the guarantee;
        # the attempts are the cost of keeping it.
        assert len(outbox) == 1
        assert outbox.claim("worker-final", lease=0).provider_message_id == "PERSISTENT"


class TestTheDrainUsesTheClaim:
    def test_two_drains_in_a_row_take_different_documents(self, tmp_path):
        from app import __main__ as entry
        from tests.intake_support import RecordingCms

        class FakeContainer:
            cms = RecordingCms()
            intake_outbox = InMemoryOutbox()

        FakeContainer.intake_outbox.add(arrival("A"))
        FakeContainer.intake_outbox.add(arrival("B"))

        entry._drain_outbox(FakeContainer)
        entry._drain_outbox(FakeContainer)

        sent = [r["provider_message_id"] for r in FakeContainer.cms.registered]

        # One pass, one document, and never the same one twice.
        assert sorted(sent) == ["A", "B"]
        assert len(FakeContainer.intake_outbox) == 0

    def test_the_default_lease_is_long_enough_to_be_useful(self):
        # Long enough to outlast a slow registration carrying a large file to a
        # struggling CMS; short enough that a crash does not strand a client's
        # document for an afternoon.
        assert 60 <= CLAIM_LEASE_SECONDS <= 900
