diff --git a/docs/getting-started/installation.md b/docs/getting-started/installation.md index 574120f..88371c6 100644 --- a/docs/getting-started/installation.md +++ b/docs/getting-started/installation.md @@ -24,3 +24,10 @@ MongoDB repositories require: ```bash yarn add mongodb ``` + +## TypeScript Resolution + +The package publishes ESM, CommonJS and declaration files for every public +subpath. Modern projects should prefer `moduleResolution: "NodeNext"` or +`"Bundler"`, but declaration mappings are also provided for projects still using +classic `moduleResolution: "node"`. diff --git a/docs/guides/adapters.md b/docs/guides/adapters.md index dc5e825..ab44246 100644 --- a/docs/guides/adapters.md +++ b/docs/guides/adapters.md @@ -27,3 +27,52 @@ export default class MyPublisher implements DomainEventPublisher { If an adapter needs a third-party dependency, expose it through a subpath and mark that dependency as an optional peer dependency. + +## Message Bus Hooks + +Message bus adapters can expose publisher hooks so applications can attach +replicated publishers, websocket notifications, tracing or auditing without +wrapping the adapter in an application-local class. + +```ts +import AmqpMessageBusAdapter from '@haskou/ddd-kernel/adapters/pubsub/amqp'; + +const messageBus = new AmqpMessageBusAdapter({ + publisherHookErrorPolicy: { + handleAfterPublishError(error, context) { + logger.error( + `Post-publish hook failed for ${context.topic}: ${String(error)}`, + ); + }, + shouldFailAfterPublish() { + return false; + }, + }, + publisherHooks: [ + { + afterPublish: async ({ domainEvent, message }) => { + await websocketPublisher.publish(domainEvent ?? message); + }, + }, + ], +}); +``` + +Custom generic adapters should implement the `MessageBus` contract. Domain-event +adapters should implement `DomainMessageBus`. Both can delegate hook execution +through `PublisherHookPipeline`: + +```ts +import { + PublisherHookPipeline, + type PublisherHook, +} from '@haskou/ddd-kernel/adapters/pubsub'; + +export default class CustomMessageBus { + private readonly hooks = new PublisherHookPipeline(); + + public registerPublisherHooks(...hooks: PublisherHook[]) { + this.hooks.register(...hooks); + } +} +``` diff --git a/docs/guides/amqp-pubsub.md b/docs/guides/amqp-pubsub.md index 004015e..aa46578 100644 --- a/docs/guides/amqp-pubsub.md +++ b/docs/guides/amqp-pubsub.md @@ -16,6 +16,7 @@ new AmqpMessageBusAdapter({ dsn: 'amqp://localhost', exchange: 'users-service', maxRetries: 3, + publisherHooks: [replicatedPublisherHook], retryDelayInMilliseconds: 1000, serviceName: 'users-service', }); @@ -30,3 +31,7 @@ Environment variables: Failed messages are sent to `_dlx`. Use `consumeDlx` to retry failed messages. + +`publisherHooks` run around each domain event published by the adapter. Use them +for transport-adjacent fan-out such as websocket updates, replicated-state +publishers, tracing or audit logs. diff --git a/docs/reference/consumer.md b/docs/reference/consumer.md index edeca78..b738b5c 100644 --- a/docs/reference/consumer.md +++ b/docs/reference/consumer.md @@ -25,11 +25,44 @@ correlation IDs around handler execution: ```ts kernel.registerConsumerMiddleware({ - async handle(event, next) { + async handle(event, next, context) { + logger.info(`Handling ${context.eventName}`); await next(); }, }); ``` -Middleware receives the event and a `next` callback. The kernel does not include -a full outbox or idempotency implementation. +Middleware receives the event, the next pipeline callback and a +`ConsumerExecutionContext` containing queue, exchange, event id, correlation id +and causation id. Transport adapters can also attach metadata, such as AMQP +headers or retry counts, to `context.metadata`. + +## Built-in Middleware + +The pub/sub adapter package includes small middleware implementations for common +consumer concerns. They are intentionally infrastructure-level primitives, not a +full outbox implementation. + +```ts +import { + CorrelationConsumerMiddleware, + IdempotencyConsumerMiddleware, + InMemoryIdempotencyStore, + RetryConsumerMiddleware, +} from '@haskou/ddd-kernel/adapters/pubsub'; + +kernel.registerConsumerMiddleware( + new CorrelationConsumerMiddleware(), + new IdempotencyConsumerMiddleware({ + store: new InMemoryIdempotencyStore(), + }), + new RetryConsumerMiddleware({ + maxAttempts: 3, + }), +); +``` + +Use a custom `IdempotencyStore` for durable idempotency. Prefer stores that +implement atomic `claim`, `commit` and `release` methods so duplicate messages +cannot pass a non-atomic `has`/`mark` check concurrently. The in-memory store is +only useful for tests and single-process applications. diff --git a/docs/reference/express-kernel-server.md b/docs/reference/express-kernel-server.md index 2d7ee40..92d761e 100644 --- a/docs/reference/express-kernel-server.md +++ b/docs/reference/express-kernel-server.md @@ -11,3 +11,56 @@ await server.run(); ``` Routes are registered with `kernel.registerRoutes(RouteClass)`. + +## External Controllers + +Applications can add controllers at the server boundary without registering them +on the kernel: + +```ts +const server = new ExpressKernelServer({ + controllers: [HealthController], + kernel, + port: 3000, +}); +``` + +`controllers` are merged with `kernel.getRoutes()` before +`routing-controllers` is configured. + +## HTTP Middleware And Hooks + +Use middleware arrays for normal Express middleware and hooks for integrations +that need direct app access, such as Swagger or static assets: + +```ts +const server = new ExpressKernelServer({ + kernel, + hooks: [ + { phase: 'beforeControllers', handle: setupTracing }, + { phase: 'beforeErrors', handle: setupSwagger }, + { phase: 'beforeErrors', handle: setupStaticAssets }, + ], + middlewares: [requestIdMiddleware], + preControllerMiddlewares: [authenticationMiddleware], + postControllerMiddlewares: [notFoundMiddleware], +}); +``` + +Hook order is: + +1. `middlewares` +2. `preControllerMiddlewares` +3. `beforeControllersHooks` +4. `hooks` with `phase: 'beforeControllers'` +5. `routing-controllers` +6. `postControllerMiddlewares` +7. `afterControllersHooks` +8. `hooks` with `phase: 'afterControllers'` +9. `swaggerHooks` +10. `staticHooks` +11. `hooks` with `phase: 'beforeErrors'` +12. `errorHandlers` + +`swaggerHooks` and `staticHooks` remain available for compatibility. New +integrations should use `hooks` with an explicit phase. diff --git a/docs/reference/scheduler.md b/docs/reference/scheduler.md index 41e2e98..4fbee4a 100644 --- a/docs/reference/scheduler.md +++ b/docs/reference/scheduler.md @@ -16,9 +16,12 @@ Register scheduler classes with `kernel.registerSchedulers(...)`. ## Error Policy -Schedulers accept a `SchedulerErrorPolicy`: +Schedulers accept a `SchedulerErrorPolicy` exported from +`@haskou/ddd-kernel/scheduler`: ```ts +import type { SchedulerErrorPolicy } from '@haskou/ddd-kernel/scheduler'; + class ReplicationScheduler extends Scheduler { constructor(errorPolicy: SchedulerErrorPolicy) { super(errorPolicy); @@ -35,3 +38,19 @@ interface SchedulerErrorPolicy { handle(error: unknown, scheduler: Scheduler): Promise | void; } ``` + +Use `shouldSkip` for domain-specific transient states that should not be logged +as scheduler failures, for example replicated state that is not ready yet: + +```ts +const policy: SchedulerErrorPolicy = { + shouldSkip(error) { + return error instanceof ReplicatedStateNotReadyError; + }, + handle(error, scheduler) { + logger.error(`${scheduler.getProcessName()} failed: ${String(error)}`); + }, +}; +``` + +The default policy never skips and wraps failures in `ScheduledExecutionError`. diff --git a/package.json b/package.json index cdfbf51..b56c842 100644 --- a/package.json +++ b/package.json @@ -6,6 +6,85 @@ "main": "./dist/index.cjs", "module": "./dist/index.js", "types": "./dist/index.d.ts", + "typesVersions": { + "*": { + "adapters": [ + "dist/adapters/index.d.ts" + ], + "adapters/db": [ + "dist/adapters/db/index.d.ts" + ], + "adapters/db/in-memory": [ + "dist/adapters/db/in-memory/index.d.ts" + ], + "adapters/db/mongo": [ + "dist/adapters/db/mongo/index.d.ts" + ], + "adapters/kernel": [ + "dist/adapters/kernel/index.d.ts" + ], + "adapters/kernel/console": [ + "dist/adapters/kernel/console/index.d.ts" + ], + "adapters/pubsub": [ + "dist/adapters/pubsub/index.d.ts" + ], + "adapters/pubsub/amqp": [ + "dist/adapters/pubsub/amqp/index.d.ts" + ], + "adapters/pubsub/in-memory": [ + "dist/adapters/pubsub/in-memory/index.d.ts" + ], + "adapters/ui": [ + "dist/adapters/ui/index.d.ts" + ], + "adapters/ui/express": [ + "dist/adapters/ui/express/index.d.ts" + ], + "adapters/ui/routes": [ + "dist/adapters/ui/routes/index.d.ts" + ], + "contracts": [ + "dist/contracts/index.d.ts" + ], + "contracts/db": [ + "dist/contracts/db/index.d.ts" + ], + "contracts/kernel": [ + "dist/contracts/kernel/index.d.ts" + ], + "contracts/pubsub": [ + "dist/contracts/pubsub/index.d.ts" + ], + "contracts/ui": [ + "dist/contracts/ui/index.d.ts" + ], + "dependency-injection": [ + "dist/infrastructure/dependency-injection/index.d.ts" + ], + "domain": [ + "dist/domain/index.d.ts" + ], + "errors": [ + "dist/errors/index.d.ts" + ], + "express": [ + "dist/adapters/ui/express/index.d.ts" + ], + "lifecycle": [ + "dist/infrastructure/lifecycle/index.d.ts" + ], + "logs": [ + "dist/infrastructure/logs/index.d.ts" + ], + "scheduler": [ + "dist/infrastructure/scheduler/index.d.ts" + ], + "websocket": [ + "dist/infrastructure/websocket/index.d.ts" + ] + } + }, "publishConfig": { "access": "public" }, diff --git a/src/Kernel.ts b/src/Kernel.ts index 015f849..6c98882 100644 --- a/src/Kernel.ts +++ b/src/Kernel.ts @@ -61,6 +61,10 @@ export class Kernel { return Kernel.getActiveKernel().logger; } + public static get active(): Kernel { + return Kernel.getActiveKernel(); + } + public static get rootDirectory(): string { return process.cwd(); } diff --git a/src/adapters/pubsub/Consumer.ts b/src/adapters/pubsub/Consumer.ts index 6e24d3f..4f1cc35 100644 --- a/src/adapters/pubsub/Consumer.ts +++ b/src/adapters/pubsub/Consumer.ts @@ -1,27 +1,36 @@ -import type { ConsumerMiddleware } from '../../contracts/index.js'; import type { DomainEventConsumer } from '../../domain/DomainEventConsumer.js'; -import type { DomainEvent } from '../../domain/index.js'; +import type { + DomainEvent, + DomainEventConsumerContext, +} from '../../domain/index.js'; import { Kernel } from '../../Kernel.js'; +import { ConsumerMiddlewarePipeline } from './ConsumerMiddlewarePipeline.js'; export abstract class Consumer { constructor(private readonly consumer: DomainEventConsumer) {} private async runMiddleware( event: DomainEvent, - middlewares: readonly ConsumerMiddleware[], - index: number, + consumerContext?: DomainEventConsumerContext, ): Promise { - const middleware = middlewares[index]; - - if (!middleware) { - await this.handler(event); - - return; - } - - await middleware.handle(event, () => - this.runMiddleware(event, middlewares, index + 1), + const pipeline = new ConsumerMiddlewarePipeline(Kernel.consumerMiddleware); + const metadata = consumerContext?.metadata ?? {}; + + await pipeline.execute( + event, + { + causationId: event.getCausationId(), + correlationId: event.getCorrelationId(), + eventId: event.eventId, + eventName: this.eventName, + exchange: this.exchange, + kernel: Kernel.active, + metadata, + queueName: this.queueName, + rawMessage: metadata.rawMessage, + }, + () => this.handler(event), ); } @@ -41,7 +50,7 @@ export abstract class Consumer { this.eventName, this.domainEvent, this.exchange, - (event) => this.runMiddleware(event, Kernel.consumerMiddleware, 0), + (event, context) => this.runMiddleware(event, context), ); } diff --git a/src/adapters/pubsub/ConsumerMiddlewarePipeline.ts b/src/adapters/pubsub/ConsumerMiddlewarePipeline.ts new file mode 100644 index 0000000..7fe83d2 --- /dev/null +++ b/src/adapters/pubsub/ConsumerMiddlewarePipeline.ts @@ -0,0 +1,38 @@ +import type { + ConsumerExecutionContext, + ConsumerMiddleware, +} from '../../contracts/index.js'; +import type { DomainEvent } from '../../domain/index.js'; + +export class ConsumerMiddlewarePipeline { + constructor(private readonly middlewares: readonly ConsumerMiddleware[]) {} + + private async run( + event: DomainEvent, + context: ConsumerExecutionContext, + handler: () => Promise, + index: number, + ): Promise { + const middleware = this.middlewares[index]; + + if (!middleware) { + await handler(); + + return; + } + + await middleware.handle( + event, + () => this.run(event, context, handler, index + 1), + context, + ); + } + + public async execute( + event: DomainEvent, + context: ConsumerExecutionContext, + handler: () => Promise, + ): Promise { + await this.run(event, context, handler, 0); + } +} diff --git a/src/adapters/pubsub/CorrelationConsumerMiddleware.ts b/src/adapters/pubsub/CorrelationConsumerMiddleware.ts new file mode 100644 index 0000000..4551f3a --- /dev/null +++ b/src/adapters/pubsub/CorrelationConsumerMiddleware.ts @@ -0,0 +1,36 @@ +import type { + ConsumerExecutionContext, + ConsumerMiddleware, + ConsumerNext, +} from '../../contracts/index.js'; +import type { DomainEvent } from '../../domain/index.js'; +import type { CorrelationConsumerMiddlewareOptions } from './CorrelationConsumerMiddlewareOptions.js'; + +export class CorrelationConsumerMiddleware implements ConsumerMiddleware { + constructor( + private readonly options: CorrelationConsumerMiddlewareOptions = {}, + ) {} + + public async handle( + event: DomainEvent, + next: ConsumerNext, + context: ConsumerExecutionContext, + ): Promise { + const correlationId = + this.options.correlationId?.(event, context) ?? context.correlationId; + const causationId = + this.options.causationId?.(event, context) ?? context.causationId; + + if (correlationId) { + event.withCorrelationId(correlationId); + } + + if (causationId) { + event.withCausationId(causationId); + } + + await next(); + } +} + +export default CorrelationConsumerMiddleware; diff --git a/src/adapters/pubsub/CorrelationConsumerMiddlewareOptions.ts b/src/adapters/pubsub/CorrelationConsumerMiddlewareOptions.ts new file mode 100644 index 0000000..7cace69 --- /dev/null +++ b/src/adapters/pubsub/CorrelationConsumerMiddlewareOptions.ts @@ -0,0 +1,13 @@ +import type { ConsumerExecutionContext } from '../../contracts/index.js'; +import type { DomainEvent } from '../../domain/index.js'; + +export interface CorrelationConsumerMiddlewareOptions { + readonly causationId?: ( + event: DomainEvent, + context: ConsumerExecutionContext, + ) => string | undefined; + readonly correlationId?: ( + event: DomainEvent, + context: ConsumerExecutionContext, + ) => string | undefined; +} diff --git a/src/adapters/pubsub/DefaultPublisherHookErrorPolicy.ts b/src/adapters/pubsub/DefaultPublisherHookErrorPolicy.ts new file mode 100644 index 0000000..1c85442 --- /dev/null +++ b/src/adapters/pubsub/DefaultPublisherHookErrorPolicy.ts @@ -0,0 +1,26 @@ +import type { + PublishContext, + PublisherHookErrorPolicy, +} from '../../contracts/index.js'; + +export class DefaultPublisherHookErrorPolicy implements PublisherHookErrorPolicy { + public handleAfterPublishError( + error: unknown, + context: PublishContext, + ): void { + void error; + void context; + } + + public shouldFailAfterPublish( + error: unknown, + context: PublishContext, + ): boolean { + void error; + void context; + + return false; + } +} + +export default DefaultPublisherHookErrorPolicy; diff --git a/src/adapters/pubsub/IdempotencyConsumerMiddleware.ts b/src/adapters/pubsub/IdempotencyConsumerMiddleware.ts new file mode 100644 index 0000000..f460eeb --- /dev/null +++ b/src/adapters/pubsub/IdempotencyConsumerMiddleware.ts @@ -0,0 +1,59 @@ +import type { + ConsumerExecutionContext, + ConsumerMiddleware, + ConsumerNext, +} from '../../contracts/index.js'; +import type { DomainEvent } from '../../domain/index.js'; +import type { IdempotencyConsumerMiddlewareOptions } from './IdempotencyConsumerMiddlewareOptions.js'; + +export class IdempotencyConsumerMiddleware implements ConsumerMiddleware { + constructor(private readonly options: IdempotencyConsumerMiddlewareOptions) {} + + private async handleClaimedKey( + key: string, + next: ConsumerNext, + ): Promise { + const claimed = await this.options.store.claim?.(key); + + if (!claimed) { + return; + } + + try { + await next(); + await (this.options.store.commit?.(key) ?? this.options.store.mark(key)); + } catch (error: unknown) { + await this.options.store.release?.(key); + + throw error; + } + } + + private async handleLegacyKey( + key: string, + next: ConsumerNext, + ): Promise { + if (await this.options.store.has(key)) { + return; + } + + await next(); + await this.options.store.mark(key); + } + + public async handle( + event: DomainEvent, + next: ConsumerNext, + context: ConsumerExecutionContext, + ): Promise { + const key = this.options.key?.(event, context) ?? context.eventId; + + if (this.options.store.claim) { + await this.handleClaimedKey(key, next); + } else { + await this.handleLegacyKey(key, next); + } + } +} + +export default IdempotencyConsumerMiddleware; diff --git a/src/adapters/pubsub/IdempotencyConsumerMiddlewareOptions.ts b/src/adapters/pubsub/IdempotencyConsumerMiddlewareOptions.ts new file mode 100644 index 0000000..7d04f9d --- /dev/null +++ b/src/adapters/pubsub/IdempotencyConsumerMiddlewareOptions.ts @@ -0,0 +1,13 @@ +import type { + ConsumerExecutionContext, + IdempotencyStore, +} from '../../contracts/index.js'; +import type { DomainEvent } from '../../domain/index.js'; + +export interface IdempotencyConsumerMiddlewareOptions { + readonly key?: ( + event: DomainEvent, + context: ConsumerExecutionContext, + ) => string; + readonly store: IdempotencyStore; +} diff --git a/src/adapters/pubsub/InMemoryIdempotencyStore.ts b/src/adapters/pubsub/InMemoryIdempotencyStore.ts new file mode 100644 index 0000000..1201dcf --- /dev/null +++ b/src/adapters/pubsub/InMemoryIdempotencyStore.ts @@ -0,0 +1,35 @@ +import type { IdempotencyStore } from '../../contracts/index.js'; + +export class InMemoryIdempotencyStore implements IdempotencyStore { + private readonly claimedKeys = new Set(); + private readonly handledKeys = new Set(); + + public claim(key: string): boolean { + if (this.handledKeys.has(key) || this.claimedKeys.has(key)) { + return false; + } + + this.claimedKeys.add(key); + + return true; + } + + public commit(key: string): void { + this.claimedKeys.delete(key); + this.handledKeys.add(key); + } + + public release(key: string): void { + this.claimedKeys.delete(key); + } + + public has(key: string): boolean { + return this.handledKeys.has(key); + } + + public mark(key: string): void { + this.handledKeys.add(key); + } +} + +export default InMemoryIdempotencyStore; diff --git a/src/adapters/pubsub/PublisherHookPipeline.ts b/src/adapters/pubsub/PublisherHookPipeline.ts new file mode 100644 index 0000000..6fecddf --- /dev/null +++ b/src/adapters/pubsub/PublisherHookPipeline.ts @@ -0,0 +1,77 @@ +import type { + PublishContext, + PublisherHook, + PublisherHookErrorPolicy, +} from '../../contracts/index.js'; + +import { DefaultPublisherHookErrorPolicy } from './DefaultPublisherHookErrorPolicy.js'; + +export class PublisherHookPipeline { + private readonly hooks: PublisherHook[] = []; + + constructor( + hooks: readonly PublisherHook[] = [], + private readonly errorPolicy: PublisherHookErrorPolicy = new DefaultPublisherHookErrorPolicy(), + ) { + this.hooks.push(...hooks); + } + + private async runAfterPublishHooks(context: PublishContext): Promise { + for (const hook of this.hooks) { + await this.runAfterPublishHook(hook, context); + } + } + + private async runAfterPublishHook( + hook: PublisherHook, + context: PublishContext, + ): Promise { + try { + await hook.afterPublish?.(context); + } catch (error: unknown) { + await this.errorPolicy.handleAfterPublishError(error, context); + + if (this.errorPolicy.shouldFailAfterPublish(error, context)) { + throw error; + } + } + } + + private async runBeforePublishHooks(context: PublishContext): Promise { + for (const hook of this.hooks) { + await hook.beforePublish?.(context); + } + } + + private async runPublishErrorHooks( + error: unknown, + context: PublishContext, + ): Promise { + for (const hook of this.hooks) { + await hook.onPublishError?.(error, context); + } + } + + public register(...hooks: PublisherHook[]): void { + this.hooks.push(...hooks); + } + + public async run( + context: PublishContext, + publish: () => Promise | T, + ): Promise { + await this.runBeforePublishHooks(context); + + try { + const result = await publish(); + + await this.runAfterPublishHooks(context); + + return result; + } catch (error: unknown) { + await this.runPublishErrorHooks(error, context); + + throw error; + } + } +} diff --git a/src/adapters/pubsub/RetryConsumerMiddleware.ts b/src/adapters/pubsub/RetryConsumerMiddleware.ts new file mode 100644 index 0000000..7509275 --- /dev/null +++ b/src/adapters/pubsub/RetryConsumerMiddleware.ts @@ -0,0 +1,53 @@ +import type { + ConsumerExecutionContext, + ConsumerMiddleware, + ConsumerNext, +} from '../../contracts/index.js'; +import type { DomainEvent } from '../../domain/index.js'; +import type { RetryConsumerMiddlewareOptions } from './RetryConsumerMiddlewareOptions.js'; + +export class RetryConsumerMiddleware implements ConsumerMiddleware { + constructor(private readonly options: RetryConsumerMiddlewareOptions) {} + + private async delay( + attempt: number, + error: unknown, + context: ConsumerExecutionContext, + ): Promise { + const delayInMilliseconds = + typeof this.options.delay === 'function' + ? await this.options.delay(attempt, error, context) + : (this.options.delay ?? 0); + + if (delayInMilliseconds > 0) { + await new Promise((resolve) => setTimeout(resolve, delayInMilliseconds)); + } + } + + public async handle( + event: DomainEvent, + next: ConsumerNext, + context: ConsumerExecutionContext, + ): Promise { + for (let attempt = 1; attempt <= this.options.maxAttempts; attempt++) { + try { + await next(); + + return; + } catch (error: unknown) { + const canRetry = + attempt < this.options.maxAttempts && + (await (this.options.shouldRetry?.(error, attempt, context) ?? true)); + + if (!canRetry) { + throw error; + } + + await this.options.onRetry?.(error, attempt, context); + await this.delay(attempt, error, context); + } + } + } +} + +export default RetryConsumerMiddleware; diff --git a/src/adapters/pubsub/RetryConsumerMiddlewareOptions.ts b/src/adapters/pubsub/RetryConsumerMiddlewareOptions.ts new file mode 100644 index 0000000..5b70aa8 --- /dev/null +++ b/src/adapters/pubsub/RetryConsumerMiddlewareOptions.ts @@ -0,0 +1,16 @@ +import type { + ConsumerExecutionContext, + RetryDelayResolver, + RetryPredicate, +} from '../../contracts/index.js'; + +export interface RetryConsumerMiddlewareOptions { + readonly delay?: number | RetryDelayResolver; + readonly maxAttempts: number; + readonly onRetry?: ( + error: unknown, + attempt: number, + context: ConsumerExecutionContext, + ) => Promise | void; + readonly shouldRetry?: RetryPredicate; +} diff --git a/src/adapters/pubsub/amqp/AmqpMessageBusAdapter.ts b/src/adapters/pubsub/amqp/AmqpMessageBusAdapter.ts index ce90763..389c4ab 100644 --- a/src/adapters/pubsub/amqp/AmqpMessageBusAdapter.ts +++ b/src/adapters/pubsub/amqp/AmqpMessageBusAdapter.ts @@ -8,10 +8,12 @@ import amqplib, { } from 'amqplib'; import { randomUUID } from 'node:crypto'; +import type { PublisherHook } from '../../../contracts/index.js'; import type { Constructor, DomainEvent, DomainEventConsumer, + DomainMessageBus, DomainEventPublisher, } from '../../../domain/index.js'; import type { AmqpMessage } from './AmqpMessage.js'; @@ -19,20 +21,26 @@ import type { AmqpMessageBusAdapterOptions } from './AmqpMessageBusAdapterOption import type { ConsumerContext } from './ConsumerContext.js'; import { Kernel } from '../../../Kernel.js'; +import { PublisherHookPipeline } from '../PublisherHookPipeline.js'; import { InvalidDomainEventError } from './InvalidDomainEventError.js'; import { NoFailedMessagesError } from './NoFailedMessagesError.js'; export default class AmqpMessageBusAdapter - implements DomainEventConsumer, DomainEventPublisher + implements DomainEventConsumer, DomainEventPublisher, DomainMessageBus { private channelInstance: Channel | undefined; private connection: ChannelModel | undefined; private readonly delayConsumers: string[] = []; private exchange: string; + private readonly publisherHookPipeline: PublisherHookPipeline; constructor(private readonly options: AmqpMessageBusAdapterOptions = {}) { this.exchange = options.exchange ?? options.serviceName ?? process.env.SERVICE_NAME ?? ''; + this.publisherHookPipeline = new PublisherHookPipeline( + options.publisherHooks, + options.publisherHookErrorPolicy, + ); } private get dsn(): string { @@ -97,7 +105,13 @@ export default class AmqpMessageBusAdapter message, ); - await context.handler(domainEvent); + await context.handler(domainEvent, { + metadata: { + headers: msg.properties.headers ?? {}, + rawMessage: msg, + retries: Number(msg.properties.headers?.retries ?? 0), + }, + }); } catch (error) { await this.handleError(msg, message, context, error); } @@ -450,15 +464,39 @@ export default class AmqpMessageBusAdapter const channel = await this.channel(); for (const event of domainEvents) { - channel.publish( - this.exchange, - event.eventName(), - Buffer.from(event.decode()), - this.opts(event), + await this.publisherHookPipeline.run( + { + domainEvent: event, + message: { + metadata: { + causationId: event.getCausationId(), + correlationId: event.getCorrelationId(), + eventId: event.eventId, + }, + name: event.eventName(), + payload: event.attributes, + }, + metadata: { + eventId: event.eventId, + exchange: this.exchange, + }, + topic: event.eventName(), + }, + () => + channel.publish( + this.exchange, + event.eventName(), + Buffer.from(event.decode()), + this.opts(event), + ), ); } } + public registerPublisherHooks(...hooks: PublisherHook[]): void { + this.publisherHookPipeline.register(...hooks); + } + public async close(): Promise { await this.channelInstance?.close(); await this.connection?.close(); diff --git a/src/adapters/pubsub/amqp/AmqpMessageBusAdapterOptions.ts b/src/adapters/pubsub/amqp/AmqpMessageBusAdapterOptions.ts index 084c3f7..0b78792 100644 --- a/src/adapters/pubsub/amqp/AmqpMessageBusAdapterOptions.ts +++ b/src/adapters/pubsub/amqp/AmqpMessageBusAdapterOptions.ts @@ -1,3 +1,7 @@ +import type { + PublisherHook, + PublisherHookErrorPolicy, +} from '../../../contracts/index.js'; import type { Log } from '../../../infrastructure/logs/index.js'; export interface AmqpMessageBusAdapterOptions { @@ -5,6 +9,8 @@ export interface AmqpMessageBusAdapterOptions { readonly exchange?: string; readonly logger?: Log; readonly maxRetries?: number; + readonly publisherHookErrorPolicy?: PublisherHookErrorPolicy; + readonly publisherHooks?: PublisherHook[]; readonly retryDelayInMilliseconds?: number; readonly serviceName?: string; } diff --git a/src/adapters/pubsub/amqp/DomainEventHandler.ts b/src/adapters/pubsub/amqp/DomainEventHandler.ts index 0be849d..bdb7420 100644 --- a/src/adapters/pubsub/amqp/DomainEventHandler.ts +++ b/src/adapters/pubsub/amqp/DomainEventHandler.ts @@ -1,3 +1,9 @@ -import type { DomainEvent } from '../../../domain/index.js'; +import type { + DomainEvent, + DomainEventConsumerContext, +} from '../../../domain/index.js'; -export type DomainEventHandler = (event: DomainEvent) => Promise; +export type DomainEventHandler = ( + event: DomainEvent, + context?: DomainEventConsumerContext, +) => Promise; diff --git a/src/adapters/pubsub/in-memory/InMemoryEventBus.ts b/src/adapters/pubsub/in-memory/InMemoryEventBus.ts index f55fd09..cef6322 100644 --- a/src/adapters/pubsub/in-memory/InMemoryEventBus.ts +++ b/src/adapters/pubsub/in-memory/InMemoryEventBus.ts @@ -2,15 +2,30 @@ import type { DomainEvent, HandlerContext, MessageHandler, + PublisherHook, + PublisherHookErrorPolicy, } from '../../../contracts/index.js'; +import { PublisherHookPipeline } from '../PublisherHookPipeline.js'; + export class InMemoryEventBus { private readonly handlers = new Map< string, MessageHandler[] >(); - constructor(private readonly context: HandlerContext) {} + private readonly publisherHookPipeline: PublisherHookPipeline; + + constructor( + private readonly context: HandlerContext, + publisherHooks: readonly PublisherHook[] = [], + publisherHookErrorPolicy?: PublisherHookErrorPolicy, + ) { + this.publisherHookPipeline = new PublisherHookPipeline( + publisherHooks, + publisherHookErrorPolicy, + ); + } public subscribe( name: TEvent['name'], @@ -27,8 +42,17 @@ export class InMemoryEventBus { ): Promise { const handlers = this.handlers.get(event.name) ?? []; - for (const handler of handlers) { - await handler(event, this.context); - } + await this.publisherHookPipeline.run( + { message: event, metadata: event.metadata ?? {}, topic: event.name }, + async () => { + for (const handler of handlers) { + await handler(event, this.context); + } + }, + ); + } + + public registerPublisherHooks(...hooks: PublisherHook[]): void { + this.publisherHookPipeline.register(...hooks); } } diff --git a/src/adapters/pubsub/in-memory/InMemoryPubSub.ts b/src/adapters/pubsub/in-memory/InMemoryPubSub.ts index 9e050e4..f8205f9 100644 --- a/src/adapters/pubsub/in-memory/InMemoryPubSub.ts +++ b/src/adapters/pubsub/in-memory/InMemoryPubSub.ts @@ -2,16 +2,31 @@ import type { HandlerContext, Message, MessageHandler, + PublisherHook, + PublisherHookErrorPolicy, Subscription, } from '../../../contracts/index.js'; +import { PublisherHookPipeline } from '../PublisherHookPipeline.js'; + export class InMemoryPubSub { private readonly consumers = new Map< string, Set> >(); - constructor(private readonly context: HandlerContext) {} + private readonly publisherHookPipeline: PublisherHookPipeline; + + constructor( + private readonly context: HandlerContext, + publisherHooks: readonly PublisherHook[] = [], + publisherHookErrorPolicy?: PublisherHookErrorPolicy, + ) { + this.publisherHookPipeline = new PublisherHookPipeline( + publisherHooks, + publisherHookErrorPolicy, + ); + } public async publish( topic: string, @@ -19,9 +34,14 @@ export class InMemoryPubSub { ): Promise { const consumers = this.consumers.get(topic) ?? new Set(); - for (const consumer of consumers) { - await consumer(message, this.context); - } + await this.publisherHookPipeline.run( + { message, metadata: message.metadata ?? {}, topic }, + async () => { + for (const consumer of consumers) { + await consumer(message, this.context); + } + }, + ); } public subscribe( @@ -43,4 +63,8 @@ export class InMemoryPubSub { }, }); } + + public registerPublisherHooks(...hooks: PublisherHook[]): void { + this.publisherHookPipeline.register(...hooks); + } } diff --git a/src/adapters/pubsub/index.ts b/src/adapters/pubsub/index.ts index e9c602a..641ce30 100644 --- a/src/adapters/pubsub/index.ts +++ b/src/adapters/pubsub/index.ts @@ -1,4 +1,14 @@ export * from './Consumer.js'; +export * from './ConsumerMiddlewarePipeline.js'; +export * from './CorrelationConsumerMiddleware.js'; +export * from './CorrelationConsumerMiddlewareOptions.js'; +export * from './DefaultPublisherHookErrorPolicy.js'; +export * from './IdempotencyConsumerMiddleware.js'; +export * from './IdempotencyConsumerMiddlewareOptions.js'; +export * from './InMemoryIdempotencyStore.js'; +export * from './PublisherHookPipeline.js'; +export * from './RetryConsumerMiddleware.js'; +export * from './RetryConsumerMiddlewareOptions.js'; export * from './amqp/index.js'; export * from './in-memory/index.js'; export { default } from './Consumer.js'; diff --git a/src/adapters/ui/express/ExpressAppHook.ts b/src/adapters/ui/express/ExpressAppHook.ts new file mode 100644 index 0000000..1fab89c --- /dev/null +++ b/src/adapters/ui/express/ExpressAppHook.ts @@ -0,0 +1,3 @@ +import type { HttpApp } from './HttpApp.js'; + +export type ExpressAppHook = (app: HttpApp) => Promise | void; diff --git a/src/adapters/ui/express/ExpressController.ts b/src/adapters/ui/express/ExpressController.ts new file mode 100644 index 0000000..6953a5a --- /dev/null +++ b/src/adapters/ui/express/ExpressController.ts @@ -0,0 +1 @@ +export type ExpressController = abstract new (...args: unknown[]) => unknown; diff --git a/src/adapters/ui/express/ExpressHookPhase.ts b/src/adapters/ui/express/ExpressHookPhase.ts new file mode 100644 index 0000000..881abf3 --- /dev/null +++ b/src/adapters/ui/express/ExpressHookPhase.ts @@ -0,0 +1,4 @@ +export type ExpressHookPhase = + | 'afterControllers' + | 'beforeControllers' + | 'beforeErrors'; diff --git a/src/adapters/ui/express/ExpressKernelServer.ts b/src/adapters/ui/express/ExpressKernelServer.ts index 590966c..4331757 100644 --- a/src/adapters/ui/express/ExpressKernelServer.ts +++ b/src/adapters/ui/express/ExpressKernelServer.ts @@ -1,6 +1,8 @@ -import type { ErrorRequestHandler } from 'express'; - -import { createExpressServer } from 'routing-controllers'; +import express, { + type ErrorRequestHandler, + type RequestHandler, +} from 'express'; +import { useExpressServer } from 'routing-controllers'; import type { ExpressKernelServerOptions } from './ExpressKernelServerOptions.js'; import type { HttpApp } from './HttpApp.js'; @@ -30,6 +32,35 @@ export class ExpressKernelServer { }; } + private async runHooks( + hooks: readonly ((app: HttpApp) => Promise | void)[] | undefined, + app: HttpApp, + ): Promise { + for (const hook of hooks ?? []) { + await hook(app); + } + } + + private async runPhaseHooks( + phase: 'afterControllers' | 'beforeControllers' | 'beforeErrors', + app: HttpApp, + ): Promise { + for (const hook of this.options.hooks ?? []) { + if (hook.phase === phase) { + await hook.handle(app); + } + } + } + + private registerMiddlewares( + app: HttpApp, + middlewares: readonly RequestHandler[] | undefined, + ): void { + for (const middleware of middlewares ?? []) { + app.use(middleware); + } + } + public get app(): HttpApp { if (!this.appInstance) { throw new Error('HTTP server is not running.'); @@ -66,16 +97,32 @@ export class ExpressKernelServer { }); } - public run(): Promise { - const app = createExpressServer({ - controllers: this.options.kernel.getRoutes(), - routePrefix: this.options.routePrefix, - }) as HttpApp; - - for (const middleware of this.options.middlewares ?? []) { - app.use(middleware); + public async run(): Promise { + if (this.serverInstance) { + throw new Error('HTTP server is already running.'); } + const controllers = [ + ...this.options.kernel.getRoutes(), + ...(this.options.controllers ?? []), + ]; + const app = express() as HttpApp; + + this.registerMiddlewares(app, this.options.middlewares); + this.registerMiddlewares(app, this.options.preControllerMiddlewares); + await this.runHooks(this.options.beforeControllersHooks, app); + await this.runPhaseHooks('beforeControllers', app); + useExpressServer(app, { + controllers, + routePrefix: this.options.routePrefix, + }); + this.registerMiddlewares(app, this.options.postControllerMiddlewares); + await this.runHooks(this.options.afterControllersHooks, app); + await this.runPhaseHooks('afterControllers', app); + await this.runHooks(this.options.swaggerHooks, app); + await this.runHooks(this.options.staticHooks, app); + await this.runPhaseHooks('beforeErrors', app); + this.registerErrorHandlers(app); this.appInstance = app; diff --git a/src/adapters/ui/express/ExpressKernelServerOptions.ts b/src/adapters/ui/express/ExpressKernelServerOptions.ts index d4cbbb5..1e5f8b8 100644 --- a/src/adapters/ui/express/ExpressKernelServerOptions.ts +++ b/src/adapters/ui/express/ExpressKernelServerOptions.ts @@ -1,11 +1,28 @@ import type { ErrorRequestHandler, RequestHandler } from 'express'; import type { Kernel } from '../../../Kernel.js'; +import type { ExpressAppHook } from './ExpressAppHook.js'; +import type { ExpressController } from './ExpressController.js'; +import type { ExpressPhaseHook } from './ExpressPhaseHook.js'; export interface ExpressKernelServerOptions { + readonly afterControllersHooks?: ExpressAppHook[]; + readonly beforeControllersHooks?: ExpressAppHook[]; + readonly controllers?: ExpressController[]; readonly errorHandlers?: ErrorRequestHandler[]; + readonly hooks?: ExpressPhaseHook[]; readonly kernel: Kernel; readonly middlewares?: RequestHandler[]; + readonly postControllerMiddlewares?: RequestHandler[]; + readonly preControllerMiddlewares?: RequestHandler[]; readonly port?: number; readonly routePrefix?: string; + /** + * @deprecated Prefer `hooks` with `phase: 'beforeErrors'`. + */ + readonly staticHooks?: ExpressAppHook[]; + /** + * @deprecated Prefer `hooks` with `phase: 'beforeErrors'`. + */ + readonly swaggerHooks?: ExpressAppHook[]; } diff --git a/src/adapters/ui/express/ExpressPhaseHook.ts b/src/adapters/ui/express/ExpressPhaseHook.ts new file mode 100644 index 0000000..77cc725 --- /dev/null +++ b/src/adapters/ui/express/ExpressPhaseHook.ts @@ -0,0 +1,7 @@ +import type { ExpressAppHook } from './ExpressAppHook.js'; +import type { ExpressHookPhase } from './ExpressHookPhase.js'; + +export interface ExpressPhaseHook { + readonly handle: ExpressAppHook; + readonly phase: ExpressHookPhase; +} diff --git a/src/adapters/ui/express/index.ts b/src/adapters/ui/express/index.ts index cf5fbf6..9df26a1 100644 --- a/src/adapters/ui/express/index.ts +++ b/src/adapters/ui/express/index.ts @@ -1,5 +1,9 @@ +export * from './ExpressAppHook.js'; +export * from './ExpressController.js'; +export * from './ExpressHookPhase.js'; export * from './ExpressKernelServer.js'; export * from './ExpressKernelServerOptions.js'; +export * from './ExpressPhaseHook.js'; export * from './HttpApp.js'; export * from './HttpServer.js'; export * from './RoutePrefix.js'; diff --git a/src/contracts/kernel/ConsumerExecutionContext.ts b/src/contracts/kernel/ConsumerExecutionContext.ts new file mode 100644 index 0000000..65d6032 --- /dev/null +++ b/src/contracts/kernel/ConsumerExecutionContext.ts @@ -0,0 +1,13 @@ +import type { Kernel } from '../../Kernel.js'; + +export interface ConsumerExecutionContext { + readonly causationId?: string; + readonly correlationId?: string; + readonly eventId: string; + readonly eventName: string; + readonly exchange: string; + readonly kernel: Kernel; + readonly metadata: Readonly>; + readonly rawMessage?: unknown; + readonly queueName: string; +} diff --git a/src/contracts/kernel/ConsumerMiddleware.ts b/src/contracts/kernel/ConsumerMiddleware.ts index 5f8d648..fa2ba96 100644 --- a/src/contracts/kernel/ConsumerMiddleware.ts +++ b/src/contracts/kernel/ConsumerMiddleware.ts @@ -1,5 +1,11 @@ import type { DomainEvent } from '../../domain/index.js'; +import type { ConsumerExecutionContext } from './ConsumerExecutionContext.js'; +import type { ConsumerNext } from './ConsumerNext.js'; export interface ConsumerMiddleware { - handle(event: DomainEvent, next: () => Promise): Promise; + handle( + event: DomainEvent, + next: ConsumerNext, + context: ConsumerExecutionContext, + ): Promise; } diff --git a/src/contracts/kernel/ConsumerNext.ts b/src/contracts/kernel/ConsumerNext.ts new file mode 100644 index 0000000..ace0f5d --- /dev/null +++ b/src/contracts/kernel/ConsumerNext.ts @@ -0,0 +1 @@ +export type ConsumerNext = () => Promise; diff --git a/src/contracts/kernel/IdempotencyStore.ts b/src/contracts/kernel/IdempotencyStore.ts new file mode 100644 index 0000000..df6b3b7 --- /dev/null +++ b/src/contracts/kernel/IdempotencyStore.ts @@ -0,0 +1,7 @@ +export interface IdempotencyStore { + claim?(key: string): Promise | boolean; + commit?(key: string): Promise | void; + release?(key: string): Promise | void; + has(key: string): Promise | boolean; + mark(key: string): Promise | void; +} diff --git a/src/contracts/kernel/RetryDelayResolver.ts b/src/contracts/kernel/RetryDelayResolver.ts new file mode 100644 index 0000000..f9bb7b3 --- /dev/null +++ b/src/contracts/kernel/RetryDelayResolver.ts @@ -0,0 +1,7 @@ +import type { ConsumerExecutionContext } from './ConsumerExecutionContext.js'; + +export type RetryDelayResolver = ( + attempt: number, + error: unknown, + context: ConsumerExecutionContext, +) => number | Promise; diff --git a/src/contracts/kernel/RetryPredicate.ts b/src/contracts/kernel/RetryPredicate.ts new file mode 100644 index 0000000..3188cf5 --- /dev/null +++ b/src/contracts/kernel/RetryPredicate.ts @@ -0,0 +1,7 @@ +import type { ConsumerExecutionContext } from './ConsumerExecutionContext.js'; + +export type RetryPredicate = ( + error: unknown, + attempt: number, + context: ConsumerExecutionContext, +) => boolean | Promise; diff --git a/src/contracts/kernel/index.ts b/src/contracts/kernel/index.ts index d22afe5..93f32c3 100644 --- a/src/contracts/kernel/index.ts +++ b/src/contracts/kernel/index.ts @@ -1,6 +1,11 @@ export * from './ConsumerMiddleware.js'; +export * from './ConsumerExecutionContext.js'; +export * from './ConsumerNext.js'; export * from './HandlerContext.js'; +export * from './IdempotencyStore.js'; export * from './KernelLogger.js'; export * from './KernelMiddleware.js'; +export * from './RetryDelayResolver.js'; +export * from './RetryPredicate.js'; export * from './ServiceResolver.js'; export * from './ShutdownHook.js'; diff --git a/src/contracts/pubsub/MessageBus.ts b/src/contracts/pubsub/MessageBus.ts new file mode 100644 index 0000000..dde86e9 --- /dev/null +++ b/src/contracts/pubsub/MessageBus.ts @@ -0,0 +1,16 @@ +import type { Message } from './Message.js'; +import type { MessageHandler } from './MessageHandler.js'; +import type { PublisherHook } from './PublisherHook.js'; +import type { Subscription } from './Subscription.js'; + +export interface MessageBus { + publish( + topic: string, + message: TMessage, + ): Promise; + registerPublisherHooks(...hooks: PublisherHook[]): void; + subscribe( + topic: string, + consumer: MessageHandler, + ): Promise; +} diff --git a/src/contracts/pubsub/PublishContext.ts b/src/contracts/pubsub/PublishContext.ts new file mode 100644 index 0000000..6fcd633 --- /dev/null +++ b/src/contracts/pubsub/PublishContext.ts @@ -0,0 +1,12 @@ +import type { DomainEvent } from '../../domain/DomainEvent.js'; +import type { DomainEvent as ContractDomainEvent } from './DomainEvent.js'; +import type { Message } from './Message.js'; + +export interface PublishContext< + TMessage extends Message | ContractDomainEvent = Message, +> { + readonly domainEvent?: DomainEvent; + readonly message: TMessage; + readonly metadata: Readonly>; + readonly topic: string; +} diff --git a/src/contracts/pubsub/PublisherHook.ts b/src/contracts/pubsub/PublisherHook.ts new file mode 100644 index 0000000..b82a0f3 --- /dev/null +++ b/src/contracts/pubsub/PublisherHook.ts @@ -0,0 +1,10 @@ +import type { PublishContext } from './PublishContext.js'; + +export interface PublisherHook { + afterPublish?(context: PublishContext): Promise | void; + beforePublish?(context: PublishContext): Promise | void; + onPublishError?( + error: unknown, + context: PublishContext, + ): Promise | void; +} diff --git a/src/contracts/pubsub/PublisherHookErrorPolicy.ts b/src/contracts/pubsub/PublisherHookErrorPolicy.ts new file mode 100644 index 0000000..9733669 --- /dev/null +++ b/src/contracts/pubsub/PublisherHookErrorPolicy.ts @@ -0,0 +1,9 @@ +import type { PublishContext } from './PublishContext.js'; + +export interface PublisherHookErrorPolicy { + handleAfterPublishError( + error: unknown, + context: PublishContext, + ): Promise | void; + shouldFailAfterPublish(error: unknown, context: PublishContext): boolean; +} diff --git a/src/contracts/pubsub/index.ts b/src/contracts/pubsub/index.ts index 0a1a397..54ecf99 100644 --- a/src/contracts/pubsub/index.ts +++ b/src/contracts/pubsub/index.ts @@ -3,7 +3,11 @@ export * from './DomainEvent.js'; export * from './EventBus.js'; export * from './EventRegistration.js'; export * from './Message.js'; +export * from './MessageBus.js'; export * from './MessageHandler.js'; export * from './MessageMetadata.js'; export * from './PubSub.js'; +export * from './PublishContext.js'; +export * from './PublisherHook.js'; +export * from './PublisherHookErrorPolicy.js'; export * from './Subscription.js'; diff --git a/src/domain/DomainEventConsumer.ts b/src/domain/DomainEventConsumer.ts index 847fa09..859b751 100644 --- a/src/domain/DomainEventConsumer.ts +++ b/src/domain/DomainEventConsumer.ts @@ -1,4 +1,5 @@ import type { DomainEvent } from './DomainEvent.js'; +import type { DomainEventConsumerContext } from './DomainEventConsumerContext.js'; export abstract class DomainEventConsumer { public abstract consume( @@ -6,6 +7,9 @@ export abstract class DomainEventConsumer { bindingKey: string, domainEvent: typeof DomainEvent, exchange: string, - handler: (event: DomainEvent) => Promise, + handler: ( + event: DomainEvent, + context?: DomainEventConsumerContext, + ) => Promise, ): Promise; } diff --git a/src/domain/DomainEventConsumerContext.ts b/src/domain/DomainEventConsumerContext.ts new file mode 100644 index 0000000..af86e89 --- /dev/null +++ b/src/domain/DomainEventConsumerContext.ts @@ -0,0 +1,3 @@ +export interface DomainEventConsumerContext { + readonly metadata: Readonly>; +} diff --git a/src/domain/DomainMessageBus.ts b/src/domain/DomainMessageBus.ts new file mode 100644 index 0000000..a6d8669 --- /dev/null +++ b/src/domain/DomainMessageBus.ts @@ -0,0 +1,8 @@ +import type { MessageBus } from '../contracts/pubsub/MessageBus.js'; +import type { DomainEventConsumer } from './DomainEventConsumer.js'; +import type { DomainEventPublisher } from './DomainEventPublisher.js'; + +export interface DomainMessageBus + extends DomainEventConsumer, DomainEventPublisher { + registerPublisherHooks: MessageBus['registerPublisherHooks']; +} diff --git a/src/domain/index.ts b/src/domain/index.ts index 09e894b..da10747 100644 --- a/src/domain/index.ts +++ b/src/domain/index.ts @@ -3,7 +3,9 @@ export * from './BaseError.js'; export * from './Constructor.js'; export * from './DomainEvent.js'; export * from './DomainEventConsumer.js'; +export * from './DomainEventConsumerContext.js'; export * from './DomainEventPublisher.js'; +export * from './DomainMessageBus.js'; export * from './Event.js'; export * from './EventAttributes.js'; export * from './EventConstructor.js'; diff --git a/tests/adapters/pubsub/Consumer.test.mjs b/tests/adapters/pubsub/Consumer.test.mjs index f40a8b7..fc28912 100644 --- a/tests/adapters/pubsub/Consumer.test.mjs +++ b/tests/adapters/pubsub/Consumer.test.mjs @@ -1,7 +1,13 @@ import assert from 'node:assert/strict'; import test from 'node:test'; -import { Consumer } from '../../../dist/adapters/pubsub/index.js'; +import { + Consumer, + CorrelationConsumerMiddleware, + IdempotencyConsumerMiddleware, + InMemoryIdempotencyStore, + RetryConsumerMiddleware, +} from '../../../dist/adapters/pubsub/index.js'; import { Kernel } from '../../../dist/index.js'; import { TestDomainEvent } from '../../helpers/TestDomainEvent.mjs'; @@ -35,19 +41,32 @@ class TestConsumer extends Consumer { test('initializes the domain event consumer with metadata and middleware chain', async () => { const calls = []; const event = new TestDomainEvent('aggregate-id'); + const rawMessage = { id: 'raw-message' }; const kernel = new Kernel(); const domainEventConsumer = { consume: async (queueName, eventName, EventClass, exchange, handler) => { calls.push([queueName, eventName, EventClass, exchange]); - await handler(event); + await handler(event, { + metadata: { + rawMessage, + retries: 1, + }, + }); }, }; const consumer = new TestConsumer(domainEventConsumer, calls); Kernel.consumerMiddleware.length = 0; kernel.registerConsumerMiddleware({ - async handle(receivedEvent, next) { + async handle(receivedEvent, next, context) { calls.push(['middleware:before', receivedEvent]); + calls.push([ + 'context', + context.eventId, + context.queueName, + context.metadata.retries, + context.rawMessage, + ]); await next(); calls.push(['middleware:after', receivedEvent]); }, @@ -58,6 +77,7 @@ test('initializes the domain event consumer with metadata and middleware chain', assert.deepEqual(calls, [ ['test-queue', 'test.domain-event', TestDomainEvent, 'test-exchange'], ['middleware:before', event], + ['context', event.eventId, 'test-queue', 1, rawMessage], ['handler', event], ['middleware:after', event], ]); @@ -81,3 +101,138 @@ test('resolves legacy services through the active kernel container', () => { assert.equal(consumer.get(Service), service); }); + +test('provides correlation, idempotency and retry middleware', async () => { + const calls = []; + const event = new TestDomainEvent('aggregate-id'); + const kernel = new Kernel(); + const store = new InMemoryIdempotencyStore(); + let attempts = 0; + const domainEventConsumer = { + consume: async (queueName, eventName, EventClass, exchange, handler) => { + void queueName; + void eventName; + void EventClass; + void exchange; + await handler(event); + await handler(event); + }, + }; + const consumer = new TestConsumer(domainEventConsumer, calls); + + Kernel.consumerMiddleware.length = 0; + kernel.registerConsumerMiddleware( + new CorrelationConsumerMiddleware({ + causationId: () => 'causation-id', + correlationId: () => 'correlation-id', + }), + new IdempotencyConsumerMiddleware({ store }), + new RetryConsumerMiddleware({ + maxAttempts: 2, + onRetry: (error, attempt, context) => { + calls.push(['retry', String(error), attempt, context.eventName]); + }, + }), + { + async handle(receivedEvent, next) { + attempts++; + + if (attempts === 1) { + throw new Error('transient'); + } + + await next(); + calls.push([ + 'ids', + receivedEvent.getCorrelationId(), + receivedEvent.getCausationId(), + ]); + }, + }, + ); + + await consumer.init(); + + assert.deepEqual(calls, [ + ['retry', 'Error: transient', 1, 'test.domain-event'], + ['handler', event], + ['ids', 'correlation-id', 'causation-id'], + ]); + assert.equal(attempts, 2); +}); + +test('supports middleware defaults and retry predicates', async () => { + const event = new TestDomainEvent('aggregate-id'); + const context = { + causationId: 'context-causation-id', + correlationId: 'context-correlation-id', + eventId: 'event-id', + eventName: 'test.domain-event', + exchange: 'exchange', + kernel: new Kernel(), + metadata: {}, + queueName: 'queue', + }; + const calls = []; + const store = new InMemoryIdempotencyStore(); + + await new CorrelationConsumerMiddleware().handle( + event, + async () => calls.push(['correlation']), + context, + ); + await new IdempotencyConsumerMiddleware({ + key: () => 'custom-key', + store, + }).handle(event, async () => calls.push(['idempotency']), context); + + await assert.rejects( + () => + new RetryConsumerMiddleware({ + delay: () => 1, + maxAttempts: 2, + shouldRetry: () => false, + }).handle( + event, + async () => { + throw new Error('permanent'); + }, + context, + ), + /permanent/, + ); + + assert.equal(event.getCorrelationId(), 'context-correlation-id'); + assert.equal(event.getCausationId(), 'context-causation-id'); + assert.equal(await store.has('custom-key'), true); + assert.deepEqual(calls, [['correlation'], ['idempotency']]); +}); + +test('releases claimed idempotency keys when handlers fail', async () => { + const event = new TestDomainEvent('aggregate-id'); + const context = { + eventId: 'event-id', + eventName: 'test.domain-event', + exchange: 'exchange', + kernel: new Kernel(), + metadata: {}, + queueName: 'queue', + }; + const store = new InMemoryIdempotencyStore(); + const middleware = new IdempotencyConsumerMiddleware({ store }); + + await assert.rejects( + () => + middleware.handle( + event, + async () => { + throw new Error('failed'); + }, + context, + ), + /failed/, + ); + await middleware.handle(event, async () => {}, context); + + assert.equal(await store.has('event-id'), true); +}); diff --git a/tests/adapters/pubsub/amqp/AmqpMessageBusAdapter.test.mjs b/tests/adapters/pubsub/amqp/AmqpMessageBusAdapter.test.mjs index da79496..fda7663 100644 --- a/tests/adapters/pubsub/amqp/AmqpMessageBusAdapter.test.mjs +++ b/tests/adapters/pubsub/amqp/AmqpMessageBusAdapter.test.mjs @@ -56,7 +56,10 @@ class FakeChannel { checkQueue(queueName) { this.calls.push(['checkQueue', queueName]); - return { consumerCount: this.consumerCount, messageCount: this.messageCount }; + return { + consumerCount: this.consumerCount, + messageCount: this.messageCount, + }; } close() { @@ -134,14 +137,20 @@ const withAmqpConnect = async (channel, run) => { test('publishes domain events and closes channel resources', async () => { const channel = new FakeChannel(); - const { default: AmqpMessageBusAdapter } = await import( - '../../../../dist/adapters/pubsub/amqp/index.js' - ); + const hookCalls = []; + const { default: AmqpMessageBusAdapter } = + await import('../../../../dist/adapters/pubsub/amqp/index.js'); const adapter = new AmqpMessageBusAdapter({ dsn: 'amqp://localhost', exchange: 'domain', serviceName: 'service', }); + adapter.registerPublisherHooks({ + afterPublish: (context) => + hookCalls.push(['after', context.topic, context.domainEvent]), + beforePublish: (context) => + hookCalls.push(['before', context.topic, context.domainEvent]), + }); const event = new TestDomainEvent( 'aggregate-id', { name: 'Ada' }, @@ -156,21 +165,33 @@ test('publishes domain events and closes channel resources', async () => { assert.equal(connection.calls.length, 1); }); - assert.equal(channel.calls.some(([name]) => name === 'publish'), true); - assert.equal(channel.calls.some(([name]) => name === 'channel:close'), true); + assert.equal( + channel.calls.some(([name]) => name === 'publish'), + true, + ); + assert.equal( + channel.calls.some(([name]) => name === 'channel:close'), + true, + ); + assert.deepEqual(hookCalls, [ + ['before', 'test.domain-event', event], + ['after', 'test.domain-event', event], + ]); }); test('consumes AMQP messages and acknowledges handled events', async () => { const channel = new FakeChannel(); const handled = []; - const { default: AmqpMessageBusAdapter } = await import( - '../../../../dist/adapters/pubsub/amqp/index.js' - ); + const { default: AmqpMessageBusAdapter } = + await import('../../../../dist/adapters/pubsub/amqp/index.js'); const adapter = new AmqpMessageBusAdapter({ dsn: 'amqp://localhost', exchange: 'domain', }); - const message = createConsumeMessage(createMessage()); + const message = createConsumeMessage(createMessage(), { + retries: 2, + traceId: 'trace-id', + }); await withAmqpConnect(channel, async () => { await adapter.consume( @@ -178,23 +199,28 @@ test('consumes AMQP messages and acknowledges handled events', async () => { 'test.domain-event', TestDomainEvent, 'domain', - async (event) => handled.push(event), + async (event, context) => handled.push([event, context]), ); await channel.consumers[0](null); await channel.consumers[0](message); }); - assert.equal(handled[0] instanceof TestDomainEvent, true); + assert.equal(handled[0][0] instanceof TestDomainEvent, true); + assert.deepEqual(handled[0][1].metadata.headers, { + retries: 2, + traceId: 'trace-id', + }); + assert.equal(handled[0][1].metadata.rawMessage, message); + assert.equal(handled[0][1].metadata.retries, 2); assert.deepEqual(channel.calls.at(-1), ['ack', message]); }); test('retries failed messages and registers delayed consumers once', async () => { const channel = new FakeChannel(); const logs = []; - const { default: AmqpMessageBusAdapter } = await import( - '../../../../dist/adapters/pubsub/amqp/index.js' - ); + const { default: AmqpMessageBusAdapter } = + await import('../../../../dist/adapters/pubsub/amqp/index.js'); const adapter = new AmqpMessageBusAdapter({ exchange: 'domain', logger: { @@ -230,15 +256,17 @@ test('retries failed messages and registers delayed consumers once', async () => assert.equal(channel.calls.filter(([name]) => name === 'consume').length, 1); assert.equal(channel.calls.filter(([name]) => name === 'publish').length, 2); - assert.equal(logs.some(([, messageText]) => messageText === 'Retry # 1'), true); + assert.equal( + logs.some(([, messageText]) => messageText === 'Retry # 1'), + true, + ); }); test('sends exhausted messages to DLX and logs publish failures', async () => { const channel = new FakeChannel(); const logs = []; - const { default: AmqpMessageBusAdapter } = await import( - '../../../../dist/adapters/pubsub/amqp/index.js' - ); + const { default: AmqpMessageBusAdapter } = + await import('../../../../dist/adapters/pubsub/amqp/index.js'); const adapter = new AmqpMessageBusAdapter({ exchange: 'domain', logger: { @@ -298,15 +326,20 @@ test('consumes DLX messages with success, nack and no-message paths', async () = ); }); - assert.equal(channel.calls.some(([name]) => name === 'ack'), true); - assert.equal(channel.calls.some(([name]) => name === 'nack'), true); + assert.equal( + channel.calls.some(([name]) => name === 'ack'), + true, + ); + assert.equal( + channel.calls.some(([name]) => name === 'nack'), + true, + ); }); test('checks queue bindings for registered consumers', async () => { const channel = new FakeChannel(); - const { default: AmqpMessageBusAdapter } = await import( - '../../../../dist/adapters/pubsub/amqp/index.js' - ); + const { default: AmqpMessageBusAdapter } = + await import('../../../../dist/adapters/pubsub/amqp/index.js'); const adapter = new AmqpMessageBusAdapter({ dsn: 'amqp://localhost' }); const kernel = new Kernel(); @@ -329,9 +362,8 @@ test('uses environment retry delay and logs non-Error DLX retry failures', async const previousRetryDelay = process.env.TRANSPORT_RETRY_DELAY; const channel = new FakeChannel(); const logs = []; - const { default: AmqpMessageBusAdapter } = await import( - '../../../../dist/adapters/pubsub/amqp/index.js' - ); + const { default: AmqpMessageBusAdapter } = + await import('../../../../dist/adapters/pubsub/amqp/index.js'); const adapter = new AmqpMessageBusAdapter({ exchange: 'domain', logger: { @@ -370,9 +402,8 @@ test('uses environment retry delay and logs non-Error DLX retry failures', async test('handles missing retry headers and cancels delayed consumers on invalid retry payloads', async () => { const channel = new FakeChannel(); const logs = []; - const { default: AmqpMessageBusAdapter } = await import( - '../../../../dist/adapters/pubsub/amqp/index.js' - ); + const { default: AmqpMessageBusAdapter } = + await import('../../../../dist/adapters/pubsub/amqp/index.js'); const adapter = new AmqpMessageBusAdapter({ exchange: 'domain', logger: { @@ -400,16 +431,18 @@ test('handles missing retry headers and cancels delayed consumers on invalid ret ); await adapter.retry(null, {}, context); - assert.equal(channel.calls.some(([name]) => name === 'cancel'), true); + assert.equal( + channel.calls.some(([name]) => name === 'cancel'), + true, + ); assert.equal(logs.includes('Invalid domain event: null'), true); }); test('handles AMQP messages that fail during handler execution', async () => { const channel = new FakeChannel(); const logs = []; - const { default: AmqpMessageBusAdapter } = await import( - '../../../../dist/adapters/pubsub/amqp/index.js' - ); + const { default: AmqpMessageBusAdapter } = + await import('../../../../dist/adapters/pubsub/amqp/index.js'); const adapter = new AmqpMessageBusAdapter({ exchange: 'domain', logger: { @@ -438,9 +471,8 @@ test('handles AMQP messages that fail during handler execution', async () => { test('supports AMQP messages without occurred_on and retry publish errors', async () => { const channel = new FakeChannel(); const logs = []; - const { default: AmqpMessageBusAdapter } = await import( - '../../../../dist/adapters/pubsub/amqp/index.js' - ); + const { default: AmqpMessageBusAdapter } = + await import('../../../../dist/adapters/pubsub/amqp/index.js'); const adapter = new AmqpMessageBusAdapter({ exchange: 'domain', logger: { @@ -471,9 +503,8 @@ test('supports AMQP messages without occurred_on and retry publish errors', asyn test('uses numeric retry delay from environment', async () => { const previousRetryDelay = process.env.TRANSPORT_RETRY_DELAY; - const { default: AmqpMessageBusAdapter } = await import( - '../../../../dist/adapters/pubsub/amqp/index.js' - ); + const { default: AmqpMessageBusAdapter } = + await import('../../../../dist/adapters/pubsub/amqp/index.js'); process.env.TRANSPORT_RETRY_DELAY = '25'; @@ -493,9 +524,8 @@ test('uses numeric retry delay from environment', async () => { test('reads AMQP DSN and max retries from environment defaults', async () => { const previousDsn = process.env.TRANSPORT_DSN; const previousRetries = process.env.TRANSPORT_MAX_RETRIES; - const { default: AmqpMessageBusAdapter } = await import( - '../../../../dist/adapters/pubsub/amqp/index.js' - ); + const { default: AmqpMessageBusAdapter } = + await import('../../../../dist/adapters/pubsub/amqp/index.js'); delete process.env.TRANSPORT_MAX_RETRIES; process.env.TRANSPORT_DSN = 'amqp://environment'; @@ -531,9 +561,8 @@ test('reads AMQP DSN and max retries from environment defaults', async () => { test('handles delayed consumer messages and removes delayed consumer state', async () => { const channel = new FakeChannel(); const handled = []; - const { default: AmqpMessageBusAdapter } = await import( - '../../../../dist/adapters/pubsub/amqp/index.js' - ); + const { default: AmqpMessageBusAdapter } = + await import('../../../../dist/adapters/pubsub/amqp/index.js'); const adapter = new AmqpMessageBusAdapter({ exchange: 'domain', retryDelayInMilliseconds: 0, @@ -563,9 +592,8 @@ test('logs non-Error retry publish failures', async () => { const channel = new FakeChannel(); const logs = []; - const { default: AmqpMessageBusAdapter } = await import( - '../../../../dist/adapters/pubsub/amqp/index.js' - ); + const { default: AmqpMessageBusAdapter } = + await import('../../../../dist/adapters/pubsub/amqp/index.js'); const adapter = new AmqpMessageBusAdapter({ exchange: 'domain', logger: { @@ -595,9 +623,8 @@ test('logs non-Error retry publish failures', async () => { test('logs string publish failures when sending to DLX', async () => { const channel = new FakeChannel(); const logs = []; - const { default: AmqpMessageBusAdapter } = await import( - '../../../../dist/adapters/pubsub/amqp/index.js' - ); + const { default: AmqpMessageBusAdapter } = + await import('../../../../dist/adapters/pubsub/amqp/index.js'); const adapter = new AmqpMessageBusAdapter({ exchange: 'domain', logger: { @@ -622,22 +649,23 @@ test('logs string publish failures when sending to DLX', async () => { }); test('throws when a channel cannot be created', async () => { - const { default: AmqpMessageBusAdapter } = await import( - '../../../../dist/adapters/pubsub/amqp/index.js' - ); + const { default: AmqpMessageBusAdapter } = + await import('../../../../dist/adapters/pubsub/amqp/index.js'); const adapter = new AmqpMessageBusAdapter({ dsn: 'amqp://localhost' }); adapter.connect = async () => {}; - await assert.rejects(() => adapter.channel(), /AMQP channel could not be created/); + await assert.rejects( + () => adapter.channel(), + /AMQP channel could not be created/, + ); }); test('reconnects consumers when AMQP channel emits close or error', async () => { const channel = new FakeChannel(); const calls = []; - const { default: AmqpMessageBusAdapter } = await import( - '../../../../dist/adapters/pubsub/amqp/index.js' - ); + const { default: AmqpMessageBusAdapter } = + await import('../../../../dist/adapters/pubsub/amqp/index.js'); const adapter = new AmqpMessageBusAdapter({ dsn: 'amqp://localhost', logger: { @@ -650,7 +678,9 @@ test('reconnects consumers when AMQP channel emits close or error', async () => const kernel = new Kernel(); kernel.removeConsumers(); - kernel.registerConsumerInstances({ init: async () => calls.push(['consumer:init']) }); + kernel.registerConsumerInstances({ + init: async () => calls.push(['consumer:init']), + }); await withAmqpConnect(channel, async () => { await adapter.channel(); diff --git a/tests/adapters/pubsub/in-memory/InMemoryEventBus.test.mjs b/tests/adapters/pubsub/in-memory/InMemoryEventBus.test.mjs index 868c8a5..ad859c7 100644 --- a/tests/adapters/pubsub/in-memory/InMemoryEventBus.test.mjs +++ b/tests/adapters/pubsub/in-memory/InMemoryEventBus.test.mjs @@ -18,3 +18,31 @@ test('publishes domain events to subscribed handlers with context', async () => assert.deepEqual(calls, [[event, context]]); }); + +test('runs registered publisher hooks around domain events', async () => { + const context = { di: {}, publish: async () => {} }; + const event = { + metadata: { correlationId: 'correlation-id' }, + name: 'user.created', + }; + const calls = []; + const eventBus = new InMemoryEventBus(context); + + eventBus.registerPublisherHooks({ + afterPublish: (publishContext) => + calls.push(['after', publishContext.topic]), + beforePublish: (publishContext) => + calls.push(['before', publishContext.metadata.correlationId]), + }); + eventBus.subscribe('user.created', async (receivedEvent) => { + calls.push(['handler', receivedEvent.name]); + }); + + await eventBus.publish(event); + + assert.deepEqual(calls, [ + ['before', 'correlation-id'], + ['handler', 'user.created'], + ['after', 'user.created'], + ]); +}); diff --git a/tests/adapters/pubsub/in-memory/InMemoryPubSub.test.mjs b/tests/adapters/pubsub/in-memory/InMemoryPubSub.test.mjs index 5039394..9705a35 100644 --- a/tests/adapters/pubsub/in-memory/InMemoryPubSub.test.mjs +++ b/tests/adapters/pubsub/in-memory/InMemoryPubSub.test.mjs @@ -9,9 +9,12 @@ test('publishes messages to subscribers and supports unsubscribe', async () => { const calls = []; const pubSub = new InMemoryPubSub(context); - const subscription = await pubSub.subscribe('topic', async (receivedMessage, receivedContext) => { - calls.push([receivedMessage, receivedContext]); - }); + const subscription = await pubSub.subscribe( + 'topic', + async (receivedMessage, receivedContext) => { + calls.push([receivedMessage, receivedContext]); + }, + ); await pubSub.publish('topic', message); await pubSub.publish('other-topic', message); @@ -20,3 +23,106 @@ test('publishes messages to subscribers and supports unsubscribe', async () => { assert.deepEqual(calls, [[message, context]]); }); + +test('runs publisher hooks around published messages', async () => { + const context = { di: {}, publish: async () => {} }; + const message = { + metadata: { correlationId: 'correlation-id' }, + name: 'message', + }; + const calls = []; + const pubSub = new InMemoryPubSub(context, [ + { + afterPublish: (publishContext) => + calls.push(['after', publishContext.topic]), + beforePublish: (publishContext) => + calls.push(['before', publishContext.metadata.correlationId]), + }, + ]); + + pubSub.registerPublisherHooks({ + afterPublish: (publishContext) => + calls.push(['registered', publishContext.message.name]), + }); + await pubSub.subscribe('topic', async (receivedMessage) => { + calls.push(['consumer', receivedMessage.name]); + }); + + await pubSub.publish('topic', message); + + assert.deepEqual(calls, [ + ['before', 'correlation-id'], + ['consumer', 'message'], + ['after', 'topic'], + ['registered', 'message'], + ]); +}); + +test('runs publisher error hooks before rethrowing', async () => { + const context = { di: {}, publish: async () => {} }; + const error = new Error('publish failed'); + const calls = []; + const pubSub = new InMemoryPubSub(context, [ + { + onPublishError: (publishError, publishContext) => + calls.push([publishError, publishContext.topic]), + }, + ]); + + await pubSub.subscribe('topic', async () => { + throw error; + }); + + await assert.rejects( + () => pubSub.publish('topic', { name: 'message' }), + error, + ); + assert.deepEqual(calls, [[error, 'topic']]); +}); + +test('does not fail publish when afterPublish hooks fail by default', async () => { + const context = { di: {}, publish: async () => {} }; + const calls = []; + const pubSub = new InMemoryPubSub(context, [ + { + afterPublish: () => { + throw new Error('websocket failed'); + }, + }, + ]); + + await pubSub.subscribe('topic', async (receivedMessage) => { + calls.push(['consumer', receivedMessage.name]); + }); + + await pubSub.publish('topic', { name: 'message' }); + + assert.deepEqual(calls, [['consumer', 'message']]); +}); + +test('can fail publish when afterPublish policy asks for it', async () => { + const context = { di: {}, publish: async () => {} }; + const afterPublishError = new Error('replica failed'); + const policyCalls = []; + const pubSub = new InMemoryPubSub( + context, + [ + { + afterPublish: () => { + throw afterPublishError; + }, + }, + ], + { + handleAfterPublishError: (error, publishContext) => + policyCalls.push([error, publishContext.topic]), + shouldFailAfterPublish: () => true, + }, + ); + + await assert.rejects( + () => pubSub.publish('topic', { name: 'message' }), + afterPublishError, + ); + assert.deepEqual(policyCalls, [[afterPublishError, 'topic']]); +}); diff --git a/tests/adapters/ui/express/ExpressKernelServer.test.mjs b/tests/adapters/ui/express/ExpressKernelServer.test.mjs index 39a0a8a..47b65c1 100644 --- a/tests/adapters/ui/express/ExpressKernelServer.test.mjs +++ b/tests/adapters/ui/express/ExpressKernelServer.test.mjs @@ -106,3 +106,66 @@ test('registers default error handlers and runs without optional middleware', as await server.run(); await server.close(); }); + +test('runs configurable controller, swagger and static hooks', async () => { + class ExternalController {} + + const calls = []; + const middleware = (name) => (request, response, next) => { + void request; + void response; + calls.push(name); + next(); + }; + const server = new ExpressKernelServer({ + afterControllersHooks: [(app) => calls.push(['after', Boolean(app)])], + beforeControllersHooks: [(app) => calls.push(['before', Boolean(app)])], + controllers: [ExternalController], + hooks: [ + { + handle: (app) => calls.push(['phase:before', Boolean(app)]), + phase: 'beforeControllers', + }, + { + handle: (app) => calls.push(['phase:after', Boolean(app)]), + phase: 'afterControllers', + }, + { + handle: (app) => calls.push(['phase:errors', Boolean(app)]), + phase: 'beforeErrors', + }, + ], + kernel: { getRoutes: () => [] }, + middlewares: [middleware('base')], + port: 0, + postControllerMiddlewares: [middleware('post')], + preControllerMiddlewares: [middleware('pre')], + staticHooks: [(app) => calls.push(['static', Boolean(app)])], + swaggerHooks: [(app) => calls.push(['swagger', Boolean(app)])], + }); + + await server.run(); + await server.close(); + + assert.deepEqual(calls, [ + ['before', true], + ['phase:before', true], + ['after', true], + ['phase:after', true], + ['swagger', true], + ['static', true], + ['phase:errors', true], + ]); +}); + +test('rejects duplicate run calls while server is running', async () => { + const server = new ExpressKernelServer({ + kernel: { getRoutes: () => [] }, + port: 0, + }); + + await server.run(); + + await assert.rejects(() => server.run(), /HTTP server is already running/); + await server.close(); +}); diff --git a/tests/typescript-module-resolution.test.mjs b/tests/typescript-module-resolution.test.mjs new file mode 100644 index 0000000..4891c48 --- /dev/null +++ b/tests/typescript-module-resolution.test.mjs @@ -0,0 +1,97 @@ +import assert from 'node:assert/strict'; +import { existsSync } from 'node:fs'; +import { mkdir, symlink, writeFile } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import path from 'node:path'; +import { spawn } from 'node:child_process'; +import test from 'node:test'; + +test('exports types for TypeScript moduleResolution node consumers', async () => { + if (!existsSync(path.resolve('dist/contracts/kernel/index.d.ts'))) { + return; + } + + const temporaryDirectory = path.join( + await import('node:fs/promises').then(({ mkdtemp }) => + mkdtemp(path.join(tmpdir(), 'ddd-kernel-types-')), + ), + ); + const packageDirectory = path.resolve('.'); + const packageScopeDirectory = path.join( + temporaryDirectory, + 'node_modules', + '@haskou', + ); + + await mkdir(packageScopeDirectory, { recursive: true }); + await symlink( + packageDirectory, + path.join(packageScopeDirectory, 'ddd-kernel'), + ); + await writeFile( + path.join(temporaryDirectory, 'package.json'), + JSON.stringify({ type: 'module' }), + ); + await writeFile( + path.join(temporaryDirectory, 'index.ts'), + ` + import type { ConsumerMiddleware } from '@haskou/ddd-kernel/contracts/kernel'; + import type { MessageBus, PublisherHook } from '@haskou/ddd-kernel/contracts/pubsub'; + import type { DomainMessageBus } from '@haskou/ddd-kernel/domain'; + import type { SchedulerErrorPolicy } from '@haskou/ddd-kernel/scheduler'; + import { ExpressKernelServer } from '@haskou/ddd-kernel/adapters/ui/express'; + + const middleware: ConsumerMiddleware | undefined = undefined; + const messageBus: MessageBus | undefined = undefined; + const domainMessageBus: DomainMessageBus | undefined = undefined; + const hook: PublisherHook | undefined = undefined; + const policy: SchedulerErrorPolicy | undefined = undefined; + + void middleware; + void messageBus; + void domainMessageBus; + void hook; + void policy; + void ExpressKernelServer; + `, + ); + await writeFile( + path.join(temporaryDirectory, 'tsconfig.json'), + JSON.stringify({ + compilerOptions: { + ignoreDeprecations: '6.0', + module: 'ESNext', + moduleResolution: 'node', + noEmit: true, + skipLibCheck: true, + strict: true, + target: 'ES2022', + }, + include: ['index.ts'], + }), + ); + + const result = await new Promise((resolve) => { + const child = spawn( + process.execPath, + [ + path.resolve('node_modules/typescript/bin/tsc'), + '-p', + path.join(temporaryDirectory, 'tsconfig.json'), + ], + { cwd: temporaryDirectory }, + ); + let stderr = ''; + let stdout = ''; + + child.stderr.on('data', (chunk) => { + stderr += chunk.toString(); + }); + child.stdout.on('data', (chunk) => { + stdout += chunk.toString(); + }); + child.on('close', (code) => resolve({ code, stderr, stdout })); + }); + + assert.equal(result.code, 0, `${result.stdout}\n${result.stderr}`); +});