From 5f50d3323ab99e8a048f274e3a4ab91ae08135fb Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Halil=20=C4=B0brahim=20Kalkan?= Date: Tue, 7 Sep 2021 17:11:42 +0300 Subject: [PATCH] Add AbpDistributedEventBusOptions to DistributedEventBusBase --- .../Volo/Abp/EventBus/Kafka/KafkaDistributedEventBus.cs | 9 ++++++--- .../Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs | 9 ++++++--- .../Volo/Abp/EventBus/Rebus/RebusDistributedEventBus.cs | 9 ++++++--- .../Distributed/AbpDistributedEventBusOptions.cs | 8 ++++++++ .../Abp/EventBus/Distributed/DistributedEventBusBase.cs | 7 ++++++- 5 files changed, 32 insertions(+), 10 deletions(-) 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 c3050adb3b..b2d099286a 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 @@ -22,7 +22,6 @@ namespace Volo.Abp.EventBus.Kafka { protected AbpEventBusOptions AbpEventBusOptions { get; } protected AbpKafkaEventBusOptions AbpKafkaEventBusOptions { get; } - protected AbpDistributedEventBusOptions AbpDistributedEventBusOptions { get; } protected IKafkaMessageConsumerFactory MessageConsumerFactory { get; } protected IKafkaSerializer Serializer { get; } protected IProducerPool ProducerPool { get; } @@ -42,10 +41,14 @@ namespace Volo.Abp.EventBus.Kafka IProducerPool producerPool, IEventErrorHandler errorHandler, IOptions abpEventBusOptions) - : base(serviceScopeFactory, currentTenant, unitOfWorkManager, errorHandler) + : base( + serviceScopeFactory, + currentTenant, + unitOfWorkManager, + errorHandler, + abpDistributedEventBusOptions) { AbpKafkaEventBusOptions = abpKafkaEventBusOptions.Value; - AbpDistributedEventBusOptions = abpDistributedEventBusOptions.Value; AbpEventBusOptions = abpEventBusOptions.Value; MessageConsumerFactory = messageConsumerFactory; Serializer = serializer; 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 9891b5cf12..471bcacad6 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 @@ -26,7 +26,6 @@ namespace Volo.Abp.EventBus.RabbitMq public class RabbitMqDistributedEventBus : DistributedEventBusBase, ISingletonDependency { protected AbpRabbitMqEventBusOptions AbpRabbitMqEventBusOptions { get; } - protected AbpDistributedEventBusOptions AbpDistributedEventBusOptions { get; } protected AbpEventBusOptions AbpEventBusOptions { get; } protected IConnectionPool ConnectionPool { get; } protected IRabbitMqSerializer Serializer { get; } @@ -48,13 +47,17 @@ namespace Volo.Abp.EventBus.RabbitMq IUnitOfWorkManager unitOfWorkManager, IEventErrorHandler errorHandler, IOptions abpEventBusOptions) - : base(serviceScopeFactory, currentTenant, unitOfWorkManager, errorHandler) + : base( + serviceScopeFactory, + currentTenant, + unitOfWorkManager, + errorHandler, + distributedEventBusOptions) { ConnectionPool = connectionPool; Serializer = serializer; MessageConsumerFactory = messageConsumerFactory; AbpEventBusOptions = abpEventBusOptions.Value; - AbpDistributedEventBusOptions = distributedEventBusOptions.Value; AbpRabbitMqEventBusOptions = options.Value; HandlerFactories = new ConcurrentDictionary>(); 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 3f07cf7150..760af619f5 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 @@ -23,7 +23,6 @@ namespace Volo.Abp.EventBus.Rebus //TODO: Accessing to the List may not be thread-safe! protected ConcurrentDictionary> HandlerFactories { get; } protected ConcurrentDictionary EventTypes { get; } - protected AbpDistributedEventBusOptions AbpDistributedEventBusOptions { get; } protected AbpRebusEventBusOptions AbpRebusEventBusOptions { get; } public RebusDistributedEventBus( @@ -34,11 +33,15 @@ namespace Volo.Abp.EventBus.Rebus IOptions abpDistributedEventBusOptions, IOptions abpEventBusRebusOptions, IEventErrorHandler errorHandler) : - base(serviceScopeFactory, currentTenant, unitOfWorkManager, errorHandler) + base( + serviceScopeFactory, + currentTenant, + unitOfWorkManager, + errorHandler, + abpDistributedEventBusOptions) { Rebus = rebus; AbpRebusEventBusOptions = abpEventBusRebusOptions.Value; - AbpDistributedEventBusOptions = abpDistributedEventBusOptions.Value; HandlerFactories = new ConcurrentDictionary>(); EventTypes = new ConcurrentDictionary(); diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/AbpDistributedEventBusOptions.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/AbpDistributedEventBusOptions.cs index 342e7b77bb..2d3861da73 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/AbpDistributedEventBusOptions.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/AbpDistributedEventBusOptions.cs @@ -1,3 +1,4 @@ +using System.Collections.Generic; using Volo.Abp.Collections; namespace Volo.Abp.EventBus.Distributed @@ -5,10 +6,17 @@ namespace Volo.Abp.EventBus.Distributed public class AbpDistributedEventBusOptions { public ITypeList Handlers { get; } + + public List Outboxes { get; } public AbpDistributedEventBusOptions() { Handlers = new TypeList(); } } + + public class OutboxConfig + { + + } } \ No newline at end of file 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 94176d9f50..541373cd13 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 @@ -1,6 +1,7 @@ using System; using System.Threading.Tasks; using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Options; using Volo.Abp.MultiTenancy; using Volo.Abp.Uow; @@ -8,17 +9,21 @@ namespace Volo.Abp.EventBus.Distributed { public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventBus { + protected AbpDistributedEventBusOptions AbpDistributedEventBusOptions { get; } + protected DistributedEventBusBase( IServiceScopeFactory serviceScopeFactory, ICurrentTenant currentTenant, IUnitOfWorkManager unitOfWorkManager, - IEventErrorHandler errorHandler + IEventErrorHandler errorHandler, + IOptions abpDistributedEventBusOptions ) : base( serviceScopeFactory, currentTenant, unitOfWorkManager, errorHandler) { + AbpDistributedEventBusOptions = abpDistributedEventBusOptions.Value; } public IDisposable Subscribe(IDistributedEventHandler handler) where TEvent : class