diff --git a/src/activities/Elsa.Activities.Console/Activities/ReadLine.cs b/src/activities/Elsa.Activities.Console/Activities/ReadLine.cs index 4dedac3be..df460f820 100644 --- a/src/activities/Elsa.Activities.Console/Activities/ReadLine.cs +++ b/src/activities/Elsa.Activities.Console/Activities/ReadLine.cs @@ -33,7 +33,7 @@ namespace Elsa.Activities.Console.Activities protected override async Task OnExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken) { if (input == null) - return Halt(); + return Suspend(); var receivedInput = await input.ReadLineAsync(); return Execute(receivedInput); diff --git a/src/activities/Elsa.Activities.Http/Activities/ReceiveHttpRequest.cs b/src/activities/Elsa.Activities.Http/Activities/ReceiveHttpRequest.cs index e93f9e608..013d9efff 100644 --- a/src/activities/Elsa.Activities.Http/Activities/ReceiveHttpRequest.cs +++ b/src/activities/Elsa.Activities.Http/Activities/ReceiveHttpRequest.cs @@ -85,7 +85,7 @@ namespace Elsa.Activities.Http.Activities protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context) { - return Halt(true); + return Suspend(true); } protected override async Task OnResumeAsync(ActivityExecutionContext context, CancellationToken cancellationToken) diff --git a/src/activities/Elsa.Activities.MassTransit/Activities/ReceiveMassTransitMessage.cs b/src/activities/Elsa.Activities.MassTransit/Activities/ReceiveMassTransitMessage.cs index fdf1981a6..02e9e8bbe 100644 --- a/src/activities/Elsa.Activities.MassTransit/Activities/ReceiveMassTransitMessage.cs +++ b/src/activities/Elsa.Activities.MassTransit/Activities/ReceiveMassTransitMessage.cs @@ -30,7 +30,7 @@ namespace Elsa.Activities.MassTransit.Activities } protected override bool OnCanExecute(ActivityExecutionContext context) => context.Input.Value?.GetType() == MessageType; - protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context) => Halt(true); + protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context) => Suspend(true); protected override Task OnResumeAsync(ActivityExecutionContext context, CancellationToken cancellationToken) { diff --git a/src/activities/Elsa.Activities.Timers/Activities/CronEvent.cs b/src/activities/Elsa.Activities.Timers/Activities/CronEvent.cs index f126ab318..d6f140ef5 100644 --- a/src/activities/Elsa.Activities.Timers/Activities/CronEvent.cs +++ b/src/activities/Elsa.Activities.Timers/Activities/CronEvent.cs @@ -44,7 +44,7 @@ namespace Elsa.Activities.Timers.Activities protected override IActivityExecutionResult OnExecute(ActivityExecutionContext workflowContext) { - return Halt(); + return Suspend(); } protected override async Task OnResumeAsync(ActivityExecutionContext context, CancellationToken cancellationToken) @@ -55,7 +55,7 @@ namespace Elsa.Activities.Timers.Activities return Done(); } - return Halt(); + return Suspend(); } private async Task IsExpiredAsync(ActivityExecutionContext context, CancellationToken cancellationToken) diff --git a/src/activities/Elsa.Activities.Timers/Activities/InstantEvent.cs b/src/activities/Elsa.Activities.Timers/Activities/InstantEvent.cs index de8ec709f..8079d74a7 100644 --- a/src/activities/Elsa.Activities.Timers/Activities/InstantEvent.cs +++ b/src/activities/Elsa.Activities.Timers/Activities/InstantEvent.cs @@ -38,14 +38,14 @@ namespace Elsa.Activities.Timers.Activities { var isExpired = await IsExpiredAsync(context, cancellationToken); - return isExpired ? (IActivityExecutionResult)Done() : Halt(); + return isExpired ? (IActivityExecutionResult)Done() : Suspend(); } protected override async Task OnResumeAsync(ActivityExecutionContext context, CancellationToken cancellationToken) { var isExpired = await IsExpiredAsync(context, cancellationToken); - return isExpired ? (IActivityExecutionResult)Done() : Halt(); + return isExpired ? (IActivityExecutionResult)Done() : Suspend(); } private async Task IsExpiredAsync(ActivityExecutionContext workflowContext, CancellationToken cancellationToken) diff --git a/src/activities/Elsa.Activities.Timers/Activities/TimerEvent.cs b/src/activities/Elsa.Activities.Timers/Activities/TimerEvent.cs index 276fd20ff..d7f080df5 100644 --- a/src/activities/Elsa.Activities.Timers/Activities/TimerEvent.cs +++ b/src/activities/Elsa.Activities.Timers/Activities/TimerEvent.cs @@ -42,7 +42,7 @@ namespace Elsa.Activities.Timers.Activities protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context) { - return Halt(); + return Suspend(); } protected override async Task OnResumeAsync(ActivityExecutionContext context, CancellationToken cancellationToken) @@ -53,7 +53,7 @@ namespace Elsa.Activities.Timers.Activities return Done(); } - return Halt(); + return Suspend(); } private async Task IsExpiredAsync(ActivityExecutionContext context, CancellationToken cancellationToken) diff --git a/src/activities/Elsa.Activities.Timers/HostedServices/TimersHostedService.cs b/src/activities/Elsa.Activities.Timers/HostedServices/TimersHostedService.cs index af07f0cae..421d59754 100644 --- a/src/activities/Elsa.Activities.Timers/HostedServices/TimersHostedService.cs +++ b/src/activities/Elsa.Activities.Timers/HostedServices/TimersHostedService.cs @@ -35,13 +35,11 @@ namespace Elsa.Activities.Timers.HostedServices { try { - using (var scope = serviceProvider.CreateScope()) - { - var workflowInvoker = scope.ServiceProvider.GetRequiredService(); - await workflowInvoker.TriggerAsync(nameof(TimerEvent), Variables.Empty, stoppingToken); - await workflowInvoker.TriggerAsync(nameof(CronEvent), Variables.Empty, stoppingToken); - await workflowInvoker.TriggerAsync(nameof(InstantEvent), Variables.Empty, stoppingToken); - } + using var scope = serviceProvider.CreateScope(); + var workflowInvoker = scope.ServiceProvider.GetRequiredService(); + await workflowInvoker.TriggerAsync(nameof(TimerEvent), cancellationToken: stoppingToken); + await workflowInvoker.TriggerAsync(nameof(CronEvent), cancellationToken:stoppingToken); + await workflowInvoker.TriggerAsync(nameof(InstantEvent), cancellationToken: stoppingToken); } catch (Exception ex) { diff --git a/src/activities/Elsa.Activities.UserTask/Activities/UserTask.cs b/src/activities/Elsa.Activities.UserTask/Activities/UserTask.cs index 33e399f1c..e8091fd05 100644 --- a/src/activities/Elsa.Activities.UserTask/Activities/UserTask.cs +++ b/src/activities/Elsa.Activities.UserTask/Activities/UserTask.cs @@ -35,7 +35,7 @@ namespace Elsa.Activities.UserTask.Activities protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context) { - return Halt(true); + return Suspend(true); } protected override IActivityExecutionResult OnResume(ActivityExecutionContext context) diff --git a/src/activities/Elsa.Activities.Workflows/Activities/Correlate.cs b/src/activities/Elsa.Activities.Workflows/Activities/Correlate.cs index 9434165d7..884cd5e52 100644 --- a/src/activities/Elsa.Activities.Workflows/Activities/Correlate.cs +++ b/src/activities/Elsa.Activities.Workflows/Activities/Correlate.cs @@ -18,7 +18,7 @@ namespace Elsa.Activities.Workflows.Activities public class Correlate : Activity { [ActivityProperty(Hint = "An expression that evaluates to the value to store as the correlation ID.")] - public IWorkflowExpression ValueScriptExpression + public IWorkflowExpression Value { get => GetState>(); set => SetState(value); @@ -26,7 +26,7 @@ namespace Elsa.Activities.Workflows.Activities protected override async Task OnExecuteAsync(ActivityExecutionContext context, CancellationToken cancellationToken) { - var value = await context.EvaluateAsync(ValueScriptExpression, cancellationToken); + var value = await context.EvaluateAsync(Value, cancellationToken); context.WorkflowExecutionContext.Workflow.CorrelationId = value; return Done(); } diff --git a/src/activities/Elsa.Activities.Workflows/Activities/Signaled.cs b/src/activities/Elsa.Activities.Workflows/Activities/Signaled.cs index 42bc2076a..42218dc4c 100644 --- a/src/activities/Elsa.Activities.Workflows/Activities/Signaled.cs +++ b/src/activities/Elsa.Activities.Workflows/Activities/Signaled.cs @@ -32,7 +32,7 @@ namespace Elsa.Activities.Workflows.Activities protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context) { - return Halt(true); + return Suspend(true); } protected override IActivityExecutionResult OnResume(ActivityExecutionContext context) diff --git a/src/core/Elsa.Abstractions/Services/Extensions/WorkflowInvokerExtensions.cs b/src/core/Elsa.Abstractions/Services/Extensions/WorkflowInvokerExtensions.cs index 9435989ac..08cb4f928 100644 --- a/src/core/Elsa.Abstractions/Services/Extensions/WorkflowInvokerExtensions.cs +++ b/src/core/Elsa.Abstractions/Services/Extensions/WorkflowInvokerExtensions.cs @@ -12,7 +12,7 @@ namespace Elsa.Services.Extensions public static Task TriggerAsync( this IWorkflowRunner workflowRunner, string activityType, - Variables input, + Variable input, CancellationToken cancellationToken = default) { return workflowRunner.TriggerAsync(activityType, input, cancellationToken: cancellationToken); diff --git a/src/core/Elsa.Abstractions/Services/Models/WorkflowBlueprint.cs b/src/core/Elsa.Abstractions/Services/Models/WorkflowBlueprint.cs index 57c7f023e..9f3a9cddb 100644 --- a/src/core/Elsa.Abstractions/Services/Models/WorkflowBlueprint.cs +++ b/src/core/Elsa.Abstractions/Services/Models/WorkflowBlueprint.cs @@ -17,6 +17,8 @@ namespace Elsa.Services.Models bool isDisabled = false, string? name = default, string? description = default, + bool isLatest = false, + bool isPublished = false, IEnumerable? activities = default, IEnumerable? connections = default) { @@ -24,6 +26,8 @@ namespace Elsa.Services.Models Version = version; IsSingleton = isSingleton; IsDisabled = isDisabled; + IsLatest = isLatest; + IsPublished = isPublished; Name = name; Description = description; Activities = activities?.ToList() ?? new List(); diff --git a/src/core/Elsa.Core/Expressions/VariableExpression.cs b/src/core/Elsa.Core/Expressions/VariableExpression.cs new file mode 100644 index 000000000..fb50c5a18 --- /dev/null +++ b/src/core/Elsa.Core/Expressions/VariableExpression.cs @@ -0,0 +1,24 @@ +using System; +using Elsa.Services.Models; + +namespace Elsa.Expressions +{ + public class VariableExpression : WorkflowExpression + { + public static string ExpressionType => "Variable"; + + public VariableExpression(string variableName, Type returnType) : base(ExpressionType, returnType) + { + VariableName = variableName; + } + + public string VariableName { get; } + } + + public class VariableExpression : VariableExpression, IWorkflowExpression + { + public VariableExpression(string variableName) : base(variableName, typeof(T)) + { + } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Expressions/VariableHandler.cs b/src/core/Elsa.Core/Expressions/VariableHandler.cs new file mode 100644 index 000000000..44f0ab86e --- /dev/null +++ b/src/core/Elsa.Core/Expressions/VariableHandler.cs @@ -0,0 +1,21 @@ +using System.Threading; +using System.Threading.Tasks; +using Elsa.Services.Models; + +namespace Elsa.Expressions +{ + public class VariableHandler : IWorkflowExpressionHandler + { + public string Type => VariableExpression.ExpressionType; + + public Task EvaluateAsync( + IWorkflowExpression expression, + ActivityExecutionContext context, + CancellationToken cancellationToken) + { + var variableExpression = (VariableExpression)expression; + var result = context.GetVariable(variableExpression.VariableName); + return Task.FromResult(result); + } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs index e5311dbb4..bbc2f8fe1 100644 --- a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs +++ b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs @@ -75,6 +75,7 @@ namespace Microsoft.Extensions.DependencyInjection .TryAddProvider(ServiceLifetime.Singleton) .TryAddProvider(ServiceLifetime.Singleton) .TryAddProvider(ServiceLifetime.Singleton) + .TryAddProvider(ServiceLifetime.Singleton) .AddTransient() .AddScoped() .AddScoped() diff --git a/src/core/Elsa.Core/Messages/Handlers/PersistenceWorkflowEventHandler.cs b/src/core/Elsa.Core/Messages/Handlers/PersistenceWorkflowEventHandler.cs index 89335307c..13ced382a 100644 --- a/src/core/Elsa.Core/Messages/Handlers/PersistenceWorkflowEventHandler.cs +++ b/src/core/Elsa.Core/Messages/Handlers/PersistenceWorkflowEventHandler.cs @@ -40,8 +40,11 @@ namespace Elsa.Messages.Handlers public async Task Handle(WorkflowCompleted notification, CancellationToken cancellationToken) { - if (notification.Workflow.Blueprint.DeleteCompletedWorkflows) - await workflowInstanceStore.DeleteAsync(notification.Workflow.Id, cancellationToken); + var workflow = notification.Workflow; + var blueprint = workflow.Blueprint; + + if (blueprint.DeleteCompletedWorkflows || blueprint.PersistenceBehavior == WorkflowPersistenceBehavior.Suspended) + await workflowInstanceStore.DeleteAsync(workflow.Id, cancellationToken); } private async Task SaveWorkflowAsync(Workflow workflow, CancellationToken cancellationToken) diff --git a/src/core/Elsa.Core/Persistence/Memory/MemoryWorkflowInstanceStore.cs b/src/core/Elsa.Core/Persistence/Memory/MemoryWorkflowInstanceStore.cs index 960e738f6..b47ecb9f7 100644 --- a/src/core/Elsa.Core/Persistence/Memory/MemoryWorkflowInstanceStore.cs +++ b/src/core/Elsa.Core/Persistence/Memory/MemoryWorkflowInstanceStore.cs @@ -51,7 +51,7 @@ namespace Elsa.Persistence.Memory { var query = workflowInstances.Values.AsQueryable(); - query = query.Where(x => x.Status == WorkflowStatus.Running); + query = query.Where(x => x.Status == WorkflowStatus.Suspended); if (!string.IsNullOrWhiteSpace(correlationId)) query = query.Where(x => x.CorrelationId == correlationId); diff --git a/src/core/Elsa.Core/Services/Activity.cs b/src/core/Elsa.Core/Services/Activity.cs index 05e04408e..d978eb613 100644 --- a/src/core/Elsa.Core/Services/Activity.cs +++ b/src/core/Elsa.Core/Services/Activity.cs @@ -8,7 +8,7 @@ namespace Elsa.Services { public abstract class Activity : ActivityBase { - protected SuspendWorkflowResult Halt(bool continueOnFirstPass = false) => new SuspendWorkflowResult(continueOnFirstPass); + protected SuspendWorkflowResult Suspend(bool continueOnFirstPass = false) => new SuspendWorkflowResult(continueOnFirstPass); protected OutcomeResult Outcomes(IEnumerable names) => new OutcomeResult(names); protected OutcomeResult Outcome(string name) => Outcomes(new[] { name }); protected OutcomeResult Outcome(string name, object output) => Outcome(name, Variable.From(output)); diff --git a/src/core/Elsa.Core/Services/WorkflowFactory.cs b/src/core/Elsa.Core/Services/WorkflowFactory.cs index 24db427b9..b28619b47 100644 --- a/src/core/Elsa.Core/Services/WorkflowFactory.cs +++ b/src/core/Elsa.Core/Services/WorkflowFactory.cs @@ -36,12 +36,22 @@ namespace Elsa.Services var workflowDefinition = workflowBuilder().Build(); return CreateWorkflow(workflowDefinition, input, workflowInstance, correlationId); } - + public WorkflowBlueprint CreateWorkflowBlueprint(WorkflowDefinitionVersion definition) { var activities = CreateActivities(definition.Activities).ToList(); var connections = CreateConnections(definition.Connections, activities).ToList(); - return new WorkflowBlueprint(definition.DefinitionId, definition.Version, definition.IsSingleton, definition.IsDisabled, definition.Name, definition.Description, activities, connections); + return new WorkflowBlueprint( + definition.DefinitionId, + definition.Version, + definition.IsSingleton, + definition.IsDisabled, + definition.Name, + definition.Description, + definition.IsLatest, + definition.IsPublished, + activities, + connections); } public Workflow CreateWorkflow( @@ -52,7 +62,7 @@ namespace Elsa.Services { if (blueprint.IsDisabled) throw new InvalidOperationException("Cannot instantiate disabled workflow definitions."); - + var id = idGenerator.Generate(); var workflow = new Workflow( id, diff --git a/src/core/Elsa.Core/Services/WorkflowRunner.cs b/src/core/Elsa.Core/Services/WorkflowRunner.cs index 4581befbe..4dd315390 100644 --- a/src/core/Elsa.Core/Services/WorkflowRunner.cs +++ b/src/core/Elsa.Core/Services/WorkflowRunner.cs @@ -318,15 +318,13 @@ namespace Elsa.Services ? await ExecuteActivityAsync(workflowExecutionContext, scheduledActivity, cancellationToken) : await ResumeActivityAsync(workflowExecutionContext, scheduledActivity, cancellationToken); - workflowExecutionContext.IsFirstPass = false; - start = true; - await mediator.Publish(new ActivityExecuted(workflow, activity), cancellationToken); - if (result == null) - break; - - await result.ExecuteAsync(this, workflowExecutionContext, cancellationToken); + if (result != null) + await result.ExecuteAsync(this, workflowExecutionContext, cancellationToken); + + workflowExecutionContext.IsFirstPass = false; + start = true; } // Determine new workflow state. @@ -470,7 +468,7 @@ namespace Elsa.Services { var instances = await workflowInstanceStore.ListByStatusAsync( definition.Item1.DefinitionId, - WorkflowStatus.Running, + WorkflowStatus.Suspended, cancellationToken ); diff --git a/src/core/Elsa.Core/WorkflowBuilders/WorkflowBuilder.cs b/src/core/Elsa.Core/WorkflowBuilders/WorkflowBuilder.cs index a7efd8055..a2079fc5a 100644 --- a/src/core/Elsa.Core/WorkflowBuilders/WorkflowBuilder.cs +++ b/src/core/Elsa.Core/WorkflowBuilders/WorkflowBuilder.cs @@ -135,7 +135,7 @@ namespace Elsa.WorkflowBuilders var connections = CreateConnections(connectionDefinitions, activities); var definitionId = !string.IsNullOrWhiteSpace(Id) ? Id : idGenerator.Generate(); - return new WorkflowBlueprint(definitionId, Version, IsSingleton, IsDisabled, Name, Description, activities, connections); + return new WorkflowBlueprint(definitionId, Version, IsSingleton, IsDisabled, Name, Description, true, true, activities, connections); } private IEnumerable CreateConnections(IEnumerable connectionDefinitions, IEnumerable activities) diff --git a/src/samples/Sample02/Program.cs b/src/samples/Sample02/Program.cs index a33abfccd..519d490df 100644 --- a/src/samples/Sample02/Program.cs +++ b/src/samples/Sample02/Program.cs @@ -25,7 +25,7 @@ namespace Sample02 var workflowBuilderFactory = services.GetRequiredService>(); var workflowBuilder = workflowBuilderFactory(); var workflowBlueprint = workflowBuilder - .StartWith(x => x.Text = new LiteralExpression("Hello world!")) + .StartWith(x => x.Text = new CodeExpression(() => "Hello world!")) .Then(() => Console.WriteLine("Look, custom code!")) .Then(x => x.Text = new LiteralExpression("Goodbye cruel world...")) .Build(); diff --git a/src/samples/Sample05/RecurringWorkflow.cs b/src/samples/Sample05/RecurringWorkflow.cs index f94658170..a9cc771ad 100644 --- a/src/samples/Sample05/RecurringWorkflow.cs +++ b/src/samples/Sample05/RecurringWorkflow.cs @@ -16,7 +16,7 @@ namespace Sample05 .WithId("RecurringWorkflow") .AsSingleton() .StartWith(x => x.TimeoutScriptExpression = new LiteralExpression("00:00:05")) - .Then(x => x.Text = new JavaScriptExpression("`Trigger received. The time is: ${new Date().toISOString()}`")); + .Then(x => x.Text = new CodeExpression(context => $"Trigger received. The time is: {DateTime.UtcNow.ToLocalTime()}")); } } } \ No newline at end of file diff --git a/src/samples/Sample08/Workflows/CreateOrderWorkflow.cs b/src/samples/Sample08/Workflows/CreateOrderWorkflow.cs index 33b72f4e5..8263a32ad 100644 --- a/src/samples/Sample08/Workflows/CreateOrderWorkflow.cs +++ b/src/samples/Sample08/Workflows/CreateOrderWorkflow.cs @@ -37,7 +37,7 @@ namespace Sample08.Workflows ) // Need to ensure that the correlation ID is the same string format that is used by the WorkflowConsumer // MassTransit will always use Guid values for correlation ID's, so need to ensure the same string format is used. - .Then(activity => activity.ValueScriptExpression = new JavaScriptExpression("newGuid()")) + .Then(activity => activity.Value = new CodeExpression(() => Guid.NewGuid().ToString())) .Then(activity => { activity.Message = new JavaScriptExpression("return { correlationId: correlationId(), order: order};"); diff --git a/src/samples/Sample13/Program.cs b/src/samples/Sample13/Program.cs index 144bbe0ed..e631cc97c 100644 --- a/src/samples/Sample13/Program.cs +++ b/src/samples/Sample13/Program.cs @@ -8,7 +8,7 @@ using Microsoft.Extensions.DependencyInjection; namespace Sample13 { /// - /// A strongly-typed workflows program demonstrating scripting, and branching. + /// A strongly-typed workflows program demonstrating scripting and branching. /// internal class Program {