node-cqrs: CQRS and Event Sourcing for TypeScript

Build CQRS and Event Sourcing applications in TypeScript with aggregates, sagas, projections, dependency injection, and pluggable infrastructure.

Features

  • Plain messages: Commands and events are ordinary typed objects without decorators or generated classes.
  • Focused domain blocks: Aggregates handle commands, projections build read models, and sagas coordinate work.
  • Replaceable infrastructure: Thin interfaces let applications provide their own storage, buses, locks, views, and dispatch processors.
  • Concurrency handling: Commands are serialized per aggregate within a process; persistent event stores add optimistic concurrency for distributed writers.
  • Projection lifecycle: Restore hooks, readiness locks, event deduplication, and checkpoints are available to persistent views.
  • Selective rehydration: Aggregates can restore selected events and use optional snapshots.
  • Dispatch pipelines: Event batches pass through configurable persistence and processing pipelines with concurrency limits.

Infrastructure modules can be combined according to the deployment:

  • node-cqrs/sqlite - embedded event storage and relational or JSON views;
  • node-cqrs/mongodb - distributed event storage and document views;
  • node-cqrs/redis - distributed document projection views;
  • node-cqrs/postgresql - transactional event storage and relational or JSON views;
  • node-cqrs/rabbitmq - distributed command and event buses;
  • node-cqrs/workers - worker-thread projections for CPU-intensive handlers.

Installation

npm install node-cqrs

The built package supports Node.js 18 and later. The TypeScript examples can be executed directly with Node.js 24 or later; earlier Node.js versions require normal TypeScript compilation or a loader.

The browser bundle exposes the browser-compatible core API. Database adapters, RabbitMQ, and Node.js worker threads are server-side modules. Infrastructure modules require their documented peer dependencies.

Quick Start

This example defines one command, one event, and one read model entirely in memory:

import { AbstractAggregate, AbstractProjection, ContainerBuilder, InMemoryEventStorage } from 'node-cqrs';
import type { IContainer, IEvent, Identifier } from 'node-cqrs';

type UserRecord = {
	username: string;
};

type UserCreatedEvent = IEvent<UserRecord>;
type UsersView = Map<Identifier, UserRecord>;

class UserAggregate extends AbstractAggregate {
	createUser(payload: UserRecord) {
		this.emit('userCreated', payload);
	}
}

class UsersProjection extends AbstractProjection<UsersView> {
	constructor() {
		super({ view: new Map() });
	}

	userCreated(event: UserCreatedEvent) {
		this.view.set(event.aggregateId!, event.payload);
	}
}

interface AppContainer extends IContainer {
	usersView: UsersView;
}

const builder = new ContainerBuilder<AppContainer>();
builder.register(InMemoryEventStorage);
builder.registerAggregate(UserAggregate);
builder.registerProjection(UsersProjection, 'usersView');

const container = builder.container();
const { usersView, commandBus } = container;

const [userCreated] = await commandBus.send('createUser', undefined, {
	payload: { username: 'alice' }
});

console.log(usersView.get(userCreated.aggregateId!)); // { username: 'alice' }

InMemoryEventStorage is useful for learning and tests; its events disappear when the process exits. Choose a persistent event store from Infrastructure for an application that must survive restarts.

How It Fits Together

Commands flow through aggregates and events update projections and sagas

Domain behavior is split into three small blocks:

  • Aggregates restore write-side state, validate commands, and emit events.
  • Projections consume events and update read-side views.
  • Sagas react to events and enqueue commands for multi-step processes.

The default runtime flow is:

  1. The command bus delivers a command to an aggregate command handler.
  2. The handler restores the target aggregate and invokes its command method.
  3. Emitted events pass through the event dispatch pipeline, including configured persistence.
  4. The event bus delivers committed events to projections, sagas, and other subscribers.

Messages And Replacement Points

Commands and events are plain objects. A message needs only a type and payload; identifiers, context, aggregate versions, and saga origins are added when the workflow needs them:

type Message<TPayload> = {
	type: string;
	aggregateId?: Identifier;
	payload: TPayload;
	context?: unknown;
};

const command: Message<{ username: string }> = {
	type: 'createUser',
	payload: { username: 'alice' }
};

Library blocks are similarly narrow. For example, a projection only needs a view and three lifecycle methods:

interface Projection<TView> {
	readonly view: TView;
	subscribe(eventStore: IObservable): void | Promise<void>;
	restore(eventStore: IEventStorageReader): void | Promise<void>;
	project(event: IEvent): void | Promise<void>;
}

Applications can implement these contracts directly or extend the supplied base classes:

Contract Replace it to customize
ICommandBus Command transport and routing
IEventBus Event broadcast and worker queues
IEventStorageReader Aggregate, saga, and projection event reads
IDispatchPipelineProcessor Persistence, encoding, validation, or event augmentation
IProjection Projection routing and view ownership
IViewLocker Projection restore coordination
IEventTracker Event deduplication, projection checkpoints, and awaiting projected events; replaces the deprecated IEventLocker

The framework-free example implements the core interfaces without using the supplied aggregate or projection base classes.

Aggregates

AbstractAggregate maps public method names to command types. This aggregate handles a createUser command and emits a userCreated event:

class UserAggregate extends AbstractAggregate {
	createUser(payload: { username: string }) {
		this.emit('userCreated', { username: payload.username });
	}
}

Override static handles when command types should be declared explicitly.

Aggregate State

State is rebuilt by applying the aggregate’s historical events. Keep mutation deterministic and validate in the command method before emitting a new event:

class UserState {
	username!: string;

	userCreated(event: IEvent<{ username: string }>) {
		this.username = event.payload.username;
	}

	userRenamed(event: IEvent<{ username: string }>) {
		this.username = event.payload.username;
	}
}

class UserAggregate extends AbstractAggregate<UserState> {
	protected readonly state = new UserState();

	renameUser(payload: { username: string }) {
		if (payload.username === this.state.username)
			throw new Error('Username is unchanged');

		this.emit('userRenamed', payload);
	}
}

Constructor dependencies are resolved from the container, so domain behavior can use application services without service locators:

type UserAggregateOptions = IAggregateConstructorParams<void> & {
	authService?: AuthService;
};

class UserAggregate extends AbstractAggregate {
	readonly #authService: AuthService;

	constructor({ authService, ...options }: UserAggregateOptions) {
		super(options);
		if (!authService)
			throw new TypeError('authService is required');

		this.#authService = authService;
	}
}

interface AggregateContainer extends IContainer {
	authService: AuthService;
}

const builder = new ContainerBuilder<AggregateContainer>();
builder.register(AuthService).as('authService');
builder.registerAggregate(UserAggregate);

Projections And Views

AbstractProjection maps event types to methods in the same way:

class UsersProjection extends AbstractProjection<Map<Identifier, UserRecord>> {
	constructor() {
		super({ view: new Map() });
	}

	userCreated(event: IEvent<UserRecord>) {
		this.view.set(event.aggregateId!, event.payload);
	}
}

Override static handles to declare event types explicitly.

Expose a projection view through the typed container and wait for startup restoration before serving reads:

interface AppContainer extends IContainer {
	usersView: Map<Identifier, UserRecord>;
}

const builder = new ContainerBuilder<AppContainer>();
builder.registerProjection(UsersProjection, 'usersView');

const container = builder.container();
await Promise.all(container.restorePromises ?? []);

const usersView = container.usersView;

Persistent projection implementations provide restore locking, event deduplication, and checkpoints. Their exact transaction and retry guarantees are documented by each infrastructure module.

Awaiting Projected Events

Projections update their views asynchronously: a command resolves once its events are stored, while projections process them shortly after. Most of the time this is unnoticeable, but sometimes the next step depends on a view being up to date. For example, an API handler returns the record it has just created, a receptor sends an email using data from another projection’s view, or one projection reads another one’s view while handling the same event.

In such cases, wait until the projection has processed the events in question with eventTracker.waitFor(eventIds). Projections backed by the SQLite, PostgreSQL, MongoDB, or Redis views provide an event tracker; expose it on the container next to the view:

interface AppContainer extends IContainer {
	usersView: SqliteObjectView<UserRecord>;
	usersViewTracker: IEventTracker;
}

builder.registerProjection(UsersProjection, 'usersView')
	.exposes(p => p.eventTracker, 'usersViewTracker');

const events: IEventSet = await container.commandBus.send('createUser', undefined, { payload });
const userCreated = events.find(e => e.type === 'userCreated')!;
await container.usersViewTracker.waitFor(userCreated.id!, { timeout: 5_000 });

const user = await container.usersView.get(userCreated.aggregateId!);

waitFor accepts one or more event IDs. Pass IDs of events the projection handles: waiting for any other event lasts until the timeout. It rejects when the projection of any of the events fails, on timeout, or when the signal is aborted. It guarantees completion only: by the time it resolves, the view may already reflect subsequent events.

The timeout can be set per call with the timeout option. Defaults are stored in static properties of EventProgressTracker, used by the event trackers of the SQLite, PostgreSQL, MongoDB, and Redis views. Change them before views are created:

EventProgressTracker.DEFAULT_TIMEOUT = 30_000;          // wait timeout, in milliseconds; `0` disables it
EventProgressTracker.DEFAULT_POLL_INTERVAL = 50;        // initial interval of polling events projected by other processes, in milliseconds
EventProgressTracker.DEFAULT_MAX_POLL_INTERVAL = 1_000; // maximum polling interval the initial one grows to, in milliseconds

The in-memory and RabbitMQ event buses run all handlers of an event concurrently, so a handler waiting for another projection does not prevent that projection from processing the same event. When one projection needs data at the exact state of an event, handling that event in the projection itself is usually simpler than waiting for another one.

Events projected in the same process resolve waits immediately. Events projected by other processes sharing the same storage are found by checking the event lock table once the wait starts, then polling it with an interval doubling from DEFAULT_POLL_INTERVAL up to DEFAULT_MAX_POLL_INTERVAL while waits are pending. Projection failures reject waits only in the process where they occur; other processes keep waiting until the event is projected or the wait times out.

projection.eventTracker is null when the projection has no event tracker, for example with the default in-memory view. Its type follows the view type: it is non-nullable when the view implements IEventTracker, so exposing it under a non-nullable container alias type-checks only for such projections.

Sagas

Sagas coordinate multi-step work by handling events and producing follow-up commands:

class WelcomeEmailSaga extends AbstractSaga {
	userSignedUp(event: IEvent<{ email: string }>) {
		this.enqueue('sendWelcomeEmail', undefined, {
			email: event.payload.email
		});
	}
}

builder.registerSaga(WelcomeEmailSaga);

Starter events use event.id as the saga origin; the default dispatch pipeline assigns missing IDs automatically.

By default, a saga starts when a handled event has no origin for that saga type. Use static startsWith for explicit starter event types, static handles for additional events, and static sagaDescriptor for a stable origin key independent of the class name.

The simple saga and overlapping sagas demonstrate state restoration and origin propagation.

Event Dispatch Pipeline

Before publishing events, EventStore runs a pipeline for cross-cutting work. ContainerBuilder provides defaults that assign missing event IDs, persist events, and save snapshots when the corresponding storage processors are registered.

Extend the defaults with another processor:

builder.register(c => [
	...c.defaultEventDispatchPipeline,
	c.createInstance(AuditProcessor)
]).as('eventDispatchPipeline');

Omit the defaults to replace the pipeline entirely:

builder.register(c => [
	c.eventIdAugmenter,
	c.createInstance(CustomStorageWriter)
]).as('eventDispatchPipeline');

Pipeline registrations replace earlier registrations, so the last one wins. A replacement must include every required processor, including eventIdAugmenter when consumers need event IDs. Extend defaultEventDispatchPipeline, not eventDispatchPipeline, from its own factory to avoid a circular dependency. Named pipelines supplied through eventDispatchPipelines are also explicit and do not inherit the defaults.

eventIdAugmenter only fills in missing IDs, keeping IDs already assigned to events and IDs returned by the IIdentifierProvider as they are, whether they are strings, numbers or objects. Components that need a string key, such as saga correlation, projection locks and transport metadata, stringify the ID at their own boundary and never modify the event, so object IDs must have a stable and unique string representation. Storage modules add their own requirements: MongoDB event storage needs ObjectId-compatible IDs, SQLite event storage needs GUID-compatible ones.

Runtime Lifecycle And Guarantees

Container dependencies are resolved lazily. Access each exposed projection view during startup to create its projection, subscribe it to the event store, and start restoration. Then await restorePromises before accepting requests that depend on those views. Destructuring exposed views from the container, as in the quick start, performs that initial resolution.

Event dispatch has two stages:

  1. commandBus.send() waits for command handling and the configured dispatch pipeline, including event storage.
  2. Event-bus publication runs asynchronously after pipeline processing so command throughput is not tied to every subscriber.

Use await eventStore.drain() when a caller, test, or shutdown path must wait for all currently queued event publications. A completed command does not otherwise guarantee that every projection has finished processing its events. Configure eventPublishErrorHandler when publication failures must be logged or reported; draining waits for publication attempts but does not make subscriber handling part of the storage transaction.

Infrastructure determines distributed guarantees:

  • In-memory locks and buses coordinate only one process.
  • Persistent event stores define transaction boundaries and optimistic concurrency behavior.
  • Persistent views define restore locking, event deduplication, retries, and checkpoint semantics.
  • RabbitMQ can redeliver acknowledged-late messages, so distributed handlers should be idempotent.

Review the selected module documentation before relying on a specific failure or multi-instance behavior.

Infrastructure

Choose infrastructure by deployment need. Modules can be used independently or combined.

Need Module Deployment Peer dependency
Learning, tests, and ephemeral state node-cqrs One process -
Embedded event storage and views node-cqrs/sqlite One process better-sqlite3
Distributed event storage and document views node-cqrs/mongodb Multiple instances mongodb
Distributed document projection views node-cqrs/redis Multiple instances ioredis
Transactional event storage and relational views node-cqrs/postgresql Multiple instances pg
Distributed command and event delivery node-cqrs/rabbitmq Multiple instances amqplib
CPU-intensive projections node-cqrs/workers One application process comlink

MongoDB, Redis, and PostgreSQL support is currently experimental and has not yet been validated in production. Their APIs may change in minor versions.

Event Storage

Implementation Notes
InMemoryEventStorage Data is lost on restart; intended for learning and tests
SqliteEventStorage Embedded storage for a single application process
MongoEventStorage Distributed document event storage
PostgresqlEventStorage Distributed transactional event storage

See the SQLite example, MongoDB event-storage example, and PostgreSQL example.

Advanced: Manual Composition

The container is optional. The same components can be assembled directly:

const commandBus = new InMemoryMessageBus();
const eventBus = new InMemoryMessageBus();
const eventStorage = new InMemoryEventStorage();
const eventStore = new EventStore({
	eventStorageReader: eventStorage,
	identifierProvider: eventStorage,
	eventDispatchPipeline: [eventStorage],
	eventBus
});

const aggregateHandler = new AggregateCommandHandler({
	aggregateType: UserAggregate,
	eventStore
});
aggregateHandler.subscribe(commandBus);

const projection = new UsersProjection();
projection.subscribe(eventStore);
await projection.restore(eventStore);

const [userCreated] = await commandBus.send('createUser', undefined, {
	payload: { username: 'alice' }
});
await eventStore.drain();

console.log(projection.view.get(userCreated.aggregateId!));

OpenTelemetry

Register a tracer factory to enable spans across commands, event dispatch, projections, sagas, storage adapters, and RabbitMQ transport. Install @opentelemetry/api alongside the library:

import { trace } from '@opentelemetry/api';

builder.register(() => (name: string) => trace.getTracer(`cqrs.${name}`)).as('tracerFactory');

See the telemetry example for a complete setup with exporters.

The project was inspired by Lokad.CQRS.