< Summary

Information
Class: Elsa.Workflows.Management.Services.WorkflowReferenceUpdater
Assembly: Elsa.Workflows.Management
File(s): /home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Workflows.Management/Services/WorkflowReferenceUpdater.cs
Line coverage
97%
Covered lines: 128
Uncovered lines: 3
Coverable lines: 131
Total lines: 239
Line coverage: 97.7%
Branch coverage
87%
Covered branches: 54
Total branches: 62
Branch coverage: 87%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)100%11100%
UpdateWorkflowReferencesAsync()96.66%303098.11%
UpdateWorkflowAsync()75%161694.44%
GetOrCreateDraftAsync()50%4491.66%
FindActivities()100%1010100%

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
 14internal record UpdatedWorkflowDefinition(WorkflowDefinition Definition, WorkflowGraph NewGraph);
 15
 40716public class WorkflowReferenceUpdater(
 40717    IWorkflowDefinitionPublisher publisher,
 40718    IWorkflowDefinitionService workflowDefinitionService,
 40719    IWorkflowDefinitionStore workflowDefinitionStore,
 40720    IWorkflowReferenceGraphBuilder workflowReferenceGraphBuilder,
 40721    WorkflowDefinitionActivityDescriptorFactory workflowDefinitionActivityDescriptorFactory,
 40722    IActivityRegistry activityRegistry,
 40723    IApiSerializer serializer,
 40724    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    {
 1033        if (_isUpdating ||
 1034            referencedDefinition.Options is not { UsableAsActivity: true, AutoUpdateConsumingWorkflows: true })
 335            return new([]);
 36
 737        var referenceGraph = await workflowReferenceGraphBuilder.BuildGraphAsync(referencedDefinition.DefinitionId, canc
 38
 39        // Get all consumer (source) and referenced (target) IDs from the edges
 1140        var referencingIds = referenceGraph.Edges.Select(e => e.Source).Distinct().ToList();
 1141        var referencedIds = referenceGraph.Edges.Select(e => e.Target).Distinct().ToList();
 42
 743        if (referencingIds.Count == 0)
 444            return new([]);
 45
 346        var referencingWorkflowGraphs = (await workflowDefinitionService.FindWorkflowGraphsAsync(new()
 347            {
 348                DefinitionIds = referencingIds,
 349                VersionOptions = VersionOptions.Latest,
 350                IsReadonly = false
 351            }, cancellationToken))
 752            .ToDictionary(g => g.Workflow.Identity.DefinitionId);
 53
 354        var referencedWorkflowDefinitionList = (await workflowDefinitionStore.FindManyAsync(new()
 355        {
 356            DefinitionIds = referencedIds,
 357            VersionOptions = VersionOptions.Published,
 358            IsReadonly = false
 359        }, cancellationToken)).ToList();
 60
 361        var referencedWorkflowDefinitionsPublished = referencedWorkflowDefinitionList
 362            .GroupBy(x => x.DefinitionId)
 363            .Select(group =>
 364            {
 665                var publishedVersion = group.FirstOrDefault(x => x.IsPublished);
 366                return publishedVersion ?? group.First();
 367            })
 668            .ToDictionary(d => d.DefinitionId);
 69
 370        var initialPublicationState = new Dictionary<string, bool>();
 71
 1472        foreach (var workflowGraph in referencingWorkflowGraphs)
 473            initialPublicationState[workflowGraph.Key] = workflowGraph.Value.Workflow.Publication.IsPublished;
 74
 75        // Add the initially referenced definition
 376        referencedWorkflowDefinitionsPublished[referencedDefinition.DefinitionId] = referencedDefinition;
 77
 78        // Use the OutboundEdges lookup from the graph (Source → what it depends on)
 379        var dependencyMap = referenceGraph.OutboundEdges;
 80
 81        // Perform topological sort to ensure dependent workflows are processed in the right order
 382        var sortedWorkflowIds = referencingIds
 783            .TSort(id => dependencyMap[id], true)
 384            // Only process workflows that exist in our referencing workflows dictionary
 785            .Where(id => referencingWorkflowGraphs.ContainsKey(id))
 386            .ToList();
 87
 388        var updatedWorkflows = new Dictionary<string, UpdatedWorkflowDefinition>();
 89
 90        // Create a cache for drafts that we've already created during this operation
 391        var draftCache = new Dictionary<string, WorkflowDefinition>();
 92
 1493        foreach (var id in sortedWorkflowIds)
 94        {
 495            if (!referencingWorkflowGraphs.TryGetValue(id, out var graph) || !dependencyMap[id].Any())
 96                continue;
 97
 1698            foreach (var refId in dependencyMap[id])
 99            {
 4100                var target = referencedWorkflowDefinitionsPublished.GetValueOrDefault(refId);
 4101                if (target == null) continue;
 102
 4103                var updated = await UpdateWorkflowAsync(graph, target, draftCache, initialPublicationState, cancellation
 4104                if (updated == null) continue;
 105
 2106                graph = updated.NewGraph;
 2107                updatedWorkflows[updated.Definition.DefinitionId] = updated;
 2108                referencedWorkflowDefinitionsPublished[id] = updated.Definition;
 2109                draftCache[id] = updated.Definition;
 2110                referencingWorkflowGraphs[id] = updated.NewGraph;
 111            }
 4112        }
 113
 3114        _isUpdating = true;
 10115        foreach (var updatedWorkflow in updatedWorkflows.Values)
 116        {
 2117            var requiresPublication = initialPublicationState.GetValueOrDefault(updatedWorkflow.Definition.DefinitionId)
 2118            if (requiresPublication)
 1119                await publisher.PublishAsync(updatedWorkflow.Definition, cancellationToken);
 120            else
 1121                await publisher.SaveDraftAsync(updatedWorkflow.Definition, cancellationToken);
 122        }
 123
 3124        _isUpdating = false;
 125
 5126        return new(updatedWorkflows.Select(u => u.Value.Definition));
 10127    }
 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    {
 4137        var willTargetBePublished = initialPublicationState.GetValueOrDefault(target.DefinitionId, target.IsPublished);
 4138        if (!willTargetBePublished)
 0139            return null;
 140
 4141        var id = graph.Workflow.Identity.DefinitionId;
 4142        var latest = await workflowDefinitionStore.FindAsync(new WorkflowDefinitionFilter
 4143        {
 4144            DefinitionId = id,
 4145            VersionOptions = VersionOptions.Latest
 4146        }, cancellationToken);
 147
 4148        if (latest == null)
 0149            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.
 4153        if (latest.MaterializerName != JsonWorkflowMaterializer.MaterializerName)
 154        {
 2155            var outdatedOnLatest = FindActivities(graph.Root, target.DefinitionId)
 4156                .Any(a => a.WorkflowDefinitionVersionId != target.Id);
 157
 2158            if (outdatedOnLatest)
 159            {
 2160                logger.LogWarning(
 2161                    "Skipping reference update for workflow definition {ConsumerDefinitionId}: it is authored in a non-J
 2162                    latest.DefinitionId,
 2163                    target.DefinitionId,
 2164                    target.Version);
 165            }
 166
 2167            return null;
 168        }
 169
 2170        var draft = await GetOrCreateDraftAsync(id, draftCache, cancellationToken);
 2171        if (draft == null) return null;
 172
 2173        var newGraph = await workflowDefinitionService.MaterializeWorkflowAsync(draft, cancellationToken);
 2174        var outdated = FindActivities(newGraph.Root, target.DefinitionId)
 2175            .Where(a => a.WorkflowDefinitionVersionId != target.Id)
 2176            .ToList();
 177
 2178        if (!outdated.Any()) return null;
 179
 8180        foreach (var act in outdated)
 181        {
 2182            act.WorkflowDefinitionVersionId = target.Id;
 2183            act.Version = target.Version;
 2184            act.LatestAvailablePublishedVersionId = target.Id;
 2185            act.LatestAvailablePublishedVersion = target.Version;
 186        }
 187
 2188        if (newGraph.Root.Activity is Workflow wf)
 189        {
 2190            draft.StringData = serializer.Serialize(wf.Root);
 2191            draft.OriginalSource = null;
 192        }
 193
 2194        return new(draft, newGraph);
 4195    }
 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
 2203        if (draftCache.TryGetValue(definitionId, out var cachedDraft))
 0204            return cachedDraft;
 205
 206        // Create or get a draft for this workflow
 2207        var draft = await publisher.GetDraftAsync(definitionId, VersionOptions.Latest, cancellationToken);
 2208        if (draft == null) return null;
 209
 210        // Store the draft in the cache for potential future use
 2211        draftCache[definitionId] = draft;
 212
 213        // Get the current published version of the workflow definition.
 2214        var publishedVersion = await workflowDefinitionStore.FindAsync(
 2215            WorkflowDefinitionHandle.ByDefinitionId(definitionId, VersionOptions.Published).ToFilter(),
 2216            cancellationToken);
 217
 218        // Update the activity registry to be able to materialize the workflow.
 2219        var activityDescriptor = workflowDefinitionActivityDescriptorFactory.CreateDescriptor(draft, publishedVersion);
 2220        activityRegistry.Add(typeof(WorkflowDefinitionActivityProvider), activityDescriptor);
 221
 2222        return draft;
 2223    }
 224
 225    private static IEnumerable<WorkflowDefinitionActivity> FindActivities(ActivityNode node, string definitionId)
 226    {
 227        // Do not drill into activities that are WorkflowDefinitionActivity
 9228        if (node.Activity is WorkflowDefinitionActivity)
 2229            yield break;
 230
 26231        foreach (var child in node.Children)
 232        {
 7233            if (child.Activity is WorkflowDefinitionActivity activity && activity.WorkflowDefinitionId == definitionId)
 4234                yield return activity;
 12235            foreach (var grandChildActivity in FindActivities(child, definitionId))
 1236                yield return grandChildActivity;
 5237        }
 5238    }
 239}