< Summary

Information
Class: AsiBackbone.EntityFrameworkCore.Outbox.EfCoreGovernanceOutboxOutcomeStore
Assembly: AsiBackbone.EntityFrameworkCore
File(s): /home/runner/work/AsiBackbone/AsiBackbone/src/AsiBackbone.EntityFrameworkCore/Outbox/EfCoreGovernanceOutboxOutcomeStore.cs
Line coverage
64%
Covered lines: 70
Uncovered lines: 39
Coverable lines: 109
Total lines: 343
Line coverage: 64.2%
Branch coverage
75%
Covered branches: 15
Total branches: 20
Branch coverage: 75%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

File(s)

/home/runner/work/AsiBackbone/AsiBackbone/src/AsiBackbone.EntityFrameworkCore/Outbox/EfCoreGovernanceOutboxOutcomeStore.cs

#LineLine coverage
 1using System.Diagnostics.Metrics;
 2using AsiBackbone.Core.Emissions;
 3using AsiBackbone.Core.Outbox;
 4using Microsoft.EntityFrameworkCore;
 5using Microsoft.Extensions.Logging;
 6using Microsoft.Extensions.Logging.Abstractions;
 7
 8namespace AsiBackbone.EntityFrameworkCore.Outbox;
 9
 10/// <summary>
 11/// EF Core governance outbox store that exposes explicit caller-owned claim transition outcomes.
 12/// </summary>
 13/// <remarks>
 14/// This store preserves the existing convenience API while adding an outcome-aware contract that distinguishes applied
 15/// transitions from stale claims, terminal no-ops, concurrency losses, and missing rows. It delegates persistence to
 16/// <see cref="EfCoreGovernanceOutboxStore" /> and observes the scoped <see cref="DbContext" /> save boundary so callers
 17/// infer write ownership from the returned durable status alone.
 18/// </remarks>
 19public sealed class EfCoreGovernanceOutboxOutcomeStore : IAsiBackboneGovernanceOutboxClaimOutcomeStore
 20{
 221    private static readonly Meter Meter = new("AsiBackbone.EntityFrameworkCore.Outbox", "1.0.0");
 222    private static readonly Counter<long> ClaimTransitionCounter = Meter.CreateCounter<long>(
 223        "asibackbone.outbox.claim_transition_attempts",
 224        description: "Counts claimed outbox transition attempts by caller-visible outcome.");
 25
 226    private static readonly Action<ILogger, string, string, string, string, Exception?> LogClaimTransitionNotApplied =
 227        LoggerMessage.Define<string, string, string, string>(
 228            LogLevel.Warning,
 229            new EventId(19801, nameof(LogClaimTransitionNotApplied)),
 230            "Claimed governance outbox transition was not applied by worker {WorkerId} for entry {OutboxEntryId}. Outcom
 31
 32    private readonly DbContext dbContext;
 33    private readonly EfCoreGovernanceOutboxStore innerStore;
 34    private readonly ILogger<EfCoreGovernanceOutboxOutcomeStore> logger;
 35
 36    /// <summary>
 37    /// Initializes a new instance of the <see cref="EfCoreGovernanceOutboxOutcomeStore" /> class.
 38    /// </summary>
 39    /// <param name="dbContext">The host-owned scoped database context.</param>
 40    /// <param name="logger">The logger used for non-applied claimed-transition diagnostics.</param>
 1441    public EfCoreGovernanceOutboxOutcomeStore(
 1442        DbContext dbContext,
 1443        ILogger<EfCoreGovernanceOutboxOutcomeStore>? logger = null)
 44    {
 1445        ArgumentNullException.ThrowIfNull(dbContext);
 46
 1447        this.dbContext = dbContext;
 1448        innerStore = new EfCoreGovernanceOutboxStore(dbContext);
 1449        this.logger = logger ?? NullLogger<EfCoreGovernanceOutboxOutcomeStore>.Instance;
 1450    }
 51
 52    /// <inheritdoc />
 53    public ValueTask<GovernanceOutboxEntry> EnqueueAsync(
 54        GovernanceEmissionEnvelope envelope,
 55        CancellationToken cancellationToken = default)
 56    {
 057        return innerStore.EnqueueAsync(envelope, cancellationToken);
 58    }
 59
 60    /// <inheritdoc />
 61    public ValueTask<GovernanceOutboxEntry> SaveAsync(
 62        GovernanceOutboxEntry entry,
 63        CancellationToken cancellationToken = default)
 64    {
 1265        return innerStore.SaveAsync(entry, cancellationToken);
 66    }
 67
 68    /// <inheritdoc />
 69    public ValueTask<GovernanceOutboxEntry?> FindByOutboxEntryIdAsync(
 70        string outboxEntryId,
 71        CancellationToken cancellationToken = default)
 72    {
 073        return innerStore.FindByOutboxEntryIdAsync(outboxEntryId, cancellationToken);
 74    }
 75
 76    /// <inheritdoc />
 77    public ValueTask<IReadOnlyList<GovernanceOutboxEntry>> FindPendingAsync(
 78        int maxCount = 100,
 79        CancellationToken cancellationToken = default)
 80    {
 081        return innerStore.FindPendingAsync(maxCount, cancellationToken);
 82    }
 83
 84    /// <inheritdoc />
 85    public ValueTask<IReadOnlyList<GovernanceOutboxEntry>> FindRetryReadyAsync(
 86        DateTimeOffset utcNow,
 87        int maxCount = 100,
 88        CancellationToken cancellationToken = default)
 89    {
 090        return innerStore.FindRetryReadyAsync(utcNow, maxCount, cancellationToken);
 91    }
 92
 93    /// <inheritdoc />
 94    public ValueTask<IReadOnlyList<GovernanceOutboxClaim>> ClaimPendingAsync(
 95        GovernanceOutboxClaimRequest request,
 96        CancellationToken cancellationToken = default)
 97    {
 1298        return innerStore.ClaimPendingAsync(request, cancellationToken);
 99    }
 100
 101    /// <inheritdoc />
 102    public ValueTask<IReadOnlyList<GovernanceOutboxClaim>> ClaimRetryReadyAsync(
 103        GovernanceOutboxClaimRequest request,
 104        CancellationToken cancellationToken = default)
 105    {
 0106        return innerStore.ClaimRetryReadyAsync(request, cancellationToken);
 107    }
 108
 109    /// <inheritdoc />
 110    public async ValueTask<GovernanceOutboxEntry> MarkClaimDeliveredAsync(
 111        GovernanceOutboxClaim claim,
 112        GovernanceEmissionResult result,
 113        CancellationToken cancellationToken = default)
 114    {
 0115        return (await TryMarkClaimDeliveredAsync(claim, result, cancellationToken).ConfigureAwait(false)).Entry;
 0116    }
 117
 118    /// <inheritdoc />
 119    public async ValueTask<GovernanceOutboxEntry> MarkClaimFailedAsync(
 120        GovernanceOutboxClaim claim,
 121        GovernanceEmissionError governanceEmissionError,
 122        DateTimeOffset? nextRetryUtc = null,
 123        CancellationToken cancellationToken = default)
 124    {
 0125        return (await TryMarkClaimFailedAsync(claim, governanceEmissionError, nextRetryUtc, cancellationToken).Configure
 0126    }
 127
 128    /// <inheritdoc />
 129    public async ValueTask<GovernanceOutboxEntry> MarkClaimDeadLetteredAsync(
 130        GovernanceOutboxClaim claim,
 131        GovernanceEmissionError governanceEmissionError,
 132        string? deadLetterReason = null,
 133        CancellationToken cancellationToken = default)
 134    {
 0135        return (await TryMarkClaimDeadLetteredAsync(
 0136            claim,
 0137            governanceEmissionError,
 0138            deadLetterReason,
 0139            cancellationToken).ConfigureAwait(false)).Entry;
 0140    }
 141
 142    /// <inheritdoc />
 143    public async ValueTask<GovernanceOutboxEntry> SaveClaimAsync(
 144        GovernanceOutboxClaim claim,
 145        GovernanceOutboxEntry entry,
 146        CancellationToken cancellationToken = default)
 147    {
 0148        return (await TrySaveClaimAsync(claim, entry, cancellationToken).ConfigureAwait(false)).Entry;
 0149    }
 150
 151    /// <inheritdoc />
 152    public ValueTask<GovernanceOutboxEntry?> ReleaseClaimAsync(
 153        GovernanceOutboxClaim claim,
 154        string? reason = null,
 155        CancellationToken cancellationToken = default)
 156    {
 2157        return innerStore.ReleaseClaimAsync(claim, reason, cancellationToken);
 158    }
 159
 160    /// <inheritdoc />
 161    public ValueTask<GovernanceOutboxEntry> MarkDeliveredAsync(
 162        string outboxEntryId,
 163        GovernanceEmissionResult result,
 164        CancellationToken cancellationToken = default)
 165    {
 2166        return innerStore.MarkDeliveredAsync(outboxEntryId, result, cancellationToken);
 167    }
 168
 169    /// <inheritdoc />
 170    public ValueTask<GovernanceOutboxEntry> MarkFailedAsync(
 171        string outboxEntryId,
 172        GovernanceEmissionError governanceEmissionError,
 173        DateTimeOffset? nextRetryUtc = null,
 174        CancellationToken cancellationToken = default)
 175    {
 0176        return innerStore.MarkFailedAsync(outboxEntryId, governanceEmissionError, nextRetryUtc, cancellationToken);
 177    }
 178
 179    /// <inheritdoc />
 180    public ValueTask<GovernanceOutboxEntry> MarkDeadLetteredAsync(
 181        string outboxEntryId,
 182        GovernanceEmissionError governanceEmissionError,
 183        string? deadLetterReason = null,
 184        CancellationToken cancellationToken = default)
 185    {
 0186        return innerStore.MarkDeadLetteredAsync(outboxEntryId, governanceEmissionError, deadLetterReason, cancellationTo
 187    }
 188
 189    /// <inheritdoc />
 190    public ValueTask<GovernanceOutboxClaimTransitionResult> TryMarkClaimDeliveredAsync(
 191        GovernanceOutboxClaim claim,
 192        GovernanceEmissionResult result,
 193        CancellationToken cancellationToken = default)
 194    {
 8195        ArgumentNullException.ThrowIfNull(result);
 8196        return TryUpdateClaimAsync(
 8197            claim,
 6198            () => innerStore.MarkClaimDeliveredAsync(claim, result, cancellationToken),
 8199            cancellationToken);
 200    }
 201
 202    /// <inheritdoc />
 203    public ValueTask<GovernanceOutboxClaimTransitionResult> TryMarkClaimFailedAsync(
 204        GovernanceOutboxClaim claim,
 205        GovernanceEmissionError governanceEmissionError,
 206        DateTimeOffset? nextRetryUtc = null,
 207        CancellationToken cancellationToken = default)
 208    {
 4209        ArgumentNullException.ThrowIfNull(governanceEmissionError);
 4210        return TryUpdateClaimAsync(
 4211            claim,
 2212            () => innerStore.MarkClaimFailedAsync(claim, governanceEmissionError, nextRetryUtc, cancellationToken),
 4213            cancellationToken);
 214    }
 215
 216    /// <inheritdoc />
 217    public ValueTask<GovernanceOutboxClaimTransitionResult> TryMarkClaimDeadLetteredAsync(
 218        GovernanceOutboxClaim claim,
 219        GovernanceEmissionError governanceEmissionError,
 220        string? deadLetterReason = null,
 221        CancellationToken cancellationToken = default)
 222    {
 0223        ArgumentNullException.ThrowIfNull(governanceEmissionError);
 0224        return TryUpdateClaimAsync(
 0225            claim,
 0226            () => innerStore.MarkClaimDeadLetteredAsync(claim, governanceEmissionError, deadLetterReason, cancellationTo
 0227            cancellationToken);
 228    }
 229
 230    /// <inheritdoc />
 231    public ValueTask<GovernanceOutboxClaimTransitionResult> TrySaveClaimAsync(
 232        GovernanceOutboxClaim claim,
 233        GovernanceOutboxEntry entry,
 234        CancellationToken cancellationToken = default)
 235    {
 0236        ArgumentNullException.ThrowIfNull(entry);
 237
 0238        return !string.Equals(claim.OutboxEntryId, entry.OutboxEntryId, StringComparison.Ordinal)
 0239            ? throw new ArgumentException("Claim and entry must reference the same outbox entry ID.", nameof(entry))
 0240            : TryUpdateClaimAsync(
 0241            claim,
 0242            () => innerStore.SaveClaimAsync(claim, entry, cancellationToken),
 0243            cancellationToken);
 244    }
 245
 246    private async ValueTask<GovernanceOutboxClaimTransitionResult> TryUpdateClaimAsync(
 247        GovernanceOutboxClaim claim,
 248        Func<ValueTask<GovernanceOutboxEntry>> update,
 249        CancellationToken cancellationToken)
 250    {
 12251        ArgumentNullException.ThrowIfNull(claim);
 12252        cancellationToken.ThrowIfCancellationRequested();
 253
 12254        GovernanceOutboxEntry? currentEntry = await innerStore
 12255            .FindByOutboxEntryIdAsync(claim.OutboxEntryId, cancellationToken)
 12256            .ConfigureAwait(false);
 257
 12258        if (currentEntry is null)
 259        {
 0260            return Complete(claim, claim.Entry, GovernanceOutboxClaimTransitionOutcome.Missing);
 261        }
 262
 12263        if (currentEntry.IsDelivered || currentEntry.IsDeadLettered)
 264        {
 2265            return Complete(claim, currentEntry, GovernanceOutboxClaimTransitionOutcome.Terminal);
 266        }
 267
 10268        if (!currentEntry.IsClaimedBy(claim))
 269        {
 2270            return Complete(claim, currentEntry, GovernanceOutboxClaimTransitionOutcome.StaleClaim);
 271        }
 272
 8273        bool saveCompleted = false;
 274        void savedChangesHandler(object? sender, SavedChangesEventArgs eventArgs)
 275        {
 276            _ = sender;
 277            _ = eventArgs;
 2278            saveCompleted = true;
 2279        }
 280
 8281        dbContext.SavedChanges += savedChangesHandler;
 282
 283        GovernanceOutboxEntry returnedEntry;
 284        try
 285        {
 8286            returnedEntry = await update().ConfigureAwait(false);
 8287        }
 0288        catch (InvalidOperationException)
 289        {
 0290            GovernanceOutboxEntry? missingEntry = await innerStore
 0291                .FindByOutboxEntryIdAsync(claim.OutboxEntryId, cancellationToken)
 0292                .ConfigureAwait(false);
 293
 0294            if (missingEntry is null)
 295            {
 0296                return Complete(claim, currentEntry, GovernanceOutboxClaimTransitionOutcome.Missing);
 297            }
 298
 0299            throw;
 300        }
 301        finally
 302        {
 8303            dbContext.SavedChanges -= savedChangesHandler;
 304        }
 305
 8306        if (saveCompleted)
 307        {
 2308            return Complete(claim, returnedEntry, GovernanceOutboxClaimTransitionOutcome.Applied);
 309        }
 310
 6311        GovernanceOutboxEntry? refreshedEntry = await innerStore
 6312            .FindByOutboxEntryIdAsync(claim.OutboxEntryId, cancellationToken)
 6313            .ConfigureAwait(false);
 314
 6315        return refreshedEntry is null
 6316            ? Complete(claim, currentEntry, GovernanceOutboxClaimTransitionOutcome.Missing)
 6317            : Complete(claim, refreshedEntry, GovernanceOutboxClaimTransitionOutcome.ConcurrencyLost);
 12318    }
 319
 320    private GovernanceOutboxClaimTransitionResult Complete(
 321        GovernanceOutboxClaim claim,
 322        GovernanceOutboxEntry entry,
 323        GovernanceOutboxClaimTransitionOutcome outcome)
 324    {
 12325        ClaimTransitionCounter.Add(
 12326            1,
 12327            new KeyValuePair<string, object?>("outcome", outcome.ToString()),
 12328            new KeyValuePair<string, object?>("durable_status", entry.Status.ToString()));
 329
 12330        if (outcome is not GovernanceOutboxClaimTransitionOutcome.Applied)
 331        {
 10332            LogClaimTransitionNotApplied(
 10333                logger,
 10334                claim.WorkerId,
 10335                claim.OutboxEntryId,
 10336                outcome.ToString(),
 10337                entry.Status.ToString(),
 10338                null);
 339        }
 340
 12341        return GovernanceOutboxClaimTransitionResult.Create(entry, outcome);
 342    }
 343}

Methods/Properties

.cctor()
.ctor(Microsoft.EntityFrameworkCore.DbContext,Microsoft.Extensions.Logging.ILogger`1<AsiBackbone.EntityFrameworkCore.Outbox.EfCoreGovernanceOutboxOutcomeStore>)
EnqueueAsync(AsiBackbone.Core.Emissions.GovernanceEmissionEnvelope,System.Threading.CancellationToken)
SaveAsync(AsiBackbone.Core.Outbox.GovernanceOutboxEntry,System.Threading.CancellationToken)
FindByOutboxEntryIdAsync(System.String,System.Threading.CancellationToken)
FindPendingAsync(System.Int32,System.Threading.CancellationToken)
FindRetryReadyAsync(System.DateTimeOffset,System.Int32,System.Threading.CancellationToken)
ClaimPendingAsync(AsiBackbone.Core.Outbox.GovernanceOutboxClaimRequest,System.Threading.CancellationToken)
ClaimRetryReadyAsync(AsiBackbone.Core.Outbox.GovernanceOutboxClaimRequest,System.Threading.CancellationToken)
MarkClaimDeliveredAsync()
MarkClaimFailedAsync()
MarkClaimDeadLetteredAsync()
SaveClaimAsync()
ReleaseClaimAsync(AsiBackbone.Core.Outbox.GovernanceOutboxClaim,System.String,System.Threading.CancellationToken)
MarkDeliveredAsync(System.String,AsiBackbone.Core.Emissions.GovernanceEmissionResult,System.Threading.CancellationToken)
MarkFailedAsync(System.String,AsiBackbone.Core.Emissions.GovernanceEmissionError,System.Nullable`1<System.DateTimeOffset>,System.Threading.CancellationToken)
MarkDeadLetteredAsync(System.String,AsiBackbone.Core.Emissions.GovernanceEmissionError,System.String,System.Threading.CancellationToken)
TryMarkClaimDeliveredAsync(AsiBackbone.Core.Outbox.GovernanceOutboxClaim,AsiBackbone.Core.Emissions.GovernanceEmissionResult,System.Threading.CancellationToken)
TryMarkClaimFailedAsync(AsiBackbone.Core.Outbox.GovernanceOutboxClaim,AsiBackbone.Core.Emissions.GovernanceEmissionError,System.Nullable`1<System.DateTimeOffset>,System.Threading.CancellationToken)
TryMarkClaimDeadLetteredAsync(AsiBackbone.Core.Outbox.GovernanceOutboxClaim,AsiBackbone.Core.Emissions.GovernanceEmissionError,System.String,System.Threading.CancellationToken)
TrySaveClaimAsync(AsiBackbone.Core.Outbox.GovernanceOutboxClaim,AsiBackbone.Core.Outbox.GovernanceOutboxEntry,System.Threading.CancellationToken)
TryUpdateClaimAsync()
savedChangesHandler()
Complete(AsiBackbone.Core.Outbox.GovernanceOutboxClaim,AsiBackbone.Core.Outbox.GovernanceOutboxEntry,AsiBackbone.Core.Outbox.GovernanceOutboxClaimTransitionOutcome)