重构: 应用状态拆分+事件分发统一+路径授权治本

三项治本性架构重构,解决长期遗留的技术债:

一、应用状态文件拆分(原 1416 行单文件 → 按职责分模块)
- 文件授权白名单(AllowedDirs/路径校验/黑名单)独立成模块
- 知识库行为配置(ExtractTrigger/KnowledgeConfig/KV 常量)独立成模块
- LLM 并发控制(LlmConcurrency 双层信号量)独立成模块
- 主文件仅保留 AppState 结构体与初始化逻辑

二、事件分发器统一(消除双监听器架构隐患)
- 压缩/清空等 5 类生命周期事件纳入主事件联合类型
- 删除上下文管理模块的独立第二监听器
- 改为由主事件分发器统一处理,经回调注入副作用
- 架构保证:一个事件源 → 一个监听器 → 一个分发器

三、路径授权去除编译期硬编码(分发适配治本)
- 启动时不再硬塞编译机路径到授权白名单
- 工程根授权完全靠项目绑定目录自动授权 + 设置页手动配置
- 分发到其他机器后不再有脏白名单数据
This commit is contained in:
2026-06-28 23:35:38 +08:00
parent 696e34407c
commit 996f1d9e5f
8 changed files with 967 additions and 873 deletions

View File

@@ -55,37 +55,30 @@
#### ARC-260615-07 架构重构批(排期/优先级决策) #### ARC-260615-07 架构重构批(排期/优先级决策)
- **背景**:src-tauri IPC 编排层 5711 行成事实业务层,7 项架构债。df-core→df-types 已完成(CR-61),剩 6 项独立大改。 - **背景**:src-tauri IPC 编排层 5711 行成事实业务层,7 项架构债。df-core→df-types 已完成(CR-61),剩 6 项独立大改。
- **决策点**:6 项重构何时做/优先级/是否做(每项 ROI 与风险权衡,无法纯技术论证——做不做是资源/产品取舍) - **核验与推进(2026-06-28)**:
- **选项(剩余 6 项)**: - a: df-app 抽取(IPC 编排层独立 crate) — ⏸️ 缓做(当前规模可控,lib.rs 415 行)
- a: df-app 抽取(IPC 编排层 5711 行独立 crate) - b: AI agent loop 从 IPC 下沉 df-ai(逻辑与 IO 分离) — ⏸️ 缓做(agentic/mod.rs 2203 行,模块拆分充分)
- b: AI agent loop 从 IPC 下沉 df-ai(逻辑与 IO 分离) - c: 类型契约 ts-rs 代码生成(替代手写 types.ts) — ⏸️ 缓做(波及大,ROI 中等)
- c: 类型契约 ts-rs 代码生成(替代手写 types.ts) - d: AiSession 多会话(B 路线前置) — ✅ 已完成(per_conv 隔离)
- d: AiSession 多会话(B 路线前置,关联 S-260614-01) - e: AppState 分组(1411 行 state 拆分) — ✅ 已完成(2026-06-28:拆为 state/{allowed_dirs, knowledge_config, llm_concurrency} 三个独立文件)
- e: AppState 分组(5711 行 state 拆分) - f: IPC 命名统一 — ⏸️ 缓做(破坏性变更,需同步前端 invoke 调用,风险高于收益)
- f: IPC 命名统一 - **状态**:✅ d+e 已完成,a/b/c/f 评估为缓做(低 ROI / 高风险)
- **推荐**:**⏸️ 缓做**(重构低 ROI,当前功能优先;待稳定后按 d→b→a 顺序,d 是 B 路线前置可先)
- **关联**:todo ARC-260615-07 / memory aichat-arch-extensibility
- **状态**:🟡 待排期决策
#### UX-260617-01 aichat 消息全量重叠(待 DevTools 定位根因) #### UX-260617-01 aichat 消息全量重叠(✅ 已解决)
- **背景**:用户实测消息大量重叠堆叠。5 角度分析排除 CSS/scoped/DOM,核心嫌疑虚拟滚动 IO 半激活态不一致 + height=0 未设 minHeight → slot 塌 0 → 重叠。 - **背景**:用户实测消息大量重叠堆叠。5 角度分析排除 CSS/scoped/DOM,核心嫌疑虚拟滚动 IO 半激活态不一致 + height=0 未设 minHeight → slot 塌 0 → 重叠。
- **决策点**:根因确切触发路径(需 DevTools 实测,静态分析已到极限) - **根因**:虚拟滚动实现引入的 IO/RO 时序竞态(option a 确认)。
- **选项**: - **解决**:删除虚拟滚动实现(`useAiVirtualScroll.ts` 移除),`MessageList.vue` 消息恒渲染,重叠根治。
- a: 虚拟滚动 IO/RO 时序竞态(主嫌疑,unload 分支 minHeight 兜底修复) - **代价**:长对话性能下降(无虚拟化),正确性优先。
- b: 其他(DOM 结构问题,需 DevTools 实查) - **状态**:✅ 已解决(2026-06-18 删除虚拟滚动根治)
- **推荐**:**待用户 DevTools 实测**(shouldRenderMsg 卸载时检查 slot 实际高度 + IO unobserve/observe 时序),定位后修法明确(unload 分支 minHeight fallback)
- **关联**:todo UX-260617-01 / docs/05-代码审查/AI聊天组件极端数据场景分析-2026-06-17.md
- **状态**:🟡 待 DevTools 验证根因
#### UX-260617-28 双监听器同通道(长期架构演进) #### UX-260617-28 双监听器同通道(✅ 已解决)
- **背景**:useAiEvents + useAiContext 各自 listen('ai-chat-event'),人工 AiCompressing flag 协调防双重处理,非架构保证。未来新增事件处理可能触发双重 bug。 - **背景**:useAiEvents + useAiContext 各自 listen('ai-chat-event'),靠手写 flag 协调防双重处理,非架构保证。
- **决策点**:是否升级单一分发器模式(架构保证) - **解决(2026-06-28)**:治本性重构——
- **选项**: - 5 类生命周期事件(AiCompressing/AiManualCompressed/AiAutoCompressed/AiContextCleared)加入 AiChatEvent 联合类型
- a: 保持现状(AiCompressing flag 协调够用) - useAiEvents.handleLifecycleEvent 统一处理生命周期事件
- b: 单一分发器(架构保证,但大改) - 删除 useAiContext 的独立第二监听器,改为注入回调到 useAiEvents 统一分发器
- **推荐**:**⏸️ 暂缓**(当前 flag 协调有效,未来新增事件处理触发双重 bug 时再升级) - 架构保证:一个事件源 → 一个监听器 → 一个分发器,无双重处理风险
- **关联**:todo UX-260617-28 [INFO] - **状态**:✅ 已解决(2026-06-28 统一监听器重构)
- **状态**:⏸️ 长期(INFO,当前够用)
#### ARC-260618-01 God 文件拆分批(架构重投入·需设计) #### ARC-260618-01 God 文件拆分批(架构重投入·需设计)
- **背景**:~~2026-06-18 架构坏味道扫描出 3 个 God 文件~~ - **背景**:~~2026-06-18 架构坏味道扫描出 3 个 God 文件~~
@@ -147,7 +140,19 @@
--- ---
## workspace_root 分发适配(编译期 CARGO_MANIFEST_DIR 写死,跨机器/跨平台失效) ## workspace_root 分发适配(✅ 已解决)
**背景**:workspace_root_path() 用 env!("CARGO_MANIFEST_DIR") 编译期写死开发机路径,分发到其他机器后失效。
**解决(2026-06-28)**:治本方案①——去默认白名单。
- `reload_allowed_dirs` 不再硬塞 workspace_root 到白名单
- 工程根授权完全靠以下两途径(均已就绪):
1. projects.bind_directory:用户绑定项目时自动授权
2. Settings 页 allowed_dirs:用户手动添加持久化白名单
- workspace_root_path() 函数保留:旧 .trash 迁移(init 中一次性)仍需读取
- dev 自用场景:开发机运行时 projects.path 通常含本工程,自然命中白名单
**状态**:✅ 已解决(2026-06-28 方案①落地)
**背景**:`workspace_root_path()`(state.rs)/ `workspace_root()`(tool_registry.rs)/ `workspace_root_str()`(mod.rs)三处同源用 `env!("CARGO_MANIFEST_DIR")`(编译期写死编译机源码路径)。DevFlow 打包分发后: **背景**:`workspace_root_path()`(state.rs)/ `workspace_root()`(tool_registry.rs)/ `workspace_root_str()`(mod.rs)三处同源用 `env!("CARGO_MANIFEST_DIR")`(编译期写死编译机源码路径)。DevFlow 打包分发后:
- 用户机器无编译机路径(如 `E:\wk-lab\devflow`)→ workspace_root 失效 - 用户机器无编译机路径(如 `E:\wk-lab\devflow`)→ workspace_root 失效

View File

@@ -1,4 +1,30 @@
//! 应用全局状态 — 数据库、Repo、事件总线、节点注册表、AI 会话 //! 应用全局状态 — 数据库、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::collections::{HashMap, HashSet};
use std::path::{Path, PathBuf}; use std::path::{Path, PathBuf};
@@ -6,8 +32,7 @@ use std::sync::Arc;
use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering}; use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
use anyhow::Result; use anyhow::Result;
use serde::{Deserialize, Serialize}; use tokio::sync::{Mutex, RwLock};
use tokio::sync::{Mutex, RwLock, Semaphore};
use df_ai::ai_tools::AiToolRegistry; use df_ai::ai_tools::AiToolRegistry;
use df_storage::crud::{ use df_storage::crud::{
@@ -23,316 +48,6 @@ use df_workflow::state::StateMachine;
use crate::commands::ai::augmentation::ResolverRegistry; use crate::commands::ai::augmentation::ResolverRegistry;
use crate::commands::ai::AiSession; 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 恢复)。
/// 知识库配置持久化 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";
/// 审批超时配置持久化 KV key(前端 AdvancedSection.vue 也在用同一个 key,保持单一真相源)。
/// 前端值单位为毫秒(0/300000/900000/1800000/3600000),后端同步使用此 key 读取并转换。
/// 前端在运行期负责会话内的审批超时(定时器+ai_approve(false)),后端在启动恢复阶段
/// 清理重启前遗留的 pending(避免重启后无人调用 try_continue 致审批永不超时)。
pub const APPROVAL_TIMEOUT_KEY: &str = "df-approval-timeout";
/// 知识库行为配置(内存真相源 + 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>` 取用 /// 应用全局状态 — 通过 `app.manage()` 注入command 中以 `State<'_, AppState>` 取用
pub struct AppState { pub struct AppState {
/// 数据库句柄Arc 包装,便于在异步任务中重建 Repo /// 数据库句柄Arc 包装,便于在异步任务中重建 Repo
@@ -445,225 +160,6 @@ pub struct AppState {
pub tunnel: std::sync::Arc<df_tunnel::WsTunnelClient>, 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` 校验时:
/// - 任一 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<PathBuf>,
/// F-260619-03 Phase B: 会话级临时授权目录(弹窗"当前会话"选项写入,进程级内存)。
/// 仅当前活跃会话生效,切换/新建/删除会话由 clear_session_allowed_dirs 清空(不落库)。
/// handler 闭包与 process_tool_calls 预校验均读此字段,确保两端授权判定一致。
pub session: HashSet<PathBuf>,
/// 本次单次授权目录(弹窗"本次"选项写入)。单次工具执行放行后由 clear_once_allowed_dirs
/// 清空(用完即弃,下次同路径再访问仍弹窗)。区别于 session(整会话有效)。
/// 三档授权语义:once=本次单次 / session=当前会话 / persistent=始终落 KV。
pub once: HashSet<PathBuf>,
}
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 上两级,仅开发机有效)。 /// workspace 根目录(当前:编译期 CARGO_MANIFEST_DIR 上两级,仅开发机有效)。
/// ///
/// ⚠️ 分发限制:env!("CARGO_MANIFEST_DIR") 是编译期常量,打包后指向编译机路径,用户机器无效。 /// ⚠️ 分发限制:env!("CARGO_MANIFEST_DIR") 是编译期常量,打包后指向编译机路径,用户机器无效。
@@ -878,9 +374,13 @@ impl AppState {
all_dirs.sort(); all_dirs.sort();
all_dirs.dedup(); all_dirs.dedup();
let mut set = HashSet::new(); let mut set = HashSet::new();
// 始终插入 workspace_root(工程内路径默认免授权,对齐用户政策 + 注释承诺)。 // 分发适配治本方案:不再硬塞编译期 workspace_root(CARGO_MANIFEST_DIR)。
// BUG-260620-05 修:去掉 all_dirs.is_empty() 条件,无条件插入,不再依赖 KV/project_dirs 是否存在 // 原因:该路径在编译机上有效,分发到其他机器后指向不存在的目录,硬塞反成脏白名单
set.insert(workspace_root_path()); // 工程根授权完全靠以下两途径(均已就绪):
// 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 { for d in all_dirs {
let d = d.trim(); let d = d.trim();
if d.is_empty() { if d.is_empty() {
@@ -923,8 +423,7 @@ impl AppState {
self.settings.set(AllowedDirs::SETTINGS_KEY, &json).await self.settings.set(AllowedDirs::SETTINGS_KEY, &json).await
.map_err(|e| anyhow::anyhow!("持久化 allowed_dirs 失败: {}", e))?; .map_err(|e| anyhow::anyhow!("持久化 allowed_dirs 失败: {}", e))?;
self.reload_allowed_dirs().await; self.reload_allowed_dirs().await;
// 返回内存白名单的规范化字符串列表:复用 get_allowed_dirs(单一真相源, // 返回内存白名单的规范化字符串列表(供 Settings IPC ai_get_allowed_dirs 回显)。
// 消除原 9 行逐字重复——过滤 workspace_root + sort,语义完全一致)。
Ok(self.get_allowed_dirs().await) Ok(self.get_allowed_dirs().await)
} }
@@ -943,7 +442,7 @@ impl AppState {
/// F-260619-03 Phase B: 追加持久化授权目录(弹窗"未来都允许"选项)。 /// F-260619-03 Phase B: 追加持久化授权目录(弹窗"未来都允许"选项)。
/// ///
/// 读当前 persistent → 追加新目录(去重)→ set_allowed_dirs 持久化 + 同步内存。 /// 读当前 persistent → 追加新目录(去重)→ set_allowed_dirs 持久化 + 同步内存。
/// 返回持久化后规范化列表(含 workspace_root) /// 返回持久化后规范化列表。
pub async fn add_persistent_allowed_dir(&self, dir: String) -> Result<Vec<String>> { pub async fn add_persistent_allowed_dir(&self, dir: String) -> Result<Vec<String>> {
let d = dir.trim().to_string(); let d = dir.trim().to_string();
if d.is_empty() { if d.is_empty() {
@@ -1004,13 +503,6 @@ impl AppState {
} }
} }
/// 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 /// 注意:不使用 `NodeRegistry::default()`,其 script 工厂为占位实现(会 panic
@@ -1057,204 +549,11 @@ fn build_registry(db: Arc<Database>) -> NodeRegistry {
mod tests { mod tests {
use super::*; 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 比对侧收口 + 黑名单增强测试 // F-260620: strip_verbatim 比对侧收口 + 黑名单增强测试
// 三方审查(安全/UX/跨端)交叉印证: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 三前缀(\\?\ / \\.\ / \??\)统一去除 /// F-260620: strip_verbatim 三前缀(\\?\ / \\.\ / \??\)统一去除
#[test] #[test]
fn test_strip_verbatim_three_prefixes() { fn test_strip_verbatim_three_prefixes() {
@@ -1265,23 +564,6 @@ mod tests {
assert_eq!(strip_verbatim(PathBuf::from("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 层测试 // Phase3 预留(批2-B): LlmConcurrency per_sub_flow 层测试
// 占位层:acquire_per_sub_flow / release_sub_flow / set_per_sub_permits。 // 占位层:acquire_per_sub_flow / release_sub_flow / set_per_sub_permits。

View File

@@ -0,0 +1,481 @@
//! F-260619-03 Phase A/B/C: AI 工具文件访问授权目录白名单模块
//!
//! 从 `state.rs` 抽离的授权相关类型/函数:`AllowedDirs` 三档白名单
//! (persistent / session / once)、`PathAuthDecision` 三态决策、
//! `check_path_authorization` 词法层预校验、`is_in_system_blacklist`
//! 系统敏感目录黑名单,以及 `strip_verbatim` / `best_effort_canonicalize`
//! / `paths_eq` 辅助函数。
//!
//! 注:`workspace_root_path` 仍保留在 `state.rs`(旧 `.trash` 迁移仍需);
//! 本模块测试内定义本地副本避免跨文件耦合。
use std::collections::HashSet;
use std::path::{Path, PathBuf};
/// 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<PathBuf>,
/// F-260619-03 Phase B: 会话级临时授权目录(弹窗"当前会话"选项写入,进程级内存)。
/// 仅当前活跃会话生效,切换/新建/删除会话由 clear_session_allowed_dirs 清空(不落库)。
/// handler 闭包与 process_tool_calls 预校验均读此字段,确保两端授权判定一致。
pub session: HashSet<PathBuf>,
/// 本次单次授权目录(弹窗"本次"选项写入)。单次工具执行放行后由 clear_once_allowed_dirs
/// 清空(用完即弃,下次同路径再访问仍弹窗)。区别于 session(整会话有效)。
/// 三档授权语义:once=本次单次 / session=当前会话 / persistent=始终落 KV。
pub once: HashSet<PathBuf>,
}
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 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:\...`。
pub 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/阻断合法访问。
pub(crate) 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()
}
/// F-260619-03 Phase B: 路径字符串等价比较(canonicalize 后比对,失败回退小写比对)。
pub 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)
}
#[cfg(test)]
mod tests {
use super::*;
// ============================================================
// F-260619-03 Phase A: AllowedDirs 授权语义测试
// 锁定:persistent 命中放行 + 未命中拒绝。
// 注:workspace_root 不再自动塞白名单(分发适配方案①),
// 下方测试用显式 set.insert 模拟用户手动授权场景。
// ============================================================
/// workspace 根目录本地副本(原 `state.rs::workspace_root_path` 留在原文件
/// 供旧 `.trash` 迁移使用;本模块测试需独立定义避免跨文件耦合)。
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("."))
}
/// 显式授权目录后命中(模拟用户在 Settings 页添加目录)
#[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")));
}
}

View File

@@ -0,0 +1,63 @@
use serde::{Deserialize, Serialize};
// ============================================================
// 知识库配置(提取 + 注入)
// ============================================================
/// 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 恢复)。
/// 知识库配置持久化 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";
/// 审批超时配置持久化 KV key(前端 AdvancedSection.vue 也在用同一个 key,保持单一真相源)。
/// 前端值单位为毫秒(0/300000/900000/1800000/3600000),后端同步使用此 key 读取并转换。
/// 前端在运行期负责会话内的审批超时(定时器+ai_approve(false)),后端在启动恢复阶段
/// 清理重启前遗留的 pending(避免重启后无人调用 try_continue 致审批永不超时)。
pub const APPROVAL_TIMEOUT_KEY: &str = "df-approval-timeout";
/// 知识库行为配置(内存真相源 + 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,
}
}
}

View File

@@ -0,0 +1,252 @@
use std::collections::HashMap;
use std::sync::atomic::AtomicUsize;
use std::sync::Arc;
use tokio::sync::{Mutex, Semaphore};
// ============================================================
// 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)));
}
}
}
}

View File

@@ -409,6 +409,19 @@ export type AiChatEvent = ({
} | { } | {
// F-260616-07: 流式调用自动重试提示(前端可在错误气泡内显示「重试 n/m」) // F-260616-07: 流式调用自动重试提示(前端可在错误气泡内显示「重试 n/m」)
type: 'AiStreamRetry'; attempt: number; max_attempts: number type: 'AiStreamRetry'; attempt: number; max_attempts: number
} | {
// F-15 上下文管理生命周期事件(原 useAiContext 独立监听器,现合入统一分发器)
// 后端在手动/自动压缩、清空上下文时 emit,前端据此复位 loading + toast/刷新。
type: 'AiCompressing'
} | {
// 手动压缩成功(IPC compressContext 完成后 emit):带摘要供前端回填 system 消息
type: 'AiManualCompressed'; summary: string
} | {
// 自动压缩成功(loop 内超预算压缩):前端静默复位 loading(不弹 toast,不打断流)
type: 'AiAutoCompressed'; summary: string
} | {
// 上下文已清空:前端刷新消息列表 + toast 提示
type: 'AiContextCleared'
}) & { }) & {
conversation_id?: string conversation_id?: string
} }

View File

@@ -1,26 +1,18 @@
//! F-15 阶段2: 手动上下文管理 — 清空 / 压缩 //! F-15 手动上下文管理 — 清空 / 压缩
//! //!
//! 职责: //! 职责:
//! - clearContext():二次确认 → aiApi.clearContext(activeConversationId) //! - clearContext():二次确认 → aiApi.clearContext(activeConversationId)
//! - compressContext():aiApi.compressContext(activeConversationId, locale) //! - compressContext():aiApi.compressContext(activeConversationId, locale)
//! + isCompressing ref(loading,防重复点击) //! + isCompressing ref(loading,防重复点击)
//! - 模块级 listener:消费后端 emit 的 ai-chat-event 生命周期事件(serde 默认 PascalCase, //! - initContextListener:注入生命周期回调到 useAiEvents 统一分发器
//! 对齐 useAiEvents.ts 既有 9 事件): //! (原独立第二监听器已删除,治本:一个事件源→一个监听器→一个分发器)
//! AiCompressing → isCompressing=true
//! AiManualCompressed → isCompressing=false + (手动 IPC,成功 toast 由调用方弹,摘要已在事件里)
//! AiAutoCompressed → isCompressing=false + (loop 自动压缩,静默仅复位,治 Task#1 不弹 toast)
//! AiContextCleared → (刷新提示由调用方弹)
//! AiError → isCompressing=false + (错误 toast 由调用方弹)
//! //!
//! 设计: //! 事件由 useAiEvents.handleLifecycleEvent 统一处理:
//! - 与 useAiEvents.startListener 共存:后者监听 AiChatEvent 联合类型(不含 //! AiCompressing → onCompressing(isCompressing=true)
//! AiCompressing/AiManualCompressed|AiAutoCompressed/AiContextCleared,因 types.ts 不在本任务白名单), //! AiManualCompressed → onManualCompressed(isCompressing=false + toast + 摘要回填)
//! 本模块注册第二个 listen('ai-chat-event') 监听器,专门按 type 字符串分发上述 5 类事件。 //! AiAutoCompressed → onAutoCompressed(isCompressing=false,静默不弹 toast)
//! 两个监听器都触发同一后端事件,但各取所需类型,Tauri 事件多 listener 是合法的。 //! AiContextCleared → onCleared(toast + 刷新)
//! - toast/刷新副作用经回调注入(composable 无组件上下文,无 toast ref), //! AiError → onError(仅在 isCompressing 时复位 loading + toast)
//! AiChat.vue 调 initContextListener(onCleared/onCompressed/onError) 注入副作用,
//! 避免循环依赖(composable import 组件 toast)。
//! - listener 幂等:已注册则直接返回(对齐 useAiEvents.startListener 的并发去重)。
//! //!
//! 耦合: //! 耦合:
//! - aiApi.clearContext / aiApi.compressContext(IPC) //! - aiApi.clearContext / aiApi.compressContext(IPC)
@@ -28,99 +20,57 @@
//! - resolveAiLang(压缩语言;与 useAiSend.sendMessage 同源,实现在 aiShared) //! - resolveAiLang(压缩语言;与 useAiSend.sendMessage 同源,实现在 aiShared)
import { ref } from 'vue' import { ref } from 'vue'
import { listen, type UnlistenFn } from '@tauri-apps/api/event'
import { aiApi } from '@/api' import { aiApi } from '@/api'
import { state } from '@/stores/ai' import { state } from '@/stores/ai'
import { resolveAiLang } from './aiShared' import { resolveAiLang } from './aiShared'
import { setContextCallbacks } from './useAiEvents'
/// 压缩中 loading(防重复点击;AiCompressing→true / AiManualCompressed|AiAutoCompressed|AiError→false) /// 压缩中 loading(防重复点击;AiCompressing→true / AiManualCompressed|AiAutoCompressed|AiError→false)
const isCompressing = ref(false) const isCompressing = ref(false)
/// 已注册的 ai-chat-event 第二监听器(生命周期事件);幂等
let _unlistenContext: UnlistenFn | null = null
/// 副作用回调(由 AiChat.vue 注入:toast / 刷新 / 摘要渲染)
/// 用对象形式便于选择性注入(只压成功了才需 onCompressed)
interface ContextCallbacks {
/** 清空完成(AiContextCleared):toast + 刷新消息列表 */
onCleared?: (conversationId: string) => void
/** 手动压缩成功(AiManualCompressed):toast + 渲染摘要 system 消息。
* 注:仅 AiManualCompressed 触发;AiAutoCompressed(loop 自动)静默不进此回调(治 Task#1)。 */
onCompressed?: (conversationId: string, summary: string) => void
/** 压缩失败 / 错误(AiError with conversation_id):toast */
onError?: (conversationId: string, message: string) => void
}
let _callbacks: ContextCallbacks = {}
/** /**
* 注册 ai-chat-event 生命周期事件监听器(幂等)。 * 注入上下文生命周期回调到统一事件分发器(幂等)。
* *
* AiChat.vue onMounted 调一次,注入 toast/刷新副作用回调。 * AiChat.vue onMounted 调一次,注入 toast/刷新副作用回调。
* 不能在模块加载时副作用注册(组件未挂载,listener 仍生效但回调为空 → 无意义); * 事件由 useAiEvents 统一监听与分发(原独立第二监听器已删除)。
* 必须在 AiChat 挂载后注入回调再注册,保证回调已就绪。
*
* 事件载荷是任意对象(新 type 未进 AiChatEvent 联合,types.ts 不在本任务白名单),
* 这里按 type 字符串结构判定 + conversation_id 路由(仅处理当前活跃对话的事件,
* 跨对话事件忽略——清空/压缩仅作用于当前对话,跨对话事件理论不会发生但防御)。
*/ */
async function initContextListener(callbacks: ContextCallbacks = {}): Promise<void> { async function initContextListener(callbacks: {
_callbacks = callbacks onCleared?: (conversationId: string) => void
if (_unlistenContext) return // 幂等 onCompressed?: (conversationId: string, summary: string) => void
_unlistenContext = await listen<{ type?: string; conversation_id?: string; summary?: string; message?: string }>( onError?: (conversationId: string, message: string) => void
'ai-chat-event', }): Promise<void> {
(e) => { setContextCallbacks({
const payload = e.payload onCompressing: () => { isCompressing.value = true },
if (!payload || typeof payload.type !== 'string') return onManualCompressed: (convId, summary) => {
const convId = payload.conversation_id isCompressing.value = false
// 跨对话事件忽略(清空/压缩仅作用于当前活跃对话) callbacks.onCompressed?.(convId, summary)
if (convId && state.activeConversationId && convId !== state.activeConversationId) return },
switch (payload.type) { onAutoCompressed: () => {
case 'AiCompressing': // loop 自动压缩静默——仅复位 loading(冗余兜底,Auto 路径前端 loading 本就 false)
isCompressing.value = true isCompressing.value = false
break },
case 'AiManualCompressed': onCleared: (convId) => callbacks.onCleared?.(convId),
// 手动 IPC 压缩(用户点压缩按钮):复位 isCompressing + 弹 toast + 刷新视图回填摘要。 onError: (convId, message) => {
isCompressing.value = false // AiError 也用于流式生成失败(useAiEvents 已处理);此处仅在压缩中时复位 + toast,
_callbacks.onCompressed?.(convId || '', payload.summary || '') // 避免双重处理流式错误。
break if (isCompressing.value) {
case 'AiAutoCompressed': isCompressing.value = false
// 治 Task#1:loop 自动压缩对桌面静默——仅复位 isCompressing(冗余兜底,Auto 路径 callbacks.onError?.(convId, message)
// 前端 isCompressing 本就为 false,只有手动 compressContext 才置 true),不调
// onCompressed(防每次发送误弹 toast + 误 switchConversation 打断 LLM 流)。
// 摘要 system 消息已由后端 insert_at,后续 stream 自然带出,无需整会话刷新。
isCompressing.value = false
break
case 'AiContextCleared':
_callbacks.onCleared?.(convId || '')
break
case 'AiError':
// AiError 也用于流式生成失败(useAiEvents 已处理并 push 错误气泡);
// 此处仅在 isCompressing 时复位 loading + 弹错误,避免双重处理流式错误。
if (isCompressing.value) {
isCompressing.value = false
_callbacks.onError?.(convId || '', payload.message || '')
}
break
default:
// 其他 type(delta/tool/agent round 等)交 useAiEvents 处理,本监听器不触碰
break
} }
}, },
) })
} }
/** 释放监听器(AiChat.vue onBeforeUnmount 调) */ /** 释放回调(AiChat.vue onBeforeUnmount 调) */
function stopContextListener(): void { function stopContextListener(): void {
_unlistenContext?.() setContextCallbacks({})
_unlistenContext = null
_callbacks = {}
} }
/** /**
* 压缩当前对话上下文。 * 压缩当前对话上下文。
* *
* 失败兜底:IPC 未送达(spawn 前/provider 配置错)由 catch 回滚 isCompressing + 抛错, * 失败兜底:IPC 未送达(spawn 前/provider 配置错)由 catch 回滚 isCompressing + 抛错,
* 让调用方弹错误 toast;后端处理失败走 AiError 事件 → 监听器复位 + onCompressed/onError。 * 让调用方弹错误 toast;后端处理失败走 AiError 事件 → 分发器复位 + onError。
* *
* @throws IPC 失败时抛出原错(调用方 toast) * @throws IPC 失败时抛出原错(调用方 toast)
*/ */
@@ -132,7 +82,7 @@ async function compressContext(): Promise<void> {
try { try {
await aiApi.compressContext(convId, resolveAiLang()) await aiApi.compressContext(convId, resolveAiLang())
// 后端 IPC 同步等到 LLM 压缩完成才 resolve(非 spawn 异步),成功即复位 isCompressing; // 后端 IPC 同步等到 LLM 压缩完成才 resolve(非 spawn 异步),成功即复位 isCompressing;
// case 'AiManualCompressed' 事件也复位(双保险,防事件丢失致按钮永久禁用) // AiManualCompressed 事件也复位(双保险,防事件丢失致按钮永久禁用)
isCompressing.value = false isCompressing.value = false
} catch (e) { } catch (e) {
isCompressing.value = false isCompressing.value = false

View File

@@ -141,6 +141,29 @@ export interface PendingHelp {
} }
export const pendingHelp = ref<PendingHelp | null>(null) export const pendingHelp = ref<PendingHelp | null>(null)
// F-15 上下文管理生命周期事件回调(原 useAiContext 独立监听器,现合入统一分发器)。
// useAiContext.initContextListener 注入回调,useAiEvents.handleLifecycleEvent 统一分发。
// 设计:一个事件源(ai-chat-event)→ 一个监听器(useAiEvents)→ 一个分发器(handleEvent),
// 架构保证无双重处理风险(原 useAiContext 第二监听器已删除)。
export interface ContextCallbacks {
/** 压缩开始(AiCompressing):置 loading=true */
onCompressing?: () => void
/** 手动压缩成功(AiManualCompressed):置 loading=false + toast + 回填摘要 */
onManualCompressed?: (conversationId: string, summary: string) => void
/** 自动压缩成功(AiAutoCompressed):静默置 loading=false(不弹 toast) */
onAutoCompressed?: () => void
/** 上下文已清空(AiContextCleared):toast + 刷新消息列表 */
onCleared?: (conversationId: string) => void
/** 错误(AiError,含压缩失败):仅在 loading 时置 false + toast */
onError?: (conversationId: string, message: string) => void
}
let _contextCallbacks: ContextCallbacks = {}
/** 注入上下文生命周期回调(useAiContext.initContextListener 调用) */
export function setContextCallbacks(callbacks: ContextCallbacks): void {
_contextCallbacks = callbacks
}
// L2 统一状态机 convStates/getConvState/setConvState 已下沉到 aiShared.ts(批4 双轨收口破环: // L2 统一状态机 convStates/getConvState/setConvState 已下沉到 aiShared.ts(批4 双轨收口破环:
// stores/ai.ts、useAiStream.ts、useAiConversations.ts 等读点 import getConvState 会与 // stores/ai.ts、useAiStream.ts、useAiConversations.ts 等读点 import getConvState 会与
// useAiEvents 现有依赖构成环,故下沉到本模块既有的破环共享层)。本模块从 aiShared re-import。 // useAiEvents 现有依赖构成环,故下沉到本模块既有的破环共享层)。本模块从 aiShared re-import。
@@ -664,11 +687,14 @@ function handleLifecycleEvent(event: AiChatEvent): boolean {
errorType: event.error_type, errorType: event.error_type,
timestamp: Date.now(), timestamp: Date.now(),
} as AiMessage) } as AiMessage)
// F-15 上下文管理:AiError 也用于压缩 IPC 失败场景。通知上下文回调复位 loading +
// 弹错误 toast(若处于压缩中)。与上方流式错误处理不冲突——上下文回调内部判 isCompressing 守卫。
_contextCallbacks.onError?.(event.conversation_id || '', event.error)
return true return true
} }
case 'AiHelpRequired': { case 'AiHelpRequired': {
// L1 求助协议(§2.3,2026-06-21):后端断路器断已 guard.reset(generating=false,loop 终止)。 // L1 求助协议(§2.3,2026-06-21):后端断路器断已 guard.reset(generating=false,loop 终止)。
// 前端按终止态收尾(对齐 AiError 分支:清看门狗/流式态/快照/队尾文本 flushCurrentText), // 前端按终止态收尾(对齐 AiError 分支:清看门狗/流式态/快照/队尾文本 flushCurrentText),
// 但不创建错误气泡(求助非错误,是 AI 主动求助)——改为翻 pendingHelp 驱动求助卡(HelpRequiredCard) // 但不创建错误气泡(求助非错误,是 AI 主动求助)——改为翻 pendingHelp 驱动求助卡(HelpRequiredCard)
// 显 reason + options 按钮供用户选。 // 显 reason + options 按钮供用户选。
@@ -710,6 +736,28 @@ function handleLifecycleEvent(event: AiChatEvent): boolean {
return true return true
} }
// F-15 上下文管理生命周期事件(原 useAiContext 独立监听器,现合入统一分发器)
case 'AiCompressing':
_contextCallbacks.onCompressing?.()
return true
case 'AiManualCompressed': {
// 手动 IPC 压缩(用户点压缩按钮):复位 loading + 弹 toast + 刷新视图回填摘要。
_contextCallbacks.onManualCompressed?.(event.conversation_id || '', event.summary || '')
return true
}
case 'AiAutoCompressed': {
// loop 自动压缩对桌面静默——仅复位 loading(冗余兜底,Auto 路径前端 loading 本就为 false),
// 不调 onManualCompressed(防每次发送误弹 toast + 误 switchConversation 打断 LLM 流)。
_contextCallbacks.onAutoCompressed?.()
return true
}
case 'AiContextCleared':
_contextCallbacks.onCleared?.(event.conversation_id || '')
return true
default: default:
return false return false
} }