From dcbf14c144cfc3bbd0971c50b6df32f1c6e9caa4 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Fri, 1 Jan 2021 17:58:08 +0100 Subject: [PATCH] Fix and simplify composite activity + while/for/foreach handling --- .../Services/QuartzWorkflowScheduler.cs | 10 ++-- .../Services/WorkflowRunnerQueue.cs | 32 +++++++++--- .../ActivityResults/ScheduleWorkflowResult.cs | 1 - .../Elsa.Client/Models/WorkflowInstance.cs | 3 +- .../PostScheduleActivitiesResult.cs | 20 ------- .../CompositeActivityBlueprintExtensions.cs | 4 +- .../Models/WorkflowInstance.cs | 2 - .../Elsa.Abstractions/Services/Activity.cs | 1 - .../Services/CompositeActivity.cs | 35 +++---------- .../Models/WorkflowExecutionContext.cs | 15 ------ .../Activities/ControlFlow/IfThen/Switch.cs | 1 - .../Builders/CompositeActivityBuilder.cs | 42 ++++++++------- src/core/Elsa.Core/ElsaOptions.cs | 9 ++++ src/core/Elsa.Core/Services/WorkflowRunner.cs | 52 +++---------------- .../Display/TimersDisplayProvider.cs | 2 - .../Extensions/ServiceCollectionExtensions.cs | 2 - .../Documents/WorkflowInstanceDocument.cs | 1 - .../WorkflowDefinitionDocumentExtensions.cs | 4 +- .../Stores/YesSqlWorkflowDefinitionStore.cs | 1 - .../Activities/CountDownActivity.cs | 12 +++-- .../Workflows/CompositionWorkflow.cs | 8 ++- .../{ => Activities}/MyContainer1.cs | 2 +- .../Activities/MyContainer2.cs | 41 +++++++++++++++ .../Elsa.Samples.Timers/MyContainer2.cs | 35 ------------- .../worker/Elsa.Samples.Timers/Program.cs | 11 ++-- .../RecurringTaskWorkflow.cs | 41 --------------- .../{ => Workflows}/CancelTimerWorkflow.cs | 5 +- .../{ => Workflows}/CronTaskWorkflow.cs | 2 +- .../{ => Workflows}/OneOffWorkflow.cs | 10 ++-- .../Workflows/RecurringTaskWorkflow.cs | 25 +++++++++ 30 files changed, 171 insertions(+), 258 deletions(-) delete mode 100644 src/core/Elsa.Abstractions/ActivityResults/PostScheduleActivitiesResult.cs rename src/samples/worker/Elsa.Samples.Timers/{ => Activities}/MyContainer1.cs (97%) create mode 100644 src/samples/worker/Elsa.Samples.Timers/Activities/MyContainer2.cs delete mode 100644 src/samples/worker/Elsa.Samples.Timers/MyContainer2.cs delete mode 100644 src/samples/worker/Elsa.Samples.Timers/RecurringTaskWorkflow.cs rename src/samples/worker/Elsa.Samples.Timers/{ => Workflows}/CancelTimerWorkflow.cs (91%) rename src/samples/worker/Elsa.Samples.Timers/{ => Workflows}/CronTaskWorkflow.cs (90%) rename src/samples/worker/Elsa.Samples.Timers/{ => Workflows}/OneOffWorkflow.cs (59%) create mode 100644 src/samples/worker/Elsa.Samples.Timers/Workflows/RecurringTaskWorkflow.cs diff --git a/src/activities/Elsa.Activities.Timers.Quartz/Services/QuartzWorkflowScheduler.cs b/src/activities/Elsa.Activities.Timers.Quartz/Services/QuartzWorkflowScheduler.cs index 4bcb0e417..f13dd40cd 100644 --- a/src/activities/Elsa.Activities.Timers.Quartz/Services/QuartzWorkflowScheduler.cs +++ b/src/activities/Elsa.Activities.Timers.Quartz/Services/QuartzWorkflowScheduler.cs @@ -63,8 +63,8 @@ namespace Elsa.Activities.Timers.Quartz.Services var existingTrigger = await scheduler.GetTrigger(trigger, cancellationToken); - // if (existingTrigger != null) - // await scheduler.UnscheduleJob(existingTrigger.Key, cancellationToken); + if (existingTrigger != null) + await scheduler.UnscheduleJob(existingTrigger.Key, cancellationToken); } private async Task ScheduleJob(ITrigger trigger, CancellationToken cancellationToken) @@ -76,8 +76,8 @@ namespace Elsa.Activities.Timers.Quartz.Services var scheduler = await _schedulerFactory.GetScheduler(cancellationToken); var existingTrigger = await scheduler.GetTrigger(trigger.Key, cancellationToken); - // if (existingTrigger != null) - // await scheduler.UnscheduleJob(existingTrigger.Key, cancellationToken); + if (existingTrigger != null) + await scheduler.UnscheduleJob(existingTrigger.Key, cancellationToken); await scheduler.ScheduleJob(trigger, cancellationToken); } @@ -104,7 +104,7 @@ namespace Elsa.Activities.Timers.Quartz.Services private TriggerKey CreateTriggerKey(string? tenantId, string workflowDefinitionId, string? workflowInstanceId, string activityId) { var groupName = $"tenant:{tenantId ?? "default"}-workflow-instance:{workflowInstanceId ?? workflowDefinitionId}"; - return new TriggerKey($"activity:{activityId}:{Guid.NewGuid()}", groupName); + return new TriggerKey($"activity:{activityId}", groupName); } } } diff --git a/src/activities/Elsa.Activities.Timers.Quartz/Services/WorkflowRunnerQueue.cs b/src/activities/Elsa.Activities.Timers.Quartz/Services/WorkflowRunnerQueue.cs index aefe13f55..d71e8f959 100644 --- a/src/activities/Elsa.Activities.Timers.Quartz/Services/WorkflowRunnerQueue.cs +++ b/src/activities/Elsa.Activities.Timers.Quartz/Services/WorkflowRunnerQueue.cs @@ -9,6 +9,7 @@ using Microsoft.Extensions.Logging; namespace Elsa.Activities.Timers.Quartz.Services { + // TODO: Consider turning this into a global service to allow background, sequential execution of a given workflow instance (but allow for parallel execution of workflow instances with different definitions). Executing the same workflow instances sequentially prevents loss of data during update concurrency conflicts. public class WorkflowRunnerQueue { private readonly IBackgroundWorker _backgroundWorker; @@ -21,7 +22,7 @@ namespace Elsa.Activities.Timers.Quartz.Services _serviceProvider = serviceProvider; _logger = logger; } - + public async Task Enqueue(string workflowInstanceId, string activityId, CancellationToken cancellationToken = default) { await _backgroundWorker.ScheduleTask(async () => await RunWorkflowAsync(workflowInstanceId, activityId, cancellationToken), cancellationToken); @@ -33,19 +34,34 @@ namespace Elsa.Activities.Timers.Quartz.Services var store = scope.ServiceProvider.GetRequiredService(); var workflowInstance = await store.FindByIdAsync(workflowInstanceId, cancellationToken); - if (workflowInstance == null) - { - _logger.LogError("Could not run Workflow instance with ID {WorkflowInstanceId} because it is not in the database", workflowInstanceId); + if (!ValidatePreconditions(workflowInstanceId, workflowInstance)) return; - } - - //await context.Scheduler.UnscheduleJob(context.Trigger.Key, cancellationToken); - var workflowDefinitionId = workflowInstance.DefinitionId; + + _logger.LogDebug("Running {WorkflowInstanceId} with status {WorkflowStatus}.", workflowInstance!.WorkflowStatus); + + var workflowDefinitionId = workflowInstance!.DefinitionId; var tenantId = workflowInstance.TenantId; var workflowRegistry = scope.ServiceProvider.GetRequiredService(); var workflowBlueprint = (await workflowRegistry.GetWorkflowAsync(workflowDefinitionId, tenantId, VersionOptions.SpecificVersion(workflowInstance.Version), cancellationToken))!; var workflowRunner = scope.ServiceProvider.GetRequiredService(); await workflowRunner.RunWorkflowAsync(workflowBlueprint, workflowInstance!, activityId, cancellationToken: cancellationToken); } + + private bool ValidatePreconditions(string? workflowInstanceId, WorkflowInstance? workflowInstance) + { + if (workflowInstance == null) + { + _logger.LogError("Could not run workflow instance with ID {WorkflowInstanceId} because it does not exist.", workflowInstanceId); + return false; + } + + if (workflowInstance.WorkflowStatus != WorkflowStatus.Suspended) + { + _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); + return false; + } + + return true; + } } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Timers/ActivityResults/ScheduleWorkflowResult.cs b/src/activities/Elsa.Activities.Timers/ActivityResults/ScheduleWorkflowResult.cs index ef5bd49f3..42c578770 100644 --- a/src/activities/Elsa.Activities.Timers/ActivityResults/ScheduleWorkflowResult.cs +++ b/src/activities/Elsa.Activities.Timers/ActivityResults/ScheduleWorkflowResult.cs @@ -1,4 +1,3 @@ -using System; using System.Threading; using System.Threading.Tasks; using Elsa.Activities.Timers.Services; diff --git a/src/clients/Elsa.Client/Models/WorkflowInstance.cs b/src/clients/Elsa.Client/Models/WorkflowInstance.cs index a9b91dd55..40c0f0588 100644 --- a/src/clients/Elsa.Client/Models/WorkflowInstance.cs +++ b/src/clients/Elsa.Client/Models/WorkflowInstance.cs @@ -15,7 +15,6 @@ namespace Elsa.Client.Models Variables = new Variables(); Activities = new List(); ScheduledActivities = new Stack(); - PostScheduledActivities = new Stack(); } [DataMember(Order = 1)] public string Id { get; set; } = default!; @@ -42,6 +41,6 @@ namespace Elsa.Client.Models [DataMember(Order = 16)] public WorkflowFault? Fault { get; set; } [DataMember(Order = 17)] public Stack ScheduledActivities { get; set; } - [DataMember(Order = 18)] public Stack PostScheduledActivities { get; set; } + [DataMember(Order = 18)]public Stack ParentActivities { get; set; } } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/ActivityResults/PostScheduleActivitiesResult.cs b/src/core/Elsa.Abstractions/ActivityResults/PostScheduleActivitiesResult.cs deleted file mode 100644 index 404c6af8d..000000000 --- a/src/core/Elsa.Abstractions/ActivityResults/PostScheduleActivitiesResult.cs +++ /dev/null @@ -1,20 +0,0 @@ -using System.Collections.Generic; -using System.Linq; -using Elsa.Models; -using Elsa.Services.Models; - -namespace Elsa.ActivityResults -{ - public class PostScheduleActivitiesResult : ActivityExecutionResult - { - public PostScheduleActivitiesResult(IEnumerable activityIds, object? input = default) => - Activities = activityIds.Select(x => new ScheduledActivity(x, input)); - - public PostScheduleActivitiesResult(IEnumerable activities) => Activities = activities; - - public IEnumerable Activities { get; } - - protected override void Execute(ActivityExecutionContext activityExecutionContext) => - activityExecutionContext.WorkflowExecutionContext.PostScheduleActivities(Activities); - } -} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Extensions/CompositeActivityBlueprintExtensions.cs b/src/core/Elsa.Abstractions/Extensions/CompositeActivityBlueprintExtensions.cs index 52f76c5dd..95c4cfa1d 100644 --- a/src/core/Elsa.Abstractions/Extensions/CompositeActivityBlueprintExtensions.cs +++ b/src/core/Elsa.Abstractions/Extensions/CompositeActivityBlueprintExtensions.cs @@ -11,8 +11,7 @@ namespace Elsa { public static IEnumerable GetStartActivities(this ICompositeActivityBlueprint workflowBlueprint) { - var targetActivityIds = workflowBlueprint.Connections.Select(x => x.Target.Activity.Id).Distinct() - .ToLookup(x => x); + var targetActivityIds = workflowBlueprint.Connections.Select(x => x.Target.Activity?.Id).Distinct().ToLookup(x => x); var query = from activity in workflowBlueprint.Activities @@ -25,6 +24,7 @@ namespace Elsa public static IEnumerable GetStartActivities(this ICompositeActivityBlueprint workflowBlueprint, string activityType) => workflowBlueprint.GetStartActivities().Where(x => x.Type == activityType); public static IEnumerable GetStartActivities(this ICompositeActivityBlueprint workflowBlueprint, Type activityType) => workflowBlueprint.GetStartActivities(activityType.Name); public static IEnumerable GetStartActivities(this ICompositeActivityBlueprint workflowBlueprint) where T : IActivity => workflowBlueprint.GetStartActivities(typeof(T)); + public static IEnumerable GetEndActivities(this ICompositeActivityBlueprint workflowBlueprint) => workflowBlueprint.Activities.Where(x => !workflowBlueprint.GetOutboundConnections(x.Id).Any()); public static IActivityBlueprint? GetActivity(this ICompositeActivityBlueprint workflowBlueprint, string id) => workflowBlueprint.Activities.FirstOrDefault(x => x.Id == id); public static IEnumerable GetActivities(this ICompositeActivityBlueprint workflowBlueprint, IEnumerable ids) => workflowBlueprint.Activities.Where(x => ids.Contains(x.Id)); diff --git a/src/core/Elsa.Abstractions/Models/WorkflowInstance.cs b/src/core/Elsa.Abstractions/Models/WorkflowInstance.cs index b4313fb21..f9ce84ee1 100644 --- a/src/core/Elsa.Abstractions/Models/WorkflowInstance.cs +++ b/src/core/Elsa.Abstractions/Models/WorkflowInstance.cs @@ -13,7 +13,6 @@ namespace Elsa.Models Variables = new Variables(); Activities = new List(); ScheduledActivities = new Stack(); - PostScheduledActivities = new Stack(); ParentActivities = new Stack(); } @@ -43,7 +42,6 @@ namespace Elsa.Models public WorkflowFault? Fault { get; set; } public Stack ScheduledActivities { get; set; } - public Stack PostScheduledActivities { get; set; } public Stack ParentActivities { get; set; } } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/Activity.cs b/src/core/Elsa.Abstractions/Services/Activity.cs index 84c7b0fd7..12db78b3b 100644 --- a/src/core/Elsa.Abstractions/Services/Activity.cs +++ b/src/core/Elsa.Abstractions/Services/Activity.cs @@ -45,7 +45,6 @@ namespace Elsa.Services protected ScheduleActivitiesResult Schedule(IEnumerable activityIds, object? input) => new(activityIds, input); protected ScheduleActivitiesResult Schedule(string activityId, object? input) => Schedule(new[] { activityId }, input); protected ScheduleActivitiesResult Schedule(IEnumerable activities) => new(activities); - protected PostScheduleActivitiesResult PostSchedule(params string[] activityIds) => new(activityIds); protected CombinedResult Combine(IEnumerable results) => new(results); protected CombinedResult Combine(params IActivityExecutionResult[] results) => new(results); protected FaultResult Fault(LocalizedString message) => new(message); diff --git a/src/core/Elsa.Abstractions/Services/CompositeActivity.cs b/src/core/Elsa.Abstractions/Services/CompositeActivity.cs index 8c60c6b5a..8951db1c5 100644 --- a/src/core/Elsa.Abstractions/Services/CompositeActivity.cs +++ b/src/core/Elsa.Abstractions/Services/CompositeActivity.cs @@ -1,5 +1,4 @@ -using System.Linq; -using Elsa.ActivityResults; +using Elsa.ActivityResults; using Elsa.Builders; using Elsa.Services.Models; @@ -7,6 +6,8 @@ namespace Elsa.Services { public class CompositeActivity : Activity { + internal const string Enter = "Enter"; + public virtual void Build(ICompositeActivityBuilder activity) { } @@ -19,34 +20,14 @@ namespace Elsa.Services protected override IActivityExecutionResult OnExecute(ActivityExecutionContext context) { - if (IsScheduled) + if (!IsScheduled) { - if (HasPendingChildren(context)) - return PostSchedule(Id); - - context.WorkflowExecutionContext.WorkflowInstance.ParentActivities.Pop(); - IsScheduled = false; - return Complete(context); + IsScheduled = true; + return Outcome(Enter); } - var compositeActivityBlueprint = (ICompositeActivityBlueprint) context.ActivityBlueprint; - var startActivities = compositeActivityBlueprint.GetStartActivities().Select(x => x.Id).ToList(); - context.WorkflowExecutionContext.WorkflowInstance.ParentActivities.Push(Id); - context.WorkflowExecutionContext.PostScheduleActivity(Id); - IsScheduled = true; - return Schedule(startActivities, context.Input); - } - - protected virtual IActivityExecutionResult Complete(ActivityExecutionContext context) => Done(); - - private static bool HasPendingChildren(ActivityExecutionContext context) - { - var children = ((CompositeActivityBlueprint) context.ActivityBlueprint).Activities.Select(x => x.Id).ToList(); - var workflowInstance = context.WorkflowExecutionContext.WorkflowInstance; - //var hasPendingPostScheduledChildren = workflowInstance.PostScheduledActivities.Any(x => children.Contains(x.ActivityId)); - var hasPendingScheduledChildren = workflowInstance.ScheduledActivities.Any(x => children.Contains(x.ActivityId)); - //return hasPendingPostScheduledChildren || hasPendingScheduledChildren; - return hasPendingScheduledChildren; + IsScheduled = false; + return Done(); } } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs b/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs index e724ba99e..eb9e4e0b5 100644 --- a/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs +++ b/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs @@ -33,7 +33,6 @@ namespace Elsa.Services.Models public JsonSerializer Serializer { get; } public object? Input { get; } public bool HasScheduledActivities => WorkflowInstance.ScheduledActivities.Any(); - public bool HasPostScheduledActivities => WorkflowInstance.PostScheduledActivities.Any(); public IWorkflowFault? WorkflowFault { get; private set; } public bool IsFirstPass { get; private set; } public bool ContextHasChanged { get; set; } @@ -55,16 +54,8 @@ namespace Elsa.Services.Models ScheduleActivity(activity); } - public void PostScheduleActivities(IEnumerable activities) - { - foreach (var activity in activities) - PostScheduleActivity(activity); - } - public void ScheduleActivity(string activityId, object? input = default) => ScheduleActivity(new ScheduledActivity(activityId, input)); public void ScheduleActivity(ScheduledActivity activity) => WorkflowInstance.ScheduledActivities.Push(activity); - public void PostScheduleActivity(string activityId, object? input = default) => PostScheduleActivity(new ScheduledActivity(activityId, input)); - public void PostScheduleActivity(ScheduledActivity activity) => WorkflowInstance.PostScheduledActivities.Push(activity); public ScheduledActivity PopScheduledActivity() => WorkflowInstance.ScheduledActivities.Pop(); public ScheduledActivity PeekScheduledActivity() => WorkflowInstance.ScheduledActivities.Peek(); @@ -102,12 +93,6 @@ namespace Elsa.Services.Models public IActivityBlueprint? GetActivityBlueprintById(string id) => WorkflowBlueprint.Activities.FirstOrDefault(x => x.Id == id); public IActivityBlueprint? GetActivityBlueprintByName(string name) => WorkflowBlueprint.Activities.FirstOrDefault(x => x.Name == name); - public void SchedulePostActivity() - { - var activity = WorkflowInstance.PostScheduledActivities.Pop(); - ScheduleActivity(activity); - } - public object? GetOutputFrom(string activityName) { var activityBlueprint = GetActivityBlueprintByName(activityName)!; diff --git a/src/core/Elsa.Core/Activities/ControlFlow/IfThen/Switch.cs b/src/core/Elsa.Core/Activities/ControlFlow/IfThen/Switch.cs index 0a3c5a85e..c7aad379b 100644 --- a/src/core/Elsa.Core/Activities/ControlFlow/IfThen/Switch.cs +++ b/src/core/Elsa.Core/Activities/ControlFlow/IfThen/Switch.cs @@ -1,4 +1,3 @@ -using System; using System.Collections.Generic; using System.Linq; using Elsa.ActivityResults; diff --git a/src/core/Elsa.Core/Builders/CompositeActivityBuilder.cs b/src/core/Elsa.Core/Builders/CompositeActivityBuilder.cs index 92b98d6af..1dfe9af26 100644 --- a/src/core/Elsa.Core/Builders/CompositeActivityBuilder.cs +++ b/src/core/Elsa.Core/Builders/CompositeActivityBuilder.cs @@ -10,7 +10,7 @@ namespace Elsa.Builders { public class CompositeActivityBuilder : ActivityBuilder, ICompositeActivityBuilder { - private readonly Func _workflowBuilderFactory; + private readonly Func _compositeActivityBuilderFactory; public CompositeActivityBuilder(IServiceProvider serviceProvider) { @@ -18,7 +18,7 @@ namespace Elsa.Builders ActivityBuilders = new List(); ConnectionBuilders = new List(); - _workflowBuilderFactory = () => + _compositeActivityBuilderFactory = () => { var builder = serviceProvider.GetRequiredService(); builder.WorkflowBuilder = WorkflowBuilder; @@ -138,14 +138,13 @@ namespace Elsa.Builders activityBuilder.ActivityId = $"{activityIdPrefix}-{++index}"; activityBlueprints.AddRange(activityBuilders.Select(x => BuildActivityBlueprint(x, compositeActivityBlueprint))); - + var activityBlueprintDictionary = activityBlueprints.ToDictionary(x => x.Id); + connections.AddRange(ConnectionBuilders.Select(x => new Connection(activityBlueprintDictionary[x.Source().ActivityId], activityBlueprintDictionary[x.Target().ActivityId], x.Outcome))); + // Build composite activities. var compositeActivityBuilders = activityBuilders.Where(x => typeof(CompositeActivity).IsAssignableFrom(x.ActivityType)); BuildCompositeActivities(compositeActivityBuilders, activityBlueprints, connections, activityPropertyProviders); - var activityBlueprintDictionary = activityBlueprints.ToDictionary(x => x.Id); - - connections.AddRange(ConnectionBuilders.Select(x => new Connection(activityBlueprintDictionary[x.Source().ActivityId], activityBlueprintDictionary[x.Target().ActivityId], x.Outcome))); - + activityPropertyProviders.AddRange( activityBuilders .Select(x => (x.ActivityId, x.PropertyValueProviders)) @@ -168,21 +167,24 @@ namespace Elsa.Builders foreach (var activityBuilder in compositeActivityBuilders) { var compositeActivity = (CompositeActivity) ActivatorUtilities.CreateInstance(scope.ServiceProvider, activityBuilder.ActivityType); - var workflowBuilder = _workflowBuilderFactory(); - workflowBuilder.ActivityId = activityBuilder.ActivityId; - compositeActivity.Build(workflowBuilder); + var compositeActivityBuilder = _compositeActivityBuilderFactory(); + compositeActivityBuilder.ActivityId = activityBuilder.ActivityId; + compositeActivity.Build(compositeActivityBuilder); - var workflow = workflowBuilder.Build($"{activityBuilder.ActivityId}:activity"); - var activityDictionary = workflow.Activities.ToDictionary(x => x.Id); + var compositeActivityBlueprint = compositeActivityBuilder.Build($"{activityBuilder.ActivityId}:activity"); + var activityDictionary = compositeActivityBlueprint.Activities.ToDictionary(x => x.Id); - activityBlueprints.AddRange(workflow.Activities); - connections.AddRange(workflow.Connections.Select(x => new Connection(activityDictionary[x.Source.Activity.Id], activityDictionary[x.Target.Activity.Id], x.Source.Outcome))); - activityPropertyProviders.AddRange(workflow.ActivityPropertyProviders); - - var compositeActivityBlueprint = (ICompositeActivityBlueprint) activityBlueprints.Single(x => x.Id == activityBuilder.ActivityId); - compositeActivityBlueprint.Activities = workflow.Activities; - compositeActivityBlueprint.Connections = workflow.Connections; - compositeActivityBlueprint.ActivityPropertyProviders = workflow.ActivityPropertyProviders; + activityBlueprints.AddRange(compositeActivityBlueprint.Activities); + connections.AddRange(compositeActivityBlueprint.Connections.Select(x => new Connection(activityDictionary[x.Source.Activity.Id], activityDictionary[x.Target.Activity.Id], x.Source.Outcome))); + activityPropertyProviders.AddRange(compositeActivityBlueprint.ActivityPropertyProviders); + + compositeActivityBlueprint.Activities = compositeActivityBlueprint.Activities; + compositeActivityBlueprint.Connections = compositeActivityBlueprint.Connections; + compositeActivityBlueprint.ActivityPropertyProviders = compositeActivityBlueprint.ActivityPropertyProviders; + + // Connect the composite activity to its starting activities. + var startActivities = compositeActivityBlueprint.GetStartActivities().ToList(); + connections.AddRange(startActivities.Select(x => new Connection(compositeActivityBlueprint, x, CompositeActivity.Enter))); } } diff --git a/src/core/Elsa.Core/ElsaOptions.cs b/src/core/Elsa.Core/ElsaOptions.cs index d698427ea..5c9fe8a7f 100644 --- a/src/core/Elsa.Core/ElsaOptions.cs +++ b/src/core/Elsa.Core/ElsaOptions.cs @@ -118,6 +118,15 @@ namespace Elsa WorkflowFactory.Add(workflow.GetType(), workflow); return this; } + + public ElsaOptions AddWorkflow(Func workflow) where T: class, IWorkflow + { + Services.AddSingleton(workflow); + Services.AddSingleton(sp => sp.GetRequiredService()); + WorkflowFactory.Add(typeof(T), sp => sp.GetRequiredService()); + + return this; + } public ElsaOptions AddWorkflowsFrom() => AddWorkflowsFrom(typeof(T).Assembly); diff --git a/src/core/Elsa.Core/Services/WorkflowRunner.cs b/src/core/Elsa.Core/Services/WorkflowRunner.cs index ea64992f8..1d3529357 100644 --- a/src/core/Elsa.Core/Services/WorkflowRunner.cs +++ b/src/core/Elsa.Core/Services/WorkflowRunner.cs @@ -1,9 +1,7 @@ using System; -using System.Collections.Generic; using System.Linq; using System.Threading; using System.Threading.Tasks; -using Elsa.Activities.ControlFlow; using Elsa.ActivityProviders; using Elsa.ActivityResults; using Elsa.Builders; @@ -301,59 +299,21 @@ namespace Elsa.Services { var parentActivityBlueprint = inboundConnection.Source.Activity; - if (parentActivityBlueprint.Type == nameof(While) && inboundConnection.Source.Outcome == "Iterate") + if (inboundConnection.Source.Outcome == "Iterate") // This covers While/For/ForEach/ParallelForEach based on the convention of having an unclosed "Iterate" branch. This should be refactored by allowing activities to explicitly opt-in to being rescheduled once there are no more scheduled activities. { workflowExecutionContext.ScheduleActivity(parentActivityBlueprint.Id); break; } - - // if (parentActivityBlueprint is CompositeActivityBlueprint) - // { - // workflowExecutionContext.ScheduleActivity(parentActivityBlueprint.Id); - // break; - // } } - - // Schedule the parent activity, if any + } + + if (!workflowExecutionContext.HasScheduledActivities && workflowExecutionContext.Status == WorkflowStatus.Running) + { + // Re-schedule the parent activity, if any if (activityBlueprint.Parent != null && workflowBlueprint.GetActivity(activityBlueprint.Parent.Id) != null) { workflowExecutionContext.ScheduleActivity(activityBlueprint.Parent.Id); } - - // Find first parent that requires rescheduling for reevaluation. - // var parents = workflowBlueprint.GetInboundActivityPath(currentActivityId).ToList(); - // - // foreach (var parentId in parents) - // { - // var parentActivityBlueprint = workflowBlueprint.GetActivity(parentId)!; - // - // if (parentActivityBlueprint.Type == nameof(While) || parentActivityBlueprint is ICompositeActivityBlueprint) - // { - // workflowExecutionContext.ScheduleActivity(parentId); - // break; - // } - // } - - // if (scheduledActivity.ParentId != null) - // { - // workflowExecutionContext.ScheduleActivity(scheduledActivity.ParentId, null, null); - // } - // - // if (workflowExecutionContext.HasPostScheduledActivities) - // { - // workflowExecutionContext.SchedulePostActivity(); - // - // // // Get children of current activity's parent. - // // var parentActivity = activityBlueprint.Parent; - // // var childActivities = workflowBlueprint.Activities.Where(x => x.Parent == parentActivity).ToList(); - // // var childActivityIds = childActivities.Select(x => x.Id).ToList(); - // // var scheduledPostActivities = workflowExecutionContext.WorkflowInstance.PostScheduledActivities.Where(x => childActivityIds.Contains(x.ActivityId)).ToList(); - // // - // // foreach (var scheduledPostActivity in scheduledPostActivities) - // // workflowExecutionContext.ScheduleActivity(scheduledPostActivity); - // // - // // workflowExecutionContext.WorkflowInstance.PostScheduledActivities = new Stack(workflowExecutionContext.WorkflowInstance.PostScheduledActivities.Where(x => !childActivityIds.Contains(x.ActivityId)).Reverse().ToList()); - // } } } diff --git a/src/dashboards/blazor/ElsaDashboard.Application/Display/TimersDisplayProvider.cs b/src/dashboards/blazor/ElsaDashboard.Application/Display/TimersDisplayProvider.cs index b0e0de157..2ebceae77 100644 --- a/src/dashboards/blazor/ElsaDashboard.Application/Display/TimersDisplayProvider.cs +++ b/src/dashboards/blazor/ElsaDashboard.Application/Display/TimersDisplayProvider.cs @@ -1,7 +1,5 @@ using System.Collections.Generic; -using ElsaDashboard.Application.Activities; using ElsaDashboard.Application.Activities.Timers; -using ElsaDashboard.Extensions; using ElsaDashboard.Models; using ElsaDashboard.Services; diff --git a/src/dashboards/blazor/ElsaDashboard.Application/Extensions/ServiceCollectionExtensions.cs b/src/dashboards/blazor/ElsaDashboard.Application/Extensions/ServiceCollectionExtensions.cs index 820a2b9c0..e3a7347ae 100644 --- a/src/dashboards/blazor/ElsaDashboard.Application/Extensions/ServiceCollectionExtensions.cs +++ b/src/dashboards/blazor/ElsaDashboard.Application/Extensions/ServiceCollectionExtensions.cs @@ -1,6 +1,4 @@ using Blazored.Modal; -using ElsaDashboard.Application.Activities; -using ElsaDashboard.Application.Activities.Console; using ElsaDashboard.Application.Display; using ElsaDashboard.Application.Services; using ElsaDashboard.Extensions; diff --git a/src/persistence/Elsa.Persistence.YesSql/Documents/WorkflowInstanceDocument.cs b/src/persistence/Elsa.Persistence.YesSql/Documents/WorkflowInstanceDocument.cs index 870c0c32b..2e567dbbd 100644 --- a/src/persistence/Elsa.Persistence.YesSql/Documents/WorkflowInstanceDocument.cs +++ b/src/persistence/Elsa.Persistence.YesSql/Documents/WorkflowInstanceDocument.cs @@ -36,7 +36,6 @@ namespace Elsa.Persistence.YesSql.Documents public WorkflowFault? Fault { get; set; } public Stack ScheduledActivities { get; set; } = new(); - public Stack PostScheduledActivities { get; set; } = new(); public Stack ParentActivities { get; set; } = new(); } } \ No newline at end of file diff --git a/src/persistence/Elsa.Persistence.YesSql/Extensions/WorkflowDefinitionDocumentExtensions.cs b/src/persistence/Elsa.Persistence.YesSql/Extensions/WorkflowDefinitionDocumentExtensions.cs index 7cb2062c9..d1ee8e0f7 100644 --- a/src/persistence/Elsa.Persistence.YesSql/Extensions/WorkflowDefinitionDocumentExtensions.cs +++ b/src/persistence/Elsa.Persistence.YesSql/Extensions/WorkflowDefinitionDocumentExtensions.cs @@ -1,6 +1,4 @@ -using System; -using System.Linq.Expressions; -using Elsa.Models; +using Elsa.Models; using Elsa.Persistence.YesSql.Documents; using Elsa.Persistence.YesSql.Indexes; using YesSql; diff --git a/src/persistence/Elsa.Persistence.YesSql/Stores/YesSqlWorkflowDefinitionStore.cs b/src/persistence/Elsa.Persistence.YesSql/Stores/YesSqlWorkflowDefinitionStore.cs index 8894d18ef..3fa15ad36 100644 --- a/src/persistence/Elsa.Persistence.YesSql/Stores/YesSqlWorkflowDefinitionStore.cs +++ b/src/persistence/Elsa.Persistence.YesSql/Stores/YesSqlWorkflowDefinitionStore.cs @@ -1,4 +1,3 @@ -using System.Linq; using System.Threading; using System.Threading.Tasks; using AutoMapper; diff --git a/src/samples/console/Elsa.Samples.ProgrammaticCompositeActivitiesConsole/Activities/CountDownActivity.cs b/src/samples/console/Elsa.Samples.ProgrammaticCompositeActivitiesConsole/Activities/CountDownActivity.cs index 8939e6081..05bd13e27 100644 --- a/src/samples/console/Elsa.Samples.ProgrammaticCompositeActivitiesConsole/Activities/CountDownActivity.cs +++ b/src/samples/console/Elsa.Samples.ProgrammaticCompositeActivitiesConsole/Activities/CountDownActivity.cs @@ -1,6 +1,6 @@ -using Elsa.Activities.Console; +using System; +using Elsa.Activities.Console; using Elsa.Activities.ControlFlow; -using Elsa.ActivityResults; using Elsa.Attributes; using Elsa.Builders; using Elsa.Services; @@ -20,11 +20,13 @@ namespace Elsa.Samples.ProgrammaticCompositeActivitiesConsole.Activities .StartWith(GetInstructions) .WriteLine(context => (string)context.Input) .ReadLine() - .Finish(context => (string) context.Input); + .IfElse(context => string.Equals(context.GetInput(), "left", StringComparison.CurrentCultureIgnoreCase), ifElse => + { + ifElse.When(IfElse.True).WriteLine("We're going left"); + ifElse.When(IfElse.False).WriteLine("We're going right"); + }); } - protected override IActivityExecutionResult Complete(ActivityExecutionContext context) => Outcome(((string) context.WorkflowExecutionContext.WorkflowInstance.Output)!); - private static void GetInstructions(ActivityExecutionContext context) => context.Output = "Turn left or right?"; } } \ No newline at end of file diff --git a/src/samples/console/Elsa.Samples.ProgrammaticCompositeActivitiesConsole/Workflows/CompositionWorkflow.cs b/src/samples/console/Elsa.Samples.ProgrammaticCompositeActivitiesConsole/Workflows/CompositionWorkflow.cs index d986f3867..f23dc0c21 100644 --- a/src/samples/console/Elsa.Samples.ProgrammaticCompositeActivitiesConsole/Workflows/CompositionWorkflow.cs +++ b/src/samples/console/Elsa.Samples.ProgrammaticCompositeActivitiesConsole/Workflows/CompositionWorkflow.cs @@ -13,10 +13,8 @@ namespace Elsa.Samples.ProgrammaticCompositeActivitiesConsole.Workflows .WriteLine("Welcome to the Composite Activities demo workflow!") // A custom, composite activity - .Then(countDown => - { - countDown.When("Left").WriteLine("We're going left."); - countDown.When("Right").WriteLine("We're going right."); - }); + .Then() + .WriteLine("Done") + ; } } \ No newline at end of file diff --git a/src/samples/worker/Elsa.Samples.Timers/MyContainer1.cs b/src/samples/worker/Elsa.Samples.Timers/Activities/MyContainer1.cs similarity index 97% rename from src/samples/worker/Elsa.Samples.Timers/MyContainer1.cs rename to src/samples/worker/Elsa.Samples.Timers/Activities/MyContainer1.cs index c366ac07a..b52479da9 100644 --- a/src/samples/worker/Elsa.Samples.Timers/MyContainer1.cs +++ b/src/samples/worker/Elsa.Samples.Timers/Activities/MyContainer1.cs @@ -5,7 +5,7 @@ using Elsa.Builders; using Elsa.Services; using NodaTime; -namespace Elsa.Samples.Timers +namespace Elsa.Samples.Timers.Activities { public class MyContainer1 : CompositeActivity { diff --git a/src/samples/worker/Elsa.Samples.Timers/Activities/MyContainer2.cs b/src/samples/worker/Elsa.Samples.Timers/Activities/MyContainer2.cs new file mode 100644 index 000000000..8b9c36d15 --- /dev/null +++ b/src/samples/worker/Elsa.Samples.Timers/Activities/MyContainer2.cs @@ -0,0 +1,41 @@ +using Elsa.Activities.Console; +using Elsa.Activities.ControlFlow; +using Elsa.Activities.Timers; +using Elsa.Builders; +using Elsa.Services; +using NodaTime; + +namespace Elsa.Samples.Timers.Activities +{ + public class MyContainer2 : CompositeActivity + { + public override void Build(ICompositeActivityBuilder activity) + { + activity + .WriteLine("In 2 seconds...") + .Timer(Duration.FromSeconds(2)) + .WriteLine("The time is ripe.") + .Then(fork => fork.WithBranches("C", "D", "E"), fork => + { + fork.When("C") + .Then() + .Then("Join2"); + + fork + .When("D") + .While(true, @while => @while + .Timer(Duration.FromSeconds(1)) + .WriteLine("Timer D went off")) + .Then("Join2"); + + fork + .When("E") + .Timer(Duration.FromSeconds(15)) + .WriteLine("Timer E went off. Exiting fork.") + .Then("Join2"); + }) + .Add(join => join.WithMode(Join.JoinMode.WaitAny)).WithName("Join2").WriteLine("Container 2 Joined!") + ; + } + } +} \ No newline at end of file diff --git a/src/samples/worker/Elsa.Samples.Timers/MyContainer2.cs b/src/samples/worker/Elsa.Samples.Timers/MyContainer2.cs deleted file mode 100644 index 1d6096d07..000000000 --- a/src/samples/worker/Elsa.Samples.Timers/MyContainer2.cs +++ /dev/null @@ -1,35 +0,0 @@ -using Elsa.Activities.Console; -using Elsa.Activities.ControlFlow; -using Elsa.Activities.Timers; -using Elsa.Builders; -using Elsa.Services; -using NodaTime; - -namespace Elsa.Samples.Timers -{ - public class MyContainer2 : CompositeActivity - { - public override void Build(ICompositeActivityBuilder activity) - { - activity - .StartIn(Duration.FromSeconds(5)) - .Then(fork => fork.WithBranches("A", "B"), fork => - { - fork - .When("A") - .While(true, @while => @while - .Timer(Duration.FromSeconds(5)) - .WriteLine("Timer C went off")) - .Then("Join2"); - - fork - .When("B") - .While(true, @while => @while - .Timer(Duration.FromSeconds(5)) - .WriteLine("Timer D went off")) - .Then("Join2"); - }) - .Add(join => join.WithMode(Join.JoinMode.WaitAny)).WithName("Join2").WriteLine("Container 2 Joined!"); - } - } -} \ No newline at end of file diff --git a/src/samples/worker/Elsa.Samples.Timers/Program.cs b/src/samples/worker/Elsa.Samples.Timers/Program.cs index 659b36722..3b0aba506 100644 --- a/src/samples/worker/Elsa.Samples.Timers/Program.cs +++ b/src/samples/worker/Elsa.Samples.Timers/Program.cs @@ -1,5 +1,7 @@ using System.Threading.Tasks; using Elsa.Persistence.YesSql.Extensions; +using Elsa.Samples.Timers.Activities; +using Elsa.Samples.Timers.Workflows; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; using NodaTime; @@ -19,14 +21,15 @@ namespace Elsa.Samples.Timers .AddElsa(options => options.UseYesSqlPersistence() .AddConsoleActivities() .AddQuartzTimerActivities() - .AddWorkflow() + //.AddWorkflow() .AddActivity() .AddActivity() - //.AddWorkflow() + .AddWorkflow() //.AddWorkflow() - //.AddWorkflow(new OneOffWorkflow(SystemClock.Instance.GetCurrentInstant().Plus(Duration.FromSeconds(5)))) + //.AddWorkflow(sp => ActivatorUtilities.CreateInstance(sp, sp.GetRequiredService().GetCurrentInstant().Plus(Duration.FromSeconds(5)), sp.GetRequiredService())) ) - .StartWorkflow(); + //.StartWorkflow() + ; }); } } \ No newline at end of file diff --git a/src/samples/worker/Elsa.Samples.Timers/RecurringTaskWorkflow.cs b/src/samples/worker/Elsa.Samples.Timers/RecurringTaskWorkflow.cs deleted file mode 100644 index 8de52798b..000000000 --- a/src/samples/worker/Elsa.Samples.Timers/RecurringTaskWorkflow.cs +++ /dev/null @@ -1,41 +0,0 @@ -using System; -using System.Threading.Tasks; -using Elsa.Activities.Console; -using Elsa.Activities.ControlFlow; -using Elsa.Activities.Timers; -using Elsa.Builders; -using NodaTime; - -namespace Elsa.Samples.Timers -{ - public class RecurringTaskWorkflow : IWorkflow - { - private readonly IClock _clock; - - public RecurringTaskWorkflow(IClock clock) - { - _clock = clock; - } - - public void Build(IWorkflowBuilder workflow) - { - workflow - .WriteLine("Started") - .Then(fork => fork.WithBranches("A", "B"), fork => - { - fork - .When("A") - .Then() - .Then("Join3"); - - fork - .When("B") - .Then() - .Then("Join3"); - }) - .Add(join => join.WithMode(Join.JoinMode.WaitAny)).WithName("Join3") - .WriteLine("Workflow Joined!") - .WriteLine("Finished"); - } - } -} \ No newline at end of file diff --git a/src/samples/worker/Elsa.Samples.Timers/CancelTimerWorkflow.cs b/src/samples/worker/Elsa.Samples.Timers/Workflows/CancelTimerWorkflow.cs similarity index 91% rename from src/samples/worker/Elsa.Samples.Timers/CancelTimerWorkflow.cs rename to src/samples/worker/Elsa.Samples.Timers/Workflows/CancelTimerWorkflow.cs index 5dd332977..ffe4e96c6 100644 --- a/src/samples/worker/Elsa.Samples.Timers/CancelTimerWorkflow.cs +++ b/src/samples/worker/Elsa.Samples.Timers/Workflows/CancelTimerWorkflow.cs @@ -1,12 +1,11 @@ using System.Collections.Generic; - using Elsa.Activities.Console; using Elsa.Activities.ControlFlow; using Elsa.Activities.Timers; using Elsa.Builders; using NodaTime; -namespace Elsa.Samples.Timers +namespace Elsa.Samples.Timers.Workflows { public class CancelTimerWorkflow : IWorkflow { @@ -14,7 +13,7 @@ namespace Elsa.Samples.Timers { workflow .StartAt(SystemClock.Instance.GetCurrentInstant().Plus(Duration.FromSeconds(5))) - .WriteLine("CancelTimerWorkflow is executed") + .WriteLine("CancelTimerWorkflow is executing") .Then( activity => activity.Set(x => x.Branches, new HashSet(new[] { "Branch 1", "Branch 2" })), fork => diff --git a/src/samples/worker/Elsa.Samples.Timers/CronTaskWorkflow.cs b/src/samples/worker/Elsa.Samples.Timers/Workflows/CronTaskWorkflow.cs similarity index 90% rename from src/samples/worker/Elsa.Samples.Timers/CronTaskWorkflow.cs rename to src/samples/worker/Elsa.Samples.Timers/Workflows/CronTaskWorkflow.cs index 5563d36c6..b3fc8e131 100644 --- a/src/samples/worker/Elsa.Samples.Timers/CronTaskWorkflow.cs +++ b/src/samples/worker/Elsa.Samples.Timers/Workflows/CronTaskWorkflow.cs @@ -3,7 +3,7 @@ using Elsa.Activities.Console; using Elsa.Activities.Timers; using Elsa.Builders; -namespace Elsa.Samples.Timers +namespace Elsa.Samples.Timers.Workflows { public class CronTaskWorkflow : IWorkflow { diff --git a/src/samples/worker/Elsa.Samples.Timers/OneOffWorkflow.cs b/src/samples/worker/Elsa.Samples.Timers/Workflows/OneOffWorkflow.cs similarity index 59% rename from src/samples/worker/Elsa.Samples.Timers/OneOffWorkflow.cs rename to src/samples/worker/Elsa.Samples.Timers/Workflows/OneOffWorkflow.cs index 8e02b2cb9..911591e49 100644 --- a/src/samples/worker/Elsa.Samples.Timers/OneOffWorkflow.cs +++ b/src/samples/worker/Elsa.Samples.Timers/Workflows/OneOffWorkflow.cs @@ -3,7 +3,7 @@ using Elsa.Activities.Timers; using Elsa.Builders; using NodaTime; -namespace Elsa.Samples.Timers +namespace Elsa.Samples.Timers.Workflows { /// /// A workflow that executes only once in the near future. @@ -11,19 +11,21 @@ namespace Elsa.Samples.Timers public class OneOffWorkflow : IWorkflow { private readonly Instant _executeAt; + private readonly IClock _clock; - public OneOffWorkflow(Instant executeAt) + public OneOffWorkflow(Instant executeAt, IClock clock) { _executeAt = executeAt; + _clock = clock; } public void Build(IWorkflowBuilder workflow) { workflow .StartAt(_executeAt) - .WriteLine(context => $"Started at {context.GetService().GetCurrentInstant()}. Next event happens 3 seconds from now.") + .WriteLine(() => $"Started at {_clock.GetCurrentInstant()}. Next event happens 3 seconds from now.") .StartIn(Duration.FromSeconds(3)) - .WriteLine(context => $"Follow-up occurred at {context.GetService().GetCurrentInstant()}."); + .WriteLine(() => $"Follow-up occurred at {_clock.GetCurrentInstant()}."); } } } \ No newline at end of file diff --git a/src/samples/worker/Elsa.Samples.Timers/Workflows/RecurringTaskWorkflow.cs b/src/samples/worker/Elsa.Samples.Timers/Workflows/RecurringTaskWorkflow.cs new file mode 100644 index 000000000..f7b8eb626 --- /dev/null +++ b/src/samples/worker/Elsa.Samples.Timers/Workflows/RecurringTaskWorkflow.cs @@ -0,0 +1,25 @@ +using Elsa.Activities.Console; +using Elsa.Builders; +using Elsa.Samples.Timers.Activities; +using NodaTime; + +namespace Elsa.Samples.Timers.Workflows +{ + public class RecurringTaskWorkflow : IWorkflow + { + private readonly IClock _clock; + + public RecurringTaskWorkflow(IClock clock) + { + _clock = clock; + } + + public void Build(IWorkflowBuilder workflow) + { + workflow + .WriteLine("Started") + .Then() + .WriteLine("Finished"); + } + } +} \ No newline at end of file