数据截至 (上游 commit b21e54d6a845)
骨架:handler、队列与后端注册表
这一章讲什么: 先把「一条流水线长什么样」讲清楚——一个 handler 是什么、队列上跑什么、几条流水线怎么共享一个服务器、以及「
--tts kokoro为什么就能换掉整个语音合成」。
1. 一个 handler 就是一个线程加一个生成器
它要解决的小问题
四个环节(VAD/STT/LLM/TTS)的实现差别巨大:有的调 HTTP,有的跑 torch,有的跑 MLX。但它们在流水线里的行为必须一致:从上游拿一件东西、处理、把结果交给下游、随时能被叫停。
思路
用一个基类把「循环 + 队列 + 停止 + 计时 + 过期过滤」全部写死,子类只实现 process(),而且 process 是生成器——一次输入可以产出零个、一个或很多个输出(TTS 一句话产出上百块 PCM 就靠这个)。
结构图
queue_in queue_out
│ ▲
▼ │
┌─────────────────── BaseHandler.run() ───────────────┴──────┐
│ while not stop_event: │
│ item = queue_in.get(timeout=0.1) ← 超时是为了能退出 │
│ ├─ SESSION_END ? → on_session_end() 后原样转发 │
│ ├─ b"END" ? → 跳出循环(哨兵,防死锁) │
│ ├─ should_process_input(item) 为假 ? → 丢弃 │
│ ├─ 是 PipelineEvent ? → 原样转发(保序,不进 process) │
│ └─ for out in process(item): │
│ should_emit_output(out) 为假 ? → 丢弃 │
│ queue_out.put(output_for_queue(out, item)) │
└────────────────────────────────────────────────────────────┘
真实实现
主循环在 BaseHandler.run(src/speech_to_speech/baseHandler.py:104-166)。几个值得注意的点:
- 超时轮询而非阻塞:
self.queue_in.get(timeout=0.1)(baseHandler.py:111)。代价是空转,好处是stop_event每 100ms 一定被看到。 - 两种哨兵:
PIPELINE_END = b"END"让线程退出,SESSION_END是软重置——handler 清掉本会话状态但线程继续活着(pipeline/control.py:22、pipeline/messages.py:361)。 - 事件透传:
isinstance(item, PipelineEvent)时直接queue_out.put(baseHandler.py:142-144)。这条 3 行的分支是保序的地基:助手文本事件和它对应的音频走同一条队列,谁也超不了谁。 - 异常不杀线程:
process()抛异常只记 log(baseHandler.py:162-163),流水线继续跑。
三个可覆写的钩子
| 钩子 | 默认行为 | 谁覆写了它、为什么 |
|---|---|---|
should_process_input | 丢弃过期世代的输入 | 基类已实现取消过滤;BaseSTTHandler 再叠一层回合过滤 |
should_emit_output | 全放行 | BaseSTTHandler 用它拦住「算完才发现回合已作废」的结果 |
output_for_queue | 裸音频包成 AudioOutput | 基类实现,给音频贴上 cancel_generation / response_key 标签 |
基类版 should_process_input(baseHandler.py:55-80)做两件事:等待预取事务被认领(见 05 章),以及比对 cancel_scope.is_stale(item.cancel_generation)。
2. 队列上跑的到底是什么
八条队列
一条流水线(PipelineUnit)自带八条 queue.Queue,全部在 _build_pipeline_unit 里创建(src/speech_to_speech/s2s_pipeline.py:483-490):
客户端音频 ──► recv_audio_chunks_queue ──► [VAD]
│
spoken_prompt_queue ◄── ┘
│
▼
[STT] ──► stt_output_queue ──► [TranscriptionNotifier]
│
RealtimeService ──► text_prompt_queue ◄───────────────────────────────┘(仅转发控制消息)
▲ │
│ ▼
│ [LLM] ──► lm_response_queue ──► [LMOutputProcessor]
│ │
│ lm_processed_queue ◄─┘
│ │
│ ▼
│ [TTS] ──► send_audio_chunks_queue ──► 发送循环
│
└──── text_output_queue(旁路:VAD 事件、转写事件、工具就绪事件)────────────────► 发送循环
两条出口队列,不是一条。 send_audio_chunks_queue 是有序主路(助手文本、工具、音频、终结事件都在这),text_output_queue 是旁路(不需要等 TTS 的事件:说话开始/结束、转写增量)。发送循环每轮先读旁路再读主路(websocket_router.py:824-909 与 898-1067),因为「用户开口了」必须比「继续播上一句」优先。
类型别名
每条队列的合法载荷在 src/speech_to_speech/pipeline/queue_types.py 集中定义(如 VADOutItem、TTSInItem),避免大段 Union 在代码里到处复制。
3. handler 链是怎么拼起来的
拼装函数
_build_handlers(src/speech_to_speech/s2s_pipeline.py:348-452)返回一个列表:
[VAD] → [STT 或 AudioInputNotifier] → (可选 TranscriptionNotifier) → [LLM] → [LMOutputProcessor] → [TTS]
一处分叉值得注意:当 --stt none 时,注册表里那个后端的 capabilities.bypasses_transcription_notifier 为真(backend_registry.py:292-298),于是不插 TranscriptionNotifier,STT 位置换成 AudioInputNotifier,音频直接送进支持音频输入的 LLM。判断写在 s2s_pipeline.py:387-388:
needs_notifier = not stt_backend.spec.capabilities.bypasses_transcription_notifier
stt_queue_out: Queue[Any] = stt_output_queue if needs_notifier else text_prompt_queue
即:不需要通知器时,STT 阶段的输出队列直接改接到 LLM 的输入队列上。同一套 handler 链,靠改接线实现两种拓扑。
线程管理
ThreadManager(src/speech_to_speech/utils/thread_manager.py)极简:每个 handler 一个非守护线程,stop() 时先 stop_event.set() 再 join(timeout=5.0),超时只打警告(thread_manager.py:35-39)。没有优雅回收,靠 PIPELINE_END 哨兵和超时轮询兜底。
4. 多路并发:PipelineUnit 池
它要解决的小问题
模型很贵,不能每来一个客户端就加载一遍;但会话状态(历史、回合号)必须彼此隔离。
做法
--num_pipelines N 会构造 N 个 PipelineUnit,每个都完整加载自己的模型和 handler,共用一个 uvicorn 服务器(build_pipeline,s2s_pipeline.py:544-581)。WebSocket 连接进来时抢一个 session is None 的空闲 unit,抢不到就拒绝。
┌──────────── RealtimeServer(单 uvicorn / 单端口)────────────┐
│ claim: 找 session is None 的 unit;满了就发 session_limit_reached │
└──────┬──────────────────┬──────────────────┬─────────────────┘
▼ ▼ ▼
┌────────────┐ ┌────── ──────┐ ┌────────────┐
│ Unit 0 │ │ Unit 1 │ │ Unit N-1 │
│ 8 条队列 │ │ 8 条队列 │ │ 8 条队列 │
│ 6 个线程 │ │ 6 个线程 │ │ 6 个线程 │
│ 自己的 Chat │ │ 自己的 Chat │ │ 自己的 Chat │
└────────────┘ └────────────┘ └────────────┘
所以 --num_pipelines 4 的显存开销大约是四份模型。这是刻意的简单:没有 batch,没有模型共享,换来的是每条流水线内部完全不用加锁。
一个真实的副作用:Apple Silicon 上所有 MLX 推理走同一把全局锁(见 03 章),池大于 1 时渐进转写会疯狂抢锁失败刷屏,于是启动时直接把实时转写关掉(s2s_pipeline.py:634-640)。
会话释放的「排空」协议
断开连接时不能立刻把 unit 还给下一个人——上一个会话的半成品可能还在管道里。释放路径是:清空四条队列 → 塞一个 SESSION_END → 等它穿过整条 handler 链回到输出队列(每个 handler 都会转发它)→ 发送循环看到后 session.drained.set() → 才真正 unit.session = None。相关逻辑在 _clean_unit(websocket_router.py:213-234)和 _release_unit_after_drain(websocket_router.py:285-337),超时 SESSION_END_QUARANTINE_TIMEOUT_S = 180.0 后把 unit 标记为隔离而非复用。
5. 后端注册表:换模型为什么只用改一个参数
它要解决的小问题
七种 STT、四种 LLM、五种 TTS,每种有自己的一堆 CLI 参数和可选依赖。如果写成 if stt == "whisper": ... elif ...,参数解析和依赖检查会烂成一坨。
思路
把每个后端描述成一条声明式记录 BackendSpec(src/speech_to_speech/backend_registry.py:81-100),包含:名字、种类、参数 dataclass 类型、构造函数、参数前缀、可选依赖名、能力标志。三张注册表就是三个 dict(backend_registry.py:289-503)。
关键设计:参数前缀自动剥离
normalize_dataclass_config(backend_registry.py:122-142)把 --parakeet_tdt_model_name 这种带前缀的 CLI 参数,自动变成 handler 的 model_name= 关键字参数,并把所有 gen_* 收进 gen_kwargs。所以 handler 的 setup() 签名可以写得很干净,而 CLI 里不同后端的同名参数不打架。
# 示意,非源码:注册表如何声明一个后端
BackendSpec(
"parakeet-tdt", # --stt parakeet-tdt
"stt",
ParakeetTDTSTTHandlerArguments, # 它的 CLI 参数 dataclass
_create_parakeet, # 构造函数(拿 HandlerContext + config)
config_prefix="parakeet_tdt", # 剥掉这个前缀再传给 setup()
)
两级参数解析
因为「有哪些参数」取决于「选了哪个后端」,parse_arguments 要解析两遍(s2s_pipeline.py:170-280):
- 预解析:只认
--stt/--llm_backend/--tts/--mac-optimal-settings四个,确定选了谁(s2s_pipeline.py:195-204)。 - 正式解析:只把被选中后端的参数 dataclass 交给
HfArgumentParser。
那没被选中后端的参数怎么办?_parse_selected_cli_configs(s2s_pipeline.py:130-167)拿剩余 token 再用一个「所有未选中后端」的兼容 parser 试一遍:能认出来就打一条 warning 忽略掉,认不出来才报错。这让老的启动脚本不会因为换了后端就直接崩。
依赖缺失的报错翻译
create_backend_handler(backend_registry.py:184-193)捕获 ImportError,查 spec 的 required_extra,把裸的 ModuleNotFoundError: kokoro 翻译成 pip install "speech-to-speech[kokoro]"。小细节,但对用户体验的杠杆很大。
6. 三个命令与启动路径
| 命令 | 行为 | 实现 |
|---|---|---|
serve | 只起服务器 | run_pipeline_command("serve", ...) → build_pipeline |
talk | 只起麦克风/扬声器客户端 | cli.py:171-173 → run_realtime_audio_client |
local | 同进程 起服务器 + 客户端,走 loopback | build_local_pipeline(s2s_pipeline.py:584-615) |
local 的实现很直白:先 build_pipeline(host="127.0.0.1"),再造一个 RealtimeAudioClient 指向 ws://127.0.0.1:<port>/v1/realtime,把两边的 handler 列表拼成一个 ThreadManager(s2s_pipeline.py:599-615)。客户端也是一个 handler——它有 run() 和 stop_event,所以能被同一套线程管理器托管。
旧的 --mode 参数保留兼容:--mode realtime 映射到 serve,--mode local 映射到 local,其余值直接报错退出(cli.py:74-87)。
7. 代码地图
| 主题 | 文件路径 | 符号名 |
|---|---|---|
| handler 线程主循环 | src/speech_to_speech/baseHandler.py | BaseHandler.run, should_process_input, output_for_queue |
| 软重置 / 退出哨兵 | src/speech_to_speech/pipeline/control.py, pipeline/messages.py | SESSION_END, PIPELINE_END, is_control_message |
| 队列载荷类型 | src/speech_to_speech/pipeline/queue_types.py | VADOutItem, TTSInItem, AudioOutItem |
| handler 链拼装 | src/speech_to_speech/s2s_pipeline.py | _build_handlers, _build_pipeline_unit, build_pipeline |
| 两级参数解析 | src/speech_to_speech/s2s_pipeline.py | parse_arguments, _parse_selected_cli_configs, _mac_preset_defaults |
| 线程管理 | src/speech_to_speech/utils/thread_manager.py | ThreadManager |
| 后端注册表 | src/speech_to_speech/backend_registry.py | BackendSpec, normalize_dataclass_config, create_backend_handler |
| 流水线单元 / 池 | src/speech_to_speech/api/openai_realtime/pipeline_unit.py | PipelineUnit, SessionState |
| 服务器线程 | src/speech_to_speech/api/openai_realtime/server.py | RealtimeServer.run |
| 命令分发 | src/speech_to_speech/cli.py | parse_command, parse_talk_arguments |