Decoupled, typed event metadata + defensive data cloning for Go: emit, listen, observe, and log.
The events module provides a robust event system with synchronous and
asynchronous delivery, context-aware listeners, timeouts and cancellation,
queue overflow policies, metrics hooks, and structured logging integration
via emitterlogger.
- Typed Metadata & Defensive Cloning:
Eventtype with severity, tags, data, correlation ID, plus safe cloning and optionalClonable - Flexible Listeners: local, global, context-aware, and once-only listeners
- Sync and Async:
EmitSync, worker pool for async, and graceful shutdown - Timeouts & Cancellation: per-listener timeouts, context cancellation
- Backpressure Control: queue size, workers, and overflow policies
- Observability: error handler and metrics hook (duration, error, queue depth)
- Safe Payloads: defensive cloning for maps, slices, and optional
Clonable - Logging Bridge:
emitterloggerwires events to your logger factory
go get github.com/aatuh/pureapi-util/eventsSubpackages are imported via their paths, for example:
import evt "github.com/aatuh/pureapi-util/events/event"Register a listener and emit events:
package main
import (
"context"
"log"
"time"
evt "github.com/aatuh/pureapi-util/events/event"
)
func main() {
emitter := evt.NewEventEmitter(
evt.WithWorkerPool(2, 16),
)
emitter.RegisterListener("user.created", func(e *evt.Event) {
log.Printf("%s: %s", e.Type, e.Message)
})
if err := emitter.Emit(evt.NewEvent("user.created", "user 123")); err != nil {
log.Printf("emit failed: %v", err)
}
if err := emitter.EmitSync(evt.NewEvent("user.created", "user 456")); err != nil {
log.Printf("emit sync failed: %v", err)
}
ctx, cancel := context.WithTimeout(context.Background(), time.Second)
defer cancel()
_ = emitter.Flush(ctx)
_ = emitter.Close()
}Context-aware and once listeners:
emitter.RegisterContextListener("job.start", func(ctx context.Context, e *evt.Event) {
// Observe cancellation, deadlines, etc.
})
emitter.RegisterOnce("job.start", func(e *evt.Event) { /* only once */ })Timeouts and error handling:
emitter := evt.NewEventEmitter(
evt.WithDefaultTimeout(100*time.Millisecond),
evt.WithErrorHandler(func(e *evt.Event, id evt.ListenerID, err error) {
// capture timeouts, panics, cancellations, overflow
}),
)Metrics hook (duration, error, queue depth):
emitter := evt.NewEventEmitter(
evt.WithMetricsHook(func(e *evt.Event, id evt.ListenerID, d time.Duration, err error, depth int) {
// export metrics
}),
)Worker pool and overflow policy:
emitter := evt.NewEventEmitter(
evt.WithWorkerPool(4, 64), // 4 workers, queue capacity 64
evt.WithOverflowPolicy(evt.OverflowPolicyDropNewest),
)Logging bridge with emitterlogger:
import (
"github.com/aatuh/pureapi-util/events/emitterlogger"
evt "github.com/aatuh/pureapi-util/events/event"
)
// Your logger implementing emitterlogger.ILogger
factory := func(params ...any) emitterlogger.ILogger { return myLogger }
el := emitterlogger.NewEmitterLogger(emitter, factory)
el.Info(evt.NewEvent("app.msg", "hello world"), "serviceA")
el.Log(evt.NewEvent("app.msg", "warn")) // severity routes to the right methodevent: Core event types andDefaultEventEmitteremitterlogger: Bridge events to your logging frameworkexamples: Executable examples as tests (progressive usage)
go test -v -count 1 ./examples
# Run a specific example
go test -v -count 1 ./examples -run TestBasicEmissionAndListenersevent.NewEvent(eventType, message, ...EventOption)EventEmitter.RegisterListener/RegisterContextListener/RegisterOnce/...EventEmitter.Emit,EmitSync,EmitContextEventEmitter.Flush,Close,CloseWithTimeout- Options:
WithDefaultTimeout,WithWorkerPool,WithOverflowPolicy,WithErrorHandler,WithMetricsHook
See the examples directory for a guided, progressive tour of the API, from basic to advanced scenarios.