realtime

package
v0.6.1 Latest Latest
Warning

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

Go to latest
Published: Sep 21, 2026 License: MIT Imports: 13 Imported by: 0

Documentation

Index

Constants

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

func (h *EventHub) ActiveClients() int64

ActiveClients reports currently-connected event-WS clients.

func (*EventHub) BroadcastLog

func (h *EventHub) BroadcastLog(l LogEntry)

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

func (h *EventHub) SetMaxClients(n int)

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

func (h *EventHub) SetOriginPolicy(enforce bool, allowedHosts []string)

SetOriginPolicy configures WebSocket origin enforcement — see Hub.SetOriginPolicy. Call once at startup.

func (*EventHub) SetTopologyProvider added in v0.5.0

func (h *EventHub) SetTopologyProvider(p topology.Provider)

SetTopologyProvider wires the mode-selected topology owner that answers host questions. Call once at startup, before the hub takes connections.

func (*EventHub) Start

func (h *EventHub) Start(ctx context.Context, snapshotInterval, batchInterval time.Duration)

Start begins the periodic flush loops. Call in a goroutine.

func (*EventHub) Stop

func (h *EventHub) Stop()

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 NewHub

func NewHub(onConnectionChange func(count int)) *Hub

NewHub creates a new buffered WebSocket hub.

func (*Hub) ActiveClients

func (h *Hub) ActiveClients() int64

ActiveClients reports the count of currently-connected WebSocket clients. Updated atomically as connections are accepted and torn down.

func (*Hub) Broadcast

func (h *Hub) Broadcast(entry LogEntry)

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

func (h *Hub) SetAggregateMode(on bool)

SetAggregateMode disables per-event log and metric broadcasts. Keepalive and connection handling are unaffected.

func (*Hub) SetDevMode

func (h *Hub) SetDevMode(devMode bool)

SetDevMode controls whether cross-origin WebSocket connections are accepted. Should be true only in development environments.

func (*Hub) SetMaxClients

func (h *Hub) SetMaxClients(n int)

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

func (h *Hub) SetOriginPolicy(enforce bool, allowedHosts []string)

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

func (h *Hub) SetTopologyProvider(provider topology.Provider, floor time.Duration)

SetTopologyProvider installs the construction-time provider used only for aggregate browser refresh notifications. Legacy raw batches are unchanged.

func (*Hub) SetWSMetrics

func (h *Hub) SetWSMetrics(onMessageSent func(string), onSlowClientDrop func())

SetWSMetrics wires WebSocket metric callbacks.

func (*Hub) Stop

func (h *Hub) Stop()

Stop gracefully shuts down the hub.

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.

Jump to

Keyboard shortcuts

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