From c49afb7eec1c86a0b88ac13f12f453c4775dfdad Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Mon, 3 Aug 2026 04:24:36 +0800 Subject: [PATCH 01/13] fix(openai): bound realtime transcription state --- .../openai/realtime-transcription-provider.ts | 173 +++++++++++++++--- 1 file changed, 149 insertions(+), 24 deletions(-) diff --git a/extensions/openai/realtime-transcription-provider.ts b/extensions/openai/realtime-transcription-provider.ts index 1c99bad367e8..cd3519e76274 100644 --- a/extensions/openai/realtime-transcription-provider.ts +++ b/extensions/openai/realtime-transcription-provider.ts @@ -75,11 +75,35 @@ const OPENAI_REALTIME_TRANSCRIPTION_CONNECT_TIMEOUT_MS = 10_000; const OPENAI_REALTIME_TRANSCRIPTION_MAX_RECONNECT_ATTEMPTS = 5; const OPENAI_REALTIME_TRANSCRIPTION_RECONNECT_DELAY_MS = 1000; const OPENAI_REALTIME_TRANSCRIPTION_DEFAULT_MODEL = "gpt-4o-transcribe"; +const OPENAI_REALTIME_TRANSCRIPTION_MAX_UNRESOLVED_ITEMS = 64; +const OPENAI_REALTIME_TRANSCRIPTION_MAX_ITEM_ID_BYTES = 1024; +const OPENAI_REALTIME_TRANSCRIPTION_MAX_RETAINED_TRANSCRIPT_BYTES = 256 * 1024; +const OPENAI_REALTIME_TRANSCRIPTION_ITEM_OVERFLOW_MESSAGE = + "OpenAI realtime transcription exceeded the 64 unresolved item limit"; +const OPENAI_REALTIME_TRANSCRIPTION_IDENTITY_OVERFLOW_MESSAGE = + "OpenAI realtime transcription exceeded the 1024-byte item identity limit"; +const OPENAI_REALTIME_TRANSCRIPTION_TEXT_OVERFLOW_MESSAGE = + "OpenAI realtime transcription exceeded the 256 KiB retained transcript limit"; const OPENAI_REALTIME_TRANSCRIPTION_API_KEY_REQUIRED = "OpenAI Realtime transcription requires an OpenAI Platform API key"; const OPENAI_REALTIME_TRANSCRIPTION_API_KEY_REJECTED = "OpenAI Realtime transcription rejected the selected API key. Update or remove the active OpenAI API-key source"; +function appendedUtf8ByteLength(previous: string, appended: string): number { + const appendedBytes = Buffer.byteLength(appended, "utf8"); + if (!previous || !appended) { + return appendedBytes; + } + const previousCodeUnit = previous.charCodeAt(previous.length - 1); + const appendedCodeUnit = appended.charCodeAt(0); + const joinsSurrogatePair = + previousCodeUnit >= 0xd800 && + previousCodeUnit <= 0xdbff && + appendedCodeUnit >= 0xdc00 && + appendedCodeUnit <= 0xdfff; + return joinsSurrogatePair ? appendedBytes - 2 : appendedBytes; +} + function normalizeProviderConfig( config: RealtimeTranscriptionProviderConfig, ): OpenAIRealtimeTranscriptionProviderConfig { @@ -172,35 +196,80 @@ async function resolveOpenAIRealtimeTranscriptionAuthorization( function createOpenAIRealtimeTranscriptionSession( config: OpenAIRealtimeTranscriptionSessionConfig, ): RealtimeTranscriptionSession { - const pendingTranscripts = new Map(); + const pendingTranscripts = new Map(); const committedItemIds: string[] = []; - const committedItems = new Set(); - const previousItemIds = new Map(); - const settledItemIds = new Set(); + const committedItems = new Map(); const completedTranscripts = new Map(); + const trackedItemIds = new Set(); + const settledItemIds = new Set(); const unkeyedTranscript = "__openclaw_unkeyed_transcript__"; + let retainedTranscriptBytes = 0; const resetTranscriptionState = () => { pendingTranscripts.clear(); committedItemIds.length = 0; committedItems.clear(); - previousItemIds.clear(); - settledItemIds.clear(); completedTranscripts.clear(); + trackedItemIds.clear(); + settledItemIds.clear(); + retainedTranscriptBytes = 0; }; - const commitItem = (itemId: string, previousItemId: string | null | undefined) => { - if (committedItems.has(itemId)) { + const failTerminal = (error: Error, transport: RealtimeTranscriptionWebSocketTransport) => { + resetTranscriptionState(); + transport.closeNow(); + try { + config.onError?.(error); + } catch { + // The provider terminal owns this outcome; observer failures must not + // re-enter shared error dispatch or duplicate the terminal callback. + } + }; + + const trackItem = ( + itemId: string, + transport: RealtimeTranscriptionWebSocketTransport, + ): boolean => { + if (settledItemIds.has(itemId)) { + return false; + } + if (trackedItemIds.has(itemId)) { + return true; + } + if (trackedItemIds.size >= OPENAI_REALTIME_TRANSCRIPTION_MAX_UNRESOLVED_ITEMS) { + failTerminal(new Error(OPENAI_REALTIME_TRANSCRIPTION_ITEM_OVERFLOW_MESSAGE), transport); + return false; + } + if (Buffer.byteLength(itemId, "utf8") > OPENAI_REALTIME_TRANSCRIPTION_MAX_ITEM_ID_BYTES) { + failTerminal(new Error(OPENAI_REALTIME_TRANSCRIPTION_IDENTITY_OVERFLOW_MESSAGE), transport); + return false; + } + trackedItemIds.add(itemId); + return true; + }; + + const commitItem = ( + itemId: string, + previousItemId: string | null | undefined, + transport: RealtimeTranscriptionWebSocketTransport, + ) => { + if (settledItemIds.has(itemId) || committedItems.has(itemId) || !trackItem(itemId, transport)) { return; } - committedItems.add(itemId); - previousItemIds.set(itemId, previousItemId); + if ( + previousItemId && + Buffer.byteLength(previousItemId, "utf8") > OPENAI_REALTIME_TRANSCRIPTION_MAX_ITEM_ID_BYTES + ) { + failTerminal(new Error(OPENAI_REALTIME_TRANSCRIPTION_IDENTITY_OVERFLOW_MESSAGE), transport); + return; + } + committedItems.set(itemId, previousItemId); committedItemIds.push(itemId); const arrivalOrder = committedItemIds.splice(0); const successors = new Map(); for (const candidateId of arrivalOrder) { - const previousId = previousItemIds.get(candidateId); + const previousId = committedItems.get(candidateId); if (previousId) { successors.set(previousId, candidateId); } @@ -215,7 +284,7 @@ function createOpenAIRealtimeTranscriptionSession( } }; for (const candidateId of arrivalOrder) { - const previousId = previousItemIds.get(candidateId); + const previousId = committedItems.get(candidateId); if (previousId == null || settledItemIds.has(previousId)) { appendChain(candidateId); } @@ -231,7 +300,7 @@ function createOpenAIRealtimeTranscriptionSession( if (!itemId || !completedTranscripts.has(itemId)) { return; } - const previousItemId = previousItemIds.get(itemId); + const previousItemId = committedItems.get(itemId); if ( previousItemId && !settledItemIds.has(previousItemId) && @@ -241,27 +310,56 @@ function createOpenAIRealtimeTranscriptionSession( } committedItemIds.shift(); committedItems.delete(itemId); - previousItemIds.delete(itemId); + trackedItemIds.delete(itemId); + // The predecessor ledger is needed for ordering and duplicate rejection, + // but provider item IDs must not accumulate for the lifetime of a call. settledItemIds.add(itemId); + if (settledItemIds.size > OPENAI_REALTIME_TRANSCRIPTION_MAX_UNRESOLVED_ITEMS) { + const oldestSettledItemId = settledItemIds.values().next().value; + if (oldestSettledItemId) { + settledItemIds.delete(oldestSettledItemId); + } + } const transcript = completedTranscripts.get(itemId); completedTranscripts.delete(itemId); - pendingTranscripts.delete(itemId); + if (transcript) { + retainedTranscriptBytes -= Buffer.byteLength(transcript, "utf8"); + } if (transcript) { config.onTranscript?.(transcript); } } }; - const completeItem = (itemId: string | undefined, transcript: string | undefined) => { + const completeItem = ( + itemId: string | undefined, + transcript: string | undefined, + transport: RealtimeTranscriptionWebSocketTransport, + ) => { const key = itemId ?? unkeyedTranscript; + if (itemId && (settledItemIds.has(itemId) || completedTranscripts.has(itemId))) { + return; + } + const partialBytes = pendingTranscripts.get(key)?.bytes ?? 0; pendingTranscripts.delete(key); + retainedTranscriptBytes -= partialBytes; if (!itemId || !committedItems.has(itemId)) { + trackedItemIds.delete(key); if (transcript) { config.onTranscript?.(transcript); } return; } + const transcriptBytes = transcript ? Buffer.byteLength(transcript, "utf8") : 0; + if ( + transcriptBytes > + OPENAI_REALTIME_TRANSCRIPTION_MAX_RETAINED_TRANSCRIPT_BYTES - retainedTranscriptBytes + ) { + failTerminal(new Error(OPENAI_REALTIME_TRANSCRIPTION_TEXT_OVERFLOW_MESSAGE), transport); + return; + } completedTranscripts.set(itemId, transcript); + retainedTranscriptBytes += transcriptBytes; flushCompletedTranscripts(); }; @@ -277,32 +375,59 @@ function createOpenAIRealtimeTranscriptionSession( case "input_audio_buffer.committed": if (event.item_id) { - commitItem(event.item_id, event.previous_item_id); + commitItem(event.item_id, event.previous_item_id, transport); } return; case "conversation.item.input_audio_transcription.delta": if (event.delta) { const key = event.item_id ?? unkeyedTranscript; - const pendingTranscript = `${pendingTranscripts.get(key) ?? ""}${event.delta}`; - pendingTranscripts.set(key, pendingTranscript); - config.onPartial?.(pendingTranscript); + if (!trackItem(key, transport)) { + return; + } + const pendingTranscript = pendingTranscripts.get(key); + const previousPartial = pendingTranscript?.text ?? ""; + const deltaBytes = appendedUtf8ByteLength(previousPartial, event.delta); + if ( + deltaBytes > + OPENAI_REALTIME_TRANSCRIPTION_MAX_RETAINED_TRANSCRIPT_BYTES - retainedTranscriptBytes + ) { + failTerminal(new Error(OPENAI_REALTIME_TRANSCRIPTION_TEXT_OVERFLOW_MESSAGE), transport); + return; + } + const partial = `${previousPartial}${event.delta}`; + pendingTranscripts.set(key, { + bytes: (pendingTranscript?.bytes ?? 0) + deltaBytes, + text: partial, + }); + retainedTranscriptBytes += deltaBytes; + config.onPartial?.(partial); } return; case "conversation.item.input_audio_transcription.completed": - completeItem(event.item_id, event.transcript); + completeItem(event.item_id, event.transcript, transport); return; case "conversation.item.input_audio_transcription.failed": - completeItem(event.item_id, undefined); + if (event.item_id && settledItemIds.has(event.item_id)) { + return; + } + completeItem(event.item_id, undefined, transport); config.onError?.(new Error(readRealtimeErrorDetail(event.error))); return; - case "input_audio_buffer.speech_started": - pendingTranscripts.delete(event.item_id ?? unkeyedTranscript); + case "input_audio_buffer.speech_started": { + const key = event.item_id ?? unkeyedTranscript; + const partialBytes = pendingTranscripts.get(key)?.bytes ?? 0; + pendingTranscripts.delete(key); + retainedTranscriptBytes -= partialBytes; + if (!committedItems.has(key)) { + trackedItemIds.delete(key); + } config.onSpeechStart?.(); return; + } case "error": { const detail = readRealtimeErrorDetail(event.error); From 22db40f33780e8e8f5955308e72f8ba16a38b232 Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Mon, 3 Aug 2026 04:24:36 +0800 Subject: [PATCH 02/13] test(openai): cover realtime transcription bounds --- .../realtime-transcription-provider.test.ts | 331 +++++++++++++++++- 1 file changed, 329 insertions(+), 2 deletions(-) diff --git a/extensions/openai/realtime-transcription-provider.test.ts b/extensions/openai/realtime-transcription-provider.test.ts index 71cce622175f..7740b1bd78f2 100644 --- a/extensions/openai/realtime-transcription-provider.test.ts +++ b/extensions/openai/realtime-transcription-provider.test.ts @@ -45,6 +45,11 @@ const { FakeWebSocket, providerAuthMocks, ssrfMocks } = vi.hoisted(() => { this.readyState = MockWebSocket.CLOSED; this.emit("close", code ?? 1000, Buffer.from(reason ?? "")); } + + terminate(): void { + this.closed = true; + this.readyState = MockWebSocket.CLOSED; + } } return { @@ -83,10 +88,10 @@ function parseSent(socket: FakeWebSocketInstance): SentRealtimeEvent[] { return socket.sent.map((payload) => JSON.parse(payload) as SentRealtimeEvent); } -async function waitForFakeSocket(): Promise { +async function waitForFakeSocket(index = 0): Promise { let socket: FakeWebSocketInstance | undefined; await vi.waitFor(() => { - socket = FakeWebSocket.instances[0]; + socket = FakeWebSocket.instances[index]; if (!socket) { throw new Error("expected session to create a websocket"); } @@ -97,6 +102,23 @@ async function waitForFakeSocket(): Promise { return socket; } +function emitJson(socket: FakeWebSocketInstance, event: Record): void { + socket.emit("message", Buffer.from(JSON.stringify(event))); +} + +async function connectFakeSession( + session: { connect(): Promise }, + socketIndex = 0, +): Promise { + const connecting = session.connect(); + const socket = await waitForFakeSocket(socketIndex); + socket.readyState = FakeWebSocket.OPEN; + socket.emit("open"); + emitJson(socket, { type: "session.updated" }); + await connecting; + return socket; +} + function mockCallArg(mock: { mock: { calls: unknown[][] } }, index = 0): Record { const call = mock.mock.calls[index]; if (!call) { @@ -565,4 +587,309 @@ describe("buildOpenAIRealtimeTranscriptionProvider", () => { expect(transcripts).toEqual(["second final"]); session.close(); }); + + it("releases settled turns from the unresolved item budget", async () => { + const transcripts: string[] = []; + const onError = vi.fn(); + const provider = buildOpenAIRealtimeTranscriptionProvider(); + const session = provider.createSession({ + providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + onError, + onTranscript: (transcript) => transcripts.push(transcript), + }); + const socket = await connectFakeSession(session); + + for (let index = 0; index < 128; index += 1) { + const itemId = `item-${index}`; + emitJson(socket, { + type: "input_audio_buffer.committed", + item_id: itemId, + previous_item_id: index === 0 ? null : `item-${index - 1}`, + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: itemId, + transcript: `turn-${index}`, + }); + } + emitJson(socket, { + type: "input_audio_buffer.committed", + item_id: "item-129", + previous_item_id: "item-128", + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "item-129", + transcript: "turn-129", + }); + emitJson(socket, { + type: "input_audio_buffer.committed", + item_id: "item-128", + previous_item_id: "item-127", + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "item-128", + transcript: "turn-128", + }); + + expect(transcripts).toHaveLength(130); + expect(transcripts.slice(-2)).toEqual(["turn-128", "turn-129"]); + expect(onError).not.toHaveBeenCalled(); + expect(session.isConnected()).toBe(true); + session.close(); + }); + + it("fails once when unresolved item correlation exceeds its session bound", async () => { + const onError = vi.fn(); + const onTranscript = vi.fn(); + const provider = buildOpenAIRealtimeTranscriptionProvider(); + const session = provider.createSession({ + providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + onError, + onTranscript, + }); + const socket = await connectFakeSession(session); + + for (let index = 0; index < 65; index += 1) { + emitJson(socket, { + type: "input_audio_buffer.committed", + item_id: `item-${index}`, + previous_item_id: "missing-predecessor", + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: `item-${index}`, + transcript: `turn-${index}`, + }); + } + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "late-item", + transcript: "late transcript", + }); + + expect(onError).toHaveBeenCalledExactlyOnceWith( + expect.objectContaining({ + message: "OpenAI realtime transcription exceeded the 64 unresolved item limit", + }), + ); + expect(onTranscript).not.toHaveBeenCalled(); + expect(session.isConnected()).toBe(false); + + const reconnecting = session.connect(); + const replacementSocket = await waitForFakeSocket(1); + replacementSocket.readyState = FakeWebSocket.OPEN; + replacementSocket.emit("open"); + emitJson(replacementSocket, { type: "session.updated" }); + await reconnecting; + emitJson(replacementSocket, { + type: "input_audio_buffer.committed", + item_id: "replacement-item", + previous_item_id: null, + }); + emitJson(replacementSocket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "replacement-item", + transcript: "replacement transcript", + }); + + expect(onTranscript).toHaveBeenCalledExactlyOnceWith("replacement transcript"); + expect(onError).toHaveBeenCalledTimes(1); + session.close(); + session.close(); + }); + + it("fails once when aggregate in-progress transcript text exceeds 256 KiB", async () => { + const onError = vi.fn(); + const onPartial = vi.fn(); + const onTranscript = vi.fn(); + const provider = buildOpenAIRealtimeTranscriptionProvider(); + const session = provider.createSession({ + providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + onError, + onPartial, + onTranscript, + }); + const socket = await connectFakeSession(session); + const exactLimit = "🙂".repeat((256 * 1024) / 4); + + emitJson(socket, { + type: "conversation.item.input_audio_transcription.delta", + item_id: "item-1", + delta: exactLimit, + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.delta", + item_id: "item-1", + delta: "x", + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "item-1", + transcript: "late transcript", + }); + + expect(onPartial).toHaveBeenCalledExactlyOnceWith(exactLimit); + expect(onError).toHaveBeenCalledExactlyOnceWith( + expect.objectContaining({ + message: "OpenAI realtime transcription exceeded the 256 KiB retained transcript limit", + }), + ); + expect(onTranscript).not.toHaveBeenCalled(); + expect(session.isConnected()).toBe(false); + session.close(); + }); + + it("accounts for UTF-8 surrogate pairs split across delta frames", async () => { + const onError = vi.fn(); + const onPartial = vi.fn(); + const provider = buildOpenAIRealtimeTranscriptionProvider(); + const session = provider.createSession({ + providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + onError, + onPartial, + }); + const socket = await connectFakeSession(session); + const prefix = "x".repeat(256 * 1024 - 4); + + emitJson(socket, { + type: "conversation.item.input_audio_transcription.delta", + item_id: "item-1", + delta: prefix, + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.delta", + item_id: "item-1", + delta: "\ud83d", + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.delta", + item_id: "item-1", + delta: "\ude42", + }); + + expect(onPartial).toHaveBeenLastCalledWith(`${prefix}🙂`); + expect(onError).not.toHaveBeenCalled(); + expect(session.isConnected()).toBe(true); + session.close(); + }); + + it("ignores duplicate completion events without double-charging retained text", async () => { + const onError = vi.fn(); + const transcripts: string[] = []; + const provider = buildOpenAIRealtimeTranscriptionProvider(); + const session = provider.createSession({ + providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + onError, + onTranscript: (transcript) => transcripts.push(transcript), + }); + const socket = await connectFakeSession(session); + const secondTranscript = "x".repeat(192 * 1024); + + emitJson(socket, { + type: "input_audio_buffer.committed", + item_id: "item-2", + previous_item_id: "item-1", + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "item-2", + transcript: secondTranscript, + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "item-2", + transcript: secondTranscript, + }); + emitJson(socket, { + type: "input_audio_buffer.committed", + item_id: "item-1", + previous_item_id: null, + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "item-1", + transcript: "first", + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "item-2", + transcript: "late duplicate", + }); + + expect(transcripts).toEqual(["first", secondTranscript]); + expect(onError).not.toHaveBeenCalled(); + expect(session.isConnected()).toBe(true); + session.close(); + }); + + it("fails before retaining oversized correlation identities", async () => { + const onError = vi.fn(); + const onTranscript = vi.fn(); + const provider = buildOpenAIRealtimeTranscriptionProvider(); + const session = provider.createSession({ + providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + onError, + onTranscript, + }); + const socket = await connectFakeSession(session); + + emitJson(socket, { + type: "input_audio_buffer.committed", + item_id: "i".repeat(1025), + previous_item_id: null, + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "late-item", + transcript: "late transcript", + }); + + expect(onError).toHaveBeenCalledExactlyOnceWith( + expect.objectContaining({ + message: "OpenAI realtime transcription exceeded the 1024-byte item identity limit", + }), + ); + expect(onTranscript).not.toHaveBeenCalled(); + expect(session.isConnected()).toBe(false); + session.close(); + }); + + it("fails before retaining an oversized completed transcript", async () => { + const onError = vi.fn(); + const onTranscript = vi.fn(); + const provider = buildOpenAIRealtimeTranscriptionProvider(); + const session = provider.createSession({ + providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + onError, + onTranscript, + }); + const socket = await connectFakeSession(session); + + emitJson(socket, { + type: "input_audio_buffer.committed", + item_id: "item-1", + previous_item_id: null, + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "item-1", + transcript: "x".repeat(256 * 1024 + 1), + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "item-1", + transcript: "late transcript", + }); + + expect(onError).toHaveBeenCalledExactlyOnceWith( + expect.objectContaining({ + message: "OpenAI realtime transcription exceeded the 256 KiB retained transcript limit", + }), + ); + expect(onTranscript).not.toHaveBeenCalled(); + expect(session.isConnected()).toBe(false); + session.close(); + }); }); From 26b04e768ecb81b6fd4b93a5b8e249f18874f059 Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Mon, 3 Aug 2026 04:25:07 +0800 Subject: [PATCH 03/13] fix(openai): ignore post-terminal transcript events --- extensions/openai/realtime-transcription-provider.ts | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/extensions/openai/realtime-transcription-provider.ts b/extensions/openai/realtime-transcription-provider.ts index cd3519e76274..396b6b75ed19 100644 --- a/extensions/openai/realtime-transcription-provider.ts +++ b/extensions/openai/realtime-transcription-provider.ts @@ -230,7 +230,7 @@ function createOpenAIRealtimeTranscriptionSession( itemId: string, transport: RealtimeTranscriptionWebSocketTransport, ): boolean => { - if (settledItemIds.has(itemId)) { + if (settledItemIds.has(itemId) || completedTranscripts.has(itemId)) { return false; } if (trackedItemIds.has(itemId)) { @@ -410,7 +410,10 @@ function createOpenAIRealtimeTranscriptionSession( return; case "conversation.item.input_audio_transcription.failed": - if (event.item_id && settledItemIds.has(event.item_id)) { + if ( + event.item_id && + (settledItemIds.has(event.item_id) || completedTranscripts.has(event.item_id)) + ) { return; } completeItem(event.item_id, undefined, transport); From 53055bb3340e16b5fcec003fd02be565e3312d9e Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Mon, 3 Aug 2026 04:25:07 +0800 Subject: [PATCH 04/13] test(openai): pin first transcript terminal outcome --- .../realtime-transcription-provider.test.ts | 75 +++++++++++++++++++ 1 file changed, 75 insertions(+) diff --git a/extensions/openai/realtime-transcription-provider.test.ts b/extensions/openai/realtime-transcription-provider.test.ts index 7740b1bd78f2..48d67d3ff749 100644 --- a/extensions/openai/realtime-transcription-provider.test.ts +++ b/extensions/openai/realtime-transcription-provider.test.ts @@ -777,11 +777,13 @@ describe("buildOpenAIRealtimeTranscriptionProvider", () => { it("ignores duplicate completion events without double-charging retained text", async () => { const onError = vi.fn(); + const onPartial = vi.fn(); const transcripts: string[] = []; const provider = buildOpenAIRealtimeTranscriptionProvider(); const session = provider.createSession({ providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret onError, + onPartial, onTranscript: (transcript) => transcripts.push(transcript), }); const socket = await connectFakeSession(session); @@ -797,6 +799,16 @@ describe("buildOpenAIRealtimeTranscriptionProvider", () => { item_id: "item-2", transcript: secondTranscript, }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.delta", + item_id: "item-2", + delta: "late partial", + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.failed", + item_id: "item-2", + error: { message: "late failure" }, + }); emitJson(socket, { type: "conversation.item.input_audio_transcription.completed", item_id: "item-2", @@ -819,8 +831,71 @@ describe("buildOpenAIRealtimeTranscriptionProvider", () => { }); expect(transcripts).toEqual(["first", secondTranscript]); + expect(onPartial).not.toHaveBeenCalled(); expect(onError).not.toHaveBeenCalled(); expect(session.isConnected()).toBe(true); + + const reconnecting = session.connect(); + const replacementSocket = await waitForFakeSocket(1); + replacementSocket.readyState = FakeWebSocket.OPEN; + replacementSocket.emit("open"); + emitJson(replacementSocket, { type: "session.updated" }); + await reconnecting; + emitJson(replacementSocket, { + type: "input_audio_buffer.committed", + item_id: "item-2", + previous_item_id: null, + }); + emitJson(replacementSocket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "item-2", + transcript: "new session", + }); + + expect(transcripts).toEqual(["first", secondTranscript, "new session"]); + session.close(); + }); + + it("keeps the first failed terminal outcome when completion arrives late", async () => { + const errors: string[] = []; + const transcripts: string[] = []; + const provider = buildOpenAIRealtimeTranscriptionProvider(); + const session = provider.createSession({ + providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + onError: (error) => errors.push(error.message), + onTranscript: (transcript) => transcripts.push(transcript), + }); + const socket = await connectFakeSession(session); + + emitJson(socket, { + type: "input_audio_buffer.committed", + item_id: "item-2", + previous_item_id: "item-1", + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.failed", + item_id: "item-2", + error: { message: "second failed" }, + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "item-2", + transcript: "late second", + }); + emitJson(socket, { + type: "input_audio_buffer.committed", + item_id: "item-1", + previous_item_id: null, + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "item-1", + transcript: "first", + }); + + expect(errors).toEqual(["second failed"]); + expect(transcripts).toEqual(["first"]); + expect(session.isConnected()).toBe(true); session.close(); }); From 72b4e42b445103a06f287b95e3e57790bdcb9909 Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Mon, 3 Aug 2026 05:09:25 +0800 Subject: [PATCH 05/13] fix(openai): tombstone pre-commit terminal events --- .../openai/realtime-transcription-provider.ts | 32 ++++++++++++------- 1 file changed, 20 insertions(+), 12 deletions(-) diff --git a/extensions/openai/realtime-transcription-provider.ts b/extensions/openai/realtime-transcription-provider.ts index 396b6b75ed19..36a32dfe05c5 100644 --- a/extensions/openai/realtime-transcription-provider.ts +++ b/extensions/openai/realtime-transcription-provider.ts @@ -248,6 +248,19 @@ function createOpenAIRealtimeTranscriptionSession( return true; }; + const settleItem = (itemId: string) => { + trackedItemIds.delete(itemId); + // Keep only a bounded terminal frontier so late provider events cannot + // recreate released state, including when completion precedes commit. + settledItemIds.add(itemId); + if (settledItemIds.size > OPENAI_REALTIME_TRANSCRIPTION_MAX_UNRESOLVED_ITEMS) { + const oldestSettledItemId = settledItemIds.values().next().value; + if (oldestSettledItemId) { + settledItemIds.delete(oldestSettledItemId); + } + } + }; + const commitItem = ( itemId: string, previousItemId: string | null | undefined, @@ -310,16 +323,7 @@ function createOpenAIRealtimeTranscriptionSession( } committedItemIds.shift(); committedItems.delete(itemId); - trackedItemIds.delete(itemId); - // The predecessor ledger is needed for ordering and duplicate rejection, - // but provider item IDs must not accumulate for the lifetime of a call. - settledItemIds.add(itemId); - if (settledItemIds.size > OPENAI_REALTIME_TRANSCRIPTION_MAX_UNRESOLVED_ITEMS) { - const oldestSettledItemId = settledItemIds.values().next().value; - if (oldestSettledItemId) { - settledItemIds.delete(oldestSettledItemId); - } - } + settleItem(itemId); const transcript = completedTranscripts.get(itemId); completedTranscripts.delete(itemId); if (transcript) { @@ -337,14 +341,18 @@ function createOpenAIRealtimeTranscriptionSession( transport: RealtimeTranscriptionWebSocketTransport, ) => { const key = itemId ?? unkeyedTranscript; - if (itemId && (settledItemIds.has(itemId) || completedTranscripts.has(itemId))) { + if (itemId && !trackItem(itemId, transport)) { return; } const partialBytes = pendingTranscripts.get(key)?.bytes ?? 0; pendingTranscripts.delete(key); retainedTranscriptBytes -= partialBytes; if (!itemId || !committedItems.has(itemId)) { - trackedItemIds.delete(key); + if (itemId) { + settleItem(itemId); + } else { + trackedItemIds.delete(key); + } if (transcript) { config.onTranscript?.(transcript); } From 8f83eaccf47ce6e25f8d8dae070d9f3ac097389b Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Mon, 3 Aug 2026 05:09:26 +0800 Subject: [PATCH 06/13] test(openai): cover pre-commit terminal events --- .../realtime-transcription-provider.test.ts | 72 +++++++++++++++++++ 1 file changed, 72 insertions(+) diff --git a/extensions/openai/realtime-transcription-provider.test.ts b/extensions/openai/realtime-transcription-provider.test.ts index 48d67d3ff749..e05519882f30 100644 --- a/extensions/openai/realtime-transcription-provider.test.ts +++ b/extensions/openai/realtime-transcription-provider.test.ts @@ -856,6 +856,78 @@ describe("buildOpenAIRealtimeTranscriptionProvider", () => { session.close(); }); + it("tombstones terminal outcomes received before their commit", async () => { + const errors: string[] = []; + const partials: string[] = []; + const transcripts: string[] = []; + const provider = buildOpenAIRealtimeTranscriptionProvider(); + const session = provider.createSession({ + providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + onError: (error) => errors.push(error.message), + onPartial: (partial) => partials.push(partial), + onTranscript: (transcript) => transcripts.push(transcript), + }); + const socket = await connectFakeSession(session); + + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "item-completed", + transcript: "first", + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "item-completed", + transcript: "duplicate", + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.delta", + item_id: "item-completed", + delta: "late partial", + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.failed", + item_id: "item-completed", + error: { message: "late failure" }, + }); + emitJson(socket, { + type: "input_audio_buffer.committed", + item_id: "item-completed", + previous_item_id: null, + }); + + emitJson(socket, { + type: "conversation.item.input_audio_transcription.failed", + item_id: "item-failed", + error: { message: "first failure" }, + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "item-failed", + transcript: "late completion", + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.delta", + item_id: "item-failed", + delta: "late partial", + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.failed", + item_id: "item-failed", + error: { message: "duplicate failure" }, + }); + emitJson(socket, { + type: "input_audio_buffer.committed", + item_id: "item-failed", + previous_item_id: null, + }); + + expect(transcripts).toEqual(["first"]); + expect(errors).toEqual(["first failure"]); + expect(partials).toEqual([]); + expect(session.isConnected()).toBe(true); + session.close(); + }); + it("keeps the first failed terminal outcome when completion arrives late", async () => { const errors: string[] = []; const transcripts: string[] = []; From a0db9cdd15a1e57b51d0589c0645409de1893a89 Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Mon, 3 Aug 2026 05:10:33 +0800 Subject: [PATCH 07/13] fix(openai): preserve terminal error precedence --- .../openai/realtime-transcription-provider.ts | 14 ++++++++------ 1 file changed, 8 insertions(+), 6 deletions(-) diff --git a/extensions/openai/realtime-transcription-provider.ts b/extensions/openai/realtime-transcription-provider.ts index 36a32dfe05c5..9a1638f61564 100644 --- a/extensions/openai/realtime-transcription-provider.ts +++ b/extensions/openai/realtime-transcription-provider.ts @@ -339,10 +339,10 @@ function createOpenAIRealtimeTranscriptionSession( itemId: string | undefined, transcript: string | undefined, transport: RealtimeTranscriptionWebSocketTransport, - ) => { + ): boolean => { const key = itemId ?? unkeyedTranscript; if (itemId && !trackItem(itemId, transport)) { - return; + return false; } const partialBytes = pendingTranscripts.get(key)?.bytes ?? 0; pendingTranscripts.delete(key); @@ -356,7 +356,7 @@ function createOpenAIRealtimeTranscriptionSession( if (transcript) { config.onTranscript?.(transcript); } - return; + return true; } const transcriptBytes = transcript ? Buffer.byteLength(transcript, "utf8") : 0; if ( @@ -364,11 +364,12 @@ function createOpenAIRealtimeTranscriptionSession( OPENAI_REALTIME_TRANSCRIPTION_MAX_RETAINED_TRANSCRIPT_BYTES - retainedTranscriptBytes ) { failTerminal(new Error(OPENAI_REALTIME_TRANSCRIPTION_TEXT_OVERFLOW_MESSAGE), transport); - return; + return false; } completedTranscripts.set(itemId, transcript); retainedTranscriptBytes += transcriptBytes; flushCompletedTranscripts(); + return true; }; const handleEvent = ( @@ -424,8 +425,9 @@ function createOpenAIRealtimeTranscriptionSession( ) { return; } - completeItem(event.item_id, undefined, transport); - config.onError?.(new Error(readRealtimeErrorDetail(event.error))); + if (completeItem(event.item_id, undefined, transport)) { + config.onError?.(new Error(readRealtimeErrorDetail(event.error))); + } return; case "input_audio_buffer.speech_started": { From 26cd629d1b76f4358d5e07aa738fe114722e9043 Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Mon, 3 Aug 2026 05:10:33 +0800 Subject: [PATCH 08/13] test(openai): cover rejected terminal failures --- .../realtime-transcription-provider.test.ts | 17 ++++++----------- 1 file changed, 6 insertions(+), 11 deletions(-) diff --git a/extensions/openai/realtime-transcription-provider.test.ts b/extensions/openai/realtime-transcription-provider.test.ts index e05519882f30..c38dbaebdfb1 100644 --- a/extensions/openai/realtime-transcription-provider.test.ts +++ b/extensions/openai/realtime-transcription-provider.test.ts @@ -651,7 +651,7 @@ describe("buildOpenAIRealtimeTranscriptionProvider", () => { }); const socket = await connectFakeSession(session); - for (let index = 0; index < 65; index += 1) { + for (let index = 0; index < 64; index += 1) { emitJson(socket, { type: "input_audio_buffer.committed", item_id: `item-${index}`, @@ -664,9 +664,9 @@ describe("buildOpenAIRealtimeTranscriptionProvider", () => { }); } emitJson(socket, { - type: "conversation.item.input_audio_transcription.completed", - item_id: "late-item", - transcript: "late transcript", + type: "conversation.item.input_audio_transcription.failed", + item_id: "overflow-item", + error: { message: "provider failure" }, }); expect(onError).toHaveBeenCalledExactlyOnceWith( @@ -983,14 +983,9 @@ describe("buildOpenAIRealtimeTranscriptionProvider", () => { const socket = await connectFakeSession(session); emitJson(socket, { - type: "input_audio_buffer.committed", + type: "conversation.item.input_audio_transcription.failed", item_id: "i".repeat(1025), - previous_item_id: null, - }); - emitJson(socket, { - type: "conversation.item.input_audio_transcription.completed", - item_id: "late-item", - transcript: "late transcript", + error: { message: "provider failure" }, }); expect(onError).toHaveBeenCalledExactlyOnceWith( From 27f6005682bb364e7bcbc069ac4e4589aef35313 Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Mon, 3 Aug 2026 05:12:21 +0800 Subject: [PATCH 09/13] fix(openai): retain active predecessor satisfaction --- extensions/openai/realtime-transcription-provider.ts | 12 +++++++++++- 1 file changed, 11 insertions(+), 1 deletion(-) diff --git a/extensions/openai/realtime-transcription-provider.ts b/extensions/openai/realtime-transcription-provider.ts index 9a1638f61564..402d02a28940 100644 --- a/extensions/openai/realtime-transcription-provider.ts +++ b/extensions/openai/realtime-transcription-provider.ts @@ -250,6 +250,13 @@ function createOpenAIRealtimeTranscriptionSession( const settleItem = (itemId: string) => { trackedItemIds.delete(itemId); + // Predecessor satisfaction belongs to each active item. Recording it here + // prevents bounded tombstone eviction from invalidating admitted state. + for (const [candidateId, previousItemId] of committedItems) { + if (previousItemId === itemId) { + committedItems.set(candidateId, null); + } + } // Keep only a bounded terminal frontier so late provider events cannot // recreate released state, including when completion precedes commit. settledItemIds.add(itemId); @@ -276,7 +283,10 @@ function createOpenAIRealtimeTranscriptionSession( failTerminal(new Error(OPENAI_REALTIME_TRANSCRIPTION_IDENTITY_OVERFLOW_MESSAGE), transport); return; } - committedItems.set(itemId, previousItemId); + committedItems.set( + itemId, + previousItemId && settledItemIds.has(previousItemId) ? null : previousItemId, + ); committedItemIds.push(itemId); const arrivalOrder = committedItemIds.splice(0); From 05ddb8ab116800553fa8da2908e24ceed23e2111 Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Mon, 3 Aug 2026 05:12:21 +0800 Subject: [PATCH 10/13] test(openai): cover tombstone eviction ordering --- .../realtime-transcription-provider.test.ts | 42 +++++++++++++++++++ 1 file changed, 42 insertions(+) diff --git a/extensions/openai/realtime-transcription-provider.test.ts b/extensions/openai/realtime-transcription-provider.test.ts index c38dbaebdfb1..1c0936713ef0 100644 --- a/extensions/openai/realtime-transcription-provider.test.ts +++ b/extensions/openai/realtime-transcription-provider.test.ts @@ -928,6 +928,48 @@ describe("buildOpenAIRealtimeTranscriptionProvider", () => { session.close(); }); + it("keeps active predecessor satisfaction after tombstone eviction", async () => { + const transcripts: string[] = []; + const provider = buildOpenAIRealtimeTranscriptionProvider(); + const session = provider.createSession({ + providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + onTranscript: (transcript) => transcripts.push(transcript), + }); + const socket = await connectFakeSession(session); + + emitJson(socket, { + type: "input_audio_buffer.committed", + item_id: "root", + previous_item_id: null, + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "root", + transcript: "root transcript", + }); + emitJson(socket, { + type: "input_audio_buffer.committed", + item_id: "waiting", + previous_item_id: "root", + }); + for (let index = 0; index < 64; index += 1) { + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: `uncommitted-${index}`, + transcript: "", + }); + } + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "waiting", + transcript: "waiting transcript", + }); + + expect(transcripts).toEqual(["root transcript", "waiting transcript"]); + expect(session.isConnected()).toBe(true); + session.close(); + }); + it("keeps the first failed terminal outcome when completion arrives late", async () => { const errors: string[] = []; const transcripts: string[] = []; From 05541a2897501f338cd491eb62ba231892ece9e8 Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Mon, 3 Aug 2026 05:35:21 +0800 Subject: [PATCH 11/13] fix(openai): preserve terminal history precedence --- .../openai/realtime-transcription-provider.ts | 53 +++++++++++++------ 1 file changed, 36 insertions(+), 17 deletions(-) diff --git a/extensions/openai/realtime-transcription-provider.ts b/extensions/openai/realtime-transcription-provider.ts index 402d02a28940..2ea9ab6bd924 100644 --- a/extensions/openai/realtime-transcription-provider.ts +++ b/extensions/openai/realtime-transcription-provider.ts @@ -78,12 +78,16 @@ const OPENAI_REALTIME_TRANSCRIPTION_DEFAULT_MODEL = "gpt-4o-transcribe"; const OPENAI_REALTIME_TRANSCRIPTION_MAX_UNRESOLVED_ITEMS = 64; const OPENAI_REALTIME_TRANSCRIPTION_MAX_ITEM_ID_BYTES = 1024; const OPENAI_REALTIME_TRANSCRIPTION_MAX_RETAINED_TRANSCRIPT_BYTES = 256 * 1024; +const OPENAI_REALTIME_TRANSCRIPTION_MAX_SETTLED_ITEMS = 4096; +const OPENAI_REALTIME_TRANSCRIPTION_MAX_SETTLED_ID_BYTES = 256 * 1024; const OPENAI_REALTIME_TRANSCRIPTION_ITEM_OVERFLOW_MESSAGE = "OpenAI realtime transcription exceeded the 64 unresolved item limit"; const OPENAI_REALTIME_TRANSCRIPTION_IDENTITY_OVERFLOW_MESSAGE = "OpenAI realtime transcription exceeded the 1024-byte item identity limit"; const OPENAI_REALTIME_TRANSCRIPTION_TEXT_OVERFLOW_MESSAGE = "OpenAI realtime transcription exceeded the 256 KiB retained transcript limit"; +const OPENAI_REALTIME_TRANSCRIPTION_SETTLED_OVERFLOW_MESSAGE = + "OpenAI realtime transcription exceeded the terminal item history limit"; const OPENAI_REALTIME_TRANSCRIPTION_API_KEY_REQUIRED = "OpenAI Realtime transcription requires an OpenAI Platform API key"; const OPENAI_REALTIME_TRANSCRIPTION_API_KEY_REJECTED = @@ -204,6 +208,7 @@ function createOpenAIRealtimeTranscriptionSession( const settledItemIds = new Set(); const unkeyedTranscript = "__openclaw_unkeyed_transcript__"; let retainedTranscriptBytes = 0; + let settledItemIdBytes = 0; const resetTranscriptionState = () => { pendingTranscripts.clear(); @@ -213,6 +218,7 @@ function createOpenAIRealtimeTranscriptionSession( trackedItemIds.clear(); settledItemIds.clear(); retainedTranscriptBytes = 0; + settledItemIdBytes = 0; }; const failTerminal = (error: Error, transport: RealtimeTranscriptionWebSocketTransport) => { @@ -248,24 +254,31 @@ function createOpenAIRealtimeTranscriptionSession( return true; }; - const settleItem = (itemId: string) => { + const settleItem = ( + itemId: string, + transport: RealtimeTranscriptionWebSocketTransport, + ): boolean => { + const itemIdBytes = Buffer.byteLength(itemId, "utf8"); + if ( + settledItemIds.size >= OPENAI_REALTIME_TRANSCRIPTION_MAX_SETTLED_ITEMS || + itemIdBytes > OPENAI_REALTIME_TRANSCRIPTION_MAX_SETTLED_ID_BYTES - settledItemIdBytes + ) { + failTerminal(new Error(OPENAI_REALTIME_TRANSCRIPTION_SETTLED_OVERFLOW_MESSAGE), transport); + return false; + } trackedItemIds.delete(itemId); // Predecessor satisfaction belongs to each active item. Recording it here - // prevents bounded tombstone eviction from invalidating admitted state. + // prevents terminal-history saturation from invalidating admitted state. for (const [candidateId, previousItemId] of committedItems) { if (previousItemId === itemId) { committedItems.set(candidateId, null); } } - // Keep only a bounded terminal frontier so late provider events cannot - // recreate released state, including when completion precedes commit. + // Never evict terminal identities within a connection generation. Closing + // at the bound preserves first-terminal precedence without unbounded state. settledItemIds.add(itemId); - if (settledItemIds.size > OPENAI_REALTIME_TRANSCRIPTION_MAX_UNRESOLVED_ITEMS) { - const oldestSettledItemId = settledItemIds.values().next().value; - if (oldestSettledItemId) { - settledItemIds.delete(oldestSettledItemId); - } - } + settledItemIdBytes += itemIdBytes; + return true; }; const commitItem = ( @@ -317,11 +330,13 @@ function createOpenAIRealtimeTranscriptionSession( } }; - const flushCompletedTranscripts = () => { + const flushCompletedTranscripts = ( + transport: RealtimeTranscriptionWebSocketTransport, + ): boolean => { while (committedItemIds.length > 0) { const itemId = committedItemIds[0]; if (!itemId || !completedTranscripts.has(itemId)) { - return; + return true; } const previousItemId = committedItems.get(itemId); if ( @@ -329,11 +344,13 @@ function createOpenAIRealtimeTranscriptionSession( !settledItemIds.has(previousItemId) && !committedItems.has(previousItemId) ) { - return; + return true; } committedItemIds.shift(); committedItems.delete(itemId); - settleItem(itemId); + if (!settleItem(itemId, transport)) { + return false; + } const transcript = completedTranscripts.get(itemId); completedTranscripts.delete(itemId); if (transcript) { @@ -343,6 +360,7 @@ function createOpenAIRealtimeTranscriptionSession( config.onTranscript?.(transcript); } } + return true; }; const completeItem = ( @@ -359,7 +377,9 @@ function createOpenAIRealtimeTranscriptionSession( retainedTranscriptBytes -= partialBytes; if (!itemId || !committedItems.has(itemId)) { if (itemId) { - settleItem(itemId); + if (!settleItem(itemId, transport)) { + return false; + } } else { trackedItemIds.delete(key); } @@ -378,8 +398,7 @@ function createOpenAIRealtimeTranscriptionSession( } completedTranscripts.set(itemId, transcript); retainedTranscriptBytes += transcriptBytes; - flushCompletedTranscripts(); - return true; + return flushCompletedTranscripts(transport); }; const handleEvent = ( From 6ea9489d08e7899787e6b6aaa2e42b586a403422 Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Mon, 3 Aug 2026 05:35:22 +0800 Subject: [PATCH 12/13] test(openai): cover terminal history saturation --- .../realtime-transcription-provider.test.ts | 133 +++++++++++++++++- 1 file changed, 132 insertions(+), 1 deletion(-) diff --git a/extensions/openai/realtime-transcription-provider.test.ts b/extensions/openai/realtime-transcription-provider.test.ts index 1c0936713ef0..4c38cf62e65b 100644 --- a/extensions/openai/realtime-transcription-provider.test.ts +++ b/extensions/openai/realtime-transcription-provider.test.ts @@ -928,7 +928,7 @@ describe("buildOpenAIRealtimeTranscriptionProvider", () => { session.close(); }); - it("keeps active predecessor satisfaction after tombstone eviction", async () => { + it("keeps active predecessor satisfaction as terminal history grows", async () => { const transcripts: string[] = []; const provider = buildOpenAIRealtimeTranscriptionProvider(); const session = provider.createSession({ @@ -970,6 +970,137 @@ describe("buildOpenAIRealtimeTranscriptionProvider", () => { session.close(); }); + it("does not re-admit a terminal item after the settled frontier fills", async () => { + const onError = vi.fn(); + const onPartial = vi.fn(); + const transcripts: string[] = []; + const provider = buildOpenAIRealtimeTranscriptionProvider(); + const session = provider.createSession({ + providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + onError, + onPartial, + onTranscript: (transcript) => transcripts.push(transcript), + }); + const socket = await connectFakeSession(session); + + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "oldest", + transcript: "first", + }); + for (let index = 0; index < 4095; index += 1) { + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: `settled-${index}`, + transcript: "", + }); + } + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "oldest", + transcript: "duplicate", + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.delta", + item_id: "oldest", + delta: "late partial", + }); + + expect(transcripts).toEqual(["first"]); + expect(onPartial).not.toHaveBeenCalled(); + expect(onError).not.toHaveBeenCalled(); + expect(session.isConnected()).toBe(true); + session.close(); + }); + + it("fails visibly instead of evicting terminal item history", async () => { + const onError = vi.fn(); + const onTranscript = vi.fn(); + const provider = buildOpenAIRealtimeTranscriptionProvider(); + const session = provider.createSession({ + providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + onError, + onTranscript, + }); + const socket = await connectFakeSession(session); + + for (let index = 0; index < 4096; index += 1) { + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: `settled-${index}`, + transcript: "", + }); + } + emitJson(socket, { + type: "input_audio_buffer.committed", + item_id: "overflow-item", + previous_item_id: null, + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.failed", + item_id: "overflow-item", + error: { message: "provider failure" }, + }); + + expect(onError).toHaveBeenCalledExactlyOnceWith( + expect.objectContaining({ + message: "OpenAI realtime transcription exceeded the terminal item history limit", + }), + ); + expect(onTranscript).not.toHaveBeenCalled(); + expect(session.isConnected()).toBe(false); + + const reconnecting = session.connect(); + const replacementSocket = await waitForFakeSocket(1); + replacementSocket.readyState = FakeWebSocket.OPEN; + replacementSocket.emit("open"); + emitJson(replacementSocket, { type: "session.updated" }); + await reconnecting; + emitJson(replacementSocket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "replacement-item", + transcript: "replacement transcript", + }); + + expect(onTranscript).toHaveBeenCalledExactlyOnceWith("replacement transcript"); + expect(onError).toHaveBeenCalledTimes(1); + session.close(); + }); + + it("fails visibly when terminal item identities exceed 256 KiB", async () => { + const onError = vi.fn(); + const onTranscript = vi.fn(); + const provider = buildOpenAIRealtimeTranscriptionProvider(); + const session = provider.createSession({ + providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + onError, + onTranscript, + }); + const socket = await connectFakeSession(session); + + for (let index = 0; index < 256; index += 1) { + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: `${index.toString().padStart(4, "0")}${"i".repeat(1020)}`, + transcript: "", + }); + } + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "overflow-item", + transcript: "must not emit", + }); + + expect(onError).toHaveBeenCalledExactlyOnceWith( + expect.objectContaining({ + message: "OpenAI realtime transcription exceeded the terminal item history limit", + }), + ); + expect(onTranscript).not.toHaveBeenCalled(); + expect(session.isConnected()).toBe(false); + session.close(); + }); + it("keeps the first failed terminal outcome when completion arrives late", async () => { const errors: string[] = []; const transcripts: string[] = []; From 15195713985f5b069f53bbb769a9dc621ed4c7e9 Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Mon, 3 Aug 2026 05:48:38 +0800 Subject: [PATCH 13/13] test(openai): isolate terminal history bounds --- ...time-transcription-provider.bounds.test.ts | 227 ++++++++++++++++++ .../realtime-transcription-provider.test.ts | 131 ---------- 2 files changed, 227 insertions(+), 131 deletions(-) create mode 100644 extensions/openai/realtime-transcription-provider.bounds.test.ts diff --git a/extensions/openai/realtime-transcription-provider.bounds.test.ts b/extensions/openai/realtime-transcription-provider.bounds.test.ts new file mode 100644 index 000000000000..3bb710870803 --- /dev/null +++ b/extensions/openai/realtime-transcription-provider.bounds.test.ts @@ -0,0 +1,227 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; +import { buildOpenAIRealtimeTranscriptionProvider } 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[] = []; + + readonly listeners = new Map(); + readyState = 0; + closed = false; + + 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(): void {} + + close(code?: number, reason?: string): void { + this.closed = true; + this.readyState = MockWebSocket.CLOSED; + this.emit("close", code ?? 1000, Buffer.from(reason ?? "")); + } + + terminate(): void { + this.closed = true; + this.readyState = MockWebSocket.CLOSED; + } + } + + return { FakeWebSocket: MockWebSocket }; +}); + +vi.mock("ws", () => ({ + default: FakeWebSocket, +})); + +type FakeWebSocketInstance = InstanceType; + +async function waitForFakeSocket(index = 0): Promise { + let socket: FakeWebSocketInstance | undefined; + await vi.waitFor(() => { + socket = FakeWebSocket.instances[index]; + if (!socket) { + throw new Error("expected session to create a websocket"); + } + }); + if (!socket) { + throw new Error("expected session to create a websocket"); + } + return socket; +} + +function emitJson(socket: FakeWebSocketInstance, event: Record): void { + socket.emit("message", Buffer.from(JSON.stringify(event))); +} + +async function connectFakeSession( + session: { connect(): Promise }, + socketIndex = 0, +): Promise { + const connecting = session.connect(); + const socket = await waitForFakeSocket(socketIndex); + socket.readyState = FakeWebSocket.OPEN; + socket.emit("open"); + emitJson(socket, { type: "session.updated" }); + await connecting; + return socket; +} + +describe("OpenAI realtime transcription terminal history bounds", () => { + beforeEach(() => { + FakeWebSocket.instances = []; + vi.stubEnv("OPENAI_API_KEY", ""); + }); + + afterEach(() => { + vi.unstubAllEnvs(); + }); + + it("does not re-admit a terminal item after the settled frontier fills", async () => { + const onError = vi.fn(); + const onPartial = vi.fn(); + const transcripts: string[] = []; + const provider = buildOpenAIRealtimeTranscriptionProvider(); + const session = provider.createSession({ + providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + onError, + onPartial, + onTranscript: (transcript) => transcripts.push(transcript), + }); + const socket = await connectFakeSession(session); + + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "oldest", + transcript: "first", + }); + for (let index = 0; index < 4095; index += 1) { + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: `settled-${index}`, + transcript: "", + }); + } + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "oldest", + transcript: "duplicate", + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.delta", + item_id: "oldest", + delta: "late partial", + }); + + expect(transcripts).toEqual(["first"]); + expect(onPartial).not.toHaveBeenCalled(); + expect(onError).not.toHaveBeenCalled(); + expect(session.isConnected()).toBe(true); + session.close(); + }); + + it("fails visibly instead of evicting terminal item history", async () => { + const onError = vi.fn(); + const onTranscript = vi.fn(); + const provider = buildOpenAIRealtimeTranscriptionProvider(); + const session = provider.createSession({ + providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + onError, + onTranscript, + }); + const socket = await connectFakeSession(session); + + for (let index = 0; index < 4096; index += 1) { + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: `settled-${index}`, + transcript: "", + }); + } + emitJson(socket, { + type: "input_audio_buffer.committed", + item_id: "overflow-item", + previous_item_id: null, + }); + emitJson(socket, { + type: "conversation.item.input_audio_transcription.failed", + item_id: "overflow-item", + error: { message: "provider failure" }, + }); + + expect(onError).toHaveBeenCalledExactlyOnceWith( + expect.objectContaining({ + message: "OpenAI realtime transcription exceeded the terminal item history limit", + }), + ); + expect(onTranscript).not.toHaveBeenCalled(); + expect(session.isConnected()).toBe(false); + + const reconnecting = session.connect(); + const replacementSocket = await waitForFakeSocket(1); + replacementSocket.readyState = FakeWebSocket.OPEN; + replacementSocket.emit("open"); + emitJson(replacementSocket, { type: "session.updated" }); + await reconnecting; + emitJson(replacementSocket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "replacement-item", + transcript: "replacement transcript", + }); + + expect(onTranscript).toHaveBeenCalledExactlyOnceWith("replacement transcript"); + expect(onError).toHaveBeenCalledTimes(1); + session.close(); + }); + + it("fails visibly when terminal item identities exceed 256 KiB", async () => { + const onError = vi.fn(); + const onTranscript = vi.fn(); + const provider = buildOpenAIRealtimeTranscriptionProvider(); + const session = provider.createSession({ + providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + onError, + onTranscript, + }); + const socket = await connectFakeSession(session); + + for (let index = 0; index < 256; index += 1) { + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: `${index.toString().padStart(4, "0")}${"i".repeat(1020)}`, + transcript: "", + }); + } + emitJson(socket, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "overflow-item", + transcript: "must not emit", + }); + + expect(onError).toHaveBeenCalledExactlyOnceWith( + expect.objectContaining({ + message: "OpenAI realtime transcription exceeded the terminal item history limit", + }), + ); + expect(onTranscript).not.toHaveBeenCalled(); + expect(session.isConnected()).toBe(false); + session.close(); + }); +}); diff --git a/extensions/openai/realtime-transcription-provider.test.ts b/extensions/openai/realtime-transcription-provider.test.ts index 4c38cf62e65b..d851d160e19f 100644 --- a/extensions/openai/realtime-transcription-provider.test.ts +++ b/extensions/openai/realtime-transcription-provider.test.ts @@ -970,137 +970,6 @@ describe("buildOpenAIRealtimeTranscriptionProvider", () => { session.close(); }); - it("does not re-admit a terminal item after the settled frontier fills", async () => { - const onError = vi.fn(); - const onPartial = vi.fn(); - const transcripts: string[] = []; - const provider = buildOpenAIRealtimeTranscriptionProvider(); - const session = provider.createSession({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onError, - onPartial, - onTranscript: (transcript) => transcripts.push(transcript), - }); - const socket = await connectFakeSession(session); - - emitJson(socket, { - type: "conversation.item.input_audio_transcription.completed", - item_id: "oldest", - transcript: "first", - }); - for (let index = 0; index < 4095; index += 1) { - emitJson(socket, { - type: "conversation.item.input_audio_transcription.completed", - item_id: `settled-${index}`, - transcript: "", - }); - } - emitJson(socket, { - type: "conversation.item.input_audio_transcription.completed", - item_id: "oldest", - transcript: "duplicate", - }); - emitJson(socket, { - type: "conversation.item.input_audio_transcription.delta", - item_id: "oldest", - delta: "late partial", - }); - - expect(transcripts).toEqual(["first"]); - expect(onPartial).not.toHaveBeenCalled(); - expect(onError).not.toHaveBeenCalled(); - expect(session.isConnected()).toBe(true); - session.close(); - }); - - it("fails visibly instead of evicting terminal item history", async () => { - const onError = vi.fn(); - const onTranscript = vi.fn(); - const provider = buildOpenAIRealtimeTranscriptionProvider(); - const session = provider.createSession({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onError, - onTranscript, - }); - const socket = await connectFakeSession(session); - - for (let index = 0; index < 4096; index += 1) { - emitJson(socket, { - type: "conversation.item.input_audio_transcription.completed", - item_id: `settled-${index}`, - transcript: "", - }); - } - emitJson(socket, { - type: "input_audio_buffer.committed", - item_id: "overflow-item", - previous_item_id: null, - }); - emitJson(socket, { - type: "conversation.item.input_audio_transcription.failed", - item_id: "overflow-item", - error: { message: "provider failure" }, - }); - - expect(onError).toHaveBeenCalledExactlyOnceWith( - expect.objectContaining({ - message: "OpenAI realtime transcription exceeded the terminal item history limit", - }), - ); - expect(onTranscript).not.toHaveBeenCalled(); - expect(session.isConnected()).toBe(false); - - const reconnecting = session.connect(); - const replacementSocket = await waitForFakeSocket(1); - replacementSocket.readyState = FakeWebSocket.OPEN; - replacementSocket.emit("open"); - emitJson(replacementSocket, { type: "session.updated" }); - await reconnecting; - emitJson(replacementSocket, { - type: "conversation.item.input_audio_transcription.completed", - item_id: "replacement-item", - transcript: "replacement transcript", - }); - - expect(onTranscript).toHaveBeenCalledExactlyOnceWith("replacement transcript"); - expect(onError).toHaveBeenCalledTimes(1); - session.close(); - }); - - it("fails visibly when terminal item identities exceed 256 KiB", async () => { - const onError = vi.fn(); - const onTranscript = vi.fn(); - const provider = buildOpenAIRealtimeTranscriptionProvider(); - const session = provider.createSession({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onError, - onTranscript, - }); - const socket = await connectFakeSession(session); - - for (let index = 0; index < 256; index += 1) { - emitJson(socket, { - type: "conversation.item.input_audio_transcription.completed", - item_id: `${index.toString().padStart(4, "0")}${"i".repeat(1020)}`, - transcript: "", - }); - } - emitJson(socket, { - type: "conversation.item.input_audio_transcription.completed", - item_id: "overflow-item", - transcript: "must not emit", - }); - - expect(onError).toHaveBeenCalledExactlyOnceWith( - expect.objectContaining({ - message: "OpenAI realtime transcription exceeded the terminal item history limit", - }), - ); - expect(onTranscript).not.toHaveBeenCalled(); - expect(session.isConnected()).toBe(false); - session.close(); - }); - it("keeps the first failed terminal outcome when completion arrives late", async () => { const errors: string[] = []; const transcripts: string[] = [];