test(voice-call): cover late native consult results

This commit is contained in:
Vincent Koc
2026-07-31 04:47:25 +08:00
parent 1fac5d6938
commit 645929da80
@@ -35,7 +35,10 @@ function createRealtimeConfig(): VoiceCallRealtimeConfig {
};
}
function createBridge(close: () => void): RealtimeVoiceBridge {
function createBridge(
close: () => void,
overrides: Partial<RealtimeVoiceBridge> = {},
): 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<unknown>((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<string, unknown>;
}
).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<void>((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();
}
});
});