diff --git a/backend/internal/service/openai_gateway_passthrough.go b/backend/internal/service/openai_gateway_passthrough.go index 44c9670a3e..030932befb 100644 --- a/backend/internal/service/openai_gateway_passthrough.go +++ b/backend/internal/service/openai_gateway_passthrough.go @@ -1160,6 +1160,7 @@ func (s *OpenAIGatewayService) handleStreamingResponsePassthrough( sawDone := false sawTerminalEvent := false sawFailedEvent := false + semanticOutputSeen := false failedMessage := "" clientOutputStarted := false upstreamRequestID := strings.TrimSpace(resp.Header.Get("x-request-id")) @@ -1305,6 +1306,18 @@ func (s *OpenAIGatewayService) handleStreamingResponsePassthrough( line = "data: " + string(sanitizedData) } lineStartsClientOutput = forceFlushFailedEvent || openAIStreamDataStartsClientOutput(trimmedData, eventType) + if lineStartsClientOutput && trimmedData != "[DONE]" && !openAIStreamEventTypeIsTerminal(eventType) { + semanticOutputSeen = true + } + // OpenAI Responses streams that terminate with an empty + // response.completed (no output, no usage, no error, nothing sent + // to the client) are silent upstream refusals: fail over instead of + // recording a successful 0/0 usage turn (issue #5009). + if (eventType == "response.completed" || eventType == "response.done") && + !sawFailedEvent && !semanticOutputSeen && !clientOutputStarted && + openAIResponsesCompletedEventIsEmpty(dataBytes, usage) { + return resultWithUsage(), newOpenAIResponsesEmptyCompletedFailoverError(c, account, upstreamRequestID) + } if firstTokenMs == nil && lineStartsClientOutput && trimmedData != "[DONE]" { ms := int(time.Since(startTime).Milliseconds()) firstTokenMs = &ms diff --git a/backend/internal/service/openai_gateway_response_handling.go b/backend/internal/service/openai_gateway_response_handling.go index cda23fa07b..d8a17b19f4 100644 --- a/backend/internal/service/openai_gateway_response_handling.go +++ b/backend/internal/service/openai_gateway_response_handling.go @@ -228,6 +228,7 @@ func (s *OpenAIGatewayService) handleStreamingResponseWithReasoning(ctx context. clientDisconnected := false // 客户端断开后继续 drain 上游以收集 usage sawTerminalEvent := false sawFailedEvent := false + responsesSemanticOutputSeen := false failedMessage := "" clientOutputStarted := false upstreamRequestID := strings.TrimSpace(resp.Header.Get("x-request-id")) @@ -542,6 +543,21 @@ func (s *OpenAIGatewayService) handleStreamingResponseWithReasoning(ctx context. if guardFirstOutput { eventStartsClientOutput = eventStartsClientOutput || startsClientOutput } + if startsClientOutput && !openAIStreamEventTypeIsTerminal(eventType) { + responsesSemanticOutputSeen = true + } + // OpenAI Responses streams that terminate with an empty + // response.completed (no output, no usage, no error, nothing sent + // to the client) are silent upstream refusals: fail over instead of + // recording a successful 0/0 usage turn (issue #5009). + if account != nil && account.Platform == PlatformOpenAI && + (eventType == "response.completed" || eventType == "response.done") && + !sawFailedEvent && !responsesSemanticOutputSeen && !clientOutputStarted && + openAIResponsesCompletedEventIsEmpty(dataBytes, usage) { + sawTerminalEvent = true + streamEarlyErr = newOpenAIResponsesEmptyCompletedFailoverError(c, account, upstreamRequestID) + return + } // 写入客户端(客户端断开后继续 drain 上游) if !clientDisconnected { @@ -1023,6 +1039,32 @@ func extractOpenAIUsageFromJSONBytes(body []byte) (OpenAIUsage, bool) { return OpenAIUsage{}, false } +// openAIResponsesCompletedEventIsEmpty reports whether a response.completed / +// response.done SSE payload carries no usage, no error and no output items. +// The accumulated usage is consulted too, because OpenAI may deliver usage on +// an earlier event. An empty terminal event after a stream with no semantic +// output is treated as a silent upstream refusal (issue #5009). +func openAIResponsesCompletedEventIsEmpty(data []byte, usage *OpenAIUsage) bool { + if len(data) == 0 || !gjson.ValidBytes(data) { + return false + } + if usage != nil && (usage.InputTokens > 0 || usage.OutputTokens > 0 || + usage.ImageInputTokens > 0 || usage.ImageOutputTokens > 0 || + usage.CacheCreationInputTokens > 0 || usage.CacheReadInputTokens > 0) { + return false + } + if gjson.GetBytes(data, "usage").Exists() || gjson.GetBytes(data, "response.usage").Exists() { + return false + } + if gjson.GetBytes(data, "error").Exists() || gjson.GetBytes(data, "response.error").Exists() { + return false + } + if output := gjson.GetBytes(data, "response.output"); output.Exists() && output.IsArray() && len(output.Array()) > 0 { + return false + } + return true +} + func mergeHostedImageGenToolUsage(imageGen gjson.Result, usage *OpenAIUsage) { if !imageGen.Exists() || !imageGen.IsObject() { return diff --git a/backend/internal/service/openai_gateway_responses_empty_completed_test.go b/backend/internal/service/openai_gateway_responses_empty_completed_test.go new file mode 100644 index 0000000000..55bfda089c --- /dev/null +++ b/backend/internal/service/openai_gateway_responses_empty_completed_test.go @@ -0,0 +1,168 @@ +//go:build unit + +package service + +import ( + "context" + "errors" + "io" + "net/http" + "strings" + "testing" + + "github.com/gin-gonic/gin" + "github.com/stretchr/testify/require" +) + +// TestOpenAIResponsesEmptyCompletedFailsOver verifies that a Responses stream +// ending with an empty response.completed (no output, no usage, no error) is +// turned into a failover error instead of a successful empty reply (issue +// #5009). +func TestOpenAIResponsesEmptyCompletedFailsOver(t *testing.T) { + gin.SetMode(gin.TestMode) + + upstream := &httpUpstreamRecorder{resp: &http.Response{ + StatusCode: http.StatusOK, + Header: http.Header{"Content-Type": []string{"text/event-stream"}}, + Body: io.NopCloser(strings.NewReader( + "data: {\"type\":\"response.created\",\"response\":{\"id\":\"resp_empty\",\"object\":\"response\",\"status\":\"in_progress\"}}\n\n" + + "data: {\"type\":\"response.completed\",\"response\":{\"id\":\"resp_empty\",\"object\":\"response\",\"status\":\"completed\"}}\n\n", + )), + }} + svc := newOpenAIImageGenerationControlTestService(upstream) + c, recorder := newOpenAIImageGenerationControlTestContext(true, "codex_cli_rs/0.144.1") + account := newOpenAIImageGenerationControlTestAccount() + account.Extra = map[string]any{"openai_passthrough": true} + + body := []byte(`{ + "model":"gpt-5.6-sol", + "stream":true, + "input":[{"type":"message","role":"user","content":[{"type":"input_text","text":"continue"}]}] + }`) + + _, err := svc.Forward(context.Background(), c, account, body) + require.Error(t, err) + var failoverErr *UpstreamFailoverError + require.True(t, errors.As(err, &failoverErr), "empty completed must produce UpstreamFailoverError, got: %v", err) + require.Equal(t, http.StatusBadGateway, failoverErr.StatusCode) + require.Empty(t, recorder.Body.String(), "no empty success stream may reach the client") +} + +// TestOpenAIResponsesEmptyCompletedWithOutputSucceeds ensures streams with real +// semantic output are untouched. +func TestOpenAIResponsesEmptyCompletedWithOutputSucceeds(t *testing.T) { + gin.SetMode(gin.TestMode) + + upstream := &httpUpstreamRecorder{resp: &http.Response{ + StatusCode: http.StatusOK, + Header: http.Header{"Content-Type": []string{"text/event-stream"}}, + Body: io.NopCloser(strings.NewReader( + "data: {\"type\":\"response.created\",\"response\":{\"id\":\"resp_ok\",\"object\":\"response\",\"status\":\"in_progress\"}}\n\n" + + "data: {\"type\":\"response.output_text.delta\",\"delta\":\"hello\"}\n\n" + + "data: {\"type\":\"response.completed\",\"response\":{\"id\":\"resp_ok\",\"object\":\"response\",\"status\":\"completed\",\"usage\":{\"input_tokens\":10,\"output_tokens\":5,\"total_tokens\":15}}}\n\n", + )), + }} + svc := newOpenAIImageGenerationControlTestService(upstream) + c, recorder := newOpenAIImageGenerationControlTestContext(true, "codex_cli_rs/0.144.1") + account := newOpenAIImageGenerationControlTestAccount() + account.Extra = map[string]any{"openai_passthrough": true} + + body := []byte(`{ + "model":"gpt-5.6-sol", + "stream":true, + "input":[{"type":"message","role":"user","content":[{"type":"input_text","text":"continue"}]}] + }`) + + result, err := svc.Forward(context.Background(), c, account, body) + require.NoError(t, err) + require.NotNil(t, result) + require.Contains(t, recorder.Body.String(), "hello") + require.NotNil(t, result.Usage) + require.Equal(t, 10, result.Usage.InputTokens) + require.Equal(t, 5, result.Usage.OutputTokens) +} + +// TestOpenAIResponsesEmptyCompletedWithUsageSucceeds ensures a completed event +// carrying usage is not mistaken for a silent refusal even without output. +func TestOpenAIResponsesEmptyCompletedWithUsageSucceeds(t *testing.T) { + gin.SetMode(gin.TestMode) + + upstream := &httpUpstreamRecorder{resp: &http.Response{ + StatusCode: http.StatusOK, + Header: http.Header{"Content-Type": []string{"text/event-stream"}}, + Body: io.NopCloser(strings.NewReader( + "data: {\"type\":\"response.created\",\"response\":{\"id\":\"resp_usage\",\"object\":\"response\",\"status\":\"in_progress\"}}\n\n" + + "data: {\"type\":\"response.completed\",\"response\":{\"id\":\"resp_usage\",\"object\":\"response\",\"status\":\"completed\",\"usage\":{\"input_tokens\":3,\"output_tokens\":0,\"total_tokens\":3}}}\n\n", + )), + }} + svc := newOpenAIImageGenerationControlTestService(upstream) + c, _ := newOpenAIImageGenerationControlTestContext(true, "codex_cli_rs/0.144.1") + account := newOpenAIImageGenerationControlTestAccount() + account.Extra = map[string]any{"openai_passthrough": true} + + body := []byte(`{ + "model":"gpt-5.6-sol", + "stream":true, + "input":[{"type":"message","role":"user","content":[{"type":"input_text","text":"continue"}]}] + }`) + + result, err := svc.Forward(context.Background(), c, account, body) + require.NoError(t, err) + require.NotNil(t, result) + require.NotNil(t, result.Usage) + require.Equal(t, 3, result.Usage.InputTokens) +} + +func TestOpenAIResponsesCompletedEventIsEmpty(t *testing.T) { + cases := []struct { + name string + data string + usage *OpenAIUsage + want bool + }{ + { + name: "bare completed", + data: `{"type":"response.completed"}`, + want: true, + }, + { + name: "completed with empty output array", + data: `{"type":"response.completed","response":{"id":"r1","status":"completed","output":[]}}`, + want: true, + }, + { + name: "completed with usage", + data: `{"type":"response.completed","response":{"id":"r1","status":"completed","usage":{"input_tokens":1,"output_tokens":1}}}`, + want: false, + }, + { + name: "completed with error", + data: `{"type":"response.completed","response":{"id":"r1","status":"completed","error":{"code":"x"}}}`, + want: false, + }, + { + name: "completed with output item", + data: `{"type":"response.completed","response":{"id":"r1","status":"completed","output":[{"type":"message","id":"msg_1"}]}}`, + want: false, + }, + { + name: "accumulated usage", + data: `{"type":"response.completed"}`, + usage: &OpenAIUsage{ + InputTokens: 7, + OutputTokens: 2, + }, + want: false, + }, + { + name: "invalid json", + data: `{"type":`, + want: false, + }, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + require.Equal(t, tc.want, openAIResponsesCompletedEventIsEmpty([]byte(tc.data), tc.usage)) + }) + } +} diff --git a/backend/internal/service/openai_gateway_service_test.go b/backend/internal/service/openai_gateway_service_test.go index b3cb5fbd97..2c5bdba16b 100644 --- a/backend/internal/service/openai_gateway_service_test.go +++ b/backend/internal/service/openai_gateway_service_test.go @@ -1568,7 +1568,7 @@ func TestOpenAIStreamingTerminalAndClientCancellationDoNotQuarantineProxy(t *tes Body: &openAIStreamReadThenErrorCloser{ reader: strings.NewReader(strings.Join([]string{ "event: response.completed", - `data: {"type":"response.completed","response":{"status":"completed","output":[]}}`, + `data: {"type":"response.completed","response":{"status":"completed","output":[],"usage":{"input_tokens":5,"output_tokens":3,"total_tokens":8}}}`, "", }, "\n")), err: io.ErrUnexpectedEOF, diff --git a/backend/internal/service/openai_silent_refusal.go b/backend/internal/service/openai_silent_refusal.go index 27b71b7571..017c4c1f6d 100644 --- a/backend/internal/service/openai_silent_refusal.go +++ b/backend/internal/service/openai_silent_refusal.go @@ -16,6 +16,7 @@ const ( openAISilentRefusalErrorCode = "openai_silent_refusal" openAISilentRefusalUpstreamMessage = "OpenAI upstream returned an empty completion stream with finish_reason=stop and no usage" openAISilentRefusalClientMessage = "Upstream returned an empty completion without usage; no fallback account was available" + openAIResponsesEmptyCompletedMessage = "OpenAI upstream returned an empty response.completed stream with no output and no usage" ) type openAIChatSilentRefusalDetector struct { @@ -266,6 +267,43 @@ func newOpenAISilentRefusalFailoverError(c *gin.Context, account *Account, upstr } } +// newOpenAIResponsesEmptyCompletedFailoverError marks an empty +// response.completed terminal event as a retryable upstream anomaly. OpenAI +// Responses streams that deliver only response.created + response.completed +// with no output, no usage and no error are treated as silent upstream +// refusals rather than successful empty replies (issue #5009). +func newOpenAIResponsesEmptyCompletedFailoverError(c *gin.Context, account *Account, upstreamRequestID string) *UpstreamFailoverError { + accountID := int64(0) + accountName := "" + platform := PlatformOpenAI + if account != nil { + accountID = account.ID + accountName = account.Name + platform = account.Platform + } + + setOpsUpstreamError(c, http.StatusBadGateway, openAIResponsesEmptyCompletedMessage, "") + appendOpsUpstreamError(c, OpsUpstreamErrorEvent{ + Platform: platform, + AccountID: accountID, + AccountName: accountName, + UpstreamStatusCode: http.StatusBadGateway, + UpstreamRequestID: upstreamRequestID, + Kind: "failover", + Message: openAIResponsesEmptyCompletedMessage, + }) + + headers := http.Header{} + if strings.TrimSpace(upstreamRequestID) != "" { + headers.Set("x-request-id", strings.TrimSpace(upstreamRequestID)) + } + return &UpstreamFailoverError{ + StatusCode: http.StatusBadGateway, + ResponseBody: openAISilentRefusalErrorBody(), + ResponseHeaders: headers, + } +} + func openAISilentRefusalErrorBody() []byte { body, err := json.Marshal(map[string]any{ "error": map[string]any{