< 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
63%
Covered lines: 68
Uncovered lines: 39
Coverable lines: 107
Total lines: 343
Line coverage: 63.5%
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 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 : IGovernanceOutboxClaimOutcomeStore
 20{
 121    private static readonly Meter Meter = new("AsiBackbone.EntityFrameworkCore.Outbox", "1.0.0");
 122    private static readonly Counter<long> ClaimTransitionCounter = Meter.CreateCounter<long>(
 123        "asibackbone.outbox.claim_transition_attempts",
 124        description: "Counts claimed outbox transition attempts by caller-visible outcome.");
 25
 126    private static readonly Action<ILogger, string, string, string, string, Exception?> LogClaimTransitionNotApplied =
 127        LoggerMessage.Define<string, string, string, string>(
 128            LogLevel.Warning,
 129            new EventId(19901, nameof(LogClaimTransitionNotApplied)),
 130            "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>
 941    public EfCoreGovernanceOutboxOutcomeStore(
 942        DbContext dbContext,
 943        ILogger<EfCoreGovernanceOutboxOutcomeStore>? logger = null)
 44    {
 945        ArgumentNullException.ThrowIfNull(dbContext);
 46
 947        this.dbContext = dbContext;
 948        innerStore = new EfCoreGovernanceOutboxStore(dbContext);
 949        this.logger = logger ?? NullLogger<EfCoreGovernanceOutboxOutcomeStore>.Instance;
 950    }
 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    {
 665        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    {
 698        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    {
 1157        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    {
 1166        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    {
 4195        ArgumentNullException.ThrowIfNull(result);
 4196        return TryUpdateClaimAsync(
 4197            claim,
 4198            () => innerStore.MarkClaimDeliveredAsync(claim, result, cancellationToken),
 4199            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    {
 2209        ArgumentNullException.ThrowIfNull(governanceEmissionError);
 2210        return TryUpdateClaimAsync(
 2211            claim,
 2212            () => innerStore.MarkClaimFailedAsync(claim, governanceEmissionError, nextRetryUtc, cancellationToken),
 2213            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    {
 6251        ArgumentNullException.ThrowIfNull(claim);
 6252        cancellationToken.ThrowIfCancellationRequested();
 253
 6254        GovernanceOutboxEntry? currentEntry = await innerStore
 6255            .FindByOutboxEntryIdAsync(claim.OutboxEntryId, cancellationToken)
 6256            .ConfigureAwait(false);
 257
 6258        if (currentEntry is null)
 259        {
 0260            return Complete(claim, claim.Entry, GovernanceOutboxClaimTransitionOutcome.Missing);
 261        }
 262
 6263        if (currentEntry.IsDelivered || currentEntry.IsDeadLettered)
 264        {
 1265            return Complete(claim, currentEntry, GovernanceOutboxClaimTransitionOutcome.Terminal);
 266        }
 267
 5268        if (!currentEntry.IsClaimedBy(claim))
 269        {
 1270            return Complete(claim, currentEntry, GovernanceOutboxClaimTransitionOutcome.StaleClaim);
 271        }
 272
 4273        bool saveCompleted = false;
 274        void savedChangesHandler(object? sender, SavedChangesEventArgs eventArgs)
 275        {
 276            _ = sender;
 277            _ = eventArgs;
 278            saveCompleted = true;
 279        }
 280
 4281        dbContext.SavedChanges += savedChangesHandler;
 282
 283        GovernanceOutboxEntry returnedEntry;
 284        try
 285        {
 4286            returnedEntry = await update().ConfigureAwait(false);
 4287        }
 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        {
 4303            dbContext.SavedChanges -= savedChangesHandler;
 304        }
 305
 4306        if (saveCompleted)
 307        {
 1308            return Complete(claim, returnedEntry, GovernanceOutboxClaimTransitionOutcome.Applied);
 309        }
 310
 3311        GovernanceOutboxEntry? refreshedEntry = await innerStore
 3312            .FindByOutboxEntryIdAsync(claim.OutboxEntryId, cancellationToken)
 3313            .ConfigureAwait(false);
 314
 3315        return refreshedEntry is null
 3316            ? Complete(claim, currentEntry, GovernanceOutboxClaimTransitionOutcome.Missing)
 3317            : Complete(claim, refreshedEntry, GovernanceOutboxClaimTransitionOutcome.ConcurrencyLost);
 6318    }
 319
 320    private GovernanceOutboxClaimTransitionResult Complete(
 321        GovernanceOutboxClaim claim,
 322        GovernanceOutboxEntry entry,
 323        GovernanceOutboxClaimTransitionOutcome outcome)
 324    {
 6325        ClaimTransitionCounter.Add(
 6326            1,
 6327            new KeyValuePair<string, object?>("outcome", outcome.ToString()),
 6328            new KeyValuePair<string, object?>("durable_status", entry.Status.ToString()));
 329
 6330        if (outcome is not GovernanceOutboxClaimTransitionOutcome.Applied)
 331        {
 5332            LogClaimTransitionNotApplied(
 5333                logger,
 5334                claim.WorkerId,
 5335                claim.OutboxEntryId,
 5336                outcome.ToString(),
 5337                entry.Status.ToString(),
 5338                null);
 339        }
 340
 6341        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()
Complete(AsiBackbone.Core.Outbox.GovernanceOutboxClaim,AsiBackbone.Core.Outbox.GovernanceOutboxEntry,AsiBackbone.Core.Outbox.GovernanceOutboxClaimTransitionOutcome)