第 18 章
第 3 周:不调真实模型,怎么测一个 Agent?
读 Codex 的测试工程:wiremock 起假模型服务器、按顺序编排 SSE 剧本、断言 agent 发出的下一次请求;stream_no_completed、retry_after、websocket_fallback 三个故障注入测试各测什么;insta 快照、clippy 禁止 unwrap、CI required 汇总门禁和 PR 评审规范。最后用 Python 标准库写一个 Claude 风格的假模型服务器,给第 2 周的 agent loop 写 6 个测试并在本地跑通。
第 2 周从零写的 agent loop,有重试、有停止条件,但怎么证明它们是对的?用真实模型跑,输出每次不一样,跑一次花一次钱,而且你没法让线上 API “在第 5 个事件断开”。Codex 的做法是把模型当成一个可以编排的 HTTP 依赖:起一个假的模型服务器,按剧本吐 SSE 事件流,故障也写进剧本,然后断言 agent 发给模型的下一次请求,而不是断言模型说了什么。这一章逐段读它的测试代码,再用 Python 标准库复刻一个,给第 2 周的 loop 写测试,在本地真的跑通。
下文引用的源码都基于 openai/codex commit
7993248
(2026-10-01),路径相对仓库根目录。
讲解视频
互动演示
一个 SSE 剧本编辑器。任务和下面的 Python 测试一样:模型第 1 轮调 read_file("README.md"),第 2 轮说完。选一种故障注入到第 1 轮(正常、缺 message_stop、中途断开、429 带 Retry-After、畸形 JSON),可以拖动故障出现的位置、调 Retry-After 秒数、选"只出现一次"还是"每次都出现",也可以关掉被测 agent 的某一道防线(检查 message_stop、流层重试、SDK 的 HTTP 重试)。点"下一步"逐个事件回放,右边是 SDK 累积出的消息和 agent 的日志,跑完后下方的测试断言逐条判通过还是失败。agent 的行为按本地代码和 SDK 1.11.0 的实测行为模拟。页面底部有自动判分的练习。
为什么断言请求,而不是断言回答
agent loop 是你的代码,模型是它的外部依赖。测你的代码,就把依赖换成假的、固定住它的行为,再看你的代码对每一种行为反应得对不对。
agent 的反应体现在它发出的下一次请求里:工具结果有没有按 id 回填,重试时有没有把半截消息塞进历史,限流后隔了多久才重发。这些都是确定的,一次就能判对错。
| 断言对象 | 能测什么 | 稳定吗 |
|---|---|---|
| 模型的回答 | 模型质量(属于评测,不属于这里) | 不稳定,要用第 1 周的统计方法 |
| agent 发出的请求 | 工具结果回传、历史拼接、重试次数、退避时间、请求头 | 确定 |
| agent 的最终状态和事件 | 停止条件、错误分类、有没有把失败报成成功 | 确定 |
在 Codex 里,这是集成测试的默认写法:codex-rs/core/tests/suite/ 下有 218 个 .rs 文件,其中 207 个用到了假模型服务器相关的工具(按 start_mock_server、mount_sse、MockServer 等关键字 grep 统计)。
Codex 的假模型服务器:三件工具
测试支持代码在 codex-rs/core/tests/common/responses.rs,底层是 Rust 的 HTTP mock 库 wiremock 0.6(
codex-rs/Cargo.toml:555
)。Codex 调的是 OpenAI 的 Responses API,事件名和 Claude 不同,但套路完全通用。
① sse(events):把一串 JSON 事件拼成 SSE 文本
codex-rs/core/tests/common/responses.rs:752-766
:
/// Build an SSE stream body from a list of JSON events.
pub fn sse(events: Vec<Value>) -> String {
use std::fmt::Write as _;
let mut out = String::new();
for ev in events {
let kind = ev.get("type").and_then(|v| v.as_str()).unwrap();
writeln!(&mut out, "event: {kind}").unwrap();
if !ev.as_object().map(|o| o.len() == 1).unwrap_or(false) {
write!(&mut out, "data: {ev}\n\n").unwrap();
} else {
out.push('\n');
}
}
out
}
逐行看:
Vec<Value>是一组 JSON 值,相当于 Python 的list[dict]。let kind = ev.get("type")...unwrap():取事件的type字段。unwrap()的意思是"假设一定有,没有就直接崩"。这是测试代码,崩了就是测试失败,正合适(后面会讲生产代码为什么禁止这么写)。writeln!(&mut out, "event: {kind}"):写一行event: xxx。if !ev.as_object().map(|o| o.len() == 1)...:只要这个 JSON 对象不止type一个字段,就再写一行data: {整个 JSON}加空行;如果只有type一个字段,就不写 data 行。后面的故障注入测试正是利用这一点造出了"半个事件"。
② mount_sse_sequence:第 n 次请求拿第 n 段剧本
responses.rs:1466-1506
(节选):
pub async fn mount_sse_sequence(server: &MockServer, bodies: Vec<String>) -> ResponseMock {
// ...
fn respond(&self, _: &wiremock::Request) -> ResponseTemplate {
let call_num = self.num_calls.fetch_add(1, Ordering::SeqCst);
let missing_response_message = format!("no response for {call_num}");
let body = self
.responses
.get(call_num)
.expect(&missing_response_message);
ResponseTemplate::new(200)
.insert_header("content-type", "text/event-stream")
.set_body_string(body.clone())
}
// ...
mock.respond_with(responder)
.up_to_n_times(num_calls as u64)
.expect(num_calls as u64)
.mount(server)
.await;
fetch_add(1, ...)是一个原子计数器,相当于线程安全的n += 1,返回加之前的值。第 n 次请求就取bodies[n]。.up_to_n_times(num_calls)让这个 mock 最多匹配 N 次。wiremock 0.6 的MountedMock::matches在已匹配 N 次之后直接返回false(wiremock-0.6.5/src/mounted_mock.rs:43-47),第 N+1 个请求不再命中这个 mock,没有别的 mock 接住时 wiremock 回默认的 404,根本走不到respond,也不会被下面的ResponseMock录下(匹配器都没被调用)。所以.expect("no response for n")实际上碰不到。函数的文档注释写着 “Panics if more requests are received than bodies provided”,和实际行为不符。- 最后的
.expect(num_calls)是 wiremock 的"期望调用次数":测试结束时收到的请求少于剧本段数,测试一定红。
所以请求次数是一条断言,但两个方向不对称:少重试一次一定红;多发一个请求只会拿到 404,测试红不红取决于 agent 怎么处理这个 404(比如当成错误结束 turn,后面的断言才会失败)。下文的 Python 版把多出来的请求记进 unexpected,fixture 收尾时断言它为空,多一个也一定红。
③ ResponseMock:录下每个请求,顺便检查不变量
impl Match for ResponseMock {
fn matches(&self, request: &wiremock::Request) -> bool {
self.requests
.lock()
.unwrap()
.push(ResponsesRequest(request.clone()));
// Enforce invariant checks on every request body captured by the mock.
// Panic on orphan tool outputs or calls to catch regressions early.
validate_request_body_invariants(request);
true
}
}
wiremock 用"匹配器"决定一个请求该由哪个 mock 响应。ResponseMock 把自己挂成匹配器,每来一个请求就存一份(lock().unwrap().push(...):加锁后追加到列表),然后永远返回 true。
存下来的请求在测试结束后可以用 function_call_output_text(call_id) 之类的方法查询,取出"agent 回给模型的工具输出"(
responses.rs:77-82
)。
validate_request_body_invariants(
responses.rs:1558
起)检查请求里的工具调用和工具输出是否一一配对,孤儿输出直接 panic。效果是:任何一个用了这套 mock 的测试,都顺带测了配对规则。这正是第 2 周说的"把 messages 的规则写成 harness 的不变量断言",Codex 把它放进了测试基础设施里,每个测试都免费拿到。
另外,start_mock_server()(
responses.rs:1203-1212
)启动时会先挂一个空的 /models 响应,保证测试完全离线。
一个完整的测试:工具输出有没有原样回到模型
codex-rs/core/tests/suite/tools.rs:681-766
的 sandbox_denied_exec_command_returns_original_output,节选:
let responses = vec![
sse(vec![
ev_response_created("resp-1"),
ev_function_call(call_id, "exec_command", &serde_json::to_string(&args)?),
ev_completed("resp-1"),
]),
sse(vec![
ev_assistant_message("msg-1", "done"),
ev_completed("resp-2"),
]),
];
let mock = mount_sse_sequence(&server, responses).await;
fixture
.submit_turn_with_permission_profile(
"run a command that should be denied by the read-only sandbox",
PermissionProfile::read_only(),
)
.await?;
let output_text = mock
.function_call_output_text(call_id)
.context("shell output present")?;
// ...
assert!(
body.contains(sentinel),
"expected sentinel output from command to reach the model: {body}"
);
三步:
- 排剧本。第一段:模型调
exec_command,命令往只读沙箱外写一个文件,同时往 stderr 打印一个哨兵字符串;第二段:模型说 “done”。 - 用只读权限提交一轮。agent 真的去执行命令,沙箱真的拒绝。
- 从请求记录里取出 agent 回给模型的工具输出,断言:里面有权限拒绝的字样,哨兵字符串到了模型那里,提到了被拒的路径,不是兜底文案,退出码非 0(
:736-763)。
? 和 .context(...) 是 Rust 的错误传播:拿不到就让测试以一条可读的错误失败,而不是 panic。
和第 2 周的 loop 对照,这就是"tool_result 有没有按 id 回填、内容对不对"的测试。Codex 断言的是 Responses API 的 function_call_output,我们换成 Claude 的 tool_result 就行。
故障注入:三个测试文件各测什么
| 文件 | 注入的故障 | 断言 |
|---|---|---|
stream_no_completed.rs | 第一段流只有一个 response.output_item.done,然后正常结束,没有 response.completed | 允许 1 次流重试时,turn 最终完成,服务器恰好收到 2 个请求 |
retry_after.rs | HTTP 429 / 503 带 Retry-After: 1;流里返回限流错误;流里返回 overload | HTTP 429(由采样循环重试):至少等够 1 秒才重发,最终成功,共 2 个请求。流里的限流错误:消息里写了 “try again in 1s” 就按它等;只有响应头带 Retry-After 时用本地退避,不按头等(标着 TODO,是现状)。流里的 overload:直接终止,只有 1 个请求 |
websocket_fallback.rs | WebSocket 握手返回 426,或者握手一直失败 | 426:只试 1 次 WebSocket 就切 HTTP;一直失败:WebSocket 共 4 次后切 HTTP,之后的 turn 一直走 HTTP |
stream_no_completed:流"说了一半"就正常结束
codex-rs/core/tests/suite/stream_no_completed.rs
:
fn sse_incomplete() -> String {
responses::sse(vec![serde_json::json!({
"type": "response.output_item.done",
})])
}
// ...
let (server, _) = start_streaming_sse_server(vec![
vec![StreamingSseChunk {
gate: None,
body: incomplete_sse,
}],
vec![StreamingSseChunk {
gate: None,
body: completed_sse,
}],
])
.await;
// ...
// exercise retry path: first attempt yields incomplete stream, so allow 1 retry
request_max_retries: Some(0),
stream_max_retries: Some(1),
stream_idle_timeout_ms: Some(2000),
// ...
let requests = server.requests().await;
assert_eq!(
requests.len(),
2,
"expected retry after incomplete SSE stream"
);
sse_incomplete()构造的事件只有type一个字段,按上面sse()的规则,连 data 行都没有,只有一行event: response.output_item.done,然后流就结束了。- 这个测试没用 wiremock,用的是同目录的
start_streaming_sse_server(core/tests/common/streaming_sse.rs),一个能逐块控制发送时机的小 HTTP 服务器。第一段剧本是半截流,第二段是完整的response.created+response.completed。 - 配置把 HTTP 层重试关掉(
request_max_retries: Some(0)),只允许 1 次流层重试,空闲超时 2 秒。这样只有一条路径能让测试通过:发现流没说完,重试一次。 - 断言收到恰好 2 个请求。
被测的生产代码在
codex-rs/codex-api/src/sse/responses.rs:545-586
:
Ok(None) => {
let error = response_error.unwrap_or(ApiError::Stream(
"stream closed before response.completed".into(),
));
let _ = tx_event.send(Err(error)).await;
return;
}
Err(_) => {
let _ = tx_event
.send(Err(ApiError::Stream("idle timeout waiting for SSE".into())))
.await;
return;
}
};
// ...
let event: ResponsesStreamEvent = match serde_json::from_str(&sse.data) {
Ok(event) => event,
Err(e) => {
debug!(
// ...
"Failed to parse SSE event"
);
continue;
}
};
三种情况三种处理:流读到头(Ok(None))还没见到 completed,报流错误;两个事件之间超过空闲超时,也报流错误;某个事件的 data 解析不了,记一条 debug 日志,跳过这个事件继续读。最后这一条值得注意,下面会看到 Claude 的 Python SDK 选择了相反的做法。
两层重试,分开配置
Codex 的重试分两层,默认值在
codex-rs/model-provider-info/src/lib.rs:62-64
:
const DEFAULT_STREAM_IDLE_TIMEOUT_MS: u64 = 300_000;
const DEFAULT_STREAM_MAX_RETRIES: u64 = 5;
const DEFAULT_REQUEST_MAX_RETRIES: u64 = 4;
request_max_retries 管 HTTP 层的 5xx、超时、网络错误,不含 429(请求层的 ApiRetryConfig 写死了 retry_429: false,
model-provider-info/src/lib.rs:447-453
);带 Retry-After 的 429 和流读到一半的错误,由外面的采样循环按 stream_max_retries 重试。空闲超时默认 300 秒。测试里把一层设 0、另一层设 1,就能精确地只打开一条路径。注意两层不是互斥的:5xx 两层都会重试,次数相乘(见「Codex 的 agent loop」一章里 3 × 3 = 9 的例子)。
retry_after:429 要等够,流里的错误另算
codex-rs/core/tests/suite/retry_after.rs:345-401
,节选:
/// A 429 with server advice retries in the sampling loop and completes without a terminal error.
#[tokio::test(flavor = "current_thread")]
async fn responses_http_429_uses_retry_after() -> Result<()> {
// ...
let response_mock = responses::mount_response_sequence(
&server,
vec![
ResponseTemplate::new(429)
.insert_header("Retry-After", "1")
.set_body_json(json!({ "error": { "code": "rate_limit_exceeded" } })),
responses::sse_response(responses::sse(vec![
responses::ev_response_created("recovered"),
responses::ev_completed("recovered"),
])),
],
)
.await;
// ...
assert!(wait_for_retry(&mut telemetry, &retry).await >= Duration::from_secs(1));
// ...
assert_eq!(response_mock.requests().len(), 2);
剧本第一段是 429 加 Retry-After: 1,第二段是正常的流。注意配置里 request_max_retries 是 0,这次重试发生在采样循环里:测试断言重试事件的 layer 是 "stream"、operation 是 "sampling"(:375-383),而不是 HTTP 请求层。断言有两条:重试前实际等待不少于 1 秒(wait_for_retry 从失败请求开始计时,到下一个请求发出为止),以及恰好 2 个请求。等待时间是通过订阅 tracing 日志里的重试事件测出来的(文件开头的 RetryTelemetryLayer),没有改生产代码去暴露内部状态。
同一个文件里的另外两类测试值得对照着读:
sse_failure_uses_local_backoff_despite_retry_after(:1121):HTTP 200,流里返回限流错误(消息里没有给等待时间),只有响应头带Retry-After: 1。断言重试延迟落在 180–220 ms 之间,也就是用本地退避,不按头里的 Retry-After 等。对照紧接着的sse_rate_limit_message_uses_server_advised_retry_delay(:1275-1329):流里的限流错误消息写着 “Please try again in 1s.” 时,会按这个服务端建议等够 1 秒。所以"流里的限流用本地退避"只指建议只在响应头里、消息里没有的情况。sse_overload_with_retry_after_is_terminal(:1389):流里返回server_is_overloaded,即使带了 Retry-After、即使允许重试 2 次,也直接终止,断言只有 1 个请求,并给用户提示换模型。
同一个"Retry-After: 1",出现在不同位置,处理不同。但这不是设计上认定的正确行为:这两个测试上方都标着 // TODO(anp) respect Retry-After(
:1119
、
:1387
,整个文件有 7 处同样的 TODO),第一个测试的文档注释也写的是 “SSE failures currently retry with local backoff”。这类测试钉住的是当前行为,而这个行为已经计划要改。它的作用是让行为变化变得显式:哪天有人让流里的错误也认 Retry-After,这些测试会红,改的人必须同时改测试、在 diff 里写明,而不是悄悄变了。测开写这类"现状测试"时,最好也像这样在测试旁边注明它是现状、不是规格。
websocket_fallback:传输层降级
codex-rs/core/tests/suite/websocket_fallback.rs:33-128
。Codex 优先用 WebSocket 连模型,不行再退回 HTTP。两个测试:
- 握手请求返回 426(Upgrade Required):启动时的预热连接看到 426 就立即切 HTTP,断言 WebSocket 尝试 1 次、HTTP 请求 1 次。注释说明了为什么要特殊处理 426:否则采样循环会先重试 WebSocket 握手,再切 HTTP。
- 握手一直失败(wiremock 上根本没挂 GET 路由):断言 WebSocket 尝试 4 次(启动预热 1 次 + 首轮 1 次 + 重试 2 次)、HTTP 1 次。同文件后面的
websocket_fallback_is_sticky_across_turns断言切到 HTTP 之后,后续的 turn 不再尝试 WebSocket。
从测开的角度看,这三个文件就是一份故障清单:没说完、限流、传输层降级。每一条都同时断言了"最终结果"和"请求次数或等待时间",后者才是防回归的关键。重试多一次少一次,用户不一定看得出来,账单和限流配额看得出来。
快照测试:insta
Codex 用 insta 1.46(codex-rs/Cargo.toml:398)做快照测试。整个仓库有 1,488 个 .snap 文件,其中 TUI 1,373 个,core 84 个。core 里的快照多半拍的是发给模型的请求长什么样,例如
codex-rs/core/tests/suite/additional_context.rs:89-95
:
insta::assert_snapshot!(
"additional_context_simple_input",
context_snapshot::format_labeled_requests_snapshot(
"additional context is inserted before the user turn input.",
&[("Request", &request)],
&ContextSnapshotOptions::default().rewrite_known_segments(),
)
);
先把录下来的请求格式化成人能读的文本,rewrite_known_segments() 把权限说明之类的大段固定文本替换成 <PERMISSIONS_INSTRUCTIONS> 这样的占位符;然后和 core/tests/suite/snapshots/ 下的文件逐字比较。这个快照文件的正文长这样:
Scenario: additional context is inserted before the user turn input.
## Window 1
-- request 1 (turn; Request) --
00:message/developer:
<PERMISSIONS_INSTRUCTIONS>
01:message/developer:
<automation_info>run one</automation_info>
02:message/user:
<external_browser_info>tab one</external_browser_info>
03:message/user:
inspect the active tab
上下文的拼接顺序一变,快照就红;改动是预期的,就重新生成快照,评审时看 .snap 的 diff 就知道请求变了什么。对 agent 来说,“请求长什么样"直接决定模型看到什么,用快照守住它很划算。仓库里没找到关于 cargo insta review 流程的书面规定。
lint:为什么禁止 unwrap / expect
codex-rs/Cargo.toml:561-597
的 [workspace.lints.clippy] 把 36 条 clippy lint 设成 deny,节选:
[workspace.lints.clippy]
await_holding_invalid_type = "deny"
await_holding_lock = "deny"
disallowed_methods = "deny"
expect_used = "deny"
# ...
redundant_clone = "deny"
# ...
uninlined_format_args = "deny"
# ...
unwrap_used = "deny"
codex-rs/clippy.toml
的前两行又把测试放开了:
allow-expect-in-tests = true
allow-unwrap-in-tests = true
unwrap() 和 expect() 相当于 Python 里"假设一定成功,失败就让进程崩”。在一个长时间运行、手里攥着用户会话的 agent 里,一个没想到的空值会让整个会话崩掉;正确做法是把错误往上返回,由 loop 决定是重试、回给模型还是终止,也就是第 2 周说的"工具报错也要回结果"。测试里崩了就是测试失败,正是想要的,所以放开。
await_holding_lock 防的是持有锁时 await,这在并发执行工具时会造成死锁。disallowed_methods 配合 clippy.toml 里的清单,禁止直接调用某些方法,并在清单里写明理由(例如 SQLite 连接必须走 codex-state 的封装)。另外 codex-rs/rustfmt.toml 设了 imports_granularity = "Item",每个 use 只导一个项,减少合并冲突。
CI 门禁
仓库里的设计是让 PR 合并只依赖一个汇总检查:“CI required”。文件开头的注释说,required job 是"main 分支的规则集应当要求的"那份受版本控制的清单(
blocking-ci.yml:10-11
);GitHub 上的规则集实际配了哪些必需检查,在仓库文件里看不到,我没有核实。
.github/workflows/blocking-ci.yml:48-72
:
required:
name: CI required
# Without `always()`, GitHub skips this job after a failed dependency and a
# required check can appear successful instead of reporting the failure.
if: ${{ always() }}
needs:
- bazel
- blob-size-policy
- cargo-deny
- codespell
- repo-checks
- rust-ci
- sdk
runs-on: ubuntu-24.04
steps:
# ...
- name: Require successful dependencies
env:
NEEDS: ${{ toJSON(needs) }}
run: python3 .github/scripts/check_ci_results.py
- 这个 job 依赖 7 个子工作流。不写
if: always()的话,任何一个依赖失败,它会被跳过,而"被跳过"的必需检查可能显示成通过。注释写的就是这个坑。 - 配套脚本
.github/scripts/check_ci_results.py把needs的结果逐个检查,只有success算过,skipped、cancelled 也算失败。 - PR 阶段跑的是 Bazel 测试和 Bazel clippy(
bazel.yml),外加格式、cargo-shear、argument-comment-lint 等(rust-ci.yml)。全平台的cargo clippy ... -- -D warnings和分片 nextest 在rust-ci-full.yml里,合并到 main 之后由postmerge-ci.yml触发。 - nextest 的配置
codex-rs/.config/nextest.toml:超过 30 秒算慢、2 个周期后终止,失败自动重试 1 次;一些已知很慢的测试单独放宽超时,其中一组注明 “Do not add new tests here”。
“没跑不能被当成通过”,测开做质量门禁时也会遇到同一类问题,比如用例被 skip 了、某个分片没起来,报告上却是绿的。
PR 评审规范
.github/codex/labels/codex-rust-review.md
是让 Codex 审 PR 用的 prompt,本身就是一份评审规范。和测试相关的几条:
- PR 描述要说清动机;一个 PR 只做一件事,重构和功能分开(
:11-13)。 - 测试里对整个对象做一次
assert_eq!,不要逐字段比较(:23-25)。理由是单元测试也是"可执行的文档",逐字段比较更啰嗦、覆盖更少,文件里给了一好一坏两个例子。 - 不用
unsafe,尤其不要为了std::env::set_var()去用,这曾多次导致竞态;换别的方式传配置(:127)。上面stream_no_completed.rs里那句注释 “Configure retry behavior explicitly to avoid mutating process-wide environment variables” 就是这条规范的落地。
这份文件也有过时的地方:它建议把共享逻辑放进 codex-rs/common(:19),而这个 crate 在当前 commit 里已经不存在。这个 prompt 被哪个 workflow 消费,我没有核实。
动手:Python 版假模型服务器,给第 2 周的 loop 写测试
代码在本地 week03_Codex源码/code/mock_model_server/,四个文件:
| 文件 | 作用 | 对应 Codex |
|---|---|---|
mock_model_server.py | 假服务器,只用标准库 http.server;事件构造器;剧本队列;请求记录;配对检查 | responses.rs |
agent_loop.py | 被测对象:第 2 周的 loop 改成流式调用,加流层重试 | run_turn + 流重试 |
conftest.py | pytest fixture;一个极简的 JSON 快照断言 | start_mock_server / insta |
test_agent_loop.py | 6 个测试 | tools.rs / stream_no_completed.rs / retry_after.rs |
被测的 agent 用官方 SDK anthropic(1.11.0)把 base_url 指向假服务器。好处是 SSE 格式写错了,SDK 会直接解析失败,等于顺便验证了假服务器的格式;Claude API 的事件格式按官方文档的流式示例核对过:message_start → 若干个内容块(content_block_start、若干 content_block_delta、content_block_stop)→ message_delta(带 stop_reason)→ message_stop,中间可能夹着 ping;工具参数以 input_json_delta 的 partial_json 分段到达。
假服务器的核心
事件构造器,字段照官方文档的示例写:
def ev_tool_use_block(index: int, tool_id: str, name: str, tool_input: dict, pieces: int = 3) -> list[dict]:
"""一个 tool_use 块:参数以 input_json_delta 的 partial_json 分段到达。"""
raw = json.dumps(tool_input, ensure_ascii=False)
step = max(1, -(-len(raw) // pieces))
parts = [raw[i:i + step] for i in range(0, len(raw), step)]
return (
[{"type": "content_block_start", "index": index,
"content_block": {"type": "tool_use", "id": tool_id, "name": name, "input": {}}}]
+ [{"type": "content_block_delta", "index": index, "delta": {"type": "input_json_delta", "partial_json": p}}
for p in parts]
+ [{"type": "content_block_stop", "index": index}]
)
def tool_turn(text: str, tool_id: str, name: str, tool_input: dict, msg_id: str = "msg_tool") -> list[dict]:
"""完整的一轮:一段文本 + 一个 tool_use,stop_reason = tool_use。"""
return [ev_message_start(msg_id), *ev_text_block(0, text), ev_ping(),
*ev_tool_use_block(1, tool_id, name, tool_input), ev_message_delta("tool_use"), ev_message_stop()]
响应用 HTTP/1.1 分块编码(chunked)发送,这样能区分两种"流结束":
def _stream(self, script: SSE) -> None:
self.send_response(200)
self.send_header("content-type", "text/event-stream")
self.send_header("transfer-encoding", "chunked")
self.end_headers()
for i, ev in enumerate(script.events):
if script.cut_after is not None and i >= script.cut_after:
self.wfile.flush()
self.connection.shutdown(socket.SHUT_RDWR) # 不发结束块,直接断开
self.close_connection = True
return
chunk = encode_event(ev)
self.wfile.write(f"{len(chunk):X}\r\n".encode() + chunk + b"\r\n")
self.wfile.flush()
if script.delay_s:
time.sleep(script.delay_s)
self.wfile.write(b"0\r\n\r\n") # 分块编码的结束块:HTTP 层面"正常说完"
self.wfile.flush()
- 正常结束:最后发一个长度为 0 的结束块,HTTP 层面这个响应是完整的。把剧本里的
message_stop去掉,就得到"HTTP 正常、协议没说完"的流,对应 Codex 的stream_no_completed。 - 中途断开:
cut_after=n时发完前 n 个事件直接关 TCP,不发结束块,客户端会发现响应体不完整。
剧本队列和请求记录对应 Codex 的 mount_sse_sequence 和 ResponseMock:第 n 个请求拿第 n 段剧本,多出来的请求记进 unexpected 并回 500;每个请求都过一遍 check_tool_pairing(每个 tool_use 在下一条消息里都有同 id 的 tool_result,且 tool_result 排在文本前面),违规记进 violations。conftest.py 的 fixture 在每个测试结束时断言这两个列表都为空,所以和 Codex 一样,每个测试都顺带测了配对规则。
被测的 loop:多了一层流层重试
第 2 周的 loop 只改了两处:client 从外面传进来;调模型改成流式,并包一层流层重试。HTTP 层的 429 / 5xx 仍然交给 SDK(默认 max_retries=2),流层只管读到一半出的问题。这个分工是我们自己的设计,和 Codex 不同:Codex 的请求层不重试 429(429 交给采样循环),5xx 则两层都会重试、次数相乘。我们这边 STREAM_ERRORS 里没有 HTTP 状态错误,所以 429 / 5xx 只被 SDK 重试一层;但它含 anthropic.APIConnectionError,而拿到响应之前超时的话,SDK 会先自己重试,耗尽后抛的 APITimeoutError 正是它的子类,于是流层再重试,次数照样相乘。本地实测:服务器只接受连接不回应、timeout=0.3,max_retries=2 加流层 2 次重试,一共建立了 9 个连接(3 × 3)。这个乘法我们没有消掉,只是把它记下来;要消掉,可以在流层把 APITimeoutError 排除,或者把其中一层设成 0。
class StreamIncomplete(Exception):
"""HTTP 正常结束了,但没收到 message_stop:模型这一轮没说完。"""
# 流层可重试的错误:断开(传输层)、没说完(协议层)、畸形 data(json.JSONDecodeError 是 ValueError)
STREAM_ERRORS = (StreamIncomplete, httpx2.TransportError, anthropic.APIConnectionError, ValueError)
def stream_once(client: anthropic.Anthropic, messages: list) -> anthropic.types.Message:
saw_stop = False
with client.messages.stream(model=MODEL, max_tokens=64000, tools=TOOLS, messages=messages) as stream:
for event in stream:
if event.type == "message_stop":
saw_stop = True
if not saw_stop:
# SDK 不会因为缺 message_stop 报错,get_final_message() 照样返回半截消息
raise StreamIncomplete("stream closed before message_stop")
return stream.get_final_message()
def call_model(client, messages, max_stream_retries=2, backoff_s=0.2,
sleep: Callable[[float], None] = time.sleep, log: list | None = None):
for attempt in range(max_stream_retries + 1):
try:
return stream_once(client, messages)
except STREAM_ERRORS as e:
if log is not None:
log.append(f"attempt {attempt + 1}: {type(e).__name__}: {e}")
if attempt == max_stream_retries:
raise
sleep(backoff_s * 2 ** attempt) # 指数退避:0.2s、0.4s……
raise AssertionError("unreachable")
重试时整轮重发,历史里不留半截的 assistant 消息。这是我们的选择,不是官方唯一推荐的做法:官方流式文档的
Error recovery
一节建议保存已收到的半截内容再续写,Claude 4.5 及更早的模型把它作为 assistant 消息的开头;4.6 及之后的模型改成发一条 user 消息,里面带上半截回复,并让模型从断开处接着写。我们选整轮重发有两个理由:同一节指出 tool_use 和 thinking 块不能部分恢复,而我们断开的位置常常落在 tool_use 里;整轮重发时重发的请求体和第一次逐字相同,测试里一条等式就能断言。代价是已经生成的那部分要重新生成一遍。run_agent 本身和第 2 周一样,只是把 client.messages.create(...) 换成了 call_model(client, messages)。
写这几个测试时发现的 SDK 行为
先用假服务器直接探了一下 SDK 1.11.0 面对各种故障的表现,这几条正是测试要防的坑:
| 注入的故障 | SDK 的表现(本地实测) | agent 要自己做什么 |
|---|---|---|
| 流读到一半 TCP 断开 | 抛 httpx2.RemoteProtocolError: peer closed connection without sending complete message body (incomplete chunked read),不是 anthropic.APIConnectionError | 重试时连传输层异常一起捕获 |
正常结束但没有 message_stop | 不报错,get_final_message() 返回累积到的消息;如果连 message_delta 也没到,stop_reason 是 None | 自己记下有没有见到 message_stop |
| data 不是合法 JSON | 抛 json.JSONDecodeError(Codex 是跳过这个事件) | 归入流层重试,次数有上限 |
ping 事件的 data 是坏的 | 不报错:SDK 遇到 ping 直接跳过,不解析 data。不只是 ping:anthropic/_streaming.py:75-129 只解析白名单里的事件名(message_start、content_block_delta 等),白名单之外的事件名(error 另有处理)连 data 都不读,直接丢掉 | 写故障用例时,注入点要选在会被解析的事件上 |
429 + retry-after: 1 | 默认 max_retries=2,按 retry-after 等待后重发,请求头 x-stainless-retry-count 从 0 变成 1 | 别把 max_retries 关掉 |
关于 retry-after:SDK 1.11.0 的 _calculate_retry_timeout(anthropic/_base_client.py)先读非标准的 retry-after-ms 头,再读 retry-after(秒数或 HTTP 日期),拿到正数就照用;都没有时才用 0.5 秒起、每次翻倍、上限 8 秒、带随机抖动的指数退避。重试的状态码是 408、409、429 和 5xx,服务器也可以用 x-should-retry 头显式指定要不要重试。
6 个测试
完整代码见本地 test_agent_loop.py,这里贴两个。第一个对应 Codex 的 tools.rs,断言第二次请求里 agent 发回去的东西,按评审规范对整个对象做一次比较,再把整个请求体存成快照:
def test_tool_result_is_sent_back(server, client):
server.enqueue(
SSE(tool_turn("先读一下 README。", "toolu_01", "read_file", {"path": "README.md"})),
SSE(text_turn("这是一个评测找 Bug 的 Agent 的项目。")),
)
status, _ = run_agent(client, TASK, impls=IMPLS, **FAST)
assert status == "done"
assert len(server.requests) == 2
second = server.requests[1].body
# 整个对象一次比较,不逐字段(Codex 评审规范的要求)
assert second["messages"][1:] == [
{"role": "assistant", "content": [
{"type": "text", "text": "先读一下 README。"},
{"type": "tool_use", "id": "toolu_01", "name": "read_file", "input": {"path": "README.md"}},
]},
{"role": "user", "content": [
{"type": "tool_result", "tool_use_id": "toolu_01", "content": FILE_TEXT},
]},
]
# 整个请求体存成快照:哪天请求形状变了(多了字段、顺序变了),这里会红
assert_json_snapshot("tool_roundtrip_request_2", second)
第二个对应 retry_after.rs,断言两次请求的间隔和 SDK 的重试计数头:
def test_429_respects_retry_after(server, client):
server.enqueue(
HTTPError(429, "rate_limit_error", "Number of requests has exceeded your rate limit.",
headers={"retry-after": "1"}),
SSE(text_turn("这是一个评测找 Bug 的 Agent 的项目。")),
)
status, _ = run_agent(client, TASK, impls=IMPLS, **FAST)
assert status == "done"
assert len(server.requests) == 2
gap = server.requests[1].at - server.requests[0].at
assert gap >= 1.0, f"只等了 {gap:.3f} 秒"
assert [r.headers["x-stainless-retry-count"] for r in server.requests] == ["0", "1"]
其余四个:
test_retries_after_midstream_disconnect:发完 5 个事件断开 TCP,断言共 3 个请求、重发的请求体和第一次逐字相同(没把半截消息塞进历史),日志里是RemoteProtocolError。test_missing_message_stop_is_retried:去掉message_stop,断言重试了一次。test_429_without_http_retries_fails_fast:反例,max_retries=0时 429 直接变成anthropic.RateLimitError抛出,只有 1 个请求。test_malformed_json_gives_up_after_bounded_retries:每次 data 都是坏的,断言重试 2 次后放弃、恰好 3 个请求,不会无限重试。
本机真实输出(Python 3.13.7,anthropic 1.11.0,pytest 9.1.1):
$ .venv/bin/python -m pytest -v -p no:cacheprovider
collecting ... collected 6 items
test_agent_loop.py::test_tool_result_is_sent_back PASSED [ 16%]
test_agent_loop.py::test_retries_after_midstream_disconnect PASSED [ 33%]
test_agent_loop.py::test_missing_message_stop_is_retried PASSED [ 50%]
test_agent_loop.py::test_429_respects_retry_after PASSED [ 66%]
test_agent_loop.py::test_429_without_http_retries_fails_fast PASSED [ 83%]
test_agent_loop.py::test_malformed_json_gives_up_after_bounded_retries PASSED [100%]
============================== 6 passed in 4.37s ===============================
4 秒多里有 1 秒是 429 测试在真等 retry-after。流层退避在测试里调成了 0.01 秒(backoff_s 参数),不用真等。
测试真的能抓 bug 吗:变异验证
全绿只说明测试和代码一致,不说明测试有用。把被测代码复制一份,每次故意改坏一处,再跑一遍:
== mutation: 从 STREAM_ERRORS 去掉 httpx2.TransportError
FAILED test_agent_loop.py::test_retries_after_midstream_disconnect - httpx2.R...
1 failed, 5 passed in 4.29s
== mutation: 不检查 message_stop
FAILED test_agent_loop.py::test_missing_message_stop_is_retried - AssertionEr...
1 failed, 5 passed in 4.31s
== mutation: 工具结果截成 10 个字符
FAILED test_agent_loop.py::test_tool_result_is_sent_back - AssertionError: as...
1 failed, 5 passed in 4.41s
三处变异各让恰好 1 个测试变红,说明每个测试都守着一条具体的路径。第一处尤其值得记住:去掉 httpx2.TransportError 之后只剩 SDK 的 anthropic.APIConnectionError 兜着断连。照"SDK 的连接错误类"去写 except,看起来很合理,但流读到一半断开时根本抛不到那里,没有这个测试很难发现。
运行方式:
cd week03_Codex源码/code/mock_model_server
python3.13 -m venv .venv && .venv/bin/pip install -r requirements.txt # anthropic==1.11.0, pytest==9.1.1
.venv/bin/python -m pytest -v
测开视角:迁移到自己的 agent
- 先做假服务器,再写 agent 的测试。剧本队列、请求记录、不变量检查这三件是地基,所有测试共享。第 5 周"长任务可靠性"的故障注入会直接复用这个服务器。
- 断言请求和次数,不断言回答。工具结果按 id 回填、重试时请求体不变、等待时间不少于 Retry-After、请求次数恰好是 N。
- 按位置分类重试,并写明谁管什么。拿到响应之前和读到一半的错误分开配置、分开测。分工有不同做法:我们的 Python 版让 SDK 管 429 / 5xx、agent 管流错误,这两类不相乘,但拿到响应之前的超时仍会被两层各重试一遍;Codex 的请求层不管 429,5xx 则两层都重试、次数相乘。不管哪种,都要有测试把"一个故障最多发几个请求"钉住。Codex 里同一个 Retry-After 在响应头和流里处理不同,那几个测试钉住的是标了 TODO、计划要改的现状,作用是让行为变化显式化,不代表这样才对。
- 每条重试路径都测"故障一直不消失"。断言最终报错、请求次数有上限,防止无限重试,也防止把失败报成成功。
- 用官方 SDK 当客户端。假服务器的格式错了 SDK 会报错;SDK 自身的行为(比如抛哪个异常、缺 message_stop 时不报错)也会被测出来。升级 SDK 时这些测试就是回归保护。
- 定期做变异验证。改坏一处,看对应的测试会不会红。不会红的测试要么断言不够,要么剧本没走到那条路径(比如把坏数据注入在 ping 上)。
- 门禁只认 success。汇总检查要在依赖失败时照样运行,skipped、cancelled 都算失败。
常见错误说法
- “用真实模型跑几遍没出错,重试逻辑就没问题”:真实 API 很少恰好在你跑测试时断流,这些路径根本没被走到。要用假服务器把每条故障路径确定地跑一遍。
- “SDK 有重试,流式调用就不用管了”:SDK 的重试只管拿到响应之前;流读到一半断开、没收到
message_stop,都要 agent 自己处理。 - “捕获
anthropic.APIConnectionError就能兜住断流”:SDK 1.11.0 里流读到一半断开,抛的是httpx2.RemoteProtocolError。 - “没收到 message_stop,SDK 会报错”:不会,
get_final_message()照样返回累积到的消息。 - “集成测试要断言模型回答得对不对”:那是评测。集成测试断言的是 agent 发出去的请求和最终状态。
- “Codex 禁止 unwrap,所以测试里也不能用”:
clippy.toml明确允许测试里用unwrap/expect。 - “Codex 的 HTTP 请求层会按 Retry-After 重试 429”:请求层的
retry_429是false,429 是由采样循环(stream_max_retries)按 Retry-After 重试的,测试里的重试事件layer是"stream"。 - “Codex 测试断言流里的错误不按 Retry-After 等,说明这样设计才对”:这些测试旁边标着
TODO(anp) respect Retry-After,钉住的是计划要改的现状。
本周到这里读完了 Codex 的六个部分。第 4 周转向上下文与记忆:RAG 和长期记忆。