mirror of
https://github.com/Wei-Shaw/sub2api.git
synced 2026-09-24 16:05:44 +08:00
Merge pull request #3853 from fengshao1227/fix/messages-transport-failover
fix(messages): /v1/messages 传输层错误对齐 failover 链路
This commit is contained in:
@@ -302,18 +302,7 @@ func (s *OpenAIGatewayService) ForwardAsAnthropic(
|
|||||||
}
|
}
|
||||||
resp, err := s.httpUpstream.Do(upstreamReq, proxyURL, account.ID, account.Concurrency)
|
resp, err := s.httpUpstream.Do(upstreamReq, proxyURL, account.ID, account.Concurrency)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
safeErr := sanitizeUpstreamErrorMessage(err.Error())
|
return nil, s.handleOpenAIUpstreamTransportError(ctx, c, account, err, false)
|
||||||
setOpsUpstreamError(c, 0, safeErr, "")
|
|
||||||
appendOpsUpstreamError(c, OpsUpstreamErrorEvent{
|
|
||||||
Platform: account.Platform,
|
|
||||||
AccountID: account.ID,
|
|
||||||
AccountName: account.Name,
|
|
||||||
UpstreamStatusCode: 0,
|
|
||||||
Kind: "request_error",
|
|
||||||
Message: safeErr,
|
|
||||||
})
|
|
||||||
writeAnthropicError(c, http.StatusBadGateway, "api_error", "Upstream request failed")
|
|
||||||
return nil, fmt.Errorf("upstream request failed: %s", safeErr)
|
|
||||||
}
|
}
|
||||||
defer func() { _ = resp.Body.Close() }()
|
defer func() { _ = resp.Body.Close() }()
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,92 @@
|
|||||||
|
//go:build unit
|
||||||
|
|
||||||
|
package service
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bytes"
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
|
"net/http"
|
||||||
|
"net/http/httptest"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"github.com/gin-gonic/gin"
|
||||||
|
"github.com/stretchr/testify/require"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestForwardAsAnthropic_TransportError_ReturnsFailoverError(t *testing.T) {
|
||||||
|
gin.SetMode(gin.TestMode)
|
||||||
|
|
||||||
|
body := []byte(`{"model":"gpt-5.4","max_tokens":32,"messages":[{"role":"user","content":"hello"}],"stream":false}`)
|
||||||
|
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")
|
||||||
|
|
||||||
|
upstream := &httpUpstreamRecorder{
|
||||||
|
err: errors.New(`dial tcp 1.2.3.4:443: connect: connection refused`),
|
||||||
|
}
|
||||||
|
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), "transport error should return UpstreamFailoverError for handler failover, got: %T", err)
|
||||||
|
require.Equal(t, http.StatusBadGateway, failoverErr.StatusCode)
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestForwardAsAnthropic_TransportError_DoesNotWriteResponse(t *testing.T) {
|
||||||
|
gin.SetMode(gin.TestMode)
|
||||||
|
|
||||||
|
body := []byte(`{"model":"gpt-5.4","max_tokens":32,"messages":[{"role":"user","content":"hello"}],"stream":false}`)
|
||||||
|
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")
|
||||||
|
|
||||||
|
upstream := &httpUpstreamRecorder{
|
||||||
|
err: errors.New(`read tcp: connection reset by peer`),
|
||||||
|
}
|
||||||
|
svc := &OpenAIGatewayService{
|
||||||
|
cfg: rawChatCompletionsTestConfig(),
|
||||||
|
httpUpstream: upstream,
|
||||||
|
}
|
||||||
|
|
||||||
|
account := rawChatCompletionsTestAccount()
|
||||||
|
_, _ = svc.ForwardAsAnthropic(context.Background(), c, account, body, "", "")
|
||||||
|
|
||||||
|
require.Equal(t, http.StatusOK, rec.Code, "transport error must not write HTTP response — handler owns the response for failover")
|
||||||
|
require.Empty(t, rec.Body.String(), "response body must be empty so handler can write the correct error or failover")
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestForwardAsAnthropic_TransportError_ClientCanceled_NoFailover(t *testing.T) {
|
||||||
|
gin.SetMode(gin.TestMode)
|
||||||
|
|
||||||
|
body := []byte(`{"model":"gpt-5.4","max_tokens":32,"messages":[{"role":"user","content":"hello"}],"stream":false}`)
|
||||||
|
rec := httptest.NewRecorder()
|
||||||
|
cancelCtx, cancel := context.WithCancel(context.Background())
|
||||||
|
cancel()
|
||||||
|
c, _ := gin.CreateTestContext(rec)
|
||||||
|
c.Request = httptest.NewRequest(http.MethodPost, "/v1/messages", bytes.NewReader(body)).WithContext(cancelCtx)
|
||||||
|
c.Request.Header.Set("Content-Type", "application/json")
|
||||||
|
|
||||||
|
upstream := &httpUpstreamRecorder{
|
||||||
|
err: context.Canceled,
|
||||||
|
}
|
||||||
|
svc := &OpenAIGatewayService{
|
||||||
|
cfg: rawChatCompletionsTestConfig(),
|
||||||
|
httpUpstream: upstream,
|
||||||
|
}
|
||||||
|
|
||||||
|
account := rawChatCompletionsTestAccount()
|
||||||
|
_, err := svc.ForwardAsAnthropic(context.Background(), c, account, body, "", "")
|
||||||
|
|
||||||
|
require.Error(t, err)
|
||||||
|
var failoverErr *UpstreamFailoverError
|
||||||
|
require.False(t, errors.As(err, &failoverErr), "client-canceled transport error should NOT trigger failover")
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user