mirror of
https://github.com/Wei-Shaw/sub2api.git
synced 2026-09-24 16:05:44 +08:00
fix(compact): 终态 output 非空但缺 compaction 时从事件流补入(146 等价性收尾)
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
This commit is contained in:
@@ -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) {
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user