Skip to content

Commit a2dc057

Browse files
committed
fix(api): preserve typed stream read errors
Devin can finish the HTTP exchange and then deliver a typed Connect error through the stream body. Reclassifying it as an interruption turns malformed requests into account cooldowns. Preserve wrapped provider errors at each SSE reader and retain their classification, failover decision, and retry hint. Continue classifying unknown transport errors as unavailable so account routing can retry them.
1 parent 6664610 commit a2dc057

5 files changed

Lines changed: 96 additions & 2 deletions

File tree

CHANGELOG.md

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,8 +7,12 @@ Write each change in both `### English` and `### 中文` under `## Unreleased`.
77

88
### English
99

10+
- Preserve typed upstream stream errors through the OpenAI, Anthropic, and Responses relays so invalid Devin requests do not falsely cool accounts, while transport interruptions remain retryable
11+
1012
### 中文
1113

14+
- OpenAI、Anthropic 与 Responses 流式转发会保留上游的类型化错误,避免无效的 Devin 请求被错误地冷却账号,同时传输中断仍可重试
15+
1216
## 0.5.6 - 2026-09-17
1317

1418
### English

internal/api/chat.go

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -153,7 +153,7 @@ func relayOpenAIStream(w http.ResponseWriter, body io.Reader) (stats streamRelay
153153
if streamErr != nil {
154154
return stats, streamErr
155155
}
156-
streamErr := newStreamProviderError("upstream_stream_interrupted", "stream read error: "+err.Error(), http.StatusBadGateway)
156+
streamErr := streamReadProviderError(err)
157157
if writeErr := writeStructuredStreamError(writer, streamErr); writeErr != nil {
158158
return stats, writeErr
159159
}
@@ -376,6 +376,14 @@ func newStreamProviderError(code, message string, status int) *providers.Error {
376376
return providerErrorFromClassified(classified)
377377
}
378378

379+
func streamReadProviderError(err error) *providers.Error {
380+
var providerErr *providers.Error
381+
if errors.As(err, &providerErr) && providerErr != nil {
382+
return providerErr
383+
}
384+
return newStreamProviderError("upstream_stream_interrupted", "stream read error: "+err.Error(), http.StatusBadGateway)
385+
}
386+
379387
func providerErrorFromClassified(classified accounts.Classified) *providers.Error {
380388
failover := classified.Failover
381389
retryAfter := classified.RetryAfter

internal/api/chat_usage_test.go

Lines changed: 57 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,8 @@ package api
22

33
import (
44
"errors"
5+
"fmt"
6+
"io"
57
"net/http"
68
"net/http/httptest"
79
"strings"
@@ -17,6 +19,12 @@ import (
1719

1820
func intPtr(value int) *int { return &value }
1921

22+
func closedStreamPipe(err error) *io.PipeReader {
23+
reader, writer := io.Pipe()
24+
_ = writer.CloseWithError(err)
25+
return reader
26+
}
27+
2028
func TestWriteClassifiedErrKeepsTraeQuotaKind(t *testing.T) {
2129
recorder := httptest.NewRecorder()
2230
failover := true
@@ -247,6 +255,55 @@ func TestRelayOpenAIStreamReportsIncompleteStreamStructurally(t *testing.T) {
247255
}
248256
}
249257

258+
func TestRelayOpenAIStreamPreservesTypedReadError(t *testing.T) {
259+
recorder := httptest.NewRecorder()
260+
failover := false
261+
want := &providers.Error{
262+
Kind: accounts.KindInvalidRequest, Status: http.StatusBadRequest,
263+
Code: "invalid_argument", Type: "invalid_request_error", Message: "upstream rejected request",
264+
RetryAfter: 45 * time.Second, Failover: &failover,
265+
}
266+
_, err := relayOpenAIStream(recorder, closedStreamPipe(fmt.Errorf("Connect trailer: %w", want)))
267+
var got *providers.Error
268+
if !errors.As(err, &got) || got != want {
269+
t.Fatalf("error=%T %+v want pointer=%p", err, err, want)
270+
}
271+
if got.Kind != accounts.KindInvalidRequest || got.Status != http.StatusBadRequest || got.Code != "invalid_argument" ||
272+
got.Type != "invalid_request_error" || got.RetryAfter != 45*time.Second || got.Failover == nil || *got.Failover {
273+
t.Fatalf("provider error=%+v", got)
274+
}
275+
output := recorder.Body.String()
276+
if !strings.Contains(output, `"code":"invalid_argument"`) || !strings.Contains(output, `"retry_after":45`) || strings.Contains(output, "upstream_stream_interrupted") {
277+
t.Fatalf("structured error=%s", output)
278+
}
279+
pool := accounts.NewPool(nil, nil)
280+
pool.Upsert(accounts.Item{ID: "devin-account"})
281+
executor.NewChatExecutor(pool, "").ObserveStreamFailure("devin-account", got, "swe-2")
282+
item, _ := pool.ByID("devin-account")
283+
if item.LastKind != "" || !item.DownUntil.IsZero() {
284+
t.Fatalf("invalid request cooled account: kind=%q down=%v", item.LastKind, item.DownUntil)
285+
}
286+
}
287+
288+
func TestRelayOpenAIStreamWrapsUnknownReadError(t *testing.T) {
289+
recorder := httptest.NewRecorder()
290+
_, err := relayOpenAIStream(recorder, closedStreamPipe(errors.New("socket closed")))
291+
var got *providers.Error
292+
if !errors.As(err, &got) || got.Kind != accounts.KindUnavailable || got.Code != "upstream_stream_interrupted" || got.Status != http.StatusBadGateway {
293+
t.Fatalf("error=%T %+v", err, err)
294+
}
295+
if !strings.Contains(got.Message, "stream read error: socket closed") {
296+
t.Fatalf("message=%q", got.Message)
297+
}
298+
pool := accounts.NewPool(nil, nil)
299+
pool.Upsert(accounts.Item{ID: "devin-account"})
300+
executor.NewChatExecutor(pool, "").ObserveStreamFailure("devin-account", got, "swe-2")
301+
item, _ := pool.ByID("devin-account")
302+
if item.LastKind != accounts.KindUnavailable || item.DownUntil.IsZero() {
303+
t.Fatalf("transport interruption was not unavailable: kind=%q down=%v", item.LastKind, item.DownUntil)
304+
}
305+
}
306+
250307
func TestParseStreamUsageLineReadsWorkBuddyCredit(t *testing.T) {
251308
stats, ok := parseStreamUsageLine(`data: {"model":"hy3","usage":{"prompt_tokens":16,"completion_tokens":2,"credit":0.75}}`)
252309
if !ok || stats.Credits == nil || *stats.Credits != 0.75 {

internal/api/compat.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -524,7 +524,7 @@ func consumeOpenAIStream(body io.Reader, handle func(json.RawMessage, *streamedC
524524
frame = append(frame, line)
525525
}
526526
if err := scanner.Err(); err != nil {
527-
return stats, output, newStreamProviderError("upstream_stream_interrupted", "stream read error: "+err.Error(), http.StatusBadGateway)
527+
return stats, output, streamReadProviderError(err)
528528
}
529529
if err := flush(); err != nil {
530530
return stats, output, err

internal/api/compat_test.go

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,19 +2,44 @@ package api
22

33
import (
44
"encoding/json"
5+
"errors"
6+
"fmt"
57
"io"
68
"net/http"
79
"net/http/httptest"
810
"regexp"
911
"strings"
1012
"testing"
13+
"time"
1114

1215
"github.com/caigee-cmd/cli2api/internal/accounts"
1316
"github.com/caigee-cmd/cli2api/internal/auth"
1417
"github.com/caigee-cmd/cli2api/internal/executor"
18+
"github.com/caigee-cmd/cli2api/internal/providers"
1519
"github.com/caigee-cmd/cli2api/internal/translate"
1620
)
1721

22+
func TestCompatibilityStreamsPreserveTypedReadError(t *testing.T) {
23+
failover := false
24+
want := &providers.Error{
25+
Kind: accounts.KindInvalidRequest, Status: http.StatusBadRequest,
26+
Code: "invalid_argument", Type: "invalid_request_error", Message: "upstream rejected request",
27+
RetryAfter: 45 * time.Second, Failover: &failover,
28+
}
29+
for name, relay := range map[string]func(io.Writer, io.Reader, string, string) (streamRelayStats, error){
30+
"anthropic": relayAnthropicStream,
31+
"responses": relayResponsesStream,
32+
} {
33+
t.Run(name, func(t *testing.T) {
34+
_, err := relay(httptest.NewRecorder(), closedStreamPipe(fmt.Errorf("Connect trailer: %w", want)), "req_1", "devin/swe-2")
35+
var got *providers.Error
36+
if !errors.As(err, &got) || got != want {
37+
t.Fatalf("error=%T %+v want pointer=%p", err, err, want)
38+
}
39+
})
40+
}
41+
}
42+
1843
func newCompatibilityServer(t *testing.T, worker http.HandlerFunc) (*Server, func()) {
1944
t.Helper()
2045
upstream := httptest.NewServer(worker)

0 commit comments

Comments
 (0)