From a498de30fe29144023c74429105ea32fbb5a5940 Mon Sep 17 00:00:00 2001 From: Artur Shiriev Date: Sat, 3 Oct 2026 23:27:10 +0300 Subject: [PATCH] docs: show grouping subscribers under the FastAPI OutboxRouter --- docs/usage/fastapi.md | 49 ++++++++++++++++++++++++++++++++++++++++--- tests/test_fastapi.py | 33 +++++++++++++++++++++++++++++ 2 files changed, 79 insertions(+), 3 deletions(-) diff --git a/docs/usage/fastapi.md b/docs/usage/fastapi.md index 6001d8d..d7d8179 100644 --- a/docs/usage/fastapi.md +++ b/docs/usage/fastapi.md @@ -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 @@ -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 diff --git a/tests/test_fastapi.py b/tests/test_fastapi.py index 2103d42..ce8abb1 100644 --- a/tests/test_fastapi.py +++ b/tests/test_fastapi.py @@ -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, ) @@ -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"])