Reliability › Sagas

Sagas

A saga is a business transaction that spans services. There is no distributed database transaction between orders, inventory and payments, so a saga replaces "rollback" with compensation: every step that has a side-effect declares how to undo it.

NestLaravel ships an orchestrated saga runner in nestlaravel/kafka (NestLaravel\Kafka\Saga). The orchestrator lives inside one service (typically the one that owns the business process, e.g. orders), keeps its state in that service's database and talks to the other services through events (via the transactional outbox).

use NestLaravel\Kafka\Saga\Saga;

class="tk-c">// AppServiceProvider::boot()
Saga::define(class="tk-s">'order_checkout')
    ->step(ReserveInventory::class)->compensate(ReleaseInventory::class)
    ->step(ProcessPayment::class)->compensate(RefundPayment::class)->retries(2)
    ->step(ConfirmOrder::class);

class="tk-c">// From an Action, after the order row is created
Saga::start(class="tk-s">'order_checkout', [class="tk-s">'order_id' => class="tk-v">$order->id], correlationId: class="tk-v">$order->id);

#Writing a step

final class ProcessPayment implements SagaStep
{
    public function execute(SagaContext class="tk-v">$ctx): StepResult
    {
        app(EventBus::class)->publish(new PaymentRequested(class="tk-v">$ctx->get(class="tk-s">'order_id'), class="tk-v">$ctx->get(class="tk-s">'total')));

        class="tk-c">// Pause until payments answers. Either failure event, or 300 s of silence, triggers compensation.
        return StepResult::waitFor(
            successEvent: class="tk-s">'payments.payment.completed',
            failureEvents: [class="tk-s">'payments.payment.failed'],
            timeoutSeconds: 300,
        );
    }
}
ResultMeaning
StepResult::done($data)Finished synchronously; $data is merged into the saga context; continue.
StepResult::waitFor($ok, [$fail…], $timeout)Remote work started; pause.
StepResult::fail($reason)Business rejection: compensate now, no retries.
throwsTechnical failure: retried up to ->retries(n) (default 2), then compensate.

A compensation implements Compensation::compensate(SagaContext $ctx).

#Feeding events into the saga

Sagas resume from events. In the orchestrating service's Kafka consumer handler:

public function handle(array class="tk-v">$event): void
{
    app(SagaOrchestrator::class)->handleEvent(class="tk-v">$event[class="tk-s">'event_type'], class="tk-v">$event[class="tk-s">'payload'], class="tk-v">$event[class="tk-s">'correlation_id']);
}

Run php artisan saga:recover every minute (scheduler). It handles timeouts and crashed runners.

#What is guaranteed (each item has a test in tests/Reliability/SagaTest.php)

GuaranteeTest
Happy path pauses for the event, then completestest_happy_path_pauses_for_the_payment_event_then_completes
Failure event ⇒ completed steps compensated in reverse ordertest_payment_failed_compensates_completed_steps_in_reverse_order
Duplicate start() and duplicate events do nothing extratest_duplicate_start_and_duplicate_events_are_idempotent
Timeout ⇒ compensationtest_timeout_triggers_compensation_via_recovery
Throwing step is retried; exhausted retries compensatetest_a_throwing_step_is_retried_then_succeeds, test_exhausted_retries_compensate_earlier_steps
Business rejection never retriestest_business_rejection_compensates_immediately_without_retries
Failing compensation parks the saga in failed instead of guessingtest_compensation_that_keeps_failing_parks_the_saga_for_a_human
Step state and its outbox events commit atomicallytest_step_events_are_committed_atomically_with_saga_state
State survives restarts; crashed runner is recoveredtest_state_survives_a_process_restart_and_a_crashed_runner_is_recovered
Two workers cannot advance one saga at oncetest_two_workers_cannot_advance_the_same_saga_concurrently

#What is NOT guaranteed — read before relying on it

#Operating

SignalWhere
Instances and their history (every transition is recorded)table saga_instances, column history
Stuck in waiting past deadline_atphp artisan saga:recover handles it; alert if saga_instances has waiting rows older than expected
Needs a humanstatus = 'failed'
Countersnestlaravel_saga_total{saga, result=started, completed, compensated, timed_out or failed}

Correlation: the saga's correlation_id is put in every log line of the step and is inherited as causation/correlation by events the step publishes, so one order can be followed across services in logs and traces.

Edit this page on GitHub