From ffe2268eecd99eab1f123aa2be528be803dbb970 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Sun, 2 Aug 2026 19:04:38 -0700 Subject: [PATCH] test(xai): consolidate realtime voice fixtures (#118318) --- .../xai/realtime-voice-provider.test.ts | 999 ++++++------------ 1 file changed, 319 insertions(+), 680 deletions(-) diff --git a/extensions/xai/realtime-voice-provider.test.ts b/extensions/xai/realtime-voice-provider.test.ts index 11ac54a891c2..c8ca46f6f73f 100644 --- a/extensions/xai/realtime-voice-provider.test.ts +++ b/extensions/xai/realtime-voice-provider.test.ts @@ -39,6 +39,15 @@ const { FakeWebSocket, isProviderAuthProfileConfiguredMock, resolveApiKeyForProv } } + emitServer(event: unknown): void { + this.emit("message", Buffer.from(JSON.stringify(event))); + } + + open(): void { + this.readyState = MockWebSocket.OPEN; + this.emit("open"); + } + send(payload: string): void { this.sent.push(payload); } @@ -90,6 +99,10 @@ vi.mock("openclaw/plugin-sdk/provider-auth-runtime", () => ({ })); type FakeWebSocketInstance = InstanceType; +type TestBridgeOptions = Parameters< + ReturnType["createBridge"] +>[0]; +type TestBridge = ReturnType["createBridge"]>; type SentRealtimeEvent = { type: string; audio?: string; @@ -144,28 +157,33 @@ function requireSession(socket: FakeWebSocketInstance, index = 0): Record; } -async function openRealtimeBridge( - bridge: ReturnType["createBridge"]>, - index = 0, - conversationId?: string, -) { +function createTestBridge(options: Partial = {}): TestBridge { + return buildXaiRealtimeVoiceProvider().createBridge({ + providerConfig: { apiKey: "xai-test" }, // pragma: allowlist secret + onAudio: vi.fn(), + onClearAudio: vi.fn(), + ...options, + }); +} + +async function startRealtimeBridge(bridge: TestBridge, index = 0, conversationId?: string) { const connecting = bridge.connect(); await waitForRealtimeState(() => expect(FakeWebSocket.instances.length).toBe(index + 1)); const socket = requireSocket(index); - socket.readyState = FakeWebSocket.OPEN; - socket.emit("open"); + socket.open(); if (conversationId) { - socket.emit( - "message", - Buffer.from( - JSON.stringify({ type: "conversation.created", conversation: { id: conversationId } }), - ), - ); + socket.emitServer({ type: "conversation.created", conversation: { id: conversationId } }); } - socket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + socket.emitServer({ type: "session.updated" }); return { connecting, socket }; } +async function openRealtimeBridge(bridge: TestBridge, index = 0, conversationId?: string) { + const { connecting, socket } = await startRealtimeBridge(bridge, index, conversationId); + await connecting; + return socket; +} + describe("buildXaiRealtimeVoiceProvider", () => { beforeEach(() => { FakeWebSocket.instances = []; @@ -204,23 +222,15 @@ describe("buildXaiRealtimeVoiceProvider", () => { }); it("does not advertise continuing realtime tool results", () => { - const provider = buildXaiRealtimeVoiceProvider(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "xai-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - }); + const bridge = createTestBridge(); expect(bridge.supportsToolResultContinuation).toBe(false); }); it("requires xAI credentials for native realtime websocket bridges", async () => { - const provider = buildXaiRealtimeVoiceProvider(); - const bridge = provider.createBridge({ + const bridge = createTestBridge({ cfg: {} as never, providerConfig: { model: "grok-voice-latest" }, - onAudio: vi.fn(), - onClearAudio: vi.fn(), }); await expect(bridge.connect()).rejects.toThrow("xAI credentials missing for realtime voice"); @@ -229,19 +239,14 @@ describe("buildXaiRealtimeVoiceProvider", () => { it("coalesces concurrent connects and ignores connects after readiness", async () => { vi.stubEnv("XAI_API_KEY", "xai-env"); // pragma: allowlist secret - const bridge = buildXaiRealtimeVoiceProvider().createBridge({ - providerConfig: { apiKey: "xai-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - }); + const bridge = createTestBridge(); const firstConnect = bridge.connect(); const secondConnect = bridge.connect(); await waitForRealtimeState(() => expect(FakeWebSocket.instances).toHaveLength(1)); const socket = requireSocket(); - socket.readyState = FakeWebSocket.OPEN; - socket.emit("open"); - socket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + socket.open(); + socket.emitServer({ type: "session.updated" }); await Promise.all([firstConnect, secondConnect]); await bridge.connect(); @@ -260,11 +265,9 @@ describe("buildXaiRealtimeVoiceProvider", () => { }), ); const onClose = vi.fn(); - const bridge = buildXaiRealtimeVoiceProvider().createBridge({ + const bridge = createTestBridge({ cfg: {} as never, providerConfig: {}, - onAudio: vi.fn(), - onClearAudio: vi.fn(), onClose, }); @@ -292,11 +295,9 @@ describe("buildXaiRealtimeVoiceProvider", () => { }), ) .mockResolvedValue({ apiKey: "xai-replacement" }); // pragma: allowlist secret - const bridge = buildXaiRealtimeVoiceProvider().createBridge({ + const bridge = createTestBridge({ cfg: {} as never, providerConfig: {}, - onAudio: vi.fn(), - onClearAudio: vi.fn(), }); const canceledConnect = bridge.connect(); @@ -306,9 +307,8 @@ describe("buildXaiRealtimeVoiceProvider", () => { await waitForRealtimeState(() => expect(FakeWebSocket.instances).toHaveLength(1)); const replacementSocket = requireSocket(); - replacementSocket.readyState = FakeWebSocket.OPEN; - replacementSocket.emit("open"); - replacementSocket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + replacementSocket.open(); + replacementSocket.emitServer({ type: "session.updated" }); await Promise.all([canceledConnect, replacementConnect]); expect(bridge.isConnected()).toBe(true); @@ -323,13 +323,11 @@ describe("buildXaiRealtimeVoiceProvider", () => { const onClose = vi.fn(); const onError = vi.fn(); const onReady = vi.fn(); - const bridge = buildXaiRealtimeVoiceProvider().createBridge({ - providerConfig: { apiKey: "xai-test" }, // pragma: allowlist secret + const bridge = createTestBridge({ onAudio, onClose, onError, onReady, - onClearAudio: vi.fn(), }); const firstConnect = bridge.connect(); @@ -343,22 +341,16 @@ describe("buildXaiRealtimeVoiceProvider", () => { const replacementConnect = bridge.connect(); await waitForRealtimeState(() => expect(FakeWebSocket.instances).toHaveLength(2)); const replacementSocket = requireSocket(1); - replacementSocket.readyState = FakeWebSocket.OPEN; - replacementSocket.emit("open"); - replacementSocket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + replacementSocket.open(); + replacementSocket.emitServer({ type: "session.updated" }); await replacementConnect; staleSocket.emit("open"); - staleSocket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); - staleSocket.emit( - "message", - Buffer.from( - JSON.stringify({ - type: "response.output_audio.delta", - delta: Buffer.from("late audio").toString("base64"), - }), - ), - ); + staleSocket.emitServer({ type: "session.updated" }); + staleSocket.emitServer({ + type: "response.output_audio.delta", + delta: Buffer.from("late audio").toString("base64"), + }); staleSocket.emit("error", new Error("late socket error")); staleSocket.flushClose(); @@ -383,10 +375,7 @@ describe("buildXaiRealtimeVoiceProvider", () => { bridgeRef.current?.close(); } }); - const bridge = buildXaiRealtimeVoiceProvider().createBridge({ - providerConfig: { apiKey: "xai-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), + const bridge = createTestBridge({ onClose, onEvent, onReady, @@ -396,9 +385,8 @@ describe("buildXaiRealtimeVoiceProvider", () => { const connecting = bridge.connect(); await waitForRealtimeState(() => expect(FakeWebSocket.instances).toHaveLength(1)); const socket = requireSocket(); - socket.readyState = FakeWebSocket.OPEN; - socket.emit("open"); - socket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + socket.open(); + socket.emitServer({ type: "session.updated" }); await connecting; expect(bridge.isConnected()).toBe(false); @@ -409,16 +397,13 @@ describe("buildXaiRealtimeVoiceProvider", () => { it("uses XAI_API_KEY for default Grok realtime bridges", async () => { vi.stubEnv("XAI_API_KEY", "xai-env"); // pragma: allowlist secret - const provider = buildXaiRealtimeVoiceProvider(); - const bridge = provider.createBridge({ + const bridge = createTestBridge({ cfg: {} as never, providerConfig: { model: "grok-voice-latest", voice: "ara" }, instructions: "Speak briefly.", - onAudio: vi.fn(), - onClearAudio: vi.fn(), }); - const { socket } = await openRealtimeBridge(bridge); + const { socket } = await startRealtimeBridge(bridge); bridge.close(); const url = socket.args[0] as string; @@ -445,14 +430,9 @@ describe("buildXaiRealtimeVoiceProvider", () => { it("does not enable xAI session resumption by default", async () => { vi.stubEnv("XAI_API_KEY", "xai-env"); // pragma: allowlist secret - const provider = buildXaiRealtimeVoiceProvider(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "xai-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - }); + const bridge = createTestBridge(); - const { socket } = await openRealtimeBridge(bridge); + const { socket } = await startRealtimeBridge(bridge); bridge.close(); expect(requireSession(socket).resumption).toBeUndefined(); @@ -460,19 +440,13 @@ describe("buildXaiRealtimeVoiceProvider", () => { it("bounds pending realtime audio by aggregate bytes before session setup", async () => { vi.stubEnv("XAI_API_KEY", "xai-env"); // pragma: allowlist secret - const provider = buildXaiRealtimeVoiceProvider(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "xai-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - }); + const bridge = createTestBridge(); bridge.sendAudio(Buffer.alloc(512 * 1024, 0x7f)); bridge.sendAudio(Buffer.alloc(512 * 1024, 0x7f)); bridge.sendAudio(Buffer.alloc(1, 0x7f)); - const { connecting, socket } = await openRealtimeBridge(bridge); - await connecting; + const socket = await openRealtimeBridge(bridge); expect( parseSent(socket).filter((event) => event.type === "input_audio_buffer.append"), @@ -482,20 +456,14 @@ describe("buildXaiRealtimeVoiceProvider", () => { it("copies pending realtime audio views without retaining their backing allocation", async () => { vi.stubEnv("XAI_API_KEY", "xai-env"); // pragma: allowlist secret - const provider = buildXaiRealtimeVoiceProvider(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "xai-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - }); + const bridge = createTestBridge(); const backing = Buffer.alloc(2 * 1024 * 1024, 0x7f); const view = backing.subarray(0, 1); bridge.sendAudio(view); backing[0] = 0; - const { connecting, socket } = await openRealtimeBridge(bridge); - await connecting; + const socket = await openRealtimeBridge(bridge); expect(parseSent(socket).filter((event) => event.type === "input_audio_buffer.append")).toEqual( [{ type: "input_audio_buffer.append", audio: "fw==" }], @@ -505,12 +473,7 @@ describe("buildXaiRealtimeVoiceProvider", () => { it("drops queued realtime input on close and ignores late input until reconnect", async () => { vi.stubEnv("XAI_API_KEY", "xai-env"); // pragma: allowlist secret - const provider = buildXaiRealtimeVoiceProvider(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "xai-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - }); + const bridge = createTestBridge(); bridge.sendAudio(Buffer.from([0x01])); bridge.sendUserMessage?.("queued before close"); @@ -520,8 +483,7 @@ describe("buildXaiRealtimeVoiceProvider", () => { bridge.sendUserMessage?.("late after close"); void bridge.submitToolResult("call-after-close", { ok: true }); - const { connecting, socket } = await openRealtimeBridge(bridge); - await connecting; + const socket = await openRealtimeBridge(bridge); expect( parseSent(socket).filter( @@ -565,15 +527,11 @@ describe("buildXaiRealtimeVoiceProvider", () => { it("sends nested xAI session.update audio formats for g711 bridges", async () => { vi.stubEnv("XAI_API_KEY", "xai-env"); // pragma: allowlist secret - const provider = buildXaiRealtimeVoiceProvider(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "xai-test" }, // pragma: allowlist secret + const bridge = createTestBridge({ audioFormat: { encoding: "g711_ulaw", sampleRateHz: 8000, channels: 1 }, - onAudio: vi.fn(), - onClearAudio: vi.fn(), }); - const { socket } = await openRealtimeBridge(bridge); + const { socket } = await startRealtimeBridge(bridge); bridge.close(); const session = requireSession(socket); @@ -587,7 +545,6 @@ describe("buildXaiRealtimeVoiceProvider", () => { }); it("only forwards xAI VAD values accepted by the realtime API", async () => { - const provider = buildXaiRealtimeVoiceProvider(); const cases = [ { options: { vadThreshold: 1, silenceDurationMs: 10_001, prefixPaddingMs: -1 }, @@ -600,15 +557,13 @@ describe("buildXaiRealtimeVoiceProvider", () => { ]; for (const [index, { options, expected }] of cases.entries()) { - const bridge = provider.createBridge({ + const bridge = createTestBridge({ providerConfig: { apiKey: "xai-test", // pragma: allowlist secret ...options, }, - onAudio: vi.fn(), - onClearAudio: vi.fn(), }); - const { socket } = await openRealtimeBridge(bridge, index); + const { socket } = await startRealtimeBridge(bridge, index); bridge.close(); expect(requireSession(socket).turn_detection).toEqual(expect.objectContaining(expected)); } @@ -636,55 +591,35 @@ describe("buildXaiRealtimeVoiceProvider", () => { ...callbacks, }); - const { socket } = await openRealtimeBridge(bridge); + const { socket } = await startRealtimeBridge(bridge); bridge.close(); expect(requireSession(socket).reasoning).toEqual({ effort: "none" }); }); it("treats xAI input transcription updates as replacements until completed", 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(), + const bridge = createTestBridge({ onTranscript, }); - const { connecting, socket } = await openRealtimeBridge(bridge); - await connecting; + const socket = await openRealtimeBridge(bridge); - socket.emit( - "message", - Buffer.from( - JSON.stringify({ - type: "conversation.item.input_audio_transcription.updated", - item_id: "item_1", - transcript: "open", - }), - ), - ); - socket.emit( - "message", - Buffer.from( - JSON.stringify({ - type: "conversation.item.input_audio_transcription.updated", - item_id: "item_1", - transcript: "open claw", - }), - ), - ); - socket.emit( - "message", - Buffer.from( - JSON.stringify({ - type: "conversation.item.input_audio_transcription.completed", - item_id: "item_1", - transcript: "OpenClaw", - }), - ), - ); + socket.emitServer({ + type: "conversation.item.input_audio_transcription.updated", + item_id: "item_1", + transcript: "open", + }); + socket.emitServer({ + type: "conversation.item.input_audio_transcription.updated", + item_id: "item_1", + transcript: "open claw", + }); + socket.emitServer({ + type: "conversation.item.input_audio_transcription.completed", + item_id: "item_1", + transcript: "OpenClaw", + }); bridge.close(); expect(onTranscript).toHaveBeenCalledOnce(); @@ -692,76 +627,46 @@ describe("buildXaiRealtimeVoiceProvider", () => { }); 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(), + const bridge = createTestBridge({ onTranscript, }); - const { connecting, socket } = await openRealtimeBridge(bridge); - await connecting; + const socket = await openRealtimeBridge(bridge); - socket.emit( - "message", - Buffer.from( - JSON.stringify({ - type: "conversation.item.input_audio_transcription.delta", - item_id: "item_speech", - delta: "open claw", - }), - ), - ); + socket.emitServer({ + 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(), + const bridge = createTestBridge({ onTranscript, onError, onEvent, }); - const { connecting, socket } = await openRealtimeBridge(bridge); - await connecting; + const socket = await openRealtimeBridge(bridge); - 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", - }), - ), - ); + socket.emitServer({ + type: "conversation.item.input_audio_transcription.updated", + item_id: "item_speech", + transcript: "stale speech", + }); + socket.emitServer({ + type: "conversation.item.input_audio_transcription.failed", + item_id: "item_speech", + error: { code: "decoder_failure", message: "speech decoder exploded" }, + }); + socket.emitServer({ + type: "conversation.item.input_audio_transcription.completed", + item_id: "item_speech", + }); expect(onError).toHaveBeenCalledWith( expect.objectContaining({ message: "speech decoder exploded" }), @@ -776,42 +681,18 @@ describe("buildXaiRealtimeVoiceProvider", () => { }); it("buffers assistant transcript deltas and finalizes them when done has no text", 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(), + const bridge = createTestBridge({ onTranscript, }); - const { connecting, socket } = await openRealtimeBridge(bridge); - await connecting; + const socket = await openRealtimeBridge(bridge); - socket.emit("message", Buffer.from(JSON.stringify({ type: "response.created" }))); - socket.emit( - "message", - Buffer.from( - JSON.stringify({ - type: "response.output_audio_transcript.delta", - delta: "Hello ", - }), - ), - ); - socket.emit( - "message", - Buffer.from( - JSON.stringify({ - type: "response.output_audio_transcript.delta", - delta: "OpenClaw", - }), - ), - ); - socket.emit( - "message", - Buffer.from(JSON.stringify({ type: "response.output_audio_transcript.done" })), - ); - socket.emit("message", Buffer.from(JSON.stringify({ type: "response.done" }))); + socket.emitServer({ type: "response.created" }); + socket.emitServer({ type: "response.output_audio_transcript.delta", delta: "Hello " }); + socket.emitServer({ type: "response.output_audio_transcript.delta", delta: "OpenClaw" }); + socket.emitServer({ type: "response.output_audio_transcript.done" }); + socket.emitServer({ type: "response.done" }); bridge.close(); expect(onTranscript).toHaveBeenNthCalledWith(1, "assistant", "Hello ", false); @@ -821,27 +702,16 @@ describe("buildXaiRealtimeVoiceProvider", () => { }); 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(), + const bridge = createTestBridge({ onTranscript, }); - const { connecting, socket } = await openRealtimeBridge(bridge); - await connecting; + const socket = await openRealtimeBridge(bridge); - 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" }))); + socket.emitServer({ type: "response.created" }); + socket.emitServer({ type: "response.text.delta", delta: "draft assistant" }); + socket.emitServer({ type: "response.text.done", text: "corrected assistant" }); + socket.emitServer({ type: "response.done" }); expect(onTranscript.mock.calls).toEqual([ ["assistant", "draft assistant", false], @@ -938,20 +808,17 @@ describe("buildXaiRealtimeVoiceProvider", () => { vi.stubEnv("XAI_API_KEY", "xai-env"); // pragma: allowlist secret const onAudio = vi.fn(); const onClearAudio = vi.fn(); - const bridge = buildXaiRealtimeVoiceProvider().createBridge({ - providerConfig: { apiKey: "xai-test" }, // pragma: allowlist secret + const bridge = createTestBridge({ audioFormat: REALTIME_VOICE_AUDIO_FORMAT_PCM16_24KHZ, onAudio, onClearAudio, }); - const { socket } = await openRealtimeBridge(bridge); - const emit = (event: Record) => - socket.emit("message", Buffer.from(JSON.stringify(event))); + const { socket } = await startRealtimeBridge(bridge); - emit({ type: "response.created", response: { id: "resp_1" } }); + socket.emitServer({ type: "response.created", response: { id: "resp_1" } }); if (hasAudio) { bridge.setMediaTimestamp(1000); - emit({ + socket.emitServer({ type: "response.output_audio.delta", item_id: "item_1", delta: Buffer.from("assistant audio").toString("base64"), @@ -961,16 +828,16 @@ describe("buildXaiRealtimeVoiceProvider", () => { bridge.acknowledgeMark?.(); } if (completed) { - emit({ type: "response.done" }); + socket.emitServer({ type: "response.done" }); } if (startNewResponse) { - emit({ type: "response.created", response: { id: "resp_2" } }); + socket.emitServer({ type: "response.created", response: { id: "resp_2" } }); } bridge.setMediaTimestamp(timestamp); if (manual) { bridge.handleBargeIn?.({ audioPlaybackActive: true }); } else { - emit({ type: "input_audio_buffer.speech_started" }); + socket.emitServer({ type: "input_audio_buffer.speech_started" }); } bridge.close(); @@ -998,20 +865,14 @@ describe("buildXaiRealtimeVoiceProvider", () => { onError(error); retry = bridge.connect(); }; - const bridge = buildXaiRealtimeVoiceProvider().createBridge({ - providerConfig: { apiKey: "xai-test" }, // pragma: allowlist secret + const bridge = createTestBridge({ onAudio, onClose, onError: handleError, - onClearAudio: vi.fn(), }); - const { connecting, socket } = await openRealtimeBridge(bridge); - await connecting; - socket.emit( - "message", - Buffer.from(JSON.stringify({ type: "response.output_audio.delta", delta: "ZE==" })), - ); + const socket = await openRealtimeBridge(bridge); + socket.emitServer({ type: "response.output_audio.delta", delta: "ZE==" }); expect(onError).toHaveBeenCalledWith( expect.objectContaining({ @@ -1036,23 +897,16 @@ describe("buildXaiRealtimeVoiceProvider", () => { vi.stubEnv("XAI_API_KEY", "xai-env"); // pragma: allowlist secret const onClose = vi.fn(); const onError = vi.fn(); - const bridge = buildXaiRealtimeVoiceProvider().createBridge({ - providerConfig: { apiKey: "xai-test" }, // pragma: allowlist secret - onAudio: vi.fn(), + const bridge = createTestBridge({ onClose, onError, - onClearAudio: vi.fn(), }); const connection = bridge.connect(); await waitForRealtimeState(() => expect(FakeWebSocket.instances.length).toBe(1)); const socket = requireSocket(); - socket.readyState = FakeWebSocket.OPEN; - socket.emit("open"); - socket.emit( - "message", - Buffer.from(JSON.stringify({ type: "response.output_audio.delta", delta: "ZE==" })), - ); + socket.open(); + socket.emitServer({ type: "response.output_audio.delta", delta: "ZE==" }); await expect(connection).rejects.toThrow( "xAI realtime voice stream returned malformed base64 audio data", @@ -1065,53 +919,33 @@ describe("buildXaiRealtimeVoiceProvider", () => { it("deduplicates repeated function-call arguments done events", async () => { vi.stubEnv("XAI_API_KEY", "xai-env"); // pragma: allowlist secret - const provider = buildXaiRealtimeVoiceProvider(); const onToolCall = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "xai-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), + const bridge = createTestBridge({ onToolCall, }); - const { connecting, socket } = await openRealtimeBridge(bridge); - await connecting; + const socket = await openRealtimeBridge(bridge); - socket.emit( - "message", - Buffer.from( - JSON.stringify({ - type: "response.function_call_arguments.delta", - item_id: "item_tool_1", - name: "openclaw_agent_consult", - call_id: "call_1", - delta: JSON.stringify({ question: "delegate this" }), - }), - ), - ); - socket.emit( - "message", - Buffer.from( - JSON.stringify({ - type: "response.function_call_arguments.done", - item_id: "item_tool_1", - name: "openclaw_agent_consult", - call_id: "call_1", - }), - ), - ); - socket.emit( - "message", - Buffer.from( - JSON.stringify({ - type: "response.function_call_arguments.done", - item_id: "item_tool_1", - name: "openclaw_agent_consult", - call_id: "call_1", - arguments: JSON.stringify({ question: "delegate this" }), - }), - ), - ); + socket.emitServer({ + type: "response.function_call_arguments.delta", + item_id: "item_tool_1", + name: "openclaw_agent_consult", + call_id: "call_1", + delta: JSON.stringify({ question: "delegate this" }), + }); + socket.emitServer({ + type: "response.function_call_arguments.done", + item_id: "item_tool_1", + name: "openclaw_agent_consult", + call_id: "call_1", + }); + socket.emitServer({ + type: "response.function_call_arguments.done", + item_id: "item_tool_1", + name: "openclaw_agent_consult", + call_id: "call_1", + arguments: JSON.stringify({ question: "delegate this" }), + }); expect(onToolCall).toHaveBeenCalledTimes(1); expect(onToolCall).toHaveBeenCalledWith({ @@ -1144,41 +978,26 @@ describe("buildXaiRealtimeVoiceProvider", () => { ])( "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(), + const bridge = createTestBridge({ onToolCall, }); - const { connecting, socket } = await openRealtimeBridge(bridge); - await connecting; + const socket = await openRealtimeBridge(bridge); - 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, - }), - ), - ); + socket.emitServer({ + type: "response.function_call_arguments.delta", + item_id: "item_tool_1", + call_id: "call_1", + name: "lookup_weather", + delta, + }); + socket.emitServer({ + 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", @@ -1194,20 +1013,17 @@ describe("buildXaiRealtimeVoiceProvider", () => { async (ingress) => { const onEvent = vi.fn(); const onToolCall = vi.fn(); - const bridge = buildXaiRealtimeVoiceProvider().createBridge({ + const bridge = createTestBridge({ providerConfig: { apiKey: "xai-test", // pragma: allowlist secret sessionResumption: ingress === "resumed item replay", }, - onAudio: vi.fn(), - onClearAudio: vi.fn(), onEvent, onToolCall, }); - const { connecting, socket } = await openRealtimeBridge(bridge); - await connecting; + const socket = await openRealtimeBridge(bridge); - socket.emit("message", Buffer.from(JSON.stringify({ type: "response.created" }))); + socket.emitServer({ type: "response.created" }); const invalidEvents = ["{", "null", "[]", JSON.stringify("text"), "1", "true"].map( (rawArgs, index) => { const itemId = `item_invalid_${index}`; @@ -1233,9 +1049,9 @@ describe("buildXaiRealtimeVoiceProvider", () => { }, ); for (const event of invalidEvents) { - socket.emit("message", Buffer.from(JSON.stringify(event))); + socket.emitServer(event); } - socket.emit("message", Buffer.from(JSON.stringify(invalidEvents[0]))); + socket.emitServer(invalidEvents[0]); expect(onToolCall).not.toHaveBeenCalled(); expect( @@ -1274,13 +1090,13 @@ describe("buildXaiRealtimeVoiceProvider", () => { ); expect(parseSent(socket).filter((event) => event.type === "response.create")).toEqual([]); - socket.emit("message", Buffer.from(JSON.stringify({ type: "response.done" }))); + socket.emitServer({ type: "response.done" }); expect(parseSent(socket).filter((event) => event.type === "response.create")).toEqual([ { type: "response.create" }, ]); - socket.emit("message", Buffer.from(JSON.stringify({ type: "response.created" }))); - socket.emit("message", Buffer.from(JSON.stringify({ type: "response.done" }))); + socket.emitServer({ type: "response.created" }); + socket.emitServer({ type: "response.done" }); bridge.sendUserMessage?.("Continue after rejected tool arguments."); expect(parseSent(socket).slice(-2)).toEqual([ { @@ -1299,28 +1115,19 @@ describe("buildXaiRealtimeVoiceProvider", () => { it("treats rejected tool call arguments as terminal for the same identity", async () => { const onToolCall = vi.fn(); - const bridge = buildXaiRealtimeVoiceProvider().createBridge({ - providerConfig: { apiKey: "xai-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), + const bridge = createTestBridge({ onToolCall, }); - const { connecting, socket } = await openRealtimeBridge(bridge); - await connecting; + const socket = await openRealtimeBridge(bridge); for (const rawArgs of ['{"city":', JSON.stringify({ city: "Paris" })]) { - socket.emit( - "message", - Buffer.from( - JSON.stringify({ - type: "response.function_call_arguments.done", - item_id: "item_rejected", - call_id: "call_rejected", - name: "lookup_weather", - arguments: rawArgs, - }), - ), - ); + socket.emitServer({ + type: "response.function_call_arguments.done", + item_id: "item_rejected", + call_id: "call_rejected", + name: "lookup_weather", + arguments: rawArgs, + }); } expect(onToolCall).not.toHaveBeenCalled(); @@ -1344,30 +1151,20 @@ describe("buildXaiRealtimeVoiceProvider", () => { 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(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "xai-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), + const bridge = createTestBridge({ onToolCall: vi.fn(), }); - const { connecting, socket } = await openRealtimeBridge(bridge); - await connecting; + const socket = await openRealtimeBridge(bridge); for (const callId of ["call_1", "call_2"]) { - socket.emit( - "message", - Buffer.from( - JSON.stringify({ - type: "response.function_call_arguments.done", - item_id: `item_${callId}`, - name: "openclaw_agent_consult", - call_id: callId, - arguments: JSON.stringify({ question: callId }), - }), - ), - ); + socket.emitServer({ + type: "response.function_call_arguments.done", + item_id: `item_${callId}`, + name: "openclaw_agent_consult", + call_id: callId, + arguments: JSON.stringify({ question: callId }), + }); } await bridge.submitToolResult("call_1", { text: "first" }); @@ -1389,29 +1186,19 @@ describe("buildXaiRealtimeVoiceProvider", () => { it("does not send unsupported interim willContinue tool results to xAI", async () => { vi.stubEnv("XAI_API_KEY", "xai-env"); // pragma: allowlist secret - const provider = buildXaiRealtimeVoiceProvider(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "xai-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), + const bridge = createTestBridge({ onToolCall: vi.fn(), }); - const { connecting, socket } = await openRealtimeBridge(bridge); - await connecting; + const socket = await openRealtimeBridge(bridge); - socket.emit( - "message", - Buffer.from( - JSON.stringify({ - type: "response.function_call_arguments.done", - item_id: "item_call_1", - name: "openclaw_agent_consult", - call_id: "call_1", - arguments: JSON.stringify({ question: "call_1" }), - }), - ), - ); + socket.emitServer({ + type: "response.function_call_arguments.done", + item_id: "item_call_1", + name: "openclaw_agent_consult", + call_id: "call_1", + arguments: JSON.stringify({ question: "call_1" }), + }); await bridge.submitToolResult("call_1", { status: "working" }, { willContinue: true }); expect(parseSent(socket).filter((event) => event.type === "conversation.item.create")).toEqual( @@ -1434,43 +1221,28 @@ describe("buildXaiRealtimeVoiceProvider", () => { it("defers response.create for tool results until queued playback marks drain", async () => { vi.stubEnv("XAI_API_KEY", "xai-env"); // pragma: allowlist secret - const provider = buildXaiRealtimeVoiceProvider(); const onMark = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "xai-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), + const bridge = createTestBridge({ onToolCall: vi.fn(), onMark, }); - const { connecting, socket } = await openRealtimeBridge(bridge); - await connecting; + const socket = await openRealtimeBridge(bridge); - socket.emit("message", Buffer.from(JSON.stringify({ type: "response.created" }))); - socket.emit( - "message", - Buffer.from( - JSON.stringify({ - type: "response.output_audio.delta", - item_id: "item_audio_1", - delta: Buffer.from("assistant audio").toString("base64"), - }), - ), - ); - socket.emit("message", Buffer.from(JSON.stringify({ type: "response.done" }))); - socket.emit( - "message", - Buffer.from( - JSON.stringify({ - type: "response.function_call_arguments.done", - item_id: "item_call_1", - name: "openclaw_agent_consult", - call_id: "call_1", - arguments: JSON.stringify({ question: "call_1" }), - }), - ), - ); + socket.emitServer({ type: "response.created" }); + socket.emitServer({ + type: "response.output_audio.delta", + item_id: "item_audio_1", + delta: Buffer.from("assistant audio").toString("base64"), + }); + socket.emitServer({ type: "response.done" }); + socket.emitServer({ + type: "response.function_call_arguments.done", + item_id: "item_call_1", + name: "openclaw_agent_consult", + call_id: "call_1", + arguments: JSON.stringify({ question: "call_1" }), + }); await bridge.submitToolResult("call_1", { text: "final" }); expect(parseSent(socket).filter((event) => event.type === "response.create")).toEqual([]); @@ -1487,30 +1259,21 @@ describe("buildXaiRealtimeVoiceProvider", () => { it("preserves pending parallel tool calls across resumed reconnects", async () => { vi.useFakeTimers(); vi.stubEnv("XAI_API_KEY", "xai-env"); // pragma: allowlist secret - const provider = buildXaiRealtimeVoiceProvider(); - const bridge = provider.createBridge({ + const bridge = createTestBridge({ providerConfig: { apiKey: "xai-test", sessionResumption: true }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), onToolCall: vi.fn(), }); - const { connecting, socket: firstSocket } = await openRealtimeBridge(bridge, 0, "conv_tools"); - await connecting; + const firstSocket = await openRealtimeBridge(bridge, 0, "conv_tools"); for (const callId of ["call_1", "call_2"]) { - firstSocket.emit( - "message", - Buffer.from( - JSON.stringify({ - type: "response.function_call_arguments.done", - item_id: `item_${callId}`, - name: "openclaw_agent_consult", - call_id: callId, - arguments: JSON.stringify({ question: callId }), - }), - ), - ); + firstSocket.emitServer({ + type: "response.function_call_arguments.done", + item_id: `item_${callId}`, + name: "openclaw_agent_consult", + call_id: callId, + arguments: JSON.stringify({ question: callId }), + }); } firstSocket.close(1006, "connection lost"); @@ -1518,9 +1281,8 @@ describe("buildXaiRealtimeVoiceProvider", () => { await waitForRealtimeState(() => expect(FakeWebSocket.instances.length).toBe(2)); const secondSocket = requireSocket(1); expect(String(secondSocket.args[0])).toContain("conversation_id=conv_tools"); - secondSocket.readyState = FakeWebSocket.OPEN; - secondSocket.emit("open"); - secondSocket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + secondSocket.open(); + secondSocket.emitServer({ type: "session.updated" }); await bridge.submitToolResult("call_1", { text: "first" }); expect(parseSent(secondSocket).filter((event) => event.type === "response.create")).toEqual([]); @@ -1543,39 +1305,29 @@ describe("buildXaiRealtimeVoiceProvider", () => { it("delivers a tool call first observed in resumed item replay", async () => { vi.useFakeTimers(); vi.stubEnv("XAI_API_KEY", "xai-env"); // pragma: allowlist secret - const provider = buildXaiRealtimeVoiceProvider(); const onToolCall = vi.fn(); - const bridge = provider.createBridge({ + const bridge = createTestBridge({ providerConfig: { apiKey: "xai-test", sessionResumption: true }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), onToolCall, }); - const { connecting, socket: firstSocket } = await openRealtimeBridge(bridge, 0, "conv_replay"); - await connecting; + const firstSocket = await openRealtimeBridge(bridge, 0, "conv_replay"); firstSocket.close(1006, "connection lost"); await vi.advanceTimersByTimeAsync(1000); await waitForRealtimeState(() => expect(FakeWebSocket.instances.length).toBe(2)); const secondSocket = requireSocket(1); - secondSocket.readyState = FakeWebSocket.OPEN; - secondSocket.emit("open"); - secondSocket.emit( - "message", - Buffer.from( - JSON.stringify({ - type: "conversation.item.created", - item: { - id: "item_replayed_call", - type: "function_call", - call_id: "call_replayed", - name: "openclaw_agent_consult", - arguments: JSON.stringify({ question: "recover me" }), - }, - }), - ), - ); + secondSocket.open(); + secondSocket.emitServer({ + type: "conversation.item.created", + item: { + id: "item_replayed_call", + type: "function_call", + call_id: "call_replayed", + name: "openclaw_agent_consult", + arguments: JSON.stringify({ question: "recover me" }), + }, + }); expect(onToolCall).toHaveBeenCalledWith({ itemId: "item_replayed_call", @@ -1589,37 +1341,24 @@ describe("buildXaiRealtimeVoiceProvider", () => { it("fails closed when a tool output was not acknowledged before reconnect", async () => { vi.useFakeTimers(); vi.stubEnv("XAI_API_KEY", "xai-env"); // pragma: allowlist secret - const provider = buildXaiRealtimeVoiceProvider(); const onToolCall = vi.fn(); const onEvent = vi.fn(); const onClose = vi.fn(); - const bridge = provider.createBridge({ + const bridge = createTestBridge({ providerConfig: { apiKey: "xai-test", sessionResumption: true }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), onToolCall, onEvent, onClose, }); - const { connecting, socket: firstSocket } = await openRealtimeBridge( - bridge, - 0, - "conv_lost_output", - ); - await connecting; - firstSocket.emit( - "message", - Buffer.from( - JSON.stringify({ - type: "response.function_call_arguments.done", - item_id: "item_lost_output", - call_id: "call_lost_output", - name: "openclaw_agent_consult", - arguments: JSON.stringify({ question: "recover output" }), - }), - ), - ); + const firstSocket = await openRealtimeBridge(bridge, 0, "conv_lost_output"); + firstSocket.emitServer({ + type: "response.function_call_arguments.done", + item_id: "item_lost_output", + call_id: "call_lost_output", + name: "openclaw_agent_consult", + arguments: JSON.stringify({ question: "recover output" }), + }); await bridge.submitToolResult("call_lost_output", { text: "recovered" }); firstSocket.close(1006, "output acknowledgement lost"); @@ -1639,56 +1378,37 @@ describe("buildXaiRealtimeVoiceProvider", () => { it("does not retry a tool output acknowledged by resumed item replay", async () => { vi.useFakeTimers(); vi.stubEnv("XAI_API_KEY", "xai-env"); // pragma: allowlist secret - const provider = buildXaiRealtimeVoiceProvider(); const onToolCall = vi.fn(); - const bridge = provider.createBridge({ + const bridge = createTestBridge({ providerConfig: { apiKey: "xai-test", sessionResumption: true }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), onToolCall, }); - const { connecting, socket: firstSocket } = await openRealtimeBridge( - bridge, - 0, - "conv_saved_output", - ); - await connecting; - firstSocket.emit( - "message", - Buffer.from( - JSON.stringify({ - type: "response.function_call_arguments.done", - item_id: "item_saved_output", - call_id: "call_saved_output", - name: "openclaw_agent_consult", - arguments: JSON.stringify({ question: "saved output" }), - }), - ), - ); + const firstSocket = await openRealtimeBridge(bridge, 0, "conv_saved_output"); + firstSocket.emitServer({ + type: "response.function_call_arguments.done", + item_id: "item_saved_output", + call_id: "call_saved_output", + name: "openclaw_agent_consult", + arguments: JSON.stringify({ question: "saved output" }), + }); await bridge.submitToolResult("call_saved_output", { text: "saved" }); - firstSocket.emit( - "message", - Buffer.from( - JSON.stringify({ - type: "conversation.item.added", - item: { - id: "item_saved_result", - type: "function_call_output", - call_id: "call_saved_output", - output: JSON.stringify({ text: "saved" }), - }, - }), - ), - ); + firstSocket.emitServer({ + type: "conversation.item.added", + item: { + id: "item_saved_result", + type: "function_call_output", + call_id: "call_saved_output", + output: JSON.stringify({ text: "saved" }), + }, + }); firstSocket.close(1006, "connection lost after output acknowledgement"); await vi.advanceTimersByTimeAsync(1000); await waitForRealtimeState(() => expect(FakeWebSocket.instances.length).toBe(2)); const secondSocket = requireSocket(1); - secondSocket.readyState = FakeWebSocket.OPEN; - secondSocket.emit("open"); - secondSocket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + secondSocket.open(); + secondSocket.emitServer({ type: "session.updated" }); for (const item of [ { id: "item_saved_output", @@ -1704,10 +1424,7 @@ describe("buildXaiRealtimeVoiceProvider", () => { output: JSON.stringify({ text: "saved" }), }, ]) { - secondSocket.emit( - "message", - Buffer.from(JSON.stringify({ type: "conversation.item.created", item })), - ); + secondSocket.emitServer({ type: "conversation.item.created", item }); } await vi.advanceTimersByTimeAsync(500); @@ -1721,34 +1438,21 @@ describe("buildXaiRealtimeVoiceProvider", () => { it("queues tool results submitted while a resumed session is reconnecting", async () => { vi.useFakeTimers(); vi.stubEnv("XAI_API_KEY", "xai-env"); // pragma: allowlist secret - const provider = buildXaiRealtimeVoiceProvider(); - const bridge = provider.createBridge({ + const bridge = createTestBridge({ providerConfig: { apiKey: "xai-test", sessionResumption: true }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), onToolCall: vi.fn(), }); - const { connecting, socket: firstSocket } = await openRealtimeBridge( - bridge, - 0, - "conv_tool_queue", - ); - await connecting; + const firstSocket = await openRealtimeBridge(bridge, 0, "conv_tool_queue"); for (const callId of ["call_1", "call_2"]) { - firstSocket.emit( - "message", - Buffer.from( - JSON.stringify({ - type: "response.function_call_arguments.done", - item_id: `item_${callId}`, - name: "openclaw_agent_consult", - call_id: callId, - arguments: JSON.stringify({ question: callId }), - }), - ), - ); + firstSocket.emitServer({ + type: "response.function_call_arguments.done", + item_id: `item_${callId}`, + name: "openclaw_agent_consult", + call_id: callId, + arguments: JSON.stringify({ question: callId }), + }); } firstSocket.close(1006, "connection lost"); @@ -1756,15 +1460,14 @@ describe("buildXaiRealtimeVoiceProvider", () => { await waitForRealtimeState(() => expect(FakeWebSocket.instances.length).toBe(2)); const secondSocket = requireSocket(1); expect(String(secondSocket.args[0])).toContain("conversation_id=conv_tool_queue"); - secondSocket.readyState = FakeWebSocket.OPEN; - secondSocket.emit("open"); + secondSocket.open(); await bridge.submitToolResult("call_1", { text: "first" }); expect( parseSent(secondSocket).filter((event) => event.type === "conversation.item.create"), ).toEqual([]); - secondSocket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + secondSocket.emitServer({ type: "session.updated" }); expect(parseSent(secondSocket).filter((event) => event.type === "response.create")).toEqual([]); expect(parseSent(secondSocket).slice(-1)).toEqual([ { @@ -1795,34 +1498,25 @@ describe("buildXaiRealtimeVoiceProvider", () => { it("queues text turns submitted while a resumed session is reconnecting", async () => { vi.useFakeTimers(); vi.stubEnv("XAI_API_KEY", "xai-env"); // pragma: allowlist secret - const provider = buildXaiRealtimeVoiceProvider(); - const bridge = provider.createBridge({ + const bridge = createTestBridge({ providerConfig: { apiKey: "xai-test", sessionResumption: true }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), }); - const { connecting, socket: firstSocket } = await openRealtimeBridge( - bridge, - 0, - "conv_text_queue", - ); - await connecting; + const firstSocket = await openRealtimeBridge(bridge, 0, "conv_text_queue"); firstSocket.close(1006, "connection lost"); await vi.advanceTimersByTimeAsync(1000); await waitForRealtimeState(() => expect(FakeWebSocket.instances.length).toBe(2)); const secondSocket = requireSocket(1); expect(String(secondSocket.args[0])).toContain("conversation_id=conv_text_queue"); - secondSocket.readyState = FakeWebSocket.OPEN; - secondSocket.emit("open"); + secondSocket.open(); bridge.sendUserMessage?.("OpenClaw finished checking."); expect( parseSent(secondSocket).filter((event) => event.type === "conversation.item.create"), ).toEqual([]); - secondSocket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + secondSocket.emitServer({ type: "session.updated" }); expect(parseSent(secondSocket).slice(-2)).toEqual([ { type: "conversation.item.create", @@ -1840,23 +1534,15 @@ describe("buildXaiRealtimeVoiceProvider", () => { it("exhausts reconnect attempts when websocket opens without session setup", async () => { vi.useFakeTimers(); vi.stubEnv("XAI_API_KEY", "xai-env"); // pragma: allowlist secret - const provider = buildXaiRealtimeVoiceProvider(); const onEvent = vi.fn(); const onClose = vi.fn(); - const bridge = provider.createBridge({ + const bridge = createTestBridge({ providerConfig: { apiKey: "xai-test", sessionResumption: true }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), onEvent, onClose, }); - const { connecting, socket: firstSocket } = await openRealtimeBridge( - bridge, - 0, - "conv_reconnect", - ); - await connecting; + const firstSocket = await openRealtimeBridge(bridge, 0, "conv_reconnect"); firstSocket.close(1006, "connection lost"); @@ -1865,8 +1551,7 @@ describe("buildXaiRealtimeVoiceProvider", () => { await vi.advanceTimersByTimeAsync(delayMs); await waitForRealtimeState(() => expect(FakeWebSocket.instances.length).toBe(attempt + 1)); const socket = requireSocket(attempt); - socket.readyState = FakeWebSocket.OPEN; - socket.emit("open"); + socket.open(); socket.close(1006, "session setup failed"); } @@ -1884,26 +1569,21 @@ describe("buildXaiRealtimeVoiceProvider", () => { it("does not replay ready callbacks after reconnect", async () => { vi.useFakeTimers(); vi.stubEnv("XAI_API_KEY", "xai-env"); // pragma: allowlist secret - const provider = buildXaiRealtimeVoiceProvider(); const onReady = vi.fn(); - const bridge = provider.createBridge({ + const bridge = createTestBridge({ providerConfig: { apiKey: "xai-test", sessionResumption: true }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), onReady, }); - const { connecting, socket: firstSocket } = await openRealtimeBridge(bridge, 0, "conv_ready"); - await connecting; + const firstSocket = await openRealtimeBridge(bridge, 0, "conv_ready"); firstSocket.close(1006, "connection lost"); await vi.advanceTimersByTimeAsync(1000); await waitForRealtimeState(() => expect(FakeWebSocket.instances.length).toBe(2)); const secondSocket = requireSocket(1); expect(String(secondSocket.args[0])).toContain("conversation_id=conv_ready"); - secondSocket.readyState = FakeWebSocket.OPEN; - secondSocket.emit("open"); - secondSocket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + secondSocket.open(); + secondSocket.emitServer({ type: "session.updated" }); bridge.close(); expect(onReady).toHaveBeenCalledOnce(); @@ -1912,17 +1592,13 @@ describe("buildXaiRealtimeVoiceProvider", () => { it("cancels a pending reconnect and allows a later explicit connect", async () => { vi.useFakeTimers(); resolveApiKeyForProviderMock.mockResolvedValue({ apiKey: ["xai", "test"].join("-") }); - const provider = buildXaiRealtimeVoiceProvider(); const onError = vi.fn(); - const bridge = provider.createBridge({ + const bridge = createTestBridge({ providerConfig: { sessionResumption: true }, - onAudio: vi.fn(), - onClearAudio: vi.fn(), onError, }); - const { connecting, socket: firstSocket } = await openRealtimeBridge(bridge, 0, "conv_close"); - await connecting; + const firstSocket = await openRealtimeBridge(bridge, 0, "conv_close"); firstSocket.close(1006, "connection lost"); await vi.advanceTimersByTimeAsync(0); @@ -1939,9 +1615,8 @@ describe("buildXaiRealtimeVoiceProvider", () => { await waitForRealtimeState(() => expect(FakeWebSocket.instances.length).toBe(2)); const reconnectedSocket = requireSocket(1); expect(String(reconnectedSocket.args[0])).not.toContain("conversation_id="); - reconnectedSocket.readyState = FakeWebSocket.OPEN; - reconnectedSocket.emit("open"); - reconnectedSocket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + reconnectedSocket.open(); + reconnectedSocket.emitServer({ type: "session.updated" }); await reconnecting; expect(bridge.isConnected()).toBe(true); @@ -1954,19 +1629,12 @@ describe("buildXaiRealtimeVoiceProvider", () => { vi.useFakeTimers(); vi.stubEnv("XAI_API_KEY", "xai-env"); // pragma: allowlist secret const onError = vi.fn(); - const bridge = buildXaiRealtimeVoiceProvider().createBridge({ + const bridge = createTestBridge({ providerConfig: { apiKey: "xai-test", sessionResumption: true }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), onError, }); - const { connecting, socket: firstSocket } = await openRealtimeBridge( - bridge, - 0, - "conv_replace_retry", - ); - await connecting; + const firstSocket = await openRealtimeBridge(bridge, 0, "conv_replace_retry"); firstSocket.close(1006, "connection lost"); await vi.advanceTimersByTimeAsync(0); @@ -1976,9 +1644,8 @@ describe("buildXaiRealtimeVoiceProvider", () => { await waitForRealtimeState(() => expect(FakeWebSocket.instances).toHaveLength(2)); const replacementSocket = requireSocket(1); expect(String(replacementSocket.args[0])).toContain("conversation_id=conv_replace_retry"); - replacementSocket.readyState = FakeWebSocket.OPEN; - replacementSocket.emit("open"); - replacementSocket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + replacementSocket.open(); + replacementSocket.emitServer({ type: "session.updated" }); await replacementConnect; await vi.advanceTimersByTimeAsync(1000); @@ -1991,26 +1658,17 @@ describe("buildXaiRealtimeVoiceProvider", () => { it("enables xAI session resumption and reconnects with the created conversation id", async () => { vi.useFakeTimers(); vi.stubEnv("XAI_API_KEY", "xai-env"); // pragma: allowlist secret - const provider = buildXaiRealtimeVoiceProvider(); - const bridge = provider.createBridge({ + const bridge = createTestBridge({ providerConfig: { apiKey: "xai-test", sessionResumption: true }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), }); const connecting = bridge.connect(); await waitForRealtimeState(() => expect(FakeWebSocket.instances.length).toBe(1)); const firstSocket = requireSocket(); - firstSocket.readyState = FakeWebSocket.OPEN; - firstSocket.emit("open"); + firstSocket.open(); expect(requireSession(firstSocket).resumption).toEqual({ enabled: true }); - firstSocket.emit( - "message", - Buffer.from( - JSON.stringify({ type: "conversation.created", conversation: { id: "conv_resume" } }), - ), - ); - firstSocket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + firstSocket.emitServer({ type: "conversation.created", conversation: { id: "conv_resume" } }); + firstSocket.emitServer({ type: "session.updated" }); await connecting; firstSocket.close(1006, "connection lost"); @@ -2018,8 +1676,7 @@ describe("buildXaiRealtimeVoiceProvider", () => { await waitForRealtimeState(() => expect(FakeWebSocket.instances.length).toBe(2)); const secondSocket = requireSocket(1); expect(String(secondSocket.args[0])).toContain("conversation_id=conv_resume"); - secondSocket.readyState = FakeWebSocket.OPEN; - secondSocket.emit("open"); + secondSocket.open(); expect(requireSession(secondSocket).resumption).toEqual({ enabled: true }); bridge.close(); }); @@ -2027,19 +1684,15 @@ describe("buildXaiRealtimeVoiceProvider", () => { it("fails closed instead of reconnecting without a conversation id", async () => { vi.useFakeTimers(); vi.stubEnv("XAI_API_KEY", "xai-env"); // pragma: allowlist secret - const provider = buildXaiRealtimeVoiceProvider(); const onEvent = vi.fn(); const onClose = vi.fn(); - const bridge = provider.createBridge({ + const bridge = createTestBridge({ providerConfig: { apiKey: "xai-test", sessionResumption: true }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), onEvent, onClose, }); - const { connecting, socket } = await openRealtimeBridge(bridge); - await connecting; + const socket = await openRealtimeBridge(bridge); socket.close(1006, "connection lost"); await vi.advanceTimersByTimeAsync(1000); @@ -2057,19 +1710,14 @@ describe("buildXaiRealtimeVoiceProvider", () => { it("fails closed instead of reconnecting when xAI session resumption is disabled", async () => { vi.useFakeTimers(); vi.stubEnv("XAI_API_KEY", "xai-env"); // pragma: allowlist secret - const provider = buildXaiRealtimeVoiceProvider(); const onEvent = vi.fn(); const onClose = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "xai-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), + const bridge = createTestBridge({ onEvent, onClose, }); - const { connecting, socket } = await openRealtimeBridge(bridge, 0, "conv_default"); - await connecting; + const socket = await openRealtimeBridge(bridge, 0, "conv_default"); socket.close(1006, "connection lost"); await vi.advanceTimersByTimeAsync(1000); @@ -2087,20 +1735,15 @@ describe("buildXaiRealtimeVoiceProvider", () => { it("does not retry after startup websocket errors", async () => { vi.useFakeTimers(); vi.stubEnv("XAI_API_KEY", "xai-env"); // pragma: allowlist secret - const provider = buildXaiRealtimeVoiceProvider(); const onClose = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "xai-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), + const bridge = createTestBridge({ onClose, }); const connecting = bridge.connect(); await waitForRealtimeState(() => expect(FakeWebSocket.instances.length).toBe(1)); const socket = requireSocket(); - socket.readyState = FakeWebSocket.OPEN; - socket.emit("open"); + socket.open(); socket.emit("error", new Error("bad auth")); await expect(connecting).rejects.toThrow("bad auth"); @@ -2112,9 +1755,7 @@ describe("buildXaiRealtimeVoiceProvider", () => { it("forwards configured provider tools in session.update", async () => { vi.stubEnv("XAI_API_KEY", "xai-env"); // pragma: allowlist secret - const provider = buildXaiRealtimeVoiceProvider(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "xai-test" }, // pragma: allowlist secret + const bridge = createTestBridge({ tools: [ { type: "function", @@ -2123,11 +1764,9 @@ describe("buildXaiRealtimeVoiceProvider", () => { parameters: { type: "object", properties: {} }, }, ], - onAudio: vi.fn(), - onClearAudio: vi.fn(), }); - const { socket } = await openRealtimeBridge(bridge); + const { socket } = await startRealtimeBridge(bridge); bridge.close(); const session = requireSession(socket);