From 59bfcf0ba51b8ab9268514f52bf08702c5ae8f9b Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Tue, 3 Jun 2025 09:58:20 +0200 Subject: [PATCH] Add distributed workflow runtime implementation. Introduced `DistributedWorkflowRuntime` to support distributed workflow execution with locking mechanisms. Added new module `Elsa.Workflows.Runtime.Distributed` with key services, features, and client implementations for handling distributed bookmarks and workflow clients. Updated integration and component tests to use the new distributed runtime where relevant. --- Elsa.sln | 7 ++ .../DispatchWorkflowExtensions.cs | 2 +- .../Elsa.Testing.Shared.Integration.csproj | 1 + .../TestApplicationBuilder.cs | 1 + .../Elsa.Workflows.Runtime.Distributed.csproj | 21 +++++ ...ows.Runtime.Distributed.csproj.DotSettings | 2 + .../Extensions/ModuleExtensions.cs | 13 +++ .../Features/DistributedRuntimeFeature.cs | 39 +++++++++ .../FodyWeavers.xml | 3 + .../DistributedBookmarkQueueWorker.cs | 25 ++++++ .../Services/DistributedWorkflowClient.cs | 86 +++++++++++++++++++ .../DistributedWorkflowRuntime.Obsolete.cs | 32 +++++++ .../Services/DistributedWorkflowRuntime.cs | 36 ++++++++ .../Services/WorkflowDefinitionRefresher.cs | 2 +- .../Elsa.Workflows.ComponentTests.csproj | 2 + .../Helpers/Fixtures/WorkflowServer.cs | 24 +++++- .../DynamicEndpointTests.cs | 8 +- .../Elsa.Workflows.IntegrationTests.csproj | 1 + .../RunAsynchronousActivityOutput/Tests.cs | 6 ++ .../ContentWriters/RawStringContentTests.cs | 25 ------ 20 files changed, 304 insertions(+), 32 deletions(-) create mode 100644 src/modules/Elsa.Workflows.Runtime.Distributed/Elsa.Workflows.Runtime.Distributed.csproj create mode 100644 src/modules/Elsa.Workflows.Runtime.Distributed/Elsa.Workflows.Runtime.Distributed.csproj.DotSettings create mode 100644 src/modules/Elsa.Workflows.Runtime.Distributed/Extensions/ModuleExtensions.cs create mode 100644 src/modules/Elsa.Workflows.Runtime.Distributed/Features/DistributedRuntimeFeature.cs create mode 100644 src/modules/Elsa.Workflows.Runtime.Distributed/FodyWeavers.xml create mode 100644 src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedBookmarkQueueWorker.cs create mode 100644 src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowClient.cs create mode 100644 src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowRuntime.Obsolete.cs create mode 100644 src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowRuntime.cs diff --git a/Elsa.sln b/Elsa.sln index d2149e5f8..82e4dc8cf 100644 --- a/Elsa.sln +++ b/Elsa.sln @@ -264,6 +264,8 @@ Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "scheduling", "scheduling", EndProject Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Scheduling", "src\modules\Elsa.Scheduling\Elsa.Scheduling.csproj", "{26849C37-2ACA-4DDE-83CD-B87939465791}" EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "Elsa.Workflows.Runtime.Distributed", "src\modules\Elsa.Workflows.Runtime.Distributed\Elsa.Workflows.Runtime.Distributed.csproj", "{42CE3E10-BC73-4D29-B099-88F0D5319E57}" +EndProject Global GlobalSection(SolutionConfigurationPlatforms) = preSolution Debug|Any CPU = Debug|Any CPU @@ -452,6 +454,10 @@ Global {26849C37-2ACA-4DDE-83CD-B87939465791}.Debug|Any CPU.Build.0 = Debug|Any CPU {26849C37-2ACA-4DDE-83CD-B87939465791}.Release|Any CPU.ActiveCfg = Release|Any CPU {26849C37-2ACA-4DDE-83CD-B87939465791}.Release|Any CPU.Build.0 = Release|Any CPU + {42CE3E10-BC73-4D29-B099-88F0D5319E57}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {42CE3E10-BC73-4D29-B099-88F0D5319E57}.Debug|Any CPU.Build.0 = Debug|Any CPU + {42CE3E10-BC73-4D29-B099-88F0D5319E57}.Release|Any CPU.ActiveCfg = Release|Any CPU + {42CE3E10-BC73-4D29-B099-88F0D5319E57}.Release|Any CPU.Build.0 = Release|Any CPU EndGlobalSection GlobalSection(SolutionProperties) = preSolution HideSolutionNode = FALSE @@ -530,6 +536,7 @@ Global {0125BC8D-C837-414C-BB53-FCEA62E3060F} = {B08B4E00-C2AB-48F3-8389-449F42AEF179} {7C038270-8BF7-4911-AC24-54D7BBE9BBF2} = {5BA4A8FA-F7F4-45B3-AEC8-8886D35AAC79} {26849C37-2ACA-4DDE-83CD-B87939465791} = {7C038270-8BF7-4911-AC24-54D7BBE9BBF2} + {42CE3E10-BC73-4D29-B099-88F0D5319E57} = {B08B4E00-C2AB-48F3-8389-449F42AEF179} EndGlobalSection GlobalSection(ExtensibilityGlobals) = postSolution SolutionGuid = {D4B5CEAA-7D70-4FCB-A68E-B03FBE5E0E5E} diff --git a/src/common/Elsa.Testing.Shared.Integration/DispatchWorkflowExtensions.cs b/src/common/Elsa.Testing.Shared.Integration/DispatchWorkflowExtensions.cs index fcd9264a4..14ed38c3f 100644 --- a/src/common/Elsa.Testing.Shared.Integration/DispatchWorkflowExtensions.cs +++ b/src/common/Elsa.Testing.Shared.Integration/DispatchWorkflowExtensions.cs @@ -65,7 +65,7 @@ public static class DispatchWorkflowExtensions dispatchWorkflowResponse.ThrowIfFailed(); // Wait for the workflow to complete, and then return the WorkflowFinished notification. - var signaled = await semaphore.WaitAsync(timeout ?? TimeSpan.FromSeconds(5)); + var signaled = await semaphore.WaitAsync(timeout ?? TimeSpan.FromSeconds(50000)); return signaled ? workflowFinishedRecord : null; } finally diff --git a/src/common/Elsa.Testing.Shared.Integration/Elsa.Testing.Shared.Integration.csproj b/src/common/Elsa.Testing.Shared.Integration/Elsa.Testing.Shared.Integration.csproj index 75336573d..605e97744 100644 --- a/src/common/Elsa.Testing.Shared.Integration/Elsa.Testing.Shared.Integration.csproj +++ b/src/common/Elsa.Testing.Shared.Integration/Elsa.Testing.Shared.Integration.csproj @@ -17,6 +17,7 @@ + diff --git a/src/common/Elsa.Testing.Shared.Integration/TestApplicationBuilder.cs b/src/common/Elsa.Testing.Shared.Integration/TestApplicationBuilder.cs index 14154f166..ac382f618 100644 --- a/src/common/Elsa.Testing.Shared.Integration/TestApplicationBuilder.cs +++ b/src/common/Elsa.Testing.Shared.Integration/TestApplicationBuilder.cs @@ -35,6 +35,7 @@ public class TestApplicationBuilder _configureElsa += elsa => elsa .AddActivitiesFrom() + .UseScheduling() .UseCSharp() .UseJavaScript() .UseLiquid() diff --git a/src/modules/Elsa.Workflows.Runtime.Distributed/Elsa.Workflows.Runtime.Distributed.csproj b/src/modules/Elsa.Workflows.Runtime.Distributed/Elsa.Workflows.Runtime.Distributed.csproj new file mode 100644 index 000000000..0439f66b4 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime.Distributed/Elsa.Workflows.Runtime.Distributed.csproj @@ -0,0 +1,21 @@ + + + + + Provides distributed workflow runtime functionality. + + elsa extensions module workflows distributed runtime + + + + + + + + + + + + + + \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime.Distributed/Elsa.Workflows.Runtime.Distributed.csproj.DotSettings b/src/modules/Elsa.Workflows.Runtime.Distributed/Elsa.Workflows.Runtime.Distributed.csproj.DotSettings new file mode 100644 index 000000000..d3ee2fe40 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime.Distributed/Elsa.Workflows.Runtime.Distributed.csproj.DotSettings @@ -0,0 +1,2 @@ + + True \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime.Distributed/Extensions/ModuleExtensions.cs b/src/modules/Elsa.Workflows.Runtime.Distributed/Extensions/ModuleExtensions.cs new file mode 100644 index 000000000..9372b5e4a --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime.Distributed/Extensions/ModuleExtensions.cs @@ -0,0 +1,13 @@ +using Elsa.Workflows.Runtime.Distributed.Features; +using Elsa.Workflows.Runtime.Features; + +namespace Elsa.Workflows.Runtime.Distributed.Extensions; + +public static class ModuleExtensions +{ + public static WorkflowRuntimeFeature UseDistributedRuntime(this WorkflowRuntimeFeature feature, Action? configure = null) + { + feature.Module.Configure(configure); + return feature; + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime.Distributed/Features/DistributedRuntimeFeature.cs b/src/modules/Elsa.Workflows.Runtime.Distributed/Features/DistributedRuntimeFeature.cs new file mode 100644 index 000000000..0e05eca4a --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime.Distributed/Features/DistributedRuntimeFeature.cs @@ -0,0 +1,39 @@ +using Elsa.Extensions; +using Elsa.Features.Abstractions; +using Elsa.Features.Attributes; +using Elsa.Features.Services; +using Elsa.Workflows.Runtime.Features; +using Elsa.Workflows.Runtime.Handlers; +using Microsoft.Extensions.DependencyInjection; + +namespace Elsa.Workflows.Runtime.Distributed.Features; + +/// +/// Installs and configures workflow runtime features. +/// +[DependsOn(typeof(WorkflowRuntimeFeature))] +public class DistributedRuntimeFeature : FeatureBase +{ + /// + public DistributedRuntimeFeature(IModule module) : base(module) + { + } + + public override void Configure() + { + Module.UseWorkflowRuntime(runtime => + { + runtime.WorkflowRuntime = sp => sp.GetRequiredService(); + runtime.BookmarkQueueWorker = sp => sp.GetRequiredService(); + }); + } + + /// + public override void Apply() + { + Services + .AddScoped() + .AddScoped() + .AddCommandHandler(); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime.Distributed/FodyWeavers.xml b/src/modules/Elsa.Workflows.Runtime.Distributed/FodyWeavers.xml new file mode 100644 index 000000000..00e1d9a1c --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime.Distributed/FodyWeavers.xml @@ -0,0 +1,3 @@ + + + \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedBookmarkQueueWorker.cs b/src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedBookmarkQueueWorker.cs new file mode 100644 index 000000000..7149728fd --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedBookmarkQueueWorker.cs @@ -0,0 +1,25 @@ +using Medallion.Threading; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging; + +namespace Elsa.Workflows.Runtime.Distributed; + +public class DistributedBookmarkQueueWorker( + IDistributedLockProvider distributedLockProvider, + IBookmarkQueueSignaler signaler, + IServiceScopeFactory scopeFactory, + ILogger logger) : BookmarkQueueWorker(signaler, scopeFactory, logger) +{ + protected override async Task ProcessAsync(CancellationToken cancellationToken) + { + await using var handle = await distributedLockProvider.TryAcquireLockAsync(nameof(DistributedBookmarkQueueWorker), TimeSpan.Zero, cancellationToken); + + if (handle == null) + { + logger.LogInformation("Could not acquire lock for distributed bookmark queue worker. This is usually an indication that another application instance is already processing."); + return; + } + + await base.ProcessAsync(cancellationToken); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowClient.cs b/src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowClient.cs new file mode 100644 index 000000000..251ac52dd --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowClient.cs @@ -0,0 +1,86 @@ +using Elsa.Common.DistributedHosting; +using Elsa.Workflows.Runtime.Messages; +using Elsa.Workflows.State; +using Medallion.Threading; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Options; + +namespace Elsa.Workflows.Runtime.Distributed; + +public class DistributedWorkflowClient( + string workflowInstanceId, + IDistributedLockProvider distributedLockProvider, + IOptions distributedLockingOptions, + IServiceProvider serviceProvider) + : IWorkflowClient +{ + private readonly LocalWorkflowClient _localWorkflowClient = ActivatorUtilities.CreateInstance(serviceProvider, workflowInstanceId); + + public string WorkflowInstanceId => workflowInstanceId; + + public async Task CreateInstanceAsync(CreateWorkflowInstanceRequest request, CancellationToken cancellationToken = default) + { + return await _localWorkflowClient.CreateInstanceAsync(request, cancellationToken); + } + + public async Task RunInstanceAsync(RunWorkflowInstanceRequest request, CancellationToken cancellationToken = default) + { + var result = await WithLockAsync(async () => await _localWorkflowClient.RunInstanceAsync(request, cancellationToken)); + return result; + } + + public async Task CreateAndRunInstanceAsync(CreateAndRunWorkflowInstanceRequest request, CancellationToken cancellationToken = default) + { + var createRequest = new CreateWorkflowInstanceRequest + { + Properties = request.Properties, + CorrelationId = request.CorrelationId, + Name = request.Name, + Input = request.Input, + WorkflowDefinitionHandle = request.WorkflowDefinitionHandle, + ParentId = request.ParentId + }; + var workflowInstance = await _localWorkflowClient.CreateInstanceInternalAsync(createRequest, cancellationToken); + + // We need to lock newly created workflow instances too, because it might dispatch child workflows that attempt to resume the parent workflow. + // For example, when using a DispatchWorkflow activity configured to wait for the dispatched workflow to complete. + return await WithLockAsync(async () => await _localWorkflowClient.RunInstanceAsync(workflowInstance, new() + { + Input = request.Input, + Variables = request.Variables, + Properties = request.Properties, + TriggerActivityId = request.TriggerActivityId, + ActivityHandle = request.ActivityHandle, + IncludeWorkflowOutput = request.IncludeWorkflowOutput + }, cancellationToken)); + } + + public async Task CancelAsync(CancellationToken cancellationToken = default) + { + await _localWorkflowClient.CancelAsync(cancellationToken); + } + + public async Task ExportStateAsync(CancellationToken cancellationToken = default) + { + return await _localWorkflowClient.ExportStateAsync(cancellationToken); + } + + public async Task ImportStateAsync(WorkflowState workflowState, CancellationToken cancellationToken = default) + { + await _localWorkflowClient.ImportStateAsync(workflowState, cancellationToken); + } + + public async Task InstanceExistsAsync(CancellationToken cancellationToken = default) + { + return await _localWorkflowClient.InstanceExistsAsync(cancellationToken); + } + + private async Task WithLockAsync(Func> func) + { + var lockKey = $"workflow-instance:{WorkflowInstanceId}"; + var lockTimeout = distributedLockingOptions.Value.LockAcquisitionTimeout; + await using var @lock = await distributedLockProvider.AcquireLockAsync(lockKey, lockTimeout); + var result = await func(); + return result; + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowRuntime.Obsolete.cs b/src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowRuntime.Obsolete.cs new file mode 100644 index 000000000..94aebcf83 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowRuntime.Obsolete.cs @@ -0,0 +1,32 @@ +using Elsa.Workflows.Runtime.Entities; +using Elsa.Workflows.Runtime.Filters; +using Elsa.Workflows.Runtime.Matches; +using Elsa.Workflows.Runtime.Options; +using Elsa.Workflows.Runtime.Parameters; +using Elsa.Workflows.Runtime.Params; +using Elsa.Workflows.Runtime.Requests; +using Elsa.Workflows.Runtime.Results; +using Elsa.Workflows.State; + +namespace Elsa.Workflows.Runtime.Distributed; + +public partial class DistributedWorkflowRuntime +{ + private readonly Lazy _obsoleteApi; + private ObsoleteWorkflowRuntime ObsoleteApi => _obsoleteApi.Value; + + public Task CanStartWorkflowAsync(string definitionId, StartWorkflowRuntimeParams? options = null) => ObsoleteApi.CanStartWorkflowAsync(definitionId, options); + public Task StartWorkflowAsync(string definitionId, StartWorkflowRuntimeParams? options = null) => ObsoleteApi.StartWorkflowAsync(definitionId, options); + public Task> StartWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions? options = null) => ObsoleteApi.StartWorkflowsAsync(activityTypeName, bookmarkPayload, options); + public Task TryStartWorkflowAsync(string definitionId, StartWorkflowRuntimeParams? options = null) => ObsoleteApi.TryStartWorkflowAsync(definitionId, options); + public Task ResumeWorkflowAsync(string workflowInstanceId, ResumeWorkflowRuntimeParams? options = null) => ObsoleteApi.ResumeWorkflowAsync(workflowInstanceId, options); + public Task> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions? options = null) => ObsoleteApi.ResumeWorkflowsAsync(activityTypeName, bookmarkPayload, options); + public Task TriggerWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsOptions? options = null) => ObsoleteApi.TriggerWorkflowsAsync(activityTypeName, bookmarkPayload, options); + public Task ExecuteWorkflowAsync(WorkflowMatch match, ExecuteWorkflowParams? options = null) => ObsoleteApi.ExecuteWorkflowAsync(match, options); + public Task CancelWorkflowAsync(string workflowInstanceId, CancellationToken cancellationToken = default) => ObsoleteApi.CancelWorkflowAsync(workflowInstanceId, cancellationToken); + public Task> FindWorkflowsAsync(WorkflowsFilter filter, CancellationToken cancellationToken = default) => ObsoleteApi.FindWorkflowsAsync(filter, cancellationToken); + public Task ExportWorkflowStateAsync(string workflowInstanceId, CancellationToken cancellationToken = default) => ObsoleteApi.ExportWorkflowStateAsync(workflowInstanceId, cancellationToken); + public Task ImportWorkflowStateAsync(WorkflowState workflowState, CancellationToken cancellationToken = default) => ObsoleteApi.ImportWorkflowStateAsync(workflowState, cancellationToken); + public Task UpdateBookmarkAsync(StoredBookmark bookmark, CancellationToken cancellationToken = default) => ObsoleteApi.UpdateBookmarkAsync(bookmark, cancellationToken); + public Task CountRunningWorkflowsAsync(CountRunningWorkflowsRequest request, CancellationToken cancellationToken = default) => ObsoleteApi.CountRunningWorkflowsAsync(request, cancellationToken); +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowRuntime.cs b/src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowRuntime.cs new file mode 100644 index 000000000..6c932fe69 --- /dev/null +++ b/src/modules/Elsa.Workflows.Runtime.Distributed/Services/DistributedWorkflowRuntime.cs @@ -0,0 +1,36 @@ +using Microsoft.Extensions.DependencyInjection; + +namespace Elsa.Workflows.Runtime.Distributed; + +/// +/// Represents a distributed workflow runtime that can create instances connected to a workflow instance. +/// +public partial class DistributedWorkflowRuntime : IWorkflowRuntime +{ + private readonly IServiceProvider _serviceProvider; + private readonly IIdentityGenerator _identityGenerator; + + /// + /// Represents a distributed workflow runtime that can create instances connected to a workflow instance. + /// + public DistributedWorkflowRuntime(IServiceProvider serviceProvider, IIdentityGenerator identityGenerator) + { + _serviceProvider = serviceProvider; + _identityGenerator = identityGenerator; + _obsoleteApi = new(() => ObsoleteWorkflowRuntime.Create(serviceProvider, CreateClientAsync)); + } + + /// + public async ValueTask CreateClientAsync(CancellationToken cancellationToken = default) + { + return await CreateClientAsync(null, cancellationToken); + } + + /// + public ValueTask CreateClientAsync(string? workflowInstanceId, CancellationToken cancellationToken = default) + { + workflowInstanceId ??= _identityGenerator.GenerateId(); + var client = (IWorkflowClient)ActivatorUtilities.CreateInstance(_serviceProvider, typeof(DistributedWorkflowClient), workflowInstanceId); + return new(client); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Services/WorkflowDefinitionRefresher.cs b/src/modules/Elsa.Workflows.Runtime/Services/WorkflowDefinitionRefresher.cs index f263229e9..6109c4400 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/WorkflowDefinitionRefresher.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/WorkflowDefinitionRefresher.cs @@ -46,7 +46,7 @@ public class WorkflowDefinitionsRefresher(IWorkflowDefinitionStore store, ITrigg var processedWorkflowDefinitionIds = processedWorkflowDefinitions.Select(x => x.DefinitionId).ToList(); var notification = new WorkflowDefinitionsRefreshed(processedWorkflowDefinitionIds); await notificationSender.SendAsync(notification, cancellationToken); - return new RefreshWorkflowDefinitionsResponse(processedWorkflowDefinitionIds, request.DefinitionIds?.Except(processedWorkflowDefinitionIds)?.ToList() ?? []); + return new(processedWorkflowDefinitionIds, request.DefinitionIds?.Except(processedWorkflowDefinitionIds)?.ToList() ?? []); } private async Task IndexWorkflowTriggersAsync(IEnumerable definitions, CancellationToken cancellationToken) diff --git a/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj b/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj index 7a1d27571..03a554d2a 100644 --- a/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj +++ b/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj @@ -22,6 +22,8 @@ + + diff --git a/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/WorkflowServer.cs b/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/WorkflowServer.cs index 1dfe08fd3..eb9c4ffad 100644 --- a/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/WorkflowServer.cs +++ b/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/WorkflowServer.cs @@ -1,3 +1,4 @@ +using System.Reflection; using Elsa.Caching; using Elsa.Extensions; using Elsa.Identity.Providers; @@ -6,6 +7,8 @@ using Elsa.Workflows.ComponentTests.Decorators; using Elsa.Workflows.ComponentTests.Materializers; using Elsa.Workflows.ComponentTests.WorkflowProviders; using Elsa.Workflows.Management; +using Elsa.Workflows.Runtime.Distributed.Extensions; +using FluentStorage; using JetBrains.Annotations; using Microsoft.AspNetCore.Hosting; using Microsoft.AspNetCore.Mvc.Testing; @@ -46,9 +49,28 @@ public class WorkflowServer(Infrastructure infrastructure, string url) : WebAppl elsa.AddWorkflowsFrom(); elsa.AddActivitiesFrom(); elsa.UseDefaultAuthentication(defaultAuthentication => defaultAuthentication.UseAdminApiKey()); + elsa.UseFluentStorageProvider(sp => + { + var assemblyLocation = Assembly.GetExecutingAssembly().Location; + var assemblyDirectory = Path.GetDirectoryName(assemblyLocation)!; + var workflowsDirectorySegments = new[] + { + assemblyDirectory, "Scenarios" + }; + var workflowsDirectory = Path.Join(workflowsDirectorySegments); + return StorageFactory.Blobs.DirectoryFiles(workflowsDirectory); + }); elsa.UseIdentity(); elsa.UseWorkflowManagement(); - elsa.UseWorkflowRuntime(); + elsa.UseWorkflowRuntime(runtime => runtime.UseDistributedRuntime()); + elsa.UseJavaScript(options => + { + options.AllowClrAccess = true; + options.ConfigureEngine(engine => + { + engine.SetValue("getStaticValue", () => StaticValueHolder.Value); + }); + }); elsa.UseHttp(); }; } diff --git a/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionRefresh/DynamicEndpointTests.cs b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionRefresh/DynamicEndpointTests.cs index 32484c686..f1192f918 100644 --- a/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionRefresh/DynamicEndpointTests.cs +++ b/test/component/Elsa.Workflows.ComponentTests/Scenarios/WorkflowDefinitionRefresh/DynamicEndpointTests.cs @@ -22,13 +22,13 @@ public class DynamicEndpointTests : AppComponentTest { var client = WorkflowServer.CreateHttpWorkflowClient(); - var firstResponse = await client.SendAsync(new HttpRequestMessage(HttpMethod.Get, "first-value")); + var firstResponse = await client.SendAsync(new(HttpMethod.Get, "first-value")); StaticValueHolder.Value = "second-value"; - var _ = await _workflowDefinitionsRefresher.RefreshWorkflowDefinitionsAsync( - new Runtime.Requests.RefreshWorkflowDefinitionsRequest() { DefinitionIds = ["f69f061159adc3ae"] }, CancellationToken.None); + _ = await _workflowDefinitionsRefresher.RefreshWorkflowDefinitionsAsync( + new() { DefinitionIds = ["f69f061159adc3ae"] }, CancellationToken.None); - var secondResponse = await client.SendAsync(new HttpRequestMessage(HttpMethod.Get, "second-value")); + var secondResponse = await client.SendAsync(new(HttpMethod.Get, "second-value")); Assert.Equal(HttpStatusCode.OK, firstResponse.StatusCode); Assert.Equal(HttpStatusCode.OK, secondResponse.StatusCode); diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Elsa.Workflows.IntegrationTests.csproj b/test/integration/Elsa.Workflows.IntegrationTests/Elsa.Workflows.IntegrationTests.csproj index d52f2709c..eb322b490 100644 --- a/test/integration/Elsa.Workflows.IntegrationTests/Elsa.Workflows.IntegrationTests.csproj +++ b/test/integration/Elsa.Workflows.IntegrationTests/Elsa.Workflows.IntegrationTests.csproj @@ -10,6 +10,7 @@ + diff --git a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/RunAsynchronousActivityOutput/Tests.cs b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/RunAsynchronousActivityOutput/Tests.cs index d454b5549..ffe94bfab 100644 --- a/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/RunAsynchronousActivityOutput/Tests.cs +++ b/test/integration/Elsa.Workflows.IntegrationTests/Scenarios/RunAsynchronousActivityOutput/Tests.cs @@ -3,6 +3,7 @@ using Elsa.Testing.Shared; using Elsa.Workflows.Activities; using Elsa.Workflows.IntegrationTests.Scenarios.RunAsynchronousActivityOutput.Activities; using Elsa.Workflows.Memory; +using Elsa.Workflows.Runtime.Distributed; using Elsa.Workflows.Runtime.Stores; using Microsoft.Extensions.DependencyInjection; @@ -115,11 +116,16 @@ public class Tests // Act var workflowFinishedRecord = await workflow.DispatchWorkflowAndRunToCompletion( + configureServices: services => + { + services.AddScoped(); + }, configureElsa: elsa => { elsa.UseWorkflowRuntime(workflowRuntime => { workflowRuntime.ActivityExecutionLogStore = sp => activityExecutionStore; + workflowRuntime.WorkflowRuntime = sp => sp.GetRequiredService(); }); }); diff --git a/test/unit/Elsa.Http.UnitTests/ContentWriters/RawStringContentTests.cs b/test/unit/Elsa.Http.UnitTests/ContentWriters/RawStringContentTests.cs index 308f22c37..da7c75768 100644 --- a/test/unit/Elsa.Http.UnitTests/ContentWriters/RawStringContentTests.cs +++ b/test/unit/Elsa.Http.UnitTests/ContentWriters/RawStringContentTests.cs @@ -27,29 +27,4 @@ public class RawStringContentTests Assert.Equal(contentType, rawContent.Headers.ContentType?.MediaType); Assert.Null(rawContent.Headers.ContentType?.CharSet); } - - /// - /// Tests that the content type with parameters is preserved exactly as provided. - /// - [Fact] - public void ContentType_WithParameters_ShouldPreserveParameters() - { - // Arrange - const string contentType = "application/json; custom-param=value"; - var expectedMediaType = new MediaTypeHeaderValue(contentType); - const string content = "{\"test\": \"value\"}"; - - // Act - var rawContent = new RawStringContent(content, Encoding.UTF8, contentType); - - // Assert - Assert.Equal(expectedMediaType.MediaType, rawContent.Headers.ContentType?.MediaType); - Assert.Equal(expectedMediaType.Parameters.Count(), rawContent.Headers.ContentType?.Parameters.Count()); - - var expectedParam = expectedMediaType.Parameters.First(); - var actualParam = rawContent.Headers.ContentType?.Parameters.First(); - - Assert.Equal(expectedParam.Name, actualParam?.Name); - Assert.Equal(expectedParam.Value, actualParam?.Value); - } } \ No newline at end of file