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 b747a858d2..ad857ad967 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 @@ -6,7 +6,6 @@ using System.Threading.Tasks; using Confluent.Kafka; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Options; -using Volo.Abp.Data; using Volo.Abp.DependencyInjection; using Volo.Abp.EventBus.Distributed; using Volo.Abp.Guids; @@ -22,7 +21,6 @@ namespace Volo.Abp.EventBus.Kafka [ExposeServices(typeof(IDistributedEventBus), typeof(KafkaDistributedEventBus))] public class KafkaDistributedEventBus : DistributedEventBusBase, ISingletonDependency { - protected AbpEventBusOptions AbpEventBusOptions { get; } protected AbpKafkaEventBusOptions AbpKafkaEventBusOptions { get; } protected IKafkaMessageConsumerFactory MessageConsumerFactory { get; } protected IKafkaSerializer Serializer { get; } @@ -30,7 +28,6 @@ namespace Volo.Abp.EventBus.Kafka protected ConcurrentDictionary> HandlerFactories { get; } protected ConcurrentDictionary EventTypes { get; } protected IKafkaMessageConsumer Consumer { get; private set; } - protected string DeadLetterTopicName { get; } public KafkaDistributedEventBus( IServiceScopeFactory serviceScopeFactory, @@ -41,26 +38,20 @@ namespace Volo.Abp.EventBus.Kafka IOptions abpDistributedEventBusOptions, IKafkaSerializer serializer, IProducerPool producerPool, - IEventErrorHandler errorHandler, - IOptions abpEventBusOptions, IGuidGenerator guidGenerator, IClock clock) : base( serviceScopeFactory, currentTenant, unitOfWorkManager, - errorHandler, abpDistributedEventBusOptions, guidGenerator, clock) { AbpKafkaEventBusOptions = abpKafkaEventBusOptions.Value; - AbpEventBusOptions = abpEventBusOptions.Value; MessageConsumerFactory = messageConsumerFactory; Serializer = serializer; ProducerPool = producerPool; - DeadLetterTopicName = - AbpEventBusOptions.DeadLetterName ?? AbpKafkaEventBusOptions.TopicName + "_dead_letter"; HandlerFactories = new ConcurrentDictionary>(); EventTypes = new ConcurrentDictionary(); @@ -70,7 +61,6 @@ namespace Volo.Abp.EventBus.Kafka { Consumer = MessageConsumerFactory.Create( AbpKafkaEventBusOptions.TopicName, - DeadLetterTopicName, AbpKafkaEventBusOptions.GroupId, AbpKafkaEventBusOptions.ConnectionName); Consumer.OnMessageReceived(ProcessEventAsync); @@ -88,12 +78,12 @@ namespace Volo.Abp.EventBus.Kafka } string messageId = null; - + if (message.Headers.TryGetLastBytes("messageId", out var messageIdBytes)) { messageId = System.Text.Encoding.UTF8.GetString(messageIdBytes); } - + if (await AddToInboxAsync(messageId, eventName, eventType, message.Value)) { return; @@ -101,18 +91,7 @@ namespace Volo.Abp.EventBus.Kafka var eventData = Serializer.Deserialize(message.Value, eventType); - await TriggerHandlersAsync(eventType, eventData, errorContext => - { - var retryAttempt = 0; - if (message.Headers.TryGetLastBytes(EventErrorHandlerBase.RetryAttemptKey, out var retryAttemptBytes)) - { - retryAttempt = Serializer.Deserialize(retryAttemptBytes); - } - - errorContext.EventData = Serializer.Deserialize(message.Value, eventType); - errorContext.SetProperty(EventErrorHandlerBase.HeadersKey, message.Headers); - errorContext.SetProperty(EventErrorHandlerBase.RetryAttemptKey, retryAttempt); - }); + await TriggerHandlersAsync(eventType, eventData); } public override IDisposable Subscribe(Type eventType, IEventHandlerFactory factory) @@ -226,7 +205,7 @@ namespace Volo.Abp.EventBus.Kafka { return; } - + var eventData = Serializer.Deserialize(incomingEvent.EventData, eventType); var exceptions = new List(); await TriggerHandlersAsync(eventType, eventData, exceptions, inboxConfig); @@ -252,11 +231,6 @@ namespace Volo.Abp.EventBus.Kafka ); } - public virtual async Task PublishToDeadLetterAsync(Type eventType, object eventData, Headers headers, Dictionary headersArguments) - { - await PublishAsync(DeadLetterTopicName, eventType, eventData, headers, headersArguments); - } - private Task PublishAsync(string topicName, Type eventType, object eventData, Headers headers, Dictionary headersArguments) { var eventName = EventNameAttribute.GetNameOrDefault(eventType); @@ -264,7 +238,7 @@ namespace Volo.Abp.EventBus.Kafka return PublishAsync(topicName, eventName, body, headers, headersArguments); } - + private async Task PublishAsync(string topicName, string eventName, byte[] body, Headers headers, Dictionary headersArguments) { var producer = ProducerPool.Get(AbpKafkaEventBusOptions.ConnectionName); diff --git a/framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaEventErrorHandler.cs b/framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaEventErrorHandler.cs deleted file mode 100644 index aee21f75a9..0000000000 --- a/framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaEventErrorHandler.cs +++ /dev/null @@ -1,53 +0,0 @@ -using System; -using System.Collections.Generic; -using System.Linq; -using System.Threading.Tasks; -using Confluent.Kafka; -using Microsoft.Extensions.Logging; -using Microsoft.Extensions.Logging.Abstractions; -using Microsoft.Extensions.Options; -using Volo.Abp.Data; -using Volo.Abp.DependencyInjection; - -namespace Volo.Abp.EventBus.Kafka -{ - public class KafkaEventErrorHandler : EventErrorHandlerBase, ISingletonDependency - { - protected ILogger Logger { get; set; } - - public KafkaEventErrorHandler( - IOptions options) : base(options) - { - Logger = NullLogger.Instance; - } - - protected override async Task RetryAsync(EventExecutionErrorContext context) - { - if (Options.RetryStrategyOptions.IntervalMillisecond > 0) - { - await Task.Delay(Options.RetryStrategyOptions.IntervalMillisecond); - } - - context.TryGetRetryAttempt(out var retryAttempt); - - await context.EventBus.As().PublishAsync( - context.EventType, - context.EventData, - context.GetProperty(HeadersKey).As(), - new Dictionary {{RetryAttemptKey, ++retryAttempt}}); - } - - protected override async Task MoveToDeadLetterAsync(EventExecutionErrorContext context) - { - Logger.LogException( - context.Exceptions.Count == 1 ? context.Exceptions.First() : new AggregateException(context.Exceptions), - LogLevel.Error); - - await context.EventBus.As().PublishToDeadLetterAsync( - context.EventType, - context.EventData, - context.GetProperty(HeadersKey).As(), - new Dictionary {{"exceptions", context.Exceptions.Select(x => x.ToString()).ToList()}}); - } - } -} 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 ceba890234..3188b9bcf1 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 @@ -7,7 +7,6 @@ using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Options; using RabbitMQ.Client; using RabbitMQ.Client.Events; -using Volo.Abp.Data; using Volo.Abp.DependencyInjection; using Volo.Abp.EventBus.Distributed; using Volo.Abp.Guids; @@ -26,7 +25,6 @@ namespace Volo.Abp.EventBus.RabbitMq public class RabbitMqDistributedEventBus : DistributedEventBusBase, ISingletonDependency { protected AbpRabbitMqEventBusOptions AbpRabbitMqEventBusOptions { get; } - protected AbpEventBusOptions AbpEventBusOptions { get; } protected IConnectionPool ConnectionPool { get; } protected IRabbitMqSerializer Serializer { get; } @@ -45,15 +43,12 @@ namespace Volo.Abp.EventBus.RabbitMq IRabbitMqMessageConsumerFactory messageConsumerFactory, ICurrentTenant currentTenant, IUnitOfWorkManager unitOfWorkManager, - IEventErrorHandler errorHandler, - IOptions abpEventBusOptions, IGuidGenerator guidGenerator, IClock clock) : base( - serviceScopeFactory, + serviceScopeFactory, currentTenant, unitOfWorkManager, - errorHandler, distributedEventBusOptions, guidGenerator, clock) @@ -61,7 +56,6 @@ namespace Volo.Abp.EventBus.RabbitMq ConnectionPool = connectionPool; Serializer = serializer; MessageConsumerFactory = messageConsumerFactory; - AbpEventBusOptions = abpEventBusOptions.Value; AbpRabbitMqEventBusOptions = options.Value; HandlerFactories = new ConcurrentDictionary>(); @@ -70,21 +64,17 @@ namespace Volo.Abp.EventBus.RabbitMq public void Initialize() { - const string suffix = "_dead_letter"; - Consumer = MessageConsumerFactory.Create( new ExchangeDeclareConfiguration( AbpRabbitMqEventBusOptions.ExchangeName, type: "direct", - durable: true, - deadLetterExchangeName: AbpRabbitMqEventBusOptions.ExchangeName + suffix + durable: true ), new QueueDeclareConfiguration( AbpRabbitMqEventBusOptions.ClientName, durable: true, exclusive: false, - autoDelete: false, - AbpEventBusOptions.DeadLetterName ?? AbpRabbitMqEventBusOptions.ClientName + suffix + autoDelete: false ), AbpRabbitMqEventBusOptions.ConnectionName ); @@ -104,27 +94,15 @@ namespace Volo.Abp.EventBus.RabbitMq } var eventBytes = ea.Body.ToArray(); - + if (await AddToInboxAsync(ea.BasicProperties.MessageId, eventName, eventType, eventBytes)) { return; } - - var eventData = Serializer.Deserialize(eventBytes, eventType); - await TriggerHandlersAsync(eventType, eventData, errorContext => - { - var retryAttempt = 0; - if (ea.BasicProperties.Headers != null && - ea.BasicProperties.Headers.ContainsKey(EventErrorHandlerBase.RetryAttemptKey)) - { - retryAttempt = (int)ea.BasicProperties.Headers[EventErrorHandlerBase.RetryAttemptKey]; - } + var eventData = Serializer.Deserialize(eventBytes, eventType); - errorContext.EventData = Serializer.Deserialize(eventBytes, eventType); - errorContext.SetProperty(EventErrorHandlerBase.HeadersKey, ea.BasicProperties); - errorContext.SetProperty(EventErrorHandlerBase.RetryAttemptKey, retryAttempt); - }); + await TriggerHandlersAsync(eventType, eventData); } public override IDisposable Subscribe(Type eventType, IEventHandlerFactory factory) @@ -226,7 +204,7 @@ namespace Volo.Abp.EventBus.RabbitMq { return; } - + var eventData = Serializer.Deserialize(incomingEvent.EventData, eventType); var exceptions = new List(); await TriggerHandlersAsync(eventType, eventData, exceptions, inboxConfig); @@ -235,7 +213,7 @@ namespace Volo.Abp.EventBus.RabbitMq ThrowOriginalExceptions(eventType, exceptions); } } - + protected override byte[] Serialize(object eventData) { return Serializer.Serialize(eventData); @@ -248,7 +226,7 @@ namespace Volo.Abp.EventBus.RabbitMq return PublishAsync(eventName, body, properties, headersArguments); } - + protected Task PublishAsync( string eventName, byte[] body, @@ -274,7 +252,7 @@ namespace Volo.Abp.EventBus.RabbitMq { properties.MessageId = (eventId ?? GuidGenerator.Create()).ToString("N"); } - + SetEventMessageHeaders(properties, headersArguments); channel.BasicPublish( diff --git a/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqEventErrorHandler.cs b/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqEventErrorHandler.cs deleted file mode 100644 index e8848bee72..0000000000 --- a/framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqEventErrorHandler.cs +++ /dev/null @@ -1,47 +0,0 @@ -using System; -using System.Collections.Generic; -using System.Linq; -using System.Threading.Tasks; -using Microsoft.Extensions.Options; -using RabbitMQ.Client; -using Volo.Abp.Data; -using Volo.Abp.DependencyInjection; - -namespace Volo.Abp.EventBus.RabbitMq -{ - public class RabbitMqEventErrorHandler : EventErrorHandlerBase, ISingletonDependency - { - public RabbitMqEventErrorHandler( - IOptions options) - : base(options) - { - } - - protected override async Task RetryAsync(EventExecutionErrorContext context) - { - if (Options.RetryStrategyOptions.IntervalMillisecond > 0) - { - await Task.Delay(Options.RetryStrategyOptions.IntervalMillisecond); - } - - context.TryGetRetryAttempt(out var retryAttempt); - - await context.EventBus.As().PublishAsync( - context.EventType, - context.EventData, - context.GetProperty(HeadersKey).As(), - new Dictionary - { - {RetryAttemptKey, ++retryAttempt}, - {"exceptions", context.Exceptions.Select(x => x.ToString()).ToList()} - }); - } - - protected override Task MoveToDeadLetterAsync(EventExecutionErrorContext context) - { - ThrowOriginalExceptions(context); - - return Task.CompletedTask; - } - } -} diff --git a/framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/AbpEventBusRebusModule.cs b/framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/AbpEventBusRebusModule.cs index a13f964d3f..1ad835c122 100644 --- a/framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/AbpEventBusRebusModule.cs +++ b/framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/AbpEventBusRebusModule.cs @@ -1,6 +1,5 @@ using Microsoft.Extensions.DependencyInjection; using Rebus.Handlers; -using Rebus.Retry.Simple; using Rebus.ServiceProvider; using Volo.Abp.Modularity; @@ -12,7 +11,6 @@ namespace Volo.Abp.EventBus.Rebus { public override void ConfigureServices(ServiceConfigurationContext context) { - var abpEventBusOptions = context.Services.ExecutePreConfiguredActions(); var options = context.Services.ExecutePreConfiguredActions();; context.Services.AddTransient(typeof(IHandleMessages<>), typeof(RebusDistributedEventHandlerAdapter<>)); @@ -24,14 +22,6 @@ namespace Volo.Abp.EventBus.Rebus context.Services.AddRebus(configure => { - if (abpEventBusOptions.RetryStrategyOptions != null) - { - configure.Options(b => - b.SimpleRetryStrategy( - errorQueueAddress: abpEventBusOptions.DeadLetterName ?? options.InputQueueName + "_dead_letter", - maxDeliveryAttempts: abpEventBusOptions.RetryStrategyOptions.MaxRetryAttempts)); - } - options.Configurer?.Invoke(configure); return configure; }); 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 1cb380679f..c0d75330cb 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 @@ -36,7 +36,6 @@ namespace Volo.Abp.EventBus.Rebus IBus rebus, IOptions abpDistributedEventBusOptions, IOptions abpEventBusRebusOptions, - IEventErrorHandler errorHandler, IRebusSerializer serializer, IGuidGenerator guidGenerator, IClock clock) : @@ -44,7 +43,6 @@ namespace Volo.Abp.EventBus.Rebus serviceScopeFactory, currentTenant, unitOfWorkManager, - errorHandler, abpDistributedEventBusOptions, guidGenerator, clock) diff --git a/framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusEventErrorHandler.cs b/framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusEventErrorHandler.cs deleted file mode 100644 index 8fe6a53dbb..0000000000 --- a/framework/src/Volo.Abp.EventBus.Rebus/Volo/Abp/EventBus/Rebus/RebusEventErrorHandler.cs +++ /dev/null @@ -1,32 +0,0 @@ -using System.Threading.Tasks; -using Microsoft.Extensions.Options; -using Volo.Abp.DependencyInjection; - -namespace Volo.Abp.EventBus.Rebus -{ - /// - /// Rebus will automatic retries and error handling: https://github.com/rebus-org/Rebus/wiki/Automatic-retries-and-error-handling - /// - public class RebusEventErrorHandler : EventErrorHandlerBase, ISingletonDependency - { - public RebusEventErrorHandler( - IOptions options) - : base(options) - { - } - - protected override Task RetryAsync(EventExecutionErrorContext context) - { - ThrowOriginalExceptions(context); - - return Task.CompletedTask; - } - - protected override Task MoveToDeadLetterAsync(EventExecutionErrorContext context) - { - ThrowOriginalExceptions(context); - - return Task.CompletedTask; - } - } -} diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusModule.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusModule.cs index bd86f14703..3bbdf3d72b 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusModule.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusModule.cs @@ -25,14 +25,6 @@ namespace Volo.Abp.EventBus AddEventHandlers(context.Services); } - public override void ConfigureServices(ServiceConfigurationContext context) - { - Configure(options => - { - context.Services.ExecutePreConfiguredActions(options); - }); - } - private static void AddEventHandlers(IServiceCollection services) { var localHandlers = new List(); diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusOptions.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusOptions.cs deleted file mode 100644 index 39631c7e18..0000000000 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusOptions.cs +++ /dev/null @@ -1,22 +0,0 @@ -using System; - -namespace Volo.Abp.EventBus -{ - public class AbpEventBusOptions - { - public bool EnabledErrorHandle { get; set; } - - public Func ErrorHandleSelector { get; set; } - - public string DeadLetterName { get; set; } - - public AbpEventBusRetryStrategyOptions RetryStrategyOptions { get; set; } - - public void UseRetryStrategy(Action action = null) - { - EnabledErrorHandle = true; - RetryStrategyOptions = new AbpEventBusRetryStrategyOptions(); - action?.Invoke(RetryStrategyOptions); - } - } -} diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusRetryStrategyOptions.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusRetryStrategyOptions.cs deleted file mode 100644 index 4b5b722e96..0000000000 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusRetryStrategyOptions.cs +++ /dev/null @@ -1,9 +0,0 @@ -namespace Volo.Abp.EventBus -{ - public class AbpEventBusRetryStrategyOptions - { - public int IntervalMillisecond { get; set; } = 3000; - - public int MaxRetryAttempts { get; set; } = 3; - } -} 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 70c5bd5533..6745327f7d 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 @@ -20,15 +20,13 @@ namespace Volo.Abp.EventBus.Distributed IServiceScopeFactory serviceScopeFactory, ICurrentTenant currentTenant, IUnitOfWorkManager unitOfWorkManager, - IEventErrorHandler errorHandler, IOptions abpDistributedEventBusOptions, IGuidGenerator guidGenerator, IClock clock ) : base( serviceScopeFactory, currentTenant, - unitOfWorkManager, - errorHandler) + unitOfWorkManager) { GuidGenerator = guidGenerator; Clock = clock; @@ -84,7 +82,7 @@ namespace Volo.Abp.EventBus.Distributed OutgoingEventInfo outgoingEvent, OutboxConfig outboxConfig ); - + public abstract Task ProcessFromInboxAsync( IncomingEventInfo incomingEvent, InboxConfig inboxConfig); @@ -144,7 +142,7 @@ namespace Volo.Abp.EventBus.Distributed continue; } } - + await eventInbox.EnqueueAsync( new IncomingEventInfo( GuidGenerator.Create(), @@ -163,4 +161,4 @@ namespace Volo.Abp.EventBus.Distributed protected abstract byte[] Serialize(object eventData); } -} \ 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 694166b1bc..46d6ecce3c 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs @@ -22,18 +22,14 @@ namespace Volo.Abp.EventBus protected IUnitOfWorkManager UnitOfWorkManager { get; } - protected IEventErrorHandler ErrorHandler { get; } - protected EventBusBase( IServiceScopeFactory serviceScopeFactory, ICurrentTenant currentTenant, - IUnitOfWorkManager unitOfWorkManager, - IEventErrorHandler errorHandler) + IUnitOfWorkManager unitOfWorkManager) { ServiceScopeFactory = serviceScopeFactory; CurrentTenant = currentTenant; UnitOfWorkManager = unitOfWorkManager; - ErrorHandler = errorHandler; } /// @@ -120,7 +116,7 @@ namespace Volo.Abp.EventBus protected abstract void AddToUnitOfWork(IUnitOfWork unitOfWork, UnitOfWorkEventRecord eventRecord); - public virtual async Task TriggerHandlersAsync(Type eventType, object eventData, Action onErrorAction = null) + public virtual async Task TriggerHandlersAsync(Type eventType, object eventData) { var exceptions = new List(); @@ -128,9 +124,7 @@ namespace Volo.Abp.EventBus if (exceptions.Any()) { - var context = new EventExecutionErrorContext(exceptions, eventType, this); - onErrorAction?.Invoke(context); - await ErrorHandler.HandleAsync(context); + ThrowOriginalExceptions(eventType, exceptions); } } @@ -162,7 +156,7 @@ namespace Volo.Abp.EventBus } } } - + protected void ThrowOriginalExceptions(Type eventType, List exceptions) { if (exceptions.Count == 1) diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventErrorHandlerBase.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventErrorHandlerBase.cs deleted file mode 100644 index ba4527fc70..0000000000 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventErrorHandlerBase.cs +++ /dev/null @@ -1,76 +0,0 @@ -using System; -using System.Threading.Tasks; -using Microsoft.Extensions.Options; - -namespace Volo.Abp.EventBus -{ - public abstract class EventErrorHandlerBase : IEventErrorHandler - { - public const string HeadersKey = "headers"; - public const string RetryAttemptKey = "retryAttempt"; - - protected AbpEventBusOptions Options { get; } - - protected EventErrorHandlerBase(IOptions options) - { - Options = options.Value; - } - - public virtual async Task HandleAsync(EventExecutionErrorContext context) - { - if (!await ShouldHandleAsync(context)) - { - ThrowOriginalExceptions(context); - } - - if (await ShouldRetryAsync(context)) - { - await RetryAsync(context); - return; - } - - await MoveToDeadLetterAsync(context); - } - - protected abstract Task RetryAsync(EventExecutionErrorContext context); - - protected abstract Task MoveToDeadLetterAsync(EventExecutionErrorContext context); - - protected virtual Task ShouldHandleAsync(EventExecutionErrorContext context) - { - if (!Options.EnabledErrorHandle) - { - return Task.FromResult(false); - } - - return Task.FromResult(Options.ErrorHandleSelector == null || Options.ErrorHandleSelector.Invoke(context.EventType)); - } - - protected virtual Task ShouldRetryAsync(EventExecutionErrorContext context) - { - if (Options.RetryStrategyOptions == null) - { - return Task.FromResult(false); - } - - if (!context.TryGetRetryAttempt(out var retryAttempt)) - { - return Task.FromResult(false); - } - - return Task.FromResult(Options.RetryStrategyOptions.MaxRetryAttempts > retryAttempt); - } - - protected virtual void ThrowOriginalExceptions(EventExecutionErrorContext context) - { - if (context.Exceptions.Count == 1) - { - context.Exceptions[0].ReThrow(); - } - - throw new AggregateException( - "More than one error has occurred while triggering the event: " + context.EventType, - context.Exceptions); - } - } -} diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventExecutionErrorContext.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventExecutionErrorContext.cs deleted file mode 100644 index e192f61cbe..0000000000 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventExecutionErrorContext.cs +++ /dev/null @@ -1,38 +0,0 @@ -using System; -using System.Collections.Generic; -using Volo.Abp.Data; -using Volo.Abp.ObjectExtending; - -namespace Volo.Abp.EventBus -{ - public class EventExecutionErrorContext : ExtensibleObject - { - public IReadOnlyList Exceptions { get; } - - public object EventData { get; set; } - - public Type EventType { get; } - - public IEventBus EventBus { get; } - - public EventExecutionErrorContext(List exceptions, Type eventType, IEventBus eventBus) - { - Exceptions = exceptions; - EventType = eventType; - EventBus = eventBus; - } - - public bool TryGetRetryAttempt(out int retryAttempt) - { - retryAttempt = 0; - if (!this.HasProperty(EventErrorHandlerBase.RetryAttemptKey)) - { - return false; - } - - retryAttempt = this.GetProperty(EventErrorHandlerBase.RetryAttemptKey); - return true; - - } - } -} diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/IEventErrorHandler.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/IEventErrorHandler.cs deleted file mode 100644 index f1b4a40f15..0000000000 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/IEventErrorHandler.cs +++ /dev/null @@ -1,9 +0,0 @@ -using System.Threading.Tasks; - -namespace Volo.Abp.EventBus -{ - public interface IEventErrorHandler - { - Task HandleAsync(EventExecutionErrorContext context); - } -} 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 77bf5b43bd..4da09bf608 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 @@ -7,11 +7,9 @@ using System.Collections.Generic; using System.Linq; using System.Threading.Tasks; using Microsoft.Extensions.DependencyInjection; -using Volo.Abp.Data; using Volo.Abp.DependencyInjection; using Volo.Abp.MultiTenancy; using Volo.Abp.Threading; -using Volo.Abp.Json; using Volo.Abp.Uow; namespace Volo.Abp.EventBus.Local @@ -35,9 +33,8 @@ namespace Volo.Abp.EventBus.Local IOptions options, IServiceScopeFactory serviceScopeFactory, ICurrentTenant currentTenant, - IUnitOfWorkManager unitOfWorkManager, - IEventErrorHandler errorHandler) - : base(serviceScopeFactory, currentTenant, unitOfWorkManager, errorHandler) + IUnitOfWorkManager unitOfWorkManager) + : base(serviceScopeFactory, currentTenant, unitOfWorkManager) { Options = options.Value; Logger = NullLogger.Instance; @@ -134,11 +131,7 @@ namespace Volo.Abp.EventBus.Local public virtual async Task PublishAsync(LocalEventMessage localEventMessage) { - await TriggerHandlersAsync(localEventMessage.EventType, localEventMessage.EventData, errorContext => - { - errorContext.EventData = localEventMessage.EventData; - errorContext.SetProperty(nameof(LocalEventMessage.MessageId), localEventMessage.MessageId); - }); + await TriggerHandlersAsync(localEventMessage.EventType, localEventMessage.EventData); } protected override IEnumerable GetHandlerFactories(Type eventType) diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventErrorHandler.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventErrorHandler.cs deleted file mode 100644 index c08bca019f..0000000000 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventErrorHandler.cs +++ /dev/null @@ -1,60 +0,0 @@ -using System; -using System.Collections.Generic; -using System.Threading.Tasks; -using Microsoft.Extensions.Options; -using Volo.Abp.Data; -using Volo.Abp.DependencyInjection; - -namespace Volo.Abp.EventBus.Local -{ - [ExposeServices(typeof(LocalEventErrorHandler), typeof(IEventErrorHandler))] - public class LocalEventErrorHandler : EventErrorHandlerBase, ISingletonDependency - { - protected Dictionary RetryTracking { get; } - - public LocalEventErrorHandler( - IOptions options) - : base(options) - { - RetryTracking = new Dictionary(); - } - - protected override async Task RetryAsync(EventExecutionErrorContext context) - { - if (Options.RetryStrategyOptions.IntervalMillisecond > 0) - { - await Task.Delay(Options.RetryStrategyOptions.IntervalMillisecond); - } - - var messageId = context.GetProperty(nameof(LocalEventMessage.MessageId)); - - context.TryGetRetryAttempt(out var retryAttempt); - RetryTracking[messageId] = ++retryAttempt; - - await context.EventBus.As().PublishAsync(new LocalEventMessage(messageId, context.EventData, context.EventType)); - - RetryTracking.Remove(messageId); - } - - protected override Task MoveToDeadLetterAsync(EventExecutionErrorContext context) - { - ThrowOriginalExceptions(context); - - return Task.CompletedTask; - } - - protected override async Task ShouldRetryAsync(EventExecutionErrorContext context) - { - var messageId = context.GetProperty(nameof(LocalEventMessage.MessageId)); - context.SetProperty(RetryAttemptKey, RetryTracking.GetOrDefault(messageId)); - - if (await base.ShouldRetryAsync(context)) - { - return true; - } - - RetryTracking.Remove(messageId); - return false; - } - } -} diff --git a/framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/KafkaMessageConsumer.cs b/framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/KafkaMessageConsumer.cs index d1bad00533..afed026402 100644 --- a/framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/KafkaMessageConsumer.cs +++ b/framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/KafkaMessageConsumer.cs @@ -39,8 +39,6 @@ namespace Volo.Abp.Kafka protected string TopicName { get; private set; } - protected string DeadLetterTopicName { get; private set; } - public KafkaMessageConsumer( IConsumerPool consumerPool, IExceptionNotifier exceptionNotifier, @@ -64,15 +62,12 @@ namespace Volo.Abp.Kafka public virtual void Initialize( [NotNull] string topicName, - [NotNull] string deadLetterTopicName, [NotNull] string groupId, string connectionName = null) { Check.NotNull(topicName, nameof(topicName)); - Check.NotNull(deadLetterTopicName, nameof(deadLetterTopicName)); Check.NotNull(groupId, nameof(groupId)); TopicName = topicName; - DeadLetterTopicName = deadLetterTopicName; ConnectionName = connectionName ?? KafkaConnections.DefaultConnectionName; GroupId = groupId; Timer.Start(); @@ -94,30 +89,18 @@ namespace Volo.Abp.Kafka { using (var adminClient = new AdminClientBuilder(Options.Connections.GetOrDefault(ConnectionName)).Build()) { - var topics = new List + var topic = new TopicSpecification { - new() - { - Name = TopicName, - NumPartitions = 1, - ReplicationFactor = 1 - }, - new() - { - Name = DeadLetterTopicName, - NumPartitions = 1, - ReplicationFactor = 1 - } + Name = TopicName, + NumPartitions = 1, + ReplicationFactor = 1 }; - topics.ForEach(topic => - { - Options.ConfigureTopic?.Invoke(topic); - }); + Options.ConfigureTopic?.Invoke(topic); try { - await adminClient.CreateTopicsAsync(topics); + await adminClient.CreateTopicsAsync(new[] {topic}); } catch (CreateTopicsException e) { diff --git a/framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/KafkaMessageConsumerFactory.cs b/framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/KafkaMessageConsumerFactory.cs index 4a22fd04f6..68d1162b7f 100644 --- a/framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/KafkaMessageConsumerFactory.cs +++ b/framework/src/Volo.Abp.Kafka/Volo/Abp/Kafka/KafkaMessageConsumerFactory.cs @@ -21,7 +21,7 @@ namespace Volo.Abp.Kafka string connectionName = null) { var consumer = ServiceScope.ServiceProvider.GetRequiredService(); - consumer.Initialize(topicName, deadLetterTopicName, groupId, connectionName); + consumer.Initialize(topicName, groupId, connectionName); return consumer; } diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ExchangeDeclareConfiguration.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ExchangeDeclareConfiguration.cs index b9e762abbe..8ea919484a 100644 --- a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ExchangeDeclareConfiguration.cs +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/ExchangeDeclareConfiguration.cs @@ -6,8 +6,6 @@ namespace Volo.Abp.RabbitMQ { public string ExchangeName { get; } - public string DeadLetterExchangeName { get; set; } - public string Type { get; } public bool Durable { get; set; } @@ -20,11 +18,9 @@ namespace Volo.Abp.RabbitMQ string exchangeName, string type, bool durable = false, - bool autoDelete = false, - string deadLetterExchangeName = null) + bool autoDelete = false) { ExchangeName = exchangeName; - DeadLetterExchangeName = deadLetterExchangeName; Type = type; Durable = durable; AutoDelete = autoDelete; diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/QueueDeclareConfiguration.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/QueueDeclareConfiguration.cs index b84f08ec42..211dc3d7b2 100644 --- a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/QueueDeclareConfiguration.cs +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/QueueDeclareConfiguration.cs @@ -8,8 +8,6 @@ namespace Volo.Abp.RabbitMQ { [NotNull] public string QueueName { get; } - public string DeadLetterQueueName { get; set; } - public bool Durable { get; set; } public bool Exclusive { get; set; } @@ -22,11 +20,9 @@ namespace Volo.Abp.RabbitMQ [NotNull] string queueName, bool durable = true, bool exclusive = false, - bool autoDelete = false, - string deadLetterQueueName = null) + bool autoDelete = false) { QueueName = queueName; - DeadLetterQueueName = deadLetterQueueName; Durable = durable; Exclusive = exclusive; AutoDelete = autoDelete; diff --git a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/RabbitMqMessageConsumer.cs b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/RabbitMqMessageConsumer.cs index 671445e00d..a0c1251c6c 100644 --- a/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/RabbitMqMessageConsumer.cs +++ b/framework/src/Volo.Abp.RabbitMQ/Volo/Abp/RabbitMQ/RabbitMqMessageConsumer.cs @@ -157,29 +157,7 @@ namespace Volo.Abp.RabbitMQ arguments: Exchange.Arguments ); - if (!Exchange.DeadLetterExchangeName.IsNullOrWhiteSpace() && - !Queue.DeadLetterQueueName.IsNullOrWhiteSpace()) - { - Channel.ExchangeDeclare( - Exchange.DeadLetterExchangeName, - Exchange.Type, - Exchange.Durable, - Exchange.AutoDelete - ); - - Channel.QueueDeclare( - Queue.DeadLetterQueueName, - Queue.Durable, - Queue.Exclusive, - Queue.AutoDelete); - - Queue.Arguments["x-dead-letter-exchange"] = Exchange.DeadLetterExchangeName; - Queue.Arguments["x-dead-letter-routing-key"] = Queue.DeadLetterQueueName; - - Channel.QueueBind(Queue.DeadLetterQueueName, Exchange.DeadLetterExchangeName, Queue.DeadLetterQueueName); - } - - var result = Channel.QueueDeclare( + Channel.QueueDeclare( queue: Queue.QueueName, durable: Queue.Durable, exclusive: Queue.Exclusive, @@ -202,11 +180,8 @@ namespace Volo.Abp.RabbitMQ operationInterruptedException.ShutdownReason.ReplyCode == 406 && operationInterruptedException.Message.Contains("arg 'x-dead-letter-exchange'")) { - Exchange.DeadLetterExchangeName = null; - Queue.DeadLetterQueueName = null; - Queue.Arguments.Remove("x-dead-letter-exchange"); - Queue.Arguments.Remove("x-dead-letter-routing-key"); - Logger.LogWarning("Unable to bind the dead letter queue to an existing queue. You can delete the queue or add policy. See: https://www.rabbitmq.com/parameters.html"); + Logger.LogException(ex, LogLevel.Warning); + await ExceptionNotifier.NotifyAsync(ex, logLevel: LogLevel.Warning); } Logger.LogException(ex, LogLevel.Warning); @@ -229,8 +204,13 @@ namespace Volo.Abp.RabbitMQ { try { - Channel.BasicReject(basicDeliverEventArgs.DeliveryTag, false); + Channel.BasicNack( + basicDeliverEventArgs.DeliveryTag, + multiple: false, + requeue: true + ); } + // ReSharper disable once EmptyGeneralCatchClause catch { } Logger.LogException(ex); diff --git a/framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/EventBusTestModule.cs b/framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/EventBusTestModule.cs index f260fecbea..9a4258c621 100644 --- a/framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/EventBusTestModule.cs +++ b/framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/EventBusTestModule.cs @@ -5,17 +5,5 @@ namespace Volo.Abp.EventBus [DependsOn(typeof(AbpEventBusModule))] public class EventBusTestModule : AbpModule { - public override void PreConfigureServices(ServiceConfigurationContext context) - { - PreConfigure(options => - { - options.UseRetryStrategy(retryStrategyOptions => - { - retryStrategyOptions.IntervalMillisecond = 0; - }); - - options.ErrorHandleSelector = type => type == typeof(MyExceptionHandleEventData); - }); - } } } diff --git a/framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/Local/EventBus_Exception_Handler_Tests.cs b/framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/Local/EventBus_Exception_Handler_Tests.cs deleted file mode 100644 index 3ed0ff19f9..0000000000 --- a/framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/Local/EventBus_Exception_Handler_Tests.cs +++ /dev/null @@ -1,72 +0,0 @@ -using System; -using System.Threading.Tasks; -using Shouldly; -using Xunit; - -namespace Volo.Abp.EventBus.Local -{ - public class EventBus_Exception_Handler_Tests : EventBusTestBase - { - [Fact] - public async Task Should_Not_Handle_Exception() - { - var retryAttempt = 0; - LocalEventBus.Subscribe(eventData => - { - retryAttempt++; - throw new Exception("This exception is intentionally thrown!"); - }); - - var appException = await Assert.ThrowsAsync(async () => - { - await LocalEventBus.PublishAsync(new MySimpleEventData(1)); - }); - - retryAttempt.ShouldBe(1); - appException.Message.ShouldBe("This exception is intentionally thrown!"); - } - - [Fact] - public async Task Should_Handle_Exception() - { - var retryAttempt = 0; - LocalEventBus.Subscribe(eventData => - { - eventData.Value.ShouldBe(0); - retryAttempt++; - if (retryAttempt < 2) - { - throw new Exception("This exception is intentionally thrown!"); - } - - return Task.CompletedTask; - - }); - - await LocalEventBus.PublishAsync(new MyExceptionHandleEventData(0)); - retryAttempt.ShouldBe(2); - } - - [Fact] - public async Task Should_Throw_Exception_After_Error_Handle() - { - var retryAttempt = 0; - LocalEventBus.Subscribe(eventData => - { - eventData.Value.ShouldBe(0); - - retryAttempt++; - - throw new Exception("This exception is intentionally thrown!"); - }); - - var appException = await Assert.ThrowsAsync(async () => - { - await LocalEventBus.PublishAsync(new MyExceptionHandleEventData(0)); - }); - - retryAttempt.ShouldBe(4); - appException.Message.ShouldBe("This exception is intentionally thrown!"); - } - } -} diff --git a/framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/MyExceptionHandleEventData.cs b/framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/MyExceptionHandleEventData.cs deleted file mode 100644 index f490d58211..0000000000 --- a/framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/MyExceptionHandleEventData.cs +++ /dev/null @@ -1,12 +0,0 @@ -namespace Volo.Abp.EventBus -{ - public class MyExceptionHandleEventData - { - public int Value { get; set; } - - public MyExceptionHandleEventData(int value) - { - Value = value; - } - } -}