diff --git a/README.md b/README.md index 54a8220..fce208d 100644 --- a/README.md +++ b/README.md @@ -260,6 +260,29 @@ for await (const row of client.runEventQlQuery(` controller.abort(); ``` +#### Detecting a Stalled Connection + +While a query is running, EventSourcingDB sends a heartbeat every second whenever there is no row to send. If neither a row nor a heartbeat arrives for 30 seconds, for example because a proxy keeps the connection open but no longer passes anything through, the client SDK closes the connection and the iterator throws a `HeartbeatTimeoutError`: + +```typescript +import { HeartbeatTimeoutError } from 'eventsourcingdb'; + +try { + for await (const row of client.runEventQlQuery(` + FROM e IN events + PROJECT INTO e + `)) { + // ... + } +} catch (error) { + if (error instanceof HeartbeatTimeoutError) { + // ... + } +} +``` + +*Note that the time you spend handling a row does not count towards the 30 seconds, and that aborting the query does not cause a `HeartbeatTimeoutError`.* + ### Observing Events To observe all events of a subject, call the `observeEvents` function with the subject as the first argument and an options object as the second argument. Set the `recursive` option to `false`. This ensures that only events of the given subject are returned, not events of nested subjects. @@ -344,6 +367,28 @@ for await (const event of client.observeEvents('/books/42', { controller.abort(); ``` +#### Detecting a Stalled Connection + +While observing, EventSourcingDB sends a heartbeat every second whenever there is no event to send. If neither an event nor a heartbeat arrives for 30 seconds, for example because a proxy keeps the connection open but no longer passes anything through, the client SDK closes the connection and the iterator throws a `HeartbeatTimeoutError`: + +```typescript +import { HeartbeatTimeoutError } from 'eventsourcingdb'; + +try { + for await (const event of client.observeEvents('/books/42', { + recursive: false + })) { + // ... + } +} catch (error) { + if (error instanceof HeartbeatTimeoutError) { + // ... + } +} +``` + +*Note that the time you spend handling an event does not count towards the 30 seconds, and that aborting observing does not cause a `HeartbeatTimeoutError`.* + ### Registering an Event Schema To register an event schema, call the `registerEventSchema` function and hand over an event type and the desired schema: diff --git a/src/Client.ts b/src/Client.ts index cf81ce7..e885a7c 100644 --- a/src/Client.ts +++ b/src/Client.ts @@ -2,6 +2,7 @@ import { convertCloudEventToEvent } from './convertCloudEventToEvent.js'; import type { Event } from './Event.js'; import type { EventCandidate } from './EventCandidate.js'; import type { EventType } from './EventType.js'; +import { heartbeatTimeout } from './heartbeatTimeout.js'; import { isCloudEvent } from './isCloudEvent.js'; import { isValidServerHeader } from './isValidServerHeader.js'; import { readNdJsonStream } from './ndjson/readNdJsonStream.js'; @@ -277,7 +278,11 @@ class Client { throw new Error('Failed to run EventQL query.'); } - for await (const line of readNdJsonStream(response.body, combinedSignal)) { + for await (const line of readNdJsonStream( + response.body, + combinedSignal, + heartbeatTimeout.milliseconds, + )) { if (isStreamHeartbeat(line)) { continue; } @@ -360,7 +365,11 @@ class Client { throw new Error('Failed to observe events.'); } - for await (const line of readNdJsonStream(response.body, combinedSignal)) { + for await (const line of readNdJsonStream( + response.body, + combinedSignal, + heartbeatTimeout.milliseconds, + )) { if (isStreamHeartbeat(line)) { continue; } diff --git a/src/HeartbeatTimeoutError.ts b/src/HeartbeatTimeoutError.ts new file mode 100644 index 0000000..ea81ff1 --- /dev/null +++ b/src/HeartbeatTimeoutError.ts @@ -0,0 +1,8 @@ +class HeartbeatTimeoutError extends Error { + public constructor() { + super('No event and no heartbeat arrived for 30 seconds.'); + this.name = 'HeartbeatTimeoutError'; + } +} + +export { HeartbeatTimeoutError }; diff --git a/src/heartbeatTimeout.test.ts b/src/heartbeatTimeout.test.ts new file mode 100644 index 0000000..a3931d9 --- /dev/null +++ b/src/heartbeatTimeout.test.ts @@ -0,0 +1,289 @@ +import assert from 'node:assert/strict'; +import { createServer, type Server, type ServerResponse } from 'node:http'; +import { afterEach, beforeEach, suite, test } from 'node:test'; +import { setTimeout as sleep } from 'node:timers/promises'; +import { Client } from './Client.js'; +import type { Event } from './Event.js'; +import { HeartbeatTimeoutError } from './HeartbeatTimeoutError.js'; +import { heartbeatTimeout } from './heartbeatTimeout.js'; + +const timeout = 250; +const heartbeatInterval = 50; + +const heartbeatLine = '{"type":"heartbeat","payload":{}}\n'; + +const getEventLine = (id: string): string => + `${JSON.stringify({ + type: 'event', + payload: { + specversion: '1.0', + id, + time: '2026-10-03T12:00:00.000Z', + source: 'https://www.eventsourcingdb.io', + subject: '/test', + type: 'io.eventsourcingdb.test', + datacontenttype: 'application/json', + data: { value: 23 }, + hash: 'hash', + predecessorhash: 'predecessorhash', + signature: null, + }, + })}\n`; + +const getRowLine = (value: number): string => + `${JSON.stringify({ + type: 'row', + payload: { value }, + })}\n`; + +const sendHeartbeats = (response: ServerResponse): void => { + const interval = setInterval(() => { + response.write(heartbeatLine); + }, heartbeatInterval); + + response.on('close', () => { + clearInterval(interval); + }); +}; + +suite('heartbeatTimeout', { timeout: 30_000 }, () => { + const originalTimeout = heartbeatTimeout.milliseconds; + + let server: Server; + let client: Client; + let respond: (response: ServerResponse) => void; + + beforeEach(async () => { + heartbeatTimeout.milliseconds = timeout; + + server = createServer((_request, response) => { + response.writeHead(200, { + server: 'EventSourcingDB/test', + 'content-type': 'application/x-ndjson', + }); + respond(response); + }); + + await new Promise(resolve => { + server.listen(0, 'localhost', resolve); + }); + + const address = server.address(); + if (address === null || typeof address === 'string') { + throw new Error('Failed to get server address.'); + } + + client = new Client(new URL(`http://localhost:${address.port}/`), 'secret'); + }); + + afterEach(async () => { + heartbeatTimeout.milliseconds = originalTimeout; + + server.closeAllConnections(); + await new Promise(resolve => { + server.close(() => resolve()); + }); + }); + + suite('observeEvents', () => { + test('ends with a heartbeat timeout error if neither an event nor a heartbeat arrives in time.', { + timeout: 5000, + }, async (): Promise => { + const connectionClosed = Promise.withResolvers(); + respond = (response: ServerResponse): void => { + response.on('close', () => connectionClosed.resolve()); + response.write(heartbeatLine); + }; + + // The signal is never aborted, so only the timeout can close the connection. + const controller = new AbortController(); + const startedAt = performance.now(); + + await assert.rejects( + async () => { + for await (const _event of client.observeEvents( + '/', + { recursive: true }, + controller.signal, + )) { + // Intentionally left blank. + } + }, + error => { + assert.ok(error instanceof HeartbeatTimeoutError); + assert.equal(error.name, 'HeartbeatTimeoutError'); + assert.equal(error.message, 'No event and no heartbeat arrived for 30 seconds.'); + return true; + }, + ); + + const elapsed = performance.now() - startedAt; + assert.ok(elapsed >= timeout - 10); + assert.ok(elapsed < timeout * 4); + + await connectionClosed.promise; + }); + + test('does not end as long as heartbeats arrive in time.', async (): Promise => { + respond = (response: ServerResponse): void => { + sendHeartbeats(response); + setTimeout(() => { + response.end(getEventLine('0')); + }, timeout * 4); + }; + + const eventsObserved: Event[] = []; + for await (const event of client.observeEvents('/', { recursive: true })) { + eventsObserved.push(event); + } + + assert.equal(eventsObserved.length, 1); + }); + + test('observes events that arrive in time.', async (): Promise => { + respond = (response: ServerResponse): void => { + response.write(heartbeatLine); + response.write(getEventLine('0')); + setTimeout(() => { + response.end(getEventLine('1')); + }, timeout / 2); + }; + + const eventsObserved: Event[] = []; + for await (const event of client.observeEvents('/', { recursive: true })) { + eventsObserved.push(event); + } + + assert.equal(eventsObserved.length, 2); + assert.equal(eventsObserved[0]?.id, '0'); + assert.equal(eventsObserved[1]?.id, '1'); + }); + + test('does not count the time the caller spends on an event.', async (): Promise => { + respond = (response: ServerResponse): void => { + response.write(getEventLine('0')); + sendHeartbeats(response); + setTimeout(() => { + response.end(getEventLine('1')); + }, timeout * 4); + }; + + const eventsObserved: Event[] = []; + for await (const event of client.observeEvents('/', { recursive: true })) { + eventsObserved.push(event); + await sleep(timeout * 2); + } + + assert.equal(eventsObserved.length, 2); + }); + + test('ends without an error if the caller aborts.', async (): Promise => { + respond = (response: ServerResponse): void => { + response.write(heartbeatLine); + }; + + const controller = new AbortController(); + setTimeout(() => { + controller.abort(); + }, timeout / 2); + + let didObserveEvents = false; + for await (const _event of client.observeEvents( + '/', + { recursive: true }, + controller.signal, + )) { + didObserveEvents = true; + } + + assert.equal(didObserveEvents, false); + }); + }); + + suite('runEventQlQuery', () => { + test('ends with a heartbeat timeout error if neither a row nor a heartbeat arrives in time.', { + timeout: 5000, + }, async (): Promise => { + const connectionClosed = Promise.withResolvers(); + respond = (response: ServerResponse): void => { + response.on('close', () => connectionClosed.resolve()); + response.write(heartbeatLine); + }; + + const startedAt = performance.now(); + + await assert.rejects( + async () => { + for await (const _row of client.runEventQlQuery('FROM e IN events PROJECT INTO e')) { + // Intentionally left blank. + } + }, + error => { + assert.ok(error instanceof HeartbeatTimeoutError); + assert.equal(error.message, 'No event and no heartbeat arrived for 30 seconds.'); + return true; + }, + ); + + const elapsed = performance.now() - startedAt; + assert.ok(elapsed >= timeout - 10); + assert.ok(elapsed < timeout * 4); + + await connectionClosed.promise; + }); + + test('does not end as long as heartbeats arrive in time.', async (): Promise => { + respond = (response: ServerResponse): void => { + sendHeartbeats(response); + setTimeout(() => { + response.end(getRowLine(23)); + }, timeout * 4); + }; + + const rowsRead: unknown[] = []; + for await (const row of client.runEventQlQuery('FROM e IN events PROJECT INTO e')) { + rowsRead.push(row); + } + + assert.deepEqual(rowsRead, [{ value: 23 }]); + }); + + test('reads rows that arrive in time.', async (): Promise => { + respond = (response: ServerResponse): void => { + response.write(heartbeatLine); + response.write(getRowLine(23)); + setTimeout(() => { + response.end(getRowLine(42)); + }, timeout / 2); + }; + + const rowsRead: unknown[] = []; + for await (const row of client.runEventQlQuery('FROM e IN events PROJECT INTO e')) { + rowsRead.push(row); + } + + assert.deepEqual(rowsRead, [{ value: 23 }, { value: 42 }]); + }); + + test('ends without an error if the caller aborts.', async (): Promise => { + respond = (response: ServerResponse): void => { + response.write(heartbeatLine); + }; + + const controller = new AbortController(); + setTimeout(() => { + controller.abort(); + }, timeout / 2); + + let didReadRows = false; + for await (const _row of client.runEventQlQuery( + 'FROM e IN events PROJECT INTO e', + controller.signal, + )) { + didReadRows = true; + } + + assert.equal(didReadRows, false); + }); + }); +}); diff --git a/src/heartbeatTimeout.ts b/src/heartbeatTimeout.ts new file mode 100644 index 0000000..27e712a --- /dev/null +++ b/src/heartbeatTimeout.ts @@ -0,0 +1,7 @@ +// A stream that carries heartbeats ends with a HeartbeatTimeoutError if no line +// arrives for this long. The package does not export this, only tests shorten it. +const heartbeatTimeout = { + milliseconds: 30_000, +}; + +export { heartbeatTimeout }; diff --git a/src/index.ts b/src/index.ts index ddb0d26..f147f80 100644 --- a/src/index.ts +++ b/src/index.ts @@ -3,6 +3,7 @@ import { Container } from './Container.js'; import type { Event } from './Event.js'; import type { EventCandidate } from './EventCandidate.js'; import type { EventType } from './EventType.js'; +import { HeartbeatTimeoutError } from './HeartbeatTimeoutError.js'; import { isEventQlQueryTrue } from './isEventQlQueryTrue.js'; import { isSubjectOnEventId } from './isSubjectOnEventId.js'; import { isSubjectPopulated } from './isSubjectPopulated.js'; @@ -22,6 +23,7 @@ export type { export { Client, Container, + HeartbeatTimeoutError, isEventQlQueryTrue, isSubjectOnEventId, isSubjectPopulated, diff --git a/src/ndjson/readNdJsonStream.ts b/src/ndjson/readNdJsonStream.ts index 1ca4e08..9fc8465 100644 --- a/src/ndjson/readNdJsonStream.ts +++ b/src/ndjson/readNdJsonStream.ts @@ -1,11 +1,37 @@ +import { HeartbeatTimeoutError } from '../HeartbeatTimeoutError.js'; + const readNdJsonStream = async function* ( stream: ReadableStream, signal: AbortSignal, + heartbeatTimeoutInMilliseconds?: number, ): AsyncGenerator, void, void> { const reader = stream.getReader(); const decoder = new TextDecoder('utf-8'); let buffer = ''; + let heartbeatTimer: ReturnType | undefined; + let hasHeartbeatTimedOut = false; + + // The timer only runs while waiting for the next line, not while the caller + // handles a line, so a slow caller does not cause a heartbeat timeout. + const startHeartbeatTimer = (): void => { + if (heartbeatTimeoutInMilliseconds === undefined || heartbeatTimer !== undefined) { + return; + } + + heartbeatTimer = setTimeout(() => { + hasHeartbeatTimedOut = true; + reader.cancel().catch(() => { + // Intentionally left blank. + }); + }, heartbeatTimeoutInMilliseconds); + }; + + const stopHeartbeatTimer = (): void => { + clearTimeout(heartbeatTimer); + heartbeatTimer = undefined; + }; + const onAbort = (): void => { reader.cancel().catch(() => { // Intentionally left blank. @@ -23,6 +49,8 @@ const readNdJsonStream = async function* ( try { while (!signal.aborted) { + startHeartbeatTimer(); + // biome-ignore lint/performance/noAwaitInLoops: Awaiting the result is fine here, although we are in a loop. const { done, value } = await reader.read(); if (done) { @@ -37,13 +65,19 @@ const readNdJsonStream = async function* ( buffer = buffer.slice(index + 1); if (line) { + stopHeartbeatTimer(); yield JSON.parse(line); } index = buffer.indexOf('\n'); } } + + if (hasHeartbeatTimedOut) { + throw new HeartbeatTimeoutError(); + } } finally { + stopHeartbeatTimer(); signal.removeEventListener('abort', onAbort); await reader.cancel().catch(() => { // Intentionally left blank.