From 1a70cfd18cbb460024dfdce970ab53ed05f1d5d0 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Thu, 25 Jul 2019 17:13:18 +0200 Subject: [PATCH] Implemented workflow correlation (#54) --- Samples.sln | 7 ++ samples/Sample07/DocumentApprovalWorkflow.cs | 4 +- samples/Sample11/CorrelationWorkflow.cs | 20 +++++ samples/Sample11/Program.cs | 64 +++++++++++++++ samples/Sample11/Sample11.csproj | 13 +++ .../Activities/SignalEvent.cs | 30 ------- .../Extensions/ServiceCollectionExtensions.cs | 3 +- .../Handlers/SignalRequestHandler.cs | 5 +- .../Handlers/TriggerRequestHandler.cs | 2 +- .../Consumers/WorkflowConsumer.cs | 62 +++++++------- .../WorkflowInstanceStoreExtensions.cs | 22 +++++ .../Elsa.Abstractions/Models/Variables.cs | 70 +++++++++------- .../Models/WorkflowInstance.cs | 1 + .../Persistence/IWorkflowInstanceStore.cs | 14 +--- .../Services/IWorkflowInvoker.cs | 25 +++--- .../Services/Models/Workflow.cs | 15 ++-- .../Activities/Primitives/Correlate.cs | 36 +++++++++ .../Activities/Primitives/SignalEvent.cs | 45 +++++++++++ .../Extensions/ServiceCollectionExtensions.cs | 11 ++- .../WorkflowInstanceCollectionExtensions.cs | 1 - .../Memory/MemoryWorkflowInstanceStore.cs | 18 ++++- .../CommonScriptEngineConfigurator.cs | 1 + .../Elsa.Core/Services/WorkflowFactory.cs | 7 +- .../Elsa.Core/Services/WorkflowInvoker.cs | 80 ++++++++++++++----- .../Documents/WorkflowInstanceDocument.cs | 1 + .../Extensions/StoreFactory.cs | 6 -- .../Indexes/WorkflowInstanceIndex.cs | 7 +- .../Services/YesSqlWorkflowInstanceStore.cs | 19 ++++- .../StartupTasks/StoreInitializationTask.cs | 2 + 29 files changed, 422 insertions(+), 169 deletions(-) create mode 100644 samples/Sample11/CorrelationWorkflow.cs create mode 100644 samples/Sample11/Program.cs create mode 100644 samples/Sample11/Sample11.csproj delete mode 100644 src/activities/Elsa.Activities.Http/Activities/SignalEvent.cs create mode 100644 src/core/Elsa.Abstractions/Extensions/WorkflowInstanceStoreExtensions.cs create mode 100644 src/core/Elsa.Core/Activities/Primitives/Correlate.cs create mode 100644 src/core/Elsa.Core/Activities/Primitives/SignalEvent.cs diff --git a/Samples.sln b/Samples.sln index 2029a0268..947951563 100644 --- a/Samples.sln +++ b/Samples.sln @@ -59,6 +59,8 @@ Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Sample10", "samples\Sample1 EndProject Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.AutoMapper.Extensions", "src\core\Elsa.AutoMapper.Extensions\Elsa.AutoMapper.Extensions.csproj", "{F7633BE5-0335-4148-8FB7-EAED9A2276A2}" EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Sample11", "samples\Sample11\Sample11.csproj", "{0C3BB425-2FED-4ADF-8DBE-D68DF0CA5A41}" +EndProject Global GlobalSection(SolutionConfigurationPlatforms) = preSolution Debug|Any CPU = Debug|Any CPU @@ -141,6 +143,10 @@ Global {F7633BE5-0335-4148-8FB7-EAED9A2276A2}.Debug|Any CPU.Build.0 = Debug|Any CPU {F7633BE5-0335-4148-8FB7-EAED9A2276A2}.Release|Any CPU.ActiveCfg = Release|Any CPU {F7633BE5-0335-4148-8FB7-EAED9A2276A2}.Release|Any CPU.Build.0 = Release|Any CPU + {0C3BB425-2FED-4ADF-8DBE-D68DF0CA5A41}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {0C3BB425-2FED-4ADF-8DBE-D68DF0CA5A41}.Debug|Any CPU.Build.0 = Debug|Any CPU + {0C3BB425-2FED-4ADF-8DBE-D68DF0CA5A41}.Release|Any CPU.ActiveCfg = Release|Any CPU + {0C3BB425-2FED-4ADF-8DBE-D68DF0CA5A41}.Release|Any CPU.Build.0 = Release|Any CPU EndGlobalSection GlobalSection(SolutionProperties) = preSolution HideSolutionNode = FALSE @@ -168,6 +174,7 @@ Global {905C1067-2119-443D-B5FE-DA8E027FAD62} = {600C4AE0-0585-4181-9CEC-2BB43430C713} {E0076858-353F-4302-9690-920204E38BF6} = {5E5E1E84-DDBC-40D6-B891-0D563A15A44A} {F7633BE5-0335-4148-8FB7-EAED9A2276A2} = {35F44BE9-13D0-417C-A01A-F0787BEE6DC3} + {0C3BB425-2FED-4ADF-8DBE-D68DF0CA5A41} = {5E5E1E84-DDBC-40D6-B891-0D563A15A44A} EndGlobalSection GlobalSection(ExtensibilityGlobals) = postSolution SolutionGuid = {8B0975FD-7050-48B0-88C5-48C33378E158} diff --git a/samples/Sample07/DocumentApprovalWorkflow.cs b/samples/Sample07/DocumentApprovalWorkflow.cs index bb2e326a1..54b5d570d 100644 --- a/samples/Sample07/DocumentApprovalWorkflow.cs +++ b/samples/Sample07/DocumentApprovalWorkflow.cs @@ -59,12 +59,12 @@ namespace Sample07 { fork .When("Approve") - .Then(activity => activity.SignalName = "approve") + .Then(activity => activity.Signal = new PlainTextExpression("approve")) .Then("join-signals"); fork .When("Reject") - .Then(activity => activity.SignalName = "reject") + .Then(activity => activity.Signal = new PlainTextExpression("reject")) .Then("join-signals"); } ) diff --git a/samples/Sample11/CorrelationWorkflow.cs b/samples/Sample11/CorrelationWorkflow.cs new file mode 100644 index 000000000..89c3ca5fa --- /dev/null +++ b/samples/Sample11/CorrelationWorkflow.cs @@ -0,0 +1,20 @@ +using Elsa.Activities.Console.Activities; +using Elsa.Activities.Primitives; +using Elsa.Expressions; +using Elsa.Services; +using Elsa.Services.Models; + +namespace Sample11 +{ + public class CorrelationWorkflow : IWorkflow + { + public void Build(IWorkflowBuilder builder) + { + builder + .StartWith(activity => activity.TextExpression = new JavaScriptExpression("`Workflow started with correlation ID \"${correlationId()}\".`")) + .Then(activity => activity.Signal = new PlainTextExpression("Proceed")) + .Then(activity => activity.TextExpression = new JavaScriptExpression("`Signal received for workflow with correlation ID: \"${correlationId()}\"`")) + .Then(activity => activity.TextExpression = new PlainTextExpression("Workflow finished.")); + } + } +} \ No newline at end of file diff --git a/samples/Sample11/Program.cs b/samples/Sample11/Program.cs new file mode 100644 index 000000000..b7b2996d5 --- /dev/null +++ b/samples/Sample11/Program.cs @@ -0,0 +1,64 @@ +using System; +using System.Linq; +using System.Threading.Tasks; +using Elsa.Activities.Console.Extensions; +using Elsa.Activities.Primitives; +using Elsa.Extensions; +using Elsa.Models; +using Elsa.Persistence.Memory; +using Elsa.Runtime; +using Elsa.Services; +using Microsoft.Extensions.DependencyInjection; + +namespace Sample11 +{ + /// + /// Demonstrates workflow correlation. + /// + public class Program + { + public static async Task Main(string[] args) + { + var services = BuildServices(); + var registry = services.GetService(); + var workflowDefinition = registry.RegisterWorkflow(); + var invoker = services.GetRequiredService(); + + Console.WriteLine("How many workflow instances should be started? Enter a number:"); + var instanceCount = int.Parse(Console.ReadLine()); + + for (var i = 0; i < instanceCount; i++) + { + await invoker.InvokeAsync(workflowDefinition, correlationId: $"document {i + 1}"); + } + + var retry = true; + + while (retry) + { + Console.WriteLine(); + Console.WriteLine("Now, enter the correlation ID of the workflow to resume:"); + var correlationId = Console.ReadLine(); + + // Resume one workflow using the specified correlation ID. + var triggeredExecutionContexts = (await invoker.TriggerAsync(nameof(SignalEvent), new Variables { ["Signal"] = "Proceed" }, correlationId)).ToList(); + + Console.WriteLine("{0} workflow was resumed. Would you like to trigger another?", triggeredExecutionContexts.Count); + retry = string.Equals("y", Console.ReadLine(), StringComparison.OrdinalIgnoreCase); + } + + Console.WriteLine("Bye!"); + } + + private static IServiceProvider BuildServices() + { + return new ServiceCollection() + .AddWorkflows() + .AddStartupRunner() + .AddConsoleActivities() + .AddMemoryWorkflowDefinitionStore() + .AddMemoryWorkflowInstanceStore() + .BuildServiceProvider(); + } + } +} \ No newline at end of file diff --git a/samples/Sample11/Sample11.csproj b/samples/Sample11/Sample11.csproj new file mode 100644 index 000000000..aa2f1d875 --- /dev/null +++ b/samples/Sample11/Sample11.csproj @@ -0,0 +1,13 @@ + + + + Exe + netcoreapp2.2 + + + + + + + + diff --git a/src/activities/Elsa.Activities.Http/Activities/SignalEvent.cs b/src/activities/Elsa.Activities.Http/Activities/SignalEvent.cs deleted file mode 100644 index 268a4b521..000000000 --- a/src/activities/Elsa.Activities.Http/Activities/SignalEvent.cs +++ /dev/null @@ -1,30 +0,0 @@ -using Elsa.Results; -using Elsa.Services; -using Elsa.Services.Models; - -namespace Elsa.Activities.Http.Activities -{ - public class SignalEvent : Activity - { - public string SignalName - { - get => GetState(); - set => SetState(value); - } - - protected override bool OnCanExecute(WorkflowExecutionContext context) - { - return context.Workflow.Input.ContainsKey("signal") && (string) context.Workflow.Input["signal"] == SignalName; - } - - protected override ActivityExecutionResult OnExecute(WorkflowExecutionContext context) - { - return Halt(); - } - - protected override ActivityExecutionResult OnResume(WorkflowExecutionContext context) - { - return Done(); - } - } -} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Http/Extensions/ServiceCollectionExtensions.cs b/src/activities/Elsa.Activities.Http/Extensions/ServiceCollectionExtensions.cs index 1253feeab..7557513c5 100644 --- a/src/activities/Elsa.Activities.Http/Extensions/ServiceCollectionExtensions.cs +++ b/src/activities/Elsa.Activities.Http/Extensions/ServiceCollectionExtensions.cs @@ -22,8 +22,7 @@ namespace Elsa.Activities.Http.Extensions services .AddActivity() .AddActivity() - .AddActivity() - .AddActivity(); + .AddActivity(); services .AddSingleton() diff --git a/src/activities/Elsa.Activities.Http/RequestHandlers/Handlers/SignalRequestHandler.cs b/src/activities/Elsa.Activities.Http/RequestHandlers/Handlers/SignalRequestHandler.cs index 712d0ecb9..8f70946bd 100644 --- a/src/activities/Elsa.Activities.Http/RequestHandlers/Handlers/SignalRequestHandler.cs +++ b/src/activities/Elsa.Activities.Http/RequestHandlers/Handlers/SignalRequestHandler.cs @@ -6,6 +6,7 @@ using CSharpFunctionalExtensions; using Elsa.Activities.Http.Activities; using Elsa.Activities.Http.Models; using Elsa.Activities.Http.Services; +using Elsa.Activities.Primitives; using Elsa.Extensions; using Elsa.Models; using Elsa.Persistence; @@ -94,12 +95,12 @@ namespace Elsa.Activities.Http.RequestHandlers.Handlers var input = new Variables { - ["signal"] = signal.Name + ["Signal"] = signal.Name }; var workflowDefinition = workflowRegistry.GetById(workflowInstance.DefinitionId); var workflow = workflowFactory.CreateWorkflow(workflowDefinition, input, workflowInstance); - var blockingSignalActivities = workflow.BlockingActivities.Where(x => x is SignalEvent).Cast().Where(x => x.SignalName == signal.Name).ToList(); + var blockingSignalActivities = workflow.BlockingActivities.ToList(); await workflowInvoker.ResumeAsync(workflow, blockingSignalActivities, cancellationToken); if (!httpContext.Response.HasStarted) diff --git a/src/activities/Elsa.Activities.Http/RequestHandlers/Handlers/TriggerRequestHandler.cs b/src/activities/Elsa.Activities.Http/RequestHandlers/Handlers/TriggerRequestHandler.cs index 1b582bc12..9e42ec24e 100644 --- a/src/activities/Elsa.Activities.Http/RequestHandlers/Handlers/TriggerRequestHandler.cs +++ b/src/activities/Elsa.Activities.Http/RequestHandlers/Handlers/TriggerRequestHandler.cs @@ -41,7 +41,7 @@ namespace Elsa.Activities.Http.RequestHandlers.Handlers var requestPath = new Uri(httpContext.Request.Path.ToString(), UriKind.Relative); var method = httpContext.Request.Method; var workflowsToStart = Filter(registry.ListByStartActivity(nameof(HttpRequestEvent)), requestPath, method).ToList(); - var workflowsToResume = Filter(await workflowInstanceStore.ListByBlockingActivityAsync(cancellationToken), requestPath, method).ToList(); + var workflowsToResume = Filter(await workflowInstanceStore.ListByBlockingActivityAsync(cancellationToken: cancellationToken), requestPath, method).ToList(); if (!workflowsToStart.Any() && !workflowsToResume.Any()) { diff --git a/src/activities/Elsa.Activities.MassTransit/Consumers/WorkflowConsumer.cs b/src/activities/Elsa.Activities.MassTransit/Consumers/WorkflowConsumer.cs index 799a455db..b7392ece5 100644 --- a/src/activities/Elsa.Activities.MassTransit/Consumers/WorkflowConsumer.cs +++ b/src/activities/Elsa.Activities.MassTransit/Consumers/WorkflowConsumer.cs @@ -1,31 +1,33 @@ -using System.Threading.Tasks; -using Elsa.Activities.MassTransit.Activities; -using Elsa.Models; -using Elsa.Services; -using MassTransit; - -namespace Elsa.Activities.MassTransit.Consumers -{ - public class WorkflowConsumer : IConsumer where T : class - { - private readonly IWorkflowInvoker workflowInvoker; - - public WorkflowConsumer(IWorkflowInvoker workflowInvoker) - { - this.workflowInvoker = workflowInvoker; - } - - public async Task Consume(ConsumeContext context) - { - var message = context.Message; - var activityType = nameof(ReceiveMassTransitMessage); - var input = new Variables { ["message"] = message }; - - await workflowInvoker.TriggerAsync( - activityType, - input, - x => ReceiveMassTransitMessage.GetMessageType(x) == message.GetType(), - context.CancellationToken); - } - } +using System.Threading.Tasks; +using Elsa.Activities.MassTransit.Activities; +using Elsa.Models; +using Elsa.Services; +using MassTransit; + +namespace Elsa.Activities.MassTransit.Consumers +{ + public class WorkflowConsumer : IConsumer where T : class + { + private readonly IWorkflowInvoker workflowInvoker; + + public WorkflowConsumer(IWorkflowInvoker workflowInvoker) + { + this.workflowInvoker = workflowInvoker; + } + + public async Task Consume(ConsumeContext context) + { + var message = context.Message; + var activityType = nameof(ReceiveMassTransitMessage); + var input = new Variables { ["message"] = message }; + var correlationId = context.CorrelationId?.ToString(); + + await workflowInvoker.TriggerAsync( + activityType, + input, + correlationId, + x => ReceiveMassTransitMessage.GetMessageType(x) == message.GetType(), + context.CancellationToken); + } + } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Extensions/WorkflowInstanceStoreExtensions.cs b/src/core/Elsa.Abstractions/Extensions/WorkflowInstanceStoreExtensions.cs new file mode 100644 index 000000000..8bb6de37b --- /dev/null +++ b/src/core/Elsa.Abstractions/Extensions/WorkflowInstanceStoreExtensions.cs @@ -0,0 +1,22 @@ +using System.Collections.Generic; +using System.Linq; +using System.Threading; +using System.Threading.Tasks; +using Elsa.Models; +using Elsa.Persistence; +using Elsa.Services.Models; + +namespace Elsa.Extensions +{ + public static class WorkflowInstanceStoreExtensions + { + public static async Task> ListByBlockingActivityAsync( + this IWorkflowInstanceStore store, + string correlationId = default, + CancellationToken cancellationToken = default) where TActivity : IActivity + { + var items = await store.ListByBlockingActivityAsync(nameof(TActivity), correlationId, cancellationToken); + return items.Select(x => (x.Item1, x.Item2)); + } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Models/Variables.cs b/src/core/Elsa.Abstractions/Models/Variables.cs index 780113b93..39c6c8a9c 100644 --- a/src/core/Elsa.Abstractions/Models/Variables.cs +++ b/src/core/Elsa.Abstractions/Models/Variables.cs @@ -1,31 +1,41 @@ -using System.Collections.Generic; - -namespace Elsa.Models -{ - public class Variables : Dictionary - { - public static readonly Variables Empty = new Variables(); - - public Variables() - { - } - - public Variables(Variables other) - { - foreach (var variable in other) - { - this[variable.Key] = variable.Value; - } - } - - public object GetVariable(string name) - { - return ContainsKey(name) ? this[name] : null; - } - - public T GetVariable(string name) - { - return ContainsKey(name) ? (T)this[name] : default(T); - } - } +using System.Collections.Generic; + +namespace Elsa.Models +{ + public class Variables : Dictionary + { + public static readonly Variables Empty = new Variables(); + + public Variables() + { + } + + public Variables(Variables other) + { + foreach (var variable in other) + { + this[variable.Key] = variable.Value; + } + } + + public object GetVariable(string name) + { + return ContainsKey(name) ? this[name] : null; + } + + public T GetVariable(string name) + { + return ContainsKey(name) ? (T)this[name] : default(T); + } + + public bool HasVariable(string name, object value) + { + return GetVariable(name) == value; + } + + public bool HasVariable(string name) + { + return ContainsKey(name); + } + } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Models/WorkflowInstance.cs b/src/core/Elsa.Abstractions/Models/WorkflowInstance.cs index 8d3dda76c..16993e9bf 100644 --- a/src/core/Elsa.Abstractions/Models/WorkflowInstance.cs +++ b/src/core/Elsa.Abstractions/Models/WorkflowInstance.cs @@ -8,6 +8,7 @@ namespace Elsa.Models public string Id { get; set; } public string DefinitionId { get; set; } public WorkflowStatus Status { get; set; } + public string CorrelationId { get; set; } public Instant CreatedAt { get; set; } public Instant? StartedAt { get; set; } public Instant? HaltedAt { get; set; } diff --git a/src/core/Elsa.Abstractions/Persistence/IWorkflowInstanceStore.cs b/src/core/Elsa.Abstractions/Persistence/IWorkflowInstanceStore.cs index c0db79cff..8963a0aa3 100644 --- a/src/core/Elsa.Abstractions/Persistence/IWorkflowInstanceStore.cs +++ b/src/core/Elsa.Abstractions/Persistence/IWorkflowInstanceStore.cs @@ -1,9 +1,7 @@ using System.Collections.Generic; -using System.Linq; using System.Threading; using System.Threading.Tasks; using Elsa.Models; -using Elsa.Services.Models; using WorkflowInstance = Elsa.Models.WorkflowInstance; namespace Elsa.Persistence @@ -12,19 +10,11 @@ namespace Elsa.Persistence { Task SaveAsync(WorkflowInstance instance, CancellationToken cancellationToken = default); Task GetByIdAsync(string id, CancellationToken cancellationToken = default); + Task GetByCorrelationIdAsync(string correlationId, CancellationToken cancellationToken = default); Task> ListByDefinitionAsync(string definitionId, CancellationToken cancellationToken = default); Task> ListAllAsync(CancellationToken cancellationToken = default); - Task> ListByBlockingActivityAsync(string activityType, 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); } - - public static class WorkflowInstanceStoreExtensions - { - public static async Task> ListByBlockingActivityAsync(this IWorkflowInstanceStore store, CancellationToken cancellationToken) where TActivity : IActivity - { - var items = await store.ListByBlockingActivityAsync(nameof(TActivity), cancellationToken); - return items.Select(x => (x.Item1, x.Item2)); - } - } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/IWorkflowInvoker.cs b/src/core/Elsa.Abstractions/Services/IWorkflowInvoker.cs index 55e6190ce..3494a1bde 100644 --- a/src/core/Elsa.Abstractions/Services/IWorkflowInvoker.cs +++ b/src/core/Elsa.Abstractions/Services/IWorkflowInvoker.cs @@ -15,35 +15,36 @@ namespace Elsa.Services IEnumerable startActivityIds = default, CancellationToken cancellationToken = default ); - + Task InvokeAsync( WorkflowDefinition workflowDefinition, - Variables input = null, - WorkflowInstance workflowInstance = null, + Variables input = default, + WorkflowInstance workflowInstance = default, IEnumerable startActivityIds = default, + string correlationId = default, CancellationToken cancellationToken = default ); - + Task InvokeAsync( - WorkflowInstance workflowInstance = null, - Variables input = null, + WorkflowInstance workflowInstance = default, + Variables input = default, IEnumerable startActivityIds = default, CancellationToken cancellationToken = default - ) where T:IWorkflow, new(); + ) where T : IWorkflow, new(); Task InvokeAsync( WorkflowInstance workflowInstance, - Variables input = null, + Variables input = default, IEnumerable startActivityIds = default, CancellationToken cancellationToken = default ); - + /// /// Starts new workflows that start with the specified activity name and resumes halted workflows that are blocked on activities with the specified activity name. /// - Task TriggerAsync( - string activityType, - Variables input, + Task> TriggerAsync(string activityType, + Variables input = default, + string correlationId = default, Func activityStatePredicate = default, CancellationToken cancellationToken = default); } diff --git a/src/core/Elsa.Abstractions/Services/Models/Workflow.cs b/src/core/Elsa.Abstractions/Services/Models/Workflow.cs index 6fe9eb648..8fa605896 100644 --- a/src/core/Elsa.Abstractions/Services/Models/Workflow.cs +++ b/src/core/Elsa.Abstractions/Services/Models/Workflow.cs @@ -1,4 +1,5 @@ -using System.Collections.Generic; +using System; +using System.Collections.Generic; using System.Linq; using Elsa.Comparers; using Elsa.Models; @@ -14,15 +15,14 @@ namespace Elsa.Services.Models string definitionId, IEnumerable activities, IEnumerable connections, - Variables input = null, - WorkflowInstance workflowInstance = null) : this() + Variables input = default, + string correlationId = default) : this() { Id = id; DefinitionId = definitionId; Activities = activities.ToList(); Connections = connections.ToList(); Input = new Variables(input ?? Variables.Empty); - Initialize(workflowInstance); } public Workflow() @@ -34,6 +34,7 @@ namespace Elsa.Services.Models public string Id { get; set; } public string DefinitionId { get; } + public string CorrelationId { get; set; } public WorkflowStatus Status { get; set; } public Instant CreatedAt { get; set; } public Instant? StartedAt { get; set; } @@ -55,6 +56,7 @@ namespace Elsa.Services.Models { Id = Id, DefinitionId = DefinitionId, + CorrelationId = CorrelationId, Status = Status, CreatedAt = CreatedAt, StartedAt = StartedAt, @@ -68,14 +70,15 @@ namespace Elsa.Services.Models }; } - private void Initialize(WorkflowInstance instance) + public void Initialize(WorkflowInstance instance) { if(instance == null) - return; + throw new ArgumentNullException(nameof(instance)); var activityLookup = Activities.ToDictionary(x => x.Id); Id = instance.Id; + CorrelationId = instance.CorrelationId; Status = instance.Status; CreatedAt = instance.CreatedAt; StartedAt = instance.StartedAt; diff --git a/src/core/Elsa.Core/Activities/Primitives/Correlate.cs b/src/core/Elsa.Core/Activities/Primitives/Correlate.cs new file mode 100644 index 000000000..f65530e5e --- /dev/null +++ b/src/core/Elsa.Core/Activities/Primitives/Correlate.cs @@ -0,0 +1,36 @@ +using System.Threading; +using System.Threading.Tasks; +using Elsa.Expressions; +using Elsa.Extensions; +using Elsa.Results; +using Elsa.Services; +using Elsa.Services.Models; + +namespace Elsa.Activities.Primitives +{ + /// + /// Sets the CorrelationId of the workflow to a given value. + /// + public class Correlate : Activity + { + private readonly IWorkflowExpressionEvaluator expressionEvaluator; + + public Correlate(IWorkflowExpressionEvaluator expressionEvaluator) + { + this.expressionEvaluator = expressionEvaluator; + } + + public WorkflowExpression Expression + { + get => GetState>(); + set => SetState(value); + } + + protected override async Task OnExecuteAsync(WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) + { + var value = await expressionEvaluator.EvaluateAsync(Expression, workflowContext, cancellationToken); + workflowContext.Workflow.CorrelationId = value?.ToString(); + return Done(); + } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Activities/Primitives/SignalEvent.cs b/src/core/Elsa.Core/Activities/Primitives/SignalEvent.cs new file mode 100644 index 000000000..9603ed9cb --- /dev/null +++ b/src/core/Elsa.Core/Activities/Primitives/SignalEvent.cs @@ -0,0 +1,45 @@ +using System.Threading; +using System.Threading.Tasks; +using Elsa.Expressions; +using Elsa.Extensions; +using Elsa.Results; +using Elsa.Services; +using Elsa.Services.Models; + +namespace Elsa.Activities.Primitives +{ + /// + /// Halts workflow execution until the specified signal is received. + /// + public class SignalEvent : Activity + { + private readonly IWorkflowExpressionEvaluator expressionEvaluator; + + public SignalEvent(IWorkflowExpressionEvaluator expressionEvaluator) + { + this.expressionEvaluator = expressionEvaluator; + } + + public WorkflowExpression Signal + { + get => GetState>(); + set => SetState(value); + } + + protected override async Task OnCanExecuteAsync(WorkflowExecutionContext context, CancellationToken cancellationToken) + { + var signal = await expressionEvaluator.EvaluateAsync(Signal, context, cancellationToken); + return context.Workflow.Input.HasVariable("Signal", signal); + } + + protected override ActivityExecutionResult OnExecute(WorkflowExecutionContext context) + { + return Halt(true); + } + + protected override ActivityExecutionResult OnResume(WorkflowExecutionContext context) + { + return Done(); + } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Extensions/ServiceCollectionExtensions.cs b/src/core/Elsa.Core/Extensions/ServiceCollectionExtensions.cs index d590346dd..b784c9faa 100644 --- a/src/core/Elsa.Core/Extensions/ServiceCollectionExtensions.cs +++ b/src/core/Elsa.Core/Extensions/ServiceCollectionExtensions.cs @@ -2,7 +2,6 @@ using System; using Elsa.Activities.ControlFlow; using Elsa.Activities.Primitives; using Elsa.Expressions; -using Elsa.Persistence; using Elsa.Scripting; using Elsa.Serialization; using Elsa.Serialization.Formatters; @@ -49,7 +48,8 @@ namespace Elsa.Extensions .AddTransient() .AddSingleton>(sp => sp.GetRequiredService) .AddSingleton() - .AddPrimitiveActivities(); + .AddPrimitiveActivities() + .AddControlFlowActivities(); } public static IServiceCollection AddActivity(this IServiceCollection services) @@ -64,6 +64,13 @@ namespace Elsa.Extensions { return services .AddActivity() + .AddActivity() + .AddActivity(); + } + + private static IServiceCollection AddControlFlowActivities(this IServiceCollection services) + { + return services .AddActivity() .AddActivity() .AddActivity() diff --git a/src/core/Elsa.Core/Extensions/WorkflowInstanceCollectionExtensions.cs b/src/core/Elsa.Core/Extensions/WorkflowInstanceCollectionExtensions.cs index d17d827e3..876da9b61 100644 --- a/src/core/Elsa.Core/Extensions/WorkflowInstanceCollectionExtensions.cs +++ b/src/core/Elsa.Core/Extensions/WorkflowInstanceCollectionExtensions.cs @@ -1,7 +1,6 @@ using System.Collections.Generic; using System.Linq; using Elsa.Models; -using Elsa.Services.Models; namespace Elsa.Extensions { diff --git a/src/core/Elsa.Core/Persistence/Memory/MemoryWorkflowInstanceStore.cs b/src/core/Elsa.Core/Persistence/Memory/MemoryWorkflowInstanceStore.cs index da581b8ed..c96ac6af4 100644 --- a/src/core/Elsa.Core/Persistence/Memory/MemoryWorkflowInstanceStore.cs +++ b/src/core/Elsa.Core/Persistence/Memory/MemoryWorkflowInstanceStore.cs @@ -1,4 +1,3 @@ -using System; using System.Collections.Generic; using System.Linq; using System.Threading; @@ -24,6 +23,12 @@ namespace Elsa.Persistence.Memory return Task.FromResult(instance); } + public Task GetByCorrelationIdAsync(string correlationId, CancellationToken cancellationToken = default) + { + var instance = workflowInstances.Values.FirstOrDefault(x => x.CorrelationId == correlationId); + return Task.FromResult(instance); + } + public Task> ListByDefinitionAsync(string definitionId, CancellationToken cancellationToken) { var workflows = workflowInstances.Values.Where(x => x.DefinitionId == definitionId); @@ -36,21 +41,26 @@ namespace Elsa.Persistence.Memory return Task.FromResult(workflows); } - public Task> ListByBlockingActivityAsync(string activityType, CancellationToken cancellationToken) + public Task> ListByBlockingActivityAsync(string activityType, string correlationId = default, CancellationToken cancellationToken = default) { var query = workflowInstances.Values.GetBlockingActivities().Where(x => x.Item2.TypeName == activityType); + if (!string.IsNullOrWhiteSpace(correlationId)) + query = query.Where(x => x.Item1.CorrelationId == correlationId); + return Task.FromResult(query); } public Task> ListByStatusAsync(string definitionId, WorkflowStatus status, CancellationToken cancellationToken) { - throw new NotImplementedException(); + var query = workflowInstances.Values.Where(x => x.DefinitionId == definitionId && x.Status == status); + return Task.FromResult(query); } public Task> ListByStatusAsync(WorkflowStatus status, CancellationToken cancellationToken) { - throw new NotImplementedException(); + var query = workflowInstances.Values.Where(x => x.Status == status); + return Task.FromResult(query); } } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Scripting/CommonScriptEngineConfigurator.cs b/src/core/Elsa.Core/Scripting/CommonScriptEngineConfigurator.cs index f787d6d34..ce95fca74 100644 --- a/src/core/Elsa.Core/Scripting/CommonScriptEngineConfigurator.cs +++ b/src/core/Elsa.Core/Scripting/CommonScriptEngineConfigurator.cs @@ -14,6 +14,7 @@ namespace Elsa.Scripting engine.SetValue("input", (Func) (name => workflow.Input.GetVariable(name))); engine.SetValue("variable", (Func) (name => context.CurrentScope.GetVariable(name))); engine.SetValue("lastResult", (Func) (name => context.CurrentScope.LastResult)); + engine.SetValue("correlationId", (Func) (name => context.Workflow.CorrelationId)); foreach (var variable in workflowExecutionContext.CurrentScope.Variables) { diff --git a/src/core/Elsa.Core/Services/WorkflowFactory.cs b/src/core/Elsa.Core/Services/WorkflowFactory.cs index ac471a6a6..2e53b9fcc 100644 --- a/src/core/Elsa.Core/Services/WorkflowFactory.cs +++ b/src/core/Elsa.Core/Services/WorkflowFactory.cs @@ -32,7 +32,12 @@ namespace Elsa.Services var activities = CreateActivities(definition.Activities).ToList(); var connections = CreateConnections(definition.Connections, activities); var id = idGenerator.Generate(); - return new Workflow(id, definition.Id, activities, connections, input, workflowInstance); + var workflow = new Workflow(id, definition.Id, activities, connections, input); + + if(workflowInstance != null) + workflow.Initialize(workflowInstance); + + return workflow; } private IEnumerable CreateConnections(IEnumerable connectionBlueprints, IEnumerable activities) diff --git a/src/core/Elsa.Core/Services/WorkflowInvoker.cs b/src/core/Elsa.Core/Services/WorkflowInvoker.cs index e0baaf91d..b1a257ada 100644 --- a/src/core/Elsa.Core/Services/WorkflowInvoker.cs +++ b/src/core/Elsa.Core/Services/WorkflowInvoker.cs @@ -60,20 +60,23 @@ namespace Elsa.Services public Task InvokeAsync( WorkflowDefinition workflowDefinition, - Variables input = null, - WorkflowInstance workflowInstance = null, - IEnumerable startActivityIds = default, + Variables input = default, + WorkflowInstance workflowInstance = default, + IEnumerable startActivityIds = default, + string correlationId = default, CancellationToken cancellationToken = default) { var workflow = workflowFactory.CreateWorkflow(workflowDefinition, input, workflowInstance); var startActivities = workflow.Activities.Find(startActivityIds); + + workflow.CorrelationId = correlationId; return InvokeAsync(workflow, startActivities, cancellationToken); } public Task InvokeAsync( WorkflowInstance workflowInstance = null, Variables input = null, - IEnumerable startActivityIds = default, + IEnumerable startActivityIds = default, CancellationToken cancellationToken = default) where T : IWorkflow, new() { var workflow = workflowFactory.CreateWorkflow(input, workflowInstance); @@ -81,19 +84,23 @@ namespace Elsa.Services return InvokeAsync(workflow, startActivities, cancellationToken); } - public Task InvokeAsync(WorkflowInstance workflowInstance, Variables input = null, IEnumerable startActivityIds = default, CancellationToken cancellationToken = default) + public Task InvokeAsync( + WorkflowInstance workflowInstance, + Variables input = null, + IEnumerable startActivityIds = default, + CancellationToken cancellationToken = default) { var definition = workflowRegistry.GetById(workflowInstance.DefinitionId); - return InvokeAsync(definition, input, workflowInstance, startActivityIds, cancellationToken); + return InvokeAsync(definition, input, workflowInstance, startActivityIds, workflowInstance.CorrelationId, cancellationToken); } - - public async Task TriggerAsync( - string activityType, - Variables input, + + public async Task> TriggerAsync(string activityType, + Variables input = default, + string correlationId = default, Func activityStatePredicate = default, CancellationToken cancellationToken = default) { - var workflowInstances = await workflowInstanceStore.ListByBlockingActivityAsync(activityType, cancellationToken).ToListAsync(); + var workflowInstances = await workflowInstanceStore.ListByBlockingActivityAsync(activityType, correlationId, cancellationToken).ToListAsync(); var workflowDefinitions = workflowRegistry.ListByStartActivity(activityType).ToList(); if (activityStatePredicate != null) @@ -102,29 +109,55 @@ namespace Elsa.Services if (activityStatePredicate != null) workflowInstances = workflowInstances.Where(x => activityStatePredicate(x.Item2.State)).ToList(); - await StartWorkflowsAsync(workflowDefinitions, input, cancellationToken); - await ResumeWorkflowsAsync(workflowInstances, input, cancellationToken); + var startedExecutionContexts = await StartWorkflowsAsync(workflowDefinitions, input, correlationId, cancellationToken); + var resumedExecutionContexts = await ResumeWorkflowsAsync(workflowInstances, input, correlationId, cancellationToken); + + return startedExecutionContexts.Concat(resumedExecutionContexts); } - - private async Task StartWorkflowsAsync(IEnumerable<(WorkflowDefinition, ActivityDefinition)> workflowDefinitions, Variables variables, CancellationToken cancellationToken1) + + private async Task> StartWorkflowsAsync( + IEnumerable<(WorkflowDefinition, ActivityDefinition)> workflowDefinitions, + Variables variables, + string correlationId, + CancellationToken cancellationToken1) { + var executionContexts = new List(); + foreach (var (workflowDefinition, activityDefinition) in workflowDefinitions) { var startActivities = workflowDefinition.Activities.Where(x => x.Id == activityDefinition.Id).Select(x => x.Id); - await InvokeAsync(workflowDefinition, variables, startActivityIds: startActivities, cancellationToken: cancellationToken1); + var executionContext = await InvokeAsync( + workflowDefinition, + variables, + startActivityIds: startActivities, + correlationId: correlationId, + cancellationToken: cancellationToken1 + ); + executionContexts.Add(executionContext); } + + return executionContexts; } - - private async Task ResumeWorkflowsAsync(IEnumerable<(WorkflowInstance, ActivityInstance)> workflowInstances, Variables input, CancellationToken cancellationToken) + + private async Task> ResumeWorkflowsAsync( + IEnumerable<(WorkflowInstance, ActivityInstance)> workflowInstances, + Variables input, + string correlationId, + CancellationToken cancellationToken) { + var executionContexts = new List(); + foreach (var (workflowInstance, startActivityInstance) in workflowInstances) { var workflowDefinition = workflowRegistry.GetById(workflowInstance.DefinitionId); workflowInstance.Status = WorkflowStatus.Resuming; - await InvokeAsync(workflowDefinition, input, workflowInstance, new[] { startActivityInstance.Id }, cancellationToken); + var executionContext = await InvokeAsync(workflowDefinition, input, workflowInstance, new[] { startActivityInstance.Id }, correlationId, cancellationToken); + executionContexts.Add(executionContext); } + + return executionContexts; } private async Task ExecuteWorkflowAsync(WorkflowExecutionContext workflowExecutionContext, CancellationToken cancellationToken) @@ -210,7 +243,7 @@ namespace Elsa.Services return null; } - + private async Task ExecuteActivityHaltedAsync(WorkflowExecutionContext workflowContext, IActivity activity, CancellationToken cancellationToken) { return await ExecuteActivityAsync(workflowContext, activity, () => activityInvoker.HaltedAsync(workflowContext, activity, cancellationToken), cancellationToken); @@ -222,11 +255,14 @@ namespace Elsa.Services workflowContext.Fault(activity, ex); } - private async Task CreateWorkflowExecutionContextAsync(Workflow workflow, IEnumerable startActivities, CancellationToken cancellationToken) + private async Task CreateWorkflowExecutionContextAsync( + Workflow workflow, + IEnumerable startActivities, + CancellationToken cancellationToken) { var workflowExecutionContext = new WorkflowExecutionContext(workflow, clock, serviceProvider); var startActivityList = startActivities?.ToList() ?? workflow.GetStartActivities().Take(1).ToList(); - + await workflowExecutionContext.ScheduleActivitiesAsync(startActivityList); if (workflowExecutionContext.HasScheduledActivities) diff --git a/src/persistence/Elsa.Persistence.YesSql/Documents/WorkflowInstanceDocument.cs b/src/persistence/Elsa.Persistence.YesSql/Documents/WorkflowInstanceDocument.cs index df286783b..650562a36 100644 --- a/src/persistence/Elsa.Persistence.YesSql/Documents/WorkflowInstanceDocument.cs +++ b/src/persistence/Elsa.Persistence.YesSql/Documents/WorkflowInstanceDocument.cs @@ -10,6 +10,7 @@ namespace Elsa.Persistence.YesSql.Documents public string WorkflowInstanceId { get; set; } public string DefinitionId { get; set; } public WorkflowStatus Status { get; set; } + public string CorrelationId { get; set; } public Instant CreatedAt { get; set; } public Instant? StartedAt { get; set; } public Instant? HaltedAt { get; set; } diff --git a/src/persistence/Elsa.Persistence.YesSql/Extensions/StoreFactory.cs b/src/persistence/Elsa.Persistence.YesSql/Extensions/StoreFactory.cs index 3db695990..39355a638 100644 --- a/src/persistence/Elsa.Persistence.YesSql/Extensions/StoreFactory.cs +++ b/src/persistence/Elsa.Persistence.YesSql/Extensions/StoreFactory.cs @@ -1,13 +1,7 @@ using System; using System.Data; -using System.Linq; -using System.Reflection; using Elsa.Persistence.YesSql.Options; -using Elsa.Persistence.YesSql.Services; -using Elsa.Persistence.YesSql.StartupTasks; -using Elsa.Runtime; using Microsoft.Extensions.DependencyInjection; -using Microsoft.Extensions.DependencyInjection.Extensions; using Microsoft.Extensions.Options; using YesSql; using YesSql.Indexes; diff --git a/src/persistence/Elsa.Persistence.YesSql/Indexes/WorkflowInstanceIndex.cs b/src/persistence/Elsa.Persistence.YesSql/Indexes/WorkflowInstanceIndex.cs index 14e50f764..33ac86134 100644 --- a/src/persistence/Elsa.Persistence.YesSql/Indexes/WorkflowInstanceIndex.cs +++ b/src/persistence/Elsa.Persistence.YesSql/Indexes/WorkflowInstanceIndex.cs @@ -1,7 +1,5 @@ using System.Linq; using Elsa.Models; -using Elsa.Services.Extensions; -using YesSql; using YesSql.Indexes; namespace Elsa.Persistence.YesSql.Indexes @@ -10,6 +8,7 @@ namespace Elsa.Persistence.YesSql.Indexes { public string WorkflowInstanceId { get; set; } public string WorkflowDefinitionId { get; set; } + public string CorrelationId { get; set; } public WorkflowStatus WorkflowStatus { get; set; } } @@ -28,7 +27,8 @@ namespace Elsa.Persistence.YesSql.Indexes workflowInstance => new WorkflowInstanceIndex { WorkflowDefinitionId = workflowInstance.Id, - WorkflowStatus = workflowInstance.Status + WorkflowStatus = workflowInstance.Status, + CorrelationId = workflowInstance.CorrelationId } ); @@ -40,6 +40,7 @@ namespace Elsa.Persistence.YesSql.Indexes { WorkflowInstanceId = workflowInstance.Id, WorkflowDefinitionId = workflowInstance.Id, + CorrelationId = workflowInstance.CorrelationId, ActivityId = activity.ActivityId, ActivityType = activity.ActivityType } diff --git a/src/persistence/Elsa.Persistence.YesSql/Services/YesSqlWorkflowInstanceStore.cs b/src/persistence/Elsa.Persistence.YesSql/Services/YesSqlWorkflowInstanceStore.cs index a80e3ab5d..6d8290f38 100644 --- a/src/persistence/Elsa.Persistence.YesSql/Services/YesSqlWorkflowInstanceStore.cs +++ b/src/persistence/Elsa.Persistence.YesSql/Services/YesSqlWorkflowInstanceStore.cs @@ -1,4 +1,3 @@ -using System.Collections; using System.Collections.Generic; using System.Threading; using System.Threading.Tasks; @@ -41,6 +40,15 @@ namespace Elsa.Persistence.YesSql.Services } } + public async Task GetByCorrelationIdAsync(string correlationId, CancellationToken cancellationToken = default) + { + using (var session = sessionProvider.GetSession()) + { + var document = await session.Query(x => x.CorrelationId == correlationId).FirstOrDefaultAsync(); + return mapper.Map(document); + } + } + public async Task> ListByDefinitionAsync(string definitionId, CancellationToken cancellationToken) { using (var session = sessionProvider.GetSession()) @@ -59,11 +67,16 @@ namespace Elsa.Persistence.YesSql.Services } } - public async Task> ListByBlockingActivityAsync(string activityType, CancellationToken cancellationToken) + public async Task> ListByBlockingActivityAsync(string activityType, string correlationId = default, CancellationToken cancellationToken = default) { using (var session = sessionProvider.GetSession()) { - var documents = await session.Query(x => x.ActivityType == activityType).ListAsync(); + var query = session.Query(x => x.ActivityType == activityType); + + if (string.IsNullOrWhiteSpace(correlationId)) + query = query.Where(x => x.CorrelationId == correlationId); + + var documents = await query.ListAsync(); var instances = mapper.Map>(documents); return instances.GetBlockingActivities(); diff --git a/src/persistence/Elsa.Persistence.YesSql/StartupTasks/StoreInitializationTask.cs b/src/persistence/Elsa.Persistence.YesSql/StartupTasks/StoreInitializationTask.cs index 996399d91..055e43540 100644 --- a/src/persistence/Elsa.Persistence.YesSql/StartupTasks/StoreInitializationTask.cs +++ b/src/persistence/Elsa.Persistence.YesSql/StartupTasks/StoreInitializationTask.cs @@ -42,11 +42,13 @@ namespace Elsa.Persistence.YesSql.StartupTasks .CreateMapIndexTable(nameof(WorkflowInstanceIndex), table => table .Column("WorkflowInstanceId") .Column("WorkflowDefinitionId") + .Column("CorrelationId") .Column("WorkflowStatus") ) .CreateMapIndexTable(nameof(WorkflowInstanceBlockingActivitiesIndex), table => table .Column("WorkflowInstanceId") .Column("WorkflowDefinitionId") + .Column("CorrelationId") .Column("WorkflowStatus") .Column("ActivityId") .Column("ActivityType")