Task抽象与生命周期
Codex Core 中的 Task 不是一次模型请求。RegularTask 可以包含多轮 sampling 与工具调用,CompactTask 运行一次 压缩工作流,ReviewTask 管理子会话,UserShellCommandTask 执行独立命令。它们共同实现 SessionTask,由 Session 注册为一个 RunningTask,并共享终态、取消、持久化屏障和指标收尾。
本文接续 Turn主循环与退出条件:前文解释 RegularTask 内部为什么继续或 返回,本文解释返回值由谁接收、active slot 如何释放、外部 Interrupt/Replaced 如何抢占 owner。还建议先了解 Session核心数据结构,避免把 ActiveTurn、TurnState 与 TurnContext 当成同一个对象。
本文不展开四种具体 Task 的业务算法,也不重复 Session 整体关闭顺序。目标是让读者能从 spawn_task() 追踪到 Tokio task、RunningTask、on_task_finished() 或 handle_task_abort(),并判断一个 terminal event 究竟由正常完成 owner 还是 abort owner 发出。
1. Task层级
TurnContext 固定本 Turn 的配置;SessionTask 定义工作流;RunningTask 是 Session 保存的运行句柄; TurnState 保存审批、权限请求、pending input 等可变状态。ActiveTurn 只是把“运行句柄”和“Turn 内状态” 放进同一 active slot,它本身不是后台任务。
源码位置:codex-rs/core/src/state/turn.rs,符号 ActiveTurn、TaskKind 与 RunningTask。
/// Metadata about the currently running turn.
pub(crate) struct ActiveTurn {
// None既可能表示尚未安装task的reservation,也可能表示task已被finish owner取走。
pub(crate) task: Option<RunningTask>,
// TurnState独立放在Arc中,使输入队列与收尾逻辑能在RunningTask被取走后继续处理。
pub(crate) turn_state: Arc<Mutex<TurnState>>,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) enum TaskKind {
Regular,
Review,
Compact,
}
pub(crate) struct RunningTask {
// done用于abort grace等待;它不是Task完成结果的传输channel。
pub(crate) done: Arc<Notify>,
pub(crate) kind: TaskKind,
// 类型擦除后的实现保留到abort hook完成,不能只保存JoinHandle。
pub(crate) task: Arc<dyn AnySessionTask>,
pub(crate) cancellation_token: CancellationToken,
// RunningTask被意外drop时AbortOnDropHandle仍会终止Tokio task。
pub(crate) handle: AbortOnDropHandle<()>,
pub(crate) turn_context: Arc<TurnContext>,
pub(crate) _agent_execution_guard: Option<AgentExecutionGuard>,
// Timer recorded when the task drops to capture the full turn duration.
pub(crate) _timer: Option<codex_otel::Timer>,
}TaskKind 只有三类,但 SessionTask 有四个当前实现:UserShellCommandTask 的 kind 是 Regular。因此 kind 是调度/遥测分类,不是 Rust 具体类型的完整判别枚举;调试时还要结合 span_name()。
2. SessionTask
SessionTask::run 使用 RPITIT 风格返回 future,不适合直接作为 trait object 方法调用;而 RunningTask 又必须 在同一个字段中保存不同实现。AnySessionTask 将 future 装箱为 BoxFuture,完成类型擦除,同时保留可动态 调用的 abort()。
源码位置:codex-rs/core/src/tasks/mod.rs,符号 SessionTask、AnySessionTask 及 blanket impl。
pub(crate) type SessionTaskResult = CodexResult<Option<String>>;
pub(crate) trait SessionTask: Send + Sync + 'static {
fn kind(&self) -> TaskKind;
fn span_name(&self) -> &'static str;
fn run(
self: Arc<Self>,
session: Arc<Session>,
ctx: Arc<TurnContext>,
input: Vec<TurnInput>,
cancellation_token: CancellationToken,
) -> impl std::future::Future<Output = SessionTaskResult> + Send;
fn abort(
&self,
session: Arc<Session>,
ctx: Arc<TurnContext>,
) -> impl std::future::Future<Output = ()> + Send {
// 默认abort为空;只有需要补充业务清理的Task覆盖它。
async move {
let _ = (session, ctx);
}
}
}
pub(crate) trait AnySessionTask: Send + Sync + 'static {
fn kind(&self) -> TaskKind;
fn span_name(&self) -> &'static str;
fn run(
self: Arc<Self>,
session: Arc<Session>,
ctx: Arc<TurnContext>,
input: Vec<TurnInput>,
cancellation_token: CancellationToken,
) -> BoxFuture<'static, SessionTaskResult>;
fn abort<'a>(
&'a self,
session: Arc<Session>,
ctx: Arc<TurnContext>,
) -> BoxFuture<'a, ()>;
}
impl<T> AnySessionTask for T
where
T: SessionTask,
{
fn kind(&self) -> TaskKind {
SessionTask::kind(self)
}
fn span_name(&self) -> &'static str {
SessionTask::span_name(self)
}
fn run(
self: Arc<Self>,
session: Arc<Session>,
ctx: Arc<TurnContext>,
input: Vec<TurnInput>,
cancellation_token: CancellationToken,
) -> BoxFuture<'static, SessionTaskResult> {
// blanket impl只做future装箱,不改变具体Task的返回和取消语义。
Box::pin(SessionTask::run(
self,
session,
ctx,
input,
cancellation_token,
))
}
fn abort<'a>(
&'a self,
session: Arc<Session>,
ctx: Arc<TurnContext>,
) -> BoxFuture<'a, ()> {
Box::pin(SessionTask::abort(self, session, ctx))
}
}SessionTaskResult 的三种语义是:Ok(Some(message)) 正常完成并带最终文本;Ok(None) 正常完成但没有最终 文本;Err(CodexErr::TurnAborted) 请求 aborted lifecycle。其他 Err 会在统一收尾中转换成 ErrorEvent, 随后仍使用 TurnComplete envelope。
3. 实现与TaskKind
| 实现 | TaskKind | span_name | run()核心职责 | 自定义abort() |
|---|---|---|---|---|
| RegularTask | Regular | session_task.turn | 多轮模型与工具主循环 | 否 |
| CompactTask | Compact | session_task.compact | local/remote/token-budget压缩 | 否 |
| ReviewTask | Review | session_task.review | 启动review子会话并转发结果 | 是,退出review mode |
| UserShellCommandTask | Regular | session_task.user_shell | standalone shell turn | 否 |
源码位置:codex-rs/core/src/tasks/review.rs,符号 ReviewTask 的 run 与 abort。
impl SessionTask for ReviewTask {
fn kind(&self) -> TaskKind {
TaskKind::Review
}
fn span_name(&self) -> &'static str {
"session_task.review"
}
async fn run(
self: Arc<Self>,
session: Arc<Session>,
ctx: Arc<TurnContext>,
input: Vec<TurnInput>,
cancellation_token: CancellationToken,
) -> SessionTaskResult {
// ...收集UserInput并启动review子会话
let output = match start_review_conversation(
session.clone(),
ctx.clone(),
user_input,
cancellation_token.clone(),
)
.await
{
Some(receiver) => {
process_review_events(session.clone(), ctx.clone(), receiver).await
}
None => None,
};
if !cancellation_token.is_cancelled() {
// 正常路径退出review mode;取消路径留给abort hook做同一清理,避免重复事件。
exit_review_mode(Arc::clone(&session), output.clone(), ctx.clone()).await;
}
Ok(None)
}
async fn abort(&self, session: Arc<Session>, ctx: Arc<TurnContext>) {
// review是当前唯一覆盖abort的实现,保证被替换/中断时仍发退出review状态。
exit_review_mode(session, /*review_output*/ None, ctx).await;
}
}CompactTask 接收 cancellation token 却不直接读取它;它调用的压缩实现负责传播 TurnAborted,超出 grace 时 外层仍会 abort Tokio handle。UserShellCommandTask 则把 token 传入执行函数。这个差异说明 trait 只规定取消 能力与最终强制手段,不要求每个实现用同一种轮询写法。
4. spawn_task
spawn_task() 是普通替换入口:先用 Replaced 中止旧 Task,清除 connector selection,再调用 start_task()。旧 Task 的 terminal event 和清理由 abort owner 完成后,新 Task 才开始安装。
当前 start_task() 还会在安装 RunningTask 前把 pending input 合并进 TurnState,记录 turn-start token baseline,并创建带 turn tracing span 的 Tokio task。RunningTask 中的取消 token、AbortOnDropHandle、计时器 和 diagnostics guard 因而覆盖完整 Turn 生命周期,而不只是模型 stream 的局部执行。
源码位置:codex-rs/core/src/tasks/mod.rs,符号 Session::spawn_task。
pub async fn spawn_task<T: SessionTask>(
self: &Arc<Self>,
turn_context: Arc<TurnContext>,
input: Vec<TurnInput>,
task: T,
) {
// 新任务不会与旧RunningTask并存;Replaced与用户Interrupt是不同abort reason。
self.abort_all_tasks(TurnAbortReason::Replaced).await;
self.clear_connector_selection().await;
self.start_task(
turn_context,
input,
task,
MailboxParentProvenance::Ignore,
)
.await;
}合成 mailbox Turn 不调用这个固定 provenance,而是用 MailboxParentProvenance::Attribute 将触发消息的父 Turn 写入 metadata。Task lifecycle 因而也承担 lineage 安装点,但不会把 mailbox 内容复制进 TurnContext 字段。
5. start_task事务
start_task() 先创建类型擦除对象、独立 cancellation token 和 done Notify,然后取得或建立 ActiveTurn 的 TurnState,将 Session 级 pending input 转移进去并发送 extension turn-start lifecycle。第二次取得 active lock 后,才 spawn Tokio future 并把 RunningTask 写入 slot。
源码位置:codex-rs/core/src/tasks/mod.rs,符号 Session::start_task 的准备阶段。
let task: Arc<dyn AnySessionTask> = Arc::new(task);
let task_kind = task.kind();
let span_name = task.span_name();
let started_at = Instant::now();
let turn_started_at_unix_ms = turn_context
.turn_timing_state
.mark_turn_started(started_at)
.await;
turn_context
.turn_metadata_state
.set_turn_started_at_unix_ms(turn_started_at_unix_ms);
let token_usage_at_turn_start = self.total_token_usage().await.unwrap_or_default();
// root token归RunningTask所有;传给实现的是child,外部abort只需cancel root。
let cancellation_token = CancellationToken::new();
let done = Arc::new(Notify::new());
let (pending_items, parent_turn_id) =
self.input_queue.get_pending_input(&self.active_turn).await;
if let (MailboxParentProvenance::Attribute, Some(id)) =
(mailbox_parent_provenance, parent_turn_id)
{
turn_context.turn_metadata_state.set_parent_turn_id(id);
}
let turn_state = {
let mut active = self.active_turn.lock().await;
let turn = active.get_or_insert_with(ActiveTurn::default);
// ActiveTurn可先作为reservation存在,但同一slot不能已有RunningTask。
debug_assert!(turn.task.is_none());
Arc::clone(&turn.turn_state)
};
turn_state.lock().await.token_usage_at_turn_start = token_usage_at_turn_start.clone();
self.input_queue
.extend_pending_input_for_turn_state(turn_state.as_ref(), pending_items)
.await;
self.emit_turn_start_lifecycle(turn_context.as_ref(), &token_usage_at_turn_start)
.await;这里的 turn-start lifecycle 是 extension 回调,不等同于协议 TurnStarted。RegularTask 和 standalone UserShell 在各自 run() 中发送协议事件;Compact 有意不在压缩前注入普通 Turn 上下文。两层事件不能用同一 名字推断顺序。
6. Tokio闭包
spawned future 先运行具体 Task,随后 flush rollout;只有 root token 未取消时才进入 on_task_finished()。 无论正常、错误还是已被取消,闭包最后都 notify_waiters(),供 abort owner 的 grace wait 观察。
源码位置:codex-rs/core/src/tasks/mod.rs,符号 Session::start_task 的 spawn 与注册阶段。
let done_clone = Arc::clone(&done);
let session = Arc::clone(self);
let ctx = Arc::clone(&turn_context);
let task_for_run = Arc::clone(&task);
let task_input = input;
let task_cancellation_token = cancellation_token.child_token();
let handle = tokio::spawn(
async move {
let ctx_for_finish = Arc::clone(&ctx);
let task_result = task_for_run
.run(
Arc::clone(&session),
ctx,
task_input,
// 再派生child,具体Task不能取消RunningTask root token。
task_cancellation_token.child_token(),
)
.instrument(trace_span!("session_task.run"))
.await;
let sess = Arc::clone(&session);
// 普通items先flush,terminal event稍后还会有第二个flush barrier。
if let Err(err) = sess.flush_rollout().await {
warn!("failed to flush rollout before completing turn: {err}");
sess.send_event(
ctx_for_finish.as_ref(),
EventMsg::Warning(WarningEvent {
message: format!(
"Failed to save the conversation transcript; Codex will continue retrying. Error: {err}"
),
}),
)
.await;
}
if !task_cancellation_token.is_cancelled() {
// 外部abort取消root后由abort owner发终态;此处必须跳过,防止双terminal event。
sess.on_task_finished(Arc::clone(&ctx_for_finish), task_result)
.await;
}
done_clone.notify_waiters();
}
.instrument(task_span),
);
let timer = turn_context
.session_telemetry
.start_timer(TURN_E2E_DURATION_METRIC, &[])
.ok();
let running_task = RunningTask {
done,
handle: AbortOnDropHandle::new(handle),
kind: task_kind,
task,
cancellation_token,
turn_context: Arc::clone(&turn_context),
_agent_execution_guard: agent_execution_guard,
_timer: timer,
};
// RunningTask写入active slot后,Session才拥有取消与终止该future的完整句柄集合。
turn.task = Some(running_task);第一道屏障是 task body 后的 rollout flush,保证普通历史尽量先持久化;第二道是 terminal event 发送后的 flush,保证客户端收到终态后重读 rollout 能看到对应事件。done 不是持久化屏障,它只表示 spawn closure 已经走到最后。
7. 正常完成 Owner
on_task_finished() 先把结果归一化,再从 active_turn.task 中 take() RunningTask 并 detach handle。若 task 已经被外部 abort owner 取走,这里得到 None 并立即返回,不再发送任何终态。
源码位置:codex-rs/core/src/tasks/mod.rs,符号 Session::on_task_finished 的 owner 获取。
let (last_agent_message, abort_reason) = match task_result {
Ok(last_agent_message) => (last_agent_message, None),
Err(err) if matches!(err.details(), CodexErrorDetails::TurnAborted) => {
// Task自发返回TurnAborted时,在正常finish owner中归约为Interrupted。
(None, Some(TurnAbortReason::Interrupted))
}
Err(err) => {
warn!(%err, "session task returned an unexpected error");
self.emit_turn_error_lifecycle(
turn_context.as_ref(),
err.to_codex_protocol_error(),
)
.await;
self.track_turn_codex_error(turn_context.as_ref(), &err);
self.send_event(
turn_context.as_ref(),
EventMsg::Error(err.to_error_event(/*message_prefix*/ None)),
)
.await;
// 普通错误没有abort reason;ErrorEvent进入terminal_error后仍生成TurnComplete。
(None, None)
}
};
turn_context
.turn_metadata_state
.cancel_git_enrichment_task();
let turn_state = {
let mut active = self.active_turn.lock().await;
active.as_mut().and_then(|active_turn| {
// take成功者成为finish owner;abort路径若先take,这里返回None。
let task = active_turn.task.take()?;
task.handle.detach();
Some(Arc::clone(&active_turn.turn_state))
})
};
let Some(turn_state) = turn_state else {
return;
};detach 很重要:正常完成后 RunningTask 即将 drop,如果不 detach,AbortOnDropHandle 会对已经完成或正在收尾的 Tokio task 再调用 abort。owner 获取在 active lock 内完成,但持久化、metrics 和 extension hooks 都在锁外, 避免长时间占用 Session 的 active slot mutex。
8. Active slot
正常 owner 会取出 TurnState 中尚未消费的 input,经 hook 检查后写入 history;然后生成 metrics 与 terminal event。最后只有 active slot 仍指向同一 turn_state 且 task=None 时才清空 ActiveTurn。
源码位置:codex-rs/core/src/tasks/mod.rs,符号 Session::on_task_finished 的终态与清理阶段。
let pending_input = self
.input_queue
.take_pending_input_for_turn_state(turn_state.as_ref())
.await;
if !pending_input.is_empty() {
for pending_input_item in pending_input {
let hook_outcome =
inspect_pending_input(self, &turn_context, &pending_input_item).await;
if hook_outcome.should_stop {
record_additional_contexts(
self,
&turn_context,
hook_outcome.additional_contexts,
)
.await;
} else {
// Task返回边界遗留的输入仍进入history,不会因RunningTask结束而静默丢失。
record_pending_input(
self,
&turn_context,
pending_input_item,
hook_outcome.additional_contexts,
)
.await;
}
}
}
// ...计算turn token/tool/memory metrics与timing
let event = if let Some(reason) = abort_reason {
self.emit_turn_abort_lifecycle(
reason.clone(),
turn_context.extension_data.as_ref(),
)
.await;
EventMsg::TurnAborted(TurnAbortedEvent {
turn_id: Some(turn_context.sub_id.clone()),
reason,
started_at,
completed_at,
duration_ms,
})
} else {
let time_to_first_token_ms = turn_context
.turn_timing_state
.time_to_first_token_ms()
.await;
let error = turn_context.terminal_error.lock().await.clone();
self.emit_turn_stop_lifecycle(turn_context.extension_data.as_ref())
.await;
EventMsg::TurnComplete(TurnCompleteEvent {
turn_id: turn_context.sub_id.clone(),
last_agent_message,
error,
started_at,
completed_at,
duration_ms,
time_to_first_token_ms,
})
};
self.send_event(turn_context.as_ref(), event).await;
let cleared_active_turn = {
let mut active = self.active_turn.lock().await;
if let Some(active_turn) = active.as_ref()
&& active_turn.task.is_none()
// 指针身份防止旧finish逻辑清掉后来安装、但恰好也task=None的新ActiveTurn。
&& Arc::ptr_eq(&active_turn.turn_state, &turn_state)
{
*active = None;
true
} else {
false
}
};
if cleared_active_turn {
self.emit_thread_idle_lifecycle_if_idle().await;
}
if let Err(err) = self.flush_rollout().await {
warn!("failed to flush rollout after emitting terminal turn event: {err}");
}
if cleared_active_turn {
// 只有本owner真正清空slot后,trigger mailbox才有资格启动下一Task。
self.maybe_start_turn_for_pending_work().await;
}Arc::ptr_eq 比较的是 TurnState owner,而不是 Turn ID 字符串。即使未来存在 ID 重用或新 reservation,旧收尾也 不能仅凭相等字符串删除新 slot。
9. 外部 Abort Owner
abort_all_tasks() 先把整个 ActiveTurn 从 mutex 中 take(),再从中取 RunningTask。这样 spawn closure 即使 同时结束,on_task_finished() 也无法再次取得 task owner。随后先 cancel token,最多等待 100ms done;无论 是否协作完成,都 abort JoinHandle,再调用具体 Task 的 abort hook。
源码位置:codex-rs/core/src/tasks/mod.rs,符号 Session::handle_task_abort。
async fn handle_task_abort(
self: &Arc<Self>,
task: RunningTask,
reason: TurnAbortReason,
) {
let sub_id = task.turn_context.sub_id.clone();
if task.cancellation_token.is_cancelled() {
// root已取消表示已有abort owner,不重复发终态。
return;
}
task.cancellation_token.cancel();
task.turn_context
.turn_metadata_state
.cancel_git_enrichment_task();
let session_task = task.task;
select! {
_ = task.done.notified() => {},
_ = tokio::time::sleep(
Duration::from_millis(GRACEFULL_INTERRUPTION_TIMEOUT_MS)
) => {
warn!(
"task {sub_id} didn't complete gracefully after {}ms",
GRACEFULL_INTERRUPTION_TIMEOUT_MS
);
}
}
// 即使body已协作返回也abort handle,封住closure尚未结束的flush/finish窗口。
task.handle.abort();
session_task
.abort(Arc::clone(self), Arc::clone(&task.turn_context))
.await;
if reason == TurnAbortReason::Interrupted
&& let Some(marker) = interrupted_turn_history_marker(
InterruptedTurnHistoryMarker::from_config_and_version(
task.turn_context.config.as_ref(),
task.turn_context.multi_agent_version,
),
)
{
self.record_conversation_items(
task.turn_context.as_ref(),
std::slice::from_ref(&marker),
)
.await;
// marker先flush,客户端收到TurnAborted后同步重读即可看到中断事实。
if let Err(err) = self.flush_rollout().await {
warn!(
"failed to flush interrupted-turn marker before emitting TurnAborted: {err}"
);
}
}
let started_at = task
.turn_context
.turn_timing_state
.started_at_unix_secs()
.await;
let (completed_at, duration_ms, profile) = task
.turn_context
.turn_timing_state
.complete_profile_and_duration_ms()
.await;
self.services
.analytics_events_client
.track_turn_profile(TurnProfileFact {
turn_id: task.turn_context.sub_id.clone(),
profile,
});
let event = EventMsg::TurnAborted(TurnAbortedEvent {
turn_id: Some(task.turn_context.sub_id.clone()),
reason,
started_at,
completed_at,
duration_ms,
});
self.send_event(task.turn_context.as_ref(), event).await;
// ...清guardian circuit breaker
if let Err(err) = self.flush_rollout().await {
warn!("failed to flush rollout after emitting terminal turn event: {err}");
}
}Interrupted 才写模型可见 marker;Replaced 不写,因为替换不是用户要求模型理解的中断上下文。两者都发送 TurnAborted,只是 reason 不同。abort_all_tasks() 在 handler 返回后才清 pending waiters,避免正在等待的 审批先被 channel drop 转成模型可见 rejection,抢在 TurnAborted 之前。
10. 生命周期状态
这张状态图是多个 owner 字段与返回路径的组合模型,不是源码中的 enum。真正的不变量是:同一时刻只有一个 路径能从 active slot 取走 RunningTask;正常 owner 发送 TurnComplete/自发 abort,外部 abort owner 发送 Interrupt/Replaced 的 TurnAborted。
11. 生命周期测试
11.1 正常完成flush
源码位置:codex-rs/core/src/session/tests.rs,测试 turn_complete_flushes_terminal_event_after_delivery。
let (mut sess, tc, rx) = make_session_and_context_with_rx().await;
let store = attach_in_memory_thread_store(
Arc::get_mut(&mut sess).expect("session should be uniquely owned"),
)
.await;
sess.spawn_task(Arc::clone(&tc), input, CompletingTask).await;
let event = recv_terminal_event(
&rx,
TerminalEventKind::TurnComplete,
)
.await;
assert!(matches!(event.msg, EventMsg::TurnComplete(_)));
// 第一次flush位于task body之后,第二次位于TurnComplete发送之后。
let calls = wait_for_flush_count(&store, /*expected_flushes*/ 2).await;
assert_eq!(2, calls.flush_thread);这个测试使用立即返回 Ok(None) 的 CompletingTask,排除了模型和工具噪声,直接证明统一 spawn/finish plumbing 建立两道持久化屏障。
该测试不等于覆盖真实模型、工具调用和并发输入路径;它只证明正常完成分支中的终态事件与两次 flush 顺序。
11.2 Interrupted
源码位置:codex-rs/core/src/session/tests.rs,测试 turn_aborted_flushes_terminal_event_after_delivery。
sess.spawn_task(
Arc::clone(&tc),
input,
NeverEndingTask {
kind: TaskKind::Regular,
listen_to_cancellation_token: true,
},
)
.await;
let abort_task = tokio::spawn({
let sess = Arc::clone(&sess);
async move {
sess.abort_all_tasks(TurnAbortReason::Interrupted).await;
}
});
let event = recv_terminal_event(&rx, TerminalEventKind::TurnAborted).await;
match event.msg {
EventMsg::TurnAborted(e) => {
assert_eq!(TurnAbortReason::Interrupted, e.reason)
}
other => panic!("unexpected event: {other:?}"),
}
abort_task.await.expect("abort task should finish");
// task body flush、interrupted marker flush、terminal event flush各一次。
let calls = wait_for_flush_count(&store, /*expected_flushes*/ 3).await;
assert_eq!(3, calls.flush_thread);它同时覆盖协作取消不会让 spawn closure 再发 TurnComplete:测试只接受 TurnAborted,并精确观察三个 barrier。
11.3 不协作的Task
源码位置:codex-rs/core/src/session/tests.rs,测试 abort_regular_task_emits_marker_before_turn_aborted。
sess.spawn_task(
Arc::clone(&tc),
input,
NeverEndingTask {
kind: TaskKind::Regular,
// fixture故意不监听token,迫使abort路径走100ms grace再abort handle。
listen_to_cancellation_token: false,
},
)
.await;
sess.abort_all_tasks(TurnAbortReason::Interrupted).await;
let marker_evt = tokio::time::timeout(
std::time::Duration::from_secs(2),
rx.recv(),
)
.await
.expect("timeout waiting for marker event")
.expect("event");
assert!(matches!(marker_evt.msg, EventMsg::RawResponseItem(_)));
let evt = tokio::time::timeout(
std::time::Duration::from_secs(2),
rx.recv(),
)
.await
.expect("timeout waiting for event")
.expect("event");
match evt.msg {
EventMsg::TurnAborted(e) => {
assert_eq!(TurnAbortReason::Interrupted, e.reason)
}
other => panic!("unexpected event: {other:?}"),
}
// abort owner只发一个terminal event,超时强杀不会制造尾随完成事件。
assert!(rx.try_recv().is_err());11.4 空槽清理
源码位置:codex-rs/core/src/session/tests.rs,测试 abort_empty_active_turn_preserves_pending_input。
let turn_state = {
let mut active = sess.active_turn.lock().await;
let active_turn = active.get_or_insert_with(ActiveTurn::default);
Arc::clone(&active_turn.turn_state)
};
sess.input_queue
.extend_pending_input_for_turn_state(
turn_state.as_ref(),
vec![TurnInput::ResponseItem(pending_item.clone())],
)
.await;
// ActiveTurn存在但task=None,Replaced不应把它误当作正在运行的Task清理输入。
sess.abort_all_tasks(TurnAbortReason::Replaced).await;
assert!(sess.active_turn.lock().await.is_none());
assert_eq!(
sess.input_queue
.take_pending_input_for_turn_state(turn_state.as_ref())
.await,
vec![TurnInput::ResponseItem(pending_item)]
);这个测试区分了 ActiveTurn reservation 与 RunningTask:abort_all_tasks 可以移除空 reservation,但只有确实 取得 task 时才执行 pending clear 与 abort lifecycle。
12. Owner定位
遇到 Task 生命周期问题时,可按下面顺序定位:
- 先看
active_turn.task是否仍存在;ActiveTurn存在不代表后台 future 已安装。 - 对照
TaskKind与span_name;UserShell 的 kind 也是 Regular,不能只按枚举判断实现。 - 正常完成缺 terminal event 时,检查 task body 后 flush、root token 是否已取消、
on_task_finished是否成功take()task。 - 出现双 terminal event 时,检查外部 abort 是否先从 slot 取走 owner,以及 spawn closure 是否在 token 取消后 仍调用正常 finish。
- Interrupt 后 marker 缺失时,区分 reason 是否为 Replaced,并检查 marker flush 是否失败。
- 旧 Task 收尾影响新 Task 时,检查 active slot 清理是否使用
Arc::ptr_eq(turn_state)身份门禁。
可以用 NeverEndingTask 的两个模式做只读推演:监听 token 时应在 grace 内 notify;不监听时应超时并 abort handle。两条路径之后都只能看到 marker(若启用)和一个 TurnAborted。普通用户输入如何进入 RegularTask 并 消费 startup prewarm,由 RegularTask完整流程 继续解释;Compact、Review 和 UserShell 的业务差异留给各自专题。
可以用下面的只读搜索把本文的 task lifecycle 主线落回源码:
rg -n "SessionTask|AnySessionTask|RunningTask|abort|join" codex-rs/core/src