diff --git a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/InboxProcessor.cs b/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/InboxProcessor.cs index 4725e85d51..b1e6a1492b 100644 --- a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/InboxProcessor.cs +++ b/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/InboxProcessor.cs @@ -85,7 +85,7 @@ namespace Volo.Abp.EventBus.Boxes { await DistributedEventBus .AsRawEventPublisher() - .ProcessRawAsync(waitingEvent.EventName, waitingEvent.EventData); + .ProcessRawAsync(InboxConfig, waitingEvent.EventName, waitingEvent.EventData); /* await DistributedEventBus diff --git a/framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaDistributedEventBus.cs b/framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaDistributedEventBus.cs index ae20ab1dd0..e587dc9749 100644 --- a/framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaDistributedEventBus.cs +++ b/framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaDistributedEventBus.cs @@ -211,7 +211,7 @@ namespace Volo.Abp.EventBus.Kafka ); } - public override async Task ProcessRawAsync(string eventName, byte[] eventDataBytes) + public override async Task ProcessRawAsync(InboxConfig inboxConfig, string eventName, byte[] eventDataBytes) { var eventType = EventTypes.GetOrDefault(eventName); if (eventType == null) @@ -221,7 +221,7 @@ namespace Volo.Abp.EventBus.Kafka var eventData = Serializer.Deserialize(eventDataBytes, eventType); var exceptions = new List(); - await TriggerHandlersAsync(eventType, eventData, exceptions); + await TriggerHandlersAsync(eventType, eventData, exceptions, inboxConfig); if (exceptions.Any()) { ThrowOriginalExceptions(eventType, exceptions); diff --git a/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs b/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs index d3d74a048c..eb171f7ffb 100644 --- a/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs +++ b/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs @@ -217,7 +217,7 @@ namespace Volo.Abp.EventBus.RabbitMq return PublishAsync(eventName, eventData, null, eventId: eventId); } - public override async Task ProcessRawAsync(string eventName, byte[] eventDataBytes) + public override async Task ProcessRawAsync(InboxConfig inboxConfig, string eventName, byte[] eventDataBytes) { //TODO: We have a duplication in logic and also with the kafka side! @@ -229,7 +229,7 @@ namespace Volo.Abp.EventBus.RabbitMq var eventData = Serializer.Deserialize(eventDataBytes, eventType); var exceptions = new List(); - await TriggerHandlersAsync(eventType, eventData, exceptions); + await TriggerHandlersAsync(eventType, eventData, exceptions, inboxConfig); if (exceptions.Any()) { ThrowOriginalExceptions(eventType, exceptions); diff --git a/framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusDistributedEventBus.cs b/framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusDistributedEventBus.cs index bbe80d9866..af3d148ecf 100644 --- a/framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusDistributedEventBus.cs +++ b/framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusDistributedEventBus.cs @@ -194,7 +194,7 @@ namespace Volo.Abp.EventBus.Rebus throw new NotImplementedException(); } - public override Task ProcessRawAsync(string eventName, byte[] eventDataBytes) + public override Task ProcessRawAsync(InboxConfig inboxConfig, string eventName, byte[] eventDataBytes) { /* TODO: IMPLEMENT! */ throw new NotImplementedException(); diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/DistributedEventBusBase.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/DistributedEventBusBase.cs index a436294e88..8de449058b 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/DistributedEventBusBase.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/DistributedEventBusBase.cs @@ -81,7 +81,7 @@ namespace Volo.Abp.EventBus.Distributed } public abstract Task PublishRawAsync(Guid eventId, string eventName, byte[] eventData); - public abstract Task ProcessRawAsync(string eventName, byte[] eventDataBytes); + public abstract Task ProcessRawAsync(InboxConfig inboxConfig, string eventName, byte[] eventDataBytes); private async Task AddToOutboxAsync(Type eventType, object eventData) { diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/IRawEventPublisher.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/IRawEventPublisher.cs index 5f96b818b7..4505239ade 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/IRawEventPublisher.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/IRawEventPublisher.cs @@ -11,8 +11,8 @@ namespace Volo.Abp.EventBus.Distributed byte[] eventData); Task ProcessRawAsync( + InboxConfig inboxConfig, string eventName, - byte[] eventDataBytes - ); + byte[] eventDataBytes); } } \ No newline at end of file diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs index c4a61b24a6..694166b1bc 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs @@ -134,7 +134,7 @@ namespace Volo.Abp.EventBus } } - protected virtual async Task TriggerHandlersAsync(Type eventType, object eventData , List exceptions) + protected virtual async Task TriggerHandlersAsync(Type eventType, object eventData, List exceptions, InboxConfig inboxConfig = null) { await new SynchronizationContextRemover(); @@ -142,7 +142,7 @@ namespace Volo.Abp.EventBus { foreach (var handlerFactory in handlerFactories.EventHandlerFactories) { - await TriggerHandlerAsync(handlerFactory, handlerFactories.EventType, eventData, exceptions); + await TriggerHandlerAsync(handlerFactory, handlerFactories.EventType, eventData, exceptions, inboxConfig); } } @@ -199,7 +199,8 @@ namespace Volo.Abp.EventBus protected abstract IEnumerable GetHandlerFactories(Type eventType); - protected virtual async Task TriggerHandlerAsync(IEventHandlerFactory asyncHandlerFactory, Type eventType, object eventData, List exceptions) + protected virtual async Task TriggerHandlerAsync(IEventHandlerFactory asyncHandlerFactory, Type eventType, + object eventData, List exceptions, InboxConfig inboxConfig = null) { using (var eventHandlerWrapper = asyncHandlerFactory.GetHandler()) { @@ -207,6 +208,12 @@ namespace Volo.Abp.EventBus { var handlerType = eventHandlerWrapper.EventHandler.GetType(); + if (inboxConfig?.HandlerSelector != null && + !inboxConfig.HandlerSelector(handlerType)) + { + return; + } + using (CurrentTenant.Change(GetEventDataTenantId(eventData))) { if (ReflectionHelper.IsAssignableToGenericType(handlerType, typeof(ILocalEventHandler<>)))