Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 9 additions & 2 deletions core/providers/anthropic/anthropic.go
Original file line number Diff line number Diff line change
Expand Up @@ -997,8 +997,12 @@ func HandleAnthropicChatCompletionStreaming(
continue
}
var event AnthropicStreamEvent
if err := sonic.Unmarshal([]byte(eventData), &event); err != nil {
logger.Warn("Failed to parse message_start event: %v", err)
// Per-event decode -> "response-parse" (Serialization) stream phase.
parseStart := time.Now()
umErr := sonic.Unmarshal([]byte(eventData), &event)
schemas.AddStreamParse(ctx, time.Since(parseStart))
if umErr != nil {
logger.Warn("Failed to parse message_start event: %v", umErr)
continue
}
if event.Type == AnthropicStreamEventTypeMessageStart && event.Message != nil && event.Message.ID != "" {
Expand Down Expand Up @@ -1172,7 +1176,10 @@ func HandleAnthropicChatCompletionStreaming(
}
}

// Per-event mapping -> "convertor" (Convertor) stream phase.
convStart := time.Now()
response, bifrostErr, isLastChunk := event.ToBifrostChatCompletionStream(ctx, structuredOutputToolName, streamState)
schemas.AddStreamConvert(ctx, time.Since(convStart))
if bifrostErr != nil {
ctx.SetValue(schemas.BifrostContextKeyStreamEndIndicator, true)
providerUtils.ProcessAndSendBifrostError(ctx, postHookRunner, bifrostErr, responseChan, logger, postHookSpanFinalizer)
Expand Down
41 changes: 29 additions & 12 deletions core/providers/bedrock/bedrock.go
Original file line number Diff line number Diff line change
Expand Up @@ -1204,13 +1204,16 @@ func (provider *BedrockProvider) TextCompletionStream(ctx *schemas.BifrostContex
}
}

// Parse the chunk payload
// Parse the chunk payload. Per-event decode -> "response-parse" (Serialization) stream phase.
var chunkPayload struct {
Bytes []byte `json:"bytes"`
}
if err := sonic.Unmarshal(message.Payload, &chunkPayload); err != nil {
provider.logger.Debug("Failed to parse JSON from event buffer: %v, data: %s", err, string(message.Payload))
providerUtils.ProcessAndSendError(ctx, postHookRunner, err, responseChan, provider.logger, postHookSpanFinalizer)
parseStart := time.Now()
umErr := sonic.Unmarshal(message.Payload, &chunkPayload)
schemas.AddStreamParse(ctx, time.Since(parseStart))
if umErr != nil {
provider.logger.Debug("Failed to parse JSON from event buffer: %v, data: %s", umErr, string(message.Payload))
providerUtils.ProcessAndSendError(ctx, postHookRunner, umErr, responseChan, provider.logger, postHookSpanFinalizer)
return
}

Expand Down Expand Up @@ -1575,11 +1578,15 @@ func (provider *BedrockProvider) ChatCompletionStream(ctx *schemas.BifrostContex
}
}

// Converse API path: parse Bedrock Converse-specific stream events
// Converse API path: parse Bedrock Converse-specific stream events.
// Per-event decode -> "response-parse" (Serialization) stream phase.
var streamEvent BedrockStreamEvent
if err := sonic.Unmarshal(message.Payload, &streamEvent); err != nil {
provider.logger.Debug("Failed to parse JSON from event buffer: %v, data: %s", err, string(message.Payload))
providerUtils.ProcessAndSendError(ctx, postHookRunner, err, responseChan, provider.logger, postHookSpanFinalizer)
parseStart := time.Now()
umErr := sonic.Unmarshal(message.Payload, &streamEvent)
schemas.AddStreamParse(ctx, time.Since(parseStart))
if umErr != nil {
provider.logger.Debug("Failed to parse JSON from event buffer: %v, data: %s", umErr, string(message.Payload))
providerUtils.ProcessAndSendError(ctx, postHookRunner, umErr, responseChan, provider.logger, postHookSpanFinalizer)
return
}

Expand Down Expand Up @@ -1696,7 +1703,10 @@ func (provider *BedrockProvider) ChatCompletionStream(ctx *schemas.BifrostContex
}
}

// Per-event mapping -> "convertor" (Convertor) stream phase.
convStart := time.Now()
response, bifrostErr, _ := streamEvent.ToBifrostChatCompletionStream(streamState)
schemas.AddStreamConvert(ctx, time.Since(convStart))
if bifrostErr != nil {
ctx.SetValue(schemas.BifrostContextKeyStreamEndIndicator, true)
providerUtils.ProcessAndSendBifrostError(ctx, postHookRunner, bifrostErr, responseChan, provider.logger, postHookSpanFinalizer)
Expand Down Expand Up @@ -1994,11 +2004,15 @@ func (provider *BedrockProvider) ResponsesStream(ctx *schemas.BifrostContext, po
}
}

// Converse API path: parse Bedrock Converse-specific stream events
// Converse API path: parse Bedrock Converse-specific stream events.
// Per-event decode -> "response-parse" (Serialization) stream phase.
var streamEvent BedrockStreamEvent
if err := sonic.Unmarshal(message.Payload, &streamEvent); err != nil {
provider.logger.Debug("Failed to parse JSON from event buffer: %v, data: %s", err, string(message.Payload))
providerUtils.ProcessAndSendError(ctx, postHookRunner, err, responseChan, provider.logger, postHookSpanFinalizer)
parseStart := time.Now()
umErr := sonic.Unmarshal(message.Payload, &streamEvent)
schemas.AddStreamParse(ctx, time.Since(parseStart))
if umErr != nil {
provider.logger.Debug("Failed to parse JSON from event buffer: %v, data: %s", umErr, string(message.Payload))
providerUtils.ProcessAndSendError(ctx, postHookRunner, umErr, responseChan, provider.logger, postHookSpanFinalizer)
return
}

Expand Down Expand Up @@ -2062,7 +2076,10 @@ func (provider *BedrockProvider) ResponsesStream(ctx *schemas.BifrostContext, po
}
}

// Per-event mapping -> "convertor" (Convertor) stream phase.
convStart := time.Now()
responses, bifrostErr, _ := streamEvent.ToBifrostResponsesStream(chunkIndex, streamState)
schemas.AddStreamConvert(ctx, time.Since(convStart))
if bifrostErr != nil {
ctx.SetValue(schemas.BifrostContextKeyStreamEndIndicator, true)
providerUtils.ProcessAndSendBifrostError(ctx, postHookRunner, bifrostErr, responseChan, provider.logger, postHookSpanFinalizer)
Expand Down
24 changes: 18 additions & 6 deletions core/providers/cohere/cohere.go
Original file line number Diff line number Diff line change
Expand Up @@ -568,10 +568,13 @@ func (provider *CohereProvider) ChatCompletionStream(ctx *schemas.BifrostContext

eventData := string(data)

// Parse the unified streaming event
// Parse the unified streaming event. Per-event decode -> "response-parse" (Serialization) stream phase.
var event CohereStreamEvent
if err := sonic.Unmarshal(data, &event); err != nil {
provider.logger.Warn("Failed to parse stream event: %v", err)
parseStart := time.Now()
umErr := sonic.Unmarshal(data, &event)
schemas.AddStreamParse(ctx, time.Since(parseStart))
if umErr != nil {
provider.logger.Warn("Failed to parse stream event: %v", umErr)
continue
}

Expand All @@ -580,7 +583,10 @@ func (provider *CohereProvider) ChatCompletionStream(ctx *schemas.BifrostContext
responseID = *event.ID
}

// Per-event mapping -> "convertor" (Convertor) stream phase.
convStart := time.Now()
response, bifrostErr, isLastChunk := event.ToBifrostChatCompletionStream()
schemas.AddStreamConvert(ctx, time.Since(convStart))
if bifrostErr != nil {
ctx.SetValue(schemas.BifrostContextKeyStreamEndIndicator, true)
providerUtils.ProcessAndSendBifrostError(ctx, postHookRunner, bifrostErr, responseChan, provider.logger, postHookSpanFinalizer)
Expand Down Expand Up @@ -855,17 +861,23 @@ func (provider *CohereProvider) ResponsesStream(ctx *schemas.BifrostContext, pos

eventData := string(data)

// Parse the unified streaming event
// Parse the unified streaming event. Per-event decode -> "response-parse" (Serialization) stream phase.
var event CohereStreamEvent
if err := sonic.Unmarshal(data, &event); err != nil {
provider.logger.Warn("Failed to parse stream event: %v", err)
parseStart := time.Now()
umErr := sonic.Unmarshal(data, &event)
schemas.AddStreamParse(ctx, time.Since(parseStart))
if umErr != nil {
provider.logger.Warn("Failed to parse stream event: %v", umErr)
continue
}

// Note: response.created and response.in_progress are now emitted by ToBifrostResponsesStream
// from the message_start event, so we don't need to call them manually here

// Per-event mapping -> "convertor" (Convertor) stream phase.
convStart := time.Now()
responses, bifrostErr, isLastChunk := event.ToBifrostResponsesStream(chunkIndex, streamState)
schemas.AddStreamConvert(ctx, time.Since(convStart))
if bifrostErr != nil {
ctx.SetValue(schemas.BifrostContextKeyStreamEndIndicator, true)
providerUtils.ProcessAndSendBifrostError(ctx, postHookRunner, bifrostErr, responseChan, provider.logger, postHookSpanFinalizer)
Expand Down
16 changes: 12 additions & 4 deletions core/providers/gemini/gemini.go
Original file line number Diff line number Diff line change
Expand Up @@ -559,8 +559,10 @@ func HandleGeminiChatCompletionStream(
providerUtils.ProcessAndSendError(ctx, postHookRunner, readErr, responseChan, logger, postHookSpanFinalizer)
return
}
// Process chunk using shared function
// Process chunk using shared function. Per-event decode -> "response-parse" (Serialization) stream phase.
parseStart := time.Now()
geminiResponse, err := processGeminiStreamChunk(eventData)
schemas.AddStreamParse(ctx, time.Since(parseStart))
if err != nil {
if strings.Contains(err.Error(), "gemini api error") {
// Handle API error
Expand All @@ -581,8 +583,10 @@ func HandleGeminiChatCompletionStream(
modelName = geminiResponse.ModelVersion
}

// Convert to Bifrost stream response
// Convert to Bifrost stream response. Per-event mapping -> "convertor" (Convertor) stream phase.
convStart := time.Now()
response, bifrostErr, isLastChunk := geminiResponse.ToBifrostChatCompletionStream(streamState)
schemas.AddStreamConvert(ctx, time.Since(convStart))
if bifrostErr != nil {
ctx.SetValue(schemas.BifrostContextKeyStreamEndIndicator, true)
providerUtils.ProcessAndSendBifrostError(ctx, postHookRunner, providerUtils.EnrichError(ctx, bifrostErr, jsonBody, nil, sendBackRawRequest, sendBackRawResponse, latency), responseChan, logger, postHookSpanFinalizer)
Expand Down Expand Up @@ -1072,8 +1076,10 @@ func HandleGeminiResponsesStream(
return
}

// Process chunk using shared function
// Process chunk using shared function. Per-event decode -> "response-parse" (Serialization) stream phase.
parseStart := time.Now()
geminiResponse, err := processGeminiStreamChunk(eventData)
schemas.AddStreamParse(ctx, time.Since(parseStart))
if err != nil {
if strings.Contains(err.Error(), "gemini api error") {
// Handle API error
Expand All @@ -1095,8 +1101,10 @@ func HandleGeminiResponsesStream(
}
}

// Convert to Bifrost responses stream response
// Convert to Bifrost responses stream response. Per-event mapping -> "convertor" (Convertor) stream phase.
convStart := time.Now()
responses, bifrostErr := geminiResponse.ToBifrostResponsesStream(sequenceNumber, streamState)
schemas.AddStreamConvert(ctx, time.Since(convStart))
if bifrostErr != nil {
ctx.SetValue(schemas.BifrostContextKeyStreamEndIndicator, true)
providerUtils.ProcessAndSendBifrostError(ctx, postHookRunner, providerUtils.EnrichError(ctx, bifrostErr, jsonBody, nil, sendBackRawRequest, sendBackRawResponse), responseChan, logger, postHookSpanFinalizer)
Expand Down
57 changes: 47 additions & 10 deletions core/providers/openai/openai.go
Original file line number Diff line number Diff line change
Expand Up @@ -618,7 +618,10 @@ func HandleOpenAITextCompletionStreaming(
jsonData := string(data)
var response schemas.BifrostTextCompletionResponse
if customResponseHandler != nil {
// Custom handler decodes the raw event itself -> time as "response-parse" stream phase.
parseStart := time.Now()
rawRequest, rawResponse, handlerErr := customResponseHandler([]byte(jsonData), &response, nil, sendBackRawRequest, sendBackRawResponse)
schemas.AddStreamParse(ctx, time.Since(parseStart))
if handlerErr != nil {
// TODO fix this
if sendBackRawRequest {
Expand Down Expand Up @@ -646,9 +649,14 @@ func HandleOpenAITextCompletionStreaming(
}
}

// Parse into bifrost response
if err := sonic.UnmarshalString(jsonData, &response); err != nil {
logger.Warn("Failed to parse stream response: %v", err)
// Parse into bifrost response. Timed as the "response-parse" stream phase
// (per-event JSON decode) so it lands in Serialization like unary/Anthropic,
// instead of folding into core/provider-internal.
parseStart := time.Now()
umErr := sonic.UnmarshalString(jsonData, &response)
schemas.AddStreamParse(ctx, time.Since(parseStart))
Comment thread
coderabbitai[bot] marked this conversation as resolved.
if umErr != nil {
logger.Warn("Failed to parse stream response: %v", umErr)
Comment thread
coderabbitai[bot] marked this conversation as resolved.
continue
}
}
Expand All @@ -659,7 +667,11 @@ func HandleOpenAITextCompletionStreaming(
}

if postResponseConverter != nil {
if converted := postResponseConverter(&response); converted != nil {
// Per-event mapping -> "convertor" (Convertor) stream phase.
convStart := time.Now()
converted := postResponseConverter(&response)
schemas.AddStreamConvert(ctx, time.Since(convStart))
Comment thread
coderabbitai[bot] marked this conversation as resolved.
if converted != nil {
response = *converted
} else {
logger.Warn("postResponseConverter returned nil; leaving chunk unmodified")
Expand Down Expand Up @@ -1263,7 +1275,10 @@ func HandleOpenAIChatCompletionStreaming(
// Parse into bifrost response
var response schemas.BifrostChatResponse
if customResponseHandler != nil {
// Custom handler decodes the raw event itself -> time as "response-parse" stream phase.
parseStart := time.Now()
rawRequest, rawResponse, handlerErr := customResponseHandler([]byte(jsonData), &response, nil, sendBackRawRequest, sendBackRawResponse)
schemas.AddStreamParse(ctx, time.Since(parseStart))
if handlerErr != nil {
if sendBackRawRequest {
handlerErr.ExtraFields.RawRequest = rawRequest
Expand All @@ -1276,8 +1291,12 @@ func HandleOpenAIChatCompletionStreaming(
return
}
} else {
if err := sonic.UnmarshalString(jsonData, &response); err != nil {
logger.Warn("Failed to parse stream response: %v", err)
// Per-event decode -> "response-parse" (Serialization) stream phase.
parseStart := time.Now()
umErr := sonic.UnmarshalString(jsonData, &response)
schemas.AddStreamParse(ctx, time.Since(parseStart))
if umErr != nil {
logger.Warn("Failed to parse stream response: %v", umErr)
continue
}
}
Expand Down Expand Up @@ -1318,7 +1337,10 @@ func HandleOpenAIChatCompletionStreaming(
}
}

// Per-event mapping (chat->responses) -> "convertor" (Convertor) stream phase.
convStart := time.Now()
spreadResponses := response.ToBifrostResponsesStreamResponse(responsesStreamState)
schemas.AddStreamConvert(ctx, time.Since(convStart))
for _, response := range spreadResponses {
if response.Type == schemas.ResponsesStreamResponseTypeError {
bifrostErr := &schemas.BifrostError{
Expand Down Expand Up @@ -1361,7 +1383,11 @@ func HandleOpenAIChatCompletionStreaming(
}
} else {
if postResponseConverter != nil {
if converted := postResponseConverter(&response); converted != nil {
// Per-event mapping -> "convertor" (Convertor) stream phase.
convStart := time.Now()
converted := postResponseConverter(&response)
schemas.AddStreamConvert(ctx, time.Since(convStart))
if converted != nil {
response = *converted
} else {
logger.Warn("postResponseConverter returned nil; leaving chunk unmodified")
Expand Down Expand Up @@ -1943,7 +1969,10 @@ func HandleOpenAIResponsesStreaming(
// Parse into bifrost response
var response schemas.BifrostResponsesStreamResponse
if customResponseHandler != nil {
// Custom handler decodes the raw event itself -> time as "response-parse" stream phase.
parseStart := time.Now()
rawRequest, rawResponse, bifrostErr := customResponseHandler([]byte(jsonData), &response, nil, sendBackRawRequest, sendBackRawResponse)
schemas.AddStreamParse(ctx, time.Since(parseStart))
if bifrostErr != nil {
if sendBackRawRequest {
bifrostErr.ExtraFields.RawRequest = rawRequest
Expand All @@ -1959,13 +1988,21 @@ func HandleOpenAIResponsesStreaming(
response.ExtraFields.RawResponse = jsonData
}
} else {
if err := sonic.UnmarshalString(jsonData, &response); err != nil {
logger.Warn("Failed to parse stream response: %v", err)
// Per-event decode -> "response-parse" (Serialization) stream phase.
parseStart := time.Now()
umErr := sonic.UnmarshalString(jsonData, &response)
schemas.AddStreamParse(ctx, time.Since(parseStart))
if umErr != nil {
logger.Warn("Failed to parse stream response: %v", umErr)
continue
}

if postResponseConverter != nil {
if converted := postResponseConverter(&response); converted != nil {
// Per-event mapping -> "convertor" (Convertor) stream phase.
convStart := time.Now()
converted := postResponseConverter(&response)
schemas.AddStreamConvert(ctx, time.Since(convStart))
if converted != nil {
response = *converted
} else {
logger.Warn("postResponseConverter returned nil; leaving chunk unmodified")
Expand Down
Loading
Loading