数据截至 (上游 commit fc74d079a18c)
任务执行引擎:纯解释器与 durable 后端(主线)
30 秒导读: v3 的执行引擎由两层组成:一个纯解释器
interpret()负责全部控制流 (顺序/并行/分支/循环/计划/agent 循环),它不 import 任何执行引擎;一个Env 协议承接 全部副作用(调工具、调模型、跑子流程、等人)。Temporal 后端给每个 Env handler 配一个 activity 或子工作流,于是整棵 IR 的求值都落在 Temporal 历史里,崩了重放即可续跑; DBOS 后端用 dbos-transact 的 step 做同样的事。这是全库最核心的一章。
本章聚焦"引擎怎么转"。IR 与数据模型见 01;怎么用 Python 写出 这些流程见 03;callTool 的工具面细节见 04。
旧版对照(已移除): v1 的
TaskExecutionWorkflow一步一工作流、STEP_TO_ACTIVITY映射表、transitions状态机表已全部不存在。v3 是一个工作流跑完整棵 IR 树,状态不再 落独立的 transitions 表,而是从 Temporal 历史派生投影。
1. 这是什么(零基础也能懂)
一句话定义: 执行引擎 = 纯解释器(控制流)+ 注入的 Env(副作用)+ durable 后端 (Temporal/DBOS,把 Env 的每次调用变成可恢复的检查点)。
它解决什么问题。 agent 流程的难点和旧版一样:跑 到一半进程挂了要能续跑、模型/工具调用 很贵不能盲目重跑、执行过程要可解释。v3 的解法升级为:把"流程结构"提前编译冻结,运行期 只剩一棵确定性可重放的求值。
Temporal 是什么(一句话): 你把业务写成一个普通 async 函数(工作流),Temporal 记录 函数做过的每次外部调用(activity)及结果;进程挂掉重启后它重放这段历史、跳过已完成的 调用,让函数从上次位置继续。
一句话直觉: 把引擎想成"带行车记录仪的自动驾驶"——路线(IR)出发前就冻结了;开车时 每一步操作都被记录仪(Temporal 历史 + 投影事件)拍下;车熄火了,换辆车按录像接着开。
2. 顶层全景(它大概怎么转)
2.1 一次 run 从发起到落账
怎么读这张图: 左列是控制平面(HTTP/CLI),中列是 Temporal 服务端,右列是 worker 进程。 箭头是控制流;虚线框内是引擎内部。
控制平面(server/routes/runs.py) Temporal 服务端 worker 进程
───────────────────────────── ──────────────── ─────────────────────
POST /runs {release, pipeline, input}
│ 校验 release/幂等键 → runs 表插一行(submitting)
▼
TemporalGateway.start ──────────▶ task queue ──────────▶ FlowWorkflow.run()
│ 返回 workflow id │ FlowInput{flow_json,
▼ │ 状态推进 accepted→running │ manifest_json, policy…}
runs.status=running ▼
interpret(Node 树, Env)
GET /runs/{id}/events ◀── 投影事件批量落库 ◀────────── 每节点:Planned→(调用)→Did/Failed
GET /runs/{id}/gates ◀── openGates 查询 ◀────────── human gate 等待时
POST /runs/{id}/signals/human ─submitHuman 信号─────▶ 释放等待,继续求值
GET /runs/{id}/result ◀── 终态回填(reconcile)◀────── run 结束 / 失败
部件一句话职责:
| 部件 | 干什么 | 在哪个文件 |
|---|---|---|
interpret | 纯解释器:走 IR,产 Result + 投影事件 | julep/execution/interpreter.py:209 |
Env(协议) | 副作用接口:run_call/invoke_reasoner/run_sub/… | julep/execution/interpreter.py:155 |
InMemoryEnv | 测试/dry_run 用的内存实现 | julep/execution/interpreter.py:821 |
FlowWorkflow | Temporal 流程工作流(一个 run 一个) | julep/execution/harness.py:1189 |
AgentWorkflow | agent 循环专用工作流(可独立截断) | julep/execution/harness.py:2415 |
SessionWorkflow | 会话轮次工作流 | julep/execution/harness.py:1759 |
callTool / invokeReasoner | 两大效果 activity(后端中立) | julep/execution/effects.py:909 / :1138 |
build_worker | 组装 worker:注册全部 workflow+activity | julep/execution/worker.py:193 |
| DBOS 后端 | 同一 IR 跑在 dbos-transact 上 | julep/execution/dbos_backend.py:708 |
2.2 主线走一遍(高层,不进代码)
- 发起:控制平面按 release 找到冻结的
flow_json + manifest_json,组装FlowInput(julep/execution/harness.py:494),start_workflow;runs 表插一行。 - 进入工作流:
FlowWorkflow.run()(julep/execution/harness.py:1496)解析 policy/manifest, 若带 run secrets 必须先过 MCP preflight(activitypreflightMcp,julep/execution/harness.py:217)。 - 求值:构造 Temporal 版 Env(每个 handler =
workflow.execute_activity(...)), 调interpret(flow, input, env)。每个节点激活先发Planned,成功发Did,失败发Failed。 - 副作用:PRIM 叶子的
callTool调 MCP 或自家 HTTP 服务;ThinkStep走invokeReasoner调 LLM;SubStep起子FlowWorkflow;APP起子AgentWorkflow。 - 观测外送:投影事件按批从工作流内部经 activity
persist_projection_batch落库 (egress 循环julep/execution/harness.py:1309),大值先putBlob。 - 收尾:run 结束,控制平面 reconcile 把 runs 表推到终态。