From 6e561d7238aa64bcd7299a38db36a357654af120 Mon Sep 17 00:00:00 2001 From: Artur Shiriev Date: Sat, 3 Oct 2026 13:15:32 +0300 Subject: [PATCH 1/2] docs: fix broken examples and wrong behaviour claims - relay.md: standalone router example uses faststream.kafka.KafkaRouter and faststream_outbox.OutboxRouter; OutboxResponse example passes a session - messaging-service.md: reference database_engine inside the class body - fastapi.md: dispose the engine in a lifespan (no add_event_handler) - troubleshooting.md: config-error path logs and leaves the row to lease expiry (no nack); real ImportError text; concurrent drain; lease_lost message text alongside the event extra; drop stale history - instrumentation-seams.md: four bus-invisible events; drop duplicate line - setup-prometheus-opentelemetry.md: OTel meters go to the global OTel meter provider; duplicate-collector error fires on second collector; consume vs publish destination attribute - dlq.md, schema-validation.md, alembic.md: no enum column; document pg_catalog probes and check_autovacuum; partitioned DLQ last_exception is VARCHAR so validate_schema passes - index.md, how-it-works.md: standalone and relay as two options; index lists performance.md - first-outbox-app.md: Python 3.11+, current version in sample output - README: faststream_outbox.fastapi.OutboxRouter - Replace private internals in testing, observability, dlq, router, comparison, troubleshooting and subscriber pages - Quote the validate extra in the validate_schema ImportError hint --- README.md | 2 +- docs/concepts/comparison.md | 13 ++-- docs/concepts/instrumentation-seams.md | 5 +- docs/concepts/performance.md | 9 ++- docs/index.md | 13 +++- docs/introduction/how-it-works.md | 15 +++-- docs/operations/alembic.md | 2 +- docs/operations/troubleshooting.md | 70 +++++++++++--------- docs/tutorials/add-kafka-relay.md | 12 ---- docs/tutorials/first-outbox-app.md | 11 ++- docs/usage/dlq.md | 13 ++-- docs/usage/fastapi.md | 22 +++++- docs/usage/messaging-service.md | 2 +- docs/usage/observability.md | 15 ++--- docs/usage/relay.md | 23 +++++-- docs/usage/router.md | 13 ++-- docs/usage/schema-validation.md | 16 +++++ docs/usage/setup-prometheus-opentelemetry.md | 20 +++--- docs/usage/subscriber.md | 11 ++- docs/usage/testing.md | 40 ++++++----- faststream_outbox/schema_validation.py | 2 +- tests/test_unit.py | 2 +- 22 files changed, 186 insertions(+), 145 deletions(-) diff --git a/README.md b/README.md index 12fdeeb..e8ad3c3 100644 --- a/README.md +++ b/README.md @@ -108,7 +108,7 @@ See [How it works](https://faststream-outbox.modern-python.org/introduction/how- ## Optional extras - `faststream-outbox[asyncpg]` — asyncpg driver (enables `LISTEN/NOTIFY` for ~10ms idle dispatch) -- `faststream-outbox[fastapi]` — FastAPI integration via `OutboxRouter` +- `faststream-outbox[fastapi]` — FastAPI integration via `faststream_outbox.fastapi.OutboxRouter` - `faststream-outbox[validate]` — Alembic for `broker.validate_schema()` - `faststream-outbox[prometheus]` — Prometheus metrics adapter - `faststream-outbox[opentelemetry]` — OpenTelemetry metrics adapter diff --git a/docs/concepts/comparison.md b/docs/concepts/comparison.md index 62d5d76..8af77dc 100644 --- a/docs/concepts/comparison.md +++ b/docs/concepts/comparison.md @@ -30,8 +30,8 @@ production system: - the **retry-strategy hierarchy** (`ExponentialRetry`, `NoRetry`, …) enforcing `max_attempts` and `max_total_delay_seconds` uniformly; - **`validate_schema()`** via Alembic's `autogenerate.compare_metadata`; -- **drain semantics** on stop, with the `running` / `_stopping` two-flag - dance and parallel-gathered subscriber shutdown; +- **drain semantics** on stop: new fetches stop while in-flight handlers + finish, and subscribers drain concurrently; - **`LISTEN/NOTIFY`** short-circuit on top of polling, with NOTIFY suppression on future-dated rows and `timer_id` conflict no-ops; - **`timer_id` dedup** via a partial unique index plus @@ -39,8 +39,8 @@ production system: - the **DLQ atomicity CTE** that rolls back the DELETE when the DLQ insert fails. -None of these is hard individually; in aggregate they decide whether the -outbox survives the second year of production load. +None of these is hard individually. Together they make up most of the work +past an MVP. You also pick up the [Subscriber](../usage/subscriber.md), [Publisher](../usage/publisher.md), [Dead-letter queue](../usage/dlq.md), and @@ -74,9 +74,8 @@ when the **async-Python logical-replication tooling gap** is too thin to lean on: there is no async-native logical-decoding client comparable to Debezium's JVM connectors, so a Python CDC path means either running the JVM stack alongside your app or driving `pg_recvlogical` / a thin -`psycopg` replication-protocol wrapper yourself. That is the load-bearing -point for this project — a 2026-05-07 reassessment confirmed the gap had -not closed sufficiently to make CDC the recommended path here. +`psycopg` replication-protocol wrapper yourself. That gap is why this +project does not recommend CDC as the default path. **TL;DR.** Pick CDC when you already need WAL capture or have producers outside your control. Pick this when you own the producer and want diff --git a/docs/concepts/instrumentation-seams.md b/docs/concepts/instrumentation-seams.md index 2ab1082..01f0f25 100644 --- a/docs/concepts/instrumentation-seams.md +++ b/docs/concepts/instrumentation-seams.md @@ -39,7 +39,7 @@ This is the "spans + bus parity" mode the native middleware ## What the middleware seam *can't* observe -Three events fire **outside** the handler invocation, with no +Four events fire **outside** the handler invocation, with no `StreamMessage` in scope: - **`fetched` ticks (including empty fetches).** Emitted by the fetch @@ -57,6 +57,8 @@ Three events fire **outside** the handler invocation, with no the `max_deliveries` ceiling and was dropped *without invoking the handler*. No handler call = no `consume_scope`. The middleware has nothing to wrap. +- **`drain_timeout`.** Fired during `stop()` when the shutdown drain + overruns `graceful_timeout`. There is no handler scope at all. ## What the recorder seam observes naturally @@ -107,7 +109,6 @@ The "Both seams together" recipe in [Setup Prometheus and OpenTelemetry wires the recommended layout: native middleware on the broker, plus a `metrics_recorder` for the outbox-internal events. -This isn't redundancy — each seam fires for events the other can't see. A service that registers only the middleware seam loses every `lease_lost`, `fetched`, and `max_deliveries`-terminal signal. A service that registers only the recorder seam loses tracing. diff --git a/docs/concepts/performance.md b/docs/concepts/performance.md index 16c6cd5..f374b03 100644 --- a/docs/concepts/performance.md +++ b/docs/concepts/performance.md @@ -52,11 +52,10 @@ still emits one NOTIFY), and the batched-flush win is largest at low `max_worker ### Producer: NOTIFY dedup (automatic) -`broker.publish` / `publish_batch` used to emit one `SELECT pg_notify(...)` per -call, so N publishes to the same queue in one transaction cost N NOTIFY -round-trips — N−1 of them pure waste, since Postgres already coalesces identical -notifications per transaction at delivery. Since 0.11.0 the producer emits **one -`pg_notify` per (transaction, queue)**. It is default-on, has no knob, and is +`broker.publish` / `publish_batch` emit **one `pg_notify` per (transaction, +queue)**, so N publishes to the same queue in one transaction cost one NOTIFY +round-trip. Postgres already coalesces identical notifications per transaction +at delivery, so extra NOTIFYs would be pure waste. It is default-on, has no knob, and is behavior-preserving: the subscriber still gets its wake, fired inline at the first publish. You do not configure this — it just makes bulk publishing cheaper. See [How it works](../introduction/how-it-works.md) for the write path. diff --git a/docs/index.md b/docs/index.md index dddf8b6..b7275e5 100644 --- a/docs/index.md +++ b/docs/index.md @@ -12,8 +12,15 @@ integration for the **transactional outbox pattern** — a Postgres table is the message queue. A producer writes a domain entity and an outbox row in the *same* SQLAlchemy transaction; a subscriber polls the table with `FOR UPDATE SKIP LOCKED`, runs the handler, and deletes the row on -success. The table *is* the queue — no separate message bus, no relay -process, no Kafka. +success. + +There are two ways to use it: + +- Standalone: the table is the queue. Subscribers process rows in place, + with no separate message bus. +- Relay: a subscriber forwards each row to Kafka, RabbitMQ, NATS or Redis + with a stacked publisher decorator. See + [Relay to a foreign broker](usage/relay.md). ## Use it when @@ -72,6 +79,8 @@ one-line summary. - [How it works](introduction/how-it-works.md) — two-loop subscriber, lease-token invariant, at-least-once semantics, opt-in DLQ on terminal failure. +- [Performance](concepts/performance.md): round-trips and table churn, + and the levers that address each. - [Comparison](concepts/comparison.md) — vs writing your own, vs CDC, vs Kafka transactions, vs `LISTEN/NOTIFY`, vs Celery, vs FastStream foreign-broker direct. diff --git a/docs/introduction/how-it-works.md b/docs/introduction/how-it-works.md index f92a8c2..0dba312 100644 --- a/docs/introduction/how-it-works.md +++ b/docs/introduction/how-it-works.md @@ -17,13 +17,16 @@ The outbox solves this by collapsing both writes into a single database transaction. Instead of publishing to a broker, you `INSERT` a row into an `outbox` table on the same `AsyncSession` that holds your domain write. A separate process polls the table and forwards rows to their consumers. The -row commits with your domain write or rolls back with it — atomicity is -free. +row commits or rolls back with your domain write. -`faststream-outbox` collapses the third "separate process" into the -subscriber itself: the same Postgres table holds the queue, and the -subscriber's polling loop *is* the consumer. No relay process, no Kafka, no -Rabbit. +In `faststream-outbox` that separate process is an outbox subscriber, and +you choose what it does with each row: + +- Standalone: the Postgres table is the queue and the subscriber's handler + is the consumer. No other message bus is involved. +- Relay: the subscriber forwards each row to Kafka, RabbitMQ, NATS or Redis + through a foreign-broker publisher decorator. See + [Relay to a foreign broker](../usage/relay.md). *See [Comparison](../concepts/comparison.md) for when CDC or Kafka transactions are the better fit.* diff --git a/docs/operations/alembic.md b/docs/operations/alembic.md index 1f77234..38ce059 100644 --- a/docs/operations/alembic.md +++ b/docs/operations/alembic.md @@ -367,7 +367,7 @@ op.execute(""" created_at TIMESTAMPTZ NOT NULL, failed_at TIMESTAMPTZ NOT NULL DEFAULT now(), failure_reason VARCHAR(64) NOT NULL, - last_exception TEXT, + last_exception VARCHAR, timer_id VARCHAR(255), PRIMARY KEY (id, failed_at) ) PARTITION BY RANGE (failed_at); diff --git a/docs/operations/troubleshooting.md b/docs/operations/troubleshooting.md index 419df05..60ff4c6 100644 --- a/docs/operations/troubleshooting.md +++ b/docs/operations/troubleshooting.md @@ -15,22 +15,25 @@ and a link into the reference page that owns the underlying design. | [Rolling deploy leaks rows](#rolling-deploy-leaks-rows) | `graceful_timeout` < handler P99, or k8s grace too short | | [`activate_in` / `activate_at` fires immediately in tests](#activate_in-activate_at-fires-immediately-in-tests) | `TestOutboxBroker(run_loops=False)` ignores scheduling | | [`AckPolicy.ACK_FIRST` raises `ValueError` at registration](#ackpolicyack_first-raises-valueerror-at-registration) | By design (would defeat outbox reliability) | -| [`OutboxResponse(...)` + foreign-publisher decorator gets nacked](#outboxresponse-foreign-publisher-decorator-gets-nacked) | By design (dual-fire footgun) | +| [`OutboxResponse(...)` + foreign-publisher decorator logs a configuration error](#outboxresponse-foreign-publisher-decorator-config-error) | By design (dual-fire footgun) | | [Chained `OutboxResponse` retries after the handler "succeeded"](#outboxresponse-relay-publish-failure) | Follow-on publish fails post-handler; nacks the inbound row | | [`validate_schema()` raises `ImportError`](#validate_schema-raises-importerror) | `[validate]` extra not installed | ## `event=lease_lost` recurring in logs { #event-lease_lost-recurring-in-logs } -**Symptom.** WARNING-level logs with structured field -`event=lease_lost`, typically with `phase=terminal` or `phase=retry`, -one per affected row. +**Symptom.** WARNING-level logs with the message text +`lease expired before terminal write` or `lease expired before retry write`, +one per affected row. The record also carries `event=lease_lost` and +`phase=terminal` / `phase=retry` as structured extras, visible when your log +formatter renders extras (for example a JSON formatter). **Likely cause.** The subscriber's `lease_ttl_seconds` is shorter than the handler's P99 duration. A handler took longer than the lease, another fetch reclaimed the row mid-flight, and the original handler's terminal `DELETE` / `UPDATE` matched zero rows. -**Diagnose.** Grep for `event=lease_lost` over the last hour and +**Diagnose.** Grep for `lease expired before` (or, with an +extras-rendering formatter, `event=lease_lost`) over the last hour and compare the rate against `dispatched`. A non-zero baseline rate (rather than occasional spikes) confirms TTL is the issue. @@ -58,8 +61,8 @@ DB (the `[validate]` extra is required). It will surface missing columns / indexes on the DLQ table. A frequent cause on older deployments is a hand-written DLQ migration missing the `timer_id` column — `validate_schema()` reports it as a missing column on the DLQ -table. (The [Alembic guide](../operations/alembic.md#adding-the-dlq-after-the-fact) -now includes it; pre-fix migrations may not.) +table. The [Alembic guide](../operations/alembic.md#adding-the-dlq-after-the-fact) +includes it. **Fix.** Bring the DLQ schema up to spec (apply the missing migration, or rename / drop the drifted column / index). After the schema is @@ -88,9 +91,8 @@ genuinely waiting to fire. **Diagnose.** Inspect a stuck row's `next_attempt_at` — if it's in the future, the row is correctly waiting. Otherwise check whether a -subscriber is registered: walk `broker.subscribers` (the *property*, -which covers router-attached subscribers — `broker._subscribers` will -miss them). +subscriber is registered: walk `broker.subscribers`, which covers +router-attached subscribers too. **Fix.** Register the subscriber, or adjust the producer's `activate_*` arg if the future date was unintentional. @@ -165,13 +167,13 @@ or the worker crashed between the handler's external side effect and the terminal `DELETE`. Both are at-least-once-delivery edge cases. **Diagnose.** Cross-reference handler-side logs (the side effect) -with `event=lease_lost` logs. Matching row IDs confirm TTL is too +with `lease expired before` warnings (`event=lease_lost` with an +extras-rendering formatter). Matching row IDs confirm TTL is too short. Crash-induced duplicates correlate with worker-process restarts. -**Fix.** Two layers: (a) make handlers idempotent — this is a -contract of the outbox pattern, not a knob, and (b) tune -`lease_ttl_seconds` above handler P99 so healthy handlers don't +**Fix.** Handlers must be idempotent; delivery is at-least-once. Also +tune `lease_ttl_seconds` above handler P99 so healthy handlers don't race their lease. **Reference.** [How it works § At-least-once @@ -188,16 +190,16 @@ nominally healthy. Drain duration appears longer than expected. **Likely cause.** Either the broker's `graceful_timeout` is shorter than the in-flight handler's remaining work, or Kubernetes `terminationGracePeriodSeconds` is shorter than the broker's -`graceful_timeout` × parallel-drain factor — `SIGKILL` arrives mid- -drain. +`graceful_timeout` (subscribers drain concurrently, so a clean shutdown +takes about one `graceful_timeout`), and `SIGKILL` arrives mid-drain. **Diagnose.** Time a clean shutdown locally (`docker compose kill -s SIGTERM application`) and compare to your k8s grace period. Look for log lines indicating drain abandonment. **Fix.** Raise `graceful_timeout` past handler P99 + margin. Raise -`terminationGracePeriodSeconds` past `graceful_timeout` + buffer for -parallel-subscriber drain. The `dispatch_one` shutdown-race guard is +`terminationGracePeriodSeconds` past `graceful_timeout` plus a +buffer. The `dispatch_one` shutdown-race guard is always on; you don't need to opt into it. **Reference.** [Production checklist § Drain & @@ -218,7 +220,7 @@ sync mode, expected immediate firing. **Fix.** Opt into `TestOutboxBroker(broker, run_loops=True)` for tests that need scheduled delivery to actually wait. Loop mode runs -the real `_fetch_loop` / `_worker_loop` against the fake client. +the real fetch and worker loops against the fake client. **Reference.** [Testing § Loop-driven mode](../usage/testing.md#loop-driven-mode), [Timers § Test broker @@ -244,18 +246,18 @@ exception via the configured retry strategy), or **Reference.** [Subscriber § Ack policy](../usage/subscriber.md#ack-policy). -## `OutboxResponse(...)` + foreign-publisher decorator gets nacked { #outboxresponse-foreign-publisher-decorator-gets-nacked } +## `OutboxResponse(...)` + foreign-publisher decorator logs a configuration error { #outboxresponse-foreign-publisher-decorator-config-error } **Symptom.** A handler with both `@kafka_pub` and an -`OutboxResponse(...)` return value gets nacked on every dispatch, with -a `_OutboxConfigError` logged. +`OutboxResponse(...)` return value logs an ERROR on every dispatch: +`Outbox configuration error (fix required; row left to lease-expiry retry)`. **Likely cause.** By design. The combination would both insert a row -into the outbox *and* publish to Kafka — a dual-fire that doubles -delivery. The subscriber refuses the chain composition by raising -`_OutboxConfigError` via the `process_message` override; it rides the -normal nack path so the row is retried (and logged) until the -configuration is fixed. +into the outbox *and* publish to Kafka, a dual-fire that doubles +delivery. The subscriber refuses the chain composition after the +handler returns. The worker logs the error and moves on without +nacking, so the retry strategy is not consulted. The row's lease expires, a later fetch reclaims +it, and the cycle repeats until the configuration is fixed. **Diagnose.** Inspect the handler decorator stack and return type. @@ -277,12 +279,11 @@ handler's work. **Likely cause.** The follow-on `OutboxResponse` row is published **after** the handler returns, inside the same consume scope — so a failure there (e.g. a DB error on the follow-on insert) unwinds through the -`AcknowledgementMiddleware` and nacks the inbound row. There is currently +`AcknowledgementMiddleware` and nacks the inbound row. There is **no distinct signal** separating "handler OK, relay-publish failed" from -an ordinary handler exception (F5-03): the metric reads as a normal +an ordinary handler exception: the metric reads as a normal `nacked_retried`/`retry_terminal`, and the ERROR log shows the publish -exception rather than a handler one. (The most common trigger — header -propagation re-encoding conflicts — was removed in #85.) +exception rather than a handler one. **Diagnose.** Read the logged exception: a `sqlalchemy`/`asyncpg` error or an envelope `ValueError` naming `content-type`/`correlation_id` points at @@ -296,8 +297,11 @@ pass a deterministic `timer_id` so a redelivery's insert is a no-op. ## `validate_schema()` raises `ImportError` { #validate_schema-raises-importerror } -**Symptom.** Calling `await broker.validate_schema()` raises -`ImportError("requires alembic")`. +**Symptom.** Calling `await broker.validate_schema()` raises: + +```text +ImportError: validate_schema() requires alembic. Install with `pip install 'faststream-outbox[validate]'`. +``` **Likely cause.** The `[validate]` extra isn't installed. Alembic is an optional dependency by design — every other code path works diff --git a/docs/tutorials/add-kafka-relay.md b/docs/tutorials/add-kafka-relay.md index 4ac0bf4..eff8181 100644 --- a/docs/tutorials/add-kafka-relay.md +++ b/docs/tutorials/add-kafka-relay.md @@ -220,18 +220,6 @@ publisher decorator → Kafka topic. Press `Ctrl-C` to stop the consumer. ## What about Kafka downtime? - - If Kafka were unavailable when the outbox subscriber dispatched a row, the foreign publish would raise, the outbox row would be nacked, and the configured `retry_strategy` would reschedule it. The next dispatch diff --git a/docs/tutorials/first-outbox-app.md b/docs/tutorials/first-outbox-app.md index 40ee452..b6097c5 100644 --- a/docs/tutorials/first-outbox-app.md +++ b/docs/tutorials/first-outbox-app.md @@ -41,7 +41,7 @@ Installed 24 packages in 37ms + click==8.4.1 + fast-depends==3.0.8 + faststream==0.7.1 - + faststream-outbox==0.8.0 + + faststream-outbox==0.14.1 + greenlet==3.5.1 + idna==3.18 + mako==1.3.12 @@ -61,8 +61,8 @@ Installed 24 packages in 37ms ``` Your exact pinned versions will differ; that is fine. The Python version -line will reflect whatever `uv` resolves on your machine — 3.13 or 3.14 -are both fine. +line will reflect whatever `uv` resolves on your machine; any Python 3.11 +or newer is fine. ## Step 2: Start Postgres @@ -302,9 +302,8 @@ The interesting property is what happened *inside* `publish_one`: the `broker.publish` call inserted a row into the outbox table through the session you opened. `session.begin()` committed it. If that commit had rolled back — say, because a domain write on the same session -failed — the outbox row would have rolled back with it. There is no -universe where the row exists but the domain write doesn't, or vice -versa. That atomicity is the whole point. +failed — the outbox row would have rolled back with it. The row and +the domain write commit or roll back together. That atomicity is the whole point. ## Clean up diff --git a/docs/usage/dlq.md b/docs/usage/dlq.md index fc1d060..d5cd077 100644 --- a/docs/usage/dlq.md +++ b/docs/usage/dlq.md @@ -75,8 +75,8 @@ same transaction, so the constraint would be unsatisfiable. There is also no ## Atomicity -When `dlq_table` is configured, `OutboxClient.delete_with_lease` switches to -a single CTE statement: +When `dlq_table` is configured, the terminal delete becomes a single CTE +statement: ```sql WITH deleted AS ( @@ -99,7 +99,7 @@ Two operator-visible properties fall out of this shape: lease-token guard documented in [Subscriber](./subscriber.md) is preserved. - **DLQ-write failure rolls back the DELETE.** If the INSERT fails - (column mismatch, disk full, ENUM violation), the whole statement + (column mismatch, disk full, a violated constraint), the whole statement rolls back. The outbox row stays leased and is reclaimed when the lease expires. Misconfiguration surfaces as outbox growth plus `lease_lost` spikes rather than silent audit loss. @@ -109,8 +109,7 @@ round-trip per terminal flush, same cost as the no-DLQ path. ## `last_exception` truncation -The serialized exception (`repr(exc)`) is bounded at 8 KiB by -`_LAST_EXCEPTION_MAX_CHARS` in `faststream_outbox/subscriber/usecase.py`. +The serialized exception (`repr(exc)`) is bounded at 8 KiB. Anything longer is truncated and `…[truncated]` appended. Rationale: some exceptions carry MB-scale payloads — pydantic validation @@ -149,8 +148,8 @@ opt-in install + `/health` pattern. ## Metric: `dlq_written` -`_flush_terminal` emits a `dlq_written` recorder event after the CTE -commits successfully. Skipped on the lease-lost path (no audit row was +The subscriber emits a `dlq_written` recorder event after the terminal +flush with the DLQ insert commits successfully. Skipped on the lease-lost path (no audit row was written, so nothing to count). Tags: diff --git a/docs/usage/fastapi.md b/docs/usage/fastapi.md index 14498b9..e4176aa 100644 --- a/docs/usage/fastapi.md +++ b/docs/usage/fastapi.md @@ -165,6 +165,22 @@ in handlers for dependencies. ## Engine ownership -The caller owns the `AsyncEngine`. `OutboxBroker` does **not** close it — -typically your FastAPI app does, via `app.add_event_handler("shutdown", -engine.dispose)` or its lifespan context manager. +The caller owns the `AsyncEngine`. `OutboxBroker` does **not** close it. +Dispose it in your app's lifespan: + +```python +from collections.abc import AsyncIterator +from contextlib import asynccontextmanager + +from fastapi import FastAPI + + +@asynccontextmanager +async def lifespan(app: FastAPI) -> AsyncIterator[None]: + yield + await engine.dispose() + + +app = FastAPI(lifespan=lifespan) +app.include_router(outbox_router) +``` diff --git a/docs/usage/messaging-service.md b/docs/usage/messaging-service.md index fea7892..6092556 100644 --- a/docs/usage/messaging-service.md +++ b/docs/usage/messaging-service.md @@ -70,7 +70,7 @@ class Resources(Group): outbox_broker = providers.Factory( scope=Scope.APP, creator=lambda engine: OutboxBroker(engine, outbox_table=OUTBOX_TABLE), - kwargs={"engine": Resources.database_engine}, + kwargs={"engine": database_engine}, ) ``` diff --git a/docs/usage/observability.md b/docs/usage/observability.md index a858096..357f0a6 100644 --- a/docs/usage/observability.md +++ b/docs/usage/observability.md @@ -14,14 +14,13 @@ catalog, and the operator PromQL playbook. MetricsRecorder = Callable[[str, Mapping[str, Any]], None] ``` -The default (`_noop_recorder`) lets instrumentation sites call -unconditionally. The recorder threads through `OutboxBrokerConfig` to: +The default is a no-op recorder. The broker passes the recorder to: -- The subscriber's seven core events via `OutboxSubscriber._emit_metric` +- Every subscriber, which emits seven core events (`fetched`, `dispatched`, `acked`, `nacked_retried`, `nacked_terminal`, `lease_lost`, `drain_timeout`), plus a conditional `dlq_written` when a DLQ is configured -- The producer's single event (`published`) via `OutboxProducer._emit_metric` +- The producer, which emits a single event (`published`) ### Bare seam @@ -147,11 +146,11 @@ dlq_written](./dlq.md#metric-dlq_written). ## Test broker note -`TestOutboxBroker` patches `broker.publish` directly via -`mock.patch.object`, bypassing `_basic_publish` — so middleware-registered +`TestOutboxBroker` replaces `broker.publish` with a patched version that +skips the middleware publish path, so middleware-registered **publish-scope** metrics do **not** fire in test mode. Middleware -**consume-scope** metrics still fire (because `dispatch_one` calls -`self.consume()` which walks the middleware stack normally). +**consume-scope** metrics still fire, because handlers are still consumed +through the normal middleware stack. The recorder-seam `published` event provides synthetic publish-side coverage in test mode via `FakeOutboxProducer`. The synthetic events use diff --git a/docs/usage/relay.md b/docs/usage/relay.md index 567eaa6..ea90b0b 100644 --- a/docs/usage/relay.md +++ b/docs/usage/relay.md @@ -159,11 +159,26 @@ that `broker.include_router(router)` must happen *before* the brokers start. Inside `FastAPI(..., lifespan=...)` the include happens during app construction (before lifespan), so it's automatic. -For the standalone (non-FastAPI) lifecycle, the order is: +For the standalone (non-FastAPI) lifecycle, use the plain broker routers +(`faststream.kafka.KafkaRouter` and `faststream_outbox.OutboxRouter`; the +FastAPI routers above only mount on a FastAPI app) and include them before +starting: ```python -# kafka_router / outbox_router are constructed exactly as in the FastAPI -# example above (KafkaRouter(...) / OutboxRouter(...)). +from faststream.kafka import KafkaRouter +from faststream_outbox import OutboxRouter + +kafka_router = KafkaRouter() +outbox_router = OutboxRouter() +publisher_kafka = kafka_router.publisher("kafka_topic") + + +@publisher_kafka +@outbox_router.subscriber("outbox_queue") +async def relay(body: dict) -> dict: + return body + + broker_kafka.include_router(kafka_router) broker_outbox.include_router(outbox_router) # then start @@ -180,7 +195,7 @@ from faststream_outbox import OutboxResponse @publisher_kafka @broker_outbox.subscriber("outbox_queue") async def relay(body: dict) -> OutboxResponse: - return OutboxResponse(body=body, queue="next_queue", session=...) # rejected at dispatch + return OutboxResponse(body=body, queue="next_queue", session=session) # rejected at dispatch ``` This would both insert a row into the outbox AND publish to Kafka. The diff --git a/docs/usage/router.md b/docs/usage/router.md index 4f4f7c7..c3c1ccb 100644 --- a/docs/usage/router.md +++ b/docs/usage/router.md @@ -79,18 +79,13 @@ All `@broker.subscriber` options (`max_workers`, `retry_strategy`, ## Gotcha: walking every subscriber Subscribers registered via `OutboxRouter` (then -`broker.include_router(router)`) live on the router, not on -`broker._subscribers`. If you need to introspect every subscriber on a -broker — counting active queues, asserting on schema, etc. — walk -`broker.subscribers` (the property): +`broker.include_router(router)`) live on the router. If you need to +introspect every subscriber on a broker (counting active queues, +asserting on schema, etc.), walk `broker.subscribers`: ```python for sub in broker.subscribers: ... ``` -The property iterates `[*broker._subscribers, -*(s for r in broker.routers for s in r.subscribers)]`, so it covers both -inline and router-attached subscribers. The bare `broker._subscribers` -(a `WeakSet`, set by the upstream `Registrator`) holds only the inline -ones and will silently miss everything attached via a router. +The property covers both inline and router-attached subscribers. diff --git a/docs/usage/schema-validation.md b/docs/usage/schema-validation.md index d746954..7311db0 100644 --- a/docs/usage/schema-validation.md +++ b/docs/usage/schema-validation.md @@ -11,6 +11,11 @@ of truth; the validator never duplicates the schema declaration. When the broker was constructed with a `dlq_table`, `validate_schema()` runs a second pass over the DLQ table the same way. +Alembic's comparison cannot see partial-index predicates or check +constraints, so the validator also queries `pg_catalog` (`pg_index`, +`pg_constraint`) directly and reports a missing or drifted partial-index +predicate or lease check constraint. + ## Install Alembic is an **optional dependency**: @@ -33,6 +38,17 @@ Raises `RuntimeError` if the live table is missing what the broker needs — absent table, missing columns, mismatched column types, flipped nullability, missing partial indexes. +Pass `check_autovacuum=True` to also check that the outbox table carries +the recommended autovacuum reloptions (see +[Alembic migrations](../operations/alembic.md), `outbox_autovacuum_ddl`). +An untuned table raises a `RuntimeError` whose message starts with +`Outbox autovacuum not tuned:`, separate from any schema mismatch. The +default (`False`) skips this check. + +```python +await broker.validate_schema(check_autovacuum=True) +``` + Extras are intentionally ignored: the validator only flags **missing** schema (`add_*` / `modify_*` ops). `remove_*` ops are silently dropped so you can attach your own audit columns or additional indexes without the diff --git a/docs/usage/setup-prometheus-opentelemetry.md b/docs/usage/setup-prometheus-opentelemetry.md index 0ff08b8..291626f 100644 --- a/docs/usage/setup-prometheus-opentelemetry.md +++ b/docs/usage/setup-prometheus-opentelemetry.md @@ -160,9 +160,9 @@ pip install 'faststream-outbox[opentelemetry,prometheus]' \ Here OpenTelemetry supplies **spans** (exported to OTLP) and Prometheus supplies **all metrics** (two registries scraped over HTTP). The OTel -middleware runs span-only — its meters would otherwise land on the global -Prometheus registry, which neither endpoint below exposes, so they are left -off to avoid a dead, unscraped meter path. +middleware runs span-only: it gets no `meter_provider`, so its meters go to +the global OpenTelemetry meter provider, which is a no-op unless you set one. +Neither endpoint below exposes them. ```python # app.py — run with `OTEL_EXPORTER_OTLP_ENDPOINT=http://localhost:4317 \ @@ -191,7 +191,7 @@ trace.set_tracer_provider(tracer_provider) # Two registries: the middleware and the recorder both define the same # faststream_* consume/publish collectors, so sharing one registry raises -# "Duplicated timeseries in CollectorRegistry" at broker construction. +# "Duplicated timeseries in CollectorRegistry" when the second one is created. MIDDLEWARE_REGISTRY = CollectorRegistry() RECORDER_REGISTRY = CollectorRegistry() @@ -236,15 +236,17 @@ contributes spans only, per the note above.) **The two seams overlap on consume/publish series.** Both the middleware and the recorder emit the same `faststream_received_*` / `faststream_published_*` collectors, which is why they must live on **separate registries** (above) — -sharing one raises `Duplicated timeseries in CollectorRegistry` at broker -construction, and summing across both double-counts every consume and +sharing one raises `Duplicated timeseries in CollectorRegistry` as soon as +the second of them is created, and summing across both double-counts every consume and publish. Treat the middleware as the source of truth for consume/publish; the recorder's unique value is the outbox-internal events the middleware can't see (`fetched`, `lease_lost`, terminal reasons, `dlq_written`). The providers set `messaging.system = "outbox"`, matching the -recorder-seam adapters. The OTel provider maps `row.id → +recorder-seam adapters. On consume spans the OTel provider maps `row.id → messaging.message.id`, `row.queue → messaging.destination_publish.name`, -`correlation_id → messaging.message.conversation_id`, `len(payload) → -messaging.message.payload_size_bytes`, and `len(cmd.batch_bodies) → +`correlation_id → messaging.message.conversation_id`, and `len(payload) → +messaging.message.payload_size_bytes`. On publish spans it maps the target +`queue → messaging.destination.name`, `correlation_id → +messaging.message.conversation_id`, and `len(cmd.batch_bodies) → messaging.batch.message_count` when >1. diff --git a/docs/usage/subscriber.md b/docs/usage/subscriber.md index 5071222..a8683f7 100644 --- a/docs/usage/subscriber.md +++ b/docs/usage/subscriber.md @@ -108,12 +108,11 @@ Per-subscriber knobs, passed to `@broker.subscriber("…", …)`: async def handle_urgent(body: dict) -> None: ... ``` -`OutboxSubscriberConfig.__post_init__` (in `subscriber/config.py`) warns or -raises on likely-wrong combinations (`lease_ttl_seconds <= max_fetch_interval`, -`max_deliveries` without retry, `min_fetch_interval > max_fetch_interval`, -etc.). Validation lives on the config, not the factory, so **every** -construction path — `@broker.subscriber`, `@router.subscriber`, direct -construction — is checked. +Subscriber options are validated when the subscriber is created: likely-wrong +combinations (`lease_ttl_seconds <= max_fetch_interval`, `max_deliveries` +without retry, `min_fetch_interval > max_fetch_interval`, etc.) warn or +raise. **Every** construction path (`@broker.subscriber`, `@router.subscriber`, +`OutboxRoute`) is checked. The table above lists the outbox-specific knobs. The standard FastStream subscriber kwargs pass through unchanged too: `dependencies`, `parser`, diff --git a/docs/usage/testing.md b/docs/usage/testing.md index 0ba15ba..190cc12 100644 --- a/docs/usage/testing.md +++ b/docs/usage/testing.md @@ -92,8 +92,8 @@ async def test_publisher() -> None: ``` `broker.publisher("q").publish(...)` lands rows in the same fake store as -`broker.publish(queue="q", ...)` — the test broker swaps the producer slot -for a `FakeOutboxProducer` via the FastStream `_basic_publish` flow. The one +`broker.publish(queue="q", ...)`, because the test broker swaps in a fake +producer that both paths go through. The one difference in tests is the `session`: `broker.publish` is patched to make it optional, but the publisher path is not, so pass a (mock) `AsyncSession` as shown. @@ -101,7 +101,7 @@ shown. ## Loop-driven mode For tests that exercise real polling semantics — retry rescheduling, lease -expiry / reclaim, `_fetch_loop` error recovery, or honoring `activate_in` +expiry / reclaim, fetch-loop error recovery, or honoring `activate_in` delays — opt in with `run_loops=True`: ```python @@ -134,9 +134,8 @@ assert received == [{"order_id": 1}] `feed(queue, payload, *, headers=None, next_attempt_at=None, timer_id=None)` inserts a row straight into the in-memory store and (in loop mode) wakes the fetch loop like a production NOTIFY would. In loop mode, the real -`_fetch_loop` / `_worker_loop` run against the fake client. Subscribers without registered handlers are skipped in -`_fake_start` (mirrors `OutboxSubscriber.start`'s `if not self.calls: -return`). +fetch and worker loops run against the fake client. Subscribers without +registered handlers are not started, matching production. ## Notes @@ -154,24 +153,23 @@ return`). to another worker, but tests will only invoke the handler once. Idempotency must be verified separately. Use `run_loops=True` for tests that need to observe lease-expiry behavior. -- **`FakeOutboxClient.validate_schema()` raises `NotImplementedError`** — - there is no real DB to validate against, and a silent pass would let - users ship broken schemas while their tests stay green. Tests that need - real schema validation must construct an `OutboxClient(real_engine, - table)` against the same DSN the migrations ran against. +- **`broker.validate_schema()` raises `NotImplementedError` under + `TestOutboxBroker`**: there is no real DB to validate against, and a + silent pass would let users ship broken schemas while their tests stay + green. Tests that need real schema validation must call + `validate_schema()` on an `OutboxBroker(real_engine, outbox_table=...)` + outside the test broker, against the same DSN the migrations ran against. ## Limitations of the fake broker -`TestOutboxBroker._fake_start` deliberately **skips the parent's -publisher-iteration loop** (the one that calls -`create_publisher_fake_subscriber`). FastStream's publisher-spy -infrastructure mocks the registered handler to forward -`publisher.publish()` calls — which conflicts with the outbox's real -dispatch path (the fake producer already lands rows in the fake client -*and* drives the real handler via `_sync_dispatch`). - -If you need FastStream's publisher-mock semantics for an outbox test, -swap that override out before re-using the parent's `_fake_start`. +`TestOutboxBroker` does not install FastStream's publisher spies (the +fake subscribers that `TestKafkaBroker` and friends attach to each +publisher). Those spies mock the registered handler to forward +`publisher.publish()` calls, which would conflict with the outbox's real +dispatch path: the fake producer already lands rows in the fake client +*and* drives the real handler. Publisher-mock assertions such as +`publisher.mock.assert_called_once_with(...)` are therefore not available; +assert on handler side effects or `fake_client.rows` instead. ## pytest-asyncio configuration diff --git a/faststream_outbox/schema_validation.py b/faststream_outbox/schema_validation.py index a76816a..434648c 100644 --- a/faststream_outbox/schema_validation.py +++ b/faststream_outbox/schema_validation.py @@ -275,7 +275,7 @@ def _run_validate( this flag. """ if not is_alembic_installed: - msg = "validate_schema() requires alembic. Install with `pip install faststream-outbox[validate]`." + msg = "validate_schema() requires alembic. Install with `pip install 'faststream-outbox[validate]'`." raise ImportError(msg) # Isolated MetaData containing ONLY the canonical table, so the user's diff --git a/tests/test_unit.py b/tests/test_unit.py index 9229bc0..c9fe7a9 100644 --- a/tests/test_unit.py +++ b/tests/test_unit.py @@ -1426,7 +1426,7 @@ def test_validate_schema_sync_raises_when_alembic_missing() -> None: # simulate "not installed" by flipping the boolean the function checks. with ( patch("faststream_outbox.schema_validation.is_alembic_installed", new=False), - pytest.raises(ImportError, match=r"pip install faststream-outbox\[validate\]"), + pytest.raises(ImportError, match=r"pip install 'faststream-outbox\[validate\]'"), ): _validate_schema_sync(MagicMock(), t) From 6bc58326b5d94de6c6cff090772f17d900c8aa4c Mon Sep 17 00:00:00 2001 From: Artur Shiriev Date: Sat, 3 Oct 2026 13:16:25 +0300 Subject: [PATCH 2/2] docs: validate the schema in tests through OutboxBroker, not the private client --- docs/usage/schema-validation.md | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/docs/usage/schema-validation.md b/docs/usage/schema-validation.md index 7311db0..d8d623b 100644 --- a/docs/usage/schema-validation.md +++ b/docs/usage/schema-validation.md @@ -125,6 +125,7 @@ asyncio.run(main()) is no real DB to validate against, and a silent pass would let users ship broken schemas while their `TestOutboxBroker`-backed tests stay green. -Tests that need real schema validation must construct an -`OutboxClient(real_engine, table)` against the same DSN the migrations -ran against. See [Testing](./testing.md). +Tests that need real schema validation must build an +`OutboxBroker(engine=real_engine, outbox_table=table)` against the same +DSN the migrations ran against and call `await broker.validate_schema()`. +See [Testing](./testing.md).