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 f5f1d7010e..2b0a448f74 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 @@ -26,7 +26,7 @@ namespace Volo.Abp.EntityFrameworkCore.DistributedEvents [UnitOfWork] public virtual async Task EnqueueAsync(IncomingEventInfo incomingEvent) { - var dbContext = await GetDbContextAsync(); + var dbContext = await DbContextProvider.GetDbContextAsync(); dbContext.IncomingEvents.Add( new IncomingEventRecord(incomingEvent) @@ -36,7 +36,7 @@ namespace Volo.Abp.EntityFrameworkCore.DistributedEvents [UnitOfWork] public virtual async Task> GetWaitingEventsAsync(int maxCount) { - var dbContext = await GetDbContextAsync(); + var dbContext = await DbContextProvider.GetDbContextAsync(); var outgoingEventRecords = await dbContext .IncomingEvents @@ -55,7 +55,7 @@ namespace Volo.Abp.EntityFrameworkCore.DistributedEvents public async Task MarkAsProcessedAsync(Guid id) { //TODO: Optimize? - var dbContext = await GetDbContextAsync(); + var dbContext = await DbContextProvider.GetDbContextAsync(); var incomingEvent = await dbContext.IncomingEvents.FindAsync(id); if (incomingEvent != null) { @@ -66,18 +66,21 @@ namespace Volo.Abp.EntityFrameworkCore.DistributedEvents [UnitOfWork] public async Task ExistsByMessageIdAsync(string messageId) { - var dbContext = await GetDbContextAsync(); + //TODO: Optimize + var dbContext = await DbContextProvider.GetDbContextAsync(); return await dbContext.IncomingEvents.AnyAsync(x => x.MessageId == messageId); } - private async Task GetDbContextAsync() - { - return (IHasEventInbox)await DbContextProvider.GetDbContextAsync(); - } - - public Task DeleteOldEventsAsync() + [UnitOfWork] + public async Task DeleteOldEventsAsync() { - throw new NotImplementedException(); + //TODO: Optimize + var dbContext = await DbContextProvider.GetDbContextAsync(); + var timeToKeepEvents = Clock.Now.AddHours(-2); //TODO: Config? + var oldEvents = await dbContext.IncomingEvents + .Where(x => x.Processed && x.CreationTime < timeToKeepEvents) + .ToListAsync(); + dbContext.IncomingEvents.RemoveRange(oldEvents); } } } \ No newline at end of file diff --git a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/InboxProcessor.cs b/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/InboxProcessor.cs index b1e6a1492b..bba97fe3d6 100644 --- a/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/InboxProcessor.cs +++ b/framework/src/Volo.Abp.EventBus.Boxes/Volo/Abp/EventBus/Boxes/InboxProcessor.cs @@ -8,6 +8,7 @@ using Microsoft.Extensions.Logging.Abstractions; using Volo.Abp.DependencyInjection; using Volo.Abp.EventBus.Distributed; using Volo.Abp.Threading; +using Volo.Abp.Timing; using Volo.Abp.Uow; namespace Volo.Abp.EventBus.Boxes @@ -19,9 +20,12 @@ namespace Volo.Abp.EventBus.Boxes protected IDistributedEventBus DistributedEventBus { get; } protected IDistributedLockProvider DistributedLockProvider { get; } protected IUnitOfWorkManager UnitOfWorkManager { get; } + protected IClock Clock { get; } protected IEventInbox Inbox { get; private set; } protected InboxConfig InboxConfig { get; private set; } + protected DateTime? LastCleanTime { get; set; } + protected string DistributedLockName => "Inbox_" + InboxConfig.Name; public ILogger Logger { get; set; } @@ -30,13 +34,15 @@ namespace Volo.Abp.EventBus.Boxes AbpTimer timer, IDistributedEventBus distributedEventBus, IDistributedLockProvider distributedLockProvider, - IUnitOfWorkManager unitOfWorkManager) + IUnitOfWorkManager unitOfWorkManager, + IClock clock) { ServiceProvider = serviceProvider; Timer = timer; DistributedEventBus = distributedEventBus; DistributedLockProvider = distributedLockProvider; UnitOfWorkManager = unitOfWorkManager; + Clock = clock; Timer.Period = 2000; //TODO: Config? Timer.Elapsed += TimerOnElapsed; Logger = NullLogger.Instance; @@ -69,6 +75,8 @@ namespace Volo.Abp.EventBus.Boxes { Logger.LogDebug("Obtained the distributed lock: " + DistributedLockName); + await DeleteOldEventsAsync(); + while (true) { var waitingEvents = await Inbox.GetWaitingEventsAsync(1000); //TODO: Config? @@ -87,12 +95,8 @@ namespace Volo.Abp.EventBus.Boxes .AsRawEventPublisher() .ProcessRawAsync(InboxConfig, waitingEvent.EventName, waitingEvent.EventData); - /* - await DistributedEventBus - .AsRawEventPublisher() - .PublishRawAsync(waitingEvent.Id, waitingEvent.EventName, waitingEvent.EventData); - */ await Inbox.MarkAsProcessedAsync(waitingEvent.Id); + await uow.CompleteAsync(); } @@ -107,5 +111,17 @@ namespace Volo.Abp.EventBus.Boxes } } } + + protected virtual async Task DeleteOldEventsAsync() + { + if (LastCleanTime != null && LastCleanTime > Clock.Now.AddHours(6)) //TODO: Config? + { + return; + } + + await Inbox.DeleteOldEventsAsync(); + + LastCleanTime = DateTime.Now; + } } } \ No newline at end of file