← back to Unclaimed Property Platform

tests/test_cycle10_outbox_worker.py

104 lines

"""Cycle 9 test — outbox worker: at-least-once delivery + bounded retry.

Run:  python -m tests.test_cycle10_outbox_worker

Proves:
  1. a transient state-API failure leaves the row pending (claim not advanced), and a later
     drain succeeds -> claim SUBMITTED_TO_STATE, outbox 'sent';
  2. re-draining after success is a no-op (idempotent);
  3. permanent failure exhausts retries -> outbox 'failed', claim NOT lost (stays SUBMITTING).
"""
from __future__ import annotations

import sys
import tempfile
from pathlib import Path
from uuid import uuid4

sys.path.insert(0, str(Path(__file__).resolve().parents[1]))

from services.claims.claim_workflow import (
    ClaimStatus, queue_state_submission, transition_claim,
)
from services.claims.outbox_worker import OutboxWorker
from services.claims.sqlite_claim_repo import SqliteClaimRepository
from services.common.sqlite_repo import SqliteRepository


class FlakyAdapter:
    def __init__(self, fail_times: int) -> None:
        self.calls = 0
        self.fail_times = fail_times

    def submit_claim(self, claim, idempotency_key) -> str:
        self.calls += 1
        if self.calls <= self.fail_times:
            raise RuntimeError("state API timeout")
        return "STATE-CASE-9"


class AlwaysFail:
    def submit_claim(self, claim, idempotency_key) -> str:
        raise RuntimeError("state API permanently down")


def ok(msg: str) -> None:
    print(f"  ✓ {msg}")


def _ready_claim(repo: SqliteClaimRepository):
    c = repo.create_claim("SAMPLE", public_property_reference="ref", claimant_id=uuid4())
    for target, key in [(ClaimStatus.IDENTITY_PENDING, "k1"),
                        (ClaimStatus.EVIDENCE_PENDING, "k2"),
                        (ClaimStatus.READY_FOR_SUBMISSION, "k3")]:
        transition_claim(repo, c.claim_id, target, actor_id="worker", idempotency_key=key)
    queue_state_submission(repo, c.claim_id, actor_id="worker", idempotency_key="sub")
    return c


def main() -> int:
    tmp = Path(tempfile.mkdtemp(prefix="upp-cycle10-"))
    base = SqliteRepository(str(tmp / "p.db"))
    repo = SqliteClaimRepository(base.conn)

    print("1) transient failure -> retry -> success")
    c1 = _ready_claim(repo)
    assert repo.get_for_update(c1.claim_id).status == ClaimStatus.SUBMITTING
    worker = OutboxWorker(repo, FlakyAdapter(fail_times=1), max_attempts=3)
    r1 = worker.drain_once()
    assert r1["retried"] == 1 and r1["sent"] == 0, r1
    assert repo.get_for_update(c1.claim_id).status == ClaimStatus.SUBMITTING, "must not advance on failure"
    ok(f"drain#1 adapter failed -> retried=1, claim still SUBMITTING")
    r2 = worker.drain_once()
    assert r2["sent"] == 1, r2
    fin = repo.get_for_update(c1.claim_id)
    assert fin.status == ClaimStatus.SUBMITTED_TO_STATE and fin.state_case_id == "STATE-CASE-9", fin
    ok("drain#2 adapter recovered -> SUBMITTED_TO_STATE, outbox sent")

    print("2) re-drain after success is a no-op")
    r3 = worker.drain_once()
    assert r3 == {"sent": 0, "failed": 0, "retried": 0}, r3
    ok("no pending rows -> nothing re-delivered")

    print("3) permanent failure exhausts retries -> failed, claim not lost")
    c2 = _ready_claim(repo)
    w2 = OutboxWorker(repo, AlwaysFail(), max_attempts=2)
    a = w2.drain_once()   # attempt 1 -> pending
    assert a["retried"] == 1, a
    b = w2.drain_once()   # attempt 2 -> failed
    assert b["failed"] == 1, b
    # find c2's outbox row state
    oid = [row for row in base.conn.execute(
        "SELECT outbox_event_id FROM outbox_event WHERE aggregate_id=?",
        (str(c2.claim_id),))][0]["outbox_event_id"]
    assert repo.outbox_row(oid)["delivery_state"] == "failed"
    assert repo.get_for_update(c2.claim_id).status == ClaimStatus.SUBMITTING, "claim must not be lost"
    ok("2 attempts exhausted -> outbox 'failed', claim still SUBMITTING (recoverable)")

    print("\nALL CYCLE-9 OUTBOX-WORKER ASSERTIONS PASSED ✅")
    return 0


if __name__ == "__main__":
    raise SystemExit(main())