go-saga-axon

module
v1.0.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Jul 11, 2026 License: MIT

README ΒΆ

πŸ”± saga-axon

A production-grade DAG-based saga orchestration engine for Go. Coordinates distributed transactions across microservices with automatic compensation, timeout monitoring, CQRS event sourcing, and split idempotency β€” built to never lose a transaction.

Go Redis RabbitMQ Kafka License


How It Works

                            β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
                            β”‚                    saga-axon engine                       β”‚
                            β”‚                                                          β”‚
  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”   START      β”‚  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”    β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”    β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”   β”‚
  β”‚ CDC-Axon │─────────────▢│  β”‚ Listener  │───▢│Lifecycle │───▢│ Command         β”‚   β”‚
  β”‚ (outbox) β”‚              β”‚  β”‚ (RabbitMQ)β”‚    β”‚ Engine   β”‚    β”‚ Dispatcher      β”‚   β”‚
  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜              β”‚  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜    β””β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”˜    β””β”€β”€β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”€β”€β”˜   β”‚
                            β”‚                        β”‚                    β”‚            β”‚
  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”   SUCCESS    β”‚                   β”Œβ”€β”€β”€β”€β–Όβ”€β”€β”€β”€β”€β”             β”‚            β”‚
  β”‚ Worker   │─────────────▢│                   β”‚  Redis   β”‚             β”‚            β”‚
  β”‚ Service  β”‚              β”‚                   β”‚  State   β”‚             β”‚            β”‚
  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜              β”‚                   β””β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”˜             β”‚            β”‚
                            β”‚                        β”‚                    β”‚            β”‚
  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”   FAILURE    β”‚                   β”Œβ”€β”€β”€β”€β–Όβ”€β”€β”€β”€β”€β”    β”Œβ”€β”€β”€β”€β”€β”€β”€β–Όβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”  β”‚
  β”‚ Worker   │─────────────▢│                   β”‚  Kafka   β”‚    β”‚ RabbitMQ/SQS/    β”‚  β”‚
  β”‚ Service  β”‚              β”‚                   β”‚  CQRS    β”‚    β”‚ Any Dispatcher   β”‚  β”‚
  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜              β”‚                   β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜    β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜  β”‚
                            β”‚                                                          β”‚
                            β”‚  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”    β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”    β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”   β”‚
                            β”‚  β”‚ Timeout   │───▢│Compensate│───▢│ Alerter         β”‚   β”‚
                            β”‚  β”‚ Monitor   β”‚    β”‚ Engine   β”‚    β”‚ (PD/Slack/CW)   β”‚   β”‚
                            β”‚  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜    β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜    β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜   β”‚
                            β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜

What is saga-axon?

In distributed systems, a single business operation (e.g. "place order") spans multiple services β€” payments, inventory, shipping, notifications. If any step fails, you need to undo everything that already succeeded. That's the saga pattern.

saga-axon is the orchestrator that:

  1. Models your workflow as a DAG β€” steps can run in parallel (fan-out) and converge (fan-in)
  2. Dispatches commands to worker services via any message broker
  3. Tracks state in Redis with per-workflow locking
  4. Compensates automatically β€” on failure, ALL succeeded steps get rollback commands
  5. Handles timeouts β€” steps that take too long are treated as failures
  6. Survives crashes β€” CQRS event log enables state rebuild from Kafka
  7. Deduplicates events β€” split idempotency ensures exactly-once processing semantics

Key Features

Feature Description
DAG Workflows Steps with fan-out, fan-in, and arbitrary dependency graphs
Automatic Compensation On failure, all succeeded steps receive compensation commands
Timeout Monitoring Per-step max duration with jittered polling
Split Idempotency Read-only check β†’ process β†’ mark. No stuck sagas on dispatch failure
CQRS Event Sourcing Every event logged to Kafka. State rebuildable from event log
Buffered Failover If Kafka is down, events buffer to Redis and flush when available
Late SUCCESS Handling If a step reports SUCCESS during rollback, it gets compensated immediately
Pluggable Everything Bring your own broker, state store, event logger, alerter
Structured Logging JSON in production, text in dev β€” ready for Loki/CloudWatch/Datadog
Graceful Shutdown WaitGroup-based drain of in-flight events

Installation

go get github.com/Srajan-Sanjay-Saxena/go-saga-axon

Core Concepts

The Envelope Protocol

Every message flowing through saga-axon follows a strict envelope format:

type Envelope struct {
    WorkflowID  string             `json:"workflow_id"`
    MessageType EnvelopMessageType `json:"message_type"` // START | COMMAND | SUCCESS | FAILURE | COMPENSATE
    StepName    string             `json:"step_name"`
    Payload     json.RawMessage    `json:"payload"`
    Timestamp   time.Time          `json:"timestamp"`
}
Message Flow
 CDC/Trigger                    saga-axon                     Worker Services
─────────────────────────────────────────────────────────────────────────────
     β”‚                              β”‚                              β”‚
     │──── START ──────────────────▢│                              β”‚
     β”‚                              │──── COMMAND (step_a) ───────▢│
     β”‚                              │──── COMMAND (step_b) ───────▢│  (parallel roots)
     β”‚                              β”‚                              β”‚
     β”‚                              │◀─── SUCCESS (step_a) ────────│
     β”‚                              │◀─── SUCCESS (step_b) ────────│
     β”‚                              β”‚                              β”‚
     β”‚                              │──── COMMAND (step_c) ───────▢│  (fan-in: depends on a+b)
     β”‚                              β”‚                              β”‚
     β”‚                              │◀─── FAILURE (step_c) ────────│
     β”‚                              β”‚                              β”‚
     β”‚                              │──── COMPENSATE (step_a) ────▢│  (rollback all succeeded)
     β”‚                              │──── COMPENSATE (step_b) ────▢│
     β”‚                              β”‚                              β”‚
Workflow State Machine
  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”     all steps done     β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
  β”‚ RUNNING │────────────────────────▢│ COMPLETED β”‚
  β””β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”˜                         β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
       β”‚
       β”‚ any step fails
       β–Ό
  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”   compensation fails   β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”
  β”‚ROLLING_BACK│────────────────────────▢│ FAILED β”‚
  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜                         β””β”€β”€β”€β”€β”€β”€β”€β”€β”˜

Defining Workflows (DAG)

package main

import (
    "time"
    "go-saga-axon/saga/route"
)

func OrderWorkflow(dispatcher contract.Dispatcher) *route.Workflow {
    w := route.NewWorkflow("order-placement")

    // Root steps β€” dispatched immediately on START (run in parallel)
    reserveInventory := w.AddStep("reserve_inventory").
        Command(dispatcher, "inventory.reserve").
        Compensation(dispatcher, "inventory.release").
        Timeout(30 * time.Second)

    chargePayment := w.AddStep("charge_payment").
        Command(dispatcher, "payment.charge").
        Compensation(dispatcher, "payment.refund").
        Timeout(15 * time.Second)

    // Fan-in β€” only dispatched when BOTH parents succeed
    shipOrder := w.AddStep("ship_order").
        Command(dispatcher, "shipping.create").
        Compensation(dispatcher, "shipping.cancel").
        Timeout(60 * time.Second).
        DependsOn(reserveInventory, chargePayment)

    // Leaf step β€” no compensation needed (fire-and-forget notification)
    _ = w.AddStep("send_confirmation").
        Command(dispatcher, "notification.send").
        Timeout(10 * time.Second).
        DependsOn(shipOrder)

    // Monitor every 20 seconds for timed-out steps
    w.SetMonitorInterval(20 * time.Second)

    return w
}

This produces the following DAG:

  reserve_inventory ──┐
                      β”œβ”€β”€β–Ά ship_order ──▢ send_confirmation
  charge_payment β”€β”€β”€β”€β”€β”˜

If ship_order fails β†’ reserve_inventory and charge_payment both get compensated automatically.


Full Production Scaffold

RabbitMQ (via goRabbit-axon) + Kafka (CQRS event log) + Redis (state store)

package main

import (
    "context"
    "os"
    "os/signal"
    "syscall"
    "time"

    "github.com/IBM/sarama"
    "github.com/redis/go-redis/v9"

    "github.com/Srajan-Sanjay-Saxena/goRabbit-axon/v2/breaker"
    connPool "github.com/Srajan-Sanjay-Saxena/goRabbit-axon/v2/connection/connectionPool"
    singleConn "github.com/Srajan-Sanjay-Saxena/goRabbit-axon/v2/connection/singleConnection"
    "github.com/Srajan-Sanjay-Saxena/goRabbit-axon/v2/exchange"

    "go-saga-axon/contract"
    "go-saga-axon/core"
    "go-saga-axon/engine"
    "go-saga-axon/handlers"
    "go-saga-axon/saga/route"
    "go-saga-axon/store"
)

func main() {
    ctx, cancel := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
    defer cancel()

    // ─── 1. Redis (State Store + Idempotency + Tracking) ───────────────────────
    redisClient := redis.NewClient(&redis.Options{
        Addr:         "localhost:6379",
        Password:     os.Getenv("REDIS_PASSWORD"),
        DB:           0,
        PoolSize:     20,
        MinIdleConns: 5,
    })
    stateStore := NewRedisStateStore(redisClient)

    // ─── 2. RabbitMQ (Command Dispatch + Event Listening) ──────────────────────
    rabbitConn := connPool.NewConnectionPool(
        os.Getenv("RABBITMQ_URL"), // amqp://user:pass@localhost:5672/
        connPool.PoolOptions{ConnSize: 3, ChanPerConn: 5},
        singleConn.DefaultOptions(),
        nil,
    )
    rabbitConn.AddBreaker(breaker.CircuitBreakerOptions{
        Threshold:   5,
        Timeout:     30 * time.Second,
    })
    if err := rabbitConn.Connect(ctx); err != nil {
        panic("rabbitmq connect: " + err.Error())
    }
    defer rabbitConn.Shutdown()

    // Declare exchange + queues
    declareInfrastructure(ctx, rabbitConn)

    // ─── 3. Kafka (CQRS Event Logger) ─────────────────────────────────────────
    kafkaLogger, err := NewKafkaEventLogger(
        []string{"localhost:9092"},
        "saga-events",
    )
    if err != nil {
        panic("kafka connect: " + err.Error())
    }
    defer kafkaLogger.Close()

    // ─── 4. Alerter (PagerDuty / Slack / CloudWatch) ──────────────────────────
    alerter := NewSlackAlerter(os.Getenv("SLACK_WEBHOOK_URL"))

    // ─── 5. Logger ────────────────────────────────────────────────────────────
    logger := core.NewDebugLogger(core.Production) // JSON output for Loki/CloudWatch

    // ─── 6. Persistence Layer ─────────────────────────────────────────────────
    persistence := store.NewPersistence(stateStore).
        WithCQRS(kafkaLogger).
        WithAlerter(alerter).
        WithLogger(logger)

    // ─── 7. Command Dispatcher (with exponential backoff) ─────────────────────
    dispatcher := NewRabbitDispatcher(rabbitConn)
    cmdDispatcher := handlers.NewCommandDispatcher(5). // max 5 retries
        WithAlerter(alerter).
        WithLogger(logger)

    // ─── 8. Define Workflow ───────────────────────────────────────────────────
    workflow := OrderWorkflow(dispatcher)

    // ─── 9. Wire Engine ───────────────────────────────────────────────────────
    listener := NewRabbitListener(rabbitConn, "saga.inbound")

    tg := core.NewTriggerGroup("order-triggers").
        AddListener(listener, workflow)

    eng := engine.NewEngine().
        WithPersistence(persistence).
        WithCommandDispatcher(cmdDispatcher).
        WithDebugLogger(logger).
        WithAlerter(alerter).
        AddTriggerGroup(tg)

    // ─── 10. Start ────────────────────────────────────────────────────────────
    eng.Start(ctx)

    // ─── 11. Graceful Shutdown ────────────────────────────────────────────────
    shutdownCtx, shutdownCancel := context.WithTimeout(context.Background(), 30*time.Second)
    defer shutdownCancel()
    eng.Shutdown(shutdownCtx)
}

func declareInfrastructure(ctx context.Context, conn helpers.IRabbitConnection) {
    ex := exchange.NewRabbitExchange(exchange.RabbitExchangeConfig{
        Name:    "saga.exchange",
        Type:    exchange.Topic,
        Durable: true,
    })
    ex.CreateExchange(ctx, conn)

    queues := []struct{ name, key string }{
        {"saga.inbound", "saga.events.#"},
        {"inventory.commands", "inventory.#"},
        {"payment.commands", "payment.#"},
        {"shipping.commands", "shipping.#"},
        {"notification.commands", "notification.#"},
    }
    for _, q := range queues {
        ex.CreateQueue(ctx, conn, exchange.RabbitQueueConfig{
            Name:       q.name,
            BindingKey: q.key,
            Durable:    true,
        })
    }
}

Implementing Adapters

saga-axon is fully pluggable. You implement 4 interfaces and bring any infrastructure:

StateStore (Redis)
type RedisStateStore struct {
    client *redis.Client
}

func NewRedisStateStore(client *redis.Client) *RedisStateStore {
    return &RedisStateStore{client: client}
}

func stateKey(workflowID string) string {
    return "state:{" + workflowID + "}"
}

func (r *RedisStateStore) SaveState(ctx context.Context, workflowID string, state *data.WorkflowState) error {
    b, _ := json.Marshal(state)
    return r.client.Set(ctx, stateKey(workflowID), b, 0).Err()
}

func (r *RedisStateStore) LoadState(ctx context.Context, workflowID string) (*data.WorkflowState, error) {
    val, err := r.client.Get(ctx, stateKey(workflowID)).Result()
    if err == redis.Nil {
        return nil, &core.EngineError{Kind: core.StateNotFound, Message: "not found"}
    }
    if err != nil {
        return nil, err
    }
    var state data.WorkflowState
    json.Unmarshal([]byte(val), &state)
    return &state, nil
}

func (r *RedisStateStore) LPush(ctx context.Context, key, value string) error {
    return r.client.LPush(ctx, key, value).Err()
}

func (r *RedisStateStore) LPop(ctx context.Context, key string) (string, error) {
    val, err := r.client.LPop(ctx, key).Result()
    if err == redis.Nil { return "", nil }
    return val, err
}

func (r *RedisStateStore) SAdd(ctx context.Context, key, member string) error {
    return r.client.SAdd(ctx, key, member).Err()
}

func (r *RedisStateStore) SMembers(ctx context.Context, key string) ([]string, error) {
    return r.client.SMembers(ctx, key).Result()
}

func (r *RedisStateStore) SRem(ctx context.Context, key, member string) error {
    return r.client.SRem(ctx, key, member).Err()
}

func (r *RedisStateStore) Eval(ctx context.Context, script string, keys []string, args ...any) error {
    return r.client.Eval(ctx, script, keys, args...).Err()
}

func (r *RedisStateStore) EvalInt(ctx context.Context, script string, keys []string, args ...any) (int, error) {
    return r.client.Eval(ctx, script, keys, args...).Int()
}
Dispatcher (RabbitMQ via goRabbit-axon)
type RabbitDispatcher struct {
    conn helpers.IRabbitConnection
}

func NewRabbitDispatcher(conn helpers.IRabbitConnection) *RabbitDispatcher {
    return &RabbitDispatcher{conn: conn}
}

func (d *RabbitDispatcher) Dispatch(ctx context.Context, route string, payload []byte) error {
    pub := producer.NewProducer(producer.ProducerConfig{
        ExchangeName: "saga.exchange",
        RoutingKey:   route,
    })
    if err := pub.GetChannel(ctx, d.conn, nil); err != nil {
        return err
    }
    return pub.Publish(ctx, payload, producer.RabbitMqPublisherConfig{Persistent: true})
}
Listener (RabbitMQ via goRabbit-axon)
type RabbitListener struct {
    conn      helpers.IRabbitConnection
    queueName string
}

func NewRabbitListener(conn helpers.IRabbitConnection, queueName string) *RabbitListener {
    return &RabbitListener{conn: conn, queueName: queueName}
}

func (l *RabbitListener) Listen(ctx context.Context, handler func([]byte) error, logger contract.Logger) error {
    cons := consumer.NewConsumer(consumer.ConsumerConfig{
        QueueName: l.queueName,
        Prefetch:  10,
        AutoAck:   false,
        Handler: func(ctx context.Context, msg amqp.Delivery) error {
            if err := handler(msg.Body); err != nil {
                // Check error type for routing decision
                var engineErr *core.EngineError
                if errors.As(err, &engineErr) && engineErr.Kind == core.PoisonPill {
                    msg.Ack(false) // discard poison pills
                    return nil
                }
                return err // nack + requeue for retryable errors
            }
            return nil
        },
    })
    if err := cons.GetChannel(ctx, l.conn); err != nil {
        return err
    }
    return cons.Consume(ctx)
}
EventLogger (Kafka)
type KafkaEventLogger struct {
    producer sarama.SyncProducer
    consumer sarama.Consumer
    topic    string
}

func NewKafkaEventLogger(brokers []string, topic string) (*KafkaEventLogger, error) {
    config := sarama.NewConfig()
    config.Producer.Return.Successes = true
    config.Producer.RequiredAcks = sarama.WaitForAll

    prod, err := sarama.NewSyncProducer(brokers, config)
    if err != nil {
        return nil, err
    }
    cons, err := sarama.NewConsumer(brokers, config)
    if err != nil {
        prod.Close()
        return nil, err
    }
    return &KafkaEventLogger{producer: prod, consumer: cons, topic: topic}, nil
}

func (k *KafkaEventLogger) LogEvent(ctx context.Context, workflowID string, event *data.Envelope) error {
    b, _ := json.Marshal(event)
    _, _, err := k.producer.SendMessage(&sarama.ProducerMessage{
        Topic: k.topic,
        Key:   sarama.StringEncoder(workflowID),
        Value: sarama.ByteEncoder(b),
    })
    return err
}

func (k *KafkaEventLogger) LoadEvents(ctx context.Context, workflowID string) ([]*data.Envelope, error) {
    partitions, _ := k.consumer.Partitions(k.topic)
    var events []*data.Envelope
    for _, p := range partitions {
        pc, _ := k.consumer.ConsumePartition(k.topic, p, sarama.OffsetOldest)
        timeout := time.After(3 * time.Second)
        for {
            select {
            case msg := <-pc.Messages():
                if string(msg.Key) == workflowID {
                    var env data.Envelope
                    json.Unmarshal(msg.Value, &env)
                    events = append(events, &env)
                }
            case <-timeout:
                pc.Close()
                goto next
            }
        }
    next:
    }
    return events, nil
}

Alternative Adapters

Dispatcher: AWS SQS
type SQSDispatcher struct {
    client *sqs.Client
    queues map[string]string // route β†’ queue URL mapping
}

func NewSQSDispatcher(client *sqs.Client, queues map[string]string) *SQSDispatcher {
    return &SQSDispatcher{client: client, queues: queues}
}

func (d *SQSDispatcher) Dispatch(ctx context.Context, route string, payload []byte) error {
    queueURL, ok := d.queues[route]
    if !ok {
        return fmt.Errorf("no queue mapped for route: %s", route)
    }
    _, err := d.client.SendMessage(ctx, &sqs.SendMessageInput{
        QueueUrl:    &queueURL,
        MessageBody: aws.String(string(payload)),
        MessageAttributes: map[string]types.MessageAttributeValue{
            "route": {DataType: aws.String("String"), StringValue: &route},
        },
    })
    return err
}
Dispatcher: Redis Streams (BullMQ-style)
type RedisStreamDispatcher struct {
    client *redis.Client
}

func NewRedisStreamDispatcher(client *redis.Client) *RedisStreamDispatcher {
    return &RedisStreamDispatcher{client: client}
}

func (d *RedisStreamDispatcher) Dispatch(ctx context.Context, route string, payload []byte) error {
    return d.client.XAdd(ctx, &redis.XAddArgs{
        Stream: "saga:" + route,
        Values: map[string]interface{}{"payload": string(payload)},
    }).Err()
}
EventLogger: DynamoDB (Serverless CQRS)
type DynamoEventLogger struct {
    client *dynamodb.Client
    table  string
}

func (d *DynamoEventLogger) LogEvent(ctx context.Context, workflowID string, event *data.Envelope) error {
    b, _ := json.Marshal(event)
    _, err := d.client.PutItem(ctx, &dynamodb.PutItemInput{
        TableName: &d.table,
        Item: map[string]ddbtypes.AttributeValue{
            "workflow_id": &ddbtypes.AttributeValueMemberS{Value: workflowID},
            "timestamp":  &ddbtypes.AttributeValueMemberS{Value: event.Timestamp.Format(time.RFC3339Nano)},
            "event":      &ddbtypes.AttributeValueMemberS{Value: string(b)},
        },
    })
    return err
}

func (d *DynamoEventLogger) LoadEvents(ctx context.Context, workflowID string) ([]*data.Envelope, error) {
    out, err := d.client.Query(ctx, &dynamodb.QueryInput{
        TableName:              &d.table,
        KeyConditionExpression: aws.String("workflow_id = :wid"),
        ExpressionAttributeValues: map[string]ddbtypes.AttributeValue{
            ":wid": &ddbtypes.AttributeValueMemberS{Value: workflowID},
        },
    })
    if err != nil {
        return nil, err
    }
    var events []*data.Envelope
    for _, item := range out.Items {
        var env data.Envelope
        json.Unmarshal([]byte(item["event"].(*ddbtypes.AttributeValueMemberS).Value), &env)
        events = append(events, &env)
    }
    return events, nil
}
Alerter: Slack Webhook
type SlackAlerter struct {
    webhookURL string
    client     *http.Client
}

func NewSlackAlerter(webhookURL string) *SlackAlerter {
    return &SlackAlerter{webhookURL: webhookURL, client: &http.Client{Timeout: 5 * time.Second}}
}

func (s *SlackAlerter) Alert(ctx context.Context, severity contract.AlertSeverity, err error, metadata map[string]any) error {
    emoji := map[contract.AlertSeverity]string{
        contract.Info: "ℹ️", contract.Warning: "⚠️", contract.Critical: "🚨",
    }
    text := fmt.Sprintf("%s *saga-axon alert*\n```%v```\nMetadata: %v", emoji[severity], err, metadata)
    body, _ := json.Marshal(map[string]string{"text": text})
    req, _ := http.NewRequestWithContext(ctx, "POST", s.webhookURL, bytes.NewReader(body))
    req.Header.Set("Content-Type", "application/json")
    _, reqErr := s.client.Do(req)
    return reqErr
}
Alerter: AWS CloudWatch
type CloudWatchAlerter struct {
    client    *cloudwatch.Client
    namespace string
}

func (a *CloudWatchAlerter) Alert(ctx context.Context, severity contract.AlertSeverity, err error, metadata map[string]any) error {
    _, putErr := a.client.PutMetricData(ctx, &cloudwatch.PutMetricDataInput{
        Namespace: &a.namespace,
        MetricData: []cwtypes.MetricDatum{{
            MetricName: aws.String("SagaAlert"),
            Value:      aws.Float64(1),
            Dimensions: []cwtypes.Dimension{
                {Name: aws.String("Severity"), Value: aws.String(fmt.Sprintf("%d", severity))},
                {Name: aws.String("Reason"), Value: aws.String(fmt.Sprintf("%v", metadata["reason"]))},
            },
        }},
    })
    return putErr
}

CDC Middleware (Outbox Pattern)

saga-axon integrates with cdc-axon for the outbox pattern. Raw CDC events (database row changes) get translated into START envelopes via middleware:

// CDC-Axon captures outbox INSERT β†’ publishes raw row to RabbitMQ
// saga-axon's CDCMiddleware translates it into a START envelope

cdcMiddleware := middleware.NewCDCMiddleware("order_id") // field containing workflow_id

tg := core.NewTriggerGroup("cdc-triggers").
    AddCDCListener(
        NewRabbitListener(rabbitConn, "cdc.outbox.events"),
        OrderWorkflow(dispatcher),
        cdcMiddleware,
    )

What happens:

  1. Your service writes to the outbox table: INSERT INTO outbox (order_id, event_type, payload) VALUES (...)
  2. CDC-Axon captures the WAL event and publishes the raw row to RabbitMQ
  3. saga-axon's CDC middleware extracts order_id as workflow_id, wraps it in a START envelope
  4. The engine processes it like any other START event

Raw CDC payload:

{"order_id": "ord-123", "event_type": "ORDER_CREATED", "payload": {"items": [...]}}

After middleware translation:

{"workflow_id": "ord-123", "message_type": "START", "payload": {"order_id": "ord-123", ...}}

Observability & Monitoring

Structured Logging

saga-axon uses Go's log/slog with three modes:

// Production β€” JSON to stdout, WARN+ level (for Loki, CloudWatch Logs, Datadog)
logger := core.NewDebugLogger(core.Production)

// Development β€” colored text to stdout, DEBUG level
logger := core.NewDebugLogger(core.Development)

// Testing β€” discards all output
logger := core.NewDebugLogger(core.Testing)
Production JSON Output

In production mode, every log line is a JSON object β€” ready for ingestion by any log aggregator:

{"time":"2025-01-15T03:17:42.123Z","level":"WARN","msg":"step timeout exceeded","workflow_id":"ord-abc123","step":"charge_payment"}
{"time":"2025-01-15T03:17:42.456Z","level":"WARN","msg":"saga failing, starting rollback","workflow_id":"ord-abc123","step":"charge_payment"}
{"time":"2025-01-15T03:17:42.789Z","level":"WARN","msg":"compensating step","workflow_id":"ord-abc123","step":"reserve_inventory","route":"inventory.release"}
{"time":"2025-01-15T03:17:43.012Z","level":"ERROR","msg":"compensation_dispatch_failed","workflow_id":"ord-abc123","step":"reserve_inventory"}
Grafana Loki Integration

Ship JSON logs from stdout to Loki via Promtail or Grafana Agent:

# promtail-config.yaml
scrape_configs:
  - job_name: saga-axon
    static_configs:
      - targets: [localhost]
        labels:
          job: saga-axon
          __path__: /var/log/saga-axon/*.log
    pipeline_stages:
      - json:
          expressions:
            level: level
            workflow_id: workflow_id
            step: step
      - labels:
          level:
          workflow_id:

Query in Grafana:

# All errors
{job="saga-axon"} | json | level="ERROR"

# Track a specific workflow
{job="saga-axon"} | json | workflow_id="ord-abc123"

# All compensation events
{job="saga-axon"} |= "compensating"

# Timeout alerts
{job="saga-axon"} |= "step timeout exceeded"
AWS CloudWatch Logs

For ECS/EKS deployments, stdout JSON is automatically captured by CloudWatch:

// CloudWatch Insights query
fields @timestamp, @message
| filter workflow_id = "ord-abc123"
| sort @timestamp asc

// Alert on compensation failures
fields @timestamp, @message
| filter @message like /compensation_dispatch_failed/
| stats count() as failures by bin(5m)
Datadog Log Pipeline
# datadog-agent config
logs:
  - type: file
    path: /var/log/saga-axon/*.log
    service: saga-axon
    source: go
    sourcecategory: saga

Datadog auto-parses JSON logs. Create monitors on:

  • @level:ERROR β†’ PagerDuty alert
  • @msg:"step timeout exceeded" β†’ Slack notification
  • @msg:"compensation_dispatch_failed" β†’ Critical incident
Alerter Integration

The alerter fires on critical engine events β€” not just logs. Use it for real-time incident response:

β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
β”‚                        Alert Routing                                  β”‚
β”œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€
β”‚                                                                      β”‚
β”‚  contract.Info     β†’  Slack #saga-info channel                       β”‚
β”‚  contract.Warning  β†’  Slack #saga-alerts + Datadog event             β”‚
β”‚  contract.Critical β†’  PagerDuty + Slack #incidents + CloudWatch      β”‚
β”‚                                                                      β”‚
β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜

Example composite alerter:

type CompositeAlerter struct {
    slack      *SlackAlerter
    pagerduty  *PagerDutyAlerter
    cloudwatch *CloudWatchAlerter
}

func (a *CompositeAlerter) Alert(ctx context.Context, severity contract.AlertSeverity, err error, metadata map[string]any) error {
    // Always send to Slack
    a.slack.Alert(ctx, severity, err, metadata)

    // CloudWatch metric for all severities
    a.cloudwatch.Alert(ctx, severity, err, metadata)

    // PagerDuty only for critical
    if severity == contract.Critical {
        a.pagerduty.Alert(ctx, severity, err, metadata)
    }
    return nil
}

Error Semantics

saga-axon uses typed errors to drive listener behavior:

Error Kind Listener Action When
Retryable Nack + requeue Redis down, broker unreachable, transient failures
PoisonPill Ack + discard Malformed payload, missing workflow_id, unknown step
StateNotFound Triggers CQRS rebuild State evicted from Redis, first load after crash
// Your listener should route errors:
func (l *RabbitListener) Listen(ctx context.Context, handler func([]byte) error, logger contract.Logger) error {
    cons := consumer.NewConsumer(consumer.ConsumerConfig{
        QueueName: l.queueName,
        Prefetch:  10,
        AutoAck:   false,
        Handler: func(ctx context.Context, msg amqp.Delivery) error {
            err := handler(msg.Body)
            if err == nil {
                msg.Ack(false)
                return nil
            }

            var engineErr *core.EngineError
            if errors.As(err, &engineErr) {
                switch engineErr.Kind {
                case core.PoisonPill:
                    msg.Ack(false) // discard β€” will never succeed
                    return nil
                case core.Retryable:
                    msg.Nack(false, true) // requeue β€” try again later
                    return nil
                }
            }
            msg.Nack(false, true) // unknown error β€” requeue
            return nil
        },
    })
    cons.GetChannel(ctx, l.conn)
    return cons.Consume(ctx)
}

Worker Service Example

Your microservices receive COMMAND envelopes and respond with SUCCESS or FAILURE:

package main

import (
    "context"
    "encoding/json"
    "time"

    "go-saga-axon/data"
)

// Worker listens on "inventory.reserve" queue, processes commands,
// and publishes SUCCESS/FAILURE back to the saga inbound queue.
func main() {
    ctx := context.Background()

    // ... setup rabbitConn, declare queues ...

    listener := NewRabbitListener(rabbitConn, "inventory.commands")
    dispatcher := NewRabbitDispatcher(rabbitConn)

    listener.Listen(ctx, func(payload []byte) error {
        var env data.Envelope
        json.Unmarshal(payload, &env)

        // Your business logic here
        err := reserveInventory(ctx, env.Payload)

        // Respond to saga engine
        msgType := data.SUCCESS
        if err != nil {
            msgType = data.FAILURE
        }

        response, _ := json.Marshal(data.Envelope{
            WorkflowID:  env.WorkflowID,
            MessageType: msgType,
            StepName:    env.StepName,
            Payload:     env.Payload,
            Timestamp:   time.Now(),
        })

        return dispatcher.Dispatch(ctx, "saga.events.response", response)
    }, nil)
}

End-to-End Flow with CDC-Axon

β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
β”‚                              FULL PIPELINE                                         β”‚
β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜

  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”         β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”         β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”         β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
  β”‚ API     β”‚  INSERT  β”‚ Postgres β”‚   WAL    β”‚ CDC-Axon β”‚ PUBLISH  β”‚  RabbitMQ   β”‚
  β”‚ Service │────────▢│  Outbox  │────────▢│  Engine  │────────▢│  Queue      β”‚
  β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜         β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜         β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜         β””β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”˜
                                                                        β”‚
                                                                        β–Ό
  β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”     β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”     β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”     β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”
  β”‚  Worker     │◀────│ RabbitMQ │◀────│ saga-axon│◀────│ CDC Middleware       β”‚
  β”‚  Services   β”‚     β”‚ Commands β”‚     β”‚  Engine  β”‚     β”‚ (raw row β†’ Envelope)β”‚
  β””β”€β”€β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”€β”˜     β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜     β””β”€β”€β”€β”€β”¬β”€β”€β”€β”€β”€β”˜     β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜
         β”‚                                    β”‚
         β”‚ SUCCESS/FAILURE                    β”‚ State + Events
         β”‚                                    β–Ό
         └───────────────────────────▢ β”Œβ”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”     β”Œβ”€β”€β”€β”€β”€β”€β”€β”
                                       β”‚  Redis   β”‚     β”‚ Kafka β”‚
                                       β”‚  State   β”‚     β”‚ CQRS  β”‚
                                       β””β”€β”€β”€β”€β”€β”€β”€β”€β”€β”€β”˜     β””β”€β”€β”€β”€β”€β”€β”€β”˜

Project Structure

go-saga-axon/
β”œβ”€β”€ contract/
β”‚   └── interfaces.go          # Listener, Dispatcher, StateStore, EventLogger, Alerter
β”œβ”€β”€ core/
β”‚   β”œβ”€β”€ errors.go              # Retryable, PoisonPill, StateNotFound
β”‚   β”œβ”€β”€ logger.go              # Production (JSON) / Development (text) / Testing (discard)
β”‚   └── router.go              # TriggerGroup β€” binds listeners to workflows
β”œβ”€β”€ data/
β”‚   └── envelope.go            # Envelope, WorkflowState, StepState, constants
β”œβ”€β”€ engine/
β”‚   β”œβ”€β”€ engine.go              # Engine struct, builder, Start(), Shutdown()
β”‚   β”œβ”€β”€ lifecycle.go           # processEvent, handleStart, handleSuccess, handleFailure
β”‚   β”œβ”€β”€ compensate.go          # dispatch, compensateAll, compensateStep
β”‚   β”œβ”€β”€ helpers.go             # debug, alert, workflowLock, allStepsDone, allParentsDone
β”‚   └── monitor.go             # Timeout monitor β€” per-workflow goroutine
β”œβ”€β”€ handlers/
β”‚   └── command.dispatcher.go  # Exponential backoff dispatch wrapper
β”œβ”€β”€ saga/
β”‚   β”œβ”€β”€ middleware/
β”‚   β”‚   └── middleware.go      # CDCMiddleware β€” translates raw CDC to Envelope
β”‚   └── route/
β”‚       β”œβ”€β”€ workflow.go        # Workflow DAG blueprint
β”‚       └── step.go            # Step with Command, Compensation, Timeout, DependsOn
β”œβ”€β”€ store/
β”‚   β”œβ”€β”€ persistence.go         # Idempotency (split check/mark), tracking
β”‚   β”œβ”€β”€ state.store.go         # LoadWorkflowState (retry + CQRS rebuild), SaveWorkflowState
β”‚   └── cqrs.store.go          # LogEvent (retry β†’ buffer), StartFlusher, flushLoop
β”œβ”€β”€ integration/               # Testcontainers-based integration tests
β”‚   β”œβ”€β”€ adapters.go            # Redis, RabbitMQ, Kafka adapter implementations
β”‚   β”œβ”€β”€ setup.go               # Container helpers, worker simulator
β”‚   β”œβ”€β”€ happy_path_test.go
β”‚   β”œβ”€β”€ compensation_test.go
β”‚   β”œβ”€β”€ timeout_test.go
β”‚   β”œβ”€β”€ idempotency_test.go
β”‚   β”œβ”€β”€ cqrs_test.go
β”‚   β”œβ”€β”€ cdc_test.go
β”‚   └── edge_cases_test.go
β”œβ”€β”€ go.mod
β”œβ”€β”€ resilience.md              # 27 documented error scenarios
└── README.md

Resilience

saga-axon handles 27 documented error scenarios including:

  • Redis read/write failures with exponential backoff
  • CQRS log failures with automatic Redis buffering + background flusher
  • Dispatch failures with retry and smart error propagation
  • Compensation dispatch failures with critical alerts
  • Step timeouts with synthetic FAILURE injection
  • Late-arriving SUCCESS during rollback (compensated immediately)
  • Concurrent events on same workflow (per-workflow mutex)
  • State eviction from Redis (CQRS replay from Kafka)
  • Split idempotency preventing stuck sagas on partial failures

See resilience.md for the full breakdown.


Testing

# Unit tests
go test ./engine/... ./store/... ./core/... ./handlers/...

# Integration tests (requires Docker)
go test ./integration/... -v -timeout 120s

Integration tests use testcontainers-go to spin up real Redis, RabbitMQ, and Kafka instances.


Built With


Author

Srajan Saxena β€” IIT (BHU) Varanasi, Computer Science


License

MIT β€” Copyright Β© 2025 Srajan Saxena. See LICENSE for details.

Directories ΒΆ

Path Synopsis
saga

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL