From 3f5cac76c550d4e08f4ff9e77b6fe28eec9dfc3a Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Tue, 15 Jul 2025 13:59:05 +0200 Subject: [PATCH] Refactors workflow reference updates (#6792) * Refactor database extensions and support migrations for V3.6 Remove `DatabaseFacadeExtensions` and introduce `IWorkflowReferenceQuery` with its default implementation. Implement database schema updates for PostgreSQL, MySQL, and Oracle to enhance compatibility with the V3.6 data structure. * Remove commented-out code and standardize null default assignment in `IWorkflowDefinitionStore` interface * Add XML documentation for `DefaultWorkflowReferenceQuery` detailing its purpose and dependencies * Refactor `WorkflowReferenceUpdater` to support recursive dependency resolution, prevent concurrent updates, and improve reference consistency. * Simplify `WorkflowReferenceUpdater` by removing topological sorting and redundant dependencies handling. * Refactor `WorkflowReferenceUpdater` to streamline reference updates, remove redundant logic, and enhance dependency resolution efficiency. * Refactor `WorkflowReferenceUpdater` to use `HashSet` for updated workflows, reducing potential duplication and improving performance. * Refactor `WorkflowReferenceUpdater` to introduce topological sorting for correct processing order, improve dependency resolution, and enhance clarity with updated records and comments. * Introduce `WorkflowDefinitionActivityDescriptorFactory` to simplify `WorkflowDefinitionActivity` descriptor creation and refactor existing components for modularity, clarity, and efficiency. * Update `WorkflowReferenceUpdater` to use `VersionOptions.Latest` instead of `VersionOptions.LatestOrPublished` for workflow reference resolution. * Refactor `WorkflowReferenceUpdater` to improve workflow dependency resolution by handling publication states, caching drafts more efficiently, and introducing distinct processing for latest and published versions. * Refactor workflow publication logic and update SQLite configuration. Removed unused draft publication logic to simplify workflow reference updates. Updated SQLite persistence configuration in `Elsa.Server.Agents.Web` to use explicit connection strings for improved clarity and maintainability. * Remove commented-out legacy code in `WorkflowReferenceUpdater` to improve clarity and maintainability. * Fix formatting by adding a missing newline at EOF in `Directory.Build.props`. * Prevent infinite recursion in `GetReferencingWorkflowDefinitionIdsAsync` by introducing visited ID tracking. Fix formatting inconsistencies in `WorkflowReferenceUpdater`. * Update `WorkflowReferenceUpdater` to use `NewGraph` instead of materializing workflows for referencing workflow graphs --- .../ContentWriters/RawStringContent.cs | 4 +- ...flowDefinitionActivityDescriptorFactory.cs | 120 ++++++++ .../WorkflowDefinitionActivityProvider.cs | 136 +-------- .../Features/WorkflowManagementFeature.cs | 1 + .../Services/DefaultWorkflowReferenceQuery.cs | 9 + .../Services/WorkflowReferenceUpdater.cs | 288 +++++++++++++----- .../Activities/RunTask.cs | 10 + 7 files changed, 352 insertions(+), 216 deletions(-) create mode 100644 src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivityDescriptorFactory.cs diff --git a/src/modules/Elsa.Http/ContentWriters/RawStringContent.cs b/src/modules/Elsa.Http/ContentWriters/RawStringContent.cs index 358131ddf..c27d476e6 100644 --- a/src/modules/Elsa.Http/ContentWriters/RawStringContent.cs +++ b/src/modules/Elsa.Http/ContentWriters/RawStringContent.cs @@ -34,9 +34,9 @@ public class RawStringContent : HttpContent /// protected override async Task SerializeToStreamAsync(Stream stream, TransportContext? context, CancellationToken cancellationToken) { - using var writer = new StreamWriter(stream, _encoding, leaveOpen: true); + await using var writer = new StreamWriter(stream, _encoding, leaveOpen: true); await writer.WriteAsync(_content.AsMemory(), cancellationToken); - await writer.FlushAsync(); + await writer.FlushAsync(cancellationToken); } /// diff --git a/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivityDescriptorFactory.cs b/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivityDescriptorFactory.cs new file mode 100644 index 000000000..cb468def0 --- /dev/null +++ b/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivityDescriptorFactory.cs @@ -0,0 +1,120 @@ +using Elsa.Extensions; +using Elsa.Workflows.Management.Entities; +using Elsa.Workflows.Models; +using Humanizer; + +namespace Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity; + +public class WorkflowDefinitionActivityDescriptorFactory(IActivityFactory activityFactory) +{ + public ActivityDescriptor CreateDescriptor(WorkflowDefinition definition, WorkflowDefinition? latestPublishedDefinition = null) + { + var typeName = definition.Name!.Pascalize(); + + var ports = definition.Outcomes.Select(outcome => new Port + { + Name = outcome, + DisplayName = outcome, + IsBrowsable = true, + Type = PortType.Flow + }).ToList(); + + var rootPort = new Port + { + Name = nameof(WorkflowDefinitionActivity.Root), + DisplayName = "Root", + IsBrowsable = false, + Type = PortType.Embedded + }; + + ports.Insert(0, rootPort); + + return new() + { + TypeName = typeName, + Name = typeName, + Version = definition.Version, + DisplayName = definition.Name, + Description = definition.Description, + Category = definition.Options.ActivityCategory ?? "Workflows", + Kind = ActivityKind.Action, + IsBrowsable = definition.IsPublished, + Inputs = DescribeInputs(definition).ToList(), + Outputs = DescribeOutputs(definition).ToList(), + Ports = ports, + CustomProperties = + { + ["RootType"] = nameof(WorkflowDefinitionActivity), + ["WorkflowDefinitionId"] = definition.DefinitionId, + ["WorkflowDefinitionVersionId"] = definition.Id + }, + ConstructionProperties = new Dictionary + { + [nameof(WorkflowDefinitionActivity.WorkflowDefinitionId)] = definition.DefinitionId, + [nameof(WorkflowDefinitionActivity.WorkflowDefinitionVersionId)] = definition.Id, + [nameof(WorkflowDefinitionActivity.Version)] = definition.Version, + }, + Constructor = context => + { + var activity = (WorkflowDefinitionActivity)activityFactory.Create(typeof(WorkflowDefinitionActivity), context); + activity.Type = typeName; + activity.WorkflowDefinitionId = definition.DefinitionId; + activity.WorkflowDefinitionVersionId = definition.Id; + activity.Version = definition.Version; + activity.LatestAvailablePublishedVersion = latestPublishedDefinition?.Version ?? definition.Version; + activity.LatestAvailablePublishedVersionId = latestPublishedDefinition?.Id ?? definition.Id; + + return activity; + } + }; + } + + private static IEnumerable DescribeInputs(WorkflowDefinition definition) + { + var inputs = definition.Inputs.Select(inputDefinition => + { + var nakedType = inputDefinition.Type; + var inputName = inputDefinition.Name; + var safeInputName = PropertyNameHelper.GetSafePropertyName(typeof(WorkflowDefinitionActivity), inputName); + + return new InputDescriptor + { + Type = nakedType, + IsWrapped = true, + ValueGetter = activity => activity.SyntheticProperties.GetValueOrDefault(safeInputName), + ValueSetter = (activity, value) => activity.SyntheticProperties[safeInputName] = value!, + Name = safeInputName, + DisplayName = inputDefinition.DisplayName, + Description = inputDefinition.Description, + Category = inputDefinition.Category, + UIHint = inputDefinition.UIHint, + StorageDriverType = inputDefinition.StorageDriverType, + IsSynthetic = true + }; + }); + + foreach (var input in inputs) + yield return input; + } + + private static IEnumerable DescribeOutputs(WorkflowDefinition definition) + { + return definition.Outputs.Select(outputDefinition => + { + var nakedType = outputDefinition.Type; + var outputName = outputDefinition.Name; + var safeOutputName = PropertyNameHelper.GetSafePropertyName(typeof(WorkflowDefinitionActivity), outputName); + + return new OutputDescriptor + { + Type = nakedType, + ValueGetter = activity => activity.SyntheticProperties.GetValueOrDefault(safeOutputName), + ValueSetter = (activity, value) => activity.SyntheticProperties[safeOutputName] = value!, + Name = safeOutputName, + DisplayName = outputDefinition.DisplayName, + Description = outputDefinition.Description, + IsSynthetic = true + }; + }); + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivityProvider.cs b/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivityProvider.cs index dbf2cd3d8..4c0699908 100644 --- a/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivityProvider.cs +++ b/src/modules/Elsa.Workflows.Management/Activities/WorkflowDefinitionActivity/WorkflowDefinitionActivityProvider.cs @@ -1,18 +1,14 @@ using Elsa.Common.Models; -using Elsa.Extensions; using Elsa.Workflows.Management.Entities; using Elsa.Workflows.Management.Filters; using Elsa.Workflows.Models; -using Elsa.Workflows.Serialization.Converters; -using Elsa.Workflows.Serialization.Helpers; -using Humanizer; namespace Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity; /// /// Provides activity descriptors based on s stored in the database. /// -public class WorkflowDefinitionActivityProvider(IWorkflowDefinitionStore store, IActivityFactory activityFactory, ActivityWriter activityWriter) : IActivityProvider +public class WorkflowDefinitionActivityProvider(IWorkflowDefinitionStore store, WorkflowDefinitionActivityDescriptorFactory workflowDefinitionActivityDescriptorFactory) : IActivityProvider { /// public async ValueTask> GetDescriptorsAsync(CancellationToken cancellationToken = default) @@ -22,24 +18,9 @@ public class WorkflowDefinitionActivityProvider(IWorkflowDefinitionStore store, UsableAsActivity = true, VersionOptions = VersionOptions.All }; - - var allDescriptors = new List(); - var currentPage = 0; - const int pageSize = 100; - - while (true) - { - var pageArgs = PageArgs.FromPage(currentPage++, pageSize); - var pageOfDefinitions = await store.FindManyAsync(filter, pageArgs, cancellationToken); - var descriptors = CreateDescriptors(pageOfDefinitions.Items).ToList(); - - allDescriptors.AddRange(descriptors); - - if (allDescriptors.Count >= pageOfDefinitions.TotalCount) - break; - } - - return allDescriptors; + + var definitions = (await store.FindManyAsync(filter, cancellationToken)).ToList(); + return CreateDescriptors(definitions).ToList(); } private IEnumerable CreateDescriptors(ICollection definitions) @@ -49,116 +30,9 @@ public class WorkflowDefinitionActivityProvider(IWorkflowDefinitionStore store, private ActivityDescriptor CreateDescriptor(WorkflowDefinition definition, ICollection allDefinitions) { - var typeName = definition.Name!.Pascalize(); - var latestPublishedVersion = allDefinitions .Where(x => x.DefinitionId == definition.DefinitionId && x.IsPublished) .MaxBy(x => x.Version); - - var ports = definition.Outcomes.Select(outcome => new Port - { - Name = outcome, - DisplayName = outcome, - IsBrowsable = true, - Type = PortType.Flow - }).ToList(); - - var rootPort = new Port - { - Name = nameof(WorkflowDefinitionActivity.Root), - DisplayName = "Root", - IsBrowsable = false, - Type = PortType.Embedded - }; - - ports.Insert(0, rootPort); - - return new() - { - TypeName = typeName, - Name = typeName, - Version = definition.Version, - DisplayName = definition.Name, - Description = definition.Description, - Category = definition.Options.ActivityCategory ?? "Workflows", - Kind = ActivityKind.Action, - IsBrowsable = definition.IsPublished, - Inputs = DescribeInputs(definition).ToList(), - Outputs = DescribeOutputs(definition).ToList(), - Ports = ports, - CustomProperties = - { - ["RootType"] = nameof(WorkflowDefinitionActivity), - ["WorkflowDefinitionId"] = definition.DefinitionId, - ["WorkflowDefinitionVersionId"] = definition.Id - }, - ConstructionProperties = new Dictionary - { - [nameof(WorkflowDefinitionActivity.WorkflowDefinitionId)] = definition.DefinitionId, - [nameof(WorkflowDefinitionActivity.WorkflowDefinitionVersionId)] = definition.Id, - [nameof(WorkflowDefinitionActivity.Version)] = definition.Version, - }, - Constructor = context => - { - var activity = (WorkflowDefinitionActivity)activityFactory.Create(typeof(WorkflowDefinitionActivity), context); - activity.Type = typeName; - activity.WorkflowDefinitionId = definition.DefinitionId; - activity.WorkflowDefinitionVersionId = definition.Id; - activity.Version = definition.Version; - activity.LatestAvailablePublishedVersion = latestPublishedVersion?.Version ?? 0; - activity.LatestAvailablePublishedVersionId = latestPublishedVersion?.Id; - - return activity; - } - }; - } - - private static IEnumerable DescribeInputs(WorkflowDefinition definition) - { - var inputs = definition.Inputs.Select(inputDefinition => - { - var nakedType = inputDefinition.Type; - var inputName = inputDefinition.Name; - var safeInputName = PropertyNameHelper.GetSafePropertyName(typeof(WorkflowDefinitionActivity), inputName); - - return new InputDescriptor - { - Type = nakedType, - IsWrapped = true, - ValueGetter = activity => activity.SyntheticProperties.GetValueOrDefault(safeInputName), - ValueSetter = (activity, value) => activity.SyntheticProperties[safeInputName] = value!, - Name = safeInputName, - DisplayName = inputDefinition.DisplayName, - Description = inputDefinition.Description, - Category = inputDefinition.Category, - UIHint = inputDefinition.UIHint, - StorageDriverType = inputDefinition.StorageDriverType, - IsSynthetic = true - }; - }); - - foreach (var input in inputs) - yield return input; - } - - private static IEnumerable DescribeOutputs(WorkflowDefinition definition) - { - return definition.Outputs.Select(outputDefinition => - { - var nakedType = outputDefinition.Type; - var outputName = outputDefinition.Name; - var safeOutputName = PropertyNameHelper.GetSafePropertyName(typeof(WorkflowDefinitionActivity), outputName); - - return new OutputDescriptor - { - Type = nakedType, - ValueGetter = activity => activity.SyntheticProperties.GetValueOrDefault(safeOutputName), - ValueSetter = (activity, value) => activity.SyntheticProperties[safeOutputName] = value!, - Name = safeOutputName, - DisplayName = outputDefinition.DisplayName, - Description = outputDefinition.Description, - IsSynthetic = true - }; - }); + return workflowDefinitionActivityDescriptorFactory.CreateDescriptor(definition, latestPublishedVersion); } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs b/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs index 30f7c7f18..589b2e847 100644 --- a/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs +++ b/src/modules/Elsa.Workflows.Management/Features/WorkflowManagementFeature.cs @@ -226,6 +226,7 @@ public class WorkflowManagementFeature(IModule module) : FeatureBase(module) .AddMemoryStore() .AddActivityProvider() .AddActivityProvider() + .AddScoped() .AddScoped() .AddScoped() .AddScoped() diff --git a/src/modules/Elsa.Workflows.Management/Services/DefaultWorkflowReferenceQuery.cs b/src/modules/Elsa.Workflows.Management/Services/DefaultWorkflowReferenceQuery.cs index b0bf0ffab..b3491a2d7 100644 --- a/src/modules/Elsa.Workflows.Management/Services/DefaultWorkflowReferenceQuery.cs +++ b/src/modules/Elsa.Workflows.Management/Services/DefaultWorkflowReferenceQuery.cs @@ -5,6 +5,15 @@ using Elsa.Workflows.Models; namespace Elsa.Workflows.Management.Services; +/// +/// Represents the default implementation of that queries workflows +/// referencing a specific workflow definition. +/// +/// +/// This class is designed to identify workflows that are dependent on a given workflow definition. +/// It leverages services such as and +/// to locate, graph, and analyze workflows that include the specified workflow definition. +/// public class DefaultWorkflowReferenceQuery(IWorkflowDefinitionService workflowDefinitionService, IWorkflowDefinitionStore workflowDefinitionStore) : IWorkflowReferenceQuery { public async Task> ExecuteAsync(string workflowDefinitionId, CancellationToken cancellationToken = default) diff --git a/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs b/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs index 988fdc7bf..a7eafc1a6 100644 --- a/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs +++ b/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs @@ -1,107 +1,229 @@ +using System.Runtime.CompilerServices; using Elsa.Common.Models; +using Elsa.Extensions; using Elsa.Workflows.Activities; using Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity; using Elsa.Workflows.Management.Entities; -using Elsa.Workflows.Management.Filters; using Elsa.Workflows.Management.Models; using Elsa.Workflows.Models; namespace Elsa.Workflows.Management.Services; -/// +internal record WorkflowReferences(string ReferencedDefinitionId, ICollection ReferencingDefinitionIds); + +internal record UpdatedWorkflowDefinition(WorkflowDefinition Definition, WorkflowGraph NewGraph); + public class WorkflowReferenceUpdater( IWorkflowDefinitionPublisher publisher, IWorkflowDefinitionService workflowDefinitionService, + IWorkflowDefinitionStore workflowDefinitionStore, IWorkflowReferenceQuery workflowReferenceQuery, - IApiSerializer serializer) : IWorkflowReferenceUpdater + WorkflowDefinitionActivityDescriptorFactory workflowDefinitionActivityDescriptorFactory, + IActivityRegistry activityRegistry, + IApiSerializer serializer) + : IWorkflowReferenceUpdater { - /// - public async Task UpdateWorkflowReferencesAsync(WorkflowDefinition referencedDefinition, CancellationToken cancellationToken = default) + private bool _isUpdating; + + public async Task UpdateWorkflowReferencesAsync( + WorkflowDefinition referencedDefinition, + CancellationToken cancellationToken = default) { - // Skip if the published workflow definition is not usable as an activity or does not auto-update consuming workflows. - if (referencedDefinition.Options is not { UsableAsActivity: true, AutoUpdateConsumingWorkflows: true }) + if (_isUpdating || + referencedDefinition.Options is not { UsableAsActivity: true, AutoUpdateConsumingWorkflows: true }) return new([]); - // Find all workflow definitions that reference the updated workflow definition. - var referencedDefinitionIds = (await workflowReferenceQuery.ExecuteAsync(referencedDefinition.DefinitionId, cancellationToken)).ToList(); - var filter = new WorkflowDefinitionFilter - { - DefinitionIds = referencedDefinitionIds, - VersionOptions = VersionOptions.Latest, - IsReadonly = false - }; - var consumingWorkflowGraphs = await workflowDefinitionService.FindWorkflowGraphsAsync(filter, cancellationToken); + var allWorkflowReferences = await GetReferencingWorkflowDefinitionIdsAsync(referencedDefinition.DefinitionId, cancellationToken).ToListAsync(cancellationToken); + var filteredWorkflowReferences = allWorkflowReferences + .Where(r => r.ReferencingDefinitionIds.Any()) + .DistinctBy(r => r.ReferencedDefinitionId) + .ToList(); - // Update consuming workflows. - var updatedWorkflows = new List(); - foreach (var workflowGraph in consumingWorkflowGraphs) - { - var newDefinition = await UpdateConsumingWorkflowAsync(workflowGraph, referencedDefinition, cancellationToken); + var referencingIds = filteredWorkflowReferences.SelectMany(r => r.ReferencingDefinitionIds).Distinct().ToList(); + var referencedIds = filteredWorkflowReferences.Select(r => r.ReferencedDefinitionId).Distinct().ToList(); - if (newDefinition != null) - updatedWorkflows.Add(newDefinition); - } - - return new(updatedWorkflows); - } - - private async Task UpdateConsumingWorkflowAsync(WorkflowGraph workflowGraph, WorkflowDefinition definition, CancellationToken cancellationToken) - { - // Create a new version of the published workflow definition or get the existing draft. - var consumerDefinitionId = workflowGraph.Workflow.Identity.DefinitionId; - var originalVersionIsPublished = workflowGraph.Workflow.Publication.IsPublished; - var newVersion = await publisher.GetDraftAsync(consumerDefinitionId, VersionOptions.LatestOrPublished, cancellationToken); - - // This is null in case the definition no longer exists in the store. - if (newVersion == null) - return null; - - // Materialize the draft to find all workflow definition activities that use the updated workflow definition. - var newWorkflowGraph = await workflowDefinitionService.MaterializeWorkflowAsync(newVersion, cancellationToken); - var outdatedWorkflowDefinitionActivities = FindOutdatedWorkflowDefinitionActivities(newWorkflowGraph, definition).ToList(); - - // Skip if the new version of the published workflow definition is not used in the workflow or if the activity is already up to date. - if (outdatedWorkflowDefinitionActivities.Count == 0) - return null; - - // Update the consuming workflow graph to use the new version of the published workflow definition. - foreach (var workflowDefinitionActivity in outdatedWorkflowDefinitionActivities) - { - workflowDefinitionActivity.WorkflowDefinitionVersionId = definition.Id; - workflowDefinitionActivity.LatestAvailablePublishedVersionId = definition.Id; - workflowDefinitionActivity.Version = definition.Version; - } - - // Update the new version of the published workflow definition. - if (newWorkflowGraph.Root.Activity is Workflow newWorkflow) - newVersion.StringData = serializer.Serialize(newWorkflow.Root); - - // If the draft is new, publish it. - if (originalVersionIsPublished) - await publisher.PublishAsync(newVersion, cancellationToken); - else - await publisher.SaveDraftAsync(newVersion, cancellationToken); - - return newVersion; - } - - private IEnumerable FindOutdatedWorkflowDefinitionActivities(WorkflowGraph workflowGraph, WorkflowDefinition updatedDefinition) - { - return FindWorkflowActivityDefinitionActivityNodes(workflowGraph.Root) - .Where(x => x.WorkflowDefinitionId == updatedDefinition.DefinitionId && x.WorkflowDefinitionVersionId != updatedDefinition.Id); - } - - private IEnumerable FindWorkflowActivityDefinitionActivityNodes(ActivityNode parent) - { - foreach (var child in parent.Children) - { - if (child.Activity is WorkflowDefinitionActivity workflowDefinitionActivity) - yield return workflowDefinitionActivity; - else + var referencingWorkflowGraphs = (await workflowDefinitionService.FindWorkflowGraphsAsync(new() { - foreach (var grandChild in FindWorkflowActivityDefinitionActivityNodes(child)) - yield return grandChild; + DefinitionIds = referencingIds, + VersionOptions = VersionOptions.Latest, + IsReadonly = false + }, cancellationToken)) + .ToDictionary(g => g.Workflow.Identity.DefinitionId); + + var referencedWorkflowDefinitionList = (await workflowDefinitionStore.FindManyAsync(new() + { + DefinitionIds = referencedIds, + VersionOptions = VersionOptions.Published, + IsReadonly = false + }, cancellationToken)).ToList(); + + var referencedWorkflowDefinitionsPublished = referencedWorkflowDefinitionList + .GroupBy(x => x.DefinitionId) + .Select(group => + { + var publishedVersion = group.FirstOrDefault(x => x.IsPublished); + return publishedVersion ?? group.First(); + }) + .ToDictionary(d => d.DefinitionId); + + var initialPublicationState = new Dictionary(); + + foreach (var workflowGraph in referencingWorkflowGraphs) + initialPublicationState[workflowGraph.Key] = workflowGraph.Value.Workflow.Publication.IsPublished; + + // Add the initially referenced definition + referencedWorkflowDefinitionsPublished[referencedDefinition.DefinitionId] = referencedDefinition; + + // Build dependency map for topological sorting + var dependencyMap = filteredWorkflowReferences + .SelectMany(r => r.ReferencingDefinitionIds.Select(id => (id, r.ReferencedDefinitionId))) + .ToLookup(x => x.id, x => x.ReferencedDefinitionId); + + // Perform topological sort to ensure dependent workflows are processed in the right order + var sortedWorkflowIds = referencingIds + .TSort(id => dependencyMap[id], true) + // Only process workflows that exist in our referencing workflows dictionary + .Where(id => referencingWorkflowGraphs.ContainsKey(id)) + .ToList(); + + var updatedWorkflows = new Dictionary(); + + // Create a cache for drafts that we've already created during this operation + var draftCache = new Dictionary(); + + foreach (var id in sortedWorkflowIds) + { + if (!referencingWorkflowGraphs.TryGetValue(id, out var graph) || !dependencyMap[id].Any()) + continue; + + foreach (var refId in dependencyMap[id]) + { + var target = referencedWorkflowDefinitionsPublished.GetValueOrDefault(refId); + if (target == null) continue; + + var updated = await UpdateWorkflowAsync(graph, target, draftCache, initialPublicationState, cancellationToken); + if (updated == null) continue; + + graph = updated.NewGraph; + updatedWorkflows[updated.Definition.DefinitionId] = updated; + referencedWorkflowDefinitionsPublished[id] = updated.Definition; + draftCache[id] = updated.Definition; + referencingWorkflowGraphs[id] = updated.NewGraph; } } + + _isUpdating = true; + foreach (var updatedWorkflow in updatedWorkflows.Values) + { + var requiresPublication = initialPublicationState.GetValueOrDefault(updatedWorkflow.Definition.DefinitionId); + if (requiresPublication) + await publisher.PublishAsync(updatedWorkflow.Definition, cancellationToken); + else + await publisher.SaveDraftAsync(updatedWorkflow.Definition, cancellationToken); + } + + _isUpdating = false; + + return new(updatedWorkflows.Select(u => u.Value.Definition)); + } + + private async IAsyncEnumerable GetReferencingWorkflowDefinitionIdsAsync( + string definitionId, + [EnumeratorCancellation] CancellationToken cancellationToken, + HashSet? visitedIds = null) + { + visitedIds ??= new(); + + // If we've already processed this definition ID, skip it to prevent infinite recursion. + if (!visitedIds.Add(definitionId)) + yield break; + + var refs = (await workflowReferenceQuery.ExecuteAsync(definitionId, cancellationToken)).ToList(); + yield return new(definitionId, refs); + + foreach (var id in refs) + { + await foreach (var child in GetReferencingWorkflowDefinitionIdsAsync(id, cancellationToken, visitedIds)) + yield return child; + } + } + + private async Task UpdateWorkflowAsync( + WorkflowGraph graph, + WorkflowDefinition target, + Dictionary draftCache, + Dictionary initialPublicationState, + CancellationToken cancellationToken) + { + var willTargetBePublished = initialPublicationState.GetValueOrDefault(target.DefinitionId, target.IsPublished); + if (!willTargetBePublished) + return null; + + var id = graph.Workflow.Identity.DefinitionId; + var draft = await GetOrCreateDraftAsync(id, draftCache, cancellationToken); + if (draft == null) return null; + + var newGraph = await workflowDefinitionService.MaterializeWorkflowAsync(draft, cancellationToken); + var outdated = FindActivities(newGraph.Root, target.DefinitionId) + .Where(a => a.WorkflowDefinitionVersionId != target.Id) + .ToList(); + + if (!outdated.Any()) return null; + + foreach (var act in outdated) + { + act.WorkflowDefinitionVersionId = target.Id; + act.Version = target.Version; + act.LatestAvailablePublishedVersionId = target.Id; + act.LatestAvailablePublishedVersion = target.Version; + } + + if (newGraph.Root.Activity is Workflow wf) + draft.StringData = serializer.Serialize(wf.Root); + + return new(draft, newGraph); + } + + private async Task GetOrCreateDraftAsync( + string definitionId, + Dictionary draftCache, + CancellationToken cancellationToken) + { + // Check if we already have a draft for this workflow + if (draftCache.TryGetValue(definitionId, out var cachedDraft)) + return cachedDraft; + + // Create or get a draft for this workflow + var draft = await publisher.GetDraftAsync(definitionId, VersionOptions.Latest, cancellationToken); + if (draft == null) return null; + + // Store the draft in the cache for potential future use + draftCache[definitionId] = draft; + + // Get the current published version of the workflow definition. + var publishedVersion = await workflowDefinitionStore.FindAsync( + WorkflowDefinitionHandle.ByDefinitionId(definitionId, VersionOptions.Published).ToFilter(), + cancellationToken); + + // Update the activity registry to be able to materialize the workflow. + var activityDescriptor = workflowDefinitionActivityDescriptorFactory.CreateDescriptor(draft, publishedVersion); + activityRegistry.Add(typeof(WorkflowDefinitionActivityProvider), activityDescriptor); + + return draft; + } + + private static IEnumerable FindActivities(ActivityNode node, string definitionId) + { + // Do not drill into activities that are WorkflowDefinitionActivity + if (node.Activity is WorkflowDefinitionActivity) + yield break; + + foreach (var child in node.Children) + { + if (child.Activity is WorkflowDefinitionActivity activity && activity.WorkflowDefinitionId == definitionId) + yield return activity; + foreach (var grandChildActivity in FindActivities(child, definitionId)) + yield return grandChildActivity; + } } } \ No newline at end of file diff --git a/src/modules/Elsa.Workflows.Runtime/Activities/RunTask.cs b/src/modules/Elsa.Workflows.Runtime/Activities/RunTask.cs index fa55186a1..b8c2a66a9 100644 --- a/src/modules/Elsa.Workflows.Runtime/Activities/RunTask.cs +++ b/src/modules/Elsa.Workflows.Runtime/Activities/RunTask.cs @@ -38,6 +38,12 @@ public class RunTask : Activity [Input(Description = "Any additional parameters to send to the task.")] public Input?> Payload { get; set; } = null!; + /// + /// The ID of the task that was requested to run. + /// + [Output(Description = "The ID of the task that was requested to run.")] + public Output TaskId { get; set; } = null!; + /// [JsonConstructor] private RunTask(string? source = null, int? line = null) : base(source, line) @@ -95,6 +101,10 @@ public class RunTask : Activity var runTaskRequest = new RunTaskRequest(context, taskId, taskName, taskParams); var dispatcher = context.GetRequiredService(); + // Set the task ID output. + TaskId.Set(context, taskId); + + // Dispatch the task request. await dispatcher.DispatchAsync(runTaskRequest, context.CancellationToken); }