From 577a1ab7f96fc3d220ca89c068988a479ac7b1aa Mon Sep 17 00:00:00 2001 From: minixalpha Date: Sun, 20 Sep 2026 22:11:10 +0800 Subject: [PATCH 1/2] fix: recover interrupted model response streams Retry response-body read errors, read timeouts, and remote protocol errors at most twice (1s, 2s) before committing a reply to history or running any of its tools. Handle both httpx and httpx2 transport families; the CLI previously only caught httpx.HTTPError, so SDK 1.5.0 failures on httpx2 ended the run. Retries reuse the last complete conversation, never execute tools from an interrupted attempt, and do not consume the turn budget. A 300-second window starting at the first interruption bounds new retry scheduling, while in-flight requests keep the SDK timeout and external deadline. Record each attempt as a journal v3 `model.failed` event and project it as an incomplete ATIF step. Reconcile generation costs for interrupted attempts and flag token and cost totals as partial when usage is missing. --- .github/workflows/ci.yml | 2 + .../reports/stream-recovery-20260915.md | 139 +++++++++++ .../harbor/tests/test_atif_compatibility.py | 29 +++ docs/dev_docs/README.md | 6 + docs/dev_docs/en/event-journal-protocol-v2.md | 5 +- docs/dev_docs/en/event-journal-protocol-v3.md | 59 +++++ .../zh-CN/event-journal-protocol-v2.md | 3 +- .../zh-CN/event-journal-protocol-v3.md | 48 ++++ docs/user_docs/en/cli_reference.md | 15 ++ docs/user_docs/zh-CN/cli_reference.md | 12 + src/nanopycodeagent/agent.py | 133 ++++++++--- src/nanopycodeagent/atif.py | 52 +++- src/nanopycodeagent/event_journal.py | 23 +- src/nanopycodeagent/transport.py | 18 ++ tests/test_agent_events.py | 2 + tests/test_cli.py | 2 + tests/test_event_journal.py | 4 +- tests/test_stream_recovery.py | 226 ++++++++++++++++++ tests/test_truncation.py | 4 +- 19 files changed, 733 insertions(+), 49 deletions(-) create mode 100644 benchmarks/harbor/reports/stream-recovery-20260915.md create mode 100644 docs/dev_docs/en/event-journal-protocol-v3.md create mode 100644 docs/dev_docs/zh-CN/event-journal-protocol-v3.md create mode 100644 src/nanopycodeagent/transport.py create mode 100644 tests/test_stream_recovery.py diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index ce2644b..5f96079 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -39,3 +39,5 @@ jobs: # --frozen installs exactly what uv.lock pins and fails if the lock has # drifted from pyproject.toml, so CI tests the declared dependencies. run: uv run --frozen pytest + - name: Test the httpx2 SDK transport + run: uv run --with anthropic==1.5.0 pytest diff --git a/benchmarks/harbor/reports/stream-recovery-20260915.md b/benchmarks/harbor/reports/stream-recovery-20260915.md new file mode 100644 index 0000000..b274067 --- /dev/null +++ b/benchmarks/harbor/reports/stream-recovery-20260915.md @@ -0,0 +1,139 @@ +# Interrupted response recovery: implementation and targeted validation + +## Status + +The transport compatibility fix, bounded stream recovery, offline tests, and +independent review are complete. Live targeted trials are prepared and await +explicit approval for OpenRouter data transfer and API charges. No live model +requests have been made by this experiment yet. + +## Problem and implementation + +The September 12 and September 14 experiments recorded four interrupted model +responses across three distinct tasks. A `RemoteProtocolError` during response +iteration ended the whole agent run. SDK 1.5.0 used `httpx2`, while the CLI +caught `httpx.HTTPError`; these exception families do not share that base class. + +The change handles both HTTP families and retries response-body `ReadError`, +`ReadTimeout`, and `RemoteProtocolError` at most twice, with delays of 1 and 2 +seconds. It retains the last complete conversation and executes tools only +after a complete reply is received. Interrupted replies do not consume the +completed-reply turn budget. Authentication, invalid requests, programming +errors, and user cancellation are not retried by this loop. Pre-stream errors +retain the SDK retry policy without a second outer retry layer. + +No further retry is scheduled beyond 300 seconds after a reply's first +interruption. This scheduling window does not cancel an in-flight request; +the SDK timeout and external supervisor deadline remain independent limits. + +Journal v3 adds `model.failed`, with a separate ID per attempt, the original +error, retry intent, duration, and generation ID captured before reading the +body. ATIF keeps failed attempts in chronological order. Interrupted-generation +billing can be reconciled without treating missing token usage as complete. +Readers continue to support v1/v2 journals. + +## Automated and offline evidence + +| Validation | Result | +| --- | --- | +| Complete suite, locked SDK 0.112.0 / httpx | 292 passed, 4 skipped | +| Complete suite, SDK 1.5.0 / httpx2 2.12.0 | 296 passed | +| Harbor adapter and official ATIF compatibility tests | 24 passed | +| Historical interrupted response replay, SDK 1.5.0 | 4/4 recovered | + +The locked environment skips four tests that explicitly require the absent +`httpx2` package. The SDK 1.5.0 run exercises both exception families. Both +full suites retain the existing Pydantic warning in the malformed-tool-input +test. CI now runs both the locked dependencies and an SDK 1.5.0 environment. + +The real SDK transport test interrupts a response after complete-looking tool +JSON arrives, verifies that the discarded tool is never executed, and checks +that an earlier append operation occurs only once. Further tests cover retry +exhaustion, the scheduling window, preserved conversation history, interrupts, +permanent errors, billing reconciliation, and journal version rejection. + +Offline replay consumes the original captured response bytes through SDK +1.5.0, raises the recorded transport error at EOF, and supplies a synthetic +successful reply to the retry. All four cases closed the interrupted response, +retried the same request, preserved the failed attempt, and executed no tools: + +| Historical trial | Captured bytes | Outcome | +| --- | ---: | --- | +| `torch-tensor-parallelism__EEWJenL` | 280,943 | Recovered | +| `schemelike-metacircular-eval__Xuw2yU4` | 33,440 | Recovered | +| `schemelike-metacircular-eval__CJJvMtU` | 1,057,796 | Recovered | +| `llm-inference-batching-scheduler__S95o8Jv` | 34,318 | Recovered | + +These replays make no external requests and provide no new task score. They +verify response handling; the retry's success is deliberately synthetic. + +## Independent review + +A Codex agent named `stream-reviewer` ran in a sibling Herdr split using the +[review-agent skill](/home/minix/.codex/skills/.system/review-agent/SKILL.md). +It reviewed the complete uncommitted diff and new files, relevant call sites, +tests, retry boundaries, exception compatibility, journals, ATIF, and costs. +Its final result was **No findings**. It independently reported 76 passed / +4 skipped with locked dependencies and 80 passed with cached SDK 1.5.0. +The review split was closed before creating the benchmark split. + +## Targeted live experiment + +The selected tasks are every distinct task explicitly classified as +`RemoteProtocolError` in the committed historical result summaries: + +| Task | Relevant prior attempt | Prior result | +| --- | --- | --- | +| `schemelike-metacircular-eval` | September 14, request 18 | Reward 0; 0/63 subcases; no `eval.scm` | +| `llm-inference-batching-scheduler` | September 14, request 3 | Reward 0; interrupted model stream | +| `torch-tensor-parallelism` | September 12, request 1 | No official score; interrupted stream and later verifier dependency timeout | + +Prepared job: `tb21-stream-recovery-targeted3-20260915`. + +- Model: `openrouter/deepseek/deepseek-v4-flash-0731`. +- Dataset: `terminal-bench/terminal-bench-2-1`, pinned to + `sha256:7d7bdc1cbedad549fc1140404bd4dc45e5fd0ea7c4186773687d177ad3a0699a`. +- Preserve each task reference and cached image digest from the baseline. +- 50 completed replies, 65,536 tokens per reply, concurrency 2, one attempt per + task, no automatic whole-task Harbor retries. +- 3,600 seconds per agent run, 120 seconds of finalization grace, 3,780 seconds + for the outer Harbor deadline; unchanged native verifier limits. +- SDK 1.5.0, httpx 0.28.1, httpx2 2.12.0, Pydantic 2.13.5. +- Reuse the cached uv bootstrap and passive HTTP recorder; periodic stack + dumping stays disabled. Live trials do not inject transport faults. +- Install a wheel verified against every current package source file. Preserve + the uncommitted source snapshot and file hashes so results can be compared + with the final commit without pretending the snapshot was already committed. + +Local reproduction command from the repository root: + +```bash +PYTHONPATH=src benchmarks/harbor/.venv/bin/python \ + jobs/tb21-stream-recovery-targeted3-20260915-record/workflow/run.py +``` + +The runner refuses to overwrite an existing job, checks the wheel and image +digests, loads the existing endpoint credentials without writing them into the +report, and saves a source manifest, raw logs, HTTP records, journals, +trajectories, checkpoints, and verifier output under `jobs/`. + +## Interpretation limits + +The recovery fix cannot identify which upstream component closed the original +connections and does not establish that the model can solve these tasks. The +September 12 Scheme interpreter passed only 4/63 subcases before its stream +failed. Sampling, routing, cache state, and provider backends are uncontrolled. +This targeted experiment is not a full benchmark score or a controlled measure +of pass-rate improvement. A live run without an interruption does not exercise +the retry branch. + +## Evidence locations + +- [September 12 results](../results/tb21-generation-budget-65536-20260912.json) +- [September 14 results](../results/tb21-tool-input-recovery-20260914.json) +- [Journal v3 protocol](../../../docs/dev_docs/en/event-journal-protocol-v3.md) +- Local record: `jobs/tb21-stream-recovery-targeted3-20260915-record/` +- Local offline replay: `offline_replay.py` and `offline-replay.json` in that record + +Raw task/model data remains in the Git-ignored local job directories. Bilingual +CLI documentation and the v2/v3 protocol sources and English versions are in sync. diff --git a/benchmarks/harbor/tests/test_atif_compatibility.py b/benchmarks/harbor/tests/test_atif_compatibility.py index 95355b9..af10875 100644 --- a/benchmarks/harbor/tests/test_atif_compatibility.py +++ b/benchmarks/harbor/tests/test_atif_compatibility.py @@ -16,6 +16,35 @@ def test_projector_output_passes_harbor_atif_validator(): assert validator.validate(trajectory), validator.get_errors() +@pytest.mark.parametrize("resolved", [False, True]) +def test_recovered_stream_v3_passes_harbor_atif_validator(tmp_path, resolved): + fixture = Path(__file__).parent / "fixtures" / "atif-journal-v1.jsonl" + with EventJournal.create("run-recovered", directory=tmp_path) as journal: + for entry in EventJournal.replay(fixture): + if entry.type == "model.started": + journal.append(NativeEvent("model.started", entry.payload | { + "model_call_id": "failed-attempt", + })) + journal.append(NativeEvent("model.failed", { + "model_call_id": "failed-attempt", "error_type": "RemoteProtocolError", + "message": "interrupted", "generation_id": "gen-interrupted", + "duration_ms": 10, "will_retry": True, "retry_delay_seconds": 1, + "source_timestamp": entry.payload["source_timestamp"], + })) + if entry.type == "run.completed" and resolved: + journal.append(NativeEvent("model.cost_resolved", { + "generation_id": "gen-interrupted", "amount": "0.02", + "currency": "USD", "source": "test", + "source_timestamp": entry.payload["source_timestamp"], + })) + journal.append(NativeEvent(entry.type, entry.payload)) + trajectory = project_atif(EventJournal.replay(journal.path)) + validator = TrajectoryValidator() + assert validator.validate(trajectory), validator.get_errors() + assert trajectory["steps"][1]["extra"]["incomplete"] is True + assert trajectory["final_metrics"]["extra"]["usage_complete"] is False + + @pytest.mark.parametrize("content", [ [{"type": "text", "text": "Partial answer"}], [{"type": "extension", "namespace": "anthropic", "source_type": "thinking", diff --git a/docs/dev_docs/README.md b/docs/dev_docs/README.md index 24a0bd1..8892f2c 100644 --- a/docs/dev_docs/README.md +++ b/docs/dev_docs/README.md @@ -30,3 +30,9 @@ They complement, rather than replace: - Event Journal Protocol v1: [English](en/event-journal-protocol-v1.md) | [Chinese](zh-CN/event-journal-protocol-v1.md) +- Event Journal Protocol v2 (response truncation): + [English](en/event-journal-protocol-v2.md) | + [Chinese](zh-CN/event-journal-protocol-v2.md) +- Event Journal Protocol v3 (current; failed model attempts): + [English](en/event-journal-protocol-v3.md) | + [Chinese](zh-CN/event-journal-protocol-v3.md) diff --git a/docs/dev_docs/en/event-journal-protocol-v2.md b/docs/dev_docs/en/event-journal-protocol-v2.md index 1fe040e..18fb2ea 100644 --- a/docs/dev_docs/en/event-journal-protocol-v2.md +++ b/docs/dev_docs/en/event-journal-protocol-v2.md @@ -4,8 +4,9 @@ > [`../zh-CN/event-journal-protocol-v2.md`](../zh-CN/event-journal-protocol-v2.md). > Do not edit by hand. -v2 is implemented and is the internal Journal protocol used by the current -writer, with `schema_version = 2`. This document defines all changes relative +v2 is implemented with `schema_version = 2`. The current writer has moved to +[v3](event-journal-protocol-v3.md); this document preserves the v2 protocol. +This document defines all changes relative to [v1](event-journal-protocol-v1.md). Envelope, event types, fields, validation, ordering, persistence, and projection rules not listed here follow v1. Public trajectories remain ATIF-v1.7. diff --git a/docs/dev_docs/en/event-journal-protocol-v3.md b/docs/dev_docs/en/event-journal-protocol-v3.md new file mode 100644 index 0000000..edeeffc --- /dev/null +++ b/docs/dev_docs/en/event-journal-protocol-v3.md @@ -0,0 +1,59 @@ +# Event Journal Implementation Protocol v3 + +> Generated from the Chinese source +> [`../zh-CN/event-journal-protocol-v3.md`](../zh-CN/event-journal-protocol-v3.md). +> Do not edit by hand. + +The current writer emits `schema_version = 3`. Readers and the ATIF projector +continue to accept v1 and v2. Public trajectories remain ATIF-v1.7. +All contracts from [v2](event-journal-protocol-v2.md) still apply except for the +additional event and projection behavior described here. + +## Failed model attempts + +`model.failed` terminates an API attempt without claiming a complete message or +complete usage. It has these required payload fields: + +| Field | Type | Meaning | +| --- | --- | --- | +| `model_call_id` | nonempty string | The corresponding `model.started` identifier. | +| `error_type` | nonempty string | Original exception class name. | +| `message` | string | Original exception message. | +| `generation_id` | nonempty string or null | Provider header captured when the stream opens. | +| `duration_ms` | nonnegative number | Attempt duration including stream consumption. | +| `will_retry` | boolean | Whether the agent scheduled another attempt. | +| `retry_delay_seconds` | nonnegative number | Scheduled delay; zero when no retry is scheduled. | +| `source_timestamp` | RFC 3339 UTC or null | Event time. | + +The event is emitted for caught SDK and transport errors, including final +failures. Unexpected exceptions and user interrupts retain the previous +`model.started` / `run.failed` representation. Each retry gets a new +`model_call_id`; it is not another completed reply and does not consume the +turn budget. Retries retain the last committed conversation and never execute +tools from the interrupted attempt. `will_retry` records intent: cancellation +or a journal failure may prevent the next attempt from starting. + +The recovery policy retries only response-body read errors, read timeouts, +and remote protocol errors, with delays of 1 and 2 seconds. No new retry is +scheduled after a 300-second window starting at the first interruption of the +reply. An in-flight request can outlast that scheduling window; SDK timeouts +and an external supervisor's deadline remain independent limits. Failures +before the stream opens retain SDK retries without another outer retry layer. + +## Projection and costs + +Failed attempts become chronological agent steps, containing any visible text +deltas, `llm_call_count = 1`, and `extra.incomplete = true`. Error and retry +details remain in `extra`. They contain no tool calls or observations because +no tools from that response were executed. A later successful retry may end +the run normally without making the earlier attempt complete. + +Generation IDs are captured before reading the response body so finalization +can reconcile interrupted generations too. Resolved costs contribute to the +total; unresolved attempts keep cost totals explicitly partial. Missing token +usage remains unknown even when billing is resolved. Legacy incomplete model +starts still project as before. + +v1/v2 records containing `model.failed` are rejected. Old readers reject v3 +rather than silently dropping failed attempts. Historical journals are not +rewritten. diff --git a/docs/dev_docs/zh-CN/event-journal-protocol-v2.md b/docs/dev_docs/zh-CN/event-journal-protocol-v2.md index 5821cae..f5e3c45 100644 --- a/docs/dev_docs/zh-CN/event-journal-protocol-v2.md +++ b/docs/dev_docs/zh-CN/event-journal-protocol-v2.md @@ -3,7 +3,8 @@ > 本文件为**中文源文件**(source of truth);英文版 > [`../en/event-journal-protocol-v2.md`](../en/event-journal-protocol-v2.md) 由其生成。 -v2 已实现,是当前 writer 使用的内部 Journal 协议,`schema_version = 2`。 +v2 已实现,使用 `schema_version = 2`。当前 writer 已升级到 +[v3](../en/event-journal-protocol-v3.md),本文保留 v2 的协议定义。 本文完整定义相对 [v1](event-journal-protocol-v1.md) 的变化;未列出的 envelope、 事件类型、字段、校验、排序、持久化与投影规则沿用 v1。公开 trajectory 仍为 ATIF-v1.7。 diff --git a/docs/dev_docs/zh-CN/event-journal-protocol-v3.md b/docs/dev_docs/zh-CN/event-journal-protocol-v3.md new file mode 100644 index 0000000..c91718f --- /dev/null +++ b/docs/dev_docs/zh-CN/event-journal-protocol-v3.md @@ -0,0 +1,48 @@ +# Event Journal 实现协议 v3 + +> 本文件为中文源;英文版本 +> [`../en/event-journal-protocol-v3.md`](../en/event-journal-protocol-v3.md) 由其生成。 + +当前 writer 写入 `schema_version = 3`。Reader 和 ATIF projector 继续兼容 +v1、v2;公开轨迹仍为 ATIF-v1.7。除下述新增事件和投影行为外, +[v2](event-journal-protocol-v2.md) 的其他契约保持有效。 + +## 模型尝试失败 + +`model.failed` 结束一次 API 尝试,不宣称收到完整消息或完整用量。必填字段如下: + +| 字段 | 类型 | 含义 | +| --- | --- | --- | +| `model_call_id` | 非空字符串 | 对应 `model.started` 的标识。 | +| `error_type` | 非空字符串 | 原始异常类名。 | +| `message` | 字符串 | 原始异常消息。 | +| `generation_id` | 非空字符串或 null | 响应流打开时取得的服务商响应头。 | +| `duration_ms` | 非负数 | 包含响应流读取的尝试耗时。 | +| `will_retry` | 布尔值 | 是否已安排另一次尝试。 | +| `retry_delay_seconds` | 非负数 | 安排的等待时长;不重试时为零。 | +| `source_timestamp` | RFC 3339 UTC 或 null | 事件时间。 | + +捕获到 SDK 或传输错误时写入该事件,包括最终失败。意外异常和用户中断沿用 +`model.started` / `run.failed` 的表示。每次重试分配新的 `model_call_id`; +重试不是已完成回复,不消耗轮数预算。重试保留最后一次已提交的对话,绝不执行 +中断回复中的工具。`will_retry` 记录安排意图;取消或 Journal 写入失败仍可能 +阻止下一次尝试启动。 + +恢复策略只重试响应体读取错误、读取超时和远端协议错误,等待时间分别为 1 秒和 +2 秒。从本轮首次中断开始,超过 300 秒窗口后不再安排新重试。进行中的请求 +可以超过这个调度窗口;SDK 超时和外部监督器的截止时间仍是独立限制。响应流 +打开前的失败沿用 SDK 的重试,不再叠加外层重试。 + +## 投影与费用 + +失败尝试按顺序成为 agent step,保留可见文本增量,设置 `llm_call_count = 1` +和 `extra.incomplete = true`,错误及重试信息写入 `extra`。这些 step 不包含 +工具调用或 observation,因为没有执行该回复的工具。后续重试成功可以使 run +正常结束,但不会使此前的失败尝试变为完整。 + +在读取响应体前取得 generation ID,使收尾阶段也能查询中断生成的费用。已确认 +费用计入总额;存在未确认的尝试时,费用总额显式保持不完整。即使账单已经补全, +缺失的 token 用量仍然未知。旧版未完成的模型开始事件沿用原来的投影方式。 + +包含 `model.failed` 的 v1/v2 记录会被拒绝。旧 reader 拒绝 v3,避免静默丢弃 +失败尝试。历史 Journal 不会被改写。 diff --git a/docs/user_docs/en/cli_reference.md b/docs/user_docs/en/cli_reference.md index 679ea44..a45318d 100644 --- a/docs/user_docs/en/cli_reference.md +++ b/docs/user_docs/en/cli_reference.md @@ -95,6 +95,21 @@ Headless mode still exits `0`; it does not retry or continue automatically. Interactive mode returns to `You>` with the partial text and a truncation notice in conversation history, so a later user message can continue the conversation. +## Interrupted responses + +Both modes retry response-body read errors, read timeouts, and remote protocol +errors at most twice, waiting 1 and 2 seconds. Retries use the same conversation; +completed tools are not rerun, and partial replies and tool arguments are not +added to request history. Failed attempts do not consume `--max-turns`, but may +still incur provider charges. Each attempt remains in the Journal and trajectory. + +No new retry is scheduled beyond 300 seconds after the first interruption of +that reply. This window does not cancel an in-flight request: the SDK timeout +and any external run deadline still apply. Failures before a response stream +opens remain subject to the SDK's own retry policy. Authentication errors, +invalid requests, programming errors, and user interrupts are not retried by +this recovery loop. Exhausted transport failures exit headless mode with `1`. + ## Output channels During a headless run, stdout carries the streamed model text plus echoed tool diff --git a/docs/user_docs/zh-CN/cli_reference.md b/docs/user_docs/zh-CN/cli_reference.md index ebd03e3..062758c 100644 --- a/docs/user_docs/zh-CN/cli_reference.md +++ b/docs/user_docs/zh-CN/cli_reference.md @@ -83,6 +83,18 @@ nanoPyCodeAgent -p "fix the failing tests" `response_truncated`。Headless 模式仍退出 `0`,不会自动重试或续写。交互模式会 返回 `You>`,会话历史中保留部分文本和截断提示,用户可以在后续消息中继续对话。 +## 响应中断 + +两种模式都会对响应体读取错误、读取超时及远端协议错误最多重试两次,分别等待 +1 秒和 2 秒。重试使用同一份对话;已完成的工具不会重新执行,未完成的回复和 +工具参数不会加入请求历史。失败尝试不消耗 `--max-turns`,但仍可能产生服务商费用。 +每次尝试都会保留在 Journal 和轨迹中。 + +从该轮首次中断起超过 300 秒后,不再安排新的重试。这个窗口不会取消进行中的 +请求:SDK 超时和外部执行期限仍然有效。响应流打开前的失败由 SDK 自身的重试 +策略处理。这个恢复循环不会重试认证失败、无效请求、程序错误或用户中断。 +传输重试耗尽时,headless 模式退出码为 `1`。 + ## 输出通道 Headless run 期间,stdout 包含流式模型文本以及回显的工具调用和工具结果。启动 banner、 diff --git a/src/nanopycodeagent/agent.py b/src/nanopycodeagent/agent.py index 799172b..cd6d9e3 100644 --- a/src/nanopycodeagent/agent.py +++ b/src/nanopycodeagent/agent.py @@ -11,12 +11,10 @@ tool to run shell commands; every call and its output are echoed to the terminal as they happen. -The interactive loop handles only the happy path: anything unexpected — a -network error, a Ctrl-C mid-turn — crashes the session, and restarting it is -the recovery. That trade keeps the core flow readable; the hardened variant -it replaced is preserved at the ``hardened-agent-loop`` tag. A headless run -has no one to restart it, so it catches API errors, reports them verbatim, -and turns them into an exit code. +Both modes retry interrupted response streams before changing conversation +history or running tools. Other unexpected failures and Ctrl-C mid-turn still +end an interactive session. Headless mode catches exhausted API/transport +errors, reports them verbatim, and turns them into an exit code. """ import json @@ -39,7 +37,6 @@ pass import anthropic -import httpx from anthropic.types import MessageParam, ToolResultBlockParam, ToolUseBlock from .atif import project_atif, write_atif @@ -63,12 +60,18 @@ from .settings import DEFAULT_MAX_TOKENS, load_settings_env, resolve_max_tokens from .terminal import Spinner, print_tool_output, print_tool_use from .tool_validation import TOOLS, tool_input_error +from .transport import HTTP_ERRORS, RETRYABLE_STREAM_ERRORS from .write_tool import content_preview, run_write # The model used when ANTHROPIC_MODEL is set in neither the environment nor # the config file. DEFAULT_MODEL = "claude-sonnet-4-6" +# The SDK retries failures before streaming begins. Recover interrupted response +# bodies here, before committing a reply to history or executing any of its tools. +STREAM_RETRY_DELAYS = (1.0, 2.0) +STREAM_RETRY_WINDOW_SECONDS = 300.0 + _TRUNCATION_NOTICE = ( "[response truncated: reached max_tokens; stopped without finishing the task. " "Tool calls from this response were not executed.]" @@ -205,6 +208,14 @@ def __call__(self, event: NativeEvent) -> None: print() if event.payload["stop_reason"] == "max_tokens": print(_TRUNCATION_NOTICE, file=sys.stderr) + elif event.type == "model.failed" and event.payload["will_retry"]: + if str(event.payload["model_call_id"]) in self._model_calls_with_text: + print() + print( + f"[response interrupted: {event.payload['error_type']}; " + f"retrying in {event.payload['retry_delay_seconds']:g}s]", + file=sys.stderr, + ) elif event.type == "tool.started": tool_name = str(event.payload["tool_name"]) arguments = event.payload["input"] @@ -480,6 +491,8 @@ def _run_model_loop( ) -> RunOutcome: """Run model replies and tool calls for an already-started Agent Run.""" turns = 0 + retries = 0 + retry_deadline = None while True: # A spinner marks the wait for the reply; the first streamed # token replaces it with the reply prefix. A tool-only reply @@ -496,35 +509,73 @@ def _run_model_loop( model_started_ns = time.perf_counter_ns() # Stream the reply so text shows up as it is generated, then grab # the accumulated message for the conversation history. - with Spinner() as spinner, client.messages.stream( - model=model, - max_tokens=max_tokens, - system=system, - tools=TOOLS, - messages=messages, - ) as stream: - input_json: dict[int, list[str]] = {} - for event in stream: - if event.type == "text": - spinner.stop() - emitter.emit( - "model.output_delta", - { - "model_call_id": model_call_id, - "delta": event.text, - "source_timestamp": utc_now(), - }, - ) - elif ( - event.type == "content_block_delta" - and event.delta.type == "input_json_delta" - ): - input_json.setdefault(event.index, []).append( - event.delta.partial_json - ) - message = stream.get_final_message() - generation_id = _response_header(stream, "x-generation-id") - model_completed_ns = time.perf_counter_ns() + generation_id = None + stream_entered = False + try: + with Spinner() as spinner, client.messages.stream( + model=model, + max_tokens=max_tokens, + system=system, + tools=TOOLS, + messages=messages, + ) as stream: + stream_entered = True + generation_id = _response_header(stream, "x-generation-id") + input_json: dict[int, list[str]] = {} + for event in stream: + if event.type == "text": + spinner.stop() + emitter.emit( + "model.output_delta", + { + "model_call_id": model_call_id, + "delta": event.text, + "source_timestamp": utc_now(), + }, + ) + elif ( + event.type == "content_block_delta" + and event.delta.type == "input_json_delta" + ): + input_json.setdefault(event.index, []).append( + event.delta.partial_json + ) + message = stream.get_final_message() + model_completed_ns = time.perf_counter_ns() + except (anthropic.APIError, *HTTP_ERRORS) as exc: + now = time.monotonic() + if retry_deadline is None: + retry_deadline = now + STREAM_RETRY_WINDOW_SECONDS + delay = STREAM_RETRY_DELAYS[retries] if retries < len(STREAM_RETRY_DELAYS) else 0 + will_retry = ( + stream_entered + and isinstance(exc, RETRYABLE_STREAM_ERRORS) + and retries < len(STREAM_RETRY_DELAYS) + and now + delay <= retry_deadline + ) + emitter.emit( + "model.failed", + { + "model_call_id": model_call_id, + "error_type": type(exc).__name__, + "message": str(exc), + "generation_id": generation_id, + "duration_ms": (time.perf_counter_ns() - model_started_ns) / 1_000_000, + "will_retry": will_retry, + "retry_delay_seconds": delay if will_retry else 0, + "source_timestamp": utc_now(), + }, + ) + if not will_retry: + raise + # The window limits when another retry may start. An in-flight + # attempt retains the SDK timeout and any external run deadline. + time.sleep(delay) + retries += 1 + continue + + retries = 0 + retry_deadline = None content = _native_content_blocks(message.content) input_errors: dict[str, str] = {} @@ -637,10 +688,14 @@ def _reconcile_costs( if entry.type == "model.cost_resolved" } for entry in entries: - if entry.type != "model.completed": + if entry.type not in {"model.completed", "model.failed"}: continue generation_id = entry.payload.get("generation_id") - cost = entry.payload.get("cost") + cost = ( + pending_cost(generation_id) + if entry.type == "model.failed" + else entry.payload.get("cost") + ) if ( not isinstance(generation_id, str) or generation_id in already_resolved @@ -754,7 +809,7 @@ def run_headless( reply_prefix="", trajectory_path=trajectory_path, ) - except (anthropic.APIError, httpx.HTTPError) as exc: + except (anthropic.APIError, *HTTP_ERRORS) as exc: # Printed verbatim on purpose: a harness classifies a failed run by # pattern-matching this text (rate limit, overloaded, context length, # …) to decide whether retrying is worth anything. Rewording it, or diff --git a/src/nanopycodeagent/atif.py b/src/nanopycodeagent/atif.py index a3b7ab1..af58cd1 100644 --- a/src/nanopycodeagent/atif.py +++ b/src/nanopycodeagent/atif.py @@ -298,6 +298,53 @@ def project_atif(entries: Sequence[JournalEntry]) -> JsonObject: model_call_id = payload["model_call_id"] assert isinstance(model_call_id, str) model_deltas.setdefault(model_call_id, []).append(entry) + elif entry.type == "model.failed": + model_call_id = str(payload["model_call_id"]) + started = model_starts[model_call_id] + completed_model_calls.add(model_call_id) + timestamp, timestamp_source = _timestamp(entry) + started_at, started_at_source = _timestamp(started) + deltas = model_deltas.get(model_call_id, []) + extra: JsonObject = { + **payload, + "incomplete": True, + "started_at": started_at, + "started_at_source": started_at_source, + "timestamp_source": timestamp_source, + } + _add_journal_truncation(extra, entry) + truncated_deltas = [ + {"journal_seq": delta.seq, "truncation": delta.truncation} + for delta in deltas if delta.truncation is not None + ] + if truncated_deltas: + extra["journal_truncations"] = truncated_deltas + step: JsonObject = { + "step_id": len(steps) + 1, + "timestamp": timestamp, + "source": "agent", + "model_name": started.payload["model"], + "message": "".join(str(delta.payload["delta"]) for delta in deltas), + "llm_call_count": 1, + "extra": extra, + } + generation_id = payload["generation_id"] + assert generation_id is None or isinstance(generation_id, str) + resolved_cost = resolved_costs.get(generation_id) + amount = None + if resolved_cost is not None: + amount = Decimal(str(resolved_cost["amount"])) + step["metrics"] = { + "cost_usd": float(amount), + "extra": { + "cost_source": resolved_cost["source"], + "generation_id": generation_id, + }, + } + # An interrupted generation may still be billed. Its missing usage + # and unresolved cost must survive a successful retry. + cost_states.append((generation_id, amount)) + steps.append(step) elif entry.type == "model.completed": model_call_id = payload["model_call_id"] assert isinstance(model_call_id, str) @@ -447,7 +494,10 @@ def project_atif(entries: Sequence[JournalEntry]) -> JsonObject: llm_steps = [step for step in steps if step.get("llm_call_count") == 1] final_metrics = trajectory["final_metrics"] assert isinstance(final_metrics, dict) - if llm_steps and len(step_metrics) == len(llm_steps): + if llm_steps and len(step_metrics) == len(llm_steps) and all( + all(field in metrics for field in ("prompt_tokens", "completion_tokens", "cached_tokens")) + for metrics in step_metrics + ): final_metrics["total_prompt_tokens"] = sum( int(metrics["prompt_tokens"]) for metrics in step_metrics ) diff --git a/src/nanopycodeagent/event_journal.py b/src/nanopycodeagent/event_journal.py index edac9d0..ae4d5c7 100644 --- a/src/nanopycodeagent/event_journal.py +++ b/src/nanopycodeagent/event_journal.py @@ -19,8 +19,8 @@ from . import settings -SCHEMA_VERSION = 2 -SUPPORTED_SCHEMA_VERSIONS = frozenset({1, 2}) +SCHEMA_VERSION = 3 +SUPPORTED_SCHEMA_VERSIONS = frozenset({1, 2, 3}) DEFAULT_MAX_STRING_CHARS = 100_000 type RunOutcome = Literal["completed", "max_turns_exhausted", "response_truncated"] @@ -32,6 +32,7 @@ "model.started", "model.output_delta", "model.completed", + "model.failed", "model.cost_resolved", "tool.started", "tool.completed", @@ -91,6 +92,10 @@ "model.cost_resolved": frozenset( {"generation_id", "amount", "currency", "source", "source_timestamp"} ), + "model.failed": frozenset( + {"model_call_id", "error_type", "message", "generation_id", "duration_ms", + "will_retry", "retry_delay_seconds", "source_timestamp"} + ), "tool.started": frozenset( {"tool_call_id", "tool_name", "input", "source_timestamp"} ), @@ -332,6 +337,18 @@ def _validate_native_payload(event_type: str, payload: JsonObject) -> None: elif event_type == "model.output_delta": if not isinstance(payload["delta"], str): raise ValueError("model.output_delta.delta must be a string") + elif event_type == "model.failed": + _require_string(payload, "error_type", event_type) + if not isinstance(payload["message"], str): + raise ValueError("model.failed.message must be a string") + generation_id = payload["generation_id"] + if generation_id is not None and (not isinstance(generation_id, str) or not generation_id): + raise ValueError("model.failed.generation_id must be a string or null") + if not isinstance(payload["will_retry"], bool): + raise ValueError("model.failed.will_retry must be a boolean") + delay = payload["retry_delay_seconds"] + if not isinstance(delay, int | float) or isinstance(delay, bool) or delay < 0: + raise ValueError("model.failed.retry_delay_seconds must be non-negative") elif event_type == "model.completed": _require_string(payload, "model", event_type) if not isinstance(payload["content"], list): @@ -531,6 +548,8 @@ def from_dict(cls, value: Mapping[str, object]) -> JournalEntry: if truncation is not None: _validate_truncation(truncation) event = NativeEvent(event_type, payload) + if schema_version < 3 and event.type == "model.failed": + raise ValueError("model.failed requires Journal Entry schema 3") if ( schema_version == 1 and event.type == "run.completed" diff --git a/src/nanopycodeagent/transport.py b/src/nanopycodeagent/transport.py new file mode 100644 index 0000000..be4af14 --- /dev/null +++ b/src/nanopycodeagent/transport.py @@ -0,0 +1,18 @@ +"""Transport exceptions across the HTTP families supported by the SDK.""" + +import httpx + +try: + import httpx2 +except ImportError: + # Older Anthropic SDKs use httpx and do not install httpx2. + _http_modules = (httpx,) +else: + _http_modules = (httpx, httpx2) + +HTTP_ERRORS = tuple(module.HTTPError for module in _http_modules) +RETRYABLE_STREAM_ERRORS = tuple( + error + for module in _http_modules + for error in (module.ReadError, module.ReadTimeout, module.RemoteProtocolError) +) diff --git a/tests/test_agent_events.py b/tests/test_agent_events.py index bdfb866..723a592 100644 --- a/tests/test_agent_events.py +++ b/tests/test_agent_events.py @@ -293,6 +293,7 @@ def test_failed_tool_events_project_the_existing_tool_output( def test_interrupted_model_stream_preserves_partial_stdout_and_records_failure( monkeypatch, capsys ): + monkeypatch.setattr(agent, "STREAM_RETRY_DELAYS", ()) class DisconnectingStream(FakeStream): @property def text_stream(self): @@ -322,6 +323,7 @@ def _chunks(): "user.message", "model.started", "model.output_delta", + "model.failed", "run.failed", ] failure = entries[-1].payload diff --git a/tests/test_cli.py b/tests/test_cli.py index db5df0c..3e61177 100644 --- a/tests/test_cli.py +++ b/tests/test_cli.py @@ -238,6 +238,7 @@ def stream(self, **kwargs): def test_stream_transport_error_is_reported_verbatim_and_exits_non_zero( monkeypatch, capsys ): + monkeypatch.setattr("nanopycodeagent.agent.STREAM_RETRY_DELAYS", ()) class DisconnectingStream(FakeStream): @property def text_stream(self): @@ -265,6 +266,7 @@ def _gen(): def test_failed_headless_run_writes_partial_atif_trajectory( monkeypatch, tmp_path, capsys ): + monkeypatch.setattr("nanopycodeagent.agent.STREAM_RETRY_DELAYS", ()) class DisconnectingStream(FakeStream): @property def text_stream(self): diff --git a/tests/test_event_journal.py b/tests/test_event_journal.py index 3b90cab..4b5298a 100644 --- a/tests/test_event_journal.py +++ b/tests/test_event_journal.py @@ -54,7 +54,7 @@ def test_journal_entry_wraps_the_native_event_with_ordering_metadata(tmp_path): }, } assert entry.to_dict() == { - "schema_version": 2, + "schema_version": 3, "run_id": "run-123", "seq": 1, "recorded_at": "2026-08-23T08:00:01.420Z", @@ -336,7 +336,7 @@ def test_native_event_contract_rejects_non_json_values(content): ) -@pytest.mark.parametrize("schema_version", [True, 0, 3, "2"]) +@pytest.mark.parametrize("schema_version", [True, 0, 4, "2"]) def test_journal_entry_rejects_unsupported_schema_version(schema_version): with pytest.raises(ValueError, match="unsupported Journal Entry schema"): JournalEntry.from_dict( diff --git a/tests/test_stream_recovery.py b/tests/test_stream_recovery.py new file mode 100644 index 0000000..0f288a9 --- /dev/null +++ b/tests/test_stream_recovery.py @@ -0,0 +1,226 @@ +"""Recover response-body failures without replaying tools or partial replies.""" + +import json +from types import SimpleNamespace + +import anthropic +import httpx +import pytest + +from nanopycodeagent import agent, settings +from nanopycodeagent.event_journal import EventJournal, JournalEntry + +from helpers import FakeClient, FakeMessages, FakeStream, patch_client, sdk_http_module, text_block + + +def entries(): + path, = (settings.SETTINGS_PATH.parent / "journals").glob("*.jsonl") + return EventJournal.replay(path) + + +class BrokenStream(FakeStream): + def __init__(self, error): + super().__init__([]) + self.error = error + self.closed = False + + def __iter__(self): + yield SimpleNamespace(type="text", text="partial") + raise self.error + + def __exit__(self, *args): + self.closed = True + + +@pytest.mark.parametrize("family", ["httpx", "httpx2"]) +@pytest.mark.parametrize("error_name", ["ReadError", "ReadTimeout", "RemoteProtocolError"]) +def test_transport_families_recover_and_preserve_failed_attempt( + monkeypatch, tmp_path, family, error_name +): + http = pytest.importorskip(family) + sleeps = [] + monkeypatch.setattr(agent.time, "sleep", sleeps.append) + broken = BrokenStream(getattr(http, error_name)("disconnected")) + messages = FakeMessages([broken, [text_block("done")]]) + patch_client(monkeypatch, FakeClient(messages)) + trajectory_path = tmp_path / "trajectory.json" + assert agent.run_headless("task", max_turns=1, trajectory_path=trajectory_path) == 0 + assert messages.calls[0] == messages.calls[1] + assert broken.closed + assert sleeps == [1.0] + failed = [e for e in entries() if e.type == "model.failed"] + assert len(failed) == 1 and failed[0].payload["will_retry"] is True + trajectory = json.loads(trajectory_path.read_text()) + assert [s["message"] for s in trajectory["steps"]] == ["task", "partial", "done"] + assert trajectory["steps"][1]["extra"]["incomplete"] is True + assert trajectory["extra"]["terminal"]["outcome"] == "completed" + assert trajectory["final_metrics"]["extra"]["usage_complete"] is False + assert trajectory["final_metrics"]["extra"]["cost_is_partial"] is True + + +@pytest.mark.parametrize("family", ["httpx", "httpx2"]) +def test_exhausted_retries_exit_cleanly_and_keep_all_attempts(monkeypatch, capsys, family): + http = pytest.importorskip(family) + sleeps = [] + monkeypatch.setattr(agent.time, "sleep", sleeps.append) + streams = [BrokenStream(http.RemoteProtocolError("disconnected")) for _ in range(3)] + messages = FakeMessages(streams) + patch_client(monkeypatch, FakeClient(messages)) + assert agent.run_headless("task") == 1 + assert len(messages.calls) == 3 + assert all(stream.closed for stream in streams) + assert sleeps == [1.0, 2.0] + assert "API error: disconnected" in capsys.readouterr().err + failures = [e for e in entries() if e.type == "model.failed"] + assert [e.payload["will_retry"] for e in failures] == [True, True, False] + assert len({e.payload["model_call_id"] for e in failures}) == 3 + assert entries()[-1].type == "run.failed" + + +def test_retry_window_stops_new_attempts_after_slow_retry(monkeypatch): + clock = iter([0.0, 301.0]) + sleeps = [] + monkeypatch.setattr(agent, "time", SimpleNamespace( + perf_counter_ns=agent.time.perf_counter_ns, + monotonic=lambda: next(clock), sleep=sleeps.append, + )) + messages = FakeMessages([BrokenStream(httpx.ReadError("broken")) for _ in range(2)]) + patch_client(monkeypatch, FakeClient(messages)) + assert agent.run_headless("task") == 1 + assert len(messages.calls) == 2 + assert sleeps == [1.0] + + +@pytest.mark.parametrize("error", [ValueError("bug"), KeyboardInterrupt()]) +def test_programming_errors_and_interrupts_are_not_retried(monkeypatch, error): + messages = FakeMessages([BrokenStream(error)]) + patch_client(monkeypatch, FakeClient(messages)) + with pytest.raises(type(error)): + agent.run_headless("task") + assert len(messages.calls) == 1 + + +def test_permanent_http_errors_are_not_retried(monkeypatch): + messages = FakeMessages([BrokenStream(httpx.LocalProtocolError("bad request"))]) + patch_client(monkeypatch, FakeClient(messages)) + assert agent.run_headless("task") == 1 + assert len(messages.calls) == 1 + + +def test_sdk_stream_retry_discards_partial_tool_and_does_not_replay_work(monkeypatch, tmp_path): + http = sdk_http_module() + requests = [] + closed = [] + monkeypatch.setattr(agent.time, "sleep", lambda _: None) + monkeypatch.setattr(agent, "resolve_generation_cost", lambda *args, **kwargs: None) + target = tmp_path / "answer.txt" + + def wire(events): + return "".join(f"event: {e['type']}\ndata: {json.dumps(e)}\n\n" for e in events).encode() + + def start(number): + return {"type": "message_start", "message": { + "id": f"msg-{number}", "type": "message", "role": "assistant", + "model": "test-model", "content": [], "stop_reason": None, + "stop_sequence": None, "usage": {"input_tokens": 10, "output_tokens": 0}, + }} + + def finish(reason): + return [ + {"type": "content_block_stop", "index": 0}, + {"type": "message_delta", "delta": {"stop_reason": reason, "stop_sequence": None}, + "usage": {"output_tokens": 5}}, + {"type": "message_stop"}, + ] + + class Body(http.SyncByteStream): + def __init__(self, data, interrupt): + self.data, self.interrupt = data, interrupt + + def __iter__(self): + yield self.data + if self.interrupt: + raise http.RemoteProtocolError("incomplete chunked read") + + def close(self): + closed.append(self.interrupt) + + def respond(request): + requests.append(json.loads(request.content)) + number = len(requests) + # First execute one append, then interrupt the next reply after its + # complete-looking tool JSON. Reissuing that reply must not append twice. + if number == 1: + block = {"type": "tool_use", "id": "append", "name": "bash", "input": {}} + args = {"command": f"printf x >> '{target}'"} + elif number in (2, 3): + block = {"type": "tool_use", "id": "discard" if number == 2 else "read", "name": "read", "input": {}} + args = {"path": str(target)} + else: + block = {"type": "text", "text": "done"} + events = [start(number), {"type": "content_block_start", "index": 0, "content_block": block}] + if number < 4: + events.append({"type": "content_block_delta", "index": 0, + "delta": {"type": "input_json_delta", "partial_json": json.dumps(args)}}) + if number != 2: + events += finish("tool_use" if number < 4 else "end_turn") + return http.Response(200, headers={"content-type": "text/event-stream", "x-generation-id": f"gen-{number}"}, + stream=Body(wire(events), number == 2)) + + with anthropic.Anthropic(api_key="test", base_url="https://example.test", + http_client=http.Client(transport=http.MockTransport(respond))) as client: + patch_client(monkeypatch, client) + trajectory_path = tmp_path / "trajectory.json" + assert agent.run_headless("task", max_turns=3, trajectory_path=trajectory_path) == 0 + assert len(requests) == 4 + assert requests[1] == requests[2] + assert target.read_text() == "x" + assert closed.count(True) == 1 + tools = [e.payload["tool_call_id"] for e in entries() if e.type == "tool.started"] + assert tools == ["append", "read"] + failed = next(e for e in entries() if e.type == "model.failed") + assert failed.payload["generation_id"] == "gen-2" + trajectory = json.loads(trajectory_path.read_text()) + assert [s["extra"].get("incomplete", False) for s in trajectory["steps"]] == [False, False, True, False, False] + + +def test_interrupted_generation_cost_is_reconciled_without_inventing_usage(monkeypatch, tmp_path): + monkeypatch.setattr(agent.time, "sleep", lambda _: None) + broken = BrokenStream(httpx.ReadError("disconnected")) + broken.response.headers = {"x-generation-id": "gen-broken"} + reply = FakeStream([text_block("done")], usage=SimpleNamespace(input_tokens=5, output_tokens=2), + response_headers={"x-generation-id": "gen-done"}) + messages = FakeMessages([broken, reply]) + patch_client(monkeypatch, FakeClient(messages, base_url="https://openrouter.ai/api")) + looked_up = [] + + def resolve(base_url, generation_id, credential, **kwargs): + looked_up.append(generation_id) + return {"generation_id": generation_id, "amount": "0.01", "currency": "USD", "source": "test"} + + monkeypatch.setattr(agent, "resolve_generation_cost", resolve) + path = tmp_path / "trajectory.json" + assert agent.run_headless("task", trajectory_path=path) == 0 + assert looked_up == ["gen-broken", "gen-done"] + final = json.loads(path.read_text())["final_metrics"] + assert final["total_cost_usd"] == 0.02 + assert final["extra"]["usage_complete"] is False + assert "total_prompt_tokens" not in final + + +@pytest.mark.parametrize("schema", [1, 2, 3]) +def test_failed_attempt_requires_v3_schema(schema): + record = { + "schema_version": schema, "run_id": "run-1", "seq": 1, + "recorded_at": "2026-09-15T00:00:00.000Z", "type": "model.failed", + "payload": { + "model_call_id": "model-1", "error_type": "ReadError", "message": "broken", + "generation_id": None, "duration_ms": 12, "will_retry": True, + "retry_delay_seconds": 1, "source_timestamp": None, + }, + } + if schema < 3: + with pytest.raises(ValueError, match="requires Journal Entry schema 3"): + JournalEntry.from_dict(record) + else: + assert JournalEntry.from_dict(record).to_dict() == record diff --git a/tests/test_truncation.py b/tests/test_truncation.py index 68e0a74..4b96707 100644 --- a/tests/test_truncation.py +++ b/tests/test_truncation.py @@ -75,7 +75,7 @@ def resolve(base_url, generation_id, credential, **kwargs): assert reconciled == ["gen-truncated"] entries = _journal_entries() - assert all(entry.schema_version == 2 for entry in entries) + assert all(entry.schema_version == 3 for entry in entries) assert entries[-1].type == "run.completed" assert entries[-1].payload["outcome"] == "response_truncated" completed = next(entry for entry in entries if entry.type == "model.completed") @@ -182,7 +182,7 @@ def respond(request): assert "observation" not in trajectory["steps"][1] -@pytest.mark.parametrize("schema_version", [1, 2]) +@pytest.mark.parametrize("schema_version", [1, 2, 3]) @pytest.mark.parametrize("outcome", ["completed", "max_turns_exhausted", "response_truncated"]) def test_journal_outcome_versions(schema_version, outcome): record = { From 1ce9c8bc3d43d0ce4ff42ceae13043d8e181627d Mon Sep 17 00:00:00 2001 From: minixalpha Date: Fri, 25 Sep 2026 20:56:30 +0800 Subject: [PATCH 2/2] docs: record stream recovery benchmark results Add the full 20-task live retest to the 0.8.x dev notes, with its changelog entry and the regenerated English dev notes. --- docs/changelogs/0.8.x.md | 6 +++ docs/dev_notes/en/0.8.x.md | 92 +++++++++++++++++++++++++++++++++++ docs/dev_notes/zh-CN/0.8.x.md | 74 ++++++++++++++++++++++++++++ 3 files changed, 172 insertions(+) diff --git a/docs/changelogs/0.8.x.md b/docs/changelogs/0.8.x.md index e3fdea8..d3318ce 100644 --- a/docs/changelogs/0.8.x.md +++ b/docs/changelogs/0.8.x.md @@ -20,6 +20,12 @@ All notable changes in the **0.8.x** release series are documented here. completion. Preserve partial output, usage, and cost accounting; skip tools from the truncated reply and keep subsequent interactive requests valid. New journals use schema v2, with v1 replay still supported. +- Recover interrupted model response streams. Retry response-body read errors, + read timeouts, and remote protocol errors up to twice (1s, 2s) before + committing a reply to history or running any of its tools, handling both the + `httpx` and `httpx2` transport families. Record each attempt as a Journal v3 + `model.failed` event, project it as an incomplete ATIF step, and mark token + and cost totals as partial when usage is missing. ### Changed - Increase the default per-reply generation limit from 8192 to 32768 tokens. diff --git a/docs/dev_notes/en/0.8.x.md b/docs/dev_notes/en/0.8.x.md index 5363167..9ae1836 100644 --- a/docs/dev_notes/en/0.8.x.md +++ b/docs/dev_notes/en/0.8.x.md @@ -916,6 +916,98 @@ the scheduler whose verification resumed after the pause, total experiment model cost is **$0.705181383**. Complete billing does not imply complete native responses or usage. +#### Full 20-task retest of stream recovery + +**Run on 2026-09-21 to 22.** On commit +`577a1ab7f96fc3d220ca89c068988a479ac7b1aa` of the `fix/stream-recovery` branch, we +reran the original 20 tasks to check whether real model responses were interrupted +and whether an interruption could recover within the same run. The run started at +22:06 Beijing time on 09-21 and ended at 01:52 on 09-22, taking about **3 hours +46 minutes**; no transport faults were injected. + +**Configuration and recording.** Harbor 0.21.0 and Terminal-Bench 2.1, keeping the +original 20 task refs and verifying the cached image digests before the run. The +model was still `deepseek/deepseek-v4-flash-0731` on OpenRouter. Containers +installed the wheel built from this branch, and every Python source file in the +package was checked against the current commit; the running version was +`0.8.1.dev28+g577a1ab7f`. Dependencies were pinned to Anthropic SDK 1.5.0, HTTPX +0.28.1, HTTPX2 2.12.0, and Pydantic 2.13.5. + +Each reply was capped at 65536 tokens, with at most 50 complete replies per task +and concurrency 2; each task was attempted once, with Harbor whole-task automatic +retries set to 0. The agent was capped at 3600 seconds, with 120 seconds of +finalization grace and a 3780-second outer Harbor limit, a 3600-second install +cap, and the official verifier left at its original limit. We reused the earlier +dependency cache and passive HTTP recorder; the container did only lightweight +event counting, and full Journal / HTTP records were collected to the host for +analysis. The live counter was corrected once to accept JSON without spaces and +restarted separately; the evaluation was not restarted and the agent, tasks, and +budget were not changed. + +| Metric | This run | +| --- | --- | +| Finished trials | 20/20 | +| Actually entered agent execution | 18 tasks; the other 2 failed during install and made no model calls | +| Official scores | reward 1: **11 tasks (55%)**; reward 0: **6 tasks**; unscored: **3 tasks** | +| Completed normally and passed | **8 tasks (40%)**: native terminal state `completed`, no Harbor exception, execution limit not reached | +| Journal model calls | **576 started, 576 fully returned** | +| `model.failed` / interrupted HTTP response bodies | **0 / 0** | +| Stream recovery triggered / succeeded | **0 / 0**; the recovery branch was not exercised this run | + +Per-task results follow. "Turns exhausted" means the native terminal state was +`max_turns_exhausted`; whether the official verifier judges the artifact as passing +is a separate matter. + +| Task | Reward | Execution or verification result | +| --- | --- | --- | +| `write-compressor` | 1 | Completed normally and passed | +| `torch-tensor-parallelism` | Unscored | Agent completed normally; verifier exceeded the original 900-second limit | +| `schemelike-metacircular-eval` | 1 | 50 turns exhausted; artifact passed | +| `kv-store-grpc` | 0 | Agent completed normally, artifact failed | +| `pypi-server` | 1 | Completed normally and passed | +| `dna-assembly` | 0 | Agent completed normally, artifact failed | +| `torch-pipeline-parallelism` | Unscored | Install used Python 3.12.3, below the project's required Python >=3.13 | +| `qemu-alpine-ssh` | Unscored | Install downloaded Python 3.14.0 and still received HTTP 504 after three retries | +| `openssl-selfsigned-cert` | 1 | Completed normally and passed | +| `regex-chess` | 1 | 50 turns exhausted; artifact passed | +| `log-summary-date-ranges` | 0 | Agent completed normally, artifact failed | +| `model-extraction-relu-logits` | 1 | Completed normally and passed | +| `path-tracing` | 1 | Completed normally and passed | +| `regex-log` | 1 | Completed normally and passed | +| `caffe-cifar-10` | 0 | 50 turns exhausted, artifact failed | +| `mteb-leaderboard` | 0 | 50 turns exhausted, artifact failed | +| `llm-inference-batching-scheduler` | 1 | 50 turns exhausted; artifact passed | +| `pytorch-model-recovery` | 1 | Completed normally and passed | +| `circuit-fibsqrt` | 0 | 50 turns exhausted, artifact failed | +| `merge-diff-arc-agi-task` | 1 | Completed normally and passed | + +Both install failures were recorded by Harbor as `NonZeroAgentExitCodeError` and +happened before the agent started; the tensor exception was +`VerifierTimeoutError`. These are not model stream interruptions, and the 20 +finished trials must not be written as "all 20 tasks completed model execution". +Likewise, Scheme, regex-chess, and scheduler among the 11 reward-1 results passed +after exhausting turns and do not count toward the 8 completed-and-passed. + +**Evidence boundary for recovery.** All 576 model calls this run returned +completely, so no real stream interruption was observed; we can only record "no +stream interruption occurred this run", and cannot claim "recovery succeeded after +a live interruption". At the start of this run we separately ran +`tests/test_stream_recovery.py`, which reported **17 passed, 0 failed, 0 skipped**, +covering both HTTP exception families, retry exhaustion and the time window, +discarding tool calls from an interrupted reply, and preserving existing tool +results. These automated tests validate the recovery mechanism but cannot replace +a real interruption sample. Sampling, routing, caching, and server state were +uncontrolled; the score change from the earlier 20 tasks also cannot be attributed +to stream recovery. + +The local job is `jobs/tb21-stream-recovery-pilot20-20260921/` and the experiment +record directory is `jobs/tb21-stream-recovery-pilot20-20260921-record/`. There, +`manifest.json` stores source, wheel, task, and image information; `summary.json` +stores per-task statistics; `report.md` is the result summary; and +`stream-recovery-tests.json` stores the dedicated test result. Raw Journals, HTTP +records, trajectories, and verifier output are kept in each trial directory; these +raw data live in the Git-ignored `jobs/`. + #### Stopping execution on timeout Timeout handling must ensure that the agent inside the container actually stops before verification begins. In this pilot, Harbor had already classified the task as timed out while the agent continued calling the model. diff --git a/docs/dev_notes/zh-CN/0.8.x.md b/docs/dev_notes/zh-CN/0.8.x.md index 9669e42..9cea304 100644 --- a/docs/dev_notes/zh-CN/0.8.x.md +++ b/docs/dev_notes/zh-CN/0.8.x.md @@ -1287,6 +1287,80 @@ Journal 投影一致,1 份是上述明确标注的重建文件;全部 Journa 完整 20 题费用为 **$0.462175819**;连同定向组、缓存补测和恢复验证的 scheduler, 本次实验模型费用合计 **$0.705181383**。费用完整不等于原生响应或用量完整。 +#### 流中断恢复的完整 20 题复测 + +**2026-09-21~22 实跑。** 在 `fix/stream-recovery` 分支提交 +`577a1ab7f96fc3d220ca89c068988a479ac7b1aa` 上重新运行原始 20 题,检查真实模型 +响应是否中断,以及中断后能否在同一次 run 内恢复。北京时间 09-21 22:06 开始, +09-22 01:52 结束,总耗时约 **3 小时 46 分钟**;没有人为注入传输故障。 + +**配置与记录。** 使用 Harbor 0.21.0、Terminal-Bench 2.1,保留原始 20 题的 +task ref,并在运行前核对缓存镜像摘要。模型仍为 OpenRouter 上的 +`deepseek/deepseek-v4-flash-0731`。容器安装本分支构建的 wheel,逐个核对包内 +Python 源文件与当前提交一致;运行版本为 `0.8.1.dev28+g577a1ab7f`。依赖固定为 +Anthropic SDK 1.5.0、HTTPX 0.28.1、HTTPX2 2.12.0、Pydantic 2.13.5。 + +每次回复上限 65536 tokens,每题最多 50 次完整回复,并发 2;每题只尝试一次, +Harbor 整题自动重试为 0。agent 执行上限 3600 秒、收尾宽限 120 秒、Harbor +外层上限 3780 秒,安装上限 3600 秒,官方 verifier 保留原始时限。沿用此前的 +依赖缓存和被动 HTTP 记录器;容器内仅做轻量事件计数,完整 Journal/HTTP +记录收集到宿主机后分析。实时计数器曾修正为兼容无空格的 JSON 并单独重启, +没有重启评测或修改 agent、题目和预算。 + +| 指标 | 本次结果 | +| --- | --- | +| 已结束的 trial | 20/20 | +| 实际进入 agent 执行 | 18 题;另外 2 题安装失败,未调用模型 | +| 官方评分 | reward 1:**11 题(55%)**;reward 0:**6 题**;未评分:**3 题** | +| 正常完成且通过 | **8 题(40%)**:原生终态为 `completed`,无 Harbor 异常,未触及执行时限 | +| Journal 模型调用 | **576 次启动,576 次完整返回** | +| `model.failed` / HTTP 响应体中断 | **0 / 0** | +| 触发流恢复 / 恢复成功 | **0 / 0**;本轮没有触发恢复分支 | + +逐题结果如下。“轮数耗尽”指原生终态为 `max_turns_exhausted`;它与官方 verifier +是否判定产物通过是两件事。 + +| 题目 | Reward | 执行或验证结果 | +| --- | --- | --- | +| `write-compressor` | 1 | 正常完成并通过 | +| `torch-tensor-parallelism` | 未评分 | agent 正常完成;verifier 超过原始 900 秒时限 | +| `schemelike-metacircular-eval` | 1 | 50 轮耗尽;产物通过 | +| `kv-store-grpc` | 0 | agent 正常完成,产物未通过 | +| `pypi-server` | 1 | 正常完成并通过 | +| `dna-assembly` | 0 | agent 正常完成,产物未通过 | +| `torch-pipeline-parallelism` | 未评分 | 安装阶段使用 Python 3.12.3,不满足项目要求的 Python >=3.13 | +| `qemu-alpine-ssh` | 未评分 | 安装阶段下载 Python 3.14.0,三次重试后仍收到 HTTP 504 | +| `openssl-selfsigned-cert` | 1 | 正常完成并通过 | +| `regex-chess` | 1 | 50 轮耗尽;产物通过 | +| `log-summary-date-ranges` | 0 | agent 正常完成,产物未通过 | +| `model-extraction-relu-logits` | 1 | 正常完成并通过 | +| `path-tracing` | 1 | 正常完成并通过 | +| `regex-log` | 1 | 正常完成并通过 | +| `caffe-cifar-10` | 0 | 50 轮耗尽,产物未通过 | +| `mteb-leaderboard` | 0 | 50 轮耗尽,产物未通过 | +| `llm-inference-batching-scheduler` | 1 | 50 轮耗尽;产物通过 | +| `pytorch-model-recovery` | 1 | 正常完成并通过 | +| `circuit-fibsqrt` | 0 | 50 轮耗尽,产物未通过 | +| `merge-diff-arc-agi-task` | 1 | 正常完成并通过 | + +两次安装失败均被 Harbor 记为 `NonZeroAgentExitCodeError`,发生在 agent 启动前; +tensor 的异常为 `VerifierTimeoutError`。这些不是模型流中断,不能把 20 个 trial +结束写成“20 题都完成了模型执行”。同样,11 次 reward 1 中的 Scheme、regex-chess、 +scheduler 是轮数耗尽后产物通过,不计入 8 次正常完成且通过。 + +**恢复能力的证据边界。** 本轮 576 次模型调用均完整返回,没有观察到真实流中断, +因此只能记录“此次未出现流中断”,不能据此声称“实跑中断后恢复成功”。本轮启动时 +另行运行 `tests/test_stream_recovery.py`,结果为 **17 passed,0 failed,0 skipped**, +覆盖两种 HTTP 异常族、重试耗尽和时间窗口、丢弃中断回复中的工具调用,以及保留 +已有工具结果等行为。这些自动化测试提供恢复机制的验证,但不能替代真实中断样本。 +采样、路由、缓存和服务端状态未受控;与此前 20 题的分数变化也不能归因于流恢复。 + +本地 job 为 `jobs/tb21-stream-recovery-pilot20-20260921/`,实验记录目录为 +`jobs/tb21-stream-recovery-pilot20-20260921-record/`。其中 `manifest.json` 保存源码、 +wheel、任务与镜像信息,`summary.json` 保存逐题统计,`report.md` 为结果摘要, +`stream-recovery-tests.json` 保存专项测试结果。原始 Journal、HTTP 记录、trajectory +和 verifier 输出保留在各 trial 目录;这些原始数据位于 Git 忽略的 `jobs/` 中。 + #### 超时停止 “超时停止”指的是:任务达到规定的运行时间后,要确保容器里的 agent 真正停止