mirror of
https://github.com/Wei-Shaw/sub2api.git
synced 2026-09-24 16:05:44 +08:00
Merge pull request #4395 from wp-a/fix/openai-body-limit-failover
[codex] fail over account-specific OpenAI body limits
This commit is contained in:
@@ -0,0 +1,59 @@
|
||||
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, 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")
|
||||
}
|
||||
|
||||
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",
|
||||
}
|
||||
}
|
||||
@@ -2085,6 +2085,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)
|
||||
|
||||
@@ -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 上游端点。
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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
|
||||
@@ -590,12 +593,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(
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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)
|
||||
})
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user