From bea9a5819ae6c8b6e77f3081ff35ec2068efa7ea Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Sun, 6 Jan 2019 22:42:59 +0100 Subject: [PATCH] Incremental work on execution logging --- .../Elsa.Abstractions/ActivityDriverBase.cs | 5 ++ .../Exceptions/WorkflowException.cs | 11 +++ src/core/Elsa.Abstractions/IActivityDriver.cs | 5 ++ .../Elsa.Abstractions/IActivityInvoker.cs | 1 + .../Models/ActivityExecutionContext.cs | 5 ++ .../Elsa.Abstractions/Models/ActivityFault.cs | 12 ++++ src/core/Elsa.Abstractions/Models/LogEntry.cs | 16 +++++ src/core/Elsa.Abstractions/Models/Workflow.cs | 2 + .../Models/WorkflowExecutionContext.cs | 15 +++++ src/core/Elsa.Core/ActivityInvoker.cs | 47 +++++++++++-- src/core/Elsa.Core/Elsa.Core.csproj | 1 + .../Extensions/ServiceCollectionExtensions.cs | 12 ++-- src/core/Elsa.Core/Results/HaltResult.cs | 4 +- .../Results/TriggerEndpointsResult.cs | 7 +- .../LocalizedStringConverter.cs | 0 .../Tokenizers/WorkflowTokenizer.cs | 35 ++++++++-- src/core/Elsa.Core/WorkflowInvoker.cs | 67 ++++++++++++++----- .../IWorkflowStore.cs | 1 + .../Extensions/ServiceCollectionExtensions.cs | 2 +- .../FileSystemWorkflowStore.cs | 20 ++++-- .../InMemoryWorkflowStore.cs | 8 +++ src/core/Elsa.Runtime/WorkflowHost.cs | 9 +-- .../Services/FileSystemWorkflowStore.cs | 17 ++++- 23 files changed, 252 insertions(+), 50 deletions(-) create mode 100644 src/core/Elsa.Abstractions/Exceptions/WorkflowException.cs create mode 100644 src/core/Elsa.Abstractions/Models/ActivityFault.cs create mode 100644 src/core/Elsa.Abstractions/Models/LogEntry.cs rename src/core/Elsa.Core/Serialization/{Json => Converters}/LocalizedStringConverter.cs (100%) diff --git a/src/core/Elsa.Abstractions/ActivityDriverBase.cs b/src/core/Elsa.Abstractions/ActivityDriverBase.cs index fd16b1a39..c03707146 100644 --- a/src/core/Elsa.Abstractions/ActivityDriverBase.cs +++ b/src/core/Elsa.Abstractions/ActivityDriverBase.cs @@ -10,17 +10,22 @@ namespace Elsa public string ActivityType => typeof(T).Name; public Task CanExecuteAsync(ActivityExecutionContext activityContext, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) => OnCanExecuteAsync(activityContext, workflowContext, cancellationToken); public Task ExecuteAsync(ActivityExecutionContext activityContext, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) => OnExecuteAsync(activityContext, workflowContext, cancellationToken); + public Task HaltedAsync(ActivityExecutionContext activityContext, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) => OnHaltedAsync(activityContext, workflowContext, cancellationToken); public Task ResumeAsync(ActivityExecutionContext activityContext, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) => OnResumeAsync(activityContext, workflowContext, cancellationToken); protected virtual Task OnCanExecuteAsync(ActivityExecutionContext activityContext, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) => OnCanExecuteAsync((T) activityContext.Activity, workflowContext, cancellationToken); protected virtual Task OnCanExecuteAsync(T activity, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) => Task.FromResult(OnCanExecute(activity, workflowContext)); protected virtual Task OnExecuteAsync(ActivityExecutionContext activityContext, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) => OnExecuteAsync((T) activityContext.Activity, workflowContext, cancellationToken); protected virtual Task OnExecuteAsync(T activity, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) => Task.FromResult(OnExecute(activity, workflowContext)); + protected virtual Task OnHaltedAsync(ActivityExecutionContext activityContext, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) => OnHaltedAsync((T) activityContext.Activity, workflowContext, cancellationToken); + protected virtual Task OnHaltedAsync(T activity, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) => Task.FromResult(OnHalted(activity, workflowContext)); protected virtual Task OnResumeAsync(ActivityExecutionContext activityContext, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) => Task.FromResult(OnResume(activityContext, workflowContext)); protected virtual Task OnResumeAsync(T activity, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) => Task.FromResult(OnResume(activity, workflowContext)); protected virtual bool OnCanExecute(ActivityExecutionContext activityContext, WorkflowExecutionContext workflowContext) => OnCanExecute((T) activityContext.Activity, workflowContext); protected virtual bool OnCanExecute(T activity, WorkflowExecutionContext workflowContext) => true; protected virtual ActivityExecutionResult OnExecute(ActivityExecutionContext activityContext, WorkflowExecutionContext workflowContext) => OnExecute((T) activityContext.Activity, workflowContext); protected virtual ActivityExecutionResult OnExecute(T activity, WorkflowExecutionContext workflowContext) => Noop(); + protected virtual ActivityExecutionResult OnHalted(ActivityExecutionContext activityContext, WorkflowExecutionContext workflowContext) => OnHalted((T) activityContext.Activity, workflowContext); + protected virtual ActivityExecutionResult OnHalted(T activity, WorkflowExecutionContext workflowContext) => Noop(); protected virtual ActivityExecutionResult OnResume(ActivityExecutionContext activityContext, WorkflowExecutionContext workflowContext) => OnResume((T) activityContext.Activity, workflowContext); protected virtual ActivityExecutionResult OnResume(T activity, WorkflowExecutionContext workflowContext) => Noop(); protected NoopResult Noop() => new NoopResult(); diff --git a/src/core/Elsa.Abstractions/Exceptions/WorkflowException.cs b/src/core/Elsa.Abstractions/Exceptions/WorkflowException.cs new file mode 100644 index 000000000..377f06f49 --- /dev/null +++ b/src/core/Elsa.Abstractions/Exceptions/WorkflowException.cs @@ -0,0 +1,11 @@ +using System; + +namespace Elsa.Exceptions +{ + public class WorkflowException : Exception + { + public WorkflowException(string message) : base(message) + { + } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/IActivityDriver.cs b/src/core/Elsa.Abstractions/IActivityDriver.cs index 477f17997..1654573c6 100644 --- a/src/core/Elsa.Abstractions/IActivityDriver.cs +++ b/src/core/Elsa.Abstractions/IActivityDriver.cs @@ -26,5 +26,10 @@ namespace Elsa /// Resumes the specified activity. /// Task ResumeAsync(ActivityExecutionContext activityContext, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken); + + /// + /// Invoked when the workflow is halted. + /// + Task HaltedAsync(ActivityExecutionContext activityContext, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken); } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/IActivityInvoker.cs b/src/core/Elsa.Abstractions/IActivityInvoker.cs index 071eb76e3..6e494dca0 100644 --- a/src/core/Elsa.Abstractions/IActivityInvoker.cs +++ b/src/core/Elsa.Abstractions/IActivityInvoker.cs @@ -9,5 +9,6 @@ namespace Elsa { Task ExecuteAsync(WorkflowExecutionContext workflowContext, IActivity activity, CancellationToken cancellationToken = default); Task ResumeAsync(WorkflowExecutionContext workflowContext, IActivity activity, CancellationToken cancellationToken = default); + Task HaltedAsync(WorkflowExecutionContext workflowContext, IActivity activity, CancellationToken cancellationToken = default); } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Models/ActivityExecutionContext.cs b/src/core/Elsa.Abstractions/Models/ActivityExecutionContext.cs index 98aac6fd7..c82d6f786 100644 --- a/src/core/Elsa.Abstractions/Models/ActivityExecutionContext.cs +++ b/src/core/Elsa.Abstractions/Models/ActivityExecutionContext.cs @@ -5,8 +5,13 @@ public ActivityExecutionContext(IActivity activity) { Activity = activity; + LogEntry = new LogEntry + { + ActivityId = activity.Id + }; } public IActivity Activity { get; } + public LogEntry LogEntry { get; set; } } } diff --git a/src/core/Elsa.Abstractions/Models/ActivityFault.cs b/src/core/Elsa.Abstractions/Models/ActivityFault.cs new file mode 100644 index 000000000..ebce4ecaa --- /dev/null +++ b/src/core/Elsa.Abstractions/Models/ActivityFault.cs @@ -0,0 +1,12 @@ +namespace Elsa.Models +{ + public class ActivityFault + { + public ActivityFault(string message) + { + Message = message; + } + + public string Message { get; set; } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Models/LogEntry.cs b/src/core/Elsa.Abstractions/Models/LogEntry.cs new file mode 100644 index 000000000..924e4dcfc --- /dev/null +++ b/src/core/Elsa.Abstractions/Models/LogEntry.cs @@ -0,0 +1,16 @@ +using System.Collections.Generic; +using NodaTime; + +namespace Elsa.Models +{ + public class LogEntry + { + public string ActivityId { get; set; } + public Instant? ExecutedAt { get; set; } + public Instant? HaltedAt { get; set; } + public Instant? ResumedAt { get; set; } + public Instant? CompletedAt { get; set; } + public IList TriggeredEndpoints { get; set; } + public ActivityFault Fault { get; set; } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Models/Workflow.cs b/src/core/Elsa.Abstractions/Models/Workflow.cs index 40c528cc1..c5b8afe17 100644 --- a/src/core/Elsa.Abstractions/Models/Workflow.cs +++ b/src/core/Elsa.Abstractions/Models/Workflow.cs @@ -20,6 +20,7 @@ namespace Elsa.Models Scopes = new Stack(new[] { CurrentScope }); Arguments = new Variables(); BlockingActivities = new List(); + ExecutionLog = new List(); Metadata = new WorkflowMetadata(); } @@ -34,6 +35,7 @@ namespace Elsa.Models public WorkflowExecutionScope CurrentScope { get; set; } public Variables Arguments { get; set; } public IList BlockingActivities { get; set; } + public IList ExecutionLog { get; set; } public WorkflowMetadata Metadata { get; set; } public WorkflowFault Fault { get; set; } } diff --git a/src/core/Elsa.Abstractions/Models/WorkflowExecutionContext.cs b/src/core/Elsa.Abstractions/Models/WorkflowExecutionContext.cs index dfe633950..1bdeae203 100644 --- a/src/core/Elsa.Abstractions/Models/WorkflowExecutionContext.cs +++ b/src/core/Elsa.Abstractions/Models/WorkflowExecutionContext.cs @@ -1,6 +1,7 @@ using System; using System.Collections.Generic; using System.Linq; +using System.Threading.Tasks; using NodaTime; namespace Elsa.Models @@ -8,18 +9,22 @@ namespace Elsa.Models public class WorkflowExecutionContext { private readonly Stack scheduledActivities; + private readonly Stack scheduledHaltingActivities; public WorkflowExecutionContext(Workflow workflow) { Workflow = workflow; IsFirstPass = true; scheduledActivities = new Stack(); + scheduledHaltingActivities = new Stack(); } public Workflow Workflow { get; } public bool HasScheduledActivities => scheduledActivities.Any(); + public bool HasScheduledHaltingActivities => scheduledHaltingActivities.Any(); public bool IsFirstPass { get; set; } public IActivity CurrentActivity { get; private set; } + public LogEntry CurrentLogEntry => Workflow.ExecutionLog.LastOrDefault(); public WorkflowExecutionScope CurrentScope { @@ -51,6 +56,16 @@ namespace Elsa.Models CurrentActivity = scheduledActivities.Pop(); return CurrentActivity; } + + public void ScheduleHaltingActivity(IActivity activity) + { + scheduledHaltingActivities.Push(activity); + } + + public IActivity PopScheduledHaltingActivity() + { + return scheduledHaltingActivities.Pop(); + } public void SetLastResult(object value) { diff --git a/src/core/Elsa.Core/ActivityInvoker.cs b/src/core/Elsa.Core/ActivityInvoker.cs index 097a6ad82..ecb0cb64d 100644 --- a/src/core/Elsa.Core/ActivityInvoker.cs +++ b/src/core/Elsa.Core/ActivityInvoker.cs @@ -3,39 +3,76 @@ using System.Collections.Generic; using System.Linq; using System.Threading; using System.Threading.Tasks; +using Elsa.Exceptions; using Elsa.Models; using Elsa.Results; using Microsoft.Extensions.Logging; +using NodaTime; namespace Elsa { public class ActivityInvoker : IActivityInvoker { private readonly IActivityDriverRegistry driverRegistry; + private readonly IClock clock; + private readonly ILogger logger; - public ActivityInvoker(IActivityDriverRegistry driverRegistry) + public ActivityInvoker(IActivityDriverRegistry driverRegistry, IClock clock, ILogger logger) { this.driverRegistry = driverRegistry; + this.clock = clock; + this.logger = logger; } public async Task ExecuteAsync(WorkflowExecutionContext workflowContext, IActivity activity, CancellationToken cancellationToken = default) { - return await InvokeAsync(workflowContext, activity, (context, driver) => driver.ExecuteAsync(context, workflowContext, cancellationToken)); + return await InvokeAsync(workflowContext, activity, (context, driver) => + { + workflowContext.Workflow.ExecutionLog.Add(context.LogEntry); + context.LogEntry.ExecutedAt = clock.GetCurrentInstant(); + return driver.ExecuteAsync(context, workflowContext, cancellationToken); + }); } public async Task ResumeAsync(WorkflowExecutionContext workflowContext, IActivity activity, CancellationToken cancellationToken = default) { - return await InvokeAsync(workflowContext, activity, (context, driver) => driver.ResumeAsync(context, workflowContext, cancellationToken)); + return await InvokeAsync(workflowContext, activity, (context, driver) => + { + context.LogEntry.ResumedAt = clock.GetCurrentInstant(); + return driver.ResumeAsync(context, workflowContext, cancellationToken); + }); } - private Task InvokeAsync( + public async Task HaltedAsync(WorkflowExecutionContext workflowContext, IActivity activity, CancellationToken cancellationToken = default) + { + return await InvokeAsync(workflowContext, activity, (context, driver) => + { + context.LogEntry.HaltedAt = clock.GetCurrentInstant(); + return driver.HaltedAsync(context, workflowContext, cancellationToken); + }); + } + + private async Task InvokeAsync( WorkflowExecutionContext workflowContext, IActivity activity, Func> invokeAction) { var activityContext = workflowContext.CreateActivityExecutionContext(activity); var driver = driverRegistry.GetDriver(activity.Name); - return invokeAction(activityContext, driver); + + try + { + if (driver == null) + throw new WorkflowException($"No driver found for activity {activity.Name}"); + + return await invokeAction(activityContext, driver); + } + catch (Exception e) + { + logger.LogError(e, "Error while invoking activity {ActivityId} of workflow {WorkflowId}", activity.Id, workflowContext.Workflow.Metadata.Id); + activityContext.LogEntry.Fault = new ActivityFault(e.Message); + return new FaultWorkflowResult(activityContext.LogEntry.Fault.Message, clock.GetCurrentInstant()); + } } } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Elsa.Core.csproj b/src/core/Elsa.Core/Elsa.Core.csproj index 393d1c520..129106231 100644 --- a/src/core/Elsa.Core/Elsa.Core.csproj +++ b/src/core/Elsa.Core/Elsa.Core.csproj @@ -29,6 +29,7 @@ + \ 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 00753ad5b..0e082079c 100644 --- a/src/core/Elsa.Core/Extensions/ServiceCollectionExtensions.cs +++ b/src/core/Elsa.Core/Extensions/ServiceCollectionExtensions.cs @@ -15,16 +15,16 @@ namespace Elsa.Extensions { services.TryAddSingleton(); services.TryAddSingleton(SystemClock.Instance); - services.TryAddScoped(); - services.TryAddScoped(); - services.AddScoped(); - services.TryAddScoped(); + services.TryAddSingleton(); + services.TryAddSingleton(); + services.TryAddSingleton(); + services.TryAddSingleton(); + services.TryAddSingleton(); + services.TryAddSingleton(); services.AddActivityDescriptors(); services.AddSingleton(); services.AddSingleton(); services.AddSingleton(); - services.TryAddSingleton(); - services.TryAddSingleton(); services.AddSingleton(); services.AddSingleton(); services.AddSingleton(); diff --git a/src/core/Elsa.Core/Results/HaltResult.cs b/src/core/Elsa.Core/Results/HaltResult.cs index d8ebb990e..0a65121f2 100644 --- a/src/core/Elsa.Core/Results/HaltResult.cs +++ b/src/core/Elsa.Core/Results/HaltResult.cs @@ -19,9 +19,10 @@ namespace Elsa.Results public override async Task ExecuteAsync(IWorkflowInvoker invoker, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) { + var activity = workflowContext.CurrentActivity; + if (workflowContext.IsFirstPass) { - var activity = workflowContext.CurrentActivity; var result = await invoker.ActivityInvoker.ResumeAsync(workflowContext, activity, cancellationToken); workflowContext.IsFirstPass = false; @@ -29,6 +30,7 @@ namespace Elsa.Results } else { + workflowContext.ScheduleHaltingActivity(activity); workflowContext.Halt(instant); } } diff --git a/src/core/Elsa.Core/Results/TriggerEndpointsResult.cs b/src/core/Elsa.Core/Results/TriggerEndpointsResult.cs index 58a6cba2b..a9142bf5f 100644 --- a/src/core/Elsa.Core/Results/TriggerEndpointsResult.cs +++ b/src/core/Elsa.Core/Results/TriggerEndpointsResult.cs @@ -1,4 +1,5 @@ using System.Collections.Generic; +using System.Linq; using Elsa.Models; namespace Elsa.Results @@ -10,15 +11,17 @@ namespace Elsa.Results { public TriggerEndpointsResult(IEnumerable endpointNames) { - EndpointNames = endpointNames; + EndpointNames = endpointNames.ToList(); } - public IEnumerable EndpointNames { get; } + public IReadOnlyList EndpointNames { get; } protected override void Execute(IWorkflowInvoker invoker, WorkflowExecutionContext workflowContext) { var currentActivity = workflowContext.CurrentActivity; + workflowContext.CurrentLogEntry.TriggeredEndpoints = EndpointNames.ToList(); + foreach (var endpointName in EndpointNames) { workflowContext.ScheduleNextActivities(workflowContext, new SourceEndpoint(currentActivity, endpointName)); diff --git a/src/core/Elsa.Core/Serialization/Json/LocalizedStringConverter.cs b/src/core/Elsa.Core/Serialization/Converters/LocalizedStringConverter.cs similarity index 100% rename from src/core/Elsa.Core/Serialization/Json/LocalizedStringConverter.cs rename to src/core/Elsa.Core/Serialization/Converters/LocalizedStringConverter.cs diff --git a/src/core/Elsa.Core/Serialization/Tokenizers/WorkflowTokenizer.cs b/src/core/Elsa.Core/Serialization/Tokenizers/WorkflowTokenizer.cs index bba8232a0..082a9cdf1 100644 --- a/src/core/Elsa.Core/Serialization/Tokenizers/WorkflowTokenizer.cs +++ b/src/core/Elsa.Core/Serialization/Tokenizers/WorkflowTokenizer.cs @@ -60,7 +60,8 @@ namespace Elsa.Serialization.Tokenizers { "haltedAt", SerializeInstant(workflow.HaltedAt) }, { "activities", await SerializeActivitiesAsync(context, cancellationToken) }, { "connections", SerializeConnections(context, workflow) }, - { "haltedActivities", SerializeHaltedActivities(context, workflow) }, + { "blockingActivities", SerializeBlockingActivities(context, workflow) }, + { "executionLog", SerializeExecutionLog(workflow) }, { "scopes", SerializeScopes(context, scopeIdLookup) }, { "currentScope", SerializeCurrentScope(workflow, scopeIdLookup) }, }; @@ -88,7 +89,8 @@ namespace Elsa.Serialization.Tokenizers HaltedAt = DeserializeInstant(token["haltedAt"]), Activities = activityDictionary.Values.ToList(), Connections = DeserializeConnections(token, activityDictionary).ToList(), - BlockingActivities = DeserializeHaltedActivities(token, activityDictionary).ToList(), + BlockingActivities = DeserializeBlockingActivities(token, activityDictionary).ToList(), + ExecutionLog = DeserializeExecutionLog(token).ToList(), Scopes = new Stack(scopeLookup.Values) }; @@ -160,7 +162,7 @@ namespace Elsa.Serialization.Tokenizers return scopeModels; } - private JArray SerializeHaltedActivities(WorkflowTokenizationContext context, Workflow workflow) + private JArray SerializeBlockingActivities(WorkflowTokenizationContext context, Workflow workflow) { var haltedActivityModels = new JArray(); @@ -206,6 +208,18 @@ namespace Elsa.Serialization.Tokenizers return activityModels; } + + private JArray SerializeExecutionLog(Workflow workflow) + { + var models = new JArray(); + + foreach (var entry in workflow.ExecutionLog) + { + models.Add(JToken.FromObject(entry, jsonSerializer)); + } + + return models; + } private IDictionary DeserializeScopes(JToken token, WorkflowTokenizationContext context) { @@ -236,9 +250,9 @@ namespace Elsa.Serialization.Tokenizers return scopeLookup; } - private IEnumerable DeserializeHaltedActivities(JToken token, IDictionary activityDictionary) + private IEnumerable DeserializeBlockingActivities(JToken token, IDictionary activityDictionary) { - var haltedActivityModels = (JArray) token["haltedActivities"] ?? new JArray(); + var haltedActivityModels = (JArray) token["blockingActivities"] ?? new JArray(); foreach (var haltedActivityModel in haltedActivityModels) { @@ -247,6 +261,17 @@ namespace Elsa.Serialization.Tokenizers yield return activity; } } + + private IEnumerable DeserializeExecutionLog(JToken token) + { + var models = (JArray) token["executionLog"] ?? new JArray(); + + foreach (var model in models) + { + var entry = JsonConvert.DeserializeObject(model.ToString(), new JsonSerializerSettings().ConfigureForNodaTime(DateTimeZoneProviders.Tzdb)); + yield return entry; + } + } private IEnumerable DeserializeConnections(JToken token, IDictionary activityDictionary) { diff --git a/src/core/Elsa.Core/WorkflowInvoker.cs b/src/core/Elsa.Core/WorkflowInvoker.cs index 08a3cef08..57c4b9c9c 100644 --- a/src/core/Elsa.Core/WorkflowInvoker.cs +++ b/src/core/Elsa.Core/WorkflowInvoker.cs @@ -3,7 +3,9 @@ using System.Collections.Generic; using System.Linq; using System.Threading; using System.Threading.Tasks; +using Elsa.Extensions; using Elsa.Models; +using Elsa.Persistence; using Elsa.Results; using Microsoft.Extensions.Logging; using NodaTime; @@ -12,15 +14,18 @@ namespace Elsa { public class WorkflowInvoker : IWorkflowInvoker { + private readonly IWorkflowStore workflowStore; private readonly IClock clock; private readonly ILogger logger; - + public WorkflowInvoker( IActivityInvoker activityInvoker, + IWorkflowStore workflowStore, IClock clock, ILogger logger) { ActivityInvoker = activityInvoker; + this.workflowStore = workflowStore; this.clock = clock; this.logger = logger; } @@ -33,34 +38,56 @@ namespace Elsa var workflowExecutionContext = new WorkflowExecutionContext(workflow); var isResuming = workflowExecutionContext.Workflow.Status == WorkflowStatus.Resuming; + // If a start activity was provided, remove it from the blocking activities list. If not start activity was provided, pick the first one that has no inbound connections. if (startActivity != null) workflow.BlockingActivities.Remove(startActivity); else - startActivity = workflow.Activities.First(); + startActivity = workflow.GetStartActivities().FirstOrDefault(); if (!isResuming) workflow.StartedAt = clock.GetCurrentInstant(); - + workflowExecutionContext.Workflow.Status = WorkflowStatus.Executing; - workflowExecutionContext.ScheduleActivity(startActivity); - + + if (startActivity != null) + workflowExecutionContext.ScheduleActivity(startActivity); + + // Keep executing activities as long as there are any scheduled. while (workflowExecutionContext.HasScheduledActivities) { var currentActivity = workflowExecutionContext.PopScheduledActivity(); var result = await ExecuteActivityAsync(workflowExecutionContext, currentActivity, isResuming, cancellationToken); - - if(result == null) + + if (result == null) break; - + await result.ExecuteAsync(this, workflowExecutionContext, cancellationToken); workflowExecutionContext.IsFirstPass = false; isResuming = false; } - if(workflowExecutionContext.Workflow.Status != WorkflowStatus.Halted) + // Any other status than Halted means the workflow has ended (either because it reached the final activity, was aborted or has faulted). + if (workflowExecutionContext.Workflow.Status != WorkflowStatus.Halted) + { workflowExecutionContext.Finish(clock.GetCurrentInstant()); - + } + else + { + // Persist workflow before executing the halted activities. + await workflowStore.SaveAsync(workflow, cancellationToken); + + // Invoke Halted event on activity drivers that halted the workflow. + while (workflowExecutionContext.HasScheduledHaltingActivities) + { + var currentActivity = workflowExecutionContext.PopScheduledHaltingActivity(); + var result = await ExecuteActivityHaltedAsync(workflowExecutionContext, currentActivity, cancellationToken); + + await result.ExecuteAsync(this, workflowExecutionContext, cancellationToken); + } + } + + await workflowStore.SaveAsync(workflow, cancellationToken); return workflowExecutionContext; } @@ -71,17 +98,27 @@ namespace Elsa } private async Task ExecuteActivityAsync(WorkflowExecutionContext workflowContext, IActivity activity, bool isResuming, CancellationToken cancellationToken) + { + return await ExecuteActivityAsync(workflowContext, activity, () => ExecuteOrResumeActivityAsync(workflowContext, activity, isResuming, cancellationToken), cancellationToken); + } + + private async Task ExecuteActivityHaltedAsync(WorkflowExecutionContext workflowContext, IActivity activity, CancellationToken cancellationToken) + { + return await ExecuteActivityAsync(workflowContext, activity, () => ActivityInvoker.HaltedAsync(workflowContext, activity, cancellationToken), cancellationToken); + } + + private async Task ExecuteActivityAsync(WorkflowExecutionContext workflowContext, IActivity activity, Func> executeAction, CancellationToken cancellationToken) { try { if (cancellationToken.IsCancellationRequested) { workflowContext.Workflow.Status = WorkflowStatus.Aborted; + workflowContext.Workflow.FinishedAt = clock.GetCurrentInstant(); return null; } - - return await ExecuteOrResumeActivityAsync(workflowContext, activity, isResuming, cancellationToken); - + + return await executeAction(); } catch (Exception ex) { @@ -94,7 +131,7 @@ namespace Elsa private void FaultWorkflow(WorkflowExecutionContext workflowContext, IActivity activity, Exception ex) { logger.LogError( - ex, + ex, "An unhandled error occurred while executing an activity. Putting the workflow in the faulted state." ); workflowContext.Fault(ex, activity, clock.GetCurrentInstant()); @@ -107,4 +144,4 @@ namespace Elsa : await ActivityInvoker.ExecuteAsync(workflowContext, activity, cancellationToken); } } -} +} \ No newline at end of file diff --git a/src/core/Elsa.Persistence.Abstractions/IWorkflowStore.cs b/src/core/Elsa.Persistence.Abstractions/IWorkflowStore.cs index b0449b5dd..b71ac4957 100644 --- a/src/core/Elsa.Persistence.Abstractions/IWorkflowStore.cs +++ b/src/core/Elsa.Persistence.Abstractions/IWorkflowStore.cs @@ -12,6 +12,7 @@ namespace Elsa.Persistence Task GetAsync(ISpecification specification, CancellationToken cancellationToken); Task AddAsync(Workflow value, CancellationToken cancellationToken); Task UpdateAsync(Workflow value, CancellationToken cancellationToken); + Task SaveAsync(Workflow value, CancellationToken cancellationToken); } public static class WorkflowStoreExtensions diff --git a/src/core/Elsa.Persistence.FileSystem/Extensions/ServiceCollectionExtensions.cs b/src/core/Elsa.Persistence.FileSystem/Extensions/ServiceCollectionExtensions.cs index ac55aa775..119233105 100644 --- a/src/core/Elsa.Persistence.FileSystem/Extensions/ServiceCollectionExtensions.cs +++ b/src/core/Elsa.Persistence.FileSystem/Extensions/ServiceCollectionExtensions.cs @@ -11,7 +11,7 @@ namespace Elsa.Persistence.FileSystem.Extensions public static IServiceCollection AddWorkflowsFileSystemPersistence(this IServiceCollection services, IConfiguration configuration) { services.TryAddSingleton(); - services.AddScoped(); + services.AddSingleton(); services.Configure(configuration); return services; diff --git a/src/core/Elsa.Persistence.FileSystem/FileSystemWorkflowStore.cs b/src/core/Elsa.Persistence.FileSystem/FileSystemWorkflowStore.cs index 95c1ab793..3eec37cdf 100644 --- a/src/core/Elsa.Persistence.FileSystem/FileSystemWorkflowStore.cs +++ b/src/core/Elsa.Persistence.FileSystem/FileSystemWorkflowStore.cs @@ -21,9 +21,9 @@ namespace Elsa.Persistence.FileSystem private readonly string format; public FileSystemWorkflowStore( - IOptions options, - IFileSystem fileSystem, - IIdGenerator idGenerator, + IOptions options, + IFileSystem fileSystem, + IIdGenerator idGenerator, IWorkflowSerializer workflowSerializer, IClock clock) { @@ -49,10 +49,20 @@ namespace Elsa.Persistence.FileSystem return query.Distinct().FirstOrDefault(); } + public async Task SaveAsync(Workflow value, CancellationToken cancellationToken) + { + if(!fileSystem.Path.HasExtension(value.Metadata.Id)) + await AddAsync(value, cancellationToken); + else + await UpdateAsync(value, cancellationToken); + } + public async Task AddAsync(Workflow value, CancellationToken cancellationToken) { - var fileExtension = $".{format.ToLower()}"; - var id = $"{idGenerator.Generate()}{fileExtension}"; + var id = string.IsNullOrWhiteSpace(value.Metadata.Id) ? idGenerator.Generate() : value.Metadata.Id; + + if (!fileSystem.Path.HasExtension(id)) + id = $"{id}.{format.ToLower()}"; value.Metadata.Id = id; value.CreatedAt = clock.GetCurrentInstant(); diff --git a/src/core/Elsa.Persistence.InMemory/InMemoryWorkflowStore.cs b/src/core/Elsa.Persistence.InMemory/InMemoryWorkflowStore.cs index a135361db..faf2daf58 100644 --- a/src/core/Elsa.Persistence.InMemory/InMemoryWorkflowStore.cs +++ b/src/core/Elsa.Persistence.InMemory/InMemoryWorkflowStore.cs @@ -48,5 +48,13 @@ namespace Elsa.Persistence.InMemory { return Task.CompletedTask; } + + public async Task SaveAsync(Workflow value, CancellationToken cancellationToken) + { + if (string.IsNullOrWhiteSpace(value.Metadata.Id)) + await AddAsync(value, cancellationToken); + else + await UpdateAsync(value, cancellationToken); + } } } \ No newline at end of file diff --git a/src/core/Elsa.Runtime/WorkflowHost.cs b/src/core/Elsa.Runtime/WorkflowHost.cs index 9b1ecfe2d..5cc0c165e 100644 --- a/src/core/Elsa.Runtime/WorkflowHost.cs +++ b/src/core/Elsa.Runtime/WorkflowHost.cs @@ -31,17 +31,12 @@ namespace Elsa.Runtime public async Task StartWorkflowAsync(Workflow workflow, IActivity startActivity, Variables arguments, CancellationToken cancellationToken) { var workflowInstance = await workflowSerializer.DeriveAsync(workflow, cancellationToken); - var workflowContext = await invoker.InvokeAsync(workflowInstance, startActivity, arguments, cancellationToken); - - await workflowStore.AddAsync(workflowInstance, cancellationToken); - return workflowContext; + return await invoker.InvokeAsync(workflowInstance, startActivity, arguments, cancellationToken); } public async Task ResumeWorkflowAsync(Workflow workflow, IActivity activity, Variables arguments, CancellationToken cancellationToken) { - var workflowContext = await invoker.ResumeAsync(workflow, activity, arguments, cancellationToken); - await workflowStore.UpdateAsync(workflow, cancellationToken); - return workflowContext; + return await invoker.ResumeAsync(workflow, activity, arguments, cancellationToken); } private async Task StartNewWorkflowsAsync(string activityName, Variables arguments, CancellationToken cancellationToken) diff --git a/src/web/modules/Elsa.Web.Persistence.FileSystem/Services/FileSystemWorkflowStore.cs b/src/web/modules/Elsa.Web.Persistence.FileSystem/Services/FileSystemWorkflowStore.cs index d4de642c7..bdbbf4269 100644 --- a/src/web/modules/Elsa.Web.Persistence.FileSystem/Services/FileSystemWorkflowStore.cs +++ b/src/web/modules/Elsa.Web.Persistence.FileSystem/Services/FileSystemWorkflowStore.cs @@ -1,4 +1,5 @@ using System.Collections.Generic; +using System.IO; using System.Linq; using System.Threading; using System.Threading.Tasks; @@ -49,11 +50,13 @@ namespace Elsa.Web.Persistence.FileSystem.Services public async Task AddAsync(Workflow value, CancellationToken cancellationToken) { - var fileExtension = $".{Format.ToLower()}"; - var id = $"{idGenerator.Generate()}{fileExtension}"; + var id = string.IsNullOrWhiteSpace(value.Metadata.Id) ? idGenerator.Generate() : value.Metadata.Id; + + if (!Path.HasExtension(id)) + id = $"{id}.{Format.ToLower()}"; - value.CreatedAt = clock.GetCurrentInstant(); value.Metadata.Id = id; + value.CreatedAt = clock.GetCurrentInstant(); await UpdateAsync(value, cancellationToken); } @@ -69,6 +72,14 @@ namespace Elsa.Web.Persistence.FileSystem.Services await fileStore.CreateFileFromStream(path, stream, true); } } + + public async Task SaveAsync(Workflow value, CancellationToken cancellationToken) + { + if(!Path.HasExtension(value.Metadata.Id)) + await AddAsync(value, cancellationToken); + else + await UpdateAsync(value, cancellationToken); + } private async Task> ListAsync(CancellationToken cancellationToken) {