Files
DevFlow/src-tauri/src/state.rs
绝尘 7c708a6e01 优化: F-09 batch5调整 删global会话级限流(用户决策不设并发会话上限)
- agentic.rs: 删loop入口acquire_global(会话数不限) + per_conv保留(单对话内限流permits=2)
- state.rs: global回归LLM调用并发限流原义(5处单次调用点+stream_llm防provider429),非会话数上限
用户决策"不设并发会话上限",token暴增接受;主代兜底cargo 0+test98+grep印证
2026-06-19 02:39:39 +08:00

417 lines
21 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;
use std::path::Path;
use std::sync::Arc;
use std::sync::atomic::AtomicUsize;
use anyhow::Result;
use serde::{Deserialize, Serialize};
use tokio::sync::{Mutex, Semaphore};
use df_ai::ai_tools::AiToolRegistry;
use df_storage::crud::{
AiConversationRepo, AiProviderRepo, AiToolExecutionRepo, 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::AiSession;
// ============================================================
// 知识库配置(提取 + 注入)
// ============================================================
/// AI 提炼触发方式
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[serde(rename_all = "snake_case")]
pub enum ExtractTrigger {
/// 对话正常完成时(默认)
OnComplete,
/// 对话闲置 N 秒后
OnIdle,
/// 仅手动按钮触发
ManualOnly,
}
/// 知识库行为配置(存 AppState 内存,前后端通过 IPC 读写)
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct KnowledgeConfig {
/// 提炼总开关,默认 true
pub auto_extract: bool,
/// 提炼触发方式,默认 OnComplete
pub trigger_mode: ExtractTrigger,
/// 最少消息数守卫(防闲聊噪音),默认 4
pub min_messages: u32,
/// 闲置触发超时(ms),默认 30000
pub idle_timeout_ms: u64,
/// 聊天时自动注入相关知识开关,默认 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,
idle_timeout_ms: 30_000,
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)),非运行时强约束。
#[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>>>>,
}
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())),
}
}
/// 取全局 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();
}
/// 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
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
pub ai_tool_executions: AiToolExecutionRepo,
/// AI 工具注册表
pub ai_tools: Arc<AiToolRegistry>,
/// AI 会话状态
pub ai_session: Arc<Mutex<AiSession>>,
// ── 知识库 ──
/// 知识库 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>>>,
}
impl AppState {
/// 初始化应用状态:打开(或创建)数据库并执行迁移,构建各 Repo 与节点注册表
pub async fn init(db_path: &Path) -> Result<Self> {
let db = Arc::new(Database::open(db_path).await?);
let ai_tools = Arc::new(crate::commands::ai::build_ai_tool_registry(&db));
// 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),
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_tool_executions: AiToolExecutionRepo::new(&db),
ai_session: Arc::new(Mutex::new(AiSession::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())),
db,
event_bus: EventBus::new(),
registry,
ai_tools,
};
// 启动恢复:重启前卡 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;
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;
}
}
/// 构建节点注册表 — 注册内置节点
///
/// 注意:不使用 `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_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_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
}