OpenHuman Flows 运行架构:tinyflows 如何把保存的工作流编译为 tinyagents 状态图并执行
本文基于 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.rs、compiler.rs、engine.rs、expr.rs、caps/、nodes/ 等模块即引擎本体的代码布局。tinyflows 已发布到 crates.io,并且刻意做到"无持久化、无宿主绑定"(persistence-free、vendor-free),OpenHuman 是第一个下游宿主,一切真实依赖都通过接缝注入。
运行流水线:validate → compile → run → outcome
一个 flow 在 wire 上就是一张 WorkflowGraph——由带类型的 Node 与 Edge 组成的有向图,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
- validate(结构检查):对原始图做结构校验——节点 id 唯一、所有被引用的节点必须存在、以及恰好一个 trigger 节点(0 个报
MissingTrigger,多于 1 个报MultipleTriggers)。这正是后文"单一触发器不变式"的落点。 - compile(降级):先跑 validate,再把通过校验的图降级成一张全新的 tinyagents 状态图。每个 tinyflows 节点变成一个 tinyagents 图节点,每条 tinyflows 边变成图边,条件路由、并行路由与 merge barrier(汇合屏障)都在 tinyagents 图层表达。图在每次运行时重建,保证编译与宿主状态完全解耦。
- run(驱动执行):宿主侧的
flows_runRPC(入口在 src/openhuman/flows/ops.rs)调用engine::run_with_checkpointer_journaled_observed,驱动编译产物跑完,把每个节点的输出折叠进 run state,并在审批门禁(approval gate)处暂停。 - 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::SqliteCheckpointer 调 engine::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-out(split_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 payload(run.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::agent是Option,没有 agent 注册表的宿主可以留None,agent节点回退到裸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→ 构建一个完整 harnessAgent(Agent::from_config_for_agent),把节点的请求通过run_single跑完。这与 builder / scout 以及 cron / subconscious 任务用的是同一个入口,内部驱动run_turn_via_tinyagents_shared(src/openhuman/agent/tinyagents/mod.rs),即 tinyagentsAgentHarness的 tool-call 循环。该定义的ToolScope、sandbox_mode、max_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_resume(src/openhuman/flows/ops.rs)通过 with_origin 给整个引擎 future 套上一个 TrustedAutomation { source: Workflow { require_approval } } origin。在每个执行型节点分发前,接缝用该节点的 CommandClass(http_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 的重复失败中间件识别为"永久性拒绝、不要重试"),节点根本不分发; - 决策为
Prompt(supervised层级)→ 由 tier.rs 的gate_call_for_tier升级为强制ApprovalGate往返。源码注释(Codex P1 修复说明)解释了为什么必须显式升级:ApprovalGate::intercept_audited对Workflow { 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_tool,caps/ops.rs):slug 只有在解析为"已知的、被策展的、在范围内、已连接"的 toolkit 动作时才放行——比通用 agent tool-call 路径更严格,因为 flow 的 slug 是作者手敲的自由字符串,而非后端实时发现的产物。配套的 classify_composio_action_for_tier(caps/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 定义自身的 ToolScope、sandbox_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 域而非引擎处理。FlowTriggerSubscriber(src/openhuman/flows/bus.rs)是 trigger → run 的桥:它订阅保存流单个 trigger 节点可绑定的归一化领域事件,每次匹配都调 flows::ops::flows_run:
DomainEvent::FlowScheduleTick——一个flow类型 cron 任务触发;加载被点名的 flow,复查其scheduletrigger 是否仍存活,然后以空 trigger payload 分发。DomainEvent::ComposioTriggerReceived——每个启用了绑定到该toolkit/trigger_slug的app_eventtrigger 的 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_agent → run_single + 作用域化 origin,与 flow 自己的 agent 节点同一个模式:
| Agent | 注册 id | 入口 | 工具带(tool belt) |
|---|---|---|---|
| Builder(copilot) | workflow_builder |
flows_build(ops.rs) |
propose_workflow / revise_workflow(仅校验)、dry_run_workflow(编译 + 对 mock 跑)、save_workflow、run_workflow、catalog/连接读取(builder_tools.rs) |
| Scout(发现) | flow_discovery |
flows_discover(ops.rs) |
suggest_workflows(discovery_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.rs、compiler.rs、engine.rs、expr.rs、caps/、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/tinyflows 与 vendor/tinyagents 为 git submodule,需初始化子模块后才能阅读引擎源码;文档中引用的引擎内部文件(如 vendor/tinyflows/src/caps/mod.rs)以子模块实际检出版本为准。
atomcodeClaude Code 的开源替代方案。连接任意大模型,编辑代码,运行命令,自动验证 — 全自动执行。用 Rust 构建,极致性能。 | An open-source alternative to Claude Code. Connect any LLM, edit code, run commands, and verify changes — autonomously. Built in Rust for speed. Get StartedRust0631
MiniCPM5-2BMiniCPM5-2B 是一款面向端侧、本地部署和资源受限场景的 2B 稠密 Transformer,能够达到同尺寸开源模型 SOTA 水平。Markdown00
video-shotcraftAI宣传片skill,使用 Remotion 制作电影级产品视频:提供106 张镜头配方卡和可复用的视频魔板。适用于 Claude Code 与 Codex以及所有其他智能体Markdown00
HivisionIDPhotos⚡️HivisionIDPhotos: a lightweight and efficient AI ID photos tools. 一个轻量级的AI证件照制作算法。Python09
DragonOSDragonOS is an operating system developed from scratch using Rust, with Linux compatibility. It is designed for **Serverless** scenarios. 使用Rust从0自研内核,具有Linux兼容性的操作系统,面向云计算Serverless场景而设计。Rust00
Spark-X2.5-1.7BSpark-X2.5-1.7B 旨在让强大的 AI 更加实用、高效且易于获取。这些模型在广泛的日常任务中表现出色,涵盖对话、写作、翻译、推理、编程、工具调用和智能体工作流,并在同等规模的开源模型中取得领先结果。Spark-X2.5 将面向效率的架构与最高 1M tokens 的原生上下文窗口相结合,并支持 200 多种语言。Python00