数据截至 (上游 commit c49982eb3aea)
生命周期与打断:Quest 这个 async 版 RAII
30 秒导读: 语音对话是一堆同时跑着的后台任务——一个在听(STT)、一个在想(LLM)、一个在说(TTS)。用户随时可能插嘴打断。本章讲 Unmute 怎么用一个叫
Quest的小抽象,把每条后台任务包成「有开场、有正戏、有谢幕」的单元,并保证打断时旧任务被干净地取消、不会污染新的对话。
本章只讲并发资源的生命周期与打断。对话状态机怎么判断轮到谁说话,见 02-stt-vad-turn-taking;逐词喂 TTS、按真实时间放音,见 03-tts-realtime-queue。这里假设你已经知道「什么时候该打断」,只关心「打断这个动作怎么干净地执行」。
1. 这是什么(零基础也能懂)
1.1 先看要解决的麻烦
想象你在跟语音助手对话。它正说到一半,你插了一句嘴。这一瞬间,系统里其实有好几件事正在半途中:
- LLM 还在一个词一个词地往外吐生成结果;
- TTS 服务正把这些词合成成音频、往你耳朵里推;
- 已经合成好的音频还排在播放队列里等着放。
如果只是「叫 LLM 停下」,你会听到:旧音频还在放几百毫秒、旧队列里的残余音频盖在新回答上、 后台那个旧的 TTS 任务可能还在悄悄往队列里塞东西。打断做不干净,体验就是「它没在听我说话」。
问题的本质是:每条后台任务都占着资源(一个 WebSocket 连接、一个后台协程、一段缓冲),打断时必须把这些资源成对地、按顺序地释放掉,不能只是「让它别跑了」。
1.2 一句话定义
Quest[T] 是一个 async 版的 RAII——把「获取资源 → 使用资源 → 释放资源」三步绑成一个对象,交给 async with 托管,保证退出时一定会跑清理。
RAII(Resource Acquisition Is Initialization,资源获取即初始化): C++/Rust 里的经典手法——对象一构造就拿到资源,一析构就自动还回去,靠作用域保证「拿了必还」。Python 没有确定性析构,所以 Unmute 用
async with+ 一个手写的类来模拟这个保证。
源码文件开头的注释把心态说得很直白:
"A desperate attempt at having some kind of RAII in Python." ——
unmute/quest_manager.py:1
翻译:「在 Python 里搞 RAII 的一次绝望尝试。」作者自己都承认这是权宜之计(还留了句「未来也许能用 TaskGroup + try/finally 做得更简洁」)。但它有效,值得学。
1.3 一句话直觉
把 Quest 想成一个可 召回的外派任务:
- 派它出去 =
__aenter__,它开始干活; - 召它回来 =
__aexit__/remove,它先做好收尾(关连接),再被取消。
而 QuestManager 是任务调度台,规矩只有一条:同名任务只能有一个。你派一个新的 "tts" 任务,调度台会先把旧的 "tts" 召回,再让新的上岗。这条「同名顶替」规则,就是整个打断机制的地基。
2. 顶层全景(它大概怎么转)
2.1 三个具名任务挂在 handler 上
UnmuteHandler 是一次对话连接的总管。它内部就挂着一个 QuestManager,里面最多同时活着三个具名 quest:
| quest 名 | 干什么 | 谁的 init | 在哪注册 |
|---|---|---|---|
"stt" | 听:流式转写用户语音 | find_instance("stt", ...) | unmute_handler.py:432 |
"tts" | 说:把词合成成音频 | 带指数退避的 find_instance("tts", ...) | unmute_handler.py:506 |
"llm" | 想:调 LLM 生成回答 | 无(from_run_step) | unmute_handler.py:181-182 |
stt 在连接建立时就启动、活整场;tts 和 llm 每个回合起一对、回合结束或被打断时销毁。
2.2 一张图:一次打断里发生了什么
下面这张图从上到下是打断时的动作顺序。读法:左边是触发,中间是 interrupt_bot 干的四件事,右边是每件事清理掉的东西。
用户插嘴
(STT 收到词 / VAD 判定打断)
│
▼
┌─────────────────────────────────────────────────────┐
│ interrupt_bot() unmute_handler.py:583 │
│ │
│ ① 记一个打断标记 ────────────► chat_history 加 │
│ add_chat_message_delta( INTERRUPTION_CHAR │
│ INTERRUPTION_CHAR) (让状态机知道被打断) │
│ │
│ ② self._clear_queue() ────────► 清 FastRTC 内部 │
│ 已排好的播放缓冲 │
│ │
│ ③ self.output_queue = Queue() ─► 换一个全新空队列; │
│ 旧 worker 再往旧队列 │
│ 塞东西也污染不到新的 │
│ │
│ ④ remove("tts") / remove("llm")► 触发各自 close(关 │
│ WebSocket)后 cancel │
└─────────────────────────────────────────────────────┘
│
▼
旧 TTS/LLM 资源已释放,output_queue 干净,可以开新回合
四件事的顺序是有讲究的(见 §4),不是随便排的。
2.3 主线走一遍(不进代码)
一次正常回合 + 一次打断,高层是这样:
- 回合开始:
_generate_response起一个"llm"quest;llm任务内部又通过start_up_tts起一个"tts"quest。两者都注册进QuestManager。 - 正常跑: LLM 逐词产出 → 喂给 TTS → TTS 合成音频 → 进
output_queue→ FastRTC 播放。 - 用户插嘴: STT 收到用户的词,发现当前是
bot_speaking,调interrupt_bot。 - 打断执行:
interrupt_bot换队列、清缓冲、remove("tts")+remove("llm")。 - 干净收场: 两个 quest 的
close(关闭 TTS/LLM 的 WebSocket)先跑完,任务再被 cancel。新回合可以开始了。
3. 核心原理
3.1 Quest:init / run / close 三段式
它要解决的小问题: 每条后台任务都要「先建连接、再干活、最后关连接」,而且要保证无论正常结束还是被取消,关连接那步都会跑。
思路: 把这三步塞进一个对象的三个字段,再让它当 async context manager 用。进入时启动 run,退出时保证跑 close。
Quest 的三段就是构造函数的三个可调用参数:
# 示意,非源码:三段式的形状
Quest(
name="tts",
init = 建立到 TTS 服务的连接,返回一个 tts 客户端, # → T
run = 拿这个客户端跑主循环(逐帧收音频), # 用 T
close = 关掉这个客户端的 WebSocket, # 清理 T
)
真实定 义在 unmute/quest_manager.py:24-46,类型是 Quest[T]——init 返回 T,run(T) 和 close(T) 都吃这个 T。三个回调的签名见 quest_manager.py:34-40(init: () -> Awaitable[T]、run: (T) -> Awaitable[None]、close: (T) -> Awaitable[None] | None)。
内部机制 —— 用一个 Future 传递 init 的产物。 Quest 里藏了个 self._data: asyncio.Future[T](quest_manager.py:46)。_run 方法先跑 init(),把结果塞进这个 future:
# 真实源码节选:unmute/quest_manager.py:65-75 Quest._run
async def _run(self):
try:
data = await self.init()
except Exception as exc:
self._data.set_exception(exc) # init 失败也记进 future
raise
else:
self._data.set_result(data) # init 成功,产物可被别人 await
await self.run(data) # 再进主循环
这个 future 是别的协程等 init 结果的地方。比如 start_up_stt 里 await quest.get()(unmute_handler.py:434)就是在等 STT 连接建好再往下走;get() 只是 await self._data(quest_manager.py:58-59)。还有个 get_nowait()(quest_manager.py:61-63),future 没完成就返回 None——handler 的 stt/tts 属性就靠它拿到「已就绪的客户端,否则 None」(unmute_handler.py:129、137)。
RAII 语义 —— 进入即启动,退出即清理。 两个魔术方法把生命周期绑到 async with:
# 真实源码节选:unmute/quest_manager.py:77-82
async def __aenter__(self) -> asyncio.Future[None]:
self.task = asyncio.create_task(self._run()) # 进入:后台起 task
return asyncio.ensure_future(self.task)
async def __aexit__(self, *exc: Any):
await self.remove() # 退出:一定跑 remove
注意 __aenter__ 返回的是那个后台 task 的 future,不是 Quest 自己——这样调用方可以对它 add_done_callback,在任务结束时收到通知(见 §3.2)。
close 的严谨处理 —— 只在 init 成功后才 close。 remove 是收尾的核心。它的讲究在于:如果 init 根本没成功,就没有资源可关,不能瞎调 close:
# 真实源码节选:unmute/quest_manager.py:84-98 Quest.remove
async def remove(self):
assert self.task is not None
try:
if self.close is not None:
try:
# 只有 init 成功了(future 完成且无异常)才 close
if self._data.done() and self._data.exception() is None:
await self.close(await self.get())
except asyncio.CancelledError:
pass
self.close = None # 关一次就置空,避免重复关
finally:
self.task.cancel() # 无论如何,最后取消 task
三个细节值得记:
self._data.done() and self._data.exception() is None(quest_manager.py:91):只在 init 确实产出了资源时才 close,避免「init 一半就死、却去关一个没建好的连接」的混乱状态。self.close = None(quest_manager.py:95):close 跑过就清空,remove被重入也不会关第二次。finally: self.task.cancel()(quest_manager.py:96-98):cancel 永远会执行——这就是 RAII 保证的落点,close 出不出错都不影响任务被取消。
from_run_step —— 无资源任务的快捷方式。 "llm" quest 没有需要建/关的连接,它只是一段要跑的逻辑。from_run_step 就是给这种情况的糖:init 返回 None,close 缺省:
# 真实源码节选:unmute/quest_manager.py:48-56 Quest.from_run_step
@staticmethod
def from_run_step(name, run) -> "Quest[None]":
async def _init() -> None:
return None
async def _run(_x: None) -> None:
await run()
return Quest(name, _init, _run) # 没有 close
_generate_response 就这么起 LLM 任务:Quest.from_run_step("llm", self._generate_response_task)(unmute_handler.py:181)。
3.2 QuestManager:按名字去重,新的顶掉旧的
它要解决的小问题: 每类任务同一时刻只该有一个。第二个 "tts" 上岗前,第一个必须先干净退场。
思路: 一个 dict[str, Quest] 按名字存活着的任务(quest_manager.py:107)。add 一个新 quest 时,如果同名的已存在,先关旧的,再登记新的。
# 真实源码节选:unmute/quest_manager.py:115-127 QuestManager.add
async def add(self, quest: Quest[T]) -> Quest[T]:
name = quest.name
try:
old = self.quests[name]
except KeyError:
pass
else:
await old.__aexit__(None) # 同名旧 quest 先谢幕(跑 close+cancel)
self.quests[name] = quest
future = await quest.__aenter__() # 新 quest 上岗
future.add_done_callback(
partial(self._one_is_done, name, self._future)
)
return quest
await old.__aexit__(None)(quest_manager.py:123)是「同名顶替」的关键一行——它等价于 await old.remove(),同步地等旧任务清理完再继续。这保证不会出现「新旧两个 TTS 连接同时活着抢队列」。
异常怎么冒泡上来。 add 给新任务的完成 future 挂了个回调 _one_is_done(quest_manager.py:126)。任务结束时:
# 真实源码节选:unmute/quest_manager.py:136-148 _one_is_done
@staticmethod
def _one_is_done(name, agg_future, future):
try:
future.result()
except asyncio.CancelledError:
pass # 被打断而取消 —— 正常,吞掉
except Exception as exc:
if not agg_future.done():
agg_future.set_exception(exc) # 真出错 —— 冒泡给聚合 future
区别对待很重要:被 cancel 是打断的正常结果,静默处理(quest_manager.py:141-142);真异常才写进那个聚合 future(_future),让 wait() 的等待方感知到(quest_manager.py:111-113)。
管理器自己也是 RAII。 QuestManager.__aexit__(quest_manager.py:156-178)在整个连接关闭时,遍历所有还活着的 quest 逐个 remove,并特意吞掉三类预期内的关闭异常——MissingServiceAtCapacity、MissingServiceTimeout、WebSocketClosedError(quest_manager.py:167-172),因为「关闭时服务已经不在了」是正常的,不该报错。
3.3 handler 上的三个具名 quest
它要解决的小问题: 把上面两个抽象接到真实的语音管线上。
stt、tts、llm 三个 quest 各有各的挂法,差别正好体现了 Quest 的灵活:
| quest | init 做什么 | run 做什么 | close 做什么 | 何时起 |
|---|---|---|---|---|
stt | find_instance("stt", SpeechToText) | _stt_loop(收转写) | stt.shutdown() | 连接建立时 start_up_stt |
tts | 带退避的 find_instance("tts", ...) | _tts_loop(收音频) | tts.shutdown() | 每回合 start_up_tts |
llm | 无(from_run_step) | _generate_response_task | 无 | 每回合 _generate_response |
start_up_stt(unmute_handler.py:422-434)在建好 quest 后还 await quest.get()——等 STT 连接真就绪才返回,注释说得明白「We want to be sure to have the STT before starting anything」(unmute_handler.py:433)。
handler 暴露的 stt/tts 属性(unmute_handler.py:123-137)不直接存客户端,而是每次从 quest_manager.quests["stt"] 取出 quest 再 get_nowait()。好处是:quest 一旦被 remove 顶替,属性自动反映最新状态,不用手动同步一个额外的引用。
4. 打断的执行:interrupt_bot 的四步
上面铺垫的一切,都是为了让 interrupt_bot(unmute_handler.py:583-607)这几行能干净地跑。这节把四步逐一拆开,重点是为什么是这个顺序。
先看前置检查:只有 bot_speaking 时打断才有意义,否则直接抛错(unmute_handler.py:584-588)。
第 ① 步:记打断标记
await self.add_chat_message_delta(INTERRUPTION_CHAR, "assistant") # :590
往 assistant 的消息里塞一个 INTERRUPTION_CHAR。这不是清理动作,而是给状态机留证据:这条回答是被打断的、没说完。状态机据此把对话推进到「轮到用户」。
第 ② 步:清 FastRTC 的播放缓冲
if self._clear_queue is not None:
self._clear_queue() # :592-595
output_queue 只是 handler 自己的队列;音频从这里被 emit 取走后,还会进 FastRTC 内部的一层播放缓冲。光清自己的队列不够——已经交给 FastRTC 的那批音频还会继续放出来。_clear_queue(FastRTC 提供)就是清这层的。源码注释:「Clear any audio queued up by FastRTC's emit()」(unmute_handler.py:593)。