From 000f6dc655849f03b95d1244c9c5decf4c2bb406 Mon Sep 17 00:00:00 2001 From: shaw Date: Fri, 10 Jul 2026 10:25:09 +0800 Subject: [PATCH] =?UTF-8?q?fix(compact):=20=E7=BB=88=E6=80=81=20output=20?= =?UTF-8?q?=E9=9D=9E=E7=A9=BA=E4=BD=86=E7=BC=BA=20compaction=20=E6=97=B6?= =?UTF-8?q?=E4=BB=8E=E4=BA=8B=E4=BB=B6=E6=B5=81=E8=A1=A5=E5=85=A5=EF=BC=88?= =?UTF-8?q?146=20=E7=AD=89=E4=BB=B7=E6=80=A7=E6=94=B6=E5=B0=BE=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 146 纯流式透传下 Codex 只从 raw output_item.done 收集 item,无论终态 output 形态如何都能拿到 compaction;SSE→JSON 提取链路此前仅在终态 output 为空时修补。补齐最后一处不等价:compact 请求终态 output 非空但缺 compaction、事件流中存在(done 优先 added 兜底)时以 raw JSON 追加。 门控:仅 compact 路径 + 仅缺失时触发,已含 compaction 或非 compact 请求逐字节不变。 Refs #3887 #3777 --- .../openai_compact_stream_bridge_test.go | 60 +++++++++++++++++ .../service/openai_gateway_passthrough.go | 1 + .../openai_gateway_response_handling.go | 66 +++++++++++++++++++ 3 files changed, 127 insertions(+) diff --git a/backend/internal/service/openai_compact_stream_bridge_test.go b/backend/internal/service/openai_compact_stream_bridge_test.go index 32b86380a9..bf3416192e 100644 --- a/backend/internal/service/openai_compact_stream_bridge_test.go +++ b/backend/internal/service/openai_compact_stream_bridge_test.go @@ -422,6 +422,66 @@ func TestReconstructResponseOutputFromSSE_MixedDoneAndCompactionAdded(t *testing require.Equal(t, "final", items[0].Get("encrypted_content").String()) } +// 上游不一致形态:终态 output 非空(含 message)但 compaction 只在 raw +// output_item.done 中。146 纯流式透传下 Codex 直接读事件流能拿到 compaction, +// SSE→JSON 提取必须补入等价结果。 +func TestHandleSSEToJSON_CompactSupplementsMissingCompactionIntoNonEmptyOutput(t *testing.T) { + svc := newCompactBridgeTestService() + c, rec := newCompactBridgeTestContext(t, true) + upstreamSSE := strings.Join([]string{ + `data: {"type":"response.output_item.done","output_index":0,"item":{"id":"cmp_sup","type":"compaction","encrypted_content":"supplement"}}`, + ``, + `data: {"type":"response.completed","response":{"id":"resp_sup","object":"response","status":"completed","output":[{"id":"msg_sup","type":"message","role":"assistant","content":[{"type":"output_text","text":"note"}]}],"usage":{"input_tokens":2,"output_tokens":1,"total_tokens":3}}}`, + ``, + }, "\n") + resp := &http.Response{ + StatusCode: http.StatusOK, + Header: http.Header{"Content-Type": []string{"text/event-stream"}}, + Body: io.NopCloser(strings.NewReader(upstreamSSE)), + } + + result, err := svc.handleNonStreamingResponse(context.Background(), resp, c, &Account{ID: 1, Type: AccountTypeOAuth}, "gpt-5.5", "gpt-5.5") + require.NoError(t, err) + require.NotNil(t, result) + + events := parseCompactBridgeSSE(t, rec.Body.String()) + require.Len(t, events, 3) + itemTypes := []string{ + gjson.Get(events[0][1], "item.type").String(), + gjson.Get(events[1][1], "item.type").String(), + } + require.Contains(t, itemTypes, "compaction") + require.Contains(t, itemTypes, "message") + require.Equal(t, "response.completed", events[2][0]) + require.Len(t, gjson.Get(events[2][1], "response.output").Array(), 2) +} + +// 补全逻辑的门控:非 compact 请求原样返回;终态已含 compaction 不重复补入。 +func TestSupplementCompactionItemFromSSE_Gating(t *testing.T) { + bodyText := `data: {"type":"response.output_item.done","item":{"id":"cmp_g","type":"compaction","encrypted_content":"g"}}` + "\n" + + // 非 compact 路径:不补入。 + gin.SetMode(gin.TestMode) + rec := httptest.NewRecorder() + c, _ := gin.CreateTestContext(rec) + c.Request = httptest.NewRequest(http.MethodPost, "/v1/responses", nil) + finalResponse := []byte(`{"id":"r1","output":[{"type":"message"}]}`) + require.Equal(t, string(finalResponse), string(supplementCompactionItemFromSSE(c, finalResponse, bodyText))) + + // compact 路径 + 终态已含 compaction:不重复补入。 + c2, _ := newCompactBridgeTestContext(t, false) + already := []byte(`{"id":"r2","output":[{"type":"compaction","encrypted_content":"x"}]}`) + require.Equal(t, string(already), string(supplementCompactionItemFromSSE(c2, already, bodyText))) + + // compact 路径 + 终态非空缺 compaction:补入到末尾。 + missing := []byte(`{"id":"r3","output":[{"type":"message"}]}`) + patched := supplementCompactionItemFromSSE(c2, missing, bodyText) + items := gjson.GetBytes(patched, "output").Array() + require.Len(t, items, 2) + require.Equal(t, "compaction", items[1].Get("type").String()) + require.Equal(t, "g", items[1].Get("encrypted_content").String()) +} + // 非 compaction 的 output_item.added 不参与回退收集(added 阶段的 message // 通常是空壳),仍走 delta 重建。 func TestReconstructResponseOutputFromSSE_NonCompactionAddedStillUsesDeltas(t *testing.T) { diff --git a/backend/internal/service/openai_gateway_passthrough.go b/backend/internal/service/openai_gateway_passthrough.go index f4a456d7fc..12eb248fcf 100644 --- a/backend/internal/service/openai_gateway_passthrough.go +++ b/backend/internal/service/openai_gateway_passthrough.go @@ -1127,6 +1127,7 @@ func (s *OpenAIGatewayService) handlePassthroughSSEToJSON(resp *http.Response, c } } } + finalResponse = supplementCompactionItemFromSSE(c, finalResponse, bodyText) body = finalResponse if originalModel != "" && mappedModel != "" && originalModel != mappedModel { body = s.replaceModelInResponseBody(body, mappedModel, originalModel) diff --git a/backend/internal/service/openai_gateway_response_handling.go b/backend/internal/service/openai_gateway_response_handling.go index 9c9ea0c1a3..45c948f56f 100644 --- a/backend/internal/service/openai_gateway_response_handling.go +++ b/backend/internal/service/openai_gateway_response_handling.go @@ -878,6 +878,7 @@ func (s *OpenAIGatewayService) handleSSEToJSON(resp *http.Response, c *gin.Conte } } } + finalResponse = supplementCompactionItemFromSSE(c, finalResponse, bodyText) body = finalResponse if originalModel != mappedModel { body = s.replaceModelInResponseBody(body, mappedModel, originalModel) @@ -1156,6 +1157,71 @@ func isResponsesCompactionItemType(itemType string) bool { } } +// supplementCompactionItemFromSSE 保证 compact 请求的终态 output 携带 +// compaction item:终态 output 非空但缺失 compaction、而原始事件流的 +// output_item.done(或 added)中存在时(上游不一致形态),以 raw JSON 补入。 +// Codex remote compact v2 只从 output_item.done 收集 item 且要求恰好一个 +// compaction item——纯流式透传(v0.1.146)下客户端直接读事件流天然拿得到, +// SSE→JSON 提取链路必须给出等价结果。非 compact 请求原样返回。 +func supplementCompactionItemFromSSE(c *gin.Context, finalResponse []byte, bodyText string) []byte { + if !isOpenAIResponsesCompactPath(c) { + return finalResponse + } + if len(gjson.GetBytes(finalResponse, "output").Array()) == 0 { + // 空 output 由 reconstructResponseOutputFromSSE 整体修补,不在此处理。 + return finalResponse + } + if responsesOutputHasCompactionItem(finalResponse) { + return finalResponse + } + item, found := findRawCompactionItemFromSSE(bodyText) + if !found { + return finalResponse + } + patched, err := sjson.SetRawBytes(finalResponse, "output.-1", item) + if err != nil { + return finalResponse + } + return patched +} + +// responsesOutputHasCompactionItem reports whether the response JSON already +// carries a compaction item in its output array. +func responsesOutputHasCompactionItem(response []byte) bool { + for _, item := range gjson.GetBytes(response, "output").Array() { + if isResponsesCompactionItemType(item.Get("type").String()) { + return true + } + } + return false +} + +// findRawCompactionItemFromSSE 从原始 SSE 事件流中提取第一个 compaction 类 +// item 的 raw JSON:output_item.done 优先,output_item.added 兜底。 +func findRawCompactionItemFromSSE(bodyText string) (json.RawMessage, bool) { + var found json.RawMessage + pick := func(eventType string) { + forEachOpenAISSEDataPayload(bodyText, func(data []byte) { + if found != nil { + return + } + if strings.TrimSpace(gjson.GetBytes(data, "type").String()) != eventType { + return + } + item := gjson.GetBytes(data, "item") + if !item.IsObject() || !isResponsesCompactionItemType(item.Get("type").String()) { + return + } + found = json.RawMessage(item.Raw) + }) + } + pick("response.output_item.done") + if found == nil { + pick("response.output_item.added") + } + return found, found != nil +} + // reconstructResponseOutputFromSSE scans raw SSE body text and returns a // JSON-encoded output array for a terminal event whose output is empty. // Raw output_item.done items are preferred: per the Responses protocol they