← 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())