Documentation
¶
Index ¶
- func EnablePersistence(db *sql.DB)
- func EventStatsHandler(w http.ResponseWriter, r *http.Request)
- func SSEHandler(w http.ResponseWriter, r *http.Request)
- func SetCurrentSpan(span *Span)
- func SetCurrentTrace(trace *Trace)
- func SetEventBusDB(db *sql.DB)
- func WithSpan(ctx context.Context, span *Span) context.Context
- func WithTrace(ctx context.Context, trace *Trace) context.Context
- func WithTraceAndSpan(ctx context.Context, trace *Trace, span *Span) context.Context
- type BusConfig
- type Event
- func (e *Event) ToJSON() (string, error)
- func (e *Event) ToSSE() (string, error)
- func (e *Event) WithData(key string, value interface{}) *Event
- func (e *Event) WithDuration(startTime time.Time) *Event
- func (e *Event) WithError(err error) *Event
- func (e *Event) WithMetadata(key string, value interface{}) *Event
- func (e *Event) WithSession(sessionID string) *Event
- func (e *Event) WithSpan(spanID string) *Event
- func (e *Event) WithTask(taskID string) *Event
- func (e *Event) WithThread(threadID string) *Event
- func (e *Event) WithTrace(traceID string) *Event
- type EventBus
- func (eb *EventBus) GetRecentEvents(count int, filters SubscriptionFilters) []*Event
- func (eb *EventBus) GetStats() map[string]interface{}
- func (eb *EventBus) IsRunning() bool
- func (eb *EventBus) Publish(event *Event) error
- func (eb *EventBus) PublishWithContext(ctx context.Context, event *Event) error
- func (eb *EventBus) Start() error
- func (eb *EventBus) Stop()
- func (eb *EventBus) Subscribe(filters SubscriptionFilters) (*Subscriber, error)
- func (eb *EventBus) SubscribeFunc(filters SubscriptionFilters, callback EventCallback) (string, error)
- func (eb *EventBus) SubscriberCount() int
- func (eb *EventBus) Unsubscribe(subID string) error
- func (eb *EventBus) UnsubscribeFunc(subID string) error
- type EventCallback
- type EventCategory
- type EventLevel
- type EventType
- type FuncSubscriber
- type RingBuffer
- func (rb *RingBuffer) Add(event *Event)
- func (rb *RingBuffer) Capacity() int
- func (rb *RingBuffer) Clear()
- func (rb *RingBuffer) ForEach(fn func(*Event) bool)
- func (rb *RingBuffer) ForEachReverse(fn func(*Event) bool)
- func (rb *RingBuffer) GetAll() []*Event
- func (rb *RingBuffer) GetLast(n int) []*Event
- func (rb *RingBuffer) GetSince(eventID string) []*Event
- func (rb *RingBuffer) RemoveOlderThan(duration time.Duration) int
- func (rb *RingBuffer) Size() int
- type Span
- func (s *Span) End()
- func (s *Span) RecordError(err error) *Span
- func (s *Span) SetStatus(status SpanStatus) *Span
- func (s *Span) StartChild(name string, kind SpanKind) *Span
- func (s *Span) ToJSON() (string, error)
- func (s *Span) WithAttribute(key string, value interface{}) *Span
- func (s *Span) WithCategory(category EventCategory) *Span
- func (s *Span) WithDetailedTokenUsage(detail interface{ ... }) *Span
- func (s *Span) WithProviderModel(provider, model string) *Span
- func (s *Span) WithSessionID(sessionID string) *Span
- func (s *Span) WithTaskID(taskID string) *Span
- func (s *Span) WithThreadID(threadID string) *Span
- func (s *Span) WithTokenUsage(inputTokens, outputTokens int) *Span
- type SpanKind
- type SpanStatus
- type Subscriber
- type SubscriptionFilters
- type TelemetryConfig
- type TelemetryPayload
- type TelemetrySubscriber
- type Trace
- type TraceSource
- type TraceStatus
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func EnablePersistence ¶
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 ¶
SetEventBusDB sets the database connection for the global event bus
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) WithDuration ¶
WithDuration adds duration in milliseconds
func (*Event) WithMetadata ¶
WithMetadata adds metadata fields to the event
func (*Event) WithSession ¶
WithSession adds session context to the event
func (*Event) WithThread ¶
WithThread adds thread context to the event
type EventBus ¶
type EventBus struct {
// contains filtered or unexported fields
}
EventBus manages event publishing and subscriptions
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) Publish ¶
Publish sends an event to all matching subscribers (backwards compatible)
func (*EventBus) PublishWithContext ¶
PublishWithContext sends an event with trace context from the provided context
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
SubscriberCount returns the number of active subscribers
func (*EventBus) Unsubscribe ¶
Unsubscribe removes a subscription
func (*EventBus) UnsubscribeFunc ¶
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) Capacity ¶
func (rb *RingBuffer) Capacity() int
Capacity returns the maximum capacity of 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 SpanFromContext ¶
SpanFromContext extracts the span from context, or nil if not present
func (*Span) RecordError ¶
func (*Span) SetStatus ¶
func (s *Span) SetStatus(status SpanStatus) *Span
func (*Span) StartChild ¶
StartChild creates a child span from this span
func (*Span) WithAttribute ¶
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 (*Span) WithSessionID ¶
func (*Span) WithTaskID ¶
func (*Span) WithTokenUsage ¶
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 TraceFromContext ¶
TraceFromContext extracts the trace from context, or nil if not present
func (*Trace) UpdateThreadID ¶
UpdateThreadID updates the trace's thread ID
func (*Trace) WithParent ¶
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" )