Browse Source

Add AbpDistributedEventBusOptions to DistributedEventBusBase

pull/10008/head
Halil İbrahim Kalkan 5 years ago
parent
commit
5f50d3323a
  1. 9
      framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaDistributedEventBus.cs
  2. 9
      framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqDistributedEventBus.cs
  3. 9
      framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusDistributedEventBus.cs
  4. 8
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/AbpDistributedEventBusOptions.cs
  5. 7
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/DistributedEventBusBase.cs

9
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 AbpEventBusOptions AbpEventBusOptions { get; }
protected AbpKafkaEventBusOptions AbpKafkaEventBusOptions { get; } protected AbpKafkaEventBusOptions AbpKafkaEventBusOptions { get; }
protected AbpDistributedEventBusOptions AbpDistributedEventBusOptions { get; }
protected IKafkaMessageConsumerFactory MessageConsumerFactory { get; } protected IKafkaMessageConsumerFactory MessageConsumerFactory { get; }
protected IKafkaSerializer Serializer { get; } protected IKafkaSerializer Serializer { get; }
protected IProducerPool ProducerPool { get; } protected IProducerPool ProducerPool { get; }
@ -42,10 +41,14 @@ namespace Volo.Abp.EventBus.Kafka
IProducerPool producerPool, IProducerPool producerPool,
IEventErrorHandler errorHandler, IEventErrorHandler errorHandler,
IOptions<AbpEventBusOptions> abpEventBusOptions) IOptions<AbpEventBusOptions> abpEventBusOptions)
: base(serviceScopeFactory, currentTenant, unitOfWorkManager, errorHandler) : base(
serviceScopeFactory,
currentTenant,
unitOfWorkManager,
errorHandler,
abpDistributedEventBusOptions)
{ {
AbpKafkaEventBusOptions = abpKafkaEventBusOptions.Value; AbpKafkaEventBusOptions = abpKafkaEventBusOptions.Value;
AbpDistributedEventBusOptions = abpDistributedEventBusOptions.Value;
AbpEventBusOptions = abpEventBusOptions.Value; AbpEventBusOptions = abpEventBusOptions.Value;
MessageConsumerFactory = messageConsumerFactory; MessageConsumerFactory = messageConsumerFactory;
Serializer = serializer; Serializer = serializer;

9
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 public class RabbitMqDistributedEventBus : DistributedEventBusBase, ISingletonDependency
{ {
protected AbpRabbitMqEventBusOptions AbpRabbitMqEventBusOptions { get; } protected AbpRabbitMqEventBusOptions AbpRabbitMqEventBusOptions { get; }
protected AbpDistributedEventBusOptions AbpDistributedEventBusOptions { get; }
protected AbpEventBusOptions AbpEventBusOptions { get; } protected AbpEventBusOptions AbpEventBusOptions { get; }
protected IConnectionPool ConnectionPool { get; } protected IConnectionPool ConnectionPool { get; }
protected IRabbitMqSerializer Serializer { get; } protected IRabbitMqSerializer Serializer { get; }
@ -48,13 +47,17 @@ namespace Volo.Abp.EventBus.RabbitMq
IUnitOfWorkManager unitOfWorkManager, IUnitOfWorkManager unitOfWorkManager,
IEventErrorHandler errorHandler, IEventErrorHandler errorHandler,
IOptions<AbpEventBusOptions> abpEventBusOptions) IOptions<AbpEventBusOptions> abpEventBusOptions)
: base(serviceScopeFactory, currentTenant, unitOfWorkManager, errorHandler) : base(
serviceScopeFactory,
currentTenant,
unitOfWorkManager,
errorHandler,
distributedEventBusOptions)
{ {
ConnectionPool = connectionPool; ConnectionPool = connectionPool;
Serializer = serializer; Serializer = serializer;
MessageConsumerFactory = messageConsumerFactory; MessageConsumerFactory = messageConsumerFactory;
AbpEventBusOptions = abpEventBusOptions.Value; AbpEventBusOptions = abpEventBusOptions.Value;
AbpDistributedEventBusOptions = distributedEventBusOptions.Value;
AbpRabbitMqEventBusOptions = options.Value; AbpRabbitMqEventBusOptions = options.Value;
HandlerFactories = new ConcurrentDictionary<Type, List<IEventHandlerFactory>>(); HandlerFactories = new ConcurrentDictionary<Type, List<IEventHandlerFactory>>();

9
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<IEventHandlerFactory> may not be thread-safe! //TODO: Accessing to the List<IEventHandlerFactory> may not be thread-safe!
protected ConcurrentDictionary<Type, List<IEventHandlerFactory>> HandlerFactories { get; } protected ConcurrentDictionary<Type, List<IEventHandlerFactory>> HandlerFactories { get; }
protected ConcurrentDictionary<string, Type> EventTypes { get; } protected ConcurrentDictionary<string, Type> EventTypes { get; }
protected AbpDistributedEventBusOptions AbpDistributedEventBusOptions { get; }
protected AbpRebusEventBusOptions AbpRebusEventBusOptions { get; } protected AbpRebusEventBusOptions AbpRebusEventBusOptions { get; }
public RebusDistributedEventBus( public RebusDistributedEventBus(
@ -34,11 +33,15 @@ namespace Volo.Abp.EventBus.Rebus
IOptions<AbpDistributedEventBusOptions> abpDistributedEventBusOptions, IOptions<AbpDistributedEventBusOptions> abpDistributedEventBusOptions,
IOptions<AbpRebusEventBusOptions> abpEventBusRebusOptions, IOptions<AbpRebusEventBusOptions> abpEventBusRebusOptions,
IEventErrorHandler errorHandler) : IEventErrorHandler errorHandler) :
base(serviceScopeFactory, currentTenant, unitOfWorkManager, errorHandler) base(
serviceScopeFactory,
currentTenant,
unitOfWorkManager,
errorHandler,
abpDistributedEventBusOptions)
{ {
Rebus = rebus; Rebus = rebus;
AbpRebusEventBusOptions = abpEventBusRebusOptions.Value; AbpRebusEventBusOptions = abpEventBusRebusOptions.Value;
AbpDistributedEventBusOptions = abpDistributedEventBusOptions.Value;
HandlerFactories = new ConcurrentDictionary<Type, List<IEventHandlerFactory>>(); HandlerFactories = new ConcurrentDictionary<Type, List<IEventHandlerFactory>>();
EventTypes = new ConcurrentDictionary<string, Type>(); EventTypes = new ConcurrentDictionary<string, Type>();

8
framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/AbpDistributedEventBusOptions.cs

@ -1,3 +1,4 @@
using System.Collections.Generic;
using Volo.Abp.Collections; using Volo.Abp.Collections;
namespace Volo.Abp.EventBus.Distributed namespace Volo.Abp.EventBus.Distributed
@ -5,10 +6,17 @@ namespace Volo.Abp.EventBus.Distributed
public class AbpDistributedEventBusOptions public class AbpDistributedEventBusOptions
{ {
public ITypeList<IEventHandler> Handlers { get; } public ITypeList<IEventHandler> Handlers { get; }
public List<OutboxConfig> Outboxes { get; }
public AbpDistributedEventBusOptions() public AbpDistributedEventBusOptions()
{ {
Handlers = new TypeList<IEventHandler>(); Handlers = new TypeList<IEventHandler>();
} }
} }
public class OutboxConfig
{
}
} }

7
framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/DistributedEventBusBase.cs

@ -1,6 +1,7 @@
using System; using System;
using System.Threading.Tasks; using System.Threading.Tasks;
using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Options;
using Volo.Abp.MultiTenancy; using Volo.Abp.MultiTenancy;
using Volo.Abp.Uow; using Volo.Abp.Uow;
@ -8,17 +9,21 @@ namespace Volo.Abp.EventBus.Distributed
{ {
public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventBus public abstract class DistributedEventBusBase : EventBusBase, IDistributedEventBus
{ {
protected AbpDistributedEventBusOptions AbpDistributedEventBusOptions { get; }
protected DistributedEventBusBase( protected DistributedEventBusBase(
IServiceScopeFactory serviceScopeFactory, IServiceScopeFactory serviceScopeFactory,
ICurrentTenant currentTenant, ICurrentTenant currentTenant,
IUnitOfWorkManager unitOfWorkManager, IUnitOfWorkManager unitOfWorkManager,
IEventErrorHandler errorHandler IEventErrorHandler errorHandler,
IOptions<AbpDistributedEventBusOptions> abpDistributedEventBusOptions
) : base( ) : base(
serviceScopeFactory, serviceScopeFactory,
currentTenant, currentTenant,
unitOfWorkManager, unitOfWorkManager,
errorHandler) errorHandler)
{ {
AbpDistributedEventBusOptions = abpDistributedEventBusOptions.Value;
} }
public IDisposable Subscribe<TEvent>(IDistributedEventHandler<TEvent> handler) where TEvent : class public IDisposable Subscribe<TEvent>(IDistributedEventHandler<TEvent> handler) where TEvent : class

Loading…
Cancel
Save