mirror of
https://github.com/Wei-Shaw/sub2api.git
synced 2026-10-07 13:18:18 +08:00
fix(openai): fail over empty response.completed streams instead of recording 0/0 success
A Responses stream that ends with response.completed but carries no output, no usage and no error (only response.created + response.completed observed) is a silent upstream refusal. It was previously treated as a successful turn, recording 0 input / 0 output tokens without triggering failover, so clients saw an empty reply and usage was never billed. Detect the empty terminal event in both the standard and passthrough streaming paths while nothing has been written to the client, and return an UpstreamFailoverError so the next account is tried; when failover is exhausted the client gets an explicit Responses-format error instead of an empty success. Streams with any semantic output, usage or error are untouched. Fixes #5009.
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -216,6 +216,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"))
|
||||
@@ -524,6 +525,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 {
|
||||
@@ -995,6 +1011,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
|
||||
|
||||
@@ -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))
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -1531,7 +1531,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,
|
||||
|
||||
@@ -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{
|
||||
|
||||
Reference in New Issue
Block a user