数据截至 (上游 commit b21e54d6a845)
Realtime 协议层:状态机、发送循环与传输
这一章讲什么: 前面所有章讲的是「怎么算出内容」,这一章讲「算出来的东西怎么、以及能不能发给客户端」。核心只有两个对象:
RealtimeService(状态机)和_send_loop_for(出口总闸)。
1. 线程世界与 asyncio 世界的接缝
┌──── 线程世界(6 个 handler 线程,阻塞式)────┐
│ VAD / STT / LLM / LMProc / TTS │
└────────────┬─────────────────────┬───────────┘
│ 放队列 │ 放队列
text_output_queue send_audio_chunks_queue
│ │
┌────────────▼─────────────────────▼───────────┐
│ _send_loop_for(unit) —— asyncio 协程 │
│ 每轮:get_nowait 两条队列 + await sleep(0.01)│
└────────────┬─────────────────────────────────┘
│ 调用(同步)
┌────────────▼─────────────────────────────────┐
│ RealtimeService —— 纯同步状态机,不 await │
└────────────┬─────────────────────────────────┘
│ 产出 ServerEvent 列表
┌────────────▼─────────────────────────────────┐
│ SessionTransport(WebSocket / WebRTC) │
└──────────────────────────────────────────────┘
接缝的形状是「非阻塞轮询 + 10 ms 让步」:发送循环用 get_nowait() 试两条队列,拿不到就 await asyncio.sleep(0.01)(websocket_router.py:1083)。粗糙但有效——它让整个协议层不需要任何跨线程的 async 原语。
代价是协议层的所有查询都不能阻塞。这就是为什么 SpeculativeTurnTracker 要提供一整套 try_* 非阻塞变体(见 02 章):发送循环只能用 try_is_latest_after_reopen_grace,拿到 None 就把消息放回去下一轮再试。
2. RealtimeService:一个连接的全部状态
ConnState 是什么
ConnState(src/speech_to_speech/api/openai_realtime/service.py:179-283)是一个 Pydantic 模型,装着一次连接的所有可变状态。它很大(约 50 个字段),但可以按用途分成六组:
| 组 | 代表字段 | 用途 |
|---|---|---|
| 身份 | session_id, conversation_id, runtime_config | 协议 ID 与会话配置(含 Chat) |
| 响应生命周期 | in_response, response_pending, pending_response_keys, closed_response_keys | 「现在有没有响应在跑/排队」 |
| 输出项编号 | current_item_id, content_index, next_output_index, pending_text_outputs | 协议要求的 item/output 编号 |
| 输入转写 | input_item_by_turn_revision, input_items | 把乱序到达的转写路由回正确的 item |
| 投机回合 | speculative_user_turn_id/revision, speculative_audio_duration_s | 记住最近一次用户回合 |
| 预取 | tool_followup_prefetch_request, generation_done_tool_calls | 见 05 章 |
响应键墓碑
三个方法定义了「响应键」的生命周期(service.py:259-283):
| 方法 | 语义 |
|---|---|
mark_response_pending(key) | 排队中,还没有任何输出 |
clear_pending_response(key=None) | 清一个;None 表示取消时清全部 |
close_response_key(key) | 墓碑化:之后带这个键的输出一律视为过期 |
墓碑集合 closed_response_keys 限长 128(service.py:281-282)。这是防「取消之后,已经在管道深处的输出跑出来污染下一个响应」的关键。
handler 拆分
RealtimeService 本体只做路由,实际逻辑在四个子 handler(service.py:310-313):
| 子 handler | 管什么 |
|---|---|
AudioHandler | 入站音频解码/切块、出站音频编码、speech_started/stopped |
SessionHandler | session.update 深合并、session.created/updated |
ResponseHandler | 响应生 命周期、预取、助手输出序列化(最大的一个) |
ConversationHandler | 会话条目增删、转写事件 |
流水线事件的分发表在 _pipeline_dispatch(service.py:315-325),按事件类型直接查函数。
入站音频的规格化
append_pcm(handlers/audio.py:122-149)做三件事:重采样到 16 kHz、和上一次的余数拼接、切成 1024 字节(512 采样 × 2 字节)的块。不足一块的余数存进 st.audio_remainder 等下次。
512 采样是 Silero VAD 在 16 kHz 下的固定窗口大小,所以这个数字不能随便改。
3. 发送循环:四道闸门
这是整个协议层最重要的函数(websocket_router.py:806-1091)。它的每一轮做两件事:先处理一条旁路事件,再处理一条主路输出。
旁路优先的原因
旁路队列里最重要的东西是 SpeechStartedEvent。它必须先于主路的音频被处理,否则用户开口后还会继续听到几百毫秒不该有的声音。
四道闸门(主路)
一条音频/事件要真正发出去,必须依次通过:
从 output_queue 取出一条
│
① _response_key_output_is_blocked ?
│ 未认领的预取 / response.created 还没发完
│ → 塞回 session.pending_output_item,sleep 10ms,下一轮再试
▼
② _generation_is_discardable(cancel_generation) ?
│ 世代过期,或处于丢弃窗口且不是当前世代 → 直接丢
▼
③ _response_key_is_obsolete(response_key) ?
│ 这个响应键已被墓碑化 → 丢,并做一次清理
▼
④ (投机回合)dispatch 时 try_* 返回 None ?
│ 重开候选未决 → 推迟
▼
真正发给 transport
闸门①和②的区别值得说清楚:①是「还不该发」(可能稍后就该发),②是「永远不该发」。所以①把消息存回去,②直接丢。
音频批量合并
通过闸门的音频不是一块一块发,而是攒到 MAX_AUDIO_BATCH_BYTES = 6400 字节(websocket_router.py:69,即 16 kHz 下 200 ms)再发(:1022-1053)。合并循环遇到四种情况会停下并把那一条存进 pending_output_item:
- 碰到
PIPELINE_END/ 音频终结哨兵 / 事件 /SESSION_END; - 碰到不同
response_key的音频(不能把两个响应的音频拼一起)。
批量的收益:WebSocket 帧数和 base64 编码次数减少一个数量级。
三种终结路径
音频终结哨兵 AUDIO_RESPONSE_DONE 到达时有三条分支(:944-995):
| 情况 | 动作 |
|---|---|
cleanup_only=True | 这是一个被投机作废的响应的生命周期清理:关掉对应响应或只清墓碑,不发协议事件 |
| 世代已过期 | 关响应键、response_done(gen)、恢复收听,不发 response.done |
| 正常 | 发 response.done、清 pending、清 response_playing、恢复收听 |
4. 传输抽象:WebSocket 与 WebRTC
接口
SessionTransport(src/speech_to_speech/api/openai_realtime/transports.py:29-57)只有四个方法:send_events、send_audio_chunk、discard_pending_audio、close。发送循环只认这个接口。
discard_pending_audio 是专为 WebRTC 存在的:WebSocket 发出去就没了,而 WebRTC 会在服务端缓冲未播音频,打断时必须把它冲掉。WebSocket 实现是空操作(transports.py:106-109)。
WebRTC 的差异
| 方面 | WebSocket | WebRTC |
|---|---|---|
| 握手 | ws://.../v1/realtime | POST /v1/realtime/calls(SDP offer → answer,201 + Location 头) |
| 音频上行 | input_audio_buffer.append 事件 | RTP 媒体轨(Opus 48 kHz) |
| 音频下行 | response.output_audio.delta | RTP 轨,20 ms 帧,空闲发静音 |
| JSON 事件 | 同一条 WS | oai-events data channel |
session.created | 连接时发 | data channel 打开时才发 |
input_audio_buffer.append | 支持 | 拒绝(invalid_event_for_transport) |
output_audio_buffer.clear | 不支持 | 支持 |
实现在 webrtc_session.py:PcmResampler(:71-98,有状态的重采样器)、PipelineAudioTrack(:100-154,自己节流成 20 ms 帧)、WebRTCSession(:156 起)。需要 webrtc extra(aiortc)。
ICE 服务器通过环境变量 SPEECH_TO_SPEECH_ICE_SERVERS 配置(rtc_configuration_from_env,webrtc_session.py:51),值是 JSON 列表。项目 README 提醒:对称 NAT 或没暴露 UDP 的容器环境需要自备 TURN。
5. 会话池与拒绝
路由拿到连接后调 _claim_unit(websocket_router.py:527-540)找第一个 session is None 的 unit。找不到就发 session_limit_reached 错误并断开。最大并发会话数 = --num_pipelines,没有排队。
两个运维端点:
| 端点 | 内容 |
|---|---|
GET /v1/usage | token、音频时长、响应数、错误分类计数(usage_endpoint,:589) |
GET /v1/pool | 每个 unit 的占用状态,能看出「卡住」的 unit(pool_endpoint,:613) |
/v1/pool 之所以有价值,是因为释放路径里存在「排空超时后隔离」的状态(见 01 章),需要一个观测口。
6. LLM 反向代理
它是什么
--enable_llm_proxy 会把服务器配置的那个远端 LLM,再暴露成一个普通的 OpenAI 端点:
| 后端 | 暴露的路径 |
|---|---|
chat-completions | POST /v1/chat/completions |
responses-api | POST /v1/responses |
(src/speech_to_speech/api/openai_realtime/llm_proxy.py:28-31)
为什么要有它
客户端常常需要做「副业任务」:给对话生成标题、做摘要、跑后台 agent。这些不该跟语音抢那条流水线,也不该让客户端自己拿一份 API key。模块 docstring 说明:代理请求完全不碰流水线的队列和取消域,所以和语音对话完全并发,永远不会被新的说话打断。
安全边界(必须读)
模块 docstring 和 README 都明确写了:服务器自身不做任何认证和限流。它假定运行在可信网络,或者前面有一个网关负责访问控制。上游的真 API key 由服务器持有,永不下发给客户端;请求里的 model 字段一律被覆写成服务器配置的 --model_name。
后端不支持代理时返回 501(s2s_pipeline.py:324-329 在启动时也会直接拒绝这种组合)。
用量统计
LLMProxyUsage(llm_proxy.py:43-104)单独计数,且 429 独立成桶不计入 4xx——注释说理由是「让被打爆的客户端一眼可见」。record_token_payload 兼容三种 usage 形状(chat 的 prompt_tokens、responses 的 input_tokens、responses 流式的 response.usage),record_sse_event 负责从 SSE 流里逐事件抠出来。
7. 支持的协议事件(速查)
以下依据项目自带的英文事件表(src/speech_to_speech/api/openai_realtime/README.md)与 service.py:86-129 的类型映射。
客户端 → 服务端:
| 事件 | 作用 |
|---|---|
input_audio_buffer.append | 送 base64 PCM |
input_audio_buffer.commit | 提交缓冲(空缓冲会报错) |
output_audio_buffer.clear | 清未播音频(仅 WebRTC) |
session.update | 深合并会话配置 |
conversation.item.create | 注入文本或工具输出,不触发生成 |
response.create | 触发生成,可带每响应覆盖 |
response.cancel | 取消当前/排队响应 |
服务端 → 客户端(节选):session.created/updated、error、input_audio_buffer.speech_started/stopped、conversation.item.input_audio_transcription.delta/completed、response.created、response.output_audio.delta/done、response.output_audio_transcript.delta/done、response.function_call_arguments.done、response.done。
一条兼容性提示(项目 README 的 “Transcript event compatibility” 段):助手字幕现在按 delta 流式发,done 只做终结;把每个块级 done 当增量渲染的老客户端需要改成消费 delta。
8. 代码地图
| 主题 | 文件路径 | 符号名 |
|---|---|---|
| 协议状态机 | src/speech_to_speech/api/openai_realtime/service.py | RealtimeService, ConnState, _pipeline_dispatch, parse_client_event |
| STT→LLM 桥接 | 同上 | _on_transcription_completed, _on_audio_input_completed |
| 用量与错误 | 同上 | UsageMetrics, GlobalUsageMetrics, _on_token_usage |
| 发送循环(四道闸门) | src/speech_to_speech/api/openai_realtime/websocket_router.py | _send_loop_for, _response_key_output_is_blocked, _generation_is_discardable, _response_key_is_obsolete |
| 队列清理与会话释放 | 同上 | _flush_queue, _clean_unit, _release_unit_after_drain, SESSION_END_QUARANTINE_TIMEOUT_S |
| 应用与路由 | 同上 | create_app, realtime_endpoint, webrtc_calls_endpoint, usage_endpoint, pool_endpoint |
| 音频编解码/切块 | src/speech_to_speech/api/openai_realtime/handlers/audio.py | append_pcm, encode_audio_chunk, on_speech_started, on_speech_stopped |
| 响应生命周期 | src/speech_to_speech/api/openai_realtime/handlers/response.py | handle_response_create, finish_response, on_assistant_output, _build_response |
| 会话配置合并 | src/speech_to_speech/api/openai_realtime/runtime_config.py | RuntimeConfig, _apply_update, interrupt_response_enabled |
| 传输抽象 | src/speech_to_speech/api/openai_realtime/transports.py | SessionTransport, WebSocketTransport |
| WebRTC | src/speech_to_speech/api/openai_realtime/webrtc_session.py | WebRTCSession, PipelineAudioTrack, PcmResampler, rtc_configuration_from_env |
| LLM 代理 | src/speech_to_speech/api/openai_realtime/llm_proxy.py | LLMProxyConfig, LLMProxyUsage |
| 官方事件表 | src/speech_to_speech/api/openai_realtime/README.md | — |