diff --git a/extensions/openai/realtime-voice-provider.test.ts b/extensions/openai/realtime-voice-provider.test.ts index 151cc2f6801a..51705a4088ea 100644 --- a/extensions/openai/realtime-voice-provider.test.ts +++ b/extensions/openai/realtime-voice-provider.test.ts @@ -1,6 +1,10 @@ // Openai tests cover realtime voice provider plugin behavior. import { REALTIME_VOICE_AUDIO_FORMAT_PCM16_24KHZ } from "openclaw/plugin-sdk/realtime-voice"; -import type { RealtimeVoiceBridge, RealtimeVoiceTool } from "openclaw/plugin-sdk/realtime-voice"; +import type { + RealtimeVoiceBridge, + RealtimeVoiceBridgeCreateRequest, + RealtimeVoiceTool, +} from "openclaw/plugin-sdk/realtime-voice"; import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { buildOpenAIRealtimeVoiceProvider } from "./realtime-voice-provider.js"; @@ -192,6 +196,57 @@ function parseSent(socket: FakeWebSocketInstance): SentRealtimeEvent[] { return socket.sent.map((payload: string) => JSON.parse(payload) as SentRealtimeEvent); } +function createNativeBridge( + overrides: Partial = {}, +): RealtimeVoiceBridge { + return buildOpenAIRealtimeVoiceProvider().createBridge({ + providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + onAudio: vi.fn(), + onClearAudio: vi.fn(), + ...overrides, + }); +} + +function requireSocket(index = 0): FakeWebSocketInstance { + const socket = FakeWebSocket.instances[index]; + if (!socket) { + throw new Error("expected bridge to create a websocket"); + } + return socket; +} + +function beginBridgeConnection( + bridge: RealtimeVoiceBridge, + socketIndex = 0, +): { connecting: Promise; socket: FakeWebSocketInstance } { + const connecting = bridge.connect(); + return { connecting, socket: requireSocket(socketIndex) }; +} + +function openSocket(socket: FakeWebSocketInstance): void { + socket.readyState = FakeWebSocket.OPEN; + socket.emit("open"); +} + +function emitServerEvent(socket: FakeWebSocketInstance, event: Record): void { + socket.emit("message", Buffer.from(JSON.stringify(event))); +} + +function emitSessionUpdated(socket: FakeWebSocketInstance): void { + emitServerEvent(socket, { type: "session.updated" }); +} + +async function connectReadyBridge( + bridge: RealtimeVoiceBridge, + socketIndex = 0, +): Promise { + const { connecting, socket } = beginBridgeConnection(bridge, socketIndex); + openSocket(socket); + emitSessionUpdated(socket); + await connecting; + return socket; +} + function expectedResponseCreateEvent() { return expect.objectContaining({ type: "response.create", @@ -1458,32 +1513,23 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("waits for session.updated before draining audio and firing onReady", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); const onReady = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + const bridge = createNativeBridge({ instructions: "Be helpful.", language: "de", - onAudio: vi.fn(), - onClearAudio: vi.fn(), onReady, }); - const connecting = bridge.connect(); + const { connecting, socket } = beginBridgeConnection(bridge); let connectResolved = false; void connecting.then(() => { connectResolved = true; }); - const socket = FakeWebSocket.instances[0]; - if (!socket) { - throw new Error("expected bridge to create a websocket"); - } - socket.readyState = FakeWebSocket.OPEN; - socket.emit("open"); + openSocket(socket); await Promise.resolve(); bridge.sendAudio(Buffer.from("before-ready")); - socket.emit("message", Buffer.from(JSON.stringify({ type: "session.created" }))); + emitServerEvent(socket, { type: "session.created" }); expect(connectResolved).toBe(false); expect(onReady).not.toHaveBeenCalled(); @@ -1507,7 +1553,7 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { expect(session).not.toHaveProperty("temperature"); expect(bridge.isConnected()).toBe(false); - socket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + emitSessionUpdated(socket); await connecting; expect(connectResolved).toBe(true); @@ -1520,25 +1566,15 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("bounds queued audio by aggregate bytes before session readiness", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - }); - 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"); + const bridge = createNativeBridge(); + const { connecting, socket } = beginBridgeConnection(bridge); + openSocket(socket); await Promise.resolve(); bridge.sendAudio(Buffer.alloc(512 * 1024, 0x01)); bridge.sendAudio(Buffer.alloc(512 * 1024, 0x02)); bridge.sendAudio(Buffer.from("overflow")); - socket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + emitSessionUpdated(socket); await connecting; const audioEvents = parseSent(socket).filter( @@ -1552,14 +1588,8 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("discards audio closed before the first connection and reconnects fresh", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); const onClose = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - onClose, - }); + const bridge = createNativeBridge({ onClose }); bridge.sendAudio(Buffer.from("queued-before-connect")); bridge.close(); @@ -1569,14 +1599,9 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { expect(FakeWebSocket.instances).toHaveLength(0); expect(onClose).not.toHaveBeenCalled(); - const connecting = bridge.connect(); - const socket = FakeWebSocket.instances[0]; - if (!socket) { - throw new Error("expected bridge to connect"); - } - socket.readyState = FakeWebSocket.OPEN; - socket.emit("open"); - socket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + const { connecting, socket } = beginBridgeConnection(bridge); + openSocket(socket); + emitSessionUpdated(socket); await connecting; expect( @@ -1589,19 +1614,9 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("does not carry queued audio across terminal close and explicit reconnect", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - }); - const firstConnect = bridge.connect(); - const firstSocket = FakeWebSocket.instances[0]; - if (!firstSocket) { - throw new Error("expected bridge to create a websocket"); - } - firstSocket.readyState = FakeWebSocket.OPEN; - firstSocket.emit("open"); + const bridge = createNativeBridge(); + const { connecting: firstConnect, socket: firstSocket } = beginBridgeConnection(bridge); + openSocket(firstSocket); await Promise.resolve(); bridge.sendAudio(Buffer.from("queued-before-close")); @@ -1609,14 +1624,9 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { await firstConnect; bridge.sendAudio(Buffer.from("sent-after-close")); - const reconnecting = bridge.connect(); - const secondSocket = FakeWebSocket.instances[1]; - if (!secondSocket) { - throw new Error("expected bridge to reconnect"); - } - secondSocket.readyState = FakeWebSocket.OPEN; - secondSocket.emit("open"); - secondSocket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + const { connecting: reconnecting, socket: secondSocket } = beginBridgeConnection(bridge, 1); + openSocket(secondSocket); + emitSessionUpdated(secondSocket); await reconnecting; expect( @@ -1626,25 +1636,15 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("shares an in-flight connection until session readiness", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); const onReady = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - onReady, - }); + const bridge = createNativeBridge({ onReady }); const firstConnect = bridge.connect(); const secondConnect = bridge.connect(); - const socket = FakeWebSocket.instances[0]; - if (!socket) { - throw new Error("expected bridge to create a websocket"); - } + const socket = requireSocket(); expect(FakeWebSocket.instances).toHaveLength(1); - socket.readyState = FakeWebSocket.OPEN; - socket.emit("open"); - socket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + openSocket(socket); + emitSessionUpdated(socket); await Promise.all([firstConnect, secondConnect]); expect(onReady).toHaveBeenCalledOnce(); @@ -1652,31 +1652,22 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("suppresses auto responses before draining queued initial greeting audio", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); const bridgeRef: { current?: RealtimeVoiceBridge } = {}; const onReady = vi.fn(() => { bridgeRef.current?.triggerGreeting?.("Say exactly: hello from explicit speech."); }); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + const bridge = createNativeBridge({ instructions: "Be helpful.", - onAudio: vi.fn(), - onClearAudio: vi.fn(), onReady, }); bridgeRef.current = bridge; - const connecting = bridge.connect(); - const socket = FakeWebSocket.instances[0]; - if (!socket) { - throw new Error("expected bridge to create a websocket"); - } + const { connecting, socket } = beginBridgeConnection(bridge); - socket.readyState = FakeWebSocket.OPEN; - socket.emit("open"); + openSocket(socket); await Promise.resolve(); bridge.sendAudio(Buffer.from("before-ready")); - socket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + emitSessionUpdated(socket); await connecting; const sent = parseSent(socket); @@ -1712,7 +1703,7 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { expect(sent.filter((event) => event.type === "response.create")).toHaveLength(1); expect(onReady).toHaveBeenCalledTimes(1); - socket.emit("message", Buffer.from(JSON.stringify({ type: "response.done" }))); + emitServerEvent(socket, { type: "response.done" }); expectRecordFields( requireNestedRecord(parseSent(socket).at(-1)?.session, ["audio", "input", "turn_detection"]), @@ -1725,9 +1716,7 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("omits unsupported OpenAI tool names from GA session updates", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + const bridge = createNativeBridge({ tools: [ createRealtimeTool("1_lookup"), createRealtimeTool("calendar.lookup:next"), @@ -1737,45 +1726,26 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { createMalformedToolName(42), createUnreadableToolName(), ], - onAudio: vi.fn(), - onClearAudio: vi.fn(), }); - const connecting = bridge.connect(); - const socket = FakeWebSocket.instances[0]; - if (!socket) { - throw new Error("expected bridge to create a websocket"); - } + const { connecting, socket } = beginBridgeConnection(bridge); - socket.readyState = FakeWebSocket.OPEN; - socket.emit("open"); + openSocket(socket); const tools = requireSession(socket).tools as Array<{ name?: string }>; expect(tools.map((tool) => tool.name)).toEqual(["1_lookup", "x".repeat(65)]); - socket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + emitSessionUpdated(socket); await connecting; }); it("rotates realtime bridges on provider max-duration events without reporting an error", async () => { vi.useFakeTimers(); - 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 firstSocket = FakeWebSocket.instances[0]; - if (!firstSocket) { - throw new Error("expected bridge to create a websocket"); - } + const bridge = createNativeBridge({ onError, onEvent }); + const { connecting, socket: firstSocket } = beginBridgeConnection(bridge); - firstSocket.readyState = FakeWebSocket.OPEN; - firstSocket.emit("open"); - firstSocket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + openSocket(firstSocket); + emitSessionUpdated(firstSocket); await connecting; firstSocket.emit( @@ -1803,13 +1773,9 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { await vi.advanceTimersByTimeAsync(1000); await vi.waitFor(() => expect(FakeWebSocket.instances).toHaveLength(2)); - const secondSocket = FakeWebSocket.instances[1]; - if (!secondSocket) { - throw new Error("expected bridge to reconnect"); - } - secondSocket.readyState = FakeWebSocket.OPEN; - secondSocket.emit("open"); - secondSocket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + const secondSocket = requireSocket(1); + openSocket(secondSocket); + emitSessionUpdated(secondSocket); await vi.waitFor(() => expect(onEvent).toHaveBeenCalledWith({ @@ -1832,23 +1798,12 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { it("cancels a pending reconnect and allows a later explicit connect", async () => { vi.useFakeTimers(); - const provider = buildOpenAIRealtimeVoiceProvider(); const onError = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - onError, - }); - const connecting = bridge.connect(); - const socket = FakeWebSocket.instances[0]; - if (!socket) { - throw new Error("expected bridge to create a websocket"); - } + const bridge = createNativeBridge({ onError }); + const { connecting, socket } = beginBridgeConnection(bridge); - socket.readyState = FakeWebSocket.OPEN; - socket.emit("open"); - socket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + openSocket(socket); + emitSessionUpdated(socket); await connecting; socket.readyState = FakeWebSocket.CLOSED; @@ -1863,14 +1818,12 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { expect(FakeWebSocket.instances).toHaveLength(1); expect(onError).not.toHaveBeenCalled(); - const reconnecting = bridge.connect(); - const reconnectedSocket = FakeWebSocket.instances[1]; - if (!reconnectedSocket) { - throw new Error("expected bridge to reconnect after close"); - } - reconnectedSocket.readyState = FakeWebSocket.OPEN; - reconnectedSocket.emit("open"); - reconnectedSocket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + const { connecting: reconnecting, socket: reconnectedSocket } = beginBridgeConnection( + bridge, + 1, + ); + openSocket(reconnectedSocket); + emitSessionUpdated(reconnectedSocket); await reconnecting; expect(bridge.isConnected()).toBe(true); @@ -1881,26 +1834,18 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { it("ignores late events from a socket replaced by reconnect", async () => { vi.useFakeTimers(); - const provider = buildOpenAIRealtimeVoiceProvider(); const onAudio = vi.fn(); const onClose = vi.fn(); const onError = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + const bridge = createNativeBridge({ onAudio, - onClearAudio: vi.fn(), onClose, onError, }); - const connecting = bridge.connect(); - const firstSocket = FakeWebSocket.instances[0]; - if (!firstSocket) { - throw new Error("expected bridge to create a websocket"); - } + const { connecting, socket: firstSocket } = beginBridgeConnection(bridge); - firstSocket.readyState = FakeWebSocket.OPEN; - firstSocket.emit("open"); - firstSocket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + openSocket(firstSocket); + emitSessionUpdated(firstSocket); await connecting; firstSocket.readyState = FakeWebSocket.CLOSED; @@ -1919,16 +1864,12 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { expect(onError).not.toHaveBeenCalled(); await vi.advanceTimersByTimeAsync(1000); - const secondSocket = FakeWebSocket.instances[1]; - if (!secondSocket) { - throw new Error("expected bridge to reconnect"); - } - secondSocket.readyState = FakeWebSocket.OPEN; - secondSocket.emit("open"); - secondSocket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + const secondSocket = requireSocket(1); + openSocket(secondSocket); + emitSessionUpdated(secondSocket); await vi.waitFor(() => expect(bridge.isConnected()).toBe(true)); - firstSocket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + emitSessionUpdated(firstSocket); firstSocket.emit("error", new Error("late socket failure")); firstSocket.emit("close", 1006, Buffer.from("late socket close")); await vi.advanceTimersByTimeAsync(0); @@ -1943,27 +1884,14 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { it("exhausts retries when sockets open but never become provider-ready", async () => { vi.useFakeTimers(); - const provider = buildOpenAIRealtimeVoiceProvider(); const onClose = vi.fn(); 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(), - onClose, - onError, - onEvent, - }); - const connecting = bridge.connect(); - const firstSocket = FakeWebSocket.instances[0]; - if (!firstSocket) { - throw new Error("expected bridge to create a websocket"); - } + const bridge = createNativeBridge({ onClose, onError, onEvent }); + const { connecting, socket: firstSocket } = beginBridgeConnection(bridge); - firstSocket.readyState = FakeWebSocket.OPEN; - firstSocket.emit("open"); - firstSocket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + openSocket(firstSocket); + emitSessionUpdated(firstSocket); await connecting; firstSocket.readyState = FakeWebSocket.CLOSED; @@ -1978,12 +1906,8 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }), ); await vi.advanceTimersByTimeAsync(1000 * 2 ** (attempt - 1)); - const retrySocket = FakeWebSocket.instances[attempt]; - if (!retrySocket) { - throw new Error(`expected reconnect socket ${attempt}`); - } - retrySocket.readyState = FakeWebSocket.OPEN; - retrySocket.emit("open"); + const retrySocket = requireSocket(attempt); + openSocket(retrySocket); retrySocket.emit( "message", Buffer.from( @@ -2010,8 +1934,7 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("keeps Azure deployment bridges on deployment-compatible session payloads", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); - const bridge = provider.createBridge({ + const bridge = createNativeBridge({ providerConfig: { apiKey: "sk-test", // pragma: allowlist secret azureEndpoint: "https://example.openai.azure.com/", @@ -2026,21 +1949,14 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { createRealtimeTool("calendar.lookup:next"), createRealtimeTool("x".repeat(65)), ], - onAudio: vi.fn(), - onClearAudio: vi.fn(), }); - const connecting = bridge.connect(); - const socket = FakeWebSocket.instances[0]; - if (!socket) { - throw new Error("expected bridge to create a websocket"); - } + const { connecting, socket } = beginBridgeConnection(bridge); expect(socket.args[0]).toBe( "wss://example.openai.azure.com/openai/realtime?api-version=2024-10-01-preview&deployment=realtime-prod", ); - socket.readyState = FakeWebSocket.OPEN; - socket.emit("open"); + openSocket(socket); await Promise.resolve(); const session = requireSession(socket); @@ -2065,7 +1981,7 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { const tools = session.tools as Array<{ name?: string }>; expect(tools.map((tool) => tool.name)).toEqual(["1_lookup"]); - socket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + emitSessionUpdated(socket); await connecting; bridge.triggerGreeting?.("Say hello."); @@ -2085,7 +2001,7 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { expectedResponseCreateEvent(), ]); - socket.emit("message", Buffer.from(JSON.stringify({ type: "response.done" }))); + emitServerEvent(socket, { type: "response.done" }); expect(parseSent(socket).at(-1)).toEqual({ type: "session.update", session: { @@ -2101,20 +2017,10 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("rejects connection when session configuration fails before readiness", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - }); - const connecting = bridge.connect(); - const socket = FakeWebSocket.instances[0]; - if (!socket) { - throw new Error("expected bridge to create a websocket"); - } + const bridge = createNativeBridge(); + const { connecting, socket } = beginBridgeConnection(bridge); - socket.readyState = FakeWebSocket.OPEN; - socket.emit("open"); + openSocket(socket); socket.emit( "message", Buffer.from( @@ -2130,24 +2036,12 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("treats pre-ready auth errors as a single startup failure", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); const onError = vi.fn(); const onClose = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - onError, - onClose, - }); - const connecting = bridge.connect(); - const socket = FakeWebSocket.instances[0]; - if (!socket) { - throw new Error("expected bridge to create a websocket"); - } + const bridge = createNativeBridge({ onError, onClose }); + const { connecting, socket } = beginBridgeConnection(bridge); - socket.readyState = FakeWebSocket.OPEN; - socket.emit("open"); + openSocket(socket); socket.emit( "message", Buffer.from( @@ -2175,20 +2069,10 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("normalizes structured direct OpenAI startup auth errors", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - }); - const connecting = bridge.connect(); - const socket = FakeWebSocket.instances[0]; - if (!socket) { - throw new Error("expected bridge to create a websocket"); - } + const bridge = createNativeBridge(); + const { connecting, socket } = beginBridgeConnection(bridge); - socket.readyState = FakeWebSocket.OPEN; - socket.emit("open"); + openSocket(socket); socket.emit( "message", Buffer.from( @@ -2208,17 +2092,8 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("normalizes direct OpenAI socket handshake auth errors", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - }); - const connecting = bridge.connect(); - const socket = FakeWebSocket.instances[0]; - if (!socket) { - throw new Error("expected bridge to create a websocket"); - } + const bridge = createNativeBridge(); + const { connecting, socket } = beginBridgeConnection(bridge); socket.emit("error", new Error("Unexpected server response: 401")); @@ -2243,17 +2118,10 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }, ], ])("preserves %s startup auth errors", async (_label, providerConfig) => { - const provider = buildOpenAIRealtimeVoiceProvider(); - const bridge = provider.createBridge({ + const bridge = createNativeBridge({ providerConfig, - onAudio: vi.fn(), - onClearAudio: vi.fn(), }); - const connecting = bridge.connect(); - const socket = FakeWebSocket.instances[0]; - if (!socket) { - throw new Error("expected bridge to create a websocket"); - } + const { connecting, socket } = beginBridgeConnection(bridge); socket.emit("error", new Error("Unexpected server response: 401")); @@ -2262,23 +2130,12 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("keeps a retried connection ready after delayed startup failure close", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); const onClose = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - onClose, - }); - const failedConnect = bridge.connect(); - const failedSocket = FakeWebSocket.instances[0]; - if (!failedSocket) { - throw new Error("expected bridge to create a websocket"); - } + const bridge = createNativeBridge({ onClose }); + const { connecting: failedConnect, socket: failedSocket } = beginBridgeConnection(bridge); failedSocket.deferClose = true; - failedSocket.readyState = FakeWebSocket.OPEN; - failedSocket.emit("open"); + openSocket(failedSocket); failedSocket.emit( "message", Buffer.from( @@ -2292,14 +2149,9 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { await expect(failedConnect).rejects.toThrow(OPENAI_REALTIME_REJECTED_KEY_MESSAGE); expect(failedSocket.deferredClose).toBeDefined(); - const retryConnect = bridge.connect(); - const retrySocket = FakeWebSocket.instances[1]; - if (!retrySocket) { - throw new Error("expected bridge retry to create a websocket"); - } - retrySocket.readyState = FakeWebSocket.OPEN; - retrySocket.emit("open"); - retrySocket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + const { connecting: retryConnect, socket: retrySocket } = beginBridgeConnection(bridge, 1); + openSocket(retrySocket); + emitSessionUpdated(retrySocket); await retryConnect; expect(bridge.isConnected()).toBe(true); @@ -2309,20 +2161,10 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("rejects connection when the socket closes before session readiness", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - }); - const connecting = bridge.connect(); - const socket = FakeWebSocket.instances[0]; - if (!socket) { - throw new Error("expected bridge to create a websocket"); - } + const bridge = createNativeBridge(); + const { connecting, socket } = beginBridgeConnection(bridge); - socket.readyState = FakeWebSocket.OPEN; - socket.emit("open"); + openSocket(socket); socket.close(1006, "session closed"); await expect(connecting).rejects.toThrow("OpenAI realtime connection closed before ready"); @@ -2331,19 +2173,9 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { it("does not report startup timeout shutdown as a clean close", async () => { vi.useFakeTimers(); - const provider = buildOpenAIRealtimeVoiceProvider(); const onClose = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - onClose, - }); - const connecting = bridge.connect(); - const socket = FakeWebSocket.instances[0]; - if (!socket) { - throw new Error("expected bridge to create a websocket"); - } + const bridge = createNativeBridge({ onClose }); + const { connecting, socket } = beginBridgeConnection(bridge); const timeoutAssertion = expect(connecting).rejects.toThrow( "OpenAI realtime connection timeout", @@ -2357,22 +2189,13 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("can disable automatic audio turn responses for agent-routed voice loops", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + const bridge = createNativeBridge({ autoRespondToAudio: false, - onAudio: vi.fn(), - onClearAudio: vi.fn(), }); - const connecting = bridge.connect(); - const socket = FakeWebSocket.instances[0]; - if (!socket) { - throw new Error("expected bridge to create a websocket"); - } + const { connecting, socket } = beginBridgeConnection(bridge); - socket.readyState = FakeWebSocket.OPEN; - socket.emit("open"); - socket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + openSocket(socket); + emitSessionUpdated(socket); await connecting; expectRecordFields( @@ -2386,24 +2209,11 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("can disable realtime response interruption while keeping audio responses enabled", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + const bridge = createNativeBridge({ autoRespondToAudio: true, interruptResponseOnInputAudio: false, - onAudio: vi.fn(), - onClearAudio: vi.fn(), }); - 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; + const socket = await connectReadyBridge(bridge); expectRecordFields( requireNestedRecord(requireSession(socket), ["audio", "input", "turn_detection"]), @@ -2416,26 +2226,15 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("does not locally clear playback on speech-start events when input interruption is disabled", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); const onAudio = vi.fn(); const onClearAudio = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + const bridge = createNativeBridge({ autoRespondToAudio: true, interruptResponseOnInputAudio: false, onAudio, onClearAudio, }); - 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; + const socket = await connectReadyBridge(bridge); socket.emit( "message", @@ -2463,25 +2262,14 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("keeps assistant playback active on server VAD when automatic audio responses are disabled", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); const onAudio = vi.fn(); const onClearAudio = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + const bridge = createNativeBridge({ autoRespondToAudio: false, onAudio, onClearAudio, }); - 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; + const socket = await connectReadyBridge(bridge); socket.emit( "message", @@ -2509,24 +2297,10 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("can request PCM16 24 kHz realtime audio for Chrome command-pair bridges", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + const bridge = createNativeBridge({ audioFormat: REALTIME_VOICE_AUDIO_FORMAT_PCM16_24KHZ, - onAudio: vi.fn(), - onClearAudio: vi.fn(), }); - - 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; + const socket = await connectReadyBridge(bridge); const session = requireSession(socket); expect(requireNestedRecord(session, ["audio", "input", "format"])).toEqual({ @@ -2540,19 +2314,9 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("settles cleanly when closed before the websocket opens", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); const onClose = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - onClose, - }); - const connecting = bridge.connect(); - const socket = FakeWebSocket.instances[0]; - if (!socket) { - throw new Error("expected bridge to create a websocket"); - } + const bridge = createNativeBridge({ onClose }); + const { connecting, socket } = beginBridgeConnection(bridge); bridge.close(); bridge.close(); @@ -2565,25 +2329,14 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("truncates externally interrupted playback after an immediate mark acknowledgement", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); const onAudio = vi.fn(); const onClearAudio = vi.fn(); - const bridge: ReturnType = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + const bridge = createNativeBridge({ onAudio, onClearAudio, onMark: () => bridge.acknowledgeMark(), }); - 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; + const socket = await connectReadyBridge(bridge); bridge.setMediaTimestamp(1000); socket.emit( @@ -2618,25 +2371,13 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("preserves FIFO playback acknowledgements after sustained output", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); const onClearAudio = vi.fn(); const onMark = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), + const bridge = createNativeBridge({ onClearAudio, onMark, }); - 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; + const socket = await connectReadyBridge(bridge); bridge.setMediaTimestamp(1000); for (let index = 0; index < 300; index += 1) { @@ -2697,24 +2438,9 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("treats a later named mark as cumulative playback progress", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); const onMark = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - onMark, - }); - 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; + const bridge = createNativeBridge({ onMark }); + const socket = await connectReadyBridge(bridge); bridge.setMediaTimestamp(1000); for (let index = 0; index < 3; index += 1) { @@ -2745,25 +2471,13 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("forwards current realtime output audio events", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); const onAudio = vi.fn(); const onTranscript = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + const bridge = createNativeBridge({ onAudio, - 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; + const socket = await connectReadyBridge(bridge); const audio = Buffer.from("assistant audio"); socket.emit( @@ -2795,25 +2509,10 @@ 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; + const bridge = createNativeBridge({ onError, onEvent }); + const socket = await connectReadyBridge(bridge); socket.emit( "message", @@ -2838,23 +2537,9 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); 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; + const bridge = createNativeBridge({ onTranscript }); + const socket = await connectReadyBridge(bridge); socket.emit( "message", @@ -2875,26 +2560,15 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { ["invalid alphabet", "not-base64!"], ["non-canonical pad bits", "ZE=="], ])("terminates the session for %s in output audio", async (_scenario, delta) => { - const provider = buildOpenAIRealtimeVoiceProvider(); const onAudio = vi.fn(); const onError = vi.fn(); const onClose = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + const bridge = createNativeBridge({ onAudio, - onClearAudio: vi.fn(), onError, onClose, }); - 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; + const socket = await connectReadyBridge(bridge); socket.emit( "message", @@ -2921,25 +2595,13 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("forwards Codex-compatible legacy realtime audio and transcript events", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); const onAudio = vi.fn(); const onTranscript = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + const bridge = createNativeBridge({ onAudio, - 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; + const socket = await connectReadyBridge(bridge); const audio = Buffer.from("legacy assistant audio"); socket.emit( @@ -2988,26 +2650,10 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("emits tool calls from realtime conversation item done events", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); const onToolCall = vi.fn(); const onEvent = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - onToolCall, - 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; + const bridge = createNativeBridge({ onToolCall, onEvent }); + const socket = await connectReadyBridge(bridge); socket.emit( "message", @@ -3039,24 +2685,9 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("deduplicates tool calls reported by arguments done and item done events", async () => { - 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; + const bridge = createNativeBridge({ onToolCall }); + const socket = await connectReadyBridge(bridge); socket.emit( "message", @@ -3128,23 +2759,9 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { ])( "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; + const bridge = createNativeBridge({ onToolCall }); + const socket = await connectReadyBridge(bridge); socket.emit( "message", @@ -3181,24 +2798,9 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { ); it("creates an explicit user item and response for manual speech", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); const onEvent = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - 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; + const bridge = createNativeBridge({ onEvent }); + const socket = await connectReadyBridge(bridge); bridge.triggerGreeting?.("Say exactly: hello from explicit speech."); @@ -3229,7 +2831,7 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { expect(onEvent).toHaveBeenCalledWith({ direction: "client", type: "conversation.item.create" }); expect(onEvent).toHaveBeenCalledWith({ direction: "client", type: "response.create" }); - socket.emit("message", Buffer.from(JSON.stringify({ type: "response.done" }))); + emitServerEvent(socket, { type: "response.done" }); expectRecordFields( requireNestedRecord(parseSent(socket).at(-1)?.session, ["audio", "input", "turn_detection"]), @@ -3242,22 +2844,8 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("defers manual response.create while a realtime response is active", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - }); - 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; + const bridge = createNativeBridge(); + const socket = await connectReadyBridge(bridge); socket.emit( "message", Buffer.from(JSON.stringify({ type: "response.created", response: { id: "resp_1" } })), @@ -3276,30 +2864,15 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }, ]); - socket.emit("message", Buffer.from(JSON.stringify({ type: "response.done" }))); + emitServerEvent(socket, { type: "response.done" }); expect(parseSent(socket).slice(-1)).toEqual([expectedResponseCreateEvent()]); }); it("restores automatic audio responses when a manual response is rejected", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); const onError = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - onError, - }); - 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; + const bridge = createNativeBridge({ onError }); + const socket = await connectReadyBridge(bridge); bridge.triggerGreeting?.("Say exactly: hello from explicit speech."); @@ -3344,24 +2917,9 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("keeps automatic audio suppressed for unrelated errors during a manual response", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); const onError = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - onError, - }); - 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; + const bridge = createNativeBridge({ onError }); + const socket = await connectReadyBridge(bridge); bridge.triggerGreeting?.("Say exactly: hello from explicit speech."); const sessionUpdatesBeforeError = parseSent(socket).filter( @@ -3383,7 +2941,7 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { sessionUpdatesBeforeError.length, ); - socket.emit("message", Buffer.from(JSON.stringify({ type: "response.done" }))); + emitServerEvent(socket, { type: "response.done" }); expectRecordFields( requireNestedRecord(parseSent(socket).at(-1)?.session, ["audio", "input", "turn_detection"]), @@ -3396,24 +2954,9 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("flushes a queued manual response after the prior request is rejected", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); const onError = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - onError, - }); - 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; + const bridge = createNativeBridge({ onError }); + const socket = await connectReadyBridge(bridge); bridge.triggerGreeting?.("Say exactly: first greeting."); const firstResponseCreate = parseSent(socket).findLast( @@ -3449,7 +2992,7 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { ); expect(onError).toHaveBeenCalledWith(new Error("bad response request")); - socket.emit("message", Buffer.from(JSON.stringify({ type: "response.done" }))); + emitServerEvent(socket, { type: "response.done" }); expectRecordFields( requireNestedRecord(parseSent(socket).at(-1)?.session, ["audio", "input", "turn_detection"]), @@ -3462,22 +3005,8 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("does not request a realtime response for continuing tool results", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - }); - 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; + const bridge = createNativeBridge(); + const socket = await connectReadyBridge(bridge); void bridge.submitToolResult("call_1", { status: "working" }, { willContinue: true }); @@ -3511,28 +3040,14 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { "message", Buffer.from(JSON.stringify({ type: "response.created", response: { id: "resp_2" } })), ); - socket.emit("message", Buffer.from(JSON.stringify({ type: "response.done" }))); + emitServerEvent(socket, { type: "response.done" }); expect(parseSent(socket).filter((event) => event.type === "response.create")).toHaveLength(1); }); it("does not request a realtime response for suppressed tool results", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - }); - 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; + const bridge = createNativeBridge(); + const socket = await connectReadyBridge(bridge); void bridge.submitToolResult( "call_1", @@ -3554,31 +3069,16 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("does not flush deferred response.create while a tool result is still continuing", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); const onError = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - onError, - }); - 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; + const bridge = createNativeBridge({ onError }); + const socket = await connectReadyBridge(bridge); socket.emit( "message", Buffer.from(JSON.stringify({ type: "response.created", response: { id: "resp_1" } })), ); void bridge.submitToolResult("call_1", { status: "working" }, { willContinue: true }); - socket.emit("message", Buffer.from(JSON.stringify({ type: "response.done" }))); + emitServerEvent(socket, { type: "response.done" }); expect(onError).not.toHaveBeenCalled(); expect(parseSent(socket).filter((event) => event.type === "response.create")).toEqual([]); @@ -3600,22 +3100,8 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("drains deferred response.create after response.cancelled", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - }); - 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; + const bridge = createNativeBridge(); + const socket = await connectReadyBridge(bridge); socket.emit( "message", Buffer.from(JSON.stringify({ type: "response.created", response: { id: "resp_1" } })), @@ -3628,24 +3114,9 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("does not send duplicate response.cancel while cancellation is pending", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); const onEvent = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - 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; + const bridge = createNativeBridge({ onEvent }); + const socket = await connectReadyBridge(bridge); socket.emit( "message", Buffer.from(JSON.stringify({ type: "response.created", response: { id: "resp_1" } })), @@ -3680,25 +3151,13 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("ignores zero-length playback barge-in without clearing audio", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); const onClearAudio = vi.fn(); const onEvent = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), + const bridge = createNativeBridge({ onClearAudio, 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; + const socket = await connectReadyBridge(bridge); bridge.setMediaTimestamp(1000); socket.emit( "message", @@ -3730,25 +3189,13 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("force-cancels zero-length playback barge-in for agent handoff fallback", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); const onClearAudio = vi.fn(); const onEvent = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), + const bridge = createNativeBridge({ onClearAudio, 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; + const socket = await connectReadyBridge(bridge); bridge.setMediaTimestamp(1000); socket.emit( "message", @@ -3785,26 +3232,15 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("allows immediate playback barge-in when the minimum audio window is zero", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); const onClearAudio = vi.fn(); - const bridge = provider.createBridge({ + const bridge = createNativeBridge({ providerConfig: { apiKey: "sk-test", // pragma: allowlist secret minBargeInAudioEndMs: 0, }, - onAudio: vi.fn(), onClearAudio, }); - 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; + const socket = await connectReadyBridge(bridge); bridge.setMediaTimestamp(1000); socket.emit( "message", @@ -3836,24 +3272,9 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("drains deferred response.create after a no-active-response cancellation error", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); const onError = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - onError, - }); - 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; + const bridge = createNativeBridge({ onError }); + const socket = await connectReadyBridge(bridge); socket.emit( "message", Buffer.from(JSON.stringify({ type: "response.created", response: { id: "resp_1" } })), @@ -3885,24 +3306,9 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("ignores a stale cancellation error after a newer manual response starts", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); const onError = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - onError, - }); - 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; + const bridge = createNativeBridge({ onError }); + const socket = await connectReadyBridge(bridge); bridge.setMediaTimestamp(1000); socket.emit( "message", @@ -3928,7 +3334,7 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { throw new Error("expected response.cancel event id"); } void bridge.submitToolResult("call_1", { text: "done" }); - socket.emit("message", Buffer.from(JSON.stringify({ type: "response.done" }))); + emitServerEvent(socket, { type: "response.done" }); const sessionUpdateCount = parseSent(socket).filter( (event) => event.type === "session.update", ).length; @@ -3952,7 +3358,7 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { ); expect(parseSent(socket).at(-1)).toEqual(expectedResponseCreateEvent()); - socket.emit("message", Buffer.from(JSON.stringify({ type: "response.done" }))); + emitServerEvent(socket, { type: "response.done" }); expectRecordFields( requireNestedRecord(parseSent(socket).at(-1)?.session, ["audio", "input", "turn_detection"]), "restored turn detection", @@ -3965,22 +3371,8 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { it("resets deferred response guards after websocket reconnect", async () => { vi.useFakeTimers(); - const provider = buildOpenAIRealtimeVoiceProvider(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - }); - 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; + const bridge = createNativeBridge(); + const socket = await connectReadyBridge(bridge); socket.emit( "message", Buffer.from(JSON.stringify({ type: "response.created", response: { id: "resp_1" } })), @@ -3991,14 +3383,9 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { socket.emit("close", 1006, Buffer.from("transient drop")); await vi.advanceTimersByTimeAsync(1000); - const reconnectedSocket = FakeWebSocket.instances[1]; - if (!reconnectedSocket) { - throw new Error("expected bridge to reconnect"); - } - - reconnectedSocket.readyState = FakeWebSocket.OPEN; - reconnectedSocket.emit("open"); - reconnectedSocket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + const reconnectedSocket = requireSocket(1); + openSocket(reconnectedSocket); + emitSessionUpdated(reconnectedSocket); bridge.sendUserMessage?.("Say hello after reconnect."); expect(parseSent(reconnectedSocket).slice(-3)).toEqual([ @@ -4016,24 +3403,9 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }); it("turns active-response errors into a deferred response.create retry", async () => { - const provider = buildOpenAIRealtimeVoiceProvider(); const onError = vi.fn(); - const bridge = provider.createBridge({ - providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret - onAudio: vi.fn(), - onClearAudio: vi.fn(), - onError, - }); - 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; + const bridge = createNativeBridge({ onError }); + const socket = await connectReadyBridge(bridge); void bridge.submitToolResult("call_1", { text: "done" }); const responseCreateEvent = parseSent(socket).findLast( @@ -4065,12 +3437,12 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { }, ); - socket.emit("message", Buffer.from(JSON.stringify({ type: "response.done" }))); + emitServerEvent(socket, { type: "response.done" }); expect(onError).not.toHaveBeenCalled(); expect(parseSent(socket).slice(-1)).toEqual([expectedResponseCreateEvent()]); - socket.emit("message", Buffer.from(JSON.stringify({ type: "response.done" }))); + emitServerEvent(socket, { type: "response.done" }); expectRecordFields( requireNestedRecord(parseSent(socket).at(-1)?.session, ["audio", "input", "turn_detection"]),