From 8423fde504e83b2cb7bf2b18e43e0c8381c6fc37 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Sat, 17 Oct 2020 12:49:41 +0200 Subject: [PATCH] Incremental work on trigger API --- .../Elsa.Activities.Http.csproj | 5 +- .../Extensions/ServiceCollectionExtensions.cs | 3 + .../Middleware/HttpRequestMiddleware.cs | 26 ++++++-- .../Triggers/ReceiveHttpRequestTrigger.cs | 59 +++++++++++++++++++ .../Elsa.Activities.MassTransit.csproj | 8 +-- .../Elsa.Abstractions.csproj | 6 +- .../Models/WorkflowExecutionScope.cs | 1 + .../Handlers/PrimitiveValueHandler.cs | 2 +- .../{Runtime => Services}/IStartupRunner.cs | 2 +- .../{Runtime => Services}/IStartupTask.cs | 2 +- .../Services/IWorkflowHost.cs | 3 + .../Models/WorkflowExecutionContext.cs | 3 - .../Elsa.Abstractions/Triggers/ITrigger.cs | 7 +++ .../Triggers/ITriggerProvider.cs | 14 +++++ .../Triggers/IWorkflowSelector.cs | 16 +++++ .../Triggers/TriggerProvider.cs | 14 +++++ .../Triggers/WorkflowSelectorExtensions.cs | 15 +++++ .../Triggers/WorkflowSelectorResult.cs | 16 +++++ .../Data/Services/DataMigrationsRunner.cs | 1 + .../Data/Services/DatabaseInitializer.cs | 1 + src/core/Elsa.Core/Elsa.Core.csproj | 8 ++- .../ElsaServiceCollectionExtensions.cs | 4 +- .../Runtime/ServiceCollectionExtensions.cs | 1 + src/core/Elsa.Core/Runtime/StartupRunner.cs | 1 + .../Runtime/StartupRunnerHostedService.cs | 1 + src/core/Elsa.Core/Services/WorkflowHost.cs | 32 ++++------ .../StartupTasks/StartServiceBusTask.cs | 1 + .../Elsa.Core/Triggers/WorkflowSelector.cs | 43 ++++++++++++++ .../Elsa.DistributedLocking.AzureBlob.csproj | 2 +- .../Elsa.Samples.DistributedLock.csproj | 2 +- .../Elsa.Samples.HelloWorldConsole.csproj | 2 +- .../Elsa.Samples.Serialization.csproj | 2 +- .../Elsa.Server.GraphQL.csproj | 6 +- .../Elsa.Server.Host/Elsa.Server.Host.csproj | 6 +- .../Elsa.Testing.Shared.csproj | 4 +- .../Elsa.Core.UnitTests.csproj | 6 +- 36 files changed, 264 insertions(+), 61 deletions(-) create mode 100644 src/activities/Elsa.Activities.Http/Triggers/ReceiveHttpRequestTrigger.cs rename src/core/Elsa.Abstractions/{Runtime => Services}/IStartupRunner.cs (88%) rename src/core/Elsa.Abstractions/{Runtime => Services}/IStartupTask.cs (88%) create mode 100644 src/core/Elsa.Abstractions/Triggers/ITrigger.cs create mode 100644 src/core/Elsa.Abstractions/Triggers/ITriggerProvider.cs create mode 100644 src/core/Elsa.Abstractions/Triggers/IWorkflowSelector.cs create mode 100644 src/core/Elsa.Abstractions/Triggers/TriggerProvider.cs create mode 100644 src/core/Elsa.Abstractions/Triggers/WorkflowSelectorExtensions.cs create mode 100644 src/core/Elsa.Abstractions/Triggers/WorkflowSelectorResult.cs create mode 100644 src/core/Elsa.Core/Triggers/WorkflowSelector.cs diff --git a/src/activities/Elsa.Activities.Http/Elsa.Activities.Http.csproj b/src/activities/Elsa.Activities.Http/Elsa.Activities.Http.csproj index a58f90304..b89359848 100644 --- a/src/activities/Elsa.Activities.Http/Elsa.Activities.Http.csproj +++ b/src/activities/Elsa.Activities.Http/Elsa.Activities.Http.csproj @@ -36,11 +36,12 @@ - + - + + \ 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 95e46d978..812bf84f8 100644 --- a/src/activities/Elsa.Activities.Http/Extensions/ServiceCollectionExtensions.cs +++ b/src/activities/Elsa.Activities.Http/Extensions/ServiceCollectionExtensions.cs @@ -5,8 +5,10 @@ using Elsa.Activities.Http.Options; using Elsa.Activities.Http.Parsers; using Elsa.Activities.Http.RequestHandlers.Handlers; using Elsa.Activities.Http.Services; +using Elsa.Activities.Http.Triggers; using Elsa.Data; using Elsa.Extensions; +using Elsa.Triggers; using Microsoft.AspNetCore.Http; using Microsoft.AspNetCore.Mvc.Infrastructure; using Microsoft.Extensions.DependencyInjection.Extensions; @@ -38,6 +40,7 @@ namespace Microsoft.Extensions.DependencyInjection .AddSingleton() .AddSingleton() .AddSingleton() + .AddSingleton() .AddIndexProvider() .AddDataMigration() .AddHttpContextAccessor() diff --git a/src/activities/Elsa.Activities.Http/Middleware/HttpRequestMiddleware.cs b/src/activities/Elsa.Activities.Http/Middleware/HttpRequestMiddleware.cs index f0e3547be..6c00179eb 100644 --- a/src/activities/Elsa.Activities.Http/Middleware/HttpRequestMiddleware.cs +++ b/src/activities/Elsa.Activities.Http/Middleware/HttpRequestMiddleware.cs @@ -1,24 +1,38 @@ +using System.Linq; +using System.Threading; using System.Threading.Tasks; using Elsa.Activities.Http.Services; +using Elsa.Activities.Http.Triggers; +using Elsa.Services; +using Elsa.Triggers; using Microsoft.AspNetCore.Http; namespace Elsa.Activities.Http.Middleware { - public class RequestHandlerMiddleware where THandler : IRequestHandler + public class ReceiveHttpRequestMiddleware { private readonly RequestDelegate _next; - public RequestHandlerMiddleware(RequestDelegate next) + public ReceiveHttpRequestMiddleware(RequestDelegate next) { _next = next; } - public async Task InvokeAsync(HttpContext httpContext, THandler handler) + public async Task InvokeAsync(HttpContext httpContext, IWorkflowSelector workflowSelector, IWorkflowHost workflowHost, CancellationToken cancellationToken) { - var result = await handler.HandleRequestAsync(); + var path = httpContext.Request.Path; + var method = httpContext.Request.Method; + + var results = await workflowSelector.SelectWorkflowsAsync( + x => x.Path == path && x.Method == null || x.Method == method, + cancellationToken).ToListAsync(cancellationToken); - if (result != null && !httpContext.Response.HasStarted) - await result.ExecuteResultAsync(httpContext, _next); + var result = results.FirstOrDefault(); + + if (result == null) + await _next(httpContext); + else + await workflowHost.RunWorkflowAsync(result.WorkflowBlueprint, result.ActivityId, cancellationToken: cancellationToken); } } } \ No newline at end of file diff --git a/src/activities/Elsa.Activities.Http/Triggers/ReceiveHttpRequestTrigger.cs b/src/activities/Elsa.Activities.Http/Triggers/ReceiveHttpRequestTrigger.cs new file mode 100644 index 000000000..453f95630 --- /dev/null +++ b/src/activities/Elsa.Activities.Http/Triggers/ReceiveHttpRequestTrigger.cs @@ -0,0 +1,59 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Runtime.CompilerServices; +using System.Security.Cryptography.Xml; +using System.Threading; +using System.Threading.Tasks; +using Elsa.Expressions; +using Elsa.Services; +using Elsa.Services.Models; +using Elsa.Triggers; +using Microsoft.AspNetCore.Http; +using Microsoft.Extensions.DependencyInjection; + +namespace Elsa.Activities.Http.Triggers +{ + public class ReceiveHttpRequestTrigger : ITrigger + { + public PathString Path { get; set; } + public string? Method { get; set; } + public string ActivityId { get; set; } + } + + public class ReceiveHttpRequestTriggerProvider : TriggerProvider + { + private readonly IServiceProvider _serviceProvider; + private readonly IWorkflowFactory _workflowFactory; + + public ReceiveHttpRequestTriggerProvider(IServiceProvider serviceProvider, IWorkflowFactory workflowFactory) + { + _serviceProvider = serviceProvider; + _workflowFactory = workflowFactory; + } + + public override async IAsyncEnumerable GetTriggersAsync(IWorkflowBlueprint workflowBlueprint, [EnumeratorCancellation] CancellationToken cancellationToken = default) + { + var httpRequestActivities = workflowBlueprint + .GetStartActivities() + .Where(x => x.Type == nameof(ReceiveHttpRequest)) + .ToList(); + + using var scope = _serviceProvider.CreateScope(); + var workflowInstance = await _workflowFactory.InstantiateAsync(workflowBlueprint, cancellationToken: cancellationToken); + var workflowExecutionContext = new WorkflowExecutionContext(scope.ServiceProvider, workflowBlueprint, workflowInstance); + + foreach (var activity in httpRequestActivities) + { + var activityExecutionContext = new ActivityExecutionContext(workflowExecutionContext, scope.ServiceProvider, activity); + + yield return new ReceiveHttpRequestTrigger + { + ActivityId = activity.Id, + Path = (PathString)(await workflowBlueprint.ActivityPropertyProviders.GetProvider(activity.Id, nameof(ReceiveHttpRequest.Path))!.GetValueAsync(activityExecutionContext, cancellationToken))!, + Method = (string)(await workflowBlueprint.ActivityPropertyProviders.GetProvider(activity.Id, nameof(ReceiveHttpRequest.Method))!.GetValueAsync(activityExecutionContext, cancellationToken))!, + }; + } + } + } +} \ No newline at end of file diff --git a/src/activities/Elsa.Activities.MassTransit/Elsa.Activities.MassTransit.csproj b/src/activities/Elsa.Activities.MassTransit/Elsa.Activities.MassTransit.csproj index 40fd1efb7..060b40f00 100644 --- a/src/activities/Elsa.Activities.MassTransit/Elsa.Activities.MassTransit.csproj +++ b/src/activities/Elsa.Activities.MassTransit/Elsa.Activities.MassTransit.csproj @@ -29,10 +29,10 @@ - - - - + + + + diff --git a/src/core/Elsa.Abstractions/Elsa.Abstractions.csproj b/src/core/Elsa.Abstractions/Elsa.Abstractions.csproj index 3b4482bde..7713b8a99 100644 --- a/src/core/Elsa.Abstractions/Elsa.Abstractions.csproj +++ b/src/core/Elsa.Abstractions/Elsa.Abstractions.csproj @@ -27,7 +27,7 @@ - + @@ -35,10 +35,10 @@ - + - + diff --git a/src/core/Elsa.Abstractions/Models/WorkflowExecutionScope.cs b/src/core/Elsa.Abstractions/Models/WorkflowExecutionScope.cs index 436f39b45..e3e717350 100644 --- a/src/core/Elsa.Abstractions/Models/WorkflowExecutionScope.cs +++ b/src/core/Elsa.Abstractions/Models/WorkflowExecutionScope.cs @@ -4,6 +4,7 @@ namespace Elsa.Models { public WorkflowExecutionScope() { + Variables = new Variables(); } public WorkflowExecutionScope(Variables variables) diff --git a/src/core/Elsa.Abstractions/Serialization/Handlers/PrimitiveValueHandler.cs b/src/core/Elsa.Abstractions/Serialization/Handlers/PrimitiveValueHandler.cs index 524df2dab..ddfcdb905 100644 --- a/src/core/Elsa.Abstractions/Serialization/Handlers/PrimitiveValueHandler.cs +++ b/src/core/Elsa.Abstractions/Serialization/Handlers/PrimitiveValueHandler.cs @@ -17,7 +17,7 @@ namespace Elsa.Serialization.Handlers { var valueToken = token["Value"]; - return valueToken == null ? null : ParseValue(valueToken); + return valueToken == null ? null! : ParseValue(valueToken); } public virtual void Serialize(JsonWriter writer, JsonSerializer serializer, Type type, JToken token, object? value) diff --git a/src/core/Elsa.Abstractions/Runtime/IStartupRunner.cs b/src/core/Elsa.Abstractions/Services/IStartupRunner.cs similarity index 88% rename from src/core/Elsa.Abstractions/Runtime/IStartupRunner.cs rename to src/core/Elsa.Abstractions/Services/IStartupRunner.cs index ee26e6981..52089c173 100644 --- a/src/core/Elsa.Abstractions/Runtime/IStartupRunner.cs +++ b/src/core/Elsa.Abstractions/Services/IStartupRunner.cs @@ -1,7 +1,7 @@ using System.Threading; using System.Threading.Tasks; -namespace Elsa.Runtime +namespace Elsa.Services { public interface IStartupRunner { diff --git a/src/core/Elsa.Abstractions/Runtime/IStartupTask.cs b/src/core/Elsa.Abstractions/Services/IStartupTask.cs similarity index 88% rename from src/core/Elsa.Abstractions/Runtime/IStartupTask.cs rename to src/core/Elsa.Abstractions/Services/IStartupTask.cs index bf49f8e42..f1186a3c4 100644 --- a/src/core/Elsa.Abstractions/Runtime/IStartupTask.cs +++ b/src/core/Elsa.Abstractions/Services/IStartupTask.cs @@ -1,7 +1,7 @@ using System.Threading; using System.Threading.Tasks; -namespace Elsa.Runtime +namespace Elsa.Services { public interface IStartupTask { diff --git a/src/core/Elsa.Abstractions/Services/IWorkflowHost.cs b/src/core/Elsa.Abstractions/Services/IWorkflowHost.cs index 33b680f14..626d47a6e 100644 --- a/src/core/Elsa.Abstractions/Services/IWorkflowHost.cs +++ b/src/core/Elsa.Abstractions/Services/IWorkflowHost.cs @@ -1,7 +1,10 @@ +using System; +using System.Collections.Generic; using System.Threading; using System.Threading.Tasks; using Elsa.Models; using Elsa.Services.Models; +using Elsa.Triggers; namespace Elsa.Services { diff --git a/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs b/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs index 135eea7d6..7855426dc 100644 --- a/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs +++ b/src/core/Elsa.Abstractions/Services/Models/WorkflowExecutionContext.cs @@ -12,7 +12,6 @@ namespace Elsa.Services.Models public class WorkflowExecutionContext { public WorkflowExecutionContext( - IExpressionEvaluator expressionEvaluator, IServiceProvider serviceProvider, IWorkflowBlueprint workflowBlueprint, WorkflowInstance workflowInstance @@ -21,7 +20,6 @@ namespace Elsa.Services.Models ServiceProvider = serviceProvider; WorkflowBlueprint = workflowBlueprint; WorkflowInstance = workflowInstance; - ExpressionEvaluator = expressionEvaluator; ExecutionLog = new List(workflowInstance.ExecutionLog); IsFirstPass = true; } @@ -53,7 +51,6 @@ namespace Elsa.Services.Models public void ScheduleActivity(ScheduledActivity activity) => WorkflowInstance.ScheduledActivities.Push(activity); public ScheduledActivity PopScheduledActivity() => WorkflowInstance.ScheduledActivities.Pop(); public ScheduledActivity PeekScheduledActivity() => WorkflowInstance.ScheduledActivities.Peek(); - public IExpressionEvaluator ExpressionEvaluator { get; } public string? CorrelationId { get; set; } public bool DeleteCompletedInstances { get; set; } public ICollection ExecutionLog { get; } diff --git a/src/core/Elsa.Abstractions/Triggers/ITrigger.cs b/src/core/Elsa.Abstractions/Triggers/ITrigger.cs new file mode 100644 index 000000000..eb93fbb47 --- /dev/null +++ b/src/core/Elsa.Abstractions/Triggers/ITrigger.cs @@ -0,0 +1,7 @@ +namespace Elsa.Triggers +{ + public interface ITrigger + { + string ActivityId { get; set; } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Triggers/ITriggerProvider.cs b/src/core/Elsa.Abstractions/Triggers/ITriggerProvider.cs new file mode 100644 index 000000000..e3c48e028 --- /dev/null +++ b/src/core/Elsa.Abstractions/Triggers/ITriggerProvider.cs @@ -0,0 +1,14 @@ +using System; +using System.Collections.Generic; +using System.Threading; +using System.Threading.Tasks; +using Elsa.Services.Models; + +namespace Elsa.Triggers +{ + public interface ITriggerProvider + { + Type ForType(); + IAsyncEnumerable GetTriggersAsync(IWorkflowBlueprint workflowBlueprint, CancellationToken cancellationToken = default); + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Triggers/IWorkflowSelector.cs b/src/core/Elsa.Abstractions/Triggers/IWorkflowSelector.cs new file mode 100644 index 000000000..bc3ea8bff --- /dev/null +++ b/src/core/Elsa.Abstractions/Triggers/IWorkflowSelector.cs @@ -0,0 +1,16 @@ +using System; +using System.Collections.Generic; +using System.Runtime.CompilerServices; +using System.Threading; +using Elsa.Services.Models; + +namespace Elsa.Triggers +{ + public interface IWorkflowSelector + { + IAsyncEnumerable SelectWorkflowsAsync( + Type triggerType, + Func evaluate, + CancellationToken cancellationToken = default); + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Triggers/TriggerProvider.cs b/src/core/Elsa.Abstractions/Triggers/TriggerProvider.cs new file mode 100644 index 000000000..935b977be --- /dev/null +++ b/src/core/Elsa.Abstractions/Triggers/TriggerProvider.cs @@ -0,0 +1,14 @@ +using System; +using System.Collections.Generic; +using System.Threading; +using System.Threading.Tasks; +using Elsa.Services.Models; + +namespace Elsa.Triggers +{ + public abstract class TriggerProvider : ITriggerProvider where T : ITrigger + { + public Type ForType() => typeof(T); + public abstract IAsyncEnumerable GetTriggersAsync(IWorkflowBlueprint workflowBlueprint, CancellationToken cancellationToken = default); + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Triggers/WorkflowSelectorExtensions.cs b/src/core/Elsa.Abstractions/Triggers/WorkflowSelectorExtensions.cs new file mode 100644 index 000000000..b761428c3 --- /dev/null +++ b/src/core/Elsa.Abstractions/Triggers/WorkflowSelectorExtensions.cs @@ -0,0 +1,15 @@ +using System; +using System.Collections.Generic; +using System.Threading; + +namespace Elsa.Triggers +{ + public static class WorkflowSelectorExtensions + { + public static IAsyncEnumerable SelectWorkflowsAsync( + this IWorkflowSelector workflowSelector, + Func evaluate, + CancellationToken cancellationToken = default) where T : ITrigger => + workflowSelector.SelectWorkflowsAsync(typeof(T), t => evaluate((T)t), cancellationToken); + } +} \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Triggers/WorkflowSelectorResult.cs b/src/core/Elsa.Abstractions/Triggers/WorkflowSelectorResult.cs new file mode 100644 index 000000000..6d3558f59 --- /dev/null +++ b/src/core/Elsa.Abstractions/Triggers/WorkflowSelectorResult.cs @@ -0,0 +1,16 @@ +using Elsa.Services.Models; + +namespace Elsa.Triggers +{ + public class WorkflowSelectorResult + { + public WorkflowSelectorResult(IWorkflowBlueprint workflowBlueprint, string activityId) + { + WorkflowBlueprint = workflowBlueprint; + ActivityId = activityId; + } + + public IWorkflowBlueprint WorkflowBlueprint { get; } + public string ActivityId { get; } + } +} \ No newline at end of file diff --git a/src/core/Elsa.Core/Data/Services/DataMigrationsRunner.cs b/src/core/Elsa.Core/Data/Services/DataMigrationsRunner.cs index 5beeb6d08..9066cd947 100644 --- a/src/core/Elsa.Core/Data/Services/DataMigrationsRunner.cs +++ b/src/core/Elsa.Core/Data/Services/DataMigrationsRunner.cs @@ -1,6 +1,7 @@ using System.Threading; using System.Threading.Tasks; using Elsa.Runtime; +using Elsa.Services; namespace Elsa.Data.Services { diff --git a/src/core/Elsa.Core/Data/Services/DatabaseInitializer.cs b/src/core/Elsa.Core/Data/Services/DatabaseInitializer.cs index 3edf8a4b7..497a4c176 100644 --- a/src/core/Elsa.Core/Data/Services/DatabaseInitializer.cs +++ b/src/core/Elsa.Core/Data/Services/DatabaseInitializer.cs @@ -1,6 +1,7 @@ using System.Threading; using System.Threading.Tasks; using Elsa.Runtime; +using Elsa.Services; using YesSql; namespace Elsa.Data.Services diff --git a/src/core/Elsa.Core/Elsa.Core.csproj b/src/core/Elsa.Core/Elsa.Core.csproj index 3e94bfe8e..e44a70989 100644 --- a/src/core/Elsa.Core/Elsa.Core.csproj +++ b/src/core/Elsa.Core/Elsa.Core.csproj @@ -28,11 +28,12 @@ - + - - + + + @@ -55,6 +56,7 @@ + diff --git a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs index 0b5f2baaa..22f1c9fa1 100644 --- a/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs +++ b/src/core/Elsa.Core/Extensions/ElsaServiceCollectionExtensions.cs @@ -21,6 +21,7 @@ using Elsa.Serialization; using Elsa.Serialization.Formatters; using Elsa.Services; using Elsa.StartupTasks; +using Elsa.Triggers; using Elsa.WorkflowProviders; using MediatR; using Microsoft.Extensions.DependencyInjection.Extensions; @@ -34,7 +35,7 @@ namespace Microsoft.Extensions.DependencyInjection { public static IServiceCollection AddElsaCore( this IServiceCollection services, - Action configure = default!) + Action? configure = default) { var options = new ElsaOptions(services); configure?.Invoke(options); @@ -102,6 +103,7 @@ namespace Microsoft.Extensions.DependencyInjection .AddSingleton() .AddScoped() .AddSingleton() + .AddSingleton() .AddScoped() .AddScoped() .AddIndexProvider() diff --git a/src/core/Elsa.Core/Runtime/ServiceCollectionExtensions.cs b/src/core/Elsa.Core/Runtime/ServiceCollectionExtensions.cs index fff0fe559..d5334cabb 100644 --- a/src/core/Elsa.Core/Runtime/ServiceCollectionExtensions.cs +++ b/src/core/Elsa.Core/Runtime/ServiceCollectionExtensions.cs @@ -1,3 +1,4 @@ +using Elsa.Services; using Microsoft.Extensions.DependencyInjection; namespace Elsa.Runtime diff --git a/src/core/Elsa.Core/Runtime/StartupRunner.cs b/src/core/Elsa.Core/Runtime/StartupRunner.cs index 92a9d1389..390f07799 100644 --- a/src/core/Elsa.Core/Runtime/StartupRunner.cs +++ b/src/core/Elsa.Core/Runtime/StartupRunner.cs @@ -1,6 +1,7 @@ using System; using System.Threading; using System.Threading.Tasks; +using Elsa.Services; using Microsoft.Extensions.DependencyInjection; namespace Elsa.Runtime diff --git a/src/core/Elsa.Core/Runtime/StartupRunnerHostedService.cs b/src/core/Elsa.Core/Runtime/StartupRunnerHostedService.cs index 4f43c8999..f39560396 100644 --- a/src/core/Elsa.Core/Runtime/StartupRunnerHostedService.cs +++ b/src/core/Elsa.Core/Runtime/StartupRunnerHostedService.cs @@ -1,5 +1,6 @@ using System.Threading; using System.Threading.Tasks; +using Elsa.Services; using Microsoft.Extensions.Hosting; namespace Elsa.Runtime diff --git a/src/core/Elsa.Core/Services/WorkflowHost.cs b/src/core/Elsa.Core/Services/WorkflowHost.cs index 9bd8ca3de..c0db1b264 100644 --- a/src/core/Elsa.Core/Services/WorkflowHost.cs +++ b/src/core/Elsa.Core/Services/WorkflowHost.cs @@ -1,5 +1,7 @@ using System; +using System.Collections.Generic; using System.Linq; +using System.Runtime.CompilerServices; using System.Threading; using System.Threading.Tasks; using Elsa.ActivityResults; @@ -9,51 +11,41 @@ using Elsa.Expressions; using Elsa.Messaging.Domain; using Elsa.Models; using Elsa.Services.Models; +using Elsa.Triggers; using MediatR; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; +using Open.Linq.AsyncExtensions; namespace Elsa.Services { public class WorkflowHost : IWorkflowHost { - private delegate ValueTask ActivityOperation( - ActivityExecutionContext activityExecutionContext, - IActivity activity, - CancellationToken cancellationToken); + private delegate ValueTask ActivityOperation(ActivityExecutionContext activityExecutionContext, IActivity activity, CancellationToken cancellationToken); - private static readonly ActivityOperation Execute = (context, activity, cancellationToken) => - activity.ExecuteAsync(context, cancellationToken); + private static readonly ActivityOperation Execute = (context, activity, cancellationToken) => activity.ExecuteAsync(context, cancellationToken); + private static readonly ActivityOperation Resume = (context, activity, cancellationToken) => activity.ResumeAsync(context, cancellationToken); - private static readonly ActivityOperation Resume = (context, activity, cancellationToken) => - activity.ResumeAsync(context, cancellationToken); - - private readonly IWorkflowDefinitionManager _workflowDefinitionManager; - private readonly IWorkflowInstanceManager _workflowInstanceManager; private readonly IWorkflowRegistry _workflowRegistry; private readonly IWorkflowFactory _workflowFactory; - private readonly IExpressionEvaluator _expressionEvaluator; private readonly IMediator _mediator; private readonly IServiceProvider _serviceProvider; + private readonly IEnumerable _triggerProviders; private readonly ILogger _logger; public WorkflowHost( - IWorkflowDefinitionManager workflowDefinitionManager, - IWorkflowInstanceManager workflowInstanceManager, IWorkflowRegistry workflowRegistry, IWorkflowFactory workflowFactory, - IExpressionEvaluator expressionEvaluator, IMediator mediator, IServiceProvider serviceProvider, + IEnumerable triggerProviders, ILogger logger) { - _workflowDefinitionManager = workflowDefinitionManager; - _workflowInstanceManager = workflowInstanceManager; _workflowRegistry = workflowRegistry; _workflowFactory = workflowFactory; - _expressionEvaluator = expressionEvaluator; _mediator = mediator; _serviceProvider = serviceProvider; + _triggerProviders = triggerProviders; _logger = logger; } @@ -72,7 +64,6 @@ namespace Elsa.Services return await RunWorkflowAsync(workflowBlueprint, workflowInstance, activityId, input, cancellationToken); } - public async ValueTask RunWorkflowAsync( WorkflowInstance workflowInstance, string? activityId = default, @@ -236,12 +227,11 @@ namespace Elsa.Services workflowExecutionContext.Complete(); } - private WorkflowExecutionContext CreateWorkflowExecutionContext( + private static WorkflowExecutionContext CreateWorkflowExecutionContext( IWorkflowBlueprint workflowBlueprint, WorkflowInstance workflowInstance, IServiceScope serviceScope) => new WorkflowExecutionContext( - _expressionEvaluator, serviceScope.ServiceProvider, workflowBlueprint, workflowInstance); diff --git a/src/core/Elsa.Core/StartupTasks/StartServiceBusTask.cs b/src/core/Elsa.Core/StartupTasks/StartServiceBusTask.cs index e9b4b4ef5..989cde899 100644 --- a/src/core/Elsa.Core/StartupTasks/StartServiceBusTask.cs +++ b/src/core/Elsa.Core/StartupTasks/StartServiceBusTask.cs @@ -3,6 +3,7 @@ using System.Threading.Tasks; using Elsa.Messages; using Elsa.Messaging.Distributed; using Elsa.Runtime; +using Elsa.Services; using Rebus.Bus; namespace Elsa.StartupTasks diff --git a/src/core/Elsa.Core/Triggers/WorkflowSelector.cs b/src/core/Elsa.Core/Triggers/WorkflowSelector.cs new file mode 100644 index 000000000..5d68fd56e --- /dev/null +++ b/src/core/Elsa.Core/Triggers/WorkflowSelector.cs @@ -0,0 +1,43 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Runtime.CompilerServices; +using System.Threading; +using System.Threading.Tasks; +using Elsa.Services; +using Elsa.Services.Models; +using Open.Linq.AsyncExtensions; + +namespace Elsa.Triggers +{ + public class WorkflowSelector : IWorkflowSelector + { + private readonly IWorkflowRegistry _workflowRegistry; + private readonly IEnumerable _triggerProviders; + + public WorkflowSelector(IWorkflowRegistry workflowRegistry, IEnumerable triggerProviders) + { + _workflowRegistry = workflowRegistry; + _triggerProviders = triggerProviders; + } + + public async IAsyncEnumerable SelectWorkflowsAsync( + Type triggerType, + Func evaluate, + [EnumeratorCancellation] CancellationToken cancellationToken = default) + { + var providers = _triggerProviders.Where(x => x.ForType() == triggerType).ToList(); + var workflowBlueprints = await _workflowRegistry.GetWorkflowsAsync(cancellationToken).ToList(); + + foreach (var provider in providers) + foreach (var workflowBlueprint in workflowBlueprints) + { + var triggers = provider.GetTriggersAsync(workflowBlueprint, cancellationToken); + + await foreach (var trigger in triggers.WithCancellation(cancellationToken)) + if (evaluate(trigger)) + yield return new WorkflowSelectorResult(workflowBlueprint, trigger.ActivityId); + } + } + } +} \ No newline at end of file diff --git a/src/providers/Elsa.DistributedLocking.AzureBlob/Elsa.DistributedLocking.AzureBlob.csproj b/src/providers/Elsa.DistributedLocking.AzureBlob/Elsa.DistributedLocking.AzureBlob.csproj index ad4a38723..2f68cc4ff 100644 --- a/src/providers/Elsa.DistributedLocking.AzureBlob/Elsa.DistributedLocking.AzureBlob.csproj +++ b/src/providers/Elsa.DistributedLocking.AzureBlob/Elsa.DistributedLocking.AzureBlob.csproj @@ -5,7 +5,7 @@ - + diff --git a/src/samples/Elsa.Samples.DistributedLock/Elsa.Samples.DistributedLock.csproj b/src/samples/Elsa.Samples.DistributedLock/Elsa.Samples.DistributedLock.csproj index f5f36fc44..dcee68e21 100644 --- a/src/samples/Elsa.Samples.DistributedLock/Elsa.Samples.DistributedLock.csproj +++ b/src/samples/Elsa.Samples.DistributedLock/Elsa.Samples.DistributedLock.csproj @@ -6,7 +6,7 @@ - + diff --git a/src/samples/Elsa.Samples.HelloWorldConsole/Elsa.Samples.HelloWorldConsole.csproj b/src/samples/Elsa.Samples.HelloWorldConsole/Elsa.Samples.HelloWorldConsole.csproj index bea637f74..7fbd68615 100644 --- a/src/samples/Elsa.Samples.HelloWorldConsole/Elsa.Samples.HelloWorldConsole.csproj +++ b/src/samples/Elsa.Samples.HelloWorldConsole/Elsa.Samples.HelloWorldConsole.csproj @@ -6,7 +6,7 @@ - + diff --git a/src/samples/Elsa.Samples.Serialization/Elsa.Samples.Serialization.csproj b/src/samples/Elsa.Samples.Serialization/Elsa.Samples.Serialization.csproj index 5dd880a86..e69d3bb4f 100644 --- a/src/samples/Elsa.Samples.Serialization/Elsa.Samples.Serialization.csproj +++ b/src/samples/Elsa.Samples.Serialization/Elsa.Samples.Serialization.csproj @@ -6,7 +6,7 @@ - + diff --git a/src/server/Elsa.Server.GraphQL/Elsa.Server.GraphQL.csproj b/src/server/Elsa.Server.GraphQL/Elsa.Server.GraphQL.csproj index 794af7533..461975920 100644 --- a/src/server/Elsa.Server.GraphQL/Elsa.Server.GraphQL.csproj +++ b/src/server/Elsa.Server.GraphQL/Elsa.Server.GraphQL.csproj @@ -23,9 +23,9 @@ - - - + + + diff --git a/src/server/Elsa.Server.Host/Elsa.Server.Host.csproj b/src/server/Elsa.Server.Host/Elsa.Server.Host.csproj index e87297638..d117e9cd8 100644 --- a/src/server/Elsa.Server.Host/Elsa.Server.Host.csproj +++ b/src/server/Elsa.Server.Host/Elsa.Server.Host.csproj @@ -10,9 +10,9 @@ - - - + + + diff --git a/test/shared/Elsa.Testing.Shared/Elsa.Testing.Shared.csproj b/test/shared/Elsa.Testing.Shared/Elsa.Testing.Shared.csproj index 49f87c3d8..d3466b73f 100644 --- a/test/shared/Elsa.Testing.Shared/Elsa.Testing.Shared.csproj +++ b/test/shared/Elsa.Testing.Shared/Elsa.Testing.Shared.csproj @@ -18,8 +18,8 @@ - - + + diff --git a/test/unit/Elsa.Core.UnitTests/Elsa.Core.UnitTests.csproj b/test/unit/Elsa.Core.UnitTests/Elsa.Core.UnitTests.csproj index 65c9b986c..01220f479 100644 --- a/test/unit/Elsa.Core.UnitTests/Elsa.Core.UnitTests.csproj +++ b/test/unit/Elsa.Core.UnitTests/Elsa.Core.UnitTests.csproj @@ -9,10 +9,10 @@ - + - - + + all