跳到主要内容

数据截至 (上游 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
FlowWorkflowTemporal 流程工作流(一个 run 一个)julep/execution/harness.py:1189
AgentWorkflowagent 循环专用工作流(可独立截断)julep/execution/harness.py:2415
SessionWorkflow会话轮次工作流julep/execution/harness.py:1759
callTool / invokeReasoner两大效果 activity(后端中立)julep/execution/effects.py:909 / :1138
build_worker组装 worker:注册全部 workflow+activityjulep/execution/worker.py:193
DBOS 后端同一 IR 跑在 dbos-transact 上julep/execution/dbos_backend.py:708

2.2 主线走一遍(高层,不进代码)

  1. 发起:控制平面按 release 找到冻结的 flow_json + manifest_json,组装 FlowInput (julep/execution/harness.py:494),start_workflow;runs 表插一行。
  2. 进入工作流:FlowWorkflow.run()(julep/execution/harness.py:1496)解析 policy/manifest, 若带 run secrets 必须先过 MCP preflight(activity preflightMcp,julep/execution/harness.py:217)。
  3. 求值:构造 Temporal 版 Env(每个 handler = workflow.execute_activity(...)), 调 interpret(flow, input, env)。每个节点激活先发 Planned,成功发 Did,失败发 Failed
  4. 副作用:PRIM 叶子的 callTool 调 MCP 或自家 HTTP 服务;ThinkStepinvokeReasoner 调 LLM;SubStep 起子 FlowWorkflow;APP 起子 AgentWorkflow
  5. 观测外送:投影事件按批从工作流内部经 activity persist_projection_batch 落库 (egress 循环 julep/execution/harness.py:1309),大值先 putBlob
  6. 收尾:run 结束,控制平面 reconcile 把 runs 表推到终态。

3. 核心机制一:纯解释器与 Env 协议

它要解决的小问题: 控制流逻辑(并行汇合、竞速取消、分支谓词、有界迭代)既要在 Temporal 上跑,又要在单测里不用 Temporal 就能验证——怎么分?

思路: 把"走树"和"做事"彻底分开。interpret()(julep/execution/interpreter.py:209) 只做前者;后者全部收敛到 Env 协议(julep/execution/interpreter.py:155):

# 示意,非源码:Env 协议的关键成员(真源码是 Protocol 类)
class Env(Protocol):
manifest: ToolManifest # 冻结的工具清单
emitter: ProjectionEmitter # 投影事件出口
def next_cid(self, node_id): ... # 确定性激活 id 计数器
def get_pure(self, name): ... # 按名取注册纯函数
async def run_call(self, node, value, cid): ... # 调工具
async def invoke_reasoner(self, reasoner, value, cid, timeout_s, batchable): ...
async def run_sub(self, ref, contract, value, cid, node_id): ... # 子流程
async def run_agent(self, controller, value, cid, config): ... # agent 循环
async def compile_plan(self, controller, value, cid): ... # 运行期计划
async def human_gate(self, value, cid, timeout_s): ...
async def sleep(self, seconds, cid): ...
def gather(...) / race_first(...) # 确定性并发原语

确定性契约写死在模块 docstring(julep/execution/interpreter.py:1):解释器不引入 墙钟、随机数、环境 IO;cid 来自确定性计数器;arr/alt 谓词和 reducer 必须是具名纯函数。 "任何非确定的东西都该在 activity 里,躲在 Env 边界后面。"

解释器怎么分派 11 种算子——_eval(julep/execution/interpreter.py:295)是一个直白的 if 链,每种 Op 一段:

Op解释器行为
IDENT原值返回
ARRenv.get_pure(name)(value, **args)
PRIM交给 _eval_prim(§4)
SEQ左结果作右输入,事件因果链相接
PAR_eval_par:两分支同输入并发,按 Merge 策略汇合(:463)
EACH_eval_each:逐项跑 body,收集输出(:504)
ALT谓词或 select 选分支,只解释选中者
ITER_UP_TO最多 bound 轮,收敛谓词为真提前停
LOOP无限循环直到 SessionClosed;split 模式把中间输出从通道 emit 出去
EVAL_PLAN有烘焙计划直接跑;否则先 env.compile_plan(运行期计划)再解释
APPenv.run_agent(...) 整个委托给 AgentWorkflow

投影事件在哪发。 interpret() 每次进入先 env.emitter.plan(...)Planned (带上游 causes);求值成功发 Did(带成本与 value_ref);任何异常路径都保证发 Failed 再重抛(julep/execution/interpreter.py:209 的 try/except 结构)。 投影在解释器里是"顺手记账",不是数据来源——同一份历史重放会重导出同样的事件。


4. 核心机制二:重试代数(契约驱动)

它要解决的小问题: 工具调用失败,重试还是不重试?盲目重试一个"扣款"类工具会重复扣钱。

v3 的答案:把决定权交给工具契约,且只认"被断言的"契约。 _eval_prim 里 call 叶子的 重试参数来自三个函数(julep/execution/interpreter.py:100 起的 call_contract 等):

# 示意,非源码:重试资格判定链(真源码为 _retry_attempts_for_call 等)
if node.ann.max_attempts is None: attempts = 1 # 没写就不重试
elif not call_contract_asserted(node): attempts = 1 # 契约非断言 → 不重试
elif not contract_allows_retry(contract): attempts = 1 # 契约说不能重试 → 不重试
else: attempts = max(1, ann.max_attempts)

三层含义:

  1. 显式 opt-in:重试次数写在节点的 Ann.max_attempts 上,默认不重试。
  2. 断言才可信:契约来自 capability manifest(Deployment 冻结后 frozen_hash 可查) 才算"断言";MCP 自己标的 idempotent 不算(call_contract_asserted, julep/execution/interpreter.py:126)。
  3. 策略错误永不重试:POLICY_ERRORS(julep/errors.py:208:CapabilityDenied、 PlanRejected、ValidationError、FreezeError、PureDriftError)是"定案的决定",重试没有意义 (julep/execution/interpreter.py:430 的 except 首位直接 raise)。

重试间隔/退避同样来自 Ann(retry_interval_s/backoff_rate),等待用 env.sleep (在 Temporal 里是持久定时器,不占 worker)。Reasoner 叶子不在这套代数里——LLM 调用 的弹性属于注入的 LlmCaller(消费方的重试/降级栈),DBOS 后端的 docstring 明确 "Reasoner steps never retry"(julep/execution/dbos_backend.py:1)。

Temporal 侧还有一层 activity 级 RetryPolicy,由 _retry_policy_for (julep/execution/harness.py:436)从冻结契约推导:read/天然幂等工具宽裕重试; 非幂等写谨慎,且只在 callTool 发出的 Idempotency-Key 保护下重试(见第 04 章)。


5. 核心机制三:Temporal harness——工作流怎么把效果变 durable

5.1 FlowWorkflow:一个 run 一个工作流

FlowWorkflow(julep/execution/harness.py:1189)是 v3 唯一的流程工作流类型——"注册" 不是代码生成,部署物只是数据(julep/deploy.py:1 docstring)。它自带三类交互面:

  • 信号:submitHuman(:1203,按激活 id 投递人工决定)、submitReasonerResult (批量 reasoner 结果回填)。
  • 查询:openGates(:1244,当前等人审批的激活 id 列表——驱动审批 UI 的精确清单)、 projection(:1256,只读的 pomset 快照)。
  • 等待:_await_humanworkflow.wait_condition 挂起(可超时),期间不占 worker。

run()(:1496)的骨架:解析 FlowInput → (ref-only 输入先经 resolveSubflow activity 取回 冻结物)→ 若带 run secrets 且无已完成 preflight 状态则拒绝 → 可选 preflightMcp → 构造 Env + 投影 emitter → interpret(...) → 投影收尾/落账。

大 payload 的处理:投影事件与大值通过 putBlob(julep/execution/effects.py:697) 以 canonical JSON 落 BlobStore,事件里只留内容引用;Temporal 原生 payload 另有 ClaimCheckCodec 自动把超阈值 payload 换成 blob 引用(§6.3)。

5.2 continue-as-new:让历史有界的两处

v1 用"每步一个子工作流"控历史;v3 的 IR 树在一次 run 里整体求值,但有边界处的截断:

  • Flow 链式续跑:当结果是"续体"(is_continuation,见 julep/continuation.py), run() 末尾 workflow.continue_as_new(...)(julep/execution/harness.py:1708)——同一份 冻结流程、新输入、累计的 call_counts 让 maxCalls 预算横跨整条链、principal 保持租户 身份、segment_seq +1 供轨迹缝合。
  • Agent 轮次截断:AgentWorkflow._run 里每当 should_continue_as_new (julep/agent_loop.py:592:本段轮数超过 continue_as_new_after)就用 continue-as-new 换新段(:3417 附近),状态经续体封套过界。Agent 故意是独立工作流,这样它的截断 只影响自己,不牵连父流程历史(julep/execution/harness.py:1 docstring)。
  • Session 轮边界截断:SessionWorkflow._should_continue_as_new(:2141)只在 轮边界(quiesce)截断,:2159 的注释写明"绝不能在 body 中途"。

5.3 AgentWorkflow:有界控制器循环的 durable 形态

AgentWorkflow(julep/execution/harness.py:2415)的 run(:2498)把 agent_loop 的纯控制逻辑(第 01 章 §6.2)接上 activity:每轮 invokeReasoner 问控制器 → interpret_reasoner_reply(julep/agent_loop.py:138)归一成 RoundActioncall 则走 callToolsub 则起子 FlowWorkflow、finish/escalate 则收尾。 带 use_session_store 时强校验 session_id == workflow_id(1:1,Temporal 的 "一个 workflow id 只能有一个运行中的执行"就是会话存储的互斥锁,:2498 起的校验)。


6. Worker 与活动:进程怎么装配

6.1 build_worker:一次注册全部

build_worker(julep/execution/worker.py:193)把 configure(context) (julep/execution/effects.py:430,进程级 WorkerContext :106)和一个 Temporal Worker 接起来:每个角色一个 workflow 类型(Flow/Agent/Session)+ 一页 ACTIVITIES 清单 (julep/execution/worker.py:92:callTool、invokeReasoner、compilePlan、verifyPures、 resolveSubflow、resolveAgentSpec、resolveQoS、loadState/commitState、loadValue/commitValue、 putBlob、projection 持久化等)。workflow 代码本身不持有任何环境配置——全部配置 (URL、MCP caller、LLM 客户端、能力清单)一次性注入 WorkerContext,workflow 保持可重放 (julep/execution/worker.py:1 docstring)。

沙箱放行名单是个值得看的细节:workflow 沙箱默认放行 julepwasmtime (julep/execution/worker.py:193 的 docstring)——不放行的话沙箱会重新 import 出空的 注册表,arr/注册依赖的叶子会在 WorkflowTaskFailed 里死循环;wasmtime 放行则让 bundle 源的 wasm 纯函数共享进程级 Engine。bundle 解析(BundleResolvingWorkflowRunner) 在 worker 启动时一次性完成、绝不放进 workflow 代码——重放时纯函数必须已注册 (julep/worker_store.py:1 docstring)。

6.2 两大效果 activity

callTool(julep/execution/effects.py:909)按 ToolRef 分派:

  • MCP 工具:先按冻结 schema 校验输入(不合法抛 ToolInputValidation),再走注入的 MCP caller(见第 04 章)。
  • native 工具:查 WorkerContext.tool_urls 拿 URL;能力清单开着就先检查网络出网域 (capabilities.network_allows);HTTP POST 带两个头——Idempotency-Key: cid (激活 id 即幂等键)与可选的 principal 头。

调用后 _capture_effect(:744 起)尽力而为地写轨迹(失败绝不影响结果)。 invokeReasoner(:1138)则组装 prompt(渲染器、transcript 水合、token 预算切分—— 第 05 章)、调注入的 LlmCaller、按 reply schema 校验并按 output_retries 重问。

6.3 payload codec:加密 + claim check

Temporal 客户端与 worker 之间过线的数据由 PayloadCodecChain (julep/execution/codec.py:123)处理,两个可组合的 codec:

  • AesGcmPayloadCodec(:35)——AES-GCM 载荷加密(run secrets 只存在于加密载荷里, 不落存储与投影);
  • ClaimCheckCodec(:144)——超阈值 payload 换成 BlobStore 引用(claim check), 大对象不进 Temporal 历史。

对照 v1 的"pickle+lz4 + 手工 offload"(已移除),v3 把这两件事做成了标准 codec 链, 业务代码零感知。


7. DBOS 后端:同一 IR 的第二落点

julep/execution/dbos_backend.py 用 dbos-transact(Postgres 上的轻量 durable 执行)实现了 同一套 Env:flow_workflow(:708)解释冻结 IR,每个效果叶子是 @DBOS.step,委托给同一个 effects 层(julep/execution/dbos_backend.py:1 docstring)。差异是能力上的诚实声明:

能力TemporalDBOS
race/hedge/quorum支持(可取消分支)拒绝执行(assert_dbos_executable:114——DBOS 无法取消在跑的 step,竞速语义会撒谎)
子流程子工作流(隔离历史)内联跑进父历史(换简单性)
Agent 循环AgentWorkflow + continue-as-new每 history 段一个 DBOS 工作流(run_agent_dbos:1244),状态走续体封套
会话存储路径支持Temporal-only

这个"不能安全做就不做"的姿态,和 §4 的重试代数一脉相承。


8. 本地执行:dry_run 与 LocalPipeline

编译/测试不需要任何引擎:

  • Deployment.dry_run(julep/deploy.py:865)在内存里跑:InMemoryEnv (julep/execution/interpreter.py:821)的效果是普通 callable,reasoner 用确定性假函数 (README.md:76)。
  • julep run(CLI)与 LocalPipeline(julep/local.py:168)走配置化的前台执行: 同一份配置/编译/解释器,但"无 PostgreSQL、无 Temporal、无 HTTP 控制面跳板、无 durable 重试边界"(julep/local.py:1 docstring);入口 run_local_pipeline(julep/local.py:441) 与 arun_local_pipeline(:416)。

9. 巧妙之处(可借鉴的技术)

  • 纯解释器 + Env 协议。 控制流逻辑(含并发汇合、竞速、有界迭代)完全与引擎解耦, 单测零依赖;两个生产后端只是两个 Env 实现。见 julep/execution/interpreter.py:155

  • 投影在解释器里"顺手记账"。 interpret() 每个节点三拍子(Planned→Did/Failed), 因果链用事件 id 串成 pomset;持久化是可选的 egress,重放即重导出。 见 julep/execution/interpreter.py:209

  • 重试资格三连问。 写了吗(Ann)?断言了吗(manifest)?契约允许吗(effect/幂等)? 三问不过就不重试;策略错误永不重试。见 julep/execution/interpreter.py:100-131julep/errors.py:208

  • 激活 id(cid)一身三职。 它是投影事件的激活标识、native 调用的 Idempotency-Key、 human gate 的等待键——一个确定性计数器把观测/幂等/交互串了起来。 见 julep/execution/effects.py:909

  • continue-as-new 只在边界截断。 Flow 链的续体边界、Agent 的轮次阈值、Session 的轮边界 ——三处都是"语义上安全的停顿点",而不是定时粗暴截断。 见 julep/execution/harness.py:1708julep/agent_loop.py:592


10. 边界与局限(诚实)

  • 投影不是持久层。 Temporal 历史才是;FlowWorkflow.projection 查询明确标注 "derived, not durable"(julep/execution/harness.py:1256)。
  • 轨迹捕获是尽力而为。 sink/blob 任何失败都被吞掉且不影响结果字节 (julep/execution/effects.py:726_capture_effect 的 docstring)。想要强审计, 投影 + runs 表才是主证据。
  • run secrets 路径收紧。 带 secrets 的 run 必须携带"workflow 拥有的 preflight 绑定状态", 裸 Temporal start 直接拒绝(julep/execution/harness.py:1496 起的校验)。
  • DBOS 后端是能力子集(§7 表),选它就放弃了 race 族。
  • 本地执行无 durable 保障。 LocalPipeline 只保证同构解释,不保证崩溃恢复 (julep/local.py:1)。

11. 横向对比

同 shelf 里"驱动 agent 长流程"的项目,取舍不同:

维度Julep v3(本章)内存循环型 agent 框架
流程表示冻结 IR(编译期定型)代码路径(运行期展开)
执行载体纯解释器 + Temporal/DBOS/内存三后端单进程 while 循环
崩溃恢复引擎历史重放,原生自建 checkpoint 或无
重试契约驱动代数,策略错误不重试通常一揽子指数退避
工具面防漂移内容哈希冻结 + preflight一般不做

一句话取舍:v3 把"可信执行"做成编译产物的一部分——想要它就得接受"先编译再跑"的 心智;想要随手写循环的自由,它不是那个工具。整体服务拓扑(控制平面/CLI/发布)见 06 平台架构


12. 代码地图(导航索引)

主题文件路径符号名
纯解释器julep/execution/interpreter.py:209interpret / _eval(:295) / _eval_prim(:390)
Env 协议 / 内存实现julep/execution/interpreter.py:155Env / InMemoryEnv(:821)
契约驱动重试julep/execution/interpreter.py:100call_ref_key / call_contract / call_contract_asserted
策略错误清单julep/errors.py:208POLICY_ERRORS
并行/遍历求值julep/execution/interpreter.py:463_eval_par / _eval_each(:504)
流程工作流julep/execution/harness.py:1189FlowWorkflow.run(:1496)
人闸信号/查询julep/execution/harness.py:1203submitHuman / openGates(:1244) / projection(:1256)
MCP preflight 活动julep/execution/harness.py:217preflightMcp
轨迹起止活动julep/execution/harness.py:262startTrajectory / finishTrajectory
Flow 续跑截断julep/execution/harness.py:1708workflow.continue_as_new(Flow 链)
Agent 工作流julep/execution/harness.py:2415AgentWorkflow.run(:2498)
会话工作流julep/execution/harness.py:1759SessionWorkflow / _should_continue_as_new(:2141)
轮次截断判据julep/agent_loop.py:592should_continue_as_new
activity RetryPolicyjulep/execution/harness.py:436_retry_policy_for
工作流输入julep/execution/harness.py:494FlowInput / SessionInput(:533) / AgentInput(:573)
工具调用 activityjulep/execution/effects.py:909callTool
模型调用 activityjulep/execution/effects.py:1138invokeReasoner
运行期计划 activityjulep/execution/effects.py:1212compilePlan / verifyPures(:1308)
worker 上下文julep/execution/effects.py:106WorkerContext / configure(:430) / RunPrincipal(:68)
worker 组装julep/execution/worker.py:193build_worker / ACTIVITIES(:92)
payload codecjulep/execution/codec.py:35AesGcmPayloadCodec / PayloadCodecChain(:123) / ClaimCheckCodec(:144)
DBOS 后端julep/execution/dbos_backend.py:708flow_workflow / assert_dbos_executable(:114)
本地前台执行julep/local.py:168LocalPipeline / run_local_pipeline(:441)
run 辅助julep/execution/harness.py:3573run_flow / start_flow(:3635)