入队前先判容量:满了不再做工程准备
- `direct_thread_manager` 新增只读的 `direct_pending_turn_count`(管理器 `pending_turn_count` + 模块级封装) - `enqueue_direct_codex_turn_typed` 在前置条件之后、工程准备之前先做一次容量预判,满队直接返回 `QueueFull` - 工程准备是分钟级、会在磁盘留产物的活,为注定收不下的消息先做它没有意义;权威判据仍是入队临界区里的检查 - 补一条管理器用例钉住这个只读口径:只数在队条目,认领 / 取消后立刻跟着变
This commit is contained in:
@@ -31,8 +31,8 @@ pub(crate) fn normalize_direct_client_turn_id(
|
||||
/// DirectProject 聊天命令:**只入队**,不 await、也不起回合。
|
||||
///
|
||||
/// 入队侧就是"这条消息能不能收下":`clientTurnId` 校验 → 工作流恢复 → 用户条目校验 → prompt
|
||||
/// 投影 → 前置条件 → 工程准备,任何一步失败都是**入队失败**(typed 载荷)。全部通过之后只做一件事
|
||||
/// ——把待发消息排进队列:不登记占用、不落盘、不发回合事件、不起 codex。
|
||||
/// 投影 → 前置条件 → 容量预判 → 工程准备,任何一步失败都是**入队失败**(typed 载荷)。全部通过
|
||||
/// 之后只做一件事——把待发消息排进队列:不登记占用、不落盘、不发回合事件、不起 codex。
|
||||
///
|
||||
/// 于是"这一轮跑成什么"仍然只有订阅事件一个来源:命令返回 `Ok` 只说明**入队成立**。真正的回合边界
|
||||
/// (`turn.started` / `turn.completed`)由 Thread Manager 在**放行**时写出(见 `direct_turn_dispatch`),
|
||||
@@ -65,11 +65,12 @@ pub(crate) async fn enqueue_direct_codex_turn(
|
||||
}
|
||||
|
||||
/// 入队侧:全程 typed,顺序固定,**每一步失败都还是入队失败**:
|
||||
/// `clientTurnId` 校验 → 工作流恢复 → 用户条目校验 → prompt 投影 → 前置条件 → 工程准备 → 入队。
|
||||
/// `clientTurnId` 校验 → 工作流恢复 → 用户条目校验 → prompt 投影 → 前置条件 → 容量预判 → 工程准备 → 入队。
|
||||
///
|
||||
/// 这个顺序不是风格问题:工程准备(可能建工程、铺脚手架)必须在"这条消息真的被收下"之前判定;
|
||||
/// 没通过检查的东西不进队列。入队本身只做判重与容量,不再有第二种副作用——落盘、回合事件与整轮
|
||||
/// 都在放行那一侧。
|
||||
/// 没通过检查的东西不进队列。容量放在工程准备之前是"别为注定收不下的消息先付分钟级的代价",它
|
||||
/// 与入队那一刻的容量检查同源,只是提前一次。入队本身只做判重与容量,不再有第二种副作用——落盘、
|
||||
/// 回合事件与整轮都在放行那一侧。
|
||||
async fn enqueue_direct_codex_turn_typed(
|
||||
root: &Path,
|
||||
user_item: DirectCodexUserItem,
|
||||
@@ -92,6 +93,16 @@ async fn enqueue_direct_codex_turn_typed(
|
||||
let user_prompt = direct_codex_user_item_to_prompt(root, &user_item)
|
||||
.map_err(|detail| DirectTurnError::InputRejected { detail })?;
|
||||
check_direct_turn_preconditions(root, &user_prompt, creation_type.as_deref())?;
|
||||
let thread_id = direct_thread_id_for_project(root);
|
||||
// 容量是**最便宜**的一条入队前置条件,先判:满了就直接说满,别让一条注定收不下的消息先去做工程
|
||||
// 准备——那是分钟级的活(`npm ci` + Vite 构建),还会在磁盘上留下产物。权威判据仍然是入队那一刻
|
||||
// 临界区里的容量检查(下面 `enqueue_direct_pending_turn`):这里只是快速失败,中间被别人的消息
|
||||
// 挤满时那一条照样拦得住。
|
||||
queue_has_room(direct_pending_turn_count(&thread_id)).map_err(|_| {
|
||||
DirectTurnError::QueueFull {
|
||||
limit: MAX_PENDING_DIRECT_TURNS,
|
||||
}
|
||||
})?;
|
||||
// 创建类型来自结构化用户入口;实际工程和可信脚手架由宿主复核。
|
||||
crate::environment_check::prepare_new_web_project_at(root, creation_type.as_deref())
|
||||
.await
|
||||
@@ -112,7 +123,6 @@ async fn enqueue_direct_codex_turn_typed(
|
||||
.map_err(|error| DirectTurnError::InputRejected {
|
||||
detail: error.to_string(),
|
||||
})?;
|
||||
let thread_id = direct_thread_id_for_project(root);
|
||||
match enqueue_direct_pending_turn(&thread_id, pending) {
|
||||
Ok(_) => {}
|
||||
Err(EnqueueRejection::QueueFull) => {
|
||||
|
||||
@@ -570,6 +570,14 @@ impl DirectThreadManager {
|
||||
.unwrap_or_default()
|
||||
}
|
||||
|
||||
/// 当前在队的待发消息条数(容量预判用,只读)。
|
||||
fn pending_turn_count(&self, thread_id: &str) -> usize {
|
||||
self.threads
|
||||
.get(thread_id)
|
||||
.map(|thread| Self::pending_turns(thread).len())
|
||||
.unwrap_or(0)
|
||||
}
|
||||
|
||||
/// 观察一条事件,返回它本身是否可回收。
|
||||
fn observe_event(thread: &mut ThreadState, seq: u64, event: &DirectThreadEvent) -> bool {
|
||||
match event {
|
||||
@@ -775,6 +783,17 @@ pub(crate) fn enqueue_direct_pending_turn(
|
||||
outcome
|
||||
}
|
||||
|
||||
/// 这个 thread 上**在队**的待发消息条数(只读)。
|
||||
///
|
||||
/// 给入队侧做最便宜的一次容量预判用:满了就先拒绝,不去做工程准备那种分钟级、会在磁盘上留下产物的
|
||||
/// 活。权威判据仍然是入队临界区里的 `queue_has_room`,这里只求早失败。
|
||||
pub(crate) fn direct_pending_turn_count(thread_id: &str) -> usize {
|
||||
global_direct_thread_manager()
|
||||
.lock()
|
||||
.unwrap_or_else(|poisoned| poisoned.into_inner())
|
||||
.pending_turn_count(thread_id)
|
||||
}
|
||||
|
||||
/// 取消一条待发消息:结果告诉调用方它是被取消、已经放行,还是本来就不在队里。
|
||||
pub(crate) fn remove_direct_pending_turn(
|
||||
thread_id: &str,
|
||||
@@ -1462,6 +1481,40 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
/// 入队侧拿来做容量预判的条数与权威判据同源:只数在队条目,认领 / 取消之后立刻跟着变。
|
||||
#[test]
|
||||
fn pending_turn_count_tracks_only_the_queue() {
|
||||
let mut manager = DirectThreadManager::with_limits(100, 100_000);
|
||||
assert_eq!(manager.pending_turn_count("thread-unknown"), 0);
|
||||
for index in 0..MAX_PENDING_DIRECT_TURNS {
|
||||
manager
|
||||
.enqueue_pending_turn("thread-1", pending_turn(&format!("turn-{index}")))
|
||||
.expect("enqueue");
|
||||
}
|
||||
assert_eq!(
|
||||
manager.pending_turn_count("thread-1"),
|
||||
MAX_PENDING_DIRECT_TURNS
|
||||
);
|
||||
assert_eq!(
|
||||
queue_has_room(manager.pending_turn_count("thread-1")),
|
||||
Err(EnqueueRejection::QueueFull)
|
||||
);
|
||||
|
||||
// 认领队首与取消一条在队消息都立刻反映到这条只读口径上。
|
||||
manager
|
||||
.claim_pending_turn("thread-1", FIXED_AT_MS)
|
||||
.expect("claim head");
|
||||
assert_eq!(
|
||||
manager.pending_turn_count("thread-1"),
|
||||
MAX_PENDING_DIRECT_TURNS - 1
|
||||
);
|
||||
manager.remove_pending_turn("thread-1", "turn-1", FIXED_AT_MS);
|
||||
assert_eq!(
|
||||
manager.pending_turn_count("thread-1"),
|
||||
MAX_PENDING_DIRECT_TURNS - 2
|
||||
);
|
||||
}
|
||||
|
||||
/// 放行认领的是**队首**,而且开始事件与离开队列事件同批写:中间没有第二个中间态。
|
||||
#[test]
|
||||
fn dispatch_claims_the_head_and_writes_both_events_in_one_batch() {
|
||||
|
||||
@@ -64,8 +64,9 @@
|
||||
改动点:
|
||||
|
||||
- `agent/direct_runtime/user_input.rs`:把命令主体拆成两半。
|
||||
**入队侧**:`clientTurnId` 校验 → 工作流恢复 → 用户条目校验 → prompt 投影 → 前置条件 → 工程准备 → 入队;
|
||||
任何一步失败返回 typed 入队失败。**放行侧**(第 2 步)从占用登记起。
|
||||
**入队侧**:`clientTurnId` 校验 → 工作流恢复 → 用户条目校验 → prompt 投影 → 前置条件 → 容量预判 → 工程准备 → 入队;
|
||||
任何一步失败返回 typed 入队失败。容量预判只是一次提前的快速失败,权威判据仍在入队临界区里(工程准备是分钟级、会
|
||||
在磁盘留产物的活,满了就不该先做它)。**放行侧**(第 2 步)从占用登记起。
|
||||
- 队列条目(宿主侧产物,只在内存):`PendingDirectTurn { client_turn_id, user_item: Value, prompt: String, creation_type: Option<String>, at: u64 }`。
|
||||
`prompt` 与 `creation_type` 是入队检查的产物,放行不再重算;事件里**不带** `prompt`。
|
||||
- `agent/direct_thread_manager.rs`:
|
||||
|
||||
Reference in New Issue
Block a user