Documentation
¶
Index ¶
- Constants
- func ClassifyError(err error) string
- func DecodeToolSpec(t *models.Tool) error
- func ExecuteResponseFilters(ctx context.Context, filters []*models.Filter, ...) (blocked bool, blockMessage string, err error)
- func NewRunID() string
- func PublishWithTimeout(ctx context.Context, queue MessageQueue, publisher func(context.Context) error, ...) error
- func UnwrapClientAnswer(content string) (interface{}, bool)
- type CallParams
- type ChatEvent
- type ChatMode
- type ChatResponse
- type ChatSession
- func (cs *ChatSession) ActiveRunID() string
- func (cs *ChatSession) AddDatasource(id uint) error
- func (cs *ChatSession) AddFileReference(filename, contents string)
- func (cs *ChatSession) AddPreProcessor(fn func(*models.UserMessage) error)
- func (cs *ChatSession) AddTool(id string, t models.Tool) error
- func (cs *ChatSession) AwaitingClientTools() bool
- func (cs *ChatSession) CancelRun() bool
- func (cs *ChatSession) ClientTools() []ClientToolInfo
- func (cs *ChatSession) CurrentTools() map[string]models.Tool
- func (cs *ChatSession) Errors() <-chan error
- func (cs *ChatSession) GetCurrentDatasources() []*models.Datasource
- func (cs *ChatSession) GetFileReference(filename string) (string, bool)
- func (cs *ChatSession) HandleLLMResponse(w *LLMResponseWrapper) (continued bool, err error)
- func (cs *ChatSession) HandleUserMessage(msg *models.UserMessage, docs []schema.Document, tools []llms.Tool, ...) (*llms.ContentResponse, error)
- func (cs *ChatSession) ID() string
- func (cs *ChatSession) Input() chan *models.UserMessage
- func (cs *ChatSession) NotifyStatus(status string)
- func (cs *ChatSession) OutputMessage() <-chan *ChatResponse
- func (cs *ChatSession) OutputMode() OutputMode
- func (cs *ChatSession) OutputStream() <-chan []byte
- func (cs *ChatSession) PendingClientCallIDs() []string
- func (cs *ChatSession) PreflightTokenLengthCheck(msgs []llms.MessageContent) []llms.MessageContent
- func (cs *ChatSession) RemoveDatasource(id uint)
- func (cs *ChatSession) RemoveTool(id string)
- func (cs *ChatSession) SetOutputMode(m OutputMode)
- func (cs *ChatSession) Start() error
- func (cs *ChatSession) Stop()
- func (cs *ChatSession) Stopped() bool
- func (cs *ChatSession) Subscribe(runID string, buf int) (<-chan ChatEvent, func())
- func (cs *ChatSession) SwitchOutputMode(m OutputMode) bool
- func (cs *ChatSession) TryLockRun() bool
- func (cs *ChatSession) UnlockRun()
- type ClientToolInfo
- type ContentChoiceForNATS
- type ContentResponseForNATS
- type ContextData
- type DefaultQueueFactory
- type DeferredPostgreSQLQueueFactory
- type ErrorData
- type FinishData
- type FunctionCallForNATS
- type GormChatMessageHistory
- func (h *GormChatMessageHistory) AddAIMessage(ctx context.Context, text string) error
- func (h *GormChatMessageHistory) AddAIToolCall(ctx context.Context, toolCall llms.ToolCall) error
- func (c *GormChatMessageHistory) AddMessage(ctx context.Context, mc llms.MessageContent) error
- func (h *GormChatMessageHistory) AddSystemMessage(ctx context.Context, text string) error
- func (h *GormChatMessageHistory) AddToolMessage(ctx context.Context, toolResp llms.ToolCallResponse) error
- func (h *GormChatMessageHistory) AddUserMessage(ctx context.Context, text string) error
- func (h *GormChatMessageHistory) CheckIfSessionExists(ctx context.Context) (bool, *models.ChatHistoryRecord, error)
- func (h *GormChatMessageHistory) Clear(ctx context.Context) error
- func (h *GormChatMessageHistory) GetAssociatedChat(ctx context.Context) (*models.Chat, error)
- func (h *GormChatMessageHistory) GetMemoryKey(context.Context) string
- func (h *GormChatMessageHistory) LastMessageID() uint
- func (h *GormChatMessageHistory) Messages(ctx context.Context) ([]llms.MessageContent, error)
- func (h *GormChatMessageHistory) SetMessages(ctx context.Context, messages []llms.MessageContent) error
- type GormChatMessageHistoryOption
- type InMemoryQueue
- func (q *InMemoryQueue) Close() error
- func (q *InMemoryQueue) ConsumeErrors(ctx context.Context) <-chan error
- func (q *InMemoryQueue) ConsumeLLMResponses(ctx context.Context) <-chan *LLMResponseWrapper
- func (q *InMemoryQueue) ConsumeMessages(ctx context.Context) <-chan *ChatResponse
- func (q *InMemoryQueue) ConsumeStream(ctx context.Context) <-chan []byte
- func (q *InMemoryQueue) PublishError(ctx context.Context, err error) error
- func (q *InMemoryQueue) PublishLLMResponse(ctx context.Context, resp *LLMResponseWrapper) error
- func (q *InMemoryQueue) PublishMessage(ctx context.Context, msg *ChatResponse) error
- func (q *InMemoryQueue) PublishStream(ctx context.Context, data []byte) error
- func (q *InMemoryQueue) QueueDepth() (messages, stream, errors, llmResponses int)
- type LLMDriver
- type LLMResponseWrapper
- type LLMResponseWrapperForNATS
- type MessageQueue
- type NATSConfig
- type NATSMessage
- type NATSQueue
- func (nq *NATSQueue) Close() error
- func (nq *NATSQueue) ConsumeErrors(ctx context.Context) <-chan error
- func (nq *NATSQueue) ConsumeLLMResponses(ctx context.Context) <-chan *LLMResponseWrapper
- func (nq *NATSQueue) ConsumeMessages(ctx context.Context) <-chan *ChatResponse
- func (nq *NATSQueue) ConsumeStream(ctx context.Context) <-chan []byte
- func (nq *NATSQueue) PublishError(ctx context.Context, err error) error
- func (nq *NATSQueue) PublishLLMResponse(ctx context.Context, resp *LLMResponseWrapper) error
- func (nq *NATSQueue) PublishMessage(ctx context.Context, msg *ChatResponse) error
- func (nq *NATSQueue) PublishStream(ctx context.Context, data []byte) error
- func (nq *NATSQueue) QueueDepth() (messages, stream, errors, llmResponses int)
- type NATSQueueFactory
- type OptimizedPostgreSQLQueue
- func (psq *OptimizedPostgreSQLQueue) Close() error
- func (psq *OptimizedPostgreSQLQueue) ConsumeErrors(ctx context.Context) <-chan error
- func (psq *OptimizedPostgreSQLQueue) ConsumeLLMResponses(ctx context.Context) <-chan *LLMResponseWrapper
- func (psq *OptimizedPostgreSQLQueue) ConsumeMessages(ctx context.Context) <-chan *ChatResponse
- func (psq *OptimizedPostgreSQLQueue) ConsumeStream(ctx context.Context) <-chan []byte
- func (psq *OptimizedPostgreSQLQueue) PublishError(ctx context.Context, err error) error
- func (psq *OptimizedPostgreSQLQueue) PublishLLMResponse(ctx context.Context, resp *LLMResponseWrapper) error
- func (psq *OptimizedPostgreSQLQueue) PublishMessage(ctx context.Context, msg *ChatResponse) error
- func (psq *OptimizedPostgreSQLQueue) PublishStream(ctx context.Context, data []byte) error
- func (psq *OptimizedPostgreSQLQueue) QueueDepth() (messages, stream, errors, llmResponses int)
- type OutputMode
- type PostgreSQLConfig
- type PostgreSQLMessage
- type PostgreSQLQueue
- func (psq *PostgreSQLQueue) Close() error
- func (psq *PostgreSQLQueue) ConsumeErrors(ctx context.Context) <-chan error
- func (psq *PostgreSQLQueue) ConsumeLLMResponses(ctx context.Context) <-chan *LLMResponseWrapper
- func (psq *PostgreSQLQueue) ConsumeMessages(ctx context.Context) <-chan *ChatResponse
- func (psq *PostgreSQLQueue) ConsumeStream(ctx context.Context) <-chan []byte
- func (psq *PostgreSQLQueue) PublishError(ctx context.Context, err error) error
- func (psq *PostgreSQLQueue) PublishLLMResponse(ctx context.Context, resp *LLMResponseWrapper) error
- func (psq *PostgreSQLQueue) PublishMessage(ctx context.Context, msg *ChatResponse) error
- func (psq *PostgreSQLQueue) PublishStream(ctx context.Context, data []byte) error
- func (psq *PostgreSQLQueue) QueueDepth() (messages, stream, errors, llmResponses int)
- type PostgreSQLQueueFactory
- type QueueFactory
- func CreateDefaultQueueFactory() QueueFactory
- func CreateDefaultQueueFactoryWithSharedDB(db *gorm.DB) QueueFactory
- func CreateQueueFactory(cfg config.QueueConfig) (QueueFactory, error)
- func CreateQueueFactoryFromConfig() (QueueFactory, error)
- func CreateQueueFactoryWithSharedDB(cfg config.QueueConfig, db *gorm.DB) (QueueFactory, error)
- type SharedPostgreSQLQueueFactory
- type StartData
- type StatusData
- type TextDeltaData
- type ToolCallDeltaData
- type ToolCallEndData
- type ToolCallForNATS
- type ToolCallStartData
- type ToolResultData
Constants ¶
const ( EventStart = "start" EventTextDelta = "text-delta" EventTextEnd = "text-end" EventReasoningDelta = "reasoning-delta" EventToolCallStart = "tool-call-start" EventToolCallDelta = "tool-call-delta" EventToolCallEnd = "tool-call-end" EventToolResult = "tool-result" EventStatus = "data-status" EventContext = "data-context" EventError = "data-error" EventFinish = "finish" )
Event kinds. Where the AI SDK "UI message stream" has a native chunk type the same name is used so the browser mapping is mechanical; Studio-specific payloads use the "data-" prefix, which that protocol reserves for custom parts.
const ( FinishStop = "stop" FinishToolCalls = "tool-calls" // the turn is waiting on client-side tool results FinishError = "error" FinishCancelled = "cancelled" FinishFiltered = "filtered" )
Finish reasons.
const ( ErrCodeLLMConfig = "llm_config" ErrCodeAPI = "api" ErrCodeConnection = "connection" ErrCodeAuth = "auth" ErrCodeFilter = "filter" ErrCodeSession = "session" ErrCodeTool = "tool" ErrCodeInternal = "internal" )
Error codes carried in EventError payloads. They replace the string sniffing the v1 frontend did on error text.
const ( ContextSourceRAG = "rag" ContextSourceToolDocs = "tool_docs" ContextSourceFile = "file" )
Context sources carried in EventContext payloads.
const ( MessageTypeChatResponse = "chat_response" MessageTypeStream = "stream" MessageTypeError = "error" MessageTypeLLMResponse = "llm_response" )
Message type constants
const ( PostgreSQLMessageTypeChatResponse = "chat_response" PostgreSQLMessageTypeStream = "stream" PostgreSQLMessageTypeError = "error" PostgreSQLMessageTypeLLMResponse = "llm_response" )
Message type constants
Variables ¶
This section is empty.
Functions ¶
func ClassifyError ¶
ClassifyError maps an error to one of the ErrCode* constants. The heuristics mirror what the v1 browser code did (processErrorMessage / detectErrorType) so the two UIs categorise the same failures the same way.
func DecodeToolSpec ¶
DecodeToolSpec replaces the tool's stored spec (base64, the API contract) with its plain text before the tool joins a session. Every path that attaches a tool goes through here: default tools, dependencies and tools picked in the chat window. A client tool parses its own definition and accepts either form (see Tool.ClientDefinition), so one stored as raw JSON is left as it is instead of failing the add. On error the spec is untouched.
func ExecuteResponseFilters ¶
func ExecuteResponseFilters( ctx context.Context, filters []*models.Filter, service services.ServiceInterface, responseText string, vendor string, modelName string, isStreaming bool, isChunk bool, chunkIndex int, currentBuffer string, sessionID string, userID uint, chatID uint, ) (blocked bool, blockMessage string, err error)
ExecuteResponseFilters executes response-side filters on chat LLM responses Returns whether the response should be blocked and an optional block message
func PublishWithTimeout ¶
func PublishWithTimeout(ctx context.Context, queue MessageQueue, publisher func(context.Context) error, timeout time.Duration) error
Helper function that can be used in sendStatus with proper timeout
func UnwrapClientAnswer ¶
UnwrapClientAnswer is unwrapClientAnswer for other packages (history materialisation).
Types ¶
type CallParams ¶
type ChatEvent ¶
type ChatEvent struct {
Version string `json:"v"`
RunID string `json:"run_id"`
Seq uint64 `json:"seq"`
Kind string `json:"kind"`
Data json.RawMessage `json:"data,omitempty"`
}
ChatEvent is the envelope published on the stream channel in OutputModeEvents.
func DecodeChatEvent ¶
DecodeChatEvent reports whether b is an event envelope and decodes it.
type ChatResponse ¶
type ChatResponse struct {
Payload string
}
type ChatSession ¶
type ChatSession struct {
// contains filtered or unexported fields
}
func NewChatSession ¶
func (*ChatSession) ActiveRunID ¶
func (cs *ChatSession) ActiveRunID() string
ActiveRunID returns the id of the turn being processed, or "".
func (*ChatSession) AddDatasource ¶
func (cs *ChatSession) AddDatasource(id uint) error
func (*ChatSession) AddFileReference ¶
func (cs *ChatSession) AddFileReference(filename, contents string)
func (*ChatSession) AddPreProcessor ¶
func (cs *ChatSession) AddPreProcessor(fn func(*models.UserMessage) error)
func (*ChatSession) AwaitingClientTools ¶
func (cs *ChatSession) AwaitingClientTools() bool
AwaitingClientTools reports whether the last turn is parked on client tools.
func (*ChatSession) CancelRun ¶
func (cs *ChatSession) CancelRun() bool
CancelRun aborts the in-flight LLM call of the active turn, if any. The turn still finishes normally (with FinishCancelled) once the call returns.
func (*ChatSession) ClientTools ¶
func (cs *ChatSession) ClientTools() []ClientToolInfo
ClientTools lists the client tools attached to the session.
func (*ChatSession) CurrentTools ¶
func (cs *ChatSession) CurrentTools() map[string]models.Tool
CurrentTools returns a copy of the tools attached to the session.
func (*ChatSession) Errors ¶
func (cs *ChatSession) Errors() <-chan error
func (*ChatSession) GetCurrentDatasources ¶
func (cs *ChatSession) GetCurrentDatasources() []*models.Datasource
GetCurrentDatasources returns a slice of current datasources
func (*ChatSession) GetFileReference ¶
func (cs *ChatSession) GetFileReference(filename string) (string, bool)
func (*ChatSession) HandleLLMResponse ¶
func (cs *ChatSession) HandleLLMResponse(w *LLMResponseWrapper) (continued bool, err error)
HandleLLMResponse persists and forwards one LLM reply. When the reply asked for tools it executes them and dispatches a follow-up LLM call, in which case it returns continued=true: the turn is not over and the next reply will arrive on the queue. continued=false means the turn ended with this reply.
func (*ChatSession) HandleUserMessage ¶
func (cs *ChatSession) HandleUserMessage(msg *models.UserMessage, docs []schema.Document, tools []llms.Tool, files map[string]string) (*llms.ContentResponse, error)
func (*ChatSession) ID ¶
func (cs *ChatSession) ID() string
func (*ChatSession) Input ¶
func (cs *ChatSession) Input() chan *models.UserMessage
func (*ChatSession) NotifyStatus ¶
func (cs *ChatSession) NotifyStatus(status string)
NotifyStatus sends a status message through the chat session
func (*ChatSession) OutputMessage ¶
func (cs *ChatSession) OutputMessage() <-chan *ChatResponse
func (*ChatSession) OutputMode ¶
func (cs *ChatSession) OutputMode() OutputMode
OutputMode returns the session's output mode.
func (*ChatSession) OutputStream ¶
func (cs *ChatSession) OutputStream() <-chan []byte
func (*ChatSession) PendingClientCallIDs ¶
func (cs *ChatSession) PendingClientCallIDs() []string
PendingClientCallIDs returns the ids of the parked calls.
func (*ChatSession) PreflightTokenLengthCheck ¶
func (cs *ChatSession) PreflightTokenLengthCheck(msgs []llms.MessageContent) []llms.MessageContent
func (*ChatSession) RemoveDatasource ¶
func (cs *ChatSession) RemoveDatasource(id uint)
func (*ChatSession) RemoveTool ¶
func (cs *ChatSession) RemoveTool(id string)
func (*ChatSession) SetOutputMode ¶
func (cs *ChatSession) SetOutputMode(m OutputMode)
SetOutputMode fixes the output mode. It must be called before Start().
func (*ChatSession) Start ¶
func (cs *ChatSession) Start() error
func (*ChatSession) Stop ¶
func (cs *ChatSession) Stop()
Stop ends the session: it cancels the session context (aborting any in-flight LLM call), stops the processing goroutine and closes the queue. It is idempotent and never blocks on the goroutine, so a reaper can call it on a busy session. The input channel is deliberately not closed: handlers that still hold a reference would panic on send; they get a queue-closed error from the session instead.
func (*ChatSession) Stopped ¶
func (cs *ChatSession) Stopped() bool
Stopped reports whether Stop has been called.
func (*ChatSession) Subscribe ¶
func (cs *ChatSession) Subscribe(runID string, buf int) (<-chan ChatEvent, func())
Subscribe returns a channel that receives the envelopes of runID until unsubscribe is called. Only meaningful in OutputModeEvents. Events are dropped (with a warning) if the subscriber falls more than buf behind.
func (*ChatSession) SwitchOutputMode ¶
func (cs *ChatSession) SwitchOutputMode(m OutputMode) bool
SwitchOutputMode moves an idle session to another API version and reports whether it did. A session reloaded from the database through a v1 endpoint (a tool toggle, an upload) starts in raw mode; when the v2 client then runs a turn it takes the session over here instead of being refused.
Only raw -> events is supported: the events fan-out owns the queue's stream channel once started, so a v1 SSE reader could not attach afterwards. A busy session (run in flight) is never switched.
func (*ChatSession) TryLockRun ¶
func (cs *ChatSession) TryLockRun() bool
TryLockRun claims the session for one turn. It returns false when another turn is in flight; the caller must UnlockRun when its turn has finished.
func (*ChatSession) UnlockRun ¶
func (cs *ChatSession) UnlockRun()
UnlockRun releases the claim taken by TryLockRun.
type ClientToolInfo ¶
type ClientToolInfo struct {
Name string `json:"name"`
Description string `json:"description"`
Schema map[string]interface{} `json:"schema"`
UI models.ClientToolUI `json:"ui"`
}
ClientToolInfo describes a client tool to the browser.
type ContentChoiceForNATS ¶
type ContentChoiceForNATS struct {
Content string `json:"content"`
StopReason string `json:"stop_reason,omitempty"`
ToolCalls []ToolCallForNATS `json:"tool_calls,omitempty"`
GenerationInfo map[string]any `json:"generation_info,omitempty"`
}
type ContentResponseForNATS ¶
type ContentResponseForNATS struct {
Choices []*ContentChoiceForNATS `json:"choices"`
}
type ContextData ¶
type DefaultQueueFactory ¶
type DefaultQueueFactory struct {
// contains filtered or unexported fields
}
DefaultQueueFactory creates InMemoryQueue instances
func NewDefaultQueueFactory ¶
func NewDefaultQueueFactory(defaultBufferSize int) *DefaultQueueFactory
NewDefaultQueueFactory creates a new factory with specified default buffer size
func (*DefaultQueueFactory) CreateQueue ¶
func (f *DefaultQueueFactory) CreateQueue(sessionID string, config map[string]interface{}) (MessageQueue, error)
CreateQueue creates a new InMemoryQueue with configuration overrides
type DeferredPostgreSQLQueueFactory ¶
type DeferredPostgreSQLQueueFactory struct {
// contains filtered or unexported fields
}
DeferredPostgreSQLQueueFactory creates PostgreSQL queues by connecting to the database at queue creation time
func NewDeferredPostgreSQLQueueFactory ¶
func NewDeferredPostgreSQLQueueFactory(config PostgreSQLConfig) *DeferredPostgreSQLQueueFactory
NewDeferredPostgreSQLQueueFactory creates a deferred PostgreSQL factory
func (*DeferredPostgreSQLQueueFactory) CreateQueue ¶
func (f *DeferredPostgreSQLQueueFactory) CreateQueue(sessionID string, config map[string]interface{}) (MessageQueue, error)
CreateQueue creates a new PostgreSQL queue by connecting to the database using DATABASE_URL
type FinishData ¶
type FinishData struct {
Reason string `json:"reason"`
// Database ids of the rows this turn produced, so a client can address
// them later (edit truncation, regenerate) without reloading history.
UserMessageID uint `json:"user_message_id,omitempty"`
AssistantMessageID uint `json:"assistant_message_id,omitempty"`
}
type FunctionCallForNATS ¶
NATS-safe structures for tool calls
type GormChatMessageHistory ¶
type GormChatMessageHistory struct {
DB *gorm.DB
Limit int
Session string
ChatID uint
UserID uint
// contains filtered or unexported fields
}
GormChatMessageHistory is a struct that stores chat messages using GORM
func NewGormChatMessageHistory ¶
func NewGormChatMessageHistory(db *gorm.DB, session string, chatReference uint, userID uint, systemPrompt string, options ...GormChatMessageHistoryOption) *GormChatMessageHistory
NewGormChatMessageHistory creates a new GormChatMessageHistory
func (*GormChatMessageHistory) AddAIMessage ¶
func (h *GormChatMessageHistory) AddAIMessage(ctx context.Context, text string) error
AddAIMessage adds an AIMessage to the chat message history
func (*GormChatMessageHistory) AddAIToolCall ¶
func (*GormChatMessageHistory) AddMessage ¶
func (c *GormChatMessageHistory) AddMessage(ctx context.Context, mc llms.MessageContent) error
func (*GormChatMessageHistory) AddSystemMessage ¶
func (h *GormChatMessageHistory) AddSystemMessage(ctx context.Context, text string) error
AddSystemMessage adds a system message to the chat message history
func (*GormChatMessageHistory) AddToolMessage ¶
func (h *GormChatMessageHistory) AddToolMessage(ctx context.Context, toolResp llms.ToolCallResponse) error
AddUserMessage adds a user message to the chat message history
func (*GormChatMessageHistory) AddUserMessage ¶
func (h *GormChatMessageHistory) AddUserMessage(ctx context.Context, text string) error
AddUserMessage adds a user message to the chat message history
func (*GormChatMessageHistory) CheckIfSessionExists ¶
func (h *GormChatMessageHistory) CheckIfSessionExists(ctx context.Context) (bool, *models.ChatHistoryRecord, error)
func (*GormChatMessageHistory) Clear ¶
func (h *GormChatMessageHistory) Clear(ctx context.Context) error
Clear resets messages
func (*GormChatMessageHistory) GetAssociatedChat ¶
func (*GormChatMessageHistory) GetMemoryKey ¶
func (h *GormChatMessageHistory) GetMemoryKey(context.Context) string
func (*GormChatMessageHistory) LastMessageID ¶
func (h *GormChatMessageHistory) LastMessageID() uint
LastMessageID returns the row id of the most recent message this history object inserted (0 before the first insert).
func (*GormChatMessageHistory) Messages ¶
func (h *GormChatMessageHistory) Messages(ctx context.Context) ([]llms.MessageContent, error)
Messages returns all messages stored
func (*GormChatMessageHistory) SetMessages ¶
func (h *GormChatMessageHistory) SetMessages(ctx context.Context, messages []llms.MessageContent) error
SetMessages resets chat history and bulk inserts new messages into it
type GormChatMessageHistoryOption ¶
type GormChatMessageHistoryOption func(*GormChatMessageHistory)
GormChatMessageHistoryOption is a function type for configuring GormChatMessageHistory
func WithLimit ¶
func WithLimit(limit int) GormChatMessageHistoryOption
WithLimit sets the limit for the number of messages to retrieve
type InMemoryQueue ¶
type InMemoryQueue struct {
// contains filtered or unexported fields
}
InMemoryQueue implements MessageQueue using Go channels Maintains exact compatibility with current ChatSession channel behavior
func NewInMemoryQueue ¶
func NewInMemoryQueue(sessionID string, bufferSize int) *InMemoryQueue
NewInMemoryQueue creates a new in-memory queue with specified buffer sizes
func (*InMemoryQueue) Close ¶
func (q *InMemoryQueue) Close() error
Close closes all channels and marks the queue as closed
func (*InMemoryQueue) ConsumeErrors ¶
func (q *InMemoryQueue) ConsumeErrors(ctx context.Context) <-chan error
ConsumeErrors returns the errors channel
func (*InMemoryQueue) ConsumeLLMResponses ¶
func (q *InMemoryQueue) ConsumeLLMResponses(ctx context.Context) <-chan *LLMResponseWrapper
ConsumeLLMResponses returns the LLM responses channel
func (*InMemoryQueue) ConsumeMessages ¶
func (q *InMemoryQueue) ConsumeMessages(ctx context.Context) <-chan *ChatResponse
ConsumeMessages returns the output messages channel
func (*InMemoryQueue) ConsumeStream ¶
func (q *InMemoryQueue) ConsumeStream(ctx context.Context) <-chan []byte
ConsumeStream returns the output stream channel
func (*InMemoryQueue) PublishError ¶
func (q *InMemoryQueue) PublishError(ctx context.Context, err error) error
PublishError sends an error, blocking until successful or context cancelled
func (*InMemoryQueue) PublishLLMResponse ¶
func (q *InMemoryQueue) PublishLLMResponse(ctx context.Context, resp *LLMResponseWrapper) error
PublishLLMResponse sends an LLM response, blocking until successful or context cancelled
func (*InMemoryQueue) PublishMessage ¶
func (q *InMemoryQueue) PublishMessage(ctx context.Context, msg *ChatResponse) error
PublishMessage sends a message, blocking until successful or context cancelled
func (*InMemoryQueue) PublishStream ¶
func (q *InMemoryQueue) PublishStream(ctx context.Context, data []byte) error
PublishStream sends stream data, blocking until successful or context cancelled
func (*InMemoryQueue) QueueDepth ¶
func (q *InMemoryQueue) QueueDepth() (messages, stream, errors, llmResponses int)
QueueDepth returns the current depth of all queues
type LLMResponseWrapper ¶
type LLMResponseWrapper struct {
Response *llms.ContentResponse
Opts []llms.CallOption
}
type LLMResponseWrapperForNATS ¶
type LLMResponseWrapperForNATS struct {
Response *ContentResponseForNATS `json:"response"`
}
LLMResponseWrapperForNATS is a NATS-serializable version of LLMResponseWrapper that excludes the non-serializable Opts field and uses NATS-safe tool call structures
type MessageQueue ¶
type MessageQueue interface {
// Publishing methods - all block until successful or context cancelled
// Returns error only on context cancellation or queue closure
PublishMessage(ctx context.Context, msg *ChatResponse) error
PublishStream(ctx context.Context, data []byte) error
PublishError(ctx context.Context, err error) error
PublishLLMResponse(ctx context.Context, resp *LLMResponseWrapper) error
// Consuming methods - returns channels for backward compatibility
// Channels remain open until Close() is called
ConsumeMessages(ctx context.Context) <-chan *ChatResponse
ConsumeStream(ctx context.Context) <-chan []byte
ConsumeErrors(ctx context.Context) <-chan error
ConsumeLLMResponses(ctx context.Context) <-chan *LLMResponseWrapper
// Lifecycle management
Close() error
// Metrics for monitoring (optional but useful for debugging)
QueueDepth() (messages, stream, errors, llmResponses int)
}
MessageQueue abstracts the message passing mechanism for chat sessions. All implementations must guarantee message delivery (no silent drops).
func CreateDefaultQueue ¶
func CreateDefaultQueue(sessionID string) (MessageQueue, error)
Helper function to create queue with automatic factory selection
func NewDefaultInMemoryQueue ¶
func NewDefaultInMemoryQueue(sessionID string) MessageQueue
Helper function for creating a queue with default settings (backward compatibility)
func NewDefaultNATSQueue ¶
func NewDefaultNATSQueue(sessionID, natsURL string) (MessageQueue, error)
Helper function for creating a NATS queue with default persistent settings
func NewDefaultPostgreSQLQueue ¶
func NewDefaultPostgreSQLQueue(sessionID string, db *gorm.DB) (MessageQueue, error)
Helper function for creating a PostgreSQL queue with default settings
type NATSConfig ¶
type NATSConfig struct {
URL string `json:"url"`
StorageType string `json:"storage_type"` // "memory" | "file"
RetentionPolicy string `json:"retention_policy"` // "limits" | "interest" | "workqueue"
MaxAge time.Duration `json:"max_age"`
MaxBytes int64 `json:"max_bytes"`
DurableConsumer bool `json:"durable_consumer"`
AckWait time.Duration `json:"ack_wait"`
MaxDeliver int `json:"max_deliver"`
BufferSize int `json:"buffer_size"`
FetchTimeout time.Duration `json:"fetch_timeout"` // Timeout for individual fetch operations
RetryInterval time.Duration `json:"retry_interval"` // Interval between fetch retries
MaxRetries int `json:"max_retries"` // Max retries for failed operations
// Authentication options
CredentialsFile string `json:"credentials_file"` // Optional NATS credentials file
Username string `json:"username"` // Optional username for basic auth
Password string `json:"password"` // Optional password for basic auth
Token string `json:"token"` // Optional token for token-based auth
NKeyFile string `json:"nkey_file"` // Optional NKey file path
// TLS options
TLSEnabled bool `json:"tls_enabled"` // Enable TLS connection
TLSCertFile string `json:"tls_cert_file"` // Optional client certificate file
TLSKeyFile string `json:"tls_key_file"` // Optional client key file
TLSCAFile string `json:"tls_ca_file"` // Optional CA certificate file
TLSSkipVerify bool `json:"tls_skip_verify"` // Skip TLS certificate verification
}
NATSConfig holds configuration for NATS JetStream
func DefaultNATSConfig ¶
func DefaultNATSConfig() NATSConfig
DefaultNATSConfig returns the hybrid persistent configuration (Option 3)
type NATSMessage ¶
type NATSMessage struct {
Type string `json:"type"`
SessionID string `json:"session_id"`
Timestamp time.Time `json:"timestamp"`
Data []byte `json:"data"`
}
NATSMessage wraps all message types with metadata
type NATSQueue ¶
type NATSQueue struct {
// contains filtered or unexported fields
}
NATSQueue implements MessageQueue using NATS JetStream Provides persistent message storage with automatic cleanup
func NewNATSQueue ¶
func NewNATSQueue(sessionID string, config NATSConfig) (*NATSQueue, error)
NewNATSQueue creates a new NATS-based message queue
func (*NATSQueue) ConsumeErrors ¶
ConsumeErrors returns the local channel for errors
func (*NATSQueue) ConsumeLLMResponses ¶
func (nq *NATSQueue) ConsumeLLMResponses(ctx context.Context) <-chan *LLMResponseWrapper
ConsumeLLMResponses returns the local channel for LLM responses
func (*NATSQueue) ConsumeMessages ¶
func (nq *NATSQueue) ConsumeMessages(ctx context.Context) <-chan *ChatResponse
ConsumeMessages returns the local channel for ChatResponse messages
func (*NATSQueue) ConsumeStream ¶
ConsumeStream returns the local channel for stream data
func (*NATSQueue) PublishError ¶
PublishError sends an error to NATS
func (*NATSQueue) PublishLLMResponse ¶
func (nq *NATSQueue) PublishLLMResponse(ctx context.Context, resp *LLMResponseWrapper) error
PublishLLMResponse sends an LLM response to NATS
func (*NATSQueue) PublishMessage ¶
func (nq *NATSQueue) PublishMessage(ctx context.Context, msg *ChatResponse) error
PublishMessage sends a ChatResponse message to NATS
func (*NATSQueue) PublishStream ¶
PublishStream sends stream data to NATS
func (*NATSQueue) QueueDepth ¶
QueueDepth returns the current depth of all local channels Note: This doesn't include messages pending in NATS streams
type NATSQueueFactory ¶
type NATSQueueFactory struct {
// contains filtered or unexported fields
}
NATSQueueFactory creates NATS queue instances
func NewNATSQueueFactory ¶
func NewNATSQueueFactory(config NATSConfig) *NATSQueueFactory
NewNATSQueueFactory creates a new NATS factory with specified configuration
func (*NATSQueueFactory) CreateQueue ¶
func (f *NATSQueueFactory) CreateQueue(sessionID string, config map[string]interface{}) (MessageQueue, error)
CreateQueue creates a new NATS queue with session-specific configuration
type OptimizedPostgreSQLQueue ¶
type OptimizedPostgreSQLQueue struct {
// contains filtered or unexported fields
}
OptimizedPostgreSQLQueue implements MessageQueue using PostgreSQL LISTEN/NOTIFY This version reuses the existing database connection to avoid connection exhaustion: NOTIFY goes through the application's pool, and notifications arrive on the process-wide shared listener (see queue_postgres_listener.go), so a session holds no connection of its own.
func NewOptimizedPostgreSQLQueue ¶
func NewOptimizedPostgreSQLQueue(sessionID string, db *gorm.DB, config PostgreSQLConfig) (*OptimizedPostgreSQLQueue, error)
NewOptimizedPostgreSQLQueue creates a new PostgreSQL-based message queue that reuses connections
func (*OptimizedPostgreSQLQueue) Close ¶
func (psq *OptimizedPostgreSQLQueue) Close() error
Close closes all channels and PostgreSQL connections
func (*OptimizedPostgreSQLQueue) ConsumeErrors ¶
func (psq *OptimizedPostgreSQLQueue) ConsumeErrors(ctx context.Context) <-chan error
ConsumeErrors returns the local channel for error messages
func (*OptimizedPostgreSQLQueue) ConsumeLLMResponses ¶
func (psq *OptimizedPostgreSQLQueue) ConsumeLLMResponses(ctx context.Context) <-chan *LLMResponseWrapper
ConsumeLLMResponses returns the local channel for LLM responses
func (*OptimizedPostgreSQLQueue) ConsumeMessages ¶
func (psq *OptimizedPostgreSQLQueue) ConsumeMessages(ctx context.Context) <-chan *ChatResponse
ConsumeMessages returns the local channel for ChatResponse messages
func (*OptimizedPostgreSQLQueue) ConsumeStream ¶
func (psq *OptimizedPostgreSQLQueue) ConsumeStream(ctx context.Context) <-chan []byte
ConsumeStream returns the local channel for stream data
func (*OptimizedPostgreSQLQueue) PublishError ¶
func (psq *OptimizedPostgreSQLQueue) PublishError(ctx context.Context, err error) error
PublishError sends an error via PostgreSQL NOTIFY
func (*OptimizedPostgreSQLQueue) PublishLLMResponse ¶
func (psq *OptimizedPostgreSQLQueue) PublishLLMResponse(ctx context.Context, resp *LLMResponseWrapper) error
PublishLLMResponse sends an LLM response via PostgreSQL NOTIFY
func (*OptimizedPostgreSQLQueue) PublishMessage ¶
func (psq *OptimizedPostgreSQLQueue) PublishMessage(ctx context.Context, msg *ChatResponse) error
PublishMessage sends a ChatResponse message via PostgreSQL NOTIFY
func (*OptimizedPostgreSQLQueue) PublishStream ¶
func (psq *OptimizedPostgreSQLQueue) PublishStream(ctx context.Context, data []byte) error
PublishStream sends stream data via PostgreSQL NOTIFY
func (*OptimizedPostgreSQLQueue) QueueDepth ¶
func (psq *OptimizedPostgreSQLQueue) QueueDepth() (messages, stream, errors, llmResponses int)
QueueDepth returns the current depth of all local channels
type OutputMode ¶
type OutputMode int
OutputMode selects how a session publishes its output on the queue.
OutputModeRaw is the v1 behaviour: text deltas are published as raw bytes on the stream channel, status lines as ":::system ...:::" messages and errors on the error channel. OutputModeEvents (v2) publishes every observable thing as a ChatEvent envelope on the stream channel, tagged with the run it belongs to, so a per-turn reader can follow exactly one turn.
A session's mode is fixed before Start() and a reader of the other mode must never be attached: the v1 SSE reader silently drops envelope JSON, and the v2 reader cannot interpret raw chunks.
const ( OutputModeRaw OutputMode = iota OutputModeEvents // OutputModeAny is accepted by callers that only mutate a session (add a // tool, upload a file) and never read its output, so they work for both // API versions. OutputModeAny )
type PostgreSQLConfig ¶
type PostgreSQLConfig struct {
BufferSize int // Local channel buffer size
ReconnectInterval time.Duration // Reconnection interval for listener
MaxReconnectRetries int // Maximum reconnection attempts
NotifyTimeout time.Duration // Timeout for NOTIFY operations
}
PostgreSQLConfig holds configuration for PostgreSQL queue
func DefaultPostgreSQLConfig ¶
func DefaultPostgreSQLConfig() PostgreSQLConfig
DefaultPostgreSQLConfig returns default configuration
type PostgreSQLMessage ¶
type PostgreSQLMessage struct {
Type string `json:"type"`
SessionID string `json:"session_id"`
Timestamp time.Time `json:"timestamp"`
Data []byte `json:"data"`
}
PostgreSQLMessage wraps all message types with metadata for JSON serialization
type PostgreSQLQueue ¶
type PostgreSQLQueue struct {
// contains filtered or unexported fields
}
PostgreSQLQueue implements MessageQueue using PostgreSQL LISTEN/NOTIFY Provides distributed message queuing with PostgreSQL pub/sub capabilities
func NewPostgreSQLQueue ¶
func NewPostgreSQLQueue(sessionID string, db *gorm.DB, config PostgreSQLConfig) (*PostgreSQLQueue, error)
NewPostgreSQLQueue creates a new PostgreSQL-based message queue
func (*PostgreSQLQueue) Close ¶
func (psq *PostgreSQLQueue) Close() error
Close closes all channels and PostgreSQL connections
func (*PostgreSQLQueue) ConsumeErrors ¶
func (psq *PostgreSQLQueue) ConsumeErrors(ctx context.Context) <-chan error
ConsumeErrors returns the local channel for error messages
func (*PostgreSQLQueue) ConsumeLLMResponses ¶
func (psq *PostgreSQLQueue) ConsumeLLMResponses(ctx context.Context) <-chan *LLMResponseWrapper
ConsumeLLMResponses returns the local channel for LLM responses
func (*PostgreSQLQueue) ConsumeMessages ¶
func (psq *PostgreSQLQueue) ConsumeMessages(ctx context.Context) <-chan *ChatResponse
ConsumeMessages returns the local channel for ChatResponse messages
func (*PostgreSQLQueue) ConsumeStream ¶
func (psq *PostgreSQLQueue) ConsumeStream(ctx context.Context) <-chan []byte
ConsumeStream returns the local channel for stream data
func (*PostgreSQLQueue) PublishError ¶
func (psq *PostgreSQLQueue) PublishError(ctx context.Context, err error) error
PublishError sends an error via PostgreSQL NOTIFY
func (*PostgreSQLQueue) PublishLLMResponse ¶
func (psq *PostgreSQLQueue) PublishLLMResponse(ctx context.Context, resp *LLMResponseWrapper) error
PublishLLMResponse sends an LLM response via PostgreSQL NOTIFY
func (*PostgreSQLQueue) PublishMessage ¶
func (psq *PostgreSQLQueue) PublishMessage(ctx context.Context, msg *ChatResponse) error
PublishMessage sends a ChatResponse message via PostgreSQL NOTIFY
func (*PostgreSQLQueue) PublishStream ¶
func (psq *PostgreSQLQueue) PublishStream(ctx context.Context, data []byte) error
PublishStream sends stream data via PostgreSQL NOTIFY
func (*PostgreSQLQueue) QueueDepth ¶
func (psq *PostgreSQLQueue) QueueDepth() (messages, stream, errors, llmResponses int)
QueueDepth returns the current depth of all local channels Note: This doesn't include messages pending in PostgreSQL notifications
type PostgreSQLQueueFactory ¶
type PostgreSQLQueueFactory struct {
// contains filtered or unexported fields
}
PostgreSQLQueueFactory creates PostgreSQL queue instances
func NewPostgreSQLQueueFactory ¶
func NewPostgreSQLQueueFactory(db *gorm.DB, config PostgreSQLConfig) *PostgreSQLQueueFactory
NewPostgreSQLQueueFactory creates a new PostgreSQL factory with specified configuration
func (*PostgreSQLQueueFactory) CreateQueue ¶
func (f *PostgreSQLQueueFactory) CreateQueue(sessionID string, config map[string]interface{}) (MessageQueue, error)
CreateQueue creates a new PostgreSQL queue with session-specific configuration
type QueueFactory ¶
type QueueFactory interface {
CreateQueue(sessionID string, config map[string]interface{}) (MessageQueue, error)
}
QueueFactory creates queue implementations based on configuration
func CreateDefaultQueueFactory ¶
func CreateDefaultQueueFactory() QueueFactory
CreateDefaultQueueFactory creates the default queue factory based on global configuration
func CreateDefaultQueueFactoryWithSharedDB ¶
func CreateDefaultQueueFactoryWithSharedDB(db *gorm.DB) QueueFactory
CreateDefaultQueueFactoryWithSharedDB creates the default queue factory with shared database
func CreateQueueFactory ¶
func CreateQueueFactory(cfg config.QueueConfig) (QueueFactory, error)
CreateQueueFactory creates appropriate queue factory based on configuration
func CreateQueueFactoryFromConfig ¶
func CreateQueueFactoryFromConfig() (QueueFactory, error)
CreateQueueFactoryFromConfig creates queue factory from global configuration
func CreateQueueFactoryWithSharedDB ¶
func CreateQueueFactoryWithSharedDB(cfg config.QueueConfig, db *gorm.DB) (QueueFactory, error)
CreateQueueFactoryWithSharedDB creates a queue factory with shared database connection This prevents connection exhaustion by reusing the application's connection pool
type SharedPostgreSQLQueueFactory ¶
type SharedPostgreSQLQueueFactory struct {
// contains filtered or unexported fields
}
SharedPostgreSQLQueueFactory creates PostgreSQL queues using a shared database connection This prevents connection exhaustion by reusing the application's existing connection pool
func NewSharedPostgreSQLQueueFactory ¶
func NewSharedPostgreSQLQueueFactory(db *gorm.DB, config PostgreSQLConfig) *SharedPostgreSQLQueueFactory
NewSharedPostgreSQLQueueFactory creates a shared PostgreSQL factory that reuses database connections
func (*SharedPostgreSQLQueueFactory) CreateQueue ¶
func (f *SharedPostgreSQLQueueFactory) CreateQueue(sessionID string, config map[string]interface{}) (MessageQueue, error)
CreateQueue creates a new PostgreSQL queue using the shared database connection
type StatusData ¶
type TextDeltaData ¶
type TextDeltaData struct {
Delta string `json:"delta"`
}
type ToolCallDeltaData ¶
type ToolCallEndData ¶
type ToolCallEndData struct {
ToolCallID string `json:"toolCallId"`
}
type ToolCallForNATS ¶
type ToolCallForNATS struct {
ID string `json:"id"`
Type string `json:"type"`
FunctionCall *FunctionCallForNATS `json:"function_call,omitempty"`
}