"""压测的故障注入:在崩溃、取消、撞预算、解析连击、提示词超限、环境故障这几条路径上验库有 没有守住承诺。 一百个任务顺利跑完什么都证明不了——顺利那条路上库只要不崩就算过。能证明东西的是这几条: 进程被杀在半路、调用方中途取消、目标根本完不成、模型的输出一句都解析不了、提示词撑到装不下、 环境跑到一半不能接着服务了。库在这些地方给下游的承诺(崩溃前的轨迹逐字节不变、声明绝不重放 的动作不会执行两次、取消原样穿透且资源归还、撞上限时以正确的停止原因干净收尾、解析失败不碰 环境、提示词超限时终止而不是静默截断历史、环境坏掉那一步的观察换成合成观察)只有在这里才 检验得到。 **九类里最后添的两类是在补一次全量跑留下的空白。** 193 个运行跑完之后,`polyloop.types. StopReason` 的十个取值里有两个一次都没出现过:`context_overflow` 与 `env_error`。没出现不等于 它们是对的,只等于没验过——那两条路上的承诺到那一刻为止一次都没有被检验,而它们错了的形态 恰好都不会当场炸:静默截断历史在数据里看起来跟「模型不行」一模一样,环境故障那一步的观察 没被替换掉则表现成模型收到一段来路不明的文本。所以这两类各自把那个停止原因造出来一次: GovDoc 那一路把 `max_prompt_chars` 压到刚好等于装配出来的初始提示词,AppWorld 那一路在跑完 一步之后真的把容器 `docker kill` 掉。 **崩溃续跑用 GovDoc 场景,不用 AppWorld。** GovDoc 的「环境」是一个工作区目录加一份审计 日志(`tools/soak/scenarios/govdoc.py` 里 `GovDocTools._append_audit` 写的那份),它天然活过 进程的死亡,而且给得出**环境侧自己数的实际执行次数**——那是「声明绝不重放的动作有没有被 执行两次」唯一可信的证据。AppWorld 的环境活在容器里:子进程被 SIGKILL 之后容器的归属与 清理会变复杂,而它的执行计数存在会话对象里,会随进程一起消失,于是只剩「库自己报的步数」 可用,那等于我们和我们自己对账。 **判据一律不依赖模型的确定性。** 活模型下「两次运行产出同样的轨迹」本来就不成立,断言它 只会随机变红。这里的判据全是结构不变量:字节前缀、步号连续性、意图有没有归宿、审计账去重 前后的条数、停止原因的取值、步数与上限的关系、环境侧的执行计数。 **崩溃是子进程自己在指定的写入点上 `os._exit`,不是父进程从外面抢窗口发信号。** 第一版靠 父进程轮询日志尾部再 `SIGKILL`,实测三次全都没能命中「一步完整落地之后」那个时机——库写完 步记录紧接着就写下一步的模型调用意图,中间只有内存里的装配计算,窗口窄到信号挤不进去。 **那次实测本身是一条要记住的事实**:自然发生的崩溃几乎总是落在「有意图没结果」那一态上, 而不是落在两步之间的干净边界上。但要验那个边界就不能靠碰运气。做法见 `SelfKillingStore`; 外部 SIGKILL 那条路留着兜底,没有删。 **产物照 `tools/soak/scoreboard.py` 的 sidecar 约定写出**(`.result.json` / `.events.jsonl` / `.meta.json`),`meta.fault` 填故障名,这样记分板能把它们和正常批次一起判。崩溃那两类天然 缺 sidecar——被 SIGKILL 的子进程来不及写,记分板会报「无法判定」,那是对的,不为了让它变绿 去补假数据。 ## 已知缺口 这两条是这套注入器**造不出来**的情形。写在这里是因为「没被验到」与「验过了没问题」在报告上 长得一样,而报告只列跑过的判据,列不出没跑的那些。 **一、崩溃形态只有「写入调用返回之后」这一种。** `os._exit` 是一次系统调用,它只能发生在两条 Python 语句之间;进程内的自杀不可能停在 `write(2)` 的中途。所以「一行只写了半截就断电」 「记录进了页缓存但 `fsync` 之前机器崩了」这两种形态在这里造不出来——真要造得靠外部手段(掐 电、`dm-flakey` 这类会掉写的块设备、或者在文件系统层注入故障),那是另一套东西。现在靠的是 记分板那侧的撕裂行判据:它按「有没有被换行终结」切行,末尾那段没终结的字节按契约算没发生过 (`tools/soak/scoreboard.py` 的 `_split_terminated`)。也就是说这种日志**读得对**是验过的, **写出来会怎样**没验过。 **二、取消之后容器里那次执行的下落看不见。** `AppWorldSession.n_executions` 是客户端侧的 计数,HTTP 响应返回之后才加一,而取消把那个协程掐断了。于是「代码已经发到容器、容器仍然把它 跑完并改了环境状态」这种情形留不下痕迹。容器那侧没有可查的执行计数接口(只有 `/execute`、 `/task_completed`、`/evaluate`、`/close`),补上它要改 `tools/soak/appworld.py` 那一层。 `check_env_quiet_after_cancel` 覆盖的是另外两种坏法:库在取消之后还往环境派活,或者有一次 执行在取消之后才落地。 跑法(会打真实模型、会花钱):: PYTHONUNBUFFERED=1 conda run --live-stream -n PolyLoop python -m tools.soak.faults \\ --runs-dir <目录> --budget-calls 60 \\ --data-root --govdoc-db --govdoc-corpus <语料目录> """ from __future__ import annotations import argparse import asyncio import contextlib import json import os import sys import time from collections.abc import Awaitable, Callable, Mapping, Sequence from dataclasses import dataclass, replace from enum import StrEnum from pathlib import Path import httpx from polyloop import session from polyloop.adapters import GatewayModelClient from polyloop.ports import Event, InvalidDecision, ParsedReply, RunLog, RunStore from polyloop.serialization import encode from polyloop.session import AgentDefinition, ParameterDriftError, RunRequest from polyloop.stores import RECORD_KEY, JsonlRunStore from polyloop.types import ( ActionOutcome, ActionStatus, Budget, Intent, ModelCallResult, ModelReply, ReplayPolicy, RunFinished, RunResult, RunStarted, StepCompleted, TextBlock, ) from tools.soak.appworld import CONTAINER_NAME_PREFIX, DEFAULT_PORT_BASE, AppWorldPool from tools.soak.scenarios import appworld as appworld_scenario from tools.soak.scenarios import govdoc as govdoc_scenario from tools.soak.scenarios.appworld import AppWorldParser from tools.soak.scenarios.govdoc import ( AuditTask, GovDocParser, Redactor, load_checkpoints, load_documents, read_audit_lines, ) #: 仓库根目录。子进程要在这里启动,否则 `tools.soak` 这个命名空间包 import 不到。 REPO_ROOT = Path(__file__).resolve().parents[2] FAULT_NAMES: tuple[str, ...] = ( "crash_resume_a", "crash_resume_b", "cancel_model", "cancel_env", "step_budget", "action_budget", "parse_failures", "context_overflow", "env_error", ) #: 用 GovDoc 场景做的那三类。其余用 AppWorld。 GOVDOC_FAULTS = frozenset({"crash_resume_a", "crash_resume_b", "context_overflow"}) APPWORLD_FAULTS = frozenset(FAULT_NAMES) - GOVDOC_FAULTS class FaultInjectionError(RuntimeError): """故障注入这一层装配不出来或跑不下去:缺数据、缺参数、子进程起不来。 它不表示某条判据被击穿——那种事由 `Criterion` 记录,不抛异常。 """ # --------------------------------------------------------------------------- # 一、判据与报告 # --------------------------------------------------------------------------- class CriterionStatus(StrEnum): """一条判据的判定。三档,不是两档。 第三档独立存在,不许折算成前两档中的任何一个。判不了和判过了是两回事——压成一档的话, 一次什么都没验成的跑会显示成全绿,而那正是最需要被看见的情况。理由与 `tools/soak/scoreboard.py` 的三档判定相同。 """ PASSED = "passed" BREACHED = "breached" UNDETERMINED = "undetermined" @dataclass(frozen=True, slots=True, kw_only=True) class Criterion: """一条判据的判定加它的证据。 **证据是必填的**,由构造期校验守住:一条只会说「有问题」或者只会说「没问题」的判据等于 没有判据。通过也要给证据,因为「通过」最常见的坏法是判据根本没跑到该判的东西上——比了 零个字节、数了零条审计、看了一个空列表,那些情况下证据文本会当场露馅。 """ name: str status: CriterionStatus evidence: str def __post_init__(self) -> None: if not self.name.strip(): raise ValueError("判据必须有名字") if not self.evidence.strip(): raise ValueError(f"判据 {self.name!r} 没有给证据") def describe(self) -> str: return f"[{self.status.value}] {self.name}:{self.evidence}" def passed(name: str, evidence: str) -> Criterion: return Criterion(name=name, status=CriterionStatus.PASSED, evidence=evidence) def breached(name: str, evidence: str) -> Criterion: return Criterion(name=name, status=CriterionStatus.BREACHED, evidence=evidence) def undetermined(name: str, evidence: str) -> Criterion: return Criterion(name=name, status=CriterionStatus.UNDETERMINED, evidence=evidence) @dataclass(frozen=True, slots=True, kw_only=True) class FaultReport: """一类故障跑完之后的全部结论。""" fault: str criteria: tuple[Criterion, ...] #: 不构成判定、但必须被人看见的事实:重试了几次、撕裂尾行、跳过的原因。 notes: tuple[str, ...] = () @property def breaches(self) -> tuple[Criterion, ...]: return tuple(item for item in self.criteria if item.status is CriterionStatus.BREACHED) @property def undetermineds(self) -> tuple[Criterion, ...]: return tuple(item for item in self.criteria if item.status is CriterionStatus.UNDETERMINED) def render(self) -> str: lines = [f"## {self.fault}"] lines += [f"- {item.describe()}" for item in self.criteria] lines += [f"- (note) {note}" for note in self.notes] return "\n".join(lines) # --------------------------------------------------------------------------- # 二、日志读取:逐行读,撕裂尾行按「没发生过」算 # --------------------------------------------------------------------------- @dataclass(frozen=True, slots=True, kw_only=True) class LogRead: """一份 `.jsonl` 日志按行读回来的结果。 `torn` 说的是末尾有没有一段没被换行终结的字节。判据照存储那边的规矩:**看有没有被换行 终结,不看能不能解析**(`polyloop.stores` 那个逐行追加的实现,理由在 `research-wiki/design/0011-jsonl-run-store.md`)。那段字节对应的那次写从来没有被确认过, 按契约它就是没发生。 """ payloads: tuple[Mapping[str, object], ...] torn: bool #: 已经被换行终结、却解不出 JSON 对象的行号。它们是真问题,不是撕裂尾行。 bad_lines: tuple[int, ...] def terminated_prefix(raw: bytes) -> bytes: """截到最后一个换行为止的那一段。没有换行就是空。 崩溃快照与最终日志的字节比对只能比这一段:末尾那段没有换行的字节没被确认过,而续跑会 直接追加在它后面(`JsonlRunStore` 用 `O_APPEND`,不回退文件指针),于是那一行在最终 日志里长得和崩溃时不一样——那不是承诺被破坏,是那次写从来没算数。 """ cut = raw.rfind(b"\n") return b"" if cut < 0 else raw[: cut + 1] def parse_terminated(raw: bytes) -> LogRead: """把日志字节切成一条条记录载荷。""" chunks = raw.split(b"\n") torn = bool(chunks) and bool(chunks[-1].strip()) body = chunks[:-1] if torn else chunks payloads: list[Mapping[str, object]] = [] bad: list[int] = [] for number, chunk in enumerate(body, start=1): if not chunk.strip(): continue try: value = json.loads(chunk) except (UnicodeDecodeError, json.JSONDecodeError): bad.append(number) continue if not isinstance(value, dict) or RECORD_KEY not in value: bad.append(number) continue payloads.append(value) return LogRead(payloads=tuple(payloads), torn=torn, bad_lines=tuple(bad)) def read_log(path: Path) -> LogRead: if not path.is_file(): return LogRead(payloads=(), torn=False, bad_lines=()) return parse_terminated(path.read_bytes()) def tagged(read: LogRead, tag: str) -> tuple[Mapping[str, object], ...]: return tuple(item for item in read.payloads if item.get(RECORD_KEY) == tag) def count_model_calls(read: LogRead) -> int: """这份日志里发起过几次模型调用。 数的是模型调用意图,不是步数:意图在调用发出去**之前**落盘,所以被杀在调用中途的那次 也算,而它已经花过钱了。用步数会把那次漏掉。 """ return sum(1 for item in tagged(read, "intent") if item.get("kind") == "model_call") def stop_reason_of(read: LogRead) -> str | None: """日志里结束记录带的停止原因。没有结束记录就是 None。""" finished = tagged(read, "run_finished") if not finished: return None result = finished[-1].get("result") if not isinstance(result, Mapping): return None value = result.get("stop_reason") return value if isinstance(value, str) else None def _step_payloads(read: LogRead) -> tuple[Mapping[str, object], ...]: steps: list[Mapping[str, object]] = [] for item in tagged(read, "step_completed"): step = item.get("step") if isinstance(step, Mapping): steps.append(step) return tuple(steps) def parameter_snapshot_of(read: LogRead) -> Mapping[str, object] | None: """运行开始那条记录带的参数快照。没有那条记录就是 None。 判据要从这里取这次运行实际用的参数,而不是从本文件的常量取:常量是我们**打算**用的值, 快照是库**真的**收到的值,两者不一致时该被判据看见的正是后者。 """ started = tagged(read, "run_started") if not started: return None snapshot = started[0].get("parameter_snapshot") return snapshot if isinstance(snapshot, Mapping) else None # --------------------------------------------------------------------------- # 三、判据:崩溃续跑 # --------------------------------------------------------------------------- def check_log_readable(read: LogRead) -> Criterion: """每一条被换行终结的行都解得出一个带类型标签的 JSON 对象。 解不出来的行让下面几条判据判的东西缺一块,所以它先判——不先判的话,一份少了半截的日志 会让「步号连续」这类判据在残缺数据上给出「通过」。 **一条记录都没有时报「无法判定」,不报通过。** 日志文件根本不存在时 `read_log()` 返回的 就是这种空读数,而「零条全都解得开」是一句真空成立的话:它读起来像验过了,实际那次跑 连一个字节都没落盘,后面每一条判据都建立在一份不存在的日志上。 """ if read.bad_lines: return breached( "log_lines_readable", f"第 {list(read.bad_lines)} 行已被换行终结却解不出带 {RECORD_KEY!r} 标签的 JSON 对象", ) if not read.payloads: return undetermined( "log_lines_readable", "日志里一条记录都没有(文件不存在,或者一个字节都没落盘),行的完整性无从判起", ) return passed("log_lines_readable", f"{len(read.payloads)} 条记录全部解得开") def check_crash_prefix_preserved(*, crashed: bytes, final: bytes) -> Criterion: """崩溃前已经落地的那些字节,续跑之后逐字节没变。 **比的是字节串,不是解析出来的对象。** 比对象只能说明「语义等价」,而承诺是更硬的那一 条:已经写下去的记录不许被重写、不许被改格式、不许被补字段。这两者的差别只有在字节层面 看得见。 """ prefix = terminated_prefix(crashed) if not prefix: return undetermined( "crash_prefix_preserved", "崩溃快照里没有任何被换行终结的记录,前缀比对无从谈起", ) if len(final) < len(prefix): return breached( "crash_prefix_preserved", f"最终日志只有 {len(final)} 字节,比崩溃时已终结的 {len(prefix)} 字节还短", ) head = final[: len(prefix)] if head != prefix: offset = next( ( index for index, pair in enumerate(zip(prefix, head, strict=True)) if pair[0] != pair[1] ), len(prefix), ) return breached( "crash_prefix_preserved", f"崩溃时已终结的 {len(prefix)} 字节里,第 {offset} 字节起与最终日志不同", ) return passed( "crash_prefix_preserved", f"崩溃时已终结的 {len(prefix)} 字节在最终日志里逐字节相同", ) def check_step_indices_dense(read: LogRead) -> Criterion: """`step_completed` 的步号从 0 开始、逐 1 递增、不重不跳。 重号意味着同一步被写了两遍(续跑把已经落地的一步又跑了一次),跳号意味着中间某一步的 原子写整个丢了。两者都会让重建出来的历史与不中断跑完时不一样,而那件事在轨迹里看不出来。 """ indices: list[object] = [step.get("step_idx") for step in _step_payloads(read)] if not indices: return undetermined("step_idx_dense", "日志里一条步记录都没有,步号连续性无从判起") bad = [value for value in indices if not isinstance(value, int) or isinstance(value, bool)] if bad: return breached("step_idx_dense", f"有 {len(bad)} 条步记录的 step_idx 不是整数") expected = list(range(len(indices))) if indices != expected: return breached("step_idx_dense", f"步号按文件顺序是 {indices},期望 {expected}") return passed("step_idx_dense", f"{len(indices)} 条步记录的步号是 0 到 {len(indices) - 1}") def check_intents_settled(read: LogRead) -> Criterion: """每条意图都有归宿;至多一条悬空,且必须是最后一条意图。 模型调用意图的归宿是一条同 `result_id` 的模型调用结果,动作意图的归宿是一条同 `result_id` 的步记录。**允许最后一条悬空**:崩溃点上那条意图写了、结果没写,续跑判成状态未知干净 停下之后它就永远悬在那儿,那是合法终态。中间悬空则不同——它说明有一步的执行状态被跳过 去了,而后面的步是建立在「那一步到底做没做」这个没有答案的问题上的。 """ intents = tagged(read, "intent") if not intents: return undetermined("intents_settled", "日志里一条意图都没有,归宿无从判起") model_results = {item.get("result_id") for item in tagged(read, "model_call_result")} action_results = {item.get("result_id") for item in tagged(read, "step_completed")} dangling: list[int] = [] for position, intent in enumerate(intents): settled = ( intent.get("result_id") in model_results if intent.get("kind") == "model_call" else intent.get("result_id") in action_results ) if not settled: dangling.append(position) if not dangling: return passed("intents_settled", f"{len(intents)} 条意图全部有归宿,没有悬空") last = len(intents) - 1 if dangling == [last]: kind = intents[last].get("kind") return passed( "intents_settled", f"{len(intents)} 条意图里只有最后一条(kind={kind})悬空,那是崩溃点,合法", ) return breached( "intents_settled", f"{len(intents)} 条意图里第 {dangling} 条悬空(按意图出现的次序计)," f"只有第 {last} 条允许悬空", ) def parse_audit_line(line: str) -> tuple[str, str, str] | None: """把审计账的一行切成 `(工具名, 文件名, 内容摘要)`。切不出来返回 None。 格式来自 `GovDocTools._append_audit`:三段用制表符分隔。 """ parts = line.split("\t") if len(parts) != 3 or not all(part for part in parts): return None return parts[0], parts[1], parts[2] def check_never_action_not_replayed(audit_lines: Sequence[str]) -> Criterion: """声明绝不重放的动作没有被执行两次。 **证据取自环境侧自己记的账**(工作区里的 `.write_audit.log`),不取库报的步数或动作数—— 后者是库对自己行为的陈述,用它来验库的行为就是我们和我们自己对账。真正的重放长这样: 库在状态未知时把一次已经落过盘的写又执行了一遍,于是环境的账上多出一条一模一样的记录, 而轨迹里看不出任何异常。 判据是「按 (工具, 文件名, 内容摘要) 去重前后的条数相等」。**它有一种已知的假阳性**: 模型自己把同一份内容原样写了两次,账上也会出现两条一样的记录。GovDoc 的 execute 阶段 提示词明确要求写完 evidence.md 就停下,所以这件事很少发生;真发生了要看的是轨迹里那两 步的步号——重放来自续跑,两条记录会分属崩溃前后。 """ entries = [parse_audit_line(line) for line in audit_lines] broken = [index for index, entry in enumerate(entries, start=1) if entry is None] if broken: return undetermined( "never_action_not_replayed", f"审计账第 {broken} 行不是「工具\\t文件名\\t摘要」三段,数不出实际执行次数", ) if not entries: return undetermined( "never_action_not_replayed", "审计账是空的:这次运行里没有任何有副作用的动作被执行过,去重比对无从谈起", ) kept = [entry for entry in entries if entry is not None] unique = set(kept) if len(unique) == len(kept): return passed( "never_action_not_replayed", f"审计账 {len(kept)} 条,按 (工具, 文件名, 内容摘要) 去重后仍是 {len(unique)} 条", ) counts: dict[tuple[str, str, str], int] = {} for entry in kept: counts[entry] = counts.get(entry, 0) + 1 repeated = sorted( ( f"{tool} 摘要 {digest} 出现 {times} 次" for (tool, _name, digest), times in counts.items() if times > 1 ) ) return breached( "never_action_not_replayed", f"审计账 {len(kept)} 条,去重后只剩 {len(unique)} 条:{';'.join(repeated)}", ) def check_audit_unchanged(*, before: Sequence[str], after: Sequence[str]) -> Criterion: """续跑之后环境的账一条都没多。 时机 B 下库判定状态未知、干净停下,那就不该把被打断的那个动作再执行一遍。多出一条就是 重放,少一条或者变了内容说明账被改写过。 **两边都是空的时候报「无法判定」。** 那代表续跑前后都没有任何副作用被记录过,「一条都没多」 就成了一句真空成立的话——而这一档要验的恰恰是「已经落过盘的那次副作用没有再发生一次」, 没有那次副作用就没有要验的东西。上一版实跑里这一条正是 0 比 0 通过的。 """ if not before and not after: return undetermined( "audit_unchanged_after_resume", "续跑前后审计账都是空的:这次运行里没有任何副作用被记录过,「没多出来」真空成立", ) if list(after[: len(before)]) != list(before): return breached( "audit_unchanged_after_resume", f"续跑前 {len(before)} 条审计记录在续跑后不再是原来那几条", ) if len(after) != len(before): return breached( "audit_unchanged_after_resume", f"续跑前审计账 {len(before)} 条,续跑后 {len(after)} 条,多出 {len(after) - len(before)} 条", ) return passed("audit_unchanged_after_resume", f"续跑前后审计账都是 {len(before)} 条") def check_audit_prefix_preserved(*, before: Sequence[str], after: Sequence[str]) -> Criterion: """崩溃时环境账上那几条,续跑之后逐字节还在原处。 这是日志那侧 `check_crash_prefix_preserved` 在**环境账**上的对应物,而且两者不可互相 替代:日志是库自己写的,账是环境写的。库把已经落地的一段轨迹重写掉,日志那条会报;库让 环境把已经发生过的副作用重做或抹掉,只有这条会报。 **时机 A 之前没有这条判据。** 那一档只跑「去重前后条数相等」,而一次把 `evidence.md` 的 旧条目改写成别的内容的续跑,条数不变、去重也不变,一路绿灯。 """ name = "audit_prefix_preserved" if not before: return undetermined(name, "崩溃时审计账是空的:没有已经落地的副作用可比,前缀比对无从谈起") if len(after) < len(before): return breached( name, f"崩溃时审计账 {len(before)} 条,续跑之后只剩 {len(after)} 条,已经落地的记录被抹掉了", ) for index, (old, new) in enumerate(zip(before, after, strict=False)): if old != new: return breached( name, f"崩溃时审计账的第 {index} 条在续跑之后变了内容(按行序计,正文不进报告)" ) return passed(name, f"崩溃时审计账的 {len(before)} 条在续跑之后逐字节还在原处") def check_no_cross_segment_replay(*, before: Sequence[str], after: Sequence[str]) -> Criterion: """崩溃前已经落地的那次副作用,续跑之后没有再发生一次。 与 `check_never_action_not_replayed` 的区别在**位置**:那一条在整份账上找重复,找到了也 说不出两条分别在崩溃的哪一侧;这一条只看续跑段里有没有哪一条与崩溃前的某一条完全相同, 而跨越崩溃边界的重复正是重放的签名——库在状态未知时把一次已经落过盘的写又执行了一遍。 **它与那一条共享同一个已知假阳性**:模型自己在续跑之后把同一份内容原样又写了一次,看起来 一模一样。GovDoc 的 execute 阶段提示词要求写完就停,所以很少发生;真发生了要看轨迹里那两步 的步号。这个代价是被接受的,理由与那一条相同。 """ name = "no_cross_segment_replay" if not before: return undetermined( name, "崩溃时审计账是空的:没有已经落地的副作用可能被重放,这一条无从判起" ) resumed = list(after[len(before) :]) if not resumed: return passed(name, f"续跑段里一条审计记录都没有,崩溃前那 {len(before)} 条不可能被重放过") known = {parse_audit_line(line) for line in before} repeats: list[str] = [] for line in resumed: entry = parse_audit_line(line) if entry is not None and entry in known: repeats.append(f"{entry[0]} 摘要 {entry[2]}") if repeats: return breached( name, f"续跑段的 {len(resumed)} 条记录里有 {len(repeats)} 条与崩溃前的某一条完全相同:" f"{';'.join(sorted(set(repeats)))}", ) return passed( name, f"续跑段新增 {len(resumed)} 条审计记录,没有一条与崩溃前那 {len(before)} 条中的任何一条相同", ) def check_resume_made_progress(*, crashed_steps: int, final_steps: int) -> Criterion: """时机 A 下续跑真的接着往下跑了。 一步完整落地之后崩,日志处在「上一步是完整的」这个状态,续跑该从下一步的开头接着走。 步数没长说明它没接着跑——那种「续跑」只是把旧结果读回来又写了一遍结束记录。 """ if final_steps > crashed_steps: return passed( "resume_made_progress", f"崩溃时 {crashed_steps} 步,续跑之后 {final_steps} 步", ) return breached( "resume_made_progress", f"崩溃时 {crashed_steps} 步,续跑之后仍是 {final_steps} 步,没有接着往下跑", ) # --------------------------------------------------------------------------- # 四、判据:停止原因、步数、取消、环境、提示词规模 # --------------------------------------------------------------------------- #: 提示词上限在参数快照里的键名。快照由 `RunRequest.parameter_snapshot()` 拼出来,值是十进制 #: 字符串。判据从快照读它、不硬编码:这一路的上限是按上下文现算出来的,写死一份的话判据比的 #: 就是另一个数字,而两边不一致时它照样给得出「通过」。 MAX_PROMPT_CHARS_SNAPSHOT_KEY = "request.max_prompt_chars" #: 观察模板在参数快照里的键名。算「再走一步的提示词会有多大」要用它。 OBSERVATION_TEMPLATE_SNAPSHOT_KEY = "request.observation_template" def check_stop_reason(read: LogRead, expected: str) -> Criterion: """结束记录在场,且它带的停止原因是期望的那个。 判据取自**日志里的结束记录**而不是 `run` 的返回值:取消那条路径上 `run` 原样重抛 `CancelledError`、根本不返回结果,返回值那侧什么都拿不到,而结束记录仍然在。 """ actual = stop_reason_of(read) if actual is None: return breached( f"stop_reason_is_{expected}", "日志里没有 run_finished 记录,读不到停止原因" ) if actual != expected: return breached(f"stop_reason_is_{expected}", f"停止原因是 {actual},期望 {expected}") return passed(f"stop_reason_is_{expected}", f"日志的 run_finished 里停止原因是 {actual}") def check_run_finished_present(read: LogRead) -> Criterion: """日志里有结束记录。 它由库在把结果交给调用方**之前**写下,所以它在不在与停止原因取值对不对是两件事:没有它 的话,恢复会把这次运行读成一次可以续跑的运行,而它其实已经结束了。分成两条判据是因为两者 的成因不同——一条查的是那次写有没有发生,另一条查的是写下去的取值对不对。 """ finished = tagged(read, "run_finished") if not finished: return breached("run_finished_recorded", "日志里没有 run_finished 记录") return passed("run_finished_recorded", f"日志里有 {len(finished)} 条 run_finished 记录") def check_step_count(read: LogRead, *, expected: int, name: str) -> Criterion: actual = len(_step_payloads(read)) if actual != expected: return breached(name, f"日志里 {actual} 条步记录,期望正好 {expected} 条") return passed(name, f"日志里正好 {expected} 条步记录") def check_executed_action_count(read: LogRead, *, expected: int) -> Criterion: """真正执行成功的动作数正好等于动作上限。 数的是步记录里 `action_outcome.status == "executed"` 的条数,与 `_stopping` 那侧 `actions_executed` 的口径一致:被拒绝与环境故障都不计入动作预算。 """ executed = 0 for item in tagged(read, "step_completed"): outcome = item.get("action_outcome") if isinstance(outcome, Mapping) and outcome.get("status") == "executed": executed += 1 if executed != expected: return breached( "executed_actions_equal_max_actions", f"执行成功的动作 {executed} 次,期望正好 {expected} 次", ) return passed("executed_actions_equal_max_actions", f"执行成功的动作正好 {expected} 次") def check_no_env_error_step(read: LogRead) -> Criterion: """轨迹里没有环境故障的步。 撞预算那两条要的是「干净地撞上限」。中途出过环境故障的话,步数与动作数的账仍然对得上, 但这次跑压到的已经不是预算这条路径了。 **一步都没产出时报「无法判定」**:「零步里没有 env_error」这句话恒为真,而它出现在报告上 的样子和「跑了十步都没出故障」一模一样。 """ entries = tagged(read, "step_completed") if not entries: return undetermined( "no_env_error_step", "日志里一条步记录都没有,有没有环境故障的步无从判起" ) bad: list[int] = [] for index, item in enumerate(entries): outcome = item.get("action_outcome") if isinstance(outcome, Mapping) and outcome.get("status") == "env_error": bad.append(index) if bad: return breached("no_env_error_step", f"第 {bad} 步(按步记录次序计)的动作状态是 env_error") return passed("no_env_error_step", f"{len(entries)} 步里没有 env_error") def check_at_least_one_step(read: LogRead) -> Criterion: """至少完整走过一步。 撞提示词上限那一类要靠它兜底:一步都没走的话,撞上的是「上下文本身就比上限大」,而那个 长度是装配决定的、不是循环一步步撑出来的,这一类什么都没验到。报「无法判定」而不是 「击穿」——那种情况下库做的事仍然是对的,错的是这一路的上限选得太小。 """ steps = _step_payloads(read) if not steps: return undetermined( "at_least_one_step", "日志里一条步记录都没有:提示词在第一次装配时就超了上限,循环一步都没走过", ) return passed("at_least_one_step", f"日志里 {len(steps)} 条步记录") def check_prompt_chars_monotonic(read: LogRead) -> Criterion: """`prompt_chars` 在整个运行里单调不减:历史一次都没有被截断过。 静默截断正是这一档要防住的坏法。它把提示词的前缀改掉,于是后面每一步都重新全价计费,而 在数据里看起来跟「模型不行」一模一样。历史只增不减时这一列只会往上走,截断会在它上面留下 一次下降——那是这件事在轨迹里唯一看得见的痕迹。 """ steps = _step_payloads(read) values = [step.get("prompt_chars") for step in steps] if not values: return undetermined("prompt_chars_never_shrinks", "日志里一条步记录都没有,单调性无从判起") bad = [value for value in values if not isinstance(value, int) or isinstance(value, bool)] if bad: return breached( "prompt_chars_never_shrinks", f"有 {len(bad)} 条步记录的 prompt_chars 不是整数" ) drops = [ f"第 {index} 步从 {previous} 掉到 {current}" for index, (previous, current) in enumerate(zip(values, values[1:], strict=False), start=1) if current < previous # type: ignore[operator] ] if drops: return breached( "prompt_chars_never_shrinks", f"提示词字符数出现下降(按步记录次序计):{';'.join(drops)},历史被截断过", ) return passed( "prompt_chars_never_shrinks", f"{len(values)} 步的提示词字符数从 {values[0]} 一路不减到 {values[-1]}", ) def check_prompt_reached_max_prompt_chars(read: LogRead) -> Criterion: """这次运行真的撑到了提示词上限那条线上。 **判的不是「最后一条步记录的 `prompt_chars` 超过了上限」。** 那条在一个正确的库上永远不 成立,断言它等于断言契约的反面:规模判定在调模型之前做,命中时不产生步记录也不写任何意图 (`polyloop._stopping.prompt_size_admission`),所以落进日志的每一条步记录都是通过了那一档 的,它们的 `prompt_chars` 必定不超过上限。真正超限的那一次装配没有在任何地方留下记录。 能从日志里判的是那一次装配会有多大。最后一步之后历史里多出两条消息:它的 `raw_output`, 以及套上观察模板的 `observation`(`polyloop._assembly.history_messages`)。两者的字符数加上 最后一步的 `prompt_chars`,就是下一次装配出来的规模。这个数不超过上限的话,说明库在还装得 下的时候就报了 `context_overflow`。 另一头也要判:最后一步自己的 `prompt_chars` 超过上限,说明有一次超限的装配被放行去调了 模型,那正是这一档该拦住的东西。 上限与观察模板都从参数快照里读,不从本文件的常量读——这一路的上限是按上下文现算的。 """ name = "prompt_reached_max_prompt_chars" snapshot = parameter_snapshot_of(read) if snapshot is None: return undetermined(name, "日志里没有 run_started 记录,读不到这次运行的提示词上限") limit_text = snapshot.get(MAX_PROMPT_CHARS_SNAPSHOT_KEY) template = snapshot.get(OBSERVATION_TEMPLATE_SNAPSHOT_KEY) if not isinstance(limit_text, str) or not isinstance(template, str): return undetermined( name, f"参数快照里缺 {MAX_PROMPT_CHARS_SNAPSHOT_KEY} 或 " f"{OBSERVATION_TEMPLATE_SNAPSHOT_KEY},算不出该不该撞上限", ) try: limit = int(limit_text) except ValueError: return undetermined(name, f"参数快照里的提示词上限不是一个整数:{limit_text!r}") steps = _step_payloads(read) if not steps: return undetermined(name, "日志里一条步记录都没有,撞没撞上上限无从判起") last = steps[-1] chars = last.get("prompt_chars") raw_output = last.get("raw_output") observation = last.get("observation") if ( not isinstance(chars, int) or isinstance(chars, bool) or not isinstance(raw_output, str) or not isinstance(observation, str) ): return undetermined(name, "最后一条步记录缺 prompt_chars、raw_output 或 observation") try: rendered = template.format(observation=observation) except (KeyError, IndexError, ValueError) as exc: return undetermined(name, f"观察模板渲染不了({exc}),算不出下一步的提示词规模") if chars > limit: return breached( name, f"最后一步自己的提示词就有 {chars} 字符,已经超过上限 {limit}:" "有一次超限的装配被放行去调了模型", ) projected = chars + len(raw_output) + len(rendered) if projected <= limit: return breached( name, f"最后一步的提示词 {chars} 字符,加上它的输出与套过模板的观察一共 {projected} 字符," f"仍不超过上限 {limit}:还装得下就报了 context_overflow", ) return passed( name, f"最后一步的提示词 {chars} 字符没超上限 {limit},再走一步会是 {projected} 字符,超了", ) def check_last_step_action_status(read: LogRead, *, expected: str) -> Criterion: """最后一条步记录的动作状态是期望的那个。""" name = f"last_step_action_status_is_{expected}" steps = _step_payloads(read) if not steps: return undetermined(name, "日志里一条步记录都没有,动作状态无从判起") actual = steps[-1].get("action_status") if actual != expected: return breached(name, f"最后一步的动作状态是 {actual!r},期望 {expected!r}") return passed(name, f"最后一步的动作状态是 {expected}") def check_last_observation_is_synthetic(read: LogRead, *, expected: str, name: str) -> Criterion: """最后一步的观察被换成了库自己那段合成文本,而且逐字就是装配时给的那一段。 环境故障那一档下,执行接缝返回的观察不进历史,库改填 `SyntheticObservations.env_failed` (`polyloop.session` 的 `_project_observation`)。这么定是因为观察属于「模型看得见的东西」, 那种东西必须能进参数快照;执行接缝每次现造一段的话,同一份配置跑出来的两次运行在模型看来 其实不同,而没有任何地方会报错。 期望值由调用处从装配时用的那份取,这里不抄字面量。 **证据只报长度不报正文。** 替换要是没发生,那一列里躺着的是环境(或与环境通信那一层)给的 文本,它可以是任何东西,而这份报告是要贴给人看的。 """ steps = _step_payloads(read) if not steps: return undetermined(name, "日志里一条步记录都没有,观察无从判起") last = steps[-1] flag = last.get("observation_is_synthetic") if flag is not True: return breached(name, f"最后一步的 observation_is_synthetic 是 {flag!r},期望 True") actual = last.get("observation") if not isinstance(actual, str): return breached(name, f"最后一步的 observation 不是字符串,是 {type(actual).__name__}") if actual != expected: return breached( name, f"最后一步的观察不是装配时给的那段合成观察(实际 {len(actual)} 字符," f"期望那段 {len(expected)} 字符)", ) return passed(name, f"最后一步的观察逐字是装配时给的那段合成观察({len(expected)} 字符)") def check_steps_before_last_all_executed(read: LogRead) -> Criterion: """环境坏掉之前的那些步没有被牵连:它们的动作状态都是「已执行」。 这一条守的是「环境故障只影响它发生的那一步」。前面的步被改写或被补上一个别的状态,说明 库把一次局部故障扩散到了已经落地的轨迹上,而那些步的观察是模型接下来要看的东西。 """ name = "steps_before_env_error_executed" steps = _step_payloads(read) if len(steps) < 2: return undetermined( name, f"日志里只有 {len(steps)} 条步记录,环境故障那一步之前一步都没有,这一条真空成立", ) bad = [ index for index, step in enumerate(steps[:-1]) if step.get("action_status") != "executed" ] if bad: return breached( name, f"第 {bad} 步(按步记录次序计)的动作状态不是 executed,它们在环境坏掉之前" ) return passed(name, f"环境坏掉之前的 {len(steps) - 1} 步动作状态全是 executed") def check_env_broken_after_a_full_step( *, broken: bool, executions_before_break: int, expected_after: int ) -> Criterion: """动手弄坏环境的时机正是声明的那一次:完整执行过 `expected_after` 次动作之后。 第一次执行之前就把容器打掉的话,验的是「环境起不来」而不是「跑到一半环境不能接着服务 了」——前者落在会话初始化上,根本走不到动作执行接缝,而这一类要压的正是那个接缝。 **动手时机与声明的不一致是击穿,不是无法判定。** 这一条判的是注入器自己:它对外宣称在第 `expected_after` 次执行之后动手,整类故障的结论都建立在这句话上。真在别的时刻动手的话, 这一类压到的是另一件事,而报告上仍然写着「环境故障那一类通过了」——那正是「判据验的不是 它声称要验的东西」这种坏法。判成无法判定则把这件事说成「看不清」,可它看得很清楚:注入器 自己记下了动手时它数到几。 `broken` 为假是另一回事,那一档确实什么都没发生(模型一次动作都没成功执行过),报无法判定。 """ name = "env_broken_after_a_full_step" if not broken: return undetermined( name, "这次运行结束时环境一次都没被弄坏过:动作执行没走到该动手的那一次" ) if executions_before_break != expected_after: return breached( name, f"声明的是完整执行 {expected_after} 次动作之后动手,实际在第 " f"{executions_before_break} 次之后动手:" + ( "压到的是会话初始化而不是动作执行接缝" if executions_before_break < expected_after else "比声明的晚,这一类压到的不是它声称的那个时刻" ), ) return passed(name, f"完整执行过 {executions_before_break} 次动作之后才把环境弄坏,与声明一致") def check_all_steps_parse_failed(read: LogRead) -> Criterion: """这几步全部解析失败,且一个动作都没被分发。 `action_status` 为空是「这一步压根没走到动作那一档」的形态:解析失败那一支直接跳过完成 判定与动作执行(`polyloop.session` 的 D 档)。它要是有值,说明有动作被分发过。 """ steps = _step_payloads(read) if not steps: return undetermined("all_steps_parse_failed", "日志里一条步记录都没有") bad_parse = [index for index, step in enumerate(steps) if step.get("parse_ok") is not False] bad_action = [ index for index, step in enumerate(steps) if step.get("action_status") is not None ] if bad_parse or bad_action: return breached( "all_steps_parse_failed", f"第 {bad_parse} 步的 parse_ok 不是 False,第 {bad_action} 步的 action_status 不为空", ) return passed( "all_steps_parse_failed", f"{len(steps)} 步全部 parse_ok=False 且 action_status 为空", ) def check_env_untouched(*, executions: int, source: str) -> Criterion: """环境侧一次都没被执行过。 解析失败不该碰环境,这是契约。计数来自环境侧(AppWorld 的 `n_executions` 或 GovDoc 的 审计条数),不来自轨迹。 """ if executions != 0: return breached("env_untouched", f"{source} 报环境被执行了 {executions} 次,期望 0 次") return passed("env_untouched", f"{source} 报环境一次都没被执行") def check_cancelled_raised(*, raised: BaseException | None) -> Criterion: """`await task` 抛的是 `CancelledError`,原样抛出。 吞掉之后返回一个结果,调用方的结构化并发就断了:它以为这次运行正常结束,而它其实是被 自己叫停的。 """ if raised is None: return breached("cancelled_error_propagated", "取消之后 await 正常返回了结果,没有抛异常") if not isinstance(raised, asyncio.CancelledError): return breached( "cancelled_error_propagated", f"取消之后 await 抛的是 {type(raised).__name__},不是 CancelledError", ) return passed("cancelled_error_propagated", "取消之后 await 原样抛出了 CancelledError") def check_env_quiet_after_cancel( *, dispatched_before: int, dispatched_after: int, completed_before: int, completed_after: int, waited_s: float, ) -> Criterion: """取消之后环境侧不再有动静:动作执行接缝没被再进入,也没有哪次执行又跑完了。 `dispatched_*` 数的是动作执行接缝被**进入**过几次(由 `TriggeringExecutor` 记), `completed_*` 数的是环境侧确认**跑完**过几次(`AppWorldSession.n_executions`)。前者涨了 说明库在取消之后还在往环境派活;后者涨了说明有一次执行在取消之后才落地。两个数分别对应 两种不同的坏法,所以一起判。 **已知看不见的那一半**:`n_executions` 是客户端侧的计数,`tools/soak/appworld.py` 在 HTTP 响应返回之后才加一。取消把那个协程掐断了,于是「代码已经发到容器、容器仍然把它跑完 并改了环境状态」这种情形在这两个数上都留不下痕迹。容器那侧没有可查的执行计数接口 (只有 `/execute`、`/task_completed`、`/evaluate`、`/close`),补上它要改环境层。这条缺口 写在模块 docstring 的已知缺口一节里,这里的「通过」只覆盖上面那两种坏法。 """ name = "env_quiet_after_cancel" if dispatched_after > dispatched_before: return breached( name, f"取消之后动作执行接缝又被进入了 {dispatched_after - dispatched_before} 次" f"(等了 {waited_s} 秒再数的):库在取消之后还在往环境派活", ) if completed_after > completed_before: return breached( name, f"取消之后又有 {completed_after - completed_before} 次执行在环境侧跑完了" f"(等了 {waited_s} 秒再数的)", ) if completed_after < completed_before or dispatched_after < dispatched_before: return breached( name, f"取消之后计数倒退了(派活 {dispatched_before}→{dispatched_after}," f"跑完 {completed_before}→{completed_after}),这两个数只该单调不减", ) return passed( name, f"取消之后等了 {waited_s} 秒,派活次数仍是 {dispatched_after}、" f"环境侧跑完次数仍是 {completed_after}", ) def check_events_file_intact(path: Path) -> Criterion: """事件文件的每一行都是完整的一条 JSON 对象。 取消可能落在事件写盘的途中。**撕裂尾行在这里算击穿,不像日志那侧算「没发生过」**:库有 宽限期(`RunRequest.cancel_grace_seconds`),在飞的写有机会收尾,而事件出口每条只写一次 `handle.write`。真留下半行,说明取消穿过了一次没有被保护起来的写——那种写下次也会半途而废, 而下游读事件流时看到的是一份解不开的文件。 """ name = "events_file_intact" if not path.is_file(): return undetermined(name, f"没有 {path.name}:这次运行一条事件都没发出过,完整性无从判起") raw = path.read_bytes() if not raw: return undetermined(name, f"{path.name} 是空的,完整性无从判起") chunks = raw.split(b"\n") if chunks[-1].strip(): return breached( name, f"{path.name} 末尾有 {len(chunks[-1])} 字节没有被换行终结:取消穿过了一次事件写入的中途", ) bad: list[int] = [] total = 0 for number, chunk in enumerate(chunks[:-1], start=1): if not chunk.strip(): continue total += 1 try: payload = json.loads(chunk) except (UnicodeDecodeError, json.JSONDecodeError): bad.append(number) continue if not isinstance(payload, dict): bad.append(number) if bad: return breached(name, f"{path.name} 的第 {bad} 行不是一个完整的 JSON 对象") return passed(name, f"{path.name} 的 {total} 行全都是完整的 JSON 对象") def check_lease_returned(*, borrowed: bool, timeout_s: float, pool_size: int) -> Criterion: """容器租约被归还:取消结束之后还借得到。 池满时 `lease()` 会一直阻塞,所以「借得到」等价于「空闲名额回到了满」。池的大小是 1, 于是这一借要么立刻成功、要么永远等下去,中间没有含糊地带。 """ if not borrowed: return breached( "container_lease_returned", f"取消之后再借一个容器,{timeout_s} 秒内没借到(池大小 {pool_size}),租约没还回来", ) return passed( "container_lease_returned", f"取消之后在 {timeout_s} 秒内又借到了容器(池大小 {pool_size}),租约还回来了", ) # --------------------------------------------------------------------------- # 五、子进程编排:子进程按时机自杀,外部 SIGKILL 兜底 # --------------------------------------------------------------------------- #: 子进程按时机自杀时用的退出码。父进程靠它区分「按预期崩在时机上」与「因为别的原因退出」。 #: #: 取 137 是照 128+9 那个惯例(外部 SIGKILL 的等价形态),让两条路径在日志里读起来是同一件事。 CRASH_EXIT_CODE = 137 #: 条件对上之后留给子进程自杀的窗口,秒。窗口内它还没死,父进程才兜底发 SIGKILL。 #: #: 两秒对「进程执行完 `os._exit` 系统调用」来说是天文数字——这个数不是在等一件慢事,是在给 #: 两个进程的观察顺序留一点余量:父进程是从文件里看见那条记录的,而子进程要先从写入调用返回、 #: 再走几行 Python 才到自杀那一句。取小了会偶发地抢在它前面,而抢赢的表现是「这次又是兜底 #: 命中的」,不是一个错误。 SELF_KILL_GRACE_S = 2.0 class KillTiming(StrEnum): """在哪个时机让子进程崩掉。两种时机的判据不同,不许混成一个用例。""" #: 时机 A:最后一条是 `step_completed`,一步完整落地之后崩。续跑该真的接着往下跑。 AFTER_STEP = "after_step" #: 时机 B:最后一条是一条声明绝不重放的意图,意图写了、结果还没写。续跑该判状态未知、 #: 干净停下。 AT_INTENT = "at_intent" def should_kill( read: LogRead, *, timing: KillTiming, after_steps: int, audit_is_not_empty: bool ) -> bool: """现在这份日志尾部是不是要等的那个时机。纯函数。 子进程自杀之后父进程拿它复核一遍尾部形态;外部 SIGKILL 那条兜底路径每次轮询也问它一句。 时机 B 额外要求那条意图的重放策略是 `never`。**这比「最后一条是意图」更严**,而且必须 更严:GovDoc 的 `read_document` 与 `grep_document` 声明的是 `safe`,悬在那种意图上续跑 会重放动作接着跑,停止原因不是 `resume_state_unknown`。放宽这一条,判据就会时对时错, 而错的那些次看起来只是「模型这次走了别的路」。 **`audit_is_not_empty` 必传,两种时机都要求它为真。** 它与 `SelfKillingStore` 的自杀条件 是同一条,两边必须逐字对齐:父进程这侧要是松一档(只看日志尾部),它总会在子进程走到自杀 那一行之前抢先发出信号,于是自杀路径成了永远走不到的死代码,而崩溃点落在哪儿又变回碰运气。 实测就是这么发生的——两次崩溃全是外部信号命中的,其中一次崩在审计账还是空的时候,最硬那条 判据只能真空成立。 """ if not audit_is_not_empty: return False if len(tagged(read, "step_completed")) < after_steps: return False if not read.payloads: return False last = read.payloads[-1] if timing is KillTiming.AFTER_STEP: return last.get(RECORD_KEY) == "step_completed" return last.get(RECORD_KEY) == "intent" and last.get("replay_policy") == "never" class SelfKillingStore: """包一层 `RunStore`:某一次写入落盘返回之后,按时机让子进程当场死掉。 **为什么不再靠外部 SIGKILL 抢窗口。** 实测下来时机 A 一次都没命中过,三次全都报「发信号 与子进程停笔之间又写进了记录」:库写完 `step_completed` 紧接着就写下一步的模型调用意图, 中间只有内存里的装配计算,那个窗口窄到外面的信号挤不进去。**这个观察本身值得记住**—— 它说明自然发生的崩溃几乎总是落在「有意图没结果」那一态上,而不是落在两步之间的干净边界 上。但要验时机 A 就不能靠碰运气,得让子进程自己在那一点上死。 **`os._exit` 与 SIGKILL 对磁盘的效果等价。** 它不跑 `finally`、不跑 `atexit`、不 flush 任何缓冲,直接进 `_exit(2)` 系统调用。`JsonlRunStore` 每次写完自己 `fsync`(意图与运行 开始、运行结束三处)或者靠同一文件的追加序兜(`0005` 决策五那条前缀持久性),本来就没有 未刷缓冲要指望进程退出时替它写下去,所以「有没有机会清理」在这里不影响磁盘上留下什么。 **`parameters()` 原样转发内层的**,不加自己的键:它进参数快照,而父进程续跑时用的是一个 没包过的 `JsonlRunStore`。加一个键,续跑就报参数漂移,而那是一次假故障。 **自杀条件带上「审计账已经非空」**:要验的最硬那条判据是「声明绝不重放的动作没有被执行 两次」,它数的是工作区审计账。崩在模型还没调过 `write_note` 的时候,账是空的,那条判据 真空成立——报出来是「通过」,实际什么都没验。 父进程那侧的兜底条件(`should_kill`)与这里逐字对齐,包括审计账那一项。松一档它就会每次 抢先,这个类成为死代码;那件事实测发生过一次,理由写在 `should_kill` 的 docstring 里。 """ __slots__ = ("_exit_now", "_inner", "_timing", "_workspace") def __init__( self, *, inner: RunStore, timing: KillTiming, workspace: Path, exit_now: Callable[[int], object] = os._exit, ) -> None: """Args: inner: 真正写盘的那个存储。 timing: 在哪一次写入之后死。 workspace: 数审计账用的工作区目录。 exit_now: 怎么死。**做成参数是为了能测**——默认的 `os._exit` 在测试里会把 pytest 自己一起带走,测试注入一个抛哨兵异常的替身。 """ self._inner = inner self._timing = timing self._workspace = Path(workspace) self._exit_now = exit_now def parameters(self) -> Mapping[str, str]: return self._inner.parameters() def _audit_is_not_empty(self) -> bool: return bool(read_audit_lines(self._workspace)) async def read_log(self, run_id: str) -> RunLog: return await self._inner.read_log(run_id) async def write_run_started(self, record: RunStarted) -> None: await self._inner.write_run_started(record) async def write_intent(self, record: Intent) -> None: await self._inner.write_intent(record) if ( self._timing is KillTiming.AT_INTENT and record.replay_policy is ReplayPolicy.NEVER and self._audit_is_not_empty() ): self._exit_now(CRASH_EXIT_CODE) async def write_model_call_result(self, record: ModelCallResult) -> None: await self._inner.write_model_call_result(record) async def write_step_completed(self, record: StepCompleted) -> None: await self._inner.write_step_completed(record) if self._timing is KillTiming.AFTER_STEP and self._audit_is_not_empty(): self._exit_now(CRASH_EXIT_CODE) async def write_run_finished(self, record: RunFinished) -> None: await self._inner.write_run_finished(record) @dataclass(frozen=True, slots=True, kw_only=True) class KillOutcome: """一次「起子进程、等它崩在时机上」的结果。""" #: 时机命中了没有。没命中要重试,重试若干次仍不命中要报出来——悄悄降级成另一种时机会让 #: 报告显示「验过了」而其实验的是另一件事。 hit: bool reason: str #: 崩溃之后重读日志、截到最后一个换行为止的那一段。字节比对拿它当基准。 snapshot: bytes = b"" #: 崩溃时已经完整落地的步数。 steps: int = 0 #: 这次子进程发起过几次模型调用(含崩在半路的那次),记账用。 model_calls: int = 0 #: 子进程的退出码。按预期自杀是 `CRASH_EXIT_CODE`,被外部 SIGKILL 兜底掉是 -9。 exit_code: int | None = None def _judge_crash( *, log_path: Path, workspace: Path, timing: KillTiming, after_steps: int, exit_code: int | None, how: str, ) -> KillOutcome: """崩溃之后重读一次日志,判尾部形态是不是要的那个时机。 **不管子进程是自杀的还是被外部信号杀的,都要重判一次。** 自杀那条路上判的是「自杀条件与 时机判据说的是不是同一件事」;外部信号那条路上判的是「读日志与发信号之间它有没有又写进 一条」——不重判的话,一次「本想在时机 A 杀、实际杀在时机 B」会被当成时机 A 判下去。 """ crashed = log_path.read_bytes() if log_path.is_file() else b"" after = parse_terminated(crashed) audited = bool(read_audit_lines(workspace)) if not should_kill(after, timing=timing, after_steps=after_steps, audit_is_not_empty=audited): return KillOutcome( hit=False, reason=( f"{how}之后重读日志,尾部不是时机 {timing.value} 要的形态" + ("(审计账还是空的)" if not audited else "") ), model_calls=count_model_calls(after), exit_code=exit_code, ) return KillOutcome( hit=True, reason=f"{how},命中时机 {timing.value}", snapshot=terminated_prefix(crashed), steps=len(tagged(after, "step_completed")), model_calls=count_model_calls(after), exit_code=exit_code, ) async def spawn_and_kill( *, argv: Sequence[str], log_path: Path, workspace: Path, timing: KillTiming, after_steps: int, child_log_path: Path | None = None, poll_interval_s: float = 0.002, self_kill_grace_s: float = SELF_KILL_GRACE_S, timeout_s: float = 600.0, ) -> KillOutcome: """起一个子进程,等它崩在时机上。 **主路径是子进程自己在时机上 `os._exit`**(见 `SelfKillingStore`),父进程只负责认领: 退出码等于 `CRASH_EXIT_CODE` 且日志尾部形态对得上,就算命中。 **条件对上之后先等一个自杀窗口,窗口内子进程还活着才兜底发信号。** 父子两侧的条件现在 是同一条,所以父进程看见条件成立的那一刻,子进程正走在自杀那一行上——不留窗口的话父进程 每次都抢先,自杀路径成了永远走不到的死代码。实测就是这样:两次崩溃全是外部信号命中的。 窗口只用来吸收进程退出的那点延迟,所以取值很小。 **外部 SIGKILL 那条路留着兜底**,没有删掉:子进程那侧的自杀条件万一因为别的原因没触发 (包装漏了一处、子进程卡在别的地方),窗口过后仍然会把它杀掉。**用 SIGKILL 不用 SIGTERM**,要的是没有任何清理机会的死法——SIGTERM 会走 Python 的信号处理,`finally` 有机会跑完,那验的是优雅退出而不是崩溃。 子进程的输出**写进文件而不是管道**:管道缓冲区满了子进程会阻塞在写上,而表现是「它卡住 不动了」,从外面区分不出是卡在模型调用上还是卡在一行日志上。 """ if child_log_path is None: sink: object = asyncio.subprocess.DEVNULL handle = None else: child_log_path.parent.mkdir(parents=True, exist_ok=True) handle = child_log_path.open("ab") sink = handle try: process = await asyncio.create_subprocess_exec( *argv, cwd=str(REPO_ROOT), stdout=sink, stderr=asyncio.subprocess.STDOUT ) finally: if handle is not None: handle.close() def _exited(code: int) -> KillOutcome: if code == CRASH_EXIT_CODE: return _judge_crash( log_path=log_path, workspace=workspace, timing=timing, after_steps=after_steps, exit_code=code, how=f"子进程按时机自杀(退出码 {code})", ) return KillOutcome( hit=False, reason=( f"子进程以退出码 {code} 结束,不是按时机自杀的 {CRASH_EXIT_CODE}" + (f",输出见 {child_log_path.name}" if child_log_path else "") ), model_calls=count_model_calls(read_log(log_path)), exit_code=code, ) deadline = time.monotonic() + timeout_s try: while True: if process.returncode is not None: return _exited(process.returncode) matched = should_kill( read_log(log_path), timing=timing, after_steps=after_steps, audit_is_not_empty=bool(read_audit_lines(workspace)), ) if matched: grace_deadline = time.monotonic() + self_kill_grace_s while process.returncode is None and time.monotonic() < grace_deadline: await asyncio.sleep(poll_interval_s) if process.returncode is not None: return _exited(process.returncode) process.kill() await process.wait() return _judge_crash( log_path=log_path, workspace=workspace, timing=timing, after_steps=after_steps, exit_code=process.returncode, how=f"等了 {self_kill_grace_s} 秒不见子进程自杀,父进程兜底发了 SIGKILL", ) if time.monotonic() > deadline: return KillOutcome( hit=False, reason=f"等了 {timeout_s} 秒仍没崩在时机 {timing.value} 上", model_calls=count_model_calls(read_log(log_path)), ) await asyncio.sleep(poll_interval_s) finally: if process.returncode is None: process.kill() await process.wait() # --------------------------------------------------------------------------- # 六、接缝:事件出口、计数、故意坏掉的解释器、取消触发器 # --------------------------------------------------------------------------- class JsonlEventSink: """把事件逐行写进一个 `.jsonl`。满足 `polyloop.ports.EventSink`。 `parameters()` 里**不放路径**:它进参数快照,而续跑时父进程写的是另一个文件(子进程那 份事件随 SIGKILL 留在原地),路径进快照会报一次假的参数漂移。 """ __slots__ = ("_path", "delivered", "failures") def __init__(self, path: Path | str) -> None: self._path = Path(path) self.delivered = 0 self.failures = 0 def parameters(self) -> Mapping[str, str]: return {"kind": "jsonl_events"} async def emit(self, event: Event) -> None: line = json.dumps( { "kind": event.kind.value, "run_id": event.run_id, "step_idx": event.step.step_idx, }, ensure_ascii=False, ) await asyncio.to_thread(self._append, line + "\n") self.delivered += 1 def _append(self, line: str) -> None: self._path.parent.mkdir(parents=True, exist_ok=True) with self._path.open("a", encoding="utf-8") as handle: handle.write(line) class AlwaysInvalidParser: """对任何模型输出都返回无效决策。解析失败连击那一类用它。 它不是「解析不出来」,是**声明这一步解释不了**——契约要求 `parse` 同步、不抛异常、解释 不出来时返回 `InvalidDecision`(`polyloop.testing.DecisionParserContract`),这份实现照做。 """ def parameters(self) -> Mapping[str, str]: return {"kind": "always_invalid"} def parse(self, reply: ModelReply) -> ParsedReply: return ParsedReply( history_text=reply.content, decision=InvalidDecision( explanation="故障注入:这一路的决策解释器对任何输出都判无效,用来压解析失败连击。" ), ) class PromptSizeRecordingClient: """转发模型调用,并把**每次真正发出去的那份消息**有多少字符按调用序号记下来。 没有它的话,提示词那两条判据读的全是库写进步记录的 `prompt_chars`——那是库对自己行为的 陈述。失败场景很具体:库真的把历史静默截断了,同时仍然把 `prompt_chars` 记成一路不减的 自述值,两条判据照样全绿。这与「动作那一维要数环境侧的账、不数库报的步数」是同一条原则, 这里把它落在提示词这一维上。 字符数的口径逐字照 `polyloop._assembly.prompt_chars`:所有消息的所有内容块的文本长度之和。 口径不同的话对不上账,而对不上会被判成击穿,那是一次假故障。 **`parameters()` 原样转发内层的**,不加自己的键:它进参数快照,续跑时那边用的是没包过的 客户端,加一个键就是一次假的参数漂移。 """ __slots__ = ("_inner", "chars_by_call_index") def __init__(self, *, inner: object) -> None: self._inner = inner #: 调用序号 → 那次调用真正发出去的字符数。 self.chars_by_call_index: dict[int, int] = {} def parameters(self) -> Mapping[str, str]: return self._inner.parameters() # type: ignore[attr-defined] async def call(self, call: object) -> ModelReply: self.chars_by_call_index[call.call_index] = sum( # type: ignore[attr-defined] len(block.text) for message in call.messages # type: ignore[attr-defined] for block in message.content ) return await self._inner.call(call) # type: ignore[attr-defined] def check_prompt_chars_match_what_was_sent( read: LogRead, *, observed: Mapping[int, int] ) -> Criterion: """步记录里的 `prompt_chars` 与库那一步真正发出去的字符数逐步相等。 这一条是提示词那两条判据的地基:它们都只读 `prompt_chars`,而 `prompt_chars` 是库自己写 下的数。库把历史截断了却照旧记一个单调递增的值,那两条都会通过,而这一条会当场对不上账。 `observed` 由 `PromptSizeRecordingClient` 在调用发生的那一刻记下,键是调用序号;步记录的 步号与调用序号是同一个数(`polyloop._recovery` 模块 docstring 那条对齐),所以直接按步号查。 """ name = "prompt_chars_match_what_was_sent" steps = _step_payloads(read) if not steps: return undetermined(name, "日志里一条步记录都没有,对不了账") if not observed: return undetermined(name, "一次模型调用都没被记到,对不了账") mismatches: list[str] = [] compared = 0 for step in steps: index = step.get("step_idx") if not isinstance(index, int) or isinstance(index, bool) or index not in observed: continue logged = step.get("prompt_chars") if not isinstance(logged, int) or isinstance(logged, bool): mismatches.append(f"第 {index} 步的 prompt_chars 不是整数") continue compared += 1 if logged != observed[index]: mismatches.append( f"第 {index} 步记的是 {logged} 字符,实际发出去 {observed[index]} 字符" ) if mismatches: return breached(name, ";".join(mismatches)) if compared == 0: return undetermined(name, "没有一条步记录的步号能和记下来的调用序号对上,对不了账") return passed(name, f"{compared} 步的 prompt_chars 与真正发出去的字符数逐步相等") class TriggeringModelClient: """转发模型调用,并在进入第 N 次调用时通知父协程。父协程收到通知立刻取消。 这样取消落在模型调用**中途**——那条路径上在飞的是网关连接与一次已经计过费的请求。 """ __slots__ = ("_at_call_index", "_inner", "_trigger", "calls") def __init__(self, *, inner: object, trigger: asyncio.Event, at_call_index: int) -> None: self._inner = inner self._trigger = trigger self._at_call_index = at_call_index self.calls = 0 def parameters(self) -> Mapping[str, str]: return dict(self._inner.parameters()) | {"cancel_probe": "model_call"} # type: ignore[attr-defined] async def call(self, call: object) -> ModelReply: self.calls += 1 if call.call_index >= self._at_call_index: # type: ignore[attr-defined] self._trigger.set() return await self._inner.call(call) # type: ignore[attr-defined] class TriggeringExecutor: """转发动作执行,并在进入第一次执行时通知父协程。 取消落在环境执行中途——那条路径上在飞的是容器租约与一次已经发出去的 HTTP 请求,与模型 调用那条路上的资源完全不同,所以两条各做一次。 """ __slots__ = ("_inner", "_trigger", "executions") def __init__(self, *, inner: object, trigger: asyncio.Event) -> None: self._inner = inner self._trigger = trigger self.executions = 0 def parameters(self) -> Mapping[str, str]: return dict(self._inner.parameters()) | {"cancel_probe": "env_execute"} # type: ignore[attr-defined] async def execute(self, action: object) -> object: self.executions += 1 self._trigger.set() return await self._inner.execute(action) # type: ignore[attr-defined] def container_name_for_port(port: int) -> str: """池里跑在这个宿主端口上的容器叫什么。 形状是「前缀 + 端口号」,前缀取 `tools/soak/appworld.py` 公开的那个常量,端口取池的 `ports` 属性。**不去调容器池那个同名的私有方法**:从外面 import 一个下划线开头的名字等于把它的实现 细节钉死。代价是这里拼错了没人会当场发现,所以拼错的后果被做成显式失败——`docker kill` 会 以「没有这个容器」的退出码告诉我们,那一路直接抛错,不会静默地跑成一次「环境没坏」的运行。 """ return f"{CONTAINER_NAME_PREFIX}-{port}" def should_break_env( *, executions_done: int, break_after_executions: int, already_broken: bool ) -> bool: """现在这一次动作执行之前,该不该动手把环境弄坏。纯函数。 `executions_done` 数的是**已经成功走完**的动作次数,所以 `break_after_executions` 取 1 就是 「第一步完整走完之后、第二步的执行之前动手」。取 0 会把容器打在第一次执行之前,那时压到的 是会话初始化而不是执行接缝。 弄坏过一次之后不再动手:容器已经没了,再发一次 `docker kill` 只会拿到一个「没有这个容器」 的错误,而那个错误会被当成「弄坏失败」报出来。 """ if already_broken: return False return executions_done >= break_after_executions async def _docker(*args: str) -> tuple[int, str, str]: """跑一条 docker 命令,返回 (退出码, stdout, stderr)。 `tools/soak/appworld.py` 里有一个形状相同的私有函数,这里不 import 它——理由与 `container_name_for_port` 那条相同。 """ process = await asyncio.create_subprocess_exec( "docker", *args, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE ) stdout, stderr = await process.communicate() return ( process.returncode or 0, stdout.decode(errors="replace"), stderr.decode(errors="replace"), ) async def kill_and_remove_container(name: str) -> str: """把这个容器打死并删掉,返回一句给 note 用的说明。 **两步分开发。** `docker rm -f` 自己也会先发 SIGKILL,一条命令就够,但分成两条之后 「环境从这一刻起服务不了了」与「容器有没有被留在盘上」在报告里各是一句话:第二步失败只是 留下一个已经死掉的容器要人手动删,故障本身照样注入成功。 第一步失败则直接抛错:那说明容器名拼错了或者容器根本不在,环境没被弄坏,这一路接下来跑出 来的东西什么都不是。抛出去由 `guarded()` 记成「无法判定」。 """ code, stdout, stderr = await _docker("kill", name) if code != 0: raise FaultInjectionError( f"docker kill {name} 失败(退出码 {code}):{(stderr or stdout).strip()};" "环境没有被弄坏,这一类没有意义" ) remove_code, remove_out, remove_err = await _docker("rm", "-f", name) if remove_code != 0: return ( f"容器 {name} 已被 docker kill,但随后的 docker rm -f 失败" f"(退出码 {remove_code}):{(remove_err or remove_out).strip()},要人手动删" ) return f"容器 {name} 已被 docker kill 并删除" class EnvBreakingExecutor: """转发动作执行,并在跑完指定次数之后真的把容器打死。 **弄坏的是真环境,不是一个测试替身。** 容器没了之后,后续的 `execute` 连不上那个端口, 这一路要验的就是库拿到 `ENV_ERROR` 之后做了什么。 **它还要把连不上翻译成 `ENV_ERROR`,因为场景那侧不翻译。** `tools/soak/scenarios/appworld.py` 的 `AppWorldExecutor` 只接 `AppWorldError`,而 `AppWorldSession.execute` 在容器没了时抛的是 `httpx.ConnectError`——`_post_json` 直接 `await client.post(...)`,没有把传输层的异常包起来。实测过(`docker kill` 之后 execute 抛 `httpx.ConnectError: All connection attempts failed`)。不翻译的话这次运行会整个抛出去, 根本产不出 `env_error` 这个停止原因。翻译这件事本来就归执行接缝:判据是「环境还能不能接着 服务」,由适配器判、库不判(`research-wiki/design/0007-seam-behaviour.md` 决策一)。 内层自己给出 `ENV_ERROR` 时原样透传,不覆盖——那样场景那边哪天补上了转换,这里不用跟着改。 """ __slots__ = ( "_break_after", "_break_env", "_inner", "broken", "break_note", "executions", "executions_before_break", ) def __init__( self, *, inner: object, break_after_executions: int, break_env: Callable[[], Awaitable[str]], ) -> None: """Args: inner: 真正打环境的那个执行器。 break_after_executions: 完整执行过几次动作之后动手。见 `should_break_env`。 break_env: 怎么弄坏,返回一句说明。**做成参数是为了能测**——默认那条路要起真容器。 """ self._inner = inner self._break_after = break_after_executions self._break_env = break_env self.executions = 0 self.broken = False self.executions_before_break = 0 self.break_note = "" def parameters(self) -> Mapping[str, str]: return dict(self._inner.parameters()) | {"env_break_probe": "docker_kill"} # type: ignore[attr-defined] async def execute(self, action: object) -> ActionOutcome: if should_break_env( executions_done=self.executions, break_after_executions=self._break_after, already_broken=self.broken, ): self.executions_before_break = self.executions self.break_note = await self._break_env() self.broken = True try: outcome = await self._inner.execute(action) # type: ignore[attr-defined] except httpx.HTTPError as exc: # 环境不能接着服务了。观察填异常的类名与文本:它是与环境通信那一层给的,不是库 # 合成的占位,所以 `observation_is_synthetic` 为假——库随后会把这一列换成自己的 # 那段合成观察,而这一类要验的正是那次替换。 return ActionOutcome( status=ActionStatus.ENV_ERROR, observation=f"{type(exc).__name__}: {exc}", observation_is_synthetic=False, env_reported_completion=False, observation_truncated_chars=0, ) self.executions += 1 return outcome # type: ignore[no-any-return] @dataclass(slots=True) class CallGuard: """调用数护栏:花掉多少次真实模型调用,还剩多少。 **它按每类故障的边界拦,不在模型客户端里拦。** 在客户端里抛异常的话,库会把那次调用记成 一次模型调用失败、合成一段观察接着跑,于是护栏本身变成了一次注入进来的故障,把要验的 停止原因搅乱。按边界拦的代价是可能超出一点点,收益是判据看到的运行是干净的。 """ limit: int spent: int = 0 @property def remaining(self) -> int: return max(0, self.limit - self.spent) def charge(self, calls: int) -> None: self.spent += calls def affordable(self, worst_case: int) -> bool: return self.spent + worst_case <= self.limit # --------------------------------------------------------------------------- # 七、sidecar # --------------------------------------------------------------------------- def write_sidecars( *, runs_dir: Path, run_id: str, scenario: str, fault: str, result: RunResult | None, wall_ms: int, model_calls: int, sink_failures: int, env_executions: int, task_id: str | None, phase: str | None, resumed_from_step: int | None, ) -> None: """按记分板的约定写 `.result.json` 与 `.meta.json`。 `.events.jsonl` 不在这里写——它由 `JsonlEventSink` 边跑边写,那才是它的真实形态。 **结果为空时不写 `.result.json`。** 取消那条路径上 `run` 不返回结果,硬造一份是伪造; 记分板会退回日志里的结束记录,那份是真的。 `success` 恒为 `None`:故障注入的运行没有「任务做没做成」这回事,填一个布尔值会让记分板 的成功率统计混进一批语义不同的行。 """ runs_dir.mkdir(parents=True, exist_ok=True) if result is not None: (runs_dir / f"{run_id}.result.json").write_text( json.dumps(encode(result), ensure_ascii=False), encoding="utf-8" ) meta: dict[str, object] = { "scenario": scenario, "task_id": task_id, "phase": phase, "wall_ms": wall_ms, "model_calls": model_calls, "sink_failures": sink_failures, "env_executions": env_executions, "fault": fault, "success": None, "resumed_from_step": resumed_from_step, } (runs_dir / f"{run_id}.meta.json").write_text( json.dumps(meta, ensure_ascii=False), encoding="utf-8" ) # --------------------------------------------------------------------------- # 八、GovDoc 崩溃续跑 # --------------------------------------------------------------------------- #: 崩溃续跑跑的是 execute 阶段。**选它是因为它最快走到一次有副作用的动作**:提示词让模型先 #: 读 plan.md、再按行号核实、然后 `write_note` 写 evidence.md,而 `write_note` 的重放策略是 #: `never`——那正是最硬那条判据要数的东西。plan 阶段要先自己检索出候选行号才写得成笔记, #: 多烧好几次调用。 CRASH_PHASE = "execute" #: 预置在工作区里的计划文件。真实场景里它由 plan 阶段写出来,这里直接写盘省掉那一整次运行。 #: #: **由父进程用普通文件写入放进去,不走 `write_note`**:走工具的话审计账里会先躺一条记录, #: 而那条记录不是这次运行产生的,会把去重比对的基数弄脏。 #: #: **写得极其直白,只留一处候选、只留一条待办**:两种崩溃时机的触发条件都要求审计账已经非空, #: 也就是模型必须先真的调过一次 `write_note`。计划写得散一点,模型会先检索几轮再动笔,六步的 #: 预算跑完都还没写过笔记,那一类就只能报「无法判定」。这里不改 `govdoc.py` 的阶段提示词 #: (那是场景层的东西),只改这份由本文件自己造的计划正文。 #: #: 行号是随手挑的一段正文,正文内容不在这里出现,也没有任何一个字取自真实文书。 SEEDED_PLAN_NAME = "plan.md" SEEDED_PLAN_TEXT = """# 审核计划(故障注入脚本预置) 候选证据只有一处,已经定位好了:**tender.md 第 880 到第 910 行**。 只做两件事,做完就停: 1. 用 read_document 读 tender.md 的第 880 到第 910 行。 2. 立刻用 write_note 把上一步读到的原文逐字摘录写进 evidence.md,注明文档名与行号区间。 不要再检索别的关键词,不要再读别的段落——其余部分与本审核点无关。evidence.md 写完就停下。 """ #: 崩溃续跑这一路的预算。**比 GovDoc 场景自己那份(50 步)小得多**:这里要验的是崩溃与续跑 #: 的接缝,接缝在头几步就压得到,而每多一步都是一次真实的模型调用。 #: #: 六步是留给「读一段 + 写一次笔记」的余量:崩溃的触发条件要求模型先真的调过一次 `write_note` #: (见 `SEEDED_PLAN_TEXT`),照计划走是第 2 步就写,留到六步是给它两三次走弯路的机会。跑到 #: 预算用完仍然一次笔记都没写过,那一类报「无法判定」并说明原因,不假装验过。 CRASH_BUDGET = Budget( max_steps=6, max_actions=6, max_consecutive_parse_failures=3, max_prompt_chars=400_000, ) #: 模型绑定必须在父子两侧逐字相同——它整个进参数快照,差一个键续跑就报参数漂移。 CRASH_MODEL_BINDING: Mapping[str, str] = {"scenario": "govdoc", "phase": CRASH_PHASE} #: 时机判据里「至少已经落地几步」这一项。 #: #: **取 1,不再拿它当「等模型写笔记」的代理**:等笔记这件事现在由自杀条件本身负责(审计账 #: 非空),而自杀点可能落在第 1 步的 `step_completed` 上。这里再要求两步的话,一次本来完全 #: 正确的崩溃会被判成没命中。留着 1 是为了保证崩溃前至少有一步完整轨迹可供字节比对与续跑。 CRASH_AFTER_STEPS = 1 APPWORLD_MODEL_BINDING: Mapping[str, str] = {"scenario": "appworld"} #: 撞提示词上限那一路,除提示词上限外的三个数。 #: #: 步数上限取 3 是**兜底**:正常情况下第 1 步走完、第 2 步装配时就撞上提示词那一档,走不到这 #: 里。它存在是为了万一那一档没命中时这次跑仍然会停下来,而不是一路烧到 GovDoc 场景自己那份 #: 50 步的预算。取 3 不取 2 是留一步余量给「第一次模型输出解析不了」这种意外。 CONTEXT_OVERFLOW_BUDGET_STEPS = 3 #: 提示词上限在初始提示词之上留的余量,字符。 #: #: **取值的两头都是实测出来的。** 下界是 0:`initial_prompt_chars` 算出来的数与库自己装配出来 #: 的第一步 `prompt_chars` 逐字符相同(拿 GovDoc execute 阶段实测,两边都是 2828),所以留 0 #: 也够第一步过关。上界是「一步能把提示词撑大多少」:48 次真实 GovDoc execute 运行里,从第 0 #: 步到第 1 步的增长最少 157 字符、中位数 989。余量必须小于那个最小值,否则第 2 步的装配可能 #: 仍然装得下,这一类就撞不上上限。64 落在两者中间,两边都留着一倍以上的距离。 CONTEXT_OVERFLOW_SLACK_CHARS = 64 def initial_prompt_chars(request: RunRequest) -> int: """一步都还没走时,这次运行装配出来的提示词有多少字符。 照 `polyloop._assembly.assemble` 在零步下的形态数:run 级片段、注入槽、目标级片段,历史 是空的。**自己数而不是 import 那个模块**:它是下划线开头的内部模块,从压测工具里 import 等于把它的内部结构钉死。代价是这里与库的度量口径可能漂移,所以取值不是拍脑袋来的——实测 过它与库写进第一条步记录的 `prompt_chars` 相同,而且余量的取法(`CONTEXT_OVERFLOW_SLACK_ CHARS`)本来就容得下几十个字符的偏差。 认不得的内容块类型直接抛错,不当成 0:当成 0 会让上限算小,于是第一步就撞上限、一步都走 不成,而报出来只是一句「无法判定」。 """ total = 0 messages = (*request.context.run_level, *request.context.goal_level) for message in messages: for block in message.content: if not isinstance(block, TextBlock): raise FaultInjectionError( f"上下文里有认不得的内容块类型 {type(block).__name__},数不出提示词规模" ) total += len(block.text) for entries in request.injections.values(): for entry in entries: total += len(entry.content) return total def build_context_overflow_budget(request: RunRequest) -> Budget: """把提示词上限压到「刚好装得下第一步、装不下第二步」。 **上限按这次运行的上下文现算,不写死一个数字。** 上下文的大小取决于挑中的是哪个审核点、 语料清单有多长,写死的话它今天够用、换一份数据就要么第一步就撞上限(什么都没验到),要么 宽到几步都撞不上(白烧模型调用)。现算出来的这个数在 GovDoc execute 阶段是三千字符上下。 """ return Budget( max_steps=CONTEXT_OVERFLOW_BUDGET_STEPS, max_actions=CONTEXT_OVERFLOW_BUDGET_STEPS, max_consecutive_parse_failures=CONTEXT_OVERFLOW_BUDGET_STEPS, max_prompt_chars=initial_prompt_chars(request) + CONTEXT_OVERFLOW_SLACK_CHARS, ) #: 崩溃那一刻的日志快照落盘时用的后缀。 CRASH_SNAPSHOT_SUFFIX = ".crash-snapshot" def load_govdoc_task(*, govdoc_db: Path, govdoc_corpus: Path) -> AuditTask: """装配崩溃续跑用的那一个审核任务:第一条已批准的审核点配脱敏后的招标文书。 父子两个进程各装配一次,装出来的必须一致。**它确实一致**:审核点按 id 排序取第一条, 语料是同一个文件过同一份脱敏器。就算不一致也不会静默出错——上下文与语料不进参数快照, 进快照的是预算、工具集、绑定与四个接缝上报的参数,那些是常量。 """ checkpoints = load_checkpoints(db_path=govdoc_db, limit=1) loaded = load_documents(prepared_dir=govdoc_corpus, redactor=Redactor()) return AuditTask( index=0, checkpoint=checkpoints[0], documents=tuple(document for document, _report in loaded), ) def build_crash_request(*, task: AuditTask, run_id: str, workspace: Path) -> RunRequest: """装配崩溃续跑那次运行的请求。父子两侧调的是同一个函数,不各写一遍。""" request = govdoc_scenario.build_run_request( task=task, phase=CRASH_PHASE, run_id=run_id, workspace=workspace, model_binding=CRASH_MODEL_BINDING, ) return replace(request, budget=CRASH_BUDGET) def build_crash_definition( *, model_client: object, store: JsonlRunStore, sink: JsonlEventSink ) -> AgentDefinition: return AgentDefinition( model_client=model_client, # type: ignore[arg-type] decision_parser=GovDocParser(), store=store, event_sink=sink, synthetic_observations=govdoc_scenario.SYNTHETIC_OBSERVATIONS, ) async def check_parameter_drift_detected( *, definition: AgentDefinition, request: RunRequest ) -> Criterion: """故意改一个预算数字续跑,验库确实拒绝。 这条守的是「续跑不能顺便换配置」。它不成立的后果很具体:前几步与后几步来自两份不同的 配置,而两段轨迹在文件里看起来是同一次运行,事后分不出来。 改的是 `max_steps` **加一**而不是减:万一这道守卫失效、`resume` 真的跑起来了,加一只是 多给一步余量,减到一会当场以预算耗尽收尾、把日志封死,后面真正的续跑就没得做了。 """ drifted = replace( request, budget=replace(request.budget, max_steps=request.budget.max_steps + 1) ) try: await session.resume(definition, drifted) except ParameterDriftError as exc: detail = str(exc).split(":", 1)[0] return passed( "parameter_drift_detected", f"改一个预算数字之后续跑报了 ParameterDriftError({detail})", ) except Exception as exc: # noqa: BLE001 - 任何别的异常都说明守卫走的不是这条路 return breached( "parameter_drift_detected", f"改一个预算数字之后续跑抛的是 {type(exc).__name__},期望 ParameterDriftError", ) return breached( "parameter_drift_detected", "改一个预算数字之后续跑正常返回了结果,参数漂移这道守卫没拦住", ) async def run_crash_fault( *, fault: str, timing: KillTiming, runs_dir: Path, workspace_root: Path, govdoc_db: Path, govdoc_corpus: Path, model_client: object, guard: CallGuard, attempts: int, self_kill_grace_s: float = SELF_KILL_GRACE_S, ) -> FaultReport: """一类崩溃续跑:起子进程 → 它按时机自杀 → 拷字节 → 同一个 run_id 续跑 → 逐条判。""" task = load_govdoc_task(govdoc_db=govdoc_db, govdoc_corpus=govdoc_corpus) notes: list[str] = [] for attempt in range(1, attempts + 1): if not guard.affordable(CRASH_BUDGET.max_steps + CRASH_AFTER_STEPS + 1): notes.append( f"调用数护栏只剩 {guard.remaining} 次,不够再试一轮,停在第 {attempt} 次之前" ) break run_id = f"fault-{fault}-{attempt}" workspace = workspace_root / run_id workspace.mkdir(parents=True, exist_ok=True) (workspace / SEEDED_PLAN_NAME).write_text(SEEDED_PLAN_TEXT, encoding="utf-8") log_path = runs_dir / f"{run_id}.jsonl" argv = [ sys.executable, "-m", "tools.soak.faults", "--child", "--runs-dir", str(runs_dir), "--run-id", run_id, "--workspace", str(workspace), "--timing", timing.value, "--govdoc-db", str(govdoc_db), "--govdoc-corpus", str(govdoc_corpus), ] outcome = await spawn_and_kill( argv=argv, log_path=log_path, workspace=workspace, timing=timing, after_steps=CRASH_AFTER_STEPS, child_log_path=runs_dir / f"{run_id}.child.log", self_kill_grace_s=self_kill_grace_s, ) guard.charge(outcome.model_calls) if not outcome.hit: audit = read_audit_lines(workspace) why_empty = ( ";这次跑到结束都没调过一次 write_note,审计账是空的,所以自杀条件从来没满足过" if not audit else "" ) notes.append( f"第 {attempt} 次没崩在时机上:{outcome.reason}" f"(花了 {outcome.model_calls} 次调用){why_empty}" ) write_sidecars( runs_dir=runs_dir, run_id=run_id, scenario="govdoc", fault=f"{fault}_missed", result=None, wall_ms=0, model_calls=outcome.model_calls, sink_failures=0, env_executions=len(read_audit_lines(workspace)), task_id=task.checkpoint.checkpoint_id, phase=CRASH_PHASE, resumed_from_step=None, ) continue notes.append( f"第 {attempt} 次{outcome.reason},崩溃时 {outcome.steps} 步、" f"审计账 {len(read_audit_lines(workspace))} 条" ) return await _resume_and_judge( fault=fault, timing=timing, task=task, run_id=run_id, runs_dir=runs_dir, workspace=workspace, log_path=log_path, outcome=outcome, model_client=model_client, guard=guard, notes=notes, ) return FaultReport( fault=fault, criteria=( undetermined( "crash_timing_hit", f"{attempts} 次都没能让子进程崩在时机 {timing.value} 上,这一类什么都没验成", ), ), notes=tuple(notes), ) async def _resume_and_judge( *, fault: str, timing: KillTiming, task: AuditTask, run_id: str, runs_dir: Path, workspace: Path, log_path: Path, outcome: KillOutcome, model_client: object, guard: CallGuard, notes: list[str], ) -> FaultReport: audit_before = read_audit_lines(workspace) # 崩溃那一刻的日志原样留一份在盘上,给人事后自己比。**后缀不是 `.jsonl`**:记分板按 # `*.jsonl` 枚举 run,叫那个名字的话这份快照会被当成另一次运行。 (runs_dir / f"{run_id}{CRASH_SNAPSHOT_SUFFIX}").write_bytes(outcome.snapshot) store = JsonlRunStore(directory=runs_dir) sink = JsonlEventSink(runs_dir / f"{run_id}.events.jsonl") definition = build_crash_definition(model_client=model_client, store=store, sink=sink) request = build_crash_request(task=task, run_id=run_id, workspace=workspace) criteria: list[Criterion] = [ await check_parameter_drift_detected(definition=definition, request=request) ] started = time.monotonic() result: RunResult | None = None try: result = await session.resume(definition, request) except Exception as exc: # noqa: BLE001 - 续跑本身炸了也是一条要报告的判定 criteria.append(breached("resume_completed", f"续跑抛了 {type(exc).__name__}:{exc}")) else: criteria.append( passed("resume_completed", f"续跑正常结束,停止原因 {result.stop_reason.value}") ) wall_ms = int((time.monotonic() - started) * 1000) final_bytes = log_path.read_bytes() if log_path.is_file() else b"" read = parse_terminated(final_bytes) audit_after = read_audit_lines(workspace) if read.torn: notes.append("最终日志末尾有一段没被换行终结的字节:那次写没有被确认过,按契约不算数") criteria.append(check_log_readable(read)) criteria.append(check_crash_prefix_preserved(crashed=outcome.snapshot, final=final_bytes)) criteria.append(check_step_indices_dense(read)) criteria.append(check_intents_settled(read)) criteria.append(check_never_action_not_replayed(audit_after)) criteria.append(check_audit_prefix_preserved(before=audit_before, after=audit_after)) criteria.append(check_no_cross_segment_replay(before=audit_before, after=audit_after)) if timing is KillTiming.AFTER_STEP: criteria.append( check_resume_made_progress( crashed_steps=outcome.steps, final_steps=len(tagged(read, "step_completed")) ) ) else: criteria.append(check_stop_reason(read, "resume_state_unknown")) criteria.append(check_audit_unchanged(before=audit_before, after=audit_after)) resumed_calls = count_model_calls(read) - outcome.model_calls guard.charge(max(0, resumed_calls)) write_sidecars( runs_dir=runs_dir, run_id=run_id, scenario="govdoc", fault=fault, result=result, wall_ms=wall_ms, model_calls=count_model_calls(read), sink_failures=sink.failures, env_executions=len(audit_after), task_id=task.checkpoint.checkpoint_id, phase=CRASH_PHASE, resumed_from_step=outcome.steps, ) return FaultReport(fault=fault, criteria=tuple(criteria), notes=tuple(notes)) async def run_context_overflow_fault( *, runs_dir: Path, workspace_root: Path, govdoc_db: Path, govdoc_corpus: Path, model_client: object, guard: CallGuard, ) -> FaultReport: """撞提示词上限:把上限压到刚好等于初始提示词,验库终止而不是静默截断历史。 **用 GovDoc 不用 AppWorld**:它的观察是文档片段,一步就能把提示词撑过去,而且上下文本身 只有三千字符上下,上限压到那个量级之后第一步仍然过得去。 工作区里照崩溃续跑那一路预置同一份计划,理由也相同:execute 阶段的提示词让模型先读 plan.md 再动笔,没有它第一步会浪费在一次读不到文件的工具调用上。这一路只走一步,但那一步的观察越 像真的越好——它正是把提示词撑过上限的那一段。 """ fault = "context_overflow" run_id = f"fault-{fault}" log_path = runs_dir / f"{run_id}.jsonl" task = load_govdoc_task(govdoc_db=govdoc_db, govdoc_corpus=govdoc_corpus) workspace = workspace_root / run_id workspace.mkdir(parents=True, exist_ok=True) (workspace / SEEDED_PLAN_NAME).write_text(SEEDED_PLAN_TEXT, encoding="utf-8") # 绑定沿用崩溃续跑那一份:它记的是「哪个场景、哪个阶段」,而这一路两样都相同。 base = govdoc_scenario.build_run_request( task=task, phase=CRASH_PHASE, run_id=run_id, workspace=workspace, model_binding=CRASH_MODEL_BINDING, ) budget = build_context_overflow_budget(base) request = replace(base, budget=budget) store = JsonlRunStore(directory=runs_dir) sink = JsonlEventSink(runs_dir / f"{run_id}.events.jsonl") # 提示词那几条判据全部读库自己写下的 `prompt_chars`。包一层把真正发出去的规模记下来, # 跑完对账——不对账的话,一次静默截断加一列伪造的自述值可以让它们全绿。 recorder = PromptSizeRecordingClient(inner=model_client) definition = build_crash_definition(model_client=recorder, store=store, sink=sink) started = time.monotonic() result = await session.run(definition, request) wall_ms = int((time.monotonic() - started) * 1000) read = read_log(log_path) guard.charge(count_model_calls(read)) criteria = ( check_log_readable(read), check_run_finished_present(read), check_stop_reason(read, "context_overflow"), check_at_least_one_step(read), check_prompt_chars_match_what_was_sent(read, observed=recorder.chars_by_call_index), check_prompt_chars_monotonic(read), check_prompt_reached_max_prompt_chars(read), ) notes = ( "提示词上限按这次运行的上下文现算:" f"{budget.max_prompt_chars - CONTEXT_OVERFLOW_SLACK_CHARS} 字符的初始提示词加 " f"{CONTEXT_OVERFLOW_SLACK_CHARS} 字符余量,一共 {budget.max_prompt_chars}", ) write_sidecars( runs_dir=runs_dir, run_id=run_id, scenario="govdoc", fault=fault, result=result, wall_ms=wall_ms, model_calls=count_model_calls(read), sink_failures=sink.failures, env_executions=len(read_audit_lines(workspace)), task_id=task.checkpoint.checkpoint_id, phase=CRASH_PHASE, resumed_from_step=None, ) return FaultReport(fault=fault, criteria=criteria, notes=notes) # --------------------------------------------------------------------------- # 九、AppWorld:取消、预算、解析连击、环境故障 # --------------------------------------------------------------------------- #: 撞步数上限那一路的预算。**动作上限必须比步数上限宽**,否则先撞上的是动作那一维,停止 #: 原因就成了 `action_budget`,这一条什么都没验到。 STEP_BUDGET_OVERRIDE = Budget( max_steps=3, max_actions=40, max_consecutive_parse_failures=3, max_prompt_chars=400_000 ) #: 撞动作上限那一路:动作上限压到比步数上限小,让动作那一维先耗尽。 ACTION_BUDGET_OVERRIDE = Budget( max_steps=10, max_actions=2, max_consecutive_parse_failures=3, max_prompt_chars=400_000 ) #: 解析失败连击那一路。步数上限留得比连击上限宽,让连击那一维先命中。 PARSE_FAILURE_BUDGET = Budget( max_steps=10, max_actions=10, max_consecutive_parse_failures=3, max_prompt_chars=400_000 ) #: 环境故障那一路的预算。四步是给「走一步 → 打死容器 → 第二步撞环境故障」留的余量:正常只用 #: 两步,多出来的两步是留给一次解析失手的。 ENV_ERROR_BUDGET = Budget( max_steps=4, max_actions=4, max_consecutive_parse_failures=3, max_prompt_chars=400_000 ) #: 完整执行过几次动作之后才把容器打死。见 `should_break_env`:取 0 压到的是初始化不是执行接缝。 ENV_ERROR_BREAK_AFTER_EXECUTIONS = 1 #: 环境故障那一路自己那个容器池的起始端口。 #: #: **它必须自己起一个池,不能用其余几类共用的那个。** 那个池只有一个容器(容器租约那条判据 #: 要求如此),而这一类会把容器打死;共用的话,排在它后面的每一类都会跑在一个已经没了的环境 #: 上,而报出来是一串「连不上」,看起来像 docker 挂了。 #: #: 端口取共用池那一个的下一个:共用池大小是 1,只占 `DEFAULT_PORT_BASE` 本身。 ENV_ERROR_PORT_BASE = DEFAULT_PORT_BASE + 1 #: 取消之后再借一个容器的等待上限。池满时 `lease()` 会一直阻塞,所以超时就是击穿。 LEASE_TIMEOUT_S = 120.0 #: 等取消触发器的上限。等不到说明这次运行在触发点之前就结束了。 TRIGGER_TIMEOUT_S = 300.0 #: 取消回来之后再等几秒,然后重数一次环境侧的账。见 `check_env_quiet_after_cancel`:立刻就数 #: 的话两个数只是同一个瞬间的两份拷贝,什么都验不到。三秒盖得住一次已经发出去的 `/execute` #: 在容器里跑完并返回——AppWorld 那侧单次执行的超时是一百秒,所以盖不住最慢的那一档;这一路 #: 的动作都是简短的 API 调用,实测在秒级以内。 POST_CANCEL_SETTLE_S = 3.0 def _appworld_definition( *, model_client: object, store: JsonlRunStore, sink: JsonlEventSink, parser: object ) -> AgentDefinition: return AgentDefinition( model_client=model_client, # type: ignore[arg-type] decision_parser=parser, # type: ignore[arg-type] store=store, event_sink=sink, synthetic_observations=appworld_scenario.build_synthetic_observations(), ) async def _can_borrow(pool: AppWorldPool, task_id: str, timeout_s: float) -> bool: """再借一次容器,借得到返回 True。借不到(超时)说明上一次的租约没还回来。""" try: async with asyncio.timeout(timeout_s), pool.session(task_id): return True except TimeoutError: return False async def run_cancel_fault( *, fault: str, pool: AppWorldPool, task_id: str, app_descriptions: str, runs_dir: Path, model_client: object, guard: CallGuard, ) -> FaultReport: """取消:跑到中途 `task.cancel()`,验异常穿透、结束记录留痕、容器租约归还。""" run_id = f"fault-{fault}" log_path = runs_dir / f"{run_id}.jsonl" store = JsonlRunStore(directory=runs_dir) sink = JsonlEventSink(runs_dir / f"{run_id}.events.jsonl") trigger = asyncio.Event() notes: list[str] = [] raised: BaseException | None = None env_executions = 0 quiet: Criterion | None = None started = time.monotonic() async with pool.session(task_id) as handle: request = appworld_scenario.build_run_request( run_id=run_id, session=handle, app_descriptions=app_descriptions, model_binding=APPWORLD_MODEL_BINDING, ) probe: TriggeringExecutor | None = None if fault == "cancel_env": probe = TriggeringExecutor(inner=request.action_executor, trigger=trigger) request = replace(request, action_executor=probe) # type: ignore[arg-type] client: object = model_client else: # 第 1 次(从 0 数起的第二次)调用时触发:让日志里先有一步完整的记录,取消才落在 # 「跑到中途」而不是「刚起步」。 client = TriggeringModelClient(inner=model_client, trigger=trigger, at_call_index=1) definition = _appworld_definition( model_client=client, store=store, sink=sink, parser=AppWorldParser() ) task = asyncio.create_task(session.run(definition, request)) try: async with asyncio.timeout(TRIGGER_TIMEOUT_S): await trigger.wait() except TimeoutError: notes.append(f"等了 {TRIGGER_TIMEOUT_S} 秒也没等到触发点,这次运行可能在触发前就结束了") task.cancel() # 这里接住的是**被等待的那个任务**抛出来的取消,不是本协程自己的取消——本协程从头到尾 # 没有被 cancel 过。接住它正是这一类要验的那条判据(CLAUDE.md §1.6 禁的是把自己的取消 # 吞掉)。 try: await task except asyncio.CancelledError as exc: raised = exc env_executions = handle.n_executions if probe is not None: # 取消刚回来就数一次,等一小会儿再数一次。在飞的那次执行要是没被真正掐断,它会在 # 这段等待里跑完并把计数推上去——立刻就数的话,两个数只是同一个瞬间的两份拷贝。 dispatched_before, completed_before = probe.executions, handle.n_executions await asyncio.sleep(POST_CANCEL_SETTLE_S) quiet = check_env_quiet_after_cancel( dispatched_before=dispatched_before, dispatched_after=probe.executions, completed_before=completed_before, completed_after=handle.n_executions, waited_s=POST_CANCEL_SETTLE_S, ) env_executions = handle.n_executions wall_ms = int((time.monotonic() - started) * 1000) borrowed = await _can_borrow(pool, task_id, LEASE_TIMEOUT_S) read = read_log(log_path) guard.charge(count_model_calls(read)) criteria = ( check_log_readable(read), check_cancelled_raised(raised=raised), check_stop_reason(read, "cancelled"), check_events_file_intact(runs_dir / f"{run_id}.events.jsonl"), check_lease_returned(borrowed=borrowed, timeout_s=LEASE_TIMEOUT_S, pool_size=1), *((quiet,) if quiet is not None else ()), ) write_sidecars( runs_dir=runs_dir, run_id=run_id, scenario="appworld", fault=fault, # 取消那条路上 `run` 不返回结果,`.result.json` 天然缺;记分板会退回日志里的结束记录。 result=None, wall_ms=wall_ms, model_calls=count_model_calls(read), sink_failures=sink.failures, env_executions=env_executions, task_id=task_id, phase=None, resumed_from_step=None, ) return FaultReport(fault=fault, criteria=criteria, notes=tuple(notes)) async def run_budget_fault( *, fault: str, pool: AppWorldPool, task_id: str, app_descriptions: str, runs_dir: Path, model_client: object, guard: CallGuard, ) -> FaultReport: """不可能完成的目标:验它干净地撞预算上限,而不是以别的原因结束。""" run_id = f"fault-{fault}" log_path = runs_dir / f"{run_id}.jsonl" store = JsonlRunStore(directory=runs_dir) sink = JsonlEventSink(runs_dir / f"{run_id}.events.jsonl") budget = STEP_BUDGET_OVERRIDE if fault == "step_budget" else ACTION_BUDGET_OVERRIDE notes: list[str] = [] result: RunResult | None = None env_executions = 0 started = time.monotonic() async with pool.session(task_id) as handle: request = replace( appworld_scenario.build_run_request( run_id=run_id, session=handle, app_descriptions=app_descriptions, model_binding=APPWORLD_MODEL_BINDING, ), budget=budget, ) definition = _appworld_definition( model_client=model_client, store=store, sink=sink, parser=AppWorldParser() ) result = await session.run(definition, request) env_executions = handle.n_executions wall_ms = int((time.monotonic() - started) * 1000) read = read_log(log_path) guard.charge(count_model_calls(read)) criteria = [check_log_readable(read), check_stop_reason(read, fault)] if fault == "step_budget": criteria.append( check_step_count(read, expected=budget.max_steps, name="steps_equal_max_steps") ) else: criteria.append(check_executed_action_count(read, expected=budget.max_actions)) criteria.append(check_no_env_error_step(read)) write_sidecars( runs_dir=runs_dir, run_id=run_id, scenario="appworld", fault=fault, result=result, wall_ms=wall_ms, model_calls=count_model_calls(read), sink_failures=sink.failures, env_executions=env_executions, task_id=task_id, phase=None, resumed_from_step=None, ) return FaultReport(fault=fault, criteria=tuple(criteria), notes=tuple(notes)) async def run_parse_failure_fault( *, pool: AppWorldPool, task_id: str, app_descriptions: str, runs_dir: Path, model_client: object, guard: CallGuard, ) -> FaultReport: """解析失败连击:包一层永远判无效的解释器,验它按连击上限收尾且一次都没碰环境。 这一类只花三次调用,但它守的是「解析失败不调环境」这条契约——那条错了的表现是环境状态 被一批根本没解释出来的动作改掉,而轨迹里每一步都写着解析失败。 """ fault = "parse_failures" run_id = f"fault-{fault}" log_path = runs_dir / f"{run_id}.jsonl" store = JsonlRunStore(directory=runs_dir) sink = JsonlEventSink(runs_dir / f"{run_id}.events.jsonl") result: RunResult | None = None env_executions = 0 started = time.monotonic() async with pool.session(task_id) as handle: request = replace( appworld_scenario.build_run_request( run_id=run_id, session=handle, app_descriptions=app_descriptions, model_binding=APPWORLD_MODEL_BINDING, ), budget=PARSE_FAILURE_BUDGET, ) definition = _appworld_definition( model_client=model_client, store=store, sink=sink, parser=AlwaysInvalidParser() ) result = await session.run(definition, request) env_executions = handle.n_executions wall_ms = int((time.monotonic() - started) * 1000) read = read_log(log_path) guard.charge(count_model_calls(read)) criteria = ( check_log_readable(read), check_stop_reason(read, "parse_failed_repeatedly"), check_step_count( read, expected=PARSE_FAILURE_BUDGET.max_consecutive_parse_failures, name="steps_equal_max_consecutive_parse_failures", ), check_all_steps_parse_failed(read), check_env_untouched(executions=env_executions, source="AppWorld 的 n_executions"), ) write_sidecars( runs_dir=runs_dir, run_id=run_id, scenario="appworld", fault=fault, result=result, wall_ms=wall_ms, model_calls=count_model_calls(read), sink_failures=sink.failures, env_executions=env_executions, task_id=task_id, phase=None, resumed_from_step=None, ) return FaultReport(fault=fault, criteria=criteria, notes=()) async def run_env_error_fault( *, data_root: Path, split: str, runs_dir: Path, model_client: object, guard: CallGuard, ) -> FaultReport: """环境故障:正常跑一步,把容器打死,验库换成合成观察并以 `env_error` 收尾。 **容器是真的被杀掉的,不是一个返回 `ENV_ERROR` 的替身。** 替身验不到「环境真的不在了 之后,在飞的连接、会话关闭、容器池收尾这一串还能不能干净地收场」,而那串正是这一类跑完 要留下的东西。 收尾这件事实测过一遍:容器被 `docker kill` 之后,`pool.session()` 退出时的 `/close` 会失败 ——那由 `AppWorldPool._close_quietly` 记一次账、不抛(连续三次才抛),所以这一类不会把整轮 压测带走;随后 `pool.stop()` 里的 `docker rm -f` 撞上一个已经删掉的容器,退出码非零而它本来 就不看退出码。两处都不需要改 `tools/soak/appworld.py`。 """ fault = "env_error" run_id = f"fault-{fault}" log_path = runs_dir / f"{run_id}.jsonl" store = JsonlRunStore(directory=runs_dir) sink = JsonlEventSink(runs_dir / f"{run_id}.events.jsonl") synthetic = appworld_scenario.build_synthetic_observations() notes: list[str] = [] result: RunResult | None = None env_executions = 0 started = time.monotonic() async with AppWorldPool(data_root=data_root, size=1, port_base=ENV_ERROR_PORT_BASE) as pool: task_ids = pool.list_task_ids(split) if not task_ids: raise FaultInjectionError(f"{split} 划分里一道题都没有,环境故障那一类跑不了") task_id = task_ids[0] app_descriptions = await appworld_scenario.load_app_descriptions(pool, task_id=task_id) container = container_name_for_port(pool.ports[0]) async with pool.session(task_id) as handle: request = replace( appworld_scenario.build_run_request( run_id=run_id, session=handle, app_descriptions=app_descriptions, model_binding=APPWORLD_MODEL_BINDING, ), budget=ENV_ERROR_BUDGET, ) breaker = EnvBreakingExecutor( inner=request.action_executor, break_after_executions=ENV_ERROR_BREAK_AFTER_EXECUTIONS, break_env=lambda: kill_and_remove_container(container), ) request = replace(request, action_executor=breaker) # type: ignore[arg-type] definition = _appworld_definition( model_client=model_client, store=store, sink=sink, parser=AppWorldParser() ) result = await session.run(definition, request) env_executions = handle.n_executions wall_ms = int((time.monotonic() - started) * 1000) if breaker.break_note: notes.append(breaker.break_note) read = read_log(log_path) guard.charge(count_model_calls(read)) criteria = ( check_log_readable(read), check_run_finished_present(read), check_stop_reason(read, "env_error"), check_last_step_action_status(read, expected="env_error"), check_last_observation_is_synthetic( read, expected=synthetic.env_failed, name="env_failed_observation_substituted" ), check_steps_before_last_all_executed(read), check_env_broken_after_a_full_step( broken=breaker.broken, executions_before_break=breaker.executions_before_break, expected_after=ENV_ERROR_BREAK_AFTER_EXECUTIONS, ), ) write_sidecars( runs_dir=runs_dir, run_id=run_id, scenario="appworld", fault=fault, result=result, wall_ms=wall_ms, model_calls=count_model_calls(read), sink_failures=sink.failures, # 环境侧自己数的执行次数。**打死容器之后那一次不计入**:`AppWorldSession` 在请求返回 # 之后才加一,而那一次请求根本没到过环境。 env_executions=env_executions, task_id=task_id, phase=None, resumed_from_step=None, ) return FaultReport(fault=fault, criteria=criteria, notes=tuple(notes)) # --------------------------------------------------------------------------- # 十、子进程入口 # --------------------------------------------------------------------------- async def run_child(args: argparse.Namespace) -> int: """子进程:跑一次 GovDoc execute 阶段,跑到指定时机就把自己打死。 **打死自己这件事由 `SelfKillingStore` 在写入落盘之后做**,不在这里判——判据要贴着那次写 才准,隔一层就又变成抢窗口了。它调的是 `os._exit`,所以这个函数在命中时机时根本不返回, 下面的 `client.aclose()` 与 `finally` 都不会跑。**那正是要的**:崩溃就是不给清理机会。 正常跑完也是允许的(时机一次都没触发),那时父进程会报「没崩在时机上」并重试。 **事件写在 `.child.events.jsonl`,不写记分板认的那个名字。** 事件只在「本进程里 真的走完」的步上发,而记分板按 `meta.resumed_from_step` 把续跑跳过的那一段减掉之后与 `.events.jsonl` 的条数对账。子进程的事件混进同一个文件的话,那笔账永远对不上,而对不上 的原因是我们把两个进程的事件叠在了一起,不是库出了问题。这个文件名以 `.events.jsonl` 结尾,所以记分板枚举 run 时会把它排掉,不会被当成另一次运行。 """ from polygateway import GatewayClient, GatewaySettings runs_dir = Path(args.runs_dir) workspace = Path(args.workspace) task = load_govdoc_task(govdoc_db=Path(args.govdoc_db), govdoc_corpus=Path(args.govdoc_corpus)) store = SelfKillingStore( inner=JsonlRunStore(directory=runs_dir), timing=KillTiming(args.timing), workspace=workspace, ) sink = JsonlEventSink(runs_dir / f"{args.run_id}.child.events.jsonl") client = GatewayClient.from_env() try: definition = build_crash_definition( model_client=GatewayModelClient(client=client, settings=GatewaySettings.from_env()), store=store, sink=sink, ) request = build_crash_request(task=task, run_id=args.run_id, workspace=workspace) await session.run(definition, request) finally: await client.aclose() return 0 # --------------------------------------------------------------------------- # 十一、入口 # --------------------------------------------------------------------------- def build_parser() -> argparse.ArgumentParser: parser = argparse.ArgumentParser( prog="python -m tools.soak.faults", description=( "PolyLoop 压测的故障注入:崩溃续跑、取消、撞预算、解析失败连击、撞提示词上限、环境故障" ), ) parser.add_argument("--runs-dir", required=True, help="产物目录:日志与三种 sidecar 都写在这里") parser.add_argument( "--fault", action="append", choices=FAULT_NAMES, help="要跑哪一类,可给多次。不给就全跑", ) parser.add_argument( "--budget-calls", type=int, help="真实模型调用数的护栏。每类故障开跑前检查一次,不够就跳过(父进程必填)", ) parser.add_argument("--data-root", help="AppWorld 数据根目录(跑 AppWorld 那几类时必填)") parser.add_argument("--govdoc-db", help="GovDoc 审核点 sqlite(跑崩溃续跑时必填)") parser.add_argument("--govdoc-corpus", help="GovDoc 脱敏前语料目录(跑崩溃续跑时必填)") parser.add_argument("--split", default="train", help="AppWorld 的任务划分,默认 train") parser.add_argument( "--attempts", type=int, default=3, help="崩溃续跑每类最多重试几次去命中时机,默认 3" ) parser.add_argument( "--self-kill-grace", type=float, default=SELF_KILL_GRACE_S, help=f"条件对上之后留给子进程自杀的秒数,超了才兜底发 SIGKILL,默认 {SELF_KILL_GRACE_S}", ) parser.add_argument( "--child", action="store_true", help="内部用:以子进程身份跑一次 GovDoc 运行" ) parser.add_argument("--run-id", help="内部用:子进程的运行标识") parser.add_argument("--workspace", help="内部用:子进程的工作区目录") parser.add_argument( "--timing", choices=[item.value for item in KillTiming], help="内部用:子进程在哪一次写入之后把自己打死", ) return parser def _require(args: argparse.Namespace, name: str, why: str) -> Path: value = getattr(args, name.replace("-", "_")) if not value: raise FaultInjectionError(f"缺 --{name}:{why}") path = Path(value) if not path.exists(): raise FaultInjectionError(f"--{name} 指向的东西不在:{path}") return path async def guarded(fault: str, coroutine: Awaitable[FaultReport]) -> FaultReport: """跑一类故障,装配或环境本身炸了就记成「无法判定」接着跑下一类。 **这不是吞错误**:异常的类名与文本原样进了那条判据的证据里,退出码那侧也看得见它不是 「通过」。不这么做的话,第五类故障在起容器时崩一下,前面四类已经花钱跑出来的结论会一起 丢掉,而那些结论正是这次跑的产出。 `CancelledError` 继承 `BaseException`,不在 `except Exception` 的范围里,所以人按 Ctrl-C 时整个跑照样立刻停下(CLAUDE.md §1.6)。 """ try: return await coroutine except Exception as exc: # noqa: BLE001 - 见 docstring:记下来,不吞掉 return FaultReport( fault=fault, criteria=( undetermined( "fault_ran_to_completion", f"这一类跑到一半抛了 {type(exc).__name__}:{exc}", ), ), ) async def run_all(args: argparse.Namespace) -> int: from polygateway import GatewayClient, GatewaySettings selected = tuple(args.fault) if args.fault else FAULT_NAMES runs_dir = Path(args.runs_dir) runs_dir.mkdir(parents=True, exist_ok=True) guard = CallGuard(limit=args.budget_calls) reports: list[FaultReport] = [] client = GatewayClient.from_env() try: model_client = GatewayModelClient(client=client, settings=GatewaySettings.from_env()) govdoc_selected = [name for name in selected if name in GOVDOC_FAULTS] if govdoc_selected: govdoc_db = _require(args, "govdoc-db", "GovDoc 那几类要读审核点") govdoc_corpus = _require(args, "govdoc-corpus", "GovDoc 那几类要读语料") for name in govdoc_selected: coroutine: Awaitable[FaultReport] if name == "context_overflow": if not guard.affordable(CONTEXT_OVERFLOW_BUDGET_STEPS): reports.append( _unaffordable( name, guard=guard, worst_case=CONTEXT_OVERFLOW_BUDGET_STEPS ) ) continue coroutine = run_context_overflow_fault( runs_dir=runs_dir, workspace_root=runs_dir / "workspaces", govdoc_db=govdoc_db, govdoc_corpus=govdoc_corpus, model_client=model_client, guard=guard, ) else: # 崩溃续跑那两类自己按轮次查护栏(每重试一轮都要再问一次),不在这里查。 timing = ( KillTiming.AFTER_STEP if name == "crash_resume_a" else KillTiming.AT_INTENT ) coroutine = run_crash_fault( fault=name, timing=timing, runs_dir=runs_dir, workspace_root=runs_dir / "workspaces", govdoc_db=govdoc_db, govdoc_corpus=govdoc_corpus, model_client=model_client, guard=guard, attempts=args.attempts, self_kill_grace_s=args.self_kill_grace, ) reports.append(await guarded(name, coroutine)) appworld_selected = [name for name in selected if name in APPWORLD_FAULTS] if appworld_selected: data_root = _require(args, "data-root", "AppWorld 那几类要起容器") reports.extend( await _run_appworld_faults( names=appworld_selected, data_root=data_root, split=args.split, runs_dir=runs_dir, model_client=model_client, guard=guard, ) ) finally: await client.aclose() print("\n".join(report.render() for report in reports)) print(f"\n真实模型调用:{guard.spent} / {guard.limit}") breaches = [item for report in reports for item in report.breaches] unknowns = [item for report in reports for item in report.undetermineds] print(f"击穿 {len(breaches)} 条,无法判定 {len(unknowns)} 条") return 1 if breaches else 0 #: 每一类 AppWorld 故障最坏会花掉几次真实模型调用。护栏按这个数在每类开跑前问一次。 _APPWORLD_WORST_CASE_CALLS: Mapping[str, int] = { "cancel_model": 3, "cancel_env": 3, "step_budget": STEP_BUDGET_OVERRIDE.max_steps, "action_budget": ACTION_BUDGET_OVERRIDE.max_steps, "parse_failures": PARSE_FAILURE_BUDGET.max_consecutive_parse_failures, "env_error": ENV_ERROR_BUDGET.max_steps, } def _unaffordable(fault: str, *, guard: CallGuard, worst_case: int) -> FaultReport: """护栏不够跑这一类了。**报「无法判定」不报「通过」**:它一次都没跑过。""" return FaultReport( fault=fault, criteria=( undetermined( "call_budget_available", f"调用数护栏只剩 {guard.remaining} 次,最坏要 {worst_case} 次,这一类没跑", ), ), ) async def _run_appworld_faults( *, names: Sequence[str], data_root: Path, split: str, runs_dir: Path, model_client: object, guard: CallGuard, ) -> list[FaultReport]: """AppWorld 那几类共用一个池,环境故障那一类除外。 **池大小固定为 1**:容器租约那条判据要的是「取消之后空闲名额回到满」,而池里有第二个 容器的话,借得到只说明还剩别的名额,什么都证明不了。 **环境故障那一类自己起一个池**(见 `ENV_ERROR_PORT_BASE`):它会把容器打死,共用的话排在 它后面的每一类都会跑在一个已经没了的环境上。所以共用池只在真有别的类要跑时才起——只选了 环境故障那一类时,起它纯属白等一次容器启动。 """ reports: list[FaultReport] = [] shared_names = [name for name in names if name != "env_error"] async with contextlib.AsyncExitStack() as stack: pool: AppWorldPool | None = None task_id = "" app_descriptions = "" if shared_names: pool = await stack.enter_async_context(AppWorldPool(data_root=data_root, size=1)) task_ids = pool.list_task_ids(split) if not task_ids: raise FaultInjectionError(f"{split} 划分里一道题都没有,AppWorld 那几类跑不了") task_id = task_ids[0] app_descriptions = await appworld_scenario.load_app_descriptions(pool, task_id=task_id) for name in names: worst_case = _APPWORLD_WORST_CASE_CALLS[name] if not guard.affordable(worst_case): reports.append(_unaffordable(name, guard=guard, worst_case=worst_case)) continue coroutine: Awaitable[FaultReport] if name == "env_error": coroutine = run_env_error_fault( data_root=data_root, split=split, runs_dir=runs_dir, model_client=model_client, guard=guard, ) elif pool is None: # pragma: no cover - 共用池只在有别的类要跑时才建,走不到这里 raise FaultInjectionError(f"{name} 要用共用容器池,而它没有被建起来") elif name in {"cancel_model", "cancel_env"}: coroutine = run_cancel_fault( fault=name, pool=pool, task_id=task_id, app_descriptions=app_descriptions, runs_dir=runs_dir, model_client=model_client, guard=guard, ) elif name == "parse_failures": coroutine = run_parse_failure_fault( pool=pool, task_id=task_id, app_descriptions=app_descriptions, runs_dir=runs_dir, model_client=model_client, guard=guard, ) else: coroutine = run_budget_fault( fault=name, pool=pool, task_id=task_id, app_descriptions=app_descriptions, runs_dir=runs_dir, model_client=model_client, guard=guard, ) reports.append(await guarded(name, coroutine)) return reports def main(argv: Sequence[str] | None = None) -> int: parser = build_parser() args = parser.parse_args(argv) if args.child: missing = [ name for name in ("run_id", "workspace", "timing", "govdoc_db", "govdoc_corpus") if not getattr(args, name) ] if missing: parser.error(f"--child 模式还缺 {missing}") return asyncio.run(run_child(args)) if args.budget_calls is None: parser.error("--budget-calls 是必填的:没有它这一跑可以无上限地花钱") if args.budget_calls < 1: parser.error("--budget-calls 至少是 1") return asyncio.run(run_all(args)) if __name__ == "__main__": raise SystemExit(main()) __all__ = [ "APPWORLD_FAULTS", "CONTEXT_OVERFLOW_BUDGET_STEPS", "CONTEXT_OVERFLOW_SLACK_CHARS", "CRASH_BUDGET", "CRASH_EXIT_CODE", "ENV_ERROR_BREAK_AFTER_EXECUTIONS", "ENV_ERROR_PORT_BASE", "FAULT_NAMES", "GOVDOC_FAULTS", "MAX_PROMPT_CHARS_SNAPSHOT_KEY", "OBSERVATION_TEMPLATE_SNAPSHOT_KEY", "SELF_KILL_GRACE_S", "AlwaysInvalidParser", "CallGuard", "Criterion", "CriterionStatus", "EnvBreakingExecutor", "FaultInjectionError", "FaultReport", "JsonlEventSink", "PromptSizeRecordingClient", "KillOutcome", "KillTiming", "LogRead", "SelfKillingStore", "build_context_overflow_budget", "check_all_steps_parse_failed", "check_at_least_one_step", "check_audit_prefix_preserved", "check_audit_unchanged", "check_cancelled_raised", "check_crash_prefix_preserved", "check_env_broken_after_a_full_step", "check_env_quiet_after_cancel", "check_env_untouched", "check_events_file_intact", "check_executed_action_count", "check_intents_settled", "check_last_observation_is_synthetic", "check_last_step_action_status", "check_lease_returned", "check_log_readable", "check_never_action_not_replayed", "check_no_cross_segment_replay", "check_no_env_error_step", "check_prompt_chars_monotonic", "check_prompt_reached_max_prompt_chars", "check_prompt_chars_match_what_was_sent", "check_resume_made_progress", "check_run_finished_present", "check_step_count", "check_step_indices_dense", "check_steps_before_last_all_executed", "check_stop_reason", "container_name_for_port", "count_model_calls", "initial_prompt_chars", "kill_and_remove_container", "main", "parameter_snapshot_of", "parse_audit_line", "parse_terminated", "read_log", "should_break_env", "should_kill", "stop_reason_of", "tagged", "terminated_prefix", ]