From d7c7089fa65cf19592879a61a47468aa657d4998 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=8E=8B=E5=BE=B7=E5=AE=87?= Date: Tue, 15 Sep 2026 20:04:34 +0800 Subject: [PATCH] =?UTF-8?q?=E9=87=8A=E6=94=BE=E5=B7=B2=E8=A7=A3=E5=86=B3?= =?UTF-8?q?=E7=9A=84=20DirectThread=20=E8=AF=B7=E6=B1=82=E4=BA=8B=E4=BB=B6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 请求无 ID 时立即可清理,解决事件到达后解除对应 requested 事件 pin。 --- .../src/agent/direct_thread_manager.rs | 56 ++++++++++++++++++- 1 file changed, 53 insertions(+), 3 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 a2465bab7..9f5adf8f9 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 @@ -152,6 +152,12 @@ impl DirectThreadManager { cleanable, }); Self::mark_item_events_cleanable(thread, event.item_id.as_deref()); + if matches!( + event.event_type.as_str(), + "approval.resolved" | "request.resolved" | "ask.resolved" + ) { + Self::mark_request_events_cleanable(thread, request_id(&event).as_deref()); + } self.evict(thread_id); event } @@ -266,10 +272,11 @@ impl DirectThreadManager { true } "approval.requested" | "request.requested" | "ask.requested" => { - if let Some(request_id) = request_id(event) { - thread.unresolved_requests.insert(request_id); + let request_id = request_id(event); + if let Some(request_id) = request_id.as_deref() { + thread.unresolved_requests.insert(request_id.to_string()); } - false + request_id.is_none() } "approval.resolved" | "request.resolved" | "ask.resolved" => { if let Some(request_id) = request_id(event) { @@ -316,6 +323,21 @@ impl DirectThreadManager { } } + fn mark_request_events_cleanable(thread: &mut ThreadState, resolved_request_id: Option<&str>) { + let Some(resolved_request_id) = resolved_request_id else { + return; + }; + for stored in &mut thread.events { + if matches!( + stored.event.event_type.as_str(), + "approval.requested" | "request.requested" | "ask.requested" + ) && request_id(&stored.event).as_deref() == Some(resolved_request_id) + { + stored.cleanable = true; + } + } + } + fn trim_prefix(thread: &mut ThreadState) { let min_cursor = thread .subscribers @@ -583,6 +605,34 @@ mod tests { ); } + #[test] + fn resolved_request_releases_the_original_requested_event() { + let mut manager = DirectThreadManager::with_limits(100, 100_000); + manager.append( + "thread-1", + DirectThreadRawEventDraft { + event_type: "approval.requested".to_string(), + turn_id: "turn-1".to_string(), + item_id: None, + payload: serde_json::json!({"requestId": "request-1"}), + }, + ); + let subscription = manager.subscribe("thread-1"); + manager.append( + "thread-1", + DirectThreadRawEventDraft { + event_type: "approval.resolved".to_string(), + turn_id: "turn-1".to_string(), + item_id: None, + payload: serde_json::json!({"requestId": "request-1"}), + }, + ); + manager + .consume(&subscription.subscription_id) + .expect("consume resolution"); + assert_eq!(manager.thread_debug("thread-1").unwrap().0, 0); + } + #[test] fn turn_completed_anchor_survives_empty_queue_for_new_subscriber() { let mut manager = DirectThreadManager::with_limits(100, 100_000);