< 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
98%
Covered lines: 351
Uncovered lines: 4
Coverable lines: 355
Total lines: 665
Line coverage: 98.8%
Branch coverage
90%
Covered branches: 99
Total branches: 110
Branch coverage: 90%
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%22100%
FindByOutboxEntryIdAsync()100%22100%
FindPendingAsync()100%11100%
FindRetryReadyAsync()100%11100%
ClaimPendingAsync()100%11100%
ClaimRetryReadyAsync()100%11100%
MarkDeliveredAsync()100%11100%
MarkClaimDeliveredAsync()100%11100%
MarkFailedAsync()100%11100%
MarkClaimFailedAsync()100%11100%
MarkDeadLetteredAsync()100%11100%
MarkClaimDeadLetteredAsync()100%11100%
SaveClaimAsync()100%22100%
ReleaseClaimAsync()50%6686.66%
ClaimEntriesAsync()100%66100%
TryClaimAsync()50%6688.88%
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(...)50%22100%
IsRetryReadyClaimEligible(...)93.75%1616100%
IsClaimAvailable(...)83.33%66100%
IsTerminal(...)50%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.Text.Json;
 3using AsiBackbone.Core.Emissions;
 4using AsiBackbone.Core.Entities;
 5using AsiBackbone.Core.Outbox;
 6using AsiBackbone.EntityFrameworkCore.Persistence;
 7using Microsoft.EntityFrameworkCore;
 8
 9namespace AsiBackbone.EntityFrameworkCore.Outbox;
 10
 11/// <summary>
 12/// Entity Framework Core-backed governance outbox store that persists provider-neutral emission envelopes through a hos
 13/// </summary>
 14/// <remarks>
 15/// This store provides durable local storage only. Provider delivery, telemetry export, SIEM routing, and cloud emissio
 16/// </remarks>
 17public sealed class EfCoreGovernanceOutboxStore : IAsiBackboneGovernanceOutboxClaimStore
 18{
 219    private static readonly JsonSerializerOptions JsonOptions = new(JsonSerializerDefaults.Web);
 20
 21    private readonly DbContext dbContext;
 22
 23    /// <summary>
 24    /// Initializes a new instance of the <see cref="EfCoreGovernanceOutboxStore" /> class.
 25    /// </summary>
 26    /// <param name="dbContext">The host-owned database context.</param>
 18627    public EfCoreGovernanceOutboxStore(DbContext dbContext)
 28    {
 18629        ArgumentNullException.ThrowIfNull(dbContext);
 30
 18631        this.dbContext = dbContext;
 18632    }
 33
 34    /// <inheritdoc />
 35    public async ValueTask<GovernanceOutboxEntry> EnqueueAsync(
 36        GovernanceEmissionEnvelope envelope,
 37        CancellationToken cancellationToken = default)
 38    {
 6039        ArgumentNullException.ThrowIfNull(envelope);
 6040        cancellationToken.ThrowIfCancellationRequested();
 41
 6042        var entry = GovernanceOutboxEntry.Create(envelope);
 43
 6044        _ = dbContext
 6045            .Set<AsiBackboneGovernanceOutboxEntryEntity>()
 6046            .Add(ToEntity(entry));
 47
 6048        _ = await dbContext.SaveChangesAsync(cancellationToken).ConfigureAwait(false);
 49
 6050        return entry;
 6051    }
 52
 53    /// <inheritdoc />
 54    public async ValueTask<GovernanceOutboxEntry> SaveAsync(
 55        GovernanceOutboxEntry entry,
 56        CancellationToken cancellationToken = default)
 57    {
 45658        ArgumentNullException.ThrowIfNull(entry);
 45659        cancellationToken.ThrowIfCancellationRequested();
 60
 45661        AsiBackboneGovernanceOutboxEntryEntity persistedEntity = ToEntity(entry);
 45662        AsiBackboneGovernanceOutboxEntryEntity? existingEntity = await dbContext
 45663            .Set<AsiBackboneGovernanceOutboxEntryEntity>()
 45664            .SingleOrDefaultAsync(entity => entity.OutboxEntryId == entry.OutboxEntryId, cancellationToken)
 45665            .ConfigureAwait(false);
 66
 45667        if (existingEntity is null)
 68        {
 41269            _ = dbContext
 41270                .Set<AsiBackboneGovernanceOutboxEntryEntity>()
 41271                .Add(persistedEntity);
 72        }
 73        else
 74        {
 4475            persistedEntity.Id = existingEntity.Id;
 4476            persistedEntity.ConcurrencyStamp = AsiBackboneEntity.NewConcurrencyStamp();
 4477            dbContext.Entry(existingEntity).CurrentValues.SetValues(persistedEntity);
 78        }
 79
 45680        _ = await dbContext.SaveChangesAsync(cancellationToken).ConfigureAwait(false);
 81
 44882        return entry;
 44883    }
 84
 85    /// <inheritdoc />
 86    public async ValueTask<GovernanceOutboxEntry?> FindByOutboxEntryIdAsync(
 87        string outboxEntryId,
 88        CancellationToken cancellationToken = default)
 89    {
 12490        ArgumentException.ThrowIfNullOrWhiteSpace(outboxEntryId);
 91
 12492        string normalizedOutboxEntryId = outboxEntryId.Trim();
 93
 12494        AsiBackboneGovernanceOutboxEntryEntity? entity = await OutboxEntries()
 12495            .Where(outboxEntry => outboxEntry.OutboxEntryId == normalizedOutboxEntryId)
 12496            .SingleOrDefaultAsync(cancellationToken)
 12497            .ConfigureAwait(false);
 98
 12499        return entity is null ? null : ToEntry(entity);
 124100    }
 101
 102    /// <inheritdoc />
 103    public async ValueTask<IReadOnlyList<GovernanceOutboxEntry>> FindPendingAsync(
 104        int maxCount = 100,
 105        CancellationToken cancellationToken = default)
 106    {
 18107        int normalizedMaxCount = NormalizeMaxCount(maxCount);
 108
 16109        List<AsiBackboneGovernanceOutboxEntryEntity> entities = await OutboxEntries()
 16110            .Where(outboxEntry => outboxEntry.Status == GovernanceEmissionStatus.Pending)
 16111            .OrderBy(outboxEntry => outboxEntry.CreatedUtc)
 16112            .ThenBy(outboxEntry => outboxEntry.OutboxEntryId)
 16113            .Take(normalizedMaxCount)
 16114            .ToListAsync(cancellationToken)
 16115            .ConfigureAwait(false);
 116
 16117        return ToEntries(entities);
 16118    }
 119
 120    /// <inheritdoc />
 121    public async ValueTask<IReadOnlyList<GovernanceOutboxEntry>> FindRetryReadyAsync(
 122        DateTimeOffset utcNow,
 123        int maxCount = 100,
 124        CancellationToken cancellationToken = default)
 125    {
 10126        int normalizedMaxCount = NormalizeMaxCount(maxCount);
 8127        DateTimeOffset normalizedUtcNow = utcNow.ToUniversalTime();
 128
 8129        List<AsiBackboneGovernanceOutboxEntryEntity> entities = await OutboxEntries()
 8130            .Where(outboxEntry =>
 8131                outboxEntry.Status == GovernanceEmissionStatus.Deferred ||
 8132                outboxEntry.Status == GovernanceEmissionStatus.Failed ||
 8133                outboxEntry.Status == GovernanceEmissionStatus.RetryableFailure)
 8134            .Where(outboxEntry => outboxEntry.RetryCount < outboxEntry.MaxRetryCount)
 8135            .Where(outboxEntry => outboxEntry.NextRetryUtc == null || outboxEntry.NextRetryUtc <= normalizedUtcNow)
 8136            .OrderBy(outboxEntry => outboxEntry.NextRetryUtc ?? outboxEntry.UpdatedUtc)
 8137            .ThenBy(outboxEntry => outboxEntry.OutboxEntryId)
 8138            .Take(normalizedMaxCount)
 8139            .ToListAsync(cancellationToken)
 8140            .ConfigureAwait(false);
 141
 8142        return ToEntries(entities);
 8143    }
 144
 145    /// <inheritdoc />
 146    public async ValueTask<IReadOnlyList<GovernanceOutboxClaim>> ClaimPendingAsync(
 147        GovernanceOutboxClaimRequest request,
 148        CancellationToken cancellationToken = default)
 149    {
 60150        ArgumentNullException.ThrowIfNull(request);
 60151        cancellationToken.ThrowIfCancellationRequested();
 152
 60153        List<string> candidateIds = await OutboxEntries()
 60154            .Where(outboxEntry => outboxEntry.Status == GovernanceEmissionStatus.Pending)
 60155            .Where(outboxEntry => outboxEntry.ClaimToken == null || outboxEntry.ClaimExpiresUtc == null || outboxEntry.C
 60156            .OrderBy(outboxEntry => outboxEntry.CreatedUtc)
 60157            .ThenBy(outboxEntry => outboxEntry.OutboxEntryId)
 60158            .Select(outboxEntry => outboxEntry.OutboxEntryId)
 60159            .Take(request.MaxCount)
 60160            .ToListAsync(cancellationToken)
 60161            .ConfigureAwait(false);
 162
 60163        return await ClaimEntriesAsync(candidateIds, request, IsPendingClaimEligible, cancellationToken).ConfigureAwait(
 60164    }
 165
 166    /// <inheritdoc />
 167    public async ValueTask<IReadOnlyList<GovernanceOutboxClaim>> ClaimRetryReadyAsync(
 168        GovernanceOutboxClaimRequest request,
 169        CancellationToken cancellationToken = default)
 170    {
 10171        ArgumentNullException.ThrowIfNull(request);
 8172        cancellationToken.ThrowIfCancellationRequested();
 173
 6174        List<string> candidateIds = await OutboxEntries()
 6175            .Where(outboxEntry =>
 6176                outboxEntry.Status == GovernanceEmissionStatus.Deferred ||
 6177                outboxEntry.Status == GovernanceEmissionStatus.Failed ||
 6178                outboxEntry.Status == GovernanceEmissionStatus.RetryableFailure)
 6179            .Where(outboxEntry => outboxEntry.RetryCount < outboxEntry.MaxRetryCount)
 6180            .Where(outboxEntry => outboxEntry.NextRetryUtc == null || outboxEntry.NextRetryUtc <= request.UtcNow)
 6181            .Where(outboxEntry => outboxEntry.ClaimToken == null || outboxEntry.ClaimExpiresUtc == null || outboxEntry.C
 6182            .OrderBy(outboxEntry => outboxEntry.NextRetryUtc ?? outboxEntry.UpdatedUtc)
 6183            .ThenBy(outboxEntry => outboxEntry.OutboxEntryId)
 6184            .Select(outboxEntry => outboxEntry.OutboxEntryId)
 6185            .Take(request.MaxCount)
 6186            .ToListAsync(cancellationToken)
 6187            .ConfigureAwait(false);
 188
 6189        return await ClaimEntriesAsync(candidateIds, request, IsRetryReadyClaimEligible, cancellationToken).ConfigureAwa
 6190    }
 191
 192    /// <inheritdoc />
 193    public async ValueTask<GovernanceOutboxEntry> MarkDeliveredAsync(
 194        string outboxEntryId,
 195        GovernanceEmissionResult result,
 196        CancellationToken cancellationToken = default)
 197    {
 12198        ArgumentException.ThrowIfNullOrWhiteSpace(outboxEntryId);
 12199        ArgumentNullException.ThrowIfNull(result);
 200
 12201        GovernanceOutboxEntry entry = await RequireEntryAsync(outboxEntryId, cancellationToken).ConfigureAwait(false);
 12202        GovernanceOutboxEntry updatedEntry = entry.MarkDelivered(result);
 203
 12204        return await SaveAsync(updatedEntry, cancellationToken).ConfigureAwait(false);
 8205    }
 206
 207    /// <inheritdoc />
 208    public async ValueTask<GovernanceOutboxEntry> MarkClaimDeliveredAsync(
 209        GovernanceOutboxClaim claim,
 210        GovernanceEmissionResult result,
 211        CancellationToken cancellationToken = default)
 212    {
 24213        ArgumentNullException.ThrowIfNull(claim);
 22214        ArgumentNullException.ThrowIfNull(result);
 215
 34216        return await UpdateClaimedEntryAsync(claim, entry => entry.MarkDelivered(result), cancellationToken).ConfigureAw
 18217    }
 218
 219    /// <inheritdoc />
 220    public async ValueTask<GovernanceOutboxEntry> MarkFailedAsync(
 221        string outboxEntryId,
 222        GovernanceEmissionError governanceEmissionError,
 223        DateTimeOffset? nextRetryUtc = null,
 224        CancellationToken cancellationToken = default)
 225    {
 20226        ArgumentException.ThrowIfNullOrWhiteSpace(outboxEntryId);
 20227        ArgumentNullException.ThrowIfNull(governanceEmissionError);
 228
 20229        GovernanceOutboxEntry entry = await RequireEntryAsync(outboxEntryId, cancellationToken).ConfigureAwait(false);
 20230        GovernanceOutboxEntry updatedEntry = entry.MarkFailed(governanceEmissionError, nextRetryUtc);
 231
 20232        return await SaveAsync(updatedEntry, cancellationToken).ConfigureAwait(false);
 18233    }
 234
 235    /// <inheritdoc />
 236    public async ValueTask<GovernanceOutboxEntry> MarkClaimFailedAsync(
 237        GovernanceOutboxClaim claim,
 238        GovernanceEmissionError governanceEmissionError,
 239        DateTimeOffset? nextRetryUtc = null,
 240        CancellationToken cancellationToken = default)
 241    {
 18242        ArgumentNullException.ThrowIfNull(claim);
 16243        ArgumentNullException.ThrowIfNull(governanceEmissionError);
 244
 14245        return await UpdateClaimedEntryAsync(
 14246            claim,
 6247            entry => entry.MarkFailed(governanceEmissionError, nextRetryUtc),
 14248            cancellationToken)
 14249            .ConfigureAwait(false);
 14250    }
 251
 252    /// <inheritdoc />
 253    public async ValueTask<GovernanceOutboxEntry> MarkDeadLetteredAsync(
 254        string outboxEntryId,
 255        GovernanceEmissionError governanceEmissionError,
 256        string? deadLetterReason = null,
 257        CancellationToken cancellationToken = default)
 258    {
 6259        ArgumentException.ThrowIfNullOrWhiteSpace(outboxEntryId);
 6260        ArgumentNullException.ThrowIfNull(governanceEmissionError);
 261
 6262        GovernanceOutboxEntry entry = await RequireEntryAsync(outboxEntryId, cancellationToken).ConfigureAwait(false);
 6263        GovernanceOutboxEntry updatedEntry = entry.MarkDeadLettered(governanceEmissionError, deadLetterReason);
 264
 6265        return await SaveAsync(updatedEntry, cancellationToken).ConfigureAwait(false);
 6266    }
 267
 268    /// <inheritdoc />
 269    public async ValueTask<GovernanceOutboxEntry> MarkClaimDeadLetteredAsync(
 270        GovernanceOutboxClaim claim,
 271        GovernanceEmissionError governanceEmissionError,
 272        string? deadLetterReason = null,
 273        CancellationToken cancellationToken = default)
 274    {
 10275        ArgumentNullException.ThrowIfNull(claim);
 8276        ArgumentNullException.ThrowIfNull(governanceEmissionError);
 277
 6278        return await UpdateClaimedEntryAsync(
 6279            claim,
 6280            entry => entry.MarkDeadLettered(governanceEmissionError, deadLetterReason),
 6281            cancellationToken)
 6282            .ConfigureAwait(false);
 6283    }
 284
 285    /// <inheritdoc />
 286    public async ValueTask<GovernanceOutboxEntry> SaveClaimAsync(
 287        GovernanceOutboxClaim claim,
 288        GovernanceOutboxEntry entry,
 289        CancellationToken cancellationToken = default)
 290    {
 10291        ArgumentNullException.ThrowIfNull(claim);
 8292        ArgumentNullException.ThrowIfNull(entry);
 293
 6294        return !string.Equals(claim.OutboxEntryId, entry.OutboxEntryId, StringComparison.Ordinal)
 6295            ? throw new ArgumentException("Claim and entry must reference the same outbox entry ID.", nameof(entry))
 8296            : await UpdateClaimedEntryAsync(claim, _ => entry, cancellationToken).ConfigureAwait(false);
 2297    }
 298
 299    /// <inheritdoc />
 300    public async ValueTask<GovernanceOutboxEntry?> ReleaseClaimAsync(
 301        GovernanceOutboxClaim claim,
 302        string? reason = null,
 303        CancellationToken cancellationToken = default)
 304    {
 4305        ArgumentNullException.ThrowIfNull(claim);
 4306        cancellationToken.ThrowIfCancellationRequested();
 307
 4308        AsiBackboneGovernanceOutboxEntryEntity? entity = await dbContext
 4309            .Set<AsiBackboneGovernanceOutboxEntryEntity>()
 4310            .SingleOrDefaultAsync(outboxEntry => outboxEntry.OutboxEntryId == claim.OutboxEntryId, cancellationToken)
 4311            .ConfigureAwait(false);
 312
 4313        if (entity is null)
 314        {
 0315            return null;
 316        }
 317
 4318        GovernanceOutboxEntry currentEntry = ToEntry(entity);
 4319        if (!currentEntry.IsClaimedBy(claim) || IsTerminal(currentEntry))
 320        {
 0321            return currentEntry;
 322        }
 323
 4324        GovernanceOutboxEntry releasedEntry = currentEntry.ReleaseClaim();
 4325        await ApplyEntryUpdateAsync(entity, releasedEntry, cancellationToken).ConfigureAwait(false);
 326
 4327        return releasedEntry;
 4328    }
 329
 330    private async ValueTask<IReadOnlyList<GovernanceOutboxClaim>> ClaimEntriesAsync(
 331        List<string> candidateIds,
 332        GovernanceOutboxClaimRequest request,
 333        Func<AsiBackboneGovernanceOutboxEntryEntity, DateTimeOffset, bool> isEligible,
 334        CancellationToken cancellationToken)
 335    {
 66336        List<GovernanceOutboxClaim> claims = new(Math.Min(request.MaxCount, candidateIds.Count));
 337
 904338        foreach (string candidateId in candidateIds)
 339        {
 386340            cancellationToken.ThrowIfCancellationRequested();
 341
 386342            if (claims.Count >= request.MaxCount)
 343            {
 344                break;
 345            }
 346
 386347            GovernanceOutboxClaim? claim = await TryClaimAsync(candidateId, request, isEligible, cancellationToken).Conf
 386348            if (claim is not null)
 349            {
 384350                claims.Add(claim);
 351            }
 352        }
 353
 66354        return claims;
 66355    }
 356
 357    private async ValueTask<GovernanceOutboxClaim?> TryClaimAsync(
 358        string outboxEntryId,
 359        GovernanceOutboxClaimRequest request,
 360        Func<AsiBackboneGovernanceOutboxEntryEntity, DateTimeOffset, bool> isEligible,
 361        CancellationToken cancellationToken)
 362    {
 386363        AsiBackboneGovernanceOutboxEntryEntity? entity = await dbContext
 386364            .Set<AsiBackboneGovernanceOutboxEntryEntity>()
 386365            .SingleOrDefaultAsync(outboxEntry => outboxEntry.OutboxEntryId == outboxEntryId, cancellationToken)
 386366            .ConfigureAwait(false);
 367
 386368        if (entity is null || !isEligible(entity, request.UtcNow))
 369        {
 0370            return null;
 371        }
 372
 386373        GovernanceOutboxEntry currentEntry = ToEntry(entity);
 386374        if (!currentEntry.CanBeClaimed(request.UtcNow))
 375        {
 0376            return null;
 377        }
 378
 386379        GovernanceOutboxEntry claimedEntry = currentEntry.MarkClaimed(
 386380            request.WorkerId,
 386381            claimedUtc: request.UtcNow,
 386382            leaseDuration: request.LeaseDuration);
 383
 384        try
 385        {
 386386            await ApplyEntryUpdateAsync(entity, claimedEntry, cancellationToken).ConfigureAwait(false);
 384387            return CreateClaim(claimedEntry);
 388        }
 389        catch (DbUpdateConcurrencyException exception)
 390        {
 2391            DetachEntries(exception);
 2392            return null;
 393        }
 386394    }
 395
 396    private async ValueTask<GovernanceOutboxEntry> UpdateClaimedEntryAsync(
 397        GovernanceOutboxClaim claim,
 398        Func<GovernanceOutboxEntry, GovernanceOutboxEntry> updateEntry,
 399        CancellationToken cancellationToken)
 400    {
 44401        AsiBackboneGovernanceOutboxEntryEntity entity = await RequireEntityAsync(claim.OutboxEntryId, cancellationToken)
 40402        GovernanceOutboxEntry currentEntry = ToEntry(entity);
 403
 40404        if (!currentEntry.IsClaimedBy(claim) || IsTerminal(currentEntry))
 405        {
 12406            return currentEntry;
 407        }
 408
 28409        GovernanceOutboxEntry updatedEntry = updateEntry(currentEntry);
 410
 411        try
 412        {
 28413            await ApplyEntryUpdateAsync(entity, updatedEntry, cancellationToken).ConfigureAwait(false);
 18414            return updatedEntry;
 415        }
 416        catch (DbUpdateConcurrencyException exception)
 417        {
 10418            DetachEntries(exception);
 10419            GovernanceOutboxEntry? refreshedEntry = await FindByOutboxEntryIdAsync(claim.OutboxEntryId, cancellationToke
 10420            return refreshedEntry ?? currentEntry;
 421        }
 40422    }
 423
 424    private async ValueTask ApplyEntryUpdateAsync(
 425        AsiBackboneGovernanceOutboxEntryEntity entity,
 426        GovernanceOutboxEntry entry,
 427        CancellationToken cancellationToken)
 428    {
 418429        AsiBackboneGovernanceOutboxEntryEntity persistedEntity = ToEntity(entry);
 418430        persistedEntity.Id = entity.Id;
 418431        persistedEntity.ConcurrencyStamp = AsiBackboneEntity.NewConcurrencyStamp();
 418432        dbContext.Entry(entity).CurrentValues.SetValues(persistedEntity);
 433
 418434        _ = await dbContext.SaveChangesAsync(cancellationToken).ConfigureAwait(false);
 406435    }
 436
 437    private async ValueTask<GovernanceOutboxEntry> RequireEntryAsync(
 438        string outboxEntryId,
 439        CancellationToken cancellationToken)
 440    {
 38441        GovernanceOutboxEntry? entry = await FindByOutboxEntryIdAsync(outboxEntryId, cancellationToken).ConfigureAwait(f
 442
 38443        return entry ?? throw new InvalidOperationException($"Outbox entry '{outboxEntryId.Trim()}' was not found.");
 38444    }
 445
 446    private async ValueTask<AsiBackboneGovernanceOutboxEntryEntity> RequireEntityAsync(
 447        string outboxEntryId,
 448        CancellationToken cancellationToken)
 449    {
 44450        AsiBackboneGovernanceOutboxEntryEntity? entity = await dbContext
 44451            .Set<AsiBackboneGovernanceOutboxEntryEntity>()
 44452            .SingleOrDefaultAsync(outboxEntry => outboxEntry.OutboxEntryId == outboxEntryId.Trim(), cancellationToken)
 44453            .ConfigureAwait(false);
 454
 42455        return entity ?? throw new InvalidOperationException($"Outbox entry '{outboxEntryId.Trim()}' was not found.");
 40456    }
 457
 458    private IQueryable<AsiBackboneGovernanceOutboxEntryEntity> OutboxEntries()
 459    {
 214460        return dbContext.Set<AsiBackboneGovernanceOutboxEntryEntity>().AsNoTracking();
 461    }
 462
 463    private static AsiBackboneGovernanceOutboxEntryEntity ToEntity(GovernanceOutboxEntry entry)
 464    {
 934465        GovernanceEmissionEnvelope envelope = entry.Envelope;
 934466        GovernanceEmissionPayload? payload = envelope.Payload;
 934467        GovernanceEmissionError? lastError = entry.LastError;
 468
 934469        return new AsiBackboneGovernanceOutboxEntryEntity
 934470        {
 934471            OutboxEntryId = entry.OutboxEntryId,
 934472            Status = entry.Status,
 934473            CreatedUtc = entry.CreatedUtc,
 934474            UpdatedUtc = entry.UpdatedUtc,
 934475            DeliveredUtc = entry.Status is GovernanceEmissionStatus.Delivered ? entry.UpdatedUtc : null,
 934476            RetryCount = entry.RetryCount,
 934477            MaxRetryCount = entry.MaxRetryCount,
 934478            NextRetryUtc = entry.NextRetryUtc,
 934479            ProviderName = entry.ProviderName,
 934480            ProviderRecordId = entry.ProviderRecordId,
 934481            DeadLetterReason = entry.DeadLetterReason,
 934482            LastErrorCode = lastError?.Code,
 934483            LastErrorMessage = lastError?.Message,
 934484            LastErrorIsRetryable = lastError?.IsRetryable,
 934485            LastErrorProviderName = lastError?.ProviderName,
 934486            LastErrorProviderErrorCode = lastError?.ProviderErrorCode,
 934487            MetadataJson = JsonSerializer.Serialize(entry.Metadata, JsonOptions),
 934488            ClaimOwner = entry.ClaimOwner,
 934489            ClaimToken = entry.ClaimToken,
 934490            ClaimedUtc = entry.ClaimedUtc,
 934491            ClaimExpiresUtc = entry.ClaimExpiresUtc,
 934492            ClaimAttemptCount = entry.ClaimAttemptCount,
 934493            EnvelopeId = envelope.EnvelopeId,
 934494            EnvelopeSchemaVersion = envelope.SchemaVersion,
 934495            EnvelopeEventType = envelope.EventType,
 934496            EnvelopeEventId = envelope.EventId,
 934497            EnvelopeOccurredUtc = envelope.OccurredUtc,
 934498            EnvelopeCreatedUtc = envelope.CreatedUtc,
 934499            EnvelopeCorrelationId = envelope.CorrelationId,
 934500            EnvelopeAuditResidueId = envelope.AuditResidueId,
 934501            EnvelopeLifecycleStage = envelope.LifecycleStage,
 934502            EnvelopeLifecycleStageSequence = envelope.LifecycleStageSequence,
 934503            EnvelopePolicyVersion = envelope.PolicyVersion,
 934504            EnvelopePolicyHash = envelope.PolicyHash,
 934505            EnvelopeTraceId = envelope.TraceId,
 934506            EnvelopeSpanId = envelope.SpanId,
 934507            EnvelopeParentSpanId = envelope.ParentSpanId,
 934508            EnvelopeOperationName = envelope.OperationName,
 934509            EnvelopeOutcome = envelope.Outcome,
 934510            EnvelopeActorId = envelope.ActorId,
 934511            EnvelopeEmitterStatus = envelope.EmitterStatus,
 934512            EnvelopeEmitterProvider = envelope.EmitterProvider,
 934513            EnvelopeOutboxSequence = envelope.OutboxSequence,
 934514            EnvelopeGatewayExecutionId = envelope.GatewayExecutionId,
 934515            EnvelopeDecisionStage = envelope.DecisionStage,
 934516            EnvelopeMetadataJson = JsonSerializer.Serialize(envelope.Metadata, JsonOptions),
 934517            EnvelopePayloadType = payload?.PayloadType,
 934518            EnvelopePayloadSchemaVersion = payload?.SchemaVersion,
 934519            EnvelopePayloadContentType = payload?.ContentType,
 934520            EnvelopePayloadContentHash = payload?.ContentHash,
 934521            EnvelopePayloadSizeBytes = payload?.SizeBytes,
 934522            EnvelopePayloadMetadataJson = JsonSerializer.Serialize(payload?.Metadata ?? EmptyMetadata(), JsonOptions)
 934523        };
 524    }
 525
 526    private static GovernanceOutboxEntry[] ToEntries(IEnumerable<AsiBackboneGovernanceOutboxEntryEntity> entities)
 527    {
 24528        return [.. entities.Select(ToEntry)];
 529    }
 530
 531    private static GovernanceOutboxEntry ToEntry(AsiBackboneGovernanceOutboxEntryEntity entity)
 532    {
 614533        GovernanceEmissionPayload? payload = string.IsNullOrWhiteSpace(entity.EnvelopePayloadType)
 614534            ? null
 614535            : GovernanceEmissionPayload.Create(
 614536                entity.EnvelopePayloadType,
 614537                entity.EnvelopePayloadSchemaVersion,
 614538                entity.EnvelopePayloadContentType,
 614539                entity.EnvelopePayloadContentHash,
 614540                entity.EnvelopePayloadSizeBytes,
 614541                DeserializeMetadata(entity.EnvelopePayloadMetadataJson));
 542
 614543        var envelope = GovernanceEmissionEnvelope.Create(
 614544            entity.EnvelopeEventType,
 614545            entity.EnvelopeEventId,
 614546            entity.EnvelopeOccurredUtc,
 614547            entity.EnvelopeId,
 614548            entity.EnvelopeCreatedUtc,
 614549            entity.EnvelopeSchemaVersion,
 614550            entity.EnvelopeCorrelationId,
 614551            entity.EnvelopeAuditResidueId,
 614552            entity.EnvelopeLifecycleStage,
 614553            entity.EnvelopePolicyVersion,
 614554            entity.EnvelopePolicyHash,
 614555            entity.EnvelopeTraceId,
 614556            entity.EnvelopeSpanId,
 614557            entity.EnvelopeParentSpanId,
 614558            entity.EnvelopeOperationName,
 614559            entity.EnvelopeOutcome,
 614560            entity.EnvelopeActorId,
 614561            entity.EnvelopeEmitterStatus,
 614562            entity.EnvelopeEmitterProvider,
 614563            entity.EnvelopeOutboxSequence,
 614564            entity.EnvelopeGatewayExecutionId,
 614565            entity.EnvelopeDecisionStage,
 614566            payload,
 614567            DeserializeMetadata(entity.EnvelopeMetadataJson));
 568
 614569        GovernanceEmissionError? lastError = string.IsNullOrWhiteSpace(entity.LastErrorCode) || string.IsNullOrWhiteSpac
 614570            ? null
 614571            : GovernanceEmissionError.Create(
 614572                entity.LastErrorCode,
 614573                entity.LastErrorMessage,
 614574                entity.LastErrorIsRetryable ?? false,
 614575                entity.LastErrorProviderName,
 614576                entity.LastErrorProviderErrorCode);
 577
 614578        return GovernanceOutboxEntry.Restore(
 614579            envelope,
 614580            entity.Status,
 614581            entity.OutboxEntryId,
 614582            entity.CreatedUtc,
 614583            entity.UpdatedUtc,
 614584            entity.RetryCount,
 614585            entity.MaxRetryCount,
 614586            entity.NextRetryUtc,
 614587            lastError,
 614588            entity.ProviderName,
 614589            entity.ProviderRecordId,
 614590            entity.DeadLetterReason,
 614591            DeserializeMetadata(entity.MetadataJson),
 614592            entity.ClaimOwner,
 614593            entity.ClaimToken,
 614594            entity.ClaimedUtc,
 614595            entity.ClaimExpiresUtc,
 614596            entity.ClaimAttemptCount);
 597    }
 598
 599    private static GovernanceOutboxClaim CreateClaim(GovernanceOutboxEntry entry)
 600    {
 392601        return GovernanceOutboxClaim.Create(
 392602            entry,
 392603            entry.ClaimOwner ?? throw new InvalidOperationException("Claimed entry is missing claim owner."),
 392604            entry.ClaimToken ?? throw new InvalidOperationException("Claimed entry is missing claim token."),
 392605            entry.ClaimedUtc ?? throw new InvalidOperationException("Claimed entry is missing claimed timestamp."),
 392606            entry.ClaimExpiresUtc ?? throw new InvalidOperationException("Claimed entry is missing claim expiration time
 607    }
 608
 609    private static bool IsPendingClaimEligible(AsiBackboneGovernanceOutboxEntryEntity entity, DateTimeOffset utcNow)
 610    {
 374611        return entity.Status is GovernanceEmissionStatus.Pending && IsClaimAvailable(entity, utcNow);
 612    }
 613
 614    private static bool IsRetryReadyClaimEligible(AsiBackboneGovernanceOutboxEntryEntity entity, DateTimeOffset utcNow)
 615    {
 52616        return (entity.Status is GovernanceEmissionStatus.Deferred or GovernanceEmissionStatus.Failed or GovernanceEmiss
 52617            && entity.RetryCount < entity.MaxRetryCount
 52618            && (entity.NextRetryUtc is null || entity.NextRetryUtc <= utcNow.ToUniversalTime())
 52619            && IsClaimAvailable(entity, utcNow);
 620    }
 621
 622    private static bool IsClaimAvailable(AsiBackboneGovernanceOutboxEntryEntity entity, DateTimeOffset utcNow)
 623    {
 414624        return entity.ClaimToken is null || entity.ClaimExpiresUtc is null || entity.ClaimExpiresUtc <= utcNow.ToUnivers
 625    }
 626
 627    private static bool IsTerminal(GovernanceOutboxEntry entry)
 628    {
 32629        return entry.IsDelivered || entry.IsDeadLettered;
 630    }
 631
 632    private static void DetachEntries(DbUpdateConcurrencyException exception)
 633    {
 48634        foreach (Microsoft.EntityFrameworkCore.ChangeTracking.EntityEntry entry in exception.Entries)
 635        {
 12636            entry.State = EntityState.Detached;
 637        }
 12638    }
 639
 640    private static ReadOnlyDictionary<string, string> DeserializeMetadata(string? json)
 641    {
 1266642        if (string.IsNullOrWhiteSpace(json))
 643        {
 2644            return EmptyMetadata();
 645        }
 646
 1264647        Dictionary<string, string>? metadata = JsonSerializer.Deserialize<Dictionary<string, string>>(json, JsonOptions)
 648
 1264649        return metadata is null || metadata.Count == 0
 1264650            ? EmptyMetadata()
 1264651            : new ReadOnlyDictionary<string, string>(new Dictionary<string, string>(metadata, StringComparer.Ordinal));
 652    }
 653
 654    private static ReadOnlyDictionary<string, string> EmptyMetadata()
 655    {
 1848656        return new ReadOnlyDictionary<string, string>(new Dictionary<string, string>(StringComparer.Ordinal));
 657    }
 658
 659    private static int NormalizeMaxCount(int maxCount)
 660    {
 28661        return maxCount <= 0
 28662            ? throw new ArgumentOutOfRangeException(nameof(maxCount), maxCount, "Maximum count must be greater than zero
 28663            : maxCount;
 664    }
 665}