fix(openai): send session-level beta features and probe native compaction v2

OpenAI sunset the legacy unary /responses/compact endpoint (404, #5598,
#5624), so the account "compact probe" in the admin UI kept failing even for
healthy accounts, and the beta-feature negotiation header was only attached
to compaction turns.

Beta features (codex-rs session/mod.rs build_model_client_beta_features_header
+ client.rs build_responses_headers): the header is a session-level constant
attached to every /responses request, the WS handshake and /responses/compact.
Enumerating FEATURES shows no Experimental feature is enabled by default, so a
default install sends exactly "remote_compaction_v2". Mirror that:

- OAuth requests without a client-declared header get the default shape, so we
  no longer produce a "header only on compaction turns" pattern real Codex
  never emits (#5586 chains that strip the header)
- a client-declared header is preserved as-is: non-empty without v2 means the
  user disabled the feature and the gateway must not rewrite that
- native v2 turns (compaction_trigger in body) always ensure v2 is present
- non-OAuth upstreams keep the compaction-turn-only behaviour
- the WS injection sits outside the client-header copy block so prewarm and
  turn handshakes cannot land in different pool compatibility buckets

Compact probe now exercises native v2 (streaming /responses +
compaction_trigger) instead of the dead endpoint. Success requires an actual
compaction output item — scanning output_item.done/added, the terminal
response.output[] and the whole-JSON fallback — so a 2xx that silently drops
the trigger is reported as unsupported (the "got 0 items" class, #5478,
#5648). Probe identity is now UUID-shaped and applies the account's
convergence, matching real traffic on the same endpoint.
This commit is contained in:
shaw
2026-08-15 16:35:08 +08:00
parent 8219dcfc87
commit 8ae6d8f67e
8 changed files with 428 additions and 45 deletions
@@ -300,6 +300,11 @@ func (h *OpenAIGatewayHandler) Responses(c *gin.Context) {
}
legacyCompact := service.IsOpenAIResponsesCompactPath(c)
nativeV2 := isBareOpenAIResponsesPath(c) && isOpenAIRemoteCompactionV2Request(body)
if nativeV2 {
// 原生 v2 压缩出站前补注 x-codex-beta-features: remote_compaction_v2,
// 与真实 Codex 线型一致(网关链剥头后本级负责恢复,#5586)。
service.MarkOpenAINativeCompactionV2(c)
}
// body-signal compact:上游 unary 等待期间向下游发 SSE 注释行心跳,防止
// 反向代理空闲超时掐断长压缩连接(#3887)。首拍延迟一个心跳间隔,快速
// 失败仍走 JSON+状态码链路;未标记客户端流式或间隔为 0 时是 no-op。
@@ -610,10 +610,11 @@ func (s *AccountTestService) testOpenAIAccountConnection(c *gin.Context, account
}
// Align test routing with gateway behavior: OpenAI accounts apply normal
// account model mapping, and compact mode applies compact-only mapping on top.
// account model mapping. Native remote compaction v2 rides the ordinary
// /responses wire and does NOT apply the legacy compact-only mapping
// (post-#5641 semantics: compact_model_mapping is /responses/compact-only).
testModelID = account.GetMappedModel(testModelID)
if mode == AccountTestModeCompact {
testModelID = resolveOpenAICompactForwardModel(account, testModelID)
return s.testOpenAICompactConnection(c, account, testModelID)
}
@@ -1976,8 +1977,10 @@ func (s *AccountTestService) testOpenAIChatCompletionsConnection(
return s.processOpenAIChatCompletionsStream(c, resp.Body)
}
// testOpenAICompactConnection probes /responses/compact and persists the
// resulting capability state on the account.
// testOpenAICompactConnection probes native remote compaction v2 (streaming
// /responses with a compaction_trigger input item) and persists the resulting
// capability state on the account. The legacy unary /responses/compact
// endpoint has been sunset upstream (404, #5598/#5624) and is no longer probed.
func (s *AccountTestService) testOpenAICompactConnection(c *gin.Context, account *Account, testModelID string) error {
ctx := c.Request.Context()
credentialAccount := account
@@ -2002,7 +2005,7 @@ func (s *AccountTestService) testOpenAICompactConnection(c *gin.Context, account
if authToken == "" && !credentialAccount.IsOpenAIAgentIdentity() {
return s.sendErrorAndEnd(c, "No access token available")
}
apiURL = chatgptCodexAPIURL + "/compact"
apiURL = chatgptCodexAPIURL
case account.Type == AccountTypeAPIKey:
authToken = account.GetOpenAIApiKey()
if authToken == "" {
@@ -2016,7 +2019,7 @@ func (s *AccountTestService) testOpenAICompactConnection(c *gin.Context, account
if err != nil {
return s.sendErrorAndEnd(c, fmt.Sprintf("Invalid base URL: %s", err.Error()))
}
apiURL = appendOpenAIResponsesRequestPathSuffix(buildOpenAIResponsesURL(normalizedBaseURL), "/compact")
apiURL = buildOpenAIResponsesURL(normalizedBaseURL)
default:
return s.sendErrorAndEnd(c, fmt.Sprintf("Unsupported account type: %s", account.Type))
}
@@ -2027,7 +2030,11 @@ func (s *AccountTestService) testOpenAICompactConnection(c *gin.Context, account
c.Writer.Header().Set("X-Accel-Buffering", "no")
c.Writer.Flush()
payloadBytes, _ := json.Marshal(createOpenAICompactProbePayload(testModelID))
// 原生 v2 走普通 /responses 线:OAuth 与真实转发一致做上游模型归一化。
if isOAuth {
testModelID = normalizeOpenAIModelForUpstream(credentialAccount, testModelID)
}
payloadBytes, _ := json.Marshal(createOpenAICompactProbePayload(testModelID, isOAuth))
if !agentIdentityTaskRecoveryWasTried(ctx) {
s.sendEvent(c, TestEvent{Type: "test_start", Model: testModelID})
}
@@ -2039,7 +2046,9 @@ func (s *AccountTestService) testOpenAICompactConnection(c *gin.Context, account
req = req.WithContext(WithHTTPUpstreamProfile(req.Context(), HTTPUpstreamProfileOpenAI))
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Accept", "application/json")
// v2 探测是流式请求;同时补注协商头,与真实 codex 出站线型一致。
req.Header.Set("Accept", "text/event-stream")
ensureOpenAIRemoteCompactionV2BetaFeature(req.Header)
if credentialAccount.IsOpenAIAgentIdentity() {
authHeaders, authErr := buildAgentIdentityAuthenticationHeaders(ctx, s.accountRepo, s.agentIdentityWS, &s.agentIdentityTaskMu, credentialAccount)
if authErr != nil {
@@ -2061,6 +2070,12 @@ func (s *AccountTestService) testOpenAICompactConnection(c *gin.Context, account
if isOAuth {
req.Host = "chatgpt.com"
setOpenAIChatGPTAccountHeaders(req.Header, credentialAccount)
// 指纹收敛:探测与真实转发走同一个 /responses 端点,身份也必须同构,
// 否则探测流量会以「缺 x-codex-installation-id + 非收敛 session」的
// 形态暴露在上游眼里。账号关闭收敛(off)时返回 nil,探测保持原样。
if fpIDs := resolveCodexFingerprintIDsFromRequest(account, req.Header); fpIDs != nil {
applyCodexFingerprintHeaders(req.Header, fpIDs)
}
}
// 账号级请求头覆写:测试请求与真实转发保持一致的最终头
@@ -2074,7 +2089,7 @@ func (s *AccountTestService) testOpenAICompactConnection(c *gin.Context, account
resp, err := s.httpUpstream.DoWithTLS(req, proxyURL, account.ID, account.Concurrency, s.tlsFPProfileService.ResolveTLSProfile(account))
if err != nil {
if s.accountRepo != nil {
updates := buildOpenAICompactProbeExtraUpdates(nil, nil, err, time.Now())
updates := buildOpenAICompactProbeExtraUpdates(nil, nil, err, false, time.Now())
_ = s.accountRepo.UpdateExtra(ctx, account.ID, updates)
mergeAccountExtra(account, updates)
}
@@ -2093,8 +2108,9 @@ func (s *AccountTestService) testOpenAICompactConnection(c *gin.Context, account
return s.testOpenAICompactConnection(c, account, testModelID)
}
compactionFound := openAICompactProbeFoundCompactionItem(body)
if s.accountRepo != nil {
updates := buildOpenAICompactProbeExtraUpdates(resp, body, nil, time.Now())
updates := buildOpenAICompactProbeExtraUpdates(resp, body, nil, compactionFound, time.Now())
if codexUpdates, err := extractOpenAICodexProbeUpdates(resp); err == nil && len(codexUpdates) > 0 {
updates = mergeExtraUpdates(updates, codexUpdates)
}
@@ -2116,7 +2132,11 @@ func (s *AccountTestService) testOpenAICompactConnection(c *gin.Context, account
return s.sendErrorAndEnd(c, fmt.Sprintf("API returned %d: %s", resp.StatusCode, string(body)))
}
s.sendEvent(c, TestEvent{Type: "content", Text: "Compact probe succeeded"})
if !compactionFound {
return s.sendErrorAndEnd(c, "Upstream returned 2xx without a compaction output item (native remote compaction v2 unsupported on this chain)")
}
s.sendEvent(c, TestEvent{Type: "content", Text: "Compact probe succeeded (native remote compaction v2)"})
s.sendEvent(c, TestEvent{Type: "test_complete", Success: true})
return nil
}
@@ -10,10 +10,16 @@ import (
"github.com/Wei-Shaw/sub2api/internal/config"
"github.com/gin-gonic/gin"
"github.com/google/uuid"
"github.com/stretchr/testify/require"
"github.com/tidwall/gjson"
)
// compactProbeSSESuccessBody 是原生 v2 压缩成功的最小 SSE 形态:
// output_item.done 携带 compaction item + response.completed。
const compactProbeSSESuccessBody = "data: {\"type\":\"response.output_item.done\",\"item\":{\"type\":\"compaction\",\"id\":\"cmp_probe\",\"encrypted_content\":\"blob\"}}\n\n" +
"data: {\"type\":\"response.completed\",\"response\":{\"id\":\"resp_probe\",\"output\":[]}}\n\n"
func TestAccountTestService_TestAccountConnection_OpenAICompactOAuthSuccessPersistsSupport(t *testing.T) {
gin.SetMode(gin.TestMode)
@@ -38,8 +44,8 @@ func TestAccountTestService_TestAccountConnection_OpenAICompactOAuthSuccessPersi
}
upstream := &httpUpstreamRecorder{resp: &http.Response{
StatusCode: http.StatusOK,
Header: http.Header{"Content-Type": []string{"application/json"}, "x-request-id": []string{"rid-probe"}},
Body: io.NopCloser(strings.NewReader(`{"id":"cmp_probe","status":"completed"}`)),
Header: http.Header{"Content-Type": []string{"text/event-stream"}, "x-request-id": []string{"rid-probe"}},
Body: io.NopCloser(strings.NewReader(compactProbeSSESuccessBody)),
}}
svc := &AccountTestService{
accountRepo: repo,
@@ -53,16 +59,22 @@ func TestAccountTestService_TestAccountConnection_OpenAICompactOAuthSuccessPersi
err := svc.TestAccountConnection(c, account.ID, "gpt-5.4", "", AccountTestModeCompact)
require.NoError(t, err)
require.Equal(t, chatgptCodexAPIURL+"/compact", upstream.lastReq.URL.String())
// 原生 v2:探测普通 /responses 线,不再打已下线的 /responses/compact。
require.Equal(t, chatgptCodexAPIURL, upstream.lastReq.URL.String())
require.Equal(t, "chatgpt.com", upstream.lastReq.Host)
require.Equal(t, "application/json", upstream.lastReq.Header.Get("Accept"))
require.Equal(t, codexCLIVersion, upstream.lastReq.Header.Get("Version"))
require.Equal(t, "text/event-stream", upstream.lastReq.Header.Get("Accept"))
require.Contains(t, upstream.lastReq.Header.Get("x-codex-beta-features"), "remote_compaction_v2")
require.NotEmpty(t, upstream.lastReq.Header.Get("Session_Id"))
require.Equal(t, HTTPUpstreamProfileOpenAI, HTTPUpstreamProfileFromContext(upstream.lastReq.Context()))
require.Equal(t, codexCLIUserAgent, upstream.lastReq.Header.Get("User-Agent"))
require.Equal(t, "chatgpt-acc", upstream.lastReq.Header.Get("chatgpt-account-id"))
require.Equal(t, "true", upstream.lastReq.Header.Get("x-openai-fedramp"))
require.Equal(t, "gpt-5.4", gjson.GetBytes(upstream.lastBody, "model").String())
require.True(t, gjson.GetBytes(upstream.lastBody, "stream").Bool())
require.False(t, gjson.GetBytes(upstream.lastBody, "store").Bool())
inputItems := gjson.GetBytes(upstream.lastBody, "input").Array()
require.NotEmpty(t, inputItems)
require.Equal(t, "compaction_trigger", inputItems[len(inputItems)-1].Get("type").String())
updates := <-updateCalls
require.Equal(t, true, updates["openai_compact_supported"])
@@ -114,7 +126,7 @@ func TestAccountTestService_TestAccountConnection_OpenAICompactOAuth404MarksUnsu
require.Contains(t, rec.Body.String(), `"type":"error"`)
}
func TestAccountTestService_TestAccountConnection_OpenAICompactAPIKeyUsesCompactPath(t *testing.T) {
func TestAccountTestService_TestAccountConnection_OpenAICompactAPIKeyUsesNativeResponsesPath(t *testing.T) {
gin.SetMode(gin.TestMode)
updateCalls := make(chan map[string]any, 1)
@@ -127,8 +139,10 @@ func TestAccountTestService_TestAccountConnection_OpenAICompactAPIKeyUsesCompact
Schedulable: true,
Concurrency: 1,
Credentials: map[string]any{
"api_key": "sk-test",
"base_url": "https://example.com/v1",
"api_key": "sk-test",
"base_url": "https://example.com/v1",
// post-#5641:compact_model_mapping 仅作用于 legacy /responses/compact,
// 原生 v2 探测不应用它。
"compact_model_mapping": map[string]any{"gpt-5.4": "gpt-5.4-openai-compact"},
},
}
@@ -138,8 +152,8 @@ func TestAccountTestService_TestAccountConnection_OpenAICompactAPIKeyUsesCompact
}
upstream := &httpUpstreamRecorder{resp: &http.Response{
StatusCode: http.StatusOK,
Header: http.Header{"Content-Type": []string{"application/json"}},
Body: io.NopCloser(strings.NewReader(`{"id":"cmp_probe_apikey","status":"completed"}`)),
Header: http.Header{"Content-Type": []string{"text/event-stream"}},
Body: io.NopCloser(strings.NewReader(compactProbeSSESuccessBody)),
}}
svc := &AccountTestService{
accountRepo: repo,
@@ -154,14 +168,16 @@ func TestAccountTestService_TestAccountConnection_OpenAICompactAPIKeyUsesCompact
err := svc.TestAccountConnection(c, account.ID, "gpt-5.4", "", AccountTestModeCompact)
require.NoError(t, err)
require.Equal(t, "https://example.com/v1/responses/compact", upstream.lastReq.URL.String())
require.Equal(t, "https://example.com/v1/responses", upstream.lastReq.URL.String())
requireOpenAICodexProbeHeaders(t, upstream.lastReq.Header)
require.Equal(t, "gpt-5.4-openai-compact", gjson.GetBytes(upstream.lastBody, "model").String())
require.Contains(t, upstream.lastReq.Header.Get("x-codex-beta-features"), "remote_compaction_v2")
require.Equal(t, "gpt-5.4", gjson.GetBytes(upstream.lastBody, "model").String(),
"原生 v2 探测不应用 compact_model_mapping")
updates := <-updateCalls
require.Equal(t, true, updates["openai_compact_supported"])
}
func TestAccountTestService_TestAccountConnection_OpenAICompactAPIKeyDefaultBaseURLUsesV1Path(t *testing.T) {
func TestAccountTestService_TestAccountConnection_OpenAICompactAPIKeyDefaultBaseURLUsesResponsesPath(t *testing.T) {
gin.SetMode(gin.TestMode)
updateCalls := make(chan map[string]any, 1)
@@ -183,8 +199,8 @@ func TestAccountTestService_TestAccountConnection_OpenAICompactAPIKeyDefaultBase
}
upstream := &httpUpstreamRecorder{resp: &http.Response{
StatusCode: http.StatusOK,
Header: http.Header{"Content-Type": []string{"application/json"}},
Body: io.NopCloser(strings.NewReader(`{"id":"cmp_probe_apikey_default","status":"completed"}`)),
Header: http.Header{"Content-Type": []string{"text/event-stream"}},
Body: io.NopCloser(strings.NewReader(compactProbeSSESuccessBody)),
}}
svc := &AccountTestService{
accountRepo: repo,
@@ -198,6 +214,112 @@ func TestAccountTestService_TestAccountConnection_OpenAICompactAPIKeyDefaultBase
err := svc.TestAccountConnection(c, account.ID, "gpt-5.4", "", AccountTestModeCompact)
require.NoError(t, err)
require.Equal(t, "https://api.openai.com/v1/responses/compact", upstream.lastReq.URL.String())
require.Equal(t, "https://api.openai.com/v1/responses", upstream.lastReq.URL.String())
<-updateCalls
}
func TestAccountTestService_TestAccountConnection_OpenAICompact2xxWithoutItemMarksUnsupported(t *testing.T) {
gin.SetMode(gin.TestMode)
updateCalls := make(chan map[string]any, 1)
account := Account{
ID: 5,
Name: "openai-oauth-no-item",
Platform: PlatformOpenAI,
Type: AccountTypeOAuth,
Status: StatusActive,
Schedulable: true,
Concurrency: 1,
Credentials: map[string]any{
"access_token": "oauth-token",
"chatgpt_account_id": "chatgpt-acc",
},
}
repo := &snapshotUpdateAccountRepo{
stubOpenAIAccountRepo: stubOpenAIAccountRepo{accounts: []Account{account}},
updateExtraCalls: updateCalls,
}
// 200 但流里没有 compaction item:链路吞掉了 compaction_trigger 的形态
//(#5478 的 "got 0 items"),必须判定为不支持。
noItemBody := "data: {\"type\":\"response.output_item.done\",\"item\":{\"type\":\"message\",\"id\":\"msg_1\"}}\n\n" +
"data: {\"type\":\"response.completed\",\"response\":{\"id\":\"resp_x\",\"output\":[]}}\n\n"
upstream := &httpUpstreamRecorder{resp: &http.Response{
StatusCode: http.StatusOK,
Header: http.Header{"Content-Type": []string{"text/event-stream"}},
Body: io.NopCloser(strings.NewReader(noItemBody)),
}}
svc := &AccountTestService{
accountRepo: repo,
httpUpstream: upstream,
}
rec := httptest.NewRecorder()
c, _ := gin.CreateTestContext(rec)
c.Request = httptest.NewRequest(http.MethodPost, "/api/v1/admin/accounts/5/test", bytes.NewReader(nil))
err := svc.TestAccountConnection(c, account.ID, "gpt-5.4", "", AccountTestModeCompact)
require.Error(t, err)
updates := <-updateCalls
require.Equal(t, false, updates["openai_compact_supported"])
require.Contains(t, rec.Body.String(), `"type":"error"`)
}
// 探测与真实转发走同一 /responses 端点,出站身份必须与真实 Codex 同构:
// session/thread 为 UUID、携带 x-codex-installation-id(收敛账号用收敛值)。
func TestAccountTestService_TestAccountConnection_OpenAICompactProbeIdentityMatchesRealTraffic(t *testing.T) {
gin.SetMode(gin.TestMode)
updateCalls := make(chan map[string]any, 1)
account := Account{
ID: 6,
Name: "openai-oauth-identity",
Platform: PlatformOpenAI,
Type: AccountTypeOAuth,
Status: StatusActive,
Schedulable: true,
Concurrency: 1,
Credentials: map[string]any{
"access_token": "oauth-token",
"chatgpt_account_id": "chatgpt-acc",
},
// 收敛是显式 opt-in(#5610),这里显式开启以验证探测身份与真实流量同构。
Extra: map[string]any{"codex_fingerprint_mode": "session"},
}
repo := &snapshotUpdateAccountRepo{
stubOpenAIAccountRepo: stubOpenAIAccountRepo{accounts: []Account{account}},
updateExtraCalls: updateCalls,
}
upstream := &httpUpstreamRecorder{resp: &http.Response{
StatusCode: http.StatusOK,
Header: http.Header{"Content-Type": []string{"text/event-stream"}},
Body: io.NopCloser(strings.NewReader(compactProbeSSESuccessBody)),
}}
svc := &AccountTestService{accountRepo: repo, httpUpstream: upstream}
rec := httptest.NewRecorder()
c, _ := gin.CreateTestContext(rec)
c.Request = httptest.NewRequest(http.MethodPost, "/api/v1/admin/accounts/6/test", bytes.NewReader(nil))
require.NoError(t, svc.TestAccountConnection(c, account.ID, "gpt-5.4", "", AccountTestModeCompact))
// 显式 session 收敛模式:出站身份 = 账号级收敛值
converged := resolveConvergedSessionID(&account)
require.Equal(t, converged, upstream.lastReq.Header.Get("session-id"))
require.Equal(t, converged, upstream.lastReq.Header.Get("session_id"))
require.Equal(t, resolveConvergedInstallationID(&account), upstream.lastReq.Header.Get("x-codex-installation-id"),
"真实 Codex 每个请求必带 installation-id,探测不得缺失")
require.NotContains(t, upstream.lastReq.Header.Get("session-id"), "probe_compact",
"探测标识不得是可被上游一眼识别的字面量")
<-updateCalls
}
func TestCompactProbeSessionID_IsUUIDShaped(t *testing.T) {
for _, id := range []int64{0, 1, 987654} {
got := compactProbeSessionID(id)
_, err := uuid.Parse(got)
require.NoError(t, err, "探测会话标识必须是 UUID 形态: %s", got)
}
require.Equal(t, compactProbeSessionID(7), compactProbeSessionID(7), "同账号应稳定复用同一会话")
require.NotEqual(t, compactProbeSessionID(7), compactProbeSessionID(8))
}
@@ -40,7 +40,7 @@ func TestAccountTestServiceOpenAICompactAgentIdentityUsesFreshAssertion(t *testi
upstream := &httpUpstreamRecorder{resp: &http.Response{
StatusCode: http.StatusOK,
Header: http.Header{"Content-Type": []string{"application/json"}},
Body: io.NopCloser(strings.NewReader(`{"id":"compact-agent","status":"completed"}`)),
Body: io.NopCloser(strings.NewReader(`{"id":"compact-agent","status":"completed","output":[{"type":"compaction","id":"cmp_agent_fresh","encrypted_content":"blob"}]}`)),
}}
svc := &AccountTestService{accountRepo: repo, httpUpstream: upstream}
@@ -87,7 +87,7 @@ func TestAccountTestServiceOpenAICompactAgentIdentityRecoversInvalidTaskOnce(t *
upstream := &httpUpstreamRecorder{responses: []*http.Response{
{StatusCode: http.StatusUnauthorized, Header: http.Header{"Content-Type": []string{"application/json"}}, Body: io.NopCloser(strings.NewReader(`{"error":{"code":"invalid_task_id"}}`))},
{StatusCode: http.StatusOK, Header: http.Header{"Content-Type": []string{"application/json"}}, Body: io.NopCloser(strings.NewReader(`{"id":"compact-agent","status":"completed"}`))},
{StatusCode: http.StatusOK, Header: http.Header{"Content-Type": []string{"application/json"}}, Body: io.NopCloser(strings.NewReader(`{"id":"compact-agent","status":"completed","output":[{"type":"compaction","id":"cmp_agent","encrypted_content":"blob"}]}`))},
}}
invalidator := &agentIdentityWSInvalidationRecorder{}
svc := &AccountTestService{accountRepo: repo, httpUpstream: upstream, agentIdentityWS: invalidator}
@@ -1,6 +1,114 @@
package service
import "github.com/tidwall/gjson"
import (
"net/http"
"strings"
"github.com/gin-gonic/gin"
"github.com/tidwall/gjson"
)
// openAINativeCompactionV2Key 标记本请求是原生 remote compaction v2
// (裸 /responses + stream:true + compaction_trigger),由 handler 在判定后
// 写入,供上游请求构造时补注协商头。
const openAINativeCompactionV2Key = "openai_native_compaction_v2"
const openAIRemoteCompactionV2Feature = "remote_compaction_v2"
// MarkOpenAINativeCompactionV2 由 handler 在识别出原生 v2 压缩请求时调用。
func MarkOpenAINativeCompactionV2(c *gin.Context) {
if c != nil {
c.Set(openAINativeCompactionV2Key, true)
}
}
func isOpenAINativeCompactionV2(c *gin.Context) bool {
if c == nil {
return false
}
return c.GetBool(openAINativeCompactionV2Key)
}
// ensureOpenAIRemoteCompactionV2BetaFeature 确保出站 x-codex-beta-features
// 头包含 remote_compaction_v2。真实 Codex 发送 compaction_trigger 时总会同时
// 携带该协商头(codex-rs build_model_client_beta_features_header 对该 feature
// 特判 advertise);上游或下游网关链剥掉它后,请求会在依赖该头做门控的
// 环节被降级(#5586)。这里在原生 v2 请求出站前补齐,使线型与真实 Codex
// 一致。已存在时保持原样,不重复追加。
func ensureOpenAIRemoteCompactionV2BetaFeature(h http.Header) {
if h == nil {
return
}
tokens := make([]string, 0, 4)
for _, value := range h.Values("x-codex-beta-features") {
for _, token := range strings.Split(value, ",") {
token = strings.TrimSpace(token)
if token == "" {
continue
}
if token == openAIRemoteCompactionV2Feature {
return
}
tokens = append(tokens, token)
}
}
tokens = append(tokens, openAIRemoteCompactionV2Feature)
h.Set("x-codex-beta-features", strings.Join(tokens, ","))
}
// hasOpenAICodexBetaFeaturesHeader 报告出站头里是否已存在非空的
// x-codex-beta-features(即客户端自己声明过能力集)。
func hasOpenAICodexBetaFeaturesHeader(h http.Header) bool {
if h == nil {
return false
}
for _, value := range h.Values("x-codex-beta-features") {
if strings.TrimSpace(value) != "" {
return true
}
}
return false
}
// applyOpenAICodexBetaFeatures 按真实 Codex 的会话级行为补注
// x-codex-beta-features。
//
// codex 侧规则(codex-rs:session/mod.rs build_model_client_beta_features_header
// 组装、client.rs build_responses_headers 附加):该头是**会话级常量**,挂在
// /responses、WS 握手、/responses/compact 三处的**每一个**请求上,而不是只挂
// 压缩回合。其内容是"已启用且需要 advertise 的 feature 列表",RemoteCompactionV2
// 被特判 advertise;实测默认安装下没有任何 Experimental 特性默认开启,
// 因此默认 Codex 的头值恰好就是单个 "remote_compaction_v2"。
//
// 网关据此对齐:
// - 原生 v2 压缩回合(body 带 compaction_trigger 实锤):无论账号类型都确保
// v2 在列,覆盖中间网关裁剪 token 的情形(#5586);
// - ChatGPT codex 上游(OAuth)的其余请求:客户端**未**声明该头时补成默认
// Codex 形态,消除"仅压缩回合才带该头"这种真实 Codex 不会产生的模式;
// - 客户端已声明该头:原样保留。非空但不含 v2 表示用户显式关闭了该特性,
// 网关不得替其改写能力声明;
// - 非 OAuth 上游(API Key/第三方兼容网关):不做会话级注入,只保留压缩回合
// 的那一条,避免向非 Codex 后端撒 Codex 专属头。
//
// 已知无解的歧义:用户关掉 v2 且无其他特性时,真实 Codex 同样不发该头,与"老
// 客户端"在线型上不可区分,此时按默认形态补注。该用户的 legacy 压缩端点本就
// 已被上游下线(404),不存在可回退的正确行为。
func applyOpenAICodexBetaFeatures(c *gin.Context, account *Account, h http.Header) {
if h == nil {
return
}
if isOpenAINativeCompactionV2(c) {
ensureOpenAIRemoteCompactionV2BetaFeature(h)
return
}
if account == nil || !account.IsOpenAIOAuth() {
return
}
if hasOpenAICodexBetaFeaturesHeader(h) {
return
}
h.Set("x-codex-beta-features", openAIRemoteCompactionV2Feature)
}
// HasCompactionTriggerInInput detects an input item with
// type="compaction_trigger". The handler combines this body signal with the
@@ -10,7 +10,8 @@ import (
const (
// AccountTestModeDefault drives the standard /responses connection test.
AccountTestModeDefault = "default"
// AccountTestModeCompact drives the /responses/compact compact-probe test.
// AccountTestModeCompact drives the remote-compaction probe test
// (native v2: streaming /responses with a compaction_trigger input item).
AccountTestModeCompact = "compact"
)
@@ -23,8 +24,12 @@ func normalizeAccountTestMode(mode string) string {
}
}
func createOpenAICompactProbePayload(model string) map[string]any {
return map[string]any{
// createOpenAICompactProbePayload 构造原生 remote compaction v2 探测载荷:
// 流式 /responses + input 末尾 {"type":"compaction_trigger"}。上游已下线
// legacy unary /responses/compact(v1 形态恒 404,#5598/#5624),现行 codex
// 默认协议即 v2(RemoteCompactionV2 Stable + default_enabled)。
func createOpenAICompactProbePayload(model string, isOAuth bool) map[string]any {
payload := map[string]any{
"model": strings.TrimSpace(model),
"instructions": "You are a helpful coding assistant.",
"input": []any{
@@ -33,8 +38,35 @@ func createOpenAICompactProbePayload(model string) map[string]any {
"role": "user",
"content": "Respond with OK.",
},
map[string]any{"type": "compaction_trigger"},
},
"stream": true,
}
// ChatGPT internal API 要求 store: false,与真实转发一致。
if isOAuth {
payload["store"] = false
}
return payload
}
// openAICompactProbeFoundCompactionItem 判定探测响应是否产出了 compaction
// 输出 item——v2 契约的核心(codex 缺它即 fatal "got 0 items")。三种形态都
// 认:① SSE 的 output_item.done/added(原生 v2 主形态,codex 只从这里收集
// item);② SSE 终态 response.completed 的 response.output[](部分上游只在
// 终态给出 item);③ 整体 JSON 的 output[](老网关链把请求降级成 unary)。
func openAICompactProbeFoundCompactionItem(body []byte) bool {
if len(body) == 0 {
return false
}
bodyText := string(body)
if _, found := findRawCompactionItemFromSSE(bodyText); found {
return true
}
if finalResponse, ok := extractCodexFinalResponse(bodyText); ok &&
responsesOutputHasCompactionItem(finalResponse) {
return true
}
return responsesOutputHasCompactionItem(body)
}
func shouldMarkOpenAICompactUnsupported(status int, body []byte) bool {
@@ -60,7 +92,12 @@ func shouldMarkOpenAICompactUnsupported(status int, body []byte) bool {
return false
}
func buildOpenAICompactProbeExtraUpdates(resp *http.Response, body []byte, probeErr error, now time.Time) map[string]any {
// buildOpenAICompactProbeExtraUpdates 计算探测结果的账号 extra 更新。
// compactionFound 是 v2 契约判据:HTTP 2xx 但响应无 compaction item 时同样
// 记为不支持(链路把 compaction_trigger 吞掉的形态,等价 codex 的 "got 0
// items" fatal,#5478/#5648)。极端场景(上游链只支持 legacy unary compact)
// 可用账号级 openai_compact_mode=force_on 人工覆盖。
func buildOpenAICompactProbeExtraUpdates(resp *http.Response, body []byte, probeErr error, compactionFound bool, now time.Time) map[string]any {
updates := map[string]any{
"openai_compact_checked_at": now.Format(time.RFC3339),
"openai_compact_last_status": nil,
@@ -84,10 +121,14 @@ func buildOpenAICompactProbeExtraUpdates(resp *http.Response, body []byte, probe
errMsg = "HTTP " + strconv.Itoa(resp.StatusCode)
}
errMsg = truncateString(sanitizeUpstreamErrorMessage(errMsg), 2048)
if resp.StatusCode >= 200 && resp.StatusCode < 300 {
switch {
case resp.StatusCode >= 200 && resp.StatusCode < 300 && compactionFound:
updates["openai_compact_supported"] = true
updates["openai_compact_last_error"] = ""
} else {
case resp.StatusCode >= 200 && resp.StatusCode < 300:
updates["openai_compact_supported"] = false
updates["openai_compact_last_error"] = "upstream returned 2xx without a compaction output item (native remote compaction v2 unsupported)"
default:
if shouldMarkOpenAICompactUnsupported(resp.StatusCode, body) {
updates["openai_compact_supported"] = false
}
@@ -112,9 +153,14 @@ func mergeExtraUpdates(base map[string]any, more map[string]any) map[string]any
return out
}
// compactProbeSessionID 返回探测请求使用的会话标识。真实 Codex 的
// session-id / thread-id 恒为 UUID(codex-protocol ThreadId 是 UUIDv7),
// 探测既然与真实流量走同一个 /responses 端点,标识形态就必须同构——
// 否则上游能凭 "probe_compact_5" 这类字面量一眼区分出探测流量。
// 账号级稳定派生:重复探测复用同一会话,而不是每次新开一个。
func compactProbeSessionID(accountID int64) string {
if accountID <= 0 {
return "probe_compact"
return deriveStableUUIDv4("sub2api:codex-compact-probe:v1:anonymous")
}
return "probe_compact_" + strconv.FormatInt(accountID, 10)
return deriveStableUUIDv4("sub2api:codex-compact-probe:v1:" + strconv.FormatInt(accountID, 10))
}
@@ -28,7 +28,7 @@ func TestNormalizeAccountTestMode(t *testing.T) {
func TestBuildOpenAICompactProbeExtraUpdates_SuccessMarksSupported(t *testing.T) {
now := time.Date(2026, 4, 10, 10, 0, 0, 0, time.UTC)
updates := buildOpenAICompactProbeExtraUpdates(&http.Response{StatusCode: http.StatusOK}, []byte(`{"id":"cmp_1"}`), nil, now)
updates := buildOpenAICompactProbeExtraUpdates(&http.Response{StatusCode: http.StatusOK}, []byte(`{"id":"cmp_1"}`), nil, true, now)
if got := updates["openai_compact_supported"]; got != true {
t.Fatalf("openai_compact_supported = %v, want true", got)
@@ -47,7 +47,7 @@ func TestBuildOpenAICompactProbeExtraUpdates_SuccessMarksSupported(t *testing.T)
func TestBuildOpenAICompactProbeExtraUpdates_404MarksUnsupported(t *testing.T) {
now := time.Date(2026, 4, 10, 10, 0, 0, 0, time.UTC)
body := []byte(`404 page not found`)
updates := buildOpenAICompactProbeExtraUpdates(&http.Response{StatusCode: http.StatusNotFound}, body, nil, now)
updates := buildOpenAICompactProbeExtraUpdates(&http.Response{StatusCode: http.StatusNotFound}, body, nil, false, now)
if got := updates["openai_compact_supported"]; got != false {
t.Fatalf("openai_compact_supported = %v, want false", got)
@@ -59,7 +59,7 @@ func TestBuildOpenAICompactProbeExtraUpdates_404MarksUnsupported(t *testing.T) {
func TestBuildOpenAICompactProbeExtraUpdates_502DoesNotMarkUnsupported(t *testing.T) {
now := time.Date(2026, 4, 10, 10, 0, 0, 0, time.UTC)
updates := buildOpenAICompactProbeExtraUpdates(&http.Response{StatusCode: http.StatusBadGateway}, []byte(`Upstream request failed`), nil, now)
updates := buildOpenAICompactProbeExtraUpdates(&http.Response{StatusCode: http.StatusBadGateway}, []byte(`Upstream request failed`), nil, false, now)
if _, exists := updates["openai_compact_supported"]; exists {
t.Fatalf("did not expect openai_compact_supported for 502 response")
@@ -71,7 +71,7 @@ func TestBuildOpenAICompactProbeExtraUpdates_502DoesNotMarkUnsupported(t *testin
func TestBuildOpenAICompactProbeExtraUpdates_RequestErrorDoesNotMarkUnsupported(t *testing.T) {
now := time.Date(2026, 4, 10, 10, 0, 0, 0, time.UTC)
updates := buildOpenAICompactProbeExtraUpdates(nil, nil, errors.New("dial tcp timeout"), now)
updates := buildOpenAICompactProbeExtraUpdates(nil, nil, errors.New("dial tcp timeout"), false, now)
if _, exists := updates["openai_compact_supported"]; exists {
t.Fatalf("did not expect openai_compact_supported for request error")
@@ -86,7 +86,7 @@ func TestBuildOpenAICompactProbeExtraUpdates_RequestErrorDoesNotMarkUnsupported(
func TestBuildOpenAICompactProbeExtraUpdates_NoResponseClearsLastStatus(t *testing.T) {
now := time.Date(2026, 4, 10, 10, 0, 0, 0, time.UTC)
updates := buildOpenAICompactProbeExtraUpdates(nil, nil, nil, now)
updates := buildOpenAICompactProbeExtraUpdates(nil, nil, nil, false, now)
if got, exists := updates["openai_compact_last_status"]; !exists || got != nil {
t.Fatalf("openai_compact_last_status = %v, want nil key", got)
@@ -99,7 +99,7 @@ func TestBuildOpenAICompactProbeExtraUpdates_NoResponseClearsLastStatus(t *testi
func TestBuildOpenAICompactProbeExtraUpdates_UnknownModelDoesNotMarkUnsupported(t *testing.T) {
now := time.Date(2026, 4, 10, 10, 0, 0, 0, time.UTC)
body := []byte(`{"error":{"message":"unknown model gpt-5.4-openai-compact"}}`)
updates := buildOpenAICompactProbeExtraUpdates(&http.Response{StatusCode: http.StatusBadRequest}, body, nil, now)
updates := buildOpenAICompactProbeExtraUpdates(&http.Response{StatusCode: http.StatusBadRequest}, body, nil, false, now)
if _, exists := updates["openai_compact_supported"]; exists {
t.Fatalf("did not expect openai_compact_supported for unknown-model diagnostics")
@@ -111,7 +111,7 @@ func TestBuildOpenAICompactProbeExtraUpdates_UnknownModelDoesNotMarkUnsupported(
func TestBuildOpenAICompactProbeExtraUpdates_EmptyFailureBodyFallsBackToHTTPStatus(t *testing.T) {
now := time.Date(2026, 4, 10, 10, 0, 0, 0, time.UTC)
updates := buildOpenAICompactProbeExtraUpdates(&http.Response{StatusCode: http.StatusServiceUnavailable}, nil, nil, now)
updates := buildOpenAICompactProbeExtraUpdates(&http.Response{StatusCode: http.StatusServiceUnavailable}, nil, nil, false, now)
if got := updates["openai_compact_last_status"]; got != http.StatusServiceUnavailable {
t.Fatalf("openai_compact_last_status = %v, want %d", got, http.StatusServiceUnavailable)
@@ -120,3 +120,79 @@ func TestBuildOpenAICompactProbeExtraUpdates_EmptyFailureBodyFallsBackToHTTPStat
t.Fatalf("openai_compact_last_error = %v, want HTTP 503", got)
}
}
func TestBuildOpenAICompactProbeExtraUpdates_2xxWithoutCompactionItemMarksUnsupported(t *testing.T) {
now := time.Date(2026, 4, 10, 10, 0, 0, 0, time.UTC)
updates := buildOpenAICompactProbeExtraUpdates(&http.Response{StatusCode: http.StatusOK}, []byte(`{"id":"resp_1","output":[]}`), nil, false, now)
if got := updates["openai_compact_supported"]; got != false {
t.Fatalf("openai_compact_supported = %v, want false(2xx 无 compaction item = v2 不可用)", got)
}
if got := updates["openai_compact_last_error"]; got == "" {
t.Fatalf("expected openai_compact_last_error to explain the missing compaction item")
}
}
func TestCreateOpenAICompactProbePayload_NativeV2Shape(t *testing.T) {
payload := createOpenAICompactProbePayload("gpt-5.6-sol", true)
if payload["stream"] != true {
t.Fatalf("v2 probe payload must be streaming")
}
if payload["store"] != false {
t.Fatalf("OAuth probe payload must carry store:false")
}
input, ok := payload["input"].([]any)
if !ok || len(input) != 2 {
t.Fatalf("expected 2 input items, got %v", payload["input"])
}
last, ok := input[len(input)-1].(map[string]any)
if !ok || last["type"] != "compaction_trigger" {
t.Fatalf("last input item must be compaction_trigger, got %v", input[len(input)-1])
}
apiKeyPayload := createOpenAICompactProbePayload("gpt-5.6-sol", false)
if _, has := apiKeyPayload["store"]; has {
t.Fatalf("API-key probe payload must not force store")
}
}
func TestOpenAICompactProbeFoundCompactionItem(t *testing.T) {
sseWithItem := []byte("event: response.output_item.done\ndata: {\"type\":\"response.output_item.done\",\"item\":{\"type\":\"compaction\",\"id\":\"cmp_1\",\"encrypted_content\":\"blob\"}}\n\ndata: {\"type\":\"response.completed\",\"response\":{\"id\":\"resp_1\",\"output\":[]}}\n\n")
if !openAICompactProbeFoundCompactionItem(sseWithItem) {
t.Fatalf("SSE output_item.done 携带 compaction item 应判定为支持")
}
sseAlias := []byte("data: {\"type\":\"response.output_item.done\",\"item\":{\"type\":\"compaction_summary\",\"id\":\"cmp_2\"}}\n\n")
if !openAICompactProbeFoundCompactionItem(sseAlias) {
t.Fatalf("compaction_summary 别名应判定为支持")
}
sseNoItem := []byte("data: {\"type\":\"response.output_item.done\",\"item\":{\"type\":\"message\",\"id\":\"msg_1\"}}\n\ndata: {\"type\":\"response.completed\",\"response\":{\"id\":\"resp_2\",\"output\":[]}}\n\n")
if openAICompactProbeFoundCompactionItem(sseNoItem) {
t.Fatalf("无 compaction item 的流不应判定为支持")
}
jsonWithItem := []byte(`{"id":"resp_3","output":[{"type":"compaction","id":"cmp_3"}]}`)
if !openAICompactProbeFoundCompactionItem(jsonWithItem) {
t.Fatalf("JSON output[] 携带 compaction item 应判定为支持(降级链兜底)")
}
if openAICompactProbeFoundCompactionItem(nil) {
t.Fatalf("空响应不应判定为支持")
}
}
func TestOpenAICompactProbeFoundCompactionItem_TerminalResponseOutput(t *testing.T) {
// 部分上游只在终态 response.completed 的 output[] 给出 compaction item,
// 事件流里没有 output_item.done——探针同样必须判定为支持。
sseTerminalOnly := []byte("data: {\"type\":\"response.output_item.done\",\"item\":{\"type\":\"message\",\"id\":\"msg_1\"}}\n\n" +
"data: {\"type\":\"response.completed\",\"response\":{\"id\":\"resp_t\",\"output\":[{\"type\":\"compaction\",\"id\":\"cmp_t\"}]}}\n\n")
if !openAICompactProbeFoundCompactionItem(sseTerminalOnly) {
t.Fatalf("终态 response.output 中的 compaction item 应判定为支持")
}
sseTerminalEmpty := []byte("data: {\"type\":\"response.completed\",\"response\":{\"id\":\"resp_e\",\"output\":[]}}\n\n")
if openAICompactProbeFoundCompactionItem(sseTerminalEmpty) {
t.Fatalf("终态 output 为空且无 item 事件时不应判定为支持")
}
}
@@ -99,6 +99,12 @@ func (s *OpenAIGatewayService) buildOpenAIWSHeaders(
}
}
}
// 真实 Codex 的 WS 握手同样携带会话级 x-codex-beta-features
// (client.rs build_websocket_headers 复用 build_responses_headers),
// 客户端未声明时补成默认形态,与 HTTP 出站保持一致。放在客户端头拷贝
// 之外:该头是账号/会话级属性,不依赖入站请求是否存在,也避免预热与
// 实际请求因头差异落进不同的连接池兼容分桶。
applyOpenAICodexBetaFeatures(c, account, headers)
// OAuth 账号:将 apiKeyID 混入 session 标识符,防止跨用户会话碰撞。
if account != nil && account.Type == AccountTypeOAuth {
apiKeyID := getAPIKeyIDFromContext(c)