diff --git a/crates/bindings-typescript/src/sdk/db_connection_impl.ts b/crates/bindings-typescript/src/sdk/db_connection_impl.ts index 12dfca7afbf..3554c126f12 100644 --- a/crates/bindings-typescript/src/sdk/db_connection_impl.ts +++ b/crates/bindings-typescript/src/sdk/db_connection_impl.ts @@ -1160,6 +1160,17 @@ export class DbConnectionImpl 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; diff --git a/crates/bindings-typescript/src/sdk/event_emitter.ts b/crates/bindings-typescript/src/sdk/event_emitter.ts index a3447ce4b58..8fb0e1413d2 100644 --- a/crates/bindings-typescript/src/sdk/event_emitter.ts +++ b/crates/bindings-typescript/src/sdk/event_emitter.ts @@ -1,3 +1,5 @@ +import { stdbLogger } from './logger.ts'; + // eslint-disable-next-line @typescript-eslint/no-unsafe-function-type export class EventEmitter { #events: Map> = new Map(); @@ -25,8 +27,14 @@ export class EventEmitter { 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); + } } } } diff --git a/crates/bindings-typescript/tests/db_connection.test.ts b/crates/bindings-typescript/tests/db_connection.test.ts index 9e340637eb8..0ab4c31e631 100644 --- a/crates/bindings-typescript/tests/db_connection.test.ts +++ b/crates/bindings-typescript/tests/db_connection.test.ts @@ -1,4 +1,4 @@ -import { assertType, beforeEach, describe, expect, test } from 'vitest'; +import { assertType, beforeEach, describe, expect, test, vi } from 'vitest'; import { BinaryWriter, ConnectionId, @@ -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(); + 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(); const wsAdapter = new WebsocketTestAdapter();