源码解析全文
- M-Agent Durable Run 源码解析:从公开 API 到持久化记录、恢复分支与补充反例,附 5 段动画。
- Adapter 源码解析:模型、Run Store、PayloadCodec 与 Telemetry 怎样兑现 Runtime 契约,附 5 段动画。
- Companion 源码解析:连续对话、上下文、运行前路由与独立评估,附 5 段动画。
- Testing 源码解析:故障实验怎样形成有边界的验收结论,附 5 段动画。
以下保留 Runtime 的入门主线。精校全文补充了一个重要反例:Tool Step 已经 SUCCEEDED、但 Checkpoint 尚未提交时崩溃,当前实现可能重复发送非幂等通知。已有 WAITING 测试不能覆盖所有提交前窗口;具体条件与复现脚本见 Durable Run 全文第 7.6、12.8 节。
开场:这条通知到底发送了几次
客服 Agent 查询订单、更新工单、发送处理通知。通知服务已经接受请求,但 M-Agent 还没把工具返回的结果保存下来,进程就被终止了。
新进程打开数据库,只能看到:Run 仍为 RUNNING,模型确实请求过通知工具,该 Tool Step 有一个 RUNNING Attempt,却没有 Tool Checkpoint。这只能证明一次执行尝试的身份已经被记录,不能单独证明工具是否调用过,更不能证明通知有没有发出。
直接重试可能重复发送,直接补写成功又缺少结果。在测试覆盖的 Step / Attempt 仍为 RUNNING 的窗口,M-Agent 对这类非幂等操作进入 WAITING(UNCERTAIN_NON_IDEMPOTENT),由应用查询外部记录后决定如何继续;这不是任意崩溃窗口下的保证。贯穿本文的关键问题就是:运行时已经确认了什么,还有什么不能证明?1
先看应用怎样使用它
M-Agent 是一个嵌入 Python 应用的运行库。应用注册 Agent 定义,传入模型与存储实现,再显式创建、启动或恢复 Run。Runner 不自行领取队列任务,也不负责启动服务。
下面是一个完整的最小离线例子,使用确定性模型观察创建与启动的区别。在安装匹配的 m-agent==0.5.1 后可运行;它不调用真实模型,也不演示跨进程持久化。
import asyncio
from m_agent import AgentDefinition, DefinitionRegistry, Runner
from m_agent.adapters import (
DeterministicModelAdapter,
InMemoryRunStore,
PlaintextPayloadCodec,
)
async def main():
registry = DefinitionRegistry()
registry.register(AgentDefinition.for_adapter(
definition_id="support",
version="1.0",
instructions="Answer the support request.",
model_adapter=DeterministicModelAdapter(
responses=("Your request has been received.",),
),
))
store = InMemoryRunStore(
payload_codec=PlaintextPayloadCodec(),
)
runner = Runner(registry=registry, store=store)
created = await runner.create_run(
"support", "1.0", "Please check my order.",
)
print(created.status.value)
terminal = await runner.start_run(created.run_id)
print(terminal.status.value, terminal.output)
asyncio.run(main())
预期先得到 CREATED,再得到 SUCCEEDED 和确定性回复。这是预期输出,不是本次上游运行报告。InMemoryRunStore 随进程消失,不能用于这个例子的跨进程恢复;PlaintextPayloadCodec 也不提供保密性。真实恢复需要持久化 Store、新进程注册精确定义,以及满足租约接管条件。2
阅读地图:五种记录不是同一份日志
| 对象 | 回答的问题 | 示例 |
|---|---|---|
| Run | 这次输入对应的整次执行是什么状态? | 处理客户的这次订单请求 |
| Step | 哪个逻辑工作已经完成? | 查询订单、更新工单各是一个 Tool Step |
| Attempt | 这个步骤具体尝试过几次? | 第一次中断,第二次重放 |
| Checkpoint | 哪个完整结果已经持久化,可以复用? | 已确认的订单查询结果 |
| RunUpdate | 界面此刻可以显示什么进展? | 临时输出字符或步骤开始通知 |
前三类执行记录与 Checkpoint 保存在 Run Store,参与恢复;RunUpdate 是实时通知,断线可能丢失,不参与恢复。Step 只有 CONTEXT、MODEL、TOOL 三类。重试保持同一个 step_id,创建新的 attempt_id;应用确认也可以形成一个确认结果的 Attempt,并不必然表示又调用了一次工具。3
Run 的状态与 Attempt 的状态也不同。一个 Run 可以最终成功,同时保留此前失败的 Attempt:
CREATED 已创建,尚未执行
RUNNING 已进入执行状态,也可能因崩溃而遗留
WAITING 等待外部条件或应用决定
SUCCEEDED / REJECTED / FAILED / CANCELLED 四个终态
崩溃不会自动替运行记录写入 FAILED。恢复仍要取得有效租约,不能只凭 RUNNING 就认为自己有权推进。3
create_run:先让执行意图成为事实
Runner.create_run() 依次解析精确定义、冻结 Snapshot、保存输入与可选历史,创建一个 CREATED Run。此时不调用上下文、模型或工具。
这让应用先取得稳定的 run_id,把它关联到工单或队列项,然后再启动。创建成功之后、启动之前即使进程退出,新进程面对的也是一个明确的“尚未执行”记录,而不是猜测上次是否已经调用了模型。
Snapshot 冻结定义版本、工具影响类型、重试策略、Model Contract(模型契约)与预算等执行声明;Registry 保留当前进程真正可调用的 Python 对象。二者分开,避免把代码对象或凭证序列化进数据库。恢复必须找到原 definition_id + version,不能悄悄使用最新配置。4
start_run:先取得推进权,再进入执行循环
start_run() 只接受 CREATED。Runner 校验当前 Adapter 与冻结契约,按观察到的 Run version 取得 Lease,再执行 INPUT 策略检查;允许后才把 Run 转为 RUNNING。
下面是正常路径的顺序示意,不是可执行代码。可选阶段只在定义中声明后发生:
create_run -> CREATED + frozen snapshot
start_run -> lease + INPUT policy -> RUNNING
-> Context Step:读取并保存外部上下文
-> 可选压缩:独立的 Model Step
-> Model Step:完整响应写入 Checkpoint
-> 有工具请求:按原顺序逐个执行 Tool Step
-> 每个完整工具结果写入 Checkpoint
-> 回到下一个 Model Step
-> 无工具请求:最终策略与输出契约检查
-> SUCCEEDED + final output
同一 Run 的工具按顺序执行,一个工具结果确认后才开始下一个。这项取舍限制了单 Run 的工具并行吞吐,但让恢复能够根据模型请求列表与已确认结果,从第一个缺失的工具结果继续。不同 Run 仍可由应用并发推进。5
三类 Step:为什么完整结果必须先保存
三类 Step 都需要在外部调用前留下尝试身份,完整结果落入 Checkpoint 后再推进依赖它的工作,但风险不同。
Context:保存这次 Run 看见的资料
Context Provider 读取工单、政策等外部资料,返回带来源的 Context Items。Runner 按冻结的 ContextPlan 推进,保存完整的 Stage Result,包括阶段身份、输出内容、来源、变换决定与计量。
若在 Checkpoint 之后崩溃,新进程复用这份资料,不重新查询 Provider。若在保存之前崩溃,则可能重新读取,也可能看到已经变化的外部内容。因此它冻结的是该 Run 已确认的观察结果,不是一个通用缓存。外部资料作为数据传给模型,与受信 instructions 分开。6
Model:预留一次调用额度,再等待完整响应
模型 dispatch 前,Store 原子检查并预留整个 Run 和对应用途的 Attempt 配额。预留成功后,即使取消、失败或进程退出,这次额度仍已消费;重启不能把计数归零。
Retry Policy 决定同一 Step 是否还允许再尝试,Model Execution Budget 决定整次 Run 及某个用途还剩多少调用额度。两者都允许,才有新 Model Attempt。预留也不等于远端已经收到请求,因为最后的契约或租约检查仍可能阻止 dispatch。
Checkpoint 只保存完整 ModelResponse,包括正文、工具请求和实际可用的 usage。流式 delta 不是完整响应,不能拿来恢复。模型语义压缩和输出修复同样是独立的 Model Step,不能隐藏成无成本的字符串处理。6
Tool:外部效果与本地确认之间仍有空隙
工具调用前记录 Step 与 Attempt,工具返回明确的 ToolOutcome 后,通过相应策略关口,再写完整 Checkpoint。工具异常形成失败 Attempt,不伪装成模型可见的自然语言成功结果。
关键空隙仍然存在:通知服务接受请求,与 SQLite 提交工具结果不属于同一事务。即使 Python 调用已经返回成功,只要完整结果尚未确认持久化,新进程就不能依赖那份已丢失的内存。6
SQLiteRunStore:重启后实际读取什么
Run Store 是恢复的权威来源;Trace 和实时更新都不是。SQLite 参考实现把可查询的 Metadata 与内容 Payload 分开:状态、版本、租约和错误码可查询,输入、指令、模型响应与工具结果通过 PayloadCodec 读写。
Step 的完整身份是 (run_id, step_id),Attempt 的完整身份是 (run_id, attempt_id)。不同 Run 可以使用相同的确定性步骤 ID,但不能覆盖彼此的记录。0.5.1 正是修复了共享 SQLite 时这类 Context 与压缩步骤身份冲突。7
SQLite 的 Checkpoint 写入返回前已经提交,metadata 与对应 payload 在该写入边界内落盘。Runner 的步骤与状态写入同时受 version、Lease owner 和过期时间条件保护;不是先查一下“我有锁”,就允许之后无条件写入。模型额度的检查与预留也在事务内完成。
这些保证增加了持久化与迁移成本。它们提供本地或单服务场景中的恢复参考实现,不把 SQLite 变成通用分布式后端,也不为工具的外部系统提供事务。7
resume_run:按证据选择继续方式
应用重开同一 Store、注册原定义后,显式调用 resume_run(run_id)。Runner 先判断 Run 状态与等待原因,再获取可用租约,读取 Steps、Attempts 与 Checkpoints,重建执行位置。
终态不能再次推进。创建后尚未启动的 Run 可进入首次启动路径;因旧定义缺失而等待的 Run,在补回精确定义后可以恢复;因非幂等效果不确定而等待的 Run,普通 resume_run() 只返回原 WAITING,不能绕过应用处置。8
| 持久化证据 | 恢复动作 | 要特别注意的条件 |
|---|---|---|
| Context 有完整 Checkpoint | 复用原 Stage Result | 不重新读取外部资料 |
| Context 没有完整 Checkpoint | 重新调用 Provider | 外部内容可能已经变化 |
| Model 有完整 Checkpoint | 复用响应 | 最终响应仍需完成尚未提交的输出检查 |
| Model 有未确认 Attempt、无 Checkpoint | 保留不确定历史,再决定新 Attempt 或失败 | 冻结 Retry Policy 与新调用预算都必须允许 |
| Tool 有完整 Checkpoint | 复用 ToolOutcome | 不再次调用工具 |
| Tool 还未建立执行 Attempt | 走正常首次调用路径 | 仍要通过定义、策略与租约检查 |
| 只读或幂等 Tool 有未确认 Attempt | 原 Step 下有界重放,或失败 | 需要冻结重试授权,外部幂等性由集成实现 |
| 非幂等 Tool 有未确认 Attempt | WAITING | 不能自动重发 |
“没有 Checkpoint”不是一个统一的重试理由。Model 的未确认调用可能重复计费;Context 重读可能看到新资料;非幂等 Tool 重发则可能重复改变业务状态。恢复按这三种风险分别处理。8
普通失败的自动重试只接受 TRANSIENT 并要求冻结策略允许;崩溃后的 UNCERTAIN Attempt 有专门的恢复分支。READ_ONLY 或 IDEMPOTENT 标签本身都不授予无限重试权。
resolve_run:把通知的结局交回应用
回到开场的通知。应用从独立通知记录中确认它已经发送,就可以提交 CONFIRM_STEP(result)。Runner 校验当前 WAITING 原因、目标 Step、观察版本与 Lease,再用应用提供的结果形成 ToolOutcome;通过 TOOL_OUTCOME 策略后保存 Checkpoint,继续剩余步骤。
以下是接口片段,不是独立应用。verified_result 必须来自外部查证;runner 已连接原 Store、注册精确定义,且租约条件允许处置:
from m_agent import RunStatus
from m_agent.runtime import RunResolution
async def confirm_notification(runner, run_id, verified_result):
waiting = await runner.get_run(run_id)
if (
waiting.status is not RunStatus.WAITING
or waiting.waiting_reason != "UNCERTAIN_NON_IDEMPOTENT"
):
raise ValueError("Unexpected waiting state")
return await runner.resolve_run(
run_id,
RunResolution.confirm_step(
verified_result,
waiting_step_id=waiting.waiting_step_id,
),
expected_version=waiting.version,
)
CONFIRM_STEP 不调用通知工具。它保存的是应用确认的结果;新的成功 Attempt 不意味着第二次发送。若应用查证后选择再次执行,则使用 RETRY_STEP,那才会再次调用工具。
版本检查防止“我确认的是刚才看到的等待状态,但别人已经处理过了”的冲突。确认结果还需要通过策略关口,后续模型也可能失败,所以 resolve_run() 并不无条件返回成功。
应用还可以选择 FAIL_RUN 或 CANCEL_RUN。它们结束 Run,不删除旧记录、不回滚通知。当前版本不将这两种命令的 reason 保存到 Run Store,需要操作者与处置理由审计时由应用另行记录。9
取消与流式输出:意图不能冒充执行结果
cancel_run() 表示停止继续推进,不是强杀已发出的请求。正在本 Runner 中执行的 Run 会在安全边界观察取消事件;已经调用的工具可能仍会完成,Runner 先保存实际结果,再阻止后续步骤。
如果非幂等工具以不确定结果结束,WAITING 优先于取消,避免用 CANCELLED 遮住“通知可能已经发出”。取消事件属于当前 Runner 实例,跨进程取消通道仍由应用实现。
流式输出同样区分临时视图与事实。假如 Attempt A 输出 AB 后失败,Attempt B 输出完整 ABC,界面应按新 attempt_id 替换临时文本,而不是拼成 ABABC。只有 B 的完整 ModelResponse 可以进入 Checkpoint。
subscribe_run() 发布的 MODEL_DELTA、步骤开始和完成通知是 live-only;断线不回放历史。应用需要核实时使用 get_run() 或 inspect_run() 读取 Store。10
Lease 与 version:谁有权推进,依据哪个状态
Lease 解决“当前谁能推进这个 Run”,version 解决“命令是否基于当前生命周期事实”。只做其中一项不够:版本不能单独阻止两个人同时发外部请求,租约也不能让过期的应用确认变得有效。
旧进程卡住后,Lease 过期,新 Runner 可以在应用显式恢复时接管。旧调用稍后仍可能返回,但其迟到写入会因 owner、过期时间或版本条件被拒绝。拒绝迟到提交不等于撤销旧的外部请求。
当前 Runner 不后台自动续租,也不扫描 Run 自动接管。应用要选择覆盖预期调用与提交窗口的 TTL,安排恢复触发;太短增加冲突,太长增加崩溃后的等待。11
接下来读什么
到这里,Durable Run 可以概括为一条可检查的协议:保存身份,记录尝试,确认完整结果;中断后按原声明复用或重放,证据不足则停下,把外部判断交回应用。
这只展开了 Runtime Core 的主线。Durable Run 精校全文、Adapter 全文、Companion 全文与 Testing 全文已独立收录,保留各自的源码修订、验证范围、插图和视频。
设计页解释上述取舍,实验验证对照故障窗口与测试断言,迭代记录说明各层能力的演进关系。
来源与范围
本文按用户提供的《M-Agent Durable Run 源码解析》重组,基线为 v0.5.1,链接固定到 tag 对应的提交 ae85d4d。调用片段是阅读用例,不是本次验收报告;本次未运行上游测试或真实模型端点。源码、测试断言和作者的设计解释不能互相替代。
参考资料
Runner.create_run。原文第 2 章。 ↩
Runner.start_run;_execute_steps。原文第 3、4 章。 ↩
Context / Model / Tool 步骤实现。原文第 5 章,分别定位
_run_context_stage、_reserve_model_attempt、_run_single_model_step、_run_tool_step。 ↩ ↩ ↩SQLiteRunStore;迁移与恢复策略。原文第 6 章。 ↩ ↩
resume_run 与证据重建。原文第 7、8 章及附录 G。 ↩ ↩
resolve_run;合法处置动作。原文第 9 章。 ↩