From bef978da6fa587c203905207a58ba7ead6824604 Mon Sep 17 00:00:00 2001 From: thulasiraman S Date: Fri, 28 Aug 2026 10:12:31 +0530 Subject: [PATCH] WeLe Agentic AI with Docker deployment --- scripts/t-multiutterance.mjs | 85 ++++++++++++++++++++++++++++++++++++ src/gateway/voice.js | 41 +++++++++++++++-- 2 files changed, 122 insertions(+), 4 deletions(-) create mode 100644 scripts/t-multiutterance.mjs diff --git a/scripts/t-multiutterance.mjs b/scripts/t-multiutterance.mjs new file mode 100644 index 0000000..d86a112 --- /dev/null +++ b/scripts/t-multiutterance.mjs @@ -0,0 +1,85 @@ +/* Reproduces the reported bug: first question works, later ones hang on + "Listening". + + Streams THREE spoken utterances over the real voice socket, at the same + 40 ms cadence the browser uses, with silence between them. Before the fix + the endpointer's state was corrupted by concurrent frame processing during + the first answer, so utterances 2 and 3 were never detected. +*/ +import WebSocket from 'ws'; +import jwt from 'jsonwebtoken'; +import dotenv from 'dotenv'; +import { synthesize } from '../src/speech/index.js'; + +dotenv.config(); +const URL_BASE = process.env.TEST_VOICE_URL || 'ws://localhost:4000'; +const token = jwt.sign({ id: '6a00cefe524bebd27037a968' }, process.env.CRM_JWT_SECRET, { expiresIn: '30m' }); + +const QUESTIONS = [ + 'How many leads are in the new lead stage?', + 'How many leads did we get today?', + 'Who has overdue follow ups?', +]; + +const resample = (a, from, to) => { + const r = from / to, o = new Float32Array(Math.floor(a.length / r)); + for (let i = 0; i < o.length; i++) { const p = i * r, k = Math.floor(p); o[i] = a[k] + (a[Math.min(k + 1, a.length - 1)] - a[k]) * (p - k); } + return o; +}; + +console.log('synthesising the three questions as speech…'); +const clips = []; +for (const q of QUESTIONS) { + const s = await synthesize(q, 'en'); + clips.push(resample(s.audio, s.sampling_rate, 16000)); +} + +const ws = new WebSocket(`${URL_BASE}/api/agent/voice?token=${encodeURIComponent(token)}`); +const transcripts = []; +let idleCount = 0; + +const sleep = (ms) => new Promise((r) => setTimeout(r, ms)); + +/** Send one clip at the browser's real cadence, then a second of silence. */ +async function speak(clip) { + const FRAME = 640; // 40 ms @16k + for (let i = 0; i < clip.length; i += FRAME) { + const slice = clip.subarray(i, Math.min(i + FRAME, clip.length)); + const pcm = Buffer.alloc(slice.length * 2); + for (let j = 0; j < slice.length; j++) { + const v = Math.max(-1, Math.min(1, slice[j])); + pcm.writeInt16LE(v < 0 ? v * 0x8000 : v * 0x7fff, j * 2); + } + ws.send(pcm); + await sleep(40); + } + const silence = Buffer.alloc(FRAME * 2); + for (let i = 0; i < 30; i++) { ws.send(silence); await sleep(40); } // 1.2 s +} + +ws.on('message', (d, bin) => { + if (bin) return; + const m = JSON.parse(d); + if (m.type === 'transcript') { transcripts.push(m.text); console.log(` πŸ“ ${transcripts.length}: ${JSON.stringify(m.text)}`); } + else if (m.type === 'idle') { idleCount++; console.log(` βœ” turn ${idleCount} complete`); } + else if (m.type === 'heard_nothing') console.log(' ⚠️ heard nothing'); + else if (m.type === 'error') console.log(' βœ— error:', m.message); +}); + +ws.on('open', async () => { + console.log('connected β€” streaming 3 utterances at browser cadence\n'); + for (let i = 0; i < clips.length; i++) { + console.log(`speaking #${i + 1}: ${JSON.stringify(QUESTIONS[i])}`); + await speak(clips[i]); + // Wait for this turn to finish before the next, as a person would. + const target = i + 1; + for (let w = 0; w < 120 && idleCount < target; w++) await sleep(1000); + } + + console.log(`\nRESULT: ${transcripts.length}/3 utterances detected, ${idleCount}/3 turns completed`); + console.log(transcripts.length === 3 ? 'βœ… PASS β€” later questions are heard' : '❌ FAIL β€” stuck after the first'); + ws.close(); + process.exit(transcripts.length === 3 ? 0 : 1); +}); + +ws.on('error', (e) => { console.log('socket error:', e.message); process.exit(1); }); diff --git a/src/gateway/voice.js b/src/gateway/voice.js index cc6ab6a..32c75ef 100644 --- a/src/gateway/voice.js +++ b/src/gateway/voice.js @@ -55,6 +55,14 @@ class VoiceSession { this.abort = null; this.speakSeq = 0; // rising token; stale synthesis is discarded this.narrated = new Set(); + // Audio frames are processed strictly one at a time. The endpointer holds + // recurrent VAD state plus a partial-frame buffer, and neither survives + // concurrent access β€” see onAudio(). + this.audioChain = Promise.resolve(); + // Turns are serialized too. handleUtterance() is fire-and-forget so the + // audio queue keeps flowing, which means two utterances can overlap; this + // chain guarantees one answer at a time without dropping the second. + this.turnChain = Promise.resolve(); } send(obj) { @@ -66,7 +74,22 @@ class VoiceSession { } // ── microphone ─────────────────────────────────────────────────────────── - async onAudio(data) { + /** + * The browser streams a frame every 40 ms and the socket's 'message' handler + * does not await us, so without a queue ~25 calls a second would run + * concurrently against one Endpointer β€” interleaving its `pending` buffer and + * Silero's recurrent state until it stopped detecting speech at all. That is + * exactly what made the FIRST question work and every one after it hang on + * "Listening": the state was corrupted while the first answer was running. + */ + onAudio(data) { + this.audioChain = this.audioChain + .then(() => this.processAudio(data)) + .catch((e) => logger.error(`audio frame failed: ${e.message}`)); + return this.audioChain; + } + + async processAudio(data) { // Browser sends 16 kHz mono PCM16; the models want float32 in [-1, 1]. const pcm16 = new Int16Array(data.buffer, data.byteOffset, Math.floor(data.byteLength / 2)); const pcm = new Float32Array(pcm16.length); @@ -89,8 +112,11 @@ class VoiceSession { this.send({ type: 'barge_in' }); } + // Deliberately NOT awaited: a turn takes 17-46 s, and awaiting it here + // would stall the audio queue for that whole time β€” no barge-in, and a + // backlog of frames to grind through afterwards. for (const utterance of result.utterances) { - await this.handleUtterance(utterance); + this.handleUtterance(utterance).catch((e) => logger.error(`turn failed: ${e.message}`)); } } @@ -111,7 +137,15 @@ class VoiceSession { this.replyLang = heard.lang; this.send({ type: 'transcript', text: heard.text, lang: heard.lang, detected: heard.detected, ms: heard.ms }); - await this.answer(heard.text); + + // Cut the running turn short so the new question is answered promptly + // rather than queueing behind 40 s of superseded work. + if (this.busy) this.abort?.abort(); + + this.turnChain = this.turnChain + .then(() => this.answer(heard.text)) + .catch((e) => logger.error(`turn failed: ${e.message}`)); + await this.turnChain; } // ── speaking ───────────────────────────────────────────────────────────── @@ -133,7 +167,6 @@ class VoiceSession { } async answer(question) { - if (this.busy) return; // one turn at a time this.busy = true; this.narrated.clear(); this.abort = new AbortController();