Skip to content

feat(stream): claude→responses 空流兜底 + 首 token 前自动重试 - #31

Merged
6Kmfi6HP merged 9 commits into
mainfrom
feat/stream-empty-retry
Sep 27, 2026
Merged

6Kmfi6HP merged 9 commits into
mainfrom
feat/stream-empty-retry

Conversation

@6Kmfi6HP

Copy link
Copy Markdown
Owner

背景

线上观察到 agent 在响应中间被 stream ended without completion 中断。定位根因(详见 commit):上游 SSE 在整个流里既没发 response.completed / response.incomplete,也没产生任何 text/tool delta,就给了干净 EOF——典型的「prefill 阶段 tunnel 被反代(CF/nginx) 杀掉」故障形态。

参考 ParalonCloud《Stream ended without finish_reason》给的范式:重试窗口只在「向客户端写过第一个字节之前」存在;已经成功 commit 的流不能重试,必须如实告知。原代码所有流式 handler 都先 WriteHeader(200) 再看上游,架构层面没有这个窗口——这是要补的关键基建。

参考实现:new-api #3275/#4067、litellm router retry 策略。

方案(Phase A + Phase B 一次到位)

Phase A — EOF 兜底细分(原 claudeResponsesStreamHandler EOF 分支)

把「一刀切 emitError(stream ended without completion)」细分为三档:

分支 行为
已有 text/tool 产出 合成正常 stop(ParalonCloud Rule 2,不变)
仅有 thinking 无 text reasoningFallback 兜底把思考提升为 text,agent 拿到思考不空手
完全空流 emitError("upstream ended stream before any content (empty completion)") 并 透出 retryable 信号

Phase B — 首 token 前自动重试(核心)

改动 内容
peekFirstOutput(claude_responses.go) 拿到上游 200 后先不 WriteHeader,在 stream\_first\_byte\_timeout\_ms 窗口内窥视首个完整 SSE 帧;EOF/超时/仅错误帧判 errStreamIncompleteNoCommit
claudeResponsesStreamHandler 改返回 (committed bool, err error);peek 的成功路径把消费行原样喂回主循环
claudeResponsesStreamWithRetry forwardClaudeViaResponses / probeClaudeViaResponses 拿到 retryable 错误后,用同一份 claudeReq 经 callOpenCodeEndpoint 重发(自动切到 key_pool 下一个可用 key),最多 stream\_empty\_retry\_max 次
streamReader.enableKeepalive peek 窗口用零心跳避免误产 keepalive-commit;commit 后再切到常规 15s ticker

关键不变量

  • peek 消费严格落在帧边界(空行或 EOF),绝不截断中间——否则后续主循环 bufio 会把帧拼错。
  • peek 的 reader 与主循环共享(bufio 已预读字节不能丢),通过 peekOutcome.reader 传递。
  • 不会在已 commit 后再触发重试,agent 永远要么拿到完整的「干净成功流」,要么 502 进上层 fallback,不会看到「半截 error 然后又突然变正常」。
  • prompt\_cache\_key 在重试间稳定,首次 attempt 已建立的 prefix cache 仍可命中,input tokens 接近零成本。

新增配置(config.json,默认开启,不进 admin UI)

{
  "stream_empty_retry_max": 1,
  "stream_first_byte_timeout_ms": 30000
}
  • stream_empty\_retry\_max:重试次数(0=关闭)。默认 1 是「该做的兜底」而非用户选择——上游空流无毒可重试。
  • stream\_first\_byte\_timeout\_ms:首字节看门狗,默认 30s(覆盖大多数 prefill 时间,<=0 关闭看门狗、仅 EOF 触发)。

测试

internal/app/claude_responses_empty_retry_test.go 加 8 个端到端回归,全部通过:

✓ TestClaudeResponsesStream_EmptyEOF_RetriesOnce              — 空流 → 重试成功
✓ TestClaudeResponsesStream_EmptyEOF_ExhaustsRetryThenErrors  — 2次空流 → fallback 502
✓ TestClaudeResponsesStream_EmptyEOF_RetryDisabled            — 关配置 → 不重试
✓ TestClaudeResponsesStream_PartialEOF_SynthesizesStopNoRetry — 半截流不重试
✓ TestClaudeResponsesStream_ThinkingOnlyEOF_PromotesReasoning — 仅thinking兜底
✓ TestClaudeResponsesStream_ErrorFrame_Retries                — response.failed → 重试
✓ TestClaudeResponsesStream_SilentUpstream_TimesOut           — 只发心跳不发data → 超时重试
✓ TestClaudeResponsesStream_HungUpstream_FirstByteWatchdogFires — 完全挂死 → 看门狗救人

外加不破坏既有用例(TestClaudeResponsesStream\_Mapping、TestClaudeStream\_PartialEOF\_NoFinish\_ErrorOnly 等)全过:make build / vet / fmt / test 全绿。

不做的事(明确边界)

  • 已 commit 后的 EOF/断流:不改。仍按现有 emitError 兜底(ParalonCloud Rule 2),agent 拿到的就是已有的半截;这类只能靠 agent 框架自身重试。
  • 其他流式 handler(chat→anthropic、responses→anthropic、anthropic 直通):本次不动,先把出问题的 claude→responses 一条修好;模式稳定后再说复制。
  • max\_tokens\_cap 配置:不动。它走的是 response.incomplete → stop\_reason: "max\_tokens" 的正常收尾路径,与本故障无关。
  • key pool failover 计数:空流不计入 key 失败(这是协议故障不是 key 失效)。

真实运行验证

docs/CONFIGURATION.md 里附了复现脚本:用反代在 5s 时杀 tunnel,旧网关会让 agent 中断在 stream ended without completion,新网关在第一个 upstream 200 之后的空 peek 窗口里透明换 key 重试,整轮圆满完成。

🤖 Generated with Claude Code

6Kmfi6HP and others added 9 commits September 27, 2026 22:58
修法对照 ParalonCloud「prefill 阶段 tunnel 被宰」三规则:
- R1 首字节前可重试:peekFirstOutput 窥视首个完整 SSE 帧,EOF/超时/上游
  error 帧一律按未 commit 处理,由 claudeResponsesStreamWithRetry 换 key 重发
- R2 已 commit 不重试:已有 text/tool 产出时 EOF 仍合成 stop(原行为)
- 兜底:仅 thinking 没有 text 时 reasoningFallback 提升为 text,agent 至少
  拿到思考内容而不是空手中断

config 新增 stream_empty_retry_max(默认 1)、stream_first_byte_timeout_ms
(默认 30000),仅 config.json 控制,不进 admin UI。重试经 callOpenCodeEndpoint
自动落到 key_pool 下一个 key,不会重复同一根死 pipe;input tokens 在缓存
命中场景下重试近零成本。

测试:claude_responses_empty_retry_test.go 8 用例覆盖空流/挂死/错误帧/
部分流兜底;既有 TestClaudeResponsesStream_Mapping 不回归。`make build /
vet / fmt / test` 全绿。

参考:ParalonCloud《Stream ended without finish_reason》规则、new-api #3275。

Co-Authored-By: Claude Code <noreply@anthropic.com>
Splits the retry/peek logic out of claude_responses.go so subsequent PRs
can plug in the same machinery for chat completions, responses passthrough,
and anthropic passthrough without duplicating the driver. The claude responses
path keeps byte-identical behavior via thin shims.

Co-Authored-By: Claude Code <noreply@anthropic.com>
chat.go stream branch now goes through DriveStreamWithRetry with
ChatProtocolHooks — empty EOF / first-byte timeout / upstream error
frames observed before the first client write trigger a key-pool retry
instead of being silently forwarded. The pre-retry stream loop body has
been extracted into chatStreamRunOnce; both peek and main loop share the
same streamReader (with deferred keepalive enabled after peek commit) so
that the first-byte watchdog and the SSE keepalive ticker do not fight
each other.

Defer WriteHeader(http.StatusOK) until the first frame is actually
written ("commit point"), so retry (= WriteHeader + JSON error) is still
possible after peek reports an empty/error upstream.

Also fix rawSSEReader.Close to release the source via sync.Once before
acquiring r.mu — previously a watchdog-fired Close could deadlock with
an in-flight Read that held the mutex while blocked inside ReadString.

Tests (chat_stream_retry_test.go) mirror the claude→responses suite:
empty-EOF retry once / exhausts / disabled, partial-EOF emits
upstream_truncated error frame without retry, error frame triggers
retry, silent upstream timeouts, hung upstream first-byte watchdog.

Co-Authored-By: Claude Code <noreply@anthropic.com>
…ream

responsesSSEToChatStream + forwardChatViaResponses + tests

Co-Authored-By: Claude Code <noreply@anthropic.com>
…ream

anthropicSSEToChatStream + forwardChatViaAnthropic + tests

Co-Authored-By: Claude Code <noreply@anthropic.com>
responsesStreamHandler + responsesHandler call site + tests

Co-Authored-By: Claude Code <noreply@anthropic.com>
relayResponsesStream + pipeAnthropicStream + tests

Co-Authored-By: Claude Code <noreply@anthropic.com>
…k, status passthrough

Review pass over the 6-commit retry rollout surfaced five structural bugs:

F1  chat streaming swallowed upstream 4xx/5xx (status + body) into a
    generic 502 message. Wrap callOnce with UpstreamErrorCapture so
    non-2xx responses are kept (RC not closed by driver) and replayed
    verbatim to the client.

F2  pipeAnthropicStream returned (false, errStreamIncompleteNoCommit)
    on EOF-after-peek even though WriteHeader + peeked bytes were
    already on the wire — Drive would then WriteHeader+replay again
    on the same connection. After commit, always finish with a
    synthesized message_stop; never return committed=false.

F3  anthropicSSEToChatStream had the symmetric bug via writeHeaderOnce
    firing on the first non-empty line (even a ping or content_block_start
    that does not tick sentRole). Returning (false, ...) after that
    lets the driver run a duplicate WriteHeader+replay. Switched the
    EOF branch's early-return guard from !sentRole to !wroteHeader.

F4  PeekFirstFrame leaked the streamReader goroutine when ctx was
    canceled mid-peek (readCh never receives, done never closed). Now
    explicitly reader.Close() on the ctx.Done branch like the other
    early exits.

F5  chatStreamRunOnce recorded usage-bearing chunks regardless of
    whether commit had happened. When the attempt was then aborted
    (peek-EOF path), the next attempt would re-stat the same upstream
    usage. Gate RecordChatUsage on wasCommitted.

F6  relayResponsesStream returned (false, ctx.Err()) from its lineLoop
    ctx branch even when WriteHeader had already been issued. Drive
    short-circuits on ctx errors so this didn't double-write in
    practice, but the contract violation was structural — the function
    now returns (true, ctx.Err()).

Also fixes a lower-severity case where FlushPeekedBytes errors after
WriteHeader were propagated as (false, err), for the same reason.

Co-Authored-By: Claude Code <noreply@anthropic.com>
Pre-existing latent bug amplified by the chat retry refactor's emitLine
path: each chunk was being terminated with a single \n, which is inside-
frame continuation line boundary in SSE — OpenAI's strict client (python
sdk ≥1.x) parses `data:` blocks separated by `\n\n` only and would
loudly JSONDecodeError "Extra data" once a stream had more than one
frame.

Caught by running the new end-to-end OpenAI SDK E2E test against the
real running gateway. Verified: 7/7 chat SDK tests now pass against a
live upstream (mimo-v2.6-flash).

Also tested:
- launch claude → /v1/messages (anthropic frames incl. message_start,
  content_block, message_delta, message_stop) → works
- launch codex exec → /v1/responses → works

Co-Authored-By: Claude Code <noreply@anthropic.com>
@6Kmfi6HP
6Kmfi6HP merged commit 39319dc into main Sep 27, 2026
1 check passed
@6Kmfi6HP
6Kmfi6HP deleted the feat/stream-empty-retry branch September 27, 2026 23:40
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant