mirror of
https://github.com/Wei-Shaw/sub2api.git
synced 2026-10-07 14:08:14 +08:00
fix(openai): surface Messages stream error events
This commit is contained in:
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user