From 053e0bd93220b118e997ce59edf9f062a81a7cf6 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=8E=8B=E5=BE=B7=E5=AE=87?= Date: Thu, 24 Sep 2026 21:11:15 +0800 Subject: [PATCH] =?UTF-8?q?Thread=20Manager=EF=BC=9A=E5=BE=85=E5=8F=91?= =?UTF-8?q?=E6=B6=88=E6=81=AF=E9=98=9F=E5=88=97=E7=9A=84=E5=85=A5=E9=98=9F?= =?UTF-8?q?=20/=20=E5=8F=96=E6=B6=88=20/=20=E6=94=BE=E8=A1=8C=E8=AE=A4?= =?UTF-8?q?=E9=A2=86?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - `StoredEvent` 增加非序列化的 `pending` 产物字段:它的有无就是队列成员身份,队列的先后就是事件先后,不另建队列表 - 新增 `enqueue_pending_turn`:按 `clientTurnId` 幂等(判重范围「在队 ∪ 正在跑的那一轮」),容量用 `queue_has_room` 挡在队条目 - 新增 `remove_pending_turn`:typed 结果 `Removed | AlreadyDispatched | NotFound`,在临界区里取走产物并追加 `queue.removed{cancelled}` - 新增 `claim_pending_turn`:同一临界区取队首 → 登记占用 → 追加 `turn.started` 与 `queue.removed{dispatched}`,已有未收口回合或队列为空时返回 `None` - `mark_queue_events_cleanable` 跳过仍挂着产物的条目,保证在队条目永远不会被回收 - 新增薄包装 `enqueue_direct_pending_turn` / `remove_direct_pending_turn` / `claim_direct_pending_turn`,并补 5 条单测覆盖 bootstrap 可见、幂等、容量、放行原子性与三种取消结果 --- .../src/agent/direct_thread_manager.rs | 451 +++++++++++++++++- 1 file changed, 449 insertions(+), 2 deletions(-) diff --git a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_thread_manager.rs b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_thread_manager.rs index 2e66f086c..a83e68baa 100644 --- a/apps/ai-game-creator-shell/src-tauri/src/agent/direct_thread_manager.rs +++ b/apps/ai-game-creator-shell/src-tauri/src/agent/direct_thread_manager.rs @@ -9,7 +9,9 @@ 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, }; const DEFAULT_MAX_EVENTS: usize = 8_192; @@ -28,6 +30,12 @@ struct StoredEvent { seq: u64, bytes: usize, cleanable: bool, + /// 这条事件在队期间挂着的宿主侧产物。 + /// + /// 只有 `queue.enqueued` 会带它,而且**它本身就是队列成员身份**:`Some` = 这条待发消息还在队, + /// 被取走 = 已经离开队列。队列的先后就是事件的先后,所以没有第二张队列表要对齐,也不会出现 + /// "事件说在队、队列表说不在"的中间态(见 `direct_thread_queue`)。 + pending: Option, } #[derive(Clone, Debug)] @@ -75,6 +83,16 @@ pub(crate) struct DirectActiveTurnSnapshot { pub(crate) sequence: u64, } +/// 一次放行的结果:队首那条待发消息,以及这次放行的占用身份。 +/// +/// `token` 是这一轮的占用身份(终态出口只认它),`pending` 是放行要用的**全部**输入——放行不再 +/// 重跑任何检查,也就不需要再回到调用方取参数。 +#[derive(Clone, Debug)] +pub(crate) struct DirectDispatchedTurn { + pub(crate) token: String, + pub(crate) pending: PendingDirectTurn, +} + #[derive(Clone, Debug)] struct ThreadState { next_seq: u64, @@ -142,6 +160,16 @@ impl DirectThreadManager { &mut self, thread_id: &str, event: DirectThreadEvent, + ) -> DirectThreadEvent { + self.append_inner(thread_id, event, None) + } + + /// 追加事件。`pending` 用来把宿主侧产物挂到这条事件上(只有入队会传)。 + fn append_inner( + &mut self, + thread_id: &str, + event: DirectThreadEvent, + pending: Option, ) -> DirectThreadEvent { let thread = self.threads.entry(thread_id.to_string()).or_default(); thread.next_seq = thread.next_seq.saturating_add(1); @@ -156,6 +184,7 @@ impl DirectThreadManager { seq, bytes, cleanable, + pending, }); if let Some(item_id) = event.item_id() { Self::mark_item_events_cleanable(thread, item_id, seq); @@ -331,6 +360,153 @@ impl DirectThreadManager { .is_some_and(|thread| thread.active_turn.is_some()) } + /// 在队的待发消息,队首在前。 + fn pending_turns(thread: &ThreadState) -> Vec<&PendingDirectTurn> { + thread + .events + .iter() + .skip(thread.head) + .filter_map(|stored| stored.pending.as_ref()) + .collect() + } + + /// 入队:把一条已经通过入队检查的待发消息排到队尾。 + /// + /// 幂等按 `clientTurnId`,判重范围是「**在队 ∪ 正在跑的那一轮**」:同一身份重复入队返回 + /// [`EnqueueOutcome::AlreadyKnown`],不排第二条、不发第二条事件。容量只数在队条目。 + fn enqueue_pending_turn( + &mut self, + thread_id: &str, + turn: PendingDirectTurn, + ) -> Result { + { + let thread = self.threads.entry(thread_id.to_string()).or_default(); + let already_known = Self::pending_turns(thread) + .iter() + .any(|pending| pending.client_turn_id == turn.client_turn_id) + || thread + .active_turn + .as_ref() + .is_some_and(|active| active.turn_id == turn.client_turn_id); + if already_known { + return Ok(EnqueueOutcome::AlreadyKnown); + } + queue_has_room(Self::pending_turns(thread).len())?; + } + let event = turn.enqueued_event(); + self.append_inner(thread_id, event, Some(turn)); + Ok(EnqueueOutcome::Enqueued) + } + + /// 取消一条待发消息。typed 结果告诉调用方它到底是被取消、已经放行,还是本来就不在队里。 + /// + /// 判「已经被放行」用的是正在跑的那一轮的回合身份:认领之后就再也按待发消息取消不了了。 + fn remove_pending_turn( + &mut self, + thread_id: &str, + client_turn_id: &str, + at_ms: u64, + ) -> DirectQueueRemovalOutcome { + let removed = { + let Some(thread) = self.threads.get_mut(thread_id) else { + return DirectQueueRemovalOutcome::NotFound; + }; + if thread + .active_turn + .as_ref() + .is_some_and(|active| active.turn_id == client_turn_id) + { + return DirectQueueRemovalOutcome::AlreadyDispatched; + } + let head = thread.head; + let mut removed = false; + for stored in thread.events.iter_mut().skip(head) { + if stored + .pending + .as_ref() + .is_some_and(|pending| pending.client_turn_id == client_turn_id) + { + // 取走产物即离开队列:这一幕在临界区里发生,外面看不到"事件还在、条目已经不在"。 + stored.pending = None; + removed = true; + } + } + removed + }; + if !removed { + return DirectQueueRemovalOutcome::NotFound; + } + self.append( + thread_id, + DirectThreadEvent::queue_removed( + client_turn_id.to_string(), + DirectQueueRemovalReason::Cancelled, + at_ms, + ), + ); + DirectQueueRemovalOutcome::Removed + } + + /// 放行:同一个临界区里取队首 → 登记占用 → 写 `turn.started` 与 `queue.removed{dispatched}`。 + /// + /// 返回 `None` 表示现在**不该放行**:这个 thread 已经有未收口的回合,或者队列是空的。这两件事 + /// 都由这里判,调用方(kick)不需要先探一遍再决定——探两次就是两个中间态。 + /// + /// 认领认的是队首,不是调用方点名的某一条:待发消息的顺序就是事件顺序,没有第三种来源。 + fn claim_pending_turn( + &mut self, + thread_id: &str, + started_at_ms: u64, + ) -> Option { + let pending = { + let thread = self.threads.get_mut(thread_id)?; + if thread.active_turn.is_some() { + return None; + } + let head = thread.head; + let stored = thread + .events + .iter_mut() + .skip(head) + .find(|stored| stored.pending.is_some())?; + stored.pending.take().expect("filtered on a pending turn") + }; + let token = Uuid::new_v4().to_string(); + let turn_id = pending.client_turn_id.clone(); + let project_name = std::path::Path::new(thread_id) + .file_name() + .and_then(|name| name.to_str()) + .map(str::to_string); + if let Some(thread) = self.threads.get_mut(thread_id) { + thread.active_turn = Some(ActiveDirectTurn { + token: token.clone(), + turn_id: turn_id.clone(), + project_name, + started_at: started_at_ms, + // 与"还没有任何进度事件"的状态一致:运行时给出的第一条进度会覆盖它。 + status: "accepted".to_string(), + activity: Some("request-accepted".to_string()), + updated_at: started_at_ms, + sequence: 0, + }); + } + let user_item_id = pending.user_item_id(); + self.append( + thread_id, + DirectThreadEvent::turn_started(started_at_ms) + .with_user_item_id(user_item_id.as_deref()), + ); + self.append( + thread_id, + DirectThreadEvent::queue_removed( + turn_id, + DirectQueueRemovalReason::Dispatched, + started_at_ms, + ), + ); + Some(DirectDispatchedTurn { token, pending }) + } + fn subscriber_ids(&self, thread_id: &str) -> Vec { self.threads .get(thread_id) @@ -395,6 +571,20 @@ impl DirectThreadManager { }) } + /// 当前在队的待发消息身份,队首在前。 + #[cfg(test)] + fn pending_turn_ids(&self, thread_id: &str) -> Vec { + self.threads + .get(thread_id) + .map(|thread| { + Self::pending_turns(thread) + .iter() + .map(|pending| pending.client_turn_id.clone()) + .collect() + }) + .unwrap_or_default() + } + /// 观察一条事件,返回它本身是否可回收。 fn observe_event(thread: &mut ThreadState, seq: u64, event: &DirectThreadEvent) -> bool { match event { @@ -482,7 +672,10 @@ impl DirectThreadManager { /// 一条待发消息离开队列:它和它对应的入队事件一起变成可回收。 fn mark_queue_events_cleanable(thread: &mut ThreadState, client_turn_id: &str) { for stored in &mut thread.events { - if stored.event.queue_client_turn_id() == Some(client_turn_id) { + // 仍挂着产物的条目永远不算可回收:在队条目被回收就等于队列悄悄丢了一条消息。 + if stored.pending.is_none() + && stored.event.queue_client_turn_id() == Some(client_turn_id) + { stored.cleanable = true; } } @@ -601,6 +794,56 @@ pub(crate) fn accept_direct_thread_turn( Ok(()) } +/// 入队一条待发消息(队列的成员与顺序就是事件列表,见 `direct_thread_queue`)。 +pub(crate) fn enqueue_direct_pending_turn( + thread_id: &str, + turn: PendingDirectTurn, +) -> Result { + let outcome = { + global_direct_thread_manager() + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .enqueue_pending_turn(thread_id, turn) + }; + if matches!(outcome, Ok(EnqueueOutcome::Enqueued)) { + notify_direct_thread_subscribers(thread_id); + } + outcome +} + +/// 取消一条待发消息:结果告诉调用方它是被取消、已经放行,还是本来就不在队里。 +pub(crate) fn remove_direct_pending_turn( + thread_id: &str, + client_turn_id: &str, +) -> DirectQueueRemovalOutcome { + let outcome = { + global_direct_thread_manager() + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .remove_pending_turn(thread_id, client_turn_id, direct_tool_call_now_ms()) + }; + if outcome == DirectQueueRemovalOutcome::Removed { + notify_direct_thread_subscribers(thread_id); + } + outcome +} + +/// 认领队首并登记占用(放行的原子步骤,见 [`DirectThreadManager::claim_pending_turn`])。 +/// +/// `None` = 现在不该放行(已有未收口的回合,或队列为空)。 +pub(crate) fn claim_direct_pending_turn(thread_id: &str) -> Option { + let claimed = { + global_direct_thread_manager() + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .claim_pending_turn(thread_id, direct_tool_call_now_ms()) + }; + if claimed.is_some() { + notify_direct_thread_subscribers(thread_id); + } + claimed +} + /// 运行时回填某一轮逻辑回合的进度(状态 / 活动 / 序号)。返回是否真的写进去了。 pub(crate) fn update_direct_thread_active_turn( thread_id: &str, @@ -700,7 +943,9 @@ pub(crate) fn consume_direct_thread( #[cfg(test)] mod tests { use super::*; - use crate::agent::{DirectThreadDeltaKind, DirectThreadItem, DirectThreadRequestKind}; + use crate::agent::{ + DirectThreadDeltaKind, DirectThreadItem, DirectThreadRequestKind, MAX_PENDING_DIRECT_TURNS, + }; fn message(item_id: &str) -> DirectThreadItem { DirectThreadItem::Message { @@ -1110,4 +1355,206 @@ mod tests { assert_eq!(conflict.as_deref(), Some("turn-1")); } + + fn pending_turn(client_turn_id: &str) -> PendingDirectTurn { + PendingDirectTurn::prepare( + client_turn_id.to_string(), + serde_json::from_value(serde_json::json!({ + "type": "message", + "role": "user", + "content": [{"type": "input_text", "text": format!("消息 {client_turn_id}")}], + "id": format!("direct-codex:{client_turn_id}:user"), + })) + .expect("canonical user item"), + format!("消息 {client_turn_id}"), + None, + FIXED_AT_MS, + ) + .expect("prepare pending turn") + } + + fn enqueued_ids(events: &[DirectThreadEvent]) -> Vec<&str> { + events + .iter() + .filter_map(|event| match event { + DirectThreadEvent::QueueEnqueued { client_turn_id, .. } => { + Some(client_turn_id.as_str()) + } + _ => None, + }) + .collect() + } + + /// 队列就是事件列表:新订阅者的 bootstrap 看得见当前队列,取消之后那条就从队列里消失。 + #[test] + fn new_subscribers_see_the_current_queue_until_it_is_removed() { + let mut manager = DirectThreadManager::with_limits(100, 100_000); + manager + .enqueue_pending_turn("thread-1", pending_turn("turn-1")) + .expect("enqueue head"); + manager + .enqueue_pending_turn("thread-1", pending_turn("turn-2")) + .expect("enqueue tail"); + + let bootstrap = manager.subscribe("thread-1"); + assert_eq!(enqueued_ids(&bootstrap.events), vec!["turn-1", "turn-2"]); + + let watcher = manager.subscribe("thread-1"); + let _ = manager.consume(&watcher.subscription_id).expect("drain"); + assert_eq!( + manager.remove_pending_turn("thread-1", "turn-1", FIXED_AT_MS), + DirectQueueRemovalOutcome::Removed + ); + // 离开队列的事件带 typed 原因:这条是用户取消。 + let events = manager + .consume(&watcher.subscription_id) + .expect("consume removal") + .events; + assert!( + matches!( + events.as_slice(), + [DirectThreadEvent::QueueRemoved { client_turn_id, reason, .. }] + if client_turn_id == "turn-1" && *reason == DirectQueueRemovalReason::Cancelled + ), + "{events:?}" + ); + // 队尾那条仍在,队首那条不再对新订阅者可见。 + let bootstrap = manager.subscribe("thread-1"); + assert_eq!(enqueued_ids(&bootstrap.events), vec!["turn-2"]); + } + + /// 入队幂等:判重范围是「在队 ∪ 正在跑的那一轮」,重复入队不排第二条、不发第二条事件。 + #[test] + fn enqueue_is_idempotent_per_client_turn_id() { + let mut manager = DirectThreadManager::with_limits(100, 100_000); + assert_eq!( + manager + .enqueue_pending_turn("thread-1", pending_turn("turn-1")) + .expect("enqueue"), + EnqueueOutcome::Enqueued + ); + assert_eq!( + manager + .enqueue_pending_turn("thread-1", pending_turn("turn-1")) + .expect("enqueue again"), + EnqueueOutcome::AlreadyKnown + ); + assert_eq!(manager.pending_turn_ids("thread-1"), vec!["turn-1"]); + + // 正在跑的那一轮同样算已知:它的身份已经在占用里,不该再排一条等放行。 + manager + .claim_pending_turn("thread-1", FIXED_AT_MS) + .expect("claim head"); + assert_eq!( + manager + .enqueue_pending_turn("thread-1", pending_turn("turn-1")) + .expect("enqueue while running"), + EnqueueOutcome::AlreadyKnown + ); + assert!(manager.pending_turn_ids("thread-1").is_empty()); + } + + /// 容量只数在队条目:正在跑的那一轮不占等待位。 + #[test] + fn capacity_counts_only_pending_turns() { + let mut manager = DirectThreadManager::with_limits(100, 100_000); + for index in 0..MAX_PENDING_DIRECT_TURNS { + assert_eq!( + manager + .enqueue_pending_turn("thread-1", pending_turn(&format!("turn-{index}"))) + .expect("enqueue"), + EnqueueOutcome::Enqueued + ); + } + assert_eq!( + manager.enqueue_pending_turn("thread-1", pending_turn("turn-overflow")), + Err(EnqueueRejection::QueueFull) + ); + + manager + .claim_pending_turn("thread-1", FIXED_AT_MS) + .expect("claim head"); + assert_eq!( + manager + .enqueue_pending_turn("thread-1", pending_turn("turn-overflow")) + .expect("enqueue after dispatch"), + EnqueueOutcome::Enqueued + ); + } + + /// 放行认领的是**队首**,而且开始事件与离开队列事件同批写:中间没有第二个中间态。 + #[test] + fn dispatch_claims_the_head_and_writes_both_events_in_one_batch() { + let mut manager = DirectThreadManager::with_limits(100, 100_000); + manager + .enqueue_pending_turn("thread-1", pending_turn("turn-1")) + .expect("enqueue head"); + manager + .enqueue_pending_turn("thread-1", pending_turn("turn-2")) + .expect("enqueue tail"); + let watcher = manager.subscribe("thread-1"); + let _ = manager.consume(&watcher.subscription_id).expect("drain"); + + let claimed = manager + .claim_pending_turn("thread-1", 2_000) + .expect("claim head"); + assert_eq!(claimed.pending.client_turn_id, "turn-1"); + assert_eq!(claimed.pending.prompt, "消息 turn-1"); + assert_eq!(manager.pending_turn_ids("thread-1"), vec!["turn-2"]); + + let events = manager + .consume(&watcher.subscription_id) + .expect("consume dispatch") + .events; + assert!( + matches!( + events.as_slice(), + [ + DirectThreadEvent::TurnStarted { at, .. }, + DirectThreadEvent::QueueRemoved { client_turn_id, reason, at: removed_at }, + ] if *at == Some(2_000) + && client_turn_id == "turn-1" + && *reason == DirectQueueRemovalReason::Dispatched + && *removed_at == Some(2_000) + ), + "{events:?}" + ); + + // 已经有未收口的回合:认领不再放行第二条。 + assert!(manager.claim_pending_turn("thread-1", 2_100).is_none()); + assert_eq!(manager.pending_turn_ids("thread-1"), vec!["turn-2"]); + } + + /// 取消的三种结果分得开:真的移除了、已经被放行了、本来就不在队里。 + #[test] + fn removal_tells_cancelled_dispatched_and_unknown_apart() { + let mut manager = DirectThreadManager::with_limits(100, 100_000); + assert_eq!( + manager.remove_pending_turn("thread-1", "turn-1", FIXED_AT_MS), + DirectQueueRemovalOutcome::NotFound + ); + + manager + .enqueue_pending_turn("thread-1", pending_turn("turn-1")) + .expect("enqueue"); + assert_eq!( + manager.remove_pending_turn("thread-1", "turn-1", FIXED_AT_MS), + DirectQueueRemovalOutcome::Removed + ); + assert_eq!( + manager.remove_pending_turn("thread-1", "turn-1", FIXED_AT_MS), + DirectQueueRemovalOutcome::NotFound + ); + + manager + .enqueue_pending_turn("thread-1", pending_turn("turn-2")) + .expect("enqueue"); + manager + .claim_pending_turn("thread-1", FIXED_AT_MS) + .expect("claim head"); + assert_eq!( + manager.remove_pending_turn("thread-1", "turn-2", FIXED_AT_MS), + DirectQueueRemovalOutcome::AlreadyDispatched + ); + } }