源码解析全文

以下保留 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 有未确认 AttemptWAITING不能自动重发

“没有 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。调用片段是阅读用例,不是本次验收报告;本次未运行上游测试或真实模型端点。源码、测试断言和作者的设计解释不能互相替代。

参考资料

  1. 非幂等通知恢复测试。对应原文开场。 ↩

  2. 公开导入与最小示例。 ↩

  3. 执行记录;Run 状态。原文第 1 章。 ↩ ↩

  4. Runner.create_run。原文第 2 章。 ↩

  5. Runner.start_run;_execute_steps。原文第 3、4 章。 ↩

  6. Context / Model / Tool 步骤实现。原文第 5 章,分别定位 _run_context_stage、_reserve_model_attempt、_run_single_model_step、_run_tool_step。 ↩ ↩ ↩

  7. SQLiteRunStore;迁移与恢复策略。原文第 6 章。 ↩ ↩

  8. resume_run 与证据重建。原文第 7、8 章及附录 G。 ↩ ↩

  9. resolve_run;合法处置动作。原文第 9 章。 ↩

  10. 流式输出与取消测试。原文第 10 章。 ↩

  11. 租约竞争与接管测试。原文第 11 章。 ↩