SamplerActor 架构 本节摘要:Sampler 是一个 Actor,但它的并发模型有一个有趣的取舍:命令循环是单线程的,而实际的请求处理是每请求一个并发 task。这种混合架构不是随意的——单线程命令循环让配置更新、请求注册等「状态操作」串行化,避免共享状态的并发难题;每请求并发 task 则让多个流式请求可以同时在飞,支持会话并发的实际需求。本节会拆解这种架构,讲清 SamplerActor 的命令循环、每请求 task 的生命周期、它们如何协作,以及为什么这种取舍既安全又高效。 一、SamplerActor 在系统中的位置 先定位 SamplerActor 在整个系统里的位置: 注意一个重要事实:SamplerActor 通常不是每会话一个。
本节摘要:Sampler 是一个 Actor,但它的并发模型有一个有趣的取舍:命令循环是单线程的,而实际的请求处理是每请求一个并发 task。这种混合架构不是随意的——单线程命令循环让配置更新、请求注册等「状态操作」串行化,避免共享状态的并发难题;每请求并发 task 则让多个流式请求可以同时在飞,支持会话并发的实际需求。本节会拆解这种架构,讲清 SamplerActor 的命令循环、每请求 task 的生命周期、它们如何协作,以及为什么这种取舍既安全又高效。
先定位 SamplerActor 在整个系统里的位置:
shell 层 ├── SessionActor(每会话一个) │ └── sampler_handle → SamplerActor ← 本节主角 │ sampler 层 ├── SamplerActor(单线程命令循环) │ ├── 命令通道 cmd_rx(接收 Submit/Cancel/UpdateConfig) │ ├── 活跃请求表 active_requests(注册的进行中请求) │ └── 派生的多个 request_task(每请求一个,并发跑)
注意一个重要事实:SamplerActor 通常不是每会话一个。它的设计是「一个 sampler 服务多个会话的请求」,通过每请求 task 实现并发。这种「单 Actor 多 task」的模式,是 sampler 架构的关键。
当然,具体是「全局一个 sampler」还是「每会话一个 sampler」取决于上层如何实例化,但 Actor 本身的设计是支持多请求并发的。本节聚焦 Actor 的内部架构。
SamplerActor 的核心是一个单线程的命令循环。它的主体在 xai-grok-sampler/src/actor/mod.rs,简化伪代码:
SamplerActor::run(): loop { tokio::select! { biased; # 优先:回收已完成的 task joined = tasks.join_next(), if !tasks.is_empty() => { 清理 active_requests 里对应的请求 } # 处理命令 cmd = cmd_rx.recv() => { handle_command(cmd) } } } # 退出时:取消所有在飞请求,关闭 task 集合 cancel all active_requests tasks.shutdown()
几个关键点:
biased 优先级
biased 让 select! 按声明顺序匹配,意味着「回收已完成的 task」优先于「处理新命令」。这避免了 task 回收被新命令挤压导致资源泄漏。
单线程串行处理命令
命令通道 cmd_rx 收到的命令,由这个单线程循环逐个处理。这意味着 Submit、Cancel、UpdateConfig 等命令是串行执行的——不会有两个命令同时修改 sampler 的内部状态。
串行化的价值
为什么命令处理要串行?因为它们都涉及共享状态的修改:
如果这些操作并发执行,就要给 active_requests 表与配置加锁,引入数据竞争风险。串行化让这些操作天然安全——同一时刻只有一个执行流在改状态,不需要锁。
虽然命令处理是串行的,但实际的请求处理(网络 IO、流解析)是并发的。这是通过「每请求派生一个 task」实现的:
handle_command(Submit { request, config, completion_tx }): request_id = 生成唯一 id cancel_token = CancellationToken::new() # 注册到活跃请求表 active_requests.insert(request_id, ActiveRequest { cancel_token, started_at: now(), ... }) # 派生一个独立 task 跑实际请求 let task = tokio::spawn(async move { request_task( request_id, request, config, cancel_token, event_tx, # 事件回吐通道 ).await }) tasks.add(task)
关键观察:命令处理(注册请求)是同步完成的、串行的;而实际请求(网络 IO、流解析)在派生的 task 里异步并发地跑。这两者分离,让 sampler 既安全(状态操作串行)又高效(请求并发)。
多请求在飞
因为每请求一个 task,所以可以多个请求同时在飞:
时间点 T1: 会话 A 提交请求 R1 → 派生 task1 时间点 T2: 会话 B 提交请求 R2 → 派生 task2 时间点 T3: 会话 A 提交请求 R3 → 派生 task3 现在 task1、task2、task3 同时在飞,各自独立地与模型 API 通信。
这支持了「多个会话同时用」的实际场景。如果 sampler 是「一次只能处理一个请求」,多会话就要排队,体验很差。
派生出的 request_task 是实际干活的地方。它的简化结构(在 actor/request_task.rs):
async fn request_task(request_id, request, config, cancel_token, event_tx): # ① 发起 HTTPS 流式请求 发出 StreamStarted 事件 # 根据 config.api_backend 选择对应的 client 方法 let response_stream = match config.api_backend { Responses => client.conversation_stream_responses(request), ChatCompletions => client.conversation_stream_chat_completions(request), Messages => client.conversation_stream_messages(request), } # ② 处理流式响应 let mut stream = response_stream; let mut retry_state = RetryState::new(config.retry_policy); loop { tokio::select! { biased; # 优先:响应取消 _ = cancel_token.cancelled() => { 发出 Failed(Cancelled) return } # 处理流的一个 chunk chunk = stream.next() => { match chunk { Some(Ok(raw_chunk)) => { # L2 transform:把原始 chunk 转成统一事件 let events = transform(raw_chunk); for ev in events { event_tx.send(ev).await # 发给桥 } } Some(Err(err)) => { # 遇到错误,判断是否可重试 if retry_state.can_retry(&err) { 发出 Retrying 事件 等待 backoff 重新建立 stream continue } else { 发出 Failed(err) return } } None => { # 流结束,聚合最终响应 发出 Completed(最终响应, metrics) return } } } } }
几个要点:
响应取消优先
select! 里 cancel_token.cancelled() 优先,意味着任何时候用户取消,task 都能及时响应——中断当前请求,发出 Failed(Cancelled)。这是「协作式取消」的实现。
L2 transform 在这里发生
原始的 HTTP chunk(格式取决于后端)在这里被 L2 transform 转换成统一的 SamplingEvent。下一节详谈事件,第 04 节详谈 transform。
重试在 task 内部
遇到可重试错误时,task 自己按 RetryPolicy 决定要不要重试、等多久。重试对调用方(桥)是透明的,桥只看到 Retrying 事件(可选)或最终结果。第 06 节详谈。
流结束聚合
流正常结束时,task 把累积的 chunks 聚合成最终的 ConversationResponse,发出 Completed。
Cancel 命令的处理体现了「串行命令 + 并发 task」的协作:
handle_command(Cancel { request_id }): if let Some(active) = active_requests.get(&request_id) { active.cancel_token.cancel() # 触发取消令牌 # 注意:不在这里等待 task 结束,只是发信号 }
关键观察:
Cancel 只是发信号
Cancel 命令把对应请求的 CancellationToken 置为 cancelled,然后立即返回,不等待 task 真正结束。task 在自己的 select! 里感知到 cancelled,自行清理退出。
为什么不等待
如果 Cancel 等待 task 结束,会阻塞命令循环,后续命令(包括其他会话的 Submit)都要等。发信号不等待,让命令循环保持响应。task 的清理是异步进行的,最终会被 join_next 回收。
取消的传播链
用户按取消 → SessionActor 收 Cancel 命令 → sampler_handle 发 Cancel 命令给 SamplerActor → SamplerActor 在命令循环里处理,触发 CancellationToken → request_task 在 select! 里感知,中断 HTTP 请求 → task 发出 Failed(Cancelled),退出 → SamplerActor 在 join_next 里回收 → 桥收到 Failed,退出事件循环 → shell 循环感知,结束当前调用
整条链路是「协作式」的——每一层发信号,下一层响应,没有强杀。这保证了清理的完整性(如 HTTP 连接正确关闭、metrics 正确记录)。
UpdateConfig 命令让 sampler 的配置可以在运行中更新,典型场景是模型切换:
handle_command(UpdateConfig { new_config }): self.config = new_config # 串行地替换配置 # 注意:不影响已在飞的请求(它们用的是提交时的 config 副本)
热更新的语义:
这种「新请求新配置,老请求跑完」的语义,是热更新的安全做法——不会因为配置变化让在飞请求出错。
把 SamplerActor 的架构放在一起,它的取舍清晰可见:
串行化的部分:命令处理。Submit、Cancel、UpdateConfig 等修改共享状态的命令,由单线程循环串行执行,避免锁与数据竞争。
并发的部分:请求处理。每个 Submit 派生一个独立 task,实际网络 IO 与流解析在 task 里并发进行,支持多请求在飞。
分离的关键:状态操作(快、串行)与 IO 操作(慢、并发)分离。状态操作不阻塞 IO,IO 操作不直接改状态(只通过事件回吐)。
协作的关键:取消、关停等控制信号通过 channel 与 CancellationToken 协作传递,不强行中断。
关键概念:SamplerActor 的架构是「串行控制平面 + 并发数据平面」的经典体现。控制平面(命令、状态)串行化保证安全,数据平面(请求、响应)并发化保证性能。这种模式在服务器编程里很常见,sampler 是它在 Agent 场景的具体应用。
把 SamplerActor 与第 3 章的 SessionActor 对比,能看出 Actor 设计的多样性:
| 维度 | SessionActor | SamplerActor |
|---|---|---|
| 粒度 | 每会话一个 | 通常全局/共享一个 |
| 主要工作 | 思考-行动循环(逻辑密集) | 请求编排(IO 密集) |
| 并发模型 | 单线程循环,子组件句柄 | 单线程命令循环 + 每请求 task |
| 状态量 | 大(整个会话) | 小(配置 + 活跃请求表) |
| 生命周期 | 随会话 | 通常长驻 |
两者都是 Actor,但因职责不同,内部结构差异很大。这正说明 Actor 模型不是「一种固定结构」,而是一种「状态隔离 + 消息通信」的设计原则,可以按需变化。
下一节,我们看清 SamplerConfig 这个「一次请求带什么参数」的配置,它的字段全景与各自的作用。