From 2e2d43a84f85a2598b79fa86f84934cd50d320e4 Mon Sep 17 00:00:00 2001 From: Seefs Date: Wed, 24 Jun 2026 17:31:46 +0800 Subject: [PATCH] perf: optimize Responses relay body and billing parsing --- common/body_storage.go | 20 +++ dto/openai_response.go | 143 +++++++++++++++++++++ relay/channel/api_request.go | 34 ++++- relay/channel/openai/chat_via_responses.go | 6 +- relay/channel/openai/helper.go | 4 +- relay/channel/openai/relay_responses.go | 34 ++--- relay/chat_completions_via_responses.go | 10 +- relay/claude_handler.go | 12 +- relay/common/outbound_body.go | 11 ++ relay/compatible_handler.go | 12 +- relay/outbound_body.go | 17 +++ relay/responses_handler.go | 12 +- 12 files changed, 275 insertions(+), 40 deletions(-) create mode 100644 relay/outbound_body.go diff --git a/common/body_storage.go b/common/body_storage.go index 094dbda36d3f..3246f2fd7e6e 100644 --- a/common/body_storage.go +++ b/common/body_storage.go @@ -14,6 +14,8 @@ import ( type BodyStorage interface { io.ReadSeeker io.Closer + // Open creates an independent reader for the stored body. + Open() (io.ReadCloser, error) // Bytes 获取全部内容 Bytes() ([]byte, error) // Size 获取数据大小 @@ -71,6 +73,15 @@ func (m *memoryStorage) Close() error { return nil } +func (m *memoryStorage) Open() (io.ReadCloser, error) { + m.mu.Lock() + defer m.mu.Unlock() + if atomic.LoadInt32(&m.closed) == 1 { + return nil, ErrStorageClosed + } + return io.NopCloser(bytes.NewReader(m.data)), nil +} + func (m *memoryStorage) Bytes() ([]byte, error) { m.mu.Lock() defer m.mu.Unlock() @@ -195,6 +206,15 @@ func (d *diskStorage) Close() error { return nil } +func (d *diskStorage) Open() (io.ReadCloser, error) { + d.mu.Lock() + defer d.mu.Unlock() + if atomic.LoadInt32(&d.closed) == 1 { + return nil, ErrStorageClosed + } + return os.Open(d.filePath) +} + func (d *diskStorage) Bytes() ([]byte, error) { d.mu.Lock() defer d.mu.Unlock() diff --git a/dto/openai_response.go b/dto/openai_response.go index 0e6b818dbd8b..35d7cf17087f 100644 --- a/dto/openai_response.go +++ b/dto/openai_response.go @@ -292,6 +292,149 @@ type OpenAIResponsesResponse struct { Metadata json.RawMessage `json:"metadata"` } +type ResponsesBillingInputTokenDetails struct { + CachedTokens int `json:"cached_tokens"` + AudioTokens int `json:"audio_tokens"` + ImageTokens int `json:"image_tokens"` +} + +type ResponsesBillingOutputTokenDetails struct { + ReasoningTokens int `json:"reasoning_tokens"` +} + +type ResponsesBillingUsage struct { + InputTokens int `json:"input_tokens"` + OutputTokens int `json:"output_tokens"` + TotalTokens int `json:"total_tokens"` + + InputTokensDetails *ResponsesBillingInputTokenDetails `json:"input_tokens_details"` + OutputTokensDetails *ResponsesBillingOutputTokenDetails `json:"output_tokens_details"` + CompletionTokenDetails *ResponsesBillingOutputTokenDetails `json:"completion_tokens_details"` +} + +type ResponsesBillingTool struct { + Type string `json:"type"` +} + +type ResponsesBillingOutput struct { + Type string `json:"type"` + Quality string `json:"quality"` + Size string `json:"size"` +} + +type ResponsesBillingMeta struct { + Error json.RawMessage `json:"error"` + Usage *ResponsesBillingUsage `json:"usage"` + Tools []ResponsesBillingTool `json:"tools"` + Output []ResponsesBillingOutput `json:"output"` +} + +type ResponsesBillingStreamItem struct { + Type string `json:"type"` +} + +type ResponsesBillingStreamResponse struct { + Type string `json:"type"` + Response *ResponsesBillingMeta `json:"response,omitempty"` + Delta string `json:"delta,omitempty"` + Item *ResponsesBillingStreamItem `json:"item,omitempty"` +} + +type ResponsesTranslatedStreamMeta struct { + Model string `json:"model"` + CreatedAt int `json:"created_at"` + Error json.RawMessage `json:"error"` + Usage *ResponsesBillingUsage `json:"usage"` +} + +type ResponsesTranslatedStreamItem struct { + Type string `json:"type"` + ID string `json:"id"` + CallId string `json:"call_id,omitempty"` + Name string `json:"name,omitempty"` + Arguments json.RawMessage `json:"arguments,omitempty"` +} + +func (r *ResponsesTranslatedStreamItem) ArgumentsString() string { + if r == nil { + return "" + } + return ResponsesArgumentsString(r.Arguments) +} + +type ResponsesTranslatedStreamResponse struct { + Type string `json:"type"` + Response *ResponsesTranslatedStreamMeta `json:"response,omitempty"` + Delta string `json:"delta,omitempty"` + Item *ResponsesTranslatedStreamItem `json:"item,omitempty"` + OutputIndex *int `json:"output_index,omitempty"` + ContentIndex *int `json:"content_index,omitempty"` + SummaryIndex *int `json:"summary_index,omitempty"` + ItemID string `json:"item_id,omitempty"` + Part *ResponsesReasoningSummaryPart `json:"part,omitempty"` +} + +func getOpenAIErrorFromRawMessage(errorField json.RawMessage) *types.OpenAIError { + if len(errorField) == 0 || common.GetJsonType(errorField) == "null" { + return nil + } + var decoded any + if err := common.Unmarshal(errorField, &decoded); err != nil { + return nil + } + return GetOpenAIError(decoded) +} + +func (m *ResponsesBillingMeta) GetOpenAIError() *types.OpenAIError { + if m == nil { + return nil + } + return getOpenAIErrorFromRawMessage(m.Error) +} + +func (m *ResponsesTranslatedStreamMeta) GetOpenAIError() *types.OpenAIError { + if m == nil { + return nil + } + return getOpenAIErrorFromRawMessage(m.Error) +} + +func (m *ResponsesBillingMeta) HasImageGenerationCall() bool { + if m == nil || len(m.Output) == 0 { + return false + } + for _, output := range m.Output { + if output.Type == ResponsesOutputTypeImageGenerationCall { + return true + } + } + return false +} + +func (m *ResponsesBillingMeta) GetQuality() string { + if m == nil || len(m.Output) == 0 { + return "" + } + for _, output := range m.Output { + if output.Type == ResponsesOutputTypeImageGenerationCall { + return output.Quality + } + } + return "" +} + +func (m *ResponsesBillingMeta) GetSize() string { + if m == nil || len(m.Output) == 0 { + return "" + } + for _, output := range m.Output { + if output.Type == ResponsesOutputTypeImageGenerationCall { + return output.Size + } + } + return "" +} + // GetOpenAIError 从动态错误类型中提取OpenAIError结构 func (o *OpenAIResponsesResponse) GetOpenAIError() *types.OpenAIError { return GetOpenAIError(o.Error) diff --git a/relay/channel/api_request.go b/relay/channel/api_request.go index f945a8383821..a86fb7062bb0 100644 --- a/relay/channel/api_request.go +++ b/relay/channel/api_request.go @@ -304,13 +304,36 @@ func applyHeaderOverrideToRequest(req *http.Request, headerOverride map[string]s } } +func newUpstreamRequest(method string, url string, requestBody io.Reader) (*http.Request, error) { + if replayableBody, ok := requestBody.(common.ReplayableBody); ok { + body, err := replayableBody.Open() + if err != nil { + return nil, err + } + req, err := http.NewRequest(method, url, body) + if err != nil { + _ = body.Close() + return nil, err + } + req.ContentLength = replayableBody.Size() + return req, nil + } + return http.NewRequest(method, url, requestBody) +} + +func closeUpstreamRequestBody(req *http.Request) { + if req != nil && req.Body != nil { + _ = req.Body.Close() + } +} + func DoApiRequest(a Adaptor, c *gin.Context, info *common.RelayInfo, requestBody io.Reader) (*http.Response, error) { fullRequestURL, err := a.GetRequestURL(info) if err != nil { return nil, fmt.Errorf("get request url failed: %w", err) } logger.LogDebug(c, "fullRequestURL: %s", fullRequestURL) - req, err := http.NewRequest(c.Request.Method, fullRequestURL, requestBody) + req, err := newUpstreamRequest(c.Request.Method, fullRequestURL, requestBody) if err != nil { return nil, fmt.Errorf("new request failed: %w", err) } @@ -318,12 +341,14 @@ func DoApiRequest(a Adaptor, c *gin.Context, info *common.RelayInfo, requestBody headers := req.Header err = a.SetupRequestHeader(c, &headers, info) if err != nil { + closeUpstreamRequestBody(req) return nil, fmt.Errorf("setup request header failed: %w", err) } // 在 SetupRequestHeader 之后应用 Header Override,确保用户设置优先级最高 // 这样可以覆盖默认的 Authorization header 设置 headerOverride, err := processHeaderOverride(info, c) if err != nil { + closeUpstreamRequestBody(req) return nil, err } applyHeaderOverrideToRequest(req, headerOverride) @@ -340,7 +365,7 @@ func DoFormRequest(a Adaptor, c *gin.Context, info *common.RelayInfo, requestBod return nil, fmt.Errorf("get request url failed: %w", err) } logger.LogDebug(c, "fullRequestURL: %s", fullRequestURL) - req, err := http.NewRequest(c.Request.Method, fullRequestURL, requestBody) + req, err := newUpstreamRequest(c.Request.Method, fullRequestURL, requestBody) if err != nil { return nil, fmt.Errorf("new request failed: %w", err) } @@ -350,12 +375,14 @@ func DoFormRequest(a Adaptor, c *gin.Context, info *common.RelayInfo, requestBod headers := req.Header err = a.SetupRequestHeader(c, &headers, info) if err != nil { + closeUpstreamRequestBody(req) return nil, fmt.Errorf("setup request header failed: %w", err) } // 在 SetupRequestHeader 之后应用 Header Override,确保用户设置优先级最高 // 这样可以覆盖默认的 Authorization header 设置 headerOverride, err := processHeaderOverride(info, c) if err != nil { + closeUpstreamRequestBody(req) return nil, err } applyHeaderOverrideToRequest(req, headerOverride) @@ -514,6 +541,8 @@ func doRequest(c *gin.Context, req *http.Request, info *common.RelayInfo) (*http } } + defer closeUpstreamRequestBody(req) + resp, err := client.Do(req) if err != nil { logger.LogError(c, "do request failed: "+err.Error()) @@ -527,7 +556,6 @@ func doRequest(c *gin.Context, req *http.Request, info *common.RelayInfo) (*http c.Set(common2.UpstreamRequestIdKey, upID) } - _ = req.Body.Close() _ = c.Request.Body.Close() return resp, nil } diff --git a/relay/channel/openai/chat_via_responses.go b/relay/channel/openai/chat_via_responses.go index 2c0752275daa..f3d2eb19aeae 100644 --- a/relay/channel/openai/chat_via_responses.go +++ b/relay/channel/openai/chat_via_responses.go @@ -302,7 +302,7 @@ func OaiResponsesToChatStreamHandler(c *gin.Context, info *relaycommon.RelayInfo return } - var streamResp dto.ResponsesStreamResponse + var streamResp dto.ResponsesTranslatedStreamResponse if err := common.UnmarshalJsonStr(data, &streamResp); err != nil { logger.LogError(c, "failed to unmarshal responses stream event: "+err.Error()) sr.Error(err) @@ -469,7 +469,9 @@ func OaiResponsesToChatStreamHandler(c *gin.Context, info *relaycommon.RelayInfo usage.PromptTokensDetails.ImageTokens = streamResp.Response.Usage.InputTokensDetails.ImageTokens usage.PromptTokensDetails.AudioTokens = streamResp.Response.Usage.InputTokensDetails.AudioTokens } - if streamResp.Response.Usage.CompletionTokenDetails.ReasoningTokens != 0 { + if streamResp.Response.Usage.OutputTokensDetails != nil && streamResp.Response.Usage.OutputTokensDetails.ReasoningTokens != 0 { + usage.CompletionTokenDetails.ReasoningTokens = streamResp.Response.Usage.OutputTokensDetails.ReasoningTokens + } else if streamResp.Response.Usage.CompletionTokenDetails != nil && streamResp.Response.Usage.CompletionTokenDetails.ReasoningTokens != 0 { usage.CompletionTokenDetails.ReasoningTokens = streamResp.Response.Usage.CompletionTokenDetails.ReasoningTokens } } diff --git a/relay/channel/openai/helper.go b/relay/channel/openai/helper.go index 1a01d06da6dc..2a56cb0c5307 100644 --- a/relay/channel/openai/helper.go +++ b/relay/channel/openai/helper.go @@ -202,9 +202,9 @@ func HandleFinalResponse(c *gin.Context, info *relaycommon.RelayInfo, lastStream } } -func sendResponsesStreamData(c *gin.Context, streamResponse dto.ResponsesStreamResponse, data string) { +func sendResponsesStreamData(c *gin.Context, eventType string, data string) { if data == "" { return } - helper.ResponseChunkData(c, streamResponse, data) + helper.ResponseChunkData(c, dto.ResponsesStreamResponse{Type: eventType}, data) } diff --git a/relay/channel/openai/relay_responses.go b/relay/channel/openai/relay_responses.go index 2665b8d027e9..b13c1b3f9d04 100644 --- a/relay/channel/openai/relay_responses.go +++ b/relay/channel/openai/relay_responses.go @@ -21,23 +21,23 @@ func OaiResponsesHandler(c *gin.Context, info *relaycommon.RelayInfo, resp *http defer service.CloseResponseBodyGracefully(resp) // read response body - var responsesResponse dto.OpenAIResponsesResponse + var billingMeta dto.ResponsesBillingMeta responseBody, err := io.ReadAll(resp.Body) if err != nil { return nil, types.NewOpenAIError(err, types.ErrorCodeReadResponseBodyFailed, http.StatusInternalServerError) } - err = common.Unmarshal(responseBody, &responsesResponse) + err = common.Unmarshal(responseBody, &billingMeta) if err != nil { return nil, types.NewOpenAIError(err, types.ErrorCodeBadResponseBody, http.StatusInternalServerError) } - if oaiError := responsesResponse.GetOpenAIError(); oaiError != nil && oaiError.Type != "" { + if oaiError := billingMeta.GetOpenAIError(); oaiError != nil && oaiError.Type != "" { return nil, types.WithOpenAIError(*oaiError, resp.StatusCode) } - if responsesResponse.HasImageGenerationCall() { + if billingMeta.HasImageGenerationCall() { c.Set("image_generation_call", true) - c.Set("image_generation_call_quality", responsesResponse.GetQuality()) - c.Set("image_generation_call_size", responsesResponse.GetSize()) + c.Set("image_generation_call_quality", billingMeta.GetQuality()) + c.Set("image_generation_call_size", billingMeta.GetSize()) } // 写入新的 response body @@ -45,22 +45,22 @@ func OaiResponsesHandler(c *gin.Context, info *relaycommon.RelayInfo, resp *http // compute usage usage := dto.Usage{} - if responsesResponse.Usage != nil { - usage.PromptTokens = responsesResponse.Usage.InputTokens - usage.CompletionTokens = responsesResponse.Usage.OutputTokens - usage.TotalTokens = responsesResponse.Usage.TotalTokens - if responsesResponse.Usage.InputTokensDetails != nil { - usage.PromptTokensDetails.CachedTokens = responsesResponse.Usage.InputTokensDetails.CachedTokens + if billingMeta.Usage != nil { + usage.PromptTokens = billingMeta.Usage.InputTokens + usage.CompletionTokens = billingMeta.Usage.OutputTokens + usage.TotalTokens = billingMeta.Usage.TotalTokens + if billingMeta.Usage.InputTokensDetails != nil { + usage.PromptTokensDetails.CachedTokens = billingMeta.Usage.InputTokensDetails.CachedTokens } } if info == nil || info.ResponsesUsageInfo == nil || info.ResponsesUsageInfo.BuiltInTools == nil { return &usage, nil } // 解析 Tools 用量 - for _, tool := range responsesResponse.Tools { - buildToolinfo, ok := info.ResponsesUsageInfo.BuiltInTools[common.Interface2String(tool["type"])] + for _, tool := range billingMeta.Tools { + buildToolinfo, ok := info.ResponsesUsageInfo.BuiltInTools[tool.Type] if !ok || buildToolinfo == nil { - logger.LogError(c, fmt.Sprintf("BuiltInTools not found for tool type: %v", tool["type"])) + logger.LogError(c, fmt.Sprintf("BuiltInTools not found for tool type: %v", tool.Type)) continue } buildToolinfo.CallCount++ @@ -82,13 +82,13 @@ func OaiResponsesStreamHandler(c *gin.Context, info *relaycommon.RelayInfo, resp helper.StreamScannerHandler(c, resp, info, func(data string, sr *helper.StreamResult) { // 检查当前数据是否包含 completed 状态和 usage 信息 - var streamResponse dto.ResponsesStreamResponse + var streamResponse dto.ResponsesBillingStreamResponse if err := common.UnmarshalJsonStr(data, &streamResponse); err != nil { logger.LogError(c, "failed to unmarshal stream response: "+err.Error()) sr.Error(err) return } - sendResponsesStreamData(c, streamResponse, data) + sendResponsesStreamData(c, streamResponse.Type, data) switch streamResponse.Type { case "response.completed": if streamResponse.Response != nil { diff --git a/relay/chat_completions_via_responses.go b/relay/chat_completions_via_responses.go index c47da1fab67c..9a959b321958 100644 --- a/relay/chat_completions_via_responses.go +++ b/relay/chat_completions_via_responses.go @@ -124,17 +124,19 @@ func chatCompletionsViaResponses(c *gin.Context, info *relaycommon.RelayInfo, ad return nil, types.NewError(err, types.ErrorCodeConvertRequestFailed, types.ErrOptionWithSkipRetry()) } - body, size, closer, err := relaycommon.NewOutboundJSONBody(jsonData) + body, err := relaycommon.NewReplayableOutboundJSONBody(jsonData) if err != nil { return nil, types.NewError(err, types.ErrorCodeConvertRequestFailed, types.ErrOptionWithSkipRetry()) } - defer closer.Close() jsonData = nil - info.UpstreamRequestBodySize = size + info.UpstreamRequestBodySize = body.Size() var requestBody io.Reader = body var httpResp *http.Response - resp, err := adaptor.DoRequest(c, info, requestBody) + resp, err := func() (any, error) { + defer closeReplayableOutboundBody(c, body, "outbound chat responses body") + return adaptor.DoRequest(c, info, requestBody) + }() if err != nil { return nil, types.NewOpenAIError(err, types.ErrorCodeDoRequestFailed, http.StatusInternalServerError) } diff --git a/relay/claude_handler.go b/relay/claude_handler.go index 527363205a1f..a55e8633be73 100644 --- a/relay/claude_handler.go +++ b/relay/claude_handler.go @@ -150,6 +150,7 @@ func ClaudeHelper(c *gin.Context, info *relaycommon.RelayInfo) (newAPIError *typ } var requestBody io.Reader + var outboundBody relaycommon.ReplayableBody if model_setting.GetGlobalSettings().PassThroughRequestEnabled || info.ChannelSetting.PassThroughBodyEnabled { storage, err := common.GetBodyStorage(c) if err != nil { @@ -183,19 +184,22 @@ func ClaudeHelper(c *gin.Context, info *relaycommon.RelayInfo) (newAPIError *typ } logger.LogDebug(c, "requestBody: %s", jsonData) - body, size, closer, err := relaycommon.NewOutboundJSONBody(jsonData) + body, err := relaycommon.NewReplayableOutboundJSONBody(jsonData) if err != nil { return types.NewError(err, types.ErrorCodeConvertRequestFailed, types.ErrOptionWithSkipRetry()) } - defer closer.Close() jsonData = nil - info.UpstreamRequestBodySize = size + info.UpstreamRequestBodySize = body.Size() + outboundBody = body requestBody = body } statusCodeMappingStr := c.GetString("status_code_mapping") var httpResp *http.Response - resp, err := adaptor.DoRequest(c, info, requestBody) + resp, err := func() (any, error) { + defer closeReplayableOutboundBody(c, outboundBody, "outbound claude body") + return adaptor.DoRequest(c, info, requestBody) + }() if err != nil { return types.NewOpenAIError(err, types.ErrorCodeDoRequestFailed, http.StatusInternalServerError) } diff --git a/relay/common/outbound_body.go b/relay/common/outbound_body.go index 94ef8dde1da6..2cda6ff5a1c3 100644 --- a/relay/common/outbound_body.go +++ b/relay/common/outbound_body.go @@ -6,6 +6,13 @@ import ( "github.com/QuantumNous/new-api/common" ) +type ReplayableBody interface { + io.Reader + Open() (io.ReadCloser, error) + Size() int64 + Close() error +} + // NewOutboundJSONBody wraps the already-marshaled upstream request body into a // BodyStorage. When disk cache is enabled and the payload exceeds the configured // threshold, the data is written to a temp file and the original []byte can be @@ -29,3 +36,7 @@ func NewOutboundJSONBody(data []byte) (body io.Reader, size int64, closer io.Clo } return common.ReaderOnly(storage), storage.Size(), storage, nil } + +func NewReplayableOutboundJSONBody(data []byte) (ReplayableBody, error) { + return common.CreateBodyStorage(data) +} diff --git a/relay/compatible_handler.go b/relay/compatible_handler.go index a68cfe730f60..3baf30f59428 100644 --- a/relay/compatible_handler.go +++ b/relay/compatible_handler.go @@ -93,6 +93,7 @@ func TextHelper(c *gin.Context, info *relaycommon.RelayInfo) (newAPIError *types } var requestBody io.Reader + var outboundBody relaycommon.ReplayableBody if passThroughGlobal || info.ChannelSetting.PassThroughBodyEnabled { storage, err := common.GetBodyStorage(c) @@ -175,18 +176,21 @@ func TextHelper(c *gin.Context, info *relaycommon.RelayInfo) (newAPIError *types logger.LogDebug(c, "text request body: %s", jsonData) - body, size, closer, err := relaycommon.NewOutboundJSONBody(jsonData) + body, err := relaycommon.NewReplayableOutboundJSONBody(jsonData) if err != nil { return types.NewError(err, types.ErrorCodeConvertRequestFailed, types.ErrOptionWithSkipRetry()) } - defer closer.Close() jsonData = nil - info.UpstreamRequestBodySize = size + info.UpstreamRequestBodySize = body.Size() + outboundBody = body requestBody = body } var httpResp *http.Response - resp, err := adaptor.DoRequest(c, info, requestBody) + resp, err := func() (any, error) { + defer closeReplayableOutboundBody(c, outboundBody, "outbound text body") + return adaptor.DoRequest(c, info, requestBody) + }() if err != nil { return types.NewOpenAIError(err, types.ErrorCodeDoRequestFailed, http.StatusInternalServerError) } diff --git a/relay/outbound_body.go b/relay/outbound_body.go new file mode 100644 index 000000000000..150fe270c36a --- /dev/null +++ b/relay/outbound_body.go @@ -0,0 +1,17 @@ +package relay + +import ( + "github.com/QuantumNous/new-api/logger" + relaycommon "github.com/QuantumNous/new-api/relay/common" + + "github.com/gin-gonic/gin" +) + +func closeReplayableOutboundBody(c *gin.Context, body relaycommon.ReplayableBody, label string) { + if body == nil { + return + } + if err := body.Close(); err != nil { + logger.LogError(c, "failed to close "+label+": "+err.Error()) + } +} diff --git a/relay/responses_handler.go b/relay/responses_handler.go index 010c38bba865..cbc1d55afaf2 100644 --- a/relay/responses_handler.go +++ b/relay/responses_handler.go @@ -71,6 +71,7 @@ func ResponsesHelper(c *gin.Context, info *relaycommon.RelayInfo) (newAPIError * } adaptor.Init(info) var requestBody io.Reader + var outboundBody relaycommon.ReplayableBody if model_setting.GetGlobalSettings().PassThroughRequestEnabled || info.ChannelSetting.PassThroughBodyEnabled { storage, err := common.GetBodyStorage(c) if err != nil { @@ -103,18 +104,21 @@ func ResponsesHelper(c *gin.Context, info *relaycommon.RelayInfo) (newAPIError * } logger.LogDebug(c, "requestBody: %s", jsonData) - body, size, closer, err := relaycommon.NewOutboundJSONBody(jsonData) + body, err := relaycommon.NewReplayableOutboundJSONBody(jsonData) if err != nil { return types.NewError(err, types.ErrorCodeConvertRequestFailed, types.ErrOptionWithSkipRetry()) } - defer closer.Close() jsonData = nil - info.UpstreamRequestBodySize = size + info.UpstreamRequestBodySize = body.Size() + outboundBody = body requestBody = body } var httpResp *http.Response - resp, err := adaptor.DoRequest(c, info, requestBody) + resp, err := func() (any, error) { + defer closeReplayableOutboundBody(c, outboundBody, "outbound responses body") + return adaptor.DoRequest(c, info, requestBody) + }() if err != nil { return types.NewOpenAIError(err, types.ErrorCodeDoRequestFailed, http.StatusInternalServerError) }