< Summary

Information
Class: Elsa.Connections.Services.DefaultConnectionLifecycleService
Assembly: Elsa.Connections
File(s): /home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Connections/Services/DefaultConnectionLifecycleService.cs
Line coverage
77%
Covered lines: 496
Uncovered lines: 140
Coverable lines: 636
Total lines: 1253
Line coverage: 77.9%
Branch coverage
68%
Covered branches: 330
Total branches: 481
Branch coverage: 68.6%
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%
ConnectAsync()50%141488.88%
ConnectApiKeyAsync()50%101085.71%
ConnectCoreAsync()87.5%8882.97%
ReplaceApiKeyAsync()58.82%533474.41%
ResolveForUseAsync()50%2275%
ResolveForUseAsync()100%22100%
ResolveAuthorizedCredentialAsync()92.85%282892.3%
DisconnectAsync()43.75%181679.41%
RequestTokenRevocationAsync()50%2280%
RequestInstallationUninstallAsync()50%2280%
ReconcileOffboardingAsync()63.49%2896361.53%
RefreshAsync(...)100%11100%
RefreshAsync(...)100%11100%
RefreshCoreAsync()63.33%1006077.77%
CleanupGenerationAsync(...)100%11100%
CleanupGenerationAsync(...)100%210%
CleanupGenerationCoreAsync()63.15%563876.92%
ReconcileAsync()71.79%1697875.4%
AuthorizeAsync()100%11100%
QueueOffboardingOperationAsync()62.5%624076.08%
GetOffboardingResultAsync()50%66100%
ReleaseOffboardingClaimAsync()100%210%
RecordOffboardingFailureAsync()100%4481.81%
GetOffboardingOperationId(...)100%11100%
CanUseCurrentGeneration(...)100%1212100%
ToMetadata(...)100%11100%
PushTenant(...)100%11100%
ReleaseUnstartedRefreshAsync()50%44100%
TryReleaseUnstartedRefreshAsync()100%1166.66%
Serialize(...)100%11100%
SerializeApiKey(...)100%11100%
Deserialize(...)50%3250%
IsApiKeyGenerationAsync()75%16850%
IsApiKeySourceGenerationAsync()75%10866.66%
CanPromoteApiKeyRecoveryAsync()75%44100%
IsValidEnvelope(...)78.57%141487.5%
RestoreApiKeySourceIfPlanMissingAsync()56.25%171684.21%
CleanupUnreferencedOrphanGenerationAsync()50%7666.66%
RequireRecoveryAsync()100%11100%
TryMarkRecoveryRequiredAsync()100%1160%
get_Kind()100%11100%

File(s)

/home/runner/work/elsa-core/elsa-core/src/modules/Elsa.Connections/Services/DefaultConnectionLifecycleService.cs

#LineLine coverage
 1using System.Security.Claims;
 2using System.Security.Cryptography;
 3using System.Text;
 4using System.Text.Json;
 5using Elsa.Common.Multitenancy;
 6using Elsa.Connections.Contracts;
 7using Elsa.Connections.Models;
 8using Elsa.Secrets.Contracts;
 9using Elsa.Secrets.Models;
 10
 11namespace Elsa.Connections.Services;
 12
 13/// <summary>Coordinates one-time provider refreshes through durable intent, provider-call, stage, and publish states.</
 14914public sealed class DefaultConnectionLifecycleService(
 14915    IConnectionLifecycleStore store,
 14916    IConnectionUseAuthorizer authorizer,
 14917    IManagedSecretManager secrets,
 14918    TimeProvider timeProvider,
 14919    ITenantAccessor tenantAccessor,
 14920    IConnectionCredentialProvider? provider = null,
 14921    IConnectionOffboardingProvider? offboardingProvider = null) : IConnectionLifecycleService, IStaticApiKeyLifecycleSer
 22{
 5123    private static readonly TimeSpan OperationLeaseDuration = TimeSpan.FromMinutes(2);
 5124    private static readonly JsonSerializerOptions JsonOptions = new(JsonSerializerDefaults.Web);
 5125    private static readonly ClaimsPrincipal SystemPrincipal = new(new ClaimsIdentity(
 5126        [new Claim(ClaimTypes.NameIdentifier, "elsa-connections-lifecycle"), new Claim("elsa:identity-kind", "system")],
 5127        "Elsa.Connections.Server"));
 28
 29    public async Task<ConnectionLifecycleResult> ConnectAsync(ClaimsPrincipal principal, ConnectConnectionRequest reques
 30    {
 3731        if (string.IsNullOrWhiteSpace(request.TenantId) || string.IsNullOrWhiteSpace(request.EnvironmentId) ||
 3732            string.IsNullOrWhiteSpace(request.ProviderId) || string.IsNullOrWhiteSpace(request.ProviderAccountId) ||
 3733            string.IsNullOrWhiteSpace(request.InitialCredentials.AccessToken) || string.IsNullOrWhiteSpace(request.Initi
 3734            request.InitialCredentials.AccessTokenExpiresAt <= timeProvider.GetUtcNow())
 35        {
 036            return new ConnectionLifecycleResult(false, "connection_input_invalid", null);
 37        }
 38
 3739        return await ConnectCoreAsync(principal, request.TenantId, request.EnvironmentId, request.ProviderId,
 3740            request.ProviderAccountId, Serialize(request.InitialCredentials), ConnectionCredentialKind.OAuth,
 3741            request.InitialCredentials.AccessTokenExpiresAt, cancellationToken);
 3742    }
 43
 44    public async Task<ConnectionLifecycleResult> ConnectApiKeyAsync(ClaimsPrincipal principal, ConnectApiKeyConnectionRe
 45    {
 1346        if (string.IsNullOrWhiteSpace(request.TenantId) || string.IsNullOrWhiteSpace(request.EnvironmentId) ||
 1347            string.IsNullOrWhiteSpace(request.ProviderId) || string.IsNullOrWhiteSpace(request.ProviderAccountId) ||
 1348            string.IsNullOrWhiteSpace(request.ApiKey))
 49        {
 050            return new ConnectionLifecycleResult(false, "connection_input_invalid", null);
 51        }
 52
 1353        return await ConnectCoreAsync(principal, request.TenantId, request.EnvironmentId, request.ProviderId,
 1354            request.ProviderAccountId, SerializeApiKey(request.ApiKey), ConnectionCredentialKind.ApiKey, null, cancellat
 1355    }
 56
 57    private async Task<ConnectionLifecycleResult> ConnectCoreAsync(
 58        ClaimsPrincipal principal,
 59        string tenantId,
 60        string environmentId,
 61        string providerId,
 62        string providerAccountId,
 63        string encryptedEnvelope,
 64        ConnectionCredentialKind credentialKind,
 65        DateTimeOffset? credentialExpiresAt,
 66        CancellationToken cancellationToken)
 67    {
 5068        if (!await AuthorizeAsync(principal, ConnectionUseKind.Human, tenantId, environmentId, "", "manage:connect", can
 69        {
 070            return new ConnectionLifecycleResult(false, "connection_unavailable", null);
 71        }
 72
 5073        using var tenantContext = PushTenant(tenantId);
 5074        var connectionId = Guid.NewGuid().ToString("N");
 5075        var operationId = Guid.NewGuid().ToString("N");
 5076        var secretName = ManagedSecretNames.ForGeneration(connectionId, operationId);
 5077        var connection = new IntegrationConnection
 5078        {
 5079            Id = connectionId,
 5080            TenantId = tenantId,
 5081            EnvironmentId = environmentId,
 5082            ProviderId = providerId,
 5083            ProviderAccountId = providerAccountId,
 5084            Status = ConnectionStatus.Active,
 5085            Revision = 1,
 5086            OperationId = operationId,
 5087            OperationExpectedRevision = 1,
 5088            OperationFence = 1,
 5089            OperationLeaseExpiresAt = timeProvider.GetUtcNow() + OperationLeaseDuration,
 5090            OperationStatus = CredentialOperationStatus.CredentialReceived,
 5091            PlannedSecretName = secretName,
 5092            PlannedGenerationId = operationId
 5093        };
 94
 95        // Persist owner metadata and planned generation before encrypted material so a process crash is recoverable.
 96        try
 97        {
 5098            await store.CreateAsync(connection, cancellationToken);
 5099        }
 0100        catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
 101        {
 0102            throw;
 103        }
 0104        catch (Exception)
 105        {
 0106            return new ConnectionLifecycleResult(false, "connection_create_unknown", null, connectionId);
 107        }
 108        try
 109        {
 50110            await secrets.CreateGenerationAsync(connectionId, operationId, encryptedEnvelope, cancellationToken);
 50111            if (!await store.TryRecordStagedGenerationAsync(connectionId, tenantId, environmentId, 1, operationId, 1, se
 50112                    credentialKind, credentialExpiresAt, cancellationToken) ||
 50113                !await store.TryPublishGenerationAsync(connectionId, tenantId, environmentId, 1, operationId, 1, cancell
 114            {
 1115                await TryMarkRecoveryRequiredAsync(connection, tenantId, environmentId, "connection_publish_conflict");
 1116                return new ConnectionLifecycleResult(false, "connection_publish_conflict", 1, connectionId);
 117            }
 48118        }
 0119        catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
 120        {
 0121            await TryMarkRecoveryRequiredAsync(connection, tenantId, environmentId, "connection_outcome_unknown");
 0122            throw new OperationCanceledException("Connection setup was cancelled; creation outcome is unknown.", cancell
 123        }
 124        catch (Exception)
 125        {
 1126            await TryMarkRecoveryRequiredAsync(connection, tenantId, environmentId, "connection_outcome_unknown");
 1127            return new ConnectionLifecycleResult(false, "connection_outcome_unknown", 1, connectionId);
 128        }
 129
 48130        connection.CurrentSecretName = secretName;
 48131        connection.CurrentGenerationId = operationId;
 48132        connection.OperationStatus = CredentialOperationStatus.Completed;
 48133        connection.Revision = 2;
 48134        return new ConnectionLifecycleResult(true, null, connection.Revision, connectionId, ToMetadata(connection));
 50135    }
 136
 137    public async Task<ConnectionLifecycleResult> ReplaceApiKeyAsync(
 138        ClaimsPrincipal principal,
 139        string tenantId,
 140        string environmentId,
 141        string connectionId,
 142        long expectedRevision,
 143        string apiKey,
 144        CancellationToken cancellationToken = default)
 145    {
 12146        if (string.IsNullOrWhiteSpace(apiKey) || expectedRevision <= 0 ||
 12147            !await AuthorizeAsync(principal, ConnectionUseKind.Human, tenantId, environmentId, connectionId, "manage:rot
 148        {
 0149            return new ConnectionLifecycleResult(false, "connection_unavailable", null);
 150        }
 151
 12152        using var tenantContext = PushTenant(tenantId);
 12153        var current = await store.FindAsync(connectionId, tenantId, environmentId, cancellationToken);
 12154        if (current is not { Status: ConnectionStatus.Active } || current.Revision != expectedRevision ||
 12155            !await IsApiKeyGenerationAsync(current, cancellationToken))
 156        {
 6157            return new ConnectionLifecycleResult(false, "connection_unavailable", current?.Revision, connectionId,
 6158                current == null ? null : ToMetadata(current));
 159        }
 160
 6161        var operationId = Guid.NewGuid().ToString("N");
 6162        var claimed = await store.TryClaimCredentialUpdateAsync(connectionId, tenantId, environmentId,
 6163            expectedRevision, operationId, timeProvider.GetUtcNow() + OperationLeaseDuration, cancellationToken);
 6164        if (claimed == null)
 165        {
 0166            var latest = await store.FindAsync(connectionId, tenantId, environmentId, cancellationToken);
 0167            return new ConnectionLifecycleResult(false, "connection_conflict", latest?.Revision, connectionId,
 0168                latest == null ? null : ToMetadata(latest));
 169        }
 170
 6171        var expectedOperationRevision = claimed.OperationExpectedRevision;
 6172        var fence = claimed.OperationFence;
 6173        var secretName = ManagedSecretNames.ForGeneration(connectionId, operationId);
 174        try
 175        {
 6176            if (!await store.TryAcceptCredentialUpdateAsync(connectionId, tenantId, environmentId,
 6177                    expectedOperationRevision, operationId, fence, timeProvider.GetUtcNow(), cancellationToken))
 178            {
 0179                await TryReleaseUnstartedRefreshAsync(claimed, tenantId, environmentId, "rotation_conflict");
 0180                return new ConnectionLifecycleResult(false, "rotation_conflict", expectedOperationRevision, connectionId
 181            }
 182
 6183            await secrets.CreateGenerationAsync(connectionId, operationId, SerializeApiKey(apiKey), cancellationToken);
 6184            if (!await store.TryRecordStagedGenerationAsync(connectionId, tenantId, environmentId, expectedOperationRevi
 6185                    operationId, fence, secretName, operationId, ConnectionCredentialKind.ApiKey, null, cancellationToke
 186            {
 1187                await CleanupUnreferencedOrphanGenerationAsync(connectionId, tenantId, environmentId, operationId, cance
 1188                await TryMarkRecoveryRequiredAsync(claimed, tenantId, environmentId, "rotation_publish_conflict");
 1189                return new ConnectionLifecycleResult(false, "rotation_publish_conflict", expectedOperationRevision, conn
 190            }
 191
 4192            if (!await store.TryPublishGenerationAsync(connectionId, tenantId, environmentId, expectedOperationRevision,
 4193                    operationId, fence, cancellationToken))
 194            {
 0195                await TryMarkRecoveryRequiredAsync(claimed, tenantId, environmentId, "rotation_publish_conflict");
 0196                return new ConnectionLifecycleResult(false, "rotation_publish_conflict", expectedOperationRevision, conn
 197            }
 198
 4199            var published = await store.FindAsync(connectionId, tenantId, environmentId, cancellationToken);
 4200            return published == null
 4201                ? new ConnectionLifecycleResult(false, "connection_unavailable", null)
 4202                : new ConnectionLifecycleResult(true, null, published.Revision, connectionId, ToMetadata(published));
 203        }
 0204        catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
 205        {
 0206            await TryMarkRecoveryRequiredAsync(claimed, tenantId, environmentId, "rotation_outcome_unknown");
 0207            throw new OperationCanceledException("API-key replacement outcome is unknown.", cancellationToken);
 208        }
 209        catch (Exception)
 210        {
 1211            await TryMarkRecoveryRequiredAsync(claimed, tenantId, environmentId, "rotation_outcome_unknown");
 1212            return new ConnectionLifecycleResult(false, "rotation_outcome_unknown", expectedOperationRevision, connectio
 213        }
 12214    }
 215
 216    public async Task<ConnectionAccessCredential> ResolveForUseAsync(ClaimsPrincipal principal, string tenantId, string 
 217    {
 29218        if (!await AuthorizeAsync(principal, ConnectionUseKind.Human, tenantId, environmentId, connectionId, "use", canc
 219        {
 0220            throw new ConnectionUnavailableException();
 221        }
 222
 29223        return await ResolveAuthorizedCredentialAsync(tenantId, environmentId, connectionId, cancellationToken);
 15224    }
 225
 226    public async Task<ConnectionAccessCredential> ResolveForUseAsync(string tenantId, string environmentId, string conne
 227    {
 7228        if (!await AuthorizeAsync(SystemPrincipal, ConnectionUseKind.BackgroundSystem, tenantId, environmentId, connecti
 229        {
 1230            throw new ConnectionUnavailableException();
 231        }
 232
 6233        return await ResolveAuthorizedCredentialAsync(tenantId, environmentId, connectionId, cancellationToken);
 6234    }
 235
 236    private async Task<ConnectionAccessCredential> ResolveAuthorizedCredentialAsync(string tenantId, string environmentI
 237    {
 35238        using var tenantContext = PushTenant(tenantId);
 35239        var connection = await store.FindAsync(connectionId, tenantId, environmentId, cancellationToken);
 35240        if (!CanUseCurrentGeneration(connection))
 241        {
 3242            throw new ConnectionUnavailableException();
 243        }
 244
 245        CredentialEnvelope? material;
 246        try
 247        {
 32248            var payload = await secrets.ResolveGenerationAsync(connection!.CurrentSecretName!, connection.Id, connection
 26249            material = Deserialize(payload.Value);
 26250        }
 0251        catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
 252        {
 0253            throw new OperationCanceledException("Connection access was cancelled.", cancellationToken);
 254        }
 6255        catch (Exception)
 256        {
 6257            throw new ConnectionUnavailableException();
 258        }
 259
 26260        if (!IsValidEnvelope(material) ||
 26261            (material!.Kind ?? ConnectionCredentialKind.OAuth) == ConnectionCredentialKind.OAuth &&
 26262            material.AccessTokenExpiresAt <= timeProvider.GetUtcNow())
 263        {
 4264            throw new ConnectionUnavailableException();
 265        }
 266
 22267        var validMaterial = material!;
 268
 22269        var latest = await store.FindAsync(connectionId, tenantId, environmentId, cancellationToken);
 22270        if (latest == null || latest.Status != ConnectionStatus.Active ||
 22271            latest.Revision != connection.Revision || latest.CurrentGenerationId != connection.CurrentGenerationId ||
 22272            latest.CurrentSecretName != connection.CurrentSecretName ||
 22273            latest.OperationStatus is not (CredentialOperationStatus.None or CredentialOperationStatus.Completed))
 274        {
 1275            throw new ConnectionUnavailableException();
 276        }
 277
 21278        return validMaterial.Kind == ConnectionCredentialKind.ApiKey
 21279            ? new ConnectionAccessCredential(ConnectionCredentialKind.ApiKey, validMaterial.AccessToken!, null)
 21280            : new ConnectionAccessCredential(ConnectionCredentialKind.OAuth, validMaterial.AccessToken!, validMaterial.A
 21281    }
 282
 283    public async Task<ConnectionOffboardingOperationResult> DisconnectAsync(
 284        ClaimsPrincipal principal,
 285        string tenantId,
 286        string environmentId,
 287        string connectionId,
 288        CancellationToken cancellationToken = default)
 289    {
 17290        if (!await AuthorizeAsync(principal, ConnectionUseKind.Human, tenantId, environmentId, connectionId, "manage:dis
 291        {
 0292            return new ConnectionOffboardingOperationResult(false, "connection_unavailable", null, null, null);
 293        }
 294
 17295        using var tenantContext = PushTenant(tenantId);
 17296        var connection = await store.FindAsync(connectionId, tenantId, environmentId, cancellationToken);
 17297        if (connection == null)
 298        {
 0299            return new ConnectionOffboardingOperationResult(false, "connection_unavailable", null, null, null);
 300        }
 301
 17302        var operationId = GetOffboardingOperationId(tenantId, environmentId, connectionId, ConnectionOffboardingOperatio
 17303        var existing = await store.FindOffboardingOperationAsync(operationId, tenantId, environmentId, connectionId, can
 17304        if (existing != null)
 305        {
 2306            return new ConnectionOffboardingOperationResult(true, null, operationId, existing.Status, connection.Revisio
 307        }
 308
 15309        var now = timeProvider.GetUtcNow();
 15310        var operation = new ConnectionOffboardingOperation
 15311        {
 15312            Id = operationId,
 15313            TenantId = tenantId,
 15314            EnvironmentId = environmentId,
 15315            ConnectionId = connectionId,
 15316            ProviderId = connection.ProviderId,
 15317            ProviderAccountId = connection.ProviderAccountId,
 15318            Kind = ConnectionOffboardingOperationKind.LocalDisconnect,
 15319            Status = ConnectionOffboardingOperationStatus.Completed,
 15320            Fence = 1,
 15321            CreatedAt = now,
 15322            UpdatedAt = now
 15323        };
 15324        var disconnected = await store.TryDisconnectAndRecordAsync(connectionId, tenantId, environmentId, connection.Rev
 15325        if (disconnected != null)
 326        {
 15327            return new ConnectionOffboardingOperationResult(true, null, operationId, operation.Status, disconnected.Revi
 328        }
 329
 0330        var latestOperation = await store.FindOffboardingOperationAsync(operationId, tenantId, environmentId, connection
 0331        var latestConnection = await store.FindAsync(connectionId, tenantId, environmentId, cancellationToken);
 0332        return latestOperation != null && latestConnection != null
 0333            ? new ConnectionOffboardingOperationResult(true, null, operationId, latestOperation.Status, latestConnection
 0334            : new ConnectionOffboardingOperationResult(false, "connection_conflict", operationId, null, latestConnection
 17335    }
 336
 337    public async Task<ConnectionOffboardingOperationResult> RequestTokenRevocationAsync(
 338        ClaimsPrincipal principal,
 339        string tenantId,
 340        string environmentId,
 341        string connectionId,
 342        string generationId,
 343        CancellationToken cancellationToken = default)
 344    {
 10345        if (!await AuthorizeAsync(principal, ConnectionUseKind.Human, tenantId, environmentId, connectionId, "manage:rev
 346        {
 0347            return new ConnectionOffboardingOperationResult(false, "connection_unavailable", null, null, null);
 348        }
 349
 10350        return await QueueOffboardingOperationAsync(tenantId, environmentId, connectionId,
 10351            ConnectionOffboardingOperationKind.TokenPairRevocation, generationId, cancellationToken);
 10352    }
 353
 354    public async Task<ConnectionOffboardingOperationResult> RequestInstallationUninstallAsync(
 355        ClaimsPrincipal principal,
 356        string tenantId,
 357        string environmentId,
 358        string connectionId,
 359        CancellationToken cancellationToken = default)
 360    {
 6361        if (!await AuthorizeAsync(principal, ConnectionUseKind.Human, tenantId, environmentId, connectionId, "manage:uni
 362        {
 0363            return new ConnectionOffboardingOperationResult(false, "connection_unavailable", null, null, null);
 364        }
 365
 6366        return await QueueOffboardingOperationAsync(tenantId, environmentId, connectionId,
 6367            ConnectionOffboardingOperationKind.InstallationUninstall, null, cancellationToken);
 6368    }
 369
 370    public async Task<ConnectionOffboardingOperationResult> ReconcileOffboardingAsync(
 371        string tenantId,
 372        string environmentId,
 373        string connectionId,
 374        CancellationToken cancellationToken = default)
 375    {
 16376        if (!await AuthorizeAsync(SystemPrincipal, ConnectionUseKind.BackgroundSystem, tenantId, environmentId, connecti
 377        {
 0378            return new ConnectionOffboardingOperationResult(false, "connection_unavailable", null, null, null);
 379        }
 380
 16381        using var tenantContext = PushTenant(tenantId);
 16382        var now = timeProvider.GetUtcNow();
 16383        var stableRevocationIdempotency = offboardingProvider?.SupportsStableOperationIdIdempotency(ConnectionOffboardin
 16384        var stableUninstallIdempotency = offboardingProvider?.SupportsStableOperationIdIdempotency(ConnectionOffboarding
 16385        var pending = await store.FindNextOffboardingOperationAsync(
 16386            tenantId, environmentId, connectionId, now, stableRevocationIdempotency, stableUninstallIdempotency, cancell
 16387        if (pending == null)
 388        {
 4389            var current = await store.FindAsync(connectionId, tenantId, environmentId, cancellationToken);
 4390            return new ConnectionOffboardingOperationResult(true, null, null, null, current?.Revision);
 391        }
 392
 12393        var supportsStableIdempotency = offboardingProvider?.SupportsStableOperationIdIdempotency(pending.Kind) == true;
 394
 12395        if (!supportsStableIdempotency && pending.Status == ConnectionOffboardingOperationStatus.UnknownOutcome)
 396        {
 1397            return await GetOffboardingResultAsync(pending.Id, tenantId, environmentId, connectionId, false,
 1398                "offboarding_outcome_unknown", cancellationToken);
 399        }
 400
 11401        if (!supportsStableIdempotency && pending.Status == ConnectionOffboardingOperationStatus.ProviderCallStarted &&
 11402            pending.LeaseExpiresAt <= now)
 403        {
 2404            await store.TryMarkOffboardingOutcomeUnknownIfLeaseExpiredAsync(pending.Id, tenantId, environmentId, connect
 2405                pending.Fence, now, "provider_outcome_unknown", CancellationToken.None);
 2406            return await GetOffboardingResultAsync(pending.Id, tenantId, environmentId, connectionId, false,
 2407                "offboarding_outcome_unknown", cancellationToken);
 408        }
 409
 9410        var claimed = await store.TryClaimOffboardingOperationAsync(
 9411            pending.Id, tenantId, environmentId, connectionId, pending.Fence, now, now + OperationLeaseDuration, cancell
 9412        if (claimed == null)
 413        {
 0414            return await GetOffboardingResultAsync(pending.Id, tenantId, environmentId, connectionId, false, "offboardin
 415        }
 416
 9417        if (offboardingProvider == null)
 418        {
 0419            await ReleaseOffboardingClaimAsync(claimed, tenantId, environmentId, connectionId, "offboarding_provider_una
 0420            return await GetOffboardingResultAsync(claimed.Id, tenantId, environmentId, connectionId, false, "offboardin
 421        }
 422
 9423        CredentialMaterial? credentials = null;
 9424        if (claimed.Kind == ConnectionOffboardingOperationKind.TokenPairRevocation)
 425        {
 426            try
 427            {
 7428                if (string.IsNullOrWhiteSpace(claimed.GenerationId))
 429                {
 0430                    throw new ConnectionUnavailableException();
 431                }
 432
 7433                var secretName = ManagedSecretNames.ForGeneration(connectionId, claimed.GenerationId);
 7434                var payload = await secrets.ResolveGenerationAsync(secretName, connectionId, claimed.GenerationId, cance
 7435                var envelope = Deserialize(payload.Value);
 7436                var accessTokenExpiresAt = envelope?.AccessTokenExpiresAt;
 7437                if (envelope is not { Kind: null or ConnectionCredentialKind.OAuth, AccessToken: not null, RefreshToken:
 7438                    !accessTokenExpiresAt.HasValue)
 439                {
 0440                    throw new ConnectionUnavailableException();
 441                }
 442
 7443                credentials = new CredentialMaterial(envelope.AccessToken, envelope.RefreshToken, accessTokenExpiresAt.V
 7444            }
 0445            catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
 446            {
 0447                await ReleaseOffboardingClaimAsync(claimed, tenantId, environmentId, connectionId, "offboarding_cancelle
 0448                throw new OperationCanceledException("Credential offboarding was cancelled before the provider call.", c
 449            }
 450            catch (Exception)
 451            {
 0452                await ReleaseOffboardingClaimAsync(claimed, tenantId, environmentId, connectionId, "credential_unavailab
 0453                return await GetOffboardingResultAsync(claimed.Id, tenantId, environmentId, connectionId, false, "creden
 454            }
 455        }
 456
 9457        now = timeProvider.GetUtcNow();
 9458        if (!await store.TryStartOffboardingProviderCallAsync(claimed.Id, tenantId, environmentId, connectionId, claimed
 459        {
 0460            await ReleaseOffboardingClaimAsync(claimed, tenantId, environmentId, connectionId, "offboarding_conflict");
 0461            return await GetOffboardingResultAsync(claimed.Id, tenantId, environmentId, connectionId, false, "offboardin
 462        }
 463
 464        try
 465        {
 9466            var providerResult = claimed.Kind switch
 9467            {
 7468                ConnectionOffboardingOperationKind.TokenPairRevocation when credentials != null =>
 7469                    await offboardingProvider.RevokeTokenPairAsync(claimed.ProviderId, claimed.ProviderAccountId, claime
 9470                ConnectionOffboardingOperationKind.InstallationUninstall =>
 2471                    await offboardingProvider.UninstallInstallationAsync(claimed.ProviderId, claimed.ProviderAccountId, 
 0472                _ => ConnectionOffboardingProviderResult.TerminalFailure
 9473            };
 474
 9475            now = timeProvider.GetUtcNow();
 476            switch (providerResult)
 477            {
 478                case ConnectionOffboardingProviderResult.Succeeded:
 479                    // Once the provider has confirmed success, caller cancellation must not turn that known
 480                    // result into an unknown operation. Persist the semantic outcome independently.
 7481                    await store.TryCompleteOffboardingOperationAsync(claimed.Id, tenantId, environmentId, connectionId, 
 7482                    break;
 483                case ConnectionOffboardingProviderResult.RetryableFailure:
 0484                    await RecordOffboardingFailureAsync(claimed, tenantId, environmentId, connectionId,
 0485                        ConnectionOffboardingOperationStatus.RetryScheduled, "provider_retryable_failure");
 0486                    break;
 487                case ConnectionOffboardingProviderResult.TerminalFailure:
 0488                    await RecordOffboardingFailureAsync(claimed, tenantId, environmentId, connectionId,
 0489                        ConnectionOffboardingOperationStatus.TerminalFailure, "provider_terminal_failure");
 0490                    break;
 491                case ConnectionOffboardingProviderResult.UnknownOutcome:
 2492                    await RecordOffboardingFailureAsync(claimed, tenantId, environmentId, connectionId,
 2493                        ConnectionOffboardingOperationStatus.UnknownOutcome, "provider_outcome_unknown", supportsStableI
 494                    break;
 495            }
 9496        }
 0497        catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
 498        {
 0499            await RecordOffboardingFailureAsync(claimed, tenantId, environmentId, connectionId,
 0500                ConnectionOffboardingOperationStatus.UnknownOutcome, "provider_outcome_unknown", supportsStableIdempoten
 0501            throw new OperationCanceledException("Credential offboarding was cancelled; provider outcome is unknown.", c
 502        }
 503        catch (Exception)
 504        {
 0505            await RecordOffboardingFailureAsync(claimed, tenantId, environmentId, connectionId,
 0506                ConnectionOffboardingOperationStatus.UnknownOutcome, "provider_outcome_unknown", supportsStableIdempoten
 507        }
 508
 9509        return await GetOffboardingResultAsync(claimed.Id, tenantId, environmentId, connectionId, true, null, cancellati
 15510    }
 511
 512    public Task<ConnectionLifecycleResult> RefreshAsync(ClaimsPrincipal principal, string tenantId, string environmentId
 26513        RefreshCoreAsync(principal, ConnectionUseKind.Human, tenantId, environmentId, connectionId, cancellationToken);
 514
 515    public Task<ConnectionLifecycleResult> RefreshAsync(string tenantId, string environmentId, string connectionId, Canc
 1516        RefreshCoreAsync(SystemPrincipal, ConnectionUseKind.BackgroundSystem, tenantId, environmentId, connectionId, can
 517
 518    private async Task<ConnectionLifecycleResult> RefreshCoreAsync(ClaimsPrincipal principal, ConnectionUseKind useKind,
 519    {
 27520        if (!await AuthorizeAsync(principal, useKind, tenantId, environmentId, connectionId, "manage:refresh", cancellat
 521        {
 4522            return new ConnectionLifecycleResult(false, "connection_unavailable", null);
 523        }
 524
 23525        using var tenantContext = PushTenant(tenantId);
 23526        var current = await store.FindAsync(connectionId, tenantId, environmentId, cancellationToken);
 23527        if (current == null)
 528        {
 0529            return new ConnectionLifecycleResult(false, "connection_unavailable", null);
 530        }
 23531        if (!CanUseCurrentGeneration(current))
 532        {
 4533            var code = current.Status == ConnectionStatus.Active ? "refresh_conflict" : "connection_unavailable";
 4534            return new ConnectionLifecycleResult(false, code, current.Revision, connectionId, ToMetadata(current));
 535        }
 536
 537        // Read before claiming so a static key never advances the revision through an unsupported
 538        // refresh. A failed read is still a safe, pre-provider failure and leaves the connection usable.
 539        CredentialEnvelope? currentMaterial;
 19540        var preflightReadFailed = false;
 541        try
 542        {
 19543            var payload = await secrets.ResolveGenerationAsync(current.CurrentSecretName!, current.Id, current.CurrentGe
 18544            currentMaterial = Deserialize(payload.Value);
 18545        }
 0546        catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
 547        {
 0548            throw;
 549        }
 1550        catch (Exception)
 551        {
 1552            currentMaterial = null;
 1553            preflightReadFailed = true;
 1554        }
 555
 19556        if (currentMaterial?.Kind == ConnectionCredentialKind.ApiKey && IsValidEnvelope(currentMaterial))
 557        {
 2558            return new ConnectionLifecycleResult(false, "credential_refresh_unsupported", current.Revision, connectionId
 559        }
 560
 17561        if (provider == null)
 562        {
 0563            var code = IsValidEnvelope(currentMaterial) ? "credential_provider_unavailable" : "credential_unavailable";
 0564            return new ConnectionLifecycleResult(false, code, current.Revision, connectionId, ToMetadata(current));
 565        }
 566
 17567        var operationId = Guid.NewGuid().ToString("N");
 17568        var claimed = await store.TryClaimRefreshAsync(connectionId, tenantId, environmentId, current!.Revision, operati
 17569        if (claimed == null)
 570        {
 571            // Another worker may already have claimed or completed the refresh. Return state reloaded after the
 572            // lost CAS instead of reporting the revision from this worker's stale pre-claim snapshot.
 0573            var latest = await store.FindAsync(connectionId, tenantId, environmentId, cancellationToken);
 0574            return latest == null
 0575                ? new ConnectionLifecycleResult(false, "connection_unavailable", null)
 0576                : new ConnectionLifecycleResult(false, "refresh_conflict", latest.Revision, connectionId, ToMetadata(lat
 577        }
 578
 17579        var expectedRevision = claimed.OperationExpectedRevision;
 17580        var fence = claimed.OperationFence;
 17581        var providerCallStarted = false;
 582        try
 583        {
 17584            if (preflightReadFailed)
 585            {
 1586                return await ReleaseUnstartedRefreshAsync(claimed, tenantId, environmentId, "refresh_not_started");
 587            }
 588
 16589            var oldPayload = await secrets.ResolveGenerationAsync(claimed.CurrentSecretName!, claimed.Id, claimed.Curren
 16590            var oldMaterial = Deserialize(oldPayload.Value);
 16591            if (oldMaterial?.Kind == ConnectionCredentialKind.ApiKey)
 592            {
 0593                return await ReleaseUnstartedRefreshAsync(claimed, tenantId, environmentId, "credential_refresh_unsuppor
 594            }
 595
 16596            if (oldMaterial == null || oldMaterial.Kind is not null and not ConnectionCredentialKind.OAuth ||
 16597                string.IsNullOrWhiteSpace(oldMaterial.AccessToken) || string.IsNullOrWhiteSpace(oldMaterial.RefreshToken
 598            {
 0599                return await ReleaseUnstartedRefreshAsync(claimed, tenantId, environmentId, "credential_unavailable");
 600            }
 601
 602            // Persist this edge before crossing the provider boundary. After it, no worker may replay the token.
 16603            if (!await store.TryStartProviderCallAsync(connectionId, tenantId, environmentId, expectedRevision, operatio
 604            {
 0605                return await ReleaseUnstartedRefreshAsync(claimed, tenantId, environmentId, "refresh_conflict");
 606            }
 16607            providerCallStarted = true;
 608
 16609            var refreshed = await provider.RefreshAsync(claimed.ProviderId, claimed.ProviderAccountId, oldMaterial.Refre
 15610            if (refreshed == null || string.IsNullOrWhiteSpace(refreshed.RefreshToken) || string.IsNullOrWhiteSpace(refr
 611            {
 0612                return await RequireRecoveryAsync(claimed, tenantId, environmentId, "provider_refresh_unknown");
 613            }
 614
 15615            var nextName = ManagedSecretNames.ForGeneration(connectionId, operationId);
 15616            var encryptedEnvelope = Serialize(refreshed);
 15617            await secrets.CreateGenerationAsync(connectionId, operationId, encryptedEnvelope, cancellationToken);
 618
 15619            if (!await store.TryRecordStagedGenerationAsync(connectionId, tenantId, environmentId, expectedRevision, ope
 15620                    nextName, operationId, ConnectionCredentialKind.OAuth, refreshed.AccessTokenExpiresAt, cancellationT
 621            {
 1622                return await RequireRecoveryAsync(claimed, tenantId, environmentId, "credential_stage_unknown");
 623            }
 624
 13625            if (!await store.TryPublishGenerationAsync(connectionId, tenantId, environmentId, expectedRevision, operatio
 626            {
 2627                return await RequireRecoveryAsync(claimed, tenantId, environmentId, "generation_publish_conflict");
 628            }
 629
 11630            var published = await store.FindAsync(connectionId, tenantId, environmentId, cancellationToken);
 11631            return published == null
 11632                ? new ConnectionLifecycleResult(false, "connection_unavailable", null)
 11633                : new ConnectionLifecycleResult(true, null, published.Revision, connectionId, ToMetadata(published));
 634        }
 1635        catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
 636        {
 1637            if (providerCallStarted)
 638            {
 1639                await TryMarkRecoveryRequiredAsync(claimed, tenantId, environmentId, "refresh_outcome_unknown");
 1640                throw new OperationCanceledException("Credential refresh was cancelled; provider outcome is unknown.", c
 641            }
 642
 0643            await TryReleaseUnstartedRefreshAsync(claimed, tenantId, environmentId, "refresh_not_started");
 0644            throw new OperationCanceledException("Credential refresh was cancelled before the provider call.", cancellat
 645        }
 646        catch (Exception)
 647        {
 1648            if (!providerCallStarted)
 649            {
 0650                return await ReleaseUnstartedRefreshAsync(claimed, tenantId, environmentId, "refresh_not_started");
 651            }
 652
 653            // Provider/network/persistence errors after the durable call-start edge can hide a one-time refresh-token r
 1654            await TryMarkRecoveryRequiredAsync(claimed, tenantId, environmentId, "refresh_outcome_unknown");
 1655            return new ConnectionLifecycleResult(false, "refresh_outcome_unknown", expectedRevision);
 656        }
 26657    }
 658
 659    public Task<ConnectionLifecycleResult> CleanupGenerationAsync(
 660        ClaimsPrincipal principal,
 661        string tenantId,
 662        string environmentId,
 663        string connectionId,
 664        string generationId,
 665        CancellationToken cancellationToken = default) =>
 19666        CleanupGenerationCoreAsync(principal, ConnectionUseKind.Human, tenantId, environmentId, connectionId, generation
 667
 668    public Task<ConnectionLifecycleResult> CleanupGenerationAsync(
 669        string tenantId,
 670        string environmentId,
 671        string connectionId,
 672        string generationId,
 673        CancellationToken cancellationToken = default) =>
 0674        CleanupGenerationCoreAsync(SystemPrincipal, ConnectionUseKind.BackgroundSystem, tenantId, environmentId, connect
 675
 676    private async Task<ConnectionLifecycleResult> CleanupGenerationCoreAsync(
 677        ClaimsPrincipal principal,
 678        ConnectionUseKind useKind,
 679        string tenantId,
 680        string environmentId,
 681        string connectionId,
 682        string generationId,
 683        CancellationToken cancellationToken)
 684    {
 19685        if (!await AuthorizeAsync(principal, useKind, tenantId, environmentId, connectionId, "manage:cleanup", cancellat
 686        {
 0687            return new ConnectionLifecycleResult(false, "connection_unavailable", null);
 688        }
 689
 19690        if (string.IsNullOrWhiteSpace(generationId))
 691        {
 0692            return new ConnectionLifecycleResult(false, "generation_unavailable", null);
 693        }
 694
 19695        using var tenantContext = PushTenant(tenantId);
 19696        var connection = await store.FindAsync(connectionId, tenantId, environmentId, cancellationToken);
 19697        if (connection == null)
 698        {
 0699            return new ConnectionLifecycleResult(false, "connection_unavailable", null);
 700        }
 701
 19702        var cleanup = await store.FindGenerationCleanupAsync(connectionId, tenantId, environmentId, generationId, cancel
 19703        if (cleanup?.Status == ConnectionGenerationCleanupStatus.Deleted)
 704        {
 1705            return new ConnectionLifecycleResult(true, null, connection.Revision, connectionId, ToMetadata(connection));
 706        }
 707
 18708        var now = timeProvider.GetUtcNow();
 18709        var cleanupClaim = await store.TryClaimGenerationCleanupAsync(
 18710            connectionId, tenantId, environmentId, connection.Revision, generationId, now, now + OperationLeaseDuration,
 18711        if (cleanupClaim == null)
 712        {
 9713            cleanup = await store.FindGenerationCleanupAsync(connectionId, tenantId, environmentId, generationId, cancel
 9714            connection = await store.FindAsync(connectionId, tenantId, environmentId, cancellationToken);
 9715            if (cleanup?.Status == ConnectionGenerationCleanupStatus.Deleted && connection != null)
 716            {
 0717                return new ConnectionLifecycleResult(true, null, connection.Revision, connectionId, ToMetadata(connectio
 718            }
 719
 9720            var errorCode = cleanup != null && cleanup.Status == ConnectionGenerationCleanupStatus.Deleting && cleanup.L
 9721                ? "generation_cleanup_in_progress"
 9722                : "generation_in_use";
 9723            return new ConnectionLifecycleResult(false, errorCode, connection?.Revision, connectionId, connection is nul
 724        }
 725
 9726        var name = ManagedSecretNames.ForGeneration(connectionId, generationId);
 727        try
 728        {
 729            // This host-only Secrets primitive validates the immutable owner/generation marker. The lifecycle
 730            // store claim above is the authorization and no-reference proof; raw host callers must not bypass it.
 9731            if (!await secrets.DeleteGenerationAsync(name, connectionId, generationId, cancellationToken))
 732            {
 1733                await store.CancelGenerationCleanupAsync(connectionId, tenantId, environmentId, generationId, cleanupCla
 1734                return new ConnectionLifecycleResult(false, "generation_unavailable", connection.Revision, connectionId,
 735            }
 736
 8737            if (!await store.CompleteGenerationCleanupAsync(connectionId, tenantId, environmentId, generationId, cleanup
 738            {
 0739                return new ConnectionLifecycleResult(false, "generation_cleanup_unknown", connection.Revision, connectio
 740            }
 8741        }
 0742        catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
 743        {
 744            // Keep the durable Deleting tombstone. A retry repeats only the idempotent owner-checked deletion.
 0745            throw new OperationCanceledException("Credential generation cleanup was cancelled; cleanup outcome is unknow
 746        }
 0747        catch (Exception)
 748        {
 749            // Keep the durable Deleting tombstone if the external Secrets write may have completed.
 0750            return new ConnectionLifecycleResult(false, "generation_cleanup_unknown", connection.Revision, connectionId,
 751        }
 752
 8753        var updated = await store.FindAsync(connectionId, tenantId, environmentId, cancellationToken);
 8754        return updated == null
 8755            ? new ConnectionLifecycleResult(false, "connection_unavailable", null)
 8756            : new ConnectionLifecycleResult(true, null, updated.Revision, connectionId, ToMetadata(updated));
 19757    }
 758
 759    public async Task<ConnectionLifecycleResult> ReconcileAsync(string tenantId, string environmentId, string connection
 760    {
 24761        if (!await AuthorizeAsync(SystemPrincipal, ConnectionUseKind.BackgroundSystem, tenantId, environmentId, connecti
 762        {
 0763            return new ConnectionLifecycleResult(false, "connection_unavailable", null);
 764        }
 765
 24766        using var tenantContext = PushTenant(tenantId);
 24767        var connection = await store.FindAsync(connectionId, tenantId, environmentId, cancellationToken);
 24768        if (connection == null)
 769        {
 0770            return new ConnectionLifecycleResult(false, "connection_unavailable", null);
 771        }
 772
 24773        if (connection.OperationStatus is CredentialOperationStatus.None or CredentialOperationStatus.Completed)
 774        {
 1775            return new ConnectionLifecycleResult(true, null, connection.Revision);
 776        }
 777
 23778        if (connection.OperationStatus == CredentialOperationStatus.Claimed)
 779        {
 5780            var expiredClaim = await store.TryReleaseExpiredRefreshClaimAsync(
 5781                connection.Id, tenantId, environmentId, connection.OperationId!, connection.OperationFence,
 5782                timeProvider.GetUtcNow(), "refresh_not_started", cancellationToken);
 5783            if (expiredClaim)
 784            {
 3785                var released = await store.FindAsync(connectionId, tenantId, environmentId, cancellationToken);
 3786                return new ConnectionLifecycleResult(false,
 3787                    released?.Status == ConnectionStatus.Active ? "refresh_not_started" : "connection_unavailable",
 3788                    released?.Revision, connectionId);
 789            }
 790
 2791            var latest = await store.FindAsync(connectionId, tenantId, environmentId, cancellationToken);
 2792            return new ConnectionLifecycleResult(false,
 2793                latest is { Status: ConnectionStatus.Active, OperationStatus: CredentialOperationStatus.Claimed or Crede
 2794                    ? "operation_in_progress"
 2795                    : "connection_unavailable",
 2796                latest?.Revision ?? connection.Revision, connectionId);
 797        }
 798
 18799        if (connection.OperationStatus is CredentialOperationStatus.ProviderCallStarted or CredentialOperationStatus.Cre
 800        {
 9801            var expired = await store.TryMarkRecoveryRequiredIfLeaseExpiredAsync(connection.Id, tenantId, environmentId,
 9802            if (expired)
 803            {
 5804                if (connection.OperationStatus == CredentialOperationStatus.CredentialReceived &&
 5805                    await IsApiKeySourceGenerationAsync(connection, cancellationToken))
 806                {
 1807                    var recovery = await store.FindAsync(connectionId, tenantId, environmentId, cancellationToken);
 1808                    if (recovery != null)
 809                    {
 1810                        var restored = await RestoreApiKeySourceIfPlanMissingAsync(recovery, tenantId, environmentId, ca
 1811                        if (restored != null)
 812                        {
 1813                            return restored;
 814                        }
 815                    }
 816                }
 817
 4818                return new ConnectionLifecycleResult(false, "refresh_outcome_unknown", connection.Revision + (connection
 819            }
 820
 4821            return new ConnectionLifecycleResult(false, "operation_in_progress", connection.Revision);
 822        }
 823
 9824        if (connection.OperationStatus == CredentialOperationStatus.Staged)
 825        {
 2826            if (connection.OperationLeaseExpiresAt > timeProvider.GetUtcNow())
 827            {
 1828                return new ConnectionLifecycleResult(false, "operation_in_progress", connection.Revision);
 829            }
 830
 1831            if (await store.TryPublishGenerationAsync(connectionId, tenantId, environmentId, connection.OperationExpecte
 832            {
 1833                return new ConnectionLifecycleResult(true, null, connection.OperationExpectedRevision + 1);
 834            }
 835
 0836            await TryMarkRecoveryRequiredAsync(connection, tenantId, environmentId, "generation_publish_conflict");
 0837            return new ConnectionLifecycleResult(false, "generation_publish_conflict", connection.Revision + (connection
 838        }
 839
 7840        if (connection.OperationStatus == CredentialOperationStatus.RecoveryRequired && connection.Status == ConnectionS
 7841            !string.IsNullOrWhiteSpace(connection.PlannedSecretName) && !string.IsNullOrWhiteSpace(connection.PlannedGen
 842        {
 843            try
 844            {
 5845                var payload = await secrets.ResolveGenerationAsync(connection.PlannedSecretName, connection.Id, connecti
 5846                if (Deserialize(payload.Value) is { } envelope && IsValidEnvelope(envelope) &&
 5847                    (envelope.Kind == ConnectionCredentialKind.ApiKey
 5848                        ? await CanPromoteApiKeyRecoveryAsync(connection, cancellationToken)
 5849                        : envelope.Kind is null or ConnectionCredentialKind.OAuth) &&
 5850                    await store.TryPromoteRecoveryGenerationAsync(connectionId, tenantId, environmentId, connection.Revi
 5851                        connection.OperationId!, connection.OperationFence, envelope.Kind ?? ConnectionCredentialKind.OA
 5852                        envelope.AccessTokenExpiresAt, cancellationToken))
 853                {
 5854                    return new ConnectionLifecycleResult(true, null, connection.Revision + 1);
 855                }
 0856            }
 857            catch (KeyNotFoundException)
 858            {
 0859                var restored = await RestoreApiKeySourceIfPlanMissingAsync(connection, tenantId, environmentId, cancella
 0860                if (restored != null)
 861                {
 0862                    return restored;
 863                }
 864            }
 0865            catch (Exception)
 866            {
 867                // Missing or invalid planned material is not safe to publish; retain RecoveryRequired.
 0868            }
 869        }
 870
 2871        return new ConnectionLifecycleResult(false, "recovery_required", connection.Revision);
 24872    }
 873
 874    private async Task<bool> AuthorizeAsync(ClaimsPrincipal principal, ConnectionUseKind kind, string tenantId, string e
 217875        await authorizer.AuthorizeAsync(new ConnectionUseRequest(principal, kind, tenantId, environmentId, connectionId,
 876
 877    private async Task<ConnectionOffboardingOperationResult> QueueOffboardingOperationAsync(
 878        string tenantId,
 879        string environmentId,
 880        string connectionId,
 881        ConnectionOffboardingOperationKind kind,
 882        string? generationId,
 883        CancellationToken cancellationToken)
 884    {
 16885        if (kind == ConnectionOffboardingOperationKind.TokenPairRevocation && string.IsNullOrWhiteSpace(generationId) ||
 16886            kind == ConnectionOffboardingOperationKind.InstallationUninstall && generationId != null)
 887        {
 0888            return new ConnectionOffboardingOperationResult(false, "connection_unavailable", null, null, null);
 889        }
 890
 16891        using var tenantContext = PushTenant(tenantId);
 16892        var connection = await store.FindAsync(connectionId, tenantId, environmentId, cancellationToken);
 16893        if (connection is not { Status: ConnectionStatus.Disconnected } ||
 16894            connection.OperationStatus is not (CredentialOperationStatus.None or CredentialOperationStatus.Completed))
 895        {
 3896            return new ConnectionOffboardingOperationResult(false, "connection_unavailable", null, null, connection?.Rev
 897        }
 898
 13899        var operationId = GetOffboardingOperationId(tenantId, environmentId, connectionId, kind, generationId);
 13900        var existing = await store.FindOffboardingOperationAsync(operationId, tenantId, environmentId, connectionId, can
 13901        if (existing != null)
 902        {
 1903            return new ConnectionOffboardingOperationResult(true, null, operationId, existing.Status, connection.Revisio
 904        }
 905
 12906        if (kind == ConnectionOffboardingOperationKind.TokenPairRevocation)
 907        {
 908            try
 909            {
 8910                var payload = await secrets.ResolveGenerationAsync(ManagedSecretNames.ForGeneration(connectionId, genera
 8911                if (Deserialize(payload.Value) is not { Kind: null or ConnectionCredentialKind.OAuth })
 912                {
 0913                    throw new ConnectionUnavailableException();
 914                }
 8915            }
 0916            catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
 917            {
 0918                throw;
 919            }
 0920            catch (Exception)
 921            {
 0922                return new ConnectionOffboardingOperationResult(false, "connection_unavailable", null, null, connection.
 923            }
 924        }
 925
 12926        var now = timeProvider.GetUtcNow();
 12927        var operation = new ConnectionOffboardingOperation
 12928        {
 12929            Id = operationId,
 12930            TenantId = tenantId,
 12931            EnvironmentId = environmentId,
 12932            ConnectionId = connectionId,
 12933            ProviderId = connection.ProviderId,
 12934            ProviderAccountId = connection.ProviderAccountId,
 12935            Kind = kind,
 12936            GenerationId = generationId,
 12937            Status = ConnectionOffboardingOperationStatus.Pending,
 12938            Fence = 1,
 12939            CreatedAt = now,
 12940            UpdatedAt = now
 12941        };
 942
 12943        var queued = await store.TryQueueOffboardingOperationAsync(connection.Revision, operation, cancellationToken);
 12944        if (queued != null)
 945        {
 12946            return new ConnectionOffboardingOperationResult(true, null, operationId, queued.Status, connection.Revision 
 947        }
 948
 0949        var latestOperation = await store.FindOffboardingOperationAsync(operationId, tenantId, environmentId, connection
 0950        var latestConnection = await store.FindAsync(connectionId, tenantId, environmentId, cancellationToken);
 0951        return latestOperation != null && latestConnection != null
 0952            ? new ConnectionOffboardingOperationResult(true, null, operationId, latestOperation.Status, latestConnection
 0953            : new ConnectionOffboardingOperationResult(false, "connection_conflict", operationId, null, latestConnection
 16954    }
 955
 956    private async Task<ConnectionOffboardingOperationResult> GetOffboardingResultAsync(
 957        string operationId,
 958        string tenantId,
 959        string environmentId,
 960        string connectionId,
 961        bool accepted,
 962        string? safeErrorCode,
 963        CancellationToken cancellationToken)
 964    {
 12965        var operation = await store.FindOffboardingOperationAsync(operationId, tenantId, environmentId, connectionId, ca
 11966        var connection = await store.FindAsync(connectionId, tenantId, environmentId, cancellationToken);
 11967        return new ConnectionOffboardingOperationResult(accepted, safeErrorCode, operation?.Id, operation?.Status, conne
 11968    }
 969
 970    private async Task ReleaseOffboardingClaimAsync(
 971        ConnectionOffboardingOperation operation,
 972        string tenantId,
 973        string environmentId,
 974        string connectionId,
 975        string safeErrorCode)
 976    {
 0977        var now = timeProvider.GetUtcNow();
 978        try
 979        {
 0980            await store.TryReleaseOffboardingClaimAsync(operation.Id, tenantId, environmentId, connectionId, operation.F
 0981                now, now + TimeSpan.FromSeconds(30), safeErrorCode, CancellationToken.None);
 0982        }
 0983        catch (Exception)
 984        {
 985            // An expired claim is safe to retry; the provider-call state is never advanced here.
 0986        }
 0987    }
 988
 989    private async Task RecordOffboardingFailureAsync(
 990        ConnectionOffboardingOperation operation,
 991        string tenantId,
 992        string environmentId,
 993        string connectionId,
 994        ConnectionOffboardingOperationStatus status,
 995        string safeErrorCode,
 996        bool supportsStableIdempotency = true)
 997    {
 2998        var now = timeProvider.GetUtcNow();
 2999        DateTimeOffset? nextAttemptAt = status == ConnectionOffboardingOperationStatus.RetryScheduled ||
 21000                            status == ConnectionOffboardingOperationStatus.UnknownOutcome && supportsStableIdempotency
 21001            ? now + TimeSpan.FromSeconds(30)
 21002            : null;
 1003        try
 1004        {
 21005            await store.TryRecordOffboardingFailureAsync(operation.Id, tenantId, environmentId, connectionId, operation.
 21006                status, now, nextAttemptAt, safeErrorCode, CancellationToken.None);
 21007        }
 01008        catch (Exception)
 1009        {
 1010            // The durable provider-call lease expires and reconciliation retries with the same idempotency key.
 01011        }
 21012    }
 1013
 1014    private static string GetOffboardingOperationId(
 1015        string tenantId,
 1016        string environmentId,
 1017        string connectionId,
 1018        ConnectionOffboardingOperationKind kind,
 1019        string? generationId)
 1020    {
 301021        var canonicalIdentity = JsonSerializer.SerializeToUtf8Bytes(new string?[]
 301022        {
 301023            tenantId,
 301024            environmentId,
 301025            connectionId,
 301026            kind.ToString(),
 301027            generationId
 301028        }, JsonOptions);
 301029        return Convert.ToHexString(SHA256.HashData(canonicalIdentity)).ToLowerInvariant();
 1030    }
 1031
 1032    private static bool CanUseCurrentGeneration(IntegrationConnection? connection) =>
 581033        connection is { Status: ConnectionStatus.Active, OperationStatus: CredentialOperationStatus.None or CredentialOp
 581034        !string.IsNullOrWhiteSpace(connection.CurrentSecretName) && !string.IsNullOrWhiteSpace(connection.CurrentGenerat
 1035
 961036    private static ConnectionLifecycleMetadata ToMetadata(IntegrationConnection connection) => new(
 961037        connection.Id,
 961038        connection.ProviderId,
 961039        connection.ProviderAccountId,
 961040        connection.Status,
 961041        connection.Revision,
 961042        connection.CurrentGenerationId);
 1043
 2121044    private IDisposable PushTenant(string tenantId) => tenantAccessor.PushContext(new Tenant { Id = tenantId, Name = ten
 1045
 1046    private async Task<ConnectionLifecycleResult> ReleaseUnstartedRefreshAsync(IntegrationConnection connection, string 
 1047    {
 11048        await TryReleaseUnstartedRefreshAsync(connection, tenantId, environmentId, safeErrorCode);
 11049        var latest = await store.FindAsync(connection.Id, tenantId, environmentId, CancellationToken.None);
 11050        return latest == null
 11051            ? new ConnectionLifecycleResult(false, "connection_unavailable", null)
 11052            : new ConnectionLifecycleResult(false, latest.Status == ConnectionStatus.Active ? safeErrorCode : "connectio
 11053    }
 1054
 1055    private async Task TryReleaseUnstartedRefreshAsync(IntegrationConnection connection, string tenantId, string environ
 1056    {
 1057        try
 1058        {
 11059            await store.TryReleaseUnstartedRefreshAsync(connection.Id, tenantId, environmentId,
 11060                connection.OperationId!, connection.OperationFence, safeErrorCode, CancellationToken.None);
 11061        }
 01062        catch (Exception)
 1063        {
 1064            // The active credential remains the only published generation; reconciliation can release this claim after 
 01065        }
 11066    }
 1067
 521068    private static string Serialize(CredentialMaterial material) => JsonSerializer.Serialize(
 521069        new CredentialEnvelope(ConnectionCredentialKind.OAuth, material.AccessToken, material.RefreshToken, material.Acc
 1070
 191071    private static string SerializeApiKey(string apiKey) => JsonSerializer.Serialize(
 191072        new CredentialEnvelope(ConnectionCredentialKind.ApiKey, apiKey, null, null), JsonOptions);
 1073
 1074    private static CredentialEnvelope? Deserialize(string? json)
 1075    {
 931076        if (string.IsNullOrWhiteSpace(json))
 1077        {
 01078            return null;
 1079        }
 1080
 1081        try
 1082        {
 931083            return JsonSerializer.Deserialize<CredentialEnvelope>(json, JsonOptions);
 1084        }
 01085        catch (JsonException)
 1086        {
 01087            return null;
 1088        }
 931089    }
 1090
 1091    private async Task<bool> IsApiKeyGenerationAsync(IntegrationConnection connection, CancellationToken cancellationTok
 1092    {
 101093        if (string.IsNullOrWhiteSpace(connection.CurrentSecretName) || string.IsNullOrWhiteSpace(connection.CurrentGener
 1094        {
 01095            return false;
 1096        }
 1097
 1098        try
 1099        {
 101100            var payload = await secrets.ResolveGenerationAsync(connection.CurrentSecretName, connection.Id, connection.C
 101101            var envelope = Deserialize(payload.Value);
 101102            return envelope?.Kind == ConnectionCredentialKind.ApiKey && IsValidEnvelope(envelope);
 1103        }
 01104        catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
 1105        {
 01106            throw;
 1107        }
 01108        catch (Exception)
 1109        {
 01110            return false;
 1111        }
 101112    }
 1113
 1114    private async Task<bool> IsApiKeySourceGenerationAsync(IntegrationConnection connection, CancellationToken cancellat
 1115    {
 51116        var generationId = connection.OperationSourceGenerationId;
 51117        if (string.IsNullOrWhiteSpace(generationId) || connection.CurrentGenerationId != generationId)
 1118        {
 21119            return false;
 1120        }
 1121
 1122        try
 1123        {
 31124            var name = ManagedSecretNames.ForGeneration(connection.Id, generationId);
 31125            var payload = await secrets.ResolveGenerationAsync(name, connection.Id, generationId, cancellationToken);
 31126            var envelope = Deserialize(payload.Value);
 31127            return envelope?.Kind == ConnectionCredentialKind.ApiKey && IsValidEnvelope(envelope);
 1128        }
 01129        catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
 1130        {
 01131            throw;
 1132        }
 01133        catch (Exception)
 1134        {
 01135            return false;
 1136        }
 51137    }
 1138
 1139    private async Task<bool> CanPromoteApiKeyRecoveryAsync(IntegrationConnection connection, CancellationToken cancellat
 21140        await IsApiKeySourceGenerationAsync(connection, cancellationToken) ||
 21141        string.IsNullOrWhiteSpace(connection.CurrentGenerationId) && string.IsNullOrWhiteSpace(connection.OperationSourc
 1142
 1143    private static bool IsValidEnvelope(CredentialEnvelope? envelope)
 1144    {
 451145        if (envelope == null || string.IsNullOrWhiteSpace(envelope.AccessToken))
 1146        {
 01147            return false;
 1148        }
 1149
 451150        return envelope.Kind switch
 451151        {
 141152            null or ConnectionCredentialKind.OAuth => !string.IsNullOrWhiteSpace(envelope.RefreshToken) && envelope.Acce
 301153            ConnectionCredentialKind.ApiKey => envelope.RefreshToken is null && !envelope.AccessTokenExpiresAt.HasValue,
 11154            _ => false
 451155        };
 1156    }
 1157
 1158    private async Task<ConnectionLifecycleResult?> RestoreApiKeySourceIfPlanMissingAsync(
 1159        IntegrationConnection connection,
 1160        string tenantId,
 1161        string environmentId,
 1162        CancellationToken cancellationToken)
 1163    {
 11164        if (connection.Status != ConnectionStatus.RecoveryRequired ||
 11165            connection.OperationStatus != CredentialOperationStatus.RecoveryRequired ||
 11166            string.IsNullOrWhiteSpace(connection.OperationId) ||
 11167            string.IsNullOrWhiteSpace(connection.PlannedSecretName) ||
 11168            string.IsNullOrWhiteSpace(connection.PlannedGenerationId) ||
 11169            !await IsApiKeySourceGenerationAsync(connection, cancellationToken))
 1170        {
 01171            return null;
 1172        }
 1173
 1174        try
 1175        {
 11176            await secrets.ResolveGenerationAsync(connection.PlannedSecretName, connection.Id, connection.PlannedGenerati
 01177            return null;
 1178        }
 1179        catch (KeyNotFoundException)
 1180        {
 11181            var restored = await store.TryRestoreSourceGenerationAfterMissingPlanAsync(
 11182                connection.Id, tenantId, environmentId, connection.Revision, connection.OperationId,
 11183                connection.OperationFence, connection.OperationSourceGenerationId!, "api_key_replacement_not_staged", ca
 11184            if (!restored)
 1185            {
 01186                return null;
 1187            }
 1188
 11189            var latest = await store.FindAsync(connection.Id, tenantId, environmentId, cancellationToken);
 11190            return latest == null
 11191                ? new ConnectionLifecycleResult(false, "connection_unavailable", null)
 11192                : new ConnectionLifecycleResult(true, null, latest.Revision, connection.Id, ToMetadata(latest));
 1193        }
 11194    }
 1195
 1196    private async Task CleanupUnreferencedOrphanGenerationAsync(
 1197        string connectionId,
 1198        string tenantId,
 1199        string environmentId,
 1200        string generationId,
 1201        CancellationToken cancellationToken)
 1202    {
 11203        var connection = await store.FindAsync(connectionId, tenantId, environmentId, cancellationToken);
 11204        if (connection == null)
 1205        {
 01206            return;
 1207        }
 1208
 11209        var now = timeProvider.GetUtcNow();
 11210        var claim = await store.TryClaimGenerationCleanupAsync(connectionId, tenantId, environmentId,
 11211            connection.Revision, generationId, now, now + OperationLeaseDuration, cancellationToken);
 11212        if (claim == null)
 1213        {
 01214            return;
 1215        }
 1216
 1217        try
 1218        {
 11219            if (await secrets.DeleteGenerationAsync(ManagedSecretNames.ForGeneration(connectionId, generationId), connec
 1220            {
 11221                await store.CompleteGenerationCleanupAsync(connectionId, tenantId, environmentId, generationId, claim.Fe
 1222            }
 1223            else
 1224            {
 01225                await store.CancelGenerationCleanupAsync(connectionId, tenantId, environmentId, generationId, claim.Fenc
 1226            }
 11227        }
 01228        catch (Exception)
 1229        {
 1230            // Keep the durable cleanup claim for idempotent retry; do not affect the restored current generation.
 01231        }
 11232    }
 1233
 1234    private async Task<ConnectionLifecycleResult> RequireRecoveryAsync(IntegrationConnection connection, string tenantId
 1235    {
 31236        await TryMarkRecoveryRequiredAsync(connection, tenantId, environmentId, code);
 31237        return new ConnectionLifecycleResult(false, code, connection.OperationExpectedRevision);
 31238    }
 1239
 1240    private async Task TryMarkRecoveryRequiredAsync(IntegrationConnection connection, string tenantId, string environmen
 1241    {
 1242        try
 1243        {
 91244            await store.MarkRecoveryRequiredAsync(connection.Id, tenantId, environmentId, connection.OperationId!, conne
 91245        }
 01246        catch
 1247        {
 1248            // Durable provider-call intent remains for startup reconciliation; never replay automatically.
 01249        }
 91250    }
 1251
 9021252    private sealed record CredentialEnvelope(ConnectionCredentialKind? Kind, string? AccessToken, string? RefreshToken, 
 1253}

Methods/Properties

.ctor(Elsa.Connections.Contracts.IConnectionLifecycleStore,Elsa.Connections.Contracts.IConnectionUseAuthorizer,Elsa.Secrets.Contracts.IManagedSecretManager,System.TimeProvider,Elsa.Common.Multitenancy.ITenantAccessor,Elsa.Connections.Contracts.IConnectionCredentialProvider,Elsa.Connections.Contracts.IConnectionOffboardingProvider)
.cctor()
ConnectAsync()
ConnectApiKeyAsync()
ConnectCoreAsync()
ReplaceApiKeyAsync()
ResolveForUseAsync()
ResolveForUseAsync()
ResolveAuthorizedCredentialAsync()
DisconnectAsync()
RequestTokenRevocationAsync()
RequestInstallationUninstallAsync()
ReconcileOffboardingAsync()
RefreshAsync(System.Security.Claims.ClaimsPrincipal,System.String,System.String,System.String,System.Threading.CancellationToken)
RefreshAsync(System.String,System.String,System.String,System.Threading.CancellationToken)
RefreshCoreAsync()
CleanupGenerationAsync(System.Security.Claims.ClaimsPrincipal,System.String,System.String,System.String,System.String,System.Threading.CancellationToken)
CleanupGenerationAsync(System.String,System.String,System.String,System.String,System.Threading.CancellationToken)
CleanupGenerationCoreAsync()
ReconcileAsync()
AuthorizeAsync()
QueueOffboardingOperationAsync()
GetOffboardingResultAsync()
ReleaseOffboardingClaimAsync()
RecordOffboardingFailureAsync()
GetOffboardingOperationId(System.String,System.String,System.String,Elsa.Connections.Models.ConnectionOffboardingOperationKind,System.String)
CanUseCurrentGeneration(Elsa.Connections.Models.IntegrationConnection)
ToMetadata(Elsa.Connections.Models.IntegrationConnection)
PushTenant(System.String)
ReleaseUnstartedRefreshAsync()
TryReleaseUnstartedRefreshAsync()
Serialize(Elsa.Connections.Models.CredentialMaterial)
SerializeApiKey(System.String)
Deserialize(System.String)
IsApiKeyGenerationAsync()
IsApiKeySourceGenerationAsync()
CanPromoteApiKeyRecoveryAsync()
IsValidEnvelope(Elsa.Connections.Services.DefaultConnectionLifecycleService/CredentialEnvelope)
RestoreApiKeySourceIfPlanMissingAsync()
CleanupUnreferencedOrphanGenerationAsync()
RequireRecoveryAsync()
TryMarkRecoveryRequiredAsync()
get_Kind()