更新ppt
设计哲学:并非所有功能都需要 LLM。通过工作流 + 多节点编排 + Hook 被动调用 LLM 调试错误 + SPJ 事件报错,最大化硬编码执行比例,最小化 Token 消耗。 双支柱架构:任务规划与验证(YAML 编译器 + SPJ + DAG 引擎)+ KG 驱动上下文管理(GraphRAG + 永久知识图谱 + 可溯源检索)
设计哲学:并非所有功能都需要 LLM。通过工作流 + 多节点编排 + Hook 被动调用 LLM 调试错误 + SPJ 事件报错,最大化硬编码执行比例,最小化 Token 消耗。
双支柱架构:任务规划与验证(YAML 编译器 + SPJ + DAG 引擎)+ KG 驱动上下文管理(GraphRAG + 永久知识图谱 + 可溯源检索)
可视化 AI 工作流引擎。基于 DAG 编排,集成 LightRAG 知识图谱、资源感知动态调度、对话式工作流生成、SPJ 自动验证修复。原生适配 openEuler / 鲲鹏 ARM64。
┌──────────────────────┐ │ 用户意图(对话) │ └──────────┬───────────┘ │ ┌──────────────┼──────────────┐ ▼ ▼ ▼ 简单问答 复杂任务 已有经验 (直接用LLM) (编译为YAML) (匹配Skill) │ ┌─────────────┼─────────────┐ ▼ ▼ ▼ 硬编码节点 条件LLM节点 被动修复 (零Token) (仅在需要时) (仅在出错时)
关键原则:
_system
workflowMode
# 一键安装 + 启动全部服务 git clone https://github.com/xhqyt/workflow.git && cd workflow ./setup.sh # 分步操作 ./setup.sh install # 安装 Node + Python 依赖 ./setup.sh start # 启动 RAG → Backend → Frontend ./setup.sh status # 查看各服务运行状态 ./setup.sh stop # 停止全部服务
访问 http://localhost:3000。首次使用在 Settings → Provider 配置 API Key。
http://localhost:3000
环境:Node.js 22+ / Python 3.10+ / Redis 7+ / PostgreSQL 15+(RAG 向量存储,可选)
┌──────────────────────────────────────────────────────────────────┐ │ Frontend (React + Vite + Tailwind) │ │ │ │ ┌──────────┐ ┌──────────┐ ┌─────────┐ ┌────────┐ ┌──────────┐ │ │ │Dashboard │ │ Chat + │ │Settings │ │ Memory │ │ Canvas │ │ │ │实时监控 │ │ Workflow │ │5 Tabs │ │KB CRUD │ │ DAG编辑 │ │ │ └──────────┘ └──────────┘ └─────────┘ └────────┘ └──────────┘ │ │ ↕ ↕ ↕ ↕ ↕ │ │ WS HTTP/SSE HTTP HTTP HTTP │ └──────────────────────────────────────────────────────────────────┘ │ ┌─────────────────────────────┴────────────────────────────────────┐ │ Backend (Node.js / TypeScript / Fastify) │ │ │ │ ┌─────────────────────────────────────────────────────────────┐ │ │ │ API Layer (21 routes) │ │ │ │ agents · workflows · executions · providers · models │ │ │ │ rag · knowledge-base · monitor · settings · schedules │ │ │ │ skills · clawhub · benchmark · heartbeat · llm │ │ │ └─────────────────────────────────────────────────────────────┘ │ │ │ │ │ ┌─────────────────────────────────────────────────────────────┐ │ │ │ Core Engine │ │ │ │ │ │ │ │ ┌─ Workflow Execution ───────────────────────────────────┐ │ │ │ │ │ Engine → DAG Executor → Node Executor │ │ │ │ │ │ ├─ TaskEnvelope (节点→任务封装) │ │ │ │ │ │ ├─ WorkerPool (资源感知执行池) │ │ │ │ │ │ ├─ CheckpointManager (断点续传) │ │ │ │ │ │ └─ TaskInstantiator / StateMachine / ReadyQueue │ │ │ │ │ └────────────────────────────────────────────────────────┘ │ │ │ │ │ │ │ │ ┌─ Resource Scheduling ──────────────────────────────────┐ │ │ │ │ │ ResourceCollector (5维探针) │ │ │ │ │ │ → SchedulingDecision (执行/排队/暂停/降级/拒绝) │ │ │ │ │ │ ResourceGuard (速率限制) · ResourceLedger (slot跟踪) │ │ │ │ │ │ ResourceAwareScheduler (公平性+评分) │ │ │ │ │ └────────────────────────────────────────────────────────┘ │ │ │ │ │ │ │ │ ┌─ Validation ──────────────────────────────────────────┐ │ │ │ │ │ SPJ Validator: Pre-execution Hook + Post-execution SPJ │ │ │ │ │ │ NodeRepairHook: LLM 旁路诊断 → fix/escalate/giveup │ │ │ │ │ │ NodeDiagnosis · NodeValidator · NodeGuidance │ │ │ │ │ └────────────────────────────────────────────────────────┘ │ │ │ │ │ │ │ │ ┌─ Retrieval ───────────────────────────────────────────┐ │ │ │ │ │ Permanent KG (_system) · KG Extractor · KG Optimizer │ │ │ │ │ │ Context Budget · Context Manager · Content Slicer │ │ │ │ │ │ Planning Memory · Chat Session Store │ │ │ │ │ └────────────────────────────────────────────────────────┘ │ │ │ │ │ │ │ │ ┌─ Compiler ────────────────────────────────────────────┐ │ │ │ │ │ YAML Compiler → WorkflowDefinition │ │ │ │ │ │ Skill Compiler → 展开 skill 节点为基本节点 │ │ │ │ │ │ Skill Store · Skill Atomizer │ │ │ │ │ └────────────────────────────────────────────────────────┘ │ │ │ └─────────────────────────────────────────────────────────────┘ │ │ │ │ │ ┌─────────┴─────────┐ │ │ │ PostgreSQL │ Redis │ │ │ └──────────────┴─────────┘ │ └──────────────────────────────────────────────────────────────────┘ │ HTTP Proxy ┌─────────────────────────────┴────────────────────────────────────┐ │ RAG Service (Python / FastAPI / LightRAG) │ │ │ │ ┌─────────────────────────────────────────────────────────────┐ │ │ │ LightRAG Engine (GraphRAG) │ │ │ │ ├─ Entity Extraction (LLM) │ │ │ │ ├─ Relation Building │ │ │ │ ├─ PGVector Storage │ │ │ │ └─ Graph Search (mix mode, no LLM) │ │ │ │ │ │ │ │ KB Registry (JSON file persistence) │ │ │ │ Custom KG Injection (ainsert_custom_kg) │ │ │ │ File Upload → Parse → Chunk → Embed → Store │ │ │ └─────────────────────────────────────────────────────────────┘ │ └──────────────────────────────────────────────────────────────────┘
将对话中的信息结构化为知识图谱,实现跨 Session 记忆共享,避免重复 Token 消耗。
┌─────────────────────────────────────────────────────────┐ │ Permanent KG (_system) │ │ │ │ ┌──────────┐ ┌──────────┐ ┌──────────────────────┐ │ │ │ Entities │ │Relations │ │ Chunks (conversation) │ │ │ │ │ │ │ │ │ │ │ │ 用户偏好 │ │ 偏好→项目 │ │ "用户喜欢简洁代码风格" │ │ │ │ 项目技术栈│ │ 技术栈→认证│ │ "项目用 TS+React" │ │ │ │ 认证系统 │ │ │ │ "从JWT迁移到OAuth2" │ │ │ └──────────┘ └──────────┘ └──────────────────────┘ │ │ │ │ 特性: 自动创建 · 不可删除(403) · 跨Session共享 │ └─────────────────────────────────────────────────────────┘
每次发送消息前(在 agents.ts 中): 1. estimateTokens(messages) → 是否 > 80% 预算? 2. 否 → 直接发送 3. 是 → ┌─────────────────────────────────────────────────────┐ │ a. 保留最后 8 轮原文 (recentMsgs) │ │ b. 旧消息 → kg-extractor.extractKnowledge() │ │ └─ 调用 LLM 提取: │ │ { entities: [{name, type, description}], │ │ relationships: [{src, tgt, keywords}], │ │ chunks: [{content, source_id}] } │ │ c. → permanent-kg.injectKnowledge() │ │ └─ POST /rag/custom-kg/_system │ │ └─ LightRAG: ainsert_custom_kg() │ │ d. → permanent-kg.retrieveFromKG(currentQuery) │ │ └─ POST /rag/graph-search/_system │ │ └─ mode="mix", only_need_context=True │ │ └─ 返回: chunks[] + entities[] + relationships[]│ │ e. → permanent-kg.formatKGContext(retrieved, budget) │ │ └─ 按 60/25/15 分配 chunks/entities/relations │ │ └─ 超 budget 自动截断 │ │ f. 组装: [system] + [KG记忆] + [recentMsgs] │ │ g. 校验: 若仍超预算 → 从 recent 尾部裁剪 │ └─────────────────────────────────────────────────────┘
LLM 回复后 (异步,不阻塞): evaluateUsage(replyText, kgContext) → 检测 LLM 回复中是否引用了 KG 的实体/关系/片段 → 基于关键字重叠(支持 CJK bigram 分词) → 返回 { usedEntityNames[], usedChunkIds[], relevanceScore } reinforce(usedEntities, replyText, query) → 更新实体 description,嵌入元信息: "[weight:0.85|used:3|query:用户偏好] 上次回复摘要: ..." → custom-kg 注入更新后的实体 maybePrune() (每 20 轮) → 标记从未使用的实体为 archived (weight→0.05) → 清理无用的记忆碎片
core/permanent-kg.ts
ensureSystemKG()
injectKnowledge()
retrieveFromKG()
formatKGContext(maxTokens)
core/kg-extractor.ts
extractKnowledge(messages) → {entities,relations,chunks}
core/kg-optimizer.ts
evaluateUsage()
reinforce()
maybePrune()
core/context-budget.ts
estimateTokens()
getBudget()
needsCompression()
core/context-manager.ts
sliceAndStore()
assemble()
core/content-slicer.ts
sliceContent()
core/planning-memory.ts
searchKnowledgeBase()
retrieveSimilarWorkflows()
rag_service/server.py
/rag/health
/rag/upload
/rag/query
/rag/custom-kg/{ws}
/rag/graph-search/{ws}
/rag/knowledge-bases
/rag/knowledge-bases/{ws}
/rag/knowledge-bases/{ws}/insert
/rag/knowledge-bases/{ws}/query
检索模式说明:
naive
vector
graph
hybrid
mix
graph-search
用户描述任务 (自然语言) │ ▼ LLM 分析需求 (agents.ts → compile_and_run tool) │ ▼ 生成 YAML 工作流定义 │ ▼ YAML Compiler (yaml-compiler.ts) ├─ YAML.parse → WorkflowDefinition ├─ validate nodes: 类型检查、必填字段 ├─ validate edges: 引用完整性、循环检测 └─ 返回 CompileResult { valid, definition, errors[], warnings[] } │ ▼ (如需展开 skill 节点) Skill Compiler (skill-compiler.ts) ├─ 加载 Skill 定义 (skill-store.ts) ├─ 命名空间隔离 (__{skillId}__) ├─ 展开子图节点 + 重写边 └─ 递归展开 (最大深度 5) │ ▼ POST /api/workflows → 保存 POST /api/executions → 启动执行 │ ▼ Engine.startWorkflow() → DAG Executor → Node Executor
core/yaml-compiler.ts
支持的节点类型(注册在 node-registry.ts):start, end, agent, llm, tool, code, condition, loop, parallel, subworkflow, wait, skill
node-registry.ts
start
end
agent
llm
tool
code
condition
loop
parallel
subworkflow
wait
skill
编译错误分两类:
// 编译结果 interface CompileResult { valid: boolean; definition?: WorkflowDefinition; errors: CompileIssue[]; // 阻塞性错误 warnings: CompileIssue[]; // 非阻塞性警告 }
Skill = 预编译的、带参数的可复用工作流。存储在 data/skills/{id}/skill.yml。
data/skills/{id}/skill.yml
外部 MD Skill (Claude Code / OpenClaw) → LLM 分析 → 生成 OCWE YAML → saveSkill() → 存入 data/skills/{id}/skill.yml → 用时: WorkflowDefinition 中引用 type:skill 节点 → SkillCompiler.expandSkillNodes() → 展开为基本节点
core/skill-store.ts
core/skill-compiler.ts
type: skill
core/skill-atomizer.ts
编译流程:
WorkflowDefinition (含 skill 节点) → skill-compiler.expandSkillNodes() → 1. 加载 SkillDefinition (skill-store.getSkill) → 2. 命名空间 node IDs: __{skillNodeId}__ → 3. 重写内部 edges → 4. 应用 inputMappings → 5. 重连外部 edges (parent in → entry nodes, exit nodes → parent out) → 6. 递归检查 (max depth 5, 防循环) → 返回纯基本节点的 WorkflowDefinition
旧设计(已移除):
escalate_to_workflow → planning 模式 → assess_task_clarity (Q&A 3-5轮) → finalize_plan (创建+执行) → 进入 workflowMode (会话锁定)
问题:不能退出、上下文丢失、模式切换僵硬。
新设计:
LLM 直接分析用户需求 → 如有需要,自然对话中收集信息 → 自主决定何时生成 YAML → 调用 compile_and_run(yaml, name?) → 编译 → 创建 Workflow → 异步执行 → SSE 实时推送进度 → 聊天不中断,上下文保留
工具定义(tool-calling.ts):
tool-calling.ts
{ "name": "compile_and_run", "parameters": { "yaml": "YAML 工作流定义 (nodes + edges)", "name": "可选工作流名称" } }
core/node-registry.ts
定义所有可用节点类型、配置 Schema、分类(logic/ai/control/data)。
NodeRegistry.register('llm', { category: 'ai', configSchema: { promptTemplate: { type: 'string', required: true }, systemPrompt: { type: 'string' }, model: { type: 'string' }, outputSchema: { type: 'object' }, } });
core/engine.ts
完整的生命周期管理:
startWorkflow(workflowId, {input}) ├─ validateWorkflowDefinition() ├─ createExecution → 初始化 Context → Redis 存储 ├─ CheckpointManager.startPeriodicCheckpoint() ├─ DAGExecutor.execute() │ ├─ buildGraph() → topologicalSort() → levels │ ├─ per level: SchedulingDecision → execute/queue/pause │ ├─ per node: SPJ pre-check → execute → SPJ post-check │ └─ on failure: NodeRepairHook.diagnose() → retry/escalate ├─ 事件广播: WebSocket + SSE └─ CheckpointManager.stopPeriodicCheckpoint()
支持:start / pause / resume / stop。Paused 状态持久化到 Redis + PostgreSQL。
pause
resume
stop
core/dag.ts
任务级并行调度器。核心流程:
1. buildGraph(nodes, edges) → DAGGraph ├─ 计算每个节点的 dependencies[] 和 dependents[] └─ 找到 startNodes (无依赖) 和 endNodes (无后继) 2. topologicalSort(graph) → levels[][] └─ Kahn 算法,检测循环依赖 3. per level: ├─ SchedulingDecision.decide(node, priority): │ ├─ check CPU/Memory/Disk/LLM │ ├─ queue? → 延迟重试 (5-30s) │ ├─ pause_flow? → 暂停 + checkpoint │ ├─ degrade? → 切换到备用模型 │ └─ execute_now → 执行 ├─ Promise.all(level) 并行执行同层节点 └─ 结果存入 context.variables
关键组件:ConcurrencyLimiter、DependencyResolver、TaskInstantiator、ReadyQueue、ResourceAwareScheduler
ConcurrencyLimiter
DependencyResolver
TaskInstantiator
ReadyQueue
ResourceAwareScheduler
core/executor.ts
按类型分发,支持模板变量 {{nodeId.output}} 和 {{nodeId.output.field}}。
{{nodeId.output}}
{{nodeId.output.field}}
executeCommand(cmd)
python3
node
callLlmWithTools()
engine.startWorkflow()
setTimeout
上游上下文组装 (buildUpstreamContext):
buildUpstreamContext
for each upstreamNode.output: if _sliced → 使用 KGGraphRAG 摘要 elif estimateTokens > 2000 → retrieveFromKG(摘要) else → text.slice(0, 4000)
┌─────────────────────────────────────────────────────────────┐ │ ResourceProbeManager │ │ CpuMemoryProbe │ LLMProviderProbe │ DiskProbe │ NetworkProbe│ │ 每 5s 采集一次 │ └────────────────────────┬────────────────────────────────────┘ ↓ ┌─────────────────────────────────────────────────────────────┐ │ ResourceCollector │ │ 聚合 5 维快照 → ResourceSnapshot │ │ { cpu, memory, llm: {providers}, disk, network } │ └────────────────────────┬────────────────────────────────────┘ ↓ ┌─────────────────────────────────────────────────────────────┐ │ SchedulingDecision │ │ decide(node, priority) → ScheduleAction │ │ │ │ ┌──────────────┬──────────────────────────────────────┐ │ │ │ 条件 │ 动作 │ │ │ ├──────────────┼──────────────────────────────────────┤ │ │ │ 磁盘 < 1GB │ reject (拒绝新执行) │ │ │ │ 内存 > 85% │ queue 30s / 优先级≥8 放行 │ │ │ │ CPU > 80% │ 非LLM节点 queue 5-15s │ │ │ │ LLM全部耗尽 │ pause_flow │ │ │ │ LLM < 20% │ queue 10s / 优先级≥7 放行 / degrade │ │ │ │ LLM 充足 │ execute_now / auto-assign provider │ │ │ └──────────────┴──────────────────────────────────────┘ │ │ │ │ selectFromPoolSync(usage, poolConfig) │ │ → 按序查找第一个 enabled + used < max 的 Provider │ │ → 全部满 → 返回 undefined (排队) │ │ → 高优先级 → 自动 failover 到下一个可用 Provider │ └─────────────────────────────────────────────────────────────┘ ↓ ┌─────────────────────────────────────────────────────────────┐ │ WorkerPool │ │ poll() → 取任务 → SchedulingDecision → acquire → execute │ │ → SPJ validate → checkpoint → release │ └─────────────────────────────────────────────────────────────┘
LLM 池配置(通过 Settings 前端管理,持久化在 engine-settings.json):
engine-settings.json
{ "llmPool": { "providers": [ { "id": "deepseek", "maxConcurrency": 5, "enabled": true }, { "id": "anthropic", "maxConcurrency": 3, "enabled": true }, { "id": "zhipu", "maxConcurrency": 10, "enabled": true } ] } }
任务优先级: | 级别 | 范围 | 行为 | |——|——|——| | Critical | 8-10 | 资源紧张也执行,自动 failover | | High | 7 | 优先分配,短队列等待 | | Normal | 5-6 | 正常排队 | | Low | 0-4 | 资源不足时优先让位 |
core/checkpoint.ts
Checkpoint 创建 (每次节点完成 + 定期 60s): { context: { variables: { nodeA.output: ..., nodeB.output: ... } _completedNodes: ["nodeA", "nodeB"], _resourceSnapshot: { cpu: 45%, mem: 60%, llm: "normal" } } } 恢复流程: 1. getLatestCheckpoint(executionId) → 加载快照 2. workflowState.setContext() → 恢复 Redis 3. 跳过 _completedNodes 4. SchedulingDecision.canResume() → 检查资源 5. 从中断点继续执行
core/scheduler.ts
基于 Cron 的定时任务系统:
SchedulerService
SchedulerRecovery
SchedulerOutboxRelay
core/resource-collector.ts
core/scheduling-decision.ts
core/worker-pool.ts
core/task-envelope.ts
core/task-instantiator.ts
core/task-state-machine.ts
core/dependency-resolver.ts
core/ready-queue.ts
core/concurrency-limiter.ts
core/resource-guard.ts
core/resource-ledger.ts
core/resource-probes.ts
core/resource-probe-manager.ts
core/resource-lock-manager.ts
core/resource-aware-scheduler.ts
┌─────────────────┐ │ Node Ready │ └────────┬────────┘ │ ┌────────▼────────┐ │ Pre-Hook 检查 │ ← 纯逻辑,0 Token │ inputSchema │ │ 模板变量 │ └────────┬────────┘ │ passed ┌────────▼────────┐ │ 执行节点 │ └────────┬────────┘ │ ┌────────▼────────┐ │ Post-SPJ 检查 │ ← 纯逻辑,0 Token │ outputSchema │ │ 类型检查 │ └────────┬────────┘ │ ┌──────────────┼──────────────┐ │ passed │ failed │ ▼ ▼ │ ✅ 完成 ┌──────────────┐ │ │ NodeRepairHook│ ← LLM 旁路 │ diagnose() │ 仅失败时调用 └──────┬───────┘ │ ┌──────────────┼──────────────┐ │ fix │ escalate │ giveup ▼ ▼ ▼ 修正配置重试 通知用户 标记失败
spj-validator.ts
validatePreExecution
validatePreExecution(node, context) → { passed, violations[] } 检查项: 1. inputSchema: 必需字段是否存在、类型是否正确 2. promptTemplate: {{变量}} 是否在 context 中有值 3. tool command: 模板语法是否正确 (未闭合的 {{)
节省 Token:如果上游输出不满足当前节点的输入要求,直接报错不调 LLM。
validateOutput
validateOutput(node, output) → { passed, violations[] } 检查项: 1. outputSchema 声明的字段是否存在 2. 字段类型是否匹配 (string/number/boolean/object/array) 3. required 字段是否非空 4. LLM 返回字符串但 schema 期望 object → 尝试 JSON.parse 5. 可选: LLM 语义检查 (spjPrompt)
core/node-repair.ts
仅在 SPJ 失败或节点异常时调用(旁路 LLM,不计入正常 Token 预算):
NodeRepairHook.diagnose({node, context, error, output, spjVerdict}) → 构建诊断 Prompt (node-diagnosis.ts) → 调用修复 LLM (低 temperature, 旁路) → parseDiagnosisResponse() → 返回 RepairAction: fix: { action:'fix', configOverrides, explanation, confidence } → 应用修正 → 重试执行 (最多1次) escalate: { action:'escalate', reason } → 数据/输入问题 → 通知用户 giveup: { action:'giveup', explanation } → 无法判断 → 标记失败
置信度阈值: CONFIDENCE_THRESHOLD = 0.5
CONFIDENCE_THRESHOLD = 0.5
core/spj-validator.ts
core/node-diagnosis.ts
core/node-validator.ts
core/node-guidance.ts
不区分”聊天模式”和”工作流模式”。 会话永远是一个聊天,工作流是聊天中可以调用的工具。
旧模型(已移除): Chat → escalate_to_workflow → Planning → workflowMode (锁定) 问题: 不能退出、上下文丢失 新模型: Chat (始终活跃) ├─ 日常对话 (LLM 直接回复) ├─ compile_and_run → 后台执行 (不阻塞聊天) ├─ 查看结果 (执行完成后) └─ 继续对话 (上下文保留)
activeExecutionId
draftYaml
referencedWorkflowIds
去掉了 workflowMode、planningPhase、workflowId、workflowExecutionId。
planningPhase
workflowId
workflowExecutionId
compile_and_run → SSE 推送: execution.started → 🚀 "工作流「XXX」开始执行 (5 节点) [a1b2c3d4]" node:complete → ✅ "节点完成: fetch-data (1.2s)" node:error → ❌ "节点错误: analyze — timeout" node.repair_diagnosing → 🔧 "正在诊断: analyze" workflow:complete → 🎉 "工作流执行完成!" workflow:error → 💥 "工作流执行失败: ..." Dashboard 通过 WebSocket 同步接收: workflow:started → 刷新执行列表 + 队列状态 node:complete → 刷新执行列表 workflow:complete → 刷新执行列表 + 队列状态
KB 管理:创建/编辑/删除、文件上传、文本插入、检索测试
/api/chat/sessions
/api/chat/sessions/:id/messages
/api/chat/upload
/api/workflows
/api/workflows/:id
/api/executions
{workflowId}
/api/executions/:id/state
/api/executions/:id/pause|resume|stop
/api/settings
/api/monitor/resources
/api/monitor/queue-status
/api/rag/knowledge-bases
/api/rag/knowledge-bases/:ws
/api/rag/knowledge-bases/:ws/insert
/api/rag/knowledge-bases/:ws/query
/api/rag/custom-kg/:ws
/api/rag/graph-search/:ws
backend/src/config/engine-settings.json(可通过 Settings API 热更新):
backend/src/config/engine-settings.json
{ "scheduling": { "maxConcurrentNodes": 3, "thresholds": { "cpuHighPercent": 80, "memoryHighPercent": 85, "diskLowMB": 1024, "llmTightPercent": 20 } }, "llmPool": { "providers": [ { "id": "deepseek", "maxConcurrency": 5, "enabled": true }, { "id": "zhipu", "maxConcurrency": 10, "enabled": true } ] }, "context": { "tokenBudgetRatio": 0.75, "compressThresholdRatio": 0.80, "keepRecentTurns": 8, "kgRetrievalBudget": 3000, "checkpointIntervalSec": 60 } }
setup.sh
deploy.sh
FlowForge 的 RAG 系统包含两项已投稿 SIGKDD 2027 的学术成果:
多 RAG 联合检索的多维度评价体系 — 形式化定义了”多个独立 RAG 协同检索同一目标”问题,构建了配套 benchmark,作为 FlowForge 检索效果的验证指标。
BP-free 图 RAG 检索方法 — 在 HotpotQA、2WikiMultihopQA 等国际公认的多跳问答基准数据集上,主要指标大幅超过 LightRAG、PathRAG 等 SOTA 方法。非工程封装,而是检索算法本身的创新。
MIT License
版权所有:中国计算机学会技术支持:开源发展技术委员会 京ICP备13000930号-9 京公网安备 11010802047560号
OCWE / FlowForge — 面向多智能体的操作系统级执行时
可视化 AI 工作流引擎。基于 DAG 编排,集成 LightRAG 知识图谱、资源感知动态调度、对话式工作流生成、SPJ 自动验证修复。原生适配 openEuler / 鲲鹏 ARM64。
目录
设计哲学与核心创新
核心理念:Token 效率最大化
关键原则:
五大技术创新
_systemKG→检索替代workflowMode,聊天永续,工作流异步执行快速开始
访问
http://localhost:3000。首次使用在 Settings → Provider 配置 API Key。环境:Node.js 22+ / Python 3.10+ / Redis 7+ / PostgreSQL 15+(RAG 向量存储,可选)
系统架构全景
一、检索体系 — KG / RAG / Memory
1.1 设计目标
将对话中的信息结构化为知识图谱,实现跨 Session 记忆共享,避免重复 Token 消耗。
1.2 永久知识图谱
_system1.3 上下文压缩流程
1.4 KG 动态优化
1.5 模块清单
core/permanent-kg.ts_systemKG 生命周期ensureSystemKG(),injectKnowledge(),retrieveFromKG(),formatKGContext(maxTokens)core/kg-extractor.tsextractKnowledge(messages) → {entities,relations,chunks}core/kg-optimizer.tsevaluateUsage(),reinforce(),maybePrune()core/context-budget.tsestimateTokens(),getBudget(),needsCompression()core/context-manager.tssliceAndStore(),assemble()core/content-slicer.tssliceContent()支持 JSON/Markdown/Code/Textcore/planning-memory.tssearchKnowledgeBase(),retrieveSimilarWorkflows()1.6 RAG 服务端点 (
rag_service/server.py)/rag/health/rag/upload/rag/query/rag/custom-kg/{ws}/rag/graph-search/{ws}/rag/knowledge-bases/rag/knowledge-bases/{ws}/rag/knowledge-bases/{ws}/insert/rag/knowledge-bases/{ws}/query检索模式说明:
naive: 纯向量搜索,最快vector: 语义向量搜索graph: 知识图谱遍历hybrid: 向量 + 图谱混合mix: 返回原始 chunks + entities + relations(仅graph-search)二、任务设计 — YAML 编译 / Skill / 工作流规划
2.1 从需求到执行:完整数据流
2.2 YAML 编译器 (
core/yaml-compiler.ts)支持的节点类型(注册在
node-registry.ts):start,end,agent,llm,tool,code,condition,loop,parallel,subworkflow,wait,skill编译错误分两类:
2.3 Skill 系统
Skill = 预编译的、带参数的可复用工作流。存储在
data/skills/{id}/skill.yml。core/skill-store.tscore/skill-compiler.tstype: skill→ 内联子图,递归 maxDepth=5core/skill-atomizer.ts编译流程:
2.4 对话式工作流生成(替代旧的 Planning 模式)
旧设计(已移除):
问题:不能退出、上下文丢失、模式切换僵硬。
新设计:
工具定义(
tool-calling.ts):2.5 节点注册表 (
core/node-registry.ts)定义所有可用节点类型、配置 Schema、分类(logic/ai/control/data)。
三、工作流执行 — Engine / DAG / 节点 / 调度
3.1 执行引擎 (
core/engine.ts)完整的生命周期管理:
支持:
start/pause/resume/stop。Paused 状态持久化到 Redis + PostgreSQL。3.2 DAG 执行器 (
core/dag.ts)任务级并行调度器。核心流程:
关键组件:
ConcurrencyLimiter、DependencyResolver、TaskInstantiator、ReadyQueue、ResourceAwareScheduler3.3 节点执行器 (
core/executor.ts)按类型分发,支持模板变量
{{nodeId.output}}和{{nodeId.output.field}}。startendtoolexecuteCommand(cmd)→ 沙箱校验 → spawncodepython3/node→ 收集 stdoutllm/agentcallLlmWithTools()conditionparallelloopsubworkflowengine.startWorkflow()waitsetTimeout上游上下文组装 (
buildUpstreamContext):3.4 资源感知调度系统
LLM 池配置(通过 Settings 前端管理,持久化在
engine-settings.json):任务优先级: | 级别 | 范围 | 行为 | |——|——|——| | Critical | 8-10 | 资源紧张也执行,自动 failover | | High | 7 | 优先分配,短队列等待 | | Normal | 5-6 | 正常排队 | | Low | 0-4 | 资源不足时优先让位 |
3.5 断点续传 (
core/checkpoint.ts)3.6 调度器 (
core/scheduler.ts)基于 Cron 的定时任务系统:
SchedulerService: 管理定时任务的创建/启停/触发SchedulerRecovery: 重启后恢复未完成的调度任务SchedulerOutboxRelay: 可靠事件投递3.7 模块清单
core/engine.tscore/dag.tscore/executor.tscore/resource-collector.tscore/scheduling-decision.tscore/worker-pool.tscore/task-envelope.tscore/task-instantiator.tscore/task-state-machine.tscore/checkpoint.tscore/dependency-resolver.tscore/ready-queue.tscore/concurrency-limiter.tscore/resource-guard.tscore/resource-ledger.tscore/resource-probes.tscore/resource-probe-manager.tscore/resource-lock-manager.tscore/resource-aware-scheduler.tscore/scheduler.ts四、验证与修复 — SPJ / Hook / Diagnosis
4.1 核心思想:正常路径零 LLM 开销
4.2 Pre-execution Hook (
spj-validator.ts→validatePreExecution)节省 Token:如果上游输出不满足当前节点的输入要求,直接报错不调 LLM。
4.3 Post-execution SPJ (
spj-validator.ts→validateOutput)4.4 节点修复 (
core/node-repair.ts)仅在 SPJ 失败或节点异常时调用(旁路 LLM,不计入正常 Token 预算):
置信度阈值:
CONFIDENCE_THRESHOLD = 0.54.5 模块清单
core/spj-validator.tscore/node-repair.tscore/node-diagnosis.tscore/node-validator.tscore/node-guidance.ts五、会话模型 — 统一 Chat + Workflow
5.1 设计原则
不区分”聊天模式”和”工作流模式”。 会话永远是一个聊天,工作流是聊天中可以调用的工具。
5.2 会话状态
activeExecutionIddraftYamlreferencedWorkflowIds去掉了
workflowMode、planningPhase、workflowId、workflowExecutionId。5.3 实时反馈
操作指南
聊天
工作流创建
Dashboard
Settings (5 Tabs)
Memory
KB 管理:创建/编辑/删除、文件上传、文本插入、检索测试
API 参考
聊天
/api/chat/sessions/api/chat/sessions/:id/messages/api/chat/upload工作流
/api/workflows/api/workflows/:id/api/executions{workflowId}/api/executions/:id/state/api/executions/:id/pause|resume|stop设置
/api/settings监控
/api/monitor/resources/api/monitor/queue-status知识库
/api/rag/knowledge-bases/api/rag/knowledge-bases/:ws_system403)/api/rag/knowledge-bases/:ws/insert/api/rag/knowledge-bases/:ws/query/api/rag/custom-kg/:ws/api/rag/graph-search/:ws配置参考
backend/src/config/engine-settings.json(可通过 Settings API 热更新):节点类型
start/endtoolcodellm/agentconditionparallelloopsubworkflowwaitskill技术栈
平台兼容性
setup.sh+deploy.sh)GraphRAG 学术创新
FlowForge 的 RAG 系统包含两项已投稿 SIGKDD 2027 的学术成果:
多 RAG 联合检索的多维度评价体系 — 形式化定义了”多个独立 RAG 协同检索同一目标”问题,构建了配套 benchmark,作为 FlowForge 检索效果的验证指标。
BP-free 图 RAG 检索方法 — 在 HotpotQA、2WikiMultihopQA 等国际公认的多跳问答基准数据集上,主要指标大幅超过 LightRAG、PathRAG 等 SOTA 方法。非工程封装,而是检索算法本身的创新。
对标分析
许可证
MIT License