Documentation
¶
Index ¶
- Constants
- type AggregatePublisher
- type EnginePublisher
- type EnginePublisherConfig
- type EventHub
- func (h *EventHub) ActiveClients() int64
- func (h *EventHub) BroadcastLog(l LogEntry)
- func (h *EventHub) BroadcastMetric(m MetricEntry)
- func (h *EventHub) HandleWebSocket(w http.ResponseWriter, r *http.Request)
- func (h *EventHub) NotifyRefresh()
- func (h *EventHub) SetAggregatePublisher(p AggregatePublisher, floor time.Duration)
- func (h *EventHub) SetMaxClients(n int)
- func (h *EventHub) SetOriginPolicy(enforce bool, allowedHosts []string)
- func (h *EventHub) SetTopologyProvider(p topology.Provider)
- func (h *EventHub) Start(ctx context.Context, snapshotInterval, batchInterval time.Duration)
- func (h *EventHub) Stop()
- type Hub
- func (h *Hub) ActiveClients() int64
- func (h *Hub) Broadcast(entry LogEntry)
- func (h *Hub) BroadcastMetric(entry MetricEntry)
- func (h *Hub) HandleWebSocket(w http.ResponseWriter, r *http.Request)
- func (h *Hub) Run()
- func (h *Hub) SetAggregateMode(on bool)
- func (h *Hub) SetDevMode(devMode bool)
- func (h *Hub) SetMaxClients(n int)
- func (h *Hub) SetOriginPolicy(enforce bool, allowedHosts []string)
- func (h *Hub) SetTopologyProvider(provider topology.Provider, floor time.Duration)
- func (h *Hub) SetWSMetrics(onMessageSent func(string), onSlowClientDrop func())
- func (h *Hub) Stop()
- type HubBatch
- type LiveSnapshot
- type LogEntry
- type MetricEntry
Constants ¶
const DefaultPublishFloor = 2 * time.Second
DefaultPublishFloor is the minimum spacing between aggregate publications. Frozen at 2 s in #164: fast enough to feel live, slow enough that a busy engine cannot turn every commit into a broadcast.
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type AggregatePublisher ¶ added in v0.5.0
type AggregatePublisher interface {
// Epoch identifies the process generation of Revision.
Epoch() string
// Revision is the aggregate engine's monotonic revision.
Revision() uint64
// Snapshot builds the coalesced payload for one service filter (empty
// means all services). It is summary, recent traffic, service health and
// topology — never the seven-day history.
Snapshot(ctx context.Context, service string) *LiveSnapshot
}
AggregatePublisher supplies revision-driven snapshots in aggregate mode. It is an interface so the hub needs neither the aggregate engine nor a live store to be tested.
type EnginePublisher ¶ added in v0.5.0
type EnginePublisher struct {
// contains filtered or unexported fields
}
EnginePublisher is the AggregatePublisher backed by the aggregate engine.
The payload it builds is deliberately the COALESCED one #164 froze: summary, recent traffic, service health and topology over a short trailing window. The seven-day history is never in a WebSocket message — a client that wants it asks for it over HTTP once, not on every revision bump.
func NewEnginePublisher ¶ added in v0.5.0
func NewEnginePublisher(cfg EnginePublisherConfig) *EnginePublisher
NewEnginePublisher builds a publisher over the engine. It returns nil when no engine is configured, which is what keeps the caller's wiring a one-liner.
func (*EnginePublisher) Epoch ¶ added in v0.5.0
func (p *EnginePublisher) Epoch() string
Epoch implements AggregatePublisher.
func (*EnginePublisher) Revision ¶ added in v0.5.0
func (p *EnginePublisher) Revision() uint64
Revision implements AggregatePublisher.
func (*EnginePublisher) Snapshot ¶ added in v0.5.0
func (p *EnginePublisher) Snapshot(ctx context.Context, service string) *LiveSnapshot
Snapshot implements AggregatePublisher.
type EnginePublisherConfig ¶ added in v0.5.0
type EnginePublisherConfig struct {
// Engine is the aggregate query facade. Required.
Engine *aggregate.Engine
// Topology is the same mode-selected provider injected into REST, MCP and
// GraphRAG. It is required and must be aggregate-owned.
Topology topology.Provider
// Window is the trailing range of the coalesced payload. Zero takes 15
// minutes, matching the legacy snapshot window.
Window time.Duration
// Tenant scopes every query. Empty takes storage.DefaultTenantID.
Tenant string
// Edges is IGNORED since #194 finding 15: topology edges come from the
// engine's own service-edge series, read in the same query as the nodes.
// The field remains so existing wiring compiles.
//
// Deprecated: supply nothing; QueryTopology carries the edges.
Edges func(ctx context.Context) []storage.ServiceMapEdge
}
EnginePublisherConfig configures an EnginePublisher.
type EventHub ¶
type EventHub struct {
// contains filtered or unexported fields
}
EventHub manages WebSocket clients and pushes live data snapshots filtered per-client's selected service. Debounces rapid ingestion bursts and only computes snapshots every flush interval.
func NewEventHub ¶
func NewEventHub(repo *storage.Repository, onConnect, onDisconnect func()) *EventHub
NewEventHub creates a new event notification hub.
func (*EventHub) ActiveClients ¶ added in v0.5.0
ActiveClients reports currently-connected event-WS clients.
func (*EventHub) BroadcastLog ¶
BroadcastLog adds a log entry to the real-time buffer. In aggregate mode per-event broadcasts are disabled and this is a no-op: the coalesced revision-driven snapshot is the only data message clients receive.
func (*EventHub) BroadcastMetric ¶
func (h *EventHub) BroadcastMetric(m MetricEntry)
BroadcastMetric adds a metric entry to the real-time buffer. Disabled in aggregate mode, for the same reason as BroadcastLog.
func (*EventHub) HandleWebSocket ¶
func (h *EventHub) HandleWebSocket(w http.ResponseWriter, r *http.Request)
HandleWebSocket upgrades an HTTP request to a WebSocket connection, registers it as an event client, and listens for filter messages.
The connection is scoped to exactly one tenant, taken from the principal the handshake gate authenticated. Writes are serialized through one bounded queue and one writer goroutine per client, so a stalled reader can neither block the snapshot loop nor grow without bound.
func (*EventHub) NotifyRefresh ¶
func (h *EventHub) NotifyRefresh()
notifyRefresh marks that new data has arrived. The actual snapshot happens on the next snapshotTicker flush.
func (*EventHub) SetAggregatePublisher ¶ added in v0.5.0
func (h *EventHub) SetAggregatePublisher(p AggregatePublisher, floor time.Duration)
SetAggregatePublisher switches the hub to revision-driven publication. Pass a zero floor to take DefaultPublishFloor. Call once at startup, before the hub takes connections.
func (*EventHub) SetMaxClients ¶ added in v0.5.0
SetMaxClients caps simultaneous event-WS connections. 0 disables the cap. Call once at startup, before the hub takes traffic.
func (*EventHub) SetOriginPolicy ¶ added in v0.5.0
SetOriginPolicy configures WebSocket origin enforcement — see Hub.SetOriginPolicy. Call once at startup.
func (*EventHub) SetTopologyProvider ¶ added in v0.5.0
SetTopologyProvider wires the mode-selected topology owner that answers host questions. Call once at startup, before the hub takes connections.
type Hub ¶
type Hub struct {
// contains filtered or unexported fields
}
Hub is a buffered WebSocket broadcast hub.
Instead of broadcasting each log individually (which would freeze the UI at high throughput), it buffers logs and flushes them as a JSON array when either:
- Buffer size >= maxBufferSize (default: 100)
- Flush ticker fires (default: every 500ms)
func (*Hub) ActiveClients ¶
ActiveClients reports the count of currently-connected WebSocket clients. Updated atomically as connections are accepted and torn down.
func (*Hub) Broadcast ¶
Broadcast adds a log entry to the broadcast buffer. No-op in aggregate mode.
func (*Hub) BroadcastMetric ¶
func (h *Hub) BroadcastMetric(entry MetricEntry)
BroadcastMetric adds a metric entry to the broadcast buffer. No-op in aggregate mode.
func (*Hub) HandleWebSocket ¶
func (h *Hub) HandleWebSocket(w http.ResponseWriter, r *http.Request)
HandleWebSocket is the HTTP handler that upgrades connections to WebSocket.
func (*Hub) Run ¶
func (h *Hub) Run()
Run starts the hub's main event loop. Should be called in a goroutine.
func (*Hub) SetAggregateMode ¶ added in v0.5.0
SetAggregateMode disables per-event log and metric broadcasts. Keepalive and connection handling are unaffected.
func (*Hub) SetDevMode ¶
SetDevMode controls whether cross-origin WebSocket connections are accepted. Should be true only in development environments.
func (*Hub) SetMaxClients ¶
SetMaxClients caps simultaneous WebSocket connections. 0 disables the cap (default). Configure once at startup before HandleWebSocket starts taking traffic — the cap is read concurrently from each upgrade attempt.
func (*Hub) SetOriginPolicy ¶ added in v0.5.0
SetOriginPolicy configures WebSocket origin enforcement. When enforce is true the browser Origin header must match one of allowedHosts, or the request host when the list is empty. Call once at startup.
func (*Hub) SetTopologyProvider ¶ added in v0.5.0
SetTopologyProvider installs the construction-time provider used only for aggregate browser refresh notifications. Legacy raw batches are unchanged.
func (*Hub) SetWSMetrics ¶
SetWSMetrics wires WebSocket metric callbacks.
type HubBatch ¶
type HubBatch struct {
Type string `json:"type"` // "logs" or "metrics"
Data any `json:"data"` // Slice of entries
}
HubBatch is a unified payload for WebSocket broadcasts.
type LiveSnapshot ¶
type LiveSnapshot struct {
Type string `json:"type"`
Dashboard *storage.DashboardStats `json:"dashboard"`
Traffic []storage.TrafficPoint `json:"traffic"`
Traces *storage.TracesResponse `json:"traces"`
ServiceMap *storage.ServiceMapMetrics `json:"service_map"`
Epoch string `json:"epoch,omitempty"`
Revision uint64 `json:"revision,omitempty"`
Reset bool `json:"reset,omitempty"`
Coverage string `json:"coverage,omitempty"`
CoverageNote string `json:"coverage_note,omitempty"`
Source string `json:"source,omitempty"`
Truncated bool `json:"truncated,omitempty"`
DroppedServices uint64 `json:"dropped_services,omitempty"`
DroppedOperations uint64 `json:"dropped_operations,omitempty"`
DroppedEdges uint64 `json:"dropped_edges,omitempty"`
DroppedMetrics uint64 `json:"dropped_metrics,omitempty"`
}
LiveSnapshot is the data payload pushed to all event WS clients.
The identity and coverage fields are ADDITIVE and only populated in aggregate mode; `omitempty` keeps the legacy payload unchanged.
A client replaces state by (series, window_start, revision) and RESETS wholesale when Epoch changes: the revision counter restarts at zero on every process generation, so revision alone cannot be trusted across a restart.
type LogEntry ¶
type LogEntry struct {
Tenant string `json:"-"`
ID uint `json:"id"`
TraceID string `json:"trace_id"`
SpanID string `json:"span_id"`
Severity string `json:"severity"`
Body string `json:"body"`
ServiceName string `json:"service_name"`
AttributesJSON string `json:"attributes_json"`
AIInsight string `json:"ai_insight,omitempty"`
Timestamp time.Time `json:"timestamp"`
}
LogEntry is a lightweight struct for WebSocket broadcast payloads.
Tenant is transport-only: it decides which sockets may see the entry and is never serialized, so the wire payload is unchanged from before per-tenant scoping existed.
type MetricEntry ¶
type MetricEntry struct {
Tenant string `json:"-"`
Name string `json:"name"`
ServiceName string `json:"service_name"`
Value float64 `json:"value"`
Timestamp time.Time `json:"timestamp"`
Attributes map[string]any `json:"attributes"`
}
MetricEntry represents a raw metric point for real-time visualization. Tenant is transport-only — see LogEntry.