Skip to content

Latest commit

 

History

263 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 

Repository files navigation

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:

  • 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, and OnError.
  • Back-pressure policies: DropOldest, Block, or ReturnError.
  • 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 DecisionCache interface.

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 local Dispatcher, 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)

NATS JetStream

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.

Apache Kafka

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.

Remaining official brokers

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.

AWS SQS

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.

AWS SNS

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.

Google Cloud Pub/Sub

GoEventBus.WithGCPPubSub(GoEventBus.GCPPubSubProviderConfig{
    Client:         pubsubClient,
    Topic:          "events",
    Subscription:   "billing",
    EnableOrdering: true,
})

Azure Service Bus

GoEventBus.WithAzureServiceBus(GoEventBus.AzureServiceBusProviderConfig{
    Client: sbClient,
    Queue:  "billing",
})

For topic subscriptions, set Topic and Subscription instead of Queue.

Apache Pulsar

GoEventBus.WithPulsar(GoEventBus.PulsarProviderConfig{
    Client:       pulsarClient,
    Topic:        "persistent://public/default/events",
    Subscription: "billing",
})

MQTT

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.

PostgreSQL

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.

NSQ

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.


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 EventDecision before commit,
  • handler side effects remain deferred until Commit,
  • Rollback discards intelligently selected events exactly like direct events,
  • RuleCacheSelector, JevSelector, and custom EventSelector implementations 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:

  • JevSelector
  • RuleCacheSelector
  • MemoryDecisionCache

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 EventSelector interface,
  • 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.

Releases

Packages

Used by

Contributors

Languages