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.
┌──────────────┐
│ 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:
- Known event type → enqueue it directly.
- Obvious rule → route deterministically.
- Repeated decision → reuse the cached choice.
- Ambiguous input → ask Jev to choose from a bounded candidate set.
- Execution → always stays inside GoEventBus.
That gives you intelligent routing without putting an LLM call in the Publish() hot path.
Features
- Bounded MPMC ring buffer built with atomics and cache-line padding.
- Rules → cache → Jev routing for optional intelligent event selection.
- Typed or string projections for local dispatch.
- Fan-out handlers with independent execution.
- Ordered handlers with FIFO delivery per ordering key.
- Batch handlers for bulk writes and high-throughput pipelines.
- Sync and async dispatch with a fixed worker pool.
- Middleware for per-handler cross-cutting behavior.
- Lifecycle hooks with
OnBefore,OnAfter, andOnError. - Back-pressure policies:
DropOldest,Block, orReturnError. - Dead-letter queue with replay support and panic recovery.
- Redis Streams and RabbitMQ EventStore options for cross-process delivery.
- Transactions with local buffering, transport-aware commit, and optional Rules/Cache/Jev event selection.
- Scheduling with absolute/relative timers and optional Rules/Cache/Jev event selection.
- Metrics for published, processed, and failed events.
- Pluggable decision cache through the
DecisionCacheinterface.
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:
- the JSON-serializable decision state,
- candidate keys,
- candidate descriptions.
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:
DecideAndSubscribe(...)selects a candidate and publishes it to the configured broker.PublishToProvider(ctx, event)is the direct path when the projection is already known and no decision step is needed.Consume(ctx)consumes from the configured broker and feeds events through the localDispatcher, middleware, hooks, ordering, batching, and DLQ path.
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:
- routing failures do not add an event to the transaction,
- the caller receives
EventDecisionbefore commit, - handler side effects remain deferred until
Commit, Rollbackdiscards intelligently selected events exactly like direct events,RuleCacheSelector,JevSelector, and customEventSelectorimplementations all use the same transaction API.
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:
JevSelectorRuleCacheSelectorMemoryDecisionCache
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:
- the event bus does not depend on Jev,
- deterministic workloads do not pay model latency,
- model failures do not change core dispatch semantics,
- alternative selectors can implement the same
EventSelectorinterface, - alternative caches can implement
DecisionCache, - application code can mix direct and intelligent routing in the same store.
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.