From 42515c4f07ea3b02e191d30cf97865d4e6229ef0 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Sat, 25 Jul 2026 08:11:09 -0700 Subject: [PATCH] fix(openai): prevent realtime transcript turns from reordering (#113723) * fix(openai): preserve realtime transcript order * chore: defer realtime note to release * test(openai): align heartbeat prompt contract * test(openai): align direct context contract --- extensions/openai/index.test.ts | 4 +- .../realtime-transcription-provider.test.ts | 135 ++++++++++++++++++ .../openai/realtime-transcription-provider.ts | 120 +++++++++++++++- .../test-helpers/provider-runtime-contract.ts | 2 +- 4 files changed, 252 insertions(+), 9 deletions(-) diff --git a/extensions/openai/index.test.ts b/extensions/openai/index.test.ts index 56f397b68e11..62b4f543c67f 100644 --- a/extensions/openai/index.test.ts +++ b/extensions/openai/index.test.ts @@ -472,7 +472,9 @@ describe("openai plugin", () => { expect(OPENAI_HEARTBEAT_PROMPT_OVERLAY).toContain( "Heartbeat = useful proactive progress, not chatter.", ); - expect(OPENAI_HEARTBEAT_PROMPT_OVERLAY).toContain("Wake, orient, read HEARTBEAT.md, act."); + expect(OPENAI_HEARTBEAT_PROMPT_OVERLAY).toContain( + "Wake, orient, use the provided monitor scratch, act.", + ); expect(OPENAI_HEARTBEAT_PROMPT_OVERLAY).toContain( "Assigned/ongoing work: pursue spirit with judgment.", ); diff --git a/extensions/openai/realtime-transcription-provider.test.ts b/extensions/openai/realtime-transcription-provider.test.ts index 831c86d03e41..71cce622175f 100644 --- a/extensions/openai/realtime-transcription-provider.test.ts +++ b/extensions/openai/realtime-transcription-provider.test.ts @@ -430,4 +430,139 @@ describe("buildOpenAIRealtimeTranscriptionProvider", () => { ]); session.close(); }); + + it("keeps out-of-order transcription items isolated and emits finals in commit order", async () => { + const partials: string[] = []; + const transcripts: string[] = []; + const provider = buildOpenAIRealtimeTranscriptionProvider(); + const session = provider.createSession({ + providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + onPartial: (partial) => partials.push(partial), + onTranscript: (transcript) => transcripts.push(transcript), + }); + + const connecting = session.connect(); + const socket = await waitForFakeSocket(); + socket.readyState = FakeWebSocket.OPEN; + socket.emit("open"); + socket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + await connecting; + + socket.emit( + "message", + Buffer.from( + JSON.stringify({ + type: "input_audio_buffer.committed", + item_id: "item-2", + previous_item_id: "item-1", + }), + ), + ); + socket.emit( + "message", + Buffer.from( + JSON.stringify({ + type: "input_audio_buffer.committed", + item_id: "item-1", + previous_item_id: null, + }), + ), + ); + socket.emit( + "message", + Buffer.from( + JSON.stringify({ + type: "conversation.item.input_audio_transcription.delta", + item_id: "item-1", + delta: "first partial", + }), + ), + ); + socket.emit( + "message", + Buffer.from( + JSON.stringify({ + type: "conversation.item.input_audio_transcription.delta", + item_id: "item-2", + delta: "second partial", + }), + ), + ); + socket.emit( + "message", + Buffer.from( + JSON.stringify({ + type: "conversation.item.input_audio_transcription.completed", + item_id: "item-2", + transcript: "second final", + }), + ), + ); + + expect(partials).toEqual(["first partial", "second partial"]); + expect(transcripts).toEqual([]); + + socket.emit( + "message", + Buffer.from( + JSON.stringify({ + type: "conversation.item.input_audio_transcription.completed", + item_id: "item-1", + transcript: "first final", + }), + ), + ); + + expect(transcripts).toEqual(["first final", "second final"]); + session.close(); + }); + + it("reports failed transcription items without blocking later committed turns", 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 connecting = session.connect(); + const socket = await waitForFakeSocket(); + socket.readyState = FakeWebSocket.OPEN; + socket.emit("open"); + socket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + await connecting; + + for (const itemId of ["item-1", "item-2"]) { + socket.emit( + "message", + Buffer.from(JSON.stringify({ type: "input_audio_buffer.committed", item_id: itemId })), + ); + } + socket.emit( + "message", + Buffer.from( + JSON.stringify({ + type: "conversation.item.input_audio_transcription.completed", + item_id: "item-2", + transcript: "second final", + }), + ), + ); + socket.emit( + "message", + Buffer.from( + JSON.stringify({ + type: "conversation.item.input_audio_transcription.failed", + item_id: "item-1", + error: { message: "first turn failed" }, + }), + ), + ); + + expect(errors).toEqual(["first turn failed"]); + expect(transcripts).toEqual(["second final"]); + session.close(); + }); }); diff --git a/extensions/openai/realtime-transcription-provider.ts b/extensions/openai/realtime-transcription-provider.ts index e56db8ea6a73..1c99bad367e8 100644 --- a/extensions/openai/realtime-transcription-provider.ts +++ b/extensions/openai/realtime-transcription-provider.ts @@ -45,6 +45,8 @@ type RealtimeEvent = { type: string; delta?: string; transcript?: string; + item_id?: string; + previous_item_id?: string | null; error?: unknown; }; @@ -170,7 +172,98 @@ async function resolveOpenAIRealtimeTranscriptionAuthorization( function createOpenAIRealtimeTranscriptionSession( config: OpenAIRealtimeTranscriptionSessionConfig, ): RealtimeTranscriptionSession { - let pendingTranscript = ""; + const pendingTranscripts = new Map(); + const committedItemIds: string[] = []; + const committedItems = new Set(); + const previousItemIds = new Map(); + const settledItemIds = new Set(); + const completedTranscripts = new Map(); + const unkeyedTranscript = "__openclaw_unkeyed_transcript__"; + + const resetTranscriptionState = () => { + pendingTranscripts.clear(); + committedItemIds.length = 0; + committedItems.clear(); + previousItemIds.clear(); + settledItemIds.clear(); + completedTranscripts.clear(); + }; + + const commitItem = (itemId: string, previousItemId: string | null | undefined) => { + if (committedItems.has(itemId)) { + return; + } + committedItems.add(itemId); + previousItemIds.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); + if (previousId) { + successors.set(previousId, candidateId); + } + } + const seen = new Set(); + const appendChain = (startId: string) => { + let candidateId: string | undefined = startId; + while (candidateId && !seen.has(candidateId)) { + seen.add(candidateId); + committedItemIds.push(candidateId); + candidateId = successors.get(candidateId); + } + }; + for (const candidateId of arrivalOrder) { + const previousId = previousItemIds.get(candidateId); + if (previousId == null || settledItemIds.has(previousId)) { + appendChain(candidateId); + } + } + for (const candidateId of arrivalOrder) { + appendChain(candidateId); + } + }; + + const flushCompletedTranscripts = () => { + while (committedItemIds.length > 0) { + const itemId = committedItemIds[0]; + if (!itemId || !completedTranscripts.has(itemId)) { + return; + } + const previousItemId = previousItemIds.get(itemId); + if ( + previousItemId && + !settledItemIds.has(previousItemId) && + !committedItems.has(previousItemId) + ) { + return; + } + committedItemIds.shift(); + committedItems.delete(itemId); + previousItemIds.delete(itemId); + settledItemIds.add(itemId); + const transcript = completedTranscripts.get(itemId); + completedTranscripts.delete(itemId); + pendingTranscripts.delete(itemId); + if (transcript) { + config.onTranscript?.(transcript); + } + } + }; + + const completeItem = (itemId: string | undefined, transcript: string | undefined) => { + const key = itemId ?? unkeyedTranscript; + pendingTranscripts.delete(key); + if (!itemId || !committedItems.has(itemId)) { + if (transcript) { + config.onTranscript?.(transcript); + } + return; + } + completedTranscripts.set(itemId, transcript); + flushCompletedTranscripts(); + }; const handleEvent = ( event: RealtimeEvent, @@ -182,22 +275,32 @@ function createOpenAIRealtimeTranscriptionSession( transport.markReady(); return; + case "input_audio_buffer.committed": + if (event.item_id) { + commitItem(event.item_id, event.previous_item_id); + } + return; + case "conversation.item.input_audio_transcription.delta": if (event.delta) { - pendingTranscript += event.delta; + const key = event.item_id ?? unkeyedTranscript; + const pendingTranscript = `${pendingTranscripts.get(key) ?? ""}${event.delta}`; + pendingTranscripts.set(key, pendingTranscript); config.onPartial?.(pendingTranscript); } return; case "conversation.item.input_audio_transcription.completed": - if (event.transcript) { - config.onTranscript?.(event.transcript); - } - pendingTranscript = ""; + completeItem(event.item_id, event.transcript); + return; + + case "conversation.item.input_audio_transcription.failed": + completeItem(event.item_id, undefined); + config.onError?.(new Error(readRealtimeErrorDetail(event.error))); return; case "input_audio_buffer.speech_started": - pendingTranscript = ""; + pendingTranscripts.delete(event.item_id ?? unkeyedTranscript); config.onSpeechStart?.(); return; @@ -248,6 +351,9 @@ function createOpenAIRealtimeTranscriptionSession( }); }, onOpen: (transport: RealtimeTranscriptionWebSocketTransport) => { + // A reconnect starts a new provider session. Retaining outstanding item + // state would splice pre-disconnect deltas into the first new turn. + resetTranscriptionState(); transport.sendJson({ type: "session.update", session: buildOpenAIRealtimeTranscriptionSessionPayload(config), diff --git a/src/plugin-sdk/test-helpers/provider-runtime-contract.ts b/src/plugin-sdk/test-helpers/provider-runtime-contract.ts index 2c0ebd5851c8..66e9e0cc8d11 100644 --- a/src/plugin-sdk/test-helpers/provider-runtime-contract.ts +++ b/src/plugin-sdk/test-helpers/provider-runtime-contract.ts @@ -475,7 +475,7 @@ export function describeOpenAIProviderRuntimeContract(load: ProviderRuntimeContr provider: "openai", api: "openai-responses", baseUrl: "https://api.openai.com/v1", - contextWindow: 1_000_000, + contextWindow: 1_050_000, contextTokens: 272_000, maxTokens: 128_000, mediaInput: {