From b6689c7e2b900ac000141229a67dd5feb582da74 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Wed, 1 May 2024 14:53:37 +0200 Subject: [PATCH] Fix Race Condition in WorkflowDefinitionActivity (#5315) * Refactor hash method in Hasher.cs The Hash method in the Hasher.cs file has been refactored to use more efficient code. Instead of creating a SHA256 instance, it directly hashes the data using SHA256.HashData. As a result, the private Hash method and its usage, which only has relevance to the previous approach, have been removed. * Replace GetHashCode with SHA256 for script caching Two unused namespaces, Esprima and Esprima.Ast, have been removed, and two new ones, System.Security.Cryptography and System.Text, have been added. These changes in Elsa.JavaScript's JintJavaScriptEvaluator file were made to replace the basic GetHashCode method, which was used earlier to generate a cache key for JavaScript scripts, with SHA256 hash function to ensure uniqueness and avoid possible hash collisions. This greatly increases the reliability of the caching mechanism. * Set HTTP client timeout in WorkflowServer The client timeout in WorkflowServer has been set to 1 minute for `ResrService` and `client` methods. This was done to manage long running requests and prevent timeouts. * Refactor code and add validation in ScheduleActivityAsync methods An unnecessary line of code in the ScheduleActivityAsync method under the ActivityExecutionContext file has been removed to simplify the function. Meanwhile, validation has been added to ensure that specified activities are part of the workflow. This improves code clarity and prevents potential errors caused by incorrect activity scheduling. * Refactor ActivityVisitor to use ActivityVisitorContext The ActivityVisitor class in Elsa.Workflows.Core was refactored to use an ActivityVisitorContext class. The introduction of this context class replaced the multiple hashsets that were used within function signatures, consolidating them into a single object and reducing function complexity. This change improves readability and code organization. * Refactor WorkflowDefinitionActivity to use WorkflowGraph Adjusted the WorkflowDefinitionActivity class in the Elsa.Workflows.Management module. The changes include migration from using Workflow to WorkflowGraph objects and condensing code blocks for clearer readability. Additionally, a new 'IsInitialized' field is added to eliminate potential race conditions during the graph construction process. * Add activity validation to WorkflowExecutionContext A check has been implemented in the WorkflowExecutionContextExtensions to validate that a specified activity is part of the workflow. This prevents incorrect activity references when scheduling tasks in the workflow. * Refactor to use WorkflowGraph instead of Workflow The codebase has been refactored to utilize the WorkflowGraph instead of the Workflow while running and manipulating workflows. Additional changes include restructuring WorkflowRunner and WorkflowDefinitionService classes, introducing WorkflowGraphBuilder usage, and mapping updates in WorkflowDefinitionMapper. The WorkflowHost, WorkflowDispatcher, and WorkflowExecutionContext have also been updated accordingly. * Add WorkflowGraph and related services This commit introduces the IWorkflowGraphBuilder interface, the WorkflowGraph model, and an implementation of the interface in the WorkflowGraphBuilder class. The purpose of these additions is to establish the building and structure of a workflow graph. The WorkflowGraph model also includes activity node handling and hashing capabilities. * Add tests to ensure exception when scheduling an activity not part of workflow Several unit tests were written to ensure that the right behavior is exhibited when scheduling an activity that is not part of the workflow. An exception is expected to be thrown in this case. Additionally, new service definitions, workflow definitions and workflow queries were added for more comprehensive testing. * Update WorkflowGraph class and add comment descriptions This commit updates the WorkflowGraph class by extending its descriptions and implementing an explicit mention to the Workflow reference. Additionally, more attribute descriptions have been added to increase code readability and comprehension. * Removed unnecessary brackets * Remove unused services from workflow management This commit deletes "ExpressionDescriptorRegistryPopulator.cs" and "ScopedWorkflowDefinitionLookup.cs" files from Elsa.Workflows.Management.Services. These files, containing obsolete services, are no longer used in the workflow management process. The services' registration has been removed as well from "WorkflowManagementFeature.cs". * Remove TestWorkflowDefinitionService.cs from component tests A file, TestWorkflowDefinitionService.cs, was removed under Elsa.Workflows.ComponentTests. This is part of the improvement process where inefficient or unnecessary test files are cleaned up. * Add comments to WorkflowDefinitionActivityTests This commit adds explanatory comments to the WorkflowDefinitionActivityTests file. The comments provide information about the purpose of these tests and a reference to a related issue in the project's Github repository. * Updated duplicate and missing package references --------- Co-authored-by: Raymond den Haan --- Directory.Packages.props | 1 + src/bundles/Elsa.Server.Web/Program.cs | 3 + .../AlterationHandlers/MigrateHandler.cs | 2 +- .../Services/DefaultAlterationRunner.cs | 13 +- .../Middleware/HttpWorkflowsMiddleware.cs | 14 +- .../Models/HttpWorkflowLookupResult.cs | 7 +- .../CachingHttpWorkflowLookupService.cs | 4 +- .../Services/HttpBookmarkProcessor.cs | 2 +- .../Services/HttpWorkflowLookupService.cs | 12 +- .../Endpoints/TypeDefinitions/Endpoint.cs | 11 +- .../Services/JintJavaScriptEvaluator.cs | 13 +- .../Services/MassTransitWorkflowDispatcher.cs | 11 +- .../Grains/WorkflowInstance.cs | 14 +- .../Services/ProtoActorWorkflowRuntime.cs | 5 +- .../Contexts/ActivityExecutionContext.cs | 11 +- .../Contexts/WorkflowExecutionContext.cs | 127 ++++++++-------- .../Contracts/IWorkflowGraphBuilder.cs | 15 ++ .../Contracts/IWorkflowRunner.cs | 2 + .../WorkflowExecutionContextExtensions.cs | 4 + .../Features/WorkflowsFeature.cs | 1 + .../Models/WorkflowGraph.cs | 62 ++++++++ .../Services/ActivityVisitor.cs | 36 ++++- .../Elsa.Workflows.Core/Services/Hasher.cs | 10 +- .../Services/WorkflowGraphBuilder.cs | 20 +++ .../Services/WorkflowRunner.cs | 83 +++++----- .../WorkflowDefinitionActivity.cs | 45 ++++-- .../WorkflowDefinitionActivityPortResolver.cs | 2 +- .../Contracts/IWorkflowDefinitionService.cs | 9 +- .../Mappers/WorkflowDefinitionMapper.cs | 7 +- .../CachingWorkflowDefinitionService.cs | 22 +-- .../ExpressionDescriptorRegistryPopulator.cs | 29 ---- .../Services/WorkflowDefinitionPublisher.cs | 4 +- .../Services/WorkflowDefinitionService.cs | 64 +++----- .../Services/WorkflowInstanceManager.cs | 1 - .../Contracts/ITriggerIndexer.cs | 13 -- .../Contracts/IWorkflowHost.cs | 10 +- .../Contracts/IWorkflowHostFactory.cs | 11 +- .../DefaultBackgroundActivityInvoker.cs | 6 +- ...DefaultWorkflowDefinitionStorePopulator.cs | 2 +- .../Services/DefaultWorkflowRuntime.cs | 12 +- .../Services/TriggerIndexer.cs | 26 +--- .../Services/WorkflowHost.cs | 15 +- .../Services/WorkflowHostFactory.cs | 16 +- .../Elsa.Workflows.ComponentTests.csproj | 9 ++ .../Helpers/Fixtures/WorkflowServer.cs | 2 + .../WorkflowDefinitionActivityTests.cs | 53 +++++++ .../Workflows/workflow-definition-child.json | 64 ++++++++ .../workflow-definition-grand-child.json | 67 +++++++++ .../Workflows/workflow-definition-parent.json | 142 ++++++++++++++++++ .../Elsa.Workflows.PerformanceTests.csproj | 7 - .../ScheduleActivityExecutionContextTests.cs | 30 ++++ 51 files changed, 799 insertions(+), 352 deletions(-) create mode 100644 src/modules/Elsa.Workflows.Core/Contracts/IWorkflowGraphBuilder.cs create mode 100644 src/modules/Elsa.Workflows.Core/Models/WorkflowGraph.cs create mode 100644 src/modules/Elsa.Workflows.Core/Services/WorkflowGraphBuilder.cs delete mode 100644 src/modules/Elsa.Workflows.Management/Services/ExpressionDescriptorRegistryPopulator.cs create mode 100644 test/component/Elsa.Workflows.ComponentTests/Scenarios/CachingAndWorkflowDefinitionActivity/WorkflowDefinitionActivityTests.cs create mode 100644 test/component/Elsa.Workflows.ComponentTests/Scenarios/CachingAndWorkflowDefinitionActivity/Workflows/workflow-definition-child.json create mode 100644 test/component/Elsa.Workflows.ComponentTests/Scenarios/CachingAndWorkflowDefinitionActivity/Workflows/workflow-definition-grand-child.json create mode 100644 test/component/Elsa.Workflows.ComponentTests/Scenarios/CachingAndWorkflowDefinitionActivity/Workflows/workflow-definition-parent.json create mode 100644 test/unit/Elsa.Workflows.Core.UnitTests/ScheduleActivityExecutionContextTests.cs diff --git a/Directory.Packages.props b/Directory.Packages.props index aa8a44a80..67f22ced0 100644 --- a/Directory.Packages.props +++ b/Directory.Packages.props @@ -91,6 +91,7 @@ + diff --git a/src/bundles/Elsa.Server.Web/Program.cs b/src/bundles/Elsa.Server.Web/Program.cs index 93c584525..b5647e504 100644 --- a/src/bundles/Elsa.Server.Web/Program.cs +++ b/src/bundles/Elsa.Server.Web/Program.cs @@ -1,6 +1,7 @@ using System.Text.Encodings.Web; using Elsa.Alterations.Extensions; using Elsa.Alterations.MassTransit.Extensions; +using Elsa.Caching.Options; using Elsa.Common.DistributedLocks.Noop; using Elsa.Dapper.Extensions; using Elsa.Dapper.Services; @@ -358,6 +359,8 @@ services ConfigureForTest?.Invoke(elsa); }); +services.Configure(options => options.CacheDuration = TimeSpan.FromDays(1)); + services.AddHealthChecks(); services.AddControllers(); services.AddCors(cors => cors.AddDefaultPolicy(policy => policy.AllowAnyHeader().AllowAnyMethod().AllowAnyOrigin().WithExposedHeaders("*"))); diff --git a/src/modules/Elsa.Alterations/AlterationHandlers/MigrateHandler.cs b/src/modules/Elsa.Alterations/AlterationHandlers/MigrateHandler.cs index 8acfabb2c..f5a9a0c15 100644 --- a/src/modules/Elsa.Alterations/AlterationHandlers/MigrateHandler.cs +++ b/src/modules/Elsa.Alterations/AlterationHandlers/MigrateHandler.cs @@ -30,7 +30,7 @@ public class MigrateHandler : AlterationHandlerBase } var targetWorkflow = await workflowDefinitionService.MaterializeWorkflowAsync(targetWorkflowDefinition, cancellationToken); - await context.WorkflowExecutionContext.SetWorkflowAsync(targetWorkflow); + await context.WorkflowExecutionContext.SetWorkflowGraphAsync(targetWorkflow); context.Succeed(); } diff --git a/src/modules/Elsa.Alterations/Services/DefaultAlterationRunner.cs b/src/modules/Elsa.Alterations/Services/DefaultAlterationRunner.cs index efd078714..8108407fd 100644 --- a/src/modules/Elsa.Alterations/Services/DefaultAlterationRunner.cs +++ b/src/modules/Elsa.Alterations/Services/DefaultAlterationRunner.cs @@ -61,29 +61,24 @@ public class DefaultAlterationRunner : IAlterationRunner { var log = new AlterationLog(_systemClock); var result = new RunAlterationsResult(workflowInstanceId, log); - - // Load workflow instance. var workflowState = await _workflowRuntime.ExportWorkflowStateAsync(workflowInstanceId, cancellationToken); - // If the workflow instance is not found, log an error and continue. if (workflowState == null) { log.Add($"Workflow instance with ID '{workflowInstanceId}' not found.", LogLevel.Error); return result; } - // Load workflow definition. - var workflow = await _workflowDefinitionService.FindWorkflowAsync(workflowState.DefinitionVersionId, cancellationToken); + var workflowGraph = await _workflowDefinitionService.FindWorkflowGraphAsync(workflowState.DefinitionVersionId, cancellationToken); - // If the workflow definition is not found, log an error and continue. - if (workflow == null) + if (workflowGraph == null) { log.Add($"Workflow definition with ID '{workflowState.DefinitionVersionId}' not found.", LogLevel.Error); return result; } - + // Create workflow execution context. - var workflowExecutionContext = await WorkflowExecutionContext.CreateAsync(_serviceProvider, workflow, workflowState, cancellationTokens: cancellationToken); + var workflowExecutionContext = await WorkflowExecutionContext.CreateAsync(_serviceProvider, workflowGraph, workflowState, cancellationTokens: cancellationToken); workflowExecutionContext.TransientProperties.Add(RunAlterationsMiddleware.AlterationsPropertyKey, alterations); workflowExecutionContext.TransientProperties.Add(RunAlterationsMiddleware.AlterationsLogPropertyKey, log); diff --git a/src/modules/Elsa.Http/Middleware/HttpWorkflowsMiddleware.cs b/src/modules/Elsa.Http/Middleware/HttpWorkflowsMiddleware.cs index a6db8b954..28ec0cd9f 100644 --- a/src/modules/Elsa.Http/Middleware/HttpWorkflowsMiddleware.cs +++ b/src/modules/Elsa.Http/Middleware/HttpWorkflowsMiddleware.cs @@ -87,8 +87,8 @@ public class HttpWorkflowsMiddleware(RequestDelegate next, IOptions FindWorkflowAsync(IServiceProvider serviceProvider, StoredTrigger trigger, CancellationToken cancellationToken) + private async Task FindWorkflowGraphAsync(IServiceProvider serviceProvider, StoredTrigger trigger, CancellationToken cancellationToken) { var workflowDefinitionService = serviceProvider.GetRequiredService(); var workflowDefinitionId = trigger.WorkflowDefinitionVersionId; - return await workflowDefinitionService.FindWorkflowAsync(workflowDefinitionId, cancellationToken); + return await workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionId, cancellationToken); } - private async Task StartWorkflowAsync(HttpContext httpContext, StoredTrigger trigger, Workflow workflow, IDictionary input) + private async Task StartWorkflowAsync(HttpContext httpContext, StoredTrigger trigger, WorkflowGraph workflowGraph, IDictionary input) { var serviceProvider = httpContext.RequestServices; var cancellationToken = httpContext.RequestAborted; var bookmarkPayload = trigger.GetPayload(); var workflowHostFactory = serviceProvider.GetRequiredService(); - var workflowHost = await workflowHostFactory.CreateAsync(workflow, cancellationToken); + var workflowHost = await workflowHostFactory.CreateAsync(workflowGraph, cancellationToken); if (await AuthorizeAsync(serviceProvider, httpContext, workflowHost.Workflow, bookmarkPayload, cancellationToken)) return; @@ -175,7 +175,7 @@ public class HttpWorkflowsMiddleware(RequestDelegate next, IOptions(); - var workflow = await workflowDefinitionService.FindWorkflowAsync(workflowInstance.DefinitionVersionId, cancellationToken); + var workflow = await workflowDefinitionService.FindWorkflowGraphAsync(workflowInstance.DefinitionVersionId, cancellationToken); if (workflow == null) { diff --git a/src/modules/Elsa.Http/Models/HttpWorkflowLookupResult.cs b/src/modules/Elsa.Http/Models/HttpWorkflowLookupResult.cs index 2df870fa0..64a578901 100644 --- a/src/modules/Elsa.Http/Models/HttpWorkflowLookupResult.cs +++ b/src/modules/Elsa.Http/Models/HttpWorkflowLookupResult.cs @@ -1,6 +1,9 @@ -using Elsa.Workflows.Activities; +using Elsa.Workflows.Models; using Elsa.Workflows.Runtime.Entities; namespace Elsa.Http.Models; -public record HttpWorkflowLookupResult(Workflow? Workflow, ICollection Triggers); \ No newline at end of file +/// +/// Represents the result of a workflow lookup. +/// +public record HttpWorkflowLookupResult(WorkflowGraph? WorkflowGraph, ICollection Triggers); \ No newline at end of file diff --git a/src/modules/Elsa.Http/Services/CachingHttpWorkflowLookupService.cs b/src/modules/Elsa.Http/Services/CachingHttpWorkflowLookupService.cs index 1f0061724..25169a64b 100644 --- a/src/modules/Elsa.Http/Services/CachingHttpWorkflowLookupService.cs +++ b/src/modules/Elsa.Http/Services/CachingHttpWorkflowLookupService.cs @@ -29,8 +29,8 @@ public class CachingHttpWorkflowLookupService( if (result == null) return null; - var workflow = result.Workflow!; - var changeTokenKey = cacheManager.GetWorkflowChangeTokenKey(workflow.Identity.DefinitionId); + var workflowGraph = result.WorkflowGraph!; + var changeTokenKey = cacheManager.GetWorkflowChangeTokenKey(workflowGraph.Workflow.Identity.DefinitionId); var changeToken = cache.GetToken(changeTokenKey); entry.AddExpirationToken(changeToken); diff --git a/src/modules/Elsa.Http/Services/HttpBookmarkProcessor.cs b/src/modules/Elsa.Http/Services/HttpBookmarkProcessor.cs index 01bb26c83..a13b8e888 100644 --- a/src/modules/Elsa.Http/Services/HttpBookmarkProcessor.cs +++ b/src/modules/Elsa.Http/Services/HttpBookmarkProcessor.cs @@ -80,7 +80,7 @@ public class HttpBookmarkProcessor : IHttpBookmarkProcessor continue; } - var workflow = await _workflowDefinitionService.FindWorkflowAsync( + var workflow = await _workflowDefinitionService.FindWorkflowGraphAsync( workflowState.DefinitionId, VersionOptions.SpecificVersion(workflowState.DefinitionVersion), systemCancellationToken); diff --git a/src/modules/Elsa.Http/Services/HttpWorkflowLookupService.cs b/src/modules/Elsa.Http/Services/HttpWorkflowLookupService.cs index 2eb860298..d8287e4cd 100644 --- a/src/modules/Elsa.Http/Services/HttpWorkflowLookupService.cs +++ b/src/modules/Elsa.Http/Services/HttpWorkflowLookupService.cs @@ -1,7 +1,7 @@ using Elsa.Http.Contracts; using Elsa.Http.Models; -using Elsa.Workflows.Activities; using Elsa.Workflows.Management.Contracts; +using Elsa.Workflows.Models; using Elsa.Workflows.Runtime.Contracts; using Elsa.Workflows.Runtime.Entities; using Elsa.Workflows.Runtime.Filters; @@ -25,12 +25,12 @@ public class HttpWorkflowLookupService(ITriggerStore triggerStore, IWorkflowDefi if (trigger == null) return default; - var workflow = await FindWorkflowAsync(trigger, cancellationToken); + var workflowGraph = await FindWorkflowGraphAsync(trigger, cancellationToken); - if (workflow == null) + if (workflowGraph == null) return default; - return new(workflow, triggers); + return new(workflowGraph, triggers); } private async Task> FindTriggersAsync(string bookmarkHash, CancellationToken cancellationToken) @@ -42,9 +42,9 @@ public class HttpWorkflowLookupService(ITriggerStore triggerStore, IWorkflowDefi return await triggerStore.FindManyAsync(triggerFilter, cancellationToken); } - private async Task FindWorkflowAsync(StoredTrigger trigger, CancellationToken cancellationToken) + private async Task FindWorkflowGraphAsync(StoredTrigger trigger, CancellationToken cancellationToken) { var workflowDefinitionVersionId = trigger.WorkflowDefinitionVersionId; - return await workflowDefinitionService.FindWorkflowAsync(workflowDefinitionVersionId, cancellationToken); + return await workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionVersionId, cancellationToken); } } \ No newline at end of file diff --git a/src/modules/Elsa.JavaScript/Endpoints/TypeDefinitions/Endpoint.cs b/src/modules/Elsa.JavaScript/Endpoints/TypeDefinitions/Endpoint.cs index a496405fb..3e966bf94 100644 --- a/src/modules/Elsa.JavaScript/Endpoints/TypeDefinitions/Endpoint.cs +++ b/src/modules/Elsa.JavaScript/Endpoints/TypeDefinitions/Endpoint.cs @@ -3,8 +3,8 @@ using Elsa.Abstractions; using Elsa.Common.Models; using Elsa.JavaScript.TypeDefinitions.Contracts; using Elsa.JavaScript.TypeDefinitions.Models; -using Elsa.Workflows.Activities; using Elsa.Workflows.Management.Contracts; +using Elsa.Workflows.Models; using JetBrains.Annotations; namespace Elsa.JavaScript.Endpoints.TypeDefinitions; @@ -35,15 +35,16 @@ internal class Get : ElsaEndpoint /// public override async Task HandleAsync(Request request, CancellationToken cancellationToken) { - var workflow = await GetWorkflowAsync(request.WorkflowDefinitionId, cancellationToken); + var workflowGraph = await GetWorkflowGraphAsync(request.WorkflowDefinitionId, cancellationToken); - if (workflow == null) + if (workflowGraph == null) { AddError($"Workflow definition {request.WorkflowDefinitionId} not found"); await SendErrorsAsync(cancellation: cancellationToken); return; } + var workflow = workflowGraph.Workflow; var typeDefinitionContext = new TypeDefinitionContext(workflow, request.ActivityTypeName, request.PropertyName, cancellationToken); var typeDefinitions = await _typeDefinitionService.GenerateTypeDefinitionsAsync(typeDefinitionContext); var fileName = $"elsa.{request.WorkflowDefinitionId}.d.ts"; @@ -52,9 +53,9 @@ internal class Get : ElsaEndpoint await SendBytesAsync(data, fileName, "application/x-typescript", cancellation: cancellationToken); } - private async Task GetWorkflowAsync(string workflowDefinitionId, CancellationToken cancellationToken) + private async Task GetWorkflowGraphAsync(string workflowDefinitionId, CancellationToken cancellationToken) { - return await _workflowDefinitionService.FindWorkflowAsync(workflowDefinitionId, VersionOptions.Latest, cancellationToken); + return await _workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionId, VersionOptions.Latest, cancellationToken); } } diff --git a/src/modules/Elsa.JavaScript/Services/JintJavaScriptEvaluator.cs b/src/modules/Elsa.JavaScript/Services/JintJavaScriptEvaluator.cs index 57b3921ce..90013a0b2 100644 --- a/src/modules/Elsa.JavaScript/Services/JintJavaScriptEvaluator.cs +++ b/src/modules/Elsa.JavaScript/Services/JintJavaScriptEvaluator.cs @@ -1,4 +1,6 @@ using System.Diagnostics.CodeAnalysis; +using System.Security.Cryptography; +using System.Text; using System.Text.Encodings.Web; using System.Text.Json; using System.Text.Json.Serialization; @@ -12,8 +14,6 @@ using Elsa.JavaScript.Helpers; using Elsa.JavaScript.Notifications; using Elsa.JavaScript.Options; using Elsa.Mediator.Contracts; -using Esprima; -using Esprima.Ast; using Humanizer; using Jint; using Jint.Runtime.Interop; @@ -182,7 +182,7 @@ public class JintJavaScriptEvaluator : IJavaScriptEvaluator private object? ExecuteExpressionAndGetResult(Engine engine, string expression) { - var cacheKey = "jint:script:" + expression.GetHashCode(StringComparison.Ordinal); + var cacheKey = "jint:script:" + Hash(expression); var parsedScript = _memoryCache.GetOrCreate(cacheKey, entry => { @@ -214,4 +214,11 @@ public class JintJavaScriptEvaluator : IJavaScriptEvaluator return JsonSerializer.Serialize(value, options); } + + private string Hash(string input) + { + var bytes = Encoding.UTF8.GetBytes(input); + var hash = SHA256.HashData(bytes); + return Convert.ToBase64String(hash); + } } \ No newline at end of file diff --git a/src/modules/Elsa.MassTransit/Services/MassTransitWorkflowDispatcher.cs b/src/modules/Elsa.MassTransit/Services/MassTransitWorkflowDispatcher.cs index 4da57240d..9136ad4d2 100644 --- a/src/modules/Elsa.MassTransit/Services/MassTransitWorkflowDispatcher.cs +++ b/src/modules/Elsa.MassTransit/Services/MassTransitWorkflowDispatcher.cs @@ -32,11 +32,12 @@ public class MassTransitWorkflowDispatcher( /// public async Task DispatchAsync(DispatchWorkflowDefinitionRequest request, DispatchWorkflowOptions? options = default, CancellationToken cancellationToken = default) { - var workflow = await workflowDefinitionService.FindWorkflowAsync(request.DefinitionId, request.VersionOptions, cancellationToken); + var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(request.DefinitionId, request.VersionOptions, cancellationToken); - if (workflow == null) + if (workflowGraph == null) throw new Exception($"Workflow definition with definition ID '{request.DefinitionId} and version {request.VersionOptions}' not found"); + var workflow = workflowGraph.Workflow; var createWorkflowInstanceRequest = new CreateWorkflowInstanceRequest { Workflow = workflow, @@ -106,9 +107,9 @@ public class MassTransitWorkflowDispatcher( foreach (var trigger in triggers) { - var workflow = await workflowDefinitionService.FindWorkflowAsync(trigger.WorkflowDefinitionVersionId, cancellationToken); + var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(trigger.WorkflowDefinitionVersionId, cancellationToken); - if (workflow == null) + if (workflowGraph == null) { logger.LogWarning("Workflow definition with ID '{WorkflowDefinitionId}' not found", trigger.WorkflowDefinitionVersionId); continue; @@ -116,7 +117,7 @@ public class MassTransitWorkflowDispatcher( var createWorkflowInstanceRequest = new CreateWorkflowInstanceRequest { - Workflow = workflow, + Workflow = workflowGraph.Workflow, WorkflowInstanceId = request.WorkflowInstanceId, Input = request.Input, Properties = request.Properties, diff --git a/src/modules/Elsa.ProtoActor/Grains/WorkflowInstance.cs b/src/modules/Elsa.ProtoActor/Grains/WorkflowInstance.cs index be8efa498..9ec169fbc 100644 --- a/src/modules/Elsa.ProtoActor/Grains/WorkflowInstance.cs +++ b/src/modules/Elsa.ProtoActor/Grains/WorkflowInstance.cs @@ -82,12 +82,14 @@ internal class WorkflowInstance : WorkflowInstanceBase using var scope = _scopeFactory.CreateScope(); var workflowDefinitionService = scope.ServiceProvider.GetRequiredService(); - // Load the workflow definition. - var workflow = await workflowDefinitionService.FindWorkflowAsync(_definitionId, VersionOptions.SpecificVersion(_version), cancellationToken); + // Get the workflow. + var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(_definitionId, VersionOptions.SpecificVersion(_version), cancellationToken); - if (workflow == null) + if (workflowGraph == null) throw new Exception("Workflow definition is no longer available"); + var workflow = workflowGraph.Workflow; + // Create an initial workflow state. if (_workflowState == null!) { @@ -100,7 +102,7 @@ internal class WorkflowInstance : WorkflowInstanceBase // Create a workflow host. var workflowHostFactory = scope.ServiceProvider.GetRequiredService(); - _workflowHost = await workflowHostFactory.CreateAsync(workflow, _workflowState, cancellationToken); + _workflowHost = await workflowHostFactory.CreateAsync(workflowGraph, _workflowState, cancellationToken); } /// @@ -416,7 +418,7 @@ internal class WorkflowInstance : WorkflowInstanceBase { using var scope = _scopeFactory.CreateScope(); var workflowDefinitionService = scope.ServiceProvider.GetRequiredService(); - var workflow = await workflowDefinitionService.FindWorkflowAsync(definitionId, versionOptions, cancellationToken); + var workflow = await workflowDefinitionService.FindWorkflowGraphAsync(definitionId, versionOptions, cancellationToken); if (workflow == null) throw new Exception("Specified workflow definition and version does not exist"); @@ -438,7 +440,7 @@ internal class WorkflowInstance : WorkflowInstanceBase var versionOptions = VersionOptions.SpecificVersion(workflowState.DefinitionVersion); using var scope = _scopeFactory.CreateScope(); var workflowDefinitionService = scope.ServiceProvider.GetRequiredService(); - var workflow = await workflowDefinitionService.FindWorkflowAsync(definitionId, versionOptions, cancellationToken); + var workflow = await workflowDefinitionService.FindWorkflowGraphAsync(definitionId, versionOptions, cancellationToken); if (workflow == null) throw new Exception("Specified workflow definition and version does not exist"); diff --git a/src/modules/Elsa.ProtoActor/Services/ProtoActorWorkflowRuntime.cs b/src/modules/Elsa.ProtoActor/Services/ProtoActorWorkflowRuntime.cs index 138790f78..a7a5e716d 100644 --- a/src/modules/Elsa.ProtoActor/Services/ProtoActorWorkflowRuntime.cs +++ b/src/modules/Elsa.ProtoActor/Services/ProtoActorWorkflowRuntime.cs @@ -386,14 +386,15 @@ internal class ProtoActorWorkflowRuntime : IWorkflowRuntime }; var canStartResult = await CanStartWorkflowAsync(definitionId, startOptions); - var workflow = await _workflowDefinitionService.FindWorkflowAsync(trigger.WorkflowDefinitionVersionId, cancellationToken); + var workflowGraph = await _workflowDefinitionService.FindWorkflowGraphAsync(trigger.WorkflowDefinitionVersionId, cancellationToken); - if (workflow == null) + if (workflowGraph == null) { _logger.LogWarning("Workflow version ID {DefinitionVersionId} not found", trigger.WorkflowDefinitionVersionId); continue; } + var workflow = workflowGraph.Workflow; var createWorkflowInstanceRequest = new CreateWorkflowInstanceRequest { Workflow = workflow, diff --git a/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.cs b/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.cs index d33473646..7f5fc4c8f 100644 --- a/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.cs +++ b/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.cs @@ -224,8 +224,7 @@ public partial class ActivityExecutionContext : IExecutionContext /// The options used to schedule the activity. public async ValueTask ScheduleActivityAsync(IActivity? activity, ScheduleWorkOptions? options = default) { - var activityNode = activity != null ? WorkflowExecutionContext.FindNodeByActivity(activity) : default; - await ScheduleActivityAsync(activityNode, this, options); + await ScheduleActivityAsync(activity, this, options); } /// @@ -236,7 +235,9 @@ public partial class ActivityExecutionContext : IExecutionContext /// The options used to schedule the activity. public async ValueTask ScheduleActivityAsync(IActivity? activity, ActivityExecutionContext? owner, ScheduleWorkOptions? options = default) { - var activityNode = activity != null ? WorkflowExecutionContext.FindNodeByActivity(activity) : default; + var activityNode = activity != null + ? WorkflowExecutionContext.FindNodeByActivity(activity) ?? throw new InvalidOperationException("The specified activity is not part of the workflow.") + : null; await ScheduleActivityAsync(activityNode, owner, options); } @@ -250,6 +251,10 @@ public partial class ActivityExecutionContext : IExecutionContext { if (this.GetIsBackgroundExecution()) { + // Validate that the specified activity is part of the workflow. + if (activityNode != null && !WorkflowExecutionContext.NodeActivityLookup.ContainsKey(activityNode.Activity)) + throw new InvalidOperationException("The specified activity is not part of the workflow."); + var scheduledActivity = new ScheduledActivity { ActivityNodeId = activityNode?.NodeId, diff --git a/src/modules/Elsa.Workflows.Core/Contexts/WorkflowExecutionContext.cs b/src/modules/Elsa.Workflows.Core/Contexts/WorkflowExecutionContext.cs index 59014f93e..1af3d2301 100644 --- a/src/modules/Elsa.Workflows.Core/Contexts/WorkflowExecutionContext.cs +++ b/src/modules/Elsa.Workflows.Core/Contexts/WorkflowExecutionContext.cs @@ -44,6 +44,7 @@ public partial class WorkflowExecutionContext : IExecutionContext /// private WorkflowExecutionContext( IServiceProvider serviceProvider, + WorkflowGraph workflowGraph, string id, string? correlationId, string? parentWorkflowInstanceId, @@ -75,7 +76,7 @@ public partial class WorkflowExecutionContext : IExecutionContext UpdatedAt = createdAt; CancellationTokens = cancellationTokens; Incidents = incidents.ToList(); - + WorkflowGraph = workflowGraph; var appSource = CancellationTokenSource.CreateLinkedTokenSource(CancellationTokens.ApplicationCancellationToken); _cancellationTokenSources.Add(appSource); var sysSource = CancellationTokenSource.CreateLinkedTokenSource(CancellationTokens.SystemCancellationToken); @@ -89,22 +90,22 @@ public partial class WorkflowExecutionContext : IExecutionContext /// public static async Task CreateAsync( IServiceProvider serviceProvider, - Workflow workflow, + WorkflowGraph workflowGraph, string id, - string? correlationId, - string? parentWorkflowInstanceId = default, - IDictionary? input = default, - IDictionary? properties = default, - ExecuteActivityDelegate? executeDelegate = default, - string? triggerActivityId = default, - Action? statusUpdatedCallback = default, + string? correlationId = null, + string? parentWorkflowInstanceId = null, + IDictionary? input = null, + IDictionary? properties = null, + ExecuteActivityDelegate? executeDelegate = null, + string? triggerActivityId = null, + Action? statusUpdatedCallback = null, CancellationTokens cancellationTokens = default) { var systemClock = serviceProvider.GetRequiredService(); return await CreateAsync( serviceProvider, - workflow, + workflowGraph, id, new List(), systemClock.UtcNow, @@ -124,20 +125,20 @@ public partial class WorkflowExecutionContext : IExecutionContext /// public static async Task CreateAsync( IServiceProvider serviceProvider, - Workflow workflow, + WorkflowGraph workflowGraph, WorkflowState workflowState, - string? correlationId = default, - string? parentWorkflowInstanceId = default, - IDictionary? input = default, - IDictionary? properties = default, - ExecuteActivityDelegate? executeDelegate = default, - string? triggerActivityId = default, - Action? statusUpdatedCallback = default, + string? correlationId = null, + string? parentWorkflowInstanceId = null, + IDictionary? input = null, + IDictionary? properties = null, + ExecuteActivityDelegate? executeDelegate = null, + string? triggerActivityId = null, + Action? statusUpdatedCallback = null, CancellationTokens cancellationTokens = default) { var workflowExecutionContext = await CreateAsync( serviceProvider, - workflow, + workflowGraph, workflowState.Id, workflowState.Incidents, workflowState.CreatedAt, @@ -161,22 +162,23 @@ public partial class WorkflowExecutionContext : IExecutionContext /// public static async Task CreateAsync( IServiceProvider serviceProvider, - Workflow workflow, + WorkflowGraph workflowGraph, string id, IEnumerable incidents, DateTimeOffset createdAt, - string? correlationId = default, - string? parentWorkflowInstanceId = default, - IDictionary? input = default, - IDictionary? properties = default, - ExecuteActivityDelegate? executeDelegate = default, - string? triggerActivityId = default, - Action? statusUpdatedCallback = default, + string? correlationId = null, + string? parentWorkflowInstanceId = null, + IDictionary? input = null, + IDictionary? properties = null, + ExecuteActivityDelegate? executeDelegate = null, + string? triggerActivityId = null, + Action? statusUpdatedCallback = null, CancellationTokens cancellationTokens = default) { // Setup a workflow execution context. var workflowExecutionContext = new WorkflowExecutionContext( serviceProvider, + workflowGraph, id, correlationId, parentWorkflowInstanceId, @@ -189,47 +191,31 @@ public partial class WorkflowExecutionContext : IExecutionContext statusUpdatedCallback, cancellationTokens) { - MemoryRegister = workflow.CreateRegister() + MemoryRegister = workflowGraph.Workflow.CreateRegister() }; workflowExecutionContext.ExpressionExecutionContext = new ExpressionExecutionContext(serviceProvider, workflowExecutionContext.MemoryRegister, cancellationToken: cancellationTokens.ApplicationCancellationToken); - await workflowExecutionContext.SetWorkflowAsync(workflow); + await workflowExecutionContext.SetWorkflowGraphAsync(workflowGraph); return workflowExecutionContext; } /// /// Assigns the specified workflow to this workflow execution context. /// - /// The workflow to assign. - public async Task SetWorkflowAsync(Workflow workflow) + /// The workflow graph to assign. + public async Task SetWorkflowGraphAsync(WorkflowGraph workflowGraph) { - var activityVisitor = GetRequiredService(); - var root = workflow; - var graph = await activityVisitor.VisitAsync(root, CancellationTokens.ApplicationCancellationToken); - var nodes = graph.Flatten().ToList(); + WorkflowGraph = workflowGraph; + var nodes = workflowGraph.Nodes; // Register activity types. var activityTypes = nodes.Select(x => x.Activity.GetType()).Distinct().ToList(); await ActivityRegistry.RegisterAsync(activityTypes, CancellationTokens.ApplicationCancellationToken); - var needsIdentityAssignment = nodes.Any(x => string.IsNullOrEmpty(x.Activity.Id)); - - if (needsIdentityAssignment) - { - var identityGraphService = GetRequiredService(); - identityGraphService.AssignIdentities(nodes); - } - - Workflow = workflow; - Graph = graph; - Nodes = nodes; - NodeIdLookup = nodes.ToDictionary(x => x.NodeId); - NodeHashLookup = nodes.ToDictionary(x => Hash(x.NodeId)); - NodeActivityLookup = nodes.ToDictionary(x => x.Activity); - - foreach (var activityExecutionContext in ActivityExecutionContexts) - activityExecutionContext.Activity = NodeIdLookup[activityExecutionContext.Activity.NodeId].Activity; + // Update the activity execution contexts with the actual activity instances. + foreach (var activityExecutionContext in ActivityExecutionContexts) + activityExecutionContext.Activity = workflowGraph.NodeIdLookup[activityExecutionContext.Activity.NodeId].Activity; } /// @@ -242,15 +228,20 @@ public partial class WorkflowExecutionContext : IExecutionContext /// public IActivityRegistry ActivityRegistry { get; } + /// + /// Gets the workflow graph. + /// + public WorkflowGraph WorkflowGraph { get; private set; } + /// /// The associated with the execution context. /// - public Workflow Workflow { get; private set; } = default!; + public Workflow Workflow => WorkflowGraph.Workflow; /// /// A graph of the workflow structure. /// - public ActivityNode Graph { get; private set; } = default!; + public ActivityNode Graph => WorkflowGraph.Root; /// /// The current status of the workflow. @@ -279,7 +270,7 @@ public partial class WorkflowExecutionContext : IExecutionContext /// An application-specific identifier associated with the execution context. /// public string? CorrelationId { get; set; } - + /// /// The ID of the workflow instance that triggered this instance. /// @@ -289,12 +280,12 @@ public partial class WorkflowExecutionContext : IExecutionContext /// The date and time the workflow execution context was created. /// public DateTimeOffset CreatedAt { get; set; } - + /// /// The date and time the workflow execution context was last updated. /// public DateTimeOffset UpdatedAt { get; set; } - + /// /// The date and time the workflow execution context has finished. /// @@ -308,22 +299,22 @@ public partial class WorkflowExecutionContext : IExecutionContext /// /// A flattened list of s from the . /// - public IReadOnlyCollection Nodes { get; private set; } = default!; + public IReadOnlyCollection Nodes => WorkflowGraph.Nodes.ToList(); /// /// A map between activity IDs and s in the workflow graph. /// - public IDictionary NodeIdLookup { get; private set; } = default!; + public IDictionary NodeIdLookup => WorkflowGraph.NodeIdLookup; /// /// A map between hashed activity node IDs and s in the workflow graph. /// - public IDictionary NodeHashLookup { get; private set; } = default!; + public IDictionary NodeHashLookup => WorkflowGraph.NodeHashLookup; /// /// A map between s and s in the workflow graph. /// - public IDictionary NodeActivityLookup { get; private set; } = default!; + public IDictionary NodeActivityLookup => WorkflowGraph.NodeActivityLookup; /// /// The for the execution context. @@ -559,15 +550,15 @@ public partial class WorkflowExecutionContext : IExecutionContext SubStatus = subStatus; UpdatedAt = SystemClock.UtcNow; - + if (Status == WorkflowStatus.Finished) FinishedAt = UpdatedAt; - + //For now only trigger on Cancelled, since the other statuses are handling via the host/runner if (SubStatus == WorkflowSubStatus.Cancelled && _statusUpdatedCallback is not null) _statusUpdatedCallback(this); - + if (Status == WorkflowStatus.Finished || SubStatus == WorkflowSubStatus.Suspended) { @@ -594,7 +585,10 @@ public partial class WorkflowExecutionContext : IExecutionContext var id = IdentityGenerator.GenerateId(); var activityExecutionContext = new ActivityExecutionContext(id, this, parentContext, expressionExecutionContext, activity, activityDescriptor, now, tag, SystemClock, CancellationTokens.ApplicationCancellationToken); var variablesToDeclare = options?.Variables ?? Array.Empty(); - var variableContainer = new[] { activityExecutionContext.ActivityNode }.Concat(activityExecutionContext.ActivityNode.Ancestors()).FirstOrDefault(x => x.Activity is IVariableContainer)?.Activity as IVariableContainer; + var variableContainer = new[] + { + activityExecutionContext.ActivityNode + }.Concat(activityExecutionContext.ActivityNode.Ancestors()).FirstOrDefault(x => x.Activity is IVariableContainer)?.Activity as IVariableContainer; expressionExecutionContext.TransientProperties[ExpressionExecutionContextExtensions.ActivityExecutionContextKey] = activityExecutionContext; if (variableContainer != null) @@ -674,7 +668,4 @@ public partial class WorkflowExecutionContext : IExecutionContext var currentMainStatus = GetMainStatus(SubStatus); return currentMainStatus != WorkflowStatus.Finished; } - - - private string Hash(string nodeId) => _hasher.Hash(nodeId); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Contracts/IWorkflowGraphBuilder.cs b/src/modules/Elsa.Workflows.Core/Contracts/IWorkflowGraphBuilder.cs new file mode 100644 index 000000000..97893e1f0 --- /dev/null +++ b/src/modules/Elsa.Workflows.Core/Contracts/IWorkflowGraphBuilder.cs @@ -0,0 +1,15 @@ +using Elsa.Workflows.Activities; +using Elsa.Workflows.Models; + +namespace Elsa.Workflows.Contracts; + +/// +/// Builds a workflow graph from a workflow. +/// +public interface IWorkflowGraphBuilder +{ + /// + /// Builds a workflow graph from a workflow. + /// + Task BuildAsync(Workflow workflow, CancellationToken cancellationToken = default); +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Contracts/IWorkflowRunner.cs b/src/modules/Elsa.Workflows.Core/Contracts/IWorkflowRunner.cs index fb0cc38ab..ad977a4a2 100644 --- a/src/modules/Elsa.Workflows.Core/Contracts/IWorkflowRunner.cs +++ b/src/modules/Elsa.Workflows.Core/Contracts/IWorkflowRunner.cs @@ -17,5 +17,7 @@ public interface IWorkflowRunner Task RunAsync(RunWorkflowOptions? options = default, CancellationToken cancellationToken = default) where T : WorkflowBase, new(); Task RunAsync(Workflow workflow, RunWorkflowOptions? options = default, CancellationToken cancellationToken = default); Task RunAsync(Workflow workflow, WorkflowState workflowState, RunWorkflowOptions? options = default, CancellationToken cancellationToken = default); + Task RunAsync(WorkflowGraph workflowGraph, RunWorkflowOptions? options = default, CancellationToken cancellationToken = default); + Task RunAsync(WorkflowGraph workflowGraph, WorkflowState workflowState, RunWorkflowOptions? options = default, CancellationToken cancellationToken = default); Task RunAsync(WorkflowExecutionContext workflowExecutionContext); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Extensions/WorkflowExecutionContextExtensions.cs b/src/modules/Elsa.Workflows.Core/Extensions/WorkflowExecutionContextExtensions.cs index e3cb8db38..4c159b209 100644 --- a/src/modules/Elsa.Workflows.Core/Extensions/WorkflowExecutionContextExtensions.cs +++ b/src/modules/Elsa.Workflows.Core/Extensions/WorkflowExecutionContextExtensions.cs @@ -99,6 +99,10 @@ public static class WorkflowExecutionContextExtensions ActivityExecutionContext owner, ScheduleWorkOptions? options = default) { + // Validate that the specified activity is part of the workflow. + if (!workflowExecutionContext.NodeActivityLookup.ContainsKey(activityNode.Activity)) + throw new InvalidOperationException("The specified activity is not part of the workflow."); + var scheduler = workflowExecutionContext.Scheduler; if (options?.PreventDuplicateScheduling == true) diff --git a/src/modules/Elsa.Workflows.Core/Features/WorkflowsFeature.cs b/src/modules/Elsa.Workflows.Core/Features/WorkflowsFeature.cs index a4af1a8cd..1135c85e2 100644 --- a/src/modules/Elsa.Workflows.Core/Features/WorkflowsFeature.cs +++ b/src/modules/Elsa.Workflows.Core/Features/WorkflowsFeature.cs @@ -127,6 +127,7 @@ public class WorkflowsFeature : FeatureBase .AddScoped() .AddScoped() .AddScoped() + .AddScoped() .AddScoped() .AddScoped() .AddSingleton() diff --git a/src/modules/Elsa.Workflows.Core/Models/WorkflowGraph.cs b/src/modules/Elsa.Workflows.Core/Models/WorkflowGraph.cs new file mode 100644 index 000000000..f2cfddab3 --- /dev/null +++ b/src/modules/Elsa.Workflows.Core/Models/WorkflowGraph.cs @@ -0,0 +1,62 @@ +using System.Security.Cryptography; +using System.Text; +using Elsa.Workflows.Activities; +using Elsa.Workflows.Contracts; + +namespace Elsa.Workflows.Models; + +/// +/// Hold a reference to a and a collection of instances representing the workflow. +/// +public record WorkflowGraph +{ + /// + /// Initializes a new instance of the class. + /// + public WorkflowGraph(Workflow workflow, ActivityNode root, IEnumerable nodes) + { + using var hashAlgorithm = SHA256.Create(); + Workflow = workflow; + Root = root; + Nodes = nodes.ToList(); + NodeIdLookup = Nodes.ToDictionary(x => x.NodeId); + NodeHashLookup = Nodes.ToDictionary(x => Hash(hashAlgorithm, x.NodeId)); + NodeActivityLookup = Nodes.ToDictionary(x => x.Activity); + } + + /// + /// Gets the workflow. + /// + public Workflow Workflow { get; } + + /// + /// Gets the root node. + /// + public ActivityNode Root { get; } + + /// + /// Gets a flat collection of all nodes in the workflow. + /// + public ICollection Nodes { get; } + + /// + /// Gets a lookup of nodes by their activity. + /// + public IDictionary NodeActivityLookup { get; } + + /// + /// Gets a lookup of nodes by their hash. + /// + public IDictionary NodeHashLookup { get; } + + /// + /// Gets a lookup of nodes by their ID. + /// + public IDictionary NodeIdLookup { get; } + + private static string Hash(HashAlgorithm hashAlgorithm, string input) + { + var data = hashAlgorithm.ComputeHash(Encoding.UTF8.GetBytes(input)); + return Convert.ToHexString(data); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Services/ActivityVisitor.cs b/src/modules/Elsa.Workflows.Core/Services/ActivityVisitor.cs index fd9550383..88b4e65b1 100644 --- a/src/modules/Elsa.Workflows.Core/Services/ActivityVisitor.cs +++ b/src/modules/Elsa.Workflows.Core/Services/ActivityVisitor.cs @@ -10,7 +10,7 @@ public class ActivityVisitor : IActivityVisitor private readonly IEnumerable _portResolvers; /// - /// Constructor. + /// Initializes a new instance of the class. /// public ActivityVisitor(IEnumerable portResolvers, IServiceProvider serviceProvider) { @@ -21,14 +21,26 @@ public class ActivityVisitor : IActivityVisitor /// public async Task VisitAsync(IActivity activity, CancellationToken cancellationToken = default) { - var collectedActivities = new HashSet(new[] { activity }); var graph = new ActivityNode(activity); - var collectedNodes = new HashSet(new[] { graph }); - await VisitRecursiveAsync((graph, activity), collectedActivities, collectedNodes, cancellationToken); + var collectedNodes = new HashSet(new[] + { + graph + }); + var collectedActivities = new HashSet(new[] + { + activity + }); + var visitorContext = new ActivityVisitorContext + { + CollectedActivities = collectedActivities, + CollectedNodes = collectedNodes + }; + + await VisitRecursiveAsync((graph, activity), visitorContext, cancellationToken); return graph; } - private async Task VisitRecursiveAsync((ActivityNode Node, IActivity Activity) pair, HashSet collectedActivities, HashSet collectedNodes, CancellationToken cancellationToken) + private async Task VisitRecursiveAsync((ActivityNode Node, IActivity Activity) pair, ActivityVisitorContext visitorContext, CancellationToken cancellationToken) { if (pair.Activity is IInitializable initializable) { @@ -36,10 +48,10 @@ public class ActivityVisitor : IActivityVisitor await initializable.InitializeAsync(context); } - await VisitPortsRecursiveAsync(pair, collectedActivities, collectedNodes, cancellationToken); + await VisitPortsRecursiveAsync(pair, visitorContext, cancellationToken); } - private async Task VisitPortsRecursiveAsync((ActivityNode Node, IActivity Activity) pair, HashSet collectedActivities, HashSet collectedNodes, CancellationToken cancellationToken) + private async Task VisitPortsRecursiveAsync((ActivityNode Node, IActivity Activity) pair, ActivityVisitorContext visitorContext, CancellationToken cancellationToken) { var resolver = _portResolvers.FirstOrDefault(x => x.GetSupportsActivity(pair.Activity)); @@ -47,6 +59,8 @@ public class ActivityVisitor : IActivityVisitor return; var activities = await resolver.GetActivitiesAsync(pair.Activity, cancellationToken); + var collectedActivities = visitorContext.CollectedActivities; + var collectedNodes = visitorContext.CollectedNodes; foreach (var activity in activities) { @@ -65,7 +79,13 @@ public class ActivityVisitor : IActivityVisitor childNode.Parents.Add(pair.Node); pair.Node.Children.Add(childNode); collectedActivities.Add(activity); - await VisitRecursiveAsync((childNode, activity), collectedActivities, collectedNodes, cancellationToken); + await VisitRecursiveAsync((childNode, activity), visitorContext, cancellationToken); } } + + private class ActivityVisitorContext + { + public HashSet CollectedActivities { get; set; } = new(); + public HashSet CollectedNodes { get; set; } = new(); + } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Services/Hasher.cs b/src/modules/Elsa.Workflows.Core/Services/Hasher.cs index 75ee6aa56..292abeb79 100644 --- a/src/modules/Elsa.Workflows.Core/Services/Hasher.cs +++ b/src/modules/Elsa.Workflows.Core/Services/Hasher.cs @@ -33,8 +33,8 @@ public class Hasher : IHasher /// public string Hash(string value) { - using var sha = SHA256.Create(); - return Hash(sha, value); + var data = SHA256.HashData(Encoding.UTF8.GetBytes(value)); + return Convert.ToHexString(data); } /// @@ -45,12 +45,6 @@ public class Hasher : IHasher var input = string.Join("|", strings); return Hash(input); } - - private static string Hash(HashAlgorithm hashAlgorithm, string input) - { - var data = hashAlgorithm.ComputeHash(Encoding.UTF8.GetBytes(input)); - return Convert.ToHexString(data); - } [RequiresUnreferencedCode("Calls System.Text.Json.JsonSerializer.Serialize(Object, Type, JsonSerializerOptions)")] private string Serialize(object? payload) diff --git a/src/modules/Elsa.Workflows.Core/Services/WorkflowGraphBuilder.cs b/src/modules/Elsa.Workflows.Core/Services/WorkflowGraphBuilder.cs new file mode 100644 index 000000000..e08aa480b --- /dev/null +++ b/src/modules/Elsa.Workflows.Core/Services/WorkflowGraphBuilder.cs @@ -0,0 +1,20 @@ +using Elsa.Extensions; +using Elsa.Workflows.Activities; +using Elsa.Workflows.Contracts; +using Elsa.Workflows.Models; + +namespace Elsa.Workflows.Services; + +/// +public class WorkflowGraphBuilder(IActivityVisitor activityVisitor, IIdentityGraphService identityGraphService, IServiceProvider serviceProvider) : IWorkflowGraphBuilder +{ + /// + public async Task BuildAsync(Workflow workflow, CancellationToken cancellationToken = default) + { + var graph = await activityVisitor.VisitAsync(workflow, cancellationToken); + var nodes = graph.Flatten().ToList(); + + identityGraphService.AssignIdentities(nodes); + return new WorkflowGraph(workflow, graph, nodes); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Core/Services/WorkflowRunner.cs b/src/modules/Elsa.Workflows.Core/Services/WorkflowRunner.cs index 68601ea5d..bf74d0a31 100644 --- a/src/modules/Elsa.Workflows.Core/Services/WorkflowRunner.cs +++ b/src/modules/Elsa.Workflows.Core/Services/WorkflowRunner.cs @@ -10,45 +10,28 @@ using Elsa.Workflows.State; namespace Elsa.Workflows.Services; /// -public class WorkflowRunner : IWorkflowRunner +public class WorkflowRunner( + IServiceProvider serviceProvider, + IWorkflowExecutionPipeline pipeline, + IWorkflowStateExtractor workflowStateExtractor, + IWorkflowBuilderFactory workflowBuilderFactory, + IWorkflowGraphBuilder workflowGraphBuilder, + IIdentityGenerator identityGenerator, + INotificationSender notificationSender) + : IWorkflowRunner { - private readonly IServiceProvider _serviceProvider; - private readonly IWorkflowExecutionPipeline _pipeline; - private readonly IWorkflowStateExtractor _workflowStateExtractor; - private readonly IWorkflowBuilderFactory _workflowBuilderFactory; - private readonly IIdentityGenerator _identityGenerator; - private readonly INotificationSender _notificationSender; - - /// - /// Constructor. - /// - public WorkflowRunner( - IServiceProvider serviceProvider, - IWorkflowExecutionPipeline pipeline, - IWorkflowStateExtractor workflowStateExtractor, - IWorkflowBuilderFactory workflowBuilderFactory, - IIdentityGenerator identityGenerator, - INotificationSender notificationSender) - { - _serviceProvider = serviceProvider; - _pipeline = pipeline; - _workflowStateExtractor = workflowStateExtractor; - _workflowBuilderFactory = workflowBuilderFactory; - _identityGenerator = identityGenerator; - _notificationSender = notificationSender; - } - /// public async Task RunAsync(IActivity activity, RunWorkflowOptions? options = default, CancellationToken cancellationToken = default) { var workflow = Workflow.FromActivity(activity); - return await RunAsync(workflow, options, cancellationToken); + var workflowGraph = await workflowGraphBuilder.BuildAsync(workflow, cancellationToken); + return await RunAsync(workflowGraph, options, cancellationToken); } /// public async Task RunAsync(IWorkflow workflow, RunWorkflowOptions? options = default, CancellationToken cancellationToken = default) { - var builder = _workflowBuilderFactory.CreateBuilder(); + var builder = workflowBuilderFactory.CreateBuilder(); var workflowDefinition = await builder.BuildWorkflowAsync(workflow, cancellationToken); return await RunAsync(workflowDefinition, options, cancellationToken); } @@ -65,7 +48,7 @@ public class WorkflowRunner : IWorkflowRunner RunWorkflowOptions? options = default, CancellationToken cancellationToken = default) where T : IWorkflow, new() { - var builder = _workflowBuilderFactory.CreateBuilder(); + var builder = workflowBuilderFactory.CreateBuilder(); var workflowDefinition = await builder.BuildWorkflowAsync(cancellationToken); return await RunAsync(workflowDefinition, options, cancellationToken); } @@ -75,17 +58,17 @@ public class WorkflowRunner : IWorkflowRunner RunWorkflowOptions? options = default, CancellationToken cancellationToken = default) where T : WorkflowBase, new() { - var builder = _workflowBuilderFactory.CreateBuilder(); + var builder = workflowBuilderFactory.CreateBuilder(); var workflow = await builder.BuildWorkflowAsync(cancellationToken); var result = await RunAsync(workflow, options, cancellationToken); return (TResult)result.Result!; } /// - public async Task RunAsync(Workflow workflow, RunWorkflowOptions? options = default, CancellationToken cancellationToken = default) + public async Task RunAsync(WorkflowGraph workflowGraph, RunWorkflowOptions? options = default, CancellationToken cancellationToken = default) { // Set up a workflow execution context. - var instanceId = options?.WorkflowInstanceId ?? _identityGenerator.GenerateId(); + var instanceId = options?.WorkflowInstanceId ?? identityGenerator.GenerateId(); var input = options?.Input; var properties = options?.Properties; var correlationId = options?.CorrelationId; @@ -93,8 +76,8 @@ public class WorkflowRunner : IWorkflowRunner var parentWorkflowInstanceId = options?.ParentWorkflowInstanceId; var statusUpdatedCallback = options?.StatusUpdatedCallback; var workflowExecutionContext = await WorkflowExecutionContext.CreateAsync( - _serviceProvider, - workflow, + serviceProvider, + workflowGraph, instanceId, correlationId, parentWorkflowInstanceId, @@ -111,8 +94,22 @@ public class WorkflowRunner : IWorkflowRunner return await RunAsync(workflowExecutionContext); } + /// + public async Task RunAsync(Workflow workflow, RunWorkflowOptions? options = default, CancellationToken cancellationToken = default) + { + var workflowGraph = await workflowGraphBuilder.BuildAsync(workflow, cancellationToken); + return await RunAsync(workflowGraph, options, cancellationToken); + } + /// public async Task RunAsync(Workflow workflow, WorkflowState workflowState, RunWorkflowOptions? options = default, CancellationToken cancellationToken = default) + { + var workflowGraph = await workflowGraphBuilder.BuildAsync(workflow, cancellationToken); + return await RunAsync(workflowGraph, workflowState, options, cancellationToken); + } + + /// + public async Task RunAsync(WorkflowGraph workflowGraph, WorkflowState workflowState, RunWorkflowOptions? options = default, CancellationToken cancellationToken = default) { // Create workflow execution context. var input = options?.Input; @@ -122,8 +119,8 @@ public class WorkflowRunner : IWorkflowRunner var parentWorkflowInstanceId = options?.ParentWorkflowInstanceId; var statusUpdatedCallback = options?.StatusUpdatedCallback; var workflowExecutionContext = await WorkflowExecutionContext.CreateAsync( - _serviceProvider, - workflow, + serviceProvider, + workflowGraph, workflowState, correlationId, parentWorkflowInstanceId, @@ -187,25 +184,25 @@ public class WorkflowRunner : IWorkflowRunner var applicationCancellationToken = workflowExecutionContext.CancellationTokens.ApplicationCancellationToken; var systemCancellationToken = workflowExecutionContext.CancellationTokens.SystemCancellationToken; - await _notificationSender.SendAsync(new WorkflowExecuting(workflow, workflowExecutionContext), applicationCancellationToken); + await notificationSender.SendAsync(new WorkflowExecuting(workflow, workflowExecutionContext), applicationCancellationToken); // If the status is Pending, it means the workflow is started for the first time. if (workflowExecutionContext.SubStatus == WorkflowSubStatus.Pending) { workflowExecutionContext.TransitionTo(WorkflowSubStatus.Executing); - await _notificationSender.SendAsync(new WorkflowStarted(workflow, workflowExecutionContext), applicationCancellationToken); + await notificationSender.SendAsync(new WorkflowStarted(workflow, workflowExecutionContext), applicationCancellationToken); } - await _pipeline.ExecuteAsync(workflowExecutionContext); - var workflowState = _workflowStateExtractor.Extract(workflowExecutionContext); + await pipeline.ExecuteAsync(workflowExecutionContext); + var workflowState = workflowStateExtractor.Extract(workflowExecutionContext); if (workflowState.Status == WorkflowStatus.Finished) { - await _notificationSender.SendAsync(new WorkflowFinished(workflow, workflowState, workflowExecutionContext), applicationCancellationToken); + await notificationSender.SendAsync(new WorkflowFinished(workflow, workflowState, workflowExecutionContext), applicationCancellationToken); } var result = workflow.ResultVariable?.Get(workflowExecutionContext.MemoryRegister); - await _notificationSender.SendAsync(new WorkflowExecuted(workflow, workflowState, workflowExecutionContext), systemCancellationToken); + await notificationSender.SendAsync(new WorkflowExecuted(workflow, workflowState, workflowExecutionContext), systemCancellationToken); return new RunWorkflowResult(workflowState, workflowExecutionContext.Workflow, result); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs b/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs index 79cf91274..e3572a7e5 100644 --- a/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs +++ b/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivity.cs @@ -20,6 +20,8 @@ namespace Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity; [Browsable(false)] public class WorkflowDefinitionActivity : Composite, IInitializable { + private bool IsInitialized => Root.Id != null; + /// /// The definition ID of the workflow to schedule for execution. /// @@ -74,7 +76,7 @@ public class WorkflowDefinitionActivity : Composite, IInitializable .Any(x => x.ActivityDescriptor.Outputs.Any(y => y.Name == outputDescriptor.Name)) == true ? activityExecutionContext.ParentActivityExecutionContext ?? activityExecutionContext : activityExecutionContext; - + parentActivityExecutionContext.Set(output, value, outputDescriptor.Name); } @@ -134,24 +136,35 @@ public class WorkflowDefinitionActivity : Composite, IInitializable } } - private async Task FindWorkflowAsync(IServiceProvider serviceProvider, CancellationToken cancellationToken) + private async Task FindWorkflowGraphAsync(IServiceProvider serviceProvider, CancellationToken cancellationToken) { var workflowDefinitionService = serviceProvider.GetRequiredService(); - var filter = new WorkflowDefinitionFilter { DefinitionId = WorkflowDefinitionId }; + var filter = new WorkflowDefinitionFilter + { + DefinitionId = WorkflowDefinitionId + }; if (!string.IsNullOrWhiteSpace(WorkflowDefinitionVersionId)) filter.Id = WorkflowDefinitionVersionId; else filter.VersionOptions = VersionOptions.SpecificVersion(Version); - var workflow = - await workflowDefinitionService.FindWorkflowAsync(filter, cancellationToken) - ?? (await workflowDefinitionService.FindWorkflowAsync(new WorkflowDefinitionFilter { DefinitionId = WorkflowDefinitionId, VersionOptions = VersionOptions.Published }, cancellationToken) - ?? await workflowDefinitionService.FindWorkflowAsync(new WorkflowDefinitionFilter { DefinitionId = WorkflowDefinitionId, VersionOptions = VersionOptions.Latest }, cancellationToken)); + var workflowGraph = + await workflowDefinitionService.FindWorkflowGraphAsync(filter, cancellationToken) + ?? (await workflowDefinitionService.FindWorkflowGraphAsync(new WorkflowDefinitionFilter + { + DefinitionId = WorkflowDefinitionId, + VersionOptions = VersionOptions.Published + }, cancellationToken) + ?? await workflowDefinitionService.FindWorkflowGraphAsync(new WorkflowDefinitionFilter + { + DefinitionId = WorkflowDefinitionId, + VersionOptions = VersionOptions.Latest + }, cancellationToken)); - return workflow; + return workflowGraph; } - + private ActivityDescriptor FindActivityDescriptor(IServiceProvider serviceProvider) { var activityRegistry = serviceProvider.GetRequiredService(); @@ -160,20 +173,26 @@ public class WorkflowDefinitionActivity : Composite, IInitializable async ValueTask IInitializable.InitializeAsync(InitializationContext context) { + // This is not just for efficiency, but also a necessity to avoid potential race conditions. + // Such conditions can occur when multiple threads are simultaneously creating consuming workflows, + // especially when cached workflows are being updated during the graph construction process. + if (IsInitialized) + return; + var serviceProvider = context.ServiceProvider; var cancellationToken = context.CancellationToken; - var workflow = await FindWorkflowAsync(serviceProvider, cancellationToken); + var workflowGraph = await FindWorkflowGraphAsync(serviceProvider, cancellationToken); - if (workflow == null) + if (workflowGraph == null) throw new Exception($"Could not find workflow definition with ID {WorkflowDefinitionId}."); var activityDescriptor = FindActivityDescriptor(serviceProvider); - + // Declare input and output variables. DeclareInputAsVariables(activityDescriptor, (_, variable) => Variables.Declare(variable)); DeclareOutputAsVariables(activityDescriptor, (_, variable) => Variables.Declare(variable)); // Set the root activity. - Root = workflow; + Root = workflowGraph.Workflow; } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivityPortResolver.cs b/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivityPortResolver.cs index 9cc848bc3..78062aed1 100644 --- a/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivityPortResolver.cs +++ b/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivityPortResolver.cs @@ -19,6 +19,6 @@ public class WorkflowDefinitionActivityResolver : IActivityResolver var definitionActivity = (WorkflowDefinitionActivity)activity; var root = definitionActivity.Root; - return new(new[] { root }); + return new([root]); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionService.cs b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionService.cs index da5b416af..c04f4d983 100644 --- a/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionService.cs +++ b/src/modules/Elsa.Workflows.Management/Contracts/IWorkflowDefinitionService.cs @@ -2,6 +2,7 @@ using Elsa.Common.Models; using Elsa.Workflows.Activities; using Elsa.Workflows.Management.Entities; using Elsa.Workflows.Management.Filters; +using Elsa.Workflows.Models; namespace Elsa.Workflows.Management.Contracts; @@ -13,7 +14,7 @@ public interface IWorkflowDefinitionService /// /// Constructs an executable from the specified . /// - Task MaterializeWorkflowAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default); + Task MaterializeWorkflowAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default); /// /// Looks for a by the specified definition ID and . @@ -33,15 +34,15 @@ public interface IWorkflowDefinitionService /// /// Looks for a by the specified definition ID and . /// - Task FindWorkflowAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default); + Task FindWorkflowGraphAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default); /// /// Looks for a by the specified version ID. /// - Task FindWorkflowAsync(string definitionVersionId, CancellationToken cancellationToken = default); + Task FindWorkflowGraphAsync(string definitionVersionId, CancellationToken cancellationToken = default); /// /// Looks for a by the specified . /// - Task FindWorkflowAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default); + Task FindWorkflowGraphAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Mappers/WorkflowDefinitionMapper.cs b/src/modules/Elsa.Workflows.Management/Mappers/WorkflowDefinitionMapper.cs index fbdf5777f..846e6b3d8 100644 --- a/src/modules/Elsa.Workflows.Management/Mappers/WorkflowDefinitionMapper.cs +++ b/src/modules/Elsa.Workflows.Management/Mappers/WorkflowDefinitionMapper.cs @@ -99,7 +99,8 @@ public class WorkflowDefinitionMapper /// The mapped . public async Task MapAsync(WorkflowDefinition workflowDefinition, CancellationToken cancellationToken = default) { - var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(workflowDefinition, cancellationToken); + var workflowGraph = await _workflowDefinitionService.MaterializeWorkflowAsync(workflowDefinition, cancellationToken); + var workflow = workflowGraph.Workflow; var variables = _variableDefinitionMapper.Map(workflow.Variables).ToList(); return new( @@ -120,7 +121,7 @@ public class WorkflowDefinitionMapper workflowDefinition.IsLatest, workflowDefinition.IsPublished, workflow.Options, - default, + null, workflow.Root); } @@ -151,7 +152,7 @@ public class WorkflowDefinitionMapper workflow.Publication.IsLatest, workflow.Publication.IsPublished, workflow.Options, - default, + null, workflow.Root); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Services/CachingWorkflowDefinitionService.cs b/src/modules/Elsa.Workflows.Management/Services/CachingWorkflowDefinitionService.cs index 9022683be..869bce4eb 100644 --- a/src/modules/Elsa.Workflows.Management/Services/CachingWorkflowDefinitionService.cs +++ b/src/modules/Elsa.Workflows.Management/Services/CachingWorkflowDefinitionService.cs @@ -1,8 +1,8 @@ using Elsa.Common.Models; -using Elsa.Workflows.Activities; using Elsa.Workflows.Management.Contracts; using Elsa.Workflows.Management.Entities; using Elsa.Workflows.Management.Filters; +using Elsa.Workflows.Models; using JetBrains.Annotations; using Microsoft.Extensions.Caching.Memory; @@ -15,7 +15,7 @@ namespace Elsa.Workflows.Management.Services; public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decoratedService, IWorkflowDefinitionCacheManager cacheManager) : IWorkflowDefinitionService { /// - public async Task MaterializeWorkflowAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default) + public async Task MaterializeWorkflowAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default) { return await decoratedService.MaterializeWorkflowAsync(definition, cancellationToken); } @@ -49,30 +49,30 @@ public class CachingWorkflowDefinitionService(IWorkflowDefinitionService decorat } /// - public async Task FindWorkflowAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default) + public async Task FindWorkflowGraphAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default) { var cacheKey = cacheManager.CreateWorkflowVersionCacheKey(definitionId, versionOptions); return await GetFromCacheAsync(cacheKey, - () => decoratedService.FindWorkflowAsync(definitionId, versionOptions, cancellationToken), - x => x.Identity.DefinitionId); + () => decoratedService.FindWorkflowGraphAsync(definitionId, versionOptions, cancellationToken), + x => x.Workflow.Identity.DefinitionId); } /// - public async Task FindWorkflowAsync(string definitionVersionId, CancellationToken cancellationToken = default) + public async Task FindWorkflowGraphAsync(string definitionVersionId, CancellationToken cancellationToken = default) { var cacheKey = cacheManager.CreateWorkflowVersionCacheKey(definitionVersionId); return await GetFromCacheAsync(cacheKey, - () => decoratedService.FindWorkflowAsync(definitionVersionId, cancellationToken), - x => x.Identity.DefinitionId); + () => decoratedService.FindWorkflowGraphAsync(definitionVersionId, cancellationToken), + x => x.Workflow.Identity.DefinitionId); } /// - public async Task FindWorkflowAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) + public async Task FindWorkflowGraphAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) { var cacheKey = cacheManager.CreateWorkflowFilterCacheKey(filter); return await GetFromCacheAsync(cacheKey, - () => decoratedService.FindWorkflowAsync(filter, cancellationToken), - x => x.Identity.DefinitionId); + () => decoratedService.FindWorkflowGraphAsync(filter, cancellationToken), + x => x.Workflow.Identity.DefinitionId); } private async Task GetFromCacheAsync(string cacheKey, Func> getObjectFunc, Func getChangeTokenKeyFunc) diff --git a/src/modules/Elsa.Workflows.Management/Services/ExpressionDescriptorRegistryPopulator.cs b/src/modules/Elsa.Workflows.Management/Services/ExpressionDescriptorRegistryPopulator.cs deleted file mode 100644 index 6957ef42f..000000000 --- a/src/modules/Elsa.Workflows.Management/Services/ExpressionDescriptorRegistryPopulator.cs +++ /dev/null @@ -1,29 +0,0 @@ -// using Elsa.Expressions.Contracts; -// -// namespace Elsa.Workflows.Management.Services; -// -// /// -// public class ExpressionDescriptorRegistryPopulator : IExpressionDescriptorRegistryPopulator -// { -// private readonly IEnumerable _providers; -// private readonly IExpressionDescriptorRegistry _registry; -// -// /// -// /// Initializes a new instance of the class. -// /// -// public ExpressionDescriptorRegistryPopulator(IEnumerable providers, IExpressionDescriptorRegistry registry) -// { -// _providers = providers; -// _registry = registry; -// } -// -// /// -// public async ValueTask PopulateRegistryAsync(CancellationToken cancellationToken) -// { -// foreach (var provider in _providers) -// { -// var descriptors = await provider.GetDescriptorsAsync(cancellationToken); -// _registry.AddRange(descriptors); -// } -// } -// } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionPublisher.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionPublisher.cs index dd37fe833..bdba1ab26 100644 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionPublisher.cs +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionPublisher.cs @@ -85,8 +85,8 @@ public class WorkflowDefinitionPublisher : IWorkflowDefinitionPublisher /// public async Task PublishAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default) { - var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(definition, cancellationToken); - var responses = await _requestSender.SendAsync(new ValidateWorkflowRequest(workflow), cancellationToken); + var workflowGraph = await _workflowDefinitionService.MaterializeWorkflowAsync(definition, cancellationToken); + var responses = await _requestSender.SendAsync(new ValidateWorkflowRequest(workflowGraph.Workflow), cancellationToken); var validationErrors = responses.SelectMany(r => r.ValidationErrors).ToList(); if (validationErrors.Any()) diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionService.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionService.cs index 3020faf67..578e537a7 100644 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionService.cs +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionService.cs @@ -1,75 +1,61 @@ using Elsa.Common.Models; -using Elsa.Extensions; -using Elsa.Workflows.Activities; using Elsa.Workflows.Contracts; using Elsa.Workflows.Management.Contracts; using Elsa.Workflows.Management.Entities; using Elsa.Workflows.Management.Filters; +using Elsa.Workflows.Models; namespace Elsa.Workflows.Management.Services; /// -public class WorkflowDefinitionService : IWorkflowDefinitionService +public class WorkflowDefinitionService( + IWorkflowDefinitionStore workflowDefinitionStore, + IWorkflowGraphBuilder workflowGraphBuilder, + Func> materializers) + : IWorkflowDefinitionService { - private readonly IWorkflowDefinitionStore _workflowDefinitionStore; - private readonly IActivityVisitor _activityVisitor; - private readonly IIdentityGraphService _identityGraphService; - private readonly Func> _materializers; - - /// - /// Constructor. - /// - public WorkflowDefinitionService( - IWorkflowDefinitionStore workflowDefinitionStore, - IActivityVisitor activityVisitor, - IIdentityGraphService identityGraphService, - Func> materializers) - { - _workflowDefinitionStore = workflowDefinitionStore; - _activityVisitor = activityVisitor; - _identityGraphService = identityGraphService; - _materializers = materializers; - } - /// - public async Task MaterializeWorkflowAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default) + public async Task MaterializeWorkflowAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default) { - var materializers = _materializers(); - var materializer = materializers.FirstOrDefault(x => x.Name == definition.MaterializerName); + var workflowMaterializers = materializers(); + var materializer = workflowMaterializers.FirstOrDefault(x => x.Name == definition.MaterializerName); if (materializer == null) throw new Exception("Provider not found"); var workflow = await materializer.MaterializeAsync(definition, cancellationToken); - var graph = (await _activityVisitor.VisitAsync(workflow, cancellationToken)).Flatten().ToList(); - - _identityGraphService.AssignIdentities(graph); - - return workflow; + return await workflowGraphBuilder.BuildAsync(workflow, cancellationToken); } /// public async Task FindWorkflowDefinitionAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default) { - var filter = new WorkflowDefinitionFilter { DefinitionId = definitionId, VersionOptions = versionOptions }; - return await _workflowDefinitionStore.FindAsync(filter, cancellationToken); + var filter = new WorkflowDefinitionFilter + { + DefinitionId = definitionId, + VersionOptions = versionOptions + }; + return await workflowDefinitionStore.FindAsync(filter, cancellationToken); } /// public async Task FindWorkflowDefinitionAsync(string definitionVersionId, CancellationToken cancellationToken = default) { - var filter = new WorkflowDefinitionFilter { Id = definitionVersionId }; - return await _workflowDefinitionStore.FindAsync(filter, cancellationToken); + var filter = new WorkflowDefinitionFilter + { + Id = definitionVersionId + }; + return await workflowDefinitionStore.FindAsync(filter, cancellationToken); } /// public async Task FindWorkflowDefinitionAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) { - return await _workflowDefinitionStore.FindAsync(filter, cancellationToken); + return await workflowDefinitionStore.FindAsync(filter, cancellationToken); } /// - public async Task FindWorkflowAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default) + public async Task FindWorkflowGraphAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default) { var definition = await FindWorkflowDefinitionAsync(definitionId, versionOptions, cancellationToken); @@ -80,7 +66,7 @@ public class WorkflowDefinitionService : IWorkflowDefinitionService } /// - public async Task FindWorkflowAsync(string definitionVersionId, CancellationToken cancellationToken = default) + public async Task FindWorkflowGraphAsync(string definitionVersionId, CancellationToken cancellationToken = default) { var definition = await FindWorkflowDefinitionAsync(definitionVersionId, cancellationToken); @@ -91,7 +77,7 @@ public class WorkflowDefinitionService : IWorkflowDefinitionService } /// - public async Task FindWorkflowAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) + public async Task FindWorkflowGraphAsync(WorkflowDefinitionFilter filter, CancellationToken cancellationToken = default) { var definition = await FindWorkflowDefinitionAsync(filter, cancellationToken); diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowInstanceManager.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowInstanceManager.cs index 9a98f3ddb..1312a18c0 100644 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowInstanceManager.cs +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowInstanceManager.cs @@ -8,7 +8,6 @@ using Elsa.Workflows.Management.Mappers; using Elsa.Workflows.Management.Notifications; using Elsa.Workflows.Management.Requests; using Elsa.Workflows.State; -using Exception = System.Exception; namespace Elsa.Workflows.Management.Services; diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/ITriggerIndexer.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/ITriggerIndexer.cs index 6a7f76c91..61434250f 100644 --- a/src/modules/Elsa.Workflows.Runtime/Contracts/ITriggerIndexer.cs +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/ITriggerIndexer.cs @@ -11,11 +11,6 @@ namespace Elsa.Workflows.Runtime.Contracts; /// public interface ITriggerIndexer { - /// - /// Removes triggers for the specified workflow. - /// - Task DeleteTriggersAsync(Workflow workflow, CancellationToken cancellationToken = default); - /// /// Removes triggers matching the specified filter. /// @@ -30,14 +25,6 @@ public interface ITriggerIndexer /// Indexes triggers of the specified workflow. /// Task IndexTriggersAsync(Workflow workflow, CancellationToken cancellationToken = default); - - /// - /// Returns triggers for the specified workflow definition. - /// - /// The workflow definition. - /// An optional cancellation token. - /// A collection of triggers. - Task> GetTriggersAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default); /// /// Returns triggers for the specified workflow definition. diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowHost.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowHost.cs index 167372f53..c372c7b44 100644 --- a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowHost.cs +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowHost.cs @@ -1,4 +1,5 @@ using Elsa.Workflows.Activities; +using Elsa.Workflows.Models; using Elsa.Workflows.Runtime.Parameters; using Elsa.Workflows.Runtime.Results; using Elsa.Workflows.State; @@ -11,9 +12,14 @@ namespace Elsa.Workflows.Runtime.Contracts; public interface IWorkflowHost { /// - /// The workflow definition. + /// The workflow graph. /// - Workflow Workflow { get; set; } + WorkflowGraph WorkflowGraph { get; } + + /// + /// The workflow. + /// + Workflow Workflow { get; } /// /// The workflow state. diff --git a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowHostFactory.cs b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowHostFactory.cs index 8b82315f1..1df02dd5f 100644 --- a/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowHostFactory.cs +++ b/src/modules/Elsa.Workflows.Runtime/Contracts/IWorkflowHostFactory.cs @@ -1,6 +1,5 @@ using Elsa.Common.Models; -using Elsa.Workflows.Activities; -using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Models; using Elsa.Workflows.State; namespace Elsa.Workflows.Runtime.Contracts; @@ -21,15 +20,15 @@ public interface IWorkflowHostFactory /// /// Creates a new object. /// - /// The workflow. + /// The workflow. /// The workflow state to initialize the workflow host with. /// An optional cancellation token. - Task CreateAsync(Workflow workflow, WorkflowState workflowState, CancellationToken cancellationToken = default); + Task CreateAsync(WorkflowGraph workflowGraph, WorkflowState workflowState, CancellationToken cancellationToken = default); /// /// Creates a new object. /// - /// The workflow. + /// The workflow. /// An optional cancellation token. - Task CreateAsync(Workflow workflow, CancellationToken cancellationToken = default); + Task CreateAsync(WorkflowGraph workflowGraph, CancellationToken cancellationToken = default); } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Services/DefaultBackgroundActivityInvoker.cs b/src/modules/Elsa.Workflows.Runtime/Services/DefaultBackgroundActivityInvoker.cs index 6ad32f9da..bd3da688d 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/DefaultBackgroundActivityInvoker.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/DefaultBackgroundActivityInvoker.cs @@ -58,12 +58,12 @@ public class DefaultBackgroundActivityInvoker : IBackgroundActivityInvoker if (workflowState == null) throw new Exception("Workflow state not found"); - var workflow = await _workflowDefinitionService.FindWorkflowAsync(workflowState.DefinitionId, VersionOptions.SpecificVersion(workflowState.DefinitionVersion), cancellationToken); + var workflowGraph = await _workflowDefinitionService.FindWorkflowGraphAsync(workflowState.DefinitionId, VersionOptions.SpecificVersion(workflowState.DefinitionVersion), cancellationToken); - if (workflow == null) + if (workflowGraph == null) throw new Exception("Workflow definition not found"); - var workflowExecutionContext = await WorkflowExecutionContext.CreateAsync(_serviceProvider, workflow, workflowState, cancellationTokens: cancellationToken); + var workflowExecutionContext = await WorkflowExecutionContext.CreateAsync(_serviceProvider, workflowGraph, workflowState, cancellationTokens: cancellationToken); var activityNodeId = scheduledBackgroundActivity.ActivityNodeId; var activityExecutionContext = workflowExecutionContext.ActivityExecutionContexts.First(x => x.NodeId == activityNodeId); diff --git a/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs b/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs index ecd15938e..03b89dc28 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowDefinitionStorePopulator.cs @@ -242,7 +242,7 @@ public class DefaultWorkflowDefinitionStorePopulator : IWorkflowDefinitionStoreP } } - private async Task IndexTriggersAsync(MaterializedWorkflow workflow, CancellationToken cancellationToken) => await _triggerIndexer.IndexTriggersAsync(workflow.Workflow, cancellationToken); + private async Task IndexTriggersAsync(MaterializedWorkflow materializedWorkflow, CancellationToken cancellationToken) => await _triggerIndexer.IndexTriggersAsync(materializedWorkflow.Workflow, cancellationToken); /// /// Syncs the items in the primary list with existing items in the secondary list, even when the object instances are not the same (but their IDs are). diff --git a/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowRuntime.cs b/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowRuntime.cs index 0d24944ac..05c5f4143 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowRuntime.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/DefaultWorkflowRuntime.cs @@ -140,7 +140,7 @@ public class DefaultWorkflowRuntime( var definitionId = workflowInstance.DefinitionId; var version = workflowInstance.Version; - var workflow = await workflowDefinitionService.FindWorkflowAsync( + var workflow = await workflowDefinitionService.FindWorkflowGraphAsync( definitionId, VersionOptions.SpecificVersion(version), systemCancellationToken); @@ -267,7 +267,7 @@ public class DefaultWorkflowRuntime( if (workflowState == null) throw new Exception("Workflow state not found"); - var workflow = await workflowDefinitionService.FindWorkflowAsync(workflowState.DefinitionId, VersionOptions.SpecificVersion(workflowState.DefinitionVersion), cancellationToken); + var workflow = await workflowDefinitionService.FindWorkflowGraphAsync(workflowState.DefinitionId, VersionOptions.SpecificVersion(workflowState.DefinitionVersion), cancellationToken); if (workflow == null) throw new Exception("Workflow definition not found"); @@ -385,7 +385,7 @@ public class DefaultWorkflowRuntime( if (options?.IsExistingInstance == true) { var workflowState = await LoadWorkflowStateAsync(options.InstanceId!, cancellationToken); - var workflow = await workflowDefinitionService.FindWorkflowAsync(workflowState.DefinitionVersionId, cancellationToken) ?? throw new Exception("Specified workflow definition and version does not exist"); + var workflow = await workflowDefinitionService.FindWorkflowGraphAsync(workflowState.DefinitionVersionId, cancellationToken) ?? throw new Exception("Specified workflow definition and version does not exist"); return await workflowHostFactory.CreateAsync(workflow, workflowState, cancellationToken); } @@ -465,14 +465,14 @@ public class DefaultWorkflowRuntime( InstanceId = workflowsFilter.Options.WorkflowInstanceId }; var canStartResult = await CanStartWorkflowAsync(definitionId, startOptions); - var workflow = await workflowDefinitionService.FindWorkflowAsync(definitionId, startOptions.VersionOptions, cancellationToken); + var workflowGraph = await workflowDefinitionService.FindWorkflowGraphAsync(definitionId, startOptions.VersionOptions, cancellationToken); - if (workflow == null) + if (workflowGraph == null) throw new Exception($"The workflow definition {definitionId} was not found"); var createWorkflowInstanceRequest = new CreateWorkflowInstanceRequest { - Workflow = workflow, + Workflow = workflowGraph.Workflow, CorrelationId = workflowsFilter.Options.CorrelationId, Properties = workflowsFilter.Options.Properties, Input = workflowsFilter.Options.Input, diff --git a/src/modules/Elsa.Workflows.Runtime/Services/TriggerIndexer.cs b/src/modules/Elsa.Workflows.Runtime/Services/TriggerIndexer.cs index ed07f9ac8..287d8cb03 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/TriggerIndexer.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/TriggerIndexer.cs @@ -63,21 +63,20 @@ public class TriggerIndexer : ITriggerIndexer foreach (string workflowDefinitionVersionId in workflowDefinitionVersionIds) { - var workflowDefinition = await _workflowDefinitionService.FindWorkflowDefinitionAsync(workflowDefinitionVersionId, cancellationToken); + var workflowGraph = await _workflowDefinitionService.FindWorkflowGraphAsync(workflowDefinitionVersionId, cancellationToken); - if (workflowDefinition == null) + if (workflowGraph == null) continue; - - var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(workflowDefinition, cancellationToken); - await DeleteTriggersAsync(workflow, cancellationToken); + + await DeleteTriggersAsync(workflowGraph.Workflow, cancellationToken); } } /// public async Task IndexTriggersAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default) { - var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(definition, cancellationToken); - return await IndexTriggersAsync(workflow, cancellationToken); + var workflowGraph = await _workflowDefinitionService.MaterializeWorkflowAsync(definition, cancellationToken); + return await IndexTriggersAsync(workflowGraph.Workflow, cancellationToken); } /// @@ -104,21 +103,13 @@ public class TriggerIndexer : ITriggerIndexer return indexedWorkflow; } - /// - public async Task> GetTriggersAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default) - { - var workflow = await _workflowDefinitionService.MaterializeWorkflowAsync(definition, cancellationToken); - return await GetTriggersAsync(workflow, cancellationToken); - } - /// public async Task> GetTriggersAsync(Workflow workflow, CancellationToken cancellationToken) { return await GetTriggersInternalAsync(workflow, cancellationToken).ToListAsync(cancellationToken); } - - /// - public async Task DeleteTriggersAsync(Workflow workflow, CancellationToken cancellationToken = default) + + private async Task DeleteTriggersAsync(Workflow workflow, CancellationToken cancellationToken = default) { var emptyTriggerList = new List(0); var currentTriggers = await GetCurrentTriggersAsync(workflow.Identity.DefinitionId, cancellationToken).ToList(); @@ -126,7 +117,6 @@ public class TriggerIndexer : ITriggerIndexer await _triggerStore.ReplaceAsync(diff.Removed, diff.Added, cancellationToken); var indexedWorkflow = new IndexedWorkflowTriggers(workflow, emptyTriggerList, currentTriggers, emptyTriggerList); await _notificationSender.SendAsync(new WorkflowTriggersIndexed(indexedWorkflow), cancellationToken); - return indexedWorkflow; } private async Task> GetCurrentTriggersAsync(string workflowDefinitionId, CancellationToken cancellationToken) diff --git a/src/modules/Elsa.Workflows.Runtime/Services/WorkflowHost.cs b/src/modules/Elsa.Workflows.Runtime/Services/WorkflowHost.cs index 7a26a1d9f..6e7035481 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/WorkflowHost.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/WorkflowHost.cs @@ -27,12 +27,12 @@ public class WorkflowHost : IWorkflowHost /// public WorkflowHost( IServiceScopeFactory serviceScopeFactory, - Workflow workflow, + WorkflowGraph workflowGraph, WorkflowState workflowState, IIdentityGenerator identityGenerator, ILogger logger) { - Workflow = workflow; + WorkflowGraph = workflowGraph; WorkflowState = workflowState; _serviceScopeFactory = serviceScopeFactory; _identityGenerator = identityGenerator; @@ -40,7 +40,10 @@ public class WorkflowHost : IWorkflowHost } /// - public Workflow Workflow { get; set; } + public WorkflowGraph WorkflowGraph { get; } + + /// + public Workflow Workflow => WorkflowGraph.Workflow; /// public WorkflowState WorkflowState { get; set; } @@ -90,8 +93,8 @@ public class WorkflowHost : IWorkflowHost using var scope = _serviceScopeFactory.CreateScope(); var workflowRunner = scope.ServiceProvider.GetRequiredService(); var workflowResult = @params?.IsExistingInstance == true - ? await workflowRunner.RunAsync(Workflow, WorkflowState, runOptions, cancellationToken) - : await workflowRunner.RunAsync(Workflow, runOptions, cancellationToken); + ? await workflowRunner.RunAsync(WorkflowGraph, WorkflowState, runOptions, cancellationToken) + : await workflowRunner.RunAsync(WorkflowGraph, runOptions, cancellationToken); WorkflowState = workflowResult.WorkflowState; @@ -133,7 +136,7 @@ public class WorkflowHost : IWorkflowHost using var scope = _serviceScopeFactory.CreateScope(); var workflowRunner = scope.ServiceProvider.GetRequiredService(); - var workflowResult = await workflowRunner.RunAsync(Workflow, WorkflowState, runOptions, cancellationToken); + var workflowResult = await workflowRunner.RunAsync(WorkflowGraph, WorkflowState, runOptions, cancellationToken); WorkflowState = workflowResult.WorkflowState; diff --git a/src/modules/Elsa.Workflows.Runtime/Services/WorkflowHostFactory.cs b/src/modules/Elsa.Workflows.Runtime/Services/WorkflowHostFactory.cs index 8bcd4d320..c2c89bb4d 100644 --- a/src/modules/Elsa.Workflows.Runtime/Services/WorkflowHostFactory.cs +++ b/src/modules/Elsa.Workflows.Runtime/Services/WorkflowHostFactory.cs @@ -1,8 +1,7 @@ using Elsa.Common.Models; -using Elsa.Workflows.Activities; using Elsa.Workflows.Contracts; using Elsa.Workflows.Management.Contracts; -using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Models; using Elsa.Workflows.Runtime.Contracts; using Elsa.Workflows.State; using Microsoft.Extensions.DependencyInjection; @@ -31,7 +30,7 @@ public class WorkflowHostFactory : IWorkflowHostFactory { using var scope = _serviceScopeFactory.CreateScope(); var workflowDefinitionService = scope.ServiceProvider.GetRequiredService(); - var workflow = await workflowDefinitionService.FindWorkflowAsync(definitionId, versionOptions, cancellationToken); + var workflow = await workflowDefinitionService.FindWorkflowGraphAsync(definitionId, versionOptions, cancellationToken); if(workflow == null) return default; @@ -40,22 +39,23 @@ public class WorkflowHostFactory : IWorkflowHostFactory } /// - public Task CreateAsync(Workflow workflow, WorkflowState workflowState, CancellationToken cancellationToken = default) + public Task CreateAsync(WorkflowGraph workflowGraph, WorkflowState workflowState, CancellationToken cancellationToken = default) { - var workflowHost = (IWorkflowHost)ActivatorUtilities.CreateInstance(_serviceProvider, workflow, workflowState); + var workflowHost = (IWorkflowHost)ActivatorUtilities.CreateInstance(_serviceProvider, workflowGraph, workflowState); return Task.FromResult(workflowHost); } /// - public Task CreateAsync(Workflow workflow, CancellationToken cancellationToken = default) + public Task CreateAsync(WorkflowGraph workflowGraph, CancellationToken cancellationToken = default) { + var workflow = workflowGraph.Workflow; var workflowState = new WorkflowState { Id = _identityGenerator.GenerateId(), - DefinitionId = workflow.Identity.DefinitionId, + DefinitionId =workflow.Identity.DefinitionId, DefinitionVersion = workflow.Identity.Version }; - return CreateAsync(workflow, workflowState, cancellationToken); + return CreateAsync(workflowGraph, workflowState, cancellationToken); } } \ No newline at end of file diff --git a/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj b/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj index f70f27f56..6542ae0f9 100644 --- a/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj +++ b/test/component/Elsa.Workflows.ComponentTests/Elsa.Workflows.ComponentTests.csproj @@ -47,6 +47,15 @@ Always + + Always + + + Always + + + Always + diff --git a/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/WorkflowServer.cs b/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/WorkflowServer.cs index 053d0f87f..39d6b2f4f 100644 --- a/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/WorkflowServer.cs +++ b/test/component/Elsa.Workflows.ComponentTests/Helpers/Fixtures/WorkflowServer.cs @@ -22,6 +22,7 @@ public class WorkflowServer(Infrastructure infrastructure, string url) : WebAppl { var client = CreateClient(); client.BaseAddress = new Uri(client.BaseAddress!, "/elsa/api"); + client.Timeout = TimeSpan.FromMinutes(1); return RestService.For(client, CreateRefitSettings()); } @@ -29,6 +30,7 @@ public class WorkflowServer(Infrastructure infrastructure, string url) : WebAppl { var client = CreateClient(); client.BaseAddress = new Uri(client.BaseAddress!, "/workflows/"); + client.Timeout = TimeSpan.FromMinutes(1); return client; } diff --git a/test/component/Elsa.Workflows.ComponentTests/Scenarios/CachingAndWorkflowDefinitionActivity/WorkflowDefinitionActivityTests.cs b/test/component/Elsa.Workflows.ComponentTests/Scenarios/CachingAndWorkflowDefinitionActivity/WorkflowDefinitionActivityTests.cs new file mode 100644 index 000000000..ad9bdac15 --- /dev/null +++ b/test/component/Elsa.Workflows.ComponentTests/Scenarios/CachingAndWorkflowDefinitionActivity/WorkflowDefinitionActivityTests.cs @@ -0,0 +1,53 @@ +using Elsa.Workflows.Management.Contracts; +using Elsa.Workflows.Management.Filters; +using Microsoft.Extensions.DependencyInjection; +using Xunit.Abstractions; + +namespace Elsa.Workflows.ComponentTests.Scenarios.CachingAndWorkflowDefinitionActivity; + +// Tests the behavior of the WorkflowDefinitionActivity. +// See https://github.com/elsa-workflows/elsa-core/issues/5314 +public class WorkflowDefinitionActivityTests : AppComponentTest +{ + private readonly ITestOutputHelper _testOutputHelper; + private const string GrandChildDefinitionId = "29595e7b37a4836d"; + + private readonly IWorkflowDefinitionCacheManager _workflowDefinitionCacheManager; + private readonly IWorkflowInstanceStore _workflowInstanceStore; + private readonly HttpClient _httpWorkflowClient; + + public WorkflowDefinitionActivityTests(App app, ITestOutputHelper testOutputHelper) : base(app) + { + _testOutputHelper = testOutputHelper; + _httpWorkflowClient = WorkflowServer.CreateHttpWorkflowClient(); + _workflowDefinitionCacheManager = Scope.ServiceProvider.GetRequiredService(); + _workflowInstanceStore = Scope.ServiceProvider.GetRequiredService(); + } + + [Fact] + public async Task SendHttpRequest_WhileEvictingCache_ShouldNotGenerateFaults() + { + var requestTasks = Enumerable.Range(0, 200).Select(SendRequestAsync).ToList(); + await Task.WhenAll(requestTasks); + + var filter = new WorkflowInstanceFilter + { + WorkflowSubStatus = WorkflowSubStatus.Faulted + }; + var faultedWorkflows = (await _workflowInstanceStore.FindManyAsync(filter)).ToList(); + var faultCount = faultedWorkflows.Count; + + foreach (var faultedWorkflow in faultedWorkflows) + foreach (var incident in faultedWorkflow.WorkflowState.Incidents) + _testOutputHelper.WriteLine(incident.Message); + + Assert.Equal(0, faultCount); + } + + private async Task SendRequestAsync(int index = 0) + { + var requestTask = _httpWorkflowClient.PostAsync("parent", new StringContent("{}")); + var evictionTask = _workflowDefinitionCacheManager.EvictWorkflowDefinitionAsync(GrandChildDefinitionId); + await Task.WhenAll(requestTask, evictionTask); + } +} \ No newline at end of file diff --git a/test/component/Elsa.Workflows.ComponentTests/Scenarios/CachingAndWorkflowDefinitionActivity/Workflows/workflow-definition-child.json b/test/component/Elsa.Workflows.ComponentTests/Scenarios/CachingAndWorkflowDefinitionActivity/Workflows/workflow-definition-child.json new file mode 100644 index 000000000..616c08190 --- /dev/null +++ b/test/component/Elsa.Workflows.ComponentTests/Scenarios/CachingAndWorkflowDefinitionActivity/Workflows/workflow-definition-child.json @@ -0,0 +1,64 @@ +{ + "id": "26d186d68deb1250", + "definitionId": "a3390d1f4c2594a8", + "name": "Child", + "createdAt": "2024-04-30T08:22:44.543407+00:00", + "version": 1, + "toolVersion": "3.2.0.0", + "variables": [], + "inputs": [], + "outputs": [], + "outcomes": [], + "customProperties": {}, + "isReadonly": false, + "isSystem": false, + "isLatest": true, + "isPublished": true, + "options": { + "usableAsActivity": true, + "autoUpdateConsumingWorkflows": true + }, + "root": { + "type": "Elsa.Flowchart", + "version": 1, + "id": "e6c941b6f2292fcf", + "nodeId": "Workflow2:e6c941b6f2292fcf", + "metadata": {}, + "customProperties": { + "source": "FlowchartJsonConverter.cs:45", + "notFoundConnections": [], + "canStartWorkflow": false, + "runAsynchronously": false + }, + "activities": [ + { + "workflowDefinitionId": "29595e7b37a4836d", + "workflowDefinitionVersionId": "f36dc264f14be3c6", + "latestAvailablePublishedVersion": 1, + "latestAvailablePublishedVersionId": "f36dc264f14be3c6", + "id": "377b23f4235ec26d", + "nodeId": "Workflow2:e6c941b6f2292fcf:377b23f4235ec26d", + "name": "GrandChild1", + "type": "GrandChild", + "version": 1, + "customProperties": { + "canStartWorkflow": false, + "runAsynchronously": false + }, + "metadata": { + "designer": { + "position": { + "x": -193.95703125, + "y": 23 + }, + "size": { + "width": 115.1875, + "height": 50 + } + } + } + } + ], + "connections": [] + } +} \ No newline at end of file diff --git a/test/component/Elsa.Workflows.ComponentTests/Scenarios/CachingAndWorkflowDefinitionActivity/Workflows/workflow-definition-grand-child.json b/test/component/Elsa.Workflows.ComponentTests/Scenarios/CachingAndWorkflowDefinitionActivity/Workflows/workflow-definition-grand-child.json new file mode 100644 index 000000000..83cf09870 --- /dev/null +++ b/test/component/Elsa.Workflows.ComponentTests/Scenarios/CachingAndWorkflowDefinitionActivity/Workflows/workflow-definition-grand-child.json @@ -0,0 +1,67 @@ +{ + "id": "f36dc264f14be3c6", + "definitionId": "29595e7b37a4836d", + "name": "Grand Child", + "createdAt": "2024-04-30T08:22:44.537237+00:00", + "version": 1, + "toolVersion": "3.2.0.0", + "variables": [], + "inputs": [], + "outputs": [], + "outcomes": [], + "customProperties": {}, + "isReadonly": false, + "isSystem": false, + "isLatest": true, + "isPublished": true, + "options": { + "usableAsActivity": true, + "autoUpdateConsumingWorkflows": true + }, + "root": { + "type": "Elsa.Flowchart", + "version": 1, + "id": "29044d8e2ba192c0", + "nodeId": "Workflow1:29044d8e2ba192c0", + "metadata": {}, + "customProperties": { + "source": "FlowchartJsonConverter.cs:45", + "notFoundConnections": [], + "canStartWorkflow": false, + "runAsynchronously": false + }, + "activities": [ + { + "text": { + "typeName": "String", + "expression": { + "type": "Literal", + "value": "Grand Child" + } + }, + "id": "fb00878500fde1d3", + "nodeId": "Workflow1:29044d8e2ba192c0:fb00878500fde1d3", + "name": "WriteLine1", + "type": "Elsa.WriteLine", + "version": 1, + "customProperties": { + "canStartWorkflow": false, + "runAsynchronously": false + }, + "metadata": { + "designer": { + "position": { + "x": -234.796875, + "y": -191 + }, + "size": { + "width": 139.296875, + "height": 50 + } + } + } + } + ], + "connections": [] + } +} \ No newline at end of file diff --git a/test/component/Elsa.Workflows.ComponentTests/Scenarios/CachingAndWorkflowDefinitionActivity/Workflows/workflow-definition-parent.json b/test/component/Elsa.Workflows.ComponentTests/Scenarios/CachingAndWorkflowDefinitionActivity/Workflows/workflow-definition-parent.json new file mode 100644 index 000000000..0adf7b501 --- /dev/null +++ b/test/component/Elsa.Workflows.ComponentTests/Scenarios/CachingAndWorkflowDefinitionActivity/Workflows/workflow-definition-parent.json @@ -0,0 +1,142 @@ +{ + "id": "9178b94f961b1cf", + "definitionId": "189be5173f90b1f6", + "name": "Parent", + "createdAt": "2024-04-30T08:42:17.838484+00:00", + "version": 6, + "toolVersion": "3.2.0.0", + "variables": [], + "inputs": [], + "outputs": [], + "outcomes": [], + "customProperties": { + "Elsa:WorkflowContextProviderTypes": [] + }, + "isReadonly": false, + "isSystem": false, + "isLatest": true, + "isPublished": true, + "options": { + "autoUpdateConsumingWorkflows": false + }, + "root": { + "type": "Elsa.Flowchart", + "version": 1, + "id": "96f60810318a6a89", + "nodeId": "Workflow3:96f60810318a6a89", + "metadata": {}, + "customProperties": { + "source": "FlowchartJsonConverter.cs:45", + "notFoundConnections": [], + "canStartWorkflow": false, + "runAsynchronously": false + }, + "activities": [ + { + "path": { + "typeName": "String", + "expression": { + "type": "Literal", + "value": "parent" + } + }, + "supportedMethods": { + "typeName": "String[]", + "expression": { + "type": "Object", + "value": "[\u0022POST\u0022]" + } + }, + "authorize": { + "typeName": "Boolean", + "expression": { + "type": "Literal", + "value": false + } + }, + "policy": { + "typeName": "String", + "expression": { + "type": "Literal" + } + }, + "requestTimeout": null, + "requestSizeLimit": null, + "fileSizeLimit": null, + "allowedFileExtensions": null, + "blockedFileExtensions": null, + "allowedMimeTypes": null, + "exposeRequestTooLargeOutcome": false, + "exposeFileTooLargeOutcome": false, + "exposeInvalidFileExtensionOutcome": false, + "exposeInvalidFileMimeTypeOutcome": false, + "parsedContent": null, + "files": null, + "routeData": null, + "queryStringData": null, + "headers": null, + "result": null, + "id": "7e20498ffc94dbc6", + "nodeId": "Workflow3:96f60810318a6a89:7e20498ffc94dbc6", + "name": "HttpEndpoint1", + "type": "Elsa.HttpEndpoint", + "version": 1, + "customProperties": { + "canStartWorkflow": true, + "runAsynchronously": false + }, + "metadata": { + "designer": { + "position": { + "x": -520, + "y": -120 + }, + "size": { + "width": 176.390625, + "height": 50 + } + } + } + }, + { + "workflowDefinitionId": "a3390d1f4c2594a8", + "workflowDefinitionVersionId": "26d186d68deb1250", + "latestAvailablePublishedVersion": 1, + "latestAvailablePublishedVersionId": "26d186d68deb1250", + "id": "ed6a03040f651906", + "nodeId": "Workflow3:96f60810318a6a89:ed6a03040f651906", + "name": "Child1", + "type": "Child", + "version": 1, + "customProperties": { + "canStartWorkflow": false, + "runAsynchronously": false + }, + "metadata": { + "designer": { + "position": { + "x": -236.5810546875, + "y": -120 + }, + "size": { + "width": 67.7734375, + "height": 50 + } + } + } + } + ], + "connections": [ + { + "source": { + "activity": "7e20498ffc94dbc6", + "port": "Done" + }, + "target": { + "activity": "ed6a03040f651906", + "port": "In" + } + } + ] + } +} \ No newline at end of file diff --git a/test/performance/Elsa.Workflows.PerformanceTests/Elsa.Workflows.PerformanceTests.csproj b/test/performance/Elsa.Workflows.PerformanceTests/Elsa.Workflows.PerformanceTests.csproj index 085bdddfe..c0dea99e4 100644 --- a/test/performance/Elsa.Workflows.PerformanceTests/Elsa.Workflows.PerformanceTests.csproj +++ b/test/performance/Elsa.Workflows.PerformanceTests/Elsa.Workflows.PerformanceTests.csproj @@ -9,13 +9,6 @@ true - - - - - - - diff --git a/test/unit/Elsa.Workflows.Core.UnitTests/ScheduleActivityExecutionContextTests.cs b/test/unit/Elsa.Workflows.Core.UnitTests/ScheduleActivityExecutionContextTests.cs new file mode 100644 index 000000000..3c7f48fd1 --- /dev/null +++ b/test/unit/Elsa.Workflows.Core.UnitTests/ScheduleActivityExecutionContextTests.cs @@ -0,0 +1,30 @@ +using Elsa.Testing.Shared; +using Elsa.Workflows.Activities; +using Elsa.Workflows.Contracts; +using Microsoft.Extensions.DependencyInjection; +using Xunit.Abstractions; + +namespace Elsa.Workflows.Core.UnitTests; + +public class ScheduleActivityExecutionContextTests(ITestOutputHelper testOutputHelper) +{ + private readonly IServiceProvider _serviceProvider = new TestApplicationBuilder(testOutputHelper).Build(); + + [Fact(DisplayName = "Scheduling an activity that is not part of the workflow should throw an exception")] + public async Task ScheduleActivityAsync_WithActivityNotPartOfWorkflow_ShouldThrowException() + { + await _serviceProvider.PopulateRegistriesAsync(); + var writeLineA = new WriteLine("Test"); + var writeLineB = new WriteLine("Test"); + var workflow = new Workflow + { + Root = writeLineA + }; + var workflowGraphBuilder = _serviceProvider.GetRequiredService(); + var workflowGraph = await workflowGraphBuilder.BuildAsync(workflow); + var workflowExecutionContext = await WorkflowExecutionContext.CreateAsync(_serviceProvider, workflowGraph, "test"); + var activityExecutionContext = workflowExecutionContext.CreateActivityExecutionContext(writeLineA); + + await Assert.ThrowsAsync(() => activityExecutionContext.ScheduleActivityAsync(writeLineB).AsTask()); + } +} \ No newline at end of file