diff --git a/constant/context_key.go b/constant/context_key.go index b856bc3dda14..082dda39cb74 100644 --- a/constant/context_key.go +++ b/constant/context_key.go @@ -63,6 +63,11 @@ const ( // It is not returned to end users, but can be persisted into consume/error logs for debugging. ContextKeyAdminRejectReason ContextKey = "admin_reject_reason" + // ContextKeyClaudeResponsesStreamState holds the chat-to-responses stream + // state when routing /v1/responses through a Claude-family adaptor. It + // persists across streaming chunks and into the final flush. + ContextKeyClaudeResponsesStreamState ContextKey = "claude_responses_stream_state" + // ContextKeyLanguage stores the user's language preference for i18n ContextKeyLanguage ContextKey = "language" ContextKeyIsStream ContextKey = "is_stream" diff --git a/dto/openai_response.go b/dto/openai_response.go index 2de6014f4d05..dc8f9f1d319b 100644 --- a/dto/openai_response.go +++ b/dto/openai_response.go @@ -413,17 +413,20 @@ const ( // ResponsesStreamResponse 用于处理 /v1/responses 流式响应 type ResponsesStreamResponse struct { - Type string `json:"type"` - Response *OpenAIResponsesResponse `json:"response,omitempty"` - Delta string `json:"delta,omitempty"` - Item *ResponsesOutput `json:"item,omitempty"` + Type string `json:"type"` + ResponseID string `json:"response_id,omitempty"` + Response *OpenAIResponsesResponse `json:"response,omitempty"` + Delta string `json:"delta,omitempty"` + Text string `json:"text,omitempty"` + Arguments string `json:"arguments,omitempty"` + Item *ResponsesOutput `json:"item,omitempty"` // - response.function_call_arguments.delta // - response.function_call_arguments.done - 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"` + 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 any `json:"part,omitempty"` } // GetOpenAIError 从动态错误类型中提取OpenAIError结构 diff --git a/relay/channel/claude/adaptor.go b/relay/channel/claude/adaptor.go index b8e4a0366dd7..79e9e04e26c6 100644 --- a/relay/channel/claude/adaptor.go +++ b/relay/channel/claude/adaptor.go @@ -10,6 +10,7 @@ import ( "github.com/QuantumNous/new-api/dto" "github.com/QuantumNous/new-api/relay/channel" relaycommon "github.com/QuantumNous/new-api/relay/common" + "github.com/QuantumNous/new-api/service" "github.com/QuantumNous/new-api/service/relayconvert" "github.com/QuantumNous/new-api/setting/model_setting" "github.com/QuantumNous/new-api/types" @@ -113,8 +114,11 @@ func (a *Adaptor) ConvertEmbeddingRequest(c *gin.Context, info *relaycommon.Rela } func (a *Adaptor) ConvertOpenAIResponsesRequest(c *gin.Context, info *relaycommon.RelayInfo, request dto.OpenAIResponsesRequest) (any, error) { - // TODO implement me - return nil, errors.New("not implemented") + chatReq, err := service.ResponsesRequestToChatCompletionsRequest(&request) + if err != nil { + return nil, err + } + return a.ConvertOpenAIRequest(c, info, chatReq) } func (a *Adaptor) DoRequest(c *gin.Context, info *relaycommon.RelayInfo, requestBody io.Reader) (any, error) { diff --git a/relay/channel/claude/relay-claude.go b/relay/channel/claude/relay-claude.go index 8488cc9bb724..e87604c79ffe 100644 --- a/relay/channel/claude/relay-claude.go +++ b/relay/channel/claude/relay-claude.go @@ -12,6 +12,7 @@ import ( relaycommon "github.com/QuantumNous/new-api/relay/common" "github.com/QuantumNous/new-api/relay/helper" "github.com/QuantumNous/new-api/service" + openaicompat "github.com/QuantumNous/new-api/service/openaicompat" "github.com/QuantumNous/new-api/service/relayconvert" "github.com/QuantumNous/new-api/setting/model_setting" "github.com/QuantumNous/new-api/types" @@ -42,6 +43,15 @@ func ResponseClaude2OpenAI(claudeResponse *dto.ClaudeResponse) *dto.OpenAITextRe type ClaudeResponseInfo = relayconvert.ClaudeResponseInfo +func claudeResponsesStreamState(c *gin.Context, claudeInfo *ClaudeResponseInfo) *openaicompat.ChatToResponsesStreamState { + if state, ok := common.GetContextKeyType[*openaicompat.ChatToResponsesStreamState](c, constant.ContextKeyClaudeResponsesStreamState); ok && state != nil { + return state + } + state := openaicompat.NewChatToResponsesStreamState(claudeInfo.ResponseId, claudeInfo.Created, claudeInfo.Model) + common.SetContextKey(c, constant.ContextKeyClaudeResponsesStreamState, state) + return state +} + func cacheCreationTokensForOpenAIUsage(usage *dto.Usage) int { if usage == nil { return 0 @@ -126,6 +136,20 @@ func HandleStreamResponseData(c *gin.Context, info *relaycommon.RelayInfo, claud if err != nil { logger.LogError(c, "send_stream_response_failed: "+err.Error()) } + } else if info.RelayFormat == types.RelayFormatOpenAIResponses { + response := StreamResponseClaude2OpenAI(&claudeResponse) + if !FormatClaudeResponseInfo(&claudeResponse, response, claudeInfo) { + return nil + } + streamState := claudeResponsesStreamState(c, claudeInfo) + for _, event := range streamState.HandleChatChunk(response) { + jsonData, marshalErr := common.Marshal(event) + if marshalErr != nil { + logger.LogError(c, "send_stream_response_failed: "+marshalErr.Error()) + continue + } + helper.ResponseChunkData(c, event, string(jsonData)) + } } return nil } @@ -168,6 +192,17 @@ func HandleStreamFinalResponse(c *gin.Context, info *relaycommon.RelayInfo, clau } } helper.Done(c) + } else if info.RelayFormat == types.RelayFormatOpenAIResponses { + streamState := claudeResponsesStreamState(c, claudeInfo) + for _, event := range streamState.FinalEvents(claudeInfo.Usage) { + jsonData, err := common.Marshal(event) + if err != nil { + common.SysLog("send final response failed: " + err.Error()) + continue + } + helper.ResponseChunkData(c, event, string(jsonData)) + } + helper.Done(c) } } @@ -227,6 +262,17 @@ func HandleClaudeResponseData(c *gin.Context, info *relaycommon.RelayInfo, claud if err != nil { return types.NewError(err, types.ErrorCodeBadResponseBody) } + case types.RelayFormatOpenAIResponses: + openaiResponse := ResponseClaude2OpenAI(&claudeResponse) + openaiResponse.Usage = buildOpenAIStyleUsageFromClaudeUsage(claudeInfo.Usage) + responsesResp, _, convErr := service.ChatCompletionsResponseToResponsesResponse(openaiResponse, info.UpstreamModelName) + if convErr != nil { + return types.NewError(convErr, types.ErrorCodeBadResponseBody) + } + responseData, err = common.Marshal(responsesResp) + if err != nil { + return types.NewError(err, types.ErrorCodeBadResponseBody) + } case types.RelayFormatClaude: responseData = data } diff --git a/relay/responses_handler.go b/relay/responses_handler.go index 5fa23d099623..ca750fbdc96d 100644 --- a/relay/responses_handler.go +++ b/relay/responses_handler.go @@ -75,6 +75,24 @@ func ResponsesHelper(c *gin.Context, info *relaycommon.RelayInfo) (newAPIError * return types.NewError(err, types.ErrorCodeChannelModelMappedError, types.ErrOptionWithSkipRetry()) } + // Image generation models may not be supported via /v1/responses on + // upstream proxies. Convert to /v1/chat/completions and convert the + // response back to Responses format. + if !model_setting.GetGlobalSettings().PassThroughRequestEnabled && + !info.ChannelSetting.PassThroughBodyEnabled && + shouldResponsesUseChatCompletions(info) { + adaptor := GetAdaptor(info.ApiType) + if adaptor != nil { + adaptor.Init(info) + usage, newApiErr := responsesViaChatCompletions(c, info, adaptor, request) + if newApiErr != nil { + return newApiErr + } + service.PostTextConsumeQuota(c, info, usage, nil) + return nil + } + } + adaptor := GetAdaptor(info.ApiType) if adaptor == nil { return types.NewError(fmt.Errorf("invalid api type: %d", info.ApiType), types.ErrorCodeInvalidApiType, types.ErrOptionWithSkipRetry()) diff --git a/relay/responses_via_chat_completions.go b/relay/responses_via_chat_completions.go new file mode 100644 index 000000000000..65b6ccd3141b --- /dev/null +++ b/relay/responses_via_chat_completions.go @@ -0,0 +1,359 @@ +package relay + +import ( + "bufio" + "bytes" + "io" + "net/http" + "strings" + "time" + + "github.com/QuantumNous/new-api/common" + "github.com/QuantumNous/new-api/dto" + "github.com/QuantumNous/new-api/relay/channel" + relaycommon "github.com/QuantumNous/new-api/relay/common" + relayconstant "github.com/QuantumNous/new-api/relay/constant" + "github.com/QuantumNous/new-api/relay/helper" + "github.com/QuantumNous/new-api/service" + "github.com/QuantumNous/new-api/service/openaicompat" + "github.com/QuantumNous/new-api/types" + + "github.com/gin-gonic/gin" +) + +// responsesViaChatCompletions converts a Responses API request to a Chat +// Completions request, sends it upstream via /v1/chat/completions, and +// converts the response back to Responses format. This is the inverse of +// chatCompletionsViaResponses and is used for models that do not support the +// Responses API natively (e.g. image generation models on upstream proxies). +func responsesViaChatCompletions(c *gin.Context, info *relaycommon.RelayInfo, adaptor channel.Adaptor, request *dto.OpenAIResponsesRequest) (*dto.Usage, *types.NewAPIError) { + chatReq, err := service.ResponsesRequestToChatCompletionsRequest(request) + if err != nil { + return nil, types.NewErrorWithStatusCode(err, types.ErrorCodeConvertRequestFailed, http.StatusBadRequest, types.ErrOptionWithSkipRetry()) + } + info.AppendRequestConversion(types.RelayFormatOpenAI) + + savedRelayMode := info.RelayMode + savedRequestURLPath := info.RequestURLPath + defer func() { + info.RelayMode = savedRelayMode + info.RequestURLPath = savedRequestURLPath + }() + + info.RelayMode = relayconstant.RelayModeChatCompletions + info.RequestURLPath = "/v1/chat/completions" + + convertedRequest, err := adaptor.ConvertOpenAIRequest(c, info, chatReq) + if err != nil { + return nil, types.NewError(err, types.ErrorCodeConvertRequestFailed, types.ErrOptionWithSkipRetry()) + } + relaycommon.AppendRequestConversionFromRequest(info, convertedRequest) + + jsonData, err := common.Marshal(convertedRequest) + if err != nil { + return nil, types.NewError(err, types.ErrorCodeConvertRequestFailed, types.ErrOptionWithSkipRetry()) + } + + jsonData, err = relaycommon.RemoveDisabledFields(jsonData, info.ChannelOtherSettings, info.ChannelSetting.PassThroughBodyEnabled) + if err != nil { + return nil, types.NewError(err, types.ErrorCodeConvertRequestFailed, types.ErrOptionWithSkipRetry()) + } + + if len(info.ParamOverride) > 0 { + jsonData, err = relaycommon.ApplyParamOverrideWithRelayInfo(jsonData, info) + if err != nil { + return nil, newAPIErrorFromParamOverride(err) + } + } + + var requestBody io.Reader = bytes.NewBuffer(jsonData) + + resp, err := adaptor.DoRequest(c, info, requestBody) + if err != nil { + return nil, types.NewOpenAIError(err, types.ErrorCodeDoRequestFailed, http.StatusInternalServerError) + } + if resp == nil { + return nil, types.NewOpenAIError(nil, types.ErrorCodeBadResponse, http.StatusInternalServerError) + } + + statusCodeMappingStr := c.GetString("status_code_mapping") + + httpResp := resp.(*http.Response) + isStream := strings.HasPrefix(httpResp.Header.Get("Content-Type"), "text/event-stream") + if httpResp.StatusCode != http.StatusOK { + newApiErr := service.RelayErrorHandler(c.Request.Context(), httpResp, false) + service.ResetStatusCode(newApiErr, statusCodeMappingStr) + return nil, newApiErr + } + + // The upstream responded with a chat completions response. Convert it + // back to the Responses API format before returning to the caller. + var usage *dto.Usage + var newApiErr *types.NewAPIError + if isStream && info.IsStream { + // Caller requested streaming: translate SSE chunks incrementally + usage, newApiErr = oaiChatStreamToResponsesStreamHandler(c, info, httpResp) + } else if isStream { + // Upstream is streaming but caller did not request it: buffer into JSON + usage, newApiErr = oaiChatStreamToResponsesHandler(c, info, httpResp) + } else { + usage, newApiErr = oaiChatToResponsesHandler(c, info, httpResp) + } + if newApiErr != nil { + service.ResetStatusCode(newApiErr, statusCodeMappingStr) + return nil, newApiErr + } + return usage, nil +} + +// oaiChatToResponsesHandler reads a non-streaming Chat Completions response +// and re-emits it as a Responses API response. +func oaiChatToResponsesHandler(c *gin.Context, info *relaycommon.RelayInfo, resp *http.Response) (*dto.Usage, *types.NewAPIError) { + if resp == nil || resp.Body == nil { + return nil, types.NewOpenAIError(nil, types.ErrorCodeBadResponse, http.StatusInternalServerError) + } + defer service.CloseResponseBodyGracefully(resp) + + body, err := io.ReadAll(resp.Body) + if err != nil { + return nil, types.NewOpenAIError(err, types.ErrorCodeReadResponseBodyFailed, http.StatusInternalServerError) + } + + var chatResp dto.OpenAITextResponse + if err := common.Unmarshal(body, &chatResp); err != nil { + return nil, types.NewOpenAIError(err, types.ErrorCodeBadResponseBody, http.StatusInternalServerError) + } + + responsesResp, _, err := service.ChatCompletionsResponseToResponsesResponse(&chatResp, info.UpstreamModelName) + if err != nil { + return nil, types.NewOpenAIError(err, types.ErrorCodeBadResponseBody, http.StatusInternalServerError) + } + + usage := &dto.Usage{ + PromptTokens: chatResp.Usage.PromptTokens, + CompletionTokens: chatResp.Usage.CompletionTokens, + TotalTokens: chatResp.Usage.TotalTokens, + } + if usage.TotalTokens == 0 { + usage.TotalTokens = usage.PromptTokens + usage.CompletionTokens + } + usage.PromptTokensDetails = chatResp.Usage.PromptTokensDetails + usage.CompletionTokenDetails = chatResp.Usage.CompletionTokenDetails + + responseBody, err := common.Marshal(responsesResp) + if err != nil { + return nil, types.NewOpenAIError(err, types.ErrorCodeJsonMarshalFailed, http.StatusInternalServerError) + } + + // Write the Responses API JSON back to the client using the original + // HTTP response wrapper so that content-type and status are set correctly. + service.IOCopyBytesGracefully(c, resp, responseBody) + return usage, nil +} + +// isImageGenerationModelForResponses checks if a model is an image generation +// model that should be routed through chat completions when called via the +// Responses API. Uses the upstream model name (after mapping) to catch mapped +// names like gpt-image-1.5-all. +func isImageGenerationModelForResponses(modelName string) bool { + return common.IsImageGenerationModel(modelName) +} + +// oaiChatStreamToResponsesHandler reads a streaming (SSE) Chat Completions +// response, accumulates all chunks into a single message, then emits it as a +// non-streaming Responses API JSON response. This handles upstreams that always +// return SSE even when stream was not requested (common with image generation +// proxies). +func oaiChatStreamToResponsesHandler(c *gin.Context, info *relaycommon.RelayInfo, resp *http.Response) (*dto.Usage, *types.NewAPIError) { + if resp == nil || resp.Body == nil { + return nil, types.NewOpenAIError(nil, types.ErrorCodeBadResponse, http.StatusInternalServerError) + } + defer service.CloseResponseBodyGracefully(resp) + + var ( + contentBuilder strings.Builder + toolCalls []dto.ToolCallResponse + model string + usage = &dto.Usage{} + finishReason = "stop" + ) + + scanner := bufio.NewScanner(resp.Body) + scanner.Buffer(make([]byte, 0, 1024*1024), 1024*1024) + + for scanner.Scan() { + line := scanner.Text() + if !strings.HasPrefix(line, "data:") { + continue + } + data := strings.TrimSpace(strings.TrimPrefix(line, "data:")) + if data == "[DONE]" || data == "" { + continue + } + + var chunk dto.ChatCompletionsStreamResponse + if err := common.UnmarshalJsonStr(data, &chunk); err != nil { + continue + } + + if chunk.Model != "" { + model = chunk.Model + } + + if chunk.Usage != nil { + usage = chunk.Usage + } + + for _, choice := range chunk.Choices { + if choice.Delta.Content != nil { + contentBuilder.WriteString(*choice.Delta.Content) + } + // Accumulate streaming tool calls by index + for _, tc := range choice.Delta.ToolCalls { + idx := 0 + if tc.Index != nil { + idx = *tc.Index + } + for len(toolCalls) <= idx { + toolCalls = append(toolCalls, dto.ToolCallResponse{Type: "function"}) + } + if tc.ID != "" { + toolCalls[idx].ID = tc.ID + } + if tc.Function.Name != "" { + toolCalls[idx].Function.Name = tc.Function.Name + } + toolCalls[idx].Function.Arguments += tc.Function.Arguments + } + if choice.FinishReason != nil { + finishReason = *choice.FinishReason + } + } + } + if err := scanner.Err(); err != nil { + return nil, types.NewOpenAIError(err, types.ErrorCodeReadResponseBodyFailed, http.StatusInternalServerError) + } + + if model == "" { + model = info.UpstreamModelName + } + + // Build a synthetic OpenAITextResponse from accumulated chunks + msg := dto.Message{Role: "assistant", Content: contentBuilder.String()} + if len(toolCalls) > 0 { + msg.SetToolCalls(toolCalls) + } + chatResp := &dto.OpenAITextResponse{ + Id: "chatcmpl-" + common.GetUUID(), + Object: "chat.completion", + Model: model, + Choices: []dto.OpenAITextResponseChoice{{ + Index: 0, + Message: msg, + FinishReason: finishReason, + }}, + Usage: *usage, + } + + if usage.TotalTokens == 0 { + text := contentBuilder.String() + usage = service.ResponseText2Usage(c, text, info.UpstreamModelName, info.GetEstimatePromptTokens()) + chatResp.Usage = *usage + } + + responsesResp, _, err := service.ChatCompletionsResponseToResponsesResponse(chatResp, info.UpstreamModelName) + if err != nil { + return nil, types.NewOpenAIError(err, types.ErrorCodeBadResponseBody, http.StatusInternalServerError) + } + + responseBody, err := common.Marshal(responsesResp) + if err != nil { + return nil, types.NewOpenAIError(err, types.ErrorCodeJsonMarshalFailed, http.StatusInternalServerError) + } + + service.IOCopyBytesGracefully(c, resp, responseBody) + return usage, nil +} + +// oaiChatStreamToResponsesStreamHandler reads a streaming Chat Completions +// response and translates each chunk into Responses API SSE events, emitting +// them incrementally. Use this when the caller requested stream=true. +func oaiChatStreamToResponsesStreamHandler(c *gin.Context, info *relaycommon.RelayInfo, resp *http.Response) (*dto.Usage, *types.NewAPIError) { + if resp == nil || resp.Body == nil { + return nil, types.NewOpenAIError(nil, types.ErrorCodeBadResponse, http.StatusInternalServerError) + } + defer service.CloseResponseBodyGracefully(resp) + + helper.SetEventStreamHeaders(c) + + state := openaicompat.NewChatToResponsesStreamState( + "resp_"+common.GetUUID(), + time.Now().Unix(), + info.UpstreamModelName, + ) + + var usage *dto.Usage + + scanner := bufio.NewScanner(resp.Body) + scanner.Buffer(make([]byte, 0, 1024*1024), 1024*1024) + + for scanner.Scan() { + line := scanner.Text() + if !strings.HasPrefix(line, "data:") { + continue + } + data := strings.TrimSpace(strings.TrimPrefix(line, "data:")) + if data == "[DONE]" || data == "" { + continue + } + + var chunk dto.ChatCompletionsStreamResponse + if err := common.UnmarshalJsonStr(data, &chunk); err != nil { + continue + } + + if u := state.HandleUsageChunk(&chunk); u != nil { + usage = u + } + + events := state.HandleChatChunk(&chunk) + for _, event := range events { + jsonData, err := common.Marshal(event) + if err != nil { + continue + } + helper.ResponseChunkData(c, event, string(jsonData)) + } + } + if err := scanner.Err(); err != nil { + return nil, types.NewOpenAIError(err, types.ErrorCodeReadResponseBodyFailed, http.StatusInternalServerError) + } + + if usage == nil { + text := state.OutputText.String() + usage = service.ResponseText2Usage(c, text, info.UpstreamModelName, info.GetEstimatePromptTokens()) + } + + finalEvents := state.FinalEvents(usage) + for _, event := range finalEvents { + jsonData, err := common.Marshal(event) + if err != nil { + continue + } + helper.ResponseChunkData(c, event, string(jsonData)) + } + + return usage, nil +} + +// shouldResponsesUseChatCompletions returns true when a /v1/responses request +// should be internally converted to /v1/chat/completions. Currently this +// applies to image generation models whose upstream providers do not support +// the Responses API for image generation. +func shouldResponsesUseChatCompletions(info *relaycommon.RelayInfo) bool { + modelToCheck := info.UpstreamModelName + if modelToCheck == "" { + modelToCheck = info.OriginModelName + } + return isImageGenerationModelForResponses(modelToCheck) +} diff --git a/service/openaicompat/chat_stream_to_responses_stream.go b/service/openaicompat/chat_stream_to_responses_stream.go new file mode 100644 index 000000000000..5c2844441cd0 --- /dev/null +++ b/service/openaicompat/chat_stream_to_responses_stream.go @@ -0,0 +1,521 @@ +package openaicompat + +import ( + "encoding/json" + "strings" + + "github.com/QuantumNous/new-api/common" + "github.com/QuantumNous/new-api/dto" +) + +// ChatToResponsesStreamState tracks state for converting a chat completions +// stream into Responses API SSE events. It is adaptor-agnostic: each adaptor +// converts its native chunk into a *dto.ChatCompletionsStreamResponse and +// feeds it to HandleChatChunk; the state machine emits the correct Responses +// API events in return. +type ChatToResponsesStreamState struct { + ResponseID string + CreatedAt int64 + Model string + SentCreated bool + SentInProgress bool + + MessageItemID string + MessageOutputIndex int + MessageContentIndex int + MessageItemAdded bool + MessageContentAdded bool + + NextOutputIndex int + + ReasoningOutputIndex int + ReasoningText strings.Builder + OutputText strings.Builder + ToolCallArgs map[string]string + ToolCallName map[string]string + ToolCallSent map[string]bool + ToolCallOrder []string + ToolCallOutIndex map[string]int +} + +func NewChatToResponsesStreamState(responseID string, createdAt int64, model string) *ChatToResponsesStreamState { + return &ChatToResponsesStreamState{ + ResponseID: normalizeResponsesID(responseID), + CreatedAt: createdAt, + Model: model, + MessageOutputIndex: -1, + MessageContentIndex: 0, + NextOutputIndex: 0, + ReasoningOutputIndex: -1, + ToolCallArgs: make(map[string]string), + ToolCallName: make(map[string]string), + ToolCallSent: make(map[string]bool), + ToolCallOutIndex: make(map[string]int), + } +} + +// HandleChatChunk converts one chat completions stream chunk into zero or more +// Responses API events. +func (s *ChatToResponsesStreamState) HandleChatChunk(chunk *dto.ChatCompletionsStreamResponse) []dto.ResponsesStreamResponse { + if chunk == nil || len(chunk.Choices) == 0 { + return nil + } + + if chunk.Model != "" { + s.Model = chunk.Model + } + if s.CreatedAt == 0 && chunk.Created != 0 { + s.CreatedAt = chunk.Created + } + + events := s.baseEvents() + + delta := chunk.Choices[0].Delta + + // Text content + if delta.Content != nil { + content := *delta.Content + if content != "" { + events = append(events, s.ensureMessageItemEvents()...) + events = append(events, s.ensureContentPartEvents()...) + s.OutputText.WriteString(content) + events = append(events, s.outputTextDeltaEvent(content)) + } + } + + // Reasoning content (for models that emit reasoning_content) + reasoningContent := delta.GetReasoningContent() + if reasoningContent != "" { + if s.ReasoningOutputIndex < 0 { + s.ReasoningOutputIndex = s.NextOutputIndex + s.NextOutputIndex++ + } + outIndex := s.ReasoningOutputIndex + summaryIndex := 0 + s.ReasoningText.WriteString(reasoningContent) + events = append(events, dto.ResponsesStreamResponse{ + Type: "response.reasoning_summary_text.delta", + ResponseID: s.ResponseID, + ItemID: "rs_" + strings.TrimPrefix(s.ResponseID, "resp_"), + OutputIndex: &outIndex, + SummaryIndex: &summaryIndex, + Delta: reasoningContent, + }) + } + + // Tool calls + if len(delta.ToolCalls) > 0 { + for _, call := range delta.ToolCalls { + callID := strings.TrimSpace(call.ID) + if callID == "" { + // For subsequent argument deltas, use the last known call ID + if call.Index != nil && *call.Index < len(s.ToolCallOrder) { + callID = s.ToolCallOrder[*call.Index] + } else if len(s.ToolCallOrder) > 0 { + callID = s.ToolCallOrder[len(s.ToolCallOrder)-1] + } else { + callID = "call_" + common.GetUUID() + } + } + if call.Function.Name != "" { + s.ToolCallName[callID] = call.Function.Name + } + if !s.ToolCallSent[callID] { + s.ToolCallSent[callID] = true + s.ToolCallOrder = append(s.ToolCallOrder, callID) + outIndex := s.allocOutputIndex(callID) + events = append(events, s.toolItemAddedEvent(callID, outIndex)) + } + + args := call.Function.Arguments + if args == "" { + continue + } + s.ToolCallArgs[callID] = s.ToolCallArgs[callID] + args + + events = append(events, dto.ResponsesStreamResponse{ + Type: "response.function_call_arguments.delta", + ResponseID: s.ResponseID, + ItemID: callID, + OutputIndex: s.outputIndexPtr(callID), + Delta: args, + }) + } + } + + return events +} + +// HandleUsageChunk processes a usage-only chunk (no choices). +func (s *ChatToResponsesStreamState) HandleUsageChunk(chunk *dto.ChatCompletionsStreamResponse) *dto.Usage { + if chunk == nil || chunk.Usage == nil { + return nil + } + usage := &dto.Usage{ + PromptTokens: chunk.Usage.PromptTokens, + CompletionTokens: chunk.Usage.CompletionTokens, + TotalTokens: chunk.Usage.TotalTokens, + InputTokens: chunk.Usage.PromptTokens, + OutputTokens: chunk.Usage.CompletionTokens, + } + usage.PromptTokensDetails = chunk.Usage.PromptTokensDetails + usage.CompletionTokenDetails = chunk.Usage.CompletionTokenDetails + return usage +} + +// FinalEvents emits the closing events: content done, tool calls done, and +// response.completed. +func (s *ChatToResponsesStreamState) FinalEvents(usage *dto.Usage) []dto.ResponsesStreamResponse { + events := s.baseEvents() + + // Finalize message item + if s.MessageItemAdded { + text := s.OutputText.String() + if s.MessageContentAdded { + events = append(events, s.outputTextDoneEvent(text)) + events = append(events, s.contentPartDoneEvent(text)) + } + events = append(events, s.messageItemDoneEvent(text)) + } + + // Finalize reasoning content + if s.ReasoningOutputIndex >= 0 { + outIndex := s.ReasoningOutputIndex + summaryIndex := 0 + events = append(events, dto.ResponsesStreamResponse{ + Type: "response.reasoning_summary_text.done", + ResponseID: s.ResponseID, + ItemID: "rs_" + strings.TrimPrefix(s.ResponseID, "resp_"), + OutputIndex: &outIndex, + SummaryIndex: &summaryIndex, + Text: s.ReasoningText.String(), + }) + } + + // Finalize tool calls + for _, callID := range s.ToolCallOrder { + outIndex := s.outputIndexPtr(callID) + args := s.ToolCallArgs[callID] + if args != "" { + events = append(events, dto.ResponsesStreamResponse{ + Type: "response.function_call_arguments.done", + ResponseID: s.ResponseID, + ItemID: callID, + OutputIndex: outIndex, + Arguments: args, + }) + } + events = append(events, dto.ResponsesStreamResponse{ + Type: "response.output_item.done", + ResponseID: s.ResponseID, + ItemID: callID, + OutputIndex: outIndex, + Item: &dto.ResponsesOutput{ + Type: "function_call", + ID: callID, + Status: "completed", + CallId: callID, + Name: s.ToolCallName[callID], + Arguments: json.RawMessage(args), + }, + }) + } + + // Build final output and usage + output := s.buildFinalOutput() + finalUsage := s.buildFinalUsage(usage) + + resp := &dto.OpenAIResponsesResponse{ + ID: s.ResponseID, + Object: "response", + CreatedAt: int(s.CreatedAt), + Status: []byte(`"completed"`), + Model: s.Model, + Output: output, + Usage: finalUsage, + } + events = append(events, dto.ResponsesStreamResponse{ + Type: "response.completed", + ResponseID: s.ResponseID, + Response: resp, + }) + + return events +} + +func (s *ChatToResponsesStreamState) baseEvents() []dto.ResponsesStreamResponse { + events := make([]dto.ResponsesStreamResponse, 0, 2) + if !s.SentCreated { + events = append(events, s.createdEvent()) + s.SentCreated = true + } + if !s.SentInProgress { + events = append(events, s.inProgressEvent()) + s.SentInProgress = true + } + return events +} + +func (s *ChatToResponsesStreamState) createdEvent() dto.ResponsesStreamResponse { + resp := &dto.OpenAIResponsesResponse{ + ID: s.ResponseID, + Object: "response", + CreatedAt: int(s.CreatedAt), + Status: []byte(`"in_progress"`), + Model: s.Model, + Output: []dto.ResponsesOutput{}, + } + return dto.ResponsesStreamResponse{ + Type: "response.created", + ResponseID: s.ResponseID, + Response: resp, + } +} + +func (s *ChatToResponsesStreamState) inProgressEvent() dto.ResponsesStreamResponse { + resp := &dto.OpenAIResponsesResponse{ + ID: s.ResponseID, + Object: "response", + CreatedAt: int(s.CreatedAt), + Status: []byte(`"in_progress"`), + Model: s.Model, + Output: []dto.ResponsesOutput{}, + } + return dto.ResponsesStreamResponse{ + Type: "response.in_progress", + ResponseID: s.ResponseID, + Response: resp, + } +} + +func (s *ChatToResponsesStreamState) ensureMessageItemEvents() []dto.ResponsesStreamResponse { + if s.MessageItemAdded { + return nil + } + s.MessageItemAdded = true + if s.MessageOutputIndex < 0 { + s.MessageOutputIndex = s.NextOutputIndex + s.NextOutputIndex++ + } + if s.MessageItemID == "" { + s.MessageItemID = "msg_" + common.GetUUID() + } + outIndex := s.MessageOutputIndex + return []dto.ResponsesStreamResponse{ + { + Type: "response.output_item.added", + ResponseID: s.ResponseID, + OutputIndex: &outIndex, + Item: &dto.ResponsesOutput{ + ID: s.MessageItemID, + Type: "message", + Status: "in_progress", + Role: "assistant", + Content: []dto.ResponsesOutputContent{}, + }, + }, + } +} + +func (s *ChatToResponsesStreamState) ensureContentPartEvents() []dto.ResponsesStreamResponse { + if s.MessageContentAdded { + return nil + } + s.MessageContentAdded = true + outIndex := s.MessageOutputIndex + contentIndex := s.MessageContentIndex + part := dto.ResponsesOutputContent{ + Type: "output_text", + Text: "", + Annotations: []interface{}{}, + } + return []dto.ResponsesStreamResponse{ + { + Type: "response.content_part.added", + ResponseID: s.ResponseID, + ItemID: s.MessageItemID, + OutputIndex: &outIndex, + ContentIndex: &contentIndex, + Part: &part, + }, + } +} + +func (s *ChatToResponsesStreamState) outputTextDeltaEvent(delta string) dto.ResponsesStreamResponse { + outIndex := s.MessageOutputIndex + contentIndex := s.MessageContentIndex + return dto.ResponsesStreamResponse{ + Type: "response.output_text.delta", + ResponseID: s.ResponseID, + ItemID: s.MessageItemID, + OutputIndex: &outIndex, + ContentIndex: &contentIndex, + Delta: delta, + } +} + +func (s *ChatToResponsesStreamState) outputTextDoneEvent(text string) dto.ResponsesStreamResponse { + outIndex := s.MessageOutputIndex + contentIndex := s.MessageContentIndex + return dto.ResponsesStreamResponse{ + Type: "response.output_text.done", + ResponseID: s.ResponseID, + ItemID: s.MessageItemID, + OutputIndex: &outIndex, + ContentIndex: &contentIndex, + Text: text, + } +} + +func (s *ChatToResponsesStreamState) contentPartDoneEvent(text string) dto.ResponsesStreamResponse { + outIndex := s.MessageOutputIndex + contentIndex := s.MessageContentIndex + part := dto.ResponsesOutputContent{ + Type: "output_text", + Text: text, + Annotations: []interface{}{}, + } + return dto.ResponsesStreamResponse{ + Type: "response.content_part.done", + ResponseID: s.ResponseID, + ItemID: s.MessageItemID, + OutputIndex: &outIndex, + ContentIndex: &contentIndex, + Part: &part, + } +} + +func (s *ChatToResponsesStreamState) messageItemDoneEvent(text string) dto.ResponsesStreamResponse { + outIndex := s.MessageOutputIndex + item := dto.ResponsesOutput{ + ID: s.MessageItemID, + Type: "message", + Status: "completed", + Role: "assistant", + Content: []dto.ResponsesOutputContent{ + { + Type: "output_text", + Text: text, + Annotations: []interface{}{}, + }, + }, + } + return dto.ResponsesStreamResponse{ + Type: "response.output_item.done", + ResponseID: s.ResponseID, + ItemID: s.MessageItemID, + OutputIndex: &outIndex, + Item: &item, + } +} + +func (s *ChatToResponsesStreamState) toolItemAddedEvent(callID string, outIndex int) dto.ResponsesStreamResponse { + item := dto.ResponsesOutput{ + Type: "function_call", + ID: callID, + Status: "in_progress", + CallId: callID, + Name: s.ToolCallName[callID], + } + return dto.ResponsesStreamResponse{ + Type: "response.output_item.added", + ResponseID: s.ResponseID, + ItemID: callID, + OutputIndex: &outIndex, + Item: &item, + } +} + +func (s *ChatToResponsesStreamState) allocOutputIndex(callID string) int { + if idx, ok := s.ToolCallOutIndex[callID]; ok { + return idx + } + idx := s.NextOutputIndex + s.NextOutputIndex++ + s.ToolCallOutIndex[callID] = idx + return idx +} + +func (s *ChatToResponsesStreamState) outputIndexPtr(callID string) *int { + idx, ok := s.ToolCallOutIndex[callID] + if !ok { + return nil + } + return &idx +} + +func (s *ChatToResponsesStreamState) buildFinalOutput() []dto.ResponsesOutput { + itemsByIndex := make(map[int]dto.ResponsesOutput) + if s.MessageItemAdded { + text := s.OutputText.String() + itemsByIndex[s.MessageOutputIndex] = dto.ResponsesOutput{ + ID: s.MessageItemID, + Type: "message", + Status: "completed", + Role: "assistant", + Content: []dto.ResponsesOutputContent{ + { + Type: "output_text", + Text: text, + Annotations: []interface{}{}, + }, + }, + } + } + for _, callID := range s.ToolCallOrder { + idx, ok := s.ToolCallOutIndex[callID] + if !ok { + continue + } + itemsByIndex[idx] = dto.ResponsesOutput{ + Type: "function_call", + ID: callID, + Status: "completed", + CallId: callID, + Name: s.ToolCallName[callID], + Arguments: json.RawMessage(s.ToolCallArgs[callID]), + } + } + output := make([]dto.ResponsesOutput, 0, len(itemsByIndex)) + for i := 0; i < s.NextOutputIndex; i++ { + if item, ok := itemsByIndex[i]; ok { + output = append(output, item) + } + } + return output +} + +func (s *ChatToResponsesStreamState) buildFinalUsage(usage *dto.Usage) *dto.Usage { + if usage == nil { + return &dto.Usage{} + } + final := &dto.Usage{} + *final = *usage + if final.InputTokens == 0 { + final.InputTokens = final.PromptTokens + } + if final.OutputTokens == 0 { + final.OutputTokens = final.CompletionTokens + } + if final.TotalTokens == 0 { + final.TotalTokens = final.PromptTokens + final.CompletionTokens + } + return final +} + +func normalizeResponsesID(id string) string { + id = strings.TrimSpace(id) + if id == "" { + return "resp_" + common.GetUUID() + } + if strings.HasPrefix(id, "resp_") { + return id + } + if strings.HasPrefix(id, "chatcmpl-") { + return "resp_" + strings.TrimPrefix(id, "chatcmpl-") + } + if strings.HasPrefix(id, "chatcmpl_") { + return "resp_" + strings.TrimPrefix(id, "chatcmpl_") + } + return "resp_" + id +}