Observability
Records what actually happened inside a running application โ which events fired, which listeners ran, how long they took, what failed โ and exposes it for live streaming, querying, and dashboard display. Nothing is recorded unless you attach it; a detached application pays one atomic pointer load per dispatch and nothing more.
import (
"github.com/nelthaarion/breeze/v2/events"
"github.com/nelthaarion/breeze/v2/observability"
)
col := observability.NewCollector(observability.Config{Capacity: 1000, Metrics: true})
defer col.Close()
detach := observability.AttachEvents(events.Default, col)
defer detach()
for _, sig := range col.Recent(10) {
fmt.Printf("%s took %v (%d listeners)\n", sig.Name, sig.Duration, sig.Executed)
} AttachEvents returns the detach function, so observability can be switched
on and off at runtime with no residual cost.
Why a separate package
The event bus could have recorded its own history โ it deliberately does
not. The bus must stay dependency-free: if it knew about collectors, ring
buffers and masking rules, every application using it would carry them.
Observability also outlives any one subsystem โ the router, scheduler, and
database layers need the same treatment, and a shared Signal model means
they all appear on the same timeline. The coupling surface between the two
packages is one file and one four-method interface.
The Signal model
type Signal struct {
ID, SourceID, ParentID uint64
Source Source // "events", "router", "scheduler", ...
Kind Kind // "dispatch", "listener", "request", ...
Name string
Time time.Time
Duration time.Duration
Executed int
Children int
Failed, Cancelled, Async bool
Err string
CorrelationID string
RequestID string
Spans []Span
Attrs map[string]string
} A Span is one unit of work inside a signal โ for an event dispatch, one
listener, carrying its name, duration, priority, phase, and whether it
failed, was skipped, stopped, or panicked. The model is intentionally flat
rather than a tree: one lock acquisition to store it, one JSON object to
send.
Skipped, stopped, failed
Conflating these is the most common way to make a dashboard lie:
| State | Meaning | Counts as failure |
|---|---|---|
Skipped | a filter rejected the event, or a once-listener was spent | No |
Stopped | the listener returned events.Stop | No |
Failed | the listener returned an error, or panicked | Yes |
A guard clause that stops propagation is doing its job โ if it showed up red, every healthy request would look broken.
Cost
Measured on a 12-core machine, one listener:
| Benchmark | ns/op | B/op | allocs/op |
|---|---|---|---|
Dispatch_NoObserver | 280 | 168 | 3 |
Dispatch_AfterDetach | 349 | 168 | 3 |
Dispatch_Observer | 1077 | 424 | 7 |
Dispatch_ObserverWithPayload | 2938 | 1192 | 16 |
The line that matters is the second: after detaching, the allocation count returns exactly to baseline โ the hook leaves nothing behind. Attached, a dispatch costs roughly 800 ns and 4 extra allocations; payload capture roughly triples that, which is why it is off by default.
Collector
col := observability.NewCollector(observability.Config{Capacity: 1000, Metrics: true})
col.Snapshot() // every retained signal, oldest first
col.Recent(20) // the 20 newest, newest first
col.ByID(id)
col.Len()
col.Stats() // lifetime totals
col.Dropped() // signals dropped by slow subscribers
col.Clear() // drop retained signals, keep metrics
col.Reset() // drop everything There is also a process-wide observability.Default() for applications
that want one collector without threading it through.
Querying
col.Find(observability.Query{
Source: observability.SourceEvents,
NameContains: "user",
FailedOnly: true,
SlowerThan: 10 * time.Millisecond,
Since: time.Now().Add(-time.Hour),
Limit: 50,
Newest: true,
})
col.Slowest(10) // by duration, descending
col.TopNames(10) // by count; ties broken alphabetically
col.Names() // distinct names, sorted
col.Rate(time.Minute) // signals per second over a window Every field is optional; a zero Query matches everything.
Live streaming
ch, unsubscribe := col.Stream()
defer unsubscribe()
for sig := range ch {
fmt.Println(sig.Name, sig.Duration)
} A subscriber that stops reading is dropped, never blocks, the
publisher โ the drop is counted in col.Dropped(), so a stalled dashboard
tab cannot slow down the application it is observing. Signals delivered to
subscribers are independent copies.
Metrics
m := col.MetricFor("user.created")
m.Count; m.Executed; m.Failed; m.Total; m.Min; m.Max; m.Avg; m.Last
m.FailureRate() Metrics are keyed by (Source, Name), not by name alone โ a route and an
event may both be called user.created, and their statistics must not
merge. MetricForSource(src, name) is the exact lookup; MetricFor scans
by name and returns the first match. Metrics are lifetime totals and
survive ring-buffer eviction.
The execution graph
for _, node := range col.Graph() {
fmt.Println(node.Name, node.Count)
for _, edge := range node.Edges {
fmt.Printf(" โ %s (%d calls, %d failed)\n", edge.Target, edge.Count, edge.Failed)
}
} Built from observed execution, not the registry โ the registry knows what is registered, the graph knows what actually ran, in what order, and how often it failed.
Payload capture and masking
Off by default:
detach := observability.AttachEventsWithPayload(bus, col) The payload is rendered to a short string immediately and stored as text, never as a live reference. Field names that look sensitive are masked before storage:
observability.IsSensitive("password") // true
observability.IsSensitive("api_key") // true
observability.IsSensitive("user_id") // false
observability.MaskAttrs(map[string]string{
"user_id": "42", // kept
"password": "hunter2", // โ "[MASKED]"
}) Masking runs at capture time, so a secret never reaches the ring buffer, the stream, or the dashboard โ a safety net, not a licence to put secrets in event payloads.
Configuration
type Config struct {
Capacity int // retained signals; default 1000
Metrics bool // aggregate per-name metrics
ErrSink func(error) // internal error reporting
} Extending to other subsystems
Signal is generic on purpose โ a new producer publishes directly:
col.Publish(observability.Signal{
Source: observability.SourceRouter,
Kind: observability.KindDispatch,
Name: "GET /users/:id",
Time: start,
Duration: time.Since(start),
RequestID: reqID,
}) Predefined sources: SourceEvents, SourceRouter, SourceHTTP, SourceCache, SourceDatabase, SourceScheduler, SourceWebSocket, SourceOAuth2, SourcePlugin. A RequestID shared between a router
signal and the event dispatches it triggered reconstructs the full causal
chain across subsystems.
Design decisions
- One observer per bus โ fan-out is the collector's job, which keeps the hot path a single pointer load.
- The bus never imports this package โ it defines an interface; this package implements it.
- Observer panics are contained โ recovered, reported, attributed to
<observer>, and never charged to a listener's own panic counter. - The context is never retained past a dispatch.
- In-flight state is bounded and sharded โ 16-way sharded map so concurrent dispatches don't serialise on one lock; spans per dispatch are capped at 256.