跳到主要内容

数据截至 (上游 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_calls05 章

响应键墓碑

三个方法定义了「响应键」的生命周期(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
SessionHandlersession.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_eventssend_audio_chunkdiscard_pending_audioclose。发送循环只认这个接口。

discard_pending_audio专为 WebRTC 存在的:WebSocket 发出去就没了,而 WebRTC 会在服务端缓冲未播音频,打断时必须把它冲掉。WebSocket 实现是空操作(transports.py:106-109)。

WebRTC 的差异

方面WebSocketWebRTC
握手ws://.../v1/realtimePOST /v1/realtime/calls(SDP offer → answer,201 + Location 头)
音频上行input_audio_buffer.append 事件RTP 媒体轨(Opus 48 kHz)
音频下行response.output_audio.deltaRTP 轨,20 ms 帧,空闲发静音
JSON 事件同一条 WSoai-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/usagetoken、音频时长、响应数、错误分类计数(usage_endpoint,:589)
GET /v1/pool每个 unit 的占用状态,能看出「卡住」的 unit(pool_endpoint,:613)

/v1/pool 之所以有价值,是因为释放路径里存在「排空超时后隔离」的状态(见 01 章),需要一个观测口。


6. LLM 反向代理

它是什么

--enable_llm_proxy 会把服务器配置的那个远端 LLM,再暴露成一个普通的 OpenAI 端点:

后端暴露的路径
chat-completionsPOST /v1/chat/completions
responses-apiPOST /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/updatederrorinput_audio_buffer.speech_started/stoppedconversation.item.input_audio_transcription.delta/completedresponse.createdresponse.output_audio.delta/doneresponse.output_audio_transcript.delta/doneresponse.function_call_arguments.doneresponse.done

一条兼容性提示(项目 README 的 “Transcript event compatibility” 段):助手字幕现在按 delta 流式发,done 只做终结;把每个块级 done 当增量渲染的老客户端需要改成消费 delta


8. 代码地图

主题文件路径符号名
协议状态机src/speech_to_speech/api/openai_realtime/service.pyRealtimeService, 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.pyappend_pcm, encode_audio_chunk, on_speech_started, on_speech_stopped
响应生命周期src/speech_to_speech/api/openai_realtime/handlers/response.pyhandle_response_create, finish_response, on_assistant_output, _build_response
会话配置合并src/speech_to_speech/api/openai_realtime/runtime_config.pyRuntimeConfig, _apply_update, interrupt_response_enabled
传输抽象src/speech_to_speech/api/openai_realtime/transports.pySessionTransport, WebSocketTransport
WebRTCsrc/speech_to_speech/api/openai_realtime/webrtc_session.pyWebRTCSession, PipelineAudioTrack, PcmResampler, rtc_configuration_from_env
LLM 代理src/speech_to_speech/api/openai_realtime/llm_proxy.pyLLMProxyConfig, LLMProxyUsage
官方事件表src/speech_to_speech/api/openai_realtime/README.md