Files
DevFlow/src-tauri/src/state.rs
绝尘 ed8e2fc155 新增: Phase3 src-tauri 集成 df-tunnel(跨端闭环)
Cargo.toml(df-tunnel+uuid workspace) + state.rs(tunnel 字段 Arc<WsTunnelClient>) + lib.rs spawn tunnel task:device_id(KV/UUID 持久) + relay_url(KV/默认 localhost:8080) + token 固定常量 + subscriber(ai_event_bus→send_raw_event) + on_command(remote_bridge 下行路由) + connect(失败非阻断)。打通桌面↔relay↔miniapp 全双工闭环。cargo EXIT=0。
2026-06-22 03:31:59 +08:00

1295 lines
68 KiB
Rust
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
//! 应用全局状态 — 数据库、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, ProjectRepo, ReleaseRepo,
SettingsRepo, 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<Mutex> 启动 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<String>,
/// embedding 模型名(如 embedding-3 / text-embedding-3-small)
#[serde(default)]
pub embedding_model: Option<String>,
}
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 约束。
///
/// ## globalpermits 默认 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<conv_id, Semaphore>
/// 原应用级单信号量AiSession 单例 + generating 互斥下退化)改为
/// `HashMap<String, Arc<Semaphore>>`每对话一份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<Mutex<Arc<Semaphore>>>` 双层包装——替换内层 Arc 即重建 Semaphore。
/// 已持有旧 permit 的任务不受影响permit 绑定旧 Semaphore软收敛
/// 新请求 lock 后克隆到最新 Arc、自动走新限制。旧 Semaphore 随最后 permit 释放而 drop。
/// per_conv HashMap 内每条 Arc<Semaphore> 不需运行时改 permits对话内并发上限 2 固定),故无重建需求。
///
/// ## F-260614-04: per-provider 层(可选)
/// `per_provider` 为 HashMap<provider_id, Arc<Semaphore>>。调用方(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<sub_flow_id, Arc<Semaphore>>,用于 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<Mutex<Arc<Semaphore>>>,
/// 单对话内并发上限(F-09 B 批5)HashMap<conv_id, Semaphore>,每对话 permits=2。
/// acquire 时按 conv_id 取/建release_conv 在 loop 结束 + 无 pending 时 remove。
per_conv: Arc<Mutex<HashMap<String, Arc<Semaphore>>>>,
/// 每对话内并发上限(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<AtomicUsize>,
/// F-260614-04: per-provider 信号量表。空 = 无 per-provider 限流(单 provider 路径零变化)。
/// Arc<Mutex<HashMap>>:运行时增删 provider 配置时替换/插入,acquire 时 clone Arc。
per_provider: Arc<Mutex<HashMap<String, Arc<Semaphore>>>>,
/// 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<Mutex<HashMap<String, Arc<Semaphore>>>>,
/// 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<AtomicUsize>,
}
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 取/建 Semaphorepermits=当前 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 批5conv 退出清理。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<tokio::sync::OwnedSemaphorePermit> {
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<tokio::sync::OwnedSemaphorePermit> {
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<String, usize>) {
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<Database>,
/// 灵感表 Repo
pub ideas: IdeaRepo,
/// 灵感评估历史表 Repo(追加型审计表,每次 AI 评估追加一行快照,可追溯历史)
pub idea_evaluations: IdeaEvalRepo,
/// 项目表 Repo
pub projects: ProjectRepo,
/// 任务表 Repo
pub tasks: TaskRepo,
/// 发布表 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<NodeRegistry>,
// ── 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<AiToolRegistry>,
/// Input Augmentation 层 Mention Resolver 注册表(核心设计2):
/// 按 MentionRef kind 分发(project/task/idea/skill),resolve_all 单条失败不阻断整批。
/// 启动期 build 一次性注册四 resolver(持 db Arc 访问 repo),通过 Arc 共享只读。
pub resolvers: Arc<ResolverRegistry>,
/// AI 会话状态
pub ai_session: Arc<Mutex<AiSession>>,
/// 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<Mutex<KnowledgeConfig>>,
/// 通用应用设置 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<AtomicUsize>,
// ── 流式对话失败自动重试 ──
/// 流式对话失败自动重试次数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<AtomicUsize>,
// ── 工作流执行状态 ──
/// 工作流执行 → 节点状态机注册表
///
/// run_workflow 创建执行器后注册其 state_machineStateMachine 内部 Arc<Mutex>
/// clone 共享底层 HashMap与下沉到 NodeContext.node_status 的是同一份);
/// cancel_workflow_node IPC 经 execution_id 取出后 set_cancelled
/// 直达运行中 HumanNode 的 is_cancelled 轮询。执行完成(成功/失败)后移除条目。
pub workflow_state_registry: Arc<Mutex<HashMap<String, StateMachine>>>,
// ── 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<RwLock<AllowedDirs>>,
// ── 跨端隧道(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<df_tunnel::WsTunnelClient>,
}
/// F-260619-03 Phase A/B/C: AI 工具文件访问授权目录白名单
///
/// - Phase A:`persistent`(持久化白名单,从 Settings KV `allowed_dirs` 加载,
/// JSON 数组 `["E:/wk-lab/u-abc"]`)。`resolve_workspace_path` 校验时:
/// - workspace_root 始终视为已授权(向后兼容,默认根)
/// - 任一 persistent 目录 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<PathBuf>,
/// F-260619-03 Phase B: 会话级临时授权目录(弹窗"仅本次"写入,进程级内存)。
/// 仅当前活跃会话生效,切换/新建/删除会话由 clear_session_allowed_dirs 清空(不落库)。
/// handler 闭包与 process_tool_calls 预校验均读此字段,确保两端授权判定一致。
pub session: HashSet<PathBuf>,
}
impl AllowedDirs {
/// Settings KV key:F-260619-03 Phase A 持久化白名单(JSON 字符串数组)
pub const SETTINGS_KEY: &'static str = "allowed_dirs";
/// 仅含 workspace_root 的默认白名单(向后兼容:无 allowed_dirs 时行为不变)。
/// 用于 mod.rs trust_key_for 等无白名单上下文的旧路径(零回归)。
pub fn default_with_root() -> Self {
let mut set = HashSet::new();
set.insert(workspace_root_path());
Self { persistent: set, session: HashSet::new() }
}
/// 路径是否被授权:persistent 或 session 白名单任一 starts_with 命中即放行。
///
/// **workspace_root 默认在 persistent**(default_with_root 初始插入,reload_allowed_dirs
/// KV 未配时保持),故首次访问工程根免授权(向后兼容)。用户从白名单删除工程根后,
/// 工程根也需授权(动态白名单完整语义,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))
}
}
/// 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 新建)。返回三态决策:
/// - 路径规范化(去 .. / 锚定 workspace_root)后,若命中黑名单 → `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 {
// 规范化:绝对路径原样,相对路径锚定 workspace_root(与 resolve_workspace_path_impl 一致)
let resolved = if Path::new(raw_path).is_absolute() {
PathBuf::from(raw_path)
} else {
workspace_root_path().join(raw_path)
};
// 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(|| workspace_root_path());
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 根目录(src-tauri 上两级,与 tool_registry::workspace_root 同源)。
///
/// ⚠️ 已知限制:env!("CARGO_MANIFEST_DIR") 是编译期写死编译机源码路径,打包分发后用户机器
/// 无此路径 → workspace_root 失效。运行期动态方案(current_exe/用户项目绑定)见待决策.md
/// (workspace_root 分发适配)。当前仅开发机自用有效。
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) -> Result<Self> {
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_with_root()));
let ai_tools = Arc::new(crate::commands::ai::build_ai_tool_registry(&db, &allowed_dirs));
// Input Augmentation 层(核心设计2):ResolverRegistry 启动期注册四 resolver。
// resolver 持 Arc<Database>(非 AppState,避免循环依赖:AppState 持 Arc<ResolverRegistry>),
// 在 struct 字段 `db` move 前 clone 注入,与 build_ai_tool_registry 同侧。
let resolvers = Arc::new(crate::commands::ai::build_resolver_registry(db.clone()));
// build_registry 需注入 Arc<Database>(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),
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())),
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<provider_id, global_cap>,让 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(仅 workspace_root)。
state.reload_allowed_dirs().await;
// P0(设置走查):从 Settings KV 恢复持久化知识库配置覆盖 default(防 8 项重启全丢)。
state.reload_knowledge_config().await;
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<String, usize> = 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::<KnowledgeConfig>(&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(仅 workspace_root,向后兼容:is_authorized 始终放行 workspace_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<String> = 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<String> = {
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();
// KV 未配(空)时保留 workspace_root(向后兼容:默认根免授权)
if all_dirs.is_empty() {
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);
*guard = AllowedDirs { persistent: set, session: preserved_session };
}
/// 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<String>) -> Result<Vec<String>> {
// 去空 + 去重(保留顺序,前端展示友好)+ 黑名单预校验(防持久化系统敏感目录,
// 纵深防御:即便写入,is_authorized 运行时黑名单也兜底拒,预校验保证白名单 UI 洁净)
let mut seen = HashSet::new();
let cleaned: Vec<String> = 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<String> {
// 不返回内部 workspace_root(隐式授权,is_authorized 经 persistent 放行):
// 避免把内部根目录暴露到前端用户管理列表。前后端均不再各自猜测 root 形态。
let root = workspace_root_path();
let guard = self.allowed_dirs.read().await;
let mut out: Vec<String> = guard
.persistent
.iter()
.filter(|p| **p != root)
.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<Vec<String>> {
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: 路径字符串等价比较(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 触达无审批 shellR-PD-2。DevFlow 工作流当前为纯演示
/// 功能(前端唯一构造点 ProjectDetail.vue demoDag 三步 echo无真实构建/部署脚本
/// 需求。需要脚本执行能力时新建独立 BuildNode白名单 + 项目目录锚定 + 复用 AI 工具
/// RiskLevel 审批链),而非回头启用 ScriptNode + 黑名单。
/// 详见 docs/02-架构设计/专项设计/工作流脚本执行边界-2026-06-15.md。
fn build_registry(db: Arc<Database>) -> 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<Database>,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<Database> 在此构造时注入(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 allowed = AllowedDirs::default_with_root();
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() };
// 精确命中 + 子路径 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() };
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 };
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 };
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() };
// 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() };
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 allowed = AllowedDirs::default_with_root();
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 };
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 allowed = AllowedDirs::default_with_root();
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<tokio::sync::OwnedSemaphorePermit> = 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} 个");
}
}