mirror of
https://github.com/Wei-Shaw/sub2api.git
synced 2026-10-07 07:13:44 +08:00
Merge pull request #5415 from Yuxin-Qiao/fix/openai-responses-empty-completed-failover
fix(openai): fail over empty response.completed streams instead of recording 0/0 success
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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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))
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -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,
|
||||
|
||||
@@ -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