package service import ( "bytes" "context" "errors" "io" "net/http" "net/http/httptest" "strings" "testing" "time" "github.com/Wei-Shaw/sub2api/internal/pkg/apicompat" "github.com/Wei-Shaw/sub2api/internal/pkg/tlsfingerprint" "github.com/gin-gonic/gin" "github.com/stretchr/testify/require" "github.com/tidwall/gjson" ) type openAIChatFailingWriter struct { gin.ResponseWriter failAfter int writes int } func (w *openAIChatFailingWriter) Write(p []byte) (int, error) { if w.writes >= w.failAfter { return 0, errors.New("write failed: client disconnected") } w.writes++ return w.ResponseWriter.Write(p) } type sequentialHTTPUpstreamRecorder struct { responses []*http.Response errs []error requests []*http.Request bodies [][]byte } func (u *sequentialHTTPUpstreamRecorder) Do(req *http.Request, proxyURL string, accountID int64, accountConcurrency int) (*http.Response, error) { u.requests = append(u.requests, req) if req != nil && req.Body != nil { b, _ := io.ReadAll(req.Body) u.bodies = append(u.bodies, append([]byte(nil), b...)) _ = req.Body.Close() req.Body = io.NopCloser(bytes.NewReader(b)) } if len(u.errs) > 0 { err := u.errs[0] u.errs = u.errs[1:] if err != nil { return nil, err } } if len(u.responses) > 0 { resp := u.responses[0] u.responses = u.responses[1:] return resp, nil } return nil, errors.New("no response configured") } func (u *sequentialHTTPUpstreamRecorder) DoWithTLS(req *http.Request, proxyURL string, accountID int64, accountConcurrency int, profile *tlsfingerprint.Profile) (*http.Response, error) { return u.Do(req, proxyURL, accountID, accountConcurrency) } func TestNormalizeResponsesRequestServiceTier(t *testing.T) { t.Parallel() req := &apicompat.ResponsesRequest{ServiceTier: " fast "} normalizeResponsesRequestServiceTier(req) require.Equal(t, "priority", req.ServiceTier) req.ServiceTier = "flex" normalizeResponsesRequestServiceTier(req) require.Equal(t, "flex", req.ServiceTier) // OpenAI 官方合法 tier 应被透传保留。 req.ServiceTier = "auto" normalizeResponsesRequestServiceTier(req) require.Equal(t, "auto", req.ServiceTier) req.ServiceTier = "default" normalizeResponsesRequestServiceTier(req) require.Equal(t, "default", req.ServiceTier) req.ServiceTier = "scale" normalizeResponsesRequestServiceTier(req) require.Equal(t, "scale", req.ServiceTier) // 真未知值仍被剥离。 req.ServiceTier = "turbo" normalizeResponsesRequestServiceTier(req) require.Empty(t, req.ServiceTier) } func TestNormalizeResponsesBodyServiceTier(t *testing.T) { t.Parallel() body, tier, err := normalizeResponsesBodyServiceTier([]byte(`{"model":"gpt-5.1","service_tier":"fast"}`)) require.NoError(t, err) require.Equal(t, "priority", tier) require.Equal(t, "priority", gjson.GetBytes(body, "service_tier").String()) body, tier, err = normalizeResponsesBodyServiceTier([]byte(`{"model":"gpt-5.1","service_tier":"flex"}`)) require.NoError(t, err) require.Equal(t, "flex", tier) require.Equal(t, "flex", gjson.GetBytes(body, "service_tier").String()) // OpenAI 官方 tier 直接保留在 body 中(透传上游)。 body, tier, err = normalizeResponsesBodyServiceTier([]byte(`{"model":"gpt-5.1","service_tier":"auto"}`)) require.NoError(t, err) require.Equal(t, "auto", tier) require.Equal(t, "auto", gjson.GetBytes(body, "service_tier").String()) body, tier, err = normalizeResponsesBodyServiceTier([]byte(`{"model":"gpt-5.1","service_tier":"default"}`)) require.NoError(t, err) require.Equal(t, "default", tier) require.Equal(t, "default", gjson.GetBytes(body, "service_tier").String()) body, tier, err = normalizeResponsesBodyServiceTier([]byte(`{"model":"gpt-5.1","service_tier":"scale"}`)) require.NoError(t, err) require.Equal(t, "scale", tier) require.Equal(t, "scale", gjson.GetBytes(body, "service_tier").String()) // 真未知值才会被删除。 body, tier, err = normalizeResponsesBodyServiceTier([]byte(`{"model":"gpt-5.1","service_tier":"turbo"}`)) require.NoError(t, err) require.Empty(t, tier) require.False(t, gjson.GetBytes(body, "service_tier").Exists()) } func TestForwardAsChatCompletions_UnknownModelDoesNotUseDefaultMappedModel(t *testing.T) { gin.SetMode(gin.TestMode) rec := httptest.NewRecorder() c, _ := gin.CreateTestContext(rec) body := []byte(`{"model":"gpt6","messages":[{"role":"user","content":"hello"}],"stream":false}`) c.Request = httptest.NewRequest(http.MethodPost, "/v1/chat/completions", bytes.NewReader(body)) c.Request.Header.Set("Content-Type", "application/json") upstream := &httpUpstreamRecorder{resp: &http.Response{ StatusCode: http.StatusBadRequest, Header: http.Header{"Content-Type": []string{"application/json"}, "x-request-id": []string{"rid_chat_unknown_model"}}, Body: io.NopCloser(strings.NewReader(`{"error":{"type":"invalid_request_error","message":"model not found"}}`)), }} svc := &OpenAIGatewayService{httpUpstream: upstream} account := &Account{ ID: 1, Name: "openai-oauth", Platform: PlatformOpenAI, Type: AccountTypeOAuth, Concurrency: 1, Credentials: map[string]any{ "access_token": "oauth-token", "chatgpt_account_id": "chatgpt-acc", }, } result, err := svc.ForwardAsChatCompletions(context.Background(), c, account, body, "", "gpt-5.4") require.Error(t, err) require.Nil(t, result) require.Equal(t, "gpt6", gjson.GetBytes(upstream.lastBody, "model").String()) require.NotEqual(t, "gpt-5.4", gjson.GetBytes(upstream.lastBody, "model").String()) require.Equal(t, http.StatusBadRequest, rec.Code) } func TestForwardAsChatCompletions_RequestErrorRetriesBeforeSuccess(t *testing.T) { gin.SetMode(gin.TestMode) rec := httptest.NewRecorder() c, _ := gin.CreateTestContext(rec) body := []byte(`{"model":"gpt-5.4","messages":[{"role":"user","content":"hello"}],"stream":false}`) c.Request = httptest.NewRequest(http.MethodPost, "/v1/chat/completions", bytes.NewReader(body)) c.Request.Header.Set("Content-Type", "application/json") upstreamBody := strings.Join([]string{ `data: {"type":"response.completed","response":{"id":"resp_retry","object":"response","model":"gpt-5.4","status":"completed","output":[{"type":"message","id":"msg_1","role":"assistant","status":"completed","content":[{"type":"output_text","text":"ok"}]}],"usage":{"input_tokens":3,"output_tokens":2,"total_tokens":5}}}`, "", "data: [DONE]", "", }, "\n") upstream := &sequentialHTTPUpstreamRecorder{ errs: []error{ errors.New("Post \"https://chatgpt.com/backend-api/codex/responses\": read tcp 172.18.0.4:60076->42.193.179.21:1081: read: connection reset by peer"), errors.New("connection reset by peer"), errors.New("unexpected EOF"), nil, }, responses: []*http.Response{{ StatusCode: http.StatusOK, Header: http.Header{"Content-Type": []string{"text/event-stream"}, "x-request-id": []string{"rid_retry_success"}}, Body: io.NopCloser(strings.NewReader(upstreamBody)), }}, } svc := &OpenAIGatewayService{httpUpstream: upstream} account := &Account{ ID: 1, Name: "openai-oauth", Platform: PlatformOpenAI, Type: AccountTypeOAuth, Concurrency: 1, Credentials: map[string]any{ "access_token": "oauth-token", "chatgpt_account_id": "chatgpt-acc", }, } result, err := svc.ForwardAsChatCompletions(context.Background(), c, account, body, "", "gpt-5.4") require.NoError(t, err) require.NotNil(t, result) require.Equal(t, http.StatusOK, rec.Code) require.Len(t, upstream.requests, 4) require.Len(t, upstream.bodies, 4) require.Equal(t, upstream.bodies[0], upstream.bodies[3], "retry must rebuild the same upstream body") rawEvents, ok := c.Get(OpsUpstreamErrorsKey) require.True(t, ok) events, ok := rawEvents.([]*OpsUpstreamErrorEvent) require.True(t, ok) require.Len(t, events, 3) require.Equal(t, "request_error", events[0].Kind) require.Contains(t, events[0].Message, "connection reset by peer") } func TestForwardAsChatCompletions_ClosedNetworkConnectionRetriesBeforeSuccess(t *testing.T) { gin.SetMode(gin.TestMode) rec := httptest.NewRecorder() c, _ := gin.CreateTestContext(rec) body := []byte(`{"model":"gpt-5.4","messages":[{"role":"user","content":"hello"}],"stream":false}`) c.Request = httptest.NewRequest(http.MethodPost, "/v1/chat/completions", bytes.NewReader(body)) c.Request.Header.Set("Content-Type", "application/json") upstreamBody := strings.Join([]string{ `data: {"type":"response.completed","response":{"id":"resp_closed_network_retry","object":"response","model":"gpt-5.4","status":"completed","output":[{"type":"message","id":"msg_1","role":"assistant","status":"completed","content":[{"type":"output_text","text":"ok"}]}],"usage":{"input_tokens":3,"output_tokens":2,"total_tokens":5}}}`, "", "data: [DONE]", "", }, "\n") upstream := &sequentialHTTPUpstreamRecorder{ errs: []error{ errors.New("Post \"https://chatgpt.com/backend-api/codex/responses\": use of closed network connection"), nil, }, responses: []*http.Response{{ StatusCode: http.StatusOK, Header: http.Header{"Content-Type": []string{"text/event-stream"}, "x-request-id": []string{"rid_closed_network_retry"}}, Body: io.NopCloser(strings.NewReader(upstreamBody)), }}, } svc := &OpenAIGatewayService{httpUpstream: upstream} account := &Account{ ID: 1, Name: "openai-oauth", Platform: PlatformOpenAI, Type: AccountTypeOAuth, Concurrency: 1, Credentials: map[string]any{ "access_token": "oauth-token", "chatgpt_account_id": "chatgpt-acc", }, } result, err := svc.ForwardAsChatCompletions(context.Background(), c, account, body, "", "gpt-5.4") require.NoError(t, err) require.NotNil(t, result) require.Equal(t, http.StatusOK, rec.Code) require.Len(t, upstream.requests, 2) rawEvents, ok := c.Get(OpsUpstreamErrorsKey) require.True(t, ok) events, ok := rawEvents.([]*OpsUpstreamErrorEvent) require.True(t, ok) require.Len(t, events, 1) require.Equal(t, "request_error", events[0].Kind) require.Contains(t, events[0].Message, "use of closed network connection") } func TestForwardAsChatCompletions_TLSBadRecordMACRetriesBeforeSuccess(t *testing.T) { gin.SetMode(gin.TestMode) rec := httptest.NewRecorder() c, _ := gin.CreateTestContext(rec) body := []byte(`{"model":"gpt-5.4","messages":[{"role":"user","content":"hello"}],"stream":false}`) c.Request = httptest.NewRequest(http.MethodPost, "/v1/chat/completions", bytes.NewReader(body)) c.Request.Header.Set("Content-Type", "application/json") upstreamBody := strings.Join([]string{ `data: {"type":"response.completed","response":{"id":"resp_tls_retry","object":"response","model":"gpt-5.4","status":"completed","output":[{"type":"message","id":"msg_1","role":"assistant","status":"completed","content":[{"type":"output_text","text":"ok"}]}],"usage":{"input_tokens":3,"output_tokens":2,"total_tokens":5}}}`, "", "data: [DONE]", "", }, "\n") upstream := &sequentialHTTPUpstreamRecorder{ errs: []error{ errors.New("Post \"https://chatgpt.com/backend-api/codex/responses\": local error: tls: bad record MAC"), nil, }, responses: []*http.Response{{ StatusCode: http.StatusOK, Header: http.Header{"Content-Type": []string{"text/event-stream"}, "x-request-id": []string{"rid_tls_retry"}}, Body: io.NopCloser(strings.NewReader(upstreamBody)), }}, } svc := &OpenAIGatewayService{httpUpstream: upstream} account := &Account{ ID: 1, Name: "openai-oauth", Platform: PlatformOpenAI, Type: AccountTypeOAuth, Concurrency: 1, Credentials: map[string]any{ "access_token": "oauth-token", "chatgpt_account_id": "chatgpt-acc", }, } result, err := svc.ForwardAsChatCompletions(context.Background(), c, account, body, "", "gpt-5.4") require.NoError(t, err) require.NotNil(t, result) require.Equal(t, http.StatusOK, rec.Code) require.Len(t, upstream.requests, 2) rawEvents, ok := c.Get(OpsUpstreamErrorsKey) require.True(t, ok) events, ok := rawEvents.([]*OpsUpstreamErrorEvent) require.True(t, ok) require.Len(t, events, 1) require.Equal(t, "request_error", events[0].Kind) require.Contains(t, strings.ToLower(events[0].Message), "tls: bad record mac") } func TestForwardAsChatCompletions_RequestErrorExhaustionReturnsFailover(t *testing.T) { gin.SetMode(gin.TestMode) rec := httptest.NewRecorder() c, _ := gin.CreateTestContext(rec) body := []byte(`{"model":"gpt-5.4","messages":[{"role":"user","content":"hello"}],"stream":false}`) c.Request = httptest.NewRequest(http.MethodPost, "/v1/chat/completions", bytes.NewReader(body)) c.Request.Header.Set("Content-Type", "application/json") upstream := &sequentialHTTPUpstreamRecorder{ errs: []error{ errors.New("connection reset by peer"), errors.New("connection reset by peer"), errors.New("connection reset by peer"), errors.New("connection reset by peer"), }, } svc := &OpenAIGatewayService{httpUpstream: upstream} account := &Account{ ID: 1, Name: "openai-oauth", Platform: PlatformOpenAI, Type: AccountTypeOAuth, Concurrency: 1, Credentials: map[string]any{ "access_token": "oauth-token", "chatgpt_account_id": "chatgpt-acc", }, } result, err := svc.ForwardAsChatCompletions(context.Background(), c, account, body, "", "gpt-5.4") require.Nil(t, result) var failoverErr *UpstreamFailoverError require.ErrorAs(t, err, &failoverErr) require.Equal(t, http.StatusBadGateway, failoverErr.StatusCode) require.False(t, c.Writer.Written(), "forward should not write a 502 before handler failover") require.Len(t, upstream.requests, 4) rawEvents, ok := c.Get(OpsUpstreamErrorsKey) require.True(t, ok) events, ok := rawEvents.([]*OpsUpstreamErrorEvent) require.True(t, ok) require.Len(t, events, 4) require.Equal(t, "request_error:retry_exhausted", events[3].Kind) } func TestForwardAsChatCompletions_ClientDisconnectDrainsUpstreamUsage(t *testing.T) { gin.SetMode(gin.TestMode) rec := httptest.NewRecorder() c, _ := gin.CreateTestContext(rec) c.Writer = &openAIChatFailingWriter{ResponseWriter: c.Writer, failAfter: 0} body := []byte(`{"model":"gpt-5.4","messages":[{"role":"user","content":"hello"}],"stream":true}`) c.Request = httptest.NewRequest(http.MethodPost, "/v1/chat/completions", bytes.NewReader(body)) c.Request.Header.Set("Content-Type", "application/json") upstreamBody := strings.Join([]string{ `data: {"type":"response.created","response":{"id":"resp_1","model":"gpt-5.4","status":"in_progress","output":[]}}`, "", `data: {"type":"response.output_text.delta","delta":"ok"}`, "", `data: {"type":"response.completed","response":{"id":"resp_1","object":"response","model":"gpt-5.4","status":"completed","output":[{"type":"message","id":"msg_1","role":"assistant","status":"completed","content":[{"type":"output_text","text":"ok"}]}],"usage":{"input_tokens":11,"output_tokens":5,"total_tokens":16,"input_tokens_details":{"cached_tokens":4}}}}`, "", "data: [DONE]", "", }, "\n") upstream := &httpUpstreamRecorder{resp: &http.Response{ StatusCode: http.StatusOK, Header: http.Header{"Content-Type": []string{"text/event-stream"}, "x-request-id": []string{"rid_chat_disconnect"}}, Body: io.NopCloser(strings.NewReader(upstreamBody)), }} svc := &OpenAIGatewayService{httpUpstream: upstream} account := &Account{ ID: 1, Name: "openai-oauth", Platform: PlatformOpenAI, Type: AccountTypeOAuth, Concurrency: 1, Credentials: map[string]any{ "access_token": "oauth-token", "chatgpt_account_id": "chatgpt-acc", }, } result, err := svc.ForwardAsChatCompletions(context.Background(), c, account, body, "", "gpt-5.1") require.NoError(t, err) require.NotNil(t, result) require.Equal(t, 11, result.Usage.InputTokens) require.Equal(t, 5, result.Usage.OutputTokens) require.Equal(t, 4, result.Usage.CacheReadInputTokens) } func TestForwardAsChatCompletions_TerminalUsageWithoutUpstreamCloseReturns(t *testing.T) { gin.SetMode(gin.TestMode) rec := httptest.NewRecorder() c, _ := gin.CreateTestContext(rec) c.Writer = &openAIChatFailingWriter{ResponseWriter: c.Writer, failAfter: 0} body := []byte(`{"model":"gpt-5.4","messages":[{"role":"user","content":"hello"}],"stream":true}`) c.Request = httptest.NewRequest(http.MethodPost, "/v1/chat/completions", bytes.NewReader(body)) c.Request.Header.Set("Content-Type", "application/json") upstreamBody := []byte(`data: {"type":"response.completed","response":{"id":"resp_1","object":"response","model":"gpt-5.4","status":"completed","output":[{"type":"message","id":"msg_1","role":"assistant","status":"completed","content":[{"type":"output_text","text":"ok"}]}],"usage":{"input_tokens":17,"output_tokens":8,"total_tokens":25,"input_tokens_details":{"cached_tokens":6}}}}` + "\n\n") upstreamStream := newOpenAICompatBlockingReadCloser(upstreamBody) defer func() { require.NoError(t, upstreamStream.Close()) }() upstream := &httpUpstreamRecorder{resp: &http.Response{ StatusCode: http.StatusOK, Header: http.Header{"Content-Type": []string{"text/event-stream"}, "x-request-id": []string{"rid_chat_terminal_no_close"}}, Body: upstreamStream, }} svc := &OpenAIGatewayService{httpUpstream: upstream} account := &Account{ ID: 1, Name: "openai-oauth", Platform: PlatformOpenAI, Type: AccountTypeOAuth, Concurrency: 1, Credentials: map[string]any{ "access_token": "oauth-token", "chatgpt_account_id": "chatgpt-acc", }, } type forwardResult struct { result *OpenAIForwardResult err error } resultCh := make(chan forwardResult, 1) go func() { result, err := svc.ForwardAsChatCompletions(context.Background(), c, account, body, "", "gpt-5.1") resultCh <- forwardResult{result: result, err: err} }() select { case got := <-resultCh: require.NoError(t, got.err) require.NotNil(t, got.result) require.Equal(t, 17, got.result.Usage.InputTokens) require.Equal(t, 8, got.result.Usage.OutputTokens) require.Equal(t, 6, got.result.Usage.CacheReadInputTokens) case <-time.After(time.Second): require.Fail(t, "ForwardAsChatCompletions should return after terminal usage event even if upstream keeps the connection open") } } func TestForwardAsChatCompletions_BufferedTerminalWithoutUpstreamCloseReturns(t *testing.T) { gin.SetMode(gin.TestMode) rec := httptest.NewRecorder() c, _ := gin.CreateTestContext(rec) body := []byte(`{"model":"gpt-5.4","messages":[{"role":"user","content":"hello"}],"stream":false}`) c.Request = httptest.NewRequest(http.MethodPost, "/v1/chat/completions", bytes.NewReader(body)) c.Request.Header.Set("Content-Type", "application/json") upstreamBody := []byte(`data: {"type":"response.completed","response":{"id":"resp_1","object":"response","model":"gpt-5.4","status":"completed","output":[{"type":"message","id":"msg_1","role":"assistant","status":"completed","content":[{"type":"output_text","text":"ok"}]}],"usage":{"input_tokens":17,"output_tokens":8,"total_tokens":25,"input_tokens_details":{"cached_tokens":6}}}}` + "\n\n") upstreamStream := newOpenAICompatBlockingReadCloser(upstreamBody) defer func() { require.NoError(t, upstreamStream.Close()) }() upstream := &httpUpstreamRecorder{resp: &http.Response{ StatusCode: http.StatusOK, Header: http.Header{"Content-Type": []string{"text/event-stream"}, "x-request-id": []string{"rid_chat_buffered_terminal_no_close"}}, Body: upstreamStream, }} svc := &OpenAIGatewayService{httpUpstream: upstream} account := &Account{ ID: 1, Name: "openai-oauth", Platform: PlatformOpenAI, Type: AccountTypeOAuth, Concurrency: 1, Credentials: map[string]any{ "access_token": "oauth-token", "chatgpt_account_id": "chatgpt-acc", }, } type forwardResult struct { result *OpenAIForwardResult err error } resultCh := make(chan forwardResult, 1) go func() { result, err := svc.ForwardAsChatCompletions(context.Background(), c, account, body, "", "gpt-5.1") resultCh <- forwardResult{result: result, err: err} }() select { case got := <-resultCh: require.NoError(t, got.err) require.NotNil(t, got.result) require.Equal(t, 17, got.result.Usage.InputTokens) require.Equal(t, 8, got.result.Usage.OutputTokens) require.Equal(t, 6, got.result.Usage.CacheReadInputTokens) require.Contains(t, rec.Body.String(), `"finish_reason":"stop"`) case <-time.After(time.Second): require.Fail(t, "ForwardAsChatCompletions buffered response should return after terminal usage event even if upstream keeps the connection open") } } func TestForwardAsChatCompletions_DoneSentinelWithoutTerminalReturnsError(t *testing.T) { gin.SetMode(gin.TestMode) rec := httptest.NewRecorder() c, _ := gin.CreateTestContext(rec) body := []byte(`{"model":"gpt-5.4","messages":[{"role":"user","content":"hello"}],"stream":true}`) c.Request = httptest.NewRequest(http.MethodPost, "/v1/chat/completions", bytes.NewReader(body)) c.Request.Header.Set("Content-Type", "application/json") upstreamBody := "data: [DONE]\n\n" upstream := &httpUpstreamRecorder{resp: &http.Response{ StatusCode: http.StatusOK, Header: http.Header{"Content-Type": []string{"text/event-stream"}, "x-request-id": []string{"rid_chat_missing_terminal"}}, Body: io.NopCloser(strings.NewReader(upstreamBody)), }} svc := &OpenAIGatewayService{httpUpstream: upstream} account := &Account{ ID: 1, Name: "openai-oauth", Platform: PlatformOpenAI, Type: AccountTypeOAuth, Concurrency: 1, Credentials: map[string]any{ "access_token": "oauth-token", "chatgpt_account_id": "chatgpt-acc", }, } result, err := svc.ForwardAsChatCompletions(context.Background(), c, account, body, "", "gpt-5.1") require.Error(t, err) require.Contains(t, err.Error(), "missing terminal event") require.NotNil(t, result) require.Zero(t, result.Usage.InputTokens) require.Zero(t, result.Usage.OutputTokens) } func TestForwardAsChatCompletions_UpstreamRequestIgnoresClientCancel(t *testing.T) { gin.SetMode(gin.TestMode) rec := httptest.NewRecorder() c, _ := gin.CreateTestContext(rec) reqCtx, cancel := context.WithCancel(context.Background()) body := []byte(`{"model":"gpt-5.4","messages":[{"role":"user","content":"hello"}],"stream":false}`) c.Request = httptest.NewRequest(http.MethodPost, "/v1/chat/completions", bytes.NewReader(body)).WithContext(reqCtx) c.Request.Header.Set("Content-Type", "application/json") cancel() upstreamBody := strings.Join([]string{ `data: {"type":"response.completed","response":{"id":"resp_1","object":"response","model":"gpt-5.4","status":"completed","output":[{"type":"message","id":"msg_1","role":"assistant","status":"completed","content":[{"type":"output_text","text":"ok"}]}],"usage":{"input_tokens":5,"output_tokens":2,"total_tokens":7}}}`, "", "data: [DONE]", "", }, "\n") upstream := &httpUpstreamRecorder{resp: &http.Response{ StatusCode: http.StatusOK, Header: http.Header{"Content-Type": []string{"text/event-stream"}, "x-request-id": []string{"rid_chat_ctx"}}, Body: io.NopCloser(strings.NewReader(upstreamBody)), }} svc := &OpenAIGatewayService{httpUpstream: upstream} account := &Account{ ID: 1, Name: "openai-oauth", Platform: PlatformOpenAI, Type: AccountTypeOAuth, Concurrency: 1, Credentials: map[string]any{ "access_token": "oauth-token", "chatgpt_account_id": "chatgpt-acc", }, } result, err := svc.ForwardAsChatCompletions(reqCtx, c, account, body, "", "gpt-5.1") require.NoError(t, err) require.NotNil(t, result) require.NotNil(t, upstream.lastReq) require.NoError(t, upstream.lastReq.Context().Err()) }