Skip to content

CQRS without NestJS

This guide builds a small users application on @cqrs-ddd/cqrs: decorated handlers, a pipeline around each one, and buses built by createCqrs(). There is no container: one function builds everything with new, so every dependency is visible and checked by the compiler. The repository’s api/ is a complete application built the same way.

Handlers depend on ports, which the application implements: in memory for tests, on a database in production.

export abstract class Users {
abstract add(name: string): User;
abstract find(id: string): User | undefined;
}
export abstract class Mailer {
abstract send(to: string, text: string): Promise<void>;
}

CreateUserCommand is a plain class, and commandBus.execute<CreateUserCommand, string>() types its result. Its handler adds the user and publishes an event. @UsePipeline runs it once per request id and audits it: a repeated requestId returns the first result without running the handler again.

export class CreateUserCommand {
constructor(
readonly requestId: string,
readonly name: string,
) {}
}
@CommandHandler(CreateUserCommand)
@UsePipeline(
idempotent({ keyFactory: (context) => (context.request as CreateUserCommand).requestId }),
audit({ action: 'user.create' }),
)
export class CreateUserHandler implements ICommandHandler<CreateUserCommand> {
constructor(
private readonly users: Users,
private readonly events: EventBus,
) {}
async execute(command: CreateUserCommand): Promise<string> {
const user = this.users.add(command.name);
this.events.publish(new UserCreatedEvent(user.id, user.name));
return user.id;
}
}

The application logs every handler through a global LoggingBehavior; @SkipPipeline(LoggingBehavior) leaves this one out.

@QueryHandler(GetUserQuery)
@SkipPipeline(LoggingBehavior)
export class GetUserHandler implements IQueryHandler<GetUserQuery> {
constructor(private readonly users: Users) {}
async execute(query: GetUserQuery): Promise<User | undefined> {
return this.users.find(query.id);
}
}

The event bus starts SendWelcomeMail inside publish() and does not await it, so the command does not wait for the mail. It runs in a pipeline of its own, under the command’s correlation id.

@EventsHandler(UserCreatedEvent)
export class SendWelcomeMail implements IEventHandler<UserCreatedEvent> {
constructor(private readonly mailer: Mailer) {}
async handle(event: UserCreatedEvent): Promise<void> {
await this.mailer.send(event.userId, `Welcome, ${event.name}!`);
}
}

One function builds the application: the behavior instances that need dependencies, the buses, then the handlers with their ports. register() compiles each handler’s pipeline and fails there on any configuration error.

import { createCqrs } from '@cqrs-ddd/cqrs';
import { LoggingBehavior } from '@cqrs-ddd/pipeline';
import { AuditBehavior } from '@cqrs-ddd/pipeline-audit';
import { IdempotencyBehavior, MemoryIdempotencyStore } from '@cqrs-ddd/pipeline-idempotency';
export function createApp(ports: { users: Users; mailer: Mailer; auditSink: AuditSink }) {
const cqrs = createCqrs({
behaviors: [
new IdempotencyBehavior(new MemoryIdempotencyStore()),
new AuditBehavior(ports.auditSink),
],
globalBehaviors: { before: [LoggingBehavior] },
});
cqrs.register(
new CreateUserHandler(ports.users, cqrs.eventBus),
new GetUserHandler(ports.users),
new SendWelcomeMail(ports.mailer),
);
return cqrs;
}

close() waits for the event handlers still running, so after it the welcome mail has been sent.

const app = createApp({ users: new MemoryUsers(), mailer, auditSink: new PostgresAuditSink(pool) });
const id = await app.commandBus.execute<CreateUserCommand, string>(
new CreateUserCommand('req-1', 'Ann'),
);
const user = await app.queryBus.execute<GetUserQuery, User | undefined>(
new GetUserQuery(id),
);
await app.close();

Coming from NestJS lists what changes when moving code between the two.