From c2797e6b710bd252c71faed32f18940241f93f08 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Sat, 16 Feb 2019 11:00:30 +0100 Subject: [PATCH] Improved HTTP request activities (#18) * HTTP activity workflow fixes * Removed Http Cache for now in favor of implementing caching storage --- .../Drivers/ReadLineDriver.cs | 6 +- .../Drivers/WriteLineDriver.cs | 2 +- .../Activities/HttpRequestAction.cs | 1 - .../Activities/HttpResponseAction.cs | 3 +- .../ActivityDescriptors.cs | 1 - .../Drivers/HttpRequestActionDriver.cs | 2 +- .../Drivers/HttpRequestTriggerDriver.cs | 13 ++-- .../Drivers/HttpResponseActionDriver.cs | 2 +- .../Extensions/ServiceCollectionExtensions.cs | 6 -- .../HttpWorkflowCacheInitializer.cs | 52 ---------------- .../HttpRequestTriggerMiddleware.cs | 59 +++++++++++-------- .../Services/IHttpWorkflowCache.cs | 16 ----- .../DefaultHttpWorkflowCache.cs | 51 ---------------- .../Drivers/ForEachDriver.cs | 2 +- .../Drivers/IfElseDriver.cs | 2 +- .../Drivers/SetVariableDriver.cs | 2 +- .../Elsa.Abstractions/ActivityDriverBase.cs | 2 +- .../Models/WorkflowExecutionContext.cs | 14 +++-- src/core/Elsa.Core/ActivityInvoker.cs | 2 +- src/core/Elsa.Core/Handlers/ActivityDriver.cs | 8 +-- .../Handlers/UnknownActivityDriver.cs | 9 +-- .../Elsa.Core/Results/FaultWorkflowResult.cs | 8 +-- src/core/Elsa.Core/Results/HaltResult.cs | 9 +-- src/core/Elsa.Core/WorkflowInvoker.cs | 4 +- 24 files changed, 73 insertions(+), 203 deletions(-) delete mode 100644 src/activities/Elsa.Activities.Http/Initialization/HttpWorkflowCacheInitializer.cs delete mode 100644 src/activities/Elsa.Activities.Http/Services/IHttpWorkflowCache.cs delete mode 100644 src/activities/Elsa.Activities.Http/Services/Implementations/DefaultHttpWorkflowCache.cs diff --git a/src/activities/Elsa.Activities.Console/Drivers/ReadLineDriver.cs b/src/activities/Elsa.Activities.Console/Drivers/ReadLineDriver.cs index 7658782ed..3bc3c1789 100644 --- a/src/activities/Elsa.Activities.Console/Drivers/ReadLineDriver.cs +++ b/src/activities/Elsa.Activities.Console/Drivers/ReadLineDriver.cs @@ -28,18 +28,18 @@ namespace Elsa.Activities.Console.Drivers protected override async Task OnExecuteAsync(ReadLine activity, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) { if (input == null) - return Halt(clock.GetCurrentInstant()); + return Halt(); var value = await input.ReadLineAsync(); workflowContext.SetLastResult(value); - return TriggerEndpoint("Done"); + return Endpoint("Done"); } protected override ActivityExecutionResult OnResume(ReadLine activity, WorkflowExecutionContext workflowContext) { var receivedInput = workflowContext.Workflow.Arguments[activity.ArgumentName]; workflowContext.SetLastResult(receivedInput); - return TriggerEndpoint("Done"); + return Endpoint("Done"); } } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Console/Drivers/WriteLineDriver.cs b/src/activities/Elsa.Activities.Console/Drivers/WriteLineDriver.cs index a54182a83..b39c33a0b 100644 --- a/src/activities/Elsa.Activities.Console/Drivers/WriteLineDriver.cs +++ b/src/activities/Elsa.Activities.Console/Drivers/WriteLineDriver.cs @@ -32,7 +32,7 @@ namespace Elsa.Activities.Console.Drivers { var text = await evaluator.EvaluateAsync(activity.TextExpression, workflowContext, cancellationToken); await output.WriteLineAsync(text); - return TriggerEndpoint("Done"); + return Endpoint("Done"); } } } diff --git a/src/activities/Elsa.Activities.Http/Activities/HttpRequestAction.cs b/src/activities/Elsa.Activities.Http/Activities/HttpRequestAction.cs index 46617b76d..5e161fa1f 100644 --- a/src/activities/Elsa.Activities.Http/Activities/HttpRequestAction.cs +++ b/src/activities/Elsa.Activities.Http/Activities/HttpRequestAction.cs @@ -1,6 +1,5 @@ using System; using System.Collections.Generic; -using Elsa.Activities.Http.Models; using Elsa.Expressions; using Elsa.Models; diff --git a/src/activities/Elsa.Activities.Http/Activities/HttpResponseAction.cs b/src/activities/Elsa.Activities.Http/Activities/HttpResponseAction.cs index 8f4c26263..af5294e2f 100644 --- a/src/activities/Elsa.Activities.Http/Activities/HttpResponseAction.cs +++ b/src/activities/Elsa.Activities.Http/Activities/HttpResponseAction.cs @@ -1,5 +1,4 @@ -using System.Collections.Generic; -using System.Net; +using System.Net; using Elsa.Expressions; using Elsa.Models; diff --git a/src/activities/Elsa.Activities.Http/ActivityDescriptors.cs b/src/activities/Elsa.Activities.Http/ActivityDescriptors.cs index 562a5fb3f..c19662f51 100644 --- a/src/activities/Elsa.Activities.Http/ActivityDescriptors.cs +++ b/src/activities/Elsa.Activities.Http/ActivityDescriptors.cs @@ -2,7 +2,6 @@ using System.Collections.Generic; using System.Linq; using Elsa.Activities.Http.Activities; using Elsa.Models; -using Microsoft.AspNetCore.Http; using Microsoft.Extensions.Localization; namespace Elsa.Activities.Http diff --git a/src/activities/Elsa.Activities.Http/Drivers/HttpRequestActionDriver.cs b/src/activities/Elsa.Activities.Http/Drivers/HttpRequestActionDriver.cs index 69229186e..5bb9b1431 100644 --- a/src/activities/Elsa.Activities.Http/Drivers/HttpRequestActionDriver.cs +++ b/src/activities/Elsa.Activities.Http/Drivers/HttpRequestActionDriver.cs @@ -18,7 +18,7 @@ namespace Elsa.Activities.Http.Drivers protected override async Task OnExecuteAsync(HttpRequestAction activity, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) { - return TriggerEndpoint("Done"); + return Endpoint("Done"); } } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Http/Drivers/HttpRequestTriggerDriver.cs b/src/activities/Elsa.Activities.Http/Drivers/HttpRequestTriggerDriver.cs index 11e70e121..8da284a25 100644 --- a/src/activities/Elsa.Activities.Http/Drivers/HttpRequestTriggerDriver.cs +++ b/src/activities/Elsa.Activities.Http/Drivers/HttpRequestTriggerDriver.cs @@ -9,7 +9,6 @@ using Elsa.Handlers; using Elsa.Models; using Elsa.Results; using Microsoft.AspNetCore.Http; -using Microsoft.Extensions.Localization; namespace Elsa.Activities.Http.Drivers { @@ -26,7 +25,12 @@ namespace Elsa.Activities.Http.Drivers this.expressionEvaluator = expressionEvaluator; } - protected override async Task OnExecuteAsync(HttpRequestTrigger activity, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) + protected override ActivityExecutionResult OnExecute(HttpRequestTrigger activity, WorkflowExecutionContext workflowContext) + { + return Halt(); + } + + protected override async Task OnResumeAsync(HttpRequestTrigger activity, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) { var request = httpContextAccessor.HttpContext.Request; var model = new HttpRequestModel @@ -48,9 +52,10 @@ namespace Elsa.Activities.Http.Drivers model.Content = await request.ReadBodyAsync(); } } - + workflowContext.CurrentScope.LastResult = model; - return TriggerEndpoint("Done"); + + return Endpoint("Done"); } } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Http/Drivers/HttpResponseActionDriver.cs b/src/activities/Elsa.Activities.Http/Drivers/HttpResponseActionDriver.cs index 1921f45c6..66059607c 100644 --- a/src/activities/Elsa.Activities.Http/Drivers/HttpResponseActionDriver.cs +++ b/src/activities/Elsa.Activities.Http/Drivers/HttpResponseActionDriver.cs @@ -53,7 +53,7 @@ namespace Elsa.Activities.Http.Drivers await response.WriteAsync(bodyText, cancellationToken); } - return TriggerEndpoint("Done"); + return Endpoint("Done"); } } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Http/Extensions/ServiceCollectionExtensions.cs b/src/activities/Elsa.Activities.Http/Extensions/ServiceCollectionExtensions.cs index 1980627b1..87bc758fa 100644 --- a/src/activities/Elsa.Activities.Http/Extensions/ServiceCollectionExtensions.cs +++ b/src/activities/Elsa.Activities.Http/Extensions/ServiceCollectionExtensions.cs @@ -1,8 +1,4 @@ -using AspNetCore.AsyncInitialization; using Elsa.Activities.Http.Drivers; -using Elsa.Activities.Http.Initialization; -using Elsa.Activities.Http.Services; -using Elsa.Activities.Http.Services.Implementations; using Elsa.Extensions; using Microsoft.AspNetCore.Http; using Microsoft.Extensions.DependencyInjection; @@ -21,8 +17,6 @@ namespace Elsa.Activities.Http.Extensions { services.AddHttpWorkflowDescriptors(); services.AddAsyncInitialization(); - services.TryAddSingleton(); - services.TryAddTransient(); services.TryAddSingleton(); services diff --git a/src/activities/Elsa.Activities.Http/Initialization/HttpWorkflowCacheInitializer.cs b/src/activities/Elsa.Activities.Http/Initialization/HttpWorkflowCacheInitializer.cs deleted file mode 100644 index 458e9a80f..000000000 --- a/src/activities/Elsa.Activities.Http/Initialization/HttpWorkflowCacheInitializer.cs +++ /dev/null @@ -1,52 +0,0 @@ -using System.Collections.Generic; -using System.Linq; -using System.Threading; -using System.Threading.Tasks; -using AspNetCore.AsyncInitialization; -using Elsa.Activities.Http.Activities; -using Elsa.Activities.Http.Services; -using Elsa.Extensions; -using Elsa.Persistence; -using Elsa.Persistence.Extensions; -using Elsa.Persistence.Specifications; - -namespace Elsa.Activities.Http.Initialization -{ - public class HttpWorkflowCacheInitializer : IAsyncInitializer - { - private readonly IWorkflowStore workflowStore; - private readonly IHttpWorkflowCache httpWorkflowCache; - - public HttpWorkflowCacheInitializer(IWorkflowStore workflowStore, IHttpWorkflowCache httpWorkflowCache) - { - this.workflowStore = workflowStore; - this.httpWorkflowCache = httpWorkflowCache; - } - - public async Task InitializeAsync() - { - var specification = new WorkflowStartsWithActivity(nameof(HttpRequestTrigger)).Or(new WorkflowIsBlockedOnActivity(nameof(HttpRequestTrigger))); - var workflows = await workflowStore.GetManyAsync(specification, CancellationToken.None); - foreach (var workflow in workflows) - { - var activities = new List(); - - if (workflow.IsDefinition()) - { - var startActivities = workflow.GetStartActivities().Where(x => x is HttpRequestTrigger).Cast(); - activities.AddRange(startActivities); - } - else - { - var blockingActivities = workflow.BlockingActivities.Where(x => x is HttpRequestTrigger).Cast(); - activities.AddRange(blockingActivities); - } - - foreach (var activity in activities) - { - await httpWorkflowCache.AddWorkflowAsync(activity.Path, workflow, CancellationToken.None); - } - } - } - } -} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Http/Middleware/HttpRequestTriggerMiddleware.cs b/src/activities/Elsa.Activities.Http/Middleware/HttpRequestTriggerMiddleware.cs index a82de121f..bc2b14fd0 100644 --- a/src/activities/Elsa.Activities.Http/Middleware/HttpRequestTriggerMiddleware.cs +++ b/src/activities/Elsa.Activities.Http/Middleware/HttpRequestTriggerMiddleware.cs @@ -4,9 +4,11 @@ using System.Linq; using System.Threading; using System.Threading.Tasks; using Elsa.Activities.Http.Activities; -using Elsa.Activities.Http.Services; using Elsa.Extensions; using Elsa.Models; +using Elsa.Persistence; +using Elsa.Persistence.Extensions; +using Elsa.Persistence.Specifications; using Elsa.Runtime; using Microsoft.AspNetCore.Http; @@ -21,11 +23,11 @@ namespace Elsa.Activities.Http.Middleware this.next = next; } - public async Task InvokeAsync(HttpContext context, IWorkflowHost workflowHost, IHttpWorkflowCache httpWorkflowCache) + public async Task InvokeAsync(HttpContext context, IWorkflowHost workflowHost, IWorkflowStore workflowStore) { var requestPath = new Uri(context.Request.Path.ToString(), UriKind.Relative); var cancellationToken = context.RequestAborted; - var workflows = await httpWorkflowCache.GetWorkflowsByPathAsync(requestPath, cancellationToken).ToListAsync(); + var workflows = await GetWorkflowsByPathAsync(workflowStore, requestPath, cancellationToken).ToListAsync(); if (!workflows.Any()) { @@ -33,45 +35,50 @@ namespace Elsa.Activities.Http.Middleware } else { - await InvokeWorkflows(workflowHost, workflows, requestPath, cancellationToken); + await InvokeWorkflows(workflowHost, workflows, cancellationToken); } } - private async Task InvokeWorkflows(IWorkflowHost workflowHost, IEnumerable workflows, Uri requestPath, CancellationToken cancellationToken) + private async Task>> GetWorkflowsByPathAsync(IWorkflowStore workflowStore, Uri path, CancellationToken cancellationToken) + { + var specification = new WorkflowStartsWithActivity(nameof(HttpRequestTrigger)).Or(new WorkflowIsBlockedOnActivity(nameof(HttpRequestTrigger))); + var httpWorkflows = await workflowStore.GetManyAsync(specification, CancellationToken.None); + + var query = + from workflow in httpWorkflows + let activities = workflow.IsDefinition() ? workflow.GetStartActivities() : workflow.BlockingActivities + let triggers = FilterByPath(activities, path) + select triggers.Select(x => Tuple.Create(workflow, x)); + + return query.SelectMany(x => x); + } + + private IEnumerable FilterByPath(IEnumerable activities, Uri path) + { + return activities.Where(x => x is HttpRequestTrigger trigger && trigger.Path == path).Cast(); + } + + private async Task InvokeWorkflows(IWorkflowHost workflowHost, IEnumerable> workflows, CancellationToken cancellationToken) { foreach (var workflow in workflows) { - await InvokeWorkflowAsync(workflowHost, workflow, requestPath, cancellationToken); + await InvokeWorkflowAsync(workflowHost, workflow, cancellationToken); } } - private async Task InvokeWorkflowAsync(IWorkflowHost workflowHost, Workflow workflow, Uri requestPath, CancellationToken cancellationToken) + private async Task InvokeWorkflowAsync(IWorkflowHost workflowHost, Tuple workflowTuple, CancellationToken cancellationToken) { + var workflow = workflowTuple.Item1; + var activity = workflowTuple.Item2; + if (workflow.Status == WorkflowStatus.Idle) { - await StartHttpWorkflowAsync(workflowHost, workflow, requestPath, cancellationToken); + await workflowHost.StartWorkflowAsync(workflow, activity, Variables.Empty, cancellationToken); } else if (workflow.Status == WorkflowStatus.Halted) { - await ResumeHttpWorkflowAsync(workflowHost, workflow, requestPath, cancellationToken); + await workflowHost.ResumeWorkflowAsync(workflow, activity, Variables.Empty, cancellationToken); } } - - private async Task StartHttpWorkflowAsync(IWorkflowHost workflowHost, Workflow workflow, Uri requestPath, CancellationToken cancellationToken) - { - var startActivity = GetActivityByRequestPath(workflow, requestPath); - await workflowHost.StartWorkflowAsync(workflow, startActivity, Variables.Empty, cancellationToken); - } - - private async Task ResumeHttpWorkflowAsync(IWorkflowHost workflowHost, Workflow workflow, Uri requestPath, CancellationToken cancellationToken) - { - var blockingActivity = GetActivityByRequestPath(workflow, requestPath); - await workflowHost.ResumeWorkflowAsync(workflow, blockingActivity, Variables.Empty, cancellationToken); - } - - private static HttpRequestTrigger GetActivityByRequestPath(Workflow workflow, Uri requestPath) - { - return (HttpRequestTrigger)workflow.Activities.Single(x => x is HttpRequestTrigger activity && activity.Path == requestPath); - } } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Http/Services/IHttpWorkflowCache.cs b/src/activities/Elsa.Activities.Http/Services/IHttpWorkflowCache.cs deleted file mode 100644 index ef8b82b3d..000000000 --- a/src/activities/Elsa.Activities.Http/Services/IHttpWorkflowCache.cs +++ /dev/null @@ -1,16 +0,0 @@ -using System; -using System.Collections.Generic; -using System.Threading; -using System.Threading.Tasks; -using Elsa.Models; - -namespace Elsa.Activities.Http.Services -{ - public interface IHttpWorkflowCache - { - Task AddWorkflowAsync(Uri requestPath, Workflow workflow, CancellationToken cancellationToken); - Task RemoveWorkflowAsync(Uri requestPath, Workflow workflow, CancellationToken cancellationToken); - Task> GetWorkflowsByPathAsync(Uri requestPath, CancellationToken cancellationToken); - IEnumerable GetWorkflowsByPath(Uri requestPath); - } -} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Http/Services/Implementations/DefaultHttpWorkflowCache.cs b/src/activities/Elsa.Activities.Http/Services/Implementations/DefaultHttpWorkflowCache.cs deleted file mode 100644 index 5ac736c66..000000000 --- a/src/activities/Elsa.Activities.Http/Services/Implementations/DefaultHttpWorkflowCache.cs +++ /dev/null @@ -1,51 +0,0 @@ -using System; -using System.Collections.Concurrent; -using System.Collections.Generic; -using System.Linq; -using System.Threading; -using System.Threading.Tasks; -using Elsa.Models; - -namespace Elsa.Activities.Http.Services.Implementations -{ - public class DefaultHttpWorkflowCache : IHttpWorkflowCache - { - private readonly ConcurrentDictionary> dictionary; - - public DefaultHttpWorkflowCache() - { - dictionary = new ConcurrentDictionary>(); - } - - public Task AddWorkflowAsync(Uri requestPath, Workflow workflow, CancellationToken cancellationToken) - { - var workflows = dictionary.GetOrAdd(requestPath, key => new List()); - workflows.Add(workflow); - return Task.CompletedTask; - } - - public Task RemoveWorkflowAsync(Uri requestPath, Workflow workflow, CancellationToken cancellationToken) - { - if (dictionary.TryGetValue(requestPath, out var workflows)) - { - workflows.Remove(workflow); - } - - return Task.CompletedTask; - } - - public Task> GetWorkflowsByPathAsync(Uri requestPath, CancellationToken cancellationToken) - { - return dictionary.TryGetValue(requestPath, out var workflows) - ? Task.FromResult>(workflows) - : Task.FromResult(Enumerable.Empty()); - } - - public IEnumerable GetWorkflowsByPath(Uri requestPath) - { - return dictionary.TryGetValue(requestPath, out var workflows) - ? workflows - : Enumerable.Empty(); - } - } -} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Primitives/Drivers/ForEachDriver.cs b/src/activities/Elsa.Activities.Primitives/Drivers/ForEachDriver.cs index 497a758df..7caef1c74 100644 --- a/src/activities/Elsa.Activities.Primitives/Drivers/ForEachDriver.cs +++ b/src/activities/Elsa.Activities.Primitives/Drivers/ForEachDriver.cs @@ -18,7 +18,7 @@ namespace Elsa.Activities.Primitives.Drivers protected override async Task OnExecuteAsync(ForEach activity, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) { - return TriggerEndpoint("Done"); + return Endpoint("Done"); } } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Primitives/Drivers/IfElseDriver.cs b/src/activities/Elsa.Activities.Primitives/Drivers/IfElseDriver.cs index ba9822909..674b97e1c 100644 --- a/src/activities/Elsa.Activities.Primitives/Drivers/IfElseDriver.cs +++ b/src/activities/Elsa.Activities.Primitives/Drivers/IfElseDriver.cs @@ -19,7 +19,7 @@ namespace Elsa.Activities.Primitives.Drivers protected override async Task OnExecuteAsync(IfElse activity, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) { var result = await expressionEvaluator.EvaluateAsync(activity.ConditionExpression, workflowContext, cancellationToken); - return TriggerEndpoint(result ? "True" : "False"); + return Endpoint(result ? "True" : "False"); } } diff --git a/src/activities/Elsa.Activities.Primitives/Drivers/SetVariableDriver.cs b/src/activities/Elsa.Activities.Primitives/Drivers/SetVariableDriver.cs index 00dcab149..5d6c4f21f 100644 --- a/src/activities/Elsa.Activities.Primitives/Drivers/SetVariableDriver.cs +++ b/src/activities/Elsa.Activities.Primitives/Drivers/SetVariableDriver.cs @@ -21,7 +21,7 @@ namespace Elsa.Activities.Primitives.Drivers { //var value = await expressionEvaluator.EvaluateAsync(activity.ValueExpression, workflowContext, cancellationToken); //workflowContext.CurrentScope.SetVariable(activity.VariableName, value); - return TriggerEndpoint("Done"); + return Endpoint("Done"); } } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/ActivityDriverBase.cs b/src/core/Elsa.Abstractions/ActivityDriverBase.cs index c03707146..20080c258 100644 --- a/src/core/Elsa.Abstractions/ActivityDriverBase.cs +++ b/src/core/Elsa.Abstractions/ActivityDriverBase.cs @@ -18,7 +18,7 @@ namespace Elsa 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(ActivityExecutionContext activityContext, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) => OnResumeAsync((T)activityContext.Activity, workflowContext, cancellationToken); 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; diff --git a/src/core/Elsa.Abstractions/Models/WorkflowExecutionContext.cs b/src/core/Elsa.Abstractions/Models/WorkflowExecutionContext.cs index 1bdeae203..38c2a7326 100644 --- a/src/core/Elsa.Abstractions/Models/WorkflowExecutionContext.cs +++ b/src/core/Elsa.Abstractions/Models/WorkflowExecutionContext.cs @@ -8,11 +8,13 @@ namespace Elsa.Models { public class WorkflowExecutionContext { + private readonly IClock clock; private readonly Stack scheduledActivities; private readonly Stack scheduledHaltingActivities; - public WorkflowExecutionContext(Workflow workflow) + public WorkflowExecutionContext(Workflow workflow, IClock clock) { + this.clock = clock; Workflow = workflow; IsFirstPass = true; scheduledActivities = new Stack(); @@ -72,11 +74,11 @@ namespace Elsa.Models CurrentScope.LastResult = value; } - public void Fault(Exception exception, IActivity activity, Instant instant) => Fault(exception.Message, activity, instant); + public void Fault(Exception exception, IActivity activity) => Fault(exception.Message, activity); - public void Fault(string errorMessage, IActivity activity, Instant instant) + public void Fault(string errorMessage, IActivity activity) { - Workflow.FinishedAt = instant; + Workflow.FinishedAt = clock.GetCurrentInstant(); Workflow.Fault = new WorkflowFault { Message = errorMessage, @@ -84,7 +86,7 @@ namespace Elsa.Models }; } - public void Halt(Instant instant) + public void Halt() { var activity = CurrentActivity; if (!Workflow.BlockingActivities.Contains(activity)) @@ -92,7 +94,7 @@ namespace Elsa.Models Workflow.BlockingActivities.Add(activity); } - Workflow.HaltedAt = instant; + Workflow.HaltedAt = clock.GetCurrentInstant(); Workflow.Status = WorkflowStatus.Halted; } diff --git a/src/core/Elsa.Core/ActivityInvoker.cs b/src/core/Elsa.Core/ActivityInvoker.cs index d02445b39..848921d73 100644 --- a/src/core/Elsa.Core/ActivityInvoker.cs +++ b/src/core/Elsa.Core/ActivityInvoker.cs @@ -77,7 +77,7 @@ namespace Elsa { logger.LogError(e, "Error while invoking activity {ActivityId} of workflow {WorkflowId}", activity.Id, workflowContext.Workflow.Metadata.Id); workflowContext.Workflow.AddLogEntry(activity.Id, clock.GetCurrentInstant(), e.Message, true); - return new FaultWorkflowResult(e, clock.GetCurrentInstant()); + return new FaultWorkflowResult(e); } } } diff --git a/src/core/Elsa.Core/Handlers/ActivityDriver.cs b/src/core/Elsa.Core/Handlers/ActivityDriver.cs index 745cc11c3..0376a61fd 100644 --- a/src/core/Elsa.Core/Handlers/ActivityDriver.cs +++ b/src/core/Elsa.Core/Handlers/ActivityDriver.cs @@ -7,13 +7,13 @@ namespace Elsa.Handlers { public abstract class ActivityDriver : ActivityDriverBase where T : IActivity { - protected HaltResult Halt(Instant instant) => new HaltResult(instant); + protected HaltResult Halt() => new HaltResult(); protected TriggerEndpointsResult TriggerEndpoints(IEnumerable names) => new TriggerEndpointsResult(names); - protected TriggerEndpointsResult TriggerEndpoint(string name) => TriggerEndpoints(new[] { name }); + protected TriggerEndpointsResult Endpoint(string name) => TriggerEndpoints(new[] { name }); protected ScheduleActivityResult ScheduleActivity(IActivity activity) => new ScheduleActivityResult(activity); protected ReturnValueResult SetReturnValue(object value) => new ReturnValueResult(value); protected FinishWorkflowResult Finish(Instant instant) => new FinishWorkflowResult(instant); - protected FaultWorkflowResult Fault(string errorMessage, Instant instant) => new FaultWorkflowResult(errorMessage, instant); - protected FaultWorkflowResult Fault(Exception exception, Instant instant) => new FaultWorkflowResult(exception, instant); + protected FaultWorkflowResult Fault(string errorMessage) => new FaultWorkflowResult(errorMessage); + protected FaultWorkflowResult Fault(Exception exception) => new FaultWorkflowResult(exception); } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Handlers/UnknownActivityDriver.cs b/src/core/Elsa.Core/Handlers/UnknownActivityDriver.cs index 4b1417d42..5b36ea812 100644 --- a/src/core/Elsa.Core/Handlers/UnknownActivityDriver.cs +++ b/src/core/Elsa.Core/Handlers/UnknownActivityDriver.cs @@ -7,16 +7,9 @@ namespace Elsa.Handlers { public class UnknownActivityDriver : ActivityDriver { - private readonly IClock clock; - - public UnknownActivityDriver(IClock clock) - { - this.clock = clock; - } - protected override ActivityExecutionResult OnExecute(UnknownActivity activity, WorkflowExecutionContext workflowContext) { - return Fault($"Unknown activity: {activity.Name}, ID: {activity.Id}", clock.GetCurrentInstant()); + return Fault($"Unknown activity: {activity.Name}, ID: {activity.Id}"); } } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Results/FaultWorkflowResult.cs b/src/core/Elsa.Core/Results/FaultWorkflowResult.cs index 11df456fd..9418f6a20 100644 --- a/src/core/Elsa.Core/Results/FaultWorkflowResult.cs +++ b/src/core/Elsa.Core/Results/FaultWorkflowResult.cs @@ -7,22 +7,20 @@ namespace Elsa.Results public class FaultWorkflowResult : ActivityExecutionResult { private readonly string errorMessage; - private readonly Instant instant; - public FaultWorkflowResult(Exception exception, Instant instant) : this(exception.Message, instant) + public FaultWorkflowResult(Exception exception) : this(exception.Message) { } - public FaultWorkflowResult(string errorMessage, Instant instant) + public FaultWorkflowResult(string errorMessage) { this.errorMessage = errorMessage; - this.instant = instant; } protected override void Execute(IWorkflowInvoker invoker, WorkflowExecutionContext workflowContext) { var activity = workflowContext.CurrentActivity; - workflowContext.Fault(errorMessage, activity, instant); + workflowContext.Fault(errorMessage, activity); } } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Results/HaltResult.cs b/src/core/Elsa.Core/Results/HaltResult.cs index 0a65121f2..37f18c92c 100644 --- a/src/core/Elsa.Core/Results/HaltResult.cs +++ b/src/core/Elsa.Core/Results/HaltResult.cs @@ -10,13 +10,6 @@ namespace Elsa.Results /// public class HaltResult : ActivityExecutionResult { - private readonly Instant instant; - - public HaltResult(Instant instant) - { - this.instant = instant; - } - public override async Task ExecuteAsync(IWorkflowInvoker invoker, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) { var activity = workflowContext.CurrentActivity; @@ -31,7 +24,7 @@ namespace Elsa.Results else { workflowContext.ScheduleHaltingActivity(activity); - workflowContext.Halt(instant); + workflowContext.Halt(); } } } diff --git a/src/core/Elsa.Core/WorkflowInvoker.cs b/src/core/Elsa.Core/WorkflowInvoker.cs index 57c4b9c9c..ed44066fb 100644 --- a/src/core/Elsa.Core/WorkflowInvoker.cs +++ b/src/core/Elsa.Core/WorkflowInvoker.cs @@ -35,7 +35,7 @@ namespace Elsa public async Task InvokeAsync(Workflow workflow, IActivity startActivity = default, Variables arguments = default, CancellationToken cancellationToken = default) { workflow.Arguments = arguments ?? new Variables(); - var workflowExecutionContext = new WorkflowExecutionContext(workflow); + var workflowExecutionContext = new WorkflowExecutionContext(workflow, clock); 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. @@ -134,7 +134,7 @@ namespace Elsa ex, "An unhandled error occurred while executing an activity. Putting the workflow in the faulted state." ); - workflowContext.Fault(ex, activity, clock.GetCurrentInstant()); + workflowContext.Fault(ex, activity); } private async Task ExecuteOrResumeActivityAsync(WorkflowExecutionContext workflowContext, IActivity activity, bool isResuming, CancellationToken cancellationToken)