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
49 changes: 46 additions & 3 deletions docs/usage/fastapi.md
Original file line number Diff line number Diff line change
Expand Up @@ -128,6 +128,49 @@ in HTTP routes. In an HTTP route, reach the broker via `router.broker` (as
the quickstart's `create_order` does); a `broker: OutboxBroker` annotation
there resolves as a request field and fails with a 422.

## Grouping subscribers

To define subscribers in separate modules, put them on the plain
[`faststream_outbox.OutboxRouter`](./router.md) and include that into the
FastAPI router:

```python
# myapp/orders.py
from faststream_outbox import OutboxRouter

orders = OutboxRouter()


@orders.subscriber("orders")
async def handle(
body: dict,
session: AsyncSession = Depends(get_session),
) -> None: ...
```

```python
# myapp/main.py
from faststream_outbox.fastapi import OutboxRouter

from myapp.orders import orders

router = OutboxRouter(engine, outbox_table=outbox_table)
router.include_router(orders)

app = FastAPI()
app.include_router(router)
```

`router.include_router(...)` wraps each included subscriber with the same
FastAPI bridge as `@router.subscriber`, so `Depends(...)` resolves in
them. They start with the inner broker in the FastAPI lifespan and appear
in the document served at `schema_url`.

Two things do not work. Including one FastAPI `OutboxRouter` into another
raises `TypeError` (FastStream does not support nesting `StreamRouter`s),
and calling `router.broker.include_router(orders)` skips the bridge, so
`Depends(...)` in those handlers does not resolve.

## What's intentionally not exposed

Several `OutboxBroker.__init__` arguments are intentionally not exposed
Expand All @@ -139,9 +182,9 @@ on `OutboxRouter.__init__`:
- `dependencies`: on the router signature this means FastAPI
`Depends(...)` only; the broker's FastStream `Dependant` list is the
wrong shape for this flow.
- `routers`: not forwarded through the router; its semantics through the
FastAPI lifespan are unsettled. Register subscribers directly on the
`OutboxRouter` instead.
- `routers`: subscribers passed this way would skip the FastAPI bridge, so
`Depends(...)` in them would not resolve. Use `router.include_router(...)`
instead (see [Grouping subscribers](#grouping-subscribers)).

The [DLQ](./dlq.md) and the [metrics-recorder seam](./observability.md)
are also available through the router: pass `dlq_table=` and
Expand Down
33 changes: 33 additions & 0 deletions tests/test_fastapi.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
from sqlalchemy.ext.asyncio import AsyncSession

from faststream_outbox import NoRetry, OutboxMessage, OutboxResponse, make_dlq_table, make_outbox_table
from faststream_outbox import OutboxRouter as BrokerOutboxRouter
from faststream_outbox.fastapi import (
OutboxBroker as AnnotatedOutboxBroker,
)
Expand Down Expand Up @@ -287,3 +288,35 @@ async def handle(body: dict) -> None:
await router.broker.publish({"x": 1}, queue="orders") # ty: ignore[missing-argument]

assert events # the recorder seam is live under the FastAPI router


async def test_included_broker_router_subscriber_resolves_fastapi_depends() -> None:
"""INVARIANT: a broker-level ``OutboxRouter`` included into the FastAPI router gets the FastAPI bridge.

``StreamRouter.include_router`` wraps each nested subscriber with the FastAPI compatibility
decorator before handing the router to the broker. Including it on the broker directly (or via
``OutboxBroker(routers=...)``) skips that wrap, and the handler receives the ``Depends`` object
itself instead of the resolved value.
"""
router = OutboxRouter(outbox_table=_make_outbox_table())
sub = BrokerOutboxRouter()
seen: list[tuple[dict, str]] = []

def get_tenant() -> str:
return "tenant-a"

tenant_dep = Depends(get_tenant)

@sub.subscriber("orders")
async def handle(body: dict, tenant: str = tenant_dep) -> None:
seen.append((body, tenant))

router.include_router(sub)

app = _make_app_with_router(router)
with TestClient(app) as client:
await router.broker.publish({"x": 1}, queue="orders") # ty: ignore[missing-argument]
schema = client.get("/asyncapi.json").json()

assert seen == [({"x": 1}, "tenant-a")]
assert any(name.startswith("orders") for name in schema["channels"])
Loading