diff --git a/extensions/voice-call/src/webhook/realtime-handler.lifecycle.test.ts b/extensions/voice-call/src/webhook/realtime-handler.lifecycle.test.ts index 6518248b5a80..7336f8f1118b 100644 --- a/extensions/voice-call/src/webhook/realtime-handler.lifecycle.test.ts +++ b/extensions/voice-call/src/webhook/realtime-handler.lifecycle.test.ts @@ -35,7 +35,10 @@ function createRealtimeConfig(): VoiceCallRealtimeConfig { }; } -function createBridge(close: () => void): RealtimeVoiceBridge { +function createBridge( + close: () => void, + overrides: Partial = {}, +): RealtimeVoiceBridge { return { connect: async () => {}, sendAudio: () => {}, @@ -45,6 +48,7 @@ function createBridge(close: () => void): RealtimeVoiceBridge { close, isConnected: () => true, triggerGreeting: () => {}, + ...overrides, }; } @@ -159,4 +163,134 @@ describe("RealtimeCallHandler lifecycle", () => { await server.close(); } }); + + it("releases a hung native consult and ignores its late result after stream teardown", async () => { + let onToolCall: + | ((event: { itemId: string; callId: string; name: string; args: unknown }) => void) + | undefined; + let resolveConsult: ((result: unknown) => void) | undefined; + const submitToolResult = vi.fn(); + const createBridgeForCall = vi.fn( + (request: { + onToolCall?: (event: { + itemId: string; + callId: string; + name: string; + args: unknown; + }) => void; + }) => { + onToolCall = request.onToolCall; + return createBridge(vi.fn(), { + supportsToolResultContinuation: true, + submitToolResult, + }); + }, + ); + const call: CallRecord = { + callId: "call-consult", + providerCallId: "CA-consult", + provider: "twilio", + direction: "inbound", + state: "ringing", + from: "+15550001111", + to: "+15550002222", + startedAt: Date.now(), + transcript: [], + processedEventIds: [], + }; + const handler = new RealtimeCallHandler( + createRealtimeConfig(), + { + processEvent: vi.fn(), + getCallByProviderCallId: vi.fn(() => call), + } as unknown as CallManager, + { + name: "twilio", + verifyWebhook: vi.fn(), + parseWebhookEvent: vi.fn(), + initiateCall: vi.fn(), + hangupCall: vi.fn(), + playTts: vi.fn(), + startListening: vi.fn(), + stopListening: vi.fn(), + getCallStatus: vi.fn(), + } as unknown as VoiceCallProvider, + { + id: "openai", + label: "OpenAI", + isConfigured: () => true, + createBridge: createBridgeForCall, + }, + { apiKey: "test-key" }, + "/voice/webhook", + ); + handler.registerToolHandler( + "openclaw_agent_consult", + async () => + await new Promise((resolve) => { + resolveConsult = resolve; + }), + ); + const { streamUrl } = handler.issueStreamSession(); + const server = await startUpgradeWsServer({ + urlPath: new URL(streamUrl).pathname, + onUpgrade: (request, socket, head) => { + handler.handleWebSocketUpgrade(request, socket, head); + }, + }); + const ws = await connectWs(server.url); + + try { + ws.send( + JSON.stringify({ + event: "start", + start: { streamSid: "MZ-consult", callSid: "CA-consult" }, + }), + ); + await vi.waitFor(() => { + expect(createBridgeForCall).toHaveBeenCalledTimes(1); + }); + + onToolCall?.({ + itemId: "item-consult", + callId: "tool-consult", + name: "openclaw_agent_consult", + args: { question: "Check the deployment." }, + }); + await vi.waitFor(() => { + expect(resolveConsult).toBeTypeOf("function"); + }); + + const consults = ( + handler as unknown as { + nativeConsultsInFlightByCallId: Map; + } + ).nativeConsultsInFlightByCallId; + expect(consults.size).toBe(1); + + const closed = waitForClose(ws); + ws.close(); + await closed; + await vi.waitFor(() => { + expect(consults.size).toBe(0); + }); + + resolveConsult?.({ text: "Deployment is healthy." }); + await new Promise((resolve) => { + setImmediate(resolve); + }); + expect(submitToolResult).toHaveBeenCalledTimes(1); + expect(submitToolResult).not.toHaveBeenCalledWith( + "tool-consult", + { text: "Deployment is healthy." }, + undefined, + ); + } finally { + if (ws.readyState !== WebSocket.CLOSED) { + ws.terminate(); + } + await handler.close(); + await server.close(); + } + }); });