diff --git a/framework/src/Volo.Abp.EventBus.Distributed.RabbitMQ/Volo/Abp/EventBus/Distributed/RabbitMq/RabbitMqDistributedEventBus.cs b/framework/src/Volo.Abp.EventBus.Distributed.RabbitMQ/Volo/Abp/EventBus/Distributed/RabbitMq/RabbitMqDistributedEventBus.cs index 2007a69bf5..297ac33df5 100644 --- a/framework/src/Volo.Abp.EventBus.Distributed.RabbitMQ/Volo/Abp/EventBus/Distributed/RabbitMq/RabbitMqDistributedEventBus.cs +++ b/framework/src/Volo.Abp.EventBus.Distributed.RabbitMQ/Volo/Abp/EventBus/Distributed/RabbitMq/RabbitMqDistributedEventBus.cs @@ -26,9 +26,10 @@ namespace Volo.Abp.EventBus.Distributed.RabbitMq protected DistributedEventBusOptions DistributedEventBusOptions { get; } protected IConnectionPool ConnectionPool { get; } protected IRabbitMqSerializer Serializer { get; } - protected ConcurrentDictionary> HandlerFactories { get; } //TODO: Accessing to the List may not be thread-safe! + + //TODO: Accessing to the List may not be thread-safe! + protected ConcurrentDictionary> HandlerFactories { get; } protected ConcurrentDictionary EventTypes { get; } - protected IHybridServiceScopeFactory ServiceScopeFactory { get; } protected IRabbitMqMessageConsumerFactory MessageConsumerFactory { get; } protected IRabbitMqMessageConsumer Consumer { get; } @@ -39,10 +40,10 @@ namespace Volo.Abp.EventBus.Distributed.RabbitMq IHybridServiceScopeFactory serviceScopeFactory, IOptions distributedEventBusOptions, IRabbitMqMessageConsumerFactory messageConsumerFactory) + : base(serviceScopeFactory) { ConnectionPool = connectionPool; Serializer = serializer; - ServiceScopeFactory = serviceScopeFactory; MessageConsumerFactory = messageConsumerFactory; DistributedEventBusOptions = distributedEventBusOptions.Value; RabbitMqDistributedEventBusOptions = options.Value; @@ -50,8 +51,6 @@ namespace Volo.Abp.EventBus.Distributed.RabbitMq HandlerFactories = new ConcurrentDictionary>(); EventTypes = new ConcurrentDictionary(); - Subscribe(DistributedEventBusOptions.Handlers); - Consumer = MessageConsumerFactory.Create( new ExchangeDeclareConfiguration( RabbitMqDistributedEventBusOptions.ExchangeName, @@ -67,27 +66,8 @@ namespace Volo.Abp.EventBus.Distributed.RabbitMq ); Consumer.OnMessageReceived(ProcessEventAsync); - } - protected virtual void Subscribe(ITypeList handlers) - { - foreach (var handler in handlers) - { - var interfaces = handler.GetInterfaces(); - foreach (var @interface in interfaces) - { - if (!typeof(IEventHandler).GetTypeInfo().IsAssignableFrom(@interface)) - { - continue; - } - - var genericArgs = @interface.GetGenericArguments(); - if (genericArgs.Length == 1) - { - Subscribe(genericArgs[0], new IocEventHandlerFactory(ServiceScopeFactory, handler)); - } - } - } + SubscribeHandlers(DistributedEventBusOptions.Handlers); } private async Task ProcessEventAsync(IModel channel, BasicDeliverEventArgs ea) @@ -186,11 +166,14 @@ namespace Volo.Abp.EventBus.Distributed.RabbitMq using (var channel = ConnectionPool.Get(RabbitMqDistributedEventBusOptions.ConnectionName).CreateModel()) { - //TODO: Other properties like durable? - channel.ExchangeDeclare(RabbitMqDistributedEventBusOptions.ExchangeName, ""); + channel.ExchangeDeclare( + RabbitMqDistributedEventBusOptions.ExchangeName, + "direct" + //TODO: Other properties like durable? + ); var properties = channel.CreateBasicProperties(); - properties.DeliveryMode = 2; //persistent + properties.DeliveryMode = RabbitMqConsts.DeliveryModes.Persistent; channel.BasicPublish( exchange: RabbitMqDistributedEventBusOptions.ExchangeName, 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 c700f46532..3b2da52a96 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs @@ -5,6 +5,8 @@ using System.Reflection; using System.Runtime.CompilerServices; using System.Threading; using System.Threading.Tasks; +using Volo.Abp.Collections; +using Volo.Abp.DependencyInjection; using Volo.Abp.EventBus.Distributed; using Volo.Abp.Reflection; @@ -12,6 +14,13 @@ namespace Volo.Abp.EventBus { public abstract class EventBusBase : IEventBus { + protected IHybridServiceScopeFactory ServiceScopeFactory { get; } + + protected EventBusBase(IHybridServiceScopeFactory serviceScopeFactory) + { + ServiceScopeFactory = serviceScopeFactory; + } + /// public virtual IDisposable Subscribe(Func action) where TEvent : class { @@ -122,6 +131,27 @@ namespace Volo.Abp.EventBus } } + protected virtual void SubscribeHandlers(ITypeList handlers) + { + foreach (var handler in handlers) + { + var interfaces = handler.GetInterfaces(); + foreach (var @interface in interfaces) + { + if (!typeof(IEventHandler).GetTypeInfo().IsAssignableFrom(@interface)) + { + continue; + } + + var genericArgs = @interface.GetGenericArguments(); + if (genericArgs.Length == 1) + { + Subscribe(genericArgs[0], new IocEventHandlerFactory(ServiceScopeFactory, handler)); + } + } + } + } + protected abstract IEnumerable GetHandlerFactories(Type eventType); protected virtual async Task TriggerHandlerAsync(IEventHandlerFactory asyncHandlerFactory, Type eventType, object eventData, List exceptions) diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventBus.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventBus.cs index ae1041c9e2..e8786073a7 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventBus.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventBus.cs @@ -5,9 +5,7 @@ using System; using System.Collections.Concurrent; using System.Collections.Generic; using System.Linq; -using System.Reflection; using System.Threading.Tasks; -using Volo.Abp.Collections; using Volo.Abp.DependencyInjection; using Volo.Abp.Threading; @@ -28,39 +26,16 @@ namespace Volo.Abp.EventBus.Local protected ConcurrentDictionary> HandlerFactories { get; } - protected IHybridServiceScopeFactory ServiceScopeFactory { get; } - public LocalEventBus( IOptions options, IHybridServiceScopeFactory serviceScopeFactory) + : base(serviceScopeFactory) { - ServiceScopeFactory = serviceScopeFactory; Options = options.Value; Logger = NullLogger.Instance; HandlerFactories = new ConcurrentDictionary>(); - Subscribe(Options.Handlers); - } - - public virtual void Subscribe(ITypeList handlers) - { - foreach (var handler in handlers) - { - var interfaces = handler.GetInterfaces(); - foreach (var @interface in interfaces) - { - if (!typeof(IEventHandler).GetTypeInfo().IsAssignableFrom(@interface)) - { - continue; - } - - var genericArgs = @interface.GetGenericArguments(); - if (genericArgs.Length == 1) - { - Subscribe(genericArgs[0], new IocEventHandlerFactory(ServiceScopeFactory, handler)); - } - } - } + SubscribeHandlers(Options.Handlers); } /// diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/RabbitMqConsts.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/RabbitMqConsts.cs new file mode 100644 index 0000000000..0869dfd3b9 --- /dev/null +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/RabbitMqConsts.cs @@ -0,0 +1,12 @@ +namespace Volo.Abp.RabbitMQ +{ + public static class RabbitMqConsts + { + public static class DeliveryModes + { + public const int NonPersistent = 1; + + public const int Persistent = 2; + } + } +}