chat_session

package
v2.2.1 Latest Latest
Warning

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

Go to latest
Published: Oct 2, 2026 License: AGPL-3.0 Imports: 38 Imported by: 0

Documentation

Index

Constants

View Source
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.

View Source
const (
	FinishStop      = "stop"
	FinishToolCalls = "tool-calls" // the turn is waiting on client-side tool results
	FinishError     = "error"
	FinishCancelled = "cancelled"
	FinishFiltered  = "filtered"
)

Finish reasons.

View Source
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.

View Source
const (
	ContextSourceRAG      = "rag"
	ContextSourceToolDocs = "tool_docs"
	ContextSourceFile     = "file"
)

Context sources carried in EventContext payloads.

View Source
const (
	MessageTypeChatResponse = "chat_response"
	MessageTypeStream       = "stream"
	MessageTypeError        = "error"
	MessageTypeLLMResponse  = "llm_response"
)

Message type constants

View Source
const (
	PostgreSQLMessageTypeChatResponse = "chat_response"
	PostgreSQLMessageTypeStream       = "stream"
	PostgreSQLMessageTypeError        = "error"
	PostgreSQLMessageTypeLLMResponse  = "llm_response"
)

Message type constants

Variables

This section is empty.

Functions

func ClassifyError

func ClassifyError(err error) string

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

func DecodeToolSpec(t *models.Tool) error

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 NewRunID

func NewRunID() string

NewRunID returns a fresh identifier for one user turn.

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

func UnwrapClientAnswer(content string) (interface{}, bool)

UnwrapClientAnswer is unwrapClientAnswer for other packages (history materialisation).

Types

type CallParams

type CallParams struct {
	Body       map[string]interface{} `json:"body"`
	Headers    map[string][]string    `json:"headers"`
	Parameters map[string][]string    `json:"parameters"`
}

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

func DecodeChatEvent(b []byte) (*ChatEvent, bool)

DecodeChatEvent reports whether b is an event envelope and decodes it.

type ChatMode

type ChatMode string
const (
	ChatStream  ChatMode = "stream"
	ChatMessage ChatMode = "message"
)

type ChatResponse

type ChatResponse struct {
	Payload string
}

type ChatSession

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

func NewChatSession

func NewChatSession(chat *models.Chat, mode ChatMode, db *gorm.DB, svc *services.Service, withFilters []*models.Filter, userID *uint, sessionID *string, queueFactory ...QueueFactory) (*ChatSession, error)

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) AddTool

func (cs *ChatSession) AddTool(id string, t models.Tool) 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 ContextData struct {
	Text   string `json:"text"`
	Source string `json:"source"`
}

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 ErrorData

type ErrorData struct {
	Code      string `json:"code"`
	Message   string `json:"message"`
	Detail    string `json:"detail,omitempty"`
	Retryable bool   `json:"retryable"`
}

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

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

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 (h *GormChatMessageHistory) AddAIToolCall(ctx context.Context, toolCall llms.ToolCall) error

func (*GormChatMessageHistory) AddMessage

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

Clear resets messages

func (*GormChatMessageHistory) GetAssociatedChat

func (h *GormChatMessageHistory) GetAssociatedChat(ctx context.Context) (*models.Chat, error)

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

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 LLMDriver

type LLMDriver interface {
	Call(ctx context.Context, inputs map[string]any, options ...chains.ChainCallOption) (map[string]any, error)
}

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) Close

func (nq *NATSQueue) Close() error

Close closes all channels and NATS connections

func (*NATSQueue) ConsumeErrors

func (nq *NATSQueue) ConsumeErrors(ctx context.Context) <-chan error

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

func (nq *NATSQueue) ConsumeStream(ctx context.Context) <-chan []byte

ConsumeStream returns the local channel for stream data

func (*NATSQueue) PublishError

func (nq *NATSQueue) PublishError(ctx context.Context, err error) error

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

func (nq *NATSQueue) PublishStream(ctx context.Context, data []byte) error

PublishStream sends stream data to NATS

func (*NATSQueue) QueueDepth

func (nq *NATSQueue) QueueDepth() (messages, stream, errors, llmResponses int)

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 StartData

type StartData struct {
	RunID string `json:"run_id"`
}

Event payloads.

type StatusData

type StatusData struct {
	Text  string `json:"text"`
	Level string `json:"level,omitempty"`
}

type TextDeltaData

type TextDeltaData struct {
	Delta string `json:"delta"`
}

type ToolCallDeltaData

type ToolCallDeltaData struct {
	ToolCallID string `json:"toolCallId"`
	ArgsText   string `json:"argsText"`
}

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"`
}

type ToolCallStartData

type ToolCallStartData struct {
	ToolCallID string `json:"toolCallId"`
	ToolName   string `json:"toolName"`
}

type ToolResultData

type ToolResultData struct {
	ToolCallID string `json:"toolCallId"`
	Result     string `json:"result"`
	IsError    bool   `json:"isError,omitempty"`
	Bytes      int    `json:"bytes,omitempty"`
}

Jump to

Keyboard shortcuts

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