Browse Source

Add unit tests

pull/8829/head
liangshiwei 6 years ago
parent
commit
3d2348134e
  1. 22
      framework/src/Volo.Abp.EventBus.Kafka/Volo/Abp/EventBus/Kafka/KafkaEventErrorHandler.cs
  2. 10
      framework/src/Volo.Abp.EventBus.RabbitMQ/Volo/Abp/EventBus/RabbitMq/RabbitMqEventErrorHandler.cs
  3. 1
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusOptions.cs
  4. 7
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs
  5. 12
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventErrorHandlerBase.cs
  6. 6
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventExecutionErrorContext.cs
  7. 12
      framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Local/LocalEventErrorHandler.cs
  8. 13
      framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/EventBusTestModule.cs
  9. 72
      framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/Local/EventBus_Exception_Handler_Tests.cs
  10. 12
      framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/MyExceptionHandleEventData.cs

22
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 Confluent.Kafka;
using Microsoft.Extensions.Options; using Microsoft.Extensions.Options;
using Volo.Abp.Data; using Volo.Abp.Data;
@ -13,15 +15,12 @@ namespace Volo.Abp.EventBus.Kafka
public const string RetryIndexKey = "retryIndex"; public const string RetryIndexKey = "retryIndex";
protected IKafkaSerializer Serializer { get; } protected IKafkaSerializer Serializer { get; }
protected KafkaDistributedEventBus EventBus { get; }
public KafkaEventErrorHandler( public KafkaEventErrorHandler(
IOptions<AbpEventBusOptions> options, IOptions<AbpEventBusOptions> options,
IKafkaSerializer serializer, IKafkaSerializer serializer) : base(options)
KafkaDistributedEventBus eventBus) : base(options)
{ {
Serializer = serializer; Serializer = serializer;
EventBus = eventBus;
} }
protected override async Task Retry(EventExecutionErrorContext context) protected override async Task Retry(EventExecutionErrorContext context)
@ -32,17 +31,22 @@ namespace Volo.Abp.EventBus.Kafka
} }
var headers = context.GetProperty<Headers>(HeadersKey) ?? new Headers(); var headers = context.GetProperty<Headers>(HeadersKey) ?? new Headers();
var index = Serializer.Deserialize<int>(headers.GetLastBytes(RetryIndexKey));
var index = 1;
if (headers.Any(x => x.Key == RetryIndexKey))
{
index = Serializer.Deserialize<int>(headers.GetLastBytes(RetryIndexKey));
}
headers.Remove(RetryIndexKey); headers.Remove(RetryIndexKey);
headers.Add(RetryIndexKey, Serializer.Serialize(++index)); headers.Add(RetryIndexKey, Serializer.Serialize(++index));
await EventBus.PublishAsync(context.EventType, context.EventData, headers); await context.EventBus.As<KafkaDistributedEventBus>().PublishAsync(context.EventType, context.EventData, headers);
} }
protected override async Task MoveToDeadLetter(EventExecutionErrorContext context) protected override async Task MoveToDeadLetter(EventExecutionErrorContext context)
{ {
await EventBus.PublishToDeadLetterAsync(context.EventType, context.EventData, new Headers await context.EventBus.As<KafkaDistributedEventBus>().PublishToDeadLetterAsync(context.EventType, context.EventData, new Headers
{ {
{"exceptions", Serializer.Serialize(context.Exceptions)} {"exceptions", Serializer.Serialize(context.Exceptions)}
}); });
@ -64,7 +68,7 @@ namespace Volo.Abp.EventBus.Kafka
var index = Serializer.Deserialize<int>(headers.GetLastBytes(RetryIndexKey)); var index = Serializer.Deserialize<int>(headers.GetLastBytes(RetryIndexKey));
return Options.RetryStrategyOptions.MaxRetryAttempts < index; return Options.RetryStrategyOptions.MaxRetryAttempts > index;
} }
} }
} }

10
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 HeadersKey = "headers";
public const string RetryIndexKey = "retryIndex"; public const string RetryIndexKey = "retryIndex";
protected RabbitMqDistributedEventBus EventBus { get; }
public RabbitMqEventErrorHandler( public RabbitMqEventErrorHandler(
IOptions<AbpEventBusOptions> options, IOptions<AbpEventBusOptions> options)
RabbitMqDistributedEventBus eventBus)
: base(options) : base(options)
{ {
EventBus = eventBus;
} }
protected override async Task Retry(EventExecutionErrorContext context) protected override async Task Retry(EventExecutionErrorContext context)
@ -45,7 +41,7 @@ namespace Volo.Abp.EventBus.RabbitMq
headers["exceptions"] = context.Exceptions; headers["exceptions"] = context.Exceptions;
await EventBus.PublishAsync(context.EventType, context.EventData, headers); await context.EventBus.As<RabbitMqDistributedEventBus>().PublishAsync(context.EventType, context.EventData, headers);
} }
protected override Task MoveToDeadLetter(EventExecutionErrorContext context) protected override Task MoveToDeadLetter(EventExecutionErrorContext context)
@ -76,7 +72,7 @@ namespace Volo.Abp.EventBus.RabbitMq
var index = (int) headers[RetryIndexKey]; var index = (int) headers[RetryIndexKey];
return Options.RetryStrategyOptions.MaxRetryAttempts < index; return Options.RetryStrategyOptions.MaxRetryAttempts > index;
} }
} }
} }

1
framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/AbpEventBusOptions.cs

@ -14,6 +14,7 @@ namespace Volo.Abp.EventBus
public void UseRetryStrategy(Action<AbpEventBusRetryStrategyOptions> action = null) public void UseRetryStrategy(Action<AbpEventBusRetryStrategyOptions> action = null)
{ {
EnabledErrorHandle = true;
RetryStrategyOptions = new AbpEventBusRetryStrategyOptions(); RetryStrategyOptions = new AbpEventBusRetryStrategyOptions();
action?.Invoke(RetryStrategyOptions); action?.Invoke(RetryStrategyOptions);
} }

7
framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventBusBase.cs

@ -104,7 +104,7 @@ namespace Volo.Abp.EventBus
if (exceptions.Any()) if (exceptions.Any())
{ {
var context = new EventExecutionErrorContext(exceptions, eventData, eventType); var context = new EventExecutionErrorContext(exceptions, eventData, eventType, this);
onErrorAction?.Invoke(context); onErrorAction?.Invoke(context);
await ErrorHandler.Handle(context); await ErrorHandler.Handle(context);
} }
@ -221,11 +221,6 @@ namespace Volo.Abp.EventBus
}; };
} }
protected virtual void OnErrorHandle(EventExecutionErrorContext context)
{
}
protected class EventTypeWithEventHandlerFactories protected class EventTypeWithEventHandlerFactories
{ {
public Type EventType { get; } public Type EventType { get; }

12
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; using Microsoft.Extensions.Options;
namespace Volo.Abp.EventBus namespace Volo.Abp.EventBus
@ -16,7 +17,14 @@ namespace Volo.Abp.EventBus
{ {
if (!ShouldHandle(context)) 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)) if (ShouldRetry(context))

6
framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/EventExecutionErrorContext.cs

@ -1,6 +1,5 @@
using System; using System;
using System.Collections.Generic; using System.Collections.Generic;
using Volo.Abp.EventBus.Local;
using Volo.Abp.ObjectExtending; using Volo.Abp.ObjectExtending;
namespace Volo.Abp.EventBus namespace Volo.Abp.EventBus
@ -13,11 +12,14 @@ namespace Volo.Abp.EventBus
public Type EventType { get; } public Type EventType { get; }
public EventExecutionErrorContext(List<Exception> exceptions, object eventData, Type eventType) public IEventBus EventBus { get; }
public EventExecutionErrorContext(List<Exception> exceptions, object eventData, Type eventType, IEventBus eventBus)
{ {
Exceptions = exceptions; Exceptions = exceptions;
EventData = eventData; EventData = eventData;
EventType = eventType; EventType = eventType;
EventBus = eventBus;
} }
} }
} }

12
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 public class LocalEventErrorHandler : EventErrorHandlerBase, ISingletonDependency
{ {
protected ILocalEventBus LocalEventBus { get; }
protected Dictionary<Guid, int> RetryTracking { get; } protected Dictionary<Guid, int> RetryTracking { get; }
public LocalEventErrorHandler( public LocalEventErrorHandler(
IOptions<AbpEventBusOptions> options, IOptions<AbpEventBusOptions> options)
ILocalEventBus localEventBus)
: base(options) : base(options)
{ {
LocalEventBus = localEventBus;
RetryTracking = new Dictionary<Guid, int>(); RetryTracking = new Dictionary<Guid, int>();
} }
@ -30,8 +27,9 @@ namespace Volo.Abp.EventBus.Local
var messageId = context.GetProperty<Guid>("messageId"); var messageId = context.GetProperty<Guid>("messageId");
await LocalEventBus.PublishAsync(context.EventType, await context.EventBus.As<LocalEventBus>().PublishAsync(new LocalEventMessage(messageId, context.EventData, context.EventType));
new LocalEventMessage(messageId, context.EventData, context.EventType));
RetryTracking.Remove(messageId);
} }
protected override Task MoveToDeadLetter(EventExecutionErrorContext context) protected override Task MoveToDeadLetter(EventExecutionErrorContext context)
@ -57,7 +55,7 @@ namespace Volo.Abp.EventBus.Local
var index = RetryTracking.GetOrDefault(messageId); var index = RetryTracking.GetOrDefault(messageId);
if (Options.RetryStrategyOptions.MaxRetryAttempts >= index) if (Options.RetryStrategyOptions.MaxRetryAttempts <= index)
{ {
RetryTracking.Remove(messageId); RetryTracking.Remove(messageId);
return false; return false;

13
framework/test/Volo.Abp.EventBus.Tests/Volo/Abp/EventBus/EventBusTestModule.cs

@ -5,6 +5,17 @@ namespace Volo.Abp.EventBus
[DependsOn(typeof(AbpEventBusModule))] [DependsOn(typeof(AbpEventBusModule))]
public class EventBusTestModule : AbpModule public class EventBusTestModule : AbpModule
{ {
public override void PreConfigureServices(ServiceConfigurationContext context)
{
PreConfigure<AbpEventBusOptions>(options =>
{
options.UseRetryStrategy(retryStrategyOptions =>
{
retryStrategyOptions.IntervalMillisecond = 0;
});
options.ErrorHandleSelector = type => type == typeof(MyExceptionHandleEventData);
});
}
} }
} }

72
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<MySimpleEventData>(eventData =>
{
++eventData.Value;
data = eventData;
throw new Exception("This exception is intentionally thrown!");
});
var appException = await Assert.ThrowsAsync<Exception>(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<MyExceptionHandleEventData>(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<MyExceptionHandleEventData>(eventData =>
{
++eventData.RetryAttempts;
data = eventData;
throw new Exception("This exception is intentionally thrown!");
});
var appException = await Assert.ThrowsAsync<Exception>(async () =>
{
await LocalEventBus.PublishAsync(new MyExceptionHandleEventData(0));
});
data.RetryAttempts.ShouldBe(4);
appException.Message.ShouldBe("This exception is intentionally thrown!");
}
}
}

12
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;
}
}
}
Loading…
Cancel
Save