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