跳到正文

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 之间:

  1. 开嗓即打断:代理内部加载超轻量级的 Silero VAD 模型(约 640KB),单次 32ms 音频帧推理仅需 0.1ms。在检测到人声起始的那一帧,立即向下游 FreeSWITCH 注入包含 speechStarted: true 的虚拟事件帧,实现几乎零延迟的即时打断;
  2. 静音闸门(Gate):用户不说话的静音期,代理不向上游 ASR 转发音频,不仅大幅节省上行带宽,还能直接降低按时长计费的云端 API 费用;
  3. 零业务侵入:下游保留平台调用所需的 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 状态机:

  1. ssb (Session Start Begin):发送会话初始化参数(appid, 语言, 采样率等),获取 sid(会话 ID);
  2. auw (Audio Write):将有效语音分块切片(每块 1280 字节),Base64 编码后通过 HTTP POST 写入上游;
  3. grs (Get Result):轮询拉取当前会话的流式增量识别文本;
  4. 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/ 下的代码:

  1. protocol.rs:定义下游帧结构体与 VadFlags。弄懂下游需要什么字段;
  2. config.rs:配置结构体定义、默认值与防御性校验;
  3. main.rs:服务主入口、命令行参数加载、模型预检与 Tokio 运行时初始化;
  4. downstream.rs:Axum HTTP/WS 路由、握手鉴权、SessionGuard 并发名额控制;
  5. vad.rs:Silero VAD 原生调用封装、512 整窗余数处理、起止边沿状态机;
  6. session.rs(核心!):单条会话主循环、tokio::select! 读写调度、Pre-roll 缓冲与闸门选通;
  7. upstream.rs:上游抽象分派与 MPSC 管道;
  8. upstream/iflytek.rs & upstream/resample.rs:讯飞协议状态机与降采样滤波器;
  9. 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,但二者上游能力并不相同。

最后更新于:

文档与代码在同一仓库维护,现有 Markdown 是唯一内容源。