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.
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:
- Models your workflow as a DAG β steps can run in parallel (fan-out) and converge (fan-in)
- Dispatches commands to worker services via any message broker
- Tracks state in Redis with per-workflow locking
- Compensates automatically β on failure, ALL succeeded steps get rollback commands
- Handles timeouts β steps that take too long are treated as failures
- Survives crashes β CQRS event log enables state rebuild from Kafka
- 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:
- Your service writes to the outbox table:
INSERT INTO outbox (order_id, event_type, payload) VALUES (...) - CDC-Axon captures the WAL event and publishes the raw row to RabbitMQ
- saga-axon's CDC middleware extracts
order_idasworkflow_id, wraps it in a START envelope - 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
goRabbit-axonβ Production RabbitMQ client with connection pooling and circuit breakerscdc-axonβ CDC engine for PostgreSQL WAL and MongoDB change streamsgo-redis/v9β Redis clientIBM/saramaβ Kafka clientcenkalti/backoffβ Exponential backofftestcontainers-goβ Integration testing
Author
Srajan Saxena β IIT (BHU) Varanasi, Computer Science
License
MIT β Copyright Β© 2025 Srajan Saxena. See LICENSE for details.