| | | 1 | | using Elsa.Common; |
| | | 2 | | using Elsa.Common.Multitenancy; |
| | | 3 | | using Elsa.Common.RecurringTasks; |
| | | 4 | | using Elsa.Workflows.Management.Contracts; |
| | | 5 | | using JetBrains.Annotations; |
| | | 6 | | |
| | | 7 | | namespace Elsa.Workflows.Runtime.Tasks; |
| | | 8 | | |
| | | 9 | | /// <summary> |
| | | 10 | | /// Reconciles the local workflow-as-activity registry with shared definition changes. |
| | | 11 | | /// </summary> |
| | | 12 | | [UsedImplicitly] |
| | 100 | 13 | | public class RefreshWorkflowDefinitionActivityRegistryTask( |
| | 100 | 14 | | IWorkflowDefinitionRegistryGenerationStore generationStore, |
| | 100 | 15 | | IWorkflowDefinitionActivityRegistryReconciler registryUpdater, |
| | 100 | 16 | | ITenantAccessor tenantAccessor, |
| | 100 | 17 | | TimeProvider timeProvider) : RecurringTask |
| | | 18 | | { |
| | 2 | 19 | | private static readonly TimeSpan FullReconciliationInterval = TimeSpan.FromMinutes(1); |
| | | 20 | | private long? _tenantGeneration; |
| | | 21 | | private long? _agnosticGeneration; |
| | | 22 | | private long? _lastFullReconciliationTimestamp; |
| | | 23 | | |
| | | 24 | | /// <inheritdoc /> |
| | | 25 | | public override async Task ExecuteAsync(CancellationToken cancellationToken = default) |
| | | 26 | | { |
| | | 27 | | // The default in-memory store is node-local, so polling it cannot discover remote writes. |
| | 941 | 28 | | if (!generationStore.IsShared) |
| | 159 | 29 | | return; |
| | | 30 | | |
| | 782 | 31 | | var tenantId = tenantAccessor.TenantId ?? Tenant.DefaultTenantId; |
| | 782 | 32 | | var tenantGeneration = await generationStore.GetGenerationAsync(tenantId, cancellationToken); |
| | 782 | 33 | | var agnosticGeneration = await generationStore.GetGenerationAsync(Tenant.AgnosticTenantId, cancellationToken); |
| | 782 | 34 | | var timestamp = timeProvider.GetTimestamp(); |
| | 782 | 35 | | var fullReconciliationDue = _lastFullReconciliationTimestamp == null || timeProvider.GetElapsedTime(_lastFullRec |
| | | 36 | | |
| | 782 | 37 | | if (!fullReconciliationDue && tenantGeneration == _tenantGeneration && agnosticGeneration == _agnosticGeneration |
| | 597 | 38 | | return; |
| | | 39 | | |
| | | 40 | | // Keep the observed values from before the refresh. A concurrent write is therefore seen on the next poll. |
| | 185 | 41 | | await registryUpdater.ReconcileRegistryAsync(cancellationToken); |
| | 185 | 42 | | _tenantGeneration = tenantGeneration; |
| | 185 | 43 | | _agnosticGeneration = agnosticGeneration; |
| | 185 | 44 | | _lastFullReconciliationTimestamp = timestamp; |
| | 941 | 45 | | } |
| | | 46 | | } |