Learning AI Quality 返回 KuthorX Blog II博客首页

第 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 的实测行为模拟。页面底部有自动判分的练习。

互动演示:SSE 剧本编辑器 在新标签页打开

为什么断言请求,而不是断言回答

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:录下每个请求,顺便检查不变量

responses.rs:738-750 :

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}"
    );

三步:

  1. 排剧本。第一段:模型调 exec_command,命令往只读沙箱外写一个文件,同时往 stderr 打印一个哨兵字符串;第二段:模型说 “done”。
  2. 用只读权限提交一轮。agent 真的去执行命令,沙箱真的拒绝。
  3. 从请求记录里取出 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.rsHTTP 429 / 503 带 Retry-After: 1;流里返回限流错误;流里返回 overloadHTTP 429(由采样循环重试):至少等够 1 秒才重发,最终成功,共 2 个请求。流里的限流错误:消息里写了 “try again in 1s” 就按它等;只有响应头带 Retry-After 时用本地退避,不按头等(标着 TODO,是现状)。流里的 overload:直接终止,只有 1 个请求
websocket_fallback.rsWebSocket 握手返回 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.pypytest fixture;一个极简的 JSON 快照断言start_mock_server / insta
test_agent_loop.py6 个测试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

  1. 先做假服务器,再写 agent 的测试。剧本队列、请求记录、不变量检查这三件是地基,所有测试共享。第 5 周"长任务可靠性"的故障注入会直接复用这个服务器。
  2. 断言请求和次数,不断言回答。工具结果按 id 回填、重试时请求体不变、等待时间不少于 Retry-After、请求次数恰好是 N。
  3. 按位置分类重试,并写明谁管什么。拿到响应之前和读到一半的错误分开配置、分开测。分工有不同做法:我们的 Python 版让 SDK 管 429 / 5xx、agent 管流错误,这两类不相乘,但拿到响应之前的超时仍会被两层各重试一遍;Codex 的请求层不管 429,5xx 则两层都重试、次数相乘。不管哪种,都要有测试把"一个故障最多发几个请求"钉住。Codex 里同一个 Retry-After 在响应头和流里处理不同,那几个测试钉住的是标了 TODO、计划要改的现状,作用是让行为变化显式化,不代表这样才对。
  4. 每条重试路径都测"故障一直不消失"。断言最终报错、请求次数有上限,防止无限重试,也防止把失败报成成功。
  5. 用官方 SDK 当客户端。假服务器的格式错了 SDK 会报错;SDK 自身的行为(比如抛哪个异常、缺 message_stop 时不报错)也会被测出来。升级 SDK 时这些测试就是回归保护。
  6. 定期做变异验证。改坏一处,看对应的测试会不会红。不会红的测试要么断言不够,要么剧本没走到那条路径(比如把坏数据注入在 ping 上)。
  7. 门禁只认 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 和长期记忆。