Incremental work on execution logging
This commit is contained in:
parent
0af3ff541b
commit
bea9a5819a
|
|
@ -10,17 +10,22 @@ namespace Elsa
|
|||
public string ActivityType => typeof(T).Name;
|
||||
public Task<bool> CanExecuteAsync(ActivityExecutionContext activityContext, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) => OnCanExecuteAsync(activityContext, workflowContext, cancellationToken);
|
||||
public Task<ActivityExecutionResult> ExecuteAsync(ActivityExecutionContext activityContext, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) => OnExecuteAsync(activityContext, workflowContext, cancellationToken);
|
||||
public Task<ActivityExecutionResult> HaltedAsync(ActivityExecutionContext activityContext, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) => OnHaltedAsync(activityContext, workflowContext, cancellationToken);
|
||||
public Task<ActivityExecutionResult> ResumeAsync(ActivityExecutionContext activityContext, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) => OnResumeAsync(activityContext, workflowContext, cancellationToken);
|
||||
protected virtual Task<bool> OnCanExecuteAsync(ActivityExecutionContext activityContext, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) => OnCanExecuteAsync((T) activityContext.Activity, workflowContext, cancellationToken);
|
||||
protected virtual Task<bool> OnCanExecuteAsync(T activity, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) => Task.FromResult(OnCanExecute(activity, workflowContext));
|
||||
protected virtual Task<ActivityExecutionResult> OnExecuteAsync(ActivityExecutionContext activityContext, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) => OnExecuteAsync((T) activityContext.Activity, workflowContext, cancellationToken);
|
||||
protected virtual Task<ActivityExecutionResult> OnExecuteAsync(T activity, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) => Task.FromResult(OnExecute(activity, workflowContext));
|
||||
protected virtual Task<ActivityExecutionResult> OnHaltedAsync(ActivityExecutionContext activityContext, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) => OnHaltedAsync((T) activityContext.Activity, workflowContext, cancellationToken);
|
||||
protected virtual Task<ActivityExecutionResult> OnHaltedAsync(T activity, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) => Task.FromResult(OnHalted(activity, workflowContext));
|
||||
protected virtual Task<ActivityExecutionResult> OnResumeAsync(ActivityExecutionContext activityContext, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) => Task.FromResult(OnResume(activityContext, workflowContext));
|
||||
protected virtual Task<ActivityExecutionResult> 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();
|
||||
|
|
|
|||
11
src/core/Elsa.Abstractions/Exceptions/WorkflowException.cs
Normal file
11
src/core/Elsa.Abstractions/Exceptions/WorkflowException.cs
Normal file
|
|
@ -0,0 +1,11 @@
|
|||
using System;
|
||||
|
||||
namespace Elsa.Exceptions
|
||||
{
|
||||
public class WorkflowException : Exception
|
||||
{
|
||||
public WorkflowException(string message) : base(message)
|
||||
{
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -26,5 +26,10 @@ namespace Elsa
|
|||
/// Resumes the specified activity.
|
||||
/// </summary>
|
||||
Task<ActivityExecutionResult> ResumeAsync(ActivityExecutionContext activityContext, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken);
|
||||
|
||||
/// <summary>
|
||||
/// Invoked when the workflow is halted.
|
||||
/// </summary>
|
||||
Task<ActivityExecutionResult> HaltedAsync(ActivityExecutionContext activityContext, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken);
|
||||
}
|
||||
}
|
||||
|
|
@ -9,5 +9,6 @@ namespace Elsa
|
|||
{
|
||||
Task<ActivityExecutionResult> ExecuteAsync(WorkflowExecutionContext workflowContext, IActivity activity, CancellationToken cancellationToken = default);
|
||||
Task<ActivityExecutionResult> ResumeAsync(WorkflowExecutionContext workflowContext, IActivity activity, CancellationToken cancellationToken = default);
|
||||
Task<ActivityExecutionResult> HaltedAsync(WorkflowExecutionContext workflowContext, IActivity activity, CancellationToken cancellationToken = default);
|
||||
}
|
||||
}
|
||||
|
|
@ -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; }
|
||||
}
|
||||
}
|
||||
|
|
|
|||
12
src/core/Elsa.Abstractions/Models/ActivityFault.cs
Normal file
12
src/core/Elsa.Abstractions/Models/ActivityFault.cs
Normal file
|
|
@ -0,0 +1,12 @@
|
|||
namespace Elsa.Models
|
||||
{
|
||||
public class ActivityFault
|
||||
{
|
||||
public ActivityFault(string message)
|
||||
{
|
||||
Message = message;
|
||||
}
|
||||
|
||||
public string Message { get; set; }
|
||||
}
|
||||
}
|
||||
16
src/core/Elsa.Abstractions/Models/LogEntry.cs
Normal file
16
src/core/Elsa.Abstractions/Models/LogEntry.cs
Normal file
|
|
@ -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<string> TriggeredEndpoints { get; set; }
|
||||
public ActivityFault Fault { get; set; }
|
||||
}
|
||||
}
|
||||
|
|
@ -20,6 +20,7 @@ namespace Elsa.Models
|
|||
Scopes = new Stack<WorkflowExecutionScope>(new[] { CurrentScope });
|
||||
Arguments = new Variables();
|
||||
BlockingActivities = new List<IActivity>();
|
||||
ExecutionLog = new List<LogEntry>();
|
||||
Metadata = new WorkflowMetadata();
|
||||
}
|
||||
|
||||
|
|
@ -34,6 +35,7 @@ namespace Elsa.Models
|
|||
public WorkflowExecutionScope CurrentScope { get; set; }
|
||||
public Variables Arguments { get; set; }
|
||||
public IList<IActivity> BlockingActivities { get; set; }
|
||||
public IList<LogEntry> ExecutionLog { get; set; }
|
||||
public WorkflowMetadata Metadata { get; set; }
|
||||
public WorkflowFault Fault { get; set; }
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<IActivity> scheduledActivities;
|
||||
private readonly Stack<IActivity> scheduledHaltingActivities;
|
||||
|
||||
public WorkflowExecutionContext(Workflow workflow)
|
||||
{
|
||||
Workflow = workflow;
|
||||
IsFirstPass = true;
|
||||
scheduledActivities = new Stack<IActivity>();
|
||||
scheduledHaltingActivities = new Stack<IActivity>();
|
||||
}
|
||||
|
||||
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)
|
||||
{
|
||||
|
|
|
|||
|
|
@ -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<ActivityInvoker> logger)
|
||||
{
|
||||
this.driverRegistry = driverRegistry;
|
||||
this.clock = clock;
|
||||
this.logger = logger;
|
||||
}
|
||||
|
||||
public async Task<ActivityExecutionResult> 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<ActivityExecutionResult> 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<ActivityExecutionResult> InvokeAsync(
|
||||
public async Task<ActivityExecutionResult> 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<ActivityExecutionResult> InvokeAsync(
|
||||
WorkflowExecutionContext workflowContext,
|
||||
IActivity activity,
|
||||
Func<ActivityExecutionContext, IActivityDriver, Task<ActivityExecutionResult>> 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());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -29,6 +29,7 @@
|
|||
|
||||
<ItemGroup>
|
||||
<ProjectReference Include="..\Elsa.Abstractions\Elsa.Abstractions.csproj" />
|
||||
<ProjectReference Include="..\Elsa.Persistence.Abstractions\Elsa.Persistence.Abstractions.csproj" />
|
||||
</ItemGroup>
|
||||
|
||||
</Project>
|
||||
|
|
@ -15,16 +15,16 @@ namespace Elsa.Extensions
|
|||
{
|
||||
services.TryAddSingleton<IIdGenerator, DefaultIdGenerator>();
|
||||
services.TryAddSingleton<IClock>(SystemClock.Instance);
|
||||
services.TryAddScoped<IWorkflowSerializer, WorkflowSerializer>();
|
||||
services.TryAddScoped<IWorkflowTokenizer, WorkflowTokenizer>();
|
||||
services.AddScoped<IActivityHarvester, TypedActivityHarvester>();
|
||||
services.TryAddScoped<IActivityLibrary, ActivityLibrary>();
|
||||
services.TryAddSingleton<IWorkflowSerializer, WorkflowSerializer>();
|
||||
services.TryAddSingleton<IWorkflowTokenizer, WorkflowTokenizer>();
|
||||
services.TryAddSingleton<IActivityHarvester, TypedActivityHarvester>();
|
||||
services.TryAddSingleton<IActivityLibrary, ActivityLibrary>();
|
||||
services.TryAddSingleton<ITokenFormatterProvider, TokenFormatterProvider>();
|
||||
services.TryAddSingleton<ITokenizerInvoker, TokenizerInvoker>();
|
||||
services.AddActivityDescriptors<ActivityDescriptors>();
|
||||
services.AddSingleton<ITokenFormatter, JsonTokenFormatter>();
|
||||
services.AddSingleton<ITokenFormatter, YamlTokenFormatter>();
|
||||
services.AddSingleton<ITokenFormatter, XmlTokenFormatter>();
|
||||
services.TryAddSingleton<ITokenFormatterProvider, TokenFormatterProvider>();
|
||||
services.TryAddSingleton<ITokenizerInvoker, TokenizerInvoker>();
|
||||
services.AddSingleton<ITokenizer, DefaultTokenizer>();
|
||||
services.AddSingleton<ITokenizer, ActivityTokenizer>();
|
||||
services.AddSingleton<IExpressionEvaluator, PlainTextEvaluator>();
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<string> endpointNames)
|
||||
{
|
||||
EndpointNames = endpointNames;
|
||||
EndpointNames = endpointNames.ToList();
|
||||
}
|
||||
|
||||
public IEnumerable<string> EndpointNames { get; }
|
||||
public IReadOnlyList<string> 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));
|
||||
|
|
|
|||
|
|
@ -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<WorkflowExecutionScope>(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<int, WorkflowExecutionScope> DeserializeScopes(JToken token, WorkflowTokenizationContext context)
|
||||
{
|
||||
|
|
@ -236,9 +250,9 @@ namespace Elsa.Serialization.Tokenizers
|
|||
return scopeLookup;
|
||||
}
|
||||
|
||||
private IEnumerable<IActivity> DeserializeHaltedActivities(JToken token, IDictionary<string, IActivity> activityDictionary)
|
||||
private IEnumerable<IActivity> DeserializeBlockingActivities(JToken token, IDictionary<string, IActivity> 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<LogEntry> DeserializeExecutionLog(JToken token)
|
||||
{
|
||||
var models = (JArray) token["executionLog"] ?? new JArray();
|
||||
|
||||
foreach (var model in models)
|
||||
{
|
||||
var entry = JsonConvert.DeserializeObject<LogEntry>(model.ToString(), new JsonSerializerSettings().ConfigureForNodaTime(DateTimeZoneProviders.Tzdb));
|
||||
yield return entry;
|
||||
}
|
||||
}
|
||||
|
||||
private IEnumerable<Connection> DeserializeConnections(JToken token, IDictionary<string, IActivity> activityDictionary)
|
||||
{
|
||||
|
|
|
|||
|
|
@ -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<WorkflowInvoker> 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<ActivityExecutionResult> ExecuteActivityAsync(WorkflowExecutionContext workflowContext, IActivity activity, bool isResuming, CancellationToken cancellationToken)
|
||||
{
|
||||
return await ExecuteActivityAsync(workflowContext, activity, () => ExecuteOrResumeActivityAsync(workflowContext, activity, isResuming, cancellationToken), cancellationToken);
|
||||
}
|
||||
|
||||
private async Task<ActivityExecutionResult> ExecuteActivityHaltedAsync(WorkflowExecutionContext workflowContext, IActivity activity, CancellationToken cancellationToken)
|
||||
{
|
||||
return await ExecuteActivityAsync(workflowContext, activity, () => ActivityInvoker.HaltedAsync(workflowContext, activity, cancellationToken), cancellationToken);
|
||||
}
|
||||
|
||||
private async Task<ActivityExecutionResult> ExecuteActivityAsync(WorkflowExecutionContext workflowContext, IActivity activity, Func<Task<ActivityExecutionResult>> 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);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -12,6 +12,7 @@ namespace Elsa.Persistence
|
|||
Task<Workflow> GetAsync(ISpecification<Workflow, IWorkflowSpecificationVisitor> 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
|
||||
|
|
|
|||
|
|
@ -11,7 +11,7 @@ namespace Elsa.Persistence.FileSystem.Extensions
|
|||
public static IServiceCollection AddWorkflowsFileSystemPersistence(this IServiceCollection services, IConfiguration configuration)
|
||||
{
|
||||
services.TryAddSingleton<IFileSystem, System.IO.Abstractions.FileSystem>();
|
||||
services.AddScoped<IWorkflowStore, FileSystemWorkflowStore>();
|
||||
services.AddSingleton<IWorkflowStore, FileSystemWorkflowStore>();
|
||||
services.Configure<FileSystemStoreOptions>(configuration);
|
||||
|
||||
return services;
|
||||
|
|
|
|||
|
|
@ -21,9 +21,9 @@ namespace Elsa.Persistence.FileSystem
|
|||
private readonly string format;
|
||||
|
||||
public FileSystemWorkflowStore(
|
||||
IOptions<FileSystemStoreOptions> options,
|
||||
IFileSystem fileSystem,
|
||||
IIdGenerator idGenerator,
|
||||
IOptions<FileSystemStoreOptions> 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();
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -31,17 +31,12 @@ namespace Elsa.Runtime
|
|||
public async Task<WorkflowExecutionContext> 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<WorkflowExecutionContext> 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)
|
||||
|
|
|
|||
|
|
@ -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<IEnumerable<Workflow>> ListAsync(CancellationToken cancellationToken)
|
||||
{
|
||||
|
|
|
|||
Loading…
Reference in a new issue