Browse Source

Check handlerselector for inbox event processing

pull/10008/head
Halil İbrahim Kalkan 5 years ago
parent
commit
2aa45fe208
  1. 2
      framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/InboxProcessor.cs
  2. 4
      framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaDistributedEventBus.cs
  3. 4
      framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs
  4. 2
      framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusDistributedEventBus.cs
  5. 2
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/DistributedEventBusBase.cs
  6. 4
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/IRawEventPublisher.cs
  7. 13
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs

2
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

4
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<Exception>();
await TriggerHandlersAsync(eventType, eventData, exceptions);
await TriggerHandlersAsync(eventType, eventData, exceptions, inboxConfig);
if (exceptions.Any())
{
ThrowOriginalExceptions(eventType, exceptions);

4
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<Exception>();
await TriggerHandlersAsync(eventType, eventData, exceptions);
await TriggerHandlersAsync(eventType, eventData, exceptions, inboxConfig);
if (exceptions.Any())
{
ThrowOriginalExceptions(eventType, exceptions);

2
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();

2
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<bool> AddToOutboxAsync(Type eventType, object eventData)
{

4
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);
}
}

13
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<Exception> exceptions)
protected virtual async Task TriggerHandlersAsync(Type eventType, object eventData, List<Exception> 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<EventTypeWithEventHandlerFactories> GetHandlerFactories(Type eventType);
protected virtual async Task TriggerHandlerAsync(IEventHandlerFactory asyncHandlerFactory, Type eventType, object eventData, List<Exception> exceptions)
protected virtual async Task TriggerHandlerAsync(IEventHandlerFactory asyncHandlerFactory, Type eventType,
object eventData, List<Exception> 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<>)))

Loading…
Cancel
Save