| | | 1 | | using Elsa.Workflows.Management.Activities.WorkflowDefinitionActivity; |
| | | 2 | | using Elsa.Workflows.Management.Contracts; |
| | | 3 | | using Elsa.Caching; |
| | | 4 | | using Elsa.Common.Multitenancy; |
| | | 5 | | using Elsa.Workflows.Management.Stores; |
| | | 6 | | using Elsa.Workflows.Models; |
| | | 7 | | |
| | | 8 | | namespace Elsa.Workflows.Management.Services; |
| | | 9 | | |
| | | 10 | | /// <summary> |
| | | 11 | | /// Service responsible for updating the activity registry based on activity providers. |
| | | 12 | | /// </summary> |
| | 408 | 13 | | public class WorkflowDefinitionActivityRegistryUpdater( |
| | 408 | 14 | | WorkflowDefinitionActivityProvider provider, |
| | 408 | 15 | | IActivityRegistry registry, |
| | 408 | 16 | | ICacheManager? cacheManager, |
| | 408 | 17 | | ITenantAccessor? tenantAccessor) : IWorkflowDefinitionActivityRegistryUpdater, IWorkflowDefinitionActivityRegistryRe |
| | | 18 | | { |
| | 408 | 19 | | private readonly Type _providerType = typeof(WorkflowDefinitionActivityProvider); |
| | 2 | 20 | | 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 |
| | 0 | 27 | | : this(provider, registry, null, null) |
| | | 28 | | { |
| | 0 | 29 | | } |
| | | 30 | | |
| | | 31 | | /// <inheritdoc /> |
| | | 32 | | public async Task AddToRegistry(string workflowDefinitionVersionId, CancellationToken cancellationToken) |
| | | 33 | | { |
| | 131 | 34 | | var descriptors = await provider.GetDescriptorsAsync(cancellationToken); |
| | 131 | 35 | | var descriptorToAdd = descriptors |
| | 131 | 36 | | .FirstOrDefault(d => |
| | 947 | 37 | | d.CustomProperties.TryGetValue("WorkflowDefinitionVersionId", out var val) && |
| | 947 | 38 | | val.ToString() == workflowDefinitionVersionId); |
| | | 39 | | |
| | 131 | 40 | | if (descriptorToAdd is null) |
| | 0 | 41 | | return; |
| | | 42 | | |
| | 131 | 43 | | await RegistryLock.WaitAsync(cancellationToken); |
| | | 44 | | try |
| | | 45 | | { |
| | 131 | 46 | | registry.Add(_providerType, descriptorToAdd); |
| | 131 | 47 | | Interlocked.Increment(ref RegistryMutationVersion); |
| | 131 | 48 | | } |
| | | 49 | | finally |
| | | 50 | | { |
| | 131 | 51 | | RegistryLock.Release(); |
| | | 52 | | } |
| | 131 | 53 | | } |
| | | 54 | | |
| | | 55 | | /// <inheritdoc /> |
| | | 56 | | public async Task ReconcileRegistryAsync(CancellationToken cancellationToken = default) |
| | | 57 | | { |
| | 184 | 58 | | if (cacheManager is null || tenantAccessor is null) |
| | 0 | 59 | | throw new InvalidOperationException("Registry reconciliation requires cache and tenant services."); |
| | | 60 | | |
| | 184 | 61 | | 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. |
| | 189 | 66 | | 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. |
| | 189 | 70 | | var descriptors = (await provider.GetDescriptorsAsync(cancellationToken)).ToList(); |
| | | 71 | | |
| | 188 | 72 | | 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. |
| | 188 | 77 | | if (observedMutationVersion != RegistryMutationVersion) |
| | | 78 | | { |
| | 5 | 79 | | observedMutationVersion = RegistryMutationVersion; |
| | 5 | 80 | | 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. |
| | 2242 | 86 | | foreach (var descriptor in registry.ListByProvider(_providerType).ToList()) |
| | 938 | 87 | | registry.Remove(_providerType, descriptor); |
| | | 88 | | |
| | 2248 | 89 | | foreach (var descriptor in descriptors) |
| | 941 | 90 | | registry.Add(_providerType, descriptor); |
| | | 91 | | |
| | 183 | 92 | | Interlocked.Increment(ref RegistryMutationVersion); |
| | 183 | 93 | | return; |
| | | 94 | | } |
| | | 95 | | finally |
| | | 96 | | { |
| | 188 | 97 | | RegistryLock.Release(); |
| | | 98 | | } |
| | | 99 | | } |
| | 183 | 100 | | } |
| | | 101 | | |
| | | 102 | | /// <inheritdoc /> |
| | | 103 | | public void RemoveDefinitionFromRegistry(string workflowDefinitionId) |
| | | 104 | | { |
| | 9 | 105 | | RegistryLock.Wait(); |
| | | 106 | | try |
| | | 107 | | { |
| | 9 | 108 | | var descriptorsToRemove = registry.ListByProvider(_providerType) |
| | 103 | 109 | | .Where(d => d.CustomProperties.TryGetValue("WorkflowDefinitionId", out var val) && val.ToString() == wor |
| | 9 | 110 | | .ToList(); |
| | | 111 | | |
| | 34 | 112 | | foreach (var activityDescriptor in descriptorsToRemove) |
| | 8 | 113 | | registry.Remove(_providerType, activityDescriptor); |
| | | 114 | | |
| | 9 | 115 | | Interlocked.Increment(ref RegistryMutationVersion); |
| | 9 | 116 | | } |
| | | 117 | | finally |
| | | 118 | | { |
| | 9 | 119 | | RegistryLock.Release(); |
| | 9 | 120 | | } |
| | 9 | 121 | | } |
| | | 122 | | |
| | | 123 | | /// <inheritdoc /> |
| | | 124 | | public void RemoveDefinitionVersionFromRegistry(string workflowDefinitionVersionId) |
| | | 125 | | { |
| | 1263 | 126 | | RegistryLock.Wait(); |
| | | 127 | | try |
| | | 128 | | { |
| | 1263 | 129 | | var descriptorToRemove = registry.ListByProvider(_providerType) |
| | 11572 | 130 | | .FirstOrDefault(d => d.CustomProperties.TryGetValue("WorkflowDefinitionVersionId", out var val) && val.T |
| | | 131 | | |
| | 1263 | 132 | | if (descriptorToRemove is not null) |
| | | 133 | | { |
| | 2 | 134 | | registry.Remove(_providerType, descriptorToRemove); |
| | 2 | 135 | | Interlocked.Increment(ref RegistryMutationVersion); |
| | | 136 | | } |
| | 1263 | 137 | | } |
| | | 138 | | finally |
| | | 139 | | { |
| | 1263 | 140 | | RegistryLock.Release(); |
| | 1263 | 141 | | } |
| | 1263 | 142 | | } |
| | | 143 | | } |