Refactor retention module to use recurring tasks

Replaced the CleanupHostedService with CleanupRecurringTask for better scheduling and simpler configuration. Converted several classes to use constructor-based dependency injection and streamlined some type assignments. Removed unneeded files and updated existing code to follow more concise conventions. This enhances maintainability and readability.
This commit is contained in:
Sipke Schoorstra 2024-10-29 08:36:31 +01:00
parent d899b674dc
commit d3d171f59f
18 changed files with 91 additions and 245 deletions

View file

@ -1,4 +1,5 @@
using Elsa.Common;
using Elsa.Common.RecurringTasks;
using Elsa.Common.Services;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.DependencyInjection.Extensions;
@ -44,4 +45,10 @@ public static class DependencyInjectionExtensions
{
return services.AddScoped<IRecurringTask, T>();
}
public static IServiceCollection AddRecurringTask<T>(this IServiceCollection services, TimeSpan interval) where T : class, IRecurringTask
{
services.Configure<RecurringTaskOptions>(options => options.Schedule.ConfigureTask<T>(interval));
return services.AddRecurringTask<T>();
}
}

View file

@ -9,29 +9,20 @@ namespace Elsa.Retention.CleanupStrategies;
/// <summary>
/// Deletes activity execution records.
/// </summary>
public class DeleteActivityExecutionRecordStrategy : IDeletionCleanupStrategy<ActivityExecutionRecord>
public class DeleteActivityExecutionRecordStrategy(IActivityExecutionStore store, ILogger<DeleteActivityExecutionRecordStrategy> logger) : IDeletionCleanupStrategy<ActivityExecutionRecord>
{
private readonly ILogger<DeleteActivityExecutionRecordStrategy> _logger;
private readonly IActivityExecutionStore _store;
public DeleteActivityExecutionRecordStrategy(IActivityExecutionStore store, ILogger<DeleteActivityExecutionRecordStrategy> logger)
{
_store = store;
_logger = logger;
}
public async Task Cleanup(ICollection<ActivityExecutionRecord> collection)
{
ActivityExecutionRecordFilter filter = new()
var filter = new ActivityExecutionRecordFilter()
{
Ids = collection.Select(x => x.Id).ToList()
};
long deletedRecords = await _store.DeleteManyAsync(filter);
var deletedRecords = await store.DeleteManyAsync(filter);
if (deletedRecords != collection.Count)
{
_logger.LogWarning("Expected to delete {Expected} activity execution records, actually deleted {Actual} activity execution records", collection.Count, deletedRecords);
logger.LogWarning("Expected to delete {Expected} activity execution records, actually deleted {Actual} activity execution records", collection.Count, deletedRecords);
}
}
}

View file

@ -9,29 +9,20 @@ namespace Elsa.Retention.CleanupStrategies;
/// <summary>
/// Deletes a collection of bookmarks
/// </summary>
public class DeleteBookmarkStrategy : IDeletionCleanupStrategy<StoredBookmark>
public class DeleteBookmarkStrategy(IBookmarkStore store, ILogger<DeleteBookmarkStrategy> logger) : IDeletionCleanupStrategy<StoredBookmark>
{
private readonly ILogger<DeleteBookmarkStrategy> _logger;
private readonly IBookmarkStore _store;
public DeleteBookmarkStrategy(IBookmarkStore store, ILogger<DeleteBookmarkStrategy> logger)
{
_store = store;
_logger = logger;
}
public async Task Cleanup(ICollection<StoredBookmark> collection)
{
BookmarkFilter bookmarkFilter = new()
var bookmarkFilter = new BookmarkFilter
{
BookmarkIds = collection.Select(x => x.Id).ToList()
};
long deletedRecords = await _store.DeleteAsync(bookmarkFilter);
var deletedRecords = await store.DeleteAsync(bookmarkFilter);
if (deletedRecords != collection.Count)
{
_logger.LogWarning("Expected to delete {Expected} bookmarks, actually deleted {Actual} bookmarks", collection.Count, deletedRecords);
logger.LogWarning("Expected to delete {Expected} bookmarks, actually deleted {Actual} bookmarks", collection.Count, deletedRecords);
}
}
}

View file

@ -9,29 +9,20 @@ namespace Elsa.Retention.CleanupStrategies;
/// <summary>
/// Deletes <see cref="WorkflowExecutionLogRecord" />
/// </summary>
public class DeleteWorkflowExecutionRecordStrategy : IDeletionCleanupStrategy<WorkflowExecutionLogRecord>
public class DeleteWorkflowExecutionRecordStrategy(IWorkflowExecutionLogStore store, ILogger<DeleteWorkflowExecutionRecordStrategy> logger) : IDeletionCleanupStrategy<WorkflowExecutionLogRecord>
{
private readonly ILogger<DeleteWorkflowExecutionRecordStrategy> _logger;
private readonly IWorkflowExecutionLogStore _store;
public DeleteWorkflowExecutionRecordStrategy(IWorkflowExecutionLogStore store, ILogger<DeleteWorkflowExecutionRecordStrategy> logger)
{
_store = store;
_logger = logger;
}
public async Task Cleanup(ICollection<WorkflowExecutionLogRecord> collection)
{
WorkflowExecutionLogRecordFilter filter = new()
var filter = new WorkflowExecutionLogRecordFilter()
{
Ids = collection.Select(x => x.Id).ToList()
};
long deletedRecords = await _store.DeleteManyAsync(filter);
var deletedRecords = await store.DeleteManyAsync(filter);
if (deletedRecords != collection.Count)
{
_logger.LogWarning("Expected to delete {Expected} workflow execution records, actually deleted {Actual} workflow execution records", collection.Count, deletedRecords);
logger.LogWarning("Expected to delete {Expected} workflow execution records, actually deleted {Actual} workflow execution records", collection.Count, deletedRecords);
}
}
}

View file

@ -9,27 +9,20 @@ namespace Elsa.Retention.Collectors;
/// <summary>
/// Collects all <see cref="ActivityExecutionRecord" /> related to the <see cref="WorkflowInstance" />
/// </summary>
public class ActivityExecutionRecordCollector : IRelatedEntityCollector<ActivityExecutionRecord>
public class ActivityExecutionRecordCollector(IActivityExecutionStore store) : IRelatedEntityCollector<ActivityExecutionRecord>
{
private readonly IActivityExecutionStore _store;
public ActivityExecutionRecordCollector(IActivityExecutionStore store)
{
_store = store;
}
public async IAsyncEnumerable<ICollection<ActivityExecutionRecord>> GetRelatedEntities(ICollection<WorkflowInstance> workflowInstances)
{
IEnumerable<WorkflowInstance[]> chunks = workflowInstances.Chunk(5);
var chunks = workflowInstances.Chunk(5);
foreach (WorkflowInstance[] chunk in chunks)
foreach (var chunk in chunks)
{
ActivityExecutionRecordFilter filter = new()
var filter = new ActivityExecutionRecordFilter()
{
WorkflowInstanceIds = chunk.Select(x => x.Id).ToArray()
};
IEnumerable<ActivityExecutionRecord> records = await _store.FindManyAsync(filter);
var records = await store.FindManyAsync(filter);
yield return records.ToArray();
}
}

View file

@ -9,27 +9,20 @@ namespace Elsa.Retention.Collectors;
/// <summary>
/// Collects all <see cref="StoredBookmark" /> related to the <see cref="WorkflowInstance" />
/// </summary>
public class BookmarkCollector : IRelatedEntityCollector<StoredBookmark>
public class BookmarkCollector(IBookmarkStore store) : IRelatedEntityCollector<StoredBookmark>
{
private readonly IBookmarkStore _store;
public BookmarkCollector(IBookmarkStore store)
{
_store = store;
}
public async IAsyncEnumerable<ICollection<StoredBookmark>> GetRelatedEntities(ICollection<WorkflowInstance> workflowInstances)
{
IEnumerable<WorkflowInstance[]> batches = workflowInstances.Chunk(25);
var batches = workflowInstances.Chunk(25);
foreach (WorkflowInstance[] batch in batches)
foreach (var batch in batches)
{
BookmarkFilter filter = new()
var filter = new BookmarkFilter()
{
WorkflowInstanceIds = batch.Select(x => x.Id).ToArray()
};
IEnumerable<StoredBookmark> bookmarks = await _store.FindManyAsync(filter);
var bookmarks = await store.FindManyAsync(filter);
yield return bookmarks.ToArray();
}
}

View file

@ -10,31 +10,24 @@ namespace Elsa.Retention.Collectors;
/// <summary>
/// Collects all <see cref="WorkflowExecutionLogRecord" /> related to the <see cref="WorkflowInstance" />
/// </summary>
public class WorkflowExecutionLogRecordCollector : IRelatedEntityCollector<WorkflowExecutionLogRecord>
public class WorkflowExecutionLogRecordCollector(IWorkflowExecutionLogStore store) : IRelatedEntityCollector<WorkflowExecutionLogRecord>
{
private readonly IWorkflowExecutionLogStore _store;
public WorkflowExecutionLogRecordCollector(IWorkflowExecutionLogStore store)
{
_store = store;
}
public async IAsyncEnumerable<ICollection<WorkflowExecutionLogRecord>> GetRelatedEntities(ICollection<WorkflowInstance> workflowInstances)
{
IEnumerable<WorkflowInstance[]> chunks = workflowInstances.Chunk(25);
var chunks = workflowInstances.Chunk(25);
foreach (WorkflowInstance[] chunk in chunks)
foreach (var chunk in chunks)
{
WorkflowExecutionLogRecordFilter filter = new()
var filter = new WorkflowExecutionLogRecordFilter()
{
WorkflowInstanceIds = chunk.Select(x => x.Id).ToArray()
};
PageArgs pageArgs = PageArgs.FromPage(0, 100);
var pageArgs = PageArgs.FromPage(0, 100);
while (true)
{
Page<WorkflowExecutionLogRecord> page = await _store.FindManyAsync(filter, pageArgs);
var page = await store.FindManyAsync(filter, pageArgs);
yield return page.Items.ToArray();
if (page.TotalCount <= pageArgs.Offset + page.Items.Count)

View file

@ -18,7 +18,7 @@ public interface IRelatedEntityCollector<TEntity> : IRelatedEntityCollector wher
{
async IAsyncEnumerable<ICollection<object>> IRelatedEntityCollector.GetRelatedEntitiesGeneric(ICollection<WorkflowInstance> workflowInstances)
{
await foreach (ICollection<TEntity> entity in GetRelatedEntities(workflowInstances).ConfigureAwait(false))
await foreach (var entity in GetRelatedEntities(workflowInstances).ConfigureAwait(false))
{
yield return entity.Select(x => (object)x).ToArray();
}

View file

@ -0,0 +1,2 @@
<wpf:ResourceDictionary xml:space="preserve" xmlns:x="http://schemas.microsoft.com/winfx/2006/xaml" xmlns:s="clr-namespace:System;assembly=mscorlib" xmlns:ss="urn:shemas-jetbrains-com:settings-storage-xaml" xmlns:wpf="http://schemas.microsoft.com/winfx/2006/xaml/presentation">
<s:Boolean x:Key="/Default/CodeInspection/NamespaceProvider/NamespaceFoldersToSkip/=tasks/@EntryIndexedValue">True</s:Boolean></wpf:ResourceDictionary>

View file

@ -1,15 +1,17 @@
using Elsa.Features.Services;
using Elsa.Retention.Feature;
using JetBrains.Annotations;
namespace Elsa.Retention.Extensions;
/// <summary>
/// Provides extensions to install the <see cref="RetentionFeature" /> feature.
/// Provides extensions to install the <see cref="RetentionFeature" /> feature.
/// </summary>
[UsedImplicitly]
public static class ModuleExtensions
{
/// <summary>
/// Install the <see cref="RetentionFeature" /> feature.
/// Installs the <see cref="RetentionFeature" /> feature.
/// </summary>
public static IModule UseRetention(this IModule module, Action<RetentionFeature>? configure = default)
{

View file

@ -19,7 +19,7 @@ public static class RetentionFeatureExtensions
/// <returns></returns>
public static RetentionFeature AddDeletePolicy(this RetentionFeature feature, string name, Func<IServiceProvider, RetentionWorkflowInstanceFilter> filterFactory)
{
List<IRetentionPolicy> policies = feature.Module.Properties.GetOrAdd(PoliciesKey, () => new List<IRetentionPolicy>());
var policies = feature.Module.Properties.GetOrAdd(PoliciesKey, () => new List<IRetentionPolicy>());
policies.Add(new DeletionRetentionPolicy(name, filterFactory));
return feature;
}

View file

@ -1,44 +0,0 @@
using Elsa.Workflows;
using Elsa.Workflows.Management.Filters;
using Elsa.Workflows.Management.Models;
namespace Elsa.Retention.Extensions;
public static class WorkflowInstanceFilterExtensions
{
/// <summary>
/// Clone the current filter
/// </summary>
/// <param name="filter"></param>
/// <returns></returns>
public static WorkflowInstanceFilter Clone(this WorkflowInstanceFilter filter)
{
return new WorkflowInstanceFilter
{
Id = filter.Id,
Ids = filter.Ids == null ? null : new List<string>(filter.Ids),
Version = filter.Version,
CorrelationId = filter.CorrelationId,
CorrelationIds = filter.CorrelationIds == null ? null : new List<string>(filter.CorrelationIds),
DefinitionId = filter.DefinitionId,
DefinitionIds = filter.DefinitionIds == null ? null : new List<string>(filter.DefinitionIds),
HasIncidents = filter.HasIncidents,
IsSystem = filter.IsSystem,
SearchTerm = filter.SearchTerm,
TimestampFilters = filter.TimestampFilters?.Select(x => new TimestampFilter
{
Column = x.Column,
Operator = x.Operator,
Timestamp = x.Timestamp
}).ToList(),
WorkflowStatus = filter.WorkflowStatus,
WorkflowStatuses = filter.WorkflowStatuses == null ? null : new List<WorkflowStatus>(filter.WorkflowStatuses),
DefinitionVersionId = filter.DefinitionVersionId,
DefinitionVersionIds = filter.DefinitionVersionIds == null ? null : new List<string>(filter.DefinitionVersionIds),
WorkflowSubStatus = filter.WorkflowSubStatus,
WorkflowSubStatuses = filter.WorkflowSubStatuses == null ? null : new List<WorkflowSubStatus>(filter.WorkflowSubStatuses),
ParentWorkflowInstanceIds =
filter.ParentWorkflowInstanceIds == null ? null : new List<string>(filter.ParentWorkflowInstanceIds)
};
}
}

View file

@ -1,10 +1,10 @@
using Elsa.Extensions;
using Elsa.Features.Abstractions;
using Elsa.Features.Services;
using Elsa.Retention.CleanupStrategies;
using Elsa.Retention.Collectors;
using Elsa.Retention.Contracts;
using Elsa.Retention.Extensions;
using Elsa.Retention.HostedServices;
using Elsa.Retention.Jobs;
using Elsa.Retention.Options;
using Elsa.Workflows.Runtime.Entities;
@ -44,15 +44,11 @@ public class RetentionFeature : FeatureBase
Services.AddScoped<IRelatedEntityCollector, ActivityExecutionRecordCollector>();
Services.AddScoped<IRelatedEntityCollector, WorkflowExecutionLogRecordCollector>();
foreach (IRetentionPolicy policy in this.GetPolicies())
Services.AddRecurringTask<CleanupRecurringTask>(TimeSpan.FromHours(4));
foreach (var policy in this.GetPolicies())
{
Services.AddSingleton(policy);
}
}
/// <inheritdoc cref="FeatureBase" />
public override void ConfigureHostedServices()
{
ConfigureHostedService<CleanupHostedService>();
}
}

View file

@ -1,55 +0,0 @@
using Elsa.Retention.Jobs;
using Elsa.Retention.Options;
using Medallion.Threading;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
using Microsoft.Extensions.Options;
namespace Elsa.Retention.HostedServices;
/// <summary>
/// Periodically wipes workflow instances and their execution logs.
/// </summary>
public class CleanupHostedService : BackgroundService
{
private readonly TimeSpan _interval;
private readonly ILogger<CleanupHostedService> _logger;
private readonly IServiceScopeFactory _serviceScopeFactory;
/// <summary>
/// Creates new Cleanup hosted service
/// </summary>
/// <param name="options"></param>
/// <param name="serviceScopeFactory"></param>
/// <param name="logger"></param>
public CleanupHostedService(IOptions<CleanupOptions> options, IServiceScopeFactory serviceScopeFactory, ILogger<CleanupHostedService> logger)
{
_serviceScopeFactory = serviceScopeFactory;
_logger = logger;
_interval = options.Value.SweepInterval;
}
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
while (!stoppingToken.IsCancellationRequested)
{
using IServiceScope scope = _serviceScopeFactory.CreateScope();
CleanupJob job = scope.ServiceProvider.GetRequiredService<CleanupJob>();
IDistributedLockProvider distributedLockProvider = scope.ServiceProvider.GetRequiredService<IDistributedLockProvider>();
await Task.Delay(_interval, stoppingToken);
await using IDistributedSynchronizationHandle handle = await distributedLockProvider.AcquireLockAsync(nameof(CleanupHostedService), cancellationToken: stoppingToken);
try
{
await job.ExecuteAsync(stoppingToken);
}
catch (Exception e)
{
_logger.LogError(e, "Failed to perform cleanup this time around. Next cleanup attempt will happen in {Interval}", _interval);
}
}
}
}

View file

@ -15,31 +15,15 @@ namespace Elsa.Retention.Jobs;
/// Deletes all workflow instances that match any of the defined <see cref="IRetentionPolicy" />
/// </summary>
[SuppressMessage("Trimming", "IL2055:Either the type on which the MakeGenericType is called can\'t be statically determined, or the type parameters to be used for generic arguments can\'t be statically determined.")]
public class CleanupJob
public class CleanupJob(
IWorkflowInstanceStore workflowInstanceStore,
IEnumerable<IRetentionPolicy> policies,
IOptions<CleanupOptions> options,
IServiceProvider serviceProvider,
ILogger<CleanupJob> logger)
{
private readonly ILogger _logger;
private readonly CleanupOptions _options;
private readonly IServiceProvider _serviceProvider;
private readonly IWorkflowInstanceStore _workflowInstanceStore;
/// <summary>
/// Creates a new cleanup job
/// </summary>
/// <param name="workflowInstanceStore"></param>
/// <param name="options"></param>
/// <param name="serviceProvider"></param>
/// <param name="logger"></param>
public CleanupJob(
IWorkflowInstanceStore workflowInstanceStore,
IOptions<CleanupOptions> options,
IServiceProvider serviceProvider,
ILogger<CleanupJob> logger)
{
_workflowInstanceStore = workflowInstanceStore;
_options = options.Value;
_serviceProvider = serviceProvider;
_logger = logger;
}
private readonly ILogger _logger = logger;
private readonly CleanupOptions _options = options.Value;
/// <summary>
/// Executes the cleanup job
@ -47,34 +31,28 @@ public class CleanupJob
/// <param name="cancellationToken"></param>
public async Task ExecuteAsync(CancellationToken cancellationToken = default)
{
await using AsyncServiceScope scope = _serviceProvider.CreateAsyncScope();
IEnumerable<IRetentionPolicy> policies = scope.ServiceProvider.GetServices<IRetentionPolicy>();
Dictionary<Type, object> collectors = GetServices(typeof(IRelatedEntityCollector), typeof(IRelatedEntityCollector<>));
foreach (IRetentionPolicy policy in policies)
var collectors = GetServices(typeof(IRelatedEntityCollector), typeof(IRelatedEntityCollector<>));
var deletedWorkflowInstances = 0L;
foreach (var policy in policies)
{
WorkflowInstanceFilter filter = policy.FilterFactory(scope.ServiceProvider).Build();
PageArgs pageArgs = PageArgs.FromPage(0, _options.PageSize);
long deletedWorkflowInstances = 0;
var filter = policy.FilterFactory(serviceProvider).Build();
var pageArgs = PageArgs.FromPage(0, _options.PageSize);
while (true)
{
Page<WorkflowInstance> page = await _workflowInstanceStore.FindManyAsync(filter, pageArgs, cancellationToken);
var page = await workflowInstanceStore.FindManyAsync(filter, pageArgs, cancellationToken);
if (page.Items.Count == 0)
{
break;
}
foreach (KeyValuePair<Type, object> collectorService in collectors)
foreach (var collectorService in collectors)
{
Type cleanupStrategyConcreteType = policy.CleanupStrategy.MakeGenericType(collectorService.Key);
IRelatedEntityCollector? collector = collectorService.Value as IRelatedEntityCollector;
ICleanupStrategy? cleanupService = _serviceProvider.GetService(cleanupStrategyConcreteType) as ICleanupStrategy;
var cleanupStrategyConcreteType = policy.CleanupStrategy.MakeGenericType(collectorService.Key);
var collector = collectorService.Value as IRelatedEntityCollector;
var cleanupService = serviceProvider.GetService(cleanupStrategyConcreteType) as ICleanupStrategy;
if (collector == null)
{
@ -88,13 +66,13 @@ public class CleanupJob
continue;
}
await foreach (ICollection<object> entities in collector.GetRelatedEntitiesGeneric(page.Items).WithCancellation(cancellationToken))
await foreach (var entities in collector.GetRelatedEntitiesGeneric(page.Items).WithCancellation(cancellationToken))
{
await cleanupService.Cleanup(entities);
}
}
deletedWorkflowInstances += await _workflowInstanceStore.DeleteAsync(new WorkflowInstanceFilter
deletedWorkflowInstances += await workflowInstanceStore.DeleteAsync(new WorkflowInstanceFilter
{
Ids = page.Items.Select(x => x.Id).ToArray()
}, cancellationToken);
@ -111,7 +89,7 @@ public class CleanupJob
private Dictionary<Type, object> GetServices(Type baseType, Type openType)
{
IEnumerable<object?> services = _serviceProvider.GetServices(baseType);
var services = serviceProvider.GetServices(baseType);
return services
.Where(x => x?.GetType() != null)

View file

@ -5,11 +5,6 @@ namespace Elsa.Retention.Options;
/// </summary>
public class CleanupOptions
{
/// <summary>
/// Controls how often the database is checked for workflow instances and execution log records to remove.
/// </summary>
public TimeSpan SweepInterval { get; set; } = TimeSpan.FromHours(4);
/// <summary>
/// Controls the page size of the workflow instance that are retained in a single batch
/// </summary>

View file

@ -6,16 +6,10 @@ namespace Elsa.Retention.Policies;
/// <summary>
/// A policy that will delete the workflow instance and its related entities
/// </summary>
public class DeletionRetentionPolicy : IRetentionPolicy
public class DeletionRetentionPolicy(string name, Func<IServiceProvider, RetentionWorkflowInstanceFilter> filter) : IRetentionPolicy
{
public DeletionRetentionPolicy(string name, Func<IServiceProvider, RetentionWorkflowInstanceFilter> filter)
{
Name = name;
FilterFactory = filter;
}
public string Name { get; }
public Func<IServiceProvider, RetentionWorkflowInstanceFilter> FilterFactory { get; }
public string Name { get; } = name;
public Func<IServiceProvider, RetentionWorkflowInstanceFilter> FilterFactory { get; } = filter;
public Type CleanupStrategy => typeof(IDeletionCleanupStrategy<>);
}

View file

@ -0,0 +1,19 @@
using Elsa.Common;
using Elsa.Common.RecurringTasks;
using Elsa.Retention.Jobs;
using JetBrains.Annotations;
namespace Elsa.Retention;
/// <summary>
/// Periodically deletes workflow instances and their execution logs.
/// </summary>
[SingleNodeTask]
[UsedImplicitly]
public class CleanupRecurringTask(CleanupJob job) : RecurringTask
{
public override async Task ExecuteAsync(CancellationToken stoppingToken)
{
await job.ExecuteAsync(stoppingToken);
}
}