@uscreen.de/cqrs-kit-2
v0.17.0
Published
> CQRS Starter Kit. Eventsourcing included. Some soldering required.
Downloads
213
Keywords
Readme
cqrs-kit
CQRS Starter Kit. Eventsourcing included. Some soldering required.
Alpha-Warning: work in progress, not tested in production and subject of change in api, features and options
Abstract
Goal of this module is to adopt CQRS/ES in a straight and reusable pattern to node applications. It should not provide a framework or make strong assumptions on infrastructure. That being said, the initial release will target frameworks like fastify, express and alike to use this package within their ecosystem. And it will target mongoDB and Nats as primary defaults.
To get to a production ready setup it should provied a working example within a setup of
Requirements
- Node.js 20+
- MongoDB
- NATS server
Assumptions
- storage and message adapters COULD handle their connections, pools, reconnections within this module
- storage and message adapters COULD reuse their clients as injected from the outside
- storage and message adapters COULD be injectable (basic reference implementations inlcluded)
- this module SHOULD NOT provide a framework,
- this module SHOULD NOT require a specific framework,
- this module SHOULD NOT depend on any given framework
- this module SHOULD NOT assume any filestructure
- this module SHOULD NOT implicitly setup a domain
Aggregate Snapshots
Since 0.17.0 aggregates are rehydrated from versioned snapshots: instead of replaying the full event stream on every command, the latest snapshot seeds the state and only events after it are replayed. Snapshots are a disposable read cache — the event store remains the single source of truth, the collection can be dropped at any time, and every failure path (missing, malformed, incompatible or unreadable snapshot) silently falls back to full replay.
Behavior change on upgrade: snapshotting is enabled by default. Each consumer automatically gets an additional collection (AggregateSnapshots by default) and starts writing snapshots once an aggregate exceeds the threshold. No code changes required; rollback = set the kill-switch or drop the collection.
How it works
- Read path: before handling a command, the latest snapshot for
(aggregateId, aggregateName)is loaded and validated (schemaVersionmust match the aggregate'ssnapshotSchemaVersion). Valid → state is seeded from it and only events withaggregateVersiongreater than the snapshot's are replayed. Invalid → full replay, never an error. - Write path: after a successful command, a snapshot is written once the aggregate is
snapshotThresholdevents ahead of its latest usable snapshot. Strictly fire-and-forget: never awaited, never fails the command, races between instances resolve via an idempotent upsert. - Retention: the last
snapshotRetentionsnapshots per aggregate are kept, older ones are pruned on write. - Schema evolution: there is intentionally no migration logic. Bump
snapshotSchemaVersionin the aggregate definition on every breaking change to the state shape — old snapshots are then ignored and rebuilt on the next trigger.
Options (Domain(opts))
| Option | Default | Description |
| ------------------------- | ---------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| snapshotsEnabled | true | Kill-switch. false skips the snapshot store entirely and forces full replay on every path — no deployment needed. |
| snapshotThreshold | 200 | Write a snapshot once an aggregate is this many events ahead of its latest snapshot. |
| snapshotRetention | 2 | Snapshots kept per (aggregateId, aggregateName). |
| SnapshotStoreCollection | 'AggregateSnapshots' | MongoDB collection name. |
| snapshotMetricsHandler | debuglog | Called once per rehydration with { aggregateName, aggregateId, snapshotHit, snapshotVersion, eventsReplayed, durationMs }. Wire it to your metrics backend; failures are swallowed. |
Per-aggregate options (Aggregate({...}))
| Option | Default | Description |
| ----------------------- | ------------------------ | ------------------------------------------------------------------------------------------------------ |
| snapshotSchemaVersion | 1 | Bump on every breaking change to the state shape; mismatching snapshots are discarded, never migrated. |
| snapshotThreshold | opts.snapshotThreshold | Per-aggregate override of the write trigger. |
Guardrails
Domain.getAggregate(name)returns a registered aggregate instance.instance.verifySnapshot(aggregateId)rebuilds the state twice — full replay vs. snapshot + tail — and returns{ equal, snapshotVersion, replayState, snapshotState }(or{ equal: true, skipped }without a usable snapshot). Run it in CI over a sample of real aggregates to prove snapshots never change semantics; a common source of divergence is non-deterministic apply handlers (timestamps, randomness, external lookups).Domain.SnapshotStoreexposes the connected store (findLatest,saveSnapshot,removeAggregate),nullwhen disabled.
Terminology
TBD
Roadmap
- add "saga/processmanager/story"
- add "DomainRegistry"
Changelog
0.17.0
Added
- aggregate snapshots (see Aggregate Snapshots): rehydration from versioned snapshots with automatic fallback to full replay; threshold-triggered, fire-and-forget snapshot writes; retention; per-aggregate
snapshotSchemaVersionandsnapshotThreshold - new
Domainoptions:snapshotsEnabled(kill-switch),snapshotThreshold,snapshotRetention,SnapshotStoreCollection,snapshotMetricsHandler Domain.getAggregate(name)registry andinstance.verifySnapshot(aggregateId)CI guardrailDomain.SnapshotStoreexposes the connected snapshot store (nullwhen disabled)EventStore.getEventStream(aggregateId, aggregateName, { fromVersion })replays only events after the given version
Changed
- snapshotting is enabled by default: consumers get an additional
AggregateSnapshotscollection on upgrade; disable viasnapshotsEnabled: falseor roll back by dropping the collection
0.16.0
Changed
- migrate NATS client from
natsv2 to@nats-io/transport-nodev3 - require Node 20+ (
engines.node: ">=20"); v3 expects nativeglobalThis.crypto natsStatusHandler(event)now receives v3 typed status objects instead of v2 enum-coded objectsnatsErrorHandler(error)now receives v3 error classes (RequestError,TimeoutError,NoRespondersError, etc.) with a.causechain, instead ofNatsErrorwitherror.code
0.15.1
- fixed
commit()hanging when a command handler threw a non-E11000 error
0.15.0
- upgrade packages
0.14.0
- switch to ESM
- switch to pnpm
- replace tap with native node tests
- upgrade packages
0.13.0
- added optional
queryparameter toProjection.rebuildmethod.
0.12.0
- use
util.debuglog()for debug output docs - add
Projection.onAfterhook to provide handler result as 3rd parameter, ie.onAfter: (event, payload, result) => {}
0.11.0
- return result of command.emit locally and remote
0.10.0
- add
Projection.onAfterhook - update aggregated state after event is emitted
0.9.0
- allow commands to emit multiple events
0.8.0
- add method
EventStore.removeAggregate
0.5.0
Added
- set
Projection.collectiontofalseto skip collection creation
Changed
- dep [email protected] replaced deprecated
returnOriginal: falsewithreturnDocument: 'after'
0.1.0
Changed
store.updatenow defaults to upsert:false. Replace withstore.createOrUpdatewhich defaults to upsert:true
Tools:
Those are helpfull for development but should not be added as dependencies:
npx nats-cli- listen to nats channelsnpx natsboard- monitoring dashboard
