Commit a043eef5 by CaIon

feat: implement Gemini to OpenAI chat stream conversion with state management and terminal handling

parent f3ab2cff
......@@ -12,6 +12,7 @@ import (
relaycommon "github.com/QuantumNous/new-api/relay/common"
"github.com/QuantumNous/new-api/relaykit/types"
"github.com/gin-gonic/gin"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
......@@ -86,6 +87,60 @@ func TestOaiResponsesToChatStreamHandlerConvertsSSEOrderAndUsage(t *testing.T) {
)
}
func TestOaiResponsesToChatStreamHandlerConvertsClaudeSSETerminalsAndUsage(t *testing.T) {
oldMode := gin.Mode()
gin.SetMode(gin.TestMode)
t.Cleanup(func() { gin.SetMode(oldMode) })
oldTimeout := constant.StreamingTimeout
constant.StreamingTimeout = 30
t.Cleanup(func() { constant.StreamingTimeout = oldTimeout })
body := strings.Join([]string{
`data: {"type":"response.created","response":{"id":"resp_1","model":"gpt-test","created_at":1710000000}}`,
`data: {"type":"response.output_text.delta","delta":"hello"}`,
`data: {"type":"response.completed","response":{"status":"completed","usage":{"input_tokens":2,"output_tokens":3,"total_tokens":5}}}`,
`data: [DONE]`,
``,
}, "\n")
c, recorder, resp, info := newResponsesChatTestContext(t, body, true)
info.RelayFormat = types.RelayFormatClaude
usage, err := OaiResponsesToChatStreamHandler(c, info, resp)
require.Nil(t, err)
require.NotNil(t, usage)
assert.Equal(t, 2, usage.PromptTokens)
assert.Equal(t, 3, usage.CompletionTokens)
assert.Equal(t, 5, usage.TotalTokens)
got := recorder.Body.String()
assert.Equal(t, "text/event-stream", recorder.Header().Get("Content-Type"))
assert.Equal(t, 1, strings.Count(got, "event: message_start\n"))
assert.Equal(t, 1, strings.Count(got, "event: content_block_stop\n"))
assert.Equal(t, 1, strings.Count(got, "event: message_delta\n"))
assert.Equal(t, 1, strings.Count(got, "event: message_stop\n"))
messageDeltaFrame := ""
for _, frame := range strings.Split(got, "\n\n") {
if strings.HasPrefix(frame, "event: message_delta\n") {
messageDeltaFrame = frame
break
}
}
require.NotEmpty(t, messageDeltaFrame)
assert.Contains(t, messageDeltaFrame, `"type":"message_delta"`)
assert.Contains(t, messageDeltaFrame, `"stop_reason":"end_turn"`)
assert.Contains(t, messageDeltaFrame, `"input_tokens":2`)
assert.Contains(t, messageDeltaFrame, `"output_tokens":3`)
requireOrderedSubstrings(t, got,
"event: message_start\n",
"event: content_block_stop\n",
"event: message_delta\n",
"event: message_stop\n",
)
}
func TestOaiResponsesToChatBufferedStreamHandlerReturnsJSONFromSSE(t *testing.T) {
oldMode := gin.Mode()
gin.SetMode(gin.TestMode)
......
......@@ -282,6 +282,107 @@ func StreamResponseGeminiChat2OpenAI(geminiResponse *dto.GeminiChatResponse) (*d
return &response, isStop
}
type GeminiToChatStreamState struct {
id string
created int64
sawToolCall bool
finishEmitted bool
latestUsage *dto.Usage
}
func NewGeminiToChatStreamState(id string, created int64) *GeminiToChatStreamState {
id = strings.TrimSpace(id)
if id == "" {
id = fmt.Sprintf("chatcmpl-%s", kitutil.GetUUID())
}
if created == 0 {
created = kitutil.GetTimestamp()
}
return &GeminiToChatStreamState{id: id, created: created}
}
func (s *GeminiToChatStreamState) ConvertChunk(geminiResponse *dto.GeminiChatResponse, model string, usage *dto.Usage) []*dto.ChatCompletionsStreamResponse {
if s == nil || geminiResponse == nil {
return nil
}
hasNonStopFinish := false
for _, candidate := range geminiResponse.Candidates {
if candidate.FinishReason != nil && *candidate.FinishReason != "" && *candidate.FinishReason != "STOP" {
hasNonStopFinish = true
break
}
}
response, isStop := StreamResponseGeminiChat2OpenAI(geminiResponse)
if response == nil {
return nil
}
response.Id = s.id
response.Created = s.created
response.Model = model
response.Usage = usage
if response.IsToolCall() {
s.sawToolCall = true
if !hasNonStopFinish {
for i := range response.Choices {
if response.Choices[i].FinishReason != nil && *response.Choices[i].FinishReason == types.FinishReasonToolCalls {
response.Choices[i].FinishReason = nil
}
}
}
}
if usage != nil {
s.latestUsage = usage
}
for _, choice := range response.Choices {
if choice.FinishReason != nil && *choice.FinishReason != "" {
s.finishEmitted = true
break
}
}
responses := []*dto.ChatCompletionsStreamResponse{response}
if isStop && !s.finishEmitted {
responses = append(responses, s.terminalChunk(model))
}
return responses
}
func (s *GeminiToChatStreamState) Finalize(model string) []*dto.ChatCompletionsStreamResponse {
if s == nil || s.finishEmitted {
return nil
}
return []*dto.ChatCompletionsStreamResponse{s.terminalChunk(model)}
}
func (s *GeminiToChatStreamState) Usage() *dto.Usage {
if s == nil {
return nil
}
return s.latestUsage
}
func (s *GeminiToChatStreamState) terminalChunk(model string) *dto.ChatCompletionsStreamResponse {
finishReason := types.FinishReasonStop
if s.sawToolCall {
finishReason = types.FinishReasonToolCalls
}
s.finishEmitted = true
return &dto.ChatCompletionsStreamResponse{
Id: s.id,
Object: "chat.completion.chunk",
Created: s.created,
Model: model,
Choices: []dto.ChatCompletionsStreamResponseChoice{
{
Delta: dto.ChatCompletionsStreamResponseChoiceDelta{},
FinishReason: &finishReason,
},
},
Usage: s.latestUsage,
}
}
func geminiResponseToolCall(item *dto.GeminiPart) *dto.ToolCallResponse {
argsBytes, err := kitutil.Marshal(item.FunctionCall.Arguments)
if err != nil {
......
......@@ -17,6 +17,24 @@ func generateStopBlock(index int) *dto.ClaudeResponse {
}
}
func stopOpenBlocks(state *convmeta.ClaudeConvertInfo) []*dto.ClaudeResponse {
if state == nil {
return nil
}
switch state.LastMessagesType {
case convmeta.LastMessageTypeText, convmeta.LastMessageTypeThinking:
return []*dto.ClaudeResponse{generateStopBlock(state.Index)}
case convmeta.LastMessageTypeTools:
responses := make([]*dto.ClaudeResponse, 0, state.ToolCallMaxIndexOffset+1)
for offset := 0; offset <= state.ToolCallMaxIndexOffset; offset++ {
responses = append(responses, generateStopBlock(state.ToolCallBaseIndex+offset))
}
return responses
default:
return nil
}
}
func buildClaudeUsageFromOpenAIUsage(oaiUsage *dto.Usage) *dto.ClaudeUsage {
if oaiUsage == nil {
return nil
......@@ -89,16 +107,8 @@ func StreamResponseOpenAI2Claude(openAIResponse *dto.ChatCompletionsStreamRespon
// For text/thinking, there is at most one open block at state.Index.
// For tools, OpenAI tool_calls can stream multiple parallel tool_use blocks (indexed from 0),
// so we may have multiple open blocks and must stop each one explicitly.
stopOpenBlocks := func() {
switch state.LastMessagesType {
case convmeta.LastMessageTypeText, convmeta.LastMessageTypeThinking:
claudeResponses = append(claudeResponses, generateStopBlock(state.Index))
case convmeta.LastMessageTypeTools:
base := state.ToolCallBaseIndex
for offset := 0; offset <= state.ToolCallMaxIndexOffset; offset++ {
claudeResponses = append(claudeResponses, generateStopBlock(base+offset))
}
}
appendStopOpenBlocks := func() {
claudeResponses = append(claudeResponses, stopOpenBlocks(state)...)
}
// stopOpenBlocksAndAdvance closes the currently open block(s) and advances the content block index
// to the next available slot for subsequent content_block_start events.
......@@ -109,7 +119,7 @@ func StreamResponseOpenAI2Claude(openAIResponse *dto.ChatCompletionsStreamRespon
if state.LastMessagesType == convmeta.LastMessageTypeNone {
return
}
stopOpenBlocks()
appendStopOpenBlocks()
switch state.LastMessagesType {
case convmeta.LastMessageTypeTools:
state.Index = state.ToolCallBaseIndex + state.ToolCallMaxIndexOffset + 1
......@@ -234,15 +244,17 @@ func StreamResponseOpenAI2Claude(openAIResponse *dto.ChatCompletionsStreamRespon
}
}
// 如果首块就带 finish_reason,需要立即发送停止块
// A first chunk can carry finish_reason before usage; defer terminal events until usage arrives.
if len(openAIResponse.Choices) > 0 && openAIResponse.Choices[0].FinishReason != nil && *openAIResponse.Choices[0].FinishReason != "" {
state.FinishReason = *openAIResponse.Choices[0].FinishReason
stopOpenBlocks()
oaiUsage := openAIResponse.Usage
if oaiUsage == nil {
oaiUsage = state.Usage
}
if oaiUsage != nil {
if oaiUsage == nil {
return claudeResponses
}
appendStopOpenBlocks()
claudeResponses = append(claudeResponses, &dto.ClaudeResponse{
Type: "message_delta",
Usage: buildClaudeUsageFromOpenAIUsage(oaiUsage),
......@@ -250,7 +262,6 @@ func StreamResponseOpenAI2Claude(openAIResponse *dto.ChatCompletionsStreamRespon
StopReason: kitutil.GetPointer[string](stopReasonOpenAI2Claude(state.FinishReason)),
},
})
}
claudeResponses = append(claudeResponses, &dto.ClaudeResponse{
Type: "message_stop",
})
......@@ -266,7 +277,7 @@ func StreamResponseOpenAI2Claude(openAIResponse *dto.ChatCompletionsStreamRespon
oaiUsage = state.Usage
}
if oaiUsage != nil {
stopOpenBlocks()
appendStopOpenBlocks()
stopReason := stopReasonOpenAI2Claude(state.FinishReason)
if stopReason == "" {
stopReason = "end_turn"
......@@ -403,7 +414,7 @@ func StreamResponseOpenAI2Claude(openAIResponse *dto.ChatCompletionsStreamRespon
}
if doneChunk || state.Done {
stopOpenBlocks()
appendStopOpenBlocks()
oaiUsage := openAIResponse.Usage
if oaiUsage == nil {
oaiUsage = state.Usage
......@@ -428,6 +439,34 @@ func StreamResponseOpenAI2Claude(openAIResponse *dto.ChatCompletionsStreamRespon
return claudeResponses
}
func FinalizeStreamResponseOpenAI2Claude(info convmeta.Meta) []*dto.ClaudeResponse {
if info == nil {
info = &convmeta.Values{}
}
state := info.EnsureClaudeConvertInfo()
if state.Done {
return nil
}
stopReason := stopReasonOpenAI2Claude(state.FinishReason)
if stopReason == "" {
stopReason = "end_turn"
}
responses := stopOpenBlocks(state)
responses = append(responses,
&dto.ClaudeResponse{
Type: "message_delta",
Usage: buildClaudeUsageFromOpenAIUsage(state.Usage),
Delta: &dto.ClaudeMediaMessage{
StopReason: kitutil.GetPointer[string](stopReason),
},
},
&dto.ClaudeResponse{Type: "message_stop"},
)
state.Done = true
return responses
}
func ResponseOpenAI2Claude(openAIResponse *dto.OpenAITextResponse, info convmeta.Meta) *dto.ClaudeResponse {
var stopReason string
contents := make([]dto.ClaudeMediaMessage, 0)
......
package relayconvert
import (
"context"
"errors"
"fmt"
"reflect"
"strings"
"sync"
"context"
"github.com/QuantumNous/new-api/relaykit/dto"
"github.com/QuantumNous/new-api/relaykit/relayconvert/convmeta"
geminichat "github.com/QuantumNous/new-api/relaykit/relayconvert/internal/gemini_chat"
oaichat "github.com/QuantumNous/new-api/relaykit/relayconvert/internal/oai_chat"
kitutil "github.com/QuantumNous/new-api/relaykit/relayconvert/kitutil"
"github.com/QuantumNous/new-api/relaykit/types"
)
......@@ -329,6 +331,13 @@ func FinalizeStreamResponse(c context.Context, info convmeta.Meta, state *Respon
return nil, nil
}
if state.To == types.RelayFormatClaude && info != nil {
claudeInfo := info.EnsureClaudeConvertInfo()
if claudeInfo.Usage == nil {
claudeInfo.Usage = state.Usage()
}
}
values := make([]any, 0)
var usage *dto.Usage
for i, spec := range state.specs {
......@@ -474,9 +483,6 @@ func executeStatelessStreamResponseSpec(c context.Context, info convmeta.Meta, f
var usage *dto.Usage
resultSteps := make([]ResponseStep, 0, len(steps))
for _, step := range steps {
if step.ConvertStreamChunk != nil || step.NewStreamState != nil || step.FinalizeStream != nil {
return nil, fmt.Errorf("response converter %q requires response stream state", step.ID)
}
if step.ConvertStream == nil {
return nil, fmt.Errorf("response converter %q has no stream implementation", step.ID)
}
......@@ -897,6 +903,15 @@ func convertOAIChatStreamResponseToClaudeMessages(_ context.Context, info convme
return StreamResponseOpenAI2Claude(chatResponse, info), canonicalUsageFromResponse(chatResponse), nil
}
func finalizeOAIChatStreamResponseToClaudeMessages(_ context.Context, info convmeta.Meta, _ any) ([]any, *dto.Usage, error) {
if info == nil {
info = &convmeta.Values{}
}
usage := info.EnsureClaudeConvertInfo().Usage
responses := oaichat.FinalizeStreamResponseOpenAI2Claude(info)
return streamValuesFromAny(responses), usage, nil
}
func convertOAIChatResponseToGeminiChat(_ context.Context, info convmeta.Meta, response any) (any, *dto.Usage, error) {
chatResponse, err := asOAIChatResponse(response)
if err != nil {
......@@ -955,6 +970,41 @@ func convertGeminiChatResponseToOAIChat(_ context.Context, info convmeta.Meta, r
return openAIResponse, usage, nil
}
func newGeminiChatToOAIChatStreamState(options ResponseStreamOptions) any {
return geminichat.NewGeminiToChatStreamState(options.ID, options.Created)
}
func convertGeminiChatStreamResponseChunkToOAIChat(_ context.Context, info convmeta.Meta, response any, state any) ([]any, *dto.Usage, error) {
geminiResponse, err := asGeminiChatResponse(response)
if err != nil {
return nil, nil, err
}
streamState, ok := state.(*geminichat.GeminiToChatStreamState)
if !ok || streamState == nil {
return nil, nil, errors.New("Gemini chat to OAI chat stream state is required")
}
usage := UsageFromGeminiMetadata(geminiResponse.GetUsageMetadata(), fallbackPromptTokens(info))
model := ""
if info != nil && info.HasChannelMeta() {
model = info.GetUpstreamModelName()
}
responses := streamState.ConvertChunk(geminiResponse, model, usage)
return streamValuesFromAny(responses), usage, nil
}
func finalizeGeminiChatStreamResponseToOAIChat(_ context.Context, info convmeta.Meta, state any) ([]any, *dto.Usage, error) {
streamState, ok := state.(*geminichat.GeminiToChatStreamState)
if !ok || streamState == nil {
return nil, nil, errors.New("Gemini chat to OAI chat stream state is required")
}
model := ""
if info != nil && info.HasChannelMeta() {
model = info.GetUpstreamModelName()
}
responses := streamState.Finalize(model)
return streamValuesFromAny(responses), streamState.Usage(), nil
}
func convertGeminiChatStreamResponseToOAIChat(_ context.Context, info convmeta.Meta, response any) (any, *dto.Usage, error) {
geminiResponse, err := asGeminiChatResponse(response)
if err != nil {
......
......@@ -14,7 +14,7 @@
"claude_cache_creation_1_h_tokens": 0
},
"role": "assistant",
"id": "chatcmpl-<uuid>",
"id": "stream_fixed",
"content": []
}
},
......@@ -41,6 +41,53 @@
"type": "text_delta",
"text": " world"
}
},
{
"type": "content_block_stop",
"index": 0
},
{
"type": "message_delta",
"usage": {
"input_tokens": 4,
"cache_creation_input_tokens": 0,
"cache_read_input_tokens": 0,
"output_tokens": 2,
"claude_cache_creation_5_m_tokens": 0,
"claude_cache_creation_1_h_tokens": 0,
"billing_usage": {
"source": "oai_chat",
"semantic": "openai",
"openai_usage": {
"prompt_tokens": 4,
"completion_tokens": 2,
"total_tokens": 6,
"prompt_tokens_details": {
"cached_tokens": 0,
"text_tokens": 4,
"audio_tokens": 0,
"image_tokens": 0
},
"completion_tokens_details": {
"text_tokens": 0,
"audio_tokens": 0,
"image_tokens": 0,
"reasoning_tokens": 0
},
"input_tokens": 0,
"output_tokens": 0,
"input_tokens_details": null,
"claude_cache_creation_5_m_tokens": 0,
"claude_cache_creation_1_h_tokens": 0
}
}
},
"delta": {
"stop_reason": "end_turn"
}
},
{
"type": "message_stop"
}
],
"usage": {
......
{
"events": [
{
"id": "chatcmpl-<uuid>",
"id": "stream_fixed",
"object": "chat.completion.chunk",
"created": 0,
"model": "upstream-model",
......@@ -40,7 +40,7 @@
}
},
{
"id": "chatcmpl-<uuid>",
"id": "stream_fixed",
"object": "chat.completion.chunk",
"created": 0,
"model": "upstream-model",
......@@ -92,6 +92,58 @@
"claude_cache_creation_5_m_tokens": 0,
"claude_cache_creation_1_h_tokens": 0
}
},
{
"id": "stream_fixed",
"object": "chat.completion.chunk",
"created": 0,
"model": "upstream-model",
"system_fingerprint": null,
"choices": [
{
"delta": {},
"logprobs": null,
"finish_reason": "stop",
"index": 0
}
],
"usage": {
"prompt_tokens": 4,
"completion_tokens": 2,
"total_tokens": 6,
"billing_usage": {
"source": "gemini_chat",
"semantic": "gemini",
"gemini_usage_metadata": {
"promptTokenCount": 4,
"toolUsePromptTokenCount": 0,
"candidatesTokenCount": 2,
"totalTokenCount": 6,
"thoughtsTokenCount": 0,
"cachedContentTokenCount": 0,
"promptTokensDetails": [],
"toolUsePromptTokensDetails": [],
"candidatesTokensDetails": []
}
},
"prompt_tokens_details": {
"cached_tokens": 0,
"text_tokens": 4,
"audio_tokens": 0,
"image_tokens": 0
},
"completion_tokens_details": {
"text_tokens": 0,
"audio_tokens": 0,
"image_tokens": 0,
"reasoning_tokens": 0
},
"input_tokens": 0,
"output_tokens": 0,
"input_tokens_details": null,
"claude_cache_creation_5_m_tokens": 0,
"claude_cache_creation_1_h_tokens": 0
}
}
],
"usage": {
......
......@@ -41,6 +41,53 @@
"type": "text_delta",
"text": " world"
}
},
{
"type": "content_block_stop",
"index": 0
},
{
"type": "message_delta",
"usage": {
"input_tokens": 4,
"cache_creation_input_tokens": 0,
"cache_read_input_tokens": 0,
"output_tokens": 2,
"claude_cache_creation_5_m_tokens": 0,
"claude_cache_creation_1_h_tokens": 0,
"billing_usage": {
"source": "oai_responses",
"semantic": "openai",
"openai_usage": {
"prompt_tokens": 0,
"completion_tokens": 0,
"total_tokens": 6,
"prompt_tokens_details": {
"cached_tokens": 0,
"text_tokens": 0,
"audio_tokens": 0,
"image_tokens": 0
},
"completion_tokens_details": {
"text_tokens": 0,
"audio_tokens": 0,
"image_tokens": 0,
"reasoning_tokens": 0
},
"input_tokens": 4,
"output_tokens": 2,
"input_tokens_details": null,
"claude_cache_creation_5_m_tokens": 0,
"claude_cache_creation_1_h_tokens": 0
}
}
},
"delta": {
"stop_reason": "end_turn"
}
},
{
"type": "message_stop"
}
],
"usage": {
......
......@@ -72,6 +72,7 @@ var builtinTextConverters = []TextConverterSpec{
Resp: TextResponseSide{
Convert: convertOAIChatResponseToClaudeMessages,
ConvertStream: convertOAIChatStreamResponseToClaudeMessages,
FinalizeStream: finalizeOAIChatStreamResponseToClaudeMessages,
Aliases: []string{ResponseConverterOAIChatToClaudeMessages},
},
},
......@@ -86,6 +87,9 @@ var builtinTextConverters = []TextConverterSpec{
Resp: TextResponseSide{
Convert: convertGeminiChatResponseToOAIChat,
ConvertStream: convertGeminiChatStreamResponseToOAIChat,
NewStreamState: newGeminiChatToOAIChatStreamState,
ConvertStreamChunk: convertGeminiChatStreamResponseChunkToOAIChat,
FinalizeStream: finalizeGeminiChatStreamResponseToOAIChat,
Aliases: []string{ResponseConverterGeminiChatToOAIChat},
},
},
......
......@@ -23,7 +23,7 @@ func TestLookupBuiltinTextConverters(t *testing.T) {
}{
{id: ConverterClaudeMessagesToOpenAIChat, from: types.RelayFormatClaude, to: types.RelayFormatOpenAI, quality: TextConverterQualityFair, reqDirect: true, respDirect: true, respAlias: ResponseConverterClaudeMessagesToOAIChat},
{id: ConverterOpenAIChatToClaudeMessages, from: types.RelayFormatOpenAI, to: types.RelayFormatClaude, quality: TextConverterQualityFair, reqDirect: true, respDirect: true, respAlias: ResponseConverterOAIChatToClaudeMessages},
{id: ConverterGeminiContentToOpenAIChat, from: types.RelayFormatGemini, to: types.RelayFormatOpenAI, quality: TextConverterQualityFair, reqDirect: true, respDirect: true, respAlias: ResponseConverterGeminiChatToOAIChat},
{id: ConverterGeminiContentToOpenAIChat, from: types.RelayFormatGemini, to: types.RelayFormatOpenAI, quality: TextConverterQualityFair, reqDirect: true, respDirect: true, respAlias: ResponseConverterGeminiChatToOAIChat, streamDirect: true},
{id: ConverterOpenAIChatToGeminiContent, from: types.RelayFormatOpenAI, to: types.RelayFormatGemini, quality: TextConverterQualityFair, reqDirect: true, respDirect: true, respAlias: ResponseConverterOAIChatToGeminiChat},
{id: ConverterOpenAIChatToOpenAIResponses, from: types.RelayFormatOpenAI, to: types.RelayFormatOpenAIResponses, quality: TextConverterQualityGood, reqDirect: true, respDirect: true, respAlias: ResponseConverterOAIChatToOAIResponses, streamDirect: true},
{id: ConverterOpenAIResponsesToOpenAIChat, from: types.RelayFormatOpenAIResponses, to: types.RelayFormatOpenAI, quality: TextConverterQualityGood, reqDirect: true, respDirect: true, respAlias: ResponseConverterOAIResponsesToOAIChat, streamDirect: true},
......
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or sign in to comment