stream

package
v1.33.41 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Feb 2, 2026 License: MIT Imports: 18 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func BuildAgentDiscoveryContext

func BuildAgentDiscoveryContext(agents []config.AgentInfo) string

BuildAgentDiscoveryContext builds context about available peer agents for collaboration

func BuildSystemPrompt

func BuildSystemPrompt(llmConfig config.LLMConfig, agentName, agentDescription string) string

BuildSystemPrompt builds the complete system prompt with dynamic context

func BuildSystemPromptWithConfig

func BuildSystemPromptWithConfig(llmConfig config.LLMConfig, agentName, agentDescription string, mcpConfig *config.MCPConfig, tasksConfig *config.TasksConfig, schedulerConfig *config.SchedulerConfig, system interface{}, availableAgents []config.AgentInfo) string

BuildSystemPromptWithConfig builds the complete system prompt with all configuration context

func BuildSystemPromptWithMCP

func BuildSystemPromptWithMCP(llmConfig config.LLMConfig, agentName, agentDescription string, mcpConfig *config.MCPConfig) string

BuildSystemPromptWithMCP builds the complete system prompt with dynamic context and MCP credentials

func BuildTaskManagementContext

func BuildTaskManagementContext(tasksConfig *config.TasksConfig, schedulerConfig *config.SchedulerConfig) string

BuildTaskManagementContext builds context about task management capabilities

func EstimateTokens

func EstimateTokens(text string) int

EstimateTokens estimates the number of tokens for a given text Uses a simple heuristic: ~4 characters per token (average for English text) This is a rough approximation that works across different LLM providers

func EstimateTokensFromInterface

func EstimateTokensFromInterface(content interface{}) int

EstimateTokensFromInterface estimates tokens from various content types

func EstimateToolDefinitionsTokens

func EstimateToolDefinitionsTokens(customTools []tools.ToolDefinition, builtinTools []interface{}) int

EstimateToolDefinitionsTokens estimates tokens for tool definitions

func GetMaxConcurrent

func GetMaxConcurrent() int

GetMaxConcurrent returns the max concurrent tool executions from config

func IsParallelToolsEnabled

func IsParallelToolsEnabled() bool

IsParallelToolsEnabled checks if parallel tools are enabled in config

func UnifiedStreamResponse

func UnifiedStreamResponse(w http.ResponseWriter, rawReader io.Reader, processor StreamProcessor) error

UnifiedStreamResponse processes a raw stream and outputs unified format

func UnifiedToolConversation

func UnifiedToolConversation(w http.ResponseWriter, provider Provider, processor StreamProcessor, messages []Message, toolDefs []tools.ToolDefinition, messageSaver MessageSaver, threadID string, model *string) error

UnifiedToolConversation - backwards compatible version

func UnifiedToolConversationWithBuiltins

func UnifiedToolConversationWithBuiltins(w http.ResponseWriter, provider Provider, processor StreamProcessor, messages []Message, customTools []tools.ToolDefinition, builtinTools []interface{}, messageSaver MessageSaver, threadID string, model *string, fileProcessor FileProcessor, taskID string) error

UnifiedToolConversationWithBuiltins handles a complete conversation with tool execution and message saving

func UnifiedToolConversationWithContext

func UnifiedToolConversationWithContext(ctx context.Context, w http.ResponseWriter, provider Provider, processor StreamProcessor, messages []Message, customTools []tools.ToolDefinition, builtinTools []interface{}, messageSaver MessageSaver, threadID string, requestID string, model *string, fileProcessor FileProcessor, taskID string) error

UnifiedToolConversationWithContext is like UnifiedToolConversationWithBuiltins but accepts a context for cancellation

Types

type AnthropicContentBlock

type AnthropicContentBlock struct {
	Type      string                 `json:"type"` // "text", "tool_use", "server_tool_use", "web_search_tool_result", "web_fetch_tool_result"
	ID        string                 `json:"id,omitempty"`
	Name      string                 `json:"name,omitempty"`
	Input     map[string]interface{} `json:"input,omitempty"`
	ToolUseID string                 `json:"tool_use_id,omitempty"` // For tool results
	Content   interface{}            `json:"content,omitempty"`     // Can be []WebSearchResult or WebFetchResult
}

type AnthropicDelta

type AnthropicDelta struct {
	Type        string `json:"type"`
	Text        string `json:"text"`
	PartialJSON string `json:"partial_json,omitempty"`
}

type AnthropicMessage

type AnthropicMessage struct {
	ID   string `json:"id"`
	Type string `json:"type"`
	Role string `json:"role"`
}

type AnthropicProcessor

type AnthropicProcessor struct {
	// contains filtered or unexported fields
}

func (*AnthropicProcessor) IsComplete

func (p *AnthropicProcessor) IsComplete(line string) bool

func (*AnthropicProcessor) ProcessLine

func (p *AnthropicProcessor) ProcessLine(line string) (*StreamEvent, error)

type AnthropicStreamData

type AnthropicStreamData struct {
	Type         string                `json:"type"`
	Index        int                   `json:"index,omitempty"`
	Delta        AnthropicDelta        `json:"delta,omitempty"`
	Message      AnthropicMessage      `json:"message,omitempty"`
	ContentBlock AnthropicContentBlock `json:"content_block,omitempty"`
	ToolUseID    string                `json:"tool_use_id,omitempty"` // For web_search_tool_result
	Content      interface{}           `json:"content,omitempty"`     // Can be []WebSearchResult or WebFetchResult
	Usage        *AnthropicUsage       `json:"usage,omitempty"`       // For message_delta token usage
}

type AnthropicUsage

type AnthropicUsage struct {
	InputTokens  int `json:"input_tokens"`
	OutputTokens int `json:"output_tokens"`
}

type BatchExecutor

type BatchExecutor struct {
	// contains filtered or unexported fields
}

BatchExecutor executes multiple tool calls in parallel

func NewBatchExecutor

func NewBatchExecutor(maxConcurrent int) *BatchExecutor

NewBatchExecutor creates a new batch executor

func (*BatchExecutor) ExecuteParallel

func (be *BatchExecutor) ExecuteParallel(toolCalls []ToolCall) []ToolResult

ExecuteParallel executes multiple tool calls in parallel with a concurrency limit

type CitationsInfo

type CitationsInfo struct {
	Enabled bool `json:"enabled"`
}

type CompactionResult

type CompactionResult struct {
	Summary             string    `json:"summary"`
	MessagesCompacted   int       `json:"messages_compacted"`
	OriginalTokens      int       `json:"original_tokens"`
	SummaryTokens       int       `json:"summary_tokens"`
	CompactedAt         time.Time `json:"compacted_at"`
	LastCompactedMsgID  string    `json:"last_compacted_message_id"`
	FirstCompactedMsgID string    `json:"first_compacted_message_id"`
}

CompactionResult holds the result of a compaction operation

type Compactor

type Compactor struct {
	// contains filtered or unexported fields
}

Compactor handles automatic context compaction via summarization

func NewCompactor

func NewCompactor(cfg *config.CompactionConfig, maxTokens int, db *sql.DB, anthropicKey string) *Compactor

NewCompactor creates a new compactor instance

func (*Compactor) Compact

func (c *Compactor) Compact(threadID string, messages []Message) (string, []Message, error)

Compact performs compaction on messages, returning summary and recent messages

func (*Compactor) GetExistingCompaction

func (c *Compactor) GetExistingCompaction(threadID string) (*CompactionResult, error)

GetExistingCompaction retrieves any existing compaction summary for a thread

func (*Compactor) ShouldCompact

func (c *Compactor) ShouldCompact(messages []Message) bool

ShouldCompact checks if compaction is needed based on current token count

type ContentBlock

type ContentBlock struct {
	Type string `json:"type"`           // "text", "image", "document", "tool_use", "tool_result"
	Text string `json:"text,omitempty"` // for text blocks

	// Image/Document fields
	Source *ContentSource `json:"source,omitempty"` // for image and document blocks

	// Tool fields
	ID               string                 `json:"id,omitempty"`                // for tool_use
	Name             string                 `json:"name,omitempty"`              // for tool_use
	Input            map[string]interface{} `json:"-"`                           // for tool_use - custom marshaling, see MarshalJSON
	ThoughtSignature string                 `json:"thought_signature,omitempty"` // for tool_use (Gemini 3 requires this)
	ToolUseID        string                 `json:"tool_use_id,omitempty"`       // for tool_result
	Content          interface{}            `json:"content,omitempty"`           // for tool_result (can be string or object)
	IsError          bool                   `json:"is_error,omitempty"`          // for tool_result
}

func (ContentBlock) MarshalJSON added in v1.33.41

func (c ContentBlock) MarshalJSON() ([]byte, error)

MarshalJSON implements custom JSON marshaling for ContentBlock This ensures the 'input' field is only included for tool_use blocks (Anthropic requires it)

type ContentSource

type ContentSource struct {
	Type      string `json:"type"`                 // "base64", "url", "file"
	MediaType string `json:"media_type,omitempty"` // "image/jpeg", "application/pdf", etc.
	Data      string `json:"data,omitempty"`       // For base64
	URL       string `json:"url,omitempty"`        // For URL
	FileID    string `json:"file_id,omitempty"`    // For file reference
}

type ContextManager

type ContextManager struct {
	// contains filtered or unexported fields
}

ContextManager handles context window management

func NewContextManager

func NewContextManager(maxMessages, maxTokens, keepImages int) *ContextManager

NewContextManager creates a new context manager maxTokens takes precedence over maxMessages if set (> 0)

func (*ContextManager) TruncateMessages

func (cm *ContextManager) TruncateMessages(messages []Message) []Message

TruncateMessages applies context limits to messages 1. If maxTokens is set, limits by estimated token count 2. Otherwise, limits by message count (maxMessages) 3. Strips images from older messages, keeping only last N with images IMPORTANT: Ensures tool_use/tool_result pairs stay together

type DocumentContent

type DocumentContent struct {
	Type      string         `json:"type"` // "document"
	Source    DocumentSource `json:"source"`
	Title     string         `json:"title,omitempty"`
	Citations CitationsInfo  `json:"citations,omitempty"`
}

type DocumentSource

type DocumentSource struct {
	Type      string `json:"type"`       // "text", "base64"
	MediaType string `json:"media_type"` // "text/plain", "application/pdf"
	Data      string `json:"data"`       // The actual content
}

type FileProcessor

type FileProcessor interface {
	ProcessContentBlocks(blocks []ContentBlock, threadID, messageID, source string) ([]ContentBlock, error)
	ExtractImagesFromToolResult(result interface{}, threadID, source string) (interface{}, []string, error)
	IsEnabled() bool
}

FileProcessor interface for filesystem operations

type GeminiCandidate

type GeminiCandidate struct {
	Content      *GeminiContentPart `json:"content,omitempty"`
	FinishReason string             `json:"finishReason,omitempty"` // "STOP", "MAX_TOKENS", etc.
	Index        int                `json:"index"`
}

GeminiCandidate represents a response candidate

type GeminiContentPart

type GeminiContentPart struct {
	Parts []GeminiPartResponse `json:"parts"`
	Role  string               `json:"role"` // "model"
}

GeminiContentPart represents content in the response

type GeminiFunctionCallResponse

type GeminiFunctionCallResponse struct {
	Name string                 `json:"name"`
	Args map[string]interface{} `json:"args"`
}

GeminiFunctionCallResponse represents a function call from Gemini

type GeminiPartResponse

type GeminiPartResponse struct {
	Text             string                      `json:"text,omitempty"`
	FunctionCall     *GeminiFunctionCallResponse `json:"functionCall,omitempty"`
	ThoughtSignature string                      `json:"thoughtSignature,omitempty"` // At Part level for Gemini 3
}

GeminiPartResponse represents a single part of the response

type GeminiStreamProcessor

type GeminiStreamProcessor struct {
	// contains filtered or unexported fields
}

GeminiStreamProcessor processes Gemini streaming responses

func NewGeminiStreamProcessor

func NewGeminiStreamProcessor() *GeminiStreamProcessor

NewGeminiStreamProcessor creates a new Gemini stream processor

func (*GeminiStreamProcessor) HasPendingEvents

func (p *GeminiStreamProcessor) HasPendingEvents() bool

HasPendingEvents returns true if there are queued events waiting to be emitted

func (*GeminiStreamProcessor) IsComplete

func (p *GeminiStreamProcessor) IsComplete(line string) bool

IsComplete checks if the stream is complete

func (*GeminiStreamProcessor) ProcessLine

func (p *GeminiStreamProcessor) ProcessLine(line string) (*StreamEvent, error)

ProcessLine processes a single line from the Gemini stream

func (*GeminiStreamProcessor) Reset

func (p *GeminiStreamProcessor) Reset()

Reset resets the processor state (useful for new conversations)

type GeminiStreamResponse

type GeminiStreamResponse struct {
	Candidates    []GeminiCandidate    `json:"candidates"`
	UsageMetadata *GeminiUsageMetadata `json:"usageMetadata,omitempty"`
}

GeminiStreamResponse represents a single chunk from Gemini's streaming API

type GeminiUsageMetadata

type GeminiUsageMetadata struct {
	PromptTokenCount     int `json:"promptTokenCount"`
	CandidatesTokenCount int `json:"candidatesTokenCount"`
	TotalTokenCount      int `json:"totalTokenCount"`
}

GeminiUsageMetadata contains token usage information

type Message

type Message struct {
	Role             string      `json:"role"`
	Content          interface{} `json:"content"`
	ReasoningContent string      `json:"reasoning_content,omitempty"` // For thinking models (Kimi K2, etc.) - preserved in conversation history
}

func ConvertDatabaseMessageToStreamMessage

func ConvertDatabaseMessageToStreamMessage(role string, content interface{}) (Message, error)

ConvertDatabaseMessageToStreamMessage converts a message loaded from the database into the proper format for the LLM provider

func ConvertDatabaseMessageToStreamMessageWithMetadata

func ConvertDatabaseMessageToStreamMessageWithMetadata(role string, content interface{}, metadata map[string]interface{}) (Message, error)

ConvertDatabaseMessageToStreamMessageWithMetadata converts a message loaded from the database into the proper format for the LLM provider, including reasoning_content from metadata

func ConvertMessagesFromDatabase

func ConvertMessagesFromDatabase(dbMessages []interface{}) ([]Message, error)

ConvertMessagesFromDatabase converts a slice of database messages to stream messages

func PrepareMessagesWithFullConfig

func PrepareMessagesWithFullConfig(messages []Message, llmConfig config.LLMConfig, agentName, agentDescription string, mcpConfig *config.MCPConfig, tasksConfig *config.TasksConfig, schedulerConfig *config.SchedulerConfig, system interface{}, availableAgents []config.AgentInfo) []Message

PrepareMessagesWithFullConfig prepends a system message with all configuration context

func PrepareMessagesWithSystemPrompt

func PrepareMessagesWithSystemPrompt(messages []Message, llmConfig config.LLMConfig, agentName, agentDescription string) []Message

PrepareMessagesWithSystemPrompt prepends a system message if not already present

func PrepareMessagesWithSystemPromptAndMCP

func PrepareMessagesWithSystemPromptAndMCP(messages []Message, llmConfig config.LLMConfig, agentName, agentDescription string, mcpConfig *config.MCPConfig) []Message

PrepareMessagesWithSystemPromptAndMCP prepends a system message with MCP credentials if not already present

type MessageSaver

type MessageSaver interface {
	SaveMessage(threadID, role string, content interface{}, model *string, metadata map[string]interface{}) error
	CreateThread(title string) (string, error)
}

type OpenAIChoice

type OpenAIChoice struct {
	Index        int         `json:"index"`
	Delta        OpenAIDelta `json:"delta"`
	FinishReason *string     `json:"finish_reason"`
}

type OpenAIDelta

type OpenAIDelta struct {
	Role             string           `json:"role,omitempty"`
	Content          string           `json:"content,omitempty"`
	ReasoningContent string           `json:"reasoning_content,omitempty"` // For thinking/reasoning models (Fireworks Kimi K2, etc.)
	Reasoning        string           `json:"reasoning,omitempty"`         // For thinking/reasoning models (Together AI Kimi K2, etc.)
	ToolCalls        []OpenAIToolCall `json:"tool_calls,omitempty"`
}

type OpenAIFunctionCall

type OpenAIFunctionCall struct {
	Name      string `json:"name"`
	Arguments string `json:"arguments"`
}

type OpenAIProcessor

type OpenAIProcessor struct {
	// contains filtered or unexported fields
}

func (*OpenAIProcessor) ClearAccumulatedReasoning

func (p *OpenAIProcessor) ClearAccumulatedReasoning()

ClearAccumulatedReasoning clears the accumulated reasoning content Called after the reasoning has been used for the next API call

func (*OpenAIProcessor) GetAccumulatedReasoning

func (p *OpenAIProcessor) GetAccumulatedReasoning() string

GetAccumulatedReasoning returns the accumulated reasoning content without clearing it This is used to preserve reasoning_content in conversation history for thinking models (Kimi K2, etc.)

func (*OpenAIProcessor) HasPendingEvents

func (p *OpenAIProcessor) HasPendingEvents() bool

HasPendingEvents returns true if there are queued events waiting to be emitted

func (*OpenAIProcessor) IsComplete

func (p *OpenAIProcessor) IsComplete(line string) bool

func (*OpenAIProcessor) ProcessLine

func (p *OpenAIProcessor) ProcessLine(line string) (*StreamEvent, error)

func (*OpenAIProcessor) ResetForNewTurn

func (p *OpenAIProcessor) ResetForNewTurn()

ResetForNewTurn resets the processor state for a new conversation turn This is needed when the same processor is reused across multiple iterations (e.g., after tool calls)

func (*OpenAIProcessor) SetStartInThinkMode

func (p *OpenAIProcessor) SetStartInThinkMode(start bool)

SetStartInThinkMode sets the processor to start in thinking mode This is needed for models like Together AI K2.5 that start thinking without <think> tag

type OpenAIStreamData

type OpenAIStreamData struct {
	ID      string         `json:"id"`
	Object  string         `json:"object"`
	Created int64          `json:"created"`
	Model   string         `json:"model"`
	Choices []OpenAIChoice `json:"choices"`
	Usage   *OpenAIUsage   `json:"usage,omitempty"` // Token usage in final chunk
}

type OpenAIToolCall

type OpenAIToolCall struct {
	Index    int                `json:"index"`
	ID       string             `json:"id"`
	Type     string             `json:"type"`
	Function OpenAIFunctionCall `json:"function"`
}

type OpenAIUsage

type OpenAIUsage struct {
	PromptTokens     int `json:"prompt_tokens"`
	CompletionTokens int `json:"completion_tokens"`
	TotalTokens      int `json:"total_tokens"`
}

type PendingEventChecker

type PendingEventChecker interface {
	HasPendingEvents() bool
}

PendingEventChecker is an optional interface for processors that queue events

type Provider

type Provider interface {
	GetRawStream(messages []Message, customTools []tools.ToolDefinition, builtinTools []interface{}) (io.ReadCloser, error)
}

type ReasoningAccumulator

type ReasoningAccumulator interface {
	GetAccumulatedReasoning() string
	ClearAccumulatedReasoning()
}

ReasoningAccumulator is an optional interface for processors that accumulate reasoning content Used by thinking models (Kimi K2, etc.) to preserve reasoning_content in conversation history

type StreamEvent

type StreamEvent struct {
	Type             string                 `json:"type"`                        // "content", "thinking", "start", "stop", "tool_use", "tool_result", "usage"
	Content          string                 `json:"content"`                     // actual text content (for simple string results)
	ContentBlocks    interface{}            `json:"content_blocks,omitempty"`    // for structured content (e.g., screenshot with image)
	ToolName         string                 `json:"tool_name,omitempty"`         // for tool_use events
	ToolDisplayName  string                 `json:"tool_display_name,omitempty"` // human-readable name like "Create Task"
	ToolID           string                 `json:"tool_id,omitempty"`           // for tool_use and tool_result events
	ToolInput        map[string]interface{} `json:"tool_input"`                  // for tool_use events
	ToolResult       interface{}            `json:"tool_result,omitempty"`       // for tool_result events (raw data for server tools)
	ThoughtSignature string                 `json:"thought_signature,omitempty"` // for Gemini 3 function calls (required for context)
	Timestamp        int64                  `json:"timestamp"`                   // Unix milliseconds timestamp

	// Token usage fields (for "usage" events)
	InputTokens  int `json:"input_tokens,omitempty"`
	OutputTokens int `json:"output_tokens,omitempty"`
	TotalTokens  int `json:"total_tokens,omitempty"`
}

type StreamProcessor

type StreamProcessor interface {
	ProcessLine(line string) (*StreamEvent, error)
	IsComplete(line string) bool
}

type TokenAttribution

type TokenAttribution struct {
	SystemPrompt        int `json:"system_prompt"`
	ToolDefinitions     int `json:"tool_definitions"`
	ConversationHistory int `json:"conversation_history"`
	MemoryContext       int `json:"memory_context"`
	UserMessage         int `json:"user_message"`
	EstimatedTotal      int `json:"estimated_total"`
	ActualTotal         int `json:"actual_total,omitempty"` // Filled in after LLM response
}

TokenAttribution tracks estimated token counts for each input component

func CalculateTokenAttribution

func CalculateTokenAttribution(messages []Message, customTools []tools.ToolDefinition, builtinTools []interface{}, memoryContextLen int) TokenAttribution

CalculateTokenAttribution calculates token attribution for all input components Parameters:

  • messages: the full message array (including system message)
  • customTools: custom tool definitions
  • builtinTools: builtin tool definitions
  • memoryContextLen: character length of memory context (if tracked separately)

type ToolCall

type ToolCall struct {
	ID    string                 `json:"id"`
	Name  string                 `json:"name"`
	Input map[string]interface{} `json:"input"`
}

ToolCall represents a single tool invocation

type ToolResult

type ToolResult struct {
	ID            string                 `json:"id"`
	Name          string                 `json:"name"`
	Content       string                 `json:"content"`
	ContentBlocks interface{}            `json:"content_blocks,omitempty"`
	Input         map[string]interface{} `json:"input"`
	Error         error                  `json:"-"`
	Duration      time.Duration          `json:"duration"`
}

ToolResult represents the result of a tool execution

type TurnResetter

type TurnResetter interface {
	ResetForNewTurn()
}

TurnResetter is an optional interface for processors that need to reset state between turns Used by thinking models that embed thinking in content (e.g., Together AI K2.5)

type WebFetchResult

type WebFetchResult struct {
	Type        string          `json:"type"` // "web_fetch_result"
	URL         string          `json:"url"`
	Content     DocumentContent `json:"content"`
	RetrievedAt string          `json:"retrieved_at"`
}

type WebSearchResult

type WebSearchResult struct {
	Type             string `json:"type"` // "web_search_result"
	URL              string `json:"url"`
	Title            string `json:"title"`
	EncryptedContent string `json:"encrypted_content"`
	PageAge          string `json:"page_age"`
}

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL