|
|
|
@ -1,5 +1,6 @@ |
|
|
|
using System; |
|
|
|
using System.Collections.Generic; |
|
|
|
using System.Linq; |
|
|
|
using System.Threading; |
|
|
|
using System.Threading.Tasks; |
|
|
|
using Microsoft.Extensions.DependencyInjection; |
|
|
|
@ -25,7 +26,6 @@ public class InboxProcessor : IInboxProcessor, ITransientDependency |
|
|
|
protected IEventInbox Inbox { get; private set; } = default!; |
|
|
|
protected InboxConfig InboxConfig { get; private set; } = default!; |
|
|
|
protected AbpEventBusBoxesOptions EventBusBoxesOptions { get; } |
|
|
|
protected AbpInboxProcessorOptions AbpInboxProcessorOptions { get; set; } |
|
|
|
protected DateTime? LastCleanTime { get; set; } |
|
|
|
|
|
|
|
protected string DistributedLockName { get; set; } = default!; |
|
|
|
@ -40,8 +40,7 @@ public class InboxProcessor : IInboxProcessor, ITransientDependency |
|
|
|
IAbpDistributedLock distributedLock, |
|
|
|
IUnitOfWorkManager unitOfWorkManager, |
|
|
|
IClock clock, |
|
|
|
IOptions<AbpEventBusBoxesOptions> eventBusBoxesOptions, |
|
|
|
IOptions<AbpInboxProcessorOptions> inboxProcessorOptions) |
|
|
|
IOptions<AbpEventBusBoxesOptions> eventBusBoxesOptions) |
|
|
|
{ |
|
|
|
ServiceProvider = serviceProvider; |
|
|
|
Timer = timer; |
|
|
|
@ -50,7 +49,6 @@ public class InboxProcessor : IInboxProcessor, ITransientDependency |
|
|
|
UnitOfWorkManager = unitOfWorkManager; |
|
|
|
Clock = clock; |
|
|
|
EventBusBoxesOptions = eventBusBoxesOptions.Value; |
|
|
|
AbpInboxProcessorOptions = inboxProcessorOptions.Value; |
|
|
|
Timer.Period = Convert.ToInt32(EventBusBoxesOptions.PeriodTimeSpan.TotalMilliseconds); |
|
|
|
Timer.Elapsed += TimerOnElapsed; |
|
|
|
Logger = NullLogger<InboxProcessor>.Instance; |
|
|
|
@ -107,12 +105,6 @@ public class InboxProcessor : IInboxProcessor, ITransientDependency |
|
|
|
{ |
|
|
|
Logger.LogInformation($"Start processing the incoming event with id = {waitingEvent.Id:N}"); |
|
|
|
|
|
|
|
if (waitingEvent.NextRetryTime.HasValue && waitingEvent.NextRetryTime.Value > Clock.Now) |
|
|
|
{ |
|
|
|
Logger.LogInformation($"Event with id = {waitingEvent.Id:N} is not ready to be processed yet. Next retry time: {waitingEvent.NextRetryTime.Value}"); |
|
|
|
continue; |
|
|
|
} |
|
|
|
|
|
|
|
try |
|
|
|
{ |
|
|
|
using (var uow = UnitOfWorkManager.Begin(isTransactional: true, requiresNew: true)) |
|
|
|
@ -130,41 +122,42 @@ public class InboxProcessor : IInboxProcessor, ITransientDependency |
|
|
|
} |
|
|
|
catch (Exception e) |
|
|
|
{ |
|
|
|
Logger.LogError(e, $"An error occurred while processing the incoming event with id = {waitingEvent.Id:N}"); |
|
|
|
Logger.LogError(e, $"Event with id = {waitingEvent.Id:N} processing failed."); |
|
|
|
|
|
|
|
if (AbpInboxProcessorOptions.FailurePolicy == InboxProcessorFailurePolicy.Retry) |
|
|
|
if (EventBusBoxesOptions.InboxProcessorFailurePolicy == InboxProcessorFailurePolicy.Retry) |
|
|
|
{ |
|
|
|
throw; |
|
|
|
} |
|
|
|
|
|
|
|
if (AbpInboxProcessorOptions.FailurePolicy == InboxProcessorFailurePolicy.RetryLater) |
|
|
|
if (EventBusBoxesOptions.InboxProcessorFailurePolicy == InboxProcessorFailurePolicy.RetryLater) |
|
|
|
{ |
|
|
|
using (var uow = UnitOfWorkManager.Begin(isTransactional: true, requiresNew: true)) |
|
|
|
{ |
|
|
|
if (waitingEvent.RetryCount > AbpInboxProcessorOptions.MaxRetryCount) |
|
|
|
if (waitingEvent.RetryCount > EventBusBoxesOptions.InboxProcessorMaxRetryCount) |
|
|
|
{ |
|
|
|
Logger.LogWarning($"Max retry count reached for event with id = {waitingEvent.Id:N}. Discarding the event."); |
|
|
|
Logger.LogWarning($"Event with id = {waitingEvent.Id:N} has exceeded the maximum retry count. Marking it as discarded."); |
|
|
|
|
|
|
|
await Inbox.MarkAsDiscardAsync(waitingEvent.Id); |
|
|
|
await uow.CompleteAsync(StoppingToken); |
|
|
|
continue; |
|
|
|
} |
|
|
|
|
|
|
|
Logger.LogInformation($"Retrying event with id = {waitingEvent.Id:N}. " + |
|
|
|
$"Retry count: {waitingEvent.RetryCount}, " + |
|
|
|
$"Next retry time: {GetNextRetryTime(waitingEvent.RetryCount, AbpInboxProcessorOptions.MaxRetryCount)}"); |
|
|
|
Logger.LogInformation($"Event with id = {waitingEvent.Id:N} will retry later. " + |
|
|
|
$"Current retry count: {waitingEvent.RetryCount}, " + |
|
|
|
$"Next retry time: {GetNextRetryTime(waitingEvent.RetryCount)}" + |
|
|
|
$"Max retry count: {EventBusBoxesOptions.InboxProcessorMaxRetryCount}."); |
|
|
|
|
|
|
|
await Inbox.RetryLaterAsync(waitingEvent.Id, ++waitingEvent.RetryCount, GetNextRetryTime(waitingEvent.RetryCount, AbpInboxProcessorOptions.MaxRetryCount)); |
|
|
|
await Inbox.RetryLaterAsync(waitingEvent.Id, ++waitingEvent.RetryCount, GetNextRetryTime(waitingEvent.RetryCount)); |
|
|
|
await uow.CompleteAsync(StoppingToken); |
|
|
|
} |
|
|
|
continue; |
|
|
|
} |
|
|
|
|
|
|
|
if (AbpInboxProcessorOptions.FailurePolicy == InboxProcessorFailurePolicy.Discard) |
|
|
|
if (EventBusBoxesOptions.InboxProcessorFailurePolicy == InboxProcessorFailurePolicy.Discard) |
|
|
|
{ |
|
|
|
using (var uow = UnitOfWorkManager.Begin(isTransactional: true, requiresNew: true)) |
|
|
|
{ |
|
|
|
Logger.LogInformation($"Discarding event with id = {waitingEvent.Id:N} due to an error."); |
|
|
|
Logger.LogInformation($"Event with id = {waitingEvent.Id:N} will be discarded."); |
|
|
|
|
|
|
|
await Inbox.MarkAsDiscardAsync(waitingEvent.Id); |
|
|
|
await uow.CompleteAsync(StoppingToken); |
|
|
|
@ -187,15 +180,9 @@ public class InboxProcessor : IInboxProcessor, ITransientDependency |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
protected virtual DateTime? GetNextRetryTime(int retryCount, int maxRetryCount) |
|
|
|
protected virtual DateTime? GetNextRetryTime(int retryCount) |
|
|
|
{ |
|
|
|
if (retryCount > maxRetryCount) |
|
|
|
{ |
|
|
|
return null; |
|
|
|
} |
|
|
|
|
|
|
|
var delaySeconds = 1 * (int)Math.Pow(2, retryCount - 1); |
|
|
|
|
|
|
|
return DateTime.Now.AddSeconds(delaySeconds); |
|
|
|
} |
|
|
|
|
|
|
|
|