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
1 change: 1 addition & 0 deletions core/changelog.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
- feat: add `additional_tools` message type support
- feat: force single-region config in Vertex key config
- feat: phase-scoped redaction and revealing with transient redaction data field, plus `ClearPausedStreamBuffer` for pause-accumulate stream flows
- fix: complete visible thinking items on the streaming Responses surface so reasoning_summary_text.done, output_item.done, and response.completed carry the accumulated reasoning text and signature (thanks [@fus3r](https://github.com/fus3r)!)
- fix: map web search options to Google Search grounding in the Gemini API
- fix: parse `SecretVar` JSON with `ref`/`env_var` fields even when `value` is absent
- fix: cap max reasoning effort in OpenAI
Expand Down
70 changes: 66 additions & 4 deletions core/providers/anthropic/responses.go
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,7 @@ type AnthropicResponsesStreamState struct {
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
ReasoningTextBuffers map[int]*strings.Builder // Maps output_index to accumulated reasoning text 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
Expand Down Expand Up @@ -149,6 +150,7 @@ var anthropicResponsesStreamStatePool = sync.Pool{
CompactionContentIndices: make(map[int]*schemas.CacheControl),
OutputItems: make(map[int]*schemas.ResponsesMessage),
TextBuffers: make(map[int]*strings.Builder),
ReasoningTextBuffers: make(map[int]*strings.Builder),
CurrentOutputIndex: 0,
CreatedAt: int(time.Now().Unix()),
HasEmittedCreated: false,
Expand Down Expand Up @@ -343,6 +345,11 @@ func AcquireAnthropicResponsesStreamState() *AnthropicResponsesStreamState {
} else {
clear(state.TextBuffers)
}
if state.ReasoningTextBuffers == nil {
state.ReasoningTextBuffers = make(map[int]*strings.Builder)
} else {
clear(state.ReasoningTextBuffers)
}
if state.CompactionContentIndices == nil {
state.CompactionContentIndices = make(map[int]*schemas.CacheControl)
} else {
Expand Down Expand Up @@ -436,6 +443,7 @@ func (state *AnthropicResponsesStreamState) flush() {
state.TextContentIndices = nil
state.ReasoningContentIndices = nil
state.TextBuffers = nil
state.ReasoningTextBuffers = nil
state.CompactionContentIndices = nil
state.OutputItems = nil
state.CurrentOutputIndex = 0
Expand Down Expand Up @@ -1191,6 +1199,12 @@ func (chunk *AnthropicStreamEvent) ToBifrostResponsesStream(ctx context.Context,
state.ReasoningContentIndices[*chunk.Index] = true
}

// Persist into OutputItems so content_block_stop can fold the
// accumulated text and signature into the matching output_item.done
// and response.completed carries the completed reasoning item.
itemCopy := *item
state.OutputItems[outputIndex] = &itemCopy

var responses []*schemas.BifrostResponsesStreamResponse

// Emit output_item.added
Expand Down Expand Up @@ -1395,6 +1409,16 @@ func (chunk *AnthropicStreamEvent) ToBifrostResponsesStream(ctx context.Context,
case AnthropicStreamDeltaTypeThinking:
// Reasoning/thinking content delta
if chunk.Delta.Thinking != nil && *chunk.Delta.Thinking != "" {
// Accumulate the text so content_block_stop can emit it on
// reasoning_summary_text.done and the completed reasoning item
if state.ReasoningTextBuffers == nil {
state.ReasoningTextBuffers = make(map[int]*strings.Builder)
}
if state.ReasoningTextBuffers[outputIndex] == nil {
state.ReasoningTextBuffers[outputIndex] = &strings.Builder{}
}
state.ReasoningTextBuffers[outputIndex].WriteString(*chunk.Delta.Thinking)

itemID := state.ItemIDs[outputIndex]
response := &schemas.BifrostResponsesStreamResponse{
Type: schemas.ResponsesStreamResponseTypeReasoningSummaryTextDelta,
Expand Down Expand Up @@ -1958,26 +1982,64 @@ func (chunk *AnthropicStreamEvent) ToBifrostResponsesStream(ctx context.Context,

// Check if this content index is a reasoning block
if state.ReasoningContentIndices[*chunk.Index] {
// Emit reasoning_summary_text.done (reasoning equivalent of output_text.done)
emptyText := ""
// Capture the accumulated reasoning text and signature
accReasoning := ""
if buf := state.ReasoningTextBuffers[outputIndex]; buf != nil {
accReasoning = buf.String()
delete(state.ReasoningTextBuffers, outputIndex)
}
var signature *string
if sig, ok := state.ReasoningSignatures[outputIndex]; ok && sig != "" {
sigCopy := sig
signature = &sigCopy
delete(state.ReasoningSignatures, outputIndex)
}

// Fold them into the stored item with the same shape the
// non-streaming converter produces (a reasoning_text content
// block carrying text + signature), so output_item.done and
// response.completed carry a replayable reasoning item.
if storedItem, exists := state.OutputItems[outputIndex]; exists {
textCopy := accReasoning
storedItem.Content = &schemas.ResponsesMessageContent{
ContentBlocks: []schemas.ResponsesMessageContentBlock{
{
Type: schemas.ResponsesOutputMessageContentTypeReasoning,
Text: &textCopy,
Signature: signature,
},
},
}
storedItem.ResponsesReasoning = nil
}

// Emit reasoning_summary_text.done (reasoning equivalent of
// output_text.done) with the full accumulated text
doneText := accReasoning
reasoningDoneResponse := &schemas.BifrostResponsesStreamResponse{
Type: schemas.ResponsesStreamResponseTypeReasoningSummaryTextDone,
SequenceNumber: sequenceNumber + len(responses),
OutputIndex: schemas.Ptr(outputIndex),
ContentIndex: chunk.Index,
Text: &emptyText,
Text: &doneText,
}
if itemID != "" {
reasoningDoneResponse.ItemID = &itemID
}
responses = append(responses, reasoningDoneResponse)

// Emit content_part.done for reasoning
// Emit content_part.done for reasoning with the completed part
partText := accReasoning
partDoneResponse := &schemas.BifrostResponsesStreamResponse{
Type: schemas.ResponsesStreamResponseTypeContentPartDone,
SequenceNumber: sequenceNumber + len(responses),
OutputIndex: schemas.Ptr(outputIndex),
ContentIndex: chunk.Index,
Part: &schemas.ResponsesMessageContentBlock{
Type: schemas.ResponsesOutputMessageContentTypeReasoning,
Text: &partText,
Signature: signature,
},
}
if itemID != "" {
partDoneResponse.ItemID = &itemID
Expand Down
Loading