//! 应用全局状态 — 数据库、Repo、事件总线、节点注册表、AI 会话 //! //! 子模块(按职责拆分,治本:原 1400+ 行单文件 → 按域拆分): //! - [`allowed_dirs`]:文件访问授权白名单 + 路径校验 + 黑名单 //! - [`knowledge_config`]:知识库行为配置 + 提炼触发方式 + KV key 常量 //! - [`llm_concurrency`]:LLM 调用并发控制(双层 Semaphore + 可选 per-provider 层) //! //! 本文件保留 AppState struct + init + 应用级辅助函数。 // 子模块声明 + 重新导出(外部消费方 `use crate::state::{AllowedDirs, LlmConcurrency, ...}` // 路径不变,零改动透明继承)。#[allow(unused)] 是因为部分符号在 impl 块内使用,不是 use 语句导入。 #[allow(unused)] mod allowed_dirs; #[allow(unused)] mod knowledge_config; #[allow(unused)] mod llm_concurrency; // 对外公共符号(pub use):外部 crate/模块经 `crate::state::X` 路径消费 pub use allowed_dirs::{ AllowedDirs, PathAuthDecision, check_path_authorization, is_in_system_blacklist, paths_eq, strip_verbatim, }; pub use knowledge_config::{ ExtractTrigger, KnowledgeConfig, KNOWLEDGE_CONFIG_KEY, APPROVAL_TIMEOUT_KEY, }; pub use llm_concurrency::LlmConcurrency; use std::collections::{HashMap, HashSet}; use std::path::{Path, PathBuf}; use std::sync::Arc; use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering}; use anyhow::Result; use tokio::sync::{Mutex, RwLock}; use df_ai::ai_tools::AiToolRegistry; use df_storage::crud::{ AiConversationRepo, AiMessageRepo, AiProviderRepo, AiToolExecutionRepo, IdeaEvalRepo, IdeaRepo, KnowledgeEventsRepo, KnowledgeRepo, NodeExecutionRepo, ProjectEventRepo, ProjectModuleRepo, ProjectRepo, ProjectServiceRepo, ReleaseRepo, SettingsRepo, TaskLinkRepo, TaskRepo, WorkflowRepo, ModuleDependencyRepo, }; 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; /// 应用全局状态 — 通过 `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(工程系统 V34,project_modules)。 /// 项目多工程(每个工程独立代码仓库):Monorepo 多仓库 / 微服务 / 前后端分离。 /// 记录工程元数据(路径/Git 地址/技术栈);Git 状态(分支/改动/提交)实时派生不存表。 pub project_modules: ProjectModuleRepo, /// 工程依赖关系 Repo(V35 module_dependencies,依赖图边数据)。 pub module_dependencies: ModuleDependencyRepo, /// 发布表 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, // ── 审批超时配置 ── /// 审批超时分钟数(前端 Settings 配置 → AppState 字段 → try_continue 入口读取)。 /// 待审批超过此时长后自动取消,避免用户离开后会话永久卡死。 /// 默认 15 分钟;0 表示禁用超时(不推荐,会导致卡死无自愈)。 pub approval_timeout_minutes: 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, } /// 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), project_modules: ProjectModuleRepo::new(&db), module_dependencies: ModuleDependencyRepo::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, )), approval_timeout_minutes: Arc::new(AtomicU64::new(15)), 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; // 从 Settings KV 恢复审批超时配置(默认 15 分钟,0=禁用超时)。 state.reload_approval_timeout().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), } } /// 从 Settings KV 恢复审批超时分钟数(默认 15 分钟,0=禁用)。 /// /// KV key `df-approval-timeout` 与前端 AdvancedSection.vue 共用,值为毫秒数(0/300000/900000/...)。 /// 启动时读取用户配置,失败或首次启动无配置时保持默认 15 分钟。 pub async fn reload_approval_timeout(&self) { match self.settings.get(APPROVAL_TIMEOUT_KEY).await { Ok(Some(json)) => { // 兼容两种格式:JSON 数字或字符串(纯数字字串) let ms: u64 = if let Ok(v) = serde_json::from_str::(&json) { v } else if let Ok(s) = serde_json::from_str::(&json) { s.parse().unwrap_or(900_000) } else { tracing::warn!("[APPROVAL-TIMEOUT] KV 格式非法,保持 default: {}", json); 900_000 }; // ms → minutes(向上取整,0 保留为 0 表示禁用) let mins = if ms == 0 { 0 } else { (ms + 59_999) / 60_000 }; self.approval_timeout_minutes.store(mins, Ordering::SeqCst); tracing::info!("[APPROVAL-TIMEOUT] 启动恢复: {} 分钟 (KV ms={})", mins, ms); } Ok(None) => {} // 首次启动保持默认 15 Err(e) => tracing::warn!("[APPROVAL-TIMEOUT] 读 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(CARGO_MANIFEST_DIR)。 // 原因:该路径在编译机上有效,分发到其他机器后指向不存在的目录,硬塞反成脏白名单。 // 工程根授权完全靠以下两途径(均已就绪): // 1. projects.bind_directory:用户绑定项目时自动授权(已在上方 project_dirs 合并) // 2. Settings 页 allowed_dirs:用户手动添加持久化白名单(已在上方 kv_dirs 合并) // dev 自用场景:开发机运行时 projects.path 通常含本工程,自然命中白名单。 // workspace_root_path() 函数保留:旧 .trash 迁移(init 中一次性)仍需读取。 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; // 返回内存白名单的规范化字符串列表(供 Settings IPC ai_get_allowed_dirs 回显)。 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 持久化 + 同步内存。 /// 返回持久化后规范化列表。 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(); } } /// 构建节点注册表 — 注册内置节点 /// /// 注意:不使用 `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-260620: strip_verbatim 比对侧收口 + 黑名单增强测试 // 三方审查(安全/UX/跨端)交叉印证:strip_verbatim 仅写入侧调用,比对侧遗漏致误弹窗。 // ============================================================ /// 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")); } // ============================================================ // 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} 个"); } }