数据截至 (上游 commit b21e54d6a845)
工具调用与投机预取
这一章讲什么: 前半讲「模型怎么说出一个工具调用、系统怎么把它变成结构化数据」——云端和本地是两条完全不同的路。后半讲这个项目最精巧的一段代码:在客户端还没开口要之前,就先把下一轮回答生成好,并且做到「不要的话完全没有副作用」。
1. 两条路径,一个出口
┌──────────────── 云端 API 路径 ────────────────┐
会话 tools ──► 原样传给 provider 的 tools= 参数
│
provider 直接返回结构化 function_call
│
▼
ResponseFunctionToolCall ──┐
│
┌──────────────── 本地模型路径 ──┼───────────────┐
会话 tools ──► FunctionTool.to_code_prompt() │
│ 渲染成 def name(...) 塞进系统提示 │
模型吐 出 <code>name(arg='v')</code> │
│ 正则找块 + Python AST 解析 │
▼ │
ResponseFunctionToolCall ──────────────────┤
▼
LMOutputProcessor → 协议层
response.function_call_arguments.done
两条路径最终产出同一个类型,后面所有逻辑(保序、历史、协议序列化)完全共用。
2. 本地模型路径:把 JSON Schema 变回 Python 函数签名
它要解决的小问题
开源小模型大多没有原生 function calling 训练。要让它可靠地输出工具调用,最有效的办法是给它看它最熟悉的东西:Python 函数签名。
三步走
第一步:JSON Schema → inspect.Signature。 signature_from_schema(src/speech_to_speech/LLM/tool_call/signature_from_schema.py:79-108)把参数 schema 转成真正的 Python 签名对象。类型映射由 _annotation_from_spec(:23-76)递归处理:
| Schema 构造 | 转成 |
|---|---|
const: X | Literal[X] |
enum: [a, b] | Literal[a, b] |
anyOf / oneOf | Union[...](去重后单个则不包 Union) |
allOf | 子 schema 合并后再解析 |
type: ["string", "null"] | Union[str, None] |
array + items | list[ItemType] |
默认值规则(:90-97):在 required 里且 schema 没给默认值 → 必填;schema 给了默认值 → 用它;其余 → None。
第二步:渲染进提示词。 build_tool_system_prompt(src/speech_to_speech/LLM/tool_call/tool_prompt.py:76-97)用 Jinja2 模板把每个工具的 to_code_prompt() 拼起来,并写明调用格式:
<code>function_name(required_arg='value')</code>
模板里的规则条款(tool_prompt.py:40-43)针对的都是实际踩过的坑:每个调用单独一个块、只用具名参数、字符串要加引号、可选参数宁可省略也别填 "random" / "none" / "" / null 这种占位值。
语音和纯文本各有一份模板,差别只在一条:语音版不强调「不需要开场白」,文本版明确说「直接调用,不需要前置句子」(tool_prompt.py:62)。
第三步:流式解析。 _process_printable_text(src/speech_to_speech/LLM/language_model.py:309-411)在 token 流上边收边找 <code>:
遇到 <code>:
├─ 把它之前的文本按句子冲出去(顺序不能乱)
├─ 还没等到 </code>? → 把整块留在缓冲里,等下一批 token
└─ 完整块 → extract_function_calls_from_text → AST 解析 → 校验 → 产出工具部分
纯文本模式还有一个细节(language_model.py:385-399):要逐字输出,但又不能把半个 <co 提前发出去。做法是只扣住可能构成 <code> 前缀的最长后缀,其余全部立刻发出。
AST 而不是正则
实际参数解析用的是 Python 的 ast 模块(src/speech_to_speech/LLM/tool_call/function_call.py 的 _parse_call_expr、_literal_from_ast),而不是硬撸正则。好处是嵌套的 list/dict 字面量、转义字符串都能正确解析。_split_top_level_calls(function_call.py:36)负责在一个块里切出多个顶层调用,_split_simple_calls_with_regex(:88)是它的降级路径。
解析出来的调用要过 to_realtime_function_tool_call(ctx.function_tools) 校验(language_model.py:365),不认识的工具名/参数会被 warning 掉而不是崩掉。
3. 云端路径:几乎没有代码
ResponsesApiModelHandler._build_optional_kwargs(src/speech_to_speech/LLM/responses_api_language_model.py:154-160)把会话里的 tools 原样传给 SDK。Chat Completions 那边多一层结构转换(chat_completions_language_model.py:43-95 的 _to_chat_tools / _to_chat_tool_choice),因为两个 API 的 tool 结构不同。
流式工具调用的拼装在 _iter_chat_stream_events(chat_completions_language_model.py:202-255):Chat Completions 的 tool_calls 是按 index 分片流回来的,要用一个 tool_accum: dict[int, dict[str, str]] 累积,遇到 finish 再 flush_tools() 一次性吐出完整调用。
一个协议限制:tool_choice 只支持字符串形式(auto / required / none),传对象会被拒绝并返回 tool_choice_not_supported 错误(handlers/response.py:518-523)。
4. 工具结果的往返:为什么需要两次客户端交互
项目自带的英文设计文档描述了标准流程(src/speech_to_speech/api/openai_realtime/README.md 的 “Tool result flow” 段),要点:
- 客户端执行工具,发
conversation.item.create(type: "function_call_output"); - 服务端把结果写进历史,回
conversation.item.created,但不触发生成; - 如果结果需要说出来,客户端再发一个
response.create; - 「跳个舞」这类无需口播的动作,客户端到第 2 步就可以停。
第 2 步和第 3 步分开,是为了让「静默执行」成为可能。但代价是:每次需要口播的工具都多一个客户端往返,而这个往返横在用户和下一句话之间。
下面这一节就是消除这个往返的方案。
5. 投机预取:先偷偷生成,要了才算数
直觉
绝大多数情况下,客户端拿到工具结果后一定会发 response.create。那为什么要等?
答案:工具输出一进历史就立刻开始生成,但把这个响应完全藏起来。 客户端真发 response.create 时,不是重新生成,而是把已经跑好的那个「认领」成公开响应。客户端如果做了别的(比如又塞了新消息),就把这次投机整体丢弃,像从没发生过。
工具输出到达
│
├──► maybe_start_tool_followup_prefetch()
│ 造一个带 prefetch_transaction 的 GenerateResponseRequest
│ 扔进 text_prompt_queue —— 但不开 响应、不发任何协议事件
│
▼
LLM 在后台 worker 线程里跑 ─────────────┐
│ │ 产出的每一条都被
│ │ 「输出闸门」挡在 session 私有列表里
▼ ▼
客户端发 response.create ──► claim() ──► 闸门打开,已跑好的内容瞬间涌出
或
客户端发了别的东西 ────────► discard() ─► 回滚历史 + 中断 provider 连接
难点在哪
难点不是「提前跑」,而是**「不要」时怎么擦干净**。生成过程中已经产生了三类副作用:
| 副作用 | 可撤销吗 | 处理方式 |
|---|---|---|
| 写进 Chat 的助手条目 | 可以 | 按 response_key 回滚(见 04 章) |
| 剥掉历史里的图片 | 不可以 | 推迟到认领时才执行 |
| 裁剪/摘要历史 | 不可以 | 推迟到认领时才执行 |
| 已经建立的 provider 连接 | 需要主动关 | 注册 abort 回调 |
ResponsePrefetchTransaction 的类 docstring 把这个取舍写得很清楚(src/speech_to_speech/pipeline/messages.py:229-236):生成的条目本来就是临时的可以回滚,但剥图和裁剪历史不可逆,所以未被认领的预取把这些操作「寄存」在事务里。
事务的状态机
ResponsePrefetchTransaction(pipeline/messages.py:229-331)有三个终态,一把锁保证互斥:
┌──────────── 初始:未认领、未丢弃 ────────────┐
│ │
register_abort(cb) ──► 若已丢弃立即执行 cb,否则挂起 │
complete(cleanup) ──► 若已认领立即跑 cleanup,否则寄存 │
│ │
┌────────┴────────┐ ┌───────────┴──────────┐
▼ ▼ ▼ ▼
claim() discard() claimed=True discarded=True
跑寄存的 cleanup 跑所有 abort 回调 _resolved.set() _resolved.set()
cleanup 抛异常? 丢掉寄存的 cleanup
→ 回退成 discarded
关键在 claim()(messages.py:295-317):cleanup 抛异常时,它会把 _claimed 改回 False、_discarded 置 True、set() 事件、然后重新抛出。也就是说「提交失败」被明确地变成「彻底作废」,不留中间态。
_resolved 是一个 threading.Event,给下游 handler 用 wait_until_resolved(timeout) 等待判决(messages.py:277-279)。
下游怎么等
BaseHandler.should_process_input(src/speech_to_speech/baseHandler.py:56-65)里有一段专门给 TTS 用的:
# 真实源码,baseHandler.py:56-65
if isinstance(item, TTSInput) and item.prefetch_transaction is not None:
prefetch_transaction = item.prefetch_transaction
while not self.stop_event.is_set() and not prefetch_transaction.wait_until_resolved(0.05):
pass
if not prefetch_transaction.claimed:
logger.debug("%s: dropping unclaimed prefetch input", self.__class__.__name__)
return False
注释解释了为什么闸门设在这里:LLM 可以偷偷跑,但语音合成必须等到被认领。 因为 TTS 是串行的独占资源,不能被一个可能作废的投机任务占着。
provider worker 的名额限制
预取会占一个额外的 provider 连接。为了不让投机把并发打爆,有一把 BoundedSemaphore(PREFETCH_PROVIDER_WORKER_LIMIT),常量值是 1(base_openai_compatible_language_model.py:65、202)。
抢不到名额时的处理很讲究(_iter_prefetch_events_interruptibly,:443-460):
抢不到 worker 名额
└─ 先 discard() ← 和 claim() 共用同一把事务锁,所以这个决定是原子的
├─ 事务已经被认领了? → 那它已经是公开响应了,老老实实同步跑一遍
└─ 还没被认领? → 永久不可认领,打一条 warning 放弃
注释明说了顺序理由:先 discard 让判定原子化,已认领的保持认领,还藏着的变成永久不可认领。
消费端的有界队列
预取 worker 和消费者之间是一个 Queue(maxsize=PREFETCH_STREAM_QUEUE_MAXSIZE),值为 16(:61、405)。publish()(:414-421)在队列满时循环重试并每次检查取消,而不是无限阻塞——否则取消信号进不来。
完成信号也有讲究(:471-476):worker 报完 done 后,消费者要 worker.join() 等它真的释放掉那个唯一的 provider 名额,再对外宣告完成。否则下一次预取会卡在一个时序敏感的交接上。
6. 协议层这一侧:什么时候开预取、什么时候作废
开启条件
maybe_start_tool_followup_prefetch(src/speech_to_speech/api/openai_realtime/handlers/response.py:163-193)要求:
- 当前没有别的预取在跑;
- 没有被推迟的客户端条目(
deferred_items); - 没有还在等输出的工具调用(
chat.has_pending_tool_calls()); - 存在一个「生成已逻辑完成且产出过工具调用」的源响应。
注意它用的是 st.speculative_user_turn_id —— 预取继承用户回合的身份,所以用户此时改口一样能让它作废。
作废的触发点(很多)
| 事件 | 位置 |
|---|---|
session.update(配置变了) | service.py:402-408 |
conversation.item.create(上下文变了) | service.py:489-497 |
| 新的转写完成(用户又说话了) | service.py:622 |
| 直接音频输入完成 | service.py:684 |
| VAD 打断 | handlers/audio.py:173-176 |
| 预取自身生成失败 | handlers/response.py:300-308 |
「上下文变了就作废」是这套机制正确性的基石:预取是基于某个历史快照生成的,快照一变,结果就不能用了。
认领的严格条件
_prefetch_matches(handlers/response.py:136-145)只有在 response.create 没有任何有效覆盖时才允许认领——唯一豁免的字段是 metadata,因为它不影响模型输出。任何 instructions / tools / tool_choice 覆盖都必须重新生成。
认领路径 _claim_tool_followup_prefetch(:238-288)还有一层保护:如果 claim() 抛异常(寄存的历史清理失败),它会 discard 掉并返回 None,让这次 response.create 像从没有过预取一样走正常生成。失败被完全隐藏在投机的内部。
输出闸门
is_response_output_blocked(handlers/response.py:147-155)是发送循环查询的谓词:一个 response_key 若属于「未认领的预取」或「response.created 还没发完」,它的所有输出都被挡下。挡下的旁路事件进 session.pending_text_output_items 列表(而不是原地阻塞),这样不会拖住 那个正在播音的源响应(websocket_router.py:848-854 的注释明确了这一点)。
7. 打包客户端的本地工具
talk / local 自带的麦克风客户端可以执行本地 Python 工具:--tool-module <module>。模块契约是导出 TOOLS 列表和 async 的 execute_tool(name, arguments),由 load_realtime_tool_module 加载(src/speech_to_speech/api/openai_realtime/audio_client.py,在 cli.py:147-148 和 s2s_pipeline.py:597-598 被调用)。
仓库里有一个可运行示例:examples/realtime_web_search_tool.py。
8. 代码地图
| 主题 | 文件路径 | 符号名 |
|---|---|---|
| Schema → 函数签名 | src/speech_to_speech/LLM/tool_call/signature_from_schema.py | signature_from_schema, _annotation_from_spec |
| 工具提示词模板 | src/speech_to_speech/LLM/tool_call/tool_prompt.py | build_tool_system_prompt, TOOL_PROMPT_TEMPLATE, build_block_regex |
| 工具签名渲染 | src/speech_to_speech/LLM/tool_call/function_tool.py | FunctionTool.to_code_prompt |
<code> 块解析 | src/speech_to_speech/LLM/tool_call/function_call.py | extract_function_calls_from_text, _parse_call_expr, _literal_from_ast |
| 流式提取工具块 | src/speech_to_speech/LLM/language_model.py | _process_printable_text |
| 云端 tools 传递 | src/speech_to_speech/LLM/chat_completions_language_model.py | _to_chat_tools, _iter_chat_stream_events, _tool_calls_from_accum |
| 工具调用写历史 | src/speech_to_speech/LLM/base_openai_compatible_language_model.py | _record_tool_call |
| 预取事务 | src/speech_to_speech/pipeline/messages.py | ResponsePrefetchTransaction, claim, discard, complete, register_abort |
| 预取 worker | src/speech_to_speech/LLM/base_openai_compatible_language_model.py | _iter_prefetch_events_interruptibly, _start_prefetch_worker, PREFETCH_PROVIDER_WORKER_LIMIT |
| 预取生命周期(协议侧) | src/speech_to_speech/api/openai_realtime/handlers/response.py | maybe_start_tool_followup_prefetch, discard_tool_followup_prefetch, _claim_tool_followup_prefetch, _prefetch_matches, is_response_output_blocked |
| 工具就绪旁路事件 | src/speech_to_speech/pipeline/events.py | AssistantToolCallReadyEvent |