跳到主要内容

数据截至 (上游 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:22pipeline/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-909898-1067),因为「用户开口了」必须比「继续播上一句」优先。

类型别名

每条队列的合法载荷在 src/speech_to_speech/pipeline/queue_types.py 集中定义(如 VADOutItemTTSInItem),避免大段 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):

  1. 预解析:只认 --stt / --llm_backend / --tts / --mac-optimal-settings 四个,确定选了谁(s2s_pipeline.py:195-204)。
  2. 正式解析:只把被选中后端的参数 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-173run_realtime_audio_client
local同进程起服务器 + 客户端,走 loopbackbuild_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.pyBaseHandler.run, should_process_input, output_for_queue
软重置 / 退出哨兵src/speech_to_speech/pipeline/control.py, pipeline/messages.pySESSION_END, PIPELINE_END, is_control_message
队列载荷类型src/speech_to_speech/pipeline/queue_types.pyVADOutItem, TTSInItem, AudioOutItem
handler 链拼装src/speech_to_speech/s2s_pipeline.py_build_handlers, _build_pipeline_unit, build_pipeline
两级参数解析src/speech_to_speech/s2s_pipeline.pyparse_arguments, _parse_selected_cli_configs, _mac_preset_defaults
线程管理src/speech_to_speech/utils/thread_manager.pyThreadManager
后端注册表src/speech_to_speech/backend_registry.pyBackendSpec, normalize_dataclass_config, create_backend_handler
流水线单元 / 池src/speech_to_speech/api/openai_realtime/pipeline_unit.pyPipelineUnit, SessionState
服务器线程src/speech_to_speech/api/openai_realtime/server.pyRealtimeServer.run
命令分发src/speech_to_speech/cli.pyparse_command, parse_talk_arguments