@cqrs-ddd/core
Framework-neutral building blocks for domain-driven design in TypeScript: aggregates with versioned mutations, detached domain events, a command handler that publishes those events, repository contracts, persistence lifecycle decorators, a revision-fenced repository cache, tenant-scoped cache keys and HTTP status mapping for its errors.
It depends on no framework and no ORM. MikroORM adapters are in
@cqrs-ddd/mikro-orm. It works in a plain Node service, and in a NestJS
application through structural compatibility: see Using it from NestJS.
Contents
Section titled “Contents”- Installation
- Entry points
- Aggregates
- Value rules
- Domain events
- Commands and event publication
- Write-side repositories
- Read-side repositories
- The repository cache
- Tenant-scoped cache keys
- HTTP status mapping
- Using it from NestJS
- Moving from ddd/core
- Known limits
- License
Installation
Section titled “Installation”pnpm add @cqrs-ddd/core# ornpm install @cqrs-ddd/coreRequires Node.js 22.12 or later. It installs @cqrs-ddd/uuidv7 and
@cqrs-ddd/safe-stringify, which have no dependencies. To persist with MikroORM, add
@cqrs-ddd/mikro-orm and @mikro-orm/core 7.
Published as an ES module; a CommonJS application loads it with require(). Coming from
0.3.x, see Upgrading from 0.3.x.
Entry points
Section titled “Entry points”Import from the entry point for the layer you are writing. Each one loads only what it needs.
| Entry | Holds |
|---|---|
@cqrs-ddd/core/domain |
AggregateRoot, IAggregateRoot, ApplyEventOptions, RootEntity, RootEntitySnapshot, @Mutable, getMutableFields, @ApplyMutation, IEvent, DomainEvent, RootDomainEvent, deepCloneAndFreeze, textRule, numberRule, ValueViolation, and the errors DomainException, InvalidValueException, EntityNotFoundException, ConcurrencyConflictError, TransientOperationError (with isTransientOperationError), MissingTenantContextError, UnknownMutableFieldError |
@cqrs-ddd/core/application |
BaseCommand, BaseQuery, IQueryOptions, CommandBaseHandler, the ports (IDomainEventPublisher, ICommandRepository, IQueryRepository, IWriteSideAggregateRepository, ICache, IVersionedCache, isVersionedCache), requireTenant, setTenantResolver, and the deprecated requireTenantId |
@cqrs-ddd/core/persistence |
the lifecycle decorators (@PersistedWrite, @Cache, @AcknowledgePersisted, @MapPersistenceErrors, @FromCache), QueryRepository, CommandRepository, MemoryCache and its injection token CACHE_TOKEN, cacheKey, cacheKeyTemplate, isCacheNewer, toCacheSnapshot, the mutation-barrier helpers, consoleCacheLogger, safeWarn, and the persistence dialect contract (IPersistenceDialect, setPersistenceDialect, persistenceDialect) |
@cqrs-ddd/core/http |
domainErrorHttpStatus |
@cqrs-ddd/core |
all of the above, plus the Method type |
No entry point loads an ORM or NestJS.
Aggregates
Section titled “Aggregates”An aggregate extends RootEntity<TSnapshot>. It gets a UUIDv7 id, createdAt and
updatedAt, an optimistic-concurrency version, and an event buffer. State changes go
through domain methods decorated with @ApplyMutation, which write only the fields
declared @Mutable:
import { ApplyMutation, DomainException, InvalidValueException, Mutable, RootEntity, type RootEntitySnapshot, textRule,} from '@cqrs-ddd/core/domain';
export class InvalidUsernameException extends InvalidValueException {}export class InvalidDepartmentException extends InvalidValueException {}export class EmptyUserUpdateException extends DomainException { constructor() { super('A user update needs at least one field.'); }}
export interface UserSnapshot extends Partial<RootEntitySnapshot> { readonly username: string; readonly email: string; readonly department?: string | null;}
export class User extends RootEntity<UserSnapshot> { static readonly aggregateName = 'user'; static readonly rules = { username: textRule({ field: 'username', minLength: 3, maxLength: 255, error: (violation) => new InvalidUsernameException(violation), }), department: textRule({ field: 'department', required: false, maxLength: 255, error: (violation) => new InvalidDepartmentException(violation), }), } as const;
@Mutable<string>({ normalize: (value) => User.rules.username.parse(value) }) private _username: string;
@Mutable<string | null>({ normalize: (value) => User.rules.department.parse(value) }) private _department: string | null;
readonly email: string;
private constructor(snapshot: UserSnapshot) { super(snapshot); this._username = User.rules.username.parse(snapshot.username); this._department = User.rules.department.parse(snapshot.department); this.email = snapshot.email; }
static create(username: string, email: string, department?: string | null): User { const user = new User({ username, email, department }); user.apply(new UserCreatedEvent(user)); return user; }
static fromJSON(snapshot: UserSnapshot): User { return new User(snapshot); }
get username(): string { return this._username; }
private set username(value: string) { this._username = User.rules.username.parse(value); }
get department(): string | null { return this._department; }
private set department(value: string | null) { this._department = User.rules.department.parse(value); }
@ApplyMutation<User>({ event: (user) => new UserUpdatedEvent(user) }) update(fields: { username?: string; department?: string | null }): this { if (fields.username === undefined && fields.department === undefined) { throw new EmptyUserUpdateException(); } this.applyPatch({ username: fields.username, department: fields.department }); return this; }
@ApplyMutation<User>({ event: (user) => new UserDeletedEvent(user) }) delete(): this { return this; }
toJSON(): UserSnapshot & RootEntitySnapshot { return this.freezeState({ id: this.id, username: this._username, department: this._department, email: this.email, createdAt: this.createdAt, updatedAt: this.updatedAt, version: this.version, }); }}The events are declared in Domain events. Using the aggregate:
const user = User.create('alice', 'alice@example.test');user.version; // 1user.getUncommittedEvents(); // [UserCreatedEvent]
user.update({ department: 'Research' });user.version; // 2user.getExpectedVersion(); // 1: the version the next write must find in storage
user.update({}); // throws EmptyUserUpdateException; nothing is recordeduser.update({ username: 'a' }); // throws InvalidUsernameException before any field is written-
@Mutableregisters a field forapplyPatchunder its property name without a leading underscore (_usernameis patched asusername);as: 'name'sets another key, andnormalizevalidates and converts each value. -
applyPatch(patch)validates and normalizes every supplied field before writing any of them, and ignoresundefinedvalues. An unknown key throwsUnknownMutableFieldError. -
After a successful method,
@ApplyMutationadvancesversionandupdatedAtand records the event. A method that throws records nothing. -
Validate before patching: a failure after the fields are written (later method code, a lifecycle hook, event creation) does not roll them back. Discard or reload the instance.
-
RootEntity.from(value)rehydrates an instance, a snapshot or a nullish database result, and throws aTypeErrorfor an incompatible aggregate. -
id,createdAtandupdatedAthave public getters and private setters. The ORM hydrates through the setters (see below); application code cannot assign them and changes state through domain methods and factories. Give a subclass’s own persisted fields private setters too, asusernameanddepartmentabove. -
versionis the in-memory version;getExpectedVersion()is the version last read from or written to storage, which a version-conditioned write compares against.acknowledgePersisted()moves it forward after a successful write;@PersistedWritecalls it for you.
Value rules
Section titled “Value rules”textRule(options) and numberRule(options) define the rule for one field once. The
aggregate uses the rule’s parse(value) to normalize the value, and other layers read
the rule’s limits, so a request schema cannot drift from the domain:
z.string().trim().min(User.rules.username.minLength).max(User.rules.username.maxLength);parse(value) returns the normalized value or throws the rule’s error. It is a
standalone function, and the rule is frozen.
textRule option |
Effect |
|---|---|
field |
Name reported in the violation and in the default message |
required |
false normalizes null, undefined and blank text to null; by default they are rejected |
minLength, maxLength |
Length after normalization, in UTF-16 code units like Zod’s .min() and .max() |
pattern |
Regular expression the text must match; a g or y flag is rejected |
multiline |
true accepts tab, line feed and carriage return inside the text |
error |
Builds the error for a violation; defaults to InvalidValueException |
Text is checked for type and unpaired surrogates, then put in Unicode NFC form and
trimmed. Control characters are always rejected, and so are tabs and line breaks unless
multiline is set.
numberRule option |
Effect |
|---|---|
field, required, error |
As for textRule (required: false accepts null and undefined) |
integer |
Accepts safe integers only |
min, max |
Inclusive bounds |
A number must be of type number and finite; a numeric string is rejected, not
converted. An invalid option, such as minLength above maxLength, throws a
TypeError when the rule is defined.
A broken rule throws error(violation). violation is a frozen ValueViolation:
{ field, rule }, plus limit for a length or range rule and expected for a type
rule. InvalidValueException carries it, builds the message (username must be at least 3 characters.), and does not keep the rejected value. Extend it for your own
exceptions, as in the example above, so every one shares that shape.
Domain events
Section titled “Domain events”A domain event extends RootDomainEvent<TEntity, TPayload>. It carries a UUIDv7 id,
the event-time aggregateId and aggregateVersion, and payload: a deep clone of the
aggregate’s toJSON(), recursively frozen.
import { RootDomainEvent } from '@cqrs-ddd/core/domain';
export class UserCreatedEvent extends RootDomainEvent<User, UserSnapshot> { constructor(user: User) { super(user); }}UserUpdatedEvent and UserDeletedEvent are declared the same way.
- Detached: later changes to the aggregate never reach
payload. - Read-only for ordinary access: own properties are frozen, and the mutating methods
of cloned
Date,MapandSetvalues throw. - Not a security boundary:
Object.freezedoes not cover built-in internal state, soDate.prototype.setTime.call(payload.at, 0)still changes it, and every consumer of one event shares the samepayload.event.clonePayload()returns a fresh frozen copy per call when a consumer needs its own. - Not a transport format: send an explicit, serializable message built from the fields you need.
deepCloneAndFreeze(value) is the function behind payload, also usable on its own. It
handles objects, arrays, Date, Map, Set, RegExp and circular references.
Commands and event publication
Section titled “Commands and event publication”BaseCommand and BaseQuery are plain base classes for command and query objects.
BaseCommand keeps an optional sessionPrincipal (or sessionUser) out of enumeration and toJSON(), so it
never reaches a fingerprint or a log. getUpdateFields(fields) returns the listed
fields the command carries (null counts, undefined does not), for field-level
authorization.
CommandBaseHandler runs a command and publishes the events of the aggregate it
returns:
import { BaseCommand, CommandBaseHandler, type IDomainEventPublisher, type IWriteSideAggregateRepository,} from '@cqrs-ddd/core/application';import { EntityNotFoundException } from '@cqrs-ddd/core/domain';
export class UpdateUserCommand extends BaseCommand { constructor( readonly id: string, readonly username?: string, readonly department?: string | null, ) { super(); }}
export class UpdateUserHandler extends CommandBaseHandler<UpdateUserCommand, User> { constructor( private readonly users: IWriteSideAggregateRepository<User>, eventBus: IDomainEventPublisher, ) { super(eventBus); }
async handle(command: UpdateUserCommand): Promise<User> { const user = await this.users.findById(command.id); if (!user) throw new EntityNotFoundException('User', command.id); user.update({ username: command.username, department: command.department }); await this.users.save(user); return user; }}
const handler = new UpdateUserHandler(users, { publishAll: async (events) => events.forEach((event) => emitter.emit('event', event)),});await handler.execute(new UpdateUserCommand(id, 'bob'));handle()returns the aggregate, or a result carrying it asaggregate; anyIAggregateRootqualifies.execute()calls it, publishes the buffered events once witheventBus.publishAll(events, aggregate)(the aggregate is the dispatcher context), clears the buffer and returns the result unchanged.- A rejected
handle()publishes nothing. IfpublishAll()throws, the error propagates and the events stay buffered. - If
publishAll()returns a promise,execute()waits for it after clearing the buffer. Its rejection rejects the command, although the aggregate is already persisted; the events are not published a second time. - Any object with
publishAll(events)is anIDomainEventPublisher.
AggregateRoot.commit() publishes nothing by default. It hands a copy of the
buffer to the aggregate’s own publishAll hook (apply() calls publish while
autoCommit is on), which does nothing until a publisher is connected, then clears the
buffer. Calling commit() without a connected publisher therefore drops the events; so
does autoCommit. Let CommandBaseHandler publish, or read getUncommittedEvents() and
call uncommit() yourself.
commit(dispatcherContext?) passes the context, such as { transaction }, to
publishAll and returns what it returns. The buffer is cleared once publishAll
returns, before an asynchronous publisher settles, so await aggregate.commit() waits
for the publisher and receives its error without risking a second publication; if
publishAll throws, the events stay buffered.
Write-side repositories
Section titled “Write-side repositories”A command handler loads the authoritative aggregate through
IWriteSideAggregateRepository<TEntity>.findById(id): it must read primary storage, never
a cache. @cqrs-ddd/mikro-orm’s AggregateRepository implements it for MikroORM.
save() declares its lifecycle with @PersistedWrite. With @cqrs-ddd/mikro-orm, where
AggregateRepository provides findById() and this.store:
import type { ICache } from '@cqrs-ddd/core/application';import { DomainException } from '@cqrs-ddd/core/domain';import { cacheKey, PersistedWrite } from '@cqrs-ddd/core/persistence';import { AggregateRepository, type IEntityManagerSource, mapPersistenceError, optimisticUpdate,} from '@cqrs-ddd/mikro-orm';
export class UniqueEmailException extends DomainException { constructor(readonly user: User) { super('This email is already registered.'); }}
export class UpdateUserRepository extends AggregateRepository<UserSnapshot, User, UserSnapshot> { constructor(cache: ICache<UserSnapshot>, store: IEntityManagerSource) { super(cache, store, User, User.aggregateName, User.fromJSON); }
@PersistedWrite<User>({ cache: { setKey: (user) => cacheKey(User.aggregateName, { id: user.id }), invalidateKeys: (user) => [cacheKey(User.aggregateName, { email: user.email })], }, unique: { email: (user) => new UniqueEmailException(user) }, otherwise: (error, user) => mapPersistenceError(error, `updating User ${user.id}`), }) async save(user: User): Promise<UserSnapshot> { const snapshot = user.toJSON(); await optimisticUpdate( this.store.em, User, user, { username: snapshot.username, department: snapshot.department ?? null, updatedAt: snapshot.updatedAt, }, 'User', ); return snapshot; }}optimisticUpdate adds version to the written fields; every other changed column,
updatedAt included, is listed explicitly.
@PersistedWrite applies three decorators. Stacked by hand, the order is fixed,
outermost first:
@Cache(options): after a successful write, writes the returned snapshot through with a compare-and-set (isCacheNewerunlessisNewersays otherwise), and installs barriers forinvalidateKeys. Whensave()resolvesnullorundefined, it installs barriers fordeleteKeys. Cache errors are logged and never turn a committed write into a failure.@AcknowledgePersisted({ entity: ([aggregate]) => aggregate }): advances the selected aggregate’s persisted version only after the write succeeded.@MapPersistenceErrors({ entity, unique, dialect, otherwise }): translates a unique violation into your domain error.uniqueis keyed by the entity property the constraint covers ({ email: ... }), or by the declared name of a multi-column constraint (MapPersistenceErrors<Args, Entity, typeof NAME>), so a repository never names a database constraint or column. Which constraint an error violated is read by the persistence dialect:options.dialect, else the one registered withsetPersistenceDialect()at the composition root. Core knows no database; an adapter such asMikroOrmDialectof@cqrs-ddd/mikro-ormimplementsIPersistenceDialect. A method that declaresuniquewith no dialect available throws aTypeErrorbefore it runs.otherwisetranslates the rest, for example withmapPersistenceErrorof@cqrs-ddd/mikro-orm, which turns retryable driver and network failures intoTransientOperationErrorand keeps every other error unchanged.
Deletes do not acknowledge, so they stack @Cache and @MapPersistenceErrors alone,
and resolve null so @Cache installs barriers for deleteKeys:
import { Cache, cacheKey, MapPersistenceErrors } from '@cqrs-ddd/core/persistence';import { mapPersistenceError, optimisticDelete } from '@cqrs-ddd/mikro-orm';
export class DeleteUserRepository extends AggregateRepository<UserSnapshot, User, null> { constructor(cache: ICache<UserSnapshot>, store: IEntityManagerSource) { super(cache, store, User, User.aggregateName, User.fromJSON); }
@Cache<User, null>({ deleteKeys: (user) => [ cacheKey(User.aggregateName, { id: user.id }), cacheKey(User.aggregateName, { email: user.email }), ], }) @MapPersistenceErrors<[User], User>({ entity: ([user]) => user, otherwise: (error, user) => mapPersistenceError(error, `deleting User ${user.id}`), }) async save(user: User): Promise<null> { await optimisticDelete(this.store.em, User, user, 'User'); return null; }}The write itself must be version-conditioned (WHERE id = ? AND version = expected) and
must not run inside an outer transaction, whose success would not mean durable
persistence. @cqrs-ddd/mikro-orm provides optimisticUpdate and optimisticDelete
for this.
Read-side repositories
Section titled “Read-side repositories”A query repository extends QueryRepository and decorates find() with @FromCache. The
constructor takes the cache (any ICache; see The repository cache
for why @FromCache needs an IVersionedCache) and the repository’s default hydration:
import { BaseQuery, type ICache, type IQueryOptions } from '@cqrs-ddd/core/application';import { cacheKey, FromCache, QueryRepository } from '@cqrs-ddd/core/persistence';import type { IEntityManagerSource } from '@cqrs-ddd/mikro-orm';
export class GetUserQuery extends BaseQuery { constructor( readonly id: string, options?: IQueryOptions, ) { super(options); }}
export class GetUserRepository extends QueryRepository<GetUserQuery, User | null> { constructor( cache: ICache<UserSnapshot>, private readonly store: IEntityManagerSource, ) { super(cache, { hydrateFn: (cached) => User.fromJSON(cached as UserSnapshot) }); }
@FromCache<GetUserQuery, User | null>({ keyFn: (query) => cacheKey(User.aggregateName, { id: query.id }), }) async find(query: GetUserQuery): Promise<User | null> { return this.store.em.findOne(User, { id: query.id }, query.refresh ? { refresh: true } : undefined); }}
await users.find(new GetUserQuery(id)); // may be served from the cacheawait users.find(new GetUserQuery(id, { refresh: true })); // always reads the database@FromCache options: keyFn (return null to skip the cache for that call),
hydrateFn, serializeFn, ttl in milliseconds, isNewer and logger.
- Every hit is rehydrated whenever a hydrator applies, so a hit and a miss return the
same type. A method’s own
hydrateFn(orhydrateFn: null) overrides the repository’s;serializeFnfollows the same rule. - Without a hydrator, only plain data is cached; a class instance is returned uncached.
- Only non-nullish results are stored. A barrier or a stored
nullcounts as a miss. - A query with
refresh: truebypasses the cache. - Cache errors propagate from reads: a query fails rather than guessing.
The repository cache
Section titled “The repository cache”@Cache works with any ICache (get, set, delete). @FromCache requires an
IVersionedCache, which adds a revision fence:
| Method | Contract |
|---|---|
readState(key) |
the entry’s status (hit, miss, expired), value and opaque revision |
invalidate(key) |
evicts the value and advances the revision |
tryFill(key, observedRevision, value, options?) |
writes only if the revision is still observedRevision; false otherwise |
@FromCache observes the revision before the database read and fills with tryFill, so
a read that raced an invalidation or a delete cannot write a stale snapshot back. A
rejected fill re-reads the key, returns a strictly newer snapshot if one is there, and
otherwise retries a bounded number of times before returning the database result
uncached. An adapter that implements only ICache is bypassed by @FromCache for
both reads and fills, with a warning logged once: it cannot fence a fill.
Two adapters implement IVersionedCache:
MemoryCache({ defaultTtlMs = 60_000, maxEntries = 10_000 }): in process, entries JSON-cloned on write and read, for tests and a single process.MikroOrmCache, in@cqrs-ddd/mikro-orm: one database row per key, every write a compare-and-set in its own transaction.
CacheSetOptions.isNewer(cached, incoming) returns true when the cached value is
newer and must not be overwritten; isCacheNewer compares version, then __gen, then
updatedAt. toCacheSnapshot turns a result into a detached snapshot: the cache never
stores a live aggregate.
Mutation barriers (CacheMutationBarrier) mark a deleted or invalidated key for
barrierTtl, 60 seconds by default. Size it above your longest in-flight read.
The logger option. @Cache and @FromCache accept logger: { warn(message) }
for their operational warnings, such as a failed cache write, a bypassed adapter or a
miswired repository. A NestJS Logger or console fits. Without one, warnings go to
console.warn. A logger that throws never changes a result. An adapter reports its own
warnings the same way through consoleCacheLogger(context) and safeWarn(logger, message).
Tenant-scoped cache keys
Section titled “Tenant-scoped cache keys”cacheKey(resource, conditions, tenant?) and cacheKeyTemplate(template, tenant?)
namespace every key by tenant. The tenant comes from, in order:
- the explicit
tenantargument: a tenant id string, or an object carryingtenantId, such as a request context; - for
cacheKeyTemplate, a context source’stenantId; - the resolver the application registers once at startup with
setTenantResolver(fn), a function that returns the tenant of the running request or job; - otherwise,
MissingTenantContextError.
import { AsyncLocalStorage } from 'node:async_hooks';import { setTenantResolver } from '@cqrs-ddd/core/application';
const requestTenant = new AsyncLocalStorage<string>();setTenantResolver(() => requestTenant.getStore());
// Per request or job:requestTenant.run(tenantId, () => handle(request));There is no shared namespace: a missing or empty tenant, or a resolver that returns
nothing, throws. A single-tenant deployment passes a fixed id or registers
setTenantResolver(() => 'single'). The same resolution is public as
requireTenant(purpose, source?), for any other key that must be partitioned by
tenant, such as an idempotency or rate-limit key, and for any code that needs the
current tenant and must fail when there is none:
import { requireTenant } from '@cqrs-ddd/core/application';
requireTenant('access token issuance'); // the resolver's tenantrequireTenant('users.create idempotency key', ctx); // ctx.tenantId, else the resolver'srequireTenantId(source, purpose) is deprecated: it takes the same arguments in the
reverse order and calls requireTenant.
cacheKey keys have the form ${tenant}:${resource}:v1:${sha256}, hashed over a
key-sorted serialization of [tenant, resource, conditions]. The format is frozen, so
stored entries stay addressable across releases.
HTTP status mapping
Section titled “HTTP status mapping”domainErrorHttpStatus(error) from /http maps this package’s errors to plain
{ statusCode, error, message } values, for any HTTP framework:
| Error | Status |
|---|---|
ConcurrencyConflictError |
409 |
EntityNotFoundException |
404 |
MissingTenantContextError |
500, with a generic message: a request without a tenant is a server fault |
any other DomainException, InvalidValueException included |
400 |
It returns undefined for anything else. Map your own exceptions first, then pass the
rest to it.
Using it from NestJS
Section titled “Using it from NestJS”Nothing in this package imports NestJS; the fit is structural.
- The NestJS CQRS
EventBushaspublishAll(events, dispatcherContext), so it is anIDomainEventPublisher: a handler passes its injected bus tosuper(eventBus), and the bus hands the aggregate to its configured publisher as the dispatcher context. AggregateRootsatisfies NestJS’sIAggregateRoot, and a NestJSAggregateRootsatisfies this package’sIAggregateRoot.EventPublisher.mergeObjectContext(aggregate)connects an aggregate to theEventBus:commit()then publishes with the aggregate as the context,commit({ transaction })with that context, and both return the bus’s result.CommandBaseHandler.execute(command)is the method the NestJS command bus calls, so a class decorated with@CommandHandlerextends it directly.BaseCommand,BaseQueryand the domain events are plain classes that NestJS dispatches like any other.
@CommandHandler(UpdateUserCommand)export class UpdateUserHandler extends CommandBaseHandler<UpdateUserCommand, User> { constructor( @Inject(USER_WRITE_REPOSITORY) private readonly users: IWriteSideAggregateRepository<User>, eventBus: EventBus, ) { super(eventBus); } // handle() as above}The application writes the rest of the glue:
-
providers for the repositories and the cache, for example a
useFactoryprovider that builds aMikroOrmCachefrom@cqrs-ddd/mikro-ormunderCACHE_TOKEN:{ provide: CACHE_TOKEN, useValue: new MemoryCache({ defaultTtlMs: 30_000 }) } -
the tenant resolver, registered once before any lifecycle hook can start work that reads it, for example in a module constructor;
-
an exception filter that maps its own errors, then calls
domainErrorHttpStatus; -
a
loggerfor the cache decorators, such asnew Logger('UserCache').
Moving from ddd/core
Section titled “Moving from ddd/core”The unpublished workspace package @nestjs-pipeline/ddd-core (directory ddd/core) maps
onto this package as follows. It depended on NestJS and MikroORM; this package depends on
neither, and the MikroORM parts are in @cqrs-ddd/mikro-orm.
ddd/core |
@cqrs-ddd/core |
|---|---|
import { … } from '@nestjs-pipeline/ddd-core' |
the layer entry points: @cqrs-ddd/core/domain, /application, /persistence, /http |
@Mutate() on a mutation method, which called onUpdate() |
@ApplyMutation({ event }), which also records the event, plus @Mutable fields written through applyPatch() |
DomainOutcome, RootDomainOutcome returned by handle() |
handle() returns the aggregate, or a result with an aggregate property; its buffered events are published |
CommandBaseHandler(eventBus: EventBus) from @nestjs/cqrs |
CommandBaseHandler(eventBus: IDomainEventPublisher); a NestJS EventBus still fits |
CacheableEntity, ICacheKey, entity.cacheKey (prefix + id) |
RootEntity with a static aggregateName, keys built with cacheKey(User.aggregateName, { id }), always tenant-scoped |
CacheableEntity.fromStringify(data, fromJSON) |
RootEntity.from(value) or the aggregate’s own fromJSON |
ICommandRepository.save(domainOutcome) |
ICommandRepository.save(entity), and IWriteSideAggregateRepository.findById(id) for loading |
@Cache(setKeyFn, deleteKeysFn) |
@Cache({ setKey, deleteKeys, invalidateKeys, ttl, isNewer, barrierTtl, logger }), usually through @PersistedWrite |
@FromCache(keyFn, hydrateFn), hydrating when query.hydrate is set |
@FromCache({ keyFn, hydrateFn, … }) with an IVersionedCache; hits are always rehydrated when a hydrator applies |
QueryRepository(cache) |
QueryRepository(cache, { hydrateFn, serializeFn? }) |
ICache in persistence/cache.interface |
ICache and IVersionedCache in @cqrs-ddd/core/application |
UnixTimestampType |
UnixTimestampType in @cqrs-ddd/mikro-orm |
RootEntity with an abstract afterUpdate() |
afterUpdate() is an optional hook; RootEntity extends AggregateRoot and adds version, getExpectedVersion(), the event buffer and private hydration setters |
The old imports and the new ones side by side:
import { CacheableEntity, CommandBaseHandler, Mutate, RootDomainOutcome } from '@nestjs-pipeline/ddd-core';
// @cqrs-ddd/coreimport { ApplyMutation, Mutable, RootEntity } from '@cqrs-ddd/core/domain';import { CommandBaseHandler } from '@cqrs-ddd/core/application';import { cacheKey, PersistedWrite } from '@cqrs-ddd/core/persistence';import { AggregateRepository, optimisticUpdate, UnixTimestampType } from '@cqrs-ddd/mikro-orm';Known limits
Section titled “Known limits”- Persistence and event publication are not atomic. A crash between the write and
publishAll()loses the events; durable delivery needs an outbox. - Cache maintenance after a commit is best-effort. If an invalidation fails, the
previous entry stays until it expires: the revision fence coordinates the cache with
itself, not with the database. Read with
{ refresh: true }where a decision must not rest on a cached value. - A mutation barrier protects only for
barrierTtl. It is a bounded race window, not a tombstone. - The persistence helpers reject outer transactions; decorators cannot make several writes atomic. A row that must commit with the aggregate, such as an audit record or an outbox message, therefore cannot share its transaction yet: that needs commit hooks that run acknowledgment and cache work after the commit.
@MapPersistenceErrorsmaps a unique violation only as well as the registered persistence dialect reads it;MikroOrmDialectreads PostgreSQL and SQLite errors.- The tenant resolver is process-wide: one per process, for each installed copy of this package.
- Event payloads are frozen for ordinary access only; see Domain events.
License
Section titled “License”Dual-licensed under AGPL-3.0-or-later or a Commercial License. See
LICENSE and COMMERCIAL_LICENSE.txt
at the repository root.