@cqrs-ddd/nestjs
Ο επίσημος adapter του NestJS για τα πακέτα @cqrs-ddd.
Μια εφαρμογή NestJS διατηρεί την ιδιωματική αρχιτεκτονική του NestJS: modules, dependency injection, controllers, και το επίσημο @nestjs/cqrs (CqrsModule, CommandBus, QueryBus, EventBus, @CommandHandler, @QueryHandler, @EventsHandler). Ο adapter παρέχει όσα λείπουν από το @nestjs/cqrs:
- Pipeline behaviors γύρω από handlers: κάνει compile και εκτελεί τα behaviors του
@cqrs-ddd/pipelineγύρω από κάθε instance command, query και event handler κατά το bootstrap της εφαρμογής. - Τυποποιημένο exception filtering: το
ErrorFilterμετατρέπει κάθε σφάλμα του@cqrs-ddd(domain errors, validation failures, άρνηση εξουσιοδότησης, rate limits, διενέξεις idempotency, feature flags) σε NestJSHttpExceptionμε το τυπικό response body του NestJS και τις κατάλληλες HTTP επικεφαλίδες. - Διάδοση correlation ID: το
CorrelationMiddlewareεκτελεί κάθε αίτημα με το correlation id από την επικεφαλίδαx-correlation-id(ή ένα νέο UUIDv7) και το επιστρέφει στην απάντηση. - Ασύγχρονο job context: το
JobContextModuleδιαδίδει tenants, correlation IDs και principals σε ουρές εργασίας (όπως BullMQ) και τα αποκαθιστά στους processors. - Ενσωμάτωση κύκλου ζωής aggregate: commands που επεκτείνουν το
CommandBaseHandlerδημοσιεύουν αυτόματα τα uncommitted domain events μέσω του NestJSEventBus.
Μια πλήρης λειτουργική εφαρμογή NestJS που παρουσιάζει όλα τα παραπάνω μοτίβα είναι διαθέσιμη στο nestjs-pipeline api/.
Εγκατάσταση
Ενότητα με τίτλο «Εγκατάσταση»pnpm add @cqrs-ddd/nestjs @cqrs-ddd/pipeline @cqrs-ddd/core @nestjs/cqrsPeer Dependencies
Ενότητα με τίτλο «Peer Dependencies»- Υποχρεωτικά peers:
@nestjs/common,@nestjs/coreκαι@nestjs/cqrs12, και@cqrs-ddd/pipelineκαι@cqrs-ddd/core0.5. Το ίδιο το NestJS απαιτείreflect-metadataκαιrxjs. - Προαιρετικά peers, εγκαθίστανται όταν τα χρησιμοποιεί η εφαρμογή:
@cqrs-ddd/pipeline-zod,@cqrs-ddd/pipeline-casl,@cqrs-ddd/pipeline-feature-flags,@cqrs-ddd/pipeline-rate-limitκαι@cqrs-ddd/pipeline-idempotency: τοErrorFilterμετατρέπει τα σφάλματα καθενός εξ αυτών αν είναι εγκατεστημένο.@cqrs-ddd/pipeline-correlation, για το@cqrs-ddd/nestjs/correlation.@cqrs-ddd/pipeline-job-context, για το@cqrs-ddd/nestjs/job-context.
Οποιοδήποτε άλλο πακέτο behavior (cache, resilience, audit, telemetry, dead letters) εκτελείται στα pipelines χωρίς να απαιτείται ειδική γνώση από τον adapter.
Αρχιτεκτονικός κύκλος ζωής (Architectural Lifecycle)
Ενότητα με τίτλο «Αρχιτεκτονικός κύκλος ζωής (Architectural Lifecycle)»Ο adapter λειτουργεί μέσω του PipelineBootstrap, μίας @Injectable() υπηρεσίας που καταχωρείται από το PipelineModule:
Nest Application Bootstrap │ ├─ 1. Εκτελείται το OnApplicationBootstrap lifecycle hook │ ├─ 2. Το DiscoveryService σαρώνει όλους τους καταχωρημένους providers │ ├─ 3. Εντοπίζει κλάσεις διακοσμημένες με @CommandHandler, @QueryHandler, @EventsHandler │ ├─ 4. Διαβάζει τις δηλώσεις pipeline (@UsePipeline, @SkipPipeline) μέσω του pipelineOf() │ ├─ 5. Κάνει compile το πλάνο εκτέλεσης (globalBehaviors + handler behaviors - skipped behaviors) │ ├─ 6. Επιλύει τα instances των behaviors από το Nest DI container μέσω του class token │ ├─ 7. Επικυρώνει τα static singleton scopes (αποτυγχάνει άμεσα αν handler ή behavior είναι request-scoped) │ ├─ 8. Επικυρώνει τα behavior contracts (σειρά, απαιτούμενα items, διαγνωστικοί έλεγχοι) │ └─ 9. Κάνει wrap τη μέθοδο του handler instance (execute για commands/queries, handle για events)Βασικές εγγυήσεις runtime:
- Instance wrapping, όχι μετάλλαξη prototype: Μόνο τα provider instances της τρέχουσας εφαρμογής γίνονται wrap. Τα prototypes των κλάσεων δεν τροποποιούνται ποτέ, επιτρέποντας πολλαπλές απομονωμένες εφαρμογές στην ίδια διεργασία.
- Fail-fast επικύρωση: Ελλείποντες providers behaviors, αμφίσημοι διπλότυποι providers, request-scoped handlers ή σπασμένα contracts αποτυγχάνουν άμεσα κατά την εκκίνηση της εφαρμογής πριν ο HTTP listener αρχίσει να δέχεται κλήσεις.
- Καθαρός τερματισμός: Το
OnModuleDestroyεπαναφέρει τους αρχικούς method descriptors σε όλα τα wrapped handler instances.
Ρύθμιση Pipeline
Ενότητα με τίτλο «Ρύθμιση Pipeline»Εισάγετε το PipelineModule.forRoot() στο root module ή στο core infrastructure module της εφαρμογής σας:
import { Module } from '@nestjs/common';import { CqrsModule } from '@nestjs/cqrs';import { PipelineModule } from '@cqrs-ddd/nestjs';import { logging } from '@cqrs-ddd/pipeline';import { correlationSource } from '@cqrs-ddd/pipeline-correlation';import { DeadLetterBehavior } from '@cqrs-ddd/pipeline-deadletter';import { tenantSource } from '@cqrs-ddd/pipeline-tenant';
@Module({ imports: [ CqrsModule.forRoot(), PipelineModule.forRoot({ sources: { tenantId: tenantSource, correlationId: correlationSource, }, globalBehaviors: [ { scope: 'all', before: [logging({ requestResponseLogLevel: 'log' })] }, { scope: 'commands', before: [DeadLetterBehavior] }, ], diagnostics: 'strict', }), ],})export class CorePipelineModule {}PipelineOptions
Ενότητα με τίτλο «PipelineOptions»| Επιλογή | Τύπος | Προεπιλογή | Περιγραφή |
|---|---|---|---|
globalBehaviors |
GlobalBehaviorsOptions | GlobalBehaviorsOptions[] |
κανένα | Behaviors τοποθετημένα γύρω από handlers στα επιλεγμένα scopes ('all', 'commands', 'queries', 'events'). |
sources |
ContextSources |
κανένα | Από πού τα pipelines διαβάζουν και αποκαθιστούν τα tenantId και correlationId. |
diagnostics |
'strict' | 'warn' | 'off' |
'strict' |
Λειτουργία επικύρωσης των behavior contracts: το 'strict' ρίχνει PipelineConfigurationError στην εκκίνηση, το 'warn' καταγράφει προειδοποιήσεις, το 'off' παρακάμπτει τους ελέγχους. |
Καταχώριση Behaviors με Dependency Injection
Ενότητα με τίτλο «Καταχώριση Behaviors με Dependency Injection»Τα behaviors είναι απλές κλάσεις παραμετροποιημένες με εξαρτήσεις (loggers, Redis stores, connection pools, CASL loaders). Κάθε module παρέχει το instance του behavior υπό το class token του:
import { Module, Logger } from '@nestjs/common';import { LoggingBehavior } from '@cqrs-ddd/pipeline';import { IdempotencyBehavior, RedisIdempotencyStore } from '@cqrs-ddd/pipeline-idempotency';import { CaslBehavior } from '@cqrs-ddd/pipeline-casl';import { RedisService } from '../redis/redis.service.js';import { PermissionLoader } from '../auth/permission-loader.service.js';
@Module({ providers: [ { provide: LoggingBehavior, useFactory: () => new LoggingBehavior(new Logger('Pipeline')), }, { provide: IdempotencyBehavior, inject: [RedisService], useFactory: (redis: RedisService) => new IdempotencyBehavior(new RedisIdempotencyStore(redis.client)), }, { provide: CaslBehavior, inject: [PermissionLoader], useFactory: (loader: PermissionLoader) => new CaslBehavior(loader), }, ], exports: [LoggingBehavior, IdempotencyBehavior, CaslBehavior],})export class PipelineBehaviorsModule {}Κανόνες Providers
Ενότητα με τίτλο «Κανόνες Providers»- Απαιτείται singleton scope: Κάθε behavior provider και κάθε handler που εκτελεί behaviors οφείλει να είναι στατικό singleton (
Scope.DEFAULT). Αν ένας handler ή behavior είναι request-scoped (Scope.REQUEST), τοPipelineBootstrapαποτυγχάνει άμεσα κατά την εκκίνηση. Δεδομένα αιτήματος πρέπει να ανακτώνται μέσωAsyncLocalStorageή context sources. - Μοναδικός provider ανά behavior: Ακριβώς ένα module πρέπει να παρέχει κάθε class token ενός behavior. Αν δύο modules καταχωρίσουν providers για το ίδιο token, το
PipelineBootstrapρίχνει σφάλμα αναφέροντας και τα δύο modules προς αποφυγή ασάφειας.
Handlers και δηλώσεις (Declarations)
Ενότητα με τίτλο «Handlers και δηλώσεις (Declarations)»Command, query και event handlers χρησιμοποιούν τους κλασικούς decorators του @nestjs/cqrs. Χρησιμοποιήστε τα @UsePipeline και @SkipPipeline από το @cqrs-ddd/pipeline για τη δήλωση behaviors:
import { CommandHandler, ICommandHandler } from '@nestjs/cqrs';import { UsePipeline, SkipPipeline } from '@cqrs-ddd/pipeline';import { idempotent } from '@cqrs-ddd/pipeline-idempotency';import { validated } from '@cqrs-ddd/pipeline-zod';import { requires } from '@cqrs-ddd/pipeline-casl';import { audit } from '@cqrs-ddd/pipeline-audit';import { DeadLetterBehavior } from '@cqrs-ddd/pipeline-deadletter';import { CreateOrderCommand, CreateOrderSchema } from './create-order.command.js';import { OrdersRepository } from '../orders.repository.js';import { orderIdempotencyKey } from './order-key.factory.js';
@CommandHandler(CreateOrderCommand)@UsePipeline( requires({ action: 'create', subject: 'Order' }), validated(CreateOrderSchema), idempotent({ keyFactory: orderIdempotencyKey }), audit({ action: 'order.created' }),)@SkipPipeline(DeadLetterBehavior)export class CreateOrderHandler implements ICommandHandler<CreateOrderCommand> { constructor(private readonly orders: OrdersRepository) {}
async execute(command: CreateOrderCommand): Promise<string> { const order = await this.orders.create(command); return order.id; }}Ροή εκτέλεσης behaviors
Ενότητα με τίτλο «Ροή εκτέλεσης behaviors»Όταν ένας controller καλεί commandBus.execute(command), η κλήση εισέρχεται στο compiled pipeline του CreateOrderHandler. Τα behaviors σχηματίζουν ομόκεντρα στρώματα: πρώτα τα global before behaviors (logging), έπειτα τα behaviors του ίδιου του handler κατά τη σειρά δήλωσης (authorization, validation, idempotency, audit), και τέλος τα global after behaviors. Η μέθοδος execute(command) του handler εκτελείται στο εσωτερικό με το επικυρωμένο command. Κατά την επιστροφή, κάθε behavior παραλαμβάνει το αποτέλεσμα ή το σφάλμα. Ένα behavior που ρίχνει σφάλμα (π.χ. απόρριψη authorization) διακόπτει την εκτέλεση πριν εκτελεστούν τα εσωτερικά στρώματα.
Διαχείριση σφαλμάτων με το ErrorFilter
Ενότητα με τίτλο «Διαχείριση σφαλμάτων με το ErrorFilter»Το ErrorFilter αναχαιτίζει όλα τα σφάλματα του @cqrs-ddd και τα μεταφράζει σε τυποποιημένα NestJS HttpException instances που ταιριάζουν με το προεπιλεγμένο JSON error format του NestJS ({ statusCode, message, error }).
Καταχωρίστε το καθολικά (globally) στο AppModule:
import { Module } from '@nestjs/common';import { APP_FILTER } from '@nestjs/core';import { ErrorFilter } from '@cqrs-ddd/nestjs';
@Module({ providers: [ { provide: APP_FILTER, useClass: ErrorFilter, }, ],})export class AppModule {}Πίνακας αντιστοίχισης σφαλμάτων
Ενότητα με τίτλο «Πίνακας αντιστοίχισης σφαλμάτων»| Κλάση σφάλματος | Πακέτο προέλευσης | HTTP Status | Μορφή Response Body | Πρόσθετες επικεφαλίδες |
|---|---|---|---|---|
ZodValidationError |
@cqrs-ddd/pipeline-zod |
400 Bad Request | { statusCode: 400, error: 'Bad Request', message: ['field: error message'] } |
— |
UnauthorizedActionException |
@cqrs-ddd/pipeline-casl |
403 Forbidden | { statusCode: 403, error: 'Forbidden', message, action, subject } |
— |
FeatureDisabledError |
@cqrs-ddd/pipeline-feature-flags |
403 Forbidden | { statusCode: 403, error: 'Forbidden', message, flag } |
— |
RateLimitExceededError |
@cqrs-ddd/pipeline-rate-limit |
429 Too Many Requests | { statusCode: 429, error: 'Too Many Requests', message, retryAfter } |
Retry-After: <seconds> |
IdempotencyConflictError |
@cqrs-ddd/pipeline-idempotency |
409 Conflict (in_progress, replay_scope) / 422 Unprocessable Entity (key_reuse) |
{ statusCode, error, message, idempotencyKey, reason } |
— |
EntityNotFoundException |
@cqrs-ddd/core |
404 Not Found | { statusCode: 404, error: 'Not Found', message: 'Order not found' } |
— |
ConcurrencyConflictError |
@cqrs-ddd/core (μέσω @cqrs-ddd/mikro-orm) |
409 Conflict | { statusCode: 409, error: 'Conflict', message } |
— |
MissingTenantContextError |
@cqrs-ddd/core |
500 Internal Server Error | { statusCode: 500, error: 'Internal Server Error', message: 'Internal server error' } |
— |
οποιοδήποτε άλλο DomainException |
@cqrs-ddd/core ή η εφαρμογή |
400 Bad Request | { statusCode: 400, error: 'Bad Request', message } |
— |
Μη αναγνωρίσιμα σφάλματα προωθούνται στην προεπιλεγμένη συμπεριφορά του BaseExceptionFilter του NestJS.
Προσαρμοσμένα Exception Filters
Ενότητα με τίτλο «Προσαρμοσμένα Exception Filters»Για να διαχειριστείτε ειδικά σφάλματα της εφαρμογής σας διατηρώντας παράλληλα τις μετατροπές του @cqrs-ddd, επεκτείνετε το ErrorFilter ή χρησιμοποιήστε τα toHttpException και httpAnswer:
import { Catch, ArgumentsHost } from '@nestjs/common';import { ErrorFilter } from '@cqrs-ddd/nestjs';import { CustomBillingError } from './billing.errors.js';
@Catch()export class CustomAppFilter extends ErrorFilter { override catch(exception: unknown, host: ArgumentsHost): void { if (exception instanceof CustomBillingError) { const response = host.switchToHttp().getResponse(); response.status(402).json({ statusCode: 402, message: exception.message }); return; } super.catch(exception, host); }}Middleware για Correlation ID
Ενότητα με τίτλο «Middleware για Correlation ID»Το CorrelationMiddleware αντλεί το correlation id κάθε αιτήματος από την επικεφαλίδα x-correlation-id (ή παράγει ένα UUIDv7 αν απουσιάζει), εκτελεί το αίτημα εντός αυτού στο AsyncLocalStorage, και επιστρέφει την επικεφαλίδα στην απάντηση. Αποτελεί το httpCorrelation() του @cqrs-ddd/pipeline-correlation ως NestJS middleware και λειτουργεί σε Express και Fastify.
Καταχώριση στο AppModule:
import { Module, NestModule, MiddlewareConsumer } from '@nestjs/common';import { CorrelationMiddleware } from '@cqrs-ddd/nestjs/correlation';
@Module({})export class AppModule implements NestModule { configure(consumer: MiddlewareConsumer) { consumer.apply(CorrelationMiddleware).forRoutes('*'); }}Background Jobs με το JobContextModule
Ενότητα με τίτλο «Background Jobs με το JobContextModule»Όταν ένα αίτημα εισάγει μια εργασία σε ουρά (BullMQ, SQS, κλπ.), ο worker που την εκτελεί δεν διαθέτει το αρχικό context του αιτήματος. Το JobContextModule καταχωρεί τη λίστα των tenants, τον principal resolver και τις πηγές context στο @cqrs-ddd/pipeline-job-context:
import { Module } from '@nestjs/common';import { JobContextModule } from '@cqrs-ddd/nestjs/job-context';import { correlationSource } from '@cqrs-ddd/pipeline-correlation';import { tenantSource } from '@cqrs-ddd/pipeline-tenant';import { SessionJobPrincipal } from './session-job-principal.service.js';import { AccountsModule } from './accounts/accounts.module.js';
@Module({ imports: [ JobContextModule.forRoot({ imports: [AccountsModule], principal: SessionJobPrincipal, tenants: ['tenant_primary', 'tenant_secondary'], sources: { tenantId: tenantSource, correlationId: correlationSource, }, }), ],})export class BackgroundJobModule {}Εισαγωγή στην ουρά με withJobContext, και χρήση του @InJobContext() στους processors:
import { Processor, WorkerHost } from '@nestjs/bullmq';import { InJobContext, withJobContext } from '@cqrs-ddd/pipeline-job-context';import { Job } from 'bullmq';
// Κατά την εισαγωγή στην ουρά:await queue.add('process-statement', withJobContext({ accountId: 'acc_123' }));
// Στον processor:@Processor('statements')export class StatementProcessor extends WorkerHost { @InJobContext() async process(job: Job): Promise<void> { // Το τρέχον tenant, principal και correlation ID είναι ενεργά στο AsyncLocalStorage }}Ενσωμάτωση Domain-Driven Design
Ενότητα με τίτλο «Ενσωμάτωση Domain-Driven Design»Τα aggregates και domain events του @cqrs-ddd/core ενσωματώνονται άψογα με το @nestjs/cqrs.
CommandBaseHandler με το NestJS EventBus
Ενότητα με τίτλο «CommandBaseHandler με το NestJS EventBus»Το EventBus του @nestjs/cqrs υλοποιεί το IDomainEventPublisher (publishAll(events, dispatcherContext)). Command handlers που επεκτείνουν το CommandBaseHandler δημοσιεύουν αυτόματα τα uncommitted domain events των aggregates:
import { CommandHandler } from '@nestjs/cqrs';import { EventBus } from '@nestjs/cqrs';import { CommandBaseHandler } from '@cqrs-ddd/core/application';import { OrdersRepository } from './orders.repository.js';import { Order } from './order.aggregate.js';import { SubmitOrderCommand } from './submit-order.command.js';
@CommandHandler(SubmitOrderCommand)export class SubmitOrderHandler extends CommandBaseHandler<SubmitOrderCommand, Order> { constructor( private readonly orders: OrdersRepository, eventBus: EventBus, ) { super(eventBus); }
async handle(command: SubmitOrderCommand): Promise<Order> { const order = await this.orders.findById(command.orderId); order.submit(); await this.orders.save(order); return order; // το execute() δημοσιεύει αυτόματα τα domain events μέσω του eventBus }}Πλήρες παράδειγμα παραγωγής
Ενότητα με τίτλο «Πλήρες παράδειγμα παραγωγής»Ακολουθεί ένας ολοκληρωμένος command handler εγγραφής χρήστη σε NestJS που συνδυάζει validation, idempotency, audit logging και δημοσίευση domain events:
// 1. Command και Schemaimport { z } from 'zod';
export const RegisterUserSchema = z.object({ email: z.email(), fullName: z.string().min(2),});
export class RegisterUserCommand { constructor( readonly email: string, readonly fullName: string, ) {}}
// 2. Command Handlerimport { CommandHandler, EventBus } from '@nestjs/cqrs';import { CommandBaseHandler } from '@cqrs-ddd/core/application';import { UsePipeline } from '@cqrs-ddd/pipeline';import { validated } from '@cqrs-ddd/pipeline-zod';import { idempotent, createPartitionedIdempotencyKeyFactory } from '@cqrs-ddd/pipeline-idempotency';import { audit } from '@cqrs-ddd/pipeline-audit';import { User } from '../domain/user.aggregate.js';import { UsersRepository } from '../persistence/users.repository.js';
const userKeyFactory = createPartitionedIdempotencyKeyFactory({ principal: () => 'registration', operation: (ctx) => (ctx.request as RegisterUserCommand).email,});
@CommandHandler(RegisterUserCommand)@UsePipeline( validated(RegisterUserSchema), idempotent({ keyFactory: userKeyFactory, ttl: 86_400_000 }), audit({ action: 'user.registered' }),)export class RegisterUserHandler extends CommandBaseHandler<RegisterUserCommand, User> { constructor( private readonly users: UsersRepository, eventBus: EventBus, ) { super(eventBus); }
async handle(command: RegisterUserCommand): Promise<User> { const user = User.register(command.email, command.fullName); await this.users.save(user); return user; }}
// 3. Controllerimport { Controller, Post, Body } from '@nestjs/common';import { CommandBus } from '@nestjs/cqrs';
@Controller('users')export class UsersController { constructor(private readonly commandBus: CommandBus) {}
@Post('register') async register(@Body() body: RegisterUserCommand) { const user = await this.commandBus.execute( new RegisterUserCommand(body.email, body.fullName), ); return { id: user.id, email: user.email }; }}