Μετάβαση στο περιεχόμενο

@cqrs-ddd/core

npm version License

Δομικά στοιχεία ανεξάρτητα από framework για domain-driven design σε TypeScript: aggregates με εκδοχοποιημένα mutations (versioned mutations), αποσυνδεδεμένα domain events (detached events), ένας command handler που δημοσιεύει αυτά τα events, repository contracts, decorators κύκλου ζωής persistence, ένα revision-fenced repository cache, cache keys διαχωρισμένα ανά tenant και αντιστοίχιση HTTP status για τα domain errors.

Δεν εξαρτάται από κανένα framework και κανένα ORM. Οι adapters για το MikroORM βρίσκονται στο @cqrs-ddd/mikro-orm.

Terminal window
pnpm add @cqrs-ddd/core
# ή
npm install @cqrs-ddd/core

Απαιτεί Node.js 22.12 ή νεότερο. Εγκαθιστά τα @cqrs-ddd/uuidv7 και @cqrs-ddd/safe-stringify, τα οποία δεν έχουν εξωτερικές εξαρτήσεις. Για persistence με MikroORM, προσθέστε τα @cqrs-ddd/mikro-orm και @mikro-orm/core 7.

Δημοσιεύεται ως ES module. Μια εφαρμογή CommonJS μπορεί να το φορτώσει μέσω require().

Κάντε import από το entry point που αντιστοιχεί στο layer που υλοποιείτε. Κάθε ένα φορτώνει αποκλειστικά ό,τι χρειάζεται.

Entry point Περιεχόμενο
@cqrs-ddd/core/domain AggregateRoot, IAggregateRoot, ApplyEventOptions, RootEntity, RootEntitySnapshot, @Mutable, getMutableFields, @ApplyMutation, IEvent, DomainEvent, RootDomainEvent, deepCloneAndFreeze, textRule, numberRule, ValueViolation, και τα errors DomainException, InvalidValueException, EntityNotFoundException, ConcurrencyConflictError, TransientOperationError (με isTransientOperationError), MissingTenantContextError, UnknownMutableFieldError
@cqrs-ddd/core/application BaseCommand, BaseQuery, IQueryOptions, CommandBaseHandler, τα ports (IDomainEventPublisher, ICommandRepository, IQueryRepository, IWriteSideAggregateRepository, ICache, IVersionedCache, isVersionedCache), requireTenant, setTenantResolver, και το deprecated requireTenantId
@cqrs-ddd/core/persistence οι lifecycle decorators (@PersistedWrite, @Cache, @AcknowledgePersisted, @MapPersistenceErrors, @FromCache), QueryRepository, CommandRepository, MemoryCache και το injection token CACHE_TOKEN, cacheKey, cacheKeyTemplate, isCacheNewer, toCacheSnapshot, τα mutation-barrier helpers, consoleCacheLogger, safeWarn, και το persistence dialect contract (IPersistenceDialect, setPersistenceDialect, persistenceDialect)
@cqrs-ddd/core/http domainErrorHttpStatus
@cqrs-ddd/core όλα τα παραπάνω, συν τον τύπο Method

Κανένα entry point δεν φορτώνει ORM ή framework.

Ένα aggregate επεκτείνει το RootEntity<TSnapshot>. Αποκτά αυτόματα UUIDv7 id, createdAt και updatedAt, version για optimistic concurrency, και ένα event buffer. Οι αλλαγές κατάστασης πραγματοποιούνται μέσω domain μεθόδων διακοσμημένων με @ApplyMutation, οι οποίες τροποποιούν αποκλειστικά τα πεδία που έχουν δηλωθεί ως @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,
});
}
}

Χρήση του aggregate:

const user = User.create('alice', 'alice@example.test');
user.version; // 1
user.getUncommittedEvents(); // [UserCreatedEvent]
user.update({ department: 'Research' });
user.version; // 2
user.getExpectedVersion(); // 1: η έκδοση που πρέπει να βρει η επόμενη εγγραφή στη βάση
user.update({}); // ρίχνει EmptyUserUpdateException, τίποτα δεν καταγράφεται
user.update({ username: 'a' }); // ρίχνει InvalidUsernameException πριν γραφτεί οποιοδήποτε πεδίο
  • Το @Mutable καταχωρεί ένα πεδίο για το applyPatch υπό το όνομα του property χωρίς το αρχικό underscore (το _username ενημερώνεται ως username). Το as: 'name' ορίζει διαφορετικό κλειδί, ενώ το normalize επικυρώνει και μετατρέπει κάθε τιμή.
  • Το applyPatch(patch) επικυρώνει και κανονικοποιεί κάθε παρεχόμενο πεδίο πριν γράψει οποιοδήποτε από αυτά, και αγνοεί τιμές undefined. Ένα άγνωστο κλειδί ρίχνει UnknownMutableFieldError.
  • Μετά από επιτυχή μέθοδο, το @ApplyMutation προωθεί τα version και updatedAt και καταγράφει το event. Μια μέθοδος που ρίχνει σφάλμα δεν καταγράφει τίποτα.
  • Επικυρώστε πριν την εφαρμογή του patch: μια αποτυχία αφού γραφτούν τα πεδία δεν πραγματοποιεί rollback. Απορρίψτε ή επαναφορτώστε το instance.
  • Το RootEntity.from(value) επαναφέρει ένα instance, snapshot ή nullish αποτέλεσμα βάσης δεδομένων, και ρίχνει TypeError για ασύμβατο aggregate.
  • Τα id, createdAt και updatedAt διαθέτουν public getters και private setters. Το ORM τα αρχικοποιεί μέσω των setters. Ο κώδικας εφαρμογής αλλάζει κατάσταση μόνο μέσω domain μεθόδων και factories.
  • Το version είναι η in-memory έκδοση. Το getExpectedVersion() είναι η έκδοση που διαβάστηκε ή γράφτηκε τελευταία στη βάση, την οποία συγκρίνει μια εγγραφή με condition. Το acknowledgePersisted() την προωθεί μετά από επιτυχή εγγραφή (το @PersistedWrite το καλεί αυτόματα).

Τα textRule(options) και numberRule(options) ορίζουν τον κανόνα για ένα πεδίο μία φορά. Το aggregate χρησιμοποιεί το parse(value) του κανόνα για την κανονικοποίηση, και τα υπόλοιπα layers διαβάζουν τα όρια του κανόνα ώστε το schema αιτήματος να μην αποκλίνει από το domain:

z.string().trim().min(User.rules.username.minLength).max(User.rules.username.maxLength);

Το parse(value) επιστρέφει την κανονικοποιημένη τιμή ή ρίχνει το σφάλμα του κανόνα. Είναι ανεξάρτητη συνάρτηση και ο κανόνας είναι παγωμένος (frozen).

Επιλογή textRule Αποτέλεσμα
field Όνομα που αναφέρεται στην παραβίαση και στο προεπιλεγμένο μήνυμα
required Με false κανονικοποιεί τα null, undefined και κενό κείμενο σε null. Εξ ορισμού απορρίπτονται
minLength, maxLength Μήκος μετά την κανονικοποίηση, σε UTF-16 code units
pattern Regular expression στο οποίο πρέπει να ταιριάζει το κείμενο (flags g ή y απορρίπτονται)
multiline Με true δέχεται tab, line feed και carriage return μέσα στο κείμενο
error Κατασκευάζει το error για την παραβίαση (προεπιλογή: InvalidValueException)

Το κείμενο ελέγχεται ως προς τον τύπο και τα μη ζευγαρωμένα surrogates, μετατρέπεται σε μορφή Unicode NFC και αφαιρούνται τα κενά στα άκρα (trimmed). Χαρακτήρες ελέγχου απορρίπτονται πάντοτε.

Επιλογή numberRule Αποτέλεσμα
field, required, error Όπως στο textRule (required: false δέχεται null και undefined)
integer Δέχεται μόνο ασφαλείς ακέραιους (safe integers)
min, max Συμπεριληπτικά όρια εύρους (inclusive bounds)

Ένας αριθμός πρέπει να είναι τύπου number και πεπερασμένος (finite). Αριθμητικά strings απορρίπτονται. Μια παραβίαση κανόνα ρίχνει error(violation). Το violation είναι ένα παγωμένο ValueViolation: { field, rule }, συν limit για κανόνες μήκους/εύρους και expected για κανόνες τύπου.

Ένα domain event επεκτείνει το RootDomainEvent<TEntity, TPayload>. Φέρει UUIDv7 id, τα aggregateId και aggregateVersion κατά τον χρόνο του συμβάντος, και το payload: ένα βαθύ αντίγραφο (deep clone) του toJSON() του aggregate, αναδρομικά παγωμένο (recursively frozen).

import { RootDomainEvent } from '@cqrs-ddd/core/domain';
export class UserCreatedEvent extends RootDomainEvent<User, UserSnapshot> {
constructor(user: User) {
super(user);
}
}
  • Αποσυνδεδεμένο (Detached): Μεταγενέστερες αλλαγές στο aggregate δεν επηρεάζουν ποτέ το payload.
  • Read-only για συνήθη χρήση: Τα ίδια properties είναι frozen, και οι μεταλλάσσουσες μέθοδοι των κλωνοποιημένων Date, Map και Set τιμών ρίχνουν σφάλμα.
  • Όχι όριο ασφαλείας: Το Object.freeze δεν καλύπτει εσωτερική κατάσταση native αντικειμένων, και όλοι οι καταναλωτές ενός event μοιράζονται το ίδιο payload. Το event.clonePayload() επιστρέφει ένα νέο αντίγραφο ανά κλήση όταν ένας consumer χρειάζεται αυτόνομο αντίγραφο.
  • Όχι transport format: Για διαμοιρασμό μέσω δικτύου, στείλτε ένα ρητό, serializable μήνυμα κατασκευασμένο από τα απαραίτητα πεδία.

Η συνάρτηση deepCloneAndFreeze(value) υποστηρίζει αντικείμενα, πίνακες, Date, Map, Set, RegExp και κυκλικές αναφορές.

Τα BaseCommand και BaseQuery είναι βασικές κλάσεις για αντικείμενα command και query. Το BaseCommand διατηρεί ένα προαιρετικό sessionPrincipal εκτός enumeration και toJSON(). Το getUpdateFields(fields) επιστρέφει τα πεδία που περιλαμβάνει το command για field-level authorization.

Το CommandBaseHandler εκτελεί ένα command και δημοσιεύει τα events του aggregate που επιστρέφει:

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() επιστρέφει το aggregate, ή αποτέλεσμα που το περιέχει ως aggregate. Το execute() το καλεί, δημοσιεύει τα buffered events μία φορά με το eventBus.publishAll(events, aggregate), καθαρίζει το buffer και επιστρέφει το αποτέλεσμα αναλλοίωτο.
  • Ένα αποτυχημένο handle() δεν δημοσιεύει τίποτα. Αν το publishAll() αποτύχει, το σφάλμα διαδίδεται και τα events παραμένουν στο buffer.
  • Αν το publishAll() επιστρέφει promise, το execute() το αναμένει αφού καθαρίσει το buffer.
  • Οποιοδήποτε αντικείμενο με μέθοδο publishAll(events) ικανοποιεί το IDomainEventPublisher.

Το AggregateRoot.commit() δεν δημοσιεύει τίποτα εξ ορισμού. Αφήστε το CommandBaseHandler να δημοσιεύσει, ή διαβάστε το getUncommittedEvents() και καλέστε το uncommit() χειροκίνητα.

Ένας command handler φορτώνει το έγκυρο aggregate μέσω του IWriteSideAggregateRepository<TEntity>.findById(id): οφείλει να διαβάζει την πρωτεύουσα αποθήκευση και ποτέ cache. Το AggregateRepository του @cqrs-ddd/mikro-orm το υλοποιεί για το MikroORM.

Η μέθοδος save() δηλώνει τον κύκλο ζωής της με το @PersistedWrite:

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;
}
}

Το @PersistedWrite συνδυάζει τρεις decorators με αυστηρή σειρά:

  1. @Cache(options): Μετά από επιτυχή εγγραφή, γράφει το snapshot μέσω CAS (isCacheNewer) και τοποθετεί barriers για τα invalidateKeys. Όταν το save() επιστρέφει null ή undefined, τοποθετεί barriers για τα deleteKeys.
  2. @AcknowledgePersisted: Προωθεί το persisted version του επιλεγμένου aggregate μόνο μετά από επιτυχή εγγραφή.
  3. @MapPersistenceErrors: Μεταφράζει παραβιάσεις μοναδικότητας (unique constraints) σε domain σφάλματα μέσω του ενεργού persistence dialect.

Για διαγραφές (deletes), χρησιμοποιούνται μόνο τα @Cache και @MapPersistenceErrors:

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;
}
}

Ένα query repository επεκτείνει το QueryRepository και διακοσμεί το find() με το @FromCache:

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)); // μπορεί να εξυπηρετηθεί από την cache
await users.find(new GetUserQuery(id, { refresh: true })); // διαβάζει πάντα τη βάση δεδομένων
  • Κάθε hit επαναϋδατώνεται (rehydrated) μέσω του hydrator, διασφαλίζοντας ότι hit και miss επιστρέφουν τον ίδιο τύπο.
  • Χωρίς hydrator αποθηκεύονται μόνο plain data.
  • Αποθηκεύονται μόνο μη-nullish αποτελέσματα. Ένα barrier ή αποθηκευμένο null μετρά ως miss.
  • Ένα query με refresh: true παρακάμπτει την cache.

Το @Cache λειτουργεί με οποιοδήποτε ICache (get, set, delete). Το @FromCache απαιτεί IVersionedCache, το οποίο προσθέτει revision fence:

Μέθοδος Συμβόλαιο
readState(key) η κατάσταση της εγγραφής (hit, miss, expired), η τιμή και το revision
invalidate(key) ακυρώνει την τιμή και προωθεί το revision
tryFill(key, observedRevision, value, options?) γράφει μόνο αν το revision παραμένει observedRevision, αλλιώς false

Το @FromCache καταγράφει το revision πριν την ανάγνωση από τη βάση και γεμίζει την cache με tryFill, αποτρέποντας την εγγραφή παρωχημένου snapshot όταν μια ανάγνωση συναγωνίστηκε με invalidation ή διαγραφή. Ένας adapter που υλοποιεί μόνο ICache παρακάμπτεται από το @FromCache τόσο για αναγνώσεις όσο και για fills, καταγράφοντας σχετική προειδοποίηση.

Υλοποιήσεις του IVersionedCache:

  • MemoryCache({ defaultTtlMs = 60_000, maxEntries = 10_000 }): in-process cache με JSON κλωνοποίηση, ιδανική για tests.
  • MikroOrmCache στο @cqrs-ddd/mikro-orm: μία γραμμή βάσης ανά κλειδί, με compare-and-set σε δικό του transaction.

Τα mutation barriers (CacheMutationBarrier) σηματοδοτούν ένα διαγραμμένο ή ακυρωμένο κλειδί για διάστημα barrierTtl (60 δευτερόλεπτα εξ ορισμού).

Τα cacheKey(resource, conditions, tenant?) και cacheKeyTemplate(template, tenant?) οργανώνουν κάθε κλειδί ανά tenant. Το tenant προέρχεται κατά σειρά από:

  1. Το ρητό όρισμα tenant (string ή αντικείμενο με tenantId),
  2. Για το cacheKeyTemplate, το tenantId της πηγής context,
  3. Τον resolver που καταχωρεί η εφαρμογή στην εκκίνηση με το setTenantResolver(fn),
  4. Αλλιώς, ρίχνεται MissingTenantContextError.
import { AsyncLocalStorage } from 'node:async_hooks';
import { setTenantResolver } from '@cqrs-ddd/core/application';
const requestTenant = new AsyncLocalStorage<string>();
setTenantResolver(() => requestTenant.getStore());
requestTenant.run(tenantId, () => handle(request));

Δεν υπάρχει κοινόχρηστο namespace: ένα ελλείπον tenant ρίχνει σφάλμα. Η ίδια επίλυση είναι διαθέσιμη μέσω του requireTenant(purpose, source?):

import { requireTenant } from '@cqrs-ddd/core/application';
requireTenant('access token issuance');

Η συνάρτηση domainErrorHttpStatus(error) από το /http αντιστοιχίζει τα σφάλματα αυτού του πακέτου σε τιμές { statusCode, error, message }:

Σφάλμα (Error) Status
ConcurrencyConflictError 409
EntityNotFoundException 404
MissingTenantContextError 500, με γενικό μήνυμα (αίτημα χωρίς tenant είναι server fault)
οποιοδήποτε άλλο DomainException, συμπεριλαμβανομένου του InvalidValueException 400

Επιστρέφει undefined για οτιδήποτε άλλο.

  • Η διατήρηση στη βάση και η δημοσίευση events δεν είναι ατομικές ενέργειες. Μια διακοπή μεταξύ της εγγραφής και του publishAll() χάνει τα events. Η αξιόπιστη παράδοση απαιτεί transactional outbox.
  • Η συντήρηση της cache μετά το commit είναι best-effort. Αν ένα invalidation αποτύχει, η προηγούμενη εγγραφή παραμένει μέχρι να λήξει.
  • Ένα mutation barrier προστατεύει μόνο για τη διάρκεια του barrierTtl.
  • Οι persistence helpers απορρίπτουν εξωτερικά transactions. Μια εγγραφή που πρέπει να γίνει commit μαζί με το aggregate δεν μπορεί να μοιραστεί το transaction του.
  • Το @MapPersistenceErrors αντιστοιχίζει μοναδικά σφάλματα στον βαθμό που το υποστηρίζει το δηλωμένο dialect (το MikroOrmDialect καλύπτει PostgreSQL και SQLite).
  • Ο tenant resolver είναι process-wide.

Διπλή άδεια υπό AGPL-3.0-or-later ή Commercial License. Δείτε τα LICENSE και COMMERCIAL_LICENSE.txt στη ρίζα του repository.