From 4acde73dade3d9542cf377cd2102ee9f40a07648 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Thu, 31 Oct 2024 09:09:05 +0100 Subject: [PATCH] Add bookmark queue purger functionality (#6080) * Add bookmark queue purger functionality Introduce a new class `DefaultBookmarkQueuePurger` to purge old bookmark queue items. This includes an interface `IBookmarkQueuePurger` and a recurring task `PurgeBookmarkQueueRecurringTask`, with necessary updates to `IBookmarkQueueStore` implementations and application configuration. * Update purging logic in DefaultBookmarkQueuePurger Refactor the purging operation to use a threshold date for clarity. This includes updating log messages and filter creation to enhance readability and maintainability. --- src/apps/Elsa.Server.Web/Program.cs | 5 +- .../Stores/DapperBookmarkQueueStore.cs | 12 +++++ .../Modules/Runtime/BookmarkQueueStore.cs | 13 +++++ .../Modules/Runtime/BookmarkQueueStore.cs | 12 +++++ .../Features/SecretManagementFeature.cs | 2 +- .../Contracts/IBookmarkQueuePurger.cs | 6 +++ .../Contracts/IBookmarkQueueStore.cs | 8 +++ .../Features/WorkflowRuntimeFeature.cs | 4 +- .../Filters/BookmarkQueueFilter.cs | 11 ++++ .../Services/DefaultBookmarkQueuePurger.cs | 52 +++++++++++++++++++ .../Stores/MemoryBookmarkQueueStore.cs | 6 +++ .../Tasks/PurgeBookmarkQueueRecurringTask.cs | 14 +++++ 12 files changed, 141 insertions(+), 4 deletions(-) create mode 100644 src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkQueuePurger.cs create mode 100644 src/modules/Elsa.Workflows.Runtime/Services/DefaultBookmarkQueuePurger.cs create mode 100644 src/modules/Elsa.Workflows.Runtime/Tasks/PurgeBookmarkQueueRecurringTask.cs diff --git a/src/apps/Elsa.Server.Web/Program.cs b/src/apps/Elsa.Server.Web/Program.cs index 8038fa15b..f84ef78ec 100644 --- a/src/apps/Elsa.Server.Web/Program.cs +++ b/src/apps/Elsa.Server.Web/Program.cs @@ -492,11 +492,12 @@ services // Obfuscate HTTP request headers. services.AddActivityStateFilter(); -// Configure recurring tasks. +// Optionally configure recurring tasks using alternative schedules. services.Configure(options => { options.Schedule.ConfigureTask(TimeSpan.FromSeconds(30)); - options.Schedule.ConfigureTask(TimeSpan.FromSeconds(10)); + options.Schedule.ConfigureTask(TimeSpan.FromHours(4)); + options.Schedule.ConfigureTask(TimeSpan.FromSeconds(60)); }); //services.Configure(options => options.CacheDuration = TimeSpan.FromDays(1)); diff --git a/src/modules/Elsa.Dapper/Modules/Runtime/Stores/DapperBookmarkQueueStore.cs b/src/modules/Elsa.Dapper/Modules/Runtime/Stores/DapperBookmarkQueueStore.cs index 56df885e0..9f442c404 100644 --- a/src/modules/Elsa.Dapper/Modules/Runtime/Stores/DapperBookmarkQueueStore.cs +++ b/src/modules/Elsa.Dapper/Modules/Runtime/Stores/DapperBookmarkQueueStore.cs @@ -39,12 +39,24 @@ internal class DapperBookmarkQueueStore(Store store, IP return record != null ? Map(record) : default; } + public async Task> FindManyAsync(BookmarkQueueFilter filter, CancellationToken cancellationToken = default) + { + var records = await store.FindManyAsync(q => ApplyFilter(q, filter), cancellationToken); + return Map(records); + } + public async Task> PageAsync(PageArgs pageArgs, BookmarkQueueItemOrder orderBy, CancellationToken cancellationToken = default) { var records = await store.ListAsync(pageArgs, orderBy.KeySelector.GetPropertyName(), orderBy.Direction, cancellationToken); return Map(records); } + public async Task> PageAsync(PageArgs pageArgs, BookmarkQueueFilter filter, BookmarkQueueItemOrder orderBy, CancellationToken cancellationToken = default) + { + var records = await store.FindManyAsync(q => ApplyFilter(q, filter), pageArgs, orderBy.KeySelector.GetPropertyName(), orderBy.Direction, cancellationToken); + return Map(records); + } + /// public async Task DeleteAsync(BookmarkQueueFilter filter, CancellationToken cancellationToken = default) { diff --git a/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/BookmarkQueueStore.cs b/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/BookmarkQueueStore.cs index 562a341e0..1532b2a8c 100644 --- a/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/BookmarkQueueStore.cs +++ b/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/BookmarkQueueStore.cs @@ -33,6 +33,12 @@ public class EFBookmarkQueueStore(Store return store.FindAsync(filter.Apply, OnLoadAsync, cancellationToken); } + /// + public async Task> FindManyAsync(BookmarkQueueFilter filter, CancellationToken cancellationToken = default) + { + return await store.QueryAsync(filter.Apply, OnLoadAsync, filter.TenantAgnostic, cancellationToken); + } + public async Task> PageAsync(PageArgs pageArgs, BookmarkQueueItemOrder orderBy, CancellationToken cancellationToken = default) { var count = await store.QueryAsync(queryable => queryable.OrderBy(orderBy), cancellationToken).LongCount(); @@ -40,6 +46,13 @@ public class EFBookmarkQueueStore(Store return new(results, count); } + public async Task> PageAsync(PageArgs pageArgs, BookmarkQueueFilter filter, BookmarkQueueItemOrder orderBy, CancellationToken cancellationToken = default) + { + var count = await store.QueryAsync(queryable => filter.Apply(queryable).OrderBy(orderBy), cancellationToken).LongCount(); + var results = await store.QueryAsync(queryable => filter.Apply(queryable).OrderBy(orderBy).Paginate(pageArgs), OnLoadAsync, cancellationToken).ToList(); + return new(results, count); + } + /// public async Task DeleteAsync(BookmarkQueueFilter filter, CancellationToken cancellationToken = default) { diff --git a/src/modules/Elsa.MongoDb/Modules/Runtime/BookmarkQueueStore.cs b/src/modules/Elsa.MongoDb/Modules/Runtime/BookmarkQueueStore.cs index bfb3444c9..79841ddf0 100644 --- a/src/modules/Elsa.MongoDb/Modules/Runtime/BookmarkQueueStore.cs +++ b/src/modules/Elsa.MongoDb/Modules/Runtime/BookmarkQueueStore.cs @@ -32,6 +32,11 @@ public class MongoBookmarkQueueStore(MongoDbStore mongoDbStor return await mongoDbStore.FindAsync(query => Filter(query, filter), cancellationToken); } + public async Task> FindManyAsync(BookmarkQueueFilter filter, CancellationToken cancellationToken = default) + { + return await mongoDbStore.FindManyAsync(query => Filter(query, filter), cancellationToken); + } + public async Task> PageAsync(PageArgs pageArgs, BookmarkQueueItemOrder orderBy, CancellationToken cancellationToken = default) { var results = await mongoDbStore.FindManyAsync(query => Paginate(Order(query, orderBy), pageArgs), cancellationToken); @@ -39,6 +44,13 @@ public class MongoBookmarkQueueStore(MongoDbStore mongoDbStor return Page.Of(results.ToList(), count); } + public async Task> PageAsync(PageArgs pageArgs, BookmarkQueueFilter filter, BookmarkQueueItemOrder orderBy, CancellationToken cancellationToken = default) + { + var results = await mongoDbStore.FindManyAsync(query => Paginate(Order(Filter(query, filter), orderBy), pageArgs), cancellationToken); + var count = await mongoDbStore.CountAsync(queryable => Filter(queryable, filter), cancellationToken); + return Page.Of(results.ToList(), count); + } + /// public async Task DeleteAsync(BookmarkQueueFilter filter, CancellationToken cancellationToken = default) { diff --git a/src/modules/Elsa.Secrets.Management/Features/SecretManagementFeature.cs b/src/modules/Elsa.Secrets.Management/Features/SecretManagementFeature.cs index d4924357d..e2f5e082e 100644 --- a/src/modules/Elsa.Secrets.Management/Features/SecretManagementFeature.cs +++ b/src/modules/Elsa.Secrets.Management/Features/SecretManagementFeature.cs @@ -57,7 +57,7 @@ public class SecretManagementFeature(IModule module) : FeatureBase(module) .AddScoped() .AddScoped() .AddScoped() - .AddRecurringTask() + .AddRecurringTask(TimeSpan.FromHours(4)) ; } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkQueuePurger.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkQueuePurger.cs new file mode 100644 index 000000000..155c9f205 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkQueuePurger.cs @@ -0,0 +1,6 @@ +namespace Elsa.Workflows.Runtime; + +public interface IBookmarkQueuePurger +{ + Task PurgeAsync(CancellationToken cancellationToken = default); +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkQueueStore.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkQueueStore.cs index d72764c50..f334b8247 100644 --- a/src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkQueueStore.cs +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IBookmarkQueueStore.cs @@ -21,11 +21,19 @@ public interface IBookmarkQueueStore /// Returns the first bookmark queue item matching the specified filter. Task FindAsync(BookmarkQueueFilter filter, CancellationToken cancellationToken = default); + + /// Returns a set of bookmark queue items matching the specified filter. + Task> FindManyAsync(BookmarkQueueFilter filter, CancellationToken cancellationToken = default); /// /// Returns a page of records, ordered by the specified order definition. /// Task> PageAsync(PageArgs pageArgs, BookmarkQueueItemOrder orderBy, CancellationToken cancellationToken = default); + + /// + /// Returns a page of records, filtered and ordered by the specified order definition. + /// + Task> PageAsync(PageArgs pageArgs, BookmarkQueueFilter filter, BookmarkQueueItemOrder orderBy, CancellationToken cancellationToken = default); /// /// Deletes a set of bookmark queue items matching the specified filter. diff --git a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs index fa1fe8168..2fa57173b 100644 --- a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs +++ b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs @@ -209,6 +209,7 @@ public class WorkflowRuntimeFeature : FeatureBase .AddScoped() .AddScoped() .AddScoped() + .AddScoped() .AddScoped, ActivityExecutionRecordExtractor>() .AddScoped, WorkflowExecutionLogRecordExtractor>() @@ -242,7 +243,8 @@ public class WorkflowRuntimeFeature : FeatureBase // Startup tasks, background tasks, and recurring tasks. .AddStartupTask() - .AddRecurringTask() + .AddRecurringTask(TimeSpan.FromMinutes(1)) + .AddRecurringTask(TimeSpan.FromMinutes(1)) // Distributed locking. .AddSingleton(DistributedLockProvider) diff --git a/src/modules/Elsa.Workflows.Runtime/Filters/BookmarkQueueFilter.cs b/src/modules/Elsa.Workflows.Runtime/Filters/BookmarkQueueFilter.cs index 0b55c1c60..8abc88d18 100644 --- a/src/modules/Elsa.Workflows.Runtime/Filters/BookmarkQueueFilter.cs +++ b/src/modules/Elsa.Workflows.Runtime/Filters/BookmarkQueueFilter.cs @@ -7,6 +7,9 @@ public class BookmarkQueueFilter { /// Gets or sets the ID of the bookmark queue item. public string? Id { get; set; } + + /// Gets or sets the IDs of the bookmark queue items. + public IEnumerable? Ids { get; set; } /// Gets or sets the ID of the bookmark. public string? BookmarkId { get; set; } @@ -22,17 +25,25 @@ public class BookmarkQueueFilter // The type name of the activity associated with the bookmark. public string? ActivityTypeName { get; set; } + + /// The timestamp less than which the bookmark queue item was created. + public DateTimeOffset? CreatedAtLessThan { get; set; } + + /// Gets or sets a value indicating whether the filter is tenant agnostic. + public bool TenantAgnostic { get; set; } /// Applies the filter to the specified query. public IQueryable Apply(IQueryable query) { var filter = this; if (filter.Id != null) query = query.Where(x => x.Id == filter.Id); + if (filter.Ids != null) query = query.Where(x => filter.Ids.Contains(x.Id)); if (filter.BookmarkId != null) query = query.Where(x => x.BookmarkId == filter.BookmarkId); if (filter.BookmarkHash != null) query = query.Where(x => x.StimulusHash == filter.BookmarkHash); if (filter.ActivityInstanceId != null) query = query.Where(x => x.ActivityInstanceId == filter.ActivityInstanceId); if (filter.ActivityTypeName != null) query = query.Where(x => x.ActivityTypeName == filter.ActivityTypeName); if (filter.WorkflowInstanceId != null) query = query.Where(x => x.WorkflowInstanceId == filter.WorkflowInstanceId); + if (filter.CreatedAtLessThan != null) query = query.Where(x => x.CreatedAt < filter.CreatedAtLessThan); return query; } diff --git a/src/modules/Elsa.Workflows.Runtime/Services/DefaultBookmarkQueuePurger.cs b/src/modules/Elsa.Workflows.Runtime/Services/DefaultBookmarkQueuePurger.cs new file mode 100644 index 000000000..a1ce6110c --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Services/DefaultBookmarkQueuePurger.cs @@ -0,0 +1,52 @@ +using Elsa.Common; +using Elsa.Common.Entities; +using Elsa.Common.Models; +using Elsa.Workflows.Runtime.Filters; +using Elsa.Workflows.Runtime.OrderDefinitions; +using JetBrains.Annotations; +using Microsoft.Extensions.Logging; + +namespace Elsa.Workflows.Runtime; + +[UsedImplicitly] +public class DefaultBookmarkQueuePurger(IBookmarkQueueStore store, ISystemClock systemClock, ILogger logger) : IBookmarkQueuePurger +{ + private readonly TimeSpan _ttl = TimeSpan.FromMinutes(1); + private readonly int _batchSize = 50; + + public async Task PurgeAsync(CancellationToken cancellationToken = default) + { + var currentPage = 0; + var now = systemClock.UtcNow; + var thresholdDate = now - _ttl; + + logger.LogInformation("Purging bookmark queue items older than {ThresholdDate}.", thresholdDate); + + while (true) + { + var pageArgs = PageArgs.FromPage(currentPage, _batchSize); + var filter = new BookmarkQueueFilter + { + CreatedAtLessThan = thresholdDate + }; + var order = new BookmarkQueueItemOrder(x => x.CreatedAt, OrderDirection.Ascending); + var page = await store.PageAsync(pageArgs, filter, order, cancellationToken); + var items = page.Items; + + if (items.Count == 0) + break; + + var ids = items.Select(x => x.Id).ToList(); + await store.DeleteAsync(new BookmarkQueueFilter + { + Ids = ids + }, cancellationToken); + + logger.LogInformation("Purged {Count} bookmark queue items.", items.Count); + + currentPage++; + } + + logger.LogInformation("Finished purging bookmark queue items."); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Stores/MemoryBookmarkQueueStore.cs b/src/modules/Elsa.Workflows.Runtime/Stores/MemoryBookmarkQueueStore.cs index 80f34ca2f..a38dd31ec 100644 --- a/src/modules/Elsa.Workflows.Runtime/Stores/MemoryBookmarkQueueStore.cs +++ b/src/modules/Elsa.Workflows.Runtime/Stores/MemoryBookmarkQueueStore.cs @@ -38,6 +38,12 @@ public class MemoryBookmarkQueueStore(MemoryStore store) : IB return Task.FromResult(entities); } + public Task> PageAsync(PageArgs pageArgs, BookmarkQueueFilter filter, BookmarkQueueItemOrder orderBy, CancellationToken cancellationToken = default) + { + var entities = store.Query(query => Filter(query, filter).OrderBy(orderBy)).Paginate(pageArgs); + return Task.FromResult(entities); + } + public Task> FindManyAsync(BookmarkQueueFilter filter, CancellationToken cancellationToken = default) { var entities = store.Query(query => Filter(query, filter)).AsEnumerable(); diff --git a/src/modules/Elsa.Workflows.Runtime/Tasks/PurgeBookmarkQueueRecurringTask.cs b/src/modules/Elsa.Workflows.Runtime/Tasks/PurgeBookmarkQueueRecurringTask.cs new file mode 100644 index 000000000..ca5138568 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Tasks/PurgeBookmarkQueueRecurringTask.cs @@ -0,0 +1,14 @@ +using Elsa.Common; +using JetBrains.Annotations; + +namespace Elsa.Workflows.Runtime.Tasks; + +/// Periodically purges the bookmark queue of old items. +[UsedImplicitly] +public class PurgeBookmarkQueueRecurringTask(IBookmarkQueuePurger bookmarkQueueWorker) : RecurringTask +{ + public override Task ExecuteAsync(CancellationToken stoppingToken) + { + return bookmarkQueueWorker.PurgeAsync(stoppingToken); + } +} \ No newline at end of file