From 892dfabeff7ce8d2adf8cbd95572c11bf44dfe1c Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Sun, 19 Jan 2020 21:47:38 +0100 Subject: [PATCH] Extract activity filtering to WorkflowInstanceStoreExtensions --- .../WorkflowInstanceStoreExtensions.cs | 28 ++++++++++++++++--- .../Persistence/IWorkflowInstanceStore.cs | 2 +- .../Elsa.Core/Services/WorkflowScheduler.cs | 15 ++-------- 3 files changed, 27 insertions(+), 18 deletions(-) diff --git a/src/core/Elsa.Abstractions/Extensions/WorkflowInstanceStoreExtensions.cs b/src/core/Elsa.Abstractions/Extensions/WorkflowInstanceStoreExtensions.cs index f0f0116d8..748aeceac 100644 --- a/src/core/Elsa.Abstractions/Extensions/WorkflowInstanceStoreExtensions.cs +++ b/src/core/Elsa.Abstractions/Extensions/WorkflowInstanceStoreExtensions.cs @@ -1,3 +1,4 @@ +using System; using System.Collections.Generic; using System.Linq; using System.Threading; @@ -10,13 +11,32 @@ namespace Elsa.Extensions { public static class WorkflowInstanceStoreExtensions { - public static async Task> ListByBlockingActivityAsync( + public static Task> ListByBlockingActivityAsync( this IWorkflowInstanceStore store, string? correlationId = default, - CancellationToken cancellationToken = default) where TActivity : IActivity + Func? activityStatePredicate = default, + CancellationToken cancellationToken = default) where TActivity : IActivity => + store.ListByBlockingActivityAsync(typeof(TActivity).Name, correlationId, activityStatePredicate, cancellationToken); + + public static async Task> ListByBlockingActivityAsync( + this IWorkflowInstanceStore store, + string activityType, + string? correlationId, + Func? activityStatePredicate = default, + CancellationToken cancellationToken = default) { - var items = await store.ListByBlockingActivityAsync(typeof(TActivity).Name, correlationId, cancellationToken); - return items.Select(x => (x.Item1, x.Item2)); + var tuples = await store.ListByBlockingActivityAsync(activityType, correlationId, cancellationToken); + + if (activityStatePredicate != null) + { + tuples = tuples.Where(tuple => + { + var activityInstance = tuple.WorkflowInstance.Activities.First(x => x.Id == tuple.BlockingActivity.ActivityId); + return activityStatePredicate(activityInstance.State); + }); + } + + return tuples; } } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Persistence/IWorkflowInstanceStore.cs b/src/core/Elsa.Abstractions/Persistence/IWorkflowInstanceStore.cs index 0b63892fd..33104ffe6 100644 --- a/src/core/Elsa.Abstractions/Persistence/IWorkflowInstanceStore.cs +++ b/src/core/Elsa.Abstractions/Persistence/IWorkflowInstanceStore.cs @@ -12,7 +12,7 @@ namespace Elsa.Persistence Task GetByCorrelationIdAsync(string correlationId, CancellationToken cancellationToken = default); Task> ListByDefinitionAsync(string definitionId, CancellationToken cancellationToken = default); Task> ListAllAsync(CancellationToken cancellationToken = default); - Task> ListByBlockingActivityAsync(string activityType, string? correlationId = default, CancellationToken cancellationToken = default); + Task> ListByBlockingActivityAsync(string activityType, string? correlationId = default, CancellationToken cancellationToken = default); Task> ListByStatusAsync(string definitionId, WorkflowStatus status, CancellationToken cancellationToken = default); Task> ListByStatusAsync(WorkflowStatus status, CancellationToken cancellationToken = default); Task DeleteAsync(string id, CancellationToken cancellationToken = default); diff --git a/src/core/Elsa.Core/Services/WorkflowScheduler.cs b/src/core/Elsa.Core/Services/WorkflowScheduler.cs index 68d77aa03..fabc27448 100644 --- a/src/core/Elsa.Core/Services/WorkflowScheduler.cs +++ b/src/core/Elsa.Core/Services/WorkflowScheduler.cs @@ -115,21 +115,10 @@ namespace Elsa.Services /// private async Task ScheduleSuspendedWorkflowsAsync(string activityType, object? input, string? correlationId, Func? activityStatePredicate, CancellationToken cancellationToken) { - var tuples = await workflowInstanceStore.ListByBlockingActivityAsync(activityType, correlationId, cancellationToken); - - foreach (var (workflowInstance, blockingActivity) in tuples) - { - var workflow = await workflowRegistry.GetWorkflowAsync(workflowInstance.DefinitionId, VersionOptions.SpecificVersion(workflowInstance.Version), cancellationToken); - var activity = workflow.GetActivity(blockingActivity.ActivityId); - - if (activityStatePredicate != null) - { - if (!activityStatePredicate(activity.State)) - continue; - } + var tuples = await workflowInstanceStore.ListByBlockingActivityAsync(activityType, correlationId, activityStatePredicate, cancellationToken); + foreach (var (workflowInstance, blockingActivity) in tuples) await ScheduleWorkflowAsync(workflowInstance.Id, blockingActivity.ActivityId, input, cancellationToken); - } } private async Task ScheduleWorkflowAsync(Workflow workflow, IActivity activity, object? input, string? correlationId, CancellationToken cancellationToken)