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, andmergeMap.
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
mapover 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
EMPTYto emit no command (skip). mergeMapruns 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
providersalongside command and event handlers. - No manual subscription, no
onModuleInitwiring 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
catchErrorinside the innermergeMapso 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 globalObservableof events to anObservable<ICommand>that theCommandBusdispatches automatically. - Start every saga with
ofType(...)to filter and type-narrow; usemapfor 1:1 andmergeMapfor 0/1/many commands per event. - Register sagas as
providersin aCqrsModule— 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
catchErrorso 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
- Commands, Handlers, and the Command Bus
- Queries and Read-Model Projections
- Domain Events and AggregateRoot
- Sagas for Long-Running Workflows