From 2e07237aacdf7d226489b8832cab13fa82269767 Mon Sep 17 00:00:00 2001 From: Artur Shiriev Date: Fri, 17 Jul 2026 10:42:16 +0300 Subject: [PATCH] refactor(client): route cancel_timer/fetch_unprocessed through the client seam Both broker methods built delete()/select() inline on self._outbox_table, bypassing the OutboxClient seam that already owns the table -- so table/column knowledge was split across producer, client, and broker-inline. Move them onto AbstractOutboxClient (real + fake) so the client is the single owner of that SQL; the broker methods become thin delegations. - guards live in the client adapters (matches how publish validates): the real client raises TypeError/ValueError, the fake ignores session and mirrors the limit guard; the broker delegation is a pure passthrough - fetch_unprocessed co-locates with _row_to_message, so broker.py drops the delete/select/_row_to_message imports - delete the test-broker patch machinery for both: the test broker already swaps in the fake as broker.client, so delegation reaches it (publish/publish_batch stay patched -- they drive sync dispatch) Behavior change (deliberate): the test broker's session arg goes from runtime-optional to required-but-ignored, aligning runtime with the always- required type and the documented public signature. No documented usage omitted it. Truth-home promotion in architecture/test-broker.md + docs/usage/testing.md. Full suite green (599 passed, coverage 100%), just lint + ty + check-planning + mkdocs --strict clean. Co-Authored-By: Claude Opus 4.8 (1M context) --- architecture/test-broker.md | 4 +- docs/usage/testing.md | 5 +- faststream_outbox/broker.py | 31 +------- faststream_outbox/client.py | 51 +++++++++++++ faststream_outbox/testing.py | 72 +++++++------------ ...6-07-17.01-table-ops-behind-client-seam.md | 69 ++++++++++++++++++ tests/test_fake.py | 18 ++--- tests/test_unit.py | 18 ++--- 8 files changed, 172 insertions(+), 96 deletions(-) create mode 100644 planning/changes/2026-07-17.01-table-ops-behind-client-seam.md diff --git a/architecture/test-broker.md b/architecture/test-broker.md index 5ea524d..1270670 100644 --- a/architecture/test-broker.md +++ b/architecture/test-broker.md @@ -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.` 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. @@ -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 diff --git a/docs/usage/testing.md b/docs/usage/testing.md index 2859307..350a585 100644 --- a/docs/usage/testing.md +++ b/docs/usage/testing.md @@ -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.` 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 — diff --git a/faststream_outbox/broker.py b/faststream_outbox/broker.py index 5e08622..039803c 100644 --- a/faststream_outbox/broker.py +++ b/faststream_outbox/broker.py @@ -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 @@ -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, @@ -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, diff --git a/faststream_outbox/client.py b/faststream_outbox/client.py index 9a0c5aa..0c3341a 100644 --- a/faststream_outbox/client.py +++ b/faststream_outbox/client.py @@ -35,6 +35,7 @@ tuple_, update, ) +from sqlalchemy.ext.asyncio import AsyncSession from faststream_outbox import schema_validation from faststream_outbox.message import OutboxInnerMessage @@ -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: ... @@ -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. diff --git a/faststream_outbox/testing.py b/faststream_outbox/testing.py index fdaf361..e9e1037 100644 --- a/faststream_outbox/testing.py +++ b/faststream_outbox/testing.py @@ -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 @@ -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 @@ -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. @@ -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 diff --git a/planning/changes/2026-07-17.01-table-ops-behind-client-seam.md b/planning/changes/2026-07-17.01-table-ops-behind-client-seam.md new file mode 100644 index 0000000..7f0321b --- /dev/null +++ b/planning/changes/2026-07-17.01-table-ops-behind-client-seam.md @@ -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.` 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. diff --git a/tests/test_fake.py b/tests/test_fake.py index ba7f525..4d155f9 100644 --- a/tests/test_fake.py +++ b/tests/test_fake.py @@ -407,7 +407,7 @@ async def test_fake_broker_cancel_timer_removes_row() -> None: async with test_broker: await broker.publish("x", queue="timers", timer_id="email-1") # ty: ignore[missing-argument] assert len(test_broker.fake_client.rows) == 1 - cancelled = await broker.cancel_timer(queue="timers", timer_id="email-1") # ty: ignore[missing-argument] + cancelled = await broker.cancel_timer(queue="timers", timer_id="email-1", session=object()) # ty: ignore[invalid-argument-type] assert cancelled is True assert test_broker.fake_client.rows == [] @@ -421,10 +421,10 @@ async def test_fake_broker_fetch_unprocessed_reads_fake_client() -> None: await broker.publish("a", queue="q1") # ty: ignore[missing-argument] await broker.publish("b", queue="q2") # ty: ignore[missing-argument] - all_rows = await broker.fetch_unprocessed() # ty: ignore[missing-argument] + all_rows = await broker.fetch_unprocessed(session=object()) # ty: ignore[invalid-argument-type] assert [r.queue for r in all_rows] == ["q1", "q2"] - q1_only = await broker.fetch_unprocessed(queue="q1") # ty: ignore[missing-argument] + q1_only = await broker.fetch_unprocessed(queue="q1", session=object()) # ty: ignore[invalid-argument-type] assert [r.queue for r in q1_only] == ["q1"] @@ -435,9 +435,9 @@ async def test_fake_broker_fetch_unprocessed_respects_limit() -> None: async with test_broker: for i in range(5): await broker.publish(str(i), queue="q1") # ty: ignore[missing-argument] - limited = await broker.fetch_unprocessed(limit=2) # ty: ignore[missing-argument] + limited = await broker.fetch_unprocessed(limit=2, session=object()) # ty: ignore[invalid-argument-type] assert len(limited) == 2 - all_rows = await broker.fetch_unprocessed() # ty: ignore[missing-argument] + all_rows = await broker.fetch_unprocessed(session=object()) # ty: ignore[invalid-argument-type] assert len(all_rows) == 5 @@ -448,7 +448,7 @@ async def test_fake_broker_fetch_unprocessed_rejects_non_positive_limit() -> Non async with test_broker: for bad in (0, -1): with pytest.raises(ValueError, match="limit"): - await broker.fetch_unprocessed(limit=bad) # ty: ignore[missing-argument] + await broker.fetch_unprocessed(limit=bad, session=object()) # ty: ignore[invalid-argument-type] async def test_fake_broker_router_subscriber_receives_publish() -> None: @@ -1201,7 +1201,7 @@ async def test_fake_client_feed_timer_id_different_queues_allowed() -> None: async def test_fake_client_cancel_timer_removes_unleased_row() -> None: fake = FakeOutboxClient() fake.feed(queue="q", payload=b"x", timer_id="email-1") - assert await fake.cancel_timer(queue="q", timer_id="email-1") is True + assert await fake.cancel_timer(queue="q", timer_id="email-1", session=object()) is True # ty: ignore[invalid-argument-type] assert fake.rows == [] @@ -1228,7 +1228,7 @@ async def test_fake_headers_not_shared_by_reference() -> None: async def test_fake_client_cancel_timer_unknown_returns_false() -> None: fake = FakeOutboxClient() - assert await fake.cancel_timer(queue="q", timer_id="never-existed") is False + assert await fake.cancel_timer(queue="q", timer_id="never-existed", session=object()) is False # ty: ignore[invalid-argument-type] async def test_fake_client_cancel_timer_skips_leased_row() -> None: @@ -1236,7 +1236,7 @@ async def test_fake_client_cancel_timer_skips_leased_row() -> None: fake.feed(queue="q", payload=b"x", timer_id="email-1") fake.rows[0].acquired_token = uuid.uuid4() fake.rows[0].acquired_at = _dt.datetime.now(tz=_dt.UTC) - assert await fake.cancel_timer(queue="q", timer_id="email-1") is False + assert await fake.cancel_timer(queue="q", timer_id="email-1", session=object()) is False # ty: ignore[invalid-argument-type] assert len(fake.rows) == 1 diff --git a/tests/test_unit.py b/tests/test_unit.py index dc60da0..a783c1d 100644 --- a/tests/test_unit.py +++ b/tests/test_unit.py @@ -1073,13 +1073,13 @@ async def test_broker_publish_batch_does_not_accept_timer_id() -> None: async def test_broker_cancel_timer_rejects_non_async_session() -> None: - broker = _make_broker() + broker = _make_broker(engine=MagicMock()) with pytest.raises(TypeError, match="AsyncSession"): await broker.cancel_timer(queue="orders", timer_id="x", session=object()) # ty: ignore[invalid-argument-type] async def test_broker_cancel_timer_emits_delete_with_lease_guard() -> None: - broker = _make_broker() + broker = _make_broker(engine=MagicMock()) session = AsyncMock(spec=AsyncSession) session.execute.return_value.rowcount = 1 deleted = await broker.cancel_timer(queue="orders", timer_id="email-1", session=session) @@ -1094,7 +1094,7 @@ async def test_broker_cancel_timer_emits_delete_with_lease_guard() -> None: async def test_broker_fetch_unprocessed_rejects_non_async_session() -> None: - broker = _make_broker() + broker = _make_broker(engine=MagicMock()) with pytest.raises(TypeError, match="AsyncSession"): await broker.fetch_unprocessed(session=object()) # ty: ignore[invalid-argument-type] @@ -1109,7 +1109,7 @@ def _fetch_unprocessed_session_mock() -> AsyncMock: async def test_broker_fetch_unprocessed_builds_select_all_columns() -> None: - broker = _make_broker() + broker = _make_broker(engine=MagicMock()) session = _fetch_unprocessed_session_mock() rows = await broker.fetch_unprocessed(session=session) assert rows == [] @@ -1123,7 +1123,7 @@ async def test_broker_fetch_unprocessed_builds_select_all_columns() -> None: async def test_broker_fetch_unprocessed_filters_by_queue() -> None: - broker = _make_broker() + broker = _make_broker(engine=MagicMock()) session = _fetch_unprocessed_session_mock() await broker.fetch_unprocessed(session=session, queue="orders") stmt = session.execute.await_args_list[0].args[0] @@ -1135,7 +1135,7 @@ async def test_broker_fetch_unprocessed_filters_by_queue() -> None: async def test_broker_fetch_unprocessed_applies_default_limit() -> None: # Guardrail against accidental SELECT * with no LIMIT against a backlogged table. - broker = _make_broker() + broker = _make_broker(engine=MagicMock()) session = _fetch_unprocessed_session_mock() await broker.fetch_unprocessed(session=session) stmt = session.execute.await_args_list[0].args[0] @@ -1143,7 +1143,7 @@ async def test_broker_fetch_unprocessed_applies_default_limit() -> None: async def test_broker_fetch_unprocessed_respects_explicit_limit() -> None: - broker = _make_broker() + broker = _make_broker(engine=MagicMock()) session = _fetch_unprocessed_session_mock() await broker.fetch_unprocessed(session=session, limit=5) stmt = session.execute.await_args_list[0].args[0] @@ -1151,7 +1151,7 @@ async def test_broker_fetch_unprocessed_respects_explicit_limit() -> None: async def test_broker_cancel_timer_returns_false_when_nothing_deleted() -> None: - broker = _make_broker() + broker = _make_broker(engine=MagicMock()) session = AsyncMock(spec=AsyncSession) session.execute.return_value.rowcount = 0 deleted = await broker.cancel_timer(queue="orders", timer_id="x", session=session) @@ -4044,7 +4044,7 @@ async def test_broker_publish_batch_empty_rejects_empty_queue() -> None: async def test_fetch_unprocessed_rejects_non_positive_limit() -> None: """F4-04: a non-positive limit raises rather than hitting SQL (limit=-1) or silently returning none (limit=0).""" - broker = _make_broker() + broker = _make_broker(engine=MagicMock()) for bad in (0, -1): with pytest.raises(ValueError, match="limit"): await broker.fetch_unprocessed(session=_make_session_mock(), limit=bad)