< Summary

Information
Class: Elsa.Workflows.Management.Services.UpdatedWorkflowDefinition
Assembly: Elsa.Workflows.Management
File(s): /home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs
Line coverage
100%
Covered lines: 1
Uncovered lines: 0
Coverable lines: 1
Total lines: 239
Line coverage: 100%
Branch coverage
N/A
Covered branches: 0
Total branches: 0
Branch coverage: N/A
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
get_Definition()100%11100%

File(s)

/home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs

#LineLine coverage
 1using Elsa.Common.Models;
 2using Elsa.Extensions;
 3using Elsa.Workflows.Activities;
 4using Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity;
 5using Elsa.Workflows.Management.Entities;
 6using Elsa.Workflows.Management.Filters;
 7using Elsa.Workflows.Management.Materializers;
 8using Elsa.Workflows.Management.Models;
 9using Elsa.Workflows.Models;
 10using Microsoft.Extensions.Logging;
 11
 12namespace Elsa.Workflows.Management.Services;
 13
 1814internal record UpdatedWorkflowDefinition(WorkflowDefinition Definition, WorkflowGraph NewGraph);
 15
 16public class WorkflowReferenceUpdater(
 17    IWorkflowDefinitionPublisher publisher,
 18    IWorkflowDefinitionService workflowDefinitionService,
 19    IWorkflowDefinitionStore workflowDefinitionStore,
 20    IWorkflowReferenceGraphBuilder workflowReferenceGraphBuilder,
 21    WorkflowDefinitionActivityDescriptorFactory workflowDefinitionActivityDescriptorFactory,
 22    IActivityRegistry activityRegistry,
 23    IApiSerializer serializer,
 24    ILogger<WorkflowReferenceUpdater> logger)
 25    : IWorkflowReferenceUpdater
 26{
 27    private bool _isUpdating;
 28
 29    public async Task<UpdateWorkflowReferencesResult> UpdateWorkflowReferencesAsync(
 30        WorkflowDefinition referencedDefinition,
 31        CancellationToken cancellationToken = default)
 32    {
 33        if (_isUpdating ||
 34            referencedDefinition.Options is not { UsableAsActivity: true, AutoUpdateConsumingWorkflows: true })
 35            return new([]);
 36
 37        var referenceGraph = await workflowReferenceGraphBuilder.BuildGraphAsync(referencedDefinition.DefinitionId, canc
 38
 39        // Get all consumer (source) and referenced (target) IDs from the edges
 40        var referencingIds = referenceGraph.Edges.Select(e => e.Source).Distinct().ToList();
 41        var referencedIds = referenceGraph.Edges.Select(e => e.Target).Distinct().ToList();
 42
 43        if (referencingIds.Count == 0)
 44            return new([]);
 45
 46        var referencingWorkflowGraphs = (await workflowDefinitionService.FindWorkflowGraphsAsync(new()
 47            {
 48                DefinitionIds = referencingIds,
 49                VersionOptions = VersionOptions.Latest,
 50                IsReadonly = false
 51            }, cancellationToken))
 52            .ToDictionary(g => g.Workflow.Identity.DefinitionId);
 53
 54        var referencedWorkflowDefinitionList = (await workflowDefinitionStore.FindManyAsync(new()
 55        {
 56            DefinitionIds = referencedIds,
 57            VersionOptions = VersionOptions.Published,
 58            IsReadonly = false
 59        }, cancellationToken)).ToList();
 60
 61        var referencedWorkflowDefinitionsPublished = referencedWorkflowDefinitionList
 62            .GroupBy(x => x.DefinitionId)
 63            .Select(group =>
 64            {
 65                var publishedVersion = group.FirstOrDefault(x => x.IsPublished);
 66                return publishedVersion ?? group.First();
 67            })
 68            .ToDictionary(d => d.DefinitionId);
 69
 70        var initialPublicationState = new Dictionary<string, bool>();
 71
 72        foreach (var workflowGraph in referencingWorkflowGraphs)
 73            initialPublicationState[workflowGraph.Key] = workflowGraph.Value.Workflow.Publication.IsPublished;
 74
 75        // Add the initially referenced definition
 76        referencedWorkflowDefinitionsPublished[referencedDefinition.DefinitionId] = referencedDefinition;
 77
 78        // Use the OutboundEdges lookup from the graph (Source → what it depends on)
 79        var dependencyMap = referenceGraph.OutboundEdges;
 80
 81        // Perform topological sort to ensure dependent workflows are processed in the right order
 82        var sortedWorkflowIds = referencingIds
 83            .TSort(id => dependencyMap[id], true)
 84            // Only process workflows that exist in our referencing workflows dictionary
 85            .Where(id => referencingWorkflowGraphs.ContainsKey(id))
 86            .ToList();
 87
 88        var updatedWorkflows = new Dictionary<string, UpdatedWorkflowDefinition>();
 89
 90        // Create a cache for drafts that we've already created during this operation
 91        var draftCache = new Dictionary<string, WorkflowDefinition>();
 92
 93        foreach (var id in sortedWorkflowIds)
 94        {
 95            if (!referencingWorkflowGraphs.TryGetValue(id, out var graph) || !dependencyMap[id].Any())
 96                continue;
 97
 98            foreach (var refId in dependencyMap[id])
 99            {
 100                var target = referencedWorkflowDefinitionsPublished.GetValueOrDefault(refId);
 101                if (target == null) continue;
 102
 103                var updated = await UpdateWorkflowAsync(graph, target, draftCache, initialPublicationState, cancellation
 104                if (updated == null) continue;
 105
 106                graph = updated.NewGraph;
 107                updatedWorkflows[updated.Definition.DefinitionId] = updated;
 108                referencedWorkflowDefinitionsPublished[id] = updated.Definition;
 109                draftCache[id] = updated.Definition;
 110                referencingWorkflowGraphs[id] = updated.NewGraph;
 111            }
 112        }
 113
 114        _isUpdating = true;
 115        foreach (var updatedWorkflow in updatedWorkflows.Values)
 116        {
 117            var requiresPublication = initialPublicationState.GetValueOrDefault(updatedWorkflow.Definition.DefinitionId)
 118            if (requiresPublication)
 119                await publisher.PublishAsync(updatedWorkflow.Definition, cancellationToken);
 120            else
 121                await publisher.SaveDraftAsync(updatedWorkflow.Definition, cancellationToken);
 122        }
 123
 124        _isUpdating = false;
 125
 126        return new(updatedWorkflows.Select(u => u.Value.Definition));
 127    }
 128
 129
 130    private async Task<UpdatedWorkflowDefinition?> UpdateWorkflowAsync(
 131        WorkflowGraph graph,
 132        WorkflowDefinition target,
 133        Dictionary<string, WorkflowDefinition> draftCache,
 134        Dictionary<string, bool> initialPublicationState,
 135        CancellationToken cancellationToken)
 136    {
 137        var willTargetBePublished = initialPublicationState.GetValueOrDefault(target.DefinitionId, target.IsPublished);
 138        if (!willTargetBePublished)
 139            return null;
 140
 141        var id = graph.Workflow.Identity.DefinitionId;
 142        var latest = await workflowDefinitionStore.FindAsync(new WorkflowDefinitionFilter
 143        {
 144            DefinitionId = id,
 145            VersionOptions = VersionOptions.Latest
 146        }, cancellationToken);
 147
 148        if (latest == null)
 149            return null;
 150
 151        // Decide skip before creating a draft. GetOrCreateDraftAsync clones published
 152        // consumers and registers that unsaved version as latest in the activity registry.
 153        if (latest.MaterializerName != JsonWorkflowMaterializer.MaterializerName)
 154        {
 155            var outdatedOnLatest = FindActivities(graph.Root, target.DefinitionId)
 156                .Any(a => a.WorkflowDefinitionVersionId != target.Id);
 157
 158            if (outdatedOnLatest)
 159            {
 160                logger.LogWarning(
 161                    "Skipping reference update for workflow definition {ConsumerDefinitionId}: it is authored in a non-J
 162                    latest.DefinitionId,
 163                    target.DefinitionId,
 164                    target.Version);
 165            }
 166
 167            return null;
 168        }
 169
 170        var draft = await GetOrCreateDraftAsync(id, draftCache, cancellationToken);
 171        if (draft == null) return null;
 172
 173        var newGraph = await workflowDefinitionService.MaterializeWorkflowAsync(draft, cancellationToken);
 174        var outdated = FindActivities(newGraph.Root, target.DefinitionId)
 175            .Where(a => a.WorkflowDefinitionVersionId != target.Id)
 176            .ToList();
 177
 178        if (!outdated.Any()) return null;
 179
 180        foreach (var act in outdated)
 181        {
 182            act.WorkflowDefinitionVersionId = target.Id;
 183            act.Version = target.Version;
 184            act.LatestAvailablePublishedVersionId = target.Id;
 185            act.LatestAvailablePublishedVersion = target.Version;
 186        }
 187
 188        if (newGraph.Root.Activity is Workflow wf)
 189        {
 190            draft.StringData = serializer.Serialize(wf.Root);
 191            draft.OriginalSource = null;
 192        }
 193
 194        return new(draft, newGraph);
 195    }
 196
 197    private async Task<WorkflowDefinition?> GetOrCreateDraftAsync(
 198        string definitionId,
 199        Dictionary<string, WorkflowDefinition> draftCache,
 200        CancellationToken cancellationToken)
 201    {
 202        // Check if we already have a draft for this workflow
 203        if (draftCache.TryGetValue(definitionId, out var cachedDraft))
 204            return cachedDraft;
 205
 206        // Create or get a draft for this workflow
 207        var draft = await publisher.GetDraftAsync(definitionId, VersionOptions.Latest, cancellationToken);
 208        if (draft == null) return null;
 209
 210        // Store the draft in the cache for potential future use
 211        draftCache[definitionId] = draft;
 212
 213        // Get the current published version of the workflow definition.
 214        var publishedVersion = await workflowDefinitionStore.FindAsync(
 215            WorkflowDefinitionHandle.ByDefinitionId(definitionId, VersionOptions.Published).ToFilter(),
 216            cancellationToken);
 217
 218        // Update the activity registry to be able to materialize the workflow.
 219        var activityDescriptor = workflowDefinitionActivityDescriptorFactory.CreateDescriptor(draft, publishedVersion);
 220        activityRegistry.Add(typeof(WorkflowDefinitionActivityProvider), activityDescriptor);
 221
 222        return draft;
 223    }
 224
 225    private static IEnumerable<WorkflowDefinitionActivity> FindActivities(ActivityNode node, string definitionId)
 226    {
 227        // Do not drill into activities that are WorkflowDefinitionActivity
 228        if (node.Activity is WorkflowDefinitionActivity)
 229            yield break;
 230
 231        foreach (var child in node.Children)
 232        {
 233            if (child.Activity is WorkflowDefinitionActivity activity && activity.WorkflowDefinitionId == definitionId)
 234                yield return activity;
 235            foreach (var grandChildActivity in FindActivities(child, definitionId))
 236                yield return grandChildActivity;
 237        }
 238    }
 239}

Methods/Properties

get_Definition()