From 2bbe76710152b45e101b18cd39c031391d5c6b19 Mon Sep 17 00:00:00 2001 From: Sahil Choudhary Date: Tue, 14 Jul 2026 16:16:41 +0530 Subject: [PATCH 1/4] add support for clickhouse for enterprise logstore tables --- framework/logstore/clickhouseextension.go | 84 ++++++++++++++++++++++ framework/logstore/clickhousestore_test.go | 79 ++++++++++++++++++++ 2 files changed, 163 insertions(+) create mode 100644 framework/logstore/clickhouseextension.go diff --git a/framework/logstore/clickhouseextension.go b/framework/logstore/clickhouseextension.go new file mode 100644 index 00000000000..1a6ccb7c464 --- /dev/null +++ b/framework/logstore/clickhouseextension.go @@ -0,0 +1,84 @@ +package logstore + +import ( + "context" + "fmt" + "regexp" + "strings" +) + +var clickHouseExtensionTableNamePattern = regexp.MustCompile(`^[A-Za-z_][A-Za-z0-9_]*$`) + +var clickHouseReservedTableNames = map[string]struct{}{ + strings.ToLower((AsyncJob{}).TableName()): {}, + strings.ToLower((Log{}).TableName()): {}, + strings.ToLower((MCPToolLog{}).TableName()): {}, +} + +// ClickHouseExtensionTableOptions defines the schema shape for an +// extension-owned table. All DDL fragments are trusted, code-owned values and +// must never be populated from configuration or user input. +type ClickHouseExtensionTableOptions struct { + Table string + PartitionBy string + OrderBy string + TTL string + SkipIndexes []string +} + +type clickHouseSchemaStore interface { + EnsureClickHouseTable(ctx context.Context, model any, opts ClickHouseExtensionTableOptions) error +} + +var ( + _ clickHouseSchemaStore = (*ClickHouseLogStore)(nil) + _ clickHouseSchemaStore = (*HybridLogStore)(nil) +) + +// EnsureClickHouseTable creates an extension table and reconciles newly added +// model columns. It preserves the configured cluster and replication settings. +func (s *ClickHouseLogStore) EnsureClickHouseTable(ctx context.Context, model any, opts ClickHouseExtensionTableOptions) error { + if s == nil || s.RDBLogStore == nil || s.db == nil { + return fmt.Errorf("clickhouse: logstore is not initialized") + } + if err := validateClickHouseExtensionTableOptions(opts); err != nil { + return err + } + + tableOpts := chTableOpts{ + table: opts.Table, + partitionBy: opts.PartitionBy, + orderBy: opts.OrderBy, + ttl: opts.TTL, + skipIndexes: append([]string(nil), opts.SkipIndexes...), + } + if err := clickhouseCreateTable(ctx, s.db, model, tableOpts, s.cluster); err != nil { + return fmt.Errorf("clickhouse: create extension table %s: %w", opts.Table, err) + } + return clickhouseReconcileColumns(ctx, s.db, model, opts.Table, s.cluster, s.logger) +} + +// EnsureClickHouseTable delegates extension-table schema management to a +// ClickHouse logstore wrapped by hybrid object storage. +func (h *HybridLogStore) EnsureClickHouseTable(ctx context.Context, model any, opts ClickHouseExtensionTableOptions) error { + schemaStore, ok := h.inner.(clickHouseSchemaStore) + if !ok { + return fmt.Errorf("logstore does not support ClickHouse extension tables") + } + return schemaStore.EnsureClickHouseTable(ctx, model, opts) +} + +// validateClickHouseExtensionTableOptions validates identifiers and invariants +// that can be checked without attempting to parse ClickHouse SQL expressions. +func validateClickHouseExtensionTableOptions(opts ClickHouseExtensionTableOptions) error { + if !clickHouseExtensionTableNamePattern.MatchString(opts.Table) { + return fmt.Errorf("clickhouse: invalid extension table name %q", opts.Table) + } + if _, reserved := clickHouseReservedTableNames[strings.ToLower(opts.Table)]; reserved { + return fmt.Errorf("clickhouse: extension table name %q is reserved", opts.Table) + } + if strings.TrimSpace(opts.OrderBy) == "" { + return fmt.Errorf("clickhouse: extension table order by is required") + } + return nil +} diff --git a/framework/logstore/clickhousestore_test.go b/framework/logstore/clickhousestore_test.go index b62099cf02e..32ac90f8626 100644 --- a/framework/logstore/clickhousestore_test.go +++ b/framework/logstore/clickhousestore_test.go @@ -68,6 +68,85 @@ func chTestLog(id string, ts time.Time) *Log { } } +type clickHouseExtensionTestRow struct { + ID string + Value string + CreatedAt time.Time +} + +func (clickHouseExtensionTestRow) TableName() string { return "extension_test_events" } + +func TestClickHouseEnsureExtensionTable(t *testing.T) { + store := trySetupClickHouseStore(t) + ctx := context.Background() + require.Error(t, store.EnsureClickHouseTable(ctx, &clickHouseExtensionTestRow{}, ClickHouseExtensionTableOptions{ + Table: "invalid`; DROP TABLE logs", OrderBy: "id", + })) + require.NoError(t, store.db.Exec("DROP TABLE IF EXISTS extension_test_events").Error) + t.Cleanup(func() { + _ = store.db.Exec("DROP TABLE IF EXISTS extension_test_events").Error + }) + + opts := ClickHouseExtensionTableOptions{ + Table: "extension_test_events", + PartitionBy: "toYYYYMM(created_at)", + OrderBy: "(created_at, id)", + TTL: "toDateTime(created_at) + INTERVAL 30 DAY", + SkipIndexes: []string{"INDEX idx_extension_value lower(value) TYPE bloom_filter(0.01) GRANULARITY 1"}, + } + require.NoError(t, store.EnsureClickHouseTable(ctx, &clickHouseExtensionTestRow{}, opts)) + hybrid := &HybridLogStore{inner: store} + require.NoError(t, hybrid.EnsureClickHouseTable(ctx, &clickHouseExtensionTestRow{}, opts)) + + row := clickHouseExtensionTestRow{ID: "event-1", Value: "matched", CreatedAt: time.Now().UTC()} + require.NoError(t, store.db.WithContext(ctx).Create(&row).Error) + + var count int64 + require.NoError(t, store.db.WithContext(ctx).Model(&clickHouseExtensionTestRow{}).Where("value = ?", "matched").Count(&count).Error) + assert.Equal(t, int64(1), count) + + var createQuery string + require.NoError(t, store.db.WithContext(ctx). + Raw("SELECT create_table_query FROM system.tables WHERE database = currentDatabase() AND name = ?", opts.Table). + Scan(&createQuery).Error) + assert.Contains(t, createQuery, "ReplacingMergeTree") + assert.Contains(t, createQuery, "TTL") + assert.Contains(t, createQuery, "idx_extension_value") +} + +func TestHybridEnsureClickHouseExtensionTableRejectsUnsupportedInner(t *testing.T) { + hybrid := &HybridLogStore{inner: &RDBLogStore{}} + err := hybrid.EnsureClickHouseTable(context.Background(), &clickHouseExtensionTestRow{}, ClickHouseExtensionTableOptions{ + Table: "extension_events", OrderBy: "id", + }) + require.EqualError(t, err, "logstore does not support ClickHouse extension tables") +} + +func TestValidateClickHouseExtensionTableOptions(t *testing.T) { + valid := ClickHouseExtensionTableOptions{Table: "extension_events", OrderBy: "(created_at, id)"} + require.NoError(t, validateClickHouseExtensionTableOptions(valid)) + invalidName := valid + invalidName.Table = "invalid`; DROP TABLE logs" + require.ErrorContains(t, validateClickHouseExtensionTableOptions(invalidName), "invalid extension table name") + + for _, table := range []string{"logs", "mcp_tool_logs", "async_jobs", "LOGS"} { + opts := valid + opts.Table = table + require.ErrorContains(t, validateClickHouseExtensionTableOptions(opts), "reserved") + } + + missingOrderBy := valid + missingOrderBy.OrderBy = " " + require.ErrorContains(t, validateClickHouseExtensionTableOptions(missingOrderBy), "order by is required") + + clickHouseSyntax := valid + clickHouseSyntax.TTL = "toDateTime(created_at) + INTERVAL 30 DAY TO VOLUME 'cold'" + clickHouseSyntax.SkipIndexes = []string{ + "INDEX idx_lower_value lower(value) TYPE bloom_filter(0.01) GRANULARITY 1", + } + require.NoError(t, validateClickHouseExtensionTableOptions(clickHouseSyntax)) +} + // chCountRows counts logical rows visible for an id; with the connection-level // final=1 setting, ReplacingMergeTree duplicates must collapse to one. func chCountRows(t *testing.T, db *gorm.DB, table, id string) int64 { From dd5df2bb8db2a957b3b1120b65fe749948c9f0d6 Mon Sep 17 00:00:00 2001 From: Akshay Deo Date: Tue, 14 Jul 2026 22:07:29 -0700 Subject: [PATCH 2/4] Revert "add support for clickhouse for enterprise logstore tables" (#5225) --- framework/logstore/clickhouseextension.go | 84 ---------------------- framework/logstore/clickhousestore_test.go | 79 -------------------- 2 files changed, 163 deletions(-) delete mode 100644 framework/logstore/clickhouseextension.go diff --git a/framework/logstore/clickhouseextension.go b/framework/logstore/clickhouseextension.go deleted file mode 100644 index 1a6ccb7c464..00000000000 --- a/framework/logstore/clickhouseextension.go +++ /dev/null @@ -1,84 +0,0 @@ -package logstore - -import ( - "context" - "fmt" - "regexp" - "strings" -) - -var clickHouseExtensionTableNamePattern = regexp.MustCompile(`^[A-Za-z_][A-Za-z0-9_]*$`) - -var clickHouseReservedTableNames = map[string]struct{}{ - strings.ToLower((AsyncJob{}).TableName()): {}, - strings.ToLower((Log{}).TableName()): {}, - strings.ToLower((MCPToolLog{}).TableName()): {}, -} - -// ClickHouseExtensionTableOptions defines the schema shape for an -// extension-owned table. All DDL fragments are trusted, code-owned values and -// must never be populated from configuration or user input. -type ClickHouseExtensionTableOptions struct { - Table string - PartitionBy string - OrderBy string - TTL string - SkipIndexes []string -} - -type clickHouseSchemaStore interface { - EnsureClickHouseTable(ctx context.Context, model any, opts ClickHouseExtensionTableOptions) error -} - -var ( - _ clickHouseSchemaStore = (*ClickHouseLogStore)(nil) - _ clickHouseSchemaStore = (*HybridLogStore)(nil) -) - -// EnsureClickHouseTable creates an extension table and reconciles newly added -// model columns. It preserves the configured cluster and replication settings. -func (s *ClickHouseLogStore) EnsureClickHouseTable(ctx context.Context, model any, opts ClickHouseExtensionTableOptions) error { - if s == nil || s.RDBLogStore == nil || s.db == nil { - return fmt.Errorf("clickhouse: logstore is not initialized") - } - if err := validateClickHouseExtensionTableOptions(opts); err != nil { - return err - } - - tableOpts := chTableOpts{ - table: opts.Table, - partitionBy: opts.PartitionBy, - orderBy: opts.OrderBy, - ttl: opts.TTL, - skipIndexes: append([]string(nil), opts.SkipIndexes...), - } - if err := clickhouseCreateTable(ctx, s.db, model, tableOpts, s.cluster); err != nil { - return fmt.Errorf("clickhouse: create extension table %s: %w", opts.Table, err) - } - return clickhouseReconcileColumns(ctx, s.db, model, opts.Table, s.cluster, s.logger) -} - -// EnsureClickHouseTable delegates extension-table schema management to a -// ClickHouse logstore wrapped by hybrid object storage. -func (h *HybridLogStore) EnsureClickHouseTable(ctx context.Context, model any, opts ClickHouseExtensionTableOptions) error { - schemaStore, ok := h.inner.(clickHouseSchemaStore) - if !ok { - return fmt.Errorf("logstore does not support ClickHouse extension tables") - } - return schemaStore.EnsureClickHouseTable(ctx, model, opts) -} - -// validateClickHouseExtensionTableOptions validates identifiers and invariants -// that can be checked without attempting to parse ClickHouse SQL expressions. -func validateClickHouseExtensionTableOptions(opts ClickHouseExtensionTableOptions) error { - if !clickHouseExtensionTableNamePattern.MatchString(opts.Table) { - return fmt.Errorf("clickhouse: invalid extension table name %q", opts.Table) - } - if _, reserved := clickHouseReservedTableNames[strings.ToLower(opts.Table)]; reserved { - return fmt.Errorf("clickhouse: extension table name %q is reserved", opts.Table) - } - if strings.TrimSpace(opts.OrderBy) == "" { - return fmt.Errorf("clickhouse: extension table order by is required") - } - return nil -} diff --git a/framework/logstore/clickhousestore_test.go b/framework/logstore/clickhousestore_test.go index 32ac90f8626..b62099cf02e 100644 --- a/framework/logstore/clickhousestore_test.go +++ b/framework/logstore/clickhousestore_test.go @@ -68,85 +68,6 @@ func chTestLog(id string, ts time.Time) *Log { } } -type clickHouseExtensionTestRow struct { - ID string - Value string - CreatedAt time.Time -} - -func (clickHouseExtensionTestRow) TableName() string { return "extension_test_events" } - -func TestClickHouseEnsureExtensionTable(t *testing.T) { - store := trySetupClickHouseStore(t) - ctx := context.Background() - require.Error(t, store.EnsureClickHouseTable(ctx, &clickHouseExtensionTestRow{}, ClickHouseExtensionTableOptions{ - Table: "invalid`; DROP TABLE logs", OrderBy: "id", - })) - require.NoError(t, store.db.Exec("DROP TABLE IF EXISTS extension_test_events").Error) - t.Cleanup(func() { - _ = store.db.Exec("DROP TABLE IF EXISTS extension_test_events").Error - }) - - opts := ClickHouseExtensionTableOptions{ - Table: "extension_test_events", - PartitionBy: "toYYYYMM(created_at)", - OrderBy: "(created_at, id)", - TTL: "toDateTime(created_at) + INTERVAL 30 DAY", - SkipIndexes: []string{"INDEX idx_extension_value lower(value) TYPE bloom_filter(0.01) GRANULARITY 1"}, - } - require.NoError(t, store.EnsureClickHouseTable(ctx, &clickHouseExtensionTestRow{}, opts)) - hybrid := &HybridLogStore{inner: store} - require.NoError(t, hybrid.EnsureClickHouseTable(ctx, &clickHouseExtensionTestRow{}, opts)) - - row := clickHouseExtensionTestRow{ID: "event-1", Value: "matched", CreatedAt: time.Now().UTC()} - require.NoError(t, store.db.WithContext(ctx).Create(&row).Error) - - var count int64 - require.NoError(t, store.db.WithContext(ctx).Model(&clickHouseExtensionTestRow{}).Where("value = ?", "matched").Count(&count).Error) - assert.Equal(t, int64(1), count) - - var createQuery string - require.NoError(t, store.db.WithContext(ctx). - Raw("SELECT create_table_query FROM system.tables WHERE database = currentDatabase() AND name = ?", opts.Table). - Scan(&createQuery).Error) - assert.Contains(t, createQuery, "ReplacingMergeTree") - assert.Contains(t, createQuery, "TTL") - assert.Contains(t, createQuery, "idx_extension_value") -} - -func TestHybridEnsureClickHouseExtensionTableRejectsUnsupportedInner(t *testing.T) { - hybrid := &HybridLogStore{inner: &RDBLogStore{}} - err := hybrid.EnsureClickHouseTable(context.Background(), &clickHouseExtensionTestRow{}, ClickHouseExtensionTableOptions{ - Table: "extension_events", OrderBy: "id", - }) - require.EqualError(t, err, "logstore does not support ClickHouse extension tables") -} - -func TestValidateClickHouseExtensionTableOptions(t *testing.T) { - valid := ClickHouseExtensionTableOptions{Table: "extension_events", OrderBy: "(created_at, id)"} - require.NoError(t, validateClickHouseExtensionTableOptions(valid)) - invalidName := valid - invalidName.Table = "invalid`; DROP TABLE logs" - require.ErrorContains(t, validateClickHouseExtensionTableOptions(invalidName), "invalid extension table name") - - for _, table := range []string{"logs", "mcp_tool_logs", "async_jobs", "LOGS"} { - opts := valid - opts.Table = table - require.ErrorContains(t, validateClickHouseExtensionTableOptions(opts), "reserved") - } - - missingOrderBy := valid - missingOrderBy.OrderBy = " " - require.ErrorContains(t, validateClickHouseExtensionTableOptions(missingOrderBy), "order by is required") - - clickHouseSyntax := valid - clickHouseSyntax.TTL = "toDateTime(created_at) + INTERVAL 30 DAY TO VOLUME 'cold'" - clickHouseSyntax.SkipIndexes = []string{ - "INDEX idx_lower_value lower(value) TYPE bloom_filter(0.01) GRANULARITY 1", - } - require.NoError(t, validateClickHouseExtensionTableOptions(clickHouseSyntax)) -} - // chCountRows counts logical rows visible for an id; with the connection-level // final=1 setting, ReplacingMergeTree duplicates must collapse to one. func chCountRows(t *testing.T, db *gorm.DB, table, id string) int64 { From c0909f9752156121c6c775694df6e656a6ad3860 Mon Sep 17 00:00:00 2001 From: Akshay Deo Date: Wed, 15 Jul 2026 11:07:10 +0530 Subject: [PATCH 3/4] Delete 1 Signed-off-by: Akshay Deo --- 1 | 36 ------------------------------------ 1 file changed, 36 deletions(-) delete mode 100644 1 diff --git a/1 b/1 deleted file mode 100644 index 9ff974df65a..00000000000 --- a/1 +++ /dev/null @@ -1,36 +0,0 @@ -changelogs commit - -# Please enter the commit message for your changes. Lines starting -# with '#' will be ignored, and an empty message aborts the commit. -# -# HEAD detached at 0cb488798 -# Changes to be committed: -# modified: core/changelog.md -# modified: core/version -# modified: framework/changelog.md -# modified: framework/version -# modified: plugins/compat/changelog.md -# modified: plugins/compat/version -# modified: plugins/governance/changelog.md -# modified: plugins/governance/version -# modified: plugins/jsonparser/changelog.md -# modified: plugins/jsonparser/version -# modified: plugins/logging/changelog.md -# modified: plugins/logging/version -# modified: plugins/maxim/changelog.md -# modified: plugins/maxim/version -# modified: plugins/mocker/changelog.md -# modified: plugins/mocker/version -# modified: plugins/modelcatalogresolver/changelog.md -# modified: plugins/modelcatalogresolver/version -# modified: plugins/otel/changelog.md -# modified: plugins/otel/version -# modified: plugins/prompts/changelog.md -# modified: plugins/prompts/version -# modified: plugins/semanticcache/changelog.md -# modified: plugins/semanticcache/version -# modified: plugins/telemetry/changelog.md -# modified: plugins/telemetry/version -# modified: transports/changelog.md -# modified: transports/version -# From ef380f925c8476f1675e5f157d3e5b3033c7623a Mon Sep 17 00:00:00 2001 From: fus3r Date: Mon, 13 Jul 2026 12:43:09 +0200 Subject: [PATCH 4/4] fix: complete visible thinking items on the streaming Responses surface --- core/changelog.md | 1 + core/providers/anthropic/responses.go | 70 ++- .../visiblethinkingresponses_test.go | 439 ++++++++++++++++++ 3 files changed, 506 insertions(+), 4 deletions(-) create mode 100644 core/providers/anthropic/visiblethinkingresponses_test.go diff --git a/core/changelog.md b/core/changelog.md index 6c5cdfc5757..3f752e2a8ee 100644 --- a/core/changelog.md +++ b/core/changelog.md @@ -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 diff --git a/core/providers/anthropic/responses.go b/core/providers/anthropic/responses.go index 5ac422919a2..fd56c16e7f5 100644 --- a/core/providers/anthropic/responses.go +++ b/core/providers/anthropic/responses.go @@ -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 @@ -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, @@ -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 { @@ -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 @@ -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 @@ -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, @@ -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 diff --git a/core/providers/anthropic/visiblethinkingresponses_test.go b/core/providers/anthropic/visiblethinkingresponses_test.go new file mode 100644 index 00000000000..e376b255912 --- /dev/null +++ b/core/providers/anthropic/visiblethinkingresponses_test.go @@ -0,0 +1,439 @@ +package anthropic + +import ( + "encoding/json" + "testing" + "time" + + "github.com/maximhq/bifrost/core/schemas" +) + +// visibleThinkingToolUseLifecycle returns the real Anthropic SSE shape of a +// response with extended thinking enabled: a thinking block (text deltas, then +// a signature delta) followed by a tool_use block. +func visibleThinkingToolUseLifecycle(thinkingParts []string, signature string) []*AnthropicStreamEvent { + stopReason := AnthropicStopReasonToolUse + events := []*AnthropicStreamEvent{ + { + Type: AnthropicStreamEventTypeMessageStart, + Message: &AnthropicMessageResponse{ + ID: "msg_visible_stream", + Model: "claude-sonnet-4-5-20250929", + }, + }, + { + Type: AnthropicStreamEventTypeContentBlockStart, + Index: schemas.Ptr(0), + ContentBlock: &AnthropicContentBlock{ + Type: AnthropicContentBlockTypeThinking, + }, + }, + } + for _, part := range thinkingParts { + events = append(events, &AnthropicStreamEvent{ + Type: AnthropicStreamEventTypeContentBlockDelta, + Index: schemas.Ptr(0), + Delta: &AnthropicStreamDelta{ + Type: AnthropicStreamDeltaTypeThinking, + Thinking: schemas.Ptr(part), + }, + }) + } + if signature != "" { + events = append(events, &AnthropicStreamEvent{ + Type: AnthropicStreamEventTypeContentBlockDelta, + Index: schemas.Ptr(0), + Delta: &AnthropicStreamDelta{ + Type: AnthropicStreamDeltaTypeSignature, + Signature: schemas.Ptr(signature), + }, + }) + } + events = append(events, + &AnthropicStreamEvent{Type: AnthropicStreamEventTypeContentBlockStop, Index: schemas.Ptr(0)}, + &AnthropicStreamEvent{ + Type: AnthropicStreamEventTypeContentBlockStart, + Index: schemas.Ptr(1), + ContentBlock: &AnthropicContentBlock{ + Type: AnthropicContentBlockTypeToolUse, + ID: schemas.Ptr("toolu_01"), + Name: schemas.Ptr("get_weather"), + }, + }, + &AnthropicStreamEvent{ + Type: AnthropicStreamEventTypeContentBlockDelta, + Index: schemas.Ptr(1), + Delta: &AnthropicStreamDelta{ + Type: AnthropicStreamDeltaTypeInputJSON, + PartialJSON: schemas.Ptr(`{"location":"Paris"}`), + }, + }, + &AnthropicStreamEvent{Type: AnthropicStreamEventTypeContentBlockStop, Index: schemas.Ptr(1)}, + &AnthropicStreamEvent{ + Type: AnthropicStreamEventTypeMessageDelta, + Delta: &AnthropicStreamDelta{StopReason: &stopReason}, + }, + &AnthropicStreamEvent{Type: AnthropicStreamEventTypeMessageStop}, + ) + return events +} + +// reasoningTextBlock returns the reasoning_text content block of a reasoning +// item, or nil when the item has none. +func reasoningTextBlock(item *schemas.ResponsesMessage) *schemas.ResponsesMessageContentBlock { + if item == nil || item.Content == nil { + return nil + } + for i := range item.Content.ContentBlocks { + if item.Content.ContentBlocks[i].Type == schemas.ResponsesOutputMessageContentTypeReasoning { + return &item.Content.ContentBlocks[i] + } + } + return nil +} + +// TestToBifrostResponsesStream_VisibleThinkingCompletion replays the streaming +// shape Anthropic uses for a visible thinking block and asserts the block's +// completion events carry the accumulated reasoning. Without the fold at +// content_block_stop, reasoning_summary_text.done carries an empty text, +// output_item.done degrades to an empty "message" shell that contradicts its +// own output_item.added, and response.completed omits the reasoning item. +func TestToBifrostResponsesStream_VisibleThinkingCompletion(t *testing.T) { + t.Parallel() + + const fullText = "Let me figure out which tool to call." + const signature = "sig-test-123" + responses := driveResponsesStream(t, visibleThinkingToolUseLifecycle( + []string{"Let me figure out ", "which tool to call."}, signature)) + + // reasoning_summary_text.done carries the full accumulated text + var summaryDone *schemas.BifrostResponsesStreamResponse + for _, r := range responses { + if r.Type == schemas.ResponsesStreamResponseTypeReasoningSummaryTextDone { + summaryDone = r + } + } + if summaryDone == nil { + t.Fatal("expected a reasoning_summary_text.done event") + } + if summaryDone.Text == nil || *summaryDone.Text != fullText { + t.Errorf("reasoning_summary_text.done text = %v, want %q", summaryDone.Text, fullText) + } + + // content_part.done for the reasoning block carries the completed part + var partDone *schemas.BifrostResponsesStreamResponse + for _, r := range responses { + if r.Type == schemas.ResponsesStreamResponseTypeContentPartDone && r.OutputIndex != nil && *r.OutputIndex == 0 { + partDone = r + } + } + if partDone == nil { + t.Fatal("expected a content_part.done event for the reasoning block") + } + if partDone.Part == nil || partDone.Part.Text == nil || *partDone.Part.Text != fullText { + t.Errorf("content_part.done part text = %+v, want %q", partDone.Part, fullText) + } + if partDone.Part != nil && (partDone.Part.Signature == nil || *partDone.Part.Signature != signature) { + t.Errorf("content_part.done part signature = %v, want %q", partDone.Part.Signature, signature) + } + + // output_item.done for output index 0 is a reasoning item matching its + // output_item.added, with the accumulated text and signature + var itemDone *schemas.ResponsesMessage + for _, r := range responses { + if r.Type == schemas.ResponsesStreamResponseTypeOutputItemDone && r.OutputIndex != nil && *r.OutputIndex == 0 { + itemDone = r.Item + } + } + if itemDone == nil { + t.Fatal("expected an output_item.done event for the reasoning block") + } + if itemDone.Type == nil { + t.Fatal("output_item.done item type is nil, want reasoning") + } + if *itemDone.Type != schemas.ResponsesMessageTypeReasoning { + t.Fatalf("output_item.done item type = %q, want reasoning", *itemDone.Type) + } + block := reasoningTextBlock(itemDone) + if block == nil { + t.Fatalf("output_item.done reasoning item has no reasoning_text content block: %+v", itemDone) + } + if block.Text == nil || *block.Text != fullText { + t.Errorf("output_item.done reasoning text = %v, want %q", block.Text, fullText) + } + if block.Signature == nil || *block.Signature != signature { + t.Errorf("output_item.done reasoning signature = %v, want %q", block.Signature, signature) + } + + // response.completed carries the reasoning item before the function call + var completed *schemas.BifrostResponsesResponse + for _, r := range responses { + if r.Type == schemas.ResponsesStreamResponseTypeCompleted && r.Response != nil { + completed = r.Response + } + } + if completed == nil { + t.Fatal("expected a response.completed event") + } + if len(completed.Output) != 2 { + t.Fatalf("response.completed output has %d items, want 2 (reasoning + function_call)", len(completed.Output)) + } + first := completed.Output[0] + if first.Type == nil { + t.Fatal("response.completed output[0] type is nil, want reasoning") + } + if *first.Type != schemas.ResponsesMessageTypeReasoning { + t.Fatalf("response.completed output[0] type = %q, want reasoning", *first.Type) + } + if b := reasoningTextBlock(&first); b == nil || b.Text == nil || *b.Text != fullText { + t.Errorf("response.completed reasoning item text = %+v, want %q", b, fullText) + } + if second := completed.Output[1]; second.Type == nil || *second.Type != schemas.ResponsesMessageTypeFunctionCall { + t.Errorf("response.completed output[1] type = %v, want function_call", second.Type) + } +} + +// TestToAnthropicResponsesRequest_ReplaysStreamedVisibleThinking closes the +// multi-turn loop for visible thinking: the completed reasoning item emitted +// on output_item.done is round-tripped through JSON the way a client echoes +// its history, then converted back into an Anthropic request. The replayed +// assistant turn must carry the thinking block, with its signature, before +// the tool_use block, so the reasoning survives in the echoed conversation +// instead of silently disappearing from the multi-turn history. +func TestToAnthropicResponsesRequest_ReplaysStreamedVisibleThinking(t *testing.T) { + t.Parallel() + + const fullText = "Let me figure out which tool to call." + const signature = "sig-test-123" + responses := driveResponsesStream(t, visibleThinkingToolUseLifecycle( + []string{"Let me figure out ", "which tool to call."}, signature)) + + var itemDone *schemas.ResponsesMessage + for _, r := range responses { + if r.Type == schemas.ResponsesStreamResponseTypeOutputItemDone && r.OutputIndex != nil && *r.OutputIndex == 0 { + itemDone = r.Item + } + } + if itemDone == nil { + t.Fatal("expected an output_item.done event for the reasoning block") + } + + // Client echo: the item leaves and re-enters Bifrost as JSON. + wire, err := json.Marshal(itemDone) + if err != nil { + t.Fatalf("marshal streamed item: %v", err) + } + var replayed schemas.ResponsesMessage + if err := json.Unmarshal(wire, &replayed); err != nil { + t.Fatalf("unmarshal client echo: %v", err) + } + + req := &schemas.BifrostResponsesRequest{ + Model: "claude-sonnet-4-5-20250929", + Input: []schemas.ResponsesMessage{replayed, functionCallItem()}, + } + ctx := schemas.NewBifrostContext(nil, time.Time{}) + out, err := ToAnthropicResponsesRequest(ctx, req) + if err != nil { + t.Fatalf("ToAnthropicResponsesRequest: %v", err) + } + + for _, msg := range out.Messages { + if msg.Role != AnthropicMessageRoleAssistant { + continue + } + thinkingIdx := -1 + toolUseIdx := -1 + for i, b := range msg.Content.ContentBlocks { + switch b.Type { + case AnthropicContentBlockTypeThinking: + if b.Thinking != nil && *b.Thinking == fullText { + thinkingIdx = i + if b.Signature == nil || *b.Signature != signature { + t.Errorf("replayed thinking block signature = %v, want %q", b.Signature, signature) + } + } + case AnthropicContentBlockTypeToolUse: + toolUseIdx = i + } + } + if toolUseIdx == -1 { + continue + } + if thinkingIdx == -1 { + t.Fatalf("assistant message with tool_use has no thinking block carrying the reasoning: %+v", msg.Content.ContentBlocks) + } + if thinkingIdx > toolUseIdx { + t.Fatalf("thinking (idx %d) must precede tool_use (idx %d)", thinkingIdx, toolUseIdx) + } + return + } + t.Fatalf("no assistant message with a tool_use block in the outgoing request (%d messages)", len(out.Messages)) +} + +// TestToBifrostResponsesStream_MultipleVisibleThinkingBlocks interleaves two +// thinking blocks around a tool_use block (the shape interleaved thinking +// produces) and asserts each completed reasoning item carries its own text +// and signature. +func TestToBifrostResponsesStream_MultipleVisibleThinkingBlocks(t *testing.T) { + t.Parallel() + + stopReason := AnthropicStopReasonToolUse + events := []*AnthropicStreamEvent{ + { + Type: AnthropicStreamEventTypeMessageStart, + Message: &AnthropicMessageResponse{ + ID: "msg_interleaved_stream", + Model: "claude-sonnet-4-5-20250929", + }, + }, + { + Type: AnthropicStreamEventTypeContentBlockStart, + Index: schemas.Ptr(0), + ContentBlock: &AnthropicContentBlock{ + Type: AnthropicContentBlockTypeThinking, + }, + }, + { + Type: AnthropicStreamEventTypeContentBlockDelta, + Index: schemas.Ptr(0), + Delta: &AnthropicStreamDelta{ + Type: AnthropicStreamDeltaTypeThinking, + Thinking: schemas.Ptr("first thought"), + }, + }, + { + Type: AnthropicStreamEventTypeContentBlockDelta, + Index: schemas.Ptr(0), + Delta: &AnthropicStreamDelta{ + Type: AnthropicStreamDeltaTypeSignature, + Signature: schemas.Ptr("sig-first"), + }, + }, + {Type: AnthropicStreamEventTypeContentBlockStop, Index: schemas.Ptr(0)}, + { + Type: AnthropicStreamEventTypeContentBlockStart, + Index: schemas.Ptr(1), + ContentBlock: &AnthropicContentBlock{ + Type: AnthropicContentBlockTypeToolUse, + ID: schemas.Ptr("toolu_01"), + Name: schemas.Ptr("get_weather"), + }, + }, + { + Type: AnthropicStreamEventTypeContentBlockDelta, + Index: schemas.Ptr(1), + Delta: &AnthropicStreamDelta{ + Type: AnthropicStreamDeltaTypeInputJSON, + PartialJSON: schemas.Ptr(`{"location":"Paris"}`), + }, + }, + {Type: AnthropicStreamEventTypeContentBlockStop, Index: schemas.Ptr(1)}, + { + Type: AnthropicStreamEventTypeContentBlockStart, + Index: schemas.Ptr(2), + ContentBlock: &AnthropicContentBlock{ + Type: AnthropicContentBlockTypeThinking, + }, + }, + { + Type: AnthropicStreamEventTypeContentBlockDelta, + Index: schemas.Ptr(2), + Delta: &AnthropicStreamDelta{ + Type: AnthropicStreamDeltaTypeThinking, + Thinking: schemas.Ptr("second thought"), + }, + }, + { + Type: AnthropicStreamEventTypeContentBlockDelta, + Index: schemas.Ptr(2), + Delta: &AnthropicStreamDelta{ + Type: AnthropicStreamDeltaTypeSignature, + Signature: schemas.Ptr("sig-second"), + }, + }, + {Type: AnthropicStreamEventTypeContentBlockStop, Index: schemas.Ptr(2)}, + { + Type: AnthropicStreamEventTypeMessageDelta, + Delta: &AnthropicStreamDelta{StopReason: &stopReason}, + }, + {Type: AnthropicStreamEventTypeMessageStop}, + } + + responses := driveResponsesStream(t, events) + + var completed *schemas.BifrostResponsesResponse + for _, r := range responses { + if r.Type == schemas.ResponsesStreamResponseTypeCompleted && r.Response != nil { + completed = r.Response + } + } + if completed == nil { + t.Fatal("expected a response.completed event") + } + if len(completed.Output) != 3 { + t.Fatalf("response.completed output has %d items, want 3 (reasoning + function_call + reasoning)", len(completed.Output)) + } + + wantBlocks := []struct { + text string + signature string + }{ + {"first thought", "sig-first"}, + {"", ""}, // output index 1 is the function call + {"second thought", "sig-second"}, + } + for i, want := range wantBlocks { + if want.text == "" { + continue + } + item := completed.Output[i] + if item.Type == nil { + t.Fatalf("output[%d] type is nil, want reasoning", i) + } + if *item.Type != schemas.ResponsesMessageTypeReasoning { + t.Fatalf("output[%d] type = %q, want reasoning", i, *item.Type) + } + b := reasoningTextBlock(&item) + if b == nil || b.Text == nil || *b.Text != want.text { + t.Errorf("output[%d] reasoning text = %+v, want %q", i, b, want.text) + continue + } + if b.Signature == nil || *b.Signature != want.signature { + t.Errorf("output[%d] reasoning signature = %v, want %q", i, b.Signature, want.signature) + } + } +} + +// TestToBifrostResponsesStream_VisibleThinkingWithoutSignature covers thinking +// blocks that never receive a signature delta: the completed reasoning item +// still carries the accumulated text, with no signature attached. +func TestToBifrostResponsesStream_VisibleThinkingWithoutSignature(t *testing.T) { + t.Parallel() + + const fullText = "unsigned reasoning" + responses := driveResponsesStream(t, visibleThinkingToolUseLifecycle([]string{fullText}, "")) + + var itemDone *schemas.ResponsesMessage + for _, r := range responses { + if r.Type == schemas.ResponsesStreamResponseTypeOutputItemDone && r.OutputIndex != nil && *r.OutputIndex == 0 { + itemDone = r.Item + } + } + if itemDone == nil { + t.Fatal("expected an output_item.done event for the reasoning block") + } + if itemDone.Type == nil { + t.Fatal("output_item.done item type is nil, want reasoning") + } + if *itemDone.Type != schemas.ResponsesMessageTypeReasoning { + t.Fatalf("output_item.done item type = %q, want reasoning", *itemDone.Type) + } + block := reasoningTextBlock(itemDone) + if block == nil || block.Text == nil || *block.Text != fullText { + t.Fatalf("output_item.done reasoning text = %+v, want %q", block, fullText) + } + if block.Signature != nil { + t.Errorf("output_item.done reasoning signature = %q, want none", *block.Signature) + } +}