Files
PolyLoop/src/polyloop/session/__init__.py
T
iomgaa e71dca7c35 feat(session): 落成事件出口,五个接缝全部有调用点
Event 从零字段变成 kind + run_id + model_binding + step,新增 EventKind(只有一种取值,
但第一天就带 kind,逼每个出口分发)。事件在「一步走完」原子落地之后发,只有这次进程里
真的执行过的步才发;投递失败接住、计数进 RunResult、继续跑,CancelledError 原样穿过。

契约套件那两条 xfail 关掉:一条要断言的是库发了几次、接缝自己看不到;另一条的前提是错的
——审计纪律由意图日志承担不由事件流承担,改成在 unit 层验日志里原文与改写后的文本各有
位置。_project_observation 那段说「将来靠事件流送出去」的注释一并改对。

283 passed / 15 skipped / 2 xfailed,剩下两条 xfail 是原子写与前缀持久性,没有机器兜底。

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-08-10 22:51:19 -04:00

876 lines
40 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""装配层:两个入口协程,以及它们收的两个装配对象。
这个模块把五个接缝和四个纯逻辑模块编织成一次运行。停止判定的顺序在
`research-wiki/design/0004-stopping-and-step-record.md` 决策三(A 到 H 八档),本模块的主循环
是它的唯一执行者;每一档的判定本身住在 `polyloop._stopping`。
**模块名叫 `session`,但治理单位不叫 session。** 这个名字说的是「这个模块把各部分装配到
一起」,与一次运行怎么称呼是两件事——`session` 在业界普遍指一个长期存在、可以来回对话的
东西,而这里的治理单位是有界的、一次性的(`0006` 决策一)。
**事件只在一处发出去**:一步走完、`StepCompleted` 原子落地之后。顺序不能倒过来——先发后写的话,
进程崩在两者之间会让观察者看见一步而存储里没有,而事件流的全部安全性建立在「它带的事实在存储
里另有一份」上(`0013` 决策一与决策五)。
"""
import asyncio
import json
import logging
import time
from collections.abc import Mapping
from dataclasses import dataclass
from polyloop._assembly import (
assemble,
check_observation_template,
injected_entry_ids,
prompt_chars,
)
from polyloop._recovery import CorruptLogError, ResumeAction, ResumePlan, plan_resume
from polyloop._stopping import (
RunCounters,
budget_admission,
completion_verdict,
parse_failure_admission,
prompt_size_admission,
)
from polyloop.ports import (
Action,
ActionExecutor,
DecisionParser,
Event,
EventKind,
EventSink,
FinalAnswer,
InvalidDecision,
ModelCall,
ModelClient,
RunLog,
RunStore,
)
from polyloop.tools import RegistryExecutor, ToolRegistry
from polyloop.types import (
ActionOutcome,
ActionStatus,
Budget,
Context,
Injection,
Intent,
IntentKind,
Message,
ModelCallResult,
ModelReply,
ReplayPolicy,
RunFinished,
RunResult,
RunStarted,
StepCompleted,
StepRecord,
StopReason,
SyntheticObservations,
)
logger = logging.getLogger(__name__)
class RunIdentityError(Exception):
"""运行标识和日志对不上:`run` 撞上已有日志,或者 `resume` 读不到日志。
**两条入口的失败方式不同,所以不合并成一个带开关的函数**(`0006` 决策二):`run` 撞上
已有日志要报错,`resume` 读不到日志要报错。合并之后调用方看不出自己走的是哪条,而两种
错的处置完全不同。
"""
class ParameterDriftError(Exception):
"""续跑时现算的参数快照与日志里那份对不上。
**这意味着续跑不能顺便改预算或换模型**,那不是限制而是这条守卫的全部意义:用同一个运行
标识换一份定义续跑,前几步与后几步会来自两个不同的配置而全程零报错。
"""
@dataclass(frozen=True, slots=True, kw_only=True)
class AgentDefinition:
"""跨运行不变的那一半装配,可以并发复用。
四个接缝挂在这里,因为它们随装配变而不随运行变。**它自己不持有任何一次运行的状态**——
持有了的话,并发跑同一份定义的两次运行会互相写到对方的计数里。
"""
model_client: ModelClient
decision_parser: DecisionParser
store: RunStore
event_sink: EventSink
#: 库自己合成、回填给模型看的那几段观察。动作被拒绝与环境故障两档不取执行器给的那段,
#: 取这里的(`0007` 决策二)——执行器每次现造一段文本的话,那段文本既不在参数快照里、
#: 也不受任何约束,两次运行之间它可以变而不会有任何地方报错。
synthetic_observations: SyntheticObservations
def parameter_snapshot(self) -> Mapping[str, str]:
"""向四个接缝各问一次参数再聚合。
**是方法不是字段。** 写成字段就要在构造定义之前先问一遍,而那时定义还不存在;写成
方法则每次现问,快照永远是从真实对象上读出来的**事实**而不是一份**声明**。
键带接缝名前缀,免得「哪一侧报的这个键」要靠约定记住。
"""
snapshot: dict[str, str] = {}
for prefix, seam in (
("model_client", self.model_client),
("decision_parser", self.decision_parser),
("store", self.store),
("event_sink", self.event_sink),
):
for key, value in seam.parameters().items():
snapshot[f"{prefix}.{key}"] = value
return snapshot
@dataclass(frozen=True, slots=True, kw_only=True)
class RunRequest:
"""一次运行独有的那一半装配。构造廉价:无 I/O、无网络校验、无哈希计算。"""
#: 不透明字符串,库不解析。它同时是日志的主键。
run_id: str
budget: Budget
#: 动作执行接缝。挂在请求上,因为它每次运行都不同。
action_executor: ActionExecutor
#: 本次可见的那个(子)注册表。**与 `action_executor` 同时存在不是重复**:不注册工具的
#: 项目传一个空注册表加一个环境句柄,注册了工具的项目传 `tools.executor()`。
tools: ToolRegistry
context: Context
#: 本次要贴进上下文的条目,按通道分组。
injections: Mapping[str, tuple[Injection, ...]]
#: 项目自己的标识,库不解释,原样透传给每次模型调用。它的全部键值都进参数快照。
model_binding: Mapping[str, str]
#: 必填无默认。工具的重放策略能从注册表查到,模型调用的查不到——只有调用方知道这次调用
#: 能不能重来。
model_replay_policy: ReplayPolicy
#: 观察回填历史时套的格式,必须含 `{observation}` 占位符。
observation_template: str
#: 取消进来之后,库留给自己写结束记录的秒数。不写的话,恢复读到的是一次没有结束标记的
#: 运行,会被当成可以续跑,而它其实是被人主动叫停的。
cancel_grace_seconds: float
def __post_init__(self) -> None:
check_observation_template(self.observation_template)
if self.cancel_grace_seconds < 0:
raise ValueError(f"取消宽限期不能为负:{self.cancel_grace_seconds}")
# 模型看见的 schema 来自一个注册表、实际分发走另一个,表现是「模型调了一个它看得见的
# 工具却说不存在」。执行器不是注册表派生的就一概放行——那是项目自己写执行器的情形,
# 库无从判断也不该判断(`0006` 决策三)。
if (
isinstance(self.action_executor, RegistryExecutor)
and self.action_executor.registry != self.tools
):
raise ValueError(
"action_executor 派生自另一个注册表:模型看见的 schema 与实际分发会来自两份"
f"不同的工具集(执行器持有 {self.action_executor.registry!r},本次可见 {self.tools!r}"
)
def parameter_snapshot(self) -> Mapping[str, str]:
"""请求这一侧的快照,加上向动作执行接缝问的那一次。
**上下文与注入内容不进快照。** 它们是这次运行的输入数据不是参数,进快照会让快照变成
一份数据副本,而它们可能很大。注入的**条目标识**另行进轨迹,所以「这次贴了哪几条」
事后查得到,查不到的只是正文。
"""
snapshot: dict[str, str] = {
"request.max_steps": str(self.budget.max_steps),
"request.max_actions": str(self.budget.max_actions),
"request.max_consecutive_parse_failures": str(
self.budget.max_consecutive_parse_failures
),
"request.max_prompt_chars": str(self.budget.max_prompt_chars),
"request.model_replay_policy": self.model_replay_policy.value,
"request.observation_template": self.observation_template,
"request.cancel_grace_seconds": str(self.cancel_grace_seconds),
"request.tools": ",".join(self.tools.names()),
"request.injected_entry_ids": ",".join(injected_entry_ids(self.injections)),
}
# 模型绑定必须进快照,否则给它选字符串映射的那条理由就落空了。失败场景很具体:崩溃后
# 用同一个运行标识、换一组绑定续跑,后面每一次调用被记到另一套坐标上,而两段轨迹在
# 文件里看起来是同一次运行。
for key, value in self.model_binding.items():
snapshot[f"request.binding.{key}"] = value
for key, value in self.action_executor.parameters().items():
snapshot[f"action_executor.{key}"] = value
return snapshot
def _merged_snapshot(definition: AgentDefinition, request: RunRequest) -> dict[str, str]:
return {**definition.parameter_snapshot(), **request.parameter_snapshot()}
def _model_result_id(run_id: str, call_index: int) -> str:
"""预分配的模型调用结果 ID。
**从运行标识与序号推出来,不用随机数。** 恢复时要精确地问「这个 ID 的结果条目在不在」,
推得出来就意味着重放同一步时拿到的是同一个 ID——重放沿用原 ID 这件事因此不需要额外传递。
随机 ID 还会往轨迹里塞一列每次都不同的值,让两次运行的逐字段比对多一个要排除的字段。
"""
return f"{run_id}#model#{call_index}"
def _action_result_id(run_id: str, call_index: int) -> str:
return f"{run_id}#action#{call_index}"
def _serialise_arguments(arguments: Mapping[str, object]) -> str:
"""工具参数落进轨迹的那一列。
`sort_keys` 是为了同样的参数每次落出同样的字符串——不排的话,字典顺序的差异会让两次运行
的逐字段比对报出一堆假差异。`default=str` 兜住不可序列化的取值,因为这一列是给人看的
留痕,不该因为某个参数装了个对象就让整次运行失败。
"""
return json.dumps(arguments, ensure_ascii=False, sort_keys=True, default=str)
def _project_observation(
outcome: ActionOutcome, synthetic: SyntheticObservations
) -> tuple[str, bool, int]:
"""决定这一步回填进历史的观察是哪一段(`0007` 决策二)。
未执行与环境故障两档取库合成的那段,执行器返回的 `observation` 不进历史——它是「模型看
得见的东西」,而那种东西必须能进参数快照。执行器每次现造一段文本的话,两次运行之间它可以
变而不会有任何地方报错,于是「同一份配置跑出来的两次运行」在模型看来其实不同。
**执行器那段不进历史,但它没有丢**:「一步走完」那条记录落盘的是动作执行接缝的原样返回
值,替换只发生在步记录的这一列上。被拒绝那一档下,它是日志里唯一的拒绝说明(「哪个参数
不合法」只有执行器知道),所以不必再靠事件流把它送出去——事件可丢,而这段文本是审计要的
`0013` 决策六)。
"""
if outcome.status is ActionStatus.NOT_EXECUTED:
return synthetic.action_rejected, True, 0
if outcome.status is ActionStatus.ENV_ERROR:
return synthetic.env_failed, True, 0
return (
outcome.observation,
outcome.observation_is_synthetic,
outcome.observation_truncated_chars,
)
class _Driver:
"""一次运行的全部可变状态。**每次运行一个实例,绝不跨运行复用。**
并发跑同一份定义的多次运行各持一个,所以计数、步序列、时钟互不可见。装配对象那两个是
frozen 的,共享它们没有问题。
"""
__slots__ = ("_counters", "_definition", "_event_failures", "_request", "_steps")
def __init__(self, definition: AgentDefinition, request: RunRequest) -> None:
self._definition = definition
self._request = request
self._counters = RunCounters()
self._steps: list[StepRecord] = []
#: 投递失败的次数。**住在这里而不是出口上**——出口挂在定义上、可以被并发的多次运行
#: 共用,而这个数属于一次运行的结果。
self._event_failures = 0
# -- 写入 ---------------------------------------------------------------
async def _write_step(self, step: StepRecord, outcome: ActionOutcome | None) -> None:
"""动作结果与步记录一次原子落地。
存储实现要么两者都可见、要么都不可见。分开写的话,崩在两者之间会让恢复把那一步读成
「执行完了,跳过」,那一步的历史文本就永远丢了——恢复出来的消息序列比不中断跑完时少
一轮,后面每一步都跟着偏。
"""
result_id = (
None if outcome is None else _action_result_id(self._request.run_id, step.step_idx)
)
await self._definition.store.write_step_completed(
StepCompleted(
run_id=self._request.run_id,
result_id=result_id,
action_outcome=outcome,
step=step,
)
)
self._steps.append(step)
self._counters = self._counters.with_step_appended()
await self._emit_step_finished(step)
async def _emit_step_finished(self, step: StepRecord) -> None:
"""把这一步发给事件出口。**只有走到这里的步才发**,从日志里读回来直接跳过的不发。
补发已经完成的步等于宣称一件早就发生过的事刚刚发生,而接进度表的那一侧会多出一批
重复行。观察者要补全前半段,从存储里读。
**投递失败接住、计数、继续跑**:事件是观察通道不是控制通道,一次运行不该因为进度回写
的数据库连不上就终止。接的是 `Exception` 不是 `BaseException`——`CancelledError` 必须
原样穿过(`CLAUDE.md` §1.6),在这一下把取消吞掉,取消就会晚一整步才生效。
**失败不转成一条事件从同一个出口再发一次**:那会自我喂食,一个持续失败的出口会让失败
处理路径变成递归,而递归的表现是进程卡住或栈溢出,不是一条错误日志。
"""
try:
await self._definition.event_sink.emit(
Event(
kind=EventKind.STEP_FINISHED,
run_id=self._request.run_id,
model_binding=self._request.model_binding,
step=step,
)
)
except Exception:
self._event_failures += 1
logger.exception("运行 %s%d 步的事件投递失败", self._request.run_id, step.step_idx)
async def _finish(self, stop_reason: StopReason, final_answer: str | None = None) -> RunResult:
"""写结束标记,然后返回结果。
**标记在把结果交给调用方之前写下。** 让项目自己落盘的话,「跑完了、库返回了、项目存
的时候崩了」这种情况下,重启后日志显示最后一步有结果、没有结束标记,而项目那边什么都
没有——续跑会重复执行最后一步的副作用,不续跑就丢掉一次已经花完钱的运行。歧义来自
结果跨了两个存储。
"""
result = RunResult(
run_id=self._request.run_id,
stop_reason=stop_reason,
final_answer=final_answer,
steps=tuple(self._steps),
event_delivery_failures=self._event_failures,
)
await self._definition.store.write_run_finished(
RunFinished(run_id=self._request.run_id, result=result)
)
return result
async def _drain_within_grace(self, writing: "asyncio.Future[None]") -> None:
"""等一个已经发出去的写入在宽限期内落完,等不到就放弃。
取消进来时那次写入可能还在半路。不等它就写结束记录的话,两条记录可能以相反的顺序落地,
而一份「先有结束记录、后有开始记录」的日志在恢复那边是结构上说不通的状态。
"""
try:
await asyncio.wait_for(
asyncio.shield(writing), timeout=self._request.cancel_grace_seconds
)
except TimeoutError:
writing.cancel()
logger.error(
"取消宽限期内没写完在途记录,运行 %s 的日志可能不完整", self._request.run_id
)
except Exception:
logger.exception("取消时在途记录写失败,运行 %s 的日志可能不完整", self._request.run_id)
async def start(self, snapshot: Mapping[str, str]) -> None:
"""写运行开始记录。取消正好落在这一下时,屏蔽它让它写完再补结束记录。
不屏蔽的话会留下一份**只有开始记录**的日志:`run` 因为标识已存在而拒绝,`resume` 把它
当成可以从第 0 步续跑——而它其实是被人主动叫停的。屏蔽让这条记录一定落地,取消结束
记录才有地方挂:先写结束记录再写开始记录会拼出一份结构上说不通的日志。
"""
writing: asyncio.Future[None] = asyncio.ensure_future(
self._definition.store.write_run_started(
RunStarted(run_id=self._request.run_id, parameter_snapshot=snapshot)
)
)
try:
await asyncio.shield(writing)
except asyncio.CancelledError:
await self._drain_within_grace(writing)
await self._finish_cancelled()
raise
def _assemble(self) -> tuple[Message, ...]:
return assemble(
context=self._request.context,
injections=self._request.injections,
steps=self._steps,
observation_template=self._request.observation_template,
)
def seed(self, plan: ResumePlan) -> None:
"""把从日志读回来的计数与步序列装进来。"""
self._counters = RunCounters(
steps_appended=plan.steps_appended,
actions_executed=plan.actions_executed,
consecutive_parse_failures=plan.consecutive_parse_failures,
)
self._steps = list(plan.steps)
async def settle_interrupted_tail(self, plan: ResumePlan) -> RunResult | None:
"""把最后一步之后的那次停止判定重演一遍。返回 `None` 表示那一步之后确实该接着跑。
**这一步不是可选的优化,是正确性。** 停止判定的结果只存在于结束记录里,而步记录与结束
记录是两次写;崩在两者之间,那次判定就丢了。恢复照常回到预算准入的话:一次「恰好用满
预算完成」会被改写成「预算耗尽」(两者的轨迹长度一模一样,事后分不出来),一次已经
达成目标的运行会接着往下跑,一次该以连续解析失败收尾的运行会再花一次模型调用。
判定所需的东西全都在日志里:动作结果在步记录那条原子写里,工具名在步记录上,模型回复
在模型调用结果里。这里只是把它们重新读一遍。
"""
entry = plan.last_step_completed
if entry is None:
return None
step = entry.step
if not step.parse_ok and step.parse_error is None:
# 模型调用失败那一步。它一写完就该以模型调用失败收尾。
return await self._finish(StopReason.LLM_ERROR)
if not step.parse_ok:
stop = parse_failure_admission(self._counters, self._request.budget)
return None if stop is None else await self._finish(stop)
if entry.action_outcome is not None:
spec = None if step.tool_name is None else self._request.tools.spec_for(step.tool_name)
stop = completion_verdict(
entry.action_outcome, bool(spec is not None and spec.completes_run)
)
return None if stop is None else await self._finish(stop)
# 有效决策却没有动作结果,只剩最终回答那一档。答案文本只存在于结束记录里,所以重新
# 解释那条已经存下来的回复把它找回来——那时副作用还没发生,重新解释是安全的。
if plan.last_reply is None:
raise CorruptLogError("最后一步是最终回答,却找不到它那次模型调用的回复")
decision = self._definition.decision_parser.parse(plan.last_reply).decision
if not isinstance(decision, FinalAnswer):
raise CorruptLogError(
"最后一步记成了最终回答,重新解释同一条回复却得到别的分支——"
"解释器在两次运行之间被换过,或者它不是确定性的"
)
return await self._finish(StopReason.AGENT_FINISHED, final_answer=decision.text)
async def _finish_cancelled(self) -> None:
"""取消进来时尽力写下结束标记,宽限期用完就放弃。
不写的话,恢复读到的是一次没有结束标记的运行,会被当成可以续跑,而它其实是被人主动
叫停的。**宽限期用完还没写完就放弃,不无限等待**——取消的语义是尽快停下,为了留痕而
卡住违背它。
写失败只记日志不再抛:这里已经在 `CancelledError` 的处置路径上,再抛一个异常会把取消
这件事本身盖掉,而调用方的结构化并发正等着那个 `CancelledError`。
"""
result = RunResult(
run_id=self._request.run_id,
stop_reason=StopReason.CANCELLED,
final_answer=None,
steps=tuple(self._steps),
event_delivery_failures=self._event_failures,
)
await self._drain_within_grace(
asyncio.ensure_future(
self._definition.store.write_run_finished(
RunFinished(run_id=self._request.run_id, result=result)
)
)
)
# -- 主循环 -------------------------------------------------------------
async def drive(self, plan: ResumePlan | None = None) -> RunResult:
"""走 `0004` 决策三那八档,直到某一档给出停止原因。
**`CancelledError` 不捕获吞没**:接住只为写一条结束标记,写完原样重抛。取消要能穿过
模型调用与动作执行,而 `run` 在取消时**不返回结果**,为的是不破坏调用方的结构化并发
语义(`0004` 决策二)。
"""
try:
return await self._loop(plan)
except asyncio.CancelledError:
await self._finish_cancelled()
raise
async def finish_resume_unknown(self, plan: ResumePlan) -> RunResult:
"""续跑撞上「状态未知且声明绝不重放」时的收尾。
返回一个正常结果而不是抛异常:撞上未知状态是一个可预期的正常终态,不是库自身的缺陷。
抛异常会丢掉「跑到第几步、已经花了多少、前面那些步的轨迹」,而那些信息正是项目决定
「重跑还是人工介入」时要看的。
"""
self._steps = list(plan.steps)
return await self._finish(StopReason.RESUME_STATE_UNKNOWN)
async def _loop(self, plan: ResumePlan | None) -> RunResult:
pending_reply: ModelReply | None = None
pending_failure: str | None = None
# 被打断的那一步的意图已经落过盘了,重放时不能再写一条——同一步两条同种意图会被恢复
# 判定认成「日志被并发写过」,于是下一次续跑直接拒绝,一次成功的重放反倒把日志弄坏了。
skip_model_intent = plan is not None and plan.action is ResumeAction.REDO_MODEL_CALL
skip_action_intent = plan is not None and plan.action is ResumeAction.REPLAY_LAST_ACTION
if plan is not None:
self.seed(plan)
pending_reply = plan.pending_reply
pending_failure = plan.pending_failure
# 模型调用失败那一步的步记录还没写完就断了:补上它,然后以模型调用失败收尾。只可能
# 发生在续跑的第一次迭代上,所以判在循环外面。
# **规模照现在重新装配出来的算**,那和被打断时装配出来的是同一份——历史、上下文、注入
# 内容都没变,而续跑守卫已经比对过它们了。填 0 会把那一列改写成一个假值。
if pending_failure is not None:
started = time.monotonic()
chars = prompt_chars(self._assemble())
await self._write_step(self._failed_call_step(len(self._steps), chars, started), None)
return await self._finish(StopReason.LLM_ERROR)
while True:
call_index = len(self._steps)
# A 预算准入。放在开头而不是上一次迭代的结尾:一次「恰好用满预算完成」的运行走的
# 是完成判定,放结尾它会先撞上预算上限,而两者的轨迹长度一模一样。
stop = budget_admission(self._counters, self._request.budget)
if stop is not None:
return await self._finish(stop)
started = time.monotonic()
messages = self._assemble()
chars = prompt_chars(messages)
# B 规模判定。命中时不产生步记录、也不写任何意图——模型还没被调用、没花钱、没有
# 调用标识需要对账。这是唯一一种「真的一步都没走」的终止。
stop = prompt_size_admission(chars, self._request.budget)
if stop is not None:
return await self._finish(stop)
# C 写模型调用意图,调模型。
if pending_reply is not None:
reply, pending_reply = pending_reply, None
else:
reply_or_none = await self._call_model(
call_index, messages, write_intent=not skip_model_intent
)
skip_model_intent = False
if reply_or_none is None:
step = self._failed_call_step(call_index, chars, started)
await self._write_step(step, None)
return await self._finish(StopReason.LLM_ERROR)
reply = reply_or_none
# D / E 解释决策。
parsed = self._definition.decision_parser.parse(reply)
decision = parsed.decision
if isinstance(decision, InvalidDecision):
self._counters = self._counters.with_parse_failure()
await self._write_step(
self._parse_failure_step(
call_index, parsed.history_text, decision, reply, chars, started
),
None,
)
stop = parse_failure_admission(self._counters, self._request.budget)
if stop is not None:
return await self._finish(stop)
continue # 这一支跳过完成判定:这一步没碰环境,完成信号不可能因为它改变。
self._counters = self._counters.with_parse_success()
if isinstance(decision, FinalAnswer):
await self._write_step(
self._final_answer_step(
call_index, parsed.history_text, decision, reply, chars, started
),
None,
)
return await self._finish(StopReason.AGENT_FINISHED, final_answer=decision.text)
# F 写动作意图,执行动作,无条件记步。
outcome, spec_completes_run = await self._execute(
call_index, decision, write_intent=not skip_action_intent
)
skip_action_intent = False
observation, is_synthetic, truncated = _project_observation(
outcome, self._definition.synthetic_observations
)
if outcome.status is ActionStatus.EXECUTED:
self._counters = self._counters.with_action_executed()
await self._write_step(
self._action_step(
call_index,
parsed.history_text,
decision,
reply,
outcome,
observation,
is_synthetic,
truncated,
chars,
started,
),
outcome,
)
# G 完成判定。
stop = completion_verdict(outcome, spec_completes_run)
if stop is not None:
return await self._finish(stop)
# H 回 A。
async def _call_model(
self, call_index: int, messages: tuple[Message, ...], *, write_intent: bool
) -> ModelReply | None:
"""写意图、调模型、写结果。返回 `None` 表示这次调用失败了。
**失败也必须落一条结果记录。** 只写步不写结果的话,进程在写完步、还没写运行结束时
崩溃,恢复读到「意图有、结果无」会判为状态未知走重放策略,而这次调用的状态一点都不
未知——它明确地失败过。
"""
result_id = _model_result_id(self._request.run_id, call_index)
if write_intent:
await self._definition.store.write_intent(
Intent(
run_id=self._request.run_id,
kind=IntentKind.MODEL_CALL,
call_index=call_index,
result_id=result_id,
replay_policy=self._request.model_replay_policy,
)
)
try:
reply = await self._definition.model_client.call(
ModelCall(
messages=messages,
call_index=call_index,
run_id=self._request.run_id,
result_id=result_id,
binding=self._request.model_binding,
)
)
except asyncio.CancelledError:
raise
except Exception as exc: # noqa: BLE001 — 见 docstring:失败必须落一条结果记录
await self._definition.store.write_model_call_result(
ModelCallResult(
run_id=self._request.run_id,
result_id=result_id,
reply=None,
# 存说明文本不存异常对象:异常对象没法可靠地序列化成任何一种持久形态,
# 而恢复只需要知道「失败过」以及失败的大致形态。
failure=f"{type(exc).__name__}: {exc}",
)
)
return None
await self._definition.store.write_model_call_result(
ModelCallResult(
run_id=self._request.run_id, result_id=result_id, reply=reply, failure=None
)
)
return reply
async def _execute(
self, call_index: int, action: Action, *, write_intent: bool
) -> tuple[ActionOutcome, bool]:
"""写动作意图,执行动作。返回结果与「这次执行的工具被标了完成标记吗」。
重放策略从注册表查:没有工具的动作(模型输出的是一整段代码)问不出规格,取「绝不
重放」——当成可重放而其实不是会重复执行副作用且静默,反过来只是多停一次。
"""
spec = (
None
if action.tool_call is None
else self._request.tools.spec_for(action.tool_call.name)
)
if write_intent:
await self._definition.store.write_intent(
Intent(
run_id=self._request.run_id,
kind=IntentKind.ACTION,
call_index=call_index,
result_id=_action_result_id(self._request.run_id, call_index),
replay_policy=ReplayPolicy.NEVER if spec is None else spec.replay_policy,
)
)
outcome = await self._request.action_executor.execute(action)
return outcome, bool(spec is not None and spec.completes_run)
# -- 步记录 -------------------------------------------------------------
def _base_step(self, call_index: int, chars: int, started: float) -> dict[str, object]:
return {
"step_idx": call_index,
"prompt_chars": chars,
"step_wall_ms": int((time.monotonic() - started) * 1000),
}
def _failed_call_step(
self, call_index: int, prompt_chars_used: int, started: float
) -> StepRecord:
"""模型调用失败那一步。
**`parse_error` 留空**,因为这一步压根没走到解释器。恢复那边正是靠这一点把它和解析
失败分开:解析失败必定带着回喂给模型的说明,它没有。
"""
return StepRecord(
**self._base_step(call_index, prompt_chars_used, started), # type: ignore[arg-type]
raw_output="",
content_chars=0,
thinking_chars=0,
action=None,
parse_ok=False,
parse_error=None,
observation=self._definition.synthetic_observations.model_call_failed,
observation_is_synthetic=True,
observation_truncated_chars=0,
call_id=None,
)
def _parse_failure_step(
self,
call_index: int,
history_text: str,
decision: InvalidDecision,
reply: ModelReply,
chars: int,
started: float,
) -> StepRecord:
return StepRecord(
**self._base_step(call_index, chars, started), # type: ignore[arg-type]
raw_output=history_text,
content_chars=len(reply.content),
thinking_chars=len(reply.thinking),
action=None,
parse_ok=False,
# 这段说明**就是**回喂给模型的那段观察,不是从一个固定串里取。压成一句会改掉模型
# 收到的纠错信息,它的纠错行为也就跟着变。
parse_error=decision.explanation,
observation=decision.explanation,
observation_is_synthetic=True,
observation_truncated_chars=0,
call_id=reply.call_id,
)
def _final_answer_step(
self,
call_index: int,
history_text: str,
decision: FinalAnswer,
reply: ModelReply,
chars: int,
started: float,
) -> StepRecord:
return StepRecord(
**self._base_step(call_index, chars, started), # type: ignore[arg-type]
raw_output=history_text,
content_chars=len(reply.content),
thinking_chars=len(reply.thinking),
action=None,
parse_ok=True,
parse_error=None,
# 这一支不碰环境,所以没有观察。
observation="",
observation_is_synthetic=False,
observation_truncated_chars=0,
call_id=reply.call_id,
)
def _action_step(
self,
call_index: int,
history_text: str,
action: Action,
reply: ModelReply,
outcome: ActionOutcome,
observation: str,
is_synthetic: bool,
truncated: int,
chars: int,
started: float,
) -> StepRecord:
return StepRecord(
**self._base_step(call_index, chars, started), # type: ignore[arg-type]
raw_output=history_text,
content_chars=len(reply.content),
thinking_chars=len(reply.thinking),
action=action.text,
parse_ok=True,
parse_error=None,
observation=observation,
observation_is_synthetic=is_synthetic,
observation_truncated_chars=truncated,
call_id=reply.call_id,
tool_name=None if action.tool_call is None else action.tool_call.name,
tool_arguments=(
None
if action.tool_call is None
else _serialise_arguments(action.tool_call.arguments)
),
action_status=outcome.status,
env_reported_completion=outcome.env_reported_completion,
)
async def run(definition: AgentDefinition, request: RunRequest) -> RunResult:
"""从头跑一次运行。
**这个运行标识已经有日志时直接报错**,不覆盖也不接着跑。覆盖会毁掉一次已经花完钱的运行
的留痕,接着跑是 `resume` 的事而两者的失败方式不同。
取消时**不返回结果**`CancelledError` 原样重抛——为的是不破坏调用方的结构化并发语义。
所以 `StopReason.CANCELLED` 只出现在两个地方:写进日志的那条结束记录里,以及恢复一个已
取消运行时重建出来的结果里。读代码的人容易写出「如果结果的停止原因是取消」这种永假分支。
"""
log = await definition.store.read_log(request.run_id)
if log != RunLog():
raise RunIdentityError(
f"运行标识 {request.run_id!r} 已经有日志了。要接着跑用 resume;"
"这里不覆盖,因为那会毁掉一次已经花完钱的运行的留痕"
)
driver = _Driver(definition, request)
await driver.start(_merged_snapshot(definition, request))
return await driver.drive()
async def resume(definition: AgentDefinition, request: RunRequest) -> RunResult:
"""接着跑一次被打断的运行。
**读不到日志直接报错**:那说明这个标识对应的运行从没开始过,而 `resume` 的语义是「同一次
运行接着做」。
比对参数快照,任何一项不一致直接报错。**这意味着续跑不能顺便改预算或换模型**——那不是
限制而是这条守卫的全部意义。
"""
log = await definition.store.read_log(request.run_id)
if log == RunLog():
raise RunIdentityError(
f"运行标识 {request.run_id!r} 读不到任何日志,这次运行从没开始过。要从头跑用 run"
)
plan = plan_resume(log, request.model_replay_policy)
if plan.action is ResumeAction.START_FRESH:
raise CorruptLogError("日志里有记录却判成全新运行,这个状态不该出现")
current = _merged_snapshot(definition, request)
stored = dict(log.started.parameter_snapshot) if log.started is not None else {}
drifted = sorted(
key for key in set(current) | set(stored) if current.get(key) != stored.get(key)
)
if drifted:
raise ParameterDriftError(
f"续跑的装配与日志里那份对不上:{drifted}。"
"前几步与后几步来自两个不同的配置,而两段轨迹在文件里看起来是同一次运行"
)
if plan.action is ResumeAction.ALREADY_FINISHED:
# 这次运行早就跑完了,把存下来的结果原样交回去,不重跑最后一步。
if plan.finished_result is None:
raise CorruptLogError("有结束标记却读不出结果,这条记录坏了")
return plan.finished_result
driver = _Driver(definition, request)
if plan.action is ResumeAction.STOP_UNKNOWN:
return await driver.finish_resume_unknown(plan)
if plan.action is ResumeAction.CONTINUE_AT_NEXT_STEP:
# 上一步是完整的,但它之后那次停止判定的结果只存在于结束记录里,而那条记录没写下来。
# 先把它重演一遍——不重演的话,一次已经达成目标的运行会被改写成预算耗尽或者接着往下跑。
driver.seed(plan)
settled = await driver.settle_interrupted_tail(plan)
if settled is not None:
return settled
return await driver.drive(plan)
__all__ = [
"AgentDefinition",
"ParameterDriftError",
"RunIdentityError",
"RunRequest",
"resume",
"run",
]