@cqrs-ddd/cqrs
Εκτελεί commands, queries και events μέσω των διακοσμημένων handlers τους, καθένας εκ των οποίων γίνεται compile μέσα στο δικό του pipeline από behaviors.
Τα @CommandHandler, @QueryHandler και @EventsHandler σηματοδοτούν κλάσεις handler. Τα @UsePipeline και @SkipPipeline από το @cqrs-ddd/pipeline δηλώνουν τα behaviors τους. Το createCqrs() κατασκευάζει τα runtime buses (CommandBus, QueryBus, EventBus, UnhandledExceptionBus). Δεν υπάρχει dependency-injection container ή framework: η εφαρμογή δημιουργεί τα instances των handlers με new, περνά τις απαραίτητες εξαρτήσεις άμεσα, και τα καταχωρεί.
Χρησιμοποιήστε το @cqrs-ddd/cqrs για αρχιτεκτονικές χωρίς framework (απλό Node.js, Fastify, Express, serverless). Εάν αναπτύσσετε σε NestJS, χρησιμοποιήστε το @cqrs-ddd/nestjs για να κάνετε wrap τους επίσημους handlers του @nestjs/cqrs.
Εγκατάσταση
Ενότητα με τίτλο «Εγκατάσταση»pnpm add @cqrs-ddd/cqrs @cqrs-ddd/pipelineΑπαιτεί Node.js 22.12 ή νεότερο. Οι decorators υποστηρίζουν τόσο τα standard decorators του TypeScript 5+ όσο και τα legacy experimentalDecorators.
Γρήγορο παράδειγμα
Ενότητα με τίτλο «Γρήγορο παράδειγμα»import { CommandHandler, createCqrs, type EventBus, type ICommandHandler,} from '@cqrs-ddd/cqrs';import { LoggingBehavior, UsePipeline } from '@cqrs-ddd/pipeline';import { AuditBehavior, audit } from '@cqrs-ddd/pipeline-audit';import { validated } from '@cqrs-ddd/pipeline-zod';import { z } from 'zod';
const CreateUserSchema = z.object({ email: z.string().email() });
class CreateUserCommand { constructor(readonly email: string) {}}
@CommandHandler(CreateUserCommand)@UsePipeline( validated(CreateUserSchema), audit({ action: 'user.create' }),)class CreateUserHandler implements ICommandHandler<CreateUserCommand> { constructor( private readonly users: UsersRepository, private readonly events: EventBus, ) {}
async execute(command: CreateUserCommand): Promise<string> { const user = await this.users.create(command.email); this.events.publish(new UserCreatedEvent(user.id)); return user.id; }}
// Composition Rootconst cqrs = createCqrs({ behaviors: [new AuditBehavior(postgresAuditSink)], globalBehaviors: { before: [LoggingBehavior] },});
cqrs.register(new CreateUserHandler(usersRepo, cqrs.eventBus));
// Dispatchconst userId = await cqrs.commandBus.execute<CreateUserCommand, string>( new CreateUserCommand('user@example.com'),);
await cqrs.close();Επιλογές παραμετροποίησης createCqrs
Ενότητα με τίτλο «Επιλογές παραμετροποίησης createCqrs»| Επιλογή | Τύπος | Προεπιλογή | Περιγραφή |
|---|---|---|---|
behaviors |
IPipelineBehavior[] |
[] |
Singleton instances από behaviors που παρέχονται στα pipelines. Δηλωμένα behaviors χωρίς ρητό instance κατασκευάζονται με new χωρίς ορίσματα. |
globalBehaviors |
GlobalBehaviorsOptions | GlobalBehaviorsOptions[] |
κανένα | Behaviors τοποθετημένα γύρω από κάθε handler εντός των καθορισμένων scopes ('all', 'commands', 'queries', 'events'). |
sources |
ContextSources |
κανένα | Από πού τα pipelines διαβάζουν και επαναφέρουν το tenantId και το correlationId, όπως tenantSource και correlationSource. |
diagnostics |
'strict' | 'warn' | 'off' |
'strict' |
Λειτουργία ελέγχου contract κατά το register(): το 'strict' ρίχνει PipelineConfigurationError, το 'warn' καταγράφει log, το 'off' απενεργοποιεί τους ελέγχους. |
logger |
PipelineLogger |
console |
Logger προορισμού για εκτέλεση pipeline, diagnostics και μη διαχειρίσιμα σφάλματα events. Το pinoLogger(pino) προσαρμόζει το Pino. |
bootstrapLogLevel |
LogLevel | 'none' |
'debug' |
Επίπεδο καταγραφής για τη συνοπτική γραμμή που εκτυπώνεται κατά το compile του pipeline κάθε handler. |
rethrowUnhandled |
boolean |
false |
Όταν είναι true, μη διαχειρίσιμα σφάλματα σε ασύγχρονους event handlers ρίχνονται ως uncaught exceptions της διεργασίας αντί να καταγράφονται απλώς σε log. |
Το createCqrs() επιστρέφει { commandBus, queryBus, eventBus, unhandledExceptionBus, register, close }.
Αρχιτεκτονική handlers & Invariants
Ενότητα με τίτλο «Αρχιτεκτονική handlers & Invariants»| Decorator | Interface | Μέθοδος | Invariant |
|---|---|---|---|
@CommandHandler(CommandClass) |
ICommandHandler<TCommand, TResult> |
execute(command) |
Ακριβώς ένας handler ανά κλάση command. Η καταχώριση διπλότυπου ρίχνει άμεσα σφάλμα κατά την εκκίνηση. |
@QueryHandler(QueryClass) |
IQueryHandler<TQuery, TResult> |
execute(query) |
Ακριβώς ένας handler ανά κλάση query. Η καταχώριση διπλότυπου ρίχνει άμεσα σφάλμα κατά την εκκίνηση. |
@EventsHandler(...EventClasses) |
IEventHandler<TEvent> |
handle(event) |
Οποιοσδήποτε αριθμός handlers ανά κλάση event. Οι handlers εκτελούνται παράλληλα όταν δημοσιεύεται το event. |
Κληρονομικότητα κλάσεων αιτήματος (Request Class Inheritance)
Ενότητα με τίτλο «Κληρονομικότητα κλάσεων αιτήματος (Request Class Inheritance)»Εάν ένα instance command ή query δεν έχει άμεσα καταχωρημένο handler για τον συγκεκριμένο constructor του, το bus αναζητά προς τα πάνω στην ιεραρχία prototype και δρομολογεί στον handler που είναι καταχωρημένος για την πλησιέστερη γονική κλάση.
Τα αιτήματα μπορούν να είναι οποιεσδήποτε κλάσεις. Κλάσεις που επεκτείνουν το BaseCommand ή το BaseQuery από το @cqrs-ddd/core/application φέρουν το brand REQUEST_KIND, ώστε τα behaviors να διακρίνουν ένα command από ένα query, καθώς και ένα προαιρετικό session principal εκτός των enumerable πεδίων τους.
Generic Result Typing
Ενότητα με τίτλο «Generic Result Typing»Τα buses δέχονται generic παραμέτρους τύπων στο execute():
const result = await cqrs.queryBus.execute<GetUserQuery, UserDto>(new GetUserQuery('u_1'));Σημασιολογία εκτέλεσης pipeline
Ενότητα με τίτλο «Σημασιολογία εκτέλεσης pipeline»Οι handlers δηλώνουν το pipeline τους χρησιμοποιώντας τα @UsePipeline(...entries) και @SkipPipeline(...Behaviors) από το @cqrs-ddd/pipeline:
- Compilation κατά την καταχώριση: Όταν εκτελείται το
cqrs.register(...handlers), το pipeline κάθε handler γίνεται compile σε μία βελτιστοποιημένη συνάρτηση runner μία φορά. Δεν πραγματοποιείται δυναμικό reflection μεταδεδομένων κατά το runtime dispatch. - Σειρά των behaviors: Τα global
beforebehaviors περιβάλλουν τα εξωτερικά behaviors του handler, ακολουθούμενα από τα εσωτερικά behaviors, την επιχειρησιακή μέθοδο του handler, και ταafterbehaviors. Δείτε Σειρά εκτέλεσης. - Execution context: Κάθε εκτέλεση λαμβάνει ένα
IPipelineContextπου περιέχει ταhandlerName,requestName,requestKind,correlationId,tenantId, και typeditems. - Skip isolation: Το
@SkipPipeline(LoggingBehavior)αφαιρεί το καθορισμένο global behavior για τον συγκεκριμένο handler χωρίς να επηρεάζει άλλους handlers.
Αναλυτικά τα buses
Ενότητα με τίτλο «Αναλυτικά τα buses»CommandBus
Ενότητα με τίτλο «CommandBus»- Σκοπός: Τροποποιεί την κατάσταση της εφαρμογής ή του domain.
- Dispatch: Το
commandBus.execute(command)περιμένει το pipeline του handler και επιστρέφει το αποτέλεσμα. - Ελλείπων handler: Εάν δεν έχει καταχωρηθεί handler, απορρίπτει με
CommandHandlerNotFoundException.
QueryBus
Ενότητα με τίτλο «QueryBus»- Σκοπός: Διαβάζει δεδομένα χωρίς μεταβολή της κατάστασης.
- Dispatch: Το
queryBus.execute(query)περιμένει το pipeline του handler και επιστρέφει το αποτέλεσμα του query. - Ελλείπων handler: Εάν δεν έχει καταχωρηθεί handler, απορρίπτει με
QueryHandlerNotFoundException.
EventBus
Ενότητα με τίτλο «EventBus»- Σκοπός: Δημοσιεύει domain και integration events μεταξύ bounded contexts.
- Dispatch: Τα
eventBus.publish(event)καιeventBus.publishAll(events)ειδοποιούν όλους τους καταχωρημένους event handlers. - Ασύγχρονη μη-μπλοκαριστική εκτέλεση: Το
publish()πυροδοτεί τους handlers και επιστρέφει άμεσα χωρίς να περιμένει την ολοκλήρωση των promises. Η σύγχρονη αρχικοποίηση εκτελείται πριν από την επιστροφή. - Ελλείποντες handlers: Events που δημοσιεύονται χωρίς καταχωρημένους handlers αγνοούνται με ασφάλεια.
- DDD Integration: Το
EventBusυλοποιεί τοIDomainEventPublisherαπό το@cqrs-ddd/core/application, επιτρέποντας στοCommandBaseHandlerνα δημοσιεύει άμεσα τα buffered domain events των aggregates.
UnhandledExceptionBus
Ενότητα με τίτλο «UnhandledExceptionBus»- Σκοπός: Κεντρικός χειρισμός για αποτυχίες ασύγχρονων event handlers.
- Όταν ένα promise ενός event handler απορρίπτεται, η αποτυχία δεν μπορεί να επιστραφεί στον εκδότη (publisher). Το bus εκπέμπει το σφάλμα και το event στο
UnhandledExceptionBus. - Εξ ορισμού, τα unhandled exceptions καταγράφονται στο επίπεδο
'error'. - Συνδρομή σε ειδοποιήσεις:
const subscription = cqrs.unhandledExceptionBus.subscribe(({ exception, cause }) => {alerts.captureException(exception, { extra: { event: cause } });});// Αργότερα:subscription.unsubscribe();
- Εάν έχει οριστεί
rethrowUnhandled: trueστοcreateCqrs(), οι αποτυχίες ρίχνονται ξανά ως uncaught exceptions, πυροδοτώντας τον τερματισμό της διεργασίας ή επανεκκίνηση από τον orchestrator.
Ομαλός τερματισμός (Graceful Shutdown)
Ενότητα με τίτλο «Ομαλός τερματισμός (Graceful Shutdown)»Περιμένετε πάντοτε το cqrs.close() πριν κλείσετε τις εξωτερικές υποδομές (συνδέσεις βάσης δεδομένων, message brokers, caches):
process.on('SIGTERM', async () => { // 1. Διακοπή αποδοχής νέας HTTP κυκλοφορίας await server.close();
// 2. Αναμονή για ολοκλήρωση των ασύγχρονων event handlers που εκτελούνται await cqrs.close();
// 3. Κλείσιμο συνδέσεων βάσης δεδομένων και cache await dbPool.end(); await redis.quit();});Το cqrs.close() παρακολουθεί όλες τις εκκρεμείς εκτελέσεις ασύγχρονων event handlers και ολοκληρώνεται μόνο αφού ολοκληρωθούν όλοι οι ενεργοί handlers.
Σημεία προσοχής & Οδηγίες παραγωγής
Ενότητα με τίτλο «Σημεία προσοχής & Οδηγίες παραγωγής»- In-process παράδοση μόνο: Οι handlers εκτελούνται στο event loop του Node.js. Εάν η διεργασία τερματιστεί απρόσμενα πριν εκτελεστεί ένας event handler, το event χάνεται. Για εγγυημένη, ανθεκτική παράδοση μεταξύ services, συνδυάστε τα domain events με το Transactional Outbox pattern.
- Οι handlers είναι singletons: Τα instances των handlers δημιουργούνται μία φορά και διαμοιράζονται σε όλες τις κλήσεις. Οι handlers οφείλουν να είναι stateless. State που αφορά το αίτημα (ταυτότητα χρήστη, correlation tokens, tenant IDs) πρέπει να διαβάζεται από το
AsyncLocalStorageμέσω context sources. - Χωρίς prototype patching: Πολλαπλά ανεξάρτητα CQRS runtimes μπορούν να συνυπάρχουν στην ίδια διεργασία Node.js χωρίς παρεμβολές μεταξύ τους.