夜雨聆风学习资料网

ARTICLE · 1114103

Codex 源码-Session 管理与 TurnContext

Codex 源码-Session 管理与 TurnContext

Codex 源码解析系列

第 8 讲:Session 管理与 TurnContext

基于 OpenAI Codex 源码 · 2026-10-02

💡 本讲一句话:Codex 的会话运行时要同时满足两个矛盾需求——配置随时可改(用户中途换模型、调审批策略),而正在跑的 turn 必须看到一份"冻结的世界观"。这一讲拆它的解法:Session/TurnContext/StepContext 三层快照各管一段生命周期,所有操作走单线程 submission_loop 串行分发,设置提交严格"锁内换值、锁外做副作用"。读完你会明白一个 agent 会话的骨架是怎么搭起来的。

一、Session:端点与服务的容器,不是数据仓库

进入核心运行时(第三阶段第一讲)。前几讲的登录、身份、模型 provider 都是"会话启动前的准备",而 Session 是这些准备的汇聚点:一个对话线程(thread)的完整执行体。源码里它定义在 core/src/session/session.rs,结构本身只有 40 来行——但每个字段都对应一个设计决策:

📄 codex-rs/core/src/session/session.rs (第 46-83 行,节选)

pub(crate) struct Session { // 会话本体:一个 thread 的运行时容器,全局唯一、不可克隆    pub(crate) thread_id: ThreadId, // "thread"=对话线程;session 是它的执行体,ID 贯穿事件与日志    pub(super) tx_event: Sender<Event>, // 事件出口:所有 UI 可见事件都从这条 channel 流出    pub(super) agent_status: watch::Sender<AgentStatus>, // 状态广播(PendingInit/Running…):订阅者随时取最新值    pub(super) state: Mutex<SessionState>, // 可变会话状态(历史、token 账本)——唯一需要锁的"大脑"    pub(super) thread_settings_persistence: Semaphore, // 串行化设置提交与持久化,存储 I/O 不阻塞运行时读    pub(super) features: ManagedFeatures, // 特性开关:注释明确"会话生命周期内不变"    pub(crate) active_turn: Mutex<Option<ActiveTurn>>, // 当前正在跑的 turn(至多一个)——单任务约束的落点    pub(crate) input_queue: InputQueue, // 排队中的用户输入:turn 运行中可插入(steering)    pub(crate) services: SessionServices, // 共享服务集合:exec policy、MCP、网络代理等依赖}

为什么这样设计:注意这个结构里几乎没有业务数据——没有历史消息、没有 token 计数。原因是刻意的:Session 持有的是"端点"(channel sender/receiver)和"服务句柄",它们创建后不再变化;所有会变的账本被拆进 SessionState,用一把 Mutex 圈起来。这样锁的临界区只覆盖真正可变的数据,事件发送、状态广播这些高频路径完全不碰锁。另外 thread_settings_persistence 用 Semaphore(而不是 Mutex)单独串行化"设置提交 + 持久化事件"这条链——源码注释写得很直白:让存储 I/O 不阻塞运行时状态访问。两个 Semaphore、一把 Mutex,三种同步原语各管一条互不干扰的序。

Session 只是最外层。往下还有两层快照,三者合起来才是完整的"会话世界观":

层级
生命周期
一句话职责
Session
整个会话
端点与共享服务的容器,单任务约束的持有者
TurnContext
单个 turn
冻结的世界观快照:配置、环境、权限、遥测
StepContext
单次模型请求
该次请求的设置版本 + MCP/工具路由绑定

二、spawn:一次启动要"定死"多少事

Session::spawn(core/src/session/mod.rs)是会话的出生点。入口函数本身很薄,但它的返回值类型透露了架构意图——返回 (Arc<Session>, SessionIo) 二元组:本体给内部用,I/O 端点(提交 channel + 事件 channel + 状态订阅)交给外部调用方:

📄 codex-rs/core/src/session/mod.rs (第 502-528 行)

pub(crate) fn spawn( // 对外入口:隐藏具体启动 future,调用方只拿到 BoxFuture    args: SessionSpawnArgs, // 40+ 个启动依赖打包成一个参数结构体(config、auth、models_manager…)) -> BoxFuture<'static, CodexResult<(Arc<Self>, SessionIo)>> { // 返回 (会话本体, I/O 端点) 二元组    Box::pin(async move {        let parent_trace = match args.parent_trace { // 父级 W3C trace carrier:跨进程串联分布式追踪            Some(trace) => {                if codex_otel::context_from_w3c_trace_context(&trace).is_some() { // 先校验格式合法                    Some(trace) // 合法才继承——防止脏数据污染 span 树                } else {                    warn!("ignoring invalid thread spawn trace carrier"); // 非法就降级为无父级:只告警,不失败                    None                }            }            None => None,        };        let thread_spawn_span = info_span!("thread_spawn", otel.name = "thread_spawn"); // 本次启动的根 span        if let Some(trace) = parent_trace.as_ref() {            let _ = set_parent_from_w3c_trace_context(&thread_spawn_span, trace); // 把父 trace 挂到本 span        }        Self::spawn_internal(SessionSpawnArgs { parent_trace, ..args }) // 真正初始化在 spawn_internal(40+ 步)            .instrument(thread_spawn_span) // 整个初始化过程都记在这个 span 下            .await    })}

为什么这样设计:BoxFuture<'static, _> 是刻意的"类型防火墙"——spawn_internal 的 future 捕获了 40 多个依赖,真实类型巨大且会随重构变化,用 Box 把它藏起来,调用方签名永远稳定。trace carrier 的处理则体现了 Codex 的错误哲学分层:遥测数据 fail-soft(格式不对就 warn + 丢弃),配置与权限 fail-hard(后面 spawn_internal 里会看到直接返回 Err)。追踪丢了顶多断链,权限错了是安全事故。

spawn_internal 里有一处特别值得看:base instructions(系统提示词基底)的三级优先级解析。这个决定影响整个会话的行为基线,所以放在启动时一次性定死:

📄 codex-rs/core/src/session/mod.rs (第 718-722 行)

let base_instructions = config // 三级优先级:显式配置 > 历史继承 > 模型模板    .base_instructions    .clone()    .or_else(|| conversation_history.get_base_instructions().map(|s| s.text)) // resume/fork 时沿用原会话指令,保证行为连续    .unwrap_or_else(|| model_info.get_model_instructions(config.personality)); // 都没有才用模型自带模板(按 personality 渲染)

base_instructions 三级优先级(mod.rs L718-722)

① config.base_instructions

显式配置覆盖,最高优先级——企业/用户明确指定的指令永远赢。

▼

② conversation_history.get_base_instructions()

resume/fork 时从历史里继承原会话的指令——行为连续性优先于当前配置。

▼

③ model_info.get_model_instructions(personality)

兜底:模型自带模板,按 personality 渲染——保证任何会话都有基线指令。

为什么这样设计:第②级是容易被忽略的关键:一个 resume 的会话如果突然改用"当前配置"的指令,模型行为会漂移(同样的问题昨天和今天得到不同回答)。把历史里的 base_instructions 提到显式配置之下、模型模板之上,等于承诺"恢复的会话延续原来的性格"。这个优先级链后面第 9 讲 System Prompt 组装还会再遇到——那里是逐 turn 重新渲染,这里是定死初始值。

三、SessionState:把"账本"从 Session 上拆下来

所有会随 turn 变化的数据都住在 SessionState(core/src/state/session.rs)里,被 Session 的那把 Mutex 圈住:

📄 codex-rs/core/src/state/session.rs (第 69-98 行,节选)

pub(crate) struct SessionState { // 会话级可变状态:从 Session 上拆出来的"账本区"    pub(crate) session_configuration: SessionConfiguration, // 当前生效的配置快照(可被设置更新整体替换)    pub(crate) history: ContextManager, // 对话历史 + token 信息——上下文管理的核心(第 13 讲主角)    pub(crate) latest_rate_limits: Option<RateLimitSnapshot>, // 最近一次限流快照,供 UI 展示剩余额度    previous_turn_settings: Option<PreviousTurnSettings>, // 上一个常规 turn 的设置:跨 turn 的模型/realtime 衔接依据    auto_compact_window: AutoCompactWindow, // 自动压缩窗口计数与预填状态(第 14 讲主角)    pub(crate) reasoning_effort_pin: ReasoningEffortPin, // effort 钉住机制:prewarm/采样后固定,防中途漂移    pub(crate) startup_prewarm: Option<SessionStartupPrewarmHandle>, // 启动预热句柄:首个 turn 直接复用,省冷启动    granted_permissions_by_environment_id: HashMap<String, AdditionalPermissionProfile>, // 按环境累计的已授权权限(审批通过后只增不减)}

为什么这样设计:这个结构是"会话记忆"的全部:历史、额度、压缩窗口、权限累积。注意 session_configuration 本身是 Clone 的——设置提交时先 clone 出旧值做 diff(权限画像变没变、MCP 输入变没变),再整体替换新值。"不可变快照 + 原子替换"比"逐字段可变修改"安全得多:任何持有一份旧快照的代码都不会看到半更新状态。granted_permissions_by_environment_id 的注释也值得玩味——审批通过的权限按环境累积,这是"一次授权、turn 内持续有效"的实现基础(第 19 讲展开)。

四、submission_loop:所有操作的单线程入口

Session 启动后,spawn_internal 会 spawn 一个后台任务跑 submission_loop(core/src/session/handlers.rs)。这是整个会话的"主循环"——用户输入、中断、实时语音、设置更新,全部以 Op 的形式从有界 channel 进来,在这里串行分发:

📄 codex-rs/core/src/session/handlers.rs (第 538-609 行,节选)

pub(super) async fn submission_loop( // 会话主循环:所有 Op 的唯一串行入口    sess: Arc<Session>, // 会话本体,Arc 共享给每个 handler    config: Arc<Config>, // 启动时的配置快照,供需要原始配置的分支使用    rx_sub: Receiver<Submission>, // 有界 channel;发送端全 drop → recv Err → 循环退出、会话终结) {    let mut shutdown_received = false; // Op::Shutdown 置位后跳出主循环(文件尾处理)    while let Ok(sub) = rx_sub.recv().await { // 阻塞等下一个提交——所有操作天然串行,无并发竞态        let dispatch_span = submission_dispatch_span(&sub); // 每个 op 一个 span:追踪里能看到完整分发链        let should_exit = async {            match sub.op { // 一个大 match 覆盖全部操作类型(以下节选)                Op::Interrupt => { // 用户按 Esc / UI 点停止                    interrupt(&sess).await; // 取消当前 turn:cancellation token + 清 pending waiters                    false // 中断不退出会话,继续等下一个 op                }                Op::TurnInput { request, mode, reply } => { // steering:turn 运行中插入新输入                    let result = turn_input::handle(&sess, *request, mode, sub.id.clone()).await; // 只做路由决定(入队 or 拒绝)                    let _ = reply.send(result); // oneshot 回传结果,调用方不阻塞等待执行                    false                }            };        }.instrument(dispatch_span).await;    }}

🔍 为什么用单线程 actor 而不是多线程 + 锁?

🔹 顺序即正确性:设置提交必须发生在 turn 构造之前、中断必须先于后续输入生效——单线程分发让"顺序"变成语言级保证,不用靠锁协议。🔹 channel 关闭 = 优雅退出:所有 SessionIo 的发送端被 drop 后 recv 返回 Err,循环自然结束并触发 teardown——不需要额外的 shutdown 标志位竞态。🔹 TurnInput 的 oneshot reply:调用方只等"路由决定"(入队成功/拒绝),不等执行完成——steering 输入因此是非阻塞的。

五、TurnContext:一个 turn 的冻结世界观

每个 turn 开始时,系统会把"此刻的世界"拍成一张快照——TurnContext(core/src/session/turn_context.rs)。这个 turn 里所有的模型请求、工具执行、审批判断,都基于这份快照:

📄 codex-rs/core/src/session/turn_context.rs (第 282-336 行,节选)

pub struct TurnContext { // 一个 turn 的冻结快照:turn 内所有请求共享同一份"世界观"    pub(crate) sub_id: String, // 提交 ID:事件、审批、日志都靠它对齐到这次 turn    pub config: Arc<Config>, // 本 turn 生效的配置(含 token budget 解析结果)    pub(crate) initial_settings: Arc<ResolvedStepSettings>, // 冻结的初始设置——legacy 消费者读这份    pub(super) current_settings: ArcSwap<ResolvedStepSettings>, // 可原子替换的设置快照:turn 中改模型/effort 走这里    pub(crate) environments: TurnEnvironmentSnapshot, // 本 turn 选定的执行环境(cwd、沙箱、权限)    #[deprecated(note = "use the selected turn environment cwd instead")]    pub(crate) cwd: AbsolutePathBuf, // 会话级工作目录:相对路径与沙箱策略都基于它解析    pub(crate) network: Option<NetworkProxy>, // 网络代理句柄:有则走白名单代理,无则直连/禁网    pub(crate) dynamic_tools: Vec<DynamicToolSpec>, // 动态工具定义:thread 启动时定死并持久化到 rollout    pub(crate) turn_metadata_state: Arc<TurnMetadataState>, // 遥测元数据:session/thread/sub_id、权限画像等    pub(crate) terminal_error: Arc<Mutex<Option<ErrorEvent>>>, // turn 级终态错误:首个致命错误落这里,供 UI 收尾展示}

为什么这样设计:最微妙的是 initial_settings 和 current_settings 并存。前者是 turn 开始时的冻结值(给还没迁移到新模型的 legacy 代码路径用),后者是 ArcSwap——可以原子换掉而不需要锁。这意味着:turn 跑到一半用户换了模型,current_settings 被换成新快照,但已经在飞的请求手里拿着自己捕获的旧 StepContext,完全不受影响。"读的人拿快照、写的人换指针"——无锁且一致。

构造这份快照时,make_turn_context(L790-903)要完成几个"从配置到有效值"的解析。token budget 的三级解析是典型例子:

📄 codex-rs/core/src/session/turn_context.rs (第 821-842 行,节选)

let mut per_turn_config = per_turn_config; // 拷贝一份 turn 私有配置,后续修改不影响会话级let configured_token_budget = per_turn_config.token_budget.clone(); // 先记下用户显式配置的 budget(供 sparse patch 用)let use_model_token_budget_defaults = // TokenBudget 特性开 && 用户没写任何显式设置 → 才允许吃模型默认值    per_turn_config.features.enabled(Feature::TokenBudget)        && !has_explicit_settings(&per_turn_config);per_turn_config.token_budget = resolve_token_budget( // 三级解析:显式配置 > 模型默认 > 无 budget    configured_token_budget.as_ref(),    use_model_token_budget_defaults,    model_info,);if step_settings.reasoning_effort() == Some(&ReasoningEffort::Persistent) { // Persistent effort:跨 turn 保持,需补时间提醒默认值    super::time_reminder::apply_persistent_defaults(&mut per_turn_config);}per_turn_config.service_tier = step_settings.service_tier.clone(); // service tier 以设置快照为准(已按模型能力过滤)let permission_profile = environments.permission_profile_or_else(|| { // 环境级权限优先;没选环境才回退配置里的    per_turn_config.permissions.effective_permission_profile()});

为什么这样设计:注意 use_model_token_budget_defaults 的判定条件——特性开启且用户没有任何显式设置才允许用模型默认值。这是"显式意图优先于平台默认"的一贯原则:用户写了 budget(哪怕写的是 0)就绝不覆盖,没写才让模型目录里的默认值生效。权限画像同理:environments.permission_profile_or_else(...)——turn 选定的执行环境带的权限优先于全局配置,因为"这个 turn 在哪个环境跑"比"会话默认什么"更具体。所有解析都发生在 turn 构造时一次完成,之后整个 turn 只读不查——这是快照模式的核心收益:把 N 次查询压缩成 1 次。

六、设置提交:锁内换值,锁外做副作用

用户中途改模型/审批策略时走 new_turn_with_sub_id_if(turn_context.rs L926-960)。它的协议分三步,锁的边界画得非常清楚:

📄 codex-rs/core/src/session/turn_context.rs (第 926-958 行,节选)

pub(super) async fn new_turn_with_sub_id_if( // 带准入谓词的 turn 构造:should_start 决定"要不要开这个 turn"    &self,    sub_id: String,    updates: SessionSettingsUpdate, // 本次提交携带的设置更新(模型、effort、审批策略…)    options: NewTurnContextOptions,    should_start: impl FnOnce(&SessionConfiguration, &SessionConfiguration) -> bool + Send, // (旧配置,新配置)->是否放行;必须快且无副作用) -> CodexResult<Option<(Arc<TurnContext>, ThreadSettingsSnapshot)>> {    let commit = match self.update_settings_if(updates, should_start).await { // ① 先原子提交设置(锁内完成)        Ok(Some(commit)) => commit,        Ok(None) => return Ok(None), // 谓词拒绝:不构造 turn,也不报错——静默跳过是合法路径        Err(error) => { /* …发 ErrorEvent 给 UI 并返回 InvalidRequest(节选) */ }    };    let mut configuration = commit.configuration; // ② 拿到已提交的配置快照(此时锁已释放)    if let Some(service_tier) = service_tier_for_turn { // turn 级 tier 覆盖:只改这份拷贝,不落线程设置        Arc::make_mut(&mut configuration.step_settings).service_tier = Some(service_tier);    }    let turn_context = self.new_turn_from_configuration(sub_id, configuration, options) // ③ 锁外构造 TurnContext(重活不占锁)        .await;    Ok(Some((turn_context, commit.snapshot)))}

①里的 update_settings_if(mod.rs L1760-1842)是设置提交的唯一入口,锁内/锁外的分界值得逐行看:

📄 codex-rs/core/src/session/mod.rs (第 1760-1842 行,节选)

async fn update_settings_if( // 设置提交的唯一入口:锁内校验+替换,锁外做副作用    &self,    updates: SessionSettingsUpdate, // 本次要应用的稀疏更新(模型、effort、审批策略…)    should_commit: impl FnOnce(&SessionConfiguration, &SessionConfiguration) -> bool + Send, // (旧配置,新配置)->是否放行;必须快且无副作用) -> ConstraintResult<Option<SessionSettingsCommit>> {    let notify_config_contributors = !self.services.extensions.config_contributors().is_empty(); // 有扩展订阅配置变化 → 锁外要发通知    let (commit, previous_config, new_config, permission_profile_changed, mcp_inputs_changed) = {        let mut state = self.state.lock().await; // ── 锁内段开始:临界区只做纯数据操作        let updated = match self.apply_session_settings(&state.session_configuration, &updates) { // 应用更新 + managed 约束校验            Ok(updated) => updated,            Err(err) => { warn!("rejected session settings update: {err}"); return Err(err); } // 违反企业约束 → 直接拒绝,不半提交        };        if !should_commit(&state.session_configuration, &updated) { return Ok(None); } // 准入谓词在锁内跑:看到的必是最新状态        let permission_profile_changed = /* …比较新旧权限画像(节选) */; // 变了 → 锁外要重建网络代理        state.session_configuration = updated; // 原子替换会话配置——下一把锁拿到的就是新值        let commit = SessionSettingsCommit { configuration: ..., snapshot: ... }; // 提交结果 + 持久化快照一起带出临界区    }; // ── 锁在这里释放:后面的 I/O 不再阻塞任何读状态的人    self.emit_config_changed_contributors(previous_config.as_ref(), new_config.as_ref()); // 副作用①:通知扩展配置已变    if permission_profile_changed { self.refresh_managed_network_proxy_for_current_permission_profile().await; } // 副作用②:权限变了 → 重建代理规则    if mcp_inputs_changed { self.schedule_mcp_prewarm(); } // 副作用③:MCP 输入变了 → 预热连接    Ok(Some(commit))}

为什么这样设计:这是整篇源码里"锁纪律"最严格的一段。临界区内只有:校验、diff、替换、打包 commit——全是内存操作,微秒级。而三个副作用(通知扩展、重建网络代理、MCP 预热)全部排在锁释放之后,其中两个是 .await 的网络/磁盘 I/O。如果把这些放进临界区,持锁期间任何一次网络抖动都会卡死所有读会话状态的代码——包括正在跑的 turn 的事件路径。should_commit 谓词的契约也写进了文档注释:"必须快、无副作用、不得回调 Session"——因为它在锁内执行,任何阻塞或重入都是死锁隐患。另外注意 Ok(None) 路径:谓词拒绝是合法结果而非错误(比如"已有 turn 在跑,这个更新等下一轮"),调用方静默跳过——错误和拒绝在这里被严格区分。

七、StepSettings:selected 与 resolved 的分离

设置本身又分两层,定义在 core/src/session/step_settings.rs——"用户选了什么"和"请求实际发什么"是两个类型:

📄 codex-rs/core/src/session/step_settings.rs (第 23-53 行)

pub(crate) struct StepSettings { // "选了什么":用户/客户端的原始选择,可被 sparse update 修改    pub(crate) collaboration_mode: CollaborationMode, // 模型 + reasoning effort(协作模式)    pub(crate) reasoning_summary: Option<ReasoningSummary>, // None = 跟随钉住模型的默认值    pub(crate) service_tier: Option<String>, // 归一化后的请求 tier;解析时再按模型能力过滤    pub(crate) approval_policy: Constrained<AskForApproval>, // 审批策略:Constrained 包装保证只能取 managed 允许的值}pub(crate) struct ResolvedStepSettings { // "实际发什么":选择 + 钉住的模型元数据 → 请求级有效值(不可变快照)    selected: Arc<StepSettings>, // 保留原始选择:后续 sparse patch 基于它合并,不能从 effective 值反推    pub(crate) model_info: Arc<ModelInfo>, // 钉住的模型元数据(上下文窗口、默认 effort…)    pub(crate) reasoning_summary: ReasoningSummary, // 有效值 = 配置值 or 模型默认    pub(crate) service_tier: Option<String>, // 有效 tier:按 feature + 模型支持过滤后的结果}

为什么这样设计:ResolvedStepSettings 的注释点破了关键约束:"Unset defaults and unsupported requested tiers must not be reconstructed from the effective values below"——不能从有效值反推原始选择。比如用户没设 reasoning_summary(None),有效值是模型默认的 "auto";如果下次更新时把 "auto" 当成用户的选择去合并,语义就错了(用户其实想跟随未来换的模型)。所以 selected 必须原样保留。另一个细节:Constrained<AskForApproval>>——审批策略被一个类型包装住,在类型层面保证只能取 managed requirements 允许的值,而不是靠运行时检查。

稀疏更新 apply_update(L126-153)展示了这套分离的实际收益:

📄 codex-rs/core/src/session/step_settings.rs (第 126-153 行,节选)

pub(super) async fn apply_update( // 稀疏更新:只改请求的字段,其余继承;模型没变就复用元数据(省一次目录查询)    &self,    update: &StepSettingsUpdate,    constraints: &StepSettingsConstraints<'_>,    models_manager: &dyn ModelsManager,    overrides: &ModelInfoOverrides,    personality_enabled: bool,    fast_mode_enabled: bool,) -> ConstraintResult<Self> {    let selected = self.selected.apply(update, constraints)?; // ① 合并 + 约束校验(含 auto-review 强制项)    let model_info = if selected.collaboration_mode.model() == self.selected.collaboration_mode.model() // 模型没变…        && selected.personality == self.selected.personality // …且 personality 没变    {        Arc::clone(&self.model_info) // → 直接复用旧元数据:避免重复查 models manager    } else {        Arc::new(selected.resolve_model_info(models_manager, overrides, personality_enabled).await) // 变了才重新解析 ModelInfo    };    let mut next = Self::new(Arc::new(selected), model_info, fast_mode_enabled); // ② 从新选择 + 元数据推导有效值    next.mcp_approvals_reviewer_override = update.approvals_reviewer.or(self.mcp_approvals_reviewer_override); // MCP reviewer 覆盖:新值优先,否则继承    Ok(next)}

为什么这样设计:中间那个 if/else 是性能与正确性的平衡点:模型和 personality 都没变时,Arc::clone 复用旧元数据(一次原子引用计数),省掉对 models manager 的异步查询;任何一个变了才重新解析。而 apply 内部还藏着一个安全规则——如果换到的新模型要求 auto-review(guardian),且当前不是受信 reviewer,就强制把 approvals_reviewer 改成 AutoReview(L307-319)。用户不能通过"换个模型"绕过企业的安全审查要求。

八、ActiveTurn:单任务约束与挂起等待点

Session 注释里有一句核心契约:"A session has at most 1 running task at a time"。它的落点是 ActiveTurn(core/src/state/turn.rs)——一个 Option,None 即空闲:

📄 codex-rs/core/src/state/turn.rs (第 34-109 行,节选)

pub(crate) struct ActiveTurn { // "当前正在跑的 turn"——Session 注释:一个会话至多 1 个运行任务    pub(crate) task: Option<RunningTask>, // None = 空闲;Some = 有任务在跑(Regular/Review/Compact)    pub(crate) turn_state: Arc<Mutex<TurnState>>, // turn 级可变状态:pending 审批、输入队列…}pub(crate) struct RunningTask { // 运行中任务的"档案袋":取消、等待、遥测全挂这里    pub(crate) done: Arc<Notify>, // 任务完成信号:等 turn 结束的调用方被 notify 唤醒    pub(crate) kind: TaskKind, // Regular / Review / Compact——决定收尾逻辑差异    pub(crate) task: Arc<dyn AnySessionTask>, // 具体执行体(trait object):turn loop、review、compact 都实现它    pub(crate) cancellation_token: CancellationToken, // 取消令牌:Interrupt op 从这里扇出到所有子任务    pub(crate) handle: AbortOnDropHandle<()>, // drop 即 abort——会话销毁时兜底杀掉残留任务    pub(crate) turn_context: Arc<TurnContext>, // 该 turn 的冻结快照(事件、日志靠 sub_id 对齐)}pub(crate) struct TurnState { // turn 内所有"等待外部回答"的挂起点集中登记处    pending_approvals: HashMap<String, oneshot::Sender<ReviewDecision>>, // 工具审批:key=call id,用户决定后 send 唤醒执行方    pending_user_input: HashMap<String, oneshot::Sender<AcceptedUserInputResponse>>, // request_user_input 挂起等待    pending_elicitations: HashMap<(String, RequestId), oneshot::Sender<ElicitationResponse>>, // MCP elicitation:(server, req_id) 双键定位    pub(crate) pending_input: TurnInputQueue, // steering 输入队列:turn 运行中插入、下个请求边界消费    mailbox_delivery_phase: MailboxDeliveryPhase, // CurrentTurn/NextTurn 状态机:子 agent 邮件何时并入(见源码注释)}

为什么这样设计:三个亮点。其一,AbortOnDropHandle——tokio 的"drop 即 abort"包装:即使取消逻辑有遗漏,会话对象一销毁,残留任务也会被强制杀掉,这是资源泄漏的最后防线。其二,所有挂起等待点(审批、用户输入、MCP elicitation)都登记在 TurnState 的 HashMap 里,值是 oneshot sender——中断时 clear_pending_waiters() 一把清空,所有等待方同时收到"被取消"信号,不需要逐个追踪。其三,MailboxDeliveryPhase 是个两态状态机(源码注释写得很细):turn 开始时是 CurrentTurn(子 agent 的邮件可以并入本次请求),一旦输出了用户可见的最终答案就切到 NextTurn(迟到的邮件留在队列等下一轮,避免"已经答完了又追加内容"),如果之后又有同 turn 的新工作则重新打开。这种细粒度状态机正是单任务模型能做到的——多任务并发时这类时序保证会复杂得多。

九、数据流:一次 turn 的准入路径

turn 准入路径(本讲主线)

① 提交:SessionIo.submit → tx_sub

Op 包成 Submission(带 W3C trace)进有界 channel;发送端全 drop = 会话终结。

▼

② 分发:submission_loop 串行 match

单线程 actor:Interrupt / TurnInput / UserTurn…顺序即正确性,无并发竞态。

▼

③ 提交设置:update_settings_if(锁内)

校验 managed 约束 → should_start 谓词 → 原子替换 session_configuration;副作用全部排到锁外。

▼

④ 构造快照:make_turn_context(锁外)

token budget / service tier / 权限画像一次解析定死;TurnMetadataState + skills snapshot 入 extension_data。

▼

⑤ 入位:RunningTask → active_turn

单任务约束生效;turn loop 启动(第 10 讲展开),审批/输入挂起点登记进 TurnState。

📚 系列导航

← 第 7 讲:Agent Identity / Workload Identity 与 Attestation

→ 第 9 讲:System Prompt 组装与注入

关注公众号「AI技术推荐官」获取更多源码解析内容

相关学习资料