From 3d2348134e2890b71e234bc5333141a044ad57ad Mon Sep 17 00:00:00 2001 From: liangshiwei Date: Tue, 27 Apr 2021 16:40:44 +0800 Subject: [PATCH] Add unit tests --- .../EventBus/Kafka/KafkaEventErrorHandler.cs | 22 +++--- .../RabbitMq/RabbitMqEventErrorHandler.cs | 10 +-- .../Volo/Abp/EventBus/AbpEventBusOptions.cs | 1 + .../Volo/Abp/EventBus/EventBusBase.cs | 7 +- .../Abp/EventBus/EventErrorHandlerBase.cs | 12 +++- .../EventBus/EventExecutionErrorContext.cs | 6 +- .../EventBus/Local/LocalEventErrorHandler.cs | 12 ++-- .../Volo/Abp/EventBus/EventBusTestModule.cs | 13 +++- .../Local/EventBus_Exception_Handler_Tests.cs | 72 +++++++++++++++++++ .../EventBus/MyExceptionHandleEventData.cs | 12 ++++ 10 files changed, 133 insertions(+), 34 deletions(-) create mode 100644 framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/Local/EventBus_Exception_Handler_Tests.cs create mode 100644 framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/MyExceptionHandleEventData.cs 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 index b9e46fdd6a..567c303704 100644 --- 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 @@ -1,4 +1,6 @@ -using System.Threading.Tasks; +using System; +using System.Linq; +using System.Threading.Tasks; using Confluent.Kafka; using Microsoft.Extensions.Options; using Volo.Abp.Data; @@ -13,15 +15,12 @@ namespace Volo.Abp.EventBus.Kafka public const string RetryIndexKey = "retryIndex"; protected IKafkaSerializer Serializer { get; } - protected KafkaDistributedEventBus EventBus { get; } public KafkaEventErrorHandler( IOptions options, - IKafkaSerializer serializer, - KafkaDistributedEventBus eventBus) : base(options) + IKafkaSerializer serializer) : base(options) { Serializer = serializer; - EventBus = eventBus; } protected override async Task Retry(EventExecutionErrorContext context) @@ -32,17 +31,22 @@ namespace Volo.Abp.EventBus.Kafka } var headers = context.GetProperty(HeadersKey) ?? new Headers(); - var index = Serializer.Deserialize(headers.GetLastBytes(RetryIndexKey)); + + var index = 1; + if (headers.Any(x => x.Key == RetryIndexKey)) + { + index = Serializer.Deserialize(headers.GetLastBytes(RetryIndexKey)); + } headers.Remove(RetryIndexKey); headers.Add(RetryIndexKey, Serializer.Serialize(++index)); - await EventBus.PublishAsync(context.EventType, context.EventData, headers); + await context.EventBus.As().PublishAsync(context.EventType, context.EventData, headers); } protected override async Task MoveToDeadLetter(EventExecutionErrorContext context) { - await EventBus.PublishToDeadLetterAsync(context.EventType, context.EventData, new Headers + await context.EventBus.As().PublishToDeadLetterAsync(context.EventType, context.EventData, new Headers { {"exceptions", Serializer.Serialize(context.Exceptions)} }); @@ -64,7 +68,7 @@ namespace Volo.Abp.EventBus.Kafka var index = Serializer.Deserialize(headers.GetLastBytes(RetryIndexKey)); - return Options.RetryStrategyOptions.MaxRetryAttempts < index; + return Options.RetryStrategyOptions.MaxRetryAttempts > index; } } } 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 index a17bcfcdcc..b9b0ac6497 100644 --- 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 @@ -12,14 +12,10 @@ namespace Volo.Abp.EventBus.RabbitMq public const string HeadersKey = "headers"; public const string RetryIndexKey = "retryIndex"; - protected RabbitMqDistributedEventBus EventBus { get; } - public RabbitMqEventErrorHandler( - IOptions options, - RabbitMqDistributedEventBus eventBus) + IOptions options) : base(options) { - EventBus = eventBus; } protected override async Task Retry(EventExecutionErrorContext context) @@ -45,7 +41,7 @@ namespace Volo.Abp.EventBus.RabbitMq headers["exceptions"] = context.Exceptions; - await EventBus.PublishAsync(context.EventType, context.EventData, headers); + await context.EventBus.As().PublishAsync(context.EventType, context.EventData, headers); } protected override Task MoveToDeadLetter(EventExecutionErrorContext context) @@ -76,7 +72,7 @@ namespace Volo.Abp.EventBus.RabbitMq var index = (int) headers[RetryIndexKey]; - return Options.RetryStrategyOptions.MaxRetryAttempts < index; + return Options.RetryStrategyOptions.MaxRetryAttempts > index; } } } diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusOptions.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusOptions.cs index 3ba5c2a548..39631c7e18 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusOptions.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusOptions.cs @@ -14,6 +14,7 @@ namespace Volo.Abp.EventBus 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/EventBusBase.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs index e505f30ccb..6dcaa496b4 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs @@ -104,7 +104,7 @@ namespace Volo.Abp.EventBus if (exceptions.Any()) { - var context = new EventExecutionErrorContext(exceptions, eventData, eventType); + var context = new EventExecutionErrorContext(exceptions, eventData, eventType, this); onErrorAction?.Invoke(context); await ErrorHandler.Handle(context); } @@ -221,11 +221,6 @@ namespace Volo.Abp.EventBus }; } - protected virtual void OnErrorHandle(EventExecutionErrorContext context) - { - - } - protected class EventTypeWithEventHandlerFactories { public Type EventType { get; } diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventErrorHandlerBase.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventErrorHandlerBase.cs index 0bb6430e56..b5d22389d3 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventErrorHandlerBase.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventErrorHandlerBase.cs @@ -1,4 +1,5 @@ -using System.Threading.Tasks; +using System; +using System.Threading.Tasks; using Microsoft.Extensions.Options; namespace Volo.Abp.EventBus @@ -16,7 +17,14 @@ namespace Volo.Abp.EventBus { if (!ShouldHandle(context)) { - return; + 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); } if (ShouldRetry(context)) diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventExecutionErrorContext.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventExecutionErrorContext.cs index cf40102579..29c362c1ca 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventExecutionErrorContext.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventExecutionErrorContext.cs @@ -1,6 +1,5 @@ using System; using System.Collections.Generic; -using Volo.Abp.EventBus.Local; using Volo.Abp.ObjectExtending; namespace Volo.Abp.EventBus @@ -13,11 +12,14 @@ namespace Volo.Abp.EventBus public Type EventType { get; } - public EventExecutionErrorContext(List exceptions, object eventData, Type eventType) + public IEventBus EventBus { get; } + + public EventExecutionErrorContext(List exceptions, object eventData, Type eventType, IEventBus eventBus) { Exceptions = exceptions; EventData = eventData; EventType = eventType; + EventBus = eventBus; } } } 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 index da7a0d7f71..79fd405940 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventErrorHandler.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventErrorHandler.cs @@ -9,15 +9,12 @@ namespace Volo.Abp.EventBus.Local { public class LocalEventErrorHandler : EventErrorHandlerBase, ISingletonDependency { - protected ILocalEventBus LocalEventBus { get; } protected Dictionary RetryTracking { get; } public LocalEventErrorHandler( - IOptions options, - ILocalEventBus localEventBus) + IOptions options) : base(options) { - LocalEventBus = localEventBus; RetryTracking = new Dictionary(); } @@ -30,8 +27,9 @@ namespace Volo.Abp.EventBus.Local var messageId = context.GetProperty("messageId"); - await LocalEventBus.PublishAsync(context.EventType, - new LocalEventMessage(messageId, context.EventData, context.EventType)); + await context.EventBus.As().PublishAsync(new LocalEventMessage(messageId, context.EventData, context.EventType)); + + RetryTracking.Remove(messageId); } protected override Task MoveToDeadLetter(EventExecutionErrorContext context) @@ -57,7 +55,7 @@ namespace Volo.Abp.EventBus.Local var index = RetryTracking.GetOrDefault(messageId); - if (Options.RetryStrategyOptions.MaxRetryAttempts >= index) + if (Options.RetryStrategyOptions.MaxRetryAttempts <= index) { RetryTracking.Remove(messageId); return false; 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 f514f37018..f260fecbea 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,6 +5,17 @@ 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); + }); + } } -} \ No newline at end of file +} 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 new file mode 100644 index 0000000000..7482606416 --- /dev/null +++ b/framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/Local/EventBus_Exception_Handler_Tests.cs @@ -0,0 +1,72 @@ +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() + { + MySimpleEventData data = null; + LocalEventBus.Subscribe(eventData => + { + ++eventData.Value; + data = eventData; + throw new Exception("This exception is intentionally thrown!"); + }); + + var appException = await Assert.ThrowsAsync(async () => + { + await LocalEventBus.PublishAsync(new MySimpleEventData(1)); + }); + + data.Value.ShouldBe(2); + appException.Message.ShouldBe("This exception is intentionally thrown!"); + } + + [Fact] + public async Task Should_Handle_Exception() + { + MyExceptionHandleEventData data = null; + LocalEventBus.Subscribe(eventData => + { + ++eventData.RetryAttempts; + data = eventData; + + if (eventData.RetryAttempts < 2) + { + throw new Exception("This exception is intentionally thrown!"); + } + + return Task.CompletedTask; + + }); + + await LocalEventBus.PublishAsync(new MyExceptionHandleEventData(0)); + data.RetryAttempts.ShouldBe(2); + } + + [Fact] + public async Task Should_Throw_Exception_After_Error_Handle() + { + MyExceptionHandleEventData data = null; + LocalEventBus.Subscribe(eventData => + { + ++eventData.RetryAttempts; + data = eventData; + throw new Exception("This exception is intentionally thrown!"); + }); + + var appException = await Assert.ThrowsAsync(async () => + { + await LocalEventBus.PublishAsync(new MyExceptionHandleEventData(0)); + }); + + data.RetryAttempts.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 new file mode 100644 index 0000000000..b68f60aecd --- /dev/null +++ b/framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/MyExceptionHandleEventData.cs @@ -0,0 +1,12 @@ +namespace Volo.Abp.EventBus +{ + public class MyExceptionHandleEventData + { + public int RetryAttempts { get; set; } + + public MyExceptionHandleEventData(int retryAttempts) + { + RetryAttempts = retryAttempts; + } + } +}