From 3db00d3fee37573bdf0bbfb1b5730118d7aa3cc8 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=8E=8B=E9=B9=8F?= <2829624376@qq.com> Date: Thu, 16 Jul 2026 00:19:20 +0800 Subject: [PATCH 1/2] fix(openai): fail over account-specific body limits --- .../openai_body_limit_failover_test.go | 58 ++++++++ .../handler/openai_gateway_handler.go | 11 ++ .../service/openai_gateway_cc_pipeline.go | 13 +- .../service/openai_gateway_forward.go | 12 +- .../service/openai_gateway_passthrough.go | 16 ++- .../service/openai_gateway_upstream_errors.go | 64 +++++++++ ...openai_request_body_limit_failover_test.go | 135 ++++++++++++++++++ 7 files changed, 292 insertions(+), 17 deletions(-) create mode 100644 backend/internal/handler/openai_body_limit_failover_test.go create mode 100644 backend/internal/service/openai_request_body_limit_failover_test.go diff --git a/backend/internal/handler/openai_body_limit_failover_test.go b/backend/internal/handler/openai_body_limit_failover_test.go new file mode 100644 index 0000000000..ff8970d147 --- /dev/null +++ b/backend/internal/handler/openai_body_limit_failover_test.go @@ -0,0 +1,58 @@ +package handler + +import ( + "bytes" + "encoding/json" + "net/http" + "net/http/httptest" + "strings" + "testing" + + "github.com/Wei-Shaw/sub2api/internal/service" + "github.com/gin-gonic/gin" + "github.com/stretchr/testify/require" +) + +func TestOpenAIBodyLimitFailoverExhausted_ReturnsRedactedJSON413(t *testing.T) { + gin.SetMode(gin.TestMode) + rec := httptest.NewRecorder() + c, _ := gin.CreateTestContext(rec) + c.Request = httptest.NewRequest(http.MethodPost, "/v1/responses", bytes.NewReader(nil)) + + (&OpenAIGatewayHandler{}).handleFailoverExhausted(c, bodyLimitFailoverTestError(), false) + + require.Equal(t, http.StatusRequestEntityTooLarge, rec.Code) + var envelope map[string]any + require.NoError(t, json.Unmarshal(rec.Body.Bytes(), &envelope)) + errBody := envelope["error"].(map[string]any) + require.Equal(t, "invalid_request_error", errBody["type"]) + require.Equal(t, "Request payload is too large", errBody["message"]) + require.NotContains(t, rec.Body.String(), "must-not-leak") +} + +func TestOpenAIBodyLimitFailoverExhausted_ReturnsRedactedResponsesSSE(t *testing.T) { + gin.SetMode(gin.TestMode) + rec := httptest.NewRecorder() + c, _ := gin.CreateTestContext(rec) + c.Request = httptest.NewRequest(http.MethodPost, "/v1/responses", bytes.NewReader(nil)) + + (&OpenAIGatewayHandler{}).handleFailoverExhausted(c, bodyLimitFailoverTestError(), true) + + body := rec.Body.String() + require.True(t, strings.HasPrefix(body, "event: response.failed\n")) + require.Contains(t, body, `"code":"invalid_request"`) + require.Contains(t, body, `"message":"Request payload is too large"`) + require.NotContains(t, body, "must-not-leak") +} + +func bodyLimitFailoverTestError() *service.UpstreamFailoverError { + return &service.UpstreamFailoverError{ + StatusCode: http.StatusRequestEntityTooLarge, + ResponseBody: []byte(`{"error":{"message":"proxy limit secret=must-not-leak"}}`), + Scope: service.GatewayFailureScopeAccount, + Reason: service.GatewayFailureReason("openai_request_body_too_large"), + NextAccountAction: service.NextAccountRetry, + ClientStatusCode: http.StatusRequestEntityTooLarge, + ClientMessage: "Request payload is too large", + } +} diff --git a/backend/internal/handler/openai_gateway_handler.go b/backend/internal/handler/openai_gateway_handler.go index 83c4b4ac23..2b56c12504 100644 --- a/backend/internal/handler/openai_gateway_handler.go +++ b/backend/internal/handler/openai_gateway_handler.go @@ -2077,6 +2077,17 @@ func (h *OpenAIGatewayHandler) handleFailoverExhausted(c *gin.Context, failoverE h.handleFailoverExhaustedSimple(c, http.StatusBadGateway, streamStarted) return } + if failoverErr.IsOpenAIRequestBodyTooLarge() { + service.SetOpsUpstreamError(c, http.StatusRequestEntityTooLarge, service.OpenAIRequestBodyTooLargeClientMessage, "") + h.handleStreamingAwareError( + c, + http.StatusRequestEntityTooLarge, + "invalid_request_error", + service.OpenAIRequestBodyTooLargeClientMessage, + streamStarted, + ) + return + } copyFailoverRetryAfter(c, failoverErr.ResponseHeaders) if failoverErr.IsCredentialFailure() { status, message := credentialFailoverClientResponse(failoverErr) diff --git a/backend/internal/service/openai_gateway_cc_pipeline.go b/backend/internal/service/openai_gateway_cc_pipeline.go index b38de924a3..52658b2729 100644 --- a/backend/internal/service/openai_gateway_cc_pipeline.go +++ b/backend/internal/service/openai_gateway_cc_pipeline.go @@ -115,12 +115,13 @@ func (s *OpenAIGatewayService) failoverOpenAIUpstreamHTTPError( if account.Platform != PlatformGrok { s.handleOpenAIAccountUpstreamError(ctx, account, resp.StatusCode, resp.Header, respBody, upstreamModel) } - return &UpstreamFailoverError{ - StatusCode: resp.StatusCode, - ResponseBody: respBody, - ResponseHeaders: resp.Header.Clone(), - RetryableOnSameAccount: account.IsPoolMode() && (account.IsPoolModeRetryableStatus(resp.StatusCode) || isOpenAITransientProcessingError(resp.StatusCode, upstreamMsg, respBody)), - } + return newOpenAIUpstreamFailoverError( + resp.StatusCode, + resp.Header, + respBody, + upstreamMsg, + account.IsPoolMode() && (account.IsPoolModeRetryableStatus(resp.StatusCode) || isOpenAITransientProcessingError(resp.StatusCode, upstreamMsg, respBody)), + ) } // openAIChatCompletionsTargetURL 解析账号的(非 Grok)Chat Completions 上游端点。 diff --git a/backend/internal/service/openai_gateway_forward.go b/backend/internal/service/openai_gateway_forward.go index 5259aa843c..6aa958799a 100644 --- a/backend/internal/service/openai_gateway_forward.go +++ b/backend/internal/service/openai_gateway_forward.go @@ -856,11 +856,13 @@ func (s *OpenAIGatewayService) Forward(ctx context.Context, c *gin.Context, acco }) s.handleFailoverSideEffects(ctx, resp, account, respBody, upstreamModel) - return nil, &UpstreamFailoverError{ - StatusCode: resp.StatusCode, - ResponseBody: respBody, - RetryableOnSameAccount: account.IsPoolMode() && (account.IsPoolModeRetryableStatus(resp.StatusCode) || isOpenAITransientProcessingError(resp.StatusCode, upstreamMsg, respBody)), - } + return nil, newOpenAIUpstreamFailoverError( + resp.StatusCode, + resp.Header, + respBody, + upstreamMsg, + account.IsPoolMode() && (account.IsPoolModeRetryableStatus(resp.StatusCode) || isOpenAITransientProcessingError(resp.StatusCode, upstreamMsg, respBody)), + ) } return s.handleErrorResponse(ctx, resp, c, account, body, billingModel) } diff --git a/backend/internal/service/openai_gateway_passthrough.go b/backend/internal/service/openai_gateway_passthrough.go index bbf80c97c9..445433167b 100644 --- a/backend/internal/service/openai_gateway_passthrough.go +++ b/backend/internal/service/openai_gateway_passthrough.go @@ -461,6 +461,9 @@ func shouldFailoverOpenAIPassthroughResponse(account *Account, statusCode int, r if isOpenAIContextWindowError("", responseBody) { return false } + if isOpenAIRequestBodyTooLargeError(statusCode, "", responseBody) { + return true + } switch statusCode { case http.StatusTooManyRequests, 529: return true @@ -589,12 +592,13 @@ func (s *OpenAIGatewayService) handleFailoverErrorResponsePassthrough( Detail: upstreamDetail, UpstreamResponseBody: upstreamDetail, }) - return &UpstreamFailoverError{ - StatusCode: resp.StatusCode, - ResponseBody: body, - ResponseHeaders: resp.Header.Clone(), - RetryableOnSameAccount: account.IsPoolMode() && account.IsPoolModeRetryableStatus(resp.StatusCode), - } + return newOpenAIUpstreamFailoverError( + resp.StatusCode, + resp.Header, + body, + upstreamMsg, + account.IsPoolMode() && account.IsPoolModeRetryableStatus(resp.StatusCode), + ) } func (s *OpenAIGatewayService) handleErrorResponsePassthrough( diff --git a/backend/internal/service/openai_gateway_upstream_errors.go b/backend/internal/service/openai_gateway_upstream_errors.go index d8823fe4ed..b01618c1f0 100644 --- a/backend/internal/service/openai_gateway_upstream_errors.go +++ b/backend/internal/service/openai_gateway_upstream_errors.go @@ -222,12 +222,55 @@ func (s *OpenAIGatewayService) shouldFailoverOpenAIUpstreamResponse(statusCode i if isOpenAIContextWindowError(upstreamMsg, upstreamBody) { return false } + if isOpenAIRequestBodyTooLargeError(statusCode, upstreamMsg, upstreamBody) { + return true + } if s.shouldFailoverUpstreamError(statusCode) { return true } return isOpenAITransientProcessingError(statusCode, upstreamMsg, upstreamBody) } +// OpenAIRequestBodyTooLargeClientMessage is the fixed downstream message used +// after all account-specific request body limit failovers are exhausted. +const OpenAIRequestBodyTooLargeClientMessage = "Request payload is too large" + +const openAIRequestBodyTooLargeReason = GatewayFailureReason("openai_request_body_too_large") + +func isOpenAIRequestBodyTooLargeError(statusCode int, upstreamMsg string, upstreamBody []byte) bool { + return statusCode == http.StatusRequestEntityTooLarge && !isOpenAIContextWindowError(upstreamMsg, upstreamBody) +} + +func newOpenAIUpstreamFailoverError( + statusCode int, + responseHeaders http.Header, + responseBody []byte, + upstreamMsg string, + retryableOnSameAccount bool, +) *UpstreamFailoverError { + failoverErr := &UpstreamFailoverError{ + StatusCode: statusCode, + ResponseBody: responseBody, + ResponseHeaders: responseHeaders.Clone(), + RetryableOnSameAccount: retryableOnSameAccount, + } + if isOpenAIRequestBodyTooLargeError(statusCode, upstreamMsg, responseBody) { + failoverErr.RetryableOnSameAccount = false + failoverErr.Scope = GatewayFailureScopeAccount + failoverErr.Reason = openAIRequestBodyTooLargeReason + failoverErr.NextAccountAction = NextAccountRetry + failoverErr.ClientStatusCode = http.StatusRequestEntityTooLarge + failoverErr.ClientMessage = OpenAIRequestBodyTooLargeClientMessage + } + return failoverErr +} + +// IsOpenAIRequestBodyTooLarge reports whether another account may accept the +// same request even though the selected account rejected its serialized size. +func (e *UpstreamFailoverError) IsOpenAIRequestBodyTooLarge() bool { + return e != nil && e.Reason == openAIRequestBodyTooLargeReason +} + func marshalOpenAIUpstreamJSON(v any) ([]byte, error) { var buf bytes.Buffer enc := json.NewEncoder(&buf) @@ -328,6 +371,27 @@ func (s *OpenAIGatewayService) handleErrorResponse( ) } + if isOpenAIRequestBodyTooLargeError(resp.StatusCode, upstreamMsg, body) { + appendOpsUpstreamError(c, OpsUpstreamErrorEvent{ + Platform: account.Platform, + AccountID: account.ID, + AccountName: account.Name, + UpstreamStatusCode: resp.StatusCode, + UpstreamRequestID: resp.Header.Get("x-request-id"), + Kind: "failover", + Message: upstreamMsg, + Detail: upstreamDetail, + }) + s.handleOpenAIAccountUpstreamError(ctx, account, resp.StatusCode, resp.Header, body, requestedModel...) + return nil, newOpenAIUpstreamFailoverError( + resp.StatusCode, + resp.Header, + body, + upstreamMsg, + false, + ) + } + if status, errType, errMsg, matched := applyErrorPassthroughRule( c, PlatformOpenAI, diff --git a/backend/internal/service/openai_request_body_limit_failover_test.go b/backend/internal/service/openai_request_body_limit_failover_test.go new file mode 100644 index 0000000000..6d6b5cb63e --- /dev/null +++ b/backend/internal/service/openai_request_body_limit_failover_test.go @@ -0,0 +1,135 @@ +package service + +import ( + "bytes" + "context" + "errors" + "fmt" + "net/http" + "net/http/httptest" + "strings" + "testing" + + "github.com/Wei-Shaw/sub2api/internal/config" + "github.com/gin-gonic/gin" + "github.com/stretchr/testify/require" + "github.com/tidwall/gjson" +) + +func TestOpenAIRequestBodyLimitFailover_HTTP413SwitchesAccountsBeforeWrite(t *testing.T) { + gin.SetMode(gin.TestMode) + requestBody := []byte(`{"model":"gpt-5.2","stream":false,"input":"hello"}`) + + for _, passthrough := range []bool{false, true} { + name := "native_responses" + if passthrough { + name = "api_key_passthrough" + } + t.Run(name, func(t *testing.T) { + rec := httptest.NewRecorder() + c, _ := gin.CreateTestContext(rec) + c.Request = httptest.NewRequest(http.MethodPost, "/v1/responses", bytes.NewReader(nil)) + + const upstreamBody = `{"error":{"message":"request body exceeds this account's 16MB proxy limit; secret=must-not-leak","type":"invalid_request_error"}}` + body := &passthroughCloseTrackingReadCloser{Reader: strings.NewReader(upstreamBody)} + upstream := &httpUpstreamRecorder{resp: &http.Response{ + StatusCode: http.StatusRequestEntityTooLarge, + Header: http.Header{ + "Content-Type": []string{"application/json"}, + "X-Request-Id": []string{"rid-body-limit"}, + }, + Body: body, + }} + svc := &OpenAIGatewayService{ + cfg: &config.Config{Gateway: config.GatewayConfig{ForceCodexCLI: false}}, + httpUpstream: upstream, + } + account := &Account{ + ID: 161, + Name: name, + Platform: PlatformOpenAI, + Type: AccountTypeAPIKey, + Concurrency: 1, + Credentials: map[string]any{ + "api_key": "sk-test", + "base_url": "https://api.example.test", + "pool_mode": true, + "pool_mode_retry_status_codes": []any{ + float64(http.StatusRequestEntityTooLarge), + }, + }, + Extra: map[string]any{ + "openai_passthrough": passthrough, + "openai_responses_supported": true, + }, + Status: StatusActive, + Schedulable: true, + } + + result, err := svc.Forward(context.Background(), c, account, requestBody) + + require.Nil(t, result) + var failoverErr *UpstreamFailoverError + require.ErrorAs(t, err, &failoverErr) + require.Equal(t, http.StatusRequestEntityTooLarge, failoverErr.StatusCode) + require.Equal(t, GatewayFailureScopeAccount, failoverErr.Scope) + require.Equal(t, GatewayFailureReason("openai_request_body_too_large"), failoverErr.Reason) + require.Equal(t, NextAccountRetry, failoverErr.NextAccountAction) + require.Equal(t, http.StatusRequestEntityTooLarge, failoverErr.ClientStatusCode) + require.Equal(t, "Request payload is too large", failoverErr.ClientMessage) + require.False(t, failoverErr.RetryableOnSameAccount, "a body limit requires another account, not another attempt on the same account") + require.False(t, c.Writer.Written(), "account failover must happen before downstream output is committed") + require.Empty(t, rec.Body.String()) + require.True(t, body.closed) + if passthrough { + require.Equal(t, requestBody, upstream.lastBody) + } else { + require.Equal(t, "gpt-5.2", gjson.GetBytes(upstream.lastBody, "model").String()) + require.Equal(t, "hello", gjson.GetBytes(upstream.lastBody, "input").String()) + } + }) + } +} + +func TestOpenAIRequestBodyLimitFailover_ContextWindow413DoesNotSwitchAccounts(t *testing.T) { + gin.SetMode(gin.TestMode) + requestBody := []byte(`{"model":"gpt-5.2","stream":false,"input":"hello"}`) + + for _, passthrough := range []bool{false, true} { + t.Run(fmt.Sprintf("passthrough_%t", passthrough), func(t *testing.T) { + rec := httptest.NewRecorder() + c, _ := gin.CreateTestContext(rec) + c.Request = httptest.NewRequest(http.MethodPost, "/v1/responses", bytes.NewReader(nil)) + + const upstreamBody = `{"error":{"message":"Your input exceeds the context window of this model. Please adjust your input and try again.","type":"invalid_request_error"}}` + body := &passthroughCloseTrackingReadCloser{Reader: strings.NewReader(upstreamBody)} + svc := &OpenAIGatewayService{ + cfg: &config.Config{Gateway: config.GatewayConfig{ForceCodexCLI: false}}, + httpUpstream: &httpUpstreamRecorder{resp: &http.Response{ + StatusCode: http.StatusRequestEntityTooLarge, + Header: http.Header{"Content-Type": []string{"application/json"}}, + Body: body, + }}, + } + account := &Account{ + ID: 162, Platform: PlatformOpenAI, Type: AccountTypeAPIKey, Concurrency: 1, + Credentials: map[string]any{"api_key": "sk-test", "base_url": "https://api.example.test"}, + Extra: map[string]any{ + "openai_passthrough": passthrough, + "openai_responses_supported": true, + }, + Status: StatusActive, Schedulable: true, + } + + result, err := svc.Forward(context.Background(), c, account, requestBody) + + require.Nil(t, result) + require.Error(t, err) + var failoverErr *UpstreamFailoverError + require.False(t, errors.As(err, &failoverErr), "context-window failures are deterministic request errors") + require.True(t, c.Writer.Written()) + require.Contains(t, rec.Body.String(), "exceeds the context window") + require.True(t, body.closed) + }) + } +} From ad5e2a85b7a4a516ca3bfdb62f1b53d067f368b1 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=8E=8B=E9=B9=8F?= <2829624376@qq.com> Date: Thu, 16 Jul 2026 00:52:26 +0800 Subject: [PATCH 2/2] test(openai): validate body-limit error envelope --- backend/internal/handler/openai_body_limit_failover_test.go | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/backend/internal/handler/openai_body_limit_failover_test.go b/backend/internal/handler/openai_body_limit_failover_test.go index ff8970d147..edcb5a946b 100644 --- a/backend/internal/handler/openai_body_limit_failover_test.go +++ b/backend/internal/handler/openai_body_limit_failover_test.go @@ -24,7 +24,8 @@ func TestOpenAIBodyLimitFailoverExhausted_ReturnsRedactedJSON413(t *testing.T) { require.Equal(t, http.StatusRequestEntityTooLarge, rec.Code) var envelope map[string]any require.NoError(t, json.Unmarshal(rec.Body.Bytes(), &envelope)) - errBody := envelope["error"].(map[string]any) + errBody, ok := envelope["error"].(map[string]any) + require.True(t, ok) require.Equal(t, "invalid_request_error", errBody["type"]) require.Equal(t, "Request payload is too large", errBody["message"]) require.NotContains(t, rec.Body.String(), "must-not-leak")