Extract activity filtering to WorkflowInstanceStoreExtensions
This commit is contained in:
parent
e7bf3f5ab6
commit
892dfabeff
|
|
@ -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<IEnumerable<(WorkflowInstance, BlockingActivity)>> ListByBlockingActivityAsync<TActivity>(
|
||||
public static Task<IEnumerable<(WorkflowInstance WorkflowInstance, BlockingActivity BlockingActivity)>> ListByBlockingActivityAsync<TActivity>(
|
||||
this IWorkflowInstanceStore store,
|
||||
string? correlationId = default,
|
||||
CancellationToken cancellationToken = default) where TActivity : IActivity
|
||||
Func<Variables, bool>? activityStatePredicate = default,
|
||||
CancellationToken cancellationToken = default) where TActivity : IActivity =>
|
||||
store.ListByBlockingActivityAsync(typeof(TActivity).Name, correlationId, activityStatePredicate, cancellationToken);
|
||||
|
||||
public static async Task<IEnumerable<(WorkflowInstance WorkflowInstance, BlockingActivity BLockingActivity)>> ListByBlockingActivityAsync(
|
||||
this IWorkflowInstanceStore store,
|
||||
string activityType,
|
||||
string? correlationId,
|
||||
Func<Variables, bool>? 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;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -12,7 +12,7 @@ namespace Elsa.Persistence
|
|||
Task<WorkflowInstance> GetByCorrelationIdAsync(string correlationId, CancellationToken cancellationToken = default);
|
||||
Task<IEnumerable<WorkflowInstance>> ListByDefinitionAsync(string definitionId, CancellationToken cancellationToken = default);
|
||||
Task<IEnumerable<WorkflowInstance>> ListAllAsync(CancellationToken cancellationToken = default);
|
||||
Task<IEnumerable<(WorkflowInstance, BlockingActivity)>> ListByBlockingActivityAsync(string activityType, string? correlationId = default, CancellationToken cancellationToken = default);
|
||||
Task<IEnumerable<(WorkflowInstance WorkflowInstance, BlockingActivity BlockingActivity)>> ListByBlockingActivityAsync(string activityType, string? correlationId = default, CancellationToken cancellationToken = default);
|
||||
Task<IEnumerable<WorkflowInstance>> ListByStatusAsync(string definitionId, WorkflowStatus status, CancellationToken cancellationToken = default);
|
||||
Task<IEnumerable<WorkflowInstance>> ListByStatusAsync(WorkflowStatus status, CancellationToken cancellationToken = default);
|
||||
Task DeleteAsync(string id, CancellationToken cancellationToken = default);
|
||||
|
|
|
|||
|
|
@ -115,21 +115,10 @@ namespace Elsa.Services
|
|||
/// </summary>
|
||||
private async Task ScheduleSuspendedWorkflowsAsync(string activityType, object? input, string? correlationId, Func<Variables, bool>? 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)
|
||||
|
|
|
|||
Loading…
Reference in a new issue