// Integration test for the mock speech-engine WebSocket protocol. // Spawns `speech-engine/server.py --mock` from the local .venv on a free port // and drives it with Node's native WebSocket client. Skipped when the venv is missing. import { spawn } from 'node:child_process'; import type { ChildProcessWithoutNullStreams } from 'node:child_process'; import fs from 'node:fs'; import net from 'node:net'; import path from 'node:path'; import { fileURLToPath } from 'node:url'; import { afterAll, beforeAll, describe, expect, it } from 'vitest'; const repoRoot = path.resolve(path.dirname(fileURLToPath(import.meta.url)), '..'); interface SpeechEvent { type: string; [key: string]: unknown; } interface SegmentPayload { id: string; speakerId: string; start: number; end: number; text: string; final: boolean; } function venvPython(): string | null { const candidates = process.platform === 'win32' ? [path.join(repoRoot, 'speech-engine', '.venv', 'Scripts', 'python.exe')] : [path.join(repoRoot, 'speech-engine', '.venv', 'bin', 'python')]; return candidates.find((candidate) => fs.existsSync(candidate)) ?? null; } function findFreePort(): Promise { return new Promise((resolve, reject) => { const server = net.createServer(); server.unref(); server.on('error', reject); server.listen(0, '127.0.0.1', () => { const address = server.address(); if (!address || typeof address === 'string') { reject(new Error('Could not determine a free port')); return; } server.close(() => resolve(address.port)); }); }); } async function waitForHealth(url: string, timeoutMs: number): Promise { const deadline = Date.now() + timeoutMs; for (;;) { try { const response = await fetch(url); if (response.ok) return; } catch { // Server not up yet. } if (Date.now() > deadline) throw new Error(`Speech engine did not become healthy at ${url}`); await new Promise((resolve) => setTimeout(resolve, 250)); } } function openSocket(url: string): Promise { return new Promise((resolve, reject) => { const socket = new WebSocket(url); socket.onopen = () => resolve(socket); socket.onerror = (event) => reject(new Error(`WebSocket error: ${String(event.message ?? event.type)}`)); }); } function nextEvent(socket: WebSocket, timeoutMs: number): Promise { return new Promise((resolve, reject) => { const timer = setTimeout(() => reject(new Error('Timed out waiting for a speech engine event')), timeoutMs); socket.onmessage = (event) => { clearTimeout(timer); try { resolve(JSON.parse(String(event.data)) as SpeechEvent); } catch (error) { reject(error instanceof Error ? error : new Error('Non-JSON message from speech engine')); } }; }); } function waitForClose(socket: WebSocket, timeoutMs: number): Promise { return new Promise((resolve, reject) => { const timer = setTimeout(() => reject(new Error('Timed out waiting for the socket to close')), timeoutMs); socket.onclose = (event) => { clearTimeout(timer); resolve(event.code); }; }); } const python = venvPython(); describe.skipIf(python === null)( 'mock speech protocol (server.py --mock)', () => { let server: ChildProcessWithoutNullStreams | undefined; let port = 0; const stderr: string[] = []; beforeAll(async () => { if (!python) return; port = await findFreePort(); server = spawn(python, [path.join(repoRoot, 'speech-engine', 'server.py'), '--mock', '--port', String(port)], { cwd: repoRoot, stdio: ['ignore', 'pipe', 'pipe'], }); server.stdout.on('data', () => undefined); server.stderr.on('data', (chunk) => stderr.push(String(chunk))); try { await waitForHealth(`http://127.0.0.1:${port}/health`, 30_000); } catch (error) { throw new Error( `Speech engine did not start: ${error instanceof Error ? error.message : String(error)}\n--- stderr ---\n${stderr.join('')}`, ); } }, 60_000); afterAll(() => { server?.kill(); }); const validConfig = JSON.stringify({ model: 'small', compute: 'auto', language: 'auto', diarization: true, sample_rate: 16000, }); it('reports mock mode on /health', async () => { const response = await fetch(`http://127.0.0.1:${port}/health`); expect(response.ok).toBe(true); const body = (await response.json()) as { ok: boolean; mode: string; sample_rate: number }; expect(body.ok).toBe(true); expect(body.mode).toBe('mock'); expect(body.sample_rate).toBe(16000); }); it('completes the ready handshake for a valid config', async () => { const socket = await openSocket(`ws://127.0.0.1:${port}/ws/live`); socket.send(validConfig); const ready = await nextEvent(socket, 10_000); expect(ready.type).toBe('ready'); expect(String(ready.engine)).toContain('mock'); socket.close(); }); it('emits final segments with cycling speakers for streamed PCM', async () => { const socket = await openSocket(`ws://127.0.0.1:${port}/ws/live`); socket.send(validConfig); expect((await nextEvent(socket, 10_000)).type).toBe('ready'); // Stream 5 seconds of silence (16 kHz mono int16) in 100 ms chunks; the mock // emits one final segment per 2.4 s of samples with cycling speaker labels. const totalBytes = 16000 * 5 * 2; const chunkSize = 3200; for (let offset = 0; offset < totalBytes; offset += chunkSize) { socket.send(new ArrayBuffer(chunkSize)); } const finals: SpeechEvent[] = []; while (finals.length < 2) { const event = await nextEvent(socket, 15_000); if (event.type === 'error') throw new Error(String(event.message)); if (event.type === 'final') finals.push(event); } expect(finals).toHaveLength(2); const first = finals[0].segment as SegmentPayload; const second = finals[1].segment as SegmentPayload; for (const segment of [first, second]) { expect(segment.id).toBeTruthy(); expect(segment.text.length).toBeGreaterThan(0); expect(segment.final).toBe(true); } expect(first.speakerId).toBe('speaker_1'); expect(second.speakerId).toBe('speaker_2'); expect(first.start).toBeCloseTo(0.3, 1); expect(first.end).toBeCloseTo(2.4, 1); expect(second.start).toBeCloseTo(2.7, 1); expect(second.end).toBeCloseTo(4.8, 1); expect(second.start).toBeGreaterThan(first.end); socket.send(JSON.stringify({ type: 'stop' })); const code = await waitForClose(socket, 5_000); expect(code).toBe(1000); }); it('rejects unsupported sample rates', async () => { const socket = await openSocket(`ws://127.0.0.1:${port}/ws/live`); socket.send( JSON.stringify({ model: 'small', compute: 'auto', language: 'auto', diarization: true, sample_rate: 44100 }), ); const event = await nextEvent(socket, 10_000); expect(event.type).toBe('error'); expect(String(event.message)).toContain('16 kHz'); await waitForClose(socket, 5_000); }); it('reports malformed config as an error event', async () => { const socket = await openSocket(`ws://127.0.0.1:${port}/ws/live`); socket.send('this is not json'); const event = await nextEvent(socket, 10_000); expect(event.type).toBe('error'); expect(String(event.message)).not.toBe(''); await waitForClose(socket, 5_000); }); }, );