Skip to content

Latest commit

 

History

3 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

events

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.

Key Features

  • Typed Metadata & Defensive Cloning: Event type with severity, tags, data, correlation ID, plus safe cloning and optional Clonable
  • 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: emitterlogger wires events to your logger factory

Install

go get github.com/aatuh/pureapi-util/events

Subpackages are imported via their paths, for example:

import evt "github.com/aatuh/pureapi-util/events/event"

Quick start

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()
}

Advanced usage

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 method

Package structure

  • event: Core event types and DefaultEventEmitter
  • emitterlogger: Bridge events to your logging framework
  • examples: Executable examples as tests (progressive usage)

Run examples

go test -v -count 1 ./examples

# Run a specific example
go test -v -count 1 ./examples -run TestBasicEmissionAndListeners

API highlights

  • event.NewEvent(eventType, message, ...EventOption)
  • EventEmitter.RegisterListener/RegisterContextListener/RegisterOnce/...
  • EventEmitter.Emit, EmitSync, EmitContext
  • EventEmitter.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.

events

About

Fast, ergonomic event emitter for Go: context/global/once listeners, bounded worker pool with overflow policies, graceful shutdown, per-listener timeouts, defensive data cloning, metrics/error hooks, and a logging bridge.

Topics

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages