From 2a64cc1473602e5d76c0f06ff672c63830dfd6a0 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Tue, 19 Jan 2021 20:17:14 +0100 Subject: [PATCH] Add continue running workflow startup task --- .../Models/WorkflowInstance.cs | 1 + .../Consumers/RunWorkflowInstanceConsumer.cs | 29 ++++++++++---- .../ElsaServiceCollectionExtensions.cs | 3 +- src/core/Elsa.Core/Services/WorkflowRunner.cs | 7 +++- ...ask.cs => ContinueRunningWorkflowsTask.cs} | 39 ++++++++++++------- .../WorkflowInstanceConfiguration.cs | 1 + .../EntityFrameworkWorkflowInstanceStore.cs | 7 +++- .../Documents/WorkflowInstanceDocument.cs | 2 + .../ElsaDashboard.Samples.Monolith/Startup.cs | 3 +- .../Elsa.Samples.Server.Host/Startup.cs | 3 +- 10 files changed, 64 insertions(+), 31 deletions(-) rename src/core/Elsa.Core/StartupTasks/{ResumeRunningWorkflowsTask.cs => ContinueRunningWorkflowsTask.cs} (63%) diff --git a/src/core/Elsa.Abstractions/Models/WorkflowInstance.cs b/src/core/Elsa.Abstractions/Models/WorkflowInstance.cs index 8e9135134..cba7e02b7 100644 --- a/src/core/Elsa.Abstractions/Models/WorkflowInstance.cs +++ b/src/core/Elsa.Abstractions/Models/WorkflowInstance.cs @@ -43,5 +43,6 @@ namespace Elsa.Models public WorkflowFault? Fault { get; set; } public Stack ScheduledActivities { get; set; } public Stack Scopes { get; set; } + public ScheduledActivity? CurrentActivity { get; set; } } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Consumers/RunWorkflowInstanceConsumer.cs b/src/core/Elsa.Core/Consumers/RunWorkflowInstanceConsumer.cs index b60c253c4..7ffa22dd0 100644 --- a/src/core/Elsa.Core/Consumers/RunWorkflowInstanceConsumer.cs +++ b/src/core/Elsa.Core/Consumers/RunWorkflowInstanceConsumer.cs @@ -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; } } diff --git a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs index ea39314ba..d22c57912 100644 --- a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs +++ b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs @@ -53,7 +53,8 @@ namespace Microsoft.Extensions.DependencyInjection .AddSingleton(options.DistributedLockProviderFactory) .AddSingleton(options.SignalFactory) .AddSingleton(options.StorageFactory) - .AddStartupTask(); + .AddStartupTask() + .AddStartupTask(); options .AddWorkflowsCore() diff --git a/src/core/Elsa.Core/Services/WorkflowRunner.cs b/src/core/Elsa.Core/Services/WorkflowRunner.cs index b5b99ef0e..726d925b7 100644 --- a/src/core/Elsa.Core/Services/WorkflowRunner.cs +++ b/src/core/Elsa.Core/Services/WorkflowRunner.cs @@ -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); } diff --git a/src/core/Elsa.Core/StartupTasks/ResumeRunningWorkflowsTask.cs b/src/core/Elsa.Core/StartupTasks/ContinueRunningWorkflowsTask.cs similarity index 63% rename from src/core/Elsa.Core/StartupTasks/ResumeRunningWorkflowsTask.cs rename to src/core/Elsa.Core/StartupTasks/ContinueRunningWorkflowsTask.cs index e186bea5f..5ecd9b76c 100644 --- a/src/core/Elsa.Core/StartupTasks/ResumeRunningWorkflowsTask.cs +++ b/src/core/Elsa.Core/StartupTasks/ContinueRunningWorkflowsTask.cs @@ -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. /// - 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 _logger; + private readonly ILogger _logger; - public ResumeRunningWorkflowsTask( + public ContinueRunningWorkflowsTask( IWorkflowInstanceStore workflowInstanceStore, - IWorkflowRunner workflowScheduler, + IWorkflowQueue workflowScheduler, IDistributedLockProvider distributedLockProvider, - ILogger logger) + ILogger 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 diff --git a/src/persistence/Elsa.Persistence.EntityFramework/Elsa.Persistence.EntityFramework.Core/Configuration/WorkflowInstanceConfiguration.cs b/src/persistence/Elsa.Persistence.EntityFramework/Elsa.Persistence.EntityFramework.Core/Configuration/WorkflowInstanceConfiguration.cs index 0b7e20e04..0b3088504 100644 --- a/src/persistence/Elsa.Persistence.EntityFramework/Elsa.Persistence.EntityFramework.Core/Configuration/WorkflowInstanceConfiguration.cs +++ b/src/persistence/Elsa.Persistence.EntityFramework/Elsa.Persistence.EntityFramework.Core/Configuration/WorkflowInstanceConfiguration.cs @@ -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("Data"); } } diff --git a/src/persistence/Elsa.Persistence.EntityFramework/Elsa.Persistence.EntityFramework.Core/Stores/EntityFrameworkWorkflowInstanceStore.cs b/src/persistence/Elsa.Persistence.EntityFramework/Elsa.Persistence.EntityFramework.Core/Stores/EntityFrameworkWorkflowInstanceStore.cs index 2d0ac75dc..3fb557ee5 100644 --- a/src/persistence/Elsa.Persistence.EntityFramework/Elsa.Persistence.EntityFramework.Core/Stores/EntityFrameworkWorkflowInstanceStore.cs +++ b/src/persistence/Elsa.Persistence.EntityFramework/Elsa.Persistence.EntityFramework.Core/Stores/EntityFrameworkWorkflowInstanceStore.cs @@ -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(); entity.Fault = data.Fault; + entity.CurrentActivity = data.CurrentActivity; } } } \ No newline at end of file diff --git a/src/persistence/Elsa.Persistence.YesSql/Documents/WorkflowInstanceDocument.cs b/src/persistence/Elsa.Persistence.YesSql/Documents/WorkflowInstanceDocument.cs index c83e44a48..05ebb1d0a 100644 --- a/src/persistence/Elsa.Persistence.YesSql/Documents/WorkflowInstanceDocument.cs +++ b/src/persistence/Elsa.Persistence.YesSql/Documents/WorkflowInstanceDocument.cs @@ -37,5 +37,7 @@ namespace Elsa.Persistence.YesSql.Documents public WorkflowFault? Fault { get; set; } public Stack ScheduledActivities { get; set; } = new(); + public Stack Scopes { get; set; } + public ScheduledActivity? CurrentActivity { get; set; } } } \ No newline at end of file diff --git a/src/samples/dashboard/ElsaDashboard.Samples.Monolith/Startup.cs b/src/samples/dashboard/ElsaDashboard.Samples.Monolith/Startup.cs index 44891d598..fea60d1df 100644 --- a/src/samples/dashboard/ElsaDashboard.Samples.Monolith/Startup.cs +++ b/src/samples/dashboard/ElsaDashboard.Samples.Monolith/Startup.cs @@ -40,8 +40,7 @@ namespace ElsaDashboard.Samples.Monolith services .AddElsaApiEndpoints() - .AddElsaSwagger() - .AddStartupTask(); + .AddElsaSwagger(); // Elsa Dashboard. services.AddRazorPages(); diff --git a/src/samples/server/Elsa.Samples.Server.Host/Startup.cs b/src/samples/server/Elsa.Samples.Server.Host/Startup.cs index ae009e68c..00a00d142 100644 --- a/src/samples/server/Elsa.Samples.Server.Host/Startup.cs +++ b/src/samples/server/Elsa.Samples.Server.Host/Startup.cs @@ -36,8 +36,7 @@ namespace Elsa.Samples.Server.Host services .AddElsaApiEndpoints() - .AddElsaSwagger() - .AddStartupTask(); + .AddElsaSwagger(); } public void Configure(IApplicationBuilder app)