| | | 1 | | using Elsa.Common.Multitenancy; |
| | | 2 | | using Elsa.Connections.Contracts; |
| | | 3 | | using Elsa.Connections.Models; |
| | | 4 | | using Microsoft.Extensions.DependencyInjection; |
| | | 5 | | using Microsoft.Extensions.Hosting; |
| | | 6 | | using Microsoft.Extensions.Logging; |
| | | 7 | | using Microsoft.Extensions.Options; |
| | | 8 | | |
| | | 9 | | namespace Elsa.Connections.Features; |
| | | 10 | | |
| | 10 | 11 | | public sealed class ConnectionLifecycleReconciliationWorker( |
| | 10 | 12 | | IServiceScopeFactory scopeFactory, |
| | 10 | 13 | | IOptions<ConnectionLifecycleReconciliationOptions> options, |
| | 10 | 14 | | TimeProvider timeProvider, |
| | 10 | 15 | | ILogger<ConnectionLifecycleReconciliationWorker> logger) : BackgroundService |
| | | 16 | | { |
| | | 17 | | private const int MaximumTrackedScopes = 1024; |
| | 10 | 18 | | private readonly Dictionary<(string TenantId, string EnvironmentId), ScopeCursorEntry> _cursors = new(); |
| | 10 | 19 | | private readonly LinkedList<(string TenantId, string EnvironmentId)> _scopeOrder = new(); |
| | | 20 | | |
| | | 21 | | protected override async Task ExecuteAsync(CancellationToken stoppingToken) |
| | | 22 | | { |
| | 1125 | 23 | | while (!stoppingToken.IsCancellationRequested) |
| | | 24 | | { |
| | 1125 | 25 | | if (options.Value.Enabled) |
| | | 26 | | { |
| | | 27 | | try |
| | | 28 | | { |
| | 1119 | 29 | | await ProcessNextScopeAsync(stoppingToken); |
| | 1119 | 30 | | } |
| | 0 | 31 | | catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) |
| | | 32 | | { |
| | 0 | 33 | | break; |
| | | 34 | | } |
| | 0 | 35 | | catch (Exception) |
| | | 36 | | { |
| | | 37 | | // Exceptions can contain provider payloads. Keep logs free of exception text and candidate data. |
| | 0 | 38 | | ConnectionLifecycleReconciliationMetrics.RecordFailedScan(); |
| | 0 | 39 | | logger.LogWarning("Credential lifecycle reconciliation cycle failed."); |
| | 0 | 40 | | } |
| | | 41 | | } |
| | | 42 | | |
| | | 43 | | try |
| | | 44 | | { |
| | | 45 | | // A failing scope must not lengthen the scan cadence for other scopes. |
| | 1125 | 46 | | await Task.Delay(options.Value.Interval, timeProvider, stoppingToken); |
| | 1115 | 47 | | } |
| | 10 | 48 | | catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) |
| | | 49 | | { |
| | 10 | 50 | | break; |
| | | 51 | | } |
| | | 52 | | } |
| | 10 | 53 | | } |
| | | 54 | | |
| | | 55 | | private async Task ProcessNextScopeAsync(CancellationToken cancellationToken) |
| | | 56 | | { |
| | 1119 | 57 | | await using var scope = scopeFactory.CreateAsyncScope(); |
| | 1119 | 58 | | var scopeProvider = scope.ServiceProvider.GetRequiredService<IConnectionLifecycleScopeProvider>(); |
| | 1119 | 59 | | var target = await scopeProvider.GetNextScopeAsync(cancellationToken); |
| | 1119 | 60 | | if (target == null) |
| | | 61 | | return; |
| | | 62 | | |
| | 1119 | 63 | | if (string.IsNullOrWhiteSpace(target.TenantId) || string.IsNullOrWhiteSpace(target.EnvironmentId)) |
| | | 64 | | { |
| | 1 | 65 | | ConnectionLifecycleReconciliationMetrics.RecordFailedScan(); |
| | 1 | 66 | | logger.LogWarning("Credential lifecycle scope provider returned an invalid scope."); |
| | 1 | 67 | | return; |
| | | 68 | | } |
| | | 69 | | |
| | 1118 | 70 | | var scopeKey = (target.TenantId, target.EnvironmentId); |
| | 1118 | 71 | | var cached = _cursors.GetValueOrDefault(scopeKey); |
| | 1118 | 72 | | if (cached?.FailureTimestamp is long failureTimestamp && |
| | 1118 | 73 | | timeProvider.GetElapsedTime(failureTimestamp) < cached.RetryDelay) |
| | | 74 | | return; |
| | | 75 | | |
| | 1112 | 76 | | var cursor = cached?.Cursor; |
| | 1112 | 77 | | var tenantAccessor = scope.ServiceProvider.GetRequiredService<ITenantAccessor>(); |
| | 1112 | 78 | | using var tenantContext = tenantAccessor.PushContext(new Tenant { Id = target.TenantId, Name = target.TenantId } |
| | 1112 | 79 | | var candidates = scope.ServiceProvider.GetRequiredService<IConnectionDueCandidateStore>(); |
| | | 80 | | try |
| | | 81 | | { |
| | 1112 | 82 | | var page = await candidates.FindDueCandidatesAsync( |
| | 1112 | 83 | | target.TenantId, target.EnvironmentId, timeProvider.GetUtcNow(), options.Value.BatchSize, cursor, cancel |
| | | 84 | | |
| | 1112 | 85 | | if (page.Items.Count == 0) |
| | | 86 | | { |
| | 552 | 87 | | ConnectionLifecycleReconciliationMetrics.RecordCompletedPage(); |
| | 552 | 88 | | RememberScope(scopeKey, page.NextCursor, 0, null, TimeSpan.Zero); |
| | 552 | 89 | | return; |
| | | 90 | | } |
| | | 91 | | |
| | 560 | 92 | | var failures = 0; |
| | 560 | 93 | | await Parallel.ForEachAsync(page.Items, |
| | 560 | 94 | | new ParallelOptions { MaxDegreeOfParallelism = options.Value.MaxConcurrency, CancellationToken = cancell |
| | 560 | 95 | | async (candidate, token) => |
| | 560 | 96 | | { |
| | 569 | 97 | | if (candidate.TenantId != target.TenantId || candidate.EnvironmentId != target.EnvironmentId || |
| | 569 | 98 | | string.IsNullOrWhiteSpace(candidate.ConnectionId) || candidate.DueAt > timeProvider.GetUtcNow()) |
| | 560 | 99 | | { |
| | 2 | 100 | | Interlocked.Increment(ref failures); |
| | 2 | 101 | | return; |
| | 560 | 102 | | } |
| | 560 | 103 | | |
| | 560 | 104 | | try |
| | 560 | 105 | | { |
| | 567 | 106 | | if (!await DispatchAsync(candidate, token)) |
| | 4 | 107 | | Interlocked.Increment(ref failures); |
| | 567 | 108 | | } |
| | 0 | 109 | | catch (OperationCanceledException) when (token.IsCancellationRequested) |
| | 560 | 110 | | { |
| | 0 | 111 | | throw; |
| | 560 | 112 | | } |
| | 0 | 113 | | catch (Exception) |
| | 560 | 114 | | { |
| | 560 | 115 | | // Keep exceptions and candidate identifiers out of logs; provider errors may contain credential |
| | 0 | 116 | | Interlocked.Increment(ref failures); |
| | 0 | 117 | | } |
| | 1129 | 118 | | }); |
| | | 119 | | |
| | 560 | 120 | | if (failures > 0) |
| | | 121 | | { |
| | 4 | 122 | | ConnectionLifecycleReconciliationMetrics.RecordFailedPage(); |
| | 4 | 123 | | logger.LogWarning("Credential lifecycle reconciliation page had {FailureCount} failed or out-of-scope ca |
| | 4 | 124 | | RememberFailure(scopeKey, cursor); |
| | 4 | 125 | | return; |
| | | 126 | | } |
| | | 127 | | |
| | | 128 | | // Commit the next cursor only after every candidate in the page succeeds. |
| | 556 | 129 | | ConnectionLifecycleReconciliationMetrics.RecordCompletedPage(); |
| | 556 | 130 | | RememberScope(scopeKey, page.NextCursor, 0, null, TimeSpan.Zero); |
| | 556 | 131 | | return; |
| | | 132 | | } |
| | 0 | 133 | | catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) |
| | | 134 | | { |
| | 0 | 135 | | throw; |
| | | 136 | | } |
| | 0 | 137 | | catch (Exception) |
| | | 138 | | { |
| | 0 | 139 | | ConnectionLifecycleReconciliationMetrics.RecordFailedScan(); |
| | 0 | 140 | | logger.LogWarning("Credential lifecycle reconciliation scope failed."); |
| | 0 | 141 | | RememberFailure(scopeKey, cursor); |
| | 0 | 142 | | } |
| | 1119 | 143 | | } |
| | | 144 | | |
| | | 145 | | private async Task<bool> DispatchAsync(ConnectionDueCandidate candidate, CancellationToken cancellationToken) |
| | | 146 | | { |
| | 567 | 147 | | await using var scope = scopeFactory.CreateAsyncScope(); |
| | 567 | 148 | | var tenantAccessor = scope.ServiceProvider.GetRequiredService<ITenantAccessor>(); |
| | 567 | 149 | | using var tenantContext = tenantAccessor.PushContext(new Tenant { Id = candidate.TenantId, Name = candidate.Tena |
| | 567 | 150 | | var lifecycle = scope.ServiceProvider.GetRequiredService<IConnectionLifecycleRecoveryService>(); |
| | | 151 | | |
| | 567 | 152 | | switch (candidate.Kind) |
| | | 153 | | { |
| | | 154 | | case ConnectionDueCandidateKind.OAuthRefresh: |
| | 2 | 155 | | return (await lifecycle.RefreshAsync(candidate.TenantId, candidate.EnvironmentId, candidate.ConnectionId |
| | | 156 | | case ConnectionDueCandidateKind.ExpiredConnectionOperation: |
| | | 157 | | case ConnectionDueCandidateKind.RecoveryRequired: |
| | 559 | 158 | | var reconciliation = await lifecycle.ReconcileAsync(candidate.TenantId, candidate.EnvironmentId, |
| | 559 | 159 | | candidate.ConnectionId, cancellationToken); |
| | | 160 | | // Report both an existing RecoveryRequired state and transitions that may enter it. |
| | | 161 | | // Only the established recovery_required result is handled for this scan. |
| | 559 | 162 | | if (reconciliation.SafeErrorCode is "recovery_required" or "refresh_outcome_unknown" or "generation_publ |
| | 3 | 163 | | ConnectionLifecycleReconciliationMetrics.RecordRecoveryRequired(); |
| | 559 | 164 | | return reconciliation.Succeeded || reconciliation.SafeErrorCode == "recovery_required"; |
| | | 165 | | case ConnectionDueCandidateKind.GenerationCleanup: |
| | 1 | 166 | | return (await lifecycle.CleanupGenerationAsync(candidate.TenantId, candidate.EnvironmentId, |
| | 1 | 167 | | candidate.ConnectionId, candidate.CandidateId, cancellationToken)).Succeeded; |
| | | 168 | | case ConnectionDueCandidateKind.Offboarding: |
| | 5 | 169 | | var offboarding = await lifecycle.ReconcileOffboardingAsync(candidate.TenantId, candidate.EnvironmentId, |
| | 5 | 170 | | candidate.ConnectionId, cancellationToken); |
| | | 171 | | // Report UnknownOutcome even when the service accepts the reconciliation with no error code. |
| | | 172 | | // A non-replayable unknown is handled for this scan so later candidates are not held behind it. |
| | 5 | 173 | | if (offboarding.SafeErrorCode == "offboarding_outcome_unknown" || |
| | 5 | 174 | | offboarding.Status == ConnectionOffboardingOperationStatus.UnknownOutcome) |
| | 3 | 175 | | ConnectionLifecycleReconciliationMetrics.RecordUnknownOffboarding(); |
| | 5 | 176 | | return offboarding.Accepted || offboarding.SafeErrorCode == "offboarding_outcome_unknown"; |
| | | 177 | | default: |
| | 0 | 178 | | throw new InvalidOperationException("Unsupported credential lifecycle candidate kind."); |
| | | 179 | | } |
| | 567 | 180 | | } |
| | | 181 | | |
| | | 182 | | private void RememberFailure((string TenantId, string EnvironmentId) scopeKey, string? cursor) |
| | | 183 | | { |
| | 4 | 184 | | var failures = _cursors.TryGetValue(scopeKey, out var existing) |
| | 4 | 185 | | ? Math.Min(existing.ConsecutiveFailures + 1, 30) |
| | 4 | 186 | | : 1; |
| | 4 | 187 | | RememberScope(scopeKey, cursor, failures, timeProvider.GetTimestamp(), GetRetryDelay(failures)); |
| | 4 | 188 | | } |
| | | 189 | | |
| | | 190 | | private void RememberScope((string TenantId, string EnvironmentId) scopeKey, string? cursor, int consecutiveFailures |
| | | 191 | | long? failureTimestamp, TimeSpan retryDelay) |
| | | 192 | | { |
| | 1112 | 193 | | if (_cursors.TryGetValue(scopeKey, out var existing)) |
| | | 194 | | { |
| | 1103 | 195 | | _scopeOrder.Remove(existing.Node); |
| | | 196 | | } |
| | | 197 | | else |
| | | 198 | | { |
| | 9 | 199 | | if (_cursors.Count >= MaximumTrackedScopes && _scopeOrder.First is { } oldest) |
| | | 200 | | { |
| | 0 | 201 | | _scopeOrder.RemoveFirst(); |
| | 0 | 202 | | _cursors.Remove(oldest.Value); |
| | | 203 | | } |
| | | 204 | | } |
| | | 205 | | |
| | 1112 | 206 | | var node = _scopeOrder.AddLast(scopeKey); |
| | 1112 | 207 | | _cursors[scopeKey] = new ScopeCursorEntry(cursor, node, consecutiveFailures, failureTimestamp, retryDelay); |
| | 1112 | 208 | | } |
| | | 209 | | |
| | | 210 | | private TimeSpan GetRetryDelay(int consecutiveFailures) |
| | | 211 | | { |
| | 4 | 212 | | var configured = options.Value.RetryBackoff.TotalMilliseconds; |
| | 4 | 213 | | var multiplier = Math.Pow(2, Math.Min(consecutiveFailures - 1, 20)); |
| | 4 | 214 | | return TimeSpan.FromMilliseconds(Math.Min(configured * multiplier, options.Value.MaxRetryBackoff.TotalMillisecon |
| | | 215 | | } |
| | | 216 | | |
| | 3318 | 217 | | private sealed record ScopeCursorEntry(string? Cursor, LinkedListNode<(string TenantId, string EnvironmentId)> Node, |
| | 2229 | 218 | | int ConsecutiveFailures, long? FailureTimestamp, TimeSpan RetryDelay); |
| | | 219 | | } |