首页
/ OpenHuman Flows 运行架构:tinyflows 如何把保存的工作流编译为 tinyagents 状态图并执行

OpenHuman Flows 运行架构:tinyflows 如何把保存的工作流编译为 tinyagents 状态图并执行

2026-09-09 15:25:34作者:田桥桑Industrious

本文基于 OpenHuman 仓库的架构文档 flows-on-tinyagents.md 展开。flows 域是 OpenHuman 中驱动"保存式自动化"(saved automations)的核心子系统:用户画布上搭出来的工作流,或 copilot 替用户构建的工作流,最终都降级(lower)到 tinyagents 状态图引擎上执行。读完本文,你将掌握一个 flow 从 SQLite 存储到 RunOutcome 的完整四阶段流水线、运行状态的 {run, nodes} JSON 模型与 = 表达式绑定机制、宿主注入的 capability 接缝(seam)设计,以及外层自治层级门禁与内层 Agent 工具沙箱组成的双层安全模型。

两个 crate,一套引擎

flows 域本身不包含工作流引擎。所有流程的 validate / compile / run 都由外部 crate tinyflows 完成,而它再降级到 tinyagents 状态图引擎——即 Agent Harness 中 agent 单个 turn 所用的同一个引擎。于是"一个保存的 flow 就是一张 tinyagents 图",而每个 agent 节点又是一张 tinyagents 图——图中之图(graph within a graph)。

各组件的分工如下:

Crate / 模块 角色 位置
tinyflows 宿主无关的工作流模型 + validate + compile + run。从不硬编码任何宿主;所有对外部世界的副作用都经过 capability trait 走 子模块根目录 vendor/tinyflows(git submodule,见 .gitmodules
tinyagents 两个运行时共同降级到的已发布状态图 + agent-loop harness crate 子模块根目录 vendor/tinyagents;OpenHuman 侧接缝在 src/openhuman/agent/tinyagents/
openhuman::flows 宿主:CRUD / run / resume RPC、SQLite 存储、触发器、builder / scout agent src/openhuman/flows/
openhuman::flows::tinyflows capability 接缝——在真实 OpenHuman 服务上实现 tinyflows 各 trait 的适配器 src/openhuman/flows/tinyflows/

需要注意版本与获取方式:vendor/tinyflows.gitmodules 中登记为 git submodule(upstream 为 tinyhumansai/tinyflows),当前检出中该目录为空壳,需 git submodule update --init 拉取引擎源码;文档中提到的 vendor/tinyflows/src/ 下的 model/validate.rscompiler.rsengine.rsexpr.rscaps/nodes/ 等模块即引擎本体的代码布局。tinyflows 已发布到 crates.io,并且刻意做到"无持久化、无宿主绑定"(persistence-free、vendor-free),OpenHuman 是第一个下游宿主,一切真实依赖都通过接缝注入。

运行流水线:validate → compile → run → outcome

一个 flow 在 wire 上就是一张 WorkflowGraph——由带类型的 NodeEdge 组成的有向图,JSON 序列化。运行它是固定的四阶段流水线:

flowchart LR
  subgraph host["openhuman::flows (host)"]
    store[(SQLite<br/>flow store)] --> graph[WorkflowGraph<br/>JSON]
  end
  graph --> validate["validate<br/>(结构检查)"]
  validate --> compile["compiler::compile<br/>(降级到 tinyagents)"]
  compile --> gb["tinyagents<br/>GraphBuilder"]
  gb --> run["engine::run_with_checkpointer<br/>(驱动到完成)"]
  run --> outcome["RunOutcome<br/>{ output, pending_approvals, cancelled }"]
  seam["openhuman::flows::tinyflows<br/>capability seam"] -. 宿主注入 .-> run
  1. validate(结构检查):对原始图做结构校验——节点 id 唯一、所有被引用的节点必须存在、以及恰好一个 trigger 节点(0 个报 MissingTrigger,多于 1 个报 MultipleTriggers)。这正是后文"单一触发器不变式"的落点。
  2. compile(降级):先跑 validate,再把通过校验的图降级成一张全新的 tinyagents 状态图。每个 tinyflows 节点变成一个 tinyagents 图节点,每条 tinyflows 边变成图边,条件路由、并行路由与 merge barrier(汇合屏障)都在 tinyagents 图层表达。图在每次运行时重建,保证编译与宿主状态完全解耦。
  3. run(驱动执行):宿主侧的 flows_run RPC(入口在 src/openhuman/flows/ops.rs)调用 engine::run_with_checkpointer_journaled_observed,驱动编译产物跑完,把每个节点的输出折叠进 run state,并在审批门禁(approval gate)处暂停。
  4. outcome(收尾):产出 RunOutcome { output, pending_approvals, cancelled }。宿主通过 FlowRunObserver 持久化 live steps,并把持久的图观测(durable graph observations)导出到 Langfuse——导出实现见 langfuse_export.rs

由于引擎用调用方提供的 thread_id 作为持久化状态的主键,HITL(human-in-the-loop)的持久化恢复就是拿同一个 tinyagents::graph::SqliteCheckpointerengine::resume_with_checkpointer——与 agent harness 用的是同一套 checkpointer,每个宿主只打开一次,落在 <workspace_dir>/flows/checkpoints.db。打开逻辑是 open_flow_checkpointer,与 build_capabilities 一样从 caps/mod.rs 重导出(底层持久化 schema 由 tinyflows_sqlite::checkpoint 提供,保证所有能恢复中断 run 的宿主共享同一磁盘格式,见 tinyflows/mod.rs 的模块注释)。

运行状态:一份 JSON map、一个 merge reducer、一套 {json, text, raw} 信封

整个 run 的工作记忆就是一个 serde_json::Value,布局固定为:

{
  "run": {
    "trigger": {
      "/* run 开始时种入的 trigger payload */": ""
    }
  },
  "nodes": {
    "planner": {
      "items": [ "/* … */" ]
    },
    "drafter": {
      "items": [ "/* … */" ]
    }
  }
}

随着图执行,一个 merge reducer 把每个节点的 item 输出折叠进 nodes.<id>。因为合并是"追加"而不是"后者覆盖前者",fan-in 节点能同时看到所有前驱的输出;fan-outsplit_out)节点的并行分支在屏障处汇合时互不踩踏——这与 tinyagents 图层自身 map_reduce 扇出用的是同一套 reducer 模式。

所有能力节点(agent / tool_call / http_request / code)的输出都包在一个稳定的 { json, text, raw } 信封里。下游节点永远看不到原始的 completion 形态,只能对着信封绑定:

  • =item.json.<field>——结构化对象(当生产者发出了 JSON 时);
  • =item.text——文本形态;
  • =item.raw——未加工的生产者原始输出。

= 表达式的四个作用域绑定

节点配置通过 jq/jaq = 表达式(引擎侧 vendor/tinyflows/src/expr.rs)在一个每节点 scope 上求值,scope 有四个绑定:

绑定 内容
item 第一个输入 item 的 json(直接前驱的输出)
items 所有输入 item 的 json,按边顺序排列
run run 元数据与 trigger payloadrun.trigger
nodes 所有已完成节点的输出,按 id 索引{ "<id>": { "item": …, "items": [ … ] } }

因此一个下游三跳的节点可以按 id 精确回溯某个祖先——"=nodes.planner.item.json.plan"——这正是演示工作流把 planner 的结构化计划传进 drafter prompt 的方式。"=item.text" 则是"直接前驱"的简写。缺失片段解析为 null 而非报错,非法的 jq 程序也返回 null 而不是 panic——节点接线永远不会弄崩整个 run。

capability 接缝:引擎不触碰任何真实世界

tinyflows 自身不碰任何真实依赖。LLM 调用、工具、HTTP、代码执行、持久化、子工作流查找——全部是宿主实现的 capability trait(引擎侧定义在 vendor/tinyflows/src/caps/mod.rs)。OpenHuman 侧每个 trait 配一个适配器,由 build_capabilities 在每次 run 时组装成 Capabilities 包。实现位于 caps/ops.rs

pub fn build_capabilities(config: Arc<Config>, state_namespace: impl Into<String>) -> Capabilities {
    let security = Arc::new(SecurityPolicy::from_config(
        &config.autonomy, &config.workspace_dir, &config.action_dir,
    ));
    // …
    Capabilities {
        llm: Arc::new(OpenHumanLlm { config: config.clone() }),
        tools: Arc::new(OpenHumanTools { config: config.clone(), security: security.clone() }),
        http: Arc::new(OpenHumanHttp { security, http_config, http_creds }),
        code: Arc::new(OpenHumanCode { config: config.clone(), security }),
        state: Arc::new(FlowStateStore { config: config.clone(), namespace }),
        agent: Some(Arc::new(OpenHumanAgentRunner { config })),
        shell: None,               // 尚无专用宿主适配器,能力保持不可用
        tasks: Some(Arc::new(TokioTaskRunner::new())),
        approvals: None,           // 有意走引擎自带的兼容回退,避免建第二套审批存储
        memory: Some(Arc::new(OpenHumanMemory { config, security })),
        // …
    }
}

各 trait 与适配器、被包装的宿主服务的对应关系:

tinyflows trait 支撑的节点 OpenHuman 适配器 包装的服务 源码位置
LlmProvider agent(无 agent_ref 的裸调用)、output_parser OpenHumanLlm inference provider(create_chat_provider caps/llm.rs
AgentRunner agent(带 agent_ref OpenHumanAgentRunner agent 注册表 + harness Agent caps/agent.rs
ToolInvoker tool_call OpenHumanTools Composio + 原生 oh: 工具 caps/tools/mod.rs
HttpClient http_request OpenHumanHttp HttpRequestTool(allowlist + DNS-rebind 防护) caps/http.rs
CodeRunner code OpenHumanCode 沙箱(execute_in_sandbox caps/code.rs
StateStore 可恢复/有状态 run FlowStateStore flow_state KV 表 caps/state.rs
WorkflowResolver 按 id 的 sub_workflow OpenHumanWorkflowResolver 保存流存储(load_flow_graph caps/resolver.rs

两个值得注意的设计细节(均直接来自上述源码注释):

  • AgentRunner 在 crate 中是可选的Capabilities::agentOption,没有 agent 注册表的宿主可以留 Noneagent 节点回退到裸 LlmProvider 补全。OpenHuman 永远会装配它,所以 agent 节点拿到的是真实的 agent 运行时。
  • state_namespace 决定 KV 隔离FlowStateStore 用它给键加命名空间(如 "flow:<id>"),保证两个保存流用同一个 state key 也互不读写。源码注释特别强调它不等同于 OpenHumanMemory 写入 flow 作用域 memory 的命名空间——后者由 run 的 trusted origin 独立派生,两者从不需要对齐分隔符约定。
  • 源码里 shell: None 的注释也印证了接缝纪律:在专用宿主适配器(应用宿主 autonomy/sandbox 策略)存在之前,宁可让能力不可用,也不让工作流引擎继承进程级 ambient shell 权限。

Agent 节点:图中之图

flow 的 agent 节点通过配置中一个受信任的 agent_ref 指名一个已注册的 agent kind(researcher、code_executor、某个自定义 specialist……)。OpenHumanAgentRunner::run_agent 解析该引用并按查找结果路由:

  • 存在完整的 harness AgentDefinition → 构建一个完整 harness AgentAgent::from_config_for_agent),把节点的请求通过 run_single 跑完。这与 builder / scout 以及 cron / subconscious 任务用的是同一个入口,内部驱动 run_turn_via_tinyagents_sharedsrc/openhuman/agent/tinyagents/mod.rs),即 tinyagents AgentHarness 的 tool-call 循环。该定义的 ToolScopesandbox_modemax_iterations 约束内部 turn。
  • 只有自定义 AgentRegistryEntry(没有完整定义)→ 走 persona 塑造回退:把 entry 的 system_prompt(以及节点未固定模型时的模型)前置,请求经 OpenHumanLlm::complete 单次补全,没有私有 tool loop。这保证自定义注册表 agent 不退化为回归缺陷。

由此得到嵌套结构:

flowchart TB
  subgraph flow["Flow = tinyagents 图(编译后的 WorkflowGraph)"]
    t[trigger] --> a["agent 节点<br/>agent_ref = planner"]
    a --> tc[tool_call 节点]
    tc --> tr[transform 节点]
  end
  a -.->|"run_agent → Agent::from_config_for_agent → run_single"| inner
  subgraph inner["Agent turn = 嵌套 tinyagents 图(AgentHarness)"]
    p[provider 调用] --> parse[解析 tool calls]
    parse --> tools[分发工具]
    tools --> compact[上下文守卫]
    compact --> p
    parse --> done[最终文本]
  end
  done -.->|"{ json, text, raw } 信封"| tc

两条路径都返回 JSON;agent 节点把它折叠进 { json, text, raw } 信封(output_parser 子端口会应用节点声明的 schema),所以下游节点的绑定方式无论节点跑的是完整 agent turn 还是裸补全都完全一致

安全上,agent_ref 只从受信节点配置解析,绝不从模型输出解析——被 prompt 注入的上游补全无法挑选任意 agent kind。builder 的 dry-run 也走这条路径:dry_run_workflow 把草稿编译后对着 tinyflows 的 mock capabilities(MockAgentRunner)跑一遍,带 agent_ref 的草稿在保存之前就能端到端自测。

双层安全模型:外层自治门禁 + 内层工具沙箱

一次 flow run 由两层相互独立的门禁保护:外层归 flows 运行时所有,内层归 agent harness 所有。

外层门禁:flow origin + 自治层级(autonomy tier)

flows_run / flows_resumesrc/openhuman/flows/ops.rs)通过 with_origin 给整个引擎 future 套上一个 TrustedAutomation { source: Workflow { require_approval } } origin。在每个执行型节点分发前,接缝用该节点的 CommandClasshttp_request → Network、code → Write、原生 oh: 工具 → 各自分类的类)查询用户 [autonomy] 层级的 SecurityPolicy::gate_decision。实现是 caps/tier.rs 中的 enforce_node_tier_gate

  • 决策为 Block(如 readonly 层级下的 network/code 动作)→ 直接返回 Err(EngineError::Capability) 且错误串带 POLICY_BLOCKED_MARKER 前缀(让 harness 的重复失败中间件识别为"永久性拒绝、不要重试"),节点根本不分发
  • 决策为 Promptsupervised 层级)→ 由 tier.rsgate_call_for_tier 升级为强制 ApprovalGate 往返。源码注释(Codex P1 修复说明)解释了为什么必须显式升级:ApprovalGate::intercept_auditedWorkflow { require_approval: false } 这类"预设信任根"origin 会无条件返回 Allow——若不升级,一个 supervised 用户的 http_request/code 节点会仅因保存流的 require_approval 默认为 false 就无人值守地跑起来。升级的做法是对当次调用临时套一个 require_approval: true 的 origin,强制真实停车/HITL 流程;
  • full 层级直接放行。

Composio tool_call 节点还多一道默认拒绝的curation 门禁is_curated_flow_toolcaps/ops.rs):slug 只有在解析为"已知的、被策展的、在范围内、已连接"的 toolkit 动作时才放行——比通用 agent tool-call 路径更严格,因为 flow 的 slug 是作者手敲的自由字符串,而非后端实时发现的产物。配套的 classify_composio_action_for_tiercaps/ops.rs)按效果分类:只有解析到策展目录中 ToolScope::Read 的 slug 才映射为 CommandClass::Read(唯一被所有层级直接 Allow 的类),Write/Admin/未策展/空 slug 一律映射为 CommandClass::Network,即"未人工策展为只读的动作必须照样弹审批"。源码注释强调这一分类从 tinyflows 0.3 起是承重墙而非纵深防御:integration 节点配置(含 slug)现在会先对上游/trigger 数据做 = 表达式求值再 invoke——也就是说注入数据可以影响参数,但影响不到"允许调用哪个工具"。

内层门禁:agent 定义的 ToolScope + 沙箱

agent 节点运行完整 harness Agent 时,那个 turn 运行在同一个 Workflow origin 下(引擎 future 已经处于 with_origin 内部),所以自治层级与审批门禁自动作用于内部 turn——不需要新的 origin 包装器。在此之上,agent 定义自身的 ToolScopesandbox_mode 与迭代上限约束 turn 的能力,与 chat sub-agent 完全一致。

也就是说 agent 节点被门禁了两次:flow 的外层 autonomy/origin 门禁 + agent 的内层 tool/sandbox 门禁。而 agent_ref 只认受信配置则是第三条腿:不可信的 trigger 或上游数据可以影响参数(受 curation + scope + approval 约束),但永远决定不了运行哪种 agent kind 或哪个工具身份。这些门禁行为有对应单测覆盖,例如 caps/ops_tests_part_02_tests.rs 分别验证了 ReadOnly/Supervised/Full 三种层级下 enforce_node_tier_gate 对 Network/Write 类的放行与拦截。

触发模型:图上恰一个 trigger,多触发是宿主侧的事

tinyflows 在 validate 阶段强制每张图恰好一个 trigger 节点:工作流因此只有一个入口,trigger 的 payload 即 run state 里 run.trigger 的种子。

这是模型层约束,不是产品限制。多触发当前是宿主侧职责,由 flows 域而非引擎处理。FlowTriggerSubscribersrc/openhuman/flows/bus.rs)是 trigger → run 的桥:它订阅保存流单个 trigger 节点可绑定的归一化领域事件,每次匹配都调 flows::ops::flows_run

  • DomainEvent::FlowScheduleTick——一个 flow 类型 cron 任务触发;加载被点名的 flow,复查其 schedule trigger 是否仍存活,然后以空 trigger payload 分发。
  • DomainEvent::ComposioTriggerReceived——每个启用了绑定到该 toolkit/trigger_slugapp_event trigger 的 flow 都会被分发,事件 payload 种入 run.trigger

分发按 flow_id 去重(trigger 爆发不会为同一 flow 产生重叠 run),而交互式 flows_run RPC 刻意不去重——用户要求再跑一次总是被尊重。所以"定时跑 + Gmail 消息到达时跑"表达为两个宿主订阅汇入同一张单 trigger 图。模型层的多 trigger 属于未来工作,当前接缝与存储都不依赖它。

Builder、Scout、Executor:一个共享 harness

三个触碰 flows 的 agent 全部跑在共享 agent harness 上——Agent::from_config_for_agentrun_single + 作用域化 origin,与 flow 自己的 agent 节点同一个模式:

Agent 注册 id 入口 工具带(tool belt)
Builder(copilot) workflow_builder flows_buildops.rs propose_workflow / revise_workflow(仅校验)、dry_run_workflow(编译 + 对 mock 跑)、save_workflowrun_workflow、catalog/连接读取(builder_tools.rs
Scout(发现) flow_discovery flows_discoverops.rs suggest_workflowsdiscovery_tools.rs
Executor (无——即引擎本身) flows_run / flows_resume 上文的 capability 接缝

两个 agent 都作为一等注册表 agent 住在 src/openhuman/flows/agents/ 下(每个是 agent.toml + prompt 定义,如 workflow_builder/agent.toml),跑在 reasoning 层级上,带窄而经过安全审查的工具带。builder 是全工具辅助的:它提议图(仅校验)、对 mock capabilities dry-run 自测、然后才保存并 confirm-first 地试跑真实 flow。完整的 discover → build → save → run 生命周期:

sequenceDiagram
  participant U as User
  participant S as flow_discovery (scout)
  participant B as workflow_builder (builder)
  participant St as flow store (SQLite)
  participant E as tinyflows engine

  U->>S: flows_discover
  S->>S: suggest_workflows (校验想法)
  S-->>U: FlowSuggestion[]
  U->>B: flows_build (挑选建议 / 提示词)
  B->>B: propose_workflow (校验草稿)
  B->>E: dry_run_workflow (编译 + 对 mock 跑)
  E-->>B: mock RunOutcome (自测)
  B->>St: save_workflow (已存 flow)
  B-->>U: WorkflowProposal
  U->>E: flows_run (confirm-first)
  E->>E: validate → compile → run (真实 caps)
  E-->>U: RunOutcome (live steps + Langfuse trace)

代码索引:去哪里看

路径 内容
src/openhuman/flows/ops.rs(及 ops_part_*.rs flows_run / flows_resume / flows_build / flows_discover 与 CRUD
src/openhuman/flows/bus.rs FlowTriggerSubscriber——宿主侧多触发桥
src/openhuman/flows/builder_tools.rs builder 的 propose / revise / dry_run / save / run 工具
src/openhuman/flows/agents/ workflow_builder + flow_discovery agent 定义
src/openhuman/flows/store.rs SQLite flow 存储
src/openhuman/flows/tinyflows/mod.rs 接缝模块入口,重导出 build_capabilities / open_flow_checkpointer / checkpoint_sqlite
src/openhuman/flows/tinyflows/caps/ 宿主适配器:ops.rs(组装/curation/preflight)、tier.rs(双层门禁)、agent.rs / http.rs / code.rs / llm.rs / state.rs / resolver.rs / tools/
src/openhuman/flows/tinyflows/memory_adapter.rs memory 节点的 OpenHumanMemory 适配器(复用同一 tier 门禁)
src/openhuman/flows/tinyflows/langfuse_export.rs run 结束后的持久图观测导出
src/openhuman/flows/tinyflows/observability.rs run 观测日志
vendor/tinyflows/ 引擎本体(submodule):model/validate.rscompiler.rsengine.rsexpr.rscaps/nodes/
tests/flows_lifecycle_e2e.rs flows 生命周期的 e2e 测试

相关文档

  • 本文对应的用户侧功能文档:gitbooks/features/workflows.md
  • Agent Harness(每个 agent 节点与 builder/scout 所跑的 tinyagents turn)与 Security(外层节点门禁所查询的 SecurityPolicy / autonomy tier)分别在 gitbooks/developing/architecture/ 目录下的对应文档中展开。

适用前提:本文描述的是当前仓库快照中的实现;vendor/tinyflowsvendor/tinyagents 为 git submodule,需初始化子模块后才能阅读引擎源码;文档中引用的引擎内部文件(如 vendor/tinyflows/src/caps/mod.rs)以子模块实际检出版本为准。

登录后查看全文
热门项目推荐
相关项目推荐

项目优选

收起
kernelkernel
deepin linux kernel
C
33
18
ops-transformerops-transformer
本项目是CANN提供的transformer类大模型算子库,实现网络在NPU上加速计算。
C++
1.14 K
2.76 K
pytorchpytorch
作为 Ascend for PyTorch 社区的核心组件,TorchNPU 是昇腾专为 PyTorch 打造的深度学习适配插件,使 PyTorch 框架能够直接调用昇腾 NPU,为开发者提供昇腾 AI 处理器的超强算力。
Python
860
1.35 K
docsdocs
暂无描述
Markdown
899
5.83 K
ops-nnops-nn
本项目是CANN提供的神经网络类计算算子库,实现网络在NPU上加速计算。
C++
925
1.85 K
jiuwenswarmjiuwenswarm
JiuwenSwarm 是一款基于openJiuwen开发的智能AI Agent,它能够将大语言模型的强大能力,通过你日常使用的各类通讯应用,直接延伸至你的指尖。
Python
3.84 K
1.02 K
kernelkernel
openEuler内核是openEuler操作系统的核心,既是系统性能与稳定性的基石,也是连接处理器、设备与服务的桥梁。
C
533
601
cann-learning-hubcann-learning-hub
CANN 学习中心仓,支持在线互动运行、边学边练,提供教程、示例与优化方案,一站式助力昇腾开发者快速上手。
Jupyter Notebook
1.03 K
525
ops-mathops-math
本项目是CANN提供的数学类基础计算算子库,实现网络在NPU上加速计算。
C++
1.37 K
1.46 K
AscendNPU-IRAscendNPU-IR
AscendNPU-IR是基于MLIR(Multi-Level Intermediate Representation)构建的,面向昇腾亲和算子编译时使用的中间表示,提供昇腾完备表达能力,通过编译优化提升昇腾AI处理器计算效率,支持通过生态框架使能昇腾AI处理器与深度调优
C++
548
395