Skip to content
Merged
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
45 changes: 45 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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:
Expand Down
13 changes: 11 additions & 2 deletions src/Client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -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;
}
Expand Down Expand Up @@ -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;
}
Expand Down
8 changes: 8 additions & 0 deletions src/HeartbeatTimeoutError.ts
Original file line number Diff line number Diff line change
@@ -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 };
289 changes: 289 additions & 0 deletions src/heartbeatTimeout.test.ts
Original file line number Diff line number Diff line change
@@ -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<void>(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<void>(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<void> => {
const connectionClosed = Promise.withResolvers<void>();
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<void> => {
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<void> => {
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<void> => {
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<void> => {
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<void> => {
const connectionClosed = Promise.withResolvers<void>();
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<void> => {
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<void> => {
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<void> => {
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);
});
});
});
7 changes: 7 additions & 0 deletions src/heartbeatTimeout.ts
Original file line number Diff line number Diff line change
@@ -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 };
Loading
Loading