< Summary

Information
Class: Elsa.Workflows.Management.Services.WorkflowDefinitionActivityRegistryUpdater
Assembly: Elsa.Workflows.Management
File(s): /home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Workflows.Management/Services/WorkflowDefinitionActivityRegistryUpdater.cs
Line coverage
93%
Covered lines: 56
Uncovered lines: 4
Coverable lines: 60
Total lines: 143
Line coverage: 93.3%
Branch coverage
72%
Covered branches: 16
Total branches: 22
Branch coverage: 72.7%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)100%11100%
.cctor()100%11100%
.ctor(...)100%210%
AddToRegistry()50%2290.9%
ReconcileRegistryAsync()80%101094.11%
RemoveDefinitionFromRegistry(...)75%44100%
RemoveDefinitionVersionFromRegistry(...)75%44100%

File(s)

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

#LineLine coverage
 1using Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity;
 2using Elsa.Workflows.Management.Contracts;
 3using Elsa.Caching;
 4using Elsa.Common.Multitenancy;
 5using Elsa.Workflows.Management.Stores;
 6using Elsa.Workflows.Models;
 7
 8namespace Elsa.Workflows.Management.Services;
 9
 10/// <summary>
 11/// Service responsible for updating the activity registry based on activity providers.
 12/// </summary>
 40813public class WorkflowDefinitionActivityRegistryUpdater(
 40814    WorkflowDefinitionActivityProvider provider,
 40815    IActivityRegistry registry,
 40816    ICacheManager? cacheManager,
 40817    ITenantAccessor? tenantAccessor) : IWorkflowDefinitionActivityRegistryUpdater, IWorkflowDefinitionActivityRegistryRe
 18{
 40819    private readonly Type _providerType = typeof(WorkflowDefinitionActivityProvider);
 220    private static readonly SemaphoreSlim RegistryLock = new(1, 1);
 21    private static long RegistryMutationVersion;
 22
 23    /// <summary>
 24    /// Preserves the existing constructor for hosts that only use local registry updates.
 25    /// </summary>
 26    public WorkflowDefinitionActivityRegistryUpdater(WorkflowDefinitionActivityProvider provider, IActivityRegistry regi
 027        : this(provider, registry, null, null)
 28    {
 029    }
 30
 31    /// <inheritdoc />
 32    public async Task AddToRegistry(string workflowDefinitionVersionId, CancellationToken cancellationToken)
 33    {
 13134        var descriptors = await provider.GetDescriptorsAsync(cancellationToken);
 13135        var descriptorToAdd = descriptors
 13136            .FirstOrDefault(d =>
 94737                d.CustomProperties.TryGetValue("WorkflowDefinitionVersionId", out var val) &&
 94738                val.ToString() == workflowDefinitionVersionId);
 39
 13140        if (descriptorToAdd is null)
 041            return;
 42
 13143        await RegistryLock.WaitAsync(cancellationToken);
 44        try
 45        {
 13146            registry.Add(_providerType, descriptorToAdd);
 13147            Interlocked.Increment(ref RegistryMutationVersion);
 13148        }
 49        finally
 50        {
 13151            RegistryLock.Release();
 52        }
 13153    }
 54
 55    /// <inheritdoc />
 56    public async Task ReconcileRegistryAsync(CancellationToken cancellationToken = default)
 57    {
 18458        if (cacheManager is null || tenantAccessor is null)
 059            throw new InvalidOperationException("Registry reconciliation requires cache and tenant services.");
 60
 18461        var observedMutationVersion = Interlocked.Read(ref RegistryMutationVersion);
 62        while (true)
 63        {
 64            // Do not hold the process-wide registry mutation lock across cache or store I/O.
 65            // A cache warmed on this node before a remote write would otherwise hide the new store state.
 18966            await cacheManager.TriggerTokenAsync(CachingWorkflowDefinitionStore.GetTenantReconciliationTokenKey(tenantAc
 67
 68            // Read the authoritative set before mutating the live registry. A failed or cancelled
 69            // store read must leave the currently usable descriptors in place.
 18970            var descriptors = (await provider.GetDescriptorsAsync(cancellationToken)).ToList();
 71
 18872            await RegistryLock.WaitAsync(cancellationToken);
 73            try
 74            {
 75                // A local delete or another reconciliation may have changed the registry after the
 76                // provider snapshot. Retry rather than restoring a descriptor from that stale snapshot.
 18877                if (observedMutationVersion != RegistryMutationVersion)
 78                {
 579                    observedMutationVersion = RegistryMutationVersion;
 580                    continue;
 81                }
 82
 83                // ListByProvider is tenant-aware: it exposes only the current tenant plus agnostic descriptors.
 84                // Removing that visible set first also handles an empty provider result, which the generic
 85                // ActivityRegistry.RefreshDescriptorsAsync currently does not clear.
 224286                foreach (var descriptor in registry.ListByProvider(_providerType).ToList())
 93887                    registry.Remove(_providerType, descriptor);
 88
 224889                foreach (var descriptor in descriptors)
 94190                    registry.Add(_providerType, descriptor);
 91
 18392                Interlocked.Increment(ref RegistryMutationVersion);
 18393                return;
 94            }
 95            finally
 96            {
 18897                RegistryLock.Release();
 98            }
 99        }
 183100    }
 101
 102    /// <inheritdoc />
 103    public void RemoveDefinitionFromRegistry(string workflowDefinitionId)
 104    {
 9105        RegistryLock.Wait();
 106        try
 107        {
 9108            var descriptorsToRemove = registry.ListByProvider(_providerType)
 103109                .Where(d => d.CustomProperties.TryGetValue("WorkflowDefinitionId", out var val) && val.ToString() == wor
 9110                .ToList();
 111
 34112            foreach (var activityDescriptor in descriptorsToRemove)
 8113                registry.Remove(_providerType, activityDescriptor);
 114
 9115            Interlocked.Increment(ref RegistryMutationVersion);
 9116        }
 117        finally
 118        {
 9119            RegistryLock.Release();
 9120        }
 9121    }
 122
 123    /// <inheritdoc />
 124    public void RemoveDefinitionVersionFromRegistry(string workflowDefinitionVersionId)
 125    {
 1263126        RegistryLock.Wait();
 127        try
 128        {
 1263129            var descriptorToRemove = registry.ListByProvider(_providerType)
 11572130                .FirstOrDefault(d => d.CustomProperties.TryGetValue("WorkflowDefinitionVersionId", out var val) && val.T
 131
 1263132            if (descriptorToRemove is not null)
 133            {
 2134                registry.Remove(_providerType, descriptorToRemove);
 2135                Interlocked.Increment(ref RegistryMutationVersion);
 136            }
 1263137        }
 138        finally
 139        {
 1263140            RegistryLock.Release();
 1263141        }
 1263142    }
 143}