| | | 1 | | using System.Collections.ObjectModel; |
| | | 2 | | using System.Linq.Expressions; |
| | | 3 | | using System.Text.Json; |
| | | 4 | | using AsiBackbone.Core.Emissions; |
| | | 5 | | using AsiBackbone.Core.Entities; |
| | | 6 | | using AsiBackbone.Core.Outbox; |
| | | 7 | | using AsiBackbone.EntityFrameworkCore.Persistence; |
| | | 8 | | using Microsoft.EntityFrameworkCore; |
| | | 9 | | |
| | | 10 | | namespace 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> |
| | | 18 | | public sealed class EfCoreGovernanceOutboxStore : IGovernanceOutboxClaimStore |
| | | 19 | | { |
| | 1 | 20 | | 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> |
| | 130 | 28 | | public EfCoreGovernanceOutboxStore(DbContext dbContext) |
| | | 29 | | { |
| | 130 | 30 | | ArgumentNullException.ThrowIfNull(dbContext); |
| | | 31 | | |
| | 130 | 32 | | this.dbContext = dbContext; |
| | 130 | 33 | | } |
| | | 34 | | |
| | | 35 | | /// <inheritdoc /> |
| | | 36 | | public async ValueTask<GovernanceOutboxEntry> EnqueueAsync( |
| | | 37 | | GovernanceEmissionEnvelope envelope, |
| | | 38 | | CancellationToken cancellationToken = default) |
| | | 39 | | { |
| | 46 | 40 | | ArgumentNullException.ThrowIfNull(envelope); |
| | 46 | 41 | | cancellationToken.ThrowIfCancellationRequested(); |
| | | 42 | | |
| | 46 | 43 | | var entry = GovernanceOutboxEntry.Create(envelope); |
| | | 44 | | |
| | 46 | 45 | | _ = dbContext |
| | 46 | 46 | | .Set<GovernanceOutboxEntryEntity>() |
| | 46 | 47 | | .Add(ToEntity(entry)); |
| | | 48 | | |
| | 46 | 49 | | _ = await dbContext.SaveChangesAsync(cancellationToken).ConfigureAwait(false); |
| | | 50 | | |
| | 46 | 51 | | return entry; |
| | 46 | 52 | | } |
| | | 53 | | |
| | | 54 | | /// <inheritdoc /> |
| | | 55 | | public async ValueTask<GovernanceOutboxEntry> SaveAsync( |
| | | 56 | | GovernanceOutboxEntry entry, |
| | | 57 | | CancellationToken cancellationToken = default) |
| | | 58 | | { |
| | 233 | 59 | | ArgumentNullException.ThrowIfNull(entry); |
| | 233 | 60 | | cancellationToken.ThrowIfCancellationRequested(); |
| | | 61 | | |
| | 233 | 62 | | GovernanceOutboxEntryEntity persistedEntity = ToEntity(entry); |
| | 233 | 63 | | GovernanceOutboxEntryEntity? existingEntity = await dbContext |
| | 233 | 64 | | .Set<GovernanceOutboxEntryEntity>() |
| | 233 | 65 | | .SingleOrDefaultAsync(entity => entity.OutboxEntryId == entry.OutboxEntryId, cancellationToken) |
| | 233 | 66 | | .ConfigureAwait(false); |
| | | 67 | | |
| | 233 | 68 | | if (existingEntity is null) |
| | | 69 | | { |
| | 206 | 70 | | _ = dbContext |
| | 206 | 71 | | .Set<GovernanceOutboxEntryEntity>() |
| | 206 | 72 | | .Add(persistedEntity); |
| | | 73 | | } |
| | | 74 | | else |
| | | 75 | | { |
| | 27 | 76 | | GovernanceOutboxEntry currentEntry = ToEntry(existingEntity); |
| | 27 | 77 | | if (IsTerminal(currentEntry)) |
| | | 78 | | { |
| | 2 | 79 | | return currentEntry; |
| | | 80 | | } |
| | | 81 | | |
| | 25 | 82 | | persistedEntity.Id = existingEntity.Id; |
| | 25 | 83 | | persistedEntity.ConcurrencyStamp = GovernanceEntity.NewConcurrencyStamp(); |
| | 25 | 84 | | dbContext.Entry(existingEntity).CurrentValues.SetValues(persistedEntity); |
| | | 85 | | } |
| | | 86 | | |
| | 231 | 87 | | _ = await dbContext.SaveChangesAsync(cancellationToken).ConfigureAwait(false); |
| | | 88 | | |
| | 228 | 89 | | return entry; |
| | 230 | 90 | | } |
| | | 91 | | |
| | | 92 | | /// <inheritdoc /> |
| | | 93 | | public async ValueTask<GovernanceOutboxEntry?> FindByOutboxEntryIdAsync( |
| | | 94 | | string outboxEntryId, |
| | | 95 | | CancellationToken cancellationToken = default) |
| | | 96 | | { |
| | 91 | 97 | | ArgumentException.ThrowIfNullOrWhiteSpace(outboxEntryId); |
| | | 98 | | |
| | 91 | 99 | | string normalizedOutboxEntryId = outboxEntryId.Trim(); |
| | | 100 | | |
| | 91 | 101 | | GovernanceOutboxEntryEntity? entity = await OutboxEntries() |
| | 91 | 102 | | .Where(outboxEntry => outboxEntry.OutboxEntryId == normalizedOutboxEntryId) |
| | 91 | 103 | | .SingleOrDefaultAsync(cancellationToken) |
| | 91 | 104 | | .ConfigureAwait(false); |
| | | 105 | | |
| | 91 | 106 | | return entity is null ? null : ToEntry(entity); |
| | 91 | 107 | | } |
| | | 108 | | |
| | | 109 | | /// <inheritdoc /> |
| | | 110 | | public async ValueTask<IReadOnlyList<GovernanceOutboxEntry>> FindPendingAsync( |
| | | 111 | | int maxCount = 100, |
| | | 112 | | CancellationToken cancellationToken = default) |
| | | 113 | | { |
| | 17 | 114 | | int normalizedMaxCount = NormalizeMaxCount(maxCount); |
| | | 115 | | |
| | 16 | 116 | | List<GovernanceOutboxEntryEntity> entities = await OutboxEntries() |
| | 16 | 117 | | .Where(outboxEntry => outboxEntry.Status == GovernanceEmissionStatus.Pending) |
| | 16 | 118 | | .OrderBy(outboxEntry => outboxEntry.CreatedUtc) |
| | 16 | 119 | | .ThenBy(outboxEntry => outboxEntry.OutboxEntryId) |
| | 16 | 120 | | .Take(normalizedMaxCount) |
| | 16 | 121 | | .ToListAsync(cancellationToken) |
| | 16 | 122 | | .ConfigureAwait(false); |
| | | 123 | | |
| | 16 | 124 | | return ToEntries(entities); |
| | 16 | 125 | | } |
| | | 126 | | |
| | | 127 | | /// <inheritdoc /> |
| | | 128 | | public async ValueTask<IReadOnlyList<GovernanceOutboxEntry>> FindRetryReadyAsync( |
| | | 129 | | DateTimeOffset utcNow, |
| | | 130 | | int maxCount = 100, |
| | | 131 | | CancellationToken cancellationToken = default) |
| | | 132 | | { |
| | 13 | 133 | | int normalizedMaxCount = NormalizeMaxCount(maxCount); |
| | 12 | 134 | | DateTimeOffset normalizedUtcNow = utcNow.ToUniversalTime(); |
| | | 135 | | |
| | 12 | 136 | | List<GovernanceOutboxEntryEntity> entities = await OutboxEntries() |
| | 12 | 137 | | .Where(outboxEntry => |
| | 12 | 138 | | outboxEntry.Status == GovernanceEmissionStatus.Deferred || |
| | 12 | 139 | | outboxEntry.Status == GovernanceEmissionStatus.Failed || |
| | 12 | 140 | | outboxEntry.Status == GovernanceEmissionStatus.RetryableFailure) |
| | 12 | 141 | | .Where(outboxEntry => outboxEntry.NextRetryUtc == null || outboxEntry.NextRetryUtc <= normalizedUtcNow) |
| | 12 | 142 | | .OrderBy(outboxEntry => outboxEntry.NextRetryUtc ?? outboxEntry.UpdatedUtc) |
| | 12 | 143 | | .ThenBy(outboxEntry => outboxEntry.OutboxEntryId) |
| | 12 | 144 | | .Take(normalizedMaxCount) |
| | 12 | 145 | | .ToListAsync(cancellationToken) |
| | 12 | 146 | | .ConfigureAwait(false); |
| | | 147 | | |
| | 12 | 148 | | return ToEntries(entities); |
| | 12 | 149 | | } |
| | | 150 | | |
| | | 151 | | /// <inheritdoc /> |
| | | 152 | | public async ValueTask<IReadOnlyList<GovernanceOutboxClaim>> ClaimPendingAsync( |
| | | 153 | | GovernanceOutboxClaimRequest request, |
| | | 154 | | CancellationToken cancellationToken = default) |
| | | 155 | | { |
| | 40 | 156 | | ArgumentNullException.ThrowIfNull(request); |
| | 40 | 157 | | cancellationToken.ThrowIfCancellationRequested(); |
| | | 158 | | |
| | 40 | 159 | | Expression<Func<GovernanceOutboxEntryEntity, bool>> isEligible = outboxEntry => |
| | 40 | 160 | | outboxEntry.Status == GovernanceEmissionStatus.Pending |
| | 40 | 161 | | && (outboxEntry.ClaimToken == null || outboxEntry.ClaimExpiresUtc == null || outboxEntry.ClaimExpiresUtc <= |
| | | 162 | | |
| | 40 | 163 | | IQueryable<GovernanceOutboxEntryEntity> candidates = OutboxEntries() |
| | 40 | 164 | | .Where(isEligible) |
| | 40 | 165 | | .OrderBy(outboxEntry => outboxEntry.CreatedUtc) |
| | 40 | 166 | | .ThenBy(outboxEntry => outboxEntry.OutboxEntryId) |
| | 40 | 167 | | .Take(request.MaxCount); |
| | | 168 | | |
| | 40 | 169 | | return await ClaimEntriesAsync( |
| | 40 | 170 | | isEligible, |
| | 40 | 171 | | candidates, |
| | 40 | 172 | | query => query |
| | 40 | 173 | | .OrderBy(outboxEntry => outboxEntry.CreatedUtc) |
| | 40 | 174 | | .ThenBy(outboxEntry => outboxEntry.OutboxEntryId), |
| | 40 | 175 | | request, |
| | 40 | 176 | | cancellationToken) |
| | 40 | 177 | | .ConfigureAwait(false); |
| | 40 | 178 | | } |
| | | 179 | | |
| | | 180 | | /// <inheritdoc /> |
| | | 181 | | public async ValueTask<IReadOnlyList<GovernanceOutboxClaim>> ClaimRetryReadyAsync( |
| | | 182 | | GovernanceOutboxClaimRequest request, |
| | | 183 | | CancellationToken cancellationToken = default) |
| | | 184 | | { |
| | 6 | 185 | | ArgumentNullException.ThrowIfNull(request); |
| | 5 | 186 | | cancellationToken.ThrowIfCancellationRequested(); |
| | | 187 | | |
| | 4 | 188 | | Expression<Func<GovernanceOutboxEntryEntity, bool>> isEligible = outboxEntry => |
| | 4 | 189 | | (outboxEntry.Status == GovernanceEmissionStatus.Deferred || |
| | 4 | 190 | | outboxEntry.Status == GovernanceEmissionStatus.Failed || |
| | 4 | 191 | | outboxEntry.Status == GovernanceEmissionStatus.RetryableFailure) |
| | 4 | 192 | | && (outboxEntry.NextRetryUtc == null || outboxEntry.NextRetryUtc <= request.UtcNow) |
| | 4 | 193 | | && (outboxEntry.ClaimToken == null || outboxEntry.ClaimExpiresUtc == null || outboxEntry.ClaimExpiresUtc <= |
| | | 194 | | |
| | 4 | 195 | | IQueryable<GovernanceOutboxEntryEntity> candidates = OutboxEntries() |
| | 4 | 196 | | .Where(isEligible) |
| | 4 | 197 | | .OrderBy(outboxEntry => outboxEntry.NextRetryUtc ?? outboxEntry.UpdatedUtc) |
| | 4 | 198 | | .ThenBy(outboxEntry => outboxEntry.OutboxEntryId) |
| | 4 | 199 | | .Take(request.MaxCount); |
| | | 200 | | |
| | 4 | 201 | | return await ClaimEntriesAsync( |
| | 4 | 202 | | isEligible, |
| | 4 | 203 | | candidates, |
| | 4 | 204 | | query => query |
| | 4 | 205 | | .OrderBy(outboxEntry => outboxEntry.NextRetryUtc ?? outboxEntry.UpdatedUtc) |
| | 4 | 206 | | .ThenBy(outboxEntry => outboxEntry.OutboxEntryId), |
| | 4 | 207 | | request, |
| | 4 | 208 | | cancellationToken) |
| | 4 | 209 | | .ConfigureAwait(false); |
| | 4 | 210 | | } |
| | | 211 | | |
| | | 212 | | /// <inheritdoc /> |
| | | 213 | | public async ValueTask<GovernanceOutboxEntry> MarkDeliveredAsync( |
| | | 214 | | string outboxEntryId, |
| | | 215 | | GovernanceEmissionResult result, |
| | | 216 | | CancellationToken cancellationToken = default) |
| | | 217 | | { |
| | 12 | 218 | | ArgumentException.ThrowIfNullOrWhiteSpace(outboxEntryId); |
| | 12 | 219 | | ArgumentNullException.ThrowIfNull(result); |
| | | 220 | | |
| | 12 | 221 | | GovernanceOutboxEntry entry = await RequireEntryAsync(outboxEntryId, cancellationToken).ConfigureAwait(false); |
| | 12 | 222 | | if (IsTerminal(entry)) |
| | | 223 | | { |
| | 3 | 224 | | return entry; |
| | | 225 | | } |
| | | 226 | | |
| | 9 | 227 | | GovernanceOutboxEntry updatedEntry = entry.MarkDelivered(result); |
| | | 228 | | |
| | 9 | 229 | | return await SaveAsync(updatedEntry, cancellationToken).ConfigureAwait(false); |
| | 11 | 230 | | } |
| | | 231 | | |
| | | 232 | | /// <inheritdoc /> |
| | | 233 | | public async ValueTask<GovernanceOutboxEntry> MarkClaimDeliveredAsync( |
| | | 234 | | GovernanceOutboxClaim claim, |
| | | 235 | | GovernanceEmissionResult result, |
| | | 236 | | CancellationToken cancellationToken = default) |
| | | 237 | | { |
| | 16 | 238 | | ArgumentNullException.ThrowIfNull(claim); |
| | 15 | 239 | | ArgumentNullException.ThrowIfNull(result); |
| | | 240 | | |
| | 14 | 241 | | return await UpdateClaimedEntryAsync(claim, entry => entry.MarkDelivered(result), cancellationToken).ConfigureAw |
| | 13 | 242 | | } |
| | | 243 | | |
| | | 244 | | /// <inheritdoc /> |
| | | 245 | | public async ValueTask<GovernanceOutboxEntry> MarkFailedAsync( |
| | | 246 | | string outboxEntryId, |
| | | 247 | | GovernanceEmissionError governanceEmissionError, |
| | | 248 | | DateTimeOffset? nextRetryUtc = null, |
| | | 249 | | CancellationToken cancellationToken = default) |
| | | 250 | | { |
| | 8 | 251 | | ArgumentException.ThrowIfNullOrWhiteSpace(outboxEntryId); |
| | 8 | 252 | | ArgumentNullException.ThrowIfNull(governanceEmissionError); |
| | | 253 | | |
| | 8 | 254 | | GovernanceOutboxEntry entry = await RequireEntryAsync(outboxEntryId, cancellationToken).ConfigureAwait(false); |
| | 8 | 255 | | if (IsTerminal(entry)) |
| | | 256 | | { |
| | 2 | 257 | | return entry; |
| | | 258 | | } |
| | | 259 | | |
| | 6 | 260 | | GovernanceOutboxEntry updatedEntry = entry.MarkFailed(governanceEmissionError, nextRetryUtc); |
| | | 261 | | |
| | 6 | 262 | | return await SaveAsync(updatedEntry, cancellationToken).ConfigureAwait(false); |
| | 7 | 263 | | } |
| | | 264 | | |
| | | 265 | | /// <inheritdoc /> |
| | | 266 | | public async ValueTask<GovernanceOutboxEntry> MarkClaimFailedAsync( |
| | | 267 | | GovernanceOutboxClaim claim, |
| | | 268 | | GovernanceEmissionError governanceEmissionError, |
| | | 269 | | DateTimeOffset? nextRetryUtc = null, |
| | | 270 | | CancellationToken cancellationToken = default) |
| | | 271 | | { |
| | 13 | 272 | | ArgumentNullException.ThrowIfNull(claim); |
| | 12 | 273 | | ArgumentNullException.ThrowIfNull(governanceEmissionError); |
| | | 274 | | |
| | 11 | 275 | | return await UpdateClaimedEntryAsync( |
| | 11 | 276 | | claim, |
| | 11 | 277 | | entry => entry.MarkFailed(governanceEmissionError, nextRetryUtc), |
| | 11 | 278 | | cancellationToken) |
| | 11 | 279 | | .ConfigureAwait(false); |
| | 11 | 280 | | } |
| | | 281 | | |
| | | 282 | | /// <inheritdoc /> |
| | | 283 | | public async ValueTask<GovernanceOutboxEntry> MarkDeadLetteredAsync( |
| | | 284 | | string outboxEntryId, |
| | | 285 | | GovernanceEmissionError governanceEmissionError, |
| | | 286 | | string? deadLetterReason = null, |
| | | 287 | | CancellationToken cancellationToken = default) |
| | | 288 | | { |
| | 9 | 289 | | ArgumentException.ThrowIfNullOrWhiteSpace(outboxEntryId); |
| | 9 | 290 | | ArgumentNullException.ThrowIfNull(governanceEmissionError); |
| | | 291 | | |
| | 9 | 292 | | GovernanceOutboxEntry entry = await RequireEntryAsync(outboxEntryId, cancellationToken).ConfigureAwait(false); |
| | 9 | 293 | | if (IsTerminal(entry)) |
| | | 294 | | { |
| | 2 | 295 | | return entry; |
| | | 296 | | } |
| | | 297 | | |
| | 7 | 298 | | GovernanceOutboxEntry updatedEntry = entry.MarkDeadLettered(governanceEmissionError, deadLetterReason); |
| | | 299 | | |
| | 7 | 300 | | return await SaveAsync(updatedEntry, cancellationToken).ConfigureAwait(false); |
| | 9 | 301 | | } |
| | | 302 | | |
| | | 303 | | /// <inheritdoc /> |
| | | 304 | | public async ValueTask<GovernanceOutboxEntry> MarkClaimDeadLetteredAsync( |
| | | 305 | | GovernanceOutboxClaim claim, |
| | | 306 | | GovernanceEmissionError governanceEmissionError, |
| | | 307 | | string? deadLetterReason = null, |
| | | 308 | | CancellationToken cancellationToken = default) |
| | | 309 | | { |
| | 5 | 310 | | ArgumentNullException.ThrowIfNull(claim); |
| | 4 | 311 | | ArgumentNullException.ThrowIfNull(governanceEmissionError); |
| | | 312 | | |
| | 3 | 313 | | return await UpdateClaimedEntryAsync( |
| | 3 | 314 | | claim, |
| | 3 | 315 | | entry => entry.MarkDeadLettered(governanceEmissionError, deadLetterReason), |
| | 3 | 316 | | cancellationToken) |
| | 3 | 317 | | .ConfigureAwait(false); |
| | 3 | 318 | | } |
| | | 319 | | |
| | | 320 | | /// <inheritdoc /> |
| | | 321 | | public async ValueTask<GovernanceOutboxEntry> SaveClaimAsync( |
| | | 322 | | GovernanceOutboxClaim claim, |
| | | 323 | | GovernanceOutboxEntry entry, |
| | | 324 | | CancellationToken cancellationToken = default) |
| | | 325 | | { |
| | 5 | 326 | | ArgumentNullException.ThrowIfNull(claim); |
| | 4 | 327 | | ArgumentNullException.ThrowIfNull(entry); |
| | | 328 | | |
| | 3 | 329 | | return !string.Equals(claim.OutboxEntryId, entry.OutboxEntryId, StringComparison.Ordinal) |
| | 3 | 330 | | ? throw new ArgumentException("Claim and entry must reference the same outbox entry ID.", nameof(entry)) |
| | 3 | 331 | | : await UpdateClaimedEntryAsync(claim, _ => entry, cancellationToken).ConfigureAwait(false); |
| | 1 | 332 | | } |
| | | 333 | | |
| | | 334 | | /// <inheritdoc /> |
| | | 335 | | public async ValueTask<GovernanceOutboxEntry?> ReleaseClaimAsync( |
| | | 336 | | GovernanceOutboxClaim claim, |
| | | 337 | | string? reason = null, |
| | | 338 | | CancellationToken cancellationToken = default) |
| | | 339 | | { |
| | 3 | 340 | | ArgumentNullException.ThrowIfNull(claim); |
| | 3 | 341 | | cancellationToken.ThrowIfCancellationRequested(); |
| | | 342 | | |
| | 3 | 343 | | GovernanceOutboxEntryEntity? entity = await dbContext |
| | 3 | 344 | | .Set<GovernanceOutboxEntryEntity>() |
| | 3 | 345 | | .SingleOrDefaultAsync(outboxEntry => outboxEntry.OutboxEntryId == claim.OutboxEntryId, cancellationToken) |
| | 3 | 346 | | .ConfigureAwait(false); |
| | | 347 | | |
| | 3 | 348 | | if (entity is null) |
| | | 349 | | { |
| | 0 | 350 | | return null; |
| | | 351 | | } |
| | | 352 | | |
| | 3 | 353 | | GovernanceOutboxEntry currentEntry = ToEntry(entity); |
| | 3 | 354 | | if (!currentEntry.IsClaimedBy(claim) || IsTerminal(currentEntry)) |
| | | 355 | | { |
| | 0 | 356 | | return currentEntry; |
| | | 357 | | } |
| | | 358 | | |
| | 3 | 359 | | GovernanceOutboxEntry releasedEntry = currentEntry.ReleaseClaim(); |
| | 3 | 360 | | await ApplyEntryUpdateAsync(entity, releasedEntry, cancellationToken).ConfigureAwait(false); |
| | | 361 | | |
| | 3 | 362 | | return releasedEntry; |
| | 3 | 363 | | } |
| | | 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 | | { |
| | 44 | 399 | | string claimToken = Guid.NewGuid().ToString("N"); |
| | 44 | 400 | | string concurrencyStamp = GovernanceEntity.NewConcurrencyStamp(); |
| | 44 | 401 | | IQueryable<Guid> candidateIds = candidates.Select(outboxEntry => outboxEntry.Id); |
| | | 402 | | |
| | 44 | 403 | | int claimedCount = await OutboxEntries() |
| | 44 | 404 | | .Where(isEligible) |
| | 44 | 405 | | .Where(outboxEntry => candidateIds.Contains(outboxEntry.Id)) |
| | 44 | 406 | | .ExecuteUpdateAsync( |
| | 44 | 407 | | setters => setters |
| | 44 | 408 | | .SetProperty(outboxEntry => outboxEntry.ConcurrencyStamp, concurrencyStamp) |
| | 44 | 409 | | .SetProperty(outboxEntry => outboxEntry.UpdatedUtc, request.UtcNow) |
| | 44 | 410 | | .SetProperty(outboxEntry => outboxEntry.ClaimOwner, request.WorkerId) |
| | 44 | 411 | | .SetProperty(outboxEntry => outboxEntry.ClaimToken, claimToken) |
| | 44 | 412 | | .SetProperty(outboxEntry => outboxEntry.ClaimedUtc, request.UtcNow) |
| | 44 | 413 | | .SetProperty(outboxEntry => outboxEntry.ClaimExpiresUtc, request.ClaimExpiresUtc) |
| | 44 | 414 | | .SetProperty(outboxEntry => outboxEntry.ClaimAttemptCount, outboxEntry => outboxEntry.ClaimAttemptCo |
| | 44 | 415 | | cancellationToken) |
| | 44 | 416 | | .ConfigureAwait(false); |
| | | 417 | | |
| | 44 | 418 | | if (claimedCount == 0) |
| | | 419 | | { |
| | 3 | 420 | | return []; |
| | | 421 | | } |
| | | 422 | | |
| | 41 | 423 | | List<GovernanceOutboxEntryEntity> claimedEntities = await orderClaimedEntries( |
| | 41 | 424 | | OutboxEntries().Where(outboxEntry => outboxEntry.ClaimToken == claimToken)) |
| | 41 | 425 | | .ToListAsync(cancellationToken) |
| | 41 | 426 | | .ConfigureAwait(false); |
| | | 427 | | |
| | 41 | 428 | | ReconcileTrackedEntries(claimedEntities); |
| | | 429 | | |
| | 41 | 430 | | return [.. claimedEntities.Select(ToEntry).Select(CreateClaim)]; |
| | 44 | 431 | | } |
| | | 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 | | { |
| | 41 | 454 | | if (claimedEntities.Count == 0) |
| | | 455 | | { |
| | 0 | 456 | | return; |
| | | 457 | | } |
| | | 458 | | |
| | 41 | 459 | | var claimedById = claimedEntities.ToDictionary(entity => entity.OutboxEntryId, StringComparer.Ordinal); |
| | | 460 | | |
| | 41 | 461 | | List<Microsoft.EntityFrameworkCore.ChangeTracking.EntityEntry<GovernanceOutboxEntryEntity>> staleEntries = [.. d |
| | 41 | 462 | | .ChangeTracker |
| | 41 | 463 | | .Entries<GovernanceOutboxEntryEntity>() |
| | 41 | 464 | | .Where(entry => claimedById.ContainsKey(entry.Entity.OutboxEntryId))]; |
| | | 465 | | |
| | 96 | 466 | | foreach (Microsoft.EntityFrameworkCore.ChangeTracking.EntityEntry<GovernanceOutboxEntryEntity> entry in staleEnt |
| | | 467 | | { |
| | 7 | 468 | | if (entry.State == EntityState.Unchanged) |
| | | 469 | | { |
| | 3 | 470 | | entry.State = EntityState.Detached; |
| | 3 | 471 | | continue; |
| | | 472 | | } |
| | | 473 | | |
| | 4 | 474 | | if (entry.State is EntityState.Modified or EntityState.Deleted) |
| | | 475 | | { |
| | 4 | 476 | | MergeClaimColumns(entry, claimedById[entry.Entity.OutboxEntryId]); |
| | | 477 | | } |
| | | 478 | | } |
| | 41 | 479 | | } |
| | | 480 | | |
| | | 481 | | private static void MergeClaimColumns( |
| | | 482 | | Microsoft.EntityFrameworkCore.ChangeTracking.EntityEntry<GovernanceOutboxEntryEntity> entry, |
| | | 483 | | GovernanceOutboxEntryEntity claimed) |
| | | 484 | | { |
| | 72 | 485 | | foreach ((string propertyName, object? claimedValue) in ClaimColumnValues(claimed)) |
| | | 486 | | { |
| | 32 | 487 | | Microsoft.EntityFrameworkCore.ChangeTracking.PropertyEntry property = entry.Property(propertyName); |
| | 32 | 488 | | property.OriginalValue = claimedValue; |
| | 32 | 489 | | property.CurrentValue = claimedValue; |
| | | 490 | | } |
| | 4 | 491 | | } |
| | | 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 | | { |
| | 4 | 497 | | return |
| | 4 | 498 | | [ |
| | 4 | 499 | | (nameof(GovernanceOutboxEntryEntity.ConcurrencyStamp), claimed.ConcurrencyStamp), |
| | 4 | 500 | | (nameof(GovernanceOutboxEntryEntity.Status), claimed.Status), |
| | 4 | 501 | | (nameof(GovernanceOutboxEntryEntity.UpdatedUtc), claimed.UpdatedUtc), |
| | 4 | 502 | | (nameof(GovernanceOutboxEntryEntity.ClaimOwner), claimed.ClaimOwner), |
| | 4 | 503 | | (nameof(GovernanceOutboxEntryEntity.ClaimToken), claimed.ClaimToken), |
| | 4 | 504 | | (nameof(GovernanceOutboxEntryEntity.ClaimedUtc), claimed.ClaimedUtc), |
| | 4 | 505 | | (nameof(GovernanceOutboxEntryEntity.ClaimExpiresUtc), claimed.ClaimExpiresUtc), |
| | 4 | 506 | | (nameof(GovernanceOutboxEntryEntity.ClaimAttemptCount), claimed.ClaimAttemptCount), |
| | 4 | 507 | | ]; |
| | | 508 | | } |
| | | 509 | | |
| | | 510 | | private async ValueTask<GovernanceOutboxEntry> UpdateClaimedEntryAsync( |
| | | 511 | | GovernanceOutboxClaim claim, |
| | | 512 | | Func<GovernanceOutboxEntry, GovernanceOutboxEntry> updateEntry, |
| | | 513 | | CancellationToken cancellationToken) |
| | | 514 | | { |
| | 30 | 515 | | GovernanceOutboxEntryEntity entity = await RequireEntityAsync(claim.OutboxEntryId, cancellationToken).ConfigureA |
| | 28 | 516 | | GovernanceOutboxEntry currentEntry = ToEntry(entity); |
| | | 517 | | |
| | 28 | 518 | | if (!currentEntry.IsClaimedBy(claim) || IsTerminal(currentEntry)) |
| | | 519 | | { |
| | 6 | 520 | | return currentEntry; |
| | | 521 | | } |
| | | 522 | | |
| | 22 | 523 | | GovernanceOutboxEntry updatedEntry = updateEntry(currentEntry); |
| | | 524 | | |
| | | 525 | | try |
| | | 526 | | { |
| | 22 | 527 | | await ApplyEntryUpdateAsync(entity, updatedEntry, cancellationToken).ConfigureAwait(false); |
| | 17 | 528 | | return updatedEntry; |
| | | 529 | | } |
| | | 530 | | catch (DbUpdateConcurrencyException exception) |
| | | 531 | | { |
| | 5 | 532 | | DetachEntries(exception); |
| | 5 | 533 | | GovernanceOutboxEntry? refreshedEntry = await FindByOutboxEntryIdAsync(claim.OutboxEntryId, cancellationToke |
| | 5 | 534 | | return refreshedEntry ?? currentEntry; |
| | | 535 | | } |
| | 28 | 536 | | } |
| | | 537 | | |
| | | 538 | | private async ValueTask ApplyEntryUpdateAsync( |
| | | 539 | | GovernanceOutboxEntryEntity entity, |
| | | 540 | | GovernanceOutboxEntry entry, |
| | | 541 | | CancellationToken cancellationToken) |
| | | 542 | | { |
| | 25 | 543 | | GovernanceOutboxEntryEntity persistedEntity = ToEntity(entry); |
| | 25 | 544 | | persistedEntity.Id = entity.Id; |
| | 25 | 545 | | persistedEntity.ConcurrencyStamp = GovernanceEntity.NewConcurrencyStamp(); |
| | 25 | 546 | | dbContext.Entry(entity).CurrentValues.SetValues(persistedEntity); |
| | | 547 | | |
| | 25 | 548 | | _ = await dbContext.SaveChangesAsync(cancellationToken).ConfigureAwait(false); |
| | 20 | 549 | | } |
| | | 550 | | |
| | | 551 | | private async ValueTask<GovernanceOutboxEntry> RequireEntryAsync( |
| | | 552 | | string outboxEntryId, |
| | | 553 | | CancellationToken cancellationToken) |
| | | 554 | | { |
| | 29 | 555 | | GovernanceOutboxEntry? entry = await FindByOutboxEntryIdAsync(outboxEntryId, cancellationToken).ConfigureAwait(f |
| | | 556 | | |
| | 29 | 557 | | return entry ?? throw new InvalidOperationException($"Outbox entry '{outboxEntryId.Trim()}' was not found."); |
| | 29 | 558 | | } |
| | | 559 | | |
| | | 560 | | private async ValueTask<GovernanceOutboxEntryEntity> RequireEntityAsync( |
| | | 561 | | string outboxEntryId, |
| | | 562 | | CancellationToken cancellationToken) |
| | | 563 | | { |
| | 30 | 564 | | GovernanceOutboxEntryEntity? entity = await dbContext |
| | 30 | 565 | | .Set<GovernanceOutboxEntryEntity>() |
| | 30 | 566 | | .SingleOrDefaultAsync(outboxEntry => outboxEntry.OutboxEntryId == outboxEntryId.Trim(), cancellationToken) |
| | 30 | 567 | | .ConfigureAwait(false); |
| | | 568 | | |
| | 29 | 569 | | return entity ?? throw new InvalidOperationException($"Outbox entry '{outboxEntryId.Trim()}' was not found."); |
| | 28 | 570 | | } |
| | | 571 | | |
| | | 572 | | private IQueryable<GovernanceOutboxEntryEntity> OutboxEntries() |
| | | 573 | | { |
| | 248 | 574 | | return dbContext.Set<GovernanceOutboxEntryEntity>().AsNoTracking(); |
| | | 575 | | } |
| | | 576 | | |
| | | 577 | | private static GovernanceOutboxEntryEntity ToEntity(GovernanceOutboxEntry entry) |
| | | 578 | | { |
| | 304 | 579 | | GovernanceEmissionEnvelope envelope = entry.Envelope; |
| | 304 | 580 | | GovernanceEmissionPayload? payload = envelope.Payload; |
| | 304 | 581 | | GovernanceEmissionError? lastError = entry.LastError; |
| | | 582 | | |
| | 304 | 583 | | return new GovernanceOutboxEntryEntity |
| | 304 | 584 | | { |
| | 304 | 585 | | OutboxEntryId = entry.OutboxEntryId, |
| | 304 | 586 | | Status = entry.Status, |
| | 304 | 587 | | CreatedUtc = entry.CreatedUtc, |
| | 304 | 588 | | UpdatedUtc = entry.UpdatedUtc, |
| | 304 | 589 | | DeliveredUtc = entry.Status is GovernanceEmissionStatus.Delivered ? entry.UpdatedUtc : null, |
| | 304 | 590 | | RetryCount = entry.RetryCount, |
| | 304 | 591 | | MaxRetryCount = entry.MaxRetryCount, |
| | 304 | 592 | | NextRetryUtc = entry.NextRetryUtc, |
| | 304 | 593 | | ProviderName = entry.ProviderName, |
| | 304 | 594 | | ProviderRecordId = entry.ProviderRecordId, |
| | 304 | 595 | | DeadLetterReason = entry.DeadLetterReason, |
| | 304 | 596 | | LastErrorCode = lastError?.Code, |
| | 304 | 597 | | LastErrorMessage = lastError?.Message, |
| | 304 | 598 | | LastErrorIsRetryable = lastError?.IsRetryable, |
| | 304 | 599 | | LastErrorProviderName = lastError?.ProviderName, |
| | 304 | 600 | | LastErrorProviderErrorCode = lastError?.ProviderErrorCode, |
| | 304 | 601 | | MetadataJson = JsonSerializer.Serialize(entry.Metadata, JsonOptions), |
| | 304 | 602 | | ClaimOwner = entry.ClaimOwner, |
| | 304 | 603 | | ClaimToken = entry.ClaimToken, |
| | 304 | 604 | | ClaimedUtc = entry.ClaimedUtc, |
| | 304 | 605 | | ClaimExpiresUtc = entry.ClaimExpiresUtc, |
| | 304 | 606 | | ClaimAttemptCount = entry.ClaimAttemptCount, |
| | 304 | 607 | | EnvelopeId = envelope.EnvelopeId, |
| | 304 | 608 | | EnvelopeSchemaVersion = envelope.SchemaVersion, |
| | 304 | 609 | | EnvelopeEventType = envelope.EventType, |
| | 304 | 610 | | EnvelopeEventId = envelope.EventId, |
| | 304 | 611 | | EnvelopeOccurredUtc = envelope.OccurredUtc, |
| | 304 | 612 | | EnvelopeCreatedUtc = envelope.CreatedUtc, |
| | 304 | 613 | | EnvelopeCorrelationId = envelope.CorrelationId, |
| | 304 | 614 | | EnvelopeDecisionReceiptId = envelope.DecisionReceiptId, |
| | 304 | 615 | | EnvelopeLifecycleStage = envelope.LifecycleStage, |
| | 304 | 616 | | EnvelopeLifecycleStageSequence = envelope.LifecycleStageSequence, |
| | 304 | 617 | | EnvelopePolicyVersion = envelope.PolicyVersion, |
| | 304 | 618 | | EnvelopePolicyHash = envelope.PolicyHash, |
| | 304 | 619 | | EnvelopeTraceId = envelope.TraceId, |
| | 304 | 620 | | EnvelopeSpanId = envelope.SpanId, |
| | 304 | 621 | | EnvelopeParentSpanId = envelope.ParentSpanId, |
| | 304 | 622 | | EnvelopeOperationName = envelope.OperationName, |
| | 304 | 623 | | EnvelopeOutcome = envelope.Outcome, |
| | 304 | 624 | | EnvelopeActorId = envelope.ActorId, |
| | 304 | 625 | | EnvelopeEmitterStatus = envelope.EmitterStatus, |
| | 304 | 626 | | EnvelopeEmitterProvider = envelope.EmitterProvider, |
| | 304 | 627 | | EnvelopeOutboxSequence = envelope.OutboxSequence, |
| | 304 | 628 | | EnvelopeGatewayExecutionId = envelope.GatewayExecutionId, |
| | 304 | 629 | | EnvelopeDecisionStage = envelope.DecisionStage, |
| | 304 | 630 | | EnvelopeMetadataJson = JsonSerializer.Serialize(envelope.Metadata, JsonOptions), |
| | 304 | 631 | | EnvelopePayloadType = payload?.PayloadType, |
| | 304 | 632 | | EnvelopePayloadSchemaVersion = payload?.SchemaVersion, |
| | 304 | 633 | | EnvelopePayloadContentType = payload?.ContentType, |
| | 304 | 634 | | EnvelopePayloadContentHash = payload?.ContentHash, |
| | 304 | 635 | | EnvelopePayloadSizeBytes = payload?.SizeBytes, |
| | 304 | 636 | | EnvelopePayloadMetadataJson = JsonSerializer.Serialize(payload?.Metadata ?? EmptyMetadata(), JsonOptions) |
| | 304 | 637 | | }; |
| | | 638 | | } |
| | | 639 | | |
| | | 640 | | private static GovernanceOutboxEntry[] ToEntries(IEnumerable<GovernanceOutboxEntryEntity> entities) |
| | | 641 | | { |
| | 28 | 642 | | return [.. entities.Select(ToEntry)]; |
| | | 643 | | } |
| | | 644 | | |
| | | 645 | | private static GovernanceOutboxEntry ToEntry(GovernanceOutboxEntryEntity entity) |
| | | 646 | | { |
| | 381 | 647 | | GovernanceEmissionPayload? payload = string.IsNullOrWhiteSpace(entity.EnvelopePayloadType) |
| | 381 | 648 | | ? null |
| | 381 | 649 | | : GovernanceEmissionPayload.Create( |
| | 381 | 650 | | entity.EnvelopePayloadType, |
| | 381 | 651 | | entity.EnvelopePayloadSchemaVersion, |
| | 381 | 652 | | entity.EnvelopePayloadContentType, |
| | 381 | 653 | | entity.EnvelopePayloadContentHash, |
| | 381 | 654 | | entity.EnvelopePayloadSizeBytes, |
| | 381 | 655 | | DeserializeMetadata(entity.EnvelopePayloadMetadataJson)); |
| | | 656 | | |
| | 381 | 657 | | var envelope = GovernanceEmissionEnvelope.Create( |
| | 381 | 658 | | entity.EnvelopeEventType, |
| | 381 | 659 | | entity.EnvelopeEventId, |
| | 381 | 660 | | entity.EnvelopeOccurredUtc, |
| | 381 | 661 | | entity.EnvelopeId, |
| | 381 | 662 | | entity.EnvelopeCreatedUtc, |
| | 381 | 663 | | entity.EnvelopeSchemaVersion, |
| | 381 | 664 | | entity.EnvelopeCorrelationId, |
| | 381 | 665 | | entity.EnvelopeDecisionReceiptId, |
| | 381 | 666 | | entity.EnvelopeLifecycleStage, |
| | 381 | 667 | | entity.EnvelopePolicyVersion, |
| | 381 | 668 | | entity.EnvelopePolicyHash, |
| | 381 | 669 | | entity.EnvelopeTraceId, |
| | 381 | 670 | | entity.EnvelopeSpanId, |
| | 381 | 671 | | entity.EnvelopeParentSpanId, |
| | 381 | 672 | | entity.EnvelopeOperationName, |
| | 381 | 673 | | entity.EnvelopeOutcome, |
| | 381 | 674 | | entity.EnvelopeActorId, |
| | 381 | 675 | | entity.EnvelopeEmitterStatus, |
| | 381 | 676 | | entity.EnvelopeEmitterProvider, |
| | 381 | 677 | | entity.EnvelopeOutboxSequence, |
| | 381 | 678 | | entity.EnvelopeGatewayExecutionId, |
| | 381 | 679 | | entity.EnvelopeDecisionStage, |
| | 381 | 680 | | payload, |
| | 381 | 681 | | DeserializeMetadata(entity.EnvelopeMetadataJson)); |
| | | 682 | | |
| | 381 | 683 | | GovernanceEmissionError? lastError = string.IsNullOrWhiteSpace(entity.LastErrorCode) || string.IsNullOrWhiteSpac |
| | 381 | 684 | | ? null |
| | 381 | 685 | | : GovernanceEmissionError.Create( |
| | 381 | 686 | | entity.LastErrorCode, |
| | 381 | 687 | | entity.LastErrorMessage, |
| | 381 | 688 | | entity.LastErrorIsRetryable ?? false, |
| | 381 | 689 | | entity.LastErrorProviderName, |
| | 381 | 690 | | entity.LastErrorProviderErrorCode); |
| | | 691 | | |
| | 381 | 692 | | return GovernanceOutboxEntry.Restore( |
| | 381 | 693 | | envelope, |
| | 381 | 694 | | entity.Status, |
| | 381 | 695 | | entity.OutboxEntryId, |
| | 381 | 696 | | entity.CreatedUtc, |
| | 381 | 697 | | entity.UpdatedUtc, |
| | 381 | 698 | | entity.RetryCount, |
| | 381 | 699 | | entity.MaxRetryCount, |
| | 381 | 700 | | entity.NextRetryUtc, |
| | 381 | 701 | | lastError, |
| | 381 | 702 | | entity.ProviderName, |
| | 381 | 703 | | entity.ProviderRecordId, |
| | 381 | 704 | | entity.DeadLetterReason, |
| | 381 | 705 | | DeserializeMetadata(entity.MetadataJson), |
| | 381 | 706 | | entity.ClaimOwner, |
| | 381 | 707 | | entity.ClaimToken, |
| | 381 | 708 | | entity.ClaimedUtc, |
| | 381 | 709 | | entity.ClaimExpiresUtc, |
| | 381 | 710 | | entity.ClaimAttemptCount); |
| | | 711 | | } |
| | | 712 | | |
| | | 713 | | private static GovernanceOutboxClaim CreateClaim(GovernanceOutboxEntry entry) |
| | | 714 | | { |
| | 208 | 715 | | return GovernanceOutboxClaim.Create( |
| | 208 | 716 | | entry, |
| | 208 | 717 | | entry.ClaimOwner ?? throw new InvalidOperationException("Claimed entry is missing claim owner."), |
| | 208 | 718 | | entry.ClaimToken ?? throw new InvalidOperationException("Claimed entry is missing claim token."), |
| | 208 | 719 | | entry.ClaimedUtc ?? throw new InvalidOperationException("Claimed entry is missing claimed timestamp."), |
| | 208 | 720 | | entry.ClaimExpiresUtc ?? throw new InvalidOperationException("Claimed entry is missing claim expiration time |
| | | 721 | | } |
| | | 722 | | |
| | | 723 | | internal static bool IsPendingClaimEligible(GovernanceOutboxEntryEntity entity, DateTimeOffset utcNow) |
| | | 724 | | { |
| | 0 | 725 | | return entity.Status is GovernanceEmissionStatus.Pending && IsClaimAvailable(entity, utcNow); |
| | | 726 | | } |
| | | 727 | | |
| | | 728 | | internal static bool IsRetryReadyClaimEligible(GovernanceOutboxEntryEntity entity, DateTimeOffset utcNow) |
| | | 729 | | { |
| | 20 | 730 | | return (entity.Status is GovernanceEmissionStatus.Deferred or GovernanceEmissionStatus.Failed or GovernanceEmiss |
| | 20 | 731 | | && (entity.NextRetryUtc is null || entity.NextRetryUtc <= utcNow.ToUniversalTime()) |
| | 20 | 732 | | && IsClaimAvailable(entity, utcNow); |
| | | 733 | | } |
| | | 734 | | |
| | | 735 | | private static bool IsClaimAvailable(GovernanceOutboxEntryEntity entity, DateTimeOffset utcNow) |
| | | 736 | | { |
| | 16 | 737 | | return entity.ClaimToken is null || entity.ClaimExpiresUtc is null || entity.ClaimExpiresUtc <= utcNow.ToUnivers |
| | | 738 | | } |
| | | 739 | | |
| | | 740 | | private static bool IsTerminal(GovernanceOutboxEntry entry) |
| | | 741 | | { |
| | 81 | 742 | | return entry.IsDelivered || entry.IsDeadLettered; |
| | | 743 | | } |
| | | 744 | | |
| | | 745 | | private static void DetachEntries(DbUpdateConcurrencyException exception) |
| | | 746 | | { |
| | 20 | 747 | | foreach (Microsoft.EntityFrameworkCore.ChangeTracking.EntityEntry entry in exception.Entries) |
| | | 748 | | { |
| | 5 | 749 | | entry.State = EntityState.Detached; |
| | | 750 | | } |
| | 5 | 751 | | } |
| | | 752 | | |
| | | 753 | | private static ReadOnlyDictionary<string, string> DeserializeMetadata(string? json) |
| | | 754 | | { |
| | 788 | 755 | | if (string.IsNullOrWhiteSpace(json)) |
| | | 756 | | { |
| | 1 | 757 | | return EmptyMetadata(); |
| | | 758 | | } |
| | | 759 | | |
| | 787 | 760 | | Dictionary<string, string>? metadata = JsonSerializer.Deserialize<Dictionary<string, string>>(json, JsonOptions) |
| | | 761 | | |
| | 787 | 762 | | return metadata is null || metadata.Count == 0 |
| | 787 | 763 | | ? EmptyMetadata() |
| | 787 | 764 | | : new ReadOnlyDictionary<string, string>(new Dictionary<string, string>(metadata, StringComparer.Ordinal)); |
| | | 765 | | } |
| | | 766 | | |
| | | 767 | | private static ReadOnlyDictionary<string, string> EmptyMetadata() |
| | | 768 | | { |
| | 895 | 769 | | return new ReadOnlyDictionary<string, string>(new Dictionary<string, string>(StringComparer.Ordinal)); |
| | | 770 | | } |
| | | 771 | | |
| | | 772 | | private static int NormalizeMaxCount(int maxCount) |
| | | 773 | | { |
| | 30 | 774 | | return maxCount <= 0 |
| | 30 | 775 | | ? throw new ArgumentOutOfRangeException(nameof(maxCount), maxCount, "Maximum count must be greater than zero |
| | 30 | 776 | | : maxCount; |
| | | 777 | | } |
| | | 778 | | } |