From 62c6f866d564e3bc73bf113a6bf49048a54b4d3a Mon Sep 17 00:00:00 2001 From: Danilo Alonso Date: Mon, 5 Oct 2026 14:03:40 -0400 Subject: [PATCH] fix: end the stream when replay fails instead of leaking An error status stops EventSource for good; an ended stream makes it retry and runs the normal cleanup. Null-character ids are rejected before they reach the replay buffer. Closes #8 --- API.md | 8 ++- src/event-buffer.ts | 12 +++- src/replayer.ts | 6 ++ src/sse.ts | 12 +++- src/subscription.ts | 5 ++ test/event-buffer.test.ts | 7 +++ test/replayer.test.ts | 13 ++++ test/sse.test.ts | 123 ++++++++++++++++++++++++++++++++++++++ 8 files changed, 177 insertions(+), 9 deletions(-) diff --git a/API.md b/API.md index 1f6ea5c..699cca2 100644 --- a/API.md +++ b/API.md @@ -90,8 +90,8 @@ server.sse.subscription('/chat/{room}', { | `filter` | `(path, message, opts) => boolean \| { override } \| Promise<...>` | Per-session delivery filter | | `onSubscribe` | `(session, path, params) => void \| Promise` | Fires before SSE headers are sent. Throwing a Boom error returns that HTTP error to the client. Closing the session responds with `204 No Content`. | | `onUnsubscribe` | `(session, path, params) => void` | Fires on client disconnect | -| `onReconnect` | `(session, path, params) => void \| Promise` | Fires when `Last-Event-ID` is present, after replay. Runs without delaying the response. Errors close the session gracefully. | -| `replay` | `Replayer` | Replay provider for automatic reconnection replay | +| `onReconnect` | `(session, path, params) => void \| Promise` | Fires when `Last-Event-ID` is present, after replay. Skipped when replay fails. Runs without delaying the response. Errors close the session gracefully. | +| `replay` | `Replayer` | Replay provider for automatic reconnection replay. If replay fails (`replay()` throws or returns an entry with an invalid `id`), the error is reported with `request.log()` (tags `sse`, `replay`, `error`), the stream ends, and the client reconnects after `retry`. A replayer that keeps returning an invalid entry fails every reconnect until that entry is gone; the built-in replayers' `record()` throws on one. | | `maxSessions` | `number` | Maximum concurrent sessions for this subscription. Excess connections receive a 503 response. | | `maxDuration` | `number` | Maximum connection lifetime in ms. Sessions are closed after this duration (with ±10% jitter to prevent thundering herd reconnections). A `: session expired` comment is sent before closing. | @@ -119,7 +119,7 @@ console.log(`Delivered to ${delivered} sessions`); - `'pattern'` (default) — delivers to all sessions on a matching subscription pattern (e.g. `/chat/{room}`) - `'literal'` — only delivers to sessions whose actual connected path equals `path` exactly. Useful for parameterized subscriptions where you want to target `/chat/general` but not `/chat/random`. -**Note:** Only events published with an explicit `id` are recorded by the replayer. Events without an `id` are delivered but not stored for replay. +**Note:** Only events published with an explicit `id` are recorded by the replayer. Events without an `id` are delivered but not stored for replay. `publish()` rejects an `id` containing a null character before anything is delivered or recorded. ### `server.sse.broadcast(data, opts?)` @@ -132,6 +132,8 @@ const count = await server.sse.broadcast( ); ``` +`broadcast()` rejects an `id` containing a null character before anything is delivered. + ### `server.sse.eachSession(fn, opts?)` Iterates over connected sessions. Optionally filter by subscription pattern. diff --git a/src/event-buffer.ts b/src/event-buffer.ts index 8028b9f..80ded72 100644 --- a/src/event-buffer.ts +++ b/src/event-buffer.ts @@ -1,3 +1,9 @@ +export const assertEventId = (id?: string): void => { + if (id?.includes('\0')) { + throw new Error('Event ID must not contain null characters'); + } +}; + export class EventBuffer { /** @internal */ #buffer = ''; @@ -20,9 +26,7 @@ export class EventBuffer { } id(id: string): this { - if (id.includes('\0')) { - throw new Error('Event ID must not contain null characters'); - } + assertEventId(id); this.#buffer += `id: ${id.replace(/[\r\n]/g, '')}\n`; @@ -60,6 +64,8 @@ export class EventBuffer { } push(data: unknown, event?: string, id?: string): this { + assertEventId(id); + if (event) { this.event(event); } diff --git a/src/replayer.ts b/src/replayer.ts index 1c23e37..5e1ec10 100644 --- a/src/replayer.ts +++ b/src/replayer.ts @@ -1,5 +1,7 @@ import Joi from 'joi'; +import { assertEventId } from './event-buffer.js'; + export interface ReplayEntry { data: unknown; event?: string; @@ -40,6 +42,8 @@ export class FiniteReplayer implements Replayer { } record(entry: ReplayEntry): void { + assertEventId(entry.id); + const stored: ReplayEntry = { data: entry.data, event: entry.event, @@ -90,6 +94,8 @@ export class ValidReplayer implements Replayer { } record(entry: ReplayEntry): void { + assertEventId(entry.id); + const stored: TimedEntry = { data: entry.data, event: entry.event, diff --git a/src/sse.ts b/src/sse.ts index 3432815..c04303f 100644 --- a/src/sse.ts +++ b/src/sse.ts @@ -295,10 +295,16 @@ export const SsePlugin: NamedPlugin = { const replayer = subConfig.replay; if (session.lastEventId && replayer) { - const entries = replayer.replay(session.lastEventId); + try { + for (const entry of replayer.replay(session.lastEventId)) { + session.push(entry.data, entry.event, entry.id); + } + } catch (err) { + // An error response would stop EventSource for good; an ended stream makes it retry. + request.log(['sse', 'replay', 'error'], err instanceof Error ? err : String(err)); + session.close(); - for (const entry of entries) { - session.push(entry.data, entry.event, entry.id); + return session.respond(h); } } diff --git a/src/subscription.ts b/src/subscription.ts index 7e633b3..fb31c3a 100644 --- a/src/subscription.ts +++ b/src/subscription.ts @@ -1,5 +1,6 @@ import type { Request, RouteOptions } from '@hapi/hapi'; +import { assertEventId } from './event-buffer.js'; import { Session } from './session.js'; import type { Replayer } from './replayer.js'; @@ -145,6 +146,8 @@ export class SubscriptionRegistry { data: T, opts?: { event?: string; id?: string; internal?: unknown; matchMode?: 'pattern' | 'literal' }, ): Promise { + assertEventId(opts?.id); + const matched = this.matchPath(path); if (!matched) { @@ -204,6 +207,8 @@ export class SubscriptionRegistry { } async broadcast(data: unknown, opts?: { event?: string; id?: string }): Promise { + assertEventId(opts?.id); + let delivered = 0; for (const sub of this.#subscriptions.values()) { diff --git a/test/event-buffer.test.ts b/test/event-buffer.test.ts index 5d8875a..7e54a67 100644 --- a/test/event-buffer.test.ts +++ b/test/event-buffer.test.ts @@ -188,6 +188,13 @@ describe.concurrent('EventBuffer', () => { expect(buf.read()).toBe('event: evt\ndata: data\n\n'); }); + it('push() with an invalid id leaves no partial event behind to relabel the next one', () => { + const buf = new EventBuffer(); + + expect(() => buf.push('data', 'alert', 'a\u0000b')).toThrow(/null characters/); + expect(buf.read()).toBe(''); + }); + it('event() strips newlines to prevent field injection', () => { const buf = new EventBuffer(); diff --git a/test/replayer.test.ts b/test/replayer.test.ts index 0a83e84..0cf069f 100644 --- a/test/replayer.test.ts +++ b/test/replayer.test.ts @@ -4,6 +4,12 @@ import { expect, describe, it } from 'vitest'; import { FiniteReplayer, ValidReplayer } from '../src/replayer.js'; describe.concurrent('FiniteReplayer', () => { + it('record() throws on an id containing a null character', () => { + const replayer = new FiniteReplayer({ size: 5 }); + + expect(() => replayer.record({ data: 'a', id: 'a\u0000b' })).toThrow(/null characters/); + }); + it('records and replays entries after lastEventId', () => { const replayer = new FiniteReplayer({ size: 10 }); @@ -110,6 +116,13 @@ describe.concurrent('FiniteReplayer', () => { }); describe.concurrent('ValidReplayer', () => { + it('record() throws on an id containing a null character', ({ onTestFinished }) => { + const replayer = new ValidReplayer({ ttl: 60_000 }); + onTestFinished(() => replayer.stop()); + + expect(() => replayer.record({ data: 'a', id: 'a\u0000b' })).toThrow(/null characters/); + }); + it('records and replays entries after lastEventId', ({ onTestFinished }) => { const replayer = new ValidReplayer({ ttl: 60_000 }); onTestFinished(() => replayer.stop()); diff --git a/test/sse.test.ts b/test/sse.test.ts index f0a7c45..236d312 100644 --- a/test/sse.test.ts +++ b/test/sse.test.ts @@ -9,6 +9,7 @@ import Boom from '@hapi/boom'; import { SsePlugin } from '../src/sse.js'; import { FiniteReplayer, ValidReplayer } from '../src/replayer.js'; +import type { Session } from '../src/session.js'; interface SseOptions { maxEvents?: number; @@ -4460,4 +4461,126 @@ describe.concurrent('SSE Plugin', () => { expect(() => new ValidReplayer({ ttl: 50.5 })).toThrow(/Invalid ValidReplayer options.*ttl/i); }); }); + + describe('replay failures and invalid event ids', () => { + it('ends the stream and releases the maxSessions slot when replay() throws', async ({ + onTestFinished, + expect, + }) => { + const server = Hapi.server({ port: 0 }); + onTestFinished(() => server.stop()); + await server.register({ plugin: SsePlugin, options: { retry: null, keepAlive: { interval: 20 } } }); + + let failed: Session | undefined; + let unsubscribed = 0; + let reconnected = 0; + const logged: string[] = []; + + server.events.on({ name: 'request', channels: 'app' }, (_request, event) => { + logged.push(`${event.tags.join(',')}: ${event.error instanceof Error ? event.error.message : ''}`); + }); + + server.sse.subscription('/events', { + maxSessions: 1, + replay: { + record: () => {}, + replay: () => { + throw new Error('replay store unavailable'); + }, + }, + onSubscribe: (session) => { + failed ??= session; + }, + onUnsubscribe: () => { + unsubscribed++; + }, + onReconnect: () => { + reconnected++; + }, + }); + await server.start(); + + const url = `http://localhost:${server.info.port}/events`; + const ended = await collectSse(url, { timeout: 500, headers: { 'last-event-id': '1' } }); + + expect(ended.status).toBe(200); + expect(ended.events).toEqual([]); + expect(reconnected).toBe(0); + expect(logged).toEqual(['sse,replay,error: replay store unavailable']); + await expect.poll(() => unsubscribed).toBe(1); + expect(failed!.isOpen).toBe(false); + expect(server.sse.sessionCount).toBe(0); + + const next = collectSse(url, { timeout: 1000 }); + + await expect.poll(() => server.sse.sessionCount).toBe(1); + await server.sse.publish('/events', 'after'); + + expect((await next).status).toBe(200); + }); + + it('logs a non-Error thrown by replay() as a string', async ({ onTestFinished }) => { + const server = Hapi.server({ port: 0 }); + onTestFinished(() => server.stop()); + await server.register({ plugin: SsePlugin, options: { retry: null, keepAlive: false } }); + + const logged: unknown[] = []; + + server.events.on({ name: 'request', channels: 'app' }, (_request, event) => { + logged.push(event.data); + }); + + server.sse.subscription('/events', { + replay: { + record: () => {}, + replay: () => { + throw 'store down'; + }, + }, + }); + await server.start(); + + await collectSse(`http://localhost:${server.info.port}/events`, { + timeout: 500, + headers: { 'last-event-id': '1' }, + }); + + expect(logged).toEqual(['store down']); + }); + + it('rejects a publish whose id contains a null character before recording it', async ({ + onTestFinished, + }) => { + const server = Hapi.server({ port: 0 }); + onTestFinished(() => server.stop()); + await server.register({ plugin: SsePlugin, options: { retry: null, keepAlive: false } }); + + const recorded: string[] = []; + + server.sse.subscription('/events', { + replay: { + record: (entry) => { + recorded.push(entry.id); + }, + replay: () => [], + }, + }); + + await expect(server.sse.publish('/events', 'poison', { id: 'a\u0000b' })).rejects.toThrow( + /null characters/, + ); + expect(recorded).toEqual([]); + }); + + it('rejects a broadcast whose id contains a null character even with no sessions', async ({ + onTestFinished, + }) => { + const server = Hapi.server({ port: 0 }); + onTestFinished(() => server.stop()); + await server.register({ plugin: SsePlugin }); + server.sse.subscription('/events'); + + await expect(server.sse.broadcast('poison', { id: 'a\u0000b' })).rejects.toThrow(/null characters/); + }); + }); });