Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 5 additions & 3 deletions API.md
Original file line number Diff line number Diff line change
Expand Up @@ -107,8 +107,8 @@ server.sse.subscription('/chat/{room}', {
| `filter` | `(path, message, opts) => boolean \| { override } \| Promise<...>` | Per-session delivery filter |
| `onSubscribe` | `(session, path, params) => void \| Promise<void>` | 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<void>` | 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<void>` | 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. |

Expand Down Expand Up @@ -136,7 +136,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?)`

Expand All @@ -149,6 +149,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.
Expand Down
12 changes: 9 additions & 3 deletions src/event-buffer.ts
Original file line number Diff line number Diff line change
@@ -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 = '';
Expand All @@ -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`;

Expand Down Expand Up @@ -60,6 +64,8 @@ export class EventBuffer {
}

push(data: unknown, event?: string, id?: string): this {
assertEventId(id);

if (event) {
this.event(event);
}
Expand Down
6 changes: 6 additions & 0 deletions src/replayer.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,7 @@
import Joi from 'joi';

import { assertEventId } from './event-buffer.js';

export interface ReplayEntry {
data: unknown;
event?: string;
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down
12 changes: 9 additions & 3 deletions src/sse.ts
Original file line number Diff line number Diff line change
Expand Up @@ -310,10 +310,16 @@ export const SsePlugin: NamedPlugin<SsePluginOptions> = {
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);
}
}

Expand Down
5 changes: 5 additions & 0 deletions src/subscription.ts
Original file line number Diff line number Diff line change
@@ -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';

Expand Down Expand Up @@ -145,6 +146,8 @@ export class SubscriptionRegistry {
data: T,
opts?: { event?: string; id?: string; internal?: unknown; matchMode?: 'pattern' | 'literal' },
): Promise<number> {
assertEventId(opts?.id);

const matched = this.matchPath(path);

if (!matched) {
Expand Down Expand Up @@ -204,6 +207,8 @@ export class SubscriptionRegistry {
}

async broadcast(data: unknown, opts?: { event?: string; id?: string }): Promise<number> {
assertEventId(opts?.id);

let delivered = 0;

for (const sub of this.#subscriptions.values()) {
Expand Down
7 changes: 7 additions & 0 deletions test/event-buffer.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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();

Expand Down
13 changes: 13 additions & 0 deletions test/replayer.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 });

Expand Down Expand Up @@ -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());
Expand Down
123 changes: 123 additions & 0 deletions test/sse.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -4467,4 +4468,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/);
});
});
});
Loading