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.
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.
- 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.
- 12 official EventStore broker options: Redis Streams, RabbitMQ, NATS JetStream, Kafka, SQS, SNS, GCP Pub/Sub, Azure Service Bus, Pulsar, MQTT, PostgreSQL, and NSQ.
- 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.
go get github.com/Protocol-Lattice/GoEventBusGoEventBus currently targets Go 1.23+.
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.
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.
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.
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.
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.
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.
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.
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: truewhen cache failures should fail the routing call instead.
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_jevSee examples/README.md for a compact guide to all examples.
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.
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.
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 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 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.
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.
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)
}Enable worker-pool dispatch with:
store.Async = trueThe 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.
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.
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)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)nc, err := nats.Connect(nats.DefaultURL)
if err != nil {
log.Fatal(err)
}
store := GoEventBus.NewEventStore(
&dispatcher,
1024,
GoEventBus.Block,
GoEventBus.WithNATSJetStream(GoEventBus.NATSJetStreamProviderConfig{
Conn: nc,
SubjectPrefix: "events.",
Subject: "events.>",
Durable: "billing",
Queue: "billing-workers",
}),
)The injected NATS connection remains application-owned. JetStream messages use explicit acknowledgements; handler failures are NAKed for broker redelivery.
store := GoEventBus.NewEventStore(
&dispatcher,
1024,
GoEventBus.Block,
GoEventBus.WithKafka(GoEventBus.KafkaProviderConfig{
Brokers: []string{"localhost:9092"},
Topic: "events",
Group: "billing",
}),
)Kafka consumer offsets are committed synchronously only after successful event handling. With the default hash balancer, the string projection is used as the message key, preserving ordering for a projection within a partition.
All remaining official transports use the same store-first API:
store := GoEventBus.NewEventStore(
&dispatcher,
1024,
GoEventBus.Block,
GoEventBus.WithSQS(sqsConfig), // swap for any provider below
)| Broker | EventStore option | Consume semantics |
|---|---|---|
| AWS SQS | WithSQS |
long-poll; delete only after successful handling |
| AWS SNS | WithSNS |
publish through SNS; consumption delegates to a real subscriber transport such as SQS |
| Google Cloud Pub/Sub | WithGCPPubSub |
Ack on success, Nack on decode/handler failure |
| Azure Service Bus | WithAzureServiceBus |
PeekLock; Complete on success, Abandon on failure |
| Apache Pulsar | WithPulsar |
durable subscription; Ack on success, Nack on failure |
| MQTT | WithMQTT |
QoS 0/1/2; broker ACK behavior follows MQTT/Paho semantics |
| PostgreSQL | WithPostgres |
persistent table + FOR UPDATE SKIP LOCKED; delete on success, delayed retry on failure |
| NSQ | WithNSQ |
channel-based consumption; handler errors are requeued by go-nsq |
Broker-backed events still require string projections, custom EventCodec
implementations are supported by every provider, and Drain / Close
closes provider-owned resources. Injected cloud clients/connections remain
application-owned.
store := GoEventBus.NewEventStore(
&dispatcher, 1024, GoEventBus.Block,
GoEventBus.WithSQS(GoEventBus.SQSProviderConfig{
Client: sqsClient,
QueueURL: queueURL,
}),
)FIFO queues can set FIFO, MessageGroupID, and an optional
MessageDeduplicationID callback.
sqsSubscriber, _ := GoEventBus.NewSQSProvider(sqsConfig)
store := GoEventBus.NewEventStore(
&dispatcher, 1024, GoEventBus.Block,
GoEventBus.WithSNS(GoEventBus.SNSProviderConfig{
Client: snsClient,
TopicARN: topicARN,
Subscriber: sqsSubscriber,
}),
)SNS is not modeled as a queue. Consume requires an actual subscriber
transport such as SQS, matching SNS delivery semantics.
GoEventBus.WithGCPPubSub(GoEventBus.GCPPubSubProviderConfig{
Client: pubsubClient,
Topic: "events",
Subscription: "billing",
EnableOrdering: true,
})GoEventBus.WithAzureServiceBus(GoEventBus.AzureServiceBusProviderConfig{
Client: sbClient,
Queue: "billing",
})For topic subscriptions, set Topic and Subscription instead of Queue.
GoEventBus.WithPulsar(GoEventBus.PulsarProviderConfig{
Client: pulsarClient,
Topic: "persistent://public/default/events",
Subscription: "billing",
})GoEventBus.WithMQTT(GoEventBus.MQTTProviderConfig{
Options: mqttOptions,
TopicPrefix: "events/",
QoS: 1,
})An injected MQTT client remains caller-owned. When Options is supplied,
GoEventBus creates and owns the client. Automatic reconnect and persistent
sessions are configured through Paho's ClientOptions.
GoEventBus.WithPostgres(GoEventBus.PostgresProviderConfig{
ConnectionString: os.Getenv("DATABASE_URL"),
Table: "goeventbus_events",
})The PostgreSQL provider uses a durable table-backed queue. Consumers claim one
row with FOR UPDATE SKIP LOCKED; successful processing deletes the row,
while failures increment the attempt count and delay the next delivery.
GoEventBus.WithNSQ(GoEventBus.NSQProviderConfig{
NSQDAddress: "127.0.0.1:4150",
Topic: "events",
Channel: "billing",
})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.
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 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)
}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.
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,
)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.
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.
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
func NewEventStore(
dispatcher *Dispatcher,
bufferSize uint64,
policy OverrunPolicy,
options ...EventStoreOption,
) *EventStorebufferSize 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 |
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.
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.
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.
Distributed under the MIT License. See LICENSE.
