残留回合兜底释放改成 Thread Manager 里的原子动作

- `DirectThreadManager::release_stale_turn` 在同一个临界区里做完"校验 → 解除占用 → 写 aborted 终态"
- 判据不满足(没有活动回合 / 身份对不上 / 没过启动窗口)一个字都不写,占用与事件都保持原样
- 拆成"先校验后释放"两次取锁会让并发认领插进中间,清掉新那一轮的占用并给老回合补假终态
- `direct_stale_turn_for_release` 降级成"调用原子出口 + 把结果翻成人话",文案与对外语义不变
- 兜底终态改由 Thread Manager 写,身份取被释放那一轮自己的 clientTurnId,删掉 app-server 侧的建事件助手
- 测试:管理器侧新增"校验与释放是一个动作"用例;运行时侧用例补上"拒绝时零写入 / 通过时占用与终态一起落地"
This commit is contained in:
2026-09-30 12:55:07 +08:00
parent e0ca5ad9b9
commit 51e4811344
3 changed files with 237 additions and 77 deletions
@@ -4465,14 +4465,9 @@ enum DirectCodexTurnCancelTarget {
/// [`direct_stale_turn_for_release`] 的注释;"正在跑的是另一轮"仍然保持原拒绝语义,
/// 什么都不释放。
///
/// 兜底终态带 `userItemId`:身份取 [`direct_stale_turn_for_release`] 返回的
/// clientTurnId(客户端回合身份的唯一来源),与正常路径的开口条目 id 同一份 canonical 口径。
/// 拿不到 clientTurnId 就留空——这一轮不会再有原生终态,猜一个身份会让前端把边界盖到别人身上。
fn direct_stale_cancel_turn_completed_event(client_turn_id: &str) -> DirectThreadEvent {
DirectThreadEvent::turn_completed("aborted".to_string(), direct_tool_call_now_ms())
.with_user_item_id(direct_codex_user_item_id_for_client_turn_id(client_turn_id).as_deref())
}
/// 兜底终态由 Thread Manager 在解除占用的同一个临界区里写下(带 `userItemId`,身份取被释放那一轮
/// 自己的 `clientTurnId`,与放行时 `turn.started` 同一份 canonical 口径):这一轮不会再有原生终态,
/// 猜一个身份会让前端把边界盖到别人身上。
pub(crate) fn cancel_direct_codex_turn_at(
root: &Path,
client_turn_id: Option<&str>,
@@ -4523,16 +4518,10 @@ pub(crate) fn cancel_direct_codex_turn_at(
})
}
DirectCodexTurnCancelTarget::Stale(reason) => {
// 先做前置校验:确实登记着一轮、身份对得上、过了启动窗口才允许判成残留。
// 校验与释放**同一个临界区**(确实登记着一轮、身份对得上、过了启动窗口才允许判成残留;
// 占用与兜底终态要么一起落地,要么一个字都不写)。校验完了再单独释放不行:两步之间并发
// 认领可能已经登记了新那一轮的占用,无条件解除就会清掉它。
let released = direct_stale_turn_for_release(root, client_turn_id, reason)?;
// 这一轮不会再有人替它发终态事件(执行进程已退出 / 从没进执行器),
// 兜底补一条,否则前端的"最新回合是否在跑"会永远停在运行中。走 Thread Manager 的
// 深层出口而不是裸 append:这也是**解除占用的唯一出口**,解除不了这个 thread 会一直
// 被认为是"还有没收口的回合",挡住后面的放行。
complete_direct_thread_turn(
&direct_thread_id_for_project(root),
direct_stale_cancel_turn_completed_event(&released),
);
// 这一轮不会再有人替它收尾(占用刚被兜底解除):踢一脚让队列继续。
// 少了这一脚,排在这条后面的待发消息要等到下一次用户动作才会被放行。
crate::agent::kick_direct_queue_dispatch(root);
@@ -6101,35 +6090,6 @@ mod tests {
);
}
/// 取消兜底终态也要带开口用户条目身份,且身份只有一个来源:release 返回的 clientTurnId
/// 走与正常路径同一份 canonical 口径;拿不到(空 / 空白)就留空,不猜。
#[test]
fn stale_cancel_terminal_event_keeps_opener_user_item_id_from_client_turn_id() {
let event = direct_stale_cancel_turn_completed_event("turn-0001");
assert_eq!(event.user_item_id(), Some("direct-codex:turn-0001:user"));
assert!(event.at().is_some(), "兜底终态仍要带宿主观测时间");
assert!(matches!(
event,
DirectThreadEvent::TurnCompleted { ref status, .. } if status == "aborted"
));
for missing in ["", " "] {
let event = direct_stale_cancel_turn_completed_event(missing);
assert_eq!(
event.user_item_id(),
None,
"拿不到 clientTurnId 时不得编造开口条目身份"
);
}
// canonical 口径与落盘侧同一份:`direct-codex:{clientTurnId}:user`。
assert_eq!(
direct_codex_user_item_id_for_client_turn_id(" turn-0001 ").as_deref(),
Some("direct-codex:turn-0001:user")
);
assert_eq!(direct_codex_user_item_id_for_client_turn_id(""), None);
}
/// 封口复核要求(`RepairRequired`)是控制流:这一轮还没结束,不能写出失败终态。
///
/// 反过来,真失败必须写成载荷,`kind` / `message` 从 typed 错误投影——两者以前共用
@@ -496,7 +496,8 @@ impl DirectStaleTurnReleaseReason {
}
}
/// 兜底释放前的前置校验:这一轮是不是**可以**被判成残留,返回它的 `clientTurnId`。
/// 兜底释放一个残留回合:校验与释放**在 Thread Manager 的同一个临界区里**完成,返回被释放的
/// `clientTurnId`。
///
/// 条件(三条必须同时成立,这段注释就是契约):
/// 1. 这个 thread 上确实登记着一轮(Thread Manager 的活动回合);
@@ -504,34 +505,33 @@ impl DirectStaleTurnReleaseReason {
/// 3. `reason == NeverReachedExecutor` 时,登记的年龄必须超过
/// [`DIRECT_STALE_TURN_RELEASE_MIN_AGE_MS`],排除"刚放行、还在本地准备"的正常启动窗口。
///
/// **释放本身不在这里**:调用方拿到身份后往 Thread Manager 写一条 `turn.completed{aborted}`,
/// 那才是解除占用的唯一出口(也因此天然幂等——占用已经没了就返回"没有正在运行的回合")。
/// 不满足任何一条都**什么都不写**(占用、事件都不变),返回可读原因;因此它天然幂等——占用已经
/// 没了就返回"没有正在运行的回合"。这一层只负责把 [`DirectStaleTurnRelease`] 翻成人话:判据与释放
/// 动作都在 Thread Manager 里,因为它们必须在同一个临界区里(拆成两步会让并发认领插进中间,见
/// [`release_stale_direct_turn`])。
pub(crate) fn direct_stale_turn_for_release(
root: &Path,
expected_client_turn_id: Option<&str>,
reason: DirectStaleTurnReleaseReason,
) -> Result<String, String> {
let thread_id = direct_thread_id_for_project(root);
let active = read_direct_turn_identity(&thread_id)
.ok_or_else(|| "当前项目没有正在运行的陶泥儿回合,无法终止".to_string())?;
if let Some(expected) = expected_client_turn_id
let expected = expected_client_turn_id
.map(str::trim)
.filter(|value| !value.is_empty())
{
if active.client_turn_id != expected {
return Err("正在运行的是另一条 Direct 客户端回合,已拒绝终止".to_string());
.filter(|value| !value.is_empty());
let min_age_ms = (reason == DirectStaleTurnReleaseReason::NeverReachedExecutor)
.then_some(DIRECT_STALE_TURN_RELEASE_MIN_AGE_MS);
match release_stale_direct_turn(&direct_thread_id_for_project(root), expected, min_age_ms) {
DirectStaleTurnRelease::Released(client_turn_id) => Ok(client_turn_id),
DirectStaleTurnRelease::NoActiveTurn => {
Err("当前项目没有正在运行的陶泥儿回合,无法终止".to_string())
}
}
if reason == DirectStaleTurnReleaseReason::NeverReachedExecutor {
let age_ms = direct_taonier_active_now_millis().saturating_sub(active.started_at);
if age_ms < DIRECT_STALE_TURN_RELEASE_MIN_AGE_MS {
return Err(format!(
"这一轮 Direct 客户端回合刚开始 {} 秒、还在准备中,暂不能强制释放;请稍后再试",
age_ms / 1000
));
DirectStaleTurnRelease::AnotherTurnRunning => {
Err("正在运行的是另一条 Direct 客户端回合,已拒绝终止".to_string())
}
DirectStaleTurnRelease::WithinStartupWindow { age_ms } => Err(format!(
"这一轮 Direct 客户端回合刚开始 {} 秒、还在准备中,暂不能强制释放;请稍后再试",
age_ms / 1000
)),
}
Ok(active.client_turn_id)
}
fn direct_taonier_regeneration_project_id(root: &Path) -> Result<String, String> {
@@ -5542,12 +5542,15 @@ mod tests {
assert!(direct_active_turn_id_at(root.path()).is_err());
}
/// 兜底释放的前置判据:身份对得上,且"从没进执行器"必须过启动窗口。
/// 兜底释放的前置判据:身份对得上,且"从没进执行器"必须过启动窗口;判据不过**一个字都不写**,
/// 通过则占用与兜底终态在同一个临界区里一起落地。
#[test]
fn stale_turn_release_requires_a_matching_identity_and_only_after_the_start_window() {
let root = tempfile::tempdir().expect("stale turn root");
let thread_id = direct_thread_id_for_project(root.path());
let reservation = DirectTurnReservation::accept_for_test(&thread_id, "client-turn-stale-1");
let watcher = subscribe_direct_thread(&thread_id);
let _ = consume_direct_thread(&watcher.subscription_id).expect("drain bootstrap");
// ① clientTurnId 不匹配:拒绝,且不误伤正在跑的那一轮。
let mismatch = direct_stale_turn_for_release(
@@ -5557,10 +5560,6 @@ mod tests {
)
.expect_err("another turn must not be released");
assert!(mismatch.contains("另一条"), "{mismatch}");
assert_eq!(
read_direct_turn_identity(&thread_id).map(|turn| turn.client_turn_id),
Some("client-turn-stale-1".to_string())
);
// ② 刚放行、还没进执行器:正常启动窗口内不许释放(释放等于放开并发)。
let young = direct_stale_turn_for_release(
@@ -5571,7 +5570,20 @@ mod tests {
.expect_err("a freshly dispatched turn is still starting");
assert!(young.contains("暂不能强制释放"), "{young}");
// ③ 老过窗口:判成残留,调用方拿到身份(解除占用由调用方写终态完成)。
// ① ② 两条拒绝路径都**什么都没写**:占用还在,事件流上也没有任何东西。
assert_eq!(
read_direct_turn_identity(&thread_id).map(|turn| turn.client_turn_id),
Some("client-turn-stale-1".to_string())
);
assert!(
consume_direct_thread(&watcher.subscription_id)
.expect("consume after rejected releases")
.events
.is_empty(),
"判据不过不许写任何事件"
);
// ③ 老过窗口:判成残留——占用与兜底终态在同一个临界区里一起落地。
backdate_direct_active_turn_for_test(&thread_id, DIRECT_STALE_TURN_RELEASE_MIN_AGE_MS + 1);
let released = direct_stale_turn_for_release(
root.path(),
@@ -5580,7 +5592,33 @@ mod tests {
)
.expect("stale turn is releasable");
assert_eq!(released, "client-turn-stale-1");
// ④ "执行进程已退出"不看启动窗口:刚放行也允许释放。
assert_eq!(
read_direct_turn_identity(&thread_id),
None,
"占用必须一起解除"
);
let events = consume_direct_thread(&watcher.subscription_id)
.expect("consume released terminal")
.events;
assert_eq!(events.len(), 1, "兜底释放只写一条终态:{events:?}");
let terminal = &events[0];
assert!(
matches!(
terminal,
DirectThreadEvent::TurnCompleted { status, .. } if status == "aborted"
),
"{terminal:?}"
);
assert_eq!(
terminal.user_item_id().as_deref(),
Some("direct-codex:client-turn-stale-1:user"),
"兜底终态的身份取被释放那一轮自己的 clientTurnId"
);
assert!(terminal.at().is_some(), "兜底终态仍要带宿主观测时间");
// ④ 释放之后这个 thread 立刻腾出来了(占用解除与终态同批落地):下一条能重新认领,
// 而且"执行进程已退出"不看启动窗口——刚放行也允许释放。
let _next = DirectTurnReservation::accept_for_test(&thread_id, "client-turn-stale-2");
assert_eq!(
direct_stale_turn_for_release(
root.path(),
@@ -5588,7 +5626,7 @@ mod tests {
DirectStaleTurnReleaseReason::ExecutorExited,
)
.expect("executor exited ignores the start window"),
"client-turn-stale-1"
"client-turn-stale-2"
);
drop(reservation);
}
@@ -9,9 +9,10 @@ use std::sync::{Mutex, OnceLock};
use uuid::Uuid;
use crate::agent::{
direct_tool_call_now_ms, queue_has_room, DirectQueueRemovalOutcome, DirectQueueRemovalReason,
DirectThreadConsumeResult, DirectThreadEvent, DirectThreadSubscriptionBootstrap,
EnqueueOutcome, EnqueueRejection, PendingDirectTurn,
direct_codex_user_item_id_for_client_turn_id, direct_tool_call_now_ms, queue_has_room,
DirectQueueRemovalOutcome, DirectQueueRemovalReason, DirectThreadConsumeResult,
DirectThreadEvent, DirectThreadSubscriptionBootstrap, EnqueueOutcome, EnqueueRejection,
PendingDirectTurn,
};
const DEFAULT_MAX_EVENTS: usize = 8_192;
@@ -96,6 +97,22 @@ pub(crate) struct DirectTurnIdentity {
pub(crate) started_at: u64,
}
/// 兜底释放一个残留回合的结果。
///
/// 除了 [`Self::Released`],其余三种都是**什么都没写**:占用没动、事件没发。调用方据此给用户
/// 不同的话;把"校验"和"释放"拆成两步就会丢掉这个性质(见 [`DirectThreadManager::release_stale_turn`])。
#[derive(Clone, Debug, Eq, PartialEq)]
pub(crate) enum DirectStaleTurnRelease {
/// 占用已解除,`turn.completed{aborted}` 已下发:这个 thread 现在可以放行下一条。
Released(String),
/// 这个 thread 上没有未收口的回合。
NoActiveTurn,
/// 登记的回合身份与调用方给的不一致:绝不误伤另一条回合。
AnotherTurnRunning,
/// 登记的年龄还没过启动窗口(`age_ms` 是登记到 `now_ms` 的差)。
WithinStartupWindow { age_ms: u64 },
}
/// 一次放行的结果:队首那条待发消息,以及这次放行的占用身份。
///
/// `token` 是这一轮的占用身份(终态出口只认它),`pending` 是放行要用的**全部**输入——放行不再
@@ -327,6 +344,58 @@ impl DirectThreadManager {
true
}
/// 兜底释放一个残留回合:**同一个临界区里**做完"校验 → 解除占用 → 写终态"。
///
/// 这三件事必须原子。拆成"先校验、后释放"两次取锁的话,两步之间占用可能已经被真实收口释放,
/// 队列顺势认领了下一条并登记新占用——那时再无条件解除占用,就会清掉**新那一轮**的占用,还给
/// 已经结束的老回合补一条假终态。这正是"多轮同时在跑 / 终态张冠李戴"那一类事故。
///
/// 三条前置条件与 [`DirectStaleTurnRelease`] 的三个失败分支一一对应:
/// 1. 这个 thread 上确实登记着一轮;
/// 2. 传了 `expected_client_turn_id` 时必须与登记一致;
/// 3. `min_age_ms` 有值时,登记的年龄(`now_ms - started_at`)必须不小于它——用来排除
/// "刚放行、还在本地准备"的正常启动窗口。
///
/// 不满足任何一条都**一个字都不写**(占用与事件都保持原样)。通过则终态与占用一起落地:占用
/// 一解除,队列就可以认领下一条。
fn release_stale_turn(
&mut self,
thread_id: &str,
expected_client_turn_id: Option<&str>,
min_age_ms: Option<u64>,
now_ms: u64,
) -> DirectStaleTurnRelease {
let released = {
let Some(thread) = self.threads.get_mut(thread_id) else {
return DirectStaleTurnRelease::NoActiveTurn;
};
let Some(active) = thread.active_turn.take() else {
return DirectStaleTurnRelease::NoActiveTurn;
};
if expected_client_turn_id.is_some_and(|expected| expected != active.turn_id.as_str()) {
// 身份对不上:放回去,绝不误伤另一条回合。
thread.active_turn = Some(active);
return DirectStaleTurnRelease::AnotherTurnRunning;
}
let age_ms = now_ms.saturating_sub(active.started_at);
if min_age_ms.is_some_and(|min_age_ms| age_ms < min_age_ms) {
thread.active_turn = Some(active);
return DirectStaleTurnRelease::WithinStartupWindow { age_ms };
}
active.turn_id
};
// 这一轮不会再有人替它发终态(执行进程已退出 / 从没进执行器):补一条 aborted,否则前端的
// "最新回合是否在跑"会永远停在运行中。身份用被释放那一轮自己的 `turn_id` 推导,与放行时
// `turn.started` 的口径同一份(`PendingDirectTurn::user_item_id`)。
self.append(
thread_id,
DirectThreadEvent::turn_completed("aborted".to_string(), now_ms).with_user_item_id(
direct_codex_user_item_id_for_client_turn_id(&released).as_deref(),
),
);
DirectStaleTurnRelease::Released(released)
}
/// 这个 thread 上正在跑的那一轮身份;`None` = 没有未收口的回合。
fn active_turn_identity(&self, thread_id: &str) -> Option<DirectTurnIdentity> {
self.threads.get(thread_id).and_then(|thread| {
@@ -881,6 +950,31 @@ pub(crate) fn complete_direct_thread_turn_if_reserved(
written
}
/// 兜底释放一个残留回合:校验与释放**同一个临界区**(见 [`DirectThreadManager::release_stale_turn`])。
///
/// 只有真释放了才通知订阅者:校验不过时什么都没写,不该惊动任何人。
pub(crate) fn release_stale_direct_turn(
thread_id: &str,
expected_client_turn_id: Option<&str>,
min_age_ms: Option<u64>,
) -> DirectStaleTurnRelease {
let outcome = {
global_direct_thread_manager()
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.release_stale_turn(
thread_id,
expected_client_turn_id,
min_age_ms,
direct_tool_call_now_ms(),
)
};
if matches!(outcome, DirectStaleTurnRelease::Released(_)) {
notify_direct_thread_subscribers(thread_id);
}
outcome
}
/// 这个 thread 上正在跑的那一轮身份。
///
/// 只读:占用仍然只由终态出口(深层终态 / 占用对象兜底 / 兜底释放)解除,这里不改任何状态。
@@ -1596,6 +1690,74 @@ mod tests {
assert_eq!(manager.pending_turn_ids("thread-1"), vec!["turn-2"]);
}
/// 兜底释放是**一个动作**:判据不过时占用与事件都不动,通过时占用解除与兜底终态一起落地。
#[test]
fn stale_release_validates_and_releases_as_one_action() {
let mut manager = DirectThreadManager::with_limits(100, 100_000);
manager
.enqueue_pending_turn("thread-1", pending_turn("turn-1"))
.expect("enqueue");
manager
.claim_pending_turn("thread-1", FIXED_AT_MS)
.expect("claim head");
let watcher = manager.subscribe("thread-1");
let _ = manager.consume(&watcher.subscription_id).expect("drain");
// 身份对不上:放回去,占用不动、一个字都不写。
assert_eq!(
manager.release_stale_turn("thread-1", Some("turn-other"), None, FIXED_AT_MS + 10_000),
DirectStaleTurnRelease::AnotherTurnRunning
);
// 启动窗口内:同样什么都不写。
assert_eq!(
manager.release_stale_turn("thread-1", None, Some(60_000), FIXED_AT_MS + 1),
DirectStaleTurnRelease::WithinStartupWindow { age_ms: 1 }
);
assert!(manager.turn_is_active("thread-1"));
assert!(
manager
.consume(&watcher.subscription_id)
.expect("consume rejected releases")
.events
.is_empty(),
"判据不过不许写任何事件"
);
// 通过:占用解除与兜底终态在同一个临界区里落地,身份取被释放那一轮自己的 turn_id。
assert_eq!(
manager.release_stale_turn(
"thread-1",
Some("turn-1"),
Some(60_000),
FIXED_AT_MS + 60_000
),
DirectStaleTurnRelease::Released("turn-1".to_string())
);
assert!(!manager.turn_is_active("thread-1"));
let events = manager
.consume(&watcher.subscription_id)
.expect("consume released terminal")
.events;
assert_eq!(events.len(), 1, "兜底释放只写一条终态:{events:?}");
assert!(
matches!(
events[0],
DirectThreadEvent::TurnCompleted { ref status, .. } if status == "aborted"
),
"{events:?}"
);
assert_eq!(
events[0].user_item_id().as_deref(),
Some("direct-codex:turn-1:user")
);
// 占用已经解除,再释放就是"没有活动回合"。
assert_eq!(
manager.release_stale_turn("thread-1", None, None, FIXED_AT_MS + 60_000),
DirectStaleTurnRelease::NoActiveTurn
);
}
/// 取消的三种结果分得开:真的移除了、已经被放行了、本来就不在队里。
#[test]
fn removal_tells_cancelled_dispatched_and_unknown_apart() {