mirror of
https://github.com/Wei-Shaw/sub2api.git
synced 2026-10-07 14:58:23 +08:00
fix(openai): separate native and legacy compaction routing
This commit is contained in:
@@ -63,14 +63,19 @@ func TestNormalizeOpenAIResponsesCompactRequest_RemoteV2StaysOnResponses(t *test
|
||||
require.True(t, ok)
|
||||
|
||||
require.Equal(t, "/v1/responses", c.Request.URL.Path)
|
||||
require.False(t, isOpenAIRemoteCompactPath(c))
|
||||
require.False(t, isOpenAILegacyCompactPath(c))
|
||||
require.Equal(t, body, normalized)
|
||||
require.True(t, gjson.GetBytes(normalized, "stream").Bool())
|
||||
require.True(t, gjson.GetBytes(normalized, "store").Bool())
|
||||
require.Equal(t, "pck-signal-1", gjson.GetBytes(normalized, "prompt_cache_key").String())
|
||||
require.Equal(t, "max", gjson.GetBytes(normalized, "reasoning.effort").String())
|
||||
require.Equal(t, "all_turns", gjson.GetBytes(normalized, "reasoning.context").String())
|
||||
require.True(t, requiresOpenAICompactAccount(c, normalized))
|
||||
legacyCompact := service.IsOpenAIResponsesCompactPath(c)
|
||||
nativeV2 := isBareOpenAIResponsesPath(c) && isOpenAIRemoteCompactionV2Request(normalized)
|
||||
require.False(t, legacyCompact)
|
||||
require.True(t, nativeV2)
|
||||
require.Equal(t, service.OpenAIEndpointCapabilityResponses,
|
||||
openAIResponsesRequiredCapabilityForRequest(false, nativeV2 || legacyCompact, service.PlatformOpenAI))
|
||||
|
||||
reqStream, streamOK := parseOpenAICompatibleStream(normalized)
|
||||
require.True(t, streamOK)
|
||||
@@ -87,105 +92,7 @@ func TestNormalizeOpenAIResponsesCompactRequest_RemoteV2StaysOnResponses(t *test
|
||||
func TestNormalizeOpenAIResponsesCompactRequest_RemoteV2PathAliasesStayOnResponses(t *testing.T) {
|
||||
h := &OpenAIGatewayHandler{}
|
||||
body := []byte(`{"model":"gpt-5.6-sol","stream":true,"input":[{"type":"compaction_trigger"}]}`)
|
||||
for _, path := range []string{"/v1/responses/", "/backend-api/codex/responses"} {
|
||||
t.Run(path, func(t *testing.T) {
|
||||
c := newCompactBodySignalTestContext(t, path, body)
|
||||
|
||||
normalized, ok := h.normalizeOpenAIResponsesCompactRequest(c, zap.NewNop(), body)
|
||||
require.True(t, ok)
|
||||
require.Equal(t, path, c.Request.URL.Path)
|
||||
require.Equal(t, body, normalized)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestRequiresOpenAICompactAccount_RemoteV2UsesNativeRequestSignal(t *testing.T) {
|
||||
h := &OpenAIGatewayHandler{}
|
||||
tests := []struct {
|
||||
name string
|
||||
body []byte
|
||||
betaHeader string
|
||||
wantBefore bool
|
||||
wantAfter bool
|
||||
path string
|
||||
}{
|
||||
{
|
||||
name: "native_v2_headerless",
|
||||
body: []byte(`{"model":"gpt-5.6-sol","stream":true,"input":[{"type":"compaction_trigger"}]}`),
|
||||
wantBefore: true,
|
||||
wantAfter: true,
|
||||
path: "/v1/responses",
|
||||
},
|
||||
{
|
||||
name: "native_v2_without_trigger",
|
||||
body: []byte(`{"model":"gpt-5.6-sol","stream":true,"input":[{"type":"message","role":"user","content":"hello"}]}`),
|
||||
betaHeader: "remote_compaction_v2",
|
||||
wantBefore: false,
|
||||
wantAfter: false,
|
||||
path: "/v1/responses",
|
||||
},
|
||||
{
|
||||
name: "native_v2_unrelated_header",
|
||||
body: []byte(`{"model":"gpt-5.6-sol","stream":true,"input":[{"type":"compaction_trigger"}]}`),
|
||||
betaHeader: "responses_websockets_v2",
|
||||
wantBefore: true,
|
||||
wantAfter: true,
|
||||
path: "/v1/responses",
|
||||
},
|
||||
{
|
||||
name: "native_v2_wrong_case_header",
|
||||
body: []byte(`{"model":"gpt-5.6-sol","stream":true,"input":[{"type":"compaction_trigger"}]}`),
|
||||
betaHeader: "REMOTE_COMPACTION_V2",
|
||||
wantBefore: true,
|
||||
wantAfter: true,
|
||||
path: "/v1/responses",
|
||||
},
|
||||
{
|
||||
name: "responses_subpath_with_native_signal",
|
||||
body: []byte(`{"model":"gpt-5.6-sol","stream":true,"input":[{"type":"compaction_trigger"}]}`),
|
||||
betaHeader: "remote_compaction_v2",
|
||||
wantBefore: false,
|
||||
wantAfter: false,
|
||||
path: "/v1/responses/resp_123/responses",
|
||||
},
|
||||
{
|
||||
name: "stream_false_promotes",
|
||||
body: []byte(`{"model":"gpt-5.6-sol","stream":false,"input":[{"type":"compaction_trigger"}]}`),
|
||||
wantBefore: false,
|
||||
wantAfter: true,
|
||||
path: "/v1/responses",
|
||||
},
|
||||
{
|
||||
name: "stream_absent_promotes",
|
||||
body: []byte(`{"model":"gpt-5.6-sol","input":[{"type":"compaction_trigger"}]}`),
|
||||
betaHeader: "remote_compaction_v2",
|
||||
wantBefore: false,
|
||||
wantAfter: true,
|
||||
path: "/v1/responses",
|
||||
},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
c := newCompactBodySignalTestContext(t, tt.path, tt.body)
|
||||
if tt.betaHeader != "" {
|
||||
c.Request.Header.Set("x-codex-beta-features", tt.betaHeader)
|
||||
}
|
||||
|
||||
require.Equal(t, tt.wantBefore, requiresOpenAICompactAccount(c, tt.body))
|
||||
normalized, ok := h.normalizeOpenAIResponsesCompactRequest(c, zap.NewNop(), tt.body)
|
||||
require.True(t, ok)
|
||||
require.Equal(t, tt.wantAfter, requiresOpenAICompactAccount(c, normalized))
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestRequiresOpenAICompactAccount_RemoteV2SupportsResponsesRootAliases(t *testing.T) {
|
||||
h := &OpenAIGatewayHandler{}
|
||||
body := []byte(`{"model":"gpt-5.6-sol","stream":true,"input":[{"type":"compaction_trigger"}]}`)
|
||||
|
||||
for _, path := range []string{
|
||||
"/v1/responses",
|
||||
"/v1/responses/",
|
||||
"/openai/v1/responses",
|
||||
"/responses",
|
||||
@@ -197,7 +104,132 @@ func TestRequiresOpenAICompactAccount_RemoteV2SupportsResponsesRootAliases(t *te
|
||||
normalized, ok := h.normalizeOpenAIResponsesCompactRequest(c, zap.NewNop(), body)
|
||||
require.True(t, ok)
|
||||
require.Equal(t, path, c.Request.URL.Path)
|
||||
require.True(t, requiresOpenAICompactAccount(c, normalized))
|
||||
require.Equal(t, body, normalized)
|
||||
legacyCompact := service.IsOpenAIResponsesCompactPath(c)
|
||||
nativeV2 := isBareOpenAIResponsesPath(c) && isOpenAIRemoteCompactionV2Request(normalized)
|
||||
require.False(t, legacyCompact)
|
||||
require.True(t, nativeV2)
|
||||
require.Equal(t, service.OpenAIEndpointCapabilityResponses,
|
||||
openAIResponsesRequiredCapabilityForRequest(false, nativeV2 || legacyCompact, service.PlatformOpenAI))
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestOpenAIResponsesCompactionRoutingFlags(t *testing.T) {
|
||||
h := &OpenAIGatewayHandler{}
|
||||
tests := []struct {
|
||||
name string
|
||||
body []byte
|
||||
path string
|
||||
wantLegacyBefore bool
|
||||
wantNativeBefore bool
|
||||
wantLegacyAfter bool
|
||||
wantNativeAfter bool
|
||||
wantCapabilityAfter service.OpenAIEndpointCapability
|
||||
wantPathAfter string
|
||||
wantBodyUnchanged bool
|
||||
}{
|
||||
{
|
||||
name: "native_v2_stream_trigger",
|
||||
body: []byte(`{"model":"gpt-5.6-sol","stream":true,"input":[{"type":"compaction_trigger"}]}`),
|
||||
path: "/v1/responses",
|
||||
wantLegacyBefore: false,
|
||||
wantNativeBefore: true,
|
||||
wantLegacyAfter: false,
|
||||
wantNativeAfter: true,
|
||||
wantCapabilityAfter: service.OpenAIEndpointCapabilityResponses,
|
||||
wantPathAfter: "/v1/responses",
|
||||
wantBodyUnchanged: true,
|
||||
},
|
||||
{
|
||||
name: "native_v2_without_trigger",
|
||||
body: []byte(`{"model":"gpt-5.6-sol","stream":true,"input":[{"type":"message","role":"user","content":"hello"}]}`),
|
||||
path: "/v1/responses",
|
||||
wantLegacyBefore: false,
|
||||
wantNativeBefore: false,
|
||||
wantLegacyAfter: false,
|
||||
wantNativeAfter: false,
|
||||
wantCapabilityAfter: service.OpenAIEndpointCapabilityChatCompletions,
|
||||
wantPathAfter: "/v1/responses",
|
||||
wantBodyUnchanged: true,
|
||||
},
|
||||
{
|
||||
name: "explicit_compact",
|
||||
body: []byte(`{"model":"gpt-5.6-sol","stream":true,"input":[{"type":"compaction_trigger"}]}`),
|
||||
path: "/v1/responses/compact",
|
||||
wantLegacyBefore: true,
|
||||
wantNativeBefore: false,
|
||||
wantLegacyAfter: true,
|
||||
wantNativeAfter: false,
|
||||
wantCapabilityAfter: service.OpenAIEndpointCapabilityResponses,
|
||||
wantPathAfter: "/v1/responses/compact",
|
||||
},
|
||||
{
|
||||
name: "nested_compact",
|
||||
body: []byte(`{"model":"gpt-5.6-sol","stream":true,"input":[{"type":"compaction_trigger"}]}`),
|
||||
path: "/v1/responses/compact/detail",
|
||||
wantLegacyBefore: true,
|
||||
wantNativeBefore: false,
|
||||
wantLegacyAfter: true,
|
||||
wantNativeAfter: false,
|
||||
wantCapabilityAfter: service.OpenAIEndpointCapabilityResponses,
|
||||
wantPathAfter: "/v1/responses/compact/detail",
|
||||
},
|
||||
{
|
||||
name: "responses_subpath_with_native_signal",
|
||||
body: []byte(`{"model":"gpt-5.6-sol","stream":true,"input":[{"type":"compaction_trigger"}]}`),
|
||||
path: "/v1/responses/resp_123/responses",
|
||||
wantLegacyBefore: false,
|
||||
wantNativeBefore: false,
|
||||
wantLegacyAfter: false,
|
||||
wantNativeAfter: false,
|
||||
wantCapabilityAfter: service.OpenAIEndpointCapabilityChatCompletions,
|
||||
wantPathAfter: "/v1/responses/resp_123/responses",
|
||||
wantBodyUnchanged: true,
|
||||
},
|
||||
{
|
||||
name: "stream_false_promotes",
|
||||
body: []byte(`{"model":"gpt-5.6-sol","stream":false,"input":[{"type":"compaction_trigger"}]}`),
|
||||
path: "/v1/responses",
|
||||
wantLegacyBefore: false,
|
||||
wantNativeBefore: false,
|
||||
wantLegacyAfter: true,
|
||||
wantNativeAfter: false,
|
||||
wantCapabilityAfter: service.OpenAIEndpointCapabilityResponses,
|
||||
wantPathAfter: "/v1/responses/compact",
|
||||
},
|
||||
{
|
||||
name: "stream_absent_promotes",
|
||||
body: []byte(`{"model":"gpt-5.6-sol","input":[{"type":"compaction_trigger"}]}`),
|
||||
path: "/v1/responses",
|
||||
wantLegacyBefore: false,
|
||||
wantNativeBefore: false,
|
||||
wantLegacyAfter: true,
|
||||
wantNativeAfter: false,
|
||||
wantCapabilityAfter: service.OpenAIEndpointCapabilityResponses,
|
||||
wantPathAfter: "/v1/responses/compact",
|
||||
},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
c := newCompactBodySignalTestContext(t, tt.path, tt.body)
|
||||
legacyBefore := service.IsOpenAIResponsesCompactPath(c)
|
||||
nativeBefore := isBareOpenAIResponsesPath(c) && isOpenAIRemoteCompactionV2Request(tt.body)
|
||||
require.Equal(t, tt.wantLegacyBefore, legacyBefore)
|
||||
require.Equal(t, tt.wantNativeBefore, nativeBefore)
|
||||
normalized, ok := h.normalizeOpenAIResponsesCompactRequest(c, zap.NewNop(), tt.body)
|
||||
require.True(t, ok)
|
||||
require.Equal(t, tt.wantPathAfter, c.Request.URL.Path)
|
||||
legacyAfter := service.IsOpenAIResponsesCompactPath(c)
|
||||
nativeAfter := isBareOpenAIResponsesPath(c) && isOpenAIRemoteCompactionV2Request(normalized)
|
||||
require.Equal(t, tt.wantLegacyAfter, legacyAfter)
|
||||
require.Equal(t, tt.wantNativeAfter, nativeAfter)
|
||||
require.Equal(t, tt.wantCapabilityAfter,
|
||||
openAIResponsesRequiredCapabilityForRequest(false, nativeAfter || legacyAfter, service.PlatformOpenAI))
|
||||
if tt.wantBodyUnchanged {
|
||||
require.Equal(t, tt.body, normalized)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -279,7 +311,7 @@ func TestNormalizeOpenAIResponsesCompactRequest_NoTriggerUntouched(t *testing.T)
|
||||
normalized, ok := h.normalizeOpenAIResponsesCompactRequest(c, zap.NewNop(), body)
|
||||
require.True(t, ok)
|
||||
require.Equal(t, "/v1/responses", c.Request.URL.Path)
|
||||
require.False(t, isOpenAIRemoteCompactPath(c))
|
||||
require.False(t, isOpenAILegacyCompactPath(c))
|
||||
require.Equal(t, body, normalized)
|
||||
require.True(t, gjson.GetBytes(normalized, "stream").Bool())
|
||||
}
|
||||
|
||||
@@ -105,21 +105,29 @@ func captureHandlerStructuredLog(t *testing.T) (*handlerInMemoryLogSink, func())
|
||||
}
|
||||
}
|
||||
|
||||
func TestIsOpenAIRemoteCompactPath(t *testing.T) {
|
||||
require.False(t, isOpenAIRemoteCompactPath(nil))
|
||||
func TestIsOpenAILegacyCompactPath(t *testing.T) {
|
||||
require.False(t, isOpenAILegacyCompactPath(nil))
|
||||
|
||||
gin.SetMode(gin.TestMode)
|
||||
rec := httptest.NewRecorder()
|
||||
c, _ := gin.CreateTestContext(rec)
|
||||
|
||||
c.Request = httptest.NewRequest(http.MethodPost, "/v1/responses/compact", nil)
|
||||
require.True(t, isOpenAIRemoteCompactPath(c))
|
||||
|
||||
c.Request = httptest.NewRequest(http.MethodPost, "/responses/compact/", nil)
|
||||
require.True(t, isOpenAIRemoteCompactPath(c))
|
||||
|
||||
c.Request = httptest.NewRequest(http.MethodPost, "/v1/responses", nil)
|
||||
require.False(t, isOpenAIRemoteCompactPath(c))
|
||||
for _, test := range []struct {
|
||||
path string
|
||||
want bool
|
||||
}{
|
||||
{path: "/v1/responses/compact", want: true},
|
||||
{path: "/v1/responses/compact/detail", want: true},
|
||||
{path: "/responses/compact/", want: true},
|
||||
{path: "/v1/responses", want: false},
|
||||
{path: "/openai/v1/responses", want: false},
|
||||
{path: "/responses", want: false},
|
||||
{path: "/backend-api/codex/responses", want: false},
|
||||
{path: "/v1/responses/resp_123/cancel", want: false},
|
||||
} {
|
||||
c.Request = httptest.NewRequest(http.MethodPost, test.path, nil)
|
||||
require.Equal(t, test.want, isOpenAILegacyCompactPath(c), test.path)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLogOpenAIRemoteCompactOutcome_Succeeded(t *testing.T) {
|
||||
|
||||
@@ -182,8 +182,11 @@ func openAIResponsesRequiredCapability(imageIntent bool, platform string) servic
|
||||
return service.OpenAIEndpointCapabilityChatCompletions
|
||||
}
|
||||
|
||||
func openAIResponsesRequiredCapabilityForRequest(imageIntent bool, requireCompact bool, platform string) service.OpenAIEndpointCapability {
|
||||
if requireCompact && platform == service.PlatformOpenAI {
|
||||
// openAIResponsesRequiredCapabilityForRequest returns the endpoint capability
|
||||
// required by an image or Responses request. needsResponses includes both the
|
||||
// legacy /responses/compact endpoint and native remote compaction v2.
|
||||
func openAIResponsesRequiredCapabilityForRequest(imageIntent bool, needsResponses bool, platform string) service.OpenAIEndpointCapability {
|
||||
if needsResponses && platform == service.PlatformOpenAI {
|
||||
return service.OpenAIEndpointCapabilityResponses
|
||||
}
|
||||
return openAIResponsesRequiredCapability(imageIntent, platform)
|
||||
@@ -295,6 +298,8 @@ func (h *OpenAIGatewayHandler) Responses(c *gin.Context) {
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
legacyCompact := service.IsOpenAIResponsesCompactPath(c)
|
||||
nativeV2 := isBareOpenAIResponsesPath(c) && isOpenAIRemoteCompactionV2Request(body)
|
||||
// body-signal compact:上游 unary 等待期间向下游发 SSE 注释行心跳,防止
|
||||
// 反向代理空闲超时掐断长压缩连接(#3887)。首拍延迟一个心跳间隔,快速
|
||||
// 失败仍走 JSON+状态码链路;未标记客户端流式或间隔为 0 时是 no-op。
|
||||
@@ -391,7 +396,7 @@ func (h *OpenAIGatewayHandler) Responses(c *gin.Context) {
|
||||
c.Request = c.Request.WithContext(service.WithOpenAIForwardModel(
|
||||
c.Request.Context(),
|
||||
forwardModel,
|
||||
isOpenAIRemoteCompactPath(c),
|
||||
legacyCompact,
|
||||
))
|
||||
|
||||
// 提前校验 function_call_output 是否具备可关联上下文,避免上游 400。
|
||||
@@ -436,7 +441,7 @@ func (h *OpenAIGatewayHandler) Responses(c *gin.Context) {
|
||||
if h.rejectIfCyberSessionBlocked(c, apiKey, sessionHashBody, reqModel, cyberBlockFormatResponses) {
|
||||
return
|
||||
}
|
||||
requireCompact := requiresOpenAICompactAccount(c, body)
|
||||
requireCompact := legacyCompact
|
||||
|
||||
maxAccountSwitches := h.maxAccountSwitches
|
||||
switchCount := 0
|
||||
@@ -453,7 +458,8 @@ func (h *OpenAIGatewayHandler) Responses(c *gin.Context) {
|
||||
// 仅对 OpenAI 平台生效:Grok 生图走独立的 forwardGrokResponses 路径,不应被过滤。
|
||||
// 复用前置权限与并发阶段在未修改 body 上确认的显式生图意图,避免大 tools 请求重复扫描。
|
||||
// 该判断已排除 Codex 被动 image_gen namespace,避免 CC-only 账号被误过滤(#4476)。
|
||||
requiredCapability := openAIResponsesRequiredCapabilityForRequest(imageIntent, requireCompact, requestPlatform)
|
||||
needsResponses := nativeV2 || legacyCompact
|
||||
requiredCapability := openAIResponsesRequiredCapabilityForRequest(imageIntent, needsResponses, requestPlatform)
|
||||
|
||||
// 分组利润控制:请求级装配定价上下文——pricingAt 固定本请求的
|
||||
// D 与计费高峰因子,选号、槽位终检与全部 failover 重入共用同一门与阈值。
|
||||
@@ -495,7 +501,7 @@ func (h *OpenAIGatewayHandler) Responses(c *gin.Context) {
|
||||
zap.Int("excluded_account_count", len(failedAccountIDs)),
|
||||
)
|
||||
if len(failedAccountIDs) == 0 {
|
||||
if errors.Is(err, service.ErrNoAvailableCompactAccounts) {
|
||||
if legacyCompact && errors.Is(err, service.ErrNoAvailableCompactAccounts) {
|
||||
markOpsRoutingCapacityLimitedIfNoAvailable(c, err)
|
||||
h.handleStreamingAwareError(c, http.StatusServiceUnavailable, "compact_not_supported", "No available accounts support /responses/compact", streamStarted)
|
||||
return
|
||||
@@ -751,12 +757,8 @@ func (h *OpenAIGatewayHandler) Responses(c *gin.Context) {
|
||||
}
|
||||
}
|
||||
|
||||
func isOpenAIRemoteCompactPath(c *gin.Context) bool {
|
||||
if c == nil || c.Request == nil || c.Request.URL == nil {
|
||||
return false
|
||||
}
|
||||
normalizedPath := strings.TrimRight(strings.TrimSpace(c.Request.URL.Path), "/")
|
||||
return strings.HasSuffix(normalizedPath, "/responses/compact")
|
||||
func isOpenAILegacyCompactPath(c *gin.Context) bool {
|
||||
return service.IsOpenAIResponsesCompactPath(c)
|
||||
}
|
||||
|
||||
// isBareOpenAIResponsesPath 仅匹配裸 /responses 端点(无 /compact 等子路径),
|
||||
@@ -779,22 +781,12 @@ func isOpenAIRemoteCompactionV2Request(body []byte) bool {
|
||||
return valid && stream && service.HasCompactionTriggerInInput(body)
|
||||
}
|
||||
|
||||
func requiresOpenAICompactAccount(c *gin.Context, body []byte) bool {
|
||||
if isOpenAIRemoteCompactPath(c) {
|
||||
return true
|
||||
}
|
||||
if !isBareOpenAIResponsesPath(c) {
|
||||
return false
|
||||
}
|
||||
return isOpenAIRemoteCompactionV2Request(body)
|
||||
}
|
||||
|
||||
// normalizeOpenAIResponsesCompactRequest keeps Codex remote compaction v2 on
|
||||
// its native streaming /responses wire and preserves the legacy body-signal
|
||||
// promotion for non-streaming requests.
|
||||
// 返回归一化后的 body;ok=false 表示错误响应已写出,调用方应直接 return。
|
||||
func (h *OpenAIGatewayHandler) normalizeOpenAIResponsesCompactRequest(c *gin.Context, reqLog *zap.Logger, body []byte) ([]byte, bool) {
|
||||
isCompactRequest := service.IsOpenAIResponsesCompactPathForTest(c)
|
||||
isCompactRequest := isOpenAILegacyCompactPath(c)
|
||||
if !isCompactRequest && isBareOpenAIResponsesPath(c) && service.HasCompactionTriggerInInput(body) {
|
||||
if isOpenAIRemoteCompactionV2Request(body) {
|
||||
return body, true
|
||||
@@ -825,7 +817,7 @@ func (h *OpenAIGatewayHandler) normalizeOpenAIResponsesCompactRequest(c *gin.Con
|
||||
}
|
||||
|
||||
func (h *OpenAIGatewayHandler) logOpenAIRemoteCompactOutcome(c *gin.Context, startedAt time.Time) {
|
||||
if !isOpenAIRemoteCompactPath(c) {
|
||||
if !isOpenAILegacyCompactPath(c) {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -82,8 +82,10 @@ type OpenAIAccountScheduleRequest struct {
|
||||
RequiredTransport OpenAIUpstreamTransport
|
||||
RequiredCapability OpenAIEndpointCapability
|
||||
RequiredImageCapability OpenAIImagesCapability
|
||||
RequireCompact bool
|
||||
ExcludedIDs map[int64]struct{}
|
||||
// RequireCompact is only for legacy /responses/compact capability filtering
|
||||
// and compact_model_mapping; native remote compaction v2 leaves it false.
|
||||
RequireCompact bool
|
||||
ExcludedIDs map[int64]struct{}
|
||||
}
|
||||
|
||||
type OpenAIAccountScheduleDecision struct {
|
||||
|
||||
@@ -116,6 +116,154 @@ func TestOpenAIGatewayService_SelectAccountWithScheduler_CompactRejectsExplicitl
|
||||
require.Nil(t, selection)
|
||||
}
|
||||
|
||||
func newOpenAICompactionSchedulerTestService(accounts []Account, advanced bool) *OpenAIGatewayService {
|
||||
cfg := &config.Config{}
|
||||
cfg.Gateway.Scheduling.LoadBatchEnabled = false
|
||||
svc := &OpenAIGatewayService{
|
||||
accountRepo: schedulerTestOpenAIAccountRepo{accounts: accounts},
|
||||
cache: &schedulerTestGatewayCache{},
|
||||
cfg: cfg,
|
||||
concurrencyService: NewConcurrencyService(schedulerTestConcurrencyCache{}),
|
||||
}
|
||||
if advanced {
|
||||
svc.rateLimitService = newOpenAIAdvancedSchedulerRateLimitService("true")
|
||||
}
|
||||
return svc
|
||||
}
|
||||
|
||||
func selectOpenAICompactionSchedulerTestAccount(t *testing.T, svc *OpenAIGatewayService, groupID int64, requireCompact bool) (*AccountSelectionResult, error) {
|
||||
t.Helper()
|
||||
selection, _, err := svc.SelectAccountWithSchedulerForCapability(
|
||||
context.Background(),
|
||||
&groupID,
|
||||
"",
|
||||
"",
|
||||
"gpt-5.6-sol",
|
||||
nil,
|
||||
OpenAIUpstreamTransportAny,
|
||||
OpenAIEndpointCapabilityResponses,
|
||||
requireCompact,
|
||||
false,
|
||||
false,
|
||||
)
|
||||
return selection, err
|
||||
}
|
||||
|
||||
func TestOpenAIGatewayService_SelectAccountWithScheduler_NativeCompactionIgnoresLegacyCompactProbe(t *testing.T) {
|
||||
for _, advanced := range []bool{false, true} {
|
||||
t.Run(map[bool]string{false: "legacy_scheduler", true: "advanced_scheduler"}[advanced], func(t *testing.T) {
|
||||
resetOpenAIAdvancedSchedulerSettingCacheForTest()
|
||||
svc := newOpenAICompactionSchedulerTestService([]Account{{
|
||||
ID: 71012,
|
||||
Platform: PlatformOpenAI,
|
||||
Type: AccountTypeAPIKey,
|
||||
Status: StatusActive,
|
||||
Schedulable: true,
|
||||
Concurrency: 1,
|
||||
Extra: map[string]any{
|
||||
"openai_compact_supported": false,
|
||||
"openai_responses_supported": true,
|
||||
},
|
||||
}}, advanced)
|
||||
|
||||
selection, err := selectOpenAICompactionSchedulerTestAccount(t, svc, 91007, false)
|
||||
require.NoError(t, err)
|
||||
require.NotNil(t, selection)
|
||||
require.Equal(t, int64(71012), selection.Account.ID)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestOpenAIGatewayService_SelectAccountWithScheduler_NativeCompactionAllowsForceOff(t *testing.T) {
|
||||
for _, advanced := range []bool{false, true} {
|
||||
t.Run(map[bool]string{false: "legacy_scheduler", true: "advanced_scheduler"}[advanced], func(t *testing.T) {
|
||||
resetOpenAIAdvancedSchedulerSettingCacheForTest()
|
||||
svc := newOpenAICompactionSchedulerTestService([]Account{{
|
||||
ID: 71013,
|
||||
Platform: PlatformOpenAI,
|
||||
Type: AccountTypeAPIKey,
|
||||
Status: StatusActive,
|
||||
Schedulable: true,
|
||||
Concurrency: 1,
|
||||
Extra: map[string]any{
|
||||
"openai_compact_mode": OpenAICompactModeForceOff,
|
||||
"openai_responses_supported": true,
|
||||
},
|
||||
}}, advanced)
|
||||
|
||||
selection, err := selectOpenAICompactionSchedulerTestAccount(t, svc, 91008, false)
|
||||
require.NoError(t, err)
|
||||
require.NotNil(t, selection)
|
||||
require.Equal(t, int64(71013), selection.Account.ID)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestOpenAIGatewayService_SelectAccountWithScheduler_NativeCompactionRequiresResponsesCapability(t *testing.T) {
|
||||
for _, advanced := range []bool{false, true} {
|
||||
t.Run(map[bool]string{false: "legacy_scheduler", true: "advanced_scheduler"}[advanced], func(t *testing.T) {
|
||||
resetOpenAIAdvancedSchedulerSettingCacheForTest()
|
||||
svc := newOpenAICompactionSchedulerTestService([]Account{{
|
||||
ID: 71014,
|
||||
Platform: PlatformOpenAI,
|
||||
Type: AccountTypeAPIKey,
|
||||
Status: StatusActive,
|
||||
Schedulable: true,
|
||||
Concurrency: 1,
|
||||
Extra: map[string]any{
|
||||
"openai_compact_mode": OpenAICompactModeForceOn,
|
||||
"openai_responses_supported": false,
|
||||
},
|
||||
}}, advanced)
|
||||
|
||||
selection, err := selectOpenAICompactionSchedulerTestAccount(t, svc, 91009, false)
|
||||
require.ErrorIs(t, err, ErrNoAvailableAccounts)
|
||||
require.NotErrorIs(t, err, ErrNoAvailableCompactAccounts)
|
||||
require.NotContains(t, err.Error(), "/responses/compact")
|
||||
require.Nil(t, selection)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestOpenAIGatewayService_SelectAccountWithScheduler_LegacyCompactionKeepsCompactEligibility(t *testing.T) {
|
||||
for _, advanced := range []bool{false, true} {
|
||||
t.Run(map[bool]string{false: "legacy_scheduler", true: "advanced_scheduler"}[advanced], func(t *testing.T) {
|
||||
resetOpenAIAdvancedSchedulerSettingCacheForTest()
|
||||
svc := newOpenAICompactionSchedulerTestService([]Account{
|
||||
{
|
||||
ID: 71015,
|
||||
Platform: PlatformOpenAI,
|
||||
Type: AccountTypeAPIKey,
|
||||
Status: StatusActive,
|
||||
Schedulable: true,
|
||||
Concurrency: 1,
|
||||
Extra: map[string]any{
|
||||
"openai_compact_supported": false,
|
||||
"openai_responses_supported": true,
|
||||
},
|
||||
},
|
||||
{
|
||||
ID: 71016,
|
||||
Platform: PlatformOpenAI,
|
||||
Type: AccountTypeAPIKey,
|
||||
Status: StatusActive,
|
||||
Schedulable: true,
|
||||
Concurrency: 1,
|
||||
Extra: map[string]any{
|
||||
"openai_compact_mode": OpenAICompactModeForceOff,
|
||||
"openai_responses_supported": true,
|
||||
},
|
||||
},
|
||||
}, advanced)
|
||||
|
||||
selection, err := selectOpenAICompactionSchedulerTestAccount(t, svc, 91010, true)
|
||||
require.ErrorIs(t, err, ErrNoAvailableCompactAccounts)
|
||||
require.Contains(t, err.Error(), "/responses/compact")
|
||||
require.Nil(t, selection)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// TestOpenAIGatewayService_SelectAccountWithScheduler_CompactRequiresResponsesCapability
|
||||
// prevents a confirmed Chat Completions-only API key from silently dropping the
|
||||
// compaction trigger through the Responses-to-Chat fallback.
|
||||
|
||||
@@ -10,9 +10,9 @@ type openAIForwardModel struct {
|
||||
}
|
||||
|
||||
// WithOpenAIForwardModel records the model present in the forwarded request
|
||||
// body after channel mapping and whether the legacy compact-only model mapping
|
||||
// applies. Channel restriction checks then follow the same model chain used by
|
||||
// Forward.
|
||||
// body after channel mapping and whether the legacy /responses/compact-only
|
||||
// model mapping applies. Native remote compaction v2 keeps this false, so
|
||||
// channel restriction checks follow the same model chain used by Forward.
|
||||
func WithOpenAIForwardModel(ctx context.Context, forwardModel string, useCompactModelMapping bool) context.Context {
|
||||
if ctx == nil {
|
||||
ctx = context.Background()
|
||||
|
||||
@@ -276,7 +276,9 @@ func isOpenAIEncryptedReasoningInputItem(item any) bool {
|
||||
return has
|
||||
}
|
||||
|
||||
func IsOpenAIResponsesCompactPathForTest(c *gin.Context) bool {
|
||||
// IsOpenAIResponsesCompactPath reports whether the request targets the legacy
|
||||
// /responses/compact endpoint, including its forwardable subpaths.
|
||||
func IsOpenAIResponsesCompactPath(c *gin.Context) bool {
|
||||
return isOpenAIResponsesCompactPath(c)
|
||||
}
|
||||
|
||||
|
||||
@@ -244,7 +244,7 @@ func (s *OpenAIGatewayService) SelectAccountForModelWithExclusions(ctx context.C
|
||||
}
|
||||
|
||||
// noAvailableOpenAISelectionError builds the standard "no account available" error
|
||||
// while preserving the compact-specific error when applicable.
|
||||
// while preserving the legacy /responses/compact error when applicable.
|
||||
func normalizeOpenAICompatiblePlatform(platform string) string {
|
||||
if platform == PlatformGrok {
|
||||
return PlatformGrok
|
||||
@@ -685,8 +685,8 @@ func prioritizeOpenAICompactAccounts(accounts []*Account) []*Account {
|
||||
}
|
||||
|
||||
// resolveOpenAIAccountUpstreamModelForRequest resolves the upstream model that
|
||||
// would be sent for a given request, honouring compact-only mappings when the
|
||||
// caller is on the /responses/compact path.
|
||||
// would be sent for a given request, honoring the legacy compact-only mapping
|
||||
// when the caller is on the /responses/compact path.
|
||||
func resolveOpenAIAccountUpstreamModelForRequest(account *Account, requestedModel string, requireCompact bool) string {
|
||||
// Forward checks the raw Chat Completions fallback before passthrough.
|
||||
// These API-key accounts therefore apply normal account model_mapping and
|
||||
@@ -842,7 +842,8 @@ func (s *OpenAIGatewayService) tryStickySessionHit(ctx context.Context, groupID
|
||||
// selectBestAccount selects the best account from candidates (priority + LRU).
|
||||
// Returns nil if no available account. The second return reports whether at
|
||||
// least one candidate was filtered out solely because it lacks compact support
|
||||
// (only meaningful when requireCompact=true); the third contains deterministic
|
||||
// (only meaningful when the legacy /responses/compact requireCompact flag is
|
||||
// true); the third contains deterministic
|
||||
// exclusion diagnostics for the evaluated snapshot.
|
||||
func (s *OpenAIGatewayService) selectBestAccount(ctx context.Context, groupID *int64, platform string, accounts []Account, requestedModel string, excludedIDs map[int64]struct{}, requireCompact bool, requiredCapability OpenAIEndpointCapability, preferLowUpstreamRate bool) (*Account, bool, openAISelectionFilterStats) {
|
||||
platform = normalizeOpenAICompatiblePlatform(platform)
|
||||
|
||||
@@ -398,8 +398,8 @@ func (t *accountWriteThrottle) Allow(id int64, now time.Time) bool {
|
||||
|
||||
var defaultOpenAICodexSnapshotPersistThrottle = newAccountWriteThrottle(openAICodexSnapshotPersistMinInterval)
|
||||
|
||||
// ErrNoAvailableCompactAccounts indicates the request needs /responses/compact
|
||||
// support but no compatible account is available.
|
||||
// ErrNoAvailableCompactAccounts indicates a legacy /responses/compact request
|
||||
// needs compact support but no compatible account is available.
|
||||
var ErrNoAvailableCompactAccounts = errors.New("no available accounts support /responses/compact")
|
||||
|
||||
// OpenAIGatewayService handles OpenAI API gateway operations
|
||||
|
||||
@@ -145,6 +145,34 @@ func TestOpenAIResponsesRequestPathSuffixRejectsNonConformingSubpaths(t *testing
|
||||
}
|
||||
}
|
||||
|
||||
func TestIsOpenAIResponsesCompactPathUsesLegacyEndpointShape(t *testing.T) {
|
||||
legacyPaths := []string{
|
||||
"/v1/responses/compact",
|
||||
"/v1/responses/compact/detail",
|
||||
"/responses/compact/",
|
||||
}
|
||||
for _, path := range legacyPaths {
|
||||
t.Run("legacy_"+path, func(t *testing.T) {
|
||||
c := newResponsesSuffixTestContext(t, path)
|
||||
require.True(t, IsOpenAIResponsesCompactPath(c))
|
||||
})
|
||||
}
|
||||
|
||||
nonLegacyPaths := []string{
|
||||
"/v1/responses",
|
||||
"/openai/v1/responses",
|
||||
"/responses",
|
||||
"/backend-api/codex/responses",
|
||||
"/v1/responses/resp_123/cancel",
|
||||
}
|
||||
for _, path := range nonLegacyPaths {
|
||||
t.Run("non_legacy_"+path, func(t *testing.T) {
|
||||
c := newResponsesSuffixTestContext(t, path)
|
||||
require.False(t, IsOpenAIResponsesCompactPath(c))
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestAppendOpenAIResponsesRequestPathSuffixRefusesUnsafeSuffix(t *testing.T) {
|
||||
// 调用方漏了校验时,拼接函数本身也不得把不合规片段带进上游 URL。
|
||||
require.Equal(t, chatgptCodexURL, appendOpenAIResponsesRequestPathSuffix(chatgptCodexURL, "/../../x"))
|
||||
|
||||
Reference in New Issue
Block a user