WeLe Agentic AI with Docker deployment
This commit is contained in:
@@ -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); });
|
||||
+37
-4
@@ -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();
|
||||
|
||||
Reference in New Issue
Block a user