diff --git a/backend/internal/service/openai_gateway_messages.go b/backend/internal/service/openai_gateway_messages.go index abf80de75f..f37f2997cd 100644 --- a/backend/internal/service/openai_gateway_messages.go +++ b/backend/internal/service/openai_gateway_messages.go @@ -793,7 +793,9 @@ func (s *OpenAIGatewayService) handleAnthropicStreamingResponse( return false } - isTerminalEvent := isOpenAICompatResponsesTerminalEvent(event.Type) + eventType := strings.TrimSpace(event.Type) + isBareErrorEvent := eventType == "error" + isTerminalEvent := isOpenAICompatResponsesTerminalEvent(eventType) || isBareErrorEvent if isTerminalEvent { if event.Response != nil { if id := strings.TrimSpace(event.Response.ID); id != "" { @@ -808,7 +810,7 @@ func (s *OpenAIGatewayService) handleAnthropicStreamingResponse( } // cyber_policy 致命不可重试:标记供 handler 事后记录;以 Anthropic SSE error 事件 // 回写让客户端感知并停止重试(F4),丢弃后续转换输出。 - if strings.TrimSpace(event.Type) == "response.failed" { + if eventType == "response.failed" || isBareErrorEvent { payloadBytes := []byte(payload) if hit, code, msg := detectOpenAICyberPolicy(payloadBytes); hit { MarkOpsCyberPolicy(c, CyberPolicyMark{ @@ -833,7 +835,10 @@ func (s *OpenAIGatewayService) handleAnthropicStreamingResponse( return true } message := extractOpenAISSEErrorMessage(payloadBytes) - if openAIStreamFailedEventShouldFailover(payloadBytes, message) { + // Once Anthropic output has started, switching accounts would splice + // two model streams together. Surface a proper Anthropic error event + // instead of returning a failover error that the handler cannot retry. + if !clientOutputStarted && openAIStreamFailedEventShouldFailover(payloadBytes, message) { streamFailoverErr = s.newOpenAIStreamFailoverError(c, account, false, requestID, payloadBytes, message) return true } diff --git a/backend/internal/service/openai_gateway_messages_failed_response_test.go b/backend/internal/service/openai_gateway_messages_failed_response_test.go index a03efca42d..51d5fcb7c9 100644 --- a/backend/internal/service/openai_gateway_messages_failed_response_test.go +++ b/backend/internal/service/openai_gateway_messages_failed_response_test.go @@ -77,6 +77,84 @@ func TestForwardAsAnthropic_StreamingResponseFailed_ReturnsError(t *testing.T) { require.Contains(t, err.Error(), "upstream response failed") } +func TestForwardAsAnthropic_StreamingBareErrorAfterOutputIsVisible(t *testing.T) { + gin.SetMode(gin.TestMode) + + body := []byte(`{"model":"gpt-5.4","max_tokens":32,"messages":[{"role":"user","content":"hello"}],"stream":true}`) + rec := httptest.NewRecorder() + c, _ := gin.CreateTestContext(rec) + c.Request = httptest.NewRequest(http.MethodPost, "/v1/messages", bytes.NewReader(body)) + c.Request.Header.Set("Content-Type", "application/json") + + ssePayload := strings.Join([]string{ + `data: {"type":"response.created","response":{"id":"resp_bare_error","object":"response","model":"gpt-5.4","status":"in_progress","output":[]}}`, + "", + `data: {"type":"response.output_text.delta","output_index":0,"content_index":0,"delta":"partial"}`, + "", + `event: error`, + `data: {"type":"error","error":{"type":"server_error","code":"upstream_error","message":"mixed tools failed"}}`, + "", + `data: [DONE]`, + "", + }, "\n") + upstream := &httpUpstreamRecorder{resp: &http.Response{ + StatusCode: http.StatusOK, + Header: http.Header{"Content-Type": []string{"text/event-stream"}}, + Body: io.NopCloser(strings.NewReader(ssePayload)), + }} + svc := &OpenAIGatewayService{ + cfg: rawChatCompletionsTestConfig(), + httpUpstream: upstream, + } + + account := rawChatCompletionsTestAccount() + _, err := svc.ForwardAsAnthropic(context.Background(), c, account, body, "", "") + + require.Error(t, err) + require.Contains(t, err.Error(), "upstream response failed: mixed tools failed") + clientStream := rec.Body.String() + require.Contains(t, clientStream, `"text":"partial"`) + require.Contains(t, clientStream, "event: error") + require.Contains(t, clientStream, "mixed tools failed") + require.NotContains(t, clientStream, "event: message_stop") + require.NotContains(t, err.Error(), "missing terminal event") +} + +func TestForwardAsAnthropic_StreamingBareErrorBeforeOutputFailsOver(t *testing.T) { + gin.SetMode(gin.TestMode) + + body := []byte(`{"model":"gpt-5.4","max_tokens":32,"messages":[{"role":"user","content":"hello"}],"stream":true}`) + rec := httptest.NewRecorder() + c, _ := gin.CreateTestContext(rec) + c.Request = httptest.NewRequest(http.MethodPost, "/v1/messages", bytes.NewReader(body)) + c.Request.Header.Set("Content-Type", "application/json") + + ssePayload := strings.Join([]string{ + `event: error`, + `data: {"type":"error","error":{"type":"server_error","code":"upstream_error","message":"temporary upstream failure"}}`, + "", + `data: [DONE]`, + "", + }, "\n") + upstream := &httpUpstreamRecorder{resp: &http.Response{ + StatusCode: http.StatusOK, + Header: http.Header{"Content-Type": []string{"text/event-stream"}}, + Body: io.NopCloser(strings.NewReader(ssePayload)), + }} + svc := &OpenAIGatewayService{ + cfg: rawChatCompletionsTestConfig(), + httpUpstream: upstream, + } + + account := rawChatCompletionsTestAccount() + _, err := svc.ForwardAsAnthropic(context.Background(), c, account, body, "", "") + + require.Error(t, err) + var failoverErr *UpstreamFailoverError + require.True(t, errors.As(err, &failoverErr), "pre-output retryable error must remain failover-safe: %T: %v", err, err) + require.Empty(t, rec.Body.String(), "failover path must not commit downstream output") +} + func TestForwardAsAnthropic_BufferedResponseFailed_Failover(t *testing.T) { gin.SetMode(gin.TestMode)