Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -76,7 +76,7 @@ The fetch CTE's OR is written so each disjunct **explicitly carries its partial-

There is **no `state` column**: a row is "available" iff `acquired_token IS NULL` or `acquired_at < now() - lease_ttl_seconds`. Terminal failures `DELETE` by default; opt in to audit via `dlq_table=make_dlq_table(metadata)`.

`validate_schema()` is **opt-in** (call from `/health` or startup hook, not `broker.start()`) so migrations can run against the same DB without a loop. Beyond the alembic column/index diff it also probes the live partial-index **predicates** (alembic ignores `postgresql_where`), catching a drifted or non-partial `timer_id_uq` that would otherwise break `ON CONFLICT` at publish time (S2). Alembic is optional (`faststream-outbox[validate]`); without it `validate_schema()` raises `ImportError` but every other path works.
`validate_schema()` is **opt-in** (call from `/health` or startup hook, not `broker.start()`) so migrations can run against the same DB without a loop. Beyond the alembic column/index diff it also probes the live partial-index **predicates** (alembic ignores `postgresql_where`), catching a drifted or non-partial `timer_id_uq` that would otherwise break `ON CONFLICT` at publish time (S2), **and probes `pg_constraint` for the `<table>_lease_ck` CHECK** (alembic has no check-constraint comparator), catching a missing or drifted lease pairing. Alembic is optional (`faststream-outbox[validate]`); without it `validate_schema()` raises `ImportError` but every other path works.

### Opt-in DLQ on terminal failure

Expand Down Expand Up @@ -121,7 +121,7 @@ Lease-loss logs at WARNING with `extra={"event": "lease_lost", "phase": "termina

Both `OutboxSubscriber.stop()` and `OutboxBroker.stop()` override FastStream parents. Override comments carry `# Upstream equivalent (replaced): …`.

- **Subscriber: two flags during drain.** `self.running` (FastStream's "actively dispatching") stays True for the duration of drain; `self._stopping` (new) signals "no new claims". `_fetch_inner` checks both; the worker loop only `running`. `stop()` flips `_stopping`, kicks `_notify_event`, waits up to `graceful_timeout` for `_inflight.join()`, then flips `running=False` and cancels tasks. `super().stop()` is **not** called — its `MultiLock.wait_release` would re-wait stuck handlers for another full budget (2× shutdown regression).
- **Subscriber: two flags during drain.** `self.running` (FastStream's "actively dispatching") stays True for the duration of drain; `self._stopping` (new) signals "no new claims". `_fetch_inner` checks both; the worker loop only `running`. `stop()` flips `_stopping`, kicks `_notify_event`, waits up to `graceful_timeout` for `_inflight.join()`, then flips `running=False` and cancels tasks. `graceful_timeout=None` (unbounded for `ping()`) is **clamped to a finite fallback in the drain** so one wedged handler can't hang `stop()` forever. `super().stop()` is **not** called — its `MultiLock.wait_release` would re-wait stuck handlers for another full budget (2× shutdown regression).
- **Broker: parallel-gather subscriber stop** via `asyncio.gather(..., return_exceptions=True)` — sequential N × `graceful_timeout` exceeds K8s default `terminationGracePeriodSeconds=30s` once a service has 2+ subscribers. Exceptions logged via `_log_subscriber_stop_error`, never re-raised.
- **Phase interaction.** During drain `running` stays True so the `dispatch_one` guard is dormant; after drain `running=False` is set before `task.cancel()` so workers mid-`dispatch_one` benefit from the guard. The two changes are complementary.
- **Upstream divergence flag.** If FastStream adds cleanup to `BrokerUsecase.stop`, `SubscriberUsecase.stop`, or `TasksMixin.stop`, we silently miss it. **Re-check both overrides when touching shutdown.** Regression tests in `tests/test_fake.py` (`test_drain_finishes_inflight_rows_before_returning_in_fake_mode`, `test_broker_stop_cancels_wedged_handler_within_graceful_timeout_in_fake_mode`) and the Postgres-backed `tests/test_integration.py`.
Expand Down
1 change: 1 addition & 0 deletions architecture/drain.md
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ The subscriber carries two flags during shutdown: `self.running` (FastStream's e
- `_fetch_inner`'s loop guard checks both: `while self.running and not self._stopping:`
- The worker loop only checks `running`.
- `stop()` flips `_stopping`, kicks the fetch loop awake via `_notify_event` (in case it's parked in an idle `_wait_for_notify_or_timeout`), waits up to `graceful_timeout` for `_inflight.join()`, then flips `running=False` and cancels the spawned tasks.
- **`graceful_timeout=None`** stays unbounded where FastStream uses it that way (e.g. `ping()`), but the drain wait clamps `None` to a finite fallback (`_DEFAULT_DRAIN_TIMEOUT_SECONDS = 15.0`). `anyio.move_on_after(None)` has deadline `inf`, so without the clamp a single wedged handler would make `_inflight.join()` — and thus `stop()` — never return.

**Why we skip `super().stop()`.** Its `MultiLock.wait_release(graceful_timeout)` would either return instantly (healthy path; `_inflight.join()` already waited a stricter condition) or re-wait the same stuck handlers for another full budget (wedged path; **2× shutdown regression**). The subscriber inlines `TasksMixin.stop`'s cleanup body instead. Per-subscriber shutdown bound: `graceful_timeout`.

Expand Down
59 changes: 57 additions & 2 deletions faststream_outbox/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -397,6 +397,10 @@ async def validate_schema(self) -> None:
# and later breaks the producer's ON CONFLICT arbiter. Probe the predicates
# directly against the live catalog.
errors.extend(await conn.run_sync(_validate_index_predicates_sync, self._table))
# Alembic's compare_metadata has no check-constraint comparator, so a missing
# or altered <table>_lease_ck (the half-set-lease guard) passes the diff above
# silently. Probe pg_constraint directly, mirroring the partial-index probe.
errors.extend(await conn.run_sync(_validate_check_constraints_sync, self._table))
if self._dlq_table is not None:
errors.extend(await conn.run_sync(_validate_dlq_schema_sync, self._dlq_table))
if errors:
Expand Down Expand Up @@ -493,6 +497,57 @@ def _validate_index_predicates_sync(connection: "Connection", table: "Table") ->
return errors


_CHECK_CONSTRAINT_QUERY = text(
"SELECT con.conname AS name, pg_get_constraintdef(con.oid) AS definition "
"FROM pg_constraint con "
"JOIN pg_class t ON t.oid = con.conrelid "
"JOIN pg_namespace n ON n.oid = t.relnamespace "
"WHERE t.relname = :table AND n.nspname = COALESCE(:schema, current_schema()) "
"AND con.contype = 'c'",
)

# Expected CHECK-constraint predicates, keyed by the constraint-name suffix make_outbox_table
# uses. Normalized form (lowercased, parens/whitespace collapsed, leading ``check`` stripped)
# of ``(acquired_token IS NULL) = (acquired_at IS NULL)``.
_EXPECTED_CHECK_CONSTRAINTS = {
"_lease_ck": "acquired_token is null = acquired_at is null",
}


def _validate_check_constraints_sync(connection: "Connection", table: "Table") -> list[str]:
"""
Compare the live CHECK constraint(s) against what the package expects.

Alembic's ``compare_metadata`` registers no check-constraint comparator, so a missing
or altered ``<table>_lease_ck`` — the ``(acquired_token IS NULL) = (acquired_at IS NULL)``
invariant that makes a half-set lease unrepresentable — slips through :func:`_run_validate`
entirely. Probe ``pg_constraint`` directly, mirroring :func:`_validate_index_predicates_sync`.

Flags both a **missing** constraint and one whose normalized predicate **drifted**.
"""
rows = (
connection.execute(
_CHECK_CONSTRAINT_QUERY,
{"table": table.name, "schema": table.schema},
)
.mappings()
.all()
)
live = {row["name"]: row["definition"] for row in rows}
errors: list[str] = []
for suffix, want in _EXPECTED_CHECK_CONSTRAINTS.items():
name = f"{table.name}{suffix}"
if name not in live:
errors.append(f"missing CHECK constraint {name!r} (expected '{want}')")
continue
# pg_get_constraintdef returns e.g. ``CHECK (((a IS NULL) = (b IS NULL)))``; strip the
# leading ``check`` keyword after normalizing away parens/case/whitespace.
got = _normalize_predicate(live[name]).removeprefix("check ").strip()
if got != want:
errors.append(f"CHECK constraint {name!r} has wrong predicate: expected '{want}', got '{got}'")
return errors


def _validate_schema_sync(connection: "Connection", table: "Table") -> list[str]:
"""Run the outbox-table validation pass; see :func:`_run_validate` for the diff machinery."""
return _run_validate(connection, table, make_outbox_table)
Expand Down Expand Up @@ -598,8 +653,8 @@ def _drift_entry_to_error(entry: "tuple[typing.Any, ...]", table_name: str) -> s
a CHECK (``<table>_lease_ck``) and a partial unique index (``<table>_timer_id_uq``):
the unique index surfaces as ``add_index`` above, but Alembic's ``compare_metadata``
has no check-constraint comparator, so a missing or altered CHECK never appears in
this diff at all (a known ``validate_schema`` blind spot — detecting it needs a
separate ``pg_constraint`` catalog probe).
this diff at all. That gap is covered separately by
:func:`_validate_check_constraints_sync` (a direct ``pg_constraint`` probe).
"""
op = entry[0]
if op == "add_table":
Expand Down
17 changes: 14 additions & 3 deletions faststream_outbox/response.py
Original file line number Diff line number Diff line change
Expand Up @@ -122,9 +122,14 @@ class OutboxResponse(Response):

Idiomatic FastStream shape: ``async def h(...) -> OutboxResponse``. Requires
``session=...`` for the same reason ``broker.publish`` does — the new row must
commit with the caller's domain writes. Validation (session type, activate
args mutex, tz-aware datetime) is deferred to ``OutboxPublishCommand.__init__``
on ``as_publish_command()`` so there's a single source of truth.
commit with the caller's domain writes.

The activate-args mutex + tz-aware checks run **eagerly** in ``__init__`` so a
misconfigured response raises at the ``return OutboxResponse(...)`` site. Deferring
them to ``as_publish_command()`` (dispatch time) made the error masquerade as a
handler failure and exhaust the inbound row's retry budget (audit 2026-06-14).
``OutboxPublishCommand.__init__`` re-checks on ``as_publish_command()``, so it stays
the authoritative single source of truth (these eager checks are a fail-fast mirror).

``correlation_id`` defaults to the inbound message's correlation_id when not set
— FastStream's ``SubscriberUsecase.process_message`` does the inheritance before
Expand All @@ -143,6 +148,12 @@ def __init__(
activate_at: _dt.datetime | None = None,
timer_id: str | None = None,
) -> None:
if activate_in is not None and activate_at is not None:
msg = "OutboxResponse accepts at most one of activate_in / activate_at"
raise ValueError(msg)
if activate_at is not None and activate_at.tzinfo is None:
msg = "OutboxResponse requires activate_at to be timezone-aware"
raise ValueError(msg)
super().__init__(body=body, headers=headers, correlation_id=correlation_id)
self.queue = queue
self.session = session
Expand Down
13 changes: 12 additions & 1 deletion faststream_outbox/subscriber/usecase.py
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,11 @@
# failing is treated as recovered: the next failure starts a fresh backoff
# sequence instead of inheriting the lifetime error count (B3).
_BACKOFF_RESET_THRESHOLD_SECONDS = 60.0
# Fallback drain budget when graceful_timeout is None. None means "unbounded" for
# ping(), but an unbounded drain lets a single wedged handler hang stop() forever
# (anyio.move_on_after(None) has deadline=inf), so the drain path clamps to this.
# Mirrors OutboxBroker/OutboxRouter's graceful_timeout=15.0 default.
_DEFAULT_DRAIN_TIMEOUT_SECONDS = 15.0


# Marker exception raised by programming guards inside ``process_message`` (e.g.
Expand Down Expand Up @@ -241,7 +246,13 @@ async def stop(self) -> None:
# SubscriberUsecase.stop -> faststream/_internal/endpoint/subscriber/usecase.py
self._stopping = True
self._notify_event.set()
with anyio.move_on_after(self._outer_config.graceful_timeout):
# graceful_timeout=None stays "unbounded" for ping(), but the drain must be
# strict-bound or one wedged handler hangs stop() forever — clamp None to a
# finite fallback so move_on_after always has a real deadline (audit 2026-06-14).
drain_timeout = self._outer_config.graceful_timeout
if drain_timeout is None:
drain_timeout = _DEFAULT_DRAIN_TIMEOUT_SECONDS
with anyio.move_on_after(drain_timeout):
await self._inflight.join()
self.running = False
tasks = list(self.tasks)
Expand Down
6 changes: 5 additions & 1 deletion planning/audits/2026-06-14-deep-audit-findings.md
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,8 @@ None.

### [LOW][design] `graceful_timeout=None` makes drain wait forever on a wedged handler

> **RESOLVED (2026-06-14)** — the drain path clamps `None` to a finite fallback (`_DEFAULT_DRAIN_TIMEOUT_SECONDS = 15.0`, mirroring the broker default) so `stop()` is always strict-bound; `None` stays unbounded for `ping()`. Test: `tests/test_integration.py::test_graceful_timeout_none_still_bounds_drain`.

`faststream_outbox/subscriber/usecase.py:244-245`; accepted at `broker.py:141` and `fastapi/router.py:79`, forwarded unchanged.

**Problem.** `OutboxBroker`/`OutboxRouter` accept `graceful_timeout: float | None = 15.0` and forward `None` straight through. In subscriber `stop()`, `with anyio.move_on_after(self._outer_config.graceful_timeout): await self._inflight.join()` — `anyio.move_on_after(None)` produces an *unbounded* scope (deadline `inf`). With `graceful_timeout=None` and a wedged handler, `_inflight.join()` never returns, so `running` is never flipped to False, tasks are never cancelled, and `stop()` (hence `broker.stop()`) hangs. This silently defeats the documented "strict-bound drain" contract (architecture/drain.md:15-17).
Expand All @@ -69,7 +71,7 @@ None.

### [LOW][bug] `validate_schema()` cannot detect a missing/wrong `_lease_ck` CHECK constraint; docstring falsely claims the table declares no CHECK

> **PARTIALLY RESOLVED (2026-06-14)** — the false `_drift_entry_to_error` docstring is corrected. The actual `pg_constraint` probe (a new failure mode in `validate_schema`) is **deferred to the behavior-change PR**, since it can make a DB that passed before start failing.
> **RESOLVED (2026-06-14)** — docstring corrected (safe PR), and `validate_schema()` now runs `_validate_check_constraints_sync`, a `pg_constraint` probe that flags a missing or drifted `<table>_lease_ck`. Tests: `tests/test_integration.py::test_validate_schema_fails_when_lease_check_constraint_{missing,predicate_wrong}` (+ the existing passes-for-correct-table happy path). Note: this is a new failure mode — a DB lacking the CHECK that passed `validate_schema()` before will now fail it (intended).

`faststream_outbox/client.py:587-621` (docstring at 597-599); constraint declared at `schema.py:92-95`.

Expand All @@ -81,6 +83,8 @@ None.

### [LOW][design] Handler returning `OutboxResponse` with a naive/invalid `activate_at` retries the inbound message to exhaustion

> **RESOLVED (2026-06-14)** — `OutboxResponse.__init__` now runs the activate-mutex + tz-aware checks eagerly, so a misconfigured response raises at the `return OutboxResponse(...)` site instead of masquerading as a handler failure at dispatch. `OutboxPublishCommand` re-checks on `as_publish_command()` (still the authoritative source). Tests: `tests/test_unit.py::test_outbox_response_rejects_{naive_activate_at,both_activate_args}_eagerly`.

`faststream_outbox/subscriber/usecase.py:905-919` (publish loop), `response.py:134-161` (deferred validation).

**Problem.** `OutboxResponse.__init__` does zero validation; all checks (session type, activate-args mutex, tz-awareness) defer to `OutboxPublishCommand` via `as_publish_command()`, called at `usecase.py:915` *after* the handler returned. A handler returning `OutboxResponse(..., activate_at=<naive datetime>)` raises an ordinary `ValueError` there, which propagates through the middleware stack → `_CaptureExceptionMiddleware` stashes it onto `row.last_exception` → the inbound row is nacked under `NACK_ON_ERROR` (or via `assert_state_set` under MANUAL). Because the error is deterministic, every reclaim re-raises it: the inbound message walks its full retry budget (default 10 attempts), then DLQs as `retry_terminal`.
Expand Down
70 changes: 70 additions & 0 deletions tests/test_integration.py
Original file line number Diff line number Diff line change
Expand Up @@ -1777,3 +1777,73 @@ async def _deleted() -> bool:

assert calls["n"] >= 2, "the terminal write must be retried after the connection rebuild"
assert received # handler ran (at least once; redelivery may run it again)


# --- behavior-change Lows (audit 2026-06-14) ----------------------------------


async def test_graceful_timeout_none_still_bounds_drain(
pg_engine: AsyncEngine,
outbox_table: Table,
) -> None:
"""
``graceful_timeout=None`` must not hang ``stop()`` on a wedged handler.

None stays "unbounded" for ``ping()``, but the drain clamps it to a finite fallback so
a single stuck handler can't make ``stop()`` (hence pod shutdown) hang forever. Without
the clamp ``anyio.move_on_after(None)`` has deadline=inf and this test would hang. The
module fallback is patched down so the test stays fast.
"""
broker = OutboxBroker(pg_engine, outbox_table=outbox_table, graceful_timeout=None)
started = asyncio.Event()

@broker.subscriber("orders", min_fetch_interval=0.02, max_fetch_interval=0.05)
async def handle(body: dict[str, Any]) -> None:
del body
started.set()
await asyncio.sleep(60.0) # wedged — never returns voluntarily

session_factory = async_sessionmaker(pg_engine, expire_on_commit=False)
with mock.patch("faststream_outbox.subscriber.usecase._DEFAULT_DRAIN_TIMEOUT_SECONDS", 0.3):
async with broker:
async with session_factory() as session, session.begin():
await broker.publish({"i": 0}, queue="orders", session=session)
await asyncio.wait_for(started.wait(), timeout=3.0)
start = asyncio.get_event_loop().time()
await broker.stop()
elapsed = asyncio.get_event_loop().time() - start

assert elapsed < 0.7, f"broker.stop() took {elapsed:.3f}s — graceful_timeout=None did not clamp the drain"
# Row preserved with its lease set; another replica reclaims after lease_ttl.
assert await _row_count(pg_engine, outbox_table) == 1


async def test_validate_schema_fails_when_lease_check_constraint_missing(
pg_engine: AsyncEngine,
outbox_table: Table,
) -> None:
"""A dropped ``<table>_lease_ck`` CHECK must be caught — alembic's diff can't see it (audit 2026-06-14)."""
drop_sql = f'ALTER TABLE "{outbox_table.name}" DROP CONSTRAINT "{outbox_table.name}_lease_ck"'
async with pg_engine.begin() as conn:
await conn.exec_driver_sql(drop_sql)
client = OutboxClient(pg_engine, outbox_table)
with pytest.raises(RuntimeError, match="missing CHECK constraint"):
await client.validate_schema()


async def test_validate_schema_fails_when_lease_check_constraint_predicate_wrong(
pg_engine: AsyncEngine,
outbox_table: Table,
) -> None:
"""A ``<table>_lease_ck`` with a drifted predicate must be caught by the pg_constraint probe."""
name = outbox_table.name
async with pg_engine.begin() as conn:
await conn.exec_driver_sql(f'ALTER TABLE "{name}" DROP CONSTRAINT "{name}_lease_ck"')
# Re-add under the same name with a different (wrong) predicate.
await conn.exec_driver_sql(
f'ALTER TABLE "{name}" ADD CONSTRAINT "{name}_lease_ck" '
f"CHECK (acquired_token IS NOT NULL OR acquired_at IS NULL)",
)
client = OutboxClient(pg_engine, outbox_table)
with pytest.raises(RuntimeError, match="wrong predicate"):
await client.validate_schema()
Loading