Improved HTTP request activities (#18)

* HTTP activity workflow fixes

* Removed Http Cache for now in favor of implementing caching storage
This commit is contained in:
Sipke Schoorstra 2019-02-16 11:00:30 +01:00 committed by GitHub
parent 04e95d8a38
commit c2797e6b71
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
24 changed files with 73 additions and 203 deletions

View file

@ -28,18 +28,18 @@ namespace Elsa.Activities.Console.Drivers
protected override async Task<ActivityExecutionResult> 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");
}
}
}

View file

@ -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");
}
}
}

View file

@ -1,6 +1,5 @@
using System;
using System.Collections.Generic;
using Elsa.Activities.Http.Models;
using Elsa.Expressions;
using Elsa.Models;

View file

@ -1,5 +1,4 @@
using System.Collections.Generic;
using System.Net;
using System.Net;
using Elsa.Expressions;
using Elsa.Models;

View file

@ -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

View file

@ -18,7 +18,7 @@ namespace Elsa.Activities.Http.Drivers
protected override async Task<ActivityExecutionResult> OnExecuteAsync(HttpRequestAction activity, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken)
{
return TriggerEndpoint("Done");
return Endpoint("Done");
}
}
}

View file

@ -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<ActivityExecutionResult> OnExecuteAsync(HttpRequestTrigger activity, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken)
protected override ActivityExecutionResult OnExecute(HttpRequestTrigger activity, WorkflowExecutionContext workflowContext)
{
return Halt();
}
protected override async Task<ActivityExecutionResult> 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");
}
}
}

View file

@ -53,7 +53,7 @@ namespace Elsa.Activities.Http.Drivers
await response.WriteAsync(bodyText, cancellationToken);
}
return TriggerEndpoint("Done");
return Endpoint("Done");
}
}
}

View file

@ -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<IHttpWorkflowCache, DefaultHttpWorkflowCache>();
services.TryAddTransient<IAsyncInitializer, HttpWorkflowCacheInitializer>();
services.TryAddSingleton<IHttpContextAccessor, HttpContextAccessor>();
services

View file

@ -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<HttpRequestTrigger>();
if (workflow.IsDefinition())
{
var startActivities = workflow.GetStartActivities().Where(x => x is HttpRequestTrigger).Cast<HttpRequestTrigger>();
activities.AddRange(startActivities);
}
else
{
var blockingActivities = workflow.BlockingActivities.Where(x => x is HttpRequestTrigger).Cast<HttpRequestTrigger>();
activities.AddRange(blockingActivities);
}
foreach (var activity in activities)
{
await httpWorkflowCache.AddWorkflowAsync(activity.Path, workflow, CancellationToken.None);
}
}
}
}
}

View file

@ -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<Workflow> workflows, Uri requestPath, CancellationToken cancellationToken)
private async Task<IEnumerable<Tuple<Workflow, HttpRequestTrigger>>> 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<HttpRequestTrigger> FilterByPath(IEnumerable<IActivity> activities, Uri path)
{
return activities.Where(x => x is HttpRequestTrigger trigger && trigger.Path == path).Cast<HttpRequestTrigger>();
}
private async Task InvokeWorkflows(IWorkflowHost workflowHost, IEnumerable<Tuple<Workflow, HttpRequestTrigger>> 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<Workflow, HttpRequestTrigger> 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);
}
}
}

View file

@ -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<IEnumerable<Workflow>> GetWorkflowsByPathAsync(Uri requestPath, CancellationToken cancellationToken);
IEnumerable<Workflow> GetWorkflowsByPath(Uri requestPath);
}
}

View file

@ -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<Uri, ICollection<Workflow>> dictionary;
public DefaultHttpWorkflowCache()
{
dictionary = new ConcurrentDictionary<Uri, ICollection<Workflow>>();
}
public Task AddWorkflowAsync(Uri requestPath, Workflow workflow, CancellationToken cancellationToken)
{
var workflows = dictionary.GetOrAdd(requestPath, key => new List<Workflow>());
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<IEnumerable<Workflow>> GetWorkflowsByPathAsync(Uri requestPath, CancellationToken cancellationToken)
{
return dictionary.TryGetValue(requestPath, out var workflows)
? Task.FromResult<IEnumerable<Workflow>>(workflows)
: Task.FromResult(Enumerable.Empty<Workflow>());
}
public IEnumerable<Workflow> GetWorkflowsByPath(Uri requestPath)
{
return dictionary.TryGetValue(requestPath, out var workflows)
? workflows
: Enumerable.Empty<Workflow>();
}
}
}

View file

@ -18,7 +18,7 @@ namespace Elsa.Activities.Primitives.Drivers
protected override async Task<ActivityExecutionResult> OnExecuteAsync(ForEach activity, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken)
{
return TriggerEndpoint("Done");
return Endpoint("Done");
}
}
}

View file

@ -19,7 +19,7 @@ namespace Elsa.Activities.Primitives.Drivers
protected override async Task<ActivityExecutionResult> 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");
}
}

View file

@ -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");
}
}
}

View file

@ -18,7 +18,7 @@ namespace Elsa
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(ActivityExecutionContext activityContext, WorkflowExecutionContext workflowContext, CancellationToken cancellationToken) => OnResumeAsync((T)activityContext.Activity, workflowContext, cancellationToken);
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;

View file

@ -8,11 +8,13 @@ namespace Elsa.Models
{
public class WorkflowExecutionContext
{
private readonly IClock clock;
private readonly Stack<IActivity> scheduledActivities;
private readonly Stack<IActivity> scheduledHaltingActivities;
public WorkflowExecutionContext(Workflow workflow)
public WorkflowExecutionContext(Workflow workflow, IClock clock)
{
this.clock = clock;
Workflow = workflow;
IsFirstPass = true;
scheduledActivities = new Stack<IActivity>();
@ -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;
}

View file

@ -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);
}
}
}

View file

@ -7,13 +7,13 @@ namespace Elsa.Handlers
{
public abstract class ActivityDriver<T> : ActivityDriverBase<T> where T : IActivity
{
protected HaltResult Halt(Instant instant) => new HaltResult(instant);
protected HaltResult Halt() => new HaltResult();
protected TriggerEndpointsResult TriggerEndpoints(IEnumerable<string> 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);
}
}

View file

@ -7,16 +7,9 @@ namespace Elsa.Handlers
{
public class UnknownActivityDriver : ActivityDriver<UnknownActivity>
{
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}");
}
}
}

View file

@ -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);
}
}
}

View file

@ -10,13 +10,6 @@ namespace Elsa.Results
/// </summary>
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();
}
}
}

View file

@ -35,7 +35,7 @@ namespace Elsa
public async Task<WorkflowExecutionContext> 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<ActivityExecutionResult> ExecuteOrResumeActivityAsync(WorkflowExecutionContext workflowContext, IActivity activity, bool isResuming, CancellationToken cancellationToken)