fix: 保留 remote_compaction_v2 原生 Responses 链路

This commit is contained in:
Tian Lee
2026-07-11 12:51:53 +08:00
parent e316ebf528
commit 84bb7d0709
10 changed files with 505 additions and 83 deletions
@@ -23,46 +23,61 @@ func newCompactBodySignalTestContext(t *testing.T, path string, body []byte) *gi
return c
}
// body-signal 提升后必须与 path-based compact 走同一条链路:
// path 改写、requireCompact 判定、stream/store/prompt_cache_key 归一化删除。
// 回归防护:若 stream 字段存活,Forward 会用流式 handler 解析 compact 的
// JSON 响应,导致 "stream ended before a terminal event" 的换号 failover 风暴。
func TestNormalizeOpenAIResponsesCompactRequest_BodySignalPromoted(t *testing.T) {
func TestNormalizeOpenAIResponsesCompactRequest_RemoteV2StaysOnResponses(t *testing.T) {
h := &OpenAIGatewayHandler{}
body := []byte(`{
"model":"gpt-5.5",
"model":"gpt-5.6-sol",
"stream":true,
"store":true,
"prompt_cache_key":"pck-signal-1",
"reasoning":{"effort":"max","context":"all_turns"},
"input":[
{"type":"message","role":"user","content":"hello"},
{"type":"compaction_trigger"}
]
}`)
c := newCompactBodySignalTestContext(t, "/v1/responses", body)
c.Request.Header.Set("x-codex-beta-features", "responses_websockets_v2, remote_compaction_v2, another_feature")
normalized, ok := h.normalizeOpenAIResponsesCompactRequest(c, zap.NewNop(), body)
require.True(t, ok)
require.Equal(t, "/v1/responses/compact", c.Request.URL.Path)
require.True(t, isOpenAIRemoteCompactPath(c))
require.False(t, gjson.GetBytes(normalized, "stream").Exists())
require.False(t, gjson.GetBytes(normalized, "store").Exists())
require.False(t, gjson.GetBytes(normalized, "prompt_cache_key").Exists())
require.Equal(t, "gpt-5.5", gjson.GetBytes(normalized, "model").String())
require.True(t, gjson.GetBytes(normalized, "input").IsArray())
require.Equal(t, "/v1/responses", c.Request.URL.Path)
require.False(t, isOpenAIRemoteCompactPath(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())
reqStream, streamOK := parseOpenAICompatibleStream(normalized)
require.True(t, streamOK)
require.False(t, reqStream)
require.True(t, reqStream)
seed, exists := c.Get(service.OpenAICompactSessionSeedKeyForTest())
require.True(t, exists)
require.Equal(t, "pck-signal-1", seed)
_, seedExists := c.Get(service.OpenAICompactSessionSeedKeyForTest())
require.False(t, seedExists)
_, streamMarkerExists := c.Get(service.OpenAICompactClientStreamKeyForTest())
require.False(t, streamMarkerExists)
}
func TestNormalizeOpenAIResponsesCompactRequest_BodySignalTrailingSlash(t *testing.T) {
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)
c.Request.Header.Set("x-codex-beta-features", "remote_compaction_v2")
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 TestNormalizeOpenAIResponsesCompactRequest_BodySignalTrailingSlashPromoted(t *testing.T) {
h := &OpenAIGatewayHandler{}
body := []byte(`{"model":"gpt-5.5","input":[{"type":"compaction_trigger"}]}`)
c := newCompactBodySignalTestContext(t, "/v1/responses/", body)
@@ -82,6 +97,64 @@ func TestNormalizeOpenAIResponsesCompactRequest_CodexDirectAliasPromoted(t *test
require.Equal(t, "/backend-api/codex/responses/compact", c.Request.URL.Path)
}
func TestNormalizeOpenAIResponsesCompactRequest_NonRemoteV2BodySignalPromoted(t *testing.T) {
h := &OpenAIGatewayHandler{}
tests := []struct {
name string
body []byte
betaHeader string
wantMarked bool
}{
{
name: "no_header",
body: []byte(`{"model":"gpt-5.5","stream":true,"input":[{"type":"compaction_trigger"}]}`),
wantMarked: true,
},
{
name: "unrelated_header",
body: []byte(`{"model":"gpt-5.5","stream":true,"input":[{"type":"compaction_trigger"}]}`),
betaHeader: "responses_websockets_v2",
wantMarked: true,
},
{
name: "wrong_case_header",
body: []byte(`{"model":"gpt-5.5","stream":true,"input":[{"type":"compaction_trigger"}]}`),
betaHeader: "REMOTE_COMPACTION_V2",
wantMarked: true,
},
{
name: "stream_false",
body: []byte(`{"model":"gpt-5.5","stream":false,"input":[{"type":"compaction_trigger"}]}`),
betaHeader: "remote_compaction_v2",
},
{
name: "stream_absent",
body: []byte(`{"model":"gpt-5.5","input":[{"type":"compaction_trigger"}]}`),
betaHeader: "remote_compaction_v2",
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
c := newCompactBodySignalTestContext(t, "/v1/responses", tt.body)
if tt.betaHeader != "" {
c.Request.Header.Set("x-codex-beta-features", tt.betaHeader)
}
normalized, ok := h.normalizeOpenAIResponsesCompactRequest(c, zap.NewNop(), tt.body)
require.True(t, ok)
require.Equal(t, "/v1/responses/compact", c.Request.URL.Path)
require.False(t, gjson.GetBytes(normalized, "stream").Exists())
marked, exists := c.Get(service.OpenAICompactClientStreamKeyForTest())
require.Equal(t, tt.wantMarked, exists)
if tt.wantMarked {
require.Equal(t, true, marked)
}
})
}
}
func TestNormalizeOpenAIResponsesCompactRequest_NoTriggerUntouched(t *testing.T) {
h := &OpenAIGatewayHandler{}
body := []byte(`{"model":"gpt-5.5","stream":true,"input":[{"type":"message","role":"user","content":"hello"}]}`)
@@ -99,6 +172,7 @@ func TestNormalizeOpenAIResponsesCompactRequest_PathBasedNoDoubleSuffix(t *testi
h := &OpenAIGatewayHandler{}
body := []byte(`{"model":"gpt-5.5","stream":true,"store":true,"input":[{"type":"message","role":"user","content":"hello"}]}`)
c := newCompactBodySignalTestContext(t, "/v1/responses/compact", body)
c.Request.Header.Set("x-codex-beta-features", "remote_compaction_v2")
normalized, ok := h.normalizeOpenAIResponsesCompactRequest(c, zap.NewNop(), body)
require.True(t, ok)
@@ -118,36 +192,6 @@ func TestNormalizeOpenAIResponsesCompactRequest_SubpathNotPromoted(t *testing.T)
require.Equal(t, body, normalized)
}
// 回归 #3875:body-signal 原始请求 stream:true 时必须标记 client-stream,
// 供响应写回阶段把上游 unary JSON 合成回 Codex remote compact v2 所需的 SSE。
func TestNormalizeOpenAIResponsesCompactRequest_BodySignalStreamTrueMarksClientStream(t *testing.T) {
h := &OpenAIGatewayHandler{}
body := []byte(`{"model":"gpt-5.5","stream":true,"input":[{"type":"compaction_trigger"}]}`)
c := newCompactBodySignalTestContext(t, "/v1/responses", body)
_, ok := h.normalizeOpenAIResponsesCompactRequest(c, zap.NewNop(), body)
require.True(t, ok)
marked, exists := c.Get(service.OpenAICompactClientStreamKeyForTest())
require.True(t, exists)
require.Equal(t, true, marked)
}
func TestNormalizeOpenAIResponsesCompactRequest_BodySignalStreamFalseNotMarked(t *testing.T) {
h := &OpenAIGatewayHandler{}
for name, body := range map[string][]byte{
"stream_false": []byte(`{"model":"gpt-5.5","stream":false,"input":[{"type":"compaction_trigger"}]}`),
"stream_absent": []byte(`{"model":"gpt-5.5","input":[{"type":"compaction_trigger"}]}`),
} {
c := newCompactBodySignalTestContext(t, "/v1/responses", body)
_, ok := h.normalizeOpenAIResponsesCompactRequest(c, zap.NewNop(), body)
require.True(t, ok, name)
require.Equal(t, "/v1/responses/compact", c.Request.URL.Path, name)
_, exists := c.Get(service.OpenAICompactClientStreamKeyForTest())
require.False(t, exists, "case %s 不应标记 client-stream", name)
}
}
// path-based compact(Codex v1 unary 协议)即使 body 带 stream:true 也不标记,
// 保持 JSON 写回行为不变。
func TestNormalizeOpenAIResponsesCompactRequest_PathBasedStreamTrueNotMarked(t *testing.T) {
@@ -580,21 +580,33 @@ func isBareOpenAIResponsesPath(c *gin.Context) bool {
return strings.HasSuffix(normalizedPath, "/responses")
}
// normalizeOpenAIResponsesCompactRequest 统一处理两种入站 compact 形态:
// path-based(POST /v1/responses/compact)与 Codex remote compact v2 的
// body-signal(普通 POST /v1/responses 的 input 中携带 type=compaction_trigger,
// 见 #3777)。body-signal 命中时在 stream 解析、compact body 归一化与
// requireCompact 调度判定之前改写 URL path,使后续全部链路(含 passthrough
// 分支与上游 URL 构建)与 path-based 完全一致。
func isOpenAIRemoteCompactionV2Request(c *gin.Context, body []byte) bool {
stream, valid := parseOpenAICompatibleStream(body)
if !valid || !stream || c == nil || c.Request == nil {
return false
}
for _, header := range c.Request.Header.Values("x-codex-beta-features") {
for _, feature := range strings.Split(header, ",") {
if strings.TrimSpace(feature) == "remote_compaction_v2" {
return true
}
}
}
return false
}
// normalizeOpenAIResponsesCompactRequest keeps Codex remote compaction v2 on
// its native streaming /responses wire and preserves the legacy body-signal
// promotion for clients that do not explicitly advertise that protocol.
// 返回归一化后的 body;ok=false 表示错误响应已写出,调用方应直接 return。
func (h *OpenAIGatewayHandler) normalizeOpenAIResponsesCompactRequest(c *gin.Context, reqLog *zap.Logger, body []byte) ([]byte, bool) {
isCompactRequest := service.IsOpenAIResponsesCompactPathForTest(c)
if !isCompactRequest && isBareOpenAIResponsesPath(c) && service.HasCompactionTriggerInInput(body) {
if isOpenAIRemoteCompactionV2Request(c, body) {
return body, true
}
c.Request.URL.Path = strings.TrimRight(c.Request.URL.Path, "/") + "/compact"
isCompactRequest = true
// Codex remote compact v2 的原始请求是流式 /responses:白名单归一化会删除
// stream 并让上游走 unary JSON,但客户端仍按 SSE 消费响应。记录原始
// stream 意图,响应写回阶段据此把 JSON 合成回 SSE(#3875)。
clientStream := gjson.GetBytes(body, "stream").Bool()
if clientStream {
service.MarkOpenAICompactClientStream(c)