dev.club — where best developers and top companies connect.

dev.club — where best developers and top companies connect.Invite only

Request invite

GoEventBus

A high-performance event bus for Go with deterministic rules, cached decisions, and optional Jev-powered event routing.

GoEventBus combines a bounded MPMC in-memory event bus with fan-out, ordered and batch handlers, middleware, lifecycle hooks, dead-letter handling, Redis Streams, RabbitMQ, and an optional decision layer for choosing which event type should execute.

Go Report Card

                           ┌──────────────┐
                           │    Rules     │
                           └──────┬───────┘
                                  │ no match
                                  ▼
Event state ───────────────► ┌──────────────┐
                             │    Cache     │
                             └──────┬───────┘
                                    │ miss
                                    ▼
                             ┌──────────────┐
                             │     Jev      │
                             └──────┬───────┘
                                    │
                                    ▼
                             selected event type
                                    │
                                    ▼
Subscribe ─► MPMC ring ─► Publish ─► handlers
                         ├─────────► ordered handlers
                         └─────────► batch handlers

The decision layer is optional. If you already know the event type, call Subscribe directly and GoEventBus behaves like a normal high-performance event bus.

Why GoEventBus?

Most events do not need model inference. Some do.

GoEventBus keeps those paths separate:

That gives you intelligent routing without putting an LLM call in the Publish() hot path.

Features

Installation

go get github.com/Protocol-Lattice/GoEventBus

GoEventBus currently targets Go 1.23+.


Quick Start

Plain event bus

Use the direct path when the event type is already known.

package main

import (
    "context"
    "fmt"
    "log"

    GoEventBus "github.com/Protocol-Lattice/GoEventBus"
)

type UserCreated struct {
    ID string
}

func main() {
    dispatcher := GoEventBus.Dispatcher{}

    dispatcher.Register("user.created", func(
        ctx context.Context,
        ev GoEventBus.Event,
    ) (GoEventBus.Result, error) {
        payload := ev.Data.(UserCreated)
        fmt.Println("user created:", payload.ID)
        return GoEventBus.Result{Message: "ok"}, nil
    })

    store := GoEventBus.NewEventStore(
        &dispatcher,
        1<<16,
        GoEventBus.DropOldest,
    )
    store.Async = true

    if err := store.Subscribe(context.Background(), GoEventBus.Event{
        ID:         "evt-1",
        Projection: "user.created",
        Data:       UserCreated{ID: "u-42"},
    }); err != nil {
        log.Fatal(err)
    }

    store.Publish()

    if err := store.Drain(context.Background()); err != nil {
        log.Fatal(err)
    }
}

Projection may be any comparable value for local dispatch, including a string or a typed struct.


Intelligent Event Routing

GoEventBus can choose an event type before enqueueing it.

The recommended pipeline is:

rules → cache → Jev → cache write → DecideAndSubscribe → local queue or configured broker

This keeps the fast path deterministic and only calls Jev when local logic cannot resolve the event type.

1. Define candidates

type HouseWasSold struct{}

candidates := []GoEventBus.EventCandidate{
    {
        Key:         "user_created",
        Projection:  "user.created",
        Description: "A new user account was created",
    },
    {
        Key:         "house_sold",
        Projection:  HouseWasSold{},
        Description: "A property sale was completed",
    },
    {
        Key:         "order_cancelled",
        Projection:  "order.cancelled",
        Description: "An existing order should be cancelled",
    },
}

Key is the stable identifier exposed to the decision layer. Projection is the real GoEventBus dispatcher key.

That separation lets Jev choose between simple string keys while the event bus can still dispatch to typed projections.

2. Add deterministic rules

Rules run first. The first matching rule wins.

rules := []GoEventBus.EventRule{
    {
        Name:   "explicit-cancel",
        Choice: "order_cancelled",
        Match: func(
            _ context.Context,
            state any,
            _ []GoEventBus.EventCandidate,
        ) bool {
            input, ok := state.(map[string]any)
            return ok && input["action"] == "cancel"
        },
    },
}

Rules are intentionally evaluated before cache. Adding a new rule can therefore override an older cached model decision immediately.

3. Add a decision cache

cache := GoEventBus.NewMemoryDecisionCache(5 * time.Minute)

The built-in cache is goroutine-safe and lazily expires entries.

The default cache key is a SHA-256 hash of:

Candidate ordering does not change the key. Projection values are deliberately excluded because typed projections may not be JSON-serializable.

Need a different strategy? Supply RuleCacheSelector.CacheKey.

4. Use Jev as the fallback

jev := &GoEventBus.JevSelector{
    APIKey: os.Getenv("OPENROUTER_API_KEY"),
}

The default model is:

typesafe/jev-1.13

The adapter calls OpenRouter's Decisions API using the Go standard library, so the Jev integration adds no SDK dependency.

5. Compose the router

selector := &GoEventBus.RuleCacheSelector{
    Rules:    rules,
    Cache:    cache,
    Fallback: jev,
}

Then route and deliver:

state := map[string]any{
    "message": "The property at 1 Main St was sold for $500000",
}

decision, err := store.DecideAndSubscribe(
    context.Background(),
    selector,
    state,
    GoEventBus.Event{
        ID: "evt-2",
    },
    candidates,
)
if err != nil {
    log.Fatal(err)
}

fmt.Printf(
    "choice=%s confidence=%.2f probabilities=%v\n",
    decision.Choice,
    decision.Confidence,
    decision.Probabilities,
)

store.Publish()

With a normal three-argument NewEventStore, DecideAndSubscribe enqueues the chosen event locally, so call Publish() as above.

With WithRedis(...), WithRabbitMQ(...), or WithProvider(...), the same DecideAndSubscribe call publishes the chosen event directly to the configured broker instead. The consumer side runs store.Consume(ctx) and dispatches broker events through the normal local handlers. Broker-backed candidates must use string projections.

If a rule matches, Jev is never called. If the same decision state is already cached, Jev is never called. Only a cache miss reaches the model.

Cache behavior

DecisionCache is intentionally small:

type DecisionCache interface {
    Get(context.Context, string) (EventDecision, bool, error)
    Set(context.Context, string, EventDecision) error
}

That makes it straightforward to implement Redis, distributed, persistent, or application-specific caches.

Cache failures are fail-open by default. Routing continues through the fallback selector because cache is treated as an optimization.

Set:

StrictCache: true

when cache failures should fail the routing call instead.


Examples

The repository includes runnable examples for both the classic event-bus path and the intelligent routing layer.

Example What it demonstrates
routing_rules Deterministic first-match routing with EventRule; no model or cache required
routing_cache Reusing a previous decision so the fallback selector is called only once
routing_jev Full rules → cache → Jev → GoEventBus routing with OpenRouter
transaction_jev Rules/Cache/Jev selection buffered until commit, then dispatched locally or published through the configured provider
schedule_jev Rules/Cache/Jev event selection before delayed scheduling
hello_world Minimal direct Subscribe → Publish flow
middleware Handler middleware and lifecycle behavior
goroutines-subscribe-publisher Concurrent producers and publishing
drop_oldest DropOldest back-pressure
return_error ReturnError back-pressure
handler_timeout Handler context timeout
publisher_timeout Blocking publisher timeout
fasthttp HTTP integration
redis Redis Streams configured directly on NewEventStore
rabbitmq RabbitMQ configured directly on NewEventStore

Run the new routing examples directly:

go run ./examples/routing_rules
go run ./examples/routing_cache
OPENROUTER_API_KEY=... go run ./examples/routing_jev
OPENROUTER_API_KEY=... go run ./examples/transaction_jev
OPENROUTER_API_KEY=... go run ./examples/schedule_jev

See examples/README.md for a compact guide to all examples.


Core Event Model

type Event struct {
    ID         string
    Projection interface{}
    Data       any
    Args       map[string]any // deprecated
}

Prefer Data for payloads. Args remains for backwards compatibility.

Handlers use:

type HandlerFunc func(
    context.Context,
    Event,
) (Result, error)

Register handlers with a Dispatcher:

dispatcher := GoEventBus.Dispatcher{}

dispatcher.Register(
    "order.created",
    handleBilling,
    handleAnalytics,
)

Calling Register repeatedly for the same projection appends handlers rather than replacing them.


Fan-out

Multiple handlers may subscribe to the same projection.

dispatcher.Register(
    "order.placed",
    auditLogger,
    inventoryReducer,
    notificationSender,
)

Each handler is independent. One handler returning an error does not prevent the other fan-out handlers from running.

In async mode, each regular handler invocation becomes its own worker-pool item.


Ordered Handlers

Use ordered handlers when events for the same entity must be processed sequentially while unrelated entities can still run concurrently.

type OrderEvent struct {
    OrderID  string
    Sequence int
}

store.RegisterOrdered(
    "order.event",
    func(ev GoEventBus.Event) string {
        return ev.Data.(OrderEvent).OrderID
    },
    func(
        ctx context.Context,
        ev GoEventBus.Event,
    ) (GoEventBus.Result, error) {
        return applyOrderEvent(ctx, ev.Data.(OrderEvent))
    },
)

For a given ordering key, async delivery preserves FIFO order.

Different keys remain concurrent.


Batch Handlers

Batch handlers collect pending events for a projection during a Publish cycle and deliver them in chunks.

store.RegisterBatch(
    "metrics.recorded",
    100,
    func(
        ctx context.Context,
        events []GoEventBus.Event,
    ) ([]GoEventBus.Result, error) {
        return nil, writeMetricsBatch(ctx, events)
    },
)

Regular and batch handlers may coexist on the same projection.

Multiple batch handlers may also be registered for one projection; each receives the full chunk independently.

Middleware is not applied to batch handlers. Lifecycle hooks are.


Middleware

Middleware wraps regular handlers.

store.Use(func(next GoEventBus.HandlerFunc) GoEventBus.HandlerFunc {
    return func(
        ctx context.Context,
        ev GoEventBus.Event,
    ) (GoEventBus.Result, error) {
        started := time.Now()
        result, err := next(ctx, ev)
        log.Printf(
            "projection=%v duration=%s err=%v",
            ev.Projection,
            time.Since(started),
            err,
        )
        return result, err
    }
})

Middleware is applied independently to every handler in a fan-out.


Lifecycle Hooks

Use hooks for observability without changing handler logic.

store.OnBefore(func(ctx context.Context, ev GoEventBus.Event) {
    metrics.Inc("handler.started")
})

store.OnAfter(func(
    ctx context.Context,
    ev GoEventBus.Event,
    result GoEventBus.Result,
    err error,
) {
    metrics.Inc("handler.finished")
})

store.OnError(func(
    ctx context.Context,
    ev GoEventBus.Event,
    err error,
) {
    log.Printf("event=%s err=%v", ev.ID, err)
})

For batch handlers, hooks are emitted per event.


Back-pressure

Choose how Subscribe behaves when the ring buffer is full.

Policy Behavior
DropOldest Evict the oldest queued event and accept the new one
Block Wait for capacity while respecting the caller context
ReturnError Return ErrBufferFull immediately

Example:

store := GoEventBus.NewEventStore(
    &dispatcher,
    1<<14,
    GoEventBus.Block,
)

ctx, cancel := context.WithTimeout(
    context.Background(),
    50*time.Millisecond,
)
defer cancel()

if err := store.Subscribe(ctx, event); err != nil {
    log.Println("enqueue failed:", err)
}

Async Mode

Enable worker-pool dispatch with:

store.Async = true

The store uses a fixed worker pool sized to runtime.NumCPU().

Publish() submits regular, ordered, and batch work to that pool.

Call Drain or Close when shutting down:

ctx, cancel := context.WithTimeout(
    context.Background(),
    5*time.Second,
)
defer cancel()

if err := store.Drain(ctx); err != nil {
    log.Println("drain failed:", err)
}

Once shutdown begins, new subscriptions return ErrEventStoreClosed.


External Providers

Redis Streams and RabbitMQ are optional transports configured directly on the EventStore. The original three-argument constructor still works unchanged:

store := GoEventBus.NewEventStore(
    &dispatcher,
    1024,
    GoEventBus.Block,
)

Add a broker as a fourth functional option when cross-process delivery is needed:

store := GoEventBus.NewEventStore(
    &dispatcher,
    1024,
    GoEventBus.Block,
    GoEventBus.WithRedis(redisConfig),
)

The configured provider is created lazily on the first broker operation and is owned by the store. Drain / Close closes it. The Redis provider still does not close an injected Redis client, so shared client ownership remains with the application.

Broker-backed stores integrate directly with intelligent routing:

Remote providers require a string projection because the projection becomes the broker routing name.

The Redis and RabbitMQ snippets below reuse the selector and state from the intelligent-routing section above; only the EventStore transport option changes.

Redis Streams

rdb := redis.NewClient(&redis.Options{
    Addr: "localhost:6379",
})

store := GoEventBus.NewEventStore(
    &dispatcher,
    1024,
    GoEventBus.Block,
    GoEventBus.WithRedis(GoEventBus.RedisProviderConfig{
        Client:   rdb,
        Stream:   "events",
        Group:    "billing",
        Consumer: "billing-1",
        StartID:  "0",
    }),
)

go func() {
    if err := store.Consume(ctx); err != nil &&
        !errors.Is(err, context.Canceled) &&
        !errors.Is(err, GoEventBus.ErrProviderClosed) {
        log.Println(err)
    }
}()

decision, err := store.DecideAndSubscribe(
    ctx,
    selector,
    state,
    GoEventBus.Event{
        ID:   "evt-redis-1",
        Data: order,
    },
    []GoEventBus.EventCandidate{{
        Key:         "order_created",
        Projection:  "order.created",
        Description: "A new order should be created",
    }},
)
if err != nil {
    log.Fatal(err)
}
fmt.Println("published decision:", decision.Choice)

RabbitMQ

store := GoEventBus.NewEventStore(
    &dispatcher,
    1024,
    GoEventBus.Block,
    GoEventBus.WithRabbitMQ(GoEventBus.RabbitMQProviderConfig{
        URL:        "amqp://guest:guest@localhost:5672/",
        Exchange:   "events",
        Queue:      "billing",
        BindingKey: "order.*",
        Consumer:   "billing-1",
    }),
)

go func() {
    if err := store.Consume(ctx); err != nil &&
        !errors.Is(err, context.Canceled) &&
        !errors.Is(err, GoEventBus.ErrProviderClosed) {
        log.Println(err)
    }
}()

decision, err := store.DecideAndSubscribe(
    ctx,
    selector,
    state,
    GoEventBus.Event{
        ID:   "evt-rabbit-1",
        Data: order,
    },
    []GoEventBus.EventCandidate{{
        Key:         "order_created",
        Projection:  "order.created",
        Description: "A new order should be created",
    }},
)
if err != nil {
    log.Fatal(err)
}
fmt.Println("published decision:", decision.Choice)

Providers use JSONCodec by default. Supply a custom DecodePayload when consumers need concrete payload types rather than generic JSON values.

The low-level provider API remains available for custom lifecycle management:

provider, err := GoEventBus.NewRedisProvider(redisConfig)
if err != nil {
    log.Fatal(err)
}

store := GoEventBus.NewEventStore(
    &dispatcher,
    1024,
    GoEventBus.Block,
    GoEventBus.WithProvider(provider),
)

See examples/redis and examples/rabbitmq for runnable examples.


Dead Letter Queue

Attach a DLQ to retain handler failures and recovered panics.

store.DLQ = GoEventBus.NewDeadLetterQueue()

Inspect failures:

for _, dead := range store.DLQ.Entries() {
    log.Printf(
        "event=%s attempts=%d err=%v",
        dead.Event.ID,
        dead.Attempts,
        dead.Err,
    )
}

Replay:

if err := store.DLQ.Replay(ctx, store); err != nil {
    log.Println("replay failed:", err)
}

Handler panics are recovered and converted into errors so a single handler cannot kill the worker pool.


Transactions

Transactions buffer events locally until commit, but commit delivery follows the parent EventStore transport. They are not database-style atomic transactions.

Known events can still be buffered directly:

tx := store.BeginTransaction()

tx.Publish(GoEventBus.Event{
    ID:         "evt-1",
    Projection: "order.created",
    Data:       order,
})

tx.Publish(GoEventBus.Event{
    ID:         "evt-2",
    Projection: "invoice.created",
    Data:       invoice,
})

if err := tx.Commit(ctx); err != nil {
    tx.Rollback()
    log.Fatal(err)
}

Intelligent transactions

Transactions can also use the same EventSelector pipeline as DecideAndSubscribe. Use DecideAndPublish when the event type must be chosen by deterministic rules, cache, Jev, or another selector:

tx := store.BeginTransaction()

decision, err := tx.DecideAndPublish(
    ctx,
    selector,
    map[string]any{
        "message": "Create an invoice for order 42",
    },
    GoEventBus.Event{
        ID:   "evt-3",
        Data: invoice,
    },
    []GoEventBus.EventCandidate{
        {
            Key:         "order_created",
            Projection:  "order.created",
            Description: "A new order should be created",
        },
        {
            Key:         "invoice_created",
            Projection:  "invoice.created",
            Description: "An invoice should be created",
        },
    },
)
if err != nil {
    tx.Rollback()
    log.Fatal(err)
}

fmt.Println("selected:", decision.Choice)

// No handler has run yet.
if err := tx.Commit(ctx); err != nil {
    tx.Rollback()
    log.Fatal(err)
}

DecideAndPublish performs the decision immediately and buffers only the resolved event. This means:

Transaction commit follows the parent EventStore transport. With a normal local store, Commit executes buffered regular handlers synchronously. With WithRedis(...), WithRabbitMQ(...), or WithProvider(...), Commit publishes buffered events to the configured provider instead; local handlers run only when those events are consumed back into an EventStore.

Local Commit stops at the first handler error. Broker-backed Commit stops at the first publish error and removes only the prefix already confirmed by the provider, so retrying starts from the first unpublished event. Transactions are still not database-style atomic transactions: handler side effects or broker publishes that already succeeded cannot be undone.

Rollback only discards events still buffered locally and never rewinds the shared EventStore ring.


Scheduling

Schedule a known event at a specific time:

timer := store.Schedule(
    ctx,
    time.Now().Add(10*time.Second),
    event,
)

Or after a duration:

timer := store.ScheduleAfter(
    ctx,
    30*time.Second,
    event,
)

Intelligent scheduling

Use the same Rules / Cache / Jev selector pipeline when the event type must be chosen before it is scheduled.

Schedule at an absolute time:

decision, timer, err := store.DecideAndSchedule(
    ctx,
    time.Now().Add(10*time.Second),
    selector,
    state,
    GoEventBus.Event{
        ID:   "evt-scheduled-1",
        Data: payload,
    },
    candidates,
)
if err != nil {
    log.Fatal(err)
}

fmt.Println("scheduled:", decision.Choice)

Or after a duration:

decision, timer, err := store.DecideAndScheduleAfter(
    ctx,
    30*time.Second,
    selector,
    state,
    GoEventBus.Event{
        ID:   "evt-scheduled-2",
        Data: payload,
    },
    candidates,
)
if err != nil {
    log.Fatal(err)
}

The selector is evaluated when scheduling is requested, not when the timer fires:

state
  ↓
Rules
  ↓
Cache
  ↓
Jev
  ↓
selected event
  ↓
timer
  ↓
Subscribe → Publish → handler

That means selection errors are returned immediately and no timer is created. Once scheduling succeeds, the timer contains an already-resolved event, so a later change in rules, cache contents, or Jev output does not change what will fire.

DecideAndSchedule and DecideAndScheduleAfter work with RuleCacheSelector, JevSelector, or any custom EventSelector.

The returned *time.Timer may be stopped before it fires. Past times and non-positive durations preserve the existing immediate-execution behavior and return a nil timer.

Scheduling remains a local EventStore operation. It does not publish the scheduled event to Redis/RabbitMQ even when the parent store has a broker configured.


Metrics

published, processed, failures := store.Metrics()

fmt.Printf(
    "published=%d processed=%d errors=%d\n",
    published,
    processed,
    failures,
)

These counters cover event execution. Decision-layer metrics such as cache hit rate or Jev latency are intentionally not baked into the core API yet.


API Snapshot

Event routing

type EventCandidate struct {
    Key         string
    Projection  interface{}
    Description string
}

type EventDecision struct {
    Choice        string
    Projection    interface{}
    Confidence    float64
    Probabilities map[string]float64
    Model         string
    RequestID     string
}

type EventSelector interface {
    SelectEvent(
        context.Context,
        any,
        []EventCandidate,
    ) (EventDecision, error)
}

Main implementations:

EventStore

func NewEventStore(
    dispatcher *Dispatcher,
    bufferSize uint64,
    policy OverrunPolicy,
    options ...EventStoreOption,
) *EventStore

bufferSize must be a non-zero power of two.

Important methods:

Method Purpose
Subscribe Enqueue a known event
DecideAndSubscribe Select a candidate and enqueue locally or publish to the configured broker
Publish Dispatch pending events
RegisterOrdered Preserve FIFO per ordering key
RegisterBatch Process projection events in chunks
Consume Feed configured or explicit provider events into the local store
PublishToProvider Publish an event through the provider configured on NewEventStore
Use Register middleware
OnBefore / OnAfter / OnError Register lifecycle hooks
Metrics Read event counters
Schedule / ScheduleAfter Schedule known events
DecideAndSchedule / DecideAndScheduleAfter Select with Rules/Cache/Jev, then schedule the resolved event
BeginTransaction Buffer events for commit; local stores execute handlers, provider-backed stores publish to the broker
Drain / Close Stop accepting events and finish in-flight work

Benchmarks

Repository benchmarks on Apple M-series:

go test -bench . -benchtime=3s
Benchmark Iterations ns/op
BenchmarkSubscribe 27,080,376 40.37
BenchmarkSubscribeParallel 26,418,999 38.42
BenchmarkPublish 295,661,464 3.91
BenchmarkPublishAfterPrefill 252,943,526 4.59
BenchmarkSubscribeLargePayload 1,613,017 771.5
BenchmarkPublishLargePayload 296,434,225 3.91
BenchmarkEventStore_Async 2,816,988 436.5
BenchmarkEventStore_Sync 2,638,519 428.5
BenchmarkFastHTTPSync 6,275,112 163.8
BenchmarkFastHTTPAsync 1,954,884 662.0
BenchmarkFastHTTPParallel 4,489,274 262.3

The intelligent routing path is intentionally outside these core dispatch benchmarks because rules, cache, and Jev have very different latency profiles.


Design Principles

GoEventBus keeps decision-making and execution separate.

Decision layer                    Execution layer

rules ─┐
       ├─► choice ───────────────► EventStore
cache ─┤                           ├─ ring buffer
       │                           ├─ back-pressure
Jev ───┘                           ├─ worker pool
                                   ├─ ordering
                                   ├─ batching
                                   ├─ middleware/hooks
                                   └─ DLQ

This separation means:


Contributing

Issues and pull requests are welcome.

Run the test suite with:

go test -race ./...

Integration tests:

go test -race -tags=integration -timeout=5m ./...

Quality checks used by CI include go vet and staticcheck.


License

Distributed under the MIT License. See LICENSE.

Join libs.tech

...and unlock some superpowers

GitHub

We won't share your data with anyone else.