diff --git a/ui/src/ui/realtime-talk-webrtc.test.ts b/ui/src/ui/realtime-talk-webrtc.test.ts index 612a2777412b..b392cdc5d412 100644 --- a/ui/src/ui/realtime-talk-webrtc.test.ts +++ b/ui/src/ui/realtime-talk-webrtc.test.ts @@ -64,6 +64,109 @@ function requireTalkEvent( return event as Record; } +type SentRealtimeEvent = { + type?: string; + item?: { + type?: string; + [key: string]: unknown; + }; + [key: string]: unknown; +}; + +function stubAnswerSdpFetch(): void { + vi.stubGlobal("fetch", vi.fn(async () => new Response("answer-sdp")) as unknown as typeof fetch); +} + +function createOpenAiTransport( + client: Record = {}, + callbacks: Record = {}, +): WebRtcSdpRealtimeTalkTransport { + return new WebRtcSdpRealtimeTalkTransport( + { + provider: "openai", + transport: "webrtc", + clientSecret: "client-secret-123", + }, + { + client: client as never, + sessionKey: "main", + callbacks: callbacks as never, + }, + ); +} + +function dispatchRealtimeEvent(peer: FakePeerConnection | undefined, event: unknown): void { + peer?.channel.dispatchEvent( + new MessageEvent("message", { + data: JSON.stringify(event), + }), + ); +} + +function dispatchConsultToolCall(peer: FakePeerConnection | undefined): void { + dispatchRealtimeEvent(peer, { + type: "response.function_call_arguments.done", + item_id: "item-1", + call_id: "call-1", + name: REALTIME_VOICE_AGENT_CONSULT_TOOL_NAME, + arguments: JSON.stringify({ question: "status?" }), + }); +} + +function dispatchTranscription(peer: FakePeerConnection | undefined, transcript: string): void { + dispatchRealtimeEvent(peer, { + type: "conversation.item.input_audio_transcription.completed", + item_id: "input-1", + transcript, + }); +} + +async function startActiveConsult( + request: ReturnType, + options: { responseAlreadyActive?: boolean } = {}, +): Promise<{ transport: WebRtcSdpRealtimeTalkTransport; peer: FakePeerConnection | undefined }> { + const transport = createOpenAiTransport({ + addEventListener: vi.fn(() => () => undefined), + request, + }); + + await transport.start(); + const peer = FakePeerConnection.instances[0]; + if (options.responseAlreadyActive) { + dispatchRealtimeEvent(peer, { type: "response.created" }); + } + dispatchConsultToolCall(peer); + await vi.waitFor(() => + expect(request).toHaveBeenCalledWith("talk.client.toolCall", expect.any(Object)), + ); + + return { transport, peer }; +} + +function sentRealtimeEvents(peer: FakePeerConnection | undefined): SentRealtimeEvent[] { + return ( + peer?.channel.send.mock.calls.map( + ([payload]) => JSON.parse(String(payload)) as SentRealtimeEvent, + ) ?? [] + ); +} + +function expectSpokenStatusMessage(events: SentRealtimeEvent[], message: string): void { + expect(events).toContainEqual({ + type: "conversation.item.create", + item: { + type: "message", + role: "user", + content: [ + { + type: "input_text", + text: expect.stringContaining(`Status: "${message}"`), + }, + ], + }, + }); +} + describe("WebRtcSdpRealtimeTalkTransport", () => { afterEach(() => { vi.unstubAllGlobals(); @@ -335,10 +438,7 @@ describe("WebRtcSdpRealtimeTalkTransport", () => { }); it("sends spoken active-control acknowledgements through the OpenAI data channel", async () => { - vi.stubGlobal( - "fetch", - vi.fn(async () => new Response("answer-sdp")) as unknown as typeof fetch, - ); + stubAnswerSdpFetch(); const request = vi.fn(async (method: string) => { if (method === "talk.client.toolCall") { return { runId: "run-1" }; @@ -357,76 +457,21 @@ describe("WebRtcSdpRealtimeTalkTransport", () => { } throw new Error(`unexpected request: ${method}`); }); - const transport = new WebRtcSdpRealtimeTalkTransport( - { - provider: "openai", - transport: "webrtc", - clientSecret: "client-secret-123", - }, - { - client: { - addEventListener: vi.fn(() => () => undefined), - request, - } as never, - sessionKey: "main", - callbacks: {}, - }, - ); + const { transport, peer } = await startActiveConsult(request); - await transport.start(); - const peer = FakePeerConnection.instances[0]; - peer?.channel.dispatchEvent( - new MessageEvent("message", { - data: JSON.stringify({ - type: "response.function_call_arguments.done", - item_id: "item-1", - call_id: "call-1", - name: REALTIME_VOICE_AGENT_CONSULT_TOOL_NAME, - arguments: JSON.stringify({ question: "status?" }), - }), - }), - ); - await vi.waitFor(() => - expect(request).toHaveBeenCalledWith("talk.client.toolCall", expect.any(Object)), - ); - - peer?.channel.dispatchEvent( - new MessageEvent("message", { - data: JSON.stringify({ - type: "conversation.item.input_audio_transcription.completed", - item_id: "input-1", - transcript: "status", - }), - }), - ); + dispatchTranscription(peer, "status"); await vi.waitFor(() => expect(request).toHaveBeenCalledWith("talk.client.steer", expect.any(Object)), ); - const sent = - peer?.channel.send.mock.calls.map(([payload]) => JSON.parse(String(payload))) ?? []; - expect(sent).toContainEqual({ - type: "conversation.item.create", - item: { - type: "message", - role: "user", - content: [ - { - type: "input_text", - text: expect.stringContaining('Status: "OpenClaw is working in read (running)."'), - }, - ], - }, - }); + const sent = sentRealtimeEvents(peer); + expectSpokenStatusMessage(sent, "OpenClaw is working in read (running)."); expect(sent).toContainEqual({ type: "response.create" }); transport.stop(); }); it("defers spoken active-control response creation until the active OpenAI response ends", async () => { - vi.stubGlobal( - "fetch", - vi.fn(async () => new Response("answer-sdp")) as unknown as typeof fetch, - ); + stubAnswerSdpFetch(); const request = vi.fn(async (method: string) => { if (method === "talk.client.toolCall") { return { runId: "run-1" }; @@ -445,90 +490,29 @@ describe("WebRtcSdpRealtimeTalkTransport", () => { } throw new Error(`unexpected request: ${method}`); }); - const transport = new WebRtcSdpRealtimeTalkTransport( - { - provider: "openai", - transport: "webrtc", - clientSecret: "client-secret-123", - }, - { - client: { - addEventListener: vi.fn(() => () => undefined), - request, - } as never, - sessionKey: "main", - callbacks: {}, - }, - ); + const { transport, peer } = await startActiveConsult(request, { + responseAlreadyActive: true, + }); - await transport.start(); - const peer = FakePeerConnection.instances[0]; - peer?.channel.dispatchEvent( - new MessageEvent("message", { - data: JSON.stringify({ type: "response.created" }), - }), - ); - peer?.channel.dispatchEvent( - new MessageEvent("message", { - data: JSON.stringify({ - type: "response.function_call_arguments.done", - item_id: "item-1", - call_id: "call-1", - name: REALTIME_VOICE_AGENT_CONSULT_TOOL_NAME, - arguments: JSON.stringify({ question: "status?" }), - }), - }), - ); - await vi.waitFor(() => - expect(request).toHaveBeenCalledWith("talk.client.toolCall", expect.any(Object)), - ); - - peer?.channel.dispatchEvent( - new MessageEvent("message", { - data: JSON.stringify({ - type: "conversation.item.input_audio_transcription.completed", - item_id: "input-1", - transcript: "status", - }), - }), - ); + dispatchTranscription(peer, "status"); await vi.waitFor(() => expect(request).toHaveBeenCalledWith("talk.client.steer", expect.any(Object)), ); - let sent = peer?.channel.send.mock.calls.map(([payload]) => JSON.parse(String(payload))) ?? []; + let sent = sentRealtimeEvents(peer); expect(sent).toContainEqual({ type: "response.cancel" }); - expect(sent).toContainEqual({ - type: "conversation.item.create", - item: { - type: "message", - role: "user", - content: [ - { - type: "input_text", - text: expect.stringContaining('Status: "OpenClaw is working in read (running)."'), - }, - ], - }, - }); + expectSpokenStatusMessage(sent, "OpenClaw is working in read (running)."); expect(sent.filter((event) => event.type === "response.create")).toHaveLength(0); - peer?.channel.dispatchEvent( - new MessageEvent("message", { - data: JSON.stringify({ type: "response.done", response: { status: "completed" } }), - }), - ); + dispatchRealtimeEvent(peer, { type: "response.done", response: { status: "completed" } }); - sent = peer?.channel.send.mock.calls.map(([payload]) => JSON.parse(String(payload))) ?? []; + sent = sentRealtimeEvents(peer); expect(sent.filter((event) => event.type === "response.create")).toHaveLength(1); transport.stop(); }); it("replaces stale OpenAI output with a spoken active-control steering acknowledgement", async () => { - vi.stubGlobal( - "fetch", - vi.fn(async () => new Response("answer-sdp")) as unknown as typeof fetch, - ); + stubAnswerSdpFetch(); const request = vi.fn(async (method: string) => { if (method === "talk.client.toolCall") { return { runId: "run-1" }; @@ -548,82 +532,24 @@ describe("WebRtcSdpRealtimeTalkTransport", () => { } throw new Error(`unexpected request: ${method}`); }); - const transport = new WebRtcSdpRealtimeTalkTransport( - { - provider: "openai", - transport: "webrtc", - clientSecret: "client-secret-123", - }, - { - client: { - addEventListener: vi.fn(() => () => undefined), - request, - } as never, - sessionKey: "main", - callbacks: {}, - }, - ); + const { transport, peer } = await startActiveConsult(request, { + responseAlreadyActive: true, + }); - await transport.start(); - const peer = FakePeerConnection.instances[0]; - peer?.channel.dispatchEvent( - new MessageEvent("message", { - data: JSON.stringify({ type: "response.created" }), - }), - ); - peer?.channel.dispatchEvent( - new MessageEvent("message", { - data: JSON.stringify({ - type: "response.function_call_arguments.done", - item_id: "item-1", - call_id: "call-1", - name: REALTIME_VOICE_AGENT_CONSULT_TOOL_NAME, - arguments: JSON.stringify({ question: "status?" }), - }), - }), - ); - await vi.waitFor(() => - expect(request).toHaveBeenCalledWith("talk.client.toolCall", expect.any(Object)), - ); - - peer?.channel.dispatchEvent( - new MessageEvent("message", { - data: JSON.stringify({ - type: "conversation.item.input_audio_transcription.completed", - item_id: "input-1", - transcript: "actually focus on WebUI", - }), - }), - ); + dispatchTranscription(peer, "actually focus on WebUI"); await vi.waitFor(() => expect(request).toHaveBeenCalledWith("talk.client.steer", expect.any(Object)), ); - const sent = - peer?.channel.send.mock.calls.map(([payload]) => JSON.parse(String(payload))) ?? []; + const sent = sentRealtimeEvents(peer); expect(sent).toContainEqual({ type: "response.cancel" }); - expect(sent).toContainEqual({ - type: "conversation.item.create", - item: { - type: "message", - role: "user", - content: [ - { - type: "input_text", - text: expect.stringContaining('Status: "Got it. I steered the active run."'), - }, - ], - }, - }); + expectSpokenStatusMessage(sent, "Got it. I steered the active run."); expect(sent.some((event) => event.type === "response.create")).toBe(false); transport.stop(); }); it("interrupts stale OpenAI output when active-control cancel is suppressed", async () => { - vi.stubGlobal( - "fetch", - vi.fn(async () => new Response("answer-sdp")) as unknown as typeof fetch, - ); + stubAnswerSdpFetch(); const request = vi.fn(async (method: string) => { if (method === "talk.client.toolCall") { return { runId: "run-1" }; @@ -643,59 +569,16 @@ describe("WebRtcSdpRealtimeTalkTransport", () => { } throw new Error(`unexpected request: ${method}`); }); - const transport = new WebRtcSdpRealtimeTalkTransport( - { - provider: "openai", - transport: "webrtc", - clientSecret: "client-secret-123", - }, - { - client: { - addEventListener: vi.fn(() => () => undefined), - request, - } as never, - sessionKey: "main", - callbacks: {}, - }, - ); + const { transport, peer } = await startActiveConsult(request, { + responseAlreadyActive: true, + }); - await transport.start(); - const peer = FakePeerConnection.instances[0]; - peer?.channel.dispatchEvent( - new MessageEvent("message", { - data: JSON.stringify({ type: "response.created" }), - }), - ); - peer?.channel.dispatchEvent( - new MessageEvent("message", { - data: JSON.stringify({ - type: "response.function_call_arguments.done", - item_id: "item-1", - call_id: "call-1", - name: REALTIME_VOICE_AGENT_CONSULT_TOOL_NAME, - arguments: JSON.stringify({ question: "status?" }), - }), - }), - ); - await vi.waitFor(() => - expect(request).toHaveBeenCalledWith("talk.client.toolCall", expect.any(Object)), - ); - - peer?.channel.dispatchEvent( - new MessageEvent("message", { - data: JSON.stringify({ - type: "conversation.item.input_audio_transcription.completed", - item_id: "input-1", - transcript: "cancel that", - }), - }), - ); + dispatchTranscription(peer, "cancel that"); await vi.waitFor(() => expect(request).toHaveBeenCalledWith("talk.client.steer", expect.any(Object)), ); - const sent = - peer?.channel.send.mock.calls.map(([payload]) => JSON.parse(String(payload))) ?? []; + const sent = sentRealtimeEvents(peer); expect(sent).toContainEqual({ type: "response.cancel" }); expect( sent.some(