diff --git a/src/activities/Elsa.Activities.Http/Elsa.Activities.Http.csproj b/src/activities/Elsa.Activities.Http/Elsa.Activities.Http.csproj index 87fa2d920..95718bee4 100644 --- a/src/activities/Elsa.Activities.Http/Elsa.Activities.Http.csproj +++ b/src/activities/Elsa.Activities.Http/Elsa.Activities.Http.csproj @@ -27,7 +27,6 @@ - diff --git a/src/activities/Elsa.Activities.Http/Extensions/HttpActivitiesServiceCollectionExtensions.cs b/src/activities/Elsa.Activities.Http/Extensions/HttpActivitiesServiceCollectionExtensions.cs index 09d55e2ce..19b5a06ff 100644 --- a/src/activities/Elsa.Activities.Http/Extensions/HttpActivitiesServiceCollectionExtensions.cs +++ b/src/activities/Elsa.Activities.Http/Extensions/HttpActivitiesServiceCollectionExtensions.cs @@ -1,12 +1,12 @@ using System; using Elsa.Activities.Http.Activities; using Elsa.Activities.Http.Formatters; +using Elsa.Activities.Http.Liquid; using Elsa.Activities.Http.Options; using Elsa.Activities.Http.RequestHandlers.Handlers; using Elsa.Activities.Http.Services; -using Elsa.Scripting; -using Elsa.Scripting.JavaScript; -using MediatR; +using Elsa.Extensions; +using Elsa.Scripting.Liquid.Extensions; using Microsoft.AspNetCore.Http; using Microsoft.AspNetCore.Mvc.Infrastructure; using Microsoft.Extensions.DependencyInjection; @@ -38,11 +38,12 @@ namespace Elsa.Activities.Http.Extensions .AddSingleton() .AddSingleton() .AddHttpContextAccessor() - .AddMediatR(typeof(HttpActivitiesServiceCollectionExtensions)) + .AddNotificationHandlers(typeof(HttpActivitiesServiceCollectionExtensions)) .AddDataProtection(); - + + services.AddLiquidFilter("signal_url"); + return services - .AddScoped(sp => sp.GetRequiredService().HttpContext) .AddRequestHandler() .AddRequestHandler(); } diff --git a/src/activities/Elsa.Activities.Http/Liquid/SignalUrlFilter.cs b/src/activities/Elsa.Activities.Http/Liquid/SignalUrlFilter.cs new file mode 100644 index 000000000..128b0e231 --- /dev/null +++ b/src/activities/Elsa.Activities.Http/Liquid/SignalUrlFilter.cs @@ -0,0 +1,46 @@ +using System; +using System.Threading.Tasks; +using Elsa.Activities.Http.Models; +using Elsa.Activities.Http.Services; +using Elsa.Scripting.Liquid.Services; +using Elsa.Services.Models; +using Fluid; +using Fluid.Values; + +namespace Elsa.Activities.Http.Liquid +{ + public class SignalUrlFilter : ILiquidFilter + { + private readonly ITokenService tokenService; + private readonly IAbsoluteUrlProvider absoluteUrlProvider; + + public SignalUrlFilter(ITokenService tokenService, IAbsoluteUrlProvider absoluteUrlProvider) + { + this.tokenService = tokenService; + this.absoluteUrlProvider = absoluteUrlProvider; + } + + public ValueTask ProcessAsync(FluidValue input, FilterArguments arguments, TemplateContext context) + { + var workflowContextValue = context.GetValue("WorkflowExecutionContext"); + + if (workflowContextValue.IsNil()) + throw new ArgumentException("WorkflowExecutionContext missing while invoking 'signal_url'"); + + var workflowContext = (WorkflowExecutionContext)workflowContextValue.ToObjectValue(); + var signalName = input.ToStringValue(); + var url = GenerateUrl(signalName, workflowContext); + return new ValueTask(new StringValue(url)); + } + + private string GenerateUrl(string signal, WorkflowExecutionContext workflowExecutionContext) + { + var workflowInstanceId = workflowExecutionContext.Workflow.Id; + var payload = new Signal(signal, workflowInstanceId); + var token = tokenService.CreateToken(payload); + var url = $"/workflows/signal?token={token}"; + + return absoluteUrlProvider.ToAbsoluteUrl(url).ToString(); + } + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Http/RequestHandlers/Handlers/SignalRequestHandler.cs b/src/activities/Elsa.Activities.Http/RequestHandlers/Handlers/SignalRequestHandler.cs index b564a79af..8453ccbc0 100644 --- a/src/activities/Elsa.Activities.Http/RequestHandlers/Handlers/SignalRequestHandler.cs +++ b/src/activities/Elsa.Activities.Http/RequestHandlers/Handlers/SignalRequestHandler.cs @@ -3,12 +3,11 @@ using System.Net; using System.Threading; using System.Threading.Tasks; using Elsa.Activities.Http.Models; +using Elsa.Activities.Http.RequestHandlers.Results; using Elsa.Activities.Http.Services; using Elsa.Models; using Elsa.Persistence; using Elsa.Services; -using LanguageExt; -using LanguageExt.Common; using Microsoft.AspNetCore.Http; namespace Elsa.Activities.Http.RequestHandlers.Handlers @@ -42,58 +41,39 @@ namespace Elsa.Activities.Http.RequestHandlers.Handlers public async Task HandleRequestAsync() { - await DecryptToken() - .BindAsync(GetWorkflowInstanceAsync) - .BindAsync(CheckIfExecutingAsync) - .MapAsync(ResumeWorkflowAsync); + var signal = DecryptToken(); - return default; + if (signal == null) + return new NotFoundResult(); + + var workflowInstance = await GetWorkflowInstanceAsync(signal); + + if (workflowInstance == null) + return new NotFoundResult(); + + if (!CheckIfExecuting(workflowInstance)) + return new BadRequestResult($"Cannot signal a workflow with status other than {WorkflowStatus.Executing}. Actual workflow status: {workflowInstance.Status}."); + + await ResumeWorkflowAsync(workflowInstance, signal); + + return new AcceptedResult(); } - private Either DecryptToken() + private Signal DecryptToken() { var token = httpContext.Request.Query["token"]; - if (tokenService.TryDecryptToken(token, out Signal signal)) - { - return signal; - } - - httpContext.Response.StatusCode = (int)HttpStatusCode.NotFound; - return Error.New("Invalid token"); + return tokenService.TryDecryptToken(token, out Signal signal) ? signal : default; } - private async Task> GetWorkflowInstanceAsync(Signal signal) + private async Task GetWorkflowInstanceAsync(Signal signal) => + await workflowInstanceStore.GetByIdAsync(signal.WorkflowInstanceId, cancellationToken); + + private bool CheckIfExecuting(WorkflowInstance workflowInstance) => + workflowInstance.Status == WorkflowStatus.Executing; + + private async Task ResumeWorkflowAsync(WorkflowInstance workflowInstance, Signal signal) { - var workflowInstance = - await workflowInstanceStore.GetByIdAsync(signal.WorkflowInstanceId, cancellationToken); - - if (workflowInstance != null) - return (workflowInstance, signal); - - httpContext.Response.StatusCode = (int)HttpStatusCode.NotFound; - return Error.New("Workflow not found"); - } - - private async Task> CheckIfExecutingAsync( - (WorkflowInstance, Signal) tuple) - { - var (workflowInstance, signal) = tuple; - - if (workflowInstance.Status == WorkflowStatus.Executing) - return (workflowInstance, signal); - - httpContext.Response.StatusCode = (int)HttpStatusCode.BadRequest; - await httpContext.Response.WriteAsync( - $"Cannot signal a workflow with status other than {WorkflowStatus.Executing}. Actual workflow status: {workflowInstance.Status}.", - cancellationToken); - return Error.New("Cannot resume workflow that is not executing."); - } - - private async Task ResumeWorkflowAsync((WorkflowInstance, Signal) tuple) - { - var (workflowInstance, signal) = tuple; - var input = new Variables { ["Signal"] = signal.Name @@ -107,11 +87,6 @@ namespace Elsa.Activities.Http.RequestHandlers.Handlers var workflow = workflowFactory.CreateWorkflow(workflowDefinition, input, workflowInstance); var blockingSignalActivities = workflow.BlockingActivities.ToList(); await workflowInvoker.ResumeAsync(workflow, blockingSignalActivities, cancellationToken); - - if (!httpContext.Response.HasStarted) - { - httpContext.Response.StatusCode = (int)HttpStatusCode.Accepted; - } } } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Http/RequestHandlers/Handlers/TriggerRequestHandler.cs b/src/activities/Elsa.Activities.Http/RequestHandlers/Handlers/TriggerRequestHandler.cs index d11baeb50..95e200956 100644 --- a/src/activities/Elsa.Activities.Http/RequestHandlers/Handlers/TriggerRequestHandler.cs +++ b/src/activities/Elsa.Activities.Http/RequestHandlers/Handlers/TriggerRequestHandler.cs @@ -24,16 +24,16 @@ namespace Elsa.Activities.Http.RequestHandlers.Handlers private readonly CancellationToken cancellationToken; public TriggerRequestHandler( - HttpContext httpContext, + IHttpContextAccessor httpContext, IWorkflowInvoker workflowInvoker, IWorkflowRegistry registry, IWorkflowInstanceStore workflowInstanceStore) { - this.httpContext = httpContext; + this.httpContext = httpContext.HttpContext; this.workflowInvoker = workflowInvoker; this.registry = registry; this.workflowInstanceStore = workflowInstanceStore; - cancellationToken = httpContext.RequestAborted; + cancellationToken = httpContext.HttpContext.RequestAborted; } public async Task HandleRequestAsync() diff --git a/src/activities/Elsa.Activities.Http/RequestHandlers/Results/BadRequestResult.cs b/src/activities/Elsa.Activities.Http/RequestHandlers/Results/BadRequestResult.cs index 2f7f433a8..489785154 100644 --- a/src/activities/Elsa.Activities.Http/RequestHandlers/Results/BadRequestResult.cs +++ b/src/activities/Elsa.Activities.Http/RequestHandlers/Results/BadRequestResult.cs @@ -7,10 +7,24 @@ namespace Elsa.Activities.Http.RequestHandlers.Results { public class BadRequestResult : IRequestHandlerResult { - public Task ExecuteResultAsync(HttpContext httpContext, RequestDelegate next) + public BadRequestResult() + { + } + + public BadRequestResult(string message) + { + Message = message; + } + + public string Message { get; } + + + public async Task ExecuteResultAsync(HttpContext httpContext, RequestDelegate next) { httpContext.Response.StatusCode = (int)HttpStatusCode.BadRequest; - return Task.CompletedTask; + + if(!string.IsNullOrWhiteSpace(Message)) + await httpContext.Response.WriteAsync(Message, httpContext.RequestAborted); } } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Workflows/Activities/TriggerSignal.cs b/src/activities/Elsa.Activities.Workflows/Activities/TriggerSignal.cs index cb08491e2..9c0ad86d2 100644 --- a/src/activities/Elsa.Activities.Workflows/Activities/TriggerSignal.cs +++ b/src/activities/Elsa.Activities.Workflows/Activities/TriggerSignal.cs @@ -5,7 +5,6 @@ using Elsa.Expressions; using Elsa.Extensions; using Elsa.Models; using Elsa.Results; -using Elsa.Scripting.JavaScript; using Elsa.Scripting.JavaScript.Services; using Elsa.Services; using Elsa.Services.Models; diff --git a/src/activities/Elsa.Activities.Workflows/Extensions/WorkflowInvokerExtensions.cs b/src/activities/Elsa.Activities.Workflows/Extensions/WorkflowInvokerExtensions.cs new file mode 100644 index 000000000..5c6db4482 --- /dev/null +++ b/src/activities/Elsa.Activities.Workflows/Extensions/WorkflowInvokerExtensions.cs @@ -0,0 +1,34 @@ +using System; +using System.Collections.Generic; +using System.Threading; +using System.Threading.Tasks; +using Elsa.Activities.Workflows.Activities; +using Elsa.Models; +using Elsa.Services; +using Newtonsoft.Json.Linq; + +namespace Elsa.Activities.Workflows.Extensions +{ + public static class WorkflowInvokerExtensions + { + public static async Task TriggerSignalAsync( + this IWorkflowInvoker workflowInvoker, + string signalName, + Variables input = default, + Func activityStatePredicate = null, + string correlationId = default, + CancellationToken cancellationToken = default) + { + var combinedInput = new Variables( + new Dictionary + { + ["Signal"] = signalName + }); + + if (input != null) + combinedInput.AddVariables(input); + + await workflowInvoker.TriggerAsync(nameof(Signaled), combinedInput, correlationId, activityStatePredicate, cancellationToken); + } + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Workflows/Extensions/WorkflowActivityServiceCollectionExtensions.cs b/src/activities/Elsa.Activities.Workflows/Extensions/WorkflowsServiceCollectionExtensions.cs similarity index 89% rename from src/activities/Elsa.Activities.Workflows/Extensions/WorkflowActivityServiceCollectionExtensions.cs rename to src/activities/Elsa.Activities.Workflows/Extensions/WorkflowsServiceCollectionExtensions.cs index 898cbd275..3a112a4dc 100644 --- a/src/activities/Elsa.Activities.Workflows/Extensions/WorkflowActivityServiceCollectionExtensions.cs +++ b/src/activities/Elsa.Activities.Workflows/Extensions/WorkflowsServiceCollectionExtensions.cs @@ -3,7 +3,7 @@ using Microsoft.Extensions.DependencyInjection; namespace Elsa.Activities.Workflows.Extensions { - public static class WorkflowActivityServiceCollectionExtensions + public static class WorkflowsServiceCollectionExtensions { public static IServiceCollection AddWorkflowActivities(this IServiceCollection services) { diff --git a/src/core/Elsa.Core/Extensions/WorkflowExpressionEvaluatorExtensions.cs b/src/core/Elsa.Abstractions/Extensions/WorkflowExpressionEvaluatorExtensions.cs similarity index 56% rename from src/core/Elsa.Core/Extensions/WorkflowExpressionEvaluatorExtensions.cs rename to src/core/Elsa.Abstractions/Extensions/WorkflowExpressionEvaluatorExtensions.cs index 7d9b00578..a440d7cd4 100644 --- a/src/core/Elsa.Core/Extensions/WorkflowExpressionEvaluatorExtensions.cs +++ b/src/core/Elsa.Abstractions/Extensions/WorkflowExpressionEvaluatorExtensions.cs @@ -1,16 +1,20 @@ -using System.Threading; -using System.Threading.Tasks; -using Elsa.Expressions; -using Elsa.Services; -using Elsa.Services.Models; - -namespace Elsa.Extensions -{ - public static class WorkflowExpressionEvaluatorExtensions - { - public static async Task EvaluateAsync(this IWorkflowExpressionEvaluator evaluator, IWorkflowExpression expression, WorkflowExecutionContext workflowExecutionContext, CancellationToken cancellationToken) - { - return (T)await evaluator.EvaluateAsync(expression, typeof(T), workflowExecutionContext, cancellationToken); - } - } +using System.Threading; +using System.Threading.Tasks; +using Elsa.Expressions; +using Elsa.Services; +using Elsa.Services.Models; + +namespace Elsa.Extensions +{ + public static class WorkflowExpressionEvaluatorExtensions + { + public static async Task EvaluateAsync( + this IWorkflowExpressionEvaluator evaluator, + IWorkflowExpression expression, + WorkflowExecutionContext workflowExecutionContext, + CancellationToken cancellationToken = default) + { + return (T)await evaluator.EvaluateAsync(expression, typeof(T), workflowExecutionContext, cancellationToken); + } + } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Models/Variables.cs b/src/core/Elsa.Abstractions/Models/Variables.cs index cce67a44f..f17e47fdc 100644 --- a/src/core/Elsa.Abstractions/Models/Variables.cs +++ b/src/core/Elsa.Abstractions/Models/Variables.cs @@ -10,7 +10,7 @@ namespace Elsa.Models { } - public Variables(Variables other) : this((IEnumerable>) other) + public Variables(Variables other) : this((IEnumerable>)other) { } @@ -29,7 +29,21 @@ namespace Elsa.Models public T GetVariable(string name) { - return ContainsKey(name) ? (T) this[name] : default(T); + return ContainsKey(name) ? (T)this[name] : default(T); + } + + public void AddVariable(string name, object value) + { + this[name] = value; + } + + public void AddVariables(Variables variables) => + AddVariables((IEnumerable>)variables); + + public void AddVariables(IEnumerable> variables) + { + foreach (var variable in variables) + AddVariable(variable.Key, variable.Value); } public bool HasVariable(string name, object value) diff --git a/src/core/Elsa.Abstractions/Services/IScopedWorkflowInvoker.cs b/src/core/Elsa.Abstractions/Services/IScopedWorkflowInvoker.cs deleted file mode 100644 index ce78b980a..000000000 --- a/src/core/Elsa.Abstractions/Services/IScopedWorkflowInvoker.cs +++ /dev/null @@ -1,7 +0,0 @@ -namespace Elsa.Services -{ - internal interface IScopedWorkflowInvoker : IWorkflowInvoker - { - - } -} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs b/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs index 592386793..071109c33 100644 --- a/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs +++ b/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs @@ -1,7 +1,12 @@ using System; using System.Collections.Generic; using System.Linq; +using System.Threading; +using System.Threading.Tasks; +using Elsa.Expressions; +using Elsa.Extensions; using Elsa.Models; +using Microsoft.Extensions.DependencyInjection; using NodaTime; namespace Elsa.Services.Models @@ -20,6 +25,7 @@ namespace Elsa.Services.Models IsFirstPass = true; scheduledActivities = new Stack(); scheduledHaltingActivities = new Stack(); + ExpressionEvaluator = serviceProvider.GetRequiredService(); } public Workflow Workflow { get; } @@ -53,6 +59,7 @@ namespace Elsa.Services.Models public IActivity PopScheduledActivity() => CurrentActivity = scheduledActivities.Pop(); public void ScheduleHaltingActivity(IActivity activity) => scheduledHaltingActivities.Push(activity); public IActivity PopScheduledHaltingActivity() => scheduledHaltingActivities.Pop(); + public IWorkflowExpressionEvaluator ExpressionEvaluator { get; } public void SetVariable(string name, object value) { @@ -75,6 +82,9 @@ namespace Elsa.Services.Models return scope.GetVariable(name); } + public Task EvaluateAsync(IWorkflowExpression expression, CancellationToken cancellationToken) => + ExpressionEvaluator.EvaluateAsync(expression, this, cancellationToken); + public void SetLastResult(object value) => CurrentScope.LastResult = value; public void Start() diff --git a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs index 5e089837d..6f1a8206b 100644 --- a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs +++ b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs @@ -1,6 +1,5 @@ using System; using System.Collections.Generic; -using System.Linq; using Elsa; using Elsa.Activities; using Elsa.AutoMapper.Extensions; @@ -30,8 +29,8 @@ namespace Microsoft.Extensions.DependencyInjection Action configure = null) { var configuration = new ElsaBuilder(services); - configuration.UseWorkflowsCore(); - configuration.UseMediatR(); + configuration.AddWorkflowsCore(); + configuration.AddMediatR(); configure?.Invoke(configuration); EnsurePersistence(configuration); EnsureCaching(configuration); @@ -53,51 +52,14 @@ namespace Microsoft.Extensions.DependencyInjection .AddTransient(sp => sp.GetRequiredService()); } - /// - /// Registers the specified service only if none already exists for the specified provider type. - /// - public static IServiceCollection TryAddProvider( - this IServiceCollection services, - ServiceLifetime lifetime) + private static IServiceCollection AddMediatR(this ElsaBuilder configuration) { - return services.TryAddProvider(typeof(TService), typeof(TProvider), lifetime); + return configuration.Services.AddMediatR( + mediatr => mediatr.AsSingleton(), + typeof(ElsaServiceCollectionExtensions)); } - /// - /// Registers the specified service only if none already exists for the specified provider type. - /// - public static IServiceCollection TryAddProvider( - this IServiceCollection services, - Type serviceType, - Type providerType, - ServiceLifetime lifetime) - { - var descriptor = services.FirstOrDefault( - x => x.ServiceType == serviceType && x.ImplementationType == providerType - ); - - if (descriptor == null) - { - descriptor = new ServiceDescriptor(serviceType, providerType, lifetime); - services.Add(descriptor); - } - - return services; - } - - public static IServiceCollection Replace( - this IServiceCollection services, - ServiceLifetime lifetime) - { - return services.Replace(new ServiceDescriptor(typeof(TService), typeof(TImplementation), lifetime)); - } - - private static IServiceCollection UseMediatR(this ElsaBuilder configuration) - { - return configuration.Services.AddMediatR(typeof(ElsaServiceCollectionExtensions)); - } - - private static ElsaBuilder UseWorkflowsCore(this ElsaBuilder configuration) + private static ElsaBuilder AddWorkflowsCore(this ElsaBuilder configuration) { var services = configuration.Services; services.TryAddSingleton(SystemClock.Instance); @@ -112,14 +74,13 @@ namespace Microsoft.Extensions.DependencyInjection .TryAddProvider(ServiceLifetime.Singleton) .TryAddProvider(ServiceLifetime.Singleton) .TryAddProvider(ServiceLifetime.Singleton) - .AddScoped() - .AddSingleton() + .AddTransient() + .AddScoped() .AddScoped() .AddSingleton() - .AddSingleton() + .AddTransient() .AddScoped() - .AddSingleton() - .AddScoped() + .AddScoped() .AddScoped() .AddScoped() .AddTransient() @@ -151,12 +112,6 @@ namespace Microsoft.Extensions.DependencyInjection configuration.Services.AddMemoryCache(); } - private static bool HasService(this IServiceCollection services) => - services.Any(x => x.ServiceType == typeof(T)); - - private static bool HasService(this ElsaBuilder configuration) => - configuration.Services.HasService(); - private static IServiceCollection AddPrimitiveActivities(this IServiceCollection services) { return services diff --git a/src/core/Elsa.Core/Extensions/MessageHandlerServiceCollectionExtensions.cs b/src/core/Elsa.Core/Extensions/MessageHandlerServiceCollectionExtensions.cs new file mode 100644 index 000000000..2f77ef12c --- /dev/null +++ b/src/core/Elsa.Core/Extensions/MessageHandlerServiceCollectionExtensions.cs @@ -0,0 +1,26 @@ +using System; +using System.Linq; +using System.Reflection; +using MediatR; +using MediatR.Registration; +using Microsoft.Extensions.DependencyInjection; + +namespace Elsa.Extensions +{ + public static class MessageHandlerServiceCollectionExtensions + { + public static IServiceCollection AddNotificationHandler(this IServiceCollection services) + where T : INotification + where THandler : INotificationHandler + { + return services.AddTransient(typeof(INotificationHandler), typeof(THandler)); + } + + public static IServiceCollection AddNotificationHandlers(this IServiceCollection services, params Type[] markerTypes) + { + var assemblies = markerTypes.Select(x => x.GetTypeInfo().Assembly); + ServiceRegistrar.AddMediatRClasses(services, assemblies); + return services; + } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Extensions/ServiceCollectionExtensions.cs b/src/core/Elsa.Core/Extensions/ServiceCollectionExtensions.cs new file mode 100644 index 000000000..9e93832cb --- /dev/null +++ b/src/core/Elsa.Core/Extensions/ServiceCollectionExtensions.cs @@ -0,0 +1,56 @@ +using System; +using System.Linq; +using Elsa; +using Microsoft.Extensions.DependencyInjection.Extensions; + +// ReSharper disable once CheckNamespace +namespace Microsoft.Extensions.DependencyInjection +{ + public static class ServiceCollectionExtensions + { + /// + /// Registers the specified service only if none already exists for the specified provider type. + /// + public static IServiceCollection TryAddProvider( + this IServiceCollection services, + ServiceLifetime lifetime) + { + return services.TryAddProvider(typeof(TService), typeof(TProvider), lifetime); + } + + /// + /// Registers the specified service only if none already exists for the specified provider type. + /// + public static IServiceCollection TryAddProvider( + this IServiceCollection services, + Type serviceType, + Type providerType, + ServiceLifetime lifetime) + { + var descriptor = services.FirstOrDefault( + x => x.ServiceType == serviceType && x.ImplementationType == providerType + ); + + if (descriptor == null) + { + descriptor = new ServiceDescriptor(serviceType, providerType, lifetime); + services.Add(descriptor); + } + + return services; + } + + public static IServiceCollection Replace( + this IServiceCollection services, + ServiceLifetime lifetime) + { + return services.Replace(new ServiceDescriptor(typeof(TService), typeof(TImplementation), lifetime)); + } + + public static bool HasService(this IServiceCollection services) => + services.Any(x => x.ServiceType == typeof(T)); + + public static bool HasService(this ElsaBuilder configuration) => + configuration.Services.HasService(); + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Metadata/ActivityDescriptor.cs b/src/core/Elsa.Core/Metadata/ActivityDescriptor.cs index 0da3dfb49..6c2199d62 100644 --- a/src/core/Elsa.Core/Metadata/ActivityDescriptor.cs +++ b/src/core/Elsa.Core/Metadata/ActivityDescriptor.cs @@ -14,11 +14,11 @@ namespace Elsa.Metadata public string Type { get; set; } public string DisplayName { get; set; } - public string? Description { get; set; } - public string? RuntimeDescription { get; set; } + public string Description { get; set; } + public string RuntimeDescription { get; set; } public string Category { get; set; } - public string? Icon { get; set; } - public object? Outcomes { get; set; } + public string Icon { get; set; } + public object Outcomes { get; set; } public ActivityPropertyDescriptor[] Properties { get; set; } } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Metadata/ActivityPropertyDescriptor.cs b/src/core/Elsa.Core/Metadata/ActivityPropertyDescriptor.cs index 58da3a6bc..ead00365e 100644 --- a/src/core/Elsa.Core/Metadata/ActivityPropertyDescriptor.cs +++ b/src/core/Elsa.Core/Metadata/ActivityPropertyDescriptor.cs @@ -2,7 +2,7 @@ namespace Elsa.Metadata { public class ActivityPropertyDescriptor { - public ActivityPropertyDescriptor(string name, string type, string label, string? hint = null, object? options = null) + public ActivityPropertyDescriptor(string name, string type, string label, string hint = null, object options = null) { Name = name; Type = type; @@ -14,7 +14,7 @@ namespace Elsa.Metadata public string Name { get; } public string Type { get; } public string Label { get; } - public string? Hint { get; } - public object? Options { get; } + public string Hint { get; } + public object Options { get; } } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Services/ScopedWorkflowInvoker.cs b/src/core/Elsa.Core/Services/ScopedWorkflowInvoker.cs deleted file mode 100644 index ca687e634..000000000 --- a/src/core/Elsa.Core/Services/ScopedWorkflowInvoker.cs +++ /dev/null @@ -1,457 +0,0 @@ -using System; -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 Elsa.Services.Extensions; -using Elsa.Services.Models; -using Microsoft.Extensions.Logging; -using Newtonsoft.Json.Linq; -using NodaTime; - -namespace Elsa.Services -{ - internal class ScopedWorkflowInvoker : IScopedWorkflowInvoker - { - private readonly IActivityInvoker activityInvoker; - private readonly IWorkflowFactory workflowFactory; - private readonly IWorkflowRegistry workflowRegistry; - private readonly IWorkflowInstanceStore workflowInstanceStore; - private readonly IEnumerable workflowEventHandlers; - private readonly IClock clock; - private readonly IServiceProvider serviceProvider; - private readonly ILogger logger; - - public ScopedWorkflowInvoker( - IActivityInvoker activityInvoker, - IWorkflowFactory workflowFactory, - IWorkflowRegistry workflowRegistry, - IWorkflowInstanceStore workflowInstanceStore, - IEnumerable workflowEventHandlers, - IClock clock, - IServiceProvider serviceProvider, - ILogger logger) - { - this.activityInvoker = activityInvoker; - this.workflowFactory = workflowFactory; - this.workflowRegistry = workflowRegistry; - this.workflowInstanceStore = workflowInstanceStore; - this.workflowEventHandlers = workflowEventHandlers; - this.clock = clock; - this.serviceProvider = serviceProvider; - this.logger = logger; - } - - public Task StartAsync( - Workflow workflow, - IEnumerable startActivities = default, - CancellationToken cancellationToken = default) - { - return ExecuteAsync(workflow, false, startActivities, cancellationToken); - } - - public Task StartAsync( - WorkflowDefinitionVersion workflowDefinition, - Variables input = default, - IEnumerable startActivityIds = default, - string correlationId = default, - CancellationToken cancellationToken = default) - { - var workflow = workflowFactory.CreateWorkflow(workflowDefinition, input, correlationId: correlationId); - var startActivities = workflow.Activities.Find(startActivityIds); - - return ExecuteAsync(workflow, false, startActivities, cancellationToken); - } - - public Task StartAsync( - Variables input = default, - IEnumerable startActivityIds = default, - string correlationId = default, - CancellationToken cancellationToken = default) where T : IWorkflow, new() - { - var workflow = workflowFactory.CreateWorkflow(input, correlationId: correlationId); - var startActivities = workflow.Activities.Find(startActivityIds); - - return ExecuteAsync(workflow, false, startActivities, cancellationToken); - } - - public Task ResumeAsync( - Workflow workflow, - IEnumerable startActivities = default, - CancellationToken cancellationToken = default) - { - return ExecuteAsync(workflow, true, startActivities, cancellationToken); - } - - public Task ResumeAsync( - WorkflowInstance workflowInstance, - Variables input = null, - IEnumerable startActivityIds = default, - CancellationToken cancellationToken = default) where T : IWorkflow, new() - { - var workflow = workflowFactory.CreateWorkflow(input, workflowInstance); - var startActivities = workflow.Activities.Find(startActivityIds); - return ExecuteAsync(workflow, true, startActivities, cancellationToken); - } - - public async Task ResumeAsync( - WorkflowInstance workflowInstance, - Variables input = null, - IEnumerable startActivityIds = default, - CancellationToken cancellationToken = default) - { - var definition = await workflowRegistry.GetWorkflowDefinitionAsync( - workflowInstance.DefinitionId, - VersionOptions.SpecificVersion(workflowInstance.Version), - cancellationToken); - var workflow = workflowFactory.CreateWorkflow(definition, input, workflowInstance); - return await ExecuteAsync(workflow, true, startActivityIds, cancellationToken); - } - - public async Task> TriggerAsync( - string activityType, - Variables input = default, - string correlationId = default, - Func activityStatePredicate = default, - CancellationToken cancellationToken = default) - { - var startedExecutionContexts = await StartManyAsync( - activityType, - input, - correlationId, - activityStatePredicate, - cancellationToken - ); - - var resumedExecutionContexts = await ResumeManyAsync( - activityType, - input, - correlationId, - activityStatePredicate, - cancellationToken - ); - - return startedExecutionContexts.Concat(resumedExecutionContexts); - } - - private async Task> ResumeManyAsync( - string activityType, - Variables input = default, - string correlationId = default, - Func activityStatePredicate = default, - CancellationToken cancellationToken = default) - { - var workflowInstances = await workflowInstanceStore - .ListByBlockingActivityAsync(activityType, correlationId, cancellationToken) - .ToListAsync(); - - if (activityStatePredicate != null) - workflowInstances = workflowInstances.Where(x => activityStatePredicate(x.Item2.State)).ToList(); - - return await ResumeManyAsync( - workflowInstances, - input, - cancellationToken - ); - } - - private async Task> StartManyAsync( - string activityType, - Variables input = default, - string correlationId = default, - Func activityStatePredicate = default, - CancellationToken cancellationToken = default) - { - var workflowDefinitions = await workflowRegistry.ListByStartActivityAsync(activityType, cancellationToken); - - if (activityStatePredicate != null) - workflowDefinitions = workflowDefinitions.Where(x => activityStatePredicate(x.Item2.State)); - - workflowDefinitions = await FilterRunningSingletonsAsync( - workflowDefinitions, - cancellationToken - ); - - return await StartManyAsync(workflowDefinitions, input, correlationId, cancellationToken); - } - - private Task ExecuteAsync( - Workflow workflow, - bool resume, - IEnumerable startActivityIds = default, - CancellationToken cancellationToken = default) - { - var startActivities = startActivityIds != null - ? workflow.Activities.Find(startActivityIds) - : Enumerable.Empty(); - - return ExecuteAsync(workflow, resume, startActivities, cancellationToken); - } - - private async Task ExecuteAsync( - Workflow workflow, - bool resume, - IEnumerable startActivities = default, - CancellationToken cancellationToken = default) - { - var workflowExecutionContext = await CreateWorkflowExecutionContextAsync( - workflow, - startActivities, - cancellationToken - ); - - var start = !resume; - - while (workflowExecutionContext.HasScheduledActivities) - { - var currentActivity = workflowExecutionContext.PopScheduledActivity(); - - var result = start - ? await ExecuteActivityAsync(workflowExecutionContext, currentActivity, cancellationToken) - : await ResumeActivityAsync(workflowExecutionContext, currentActivity, cancellationToken); - - if (result == null) - break; - - await result.ExecuteAsync(this, workflowExecutionContext, cancellationToken); - - workflowExecutionContext.IsFirstPass = false; - start = true; - } - - await FinalizeWorkflowExecutionAsync(workflowExecutionContext, cancellationToken); - - return workflowExecutionContext; - } - - private async Task> StartManyAsync( - IEnumerable<(WorkflowDefinitionVersion, ActivityDefinition)> workflowDefinitions, - Variables input, - string correlationId, - CancellationToken cancellationToken1) - { - var executionContexts = new List(); - - foreach (var (workflowDefinition, activityDefinition) in workflowDefinitions) - { - var startActivityIds = workflowDefinition.Activities - .Where(x => x.Id == activityDefinition.Id) - .Select(x => x.Id); - - var workflow = workflowFactory.CreateWorkflow(workflowDefinition, input, correlationId: correlationId); - - var executionContext = await ExecuteAsync( - workflow, - false, - startActivityIds, - cancellationToken1 - ); - executionContexts.Add(executionContext); - } - - return executionContexts; - } - - private async Task> ResumeManyAsync( - IEnumerable<(WorkflowInstance, ActivityInstance)> workflowInstances, - Variables input, - CancellationToken cancellationToken) - { - var executionContexts = new List(); - var workflowInstanceGroups = workflowInstances.GroupBy(x => x.Item1); - - foreach (var workflowInstanceGroup in workflowInstanceGroups) - { - var workflowInstance = workflowInstanceGroup.Key; - - var workflowDefinition = await workflowRegistry.GetWorkflowDefinitionAsync( - workflowInstance.DefinitionId, - VersionOptions.SpecificVersion(workflowInstance.Version), - cancellationToken - ); - - var workflow = workflowFactory.CreateWorkflow(workflowDefinition, input, workflowInstance); - - foreach (var activity in workflowInstanceGroup) - { - var executionContext = await ExecuteAsync( - workflow, - true, - new[] { activity.Item2.Id }, - cancellationToken - ); - - executionContexts.Add(executionContext); - } - } - - return executionContexts; - } - - private async Task FinalizeWorkflowExecutionAsync( - WorkflowExecutionContext workflowExecutionContext, - CancellationToken cancellationToken) - { - if (!workflowExecutionContext.Workflow.BlockingActivities.Any() && - workflowExecutionContext.Workflow.IsExecuting()) - { - workflowExecutionContext.Finish(); - } - else - { - // Notify event handlers that halting activities are about to be executed. - await workflowEventHandlers.InvokeAsync( - async x => await x.InvokingHaltedActivitiesAsync(workflowExecutionContext, cancellationToken), - logger - ); - - // 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); - } - } - - // Notify event handlers that workflow execution has ended. - await workflowEventHandlers.InvokeAsync( - async x => await x.WorkflowInvokedAsync(workflowExecutionContext, cancellationToken), - logger - ); - } - - private async Task ExecuteActivityAsync( - WorkflowExecutionContext workflowContext, - IActivity activity, - CancellationToken cancellationToken) - { - return await InvokeActivityAsync( - workflowContext, - activity, - () => activityInvoker.ExecuteAsync(workflowContext, activity, cancellationToken), - cancellationToken - ); - } - - private async Task ResumeActivityAsync( - WorkflowExecutionContext workflowContext, - IActivity activity, - CancellationToken cancellationToken) - { - return await InvokeActivityAsync( - workflowContext, - activity, - () => activityInvoker.ResumeAsync(workflowContext, activity, cancellationToken), - cancellationToken - ); - } - - private async Task InvokeActivityAsync( - 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 executeAction(); - } - catch (Exception ex) - { - FaultWorkflow(workflowContext, activity, ex); - } - - return null; - } - - private async Task ExecuteActivityHaltedAsync( - WorkflowExecutionContext workflowContext, - IActivity activity, - CancellationToken cancellationToken) - { - return await InvokeActivityAsync( - workflowContext, - activity, - () => activityInvoker.HaltedAsync(workflowContext, activity, cancellationToken), - cancellationToken - ); - } - - private void FaultWorkflow(WorkflowExecutionContext workflowContext, IActivity activity, Exception ex) - { - logger.LogError( - ex, - "An unhandled error occurred while executing an activity. Putting the workflow in the faulted state." - ); - workflowContext.Fault(activity, ex); - } - - private async Task CreateWorkflowExecutionContextAsync( - Workflow workflow, - IEnumerable startActivities, - CancellationToken cancellationToken) - { - var workflowExecutionContext = new WorkflowExecutionContext(workflow, clock, serviceProvider); - var startActivityList = startActivities?.ToList() ?? workflow.GetStartActivities().Take(1).ToList(); - - foreach (var startActivity in startActivityList) - { - if (await startActivity.CanExecuteAsync(workflowExecutionContext, cancellationToken)) - workflowExecutionContext.ScheduleActivity(startActivity); - } - - if (workflowExecutionContext.HasScheduledActivities) - { - workflow.BlockingActivities.RemoveWhere(startActivityList.Contains); - - if (workflowExecutionContext.Workflow.Status == WorkflowStatus.Idle) - workflowExecutionContext.Start(); - } - - return workflowExecutionContext; - } - - private async Task> FilterRunningSingletonsAsync( - IEnumerable<(WorkflowDefinitionVersion, ActivityDefinition)> workflowDefinitions, - CancellationToken cancellationToken) - { - var definitions = workflowDefinitions.ToList(); - var transients = definitions.Where(x => !x.Item1.IsSingleton).ToList(); - var singletons = definitions.Where(x => x.Item1.IsSingleton).ToList(); - var result = transients.ToList(); - - foreach (var definition in singletons) - { - var instances = await workflowInstanceStore.ListByStatusAsync( - definition.Item1.DefinitionId, - WorkflowStatus.Executing, - cancellationToken - ); - - if (!instances.Any()) - { - result.Add(definition); - } - } - - return result; - } - } -} \ No newline at end of file diff --git a/src/core/Elsa.Core/Services/WorkflowInvoker.cs b/src/core/Elsa.Core/Services/WorkflowInvoker.cs index 4a0b37fe2..3d21cbda3 100644 --- a/src/core/Elsa.Core/Services/WorkflowInvoker.cs +++ b/src/core/Elsa.Core/Services/WorkflowInvoker.cs @@ -1,21 +1,49 @@ using System; 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 Elsa.Services.Extensions; using Elsa.Services.Models; -using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; using Newtonsoft.Json.Linq; +using NodaTime; namespace Elsa.Services { - public class WorkflowInvoker : IWorkflowInvoker + internal class WorkflowInvoker : IWorkflowInvoker { + private readonly IActivityInvoker activityInvoker; + private readonly IWorkflowFactory workflowFactory; + private readonly IWorkflowRegistry workflowRegistry; + private readonly IWorkflowInstanceStore workflowInstanceStore; + private readonly IEnumerable workflowEventHandlers; + private readonly IClock clock; private readonly IServiceProvider serviceProvider; + private readonly ILogger logger; - public WorkflowInvoker(IServiceProvider serviceProvider) + public WorkflowInvoker( + IActivityInvoker activityInvoker, + IWorkflowFactory workflowFactory, + IWorkflowRegistry workflowRegistry, + IWorkflowInstanceStore workflowInstanceStore, + IEnumerable workflowEventHandlers, + IClock clock, + IServiceProvider serviceProvider, + ILogger logger) { + this.activityInvoker = activityInvoker; + this.workflowFactory = workflowFactory; + this.workflowRegistry = workflowRegistry; + this.workflowInstanceStore = workflowInstanceStore; + this.workflowEventHandlers = workflowEventHandlers; + this.clock = clock; this.serviceProvider = serviceProvider; + this.logger = logger; } public Task StartAsync( @@ -23,16 +51,7 @@ namespace Elsa.Services IEnumerable startActivities = default, CancellationToken cancellationToken = default) { - return Invoke(x => x.StartAsync(workflow, startActivities, cancellationToken)); - } - - public Task StartAsync( - Variables input = default, - IEnumerable startActivityIds = default, - string correlationId = default, - CancellationToken cancellationToken = default) where T : IWorkflow, new() - { - return Invoke(x => x.StartAsync(input, startActivityIds, correlationId, cancellationToken)); + return ExecuteAsync(workflow, false, startActivities, cancellationToken); } public Task StartAsync( @@ -42,9 +61,22 @@ namespace Elsa.Services string correlationId = default, CancellationToken cancellationToken = default) { - return Invoke( - x => x.StartAsync(workflowDefinition, input, startActivityIds, correlationId, cancellationToken) - ); + var workflow = workflowFactory.CreateWorkflow(workflowDefinition, input, correlationId: correlationId); + var startActivities = workflow.Activities.Find(startActivityIds); + + return ExecuteAsync(workflow, false, startActivities, cancellationToken); + } + + public Task StartAsync( + Variables input = default, + IEnumerable startActivityIds = default, + string correlationId = default, + CancellationToken cancellationToken = default) where T : IWorkflow, new() + { + var workflow = workflowFactory.CreateWorkflow(input, correlationId: correlationId); + var startActivities = workflow.Activities.Find(startActivityIds); + + return ExecuteAsync(workflow, false, startActivities, cancellationToken); } public Task ResumeAsync( @@ -52,45 +84,374 @@ namespace Elsa.Services IEnumerable startActivities = default, CancellationToken cancellationToken = default) { - return Invoke(x => x.ResumeAsync(workflow, startActivities, cancellationToken)); + return ExecuteAsync(workflow, true, startActivities, cancellationToken); } public Task ResumeAsync( WorkflowInstance workflowInstance, - Variables input = default, + Variables input = null, IEnumerable startActivityIds = default, - CancellationToken cancellationToken = default) - where T : IWorkflow, new() + CancellationToken cancellationToken = default) where T : IWorkflow, new() { - return Invoke(x => x.ResumeAsync(workflowInstance, input, startActivityIds, cancellationToken)); + var workflow = workflowFactory.CreateWorkflow(input, workflowInstance); + var startActivities = workflow.Activities.Find(startActivityIds); + return ExecuteAsync(workflow, true, startActivities, cancellationToken); } - public Task ResumeAsync( + public async Task ResumeAsync( WorkflowInstance workflowInstance, - Variables input = default, + Variables input = null, IEnumerable startActivityIds = default, CancellationToken cancellationToken = default) { - return Invoke(x => x.ResumeAsync(workflowInstance, input, startActivityIds, cancellationToken)); + var definition = await workflowRegistry.GetWorkflowDefinitionAsync( + workflowInstance.DefinitionId, + VersionOptions.SpecificVersion(workflowInstance.Version), + cancellationToken); + var workflow = workflowFactory.CreateWorkflow(definition, input, workflowInstance); + return await ExecuteAsync(workflow, true, startActivityIds, cancellationToken); } - public Task> TriggerAsync( + public async Task> TriggerAsync( string activityType, Variables input = default, string correlationId = default, Func activityStatePredicate = default, CancellationToken cancellationToken = default) { - return Invoke(x => x.TriggerAsync(activityType, input, correlationId, activityStatePredicate, cancellationToken)); + var startedExecutionContexts = await StartManyAsync( + activityType, + input, + correlationId, + activityStatePredicate, + cancellationToken + ); + + var resumedExecutionContexts = await ResumeManyAsync( + activityType, + input, + correlationId, + activityStatePredicate, + cancellationToken + ); + + return startedExecutionContexts.Concat(resumedExecutionContexts); } - private async Task Invoke(Func> action) + private async Task> ResumeManyAsync( + string activityType, + Variables input = default, + string correlationId = default, + Func activityStatePredicate = default, + CancellationToken cancellationToken = default) { - using (var scope = serviceProvider.CreateScope()) + var workflowInstances = await workflowInstanceStore + .ListByBlockingActivityAsync(activityType, correlationId, cancellationToken) + .ToListAsync(); + + if (activityStatePredicate != null) + workflowInstances = workflowInstances.Where(x => activityStatePredicate(x.Item2.State)).ToList(); + + return await ResumeManyAsync( + workflowInstances, + input, + cancellationToken + ); + } + + private async Task> StartManyAsync( + string activityType, + Variables input = default, + string correlationId = default, + Func activityStatePredicate = default, + CancellationToken cancellationToken = default) + { + var workflowDefinitions = await workflowRegistry.ListByStartActivityAsync(activityType, cancellationToken); + + if (activityStatePredicate != null) + workflowDefinitions = workflowDefinitions.Where(x => activityStatePredicate(x.Item2.State)); + + workflowDefinitions = await FilterRunningSingletonsAsync( + workflowDefinitions, + cancellationToken + ); + + return await StartManyAsync(workflowDefinitions, input, correlationId, cancellationToken); + } + + private Task ExecuteAsync( + Workflow workflow, + bool resume, + IEnumerable startActivityIds = default, + CancellationToken cancellationToken = default) + { + var startActivities = startActivityIds != null + ? workflow.Activities.Find(startActivityIds) + : Enumerable.Empty(); + + return ExecuteAsync(workflow, resume, startActivities, cancellationToken); + } + + private async Task ExecuteAsync( + Workflow workflow, + bool resume, + IEnumerable startActivities = default, + CancellationToken cancellationToken = default) + { + var workflowExecutionContext = await CreateWorkflowExecutionContextAsync( + workflow, + startActivities, + cancellationToken + ); + + var start = !resume; + + while (workflowExecutionContext.HasScheduledActivities) { - var invoker = scope.ServiceProvider.GetRequiredService(); - return await action(invoker); + var currentActivity = workflowExecutionContext.PopScheduledActivity(); + + var result = start + ? await ExecuteActivityAsync(workflowExecutionContext, currentActivity, cancellationToken) + : await ResumeActivityAsync(workflowExecutionContext, currentActivity, cancellationToken); + + if (result == null) + break; + + await result.ExecuteAsync(this, workflowExecutionContext, cancellationToken); + + workflowExecutionContext.IsFirstPass = false; + start = true; } + + await FinalizeWorkflowExecutionAsync(workflowExecutionContext, cancellationToken); + + return workflowExecutionContext; + } + + private async Task> StartManyAsync( + IEnumerable<(WorkflowDefinitionVersion, ActivityDefinition)> workflowDefinitions, + Variables input, + string correlationId, + CancellationToken cancellationToken1) + { + var executionContexts = new List(); + + foreach (var (workflowDefinition, activityDefinition) in workflowDefinitions) + { + var startActivityIds = workflowDefinition.Activities + .Where(x => x.Id == activityDefinition.Id) + .Select(x => x.Id); + + var workflow = workflowFactory.CreateWorkflow(workflowDefinition, input, correlationId: correlationId); + + var executionContext = await ExecuteAsync( + workflow, + false, + startActivityIds, + cancellationToken1 + ); + executionContexts.Add(executionContext); + } + + return executionContexts; + } + + private async Task> ResumeManyAsync( + IEnumerable<(WorkflowInstance, ActivityInstance)> workflowInstances, + Variables input, + CancellationToken cancellationToken) + { + var executionContexts = new List(); + var workflowInstanceGroups = workflowInstances.GroupBy(x => x.Item1); + + foreach (var workflowInstanceGroup in workflowInstanceGroups) + { + var workflowInstance = workflowInstanceGroup.Key; + + var workflowDefinition = await workflowRegistry.GetWorkflowDefinitionAsync( + workflowInstance.DefinitionId, + VersionOptions.SpecificVersion(workflowInstance.Version), + cancellationToken + ); + + var workflow = workflowFactory.CreateWorkflow(workflowDefinition, input, workflowInstance); + + foreach (var activity in workflowInstanceGroup) + { + var executionContext = await ExecuteAsync( + workflow, + true, + new[] { activity.Item2.Id }, + cancellationToken + ); + + executionContexts.Add(executionContext); + } + } + + return executionContexts; + } + + private async Task FinalizeWorkflowExecutionAsync( + WorkflowExecutionContext workflowExecutionContext, + CancellationToken cancellationToken) + { + if (!workflowExecutionContext.Workflow.BlockingActivities.Any() && + workflowExecutionContext.Workflow.IsExecuting()) + { + workflowExecutionContext.Finish(); + } + else + { + // Notify event handlers that halting activities are about to be executed. + await workflowEventHandlers.InvokeAsync( + async x => await x.InvokingHaltedActivitiesAsync(workflowExecutionContext, cancellationToken), + logger + ); + + // 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); + } + } + + // Notify event handlers that workflow execution has ended. + await workflowEventHandlers.InvokeAsync( + async x => await x.WorkflowInvokedAsync(workflowExecutionContext, cancellationToken), + logger + ); + } + + private async Task ExecuteActivityAsync( + WorkflowExecutionContext workflowContext, + IActivity activity, + CancellationToken cancellationToken) + { + return await InvokeActivityAsync( + workflowContext, + activity, + async () => await activityInvoker.ExecuteAsync(workflowContext, activity, cancellationToken), + cancellationToken + ); + } + + private async Task ResumeActivityAsync( + WorkflowExecutionContext workflowContext, + IActivity activity, + CancellationToken cancellationToken) + { + return await InvokeActivityAsync( + workflowContext, + activity, + async () => await activityInvoker.ResumeAsync(workflowContext, activity, cancellationToken), + cancellationToken + ); + } + + private async Task InvokeActivityAsync( + 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 executeAction(); + } + catch (Exception ex) + { + FaultWorkflow(workflowContext, activity, ex); + } + + return null; + } + + private async Task ExecuteActivityHaltedAsync( + WorkflowExecutionContext workflowContext, + IActivity activity, + CancellationToken cancellationToken) + { + return await InvokeActivityAsync( + workflowContext, + activity, + async () => await activityInvoker.HaltedAsync(workflowContext, activity, cancellationToken), + cancellationToken + ); + } + + private void FaultWorkflow(WorkflowExecutionContext workflowContext, IActivity activity, Exception ex) + { + logger.LogError( + ex, + "An unhandled error occurred while executing an activity. Putting the workflow in the faulted state." + ); + workflowContext.Fault(activity, ex); + } + + private async Task CreateWorkflowExecutionContextAsync( + Workflow workflow, + IEnumerable startActivities, + CancellationToken cancellationToken) + { + var workflowExecutionContext = new WorkflowExecutionContext(workflow, clock, serviceProvider); + var startActivityList = startActivities?.ToList() ?? workflow.GetStartActivities().Take(1).ToList(); + + foreach (var startActivity in startActivityList) + { + if (await startActivity.CanExecuteAsync(workflowExecutionContext, cancellationToken)) + workflowExecutionContext.ScheduleActivity(startActivity); + } + + if (workflowExecutionContext.HasScheduledActivities) + { + workflow.BlockingActivities.RemoveWhere(startActivityList.Contains); + + if (workflowExecutionContext.Workflow.Status == WorkflowStatus.Idle) + workflowExecutionContext.Start(); + } + + return workflowExecutionContext; + } + + private async Task> FilterRunningSingletonsAsync( + IEnumerable<(WorkflowDefinitionVersion, ActivityDefinition)> workflowDefinitions, + CancellationToken cancellationToken) + { + var definitions = workflowDefinitions.ToList(); + var transients = definitions.Where(x => !x.Item1.IsSingleton).ToList(); + var singletons = definitions.Where(x => x.Item1.IsSingleton).ToList(); + var result = transients.ToList(); + + foreach (var definition in singletons) + { + var instances = await workflowInstanceStore.ListByStatusAsync( + definition.Item1.DefinitionId, + WorkflowStatus.Executing, + cancellationToken + ); + + if (!instances.Any()) + { + result.Add(definition); + } + } + + return result; } } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Services/WorkflowRegistry.cs b/src/core/Elsa.Core/Services/WorkflowRegistry.cs index be96099c2..685000dd3 100644 --- a/src/core/Elsa.Core/Services/WorkflowRegistry.cs +++ b/src/core/Elsa.Core/Services/WorkflowRegistry.cs @@ -19,13 +19,13 @@ namespace Elsa.Services private readonly ISignal signal; public WorkflowRegistry( - IServiceProvider serviceProvider, IMemoryCache cache, - ISignal signal) + ISignal signal, + IServiceProvider serviceProvider) { - this.serviceProvider = serviceProvider; this.cache = cache; this.signal = signal; + this.serviceProvider = serviceProvider; } public async Task> ListByStartActivityAsync( @@ -72,8 +72,7 @@ namespace Elsa.Services }); } - private async Task> LoadWorkflowDefinitionsAsync( - CancellationToken cancellationToken) + private async Task> LoadWorkflowDefinitionsAsync(CancellationToken cancellationToken) { using var scope = serviceProvider.CreateScope(); var providers = scope.ServiceProvider.GetServices(); diff --git a/src/core/Elsa/ElsaServiceCollectionExtensions.cs b/src/core/Elsa/ElsaServiceCollectionExtensions.cs index 7bfc2ac35..0f5174069 100644 --- a/src/core/Elsa/ElsaServiceCollectionExtensions.cs +++ b/src/core/Elsa/ElsaServiceCollectionExtensions.cs @@ -3,9 +3,7 @@ using Elsa; using Elsa.Activities.ControlFlow.Extensions; using Elsa.Activities.UserTask.Extensions; using Elsa.Activities.Workflows.Extensions; -using Elsa.Scripting.JavaScript; using Elsa.Scripting.JavaScript.Extensions; -using Elsa.Scripting.Liquid; using Elsa.Scripting.Liquid.Extensions; // ReSharper disable once CheckNamespace diff --git a/src/dashboard/Elsa.Dashboard/Extensions/ActivityDefinitionListExtensions.cs b/src/dashboard/Elsa.Dashboard/Extensions/ActivityDefinitionListExtensions.cs index 9669b422b..80153c0e7 100644 --- a/src/dashboard/Elsa.Dashboard/Extensions/ActivityDefinitionListExtensions.cs +++ b/src/dashboard/Elsa.Dashboard/Extensions/ActivityDefinitionListExtensions.cs @@ -2,7 +2,6 @@ using System; using Elsa.Dashboard.Options; using Elsa.Metadata; using Elsa.Services.Models; -using Elsa.WorkflowDesigner; using Microsoft.Extensions.DependencyInjection; using Scrutor; diff --git a/src/dashboard/Elsa.Dashboard/Options/ActivityDefinitionList.cs b/src/dashboard/Elsa.Dashboard/Options/ActivityDefinitionList.cs index 1f1f26f99..f98ee40ea 100644 --- a/src/dashboard/Elsa.Dashboard/Options/ActivityDefinitionList.cs +++ b/src/dashboard/Elsa.Dashboard/Options/ActivityDefinitionList.cs @@ -1,7 +1,6 @@ using System.Collections; using System.Collections.Generic; using Elsa.Metadata; -using Elsa.WorkflowDesigner.Models; namespace Elsa.Dashboard.Options { diff --git a/src/persistence/Elsa.Persistence.MongoDb/Extensions/ServiceCollectionExtensions.cs b/src/persistence/Elsa.Persistence.MongoDb/Extensions/ServiceCollectionExtensions.cs index 23b047581..b282c1c15 100644 --- a/src/persistence/Elsa.Persistence.MongoDb/Extensions/ServiceCollectionExtensions.cs +++ b/src/persistence/Elsa.Persistence.MongoDb/Extensions/ServiceCollectionExtensions.cs @@ -25,7 +25,6 @@ namespace Elsa.Persistence.MongoDb.Extensions RegisterEnumAsStringConvention(); BsonSerializer.RegisterSerializer(new JObjectSerializer()); BsonSerializer.RegisterSerializer(new WorkflowExecutionScopeSerializer()); - //BsonSerializer.RegisterSerializer(new WorkflowInstanceSerializer()); elsaBuilder.Services .AddSingleton(sp => CreateDbClient(configuration, connectionStringName)) @@ -51,7 +50,7 @@ namespace Elsa.Persistence.MongoDb.Extensions { configuration.Services .AddMongoDbCollection("WorkflowInstances") - .Replace(ServiceLifetime.Scoped); + .AddScoped(); return configuration; } @@ -61,7 +60,7 @@ namespace Elsa.Persistence.MongoDb.Extensions { configuration.Services .AddMongoDbCollection("WorkflowDefinitions") - .Replace(ServiceLifetime.Scoped); + .AddScoped(); return configuration; } diff --git a/src/scripting/Elsa.Scripting.JavaScript/Extensions/JavaScriptServiceCollectionExtensions.cs b/src/scripting/Elsa.Scripting.JavaScript/Extensions/JavaScriptServiceCollectionExtensions.cs index 3e43bee6d..442030211 100644 --- a/src/scripting/Elsa.Scripting.JavaScript/Extensions/JavaScriptServiceCollectionExtensions.cs +++ b/src/scripting/Elsa.Scripting.JavaScript/Extensions/JavaScriptServiceCollectionExtensions.cs @@ -1,6 +1,6 @@ +using Elsa.Extensions; using Elsa.Scripting.JavaScript.Services; using Elsa.Services; -using MediatR; using Microsoft.Extensions.DependencyInjection; namespace Elsa.Scripting.JavaScript.Extensions @@ -11,7 +11,7 @@ namespace Elsa.Scripting.JavaScript.Extensions { return services .TryAddProvider(ServiceLifetime.Scoped) - .AddMediatR(typeof(JavaScriptServiceCollectionExtensions)); + .AddNotificationHandlers(typeof(JavaScriptServiceCollectionExtensions)); } } } \ No newline at end of file diff --git a/src/scripting/Elsa.Scripting.JavaScript/Services/JavaScriptExpressionEvaluator.cs b/src/scripting/Elsa.Scripting.JavaScript/Services/JavaScriptExpressionEvaluator.cs index 7d27745b2..8eec73e5c 100644 --- a/src/scripting/Elsa.Scripting.JavaScript/Services/JavaScriptExpressionEvaluator.cs +++ b/src/scripting/Elsa.Scripting.JavaScript/Services/JavaScriptExpressionEvaluator.cs @@ -1,5 +1,4 @@ using System; -using System.Collections.Generic; using System.Linq; using System.Threading; using System.Threading.Tasks; diff --git a/src/scripting/Elsa.Scripting.Liquid/Extensions/LiquidServiceCollectionExtensions.cs b/src/scripting/Elsa.Scripting.Liquid/Extensions/LiquidServiceCollectionExtensions.cs index 7eb350a39..9c422c0ce 100644 --- a/src/scripting/Elsa.Scripting.Liquid/Extensions/LiquidServiceCollectionExtensions.cs +++ b/src/scripting/Elsa.Scripting.Liquid/Extensions/LiquidServiceCollectionExtensions.cs @@ -1,8 +1,8 @@ +using Elsa.Extensions; using Elsa.Scripting.Liquid.Filters; using Elsa.Scripting.Liquid.Options; using Elsa.Scripting.Liquid.Services; using Elsa.Services; -using MediatR; using Microsoft.Extensions.DependencyInjection; namespace Elsa.Scripting.Liquid.Extensions @@ -14,7 +14,7 @@ namespace Elsa.Scripting.Liquid.Extensions return services .TryAddProvider(ServiceLifetime.Scoped) .AddMemoryCache() - .AddMediatR(typeof(LiquidServiceCollectionExtensions)) + .AddNotificationHandlers(typeof(LiquidServiceCollectionExtensions)) .AddScoped() .AddLiquidFilter("json"); } diff --git a/src/scripting/Elsa.Scripting.Liquid/Handlers/CommonLiquidContextHandler.cs b/src/scripting/Elsa.Scripting.Liquid/Handlers/CommonLiquidContextHandler.cs index 5cec19685..11fe56191 100644 --- a/src/scripting/Elsa.Scripting.Liquid/Handlers/CommonLiquidContextHandler.cs +++ b/src/scripting/Elsa.Scripting.Liquid/Handlers/CommonLiquidContextHandler.cs @@ -9,6 +9,7 @@ using Elsa.Services.Models; using Fluid; using Fluid.Values; using MediatR; +using Newtonsoft.Json.Linq; namespace Elsa.Scripting.Liquid.Handlers { @@ -17,11 +18,14 @@ namespace Elsa.Scripting.Liquid.Handlers static CommonLiquidContextHandler() { FluidValue.SetTypeMapping(x => new ObjectValue(x)); + FluidValue.SetTypeMapping(o => new ObjectValue(o)); + FluidValue.SetTypeMapping(o => FluidValue.Create(o.Value)); } - + public Task Handle(EvaluatingLiquidExpression notification, CancellationToken cancellationToken) { var context = notification.TemplateContext; + context.MemberAccessStrategy.Register((x, name) => x.GetValueAsync(name)); context.MemberAccessStrategy.Register("Input", x => new LiquidPropertyAccessor(name => ToFluidValue(x.Workflow.Input, name))); context.MemberAccessStrategy.Register("Output", x => new LiquidPropertyAccessor(name => ToFluidValue(x.Workflow.Output, name))); @@ -30,10 +34,11 @@ namespace Elsa.Scripting.Liquid.Handlers context.MemberAccessStrategy.Register, LiquidObjectAccessor>((x, activityName) => new LiquidObjectAccessor(outputKey => GetActivityOutput(x, activityName, outputKey))); context.MemberAccessStrategy.Register, object>((x, name) => x.GetValueAsync(name)); context.MemberAccessStrategy.Register((x, name) => ((IDictionary)x)[name]); - + context.MemberAccessStrategy.Register((source, name) => source[name]); + return Task.CompletedTask; } - + private Task ToFluidValue(IDictionary dictionary, string key) { return Task.FromResult(!dictionary.ContainsKey(key) ? default : FluidValue.Create(dictionary[key])); diff --git a/src/scripting/Elsa.Scripting.Liquid/Services/LiquidExpressionEvaluator.cs b/src/scripting/Elsa.Scripting.Liquid/Services/LiquidExpressionEvaluator.cs index 88ad8e2e9..d7ee3f8e1 100644 --- a/src/scripting/Elsa.Scripting.Liquid/Services/LiquidExpressionEvaluator.cs +++ b/src/scripting/Elsa.Scripting.Liquid/Services/LiquidExpressionEvaluator.cs @@ -34,6 +34,7 @@ namespace Elsa.Scripting.Liquid.Services private async Task CreateTemplateContextAsync(WorkflowExecutionContext workflowContext) { var context = new TemplateContext(); + context.SetValue("WorkflowExecutionContext", workflowContext); await mediator.Publish(new EvaluatingLiquidExpression(context, workflowContext)); context.Model = workflowContext; return context;