Consuming SSE in Node.js Permalink to this section

Part of Non-Browser SSE Clients, under Frontend Consumption & Client Patterns.

Node.js services consume Server-Sent Events constantly: following another team’s event feed, relaying a language model’s token stream, tailing a deployment, bridging a partner’s webhook-alternative stream into a queue. Node has everything needed — streaming fetch since version 18, and mature libraries on top — but the defaults of a service are not the defaults of a browser tab. This guide covers the two good approaches, authentication, what to do when events arrive faster than the service can process them, and shutting down without losing position.

Symptom & Developer Intent Permalink to this section

  • The consumer works for a while, then silently stops receiving events while the process keeps running.
  • Adding an Authorization header to the eventsource package seemed impossible.
  • Memory grows during upstream bursts because events are processed asynchronously without limit.
  • A redeploy of the consumer replays everything from the start, or skips everything published during the restart.
  • The consumer hammers the upstream with reconnects during its outage.

The intent is a consumer that stays connected, authenticates correctly, processes events at its own pace, resumes exactly across restarts and reconnects politely.

Root Cause Analysis Permalink to this section

Most failures come from treating the stream like a request. Silent death happens when a TCP connection dies without a FIN and nothing notices: no read timeout, no watchdog. Memory growth happens because an event handler that starts async work and returns immediately lets work pile up without bound. Position loss across restarts happens because the last event id lives only in memory.

A well-behaved Node.js SSE consumer Flow from the upstream stream through a parser, a bounded work queue with backpressure, handlers, and a durable checkpoint of the last processed id. A well-behaved Node.js SSE consumer Upstream stream text/event-stream bytes Parser eventsource / parser events Bounded queue pause when full pull Handlers at own pace persist Checkpoint last processed id
The checkpoint records the last id that was processed, not merely received. That is the position to resume from after a restart.

Step-by-Step Resolution Permalink to this section

Step 1 — Choose a client Permalink to this section

For most consumers, the eventsource package is the simplest: it implements the standard EventSource API, including reconnection and Last-Event-ID, on top of fetch, and accepts a custom fetch for headers:

import { EventSource } from 'eventsource';

const es = new EventSource(process.env.FEED_URL, {
  fetch: async (input, init) => fetch(input, {
    ...init,
    headers: { ...init.headers, Authorization: `Bearer ${await tokens.get()}` },   // fresh per attempt
  }),
});

When you need a POST body, precise status handling or full control over reconnection, use fetch with eventsource-parser and the reconnect loop from the fetch-based clients topic.

Step 2 — Add a silence watchdog Permalink to this section

A dead TCP connection can leave a stream “open” indefinitely. The server sends heartbeats; if nothing at all arrives for three heartbeat intervals, force a reconnect:

let lastByte = Date.now();
const touch = () => { lastByte = Date.now(); };
es.addEventListener('open', touch);
es.addEventListener('message', touch);
for (const t of EVENT_TYPES) es.addEventListener(t, touch);

setInterval(() => {
  if (Date.now() - lastByte > 3 * HEARTBEAT_MS) {
    log.warn('stream silent, reconnecting');
    reopen();                                 // close and create a new EventSource from the checkpoint
  }
}, HEARTBEAT_MS);

Comment-line heartbeats are invisible to EventSource listeners; if the upstream only sends comments, use the fetch-based approach, where every byte is observable, or ask for a named heartbeat event.

Step 3 — Process through a bounded queue Permalink to this section

import PQueue from 'p-queue';
const queue = new PQueue({ concurrency: 8 });

es.addEventListener('order.updated', async (e) => {
  if (queue.size > 1000) await queue.onSizeLessThan(500);   // simple backpressure
  queue.add(async () => {
    await handleOrder(JSON.parse(e.data));
    await checkpoint.save(e.lastEventId);                   // record only after processing
  });
});

With EventSource, awaiting inside a listener does not pause reading from the socket, so this bounds work in progress but not memory for received events. For true backpressure — stop reading when the queue is full so TCP flow control slows the upstream — read the body yourself with fetch and only call reader.read() when the queue has room.

Consumer memory during a 5-minute upstream burst Bar chart comparing peak consumer memory during an upstream burst for unbounded async handlers, a bounded concurrency queue, and a fetch-based reader that pauses reading when the queue is full. Consumer memory during a 5-minute upstream burst Unbounded async handlers ~1.9 GB Bounded concurrency ~620 MB Pause reading when full ~140 MB peak resident memory during a burst of 50,000 events
Only pausing the read propagates backpressure to the upstream. Bounding concurrency alone still accumulates received events.

Step 4 — Checkpoint the last processed id durably Permalink to this section

const checkpoint = {
  async load() { return (await redis.get('feed:orders:cursor')) ?? ''; },
  async save(id) { if (id) await redis.set('feed:orders:cursor', id); },
};

async function open() {
  const after = await checkpoint.load();
  return new EventSource(`${process.env.FEED_URL}?after=${encodeURIComponent(after)}`, { fetch: authedFetch });
}

The query parameter covers the first connection after a restart, where EventSource has no Last-Event-ID yet; within a process, the library sends the header on reconnects. Handlers must be idempotent, since events processed but not yet checkpointed at a crash will be delivered again.

Step 5 — Shut down without losing position Permalink to this section

process.on('SIGTERM', async () => {
  es.close();                                  // stop receiving
  await queue.onIdle();                        // finish in-flight work
  await checkpoint.flush?.();                  // persist the last processed id
  process.exit(0);
});

Validation & Monitoring Permalink to this section

# Kill the upstream's network path briefly and watch the consumer reconnect and resume.
sudo iptables -A OUTPUT -p tcp -d feed.example.com --dport 443 -j DROP; sleep 90
sudo iptables -D OUTPUT -p tcp -d feed.example.com --dport 443 -j DROP
# Expect: watchdog fires within 3 heartbeats, reconnect resumes after the checkpoint.
A silently dead connection caught by the watchdog Timeline of three minutes showing heartbeats every fifteen seconds, a connection that dies silently at sixty seconds, the watchdog firing after three missed heartbeats, and the consumer reconnecting from its checkpoint. A silently dead connection caught by the watchdog Heartbeats Silence Reconnected every 15 s no bytes resumed from checkpoint 0 36 72 108 144 180 seconds path dies watchdog
Without the watchdog the process would wait indefinitely on a socket that will never deliver another byte.

The watchdog’s threshold is a trade-off: too short and brief network hiccups cause needless reconnects, each of which replays a little history; too long and a dead stream goes unnoticed for minutes. Three heartbeat intervals is a sound default, and the server’s heartbeat interval should be part of the stream’s documented contract so consumers can set it correctly.

Export: stream state (connected or not), seconds since last event, events processed per second, queue depth, reconnects by cause and checkpoint lag (upstream’s latest id versus the checkpoint, if the upstream exposes it). Alert on seconds-since-last-event and on checkpoint lag growing.

Production Checklist Permalink to this section

Frequently Asked Questions Permalink to this section

Which eventsource package version supports custom headers?

Recent major versions implement the standard API on top of fetch and accept a custom fetch function, which is where headers are added. Older versions had a non-standard headers option; check the version you depend on.

Why checkpoint after processing rather than on receipt?

Checkpointing on receipt loses events that were received but not processed when the process crashed. Checkpointing after processing may repeat a few, which idempotent handlers absorb.

Can Node's built-in fetch stream SSE?

Yes. response.body is a web ReadableStream in Node 18 and later, and works with TextDecoderStream and any spec-compliant parser.

Should one process follow many upstream streams?

Yes, if they are independent: Node's event loop holds many idle streams cheaply. Give each its own checkpoint key, watchdog and queue, so a burst or failure on one feed never delays the others.

How do I relay an upstream stream to my own clients?

Parse upstream events, transform or filter them, and write them to your own SSE responses with your own ids. Propagate cancellation upstream when your last client for that stream disconnects.