From 1573c189230bb9e098d20e4b34ec974b52d363ca Mon Sep 17 00:00:00 2001 From: Hasko Date: Thu, 25 Jun 2026 16:45:42 +0200 Subject: [PATCH 1/5] =?UTF-8?q?feat(pubsub):=20=E2=9C=A8=20Add=20consumer?= =?UTF-8?q?=20middleware=20pipeline?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/reference/consumer.md | 36 +++++- src/Kernel.ts | 4 + src/adapters/pubsub/Consumer.ts | 35 +++--- .../pubsub/ConsumerMiddlewarePipeline.ts | 38 ++++++ .../pubsub/CorrelationConsumerMiddleware.ts | 36 ++++++ .../CorrelationConsumerMiddlewareOptions.ts | 13 ++ .../pubsub/IdempotencyConsumerMiddleware.ts | 28 +++++ .../IdempotencyConsumerMiddlewareOptions.ts | 13 ++ .../pubsub/InMemoryIdempotencyStore.ts | 15 +++ .../pubsub/RetryConsumerMiddleware.ts | 53 ++++++++ .../pubsub/RetryConsumerMiddlewareOptions.ts | 16 +++ src/adapters/pubsub/index.ts | 9 ++ .../kernel/ConsumerExecutionContext.ts | 12 ++ src/contracts/kernel/ConsumerMiddleware.ts | 8 +- src/contracts/kernel/ConsumerNext.ts | 1 + src/contracts/kernel/IdempotencyStore.ts | 4 + src/contracts/kernel/RetryDelayResolver.ts | 7 ++ src/contracts/kernel/RetryPredicate.ts | 7 ++ src/contracts/kernel/index.ts | 5 + tests/adapters/pubsub/Consumer.test.mjs | 118 +++++++++++++++++- 20 files changed, 435 insertions(+), 23 deletions(-) create mode 100644 src/adapters/pubsub/ConsumerMiddlewarePipeline.ts create mode 100644 src/adapters/pubsub/CorrelationConsumerMiddleware.ts create mode 100644 src/adapters/pubsub/CorrelationConsumerMiddlewareOptions.ts create mode 100644 src/adapters/pubsub/IdempotencyConsumerMiddleware.ts create mode 100644 src/adapters/pubsub/IdempotencyConsumerMiddlewareOptions.ts create mode 100644 src/adapters/pubsub/InMemoryIdempotencyStore.ts create mode 100644 src/adapters/pubsub/RetryConsumerMiddleware.ts create mode 100644 src/adapters/pubsub/RetryConsumerMiddlewareOptions.ts create mode 100644 src/contracts/kernel/ConsumerExecutionContext.ts create mode 100644 src/contracts/kernel/ConsumerNext.ts create mode 100644 src/contracts/kernel/IdempotencyStore.ts create mode 100644 src/contracts/kernel/RetryDelayResolver.ts create mode 100644 src/contracts/kernel/RetryPredicate.ts diff --git a/docs/reference/consumer.md b/docs/reference/consumer.md index edeca78..fcf94fa 100644 --- a/docs/reference/consumer.md +++ b/docs/reference/consumer.md @@ -25,11 +25,41 @@ 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. + +## 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. The in-memory store is +only useful for tests and single-process applications. 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..2c86cd4 100644 --- a/src/adapters/pubsub/Consumer.ts +++ b/src/adapters/pubsub/Consumer.ts @@ -1,27 +1,28 @@ -import type { ConsumerMiddleware } from '../../contracts/index.js'; import type { DomainEventConsumer } from '../../domain/DomainEventConsumer.js'; import type { DomainEvent } 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, - ): Promise { - const middleware = middlewares[index]; - - if (!middleware) { - await this.handler(event); - - return; - } - - await middleware.handle(event, () => - this.runMiddleware(event, middlewares, index + 1), + private async runMiddleware(event: DomainEvent): Promise { + const pipeline = new ConsumerMiddlewarePipeline(Kernel.consumerMiddleware); + + 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, + }, + () => this.handler(event), ); } @@ -41,7 +42,7 @@ export abstract class Consumer { this.eventName, this.domainEvent, this.exchange, - (event) => this.runMiddleware(event, Kernel.consumerMiddleware, 0), + (event) => this.runMiddleware(event), ); } 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/IdempotencyConsumerMiddleware.ts b/src/adapters/pubsub/IdempotencyConsumerMiddleware.ts new file mode 100644 index 0000000..930734d --- /dev/null +++ b/src/adapters/pubsub/IdempotencyConsumerMiddleware.ts @@ -0,0 +1,28 @@ +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) {} + + public async handle( + event: DomainEvent, + next: ConsumerNext, + context: ConsumerExecutionContext, + ): Promise { + const key = this.options.key?.(event, context) ?? context.eventId; + + if (await this.options.store.has(key)) { + return; + } + + await next(); + await this.options.store.mark(key); + } +} + +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..78db276 --- /dev/null +++ b/src/adapters/pubsub/InMemoryIdempotencyStore.ts @@ -0,0 +1,15 @@ +import type { IdempotencyStore } from '../../contracts/index.js'; + +export class InMemoryIdempotencyStore implements IdempotencyStore { + private readonly handledKeys = new Set(); + + 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/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/index.ts b/src/adapters/pubsub/index.ts index e9c602a..823e299 100644 --- a/src/adapters/pubsub/index.ts +++ b/src/adapters/pubsub/index.ts @@ -1,4 +1,13 @@ export * from './Consumer.js'; +export * from './ConsumerMiddlewarePipeline.js'; +export * from './CorrelationConsumerMiddleware.js'; +export * from './CorrelationConsumerMiddlewareOptions.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/contracts/kernel/ConsumerExecutionContext.ts b/src/contracts/kernel/ConsumerExecutionContext.ts new file mode 100644 index 0000000..ba8a3ff --- /dev/null +++ b/src/contracts/kernel/ConsumerExecutionContext.ts @@ -0,0 +1,12 @@ +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 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..97dda55 --- /dev/null +++ b/src/contracts/kernel/IdempotencyStore.ts @@ -0,0 +1,4 @@ +export interface IdempotencyStore { + 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/tests/adapters/pubsub/Consumer.test.mjs b/tests/adapters/pubsub/Consumer.test.mjs index f40a8b7..03a118b 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'; @@ -46,8 +52,9 @@ test('initializes the domain event consumer with metadata and middleware chain', 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]); await next(); calls.push(['middleware:after', receivedEvent]); }, @@ -58,6 +65,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'], ['handler', event], ['middleware:after', event], ]); @@ -81,3 +89,109 @@ 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']]); +}); From 9646a7936ae6c530f7e63386b90a044f1b9e3f77 Mon Sep 17 00:00:00 2001 From: Hasko Date: Thu, 25 Jun 2026 16:45:52 +0200 Subject: [PATCH 2/5] =?UTF-8?q?feat(pubsub):=20=E2=9C=A8=20Add=20publisher?= =?UTF-8?q?=20hook=20extension=20points?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/guides/adapters.md | 38 +++++ docs/guides/amqp-pubsub.md | 5 + src/adapters/pubsub/PublisherHookPipeline.ts | 38 +++++ .../pubsub/amqp/AmqpMessageBusAdapter.ts | 39 ++++- .../amqp/AmqpMessageBusAdapterOptions.ts | 2 + .../pubsub/in-memory/InMemoryEventBus.ts | 27 +++- .../pubsub/in-memory/InMemoryPubSub.ts | 27 +++- src/contracts/pubsub/MessageBus.ts | 7 + src/contracts/pubsub/PublishContext.ts | 10 ++ src/contracts/pubsub/PublisherHook.ts | 10 ++ src/contracts/pubsub/index.ts | 3 + .../amqp/AmqpMessageBusAdapter.test.mjs | 133 ++++++++++-------- .../in-memory/InMemoryEventBus.test.mjs | 28 ++++ .../pubsub/in-memory/InMemoryPubSub.test.mjs | 65 ++++++++- 14 files changed, 359 insertions(+), 73 deletions(-) create mode 100644 src/adapters/pubsub/PublisherHookPipeline.ts create mode 100644 src/contracts/pubsub/MessageBus.ts create mode 100644 src/contracts/pubsub/PublishContext.ts create mode 100644 src/contracts/pubsub/PublisherHook.ts diff --git a/docs/guides/adapters.md b/docs/guides/adapters.md index dc5e825..0c566ea 100644 --- a/docs/guides/adapters.md +++ b/docs/guides/adapters.md @@ -27,3 +27,41 @@ 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({ + publisherHooks: [ + { + afterPublish: async ({ message }) => { + await websocketPublisher.publish(message); + }, + }, + ], +}); +``` + +Custom adapters should implement the `MessageBus` contract and 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/src/adapters/pubsub/PublisherHookPipeline.ts b/src/adapters/pubsub/PublisherHookPipeline.ts new file mode 100644 index 0000000..4055176 --- /dev/null +++ b/src/adapters/pubsub/PublisherHookPipeline.ts @@ -0,0 +1,38 @@ +import type { PublishContext, PublisherHook } from '../../contracts/index.js'; + +export class PublisherHookPipeline { + private readonly hooks: PublisherHook[] = []; + + constructor(hooks: readonly PublisherHook[] = []) { + this.hooks.push(...hooks); + } + + public register(...hooks: PublisherHook[]): void { + this.hooks.push(...hooks); + } + + public async run( + context: PublishContext, + publish: () => Promise | T, + ): Promise { + for (const hook of this.hooks) { + await hook.beforePublish?.(context); + } + + try { + const result = await publish(); + + for (const hook of this.hooks) { + await hook.afterPublish?.(context); + } + + return result; + } catch (error: unknown) { + for (const hook of this.hooks) { + await hook.onPublishError?.(error, context); + } + + throw error; + } + } +} diff --git a/src/adapters/pubsub/amqp/AmqpMessageBusAdapter.ts b/src/adapters/pubsub/amqp/AmqpMessageBusAdapter.ts index ce90763..c63f9dc 100644 --- a/src/adapters/pubsub/amqp/AmqpMessageBusAdapter.ts +++ b/src/adapters/pubsub/amqp/AmqpMessageBusAdapter.ts @@ -8,6 +8,7 @@ import amqplib, { } from 'amqplib'; import { randomUUID } from 'node:crypto'; +import type { PublisherHook } from '../../../contracts/index.js'; import type { Constructor, DomainEvent, @@ -19,6 +20,7 @@ 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'; @@ -29,10 +31,14 @@ export default class AmqpMessageBusAdapter 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, + ); } private get dsn(): string { @@ -450,15 +456,38 @@ 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( + { + 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..fc2db4a 100644 --- a/src/adapters/pubsub/amqp/AmqpMessageBusAdapterOptions.ts +++ b/src/adapters/pubsub/amqp/AmqpMessageBusAdapterOptions.ts @@ -1,3 +1,4 @@ +import type { PublisherHook } from '../../../contracts/index.js'; import type { Log } from '../../../infrastructure/logs/index.js'; export interface AmqpMessageBusAdapterOptions { @@ -5,6 +6,7 @@ export interface AmqpMessageBusAdapterOptions { readonly exchange?: string; readonly logger?: Log; readonly maxRetries?: number; + readonly publisherHooks?: PublisherHook[]; readonly retryDelayInMilliseconds?: number; readonly serviceName?: string; } diff --git a/src/adapters/pubsub/in-memory/InMemoryEventBus.ts b/src/adapters/pubsub/in-memory/InMemoryEventBus.ts index f55fd09..727c73e 100644 --- a/src/adapters/pubsub/in-memory/InMemoryEventBus.ts +++ b/src/adapters/pubsub/in-memory/InMemoryEventBus.ts @@ -2,15 +2,25 @@ import type { DomainEvent, HandlerContext, MessageHandler, + PublisherHook, } 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[] = [], + ) { + this.publisherHookPipeline = new PublisherHookPipeline(publisherHooks); + } public subscribe( name: TEvent['name'], @@ -27,8 +37,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..2523f3f 100644 --- a/src/adapters/pubsub/in-memory/InMemoryPubSub.ts +++ b/src/adapters/pubsub/in-memory/InMemoryPubSub.ts @@ -2,16 +2,26 @@ import type { HandlerContext, Message, MessageHandler, + PublisherHook, 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[] = [], + ) { + this.publisherHookPipeline = new PublisherHookPipeline(publisherHooks); + } public async publish( topic: string, @@ -19,9 +29,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 +58,8 @@ export class InMemoryPubSub { }, }); } + + public registerPublisherHooks(...hooks: PublisherHook[]): void { + this.publisherHookPipeline.register(...hooks); + } } diff --git a/src/contracts/pubsub/MessageBus.ts b/src/contracts/pubsub/MessageBus.ts new file mode 100644 index 0000000..4afe13b --- /dev/null +++ b/src/contracts/pubsub/MessageBus.ts @@ -0,0 +1,7 @@ +import type { DomainEventConsumer } from '../../domain/DomainEventConsumer.js'; +import type { DomainEventPublisher } from '../../domain/DomainEventPublisher.js'; +import type { PublisherHook } from './PublisherHook.js'; + +export interface MessageBus extends DomainEventConsumer, DomainEventPublisher { + registerPublisherHooks(...hooks: PublisherHook[]): void; +} diff --git a/src/contracts/pubsub/PublishContext.ts b/src/contracts/pubsub/PublishContext.ts new file mode 100644 index 0000000..43ce260 --- /dev/null +++ b/src/contracts/pubsub/PublishContext.ts @@ -0,0 +1,10 @@ +import type { DomainEvent as ContractDomainEvent } from './DomainEvent.js'; +import type { Message } from './Message.js'; + +export interface PublishContext< + TMessage extends Message | ContractDomainEvent = Message, +> { + 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/index.ts b/src/contracts/pubsub/index.ts index 0a1a397..a119358 100644 --- a/src/contracts/pubsub/index.ts +++ b/src/contracts/pubsub/index.ts @@ -3,7 +3,10 @@ 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 './Subscription.js'; diff --git a/tests/adapters/pubsub/amqp/AmqpMessageBusAdapter.test.mjs b/tests/adapters/pubsub/amqp/AmqpMessageBusAdapter.test.mjs index da79496..f498333 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,18 @@ 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]), + beforePublish: (context) => hookCalls.push(['before', context.topic]), + }); const event = new TestDomainEvent( 'aggregate-id', { name: 'Ada' }, @@ -156,16 +163,25 @@ 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'], + ['after', 'test.domain-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', @@ -192,9 +208,8 @@ test('consumes AMQP messages and acknowledges handled events', async () => { 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 +245,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 +315,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 +351,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 +391,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 +420,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 +460,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 +492,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 +513,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 +550,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 +581,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 +612,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 +638,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 +667,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..2e139e6 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,59 @@ 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']]); +}); From 1af5d0c643a5162b38dd2993dbeb3ec449884adc Mon Sep 17 00:00:00 2001 From: Hasko Date: Thu, 25 Jun 2026 16:45:59 +0200 Subject: [PATCH 3/5] =?UTF-8?q?feat(ui):=20=E2=9C=A8=20Extend=20Express=20?= =?UTF-8?q?kernel=20server=20hooks?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/reference/express-kernel-server.md | 44 ++++++++++++++++ src/adapters/ui/express/ExpressAppHook.ts | 3 ++ src/adapters/ui/express/ExpressController.ts | 1 + .../ui/express/ExpressKernelServer.ts | 51 +++++++++++++++---- .../ui/express/ExpressKernelServerOptions.ts | 9 ++++ src/adapters/ui/express/index.ts | 2 + .../ui/express/ExpressKernelServer.test.mjs | 34 +++++++++++++ 7 files changed, 133 insertions(+), 11 deletions(-) create mode 100644 src/adapters/ui/express/ExpressAppHook.ts create mode 100644 src/adapters/ui/express/ExpressController.ts diff --git a/docs/reference/express-kernel-server.md b/docs/reference/express-kernel-server.md index 2d7ee40..c5637d4 100644 --- a/docs/reference/express-kernel-server.md +++ b/docs/reference/express-kernel-server.md @@ -11,3 +11,47 @@ 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, + middlewares: [requestIdMiddleware], + preControllerMiddlewares: [authenticationMiddleware], + postControllerMiddlewares: [notFoundMiddleware], + swaggerHooks: [(app) => setupSwagger(app)], + staticHooks: [(app) => app.use('/public', express.static('public'))], +}); +``` + +Hook order is: + +1. `middlewares` +2. `preControllerMiddlewares` +3. `beforeControllersHooks` +4. `routing-controllers` +5. `postControllerMiddlewares` +6. `afterControllersHooks` +7. `swaggerHooks` +8. `staticHooks` +9. `errorHandlers` 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/ExpressKernelServer.ts b/src/adapters/ui/express/ExpressKernelServer.ts index 590966c..c7bafa8 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,24 @@ 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 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,15 +86,24 @@ export class ExpressKernelServer { }); } - public run(): Promise { - const app = createExpressServer({ - controllers: this.options.kernel.getRoutes(), + public async run(): Promise { + 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); + useExpressServer(app, { + controllers, routePrefix: this.options.routePrefix, - }) as HttpApp; - - for (const middleware of this.options.middlewares ?? []) { - app.use(middleware); - } + }); + this.registerMiddlewares(app, this.options.postControllerMiddlewares); + await this.runHooks(this.options.afterControllersHooks, app); + await this.runHooks(this.options.swaggerHooks, app); + await this.runHooks(this.options.staticHooks, 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..6e1bca6 100644 --- a/src/adapters/ui/express/ExpressKernelServerOptions.ts +++ b/src/adapters/ui/express/ExpressKernelServerOptions.ts @@ -1,11 +1,20 @@ import type { ErrorRequestHandler, RequestHandler } from 'express'; import type { Kernel } from '../../../Kernel.js'; +import type { ExpressAppHook } from './ExpressAppHook.js'; +import type { ExpressController } from './ExpressController.js'; export interface ExpressKernelServerOptions { + readonly afterControllersHooks?: ExpressAppHook[]; + readonly beforeControllersHooks?: ExpressAppHook[]; + readonly controllers?: ExpressController[]; readonly errorHandlers?: ErrorRequestHandler[]; readonly kernel: Kernel; readonly middlewares?: RequestHandler[]; + readonly postControllerMiddlewares?: RequestHandler[]; + readonly preControllerMiddlewares?: RequestHandler[]; readonly port?: number; readonly routePrefix?: string; + readonly staticHooks?: ExpressAppHook[]; + readonly swaggerHooks?: ExpressAppHook[]; } diff --git a/src/adapters/ui/express/index.ts b/src/adapters/ui/express/index.ts index cf5fbf6..15f3913 100644 --- a/src/adapters/ui/express/index.ts +++ b/src/adapters/ui/express/index.ts @@ -1,3 +1,5 @@ +export * from './ExpressAppHook.js'; +export * from './ExpressController.js'; export * from './ExpressKernelServer.js'; export * from './ExpressKernelServerOptions.js'; export * from './HttpApp.js'; diff --git a/tests/adapters/ui/express/ExpressKernelServer.test.mjs b/tests/adapters/ui/express/ExpressKernelServer.test.mjs index 39a0a8a..867898b 100644 --- a/tests/adapters/ui/express/ExpressKernelServer.test.mjs +++ b/tests/adapters/ui/express/ExpressKernelServer.test.mjs @@ -106,3 +106,37 @@ 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], + 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], + ['after', true], + ['swagger', true], + ['static', true], + ]); +}); From e204c852030c4c01e64d72afaf46e1782ed67bb9 Mon Sep 17 00:00:00 2001 From: Hasko Date: Thu, 25 Jun 2026 16:46:05 +0200 Subject: [PATCH 4/5] =?UTF-8?q?fix(types):=20=F0=9F=90=9B=20Support=20clas?= =?UTF-8?q?sic=20TypeScript=20resolution?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/getting-started/installation.md | 7 ++ docs/reference/scheduler.md | 21 ++++- package.json | 79 +++++++++++++++++ tests/typescript-module-resolution.test.mjs | 94 +++++++++++++++++++++ 4 files changed, 200 insertions(+), 1 deletion(-) create mode 100644 tests/typescript-module-resolution.test.mjs 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/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/tests/typescript-module-resolution.test.mjs b/tests/typescript-module-resolution.test.mjs new file mode 100644 index 0000000..c21db44 --- /dev/null +++ b/tests/typescript-module-resolution.test.mjs @@ -0,0 +1,94 @@ +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 { 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 hook: PublisherHook | undefined = undefined; + const policy: SchedulerErrorPolicy | undefined = undefined; + + void middleware; + void messageBus; + 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}`); +}); From 7510b592f6ab1cf4a733882e80da4409bf23e155 Mon Sep 17 00:00:00 2001 From: Hasko Date: Thu, 25 Jun 2026 19:04:08 +0200 Subject: [PATCH 5/5] =?UTF-8?q?fix(kernel):=20=F0=9F=90=9B=20Harden=20inte?= =?UTF-8?q?gration=20extension=20points?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/guides/adapters.md | 19 ++++-- docs/reference/consumer.md | 7 ++- docs/reference/express-kernel-server.md | 25 +++++--- src/adapters/pubsub/Consumer.ts | 16 +++-- .../pubsub/DefaultPublisherHookErrorPolicy.ts | 26 ++++++++ .../pubsub/IdempotencyConsumerMiddleware.ts | 39 ++++++++++-- .../pubsub/InMemoryIdempotencyStore.ts | 20 ++++++ src/adapters/pubsub/PublisherHookPipeline.ts | 61 +++++++++++++++---- .../pubsub/amqp/AmqpMessageBusAdapter.ts | 13 +++- .../amqp/AmqpMessageBusAdapterOptions.ts | 6 +- .../pubsub/amqp/DomainEventHandler.ts | 10 ++- .../pubsub/in-memory/InMemoryEventBus.ts | 7 ++- .../pubsub/in-memory/InMemoryPubSub.ts | 7 ++- src/adapters/pubsub/index.ts | 1 + src/adapters/ui/express/ExpressHookPhase.ts | 4 ++ .../ui/express/ExpressKernelServer.ts | 18 ++++++ .../ui/express/ExpressKernelServerOptions.ts | 8 +++ src/adapters/ui/express/ExpressPhaseHook.ts | 7 +++ src/adapters/ui/express/index.ts | 2 + .../kernel/ConsumerExecutionContext.ts | 1 + src/contracts/kernel/IdempotencyStore.ts | 3 + src/contracts/pubsub/MessageBus.ts | 15 ++++- src/contracts/pubsub/PublishContext.ts | 2 + .../pubsub/PublisherHookErrorPolicy.ts | 9 +++ src/contracts/pubsub/index.ts | 1 + src/domain/DomainEventConsumer.ts | 6 +- src/domain/DomainEventConsumerContext.ts | 3 + src/domain/DomainMessageBus.ts | 8 +++ src/domain/index.ts | 2 + tests/adapters/pubsub/Consumer.test.mjs | 47 +++++++++++++- .../amqp/AmqpMessageBusAdapter.test.mjs | 25 +++++--- .../pubsub/in-memory/InMemoryPubSub.test.mjs | 47 ++++++++++++++ .../ui/express/ExpressKernelServer.test.mjs | 29 +++++++++ tests/typescript-module-resolution.test.mjs | 3 + 34 files changed, 443 insertions(+), 54 deletions(-) create mode 100644 src/adapters/pubsub/DefaultPublisherHookErrorPolicy.ts create mode 100644 src/adapters/ui/express/ExpressHookPhase.ts create mode 100644 src/adapters/ui/express/ExpressPhaseHook.ts create mode 100644 src/contracts/pubsub/PublisherHookErrorPolicy.ts create mode 100644 src/domain/DomainEventConsumerContext.ts create mode 100644 src/domain/DomainMessageBus.ts diff --git a/docs/guides/adapters.md b/docs/guides/adapters.md index 0c566ea..ab44246 100644 --- a/docs/guides/adapters.md +++ b/docs/guides/adapters.md @@ -38,18 +38,29 @@ wrapping the adapter in an application-local class. 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 ({ message }) => { - await websocketPublisher.publish(message); + afterPublish: async ({ domainEvent, message }) => { + await websocketPublisher.publish(domainEvent ?? message); }, }, ], }); ``` -Custom adapters should implement the `MessageBus` contract and delegate hook -execution through `PublisherHookPipeline`: +Custom generic adapters should implement the `MessageBus` contract. Domain-event +adapters should implement `DomainMessageBus`. Both can delegate hook execution +through `PublisherHookPipeline`: ```ts import { diff --git a/docs/reference/consumer.md b/docs/reference/consumer.md index fcf94fa..b738b5c 100644 --- a/docs/reference/consumer.md +++ b/docs/reference/consumer.md @@ -34,7 +34,8 @@ kernel.registerConsumerMiddleware({ Middleware receives the event, the next pipeline callback and a `ConsumerExecutionContext` containing queue, exchange, event id, correlation id -and causation id. +and causation id. Transport adapters can also attach metadata, such as AMQP +headers or retry counts, to `context.metadata`. ## Built-in Middleware @@ -61,5 +62,7 @@ kernel.registerConsumerMiddleware( ); ``` -Use a custom `IdempotencyStore` for durable idempotency. The in-memory store is +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 c5637d4..92d761e 100644 --- a/docs/reference/express-kernel-server.md +++ b/docs/reference/express-kernel-server.md @@ -36,11 +36,14 @@ 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], - swaggerHooks: [(app) => setupSwagger(app)], - staticHooks: [(app) => app.use('/public', express.static('public'))], }); ``` @@ -49,9 +52,15 @@ Hook order is: 1. `middlewares` 2. `preControllerMiddlewares` 3. `beforeControllersHooks` -4. `routing-controllers` -5. `postControllerMiddlewares` -6. `afterControllersHooks` -7. `swaggerHooks` -8. `staticHooks` -9. `errorHandlers` +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/src/adapters/pubsub/Consumer.ts b/src/adapters/pubsub/Consumer.ts index 2c86cd4..4f1cc35 100644 --- a/src/adapters/pubsub/Consumer.ts +++ b/src/adapters/pubsub/Consumer.ts @@ -1,5 +1,8 @@ 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'; @@ -7,8 +10,12 @@ import { ConsumerMiddlewarePipeline } from './ConsumerMiddlewarePipeline.js'; export abstract class Consumer { constructor(private readonly consumer: DomainEventConsumer) {} - private async runMiddleware(event: DomainEvent): Promise { + private async runMiddleware( + event: DomainEvent, + consumerContext?: DomainEventConsumerContext, + ): Promise { const pipeline = new ConsumerMiddlewarePipeline(Kernel.consumerMiddleware); + const metadata = consumerContext?.metadata ?? {}; await pipeline.execute( event, @@ -19,8 +26,9 @@ export abstract class Consumer { eventName: this.eventName, exchange: this.exchange, kernel: Kernel.active, - metadata: {}, + metadata, queueName: this.queueName, + rawMessage: metadata.rawMessage, }, () => this.handler(event), ); @@ -42,7 +50,7 @@ export abstract class Consumer { this.eventName, this.domainEvent, this.exchange, - (event) => this.runMiddleware(event), + (event, context) => this.runMiddleware(event, context), ); } 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 index 930734d..f460eeb 100644 --- a/src/adapters/pubsub/IdempotencyConsumerMiddleware.ts +++ b/src/adapters/pubsub/IdempotencyConsumerMiddleware.ts @@ -9,13 +9,30 @@ import type { IdempotencyConsumerMiddlewareOptions } from './IdempotencyConsumer export class IdempotencyConsumerMiddleware implements ConsumerMiddleware { constructor(private readonly options: IdempotencyConsumerMiddlewareOptions) {} - public async handle( - event: DomainEvent, + private async handleClaimedKey( + key: string, next: ConsumerNext, - context: ConsumerExecutionContext, ): Promise { - const key = this.options.key?.(event, context) ?? context.eventId; + 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; } @@ -23,6 +40,20 @@ export class IdempotencyConsumerMiddleware implements ConsumerMiddleware { 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/InMemoryIdempotencyStore.ts b/src/adapters/pubsub/InMemoryIdempotencyStore.ts index 78db276..1201dcf 100644 --- a/src/adapters/pubsub/InMemoryIdempotencyStore.ts +++ b/src/adapters/pubsub/InMemoryIdempotencyStore.ts @@ -1,8 +1,28 @@ 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); } diff --git a/src/adapters/pubsub/PublisherHookPipeline.ts b/src/adapters/pubsub/PublisherHookPipeline.ts index 4055176..6fecddf 100644 --- a/src/adapters/pubsub/PublisherHookPipeline.ts +++ b/src/adapters/pubsub/PublisherHookPipeline.ts @@ -1,12 +1,57 @@ -import type { PublishContext, PublisherHook } from '../../contracts/index.js'; +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[] = []) { + 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); } @@ -15,22 +60,16 @@ export class PublisherHookPipeline { context: PublishContext, publish: () => Promise | T, ): Promise { - for (const hook of this.hooks) { - await hook.beforePublish?.(context); - } + await this.runBeforePublishHooks(context); try { const result = await publish(); - for (const hook of this.hooks) { - await hook.afterPublish?.(context); - } + await this.runAfterPublishHooks(context); return result; } catch (error: unknown) { - for (const hook of this.hooks) { - await hook.onPublishError?.(error, context); - } + await this.runPublishErrorHooks(error, context); throw error; } diff --git a/src/adapters/pubsub/amqp/AmqpMessageBusAdapter.ts b/src/adapters/pubsub/amqp/AmqpMessageBusAdapter.ts index c63f9dc..389c4ab 100644 --- a/src/adapters/pubsub/amqp/AmqpMessageBusAdapter.ts +++ b/src/adapters/pubsub/amqp/AmqpMessageBusAdapter.ts @@ -13,6 +13,7 @@ import type { Constructor, DomainEvent, DomainEventConsumer, + DomainMessageBus, DomainEventPublisher, } from '../../../domain/index.js'; import type { AmqpMessage } from './AmqpMessage.js'; @@ -25,7 +26,7 @@ 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; @@ -38,6 +39,7 @@ export default class AmqpMessageBusAdapter options.exchange ?? options.serviceName ?? process.env.SERVICE_NAME ?? ''; this.publisherHookPipeline = new PublisherHookPipeline( options.publisherHooks, + options.publisherHookErrorPolicy, ); } @@ -103,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); } @@ -458,6 +466,7 @@ export default class AmqpMessageBusAdapter for (const event of domainEvents) { await this.publisherHookPipeline.run( { + domainEvent: event, message: { metadata: { causationId: event.getCausationId(), diff --git a/src/adapters/pubsub/amqp/AmqpMessageBusAdapterOptions.ts b/src/adapters/pubsub/amqp/AmqpMessageBusAdapterOptions.ts index fc2db4a..0b78792 100644 --- a/src/adapters/pubsub/amqp/AmqpMessageBusAdapterOptions.ts +++ b/src/adapters/pubsub/amqp/AmqpMessageBusAdapterOptions.ts @@ -1,4 +1,7 @@ -import type { PublisherHook } from '../../../contracts/index.js'; +import type { + PublisherHook, + PublisherHookErrorPolicy, +} from '../../../contracts/index.js'; import type { Log } from '../../../infrastructure/logs/index.js'; export interface AmqpMessageBusAdapterOptions { @@ -6,6 +9,7 @@ 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 727c73e..cef6322 100644 --- a/src/adapters/pubsub/in-memory/InMemoryEventBus.ts +++ b/src/adapters/pubsub/in-memory/InMemoryEventBus.ts @@ -3,6 +3,7 @@ import type { HandlerContext, MessageHandler, PublisherHook, + PublisherHookErrorPolicy, } from '../../../contracts/index.js'; import { PublisherHookPipeline } from '../PublisherHookPipeline.js'; @@ -18,8 +19,12 @@ export class InMemoryEventBus { constructor( private readonly context: HandlerContext, publisherHooks: readonly PublisherHook[] = [], + publisherHookErrorPolicy?: PublisherHookErrorPolicy, ) { - this.publisherHookPipeline = new PublisherHookPipeline(publisherHooks); + this.publisherHookPipeline = new PublisherHookPipeline( + publisherHooks, + publisherHookErrorPolicy, + ); } public subscribe( diff --git a/src/adapters/pubsub/in-memory/InMemoryPubSub.ts b/src/adapters/pubsub/in-memory/InMemoryPubSub.ts index 2523f3f..f8205f9 100644 --- a/src/adapters/pubsub/in-memory/InMemoryPubSub.ts +++ b/src/adapters/pubsub/in-memory/InMemoryPubSub.ts @@ -3,6 +3,7 @@ import type { Message, MessageHandler, PublisherHook, + PublisherHookErrorPolicy, Subscription, } from '../../../contracts/index.js'; @@ -19,8 +20,12 @@ export class InMemoryPubSub { constructor( private readonly context: HandlerContext, publisherHooks: readonly PublisherHook[] = [], + publisherHookErrorPolicy?: PublisherHookErrorPolicy, ) { - this.publisherHookPipeline = new PublisherHookPipeline(publisherHooks); + this.publisherHookPipeline = new PublisherHookPipeline( + publisherHooks, + publisherHookErrorPolicy, + ); } public async publish( diff --git a/src/adapters/pubsub/index.ts b/src/adapters/pubsub/index.ts index 823e299..641ce30 100644 --- a/src/adapters/pubsub/index.ts +++ b/src/adapters/pubsub/index.ts @@ -2,6 +2,7 @@ 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'; 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 c7bafa8..4331757 100644 --- a/src/adapters/ui/express/ExpressKernelServer.ts +++ b/src/adapters/ui/express/ExpressKernelServer.ts @@ -41,6 +41,17 @@ export class ExpressKernelServer { } } + 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, @@ -87,6 +98,10 @@ export class ExpressKernelServer { } public async run(): Promise { + if (this.serverInstance) { + throw new Error('HTTP server is already running.'); + } + const controllers = [ ...this.options.kernel.getRoutes(), ...(this.options.controllers ?? []), @@ -96,14 +111,17 @@ export class ExpressKernelServer { 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 6e1bca6..1e5f8b8 100644 --- a/src/adapters/ui/express/ExpressKernelServerOptions.ts +++ b/src/adapters/ui/express/ExpressKernelServerOptions.ts @@ -3,18 +3,26 @@ 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 15f3913..9df26a1 100644 --- a/src/adapters/ui/express/index.ts +++ b/src/adapters/ui/express/index.ts @@ -1,7 +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 index ba8a3ff..65d6032 100644 --- a/src/contracts/kernel/ConsumerExecutionContext.ts +++ b/src/contracts/kernel/ConsumerExecutionContext.ts @@ -8,5 +8,6 @@ export interface ConsumerExecutionContext { readonly exchange: string; readonly kernel: Kernel; readonly metadata: Readonly>; + readonly rawMessage?: unknown; readonly queueName: string; } diff --git a/src/contracts/kernel/IdempotencyStore.ts b/src/contracts/kernel/IdempotencyStore.ts index 97dda55..df6b3b7 100644 --- a/src/contracts/kernel/IdempotencyStore.ts +++ b/src/contracts/kernel/IdempotencyStore.ts @@ -1,4 +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/pubsub/MessageBus.ts b/src/contracts/pubsub/MessageBus.ts index 4afe13b..dde86e9 100644 --- a/src/contracts/pubsub/MessageBus.ts +++ b/src/contracts/pubsub/MessageBus.ts @@ -1,7 +1,16 @@ -import type { DomainEventConsumer } from '../../domain/DomainEventConsumer.js'; -import type { DomainEventPublisher } from '../../domain/DomainEventPublisher.js'; +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 extends DomainEventConsumer, DomainEventPublisher { +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 index 43ce260..6fcd633 100644 --- a/src/contracts/pubsub/PublishContext.ts +++ b/src/contracts/pubsub/PublishContext.ts @@ -1,9 +1,11 @@ +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/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 a119358..54ecf99 100644 --- a/src/contracts/pubsub/index.ts +++ b/src/contracts/pubsub/index.ts @@ -9,4 +9,5 @@ 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 03a118b..fc28912 100644 --- a/tests/adapters/pubsub/Consumer.test.mjs +++ b/tests/adapters/pubsub/Consumer.test.mjs @@ -41,11 +41,17 @@ 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); @@ -54,7 +60,13 @@ test('initializes the domain event consumer with metadata and middleware chain', kernel.registerConsumerMiddleware({ async handle(receivedEvent, next, context) { calls.push(['middleware:before', receivedEvent]); - calls.push(['context', context.eventId, context.queueName]); + calls.push([ + 'context', + context.eventId, + context.queueName, + context.metadata.retries, + context.rawMessage, + ]); await next(); calls.push(['middleware:after', receivedEvent]); }, @@ -65,7 +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'], + ['context', event.eventId, 'test-queue', 1, rawMessage], ['handler', event], ['middleware:after', event], ]); @@ -195,3 +207,32 @@ test('supports middleware defaults and retry predicates', async () => { 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 f498333..fda7663 100644 --- a/tests/adapters/pubsub/amqp/AmqpMessageBusAdapter.test.mjs +++ b/tests/adapters/pubsub/amqp/AmqpMessageBusAdapter.test.mjs @@ -146,8 +146,10 @@ test('publishes domain events and closes channel resources', async () => { serviceName: 'service', }); adapter.registerPublisherHooks({ - afterPublish: (context) => hookCalls.push(['after', context.topic]), - beforePublish: (context) => hookCalls.push(['before', context.topic]), + afterPublish: (context) => + hookCalls.push(['after', context.topic, context.domainEvent]), + beforePublish: (context) => + hookCalls.push(['before', context.topic, context.domainEvent]), }); const event = new TestDomainEvent( 'aggregate-id', @@ -172,8 +174,8 @@ test('publishes domain events and closes channel resources', async () => { true, ); assert.deepEqual(hookCalls, [ - ['before', 'test.domain-event'], - ['after', 'test.domain-event'], + ['before', 'test.domain-event', event], + ['after', 'test.domain-event', event], ]); }); @@ -186,7 +188,10 @@ test('consumes AMQP messages and acknowledges handled events', async () => { 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( @@ -194,14 +199,20 @@ 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]); }); diff --git a/tests/adapters/pubsub/in-memory/InMemoryPubSub.test.mjs b/tests/adapters/pubsub/in-memory/InMemoryPubSub.test.mjs index 2e139e6..9705a35 100644 --- a/tests/adapters/pubsub/in-memory/InMemoryPubSub.test.mjs +++ b/tests/adapters/pubsub/in-memory/InMemoryPubSub.test.mjs @@ -79,3 +79,50 @@ test('runs publisher error hooks before rethrowing', async () => { ); 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 867898b..47b65c1 100644 --- a/tests/adapters/ui/express/ExpressKernelServer.test.mjs +++ b/tests/adapters/ui/express/ExpressKernelServer.test.mjs @@ -121,6 +121,20 @@ test('runs configurable controller, swagger and static hooks', async () => { 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, @@ -135,8 +149,23 @@ test('runs configurable controller, swagger and static hooks', async () => { 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 index c21db44..4891c48 100644 --- a/tests/typescript-module-resolution.test.mjs +++ b/tests/typescript-module-resolution.test.mjs @@ -37,16 +37,19 @@ test('exports types for TypeScript moduleResolution node consumers', async () => ` 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;