From af34276d6020521eaebdbfb5bca01e66e9724b80 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Sat, 12 Aug 2023 22:22:40 +0200 Subject: [PATCH] Delete ActivityExecutionRecords when deleting workflow instances (#4324) * Delete ActivityExecutionRecords when deleting workflow instances * Enable Hangfire for testing, code formatting --- .../Elsa.WorkflowServer.Web/Program.cs | 10 +++++- .../DapperActivityExecutionRecordStore.cs | 12 +++++-- .../Runtime/ActivityExecutionLogStore.cs | 3 ++ .../Runtime/ActivityExecutionLogStore.cs | 6 ++++ .../Handlers/DeleteSchedules.cs | 9 +++-- .../Contracts/IActivityExecutionStore.cs | 8 +++++ .../Features/WorkflowRuntimeFeature.cs | 1 + .../Filters/ActivityExecutionRecordFilter.cs | 6 ++++ .../DeleteActivityExecutionLogRecords.cs | 35 +++++++++++++++++++ .../Stores/MemoryActivityExecutionStore.cs | 8 +++++ .../Stores/NoopActivityExecutionStore.cs | 3 ++ 11 files changed, 94 insertions(+), 7 deletions(-) create mode 100644 src/modules/Elsa.Workflows.Runtime/Handlers/DeleteActivityExecutionLogRecords.cs diff --git a/src/bundles/Elsa.WorkflowServer.Web/Program.cs b/src/bundles/Elsa.WorkflowServer.Web/Program.cs index df4b2793e..568ea02ca 100644 --- a/src/bundles/Elsa.WorkflowServer.Web/Program.cs +++ b/src/bundles/Elsa.WorkflowServer.Web/Program.cs @@ -14,6 +14,7 @@ using Proto.Persistence.Sqlite; const bool useMongoDb = false; const bool useProtoActor = false; +const bool useHangfire = true; var builder = WebApplication.CreateBuilder(args); var services = builder.Services; @@ -30,6 +31,9 @@ services if(useMongoDb) elsa.UseMongoDb(mongoDbConnectionString); + if (useHangfire) + elsa.UseHangfire(); + elsa .AddActivitiesFrom() .AddWorkflowsFrom() @@ -88,7 +92,11 @@ services runtime.UseMassTransitDispatcher(); }) .UseEnvironments(environments => environments.EnvironmentsOptions = options => configuration.GetSection("Environments").Bind(options)) - .UseScheduling() + .UseScheduling(scheduling => + { + if (useHangfire) + scheduling.UseHangfireScheduler(); + }) .UseWorkflowsApi(api => api.AddFastEndpointsAssembly()) .UseRealTimeWorkflows() .UseJavaScript() diff --git a/src/modules/Elsa.Dapper/Modules/Runtime/Stores/DapperActivityExecutionRecordStore.cs b/src/modules/Elsa.Dapper/Modules/Runtime/Stores/DapperActivityExecutionRecordStore.cs index b37d0ef0a..9bfbd5fb4 100644 --- a/src/modules/Elsa.Dapper/Modules/Runtime/Stores/DapperActivityExecutionRecordStore.cs +++ b/src/modules/Elsa.Dapper/Modules/Runtime/Stores/DapperActivityExecutionRecordStore.cs @@ -59,11 +59,17 @@ public class DapperActivityExecutionRecordStore : IActivityExecutionStore } /// - public Task CountAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default) + public async Task CountAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default) { - throw new NotImplementedException(); + return await _store.CountAsync(q => ApplyFilter(q, filter), cancellationToken); } - + + /// + public async Task DeleteManyAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default) + { + return await _store.DeleteAsync(q => ApplyFilter(q, filter), cancellationToken); + } + private static void ApplyFilter(ParameterizedQuery query, ActivityExecutionRecordFilter filter) { query diff --git a/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/ActivityExecutionLogStore.cs b/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/ActivityExecutionLogStore.cs index bb03ddcc4..349f67c7e 100644 --- a/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/ActivityExecutionLogStore.cs +++ b/src/modules/Elsa.EntityFrameworkCore/Modules/Runtime/ActivityExecutionLogStore.cs @@ -43,6 +43,9 @@ public class EFCoreActivityExecutionStore : IActivityExecutionStore /// public async Task CountAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default) => await _store.CountAsync(queryable => Filter(queryable, filter), cancellationToken); + /// + public async Task DeleteManyAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default) => await _store.DeleteWhereAsync(queryable => Filter(queryable, filter), cancellationToken); + private ValueTask OnSaveAsync(RuntimeElsaDbContext dbContext, ActivityExecutionRecord entity, CancellationToken cancellationToken) { dbContext.Entry(entity).Property("ActivityData").CurrentValue = entity.ActivityState != null ? _serializer.Serialize(entity.ActivityState) : default; diff --git a/src/modules/Elsa.MongoDb/Modules/Runtime/ActivityExecutionLogStore.cs b/src/modules/Elsa.MongoDb/Modules/Runtime/ActivityExecutionLogStore.cs index 286104c42..b729ceaa7 100644 --- a/src/modules/Elsa.MongoDb/Modules/Runtime/ActivityExecutionLogStore.cs +++ b/src/modules/Elsa.MongoDb/Modules/Runtime/ActivityExecutionLogStore.cs @@ -53,6 +53,12 @@ public class MongoActivityExecutionLogStore : IActivityExecutionStore return await _mongoDbStore.CountAsync(queryable => Filter(queryable, filter).OrderBy(x => x.StartedAt), cancellationToken); } + /// + public async Task DeleteManyAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default) + { + return await _mongoDbStore.DeleteWhereAsync(queryable => Filter(queryable, filter), x => x.Id, cancellationToken); + } + private IMongoQueryable Filter(IMongoQueryable queryable, ActivityExecutionRecordFilter filter) => (filter.Apply(queryable) as IMongoQueryable)!; diff --git a/src/modules/Elsa.Scheduling/Handlers/DeleteSchedules.cs b/src/modules/Elsa.Scheduling/Handlers/DeleteSchedules.cs index 8aea9a6f8..96e5b0be2 100644 --- a/src/modules/Elsa.Scheduling/Handlers/DeleteSchedules.cs +++ b/src/modules/Elsa.Scheduling/Handlers/DeleteSchedules.cs @@ -10,7 +10,7 @@ namespace Elsa.Scheduling.Handlers; /// /// Deletes scheduled jobs based on deleted workflow triggers and bookmarks. /// -public class DeleteSchedules : +public class DeleteSchedules : INotificationHandler, INotificationHandler, INotificationHandler, @@ -41,8 +41,11 @@ public class DeleteSchedules : } async Task INotificationHandler.HandleAsync(WorkflowDefinitionDeleting notification, CancellationToken cancellationToken) => await RemoveSchedulesAsync(new TriggerFilter { WorkflowDefinitionId = notification.DefinitionId }, cancellationToken); - async Task INotificationHandler.HandleAsync(WorkflowDefinitionsDeleting notification, CancellationToken cancellationToken) => await RemoveSchedulesAsync(new TriggerFilter{ WorkflowDefinitionIds = notification.DefinitionIds }, cancellationToken); - async Task INotificationHandler.HandleAsync(WorkflowDefinitionVersionDeleting notification, CancellationToken cancellationToken) => await RemoveSchedulesAsync(new TriggerFilter { WorkflowDefinitionVersionId = notification.WorkflowDefinition.Id }, cancellationToken); + async Task INotificationHandler.HandleAsync(WorkflowDefinitionsDeleting notification, CancellationToken cancellationToken) => await RemoveSchedulesAsync(new TriggerFilter { WorkflowDefinitionIds = notification.DefinitionIds }, cancellationToken); + + async Task INotificationHandler.HandleAsync(WorkflowDefinitionVersionDeleting notification, CancellationToken cancellationToken) => + await RemoveSchedulesAsync(new TriggerFilter { WorkflowDefinitionVersionId = notification.WorkflowDefinition.Id }, cancellationToken); + async Task INotificationHandler.HandleAsync(WorkflowDefinitionVersionsDeleting notification, CancellationToken cancellationToken) => await RemoveSchedulesAsync(new TriggerFilter { WorkflowDefinitionVersionIds = notification.Ids }, cancellationToken); private async Task RemoveSchedulesAsync(TriggerFilter filter, CancellationToken cancellationToken) diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IActivityExecutionStore.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IActivityExecutionStore.cs index 26daaea6a..77238db10 100644 --- a/src/modules/Elsa.Workflows.Runtime/Contracts/IActivityExecutionStore.cs +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IActivityExecutionStore.cs @@ -48,4 +48,12 @@ public interface IActivityExecutionStore /// An optional cancellation token. /// The number of activity execution records matching the specified filter. Task CountAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default); + + /// + /// Deletes all activity execution records matching the specified filter. + /// + /// The filter. + /// An optional cancellation token. + /// The number of deleted records. + Task DeleteManyAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs index e1b00fcf3..ab7016d5c 100644 --- a/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs +++ b/src/modules/Elsa.Workflows.Runtime/Features/WorkflowRuntimeFeature.cs @@ -204,6 +204,7 @@ public class WorkflowRuntimeFeature : FeatureBase .AddNotificationHandler() .AddNotificationHandler() .AddNotificationHandler() + .AddNotificationHandler() .AddNotificationHandler() .AddNotificationHandler() diff --git a/src/modules/Elsa.Workflows.Runtime/Filters/ActivityExecutionRecordFilter.cs b/src/modules/Elsa.Workflows.Runtime/Filters/ActivityExecutionRecordFilter.cs index 51f7495f1..62922be5a 100644 --- a/src/modules/Elsa.Workflows.Runtime/Filters/ActivityExecutionRecordFilter.cs +++ b/src/modules/Elsa.Workflows.Runtime/Filters/ActivityExecutionRecordFilter.cs @@ -12,6 +12,11 @@ public class ActivityExecutionRecordFilter /// public string? WorkflowInstanceId { get; set; } + /// + /// The IDs of the workflow instances. + /// + public ICollection? WorkflowInstanceIds { get; set; } + /// /// The ID of the activity. /// @@ -34,6 +39,7 @@ public class ActivityExecutionRecordFilter { var filter = this; if (filter.WorkflowInstanceId != null) queryable = queryable.Where(x => x.WorkflowInstanceId == filter.WorkflowInstanceId); + if (filter.WorkflowInstanceIds != null) queryable = queryable.Where(x => filter.WorkflowInstanceIds.Contains(x.WorkflowInstanceId)); if (filter.ActivityId != null) queryable = queryable.Where(x => x.ActivityId == filter.ActivityId); if (filter.ActivityIds != null && filter.ActivityIds.Any()) queryable = queryable.Where(x => filter.ActivityIds.Contains(x.ActivityId)); if (filter.Completed != null) queryable = filter.Completed == true ? queryable.Where(x => x.CompletedAt != null) : queryable.Where(x => x.CompletedAt == null); diff --git a/src/modules/Elsa.Workflows.Runtime/Handlers/DeleteActivityExecutionLogRecords.cs b/src/modules/Elsa.Workflows.Runtime/Handlers/DeleteActivityExecutionLogRecords.cs new file mode 100644 index 000000000..b5dd6a86f --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime/Handlers/DeleteActivityExecutionLogRecords.cs @@ -0,0 +1,35 @@ +using Elsa.Mediator.Contracts; +using Elsa.Workflows.Management.Notifications; +using Elsa.Workflows.Runtime.Contracts; +using Elsa.Workflows.Runtime.Filters; +using JetBrains.Annotations; + +namespace Elsa.Workflows.Runtime.Handlers; + +/// +/// Deletes activity execution log records in response to the notification. +/// +[PublicAPI] +internal class DeleteActivityExecutionLogRecords : INotificationHandler +{ + private readonly IActivityExecutionStore _store; + + /// + /// Initializes a new instance of the class. + /// + public DeleteActivityExecutionLogRecords(IActivityExecutionStore store) + { + _store = store; + } + + /// + public async Task HandleAsync(WorkflowInstancesDeleting notification, CancellationToken cancellationToken) + { + await DeleteManyAsync(new ActivityExecutionRecordFilter { WorkflowInstanceIds = notification.Ids }, cancellationToken); + } + + private async Task DeleteManyAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default) + { + await _store.DeleteManyAsync(filter, cancellationToken); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Stores/MemoryActivityExecutionStore.cs b/src/modules/Elsa.Workflows.Runtime/Stores/MemoryActivityExecutionStore.cs index c173ba49c..e0cfe23bd 100644 --- a/src/modules/Elsa.Workflows.Runtime/Stores/MemoryActivityExecutionStore.cs +++ b/src/modules/Elsa.Workflows.Runtime/Stores/MemoryActivityExecutionStore.cs @@ -57,5 +57,13 @@ public class MemoryActivityExecutionStore : IActivityExecutionStore return Task.FromResult(count); } + /// + public Task DeleteManyAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default) + { + var records = _store.Query(query => Filter(query, filter)).ToList(); + _store.DeleteMany(records, x => x.Id); + return Task.FromResult(records.LongCount()); + } + private static IQueryable Filter(IQueryable queryable, ActivityExecutionRecordFilter filter) => filter.Apply(queryable); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Stores/NoopActivityExecutionStore.cs b/src/modules/Elsa.Workflows.Runtime/Stores/NoopActivityExecutionStore.cs index b946cf21c..1a4124a39 100644 --- a/src/modules/Elsa.Workflows.Runtime/Stores/NoopActivityExecutionStore.cs +++ b/src/modules/Elsa.Workflows.Runtime/Stores/NoopActivityExecutionStore.cs @@ -24,4 +24,7 @@ public class NoopActivityExecutionStore : IActivityExecutionStore /// public Task CountAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default) => Task.FromResult(0L); + + /// + public Task DeleteManyAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default) => Task.FromResult(0L); } \ No newline at end of file