feat(sse): heal a stalled SSE stream with a heartbeat + client watchdog

An EventSource that stops delivering does not always error. A proxy that
idle-closed the connection, a laptop resumed from sleep, a tailnet reconnect:
`onerror` never fires, the header dot stays green, and every SSE-driven surface
(tab status dots, sessions created on another device, renames) freezes until the
user reloads. Nothing on the client tracked stream liveness at all.

The server already wrote a keepalive every 15s, but as an SSE `:keepalive`
COMMENT, and comments are invisible to `EventSource` by spec, so there was
nothing a client could observe.

Server:
- `sse:heartbeat` under a new Transport category in the event registry
  (155 constants now, both counts updated).
- `cleanupDeadClients()` writes that named frame (`{"t":<epoch ms>}`) instead of
  the comment. Interval, tunnel padding and dead-socket eviction are unchanged.
  The write stays per-client rather than going through `broadcast()`: the frame
  carries no session data, so it needs no multi-user owner routing.

Client:
- `computeSseStale()` in constants.js, a pure policy beside
  `computeConnectionLossUi`. Stale only when the transport believes it is
  `connected`, the device is online, and no frame has arrived for 45s (three
  missed heartbeats). The `connected`-only guard is also the loop breaker: a
  forced reconnect leaves that state immediately, so the watchdog cannot re-fire
  while one is in flight.
- The liveness stamp is applied inside `addListener` itself, so the
  `_SSE_HANDLER_MAP` wrappers and the directly-registered listeners all feed it
  from one place instead of three that can drift. The heartbeat's own listener
  is a no-op that exists only to be registered, since `EventSource` drops named
  events nobody listens for.
- A 5s watchdog forces `connectSSE()` when the policy says stale, and is cleared
  at the top of `connectSSE()` and nowhere else (its only teardown path).
  Recovery needs no new sync path: the reconnect re-runs `handleInit`, which
  already rebuilds from the server. `visibilitychange` -> visible checks too,
  riding the existing listener, since a background tab's timers are throttled
  and a wake is exactly when a stream comes back zombie.
- The forced reconnect logs one diagnostic line: if a middlebox ever strips or
  delays heartbeats, the failure mode is "silently reconnects every 45s", which
  is undebuggable from a field report without it.

Tests: `test/sse-staleness.test.ts` (node VM over constants.js, threshold
boundaries and every not-stale guard) and `test/sse-heartbeat.test.ts` (drives
`cleanupDeadClients()` with fake replies: named frame not a comment, parseable
payload, padding only with a tunnel, dead clients still evicted).

Verified end to end on an isolated instance: with the stream closed client-side
(no `onerror`), a rename sticks, an out-of-band session stays invisible, then
the watchdog reconnects on its own and it appears without a reload.

Event names are part of the stable API contract, so this is a MINOR bump.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
Codeman maintainer
2026-08-13 17:28:35 +02:00
parent d19895651d
commit c790166564
8 changed files with 434 additions and 11 deletions
+86 -3
View File
@@ -684,6 +684,17 @@ class CodemanApp {
this.maxReconnectAttempts = 10;
this.isOnline = navigator.onLine;
// SSE staleness watchdog. An EventSource that stops delivering does not
// always error (a proxy that idle-closed it, a resumed laptop), so
// `onerror` never fires and every SSE-driven surface freezes silently.
// The server heartbeats every 15s; going quiet for three of them means the
// stream is a zombie and has to be rebuilt. The decision is pure
// (computeSseStale in constants.js); these are its inputs. The threshold
// is an instance field so a browser test can shrink it.
this._sseLastMessageAt = 0;
this._sseStaleTimeoutMs = window.CodemanSseStale?.TIMEOUT_MS ?? 45000;
this._sseStaleWatchdog = null;
// Connection-loss UI (banner + full-screen overlay). The decision itself is
// pure and lives in constants.js (computeConnectionLossUi); these are just
// its inputs. `_connDownSince` is the timestamp the transport LEFT the
@@ -719,6 +730,11 @@ class CodemanApp {
window.addEventListener('pagehide', () => this._persistReliableNow());
document.addEventListener('visibilitychange', () => {
if (document.visibilityState === 'hidden') this._persistReliableNow();
// A background tab's timers are throttled, so the 5s watchdog may not
// have run for minutes, and a wake/unlock is exactly when a stream
// comes back zombie. Checking here is what makes recovery feel instant
// instead of up to a full timeout late.
else this._checkSseStale();
});
// Local echo overlay — DOM overlay positioned at the visible ❯ prompt
@@ -1421,6 +1437,14 @@ class CodemanApp {
// Clear any pending reconnect timeout to prevent duplicate connections
this._clearTimer('sseReconnectTimeout');
// Same discipline for the staleness watchdog: connectSSE() runs on every
// reconnect and is the only teardown path this page-lifetime interval has,
// so clearing it anywhere else (or not at all) stacks intervals.
if (this._sseStaleWatchdog) {
clearInterval(this._sseStaleWatchdog);
this._sseStaleWatchdog = null;
}
// Clean up existing SSE listeners before creating new connection (prevents listener accumulation)
if (this._sseListenerCleanup) {
this._sseListenerCleanup();
@@ -1448,11 +1472,20 @@ class CodemanApp {
if (this.activeSessionId) _sseParams.set('sessions', this.activeSessionId);
this.eventSource = new EventSource(`/api/events?${_sseParams.toString()}`);
// Store all event listeners for cleanup on reconnect
// Store all event listeners for cleanup on reconnect.
//
// Every handler is wrapped so ANY frame that arrives stamps the liveness
// clock the staleness watchdog reads. Doing it here (rather than at the
// three separate registration sites below) is what keeps a future
// addListener() call from silently opting out of it.
const listeners = [];
const addListener = (event, handler) => {
this.eventSource.addEventListener(event, handler);
listeners.push({ event, handler });
const stamped = (e) => {
this._sseLastMessageAt = Date.now();
handler(e);
};
this.eventSource.addEventListener(event, stamped);
listeners.push({ event, handler: stamped });
};
// Create cleanup function to remove all listeners
@@ -1467,6 +1500,10 @@ class CodemanApp {
this.eventSource.onopen = () => {
this.reconnectAttempts = 0;
// Start the liveness clock here, not at the first frame: the watchdog
// only ever fires while the status is 'connected', and this is the
// moment that becomes true.
this._sseLastMessageAt = Date.now();
this.setConnectionStatus('connected');
};
this.eventSource.onerror = () => {
@@ -1614,6 +1651,52 @@ class CodemanApp {
}
this._onSessionListMaybeChanged();
});
// Liveness heartbeat. The handler is deliberately empty: the whole point
// is the stamp inherited from addListener's wrapper. It still has to be
// REGISTERED: EventSource only dispatches named events that have a
// listener, so without this the frame arrives on the wire and is dropped
// before it can prove the stream is alive.
addListener(SSE_EVENTS.HEARTBEAT, () => {});
// Watchdog: a stream that goes quiet without erroring is invisible to
// onerror, so poll the pure staleness policy and rebuild the connection
// ourselves. 5s granularity against a 45s threshold: cheap, and it keeps
// the worst-case detection lag well under a heartbeat interval.
this._sseStaleWatchdog = setInterval(() => this._checkSseStale(), 5000);
}
/**
* Force a reconnect if the SSE stream has gone quiet while still claiming to
* be connected. Called by the 5s watchdog and on tab-visible.
*
* Recovery needs no new sync path: the reconnect re-runs `handleInit`, which
* already calls `_resetAllAppState()` and rebuilds everything from the
* server. The connection-loss UI needs nothing either: `connectSSE()` sets
* status 'connecting' (reconnectAttempts was zeroed by onopen), and the 2.5s
* grace in computeConnectionLossUi means a stream that heals in 200ms shows
* nothing at all.
*/
_checkSseStale() {
const policy = window.CodemanSseStale;
if (!policy) return;
const now = Date.now();
const stale = policy.compute({
lastMessageAt: this._sseLastMessageAt,
now,
status: this._connectionStatus,
isOnline: this.isOnline,
timeoutMs: this._sseStaleTimeoutMs,
});
if (!stale) return;
// If a middlebox ever strips or delays heartbeats, the failure mode is
// "silently reconnects every 45s", and a field report of that would be
// undebuggable without this line.
console.log(
`[SSE] stream stale: no frame for ${now - this._sseLastMessageAt}ms ` +
`(threshold ${this._sseStaleTimeoutMs}ms), forcing reconnect`
);
this.connectSSE();
}
// ═══════════════════════════════════════════════════════════════
+42
View File
@@ -379,6 +379,41 @@ function computeConnectionLossUi(input) {
};
}
// SSE staleness policy: is this stream a zombie?
//
// An EventSource that stops delivering does not always error. A proxy that
// idle-closed the connection, a laptop resumed from sleep, a tailnet
// reconnect: `onerror` never fires, the header dot stays green, and every
// SSE-driven surface (tab status dots, sessions created on another device,
// renames) freezes until the user reloads. The server writes a
// `sse:heartbeat` frame every 15s, so silence longer than three of them means
// the stream is dead even though the transport still claims otherwise.
//
// Stale ONLY when the transport believes it is 'connected': the other states
// already have the reconnect/backoff machinery running, and re-firing on top
// of them would stack reconnects. That guard is also the loop breaker: a
// forced reconnect leaves 'connected' immediately, so the watchdog cannot
// fire again while one is in flight. `navigator.onLine === false` is not
// staleness either; there is nothing to reconnect to yet.
//
// Pure: no DOM, no timers, no side effects. `now` is passed in.
const SSE_STALE_TIMEOUT_MS = 45000; // three missed 15s heartbeats
function computeSseStale(input) {
const {
lastMessageAt = null,
now = 0,
status = 'connected',
isOnline = true,
timeoutMs = SSE_STALE_TIMEOUT_MS,
} = input || {};
if (!isOnline || status !== 'connected') return false;
// No frame has ever arrived: `init` lands on connect, so this is a stream
// that has not opened yet rather than one that went quiet.
if (typeof lastMessageAt !== 'number' || !(lastMessageAt > 0)) return false;
return now - lastMessageAt >= timeoutMs;
}
if (typeof window !== 'undefined') {
window.WEBGL_FALLBACK = WEBGL_FALLBACK;
window.evaluateWebGLLongTaskTrip = evaluateWebGLLongTaskTrip;
@@ -401,6 +436,10 @@ if (typeof window !== 'undefined') {
compute: computeConnectionLossUi,
GRACE_MS: CONNECTION_LOSS_GRACE_MS,
};
window.CodemanSseStale = {
compute: computeSseStale,
TIMEOUT_MS: SSE_STALE_TIMEOUT_MS,
};
}
// Scheduler API — prioritize terminal writes over background UI updates.
@@ -514,6 +553,9 @@ const SSE_EVENTS = {
// Core
INIT: 'init',
// Transport
HEARTBEAT: 'sse:heartbeat',
// Session lifecycle
SESSION_CREATED: 'session:created',
SESSION_UPDATED: 'session:updated',
+21 -1
View File
@@ -5,8 +5,9 @@
* and referenced by the frontend (`SSE_EVENTS` in `constants.js`).
* Both files MUST be kept in sync.
*
* 154 event constants organized by category:
* 155 event constants organized by category:
* - **Core** (1): init
* - **Transport** (1): sse:heartbeat
* - **Session lifecycle** (23): created, updated, deleted, terminal, idle, working, ...
* - **Session: Ralph** (6): ralphLoopUpdate, todoUpdate, completionDetected, ...
* - **Session: Bash tools** (3): bashToolStart, bashToolEnd, bashToolsUpdate
@@ -52,6 +53,22 @@
/** Sent to each SSE client on initial connection with full app state. */
export const Init = 'init' as const;
// ─── Transport ───────────────────────────────────────────────────────────────
/**
* Liveness frame written to every SSE client every `SSE_HEARTBEAT_INTERVAL`.
* Payload: `{ t: <epoch ms> }`.
*
* Carries no application data; its only job is to be *observable*. This was a
* `:keepalive` SSE **comment**, and comments are invisible to `EventSource` by
* spec, so a stream that stopped delivering without erroring (a proxy that
* idle-closed it, a laptop resumed from sleep, a tailnet reconnect) was
* undetectable to the client: `onerror` never fires and the UI freezes until a
* reload. A named event reaches a listener, which is what lets the client's
* staleness watchdog notice the silence and force a reconnect.
*/
export const Heartbeat = 'sse:heartbeat' as const;
// ─── Session Lifecycle ───────────────────────────────────────────────────────
/** New session spawned. */
@@ -443,6 +460,9 @@ export const SseEvent = {
// Core
Init,
// Transport
Heartbeat,
// Session lifecycle
SessionCreated,
SessionUpdated,
+12 -6
View File
@@ -470,12 +470,20 @@ export class SseStreamManager {
// ========== Client Health ==========
/**
* Clean up dead SSE clients and send keep-alive comments.
* Clean up dead SSE clients and send the liveness heartbeat.
* Keep-alive prevents proxy/load-balancer timeouts on idle connections.
* Dead client cleanup prevents memory leaks from abruptly terminated connections.
*
* The heartbeat is a NAMED event, not the `:keepalive` comment it used to be:
* comments are invisible to `EventSource` by spec, so a stream that stopped
* delivering without erroring was undetectable to the client (see
* `SseEvent.Heartbeat`). Written per-client rather than through `broadcast()`
* deliberately: the frame carries no session data, so it needs no owner
* routing, and this loop is already walking every client to check its socket.
*/
cleanupDeadClients(): void {
const deadClients: FastifyReply[] = [];
const heartbeat = `event: ${SseEvent.Heartbeat}\ndata: ${JSON.stringify({ t: Date.now() })}\n\n`;
for (const [client] of this.sseClients) {
try {
@@ -484,11 +492,9 @@ export class SseStreamManager {
if (!socket || socket.destroyed || !socket.writable) {
deadClients.push(client);
} else {
// Send SSE comment as keep-alive. Only add padding when tunnel is
// active — it flushes Cloudflare proxy buffers but wastes bandwidth
// for direct/Tailscale connections.
const ka = this._isTunnelActive ? ':keepalive\n' + SSE_PADDING : ':keepalive\n\n';
client.raw.write(ka);
// Only add padding when tunnel is active: it flushes Cloudflare
// proxy buffers but wastes bandwidth for direct/Tailscale connections.
client.raw.write(this._isTunnelActive ? heartbeat + SSE_PADDING : heartbeat);
}
} catch {
// Error accessing socket means client is dead