Add continue running workflow startup task

This commit is contained in:
Sipke Schoorstra 2021-01-19 20:17:14 +01:00
parent 860c045e6f
commit 2a64cc1473
10 changed files with 64 additions and 31 deletions

View file

@ -43,5 +43,6 @@ namespace Elsa.Models
public WorkflowFault? Fault { get; set; }
public Stack<ScheduledActivity> ScheduledActivities { get; set; }
public Stack<string> Scopes { get; set; }
public ScheduledActivity? CurrentActivity { get; set; }
}
}

View file

@ -1,5 +1,6 @@
using System;
using System.Diagnostics;
using System.Linq;
using System.Threading.Tasks;
using Elsa.DistributedLock;
using Elsa.Messages;
@ -39,13 +40,13 @@ namespace Elsa.Consumers
var workflowInstanceId = message.WorkflowInstanceId;
var lockKey = workflowInstanceId;
_logger.LogDebug("Acquiring lock on workflow instance {WorkflowInstanceId}.", workflowInstanceId);
_logger.LogDebug("Acquiring lock on workflow instance {WorkflowInstanceId}", workflowInstanceId);
_stopwatch.Restart();
if (!await _distributedLockProvider.AcquireLockAsync(lockKey))
{
// Reschedule message.
_logger.LogDebug("Failed to acquire lock on workflow instance {WorkflowInstanceId}. Rescheduling message.", workflowInstanceId);
_logger.LogDebug("Failed to acquire lock on workflow instance {WorkflowInstanceId}. Rescheduling message", workflowInstanceId);
await Task.Delay(TimeSpan.FromSeconds(1));
await _commandSender.SendAsync(message);
return;
@ -55,7 +56,7 @@ namespace Elsa.Consumers
{
var workflowInstance = await _workflowInstanceManager.FindByIdAsync(message.WorkflowInstanceId);
if (!ValidatePreconditions(workflowInstanceId, workflowInstance))
if (!ValidatePreconditions(workflowInstanceId, workflowInstance, message.ActivityId))
return;
await _workflowRunner.RunWorkflowAsync(
@ -67,24 +68,36 @@ namespace Elsa.Consumers
{
await _distributedLockProvider.ReleaseLockAsync(lockKey);
_stopwatch.Stop();
_logger.LogDebug("Held lock on workflow instance {WorkflowInstanceId} for {ElapsedTime}.", workflowInstanceId, _stopwatch.Elapsed);
_logger.LogDebug("Held lock on workflow instance {WorkflowInstanceId} for {ElapsedTime}", workflowInstanceId, _stopwatch.Elapsed);
}
}
private bool ValidatePreconditions(string? workflowInstanceId, WorkflowInstance? workflowInstance)
private bool ValidatePreconditions(string? workflowInstanceId, WorkflowInstance? workflowInstance, string? activityId)
{
if (workflowInstance == null)
{
_logger.LogWarning("Could not run workflow instance with ID {WorkflowInstanceId} because it does not exist.", workflowInstanceId);
_logger.LogWarning("Could not run workflow instance with ID {WorkflowInstanceId} because it does not exist", workflowInstanceId);
return false;
}
if (workflowInstance.WorkflowStatus != WorkflowStatus.Suspended)
if (workflowInstance.WorkflowStatus != WorkflowStatus.Suspended && workflowInstance.WorkflowStatus != WorkflowStatus.Running)
{
_logger.LogWarning("Could not run workflow instance with ID {WorkflowInstanceId} because it has a status other than Suspended. Its actual status is {WorkflowStatus}", workflowInstanceId, workflowInstance.WorkflowStatus);
_logger.LogWarning("Could not run workflow instance with ID {WorkflowInstanceId} because it has a status other than Suspended or Running. Its actual status is {WorkflowStatus}", workflowInstanceId, workflowInstance.WorkflowStatus);
return false;
}
if (activityId != null)
{
var activityIsBlocking = workflowInstance.BlockingActivities.Any(x => x.ActivityId == activityId);
var activityIsScheduled = workflowInstance.ScheduledActivities.Any(x => x.ActivityId == activityId) || workflowInstance.CurrentActivity?.ActivityId == activityId;
if (!activityIsBlocking && !activityIsScheduled)
{
_logger.LogWarning("Did not run workflow {WorkflowInstanceId} for activity {ActivityId} because the workflow is not blocked on that activity nor is that activity scheduled for execution", workflowInstanceId, activityId);
return false;
}
}
return true;
}
}

View file

@ -53,7 +53,8 @@ namespace Microsoft.Extensions.DependencyInjection
.AddSingleton(options.DistributedLockProviderFactory)
.AddSingleton(options.SignalFactory)
.AddSingleton(options.StorageFactory)
.AddStartupTask<CreateSubscriptions>();
.AddStartupTask<CreateSubscriptions>()
.AddStartupTask<ContinueRunningWorkflowsTask>();
options
.AddWorkflowsCore()

View file

@ -287,10 +287,11 @@ namespace Elsa.Services
{
var scope = workflowExecutionContext.ServiceScope;
var workflowBlueprint = workflowExecutionContext.WorkflowBlueprint;
var workflowInstance = workflowExecutionContext.WorkflowInstance;
while (workflowExecutionContext.HasScheduledActivities)
{
var scheduledActivity = workflowExecutionContext.PopScheduledActivity();
var scheduledActivity = workflowInstance.CurrentActivity = workflowExecutionContext.PopScheduledActivity();
var currentActivityId = scheduledActivity.ActivityId;
var activityBlueprint = workflowBlueprint.GetActivity(currentActivityId)!;
var activityExecutionContext = new ActivityExecutionContext(scope, workflowExecutionContext, activityBlueprint, scheduledActivity.Input, cancellationToken);
@ -314,6 +315,8 @@ namespace Elsa.Services
activityOperation = Execute;
}
workflowInstance.CurrentActivity = null;
if (workflowExecutionContext.HasBlockingActivities)
workflowExecutionContext.Suspend();
@ -329,7 +332,7 @@ namespace Elsa.Services
}
catch (Exception e)
{
_logger.LogError(e, e.Message);
_logger.LogWarning(e, "An error occurred while executing activity {ActivityId}. Entering Faulted state", activity.Id);
activityExecutionContext.WorkflowExecutionContext.Fault(activity.Id, e.Message, e.StackTrace);
}

View file

@ -1,4 +1,3 @@
using System;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
@ -16,18 +15,18 @@ namespace Elsa.StartupTasks
/// If there are workflows in the Running state while the server starts, it means the workflow instance never finished execution, e.g. because the workflow host terminated.
/// This startup task resumes these workflows.
/// </summary>
public class ResumeRunningWorkflowsTask : IStartupTask
public class ContinueRunningWorkflowsTask : IStartupTask
{
private readonly IWorkflowInstanceStore _workflowInstanceStore;
private readonly IWorkflowRunner _workflowScheduler;
private readonly IWorkflowQueue _workflowScheduler;
private readonly IDistributedLockProvider _distributedLockProvider;
private readonly ILogger<ResumeRunningWorkflowsTask> _logger;
private readonly ILogger<ContinueRunningWorkflowsTask> _logger;
public ResumeRunningWorkflowsTask(
public ContinueRunningWorkflowsTask(
IWorkflowInstanceStore workflowInstanceStore,
IWorkflowRunner workflowScheduler,
IWorkflowQueue workflowScheduler,
IDistributedLockProvider distributedLockProvider,
ILogger<ResumeRunningWorkflowsTask> logger)
ILogger<ContinueRunningWorkflowsTask> logger)
{
_workflowInstanceStore = workflowInstanceStore;
_workflowScheduler = workflowScheduler;
@ -38,22 +37,34 @@ namespace Elsa.StartupTasks
public async Task ExecuteAsync(CancellationToken cancellationToken = default)
{
var lockKey = GetType().Name;
if (!await _distributedLockProvider.AcquireLockAsync(lockKey, cancellationToken))
return;
try
{
var instances = await _workflowInstanceStore.FindManyAsync(new WorkflowStatusSpecification(WorkflowStatus.Running), cancellationToken: cancellationToken).ToList();
_logger.LogInformation("Found {WorkflowInstanceCount} workflows with status 'Running'. Resuming each one of them.", instances.Count);
_logger.LogInformation("Found {WorkflowInstanceCount} workflows with status 'Running'. Resuming each one of them", instances.Count);
foreach (var instance in instances)
{
_logger.LogInformation("Resuming {WorkflowInstanceId}", instance.Id);
await _workflowScheduler.RunWorkflowAsync(
instance,
cancellationToken: cancellationToken);
var scheduledActivities = instance.ScheduledActivities;
if (instance.CurrentActivity == null || !scheduledActivities.Any())
{
_logger.LogWarning("Workflow '{WorkflowInstanceId}' was in the Running state, but has no scheduled activities nor has a currently executing one", instance.Id);
continue;
}
var scheduledActivity = instance.CurrentActivity ?? instance.ScheduledActivities.Peek();
await _workflowScheduler.EnqueueWorkflowInstance(
instance.Id,
scheduledActivity.ActivityId,
scheduledActivity.Input,
cancellationToken);
}
}
finally

View file

@ -21,6 +21,7 @@ namespace Elsa.Persistence.EntityFramework.Core.Configuration
builder.Ignore(x => x.ScheduledActivities);
builder.Ignore(x => x.Scopes);
builder.Ignore(x => x.Variables);
builder.Ignore(x => x.CurrentActivity);
builder.Property<string>("Data");
}
}

View file

@ -36,7 +36,8 @@ namespace Elsa.Persistence.EntityFramework.Core.Stores
entity.BlockingActivities,
entity.ScheduledActivities,
entity.Scopes,
entity.Fault
entity.Fault,
entity.CurrentActivity
};
var json = _contentSerializer.Serialize(data);
@ -55,7 +56,8 @@ namespace Elsa.Persistence.EntityFramework.Core.Stores
entity.BlockingActivities,
entity.ScheduledActivities,
entity.Scopes,
entity.Fault
entity.Fault,
entity.CurrentActivity
};
var json = (string)DbContext.Entry(entity).Property("Data").CurrentValue;
@ -71,6 +73,7 @@ namespace Elsa.Persistence.EntityFramework.Core.Stores
entity.ScheduledActivities = data.ScheduledActivities;
entity.Scopes = data.Scopes ?? new Stack<string>();
entity.Fault = data.Fault;
entity.CurrentActivity = data.CurrentActivity;
}
}
}

View file

@ -37,5 +37,7 @@ namespace Elsa.Persistence.YesSql.Documents
public WorkflowFault? Fault { get; set; }
public Stack<ScheduledActivity> ScheduledActivities { get; set; } = new();
public Stack<string> Scopes { get; set; }
public ScheduledActivity? CurrentActivity { get; set; }
}
}

View file

@ -40,8 +40,7 @@ namespace ElsaDashboard.Samples.Monolith
services
.AddElsaApiEndpoints()
.AddElsaSwagger()
.AddStartupTask<ResumeRunningWorkflowsTask>();
.AddElsaSwagger();
// Elsa Dashboard.
services.AddRazorPages();

View file

@ -36,8 +36,7 @@ namespace Elsa.Samples.Server.Host
services
.AddElsaApiEndpoints()
.AddElsaSwagger()
.AddStartupTask<ResumeRunningWorkflowsTask>();
.AddElsaSwagger();
}
public void Configure(IApplicationBuilder app)