Delete ActivityExecutionRecords when deleting workflow instances (#4324)
* Delete ActivityExecutionRecords when deleting workflow instances * Enable Hangfire for testing, code formatting
This commit is contained in:
parent
b829a0acef
commit
af34276d60
|
|
@ -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<Program>()
|
||||
.AddWorkflowsFrom<Program>()
|
||||
|
|
@ -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<Program>())
|
||||
.UseRealTimeWorkflows()
|
||||
.UseJavaScript()
|
||||
|
|
|
|||
|
|
@ -59,11 +59,17 @@ public class DapperActivityExecutionRecordStore : IActivityExecutionStore
|
|||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public Task<long> CountAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default)
|
||||
public async Task<long> CountAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default)
|
||||
{
|
||||
throw new NotImplementedException();
|
||||
return await _store.CountAsync(q => ApplyFilter(q, filter), cancellationToken);
|
||||
}
|
||||
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<long> DeleteManyAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default)
|
||||
{
|
||||
return await _store.DeleteAsync(q => ApplyFilter(q, filter), cancellationToken);
|
||||
}
|
||||
|
||||
private static void ApplyFilter(ParameterizedQuery query, ActivityExecutionRecordFilter filter)
|
||||
{
|
||||
query
|
||||
|
|
|
|||
|
|
@ -43,6 +43,9 @@ public class EFCoreActivityExecutionStore : IActivityExecutionStore
|
|||
/// <inheritdoc />
|
||||
public async Task<long> CountAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default) => await _store.CountAsync(queryable => Filter(queryable, filter), cancellationToken);
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<long> 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;
|
||||
|
|
|
|||
|
|
@ -53,6 +53,12 @@ public class MongoActivityExecutionLogStore : IActivityExecutionStore
|
|||
return await _mongoDbStore.CountAsync(queryable => Filter(queryable, filter).OrderBy(x => x.StartedAt), cancellationToken);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public async Task<long> DeleteManyAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default)
|
||||
{
|
||||
return await _mongoDbStore.DeleteWhereAsync<string>(queryable => Filter(queryable, filter), x => x.Id, cancellationToken);
|
||||
}
|
||||
|
||||
private IMongoQueryable<ActivityExecutionRecord> Filter(IMongoQueryable<ActivityExecutionRecord> queryable, ActivityExecutionRecordFilter filter) =>
|
||||
(filter.Apply(queryable) as IMongoQueryable<ActivityExecutionRecord>)!;
|
||||
|
||||
|
|
|
|||
|
|
@ -10,7 +10,7 @@ namespace Elsa.Scheduling.Handlers;
|
|||
/// <summary>
|
||||
/// Deletes scheduled jobs based on deleted workflow triggers and bookmarks.
|
||||
/// </summary>
|
||||
public class DeleteSchedules :
|
||||
public class DeleteSchedules :
|
||||
INotificationHandler<WorkflowDefinitionDeleting>,
|
||||
INotificationHandler<WorkflowDefinitionsDeleting>,
|
||||
INotificationHandler<WorkflowDefinitionVersionDeleting>,
|
||||
|
|
@ -41,8 +41,11 @@ public class DeleteSchedules :
|
|||
}
|
||||
|
||||
async Task INotificationHandler<WorkflowDefinitionDeleting>.HandleAsync(WorkflowDefinitionDeleting notification, CancellationToken cancellationToken) => await RemoveSchedulesAsync(new TriggerFilter { WorkflowDefinitionId = notification.DefinitionId }, cancellationToken);
|
||||
async Task INotificationHandler<WorkflowDefinitionsDeleting>.HandleAsync(WorkflowDefinitionsDeleting notification, CancellationToken cancellationToken) => await RemoveSchedulesAsync(new TriggerFilter{ WorkflowDefinitionIds = notification.DefinitionIds }, cancellationToken);
|
||||
async Task INotificationHandler<WorkflowDefinitionVersionDeleting>.HandleAsync(WorkflowDefinitionVersionDeleting notification, CancellationToken cancellationToken) => await RemoveSchedulesAsync(new TriggerFilter { WorkflowDefinitionVersionId = notification.WorkflowDefinition.Id }, cancellationToken);
|
||||
async Task INotificationHandler<WorkflowDefinitionsDeleting>.HandleAsync(WorkflowDefinitionsDeleting notification, CancellationToken cancellationToken) => await RemoveSchedulesAsync(new TriggerFilter { WorkflowDefinitionIds = notification.DefinitionIds }, cancellationToken);
|
||||
|
||||
async Task INotificationHandler<WorkflowDefinitionVersionDeleting>.HandleAsync(WorkflowDefinitionVersionDeleting notification, CancellationToken cancellationToken) =>
|
||||
await RemoveSchedulesAsync(new TriggerFilter { WorkflowDefinitionVersionId = notification.WorkflowDefinition.Id }, cancellationToken);
|
||||
|
||||
async Task INotificationHandler<WorkflowDefinitionVersionsDeleting>.HandleAsync(WorkflowDefinitionVersionsDeleting notification, CancellationToken cancellationToken) => await RemoveSchedulesAsync(new TriggerFilter { WorkflowDefinitionVersionIds = notification.Ids }, cancellationToken);
|
||||
|
||||
private async Task RemoveSchedulesAsync(TriggerFilter filter, CancellationToken cancellationToken)
|
||||
|
|
|
|||
|
|
@ -48,4 +48,12 @@ public interface IActivityExecutionStore
|
|||
/// <param name="cancellationToken">An optional cancellation token.</param>
|
||||
/// <returns>The number of activity execution records matching the specified filter.</returns>
|
||||
Task<long> CountAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default);
|
||||
|
||||
/// <summary>
|
||||
/// Deletes all activity execution records matching the specified filter.
|
||||
/// </summary>
|
||||
/// <param name="filter">The filter.</param>
|
||||
/// <param name="cancellationToken">An optional cancellation token.</param>
|
||||
/// <returns>The number of deleted records.</returns>
|
||||
Task<long> DeleteManyAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default);
|
||||
}
|
||||
|
|
@ -204,6 +204,7 @@ public class WorkflowRuntimeFeature : FeatureBase
|
|||
.AddNotificationHandler<DeleteBookmarks>()
|
||||
.AddNotificationHandler<DeleteTriggers>()
|
||||
.AddNotificationHandler<DeleteWorkflowInstances>()
|
||||
.AddNotificationHandler<DeleteActivityExecutionLogRecords>()
|
||||
.AddNotificationHandler<ReadWorkflowInboxMessage>()
|
||||
.AddNotificationHandler<DeliverWorkflowMessagesFromInbox>()
|
||||
|
||||
|
|
|
|||
|
|
@ -12,6 +12,11 @@ public class ActivityExecutionRecordFilter
|
|||
/// </summary>
|
||||
public string? WorkflowInstanceId { get; set; }
|
||||
|
||||
/// <summary>
|
||||
/// The IDs of the workflow instances.
|
||||
/// </summary>
|
||||
public ICollection<string>? WorkflowInstanceIds { get; set; }
|
||||
|
||||
/// <summary>
|
||||
/// The ID of the activity.
|
||||
/// </summary>
|
||||
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
|
||||
/// <summary>
|
||||
/// Deletes activity execution log records in response to the <see cref="WorkflowInstancesDeleting"/> notification.
|
||||
/// </summary>
|
||||
[PublicAPI]
|
||||
internal class DeleteActivityExecutionLogRecords : INotificationHandler<WorkflowInstancesDeleting>
|
||||
{
|
||||
private readonly IActivityExecutionStore _store;
|
||||
|
||||
/// <summary>
|
||||
/// Initializes a new instance of the <see cref="DeleteActivityExecutionLogRecords"/> class.
|
||||
/// </summary>
|
||||
public DeleteActivityExecutionLogRecords(IActivityExecutionStore store)
|
||||
{
|
||||
_store = store;
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
|
@ -57,5 +57,13 @@ public class MemoryActivityExecutionStore : IActivityExecutionStore
|
|||
return Task.FromResult(count);
|
||||
}
|
||||
|
||||
/// <inheritdoc />
|
||||
public Task<long> 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<ActivityExecutionRecord> Filter(IQueryable<ActivityExecutionRecord> queryable, ActivityExecutionRecordFilter filter) => filter.Apply(queryable);
|
||||
}
|
||||
|
|
@ -24,4 +24,7 @@ public class NoopActivityExecutionStore : IActivityExecutionStore
|
|||
|
||||
/// <inheritdoc />
|
||||
public Task<long> CountAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default) => Task.FromResult(0L);
|
||||
|
||||
/// <inheritdoc />
|
||||
public Task<long> DeleteManyAsync(ActivityExecutionRecordFilter filter, CancellationToken cancellationToken = default) => Task.FromResult(0L);
|
||||
}
|
||||
Loading…
Reference in a new issue