From 220df691bf09a3b3a0f42ec6425691bf2169bda0 Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Fri, 31 Jul 2026 18:05:49 +0800 Subject: [PATCH 1/4] fix(mistral): bound realtime transcript accumulation --- .../realtime-transcription-provider.ts | 65 +++++++++++++++++-- 1 file changed, 59 insertions(+), 6 deletions(-) diff --git a/extensions/mistral/realtime-transcription-provider.ts b/extensions/mistral/realtime-transcription-provider.ts index 649c3db99394..5b07a5dfaa9b 100644 --- a/extensions/mistral/realtime-transcription-provider.ts +++ b/extensions/mistral/realtime-transcription-provider.ts @@ -62,6 +62,9 @@ const MISTRAL_REALTIME_CLOSE_TIMEOUT_MS = 5_000; const MISTRAL_REALTIME_MAX_RECONNECT_ATTEMPTS = 5; const MISTRAL_REALTIME_RECONNECT_DELAY_MS = 1000; const MISTRAL_REALTIME_MAX_QUEUED_BYTES = 2 * 1024 * 1024; +const MISTRAL_REALTIME_MAX_PARTIAL_TRANSCRIPT_BYTES = 256 * 1024; +const MISTRAL_REALTIME_PARTIAL_TRANSCRIPT_OVERFLOW_MESSAGE = + "Mistral realtime transcription exceeded the 256 KiB in-progress transcript limit"; function readNestedMistralConfig(rawConfig: RealtimeTranscriptionProviderConfig) { const raw = readRecord(rawConfig); @@ -161,15 +164,53 @@ function readErrorDetail(event: MistralRealtimeTranscriptionEvent): string { return "Mistral realtime transcription error"; } +function measureTranscriptDeltaBytes(partialText: string, delta: string): number { + const previousCodeUnit = partialText.charCodeAt(partialText.length - 1); + const nextCodeUnit = delta.charCodeAt(0); + const completesSplitSurrogatePair = + previousCodeUnit >= 0xd800 && + previousCodeUnit <= 0xdbff && + nextCodeUnit >= 0xdc00 && + nextCodeUnit <= 0xdfff; + // Separate UTF-8 measurements encode split surrogates as two replacement + // characters (six bytes); the combined transcript encodes one four-byte code point. + return Buffer.byteLength(delta, "utf8") - (completesSplitSurrogatePair ? 2 : 0); +} + function createMistralRealtimeTranscriptionSession( config: MistralRealtimeTranscriptionSessionConfig, ): RealtimeTranscriptionSession { let partialText = ""; + let partialBytes = 0; + let terminal = false; + + const clearPartial = () => { + partialText = ""; + partialBytes = 0; + }; + + const failPartialOverflow = (transport: RealtimeTranscriptionWebSocketTransport) => { + if (terminal) { + return; + } + terminal = true; + clearPartial(); + transport.closeNow(); + try { + config.onError?.(new Error(MISTRAL_REALTIME_PARTIAL_TRANSCRIPT_OVERFLOW_MESSAGE)); + } catch { + // The terminal provider error already owns the outcome. Do not let an + // observer exception re-enter shared error dispatch and emit it twice. + } + }; const handleEvent = ( event: MistralRealtimeTranscriptionEvent, transport: RealtimeTranscriptionWebSocketTransport, ) => { + if (terminal) { + return; + } if (event.type === "session.created") { transport.sendJson({ type: "session.update", @@ -190,23 +231,35 @@ function createMistralRealtimeTranscriptionSession( switch (event.type) { case "transcription.text.delta": if (event.text) { + const deltaBytes = measureTranscriptDeltaBytes(partialText, event.text); + if (deltaBytes > MISTRAL_REALTIME_MAX_PARTIAL_TRANSCRIPT_BYTES - partialBytes) { + failPartialOverflow(transport); + return; + } partialText += event.text; + partialBytes += deltaBytes; config.onPartial?.(partialText); } return; case "transcription.segment": if (event.text) { config.onTranscript?.(event.text); - partialText = ""; + clearPartial(); } return; - case "transcription.done": - if (partialText.trim()) { - config.onTranscript?.(partialText); - partialText = ""; + case "transcription.done": { + terminal = true; + const transcript = partialText; + clearPartial(); + try { + if (transcript.trim()) { + config.onTranscript?.(transcript); + } + } finally { + transport.closeNow(); } - transport.closeNow(); return; + } case "error": config.onError?.(new Error(readErrorDetail(event))); From 81915d8566c7fcf11a43c60c161ec32f7b88bc15 Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Fri, 31 Jul 2026 18:05:53 +0800 Subject: [PATCH 2/4] test(mistral): cover realtime transcript overflow --- ...e-transcription-provider.lifecycle.test.ts | 178 ++++++++++++++++++ 1 file changed, 178 insertions(+) create mode 100644 extensions/mistral/realtime-transcription-provider.lifecycle.test.ts diff --git a/extensions/mistral/realtime-transcription-provider.lifecycle.test.ts b/extensions/mistral/realtime-transcription-provider.lifecycle.test.ts new file mode 100644 index 000000000000..67d2c56036d1 --- /dev/null +++ b/extensions/mistral/realtime-transcription-provider.lifecycle.test.ts @@ -0,0 +1,178 @@ +// Mistral lifecycle tests cover bounded transcript accumulation and terminal events. +import { beforeEach, describe, expect, it, vi } from "vitest"; +import { buildMistralRealtimeTranscriptionProvider } from "./realtime-transcription-provider.js"; + +const { FakeWebSocket } = vi.hoisted(() => { + type Listener = (...args: unknown[]) => void; + + class MockWebSocket { + static readonly OPEN = 1; + static readonly CLOSED = 3; + static instances: MockWebSocket[] = []; + + binaryType = "nodebuffer"; + closeCalls = 0; + readonly listeners = new Map(); + readyState = 0; + sent: string[] = []; + + constructor() { + MockWebSocket.instances.push(this); + } + + on(event: string, listener: Listener): this { + const listeners = this.listeners.get(event) ?? []; + listeners.push(listener); + this.listeners.set(event, listeners); + return this; + } + + emit(event: string, ...args: unknown[]): void { + for (const listener of this.listeners.get(event) ?? []) { + listener(...args); + } + } + + send(payload: string): void { + this.sent.push(payload); + } + + close(code?: number, reason?: string): void { + this.closeCalls += 1; + if (this.readyState === MockWebSocket.CLOSED) { + return; + } + this.readyState = MockWebSocket.CLOSED; + this.emit("close", code ?? 1000, Buffer.from(reason ?? "")); + } + } + + return { FakeWebSocket: MockWebSocket }; +}); + +vi.mock("ws", () => ({ + default: FakeWebSocket, +})); + +type FakeWebSocketInstance = InstanceType; + +function emitEvent(socket: FakeWebSocketInstance, event: unknown): void { + socket.emit("message", Buffer.from(JSON.stringify(event))); +} + +async function connectSession(callbacks: { + onError?: (error: Error) => void; + onPartial?: (partial: string) => void; + onTranscript?: (transcript: string) => void; +}) { + const session = buildMistralRealtimeTranscriptionProvider().createSession({ + providerConfig: { + apiKey: "fixture-value", + baseUrl: "ws://mistral.test", + }, + ...callbacks, + }); + const connecting = session.connect(); + let socket: FakeWebSocketInstance | undefined; + await vi.waitFor(() => { + socket = FakeWebSocket.instances[0]; + if (!socket) { + throw new Error("expected session to create a websocket"); + } + }); + if (!socket) { + throw new Error("expected session to create a websocket"); + } + socket.readyState = FakeWebSocket.OPEN; + socket.emit("open"); + emitEvent(socket, { type: "session.created" }); + await connecting; + return { session, socket }; +} + +describe("Mistral realtime transcription lifecycle", () => { + beforeEach(() => { + FakeWebSocket.instances = []; + }); + + it("preserves partial, segment, and done transcript semantics", async () => { + const errors: string[] = []; + const partials: string[] = []; + const transcripts: string[] = []; + const { session, socket } = await connectSession({ + onError: (error) => errors.push(error.message), + onPartial: (partial) => partials.push(partial), + onTranscript: (transcript) => transcripts.push(transcript), + }); + + emitEvent(socket, { type: "transcription.text.delta", text: "hel" }); + emitEvent(socket, { type: "transcription.text.delta", text: "lo" }); + emitEvent(socket, { type: "transcription.segment", text: "hello final" }); + emitEvent(socket, { type: "transcription.text.delta", text: "next" }); + emitEvent(socket, { + type: "transcription.done", + text: "provider full transcript remains ignored", + }); + + expect(partials).toEqual(["hel", "hello", "next"]); + expect(transcripts).toEqual(["hello final", "next"]); + expect(errors).toEqual([]); + expect(socket.closeCalls).toBe(1); + expect(session.isConnected()).toBe(false); + }); + + it("tracks the in-progress transcript limit as aggregate UTF-8 bytes", async () => { + const errors: string[] = []; + const transcripts: string[] = []; + const { socket } = await connectSession({ + onError: (error) => errors.push(error.message), + onTranscript: (transcript) => transcripts.push(transcript), + }); + const exactUtf8Limit = "🙂".repeat((256 * 1024) / 4); + const splitSurrogatePrefix = "x".repeat(256 * 1024 - 4); + const splitSurrogateTranscript = `${splitSurrogatePrefix}🙂`; + + emitEvent(socket, { type: "transcription.text.delta", text: exactUtf8Limit }); + emitEvent(socket, { type: "transcription.segment", text: "first segment" }); + emitEvent(socket, { + type: "transcription.text.delta", + text: `${splitSurrogatePrefix}\ud83d`, + }); + emitEvent(socket, { type: "transcription.text.delta", text: "\ude42" }); + emitEvent(socket, { type: "transcription.done" }); + + expect(errors).toEqual([]); + expect(transcripts).toEqual(["first segment", splitSurrogateTranscript]); + expect(socket.closeCalls).toBe(1); + }); + + it("fails once and ignores late terminal events after 10,000 runaway deltas", async () => { + const errors: string[] = []; + const transcripts: string[] = []; + let lastPartialLength = 0; + let partialCalls = 0; + const { session, socket } = await connectSession({ + onError: (error) => errors.push(error.message), + onPartial: (partial) => { + lastPartialLength = partial.length; + partialCalls += 1; + }, + onTranscript: (transcript) => transcripts.push(transcript), + }); + + for (let index = 0; index < 10_000; index += 1) { + emitEvent(socket, { type: "transcription.text.delta", text: "x".repeat(32) }); + } + emitEvent(socket, { type: "transcription.segment", text: "late segment" }); + emitEvent(socket, { type: "transcription.done", text: "late done" }); + + expect(errors).toEqual([ + "Mistral realtime transcription exceeded the 256 KiB in-progress transcript limit", + ]); + expect(partialCalls).toBe(8_192); + expect(lastPartialLength).toBe(256 * 1024); + expect(transcripts).toEqual([]); + expect(socket.closeCalls).toBe(1); + expect(session.isConnected()).toBe(false); + }); +}); From d351187e45aa7c971eb341ec4019ce9f24dd7b3f Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Fri, 31 Jul 2026 21:31:06 +0800 Subject: [PATCH 3/4] fix(mistral): make realtime errors terminal --- ...e-transcription-provider.lifecycle.test.ts | 178 ------------------ .../realtime-transcription-provider.test.ts | 102 ++++++++++ .../realtime-transcription-provider.ts | 13 +- 3 files changed, 110 insertions(+), 183 deletions(-) delete mode 100644 extensions/mistral/realtime-transcription-provider.lifecycle.test.ts diff --git a/extensions/mistral/realtime-transcription-provider.lifecycle.test.ts b/extensions/mistral/realtime-transcription-provider.lifecycle.test.ts deleted file mode 100644 index 67d2c56036d1..000000000000 --- a/extensions/mistral/realtime-transcription-provider.lifecycle.test.ts +++ /dev/null @@ -1,178 +0,0 @@ -// Mistral lifecycle tests cover bounded transcript accumulation and terminal events. -import { beforeEach, describe, expect, it, vi } from "vitest"; -import { buildMistralRealtimeTranscriptionProvider } from "./realtime-transcription-provider.js"; - -const { FakeWebSocket } = vi.hoisted(() => { - type Listener = (...args: unknown[]) => void; - - class MockWebSocket { - static readonly OPEN = 1; - static readonly CLOSED = 3; - static instances: MockWebSocket[] = []; - - binaryType = "nodebuffer"; - closeCalls = 0; - readonly listeners = new Map(); - readyState = 0; - sent: string[] = []; - - constructor() { - MockWebSocket.instances.push(this); - } - - on(event: string, listener: Listener): this { - const listeners = this.listeners.get(event) ?? []; - listeners.push(listener); - this.listeners.set(event, listeners); - return this; - } - - emit(event: string, ...args: unknown[]): void { - for (const listener of this.listeners.get(event) ?? []) { - listener(...args); - } - } - - send(payload: string): void { - this.sent.push(payload); - } - - close(code?: number, reason?: string): void { - this.closeCalls += 1; - if (this.readyState === MockWebSocket.CLOSED) { - return; - } - this.readyState = MockWebSocket.CLOSED; - this.emit("close", code ?? 1000, Buffer.from(reason ?? "")); - } - } - - return { FakeWebSocket: MockWebSocket }; -}); - -vi.mock("ws", () => ({ - default: FakeWebSocket, -})); - -type FakeWebSocketInstance = InstanceType; - -function emitEvent(socket: FakeWebSocketInstance, event: unknown): void { - socket.emit("message", Buffer.from(JSON.stringify(event))); -} - -async function connectSession(callbacks: { - onError?: (error: Error) => void; - onPartial?: (partial: string) => void; - onTranscript?: (transcript: string) => void; -}) { - const session = buildMistralRealtimeTranscriptionProvider().createSession({ - providerConfig: { - apiKey: "fixture-value", - baseUrl: "ws://mistral.test", - }, - ...callbacks, - }); - const connecting = session.connect(); - let socket: FakeWebSocketInstance | undefined; - await vi.waitFor(() => { - socket = FakeWebSocket.instances[0]; - if (!socket) { - throw new Error("expected session to create a websocket"); - } - }); - if (!socket) { - throw new Error("expected session to create a websocket"); - } - socket.readyState = FakeWebSocket.OPEN; - socket.emit("open"); - emitEvent(socket, { type: "session.created" }); - await connecting; - return { session, socket }; -} - -describe("Mistral realtime transcription lifecycle", () => { - beforeEach(() => { - FakeWebSocket.instances = []; - }); - - it("preserves partial, segment, and done transcript semantics", async () => { - const errors: string[] = []; - const partials: string[] = []; - const transcripts: string[] = []; - const { session, socket } = await connectSession({ - onError: (error) => errors.push(error.message), - onPartial: (partial) => partials.push(partial), - onTranscript: (transcript) => transcripts.push(transcript), - }); - - emitEvent(socket, { type: "transcription.text.delta", text: "hel" }); - emitEvent(socket, { type: "transcription.text.delta", text: "lo" }); - emitEvent(socket, { type: "transcription.segment", text: "hello final" }); - emitEvent(socket, { type: "transcription.text.delta", text: "next" }); - emitEvent(socket, { - type: "transcription.done", - text: "provider full transcript remains ignored", - }); - - expect(partials).toEqual(["hel", "hello", "next"]); - expect(transcripts).toEqual(["hello final", "next"]); - expect(errors).toEqual([]); - expect(socket.closeCalls).toBe(1); - expect(session.isConnected()).toBe(false); - }); - - it("tracks the in-progress transcript limit as aggregate UTF-8 bytes", async () => { - const errors: string[] = []; - const transcripts: string[] = []; - const { socket } = await connectSession({ - onError: (error) => errors.push(error.message), - onTranscript: (transcript) => transcripts.push(transcript), - }); - const exactUtf8Limit = "🙂".repeat((256 * 1024) / 4); - const splitSurrogatePrefix = "x".repeat(256 * 1024 - 4); - const splitSurrogateTranscript = `${splitSurrogatePrefix}🙂`; - - emitEvent(socket, { type: "transcription.text.delta", text: exactUtf8Limit }); - emitEvent(socket, { type: "transcription.segment", text: "first segment" }); - emitEvent(socket, { - type: "transcription.text.delta", - text: `${splitSurrogatePrefix}\ud83d`, - }); - emitEvent(socket, { type: "transcription.text.delta", text: "\ude42" }); - emitEvent(socket, { type: "transcription.done" }); - - expect(errors).toEqual([]); - expect(transcripts).toEqual(["first segment", splitSurrogateTranscript]); - expect(socket.closeCalls).toBe(1); - }); - - it("fails once and ignores late terminal events after 10,000 runaway deltas", async () => { - const errors: string[] = []; - const transcripts: string[] = []; - let lastPartialLength = 0; - let partialCalls = 0; - const { session, socket } = await connectSession({ - onError: (error) => errors.push(error.message), - onPartial: (partial) => { - lastPartialLength = partial.length; - partialCalls += 1; - }, - onTranscript: (transcript) => transcripts.push(transcript), - }); - - for (let index = 0; index < 10_000; index += 1) { - emitEvent(socket, { type: "transcription.text.delta", text: "x".repeat(32) }); - } - emitEvent(socket, { type: "transcription.segment", text: "late segment" }); - emitEvent(socket, { type: "transcription.done", text: "late done" }); - - expect(errors).toEqual([ - "Mistral realtime transcription exceeded the 256 KiB in-progress transcript limit", - ]); - expect(partialCalls).toBe(8_192); - expect(lastPartialLength).toBe(256 * 1024); - expect(transcripts).toEqual([]); - expect(socket.closeCalls).toBe(1); - expect(session.isConnected()).toBe(false); - }); -}); diff --git a/extensions/mistral/realtime-transcription-provider.test.ts b/extensions/mistral/realtime-transcription-provider.test.ts index 3d421c5bddd5..89a4cbeb8807 100644 --- a/extensions/mistral/realtime-transcription-provider.test.ts +++ b/extensions/mistral/realtime-transcription-provider.test.ts @@ -337,4 +337,106 @@ describe("buildMistralRealtimeTranscriptionProvider", () => { expect(onPartial.mock.calls.map(([text]) => text)).toEqual(partials); }); + + it("tracks the in-progress transcript limit as aggregate UTF-8 bytes", async () => { + const exactUtf8Limit = "🙂".repeat((256 * 1024) / 4); + const splitSurrogatePrefix = "x".repeat(256 * 1024 - 4); + const splitSurrogateTranscript = `${splitSurrogatePrefix}🙂`; + const baseUrl = await createRealtimeServer(() => {}, [ + { type: "transcription.text.delta", text: exactUtf8Limit }, + { type: "transcription.segment", text: "first segment", start: 0, end: 1 }, + { type: "transcription.text.delta", text: `${splitSurrogatePrefix}\ud83d` }, + { type: "transcription.text.delta", text: "\ude42" }, + { type: "transcription.done" }, + ]); + const onError = vi.fn(); + const onTranscript = vi.fn(); + const session = buildMistralRealtimeTranscriptionProvider().createSession({ + providerConfig: { apiKey: "fixture-value", baseUrl }, + onError, + onTranscript, + }); + + await session.connect(); + await vi.waitFor(() => { + expect(onTranscript.mock.calls.map(([text]) => text)).toEqual([ + "first segment", + splitSurrogateTranscript, + ]); + expect(session.isConnected()).toBe(false); + }); + + expect(onError).not.toHaveBeenCalled(); + }); + + it("fails once and ignores late terminal events after 10,000 runaway deltas", async () => { + const baseUrl = await createRealtimeServer(() => {}, [ + ...Array.from({ length: 10_000 }, () => ({ + type: "transcription.text.delta", + text: "x".repeat(32), + })), + { type: "transcription.segment", text: "late segment", start: 0, end: 1 }, + { type: "transcription.done", text: "late done" }, + ]); + const onError = vi.fn(); + const onTranscript = vi.fn(); + let lastPartialLength = 0; + let partialCalls = 0; + const session = buildMistralRealtimeTranscriptionProvider().createSession({ + providerConfig: { apiKey: "fixture-value", baseUrl }, + onError, + onPartial: (partial) => { + lastPartialLength = partial.length; + partialCalls += 1; + }, + onTranscript, + }); + + await session.connect(); + await vi.waitFor(() => { + expect(onError).toHaveBeenCalledExactlyOnceWith( + expect.objectContaining({ + message: + "Mistral realtime transcription exceeded the 256 KiB in-progress transcript limit", + }), + ); + expect(session.isConnected()).toBe(false); + }); + + expect(partialCalls).toBe(8_192); + expect(lastPartialLength).toBe(256 * 1024); + expect(onTranscript).not.toHaveBeenCalled(); + }); + + it("makes a ready-state provider error terminal and ignores late events", async () => { + const baseUrl = await createRealtimeServer(() => {}, [ + { type: "transcription.text.delta", text: "draft" }, + { type: "error", error: { message: "provider failed" } }, + { type: "transcription.text.delta", text: "x".repeat(256 * 1024 + 1) }, + { type: "transcription.segment", text: "late segment", start: 0, end: 1 }, + { type: "transcription.done", text: "late done" }, + ]); + const onError = vi.fn(() => { + throw new Error("observer failed"); + }); + const onPartial = vi.fn(); + const onTranscript = vi.fn(); + const session = buildMistralRealtimeTranscriptionProvider().createSession({ + providerConfig: { apiKey: "fixture-value", baseUrl }, + onError, + onPartial, + onTranscript, + }); + + await session.connect(); + await vi.waitFor(() => { + expect(onError).toHaveBeenCalledExactlyOnceWith( + expect.objectContaining({ message: "provider failed" }), + ); + expect(session.isConnected()).toBe(false); + }); + + expect(onPartial).toHaveBeenCalledExactlyOnceWith("draft"); + expect(onTranscript).not.toHaveBeenCalled(); + }); }); diff --git a/extensions/mistral/realtime-transcription-provider.ts b/extensions/mistral/realtime-transcription-provider.ts index 52e65c6d3f0a..ac7124cd604b 100644 --- a/extensions/mistral/realtime-transcription-provider.ts +++ b/extensions/mistral/realtime-transcription-provider.ts @@ -199,7 +199,7 @@ function createMistralRealtimeTranscriptionSession( config.onTranscript?.(text); }; - const failPartialOverflow = (transport: RealtimeTranscriptionWebSocketTransport) => { + const failTerminal = (error: Error, transport: RealtimeTranscriptionWebSocketTransport) => { if (terminal) { return; } @@ -207,7 +207,7 @@ function createMistralRealtimeTranscriptionSession( clearPartial(); transport.closeNow(); try { - config.onError?.(new Error(MISTRAL_REALTIME_PARTIAL_TRANSCRIPT_OVERFLOW_MESSAGE)); + config.onError?.(error); } catch { // The terminal provider error already owns the outcome. Do not let an // observer exception re-enter shared error dispatch and emit it twice. @@ -243,7 +243,10 @@ function createMistralRealtimeTranscriptionSession( if (event.text) { const deltaBytes = measureTranscriptDeltaBytes(partialText, event.text); if (deltaBytes > MISTRAL_REALTIME_MAX_PARTIAL_TRANSCRIPT_BYTES - partialBytes) { - failPartialOverflow(transport); + failTerminal( + new Error(MISTRAL_REALTIME_PARTIAL_TRANSCRIPT_OVERFLOW_MESSAGE), + transport, + ); return; } partialText += event.text; @@ -273,8 +276,8 @@ function createMistralRealtimeTranscriptionSession( return; } case "error": - config.onError?.(new Error(readErrorDetail(event))); - + failTerminal(new Error(readErrorDetail(event)), transport); + return; default: } }; From 834af948f4d9086b8c0d9a950150544d487dfe3c Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Sat, 1 Aug 2026 11:45:06 +0800 Subject: [PATCH 4/4] fix(mistral): keep terminal error branch lint-clean --- extensions/mistral/realtime-transcription-provider.ts | 1 - 1 file changed, 1 deletion(-) diff --git a/extensions/mistral/realtime-transcription-provider.ts b/extensions/mistral/realtime-transcription-provider.ts index ac7124cd604b..8fc7b5175e64 100644 --- a/extensions/mistral/realtime-transcription-provider.ts +++ b/extensions/mistral/realtime-transcription-provider.ts @@ -277,7 +277,6 @@ function createMistralRealtimeTranscriptionSession( } case "error": failTerminal(new Error(readErrorDetail(event)), transport); - return; default: } };