InteractiveFrameworks

CQRS

Keep the cats service's writes and reads apart: messages become commands and queries on their buses, an adoption is an aggregate that applies an event, an event handler keeps the notice board, and a saga turns every adoption into the next command.

What you'll learn

  • Explain CQRS: commands change state and say little, queries read and change nothing, events announce what changed
  • Dispatch through CommandBus and QueryBus, and handle with @CommandHandler and @QueryHandler classes registered as providers
  • Model a change as an AggregateRoot that applies events, and publish them with EventPublisher.mergeObjectContext() and commit()
  • React to events with @EventsHandler, and chain work with a @Saga that maps events to commands

The cats service started as a list and a lookup. Now an adoption records the adopter, refuses a cat that is already home, writes a notice for the board, and books the cat's first vet check. Next month it will send a welcome email and update the partner shelters. Put all of that in one service method and every new reaction means editing the code that adopts, while the board and the lookups want data shaped for reading, not for changing.

CQRS, Command and Query Responsibility Segregation, splits the two. Commands change state and return little. Queries read and change nothing. Events announce what a command changed, so any number of reactions can hang off them without the command knowing. @nestjs/cqrs gives each a bus and a kind of handler. CqrsModule.forRoot() in the module's imports provides the buses.

Commands and queries

A command or a query is a plain class carrying its data. Extending Command<T> or Query<T> declares what executing it returns, so execute() is typed:

export class RenewLicenceCommand extends Command<Licence> {
  constructor(readonly shopId: number) { super(); }
}

@CommandHandler(RenewLicenceCommand)
export class RenewLicenceHandler implements ICommandHandler<RenewLicenceCommand> {
  constructor(private readonly licences: LicencesRepository) {}
  async execute(command: RenewLicenceCommand): Promise<Licence> { ... }
}

commandBus.execute(new RenewLicenceCommand(7)) finds the one handler for that class and returns its result. The QueryBus works the same way with @QueryHandler and IQueryHandler. Handlers are providers, so they get dependency injection, and they must be listed in the module's providers. Nest finds them there by their decorators. One that is not listed does not exist: No handler found for the query: "GetShopQuery".

Here the buses sit behind the transport. The service's message handlers stay thin: each message becomes a command or a query, and a thrown RpcException still reaches the caller as before.

Aggregates and their events

An aggregate is the object a command changes, and it says what happened by applying events. Extend AggregateRoot:

export class Licence extends AggregateRoot {
  renew(): void {
    this.expires = nextYear();
    this.apply(new LicenceRenewedEvent(this.shopId));
  }
}

apply() only queues the event. An aggregate loaded from a repository does not know the event bus, so the command handler connects them and then commits:

const licence = this.publisher.mergeObjectContext(await this.licences.find(id));
licence.renew();
licence.commit();

EventPublisher.mergeObjectContext() gives the object a way to publish, and commit() publishes every queued event. Forget either and nothing fails: the events stay in the aggregate and no reaction ever runs.

Reacting: event handlers and sagas

@EventsHandler(LicenceRenewedEvent) on a class with a handle(event) method reacts to an event. Several handlers may react to the same event. They run after the command, outside the caller's request: an error thrown there goes to the UnhandledExceptionBus and is logged, and never reaches the caller.

A saga reacts by issuing further commands. It is a property holding a function from the stream of all events to a stream of commands, marked with @Saga():

@Injectable()
export class LicenceSagas {
  @Saga()
  renewed = (events$: Observable<unknown>): Observable<ICommand> =>
    events$.pipe(ofType(LicenceRenewedEvent), map((event) => new PrintCertificateCommand(event.shopId)));
}

ofType() keeps the events of one class, and every command the saga emits is executed by the CommandBus. The saga's class is a provider too.

Your task

The service's controller dispatches GetCatQuery, AdoptCatCommand and GetBoardQuery. CatsRepository holds Cat aggregates, ShelterBoard holds the notices and appointments, and ScheduleVetCheckCommand has its handler already.

  1. Write GetCatHandler: the cat from the repository, or an RpcException carrying { status: 404, message: 'Cat <id> not found' }.
  2. In Cat.adopt(), record the adopter and apply a CatAdoptedEvent for the cat.
  3. Write AdoptCatHandler. The cat must exist (the same 404) and have no adopter yet, or it is refused with { status: 409, message: '<name> already went home with <adopter>' }. Adopt it so that its events reach the event bus, and return it.
  4. Write CatAdoptedHandler, which adds <name> went home with <adopter> to the board's notices, and AdoptionSagas, which turns every CatAdoptedEvent into a ScheduleVetCheckCommand for that cat.
  5. Register everything you wrote in the module's providers.

When it fails

  • A 502 and the service logs No handler found for the query: "GetCatQuery": the handler is not in providers, or its decorator names another class.
  • The adoption succeeds and the board stays empty, with nothing logged: the events never left the aggregate. The handler needs both mergeObjectContext() and commit().
  • Notices appear but no appointments: the saga is not in providers, or its property lost @Saga().

Remember

  • Commands change state, queries read it, events announce what changed.
  • Handlers are providers found by their decorators; an unlisted one does not exist.
  • An aggregate applies events; mergeObjectContext() and commit() publish them.
  • Event handlers react, sagas map events to new commands; neither can answer the caller.
Stuck? Show a hint

@QueryHandler(GetCatQuery) and @CommandHandler(AdoptCatCommand) on classes implementing IQueryHandler / ICommandHandler with an execute(); in the command handler, this.publisher.mergeObjectContext(cat), cat.adopt(by), cat.commit(). In the model, this.apply(new CatAdoptedEvent(...)). @EventsHandler(CatAdoptedEvent) with handle(event). A saga is a property: @Saga() adopted = (events$) => events$.pipe(ofType(CatAdoptedEvent), map(...)). Every one of them goes in the module's providers.