import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, mock, test, } from "bun:test"; import type { StreamingTranscriber, SttStreamServerEvent, } from "../../stt/types.js"; // `mock.module` is process-global in Bun and leaks into sibling files that // run later in the same `bun test` invocation, so the STT/preflight stubs // delegate to the real implementation unless this file's tests are active // (`shellMocksActive`, toggled in beforeAll/afterAll). The real exports are // snapshotted into plain objects NOW, before the stubs register — a module // namespace is a live view, so reading the real export after the stub // installs would resolve back to the stub (infinite recursion). let shellMocksActive = false; const realSttResolveModule = { ...(await import("../../providers/speech-to-text/resolve.js")), }; const realPreflightModule = { ...(await import("../live-voice-credential-preflight.js")), }; class MockStreamingTranscriber implements StreamingTranscriber { readonly providerId = "deepgram" as const; readonly boundaryId = "daemon-streaming" as const; readonly audioChunks: number[][] = []; readonly mimeTypes: string[] = []; started = false; private onEvent: ((event: SttStreamServerEvent) => void) | null = null; async start(onEvent: (event: SttStreamServerEvent) => void): Promise { this.started = true; this.onEvent = onEvent; } sendAudio(audio: Buffer, mimeType: string): void { this.audioChunks.push([...audio]); this.mimeTypes.push(mimeType); this.onEvent?.({ type: "partial", text: `partial-${this.audioChunks.length}`, }); } stop(): void { this.onEvent?.({ type: "closed" }); } } const resolvedTranscribers: MockStreamingTranscriber[] = []; function createResolvedTranscriber(): MockStreamingTranscriber { const transcriber = new MockStreamingTranscriber(); resolvedTranscribers.push(transcriber); return transcriber; } let resolveStreamingTranscriberImpl = async () => createResolvedTranscriber(); const resolveStreamingTranscriberMock = mock( ( options?: Parameters< typeof realSttResolveModule.resolveStreamingTranscriber >[0], ) => shellMocksActive ? resolveStreamingTranscriberImpl() : realSttResolveModule.resolveStreamingTranscriber(options), ); mock.module("../../providers/speech-to-text/resolve.js", () => ({ ...realSttResolveModule, resolveStreamingTranscriber: resolveStreamingTranscriberMock, })); // The shell tests exercise WebSocket frame routing and session locking, not // credential resolution — the preflight is stubbed ready so start frames // reach the session legs mocked above. mock.module("../live-voice-credential-preflight.js", () => ({ ...realPreflightModule, resolveLiveVoiceCredentialReadiness: async () => shellMocksActive ? { status: "ready" } : realPreflightModule.resolveLiveVoiceCredentialReadiness(), })); import { CURRENT_POLICY_EPOCH } from "../../runtime/auth/policy.js"; import { mintToken } from "../../runtime/auth/token-service.js"; import { RuntimeHttpServer } from "../../runtime/http-server.js"; type JsonFrame = Record; const savedAuthEnv = { DISABLE_HTTP_AUTH: process.env.DISABLE_HTTP_AUTH, }; function mintGatewayToken(): string { return mintToken({ aud: "vellum-daemon", sub: "svc:gateway:self", scope_profile: "gateway_ingress_v1", policy_epoch: CURRENT_POLICY_EPOCH, ttlSeconds: 3600, }); } function mintActorToken(): string { return mintToken({ aud: "vellum-daemon", sub: "actor:self:user-123", scope_profile: "actor_client_v1", policy_epoch: CURRENT_POLICY_EPOCH, ttlSeconds: 3600, }); } function startFrame(conversationId = "conversation-123"): string { return JSON.stringify({ type: "start", conversationId, audio: { mimeType: "audio/pcm", sampleRate: 24_000, channels: 1, }, }); } async function waitForOpen(ws: WebSocket, timeoutMs = 2000): Promise { if (ws.readyState === WebSocket.OPEN) { return; } await new Promise((resolve, reject) => { const timer = setTimeout(() => { cleanup(); reject(new Error("Timed out waiting for WebSocket open")); }, timeoutMs); const cleanup = () => { clearTimeout(timer); ws.removeEventListener("open", onOpen); ws.removeEventListener("error", onError); }; const onOpen = () => { cleanup(); resolve(); }; const onError = () => { cleanup(); reject(new Error("WebSocket failed to open")); }; ws.addEventListener("open", onOpen); ws.addEventListener("error", onError); }); } async function waitForClose(ws: WebSocket, timeoutMs = 2000): Promise { if (ws.readyState === WebSocket.CLOSED) { return; } await new Promise((resolve, reject) => { const timer = setTimeout(() => { cleanup(); reject(new Error("Timed out waiting for WebSocket close")); }, timeoutMs); const cleanup = () => { clearTimeout(timer); ws.removeEventListener("close", onClose); ws.removeEventListener("error", onError); }; const onClose = () => { cleanup(); resolve(); }; const onError = () => { cleanup(); reject(new Error("WebSocket close failed")); }; ws.addEventListener("close", onClose); ws.addEventListener("error", onError); }); } async function waitForJsonFrame( ws: WebSocket, timeoutMs = 2000, ): Promise { await waitForOpen(ws, timeoutMs); return await new Promise((resolve, reject) => { const timer = setTimeout(() => { cleanup(); reject(new Error("Timed out waiting for WebSocket message")); }, timeoutMs); const cleanup = () => { clearTimeout(timer); ws.removeEventListener("message", onMessage); ws.removeEventListener("close", onClose); ws.removeEventListener("error", onError); }; const onMessage = (event: MessageEvent) => { cleanup(); const data = event.data; if (typeof data !== "string") { reject(new Error("Expected text WebSocket message")); return; } resolve(JSON.parse(data) as JsonFrame); }; const onClose = () => { cleanup(); reject(new Error("WebSocket closed before message")); }; const onError = () => { cleanup(); reject(new Error("WebSocket errored before message")); }; ws.addEventListener("message", onMessage); ws.addEventListener("close", onClose); ws.addEventListener("error", onError); }); } function closeClient(ws: WebSocket): void { if ( ws.readyState === WebSocket.CONNECTING || ws.readyState === WebSocket.OPEN ) { ws.close(1000, "test shutdown"); } } describe("RuntimeHttpServer live voice WebSocket shell", () => { let server: RuntimeHttpServer; let baseUrl: string; let wsBaseUrl: string; let clients: WebSocket[]; beforeAll(() => { shellMocksActive = true; }); afterAll(() => { shellMocksActive = false; }); beforeEach(async () => { delete process.env.DISABLE_HTTP_AUTH; resolveStreamingTranscriberImpl = async () => createResolvedTranscriber(); resolveStreamingTranscriberMock.mockClear(); resolvedTranscribers.length = 0; clients = []; const port = 21100 + Math.floor(Math.random() * 300); server = new RuntimeHttpServer({ port, hostname: "127.0.0.1" }); await server.start(); baseUrl = `http://127.0.0.1:${server.actualPort}`; wsBaseUrl = `ws://127.0.0.1:${server.actualPort}`; }); afterEach(async () => { for (const client of clients) { closeClient(client); } await server.stop(); if (savedAuthEnv.DISABLE_HTTP_AUTH === undefined) { delete process.env.DISABLE_HTTP_AUTH; } else { process.env.DISABLE_HTTP_AUTH = savedAuthEnv.DISABLE_HTTP_AUTH; } }); function openLiveVoiceClient(token = mintGatewayToken()): WebSocket { const ws = new WebSocket( `${wsBaseUrl}/v1/live-voice?token=${encodeURIComponent(token)}`, ); clients.push(ws); return ws; } test("rejects unauthorized upgrades before creating a WebSocket", async () => { const baseHeaders = { Upgrade: "websocket", Connection: "Upgrade", "Sec-WebSocket-Key": "dGhlIHNhbXBsZSBub25jZQ==", "Sec-WebSocket-Version": "13", }; const missingToken = await fetch(`${baseUrl}/v1/live-voice`, { headers: baseHeaders, }); expect(missingToken.status).toBe(401); const actorToken = await fetch( `${baseUrl}/v1/live-voice?token=${mintActorToken()}`, { headers: baseHeaders }, ); expect(actorToken.status).toBe(401); const externalOrigin = await fetch( `${baseUrl}/v1/live-voice?token=${mintGatewayToken()}`, { headers: { ...baseHeaders, Origin: "https://external.example.com", }, }, ); expect(externalOrigin.status).toBe(403); }); test("routes start and audio frames through the real live voice session", async () => { const ws = openLiveVoiceClient(); await waitForOpen(ws); ws.send(startFrame("conversation-ready")); const ready = await waitForJsonFrame(ws); expect(ready).toMatchObject({ type: "ready", seq: 1, conversationId: "conversation-ready", }); expect(typeof ready.sessionId).toBe("string"); // The config default is "multi", so the session carries a language into // the resolver rather than leaving it to be inferred downstream. expect(resolveStreamingTranscriberMock).toHaveBeenCalledWith({ sampleRate: 24_000, providerId: "deepgram", language: "multi", }); expect(resolvedTranscribers).toHaveLength(1); expect(resolvedTranscribers[0]?.started).toBe(true); ws.send(new Uint8Array([1, 2, 3])); const partial = await waitForJsonFrame(ws); expect(partial).toMatchObject({ type: "stt_partial", seq: 2, text: "partial-1", }); expect(resolvedTranscribers[0]?.audioChunks).toEqual([[1, 2, 3]]); expect(resolvedTranscribers[0]?.mimeTypes).toEqual(["audio/pcm"]); }); test("sends an error for malformed frames and can still start", async () => { const ws = openLiveVoiceClient(); await waitForOpen(ws); ws.send("{"); const error = await waitForJsonFrame(ws); expect(error).toMatchObject({ type: "error", seq: 1, code: "invalid_json", }); ws.send(startFrame("conversation-after-error")); const ready = await waitForJsonFrame(ws); expect(ready).toMatchObject({ type: "ready", conversationId: "conversation-after-error", }); expect(typeof ready.sessionId).toBe("string"); expect(ready.seq as number).toBeGreaterThan(error.seq as number); }); test("releases the session lock when the WebSocket closes", async () => { const first = openLiveVoiceClient(); const second = openLiveVoiceClient(); await Promise.all([waitForOpen(first), waitForOpen(second)]); first.send(startFrame("conversation-first")); const firstReady = await waitForJsonFrame(first); expect(firstReady.type).toBe("ready"); second.send(startFrame("conversation-second")); const busy = await waitForJsonFrame(second); expect(busy).toMatchObject({ type: "busy", activeSessionId: firstReady.sessionId, }); first.close(1000, "client finished"); await waitForClose(first); second.send(startFrame("conversation-second")); const secondReady = await waitForJsonFrame(second); expect(secondReady).toMatchObject({ type: "ready", conversationId: "conversation-second", }); expect(secondReady.sessionId).not.toBe(firstReady.sessionId); }); test("releases the session lock after startup STT failure without WebSocket close", async () => { let attempts = 0; resolveStreamingTranscriberImpl = async () => { attempts += 1; if (attempts === 1) { throw new Error("Deepgram credentials missing"); } return createResolvedTranscriber(); }; const ws = openLiveVoiceClient(); await waitForOpen(ws); ws.send(startFrame("conversation-failed")); // ready precedes the arm: the STT connection is established in the // background, so its failure surfaces as a post-ready error frame. const readyBeforeFailure = await waitForJsonFrame(ws); expect(readyBeforeFailure).toMatchObject({ type: "ready", conversationId: "conversation-failed", }); const error = await waitForJsonFrame(ws); expect(error).toMatchObject({ type: "error", code: "invalid_field", message: expect.stringContaining("Deepgram credentials missing"), }); ws.send(startFrame("conversation-retry")); const ready = await waitForJsonFrame(ws); expect(ready).toMatchObject({ type: "ready", conversationId: "conversation-retry", }); expect(typeof ready.sessionId).toBe("string"); expect(resolveStreamingTranscriberMock).toHaveBeenCalledTimes(2); expect(resolvedTranscribers).toHaveLength(1); expect(resolvedTranscribers[0]?.started).toBe(true); }); });