Events
The framework's internal nervous system โ a typed, reflection-free publish/subscribe bus every subsystem uses to talk to your code. It imports nothing outside the standard library and works in a program that serves no HTTP at all.
import "github.com/nelthaarion/breeze/v2/events"
type UserCreated struct {
UserID uint64
Email string
}
sub := events.On(UserCreated{}, func(ctx *events.Context, e UserCreated) error {
fmt.Printf("Welcome, user %d!\n", e.UserID)
return nil
})
defer sub.Unsubscribe()
events.Emit(UserCreated{UserID: 42, Email: "user@example.com"}) Core concepts
Events are plain Go structs โ no interfaces, no embedding, no tags.
Registering a listener is events.On(EventType{}, handler), where the
handler receives *events.Context and the event value. events.Emit dispatches synchronously: listeners run in priority order, and the first
error stops propagation and is returned to the caller.
events.EmitAsync(OrderPlaced{OrderID: 123, Total: 99.99}) // don't wait Errors from an async emit go to Config.OnError; the call returns once
listeners are scheduled.
Priorities and phases
events.On(UserCreated{}, Validate).Priority(events.PriorityHighest)
events.On(UserCreated{}, Save).Priority(events.PriorityNormal)
events.On(UserCreated{}, SendEmail).Priority(events.PriorityLow) | Priority | Value | Typical use |
|---|---|---|
PriorityHighest | 1000 | validation, security checks |
PriorityHigh | 100 | normalization, enrichment |
PriorityNormal | 0 | default |
PriorityLow | -100 | persistence |
PriorityLowest | -1000 | auditing, metrics |
Listeners are also grouped into three phases, which always run in this order regardless of priority โ priority only controls order within a phase:
events.Before(RequestStarted{}, Logger) // phase 1: prepare state
events.On(RequestStarted{}, Handler) // phase 2: ordinary listeners
events.After(RequestStarted{}, Metrics) // phase 3: cleanup, metrics, auditing Filters, once-listeners, and stopping propagation
events.On(UserCreated{}, SendPromoEmail).
Where(func(e UserCreated) bool { return e.Age >= 18 })
events.Once(AppStarted{}, InitializeCache) // runs at most once, then removed
func CheckAuth(ctx *events.Context, e RequestStarted) error {
if !authorized(e.UserID) {
return events.Stop // normal control flow, not failure
}
return nil
} A filter is evaluated before the handler runs, so a rejected event never
enters the listener. events.Stop is consumed by the dispatcher and reported
as nil โ use it instead of a sentinel error for a guard clause that halts
dispatch on purpose.
Context and metadata
func Handler(ctx *events.Context, e MyEvent) error {
_ = ctx.Time // when the dispatch started
_ = ctx.EventID // unique id for this emit
_ = ctx.EventName // registered name, or the Go type
_ = ctx.Ctx // context.Context, never nil
ctx.Set("user", user)
if u, ok := ctx.Get("user"); ok { /* ... */ }
ctx.Cancel() // cancel remaining listeners
return nil
} GetMeta[T](ctx, key) is the type-safe accessor, returning (T, bool) โ
this is how one listener passes computed state to a later one at a lower
priority. Context is pooled for sync dispatch and must not be retained
past the handler's return; async listeners get a non-pooled Context.
Middleware
func Logger(ctx *events.Context, next events.Next) error {
log.Printf("dispatching %s", ctx.EventName)
err := next()
log.Printf("finished %s: %v", ctx.EventName, err)
return err
}
events.Default.Use(Logger) Middleware runs even when an event has no listeners, which makes it suitable for tracing and observability.
Error handling and panics
By default, the first error stops the dispatch and is returned. Enable ContinueOnError to run every listener and aggregate failures into a *events.MultiError. Panics are recovered by default:
bus := events.New(events.Config{
PanicMode: events.PanicRecoverAndContinue, // default
OnPanic: func(pe *events.PanicError) {
log.Printf("panic: %v\n%s", pe.Value, pe.Stack)
},
}) Modes: PanicRecoverAndContinue (recover, report, continue), PanicRecoverAndFail (recover, report, stop dispatch, return *PanicError), PanicPropagate (re-panic โ for tests).
Async dispatch
Goroutine mode (default) spawns one goroutine per listener โ lowest latency, unbounded concurrency. Worker pool mode bounds concurrency:
bus := events.New(events.Config{
Async: events.AsyncWorkerPool,
Workers: 8,
QueueSize: 512,
})
defer bus.Close()
events.EmitAsyncBus(bus, MyEvent{}) When the queue is full, AsyncOverflow picks the policy: OverflowSpawn (default, spawn a goroutine for the rejected task) or OverflowDrop (discard it). events.EmitAsyncWaitBus combines parallel execution with a
synchronous completion point.
Naming, inspection, and metrics
events.Name[UserCreated](bus, "user.created") // stable name for dashboards/logs
info := events.Inspect[UserCreated](bus)
fmt.Println(info.ListenerCount, info.Metrics.Dispatches)
m := events.MetricsFor[UserCreated](bus)
fmt.Println(m.Dispatches, m.Failures, m.AvgDuration)
total := bus.TotalMetrics() // aggregate across all events Disable per-event metrics with events.Config{DisableMetrics: true}.
Event recorder
bus.EnableRecorder()
// ... emit events ...
for _, rec := range bus.RecorderHistory() {
fmt.Printf("%s at %s: %d listeners, %v\n", rec.Name, rec.Time, rec.Listeners, rec.Duration)
} EnableRecorderWithPayload() also captures the event value itself. The
recorder is a ring buffer (Config.RecorderSize); old entries are evicted.
Bus isolation
events.New() builds an independent bus; the package-level functions
(events.Emit, events.On, ...) operate on events.Default. Multiple
independent buses are useful for test isolation.
Framework events
Breeze emits built-in events for application lifecycle, HTTP requests, OAuth2 token refresh, WebSocket connect/disconnect, scheduler jobs, and plugin install โ subscribe the same way as any custom event:
events.On(events.RequestFinished{}, func(ctx *events.Context, e events.RequestFinished) error {
log.Printf("%s %s -> %d (%v)", e.Method, e.Route, e.Status, e.Duration)
return nil
}) Configuration
bus := events.New(events.Config{
ContinueOnError: false,
PanicMode: events.PanicRecoverAndContinue,
OnPanic: func(*events.PanicError) { /* log it */ },
OnError: func(ctx *events.Context, listener string, err error) { /* ... */ },
Async: events.AsyncWorkerPool,
Workers: runtime.NumCPU(),
QueueSize: 512,
AsyncOverflow: events.OverflowSpawn,
Metrics: true,
Recorder: false,
RecorderSize: 256,
}) Performance
Benchmarked on an 11th Gen i5-11400F:
| Scenario | ns/op | B/op | allocs/op |
|---|---|---|---|
| Emit (1 listener) | 598 | 152 | 3 |
| Emit (100 listeners) | 6,519 | 152 | 3 |
| Emit (1000 listeners) | 50,968 | 152 | 3 |
| Emit (no listeners) | 31.8 | 0 | 0 |
| Middleware x3 | 2,914 | 248 | 7 |
| Async goroutine (1 listener) | 2,407 | 249 | 3 |
Allocation is flat at 152 bytes / 3 allocs from 1 to 1,000 listeners,
because the Context is pooled and listeners execute against a pre-sorted,
copy-on-write snapshot behind an atomic.Pointer โ no locks on the read
path. Emitting an event nobody listens to costs 31.8 ns with zero
allocations: a map lookup and an atomic load.
Best practices
- Register at startup, not inside a request handler โ registration rebuilds the snapshot.
- Reserve extreme priorities for validation and security; most
listeners belong at
PriorityNormal. - Prefer
Before/AfteroverPriority(9999)โ phases are clearer. - Use filters, not no-op handlers โ a filter runs before the handler.
- Offload heavy work to
EmitAsyncor a queue; dispatch is synchronous by default. - Never retain a pooled
*Contextpast the handler's return. - Name your events โ dashboards and logs prefer the registered name over a raw Go type name.
- Use
events.Stop, not a sentinel error, to halt propagation as normal control flow.