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
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
13 changes: 6 additions & 7 deletions docs/concepts/comparison.md
Original file line number Diff line number Diff line change
Expand Up @@ -30,17 +30,17 @@ 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
`on_conflict_do_nothing`;
- 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
Expand Down Expand Up @@ -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
Expand Down
5 changes: 3 additions & 2 deletions docs/concepts/instrumentation-seams.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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

Expand Down Expand Up @@ -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.
9 changes: 4 additions & 5 deletions docs/concepts/performance.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
13 changes: 11 additions & 2 deletions docs/index.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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.
Expand Down
15 changes: 9 additions & 6 deletions docs/introduction/how-it-works.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.*
Expand Down
2 changes: 1 addition & 1 deletion docs/operations/alembic.md
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
70 changes: 37 additions & 33 deletions docs/operations/troubleshooting.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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
Expand All @@ -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 &
Expand All @@ -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
Expand All @@ -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.

Expand All @@ -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
Expand All @@ -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
Expand Down
12 changes: 0 additions & 12 deletions docs/tutorials/add-kafka-relay.md
Original file line number Diff line number Diff line change
Expand Up @@ -220,18 +220,6 @@ publisher decorator → Kafka topic. Press `Ctrl-C` to stop the consumer.

## What about Kafka downtime?

<!--
Maintainer note: the spec for this tutorial originally proposed a live
"kill Kafka, watch the retry" step. We attempted it during authoring
(Confluent cp-kafka 7.6.0, ~10s and ~20s outage windows) and could not
get the outbox subscriber's retry log lines to surface — aiokafka's
client-side reconnect absorbs short outages internally, so no outbox-
level retry fires. The plan authorized falling back to a contract-
focused callout in lieu of a fragile live demo. If you re-attempt this
in the future and find a Kafka setup that reliably surfaces the retry,
this callout can give way to a real Step 6 again.
-->

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
Expand Down
11 changes: 5 additions & 6 deletions docs/tutorials/first-outbox-app.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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

Expand Down Expand Up @@ -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

Expand Down
Loading
Loading