WeLe Agentic AI with Docker deployment
This commit is contained in:
+131
-178
@@ -1,104 +1,59 @@
|
||||
// ============================================
|
||||
// Voice channel — a WebSocket bridge between the browser and the GPU service.
|
||||
// Voice channel — speech in, speech out, in this same Node process.
|
||||
//
|
||||
// Voice is a *channel*, not a parallel product: a spoken question runs through
|
||||
// the same graph, guardrails and agents as a typed one. Only the transport and
|
||||
// the presentation differ, which is why this file contains no CRM logic.
|
||||
//
|
||||
// browser ──audio──► this ──audio──► python(:4100) ──STT──► transcript
|
||||
// │ │
|
||||
// └──────────── runTurn(graph) ◄───────────┘
|
||||
// │
|
||||
// browser ◄──audio── this ◄──audio── python(TTS) ◄──sentences───┘
|
||||
// browser ──PCM16──► Endpointer ──► transcribe() ──► runTurn(graph)
|
||||
// │ │
|
||||
// browser ◄──float32──── synthesize() ◄── sentences ◄─────┘
|
||||
//
|
||||
// The hard problem is not transport, it is that a turn takes 17–46 s. Silence
|
||||
// for that long feels broken, so the bridge speaks immediately, narrates what
|
||||
// the agents are doing, and starts reading the answer at the first sentence
|
||||
// rather than waiting for the last.
|
||||
// The hard problem is not audio, it is that a turn takes 17–46 s and that is
|
||||
// dead silence in voice. So this speaks an acknowledgement within ~1 s,
|
||||
// narrates each agent delegation aloud, then reads the answer sentence by
|
||||
// sentence as it is composed.
|
||||
// ============================================
|
||||
import { WebSocketServer, WebSocket } from 'ws';
|
||||
import { randomUUID } from 'node:crypto';
|
||||
import { principalFromToken } from './auth.js';
|
||||
import { runTurn } from '../orchestration/runner.js';
|
||||
import {
|
||||
Endpointer, transcribe, synthesize, sentences, speakable, LANGUAGES,
|
||||
} from '../speech/index.js';
|
||||
import config from '../config/index.js';
|
||||
import logger from '../utils/logger.js';
|
||||
|
||||
const VOICE_URL = process.env.VOICE_SERVICE_URL || 'ws://127.0.0.1:4100/ws/voice';
|
||||
|
||||
/** Spoken filler, per language. Said the instant a question lands. */
|
||||
/** Spoken filler, said the instant a question lands. */
|
||||
const ACK = {
|
||||
ta: ['பார்க்கிறேன்...', 'ஒரு நிமிடம், பார்க்கிறேன்.'],
|
||||
hi: ['देखता हूँ...', 'एक मिनट, देख रहा हूँ।'],
|
||||
te: ['చూస్తున్నాను...'],
|
||||
kn: ['ನೋಡುತ್ತಿದ್ದೇನೆ...'],
|
||||
ml: ['നോക്കുന്നു...'],
|
||||
mr: ['बघतो...'],
|
||||
bn: ['দেখছি...'],
|
||||
ta: ['பார்க்கிறேன்.', 'ஒரு நிமிடம், பார்க்கிறேன்.'],
|
||||
en: ['Let me check.', 'One moment, checking now.'],
|
||||
};
|
||||
|
||||
/** Progress narration, kept short — it is spoken over the user's waiting time. */
|
||||
/** Progress narration — spoken over the user's waiting time, so keep it short. */
|
||||
const NARRATE = {
|
||||
ta: { lead: 'லீட் விவரங்களைப் பார்க்கிறேன்.', analytics: 'புள்ளிவிவரங்களைச் சரிபார்க்கிறேன்.', conversation: 'உரையாடல்களைப் பார்க்கிறேன்.', default: 'தரவைச் சரிபார்க்கிறேன்.' },
|
||||
hi: { lead: 'लीड्स देख रहा हूँ।', analytics: 'आँकड़े देख रहा हूँ।', conversation: 'बातचीत देख रहा हूँ।', default: 'डेटा देख रहा हूँ।' },
|
||||
en: { lead: 'Checking the leads.', analytics: 'Pulling the numbers.', conversation: 'Looking at the conversations.', default: 'Checking the data.' },
|
||||
};
|
||||
|
||||
const pick = (arr) => arr[Math.floor(Math.random() * arr.length)];
|
||||
const NOTHING = { ta: 'பதில் கிடைக்கவில்லை.', en: 'I could not find an answer for that.' };
|
||||
const OOPS = { ta: 'மன்னிக்கவும், ஒரு பிழை ஏற்பட்டது.', en: 'Sorry, something went wrong.' };
|
||||
|
||||
function ackFor(lang) {
|
||||
return pick(ACK[lang] || ACK.en);
|
||||
}
|
||||
const pick = (a) => a[Math.floor(Math.random() * a.length)];
|
||||
const ackFor = (l) => pick(ACK[l] || ACK.en);
|
||||
const narrateFor = (l, agent) => (NARRATE[l] || NARRATE.en)[agent] || (NARRATE[l] || NARRATE.en).default;
|
||||
|
||||
function narrationFor(lang, agent) {
|
||||
const set = NARRATE[lang] || NARRATE.en;
|
||||
return set[agent] || set.default;
|
||||
}
|
||||
|
||||
/**
|
||||
* Strip block-oriented markdown before speaking. Tables and code read terribly
|
||||
* aloud, and the visual blocks are already on screen.
|
||||
*/
|
||||
export function speakable(markdown = '') {
|
||||
return markdown
|
||||
.replace(/```[\s\S]*?```/g, ' ')
|
||||
.replace(/^\s*\|.*\|\s*$/gm, ' ') // table rows
|
||||
.replace(/^\s*[-*]\s+/gm, '') // bullets
|
||||
.replace(/^#{1,6}\s*/gm, '') // headings
|
||||
.replace(/\*\*([^*]+)\*\*/g, '$1')
|
||||
.replace(/`([^`]+)`/g, '$1')
|
||||
.replace(/\[([^\]]+)\]\([^)]+\)/g, '$1')
|
||||
.replace(/₹\s?([\d,.]+)/g, 'rupees $1')
|
||||
.replace(/\s{2,}/g, ' ')
|
||||
.trim();
|
||||
}
|
||||
|
||||
/** Split into sentences so speech can start before the answer is finished. */
|
||||
export function sentences(text, max = 240) {
|
||||
const out = [];
|
||||
for (const raw of text.split(/(?<=[.!?।])\s+/)) {
|
||||
let s = raw.trim();
|
||||
if (!s) continue;
|
||||
while (s.length > max) {
|
||||
const cut = s.lastIndexOf(' ', max);
|
||||
out.push(s.slice(0, cut > 0 ? cut : max).trim());
|
||||
s = s.slice(cut > 0 ? cut : max).trim();
|
||||
}
|
||||
if (s) out.push(s);
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
class VoiceBridge {
|
||||
class VoiceSession {
|
||||
constructor(client, user) {
|
||||
this.client = client;
|
||||
this.user = user;
|
||||
this.lang = 'auto'; // what the user chose
|
||||
this.replyLang = 'ta'; // what the last utterance actually was
|
||||
this.lang = 'auto'; // what the user selected
|
||||
this.replyLang = 'ta'; // what the last utterance actually was
|
||||
this.prefer = 'ta'; // tiebreak when detection is unusable
|
||||
this.sessionId = `voice:${Date.now().toString(36)}:${Math.random().toString(36).slice(2, 8)}`;
|
||||
this.gpu = null;
|
||||
this.endpointer = new Endpointer();
|
||||
this.busy = false;
|
||||
this.abort = null;
|
||||
this.speakSeq = 0; // rising token; stale synthesis is discarded
|
||||
this.narrated = new Set();
|
||||
}
|
||||
|
||||
@@ -106,109 +61,104 @@ class VoiceBridge {
|
||||
if (this.client.readyState === WebSocket.OPEN) this.client.send(JSON.stringify(obj));
|
||||
}
|
||||
|
||||
toGpu(obj) {
|
||||
if (this.gpu?.readyState === WebSocket.OPEN) this.gpu.send(JSON.stringify(obj));
|
||||
sendAudio(buf) {
|
||||
if (this.client.readyState === WebSocket.OPEN) this.client.send(buf, { binary: true });
|
||||
}
|
||||
|
||||
async connect() {
|
||||
this.gpu = new WebSocket(VOICE_URL);
|
||||
// ── microphone ───────────────────────────────────────────────────────────
|
||||
async onAudio(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);
|
||||
for (let i = 0; i < pcm16.length; i++) pcm[i] = pcm16[i] / 32768;
|
||||
|
||||
this.gpu.on('open', () => {
|
||||
logger.info(`🎙️ voice session ${this.sessionId} → GPU service`);
|
||||
this.toGpu({ type: 'config', lang: this.lang });
|
||||
});
|
||||
let result;
|
||||
try {
|
||||
result = await this.endpointer.push(pcm);
|
||||
} catch (e) {
|
||||
logger.error(`VAD failed: ${e.message}`);
|
||||
this.send({ type: 'error', message: 'Voice input failed to initialise. Check the server logs.' });
|
||||
return;
|
||||
}
|
||||
|
||||
this.gpu.on('message', (data, isBinary) => {
|
||||
// TTS audio: pass straight through, no re-encoding.
|
||||
if (isBinary) {
|
||||
if (this.client.readyState === WebSocket.OPEN) this.client.send(data, { binary: true });
|
||||
return;
|
||||
}
|
||||
let msg;
|
||||
try { msg = JSON.parse(data.toString()); } catch { return; }
|
||||
this.onGpuMessage(msg);
|
||||
});
|
||||
if (result.started) {
|
||||
// Barge-in: the user talking wins immediately. Bumping the token drops
|
||||
// any in-flight synthesis rather than letting it arrive late.
|
||||
this.speakSeq++;
|
||||
this.abort?.abort();
|
||||
this.send({ type: 'barge_in' });
|
||||
}
|
||||
|
||||
this.gpu.on('error', (err) => {
|
||||
logger.error(`voice GPU service: ${err.message}`);
|
||||
this.send({ type: 'error', message: 'The voice service is not reachable. Start it with: npm run voice' });
|
||||
});
|
||||
|
||||
this.gpu.on('close', () => {
|
||||
this.send({ type: 'voice_service_closed' });
|
||||
this.client.close();
|
||||
});
|
||||
}
|
||||
|
||||
onGpuMessage(msg) {
|
||||
switch (msg.type) {
|
||||
case 'ready':
|
||||
this.send({ type: 'ready', session_id: this.sessionId, languages: msg.languages, sample_rate_out: msg.sample_rate_out });
|
||||
break;
|
||||
|
||||
case 'speech_start':
|
||||
// The user started talking — the GPU service already stopped speaking.
|
||||
// Tell the browser to dump whatever is still in its playback buffer,
|
||||
// and abandon any answer still being composed.
|
||||
this.send({ type: 'barge_in' });
|
||||
this.abort?.abort();
|
||||
break;
|
||||
|
||||
case 'transcript':
|
||||
// Answer in the language the person actually spoke, not the menu
|
||||
// setting — that is the whole point of auto mode.
|
||||
if (msg.lang) this.replyLang = msg.lang;
|
||||
this.send({ type: 'transcript', text: msg.text, lang: msg.lang, detected: msg.detected, confidence: msg.confidence, ms: msg.ms });
|
||||
this.handleQuestion(msg.text);
|
||||
break;
|
||||
|
||||
case 'transcript_empty':
|
||||
this.send({ type: 'heard_nothing' });
|
||||
break;
|
||||
|
||||
case 'audio_start':
|
||||
case 'audio_end':
|
||||
case 'error':
|
||||
this.send(msg);
|
||||
break;
|
||||
|
||||
default:
|
||||
break;
|
||||
for (const utterance of result.utterances) {
|
||||
await this.handleUtterance(utterance);
|
||||
}
|
||||
}
|
||||
|
||||
speak(text, id = randomUUID()) {
|
||||
const clean = speakable(text);
|
||||
if (clean) this.toGpu({ type: 'speak', text: clean, id, lang: this.replyLang });
|
||||
async handleUtterance(audio) {
|
||||
let heard;
|
||||
try {
|
||||
heard = await transcribe(audio, this.lang, this.prefer);
|
||||
} catch (e) {
|
||||
logger.error(`STT failed: ${e.message}`);
|
||||
this.send({ type: 'error', message: 'Could not transcribe that. Try again.' });
|
||||
return;
|
||||
}
|
||||
|
||||
if (!heard.text) {
|
||||
this.send({ type: 'heard_nothing' });
|
||||
return;
|
||||
}
|
||||
|
||||
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);
|
||||
}
|
||||
|
||||
async handleQuestion(text) {
|
||||
if (!text?.trim()) return;
|
||||
if (this.busy) return; // one turn at a time
|
||||
// ── speaking ─────────────────────────────────────────────────────────────
|
||||
/** Synthesise and stream one piece, unless a newer turn has superseded it. */
|
||||
async say(text, seq) {
|
||||
const clean = speakable(text);
|
||||
if (!clean || seq !== this.speakSeq) return;
|
||||
try {
|
||||
const out = await synthesize(clean, this.replyLang);
|
||||
if (!out || seq !== this.speakSeq) return; // interrupted while generating
|
||||
|
||||
this.send({ type: 'audio_start', sample_rate: out.sampling_rate });
|
||||
// Float32 straight down the socket — the playback worklet takes it as-is.
|
||||
this.sendAudio(Buffer.from(out.audio.buffer, out.audio.byteOffset, out.audio.byteLength));
|
||||
this.send({ type: 'audio_end' });
|
||||
} catch (e) {
|
||||
logger.error(`TTS failed: ${e.message}`);
|
||||
}
|
||||
}
|
||||
|
||||
async answer(question) {
|
||||
if (this.busy) return; // one turn at a time
|
||||
this.busy = true;
|
||||
this.narrated.clear();
|
||||
this.abort = new AbortController();
|
||||
const seq = ++this.speakSeq;
|
||||
|
||||
// 1. Answer the silence immediately. This is the whole trick: the pipeline
|
||||
// still takes 17–46 s, but the user hears a response in ~1 s.
|
||||
this.speak(ackFor(this.replyLang));
|
||||
// Answer the silence immediately. The pipeline still takes 17–46 s, but
|
||||
// the user hears a response in about a second.
|
||||
this.say(ackFor(this.replyLang), seq);
|
||||
this.send({ type: 'thinking' });
|
||||
|
||||
try {
|
||||
const result = await runTurn({
|
||||
sessionId: this.sessionId,
|
||||
message: text,
|
||||
message: question,
|
||||
user: this.user,
|
||||
channel: 'crm_chat', // voice users are staff; full tool access
|
||||
channel: 'crm_chat', // voice users are staff
|
||||
signal: this.abort.signal,
|
||||
onEvent: (ev) => {
|
||||
this.send(ev);
|
||||
// 2. Narrate delegations — but only once per agent, or it chatters.
|
||||
// Narrate delegations, once per agent, or it chatters.
|
||||
if (ev.type === 'step' && ev.kind === 'delegate') {
|
||||
const agent = String(ev.label || '').toLowerCase().split(' ')[0];
|
||||
if (!this.narrated.has(agent)) {
|
||||
this.narrated.add(agent);
|
||||
this.speak(narrationFor(this.replyLang, agent));
|
||||
this.say(narrateFor(this.replyLang, agent), seq);
|
||||
}
|
||||
}
|
||||
},
|
||||
@@ -216,23 +166,24 @@ class VoiceBridge {
|
||||
|
||||
this.send({ type: 'result', blocks: result.blocks, usage: result.usage });
|
||||
|
||||
// 3. Read the answer. Sentence at a time so speech starts sooner and can
|
||||
// be cut cleanly if the user interrupts.
|
||||
const answer = result.blocks?.filter((b) => b.type === 'text').map((b) => b.markdown).join(' ')
|
||||
|| result.answer || '';
|
||||
const answer = (result.blocks || [])
|
||||
.filter((b) => b.type === 'text').map((b) => b.markdown).join(' ') || result.answer || '';
|
||||
const parts = sentences(speakable(answer));
|
||||
|
||||
if (!parts.length) {
|
||||
this.speak(this.replyLang === 'ta' ? 'பதில் கிடைக்கவில்லை.' : 'I could not find an answer for that.');
|
||||
await this.say(NOTHING[this.replyLang] || NOTHING.en, seq);
|
||||
} else {
|
||||
// Sequential on purpose: parallel synthesis would race to the socket
|
||||
// and play the answer out of order.
|
||||
for (const part of parts) {
|
||||
if (this.abort.signal.aborted) break;
|
||||
this.speak(part);
|
||||
if (seq !== this.speakSeq || this.abort.signal.aborted) break;
|
||||
await this.say(part, seq);
|
||||
}
|
||||
}
|
||||
} catch (err) {
|
||||
if (err?.name !== 'AbortError') {
|
||||
logger.error(`voice turn failed: ${err.message}`);
|
||||
this.speak(this.replyLang === 'ta' ? 'மன்னிக்கவும், ஒரு பிழை ஏற்பட்டது.' : 'Sorry, something went wrong.');
|
||||
await this.say(OOPS[this.replyLang] || OOPS.en, seq);
|
||||
}
|
||||
} finally {
|
||||
this.busy = false;
|
||||
@@ -240,32 +191,33 @@ class VoiceBridge {
|
||||
}
|
||||
}
|
||||
|
||||
onClientMessage(data, isBinary) {
|
||||
if (isBinary) {
|
||||
if (this.gpu?.readyState === WebSocket.OPEN) this.gpu.send(data, { binary: true });
|
||||
return;
|
||||
}
|
||||
// ── control ──────────────────────────────────────────────────────────────
|
||||
onMessage(data, isBinary) {
|
||||
if (isBinary) return this.onAudio(data);
|
||||
|
||||
let msg;
|
||||
try { msg = JSON.parse(data.toString()); } catch { return; }
|
||||
try { msg = JSON.parse(data.toString()); } catch { return undefined; }
|
||||
|
||||
if (msg.type === 'config' && msg.lang) {
|
||||
this.lang = msg.lang;
|
||||
if (msg.lang !== 'auto') this.replyLang = msg.lang;
|
||||
this.toGpu({ type: 'config', lang: msg.lang });
|
||||
this.send({ type: 'config_ok', lang: msg.lang });
|
||||
if (msg.lang !== 'auto') this.replyLang = this.prefer = msg.lang;
|
||||
else if (msg.prefer) this.prefer = msg.prefer;
|
||||
this.endpointer.reset();
|
||||
this.send({ type: 'config_ok', lang: this.lang, prefer: this.prefer });
|
||||
} else if (msg.type === 'cancel') {
|
||||
this.speakSeq++;
|
||||
this.abort?.abort();
|
||||
this.toGpu({ type: 'cancel' });
|
||||
} else if (msg.type === 'text') {
|
||||
// Typed question while in voice mode — answered aloud like a spoken one.
|
||||
this.send({ type: 'cancelled' });
|
||||
} else if (msg.type === 'text' && msg.text) {
|
||||
this.send({ type: 'transcript', text: msg.text, lang: this.replyLang, typed: true });
|
||||
this.handleQuestion(msg.text);
|
||||
this.answer(msg.text);
|
||||
}
|
||||
return undefined;
|
||||
}
|
||||
|
||||
close() {
|
||||
this.speakSeq++;
|
||||
this.abort?.abort();
|
||||
try { this.gpu?.close(); } catch { /* already gone */ }
|
||||
}
|
||||
}
|
||||
|
||||
@@ -275,12 +227,11 @@ export function attachVoice(server) {
|
||||
|
||||
server.on('upgrade', async (req, socket, head) => {
|
||||
const url = new URL(req.url, `http://${req.headers.host}`);
|
||||
if (url.pathname !== '/api/agent/voice') return; // leave other upgrades alone
|
||||
if (url.pathname !== '/api/agent/voice') return; // leave other upgrades alone
|
||||
|
||||
// Browsers cannot set headers on a WebSocket, so the CRM token arrives as
|
||||
// a query parameter. It is the same token and the same verification.
|
||||
const token = url.searchParams.get('token');
|
||||
const user = await principalFromToken(token).catch(() => null);
|
||||
// a query parameter. Same token, same verification as every other route.
|
||||
const user = await principalFromToken(url.searchParams.get('token')).catch(() => null);
|
||||
if (!user) {
|
||||
socket.write('HTTP/1.1 401 Unauthorized\r\n\r\n');
|
||||
socket.destroy();
|
||||
@@ -288,14 +239,16 @@ export function attachVoice(server) {
|
||||
}
|
||||
|
||||
wss.handleUpgrade(req, socket, head, (client) => {
|
||||
const bridge = new VoiceBridge(client, user);
|
||||
bridge.connect();
|
||||
client.on('message', (d, bin) => bridge.onClientMessage(d, bin));
|
||||
client.on('close', () => bridge.close());
|
||||
client.on('error', () => bridge.close());
|
||||
const session = new VoiceSession(client, user);
|
||||
logger.info(`🎙️ voice session ${session.sessionId} (${user.name})`);
|
||||
session.send({ type: 'ready', session_id: session.sessionId, languages: LANGUAGES });
|
||||
|
||||
client.on('message', (d, bin) => session.onMessage(d, bin));
|
||||
client.on('close', () => session.close());
|
||||
client.on('error', () => session.close());
|
||||
});
|
||||
});
|
||||
|
||||
logger.info(` voice =ws://localhost:${config.port}/api/agent/voice → ${VOICE_URL}`);
|
||||
logger.info(` voice =ws://localhost:${config.port}/api/agent/voice (in-process, CPU)`);
|
||||
return wss;
|
||||
}
|
||||
|
||||
@@ -17,6 +17,7 @@ import { ensureDir, sweep } from './output/artifactStore.js';
|
||||
import crmApi from './tools/http/crmApi.js';
|
||||
import { describeChains } from './orchestration/llm.js';
|
||||
import { attachVoice } from './gateway/voice.js';
|
||||
import { warmup as warmSpeech, speechStatus } from './speech/index.js';
|
||||
|
||||
const app = express();
|
||||
|
||||
@@ -40,6 +41,7 @@ app.get('/health', async (_req, res) => {
|
||||
redis: redisOk ? redisMode() : 'unavailable',
|
||||
crm_api: crm.reachable ? 'reachable' : `unreachable (${crm.error || crm.status})`,
|
||||
models: describeChains(),
|
||||
speech: speechStatus(),
|
||||
uptime_s: Math.round(process.uptime()),
|
||||
});
|
||||
});
|
||||
@@ -83,6 +85,11 @@ async function start() {
|
||||
// second origin and the CRM token works unchanged.
|
||||
attachVoice(server);
|
||||
|
||||
// Speech models load lazily on the first voice turn (~10 s). Set
|
||||
// SPEECH_WARMUP=true to pay that at boot instead — worth it in production,
|
||||
// wasteful in development where most restarts never use voice.
|
||||
if (process.env.SPEECH_WARMUP === 'true') warmSpeech(['ta', 'en']);
|
||||
|
||||
const shutdown = (sig) => {
|
||||
logger.info(`${sig} — shutting down`);
|
||||
server.close(() => process.exit(0));
|
||||
|
||||
@@ -0,0 +1,169 @@
|
||||
// ============================================
|
||||
// Speech pipeline — transcribe() and synthesize().
|
||||
//
|
||||
// Whisper is multilingual and can identify the spoken language, so "auto"
|
||||
// costs nothing extra: detection and transcription are the same forward pass.
|
||||
// That matters for a WeLe agent who switches between Tamil and English inside
|
||||
// one shift and should never have to touch a language menu.
|
||||
// ============================================
|
||||
import { Tensor } from '@huggingface/transformers';
|
||||
import { getSTT, getTTS, supportsTTS } from './models.js';
|
||||
import logger from '../utils/logger.js';
|
||||
|
||||
export { LANGUAGES, warmup, speechStatus } from './models.js';
|
||||
export { Endpointer, warmupVad } from './vad.js';
|
||||
|
||||
const RATE = 16000;
|
||||
|
||||
/** Languages we can both hear and speak. */
|
||||
const SPOKEN = new Set(['ta', 'en']);
|
||||
|
||||
// Below this, trust the caller's preference over the detector. Short or noisy
|
||||
// utterances — and code-mixed "Tanglish" especially — can land either side.
|
||||
const DETECT_CONFIDENCE = Number(process.env.DETECT_CONFIDENCE ?? 0.6);
|
||||
|
||||
let detectIds = null;
|
||||
|
||||
/**
|
||||
* Identify the spoken language in ONE decoder step.
|
||||
*
|
||||
* Passing no `language` to the pipeline does NOT auto-detect — Transformers.js
|
||||
* logs "No language specified - defaulting to English" and transcribes Tamil
|
||||
* as English, producing nonsense. Whisper does emit a language token right
|
||||
* after <|startoftranscript|>, so we read that distribution directly. Measured
|
||||
* ~700 ms, and 0.998 / 1.000 confidence on clean Tamil / English.
|
||||
*
|
||||
* Reuses the pipeline's own model and processor, so nothing loads twice.
|
||||
*/
|
||||
async function detectLanguage(audio) {
|
||||
const stt = await getSTT();
|
||||
const tok = stt.tokenizer;
|
||||
|
||||
if (!detectIds) {
|
||||
const id = (t) => tok.encode(t, { add_special_tokens: false })[0];
|
||||
detectIds = { sot: id('<|startoftranscript|>'), langs: [...SPOKEN].map((c) => ({ code: c, id: id(`<|${c}|>`) })) };
|
||||
}
|
||||
|
||||
const inputs = await stt.processor(audio);
|
||||
const out = await stt.model({
|
||||
...inputs,
|
||||
decoder_input_ids: new Tensor('int64', BigInt64Array.from([BigInt(detectIds.sot)]), [1, 1]),
|
||||
});
|
||||
|
||||
const { dims, data } = out.logits;
|
||||
const row = data.slice((dims[1] - 1) * dims[2], dims[1] * dims[2]);
|
||||
const scores = detectIds.langs.map((l) => Number(row[l.id]));
|
||||
const max = Math.max(...scores);
|
||||
const exp = scores.map((v) => Math.exp(v - max));
|
||||
const sum = exp.reduce((a, b) => a + b, 0);
|
||||
const probs = exp.map((v) => v / sum);
|
||||
const best = probs.indexOf(Math.max(...probs));
|
||||
|
||||
return { lang: detectIds.langs[best].code, confidence: probs[best] };
|
||||
}
|
||||
|
||||
/**
|
||||
* @param {Float32Array} audio mono @16 kHz in [-1, 1]
|
||||
* @param {string} lang 'auto' | 'ta' | 'en'
|
||||
* @param {string} prefer used when detection is unusable
|
||||
*/
|
||||
export async function transcribe(audio, lang = 'auto', prefer = 'ta') {
|
||||
if (!audio || audio.length < RATE / 5) { // under 200 ms
|
||||
return { text: '', lang: prefer, note: 'too short' };
|
||||
}
|
||||
|
||||
const stt = await getSTT();
|
||||
const t0 = Date.now();
|
||||
|
||||
// Whisper must always be told a language — it never detects on its own here.
|
||||
let used = lang;
|
||||
let detected = null;
|
||||
let confidence = null;
|
||||
|
||||
if (lang === 'auto') {
|
||||
try {
|
||||
const d = await detectLanguage(audio);
|
||||
detected = d.lang;
|
||||
confidence = d.confidence;
|
||||
used = d.confidence >= DETECT_CONFIDENCE ? d.lang : prefer;
|
||||
if (used !== d.lang) {
|
||||
logger.info(`language ID unsure (${d.lang} @ ${d.confidence.toFixed(2)}) — using preferred ${prefer}`);
|
||||
}
|
||||
} catch (e) {
|
||||
logger.warn(`language ID failed (${e.message}) — using preferred ${prefer}`);
|
||||
used = prefer;
|
||||
}
|
||||
}
|
||||
|
||||
const result = await stt(audio, { task: 'transcribe', language: used, return_timestamps: false });
|
||||
return finish(result, used, audio, t0, detected, confidence);
|
||||
}
|
||||
|
||||
function finish(result, used, audio, t0, detected, confidence) {
|
||||
const text = (result?.text || '').trim();
|
||||
const ms = Date.now() - t0;
|
||||
const audioMs = Math.round((audio.length / RATE) * 1000);
|
||||
logger.info(`🎤 STT ${used}${detected && detected !== used ? ` (heard ${detected})` : ''}: ${audioMs}ms → ${ms}ms → ${JSON.stringify(text.slice(0, 70))}`);
|
||||
return {
|
||||
text, lang: used, detected: detected || null,
|
||||
confidence: confidence == null ? null : Number(confidence.toFixed(3)),
|
||||
ms, audio_ms: audioMs,
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Synthesise one piece of text.
|
||||
* @returns {Promise<{audio: Float32Array, sampling_rate: number}>}
|
||||
*/
|
||||
export async function synthesize(text, lang = 'ta') {
|
||||
const clean = (text || '').trim();
|
||||
if (!clean) return null;
|
||||
|
||||
const use = supportsTTS(lang) ? lang : 'en';
|
||||
const tts = await getTTS(use);
|
||||
|
||||
const t0 = Date.now();
|
||||
const out = await tts(clean);
|
||||
const ms = Date.now() - t0;
|
||||
const audioMs = (out.audio.length / out.sampling_rate) * 1000;
|
||||
logger.debug(`🔈 TTS ${use}: ${clean.length} chars → ${ms}ms for ${audioMs.toFixed(0)}ms (RTF ${(ms / audioMs).toFixed(2)}x)`);
|
||||
|
||||
return { audio: out.audio, sampling_rate: out.sampling_rate };
|
||||
}
|
||||
|
||||
/**
|
||||
* Split into speakable pieces. Short prompts reach audio sooner, and a sentence
|
||||
* boundary is a clean place to be interrupted.
|
||||
*/
|
||||
export function sentences(text, max = 200) {
|
||||
const out = [];
|
||||
for (const raw of String(text || '').split(/(?<=[.!?।])\s+/)) {
|
||||
let s = raw.trim();
|
||||
if (!s) continue;
|
||||
while (s.length > max) {
|
||||
const cut = s.lastIndexOf(' ', max);
|
||||
out.push(s.slice(0, cut > 0 ? cut : max).trim());
|
||||
s = s.slice(cut > 0 ? cut : max).trim();
|
||||
}
|
||||
if (s) out.push(s);
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
/**
|
||||
* Strip block markdown before speaking — tables and code read terribly aloud,
|
||||
* and the visual blocks are already on screen.
|
||||
*/
|
||||
export function speakable(markdown = '') {
|
||||
return markdown
|
||||
.replace(/```[\s\S]*?```/g, ' ')
|
||||
.replace(/^\s*\|.*\|\s*$/gm, ' ')
|
||||
.replace(/^\s*[-*]\s+/gm, '')
|
||||
.replace(/^#{1,6}\s*/gm, '')
|
||||
.replace(/\*\*([^*]+)\*\*/g, '$1')
|
||||
.replace(/`([^`]+)`/g, '$1')
|
||||
.replace(/\[([^\]]+)\]\([^)]+\)/g, '$1')
|
||||
.replace(/₹\s?([\d,.]+)/g, 'rupees $1')
|
||||
.replace(/\s{2,}/g, ' ')
|
||||
.trim();
|
||||
}
|
||||
@@ -0,0 +1,118 @@
|
||||
// ============================================
|
||||
// Speech models — all ONNX, all CPU, all in this Node process.
|
||||
//
|
||||
// There is no GPU and no Python. That is the whole point: the AWS host has
|
||||
// neither, and a second service was one more thing to deploy and keep alive.
|
||||
//
|
||||
// VAD Silero 2 MB endpointing
|
||||
// STT Whisper base ~80 MB Tamil + English + language detection
|
||||
// TTS MMS-TTS VITS ~114 MB per language, feed-forward
|
||||
//
|
||||
// Measured on an i7-10850H, CPU only:
|
||||
// TTS RTF 0.28x (3.5x faster than realtime)
|
||||
// STT ~1.2 s for 4 s of audio
|
||||
//
|
||||
// Two findings worth keeping:
|
||||
//
|
||||
// * VITS is feed-forward. The earlier Parler-TTS attempt was autoregressive
|
||||
// and ran at RTF ~5x — i.e. 5x SLOWER than realtime — which is why voice was
|
||||
// unusable even on a GPU. Architecture mattered far more than hardware here.
|
||||
//
|
||||
// * int8 is a trap for a model this small: dynamic quantisation made TTS 5.7x
|
||||
// SLOWER than fp32 (RTF 1.67x vs 0.28x) because the quantise/dequantise
|
||||
// overhead dominates. We ship fp32 deliberately.
|
||||
// ============================================
|
||||
import path from 'node:path';
|
||||
import { fileURLToPath } from 'node:url';
|
||||
import { pipeline, env } from '@huggingface/transformers';
|
||||
import logger from '../utils/logger.js';
|
||||
|
||||
const ROOT = path.resolve(path.dirname(fileURLToPath(import.meta.url)), '../..');
|
||||
|
||||
// Hub downloads are cached here so a container restart does not re-fetch.
|
||||
env.cacheDir = process.env.SPEECH_CACHE_DIR || path.join(ROOT, '.transformers-cache');
|
||||
// Tamil is loaded from a folder we exported ourselves — no public ONNX build
|
||||
// of mms-tts-tam exists. See scripts/export-tamil-tts.py.
|
||||
env.localModelPath = path.join(ROOT, 'assets/tts');
|
||||
|
||||
const STT_MODEL = process.env.STT_MODEL || 'onnx-community/whisper-base';
|
||||
|
||||
/** TTS voice per language. Tamil is local; English comes from the Hub. */
|
||||
const VOICES = {
|
||||
ta: { id: 'mms-tts-tam', local: true },
|
||||
en: { id: 'Xenova/mms-tts-eng', local: false },
|
||||
};
|
||||
|
||||
export const LANGUAGES = [
|
||||
{ code: 'auto', label: 'Auto-detect', native: 'Auto' },
|
||||
{ code: 'ta', label: 'Tamil', native: 'தமிழ்' },
|
||||
{ code: 'en', label: 'English', native: 'English' },
|
||||
];
|
||||
|
||||
const cache = new Map();
|
||||
let sttPromise = null;
|
||||
|
||||
/**
|
||||
* Models load on first use, not at boot. A CRM restart should not wait ~10 s
|
||||
* for speech models that most sessions never touch.
|
||||
*/
|
||||
async function loadOnce(key, build) {
|
||||
if (!cache.has(key)) {
|
||||
const t0 = Date.now();
|
||||
cache.set(key, build().then((m) => {
|
||||
logger.info(`🔊 loaded ${key} in ${((Date.now() - t0) / 1000).toFixed(1)}s`);
|
||||
return m;
|
||||
}).catch((e) => {
|
||||
cache.delete(key); // let the next attempt retry
|
||||
throw e;
|
||||
}));
|
||||
}
|
||||
return cache.get(key);
|
||||
}
|
||||
|
||||
export async function getSTT() {
|
||||
if (!sttPromise) {
|
||||
sttPromise = loadOnce(STT_MODEL, () =>
|
||||
// q8 is the right call for Whisper — unlike VITS it is big enough that
|
||||
// quantisation is a clear win.
|
||||
pipeline('automatic-speech-recognition', STT_MODEL, { dtype: 'q8' }),
|
||||
).catch((e) => { sttPromise = null; throw e; });
|
||||
}
|
||||
return sttPromise;
|
||||
}
|
||||
|
||||
export async function getTTS(lang = 'ta') {
|
||||
const voice = VOICES[lang] || VOICES.ta;
|
||||
const prev = env.allowRemoteModels;
|
||||
try {
|
||||
// Local folders must not be looked up on the Hub, and vice versa.
|
||||
env.allowRemoteModels = !voice.local;
|
||||
return await loadOnce(`tts:${voice.id}`, () =>
|
||||
pipeline('text-to-speech', voice.id, { dtype: 'fp32' }),
|
||||
);
|
||||
} finally {
|
||||
env.allowRemoteModels = prev;
|
||||
}
|
||||
}
|
||||
|
||||
export const supportsTTS = (lang) => Boolean(VOICES[lang]);
|
||||
|
||||
/** Warm the models the deployment actually expects to use. */
|
||||
export async function warmup(langs = ['ta', 'en']) {
|
||||
try {
|
||||
await getSTT();
|
||||
for (const l of langs) await getTTS(l);
|
||||
logger.info('🔊 speech models warm');
|
||||
} catch (e) {
|
||||
logger.warn(`speech warmup failed (will retry on first use): ${e.message}`);
|
||||
}
|
||||
}
|
||||
|
||||
export function speechStatus() {
|
||||
return {
|
||||
stt_model: STT_MODEL,
|
||||
tts_voices: Object.fromEntries(Object.entries(VOICES).map(([k, v]) => [k, v.id])),
|
||||
loaded: [...cache.keys()],
|
||||
languages: LANGUAGES,
|
||||
};
|
||||
}
|
||||
@@ -0,0 +1,152 @@
|
||||
// ============================================
|
||||
// Endpointing — Silero VAD via onnxruntime-node.
|
||||
//
|
||||
// Deciding turn boundaries on the server rather than in the browser keeps the
|
||||
// rule in one place for every future channel (a phone bridge has no
|
||||
// AudioWorklet), and gives the server the signal it needs for barge-in: it has
|
||||
// to know the user started talking while the assistant was still speaking.
|
||||
// ============================================
|
||||
import path from 'node:path';
|
||||
import { fileURLToPath } from 'node:url';
|
||||
import fs from 'node:fs/promises';
|
||||
import ort from 'onnxruntime-node';
|
||||
import logger from '../utils/logger.js';
|
||||
|
||||
const ROOT = path.resolve(path.dirname(fileURLToPath(import.meta.url)), '../..');
|
||||
const MODEL_URL = 'https://huggingface.co/onnx-community/silero-vad/resolve/main/onnx/model.onnx';
|
||||
const MODEL_PATH = path.join(process.env.SPEECH_CACHE_DIR || path.join(ROOT, '.transformers-cache'), 'silero-vad.onnx');
|
||||
|
||||
// Silero wants exactly 512 samples at 16 kHz (32 ms). The browser sends 40 ms
|
||||
// chunks, so audio is buffered and drained in exact frames rather than forcing
|
||||
// the client to match.
|
||||
const FRAME = 512;
|
||||
const RATE = 16000;
|
||||
const FRAME_MS = (FRAME / RATE) * 1000;
|
||||
|
||||
let sessionPromise = null;
|
||||
|
||||
async function getSession() {
|
||||
if (sessionPromise) return sessionPromise;
|
||||
sessionPromise = (async () => {
|
||||
try {
|
||||
await fs.access(MODEL_PATH);
|
||||
} catch {
|
||||
logger.info('⬇️ fetching Silero VAD (2 MB)…');
|
||||
const res = await fetch(MODEL_URL);
|
||||
if (!res.ok) throw new Error(`VAD download failed: ${res.status}`);
|
||||
await fs.mkdir(path.dirname(MODEL_PATH), { recursive: true });
|
||||
await fs.writeFile(MODEL_PATH, Buffer.from(await res.arrayBuffer()));
|
||||
}
|
||||
const s = await ort.InferenceSession.create(MODEL_PATH);
|
||||
logger.info('🎚️ Silero VAD ready');
|
||||
return s;
|
||||
})().catch((e) => { sessionPromise = null; throw e; });
|
||||
return sessionPromise;
|
||||
}
|
||||
|
||||
export const vadOptions = {
|
||||
threshold: Number(process.env.VAD_THRESHOLD ?? 0.5),
|
||||
// Trailing silence that ends a turn. Too short truncates someone who pauses
|
||||
// mid-sentence; too long makes the assistant feel sluggish.
|
||||
silenceMs: Number(process.env.VAD_SILENCE_MS ?? 700),
|
||||
// Ignore blips, so a cough or a door does not open a turn.
|
||||
minSpeechMs: Number(process.env.VAD_MIN_SPEECH_MS ?? 250),
|
||||
// Audio kept from BEFORE detection, so word onsets are not clipped.
|
||||
prefixMs: Number(process.env.VAD_PREFIX_MS ?? 300),
|
||||
maxUtteranceMs: Number(process.env.VAD_MAX_UTTERANCE_MS ?? 30000),
|
||||
};
|
||||
|
||||
/** Streaming endpointer. One instance per connection. */
|
||||
export class Endpointer {
|
||||
constructor(opts = {}) {
|
||||
this.o = { ...vadOptions, ...opts };
|
||||
this.pending = new Float32Array(0);
|
||||
this.prefixFrames = Math.max(1, Math.round(this.o.prefixMs / FRAME_MS));
|
||||
this.reset();
|
||||
}
|
||||
|
||||
reset() {
|
||||
this.speaking = false;
|
||||
this.speechMs = 0;
|
||||
this.silenceMs = 0;
|
||||
this.buffer = [];
|
||||
this.prefix = [];
|
||||
// Silero is recurrent: this 2x1x128 state carries across frames and must
|
||||
// be reset between turns or the model stays biased by the last utterance.
|
||||
this.state = new ort.Tensor('float32', new Float32Array(2 * 1 * 128), [2, 1, 128]);
|
||||
this.pending = new Float32Array(0);
|
||||
}
|
||||
|
||||
/**
|
||||
* Feed float32 mono @16k.
|
||||
* @returns {Promise<{utterances: Float32Array[], started: boolean}>}
|
||||
* `started` flips the moment speech begins — that is the barge-in signal.
|
||||
*/
|
||||
async push(pcm) {
|
||||
const session = await getSession();
|
||||
|
||||
const merged = new Float32Array(this.pending.length + pcm.length);
|
||||
merged.set(this.pending);
|
||||
merged.set(pcm, this.pending.length);
|
||||
this.pending = merged;
|
||||
|
||||
const utterances = [];
|
||||
let started = false;
|
||||
let offset = 0;
|
||||
|
||||
while (this.pending.length - offset >= FRAME) {
|
||||
const frame = this.pending.subarray(offset, offset + FRAME);
|
||||
offset += FRAME;
|
||||
|
||||
const out = await session.run({
|
||||
input: new ort.Tensor('float32', frame, [1, FRAME]),
|
||||
sr: new ort.Tensor('int64', BigInt64Array.from([BigInt(RATE)]), []),
|
||||
state: this.state,
|
||||
});
|
||||
this.state = out.stateN ?? out.state_n ?? this.state;
|
||||
const voiced = out.output.data[0] >= this.o.threshold;
|
||||
|
||||
if (!this.speaking) {
|
||||
this.prefix.push(Float32Array.from(frame));
|
||||
if (this.prefix.length > this.prefixFrames) this.prefix.shift();
|
||||
|
||||
if (voiced) {
|
||||
this.speechMs += FRAME_MS;
|
||||
if (this.speechMs >= this.o.minSpeechMs) {
|
||||
this.speaking = true;
|
||||
this.silenceMs = 0;
|
||||
this.buffer = this.prefix; // open the turn with the pre-roll
|
||||
this.prefix = [];
|
||||
started = true;
|
||||
}
|
||||
} else {
|
||||
this.speechMs = 0;
|
||||
}
|
||||
continue;
|
||||
}
|
||||
|
||||
this.buffer.push(Float32Array.from(frame));
|
||||
if (voiced) this.silenceMs = 0;
|
||||
else this.silenceMs += FRAME_MS;
|
||||
|
||||
const spokenMs = this.buffer.length * FRAME_MS;
|
||||
if (this.silenceMs >= this.o.silenceMs || spokenMs >= this.o.maxUtteranceMs) {
|
||||
utterances.push(concat(this.buffer));
|
||||
this.reset();
|
||||
}
|
||||
}
|
||||
|
||||
this.pending = this.pending.slice(offset);
|
||||
return { utterances, started };
|
||||
}
|
||||
}
|
||||
|
||||
function concat(frames) {
|
||||
const total = frames.reduce((n, f) => n + f.length, 0);
|
||||
const out = new Float32Array(total);
|
||||
let i = 0;
|
||||
for (const f of frames) { out.set(f, i); i += f.length; }
|
||||
return out;
|
||||
}
|
||||
|
||||
export const warmupVad = () => getSession().catch(() => {});
|
||||
Reference in New Issue
Block a user