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::*;