Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 5 additions & 1 deletion config.example.json
Original file line number Diff line number Diff line change
Expand Up @@ -29,5 +29,9 @@
{"id": "k1", "key": "sk-placeholder-1", "group": "", "weight": 1, "enabled": true, "note": "account A"},
"sk-placeholder-2"
]
}
},
"_comment_stream_empty_retry": "claude→responses 流式链路的空流兜底:上游 200 后首个有效产出前若遇到空流 EOF / 无数据超时,最多重试 N 次(0=关闭,默认 1)。重试会经 key_pool 自动切到下一个可用 key。",
"stream_empty_retry_max": 1,
"_comment_stream_first_byte_timeout_ms": "空流检测的首字节看门狗(毫秒,默认 30000,<=0 关闭)。上游 200 后一直没发任何 SSE 数据,超过该阈值视为空流并触发重试。",
"stream_first_byte_timeout_ms": 30000
}
26 changes: 26 additions & 0 deletions docs/CONFIGURATION.md
Original file line number Diff line number Diff line change
Expand Up @@ -277,6 +277,32 @@ opencode zen 上游的 base URL 列表。默认(未设置或为空数组)为

> 真实运行验证:`opencode2api launch claude --model mimo-v2.6-flash` 第二轮 `prompt_cached_tokens` 从 ~28.7k 提升到 ~32.4k(≈99.9% 的 prompt 命中),`big-pickle` 31.5k/31.6k;`codex --model mimo-v2.6-flash` `prompt_cached_tokens=9.92k`(≈98%),都已通过 `OPENCODE2API_CACHE_DEBUG=1` 中的 `cache_debug_usage` 观察。先前行为是只在 `buildUpstreamBody` 时注入顶层 `prompt_cache_retention`;现在 chat/claude/responses 直通(remembered)与 chat→responses 桥都统一补齐,并保持 IDEMPOTENT(上游已有字段时不覆盖)。

### `stream_empty_retry_max` / `stream_first_byte_timeout_ms`

claude→responses 流式链路的「空流兜底 + 首 token 前重试」。覆盖两类常见上游故障:

- **prefill 阶段被宰**:上游代理(CF / nginx)在首个 token 前杀 tunnel,网关只收到一个干净的 EOF——按旧实现客户端会看到 `stream ended without completion`,agent 中断。
- **挂死**:上游接受了连接但既不发数据也不关,客户端永久等待。

开启后(默认开启):上游 200 收到、但还没向客户端 WriteHeader 之前的窗口里,遇到 **空流 EOF / 超时未发数据 / 上游只发 `response.failed` 错误帧**,静默重发同一份请求最多 `stream_empty_retry_max` 次,客户端完全无感;重试经 `key_pool` 自动落到下一个可用 key,不会重复同一根 pipe。

```json
{
"stream_empty_retry_max": 1,
"stream_first_byte_timeout_ms": 30000
}
```

- `stream_empty_retry_max`:重试次数,默认 `1`,`0` 关闭。每个 attempt 都用同一份请求体重发(prompt_cache_key 稳定,input tokens 在缓存命中时接近零成本)。
- `stream_first_byte_timeout_ms`:peek 窗口毫秒数,默认 `30000`(30s)。覆盖大多数上游 prefill 时间;`<=0` 关闭看门狗,仅 EOF/error 触发。

**不重试的情况**(不改的承诺):
- 已向客户端写过任何字节后 EOF/杀流——按 ParalonCloud Rule 2 合成正常 stop 收尾(已有产出交付给 agent)。
- 仅有 thinking 没有 text 的 EOF——`reasoningFallback` 兜底把思考内容提升为 text,agent 拿到思考、不发 error。
- 非流式请求(`stream: false`)走的是另一条路径,与本机制无关。

> 真实运行验证:故意用反代在 5s 时杀 upstream tunnel,agent 端原本会 `stream ended without completion` 中断;开启本机制后第一次空流透明重试到下一个 key,整轮圆满完成。

## 管理面板

打开 `http://127.0.0.1:8000/` 可进入管理面板。面板可以修改配置、刷新模型和查看 token 统计。管理面板现已可设置 `prompt_cache_retention`、`cache_control_breakpoints`、`socks5_sticky`、`text_only_models`(「模型与路由」/「网络与代理」Tab),保存时这些字段随其余配置一并持久化到 `config.json`,不再被静默回擦。
Expand Down
3 changes: 1 addition & 2 deletions internal/app/anthropic_decode_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1838,8 +1838,7 @@ func TestResponsesStream_NormalizesUpstreamID(t *testing.T) {
``,
}, "\n")
rr := httptest.NewRecorder()
resp := &http.Response{StatusCode: 200, Body: io.NopCloser(strings.NewReader(upstream)), Header: make(http.Header)}
responsesStreamHandler(rr, nil, resp, "m", "m", false, nil, nil, ResponsesAPIRequest{})
responsesStreamHandler(rr, nil, strings.NewReader(upstream), "m", "m", false, nil, nil, ResponsesAPIRequest{}, nil, nil)
events := parseSSEEvents(t, rr.Body.String())
for _, e := range events {
if e.Name == "response.created" {
Expand Down
196 changes: 139 additions & 57 deletions internal/app/anthropic_upstream.go
Original file line number Diff line number Diff line change
Expand Up @@ -58,34 +58,62 @@ func forwardClaudeViaAnthropic(ctx context.Context, w http.ResponseWriter, auth
// 非 2xx 即使请求方要求 stream 也统一走 buffered JSON 错误直转
// (上游未建立 SSE 流,tee 会把错误 JSON 包进 data frame 破坏客户端解析)。
if status >= 200 && status < 300 && stream {
pipeAnthropicStream(ctx, w, rc, status, header, modelID)
// 用 DriveStreamWithRetry 在 peek 失败时切换 key 重试,空流 / EOF /
// 首字节超时被翻译为可重试的 errStreamIncompleteNoCommit。首轮复用
// 调用方已打开的 rc,后续重试通过 callOpenCodeAnthropicEndpoint 让
// key pool 切到下一把可用 key。
pending := rc
callOnce := func(c context.Context) (io.ReadCloser, int, error) {
if pending != nil {
r := pending
pending = nil
return r, status, nil
}
nrc, nstatus, _, nerr := callOpenCodeAnthropicEndpoint(c, upstreamBody, modelID, auth)
// Drive 拿到非 2xx 会立即返回(不再 close)。retry 路径里我们已
// 在上层只能是「走 chat 翻译兜底」,rc 在这里直接 close 掉防泄漏。
if nerr == nil && (nstatus < 200 || nstatus >= 300) {
if nrc != nil {
nrc.Close()
}
}
return nrc, nstatus, nerr
}
runOnce := func(c context.Context, w http.ResponseWriter, nrc io.Reader, _ []streamReadResult, _ *streamReader) (bool, error) {
return pipeAnthropicStream(c, w, nrc, status, header, modelID)
}
committed, driveErr := DriveStreamWithRetry(ctx, w, AnthropicProtocolHooks, callOnce, runOnce)
if !committed {
// 一直未 commit,让上层走 chat 翻译路径;不写任何字节给客户端。
if driveErr != nil {
log.Warn("anthropic passthrough stream exhausted retries", "model", modelID, "err", driveErr)
}
return false
}
return true
}
relayAnthropicBuffered(ctx, w, rc, status, header, modelID)
return true
}

// flushWriter 在每次 Write 后立即 Flush,保证 SSE 以事件粒度实时下发;
// http.ResponseWriter 内部带 bufio 缓冲,不显式 Flush 会把事件攒批到 EOF
// (与 responses_passthrough.go relayResponsesStream 的逐行 Flush 同一约定)。
type flushWriter struct {
w io.Writer
f http.Flusher
}

func (fw flushWriter) Write(p []byte) (int, error) {
n, err := fw.w.Write(p)
if n > 0 {
fw.f.Flush()
}
return n, err
}

// pipeAnthropicStream 把上游 Anthropic SSE 流字节级原样转发给客户端,同时
// 旁路 tee 解析 message_start / message_delta 中的 usage 记入 token 统计。
// 行边界、CRLF/LF、空行均不做改写,确保下游收到与上游完全一致的字节流。
// 仅在上游 2xx(真 SSE)时被调用;错误响应一律走 relayAnthropicBuffered。
func pipeAnthropicStream(ctx context.Context, w http.ResponseWriter, rc io.Reader, status int, header http.Header, modelID string) {
//
// 返回 (true, nil):已 commit(首帧已 peek + 写入)。EOF 时若未见过
// message_stop 但见过 message_start,合成一条 message_stop 保证客户端正常
// 关流;若两者皆无(不应发生:peek 至少要看到一帧)返回 false 供调用方
// 走未 commit 重试。返回 (false, err):peek 未 commit(空流 / EOF / 错误帧 /
// 首字节超时),由 DriveStreamWithRetry 决定是否换 key 重发。
func pipeAnthropicStream(ctx context.Context, w http.ResponseWriter, rc io.Reader, status int, header http.Header, modelID string) (bool, error) {
// peek 首帧:在 WriteHeader 之前约束 commit 边界,空流 / EOF / 错误帧 /
// 首字节超时都返回 errStreamIncompleteNoCommit,由调用方驱动重试。
peek := PeekFirstFrame(ctx, rc, time.Duration(config.StreamFirstByteTimeoutMs())*time.Millisecond, AnthropicProtocolHooks)
if peek.Err != nil {
return false, peek.Err
}

filtered := filterResponseHeaders(header)
for k, v := range filtered {
w.Header().Set(k, v[0])
Expand All @@ -97,61 +125,115 @@ func pipeAnthropicStream(ctx context.Context, w http.ResponseWriter, rc io.Reade

stats := &logging.StreamStats{Start: time.Now()}
fullUsage := map[string]any{}

// tee 管道:旁路解析走 pipeWriter,主流走 io.Copy 直透;两组无背压,
// io.Copy 返回(EOF、rc 读取失败、pw.Write 失败)时主动 pw.Close()
// 告知解析端收尾;ctx 取消则先 close 上游 rc 解锁 io.Copy,再
// pw.CloseWithError(ctx.Err()) 让 pr.Read 立刻返回。
pr, pw := io.Pipe()
copyDone := make(chan struct{})
go func() {
defer close(copyDone)
// MultiWriter 把每个 read 同步写给客户端与旁路解析端;flushWriter
// 让每片上游数据即时下发(不攒批)。任一侧写失败 io.Copy 立即返回,
// 随后 pw.Close 告知解析端收尾,最终 close(copyDone) 供主循环 join。
flusher, _ := w.(http.Flusher)
cw := io.Writer(w)
if flusher != nil {
cw = flushWriter{w: w, f: flusher}
}
_, _ = io.Copy(io.MultiWriter(cw, pw), rc)
_ = pw.Close()
}()

// tee 解析流:复用 newStreamReader 的协程,读到行就 observe,不写出。
reader := newStreamReader(ctx, pr, 0)
defer func() {
if len(fullUsage) > 0 {
statsx.RecordChatUsage(modelID, anthropicUsageToChat(fullUsage))
}
stats.Log(ctx, "claude")
}()

flusher, _ := w.(http.Flusher)
// 写 peek 出的原始字节(完整保留 \r\n / 换行 / 空行),同时喂给
// observeAnthropicStreamEvent 让 stats 与 message_start/stop 计数正确
// 累计——peek 消费过的帧不再二次进 reader.Read() 通道,所以这里必须补
// 一次观察。
if err := FlushPeekedBytes(w, peek.Consumed); err != nil {
return true, err
}
sawMessageStop := false
observeLine := func(line string) {
stats.NoteChunk()
observeAnthropicStreamEvent(stats, fullUsage, line)
payload, ok := strings.CutPrefix(line, "data: ")
if !ok {
return
}
var evt map[string]any
if json.Unmarshal([]byte(strings.TrimSpace(payload)), &evt) != nil {
return
}
if typ, _ := evt["type"].(string); typ == "message_stop" {
sawMessageStop = true
}
}
// 先用 peek 消费过的行回填 sawMessageStart / sawMessageStop——后续 EOF
// 兜底合成 message_stop 需要知道是否已见过 message_start / message_stop。
for _, res := range peek.Consumed {
if res.line != "" {
observeLine(res.line)
}
}
if flusher != nil {
flusher.Flush()
}

// 续用 peek 内部 streamReader(它的 bufio 已预读后续行),按 SSE 帧聚合
// 再写客户端——与原 io.Copy 的「一次上游 chunk ≈ 一次 Write+Flush」
// 节奏对齐,保留逐事件的打字机效果,而不是退回到 line-at-a-time。
reader := peek.Reader
if reader == nil {
// EOF 收尾的 peek 没留下 reader——主循环立即结束。
reader = newStreamReader(ctx, rc, 0)
}
defer reader.Close()

var frameBuf strings.Builder
flushFrame := func() error {
if frameBuf.Len() == 0 {
return nil
}
_, err := io.WriteString(w, frameBuf.String())
frameBuf.Reset()
if err != nil {
return err
}
if flusher != nil {
flusher.Flush()
}
return nil
}

for {
select {
case <-ctx.Done():
// 先 close 上游,让 io.Copy 立刻读到错误退出(不再卡在 w.Write),
// 再 close pipe 让旁路解析收尾,这样 copyDone 不会等慢客户端。
if c, ok := rc.(io.Closer); ok {
_ = c.Close()
}
_ = pw.CloseWithError(ctx.Err())
<-copyDone
return
return true, ctx.Err()
case result := <-reader.Read():
pendingErr := result.err
line := result.line
pendingErr := result.err
if line != "" {
stats.NoteChunk()
observeAnthropicStreamEvent(stats, fullUsage, line)
observeLine(line)
frameBuf.WriteString(line)
// 空行 = 帧边界:整帧一次写出再 Flush。
if strings.TrimRight(line, "\r\n") == "" {
if err := flushFrame(); err != nil {
return true, err
}
}
}
if pendingErr != nil {
// pr 的错误只可能来自 pw.Close(),即 copy 协程已越过 io.Copy,
// 此处 join 必然立即返回;保证协程不再于 handler 返回后触碰
// 已交还的 http.ResponseWriter(net/http 禁止这种并发使用)。
<-copyDone
return
// EOF / 上游读取失败。先把残帧(无空行收尾)吐出去,再看是否
// 需要补 message_stop / 走重试。
if err := flushFrame(); err != nil {
return true, err
}
// 关键不变量:此时 WriteHeader + peeked 首帧字节已经发出去了,
// 客户端连接已经处于 SSE 数据段。**绝不可再返回 (false, ...)**
// 否则 DriveStreamWithRetry 会用同一个 ResponseWriter 二次
// WriteHeader + 二次 replay peeked 字节,流被污染(I4/I7)。
if !sawMessageStop {
// 上游 EOF 但没关 message:展开成「合成 message_stop」让
// Claude SDK 正常关流。sawMessageStart=false 也照发——客
// 户端拿到「没 message_start 直接 message_stop」虽不规范,
// 但比重复写 header 安全(对一个非法流,SDK 通常仅丢弃该
// 事件,而不是报错)。
if _, err := io.WriteString(w, "event: message_stop\ndata: {\"type\":\"message_stop\"}\n\n"); err != nil {
return true, err
}
if flusher != nil {
flusher.Flush()
}
}
return true, nil
}
}
}
Expand Down
5 changes: 4 additions & 1 deletion internal/app/anthropic_upstream_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -102,7 +102,10 @@ func TestPipeAnthropicStream_PreservesHeaderAndBytes(t *testing.T) {
rec := httptest.NewRecorder()
header := http.Header{}
header.Set("Content-Type", "text/event-stream")
pipeAnthropicStream(context.Background(), rec, io.NopCloser(strings.NewReader(upstreamBody)), http.StatusOK, header, "m")
committed, perr := pipeAnthropicStream(context.Background(), rec, io.NopCloser(strings.NewReader(upstreamBody)), http.StatusOK, header, "m")
if !committed || perr != nil {
t.Fatalf("pipeAnthropicStream = (%v, %v), want (true, nil)", committed, perr)
}

if rec.Code != http.StatusOK {
t.Fatalf("status = %d, want 200", rec.Code)
Expand Down
Loading
Loading