Restored workflow persistence functionality

This commit is contained in:
Sipke Schoorstra 2020-01-18 13:11:11 +01:00
parent d1ce5ce897
commit 6a6edde7de
25 changed files with 156 additions and 116 deletions

View file

@ -40,6 +40,7 @@ namespace Elsa.Extensions
{
workflowInstance.Variables = workflowExecutionContext.Variables;
workflowInstance.ScheduledActivities = new Stack<Models.ScheduledActivity>(workflowExecutionContext.ScheduledActivities.Select(x => new Models.ScheduledActivity(x.Activity.Id, x.Input)));
workflowInstance.Activities = workflowExecutionContext.Activities.Select(x => new ActivityInstance(x.Id, x.Type, x.State, x.Output)).ToList();
workflowInstance.BlockingActivities = new HashSet<BlockingActivity>(workflowExecutionContext.BlockingActivities.Select(x => new BlockingActivity(x.Id, x.Type)), new BlockingActivityEqualityComparer());
workflowInstance.Status = workflowExecutionContext.Status;
workflowInstance.CorrelationId = workflowExecutionContext.CorrelationId;

View file

@ -1,19 +1,11 @@
using Elsa.Services;
using Elsa.Services.Models;
using MediatR;
using Elsa.Services.Models;
namespace Elsa.Messages
namespace Elsa.Messages.Domain
{
public class ActivityExecuted : INotification
public class ActivityExecuted : ActivityNotification
{
public ActivityExecuted(WorkflowExecutionContext workflowExecutionContext, ActivityExecutionContext activityExecutionContext)
public ActivityExecuted(ActivityExecutionContext activityExecutionContext) : base(activityExecutionContext)
{
WorkflowExecutionContext = workflowExecutionContext;
ActivityExecutionContext = activityExecutionContext;
}
public WorkflowExecutionContext WorkflowExecutionContext { get; }
public ActivityExecutionContext ActivityExecutionContext { get; }
public IActivity Activity => ActivityExecutionContext.Activity;
}
}

View file

@ -0,0 +1,11 @@
using Elsa.Services.Models;
namespace Elsa.Messages.Domain
{
public class ActivityExecuting : ActivityNotification
{
public ActivityExecuting(ActivityExecutionContext activityExecutionContext) : base(activityExecutionContext)
{
}
}
}

View file

@ -0,0 +1,18 @@
using Elsa.Services;
using Elsa.Services.Models;
using MediatR;
namespace Elsa.Messages.Domain
{
public abstract class ActivityNotification : INotification
{
protected ActivityNotification(ActivityExecutionContext activityExecutionContext)
{
ActivityExecutionContext = activityExecutionContext;
}
public ActivityExecutionContext ActivityExecutionContext { get; }
public WorkflowExecutionContext WorkflowExecutionContext => ActivityExecutionContext.WorkflowExecutionContext;
public IActivity Activity => ActivityExecutionContext.Activity;
}
}

View file

@ -1,14 +0,0 @@
using Elsa.Services.Models;
namespace Elsa.Messages
{
public class ExecutingActivity : WorkflowNotification
{
public ExecutingActivity(WorkflowExecutionContext workflowExecutionContext, ActivityExecutionContext activityExecutionContext) : base(workflowExecutionContext)
{
ActivityExecutionContext = activityExecutionContext;
}
public ActivityExecutionContext ActivityExecutionContext { get; }
}
}

View file

@ -1,6 +1,6 @@
using Elsa.Services.Models;
namespace Elsa.Messages
namespace Elsa.Messages.Domain
{
public class ExecutingWorkflow : WorkflowNotification
{

View file

@ -1,6 +1,6 @@
using Elsa.Services.Models;
namespace Elsa.Messages
namespace Elsa.Messages.Domain
{
/// <summary>
/// Published when a workflow transitioned into the Cancelled state.

View file

@ -1,6 +1,6 @@
using Elsa.Services.Models;
namespace Elsa.Messages
namespace Elsa.Messages.Domain
{
/// <summary>
/// Published when a workflow transitioned into the Completed state.

View file

@ -1,6 +1,6 @@
using MediatR;
namespace Elsa.Messages
namespace Elsa.Messages.Domain
{
/// <summary>
/// Published when the workflow definition store is updated.

View file

@ -1,6 +1,6 @@
using Elsa.Services.Models;
namespace Elsa.Messages
namespace Elsa.Messages.Domain
{
/// <summary>
/// Published when a burst of execution finished.

View file

@ -1,6 +1,6 @@
using Elsa.Services.Models;
namespace Elsa.Messages
namespace Elsa.Messages.Domain
{
/// <summary>
/// Published when a workflow transitioned into the Faulted state.

View file

@ -1,7 +1,7 @@
using Elsa.Services.Models;
using MediatR;
namespace Elsa.Messages
namespace Elsa.Messages.Domain
{
/// <summary>
/// Common base for workflow-related events.

View file

@ -1,6 +1,6 @@
using Elsa.Services.Models;
namespace Elsa.Messages
namespace Elsa.Messages.Domain
{
/// <summary>
/// Published when a workflow transitioned into the Suspended state.

View file

@ -2,6 +2,18 @@ namespace Elsa.Models
{
public class ActivityInstance
{
public ActivityInstance()
{
}
public ActivityInstance(string id, string type, Variables state, Variable? output)
{
Id = id;
Type = type;
State = state;
Output = output;
}
public string? Id { get; set; }
public string? Type { get; set; }
public Variables? State { get; set; }

View file

@ -9,6 +9,7 @@ namespace Elsa.Models
public WorkflowInstance()
{
Variables = new Variables();
Activities = new List<ActivityInstance>();
BlockingActivities = new HashSet<BlockingActivity>(new BlockingActivityEqualityComparer());
ExecutionLog = new List<ExecutionLogEntry>();
ScheduledActivities = new Stack<ScheduledActivity>();
@ -22,6 +23,7 @@ namespace Elsa.Models
public Instant CreatedAt { get; set; }
public Variables Variables { get; set; }
public Variable? Output { get; set; }
public ICollection<ActivityInstance> Activities { get; set; }
public HashSet<BlockingActivity> BlockingActivities { get; set; }
public ICollection<ExecutionLogEntry> ExecutionLog { get; set; }
public WorkflowFault? Fault { get; set; }

View file

@ -8,7 +8,7 @@ namespace Elsa.Services
{
Task<WorkflowExecutionContext> RunWorkflowAsync(Workflow workflow, string? activityId = default, object? input = default, string? correlationId = default, CancellationToken cancellationToken = default);
Task<WorkflowExecutionContext> RunWorkflowInstanceAsync(string workflowInstanceId, string? activityId = default, object? input = default, CancellationToken cancellationToken = default);
Task<WorkflowExecutionContext?> RunWorkflowInstanceAsync(string workflowInstanceId, string? activityId = default, object? input = default, CancellationToken cancellationToken = default);
/// <summary>
/// Run a registered workflow by its ID.

View file

@ -1,18 +0,0 @@
using Elsa.Models;
namespace Elsa.Services.Models
{
public class WorkflowExecutionScope
{
public WorkflowExecutionScope(Variables? variables = default)
{
Variables = variables ?? new Variables();
}
public Variables Variables { get; }
public void SetVariable(string variableName, object value) => Variables.SetVariable(variableName, value);
public T GetVariable<T>(string name) => Variables.GetVariable<T>(name);
public object GetVariable(string name) => Variables.GetVariable(name);
}
}

View file

@ -6,6 +6,7 @@ using Elsa.Attributes;
using Elsa.Design;
using Elsa.Extensions;
using Elsa.Messages;
using Elsa.Messages.Domain;
using Elsa.Results;
using Elsa.Services;
using Elsa.Services.Models;

View file

@ -1,57 +1,69 @@
using System;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Extensions;
using Elsa.Models;
using Elsa.Persistence;
using Elsa.Services.Models;
using MediatR;
using Microsoft.Extensions.Logging;
namespace Elsa.Messages.Handlers
namespace Elsa.Messages.Domain.Handlers
{
public class PersistenceWorkflowEventHandler :
INotificationHandler<WorkflowExecuted>,
public class PersistenceWorkflowEventHandler :
INotificationHandler<WorkflowExecuted>,
INotificationHandler<WorkflowSuspended>,
INotificationHandler<ActivityExecuted>,
INotificationHandler<WorkflowCompleted>
{
private readonly IWorkflowInstanceStore workflowInstanceStore;
private readonly ILogger logger;
public PersistenceWorkflowEventHandler(IWorkflowInstanceStore workflowInstanceStore)
public PersistenceWorkflowEventHandler(IWorkflowInstanceStore workflowInstanceStore, ILogger<PersistenceWorkflowEventHandler> logger)
{
this.workflowInstanceStore = workflowInstanceStore;
this.logger = logger;
}
public async Task Handle(WorkflowSuspended notification, CancellationToken cancellationToken)
{
if (notification.WorkflowExecutionContext.PersistenceBehavior == WorkflowPersistenceBehavior.Suspended)
await SaveWorkflowAsync(cancellationToken);
await SaveWorkflowAsync(notification.WorkflowExecutionContext, cancellationToken);
}
public async Task Handle(WorkflowExecuted notification, CancellationToken cancellationToken)
{
if (notification.WorkflowExecutionContext.PersistenceBehavior == WorkflowPersistenceBehavior.WorkflowExecuted)
await SaveWorkflowAsync(cancellationToken);
await SaveWorkflowAsync(notification.WorkflowExecutionContext, cancellationToken);
}
public async Task Handle(ActivityExecuted notification, CancellationToken cancellationToken)
{
if (notification.WorkflowExecutionContext.PersistenceBehavior == WorkflowPersistenceBehavior.ActivityExecuted)
await SaveWorkflowAsync(cancellationToken);
}
public async Task Handle(WorkflowCompleted notification, CancellationToken cancellationToken)
{
var workflow = notification.WorkflowExecutionContext;
var blueprint = workflow;
if (blueprint.DeleteCompletedInstances || blueprint.PersistenceBehavior == WorkflowPersistenceBehavior.Suspended)
await workflowInstanceStore.DeleteAsync(workflow.InstanceId, cancellationToken);
await SaveWorkflowAsync(notification.WorkflowExecutionContext, cancellationToken);
}
private async Task SaveWorkflowAsync(CancellationToken cancellationToken)
public async Task Handle(WorkflowCompleted notification, CancellationToken cancellationToken)
{
//var workflowInstance = workflow.ToInstance();
//await workflowInstanceStore.SaveAsync(workflowInstance, cancellationToken);
throw new NotImplementedException();
var workflowExecutionContext = notification.WorkflowExecutionContext;
if (workflowExecutionContext.DeleteCompletedInstances || workflowExecutionContext.PersistenceBehavior == WorkflowPersistenceBehavior.Suspended)
{
logger.LogDebug("Deleting completed workflow instance {WorkflowInstanceId}", workflowExecutionContext.InstanceId);
await workflowInstanceStore.DeleteAsync(workflowExecutionContext.InstanceId, cancellationToken);
}
}
private async Task SaveWorkflowAsync(WorkflowExecutionContext workflowExecutionContext, CancellationToken cancellationToken)
{
var workflowInstance = await workflowInstanceStore.GetByIdAsync(workflowExecutionContext.InstanceId, cancellationToken);
if (workflowInstance == null)
workflowInstance = workflowExecutionContext.CreateWorkflowInstance();
else
workflowInstance = workflowExecutionContext.UpdateWorkflowInstance(workflowInstance);
await workflowInstanceStore.SaveAsync(workflowInstance, cancellationToken);
}
}
}

View file

@ -2,6 +2,7 @@
using System.Threading;
using System.Threading.Tasks;
using Elsa.Messages;
using Elsa.Messages.Domain;
using Elsa.Models;
using MediatR;

View file

@ -1,27 +0,0 @@
using System;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Messages.Distributed;
using Microsoft.Extensions.DependencyInjection;
using Rebus.Bus;
using Rebus.ServiceProvider;
namespace Elsa.Runtime
{
public static class ServiceProviderExtensions
{
/// <summary>
/// Starts the service bus and registers workflow message consumers.
/// </summary>
public static async Task<IServiceProvider> StartElsaAsync(this IServiceProvider serviceProvider, CancellationToken cancellationToken = default)
{
var bus = serviceProvider.GetRequiredService<IBus>();
serviceProvider.UseRebus();
await bus.Subscribe<RunWorkflow>();
return serviceProvider;
}
}
}

View file

@ -7,11 +7,13 @@ using Elsa.Comparers;
using Elsa.Expressions;
using Elsa.Extensions;
using Elsa.Messages;
using Elsa.Messages.Domain;
using Elsa.Models;
using Elsa.Persistence;
using Elsa.Results;
using Elsa.Services.Models;
using MediatR;
using Microsoft.Extensions.Logging;
using ScheduledActivity = Elsa.Services.Models.ScheduledActivity;
namespace Elsa.Services
@ -30,6 +32,7 @@ namespace Elsa.Services
private readonly IIdGenerator idGenerator;
private readonly IMediator mediator;
private readonly IServiceProvider serviceProvider;
private readonly ILogger logger;
public WorkflowHost(
IWorkflowRegistry workflowRegistry,
@ -38,7 +41,8 @@ namespace Elsa.Services
IExpressionEvaluator expressionEvaluator,
IIdGenerator idGenerator,
IMediator mediator,
IServiceProvider serviceProvider)
IServiceProvider serviceProvider,
ILogger<WorkflowHost> logger)
{
this.workflowRegistry = workflowRegistry;
this.workflowInstanceStore = workflowInstanceStore;
@ -47,11 +51,19 @@ namespace Elsa.Services
this.idGenerator = idGenerator;
this.mediator = mediator;
this.serviceProvider = serviceProvider;
this.logger = logger;
}
public async Task<WorkflowExecutionContext> RunWorkflowInstanceAsync(string workflowInstanceId, string? activityId = default, object? input = default, CancellationToken cancellationToken = default)
public async Task<WorkflowExecutionContext?> RunWorkflowInstanceAsync(string workflowInstanceId, string? activityId = default, object? input = default, CancellationToken cancellationToken = default)
{
var workflowInstance = await workflowInstanceStore.GetByIdAsync(workflowInstanceId, cancellationToken);
if (workflowInstance == null)
{
logger.LogDebug("Workflow instance {WorkflowInstanceId} does not exist.", workflowInstanceId);
return null;
}
var workflow = await workflowRegistry.GetWorkflowAsync(workflowInstance.DefinitionId, VersionOptions.SpecificVersion(workflowInstance.Version), cancellationToken);
return await RunAsync(workflow, workflowInstance, activityId, input, cancellationToken);
}
@ -127,13 +139,36 @@ namespace Elsa.Services
var result = await activityOperation(activityExecutionContext, currentActivity, cancellationToken);
await result.ExecuteAsync(activityExecutionContext, cancellationToken);
await mediator.Publish(new ActivityExecuted(workflowExecutionContext, activityExecutionContext), cancellationToken);
await mediator.Publish(new ActivityExecuted(activityExecutionContext), cancellationToken);
activityOperation = Execute;
}
if (workflowExecutionContext.Status == WorkflowStatus.Running)
workflowExecutionContext.Complete();
await mediator.Publish(new WorkflowExecuted(workflowExecutionContext), cancellationToken);
var statusEvent = default(object);
switch (workflowExecutionContext.Status)
{
case WorkflowStatus.Cancelled:
statusEvent = new WorkflowCancelled(workflowExecutionContext);
break;
case WorkflowStatus.Completed:
statusEvent = new WorkflowCompleted(workflowExecutionContext);
break;
case WorkflowStatus.Faulted:
statusEvent = new WorkflowFaulted(workflowExecutionContext);
break;
case WorkflowStatus.Suspended:
statusEvent = new WorkflowSuspended(workflowExecutionContext);
break;
}
if (statusEvent != null)
await mediator.Publish(statusEvent, cancellationToken);
}
private ScheduledActivity CreateScheduledActivity(Elsa.Models.ScheduledActivity scheduledActivityModel, IDictionary<string, IActivity> activityLookup)
@ -145,11 +180,23 @@ namespace Elsa.Services
private WorkflowExecutionContext CreateWorkflowExecutionContext(Workflow workflow, WorkflowInstance workflowInstance)
{
var activityLookup = workflow.Activities.ToDictionary(x => x.Id);
var activityInstanceLookup = workflowInstance.Activities.ToDictionary(x => x.Id);
var scheduledActivities = new Stack<ScheduledActivity>(workflowInstance.ScheduledActivities.Reverse().Select(x => CreateScheduledActivity(x, activityLookup)));
var blockingActivities = new HashSet<IActivity>(workflowInstance.BlockingActivities.Select(x => activityLookup[x.ActivityId]));
var variables = workflowInstance.Variables;
var status = workflowInstance.Status;
var persistenceBehavior = workflow.PersistenceBehavior;
foreach (var activity in workflow.Activities)
{
if (!activityInstanceLookup.ContainsKey(activity.Id))
continue;
var activityInstance = activityInstanceLookup[activity.Id];
activity.State = activityInstance.State;
activity.Output = activityInstance.Output;
}
return CreateWorkflowExecutionContext(
workflowInstance.Id,
workflow.DefinitionId,

View file

@ -1,14 +1,15 @@
using System;
using System.Threading;
using System.Threading.Tasks;
using Elsa.Messages.Distributed;
using Elsa.Runtime;
using Rebus.Bus;
namespace Elsa.StartupTasks
{
public class StartServiceBusTask : IStartupTask
{
private readonly IServiceProvider serviceProvider;
public StartServiceBusTask(IServiceProvider serviceProvider) => this.serviceProvider = serviceProvider;
public Task ExecuteAsync(CancellationToken cancellationToken = default) => serviceProvider.StartElsaAsync();
private readonly IBus serviceBus;
public StartServiceBusTask(IBus serviceBus) => this.serviceBus = serviceBus;
public Task ExecuteAsync(CancellationToken cancellationToken = default) => serviceBus.Subscribe<RunWorkflow>();
}
}

View file

@ -2,6 +2,7 @@
using Elsa.Runtime;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
using NodaTime;
namespace Elsa.Samples.Timers
@ -10,13 +11,12 @@ namespace Elsa.Samples.Timers
{
private static async Task Main()
{
var host = CreateHostBuilder().UseConsoleLifetime().Build();
await host.Services.StartElsaAsync();
await host.RunAsync();
await CreateHostBuilder().UseConsoleLifetime().Build().RunAsync();
}
public static IHostBuilder CreateHostBuilder() =>
Host.CreateDefaultBuilder()
.ConfigureLogging(logging => logging.AddConsole().SetMinimumLevel(LogLevel.Debug))
.ConfigureServices((hostContext, services) =>
{
services

View file

@ -19,11 +19,11 @@ using Xunit;
namespace Elsa.Core.UnitTests
{
public class WorkflowRunnerTests
public class WorkflowHostTests
{
private readonly WorkflowHost workflowHost;
public WorkflowRunnerTests()
public WorkflowHostTests()
{
var fixture = new Fixture().Customize(new NodaTimeCustomization());
var workflowActivatorMock = new Mock<IWorkflowActivator>();
@ -44,7 +44,8 @@ namespace Elsa.Core.UnitTests
workflowExpressionEvaluatorMock.Object,
idGeneratorMock.Object,
mediatorMock.Object,
serviceProvider);
serviceProvider,
logger);
}
[Fact(DisplayName = "Can run simple workflow to completed state.")]