//! 应用全局状态 — 数据库、Repo、事件总线、节点注册表、AI 会话 use std::collections::{HashMap, HashSet}; use std::path::{Path, PathBuf}; use std::sync::Arc; use std::sync::atomic::AtomicUsize; use anyhow::Result; use serde::{Deserialize, Serialize}; use tokio::sync::{Mutex, RwLock, Semaphore}; use df_ai::ai_tools::AiToolRegistry; use df_storage::crud::{ AiConversationRepo, AiMessageRepo, AiProviderRepo, AiToolExecutionRepo, IdeaEvalRepo, IdeaRepo, KnowledgeEventsRepo, KnowledgeRepo, NodeExecutionRepo, ProjectEventRepo, ProjectRepo, ProjectServiceRepo, ReleaseRepo, SettingsRepo, TaskLinkRepo, TaskRepo, WorkflowRepo, }; use df_storage::db::Database; use df_workflow::eventbus::EventBus; use df_workflow::registry::NodeRegistry; use df_workflow::state::StateMachine; use crate::commands::ai::augmentation::ResolverRegistry; use crate::commands::ai::AiSession; // ============================================================ // 知识库配置(提取 + 注入) // ============================================================ /// AI 提炼触发方式 #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] #[serde(rename_all = "snake_case")] pub enum ExtractTrigger { /// 对话正常完成时(默认) OnComplete, /// 仅手动按钮触发 ManualOnly, } /// 知识库配置持久化 KV key(P0 设置走查-2026-06-21:原纯内存 Arc 启动 default 覆盖 /// 致 8 项配置重启全丢;save_config 落此 KV,init reload_knowledge_config 恢复)。 pub const KNOWLEDGE_CONFIG_KEY: &str = "df-knowledge-config"; /// 知识库行为配置(内存真相源 + Settings KV 持久化,前后端通过 IPC 读写) #[derive(Debug, Clone, Serialize, Deserialize)] pub struct KnowledgeConfig { /// 提炼总开关,默认 true pub auto_extract: bool, /// 提炼触发方式,默认 OnComplete pub trigger_mode: ExtractTrigger, /// 最少消息数守卫(防闲聊噪音),默认 4 pub min_messages: u32, /// 聊天时自动注入相关知识开关,默认 true pub auto_inject: bool, /// 语义检索(向量)总开关,默认 false——关闭时纯 LIKE 零外部依赖 #[serde(default)] pub vector_enabled: bool, /// embedding 用的 provider id(仅 openai_compat 类型,Anthropic 无 embed API) #[serde(default)] pub embedding_provider_id: Option, /// embedding 模型名(如 embedding-3 / text-embedding-3-small) #[serde(default)] pub embedding_model: Option, } impl Default for KnowledgeConfig { fn default() -> Self { Self { auto_extract: true, trigger_mode: ExtractTrigger::OnComplete, min_messages: 4, auto_inject: true, vector_enabled: false, embedding_provider_id: None, embedding_model: None, } } } // ============================================================ // LLM 调用并发控制(双层 Semaphore + 可选 per-provider 层) // ============================================================ /// LLM 调用并发控制 — 全局 LLM 调用限流 global + 单对话内 per_conv 双层 Semaphore + 可选 per-provider 层 /// /// 限流对象:所有真实 LLM 调用(主循环 stream_llm / 标题生成 / 知识提炼)。 /// 不限流本地工具执行(tools.execute)——本地操作无外部成本、不受 RPM 约束。 /// /// ## global(permits 默认 3,「LLM 调用并发限流」原义) /// F-09 batch5 曾把 global 改「并发会话数上限」(loop 入口持整 loop),用户决策修正「不设并发会话上限」后 /// 已删除 loop 入口 acquire。global 回归原义:由各单次 LLM 调用点(stream_llm 重试循环 / 标题 / 压缩 / /// 提炼 / 项目分析扫描)各自 `acquire_global()` 拿 permit、调用结束 Drop 释放,防 provider 429。 /// 多对话 loop 并发不限数(用户接受 token 暴增)。 /// /// ## per_conv 改 HashMap /// 原应用级单信号量(AiSession 单例 + generating 互斥下退化)改为 /// `HashMap>`:每对话一份(permits=2,主循环 + 标题 + 提炼各自限流)。 /// `run_agentic_loop` 入口 `acquire_per_conv(conv_id)` 拿 1 permit 持整个 loop 生命周期(含工具执行/审批 /// 等待/重试),防单对话内并发 LLM 调用失控。这是单对话内限流,非会话数限制,符合用户「不设上限」。 /// `acquire_per_conv(conv_id)`:lock HashMap → 无则建(permits=2)→ clone Arc → 释放 lock → acquire_owned。 /// conv 退出清理(`release_conv(conv_id)`):conv 删除时 remove 条目(防 HashMap 无限增长)。 /// /// 运行时调整:tokio Semaphore 的 permits 数构造时固定、不可增减, /// 故 `global` 用 `Arc>>` 双层包装——替换内层 Arc 即重建 Semaphore。 /// 已持有旧 permit 的任务不受影响(permit 绑定旧 Semaphore,软收敛), /// 新请求 lock 后克隆到最新 Arc、自动走新限制。旧 Semaphore 随最后 permit 释放而 drop。 /// per_conv HashMap 内每条 Arc 不需运行时改 permits(对话内并发上限 2 固定),故无重建需求。 /// /// ## F-260614-04: per-provider 层(可选) /// `per_provider` 为 HashMap>。调用方(agentic loop)经 /// `acquire_for_provider(pid)` 取额外 permit,防单 provider 被打满(限流 429)。 /// **单 provider 场景**:若未调 `set_provider_caps`,HashMap 为空, /// `acquire_for_provider` 返回 None(无限流,行为同 F-01 前)。零变化保证。 /// 全局容量 = min(sum(各 provider 上限), global_cap):由调用方在配置时约束 /// (set_provider_caps 传 min(sum, global_cap)),非运行时强约束。 /// /// ## Phase3 预留: per_sub_flow 层(批2-B,占位未接入) /// `per_sub_flow` 为 HashMap>,用于 Phase3 单对话并行多轮的 /// **子流并发上限**(单对话内并行 spawn 多条子流时,防子流数失控)。本批仅加层 + 占位, /// **不接入调用点**(接入待 Phase3 子流 spawn 落地)。 /// /// **key 约定(文档化,非强约束)**:`sub_flow_id` 采用 `"{conv_id}::{sub_id}"` 格式, /// 唯一标识某对话下的某条子流,便于 release_sub_flow 精确清理。 /// /// **三层语义**: /// - per_sub_flow = 单对话内**子流**并发上限(每条子流一份 Semaphore,permits 默认 3,预留) /// - per_conv = 单对话内**总**并发上限(permits 默认 2) /// - global = 全应用**总**并发上限(permits 默认 3) /// /// **理想配额(用户配置建议,非硬限)**: per_sub_permits ≤ per_conv_permits ≤ global_permits。 /// 实际限流由各层 Semaphore 自然保证:global Semaphore 硬限全应用并发不超 global permits /// (跨对话 sum(per_conv) 可 > global,但 global acquire_owned 自然阻塞,不超卖)。 /// /// **acquire 语义(非阻塞)**:`acquire_per_sub_flow` 用 `try_acquire_owned`(非阻塞), /// 耗尽返 None 降级串行——子流并行不阻塞主 loop(Phase3 子流 spawn 失败不拖垮整对话)。 #[derive(Clone)] pub struct LlmConcurrency { /// 全局 LLM 调用并发上限(permits 默认 3):F-09 batch5 修正后回归原义,由各单次 LLM 调用点 /// (stream_llm/标题/压缩/提炼/项目分析)各自 acquire/drop 防 429,不再由 loop 入口持整 loop。 global: Arc>>, /// 单对话内并发上限(F-09 B 批5):HashMap,每对话 permits=2。 /// acquire 时按 conv_id 取/建;release_conv 在 loop 结束 + 无 pending 时 remove。 per_conv: Arc>>>, /// 每对话内并发上限(permits 默认 2,构造时传入 new(3, 2) 的第二参)。 /// F-09 B 批5:AtomicUsize 支持运行时热改(ai_set_concurrency_config 的 per_conv_limit)。 /// 热改后**已建对话**的旧 Semaphore 不变(permits 构造时固定),**新建对话**用新值; /// 为使热改立即全量生效,set_per_conv 同时清空 HashMap 强制重建(软收敛:旧 permit 随 Drop 释放)。 per_conv_permits: Arc, /// F-260614-04: per-provider 信号量表。空 = 无 per-provider 限流(单 provider 路径零变化)。 /// Arc>:运行时增删 provider 配置时替换/插入,acquire 时 clone Arc。 per_provider: Arc>>>, /// Phase3 预留(批2-B): per-sub_flow 信号量表。key = sub_flow_id(约定 "{conv_id}::{sub_id}")。 /// 空 = 无子流并发限流(Phase3 未接入时零变化)。acquire 用 try_acquire_owned(非阻塞), /// 耗尽返 None 降级串行,子流并行不阻塞主 loop。 #[allow(dead_code)] // Phase3 子流 spawn 落地时接入调用点,本批仅占位。 per_sub_flow: Arc>>>, /// Phase3 预留(批2-B): 每子流并发上限(permits 默认 3,内部初始化,预留)。 /// AtomicUsize 支持运行时热改;set_per_sub_permits 同时清空 HashMap(软收敛,对齐 set_per_conv)。 /// new(global, per_conv) 签名保持向后兼容,per_sub_permits 内部 AtomicUsize::new(3)。 #[allow(dead_code)] // Phase3 子流 spawn 落地时接入调用点,本批仅占位。 per_sub_permits: Arc, } impl LlmConcurrency { pub fn new(global: usize, per_conv: usize) -> Self { Self { global: Arc::new(Mutex::new(Arc::new(Semaphore::new(global)))), per_conv: Arc::new(Mutex::new(HashMap::new())), per_conv_permits: Arc::new(AtomicUsize::new(per_conv)), per_provider: Arc::new(Mutex::new(HashMap::new())), // Phase3 预留(批2-B): per_sub_flow 默认 permits=3,内部初始化。 // 签名保持 new(global, per_conv) 向后兼容;Phase3 接入后用户可经 // set_per_sub_permits 热改(配置化待后续批,本批仅占位)。 per_sub_flow: Arc::new(Mutex::new(HashMap::new())), per_sub_permits: Arc::new(AtomicUsize::new(3)), } } /// 取全局 LLM 调用并发 permit:由各单次 LLM 调用点(stream_llm/标题/压缩/提炼/项目分析) /// 各自调用、调用结束 Drop 释放。F-09 batch5 修正后不再由 run_agentic_loop 入口持整 loop。 /// 重建后(set_global)新请求自动走最新 Semaphore。 pub async fn acquire_global(&self) -> tokio::sync::OwnedSemaphorePermit { let sema = self.global.lock().await.clone(); sema.acquire_owned().await.expect("llm global semaphore closed") } /// 取单对话内并发 permit(F-09 B 批5):按 conv_id 取/建 Semaphore,permits=当前 per_conv_permits(默认 2)。 /// lock HashMap → 无则建 → clone Arc → 释放 lock → acquire_owned。锁持有短(不含 await acquire)。 pub async fn acquire_per_conv(&self, conv_id: &str) -> tokio::sync::OwnedSemaphorePermit { let sema = { let mut map = self.per_conv.lock().await; let permits = self.per_conv_permits.load(std::sync::atomic::Ordering::SeqCst); map.entry(conv_id.to_string()) .or_insert_with(|| Arc::new(Semaphore::new(permits))) .clone() }; sema.acquire_owned().await.expect("llm per_conv semaphore closed") } /// F-09 B 批5:conv 退出清理。loop 结束 + 无 pending 审批时 remove 该 conv 的 Semaphore 条目。 /// **时机由调用方判断**(agentic loop 退出点):仅在确信无后续 acquire 时调用,否则误删会致 /// 该 conv 下次 acquire 重建 Semaphore(限流计数清零,非致命,但语义偏离)。 /// remove 后已持 permit 不受影响(permit 绑旧 Arc,随 Drop 释放),仅阻止新条目累积。 pub async fn release_conv(&self, conv_id: &str) { self.per_conv.lock().await.remove(conv_id); } /// 重建全局 Semaphore(软收敛:旧 permit 不回收,待其释放后新限制完全生效) pub async fn set_global(&self, permits: usize) { *self.global.lock().await = Arc::new(Semaphore::new(permits)); } /// 重建单对话 Semaphore(F-09 B 批5 后 per_conv 为 HashMap)。 /// 行为:更新 per_conv_permits(AtomicUsize) + 清空 HashMap(软收敛:旧 permit 随 Drop 释放, /// 新对话 acquire 用新 permits 值重建 Semaphore)。已建对话若仍在跑,旧 Semaphore 不变; /// 下次该 conv 新 acquire 时因 HashMap 已清空会重建为新 permits。 pub async fn set_per_conv(&self, permits: usize) { self.per_conv_permits .store(permits, std::sync::atomic::Ordering::SeqCst); self.per_conv.lock().await.clear(); } // ============================================================ // Phase3 预留(批2-B): per_sub_flow 层 — 占位未接入调用点 // ============================================================ /// Phase3 预留(批2-B): 取单子流内并发 permit(非阻塞)。 /// /// 按 `sub_flow_id` 取/建 Semaphore(permits=当前 per_sub_permits,默认 3)。 /// lock HashMap → 无则建(or_insert_with, permits 取 per_sub_permits 当前值) /// → clone Arc → 释放 lock → **try_acquire_owned**(非阻塞)。 /// /// **非阻塞语义(关键)**:用 try_acquire_owned 而非 acquire_owned——耗尽返 None /// 调用方降级串行处理,子流并行不阻塞主 loop(Phase3 子流 spawn 失败不拖垮整对话)。 /// /// - 返 `Ok(Some(permit))`:成功获取,permit 随 Drop 自动释放。 /// - 返 `Ok(None)`:子流并发已耗尽,调用方降级串行(非错误)。 /// /// **key 约定**:`sub_flow_id` 采用 `"{conv_id}::{sub_id}"` 格式(文档化,调用方组装)。 /// 锁持有短(不含 try_acquire),对齐 acquire_per_conv 模式。 #[allow(dead_code)] // Phase3 子流 spawn 落地时接入调用点,本批仅占位。 pub async fn acquire_per_sub_flow( &self, sub_flow_id: &str, ) -> Option { let sema = { let mut map = self.per_sub_flow.lock().await; let permits = self.per_sub_permits.load(std::sync::atomic::Ordering::SeqCst); map.entry(sub_flow_id.to_string()) .or_insert_with(|| Arc::new(Semaphore::new(permits))) .clone() }; // try_acquire_owned 非阻塞:Err(NoPermits) 返 None 降级串行,不阻塞主 loop。 // 对齐 acquire_global/per_conv 的 expect 前提:Semaphore 不会 close(无 close() 调用)。 sema.try_acquire_owned().ok() } /// Phase3 预留(批2-B): 子流完成清理。remove 该 sub_flow 的 Semaphore 条目(防 HashMap 无限增长)。 /// /// **对齐 release_conv 模式**:remove 后已持 permit 不受影响(permit 绑旧 Arc,随 Drop 释放), /// 仅阻止新条目累积。时机由调用方判断(Phase3 子流退出点)。 #[allow(dead_code)] // Phase3 子流 spawn 落地时接入调用点,本批仅占位。 pub async fn release_sub_flow(&self, sub_flow_id: &str) { self.per_sub_flow.lock().await.remove(sub_flow_id); } /// Phase3 预留(批2-B): 热改 per_sub_flow permits 上限。 /// /// 行为(对齐 set_per_conv):更新 per_sub_permits(AtomicUsize) + 清空 HashMap(软收敛: /// 旧 permit 随 Drop 释放,新子流 acquire 用新 permits 值重建 Semaphore)。已建子流若仍在跑, /// 旧 Semaphore 不变;下次该 sub_flow 新 acquire 时因 HashMap 已清空会重建为新 permits。 #[allow(dead_code)] // Phase3 子流 spawn 落地时接入调用点,本批仅占位。 pub async fn set_per_sub_permits(&self, permits: usize) { self.per_sub_permits .store(permits, std::sync::atomic::Ordering::SeqCst); self.per_sub_flow.lock().await.clear(); } /// F-260614-04: 取 per-provider 并发 permit(可选)。 /// /// - provider 在 `per_provider` 表中有配置 → 取其 Semaphore permit,返回 Some。 /// - provider 无配置(单 provider 场景或未 set_provider_caps)→ 返回 None,无限流。 /// /// 调用方(agentic loop)用法: /// ```ignore /// let _global_permit = llm_concurrency.acquire_global().await; /// let _per_conv_permit = llm_concurrency.acquire_per_conv(&conv_id).await; /// let _provider_permit = llm_concurrency.acquire_for_provider(&provider_id).await; /// ``` /// 三 permit 均绑 guard Drop 自动释放。None 时无 permit 需释放(行为同 F-01 前)。 pub async fn acquire_for_provider( &self, provider_id: &str, ) -> Option { let sema = { let map = self.per_provider.lock().await; map.get(provider_id).cloned() }?; // Semaphore 存在 → acquire。expect 同 acquire_global/per_conv:Semaphore 不会 close // (无 close() 调用,仅在 set_provider_caps 时替换表内 Arc,旧 Arc permit 仍有效)。 Some( sema.acquire_owned() .await .expect("llm per_provider semaphore closed"), ) } /// F-260614-04: 批量设置 per-provider 并发上限(替换整表)。 /// /// 调用方(Settings 配置热改 / 启动初始化)传入 `{ provider_id: permits }` map, /// 替换整张 per_provider 表(软收敛:已持 permit 不回收)。全局容量约束 min(sum, global_cap) /// 由调用方在构造 map 时应用(本函数不做强约束,只落表)。 /// /// 传空 map → 清空 per_provider 表(所有 provider 回退到无 per-provider 限流)。 pub async fn set_provider_caps(&self, caps: HashMap) { let mut map = self.per_provider.lock().await; map.clear(); for (pid, permits) in caps { // permits=0 等同无配置(Semaphore::new(0) 永远 acquire 不到 → 死锁), // 故 permits=0 跳过(不落表 → acquire_for_provider 返 None → 无限流)。 if permits > 0 { map.insert(pid, Arc::new(Semaphore::new(permits))); } } } } /// 应用全局状态 — 通过 `app.manage()` 注入,command 中以 `State<'_, AppState>` 取用 pub struct AppState { /// 数据库句柄(Arc 包装,便于在异步任务中重建 Repo) pub db: Arc, /// 灵感表 Repo pub ideas: IdeaRepo, /// 灵感评估历史表 Repo(追加型审计表,每次 AI 评估追加一行快照,可追溯历史) pub idea_evaluations: IdeaEvalRepo, /// 项目表 Repo pub projects: ProjectRepo, /// 任务表 Repo pub tasks: TaskRepo, /// 任务横向关联表 Repo(知识图谱 Phase 1 V29,AI 拓扑排序编排调度的基础) pub task_links: TaskLinkRepo, /// 项目统一事件流表 Repo(知识图谱 Phase 2 V30,project_events 追加型审计表)。 /// 埋点 best-effort:事件写入失败不阻断主操作(设计 §10.1),commands 层 hook/after 追加。 pub project_events: ProjectEventRepo, /// 项目基础设施配置表 Repo(知识图谱 Phase 3 V31,project_services)。 /// 记录项目依赖的基础设施(数据库/缓存/MQ/API 等),为 AI 执行任务时提供基础设施上下文 /// (设计 §2.3 + §五 add_project_service/list_project_services)。 /// D10 安全边界:不存敏感凭证,凭证走环境变量(insert/update_validated 应用层审查拒绝)。 pub project_services: ProjectServiceRepo, /// 发布表 Repo(预留:ReleaseRepo 持久化已就位·IPC/逻辑未接入·SW-260618-21 b 保留) #[allow(dead_code)] pub releases: ReleaseRepo, /// 工作流执行表 Repo pub workflows: WorkflowRepo, /// 节点执行表 Repo(预留:NodeExecutionRepo 持久化已就位·IPC/逻辑未接入·SW-260618-21 b 保留) #[allow(dead_code)] pub node_executions: NodeExecutionRepo, /// 工作流事件总线(tokio broadcast) pub event_bus: EventBus, /// 节点注册表(已注册内置节点) pub registry: Arc, // ── AI ── /// AI 提供商配置 Repo pub ai_providers: AiProviderRepo, /// AI 对话历史 Repo pub ai_conversations: AiConversationRepo, /// AI 消息 Repo(F-260619-03 消息拆分存储:ai_messages 表,读路径批次 A 切读用) pub ai_messages: AiMessageRepo, /// AI 工具执行审计 Repo pub ai_tool_executions: AiToolExecutionRepo, /// AI 工具注册表 pub ai_tools: Arc, /// Input Augmentation 层 Mention Resolver 注册表(核心设计2): /// 按 MentionRef kind 分发(project/task/idea/skill),resolve_all 单条失败不阻断整批。 /// 启动期 build 一次性注册四 resolver(持 db Arc 访问 repo),通过 Arc 共享只读。 pub resolvers: Arc, /// AI 会话状态 pub ai_session: Arc>, /// L3 AI 事件总线(pub-sub 域,独立于工作流 event_bus)。 /// /// 批2b:总线入 AppState 就绪(基建就位)。emit 点双写 + 前端订阅**留批3**: /// 现阶段无后端订阅者(工作流独立域)+ 前端经 ai-chat-event 已收 → 全双写空转, /// 待真实消费者(跨端透传/多模块订阅)出现再针对性双写(渐进接入,机制优先)。 /// /// dead_code:批2b 仅入 AppState,emit 点未双写(publish 无调用方); /// 标 allow 保留作批3 接入(零调用方≠垃圾,预留保留)。 #[allow(dead_code)] pub ai_event_bus: crate::commands::ai::event_bus::EventBus, // ── 知识库 ── /// 知识库 Repo pub knowledge: KnowledgeRepo, /// 知识生命线事件 Repo(产生/审核/引用/归档审计) pub knowledge_events: KnowledgeEventsRepo, /// 知识库行为配置(提取 + 注入) pub knowledge_config: Arc>, /// 应用数据目录(运行期确定:Windows $APPDATA/devflow, macOS ~/Library/Application Support/devflow) pub data_dir: PathBuf, /// 通用应用设置 KV Repo(前端 localStorage 迁移目标) pub settings: SettingsRepo, // ── LLM 并发控制 ── /// LLM 调用并发上限(全局 + 单对话双层 Semaphore,运行时可调) pub llm_concurrency: LlmConcurrency, // ── Agentic 循环轮次上限 ── /// Agentic 循环最大轮次(前端 Settings 数字配置 → AppState 字段 → 热改 command → /// loop 入口 load 快照透传形参;当前 loop 锁定边界,热改下次发消息生效)。 /// 默认 10,与 agentic.rs::DEFAULT_MAX_AGENT_ITERATIONS 对齐。 pub agent_max_iterations: Arc, // ── 流式对话失败自动重试 ── /// 流式对话失败自动重试次数(F-260616-07 / 决策 a1:只重试流前失败 Init Err——未输出 /// 任何 token;流中途失败 MidStream Partial 保文不重试。退避复用 retry::backoff_delay + /// is_status_retryable Fatal 分类 + 30s 总预算,详见 agentic.rs 重试循环。默认 3)。 pub agent_max_retries: Arc, // ── 工作流执行状态 ── /// 工作流执行 → 节点状态机注册表 /// /// run_workflow 创建执行器后注册其 state_machine(StateMachine 内部 Arc, /// clone 共享底层 HashMap,与下沉到 NodeContext.node_status 的是同一份); /// cancel_workflow_node IPC 经 execution_id 取出后 set_cancelled, /// 直达运行中 HumanNode 的 is_cancelled 轮询。执行完成(成功/失败)后移除条目。 pub workflow_state_registry: Arc>>, // ── F-260619-03 Phase A: AI 工具文件访问授权目录白名单 ── /// AI 文件工具(read/write/list/patch/...)可访问的授权目录池。 /// Phase A 仅持久化白名单(Settings KV `allowed_dirs` 加载), /// workspace_root 始终在白名单(向后兼容)。 /// Phase B 将扩展为 persistent + session(会话临时授权) + 弹窗挂起。 pub allowed_dirs: Arc>, // ── 跨端隧道(Phase3 Layer1)── /// 桌面端出站 WS 隧道客户端,连 df-relay(/ws/device)透传 ai_event_bus → miniapp /// + 转发 miniapp 下行指令 → remote_bridge。lib.rs setup spawn task 解析 /// device_id/relay_url 后 connect(失败非阻断,可后续手动重连)。 /// Arc 包裹:setup task(subscriber + on_command 回调)与 AppState 共享同一客户端。 pub tunnel: std::sync::Arc, } /// F-260619-03 Phase A/B/C: AI 工具文件访问授权目录白名单 /// /// - Phase A:`persistent`(持久化白名单,从 Settings KV `allowed_dirs` 加载, /// JSON 数组 `["E:/wk-lab/u-abc"]`)。`resolve_workspace_path` 校验时: /// - 任一 persistent 或 session 目录 starts_with 命中即放行 /// - 无任何授权时全部拒绝(引导用户绑定项目或授权目录) /// - Phase B:`session`(进程级会话临时授权,弹窗"仅本次"写入;切换/新建/删除会话清空)。 /// 单用户桌面应用 active_conversation_id 单全局模型,session 字段随 active 切换清空, /// 行为等价"当前活跃会话的临时授权"。handler 闭包(read lock 取快照)与 /// process_tool_calls 预校验均读此字段,两端一致。 /// - Phase C:黑名单(is_authorized 内,黑名单优先于白名单 — 命中即拒)。 #[derive(Debug, Clone, Default)] pub struct AllowedDirs { /// 持久化授权目录(Settings KV `allowed_dirs`,JSON 字符串数组)。 /// 已规范化(canonicalize 失败回退原字面量),便于 starts_with 精确比对。 pub persistent: HashSet, /// F-260619-03 Phase B: 会话级临时授权目录(弹窗"当前会话"选项写入,进程级内存)。 /// 仅当前活跃会话生效,切换/新建/删除会话由 clear_session_allowed_dirs 清空(不落库)。 /// handler 闭包与 process_tool_calls 预校验均读此字段,确保两端授权判定一致。 pub session: HashSet, /// 本次单次授权目录(弹窗"本次"选项写入)。单次工具执行放行后由 clear_once_allowed_dirs /// 清空(用完即弃,下次同路径再访问仍弹窗)。区别于 session(整会话有效)。 /// 三档授权语义:once=本次单次 / session=当前会话 / persistent=始终落 KV。 pub once: HashSet, } impl AllowedDirs { /// Settings KV key:F-260619-03 Phase A 持久化白名单(JSON 字符串数组) pub const SETTINGS_KEY: &'static str = "allowed_dirs"; /// 空白名单(方案①弱化 workspace_root:不再编译期硬塞开发机 CARGO_MANIFEST_DIR)。 /// 分发后该路径指向编译机不存在的目录,硬塞反成脏白名单。 /// 当前仅测试使用,生产环境由 reload_allowed_dirs 覆盖。 #[cfg_attr(not(test), allow(dead_code))] pub fn default_with_root() -> Self { Self { persistent: HashSet::new(), session: HashSet::new(), once: HashSet::new() } } /// 取首个 persistent 授权目录,作为相对路径锚定基准。 /// 无任何持久授权时返回 None,调用方应引导用户绑定项目。 pub fn first_persistent_dir(&self) -> Option<&PathBuf> { self.persistent.iter().next() } /// 路径是否被授权:persistent 或 session 白名单任一 starts_with 命中即放行。 /// /// **方案①弱化后 workspace_root 不再默认授权**(default_with_root 空,reload_allowed_dirs /// KV 未配时也不塞),首次访问工程根走弹窗三档授权。用户授权(always/session/once)后 /// 落白名单,后续命中放行(动态白名单完整语义,F-260619-03 收尾)。 /// /// **Phase C 黑名单优先**:即使白名单命中,若路径落入系统敏感目录(Windows /// `\Windows\System32` / `\Program Files\`;Unix `/etc /usr /bin /sbin /boot /// /dev /proc /sys`)仍拒。黑名单优先于白名单,防用户误授权系统目录。 /// /// 输入 `candidate` 理想为 canonicalize 后的真实路径(防 symlink 逃逸);调用方 /// (resolve_workspace_path_with_allowed)负责 canonicalize。但词法层预校验 /// (`check_path_authorization`)**故意不 canonicalize**(路径可能不存在,如 write_file 新建), /// 会传入未归一的词法路径 → 与白名单(canonicalize 后)形态不对齐 → 已授权路径反复误弹窗。 /// 故本函数内部 best-effort canonicalize 兜底(失败回退词法),消除所有调用方形态差异。 /// 白名单仍是唯一真相源,canonicalize 不引入新放行,顺带解析 symlink 增强(非削弱)防逃逸。 pub fn is_authorized(&self, candidate: &Path) -> bool { // Phase C: 黑名单优先(白名单命中也拒)。validate_path 已挡 .ssh/.aws 等, // 此处补系统核心目录(用户可能误把 C:\ 加入 persistent,黑名单兜底拒 System32)。 if is_in_system_blacklist(candidate) { return false; } // F-260620 比对侧收口(误弹窗核心根因):白名单 persistent/session 是 canonicalize 后形态, // candidate 来自两类调用方形态不一(handler canonicalize 后 / 预校验词法未 canonicalize)。 // 仅 strip_verbatim(去 \\?\ 前缀)不够——盘符大小写/.. /symlink/分隔符差异仍让 starts_with // 失败。best_effort_canonicalize 把 candidate 归一到真实路径(失败回退词法,不阻断合法访问), // 再 strip_verbatim 去前缀,使比对双侧形态完全对齐。e2ece2d 只收口 verbatim 前缀漏了这步。 let cand = strip_verbatim(best_effort_canonicalize(candidate)); self.persistent.iter().any(|d| cand.starts_with(d)) || self.session.iter().any(|d| cand.starts_with(d)) || self.once.iter().any(|d| cand.starts_with(d)) } } /// F-260619-03 Phase C: 路径授权决策(预校验用,process_tool_calls 分类前调)。 /// /// `check_path_authorization` 返回本枚举,process_tool_calls 据此决定: /// - `Authorized`:路径已授权 → 走原 Low/Med/High 流程 /// - `NeedsAuthorization`:路径未授权但非黑名单 → Phase B 挂起弹窗(emit AiDirAuthRequired) /// - `Denied`:路径命中黑名单 → 硬拒(工具返 Err tool_result,不挂起) #[derive(Debug, Clone, PartialEq, Eq)] pub enum PathAuthDecision { /// 已授权(workspace_root / persistent / session 命中且非黑名单) Authorized, /// 未授权且非黑名单:需 Phase B 弹窗询问用户。附带规范化后的待授权目录(父目录, /// 对齐 session_trust 目录粒度),供"仅本次/未来都允许"写入白名单。 NeedsAuthorization { dir: PathBuf }, /// 命中系统黑名单(Phase C):硬拒,工具返 Err,不弹窗。 Denied { reason: String }, } /// F-260619-03 Phase C: 系统敏感目录黑名单判定(跨平台)。 /// /// 白名单命中但黑名单命中 → 拒(防用户误把整个盘符如 `C:\` 加入 persistent 致 /// System32 被放行)。validate_path(tool_registry.rs)已挡 .ssh/.aws/AppData 等, /// 此处补系统核心目录。 /// /// - Windows(不区分大小写,按路径分隔符分段):`\Windows\System32`、`\Program Files\`、 /// `\Program Files (x86)\`、`\Windows\` 根下 System/SysWOW64 /// - Unix:`/etc`、`/usr`、`/bin`、`/sbin`、`/boot`、`/dev`、`/proc`、`/sys` pub(crate) fn is_in_system_blacklist(path: &Path) -> bool { let s = path.to_string_lossy().to_lowercase(); let seps = ['/', '\\']; // Windows:按分隔符分段判定(避免 contains 误伤 "my program files backup" 这类目录名) let segs: Vec<&str> = s.split(seps).filter(|s| !s.is_empty()).collect(); for (i, seg) in segs.iter().enumerate() { // Windows 系统目录:C:\Windows(根本身及所有子目录 win.ini/explorer.exe/hosts 等)/ // C:\Program Files / C:\Program Files (x86) / C:\ProgramData if cfg!(windows) { // windows 段命中即拒(含根本身,不再限定 system32 子目录 — 防 win.ini/hosts 等) if *seg == "windows" { return true; } // C:\Program Files / C:\Program Files (x86) if *seg == "program files" || *seg == "program files (x86)" { return true; } // C:\ProgramData(系统级应用数据) if *seg == "programdata" { return true; } } // 用户级凭据目录(.ssh/.aws/.gnupg):任意路径段命中即拒(跨平台,从 validate_path 迁移统一, // 消除 validate_path contains 与 is_in_system_blacklist 分段两套黑名单不一致) if matches!(*seg, ".ssh" | ".aws" | ".gnupg") { return true; } // Unix 系统目录(路径首段为这些即拒;Windows 上也防 Unix 风格绝对路径,防御性) if i == 0 && matches!(*seg, "etc" | "usr" | "bin" | "sbin" | "boot" | "dev" | "proc" | "sys") { return true; } } false } /// F-260619-03 Phase B/C: 路径授权预校验(供 process_tool_calls 分类前调)。 /// /// 词法层判定(不 canonicalize,因路径可能不存在 — write_file 新建)。返回三态决策: /// - 路径规范化(去 .. / 锚定首个持久授权目录)后,若命中黑名单 → `Denied` /// - 否则若 `is_authorized(规范化路径)`(persistent + session) → `Authorized` /// - 否则 → `NeedsAuthorization { dir: 父目录规范化 }`(目录粒度,对齐 session_trust) /// /// 注意:本函数只做词法层预判(防不存在路径兜底);实际执行时 handler 内 /// resolve_workspace_path_with_allowed 仍做完整校验(canonicalize symlink 防逃逸)。 /// 黑名单在两处都判(is_authorized 内 + 此处独立判),双保险。 pub fn check_path_authorization( raw_path: &str, allowed: &AllowedDirs, ) -> PathAuthDecision { // 规范化:绝对路径原样,相对路径锚定首个持久授权目录(与 resolve_workspace_path_impl 一致) let resolved = if Path::new(raw_path).is_absolute() { PathBuf::from(raw_path) } else if let Some(root) = allowed.first_persistent_dir() { root.join(raw_path) } else { // 无授权目录时直接返 NeedsAuthorization(让引导流程处理),避免锚定到编译机路径。 return PathAuthDecision::NeedsAuthorization { dir: PathBuf::from("."), }; }; // Phase C: 黑名单优先独立判定(is_authorized 内也判,此处先判便于 NeedsAuthorization // 不误把黑名单路径推到弹窗 — 黑名单路径直接硬拒不让用户"授权")。 if is_in_system_blacklist(&resolved) { return PathAuthDecision::Denied { reason: format!("路径命中系统敏感目录黑名单: {}", raw_path), }; } if allowed.is_authorized(&resolved) { return PathAuthDecision::Authorized; } // 未授权 → 取父目录作待授权目录(目录粒度,对齐 session_trust 的 dir 语义) let dir = resolved .parent() .map(|p| p.to_path_buf()) .unwrap_or_else(|| PathBuf::from(".")); PathAuthDecision::NeedsAuthorization { dir } } /// Windows canonicalize 返回 `\\?\E:\...` 扩展路径(verbatim 前缀),与词法层(不 canonicalize) /// 的 starts_with 比对不一致 → 白名单含但工具路径不匹配 → 误弹窗。strip 前缀统一为 `E:\...`。 fn strip_verbatim(p: PathBuf) -> PathBuf { let s = p.to_string_lossy().to_string(); // Windows verbatim/设备路径前缀:\\?\ (Volume canonicalize)、\\.\ (设备命名空间)、 // \??\ (对象管理器命名空间)。统一 strip 后与词法层 starts_with 比对一致。 for prefix in [r"\\?\", r"\\.\", r"\??\"] { if let Some(stripped) = s.strip_prefix(prefix) { return PathBuf::from(stripped); } } p } /// best-effort canonicalize:把任意形态路径(词法/正斜杠/含../大小写不一)归一到真实路径, /// 供 `is_authorized` 与白名单(canonicalize 后)比对侧形态对齐(误弹窗根因修复)。 /// /// - 路径存在 → `canonicalize` 成功,返回真实路径(大小写归一/.. 解析/symlink 解析/去冗余分隔符) /// - 路径不存在(write_file 新建文件)→ canonicalize 其**父目录**(通常存在)+ 拼回文件名, /// 父目录也不存在 → 回退原词法路径(交由 strip_verbatim + starts_with 尽力匹配,不阻断) /// /// 失败一律回退,绝不返回 Err——授权比对宁可降级匹配不可 panic/阻断合法访问。 fn best_effort_canonicalize(p: &Path) -> PathBuf { if let Ok(real) = std::fs::canonicalize(p) { return real; } if let Some(parent) = p.parent() { if let Ok(real_parent) = std::fs::canonicalize(parent) { if let Some(file_name) = p.file_name() { return real_parent.join(file_name); } return real_parent; } } p.to_path_buf() } /// workspace 根目录(当前:编译期 CARGO_MANIFEST_DIR 上两级,仅开发机有效)。 /// /// ⚠️ 分发限制:env!("CARGO_MANIFEST_DIR") 是编译期常量,打包后指向编译机路径,用户机器无效。 /// 目标方案(D 混合策略):workspace_root 改为运行期动态确定,从已绑定项目的 /// `projects.path` 并集 + `AllowedDirs` 推导。无绑定项目时 workspace_root 为空, /// 文件操作拒绝,引导用户绑定项目。当前保留编译期常量作为开发期兼容。 /// 已修复: /// - authz_debug → std::env::temp_dir() (chat.rs) /// - workspace_drive → std::env::current_dir() 运行时确定 (chat.rs) fn workspace_root_path() -> PathBuf { PathBuf::from(env!("CARGO_MANIFEST_DIR")) .parent() .and_then(|p| p.parent()) .map(PathBuf::from) .unwrap_or_else(|| PathBuf::from(".")) } impl AppState { /// 初始化应用状态:打开(或创建)数据库并执行迁移,构建各 Repo 与节点注册表 pub async fn init(db_path: &Path, data_dir: PathBuf) -> Result { let db = Arc::new(Database::open(db_path).await?); // F-260619-03 Phase A: build_ai_tool_registry 注入 allowed_dirs Arc, // 文件工具闭包捕获后 resolve_workspace_path_with_allowed 校验动态白名单。 // 此处用 default_with_root 占位,下方 init 尾部 reload_allowed_dirs 从 Settings KV 覆盖。 let allowed_dirs = Arc::new(RwLock::new(AllowedDirs::default())); let ai_tools = Arc::new(crate::commands::ai::build_ai_tool_registry( &db, &allowed_dirs, data_dir.clone(), )); // Input Augmentation 层(核心设计2):ResolverRegistry 启动期注册四 resolver。 // resolver 持 Arc(非 AppState,避免循环依赖:AppState 持 Arc), // 在 struct 字段 `db` move 前 clone 注入,与 build_ai_tool_registry 同侧。 let resolvers = Arc::new(crate::commands::ai::build_resolver_registry(db.clone())); // build_registry 需注入 Arc(TaskAdvanceNode 持 db)。 // 在 struct 字段 `db` move 前 clone,避免 E0382。 let registry = Arc::new(build_registry(db.clone())); let state = Self { ideas: IdeaRepo::new(&db), idea_evaluations: IdeaEvalRepo::new(&db), projects: ProjectRepo::new(&db), tasks: TaskRepo::new(&db), task_links: TaskLinkRepo::new(&db), project_events: ProjectEventRepo::new(&db), project_services: ProjectServiceRepo::new(&db), releases: ReleaseRepo::new(&db), workflows: WorkflowRepo::new(&db), node_executions: NodeExecutionRepo::new(&db), ai_providers: AiProviderRepo::new(&db), ai_conversations: AiConversationRepo::new(&db), ai_messages: AiMessageRepo::new(&db), ai_tool_executions: AiToolExecutionRepo::new(&db), ai_session: Arc::new(Mutex::new(AiSession::new())), // L3 批2b:AI 事件总线入 AppState(独立于工作流 event_bus)。emit 双写留批3。 ai_event_bus: crate::commands::ai::event_bus::EventBus::new(), // Phase3 Layer1:tunnel 客户端未连接实例,lib.rs setup 内 connect tunnel: std::sync::Arc::new(df_tunnel::WsTunnelClient::new()), knowledge: KnowledgeRepo::new(&db), knowledge_events: KnowledgeEventsRepo::new(&db), knowledge_config: Arc::new(Mutex::new(KnowledgeConfig::default())), data_dir: data_dir.clone(), settings: SettingsRepo::new(&db), llm_concurrency: LlmConcurrency::new(3, 2), agent_max_iterations: Arc::new(AtomicUsize::new( crate::commands::ai::agentic::DEFAULT_MAX_AGENT_ITERATIONS, )), agent_max_retries: Arc::new(AtomicUsize::new( crate::commands::ai::agentic::DEFAULT_MAX_AGENT_RETRIES, )), workflow_state_registry: Arc::new(Mutex::new(HashMap::new())), // F-260619-03 Phase A: 与 ai_tools registry 共享同一 Arc(构建时注入同一句柄) allowed_dirs: allowed_dirs.clone(), db, event_bus: EventBus::new(), registry, ai_tools, resolvers, }; // 启动恢复:重启前卡 pending 的工具审批(内存 pending_approvals 已丢)从审计表重建, // 使重启后待审批不丢。前端经 ai_pending_tool_calls + switchConversation 恢复 toolCard 态。 crate::commands::ai::restore_pending_approvals(&state).await; // F-260614-04c: 启动一次性初始化 per-provider caps 表。 // 根据 DB enabled providers 建 HashMap,让 agentic loop 的 // acquire_for_provider 从 None(无限流)切换到 Some(按配置限流)。disabled / weight=0 // 的 provider 不入表(其被 provider_pool::select 过滤出候选,不会被 acquire)。 // 单 provider 场景:该 provider cap=global_cap → acquire_global+acquire_for_provider // 串联,min(3,3)=3,有效上限同未配置 → 行为零变化。 state.reload_provider_caps().await; // F-260619-03 Phase A: 从 Settings KV 加载持久化授权目录白名单覆盖默认值。 // 失败(读 KV/解析 JSON 出错)不阻断启动,保持 default_with_root(方案①后为空)。 state.reload_allowed_dirs().await; // P0(设置走查):从 Settings KV 恢复持久化知识库配置覆盖 default(防 8 项重启全丢)。 state.reload_knowledge_config().await; // 迁移旧 .trash(编译期 workspace_root → 运行期 data_dir),仅一次,幂等。 let old_trash = workspace_root_path().join(".trash"); let new_trash = data_dir.join(".trash"); if old_trash.exists() && !new_trash.exists() { if let Err(e) = std::fs::rename(&old_trash, &new_trash) { // 跨盘 rename 失败(E:→C:),静默跳过,新目录自动创建 tracing::debug!( "迁移 .trash 跨盘失败(从 {:?} 到 {:?}): {} (原目录保留,新目录将自动创建)", old_trash, new_trash, e, ); } } Ok(state) } /// F-260614-04c: 从 DB enabled providers 重建 per_provider caps 并 set_provider_caps。 /// /// 启动 + provider 变更(ai_update_provider_pool / 删除 / 新增)后调用。caps 策略见 /// AppState::init 注释(本轮每 provider cap = global_cap,F-04d 配差异化上限时仅改此)。 pub async fn reload_provider_caps(&self) { let providers = match self.ai_providers.list_all().await { Ok(ps) => ps, Err(e) => { tracing::warn!( "[F-04c] 读取 providers 重建 caps 失败,保持当前 caps 表不变: {}", e ); return; } }; let global_cap = 3; // 与 LlmConcurrency::new(3, 2) 的 global 上限对齐。 let caps: HashMap = providers .into_iter() .filter(|p| p.enabled && p.weight > 0) .map(|p| (p.id, global_cap)) .collect(); self.llm_concurrency.set_provider_caps(caps).await; } /// P0(设置走查-2026-06-21):从 Settings KV 加载持久化知识库配置覆盖默认值。 /// /// 原纯内存 knowledge_config 启动 default() 覆盖致用户配置重启全丢;save_config 落 /// KNOWLEDGE_CONFIG_KEY KV,本方法启动恢复。失败(读 KV/反序列化)不阻断启动,保持 default。 pub async fn reload_knowledge_config(&self) { match self.settings.get(KNOWLEDGE_CONFIG_KEY).await { Ok(Some(json)) => match serde_json::from_str::(&json) { Ok(cfg) => *self.knowledge_config.lock().await = cfg, Err(e) => tracing::warn!( "[KNOWLEDGE-CONFIG] KV 反序列化失败,保持 default: {}", e ), }, Ok(None) => {} // 首次启动无持久化,保持 default Err(e) => tracing::warn!("[KNOWLEDGE-CONFIG] 读 KV 失败,保持 default: {}", e), } } /// F-260619-03 Phase A: 从 Settings KV `allowed_dirs`(JSON 字符串数组)加载持久化白名单。 /// /// 启动 + Settings IPC `ai_set_allowed_dirs` 写入后调用。解析失败/缺失 → 保持 /// default_with_root(方案①后为空,首次访问走弹窗授权)。 /// 每条路径尝试 canonicalize 规范化(防大小写/分隔符差异绕过);canonicalize 失败 /// (目录不存在)回退原字面量 trim(写入后再校验场景:先授权目录路径,目录暂不存在)。 pub async fn reload_allowed_dirs(&self) { // F-260619-03: 白名单 = Settings KV allowed_dirs(用户配)+ projects.bind_directory(项目绑定目录自动授权)。 // 项目绑定目录 = AI 天然可访问(用户已主动绑定),不应重复弹窗。reload 时自动合并入 persistent。 let kv_dirs: Vec = match self.settings.get(AllowedDirs::SETTINGS_KEY).await { Ok(Some(v)) => serde_json::from_str(&v).unwrap_or_default(), Ok(None) => Vec::new(), Err(e) => { tracing::warn!("[F-03A] 读取 allowed_dirs 失败: {}", e); Vec::new() } }; // 读 projects.bind_directory(所有项目的绑定目录,自动白名单) let project_dirs: Vec = { let repo = df_storage::crud::ProjectRepo::new(&self.db); match repo.list_all().await { Ok(projects) => projects.iter() .filter_map(|p| p.path.as_ref().filter(|d| !d.is_empty()).map(|d| d.clone())) .collect(), Err(e) => { tracing::warn!("[F-03B] 读取项目绑定目录失败: {}", e); Vec::new() } } }; // 合并:KV allowed_dirs + projects.bind_directory let mut all_dirs = kv_dirs; all_dirs.extend(project_dirs); // 去重 all_dirs.sort(); all_dirs.dedup(); let mut set = HashSet::new(); // 始终插入 workspace_root(工程内路径默认免授权,对齐用户政策 + 注释承诺)。 // BUG-260620-05 修:去掉 all_dirs.is_empty() 条件,无条件插入,不再依赖 KV/project_dirs 是否存在。 set.insert(workspace_root_path()); for d in all_dirs { let d = d.trim(); if d.is_empty() { continue; } let p = PathBuf::from(d); // canonicalize 成功用真实路径(去 symlink/大小写归一);失败回退原字面量 trim // (兼容"先授权目录,目录暂不存在"用例)。失败打 warn 便于排查静默授权错路径。 let normalized = strip_verbatim(std::fs::canonicalize(&p).unwrap_or_else(|_| { tracing::warn!("[allowed_dirs] canonicalize 失败(目录可能不存在),按字面量保存: {}", d); PathBuf::from(d.trim_end_matches(['/', '\\'])) })); set.insert(normalized); } // F-260619-03 Phase B: reload 时保留当前会话临时授权(session 不落库,仅内存), // 仅覆盖 persistent。用户改 Settings 不影响当前会话已临时授权的目录。 let mut guard = self.allowed_dirs.write().await; let preserved_session = std::mem::take(&mut guard.session); let preserved_once = std::mem::take(&mut guard.once); *guard = AllowedDirs { persistent: set, session: preserved_session, once: preserved_once }; } /// F-260619-03 Phase A: 写 Settings KV + 同步内存白名单(供 Settings IPC 调用)。 /// /// - 持久化:JSON 字符串数组写 `app_settings` key=`allowed_dirs` /// - 内存:reload_allowed_dirs 重新加载(规范化逻辑复用,避免双份) /// 返回持久化后的规范化路径列表(供前端回显 canonicalize 后的真实路径)。 pub async fn set_allowed_dirs(&self, dirs: Vec) -> Result> { // 去空 + 去重(保留顺序,前端展示友好)+ 黑名单预校验(防持久化系统敏感目录, // 纵深防御:即便写入,is_authorized 运行时黑名单也兜底拒,预校验保证白名单 UI 洁净) let mut seen = HashSet::new(); let cleaned: Vec = dirs .into_iter() .map(|s| s.trim().to_string()) .filter(|s| !s.is_empty()) .filter(|s| seen.insert(s.clone())) .filter(|s| !is_in_system_blacklist(&PathBuf::from(s))) .collect(); let json = serde_json::to_string(&cleaned)?; self.settings.set(AllowedDirs::SETTINGS_KEY, &json).await .map_err(|e| anyhow::anyhow!("持久化 allowed_dirs 失败: {}", e))?; self.reload_allowed_dirs().await; // 返回内存白名单的规范化字符串列表:复用 get_allowed_dirs(单一真相源, // 消除原 9 行逐字重复——过滤 workspace_root + sort,语义完全一致)。 Ok(self.get_allowed_dirs().await) } /// F-260619-03 Phase A: 读内存白名单为字符串列表(供 Settings IPC `ai_get_allowed_dirs` 回显)。 pub async fn get_allowed_dirs(&self) -> Vec { let guard = self.allowed_dirs.read().await; let mut out: Vec = guard .persistent .iter() .map(|p| p.to_string_lossy().to_string()) .collect(); out.sort(); out } /// F-260619-03 Phase B: 追加持久化授权目录(弹窗"未来都允许"选项)。 /// /// 读当前 persistent → 追加新目录(去重)→ set_allowed_dirs 持久化 + 同步内存。 /// 返回持久化后规范化列表(含 workspace_root)。 pub async fn add_persistent_allowed_dir(&self, dir: String) -> Result> { let d = dir.trim().to_string(); if d.is_empty() { return Ok(self.get_allowed_dirs().await); } let mut current = self.get_allowed_dirs().await; if !current.iter().any(|c| paths_eq(c, &d)) { current.push(d); } self.set_allowed_dirs(current).await } /// F-260619-03 Phase B: 追加会话级临时授权目录(弹窗"仅本次"选项)。 /// /// 写入 `allowed_dirs.session`(进程级内存,不落库)。handler 闭包 read lock 取快照时 /// 与 process_tool_calls 预校验读同一字段,两端授权判定一致。规范化:canonicalize /// 失败回退原字面量 trim(与 reload_allowed_dirs 一致)。 pub async fn add_session_allowed_dir(&self, dir: String) { let d = dir.trim(); if d.is_empty() { return; } let p = PathBuf::from(d); let normalized = strip_verbatim(std::fs::canonicalize(&p).unwrap_or_else(|_| { PathBuf::from(d.trim_end_matches(['/', '\\'])) })); self.allowed_dirs.write().await.session.insert(normalized); } /// F-260619-03 Phase B: 清空会话级临时授权(切换/新建/删除活跃会话时调)。 /// /// 单用户桌面应用 active_conversation_id 单全局模型,session 字段语义为 /// "当前活跃会话的临时授权",切走即清空(不跨会话继承临时授权)。 pub async fn clear_session_allowed_dirs(&self) { self.allowed_dirs.write().await.session.clear(); } /// F-260619-03 Phase B: 追加本次单次授权目录(弹窗"本次"选项)。 /// /// 写入 `allowed_dirs.once`(进程级内存,不落库)。handler 执行工具时 read lock /// 取快照读 once 放行,执行后由 clear_once_allowed_dirs 清空(本次放行即失效, /// 下次同路径仍弹窗)。规范化与 add_session_allowed_dir 一致(canonicalize 失败回退 trim)。 pub async fn add_once_allowed_dir(&self, dir: String) { let d = dir.trim(); if d.is_empty() { return; } let p = PathBuf::from(d); let normalized = strip_verbatim(std::fs::canonicalize(&p).unwrap_or_else(|_| { PathBuf::from(d.trim_end_matches(['/', '\\'])) })); self.allowed_dirs.write().await.once.insert(normalized); } /// F-260619-03 Phase B: 清空本次单次授权(工具执行后调,本次放行即失效)。 pub async fn clear_once_allowed_dirs(&self) { self.allowed_dirs.write().await.once.clear(); } } /// F-260619-03 Phase B: 路径字符串等价比较(canonicalize 后比对,失败回退小写比对)。 fn paths_eq(a: &str, b: &str) -> bool { let pa = std::fs::canonicalize(a).map(|p| p.to_string_lossy().to_string()).unwrap_or_else(|_| a.to_string()); let pb = std::fs::canonicalize(b).map(|p| p.to_string_lossy().to_string()).unwrap_or_else(|_| b.to_string()); pa.eq_ignore_ascii_case(&pb) } /// 构建节点注册表 — 注册内置节点 /// /// 注意:不使用 `NodeRegistry::default()`,其 script 工厂为占位实现(会 panic), /// 这里注册 df-nodes 提供的真实 HumanNode / AiNode。 /// /// "script" 节点不注册:ScriptNode 走 cmd /C | sh -c 执行 config.command 原始串, /// 前端可构造任意 DagDef 触达无审批 shell(R-PD-2)。DevFlow 工作流当前为纯演示 /// 功能(前端唯一构造点 ProjectDetail.vue demoDag 三步 echo),无真实构建/部署脚本 /// 需求。需要脚本执行能力时新建独立 BuildNode(白名单 + 项目目录锚定 + 复用 AI 工具 /// RiskLevel 审批链),而非回头启用 ScriptNode + 黑名单。 /// 详见 docs/02-架构设计/专项设计/工作流脚本执行边界-2026-06-15.md。 fn build_registry(db: Arc) -> NodeRegistry { let mut registry = NodeRegistry::new(); registry.register("human", |_config| { Box::new(df_nodes::human_node::HumanNode) }); // AiNode(df_nodes::ai_node, impl Node trait): // 决策 a(AiNode 自审闭环)步骤②:AiNode 持 Arc,execute 完成后若有 task_id // 则把产出落 task.output_json。工厂闭包 move 捕获 db 句柄(Arc clone 廉价)。 let ai_db = db.clone(); registry.register("ai", move |_config| { Box::new(df_nodes::ai_node::AiNode::new(ai_db.clone())) }); // AiSelfReviewNode(df_nodes::ai_self_review_node, impl Node trait): // 决策 a 步骤③:四维度自审(prompt 固定+JSON 解析兜底),写回 output_json 加 review 子字段, // review 摘要塞 NodeOutput.data 供下游 human_review 经 DAG inputs 透传。 // testing 模板 ai_self_review 节点类型对齐此注册 key。 let review_db = db.clone(); registry.register("ai_self_review", move |_config| { Box::new(df_nodes::ai_self_review_node::AiSelfReviewNode::new(review_db.clone())) }); // TaskAdvanceNode(df_nodes::task_advance_node, impl Node trait): // 推进链阶段 2 工作流联动入口 — DAG 内触发 advance_task。 // 持有 Arc 在此构造时注入(NodeRegistry::register 工厂闭包 move 捕获 db, // Arc clone 廉价),Node::execute 从 NodeContext.config 读 task_id/target_status。 // D-260616-03: 推进链/状态机/闸门走 df-nodes Node trait,非复活 df-task。 registry.register("task_advance", move |_config| { Box::new(df_nodes::task_advance_node::TaskAdvanceNode::new(db.clone())) }); registry } #[cfg(test)] mod tests { use super::*; // ============================================================ // F-260619-03 Phase A: AllowedDirs 授权语义测试 // 锁定:workspace_root 始终授权 + persistent 命中放行 + 未命中拒绝。 // ============================================================ /// 显式授权 workspace_root 后命中(方案①:default_with_root 不再预设,需显式授权) #[test] fn test_allowed_dirs_workspace_root_authorized() { let mut set = HashSet::new(); set.insert(workspace_root_path()); let allowed = AllowedDirs { persistent: set, session: HashSet::new(), once: HashSet::new() }; let root = workspace_root_path(); assert!(allowed.is_authorized(&root), "显式授权的 workspace_root 应命中"); // workspace 内子路径也应授权(starts_with workspace_root) let child = root.join("src").join("main.rs"); assert!(allowed.is_authorized(&child), "workspace_root 子路径应被授权"); } /// 自定义授权目录命中放行(模拟用户授权 E:/some/external/dir) #[test] fn test_allowed_dirs_custom_authorized() { let mut set = HashSet::new(); let custom = PathBuf::from("E:/wk-test-external-dir"); set.insert(custom.clone()); let allowed = AllowedDirs { persistent: set, session: HashSet::new(), once: HashSet::new() }; // 精确命中 + 子路径 starts_with 命中 assert!(allowed.is_authorized(&custom)); assert!(allowed.is_authorized(&custom.join("sub").join("file.txt"))); } /// 未授权目录被拒绝 #[test] fn test_allowed_dirs_unauthorized_rejected() { let mut set = HashSet::new(); set.insert(PathBuf::from("E:/wk-test-authorized")); let allowed = AllowedDirs { persistent: set, session: HashSet::new(), once: HashSet::new() }; let outside = PathBuf::from("E:/wk-test-unauthorized/file.txt"); assert!(!allowed.is_authorized(&outside), "未授权目录应被拒绝"); } /// 默认(空 persistent)只授权 workspace_root #[test] fn test_allowed_dirs_default_empty_persistent() { let allowed = AllowedDirs::default(); // 空 persistent,无 workspace_root → 任何路径都不授权 // (default_with_root 才含 workspace_root;Default 不含,用于边界测试) assert!(!allowed.is_authorized(&PathBuf::from("E:/anything"))); } /// SETTINGS_KEY 常量稳定(防 rename 致持久化数据丢失) #[test] fn test_allowed_dirs_settings_key_stable() { assert_eq!(AllowedDirs::SETTINGS_KEY, "allowed_dirs"); } // ============================================================ // F-260619-03 Phase B: session 临时授权语义测试 // ============================================================ /// Phase B: session 命中放行(persistent 未命中但 session 命中) #[test] fn test_allowed_dirs_session_authorized() { let mut session = HashSet::new(); session.insert(PathBuf::from("E:/wk-temp-session")); let allowed = AllowedDirs { persistent: HashSet::new(), session, once: HashSet::new() }; assert!(allowed.is_authorized(&PathBuf::from("E:/wk-temp-session/file.txt"))); } /// Phase B: persistent + session 任一命中即放行 #[test] fn test_allowed_dirs_persistent_or_session() { let mut persistent = HashSet::new(); persistent.insert(PathBuf::from("E:/wk-persist")); let mut session = HashSet::new(); session.insert(PathBuf::from("E:/wk-session")); let allowed = AllowedDirs { persistent, session, once: HashSet::new() }; assert!(allowed.is_authorized(&PathBuf::from("E:/wk-persist/a"))); assert!(allowed.is_authorized(&PathBuf::from("E:/wk-session/b"))); assert!(!allowed.is_authorized(&PathBuf::from("E:/wk-other/c"))); } // ============================================================ // F-260619-03 Phase C: 系统目录黑名单测试 // ============================================================ /// Phase C: Windows System32 黑名单命中(即使在白名单内也拒) #[test] fn test_blacklist_windows_system32() { let mut persistent = HashSet::new(); // 用户误把整个 C:\ 加入白名单 persistent.insert(PathBuf::from("C:\\")); let allowed = AllowedDirs { persistent, session: HashSet::new(), once: HashSet::new() }; // System32 应被黑名单拒(尽管 C:\ 在白名单) assert!(!allowed.is_authorized(&PathBuf::from("C:\\Windows\\System32\\config\\sam"))); } /// Phase C: Windows Program Files 黑名单命中 #[test] fn test_blacklist_windows_program_files() { let mut persistent = HashSet::new(); persistent.insert(PathBuf::from("C:\\")); let allowed = AllowedDirs { persistent, session: HashSet::new(), once: HashSet::new() }; assert!(!allowed.is_authorized(&PathBuf::from("C:\\Program Files\\SomeApp\\app.exe"))); } /// Phase C: 黑名单不误伤合法目录(含 "program files" 子串的自定义目录名) #[test] fn test_blacklist_no_false_positive() { // "my program files backup" 不应被拒(分段匹配,非精确段) assert!(!is_in_system_blacklist(&PathBuf::from("E:/my program files backup/x"))); // 普通工作目录不拒 assert!(!is_in_system_blacklist(&PathBuf::from("E:/wk-lab/devflow/src"))); } /// Phase C: Unix 系统目录黑名单(/etc /usr /bin 等) #[test] fn test_blacklist_unix_system_dirs() { assert!(is_in_system_blacklist(&PathBuf::from("/etc/passwd"))); assert!(is_in_system_blacklist(&PathBuf::from("/usr/bin/python"))); assert!(is_in_system_blacklist(&PathBuf::from("/proc/self/environ"))); assert!(is_in_system_blacklist(&PathBuf::from("/sys/kernel"))); // 普通用户目录不拒 assert!(!is_in_system_blacklist(&PathBuf::from("/home/user/project"))); } // ============================================================ // F-260619-03 Phase B/C: check_path_authorization 三态决策测试 // ============================================================ /// Phase B/C: 显式授权 workspace_root 后,内路径 → Authorized(方案①:不再默认) #[test] fn test_check_path_authorized_workspace() { let mut set = HashSet::new(); set.insert(workspace_root_path()); let allowed = AllowedDirs { persistent: set, session: HashSet::new(), once: HashSet::new() }; let rel = "src/main.rs"; match check_path_authorization(rel, &allowed) { PathAuthDecision::Authorized => {} other => panic!("workspace_root 内路径应 Authorized, got {:?}", other), } } /// Phase B/C: 未授权路径 → NeedsAuthorization(附父目录) #[test] fn test_check_path_needs_auth() { let allowed = AllowedDirs::default(); // 空,无 workspace_root match check_path_authorization("E:/wk-external/file.txt", &allowed) { PathAuthDecision::NeedsAuthorization { dir } => { assert_eq!(dir, PathBuf::from("E:/wk-external")); } other => panic!("未授权路径应 NeedsAuthorization, got {:?}", other), } } /// Phase B: session 命中 → Authorized(check_path_authorization 读 session 字段) #[test] fn test_check_path_session_hit() { let mut session = HashSet::new(); session.insert(PathBuf::from("E:/wk-session")); let allowed = AllowedDirs { persistent: HashSet::new(), session, once: HashSet::new() }; match check_path_authorization("E:/wk-session/sub/file.txt", &allowed) { PathAuthDecision::Authorized => {} other => panic!("session 命中应 Authorized, got {:?}", other), } } /// Phase C: 黑名单路径 → Denied(不弹窗直接拒) #[test] fn test_check_path_denied_blacklist() { let allowed = AllowedDirs::default_with_root(); match check_path_authorization("C:/Windows/System32/config/sam", &allowed) { PathAuthDecision::Denied { .. } => {} other => panic!("System32 应 Denied, got {:?}", other), } } // ============================================================ // F-260620: strip_verbatim 比对侧收口 + 黑名单增强测试 // 三方审查(安全/UX/跨端)交叉印证:strip_verbatim 仅写入侧调用,比对侧遗漏致误弹窗。 // ============================================================ /// F-260620: is_authorized 比对侧 strip_verbatim — candidate 带 \\?\ 前缀也命中白名单 /// (handler canonicalize 后路径 vs persistent 词法路径,形态不一致曾致误弹窗)。 /// persistent 用动态 workspace_root_path(非硬编码开发机路径),candidate 拼 verbatim 前缀。 #[test] fn test_is_authorized_strips_verbatim_prefix() { let mut set = HashSet::new(); set.insert(workspace_root_path()); let allowed = AllowedDirs { persistent: set, session: HashSet::new(), once: HashSet::new() }; let child = workspace_root_path().join("src").join("main.rs"); // candidate 带 verbatim 前缀(handler canonicalize 后形态)应命中白名单 let verbatim = PathBuf::from(format!("{}{}", r"\\?\", child.to_string_lossy())); assert!(allowed.is_authorized(&verbatim), "verbatim 前缀路径应命中白名单"); // 设备命名空间前缀(\\.\) let dev = PathBuf::from(format!("{}{}", r"\\.\", child.to_string_lossy())); assert!(allowed.is_authorized(&dev), "设备命名空间前缀路径应命中白名单"); } /// F-260620: strip_verbatim 三前缀(\\?\ / \\.\ / \??\)统一去除 #[test] fn test_strip_verbatim_three_prefixes() { assert_eq!(strip_verbatim(PathBuf::from(r"\\?\E:\foo")), PathBuf::from("E:\\foo")); assert_eq!(strip_verbatim(PathBuf::from(r"\\.\E:\foo")), PathBuf::from("E:\\foo")); assert_eq!(strip_verbatim(PathBuf::from(r"\??\E:\foo")), PathBuf::from("E:\\foo")); // 无前缀原样返回 assert_eq!(strip_verbatim(PathBuf::from("E:\\foo")), PathBuf::from("E:\\foo")); } /// F-260620: 黑名单增强 — Windows 根本身 / ProgramData / 用户凭据目录(.ssh/.aws/.gnupg) #[test] fn test_blacklist_windows_root_programdata_creds() { // C:\Windows 根本身(不再限定 system32 子目录 — 防 win.ini/hosts 等) assert!(is_in_system_blacklist(&PathBuf::from("C:\\Windows\\win.ini"))); assert!(is_in_system_blacklist(&PathBuf::from("C:\\Windows\\explorer.exe"))); // ProgramData assert!(is_in_system_blacklist(&PathBuf::from("C:\\ProgramData\\App\\config"))); // 用户凭据目录(从 validate_path 迁移统一,任意路径段命中即拒) assert!(is_in_system_blacklist(&PathBuf::from("C:\\Users\\user\\.ssh\\id_rsa"))); assert!(is_in_system_blacklist(&PathBuf::from("/home/user/.aws/credentials"))); assert!(is_in_system_blacklist(&PathBuf::from("/home/user/.gnupg/pubring.gpg"))); // 普通目录不误伤(devflow 自身 / 含 appdata 段的合法备份目录) assert!(!is_in_system_blacklist(&PathBuf::from("E:\\wk-lab\\devflow\\src"))); assert!(!is_in_system_blacklist(&PathBuf::from("D:\\backup\\appdata\\data"))); } // ============================================================ // Phase3 预留(批2-B): LlmConcurrency per_sub_flow 层测试 // 占位层:acquire_per_sub_flow / release_sub_flow / set_per_sub_permits。 // 不接入 ai/ 调用点(Phase3 子流 spawn 落地时接入),此处仅验证层本身语义。 // ============================================================ /// 基本 acquire/drop:首次 acquire 返 Some,Drop 后 permit 释放可再 acquire。 #[tokio::test] async fn test_per_sub_flow_basic_acquire_drop() { let lc = LlmConcurrency::new(3, 2); let sub = "conv-1::sub-1"; // 首次 acquire 返 Some(permits=3 默认)。 let p1 = lc.acquire_per_sub_flow(sub).await; assert!(p1.is_some(), "首次 acquire 应返回 Some"); // Drop 后 permit 释放,可再 acquire。 drop(p1); let p2 = lc.acquire_per_sub_flow(sub).await; assert!(p2.is_some(), "Drop 后再 acquire 应返回 Some"); } /// 耗尽返 None:permits=1 时二次 acquire 返 None 降级串行(非阻塞)。 #[tokio::test] async fn test_per_sub_flow_exhausted_returns_none() { let lc = LlmConcurrency::new(3, 2); let sub = "conv-1::sub-1"; // 先热改 per_sub_permits=1 再 acquire(热改清 HashMap,新 sub_flow 用 permits=1)。 lc.set_per_sub_permits(1).await; let p1 = lc.acquire_per_sub_flow(sub).await; assert!(p1.is_some(), "permits=1 首次 acquire 应返回 Some"); // permits=1 已被 p1 占用,二次 acquire 非阻塞返 None 降级串行。 let p2 = lc.acquire_per_sub_flow(sub).await; assert!(p2.is_none(), "permits=1 二次 acquire 应返回 None(降级串行)"); // Drop p1 后释放,再 acquire 可成功(证明 None 是因耗尽,非 Semaphore close)。 drop(p1); let p3 = lc.acquire_per_sub_flow(sub).await; assert!(p3.is_some(), "Drop 后再 acquire 应返回 Some"); } /// 不同 sub_flow_id 独立限流:A 耗尽不影响 B。 #[tokio::test] async fn test_per_sub_flow_independent_per_id() { let lc = LlmConcurrency::new(3, 2); // permits=1:A 和 B 各自独立 Semaphore。 lc.set_per_sub_permits(1).await; let pa = lc.acquire_per_sub_flow("conv-1::sub-A").await; let pb = lc.acquire_per_sub_flow("conv-1::sub-B").await; assert!(pa.is_some(), "sub-A 首次 acquire 应返回 Some"); assert!(pb.is_some(), "sub-B 首次 acquire 应返回 Some(独立限流,A 不影响 B)"); // A 二次耗尽,B 二次也耗尽(各自 permits=1)。 assert!(lc.acquire_per_sub_flow("conv-1::sub-A").await.is_none()); assert!(lc.acquire_per_sub_flow("conv-1::sub-B").await.is_none()); } /// release_sub_flow 清理 HashMap:release 后再 acquire 会重建 Semaphore(permits 重置)。 #[tokio::test] async fn test_per_sub_flow_release_clears_entry() { let lc = LlmConcurrency::new(3, 2); lc.set_per_sub_permits(1).await; let sub = "conv-1::sub-1"; let p1 = lc.acquire_per_sub_flow(sub).await; assert!(p1.is_some()); // release 清理 HashMap 条目(p1 permit 仍有效,绑旧 Arc)。 lc.release_sub_flow(sub).await; // 再 acquire:HashMap 已清空 → 重建 Semaphore(permits=1)→ 返 Some。 // 证明 release 移除了条目(否则旧 Semaphore 已耗尽会返 None)。 let p2 = lc.acquire_per_sub_flow(sub).await; assert!(p2.is_some(), "release 后再 acquire 应重建 Semaphore 返回 Some"); drop(p1); drop(p2); } /// set_per_sub_permits 软收敛:热改后新 sub_flow 用新 permits 值。 #[tokio::test] async fn test_set_per_sub_permits_soft_converge() { let lc = LlmConcurrency::new(3, 2); // 默认 permits=3:可连续 3 个 Some,第 4 个 None。 let sub = "conv-1::sub-1"; let p1 = lc.acquire_per_sub_flow(sub).await; let p2 = lc.acquire_per_sub_flow(sub).await; let p3 = lc.acquire_per_sub_flow(sub).await; let p4 = lc.acquire_per_sub_flow(sub).await; assert!(p1.is_some() && p2.is_some() && p3.is_some()); assert!(p4.is_none(), "默认 permits=3,第 4 个应返回 None"); drop(p1); drop(p2); drop(p3); drop(p4); // 热改 permits=2 + 清空 HashMap(软收敛)。 lc.set_per_sub_permits(2).await; // 新 sub_flow(HashMap 已清空)用新 permits=2:第 3 个 None。 let q1 = lc.acquire_per_sub_flow(sub).await; let q2 = lc.acquire_per_sub_flow(sub).await; let q3 = lc.acquire_per_sub_flow(sub).await; assert!(q1.is_some() && q2.is_some(), "热改 permits=2 后前两个应返回 Some"); assert!(q3.is_none(), "permits=2 第三个应返回 None(软收敛生效)"); } /// CAS 竞态:多并发 acquire 不超卖 — permits=N 时最多 N 个 Some。 /// 用 try_acquire_owned(非阻塞)模拟瞬时并发,N 个 Some 后其余全 None。 #[tokio::test] async fn test_per_sub_flow_no_oversell_under_concurrency() { let lc = LlmConcurrency::new(3, 2); let n = 5usize; lc.set_per_sub_permits(n).await; let sub = "conv-1::sub-1"; // 串行连续 acquire n+3 次:前 n 个 Some,后 3 个 None(非阻塞)。 let mut some_count = 0usize; let mut held: Vec = Vec::with_capacity(n + 3); for _ in 0..(n + 3) { if let Some(permit) = lc.acquire_per_sub_flow(sub).await { some_count += 1; held.push(permit); } } assert_eq!(some_count, n, "permits={n} 时最多 {n} 个 Some,不超卖"); // held 持有 n 个 permit,其余 acquire 全 None(已验证 some_count=n)。 // Drop 全部后,可再 acquire n 个(证明超卖未发生,Semaphore 计数正确)。 drop(held); let mut some_count2 = 0usize; for _ in 0..n { if lc.acquire_per_sub_flow(sub).await.is_some() { some_count2 += 1; } } assert_eq!(some_count2, n, "Drop 全部后应可再 acquire {n} 个"); } }