fix: preserve realtime provider event ownership (#116891)

Co-authored-by: Peter Steinberger <steipete@macos.shared>
This commit is contained in:
Peter Steinberger
2026-07-31 08:33:06 -07:00
committed by GitHub
parent a1026663ae
commit 0d209991b6
6 changed files with 396 additions and 13 deletions

View File

@@ -2794,6 +2794,83 @@ describe("buildOpenAIRealtimeVoiceProvider", () => {
);
});
it("surfaces input transcription failures with their provider error details", async () => {
const provider = buildOpenAIRealtimeVoiceProvider();
const onError = vi.fn();
const onEvent = vi.fn();
const bridge = provider.createBridge({
providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret
onAudio: vi.fn(),
onClearAudio: vi.fn(),
onError,
onEvent,
});
const connecting = bridge.connect();
const socket = FakeWebSocket.instances[0];
if (!socket) {
throw new Error("expected bridge to create a websocket");
}
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: "conversation.item.input_audio_transcription.failed",
item_id: "item_speech",
error: { code: "decoder_failure", message: "speech decoder exploded" },
}),
),
);
expect(onError).toHaveBeenCalledWith(
expect.objectContaining({ message: "speech decoder exploded" }),
);
expect(onEvent).toHaveBeenCalledWith({
direction: "server",
type: "conversation.item.input_audio_transcription.failed",
itemId: "item_speech",
detail: "speech decoder exploded",
});
});
it("preserves corrected final text from legacy realtime text events", async () => {
const provider = buildOpenAIRealtimeVoiceProvider();
const onTranscript = vi.fn();
const bridge = provider.createBridge({
providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret
onAudio: vi.fn(),
onClearAudio: vi.fn(),
onTranscript,
});
const connecting = bridge.connect();
const socket = FakeWebSocket.instances[0];
if (!socket) {
throw new Error("expected bridge to create a websocket");
}
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: "response.text.delta", delta: "draft assistant" })),
);
socket.emit(
"message",
Buffer.from(JSON.stringify({ type: "response.text.done", text: "corrected assistant" })),
);
expect(onTranscript.mock.calls).toEqual([
["assistant", "draft assistant", false],
["assistant", "corrected assistant", true],
]);
});
it.each([
["invalid alphabet", "not-base64!"],
["non-canonical pad bits", "ZE=="],
@@ -3029,6 +3106,80 @@ describe("buildOpenAIRealtimeVoiceProvider", () => {
});
});
it.each([
{
name: "corrected streamed arguments",
delta: '{"city":"draft"}',
finalArguments: '{"city":"Paris"}',
expectedArguments: { city: "Paris" },
},
{
name: "truncated streamed arguments",
delta: '{"city":',
finalArguments: '{"city":"Paris"}',
expectedArguments: { city: "Paris" },
},
{
name: "an explicitly empty completed payload",
delta: '{"city":"draft"}',
finalArguments: "",
expectedArguments: {},
},
])(
"uses authoritative completed tool arguments for $name",
async ({ delta, finalArguments, expectedArguments }) => {
const provider = buildOpenAIRealtimeVoiceProvider();
const onToolCall = vi.fn();
const bridge = provider.createBridge({
providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret
onAudio: vi.fn(),
onClearAudio: vi.fn(),
onToolCall,
});
const connecting = bridge.connect();
const socket = FakeWebSocket.instances[0];
if (!socket) {
throw new Error("expected bridge to create a websocket");
}
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: "response.function_call_arguments.delta",
item_id: "item_tool_1",
call_id: "call_1",
name: "lookup_weather",
delta,
}),
),
);
socket.emit(
"message",
Buffer.from(
JSON.stringify({
type: "response.function_call_arguments.done",
item_id: "item_tool_1",
call_id: "call_1",
name: "lookup_weather",
arguments: finalArguments,
}),
),
);
expect(onToolCall).toHaveBeenCalledWith({
itemId: "item_tool_1",
callId: "call_1",
name: "lookup_weather",
args: expectedArguments,
});
},
);
it("creates an explicit user item and response for manual speech", async () => {
const provider = buildOpenAIRealtimeVoiceProvider();
const onEvent = vi.fn();

View File

@@ -1388,6 +1388,7 @@ class OpenAIRealtimeVoiceBridge implements RealtimeVoiceBridge {
return;
case "conversation.output_transcript.delta":
case "response.text.delta":
case "response.output_text.delta":
case "response.audio_transcript.delta":
case "response.output_audio_transcript.delta":
@@ -1396,6 +1397,7 @@ class OpenAIRealtimeVoiceBridge implements RealtimeVoiceBridge {
}
return;
case "response.text.done":
case "response.output_text.done":
case "response.audio_transcript.done":
case "response.output_audio_transcript.done":
@@ -1420,6 +1422,10 @@ class OpenAIRealtimeVoiceBridge implements RealtimeVoiceBridge {
}
return;
case "conversation.item.input_audio_transcription.failed":
this.config.onError?.(new Error(readRealtimeErrorDetail(event.error)));
return;
case "response.cancelled":
case "response.done":
this.responseActive = false;
@@ -1456,7 +1462,8 @@ class OpenAIRealtimeVoiceBridge implements RealtimeVoiceBridge {
itemId: event.item_id,
callId: buffered?.callId || event.call_id,
name: buffered?.name || event.name,
rawArgs: buffered?.args || event.arguments,
// The done payload owns the final JSON; streamed chunks may be stale or incomplete.
rawArgs: event.arguments ?? buffered?.args,
});
this.toolCallBuffers.delete(key);
return;
@@ -1765,7 +1772,10 @@ class OpenAIRealtimeVoiceBridge implements RealtimeVoiceBridge {
}
private describeServerEvent(event: RealtimeEvent): string | undefined {
if (event.type === "error") {
if (
event.type === "error" ||
event.type === "conversation.item.input_audio_transcription.failed"
) {
return readRealtimeErrorDetail(event.error);
}
if (event.type === "response.done") {

View File

@@ -38,9 +38,10 @@ async function createRealtimeSttServer(params?: {
onRequest?: (url: URL, headers: Record<string, string | string[] | undefined>) => void;
onBinary?: (audio: Buffer) => void;
initialEvent?: unknown;
transcriptionEvents?: readonly Record<string, unknown>[];
}) {
const server = createServer();
const wss = new WebSocketServer({ noServer: true });
const wss = new WebSocketServer({ noServer: true, maxPayload: 16 * 1024 * 1024 });
const clients = new Set<WebSocket>();
const done = vi.fn();
let resolveDone: (() => void) | undefined;
@@ -63,22 +64,23 @@ async function createRealtimeSttServer(params?: {
: Buffer.from(data);
if (isBinary) {
params?.onBinary?.(buffer);
ws.send(
JSON.stringify({
const events = params?.transcriptionEvents ?? [
{
type: "transcript.partial",
text: "hello openclaw",
is_final: false,
speech_final: false,
}),
);
ws.send(
JSON.stringify({
},
{
type: "transcript.partial",
text: "hello openclaw final",
is_final: true,
speech_final: true,
}),
);
},
];
for (const event of events) {
ws.send(JSON.stringify(event));
}
return;
}
const event = JSON.parse(buffer.toString()) as { type?: string };
@@ -208,6 +210,30 @@ describe("xai realtime transcription provider", () => {
vi.unstubAllEnvs();
});
it("preserves identical final transcripts from separate speech turns", async () => {
const server = await createRealtimeSttServer({
transcriptionEvents: [
{ type: "transcript.partial", text: "yes", is_final: true, speech_final: true },
{ type: "transcript.partial", text: "yes", is_final: true, speech_final: true },
],
});
const onTranscript = vi.fn();
const onSpeechStart = vi.fn();
const session = buildXaiRealtimeTranscriptionProvider().createSession({
providerConfig: { apiKey: "xai-test-key", baseUrl: server.baseUrl },
onTranscript,
onSpeechStart,
});
await session.connect();
session.sendAudio(Buffer.from("two speech turns"));
await vi.waitFor(() => expect(onSpeechStart).toHaveBeenCalledTimes(2));
session.close();
await server.donePromise;
expect(onTranscript.mock.calls).toEqual([["yes"], ["yes"]]);
});
it("rejects setup errors before the stream is ready", async () => {
const server = await createRealtimeSttServer({
initialEvent: {

View File

@@ -174,6 +174,8 @@ function createXaiRealtimeTranscriptionSession(
return;
}
if (!speechStarted) {
// Dedupe final/terminal echoes within one utterance, not identical later turns.
lastTranscript = undefined;
speechStarted = true;
config.onSpeechStart?.();
}

View File

@@ -100,10 +100,16 @@ export abstract class XaiRealtimeVoiceEvents extends XaiRealtimeVoiceProtocol {
this.appendAssistantTranscriptDelta(event.delta);
}
return;
case "response.text.done":
case "response.output_text.done":
case "response.output_audio_transcript.done":
this.flushAssistantTranscript(event.transcript ?? event.text);
return;
case "conversation.item.input_audio_transcription.delta":
if (event.delta) {
this.config.onTranscript?.("user", event.delta, false);
}
return;
case "conversation.item.input_audio_transcription.updated":
if (event.transcript) {
this.inputTranscriptReplacements.set(this.inputTranscriptKey(event), event.transcript);
@@ -118,6 +124,10 @@ export abstract class XaiRealtimeVoiceEvents extends XaiRealtimeVoiceProtocol {
}
return;
}
case "conversation.item.input_audio_transcription.failed":
this.inputTranscriptReplacements.delete(this.inputTranscriptKey(event));
this.config.onError?.(new Error(readXaiRealtimeErrorDetail(event.error)));
return;
case "response.done":
this.flushAssistantTranscript();
this.responseActive = false;
@@ -146,7 +156,8 @@ export abstract class XaiRealtimeVoiceEvents extends XaiRealtimeVoiceProtocol {
itemId: event.item_id,
callId: buffered?.callId || event.call_id,
name: buffered?.name || event.name,
rawArgs: buffered?.args || event.arguments,
// The done payload owns the final JSON; streamed chunks may be stale or incomplete.
rawArgs: event.arguments ?? buffered?.args,
});
this.toolCallBuffers.delete(key);
return;
@@ -209,7 +220,10 @@ export abstract class XaiRealtimeVoiceEvents extends XaiRealtimeVoiceProtocol {
}
private describeServerEvent(event: XaiRealtimeEvent): string | undefined {
if (event.type === "error") {
if (
event.type === "error" ||
event.type === "conversation.item.input_audio_transcription.failed"
) {
return readXaiRealtimeErrorDetail(event.error);
}
if (event.type !== "response.done") {

View File

@@ -493,6 +493,90 @@ describe("buildXaiRealtimeVoiceProvider", () => {
expect(onTranscript).toHaveBeenCalledWith("user", "OpenClaw", true);
});
it("forwards standard incremental input-transcription events", async () => {
const provider = buildXaiRealtimeVoiceProvider();
const onTranscript = vi.fn();
const bridge = provider.createBridge({
providerConfig: { apiKey: "xai-test" }, // pragma: allowlist secret
onAudio: vi.fn(),
onClearAudio: vi.fn(),
onTranscript,
});
const { connecting, socket } = await openRealtimeBridge(bridge);
await connecting;
socket.emit(
"message",
Buffer.from(
JSON.stringify({
type: "conversation.item.input_audio_transcription.delta",
item_id: "item_speech",
delta: "open claw",
}),
),
);
expect(onTranscript).toHaveBeenCalledWith("user", "open claw", false);
});
it("surfaces input transcription failures and discards their stale replacement text", async () => {
const provider = buildXaiRealtimeVoiceProvider();
const onTranscript = vi.fn();
const onError = vi.fn();
const onEvent = vi.fn();
const bridge = provider.createBridge({
providerConfig: { apiKey: "xai-test" }, // pragma: allowlist secret
onAudio: vi.fn(),
onClearAudio: vi.fn(),
onTranscript,
onError,
onEvent,
});
const { connecting, socket } = await openRealtimeBridge(bridge);
await connecting;
socket.emit(
"message",
Buffer.from(
JSON.stringify({
type: "conversation.item.input_audio_transcription.updated",
item_id: "item_speech",
transcript: "stale speech",
}),
),
);
socket.emit(
"message",
Buffer.from(
JSON.stringify({
type: "conversation.item.input_audio_transcription.failed",
item_id: "item_speech",
error: { code: "decoder_failure", message: "speech decoder exploded" },
}),
),
);
socket.emit(
"message",
Buffer.from(
JSON.stringify({
type: "conversation.item.input_audio_transcription.completed",
item_id: "item_speech",
}),
),
);
expect(onError).toHaveBeenCalledWith(
expect.objectContaining({ message: "speech decoder exploded" }),
);
expect(onTranscript).not.toHaveBeenCalled();
expect(onEvent).toHaveBeenCalledWith({
direction: "server",
type: "conversation.item.input_audio_transcription.failed",
itemId: "item_speech",
detail: "speech decoder exploded",
});
});
it("buffers assistant transcript deltas and finalizes them when done has no text", async () => {
const provider = buildXaiRealtimeVoiceProvider();
const onTranscript = vi.fn();
@@ -538,6 +622,35 @@ describe("buildXaiRealtimeVoiceProvider", () => {
expect(onTranscript).toHaveBeenCalledTimes(3);
});
it("preserves corrected final text from legacy realtime text events", async () => {
const provider = buildXaiRealtimeVoiceProvider();
const onTranscript = vi.fn();
const bridge = provider.createBridge({
providerConfig: { apiKey: "xai-test" }, // pragma: allowlist secret
onAudio: vi.fn(),
onClearAudio: vi.fn(),
onTranscript,
});
const { connecting, socket } = await openRealtimeBridge(bridge);
await connecting;
socket.emit("message", Buffer.from(JSON.stringify({ type: "response.created" })));
socket.emit(
"message",
Buffer.from(JSON.stringify({ type: "response.text.delta", delta: "draft assistant" })),
);
socket.emit(
"message",
Buffer.from(JSON.stringify({ type: "response.text.done", text: "corrected assistant" })),
);
socket.emit("message", Buffer.from(JSON.stringify({ type: "response.done" })));
expect(onTranscript.mock.calls).toEqual([
["assistant", "draft assistant", false],
["assistant", "corrected assistant", true],
]);
});
it.each([
{
name: "lets server VAD own interruption before an audio item exists",
@@ -811,6 +924,73 @@ describe("buildXaiRealtimeVoiceProvider", () => {
});
});
it.each([
{
name: "corrected streamed arguments",
delta: '{"city":"draft"}',
finalArguments: '{"city":"Paris"}',
expectedArguments: { city: "Paris" },
},
{
name: "truncated streamed arguments",
delta: '{"city":',
finalArguments: '{"city":"Paris"}',
expectedArguments: { city: "Paris" },
},
{
name: "an explicitly empty completed payload",
delta: '{"city":"draft"}',
finalArguments: "",
expectedArguments: {},
},
])(
"uses authoritative completed tool arguments for $name",
async ({ delta, finalArguments, expectedArguments }) => {
const provider = buildXaiRealtimeVoiceProvider();
const onToolCall = vi.fn();
const bridge = provider.createBridge({
providerConfig: { apiKey: "xai-test" }, // pragma: allowlist secret
onAudio: vi.fn(),
onClearAudio: vi.fn(),
onToolCall,
});
const { connecting, socket } = await openRealtimeBridge(bridge);
await connecting;
socket.emit(
"message",
Buffer.from(
JSON.stringify({
type: "response.function_call_arguments.delta",
item_id: "item_tool_1",
call_id: "call_1",
name: "lookup_weather",
delta,
}),
),
);
socket.emit(
"message",
Buffer.from(
JSON.stringify({
type: "response.function_call_arguments.done",
item_id: "item_tool_1",
call_id: "call_1",
name: "lookup_weather",
arguments: finalArguments,
}),
),
);
expect(onToolCall).toHaveBeenCalledWith({
itemId: "item_tool_1",
callId: "call_1",
name: "lookup_weather",
args: expectedArguments,
});
},
);
it("waits for all parallel tool results before sending response.create", async () => {
vi.stubEnv("XAI_API_KEY", "xai-env"); // pragma: allowlist secret
const provider = buildXaiRealtimeVoiceProvider();