feat(stream): claude→responses 空流兜底 + 首 token 前自动重试 - #31
Merged
Merged
Conversation
修法对照 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>
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
背景
线上观察到 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 兜底细分(原
claudeResponsesStreamHandlerEOF 分支)把「一刀切
emitError(stream ended without completion)」细分为三档:reasoningFallback兜底把思考提升为 text,agent 拿到思考不空手emitError("upstream ended stream before any content (empty completion)")并 透出retryable信号Phase B — 首 token 前自动重试(核心)
peekFirstOutput(claude_responses.go)stream\_first\_byte\_timeout\_ms窗口内窥视首个完整 SSE 帧;EOF/超时/仅错误帧判errStreamIncompleteNoCommitclaudeResponsesStreamHandler(committed bool, err error);peek 的成功路径把消费行原样喂回主循环claudeResponsesStreamWithRetryforwardClaudeViaResponses/probeClaudeViaResponses拿到 retryable 错误后,用同一份 claudeReq 经callOpenCodeEndpoint重发(自动切到 key_pool 下一个可用 key),最多stream\_empty\_retry\_max次streamReader.enableKeepalive关键不变量
peekOutcome.reader传递。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\_Mapping、TestClaudeStream\_PartialEOF\_NoFinish\_ErrorOnly等)全过:make build / vet / fmt / test全绿。不做的事(明确边界)
max\_tokens\_cap配置:不动。它走的是response.incomplete→stop\_reason: "max\_tokens"的正常收尾路径,与本故障无关。真实运行验证
docs/CONFIGURATION.md里附了复现脚本:用反代在 5s 时杀 tunnel,旧网关会让 agent 中断在stream ended without completion,新网关在第一个upstream 200之后的空 peek 窗口里透明换 key 重试,整轮圆满完成。🤖 Generated with Claude Code