0Pricing
NestJS Enterprise Backend APIs · Lesson

Sagas for Long-Running Workflows

Coordinate multi-step processes reactively with RxJS-based CQRS sagas.

Sagas for Long-Running Workflows is a free NestJS Enterprise Backend APIs lesson on CoddyKit — lesson 4 of 4. You can read the complete lesson below for free — then practise it hands-on in the browser with a built-in code editor and a 24/7 AI tutor. It is part of the NestJS Enterprise Backend APIs learning path, one of 4 lessons in the course, and your progress syncs across the web and the CoddyKit app.

Why Sagas Exist

In an event-driven system, a single business outcome often requires many steps across multiple aggregates: place order, reserve stock, charge payment, schedule shipping. No single command handler owns that whole flow.

A Saga coordinates such a long-running workflow by listening to events and reacting with new commands. It is the glue that turns one event into the next step of a process.

  • Reactive: a saga is triggered by events, not called directly.
  • Stateless dispatcher: in NestJS, a CQRS saga maps an event stream to a command stream.
  • Decoupled: handlers stay small; the saga owns orchestration.

Sagas in @nestjs/cqrs

In @nestjs/cqrs a saga is a class method decorated with @Saga() that receives an RxJS Observable of all published events and returns an Observable<ICommand>.

Whatever commands the returned stream emits are automatically dispatched through the CommandBus. The framework subscribes to your stream for you.

  • Input type: Observable<any> (the global event stream).
  • Output type: Observable<ICommand>.
  • You shape the flow with RxJS operators like ofType, map, and mergeMap.
import { Injectable } from '@nestjs/common';
import { ICommand, ofType, Saga } from '@nestjs/cqrs';
import { Observable } from 'rxjs';
import { map } from 'rxjs/operators';
import { OrderCreatedEvent } from './events/order-created.event';
import { ReserveStockCommand } from './commands/reserve-stock.command';

@Injectable()
export class OrderSagas {
  @Saga()
  orderCreated = (events$: Observable<any>): Observable<ICommand> => {
    return events$.pipe(
      ofType(OrderCreatedEvent),
      map((event) => new ReserveStockCommand(event.orderId, event.items)),
    );
  };
}

The ofType Operator

The global stream carries every published event. You almost always start a saga by filtering it down to the event types you care about with ofType(...EventClasses).

ofType is a custom RxJS operator shipped by @nestjs/cqrs. It filters by event constructor and, crucially, narrows the TypeScript type of downstream values to that event, so event.orderId type-checks.

  • Pass one or more event classes: ofType(A, B).
  • Always filter early — never map over the raw stream blindly.

From Event to Command

The simplest saga is a one-to-one translation: each matching event produces exactly one command. map is the right operator here because it is synchronous and emits one value per input.

The pattern below reacts to a successful stock reservation by issuing a payment command — moving the workflow forward one step.

import { Injectable } from '@nestjs/common';
import { ICommand, ofType, Saga } from '@nestjs/cqrs';
import { Observable } from 'rxjs';
import { map } from 'rxjs/operators';
import { StockReservedEvent } from './events/stock-reserved.event';
import { ChargePaymentCommand } from './commands/charge-payment.command';

@Injectable()
export class PaymentSagas {
  @Saga()
  stockReserved = (events$: Observable<any>): Observable<ICommand> =>
    events$.pipe(
      ofType(StockReservedEvent),
      map((e) => new ChargePaymentCommand(e.orderId, e.amount)),
    );
}

map vs mergeMap

Use map when one event yields exactly one command synchronously. Use mergeMap (a.k.a. flatMap) when a single event must produce zero, one, or many commands, or when you need an inner Observable.

  • mergeMap + of(...) lets you emit several commands per event.
  • Return EMPTY to emit no command (skip).
  • mergeMap runs inner streams concurrently — good for independent fan-out steps.
import { Injectable } from '@nestjs/common';
import { ICommand, ofType, Saga } from '@nestjs/cqrs';
import { Observable, of, EMPTY } from 'rxjs';
import { mergeMap } from 'rxjs/operators';
import { PaymentConfirmedEvent } from './events/payment-confirmed.event';
import { ScheduleShippingCommand } from './commands/schedule-shipping.command';
import { SendReceiptCommand } from './commands/send-receipt.command';

@Injectable()
export class FulfillmentSagas {
  @Saga()
  paymentConfirmed = (events$: Observable<any>): Observable<ICommand> =>
    events$.pipe(
      ofType(PaymentConfirmedEvent),
      mergeMap((e) =>
        e.amount > 0
          ? of(
              new ScheduleShippingCommand(e.orderId),
              new SendReceiptCommand(e.orderId, e.amount),
            )
          : EMPTY,
      ),
    );
}

Registering Sagas in a Module

A saga class only runs if it is listed in the module's providers. NestJS's CqrsModule discovers every provider that exposes @Saga() methods and subscribes them to the event bus at bootstrap.

  • Import CqrsModule.
  • Add the saga class to providers alongside command and event handlers.
  • No manual subscription, no onModuleInit wiring needed.
import { Module } from '@nestjs/common';
import { CqrsModule } from '@nestjs/cqrs';
import { OrderSagas } from './sagas/order.sagas';
import { PaymentSagas } from './sagas/payment.sagas';
import { ReserveStockHandler } from './commands/reserve-stock.handler';
import { ChargePaymentHandler } from './commands/charge-payment.handler';

@Module({
  imports: [CqrsModule],
  providers: [
    OrderSagas,
    PaymentSagas,
    ReserveStockHandler,
    ChargePaymentHandler,
  ],
})
export class OrderingModule {}

Correlating Multi-Step State

Most sagas need a correlation id — typically the aggregate id (e.g. orderId) — carried on every event so steps can be tied to the same process instance.

The CQRS saga itself is stateless: it just maps events to commands. The workflow state (which step completed, what was reserved) lives in the aggregate or a dedicated read model, updated by command handlers. The saga reacts to the events those handlers emit.

  • Always include the correlation id in event payloads.
  • Keep mutable progress in an aggregate/persistence layer, not in the saga.

Compensation: The Saga's Real Job

Distributed workflows cannot use a single ACID transaction across services. Instead a saga guarantees consistency through compensating actions: if a later step fails, earlier successful steps are semantically undone.

If payment fails after stock was reserved, the saga reacts to PaymentFailedEvent by dispatching a ReleaseStockCommand. Compensation is forward-recovery, not a rollback — you issue a new command that reverses the effect.

import { Injectable } from '@nestjs/common';
import { ICommand, ofType, Saga } from '@nestjs/cqrs';
import { Observable } from 'rxjs';
import { map } from 'rxjs/operators';
import { PaymentFailedEvent } from './events/payment-failed.event';
import { ReleaseStockCommand } from './commands/release-stock.command';

@Injectable()
export class CompensationSagas {
  @Saga()
  paymentFailed = (events$: Observable<any>): Observable<ICommand> =>
    events$.pipe(
      ofType(PaymentFailedEvent),
      map((e) => new ReleaseStockCommand(e.orderId, e.items)),
    );
}

Errors Inside the Saga Stream

A saga returns one long-lived Observable. If an error escapes that stream, RxJS terminates the subscription and your saga stops reacting to all future events — a silent, system-wide failure.

Protect the stream with catchError. Recover by mapping the failure to a compensating command and continuing, rather than letting the observable die.

  • Never let exceptions propagate out of the saga pipe.
  • Prefer per-event isolation: put catchError inside the inner mergeMap so one bad event doesn't kill the outer stream.
import { Injectable, Logger } from '@nestjs/common';
import { ICommand, ofType, Saga } from '@nestjs/cqrs';
import { Observable, of } from 'rxjs';
import { mergeMap, catchError } from 'rxjs/operators';
import { ShippingRequestedEvent } from './events/shipping-requested.event';
import { NotifyOpsCommand } from './commands/notify-ops.command';
import { DispatchCarrierCommand } from './commands/dispatch-carrier.command';

@Injectable()
export class ShippingSagas {
  private readonly logger = new Logger(ShippingSagas.name);

  @Saga()
  shippingRequested = (events$: Observable<any>): Observable<ICommand> =>
    events$.pipe(
      ofType(ShippingRequestedEvent),
      mergeMap((e) =>
        of(new DispatchCarrierCommand(e.orderId)).pipe(
          catchError((err) => {
            this.logger.error(err);
            return of(new NotifyOpsCommand(e.orderId, 'dispatch failed'));
          }),
        ),
      ),
    );
}

Idempotency and At-Least-Once Delivery

When events are delivered over a broker (Kafka, RabbitMQ) the saga may see the same event more than once after a redelivery or restart. Dispatching the same command twice can double-charge or double-ship.

The fix is not in the saga's RxJS plumbing but in the command handler: make it idempotent by recording a processed key (correlation id + step) and ignoring duplicates.

  • Sagas should produce commands that are safe to retry.
  • Track (orderId, step) in a dedup table or use the aggregate's version.
  • Design for at-least-once, not exactly-once.

A Standalone RxJS Saga Pipeline

You can model the exact event-to-command mapping a saga performs using plain RxJS — no NestJS runtime required. This mirrors how ofType + map drive a workflow, and is great for reasoning about the flow in isolation.

import { from } from 'rxjs';
import { filter, map } from 'rxjs/operators';

class OrderCreated { constructor(public orderId: string) {} }
class PaymentFailed { constructor(public orderId: string) {} }
class ReserveStock { constructor(public orderId: string) {} }
class ReleaseStock { constructor(public orderId: string) {} }

const ofType =
  <T>(type: new (...a: any[]) => T) =>
  (source: any) =>
    source.pipe(filter((e: any): e is T => e instanceof type));

const events$ = from([
  new OrderCreated('A1'),
  new PaymentFailed('A1'),
]);

events$
  .pipe(
    map((e) =>
      e instanceof OrderCreated
        ? new ReserveStock(e.orderId)
        : e instanceof PaymentFailed
          ? new ReleaseStock(e.orderId)
          : null,
    ),
    filter((c): c is ReserveStock | ReleaseStock => c !== null),
  )
  .subscribe((cmd) => console.log('dispatch ->', cmd.constructor.name, cmd.orderId));

Quick Check

A teammate writes a CQRS saga whose returned Observable occasionally throws when building a command for a malformed event. After one such event in production, the saga stops reacting to all subsequent events. What is the correct fix?

Recap

You learned how to coordinate long-running workflows with RxJS-based CQRS sagas in NestJS:

  • A @Saga() maps the global Observable of events to an Observable<ICommand> that the CommandBus dispatches automatically.
  • Start every saga with ofType(...) to filter and type-narrow; use map for 1:1 and mergeMap for 0/1/many commands per event.
  • Register sagas as providers in a CqrsModule — no manual subscription.
  • Carry a correlation id on events; keep workflow state in aggregates/read models, not the saga.
  • Achieve distributed consistency via compensating commands, not transactions.
  • Guard the stream with catchError so one failure can't kill the saga, and make command handlers idempotent for at-least-once delivery.

Frequently asked questions

Is the “Sagas for Long-Running Workflows” lesson free?

Yes — the full text of “Sagas for Long-Running Workflows” is free to read here on the web, and the NestJS Enterprise Backend APIs course includes 4 lessons in total. To practise it interactively (a built-in code editor and a 24/7 AI tutor) and unlock the rest of the NestJS Enterprise Backend APIs course, upgrade to CoddyKit PRO.

What will I learn in “Sagas for Long-Running Workflows”?

Coordinate multi-step processes reactively with RxJS-based CQRS sagas. You practise NestJS Enterprise Backend APIs with hands-on code you run directly in the browser, and a 24/7 AI tutor answers your questions as you work through the lesson.

Do I need any experience to start NestJS Enterprise Backend APIs?

No prior experience is required. NestJS Enterprise Backend APIs on CoddyKit is structured for beginners through advanced learners; this is — lesson 4 of 4, so you can start here or from the beginning and move at your own pace.

How long does the “Sagas for Long-Running Workflows” lesson take?

Most CoddyKit lessons take about 5–10 minutes. Each one is bite-sized and interactive, so you make steady progress and pick up exactly where you left off across the web and the app.

Can I write and run code in this NestJS Enterprise Backend APIs lesson?

Yes. Every NestJS Enterprise Backend APIs lesson includes a built-in code editor, so you write and run real code right in your browser and get instant AI feedback — no local setup required.

All lessons in this course

  1. Commands, Handlers, and the Command Bus
  2. Queries and Read-Model Projections
  3. Domain Events and AggregateRoot
  4. Sagas for Long-Running Workflows
← Back to NestJS Enterprise Backend APIs