CellActor与执行队列
每个 Code Mode cell 不是一个裸 V8 thread。CellActor 把 runtime event、observe command、yield timer、 tool callback、notification callback 和 termination 汇聚到一个 biased select loop;CellState 用同步 mutex 作为唯一终态线性化点,决定 completion 还是 termination 胜出。actor 还必须在 dropped observer 时保存 未送达 output/pending frontier,并在真正关闭路由前 tombstone 状态。
本文承接V8Runtime初始化和CodeMode架构总览,面向理解 actor mailbox、oneshot、JoinSet 和 cancellation tree 的读者。范围是 CellActor 与 SessionRuntime admission,不展开 V8 module 细节;动态 actor 测试受 rusty_v8 archive 404 阻塞,本文把源码事实与测试边界分开记录,当前不能证明 V8 运行时的动态断言。
1. Actor创建
CellActor::prepare 建立 runtime event channel、cell command mailbox 和 initial response oneshot,再调用 spawn_runtime,返回 CellHandle、initial event future 和 actor task。
源码位置:codex-rs/code-mode-runtime/src/cell_actor/mod.rs :: CellActor::prepare
let (event_tx, event_rx) = mpsc::unbounded_channel();
let (command_tx, command_rx) = mpsc::unbounded_channel();
let (initial_response_tx, initial_response_rx) = oneshot::channel();
let (runtime_tx, runtime_control_tx, runtime_terminate_handle) = spawn_runtime(
stored_values,
runtime_request(request),
event_tx,
PendingRuntimeMode::PauseUntilResumed,
task_failure_handler.clone(),
)?;
let handle = CellHandle::new(command_tx, Arc::clone(&cell_state));2. Mailbox与观察
CellHandle 的 observe 发送 CellCommand::Observe,terminate 不经过 mailbox,而是同步进入 CellState 请求 终止。actor 同时只允许一个 active observer;已有 observer 或 termination 时返回 Busy。observe 在发送 command 前先检查 accepting_observations,但最终状态仍需 actor/CellState 再判断,因为两步之间可能竞态。
源码位置:codex-rs/code-mode-runtime/src/cell_actor/types.rs :: CellHandle、CellState::request_termination
pub(crate) fn observe(&self, mode: ObserveMode) -> CellEventFuture {
if !self.state.accepting_observations() {
return closed_event();
}
let (response_tx, response_rx) = oneshot::channel();
if self.command_tx.send(CellCommand::Observe { mode, response_tx }).is_err() {
return closed_event();
}
response_event(response_rx)
}actor 收到 Observe 后先让 CellState 路由已完成事件;仍在 Running 时才安装新 observer 和 yield timer。
2.1 两种观察模式
YieldAfter(duration) 创建 timer:时间到时返回当前 content items,但 cell 继续运行。PendingFrontier 没有 timer,它要求 runtime 暂停到当前 microtask/tool frontier,并返回 content items 加 pending_tool_call_ids。如果 frontier 已经到达但上个 receiver drop,pending_frontier_ready 会保存它,下一 observer 可立即取走。
源码位置:codex-rs/code-mode-runtime/src/cell_actor/mod.rs :: observer_timer、resume_for_observation
fn observer_timer(observer: &Observer) -> Option<std::pin::Pin<Box<tokio::time::Sleep>>> {
match observer.mode {
ObserveMode::YieldAfter(duration) => Some(Box::pin(tokio::time::sleep(duration))),
ObserveMode::PendingFrontier => None,
}
}
if *runtime_paused {
let control = match mode {
ObserveMode::YieldAfter(_) => RuntimeControlCommand::Continue,
ObserveMode::PendingFrontier => RuntimeControlCommand::Resume,
};
let _ = runtime_control_tx.send(control);
} else if matches!(mode, ObserveMode::PendingFrontier) {
let _ = runtime_tx.send(RuntimeCommand::ObservePendingFrontier);
}3. Actor循环
run_cell 的 biased select 优先 cancellation,其次 command、yield timer、runtime event、notification task 和 tool task completion。termination 会取走 command receiver,让新 observation 永久 pending/closed,随后同时 发送 data/control Terminate 并终止 isolate。yield_deadline_elapsed 还会暂时禁止读取 runtime event,确保 已到期 timer 不被持续 output 饿死。
源码位置:codex-rs/code-mode-runtime/src/cell_actor/mod.rs :: run_cell
tokio::select! {
biased;
_ = cancellation_token.cancelled(), if !termination => {
termination = true;
yield_timer = None;
drop(command_rx.take());
begin_termination(
&runtime_tx,
&runtime_control_tx,
&runtime_terminate_handle,
&cancellation_token,
);
}
maybe_command = async {
match command_rx.as_mut() {
Some(command_rx) => command_rx.recv().await,
None => std::future::pending::<Option<CellCommand>>().await,
}
} => { /* observe */ }
_ = yield_timer => { /* Yielded */ }
maybe_event = event_rx.recv(), if !yield_deadline_elapsed => { /* runtime events */ }
task_result = notification_tasks.join_next(), if !notification_tasks.is_empty() => { /* report */ }
task_result = tool_tasks.join_next(), if !tool_tasks.is_empty() => { /* report */ }
}biased 保证取消不会被高频 runtime output 永久饿死。
3.1 Runtime事件
Started 启动 yield timer;Pending 标记 V8 已暂停。若当前 observer 是 PendingFrontier,actor 发送 Pending event;否则立即发 Continue,让普通 YieldAfter 观察不因内部 quiescent frontier 停住。 ContentItem 累积输出;YieldRequested 只对 YieldAfter observer 生效;ToolCall/Notify 各自 spawn callback; Result 进入 terminal commit。ThreadPanicked 只标记 failure 已上报,channel 随后关闭时再形成终态。
3.2 未送达回填
oneshot receiver 可能在 actor send 前 drop。send_cell_event 会把未送达的原 CellEvent 返回;Yielded 通过 restore_undelivered_yield 把旧 items 放回 accumulator 前部,Pending 则恢复 items、tool IDs 和 frontier flag。completion 由 CellState 保持为 Completed,下一 observation 再投递。输出因此不会因为一次 wait future 被取消而自动丢失。
源码位置:codex-rs/code-mode-runtime/src/cell_actor/mod.rs :: send_cell_event、restore_undelivered_yield
fn send_cell_event(
response_tx: oneshot::Sender<Result<CellEvent, CellError>>,
event: CellEvent,
) -> Result<(), CellEvent> {
match response_tx.send(Ok(event)) {
Ok(()) => Ok(()),
Err(Ok(event)) => Err(event),
Err(Err(error)) => panic!("cell event delivery returned an actor error: {error:?}"),
}
}4. 终态线性化
CellState.phase 是 completion/termination 的唯一竞态裁决点。completion 只有在 Running 且 cancellation 未发生时 commit session side effects;termination 把 Running 转为 Terminating 并取消 token。若 yield_control() 在 cell 创建到 initial observer 接管之间发生,pending_initial_yield_items 与最终 event 一起缓冲;YieldAfter observer 先收到这段 yield,PendingFrontier/terminal observer 则得到 prepend 后的完整 event。
源码位置:codex-rs/code-mode-runtime/src/cell_actor/types.rs :: CellState、CellPhase
enum CellPhase {
Running,
Terminating { response_tx: oneshot::Sender<Result<CellEvent, CellError>> },
Completed {
pending_initial_yield_items: Option<Vec<OutputItem>>,
event: CellEvent,
},
CompletionClaimed(CellEvent),
Tombstone,
}
pub(crate) fn commit_completion(
&self,
event: CellEvent,
pending_initial_yield_items: Option<Vec<OutputItem>>,
commit: impl FnOnce(),
) -> CompletionCommit {
let mut phase = self.phase.lock().unwrap();
if !matches!(*phase, CellPhase::Running) || self.cancellation_token.is_cancelled() {
return CompletionCommit::Rejected(event);
}
commit();
*phase = CellPhase::Completed { pending_initial_yield_items, event };
CompletionCommit::Committed
}5. Callback队列
notification 和 tool callback 使用不同 JoinSet。正常 completion 使用 DrainNotifications:先等待通知交付, 再取消 callback token 并等待 tool tasks,因此未 await 的慢工具不会无限阻塞 cell 完成。termination 使用 Cancel,在 drain notification 前就取消所有 callback。tool panic 被转换成 ToolError command 发回 runtime, 同时报告 task failure;notification panic 只报告,不存在需要 reject 的 JS Promise。
源码位置:codex-rs/code-mode-runtime/src/cell_actor/callbacks.rs :: finish_callbacks、spawn_tool、spawn_notification
pub(super) async fn finish_callbacks(
cancellation_token: &CancellationToken,
notification_tasks: &mut JoinSet<()>,
tool_tasks: &mut JoinSet<()>,
completion: CallbackCompletion,
task_failure_handler: Option<&TaskFailureHandler>,
) {
if matches!(completion, CallbackCompletion::Cancel) {
cancellation_token.cancel();
}
drain_tasks(notification_tasks, "notification", task_failure_handler).await;
cancellation_token.cancel();
drain_tasks(tool_tasks, "tool", task_failure_handler).await;
}6. Session接入
SessionRuntime 的 execute 先检查 shutdown 并分配 cell ID;start_cell 复制当前 stored values、创建 RuntimeCellHost,然后取得 cell registry 锁并再次检查 shutdown。在持锁状态下 prepare actor、插入 handle、 注册 TaskTracker,保证任何通过第二次检查的 cell 都先进入 tracker,shutdown 才能等待完整。
源码位置:codex-rs/code-mode-runtime/src/session_runtime/mod.rs :: SessionRuntime::execute、start_cell、shutdown
let stored_values = self.inner.stored_values.lock().await.clone();
let host = Arc::new(RuntimeCellHost {
cell_id: cell_id.clone(),
inner: Arc::clone(&self.inner),
});
let mut cells = self.inner.cells.lock().await;
if self.inner.shutdown_token.is_cancelled() {
return Err(Error::ShuttingDown);
}
let cell_state = Arc::new(CellState::new(self.inner.shutdown_token.child_token()));
let (handle, initial_event, task) = CellActor::prepare(/* ... */)?;
cells.insert(cell_id.clone(), handle);
let task = self.inner.cell_tasks.spawn(task);7. 关闭屏障
actor 主循环退出后不会立即 return。它先把 CellState 设为 Tombstone,使新 observe/terminate 不再进入; 再关闭 command receiver、重复发送 runtime 终止、取消并 drain 所有 callbacks,最后 await host.closed()。 CellHost::closed 的契约要求直到 session 已无法把新请求路由到该 cell 才返回,因此这是 actor task 与 session registry removal 的资源屏障。
源码位置:
codex-rs/code-mode-runtime/src/cell_actor/mod.rs::run_cell尾部清理codex-rs/code-mode-runtime/src/cell_actor/types.rs::CellHost::closed
cell_state.tombstone();
drop(command_rx.take());
begin_termination(
&runtime_tx,
&runtime_control_tx,
&runtime_terminate_handle,
&cancellation_token,
);
finish_callbacks(
&callback_cancellation_token,
&mut notification_tasks,
&mut tool_tasks,
CallbackCompletion::Cancel,
task_failure_handler.as_ref(),
).await;
host.closed().await;8. 验证边界
CellActor tests 覆盖 queued termination、dropped Yield/Pending observation、frontier 恢复、completion buffering、 initial-yield 顺序、store commit 竞态和 callback panic;SessionRuntime tests 覆盖 shutdown/admission、cell ID 耗尽、registry lock 与 drop 竞态。它们需要构建 rusty_v8,当前环境在 archive 404 阶段阻塞,未进入断言。
源码位置:codex-rs/code-mode-runtime/src/cell_actor/tests.rs、session_runtime/tests.rs
cd codex-rs
cargo test -p codex-code-mode-runtime --lib cell_actor::tests -- --test-threads=1
cargo test -p codex-code-mode-runtime --lib cell_actor::callbacks_tests -- --test-threads=1
cargo test -p codex-code-mode-runtime --lib session_runtime::tests -- --test-threads=1这些测试未执行时,源码审阅不能证明 observer drop、termination 与 completion 的真实调度结果;尤其不能 把 biased select 的静态顺序等同于所有 runtime timing 已被动态覆盖。
9. 源码排查
rg -n "CellActor::prepare|run_cell|CellCommand::Observe" codex-rs/code-mode-runtime/src/cell_actor
rg -n "CellPhase|commit_completion|deliver_completion|request_termination" codex-rs/code-mode-runtime/src/cell_actor/types.rs
rg -n "finish_callbacks|spawn_tool|spawn_notification" codex-rs/code-mode-runtime/src/cell_actor/callbacks.rs
rg -n "pending_frontier_ready|restore_undelivered_yield|tombstone|host.closed" codex-rs/code-mode-runtime/src/cell_actorCellActor 主线是:SessionRuntime admission 在 registry 锁下创建 mailbox/runtime,actor 串行协调两类 observer、timer、runtime frontier 和 callback task,CellState 线性化 completion/termination,未送达事件 被回填,最终 tombstone 并等待 host route 关闭。下一篇SessionRuntime与状态保持 将分析跨 cell store、completion commit 和 session shutdown。
