Merge pull request #4611 from superman2003/fix/codex-models-manifest-401-unschedulable
fix(openai): mark OAuth accounts unschedulable on Codex models manifest 401
This commit is contained in:
@@ -53,6 +53,8 @@ func (h *OpenAIGatewayHandler) CodexModels(c *gin.Context) {
|
||||
h.errorResponse(c, http.StatusServiceUnavailable, "upstream_error", "No available OpenAI accounts")
|
||||
return
|
||||
}
|
||||
// 让 ops 错误日志携带实际选中的上游账号,便于定位失效账号(#4544)。
|
||||
setOpsSelectedAccount(c, account.ID, account.Platform)
|
||||
|
||||
manifest, err := h.gatewayService.FetchCodexModelsManifest(c.Request.Context(), account, c.Query("client_version"), c.GetHeader("If-None-Match"))
|
||||
if err != nil {
|
||||
|
||||
@@ -47,6 +47,7 @@ type codexModelsManifestUpstreamError struct {
|
||||
err error
|
||||
retryable bool
|
||||
statusCode int
|
||||
headers http.Header
|
||||
body []byte
|
||||
}
|
||||
|
||||
@@ -56,7 +57,13 @@ func (e *codexModelsManifestUpstreamError) Unwrap() error { return e.err }
|
||||
|
||||
// IsRetryableCodexModelsManifestError reports whether another selected account
|
||||
// may succeed without changing the request. Configuration and upstream 4xx
|
||||
// responses, except 429, are intentionally not retried.
|
||||
// responses, except 429 and ChatGPT-backend 401, are intentionally not
|
||||
// retried. A manifest 401 from the ChatGPT Codex backend reflects the selected
|
||||
// OAuth account's upstream token rather than the client request (the client's
|
||||
// own API key was already validated locally), so a different account may still
|
||||
// serve the manifest. Custom API key upstreams keep the old no-failover 401
|
||||
// behavior because their /models auth semantics are not authoritative for the
|
||||
// account.
|
||||
func IsRetryableCodexModelsManifestError(err error) bool {
|
||||
var upstreamErr *codexModelsManifestUpstreamError
|
||||
return errors.As(err, &upstreamErr) && upstreamErr.retryable
|
||||
@@ -322,6 +329,7 @@ func (s *OpenAIGatewayService) FetchCodexModelsManifest(ctx context.Context, acc
|
||||
}
|
||||
manifest, fetchErr := s.fetchCodexModelsManifestUpstream(ctx, request, ifNoneMatch)
|
||||
if !credAccount.IsOpenAIAgentIdentity() || !isAgentIdentityTaskInvalidCodexModelsError(fetchErr) {
|
||||
s.handleCodexModelsManifestAccountAuthError(ctx, account, credAccount, fetchErr)
|
||||
return manifest, fetchErr
|
||||
}
|
||||
expectedTaskID := strings.TrimSpace(credAccount.GetCredential("task_id"))
|
||||
@@ -349,6 +357,37 @@ func isAgentIdentityTaskInvalidCodexModelsError(err error) bool {
|
||||
isAgentIdentityTaskInvalidHTTPResponse(upstreamErr.statusCode, upstreamErr.body)
|
||||
}
|
||||
|
||||
// handleCodexModelsManifestAccountAuthError feeds manifest 401s from the
|
||||
// ChatGPT Codex backend into the shared upstream-error state machinery
|
||||
// (token cache invalidation, temp-unschedulable cooldown, or permanent
|
||||
// disable for token_revoked/token_invalidated). Without this, an account
|
||||
// whose OAuth token was revoked upstream stays active and schedulable and
|
||||
// keeps being selected for every subsequent /models request (#4544).
|
||||
//
|
||||
// Scope is deliberately limited to plain OAuth accounts: the manifest
|
||||
// endpoint authenticates with the same token as /responses forwarding, so a
|
||||
// 401 is authoritative for the account. Agent Identity accounts are excluded
|
||||
// because their 401s can be task-scoped and have a dedicated recovery flow,
|
||||
// and API key manifests come from custom upstreams whose /models auth may
|
||||
// diverge from their chat endpoints.
|
||||
func (s *OpenAIGatewayService) handleCodexModelsManifestAccountAuthError(ctx context.Context, account, credAccount *Account, err error) {
|
||||
if s == nil || account == nil || err == nil {
|
||||
return
|
||||
}
|
||||
if credAccount == nil || !credAccount.IsOpenAIOAuth() || credAccount.IsOpenAIAgentIdentity() {
|
||||
return
|
||||
}
|
||||
var upstreamErr *codexModelsManifestUpstreamError
|
||||
if !errors.As(err, &upstreamErr) || upstreamErr.statusCode != http.StatusUnauthorized {
|
||||
return
|
||||
}
|
||||
headers := upstreamErr.headers
|
||||
if headers == nil {
|
||||
headers = http.Header{}
|
||||
}
|
||||
s.handleOpenAIAccountUpstreamError(ctx, account, upstreamErr.statusCode, headers, upstreamErr.body)
|
||||
}
|
||||
|
||||
func (s *OpenAIGatewayService) fetchCachedAPIKeyCodexModelsManifest(ctx context.Context, request codexModelsManifestRequest, ifNoneMatch string) (*CodexModelsManifest, error) {
|
||||
if err := ctx.Err(); err != nil {
|
||||
return nil, err
|
||||
@@ -450,8 +489,10 @@ func (s *OpenAIGatewayService) fetchCodexModelsManifestUpstream(ctx context.Cont
|
||||
return nil, &codexModelsManifestUpstreamError{
|
||||
err: infraerrors.Newf(http.StatusBadGateway, "OPENAI_CODEX_MODELS_UPSTREAM_FAILED", "codex models manifest upstream error %d: %s", resp.StatusCode, message),
|
||||
statusCode: resp.StatusCode,
|
||||
headers: resp.Header.Clone(),
|
||||
body: body,
|
||||
retryable: resp.StatusCode == http.StatusTooManyRequests ||
|
||||
retryable: (resp.StatusCode == http.StatusUnauthorized && !request.useAPIKeyUpstream) ||
|
||||
resp.StatusCode == http.StatusTooManyRequests ||
|
||||
(resp.StatusCode >= http.StatusInternalServerError && resp.StatusCode < 600),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1094,6 +1094,147 @@ func TestFetchCodexModelsManifestAPIKeyRejectsBaseURLFragment(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// codexModelsAccountStateRepo records account state transitions triggered by
|
||||
// manifest upstream errors (#4544).
|
||||
type codexModelsAccountStateRepo struct {
|
||||
AccountRepository
|
||||
mu sync.Mutex
|
||||
setErrorCalls int
|
||||
lastErrorMsg string
|
||||
setTempUnschedCalls int
|
||||
lastTempReason string
|
||||
}
|
||||
|
||||
func (r *codexModelsAccountStateRepo) SetError(_ context.Context, _ int64, errorMsg string) error {
|
||||
r.mu.Lock()
|
||||
defer r.mu.Unlock()
|
||||
r.setErrorCalls++
|
||||
r.lastErrorMsg = errorMsg
|
||||
return nil
|
||||
}
|
||||
|
||||
func (r *codexModelsAccountStateRepo) SetTempUnschedulable(_ context.Context, _ int64, _ time.Time, reason string) error {
|
||||
r.mu.Lock()
|
||||
defer r.mu.Unlock()
|
||||
r.setTempUnschedCalls++
|
||||
r.lastTempReason = reason
|
||||
return nil
|
||||
}
|
||||
|
||||
func newCodexModels401TestService(repo AccountRepository) *OpenAIGatewayService {
|
||||
rateLimitService := NewRateLimitService(repo, nil, &config.Config{}, nil, nil)
|
||||
s := &OpenAIGatewayService{rateLimitService: rateLimitService}
|
||||
rateLimitService.SetAccountRuntimeBlocker(s)
|
||||
return s
|
||||
}
|
||||
|
||||
func TestFetchCodexModelsManifestOAuth401MarksAccountUnschedulable(t *testing.T) {
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
||||
w.WriteHeader(http.StatusUnauthorized)
|
||||
_, _ = w.Write([]byte(`{"detail":{"message":"invalid token"}}`))
|
||||
}))
|
||||
defer server.Close()
|
||||
|
||||
original := chatgptCodexModelsURL
|
||||
chatgptCodexModelsURL = server.URL
|
||||
defer func() { chatgptCodexModelsURL = original }()
|
||||
|
||||
repo := &codexModelsAccountStateRepo{}
|
||||
s := newCodexModels401TestService(repo)
|
||||
account := newCodexModelsTestAccount()
|
||||
account.Credentials["refresh_token"] = "test-refresh-token"
|
||||
|
||||
_, err := s.FetchCodexModelsManifest(context.Background(), account, "0.137.0", "")
|
||||
require.Error(t, err)
|
||||
require.True(t, IsRetryableCodexModelsManifestError(err), "manifest 401 should allow account failover")
|
||||
require.Equal(t, 1, repo.setTempUnschedCalls, "OAuth 401 should temp-unschedule the account")
|
||||
require.Equal(t, 0, repo.setErrorCalls)
|
||||
require.True(t, s.isOpenAIAccountRuntimeBlocked(account), "account should be runtime-blocked after manifest 401")
|
||||
}
|
||||
|
||||
func TestFetchCodexModelsManifestOAuth401TokenRevokedDisablesAccount(t *testing.T) {
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
||||
w.WriteHeader(http.StatusUnauthorized)
|
||||
_, _ = w.Write([]byte(`{"error":{"code":"token_revoked","message":"token has been revoked"}}`))
|
||||
}))
|
||||
defer server.Close()
|
||||
|
||||
original := chatgptCodexModelsURL
|
||||
chatgptCodexModelsURL = server.URL
|
||||
defer func() { chatgptCodexModelsURL = original }()
|
||||
|
||||
repo := &codexModelsAccountStateRepo{}
|
||||
s := newCodexModels401TestService(repo)
|
||||
account := newCodexModelsTestAccount()
|
||||
account.Credentials["refresh_token"] = "test-refresh-token"
|
||||
|
||||
_, err := s.FetchCodexModelsManifest(context.Background(), account, "0.137.0", "")
|
||||
require.Error(t, err)
|
||||
require.True(t, IsRetryableCodexModelsManifestError(err))
|
||||
require.Equal(t, 1, repo.setErrorCalls, "revoked token should permanently disable the account")
|
||||
require.Contains(t, repo.lastErrorMsg, "Token revoked")
|
||||
require.Equal(t, 0, repo.setTempUnschedCalls)
|
||||
}
|
||||
|
||||
func TestFetchCodexModelsManifestAgentIdentity401DoesNotDisableAccount(t *testing.T) {
|
||||
key, privateKey := newTestAgentIdentityKey(t)
|
||||
account := &Account{
|
||||
ID: 6,
|
||||
Platform: PlatformOpenAI,
|
||||
Type: AccountTypeOAuth,
|
||||
Credentials: map[string]any{
|
||||
"auth_mode": OpenAIAuthModeAgentIdentity,
|
||||
"agent_runtime_id": key.runtimeID,
|
||||
"agent_private_key": privateKey,
|
||||
"task_id": key.taskID,
|
||||
"chatgpt_account_id": "acc-agent-401",
|
||||
},
|
||||
}
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
||||
w.WriteHeader(http.StatusUnauthorized)
|
||||
_, _ = w.Write([]byte(`{"detail":"some non-task 401"}`))
|
||||
}))
|
||||
defer server.Close()
|
||||
|
||||
original := chatgptCodexModelsURL
|
||||
chatgptCodexModelsURL = server.URL
|
||||
defer func() { chatgptCodexModelsURL = original }()
|
||||
|
||||
repo := &codexModelsAccountStateRepo{}
|
||||
s := newCodexModels401TestService(repo)
|
||||
|
||||
_, err := s.FetchCodexModelsManifest(context.Background(), account, "0.137.0", "")
|
||||
require.Error(t, err)
|
||||
require.Equal(t, 0, repo.setErrorCalls, "agent identity 401s must not disable the account")
|
||||
require.Equal(t, 0, repo.setTempUnschedCalls)
|
||||
}
|
||||
|
||||
func TestFetchCodexModelsManifestAPIKey401KeepsNoFailoverAndNoDisable(t *testing.T) {
|
||||
upstream := &codexModelsHTTPUpstreamStub{do: func(_ *http.Request, _ string, _ int64, _ int) (*http.Response, error) {
|
||||
return &http.Response{
|
||||
StatusCode: http.StatusUnauthorized,
|
||||
Status: "401 Unauthorized",
|
||||
Header: make(http.Header),
|
||||
Body: io.NopCloser(strings.NewReader(`{"error":"invalid api key"}`)),
|
||||
}, nil
|
||||
}}
|
||||
|
||||
repo := &codexModelsAccountStateRepo{}
|
||||
s := newCodexModelsAPIKeyTestService(upstream)
|
||||
s.rateLimitService = NewRateLimitService(repo, nil, &config.Config{}, nil, nil)
|
||||
|
||||
_, err := s.FetchCodexModelsManifest(
|
||||
context.Background(),
|
||||
newCodexModelsAPIKeyTestAccount("https://upstream.example"),
|
||||
"0.144.0",
|
||||
"",
|
||||
)
|
||||
require.Error(t, err)
|
||||
require.False(t, IsRetryableCodexModelsManifestError(err), "custom upstream manifest 401 keeps the no-failover behavior")
|
||||
require.Equal(t, 0, repo.setErrorCalls, "custom upstream manifest 401 must not disable the account")
|
||||
require.Equal(t, 0, repo.setTempUnschedCalls)
|
||||
}
|
||||
|
||||
func TestFetchCodexModelsManifestAPIKeyUpstreamError(t *testing.T) {
|
||||
upstream := &codexModelsHTTPUpstreamStub{do: func(_ *http.Request, _ string, _ int64, _ int) (*http.Response, error) {
|
||||
return &http.Response{
|
||||
|
||||
Reference in New Issue
Block a user