Server-Sent Events (SSE)¶
SSE pushes data from the server to the browser over one long-lived HTTP connection, no polling. It's the simplest path to "one-way real time": a notification feed, a progress bar, a price ticker, live logs.
SSE vs WebSocket vs Web Push
The SDK ships three pieces: EventStream (an in-memory async queue
feeding one connection), ServerSentEvent (encodes a frame in the spec
wire format) and sse_response (wraps the stream in a
StreamingResponse with the right headers — Cache-Control: no-cache,
Connection: keep-alive, X-Accel-Buffering: no to disable nginx
buffering). Day to day you call the shortcuts EventStream.response(...) /
SSEBroker.response(channel), which wrap sse_response under the hood; reach
for raw sse_response only when you want to drive the generator by hand.
Do you need to install anything? SSE is built in
EventStream, ServerSentEvent, sse_response and the in-memory
SSEBroker are part of the core — no extra of their own, they ship
with tempest-fastapi-sdk (they depend only on starlette, which FastAPI
already pulls in). There is no [sse] extra. Only the Redis bridge
(multi-worker) needs the [cache] extra —
uv add "tempest-fastapi-sdk[cache]", which pulls in redis. Cookie/query
auth uses JWTUtils from the [auth] extra.
New in v0.91
- Backpressure — the
EventStreamqueue is now bounded (max_queue, default1000): a slow client can't grow memory without limit. Theoverflowpolicy decides what gives. See Backpressure. - Lifecycle without boilerplate —
sse_response(..., on_disconnect=),EventStream.response(...)andSSEBroker.response(channel)tear down the producer / unregister the channel on their own when the client drops. - Query-string auth for cookieless clients (
EventSource). See Authentication.
One SSE endpoint¶
Create an EventStream per request, publish from a producer, and tie the
producer's lifecycle to the client connection — if the client drops, the
producer stops.
# src/api/routers/events.py
import asyncio
from fastapi import APIRouter
from starlette.responses import StreamingResponse
from tempest_fastapi_sdk import EventStream
router = APIRouter()
@router.get("/events")
async def events() -> StreamingResponse:
"""Emit 3 SSE frames then close the stream."""
stream = EventStream(heartbeat_seconds=15.0)
async def producer() -> None:
try:
for n in range(1, 4):
await stream.publish({"n": n}, event="counter", id=str(n))
await asyncio.sleep(1)
finally:
await stream.close()
task = asyncio.create_task(producer())
# on_disconnect runs when the client drops OR the stream ends:
# that's where you cancel the producer so it doesn't leak.
return stream.response(on_disconnect=task.cancel)
Always tie the producer to the connection
SSE streams are long-lived. If the client disconnects mid-stream you
don't want the producer running forever. Pass on_disconnect= to
EventStream.response (or sse_response) — it runs in the response
generator's finally, the one place that fires on disconnect.
Start the API and watch the raw frames in your terminal — curl -N
disables buffering and prints each frame as it arrives:
event: counter
id: 1
data: {"n": 1}
event: counter
id: 2
data: {"n": 2}
event: counter
id: 3
data: {"n": 3}
Note the spec wire format: each frame is a block of field: value lines
(event, id, then data), and the blank line (\n\n) separates one
frame from the next. Because you passed a dict, data was JSON-serialized
for you.
Before v0.91: hand-rolled try/finally
Up to v0.90 you wrapped stream() in an outer generator just to get
the finally. on_disconnect= replaces that boilerplate:
Anatomy of an event¶
publish() takes the four spec fields:
import asyncio
from tempest_fastapi_sdk import EventStream
stream = EventStream()
async def main() -> None:
"""Run this example."""
await stream.publish(
{"orderId": "abc", "status": "paid"}, # data: auto-JSON
event="order_update", # event: front-end listener name
id="42", # id: becomes Last-Event-ID (resume)
retry=3000, # retry: reconnect hint (ms)
)
asyncio.run(main())
| Field | What it does |
|---|---|
data |
Payload. String/bytes go raw; any object becomes JSON. |
event |
Event name — the front listens with addEventListener(name). Without it, falls back to "message". |
id |
Becomes Last-Event-ID; the browser resends it on reconnect so you can resume. |
retry |
Suggested reconnect delay (ms). |
The data type is the exported alias SSEData
(str | bytes | Mapping | Sequence | int | float | bool | None) — the JSON
value shapes, plus raw str/bytes. To send an object that only serializes via
str() (e.g. a bare UUID), wrap it in str(...) or a dict first.
heartbeat_seconds emits a beat while the stream is idle so load-balancers
don't cut the connection. By default that beat is an SSE comment
(: keepalive), invisible to EventSource: it fires no listener, it
just keeps the socket alive. None disables the heartbeat.
A visible beat: heartbeat_event¶
A comment keeps the TCP connection alive, which is the point — but a client
that uses the beat as proof the connection is live (to reconnect, to light
a status indicator, to arm a watchdog) has nothing to listen for.
heartbeat_event swaps the frame:
from tempest_fastapi_sdk import EventStream, ServerSentEvent
stream: EventStream = EventStream(
heartbeat_seconds=15.0,
heartbeat_event=ServerSentEvent(
data={"id": None, "type": "PING", "message": "ping"},
event="ping",
),
)
Now every 15 idle seconds put a ping event on the wire, which
addEventListener("ping", ...) receives like any other.
Need to stamp something that changes per beat — a timestamp, a counter, a queue depth? Pass a callable, resolved per beat:
from datetime import UTC, datetime
from tempest_fastapi_sdk import EventStream, ServerSentEvent
def ping() -> ServerSentEvent:
"""Build a fresh heartbeat frame carrying the current time.
Returns:
ServerSentEvent: The frame this beat puts on the wire.
"""
return ServerSentEvent(
data={"type": "PING", "at": datetime.now(UTC).isoformat()},
event="ping",
)
stream: EventStream = EventStream(heartbeat_seconds=15.0, heartbeat_event=ping)
The two choices are separate now, and that was the problem
heartbeat_seconds used to decide both: whether to beat and what the
beat puts on the wire. Anyone who wanted a visible ping had to disable
the heartbeat and reimplement it outside — a periodic loop publishing into
every open channel, which is generic infrastructure back inside the
service. Now heartbeat_seconds decides when and heartbeat_event
decides what.
SSEBroker(heartbeat_event=...) forwards it to every stream it opens,
which is the only way to configure this in a service that uses the broker:
register constructs the EventStream internally, so not even a subclass
gets into the path.
Backpressure (bounded queue)¶
If a client stops reading (backgrounded tab, bad network) but the
producer keeps publishing, the EventStream queue would grow forever — a
classic memory leak. So the queue is bounded: max_queue (default
1000) plus an overflow policy that decides what to do when it fills.
from tempest_fastapi_sdk import EventStream
# Live ticker: a stale frame is worthless -> drop the oldest.
stream = EventStream(max_queue=500, overflow="drop_oldest")
overflow |
When the queue fills | Use when |
|---|---|---|
"drop_oldest" (default) |
Evict the oldest event | Live data: ticker, progress, telemetry — only the recent state matters. |
"drop_newest" |
Discard the incoming event | The start of the stream matters more than the end. |
"block" |
Hold publish() until a slot frees |
Producer dedicated to one connection and losing an event is unacceptable. |
block can stall a shared producer
With overflow="block", a slow client holds publish(). If the
same producer feeds many clients (fan-out), one bad client stalls them
all. Only use block when the producer serves one connection.
The close() sentinel is never dropped or blocked — the stream always
terminates. stream.dropped_events counts how many events were lost to
overflow, so you can surface it in metrics/logs:
import logging
from tempest_fastapi_sdk import EventStream
logger = logging.getLogger(__name__)
stream = EventStream()
if stream.dropped_events:
logger.warning("Slow SSE client: %d events dropped", stream.dropped_events)
Back to the old behavior
max_queue=0 disables the bound (unbounded queue, pre-0.91). Only do
this if you're sure the producer stops together with the connection.
Broadcast to many clients (SSEBroker)¶
EventStream is one connection. To send the same event to every
client of a channel (e.g. a user's devices, or a topic), the SDK ships
SSEBroker — a per-channel stream registry plus fan-out. The channel is
any string (a user id, a room slug...).
It takes three steps: create the broker once, keep that instance on the app (so everyone uses the same one), and inject it into endpoints.
Step 1 — create the broker and the wiring¶
SSEBroker() is a process-wide singleton: every open channel and stream
lives inside it. So it must be one instance shared across the whole app —
if each request made its own, a publish on one broker would never reach the
streams pinned to another.
# src/api/dependencies/resources.py
from fastapi import FastAPI, Request
from tempest_fastapi_sdk import SSEBroker
broker = SSEBroker()
def register_broker(app: FastAPI) -> None:
"""Store the singleton broker on app.state (call it in create_app)."""
app.state.broker = broker
def get_broker(request: Request) -> SSEBroker:
"""Return the shared broker from app.state for use in Depends()."""
return request.app.state.broker
What each part does:
broker = SSEBroker()— creates the broker at module import. Since the module is imported once, this is the same object for everyone who uses it.register_broker(app)— pins the broker onapp.state.broker. You call it once when building the app (insidecreate_appor the lifespan):register_broker(app).get_broker(request)— returnsrequest.app.state.broker. This is what endpoints receive viaDepends(get_broker), guaranteeing they all talk to the same instance.
Step 2 — the subscribe endpoint¶
The client opens GET /feed; the endpoint subscribes it to its own user
channel and returns the stream. One line does it all: broker.response(channel).
# src/api/routers/feed.py
from uuid import UUID
from fastapi import APIRouter, Depends
from starlette.responses import StreamingResponse
from tempest_fastapi_sdk import SSEBroker
from src.api.dependencies.auth import get_current_user_id
from src.api.dependencies.resources import get_broker
router = APIRouter()
@router.get("/feed")
async def feed(
user_id: UUID = Depends(get_current_user_id),
broker: SSEBroker = Depends(get_broker),
) -> StreamingResponse:
"""Subscribe the caller to their own user channel and stream it."""
return broker.response(str(user_id))
Step by step, on each GET /feed:
Depends(get_current_user_id)resolves who the client is. Their id becomes the channel name — each user gets their own, isolated from others.Depends(get_broker)hands over the shared broker (the one from Step 1).broker.response(str(user_id))does three things in a single call:- register — creates a fresh
EventStreamand subscribes it to theuser_idchannel; - stream — returns a
StreamingResponsewith the SSE headers already set (the client starts receiving); - unregister — wires an
on_disconnectthat removes this stream from the channel when the client drops.
- register — creates a fresh
Why broker.response() doesn't leak streams
The three steps (register → stream → unregister) are tied together in one
call. The unregister runs in the response generator's finally — the one
point that fires on disconnect — so there's no try/finally for you to
forget: every client that leaves cleans up its own registration. And
SSEBroker(max_queue=..., overflow=...) applies the same backpressure
policy (see above) to every stream the broker opens.
A global notice: broadcast()¶
publish(channel, ...) resolves one channel to N connections. The orthogonal
axis — one event to N channels — is broadcast():
from tempest_fastapi_sdk import SSEBroker
async def warn_maintenance(broker: SSEBroker) -> None:
"""Warn every connected client, whichever channel they are on.
Args:
broker (SSEBroker): The process-wide broker.
"""
await broker.broadcast(
{"type": "MAINTENANCE", "message": "Maintenance in 5 minutes."},
event="notice",
)
Same parameters as publish, minus the channel. In single-process mode it
walks the local channels and delivers into each. With Redis it publishes to
the reserved __broadcast__ channel — which every worker's
PSUBSCRIBE {prefix}:* already reaches, no new subscription — and each
worker's run() re-fans it to all of its local streams.
A stream subscribed to more than one channel receives the event once: the fan-out is over the set of streams, not the list of channels.
__broadcast__ is a reserved name
register("__broadcast__") and publish("__broadcast__", ...) raise
ValueError. The name is how the receiving side decides between
delivering to one channel and delivering to all of them; a stream
subscribed under it would make that decision ambiguous. The constant is
BROADCAST_CHANNEL, in tempest_fastapi_sdk.sse.
And for a "how many are connected right now" metric, local_channels() lists
the channels with at least one open stream on this worker — the
counterpart of local_subscribers(channel), with no need to touch
_channels:
from tempest_fastapi_sdk import SSEBroker
def connected(broker: SSEBroker) -> dict[str, int]:
"""Count the live streams per channel on this worker.
Args:
broker (SSEBroker): The process-wide broker.
Returns:
dict[str, int]: Channel name to local subscriber count.
"""
return {
channel: broker.local_subscribers(channel)
for channel in broker.local_channels()
}
Firing from the domain (controller)¶
Step 2 showed the subscribe side (the client joins a channel). What's left is the publish side — the broadcast itself.
What "broadcast" means here: the broker keeps, per channel, the list of
subscribed streams. When you call broker.publish("<channel>", ...), it walks
every stream on that channel and delivers the same event to each. One
publish → N clients. If the channel is a user's id and they have two tabs
open, both receive it; if nobody is subscribed, the publish is a no-op (no
error).
The controller is what fires the publish, when a business event happens
(order created, payment confirmed, new message). It orchestrates the service
(business rule) and the broker (live notification), using the id of the user to
notify as the channel.
# src/controllers/order.py
from tempest_fastapi_sdk import SSEBroker
from src.schemas import OrderCreateSchema, OrderResponseSchema
from src.services import OrderService
class OrderController:
"""Orchestrates order creation and the live seller notification."""
def __init__(self, order_service: OrderService, broker: SSEBroker) -> None:
"""Wire the order service and the SSE broker.
Args:
order_service (OrderService): Order business logic.
broker (SSEBroker): Fan-out broker for live notifications.
"""
self.order_service = order_service
self.broker = broker
async def create_order(self, data: OrderCreateSchema) -> OrderResponseSchema:
"""Create an order and notify the seller in real time.
Args:
data (OrderCreateSchema): The order creation payload.
Returns:
The created order.
"""
order = await self.order_service.create(data)
await self.broker.publish(
str(order.seller_id), # channel = id of whoever gets the notice
{"order_id": str(order.id), "total": str(order.total)},
event="order_created",
id=str(order.id),
)
return order
The heart of it is the broker.publish(...) call, argument by argument:
- 1st argument (
str(order.seller_id)) — the channel. Sends the event to every stream subscribed to that id (here, the seller's devices). This is why the subscribe endpoint uses the user's id as the channel: both sides must agree on the same string. - 2nd argument (the dict) — the payload (SSE
data). A non-string object becomes JSON automatically (see Anatomy). event="order_created"— the event name; the front listens withaddEventListener("order_created", ...).id=str(order.id)— becomesLast-Event-ID, so the client can resume from the right spot if it reconnects.
Notice the controller never touches EventStream or an HTTP response: it
just hands the event to the broker, which handles the fan-out. Publishing is
fire-and-forget.
The provider builds the controller with its service and the broker injected —
the same get_broker from Step 1:
# src/api/dependencies/controllers.py
from fastapi import Depends
from tempest_fastapi_sdk import SSEBroker
from src.api.dependencies.resources import get_broker
from src.api.dependencies.services import get_order_service
from src.controllers import OrderController
from src.services import OrderService
def get_order_controller(
order_service: OrderService = Depends(get_order_service),
broker: SSEBroker = Depends(get_broker),
) -> OrderController:
"""Build an OrderController with its service and the SSE broker."""
return OrderController(order_service, broker)
The router just receives the controller via Depends and delegates — no
business rule, no loose publish in the route:
# src/api/routers/orders.py
from fastapi import APIRouter, Depends
from src.api.dependencies.controllers import get_order_controller
from src.controllers import OrderController
from src.schemas import OrderCreateSchema, OrderResponseSchema
router = APIRouter()
@router.post("/orders")
async def create_order(
data: OrderCreateSchema,
controller: OrderController = Depends(get_order_controller),
) -> OrderResponseSchema:
"""Create an order; the seller gets a live SSE notification."""
return await controller.create_order(data)
A buyer with GET /feed open receives it instantly:
Publishing from outside the request (queue, task, webhook)
broker.publish is just a coroutine — call it from anywhere that has the
broker: a FastStream consumer, a TaskIQ task, a webhook. It only reaches
who is connected right now (the channel registration drops on
disconnect); for a durable notification, persist it in the database and
treat SSE as the live layer on top. In multi-worker mode (Redis), publish
reaches the worker the client is pinned to — see below.
Multi-worker: Redis bridge (ready, no extra code)¶
An in-memory SSEBroker lives in one worker — with --workers N a
publish only reaches the clients pinned to that process. Give the broker a
Redis client and the same broker publishes via Redis PUBLISH; a
background task (run()) PSUBSCRIBE-s and relays to each worker's local
streams. Same call site, now horizontal.
Use the SDK's AsyncRedisManager ([cache] extra) to open the connection —
it's the same managed client (connect/disconnect/health-check) that cache,
sessions and feature flags use; no raw redis.asyncio:
# src/api/app.py
import asyncio
from collections.abc import AsyncGenerator
from contextlib import asynccontextmanager
from fastapi import FastAPI
from tempest_fastapi_sdk import SSEBroker
from tempest_fastapi_sdk.cache import AsyncRedisManager
from src.core.settings import settings
cache = AsyncRedisManager(**settings.redis_kwargs())
@asynccontextmanager
async def lifespan(app: FastAPI) -> AsyncGenerator[None, None]:
"""Connect Redis, wire the cross-worker broker, run its fan-out loop."""
await cache.connect()
broker = SSEBroker(redis=cache.client, channel_prefix="sse")
app.state.broker = broker
task = asyncio.create_task(broker.run())
try:
yield
finally:
task.cancel()
await broker.aclose()
await cache.disconnect()
app = FastAPI(lifespan=lifespan)
# broker.publish(...) on any worker -> reaches ALL workers
What each part does:
AsyncRedisManager(**settings.redis_kwargs())— the SDK's managed Redis client.redis_kwargs()comes from theRedisSettingsmixin (URL +decode_responses).await cache.connect()first — before that,cache.clientraisesRuntimeError. That's why the broker is built inside the lifespan, not at module import.SSEBroker(redis=cache.client, ...)—cache.clientis the underlying rawredis.asyncio.Redis; the broker uses it forPUBLISH/PSUBSCRIBE.app.state.broker = broker— the same wiring as Step 1, so the endpoints'get_brokerstays identical.asyncio.create_task(broker.run())— background task that subscribes to Redis and relays each event to the worker's local streams. On teardown: cancel the task, close the broker, disconnect Redis.
Start simple, scale later
Without Redis, SSEBroker() already covers a single process. When you need
multiple workers/hosts, just pass the AsyncRedisManager's cache.client
and start run() in the lifespan — no endpoint changes, publish becomes
cross-process for free. AsyncRedisManager comes from the [cache] extra
(uv add "tempest-fastapi-sdk[cache]").
Multi-worker without Redis: bridge over RabbitMQ¶
The built-in bridge (SSEBroker.run()) is Redis-only — the redis= param
takes no other transport. If you already run RabbitMQ and don't want Redis
just for this, you can build the cross-worker fan-out by hand: each worker keeps
an in-memory SSEBroker() (no redis=), and a RabbitMQ subscriber in each
worker relays the event to the local broker.publish.
The piece that makes the broadcast work is the fanout exchange: every message published to it is copied to all bound queues. Each worker declares an exclusive queue (gone when the worker dies), so all of them receive every event — unlike RabbitMQ's default (work-queue), where only one consumer would get the message.
# src/api/app.py
from collections.abc import AsyncGenerator
from contextlib import asynccontextmanager
from fastapi import FastAPI
from faststream.rabbit import ExchangeType, RabbitExchange, RabbitQueue
from tempest_fastapi_sdk import SSEBroker
from tempest_fastapi_sdk.queue import MessageBroker
from src.core.settings import settings
broker = SSEBroker() # in-memory, one per worker
mq = MessageBroker.rabbitmq(settings.RABBITMQ_URL) # [queue] extra
sse_exchange = RabbitExchange("sse.fanout", type=ExchangeType.FANOUT)
worker_queue = RabbitQueue("", exclusive=True, auto_delete=True)
@mq.broker.subscriber(worker_queue, sse_exchange) # mq.broker = FastStream's RabbitBroker
async def _fan(evt: dict) -> None:
"""Relay one fanned-out event to this worker's local SSE streams."""
await broker.publish(evt["channel"], evt.get("data", ""), event=evt.get("event"))
@asynccontextmanager
async def lifespan(app: FastAPI) -> AsyncGenerator[None, None]:
await mq.connect() # opens the connection and starts the subscriber
app.state.broker = broker # the same get_broker the endpoints use resolves this
try:
yield
finally:
await mq.disconnect()
app = FastAPI(lifespan=lifespan)
What each part does:
SSEBroker()with noredis=— fan-out local to the worker (only the streams pinned to it).RabbitExchange(..., type=FANOUT)+RabbitQueue("", exclusive=True, auto_delete=True)— the exchange copies to every bound queue; each worker declares its own, so they all receive every event.@mq.broker.subscriber(worker_queue, sse_exchange)—mq.brokeris FastStream'sRabbitBroker(the escape hatch). TheMessageBroker.on(...)facade doesn't expose fanout, so we drop to the raw broker here.await mq.connect()in the lifespan starts the consumer;mq.disconnect()stops it.
The domain publishes to the exchange (not to broker.publish directly) — so
it reaches every worker:
import asyncio
from uuid import UUID
from src.db.models import OrderModel, UserModel
from src.queue import mq
user = UserModel(name="Ana", email="ana@example.com")
order = OrderModel(user_id=user.id, total=100)
sse_exchange = "sse.fanout"
user_id = UUID("2b1d0c2e-7f3a-4c56-9d18-2f9a4c5b6d70")
async def main() -> None:
"""Run this example."""
# from any worker/handler:
await mq.broker.publish(
{
"channel": str(user_id),
"event": "order_created",
"data": {"order_id": str(order.id)},
},
exchange=sse_exchange,
)
asyncio.run(main())
Redis is still the simpler path
RabbitMQ needs this fanout exchange + per-worker exclusive queue to replicate
what Redis pub/sub (SSEBroker.run()) does in one line. Prefer RabbitMQ when
it's already your infra; otherwise Redis ([cache]) is fewer moving parts.
Authentication (cookie or query string)¶
Here's the SSE gotcha: the browser's native EventSource can't send
a header. No Authorization: Bearer on the handshake. So there are two
ways to authenticate the stream.
Preferred: session cookie¶
If the front is on the same origin as the API, use an HttpOnly
cookie. The browser sends it on its own when you open with
withCredentials:
On the backend, guard the stream with your UserAuthService's
current_user_dependency() — the same service make_auth_router uses (see the
Auth flow recipe). With AUTH_TOKEN_DELIVERY set to
"cookie"/"both" it auto-derives the cookie_name from the login (the
header still wins when present), so there's nothing to wire:
# src/api/dependencies/auth.py
from src.services import auth_service # your UserAuthService instance
current_user = auth_service.current_user_dependency()
Why cookie is better
The token stays out of the URL: no leak into access logs, browser
history or the Referer header. HttpOnly also keeps the token out
of JavaScript's reach (XSS defense). Prefer this path whenever the
origin is shared.
Cookieless alternative: token in the query string¶
Without a session cookie (front on a different origin, a mobile app
opening a raw EventSource, an environment where withCredentials isn't
an option), pass the access token in the query string. As of v0.135
current_user_dependency accepts this via query_param:
# src/api/dependencies/auth.py
from src.services import auth_service # your UserAuthService instance
# Lookup order: header -> cookie -> query string.
current_user = auth_service.current_user_dependency(query_param="access_token")
# src/api/routers/feed.py
from uuid import UUID
from fastapi import APIRouter, Depends
from starlette.responses import StreamingResponse
from tempest_fastapi_sdk import SSEBroker
from src.db.models import User, UserModel
current_user = UserModel(name="Ana", email="ana@example.com")
def get_broker() -> SSEBroker:
"""Return the process-wide broker built at startup."""
return broker
broker = SSEBroker()
router = APIRouter()
@router.get("/feed")
async def feed(
user: User = Depends(current_user), # resolves the JWT from the query
broker: SSEBroker = Depends(get_broker),
) -> StreamingResponse:
"""Authenticated stream without a cookie — token comes in the URL."""
return broker.response(str(user.id))
On the front:
// The short-lived access token goes in the URL — never the refresh token.
const es = new EventSource(`/api/feed?access_token=${accessToken}`);
Query strings leak — treat the token as disposable
A token in the URL shows up in access logs, history and the
Referer header. Non-negotiable rules:
- Short-lived access token only (minutes). Never the refresh token.
- Always over TLS (HTTPS).
- Strip the value from your proxy/server log format.
- Refresh through a normal endpoint (header/cookie), not the query.
Under the hood: make_jwt_user_dependency / make_bearer_token_dependency
current_user_dependency wraps make_jwt_user_dependency (passing through
cookie_name/query_param). If you don't have a UserAuthService, call
the factory directly — make_jwt_user_dependency(jwt, load_user, query_param="access_token") —
or make_bearer_token_dependency(jwt, query_param="access_token") when you
only want the decoded claims. Same header → cookie → query order.
Aligned with tempest-react-sdk¶
tempest-react-sdk's createEventStream / useEventStream
(repo)
consumes these endpoints with built-in exponential-backoff reconnect:
import { createEventStream } from "tempest-react-sdk";
const stream = createEventStream<{ text: string }>("/feed", {
withCredentials: true, // sends the auth cookie on the handshake
namedEvents: ["notice"], // <- matches publish(event="notice")
onMessage: (m) => console.log(m.event, m.data), // data already JSON-parsed
});
// stream.close() to tear down; stream.reconnect() to force a reconnect
Heartbeat: comment vs a ping event
By default EventStream's heartbeat is a comment — EventSource
ignores it, so the react-sdk does not even need heartbeatEvents. For a
named, visible heartbeat, pass
heartbeat_event=ServerSentEvent(data="ping", event="ping") to
EventStream (or to SSEBroker, which forwards it to every stream) and
set heartbeatEvents: ["ping"] on the front end (its default). Publishing
the ping by hand becomes unnecessary.
Alignment points:
publish(event="x")↔namedEvents: ["x"]+onMessage.- non-string
databecomes JSON ↔ the react default parser decodes JSON. id=↔Last-Event-IDresent on reconnect (resume where you left off).- cookie auth ↔
withCredentials: true.
Recap¶
EventStream(one per connection) +.response()— an SSE endpoint with headers set (sse_responseis the low-level primitive underneath).- Tie the producer to the connection with
on_disconnect=(onEventStream.response,sse_responseorbroker.response) — no hand-rolledtry/finally. - Queue is bounded (
max_queue, default1000) +overflow(drop_oldest/drop_newest/block) prevents leaks from slow clients;dropped_eventscounts the discards. publish(data, event=, id=, retry=)covers the 4 spec fields; non-stringdatabecomes JSON.- Heartbeat:
heartbeat_secondsdecides when,heartbeat_eventdecides what. The default is a comment (invisible to EventSource); pass aServerSentEvent(or a callable, resolved per beat) for apingthe client listens to;Noneonheartbeat_secondsdisables it. - Broadcast =
SSEBroker;broker.response(channel)does register + response + unregister; publish from controllers/tasks/queues withbroker.publish(channel, ...); multi-worker = pass a Redis client + startbroker.run()in the lifespan. - A global notice =
broker.broadcast(...): one event to every channel, including on the other workers (reserved__broadcast__channel, already reached by the existingPSUBSCRIBE).local_channels()lists this worker's live channels. - Auth: cookie (
cookie_name+withCredentials) on the same origin; query string (query_param, short-lived access token over TLS only) for cookieless clients. tempest-react-sdkcreateEventStream/useEventStreamconsumes with reconnect;namedEvents↔publish(event=).