From c8f5e6aca92e2fc3c0e43e3535ee8a2c50275474 Mon Sep 17 00:00:00 2001 From: Mikhail Kugan <68385536+mike-k-git@users.noreply.github.com> Date: Sun, 5 Jul 2026 22:21:41 +0200 Subject: [PATCH 1/9] feat: scaffold replay endpoint --- backend/app/routers/deliveries.py | 44 ++++++++++++++++++++++++++++--- 1 file changed, 40 insertions(+), 4 deletions(-) diff --git a/backend/app/routers/deliveries.py b/backend/app/routers/deliveries.py index fb41723..b7b4ba9 100644 --- a/backend/app/routers/deliveries.py +++ b/backend/app/routers/deliveries.py @@ -1,17 +1,53 @@ +import uuid from typing import Annotated -from fastapi import APIRouter, Query -from sqlalchemy import select +from fastapi import APIRouter, HTTPException, Query, status +from sqlalchemy import and_, select, update from sqlalchemy.orm import selectinload from app.db import SessionDep from app.models import Delivery, DeliveryStatus from app.schemas import DeliveryInboxItem, DeliveryRead, EventRead -router = APIRouter(prefix="/deliveries/dead_letter", tags=["deliveries"]) +router = APIRouter(prefix="/deliveries", tags=["deliveries"]) -@router.get("", response_model=list[DeliveryInboxItem]) +@router.post("/{delivery_id}/replay", status_code=status.HTTP_202_ACCEPTED) +async def replay(delivery_id: uuid.UUID, session: SessionDep): + owner = ( + await session.execute( + update(Delivery) + .where( + and_( + Delivery.id == delivery_id, + Delivery.status == DeliveryStatus.dead_letter, + ) + ) + .values( + status=DeliveryStatus.pending, + next_attempt_at=None, + locked_by=None, + locked_until=None, + ) + .returning(Delivery.id) + ) + ).one_or_none() + + if owner is not None: + await session.commit() + return + + delivery = ( + await session.execute(select(Delivery).where(Delivery.id == delivery_id)) + ).scalar_one_or_none() + + if delivery is None: + raise HTTPException(status.HTTP_404_NOT_FOUND, detail="unknown delivery") + + raise HTTPException(status.HTTP_409_CONFLICT, detail=f"{delivery.status.value}") + + +@router.get("/dead_letter", response_model=list[DeliveryInboxItem]) async def inbox( session: SessionDep, limit: Annotated[int, Query(le=200)] = 50, From 070c7a4d8d0b0bf67a743c36ff86049215a9d347 Mon Sep 17 00:00:00 2001 From: Mikhail Kugan <68385536+mike-k-git@users.noreply.github.com> Date: Mon, 6 Jul 2026 00:05:56 +0200 Subject: [PATCH 2/9] refactor: extract fastapi dependencies --- backend/app/db.py | 13 ---- backend/app/deps.py | 23 +++++++ backend/app/routers/config.py | 2 +- backend/app/routers/deliveries.py | 2 +- backend/app/routers/destinations.py | 2 +- backend/app/routers/events.py | 2 +- backend/app/routers/ingest.py | 7 +- backend/app/routers/routes.py | 2 +- backend/tests/conftest.py | 7 +- backend/tests/test_ingest.py | 100 +++++++++++++--------------- 10 files changed, 83 insertions(+), 77 deletions(-) create mode 100644 backend/app/deps.py diff --git a/backend/app/db.py b/backend/app/db.py index 2bad4e2..c8bc59a 100644 --- a/backend/app/db.py +++ b/backend/app/db.py @@ -1,9 +1,4 @@ -from collections.abc import AsyncIterator -from typing import Annotated - -from fastapi import Depends from sqlalchemy.ext.asyncio import ( - AsyncSession, async_sessionmaker, create_async_engine, ) @@ -18,11 +13,3 @@ class Base(DeclarativeBase): pass - - -async def get_session() -> AsyncIterator[AsyncSession]: - async with AsyncSessionLocal() as session: - yield session - - -SessionDep = Annotated[AsyncSession, Depends(get_session)] diff --git a/backend/app/deps.py b/backend/app/deps.py new file mode 100644 index 0000000..2764c67 --- /dev/null +++ b/backend/app/deps.py @@ -0,0 +1,23 @@ +from collections.abc import AsyncIterator +from typing import Annotated + +from fastapi import Depends, Request +from saq import Queue +from sqlalchemy.ext.asyncio import AsyncSession + +from app.db import AsyncSessionLocal + + +def get_queue(request: Request) -> Queue: + return request.app.state.queue + + +QueueDep = Annotated[Queue, Depends(get_queue)] + + +async def get_session() -> AsyncIterator[AsyncSession]: + async with AsyncSessionLocal() as session: + yield session + + +SessionDep = Annotated[AsyncSession, Depends(get_session)] diff --git a/backend/app/routers/config.py b/backend/app/routers/config.py index 425120b..f794867 100644 --- a/backend/app/routers/config.py +++ b/backend/app/routers/config.py @@ -2,7 +2,7 @@ from sqlalchemy import select from sqlalchemy.exc import IntegrityError -from app.db import SessionDep +from app.deps import SessionDep from app.models import Source from app.schemas import SourceCreate, SourceRead diff --git a/backend/app/routers/deliveries.py b/backend/app/routers/deliveries.py index b7b4ba9..e0f0755 100644 --- a/backend/app/routers/deliveries.py +++ b/backend/app/routers/deliveries.py @@ -5,7 +5,7 @@ from sqlalchemy import and_, select, update from sqlalchemy.orm import selectinload -from app.db import SessionDep +from app.deps import SessionDep from app.models import Delivery, DeliveryStatus from app.schemas import DeliveryInboxItem, DeliveryRead, EventRead diff --git a/backend/app/routers/destinations.py b/backend/app/routers/destinations.py index 21dc007..d284ed6 100644 --- a/backend/app/routers/destinations.py +++ b/backend/app/routers/destinations.py @@ -3,7 +3,7 @@ from fastapi import APIRouter, HTTPException, status from sqlalchemy import select -from app.db import SessionDep +from app.deps import SessionDep from app.models import Destination from app.schemas import DestinationCreate, DestinationRead, DestinationUpdate diff --git a/backend/app/routers/events.py b/backend/app/routers/events.py index dc72447..65a4006 100644 --- a/backend/app/routers/events.py +++ b/backend/app/routers/events.py @@ -5,7 +5,7 @@ from sqlalchemy import select, tuple_ from sqlalchemy.orm import selectinload -from app.db import SessionDep +from app.deps import SessionDep from app.models import Delivery, DeliveryStatus, Event, Source from app.pagination import CursorError, decode_cursor, encode_cursor from app.queries import event_rollups diff --git a/backend/app/routers/ingest.py b/backend/app/routers/ingest.py index 54530fe..bb696f7 100644 --- a/backend/app/routers/ingest.py +++ b/backend/app/routers/ingest.py @@ -8,7 +8,7 @@ from sqlalchemy import select from sqlalchemy.exc import IntegrityError -from app.db import SessionDep +from app.deps import QueueDep, SessionDep from app.models import Delivery, Event, Route, Source from app.schemas import IngestAck from app.security import verify @@ -27,6 +27,7 @@ async def ingest( source_name: str, request: Request, session: SessionDep, + queue: QueueDep, response: Response, x_webhook_signature: Annotated[str | None, Header()] = None, idempotency_key: Annotated[str | None, Header()] = None, @@ -98,9 +99,7 @@ async def ingest( try: for delivery in deliveries: - await request.app.state.queue.enqueue( - "deliver", delivery_id=str(delivery.id) - ) + await queue.enqueue("deliver", delivery_id=str(delivery.id)) except Exception: logger.warning("enqueue failed for event %s; sweeper will recover", event.id) return IngestAck(event_id=event.id) diff --git a/backend/app/routers/routes.py b/backend/app/routers/routes.py index 49ec149..9f06f80 100644 --- a/backend/app/routers/routes.py +++ b/backend/app/routers/routes.py @@ -4,7 +4,7 @@ from sqlalchemy import select from sqlalchemy.exc import IntegrityError -from app.db import SessionDep +from app.deps import SessionDep from app.models import Destination, Route, Source from app.schemas import RouteCreate, RouteRead diff --git a/backend/tests/conftest.py b/backend/tests/conftest.py index 9ffb57b..42fea21 100644 --- a/backend/tests/conftest.py +++ b/backend/tests/conftest.py @@ -7,7 +7,8 @@ from httpx import ASGITransport, AsyncClient from sqlalchemy.ext.asyncio import async_sessionmaker, create_async_engine -from app.db import Base, get_session +from app.db import Base +from app.deps import get_queue, get_session from app.main import app from app.models import ( Delivery, @@ -18,6 +19,7 @@ Source, ) from app.security import sign +from tests.fakes import FakeQueue TEST_DATABASE_URL = os.environ["TEST_DATABASE_URL"] @@ -55,7 +57,10 @@ async def override_get_session(): async with maker() as session: yield session + override_get_queue = FakeQueue() + app.dependency_overrides[get_session] = override_get_session + app.dependency_overrides[get_queue] = lambda: override_get_queue transport = ASGITransport(app=app) async with AsyncClient(transport=transport, base_url="http://test") as ac: yield ac diff --git a/backend/tests/test_ingest.py b/backend/tests/test_ingest.py index bae2431..06d05a9 100644 --- a/backend/tests/test_ingest.py +++ b/backend/tests/test_ingest.py @@ -1,3 +1,4 @@ +from app.deps import get_queue from app.main import app from app.models import DeliveryStatus from tests.fakes import FakeQueue, FakeRaisingQueue @@ -70,73 +71,64 @@ async def test_ingest_enqueue(client, source, signed): await _setup_routes(client, source) raw, headers = signed({"type": "payment.succeeded"}, key="evt_1") - app.state.queue = FakeQueue() - - try: - r = await client.post(f"/ingest/{source}", content=raw, headers=headers) - assert r.status_code == 202 - event_id = r.json()["event_id"] + q = FakeQueue() + app.dependency_overrides[get_queue] = lambda: q + r = await client.post(f"/ingest/{source}", content=raw, headers=headers) + assert r.status_code == 202 + event_id = r.json()["event_id"] - r2 = await client.get(f"/events/{event_id}") - assert r2.status_code == 200 + r2 = await client.get(f"/events/{event_id}") + assert r2.status_code == 200 - event_details = r2.json() - assert event_details["id"] == event_id - assert {e["delivery_id"] for e in app.state.queue.enqueued} == { - event_details["deliveries"][0]["id"], - event_details["deliveries"][1]["id"], - } - finally: - del app.state.queue + event_details = r2.json() + assert event_details["id"] == event_id + assert {e["delivery_id"] for e in q.enqueued} == { + event_details["deliveries"][0]["id"], + event_details["deliveries"][1]["id"], + } async def test_ingest_failed_enqueue(client, source, signed): await _setup_routes(client, source) raw, headers = signed({"type": "payment.succeeded"}, key="evt_1") - app.state.queue = FakeRaisingQueue(fail_for=set(), fail_all=True) - - try: - r = await client.post(f"/ingest/{source}", content=raw, headers=headers) - assert r.status_code == 202 - event_id = r.json()["event_id"] + q = FakeRaisingQueue(fail_for=set(), fail_all=True) + app.dependency_overrides[get_queue] = lambda: q + r = await client.post(f"/ingest/{source}", content=raw, headers=headers) + assert r.status_code == 202 + event_id = r.json()["event_id"] - r2 = await client.get(f"/events/{event_id}") - assert r2.status_code == 200 + r2 = await client.get(f"/events/{event_id}") + assert r2.status_code == 200 - event_details = r2.json() - assert event_details["id"] == event_id - assert event_details["deliveries"][0]["status"] == DeliveryStatus.pending - assert event_details["deliveries"][1]["status"] == DeliveryStatus.pending - assert app.state.queue.enqueued == [] - finally: - del app.state.queue + event_details = r2.json() + assert event_details["id"] == event_id + assert event_details["deliveries"][0]["status"] == DeliveryStatus.pending + assert event_details["deliveries"][1]["status"] == DeliveryStatus.pending + assert q.enqueued == [] async def test_ingest_no_enqueue_on_duplicate(client, source, signed): await _setup_routes(client, source) raw, headers = signed({"type": "payment.succeeded"}, key="evt_1") - app.state.queue = FakeQueue() - - try: - r = await client.post(f"/ingest/{source}", content=raw, headers=headers) - assert r.status_code == 202 - event_id = r.json()["event_id"] - - r2 = await client.post(f"/ingest/{source}", content=raw, headers=headers) - assert r2.status_code == 200 - assert event_id == r2.json()["event_id"] - - r3 = await client.get(f"/events/{event_id}") - assert r3.status_code == 200 - - event_details = r3.json() - assert event_details["id"] == event_id - assert {e["delivery_id"] for e in app.state.queue.enqueued} == { - event_details["deliveries"][0]["id"], - event_details["deliveries"][1]["id"], - } - assert len(app.state.queue.enqueued) == 2 - finally: - del app.state.queue + q = FakeQueue() + app.dependency_overrides[get_queue] = lambda: q + r = await client.post(f"/ingest/{source}", content=raw, headers=headers) + assert r.status_code == 202 + event_id = r.json()["event_id"] + + r2 = await client.post(f"/ingest/{source}", content=raw, headers=headers) + assert r2.status_code == 200 + assert event_id == r2.json()["event_id"] + + r3 = await client.get(f"/events/{event_id}") + assert r3.status_code == 200 + + event_details = r3.json() + assert event_details["id"] == event_id + assert {e["delivery_id"] for e in q.enqueued} == { + event_details["deliveries"][0]["id"], + event_details["deliveries"][1]["id"], + } + assert len(q.enqueued) == 2 From 1ac23a47c38ef40afd48757bf115e73776cf31f7 Mon Sep 17 00:00:00 2001 From: Mikhail Kugan <68385536+mike-k-git@users.noreply.github.com> Date: Mon, 6 Jul 2026 00:43:30 +0200 Subject: [PATCH 3/9] feat: enqueu on replay --- backend/app/routers/deliveries.py | 17 +++++++++++++++-- 1 file changed, 15 insertions(+), 2 deletions(-) diff --git a/backend/app/routers/deliveries.py b/backend/app/routers/deliveries.py index e0f0755..a6ba8f5 100644 --- a/backend/app/routers/deliveries.py +++ b/backend/app/routers/deliveries.py @@ -1,3 +1,4 @@ +import logging import uuid from typing import Annotated @@ -5,15 +6,17 @@ from sqlalchemy import and_, select, update from sqlalchemy.orm import selectinload -from app.deps import SessionDep +from app.deps import QueueDep, SessionDep from app.models import Delivery, DeliveryStatus from app.schemas import DeliveryInboxItem, DeliveryRead, EventRead +logger = logging.getLogger(__name__) + router = APIRouter(prefix="/deliveries", tags=["deliveries"]) @router.post("/{delivery_id}/replay", status_code=status.HTTP_202_ACCEPTED) -async def replay(delivery_id: uuid.UUID, session: SessionDep): +async def replay(delivery_id: uuid.UUID, session: SessionDep, queue: QueueDep): owner = ( await session.execute( update(Delivery) @@ -35,6 +38,16 @@ async def replay(delivery_id: uuid.UUID, session: SessionDep): if owner is not None: await session.commit() + try: + await queue.enqueue( + "deliver", + delivery_id=str(delivery_id), + key=f"deliver:{delivery_id!s}", + ) + except Exception: + logger.warning( + "replay enqueue failed for %s; sweeper will recover", delivery_id + ) return delivery = ( From 7870bc7dd369664fee6692fe256ad785d42943a7 Mon Sep 17 00:00:00 2001 From: Mikhail Kugan <68385536+mike-k-git@users.noreply.github.com> Date: Mon, 6 Jul 2026 01:14:45 +0200 Subject: [PATCH 4/9] refactor(tests): consolidate fakes --- backend/tests/fakes.py | 65 ++++++++++++++++++- backend/tests/test_claim.py | 47 ++++++-------- backend/tests/test_delivery_retry.py | 95 +++++++--------------------- backend/tests/test_sweeper.py | 23 +++---- 4 files changed, 111 insertions(+), 119 deletions(-) diff --git a/backend/tests/fakes.py b/backend/tests/fakes.py index f904379..41bfdfc 100644 --- a/backend/tests/fakes.py +++ b/backend/tests/fakes.py @@ -1,8 +1,11 @@ -from typing import override +from datetime import UTC, datetime, timedelta +from typing import cast, override import httpx +from sqlalchemy import update -from app.tasks import DeliveryResult, DeliverySnapshot, SendFn +from app.models import Delivery +from app.tasks import DeliveryResult, DeliverySnapshot, SendFn, WorkerContext, deliver class FakeQueue: @@ -39,3 +42,61 @@ async def _send_fn( class FakeWorker: def __init__(self, queue) -> None: self.queue = queue + + +def ctx(*, queue=None, client=None, sessionmaker=None) -> WorkerContext: + worker_context = {} + if queue is not None: + worker_context["worker"] = FakeWorker(queue) + if client is not None: + worker_context["client"] = client + if sessionmaker is not None: + worker_context["sessionmaker"] = sessionmaker + return cast(WorkerContext, worker_context) + + +def fail_result() -> DeliveryResult: + return DeliveryResult( + success=False, + response_status=500, + error=None, + duration_ms=10, + ) + + +def ok_result() -> DeliveryResult: + return DeliveryResult( + success=True, + response_status=200, + response_body="", + error=None, + duration_ms=10, + ) + + +def reclaiming_send(sessionmaker_factory, delivery, a_calls, b_calls): + + async def a_send( + client: httpx.AsyncClient, snapshot: DeliverySnapshot + ) -> DeliveryResult: + a_calls.append(snapshot) + async with sessionmaker_factory() as s: + await s.execute( + update(Delivery) + .where(Delivery.id == delivery.id) + .values(locked_until=datetime.now(UTC) - timedelta(seconds=1)) + ) + await s.commit() + await deliver( + ctx(client=client, sessionmaker=sessionmaker_factory), + delivery_id=str(delivery.id), + send_fn=send_fn([ok_result()], b_calls), + ) + + return fail_result() + + return a_send + + +def fake_rng(a: float, b: float) -> float: + return b - a diff --git a/backend/tests/test_claim.py b/backend/tests/test_claim.py index a3f4985..d95a300 100644 --- a/backend/tests/test_claim.py +++ b/backend/tests/test_claim.py @@ -1,6 +1,5 @@ import asyncio from datetime import UTC, datetime, timedelta -from typing import cast import httpx from sqlalchemy import select @@ -8,18 +7,8 @@ from app.config import settings from app.models import Delivery, DeliveryStatus -from app.tasks import DeliveryResult, DeliverySnapshot, WorkerContext, deliver -from tests.fakes import send_fn - - -def _ctx(*, client, sessionmaker) -> WorkerContext: - return cast(WorkerContext, {"client": client, "sessionmaker": sessionmaker}) - - -def _ok_result() -> DeliveryResult: - return DeliveryResult( - success=True, response_status=200, response_body="", error=None, duration_ms=10 - ) +from app.tasks import DeliveryResult, DeliverySnapshot, deliver +from tests.fakes import ctx, ok_result, send_fn async def test_claims_pending_row(make_delivery, sessionmaker_factory): @@ -29,12 +18,12 @@ async def test_claims_pending_row(make_delivery, sessionmaker_factory): next_attempt_at=datetime.now(UTC) - timedelta(seconds=1), ) - results: list[DeliveryResult] = [_ok_result()] + results: list[DeliveryResult] = [ok_result()] calls: list[DeliverySnapshot] = [] async with httpx.AsyncClient() as worker_client: await deliver( - _ctx(client=worker_client, sessionmaker=sessionmaker_factory), + ctx(client=worker_client, sessionmaker=sessionmaker_factory), delivery_id=str(delivery.id), send_fn=send_fn(results, calls), ) @@ -66,12 +55,12 @@ async def test_claims_nextattemptat_null_pending_row( next_attempt_at=None, ) - results: list[DeliveryResult] = [_ok_result()] + results: list[DeliveryResult] = [ok_result()] calls: list[DeliverySnapshot] = [] async with httpx.AsyncClient() as worker_client: await deliver( - _ctx(client=worker_client, sessionmaker=sessionmaker_factory), + ctx(client=worker_client, sessionmaker=sessionmaker_factory), delivery_id=str(delivery.id), send_fn=send_fn(results, calls), ) @@ -102,12 +91,12 @@ async def test_does_not_claim_future_scheduled_row(make_delivery, sessionmaker_f next_attempt_at=next_attempt, ) - results: list[DeliveryResult] = [_ok_result()] + results: list[DeliveryResult] = [ok_result()] calls: list[DeliverySnapshot] = [] async with httpx.AsyncClient() as worker_client: await deliver( - _ctx(client=worker_client, sessionmaker=sessionmaker_factory), + ctx(client=worker_client, sessionmaker=sessionmaker_factory), delivery_id=str(delivery.id), send_fn=send_fn(results, calls), ) @@ -138,12 +127,12 @@ async def test_does_not_claim_already_claimed(make_delivery, sessionmaker_factor locked_until=datetime.now(UTC) + timedelta(seconds=settings.lease), ) - results: list[DeliveryResult] = [_ok_result()] + results: list[DeliveryResult] = [ok_result()] calls: list[DeliverySnapshot] = [] async with httpx.AsyncClient() as worker_client: await deliver( - _ctx(client=worker_client, sessionmaker=sessionmaker_factory), + ctx(client=worker_client, sessionmaker=sessionmaker_factory), delivery_id=str(delivery.id), send_fn=send_fn(results, calls), ) @@ -173,12 +162,12 @@ async def test_claims_orphaned_row(make_delivery, sessionmaker_factory): updated_at=datetime.now(UTC) - timedelta(seconds=settings.lease * 5), ) - results: list[DeliveryResult] = [_ok_result()] + results: list[DeliveryResult] = [ok_result()] calls: list[DeliverySnapshot] = [] async with httpx.AsyncClient() as worker_client: await deliver( - _ctx(client=worker_client, sessionmaker=sessionmaker_factory), + ctx(client=worker_client, sessionmaker=sessionmaker_factory), delivery_id=str(delivery.id), send_fn=send_fn(results, calls), ) @@ -207,20 +196,20 @@ async def test_atomic_claim_under_concurrency(make_delivery, sessionmaker_factor next_attempt_at=datetime.now(UTC), ) - results: list[DeliveryResult] = [_ok_result(), _ok_result()] + results: list[DeliveryResult] = [ok_result(), ok_result()] calls: list[DeliverySnapshot] = [] async with httpx.AsyncClient() as worker_client: - ctx = _ctx(client=worker_client, sessionmaker=sessionmaker_factory) + c = ctx(client=worker_client, sessionmaker=sessionmaker_factory) await asyncio.gather( deliver( - ctx, + c, delivery_id=str(delivery.id), send_fn=send_fn(results, calls), ), deliver( - ctx, + c, delivery_id=str(delivery.id), send_fn=send_fn(results, calls), ), @@ -252,12 +241,12 @@ async def test_claims_failed_row(make_delivery, sessionmaker_factory): updated_at=datetime.now(UTC) - past_lease, ) - results: list[DeliveryResult] = [_ok_result()] + results: list[DeliveryResult] = [ok_result()] calls: list[DeliverySnapshot] = [] async with httpx.AsyncClient() as worker_client: await deliver( - _ctx(client=worker_client, sessionmaker=sessionmaker_factory), + ctx(client=worker_client, sessionmaker=sessionmaker_factory), delivery_id=str(delivery.id), send_fn=send_fn(results, calls), ) diff --git a/backend/tests/test_delivery_retry.py b/backend/tests/test_delivery_retry.py index e67ba6c..e222fb3 100644 --- a/backend/tests/test_delivery_retry.py +++ b/backend/tests/test_delivery_retry.py @@ -1,5 +1,4 @@ from datetime import UTC, datetime, timedelta -from typing import cast import httpx from pytest import approx @@ -11,75 +10,25 @@ from app.tasks import ( DeliveryResult, DeliverySnapshot, - WorkerContext, compute_backoff, deliver, sweep, ) -from tests.fakes import FakeQueue, FakeWorker, send_fn - - -def _ctx(*, queue=None, client=None, sessionmaker=None) -> WorkerContext: - worker_context = {} - if queue is not None: - worker_context["worker"] = FakeWorker(queue) - if client is not None: - worker_context["client"] = client - if sessionmaker is not None: - worker_context["sessionmaker"] = sessionmaker - return cast(WorkerContext, worker_context) - - -def _fail_result() -> DeliveryResult: - return DeliveryResult( - success=False, - response_status=500, - error=None, - duration_ms=10, - ) - - -def _ok_result() -> DeliveryResult: - return DeliveryResult( - success=True, - response_status=200, - error=None, - duration_ms=10, - ) - - -def _reclaiming_send(sessionmaker_factory, delivery, a_calls, b_calls): - - async def a_send( - client: httpx.AsyncClient, snapshot: DeliverySnapshot - ) -> DeliveryResult: - a_calls.append(snapshot) - async with sessionmaker_factory() as s: - await s.execute( - update(Delivery) - .where(Delivery.id == delivery.id) - .values(locked_until=datetime.now(UTC) - timedelta(seconds=1)) - ) - await s.commit() - await deliver( - _ctx(client=client, sessionmaker=sessionmaker_factory), - delivery_id=str(delivery.id), - send_fn=send_fn([_ok_result()], b_calls), - ) - return _fail_result() - - return a_send - - -def _rng(a: float, b: float) -> float: - return b - a +from tests.fakes import ( + FakeQueue, + ctx, + fail_result, + fake_rng, + reclaiming_send, + send_fn, +) def test_backoff_schedule(): pre_computed_backoff = [2, 4, 8, 16, 32, 60, 60] for i in range(7): - assert compute_backoff(i, 2, 2, 60, _rng) == pre_computed_backoff[i] + assert compute_backoff(i, 2, 2, 60, fake_rng) == pre_computed_backoff[i] async def test_cap_boundary_last_retry_stays_failed( @@ -92,12 +41,12 @@ async def test_cap_boundary_last_retry_stays_failed( next_attempt_at=datetime.now(UTC) - timedelta(seconds=1), ) - results: list[DeliveryResult] = [_fail_result()] + results: list[DeliveryResult] = [fail_result()] calls: list[DeliverySnapshot] = [] async with httpx.AsyncClient() as worker_client: await deliver( - _ctx(client=worker_client, sessionmaker=sessionmaker_factory), + ctx(client=worker_client, sessionmaker=sessionmaker_factory), delivery_id=str(delivery.id), send_fn=send_fn(results, calls), ) @@ -128,12 +77,12 @@ async def test_cap_boundary_flips_to_dead_letter( next_attempt_at=datetime.now(UTC) - timedelta(seconds=1), ) - results: list[DeliveryResult] = [_fail_result()] + results: list[DeliveryResult] = [fail_result()] calls: list[DeliverySnapshot] = [] async with httpx.AsyncClient() as worker_client: await deliver( - _ctx(client=worker_client, sessionmaker=sessionmaker_factory), + ctx(client=worker_client, sessionmaker=sessionmaker_factory), delivery_id=str(delivery.id), send_fn=send_fn(results, calls), ) @@ -162,9 +111,9 @@ async def test_terminal_is_inert(make_delivery, sessionmaker_factory): ) fake_queue = FakeQueue() - ctx = _ctx(queue=fake_queue, sessionmaker=sessionmaker_factory) + c = ctx(queue=fake_queue, sessionmaker=sessionmaker_factory) - await sweep(ctx) + await sweep(c) assert len(fake_queue.enqueued) == 0 @@ -194,9 +143,9 @@ async def test_fence_drops_stale_finalize(make_delivery, sessionmaker_factory): async with httpx.AsyncClient() as client: await deliver( - _ctx(client=client, sessionmaker=sessionmaker_factory), + ctx(client=client, sessionmaker=sessionmaker_factory), delivery_id=str(delivery.id), - send_fn=_reclaiming_send(sessionmaker_factory, delivery, a_calls, b_calls), + send_fn=reclaiming_send(sessionmaker_factory, delivery, a_calls, b_calls), ) async with sessionmaker_factory() as check: @@ -226,9 +175,9 @@ async def test_fence_orphaned_double_dispatch(make_delivery, sessionmaker_factor async with httpx.AsyncClient() as client: await deliver( - _ctx(client=client, sessionmaker=sessionmaker_factory), + ctx(client=client, sessionmaker=sessionmaker_factory), delivery_id=str(delivery.id), - send_fn=_reclaiming_send(sessionmaker_factory, delivery, a_calls, b_calls), + send_fn=reclaiming_send(sessionmaker_factory, delivery, a_calls, b_calls), ) async with sessionmaker_factory() as check: @@ -267,9 +216,9 @@ async def test_missing_destination_flips_to_dead_letter( async with httpx.AsyncClient() as client: await deliver( - _ctx(client=client, sessionmaker=sessionmaker_factory), + ctx(client=client, sessionmaker=sessionmaker_factory), delivery_id=str(delivery.id), - send_fn=send_fn([_fail_result()], calls), + send_fn=send_fn([fail_result()], calls), ) async with sessionmaker_factory() as check: @@ -312,7 +261,7 @@ async def test_destination_is_not_active( async with httpx.AsyncClient() as client: await deliver( - _ctx(client=client, sessionmaker=sessionmaker_factory), + ctx(client=client, sessionmaker=sessionmaker_factory), delivery_id=str(delivery.id), send_fn=send_fn([], calls), ) diff --git a/backend/tests/test_sweeper.py b/backend/tests/test_sweeper.py index fd24b54..b29490e 100644 --- a/backend/tests/test_sweeper.py +++ b/backend/tests/test_sweeper.py @@ -1,16 +1,9 @@ from datetime import UTC, datetime, timedelta -from typing import cast from app.config import settings from app.models import DeliveryStatus -from app.tasks import WorkerContext, sweep -from tests.fakes import FakeQueue, FakeRaisingQueue, FakeWorker - - -def _ctx(*, queue, sessionmaker) -> WorkerContext: - return cast( - WorkerContext, {"worker": FakeWorker(queue), "sessionmaker": sessionmaker} - ) +from app.tasks import sweep +from tests.fakes import FakeQueue, FakeRaisingQueue, ctx async def test_sweep_mix_of_rows(make_delivery, sessionmaker_factory): @@ -63,9 +56,9 @@ async def test_sweep_mix_of_rows(make_delivery, sessionmaker_factory): ) fake_queue = FakeQueue() - ctx = _ctx(queue=fake_queue, sessionmaker=sessionmaker_factory) + c = ctx(queue=fake_queue, sessionmaker=sessionmaker_factory) - await sweep(ctx) + await sweep(c) assert {e["delivery_id"] for e in fake_queue.enqueued} == { str(pending.id), @@ -84,9 +77,9 @@ async def test_bounded_batch(make_delivery, sessionmaker_factory): ) fake_queue = FakeQueue() - ctx = _ctx(queue=fake_queue, sessionmaker=sessionmaker_factory) + c = ctx(queue=fake_queue, sessionmaker=sessionmaker_factory) - await sweep(ctx) + await sweep(c) assert len(fake_queue.enqueued) == settings.redispatch_limit @@ -104,8 +97,8 @@ async def test_failure_is_skipped(make_delivery, sessionmaker_factory): ) fake_raising_queue = FakeRaisingQueue(fail_for={str(first.id)}) - ctx = _ctx(queue=fake_raising_queue, sessionmaker=sessionmaker_factory) - await sweep(ctx) + c = ctx(queue=fake_raising_queue, sessionmaker=sessionmaker_factory) + await sweep(c) assert len(fake_raising_queue.enqueued) == 1 assert fake_raising_queue.enqueued[0]["delivery_id"] == str(second.id) From 86d065a46610d4d4099f82ec7cb462c94e05ae5b Mon Sep 17 00:00:00 2001 From: Mikhail Kugan <68385536+mike-k-git@users.noreply.github.com> Date: Mon, 6 Jul 2026 01:20:24 +0200 Subject: [PATCH 5/9] tests: migrate to QueueDep --- backend/tests/test_retry_loop.py | 183 ++++++++++++++----------------- 1 file changed, 84 insertions(+), 99 deletions(-) diff --git a/backend/tests/test_retry_loop.py b/backend/tests/test_retry_loop.py index 0f0b133..ea9a844 100644 --- a/backend/tests/test_retry_loop.py +++ b/backend/tests/test_retry_loop.py @@ -1,21 +1,20 @@ from datetime import UTC, datetime, timedelta -from typing import cast import httpx from sqlalchemy import select, update from sqlalchemy.orm import selectinload from app.config import settings +from app.deps import get_queue from app.main import app from app.models import Delivery, DeliveryStatus from app.tasks import ( DeliveryResult, DeliverySnapshot, - WorkerContext, deliver, sweep, ) -from tests.fakes import FakeQueue, FakeWorker, send_fn +from tests.fakes import FakeQueue, ctx, send_fn async def test_retry_loop(client, source, sessionmaker_factory, signed): @@ -35,107 +34,93 @@ async def test_retry_loop(client, source, sessionmaker_factory, signed): "/routes", json={"source_id": f"{src_id}", "destination_id": f"{dst}"} ) - app.state.queue = FakeQueue() - - try: - raw, headers = signed({"type": "payment.succeeded"}, key="evt_1") - r = await client.post(f"/ingest/{source}", content=raw, headers=headers) - assert r.status_code == 202 - event_id = r.json()["event_id"] - r2 = await client.get(f"/events/{event_id}") - assert r2.status_code == 200 - event_details = r2.json() - - results: list[DeliveryResult] = [ - DeliveryResult( - success=False, - response_status=401, - response_body="", - error="error", - duration_ms=10, - ), - DeliveryResult( - success=True, - response_status=200, - response_body="", - error=None, - duration_ms=10, - ), - ] - calls: list[DeliverySnapshot] = [] - - async with httpx.AsyncClient() as worker_client: - await deliver( - cast( - WorkerContext, - {"client": worker_client, "sessionmaker": sessionmaker_factory}, - ), - delivery_id=str(event_details["deliveries"][0]["id"]), - send_fn=send_fn(results, calls), - ) + q = FakeQueue() + app.dependency_overrides[get_queue] = lambda: q + raw, headers = signed({"type": "payment.succeeded"}, key="evt_1") + r = await client.post(f"/ingest/{source}", content=raw, headers=headers) + assert r.status_code == 202 + event_id = r.json()["event_id"] + r2 = await client.get(f"/events/{event_id}") + assert r2.status_code == 200 + event_details = r2.json() + + results: list[DeliveryResult] = [ + DeliveryResult( + success=False, + response_status=401, + response_body="", + error="error", + duration_ms=10, + ), + DeliveryResult( + success=True, + response_status=200, + response_body="", + error=None, + duration_ms=10, + ), + ] + calls: list[DeliverySnapshot] = [] + + async with httpx.AsyncClient() as worker_client: + await deliver( + ctx(client=worker_client, sessionmaker=sessionmaker_factory), + delivery_id=str(event_details["deliveries"][0]["id"]), + send_fn=send_fn(results, calls), + ) - async with sessionmaker_factory() as check: - updated_delivery = ( - await check.execute( - select(Delivery) - .where(Delivery.id == event_details["deliveries"][0]["id"]) - .options(selectinload(Delivery.attempts)) - ) - ).scalar_one() - - assert updated_delivery.status == DeliveryStatus.failed - assert updated_delivery.attempt_count == 1 - assert len(updated_delivery.attempts) == 1 - assert len(calls) == 1 - assert len(results) == 1 - - async with sessionmaker_factory() as check: + async with sessionmaker_factory() as check: + updated_delivery = ( await check.execute( - update(Delivery) - .values( - next_attempt_at=datetime.now(UTC) - - timedelta(seconds=settings.lease * 2) - ) + select(Delivery) .where(Delivery.id == event_details["deliveries"][0]["id"]) + .options(selectinload(Delivery.attempts)) + ) + ).scalar_one() + + assert updated_delivery.status == DeliveryStatus.failed + assert updated_delivery.attempt_count == 1 + assert len(updated_delivery.attempts) == 1 + assert len(calls) == 1 + assert len(results) == 1 + + async with sessionmaker_factory() as check: + await check.execute( + update(Delivery) + .values( + next_attempt_at=datetime.now(UTC) + - timedelta(seconds=settings.lease * 2) ) - await check.commit() + .where(Delivery.id == event_details["deliveries"][0]["id"]) + ) + await check.commit() - ctx = cast( - WorkerContext, - { - "worker": FakeWorker(app.state.queue), - "sessionmaker": sessionmaker_factory, - }, + c = ctx(queue=q, sessionmaker=sessionmaker_factory) + + await sweep(c) + + assert event_details["deliveries"][0]["id"] in { + e["delivery_id"] for e in q.enqueued + } + + async with httpx.AsyncClient() as worker_client: + await deliver( + ctx(client=worker_client, sessionmaker=sessionmaker_factory), + delivery_id=str(event_details["deliveries"][0]["id"]), + send_fn=send_fn(results, calls), ) - await sweep(ctx) - - assert event_details["deliveries"][0]["id"] in { - e["delivery_id"] for e in app.state.queue.enqueued - } - - async with httpx.AsyncClient() as worker_client: - await deliver( - cast( - WorkerContext, - {"client": worker_client, "sessionmaker": sessionmaker_factory}, - ), - delivery_id=str(event_details["deliveries"][0]["id"]), - send_fn=send_fn(results, calls), + + async with sessionmaker_factory() as check: + updated_delivery = ( + await check.execute( + select(Delivery) + .where(Delivery.id == event_details["deliveries"][0]["id"]) + .options(selectinload(Delivery.attempts)) ) + ).scalar_one() - async with sessionmaker_factory() as check: - updated_delivery = ( - await check.execute( - select(Delivery) - .where(Delivery.id == event_details["deliveries"][0]["id"]) - .options(selectinload(Delivery.attempts)) - ) - ).scalar_one() - - assert updated_delivery.status == DeliveryStatus.succeeded - assert updated_delivery.attempt_count == 2 - assert len(updated_delivery.attempts) == 2 - assert len(calls) == 2 - assert len(results) == 0 - finally: - del app.state.queue + assert updated_delivery.status == DeliveryStatus.succeeded + assert updated_delivery.attempt_count == 2 + assert len(updated_delivery.attempts) == 2 + assert len(calls) == 2 + assert len(results) == 0 From 7dda2f64480f0e71145d08cf93363fc5332e5501 Mon Sep 17 00:00:00 2001 From: Mikhail Kugan <68385536+mike-k-git@users.noreply.github.com> Date: Mon, 6 Jul 2026 01:28:07 +0200 Subject: [PATCH 6/9] tests: harden retry loop test --- backend/tests/test_retry_loop.py | 3 +++ 1 file changed, 3 insertions(+) diff --git a/backend/tests/test_retry_loop.py b/backend/tests/test_retry_loop.py index ea9a844..1690b5e 100644 --- a/backend/tests/test_retry_loop.py +++ b/backend/tests/test_retry_loop.py @@ -97,8 +97,11 @@ async def test_retry_loop(client, source, sessionmaker_factory, signed): c = ctx(queue=q, sessionmaker=sessionmaker_factory) + before = len(q.enqueued) + await sweep(c) + assert len(q.enqueued) == before + 1 assert event_details["deliveries"][0]["id"] in { e["delivery_id"] for e in q.enqueued } From a8d2d6493bd268a5895ad4333f3775d1ba13a5d7 Mon Sep 17 00:00:00 2001 From: Mikhail Kugan <68385536+mike-k-git@users.noreply.github.com> Date: Mon, 6 Jul 2026 13:09:53 +0200 Subject: [PATCH 7/9] refactor(queue): remove key from enqueue --- backend/app/routers/deliveries.py | 6 +----- backend/app/tasks.py | 4 +--- 2 files changed, 2 insertions(+), 8 deletions(-) diff --git a/backend/app/routers/deliveries.py b/backend/app/routers/deliveries.py index a6ba8f5..42a8288 100644 --- a/backend/app/routers/deliveries.py +++ b/backend/app/routers/deliveries.py @@ -39,11 +39,7 @@ async def replay(delivery_id: uuid.UUID, session: SessionDep, queue: QueueDep): if owner is not None: await session.commit() try: - await queue.enqueue( - "deliver", - delivery_id=str(delivery_id), - key=f"deliver:{delivery_id!s}", - ) + await queue.enqueue("deliver", delivery_id=str(delivery_id)) except Exception: logger.warning( "replay enqueue failed for %s; sweeper will recover", delivery_id diff --git a/backend/app/tasks.py b/backend/app/tasks.py index 438f342..56d4c47 100644 --- a/backend/app/tasks.py +++ b/backend/app/tasks.py @@ -282,9 +282,7 @@ async def sweep(ctx: WorkerContext) -> None: queue = ctx["worker"].queue for did in redispatch: try: - await queue.enqueue( - "deliver", delivery_id=str(did), key=f"deliver:{did}" - ) + await queue.enqueue("deliver", delivery_id=str(did)) except Exception: logger.exception("failed to enqueue delivery %s", did) From 80462648dbd30dde670b03d9073b48df4a222e07 Mon Sep 17 00:00:00 2001 From: Mikhail Kugan <68385536+mike-k-git@users.noreply.github.com> Date: Tue, 7 Jul 2026 00:55:36 +0200 Subject: [PATCH 8/9] tests: replay tests --- backend/tests/test_replay.py | 193 +++++++++++++++++++++++++++++++++++ 1 file changed, 193 insertions(+) create mode 100644 backend/tests/test_replay.py diff --git a/backend/tests/test_replay.py b/backend/tests/test_replay.py new file mode 100644 index 0000000..4df3345 --- /dev/null +++ b/backend/tests/test_replay.py @@ -0,0 +1,193 @@ +import uuid + +import httpx +import pytest +from sqlalchemy import func, select +from sqlalchemy.orm import selectinload + +from app.deps import get_queue +from app.main import app, settings +from app.models import Delivery, DeliveryStatus, Event +from app.tasks import deliver +from tests.fakes import FakeRaisingQueue, ctx, fail_result, ok_result, send_fn + + +async def test_replay_accepted(client, make_delivery, sessionmaker_factory): + d = await make_delivery( + status=DeliveryStatus.dead_letter, next_attempt_at=None, attempt_count=1 + ) + + r = await client.post(f"/deliveries/{d.id}/replay") + assert r.status_code == 202 + + async with sessionmaker_factory() as s: + updated_d = ( + await s.execute(select(Delivery).where(Delivery.id == d.id)) + ).scalar_one() + + assert updated_d.status == DeliveryStatus.pending + assert updated_d.next_attempt_at is None + assert updated_d.locked_by is None + assert updated_d.locked_until is None + + +@pytest.mark.parametrize( + "status", + [ + DeliveryStatus.pending, + DeliveryStatus.delivering, + DeliveryStatus.failed, + DeliveryStatus.succeeded, + ], +) +async def test_replay_conflict_status( + status, client, make_delivery, sessionmaker_factory +): + d = await make_delivery(status=status, next_attempt_at=None, attempt_count=1) + r = await client.post(f"/deliveries/{d.id}/replay") + + assert r.status_code == 409 + + async with sessionmaker_factory() as s: + updated_d = ( + await s.execute(select(Delivery).where(Delivery.id == d.id)) + ).scalar_one() + + assert updated_d.status == status + assert updated_d.next_attempt_at is None + assert updated_d.attempt_count == 1 + assert updated_d.locked_by is None + assert updated_d.locked_until is None + + +async def test_replay_unknown_delivery(client): + r = await client.post(f"/deliveries/{uuid.uuid4()}/replay") + + assert r.status_code == 404 + + +async def test_oneshot_boundary( + client, make_delivery, monkeypatch, sessionmaker_factory +): + monkeypatch.setattr(settings, "max_attempts", 3) + d = await make_delivery( + status=DeliveryStatus.dead_letter, next_attempt_at=None, attempt_count=3 + ) + r = await client.post(f"/deliveries/{d.id}/replay") + assert r.status_code == 202 + + async with httpx.AsyncClient() as w: + await deliver( + ctx(client=w, sessionmaker=sessionmaker_factory), + delivery_id=str(d.id), + send_fn=send_fn([fail_result()], []), + ) + + async with sessionmaker_factory() as s: + updated_d = ( + await s.execute( + select(Delivery) + .where(Delivery.id == d.id) + .options(selectinload(Delivery.attempts)) + ) + ).scalar_one() + + assert updated_d.status == DeliveryStatus.dead_letter + assert updated_d.attempt_count == 4 + assert len(updated_d.attempts) == 2 + assert ( + max(updated_d.attempts, key=lambda a: a.attempt_number).attempt_number == 4 + ) + + +async def test_replay_to_success(client, make_delivery, sessionmaker_factory): + d = await make_delivery( + status=DeliveryStatus.dead_letter, next_attempt_at=None, attempt_count=1 + ) + r = await client.post(f"/deliveries/{d.id}/replay") + assert r.status_code == 202 + + async with httpx.AsyncClient() as w: + await deliver( + ctx(client=w, sessionmaker=sessionmaker_factory), + delivery_id=str(d.id), + send_fn=send_fn([ok_result()], []), + ) + + async with sessionmaker_factory() as s: + updated_d = ( + await s.execute( + select(Delivery) + .where(Delivery.id == d.id) + .options(selectinload(Delivery.attempts)) + ) + ).scalar_one() + + assert updated_d.status == DeliveryStatus.succeeded + assert updated_d.attempt_count == 2 + assert ( + max(updated_d.attempts, key=lambda a: a.attempt_number).attempt_number == 2 + ) + + +async def test_routes_snapshot(client, make_delivery, sessionmaker_factory): + d = await make_delivery( + status=DeliveryStatus.dead_letter, next_attempt_at=None, attempt_count=1 + ) + + async with sessionmaker_factory() as s: + event = ( + await s.execute(select(Event).where(Event.id == d.event_id)) + ).scalar_one() + + dst = ( + await client.post("/destinations", json={"name": "d1", "url": "http://a.test"}) + ).json()["id"] + + await client.post( + "/routes", json={"source_id": str(event.source_id), "destination_id": dst} + ) + + r = await client.post(f"/deliveries/{d.id}/replay") + assert r.status_code == 202 + + async with sessionmaker_factory() as s: + updated_d = ( + await s.execute( + select(Delivery) + .where(Delivery.id == d.id) + .options(selectinload(Delivery.attempts)) + ) + ).scalar_one() + + n = await s.scalar( + select(func.count()) + .select_from(Delivery) + .where(Delivery.event_id == d.event_id) + ) + assert n == 1 + assert updated_d.status == DeliveryStatus.pending + assert updated_d.attempt_count == 1 + assert len(updated_d.attempts) == 1 + assert updated_d.attempts[0].delivery_id == d.id + + +async def test_enqueu_raises(client, make_delivery, sessionmaker_factory): + app.dependency_overrides[get_queue] = lambda: FakeRaisingQueue(set(), fail_all=True) + d = await make_delivery( + status=DeliveryStatus.dead_letter, next_attempt_at=None, attempt_count=1 + ) + r = await client.post(f"/deliveries/{d.id}/replay") + assert r.status_code == 202 + + async with sessionmaker_factory() as s: + updated_d = ( + await s.execute( + select(Delivery) + .where(Delivery.id == d.id) + .options(selectinload(Delivery.attempts)) + ) + ).scalar_one() + + assert updated_d.status == DeliveryStatus.pending + assert updated_d.attempt_count == 1 From e70044ea3d142e6dc7932ffaa354f5df7c89369c Mon Sep 17 00:00:00 2001 From: Mikhail Kugan <68385536+mike-k-git@users.noreply.github.com> Date: Tue, 7 Jul 2026 01:02:20 +0200 Subject: [PATCH 9/9] docs: update README --- README.md | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/README.md b/README.md index fec2c5a..a4a96d1 100644 --- a/README.md +++ b/README.md @@ -14,7 +14,7 @@ A handful of decisions shape the whole thing. They're all on purpose. **At-least-once, and honest about it.** Duplicates are caught at ingest with a unique constraint on `(source, idempotency_key)`. Each delivery is claimed atomically before any work starts, so two workers can't run the same one. Claims expire on a lease and carry a fencing token, so a worker that stalls, gets replaced, and wakes up later writes nothing. Delivery can still duplicate when things go wrong. That's expected, and destinations should handle it. The system doesn't claim exactly-once, because it can't. -**Retries back off, and they end.** A failed delivery is rescheduled with exponential backoff and jitter, so a struggling destination gets breathing room instead of a stampede. Attempts are capped. Whatever runs out goes to dead-letter and waits to be replayed, instead of retrying forever. And a destination you've switched off doesn't burn attempts at all. Its deliveries are held, then flow again the moment it's back on. +**Retries back off, and they end.** A failed delivery is rescheduled with exponential backoff and jitter, so a struggling destination gets breathing room instead of a stampede. Attempts are capped. Whatever runs out goes to dead-letter and waits to be replayed, instead of retrying forever. Replay buys exactly one new attempt: succeed and it's delivered, fail and it's back in the inbox right away — not off in the background running another backoff ladder. The operator stays in the loop. And a destination you've switched off doesn't burn attempts at all. Its deliveries are held, then flow again the moment it's back on. **Destinations are locked in at ingest.** When an event arrives, its list of destinations is frozen. Replay re-runs that same list. The delivery history stays an honest record of what happened. @@ -31,14 +31,14 @@ cp backend/.env.example backend/.env # set POSTGRES_* and the DSNs docker compose up --build ``` -That starts Postgres, Redis, the API on `:8000`, and the worker. From there you can configure sources, destinations, and routes. Send webhooks to `POST /ingest/{source}`. Read back the event feed, the full detail for any event including its delivery attempts, and the dead-letter inbox. +That starts Postgres, Redis, the API on `:8000`, and the worker. From there you can configure sources, destinations, and routes. Send webhooks to `POST /ingest/{source}`. Read back the event feed, the full detail for any event including its delivery attempts, and the dead-letter inbox — and send anything in it back through the pipeline with `POST /deliveries/{id}/replay`. ## Status -Still being built, but the whole backend delivery story now runs end to end. An event comes in, gets verified and stored, fans out to every matched destination, and gets delivered with real retry semantics: exponential backoff with jitter, a cap on attempts, and dead-letter at the end of the line. The claim and finalize path is fenced, so even a worker that loses its lease mid-flight can't corrupt the record. A sweeper recovers anything a lost enqueue or a crash leaves stranded. Every tunable, from the lease to the backoff curve to the attempt cap, is a validated setting instead of a constant buried in the worker. The state machine is tested end to end. +Still being built, but the backend is feature-complete: ingest with signature checks and dedupe, routing with fan-out, delivery with backoff, a cap, and dead-lettering — and now replay. `POST /deliveries/{id}/replay` sends a dead-lettered delivery back through the exact same path: same claim, same worker, same ledger. Nothing about replay is a special case, which is the point. The state machine and the replay contract are tested end to end. -Next up is replay: one click to send anything in dead-letter back through the same path. +Next up is the React dashboard. ## Planned -Failed deliveries become replayable in one click, reusing the same delivery path. A React dashboard will sit on top of the read API for inspecting payloads and replaying failures. After that, a one-command deploy to Fly.io or Railway. Further out, the hub will reshape payloads per route and sign its own outbound requests. +A React dashboard will sit on top of the read API for inspecting payloads and replaying failures in one click. After that, a one-command deploy to Fly.io or Railway. Further out, the hub will reshape payloads per route and sign its own outbound requests.