From 9a60961e68f84bc370797ff85aebd5e877114be7 Mon Sep 17 00:00:00 2001 From: tejas ghatte Date: Thu, 28 May 2026 17:14:26 +0530 Subject: [PATCH 1/2] fix: responses stream events --- core/providers/anthropic/responses.go | 74 ++++++++++++++-- core/providers/bedrock/responses.go | 89 ++++++++++++++++--- core/providers/cohere/responses.go | 121 +++++++++++++++++++++----- core/providers/gemini/responses.go | 19 +++- 4 files changed, 256 insertions(+), 47 deletions(-) diff --git a/core/providers/anthropic/responses.go b/core/providers/anthropic/responses.go index 7b795aa14b0..98c6fdbbd43 100644 --- a/core/providers/anthropic/responses.go +++ b/core/providers/anthropic/responses.go @@ -43,6 +43,7 @@ type AnthropicResponsesStreamState struct { ReasoningSignatures map[int]string // Maps output_index to reasoning signature TextContentIndices map[int]bool // Tracks which content indices are text blocks ReasoningContentIndices map[int]bool // Tracks which content indices are reasoning blocks + TextBuffers map[int]*strings.Builder // Maps output_index to accumulated text content for done events CompactionContentIndices map[int]*schemas.CacheControl // Tracks pending compaction blocks with their cache control CurrentOutputIndex int // Current output index counter MessageID *string // Message ID from message_start @@ -71,6 +72,7 @@ var anthropicResponsesStreamStatePool = sync.Pool{ ReasoningContentIndices: make(map[int]bool), CompactionContentIndices: make(map[int]*schemas.CacheControl), OutputItems: make(map[int]*schemas.ResponsesMessage), + TextBuffers: make(map[int]*strings.Builder), CurrentOutputIndex: 0, CreatedAt: int(time.Now().Unix()), HasEmittedCreated: false, @@ -147,6 +149,11 @@ func acquireAnthropicResponsesStreamState() *AnthropicResponsesStreamState { } else { clear(state.ReasoningContentIndices) } + if state.TextBuffers == nil { + state.TextBuffers = make(map[int]*strings.Builder) + } else { + clear(state.TextBuffers) + } if state.CompactionContentIndices == nil { state.CompactionContentIndices = make(map[int]*schemas.CacheControl) } else { @@ -207,6 +214,7 @@ func (state *AnthropicResponsesStreamState) flush() { state.ReasoningSignatures = nil state.TextContentIndices = nil state.ReasoningContentIndices = nil + state.TextBuffers = nil state.CompactionContentIndices = nil state.OutputItems = nil state.CurrentOutputIndex = 0 @@ -837,6 +845,12 @@ func (chunk *AnthropicStreamEvent) ToBifrostResponsesStream(ctx context.Context, } case AnthropicStreamDeltaTypeText: if chunk.Delta.Text != nil && *chunk.Delta.Text != "" { + // Accumulate text for done events + if state.TextBuffers[outputIndex] == nil { + state.TextBuffers[outputIndex] = &strings.Builder{} + } + state.TextBuffers[outputIndex].WriteString(*chunk.Delta.Text) + // Text content delta - emit output_text.delta with item ID itemID := state.ItemIDs[outputIndex] response := &schemas.BifrostResponsesStreamResponse{ @@ -1124,29 +1138,44 @@ func (chunk *AnthropicStreamEvent) ToBifrostResponsesStream(ctx context.Context, var responses []*schemas.BifrostResponsesStreamResponse itemID := state.ItemIDs[outputIndex] + // Capture accumulated text once — shared by output_text.done and output_item.done + accText := "" + if buf := state.TextBuffers[outputIndex]; buf != nil { + accText = buf.String() + } + // Check if this content index is a text block if chunk.Index != nil { if state.TextContentIndices[*chunk.Index] { - // Emit output_text.done (without accumulated text, just the event) - emptyText := "" + // Emit output_text.done with full accumulated text textDoneResponse := &schemas.BifrostResponsesStreamResponse{ Type: schemas.ResponsesStreamResponseTypeOutputTextDone, SequenceNumber: sequenceNumber + len(responses), OutputIndex: schemas.Ptr(outputIndex), ContentIndex: chunk.Index, - Text: &emptyText, + Text: &accText, } if itemID != "" { textDoneResponse.ItemID = &itemID } responses = append(responses, textDoneResponse) - // Emit content_part.done + // Emit content_part.done with full accumulated text in Part + partText := accText + part := &schemas.ResponsesMessageContentBlock{ + Type: schemas.ResponsesOutputMessageContentTypeText, + Text: &partText, + ResponsesOutputMessageContentText: &schemas.ResponsesOutputMessageContentText{ + Annotations: []schemas.ResponsesOutputMessageContentTextAnnotation{}, + LogProbs: []schemas.ResponsesOutputMessageContentTextLogProb{}, + }, + } partDoneResponse := &schemas.BifrostResponsesStreamResponse{ Type: schemas.ResponsesStreamResponseTypeContentPartDone, SequenceNumber: sequenceNumber + len(responses), OutputIndex: schemas.Ptr(outputIndex), ContentIndex: chunk.Index, + Part: part, } if itemID != "" { partDoneResponse.ItemID = &itemID @@ -1292,17 +1321,42 @@ func (chunk *AnthropicStreamEvent) ToBifrostResponsesStream(ctx context.Context, } doneItem = &copied } else { + // Build content blocks from accumulated text (captured above) + contentBlocks := []schemas.ResponsesMessageContentBlock{} + if accText != "" { + textCopy := accText + contentBlocks = []schemas.ResponsesMessageContentBlock{ + { + Type: schemas.ResponsesOutputMessageContentTypeText, + Text: &textCopy, + ResponsesOutputMessageContentText: &schemas.ResponsesOutputMessageContentText{ + Annotations: []schemas.ResponsesOutputMessageContentTextAnnotation{}, + LogProbs: []schemas.ResponsesOutputMessageContentTextLogProb{}, + }, + }, + } + } + delete(state.TextBuffers, outputIndex) doneItem = &schemas.ResponsesMessage{ Type: schemas.Ptr(schemas.ResponsesMessageTypeMessage), Role: schemas.Ptr(schemas.ResponsesInputMessageRoleAssistant), Status: &statusCompleted, Content: &schemas.ResponsesMessageContent{ - ContentBlocks: []schemas.ResponsesMessageContentBlock{}, + ContentBlocks: contentBlocks, }, } if doneItemID != "" { doneItem.ID = &doneItemID } + // Only persist synthesized items that actually have text content — reasoning + // and MCP blocks fall through here without storedItems but must not pollute + // response.completed with empty assistant message shells. + if len(contentBlocks) > 0 { + cloned := *doneItem + clonedContent := *doneItem.Content + cloned.Content = &clonedContent + state.OutputItems[outputIndex] = &cloned + } } responses = append(responses, &schemas.BifrostResponsesStreamResponse{ Type: schemas.ResponsesStreamResponseTypeOutputItemDone, @@ -5192,12 +5246,14 @@ func convertBifrostToolToAnthropic(model string, tool *schemas.ResponsesTool, pr } } - anthropicTool := &AnthropicTool{} - - if tool.Name != nil { - anthropicTool.Name = *tool.Name + // Skip tools with no name — Anthropic rejects them + if tool.Name == nil || *tool.Name == "" { + return nil } + anthropicTool := &AnthropicTool{} + anthropicTool.Name = *tool.Name + if tool.Description != nil { anthropicTool.Description = tool.Description } diff --git a/core/providers/bedrock/responses.go b/core/providers/bedrock/responses.go index 7e70ac5a358..8466fefbf80 100644 --- a/core/providers/bedrock/responses.go +++ b/core/providers/bedrock/responses.go @@ -30,6 +30,7 @@ type BedrockResponsesStreamState struct { NovaGroundingCitations map[int][]schemas.ResponsesWebSearchToolCallActionSearchSource // Collected citation sources per nova_grounding output index CompletedOutputIndices map[int]bool // Tracks which output indices have been completed AnnotationIndices map[int]int // Maps output_index to next annotation index for sequential citation numbering + TextBuffers map[int]*strings.Builder // Maps output_index to accumulated text content for done events CurrentOutputIndex int // Current output index counter MessageID *string // Message ID (generated) Model *string // Model name @@ -55,6 +56,7 @@ var bedrockResponsesStreamStatePool = sync.Pool{ NovaGroundingCitations: make(map[int][]schemas.ResponsesWebSearchToolCallActionSearchSource), CompletedOutputIndices: make(map[int]bool), AnnotationIndices: make(map[int]int), + TextBuffers: make(map[int]*strings.Builder), CurrentOutputIndex: 0, CreatedAt: int(time.Now().Unix()), HasEmittedCreated: false, @@ -123,6 +125,11 @@ func acquireBedrockResponsesStreamState() *BedrockResponsesStreamState { } else { clear(state.AnnotationIndices) } + if state.TextBuffers == nil { + state.TextBuffers = make(map[int]*strings.Builder) + } else { + clear(state.TextBuffers) + } // Reset other fields state.CurrentOutputIndex = 0 state.MessageID = nil @@ -200,6 +207,11 @@ func (state *BedrockResponsesStreamState) flush() { } else { clear(state.AnnotationIndices) } + if state.TextBuffers == nil { + state.TextBuffers = make(map[int]*strings.Builder) + } else { + clear(state.TextBuffers) + } state.CurrentOutputIndex = 0 state.MessageID = nil state.Model = nil @@ -377,22 +389,27 @@ func (chunk *BedrockStreamEvent) ToBifrostResponsesStream(sequenceNumber int, st continue } - // Emit output_text.done - emptyText := "" + prevAccText := "" + if buf := state.TextBuffers[prevOutputIndex]; buf != nil { + prevAccText = buf.String() + } + + // Emit output_text.done with accumulated text responses = append(responses, &schemas.BifrostResponsesStreamResponse{ Type: schemas.ResponsesStreamResponseTypeOutputTextDone, SequenceNumber: sequenceNumber + len(responses), OutputIndex: schemas.Ptr(prevOutputIndex), ContentIndex: &prevContentIndex, ItemID: &prevItemID, - Text: &emptyText, + Text: &prevAccText, LogProbs: []schemas.ResponsesOutputMessageContentTextLogProb{}, }) - // Emit content_part.done for text + // Emit content_part.done for text with accumulated text + prevPartText := prevAccText part := &schemas.ResponsesMessageContentBlock{ Type: schemas.ResponsesOutputMessageContentTypeText, - Text: &emptyText, + Text: &prevPartText, ResponsesOutputMessageContentText: &schemas.ResponsesOutputMessageContentText{ LogProbs: []schemas.ResponsesOutputMessageContentTextLogProb{}, Annotations: []schemas.ResponsesOutputMessageContentTextAnnotation{}, @@ -407,8 +424,22 @@ func (chunk *BedrockStreamEvent) ToBifrostResponsesStream(sequenceNumber int, st Part: part, }) - // Emit output_item.done for text + // Emit output_item.done for text with content blocks statusCompleted := "completed" + var prevContentBlocks []schemas.ResponsesMessageContentBlock + if prevAccText != "" { + prevItemText := prevAccText + prevContentBlocks = []schemas.ResponsesMessageContentBlock{ + { + Type: schemas.ResponsesOutputMessageContentTypeText, + Text: &prevItemText, + ResponsesOutputMessageContentText: &schemas.ResponsesOutputMessageContentText{ + Annotations: []schemas.ResponsesOutputMessageContentTextAnnotation{}, + LogProbs: []schemas.ResponsesOutputMessageContentTextLogProb{}, + }, + }, + } + } messageType := schemas.ResponsesMessageTypeMessage role := schemas.ResponsesInputMessageRoleAssistant doneItem := &schemas.ResponsesMessage{ @@ -416,7 +447,7 @@ func (chunk *BedrockStreamEvent) ToBifrostResponsesStream(sequenceNumber int, st Role: &role, Status: &statusCompleted, Content: &schemas.ResponsesMessageContent{ - ContentBlocks: []schemas.ResponsesMessageContentBlock{}, + ContentBlocks: prevContentBlocks, }, } if prevItemID != "" { @@ -429,6 +460,7 @@ func (chunk *BedrockStreamEvent) ToBifrostResponsesStream(sequenceNumber int, st ContentIndex: &prevContentIndex, Item: doneItem, }) + delete(state.TextBuffers, prevOutputIndex) // Mark this output index as completed state.CompletedOutputIndices[prevOutputIndex] = true @@ -784,6 +816,11 @@ func (chunk *BedrockStreamEvent) ToBifrostResponsesStream(sequenceNumber int, st // If this is a text delta for a new content block, also emit the text delta in the same batch if chunk.Delta.Text != nil && *chunk.Delta.Text != "" { text := *chunk.Delta.Text + // Accumulate text for done events + if state.TextBuffers[outputIndex] == nil { + state.TextBuffers[outputIndex] = &strings.Builder{} + } + state.TextBuffers[outputIndex].WriteString(text) itemID := state.ItemIDs[outputIndex] textDeltaResponse := &schemas.BifrostResponsesStreamResponse{ Type: schemas.ResponsesStreamResponseTypeOutputTextDelta, @@ -810,6 +847,11 @@ func (chunk *BedrockStreamEvent) ToBifrostResponsesStream(sequenceNumber int, st // Handle text delta text := *chunk.Delta.Text if text != "" { + // Accumulate text for done events + if state.TextBuffers[outputIndex] == nil { + state.TextBuffers[outputIndex] = &strings.Builder{} + } + state.TextBuffers[outputIndex].WriteString(text) itemID := state.ItemIDs[outputIndex] response := &schemas.BifrostResponsesStreamResponse{ Type: schemas.ResponsesStreamResponseTypeOutputTextDelta, @@ -1253,23 +1295,27 @@ func FinalizeBedrockStream(state *BedrockResponsesStreamState, sequenceNumber in } // end else (regular function call) } else { // This is likely a text item that needs to be closed + accText := "" + if buf := state.TextBuffers[outputIndex]; buf != nil { + accText = buf.String() + } - // Emit output_text.done (without accumulated text, just the event) - emptyText := "" + // Emit output_text.done with accumulated text responses = append(responses, &schemas.BifrostResponsesStreamResponse{ Type: schemas.ResponsesStreamResponseTypeOutputTextDone, SequenceNumber: sequenceNumber + len(responses), OutputIndex: schemas.Ptr(outputIndex), ContentIndex: &contentIndex, ItemID: &itemID, - Text: &emptyText, + Text: &accText, LogProbs: []schemas.ResponsesOutputMessageContentTextLogProb{}, }) - // Emit content_part.done for text + // Emit content_part.done for text with accumulated text + partText := accText part := &schemas.ResponsesMessageContentBlock{ Type: schemas.ResponsesOutputMessageContentTypeText, - Text: &emptyText, + Text: &partText, ResponsesOutputMessageContentText: &schemas.ResponsesOutputMessageContentText{ LogProbs: []schemas.ResponsesOutputMessageContentTextLogProb{}, Annotations: []schemas.ResponsesOutputMessageContentTextAnnotation{}, @@ -1284,8 +1330,22 @@ func FinalizeBedrockStream(state *BedrockResponsesStreamState, sequenceNumber in Part: part, }) - // Emit output_item.done for text + // Emit output_item.done for text with content blocks statusCompleted := "completed" + contentBlocks := []schemas.ResponsesMessageContentBlock{} + if accText != "" { + itemText := accText + contentBlocks = []schemas.ResponsesMessageContentBlock{ + { + Type: schemas.ResponsesOutputMessageContentTypeText, + Text: &itemText, + ResponsesOutputMessageContentText: &schemas.ResponsesOutputMessageContentText{ + Annotations: []schemas.ResponsesOutputMessageContentTextAnnotation{}, + LogProbs: []schemas.ResponsesOutputMessageContentTextLogProb{}, + }, + }, + } + } messageType := schemas.ResponsesMessageTypeMessage role := schemas.ResponsesInputMessageRoleAssistant doneItem := &schemas.ResponsesMessage{ @@ -1293,7 +1353,7 @@ func FinalizeBedrockStream(state *BedrockResponsesStreamState, sequenceNumber in Role: &role, Status: &statusCompleted, Content: &schemas.ResponsesMessageContent{ - ContentBlocks: []schemas.ResponsesMessageContentBlock{}, + ContentBlocks: contentBlocks, }, } if itemID != "" { @@ -1306,6 +1366,7 @@ func FinalizeBedrockStream(state *BedrockResponsesStreamState, sequenceNumber in ContentIndex: &contentIndex, Item: doneItem, }) + delete(state.TextBuffers, outputIndex) } // Mark this output index as completed diff --git a/core/providers/cohere/responses.go b/core/providers/cohere/responses.go index c4fd300792c..77f9eb8fad7 100644 --- a/core/providers/cohere/responses.go +++ b/core/providers/cohere/responses.go @@ -20,6 +20,7 @@ type CohereResponsesStreamState struct { ItemIDs map[int]string // Maps output_index to item ID for stable IDs ReasoningContentIndices map[int]bool // Tracks which content indices are reasoning blocks AnnotationIndexToContentIndex map[int]int // Maps annotation index to content index for citation pairing + TextBuffers map[int]*strings.Builder // Maps output_index to accumulated text content for done events CurrentOutputIndex int // Current output index counter MessageID *string // Message ID from message_start Model *string // Model name from message_start @@ -39,6 +40,7 @@ var cohereResponsesStreamStatePool = sync.Pool{ ItemIDs: make(map[int]string), ReasoningContentIndices: make(map[int]bool), AnnotationIndexToContentIndex: make(map[int]int), + TextBuffers: make(map[int]*strings.Builder), CurrentOutputIndex: 0, CreatedAt: int(time.Now().Unix()), HasEmittedCreated: false, @@ -83,6 +85,11 @@ func acquireCohereResponsesStreamState() *CohereResponsesStreamState { } else { clear(state.AnnotationIndexToContentIndex) } + if state.TextBuffers == nil { + state.TextBuffers = make(map[int]*strings.Builder) + } else { + clear(state.TextBuffers) + } // Reset other fields state.CurrentOutputIndex = 0 state.MessageID = nil @@ -135,6 +142,11 @@ func (state *CohereResponsesStreamState) flush() { } else { clear(state.AnnotationIndexToContentIndex) } + if state.TextBuffers == nil { + state.TextBuffers = make(map[int]*strings.Builder) + } else { + clear(state.TextBuffers) + } state.CurrentOutputIndex = 0 state.MessageID = nil state.Model = nil @@ -258,23 +270,27 @@ func (chunk *CohereStreamEvent) ToBifrostResponsesStream(sequenceNumber int, sta if state.ToolPlanOutputIndex != nil { outputIndex := *state.ToolPlanOutputIndex itemID := state.ItemIDs[outputIndex] + accText := "" + if buf := state.TextBuffers[outputIndex]; buf != nil { + accText = buf.String() + } - // Emit output_text.done (without accumulated text, just the event) - emptyText := "" + // Emit output_text.done with accumulated text responses = append(responses, &schemas.BifrostResponsesStreamResponse{ Type: schemas.ResponsesStreamResponseTypeOutputTextDone, SequenceNumber: sequenceNumber + len(responses), OutputIndex: schemas.Ptr(outputIndex), ContentIndex: schemas.Ptr(0), ItemID: &itemID, - Text: &emptyText, + Text: &accText, LogProbs: []schemas.ResponsesOutputMessageContentTextLogProb{}, }) - // Emit content_part.done + // Emit content_part.done with accumulated text + partText := accText part := &schemas.ResponsesMessageContentBlock{ Type: schemas.ResponsesOutputMessageContentTypeText, - Text: &emptyText, + Text: &partText, ResponsesOutputMessageContentText: &schemas.ResponsesOutputMessageContentText{ LogProbs: []schemas.ResponsesOutputMessageContentTextLogProb{}, Annotations: []schemas.ResponsesOutputMessageContentTextAnnotation{}, @@ -289,7 +305,8 @@ func (chunk *CohereStreamEvent) ToBifrostResponsesStream(sequenceNumber int, sta Part: part, }) - // Emit output_item.done + // Emit output_item.done with content blocks + itemText := accText statusCompleted := "completed" messageType := schemas.ResponsesMessageTypeMessage role := schemas.ResponsesInputMessageRoleAssistant @@ -298,7 +315,16 @@ func (chunk *CohereStreamEvent) ToBifrostResponsesStream(sequenceNumber int, sta Role: &role, Status: &statusCompleted, Content: &schemas.ResponsesMessageContent{ - ContentBlocks: []schemas.ResponsesMessageContentBlock{}, + ContentBlocks: []schemas.ResponsesMessageContentBlock{ + { + Type: schemas.ResponsesOutputMessageContentTypeText, + Text: &itemText, + ResponsesOutputMessageContentText: &schemas.ResponsesOutputMessageContentText{ + Annotations: []schemas.ResponsesOutputMessageContentTextAnnotation{}, + LogProbs: []schemas.ResponsesOutputMessageContentTextLogProb{}, + }, + }, + }, }, } if itemID != "" { @@ -311,6 +337,10 @@ func (chunk *CohereStreamEvent) ToBifrostResponsesStream(sequenceNumber int, sta ContentIndex: schemas.Ptr(0), Item: doneItem, }) + delete(state.TextBuffers, outputIndex) + if mapped, ok := state.ContentIndexToOutputIndex[0]; ok && mapped == outputIndex { + delete(state.ContentIndexToOutputIndex, 0) + } state.ToolPlanOutputIndex = nil // Mark as closed } @@ -427,6 +457,12 @@ func (chunk *CohereStreamEvent) ToBifrostResponsesStream(sequenceNumber int, sta // Handle text content delta if chunk.Delta.Message != nil && chunk.Delta.Message.Content != nil && chunk.Delta.Message.Content.CohereStreamContentObject != nil && chunk.Delta.Message.Content.CohereStreamContentObject.Text != nil && *chunk.Delta.Message.Content.CohereStreamContentObject.Text != "" { + // Accumulate text for done events + if state.TextBuffers[outputIndex] == nil { + state.TextBuffers[outputIndex] = &strings.Builder{} + } + state.TextBuffers[outputIndex].WriteString(*chunk.Delta.Message.Content.CohereStreamContentObject.Text) + // Emit output_text.delta (not reasoning_summary_text.delta for regular text) itemID := state.ItemIDs[outputIndex] response := &schemas.BifrostResponsesStreamResponse{ @@ -469,6 +505,12 @@ func (chunk *CohereStreamEvent) ToBifrostResponsesStream(sequenceNumber int, sta var responses []*schemas.BifrostResponsesStreamResponse isReasoning := state.ReasoningContentIndices[*chunk.Index] + // Grab accumulated text up front (empty string for reasoning blocks) + accText := "" + if buf := state.TextBuffers[outputIndex]; buf != nil { + accText = buf.String() + } + // Check if this content index is a reasoning block if isReasoning { // Emit reasoning_summary_text.done (reasoning equivalent of output_text.done) @@ -505,22 +547,22 @@ func (chunk *CohereStreamEvent) ToBifrostResponsesStream(sequenceNumber int, sta // Clear the reasoning content index tracking delete(state.ReasoningContentIndices, *chunk.Index) } else { - // Regular text block - emit output_text.done (without accumulated text, just the event) - emptyText := "" + // Regular text block - emit output_text.done with accumulated text responses = append(responses, &schemas.BifrostResponsesStreamResponse{ Type: schemas.ResponsesStreamResponseTypeOutputTextDone, SequenceNumber: sequenceNumber + len(responses), OutputIndex: schemas.Ptr(outputIndex), ContentIndex: chunk.Index, ItemID: &itemID, - Text: &emptyText, + Text: &accText, LogProbs: []schemas.ResponsesOutputMessageContentTextLogProb{}, }) - // Emit content_part.done + // Emit content_part.done with accumulated text + partText := accText part := &schemas.ResponsesMessageContentBlock{ Type: schemas.ResponsesOutputMessageContentTypeText, - Text: &emptyText, + Text: &partText, ResponsesOutputMessageContentText: &schemas.ResponsesOutputMessageContentText{ LogProbs: []schemas.ResponsesOutputMessageContentTextLogProb{}, Annotations: []schemas.ResponsesOutputMessageContentTextAnnotation{}, @@ -534,6 +576,7 @@ func (chunk *CohereStreamEvent) ToBifrostResponsesStream(sequenceNumber int, sta ItemID: &itemID, Part: part, }) + delete(state.TextBuffers, outputIndex) } // Emit output_item.done for all content blocks (text, reasoning, etc.) @@ -551,6 +594,20 @@ func (chunk *CohereStreamEvent) ToBifrostResponsesStream(sequenceNumber int, sta }, } } else { + contentBlocks := []schemas.ResponsesMessageContentBlock{} + if accText != "" { + itemText := accText + contentBlocks = []schemas.ResponsesMessageContentBlock{ + { + Type: schemas.ResponsesOutputMessageContentTypeText, + Text: &itemText, + ResponsesOutputMessageContentText: &schemas.ResponsesOutputMessageContentText{ + Annotations: []schemas.ResponsesOutputMessageContentTextAnnotation{}, + LogProbs: []schemas.ResponsesOutputMessageContentTextLogProb{}, + }, + }, + } + } messageType := schemas.ResponsesMessageTypeMessage role := schemas.ResponsesInputMessageRoleAssistant doneItem = &schemas.ResponsesMessage{ @@ -558,7 +615,7 @@ func (chunk *CohereStreamEvent) ToBifrostResponsesStream(sequenceNumber int, sta Role: &role, Status: &statusCompleted, Content: &schemas.ResponsesMessageContent{ - ContentBlocks: []schemas.ResponsesMessageContentBlock{}, + ContentBlocks: contentBlocks, }, } } @@ -640,6 +697,12 @@ func (chunk *CohereStreamEvent) ToBifrostResponsesStream(sequenceNumber int, sta }) } + // Accumulate tool plan text for done events + if state.TextBuffers[outputIndex] == nil { + state.TextBuffers[outputIndex] = &strings.Builder{} + } + state.TextBuffers[outputIndex].WriteString(*chunk.Delta.Message.ToolPlan) + // Emit output_text.delta (not reasoning_summary_text.delta) itemID := state.ItemIDs[outputIndex] response := &schemas.BifrostResponsesStreamResponse{ @@ -663,23 +726,27 @@ func (chunk *CohereStreamEvent) ToBifrostResponsesStream(sequenceNumber int, sta if state.ToolPlanOutputIndex != nil { outputIndex := *state.ToolPlanOutputIndex itemID := state.ItemIDs[outputIndex] + accText := "" + if buf := state.TextBuffers[outputIndex]; buf != nil { + accText = buf.String() + } - // Emit output_text.done (without accumulated text, just the event) - emptyText := "" + // Emit output_text.done with accumulated text responses = append(responses, &schemas.BifrostResponsesStreamResponse{ Type: schemas.ResponsesStreamResponseTypeOutputTextDone, SequenceNumber: sequenceNumber + len(responses), OutputIndex: schemas.Ptr(outputIndex), ContentIndex: schemas.Ptr(0), ItemID: &itemID, - Text: &emptyText, + Text: &accText, LogProbs: []schemas.ResponsesOutputMessageContentTextLogProb{}, }) - // Emit content_part.done + // Emit content_part.done with accumulated text + partText := accText part := &schemas.ResponsesMessageContentBlock{ Type: schemas.ResponsesOutputMessageContentTypeText, - Text: &emptyText, + Text: &partText, ResponsesOutputMessageContentText: &schemas.ResponsesOutputMessageContentText{ LogProbs: []schemas.ResponsesOutputMessageContentTextLogProb{}, Annotations: []schemas.ResponsesOutputMessageContentTextAnnotation{}, @@ -694,7 +761,8 @@ func (chunk *CohereStreamEvent) ToBifrostResponsesStream(sequenceNumber int, sta Part: part, }) - // Emit output_item.done + // Emit output_item.done with content blocks + itemText := accText statusCompleted := "completed" messageType := schemas.ResponsesMessageTypeMessage role := schemas.ResponsesInputMessageRoleAssistant @@ -703,7 +771,16 @@ func (chunk *CohereStreamEvent) ToBifrostResponsesStream(sequenceNumber int, sta Role: &role, Status: &statusCompleted, Content: &schemas.ResponsesMessageContent{ - ContentBlocks: []schemas.ResponsesMessageContentBlock{}, + ContentBlocks: []schemas.ResponsesMessageContentBlock{ + { + Type: schemas.ResponsesOutputMessageContentTypeText, + Text: &itemText, + ResponsesOutputMessageContentText: &schemas.ResponsesOutputMessageContentText{ + Annotations: []schemas.ResponsesOutputMessageContentTextAnnotation{}, + LogProbs: []schemas.ResponsesOutputMessageContentTextLogProb{}, + }, + }, + }, }, } if itemID != "" { @@ -716,6 +793,10 @@ func (chunk *CohereStreamEvent) ToBifrostResponsesStream(sequenceNumber int, sta ContentIndex: schemas.Ptr(0), Item: doneItem, }) + delete(state.TextBuffers, outputIndex) + if mapped, ok := state.ContentIndexToOutputIndex[0]; ok && mapped == outputIndex { + delete(state.ContentIndexToOutputIndex, 0) + } state.ToolPlanOutputIndex = nil // Mark as closed } diff --git a/core/providers/gemini/responses.go b/core/providers/gemini/responses.go index 83e7cc393f7..67b09eb1d97 100644 --- a/core/providers/gemini/responses.go +++ b/core/providers/gemini/responses.go @@ -1651,10 +1651,11 @@ func closeGeminiTextItem(state *GeminiResponsesStreamState, sequenceNumber int) LogProbs: []schemas.ResponsesOutputMessageContentTextLogProb{}, }) - // Emit content_part.done + // Emit content_part.done with accumulated text + partText := fullText part := &schemas.ResponsesMessageContentBlock{ Type: schemas.ResponsesOutputMessageContentTypeText, - Text: schemas.Ptr(""), + Text: &partText, ResponsesOutputMessageContentText: &schemas.ResponsesOutputMessageContentText{ LogProbs: []schemas.ResponsesOutputMessageContentTextLogProb{}, Annotations: []schemas.ResponsesOutputMessageContentTextAnnotation{}, @@ -1669,13 +1670,23 @@ func closeGeminiTextItem(state *GeminiResponsesStreamState, sequenceNumber int) Part: part, }) - // Emit output_item.done + // Emit output_item.done with content blocks + itemText := fullText doneItem := &schemas.ResponsesMessage{ Type: schemas.Ptr(schemas.ResponsesMessageTypeMessage), Role: schemas.Ptr(schemas.ResponsesInputMessageRoleAssistant), Status: schemas.Ptr("completed"), Content: &schemas.ResponsesMessageContent{ - ContentBlocks: []schemas.ResponsesMessageContentBlock{}, + ContentBlocks: []schemas.ResponsesMessageContentBlock{ + { + Type: schemas.ResponsesOutputMessageContentTypeText, + Text: &itemText, + ResponsesOutputMessageContentText: &schemas.ResponsesOutputMessageContentText{ + Annotations: []schemas.ResponsesOutputMessageContentTextAnnotation{}, + LogProbs: []schemas.ResponsesOutputMessageContentTextLogProb{}, + }, + }, + }, }, } if itemID != "" { From 63bdb2b559deb7c60c582d4f528998b4200755d8 Mon Sep 17 00:00:00 2001 From: Suresh Chaudhary Date: Fri, 29 May 2026 16:47:38 +0530 Subject: [PATCH 2/2] fix: placeholder for tool arguments in Anthropic responses --- core/providers/anthropic/responses.go | 12 ++ core/providers/anthropic/toolinput_test.go | 155 +++++++++++++++++++++ 2 files changed, 167 insertions(+) create mode 100644 core/providers/anthropic/toolinput_test.go diff --git a/core/providers/anthropic/responses.go b/core/providers/anthropic/responses.go index 98c6fdbbd43..6ceb8d2758e 100644 --- a/core/providers/anthropic/responses.go +++ b/core/providers/anthropic/responses.go @@ -4406,6 +4406,10 @@ func convertBifrostFunctionCallToAnthropicToolUse(ctx *schemas.BifrostContext, m } } toolUseBlock.Input = parseJSONInput(argumentsJSON) + } else { + // Anthropic requires input to always be present on tool_use blocks; + // default to an empty object for tools that take no arguments. + toolUseBlock.Input = json.RawMessage("{}") } return &toolUseBlock @@ -4574,6 +4578,10 @@ func convertBifrostMCPCallToAnthropicToolUse(msg *schemas.ResponsesMessage) *Ant // Parse arguments as JSON input if msg.ResponsesToolMessage.Arguments != nil && *msg.ResponsesToolMessage.Arguments != "" { toolUseBlock.Input = parseJSONInput(*msg.ResponsesToolMessage.Arguments) + } else { + // Anthropic requires input to always be present on tool_use blocks; + // default to an empty object for tools that take no arguments. + toolUseBlock.Input = json.RawMessage("{}") } return &toolUseBlock @@ -4620,6 +4628,10 @@ func convertBifrostMCPApprovalToAnthropicToolUse(msg *schemas.ResponsesMessage) // Parse arguments as JSON input if msg.ResponsesToolMessage.Arguments != nil && *msg.ResponsesToolMessage.Arguments != "" { toolUseBlock.Input = parseJSONInput(*msg.ResponsesToolMessage.Arguments) + } else { + // Anthropic requires input to always be present on tool_use blocks; + // default to an empty object for tools that take no arguments. + toolUseBlock.Input = json.RawMessage("{}") } return &toolUseBlock diff --git a/core/providers/anthropic/toolinput_test.go b/core/providers/anthropic/toolinput_test.go new file mode 100644 index 00000000000..63d75de4fd7 --- /dev/null +++ b/core/providers/anthropic/toolinput_test.go @@ -0,0 +1,155 @@ +package anthropic + +import ( + "context" + "testing" + + "github.com/maximhq/bifrost/core/schemas" +) + +// TestConvertBifrostFunctionCallToAnthropicToolUse_Input verifies that the +// tool_use block always carries an "input" object, defaulting to "{}" when the +// function call has nil or empty arguments (tools that take no arguments). +func TestConvertBifrostFunctionCallToAnthropicToolUse_Input(t *testing.T) { + t.Parallel() + + ctx, cancel := schemas.NewBifrostContextWithCancel(context.Background()) + defer cancel() + + tests := []struct { + name string + arguments *string + wantInput string + }{ + {name: "nil arguments", arguments: nil, wantInput: "{}"}, + {name: "empty arguments", arguments: schemas.Ptr(""), wantInput: "{}"}, + {name: "populated arguments", arguments: schemas.Ptr(`{"foo":"bar"}`), wantInput: `{"foo":"bar"}`}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + + msg := &schemas.ResponsesMessage{ + Type: schemas.Ptr(schemas.ResponsesMessageTypeFunctionCall), + ResponsesToolMessage: &schemas.ResponsesToolMessage{ + CallID: schemas.Ptr("toolu_fn_test"), + Name: schemas.Ptr("get_workspace_id"), + Arguments: tt.arguments, + }, + } + + block := convertBifrostFunctionCallToAnthropicToolUse(ctx, msg) + if block == nil { + t.Fatal("expected non-nil tool_use block") + } + if block.Type != AnthropicContentBlockTypeToolUse { + t.Errorf("block.Type = %v, want %v", block.Type, AnthropicContentBlockTypeToolUse) + } + if block.Input == nil { + t.Fatal("expected non-nil Input") + } + if string(block.Input) != tt.wantInput { + t.Errorf("Input = %s, want %s", block.Input, tt.wantInput) + } + }) + } +} + +// TestConvertBifrostMCPCallToAnthropicToolUse_Input verifies that the +// mcp_tool_use block always carries an "input" object, defaulting to "{}" when +// the MCP call has nil or empty arguments. +func TestConvertBifrostMCPCallToAnthropicToolUse_Input(t *testing.T) { + t.Parallel() + + tests := []struct { + name string + arguments *string + wantInput string + }{ + {name: "nil arguments", arguments: nil, wantInput: "{}"}, + {name: "empty arguments", arguments: schemas.Ptr(""), wantInput: "{}"}, + {name: "populated arguments", arguments: schemas.Ptr(`{"foo":"bar"}`), wantInput: `{"foo":"bar"}`}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + + msg := &schemas.ResponsesMessage{ + ID: schemas.Ptr("mcp_call_test"), + Type: schemas.Ptr(schemas.ResponsesMessageTypeMCPCall), + ResponsesToolMessage: &schemas.ResponsesToolMessage{ + Name: schemas.Ptr("maximsse-get-maxim-workspace-id"), + Arguments: tt.arguments, + ResponsesMCPToolCall: &schemas.ResponsesMCPToolCall{ + ServerLabel: "maximsse", + }, + }, + } + + block := convertBifrostMCPCallToAnthropicToolUse(msg) + if block == nil { + t.Fatal("expected non-nil mcp_tool_use block") + } + if block.Type != AnthropicContentBlockTypeMCPToolUse { + t.Errorf("block.Type = %v, want %v", block.Type, AnthropicContentBlockTypeMCPToolUse) + } + if block.Input == nil { + t.Fatal("expected non-nil Input") + } + if string(block.Input) != tt.wantInput { + t.Errorf("Input = %s, want %s", block.Input, tt.wantInput) + } + }) + } +} + +// TestConvertBifrostMCPApprovalToAnthropicToolUse_Input verifies that the +// mcp_tool_use block produced for an MCP approval request always carries an +// "input" object, defaulting to "{}" when arguments are nil or empty. +func TestConvertBifrostMCPApprovalToAnthropicToolUse_Input(t *testing.T) { + t.Parallel() + + tests := []struct { + name string + arguments *string + wantInput string + }{ + {name: "nil arguments", arguments: nil, wantInput: "{}"}, + {name: "empty arguments", arguments: schemas.Ptr(""), wantInput: "{}"}, + {name: "populated arguments", arguments: schemas.Ptr(`{"foo":"bar"}`), wantInput: `{"foo":"bar"}`}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + + msg := &schemas.ResponsesMessage{ + ID: schemas.Ptr("mcp_approval_test"), + Type: schemas.Ptr(schemas.ResponsesMessageTypeMCPApprovalRequest), + ResponsesToolMessage: &schemas.ResponsesToolMessage{ + Name: schemas.Ptr("maximsse-get-maxim-workspace-id"), + Arguments: tt.arguments, + ResponsesMCPToolCall: &schemas.ResponsesMCPToolCall{ + ServerLabel: "maximsse", + }, + }, + } + + block := convertBifrostMCPApprovalToAnthropicToolUse(msg) + if block == nil { + t.Fatal("expected non-nil mcp_tool_use block") + } + if block.Type != AnthropicContentBlockTypeMCPToolUse { + t.Errorf("block.Type = %v, want %v", block.Type, AnthropicContentBlockTypeMCPToolUse) + } + if block.Input == nil { + t.Fatal("expected non-nil Input") + } + if string(block.Input) != tt.wantInput { + t.Errorf("Input = %s, want %s", block.Input, tt.wantInput) + } + }) + } +}