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
| Term | Meaning |
|---|---|
| Saga Definition | Blueprint: trigger event, steps, compensation events |
| Saga Instance | One running saga for a specific correlation_id |
| Step | One unit of work, handled by one subscribed workflow |
| Compensation event | Event emitted to undo a completed step |
| Correlation ID | Business 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:
| Field | Description |
|---|---|
step_id | Unique ID within this saga |
trigger_event | Event the coordinator emits to start this step's workflow |
success_event | Event the workflow emits on success (advances saga to next step) |
failure_event | Event the workflow emits on failure (triggers compensation if required: true) |
compensation_event | Event emitted to the compensation workflow to undo this step |
compensation_done_event | Event the compensation workflow emits when undo is complete. Defaults to comp_{step_id}_done |
timeout_minutes | Max time before the coordinator auto-fails the step |
required | true = 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
| Status | Meaning |
|---|---|
running | Coordinator is actively advancing through steps |
completed | All steps succeeded |
compensating | A required step failed; coordinator is undoing completed steps backward |
compensated | All compensation completed successfully |
failed | Unrecoverable 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}/retryTimeout behavior
Each step has a timeout_minutes setting. If the step's workflow hasn't emitted a success or failure event within that time:
- The step is marked
failed - The coordinator starts compensation from the previous completed step
- A
saga_compensation_stalledevent 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>"Related
- Event-Driven Workflows — how the Event Bus and subscriptions work
- Unpredictable Process Patterns — choosing between sagas and other patterns
- emit_event tool — how workflows signal saga step completion