diff --git a/framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/DistributedEvents/DbContextEventInbox.cs b/framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/DistributedEvents/DbContextEventInbox.cs index 856b8612be..2e148a97c6 100644 --- a/framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/DistributedEvents/DbContextEventInbox.cs +++ b/framework/src/Volo.Abp.EntityFrameworkCore/Volo/Abp/EntityFrameworkCore/DistributedEvents/DbContextEventInbox.cs @@ -51,6 +51,7 @@ public class DbContextEventInbox : IDbContextEventInbox .IncomingEvents .AsNoTracking() .Where(x => x.Status == IncomingEventStatus.Pending) + .Where(x => x.NextRetryTime == null || x.NextRetryTime <= Clock.Now) .WhereIf(transformedFilter != null, transformedFilter!) .OrderBy(x => x.CreationTime) .Take(maxCount) diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/AbpEventBusBoxesOptions.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/AbpEventBusBoxesOptions.cs index cc91cad5df..ec7b603176 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/AbpEventBusBoxesOptions.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/AbpEventBusBoxesOptions.cs @@ -36,6 +36,16 @@ public class AbpEventBusBoxesOptions /// public TimeSpan PeriodTimeSpan { get; set; } + /// + /// Default: + /// + public InboxProcessorFailurePolicy InboxProcessorFailurePolicy { get; set; } = InboxProcessorFailurePolicy.Retry; + + /// + /// Default: 10 + /// + public int InboxProcessorMaxRetryCount { get; set; } = 10; + /// /// Default: 15 seconds /// diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/AbpInboxProcessorOptions.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/AbpInboxProcessorOptions.cs deleted file mode 100644 index 32b209a94d..0000000000 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/AbpInboxProcessorOptions.cs +++ /dev/null @@ -1,15 +0,0 @@ -namespace Volo.Abp.EventBus.Distributed; - -public class AbpInboxProcessorOptions -{ - public InboxProcessorFailurePolicy FailurePolicy { get; set; } = InboxProcessorFailurePolicy.Retry; - - /// - /// Retry intervals follow an exponential backoff strategy: - /// The delay for each retry is twice the previous one, - /// starting with an initial delay of 1 second. - /// For example: 1s, 2s, 4s, 8s, 16s, 32s ... - /// Maximum of 10 retries. - /// - public int MaxRetryCount { get; set; } = 10; -} diff --git a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/InboxProcessor.cs b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/InboxProcessor.cs index df137725b8..3b61527bc5 100644 --- a/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/InboxProcessor.cs +++ b/framework/src/Volo.Abp.EventBus/Volo/Abp/EventBus/Distributed/InboxProcessor.cs @@ -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 eventBusBoxesOptions, - IOptions inboxProcessorOptions) + IOptions 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.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); } diff --git a/framework/src/Volo.Abp.MongoDB/Volo/Abp/MongoDB/DistributedEvents/MongoDbContextEventInbox.cs b/framework/src/Volo.Abp.MongoDB/Volo/Abp/MongoDB/DistributedEvents/MongoDbContextEventInbox.cs index 112f6d7cf7..427c794542 100644 --- a/framework/src/Volo.Abp.MongoDB/Volo/Abp/MongoDB/DistributedEvents/MongoDbContextEventInbox.cs +++ b/framework/src/Volo.Abp.MongoDB/Volo/Abp/MongoDB/DistributedEvents/MongoDbContextEventInbox.cs @@ -65,6 +65,7 @@ public class MongoDbContextEventInbox : IMongoDbContextEventInb .IncomingEvents .AsQueryable() .Where(x => x.Status == IncomingEventStatus.Pending) + .Where(x => x.NextRetryTime == null || x.NextRetryTime <= Clock.Now) .WhereIf(transformedFilter != null, transformedFilter!) .OrderBy(x => x.CreationTime) .Take(maxCount)