feat(sse): per-client live subscription filter (#86)

* feat(sse): per-client live subscription filter

Lets a connected client narrow its SSE stream to a single session
without forcing an EventSource reconnect. With many open sessions
(N tabs in the UI, all generating output), this cuts terminal-event
SSE traffic by roughly Nx — we only send the actively-rendered
session's bytes instead of all of them.

The existing ?sessions= query filter only worked at connect time;
narrowing or widening it required tearing down the EventSource and
losing in-flight messages. That was acceptable when filters were set
once at page load, but the UI now flips active sessions on every
tab switch.

How it works
============

- Client generates a stable per-page UUID (`_clientId`) once at
  CodemanApp construction and includes it on the SSE URL:
    GET /api/events?clientId=<uuid>&sessions=<active-id>
- Server records a `clientId -> reply` mapping in addition to the
  existing `reply -> sessionFilter` map.
- New endpoint:
    POST /api/events/subscribe { clientId, sessions: string[] | null }
  updates the in-memory filter for the matching reply. 204 on success,
  404 if the client isn't known yet (race on first selectSession after
  reconnect — the next reconnect carries the filter via the URL).
- On every selectSession the client fires a fire-and-forget POST. No
  reconnect, no re-init, no replay buffer needed.

Behavioural change to broadcast()
=================================

The per-event session filter is removed from `broadcast()`. Previously
that path filtered lifecycle/metadata events (`session:created`,
`session:updated`, `ralph:*`, `hook:*`) by extracting a `sessionId` from
the payload. With per-client narrow filters, that meant a client
subscribed to session A would never see session:created for B and the
sidebar would silently de-sync.

The new contract:
- **Lifecycle/metadata events** (low-volume, UI-correctness critical)
  broadcast to all clients regardless of filter.
- **Terminal events** (high-volume, the actual reason for filtering)
  apply the filter in `flushSessionTerminalBatch` (already there;
  unchanged).

`extractSessionId()` was only used by the old broadcast() filter and
has been removed.

Files
=====

- src/web/sse-stream-manager.ts (+34/-29): add `sseClientsById`,
  optional `clientId` arg to addClient/removeClient cleanup, new
  `updateClientFilter()`, and the broadcast() change above.
- src/web/server.ts (+22/-3): parse `clientId` on /api/events, pass
  to `addClient`, register POST /api/events/subscribe handler.
- src/web/public/app.js (+41/-1): generate `_clientId`, build the
  EventSource URL with both clientId + active session, add
  `_updateSseSubscription()`, call it on selectSession.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>

* test(sse): update operation-lightspeed to match broadcast-all contract

Lifecycle events (session:*, case:*) now reach every connected SSE
client; only session:terminal is gated by the per-client filter.
Updates the four assertions in operation-lightspeed.test.ts that
encoded the old "filter applies to all events" contract.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
Co-authored-by: arkon <arkon.85@hotmail.com>
This commit is contained in:
aakhter
2026-05-17 05:56:53 +02:00
committed by GitHub
co-authored by Claude Opus 4.7 arkon
parent e87b03b6c2
commit 98966def03
4 changed files with 127 additions and 58 deletions
+37 -26
View File
@@ -48,6 +48,8 @@ export class SseStreamManager {
* or `null` meaning "receive all events" (backwards-compatible default).
*/
private sseClients: Map<FastifyReply, Set<string> | null> = new Map();
/** Optional client-supplied IDs → reply, for live filter updates without reconnecting */
private sseClientsById: Map<string, FastifyReply> = new Map();
/** SSE clients connecting from non-localhost (i.e. through tunnel) */
private remoteSseClients: Set<FastifyReply> = new Set();
/** Clients with backpressure — skip writes until 'drain' fires */
@@ -103,17 +105,43 @@ export class SseStreamManager {
this._isTunnelActive = active;
}
addClient(reply: FastifyReply, sessionFilter: Set<string> | null, isRemote: boolean): void {
addClient(reply: FastifyReply, sessionFilter: Set<string> | null, isRemote: boolean, clientId?: string): void {
this.sseClients.set(reply, sessionFilter);
if (isRemote) {
this.remoteSseClients.add(reply);
}
if (clientId) {
// If a previous reply registered the same id (reconnect), drop the old one.
const prev = this.sseClientsById.get(clientId);
if (prev && prev !== reply) {
this.sseClients.delete(prev);
this.remoteSseClients.delete(prev);
this.backpressuredClients.delete(prev);
}
this.sseClientsById.set(clientId, reply);
}
}
removeClient(reply: FastifyReply): void {
this.sseClients.delete(reply);
this.remoteSseClients.delete(reply);
this.backpressuredClients.delete(reply);
// Clear any clientId mappings pointing at this reply
for (const [id, r] of this.sseClientsById) {
if (r === reply) this.sseClientsById.delete(id);
}
}
/**
* Update an existing client's session subscription filter without forcing
* an SSE reconnect. Returns true if the client was found and updated.
*/
updateClientFilter(clientId: string, sessions: string[] | null): boolean {
const reply = this.sseClientsById.get(clientId);
if (!reply || !this.sseClients.has(reply)) return false;
const filter = sessions && sessions.length > 0 ? new Set(sessions) : null;
this.sseClients.set(reply, filter);
return true;
}
/** Send a single SSE event to a specific client. */
@@ -188,35 +216,18 @@ export class SseStreamManager {
console.error(`[Server] Failed to serialize SSE event "${event}":`, err);
return;
}
// Extract sessionId from event data for subscription filtering.
const eventSessionId = this.extractSessionId(event, data);
for (const [client, filter] of this.sseClients) {
// No filter (null) = receive everything. Otherwise, skip if event is
// session-scoped and the session isn't in the client's subscription set.
if (filter && eventSessionId && !filter.has(eventSessionId)) continue;
// Subscription filtering is intentionally NOT applied here. The
// `?sessions=` filter is intended to suppress only the high-volume
// terminal stream — lifecycle/metadata events (session:created,
// session:updated, ralph:*, hook:*, etc.) are needed for correct UI
// state across all sessions even when the client subscribes to a single
// active session's terminal output. Terminal events bypass this method
// entirely (see flushSessionTerminalBatch — it applies the filter).
for (const [client] of this.sseClients) {
this.sendSSEPreformatted(client, message);
}
}
/**
* Extract the session ID from an event's data payload for subscription filtering.
* Returns the sessionId string if the event is session-scoped, or null for global events.
*/
private extractSessionId(event: string, data: unknown): string | null {
if (data == null || typeof data !== 'object') return null;
const record = data as Record<string, unknown>;
// Most session-scoped events use `sessionId`
if (typeof record.sessionId === 'string') return record.sessionId;
// Session lifecycle events (session:*) use `id` from the session state object
if (typeof record.id === 'string' && event.startsWith('session:')) return record.id;
// No session ID found — treat as global event (sent to all clients)
return null;
}
// ========== Terminal Data Batching ==========
// Batch terminal data for better performance (60fps)