diff --git a/Elsa.sln b/Elsa.sln index 56353aed3..45b5d697b 100644 --- a/Elsa.sln +++ b/Elsa.sln @@ -245,6 +245,8 @@ Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Elsa.Activities.Telnyx", "s EndProject Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Elsa.Activities.Webhooks", "src\activities\Elsa.Activities.Webhooks\Elsa.Activities.Webhooks.csproj", "{2B67E954-3B04-402D-A9A7-AAAB1D6C3215}" EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Server.Orleans", "src\server\Elsa.Server.Orleans\Elsa.Server.Orleans.csproj", "{76BD888E-F0FD-40DD-B025-2ED773C1C50B}" +EndProject Global GlobalSection(SolutionConfigurationPlatforms) = preSolution Debug|Any CPU = Debug|Any CPU @@ -591,6 +593,10 @@ Global {2B67E954-3B04-402D-A9A7-AAAB1D6C3215}.Debug|Any CPU.Build.0 = Debug|Any CPU {2B67E954-3B04-402D-A9A7-AAAB1D6C3215}.Release|Any CPU.ActiveCfg = Release|Any CPU {2B67E954-3B04-402D-A9A7-AAAB1D6C3215}.Release|Any CPU.Build.0 = Release|Any CPU + {76BD888E-F0FD-40DD-B025-2ED773C1C50B}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {76BD888E-F0FD-40DD-B025-2ED773C1C50B}.Debug|Any CPU.Build.0 = Debug|Any CPU + {76BD888E-F0FD-40DD-B025-2ED773C1C50B}.Release|Any CPU.ActiveCfg = Release|Any CPU + {76BD888E-F0FD-40DD-B025-2ED773C1C50B}.Release|Any CPU.Build.0 = Release|Any CPU EndGlobalSection GlobalSection(SolutionProperties) = preSolution HideSolutionNode = FALSE @@ -708,6 +714,7 @@ Global {E3EA6449-28EC-48E4-91C4-E82DE08A7DD9} = {C865B0FD-E505-48F0-BFAF-0D4D7C1B5CA1} {93878BAD-855D-48D1-97ED-C776EC54E19E} = {B43B546E-23F3-46E8-ACB7-D04F05CDA180} {2B67E954-3B04-402D-A9A7-AAAB1D6C3215} = {B43B546E-23F3-46E8-ACB7-D04F05CDA180} + {76BD888E-F0FD-40DD-B025-2ED773C1C50B} = {468E498C-59DB-4541-8D8A-0D89DE9083AB} EndGlobalSection GlobalSection(ExtensibilityGlobals) = postSolution SolutionGuid = {8B0975FD-7050-48B0-88C5-48C33378E158} diff --git a/src/core/Elsa.Abstractions/Dispatch/ICorrelatingWorkflowDispatcher.cs b/src/core/Elsa.Abstractions/Dispatch/ICorrelatingWorkflowDispatcher.cs index 252f4ae1c..6609a4fe1 100644 --- a/src/core/Elsa.Abstractions/Dispatch/ICorrelatingWorkflowDispatcher.cs +++ b/src/core/Elsa.Abstractions/Dispatch/ICorrelatingWorkflowDispatcher.cs @@ -9,6 +9,6 @@ namespace Elsa.Dispatch /// public interface ICorrelatingWorkflowDispatcher { - Task DispatchAsync(ExecuteCorrelatedWorkflowRequest request, CancellationToken? cancellationToken = default); + Task DispatchAsync(ExecuteCorrelatedWorkflowRequest request, CancellationToken cancellationToken = default); } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Dispatch/IWorkflowDefinitionDispatcher.cs b/src/core/Elsa.Abstractions/Dispatch/IWorkflowDefinitionDispatcher.cs index 4a95cac58..2f25e3fb1 100644 --- a/src/core/Elsa.Abstractions/Dispatch/IWorkflowDefinitionDispatcher.cs +++ b/src/core/Elsa.Abstractions/Dispatch/IWorkflowDefinitionDispatcher.cs @@ -8,6 +8,6 @@ namespace Elsa.Dispatch /// public interface IWorkflowDefinitionDispatcher { - Task DispatchAsync(ExecuteWorkflowDefinitionRequest request, CancellationToken? cancellationToken = default); + Task DispatchAsync(ExecuteWorkflowDefinitionRequest request, CancellationToken cancellationToken = default); } } \ No newline at end of file diff --git a/src/core/Elsa.Abstractions/Dispatch/IWorkflowInstanceDispatcher.cs b/src/core/Elsa.Abstractions/Dispatch/IWorkflowInstanceDispatcher.cs index 806b5aee1..1158aacd4 100644 --- a/src/core/Elsa.Abstractions/Dispatch/IWorkflowInstanceDispatcher.cs +++ b/src/core/Elsa.Abstractions/Dispatch/IWorkflowInstanceDispatcher.cs @@ -8,6 +8,6 @@ namespace Elsa.Dispatch /// public interface IWorkflowInstanceDispatcher { - Task DispatchAsync(ExecuteWorkflowInstanceRequest request, CancellationToken? cancellationToken = default); + Task DispatchAsync(ExecuteWorkflowInstanceRequest request, CancellationToken cancellationToken = default); } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Dispatch/Consumers/ExecuteCorrelatedWorkflowRequestConsumer.cs b/src/core/Elsa.Core/Dispatch/Consumers/ExecuteCorrelatedWorkflowRequestConsumer.cs index 9a5d71182..3f8ca590d 100644 --- a/src/core/Elsa.Core/Dispatch/Consumers/ExecuteCorrelatedWorkflowRequestConsumer.cs +++ b/src/core/Elsa.Core/Dispatch/Consumers/ExecuteCorrelatedWorkflowRequestConsumer.cs @@ -67,7 +67,6 @@ namespace Elsa.Dispatch.Consumers } else { - // Trigger new workflow. _logger.LogDebug("No existing workflows found with correlation ID '{CorrelationId}'. Starting new workflow", correlationId); await TriggerNewWorkflowAsync(message); } diff --git a/src/core/Elsa.Core/Dispatch/QueuingWorkflowDispatcher.cs b/src/core/Elsa.Core/Dispatch/QueuingWorkflowDispatcher.cs index 4ac0d4fc3..98e27fcfe 100644 --- a/src/core/Elsa.Core/Dispatch/QueuingWorkflowDispatcher.cs +++ b/src/core/Elsa.Core/Dispatch/QueuingWorkflowDispatcher.cs @@ -11,8 +11,8 @@ namespace Elsa.Dispatch { private readonly ICommandSender _commandSender; public QueuingWorkflowDispatcher(ICommandSender commandSender) => _commandSender = commandSender; - public async Task DispatchAsync(ExecuteWorkflowInstanceRequest request, CancellationToken? cancellationToken = default) => await _commandSender.SendAsync(request); - public async Task DispatchAsync(ExecuteCorrelatedWorkflowRequest request, CancellationToken? cancellationToken = default) => await _commandSender.SendAsync(request); - public async Task DispatchAsync(ExecuteWorkflowDefinitionRequest request, CancellationToken? cancellationToken = default) => await _commandSender.SendAsync(request); + public async Task DispatchAsync(ExecuteWorkflowInstanceRequest request, CancellationToken cancellationToken = default) => await _commandSender.SendAsync(request); + public async Task DispatchAsync(ExecuteCorrelatedWorkflowRequest request, CancellationToken cancellationToken = default) => await _commandSender.SendAsync(request); + public async Task DispatchAsync(ExecuteWorkflowDefinitionRequest request, CancellationToken cancellationToken = default) => await _commandSender.SendAsync(request); } } \ No newline at end of file diff --git a/src/core/Elsa.Core/Extensions/WorkflowRegistryExtensions.cs b/src/core/Elsa.Core/Extensions/WorkflowRegistryExtensions.cs index 18773b8f4..956a8278b 100644 --- a/src/core/Elsa.Core/Extensions/WorkflowRegistryExtensions.cs +++ b/src/core/Elsa.Core/Extensions/WorkflowRegistryExtensions.cs @@ -24,33 +24,14 @@ namespace Elsa string id, VersionOptions versionOptions, CancellationToken cancellationToken = default) => - workflowRegistry.GetAsync(id, default, versionOptions, cancellationToken); - - // public static async Task> - // GetWorkflowsByStartActivityAsync( - // this IWorkflowRegistry workflowRegistry, - // CancellationToken cancellationToken = default) - // where T : IActivity - // { - // var results = await workflowRegistry.GetWorkflowsByStartActivityAsync(typeof(T).Name, cancellationToken); - // return results.Select(x => (x.Workflow, x.Activity)); - // } - - // public static async Task> GetWorkflowsByStartActivityAsync( - // this IWorkflowRegistry workflowRegistry, - // string activityType, - // CancellationToken cancellationToken = default) - // { - // var workflows = await workflowRegistry.GetWorkflowsAsync(cancellationToken).ToListAsync(cancellationToken); - // - // var query = - // from workflow in workflows - // where workflow.IsPublished - // from activity in workflow.GetStartActivities() - // where activity.Type == activityType - // select (workflow, activity); - // - // return query.Distinct(); - // } + workflowRegistry.GetWorkflowAsync(id, default, versionOptions, cancellationToken); + + public static Task GetWorkflowAsync( + this IWorkflowRegistry workflowRegistry, + string id, + string? tenantId, + VersionOptions versionOptions, + CancellationToken cancellationToken = default) => + workflowRegistry.GetAsync(id, tenantId, versionOptions, cancellationToken); } } \ No newline at end of file diff --git a/src/samples/server/Elsa.Samples.Server.Host/Elsa.Samples.Server.Host.csproj b/src/samples/server/Elsa.Samples.Server.Host/Elsa.Samples.Server.Host.csproj index bd911137c..26508551e 100644 --- a/src/samples/server/Elsa.Samples.Server.Host/Elsa.Samples.Server.Host.csproj +++ b/src/samples/server/Elsa.Samples.Server.Host/Elsa.Samples.Server.Host.csproj @@ -1,6 +1,6 @@ - - + + net5.0 latest @@ -9,17 +9,14 @@ - + + + + + + - - - - - - - - .\Workflows\HelloWorld.cs diff --git a/src/samples/server/Elsa.Samples.Server.Host/Program.cs b/src/samples/server/Elsa.Samples.Server.Host/Program.cs index 8c6c313e7..e83fd74aa 100644 --- a/src/samples/server/Elsa.Samples.Server.Host/Program.cs +++ b/src/samples/server/Elsa.Samples.Server.Host/Program.cs @@ -1,5 +1,10 @@ +using System.Net; +using Elsa.Server.Orleans.Grains.Contracts; using Microsoft.AspNetCore.Hosting; using Microsoft.Extensions.Hosting; +using Orleans; +using Orleans.Configuration; +using Orleans.Hosting; namespace Elsa.Samples.Server.Host { @@ -12,6 +17,15 @@ namespace Elsa.Samples.Server.Host public static IHostBuilder CreateHostBuilder(string[] args) => Microsoft.Extensions.Hosting.Host.CreateDefaultBuilder(args) - .ConfigureWebHostDefaults(webBuilder => { webBuilder.UseStartup(); }); + .ConfigureWebHostDefaults(webBuilder => webBuilder.UseStartup()) + .UseOrleans(siloBuilder => siloBuilder + .UseLocalhostClustering() + .Configure(options => + { + options.ClusterId = "localhost"; + options.ServiceId = "elsa-workflows"; + }) + .ConfigureApplicationParts(parts => parts.AddApplicationPart(typeof(IWorkflowDefinitionGrain).Assembly).WithReferences()) + .Configure(options => options.AdvertisedIPAddress = IPAddress.Loopback)); } } \ No newline at end of file diff --git a/src/samples/server/Elsa.Samples.Server.Host/Startup.cs b/src/samples/server/Elsa.Samples.Server.Host/Startup.cs index a5499de4f..dc9a19962 100644 --- a/src/samples/server/Elsa.Samples.Server.Host/Startup.cs +++ b/src/samples/server/Elsa.Samples.Server.Host/Startup.cs @@ -1,6 +1,7 @@ using Elsa.Activities.Telnyx.Extensions; using Elsa.Persistence.EntityFramework.Core.Extensions; using Elsa.Persistence.EntityFramework.Sqlite; +using Elsa.Server.Orleans.Extensions; using Microsoft.AspNetCore.Builder; using Microsoft.AspNetCore.Hosting; using Microsoft.Extensions.Configuration; @@ -27,6 +28,7 @@ namespace Elsa.Samples.Server.Host services .AddElsa(elsa => elsa .UseEntityFrameworkPersistence(ef => ef.UseSqlite()) + .UseOrleansDispatchers() .AddConsoleActivities() .AddHttpActivities(elsaSection.GetSection("Http").Bind) .AddEmailActivities(elsaSection.GetSection("Smtp").Bind) diff --git a/src/samples/server/Elsa.Samples.Server.Host/Workflows/HeartbeatWorkflow.cs b/src/samples/server/Elsa.Samples.Server.Host/Workflows/HeartbeatWorkflow.cs new file mode 100644 index 000000000..5a343e4c0 --- /dev/null +++ b/src/samples/server/Elsa.Samples.Server.Host/Workflows/HeartbeatWorkflow.cs @@ -0,0 +1,21 @@ +using Elsa.Activities.Console; +using Elsa.Activities.Temporal; +using Elsa.Builders; +using NodaTime; + +namespace Elsa.Samples.Server.Host.Workflows +{ + public class HeartbeatWorkflow : IWorkflow + { + private readonly IClock _clock; + public HeartbeatWorkflow(IClock clock) => _clock = clock; + + public void Build(IWorkflowBuilder builder) + { + builder + .WithDisplayName("Timer") + .Timer(Duration.FromSeconds(10)) + .WriteLine(() => $"Heartbeat at {_clock.GetCurrentInstant()}"); + } + } +} \ No newline at end of file diff --git a/src/server/Elsa.Server.Orleans/Dispatch/GrainWorkflowDispatcher.cs b/src/server/Elsa.Server.Orleans/Dispatch/GrainWorkflowDispatcher.cs new file mode 100644 index 000000000..e25fe2356 --- /dev/null +++ b/src/server/Elsa.Server.Orleans/Dispatch/GrainWorkflowDispatcher.cs @@ -0,0 +1,37 @@ +using System.Threading; +using System.Threading.Tasks; +using Elsa.Server.Orleans.Grains; +using Elsa.Dispatch; +using Elsa.Server.Orleans.Grains.Contracts; +using Orleans; + +namespace Elsa.Server.Orleans.Dispatch +{ + public class GrainWorkflowDispatcher : IWorkflowDefinitionDispatcher, IWorkflowInstanceDispatcher, ICorrelatingWorkflowDispatcher + { + private readonly IClusterClient _clusterClient; + + public GrainWorkflowDispatcher(IClusterClient clusterClient) + { + _clusterClient = clusterClient; + } + + public async Task DispatchAsync(ExecuteWorkflowDefinitionRequest request, CancellationToken cancellationToken = default) + { + var grain = _clusterClient.GetGrain(request.WorkflowDefinitionId); + await grain.ExecuteWorkflowAsync(request, cancellationToken); + } + + public async Task DispatchAsync(ExecuteWorkflowInstanceRequest request, CancellationToken cancellationToken = default) + { + var grain = _clusterClient.GetGrain(request.WorkflowInstanceId); + await grain.ExecuteWorkflowAsync(request, cancellationToken); + } + + public async Task DispatchAsync(ExecuteCorrelatedWorkflowRequest request, CancellationToken cancellationToken = default) + { + var grain = _clusterClient.GetGrain(request.CorrelationId); + await grain.ExecutedCorrelatedWorkflowAsync(request, cancellationToken); + } + } +} \ No newline at end of file diff --git a/src/server/Elsa.Server.Orleans/Elsa.Server.Orleans.csproj b/src/server/Elsa.Server.Orleans/Elsa.Server.Orleans.csproj new file mode 100644 index 000000000..61dd9f3e5 --- /dev/null +++ b/src/server/Elsa.Server.Orleans/Elsa.Server.Orleans.csproj @@ -0,0 +1,23 @@ + + + + + + + net5.0 + + Elsa is a set of workflow libraries and tools that enable lean and mean workflowing capabilities in any .NET Core application. + This package provides Orleans grains that handle workflow dispatch. + + elsa, workflows, orleans, actor model + + + + + + + + + + + diff --git a/src/server/Elsa.Server.Orleans/Extensions/ServiceCollectionExtensions.cs b/src/server/Elsa.Server.Orleans/Extensions/ServiceCollectionExtensions.cs new file mode 100644 index 000000000..67da13e93 --- /dev/null +++ b/src/server/Elsa.Server.Orleans/Extensions/ServiceCollectionExtensions.cs @@ -0,0 +1,23 @@ +using Elsa.Dispatch; +using Elsa.Server.Orleans.Dispatch; +using Microsoft.Extensions.DependencyInjection; + +namespace Elsa.Server.Orleans.Extensions +{ + public static class ServiceCollectionExtensions + { + public static ElsaOptions UseOrleansDispatchers(this ElsaOptions elsaOptions) + { + var services = elsaOptions.Services; + + services.AddSingleton(); + + elsaOptions + .UseCorrelatingWorkflowDispatcher(sp => sp.GetRequiredService()) + .UseWorkflowDefinitionDispatcher(sp => sp.GetRequiredService()) + .UseWorkflowInstanceDispatcher(sp => sp.GetRequiredService()); + + return elsaOptions; + } + } +} \ No newline at end of file diff --git a/src/server/Elsa.Server.Orleans/FodyWeavers.xml b/src/server/Elsa.Server.Orleans/FodyWeavers.xml new file mode 100644 index 000000000..00e1d9a1c --- /dev/null +++ b/src/server/Elsa.Server.Orleans/FodyWeavers.xml @@ -0,0 +1,3 @@ + + + \ No newline at end of file diff --git a/src/server/Elsa.Server.Orleans/Grains/Contracts/ICorrelatedWorkflowGrain.cs b/src/server/Elsa.Server.Orleans/Grains/Contracts/ICorrelatedWorkflowGrain.cs new file mode 100644 index 000000000..928808b22 --- /dev/null +++ b/src/server/Elsa.Server.Orleans/Grains/Contracts/ICorrelatedWorkflowGrain.cs @@ -0,0 +1,12 @@ +using System.Threading; +using System.Threading.Tasks; +using Elsa.Dispatch; +using Orleans; + +namespace Elsa.Server.Orleans.Grains.Contracts +{ + public interface ICorrelatedWorkflowGrain : IGrainWithStringKey + { + Task ExecutedCorrelatedWorkflowAsync(ExecuteCorrelatedWorkflowRequest request, CancellationToken cancellationToken = default); + } +} \ No newline at end of file diff --git a/src/server/Elsa.Server.Orleans/Grains/Contracts/IWorkflowDefinitionGrain.cs b/src/server/Elsa.Server.Orleans/Grains/Contracts/IWorkflowDefinitionGrain.cs new file mode 100644 index 000000000..cb670e322 --- /dev/null +++ b/src/server/Elsa.Server.Orleans/Grains/Contracts/IWorkflowDefinitionGrain.cs @@ -0,0 +1,12 @@ +using System.Threading; +using System.Threading.Tasks; +using Elsa.Dispatch; +using Orleans; + +namespace Elsa.Server.Orleans.Grains.Contracts +{ + public interface IWorkflowDefinitionGrain : IGrainWithStringKey + { + Task ExecuteWorkflowAsync(ExecuteWorkflowDefinitionRequest request, CancellationToken cancellationToken = default); + } +} \ No newline at end of file diff --git a/src/server/Elsa.Server.Orleans/Grains/Contracts/IWorkflowInstanceGrain.cs b/src/server/Elsa.Server.Orleans/Grains/Contracts/IWorkflowInstanceGrain.cs new file mode 100644 index 000000000..9d10a9e17 --- /dev/null +++ b/src/server/Elsa.Server.Orleans/Grains/Contracts/IWorkflowInstanceGrain.cs @@ -0,0 +1,12 @@ +using System.Threading; +using System.Threading.Tasks; +using Elsa.Dispatch; +using Orleans; + +namespace Elsa.Server.Orleans.Grains.Contracts +{ + public interface IWorkflowInstanceGrain : IGrainWithStringKey + { + Task ExecuteWorkflowAsync(ExecuteWorkflowInstanceRequest request, CancellationToken cancellationToken = default); + } +} \ No newline at end of file diff --git a/src/server/Elsa.Server.Orleans/Grains/CorrelatedWorkflowDefinitionGrain.cs b/src/server/Elsa.Server.Orleans/Grains/CorrelatedWorkflowDefinitionGrain.cs new file mode 100644 index 000000000..3d2d99c20 --- /dev/null +++ b/src/server/Elsa.Server.Orleans/Grains/CorrelatedWorkflowDefinitionGrain.cs @@ -0,0 +1,78 @@ +using System.Collections.Generic; +using System.Threading; +using System.Threading.Tasks; +using Elsa.Bookmarks; +using Elsa.Dispatch; +using Elsa.Models; +using Elsa.Persistence; +using Elsa.Persistence.Specifications; +using Elsa.Server.Orleans.Grains.Contracts; +using Elsa.Triggers; +using Microsoft.Extensions.Logging; +using Open.Linq.AsyncExtensions; +using Orleans; + +namespace Elsa.Server.Orleans.Grains +{ + public class CorrelatedWorkflowDefinitionGrain : Grain, ICorrelatedWorkflowGrain + { + private readonly IWorkflowInstanceStore _workflowInstanceStore; + private readonly IBookmarkFinder _bookmarkFinder; + private readonly ITriggerFinder _triggerFinder; + private readonly IWorkflowDefinitionDispatcher _workflowDefinitionDispatcher; + private readonly IWorkflowInstanceDispatcher _workflowInstanceDispatcher; + private readonly ILogger _logger; + + public CorrelatedWorkflowDefinitionGrain( + IWorkflowInstanceStore workflowInstanceStore, + IBookmarkFinder bookmarkFinder, + ITriggerFinder triggerFinder, + IWorkflowDefinitionDispatcher workflowDefinitionDispatcher, + IWorkflowInstanceDispatcher workflowInstanceDispatcher, + ILogger logger) + { + _workflowInstanceStore = workflowInstanceStore; + _bookmarkFinder = bookmarkFinder; + _triggerFinder = triggerFinder; + _workflowDefinitionDispatcher = workflowDefinitionDispatcher; + _workflowInstanceDispatcher = workflowInstanceDispatcher; + _logger = logger; + } + + public async Task ExecutedCorrelatedWorkflowAsync(ExecuteCorrelatedWorkflowRequest request, CancellationToken cancellationToken = default) + { + var correlationId = request.CorrelationId; + var correlatedWorkflowInstanceCount = await _workflowInstanceStore.CountAsync(new CorrelationIdSpecification(correlationId), cancellationToken); + + if (correlatedWorkflowInstanceCount > 0) + { + _logger.LogDebug("{WorkflowInstanceCount} existing workflows found with correlation ID '{CorrelationId}' will be queued for execution", correlatedWorkflowInstanceCount, correlationId); + var existingWorkflows = await _bookmarkFinder.FindBookmarksAsync(request.ActivityType, request.Bookmark, request.TenantId, cancellationToken).ToList(); + await ResumeWorkflowsAsync(existingWorkflows, request.Input, cancellationToken); + } + else + { + _logger.LogDebug("No existing workflows found with correlation ID '{CorrelationId}'. Starting new workflow", correlationId); + await StartWorkflowsAsync(request, cancellationToken); + } + } + + private async Task StartWorkflowsAsync(ExecuteCorrelatedWorkflowRequest request, CancellationToken cancellationToken) + { + var filter = request.Trigger; + var triggers = await _triggerFinder.FindTriggersAsync(request.ActivityType, filter, request.TenantId, cancellationToken); + + foreach (var trigger in triggers) + { + var workflowBlueprint = trigger.WorkflowBlueprint; + await _workflowDefinitionDispatcher.DispatchAsync(new ExecuteWorkflowDefinitionRequest(workflowBlueprint.Id, trigger.ActivityId, request.Input, request.CorrelationId, request.ContextId, workflowBlueprint.TenantId), cancellationToken); + } + } + + private async Task ResumeWorkflowsAsync(IEnumerable results, object? input, CancellationToken cancellationToken) + { + foreach (var result in results) + await _workflowInstanceDispatcher.DispatchAsync(new ExecuteWorkflowInstanceRequest(result.WorkflowInstanceId, result.ActivityId, input), cancellationToken); + } + } +} \ No newline at end of file diff --git a/src/server/Elsa.Server.Orleans/Grains/WorkflowDefinitionGrain.cs b/src/server/Elsa.Server.Orleans/Grains/WorkflowDefinitionGrain.cs new file mode 100644 index 000000000..d6cb2a396 --- /dev/null +++ b/src/server/Elsa.Server.Orleans/Grains/WorkflowDefinitionGrain.cs @@ -0,0 +1,38 @@ +using System.Threading; +using System.Threading.Tasks; +using Elsa.Dispatch; +using Elsa.Models; +using Elsa.Server.Orleans.Grains.Contracts; +using Elsa.Services; +using Microsoft.Extensions.Logging; +using Orleans; + +namespace Elsa.Server.Orleans.Grains +{ + public class WorkflowDefinitionGrain : Grain, IWorkflowDefinitionGrain + { + private readonly IWorkflowRegistry _workflowRegistry; + private readonly IWorkflowRunner _workflowRunner; + private readonly ILogger _logger; + + public WorkflowDefinitionGrain(IWorkflowRegistry workflowRegistry, IWorkflowRunner workflowRunner, ILogger logger) + { + _workflowRegistry = workflowRegistry; + _workflowRunner = workflowRunner; + _logger = logger; + } + + public async Task ExecuteWorkflowAsync(ExecuteWorkflowDefinitionRequest request, CancellationToken cancellationToken = default) + { + var workflowBlueprint = await _workflowRegistry.GetWorkflowAsync(request.WorkflowDefinitionId, request.TenantId, VersionOptions.Published, cancellationToken); + + if (workflowBlueprint == null) + { + _logger.LogWarning("No published workflow definition {WorkflowDefinitionId} found", request.WorkflowDefinitionId); + return; + } + + await _workflowRunner.RunWorkflowAsync(workflowBlueprint, request.ActivityId, request.Input, request.CorrelationId, request.ContextId, cancellationToken); + } + } +} \ No newline at end of file diff --git a/src/server/Elsa.Server.Orleans/Grains/WorkflowInstanceGrain.cs b/src/server/Elsa.Server.Orleans/Grains/WorkflowInstanceGrain.cs new file mode 100644 index 000000000..5ec7f3899 --- /dev/null +++ b/src/server/Elsa.Server.Orleans/Grains/WorkflowInstanceGrain.cs @@ -0,0 +1,38 @@ +using System.Threading; +using System.Threading.Tasks; +using Elsa.Dispatch; +using Elsa.Persistence; +using Elsa.Server.Orleans.Grains.Contracts; +using Elsa.Services; +using Microsoft.Extensions.Logging; +using Orleans; + +namespace Elsa.Server.Orleans.Grains +{ + public class WorkflowInstanceGrain : Grain, IWorkflowInstanceGrain + { + private readonly IWorkflowInstanceStore _store; + private readonly IWorkflowRunner _workflowRunner; + private readonly ILogger _logger; + + public WorkflowInstanceGrain(IWorkflowInstanceStore store, IWorkflowRunner workflowRunner, ILogger logger) + { + _store = store; + _workflowRunner = workflowRunner; + _logger = logger; + } + + public async Task ExecuteWorkflowAsync(ExecuteWorkflowInstanceRequest request, CancellationToken cancellationToken = default) + { + var workflowInstance = await _store.FindByIdAsync(request.WorkflowInstanceId, cancellationToken); + + if(workflowInstance == null) + { + _logger.LogWarning("Workflow instance {WorkflowInstanceId} not found", request.WorkflowInstanceId); + return; + } + + await _workflowRunner.RunWorkflowAsync(workflowInstance, request.ActivityId, request.Input, cancellationToken); + } + } +} \ No newline at end of file