< Summary

Information
Class: Elsa.Connections.Features.ConnectionLifecycleReconciliationWorker
Assembly: Elsa.Connections
File(s): /home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Connections/Features/ConnectionLifecycleReconciliationWorker.cs
Line coverage
83%
Covered lines: 105
Uncovered lines: 21
Coverable lines: 126
Total lines: 219
Line coverage: 83.3%
Branch coverage
93%
Covered branches: 62
Total branches: 66
Branch coverage: 93.9%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)100%11100%
ExecuteAsync()100%5460%
ProcessNextScopeAsync()100%232286.53%
<ProcessNextScopeAsync()100%161061.53%
DispatchAsync()95.45%222295.23%
RememberFailure(...)50%22100%
RememberScope(...)66.66%7675%
GetRetryDelay(...)100%11100%
get_Cursor()100%11100%
get_ConsecutiveFailures()100%11100%

File(s)

/home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Connections/Features/ConnectionLifecycleReconciliationWorker.cs

#LineLine coverage
 1using Elsa.Common.Multitenancy;
 2using Elsa.Connections.Contracts;
 3using Elsa.Connections.Models;
 4using Microsoft.Extensions.DependencyInjection;
 5using Microsoft.Extensions.Hosting;
 6using Microsoft.Extensions.Logging;
 7using Microsoft.Extensions.Options;
 8
 9namespace Elsa.Connections.Features;
 10
 1011public sealed class ConnectionLifecycleReconciliationWorker(
 1012    IServiceScopeFactory scopeFactory,
 1013    IOptions<ConnectionLifecycleReconciliationOptions> options,
 1014    TimeProvider timeProvider,
 1015    ILogger<ConnectionLifecycleReconciliationWorker> logger) : BackgroundService
 16{
 17    private const int MaximumTrackedScopes = 1024;
 1018    private readonly Dictionary<(string TenantId, string EnvironmentId), ScopeCursorEntry> _cursors = new();
 1019    private readonly LinkedList<(string TenantId, string EnvironmentId)> _scopeOrder = new();
 20
 21    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
 22    {
 112523        while (!stoppingToken.IsCancellationRequested)
 24        {
 112525            if (options.Value.Enabled)
 26            {
 27                try
 28                {
 111929                    await ProcessNextScopeAsync(stoppingToken);
 111930                }
 031                catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
 32                {
 033                    break;
 34                }
 035                catch (Exception)
 36                {
 37                    // Exceptions can contain provider payloads. Keep logs free of exception text and candidate data.
 038                    ConnectionLifecycleReconciliationMetrics.RecordFailedScan();
 039                    logger.LogWarning("Credential lifecycle reconciliation cycle failed.");
 040                }
 41            }
 42
 43            try
 44            {
 45                // A failing scope must not lengthen the scan cadence for other scopes.
 112546                await Task.Delay(options.Value.Interval, timeProvider, stoppingToken);
 111547            }
 1048            catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
 49            {
 1050                break;
 51            }
 52        }
 1053    }
 54
 55    private async Task ProcessNextScopeAsync(CancellationToken cancellationToken)
 56    {
 111957        await using var scope = scopeFactory.CreateAsyncScope();
 111958        var scopeProvider = scope.ServiceProvider.GetRequiredService<IConnectionLifecycleScopeProvider>();
 111959        var target = await scopeProvider.GetNextScopeAsync(cancellationToken);
 111960        if (target == null)
 61            return;
 62
 111963        if (string.IsNullOrWhiteSpace(target.TenantId) || string.IsNullOrWhiteSpace(target.EnvironmentId))
 64        {
 165            ConnectionLifecycleReconciliationMetrics.RecordFailedScan();
 166            logger.LogWarning("Credential lifecycle scope provider returned an invalid scope.");
 167            return;
 68        }
 69
 111870        var scopeKey = (target.TenantId, target.EnvironmentId);
 111871        var cached = _cursors.GetValueOrDefault(scopeKey);
 111872        if (cached?.FailureTimestamp is long failureTimestamp &&
 111873            timeProvider.GetElapsedTime(failureTimestamp) < cached.RetryDelay)
 74            return;
 75
 111276        var cursor = cached?.Cursor;
 111277        var tenantAccessor = scope.ServiceProvider.GetRequiredService<ITenantAccessor>();
 111278        using var tenantContext = tenantAccessor.PushContext(new Tenant { Id = target.TenantId, Name = target.TenantId }
 111279        var candidates = scope.ServiceProvider.GetRequiredService<IConnectionDueCandidateStore>();
 80        try
 81        {
 111282            var page = await candidates.FindDueCandidatesAsync(
 111283                target.TenantId, target.EnvironmentId, timeProvider.GetUtcNow(), options.Value.BatchSize, cursor, cancel
 84
 111285            if (page.Items.Count == 0)
 86            {
 55287                ConnectionLifecycleReconciliationMetrics.RecordCompletedPage();
 55288                RememberScope(scopeKey, page.NextCursor, 0, null, TimeSpan.Zero);
 55289                return;
 90            }
 91
 56092            var failures = 0;
 56093            await Parallel.ForEachAsync(page.Items,
 56094                new ParallelOptions { MaxDegreeOfParallelism = options.Value.MaxConcurrency, CancellationToken = cancell
 56095                async (candidate, token) =>
 56096                {
 56997                    if (candidate.TenantId != target.TenantId || candidate.EnvironmentId != target.EnvironmentId ||
 56998                        string.IsNullOrWhiteSpace(candidate.ConnectionId) || candidate.DueAt > timeProvider.GetUtcNow())
 56099                    {
 2100                        Interlocked.Increment(ref failures);
 2101                        return;
 560102                    }
 560103
 560104                    try
 560105                    {
 567106                        if (!await DispatchAsync(candidate, token))
 4107                            Interlocked.Increment(ref failures);
 567108                    }
 0109                    catch (OperationCanceledException) when (token.IsCancellationRequested)
 560110                    {
 0111                        throw;
 560112                    }
 0113                    catch (Exception)
 560114                    {
 560115                        // Keep exceptions and candidate identifiers out of logs; provider errors may contain credential
 0116                        Interlocked.Increment(ref failures);
 0117                    }
 1129118                });
 119
 560120            if (failures > 0)
 121            {
 4122                ConnectionLifecycleReconciliationMetrics.RecordFailedPage();
 4123                logger.LogWarning("Credential lifecycle reconciliation page had {FailureCount} failed or out-of-scope ca
 4124                RememberFailure(scopeKey, cursor);
 4125                return;
 126            }
 127
 128            // Commit the next cursor only after every candidate in the page succeeds.
 556129            ConnectionLifecycleReconciliationMetrics.RecordCompletedPage();
 556130            RememberScope(scopeKey, page.NextCursor, 0, null, TimeSpan.Zero);
 556131            return;
 132        }
 0133        catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
 134        {
 0135            throw;
 136        }
 0137        catch (Exception)
 138        {
 0139            ConnectionLifecycleReconciliationMetrics.RecordFailedScan();
 0140            logger.LogWarning("Credential lifecycle reconciliation scope failed.");
 0141            RememberFailure(scopeKey, cursor);
 0142        }
 1119143    }
 144
 145    private async Task<bool> DispatchAsync(ConnectionDueCandidate candidate, CancellationToken cancellationToken)
 146    {
 567147        await using var scope = scopeFactory.CreateAsyncScope();
 567148        var tenantAccessor = scope.ServiceProvider.GetRequiredService<ITenantAccessor>();
 567149        using var tenantContext = tenantAccessor.PushContext(new Tenant { Id = candidate.TenantId, Name = candidate.Tena
 567150        var lifecycle = scope.ServiceProvider.GetRequiredService<IConnectionLifecycleRecoveryService>();
 151
 567152        switch (candidate.Kind)
 153        {
 154            case ConnectionDueCandidateKind.OAuthRefresh:
 2155                return (await lifecycle.RefreshAsync(candidate.TenantId, candidate.EnvironmentId, candidate.ConnectionId
 156            case ConnectionDueCandidateKind.ExpiredConnectionOperation:
 157            case ConnectionDueCandidateKind.RecoveryRequired:
 559158                var reconciliation = await lifecycle.ReconcileAsync(candidate.TenantId, candidate.EnvironmentId,
 559159                    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.
 559162                if (reconciliation.SafeErrorCode is "recovery_required" or "refresh_outcome_unknown" or "generation_publ
 3163                    ConnectionLifecycleReconciliationMetrics.RecordRecoveryRequired();
 559164                return reconciliation.Succeeded || reconciliation.SafeErrorCode == "recovery_required";
 165            case ConnectionDueCandidateKind.GenerationCleanup:
 1166                return (await lifecycle.CleanupGenerationAsync(candidate.TenantId, candidate.EnvironmentId,
 1167                    candidate.ConnectionId, candidate.CandidateId, cancellationToken)).Succeeded;
 168            case ConnectionDueCandidateKind.Offboarding:
 5169                var offboarding = await lifecycle.ReconcileOffboardingAsync(candidate.TenantId, candidate.EnvironmentId,
 5170                    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.
 5173                if (offboarding.SafeErrorCode == "offboarding_outcome_unknown" ||
 5174                    offboarding.Status == ConnectionOffboardingOperationStatus.UnknownOutcome)
 3175                    ConnectionLifecycleReconciliationMetrics.RecordUnknownOffboarding();
 5176                return offboarding.Accepted || offboarding.SafeErrorCode == "offboarding_outcome_unknown";
 177            default:
 0178                throw new InvalidOperationException("Unsupported credential lifecycle candidate kind.");
 179        }
 567180    }
 181
 182    private void RememberFailure((string TenantId, string EnvironmentId) scopeKey, string? cursor)
 183    {
 4184        var failures = _cursors.TryGetValue(scopeKey, out var existing)
 4185            ? Math.Min(existing.ConsecutiveFailures + 1, 30)
 4186            : 1;
 4187        RememberScope(scopeKey, cursor, failures, timeProvider.GetTimestamp(), GetRetryDelay(failures));
 4188    }
 189
 190    private void RememberScope((string TenantId, string EnvironmentId) scopeKey, string? cursor, int consecutiveFailures
 191        long? failureTimestamp, TimeSpan retryDelay)
 192    {
 1112193        if (_cursors.TryGetValue(scopeKey, out var existing))
 194        {
 1103195            _scopeOrder.Remove(existing.Node);
 196        }
 197        else
 198        {
 9199            if (_cursors.Count >= MaximumTrackedScopes && _scopeOrder.First is { } oldest)
 200            {
 0201                _scopeOrder.RemoveFirst();
 0202                _cursors.Remove(oldest.Value);
 203            }
 204        }
 205
 1112206        var node = _scopeOrder.AddLast(scopeKey);
 1112207        _cursors[scopeKey] = new ScopeCursorEntry(cursor, node, consecutiveFailures, failureTimestamp, retryDelay);
 1112208    }
 209
 210    private TimeSpan GetRetryDelay(int consecutiveFailures)
 211    {
 4212        var configured = options.Value.RetryBackoff.TotalMilliseconds;
 4213        var multiplier = Math.Pow(2, Math.Min(consecutiveFailures - 1, 20));
 4214        return TimeSpan.FromMilliseconds(Math.Min(configured * multiplier, options.Value.MaxRetryBackoff.TotalMillisecon
 215    }
 216
 3318217    private sealed record ScopeCursorEntry(string? Cursor, LinkedListNode<(string TenantId, string EnvironmentId)> Node,
 2229218        int ConsecutiveFailures, long? FailureTimestamp, TimeSpan RetryDelay);
 219}