Skip to main content

Saga Orchestration

Coordinate multi-step business processes across independent agents — with automatic rollback when anything fails.

A Saga is a named sequence of steps where each step is handled by a separate autonomous workflow. If any required step fails, the Saga coordinator automatically walks backward through all previously completed steps and triggers their compensating actions to undo the work.


When to use Sagas​

Use a Saga when a business process:

  • Spans multiple systems or departments (HR + IT + Finance for employee onboarding)
  • Must be fully completed or fully undone (partial state is worse than no state)
  • Has steps that can fail independently (IT provisioning can fail after HR record creation)
  • Runs over minutes or hours, not milliseconds

Do not use Sagas for a single-system operation or a simple workflow with if_else branching — those are cleaner as a regular sequential workflow.


How a Saga works​

Trigger event arrives (e.g. employee_hired)
│
▼
Saga Coordinator creates a SagaInstance
│
├──▶ emits: start_hr_step ─────▶ HR Workflow runs
│ │
│ hr_record_created ◀────────────────┘
│
├──▶ emits: start_it_step ─────▶ IT Workflow runs
│ │
│ it_accounts_provisioned ◀──────────┘
│
└──▶ emits: employee_onboarding_completed

If IT provisioning fails:

            it_provisioning_failed
│
▼
Saga Coordinator starts compensation
│
├──▶ emits: delete_hr_record ─────▶ HR Compensation Workflow
│ │
│ hr_record_deleted ◀────────────────────────┘
│
└──▶ emits: saga_compensated

Saga vocabulary​

TermMeaning
Saga DefinitionBlueprint: trigger event, steps, compensation events
Saga InstanceOne running saga for a specific correlation_id
StepOne unit of work, handled by one subscribed workflow
Compensation eventEvent emitted to undo a completed step
Correlation IDBusiness key linking all events in this saga (e.g. EMP001)

Defining a Saga​

Create a Saga Definition via API. Each step declares what event starts it, what events signal success or failure, and what event triggers compensation:

curl -X POST http://localhost:8080/api/v1/sagas/definitions \
  -H "Authorization: Bearer <admin_token>" \
  -H "Content-Type: application/json" \
  -d '{
    "name": "employee_onboarding",
    "trigger_event":   "employee_hired",
    "correlation_key": "$.employee_id",
    "steps": [
      {
        "step_id":                "hr_step",
        "name":                   "Create HR Record",
        "trigger_event":          "start_hr_step",
        "success_event":          "hr_record_created",
        "failure_event":          "hr_record_failed",
        "compensation_event":     "delete_hr_record",
        "compensation_done_event":"hr_record_deleted",
        "timeout_minutes":        30,
        "required":               true
      },
      {
        "step_id":                "it_step",
        "name":                   "Provision IT Accounts",
        "trigger_event":          "start_it_step",
        "success_event":          "it_accounts_provisioned",
        "failure_event":          "it_provisioning_failed",
        "compensation_event":     "delete_it_accounts",
        "compensation_done_event":"it_accounts_deleted",
        "timeout_minutes":        60,
        "required":               true
      }
    ]
  }'

Step fields:

FieldDescription
step_idUnique ID within this saga
trigger_eventEvent the coordinator emits to start this step's workflow
success_eventEvent the workflow emits on success (advances saga to next step)
failure_eventEvent the workflow emits on failure (triggers compensation if required: true)
compensation_eventEvent emitted to the compensation workflow to undo this step
compensation_done_eventEvent the compensation workflow emits when undo is complete. Defaults to comp_{step_id}_done
timeout_minutesMax time before the coordinator auto-fails the step
requiredtrue = saga compensates on failure; false = saga skips and continues

correlation_key​

A dot-path into the trigger event's payload to extract the correlation ID. Examples:

  • $.employee_id → payload.employee_id
  • $.order.id → payload.order.id
  • $.invoice_number → payload.invoice_number

One saga instance is created per unique (saga_definition, correlation_id) pair.


Building the step workflows​

Each saga step needs two workflows: one that does the work, and one that undoes it.

Work workflow (subscribes to trigger_event)​

# Subscribe the HR workflow to start_hr_step
curl -X POST "http://localhost:8080/api/v1/agent-workflows/{hr_wf_id}/triggers/event" \
  -H "Authorization: Bearer <token>" \
  -H "Content-Type: application/json" \
  -d '{
    "name":        "On start_hr_step",
    "event_type":  "start_hr_step",
    "concurrency": "one_per_correlation"
  }'

Agent node system prompt:

You are creating an HR employee record.

The event payload (in state._event.payload) contains:
employee_id, name, department, email

Steps:
1. Call the HR system API to create the employee record
2. On success, emit success event with the new record ID:
emit_event({
"event_type": "hr_record_created",
"correlation_id": "<from input>",
"payload": "{\"employee_id\": \"...\", \"hr_record_id\": \"...\"}"
})
3. On failure, emit failure event with the error:
emit_event({
"event_type": "hr_record_failed",
"correlation_id": "<from input>",
"payload": "{\"error\": \"...\"}"
})

Compensation workflow (subscribes to compensation_event)​

# Subscribe the HR undo workflow to delete_hr_record
curl -X POST "http://localhost:8080/api/v1/agent-workflows/{hr_undo_wf_id}/triggers/event" \
  -H "Content-Type: application/json" \
  -H "Authorization: Bearer <token>" \
  -d '{"name": "Undo HR record", "event_type": "delete_hr_record", "concurrency": "one_per_correlation"}'

Agent node system prompt:

You are undoing an HR employee record creation.

The event payload contains employee_id and hr_record_id.

Steps:
1. Call the HR system API to delete the record
2. On completion (success or already-deleted), signal done:
emit_event({
"event_type": "hr_record_deleted",
"correlation_id": "<from input>",
"payload": "{\"employee_id\": \"...\"}"
})

The compensation workflow MUST emit compensation_done_event so the coordinator knows the undo is complete.


Saga lifecycle​

StatusMeaning
runningCoordinator is actively advancing through steps
completedAll steps succeeded
compensatingA required step failed; coordinator is undoing completed steps backward
compensatedAll compensation completed successfully
failedUnrecoverable error (step definition mismatch, compensation timeout)

Monitoring sagas​

# List active sagas
GET /api/v1/sagas/instances?status=running

# Full timeline for one saga
GET /api/v1/sagas/instances/{id}

Response includes the instance (current status, step index) and steps (per-step status, workflow_run_id, started_at, ended_at).

Manual controls (admin)​

# Trigger compensation manually (e.g. business decision to undo)
POST /api/v1/sagas/instances/{id}/compensate

# Re-emit the current step's trigger event (unstick a stuck step)
POST /api/v1/sagas/instances/{id}/retry

Timeout behavior​

Each step has a timeout_minutes setting. If the step's workflow hasn't emitted a success or failure event within that time:

  1. The step is marked failed
  2. The coordinator starts compensation from the previous completed step
  3. A saga_compensation_stalled event is emitted if compensation itself gets stuck

The SAGA_COMPENSATION_TIMEOUT_MINUTES environment variable (default: 60) controls how long a compensation step is allowed to run before the saga is flagged for manual investigation.


Optional steps​

Set "required": false on a step to allow the saga to continue even when that step fails:

{
  "step_id": "send_welcome_email",
  "required": false,
  ...
}

If the welcome email step fails, the saga marks it skipped and continues to the next step. Only required: true steps trigger compensation.


End-to-end test​

# 1. Fire the trigger event
curl -X POST http://localhost:8080/api/v1/events/ingest \
  -H "X-Api-Key: <key>" \
  -H "Content-Type: application/json" \
  -d '{"event_type": "employee_hired", "correlation_id": "EMP001",
       "payload": {"employee_id": "EMP001", "name": "Alice"}}'

# 2. Watch events flow
curl "http://localhost:8080/api/v1/events?correlation_id=EMP001" \
  -H "Authorization: Bearer <token>"

# 3. Check saga progress
curl "http://localhost:8080/api/v1/sagas/instances?correlation_id=EMP001" \
  -H "Authorization: Bearer <token>"

Hrida.ai is proprietary software of Zlabs Innovation. See the license for terms. © 2026 Zlabs Innovation.