From 89138309e72e374e2d074cd4e245f3a17e61af73 Mon Sep 17 00:00:00 2001 From: sebastian Date: Sun, 15 Mar 2026 17:36:25 +0100 Subject: [PATCH 01/16] [messaging] feat: introduce generic listener and hooks --- .../src/bus/in-memory-message-bus.factory.ts | 6 ++- .../src/bus/in-memory-message.bus.ts | 42 ++++++++++++------- .../src/consumer/distributed.consumer.ts | 3 +- .../src/dependency-injection/decorator.ts | 12 ++++++ .../src/dependency-injection/register.ts | 18 +++++++- .../src/listener/listener-handler.ts | 24 +++++++++++ .../src/listener/listener.registry.ts | 8 ++++ .../src/listener/messaging-listener.ts | 12 ++++++ packages/messaging/src/messaging.module.ts | 34 +++++++++------ 9 files changed, 126 insertions(+), 33 deletions(-) create mode 100644 packages/messaging/src/listener/listener-handler.ts create mode 100644 packages/messaging/src/listener/listener.registry.ts create mode 100644 packages/messaging/src/listener/messaging-listener.ts diff --git a/packages/messaging/src/bus/in-memory-message-bus.factory.ts b/packages/messaging/src/bus/in-memory-message-bus.factory.ts index 56ea2f6..fd79590 100644 --- a/packages/messaging/src/bus/in-memory-message-bus.factory.ts +++ b/packages/messaging/src/bus/in-memory-message-bus.factory.ts @@ -8,6 +8,7 @@ import { MessageBusFactory } from '../dependency-injection/decorator'; import { InMemoryMessageBus } from './in-memory-message.bus'; import { IMessageBusFactory } from './i-message-bus.factory'; import { NormalizerRegistry } from '../normalizer/normalizer.registry'; +import { ListenerHandler } from '../listener/listener-handler'; @Injectable() @MessageBusFactory(InMemoryChannel) @@ -19,7 +20,9 @@ export class InMemoryMessageBusFactory implements IMessageBusFactory { const middlewares = []; @@ -30,6 +33,24 @@ export class InMemoryMessageBus implements IMessageBus { HandlerMiddleware, ); + let messageToDispatch = + message instanceof RoutingMessage ? message.message : {}; + + if (message instanceof SealedRoutingMessage) { + const normalizerDefinition: object = + message.messageOptions instanceof DefaultMessageOptions + ? message.messageOptions.normalizer + : ObjectForwardMessageNormalizer; + + messageToDispatch = await this.normalizerRegistry + .getByName(normalizerDefinition['name']) + .denormalize(message.message, message.messageRoutingKey); + } + + const routingMessage = MessageFactory.creteRoutingFromMessage(messageToDispatch, message); + + await this.listenerHandler.handlePreMessageDispatched(routingMessage); + try { this.registry.getByRoutingKey(message.messageRoutingKey); } catch (e) { @@ -60,27 +81,16 @@ export class InMemoryMessageBus implements IMessageBus { DecoratorExtractor.extractMessageMiddleware(middleware), ), ); - const context = MiddlewareContext.createFresh(middlewareInstances); - - let messageToDispatch = - message instanceof RoutingMessage ? message.message : {}; - - if (message instanceof SealedRoutingMessage) { - const normalizerDefinition: object = - message.messageOptions instanceof DefaultMessageOptions - ? message.messageOptions.normalizer - : ObjectForwardMessageNormalizer; - messageToDispatch = await this.normalizerRegistry - .getByName(normalizerDefinition['name']) - .denormalize(message.message, message.messageRoutingKey); - } + const context = MiddlewareContext.createFresh(middlewareInstances); const response = await middlewareInstances[0].process( - MessageFactory.creteRoutingFromMessage(messageToDispatch, message), + routingMessage, context, ); + await this.listenerHandler.handlePostMessageDispatched(routingMessage); + return Promise.resolve(response); } } diff --git a/packages/messaging/src/consumer/distributed.consumer.ts b/packages/messaging/src/consumer/distributed.consumer.ts index 40620b5..54fd476 100644 --- a/packages/messaging/src/consumer/distributed.consumer.ts +++ b/packages/messaging/src/consumer/distributed.consumer.ts @@ -20,7 +20,8 @@ export class DistributedConsumer { private readonly exceptionListenerHandler: ExceptionListenerHandler, @Inject(Service.LOGGER) private readonly logger: MessagingLogger, private readonly discoveryService: DiscoveryService, - ) {} + ) { + } async run(): Promise { for (const channel of this.channelRegistry.getAll()) { diff --git a/packages/messaging/src/dependency-injection/decorator.ts b/packages/messaging/src/dependency-injection/decorator.ts index 56bf903..517c87e 100644 --- a/packages/messaging/src/dependency-injection/decorator.ts +++ b/packages/messaging/src/dependency-injection/decorator.ts @@ -1,4 +1,5 @@ import { ChannelConfig } from '../config'; +import { MessagingListenerHook } from '../listener/messaging-listener'; export const MESSAGE_HANDLER_METADATA = 'MESSAGE_HANDLER_METADATA'; export const CHANNEL_FACTORY_METADATA = 'CHANNEL_FACTORY_METADATA'; @@ -9,6 +10,7 @@ export const MESSAGING_NORMALIZER_METADATA = 'MESSAGING_NORMALIZER_METADATA'; export const MESSAGING_EXCEPTION_LISTENER_METADATA = 'MESSAGING_EXCEPTION_LISTENER_METADATA'; export const MESSAGING_MESSAGE_METADATA = 'MESSAGING_MESSAGE_METADATA'; +export const MESSAGING_LISTENER_METADATA = 'MESSAGING_LISTENER_METADATA'; export const MessageHandler = (...routingKey: string[]): ClassDecorator => { return (target) => { @@ -70,6 +72,16 @@ export const MessagingExceptionListener = (): ClassDecorator => { }; }; +export const MessagingListener = (hook: MessagingListenerHook): ClassDecorator => { + return (target) => { + Reflect.defineMetadata( + MESSAGING_LISTENER_METADATA, + `${hook}:${target.name}`, + target, + ); + }; +}; + export function DenormalizeMessage(): ParameterDecorator { return (target, propertyKey, parameterIndex) => { const paramTypes = Reflect.getMetadata( diff --git a/packages/messaging/src/dependency-injection/register.ts b/packages/messaging/src/dependency-injection/register.ts index 96d01f6..3c22d39 100644 --- a/packages/messaging/src/dependency-injection/register.ts +++ b/packages/messaging/src/dependency-injection/register.ts @@ -4,7 +4,7 @@ import { MessagingLogger } from '../logger/messaging-logger'; import { Service } from './service'; import { MESSAGE_HANDLER_METADATA, - MESSAGING_EXCEPTION_LISTENER_METADATA, + MESSAGING_EXCEPTION_LISTENER_METADATA, MESSAGING_LISTENER_METADATA, MESSAGING_MIDDLEWARE_METADATA, MESSAGING_NORMALIZER_METADATA, } from './decorator'; @@ -13,6 +13,7 @@ import { Registry } from '../shared/registry'; import { MiddlewareRegistry } from '../middleware/middleware.registry'; import { ExceptionListenerRegistry } from '../exception-listener/exception-listener.registry'; import { NormalizerRegistry } from '../normalizer/normalizer.registry'; +import { ListenerRegistry } from '../listener/listener.registry'; export const registerHandlers = ( moduleRef: ModuleRef, @@ -78,10 +79,23 @@ export const registerExceptionListener = ( ); }; +export const registerListener = ( + moduleRef: ModuleRef, + discoveryService: DiscoveryService, +) => { + register( + moduleRef, + discoveryService, + ListenerRegistry, + MESSAGING_LISTENER_METADATA, + 'MessagingListener', + ); +}; + const register = >( moduleRef: ModuleRef, discoveryService: DiscoveryService, - registryProvider: string, + registryProvider: string | Function, decoratorMetadata: string, name: string, ) => { diff --git a/packages/messaging/src/listener/listener-handler.ts b/packages/messaging/src/listener/listener-handler.ts new file mode 100644 index 0000000..d5f504a --- /dev/null +++ b/packages/messaging/src/listener/listener-handler.ts @@ -0,0 +1,24 @@ +import { Injectable } from '@nestjs/common'; +import { ListenerRegistry } from './listener.registry'; +import { MessagingListenerHook } from './messaging-listener'; +import { Message } from '../message/message'; + +@Injectable() +export class ListenerHandler { + constructor( + private readonly listenerRegistry: ListenerRegistry, + ) { + } + + async handlePreMessageDispatched(message: Message): Promise { + await this.listenerRegistry + .getAllByHook(MessagingListenerHook.PRE_MESSAGE_DISPATCHED) + .forEach((listener) => listener.on(message)); + } + + async handlePostMessageDispatched(message: Message): Promise { + await this.listenerRegistry + .getAllByHook(MessagingListenerHook.POST_MESSAGE_DISPATCHED) + .forEach((listener) => listener.on(message)); + } +} diff --git a/packages/messaging/src/listener/listener.registry.ts b/packages/messaging/src/listener/listener.registry.ts new file mode 100644 index 0000000..9601f0f --- /dev/null +++ b/packages/messaging/src/listener/listener.registry.ts @@ -0,0 +1,8 @@ +import { BaseRegistry } from '../shared/base-registry'; +import { MessagingListener, MessagingListenerHook } from './messaging-listener'; + +export class ListenerRegistry extends BaseRegistry { + getAllByHook(hook: MessagingListenerHook): MessagingListener[] { + return this.getAll().filter((listener) => listener.constructor.name.includes(hook)); + } +} diff --git a/packages/messaging/src/listener/messaging-listener.ts b/packages/messaging/src/listener/messaging-listener.ts new file mode 100644 index 0000000..89ec6fd --- /dev/null +++ b/packages/messaging/src/listener/messaging-listener.ts @@ -0,0 +1,12 @@ +import { Message } from '../message/message'; + +export type MessagingData = Message; + +export enum MessagingListenerHook { + PRE_MESSAGE_DISPATCHED = 'PRE_MESSAGE_DISPATCHED', + POST_MESSAGE_DISPATCHED = 'POST_MESSAGE_DISPATCHED', +} + +export interface MessagingListener { + on(data: MessagingData): Promise; +} diff --git a/packages/messaging/src/messaging.module.ts b/packages/messaging/src/messaging.module.ts index 420ec90..01a2ada 100644 --- a/packages/messaging/src/messaging.module.ts +++ b/packages/messaging/src/messaging.module.ts @@ -30,7 +30,7 @@ import { InMemoryChannelFactory } from './channel/factory/in-memory-channel.fact import { DistributedConsumer } from './consumer/distributed.consumer'; import { registerExceptionListener, - registerHandlers, + registerHandlers, registerListener, registerMessageNormalizers, registerMiddlewares, } from './dependency-injection/register'; @@ -44,11 +44,12 @@ import { ObjectForwardMessageNormalizer } from './normalizer/object-forward-mess import { ExceptionListenerRegistry } from './exception-listener/exception-listener.registry'; import { ExceptionListenerHandler } from './exception-listener/exception-listener-handler'; import { MessagingLogger } from './logger/messaging-logger'; +import { ListenerHandler } from './listener/listener-handler'; +import { ListenerRegistry } from './listener/listener.registry'; @Module({}) export class MessagingModule - implements OnApplicationBootstrap, OnModuleDestroy -{ + implements OnApplicationBootstrap, OnModuleDestroy { static forRoot(options: MessagingModuleOptions): DynamicModule { const channels = options.channels ?? []; @@ -146,6 +147,7 @@ export class MessagingModule messageHandlerRegistry: MessageHandlerRegistry, middlewareRegistry: MiddlewareRegistry, normalizerRegistry: NormalizerRegistry, + listenerHandler: ListenerHandler, ) => { return new InMemoryMessageBus( messageHandlerRegistry, @@ -158,12 +160,14 @@ export class MessagingModule }), ), normalizerRegistry, + listenerHandler, ); }, inject: [ Service.MESSAGE_HANDLERS_REGISTRY, Service.MIDDLEWARE_REGISTRY, Service.MESSAGE_NORMALIZERS_REGISTRY, + ListenerHandler, ], }; }; @@ -172,15 +176,15 @@ export class MessagingModule options.customLogger && typeof options.customLogger === 'function' ? { provide: Service.LOGGER, useClass: options.customLogger } : { - provide: Service.LOGGER, - useValue: - options.customLogger ?? - new NestLogger( - new NestCommonLogger(), - options.debug ?? false, - options.logging ?? true, - ), - } + provide: Service.LOGGER, + useValue: + options.customLogger ?? + new NestLogger( + new NestCommonLogger(), + options.debug ?? false, + options.logging ?? true, + ), + } ) as Provider; return { @@ -230,6 +234,8 @@ export class MessagingModule InMemoryChannelFactory, DistributedConsumer, ObjectForwardMessageNormalizer, + ListenerRegistry, + ListenerHandler, ], exports: [ Service.DEFAULT_MESSAGE_BUS, @@ -247,13 +253,15 @@ export class MessagingModule private readonly configuration: MandatoryMessagingModuleOptions, @Inject(Service.LOGGER) private readonly logger: MessagingLogger, - ) {} + ) { + } onApplicationBootstrap(): void { registerHandlers(this.moduleRef, this.discoveryService); registerMiddlewares(this.moduleRef, this.discoveryService); registerMessageNormalizers(this.moduleRef, this.discoveryService); registerExceptionListener(this.moduleRef, this.discoveryService); + registerListener(this.moduleRef, this.discoveryService); if (this.configuration.forceDisableAllConsumers ?? false) { this.logger.log( From 103306e0a0301678309915b48d6950ee7d0ca6fd Mon Sep 17 00:00:00 2001 From: sebastian Date: Sun, 15 Mar 2026 19:39:48 +0100 Subject: [PATCH 02/16] [messaging] feat: lifecycle hooks --- .../src/bus/in-memory-message-bus.factory.ts | 6 +- .../src/bus/in-memory-message.bus.ts | 25 ++++++--- .../src/dependency-injection/decorator.ts | 10 ++-- .../src/dependency-injection/register.ts | 30 +++++----- packages/messaging/src/index.ts | 1 + .../messaging-lifecycle-hook-handler.ts | 30 ++++++++++ .../messaging-lifecycle-hook-listener.ts | 13 +++++ .../messaging-lifecycle-hook.registry.ts | 18 ++++++ .../src/listener/listener-handler.ts | 24 -------- .../src/listener/listener.registry.ts | 8 --- .../src/listener/messaging-listener.ts | 12 ---- packages/messaging/src/messaging.module.ts | 18 +++--- .../messaging/src/shared/base-registry.ts | 2 +- .../unit/bus/in-memory-message.bus.spec.ts | 12 ++++ .../messaging-hook-handler.spec.ts | 56 +++++++++++++++++++ .../messaging-lifecycle-hook.registry.spec.ts | 45 +++++++++++++++ 16 files changed, 226 insertions(+), 84 deletions(-) create mode 100644 packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-handler.ts create mode 100644 packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-listener.ts create mode 100644 packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook.registry.ts delete mode 100644 packages/messaging/src/listener/listener-handler.ts delete mode 100644 packages/messaging/src/listener/listener.registry.ts delete mode 100644 packages/messaging/src/listener/messaging-listener.ts create mode 100644 packages/messaging/test/unit/lifecycle-hook/messaging-hook-handler.spec.ts create mode 100644 packages/messaging/test/unit/lifecycle-hook/messaging-lifecycle-hook.registry.spec.ts diff --git a/packages/messaging/src/bus/in-memory-message-bus.factory.ts b/packages/messaging/src/bus/in-memory-message-bus.factory.ts index fd79590..49fb749 100644 --- a/packages/messaging/src/bus/in-memory-message-bus.factory.ts +++ b/packages/messaging/src/bus/in-memory-message-bus.factory.ts @@ -8,7 +8,7 @@ import { MessageBusFactory } from '../dependency-injection/decorator'; import { InMemoryMessageBus } from './in-memory-message.bus'; import { IMessageBusFactory } from './i-message-bus.factory'; import { NormalizerRegistry } from '../normalizer/normalizer.registry'; -import { ListenerHandler } from '../listener/listener-handler'; +import { MessagingLifecycleHookHandler } from '../lifecycle-hook/messaging-lifecycle-hook-handler'; @Injectable() @MessageBusFactory(InMemoryChannel) @@ -20,7 +20,7 @@ export class InMemoryMessageBusFactory implements IMessageBusFactory { const middlewares = []; + // Execution order: channel middlewares -> message middlewares -> handler middleware. middlewares.push( ...(this.channel.config?.middlewares ?? []), ...(message.messageOptions?.middlewares ?? []), @@ -37,6 +38,7 @@ export class InMemoryMessageBus implements IMessageBus { message instanceof RoutingMessage ? message.message : {}; if (message instanceof SealedRoutingMessage) { + // Sealed messages carry raw payload and must be denormalized before dispatch. const normalizerDefinition: object = message.messageOptions instanceof DefaultMessageOptions ? message.messageOptions.normalizer @@ -47,9 +49,10 @@ export class InMemoryMessageBus implements IMessageBus { .denormalize(message.message, message.messageRoutingKey); } - const routingMessage = MessageFactory.creteRoutingFromMessage(messageToDispatch, message); - - await this.listenerHandler.handlePreMessageDispatched(routingMessage); + // Hook fired once the payload shape is ready for handler pipeline. + await this.messagingHookHandler.handleAfterMessageDenormalized( + MessageFactory.creteRoutingFromMessage(messageToDispatch, message), + ); try { this.registry.getByRoutingKey(message.messageRoutingKey); @@ -69,6 +72,7 @@ export class InMemoryMessageBus implements IMessageBus { avoidErrorsForNonExistedHandlers; } + // Missing handler can be configured as no-op for fire-and-forget scenarios. if (avoidErrorsForNonExistedHandlers) { return Promise.resolve(); } @@ -84,12 +88,19 @@ export class InMemoryMessageBus implements IMessageBus { const context = MiddlewareContext.createFresh(middlewareInstances); + // Hook around handler execution. + await this.messagingHookHandler.handleBeforeMessageHandler( + MessageFactory.creteRoutingFromMessage(messageToDispatch, message), + ); + const response = await middlewareInstances[0].process( - routingMessage, + MessageFactory.creteRoutingFromMessage(messageToDispatch, message), context, ); - await this.listenerHandler.handlePostMessageDispatched(routingMessage); + await this.messagingHookHandler.handleAfterMessageHandlerExecuted( + MessageFactory.creteRoutingFromMessage(messageToDispatch, message), + ); return Promise.resolve(response); } diff --git a/packages/messaging/src/dependency-injection/decorator.ts b/packages/messaging/src/dependency-injection/decorator.ts index 517c87e..e3d9079 100644 --- a/packages/messaging/src/dependency-injection/decorator.ts +++ b/packages/messaging/src/dependency-injection/decorator.ts @@ -1,5 +1,5 @@ import { ChannelConfig } from '../config'; -import { MessagingListenerHook } from '../listener/messaging-listener'; +import { LifecycleHook } from '../lifecycle-hook/messaging-lifecycle-hook-listener'; export const MESSAGE_HANDLER_METADATA = 'MESSAGE_HANDLER_METADATA'; export const CHANNEL_FACTORY_METADATA = 'CHANNEL_FACTORY_METADATA'; @@ -10,7 +10,7 @@ export const MESSAGING_NORMALIZER_METADATA = 'MESSAGING_NORMALIZER_METADATA'; export const MESSAGING_EXCEPTION_LISTENER_METADATA = 'MESSAGING_EXCEPTION_LISTENER_METADATA'; export const MESSAGING_MESSAGE_METADATA = 'MESSAGING_MESSAGE_METADATA'; -export const MESSAGING_LISTENER_METADATA = 'MESSAGING_LISTENER_METADATA'; +export const MESSAGING_LIFECYCLE_HOOK_METADATA = 'MESSAGING_LIFECYCLE_HOOK_METADATA'; export const MessageHandler = (...routingKey: string[]): ClassDecorator => { return (target) => { @@ -72,11 +72,11 @@ export const MessagingExceptionListener = (): ClassDecorator => { }; }; -export const MessagingListener = (hook: MessagingListenerHook): ClassDecorator => { +export const MessagingLifecycleHook = (lifecycleHook: LifecycleHook): ClassDecorator => { return (target) => { Reflect.defineMetadata( - MESSAGING_LISTENER_METADATA, - `${hook}:${target.name}`, + MESSAGING_LIFECYCLE_HOOK_METADATA, + `${lifecycleHook}:${target.name}`, target, ); }; diff --git a/packages/messaging/src/dependency-injection/register.ts b/packages/messaging/src/dependency-injection/register.ts index 3c22d39..91da24d 100644 --- a/packages/messaging/src/dependency-injection/register.ts +++ b/packages/messaging/src/dependency-injection/register.ts @@ -4,7 +4,7 @@ import { MessagingLogger } from '../logger/messaging-logger'; import { Service } from './service'; import { MESSAGE_HANDLER_METADATA, - MESSAGING_EXCEPTION_LISTENER_METADATA, MESSAGING_LISTENER_METADATA, + MESSAGING_EXCEPTION_LISTENER_METADATA, MESSAGING_LIFECYCLE_HOOK_METADATA, MESSAGING_MIDDLEWARE_METADATA, MESSAGING_NORMALIZER_METADATA, } from './decorator'; @@ -13,7 +13,7 @@ import { Registry } from '../shared/registry'; import { MiddlewareRegistry } from '../middleware/middleware.registry'; import { ExceptionListenerRegistry } from '../exception-listener/exception-listener.registry'; import { NormalizerRegistry } from '../normalizer/normalizer.registry'; -import { ListenerRegistry } from '../listener/listener.registry'; +import { MessagingLifecycleHookRegistry } from '../lifecycle-hook/messaging-lifecycle-hook.registry'; export const registerHandlers = ( moduleRef: ModuleRef, @@ -79,16 +79,16 @@ export const registerExceptionListener = ( ); }; -export const registerListener = ( +export const registerMessagingHooks = ( moduleRef: ModuleRef, discoveryService: DiscoveryService, ) => { - register( + register( moduleRef, discoveryService, - ListenerRegistry, - MESSAGING_LISTENER_METADATA, - 'MessagingListener', + MessagingLifecycleHookRegistry, + MESSAGING_LIFECYCLE_HOOK_METADATA, + 'MessagingLifecycleHook', ); }; @@ -104,24 +104,24 @@ const register = >( const logger: MessagingLogger = moduleRef.get(Service.LOGGER); const instances = discoveryService .getProviders() - .filter((messageExceptionListener) => { - if (!messageExceptionListener.metatype) { + .filter((provider) => { + if (!provider.metatype) { return false; } return Reflect.hasMetadata( decoratorMetadata, - messageExceptionListener.metatype, + provider.metatype, ); }); - instances.forEach((messageExceptionListener) => { + instances.forEach((provider) => { registry.register( - Reflect.getMetadata(decoratorMetadata, messageExceptionListener.metatype), - messageExceptionListener.instance, + Reflect.getMetadata(decoratorMetadata, provider.metatype), + provider.instance, ); - if (!exceptions.includes(messageExceptionListener.name)) { - logger.log(`${name} [${messageExceptionListener.name}] was registered`); + if (!exceptions.includes(provider.name)) { + logger.log(`${name} [${provider.name}] was registered`); } }); }; diff --git a/packages/messaging/src/index.ts b/packages/messaging/src/index.ts index a750387..8e1457b 100644 --- a/packages/messaging/src/index.ts +++ b/packages/messaging/src/index.ts @@ -31,3 +31,4 @@ export * from './exception/handlers.exception'; export * from './logger/log'; export * from './logger/messaging-logger'; export * from './logger/nest-logger'; +export * from './lifecycle-hook/messaging-lifecycle-hook-listener' diff --git a/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-handler.ts b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-handler.ts new file mode 100644 index 0000000..700a8ea --- /dev/null +++ b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-handler.ts @@ -0,0 +1,30 @@ +import { Injectable } from '@nestjs/common'; +import { MessagingLifecycleHookRegistry } from './messaging-lifecycle-hook.registry'; +import { LifecycleHook } from './messaging-lifecycle-hook-listener'; +import { RoutingMessage } from '../message/routing-message'; + +@Injectable() +export class MessagingLifecycleHookHandler { + constructor( + private readonly messagingHookRegistry: MessagingLifecycleHookRegistry, + ) { + } + + async handleAfterMessageDenormalized(message: RoutingMessage): Promise { + await this.messagingHookRegistry + .getAllByHook(LifecycleHook.AFTER_MESSAGE_DENORMALIZED) + .forEach((listener) => listener.on(message)); + } + + async handleBeforeMessageHandler(message: RoutingMessage): Promise { + await this.messagingHookRegistry + .getAllByHook(LifecycleHook.BEFORE_MESSAGE_HANDLER) + .forEach((listener) => listener.on(message)); + } + + async handleAfterMessageHandlerExecuted(message: RoutingMessage): Promise { + await this.messagingHookRegistry + .getAllByHook(LifecycleHook.AFTER_MESSAGE_HANDLER_EXECUTED) + .forEach((listener) => listener.on(message)); + } +} diff --git a/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-listener.ts b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-listener.ts new file mode 100644 index 0000000..532bdb7 --- /dev/null +++ b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-listener.ts @@ -0,0 +1,13 @@ +import { RoutingMessage } from '../message/routing-message'; + +export type MessagingData = RoutingMessage; + +export enum LifecycleHook { + AFTER_MESSAGE_DENORMALIZED = 'AFTER_MESSAGE_DENORMALIZED', + BEFORE_MESSAGE_HANDLER = 'BEFORE_MESSAGE_HANDLER', + AFTER_MESSAGE_HANDLER_EXECUTED = 'AFTER_MESSAGE_HANDLER_EXECUTED', +} + +export interface MessagingLifecycleHookListener { + on(data: MessagingData): Promise; +} diff --git a/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook.registry.ts b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook.registry.ts new file mode 100644 index 0000000..4a9cd71 --- /dev/null +++ b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook.registry.ts @@ -0,0 +1,18 @@ +import { BaseRegistry } from '../shared/base-registry'; +import { MessagingLifecycleHookListener, LifecycleHook } from './messaging-lifecycle-hook-listener'; + +export class MessagingLifecycleHookRegistry extends BaseRegistry { + getAllByHook(hook: LifecycleHook): MessagingLifecycleHookListener[] { + const result: MessagingLifecycleHookListener[] = []; + + for (const [key, listener] of this.registry.entries()) { + //Because pattern is `${hook}:${target.name}`, + const name = key.split(':')[0]; + if (name.includes(hook)) { + result.push(listener); + } + } + + return result; + } +} diff --git a/packages/messaging/src/listener/listener-handler.ts b/packages/messaging/src/listener/listener-handler.ts deleted file mode 100644 index d5f504a..0000000 --- a/packages/messaging/src/listener/listener-handler.ts +++ /dev/null @@ -1,24 +0,0 @@ -import { Injectable } from '@nestjs/common'; -import { ListenerRegistry } from './listener.registry'; -import { MessagingListenerHook } from './messaging-listener'; -import { Message } from '../message/message'; - -@Injectable() -export class ListenerHandler { - constructor( - private readonly listenerRegistry: ListenerRegistry, - ) { - } - - async handlePreMessageDispatched(message: Message): Promise { - await this.listenerRegistry - .getAllByHook(MessagingListenerHook.PRE_MESSAGE_DISPATCHED) - .forEach((listener) => listener.on(message)); - } - - async handlePostMessageDispatched(message: Message): Promise { - await this.listenerRegistry - .getAllByHook(MessagingListenerHook.POST_MESSAGE_DISPATCHED) - .forEach((listener) => listener.on(message)); - } -} diff --git a/packages/messaging/src/listener/listener.registry.ts b/packages/messaging/src/listener/listener.registry.ts deleted file mode 100644 index 9601f0f..0000000 --- a/packages/messaging/src/listener/listener.registry.ts +++ /dev/null @@ -1,8 +0,0 @@ -import { BaseRegistry } from '../shared/base-registry'; -import { MessagingListener, MessagingListenerHook } from './messaging-listener'; - -export class ListenerRegistry extends BaseRegistry { - getAllByHook(hook: MessagingListenerHook): MessagingListener[] { - return this.getAll().filter((listener) => listener.constructor.name.includes(hook)); - } -} diff --git a/packages/messaging/src/listener/messaging-listener.ts b/packages/messaging/src/listener/messaging-listener.ts deleted file mode 100644 index 89ec6fd..0000000 --- a/packages/messaging/src/listener/messaging-listener.ts +++ /dev/null @@ -1,12 +0,0 @@ -import { Message } from '../message/message'; - -export type MessagingData = Message; - -export enum MessagingListenerHook { - PRE_MESSAGE_DISPATCHED = 'PRE_MESSAGE_DISPATCHED', - POST_MESSAGE_DISPATCHED = 'POST_MESSAGE_DISPATCHED', -} - -export interface MessagingListener { - on(data: MessagingData): Promise; -} diff --git a/packages/messaging/src/messaging.module.ts b/packages/messaging/src/messaging.module.ts index 01a2ada..63f8da6 100644 --- a/packages/messaging/src/messaging.module.ts +++ b/packages/messaging/src/messaging.module.ts @@ -30,7 +30,7 @@ import { InMemoryChannelFactory } from './channel/factory/in-memory-channel.fact import { DistributedConsumer } from './consumer/distributed.consumer'; import { registerExceptionListener, - registerHandlers, registerListener, + registerHandlers, registerMessagingHooks, registerMessageNormalizers, registerMiddlewares, } from './dependency-injection/register'; @@ -44,8 +44,8 @@ import { ObjectForwardMessageNormalizer } from './normalizer/object-forward-mess import { ExceptionListenerRegistry } from './exception-listener/exception-listener.registry'; import { ExceptionListenerHandler } from './exception-listener/exception-listener-handler'; import { MessagingLogger } from './logger/messaging-logger'; -import { ListenerHandler } from './listener/listener-handler'; -import { ListenerRegistry } from './listener/listener.registry'; +import { MessagingLifecycleHookHandler } from './lifecycle-hook/messaging-lifecycle-hook-handler'; +import { MessagingLifecycleHookRegistry } from './lifecycle-hook/messaging-lifecycle-hook.registry'; @Module({}) export class MessagingModule @@ -147,7 +147,7 @@ export class MessagingModule messageHandlerRegistry: MessageHandlerRegistry, middlewareRegistry: MiddlewareRegistry, normalizerRegistry: NormalizerRegistry, - listenerHandler: ListenerHandler, + messagingHookHandler: MessagingLifecycleHookHandler, ) => { return new InMemoryMessageBus( messageHandlerRegistry, @@ -160,14 +160,14 @@ export class MessagingModule }), ), normalizerRegistry, - listenerHandler, + messagingHookHandler, ); }, inject: [ Service.MESSAGE_HANDLERS_REGISTRY, Service.MIDDLEWARE_REGISTRY, Service.MESSAGE_NORMALIZERS_REGISTRY, - ListenerHandler, + MessagingLifecycleHookHandler, ], }; }; @@ -234,8 +234,8 @@ export class MessagingModule InMemoryChannelFactory, DistributedConsumer, ObjectForwardMessageNormalizer, - ListenerRegistry, - ListenerHandler, + MessagingLifecycleHookRegistry, + MessagingLifecycleHookHandler, ], exports: [ Service.DEFAULT_MESSAGE_BUS, @@ -261,7 +261,7 @@ export class MessagingModule registerMiddlewares(this.moduleRef, this.discoveryService); registerMessageNormalizers(this.moduleRef, this.discoveryService); registerExceptionListener(this.moduleRef, this.discoveryService); - registerListener(this.moduleRef, this.discoveryService); + registerMessagingHooks(this.moduleRef, this.discoveryService); if (this.configuration.forceDisableAllConsumers ?? false) { this.logger.log( diff --git a/packages/messaging/src/shared/base-registry.ts b/packages/messaging/src/shared/base-registry.ts index eb9dc3d..45e0096 100644 --- a/packages/messaging/src/shared/base-registry.ts +++ b/packages/messaging/src/shared/base-registry.ts @@ -2,7 +2,7 @@ import { Registry } from './registry'; import { MessagingException } from '../exception/messaging.exception'; export abstract class BaseRegistry implements Registry { - private registry: Map = new Map(); + protected registry: Map = new Map(); register(name: string, middleware: T): void { if (this.registry.has(name)) { diff --git a/packages/messaging/test/unit/bus/in-memory-message.bus.spec.ts b/packages/messaging/test/unit/bus/in-memory-message.bus.spec.ts index c35e29c..338abba 100644 --- a/packages/messaging/test/unit/bus/in-memory-message.bus.spec.ts +++ b/packages/messaging/test/unit/bus/in-memory-message.bus.spec.ts @@ -8,17 +8,25 @@ import { MessageHandlerRegistry } from '../../../src/handler/message-handler.reg import { MiddlewareRegistry } from '../../../src/middleware/middleware.registry'; import { InMemoryChannel } from '../../../src/channel/in-memory.channel'; import { NormalizerRegistry } from '../../../src/normalizer/normalizer.registry'; +import { MessagingLifecycleHookHandler } from '../../../src/lifecycle-hook/messaging-lifecycle-hook-handler'; describe('InMemoryMessageBus', () => { let handlerRegistry: MessageHandlerRegistry; let middlewareRegistry: MiddlewareRegistry; let normalizerRegistry: NormalizerRegistry; + let messagingHookHandler: MessagingLifecycleHookHandler; let defaultMiddleware: Middleware; beforeEach(async () => { handlerRegistry = new MessageHandlerRegistry(); middlewareRegistry = new MiddlewareRegistry(); normalizerRegistry = new NormalizerRegistry(); + messagingHookHandler = { + handleAfterConsumerDispatchMessage: jest.fn(), + handleAfterMessageDenormalized: jest.fn(), + handleBeforeMessageHandler: jest.fn(), + handleAfterMessageHandlerExecuted: jest.fn(), + } as unknown as MessagingLifecycleHookHandler; defaultMiddleware = { process: jest.fn().mockImplementation(() => { @@ -37,6 +45,7 @@ describe('InMemoryMessageBus', () => { name: 'example.bus', }), normalizerRegistry, + messagingHookHandler, ); await subjectUnderTest.dispatch( @@ -52,6 +61,7 @@ describe('InMemoryMessageBus', () => { name: 'default.bus', }), normalizerRegistry, + messagingHookHandler, ); await subjectUnderTest.dispatch( @@ -68,6 +78,7 @@ describe('InMemoryMessageBus', () => { avoidErrorsForNotExistedHandlers: false, }), normalizerRegistry, + messagingHookHandler, ); await expect( @@ -94,6 +105,7 @@ describe('InMemoryMessageBus', () => { avoidErrorsForNotExistedHandlers: false, }), normalizerRegistry, + messagingHookHandler, ); const response = await subjectUnderTest.dispatch( diff --git a/packages/messaging/test/unit/lifecycle-hook/messaging-hook-handler.spec.ts b/packages/messaging/test/unit/lifecycle-hook/messaging-hook-handler.spec.ts new file mode 100644 index 0000000..bf8bba4 --- /dev/null +++ b/packages/messaging/test/unit/lifecycle-hook/messaging-hook-handler.spec.ts @@ -0,0 +1,56 @@ +import { MessagingLifecycleHookHandler } from '../../../src/lifecycle-hook/messaging-lifecycle-hook-handler'; +import { MessagingLifecycleHookRegistry } from '../../../src/lifecycle-hook/messaging-lifecycle-hook.registry'; +import { + LifecycleHook, + MessagingLifecycleHookListener, +} from '../../../src/lifecycle-hook/messaging-lifecycle-hook-listener'; +import { RoutingMessage } from '../../../src/message/routing-message'; + +describe('MessagingHookHandler', () => { + let hookRegistry: MessagingLifecycleHookRegistry; + let handler: MessagingLifecycleHookHandler; + let listener: MessagingLifecycleHookListener; + const message = new RoutingMessage({ title: 'hello' }, 'test.key'); + + beforeEach(() => { + hookRegistry = { + getAllByHook: jest.fn(), + } as unknown as MessagingLifecycleHookRegistry; + + handler = new MessagingLifecycleHookHandler(hookRegistry); + listener = { on: jest.fn().mockResolvedValue(undefined) }; + }); + + test('should execute listeners for AFTER_MESSAGE_DENORMALIZED hook', async () => { + (hookRegistry.getAllByHook as jest.Mock).mockReturnValue([listener]); + + await handler.handleAfterMessageDenormalized(message); + + expect(hookRegistry.getAllByHook).toHaveBeenCalledWith( + LifecycleHook.AFTER_MESSAGE_DENORMALIZED, + ); + expect(listener.on).toHaveBeenCalledWith(message); + }); + + test('should execute listeners for BEFORE_MESSAGE_HANDLER hook', async () => { + (hookRegistry.getAllByHook as jest.Mock).mockReturnValue([listener]); + + await handler.handleBeforeMessageHandler(message); + + expect(hookRegistry.getAllByHook).toHaveBeenCalledWith( + LifecycleHook.BEFORE_MESSAGE_HANDLER, + ); + expect(listener.on).toHaveBeenCalledWith(message); + }); + + test('should execute listeners for AFTER_MESSAGE_HANDLER_EXECUTED hook', async () => { + (hookRegistry.getAllByHook as jest.Mock).mockReturnValue([listener]); + + await handler.handleAfterMessageHandlerExecuted(message); + + expect(hookRegistry.getAllByHook).toHaveBeenCalledWith( + LifecycleHook.AFTER_MESSAGE_HANDLER_EXECUTED, + ); + expect(listener.on).toHaveBeenCalledWith(message); + }); +}); diff --git a/packages/messaging/test/unit/lifecycle-hook/messaging-lifecycle-hook.registry.spec.ts b/packages/messaging/test/unit/lifecycle-hook/messaging-lifecycle-hook.registry.spec.ts new file mode 100644 index 0000000..7d302e7 --- /dev/null +++ b/packages/messaging/test/unit/lifecycle-hook/messaging-lifecycle-hook.registry.spec.ts @@ -0,0 +1,45 @@ +import { MessagingLifecycleHookRegistry } from '../../../src/lifecycle-hook/messaging-lifecycle-hook.registry'; +import { + LifecycleHook, + MessagingLifecycleHookListener, +} from '../../../src/lifecycle-hook/messaging-lifecycle-hook-listener'; + +describe('MessagingLifecycleHookRegistry', () => { + let registry: MessagingLifecycleHookRegistry; + let beforeListener: MessagingLifecycleHookListener; + let afterDenormalizedListener: MessagingLifecycleHookListener; + + beforeEach(() => { + registry = new MessagingLifecycleHookRegistry(); + beforeListener = { on: jest.fn() } as unknown as MessagingLifecycleHookListener; + afterDenormalizedListener = { + on: jest.fn(), + } as unknown as MessagingLifecycleHookListener; + }); + + test('should return listeners only for selected lifecycle hook', () => { + registry.register( + `${LifecycleHook.BEFORE_MESSAGE_HANDLER}:BeforeListener`, + beforeListener, + ); + registry.register( + `${LifecycleHook.AFTER_MESSAGE_DENORMALIZED}:AfterDenormalizedListener`, + afterDenormalizedListener, + ); + + expect(registry.getAllByHook(LifecycleHook.BEFORE_MESSAGE_HANDLER)).toEqual([ + beforeListener, + ]); + }); + + test('should return empty array when no listeners are registered for given lifecycle hook', () => { + registry.register( + `${LifecycleHook.BEFORE_MESSAGE_HANDLER}:BeforeListener`, + beforeListener, + ); + + expect( + registry.getAllByHook(LifecycleHook.AFTER_MESSAGE_HANDLER_EXECUTED), + ).toEqual([]); + }); +}); From bb74908b9bd5b56bed09cbea6ad6a1cadd4b94b3 Mon Sep 17 00:00:00 2001 From: sebastian Date: Sun, 15 Mar 2026 19:40:04 +0100 Subject: [PATCH 03/16] [messaging] feat: lifecycle hooks --- packages/messaging/package.json | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/messaging/package.json b/packages/messaging/package.json index 44b0fee..6bff5d1 100644 --- a/packages/messaging/package.json +++ b/packages/messaging/package.json @@ -1,6 +1,6 @@ { "name": "@nestjstools/messaging", - "version": "4.1.1", + "version": "4.2.0", "description": "Simplifies asynchronous and synchronous message handling with support for buses, handlers, channels, and consumers. Build scalable, decoupled applications with ease and reliability.", "author": "Sebastian Iwanczyszyn", "license": "MIT", From 1b999e099890698f36a8b7640c2fcb91be232aee Mon Sep 17 00:00:00 2001 From: sebastian Date: Sun, 15 Mar 2026 19:40:33 +0100 Subject: [PATCH 04/16] [messaging] feat: lifecycle hooks --- packages/messaging/package.json | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/messaging/package.json b/packages/messaging/package.json index 6bff5d1..23e7349 100644 --- a/packages/messaging/package.json +++ b/packages/messaging/package.json @@ -1,6 +1,6 @@ { "name": "@nestjstools/messaging", - "version": "4.2.0", + "version": "4.2.0-beta.0", "description": "Simplifies asynchronous and synchronous message handling with support for buses, handlers, channels, and consumers. Build scalable, decoupled applications with ease and reliability.", "author": "Sebastian Iwanczyszyn", "license": "MIT", From 5b71beb5b43b2ad2e95b8fccea892e091a50862b Mon Sep 17 00:00:00 2001 From: sebastian Date: Tue, 17 Mar 2026 17:48:43 +0100 Subject: [PATCH 05/16] [messaging] feat: lifecycle hooks --- .../messaging/src/bus/consumer.message-bus.ts | 14 +++++++++- .../src/consumer/distributed.consumer.ts | 6 +++- .../messaging-lifecycle-hook-handler.ts | 13 +++++++-- .../messaging-lifecycle-hook-listener.ts | 28 +++++++++++++++++-- .../messaging-hook-handler.spec.ts | 8 +++--- 5 files changed, 58 insertions(+), 11 deletions(-) diff --git a/packages/messaging/src/bus/consumer.message-bus.ts b/packages/messaging/src/bus/consumer.message-bus.ts index 648405c..d7ca6ca 100644 --- a/packages/messaging/src/bus/consumer.message-bus.ts +++ b/packages/messaging/src/bus/consumer.message-bus.ts @@ -11,6 +11,8 @@ import { ConsumerDispatchedMessageError } from '../consumer/consumer-dispatched- import { HandlersException } from '../exception/handlers.exception'; import { ExceptionListenerHandler } from '../exception-listener/exception-listener-handler'; import { ExceptionContext } from '../exception-listener/exception-context'; +import { MessagingLifecycleHookHandler } from '../lifecycle-hook/messaging-lifecycle-hook-handler'; +import { DetailedConsumerMessage } from '../lifecycle-hook/messaging-lifecycle-hook-listener'; export class ConsumerMessageBus { constructor( @@ -19,7 +21,9 @@ export class ConsumerMessageBus { private readonly logger: MessagingLogger, private readonly consumer: IMessagingConsumer, private readonly exceptionListenerHandler: ExceptionListenerHandler, - ) {} + private readonly messagingHookHandler: MessagingLifecycleHookHandler, + ) { + } async dispatch(consumerMessage: ConsumerMessage): Promise { try { @@ -74,6 +78,14 @@ export class ConsumerMessageBus { consumerMessage.routingKey, ), ); + + await this.messagingHookHandler.handleOnFailedMessageConsumer( + DetailedConsumerMessage.fromConsumerMessage( + consumerMessage, + this.channel.config.name, + this.channel.constructor.name, + ), + ); } return Promise.resolve(); diff --git a/packages/messaging/src/consumer/distributed.consumer.ts b/packages/messaging/src/consumer/distributed.consumer.ts index 54fd476..41015c7 100644 --- a/packages/messaging/src/consumer/distributed.consumer.ts +++ b/packages/messaging/src/consumer/distributed.consumer.ts @@ -1,4 +1,4 @@ -import { Inject } from '@nestjs/common'; +import { Inject, Injectable } from '@nestjs/common'; import { DiscoveryService } from '@nestjs/core'; import { Service } from '../dependency-injection/service'; import { IMessageBus } from '../bus/i-message-bus'; @@ -9,7 +9,9 @@ import { MESSAGE_CONSUMER_METADATA } from '../dependency-injection/decorator'; import { IMessagingConsumer } from './i-messaging-consumer'; import { ExceptionListenerHandler } from '../exception-listener/exception-listener-handler'; import { ConsumerMessageBus } from '../bus/consumer.message-bus'; +import { MessagingLifecycleHookHandler } from '../lifecycle-hook/messaging-lifecycle-hook-handler'; +@Injectable() export class DistributedConsumer { constructor( @Inject(Service.DEFAULT_MESSAGE_BUS) @@ -20,6 +22,7 @@ export class DistributedConsumer { private readonly exceptionListenerHandler: ExceptionListenerHandler, @Inject(Service.LOGGER) private readonly logger: MessagingLogger, private readonly discoveryService: DiscoveryService, + private readonly messagingLifecycleHookHandler: MessagingLifecycleHookHandler, ) { } @@ -64,6 +67,7 @@ export class DistributedConsumer { this.logger, consumer, this.exceptionListenerHandler, + this.messagingLifecycleHookHandler, ); await consumer.consume(dispatcher, channel); diff --git a/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-handler.ts b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-handler.ts index 700a8ea..0b7378d 100644 --- a/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-handler.ts +++ b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-handler.ts @@ -2,6 +2,7 @@ import { Injectable } from '@nestjs/common'; import { MessagingLifecycleHookRegistry } from './messaging-lifecycle-hook.registry'; import { LifecycleHook } from './messaging-lifecycle-hook-listener'; import { RoutingMessage } from '../message/routing-message'; +import { ConsumerMessage } from '../consumer/consumer-message'; @Injectable() export class MessagingLifecycleHookHandler { @@ -13,18 +14,24 @@ export class MessagingLifecycleHookHandler { async handleAfterMessageDenormalized(message: RoutingMessage): Promise { await this.messagingHookRegistry .getAllByHook(LifecycleHook.AFTER_MESSAGE_DENORMALIZED) - .forEach((listener) => listener.on(message)); + .forEach((listener) => listener.hook(message)); } async handleBeforeMessageHandler(message: RoutingMessage): Promise { await this.messagingHookRegistry .getAllByHook(LifecycleHook.BEFORE_MESSAGE_HANDLER) - .forEach((listener) => listener.on(message)); + .forEach((listener) => listener.hook(message)); } async handleAfterMessageHandlerExecuted(message: RoutingMessage): Promise { await this.messagingHookRegistry .getAllByHook(LifecycleHook.AFTER_MESSAGE_HANDLER_EXECUTED) - .forEach((listener) => listener.on(message)); + .forEach((listener) => listener.hook(message)); + } + + async handleOnFailedMessageConsumer(message: ConsumerMessage): Promise { + await this.messagingHookRegistry + .getAllByHook(LifecycleHook.ON_FAILED_MESSAGE_CONSUMER) + .forEach((listener) => listener.hook(message)); } } diff --git a/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-listener.ts b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-listener.ts index 532bdb7..0f97010 100644 --- a/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-listener.ts +++ b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-listener.ts @@ -1,13 +1,37 @@ import { RoutingMessage } from '../message/routing-message'; +import { ConsumerMessage } from '../consumer/consumer-message'; -export type MessagingData = RoutingMessage; +export type MessagingData = RoutingMessage | ConsumerMessage | DetailedConsumerMessage; + +export class DetailedConsumerMessage extends ConsumerMessage { + constructor( + public readonly message: object | string, + public readonly routingKey: string, + public readonly metadata: Record = {}, + public readonly channelName: string, + public readonly channelType: string, + ) { + super(message, routingKey, metadata); + } + + static fromConsumerMessage(consumerMessage: ConsumerMessage, channelName: string, channelType: string): DetailedConsumerMessage { + return new DetailedConsumerMessage( + consumerMessage.message, + consumerMessage.routingKey, + consumerMessage.metadata, + channelName, + channelType, + ); + } +} export enum LifecycleHook { AFTER_MESSAGE_DENORMALIZED = 'AFTER_MESSAGE_DENORMALIZED', BEFORE_MESSAGE_HANDLER = 'BEFORE_MESSAGE_HANDLER', AFTER_MESSAGE_HANDLER_EXECUTED = 'AFTER_MESSAGE_HANDLER_EXECUTED', + ON_FAILED_MESSAGE_CONSUMER = 'ON_FAILED_MESSAGE_CONSUMER', } export interface MessagingLifecycleHookListener { - on(data: MessagingData): Promise; + hook(data: MessagingData): Promise; } diff --git a/packages/messaging/test/unit/lifecycle-hook/messaging-hook-handler.spec.ts b/packages/messaging/test/unit/lifecycle-hook/messaging-hook-handler.spec.ts index bf8bba4..6ea32ad 100644 --- a/packages/messaging/test/unit/lifecycle-hook/messaging-hook-handler.spec.ts +++ b/packages/messaging/test/unit/lifecycle-hook/messaging-hook-handler.spec.ts @@ -18,7 +18,7 @@ describe('MessagingHookHandler', () => { } as unknown as MessagingLifecycleHookRegistry; handler = new MessagingLifecycleHookHandler(hookRegistry); - listener = { on: jest.fn().mockResolvedValue(undefined) }; + listener = { hook: jest.fn().mockResolvedValue(undefined) }; }); test('should execute listeners for AFTER_MESSAGE_DENORMALIZED hook', async () => { @@ -29,7 +29,7 @@ describe('MessagingHookHandler', () => { expect(hookRegistry.getAllByHook).toHaveBeenCalledWith( LifecycleHook.AFTER_MESSAGE_DENORMALIZED, ); - expect(listener.on).toHaveBeenCalledWith(message); + expect(listener.hook).toHaveBeenCalledWith(message); }); test('should execute listeners for BEFORE_MESSAGE_HANDLER hook', async () => { @@ -40,7 +40,7 @@ describe('MessagingHookHandler', () => { expect(hookRegistry.getAllByHook).toHaveBeenCalledWith( LifecycleHook.BEFORE_MESSAGE_HANDLER, ); - expect(listener.on).toHaveBeenCalledWith(message); + expect(listener.hook).toHaveBeenCalledWith(message); }); test('should execute listeners for AFTER_MESSAGE_HANDLER_EXECUTED hook', async () => { @@ -51,6 +51,6 @@ describe('MessagingHookHandler', () => { expect(hookRegistry.getAllByHook).toHaveBeenCalledWith( LifecycleHook.AFTER_MESSAGE_HANDLER_EXECUTED, ); - expect(listener.on).toHaveBeenCalledWith(message); + expect(listener.hook).toHaveBeenCalledWith(message); }); }); From f773ccfbbddc4ed08527998bc030864f73ad0e51 Mon Sep 17 00:00:00 2001 From: sebastian Date: Tue, 17 Mar 2026 21:39:01 +0100 Subject: [PATCH 06/16] [messaging] feat: lifecycle hooks --- packages/messaging/src/bus/in-memory-message.bus.ts | 1 - 1 file changed, 1 deletion(-) diff --git a/packages/messaging/src/bus/in-memory-message.bus.ts b/packages/messaging/src/bus/in-memory-message.bus.ts index 3d2836b..ae8f23b 100644 --- a/packages/messaging/src/bus/in-memory-message.bus.ts +++ b/packages/messaging/src/bus/in-memory-message.bus.ts @@ -27,7 +27,6 @@ export class InMemoryMessageBus implements IMessageBus { async dispatch(message: Message): Promise { const middlewares = []; - // Execution order: channel middlewares -> message middlewares -> handler middleware. middlewares.push( ...(this.channel.config?.middlewares ?? []), ...(message.messageOptions?.middlewares ?? []), From f63954a82130c8ce3164e0d114280ec083a727f6 Mon Sep 17 00:00:00 2001 From: sebastian Date: Wed, 18 Mar 2026 18:50:40 +0100 Subject: [PATCH 07/16] [messaging] feat: lifecycle hooks --- packages/messaging/package.json | 2 +- .../messaging/src/bus/consumer.message-bus.ts | 3 +- .../src/bus/distributed-message.bus.ts | 23 +++++++++++ .../src/bus/in-memory-message-bus.factory.ts | 3 +- .../src/bus/in-memory-message.bus.ts | 3 +- .../src/consumer/distributed.consumer.ts | 3 +- .../src/dependency-injection/decorator.ts | 7 +++- .../src/dependency-injection/register.ts | 20 ++++------ packages/messaging/src/index.ts | 2 +- .../messaging-lifecycle-hook-handler.ts | 28 +++++++++++-- .../messaging-lifecycle-hook-listener.ts | 32 ++++++++++++++- .../messaging-lifecycle-hook.registry.ts | 5 ++- packages/messaging/src/messaging.module.ts | 30 +++++++------- .../unit/bus/consumer.message-bus.spec.ts | 39 ++++++++++++++++--- .../consumer/distributed.consumer.spec.ts | 16 ++++++++ .../messaging-lifecycle-hook.registry.spec.ts | 10 +++-- 16 files changed, 172 insertions(+), 54 deletions(-) diff --git a/packages/messaging/package.json b/packages/messaging/package.json index 23e7349..e0e4bd3 100644 --- a/packages/messaging/package.json +++ b/packages/messaging/package.json @@ -1,6 +1,6 @@ { "name": "@nestjstools/messaging", - "version": "4.2.0-beta.0", + "version": "4.2.0-beta.1", "description": "Simplifies asynchronous and synchronous message handling with support for buses, handlers, channels, and consumers. Build scalable, decoupled applications with ease and reliability.", "author": "Sebastian Iwanczyszyn", "license": "MIT", diff --git a/packages/messaging/src/bus/consumer.message-bus.ts b/packages/messaging/src/bus/consumer.message-bus.ts index d7ca6ca..ee879e6 100644 --- a/packages/messaging/src/bus/consumer.message-bus.ts +++ b/packages/messaging/src/bus/consumer.message-bus.ts @@ -22,8 +22,7 @@ export class ConsumerMessageBus { private readonly consumer: IMessagingConsumer, private readonly exceptionListenerHandler: ExceptionListenerHandler, private readonly messagingHookHandler: MessagingLifecycleHookHandler, - ) { - } + ) {} async dispatch(consumerMessage: ConsumerMessage): Promise { try { diff --git a/packages/messaging/src/bus/distributed-message.bus.ts b/packages/messaging/src/bus/distributed-message.bus.ts index 6dd37d1..41bbe67 100644 --- a/packages/messaging/src/bus/distributed-message.bus.ts +++ b/packages/messaging/src/bus/distributed-message.bus.ts @@ -5,12 +5,15 @@ import { MessageBusCollection } from './message-bus.collection'; import { RoutingMessage } from '../message/routing-message'; import { MessageFactory } from '../message/message.factory'; import { NormalizerRegistry } from '../normalizer/normalizer.registry'; +import { MessagingLifecycleHookHandler } from '../lifecycle-hook/messaging-lifecycle-hook-handler'; +import { MessageBusMessage } from '../lifecycle-hook/messaging-lifecycle-hook-listener'; @Injectable() export class DistributedMessageBus implements IMessageBus { constructor( private messageBusCollection: MessageBusCollection, private normalizerRegistry: NormalizerRegistry, + private messagingLifecycleHookHandler: MessagingLifecycleHookHandler, ) {} async dispatch(message: RoutingMessage): Promise { @@ -20,12 +23,32 @@ export class DistributedMessageBus implements IMessageBus { const response = []; for (const collection of this.messageBusCollection.getAll()) { + await this.messagingLifecycleHookHandler.handleBeforeMessageNormalization( + MessageBusMessage.fromMessage( + message.message, + message.messageRoutingKey, + collection.channel.config.name, + collection.channel.constructor.name, + ), + ); + const normalizedMessage = await this.normalizerRegistry .getByName(collection.channel.config.normalizer.name) .normalize(message.message, message.messageRoutingKey); + + await this.messagingLifecycleHookHandler.handleAfterMessageNormalization( + MessageBusMessage.fromMessage( + message.message, + message.messageRoutingKey, + collection.channel.config.name, + collection.channel.constructor.name, + ), + ); + const handlerResponse = await collection.messageBus.dispatch( MessageFactory.creteSealedFromMessage(normalizedMessage, message), ); + if (handlerResponse) { response.push(handlerResponse); } diff --git a/packages/messaging/src/bus/in-memory-message-bus.factory.ts b/packages/messaging/src/bus/in-memory-message-bus.factory.ts index 49fb749..b585001 100644 --- a/packages/messaging/src/bus/in-memory-message-bus.factory.ts +++ b/packages/messaging/src/bus/in-memory-message-bus.factory.ts @@ -21,8 +21,7 @@ export class InMemoryMessageBusFactory implements IMessageBusFactory { const middlewares = []; diff --git a/packages/messaging/src/consumer/distributed.consumer.ts b/packages/messaging/src/consumer/distributed.consumer.ts index 41015c7..a7309e2 100644 --- a/packages/messaging/src/consumer/distributed.consumer.ts +++ b/packages/messaging/src/consumer/distributed.consumer.ts @@ -23,8 +23,7 @@ export class DistributedConsumer { @Inject(Service.LOGGER) private readonly logger: MessagingLogger, private readonly discoveryService: DiscoveryService, private readonly messagingLifecycleHookHandler: MessagingLifecycleHookHandler, - ) { - } + ) {} async run(): Promise { for (const channel of this.channelRegistry.getAll()) { diff --git a/packages/messaging/src/dependency-injection/decorator.ts b/packages/messaging/src/dependency-injection/decorator.ts index e3d9079..e00923d 100644 --- a/packages/messaging/src/dependency-injection/decorator.ts +++ b/packages/messaging/src/dependency-injection/decorator.ts @@ -10,7 +10,8 @@ export const MESSAGING_NORMALIZER_METADATA = 'MESSAGING_NORMALIZER_METADATA'; export const MESSAGING_EXCEPTION_LISTENER_METADATA = 'MESSAGING_EXCEPTION_LISTENER_METADATA'; export const MESSAGING_MESSAGE_METADATA = 'MESSAGING_MESSAGE_METADATA'; -export const MESSAGING_LIFECYCLE_HOOK_METADATA = 'MESSAGING_LIFECYCLE_HOOK_METADATA'; +export const MESSAGING_LIFECYCLE_HOOK_METADATA = + 'MESSAGING_LIFECYCLE_HOOK_METADATA'; export const MessageHandler = (...routingKey: string[]): ClassDecorator => { return (target) => { @@ -72,7 +73,9 @@ export const MessagingExceptionListener = (): ClassDecorator => { }; }; -export const MessagingLifecycleHook = (lifecycleHook: LifecycleHook): ClassDecorator => { +export const MessagingLifecycleHook = ( + lifecycleHook: LifecycleHook, +): ClassDecorator => { return (target) => { Reflect.defineMetadata( MESSAGING_LIFECYCLE_HOOK_METADATA, diff --git a/packages/messaging/src/dependency-injection/register.ts b/packages/messaging/src/dependency-injection/register.ts index 91da24d..711abb3 100644 --- a/packages/messaging/src/dependency-injection/register.ts +++ b/packages/messaging/src/dependency-injection/register.ts @@ -4,7 +4,8 @@ import { MessagingLogger } from '../logger/messaging-logger'; import { Service } from './service'; import { MESSAGE_HANDLER_METADATA, - MESSAGING_EXCEPTION_LISTENER_METADATA, MESSAGING_LIFECYCLE_HOOK_METADATA, + MESSAGING_EXCEPTION_LISTENER_METADATA, + MESSAGING_LIFECYCLE_HOOK_METADATA, MESSAGING_MIDDLEWARE_METADATA, MESSAGING_NORMALIZER_METADATA, } from './decorator'; @@ -102,18 +103,13 @@ const register = >( const exceptions = [DEFAULT_NORMALIZER, DEFAULT_MIDDLEWARE]; const registry: Registry = moduleRef.get(registryProvider); const logger: MessagingLogger = moduleRef.get(Service.LOGGER); - const instances = discoveryService - .getProviders() - .filter((provider) => { - if (!provider.metatype) { - return false; - } + const instances = discoveryService.getProviders().filter((provider) => { + if (!provider.metatype) { + return false; + } - return Reflect.hasMetadata( - decoratorMetadata, - provider.metatype, - ); - }); + return Reflect.hasMetadata(decoratorMetadata, provider.metatype); + }); instances.forEach((provider) => { registry.register( diff --git a/packages/messaging/src/index.ts b/packages/messaging/src/index.ts index 8e1457b..54afd0d 100644 --- a/packages/messaging/src/index.ts +++ b/packages/messaging/src/index.ts @@ -31,4 +31,4 @@ export * from './exception/handlers.exception'; export * from './logger/log'; export * from './logger/messaging-logger'; export * from './logger/nest-logger'; -export * from './lifecycle-hook/messaging-lifecycle-hook-listener' +export * from './lifecycle-hook/messaging-lifecycle-hook-listener'; diff --git a/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-handler.ts b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-handler.ts index 0b7378d..162498e 100644 --- a/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-handler.ts +++ b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-handler.ts @@ -1,6 +1,9 @@ import { Injectable } from '@nestjs/common'; import { MessagingLifecycleHookRegistry } from './messaging-lifecycle-hook.registry'; -import { LifecycleHook } from './messaging-lifecycle-hook-listener'; +import { + LifecycleHook, + MessageBusMessage, +} from './messaging-lifecycle-hook-listener'; import { RoutingMessage } from '../message/routing-message'; import { ConsumerMessage } from '../consumer/consumer-message'; @@ -8,8 +11,7 @@ import { ConsumerMessage } from '../consumer/consumer-message'; export class MessagingLifecycleHookHandler { constructor( private readonly messagingHookRegistry: MessagingLifecycleHookRegistry, - ) { - } + ) {} async handleAfterMessageDenormalized(message: RoutingMessage): Promise { await this.messagingHookRegistry @@ -23,7 +25,9 @@ export class MessagingLifecycleHookHandler { .forEach((listener) => listener.hook(message)); } - async handleAfterMessageHandlerExecuted(message: RoutingMessage): Promise { + async handleAfterMessageHandlerExecuted( + message: RoutingMessage, + ): Promise { await this.messagingHookRegistry .getAllByHook(LifecycleHook.AFTER_MESSAGE_HANDLER_EXECUTED) .forEach((listener) => listener.hook(message)); @@ -34,4 +38,20 @@ export class MessagingLifecycleHookHandler { .getAllByHook(LifecycleHook.ON_FAILED_MESSAGE_CONSUMER) .forEach((listener) => listener.hook(message)); } + + async handleBeforeMessageNormalization( + message: MessageBusMessage, + ): Promise { + await this.messagingHookRegistry + .getAllByHook(LifecycleHook.BEFORE_MESSAGE_NORMALIZATION) + .forEach((listener) => listener.hook(message)); + } + + async handleAfterMessageNormalization( + message: MessageBusMessage, + ): Promise { + await this.messagingHookRegistry + .getAllByHook(LifecycleHook.AFTER_MESSAGE_NORMALIZATION) + .forEach((listener) => listener.hook(message)); + } } diff --git a/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-listener.ts b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-listener.ts index 0f97010..c7149dc 100644 --- a/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-listener.ts +++ b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-listener.ts @@ -1,7 +1,11 @@ import { RoutingMessage } from '../message/routing-message'; import { ConsumerMessage } from '../consumer/consumer-message'; -export type MessagingData = RoutingMessage | ConsumerMessage | DetailedConsumerMessage; +export type MessagingData = + | RoutingMessage + | ConsumerMessage + | DetailedConsumerMessage + | MessageBusMessage; export class DetailedConsumerMessage extends ConsumerMessage { constructor( @@ -14,7 +18,11 @@ export class DetailedConsumerMessage extends ConsumerMessage { super(message, routingKey, metadata); } - static fromConsumerMessage(consumerMessage: ConsumerMessage, channelName: string, channelType: string): DetailedConsumerMessage { + static fromConsumerMessage( + consumerMessage: ConsumerMessage, + channelName: string, + channelType: string, + ): DetailedConsumerMessage { return new DetailedConsumerMessage( consumerMessage.message, consumerMessage.routingKey, @@ -25,7 +33,27 @@ export class DetailedConsumerMessage extends ConsumerMessage { } } +export class MessageBusMessage { + constructor( + public readonly message: T, + public readonly routingKey: string, + public readonly channelName: string, + public readonly channelType: string, + ) {} + + static fromMessage( + message: T, + routingKey: string, + channelName: string, + channelType: string, + ): MessageBusMessage { + return new MessageBusMessage(message, routingKey, channelName, channelType); + } +} + export enum LifecycleHook { + BEFORE_MESSAGE_NORMALIZATION = 'BEFORE_MESSAGE_NORMALIZATION', + AFTER_MESSAGE_NORMALIZATION = 'AFTER_MESSAGE_NORMALIZATION', AFTER_MESSAGE_DENORMALIZED = 'AFTER_MESSAGE_DENORMALIZED', BEFORE_MESSAGE_HANDLER = 'BEFORE_MESSAGE_HANDLER', AFTER_MESSAGE_HANDLER_EXECUTED = 'AFTER_MESSAGE_HANDLER_EXECUTED', diff --git a/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook.registry.ts b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook.registry.ts index 4a9cd71..b2e60b9 100644 --- a/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook.registry.ts +++ b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook.registry.ts @@ -1,5 +1,8 @@ import { BaseRegistry } from '../shared/base-registry'; -import { MessagingLifecycleHookListener, LifecycleHook } from './messaging-lifecycle-hook-listener'; +import { + MessagingLifecycleHookListener, + LifecycleHook, +} from './messaging-lifecycle-hook-listener'; export class MessagingLifecycleHookRegistry extends BaseRegistry { getAllByHook(hook: LifecycleHook): MessagingLifecycleHookListener[] { diff --git a/packages/messaging/src/messaging.module.ts b/packages/messaging/src/messaging.module.ts index 63f8da6..e5eb491 100644 --- a/packages/messaging/src/messaging.module.ts +++ b/packages/messaging/src/messaging.module.ts @@ -30,7 +30,8 @@ import { InMemoryChannelFactory } from './channel/factory/in-memory-channel.fact import { DistributedConsumer } from './consumer/distributed.consumer'; import { registerExceptionListener, - registerHandlers, registerMessagingHooks, + registerHandlers, + registerMessagingHooks, registerMessageNormalizers, registerMiddlewares, } from './dependency-injection/register'; @@ -49,7 +50,8 @@ import { MessagingLifecycleHookRegistry } from './lifecycle-hook/messaging-lifec @Module({}) export class MessagingModule - implements OnApplicationBootstrap, OnModuleDestroy { + implements OnApplicationBootstrap, OnModuleDestroy +{ static forRoot(options: MessagingModuleOptions): DynamicModule { const channels = options.channels ?? []; @@ -111,6 +113,7 @@ export class MessagingModule busFactory: CompositeMessageBusFactory, logger: MessagingLogger, normalizerRegistry: NormalizerRegistry, + messagingLifecycleHookHandler: MessagingLifecycleHookHandler, ) => { const messageBusCollection = new MessageBusCollection(); @@ -125,6 +128,7 @@ export class MessagingModule const messageBus = new DistributedMessageBus( messageBusCollection, normalizerRegistry, + messagingLifecycleHookHandler, ); logger.log(`MessageBus [${bus.name}] was created successfully`); @@ -136,6 +140,7 @@ export class MessagingModule CompositeMessageBusFactory, Service.LOGGER, Service.MESSAGE_NORMALIZERS_REGISTRY, + MessagingLifecycleHookHandler, ], })); }; @@ -176,15 +181,15 @@ export class MessagingModule options.customLogger && typeof options.customLogger === 'function' ? { provide: Service.LOGGER, useClass: options.customLogger } : { - provide: Service.LOGGER, - useValue: - options.customLogger ?? - new NestLogger( - new NestCommonLogger(), - options.debug ?? false, - options.logging ?? true, - ), - } + provide: Service.LOGGER, + useValue: + options.customLogger ?? + new NestLogger( + new NestCommonLogger(), + options.debug ?? false, + options.logging ?? true, + ), + } ) as Provider; return { @@ -253,8 +258,7 @@ export class MessagingModule private readonly configuration: MandatoryMessagingModuleOptions, @Inject(Service.LOGGER) private readonly logger: MessagingLogger, - ) { - } + ) {} onApplicationBootstrap(): void { registerHandlers(this.moduleRef, this.discoveryService); diff --git a/packages/messaging/test/unit/bus/consumer.message-bus.spec.ts b/packages/messaging/test/unit/bus/consumer.message-bus.spec.ts index 48d3234..98a1997 100644 --- a/packages/messaging/test/unit/bus/consumer.message-bus.spec.ts +++ b/packages/messaging/test/unit/bus/consumer.message-bus.spec.ts @@ -1,4 +1,3 @@ -import { Logger } from '@nestjs/common'; import { ConsumerMessageBus } from '../../../src'; import { IMessageBus } from '../../../src'; import { TestChannel } from '../../support/test.channel'; @@ -13,26 +12,47 @@ import { ObjectForwardMessageNormalizer } from '../../../src/normalizer/object-f import { ConsumerDispatchedMessageError } from '../../../src'; import { HandlerError, HandlersException } from '../../../src'; import { ExceptionContext } from '../../../src'; +import { MessagingLifecycleHookHandler } from '../../../src/lifecycle-hook/messaging-lifecycle-hook-handler'; +import { Logger } from '@nestjs/common'; describe('ConsumerMessageBus', () => { let messageBus: IMessageBus; + let messageBusDispatchMock: jest.Mock; let logger: SpyLogger; let consumer: IMessagingConsumer; + let consumerOnErrorMock: jest.Mock; let exceptionListenerHandler: ExceptionListenerHandler; + let exceptionHandlerMock: jest.Mock; + let messagingLifecycleHookHandler: MessagingLifecycleHookHandler; + let failedConsumerHookMock: jest.Mock; let channel: TestChannel; beforeEach(() => { + jest.clearAllMocks(); + + messageBusDispatchMock = jest.fn().mockResolvedValue(undefined); messageBus = { - dispatch: jest.fn().mockResolvedValue(undefined), + dispatch: messageBusDispatchMock, } as unknown as IMessageBus; logger = new SpyLogger(new Logger(), false, false); + consumerOnErrorMock = jest.fn().mockResolvedValue(undefined); consumer = { consume: jest.fn(), - onError: jest.fn().mockResolvedValue(undefined), + onError: consumerOnErrorMock, } as unknown as IMessagingConsumer; + exceptionHandlerMock = jest.fn().mockResolvedValue(undefined); exceptionListenerHandler = { - handleError: jest.fn().mockResolvedValue(undefined), + handleError: exceptionHandlerMock, } as unknown as ExceptionListenerHandler; + failedConsumerHookMock = jest.fn().mockResolvedValue(undefined); + messagingLifecycleHookHandler = { + handleAfterMessageDenormalized: jest.fn().mockResolvedValue(undefined), + handleBeforeMessageHandler: jest.fn().mockResolvedValue(undefined), + handleAfterMessageHandlerExecuted: jest.fn().mockResolvedValue(undefined), + handleOnFailedMessageConsumer: failedConsumerHookMock, + handleBeforeMessageNormalization: jest.fn().mockResolvedValue(undefined), + handleAfterMessageNormalization: jest.fn().mockResolvedValue(undefined), + } as unknown as MessagingLifecycleHookHandler; channel = new TestChannel(new InMemoryChannelConfig({ name: 'ds' })); }); @@ -43,6 +63,7 @@ describe('ConsumerMessageBus', () => { logger, consumer, exceptionListenerHandler, + messagingLifecycleHookHandler, ); await subjectUnderTest.dispatch( @@ -72,8 +93,9 @@ describe('ConsumerMessageBus', () => { it('should call onError and exception listener when message bus dispatch throws regular error', async () => { const error = new Error('boom'); + messageBusDispatchMock = jest.fn().mockRejectedValue(error); messageBus = { - dispatch: jest.fn().mockRejectedValue(error), + dispatch: messageBusDispatchMock, } as unknown as IMessageBus; const subjectUnderTest = new ConsumerMessageBus( @@ -82,6 +104,7 @@ describe('ConsumerMessageBus', () => { logger, consumer, exceptionListenerHandler, + messagingLifecycleHookHandler, ); const consumerMessage = new ConsumerMessage({ status: 'fail' }, 'rk.fail'); @@ -111,6 +134,7 @@ describe('ConsumerMessageBus', () => { }, }, }); + expect(failedConsumerHookMock).toHaveBeenCalledTimes(1); }); it('should not log error when dispatch throws HandlersException', async () => { @@ -118,8 +142,9 @@ describe('ConsumerMessageBus', () => { new HandlerError('MyHandler', new Error('handler failed')), ]); + messageBusDispatchMock = jest.fn().mockRejectedValue(handlersError); messageBus = { - dispatch: jest.fn().mockRejectedValue(handlersError), + dispatch: messageBusDispatchMock, } as unknown as IMessageBus; const subjectUnderTest = new ConsumerMessageBus( @@ -128,6 +153,7 @@ describe('ConsumerMessageBus', () => { logger, consumer, exceptionListenerHandler, + messagingLifecycleHookHandler, ); await subjectUnderTest.dispatch( @@ -144,5 +170,6 @@ describe('ConsumerMessageBus', () => { ), ); expect(logger.getLogs().find((log) => log.type === 'ERROR')).toBeFalsy(); + expect(failedConsumerHookMock).toHaveBeenCalledTimes(1); }); }); diff --git a/packages/messaging/test/unit/consumer/distributed.consumer.spec.ts b/packages/messaging/test/unit/consumer/distributed.consumer.spec.ts index 1f2bdea..e226ba4 100644 --- a/packages/messaging/test/unit/consumer/distributed.consumer.spec.ts +++ b/packages/messaging/test/unit/consumer/distributed.consumer.spec.ts @@ -11,6 +11,7 @@ import { TestChannel } from '../../support/test.channel'; import { IMessageBus } from '../../../src'; import { ExceptionListenerHandler } from '../../../src/exception-listener/exception-listener-handler'; import { ConsumerMessageBus } from '../../../src'; +import { MessagingLifecycleHookHandler } from '../../../src/lifecycle-hook/messaging-lifecycle-hook-handler'; describe('DistributedConsumer', () => { let subjectUnderTest: DistributedConsumer; @@ -20,6 +21,7 @@ describe('DistributedConsumer', () => { let discoveryService: DiscoveryService; let consumeMock: jest.Mock; let onErrorMock: jest.Mock; + let messagingLifecycleHookHandler: MessagingLifecycleHookHandler; beforeEach(() => { logger = new SpyLogger(new Logger(), false, false); @@ -30,6 +32,15 @@ describe('DistributedConsumer', () => { dispatch: jest.fn(), } as unknown as IMessageBus; + messagingLifecycleHookHandler = { + handleAfterMessageDenormalized: jest.fn().mockResolvedValue(undefined), + handleBeforeMessageHandler: jest.fn().mockResolvedValue(undefined), + handleAfterMessageHandlerExecuted: jest.fn().mockResolvedValue(undefined), + handleOnFailedMessageConsumer: jest.fn().mockResolvedValue(undefined), + handleBeforeMessageNormalization: jest.fn().mockResolvedValue(undefined), + handleAfterMessageNormalization: jest.fn().mockResolvedValue(undefined), + } as unknown as MessagingLifecycleHookHandler; + consumeMock = jest.fn().mockResolvedValue(undefined); onErrorMock = jest.fn().mockResolvedValue(undefined); @@ -61,6 +72,7 @@ describe('DistributedConsumer', () => { exceptionListenerHandler, logger, discoveryService, + messagingLifecycleHookHandler, ); await subjectUnderTest.run(); @@ -92,6 +104,7 @@ describe('DistributedConsumer', () => { exceptionListenerHandler, logger, discoveryService, + messagingLifecycleHookHandler, ); await subjectUnderTest.run(); @@ -122,6 +135,7 @@ describe('DistributedConsumer', () => { exceptionListenerHandler, logger, discoveryService, + messagingLifecycleHookHandler, ); await subjectUnderTest.run(); @@ -145,6 +159,7 @@ describe('DistributedConsumer', () => { exceptionListenerHandler, logger, discoveryService, + messagingLifecycleHookHandler, ); await expect(subjectUnderTest.run()).rejects.toThrow( @@ -184,6 +199,7 @@ describe('DistributedConsumer', () => { exceptionListenerHandler, logger, discoveryService, + messagingLifecycleHookHandler, ); await expect(subjectUnderTest.run()).rejects.toThrow( diff --git a/packages/messaging/test/unit/lifecycle-hook/messaging-lifecycle-hook.registry.spec.ts b/packages/messaging/test/unit/lifecycle-hook/messaging-lifecycle-hook.registry.spec.ts index 7d302e7..5efa38a 100644 --- a/packages/messaging/test/unit/lifecycle-hook/messaging-lifecycle-hook.registry.spec.ts +++ b/packages/messaging/test/unit/lifecycle-hook/messaging-lifecycle-hook.registry.spec.ts @@ -11,7 +11,9 @@ describe('MessagingLifecycleHookRegistry', () => { beforeEach(() => { registry = new MessagingLifecycleHookRegistry(); - beforeListener = { on: jest.fn() } as unknown as MessagingLifecycleHookListener; + beforeListener = { + on: jest.fn(), + } as unknown as MessagingLifecycleHookListener; afterDenormalizedListener = { on: jest.fn(), } as unknown as MessagingLifecycleHookListener; @@ -27,9 +29,9 @@ describe('MessagingLifecycleHookRegistry', () => { afterDenormalizedListener, ); - expect(registry.getAllByHook(LifecycleHook.BEFORE_MESSAGE_HANDLER)).toEqual([ - beforeListener, - ]); + expect(registry.getAllByHook(LifecycleHook.BEFORE_MESSAGE_HANDLER)).toEqual( + [beforeListener], + ); }); test('should return empty array when no listeners are registered for given lifecycle hook', () => { From b251e02e44d727c9ca53fa4981c4ecb039261528 Mon Sep 17 00:00:00 2001 From: sebastian Date: Wed, 18 Mar 2026 23:05:24 +0100 Subject: [PATCH 08/16] [messaging] feat: lifecycle hooks --- README.md | 2 + .../messaging/src/bus/consumer.message-bus.ts | 4 +- .../src/bus/distributed-message.bus.ts | 6 +- .../src/bus/in-memory-message.bus.ts | 19 ++++++- .../messaging-lifecycle-hook-handler.ts | 24 +++----- .../messaging-lifecycle-hook-listener.ts | 56 +++++++++---------- .../messaging-hook-handler.spec.ts | 30 ++++++---- .../messaging-lifecycle-hook.registry.spec.ts | 2 +- 8 files changed, 75 insertions(+), 68 deletions(-) diff --git a/README.md b/README.md index 1098779..da99c73 100644 --- a/README.md +++ b/README.md @@ -219,6 +219,8 @@ flowchart LR MW --> H ``` +--- + ### Exception handling flow: ```mermaid diff --git a/packages/messaging/src/bus/consumer.message-bus.ts b/packages/messaging/src/bus/consumer.message-bus.ts index ee879e6..f76036c 100644 --- a/packages/messaging/src/bus/consumer.message-bus.ts +++ b/packages/messaging/src/bus/consumer.message-bus.ts @@ -12,7 +12,7 @@ import { HandlersException } from '../exception/handlers.exception'; import { ExceptionListenerHandler } from '../exception-listener/exception-listener-handler'; import { ExceptionContext } from '../exception-listener/exception-context'; import { MessagingLifecycleHookHandler } from '../lifecycle-hook/messaging-lifecycle-hook-handler'; -import { DetailedConsumerMessage } from '../lifecycle-hook/messaging-lifecycle-hook-listener'; +import { HookMessage } from '../lifecycle-hook/messaging-lifecycle-hook-listener'; export class ConsumerMessageBus { constructor( @@ -79,7 +79,7 @@ export class ConsumerMessageBus { ); await this.messagingHookHandler.handleOnFailedMessageConsumer( - DetailedConsumerMessage.fromConsumerMessage( + HookMessage.fromConsumerMessage( consumerMessage, this.channel.config.name, this.channel.constructor.name, diff --git a/packages/messaging/src/bus/distributed-message.bus.ts b/packages/messaging/src/bus/distributed-message.bus.ts index 41bbe67..1e257dd 100644 --- a/packages/messaging/src/bus/distributed-message.bus.ts +++ b/packages/messaging/src/bus/distributed-message.bus.ts @@ -6,7 +6,7 @@ import { RoutingMessage } from '../message/routing-message'; import { MessageFactory } from '../message/message.factory'; import { NormalizerRegistry } from '../normalizer/normalizer.registry'; import { MessagingLifecycleHookHandler } from '../lifecycle-hook/messaging-lifecycle-hook-handler'; -import { MessageBusMessage } from '../lifecycle-hook/messaging-lifecycle-hook-listener'; +import { HookMessage } from '../lifecycle-hook/messaging-lifecycle-hook-listener'; @Injectable() export class DistributedMessageBus implements IMessageBus { @@ -24,7 +24,7 @@ export class DistributedMessageBus implements IMessageBus { const response = []; for (const collection of this.messageBusCollection.getAll()) { await this.messagingLifecycleHookHandler.handleBeforeMessageNormalization( - MessageBusMessage.fromMessage( + HookMessage.fromMessage( message.message, message.messageRoutingKey, collection.channel.config.name, @@ -37,7 +37,7 @@ export class DistributedMessageBus implements IMessageBus { .normalize(message.message, message.messageRoutingKey); await this.messagingLifecycleHookHandler.handleAfterMessageNormalization( - MessageBusMessage.fromMessage( + HookMessage.fromMessage( message.message, message.messageRoutingKey, collection.channel.config.name, diff --git a/packages/messaging/src/bus/in-memory-message.bus.ts b/packages/messaging/src/bus/in-memory-message.bus.ts index af3961f..257c34d 100644 --- a/packages/messaging/src/bus/in-memory-message.bus.ts +++ b/packages/messaging/src/bus/in-memory-message.bus.ts @@ -14,6 +14,7 @@ import { RoutingMessage } from '../message/routing-message'; import { NormalizerRegistry } from '../normalizer/normalizer.registry'; import { DefaultMessageOptions } from '../message/default-message-options'; import { MessagingLifecycleHookHandler } from '../lifecycle-hook/messaging-lifecycle-hook-handler'; +import { HookMessage } from '../lifecycle-hook/messaging-lifecycle-hook-listener'; export class InMemoryMessageBus implements IMessageBus { constructor( @@ -49,7 +50,11 @@ export class InMemoryMessageBus implements IMessageBus { // Hook fired once the payload shape is ready for handler pipeline. await this.messagingHookHandler.handleAfterMessageDenormalized( - MessageFactory.creteRoutingFromMessage(messageToDispatch, message), + HookMessage.fromRoutingMessage( + MessageFactory.creteRoutingFromMessage(messageToDispatch, message), + this.channel.config.name, + this.channel.constructor.name, + ), ); try { @@ -88,7 +93,11 @@ export class InMemoryMessageBus implements IMessageBus { // Hook around handler execution. await this.messagingHookHandler.handleBeforeMessageHandler( - MessageFactory.creteRoutingFromMessage(messageToDispatch, message), + HookMessage.fromRoutingMessage( + MessageFactory.creteRoutingFromMessage(messageToDispatch, message), + this.channel.config.name, + this.channel.constructor.name, + ), ); const response = await middlewareInstances[0].process( @@ -97,7 +106,11 @@ export class InMemoryMessageBus implements IMessageBus { ); await this.messagingHookHandler.handleAfterMessageHandlerExecuted( - MessageFactory.creteRoutingFromMessage(messageToDispatch, message), + HookMessage.fromRoutingMessage( + MessageFactory.creteRoutingFromMessage(messageToDispatch, message), + this.channel.config.name, + this.channel.constructor.name, + ), ); return Promise.resolve(response); diff --git a/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-handler.ts b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-handler.ts index 162498e..91da0cf 100644 --- a/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-handler.ts +++ b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-handler.ts @@ -2,10 +2,8 @@ import { Injectable } from '@nestjs/common'; import { MessagingLifecycleHookRegistry } from './messaging-lifecycle-hook.registry'; import { LifecycleHook, - MessageBusMessage, + HookMessage, } from './messaging-lifecycle-hook-listener'; -import { RoutingMessage } from '../message/routing-message'; -import { ConsumerMessage } from '../consumer/consumer-message'; @Injectable() export class MessagingLifecycleHookHandler { @@ -13,43 +11,37 @@ export class MessagingLifecycleHookHandler { private readonly messagingHookRegistry: MessagingLifecycleHookRegistry, ) {} - async handleAfterMessageDenormalized(message: RoutingMessage): Promise { + async handleAfterMessageDenormalized(message: HookMessage): Promise { await this.messagingHookRegistry .getAllByHook(LifecycleHook.AFTER_MESSAGE_DENORMALIZED) .forEach((listener) => listener.hook(message)); } - async handleBeforeMessageHandler(message: RoutingMessage): Promise { + async handleBeforeMessageHandler(message: HookMessage): Promise { await this.messagingHookRegistry .getAllByHook(LifecycleHook.BEFORE_MESSAGE_HANDLER) .forEach((listener) => listener.hook(message)); } - async handleAfterMessageHandlerExecuted( - message: RoutingMessage, - ): Promise { + async handleAfterMessageHandlerExecuted(message: HookMessage): Promise { await this.messagingHookRegistry - .getAllByHook(LifecycleHook.AFTER_MESSAGE_HANDLER_EXECUTED) + .getAllByHook(LifecycleHook.AFTER_MESSAGE_HANDLER_EXECUTION) .forEach((listener) => listener.hook(message)); } - async handleOnFailedMessageConsumer(message: ConsumerMessage): Promise { + async handleOnFailedMessageConsumer(message: HookMessage): Promise { await this.messagingHookRegistry .getAllByHook(LifecycleHook.ON_FAILED_MESSAGE_CONSUMER) .forEach((listener) => listener.hook(message)); } - async handleBeforeMessageNormalization( - message: MessageBusMessage, - ): Promise { + async handleBeforeMessageNormalization(message: HookMessage): Promise { await this.messagingHookRegistry .getAllByHook(LifecycleHook.BEFORE_MESSAGE_NORMALIZATION) .forEach((listener) => listener.hook(message)); } - async handleAfterMessageNormalization( - message: MessageBusMessage, - ): Promise { + async handleAfterMessageNormalization(message: HookMessage): Promise { await this.messagingHookRegistry .getAllByHook(LifecycleHook.AFTER_MESSAGE_NORMALIZATION) .forEach((listener) => listener.hook(message)); diff --git a/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-listener.ts b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-listener.ts index c7149dc..431b1dc 100644 --- a/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-listener.ts +++ b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-listener.ts @@ -1,53 +1,47 @@ -import { RoutingMessage } from '../message/routing-message'; import { ConsumerMessage } from '../consumer/consumer-message'; +import { RoutingMessage } from '../message/routing-message'; -export type MessagingData = - | RoutingMessage - | ConsumerMessage - | DetailedConsumerMessage - | MessageBusMessage; - -export class DetailedConsumerMessage extends ConsumerMessage { +export class HookMessage { constructor( - public readonly message: object | string, + public readonly message: T, public readonly routingKey: string, - public readonly metadata: Record = {}, public readonly channelName: string, public readonly channelType: string, - ) { - super(message, routingKey, metadata); + ) {} + + static fromMessage( + message: T, + routingKey: string, + channelName: string, + channelType: string, + ): HookMessage { + return new HookMessage(message, routingKey, channelName, channelType); } static fromConsumerMessage( consumerMessage: ConsumerMessage, channelName: string, channelType: string, - ): DetailedConsumerMessage { - return new DetailedConsumerMessage( + ): HookMessage { + return new HookMessage( consumerMessage.message, consumerMessage.routingKey, - consumerMessage.metadata, channelName, channelType, ); } -} -export class MessageBusMessage { - constructor( - public readonly message: T, - public readonly routingKey: string, - public readonly channelName: string, - public readonly channelType: string, - ) {} - - static fromMessage( - message: T, - routingKey: string, + static fromRoutingMessage( + routingMessage: RoutingMessage, channelName: string, channelType: string, - ): MessageBusMessage { - return new MessageBusMessage(message, routingKey, channelName, channelType); + ): HookMessage { + return new HookMessage( + routingMessage.message, + routingMessage.messageRoutingKey, + channelName, + channelType, + ); } } @@ -56,10 +50,10 @@ export enum LifecycleHook { AFTER_MESSAGE_NORMALIZATION = 'AFTER_MESSAGE_NORMALIZATION', AFTER_MESSAGE_DENORMALIZED = 'AFTER_MESSAGE_DENORMALIZED', BEFORE_MESSAGE_HANDLER = 'BEFORE_MESSAGE_HANDLER', - AFTER_MESSAGE_HANDLER_EXECUTED = 'AFTER_MESSAGE_HANDLER_EXECUTED', + AFTER_MESSAGE_HANDLER_EXECUTION = 'AFTER_MESSAGE_HANDLER_EXECUTION', ON_FAILED_MESSAGE_CONSUMER = 'ON_FAILED_MESSAGE_CONSUMER', } export interface MessagingLifecycleHookListener { - hook(data: MessagingData): Promise; + hook(message: HookMessage): Promise; } diff --git a/packages/messaging/test/unit/lifecycle-hook/messaging-hook-handler.spec.ts b/packages/messaging/test/unit/lifecycle-hook/messaging-hook-handler.spec.ts index 6ea32ad..d6d3e5f 100644 --- a/packages/messaging/test/unit/lifecycle-hook/messaging-hook-handler.spec.ts +++ b/packages/messaging/test/unit/lifecycle-hook/messaging-hook-handler.spec.ts @@ -1,13 +1,13 @@ import { MessagingLifecycleHookHandler } from '../../../src/lifecycle-hook/messaging-lifecycle-hook-handler'; import { MessagingLifecycleHookRegistry } from '../../../src/lifecycle-hook/messaging-lifecycle-hook.registry'; import { + HookMessage, LifecycleHook, - MessagingLifecycleHookListener, -} from '../../../src/lifecycle-hook/messaging-lifecycle-hook-listener'; -import { RoutingMessage } from '../../../src/message/routing-message'; + MessagingLifecycleHookListener, RoutingMessage, +} from '../../../src'; describe('MessagingHookHandler', () => { - let hookRegistry: MessagingLifecycleHookRegistry; + let hookRegistry: jest.Mocked; let handler: MessagingLifecycleHookHandler; let listener: MessagingLifecycleHookListener; const message = new RoutingMessage({ title: 'hello' }, 'test.key'); @@ -15,7 +15,7 @@ describe('MessagingHookHandler', () => { beforeEach(() => { hookRegistry = { getAllByHook: jest.fn(), - } as unknown as MessagingLifecycleHookRegistry; + } as unknown as jest.Mocked; handler = new MessagingLifecycleHookHandler(hookRegistry); listener = { hook: jest.fn().mockResolvedValue(undefined) }; @@ -24,33 +24,39 @@ describe('MessagingHookHandler', () => { test('should execute listeners for AFTER_MESSAGE_DENORMALIZED hook', async () => { (hookRegistry.getAllByHook as jest.Mock).mockReturnValue([listener]); - await handler.handleAfterMessageDenormalized(message); + const hookMessage = HookMessage.fromRoutingMessage(message, 'example', 'example'); + + await handler.handleAfterMessageDenormalized(hookMessage); expect(hookRegistry.getAllByHook).toHaveBeenCalledWith( LifecycleHook.AFTER_MESSAGE_DENORMALIZED, ); - expect(listener.hook).toHaveBeenCalledWith(message); + expect(listener.hook).toHaveBeenCalledWith(hookMessage); }); test('should execute listeners for BEFORE_MESSAGE_HANDLER hook', async () => { (hookRegistry.getAllByHook as jest.Mock).mockReturnValue([listener]); - await handler.handleBeforeMessageHandler(message); + const hookMessage = HookMessage.fromRoutingMessage(message, 'example', 'example'); + + await handler.handleBeforeMessageHandler(hookMessage); expect(hookRegistry.getAllByHook).toHaveBeenCalledWith( LifecycleHook.BEFORE_MESSAGE_HANDLER, ); - expect(listener.hook).toHaveBeenCalledWith(message); + expect(listener.hook).toHaveBeenCalledWith(hookMessage); }); test('should execute listeners for AFTER_MESSAGE_HANDLER_EXECUTED hook', async () => { (hookRegistry.getAllByHook as jest.Mock).mockReturnValue([listener]); - await handler.handleAfterMessageHandlerExecuted(message); + const hookMessage = HookMessage.fromRoutingMessage(message, 'example', 'example'); + + await handler.handleAfterMessageHandlerExecuted(hookMessage); expect(hookRegistry.getAllByHook).toHaveBeenCalledWith( - LifecycleHook.AFTER_MESSAGE_HANDLER_EXECUTED, + LifecycleHook.AFTER_MESSAGE_HANDLER_EXECUTION, ); - expect(listener.hook).toHaveBeenCalledWith(message); + expect(listener.hook).toHaveBeenCalledWith(hookMessage); }); }); diff --git a/packages/messaging/test/unit/lifecycle-hook/messaging-lifecycle-hook.registry.spec.ts b/packages/messaging/test/unit/lifecycle-hook/messaging-lifecycle-hook.registry.spec.ts index 5efa38a..9fdc3db 100644 --- a/packages/messaging/test/unit/lifecycle-hook/messaging-lifecycle-hook.registry.spec.ts +++ b/packages/messaging/test/unit/lifecycle-hook/messaging-lifecycle-hook.registry.spec.ts @@ -41,7 +41,7 @@ describe('MessagingLifecycleHookRegistry', () => { ); expect( - registry.getAllByHook(LifecycleHook.AFTER_MESSAGE_HANDLER_EXECUTED), + registry.getAllByHook(LifecycleHook.AFTER_MESSAGE_HANDLER_EXECUTION), ).toEqual([]); }); }); From cb074acc0782082047c8db1eb1f5d2751c742ae5 Mon Sep 17 00:00:00 2001 From: sebastian Date: Wed, 18 Mar 2026 23:05:52 +0100 Subject: [PATCH 09/16] [messaging] feat: lifecycle hooks --- .../messaging-hook-handler.spec.ts | 21 +++++++++++++++---- 1 file changed, 17 insertions(+), 4 deletions(-) diff --git a/packages/messaging/test/unit/lifecycle-hook/messaging-hook-handler.spec.ts b/packages/messaging/test/unit/lifecycle-hook/messaging-hook-handler.spec.ts index d6d3e5f..b288eba 100644 --- a/packages/messaging/test/unit/lifecycle-hook/messaging-hook-handler.spec.ts +++ b/packages/messaging/test/unit/lifecycle-hook/messaging-hook-handler.spec.ts @@ -3,7 +3,8 @@ import { MessagingLifecycleHookRegistry } from '../../../src/lifecycle-hook/mess import { HookMessage, LifecycleHook, - MessagingLifecycleHookListener, RoutingMessage, + MessagingLifecycleHookListener, + RoutingMessage, } from '../../../src'; describe('MessagingHookHandler', () => { @@ -24,7 +25,11 @@ describe('MessagingHookHandler', () => { test('should execute listeners for AFTER_MESSAGE_DENORMALIZED hook', async () => { (hookRegistry.getAllByHook as jest.Mock).mockReturnValue([listener]); - const hookMessage = HookMessage.fromRoutingMessage(message, 'example', 'example'); + const hookMessage = HookMessage.fromRoutingMessage( + message, + 'example', + 'example', + ); await handler.handleAfterMessageDenormalized(hookMessage); @@ -37,7 +42,11 @@ describe('MessagingHookHandler', () => { test('should execute listeners for BEFORE_MESSAGE_HANDLER hook', async () => { (hookRegistry.getAllByHook as jest.Mock).mockReturnValue([listener]); - const hookMessage = HookMessage.fromRoutingMessage(message, 'example', 'example'); + const hookMessage = HookMessage.fromRoutingMessage( + message, + 'example', + 'example', + ); await handler.handleBeforeMessageHandler(hookMessage); @@ -50,7 +59,11 @@ describe('MessagingHookHandler', () => { test('should execute listeners for AFTER_MESSAGE_HANDLER_EXECUTED hook', async () => { (hookRegistry.getAllByHook as jest.Mock).mockReturnValue([listener]); - const hookMessage = HookMessage.fromRoutingMessage(message, 'example', 'example'); + const hookMessage = HookMessage.fromRoutingMessage( + message, + 'example', + 'example', + ); await handler.handleAfterMessageHandlerExecuted(hookMessage); From 764578e4f293da8691a690646a7be1e76cb51789 Mon Sep 17 00:00:00 2001 From: sebastian Date: Wed, 18 Mar 2026 23:06:24 +0100 Subject: [PATCH 10/16] [messaging] feat: lifecycle hooks --- packages/messaging/package.json | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/messaging/package.json b/packages/messaging/package.json index e0e4bd3..e0373ff 100644 --- a/packages/messaging/package.json +++ b/packages/messaging/package.json @@ -1,6 +1,6 @@ { "name": "@nestjstools/messaging", - "version": "4.2.0-beta.1", + "version": "4.2.0-beta.2", "description": "Simplifies asynchronous and synchronous message handling with support for buses, handlers, channels, and consumers. Build scalable, decoupled applications with ease and reliability.", "author": "Sebastian Iwanczyszyn", "license": "MIT", From fad7cec229fa74b7195cefc415f5dfdb83c39bd0 Mon Sep 17 00:00:00 2001 From: sebastian Date: Wed, 18 Mar 2026 23:19:49 +0100 Subject: [PATCH 11/16] [messaging] feat: lifecycle hooks --- packages/messaging/package.json | 2 +- .../messaging/src/bus/consumer.message-bus.ts | 12 ++++++- .../src/bus/distributed-message.bus.ts | 13 ++++---- .../src/bus/in-memory-message.bus.ts | 9 ++--- .../messaging-lifecycle-hook-handler.ts | 6 ++++ .../messaging-lifecycle-hook-listener.ts | 33 +++++++++++-------- 6 files changed, 46 insertions(+), 29 deletions(-) diff --git a/packages/messaging/package.json b/packages/messaging/package.json index e0373ff..414fdd9 100644 --- a/packages/messaging/package.json +++ b/packages/messaging/package.json @@ -1,6 +1,6 @@ { "name": "@nestjstools/messaging", - "version": "4.2.0-beta.2", + "version": "4.2.0-beta.3", "description": "Simplifies asynchronous and synchronous message handling with support for buses, handlers, channels, and consumers. Build scalable, decoupled applications with ease and reliability.", "author": "Sebastian Iwanczyszyn", "license": "MIT", diff --git a/packages/messaging/src/bus/consumer.message-bus.ts b/packages/messaging/src/bus/consumer.message-bus.ts index f76036c..5f9479b 100644 --- a/packages/messaging/src/bus/consumer.message-bus.ts +++ b/packages/messaging/src/bus/consumer.message-bus.ts @@ -13,6 +13,7 @@ import { ExceptionListenerHandler } from '../exception-listener/exception-listen import { ExceptionContext } from '../exception-listener/exception-context'; import { MessagingLifecycleHookHandler } from '../lifecycle-hook/messaging-lifecycle-hook-handler'; import { HookMessage } from '../lifecycle-hook/messaging-lifecycle-hook-listener'; +import { MessageFactory } from '../message/message.factory'; export class ConsumerMessageBus { constructor( @@ -22,7 +23,8 @@ export class ConsumerMessageBus { private readonly consumer: IMessagingConsumer, private readonly exceptionListenerHandler: ExceptionListenerHandler, private readonly messagingHookHandler: MessagingLifecycleHookHandler, - ) {} + ) { + } async dispatch(consumerMessage: ConsumerMessage): Promise { try { @@ -49,6 +51,14 @@ export class ConsumerMessageBus { ), ); + await this.messagingHookHandler.handleAfterMessageDenormalized( + HookMessage.fromSealedRoutingMessage( + routingMessage, + this.channel.config.name, + this.channel.constructor.name, + ), + ); + await this.messageBus.dispatch(routingMessage); } catch (e) { await this.consumer.onError( diff --git a/packages/messaging/src/bus/distributed-message.bus.ts b/packages/messaging/src/bus/distributed-message.bus.ts index 1e257dd..b61e39a 100644 --- a/packages/messaging/src/bus/distributed-message.bus.ts +++ b/packages/messaging/src/bus/distributed-message.bus.ts @@ -14,7 +14,8 @@ export class DistributedMessageBus implements IMessageBus { private messageBusCollection: MessageBusCollection, private normalizerRegistry: NormalizerRegistry, private messagingLifecycleHookHandler: MessagingLifecycleHookHandler, - ) {} + ) { + } async dispatch(message: RoutingMessage): Promise { if (!(message instanceof RoutingMessage)) { @@ -24,9 +25,8 @@ export class DistributedMessageBus implements IMessageBus { const response = []; for (const collection of this.messageBusCollection.getAll()) { await this.messagingLifecycleHookHandler.handleBeforeMessageNormalization( - HookMessage.fromMessage( - message.message, - message.messageRoutingKey, + HookMessage.fromRoutingMessage( + message, collection.channel.config.name, collection.channel.constructor.name, ), @@ -37,9 +37,8 @@ export class DistributedMessageBus implements IMessageBus { .normalize(message.message, message.messageRoutingKey); await this.messagingLifecycleHookHandler.handleAfterMessageNormalization( - HookMessage.fromMessage( - message.message, - message.messageRoutingKey, + HookMessage.fromRoutingMessage( + message, collection.channel.config.name, collection.channel.constructor.name, ), diff --git a/packages/messaging/src/bus/in-memory-message.bus.ts b/packages/messaging/src/bus/in-memory-message.bus.ts index 257c34d..cb21caf 100644 --- a/packages/messaging/src/bus/in-memory-message.bus.ts +++ b/packages/messaging/src/bus/in-memory-message.bus.ts @@ -23,7 +23,8 @@ export class InMemoryMessageBus implements IMessageBus { private channel: InMemoryChannel, private normalizerRegistry: NormalizerRegistry, private messagingHookHandler: MessagingLifecycleHookHandler, - ) {} + ) { + } async dispatch(message: Message): Promise { const middlewares = []; @@ -52,8 +53,6 @@ export class InMemoryMessageBus implements IMessageBus { await this.messagingHookHandler.handleAfterMessageDenormalized( HookMessage.fromRoutingMessage( MessageFactory.creteRoutingFromMessage(messageToDispatch, message), - this.channel.config.name, - this.channel.constructor.name, ), ); @@ -95,8 +94,6 @@ export class InMemoryMessageBus implements IMessageBus { await this.messagingHookHandler.handleBeforeMessageHandler( HookMessage.fromRoutingMessage( MessageFactory.creteRoutingFromMessage(messageToDispatch, message), - this.channel.config.name, - this.channel.constructor.name, ), ); @@ -108,8 +105,6 @@ export class InMemoryMessageBus implements IMessageBus { await this.messagingHookHandler.handleAfterMessageHandlerExecuted( HookMessage.fromRoutingMessage( MessageFactory.creteRoutingFromMessage(messageToDispatch, message), - this.channel.config.name, - this.channel.constructor.name, ), ); diff --git a/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-handler.ts b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-handler.ts index 91da0cf..a4f6e6a 100644 --- a/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-handler.ts +++ b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-handler.ts @@ -46,4 +46,10 @@ export class MessagingLifecycleHookHandler { .getAllByHook(LifecycleHook.AFTER_MESSAGE_NORMALIZATION) .forEach((listener) => listener.hook(message)); } + + async handleOnConsumerHandledMessage(message: HookMessage): Promise { + await this.messagingHookRegistry + .getAllByHook(LifecycleHook.ON_CONSUMER_HANDLED_MESSAGE) + .forEach((listener) => listener.hook(message)); + } } diff --git a/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-listener.ts b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-listener.ts index 431b1dc..be44c91 100644 --- a/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-listener.ts +++ b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-listener.ts @@ -1,21 +1,14 @@ import { ConsumerMessage } from '../consumer/consumer-message'; import { RoutingMessage } from '../message/routing-message'; +import { SealedRoutingMessage } from '../message/sealed-routing-message'; export class HookMessage { constructor( public readonly message: T, public readonly routingKey: string, - public readonly channelName: string, - public readonly channelType: string, - ) {} - - static fromMessage( - message: T, - routingKey: string, - channelName: string, - channelType: string, - ): HookMessage { - return new HookMessage(message, routingKey, channelName, channelType); + public readonly channelName?: string, + public readonly channelType?: string, + ) { } static fromConsumerMessage( @@ -33,8 +26,21 @@ export class HookMessage { static fromRoutingMessage( routingMessage: RoutingMessage, - channelName: string, - channelType: string, + channelName?: string, + channelType?: string, + ): HookMessage { + return new HookMessage( + routingMessage.message, + routingMessage.messageRoutingKey, + channelName, + channelType, + ); + } + + static fromSealedRoutingMessage( + routingMessage: SealedRoutingMessage, + channelName?: string, + channelType?: string, ): HookMessage { return new HookMessage( routingMessage.message, @@ -48,6 +54,7 @@ export class HookMessage { export enum LifecycleHook { BEFORE_MESSAGE_NORMALIZATION = 'BEFORE_MESSAGE_NORMALIZATION', AFTER_MESSAGE_NORMALIZATION = 'AFTER_MESSAGE_NORMALIZATION', + ON_CONSUMER_HANDLED_MESSAGE = 'ON_CONSUMER_HANDLED_MESSAGE', AFTER_MESSAGE_DENORMALIZED = 'AFTER_MESSAGE_DENORMALIZED', BEFORE_MESSAGE_HANDLER = 'BEFORE_MESSAGE_HANDLER', AFTER_MESSAGE_HANDLER_EXECUTION = 'AFTER_MESSAGE_HANDLER_EXECUTION', From 139fa5d3dcfe5a061ad71375847f33592422d44c Mon Sep 17 00:00:00 2001 From: sebastian Date: Wed, 18 Mar 2026 23:21:37 +0100 Subject: [PATCH 12/16] [messaging] feat: lifecycle hooks --- packages/messaging/package.json | 2 +- packages/messaging/src/bus/consumer.message-bus.ts | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/packages/messaging/package.json b/packages/messaging/package.json index 414fdd9..f756bba 100644 --- a/packages/messaging/package.json +++ b/packages/messaging/package.json @@ -1,6 +1,6 @@ { "name": "@nestjstools/messaging", - "version": "4.2.0-beta.3", + "version": "4.2.0-beta.4", "description": "Simplifies asynchronous and synchronous message handling with support for buses, handlers, channels, and consumers. Build scalable, decoupled applications with ease and reliability.", "author": "Sebastian Iwanczyszyn", "license": "MIT", diff --git a/packages/messaging/src/bus/consumer.message-bus.ts b/packages/messaging/src/bus/consumer.message-bus.ts index 5f9479b..c6d0dd0 100644 --- a/packages/messaging/src/bus/consumer.message-bus.ts +++ b/packages/messaging/src/bus/consumer.message-bus.ts @@ -51,7 +51,7 @@ export class ConsumerMessageBus { ), ); - await this.messagingHookHandler.handleAfterMessageDenormalized( + await this.messagingHookHandler.handleOnConsumerHandledMessage( HookMessage.fromSealedRoutingMessage( routingMessage, this.channel.config.name, From 3693c3e4165184beae9ce4bebba397d855fcf050 Mon Sep 17 00:00:00 2001 From: sebastian Date: Wed, 18 Mar 2026 23:24:07 +0100 Subject: [PATCH 13/16] [messaging] feat: lifecycle hooks --- packages/messaging/src/bus/consumer.message-bus.ts | 3 +-- packages/messaging/src/bus/distributed-message.bus.ts | 3 +-- packages/messaging/src/bus/in-memory-message.bus.ts | 3 +-- .../src/lifecycle-hook/messaging-lifecycle-hook-listener.ts | 3 +-- 4 files changed, 4 insertions(+), 8 deletions(-) diff --git a/packages/messaging/src/bus/consumer.message-bus.ts b/packages/messaging/src/bus/consumer.message-bus.ts index c6d0dd0..583f492 100644 --- a/packages/messaging/src/bus/consumer.message-bus.ts +++ b/packages/messaging/src/bus/consumer.message-bus.ts @@ -23,8 +23,7 @@ export class ConsumerMessageBus { private readonly consumer: IMessagingConsumer, private readonly exceptionListenerHandler: ExceptionListenerHandler, private readonly messagingHookHandler: MessagingLifecycleHookHandler, - ) { - } + ) {} async dispatch(consumerMessage: ConsumerMessage): Promise { try { diff --git a/packages/messaging/src/bus/distributed-message.bus.ts b/packages/messaging/src/bus/distributed-message.bus.ts index b61e39a..381d41d 100644 --- a/packages/messaging/src/bus/distributed-message.bus.ts +++ b/packages/messaging/src/bus/distributed-message.bus.ts @@ -14,8 +14,7 @@ export class DistributedMessageBus implements IMessageBus { private messageBusCollection: MessageBusCollection, private normalizerRegistry: NormalizerRegistry, private messagingLifecycleHookHandler: MessagingLifecycleHookHandler, - ) { - } + ) {} async dispatch(message: RoutingMessage): Promise { if (!(message instanceof RoutingMessage)) { diff --git a/packages/messaging/src/bus/in-memory-message.bus.ts b/packages/messaging/src/bus/in-memory-message.bus.ts index cb21caf..dd2b7a5 100644 --- a/packages/messaging/src/bus/in-memory-message.bus.ts +++ b/packages/messaging/src/bus/in-memory-message.bus.ts @@ -23,8 +23,7 @@ export class InMemoryMessageBus implements IMessageBus { private channel: InMemoryChannel, private normalizerRegistry: NormalizerRegistry, private messagingHookHandler: MessagingLifecycleHookHandler, - ) { - } + ) {} async dispatch(message: Message): Promise { const middlewares = []; diff --git a/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-listener.ts b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-listener.ts index be44c91..3520e57 100644 --- a/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-listener.ts +++ b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-listener.ts @@ -8,8 +8,7 @@ export class HookMessage { public readonly routingKey: string, public readonly channelName?: string, public readonly channelType?: string, - ) { - } + ) {} static fromConsumerMessage( consumerMessage: ConsumerMessage, From faaa9753c4c7e7ca0a86e0dfcf44456ddd389e9b Mon Sep 17 00:00:00 2001 From: sebastian Date: Thu, 19 Mar 2026 08:51:04 +0100 Subject: [PATCH 14/16] [messaging] feat: lifecycle hooks --- .../test/unit/bus/consumer.message-bus.spec.ts | 16 ++++------------ 1 file changed, 4 insertions(+), 12 deletions(-) diff --git a/packages/messaging/test/unit/bus/consumer.message-bus.spec.ts b/packages/messaging/test/unit/bus/consumer.message-bus.spec.ts index 98a1997..ec39a91 100644 --- a/packages/messaging/test/unit/bus/consumer.message-bus.spec.ts +++ b/packages/messaging/test/unit/bus/consumer.message-bus.spec.ts @@ -14,6 +14,7 @@ import { HandlerError, HandlersException } from '../../../src'; import { ExceptionContext } from '../../../src'; import { MessagingLifecycleHookHandler } from '../../../src/lifecycle-hook/messaging-lifecycle-hook-handler'; import { Logger } from '@nestjs/common'; +import { MessagingLifecycleHookRegistry } from '../../../src/lifecycle-hook/messaging-lifecycle-hook.registry'; describe('ConsumerMessageBus', () => { let messageBus: IMessageBus; @@ -24,7 +25,6 @@ describe('ConsumerMessageBus', () => { let exceptionListenerHandler: ExceptionListenerHandler; let exceptionHandlerMock: jest.Mock; let messagingLifecycleHookHandler: MessagingLifecycleHookHandler; - let failedConsumerHookMock: jest.Mock; let channel: TestChannel; beforeEach(() => { @@ -44,15 +44,9 @@ describe('ConsumerMessageBus', () => { exceptionListenerHandler = { handleError: exceptionHandlerMock, } as unknown as ExceptionListenerHandler; - failedConsumerHookMock = jest.fn().mockResolvedValue(undefined); - messagingLifecycleHookHandler = { - handleAfterMessageDenormalized: jest.fn().mockResolvedValue(undefined), - handleBeforeMessageHandler: jest.fn().mockResolvedValue(undefined), - handleAfterMessageHandlerExecuted: jest.fn().mockResolvedValue(undefined), - handleOnFailedMessageConsumer: failedConsumerHookMock, - handleBeforeMessageNormalization: jest.fn().mockResolvedValue(undefined), - handleAfterMessageNormalization: jest.fn().mockResolvedValue(undefined), - } as unknown as MessagingLifecycleHookHandler; + + const registry = new MessagingLifecycleHookRegistry(); + messagingLifecycleHookHandler = new MessagingLifecycleHookHandler(registry); channel = new TestChannel(new InMemoryChannelConfig({ name: 'ds' })); }); @@ -134,7 +128,6 @@ describe('ConsumerMessageBus', () => { }, }, }); - expect(failedConsumerHookMock).toHaveBeenCalledTimes(1); }); it('should not log error when dispatch throws HandlersException', async () => { @@ -170,6 +163,5 @@ describe('ConsumerMessageBus', () => { ), ); expect(logger.getLogs().find((log) => log.type === 'ERROR')).toBeFalsy(); - expect(failedConsumerHookMock).toHaveBeenCalledTimes(1); }); }); From 9e5f2a9e159c260b5c7c65e489f358bcb527a66a Mon Sep 17 00:00:00 2001 From: sebastian Date: Thu, 19 Mar 2026 08:55:38 +0100 Subject: [PATCH 15/16] [messaging] feat: lifecycle hooks --- packages/messaging/src/dependency-injection/register.ts | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/packages/messaging/src/dependency-injection/register.ts b/packages/messaging/src/dependency-injection/register.ts index 711abb3..a4a17a4 100644 --- a/packages/messaging/src/dependency-injection/register.ts +++ b/packages/messaging/src/dependency-injection/register.ts @@ -96,7 +96,10 @@ export const registerMessagingHooks = ( const register = >( moduleRef: ModuleRef, discoveryService: DiscoveryService, - registryProvider: string | Function, + registryProvider: + | string + | symbol + | (abstract new (...args: unknown[]) => unknown), decoratorMetadata: string, name: string, ) => { From a9ba96d6c84ffcf520702454ec0c2fbd0df4371b Mon Sep 17 00:00:00 2001 From: sebastian Date: Thu, 19 Mar 2026 08:55:51 +0100 Subject: [PATCH 16/16] [messaging] feat: lifecycle hooks --- packages/messaging/package.json | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/messaging/package.json b/packages/messaging/package.json index f756bba..6bff5d1 100644 --- a/packages/messaging/package.json +++ b/packages/messaging/package.json @@ -1,6 +1,6 @@ { "name": "@nestjstools/messaging", - "version": "4.2.0-beta.4", + "version": "4.2.0", "description": "Simplifies asynchronous and synchronous message handling with support for buses, handlers, channels, and consumers. Build scalable, decoupled applications with ease and reliability.", "author": "Sebastian Iwanczyszyn", "license": "MIT",