外观
Java/TS 开发者源码导读与实战指南 (stream-vad-proxy)
本文档专为具有 Java 或 TypeScript / Node.js 开发背景的工程师编写,旨在帮助你快速理解 onnx-platform/stream-vad-proxy(前置 VAD 流式代理)的核心架构、设计意图与 Rust 源码实现。
1. 为什么需要这个服务?(业务背景)
在典型的 AI 语音对话(如智能外呼、AI 坐席)场景中,系统架构如下:
text
[用户电话] <--- SIP RTP ---> [FreeSWITCH] <--- WebSocket ---> [ASR 语音识别]
|
v
[AI 业务决策 / 播报 TTS]业务痛点:打断延迟过高
当 AI 正在向用户播放语音话术时,用户突然开口说话(例如打断说“不需要,谢谢”),系统需要以最快速度打断(Barge-In)当前的 TTS 播放,并倾听用户说话。
- 对需要本代理补齐能力的厂商接口,首个识别文本通常晚于本地 VAD 边沿;具体延迟取决于网络和厂商模型。
- 本地
sherpa-asr-server与部分云网关具备自己的 VAD / 断句能力;本代理只补齐既无开嗓事件、也无服务端断句的厂商接口。
解决方案:前置透明 VAD 代理
stream-vad-proxy 作为一个“透明中间件”,串接在 FreeSWITCH mod_audio_fork 与上游 ASR 之间:
- 开嗓即打断:代理内部加载超轻量级的 Silero VAD 模型(约 640KB),单次 32ms 音频帧推理仅需 0.1ms。在检测到人声起始的那一帧,立即向下游 FreeSWITCH 注入包含
speechStarted: true的虚拟事件帧,实现几乎零延迟的即时打断; - 静音闸门(Gate):用户不说话的静音期,代理不向上游 ASR 转发音频,不仅大幅节省上行带宽,还能直接降低按时长计费的云端 API 费用;
- 零业务侵入:下游保留平台调用所需的 WebSocket 握手、metadata、PCM16、零长 flush 与事件外形,接入只需将 FreeSWITCH 的推流地址
audioFork.wsUrl改指到代理端口(默认10093);这不承诺与任一 ASR 的全部能力一致。
2. Java / TS 开发者的 Rust 心智模型转换
Rust 没有虚拟机(JVM)也没有单线程事件循环(Node.js EventLoop),但其设计模式与 Java/TS 有强烈的对应关系:
+-----------------------------------------------------------------------------------------------+
| 概念域 | Java 对应物 | TypeScript 对应物 | Rust (当前项目具体体现) |
+-----------------------------------------------------------------------------------------------+
| 堆对象与所有权 | 对象指针引用 (GC 垃圾回收) | 变量引用 (V8 GC 回收) | Arc<T> (原子引用计数指针)|
| 跨线程安全共享 | synchronized / ReentrantLock | 单线程无锁概念 | std::sync::Mutex / Arc |
| 异步并发模型 | Virtual Threads (Loom) / 线程池| Promise / libuv 循环 | Tokio 协程 (Task) |
| 消息队列解耦 | LinkedBlockingQueue / Disruptor| EventEmitter / RxJS | tokio::sync::mpsc 消息信道|
| 资源自动释放 | try-with-resources / AutoClose| try ... finally | impl Drop 析构函数 (RAII)|
| 错误传递方式 | throw Exception + try/catch | throw Error + try/catch| Result<T, E> 与 ? 语法糖 |
| 异步分支竞争 | CompletableFuture.anyOf() | Promise.race() | tokio::select! 宏 |
+-----------------------------------------------------------------------------------------------+关键重点 1:Tokio 运行时与 spawn_blocking
- 在 Node.js 中,你不能在主线程执行同步文件 IO 或复杂加解密,否则整个 Event Loop 会卡死;
- 在 Netty / Spring WebFlux 中,CPU 密集型计算必须扔给独立的
EventExecutorGroup,不能占用 IO Worker; - 在 Rust Tokio 中亦然:Tokio 负责网络 IO 的 Worker 线程数量默认等于 CPU 核心数。Silero VAD 的模型推理涉及 ONNX 原生 C 库计算。在建立会话加载模型时,代码显式使用
tokio::task::spawn_blocking将其交由 Tokio 的专用阻塞线程池执行,坚决不阻塞网络事件泵。
关键重点 2:RAII 与 SessionGuard
在 Java 中限制并发连接数,我们可能会用 Semaphore 配合 try { ... } finally { semaphore.release(); }。 在 Rust 中,downstream.rs 定义了 SessionGuard:
rust
struct SessionGuard {
active_sessions: Arc<AtomicUsize>,
}
impl Drop for SessionGuard {
fn drop(&mut self) {
self.active_sessions.fetch_sub(1, Ordering::AcqRel);
}
}当 WebSocket 连接断开、协程发生 Panic 或正常 return 退出其所在作用域时,Rust 编译器保证自动调用 drop(),原子计数安全递减,绝对不会发生“连接断开却泄漏配额”的问题。
3. 系统核心数据流与时序全景
下图展示单条呼叫链路从握手到推流再到关闭的完整数据流向:
text
[下游 FreeSWITCH] [stream-vad-proxy (session.rs)] [上游 ASR 服务]
| | |
|----- 1. GET /audio (WS 握手) ------>| |
|<---- 2. 101 Switching Protocols ---| |
| |----- 3. 异步连接上游 ASR ----------->|
| |<---- 4. 上游握手完成 ----------------|
| | |
|----- 5. 首帧 JSON 元数据 ---------->| (捕获 channel_uuid, 缓存 metadata) |
| | |
|===== 6. 静音音频流 (PCM16 LE) =====>| |
| | [Silero VAD: 未检测到说话] |
| | [存入 1秒 Pre-roll 环形缓冲队列] |
| | (闸门闭合: 不向上游转发音频) |
| | |
|===== 7. 用户开口发声音频 ==========>| |
| | [Silero VAD: 检测到上升沿!] |
|<---- 8. 即时注入 speechStarted 帧 --| (FreeSWITCH 秒级打断 TTS 播报!) |
| | |
| | [冲刷 Pre-roll 缓冲 + 选通当前音频] |
| |===== 9. 将缓冲与实时音频发给上游 ===>|
| | |
| |<==== 10. 上游返回 Partial 文本 =====|
|<---- 11. 透传 Partial 识别文本 -----| (回显首帧 metadata 原文) |
| | |
|===== 12. 用户停止发音 (静音达到阈值)=>| |
| | [Silero VAD: 检测到下降沿完成] |
|<---- 13. 下发 segmentCompleted 帧 --| (通知本句结束) |
| | |
| |<==== 14. 上游返回 Final 终态文本 ===|
|<---- 15. 透传 Final 识别文本 -------| |
| | |
|----- 16. WebSocket 关闭帧 --------->| |
| |-- 17. 排空残余音频促上游定稿 ------>|
|<---- 18. 确认关闭 1000 -------------|<-- 18. 上游连接正常释放 ------------|4. 关键核心机制与设计考量
4.1 异步读写分离与 MPSC 解耦
在 session.rs 中,WebSocket 连接被拆分为两半:
sink(负责向下游写数据帧);stream(负责从下游读音频二进制帧)。
上游 ASR 的网络读写由一个独立的后台异步 Task 处理,并通过 tokio::sync::mpsc::channel(128) 有界信道与会话主协程通信。 主会话通过 tokio::select! 宏同时监听下游入站音频与上游回传事件:
- 只要有音频进站,立刻跑 VAD 并通过 Channel 发给上游;
- 只要上游出字,立刻通过 Sink 刷给下游;
- 双向异步,完全没有阻塞点。
4.2 Silero VAD 整窗切分机制 (vad.rs)
- 模型约束:Silero VAD v4 模型的 ONNX 计算图严格要求每次输入的音频样本数恰好为 512 个采样点(对应 16kHz 采样率下的 32 毫秒)。
- 缓冲对齐:FreeSWITCH
mod_audio_fork推送音频时,通常按 RTP 节奏发送(如每 20ms 一包,即 320 个采样点 = 640 字节)。320 不能被 512 整除! - 处理方式:
vad.rs内部维护了pending: Vec<f32>暂存区。入站音频先追加进暂存区,每次切出完整的 512 采样送进模型推理,不足 512 的余数保留到下一包音频到来,确保模型推理永远对齐整窗。
4.3 Pre-roll 前置缓冲(为什么不会“吃字头”?)
用户从生理上开始发音,到 VAD 概率累积超过门限(例如 0.35)判定为 speech_started,通常需要经历 60ms ~ 120ms 的爬升期。
- 如果在判定说话的那一刻才开始给上游发音频,前 100ms 的辅音/声母就会丢失(即俗称的“吞字”、“吃首字”);
- 解决方案:代理维护一个容量为 1 秒(16000 采样 = 32000 字节)的
pre_roll: VecDeque<u8>环形缓冲队列。静音期间音频在队列中滚动更新;一旦检测到开嗓,先把 Pre-roll 缓冲里的音频全量发给上游,再接着发后续实时音频,上游 ASR 就能完整听到字头!
4.4 讯飞私有云协议适配器 (upstream/iflytek.rs)
对于当前科大讯飞私有云实时识别(IAT)上游,代理使用厂商自定义的 HTTP JSON-RPC 状态机:
ssb(Session Start Begin):发送会话初始化参数(appid, 语言, 采样率等),获取sid(会话 ID);auw(Audio Write):将有效语音分块切片(每块 1280 字节),Base64 编码后通过 HTTP POST 写入上游;grs(Get Result):轮询拉取当前会话的流式增量识别文本;sse(Session End):在 VAD 静音切句或通话结束时发送定稿终止信号。
代理把讯飞这一整套复杂的 RPC 状态机完全封装在 iflytek.rs 内部,对外依然表现为统一的 ASR 事件,下游完全无感。
4.5 抗混叠重采样 (upstream/resample.rs)
部分厂商私有云(如特定的讯飞私有化集群)仅支持 8kHz 电话音频输入,而 FreeSWITCH mod_audio_fork 下发的是 16kHz PCM:
- 不能简单地“隔一个点扔一个点”(抽取下采样),因为高频成分会折叠回低频区形成频谱混叠(Aliasing),导致识别出严重杂音;
resample.rs纯 Rust 手写了带有 Kaiser 窗的 Sinc FIR 低通滤波器,先滤除 4kHz 以上的高频分量,再进行 2 倍抽取下采样,保证重采样后音质纯净。
5. 源码目录与阅读推荐顺序
建议按以下顺序阅读 onnx-platform/stream-vad-proxy/src/ 下的代码:
protocol.rs:定义下游帧结构体与VadFlags。弄懂下游需要什么字段;config.rs:配置结构体定义、默认值与防御性校验;main.rs:服务主入口、命令行参数加载、模型预检与 Tokio 运行时初始化;downstream.rs:Axum HTTP/WS 路由、握手鉴权、SessionGuard并发名额控制;vad.rs:Silero VAD 原生调用封装、512 整窗余数处理、起止边沿状态机;session.rs(核心!):单条会话主循环、tokio::select!读写调度、Pre-roll 缓冲与闸门选通;upstream.rs:上游抽象分派与 MPSC 管道;upstream/iflytek.rs&upstream/resample.rs:讯飞协议状态机与降采样滤波器;recording.rs&logging.rs:异步录音落盘与按天滚动日志。
6. 本地排障与实战验证
健康检查探针
bash
curl http://127.0.0.1:10093/health返回中关注:
profile: 当前激活的 Profile 名称(默认"iflytek");provider: 当前生效的上游适配器实现(如"iflytek-iat")。
推流联调测试
可以使用平台现有的 Node.js 测试客户端对代理推流:
bash
bun run onnx-platform/sherpa-asr-server/request-ws.js --file <audio.wav> --port 10093- 把测试客户端指向代理(
--port 10093),观察第一个带speechStarted: true的空文本事件是否早于首个非空 partial;需要单独验证 sherpa-asr-server 时再把端口改为10096,但二者上游能力并不相同。