events

package
v1.55.0 Latest Latest
Warning

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

Go to latest
Published: Mar 20, 2026 License: MIT Imports: 13 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func EnablePersistence

func EnablePersistence(db *sql.DB)

EnablePersistence enables database persistence for events

func EventStatsHandler

func EventStatsHandler(w http.ResponseWriter, r *http.Request)

EventStatsHandler returns event bus statistics

func SSEHandler

func SSEHandler(w http.ResponseWriter, r *http.Request)

SSEHandler handles Server-Sent Events endpoint

func SetCurrentSpan

func SetCurrentSpan(span *Span)

func SetCurrentTrace

func SetCurrentTrace(trace *Trace)

func SetEventBusDB

func SetEventBusDB(db *sql.DB)

SetEventBusDB sets the database connection for the global event bus

func WithSpan

func WithSpan(ctx context.Context, span *Span) context.Context

WithSpan returns a new context with the span attached

func WithTrace

func WithTrace(ctx context.Context, trace *Trace) context.Context

WithTrace returns a new context with the trace attached

func WithTraceAndSpan

func WithTraceAndSpan(ctx context.Context, trace *Trace, span *Span) context.Context

WithTraceAndSpan is a convenience function to attach both trace and span to context

Types

type BusConfig

type BusConfig struct {
	BufferSize        int           // Size of the ring buffer
	SubscriberBuffer  int           // Default buffer size for subscribers
	RetentionPeriod   time.Duration // How long to keep events
	EnablePersistence bool          // Whether to persist events to database
}

BusConfig contains configuration for the event bus

func DefaultBusConfig

func DefaultBusConfig() BusConfig

DefaultBusConfig returns default configuration

type Event

type Event struct {
	ID         string                 `json:"id"`
	Timestamp  time.Time              `json:"timestamp"`
	Category   EventCategory          `json:"category"`
	Type       EventType              `json:"type"`
	Level      EventLevel             `json:"level"`
	ThreadID   string                 `json:"thread_id,omitempty"`
	TaskID     string                 `json:"task_id,omitempty"`
	SessionID  string                 `json:"session_id,omitempty"`
	TraceID    string                 `json:"trace_id,omitempty"` // Link to trace
	SpanID     string                 `json:"span_id,omitempty"`  // Link to span
	Data       map[string]interface{} `json:"data,omitempty"`
	Metadata   map[string]interface{} `json:"metadata,omitempty"`
	DurationMS int64                  `json:"duration_ms,omitempty"`
	Error      string                 `json:"error,omitempty"`
}

Event represents a single event in the system

func NewEvent

func NewEvent(category EventCategory, eventType EventType, level EventLevel) *Event

NewEvent creates a new event with generated ID and current timestamp

func (*Event) ToJSON

func (e *Event) ToJSON() (string, error)

ToJSON converts event to JSON string

func (*Event) ToSSE

func (e *Event) ToSSE() (string, error)

ToSSE formats the event for Server-Sent Events

func (*Event) WithData

func (e *Event) WithData(key string, value interface{}) *Event

WithData adds data fields to the event

func (*Event) WithDuration

func (e *Event) WithDuration(startTime time.Time) *Event

WithDuration adds duration in milliseconds

func (*Event) WithError

func (e *Event) WithError(err error) *Event

WithError adds error information to the event

func (*Event) WithMetadata

func (e *Event) WithMetadata(key string, value interface{}) *Event

WithMetadata adds metadata fields to the event

func (*Event) WithSession

func (e *Event) WithSession(sessionID string) *Event

WithSession adds session context to the event

func (*Event) WithSpan

func (e *Event) WithSpan(spanID string) *Event

WithSpan adds span context to the event

func (*Event) WithTask

func (e *Event) WithTask(taskID string) *Event

WithTask adds task context to the event

func (*Event) WithThread

func (e *Event) WithThread(threadID string) *Event

WithThread adds thread context to the event

func (*Event) WithTrace

func (e *Event) WithTrace(traceID string) *Event

WithTrace adds trace context to the event

type EventBus

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

EventBus manages event publishing and subscriptions

func GetEventBus

func GetEventBus() *EventBus

GetEventBus returns the global event bus instance

func NewEventBus

func NewEventBus(config BusConfig) *EventBus

NewEventBus creates a new event bus

func (*EventBus) GetRecentEvents

func (eb *EventBus) GetRecentEvents(count int, filters SubscriptionFilters) []*Event

GetRecentEvents returns recent events from the ring buffer Uses efficient iteration to avoid allocating the entire buffer

func (*EventBus) GetStats

func (eb *EventBus) GetStats() map[string]interface{}

GetStats returns statistics about the event bus

func (*EventBus) IsRunning added in v1.37.3

func (eb *EventBus) IsRunning() bool

IsRunning returns whether the event bus is running

func (*EventBus) Publish

func (eb *EventBus) Publish(event *Event) error

Publish sends an event to all matching subscribers (backwards compatible)

func (*EventBus) PublishWithContext

func (eb *EventBus) PublishWithContext(ctx context.Context, event *Event) error

PublishWithContext sends an event with trace context from the provided context

func (*EventBus) Start

func (eb *EventBus) Start() error

Start begins the event bus processing

func (*EventBus) Stop

func (eb *EventBus) Stop()

Stop gracefully stops the event bus

func (*EventBus) Subscribe

func (eb *EventBus) Subscribe(filters SubscriptionFilters) (*Subscriber, error)

Subscribe creates a new subscription with filters

func (*EventBus) SubscribeFunc

func (eb *EventBus) SubscribeFunc(filters SubscriptionFilters, callback EventCallback) (string, error)

SubscribeFunc creates a callback-based subscription (non-blocking, fire-and-forget) Returns the subscriber ID for later unsubscription

func (*EventBus) SubscriberCount added in v1.37.3

func (eb *EventBus) SubscriberCount() int

SubscriberCount returns the number of active subscribers

func (*EventBus) Unsubscribe

func (eb *EventBus) Unsubscribe(subID string) error

Unsubscribe removes a subscription

func (*EventBus) UnsubscribeFunc

func (eb *EventBus) UnsubscribeFunc(subID string) error

UnsubscribeFunc removes a callback-based subscription

type EventCallback

type EventCallback func(event *Event)

EventCallback is a function that handles events

type EventCategory

type EventCategory string

EventCategory represents the category of an event

const (
	CategorySystem     EventCategory = "SYSTEM"
	CategoryChat       EventCategory = "CHAT"
	CategoryTool       EventCategory = "TOOL"
	CategoryDatabase   EventCategory = "DATABASE"
	CategoryMCP        EventCategory = "MCP"
	CategoryScheduler  EventCategory = "SCHEDULER"
	CategoryLLM        EventCategory = "LLM"
	CategoryTask       EventCategory = "TASK"
	CategoryMemory     EventCategory = "MEMORY"
	CategoryAgent      EventCategory = "AGENT"
	CategoryError      EventCategory = "ERROR"
	CategoryFile       EventCategory = "FILE"
	CategoryReflection EventCategory = "REFLECTION"
)

type EventLevel

type EventLevel string

EventLevel represents the severity level of an event

const (
	LevelDebug EventLevel = "debug"
	LevelInfo  EventLevel = "info"
	LevelWarn  EventLevel = "warn"
	LevelError EventLevel = "error"
)

type EventType

type EventType string

EventType represents specific event types within categories

const (
	// System events
	TypeSystemStartup  EventType = "system_startup"
	TypeSystemShutdown EventType = "system_shutdown"
	TypeConfigChange   EventType = "config_change"

	// Chat events
	TypeChatMessageReceived  EventType = "message_received"
	TypeChatResponseStarted  EventType = "response_started"
	TypeChatStreamChunk      EventType = "stream_chunk"
	TypeChatResponseComplete EventType = "response_complete"
	TypeChatError            EventType = "chat_error"

	// Tool events
	TypeToolInvocation EventType = "tool_invocation"
	TypeToolResult     EventType = "tool_result"
	TypeToolError      EventType = "tool_error"

	// Database events
	TypeDBQuery  EventType = "db_query"
	TypeDBInsert EventType = "db_insert"
	TypeDBUpdate EventType = "db_update"
	TypeDBDelete EventType = "db_delete"
	TypeDBError  EventType = "db_error"

	// MCP events
	TypeMCPServerConnect EventType = "mcp_server_connect"
	TypeMCPToolDiscovery EventType = "mcp_tool_discovery"
	TypeMCPToolExecution EventType = "mcp_tool_execution"
	TypeMCPError         EventType = "mcp_error"

	// Scheduler events
	TypeSchedulerTaskStart    EventType = "scheduler_task_start"
	TypeSchedulerTaskComplete EventType = "scheduler_task_complete"
	TypeSchedulerTaskFailed   EventType = "scheduler_task_failed"

	// LLM events
	TypeLLMRequest  EventType = "llm_request"
	TypeLLMResponse EventType = "llm_response"
	TypeLLMTokens   EventType = "llm_tokens"
	TypeLLMError    EventType = "llm_error"

	// Task events
	TypeTaskCreated   EventType = "task_created"
	TypeTaskUpdated   EventType = "task_updated"
	TypeTaskExecuted  EventType = "task_executed"
	TypeTaskDeleted   EventType = "task_deleted"
	TypeTaskFailed    EventType = "task_failed"
	TypeTaskCompleted EventType = "task_completed" // Async task finished, result reported back
	TypeTaskListed    EventType = "task_listed"

	// File events
	TypeFileUploaded EventType = "file_uploaded"
	TypeFileDeleted  EventType = "file_deleted"
	TypeFileExpired  EventType = "file_expired"
	TypeFileIngested EventType = "file_ingested"
)

type FuncSubscriber

type FuncSubscriber struct {
	ID       string
	Filters  SubscriptionFilters
	Callback EventCallback
	Active   bool
	// contains filtered or unexported fields
}

FuncSubscriber wraps a callback function as a subscriber

type RingBuffer

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

RingBuffer is a circular buffer for storing events

func NewRingBuffer

func NewRingBuffer(capacity int) *RingBuffer

NewRingBuffer creates a new ring buffer with the specified capacity

func (*RingBuffer) Add

func (rb *RingBuffer) Add(event *Event)

Add adds an event to the buffer

func (*RingBuffer) Capacity

func (rb *RingBuffer) Capacity() int

Capacity returns the maximum capacity of the buffer

func (*RingBuffer) Clear

func (rb *RingBuffer) Clear()

Clear removes all events from the buffer

func (*RingBuffer) ForEach

func (rb *RingBuffer) ForEach(fn func(*Event) bool)

ForEach iterates over all events in chronological order The callback function receives each event and should return true to continue, false to stop

func (*RingBuffer) ForEachReverse

func (rb *RingBuffer) ForEachReverse(fn func(*Event) bool)

ForEachReverse iterates over all events in reverse chronological order (newest first) The callback function receives each event and should return true to continue, false to stop

func (*RingBuffer) GetAll

func (rb *RingBuffer) GetAll() []*Event

GetAll returns all events in the buffer in chronological order

func (*RingBuffer) GetLast

func (rb *RingBuffer) GetLast(n int) []*Event

GetLast returns the last N events from the buffer

func (*RingBuffer) GetSince

func (rb *RingBuffer) GetSince(eventID string) []*Event

GetSince returns all events since the specified event ID

func (*RingBuffer) RemoveOlderThan

func (rb *RingBuffer) RemoveOlderThan(duration time.Duration) int

RemoveOlderThan removes events older than the specified duration

func (*RingBuffer) Size

func (rb *RingBuffer) Size() int

Size returns the current number of events in the buffer

type Span

type Span struct {
	ID                  string                 `json:"id"`
	TraceID             string                 `json:"trace_id"`
	ParentSpanID        string                 `json:"parent_span_id,omitempty"`
	Name                string                 `json:"name"`
	Kind                SpanKind               `json:"kind"`
	StartTime           time.Time              `json:"start_time"`
	EndTime             *time.Time             `json:"end_time,omitempty"`
	DurationMS          int64                  `json:"duration_ms,omitempty"`
	Status              SpanStatus             `json:"status"`
	StatusMessage       string                 `json:"status_message,omitempty"`
	Category            EventCategory          `json:"category,omitempty"`
	ThreadID            string                 `json:"thread_id,omitempty"`
	TaskID              string                 `json:"task_id,omitempty"`
	SessionID           string                 `json:"session_id,omitempty"`
	AgentID             string                 `json:"agent_id,omitempty"`
	Attributes          map[string]interface{} `json:"attributes"`
	InputData           string                 `json:"input_data,omitempty"`
	OutputData          string                 `json:"output_data,omitempty"`
	TokenUsageIn        int                    `json:"token_usage_input,omitempty"`
	TokenUsageOut       int                    `json:"token_usage_output,omitempty"`
	CacheCreationTokens int                    `json:"cache_creation_tokens,omitempty"`
	CacheReadTokens     int                    `json:"cache_read_tokens,omitempty"`
	ReasoningTokens     int                    `json:"reasoning_tokens,omitempty"`
	Provider            string                 `json:"provider,omitempty"`
	Model               string                 `json:"model,omitempty"`

	// Not persisted, used for building tree
	Events     []*Event `json:"events,omitempty"`
	ChildSpans []*Span  `json:"child_spans,omitempty"`
	// contains filtered or unexported fields
}

Span represents a unit of work within a trace

func GetCurrentSpan

func GetCurrentSpan() *Span

func GetSpan

func GetSpan(db *sql.DB, spanID string) (*Span, error)

func GetSpansForTrace

func GetSpansForTrace(db *sql.DB, traceID string) ([]*Span, error)

func SpanFromContext

func SpanFromContext(ctx context.Context) *Span

SpanFromContext extracts the span from context, or nil if not present

func (*Span) End

func (s *Span) End()

End completes the span and persists it

func (*Span) RecordError

func (s *Span) RecordError(err error) *Span

func (*Span) SetStatus

func (s *Span) SetStatus(status SpanStatus) *Span

func (*Span) StartChild

func (s *Span) StartChild(name string, kind SpanKind) *Span

StartChild creates a child span from this span

func (*Span) ToJSON

func (s *Span) ToJSON() (string, error)

Helper to convert span to JSON

func (*Span) WithAttribute

func (s *Span) WithAttribute(key string, value interface{}) *Span

func (*Span) WithCategory

func (s *Span) WithCategory(category EventCategory) *Span

func (*Span) WithDetailedTokenUsage added in v1.44.0

func (s *Span) WithDetailedTokenUsage(detail interface {
	GetCacheCreation() int
	GetCacheRead() int
	GetReasoning() int
}) *Span

WithDetailedTokenUsage sets cache and reasoning token counts. Accepts a struct with the same field layout as stream.TokenUsageResult to avoid circular imports.

func (*Span) WithProviderModel added in v1.44.0

func (s *Span) WithProviderModel(provider, model string) *Span

func (*Span) WithSessionID

func (s *Span) WithSessionID(sessionID string) *Span

func (*Span) WithTaskID

func (s *Span) WithTaskID(taskID string) *Span

func (*Span) WithThreadID

func (s *Span) WithThreadID(threadID string) *Span

Span builder methods

func (*Span) WithTokenUsage

func (s *Span) WithTokenUsage(inputTokens, outputTokens int) *Span

type SpanKind

type SpanKind string

SpanKind represents the type of span

const (
	SpanKindInternal SpanKind = "internal"
	SpanKindServer   SpanKind = "server" // Incoming request
	SpanKindClient   SpanKind = "client" // Outgoing request
	SpanKindTool     SpanKind = "tool"   // Tool execution
	SpanKindLLM      SpanKind = "llm"    // LLM API call
	SpanKindAgent    SpanKind = "agent"  // Agent-to-agent call
)

type SpanStatus

type SpanStatus string

SpanStatus represents the completion status of a span

const (
	SpanStatusOK      SpanStatus = "ok"
	SpanStatusError   SpanStatus = "error"
	SpanStatusRunning SpanStatus = "running"
)

type Subscriber

type Subscriber struct {
	ID         string
	Channel    chan *Event
	Filters    SubscriptionFilters
	BufferSize int
	Active     bool
	// contains filtered or unexported fields
}

Subscriber represents an event subscriber

type SubscriptionFilters

type SubscriptionFilters struct {
	Categories []EventCategory
	Types      []EventType
	Levels     []EventLevel
	ThreadID   string
	SessionID  string
}

SubscriptionFilters defines filtering options for subscribers

type TelemetryConfig

type TelemetryConfig struct {
	Endpoint      string
	APIKey        string
	BatchSize     int // Default: 10
	FlushInterval int // Seconds, default: 30
}

TelemetryConfig holds configuration for the telemetry subscriber

type TelemetryPayload

type TelemetryPayload struct {
	AgentID string   `json:"agent_id"`
	Events  []*Event `json:"events"`
	SentAt  string   `json:"sent_at"`
}

TelemetryPayload is the payload sent to the telemetry endpoint

type TelemetrySubscriber

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

TelemetrySubscriber listens to events and sends them to a remote endpoint

func NewTelemetrySubscriber

func NewTelemetrySubscriber(cfg TelemetryConfig) *TelemetrySubscriber

NewTelemetrySubscriber creates a new telemetry subscriber

func (*TelemetrySubscriber) OnEvent

func (t *TelemetrySubscriber) OnEvent(event *Event)

OnEvent handles incoming events - called by the event bus

func (*TelemetrySubscriber) Stop

func (t *TelemetrySubscriber) Stop()

Stop gracefully stops the telemetry subscriber

type Trace

type Trace struct {
	ID            string                 `json:"id"`
	ParentTraceID string                 `json:"parent_trace_id,omitempty"`
	Source        TraceSource            `json:"source"`
	Name          string                 `json:"name"`
	ThreadID      string                 `json:"thread_id,omitempty"`
	ServiceName   string                 `json:"service_name"`
	StartTime     time.Time              `json:"start_time"`
	EndTime       *time.Time             `json:"end_time,omitempty"`
	DurationMS    int64                  `json:"duration_ms,omitempty"`
	Status        TraceStatus            `json:"status"`
	RootSpanID    string                 `json:"root_span_id,omitempty"`
	SpanCount     int                    `json:"span_count"`
	ErrorCount    int                    `json:"error_count"`
	Metadata      map[string]interface{} `json:"metadata"`
	// contains filtered or unexported fields
}

Trace represents a complete request/operation flow

func GetCurrentTrace

func GetCurrentTrace() *Trace

func GetTrace

func GetTrace(db *sql.DB, traceID string) (*Trace, error)

Query helpers

func NewTrace

func NewTrace(name, threadID string, db *sql.DB) *Trace

NewTrace creates a new trace

func TraceFromContext

func TraceFromContext(ctx context.Context) *Trace

TraceFromContext extracts the trace from context, or nil if not present

func (*Trace) Finish

func (t *Trace) Finish()

Finish completes the trace

func (*Trace) StartSpan

func (t *Trace) StartSpan(name string, kind SpanKind) *Span

StartSpan creates a new root span for this trace

func (*Trace) ToJSON

func (t *Trace) ToJSON() (string, error)

Helper to convert trace to JSON

func (*Trace) UpdateThreadID

func (t *Trace) UpdateThreadID(threadID string)

UpdateThreadID updates the trace's thread ID

func (*Trace) WithParent

func (t *Trace) WithParent(parentTraceID string) *Trace

WithParent sets the parent trace ID for hierarchical tracing

func (*Trace) WithSource

func (t *Trace) WithSource(source TraceSource) *Trace

WithSource sets the trace source

type TraceSource

type TraceSource string

TraceSource represents how the trace was initiated

const (
	TraceSourceChat      TraceSource = "chat"       // Direct /chat request
	TraceSourceWebhook   TraceSource = "webhook"    // Webhook trigger
	TraceSourceTask      TraceSource = "task"       // Scheduled task
	TraceSourceSubtask   TraceSource = "subtask"    // Spawned subtask
	TraceSourceAgentCall TraceSource = "agent_call" // Agent-to-agent call
)

type TraceStatus

type TraceStatus string

TraceStatus represents the completion status of a trace

const (
	TraceStatusOK         TraceStatus = "ok"
	TraceStatusError      TraceStatus = "error"
	TraceStatusUnfinished TraceStatus = "unfinished"
)

Jump to

Keyboard shortcuts

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