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/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", diff --git a/packages/messaging/src/bus/consumer.message-bus.ts b/packages/messaging/src/bus/consumer.message-bus.ts index 648405c..583f492 100644 --- a/packages/messaging/src/bus/consumer.message-bus.ts +++ b/packages/messaging/src/bus/consumer.message-bus.ts @@ -11,6 +11,9 @@ 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 { HookMessage } from '../lifecycle-hook/messaging-lifecycle-hook-listener'; +import { MessageFactory } from '../message/message.factory'; export class ConsumerMessageBus { constructor( @@ -19,6 +22,7 @@ export class ConsumerMessageBus { private readonly logger: MessagingLogger, private readonly consumer: IMessagingConsumer, private readonly exceptionListenerHandler: ExceptionListenerHandler, + private readonly messagingHookHandler: MessagingLifecycleHookHandler, ) {} async dispatch(consumerMessage: ConsumerMessage): Promise { @@ -46,6 +50,14 @@ export class ConsumerMessageBus { ), ); + await this.messagingHookHandler.handleOnConsumerHandledMessage( + HookMessage.fromSealedRoutingMessage( + routingMessage, + this.channel.config.name, + this.channel.constructor.name, + ), + ); + await this.messageBus.dispatch(routingMessage); } catch (e) { await this.consumer.onError( @@ -74,6 +86,14 @@ export class ConsumerMessageBus { consumerMessage.routingKey, ), ); + + await this.messagingHookHandler.handleOnFailedMessageConsumer( + HookMessage.fromConsumerMessage( + consumerMessage, + this.channel.config.name, + this.channel.constructor.name, + ), + ); } return Promise.resolve(); diff --git a/packages/messaging/src/bus/distributed-message.bus.ts b/packages/messaging/src/bus/distributed-message.bus.ts index 6dd37d1..381d41d 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 { HookMessage } 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,30 @@ export class DistributedMessageBus implements IMessageBus { const response = []; for (const collection of this.messageBusCollection.getAll()) { + await this.messagingLifecycleHookHandler.handleBeforeMessageNormalization( + HookMessage.fromRoutingMessage( + message, + 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( + HookMessage.fromRoutingMessage( + message, + 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 56ea2f6..b585001 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 { MessagingLifecycleHookHandler } from '../lifecycle-hook/messaging-lifecycle-hook-handler'; @Injectable() @MessageBusFactory(InMemoryChannel) @@ -19,6 +20,7 @@ export class InMemoryMessageBusFactory implements IMessageBusFactory { @@ -30,6 +33,28 @@ export class InMemoryMessageBus implements IMessageBus { HandlerMiddleware, ); + let messageToDispatch = + 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 + : ObjectForwardMessageNormalizer; + + messageToDispatch = await this.normalizerRegistry + .getByName(normalizerDefinition['name']) + .denormalize(message.message, message.messageRoutingKey); + } + + // Hook fired once the payload shape is ready for handler pipeline. + await this.messagingHookHandler.handleAfterMessageDenormalized( + HookMessage.fromRoutingMessage( + MessageFactory.creteRoutingFromMessage(messageToDispatch, message), + ), + ); + try { this.registry.getByRoutingKey(message.messageRoutingKey); } catch (e) { @@ -48,6 +73,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(); } @@ -60,27 +86,27 @@ 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; + const context = MiddlewareContext.createFresh(middlewareInstances); - messageToDispatch = await this.normalizerRegistry - .getByName(normalizerDefinition['name']) - .denormalize(message.message, message.messageRoutingKey); - } + // Hook around handler execution. + await this.messagingHookHandler.handleBeforeMessageHandler( + HookMessage.fromRoutingMessage( + MessageFactory.creteRoutingFromMessage(messageToDispatch, message), + ), + ); const response = await middlewareInstances[0].process( MessageFactory.creteRoutingFromMessage(messageToDispatch, message), context, ); + await this.messagingHookHandler.handleAfterMessageHandlerExecuted( + HookMessage.fromRoutingMessage( + MessageFactory.creteRoutingFromMessage(messageToDispatch, message), + ), + ); + return Promise.resolve(response); } } diff --git a/packages/messaging/src/consumer/distributed.consumer.ts b/packages/messaging/src/consumer/distributed.consumer.ts index 40620b5..a7309e2 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, ) {} async run(): Promise { @@ -63,6 +66,7 @@ export class DistributedConsumer { this.logger, consumer, this.exceptionListenerHandler, + this.messagingLifecycleHookHandler, ); await consumer.consume(dispatcher, channel); diff --git a/packages/messaging/src/dependency-injection/decorator.ts b/packages/messaging/src/dependency-injection/decorator.ts index 56bf903..e00923d 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 { LifecycleHook } from '../lifecycle-hook/messaging-lifecycle-hook-listener'; export const MESSAGE_HANDLER_METADATA = 'MESSAGE_HANDLER_METADATA'; export const CHANNEL_FACTORY_METADATA = 'CHANNEL_FACTORY_METADATA'; @@ -9,6 +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 MessageHandler = (...routingKey: string[]): ClassDecorator => { return (target) => { @@ -70,6 +73,18 @@ export const MessagingExceptionListener = (): ClassDecorator => { }; }; +export const MessagingLifecycleHook = ( + lifecycleHook: LifecycleHook, +): ClassDecorator => { + return (target) => { + Reflect.defineMetadata( + MESSAGING_LIFECYCLE_HOOK_METADATA, + `${lifecycleHook}:${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..a4a17a4 100644 --- a/packages/messaging/src/dependency-injection/register.ts +++ b/packages/messaging/src/dependency-injection/register.ts @@ -5,6 +5,7 @@ import { Service } from './service'; import { MESSAGE_HANDLER_METADATA, MESSAGING_EXCEPTION_LISTENER_METADATA, + MESSAGING_LIFECYCLE_HOOK_METADATA, MESSAGING_MIDDLEWARE_METADATA, MESSAGING_NORMALIZER_METADATA, } from './decorator'; @@ -13,6 +14,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 { MessagingLifecycleHookRegistry } from '../lifecycle-hook/messaging-lifecycle-hook.registry'; export const registerHandlers = ( moduleRef: ModuleRef, @@ -78,36 +80,47 @@ export const registerExceptionListener = ( ); }; +export const registerMessagingHooks = ( + moduleRef: ModuleRef, + discoveryService: DiscoveryService, +) => { + register( + moduleRef, + discoveryService, + MessagingLifecycleHookRegistry, + MESSAGING_LIFECYCLE_HOOK_METADATA, + 'MessagingLifecycleHook', + ); +}; + const register = >( moduleRef: ModuleRef, discoveryService: DiscoveryService, - registryProvider: string, + registryProvider: + | string + | symbol + | (abstract new (...args: unknown[]) => unknown), decoratorMetadata: string, name: string, ) => { const exceptions = [DEFAULT_NORMALIZER, DEFAULT_MIDDLEWARE]; const registry: Registry = moduleRef.get(registryProvider); const logger: MessagingLogger = moduleRef.get(Service.LOGGER); - const instances = discoveryService - .getProviders() - .filter((messageExceptionListener) => { - if (!messageExceptionListener.metatype) { - return false; - } + const instances = discoveryService.getProviders().filter((provider) => { + if (!provider.metatype) { + return false; + } - return Reflect.hasMetadata( - decoratorMetadata, - messageExceptionListener.metatype, - ); - }); + return Reflect.hasMetadata(decoratorMetadata, 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..54afd0d 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..a4f6e6a --- /dev/null +++ b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-handler.ts @@ -0,0 +1,55 @@ +import { Injectable } from '@nestjs/common'; +import { MessagingLifecycleHookRegistry } from './messaging-lifecycle-hook.registry'; +import { + LifecycleHook, + HookMessage, +} from './messaging-lifecycle-hook-listener'; + +@Injectable() +export class MessagingLifecycleHookHandler { + constructor( + private readonly messagingHookRegistry: MessagingLifecycleHookRegistry, + ) {} + + async handleAfterMessageDenormalized(message: HookMessage): Promise { + await this.messagingHookRegistry + .getAllByHook(LifecycleHook.AFTER_MESSAGE_DENORMALIZED) + .forEach((listener) => listener.hook(message)); + } + + async handleBeforeMessageHandler(message: HookMessage): Promise { + await this.messagingHookRegistry + .getAllByHook(LifecycleHook.BEFORE_MESSAGE_HANDLER) + .forEach((listener) => listener.hook(message)); + } + + async handleAfterMessageHandlerExecuted(message: HookMessage): Promise { + await this.messagingHookRegistry + .getAllByHook(LifecycleHook.AFTER_MESSAGE_HANDLER_EXECUTION) + .forEach((listener) => listener.hook(message)); + } + + async handleOnFailedMessageConsumer(message: HookMessage): Promise { + await this.messagingHookRegistry + .getAllByHook(LifecycleHook.ON_FAILED_MESSAGE_CONSUMER) + .forEach((listener) => listener.hook(message)); + } + + async handleBeforeMessageNormalization(message: HookMessage): Promise { + await this.messagingHookRegistry + .getAllByHook(LifecycleHook.BEFORE_MESSAGE_NORMALIZATION) + .forEach((listener) => listener.hook(message)); + } + + async handleAfterMessageNormalization(message: HookMessage): Promise { + await this.messagingHookRegistry + .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 new file mode 100644 index 0000000..3520e57 --- /dev/null +++ b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook-listener.ts @@ -0,0 +1,65 @@ +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 fromConsumerMessage( + consumerMessage: ConsumerMessage, + channelName: string, + channelType: string, + ): HookMessage { + return new HookMessage( + consumerMessage.message, + consumerMessage.routingKey, + channelName, + channelType, + ); + } + + static fromRoutingMessage( + routingMessage: RoutingMessage, + 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, + routingMessage.messageRoutingKey, + channelName, + channelType, + ); + } +} + +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', + ON_FAILED_MESSAGE_CONSUMER = 'ON_FAILED_MESSAGE_CONSUMER', +} + +export interface MessagingLifecycleHookListener { + hook(message: HookMessage): 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..b2e60b9 --- /dev/null +++ b/packages/messaging/src/lifecycle-hook/messaging-lifecycle-hook.registry.ts @@ -0,0 +1,21 @@ +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/messaging.module.ts b/packages/messaging/src/messaging.module.ts index 420ec90..e5eb491 100644 --- a/packages/messaging/src/messaging.module.ts +++ b/packages/messaging/src/messaging.module.ts @@ -31,6 +31,7 @@ import { DistributedConsumer } from './consumer/distributed.consumer'; import { registerExceptionListener, registerHandlers, + registerMessagingHooks, registerMessageNormalizers, registerMiddlewares, } from './dependency-injection/register'; @@ -44,6 +45,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 { MessagingLifecycleHookHandler } from './lifecycle-hook/messaging-lifecycle-hook-handler'; +import { MessagingLifecycleHookRegistry } from './lifecycle-hook/messaging-lifecycle-hook.registry'; @Module({}) export class MessagingModule @@ -110,6 +113,7 @@ export class MessagingModule busFactory: CompositeMessageBusFactory, logger: MessagingLogger, normalizerRegistry: NormalizerRegistry, + messagingLifecycleHookHandler: MessagingLifecycleHookHandler, ) => { const messageBusCollection = new MessageBusCollection(); @@ -124,6 +128,7 @@ export class MessagingModule const messageBus = new DistributedMessageBus( messageBusCollection, normalizerRegistry, + messagingLifecycleHookHandler, ); logger.log(`MessageBus [${bus.name}] was created successfully`); @@ -135,6 +140,7 @@ export class MessagingModule CompositeMessageBusFactory, Service.LOGGER, Service.MESSAGE_NORMALIZERS_REGISTRY, + MessagingLifecycleHookHandler, ], })); }; @@ -146,6 +152,7 @@ export class MessagingModule messageHandlerRegistry: MessageHandlerRegistry, middlewareRegistry: MiddlewareRegistry, normalizerRegistry: NormalizerRegistry, + messagingHookHandler: MessagingLifecycleHookHandler, ) => { return new InMemoryMessageBus( messageHandlerRegistry, @@ -158,12 +165,14 @@ export class MessagingModule }), ), normalizerRegistry, + messagingHookHandler, ); }, inject: [ Service.MESSAGE_HANDLERS_REGISTRY, Service.MIDDLEWARE_REGISTRY, Service.MESSAGE_NORMALIZERS_REGISTRY, + MessagingLifecycleHookHandler, ], }; }; @@ -230,6 +239,8 @@ export class MessagingModule InMemoryChannelFactory, DistributedConsumer, ObjectForwardMessageNormalizer, + MessagingLifecycleHookRegistry, + MessagingLifecycleHookHandler, ], exports: [ Service.DEFAULT_MESSAGE_BUS, @@ -254,6 +265,7 @@ export class MessagingModule registerMiddlewares(this.moduleRef, this.discoveryService); registerMessageNormalizers(this.moduleRef, this.discoveryService); registerExceptionListener(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/consumer.message-bus.spec.ts b/packages/messaging/test/unit/bus/consumer.message-bus.spec.ts index 48d3234..ec39a91 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,41 @@ 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'; +import { MessagingLifecycleHookRegistry } from '../../../src/lifecycle-hook/messaging-lifecycle-hook.registry'; 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 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; + + const registry = new MessagingLifecycleHookRegistry(); + messagingLifecycleHookHandler = new MessagingLifecycleHookHandler(registry); channel = new TestChannel(new InMemoryChannelConfig({ name: 'ds' })); }); @@ -43,6 +57,7 @@ describe('ConsumerMessageBus', () => { logger, consumer, exceptionListenerHandler, + messagingLifecycleHookHandler, ); await subjectUnderTest.dispatch( @@ -72,8 +87,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 +98,7 @@ describe('ConsumerMessageBus', () => { logger, consumer, exceptionListenerHandler, + messagingLifecycleHookHandler, ); const consumerMessage = new ConsumerMessage({ status: 'fail' }, 'rk.fail'); @@ -118,8 +135,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 +146,7 @@ describe('ConsumerMessageBus', () => { logger, consumer, exceptionListenerHandler, + messagingLifecycleHookHandler, ); await subjectUnderTest.dispatch( 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/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-hook-handler.spec.ts b/packages/messaging/test/unit/lifecycle-hook/messaging-hook-handler.spec.ts new file mode 100644 index 0000000..b288eba --- /dev/null +++ b/packages/messaging/test/unit/lifecycle-hook/messaging-hook-handler.spec.ts @@ -0,0 +1,75 @@ +import { MessagingLifecycleHookHandler } from '../../../src/lifecycle-hook/messaging-lifecycle-hook-handler'; +import { MessagingLifecycleHookRegistry } from '../../../src/lifecycle-hook/messaging-lifecycle-hook.registry'; +import { + HookMessage, + LifecycleHook, + MessagingLifecycleHookListener, + RoutingMessage, +} from '../../../src'; + +describe('MessagingHookHandler', () => { + let hookRegistry: jest.Mocked; + let handler: MessagingLifecycleHookHandler; + let listener: MessagingLifecycleHookListener; + const message = new RoutingMessage({ title: 'hello' }, 'test.key'); + + beforeEach(() => { + hookRegistry = { + getAllByHook: jest.fn(), + } as unknown as jest.Mocked; + + handler = new MessagingLifecycleHookHandler(hookRegistry); + listener = { hook: jest.fn().mockResolvedValue(undefined) }; + }); + + test('should execute listeners for AFTER_MESSAGE_DENORMALIZED hook', async () => { + (hookRegistry.getAllByHook as jest.Mock).mockReturnValue([listener]); + + const hookMessage = HookMessage.fromRoutingMessage( + message, + 'example', + 'example', + ); + + await handler.handleAfterMessageDenormalized(hookMessage); + + expect(hookRegistry.getAllByHook).toHaveBeenCalledWith( + LifecycleHook.AFTER_MESSAGE_DENORMALIZED, + ); + expect(listener.hook).toHaveBeenCalledWith(hookMessage); + }); + + test('should execute listeners for BEFORE_MESSAGE_HANDLER hook', async () => { + (hookRegistry.getAllByHook as jest.Mock).mockReturnValue([listener]); + + const hookMessage = HookMessage.fromRoutingMessage( + message, + 'example', + 'example', + ); + + await handler.handleBeforeMessageHandler(hookMessage); + + expect(hookRegistry.getAllByHook).toHaveBeenCalledWith( + LifecycleHook.BEFORE_MESSAGE_HANDLER, + ); + expect(listener.hook).toHaveBeenCalledWith(hookMessage); + }); + + 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', + ); + + await handler.handleAfterMessageHandlerExecuted(hookMessage); + + expect(hookRegistry.getAllByHook).toHaveBeenCalledWith( + LifecycleHook.AFTER_MESSAGE_HANDLER_EXECUTION, + ); + 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 new file mode 100644 index 0000000..9fdc3db --- /dev/null +++ b/packages/messaging/test/unit/lifecycle-hook/messaging-lifecycle-hook.registry.spec.ts @@ -0,0 +1,47 @@ +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_EXECUTION), + ).toEqual([]); + }); +});