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 architecture/test-broker.md
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@ User-facing: `docs/usage/testing.md`. Invariant summary: `CLAUDE.md` § Test bro

### Sync (default, `run_loops=False`)

`broker.publish` synchronously routes through `OutboxSubscriber.dispatch_one` — matches the FastStream test-broker idiom (`TestKafkaBroker` / `TestRabbitBroker`). The handler runs before `publish` returns; no background loops. `broker.publish_batch`, `cancel_timer`, and `fetch_unprocessed` are also patched to operate on the fake client (the `session` argument is **ignored** — `del session` — which diverges from production's `isinstance(session, AsyncSession)` `TypeError`; tests needing the session contract must use a real `OutboxClient`). `OutboxResponse` is *not* faked, so its eager session/queue/activate validation still fires under the test broker.
`broker.publish` synchronously routes through `OutboxSubscriber.dispatch_one` — matches the FastStream test-broker idiom (`TestKafkaBroker` / `TestRabbitBroker`). The handler runs before `publish` returns; no background loops. `broker.publish_batch` is also patched to operate on the fake client. `cancel_timer` and `fetch_unprocessed` are **not** patched — they are methods on `AbstractOutboxClient`, so `broker.<method>` delegates through the `FakeOutboxClient` the test broker swaps in as `broker.config.broker_config.client`; the fake requires the `session` argument (matching the client signature) but **ignores** it — `del session` — which diverges from production's `isinstance(session, AsyncSession)` `TypeError`; tests needing the session contract must use a real `OutboxClient`). `OutboxResponse` is *not* faked, so its eager session/queue/activate validation still fires under the test broker.

`_sync_dispatch` claims the just-fed row via the shared `_claim_fake_row` (the same lease + `deliveries_count++` mechanics `FakeOutboxClient.fetch` uses), so the `max_deliveries` boundary runs on one path; only the eligibility gate (`next_attempt_at <= now`) lives in `fetch`, which is why sync mode fires future-dated rows immediately.

Expand All @@ -32,7 +32,7 @@ Spins up the real `_fetch_loop` / `_worker_loop` against the fake client. Requir

`FakeOutboxClient` re-implements the outbox rules in Python because there is no in-process Postgres — eligibility, lease cutoff, retry timing, the NULL-token guard, and the DLQ projection all exist twice (SQL in `OutboxClient`, Python here). The two can't share an implementation (one runs in the database, one in the process), so `tests/test_client_contract.py` couples them by **behaviour** instead: one parametrized scenario module asserts the shared `AbstractOutboxClient` surface (`fetch` / `delete_with_lease` / `mark_pending_with_lease` + DLQ) against *both* adapters — the fake everywhere, real Postgres auto-skipped when unreachable. A per-adapter harness hides substrate differences (how a row is seeded, which connection a terminal write needs); scheduling is seeded as server-side `make_interval` offsets so the comparison is clock-skew-free; expectations are hand-specified so neither adapter passes trivially against itself.

What it pins is *structural* drift (eligibility states, FIFO selection under contention, the token guard, the DLQ projection). What it deliberately cannot pin: cross-host DB-vs-worker clock skew — an in-process test can't manufacture it — so the real client's server-side clock authority stays a documented invariant. `cancel_timer` and `timer_id` insert-dedup are broker/producer concerns (not on the client interface), so they remain covered in `test_integration.py` / `test_fake.py`, not the contract suite. The pure helpers that *can* be shared — the DLQ projection (`_DLQ_PROJECTION` in `schema.py`) and the activate-args resolution (`_scheduling.py`) — are extracted so the fake consumes the same code as production rather than a parallel copy.
What it pins is *structural* drift (eligibility states, FIFO selection under contention, the token guard, the DLQ projection). What it deliberately cannot pin: cross-host DB-vs-worker clock skew — an in-process test can't manufacture it — so the real client's server-side clock authority stays a documented invariant. `timer_id` insert-dedup is a broker/producer concern (not on the client interface). `cancel_timer` and `fetch_unprocessed` are now on the client interface, but their cross-adapter parity is covered by `test_integration.py` / `test_fake.py` rather than the parametrized contract suite (which pins the fetch/terminal/DLQ surface); folding them into the contract suite is a possible future tightening. The pure helpers that *can* be shared — the DLQ projection (`_DLQ_PROJECTION` in `schema.py`) and the activate-args resolution (`_scheduling.py`) — are extracted so the fake consumes the same code as production rather than a parallel copy.

## `_fake_start` skips the parent publisher-iteration loop

Expand Down
5 changes: 3 additions & 2 deletions docs/usage/testing.md
Original file line number Diff line number Diff line change
Expand Up @@ -144,8 +144,9 @@ return`).
immediately. The intended firing time is preserved on the harness's
`fake_client.rows[i].next_attempt_at` for assertions. Use
`run_loops=True` if you need scheduled delivery to actually wait.
- **`cancel_timer` and `fetch_unprocessed` are patched** to operate on the
fake client. The `session` argument is ignored in tests.
- **`cancel_timer` and `fetch_unprocessed` run against the fake client**
(the test broker swaps its client in, so `broker.<method>` reaches it).
The `session` argument is still required but is ignored in tests.
- **The fake producer uses the same envelope format as the real one**, so
all serialization paths are exercised.
- **`lease_ttl_seconds` and re-delivery are not simulated** in sync mode —
Expand Down
31 changes: 3 additions & 28 deletions faststream_outbox/broker.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,11 +26,10 @@
from faststream.exceptions import IncorrectState
from faststream.specification.schema import BrokerSpec
from faststream.specification.schema.extra import Tag, TagDict
from sqlalchemy import delete, select
from sqlalchemy.ext.asyncio import AsyncSession
from typing_extensions import override

from faststream_outbox.client import AbstractOutboxClient, OutboxClient, _row_to_message
from faststream_outbox.client import AbstractOutboxClient, OutboxClient
from faststream_outbox.configs import LastExceptionRenderer, OutboxBrokerConfig
from faststream_outbox.message import OutboxInnerMessage
from faststream_outbox.metrics import MetricsRecorder, _noop_recorder
Expand Down Expand Up @@ -484,17 +483,7 @@ async def cancel_timer(
``True`` only becomes durable when your transaction commits. If you branch on it and
then roll back, the cancellation rolls back with you.
"""
if not isinstance(session, AsyncSession):
msg = "broker.cancel_timer requires an sqlalchemy.ext.asyncio.AsyncSession"
raise TypeError(msg)
t = self._outbox_table
stmt = delete(t).where(
t.c.queue == queue,
t.c.timer_id == timer_id,
t.c.acquired_token.is_(None),
)
result = await session.execute(stmt)
return (result.rowcount or 0) > 0 # ty: ignore[unresolved-attribute]
return await self.client.cancel_timer(queue=queue, timer_id=timer_id, session=session)

async def fetch_unprocessed(
self,
Expand All @@ -514,21 +503,7 @@ async def fetch_unprocessed(
contract as :meth:`publish`); does not acquire a lease and does not mutate
row state, so it is safe to call alongside running subscribers.
"""
if not isinstance(session, AsyncSession):
msg = "broker.fetch_unprocessed requires an sqlalchemy.ext.asyncio.AsyncSession"
raise TypeError(msg)
if limit < 1:
# F4-04: a non-positive limit otherwise hits SQL (LIMIT -1 → DB error) or
# silently returns nothing (LIMIT 0); reject it up front, consistently with
# the fake.
msg = f"limit must be >= 1, got {limit}"
raise ValueError(msg)
t = self._outbox_table
stmt = select(*t.c).order_by(t.c.id).limit(limit)
if queue is not None:
stmt = stmt.where(t.c.queue == queue)
result = await session.execute(stmt)
return [_row_to_message(dict(row)) for row in result.mappings().all()]
return await self.client.fetch_unprocessed(session=session, queue=queue, limit=limit)

async def request(
self,
Expand Down
51 changes: 51 additions & 0 deletions faststream_outbox/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@
tuple_,
update,
)
from sqlalchemy.ext.asyncio import AsyncSession

from faststream_outbox import schema_validation
from faststream_outbox.message import OutboxInnerMessage
Expand Down Expand Up @@ -120,6 +121,18 @@ async def mark_pending_with_lease(
last_attempt_at: _dt.datetime,
) -> bool: ...

@abc.abstractmethod
async def cancel_timer(self, *, queue: str, timer_id: str, session: "AsyncSession") -> bool: ...

@abc.abstractmethod
async def fetch_unprocessed(
self,
*,
session: "AsyncSession",
queue: str | None = None,
limit: int = 1000,
) -> list[OutboxInnerMessage]: ...

@abc.abstractmethod
async def validate_schema(self, *, check_autovacuum: bool = False) -> None: ...

Expand Down Expand Up @@ -419,6 +432,44 @@ async def mark_pending_with_lease(
result = await conn.execute(stmt, {"delay": max(0.0, delay_seconds)})
return (result.rowcount or 0) > 0

async def cancel_timer(self, *, queue: str, timer_id: str, session: "AsyncSession") -> bool:
"""Delete a not-yet-leased ``(queue, timer_id)`` row on the caller's session."""
if not isinstance(session, AsyncSession):
msg = "OutboxClient.cancel_timer requires an sqlalchemy.ext.asyncio.AsyncSession"
raise TypeError(msg)
t = self._table
stmt = delete(t).where(
t.c.queue == queue,
t.c.timer_id == timer_id,
t.c.acquired_token.is_(None),
)
result = await session.execute(stmt)
return (result.rowcount or 0) > 0 # ty: ignore[unresolved-attribute]

async def fetch_unprocessed(
self,
*,
session: "AsyncSession",
queue: str | None = None,
limit: int = 1000,
) -> list[OutboxInnerMessage]:
"""Read up to *limit* rows (optionally filtered by *queue*) on the caller's session."""
if not isinstance(session, AsyncSession):
msg = "OutboxClient.fetch_unprocessed requires an sqlalchemy.ext.asyncio.AsyncSession"
raise TypeError(msg)
if limit < 1:
# F4-04: a non-positive limit otherwise hits SQL (LIMIT -1 → DB error) or
# silently returns nothing (LIMIT 0); reject it up front, consistently with
# the fake.
msg = f"limit must be >= 1, got {limit}"
raise ValueError(msg)
t = self._table
stmt = select(*t.c).order_by(t.c.id).limit(limit)
if queue is not None:
stmt = stmt.where(t.c.queue == queue)
result = await session.execute(stmt)
return [_row_to_message(dict(row)) for row in result.mappings().all()]

async def validate_schema(self, *, check_autovacuum: bool = False) -> None:
"""Validate that the database table(s) match the package's expected columns.

Expand Down
72 changes: 26 additions & 46 deletions faststream_outbox/testing.py
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,8 @@
if typing.TYPE_CHECKING:
from collections.abc import Iterator, Sequence

from sqlalchemy.ext.asyncio import AsyncSession

from faststream_outbox.subscriber.usecase import OutboxSubscriber


Expand Down Expand Up @@ -227,14 +229,37 @@ async def mark_pending_with_lease(
return True
return False

async def cancel_timer(self, *, queue: str, timer_id: str) -> bool:
async def cancel_timer(self, *, queue: str, timer_id: str, session: "AsyncSession") -> bool:
"""Mirror :meth:`OutboxBroker.cancel_timer` — drop a not-yet-leased timer row."""
del session # ignored (no real DB), matching the fake publish contract
for i, row in enumerate(self._rows):
if row.queue == queue and row.timer_id == timer_id and row.acquired_token is None:
del self._rows[i]
return True
return False

async def fetch_unprocessed(
self,
*,
session: "AsyncSession",
queue: str | None = None,
limit: int = 1000,
) -> list[OutboxInnerMessage]:
"""Mirror :meth:`OutboxBroker.fetch_unprocessed` — read rows from the in-memory store."""
del session # ignored (no real DB), matching the fake publish contract
if limit < 1:
# Mirror the real client's validation so a non-positive limit can't
# silently mis-slice (rows[:-1]) under the test broker (F4-04).
msg = f"limit must be >= 1, got {limit}"
raise ValueError(msg)
rows = sorted(self.rows, key=lambda r: r.id)
if queue is not None:
rows = [r for r in rows if r.queue == queue]
# Mirror production's ``limit`` (default 1000) so a valid
# ``broker.fetch_unprocessed(..., limit=N)`` call doesn't TypeError under
# the test broker (B16).
return [_to_inner(r) for r in rows[:limit]]

async def validate_schema(self, *, check_autovacuum: bool = False) -> None:
# Silently passing here would give tests false confidence — a user calling
# ``broker.validate_schema()`` against ``TestOutboxBroker`` would see a green
Expand Down Expand Up @@ -532,47 +557,6 @@ async def fake_publish_batch(
return fake_publish_batch


def _build_fake_cancel_timer(
fake_client: FakeOutboxClient,
) -> typing.Callable[..., typing.Awaitable[bool]]:
async def fake_cancel_timer(
*,
queue: str,
timer_id: str,
session: typing.Any = None,
) -> bool:
del session
return await fake_client.cancel_timer(queue=queue, timer_id=timer_id)

return fake_cancel_timer


def _build_fake_fetch_unprocessed(
fake_client: FakeOutboxClient,
) -> typing.Callable[..., typing.Awaitable[list[OutboxInnerMessage]]]:
async def fake_fetch_unprocessed(
*,
session: typing.Any = None,
queue: str | None = None,
limit: int = 1000,
) -> list[OutboxInnerMessage]:
del session
if limit < 1:
# Mirror the real broker's validation so a non-positive limit can't
# silently mis-slice (rows[:-1]) under the test broker (F4-04).
msg = f"limit must be >= 1, got {limit}"
raise ValueError(msg)
rows = sorted(fake_client.rows, key=lambda r: r.id)
if queue is not None:
rows = [r for r in rows if r.queue == queue]
# Mirror production's ``limit`` (default 1000) so a valid
# ``broker.fetch_unprocessed(..., limit=N)`` call doesn't TypeError under
# the test broker (B16).
return [_to_inner(r) for r in rows[:limit]]

return fake_fetch_unprocessed


class TestOutboxBroker(TestBroker[OutboxBroker, OutboxBroker]): # ty: ignore[invalid-type-arguments]
"""Test harness for ``OutboxBroker``. Two dispatch modes.

Expand Down Expand Up @@ -657,14 +641,10 @@ def _patch_broker(self, broker: OutboxBroker) -> "Iterator[None]":
serializer = broker.config.broker_config.fd_config._serializer # noqa: SLF001
fake_publish = _build_fake_publish(self.fake_client, broker, serializer, run_loops=self.run_loops)
fake_publish_batch = _build_fake_publish_batch(self.fake_client, broker, serializer, run_loops=self.run_loops)
fake_cancel_timer = _build_fake_cancel_timer(self.fake_client)
fake_fetch_unprocessed = _build_fake_fetch_unprocessed(self.fake_client)
try:
with (
mock.patch.object(broker, "publish", new=fake_publish),
mock.patch.object(broker, "publish_batch", new=fake_publish_batch),
mock.patch.object(broker, "cancel_timer", new=fake_cancel_timer),
mock.patch.object(broker, "fetch_unprocessed", new=fake_fetch_unprocessed),
super()._patch_broker(broker),
):
yield
Expand Down
69 changes: 69 additions & 0 deletions planning/changes/2026-07-17.01-table-ops-behind-client-seam.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,69 @@
---
summary: Moved cancel_timer/fetch_unprocessed off the broker's inline SQL onto AbstractOutboxClient; broker methods now delegate. Deletes the test-broker patch machinery for both; tightens the test broker's session arg from runtime-optional to required-but-ignored.
---

# Change: Route cancel_timer/fetch_unprocessed through the client seam

**Lane:** lightweight. A mostly behavior-preserving refactor that reshapes an
internal interface (`AbstractOutboxClient`) and carries one small, deliberate
behavior change on the test broker (session arg) — recorded below so a future
reader knows it was intended, not incidental.

## Goal

Make the client the single owner of outbox-table SQL. `broker.cancel_timer` and
`broker.fetch_unprocessed` previously built `delete(t)` / `select(*t.c)` inline on
`self._outbox_table`, bypassing the `OutboxClient` seam that already owns the
table — so table/column knowledge was scattered across producer (insert+NOTIFY),
client (lease/DLQ), and broker-inline. This concentrates it behind one seam.

## Approach

- Added `cancel_timer(*, queue, timer_id, session)` and
`fetch_unprocessed(*, session, queue=None, limit=1000)` to
`AbstractOutboxClient`; implemented on `OutboxClient` (SQL + both input guards
moved verbatim from the broker) and on `FakeOutboxClient`.
- **Guards live in the client adapters** (design decision): the broker delegation
is a pure passthrough. Matches how `publish` validates (client/command layer,
not broker) and keeps the whole operation behind the seam, so a directly-
constructed `OutboxClient` enforces the session/limit contract itself. The real
client raises `TypeError` (non-`AsyncSession`) / `ValueError` (`limit < 1`); the
fake ignores `session` (`del session`) and mirrors the limit `ValueError`.
- `fetch_unprocessed` co-locates with `_row_to_message` (already in `client.py`),
so `broker.py` stops importing `_row_to_message`, `delete`, and `select`.
- **Deleted test-broker patch machinery.** The test broker already swaps
`broker.config.broker_config.client = fake_client`, and `broker.client` reads
that field — so the delegating `broker.<method>` reaches the fake with no patch.
Removed `_build_fake_cancel_timer` / `_build_fake_fetch_unprocessed` and their
two `mock.patch.object` lines in `_patch_broker` (publish/publish_batch stay
patched — they drive sync dispatch and can't just delegate).

## Behavior change (deliberate)

The test broker's `cancel_timer` / `fetch_unprocessed` previously accepted a
*runtime-optional* `session` (the deleted patch defaulted it to `None`). It is now
**required-but-ignored**, matching the `AbstractOutboxClient` signature. This
aligns runtime with what was always true at the type level (the old test calls
carried `# ty: ignore[missing-argument]`) and with the documented public signature
(`docs/usage/timers.md`: `cancel_timer(*, queue, timer_id, session)`). No
documented user pattern omitted `session`, so no supported usage breaks; internal
tests that omitted it now pass a sentinel session.

## Files

- `faststream_outbox/client.py` — two abstract methods + `OutboxClient` impls.
- `faststream_outbox/broker.py` — both methods now one-line delegations; unused
imports pruned.
- `faststream_outbox/testing.py` — `FakeOutboxClient` gains `fetch_unprocessed`,
`cancel_timer` signature aligned; two builders + two patch lines deleted.
- `tests/test_unit.py`, `tests/test_fake.py` — call sites give the broker a client
/ pass a sentinel session; SQL assertions unchanged.
- `architecture/test-broker.md`, `docs/usage/testing.md` — truth-home promotion:
the two methods now delegate through the swapped fake client, session required.

## Verification

- [x] `just test` — full suite green: 599 passed, Postgres 17, coverage 100.00%.
- [x] `just lint` — clean (ruff + ty).
- [x] `just check-planning` — passes.
- [x] `mkdocs build --strict` — clean.
Loading