Skip to content
Open
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
4 changes: 2 additions & 2 deletions internal/runtime/executor/codex_executor_execute.go
Original file line number Diff line number Diff line change
Expand Up @@ -120,7 +120,7 @@ func (e *CodexExecutor) Execute(ctx context.Context, auth *cliproxyauth.Auth, re
}
helps.AppendAPIResponseChunk(ctx, e.cfg, b)
helps.LogWithRequestID(ctx).Debugf("request error, error status: %d, error message: %s", httpResp.StatusCode, helps.SummarizeErrorBody(httpResp.Header.Get("Content-Type"), b))
err = newCodexStatusErr(httpResp.StatusCode, b)
err = newCodexStatusErrWithHeaders(httpResp.StatusCode, b, httpResp.Header)
return resp, err
}
data, errRead := io.ReadAll(httpResp.Body)
Expand Down Expand Up @@ -274,7 +274,7 @@ func (e *CodexExecutor) executeCompact(ctx context.Context, auth *cliproxyauth.A
b = applyCodexIdentityConfuseResponsePayload(b, identityState)
helps.AppendAPIResponseChunk(ctx, e.cfg, b)
helps.LogWithRequestID(ctx).Debugf("request error, error status: %d, error message: %s", httpResp.StatusCode, helps.SummarizeErrorBody(httpResp.Header.Get("Content-Type"), b))
err = newCodexStatusErr(httpResp.StatusCode, b)
err = newCodexStatusErrWithHeaders(httpResp.StatusCode, b, httpResp.Header)
return resp, err
}
data, err := io.ReadAll(httpResp.Body)
Expand Down
64 changes: 64 additions & 0 deletions internal/runtime/executor/codex_executor_retry_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,70 @@ func TestParseCodexRetryAfter(t *testing.T) {
})
}

func TestNewCodexStatusErrWithHeadersPreservesRetryMetadata(t *testing.T) {
t.Run("delta seconds", func(t *testing.T) {
headers := http.Header{
"Retry-After": {"7"},
"X-Request-Id": {"req-test"},
}
err := newCodexStatusErrWithHeaders(
http.StatusServiceUnavailable,
[]byte(`{"error":{"type":"server_error","message":"overloaded"}}`),
headers,
)

if retryAfter := err.RetryAfter(); retryAfter == nil || *retryAfter != 7*time.Second {
t.Fatalf("RetryAfter() = %v, want 7s", retryAfter)
}
if got := err.Headers().Get("X-Request-Id"); got != "req-test" {
t.Fatalf("Headers().Get(X-Request-Id) = %q, want req-test", got)
}

headers.Set("X-Request-Id", "mutated")
if got := err.Headers().Get("X-Request-Id"); got != "req-test" {
t.Fatalf("stored headers changed with source mutation: %q", got)
}
})

t.Run("http date", func(t *testing.T) {
retryAt := time.Now().Add(10 * time.Second).UTC().Truncate(time.Second)
err := newCodexStatusErrWithHeaders(
http.StatusServiceUnavailable,
[]byte(`{"error":{"type":"server_error"}}`),
http.Header{"Retry-After": {retryAt.Format(http.TimeFormat)}},
)

retryAfter := err.RetryAfter()
if retryAfter == nil || *retryAfter < 8*time.Second || *retryAfter > 10*time.Second {
t.Fatalf("RetryAfter() = %v, want approximately 10s", retryAfter)
}
})

t.Run("overflowing delta seconds", func(t *testing.T) {
err := newCodexStatusErrWithHeaders(
http.StatusServiceUnavailable,
[]byte(`{"error":{"type":"server_error"}}`),
http.Header{"Retry-After": {"9999999999"}},
)

if retryAfter := err.RetryAfter(); retryAfter != nil {
t.Fatalf("RetryAfter() = %v, want nil for overflowing delta seconds", *retryAfter)
}
})

t.Run("body retry metadata wins", func(t *testing.T) {
err := newCodexStatusErrWithHeaders(
http.StatusTooManyRequests,
[]byte(`{"error":{"type":"usage_limit_reached","resets_in_seconds":120}}`),
http.Header{"Retry-After": {"7"}},
)

if retryAfter := err.RetryAfter(); retryAfter == nil || *retryAfter != 120*time.Second {
t.Fatalf("RetryAfter() = %v, want body-derived 2m", retryAfter)
}
})
}

func TestNewCodexStatusErrTreatsCapacityAsRetryableRateLimit(t *testing.T) {
body := []byte(`{"error":{"message":"Selected model is at capacity. Please try a different model."}}`)

Expand Down
2 changes: 1 addition & 1 deletion internal/runtime/executor/codex_executor_stream.go
Original file line number Diff line number Diff line change
Expand Up @@ -123,7 +123,7 @@ func (e *CodexExecutor) ExecuteStream(ctx context.Context, auth *cliproxyauth.Au
}
helps.AppendAPIResponseChunk(ctx, e.cfg, data)
helps.LogWithRequestID(ctx).Debugf("request error, error status: %d, error message: %s", httpResp.StatusCode, helps.SummarizeErrorBody(httpResp.Header.Get("Content-Type"), data))
err = newCodexStatusErr(httpResp.StatusCode, data)
err = newCodexStatusErrWithHeaders(httpResp.StatusCode, data, httpResp.Header)
return nil, err
}
out := make(chan cliproxyexecutor.StreamChunk)
Expand Down
56 changes: 56 additions & 0 deletions internal/runtime/executor/codex_executor_stream_output_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,62 @@ import (
"github.com/tidwall/gjson"
)

func TestCodexExecutorHTTPErrorPreservesUpstreamHeaders(t *testing.T) {
for _, stream := range []bool{false, true} {
t.Run(map[bool]string{false: "execute", true: "execute_stream"}[stream], func(t *testing.T) {
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
w.Header().Set("Retry-After", "7")
w.Header().Set("X-Request-Id", "req-test")
w.WriteHeader(http.StatusServiceUnavailable)
_, _ = w.Write([]byte(`{"error":{"type":"server_error","message":"overloaded"}}`))
}))
defer server.Close()

executor := NewCodexExecutor(&config.Config{})
auth := &cliproxyauth.Auth{Attributes: map[string]string{
"base_url": server.URL,
"api_key": "test",
}}
req := cliproxyexecutor.Request{
Model: "gpt-test",
Payload: []byte(`{"model":"gpt-test","input":"hello"}`),
}
opts := cliproxyexecutor.Options{
SourceFormat: sdktranslator.FromString("openai-response"),
Stream: stream,
}

var err error
if stream {
_, err = executor.ExecuteStream(context.Background(), auth, req, opts)
} else {
_, err = executor.Execute(context.Background(), auth, req, opts)
}
if err == nil {
t.Fatal("expected upstream 503 error")
}
if got := statusCodeFromTestError(t, err); got != http.StatusServiceUnavailable {
t.Fatalf("status code = %d, want %d", got, http.StatusServiceUnavailable)
}
retryErr, ok := err.(interface{ RetryAfter() *time.Duration })
if !ok {
t.Fatalf("error %T does not expose RetryAfter()", err)
}
if retryAfter := retryErr.RetryAfter(); retryAfter == nil || *retryAfter != 7*time.Second {
t.Fatalf("RetryAfter() = %v, want 7s", retryAfter)
}
headerErr, ok := err.(interface{ Headers() http.Header })
if !ok {
t.Fatalf("error %T does not expose Headers()", err)
}
if got := headerErr.Headers().Get("X-Request-Id"); got != "req-test" {
t.Fatalf("X-Request-Id = %q, want req-test", got)
}
})
}
}

func TestCodexExecutorExecute_EmptyStreamCompletionOutputUsesOutputItemDone(t *testing.T) {
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "text/event-stream")
Expand Down
39 changes: 39 additions & 0 deletions internal/runtime/executor/codex_executor_terminal.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import (
"bytes"
"net/http"
"sort"
"strconv"
"strings"
"time"

Expand Down Expand Up @@ -261,6 +262,44 @@ func newCodexStatusErr(statusCode int, body []byte) statusErr {
return err
}

func newCodexStatusErrWithHeaders(statusCode int, body []byte, headers http.Header) statusErrWithHeaders {
err := newCodexStatusErr(statusCode, body)
if err.retryAfter == nil {
err.retryAfter = parseRetryAfterHeader(headers.Get("Retry-After"), time.Now())
}
return statusErrWithHeaders{
statusErr: err,
headers: headers.Clone(),
}
}

func parseRetryAfterHeader(value string, now time.Time) *time.Duration {
value = strings.TrimSpace(value)
if value == "" {
return nil
}
if seconds, err := strconv.ParseInt(value, 10, 64); err == nil {
if seconds < 0 {
return nil
}
const maxRetryAfterSeconds = int64(time.Duration(1<<63-1) / time.Second)
if seconds > maxRetryAfterSeconds {
return nil
}
delay := time.Duration(seconds) * time.Second
return &delay
}
retryAt, err := http.ParseTime(value)
if err != nil {
return nil
}
delay := retryAt.Sub(now)
if delay < 0 {
delay = 0
}
return &delay
}

func classifyCodexStatusError(statusCode int, body []byte) []byte {
code, errType, ok := codexStatusErrorClassification(statusCode, body)
if !ok {
Expand Down
8 changes: 4 additions & 4 deletions internal/runtime/executor/codex_openai_images.go
Original file line number Diff line number Diff line change
Expand Up @@ -140,7 +140,7 @@ func (e *CodexExecutor) executeOpenAIImage(ctx context.Context, auth *cliproxyau
helps.AppendAPIResponseChunk(ctx, e.cfg, data)
if httpResp.StatusCode < 200 || httpResp.StatusCode >= 300 {
helps.LogWithRequestID(ctx).Debugf("request error, error status: %d, error message: %s", httpResp.StatusCode, helps.SummarizeErrorBody(httpResp.Header.Get("Content-Type"), data))
err = newCodexStatusErr(httpResp.StatusCode, data)
err = newCodexStatusErrWithHeaders(httpResp.StatusCode, data, httpResp.Header)
return resp, err
}

Expand Down Expand Up @@ -234,7 +234,7 @@ func (e *CodexExecutor) executeOpenAIImageStream(ctx context.Context, auth *clip
data = applyCodexIdentityConfuseResponsePayload(data, identityState)
helps.AppendAPIResponseChunk(ctx, e.cfg, data)
helps.LogWithRequestID(ctx).Debugf("request error, error status: %d, error message: %s", httpResp.StatusCode, helps.SummarizeErrorBody(httpResp.Header.Get("Content-Type"), data))
err = newCodexStatusErr(httpResp.StatusCode, data)
err = newCodexStatusErrWithHeaders(httpResp.StatusCode, data, httpResp.Header)
return nil, err
}

Expand Down Expand Up @@ -367,7 +367,7 @@ func (e *CodexExecutor) executeDirectOpenAIImage(ctx context.Context, auth *clip
helps.AppendAPIResponseChunk(ctx, e.cfg, data)
if httpResp.StatusCode < 200 || httpResp.StatusCode >= 300 {
helps.LogWithRequestID(ctx).Debugf("request error, error status: %d, error message: %s", httpResp.StatusCode, helps.SummarizeErrorBody(httpResp.Header.Get("Content-Type"), data))
err = newCodexStatusErr(httpResp.StatusCode, data)
err = newCodexStatusErrWithHeaders(httpResp.StatusCode, data, httpResp.Header)
return resp, err
}

Expand Down Expand Up @@ -425,7 +425,7 @@ func (e *CodexExecutor) executeDirectOpenAIImageStream(ctx context.Context, auth
data = applyCodexIdentityConfuseResponsePayload(data, identityState)
helps.AppendAPIResponseChunk(ctx, e.cfg, data)
helps.LogWithRequestID(ctx).Debugf("request error, error status: %d, error message: %s", httpResp.StatusCode, helps.SummarizeErrorBody(httpResp.Header.Get("Content-Type"), data))
err = newCodexStatusErr(httpResp.StatusCode, data)
err = newCodexStatusErrWithHeaders(httpResp.StatusCode, data, httpResp.Header)
return nil, err
}

Expand Down
4 changes: 1 addition & 3 deletions sdk/api/handlers/handlers_execution.go
Original file line number Diff line number Diff line change
Expand Up @@ -310,9 +310,7 @@ func executionErrorMessage(err error) *interfaces.ErrorMessage {
}
var addon http.Header
if he, ok := err.(interface{ Headers() http.Header }); ok && he != nil {
if hdr := he.Headers(); hdr != nil {
addon = hdr.Clone()
}
addon = FilterUpstreamHeaders(he.Headers())
}
return &interfaces.ErrorMessage{StatusCode: status, Error: err, Addon: addon}
}
Expand Down
91 changes: 91 additions & 0 deletions sdk/api/handlers/handlers_execution_error_headers_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,91 @@
package handlers

import (
"net/http"
"net/http/httptest"
"testing"

"github.com/gin-gonic/gin"
sdkconfig "github.com/router-for-me/CLIProxyAPI/v7/sdk/config"
)

type executionHeaderError struct {
headers http.Header
}

func (e executionHeaderError) Error() string { return "upstream unavailable" }
func (e executionHeaderError) StatusCode() int { return http.StatusServiceUnavailable }
func (e executionHeaderError) Headers() http.Header {
return e.headers.Clone()
}

func newExecutionHeaderError() executionHeaderError {
return executionHeaderError{headers: http.Header{
"Retry-After": {"120"},
"X-Request-Id": {"request-123"},
"Content-Length": {"999"},
"Set-Cookie": {"session=secret"},
"Connection": {"X-Connection-Only"},
"X-Connection-Only": {"secret"},
"X-Litellm-Trace": {"gateway"},
}}
}

func TestExecutionErrorMessageFiltersUpstreamHeaders(t *testing.T) {
err := newExecutionHeaderError()

msg := executionErrorMessage(err)
if msg == nil || msg.Error == nil || msg.Error.Error() != err.Error() {
t.Fatalf("error message = %#v, want original error", msg)
}
if msg.StatusCode != http.StatusServiceUnavailable {
t.Fatalf("status = %d, want %d", msg.StatusCode, http.StatusServiceUnavailable)
}
if got := msg.Addon.Get("Retry-After"); got != "120" {
t.Fatalf("Retry-After = %q, want 120", got)
}
if got := msg.Addon.Get("X-Request-Id"); got != "request-123" {
t.Fatalf("X-Request-Id = %q, want request-123", got)
}
for _, key := range []string{"Content-Length", "Set-Cookie", "Connection", "X-Connection-Only", "X-Litellm-Trace"} {
if values := msg.Addon.Values(key); len(values) != 0 {
t.Fatalf("%s leaked through filtered error headers: %v", key, values)
}
}
}

func TestExecutionErrorHeadersRespectPassthroughToggle(t *testing.T) {
for _, enabled := range []bool{false, true} {
t.Run(map[bool]string{false: "disabled", true: "enabled"}[enabled], func(t *testing.T) {
gin.SetMode(gin.TestMode)
recorder := httptest.NewRecorder()
ctx, _ := gin.CreateTestContext(recorder)
ctx.Request = httptest.NewRequest(http.MethodPost, "/v1/messages", nil)

handler := NewBaseAPIHandlers(&sdkconfig.SDKConfig{PassthroughHeaders: enabled}, nil)
handler.WriteErrorResponse(ctx, executionErrorMessage(newExecutionHeaderError()))

if recorder.Code != http.StatusServiceUnavailable {
t.Fatalf("status = %d, want %d", recorder.Code, http.StatusServiceUnavailable)
}
wantRetryAfter := ""
wantRequestID := ""
if enabled {
wantRetryAfter = "120"
wantRequestID = "request-123"
}
if got := recorder.Header().Get("Retry-After"); got != wantRetryAfter {
t.Fatalf("Retry-After = %q, want %q", got, wantRetryAfter)
}
if got := recorder.Header().Get("X-Request-Id"); got != wantRequestID {
t.Fatalf("X-Request-Id = %q, want %q", got, wantRequestID)
}
if got := recorder.Header().Get("Set-Cookie"); got != "" {
t.Fatalf("Set-Cookie leaked: %q", got)
}
if got := recorder.Header().Get("Content-Length"); got == "999" {
t.Fatalf("upstream Content-Length leaked: %q", got)
}
})
}
}
2 changes: 1 addition & 1 deletion sdk/api/handlers/openai/openai_responses_websocket.go
Original file line number Diff line number Diff line change
Expand Up @@ -436,7 +436,7 @@ func (h *OpenAIResponsesAPIHandler) ResponsesWebsocket(c *gin.Context) {
if errMsg != nil {
h.LoggingAPIResponseError(context.WithValue(context.Background(), "gin", c), errMsg)
markAPIResponseTimestamp(c)
errorPayload, errWrite := writeResponsesWebsocketError(writer, wsTimelineLog, errMsg)
errorPayload, errWrite := writeResponsesWebsocketError(writer, wsTimelineLog, errMsg, handlers.PassthroughHeadersEnabled(h.Cfg))
log.Infof(
"responses websocket: downstream_out id=%s type=%d event=%s payload=%s",
passthroughSessionID,
Expand Down
Loading
Loading