Fanning Out SSE with Durable Objects Permalink to this section
Part of Edge & Serverless SSE Deployment, under Backend Stream Generation & Connection Management.
A Cloudflare Worker can stream a Server-Sent Events response, but each request runs in isolation: there is no shared memory in which one request’s publish can reach another request’s stream. Durable Objects solve exactly that. A Durable Object is a single-threaded, addressable instance with its own storage; every request routed to the same object id reaches the same instance. Make one object per room, document or channel, have it hold the open streams, and it becomes a natural fan-out point with ordering and replay built in. This guide builds that design.
Symptom & Developer Intent Permalink to this section
- A Worker streams events, but events published by one request never reach clients connected through other requests.
- Polling KV or a database from every stream is slow and expensive.
- Messages for a room arrive out of order when published from several places.
- Reconnecting clients miss events published while they were away.
The intent is per-channel fan-out at the edge with a single order of events, replay for reconnecting clients, and predictable cost.
Root Cause Analysis Permalink to this section
Workers are stateless per request. To fan out, every stream for a channel must be attached to something that also receives that channel’s publishes. A Durable Object is that something: routing by name (idFromName('room:42')) guarantees that all streams and all publishes for room 42 meet in one instance, which processes requests one at a time and therefore defines a single order.
One important constraint: Durable Objects offer a hibernation API for WebSockets, which lets an idle object be evicted from memory while connections stay open. SSE responses are ordinary streaming HTTP responses and do not benefit from it — an object holding open SSE streams stays active, and active duration is billed. For rooms that are idle most of the time with many connected clients, factor that into the cost model.
Step-by-Step Resolution Permalink to this section
Step 1 — Route requests to the channel’s object Permalink to this section
// worker.ts
export default {
async fetch(req: Request, env: Env): Promise<Response> {
const url = new URL(req.url);
const m = url.pathname.match(/^\/rooms\/([\w-]+)\/(events|publish)$/);
if (!m) return new Response('not found', { status: 404 });
// authenticate here, before touching the object
const stub = env.ROOMS.get(env.ROOMS.idFromName(m[1]));
return stub.fetch(req); // same object for every request about this room
},
};
Step 2 — Hold one writer per stream inside the object Permalink to this section
// room.ts
export class Room {
private writers = new Set<WritableStreamDefaultWriter<Uint8Array>>();
private seq = 0;
private enc = new TextEncoder();
constructor(private state: DurableObjectState, private env: Env) {
state.blockConcurrencyWhile(async () => {
this.seq = (await state.storage.get<number>('seq')) ?? 0;
});
}
async fetch(req: Request): Promise<Response> {
const path = new URL(req.url).pathname;
return path.endsWith('/publish') ? this.publish(req) : this.subscribe(req);
}
private async subscribe(req: Request): Promise<Response> {
const { readable, writable } = new TransformStream<Uint8Array, Uint8Array>();
const writer = writable.getWriter();
const after = Number(req.headers.get('Last-Event-ID') ?? 0);
writer.write(this.enc.encode('retry: 2000\n\n'));
for (const [, frame] of await this.state.storage.list<string>({ prefix: 'e:', start: key(after + 1) })) {
writer.write(this.enc.encode(frame)); // replay from the object's own storage
}
this.writers.add(writer);
writer.closed.catch(() => {}).finally(() => this.writers.delete(writer));
return new Response(readable, {
headers: { 'Content-Type': 'text/event-stream', 'Cache-Control': 'no-cache' },
});
}
private async publish(req: Request): Promise<Response> {
const body = await req.text();
const id = ++this.seq;
const frame = `id: ${id}\nevent: message\ndata: ${body}\n\n`;
await this.state.storage.put({ seq: id, [key(id)]: frame }); // durable before broadcast
this.broadcast(frame);
await this.trim(id);
return new Response(null, { status: 204 });
}
private broadcast(frame: string) {
const bytes = this.enc.encode(frame);
for (const w of this.writers) {
w.write(bytes).catch(() => this.writers.delete(w)); // failed write: client gone
}
}
private async trim(id: number) {
if (id % 100 === 0) { // keep the last 1,000 events
const old = await this.state.storage.list({ prefix: 'e:', end: key(id - 1000) });
await this.state.storage.delete([...old.keys()]);
}
}
}
const key = (n: number) => `e:${String(n).padStart(12, '0')}`; // lexicographic = numeric order
Because the object processes one request at a time, sequence numbers are allocated without races, and every stream sees events in the same order. Storing each event before broadcasting makes replay complete: a client that reconnects receives everything after its id from storage.
Step 3 — Send heartbeats from the object Permalink to this section
Use an alarm or a timer inside the object to write a comment to all writers every 15–20 seconds while any stream is open. This keeps intermediaries from closing idle streams and exposes clients that have gone, whose writes then fail and are removed.
Clients that stop reading are the one thing the object must not let accumulate. A TransformStream buffers written chunks until the readable side is consumed, so a stalled client’s writer keeps growing. Check writer.desiredSize: when it falls below zero, the client is behind; when it stays negative past a threshold, close that writer so the client reconnects and replays from storage.
Step 4 — Shard hot channels Permalink to this section
One object is single-threaded. A channel with tens of thousands of subscribers or a very high publish rate can exceed what one instance should handle. Split it: a coordinator object sequences publishes and forwards them to N fan-out objects, each holding a slice of the subscribers (room:42:shard:0 … :N-1, chosen by a hash of the client id).
Validation & Monitoring Permalink to this section
# Two clients, one publish: both receive id n, in the same order.
curl -sN https://edge.example.com/rooms/42/events > a.log &
curl -sN https://edge.example.com/rooms/42/events > b.log &
curl -s -X POST -d '{"text":"hello"}' https://edge.example.com/rooms/42/publish
sleep 1; grep '^id:' a.log b.log
# Replay: reconnect with an older id.
curl -sN -H 'Last-Event-ID: 50' https://edge.example.com/rooms/42/events | grep -m5 '^id:'
Use Workers analytics and object-level logging to track open writers per object, publishes per second, broadcast write failures and storage size. A growing writer count with no matching client traffic indicates writers not being removed on failure.
Production Checklist Permalink to this section
Frequently Asked Questions Permalink to this section
Can Durable Objects hibernate while SSE streams are open?
The hibernation API applies to WebSockets. An object holding open SSE responses stays active, so idle streams incur active duration. If that cost matters for mostly idle rooms, consider WebSockets with hibernation for those rooms.
How many subscribers can one object hold?
Thousands comfortably, depending on publish rate and payload size, since broadcast cost is one write per subscriber on a single thread. Beyond that, shard subscribers across several fan-out objects.
Is ordering guaranteed across channels?
Only within a channel, because each channel is one object. If a feature needs a single order across channels, route those publishes through one sequencing object.
Where should authentication happen?
In the Worker, before forwarding to the object. Pass the authenticated identity to the object in a header the Worker sets, so the object can filter or authorise per event if needed.