< Summary

Information
Class: AsiBackbone.Core.Outbox.GovernanceOutboxDrain
Assembly: AsiBackbone.Core
File(s): /home/runner/work/AsiBackbone/AsiBackbone/src/AsiBackbone.Core/Outbox/GovernanceOutboxDrain.cs
Line coverage
99%
Covered lines: 364
Uncovered lines: 3
Coverable lines: 367
Total lines: 702
Line coverage: 99.1%
Branch coverage
93%
Covered branches: 118
Total branches: 126
Branch coverage: 93.6%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)100%88100%
.cctor()100%11100%
DrainAsync()100%1010100%
DrainClaimedAsync()90%1010100%
ClaimPageAsync()100%22100%
DrainEntriesAsync()100%44100%
DrainClaimsAsync()87.5%8894.44%
MergeEntries(...)91.67%121293.33%
MergeClaims(...)100%1212100%
DrainEntryAsync()100%11100%
DrainClaimAsync()100%2296.55%
ReleaseClaimLeaseAsync()100%11100%
ApplyEmissionResultAsync()90%2020100%
ApplyEmissionResultAsync()90%2020100%
ApplyFailureAsync()100%22100%
ApplyClaimFailureAsync()100%22100%
ShouldDeadLetter(...)100%22100%
ShouldDeadLetterForClaimAttempts(...)100%22100%
CreateMaxRetryError(...)100%11100%
LogEmissionException(...)100%11100%
LogClaimAttemptsExceeded(...)100%11100%
CreateExceptionError(...)100%11100%
GetRetryUtc(...)100%11100%
GetDeferredUtc(...)100%11100%
ResolveOptions(...)100%66100%
ResolveEmitterProvider(...)75%44100%

File(s)

/home/runner/work/AsiBackbone/AsiBackbone/src/AsiBackbone.Core/Outbox/GovernanceOutboxDrain.cs

#LineLine coverage
 1using AsiBackbone.Core.Emissions;
 2using Microsoft.Extensions.Logging;
 3using Microsoft.Extensions.Logging.Abstractions;
 4using Microsoft.Extensions.Options;
 5
 6namespace AsiBackbone.Core.Outbox;
 7
 8/// <summary>
 9/// Drains provider-neutral outbox entries through a configured governance emitter.
 10/// </summary>
 11/// <remarks>
 12/// This drain path is provider-neutral. It is suitable for tests, samples, local validation, and host-owned workers tha
 13/// </remarks>
 14/// <remarks>
 15/// Initializes a new instance of the <see cref="GovernanceOutboxDrain" /> class.
 16/// </remarks>
 17/// <param name="outboxStore">The provider-neutral outbox store.</param>
 18/// <param name="emitter">The provider-neutral governance emitter.</param>
 19/// <param name="logger">The logger used to record local operational diagnostics for drain failures.</param>
 20/// <param name="outboxOptions">The provider-neutral retry, poison-message, and claim options used by the drain.</param>
 21/// <param name="timeProvider">The clock used for the drain time and per-page claim-lease time when <see cref="DrainAsyn
 9522public sealed class GovernanceOutboxDrain(
 9523    IGovernanceOutboxStore outboxStore,
 9524    IGovernanceEmitter emitter,
 9525    ILogger<GovernanceOutboxDrain>? logger = null,
 9526    IOptions<GovernanceOutboxOptions>? outboxOptions = null,
 9527    TimeProvider? timeProvider = null)
 28{
 229    private static readonly Action<ILogger, string, int, string, DateTimeOffset, string?, string?, Exception?> LogGovern
 230        LogLevel.Warning,
 231        new EventId(19701, nameof(LogGovernanceEmissionException)),
 232        "Governance outbox emission threw an exception for outbox entry {OutboxEntryId} on attempt {AttemptCount}. Emitt
 33
 234    private static readonly Action<ILogger, string, int, int, string?, Exception?> LogGovernanceClaimAttemptsExceeded = 
 235        LogLevel.Warning,
 236        new EventId(19702, nameof(LogGovernanceClaimAttemptsExceeded)),
 237        "Governance outbox entry {OutboxEntryId} was claimed {ClaimAttemptCount} times without reaching a terminal state
 38
 239    private static readonly Action<ILogger, string, string, Exception?> LogGovernanceClaimReleaseFailure = LoggerMessage
 240        LogLevel.Warning,
 241        new EventId(19703, nameof(LogGovernanceClaimReleaseFailure)),
 242        "Governance outbox claim release failed during cancellation for outbox entry {OutboxEntryId} owned by worker {Cl
 43
 9944    private readonly IGovernanceOutboxStore outboxStore = outboxStore ?? throw new ArgumentNullException(nameof(outboxSt
 9845    private readonly IGovernanceEmitter emitter = emitter ?? throw new ArgumentNullException(nameof(emitter));
 9746    private readonly ILogger<GovernanceOutboxDrain> logger = logger ?? NullLogger<GovernanceOutboxDrain>.Instance;
 9747    private readonly GovernanceOutboxOptions retryOptions = ResolveOptions(outboxOptions);
 9548    private readonly TimeProvider timeProvider = timeProvider ?? TimeProvider.System;
 49
 50    /// <summary>
 51    /// Drains pending and retry-ready outbox entries through the configured emitter.
 52    /// </summary>
 53    /// <param name="utcNow">The UTC timestamp used for retry-ready checks. When omitted, the configured <see cref="Time
 54    /// <param name="maxCount">The maximum number of entries to drain.</param>
 55    /// <param name="cancellationToken">A cancellation token.</param>
 56    /// <returns>The updated outbox entries that were attempted by the drain.</returns>
 57    public async ValueTask<IReadOnlyList<GovernanceOutboxEntry>> DrainAsync(
 58        DateTimeOffset? utcNow = null,
 59        int maxCount = 100,
 60        CancellationToken cancellationToken = default)
 61    {
 10262        if (maxCount <= 0)
 63        {
 164            throw new ArgumentOutOfRangeException(nameof(maxCount), maxCount, "Maximum count must be greater than zero."
 65        }
 66
 10167        cancellationToken.ThrowIfCancellationRequested();
 68
 10169        DateTimeOffset drainUtc = (utcNow ?? timeProvider.GetUtcNow()).ToUniversalTime();
 70
 10171        if (retryOptions.UseClaimLeases)
 72        {
 73            // A caller-supplied timestamp is honored for every page so deterministic tests stay deterministic;
 74            // otherwise each page is leased from a fresh reading taken when that page is claimed.
 4875            Func<DateTimeOffset> claimClock = utcNow.HasValue
 4876                ? () => drainUtc
 4877                : timeProvider.GetUtcNow;
 78
 4879            return await DrainClaimedAsync(claimClock, maxCount, cancellationToken).ConfigureAwait(false);
 80        }
 81
 5382        IReadOnlyList<GovernanceOutboxEntry> pendingEntries = await outboxStore
 5383            .FindPendingAsync(maxCount, cancellationToken)
 5384            .ConfigureAwait(false);
 85
 4986        if (pendingEntries.Count >= maxCount)
 87        {
 588            return await DrainEntriesAsync(pendingEntries, drainUtc, cancellationToken).ConfigureAwait(false);
 89        }
 90
 4491        IReadOnlyList<GovernanceOutboxEntry> retryReadyEntries = await outboxStore
 4492            .FindRetryReadyAsync(drainUtc, maxCount - pendingEntries.Count, cancellationToken)
 4493            .ConfigureAwait(false);
 94
 4495        IReadOnlyList<GovernanceOutboxEntry> entriesToDrain = MergeEntries(pendingEntries, retryReadyEntries, maxCount);
 4496        return await DrainEntriesAsync(entriesToDrain, drainUtc, cancellationToken).ConfigureAwait(false);
 8897    }
 98
 99    private async ValueTask<IReadOnlyList<GovernanceOutboxEntry>> DrainClaimedAsync(
 100        Func<DateTimeOffset> claimClock,
 101        int maxCount,
 102        CancellationToken cancellationToken)
 103    {
 48104        if (outboxStore is not IGovernanceOutboxClaimStore claimStore)
 105        {
 3106            throw new InvalidOperationException("Claim leases are enabled, but the configured outbox store does not impl
 107        }
 108
 45109        string workerId = retryOptions.ClaimWorkerId ?? throw new InvalidOperationException("ClaimWorkerId is required w
 110
 45111        List<GovernanceOutboxEntry> drainedEntries = [];
 45112        int remainingCount = maxCount;
 113
 114        // Entries are claimed a page at a time rather than leasing the whole batch at once, so a slow emitter
 115        // cannot exhaust a single lease across the batch and leave later entries reclaimable while in flight.
 66116        while (remainingCount > 0)
 117        {
 56118            cancellationToken.ThrowIfCancellationRequested();
 119
 56120            int pageSize = Math.Min(retryOptions.ClaimPageSize, remainingCount);
 56121            DateTimeOffset pageUtc = claimClock().ToUniversalTime();
 122
 56123            IReadOnlyList<GovernanceOutboxClaim> pageClaims = await ClaimPageAsync(
 56124                claimStore,
 56125                workerId,
 56126                pageUtc,
 56127                pageSize,
 56128                cancellationToken)
 56129                .ConfigureAwait(false);
 130
 56131            if (pageClaims.Count == 0)
 132            {
 133                break;
 134            }
 135
 52136            IReadOnlyList<GovernanceOutboxEntry> pageEntries = await DrainClaimsAsync(
 52137                claimStore,
 52138                pageClaims,
 52139                pageUtc,
 52140                cancellationToken)
 52141                .ConfigureAwait(false);
 142
 47143            drainedEntries.AddRange(pageEntries);
 47144            remainingCount -= pageClaims.Count;
 145
 47146            if (pageClaims.Count < pageSize)
 147            {
 148                break;
 149            }
 21150        }
 151
 40152        return drainedEntries;
 40153    }
 154
 155    private async ValueTask<IReadOnlyList<GovernanceOutboxClaim>> ClaimPageAsync(
 156        IGovernanceOutboxClaimStore claimStore,
 157        string workerId,
 158        DateTimeOffset pageUtc,
 159        int pageSize,
 160        CancellationToken cancellationToken)
 161    {
 56162        var pendingRequest = GovernanceOutboxClaimRequest.Create(
 56163            workerId,
 56164            pageUtc,
 56165            retryOptions.ClaimLeaseDuration,
 56166            pageSize);
 56167        IReadOnlyList<GovernanceOutboxClaim> pendingClaims = await claimStore
 56168            .ClaimPendingAsync(pendingRequest, cancellationToken)
 56169            .ConfigureAwait(false);
 170
 56171        if (pendingClaims.Count >= pageSize)
 172        {
 25173            return pendingClaims;
 174        }
 175
 31176        var retryRequest = GovernanceOutboxClaimRequest.Create(
 31177            workerId,
 31178            pageUtc,
 31179            retryOptions.ClaimLeaseDuration,
 31180            pageSize - pendingClaims.Count);
 31181        IReadOnlyList<GovernanceOutboxClaim> retryReadyClaims = await claimStore
 31182            .ClaimRetryReadyAsync(retryRequest, cancellationToken)
 31183            .ConfigureAwait(false);
 184
 31185        return MergeClaims(pendingClaims, retryReadyClaims, pageSize);
 56186    }
 187
 188    private async ValueTask<IReadOnlyList<GovernanceOutboxEntry>> DrainEntriesAsync(
 189        IReadOnlyList<GovernanceOutboxEntry> entriesToDrain,
 190        DateTimeOffset drainUtc,
 191        CancellationToken cancellationToken)
 192    {
 49193        if (entriesToDrain.Count == 0)
 194        {
 12195            return Array.Empty<GovernanceOutboxEntry>();
 196        }
 197
 37198        List<GovernanceOutboxEntry> updatedEntries = new(entriesToDrain.Count);
 199
 151200        foreach (GovernanceOutboxEntry entry in entriesToDrain)
 201        {
 39202            cancellationToken.ThrowIfCancellationRequested();
 39203            GovernanceOutboxEntry updatedEntry = await DrainEntryAsync(entry, drainUtc, cancellationToken).ConfigureAwai
 38204            updatedEntries.Add(updatedEntry);
 205        }
 206
 36207        return updatedEntries;
 48208    }
 209
 210    private async ValueTask<IReadOnlyList<GovernanceOutboxEntry>> DrainClaimsAsync(
 211        IGovernanceOutboxClaimStore claimStore,
 212        IReadOnlyList<GovernanceOutboxClaim> claimsToDrain,
 213        DateTimeOffset drainUtc,
 214        CancellationToken cancellationToken)
 215    {
 52216        if (claimsToDrain.Count == 0)
 217        {
 0218            return Array.Empty<GovernanceOutboxEntry>();
 219        }
 220
 52221        List<GovernanceOutboxEntry> updatedEntries = new(claimsToDrain.Count);
 222
 396223        for (int claimIndex = 0; claimIndex < claimsToDrain.Count; claimIndex++)
 224        {
 151225            GovernanceOutboxClaim claim = claimsToDrain[claimIndex];
 151226            bool drainAttempted = false;
 227
 228            try
 229            {
 151230                cancellationToken.ThrowIfCancellationRequested();
 150231                drainAttempted = true;
 150232                GovernanceOutboxEntry updatedEntry = await DrainClaimAsync(claimStore, claim, drainUtc, cancellationToke
 146233                updatedEntries.Add(updatedEntry);
 146234            }
 5235            catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
 236            {
 5237                int releaseStartIndex = drainAttempted ? claimIndex + 1 : claimIndex;
 12238                for (int releaseIndex = releaseStartIndex; releaseIndex < claimsToDrain.Count; releaseIndex++)
 239                {
 1240                    await ReleaseClaimLeaseAsync(claimStore, claimsToDrain[releaseIndex]).ConfigureAwait(false);
 241                }
 242
 5243                throw;
 244            }
 245        }
 246
 47247        return updatedEntries;
 47248    }
 249
 250    private static IReadOnlyList<GovernanceOutboxEntry> MergeEntries(
 251        IReadOnlyList<GovernanceOutboxEntry> pendingEntries,
 252        IReadOnlyList<GovernanceOutboxEntry> retryReadyEntries,
 253        int maxCount)
 254    {
 44255        if (pendingEntries.Count == 0)
 256        {
 33257            return retryReadyEntries;
 258        }
 259
 11260        if (retryReadyEntries.Count == 0)
 261        {
 10262            return pendingEntries;
 263        }
 264
 1265        var entriesToDrain = new List<GovernanceOutboxEntry>(Math.Min(maxCount, pendingEntries.Count + retryReadyEntries
 1266        var existingEntryIds = new HashSet<string>(pendingEntries.Count + retryReadyEntries.Count, StringComparer.Ordina
 267
 4268        foreach (GovernanceOutboxEntry pendingEntry in pendingEntries)
 269        {
 1270            _ = existingEntryIds.Add(pendingEntry.OutboxEntryId);
 1271            entriesToDrain.Add(pendingEntry);
 272        }
 273
 6274        foreach (GovernanceOutboxEntry retryReadyEntry in retryReadyEntries)
 275        {
 2276            if (entriesToDrain.Count >= maxCount)
 277            {
 0278                break;
 279            }
 280
 2281            if (existingEntryIds.Add(retryReadyEntry.OutboxEntryId))
 282            {
 1283                entriesToDrain.Add(retryReadyEntry);
 284            }
 285        }
 286
 1287        return entriesToDrain;
 288    }
 289
 290    private static IReadOnlyList<GovernanceOutboxClaim> MergeClaims(
 291        IReadOnlyList<GovernanceOutboxClaim> pendingClaims,
 292        IReadOnlyList<GovernanceOutboxClaim> retryReadyClaims,
 293        int maxCount)
 294    {
 38295        if (pendingClaims.Count == 0)
 296        {
 9297            return retryReadyClaims;
 298        }
 299
 29300        if (retryReadyClaims.Count == 0)
 301        {
 24302            return pendingClaims;
 303        }
 304
 5305        var claimsToDrain = new List<GovernanceOutboxClaim>(Math.Min(maxCount, pendingClaims.Count + retryReadyClaims.Co
 5306        var existingEntryIds = new HashSet<string>(pendingClaims.Count + retryReadyClaims.Count, StringComparer.Ordinal)
 307
 22308        foreach (GovernanceOutboxClaim pendingClaim in pendingClaims)
 309        {
 6310            _ = existingEntryIds.Add(pendingClaim.OutboxEntryId);
 6311            claimsToDrain.Add(pendingClaim);
 312        }
 313
 39314        foreach (GovernanceOutboxClaim retryReadyClaim in retryReadyClaims)
 315        {
 15316            if (claimsToDrain.Count >= maxCount)
 317            {
 1318                break;
 319            }
 320
 14321            if (existingEntryIds.Add(retryReadyClaim.OutboxEntryId))
 322            {
 8323                claimsToDrain.Add(retryReadyClaim);
 324            }
 325        }
 326
 5327        return claimsToDrain;
 328    }
 329
 330    private async ValueTask<GovernanceOutboxEntry> DrainEntryAsync(
 331        GovernanceOutboxEntry entry,
 332        DateTimeOffset drainUtc,
 333        CancellationToken cancellationToken)
 334    {
 335        GovernanceEmissionResult result;
 336
 337        try
 338        {
 39339            result = await emitter.EmitAsync(entry.Envelope, cancellationToken).ConfigureAwait(false);
 36340        }
 1341        catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
 342        {
 1343            throw;
 344        }
 2345        catch (Exception ex)
 346        {
 2347            DateTimeOffset nextRetryUtc = GetRetryUtc(drainUtc);
 2348            LogEmissionException(entry, nextRetryUtc, ex);
 2349            GovernanceEmissionError governanceEmissionError = CreateExceptionError(ex);
 350
 2351            return await ApplyFailureAsync(
 2352                entry,
 2353                governanceEmissionError,
 2354                nextRetryUtc,
 2355                cancellationToken)
 2356                .ConfigureAwait(false);
 357        }
 358
 36359        return await ApplyEmissionResultAsync(entry, result, drainUtc, cancellationToken).ConfigureAwait(false);
 38360    }
 361
 362    private async ValueTask<GovernanceOutboxEntry> DrainClaimAsync(
 363        IGovernanceOutboxClaimStore claimStore,
 364        GovernanceOutboxClaim claim,
 365        DateTimeOffset drainUtc,
 366        CancellationToken cancellationToken)
 367    {
 368        try
 369        {
 370            // An emitter that hangs or is killed mid-emission leaves the entry claimed but never failed, so its retry
 371            // count does not advance and the retry-based poison-message policy never fires. The claim count does
 372            // advance on every reclaim, so it is the only signal that bounds that loop. Checked before emission so a
 373            // repeatedly reclaimed entry is not handed to the emitter again.
 150374            if (ShouldDeadLetterForClaimAttempts(claim.Entry))
 375            {
 2376                var claimExhaustedError = GovernanceEmissionError.Create(
 2377                    retryOptions.MaxClaimAttemptsReasonCode,
 2378                    retryOptions.MaxClaimAttemptsReasonMessage);
 379
 2380                LogClaimAttemptsExceeded(claim.Entry);
 381
 2382                return await claimStore.MarkClaimDeadLetteredAsync(
 2383                    claim,
 2384                    claimExhaustedError,
 2385                    retryOptions.MaxClaimAttemptsReasonMessage,
 2386                    cancellationToken)
 2387                    .ConfigureAwait(false);
 388            }
 389
 148390            GovernanceEmissionResult result = await emitter.EmitAsync(claim.Entry.Envelope, cancellationToken).Configure
 142391            return await ApplyEmissionResultAsync(claimStore, claim, result, drainUtc, cancellationToken).ConfigureAwait
 392        }
 4393        catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
 394        {
 4395            await ReleaseClaimLeaseAsync(claimStore, claim).ConfigureAwait(false);
 4396            throw;
 0397        }
 4398        catch (Exception ex)
 399        {
 4400            DateTimeOffset nextRetryUtc = GetRetryUtc(drainUtc);
 4401            LogEmissionException(claim.Entry, nextRetryUtc, ex);
 4402            GovernanceEmissionError governanceEmissionError = CreateExceptionError(ex);
 403
 4404            return await ApplyClaimFailureAsync(
 4405                claimStore,
 4406                claim,
 4407                governanceEmissionError,
 4408                nextRetryUtc,
 4409                cancellationToken)
 4410                .ConfigureAwait(false);
 411        }
 146412    }
 413
 414    private async ValueTask ReleaseClaimLeaseAsync(
 415        IGovernanceOutboxClaimStore claimStore,
 416        GovernanceOutboxClaim claim)
 417    {
 418        try
 419        {
 5420            _ = await claimStore.ReleaseClaimAsync(
 5421                claim,
 5422                reason: "drain canceled; releasing active claim",
 5423                cancellationToken: CancellationToken.None)
 5424                .ConfigureAwait(false);
 4425        }
 1426        catch (Exception ex)
 427        {
 1428            LogGovernanceClaimReleaseFailure(logger, claim.OutboxEntryId, claim.WorkerId, ex);
 429
 430            // Best-effort release must not block shutdown or surface a secondary failure while the drain is already
 431            // aborting. The caller is exiting and claim release is idempotent, so a transient storage failure should no
 432            // mask the original cancellation or leave the remaining page of leases stuck until the normal lease expiry.
 1433        }
 5434    }
 435
 436    private async ValueTask<GovernanceOutboxEntry> ApplyEmissionResultAsync(
 437        GovernanceOutboxEntry entry,
 438        GovernanceEmissionResult result,
 439        DateTimeOffset drainUtc,
 440        CancellationToken cancellationToken)
 441    {
 36442        ArgumentNullException.ThrowIfNull(result);
 443
 36444        if (result.IsSuccess)
 445        {
 10446            return await outboxStore.MarkDeliveredAsync(entry.OutboxEntryId, result, cancellationToken).ConfigureAwait(f
 447        }
 448
 26449        if (result.Status is GovernanceEmissionStatus.DeadLettered)
 450        {
 1451            GovernanceEmissionError governanceEmissionError = result.Error ?? GovernanceEmissionError.Create(
 1452                "emission.deadlettered",
 1453                "Governance emission returned a dead-lettered result.",
 1454                providerName: result.ProviderName);
 455
 1456            return await outboxStore.MarkDeadLetteredAsync(
 1457                entry.OutboxEntryId,
 1458                governanceEmissionError,
 1459                governanceEmissionError.Message,
 1460                cancellationToken)
 1461                .ConfigureAwait(false);
 462        }
 463
 25464        if (result.Status is GovernanceEmissionStatus.Deferred or GovernanceEmissionStatus.Pending)
 465        {
 3466            GovernanceEmissionError? governanceEmissionError = result.Error ?? (result.Status is GovernanceEmissionStatu
 3467                ? GovernanceEmissionError.Create(
 3468                    "emission.pending",
 3469                    "Governance emission remained pending after the outbox drain attempt.",
 3470                    isRetryable: true,
 3471                    providerName: result.ProviderName)
 3472                : null);
 473
 3474            GovernanceOutboxEntry deferredEntry = entry.MarkDeferred(
 3475                governanceEmissionError,
 3476                result.RetryAfterUtc ?? GetDeferredUtc(drainUtc),
 3477                drainUtc);
 478
 3479            return await outboxStore.SaveAsync(deferredEntry, cancellationToken).ConfigureAwait(false);
 480        }
 481
 22482        GovernanceEmissionError failure = result.Error ?? GovernanceEmissionError.Create(
 22483            "emission.failed",
 22484            "Governance emission returned a failed result without provider-neutral error details.",
 22485            isRetryable: result.ShouldRetry,
 22486            providerName: result.ProviderName);
 487
 22488        return await ApplyFailureAsync(
 22489            entry,
 22490            failure,
 22491            result.RetryAfterUtc,
 22492            cancellationToken)
 22493            .ConfigureAwait(false);
 36494    }
 495
 496    private async ValueTask<GovernanceOutboxEntry> ApplyEmissionResultAsync(
 497        IGovernanceOutboxClaimStore claimStore,
 498        GovernanceOutboxClaim claim,
 499        GovernanceEmissionResult result,
 500        DateTimeOffset drainUtc,
 501        CancellationToken cancellationToken)
 502    {
 142503        ArgumentNullException.ThrowIfNull(result);
 504
 142505        if (result.IsSuccess)
 506        {
 121507            return await claimStore.MarkClaimDeliveredAsync(claim, result, cancellationToken).ConfigureAwait(false);
 508        }
 509
 21510        if (result.Status is GovernanceEmissionStatus.DeadLettered)
 511        {
 3512            GovernanceEmissionError governanceEmissionError = result.Error ?? GovernanceEmissionError.Create(
 3513                "emission.deadlettered",
 3514                "Governance emission returned a dead-lettered result.",
 3515                providerName: result.ProviderName);
 516
 3517            return await claimStore.MarkClaimDeadLetteredAsync(
 3518                claim,
 3519                governanceEmissionError,
 3520                governanceEmissionError.Message,
 3521                cancellationToken)
 3522                .ConfigureAwait(false);
 523        }
 524
 18525        if (result.Status is GovernanceEmissionStatus.Deferred or GovernanceEmissionStatus.Pending)
 526        {
 8527            GovernanceEmissionError? governanceEmissionError = result.Error ?? (result.Status is GovernanceEmissionStatu
 8528                ? GovernanceEmissionError.Create(
 8529                    "emission.pending",
 8530                    "Governance emission remained pending after the outbox drain attempt.",
 8531                    isRetryable: true,
 8532                    providerName: result.ProviderName)
 8533                : null);
 534
 8535            GovernanceOutboxEntry deferredEntry = claim.Entry.MarkDeferred(
 8536                governanceEmissionError,
 8537                result.RetryAfterUtc ?? GetDeferredUtc(drainUtc),
 8538                drainUtc);
 539
 8540            return await claimStore.SaveClaimAsync(claim, deferredEntry, cancellationToken).ConfigureAwait(false);
 541        }
 542
 10543        GovernanceEmissionError failure = result.Error ?? GovernanceEmissionError.Create(
 10544            "emission.failed",
 10545            "Governance emission returned a failed result without provider-neutral error details.",
 10546            isRetryable: result.ShouldRetry,
 10547            providerName: result.ProviderName);
 548
 10549        return await ApplyClaimFailureAsync(
 10550            claimStore,
 10551            claim,
 10552            failure,
 10553            result.RetryAfterUtc,
 10554            cancellationToken)
 10555            .ConfigureAwait(false);
 141556    }
 557
 558    private async ValueTask<GovernanceOutboxEntry> ApplyFailureAsync(
 559        GovernanceOutboxEntry entry,
 560        GovernanceEmissionError failure,
 561        DateTimeOffset? nextRetryUtc,
 562        CancellationToken cancellationToken)
 563    {
 24564        if (ShouldDeadLetter(entry))
 565        {
 2566            GovernanceEmissionError deadLetterError = CreateMaxRetryError(failure);
 2567            return await outboxStore.MarkDeadLetteredAsync(
 2568                entry.OutboxEntryId,
 2569                deadLetterError,
 2570                retryOptions.DeadLetterReasonMessage,
 2571                cancellationToken)
 2572                .ConfigureAwait(false);
 573        }
 574
 22575        return await outboxStore.MarkFailedAsync(
 22576            entry.OutboxEntryId,
 22577            failure,
 22578            nextRetryUtc,
 22579            cancellationToken)
 22580            .ConfigureAwait(false);
 24581    }
 582
 583    private async ValueTask<GovernanceOutboxEntry> ApplyClaimFailureAsync(
 584        IGovernanceOutboxClaimStore claimStore,
 585        GovernanceOutboxClaim claim,
 586        GovernanceEmissionError failure,
 587        DateTimeOffset? nextRetryUtc,
 588        CancellationToken cancellationToken)
 589    {
 14590        if (ShouldDeadLetter(claim.Entry))
 591        {
 1592            GovernanceEmissionError deadLetterError = CreateMaxRetryError(failure);
 1593            return await claimStore.MarkClaimDeadLetteredAsync(
 1594                claim,
 1595                deadLetterError,
 1596                retryOptions.DeadLetterReasonMessage,
 1597                cancellationToken)
 1598                .ConfigureAwait(false);
 599        }
 600
 13601        return await claimStore.MarkClaimFailedAsync(
 13602            claim,
 13603            failure,
 13604            nextRetryUtc,
 13605            cancellationToken)
 13606            .ConfigureAwait(false);
 14607    }
 608
 609    private bool ShouldDeadLetter(GovernanceOutboxEntry entry)
 610    {
 38611        return retryOptions.DeadLetterOnMaxRetryAttempts
 38612            && entry.RetryCount + 1 >= retryOptions.MaxRetryAttempts;
 613    }
 614
 615    private bool ShouldDeadLetterForClaimAttempts(GovernanceOutboxEntry entry)
 616    {
 617        // The store has already stamped this claim, so ClaimAttemptCount includes the current attempt.
 150618        return retryOptions.DeadLetterOnMaxClaimAttempts
 150619            && entry.ClaimAttemptCount > retryOptions.MaxClaimAttempts;
 620    }
 621
 622    private GovernanceEmissionError CreateMaxRetryError(GovernanceEmissionError failure)
 623    {
 3624        return GovernanceEmissionError.Create(
 3625            retryOptions.DeadLetterReasonCode,
 3626            retryOptions.DeadLetterReasonMessage,
 3627            providerName: failure.ProviderName,
 3628            providerErrorCode: failure.Code);
 629    }
 630
 631    private void LogEmissionException(GovernanceOutboxEntry entry, DateTimeOffset nextRetryUtc, Exception exception)
 632    {
 6633        LogGovernanceEmissionException(
 6634            logger,
 6635            entry.OutboxEntryId,
 6636            entry.RetryCount + 1,
 6637            ResolveEmitterProvider(entry),
 6638            nextRetryUtc,
 6639            entry.Envelope.CorrelationId,
 6640            entry.Envelope.DecisionReceiptId,
 6641            exception);
 6642    }
 643
 644    private void LogClaimAttemptsExceeded(GovernanceOutboxEntry entry)
 645    {
 2646        LogGovernanceClaimAttemptsExceeded(
 2647            logger,
 2648            entry.OutboxEntryId,
 2649            entry.ClaimAttemptCount,
 2650            retryOptions.MaxClaimAttempts,
 2651            entry.Envelope.CorrelationId,
 2652            null);
 2653    }
 654
 655    private static GovernanceEmissionError CreateExceptionError(Exception exception)
 656    {
 6657        return GovernanceEmissionError.Create(
 6658            "emission.exception",
 6659            $"Governance emission threw {exception.GetType().Name} during outbox drain.",
 6660            isRetryable: true,
 6661            providerErrorCode: exception.GetType().FullName);
 662    }
 663
 664    private DateTimeOffset GetRetryUtc(DateTimeOffset drainUtc)
 665    {
 6666        return drainUtc.Add(retryOptions.RetryDelay);
 667    }
 668
 669    private DateTimeOffset GetDeferredUtc(DateTimeOffset drainUtc)
 670    {
 6671        return drainUtc.Add(retryOptions.DeferredDelay);
 672    }
 673
 674    private static GovernanceOutboxOptions ResolveOptions(IOptions<GovernanceOutboxOptions>? options)
 675    {
 97676        GovernanceOutboxOptions resolved = options?.Value ?? new GovernanceOutboxOptions();
 97677        resolved.Validate();
 678
 95679        return new GovernanceOutboxOptions
 95680        {
 95681            RetryDelay = resolved.RetryDelay,
 95682            DeferredDelay = resolved.DeferredDelay,
 95683            MaxRetryAttempts = resolved.MaxRetryAttempts,
 95684            DeadLetterOnMaxRetryAttempts = resolved.DeadLetterOnMaxRetryAttempts,
 95685            DeadLetterReasonCode = resolved.DeadLetterReasonCode.Trim(),
 95686            DeadLetterReasonMessage = resolved.DeadLetterReasonMessage.Trim(),
 95687            UseClaimLeases = resolved.UseClaimLeases,
 95688            ClaimWorkerId = string.IsNullOrWhiteSpace(resolved.ClaimWorkerId) ? null : resolved.ClaimWorkerId.Trim(),
 95689            ClaimLeaseDuration = resolved.ClaimLeaseDuration,
 95690            ClaimPageSize = resolved.ClaimPageSize,
 95691            MaxClaimAttempts = resolved.MaxClaimAttempts,
 95692            DeadLetterOnMaxClaimAttempts = resolved.DeadLetterOnMaxClaimAttempts,
 95693            MaxClaimAttemptsReasonCode = resolved.MaxClaimAttemptsReasonCode.Trim(),
 95694            MaxClaimAttemptsReasonMessage = resolved.MaxClaimAttemptsReasonMessage.Trim()
 95695        };
 696    }
 697
 698    private static string ResolveEmitterProvider(GovernanceOutboxEntry entry)
 699    {
 6700        return entry.Envelope.EmitterProvider ?? entry.ProviderName ?? "unspecified";
 701    }
 702}

Methods/Properties

.ctor(AsiBackbone.Core.Outbox.IGovernanceOutboxStore,AsiBackbone.Core.Emissions.IGovernanceEmitter,Microsoft.Extensions.Logging.ILogger`1<AsiBackbone.Core.Outbox.GovernanceOutboxDrain>,Microsoft.Extensions.Options.IOptions`1<AsiBackbone.Core.Outbox.GovernanceOutboxOptions>,System.TimeProvider)
.cctor()
DrainAsync()
DrainClaimedAsync()
ClaimPageAsync()
DrainEntriesAsync()
DrainClaimsAsync()
MergeEntries(System.Collections.Generic.IReadOnlyList`1<AsiBackbone.Core.Outbox.GovernanceOutboxEntry>,System.Collections.Generic.IReadOnlyList`1<AsiBackbone.Core.Outbox.GovernanceOutboxEntry>,System.Int32)
MergeClaims(System.Collections.Generic.IReadOnlyList`1<AsiBackbone.Core.Outbox.GovernanceOutboxClaim>,System.Collections.Generic.IReadOnlyList`1<AsiBackbone.Core.Outbox.GovernanceOutboxClaim>,System.Int32)
DrainEntryAsync()
DrainClaimAsync()
ReleaseClaimLeaseAsync()
ApplyEmissionResultAsync()
ApplyEmissionResultAsync()
ApplyFailureAsync()
ApplyClaimFailureAsync()
ShouldDeadLetter(AsiBackbone.Core.Outbox.GovernanceOutboxEntry)
ShouldDeadLetterForClaimAttempts(AsiBackbone.Core.Outbox.GovernanceOutboxEntry)
CreateMaxRetryError(AsiBackbone.Core.Emissions.GovernanceEmissionError)
LogEmissionException(AsiBackbone.Core.Outbox.GovernanceOutboxEntry,System.DateTimeOffset,System.Exception)
LogClaimAttemptsExceeded(AsiBackbone.Core.Outbox.GovernanceOutboxEntry)
CreateExceptionError(System.Exception)
GetRetryUtc(System.DateTimeOffset)
GetDeferredUtc(System.DateTimeOffset)
ResolveOptions(Microsoft.Extensions.Options.IOptions`1<AsiBackbone.Core.Outbox.GovernanceOutboxOptions>)
ResolveEmitterProvider(AsiBackbone.Core.Outbox.GovernanceOutboxEntry)