From b18740405c745897a992d2f730248ece94884e9e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=BB=9D=E5=B0=98?= <237809796@qq.com> Date: Thu, 2 Jul 2026 15:47:37 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E5=A4=8D:=20run=5Fcommand=20=E7=A9=BA?= =?UTF-8?q?=E8=B7=AF=E5=BE=84=E6=8A=A5=E9=94=99=20+=20=E9=9A=A7=E9=81=93?= =?UTF-8?q?=E6=97=A5=E5=BF=97=E6=B4=AA=E6=B0=B4=20+=20=E5=90=AF=E5=8A=A8?= =?UTF-8?q?=E6=AE=8B=E7=95=99=E6=B8=85=E7=90=86?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - tool_registry: run_command working_dir 空字符串 "" → None,修 Windows os error 123 - lib: tunnel subscriber 从频率压制(suppress_until)改为连接状态感知(is_connected), 断开时静默丢弃事件,重连后恢复透传,记 INFO 状态变迁(治本) - conversation_repo: 新增 cleanup_stale_pending(),超 24h 残留 pending 标记 interrupted - restore: 启动时先清理过期 pending 再恢复审批 - state: 实现 cleanup_orphan_pending_messages(),启动时清理对应已决/超时 tool 的 __PENDING__ 占位消息 --- .../df-storage/src/crud/conversation_repo.rs | 19 ++++++++++++++ src-tauri/src/commands/ai/audit/restore.rs | 5 ++++ src-tauri/src/commands/ai/tool_registry.rs | 3 ++- src-tauri/src/lib.rs | 18 ++++++++++++- src-tauri/src/state.rs | 25 +++++++++++++++++++ 5 files changed, 68 insertions(+), 2 deletions(-) diff --git a/crates/df-storage/src/crud/conversation_repo.rs b/crates/df-storage/src/crud/conversation_repo.rs index a96f48c..0999bc4 100644 --- a/crates/df-storage/src/crud/conversation_repo.rs +++ b/crates/df-storage/src/crud/conversation_repo.rs @@ -264,6 +264,25 @@ impl AiToolExecutionRepo { .map_err(storage_err)? } + /// 清理超期的残留 pending 工具调用(旧会话遗留)。 + /// + /// `max_age_secs`: 超过此秒数的 pending 记录被标记为 interrupted(不硬删,保留审计痕迹)。 + pub async fn cleanup_stale_pending(&self, max_age_secs: u64) -> Result { + let conn = self.conn.clone(); + let cutoff_ms = (df_types::now_millis() as i64 - (max_age_secs as i64 * 1000)).to_string(); + let affected = tokio::task::spawn_blocking(move || { + let guard = conn.blocking_lock(); + guard.execute( + "UPDATE ai_tool_executions SET status = 'interrupted' \ + WHERE status = 'pending' AND CAST(requested_at AS INTEGER) < ?1", + params![cutoff_ms], + ).map_err(storage_err) + }) + .await + .map_err(storage_err)??; + Ok(affected as u64) + } + /// 审批历史面板分页查询:按 requested_at 倒序(最新在前),limit 默认 50。 /// /// 与 list_pending 同理走专用 SELECT,绕过通用 query 宏(后者硬编码 diff --git a/src-tauri/src/commands/ai/audit/restore.rs b/src-tauri/src/commands/ai/audit/restore.rs index 81f05e3..a512fd7 100644 --- a/src-tauri/src/commands/ai/audit/restore.rs +++ b/src-tauri/src/commands/ai/audit/restore.rs @@ -24,6 +24,11 @@ use super::risk_from_str; /// 语义而非路由键)。conversation_id=None 的无主审批(R-9)不建 per_conv(无 conv_id 可挂), /// 仍进 pending_approvals 单层表,后续审批按 tool_call_id 路由,不影响正确性。 pub async fn restore_pending_approvals(state: &AppState) { + // 清理超 24 小时的残留 pending(旧会话遗留,不再有意义) + if let Err(e) = state.ai_tool_executions.cleanup_stale_pending(86400).await { + tracing::warn!("清理过期 pending 审批失败(非阻断): {}", e); + } + let pending = match state.ai_tool_executions.list_pending().await { Ok(v) => v, Err(e) => { diff --git a/src-tauri/src/commands/ai/tool_registry.rs b/src-tauri/src/commands/ai/tool_registry.rs index 05fad27..b2557e2 100644 --- a/src-tauri/src/commands/ai/tool_registry.rs +++ b/src-tauri/src/commands/ai/tool_registry.rs @@ -2680,7 +2680,8 @@ fn register_file_tools( let request = ShellRequest { command: command.clone(), - working_dir: Some(working_dir.clone()), + // 空字符串 → None(空路径是非法 current_dir,Windows 报 os error 123) + working_dir: if working_dir.is_empty() { None } else { Some(working_dir.clone()) }, env: HashMap::new(), timeout_secs: Some(timeout_secs), shell_type: Default::default(), diff --git a/src-tauri/src/lib.rs b/src-tauri/src/lib.rs index b75c3db..56cb528 100644 --- a/src-tauri/src/lib.rs +++ b/src-tauri/src/lib.rs @@ -186,13 +186,29 @@ pub fn run() { let token = "devflow-relay-default-token".to_string(); // subscriber task:subscribe ai_event_bus → tunnel.send_raw_event 透传 miniapp + // 治本:先查 is_connected()(AtomicBool 无锁读),未连接时静默丢弃事件, + // 避免每次 send_raw_event 都失败并记 WARN(旧实现用 suppress_until 降频,是治标)。 let mut rx = state.ai_event_bus.subscribe(); let tunnel_for_sub = state.tunnel.clone(); tauri::async_runtime::spawn(async move { tracing::info!("[tunnel-sub] subscriber task 启动,透传 ai_event_bus → relay"); + let mut was_connected = false; while let Ok(value) = rx.recv().await { + if !tunnel_for_sub.is_connected() { + if was_connected { + tracing::info!("[tunnel-sub] tunnel 已断开,暂停透传"); + was_connected = false; + } + continue; + } + if !was_connected { + tracing::info!("[tunnel-sub] tunnel 已重连,恢复透传"); + was_connected = true; + } if let Err(e) = tunnel_for_sub.send_raw_event(value).await { - tracing::warn!("[tunnel-sub] send_raw_event 失败(可能未连接): {}", e); + // 连接刚断(查询与发送间窗口),记一条 DEBUG 而非 WARN + tracing::debug!("[tunnel-sub] send_raw_event 失败(连接瞬断): {}", e); + was_connected = false; } } tracing::info!("[tunnel-sub] subscriber task 退出(rx 关闭)"); diff --git a/src-tauri/src/state.rs b/src-tauri/src/state.rs index 5c28e4b..18e8d05 100644 --- a/src-tauri/src/state.rs +++ b/src-tauri/src/state.rs @@ -264,6 +264,11 @@ impl AppState { // 从 Settings KV 恢复审批超时配置(默认 15 分钟,0=禁用超时)。 state.reload_approval_timeout().await; + // 启动时清理残留 __PENDING__ 占位消息(对应 tool_execution 已非 pending 状态) + if let Err(e) = cleanup_orphan_pending_messages(&state.db).await { + tracing::warn!("清理孤儿 __PENDING__ 消息失败(非阻断): {}", e); + } + // 迁移旧 .trash(编译期 workspace_root → 运行期 data_dir),仅一次,幂等。 let old_trash = workspace_root_path().join(".trash"); let new_trash = data_dir.join(".trash"); @@ -554,6 +559,26 @@ fn build_registry(db: Arc) -> NodeRegistry { registry } +/// 启动时清理孤儿 `__PENDING__` 占位消息。 +/// +/// 对应 `tool_execution` 已非 pending(已审批/已超时/已中断)的占位消息是残留垃圾, +/// 不删会导致前端展示冗余的 pending 态工具卡。 +pub async fn cleanup_orphan_pending_messages(db: &Database) -> anyhow::Result { + let conn_arc = db.conn(); + let guard = conn_arc.lock().await; + let deleted = guard.execute( + "DELETE FROM ai_messages \ + WHERE (content = '__PENDING__' OR content LIKE '%__PENDING__%') \ + AND tool_call_id IS NOT NULL \ + AND tool_call_id NOT IN (\ + SELECT tool_call_id FROM ai_tool_executions WHERE status = 'pending'\ + )", + [], + )?; + tracing::info!("清理孤儿 __PENDING__ 消息: {} 条", deleted); + Ok(deleted as u64) +} + #[cfg(test)] mod tests { use super::*;