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 Durable Object per channel Worker requests for subscribing and publishing are routed by channel name to a single Durable Object, which holds a writer for each open stream and broadcasts to all of them. One Durable Object per channel Workers (edge) route by room Durable Object room:42 stub.fetch Client A SSE writer Client B SSE writer Client C SSE writer writer.write
The object is both the rendezvous and the sequencer. Every client of room 42 reads from, and every publisher writes to, the same instance.

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.

Publish and fan-out inside the object Sequence diagram of a publisher request reaching the Durable Object, the object assigning sequence 57, storing the frame, and writing it to each open stream's writer. Publish and fan-out inside the object Publisher Room object Storage Streams POST /publish seq = 57 put e:57, seq write frame 57 to each writer 204
Storage write, then broadcast. A client that reconnects a moment later finds event 57 in storage even if its stream missed the live write.

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:'
When a Durable Object per channel fits Matrix comparing Durable Objects, a Redis-backed server tier and a managed pub/sub service on ordering, replay, operational effort and cost for idle connections. When a Durable Object per channel fits Approach Per-channel order Replay Ops effort Idle-stream cost Durable Object per channel built in object storage none active duration Redis + server tier by design Streams run servers low Managed pub/sub varies varies low per message strength neutral weakness
Durable Objects give ordering and replay with no servers to run. The trade-off is duration billing while streams are open.

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.