Skip to content

Commit e635cdb

Browse files
fix(realtime-transcription): bound inbound websocket payloads (#102443)
* fix(realtime-transcription): bound inbound websocket payloads * test(realtime-transcription): tighten payload cap proof --------- Co-authored-by: Peter Steinberger <[email protected]>
1 parent 29d2a1e commit e635cdb

2 files changed

Lines changed: 69 additions & 1 deletion

File tree

src/realtime-transcription/websocket-session.test.ts

Lines changed: 64 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,10 @@ import type { AddressInfo } from "node:net";
44
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
55
import type WebSocket from "ws";
66
import { WebSocketServer } from "ws";
7-
import { createRealtimeTranscriptionWebSocketSession } from "./websocket-session.js";
7+
import {
8+
createRealtimeTranscriptionWebSocketSession,
9+
REALTIME_TRANSCRIPTION_WS_MAX_PAYLOAD_BYTES,
10+
} from "./websocket-session.js";
811

912
let cleanup: (() => Promise<void>) | undefined;
1013

@@ -403,4 +406,64 @@ describe("createRealtimeTranscriptionWebSocketSession", () => {
403406
expect(closeError).toBeInstanceOf(Error);
404407
expect(closeError.message).toBe("test realtime transcription connection closed before ready");
405408
});
409+
410+
it("delivers a legitimate large inbound message below the payload cap", async () => {
411+
// Well above any real transcript message yet far under the 16 MiB cap: proves
412+
// the bound does not reject legitimate large provider traffic.
413+
const largeText = "x".repeat(2 * 1024 * 1024);
414+
const server = await createRealtimeServer({
415+
initialEvent: { type: "transcript", text: largeText },
416+
});
417+
const onMessage = vi.fn();
418+
const session = createRealtimeTranscriptionWebSocketSession<{ type?: string; text?: string }>({
419+
providerId: "test",
420+
callbacks: {},
421+
url: server.url,
422+
readyOnOpen: true,
423+
onMessage,
424+
sendAudio: (audio, transport) => {
425+
transport.sendBinary(audio);
426+
},
427+
});
428+
429+
await session.connect();
430+
await vi.waitFor(() => {
431+
expect(onMessage).toHaveBeenCalledTimes(1);
432+
});
433+
const event = requireFirstMockArg(onMessage, "large inbound message");
434+
expect(event).toEqual({ type: "transcript", text: largeText });
435+
session.close();
436+
});
437+
438+
it("drops an oversized inbound message before it reaches the provider parser", async () => {
439+
// ws rejects a message above maxPayload with an error + 1009 close, so an
440+
// oversized upstream message never reaches onMessage/JSON parse.
441+
const oversized = "x".repeat(REALTIME_TRANSCRIPTION_WS_MAX_PAYLOAD_BYTES + 1);
442+
const server = await createRealtimeServer({ initialText: oversized });
443+
const onError = vi.fn();
444+
const onMessage = vi.fn(() => {
445+
throw new Error("oversized frame should not reach provider handler");
446+
});
447+
const session = createRealtimeTranscriptionWebSocketSession({
448+
providerId: "test",
449+
callbacks: { onError },
450+
url: server.url,
451+
readyOnOpen: true,
452+
onMessage,
453+
sendAudio: (audio, transport) => {
454+
transport.sendBinary(audio);
455+
},
456+
});
457+
458+
await session.connect();
459+
await vi.waitFor(() => {
460+
expect(onError).toHaveBeenCalledTimes(1);
461+
});
462+
expect(onMessage).not.toHaveBeenCalled();
463+
const overflowError = requireFirstMockArg(onError, "oversized inbound message error");
464+
expect(overflowError).toBeInstanceOf(Error);
465+
expect(overflowError).toHaveProperty("code", "WS_ERR_UNSUPPORTED_MESSAGE_LENGTH");
466+
expect(overflowError.message).toMatch(/max payload/i);
467+
session.close();
468+
});
406469
});

src/realtime-transcription/websocket-session.ts

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -50,6 +50,10 @@ const DEFAULT_CLOSE_TIMEOUT_MS = 5_000;
5050
const DEFAULT_MAX_RECONNECT_ATTEMPTS = 5;
5151
const DEFAULT_RECONNECT_DELAY_MS = 1000;
5252
const DEFAULT_MAX_QUEUED_BYTES = 2 * 1024 * 1024;
53+
// Bound inbound messages before ws buffers them for JSON parsing. The 16 MiB cap
54+
// matches realtime voice; ws rejects larger messages with close 1009 before
55+
// they reach onMessage, replacing its 100 MiB client default.
56+
export const REALTIME_TRANSCRIPTION_WS_MAX_PAYLOAD_BYTES = 16 * 1024 * 1024;
5357

5458
function rawWsDataToBuffer(data: RawData): Buffer {
5559
if (Buffer.isBuffer(data)) {
@@ -249,6 +253,7 @@ class WebSocketRealtimeTranscriptionSession<Event> implements RealtimeTranscript
249253
try {
250254
this.ws = new WebSocket(this.currentUrl, {
251255
headers: connection.headers,
256+
maxPayload: REALTIME_TRANSCRIPTION_WS_MAX_PAYLOAD_BYTES,
252257
...(proxyAgent ? { agent: proxyAgent } : {}),
253258
});
254259
} catch (error) {

0 commit comments

Comments
 (0)