Skip to content
Draft
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
11 changes: 11 additions & 0 deletions crates/bindings-typescript/src/sdk/db_connection_impl.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1160,6 +1160,17 @@ export class DbConnectionImpl<RemoteModule extends UntypedRemoteModule>
this.#processMessage(data);
}
}
} catch (e) {
// A message that could not be decoded or applied leaves the client cache
// inconsistent, and rethrowing from a WebSocket listener crashes Node
// hosts. Fail the connection instead: drop what is queued and close, so
// `onDisconnect` receives the error. (Callback errors never get here;
// `EventEmitter` logs them.)
stdbLogger('error', 'Failed to process a server message', e);
this.#connectionError = toError(e);
this.#inboundQueueOffset = this.#inboundQueue.length;
this.isActive = false;
this.ws?.close();
} finally {
if (this.#inboundQueueOffset >= this.#inboundQueue.length) {
this.#inboundQueue.length = 0;
Expand Down
10 changes: 9 additions & 1 deletion crates/bindings-typescript/src/sdk/event_emitter.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
import { stdbLogger } from './logger.ts';

// eslint-disable-next-line @typescript-eslint/no-unsafe-function-type
export class EventEmitter<Key, Callback extends Function = Function> {
#events: Map<Key, Set<Callback>> = new Map();
Expand Down Expand Up @@ -25,8 +27,14 @@ export class EventEmitter<Key, Callback extends Function = Function> {
return;
}

// Like the C# SDK, a throwing callback is logged and does not stop the
// others or escape into the WebSocket listener, which would crash Node.
for (const callback of callbacks) {
callback(...args);
try {
callback(...args);
} catch (e) {
stdbLogger('error', 'A callback threw', e);
}
}
}
}
81 changes: 80 additions & 1 deletion crates/bindings-typescript/tests/db_connection.test.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
import { assertType, beforeEach, describe, expect, test } from 'vitest';
import { assertType, beforeEach, describe, expect, test, vi } from 'vitest';
import {
BinaryWriter,
ConnectionId,
Expand Down Expand Up @@ -213,6 +213,85 @@ describe('DbConnection', () => {
expect(client.isActive).toBe(false);
});

test('logs an error thrown by a callback and keeps the connection open', async () => {
const wsAdapter = new WebsocketTestAdapter();
const onDisconnect = vi.fn();
const secondOnConnect = vi.fn();

const client = DbConnection.builder()
.withUri('ws://127.0.0.1:1234')
.withDatabaseName('db')
.withWSFn(wsAdapter.openWebSocket)
.onConnect(() => {
throw new Error('callback failed');
})
.onConnect(secondOnConnect)
.onDisconnect(onDisconnect)
.build();

await client['wsPromise'];
wsAdapter.acceptConnection();
// Node rethrows errors from WebSocket listeners on the next tick, which
// crashes the host process.
expect(() =>
wsAdapter.sendToClient(
ServerMessage.InitialConnection({
identity: anIdentity,
token: 'a-token',
connectionId: ConnectionId.random(),
})
)
).not.toThrow();

expect(secondOnConnect).toHaveBeenCalledOnce();
expect(client.isActive).toBe(true);
expect(wsAdapter.closed).toBe(false);
expect(onDisconnect).not.toHaveBeenCalled();
});

test('reports a row it cannot decode through onDisconnect instead of throwing', async () => {
const onDisconnectPromise = new Deferred<Error | undefined>();
const wsAdapter = new WebsocketTestAdapter();

const client = DbConnection.builder()
.withUri('ws://127.0.0.1:1234')
.withDatabaseName('db')
.withWSFn(wsAdapter.openWebSocket)
.onDisconnect((_ctx, error) => onDisconnectPromise.resolve(error))
.build();

await client['wsPromise'];
wsAdapter.acceptConnection();
wsAdapter.sendToClient(
ServerMessage.InitialConnection({
identity: anIdentity,
token: 'a-token',
connectionId: ConnectionId.random(),
})
);
expect(client.isActive).toBe(true);

// A `player` row one byte short, as when the client's bindings are out of
// step with the module's schema.
const truncatedPlayer = encodePlayer({
id: 1,
userId: anIdentity,
name: 'drogus',
location: { x: 0, y: 0 },
}).slice(0, -1);
expect(() =>
wsAdapter.sendToClient(
ServerMessage.TransactionUpdate({
querySets: [makeQuerySetUpdate(0, 'player', truncatedPlayer)],
})
)
).not.toThrow();

expect(await onDisconnectPromise.promise).toBeInstanceOf(RangeError);
expect(client.isActive).toBe(false);
expect(wsAdapter.closed).toBe(true);
});

test('marks disconnect as requested when disconnect() is called', async () => {
const onDisconnectPromise = new Deferred<void>();
const wsAdapter = new WebsocketTestAdapter();
Expand Down
Loading