< Summary

Information
Class: AsiBackbone.EntityFrameworkCore.Outbox.EfCoreGovernanceOutboxStore
Assembly: AsiBackbone.EntityFrameworkCore
File(s): /home/runner/work/AsiBackbone/AsiBackbone/src/AsiBackbone.EntityFrameworkCore/Outbox/EfCoreGovernanceOutboxStore.cs
Line coverage
99%
Covered lines: 399
Uncovered lines: 4
Coverable lines: 403
Total lines: 778
Line coverage: 99%
Branch coverage
91%
Covered branches: 108
Total branches: 118
Branch coverage: 91.5%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.cctor()100%11100%
.ctor(...)100%11100%
EnqueueAsync()100%11100%
SaveAsync()100%44100%
FindByOutboxEntryIdAsync()100%22100%
FindPendingAsync()100%11100%
FindRetryReadyAsync()100%11100%
ClaimPendingAsync()100%11100%
ClaimRetryReadyAsync()100%11100%
MarkDeliveredAsync()100%22100%
MarkClaimDeliveredAsync()100%11100%
MarkFailedAsync()100%22100%
MarkClaimFailedAsync()100%11100%
MarkDeadLetteredAsync()100%22100%
MarkClaimDeadLetteredAsync()100%11100%
SaveClaimAsync()100%22100%
ReleaseClaimAsync()50%6686.67%
ClaimEntriesAsync()100%22100%
ReconcileTrackedEntries(...)80%101092.86%
MergeClaimColumns(...)100%22100%
ClaimColumnValues(...)100%11100%
UpdateClaimedEntryAsync()100%66100%
ApplyEntryUpdateAsync()100%11100%
RequireEntryAsync()50%22100%
RequireEntityAsync()100%22100%
OutboxEntries()100%11100%
ToEntity(...)100%2626100%
ToEntries(...)100%11100%
ToEntry(...)100%66100%
CreateClaim(...)100%88100%
IsPendingClaimEligible(...)0%620%
IsRetryReadyClaimEligible(...)92.86%1414100%
IsClaimAvailable(...)83.33%66100%
IsTerminal(...)100%22100%
DetachEntries(...)100%22100%
DeserializeMetadata(...)100%66100%
EmptyMetadata()100%11100%
NormalizeMaxCount(...)100%22100%

File(s)

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

#LineLine coverage
 1using System.Collections.ObjectModel;
 2using System.Linq.Expressions;
 3using System.Text.Json;
 4using AsiBackbone.Core.Emissions;
 5using AsiBackbone.Core.Entities;
 6using AsiBackbone.Core.Outbox;
 7using AsiBackbone.EntityFrameworkCore.Persistence;
 8using Microsoft.EntityFrameworkCore;
 9
 10namespace AsiBackbone.EntityFrameworkCore.Outbox;
 11
 12/// <summary>
 13/// Entity Framework Core-backed outbox store that persists provider-neutral emission envelopes through a host-owned <se
 14/// </summary>
 15/// <remarks>
 16/// This store provides durable local storage only. Provider delivery, telemetry export, SIEM routing, and cloud emissio
 17/// </remarks>
 18public sealed class EfCoreGovernanceOutboxStore : IGovernanceOutboxClaimStore
 19{
 120    private static readonly JsonSerializerOptions JsonOptions = new(JsonSerializerDefaults.Web);
 21
 22    private readonly DbContext dbContext;
 23
 24    /// <summary>
 25    /// Initializes a new instance of the <see cref="EfCoreGovernanceOutboxStore" /> class.
 26    /// </summary>
 27    /// <param name="dbContext">The host-owned database context.</param>
 13028    public EfCoreGovernanceOutboxStore(DbContext dbContext)
 29    {
 13030        ArgumentNullException.ThrowIfNull(dbContext);
 31
 13032        this.dbContext = dbContext;
 13033    }
 34
 35    /// <inheritdoc />
 36    public async ValueTask<GovernanceOutboxEntry> EnqueueAsync(
 37        GovernanceEmissionEnvelope envelope,
 38        CancellationToken cancellationToken = default)
 39    {
 4640        ArgumentNullException.ThrowIfNull(envelope);
 4641        cancellationToken.ThrowIfCancellationRequested();
 42
 4643        var entry = GovernanceOutboxEntry.Create(envelope);
 44
 4645        _ = dbContext
 4646            .Set<GovernanceOutboxEntryEntity>()
 4647            .Add(ToEntity(entry));
 48
 4649        _ = await dbContext.SaveChangesAsync(cancellationToken).ConfigureAwait(false);
 50
 4651        return entry;
 4652    }
 53
 54    /// <inheritdoc />
 55    public async ValueTask<GovernanceOutboxEntry> SaveAsync(
 56        GovernanceOutboxEntry entry,
 57        CancellationToken cancellationToken = default)
 58    {
 23359        ArgumentNullException.ThrowIfNull(entry);
 23360        cancellationToken.ThrowIfCancellationRequested();
 61
 23362        GovernanceOutboxEntryEntity persistedEntity = ToEntity(entry);
 23363        GovernanceOutboxEntryEntity? existingEntity = await dbContext
 23364            .Set<GovernanceOutboxEntryEntity>()
 23365            .SingleOrDefaultAsync(entity => entity.OutboxEntryId == entry.OutboxEntryId, cancellationToken)
 23366            .ConfigureAwait(false);
 67
 23368        if (existingEntity is null)
 69        {
 20670            _ = dbContext
 20671                .Set<GovernanceOutboxEntryEntity>()
 20672                .Add(persistedEntity);
 73        }
 74        else
 75        {
 2776            GovernanceOutboxEntry currentEntry = ToEntry(existingEntity);
 2777            if (IsTerminal(currentEntry))
 78            {
 279                return currentEntry;
 80            }
 81
 2582            persistedEntity.Id = existingEntity.Id;
 2583            persistedEntity.ConcurrencyStamp = GovernanceEntity.NewConcurrencyStamp();
 2584            dbContext.Entry(existingEntity).CurrentValues.SetValues(persistedEntity);
 85        }
 86
 23187        _ = await dbContext.SaveChangesAsync(cancellationToken).ConfigureAwait(false);
 88
 22889        return entry;
 23090    }
 91
 92    /// <inheritdoc />
 93    public async ValueTask<GovernanceOutboxEntry?> FindByOutboxEntryIdAsync(
 94        string outboxEntryId,
 95        CancellationToken cancellationToken = default)
 96    {
 9197        ArgumentException.ThrowIfNullOrWhiteSpace(outboxEntryId);
 98
 9199        string normalizedOutboxEntryId = outboxEntryId.Trim();
 100
 91101        GovernanceOutboxEntryEntity? entity = await OutboxEntries()
 91102            .Where(outboxEntry => outboxEntry.OutboxEntryId == normalizedOutboxEntryId)
 91103            .SingleOrDefaultAsync(cancellationToken)
 91104            .ConfigureAwait(false);
 105
 91106        return entity is null ? null : ToEntry(entity);
 91107    }
 108
 109    /// <inheritdoc />
 110    public async ValueTask<IReadOnlyList<GovernanceOutboxEntry>> FindPendingAsync(
 111        int maxCount = 100,
 112        CancellationToken cancellationToken = default)
 113    {
 17114        int normalizedMaxCount = NormalizeMaxCount(maxCount);
 115
 16116        List<GovernanceOutboxEntryEntity> entities = await OutboxEntries()
 16117            .Where(outboxEntry => outboxEntry.Status == GovernanceEmissionStatus.Pending)
 16118            .OrderBy(outboxEntry => outboxEntry.CreatedUtc)
 16119            .ThenBy(outboxEntry => outboxEntry.OutboxEntryId)
 16120            .Take(normalizedMaxCount)
 16121            .ToListAsync(cancellationToken)
 16122            .ConfigureAwait(false);
 123
 16124        return ToEntries(entities);
 16125    }
 126
 127    /// <inheritdoc />
 128    public async ValueTask<IReadOnlyList<GovernanceOutboxEntry>> FindRetryReadyAsync(
 129        DateTimeOffset utcNow,
 130        int maxCount = 100,
 131        CancellationToken cancellationToken = default)
 132    {
 13133        int normalizedMaxCount = NormalizeMaxCount(maxCount);
 12134        DateTimeOffset normalizedUtcNow = utcNow.ToUniversalTime();
 135
 12136        List<GovernanceOutboxEntryEntity> entities = await OutboxEntries()
 12137            .Where(outboxEntry =>
 12138                outboxEntry.Status == GovernanceEmissionStatus.Deferred ||
 12139                outboxEntry.Status == GovernanceEmissionStatus.Failed ||
 12140                outboxEntry.Status == GovernanceEmissionStatus.RetryableFailure)
 12141            .Where(outboxEntry => outboxEntry.NextRetryUtc == null || outboxEntry.NextRetryUtc <= normalizedUtcNow)
 12142            .OrderBy(outboxEntry => outboxEntry.NextRetryUtc ?? outboxEntry.UpdatedUtc)
 12143            .ThenBy(outboxEntry => outboxEntry.OutboxEntryId)
 12144            .Take(normalizedMaxCount)
 12145            .ToListAsync(cancellationToken)
 12146            .ConfigureAwait(false);
 147
 12148        return ToEntries(entities);
 12149    }
 150
 151    /// <inheritdoc />
 152    public async ValueTask<IReadOnlyList<GovernanceOutboxClaim>> ClaimPendingAsync(
 153        GovernanceOutboxClaimRequest request,
 154        CancellationToken cancellationToken = default)
 155    {
 40156        ArgumentNullException.ThrowIfNull(request);
 40157        cancellationToken.ThrowIfCancellationRequested();
 158
 40159        Expression<Func<GovernanceOutboxEntryEntity, bool>> isEligible = outboxEntry =>
 40160            outboxEntry.Status == GovernanceEmissionStatus.Pending
 40161            && (outboxEntry.ClaimToken == null || outboxEntry.ClaimExpiresUtc == null || outboxEntry.ClaimExpiresUtc <= 
 162
 40163        IQueryable<GovernanceOutboxEntryEntity> candidates = OutboxEntries()
 40164            .Where(isEligible)
 40165            .OrderBy(outboxEntry => outboxEntry.CreatedUtc)
 40166            .ThenBy(outboxEntry => outboxEntry.OutboxEntryId)
 40167            .Take(request.MaxCount);
 168
 40169        return await ClaimEntriesAsync(
 40170            isEligible,
 40171            candidates,
 40172            query => query
 40173                .OrderBy(outboxEntry => outboxEntry.CreatedUtc)
 40174                .ThenBy(outboxEntry => outboxEntry.OutboxEntryId),
 40175            request,
 40176            cancellationToken)
 40177            .ConfigureAwait(false);
 40178    }
 179
 180    /// <inheritdoc />
 181    public async ValueTask<IReadOnlyList<GovernanceOutboxClaim>> ClaimRetryReadyAsync(
 182        GovernanceOutboxClaimRequest request,
 183        CancellationToken cancellationToken = default)
 184    {
 6185        ArgumentNullException.ThrowIfNull(request);
 5186        cancellationToken.ThrowIfCancellationRequested();
 187
 4188        Expression<Func<GovernanceOutboxEntryEntity, bool>> isEligible = outboxEntry =>
 4189            (outboxEntry.Status == GovernanceEmissionStatus.Deferred ||
 4190                outboxEntry.Status == GovernanceEmissionStatus.Failed ||
 4191                outboxEntry.Status == GovernanceEmissionStatus.RetryableFailure)
 4192            && (outboxEntry.NextRetryUtc == null || outboxEntry.NextRetryUtc <= request.UtcNow)
 4193            && (outboxEntry.ClaimToken == null || outboxEntry.ClaimExpiresUtc == null || outboxEntry.ClaimExpiresUtc <= 
 194
 4195        IQueryable<GovernanceOutboxEntryEntity> candidates = OutboxEntries()
 4196            .Where(isEligible)
 4197            .OrderBy(outboxEntry => outboxEntry.NextRetryUtc ?? outboxEntry.UpdatedUtc)
 4198            .ThenBy(outboxEntry => outboxEntry.OutboxEntryId)
 4199            .Take(request.MaxCount);
 200
 4201        return await ClaimEntriesAsync(
 4202            isEligible,
 4203            candidates,
 4204            query => query
 4205                .OrderBy(outboxEntry => outboxEntry.NextRetryUtc ?? outboxEntry.UpdatedUtc)
 4206                .ThenBy(outboxEntry => outboxEntry.OutboxEntryId),
 4207            request,
 4208            cancellationToken)
 4209            .ConfigureAwait(false);
 4210    }
 211
 212    /// <inheritdoc />
 213    public async ValueTask<GovernanceOutboxEntry> MarkDeliveredAsync(
 214        string outboxEntryId,
 215        GovernanceEmissionResult result,
 216        CancellationToken cancellationToken = default)
 217    {
 12218        ArgumentException.ThrowIfNullOrWhiteSpace(outboxEntryId);
 12219        ArgumentNullException.ThrowIfNull(result);
 220
 12221        GovernanceOutboxEntry entry = await RequireEntryAsync(outboxEntryId, cancellationToken).ConfigureAwait(false);
 12222        if (IsTerminal(entry))
 223        {
 3224            return entry;
 225        }
 226
 9227        GovernanceOutboxEntry updatedEntry = entry.MarkDelivered(result);
 228
 9229        return await SaveAsync(updatedEntry, cancellationToken).ConfigureAwait(false);
 11230    }
 231
 232    /// <inheritdoc />
 233    public async ValueTask<GovernanceOutboxEntry> MarkClaimDeliveredAsync(
 234        GovernanceOutboxClaim claim,
 235        GovernanceEmissionResult result,
 236        CancellationToken cancellationToken = default)
 237    {
 16238        ArgumentNullException.ThrowIfNull(claim);
 15239        ArgumentNullException.ThrowIfNull(result);
 240
 14241        return await UpdateClaimedEntryAsync(claim, entry => entry.MarkDelivered(result), cancellationToken).ConfigureAw
 13242    }
 243
 244    /// <inheritdoc />
 245    public async ValueTask<GovernanceOutboxEntry> MarkFailedAsync(
 246        string outboxEntryId,
 247        GovernanceEmissionError governanceEmissionError,
 248        DateTimeOffset? nextRetryUtc = null,
 249        CancellationToken cancellationToken = default)
 250    {
 8251        ArgumentException.ThrowIfNullOrWhiteSpace(outboxEntryId);
 8252        ArgumentNullException.ThrowIfNull(governanceEmissionError);
 253
 8254        GovernanceOutboxEntry entry = await RequireEntryAsync(outboxEntryId, cancellationToken).ConfigureAwait(false);
 8255        if (IsTerminal(entry))
 256        {
 2257            return entry;
 258        }
 259
 6260        GovernanceOutboxEntry updatedEntry = entry.MarkFailed(governanceEmissionError, nextRetryUtc);
 261
 6262        return await SaveAsync(updatedEntry, cancellationToken).ConfigureAwait(false);
 7263    }
 264
 265    /// <inheritdoc />
 266    public async ValueTask<GovernanceOutboxEntry> MarkClaimFailedAsync(
 267        GovernanceOutboxClaim claim,
 268        GovernanceEmissionError governanceEmissionError,
 269        DateTimeOffset? nextRetryUtc = null,
 270        CancellationToken cancellationToken = default)
 271    {
 13272        ArgumentNullException.ThrowIfNull(claim);
 12273        ArgumentNullException.ThrowIfNull(governanceEmissionError);
 274
 11275        return await UpdateClaimedEntryAsync(
 11276            claim,
 11277            entry => entry.MarkFailed(governanceEmissionError, nextRetryUtc),
 11278            cancellationToken)
 11279            .ConfigureAwait(false);
 11280    }
 281
 282    /// <inheritdoc />
 283    public async ValueTask<GovernanceOutboxEntry> MarkDeadLetteredAsync(
 284        string outboxEntryId,
 285        GovernanceEmissionError governanceEmissionError,
 286        string? deadLetterReason = null,
 287        CancellationToken cancellationToken = default)
 288    {
 9289        ArgumentException.ThrowIfNullOrWhiteSpace(outboxEntryId);
 9290        ArgumentNullException.ThrowIfNull(governanceEmissionError);
 291
 9292        GovernanceOutboxEntry entry = await RequireEntryAsync(outboxEntryId, cancellationToken).ConfigureAwait(false);
 9293        if (IsTerminal(entry))
 294        {
 2295            return entry;
 296        }
 297
 7298        GovernanceOutboxEntry updatedEntry = entry.MarkDeadLettered(governanceEmissionError, deadLetterReason);
 299
 7300        return await SaveAsync(updatedEntry, cancellationToken).ConfigureAwait(false);
 9301    }
 302
 303    /// <inheritdoc />
 304    public async ValueTask<GovernanceOutboxEntry> MarkClaimDeadLetteredAsync(
 305        GovernanceOutboxClaim claim,
 306        GovernanceEmissionError governanceEmissionError,
 307        string? deadLetterReason = null,
 308        CancellationToken cancellationToken = default)
 309    {
 5310        ArgumentNullException.ThrowIfNull(claim);
 4311        ArgumentNullException.ThrowIfNull(governanceEmissionError);
 312
 3313        return await UpdateClaimedEntryAsync(
 3314            claim,
 3315            entry => entry.MarkDeadLettered(governanceEmissionError, deadLetterReason),
 3316            cancellationToken)
 3317            .ConfigureAwait(false);
 3318    }
 319
 320    /// <inheritdoc />
 321    public async ValueTask<GovernanceOutboxEntry> SaveClaimAsync(
 322        GovernanceOutboxClaim claim,
 323        GovernanceOutboxEntry entry,
 324        CancellationToken cancellationToken = default)
 325    {
 5326        ArgumentNullException.ThrowIfNull(claim);
 4327        ArgumentNullException.ThrowIfNull(entry);
 328
 3329        return !string.Equals(claim.OutboxEntryId, entry.OutboxEntryId, StringComparison.Ordinal)
 3330            ? throw new ArgumentException("Claim and entry must reference the same outbox entry ID.", nameof(entry))
 3331            : await UpdateClaimedEntryAsync(claim, _ => entry, cancellationToken).ConfigureAwait(false);
 1332    }
 333
 334    /// <inheritdoc />
 335    public async ValueTask<GovernanceOutboxEntry?> ReleaseClaimAsync(
 336        GovernanceOutboxClaim claim,
 337        string? reason = null,
 338        CancellationToken cancellationToken = default)
 339    {
 3340        ArgumentNullException.ThrowIfNull(claim);
 3341        cancellationToken.ThrowIfCancellationRequested();
 342
 3343        GovernanceOutboxEntryEntity? entity = await dbContext
 3344            .Set<GovernanceOutboxEntryEntity>()
 3345            .SingleOrDefaultAsync(outboxEntry => outboxEntry.OutboxEntryId == claim.OutboxEntryId, cancellationToken)
 3346            .ConfigureAwait(false);
 347
 3348        if (entity is null)
 349        {
 0350            return null;
 351        }
 352
 3353        GovernanceOutboxEntry currentEntry = ToEntry(entity);
 3354        if (!currentEntry.IsClaimedBy(claim) || IsTerminal(currentEntry))
 355        {
 0356            return currentEntry;
 357        }
 358
 3359        GovernanceOutboxEntry releasedEntry = currentEntry.ReleaseClaim();
 3360        await ApplyEntryUpdateAsync(entity, releasedEntry, cancellationToken).ConfigureAwait(false);
 361
 3362        return releasedEntry;
 3363    }
 364
 365    /// <summary>
 366    /// Claims up to one batch of candidate rows in a single set-based update and returns the rows this call won.
 367    /// </summary>
 368    /// <remarks>
 369    /// <para>
 370    /// <paramref name="isEligible" /> is applied twice: once inside <paramref name="candidates" />, which chooses and
 371    /// orders the batch, and again directly to each row the update writes. The second application is what makes the
 372    /// claim safe under row-locking concurrency (issue #823).
 373    /// </para>
 374    /// <para>
 375    /// A claim statement chooses its candidate set before it can lock the chosen rows, so two workers starting
 376    /// together can choose the same rows. The first worker locks and claims them; the second waits on those row locks
 377    /// and resumes after the first commits. If eligibility were checked only inside the candidate query, the second
 378    /// worker could resume with a candidate set read before the first claim became visible, overwrite the first
 379    /// worker's claim token, and leave both workers believing they hold the same entries. PostgreSQL at read committed
 380    /// re-evaluates the update's own predicate against the newly committed row version before writing it, but not the
 381    /// already-materialized candidate set. SQL Server evaluates the update target's predicate against the committed
 382    /// row after it acquires the update lock. Stating eligibility on the update target therefore makes the waiting
 383    /// worker skip rows that another worker claimed while it waited.
 384    /// </para>
 385    /// <para>
 386    /// The consequence is liveness, not correctness: a worker that loses a row to another worker receives a smaller
 387    /// batch, possibly an empty one, and claims again on its next pass. The expression is provider-neutral LINQ; no
 388    /// lock hints or provider-specific SQL are involved. The behavior is verified against SQL Server and PostgreSQL by
 389    /// the opt-in provider contention tests.
 390    /// </para>
 391    /// </remarks>
 392    private async ValueTask<IReadOnlyList<GovernanceOutboxClaim>> ClaimEntriesAsync(
 393        Expression<Func<GovernanceOutboxEntryEntity, bool>> isEligible,
 394        IQueryable<GovernanceOutboxEntryEntity> candidates,
 395        Func<IQueryable<GovernanceOutboxEntryEntity>, IOrderedQueryable<GovernanceOutboxEntryEntity>> orderClaimedEntrie
 396        GovernanceOutboxClaimRequest request,
 397        CancellationToken cancellationToken)
 398    {
 44399        string claimToken = Guid.NewGuid().ToString("N");
 44400        string concurrencyStamp = GovernanceEntity.NewConcurrencyStamp();
 44401        IQueryable<Guid> candidateIds = candidates.Select(outboxEntry => outboxEntry.Id);
 402
 44403        int claimedCount = await OutboxEntries()
 44404            .Where(isEligible)
 44405            .Where(outboxEntry => candidateIds.Contains(outboxEntry.Id))
 44406            .ExecuteUpdateAsync(
 44407                setters => setters
 44408                    .SetProperty(outboxEntry => outboxEntry.ConcurrencyStamp, concurrencyStamp)
 44409                    .SetProperty(outboxEntry => outboxEntry.UpdatedUtc, request.UtcNow)
 44410                    .SetProperty(outboxEntry => outboxEntry.ClaimOwner, request.WorkerId)
 44411                    .SetProperty(outboxEntry => outboxEntry.ClaimToken, claimToken)
 44412                    .SetProperty(outboxEntry => outboxEntry.ClaimedUtc, request.UtcNow)
 44413                    .SetProperty(outboxEntry => outboxEntry.ClaimExpiresUtc, request.ClaimExpiresUtc)
 44414                    .SetProperty(outboxEntry => outboxEntry.ClaimAttemptCount, outboxEntry => outboxEntry.ClaimAttemptCo
 44415                cancellationToken)
 44416            .ConfigureAwait(false);
 417
 44418        if (claimedCount == 0)
 419        {
 3420            return [];
 421        }
 422
 41423        List<GovernanceOutboxEntryEntity> claimedEntities = await orderClaimedEntries(
 41424            OutboxEntries().Where(outboxEntry => outboxEntry.ClaimToken == claimToken))
 41425            .ToListAsync(cancellationToken)
 41426            .ConfigureAwait(false);
 427
 41428        ReconcileTrackedEntries(claimedEntities);
 429
 41430        return [.. claimedEntities.Select(ToEntry).Select(CreateClaim)];
 44431    }
 432
 433    /// <summary>
 434    /// Reconciles tracked instances of rows that a set-based update has just rewritten.
 435    /// </summary>
 436    /// <remarks>
 437    /// <c>ExecuteUpdateAsync</c> writes to the database without passing through the change tracker. An instance the
 438    /// host-owned context already tracks, such as one added by <see cref="EnqueueAsync" />, keeps its pre-claim values,
 439    /// and identity resolution returns that stale instance from later tracking queries. The claim then appears not to
 440    /// be held, so claim-scoped transitions return without applying.
 441    /// <para>
 442    /// An unchanged instance is detached, so the next query loads the claimed row. A modified or deleted instance holds
 443    /// unsaved host work that detaching would discard, so the claim-governed columns are resolved from the persisted
 444    /// claimed row instead: the columns the claim wrote, plus <see cref="GovernanceOutboxEntryEntity.Status" />, which
 445    /// the claim was taken against. Each is set as both the original and the current value, overriding any pending host
 446    /// change to that column. A pending host status or claim-field change would otherwise let the tracked instance look
 447    /// terminal or unclaimed, so the claim-scoped transition would return without recording delivery or clearing the
 448    /// lease. The host's other pending changes, and a pending deletion, are kept and save against the claimed row.
 449    /// </para>
 450    /// </remarks>
 451    /// <param name="claimedEntities">The claimed rows, read without tracking after the update.</param>
 452    private void ReconcileTrackedEntries(List<GovernanceOutboxEntryEntity> claimedEntities)
 453    {
 41454        if (claimedEntities.Count == 0)
 455        {
 0456            return;
 457        }
 458
 41459        var claimedById = claimedEntities.ToDictionary(entity => entity.OutboxEntryId, StringComparer.Ordinal);
 460
 41461        List<Microsoft.EntityFrameworkCore.ChangeTracking.EntityEntry<GovernanceOutboxEntryEntity>> staleEntries = [.. d
 41462            .ChangeTracker
 41463            .Entries<GovernanceOutboxEntryEntity>()
 41464            .Where(entry => claimedById.ContainsKey(entry.Entity.OutboxEntryId))];
 465
 96466        foreach (Microsoft.EntityFrameworkCore.ChangeTracking.EntityEntry<GovernanceOutboxEntryEntity> entry in staleEnt
 467        {
 7468            if (entry.State == EntityState.Unchanged)
 469            {
 3470                entry.State = EntityState.Detached;
 3471                continue;
 472            }
 473
 4474            if (entry.State is EntityState.Modified or EntityState.Deleted)
 475            {
 4476                MergeClaimColumns(entry, claimedById[entry.Entity.OutboxEntryId]);
 477            }
 478        }
 41479    }
 480
 481    private static void MergeClaimColumns(
 482        Microsoft.EntityFrameworkCore.ChangeTracking.EntityEntry<GovernanceOutboxEntryEntity> entry,
 483        GovernanceOutboxEntryEntity claimed)
 484    {
 72485        foreach ((string propertyName, object? claimedValue) in ClaimColumnValues(claimed))
 486        {
 32487            Microsoft.EntityFrameworkCore.ChangeTracking.PropertyEntry property = entry.Property(propertyName);
 32488            property.OriginalValue = claimedValue;
 32489            property.CurrentValue = claimedValue;
 490        }
 4491    }
 492
 493    // The claim-governed columns: those written by the claim update in ClaimEntriesAsync, which must stay aligned with
 494    // that update, plus Status, which the claim eligibility query was evaluated against.
 495    private static (string PropertyName, object? Value)[] ClaimColumnValues(GovernanceOutboxEntryEntity claimed)
 496    {
 4497        return
 4498        [
 4499            (nameof(GovernanceOutboxEntryEntity.ConcurrencyStamp), claimed.ConcurrencyStamp),
 4500            (nameof(GovernanceOutboxEntryEntity.Status), claimed.Status),
 4501            (nameof(GovernanceOutboxEntryEntity.UpdatedUtc), claimed.UpdatedUtc),
 4502            (nameof(GovernanceOutboxEntryEntity.ClaimOwner), claimed.ClaimOwner),
 4503            (nameof(GovernanceOutboxEntryEntity.ClaimToken), claimed.ClaimToken),
 4504            (nameof(GovernanceOutboxEntryEntity.ClaimedUtc), claimed.ClaimedUtc),
 4505            (nameof(GovernanceOutboxEntryEntity.ClaimExpiresUtc), claimed.ClaimExpiresUtc),
 4506            (nameof(GovernanceOutboxEntryEntity.ClaimAttemptCount), claimed.ClaimAttemptCount),
 4507        ];
 508    }
 509
 510    private async ValueTask<GovernanceOutboxEntry> UpdateClaimedEntryAsync(
 511        GovernanceOutboxClaim claim,
 512        Func<GovernanceOutboxEntry, GovernanceOutboxEntry> updateEntry,
 513        CancellationToken cancellationToken)
 514    {
 30515        GovernanceOutboxEntryEntity entity = await RequireEntityAsync(claim.OutboxEntryId, cancellationToken).ConfigureA
 28516        GovernanceOutboxEntry currentEntry = ToEntry(entity);
 517
 28518        if (!currentEntry.IsClaimedBy(claim) || IsTerminal(currentEntry))
 519        {
 6520            return currentEntry;
 521        }
 522
 22523        GovernanceOutboxEntry updatedEntry = updateEntry(currentEntry);
 524
 525        try
 526        {
 22527            await ApplyEntryUpdateAsync(entity, updatedEntry, cancellationToken).ConfigureAwait(false);
 17528            return updatedEntry;
 529        }
 530        catch (DbUpdateConcurrencyException exception)
 531        {
 5532            DetachEntries(exception);
 5533            GovernanceOutboxEntry? refreshedEntry = await FindByOutboxEntryIdAsync(claim.OutboxEntryId, cancellationToke
 5534            return refreshedEntry ?? currentEntry;
 535        }
 28536    }
 537
 538    private async ValueTask ApplyEntryUpdateAsync(
 539        GovernanceOutboxEntryEntity entity,
 540        GovernanceOutboxEntry entry,
 541        CancellationToken cancellationToken)
 542    {
 25543        GovernanceOutboxEntryEntity persistedEntity = ToEntity(entry);
 25544        persistedEntity.Id = entity.Id;
 25545        persistedEntity.ConcurrencyStamp = GovernanceEntity.NewConcurrencyStamp();
 25546        dbContext.Entry(entity).CurrentValues.SetValues(persistedEntity);
 547
 25548        _ = await dbContext.SaveChangesAsync(cancellationToken).ConfigureAwait(false);
 20549    }
 550
 551    private async ValueTask<GovernanceOutboxEntry> RequireEntryAsync(
 552        string outboxEntryId,
 553        CancellationToken cancellationToken)
 554    {
 29555        GovernanceOutboxEntry? entry = await FindByOutboxEntryIdAsync(outboxEntryId, cancellationToken).ConfigureAwait(f
 556
 29557        return entry ?? throw new InvalidOperationException($"Outbox entry '{outboxEntryId.Trim()}' was not found.");
 29558    }
 559
 560    private async ValueTask<GovernanceOutboxEntryEntity> RequireEntityAsync(
 561        string outboxEntryId,
 562        CancellationToken cancellationToken)
 563    {
 30564        GovernanceOutboxEntryEntity? entity = await dbContext
 30565            .Set<GovernanceOutboxEntryEntity>()
 30566            .SingleOrDefaultAsync(outboxEntry => outboxEntry.OutboxEntryId == outboxEntryId.Trim(), cancellationToken)
 30567            .ConfigureAwait(false);
 568
 29569        return entity ?? throw new InvalidOperationException($"Outbox entry '{outboxEntryId.Trim()}' was not found.");
 28570    }
 571
 572    private IQueryable<GovernanceOutboxEntryEntity> OutboxEntries()
 573    {
 248574        return dbContext.Set<GovernanceOutboxEntryEntity>().AsNoTracking();
 575    }
 576
 577    private static GovernanceOutboxEntryEntity ToEntity(GovernanceOutboxEntry entry)
 578    {
 304579        GovernanceEmissionEnvelope envelope = entry.Envelope;
 304580        GovernanceEmissionPayload? payload = envelope.Payload;
 304581        GovernanceEmissionError? lastError = entry.LastError;
 582
 304583        return new GovernanceOutboxEntryEntity
 304584        {
 304585            OutboxEntryId = entry.OutboxEntryId,
 304586            Status = entry.Status,
 304587            CreatedUtc = entry.CreatedUtc,
 304588            UpdatedUtc = entry.UpdatedUtc,
 304589            DeliveredUtc = entry.Status is GovernanceEmissionStatus.Delivered ? entry.UpdatedUtc : null,
 304590            RetryCount = entry.RetryCount,
 304591            MaxRetryCount = entry.MaxRetryCount,
 304592            NextRetryUtc = entry.NextRetryUtc,
 304593            ProviderName = entry.ProviderName,
 304594            ProviderRecordId = entry.ProviderRecordId,
 304595            DeadLetterReason = entry.DeadLetterReason,
 304596            LastErrorCode = lastError?.Code,
 304597            LastErrorMessage = lastError?.Message,
 304598            LastErrorIsRetryable = lastError?.IsRetryable,
 304599            LastErrorProviderName = lastError?.ProviderName,
 304600            LastErrorProviderErrorCode = lastError?.ProviderErrorCode,
 304601            MetadataJson = JsonSerializer.Serialize(entry.Metadata, JsonOptions),
 304602            ClaimOwner = entry.ClaimOwner,
 304603            ClaimToken = entry.ClaimToken,
 304604            ClaimedUtc = entry.ClaimedUtc,
 304605            ClaimExpiresUtc = entry.ClaimExpiresUtc,
 304606            ClaimAttemptCount = entry.ClaimAttemptCount,
 304607            EnvelopeId = envelope.EnvelopeId,
 304608            EnvelopeSchemaVersion = envelope.SchemaVersion,
 304609            EnvelopeEventType = envelope.EventType,
 304610            EnvelopeEventId = envelope.EventId,
 304611            EnvelopeOccurredUtc = envelope.OccurredUtc,
 304612            EnvelopeCreatedUtc = envelope.CreatedUtc,
 304613            EnvelopeCorrelationId = envelope.CorrelationId,
 304614            EnvelopeDecisionReceiptId = envelope.DecisionReceiptId,
 304615            EnvelopeLifecycleStage = envelope.LifecycleStage,
 304616            EnvelopeLifecycleStageSequence = envelope.LifecycleStageSequence,
 304617            EnvelopePolicyVersion = envelope.PolicyVersion,
 304618            EnvelopePolicyHash = envelope.PolicyHash,
 304619            EnvelopeTraceId = envelope.TraceId,
 304620            EnvelopeSpanId = envelope.SpanId,
 304621            EnvelopeParentSpanId = envelope.ParentSpanId,
 304622            EnvelopeOperationName = envelope.OperationName,
 304623            EnvelopeOutcome = envelope.Outcome,
 304624            EnvelopeActorId = envelope.ActorId,
 304625            EnvelopeEmitterStatus = envelope.EmitterStatus,
 304626            EnvelopeEmitterProvider = envelope.EmitterProvider,
 304627            EnvelopeOutboxSequence = envelope.OutboxSequence,
 304628            EnvelopeGatewayExecutionId = envelope.GatewayExecutionId,
 304629            EnvelopeDecisionStage = envelope.DecisionStage,
 304630            EnvelopeMetadataJson = JsonSerializer.Serialize(envelope.Metadata, JsonOptions),
 304631            EnvelopePayloadType = payload?.PayloadType,
 304632            EnvelopePayloadSchemaVersion = payload?.SchemaVersion,
 304633            EnvelopePayloadContentType = payload?.ContentType,
 304634            EnvelopePayloadContentHash = payload?.ContentHash,
 304635            EnvelopePayloadSizeBytes = payload?.SizeBytes,
 304636            EnvelopePayloadMetadataJson = JsonSerializer.Serialize(payload?.Metadata ?? EmptyMetadata(), JsonOptions)
 304637        };
 638    }
 639
 640    private static GovernanceOutboxEntry[] ToEntries(IEnumerable<GovernanceOutboxEntryEntity> entities)
 641    {
 28642        return [.. entities.Select(ToEntry)];
 643    }
 644
 645    private static GovernanceOutboxEntry ToEntry(GovernanceOutboxEntryEntity entity)
 646    {
 381647        GovernanceEmissionPayload? payload = string.IsNullOrWhiteSpace(entity.EnvelopePayloadType)
 381648            ? null
 381649            : GovernanceEmissionPayload.Create(
 381650                entity.EnvelopePayloadType,
 381651                entity.EnvelopePayloadSchemaVersion,
 381652                entity.EnvelopePayloadContentType,
 381653                entity.EnvelopePayloadContentHash,
 381654                entity.EnvelopePayloadSizeBytes,
 381655                DeserializeMetadata(entity.EnvelopePayloadMetadataJson));
 656
 381657        var envelope = GovernanceEmissionEnvelope.Create(
 381658            entity.EnvelopeEventType,
 381659            entity.EnvelopeEventId,
 381660            entity.EnvelopeOccurredUtc,
 381661            entity.EnvelopeId,
 381662            entity.EnvelopeCreatedUtc,
 381663            entity.EnvelopeSchemaVersion,
 381664            entity.EnvelopeCorrelationId,
 381665            entity.EnvelopeDecisionReceiptId,
 381666            entity.EnvelopeLifecycleStage,
 381667            entity.EnvelopePolicyVersion,
 381668            entity.EnvelopePolicyHash,
 381669            entity.EnvelopeTraceId,
 381670            entity.EnvelopeSpanId,
 381671            entity.EnvelopeParentSpanId,
 381672            entity.EnvelopeOperationName,
 381673            entity.EnvelopeOutcome,
 381674            entity.EnvelopeActorId,
 381675            entity.EnvelopeEmitterStatus,
 381676            entity.EnvelopeEmitterProvider,
 381677            entity.EnvelopeOutboxSequence,
 381678            entity.EnvelopeGatewayExecutionId,
 381679            entity.EnvelopeDecisionStage,
 381680            payload,
 381681            DeserializeMetadata(entity.EnvelopeMetadataJson));
 682
 381683        GovernanceEmissionError? lastError = string.IsNullOrWhiteSpace(entity.LastErrorCode) || string.IsNullOrWhiteSpac
 381684            ? null
 381685            : GovernanceEmissionError.Create(
 381686                entity.LastErrorCode,
 381687                entity.LastErrorMessage,
 381688                entity.LastErrorIsRetryable ?? false,
 381689                entity.LastErrorProviderName,
 381690                entity.LastErrorProviderErrorCode);
 691
 381692        return GovernanceOutboxEntry.Restore(
 381693            envelope,
 381694            entity.Status,
 381695            entity.OutboxEntryId,
 381696            entity.CreatedUtc,
 381697            entity.UpdatedUtc,
 381698            entity.RetryCount,
 381699            entity.MaxRetryCount,
 381700            entity.NextRetryUtc,
 381701            lastError,
 381702            entity.ProviderName,
 381703            entity.ProviderRecordId,
 381704            entity.DeadLetterReason,
 381705            DeserializeMetadata(entity.MetadataJson),
 381706            entity.ClaimOwner,
 381707            entity.ClaimToken,
 381708            entity.ClaimedUtc,
 381709            entity.ClaimExpiresUtc,
 381710            entity.ClaimAttemptCount);
 711    }
 712
 713    private static GovernanceOutboxClaim CreateClaim(GovernanceOutboxEntry entry)
 714    {
 208715        return GovernanceOutboxClaim.Create(
 208716            entry,
 208717            entry.ClaimOwner ?? throw new InvalidOperationException("Claimed entry is missing claim owner."),
 208718            entry.ClaimToken ?? throw new InvalidOperationException("Claimed entry is missing claim token."),
 208719            entry.ClaimedUtc ?? throw new InvalidOperationException("Claimed entry is missing claimed timestamp."),
 208720            entry.ClaimExpiresUtc ?? throw new InvalidOperationException("Claimed entry is missing claim expiration time
 721    }
 722
 723    internal static bool IsPendingClaimEligible(GovernanceOutboxEntryEntity entity, DateTimeOffset utcNow)
 724    {
 0725        return entity.Status is GovernanceEmissionStatus.Pending && IsClaimAvailable(entity, utcNow);
 726    }
 727
 728    internal static bool IsRetryReadyClaimEligible(GovernanceOutboxEntryEntity entity, DateTimeOffset utcNow)
 729    {
 20730        return (entity.Status is GovernanceEmissionStatus.Deferred or GovernanceEmissionStatus.Failed or GovernanceEmiss
 20731            && (entity.NextRetryUtc is null || entity.NextRetryUtc <= utcNow.ToUniversalTime())
 20732            && IsClaimAvailable(entity, utcNow);
 733    }
 734
 735    private static bool IsClaimAvailable(GovernanceOutboxEntryEntity entity, DateTimeOffset utcNow)
 736    {
 16737        return entity.ClaimToken is null || entity.ClaimExpiresUtc is null || entity.ClaimExpiresUtc <= utcNow.ToUnivers
 738    }
 739
 740    private static bool IsTerminal(GovernanceOutboxEntry entry)
 741    {
 81742        return entry.IsDelivered || entry.IsDeadLettered;
 743    }
 744
 745    private static void DetachEntries(DbUpdateConcurrencyException exception)
 746    {
 20747        foreach (Microsoft.EntityFrameworkCore.ChangeTracking.EntityEntry entry in exception.Entries)
 748        {
 5749            entry.State = EntityState.Detached;
 750        }
 5751    }
 752
 753    private static ReadOnlyDictionary<string, string> DeserializeMetadata(string? json)
 754    {
 788755        if (string.IsNullOrWhiteSpace(json))
 756        {
 1757            return EmptyMetadata();
 758        }
 759
 787760        Dictionary<string, string>? metadata = JsonSerializer.Deserialize<Dictionary<string, string>>(json, JsonOptions)
 761
 787762        return metadata is null || metadata.Count == 0
 787763            ? EmptyMetadata()
 787764            : new ReadOnlyDictionary<string, string>(new Dictionary<string, string>(metadata, StringComparer.Ordinal));
 765    }
 766
 767    private static ReadOnlyDictionary<string, string> EmptyMetadata()
 768    {
 895769        return new ReadOnlyDictionary<string, string>(new Dictionary<string, string>(StringComparer.Ordinal));
 770    }
 771
 772    private static int NormalizeMaxCount(int maxCount)
 773    {
 30774        return maxCount <= 0
 30775            ? throw new ArgumentOutOfRangeException(nameof(maxCount), maxCount, "Maximum count must be greater than zero
 30776            : maxCount;
 777    }
 778}

Methods/Properties

.cctor()
.ctor(Microsoft.EntityFrameworkCore.DbContext)
EnqueueAsync()
SaveAsync()
FindByOutboxEntryIdAsync()
FindPendingAsync()
FindRetryReadyAsync()
ClaimPendingAsync()
ClaimRetryReadyAsync()
MarkDeliveredAsync()
MarkClaimDeliveredAsync()
MarkFailedAsync()
MarkClaimFailedAsync()
MarkDeadLetteredAsync()
MarkClaimDeadLetteredAsync()
SaveClaimAsync()
ReleaseClaimAsync()
ClaimEntriesAsync()
ReconcileTrackedEntries(System.Collections.Generic.List`1<AsiBackbone.EntityFrameworkCore.Persistence.GovernanceOutboxEntryEntity>)
MergeClaimColumns(Microsoft.EntityFrameworkCore.ChangeTracking.EntityEntry`1<AsiBackbone.EntityFrameworkCore.Persistence.GovernanceOutboxEntryEntity>,AsiBackbone.EntityFrameworkCore.Persistence.GovernanceOutboxEntryEntity)
ClaimColumnValues(AsiBackbone.EntityFrameworkCore.Persistence.GovernanceOutboxEntryEntity)
UpdateClaimedEntryAsync()
ApplyEntryUpdateAsync()
RequireEntryAsync()
RequireEntityAsync()
OutboxEntries()
ToEntity(AsiBackbone.Core.Outbox.GovernanceOutboxEntry)
ToEntries(System.Collections.Generic.IEnumerable`1<AsiBackbone.EntityFrameworkCore.Persistence.GovernanceOutboxEntryEntity>)
ToEntry(AsiBackbone.EntityFrameworkCore.Persistence.GovernanceOutboxEntryEntity)
CreateClaim(AsiBackbone.Core.Outbox.GovernanceOutboxEntry)
IsPendingClaimEligible(AsiBackbone.EntityFrameworkCore.Persistence.GovernanceOutboxEntryEntity,System.DateTimeOffset)
IsRetryReadyClaimEligible(AsiBackbone.EntityFrameworkCore.Persistence.GovernanceOutboxEntryEntity,System.DateTimeOffset)
IsClaimAvailable(AsiBackbone.EntityFrameworkCore.Persistence.GovernanceOutboxEntryEntity,System.DateTimeOffset)
IsTerminal(AsiBackbone.Core.Outbox.GovernanceOutboxEntry)
DetachEntries(Microsoft.EntityFrameworkCore.DbUpdateConcurrencyException)
DeserializeMetadata(System.String)
EmptyMetadata()
NormalizeMaxCount(System.Int32)