removed the voice agent

This commit is contained in:
2026-08-29 14:01:56 +05:30
parent bef978da6f
commit 8ae5e7ebb9
27 changed files with 7 additions and 2645 deletions
-18
View File
@@ -55,21 +55,3 @@ RATE_LIMIT_MAX=40
# --- Artifacts --- # --- Artifacts ---
ARTIFACT_DIR=./storage/artifacts ARTIFACT_DIR=./storage/artifacts
ARTIFACT_TTL_HOURS=72 ARTIFACT_TTL_HOURS=72
# --- Voice (speech-to-speech, CPU, in-process) ---
# Models are ONNX via Transformers.js — no GPU, no Python, no second service.
# STT onnx-community/whisper-tiny.en English only; half the RAM of base
# TTS Xenova/mms-tts-eng from the Hub
# Tamil is built and tested (assets/tts/mms-tts-tam, exported locally — no
# public ONNX exists). Enable it with VOICE_LANGUAGES=en,ta and a multilingual
# STT_MODEL, but budget ~200 MB more resident for the extra voice.
VOICE_LANGUAGES=en
SPEECH_WARMUP=false
STT_MODEL=onnx-community/whisper-tiny.en
# Below this confidence, the user's preferred language beats the detector.
DETECT_CONFIDENCE=0.6
# Endpointing
VAD_SILENCE_MS=700
VAD_MIN_SPEECH_MS=250
VAD_PREFIX_MS=300
-1
View File
@@ -55,7 +55,6 @@ __pycache__/
/huggingface/ /huggingface/
/models/ /models/
.transformers-cache/ .transformers-cache/
*.onnx
*.safetensors *.safetensors
*.ckpt *.ckpt
*.pt *.pt
+6 -11
View File
@@ -1,13 +1,8 @@
# ============================================ # ============================================
# WeLe Agentic AI — production image # WeLe Agentic AI — production image
# #
# Speech runs in this same process: ONNX on CPU via Transformers.js. No GPU, # Text-only agent service. node:20-slim rather than alpine: several native
# no Python, no second service. # dependencies ship glibc binaries that will not load against musl.
#
# node:20-slim, NOT alpine. onnxruntime-node ships glibc binaries and Alpine is
# musl, so the native module fails at load with:
# Error loading shared library ld-linux-x86-64.so.2 (needed by libonnxruntime.so.1)
# `sharp`, pulled in by Transformers.js, has the same constraint.
# ============================================ # ============================================
FROM node:20-slim AS deps FROM node:20-slim AS deps
@@ -28,10 +23,10 @@ COPY --from=deps /app/node_modules ./node_modules
COPY package.json ./ COPY package.json ./
COPY src/ ./src/ COPY src/ ./src/
# Both of these get a named volume mounted over them in compose. Docker seeds a # A named volume is mounted over this in compose. Docker seeds a NEW volume
# NEW volume from the image path — including ownership — so they must exist here # from the image path — including ownership — so it must exist here owned by
# owned by `node`, or the unprivileged process gets EACCES on first write. # `node`, or the unprivileged process gets EACCES on first write.
RUN mkdir -p /app/storage/artifacts /app/.transformers-cache && chown -R node:node /app/storage /app/.transformers-cache RUN mkdir -p /app/storage/artifacts && chown -R node:node /app/storage
USER node USER node
EXPOSE 4000 EXPOSE 4000
-3
View File
@@ -1,3 +0,0 @@
{
"<unk>": 58
}
-82
View File
@@ -1,82 +0,0 @@
{
"activation_dropout": 0.1,
"architectures": [
"VitsModel"
],
"attention_dropout": 0.1,
"depth_separable_channels": 2,
"depth_separable_num_layers": 3,
"dtype": "float32",
"duration_predictor_dropout": 0.5,
"duration_predictor_filter_channels": 256,
"duration_predictor_flow_bins": 10,
"duration_predictor_kernel_size": 3,
"duration_predictor_num_flows": 4,
"duration_predictor_tail_bound": 5.0,
"ffn_dim": 768,
"ffn_kernel_size": 3,
"flow_size": 192,
"hidden_act": "relu",
"hidden_dropout": 0.1,
"hidden_size": 192,
"initializer_range": 0.02,
"layer_norm_eps": 1e-05,
"layerdrop": 0.1,
"leaky_relu_slope": 0.1,
"model_type": "vits",
"noise_scale": 0.667,
"noise_scale_duration": 0.8,
"num_attention_heads": 2,
"num_hidden_layers": 6,
"num_speakers": 1,
"posterior_encoder_num_wavenet_layers": 16,
"prior_encoder_num_flows": 4,
"prior_encoder_num_wavenet_layers": 4,
"resblock_dilation_sizes": [
[
1,
3,
5
],
[
1,
3,
5
],
[
1,
3,
5
]
],
"resblock_kernel_sizes": [
3,
7,
11
],
"sampling_rate": 16000,
"speaker_embedding_size": 0,
"speaking_rate": 1.0,
"spectrogram_bins": 513,
"transformers_version": "4.57.3",
"upsample_initial_channel": 512,
"upsample_kernel_sizes": [
16,
16,
4,
4
],
"upsample_rates": [
8,
8,
2,
2
],
"use_bias": true,
"use_stochastic_duration_prediction": true,
"vocab_size": 58,
"wavenet_dilation_rate": 1,
"wavenet_dropout": 0.0,
"wavenet_kernel_size": 5,
"window_size": 4
}
@@ -1,4 +0,0 @@
{
"per_channel": false,
"reduce_range": false
}
@@ -1,4 +0,0 @@
{
"pad_token": "3",
"unk_token": "<unk>"
}
-115
View File
@@ -1,115 +0,0 @@
{
"version": "1.0",
"truncation": null,
"padding": null,
"added_tokens": [
{
"id": 58,
"content": "<unk>",
"single_word": false,
"lstrip": false,
"rstrip": false,
"normalized": false,
"special": true
}
],
"normalizer": {
"type": "Sequence",
"normalizers": [
{
"type": "Lowercase"
},
{
"type": "Replace",
"pattern": {
"Regex": "[^012345679 '_aஅஆஇஈஉஊஎஏஐஒஓகஙசஜஞடணதநனபமயரறலளழவஷஸஹாிீுூெேைொோௌ்]"
},
"content": ""
},
{
"type": "Strip",
"strip_left": true,
"strip_right": true
},
{
"type": "Replace",
"pattern": {
"Regex": "(?=.)|(?<!^)$"
},
"content": "3"
}
]
},
"pre_tokenizer": {
"type": "Split",
"pattern": {
"Regex": ""
},
"behavior": "Isolated",
"invert": false
},
"post_processor": null,
"decoder": null,
"model": {
"vocab": {
"0": 47,
"1": 44,
"2": 23,
"3": 0,
"4": 54,
"5": 57,
"6": 36,
"7": 14,
"9": 31,
" ": 7,
"'": 13,
"_": 4,
"a": 15,
"அ": 1,
"ஆ": 45,
"இ": 38,
"ஈ": 2,
"உ": 3,
"ஊ": 11,
"எ": 37,
"ஏ": 16,
"ஐ": 52,
"ஒ": 27,
"ஓ": 49,
"க": 6,
"ங": 50,
"ச": 30,
"ஜ": 53,
"ஞ": 29,
"ட": 22,
"ண": 48,
"த": 41,
"ந": 5,
"ன": 35,
"ப": 46,
"ம": 26,
"ய": 39,
"ர": 25,
"ற": 28,
"ல": 21,
"ள": 43,
"ழ": 24,
"வ": 17,
"ஷ": 55,
"ஸ": 33,
"ஹ": 19,
"ா": 9,
"ி": 32,
"ீ": 12,
"ு": 51,
"ூ": 20,
"ெ": 10,
"ே": 8,
"ை": 34,
"ொ": 56,
"ோ": 42,
"ௌ": 40,
"்": 18
}
}
}
@@ -1,31 +0,0 @@
{
"add_blank": true,
"added_tokens_decoder": {
"0": {
"content": "3",
"lstrip": false,
"normalized": false,
"rstrip": false,
"single_word": false,
"special": true
},
"58": {
"content": "<unk>",
"lstrip": false,
"normalized": false,
"rstrip": false,
"single_word": false,
"special": true
}
},
"clean_up_tokenization_spaces": true,
"extra_special_tokens": {},
"is_uroman": false,
"language": "tam",
"model_max_length": 1000000000000000019884624838656,
"normalize": true,
"pad_token": "3",
"phonemize": false,
"tokenizer_class": "VitsTokenizer",
"unk_token": "<unk>"
}
-60
View File
@@ -1,60 +0,0 @@
{
" ": 7,
"'": 13,
"0": 47,
"1": 44,
"2": 23,
"3": 0,
"4": 54,
"5": 57,
"6": 36,
"7": 14,
"9": 31,
"_": 4,
"a": 15,
"அ": 1,
"ஆ": 45,
"இ": 38,
"ஈ": 2,
"உ": 3,
"ஊ": 11,
"எ": 37,
"ஏ": 16,
"ஐ": 52,
"ஒ": 27,
"ஓ": 49,
"க": 6,
"ங": 50,
"ச": 30,
"ஜ": 53,
"ஞ": 29,
"ட": 22,
"ண": 48,
"த": 41,
"ந": 5,
"ன": 35,
"ப": 46,
"ம": 26,
"ய": 39,
"ர": 25,
"ற": 28,
"ல": 21,
"ள": 43,
"ழ": 24,
"வ": 17,
"ஷ": 55,
"ஸ": 33,
"ஹ": 19,
"ா": 9,
"ி": 32,
"ீ": 12,
"ு": 51,
"ூ": 20,
"ெ": 10,
"ே": 8,
"ை": 34,
"ொ": 56,
"ோ": 42,
"ௌ": 40,
"்": 18
}
+1 -17
View File
@@ -28,12 +28,6 @@ services:
- NODE_ENV=production - NODE_ENV=production
- PORT=4000 - PORT=4000
# English-only keeps one TTS voice resident (~200 MB). Tamil is built and
# tested — VOICE_LANGUAGES=en,ta plus a multilingual STT_MODEL enables it,
# at roughly 200 MB more.
- VOICE_LANGUAGES=en
- STT_MODEL=onnx-community/whisper-tiny.en
# Reuse the CRM's Redis by service name on the shared network. # Reuse the CRM's Redis by service name on the shared network.
# Keys are namespaced with REDIS_PREFIX, so the two never collide. # Keys are namespaced with REDIS_PREFIX, so the two never collide.
- REDIS_ENABLED=true - REDIS_ENABLED=true
@@ -56,22 +50,13 @@ services:
volumes: volumes:
# Generated xlsx/pdf/pptx survive rebuilds; swept on a TTL by the app. # Generated xlsx/pdf/pptx survive rebuilds; swept on a TTL by the app.
- artifacts:/app/storage/artifacts - artifacts:/app/storage/artifacts
# Speech models are fetched from HuggingFace on first use (~200 MB).
# Without this they re-download on every restart and the first voice turn
# after a deploy stalls for a minute.
- speech_cache:/app/.transformers-cache
# A 2 vCPU / 3.7 GB host already runs the CRM, chat-service, Redis and # A 2 vCPU / 3.7 GB host already runs the CRM, chat-service, Redis and
# Milvus. Capping this container keeps a runaway turn from starving them. # Milvus. Capping this container keeps a runaway turn from starving them.
#
# 1 GB, not 768 MB: the speech models are resident once voice is used —
# measured 729 MB (Whisper tiny.en 415 MB + MMS-TTS English 203 MB + VAD
# 28 MB + the agent itself). 768 MB left no headroom, and an OOM kill takes
# text chat down with voice. Text-only sessions stay near 80 MB.
deploy: deploy:
resources: resources:
limits: limits:
memory: 1024M memory: 768M
logging: logging:
driver: json-file driver: json-file
@@ -90,4 +75,3 @@ networks:
volumes: volumes:
artifacts: artifacts:
speech_cache:
-977
View File
File diff suppressed because it is too large Load Diff
-3
View File
@@ -12,7 +12,6 @@
"author": "WeLe EdTech", "author": "WeLe EdTech",
"license": "ISC", "license": "ISC",
"dependencies": { "dependencies": {
"@huggingface/transformers": "^4.2.0",
"@langchain/core": "^1.1.18", "@langchain/core": "^1.1.18",
"@langchain/langgraph": "^1.1.0", "@langchain/langgraph": "^1.1.0",
"@langchain/openai": "^1.5.10", "@langchain/openai": "^1.5.10",
@@ -27,12 +26,10 @@
"ioredis": "^5.10.1", "ioredis": "^5.10.1",
"jsonwebtoken": "^9.0.2", "jsonwebtoken": "^9.0.2",
"mongoose": "^9.2.3", "mongoose": "^9.2.3",
"onnxruntime-node": "^1.24.3",
"pdfkit": "^0.17.2", "pdfkit": "^0.17.2",
"pptxgenjs": "^4.0.1", "pptxgenjs": "^4.0.1",
"uuid": "^13.0.0", "uuid": "^13.0.0",
"winston": "^3.19.0", "winston": "^3.19.0",
"ws": "^8.21.3",
"zod": "^3.25.76" "zod": "^3.25.76"
} }
} }
-66
View File
@@ -1,66 +0,0 @@
/* Generate the tokenizer.json that Transformers.js needs for the exported
Tamil VITS model.
`save_pretrained` does not emit one: VitsTokenizer is a "slow" tokenizer with
no fast counterpart, so Python writes vocab.json + tokenizer_config.json and
nothing else. Transformers.js only reads tokenizer.json, so we synthesise it
from the exported vocab, mirroring the structure of the working English
model (Xenova/mms-tts-eng) exactly.
The four normalizer steps, in order:
1. Lowercase — no-op for Tamil, matters for embedded Latin/digits
2. Replace — drop every character outside the vocab
3. Strip — trim surrounding whitespace
4. Replace — insert the blank token between every character,
which is what `add_blank: true` means for VITS.
Omit this and the audio comes out garbled.
*/
import fs from 'node:fs';
import path from 'node:path';
const DIR = 'assets/tts/mms-tts-tam';
const vocab = JSON.parse(fs.readFileSync(path.join(DIR, 'vocab.json'), 'utf8'));
const cfg = JSON.parse(fs.readFileSync(path.join(DIR, 'tokenizer_config.json'), 'utf8'));
// The blank/pad token is whichever character maps to id 0.
const blank = Object.keys(vocab).find((k) => vocab[k] === 0);
const unk = cfg.unk_token ?? '<unk>';
const unkId = vocab[unk] ?? Object.keys(vocab).length;
// Character class of everything we keep. Escape the regex metacharacters that
// are still special inside a negated class.
const escaped = Object.keys(vocab)
.filter((c) => c !== unk)
.map((c) => (']\\^-'.includes(c) ? '\\' + c : c))
.join('');
const tokenizer = {
version: '1.0',
truncation: null,
padding: null,
added_tokens: [{
id: unkId, content: unk,
single_word: false, lstrip: false, rstrip: false, normalized: false, special: true,
}],
normalizer: {
type: 'Sequence',
normalizers: [
{ type: 'Lowercase' },
{ type: 'Replace', pattern: { Regex: `[^${escaped}]` }, content: '' },
{ type: 'Strip', strip_left: true, strip_right: true },
...(cfg.add_blank ? [{ type: 'Replace', pattern: { Regex: '(?=.)|(?<!^)$' }, content: blank }] : []),
],
},
pre_tokenizer: { type: 'Split', pattern: { Regex: '' }, behavior: 'Isolated', invert: false },
post_processor: null,
decoder: null,
model: { vocab },
};
fs.writeFileSync(path.join(DIR, 'tokenizer.json'), JSON.stringify(tokenizer, null, 1));
console.log(`wrote ${DIR}/tokenizer.json`);
console.log(` vocab ${Object.keys(vocab).length} tokens`);
console.log(` blank token ${JSON.stringify(blank)} (id 0)`);
console.log(` unk ${JSON.stringify(unk)} (id ${unkId})`);
console.log(` add_blank ${cfg.add_blank}`);
-99
View File
@@ -1,99 +0,0 @@
"""One-time export of facebook/mms-tts-tam to ONNX.
No public ONNX build of Tamil MMS-TTS exists, so we make one. This runs ONCE on
a workstation; the committed artefact is what ships. The service itself is pure
JavaScript and never needs Python or this script.
The output must match the contract Transformers.js expects for VITS, taken from
the working English model (Xenova/mms-tts-eng):
inputs : input_ids, attention_mask
outputs: waveform, spectrogram
Layout produced (mirrors the HF repo so Transformers.js can load the folder):
assets/tts/mms-tts-tam/
config.json, tokenizer.json, vocab.json, …
onnx/model.onnx fp32
onnx/model_quantized.onnx int8 ← what we ship
"""
from __future__ import annotations
import json
import shutil
import sys
from pathlib import Path
import torch
from transformers import AutoTokenizer, VitsModel
MODEL = "facebook/mms-tts-tam"
OUT = Path("assets/tts/mms-tts-tam")
ONNX_DIR = OUT / "onnx"
ONNX_DIR.mkdir(parents=True, exist_ok=True)
print(f"loading {MODEL} …")
model = VitsModel.from_pretrained(MODEL).eval()
tok = AutoTokenizer.from_pretrained(MODEL)
# VITS has a stochastic duration predictor. Exporting with noise left on bakes
# RandomNormalLike nodes into the graph, which is fine and keeps prosody
# natural — but seed it so this export is reproducible.
torch.manual_seed(0)
class Exportable(torch.nn.Module):
"""Return only (waveform, spectrogram), in that order — the JS side indexes
outputs by name, but a plain tuple keeps the exported graph simple."""
def __init__(self, m: VitsModel) -> None:
super().__init__()
self.m = m
def forward(self, input_ids: torch.Tensor, attention_mask: torch.Tensor):
out = self.m(input_ids=input_ids, attention_mask=attention_mask)
return out.waveform, out.spectrogram
sample = tok("வணக்கம், இது ஒரு சோதனை.", return_tensors="pt")
fp32 = ONNX_DIR / "model.onnx"
print("exporting to ONNX …")
torch.onnx.export(
Exportable(model),
(sample["input_ids"], sample["attention_mask"]),
str(fp32),
input_names=["input_ids", "attention_mask"],
output_names=["waveform", "spectrogram"],
dynamic_axes={
"input_ids": {0: "batch", 1: "sequence"},
"attention_mask": {0: "batch", 1: "sequence"},
"waveform": {0: "batch", 1: "samples"},
"spectrogram": {0: "batch", 2: "frames"},
},
opset_version=17,
do_constant_folding=True,
)
print(f" fp32: {fp32.stat().st_size / 1e6:.1f} MB")
# ── int8 ────────────────────────────────────────────────────────────────────
try:
from onnxruntime.quantization import QuantType, quantize_dynamic
q = ONNX_DIR / "model_quantized.onnx"
quantize_dynamic(str(fp32), str(q), weight_type=QuantType.QUInt8)
print(f" int8: {q.stat().st_size / 1e6:.1f} MB")
except Exception as e: # noqa: BLE001
print(f" quantisation skipped: {e}")
# ── tokenizer + config, so the folder loads standalone ──────────────────────
tok.save_pretrained(OUT)
model.config.to_json_file(OUT / "config.json")
# Transformers.js reads this to pick a default dtype.
(OUT / "quantize_config.json").write_text(json.dumps({"per_channel": False, "reduce_range": False}, indent=2))
print("\nwrote:")
for p in sorted(OUT.rglob("*")):
if p.is_file():
print(f" {p.relative_to(OUT)} ({p.stat().st_size / 1e6:.2f} MB)")
-21
View File
@@ -1,21 +0,0 @@
/* Does auto mode now route Tamil to Tamil instead of silently using English? */
import { synthesize, transcribe } from '../src/speech/index.js';
const resample = (a, from, to) => {
const r = from / to, out = new Float32Array(Math.floor(a.length / r));
for (let i = 0; i < out.length; i++) { const p = i * r, k = Math.floor(p); out[i] = a[k] + (a[Math.min(k + 1, a.length - 1)] - a[k]) * (p - k); }
return out;
};
for (const [lang, text] of [
['ta', 'மூவாயிரம் நானூற்று இருபத்தேழு புதிய லீட்கள் உள்ளன.'],
['en', 'How many new leads did we receive today?'],
]) {
const spoken = await synthesize(text, lang);
const audio = resample(spoken.audio, spoken.sampling_rate, 16000);
const r = await transcribe(audio, 'auto', 'ta');
const ok = r.lang === lang ? 'PASS' : 'FAIL';
console.log(`${ok} spoke ${lang} → routed ${r.lang} (detected ${r.detected} @ ${r.confidence}) ${r.ms}ms`);
console.log(` ${JSON.stringify(r.text.slice(0, 80))}`);
}
process.exit(0);
-62
View File
@@ -1,62 +0,0 @@
/* Can we get real language detection out of Whisper in Transformers.js? */
import { AutoProcessor, WhisperForConditionalGeneration, Tensor, env } from '@huggingface/transformers';
import { synthesize } from '../src/speech/index.js';
env.cacheDir = './.transformers-cache';
const MODEL = 'onnx-community/whisper-base';
const processor = await AutoProcessor.from_pretrained(MODEL);
const model = await WhisperForConditionalGeneration.from_pretrained(MODEL, { dtype: 'q8' });
const tok = processor.tokenizer;
// Whisper emits one language token right after <|startoftranscript|>. Reading
// that distribution is a single decoder step — far cheaper than transcribing
// twice to see which language "looks better".
const id = (t) => tok.encode(t, { add_special_tokens: false })[0];
const sot = id('<|startoftranscript|>');
const CANDIDATES = ['en', 'ta'];
const langIds = CANDIDATES.map((c) => id(`<|${c}|>`));
console.log('sot:', sot, '| language token ids:', JSON.stringify(Object.fromEntries(CANDIDATES.map((c, i) => [c, langIds[i]]))));
function resample(a, from, to) {
const r = from / to;
const out = new Float32Array(Math.floor(a.length / r));
for (let i = 0; i < out.length; i++) {
const p = i * r, k = Math.floor(p);
out[i] = a[k] + (a[Math.min(k + 1, a.length - 1)] - a[k]) * (p - k);
}
return out;
}
for (const [lang, text] of [
['ta', 'மூவாயிரம் நானூற்று இருபத்தேழு புதிய லீட்கள் உள்ளன.'],
['en', 'There are three thousand four hundred and twenty seven new leads.'],
]) {
const spoken = await synthesize(text, lang);
const audio = resample(spoken.audio, spoken.sampling_rate, 16000);
const inputs = await processor(audio);
const t0 = Date.now();
const out = await model({
...inputs,
decoder_input_ids: new Tensor('int64', BigInt64Array.from([BigInt(sot)]), [1, 1]),
});
const ms = Date.now() - t0;
const logits = out.logits;
const last = logits.dims[1] - 1;
const vocab = logits.dims[2];
const row = logits.data.slice(last * vocab, (last + 1) * vocab);
const scores = langIds.map((id) => Number(row[id]));
const max = Math.max(...scores);
const exp = scores.map((s) => Math.exp(s - max));
const sum = exp.reduce((a, b) => a + b, 0);
const probs = exp.map((e) => e / sum);
const best = probs.indexOf(Math.max(...probs));
console.log(`spoken ${lang} → detected ${CANDIDATES[best]} `
+ `(${CANDIDATES.map((c, i) => `${c} ${probs[i].toFixed(3)}`).join(', ')}) in ${ms}ms `
+ `${CANDIDATES[best] === lang ? '✅' : '❌'}`);
}
process.exit(0);
-46
View File
@@ -1,46 +0,0 @@
/* English-only footprint: does it fit the 768 MB container cap on AWS?
Loads exactly what an English-only deployment needs and reports RSS after
each stage, so the answer is measured rather than estimated.
*/
const mb = () => Math.round(process.memoryUsage().rss / 1048576);
const step = (label) => console.log(` ${label.padEnd(34)} RSS ${String(mb()).padStart(4)} MB`);
step('baseline (node + agent code)');
const { synthesize, transcribe, Endpointer } = await import('../src/speech/index.js');
step('after importing speech module');
// VAD
const ep = new Endpointer();
await ep.push(new Float32Array(16000));
step('+ Silero VAD');
// TTS English
const spoken = await synthesize('There are three thousand four hundred and twenty seven new leads.', 'en');
step('+ MMS-TTS English');
// STT
const resample = (a, from, to) => {
const r = from / to, out = new Float32Array(Math.floor(a.length / r));
for (let i = 0; i < out.length; i++) { const p = i * r, k = Math.floor(p); out[i] = a[k] + (a[Math.min(k + 1, a.length - 1)] - a[k]) * (p - k); }
return out;
};
const audio = resample(spoken.audio, spoken.sampling_rate, 16000);
const heard = await transcribe(audio, 'en', 'en');
step('+ Whisper base (STT)');
// Steady state: a few turns, to see whether it keeps growing.
for (let i = 0; i < 3; i++) {
await synthesize('Checking the leads now.', 'en');
await transcribe(audio, 'en', 'en');
}
step('after 3 more turns');
console.log(`\n transcript: ${JSON.stringify(heard.text)}`);
console.log(` TTS rate : ${spoken.sampling_rate} Hz`);
const peak = mb();
const CAP = 768;
console.log(`\n peak ${peak} MB vs ${CAP} MB container cap → ${peak < CAP * 0.8 ? 'FITS ✅' : peak < CAP ? 'TIGHT ⚠️' : 'EXCEEDS ❌'}`);
process.exit(0);
-85
View File
@@ -1,85 +0,0 @@
/* 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); });
-44
View File
@@ -1,44 +0,0 @@
/* Do Whisper (STT) and MMS-TTS (TTS) actually run in Node on CPU? */
import { pipeline, env } from '@huggingface/transformers';
env.cacheDir = './.transformers-cache';
const t = (t0) => `${((performance.now() - t0) / 1000).toFixed(1)}s`;
// ── TTS: MMS-TTS Tamil (VITS, 36M, feed-forward) ────────────────────────────
console.log('[1/2] loading MMS-TTS Tamil…');
let t0 = performance.now();
const tts = await pipeline('text-to-speech', 'Xenova/mms-tts-eng', { dtype: 'fp32' });
console.log(` loaded in ${t(t0)}`);
const TA = 'There are three thousand four hundred leads in the new lead stage.';
await tts(TA); // warm
for (const [label, text] of [['short', TA], ['long', TA + ' ' + TA + ' ' + TA]]) {
t0 = performance.now();
const out = await tts(text);
const ms = performance.now() - t0;
const audioMs = (out.audio.length / out.sampling_rate) * 1000;
console.log(
` ${label.padEnd(5)} ${String(text.length).padStart(3)} chars → ${ms.toFixed(0)}ms `
+ `for ${audioMs.toFixed(0)}ms audio @ ${out.sampling_rate}Hz → RTF ${(ms / audioMs).toFixed(2)}x`,
);
}
// ── STT: Whisper (multilingual — Tamil, Hindi, English + detection) ─────────
console.log('\n[2/2] loading Whisper base…');
t0 = performance.now();
const stt = await pipeline('automatic-speech-recognition', 'onnx-community/whisper-base', { dtype: 'q8' });
console.log(` loaded in ${t(t0)}`);
// 4 s of quiet noise — proves the graph runs and times it.
const audio = Float32Array.from({ length: 16000 * 4 }, () => (Math.random() - 0.5) * 0.02);
t0 = performance.now();
const r = await stt(audio, { language: 'ta', task: 'transcribe' });
console.log(` 4000ms audio → ${(performance.now() - t0).toFixed(0)}ms → ${JSON.stringify(r.text).slice(0, 60)}`);
t0 = performance.now();
const r2 = await stt(audio, { language: 'en', task: 'transcribe' });
console.log(` english pass → ${(performance.now() - t0).toFixed(0)}ms → ${JSON.stringify(r2.text).slice(0, 60)}`);
console.log(`\nRSS ${(process.memoryUsage().rss / 1e9).toFixed(2)} GB`);
process.exit(0);
-55
View File
@@ -1,55 +0,0 @@
/* Does the locally-exported Tamil ONNX load and speak through Transformers.js? */
import { pipeline, env } from '@huggingface/transformers';
import fs from 'node:fs';
// Load from the local folder, not the Hub.
env.allowRemoteModels = false;
env.localModelPath = './assets/tts';
for (const dtype of ['q8', 'fp32']) {
try {
const t0 = performance.now();
const tts = await pipeline('text-to-speech', 'mms-tts-tam', { dtype });
const load = performance.now() - t0;
const TEXT = 'மூவாயிரம் நானூற்று இருபத்தேழு புதிய லீட்கள் உள்ளன.';
await tts(TEXT); // warm
const t1 = performance.now();
const out = await tts(TEXT);
const ms = performance.now() - t1;
const audioMs = (out.audio.length / out.sampling_rate) * 1000;
console.log(
`${dtype.padEnd(5)} load ${(load / 1000).toFixed(1)}s | `
+ `${ms.toFixed(0)}ms for ${audioMs.toFixed(0)}ms audio @ ${out.sampling_rate}Hz | `
+ `RTF ${(ms / audioMs).toFixed(2)}x`,
);
// Non-silent output is the real proof the graph is wired correctly.
const peak = out.audio.reduce((m, v) => Math.max(m, Math.abs(v)), 0);
console.log(` samples ${out.audio.length}, peak amplitude ${peak.toFixed(3)} ${peak > 0.01 ? '✅ audible' : '⚠️ SILENT'}`);
if (dtype === 'q8') {
const wav = toWav(out.audio, out.sampling_rate);
fs.writeFileSync('scripts/tamil-sample.wav', wav);
console.log(' wrote scripts/tamil-sample.wav — play it to judge quality');
}
} catch (e) {
console.log(`${dtype.padEnd(5)} FAILED: ${e.message.slice(0, 160)}`);
}
}
function toWav(samples, rate) {
const buf = Buffer.alloc(44 + samples.length * 2);
buf.write('RIFF', 0); buf.writeUInt32LE(36 + samples.length * 2, 4); buf.write('WAVE', 8);
buf.write('fmt ', 12); buf.writeUInt32LE(16, 16); buf.writeUInt16LE(1, 20); buf.writeUInt16LE(1, 22);
buf.writeUInt32LE(rate, 24); buf.writeUInt32LE(rate * 2, 28); buf.writeUInt16LE(2, 32); buf.writeUInt16LE(16, 34);
buf.write('data', 36); buf.writeUInt32LE(samples.length * 2, 40);
for (let i = 0; i < samples.length; i++) {
const s = Math.max(-1, Math.min(1, samples[i]));
buf.writeInt16LE(s < 0 ? s * 0x8000 : s * 0x7fff, 44 + i * 2);
}
return buf;
}
process.exit(0);
-72
View File
@@ -1,72 +0,0 @@
/* End-to-end check of the in-process speech pipeline: TTS → VAD → STT.
Synthesising a sentence and feeding that audio back through the endpointer
and recogniser exercises every stage with real speech, which a noise buffer
cannot do — silence never opens a VAD turn.
*/
import { synthesize, transcribe, sentences, speakable, Endpointer } from '../src/speech/index.js';
const say = (m) => console.log(m);
// ── 1. TTS both languages ───────────────────────────────────────────────────
const CASES = [
['ta', 'மூவாயிரம் நானூற்று இருபத்தேழு புதிய லீட்கள் உள்ளன.'],
['en', 'There are three thousand four hundred and twenty seven new leads.'],
];
const rendered = {};
for (const [lang, text] of CASES) {
const t0 = Date.now();
await synthesize(text, lang); // warm
const t1 = Date.now();
const out = await synthesize(text, lang);
const ms = Date.now() - t1;
const audioMs = (out.audio.length / out.sampling_rate) * 1000;
rendered[lang] = out;
say(`TTS ${lang} warm ${((t1 - t0) / 1000).toFixed(1)}s | ${ms}ms for ${audioMs.toFixed(0)}ms `
+ `@${out.sampling_rate}Hz | RTF ${(ms / audioMs).toFixed(2)}x`);
}
// ── 2. VAD: does synthesised speech open and close a turn? ──────────────────
function resample(audio, from, to) {
if (from === to) return audio;
const ratio = from / to;
const out = new Float32Array(Math.floor(audio.length / ratio));
for (let i = 0; i < out.length; i++) {
const p = i * ratio;
const a = Math.floor(p);
out[i] = audio[a] + (audio[Math.min(a + 1, audio.length - 1)] - audio[a]) * (p - a);
}
return out;
}
const ep = new Endpointer();
const speech = resample(rendered.en.audio, rendered.en.sampling_rate, 16000);
// Speech, then a second of silence so the endpointer closes the turn.
const withTail = new Float32Array(speech.length + 16000);
withTail.set(speech);
let started = false;
let captured = null;
for (let i = 0; i < withTail.length; i += 640) { // 40 ms chunks, as the browser sends
const { utterances, started: s } = await ep.push(withTail.subarray(i, Math.min(i + 640, withTail.length)));
if (s) started = true;
if (utterances.length) { captured = utterances[0]; break; }
}
say(`VAD speech detected: ${started ? 'yes' : 'NO'} | turn closed: ${captured ? 'yes' : 'NO'}`
+ (captured ? ` | captured ${(captured.length / 16000).toFixed(2)}s` : ''));
// ── 3. STT on that captured audio ───────────────────────────────────────────
if (captured) {
for (const lang of ['en', 'auto']) {
const r = await transcribe(captured, lang, 'en');
say(`STT ${lang.padEnd(4)} ${r.ms}ms → lang=${r.lang}${r.detected ? ` (heard ${r.detected})` : ''} → ${JSON.stringify(r.text.slice(0, 70))}`);
}
}
// ── 4. Text shaping ─────────────────────────────────────────────────────────
const md = '## Leads\n\n**3,427** in `new_lead`.\n\n| a | b |\n|---|---|\n| 1 | 2 |\n\n- Only 12% contacted.\nNext step is triage.';
say(`\nspeakable: ${JSON.stringify(speakable(md))}`);
say(`sentences: ${JSON.stringify(sentences(speakable(md)))}`);
say(`\nRSS ${(process.memoryUsage().rss / 1e9).toFixed(2)} GB`);
process.exit(0);
-287
View File
@@ -1,287 +0,0 @@
// ============================================
// 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 ──PCM16──► Endpointer ──► transcribe() ──► runTurn(graph)
// │ │
// browser ◄──float32──── synthesize() ◄── sentences ◄─────┘
//
// 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 { 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';
/** Spoken filler, said the instant a question lands. */
const ACK = {
ta: ['பார்க்கிறேன்.', 'ஒரு நிமிடம், பார்க்கிறேன்.'],
en: ['Let me check.', 'One moment, checking now.'],
};
/** Progress narration — spoken over the user's waiting time, so keep it short. */
const NARRATE = {
ta: { lead: 'லீட் விவரங்களைப் பார்க்கிறேன்.', analytics: 'புள்ளிவிவரங்களைச் சரிபார்க்கிறேன்.', conversation: 'உரையாடல்களைப் பார்க்கிறேன்.', default: 'தரவைச் சரிபார்க்கிறேன்.' },
en: { lead: 'Checking the leads.', analytics: 'Pulling the numbers.', conversation: 'Looking at the conversations.', default: 'Checking the data.' },
};
const NOTHING = { ta: 'பதில் கிடைக்கவில்லை.', en: 'I could not find an answer for that.' };
const OOPS = { ta: 'மன்னிக்கவும், ஒரு பிழை ஏற்பட்டது.', en: 'Sorry, something went wrong.' };
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;
class VoiceSession {
constructor(client, user) {
this.client = client;
this.user = user;
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.endpointer = new Endpointer();
this.busy = false;
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) {
if (this.client.readyState === WebSocket.OPEN) this.client.send(JSON.stringify(obj));
}
sendAudio(buf) {
if (this.client.readyState === WebSocket.OPEN) this.client.send(buf, { binary: true });
}
// ── microphone ───────────────────────────────────────────────────────────
/**
* 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);
for (let i = 0; i < pcm16.length; i++) pcm[i] = pcm16[i] / 32768;
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;
}
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' });
}
// 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) {
this.handleUtterance(utterance).catch((e) => logger.error(`turn failed: ${e.message}`));
}
}
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 });
// 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 ─────────────────────────────────────────────────────────────
/** 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) {
this.busy = true;
this.narrated.clear();
this.abort = new AbortController();
const seq = ++this.speakSeq;
// 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: question,
user: this.user,
channel: 'crm_chat', // voice users are staff
signal: this.abort.signal,
onEvent: (ev) => {
this.send(ev);
// 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.say(narrateFor(this.replyLang, agent), seq);
}
}
},
});
this.send({ type: 'result', blocks: result.blocks, usage: result.usage });
const answer = (result.blocks || [])
.filter((b) => b.type === 'text').map((b) => b.markdown).join(' ') || result.answer || '';
const parts = sentences(speakable(answer));
if (!parts.length) {
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 (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}`);
await this.say(OOPS[this.replyLang] || OOPS.en, seq);
}
} finally {
this.busy = false;
this.send({ type: 'idle' });
}
}
// ── control ──────────────────────────────────────────────────────────────
onMessage(data, isBinary) {
if (isBinary) return this.onAudio(data);
let msg;
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 = 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.send({ type: 'cancelled' });
} else if (msg.type === 'text' && msg.text) {
this.send({ type: 'transcript', text: msg.text, lang: this.replyLang, typed: true });
this.answer(msg.text);
}
return undefined;
}
close() {
this.speakSeq++;
this.abort?.abort();
}
}
/** Attach the voice WebSocket to the HTTP server. */
export function attachVoice(server) {
const wss = new WebSocketServer({ noServer: true });
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
// Browsers cannot set headers on a WebSocket, so the CRM token arrives as
// 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();
return;
}
wss.handleUpgrade(req, socket, head, (client) => {
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 (in-process, CPU)`);
return wss;
}
-12
View File
@@ -16,8 +16,6 @@ import agentRoutes, { artifactRouter } from './gateway/routes.js';
import { ensureDir, sweep } from './output/artifactStore.js'; import { ensureDir, sweep } from './output/artifactStore.js';
import crmApi from './tools/http/crmApi.js'; import crmApi from './tools/http/crmApi.js';
import { describeChains } from './orchestration/llm.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(); const app = express();
@@ -41,7 +39,6 @@ app.get('/health', async (_req, res) => {
redis: redisOk ? redisMode() : 'unavailable', redis: redisOk ? redisMode() : 'unavailable',
crm_api: crm.reachable ? 'reachable' : `unreachable (${crm.error || crm.status})`, crm_api: crm.reachable ? 'reachable' : `unreachable (${crm.error || crm.status})`,
models: describeChains(), models: describeChains(),
speech: speechStatus(),
uptime_s: Math.round(process.uptime()), uptime_s: Math.round(process.uptime()),
}); });
}); });
@@ -81,15 +78,6 @@ async function start() {
logger.info(` CRM API =${config.crmApi.base}`); logger.info(` CRM API =${config.crmApi.base}`);
}); });
// Voice is a WebSocket upgrade on the same port, so the browser needs no
// 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) => { const shutdown = (sig) => {
logger.info(`${sig} — shutting down`); logger.info(`${sig} — shutting down`);
server.close(() => process.exit(0)); server.close(() => process.exit(0));
-181
View File
@@ -1,181 +0,0 @@
// ============================================
// 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, sttIsEnglishOnly, defaultLanguage } from './models.js';
import logger from '../utils/logger.js';
export { LANGUAGES, warmup, speechStatus, defaultLanguage } 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 = defaultLanguage()) {
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;
// An English-only checkpoint has no language tokens to read, and passing
// `language` to it is rejected — so "auto" simply means English there.
if (lang === 'auto' && sttIsEnglishOnly()) {
used = 'en';
} else 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;
}
}
// An English-only checkpoint rejects BOTH `task` and `language` — it has no
// other mode to select. Multilingual builds require the language, since they
// silently default to English otherwise.
const opts = { return_timestamps: false };
if (!sttIsEnglishOnly()) {
opts.task = 'transcribe';
opts.language = used;
}
const result = await stt(audio, opts);
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 = defaultLanguage()) {
const clean = (text || '').trim();
if (!clean) return null;
const use = supportsTTS(lang) ? lang : defaultLanguage();
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();
}
-137
View File
@@ -1,137 +0,0 @@
// ============================================
// 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');
// whisper-tiny.en by default: measured, Whisper is the memory hog, not TTS.
// base cost 554 MB of an 879 MB total, which overran the 768 MB container cap.
// The .en build is half the size and, being English-only, cannot detect a
// language — which is fine when VOICE_LANGUAGES is just `en`.
const STT_MODEL = process.env.STT_MODEL || 'onnx-community/whisper-tiny.en';
/** True when the STT checkpoint is English-only and cannot identify languages. */
export const sttIsEnglishOnly = () => /\.en$/.test(STT_MODEL);
/** TTS voice per language. Tamil is exported locally; English is on the Hub. */
const ALL_VOICES = {
en: { id: 'Xenova/mms-tts-eng', local: false, label: 'English', native: 'English' },
ta: { id: 'mms-tts-tam', local: true, label: 'Tamil', native: 'தமிழ்' },
};
// Each extra language is a further ~200 MB resident. Enable only what the
// deployment actually speaks — English alone on the current AWS box.
const ENABLED = (process.env.VOICE_LANGUAGES || 'en')
.split(',').map((s) => s.trim()).filter((c) => ALL_VOICES[c]);
const VOICES = Object.fromEntries(ENABLED.map((c) => [c, ALL_VOICES[c]]));
export const LANGUAGES = [
// Auto-detect is only offered when there is a choice to make AND the STT
// model can actually detect — offering it otherwise is a lie.
...(ENABLED.length > 1 && !sttIsEnglishOnly()
? [{ code: 'auto', label: 'Auto-detect', native: 'Auto' }] : []),
...ENABLED.map((c) => ({ code: c, label: ALL_VOICES[c].label, native: ALL_VOICES[c].native })),
];
export const defaultLanguage = () => (LANGUAGES[0]?.code || 'en');
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) {
const voice = VOICES[lang] || VOICES[ENABLED[0]];
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 = ENABLED) {
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,
english_only_stt: sttIsEnglishOnly(),
tts_voices: Object.fromEntries(Object.entries(VOICES).map(([k, v]) => [k, v.id])),
loaded: [...cache.keys()],
languages: LANGUAGES,
};
}
-152
View File
@@ -1,152 +0,0 @@
// ============================================
// 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(() => {});