| | | 1 | | using System.Runtime.CompilerServices; |
| | | 2 | | using System.Security.Cryptography; |
| | | 3 | | using System.Text; |
| | | 4 | | using Microsoft.EntityFrameworkCore; |
| | | 5 | | using Microsoft.EntityFrameworkCore.Infrastructure; |
| | | 6 | | using Microsoft.EntityFrameworkCore.Metadata; |
| | | 7 | | using Microsoft.EntityFrameworkCore.Storage; |
| | | 8 | | using Microsoft.Extensions.Options; |
| | | 9 | | using ProjectTemplate.Infrastructure.Data.Entities; |
| | | 10 | | |
| | | 11 | | namespace ProjectTemplate.Infrastructure.Data.Auditing; |
| | | 12 | | |
| | 19 | 13 | | public sealed class ApplicationAuditReconciler( |
| | 19 | 14 | | ApplicationDbContext dbContext, |
| | 19 | 15 | | IApplicationMutationManifestVerifier manifestVerifier, |
| | 19 | 16 | | IOptions<ApplicationAuditReconciliationOptions> options, |
| | 19 | 17 | | ApplicationAuditReconciliationMetrics metrics, |
| | 19 | 18 | | TimeProvider timeProvider) |
| | | 19 | | : IApplicationAuditReconciler |
| | | 20 | | { |
| | 19 | 21 | | private readonly ApplicationDbContext _dbContext = |
| | 19 | 22 | | dbContext ?? throw new ArgumentNullException(nameof(dbContext)); |
| | 19 | 23 | | private readonly IApplicationMutationManifestVerifier _manifestVerifier = |
| | 19 | 24 | | manifestVerifier ?? throw new ArgumentNullException(nameof(manifestVerifier)); |
| | 19 | 25 | | private readonly ApplicationAuditReconciliationOptions _options = |
| | 19 | 26 | | options?.Value ?? throw new ArgumentNullException(nameof(options)); |
| | 19 | 27 | | private readonly ApplicationAuditReconciliationMetrics _metrics = |
| | 19 | 28 | | metrics ?? throw new ArgumentNullException(nameof(metrics)); |
| | 19 | 29 | | private readonly TimeProvider _timeProvider = |
| | 19 | 30 | | timeProvider ?? throw new ArgumentNullException(nameof(timeProvider)); |
| | | 31 | | |
| | | 32 | | public async Task<ApplicationAuditReconciliationSummary> ReconcileAsync( |
| | | 33 | | CancellationToken cancellationToken = default) |
| | | 34 | | { |
| | 20 | 35 | | cancellationToken.ThrowIfCancellationRequested(); |
| | 20 | 36 | | if (!_options.Enabled) |
| | | 37 | | { |
| | 1 | 38 | | return ApplicationAuditReconciliationMetrics.DisabledSummary; |
| | | 39 | | } |
| | | 40 | | |
| | 19 | 41 | | DateTime now = UtcNow(); |
| | 19 | 42 | | List<string> batchIds = await _dbContext.AuditRecords |
| | 19 | 43 | | .AsNoTracking() |
| | 19 | 44 | | .Where(record => record.MutationBatchId != string.Empty) |
| | 19 | 45 | | .GroupBy(record => record.MutationBatchId) |
| | 19 | 46 | | .OrderByDescending(group => group.Max(record => record.ModifiedOnUtc)) |
| | 19 | 47 | | .Select(group => group.Key) |
| | 19 | 48 | | .Take(_options.MaximumBatchesPerRun) |
| | 19 | 49 | | .ToListAsync(cancellationToken) |
| | 19 | 50 | | .ConfigureAwait(false); |
| | | 51 | | |
| | 19 | 52 | | List<AuditRecord> auditRecords = await _dbContext.AuditRecords |
| | 19 | 53 | | .AsNoTracking() |
| | 19 | 54 | | .Where(record => batchIds.Contains(record.MutationBatchId)) |
| | 19 | 55 | | .ToListAsync(cancellationToken) |
| | 19 | 56 | | .ConfigureAwait(false); |
| | | 57 | | |
| | | 58 | | // Records without a batch id are not bounded by the batch selection above, so cap them separately; |
| | | 59 | | // otherwise every such record ever written would be loaded on every reconciliation pass. |
| | 19 | 60 | | List<AuditRecord> malformedRecords = await _dbContext.AuditRecords |
| | 19 | 61 | | .AsNoTracking() |
| | 19 | 62 | | .Where(record => record.MutationBatchId == string.Empty) |
| | 19 | 63 | | .OrderByDescending(record => record.ModifiedOnUtc) |
| | 19 | 64 | | .ThenByDescending(record => record.Id) |
| | 19 | 65 | | .Take(_options.MaximumMalformedRecordsPerRun) |
| | 19 | 66 | | .ToListAsync(cancellationToken) |
| | 19 | 67 | | .ConfigureAwait(false); |
| | 19 | 68 | | auditRecords.AddRange(malformedRecords); |
| | | 69 | | |
| | 19 | 70 | | List<ApplicationAuditCompletionOutboxEntry> completionEntries = await _dbContext |
| | 19 | 71 | | .ApplicationAuditCompletionOutboxEntries |
| | 19 | 72 | | .AsNoTracking() |
| | 19 | 73 | | .Where(entry => batchIds.Contains(entry.MutationBatchId) || |
| | 19 | 74 | | entry.Status != ApplicationAuditCompletionOutboxStatuses.Delivered) |
| | 19 | 75 | | .OrderByDescending(entry => entry.CreatedUtc) |
| | 19 | 76 | | .Take(_options.MaximumBatchesPerRun * 2) |
| | 19 | 77 | | .ToListAsync(cancellationToken) |
| | 19 | 78 | | .ConfigureAwait(false); |
| | | 79 | | |
| | 19 | 80 | | List<ApplicationAuditReconciliationCandidate> candidates = BuildCandidates( |
| | 19 | 81 | | auditRecords, |
| | 19 | 82 | | completionEntries, |
| | 19 | 83 | | now); |
| | | 84 | | |
| | 19 | 85 | | await PersistCandidatesAsync(candidates, batchIds, completionEntries, now, cancellationToken) |
| | 19 | 86 | | .ConfigureAwait(false); |
| | | 87 | | |
| | 17 | 88 | | ApplicationAuditReconciliationSummary summary = await GetSummaryCoreAsync(now, cancellationToken) |
| | 17 | 89 | | .ConfigureAwait(false); |
| | 17 | 90 | | _metrics.Update(summary); |
| | 17 | 91 | | return summary; |
| | 18 | 92 | | } |
| | | 93 | | |
| | | 94 | | public async Task<ApplicationAuditReconciliationSummary> GetSummaryAsync( |
| | | 95 | | CancellationToken cancellationToken = default) |
| | | 96 | | { |
| | 2 | 97 | | cancellationToken.ThrowIfCancellationRequested(); |
| | 2 | 98 | | if (!_options.Enabled) |
| | | 99 | | { |
| | 1 | 100 | | return ApplicationAuditReconciliationMetrics.DisabledSummary; |
| | | 101 | | } |
| | | 102 | | |
| | 1 | 103 | | ApplicationAuditReconciliationSummary summary = await GetSummaryCoreAsync( |
| | 1 | 104 | | _metrics.LastRunUtc, |
| | 1 | 105 | | cancellationToken) |
| | 1 | 106 | | .ConfigureAwait(false); |
| | 1 | 107 | | _metrics.Update(summary); |
| | 1 | 108 | | return summary; |
| | 2 | 109 | | } |
| | | 110 | | |
| | | 111 | | public async Task<IReadOnlyList<ApplicationAuditReconciliationFindingItem>> QueryFindingsAsync( |
| | | 112 | | ApplicationAuditReconciliationQuery request, |
| | | 113 | | CancellationToken cancellationToken = default) |
| | | 114 | | { |
| | 1 | 115 | | ArgumentNullException.ThrowIfNull(request); |
| | 1 | 116 | | cancellationToken.ThrowIfCancellationRequested(); |
| | | 117 | | |
| | 1 | 118 | | IQueryable<ApplicationAuditReconciliationFinding> query = _dbContext |
| | 1 | 119 | | .ApplicationAuditReconciliationFindings |
| | 1 | 120 | | .AsNoTracking(); |
| | | 121 | | |
| | 1 | 122 | | if (!string.IsNullOrWhiteSpace(request.ReasonCode)) |
| | | 123 | | { |
| | 1 | 124 | | string reasonCode = request.ReasonCode.Trim(); |
| | 1 | 125 | | query = query.Where(finding => finding.ReasonCode == reasonCode); |
| | | 126 | | } |
| | | 127 | | |
| | 1 | 128 | | if (!string.IsNullOrWhiteSpace(request.Severity)) |
| | | 129 | | { |
| | 1 | 130 | | string severity = request.Severity.Trim(); |
| | 1 | 131 | | query = query.Where(finding => finding.Severity == severity); |
| | | 132 | | } |
| | | 133 | | |
| | 1 | 134 | | if (!string.IsNullOrWhiteSpace(request.MutationBatchId)) |
| | | 135 | | { |
| | 1 | 136 | | string batchId = request.MutationBatchId.Trim(); |
| | 1 | 137 | | query = query.Where(finding => finding.MutationBatchId == batchId); |
| | | 138 | | } |
| | | 139 | | |
| | 1 | 140 | | if (!string.IsNullOrWhiteSpace(request.RemediationStatus)) |
| | | 141 | | { |
| | 1 | 142 | | string status = request.RemediationStatus.Trim(); |
| | 1 | 143 | | query = query.Where(finding => finding.RemediationStatus == status); |
| | | 144 | | } |
| | | 145 | | |
| | 1 | 146 | | int maximumResults = Math.Clamp(request.MaximumResults, 1, 500); |
| | 1 | 147 | | return await query |
| | 1 | 148 | | .OrderByDescending(finding => finding.LastObservedUtc) |
| | 1 | 149 | | .Take(maximumResults) |
| | 1 | 150 | | .Select(finding => new ApplicationAuditReconciliationFindingItem( |
| | 1 | 151 | | finding.Id, |
| | 1 | 152 | | finding.SchemaVersion, |
| | 1 | 153 | | finding.FindingKey, |
| | 1 | 154 | | finding.ReasonCode, |
| | 1 | 155 | | finding.Severity, |
| | 1 | 156 | | finding.MutationBatchId, |
| | 1 | 157 | | finding.Destination, |
| | 1 | 158 | | finding.Guidance, |
| | 1 | 159 | | finding.RemediationStatus, |
| | 1 | 160 | | finding.FirstObservedUtc, |
| | 1 | 161 | | finding.LastObservedUtc, |
| | 1 | 162 | | finding.ResolvedUtc)) |
| | 1 | 163 | | .ToListAsync(cancellationToken) |
| | 1 | 164 | | .ConfigureAwait(false); |
| | 1 | 165 | | } |
| | | 166 | | |
| | | 167 | | public async Task<ApplicationAuditReconciliationRemediationItem> RecordRemediationAsync( |
| | | 168 | | Guid findingId, |
| | | 169 | | ApplicationAuditReconciliationRemediationRequest request, |
| | | 170 | | CancellationToken cancellationToken = default) |
| | | 171 | | { |
| | 5 | 172 | | ArgumentNullException.ThrowIfNull(request); |
| | 5 | 173 | | cancellationToken.ThrowIfCancellationRequested(); |
| | | 174 | | |
| | 5 | 175 | | string actionCode = NormalizeRequired(request.ActionCode, 64, nameof(request.ActionCode)); |
| | 4 | 176 | | string actorId = NormalizeRequired(request.ActorId, 256, nameof(request.ActorId)); |
| | 4 | 177 | | string? evidenceReference = NormalizeOptional(request.EvidenceReference, 256); |
| | | 178 | | |
| | | 179 | | // When the caller already owns an EF Core transaction, join it. The caller decides whether the |
| | | 180 | | // remediation commits, and the execution strategy is bypassed because retrying strategies cannot |
| | | 181 | | // replay work inside a caller-owned transaction. |
| | 4 | 182 | | if (_dbContext.Database.CurrentTransaction is not null) |
| | | 183 | | { |
| | 1 | 184 | | return await RecordRemediationCoreAsync( |
| | 1 | 185 | | findingId, |
| | 1 | 186 | | actionCode, |
| | 1 | 187 | | actorId, |
| | 1 | 188 | | evidenceReference, |
| | 1 | 189 | | request.ResolveFinding, |
| | 1 | 190 | | cancellationToken) |
| | 1 | 191 | | .ConfigureAwait(false); |
| | | 192 | | } |
| | | 193 | | |
| | | 194 | | // Otherwise own the transaction inside the execution strategy so a retrying provider replays the |
| | | 195 | | // read, the guarded update, and the insert together. |
| | 3 | 196 | | IExecutionStrategy executionStrategy = _dbContext.Database.CreateExecutionStrategy(); |
| | | 197 | | |
| | 3 | 198 | | return await executionStrategy.ExecuteAsync( |
| | 3 | 199 | | async strategyCancellationToken => |
| | 3 | 200 | | { |
| | 3 | 201 | | await using IDbContextTransaction transaction = await _dbContext.Database |
| | 3 | 202 | | .BeginTransactionAsync(strategyCancellationToken) |
| | 3 | 203 | | .ConfigureAwait(false); |
| | 3 | 204 | | |
| | 3 | 205 | | ApplicationAuditReconciliationRemediationItem remediation = await RecordRemediationCoreAsync( |
| | 3 | 206 | | findingId, |
| | 3 | 207 | | actionCode, |
| | 3 | 208 | | actorId, |
| | 3 | 209 | | evidenceReference, |
| | 3 | 210 | | request.ResolveFinding, |
| | 3 | 211 | | strategyCancellationToken) |
| | 3 | 212 | | .ConfigureAwait(false); |
| | 3 | 213 | | |
| | 3 | 214 | | await transaction.CommitAsync(strategyCancellationToken).ConfigureAwait(false); |
| | 3 | 215 | | return remediation; |
| | 3 | 216 | | }, |
| | 3 | 217 | | cancellationToken) |
| | 3 | 218 | | .ConfigureAwait(false); |
| | 2 | 219 | | } |
| | | 220 | | |
| | | 221 | | private async Task<ApplicationAuditReconciliationRemediationItem> RecordRemediationCoreAsync( |
| | | 222 | | Guid findingId, |
| | | 223 | | string actionCode, |
| | | 224 | | string actorId, |
| | | 225 | | string? evidenceReference, |
| | | 226 | | bool resolveFinding, |
| | | 227 | | CancellationToken cancellationToken) |
| | | 228 | | { |
| | 4 | 229 | | ApplicationAuditReconciliationFinding finding = await _dbContext |
| | 4 | 230 | | .ApplicationAuditReconciliationFindings |
| | 4 | 231 | | .AsNoTracking() |
| | 4 | 232 | | .SingleOrDefaultAsync(item => item.Id == findingId, cancellationToken) |
| | 4 | 233 | | .ConfigureAwait(false) |
| | 4 | 234 | | ?? throw new KeyNotFoundException($"Audit reconciliation finding '{findingId}' was not found."); |
| | | 235 | | |
| | 4 | 236 | | DateTime now = UtcNow(); |
| | 4 | 237 | | string remediationStatus = resolveFinding |
| | 4 | 238 | | ? ApplicationAuditReconciliationRemediationStatuses.Resolved |
| | 4 | 239 | | : ApplicationAuditReconciliationRemediationStatuses.Acknowledged; |
| | 4 | 240 | | DateTime? resolvedUtc = resolveFinding ? now : null; |
| | 4 | 241 | | string expectedFindingConcurrencyStamp = finding.ConcurrencyStamp; |
| | 4 | 242 | | string nextFindingConcurrencyStamp = Guid.NewGuid().ToString("N"); |
| | | 243 | | |
| | | 244 | | // Update the finding first, guarded by the stamp that was read. If a reconciliation run or another |
| | | 245 | | // remediation changed the finding in the meantime, no row matches and nothing is written. |
| | 4 | 246 | | int updatedFindingCount = await _dbContext.ApplicationAuditReconciliationFindings |
| | 4 | 247 | | .Where(item => item.Id == findingId && item.ConcurrencyStamp == expectedFindingConcurrencyStamp) |
| | 4 | 248 | | .ExecuteUpdateAsync( |
| | 4 | 249 | | setters => setters |
| | 4 | 250 | | .SetProperty(item => item.RemediationStatus, remediationStatus) |
| | 4 | 251 | | .SetProperty(item => item.ResolvedUtc, resolvedUtc) |
| | 4 | 252 | | .SetProperty(item => item.ConcurrencyStamp, nextFindingConcurrencyStamp), |
| | 4 | 253 | | cancellationToken) |
| | 4 | 254 | | .ConfigureAwait(false); |
| | | 255 | | |
| | 4 | 256 | | if (updatedFindingCount != 1) |
| | | 257 | | { |
| | 1 | 258 | | throw new DbUpdateConcurrencyException( |
| | 1 | 259 | | $"Audit reconciliation finding '{findingId}' was modified after it was read. " + |
| | 1 | 260 | | "No remediation was recorded. Reload the finding, verify its current state, and retry."); |
| | | 261 | | } |
| | | 262 | | |
| | 3 | 263 | | var remediationId = Guid.NewGuid(); |
| | 3 | 264 | | string remediationConcurrencyStamp = Guid.NewGuid().ToString("N"); |
| | | 265 | | |
| | 3 | 266 | | _ = await InsertAsync<ApplicationAuditReconciliationRemediation>( |
| | 3 | 267 | | [ |
| | 3 | 268 | | (nameof(ApplicationAuditReconciliationRemediation.Id), remediationId), |
| | 3 | 269 | | (nameof(ApplicationAuditReconciliationRemediation.FindingId), findingId), |
| | 3 | 270 | | (nameof(ApplicationAuditReconciliationRemediation.MutationBatchId), finding.MutationBatchId), |
| | 3 | 271 | | (nameof(ApplicationAuditReconciliationRemediation.ActionCode), actionCode), |
| | 3 | 272 | | (nameof(ApplicationAuditReconciliationRemediation.ActorId), actorId), |
| | 3 | 273 | | (nameof(ApplicationAuditReconciliationRemediation.EvidenceReference), evidenceReference), |
| | 3 | 274 | | (nameof(ApplicationAuditReconciliationRemediation.RecordedUtc), now), |
| | 3 | 275 | | (nameof(ApplicationAuditReconciliationRemediation.ConcurrencyStamp), remediationConcurrencyStamp) |
| | 3 | 276 | | ], |
| | 3 | 277 | | uniqueProperty: null, |
| | 3 | 278 | | cancellationToken) |
| | 3 | 279 | | .ConfigureAwait(false); |
| | | 280 | | |
| | 2 | 281 | | return new( |
| | 2 | 282 | | remediationId, |
| | 2 | 283 | | findingId, |
| | 2 | 284 | | finding.MutationBatchId, |
| | 2 | 285 | | actionCode, |
| | 2 | 286 | | actorId, |
| | 2 | 287 | | evidenceReference, |
| | 2 | 288 | | now); |
| | 2 | 289 | | } |
| | | 290 | | |
| | | 291 | | private List<ApplicationAuditReconciliationCandidate> BuildCandidates( |
| | | 292 | | IReadOnlyCollection<AuditRecord> auditRecords, |
| | | 293 | | IReadOnlyCollection<ApplicationAuditCompletionOutboxEntry> completionEntries, |
| | | 294 | | DateTime now) |
| | | 295 | | { |
| | 19 | 296 | | var candidates = new Dictionary<string, ApplicationAuditReconciliationCandidate>(StringComparer.Ordinal); |
| | 19 | 297 | | ILookup<string, AuditRecord> auditBatches = auditRecords |
| | 19 | 298 | | .Where(record => !string.IsNullOrWhiteSpace(record.MutationBatchId)) |
| | 19 | 299 | | .ToLookup(record => record.MutationBatchId, StringComparer.Ordinal); |
| | 19 | 300 | | ILookup<string, ApplicationAuditCompletionOutboxEntry> completionBatches = completionEntries |
| | 19 | 301 | | .Where(entry => !string.IsNullOrWhiteSpace(entry.MutationBatchId)) |
| | 19 | 302 | | .ToLookup(entry => entry.MutationBatchId, StringComparer.Ordinal); |
| | | 303 | | |
| | 42 | 304 | | foreach (AuditRecord malformed in auditRecords.Where(record => string.IsNullOrWhiteSpace(record.MutationBatchId) |
| | | 305 | | { |
| | 2 | 306 | | Add(candidates, Candidate( |
| | 2 | 307 | | ApplicationAuditReconciliationReasonCodes.MalformedCorrelation, |
| | 2 | 308 | | ApplicationAuditReconciliationSeverities.Error, |
| | 2 | 309 | | $"missing-{malformed.Id:N}", |
| | 2 | 310 | | null, |
| | 2 | 311 | | "Retain the row, investigate the originating save path, and append remediation evidence.")); |
| | | 312 | | } |
| | | 313 | | |
| | 76 | 314 | | foreach (IGrouping<string, AuditRecord> batch in auditBatches) |
| | | 315 | | { |
| | 19 | 316 | | List<AuditRecord> records = [.. batch]; |
| | 19 | 317 | | List<ApplicationAuditCompletionOutboxEntry> completions = [.. completionBatches[batch.Key]]; |
| | 19 | 318 | | DateTime newestRecordUtc = records.Max(record => record.ModifiedOnUtc); |
| | | 319 | | |
| | 19 | 320 | | if (completions.Count == 0 && newestRecordUtc <= now - _options.CompletionGracePeriod) |
| | | 321 | | { |
| | 14 | 322 | | Add(candidates, Candidate( |
| | 14 | 323 | | ApplicationAuditReconciliationReasonCodes.MissingCompletion, |
| | 14 | 324 | | ApplicationAuditReconciliationSeverities.Critical, |
| | 14 | 325 | | batch.Key, |
| | 14 | 326 | | null, |
| | 14 | 327 | | "Preserve the audit batch, investigate transaction completion, and append an operator remediation re |
| | | 328 | | } |
| | | 329 | | |
| | 50 | 330 | | foreach (ApplicationAuditCompletionOutboxEntry completion in completions) |
| | | 331 | | { |
| | 6 | 332 | | if (completion.AuditRecordCount != records.Count) |
| | | 333 | | { |
| | 1 | 334 | | Add(candidates, Candidate( |
| | 1 | 335 | | ApplicationAuditReconciliationReasonCodes.AuditRecordCountMismatch, |
| | 1 | 336 | | ApplicationAuditReconciliationSeverities.Critical, |
| | 1 | 337 | | batch.Key, |
| | 1 | 338 | | completion.Destination, |
| | 1 | 339 | | "Do not rewrite audit rows; compare retained evidence with the originating transaction and docum |
| | | 340 | | } |
| | | 341 | | |
| | 6 | 342 | | ApplicationMutationAuditReceipt receipt = ToReceipt(completion); |
| | 6 | 343 | | if (!_manifestVerifier.Verify(receipt, records)) |
| | | 344 | | { |
| | 2 | 345 | | Add(candidates, Candidate( |
| | 2 | 346 | | ApplicationAuditReconciliationReasonCodes.ManifestVerificationFailed, |
| | 2 | 347 | | ApplicationAuditReconciliationSeverities.Critical, |
| | 2 | 348 | | batch.Key, |
| | 2 | 349 | | completion.Destination, |
| | 2 | 350 | | "Quarantine downstream use of the batch, preserve all records, and investigate unauthorized or i |
| | | 351 | | } |
| | | 352 | | } |
| | | 353 | | |
| | 19 | 354 | | if (records.Any(record => record.State == "Added" && string.IsNullOrWhiteSpace(record.KeyValues))) |
| | | 355 | | { |
| | 0 | 356 | | Add(candidates, Candidate( |
| | 0 | 357 | | ApplicationAuditReconciliationReasonCodes.IncompleteGeneratedValues, |
| | 0 | 358 | | ApplicationAuditReconciliationSeverities.Error, |
| | 0 | 359 | | batch.Key, |
| | 0 | 360 | | null, |
| | 0 | 361 | | "Verify generated keys in the business database and append remediation evidence without modifying th |
| | | 362 | | } |
| | | 363 | | |
| | 19 | 364 | | if (HasMalformedCorrelation(records)) |
| | | 365 | | { |
| | 0 | 366 | | Add(candidates, Candidate( |
| | 0 | 367 | | ApplicationAuditReconciliationReasonCodes.MalformedCorrelation, |
| | 0 | 368 | | ApplicationAuditReconciliationSeverities.Warning, |
| | 0 | 369 | | batch.Key, |
| | 0 | 370 | | null, |
| | 0 | 371 | | "Investigate inconsistent correlation metadata and preserve the original records as evidence.")); |
| | | 372 | | } |
| | | 373 | | } |
| | | 374 | | |
| | 48 | 375 | | foreach (IGrouping<string, ApplicationAuditCompletionOutboxEntry> batch in completionBatches) |
| | | 376 | | { |
| | 5 | 377 | | if (!auditBatches.Contains(batch.Key)) |
| | | 378 | | { |
| | 0 | 379 | | foreach (ApplicationAuditCompletionOutboxEntry completion in batch) |
| | | 380 | | { |
| | 0 | 381 | | Add(candidates, Candidate( |
| | 0 | 382 | | ApplicationAuditReconciliationReasonCodes.MissingAuditBatch, |
| | 0 | 383 | | ApplicationAuditReconciliationSeverities.Critical, |
| | 0 | 384 | | batch.Key, |
| | 0 | 385 | | completion.Destination, |
| | 0 | 386 | | "Preserve the completion record and investigate missing or externally stored audit evidence.")); |
| | | 387 | | } |
| | | 388 | | } |
| | | 389 | | } |
| | | 390 | | |
| | 29 | 391 | | foreach (IGrouping<(string MutationBatchId, string Destination), ApplicationAuditCompletionOutboxEntry> duplicat |
| | 19 | 392 | | completionEntries.GroupBy(entry => (entry.MutationBatchId, entry.Destination))) |
| | | 393 | | { |
| | 5 | 394 | | if (duplicate.Count() > 1) |
| | | 395 | | { |
| | 1 | 396 | | Add(candidates, Candidate( |
| | 1 | 397 | | ApplicationAuditReconciliationReasonCodes.DuplicateCompletion, |
| | 1 | 398 | | ApplicationAuditReconciliationSeverities.Critical, |
| | 1 | 399 | | duplicate.Key.MutationBatchId, |
| | 1 | 400 | | duplicate.Key.Destination, |
| | 1 | 401 | | "Preserve all records, stop dispatch for the destination, and investigate uniqueness or migration dr |
| | | 402 | | } |
| | | 403 | | } |
| | | 404 | | |
| | 50 | 405 | | foreach (ApplicationAuditCompletionOutboxEntry entry in completionEntries) |
| | | 406 | | { |
| | 6 | 407 | | AddDeliveryCandidate(candidates, entry, now); |
| | | 408 | | } |
| | | 409 | | |
| | 19 | 410 | | return [.. candidates.Values]; |
| | | 411 | | } |
| | | 412 | | |
| | | 413 | | private void AddDeliveryCandidate( |
| | | 414 | | IDictionary<string, ApplicationAuditReconciliationCandidate> candidates, |
| | | 415 | | ApplicationAuditCompletionOutboxEntry entry, |
| | | 416 | | DateTime now) |
| | | 417 | | { |
| | 6 | 418 | | if ((entry.Status == ApplicationAuditCompletionOutboxStatuses.Pending || |
| | 6 | 419 | | entry.Status == ApplicationAuditCompletionOutboxStatuses.Deferred) && |
| | 6 | 420 | | entry.CreatedUtc <= now - _options.StalePendingThreshold) |
| | | 421 | | { |
| | 1 | 422 | | Add(candidates, Candidate( |
| | 1 | 423 | | ApplicationAuditReconciliationReasonCodes.StalePending, |
| | 1 | 424 | | ApplicationAuditReconciliationSeverities.Warning, |
| | 1 | 425 | | entry.MutationBatchId, |
| | 1 | 426 | | entry.Destination, |
| | 1 | 427 | | "Verify dispatcher availability and destination registration before retrying delivery.")); |
| | | 428 | | } |
| | 5 | 429 | | else if (entry.Status == ApplicationAuditCompletionOutboxStatuses.RetryableFailure && |
| | 5 | 430 | | entry.NextAttemptUtc <= now - _options.StaleRetryReadyThreshold) |
| | | 431 | | { |
| | 0 | 432 | | Add(candidates, Candidate( |
| | 0 | 433 | | ApplicationAuditReconciliationReasonCodes.StaleRetryReady, |
| | 0 | 434 | | ApplicationAuditReconciliationSeverities.Error, |
| | 0 | 435 | | entry.MutationBatchId, |
| | 0 | 436 | | entry.Destination, |
| | 0 | 437 | | "Inspect destination availability and retry policy; preserve prior attempt diagnostics.")); |
| | | 438 | | } |
| | 5 | 439 | | else if (entry.Status == ApplicationAuditCompletionOutboxStatuses.Failed) |
| | | 440 | | { |
| | 0 | 441 | | Add(candidates, Candidate( |
| | 0 | 442 | | ApplicationAuditReconciliationReasonCodes.DeliveryFailed, |
| | 0 | 443 | | ApplicationAuditReconciliationSeverities.Error, |
| | 0 | 444 | | entry.MutationBatchId, |
| | 0 | 445 | | entry.Destination, |
| | 0 | 446 | | "Investigate the terminal delivery failure and append operator remediation evidence.")); |
| | | 447 | | } |
| | 5 | 448 | | else if (entry.Status == ApplicationAuditCompletionOutboxStatuses.DeadLettered) |
| | | 449 | | { |
| | 0 | 450 | | Add(candidates, Candidate( |
| | 0 | 451 | | ApplicationAuditReconciliationReasonCodes.DeadLettered, |
| | 0 | 452 | | ApplicationAuditReconciliationSeverities.Critical, |
| | 0 | 453 | | entry.MutationBatchId, |
| | 0 | 454 | | entry.Destination, |
| | 0 | 455 | | "Review the dead letter, preserve diagnostics, correct the destination, and explicitly requeue only unde |
| | | 456 | | } |
| | 5 | 457 | | } |
| | | 458 | | |
| | | 459 | | private async Task PersistCandidatesAsync( |
| | | 460 | | IReadOnlyCollection<ApplicationAuditReconciliationCandidate> candidates, |
| | | 461 | | IReadOnlyCollection<string> auditBatchIds, |
| | | 462 | | IReadOnlyCollection<ApplicationAuditCompletionOutboxEntry> completionEntries, |
| | | 463 | | DateTime now, |
| | | 464 | | CancellationToken cancellationToken) |
| | | 465 | | { |
| | | 466 | | // When the caller already owns an EF Core transaction, join it and let the caller decide whether |
| | | 467 | | // the reconciliation writes commit. |
| | 19 | 468 | | if (_dbContext.Database.CurrentTransaction is not null) |
| | | 469 | | { |
| | 1 | 470 | | await PersistCandidatesCoreAsync(candidates, auditBatchIds, completionEntries, now, cancellationToken) |
| | 1 | 471 | | .ConfigureAwait(false); |
| | 1 | 472 | | return; |
| | | 473 | | } |
| | | 474 | | |
| | | 475 | | // Otherwise persist every finding change in one transaction inside the execution strategy, so a run |
| | | 476 | | // applies all of its inserts and guarded updates or none of them, and a retrying provider replays |
| | | 477 | | // the reads together with the writes. |
| | 18 | 478 | | IExecutionStrategy executionStrategy = _dbContext.Database.CreateExecutionStrategy(); |
| | | 479 | | |
| | 18 | 480 | | await executionStrategy.ExecuteAsync( |
| | 18 | 481 | | async strategyCancellationToken => |
| | 18 | 482 | | { |
| | 18 | 483 | | await using IDbContextTransaction transaction = await _dbContext.Database |
| | 18 | 484 | | .BeginTransactionAsync(strategyCancellationToken) |
| | 18 | 485 | | .ConfigureAwait(false); |
| | 18 | 486 | | |
| | 18 | 487 | | await PersistCandidatesCoreAsync( |
| | 18 | 488 | | candidates, |
| | 18 | 489 | | auditBatchIds, |
| | 18 | 490 | | completionEntries, |
| | 18 | 491 | | now, |
| | 18 | 492 | | strategyCancellationToken) |
| | 18 | 493 | | .ConfigureAwait(false); |
| | 18 | 494 | | |
| | 18 | 495 | | await transaction.CommitAsync(strategyCancellationToken).ConfigureAwait(false); |
| | 18 | 496 | | }, |
| | 18 | 497 | | cancellationToken) |
| | 18 | 498 | | .ConfigureAwait(false); |
| | 17 | 499 | | } |
| | | 500 | | |
| | | 501 | | private async Task PersistCandidatesCoreAsync( |
| | | 502 | | IReadOnlyCollection<ApplicationAuditReconciliationCandidate> candidates, |
| | | 503 | | IReadOnlyCollection<string> auditBatchIds, |
| | | 504 | | IReadOnlyCollection<ApplicationAuditCompletionOutboxEntry> completionEntries, |
| | | 505 | | DateTime now, |
| | | 506 | | CancellationToken cancellationToken) |
| | | 507 | | { |
| | 19 | 508 | | string[] keys = [.. candidates.Select(candidate => candidate.FindingKey)]; |
| | 19 | 509 | | string[] scopeBatchIds = [.. auditBatchIds |
| | 19 | 510 | | .Concat(completionEntries.Select(entry => entry.MutationBatchId)) |
| | 19 | 511 | | .Where(batchId => !string.IsNullOrWhiteSpace(batchId)) |
| | 19 | 512 | | .Distinct(StringComparer.Ordinal)]; |
| | | 513 | | |
| | | 514 | | // Read the active and the resolvable findings in one round trip. Each row's stamp guards its update. |
| | 19 | 515 | | List<ApplicationAuditReconciliationFinding> existing = await _dbContext |
| | 19 | 516 | | .ApplicationAuditReconciliationFindings |
| | 19 | 517 | | .AsNoTracking() |
| | 19 | 518 | | .Where(finding => keys.Contains(finding.FindingKey) || |
| | 19 | 519 | | (scopeBatchIds.Contains(finding.MutationBatchId) && |
| | 19 | 520 | | finding.RemediationStatus != ApplicationAuditReconciliationRemediationStatuses.Resolved)) |
| | 19 | 521 | | .ToListAsync(cancellationToken) |
| | 19 | 522 | | .ConfigureAwait(false); |
| | 19 | 523 | | var existingByKey = existing |
| | 19 | 524 | | .ToDictionary(finding => finding.FindingKey, StringComparer.Ordinal); |
| | 19 | 525 | | var activeKeys = new HashSet<string>(keys, StringComparer.Ordinal); |
| | 19 | 526 | | var scopeBatchIdSet = new HashSet<string>(scopeBatchIds, StringComparer.Ordinal); |
| | | 527 | | |
| | 76 | 528 | | foreach (ApplicationAuditReconciliationCandidate candidate in candidates) |
| | | 529 | | { |
| | 20 | 530 | | if (existingByKey.TryGetValue(candidate.FindingKey, out ApplicationAuditReconciliationFinding? finding)) |
| | | 531 | | { |
| | 2 | 532 | | Guid findingId = finding.Id; |
| | 2 | 533 | | string expectedConcurrencyStamp = finding.ConcurrencyStamp; |
| | 2 | 534 | | string nextConcurrencyStamp = Guid.NewGuid().ToString("N"); |
| | 2 | 535 | | int updatedCount = await _dbContext.ApplicationAuditReconciliationFindings |
| | 2 | 536 | | .Where(item => item.Id == findingId && item.ConcurrencyStamp == expectedConcurrencyStamp) |
| | 2 | 537 | | .ExecuteUpdateAsync( |
| | 2 | 538 | | setters => setters |
| | 2 | 539 | | .SetProperty(item => item.Severity, candidate.Severity) |
| | 2 | 540 | | .SetProperty(item => item.Guidance, candidate.Guidance) |
| | 2 | 541 | | .SetProperty(item => item.LastObservedUtc, now) |
| | 2 | 542 | | .SetProperty(item => item.RemediationStatus, ApplicationAuditReconciliationRemediationStatus |
| | 2 | 543 | | .SetProperty(item => item.ResolvedUtc, (DateTime?)null) |
| | 2 | 544 | | .SetProperty(item => item.ConcurrencyStamp, nextConcurrencyStamp), |
| | 2 | 545 | | cancellationToken) |
| | 2 | 546 | | .ConfigureAwait(false); |
| | | 547 | | |
| | 2 | 548 | | EnsureSingleRowWritten(updatedCount, candidate.FindingKey); |
| | | 549 | | } |
| | | 550 | | else |
| | | 551 | | { |
| | | 552 | | // Insert only when no row holds the key, so a run that interleaved and inserted the same finding |
| | | 553 | | // surfaces as a concurrency conflict (and a rollback) rather than a unique index violation. |
| | 18 | 554 | | int insertedCount = await InsertAsync<ApplicationAuditReconciliationFinding>( |
| | 18 | 555 | | [ |
| | 18 | 556 | | (nameof(ApplicationAuditReconciliationFinding.Id), Guid.NewGuid()), |
| | 18 | 557 | | (nameof(ApplicationAuditReconciliationFinding.SchemaVersion), ApplicationAuditReconciliation |
| | 18 | 558 | | (nameof(ApplicationAuditReconciliationFinding.FindingKey), candidate.FindingKey), |
| | 18 | 559 | | (nameof(ApplicationAuditReconciliationFinding.ReasonCode), candidate.ReasonCode), |
| | 18 | 560 | | (nameof(ApplicationAuditReconciliationFinding.Severity), candidate.Severity), |
| | 18 | 561 | | (nameof(ApplicationAuditReconciliationFinding.MutationBatchId), candidate.MutationBatchId), |
| | 18 | 562 | | (nameof(ApplicationAuditReconciliationFinding.Destination), candidate.Destination), |
| | 18 | 563 | | (nameof(ApplicationAuditReconciliationFinding.Guidance), candidate.Guidance), |
| | 18 | 564 | | (nameof(ApplicationAuditReconciliationFinding.RemediationStatus), ApplicationAuditReconcilia |
| | 18 | 565 | | (nameof(ApplicationAuditReconciliationFinding.FirstObservedUtc), now), |
| | 18 | 566 | | (nameof(ApplicationAuditReconciliationFinding.LastObservedUtc), now), |
| | 18 | 567 | | (nameof(ApplicationAuditReconciliationFinding.ResolvedUtc), null), |
| | 18 | 568 | | (nameof(ApplicationAuditReconciliationFinding.ConcurrencyStamp), Guid.NewGuid().ToString("N" |
| | 18 | 569 | | ], |
| | 18 | 570 | | uniqueProperty: nameof(ApplicationAuditReconciliationFinding.FindingKey), |
| | 18 | 571 | | cancellationToken) |
| | 18 | 572 | | .ConfigureAwait(false); |
| | | 573 | | |
| | 18 | 574 | | EnsureSingleRowWritten(insertedCount, candidate.FindingKey); |
| | | 575 | | } |
| | 18 | 576 | | } |
| | | 577 | | |
| | 34 | 578 | | foreach (ApplicationAuditReconciliationFinding finding in existing.Where(finding => |
| | 17 | 579 | | !activeKeys.Contains(finding.FindingKey) && |
| | 17 | 580 | | scopeBatchIdSet.Contains(finding.MutationBatchId) && |
| | 17 | 581 | | finding.RemediationStatus != ApplicationAuditReconciliationRemediationStatuses.Resolved)) |
| | | 582 | | { |
| | 0 | 583 | | Guid findingId = finding.Id; |
| | 0 | 584 | | string expectedConcurrencyStamp = finding.ConcurrencyStamp; |
| | 0 | 585 | | string nextConcurrencyStamp = Guid.NewGuid().ToString("N"); |
| | 0 | 586 | | int resolvedCount = await _dbContext.ApplicationAuditReconciliationFindings |
| | 0 | 587 | | .Where(item => item.Id == findingId && item.ConcurrencyStamp == expectedConcurrencyStamp) |
| | 0 | 588 | | .ExecuteUpdateAsync( |
| | 0 | 589 | | setters => setters |
| | 0 | 590 | | .SetProperty(item => item.RemediationStatus, ApplicationAuditReconciliationRemediationStatuses.R |
| | 0 | 591 | | .SetProperty(item => item.ResolvedUtc, (DateTime?)now) |
| | 0 | 592 | | .SetProperty(item => item.ConcurrencyStamp, nextConcurrencyStamp), |
| | 0 | 593 | | cancellationToken) |
| | 0 | 594 | | .ConfigureAwait(false); |
| | | 595 | | |
| | 0 | 596 | | EnsureSingleRowWritten(resolvedCount, finding.FindingKey); |
| | 0 | 597 | | } |
| | 17 | 598 | | } |
| | | 599 | | |
| | | 600 | | // Inserts one row without going through SaveChanges (so the save pipeline does not audit reconciliation |
| | | 601 | | // bookkeeping). Table and column names come from the EF Core model and are quoted by the provider's |
| | | 602 | | // ISqlGenerationHelper, so the statement follows entity configuration and works on any relational provider |
| | | 603 | | // that accepts INSERT ... SELECT without FROM (SQL Server, SQLite, PostgreSQL). Values are always parameters. |
| | | 604 | | // When uniqueProperty is set, the row is written only if no row already holds that property's value. |
| | | 605 | | private Task<int> InsertAsync<TEntity>( |
| | | 606 | | IReadOnlyList<(string Property, object? Value)> values, |
| | | 607 | | string? uniqueProperty, |
| | | 608 | | CancellationToken cancellationToken) |
| | | 609 | | where TEntity : class |
| | | 610 | | { |
| | 21 | 611 | | IEntityType entityType = _dbContext.Model.FindEntityType(typeof(TEntity)) |
| | 21 | 612 | | ?? throw new InvalidOperationException($"'{typeof(TEntity).Name}' is not part of the EF Core model."); |
| | 21 | 613 | | string tableName = entityType.GetTableName() |
| | 21 | 614 | | ?? throw new InvalidOperationException($"'{typeof(TEntity).Name}' is not mapped to a table."); |
| | 21 | 615 | | var table = StoreObjectIdentifier.Table(tableName, entityType.GetSchema()); |
| | 21 | 616 | | ISqlGenerationHelper sqlGenerationHelper = _dbContext.GetService<ISqlGenerationHelper>(); |
| | | 617 | | |
| | | 618 | | string Column(string propertyName) |
| | | 619 | | { |
| | | 620 | | IProperty property = entityType.FindProperty(propertyName) |
| | | 621 | | ?? throw new InvalidOperationException($"'{typeof(TEntity).Name}.{propertyName}' is not mapped."); |
| | | 622 | | return EscapeFormatBraces(sqlGenerationHelper.DelimitIdentifier(property.GetColumnName(table) |
| | | 623 | | ?? throw new InvalidOperationException($"'{typeof(TEntity).Name}.{propertyName}' has no column."))); |
| | | 624 | | } |
| | | 625 | | |
| | 21 | 626 | | string delimitedTable = EscapeFormatBraces(sqlGenerationHelper.DelimitIdentifier(table.Name, table.Schema)); |
| | 21 | 627 | | var parameters = values.Select(value => value.Value).ToList(); |
| | 21 | 628 | | StringBuilder sql = new StringBuilder() |
| | 21 | 629 | | .Append("INSERT INTO ").Append(delimitedTable) |
| | 21 | 630 | | .Append(" (").AppendJoin(", ", values.Select(value => Column(value.Property))).Append(") SELECT ") |
| | 21 | 631 | | .AppendJoin(", ", values.Select((_, index) => $"{{{index}}}")); |
| | | 632 | | |
| | 21 | 633 | | if (uniqueProperty is not null) |
| | | 634 | | { |
| | 18 | 635 | | object? uniqueValue = values.Single(value => value.Property == uniqueProperty).Value; |
| | 18 | 636 | | _ = sql.Append(" WHERE NOT EXISTS (SELECT 1 FROM ").Append(delimitedTable) |
| | 18 | 637 | | .Append(" WHERE ").Append(Column(uniqueProperty)).Append(" = {").Append(parameters.Count).Append("})"); |
| | 18 | 638 | | parameters.Add(uniqueValue); |
| | | 639 | | } |
| | | 640 | | |
| | | 641 | | // Only model-derived identifiers are part of the format string; every value is a format argument, which |
| | | 642 | | // ExecuteSqlAsync binds as a DbParameter. |
| | 21 | 643 | | return _dbContext.Database.ExecuteSqlAsync( |
| | 21 | 644 | | FormattableStringFactory.Create(sql.ToString(), [.. parameters]), |
| | 21 | 645 | | cancellationToken); |
| | | 646 | | } |
| | | 647 | | |
| | | 648 | | private static string EscapeFormatBraces(string identifier) |
| | | 649 | | { |
| | 297 | 650 | | return identifier.Replace("{", "{{", StringComparison.Ordinal).Replace("}", "}}", StringComparison.Ordinal); |
| | | 651 | | } |
| | | 652 | | |
| | | 653 | | private static void EnsureSingleRowWritten(int affectedCount, string findingKey) |
| | | 654 | | { |
| | 20 | 655 | | if (affectedCount != 1) |
| | | 656 | | { |
| | 2 | 657 | | throw new DbUpdateConcurrencyException( |
| | 2 | 658 | | $"Audit reconciliation finding '{findingKey}' was modified by another writer after it was read. " + |
| | 2 | 659 | | "No reconciliation changes were committed; the next reconciliation run re-evaluates the finding."); |
| | | 660 | | } |
| | 18 | 661 | | } |
| | | 662 | | |
| | | 663 | | private async Task<ApplicationAuditReconciliationSummary> GetSummaryCoreAsync( |
| | | 664 | | DateTime? lastRunUtc, |
| | | 665 | | CancellationToken cancellationToken) |
| | | 666 | | { |
| | 18 | 667 | | IQueryable<ApplicationAuditReconciliationFinding> open = _dbContext |
| | 18 | 668 | | .ApplicationAuditReconciliationFindings |
| | 18 | 669 | | .AsNoTracking() |
| | 18 | 670 | | .Where(finding => finding.RemediationStatus != ApplicationAuditReconciliationRemediationStatuses.Resolved); |
| | | 671 | | |
| | 18 | 672 | | long openCount = await open.LongCountAsync(cancellationToken).ConfigureAwait(false); |
| | 18 | 673 | | long errorCount = await open.LongCountAsync( |
| | 18 | 674 | | finding => finding.Severity == ApplicationAuditReconciliationSeverities.Error, |
| | 18 | 675 | | cancellationToken).ConfigureAwait(false); |
| | 18 | 676 | | long criticalCount = await open.LongCountAsync( |
| | 18 | 677 | | finding => finding.Severity == ApplicationAuditReconciliationSeverities.Critical, |
| | 18 | 678 | | cancellationToken).ConfigureAwait(false); |
| | 18 | 679 | | long manifestFailures = await open.LongCountAsync( |
| | 18 | 680 | | finding => finding.ReasonCode == ApplicationAuditReconciliationReasonCodes.ManifestVerificationFailed, |
| | 18 | 681 | | cancellationToken).ConfigureAwait(false); |
| | 18 | 682 | | long missingCompletion = await open.LongCountAsync( |
| | 18 | 683 | | finding => finding.ReasonCode == ApplicationAuditReconciliationReasonCodes.MissingCompletion, |
| | 18 | 684 | | cancellationToken).ConfigureAwait(false); |
| | 18 | 685 | | long staleDelivery = await open.LongCountAsync( |
| | 18 | 686 | | finding => finding.ReasonCode == ApplicationAuditReconciliationReasonCodes.StalePending || |
| | 18 | 687 | | finding.ReasonCode == ApplicationAuditReconciliationReasonCodes.StaleRetryReady || |
| | 18 | 688 | | finding.ReasonCode == ApplicationAuditReconciliationReasonCodes.DeliveryFailed, |
| | 18 | 689 | | cancellationToken).ConfigureAwait(false); |
| | 18 | 690 | | long deadLetters = await open.LongCountAsync( |
| | 18 | 691 | | finding => finding.ReasonCode == ApplicationAuditReconciliationReasonCodes.DeadLettered, |
| | 18 | 692 | | cancellationToken).ConfigureAwait(false); |
| | | 693 | | |
| | 18 | 694 | | return new( |
| | 18 | 695 | | true, |
| | 18 | 696 | | lastRunUtc, |
| | 18 | 697 | | openCount, |
| | 18 | 698 | | errorCount, |
| | 18 | 699 | | criticalCount, |
| | 18 | 700 | | manifestFailures, |
| | 18 | 701 | | missingCompletion, |
| | 18 | 702 | | staleDelivery, |
| | 18 | 703 | | deadLetters); |
| | 18 | 704 | | } |
| | | 705 | | |
| | | 706 | | private static bool HasMalformedCorrelation(IReadOnlyCollection<AuditRecord> records) |
| | | 707 | | { |
| | 19 | 708 | | return records.Any(record => |
| | 19 | 709 | | IsWhitespaceOnly(record.OperationExecutionId) || |
| | 19 | 710 | | IsWhitespaceOnly(record.ExecutionAttemptId) || |
| | 19 | 711 | | IsWhitespaceOnly(record.DecisionAuditRecordId) || |
| | 19 | 712 | | IsWhitespaceOnly(record.CorrelationId) || |
| | 19 | 713 | | IsWhitespaceOnly(record.TraceId)) || |
| | 19 | 714 | | HasMultipleValues(records.Select(record => record.OperationExecutionId)) || |
| | 19 | 715 | | HasMultipleValues(records.Select(record => record.ExecutionAttemptId)) || |
| | 19 | 716 | | HasMultipleValues(records.Select(record => record.DecisionAuditRecordId)) || |
| | 19 | 717 | | HasMultipleValues(records.Select(record => record.CorrelationId)) || |
| | 19 | 718 | | HasMultipleValues(records.Select(record => record.TraceId)); |
| | | 719 | | } |
| | | 720 | | |
| | | 721 | | private static bool HasMultipleValues(IEnumerable<string?> values) |
| | | 722 | | { |
| | 95 | 723 | | return values |
| | 95 | 724 | | .Where(value => !string.IsNullOrWhiteSpace(value)) |
| | 95 | 725 | | .Select(value => value!.Trim()) |
| | 95 | 726 | | .Distinct(StringComparer.Ordinal) |
| | 95 | 727 | | .Skip(1) |
| | 95 | 728 | | .Any(); |
| | | 729 | | } |
| | | 730 | | |
| | | 731 | | private static bool IsWhitespaceOnly(string? value) |
| | | 732 | | { |
| | 95 | 733 | | return value is not null && string.IsNullOrWhiteSpace(value); |
| | | 734 | | } |
| | | 735 | | |
| | | 736 | | private static ApplicationMutationAuditReceipt ToReceipt(ApplicationAuditCompletionOutboxEntry entry) |
| | | 737 | | { |
| | 6 | 738 | | return new( |
| | 6 | 739 | | entry.MutationBatchId, |
| | 6 | 740 | | entry.AuditRecordCount, |
| | 6 | 741 | | entry.PersistenceOutcome, |
| | 6 | 742 | | new DateTimeOffset(DateTime.SpecifyKind(entry.ReceiptCompletedUtc, DateTimeKind.Utc)), |
| | 6 | 743 | | entry.MutationManifestHash, |
| | 6 | 744 | | entry.MutationManifestAlgorithm, |
| | 6 | 745 | | entry.MutationManifestSchemaVersion, |
| | 6 | 746 | | entry.OperationExecutionId, |
| | 6 | 747 | | entry.ExecutionAttemptId, |
| | 6 | 748 | | entry.DecisionAuditRecordId, |
| | 6 | 749 | | entry.CorrelationId, |
| | 6 | 750 | | entry.TraceId); |
| | | 751 | | } |
| | | 752 | | |
| | | 753 | | private static ApplicationAuditReconciliationCandidate Candidate( |
| | | 754 | | string reasonCode, |
| | | 755 | | string severity, |
| | | 756 | | string batchId, |
| | | 757 | | string? destination, |
| | | 758 | | string guidance) |
| | | 759 | | { |
| | 21 | 760 | | string normalizedBatchId = NormalizeRequired(batchId, 64, nameof(batchId)); |
| | 21 | 761 | | string? normalizedDestination = NormalizeOptional(destination, 128); |
| | 21 | 762 | | string identity = $"{reasonCode}\n{normalizedBatchId}\n{normalizedDestination}"; |
| | 21 | 763 | | string key = Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(identity))); |
| | 21 | 764 | | return new(key, reasonCode, severity, normalizedBatchId, normalizedDestination, guidance); |
| | | 765 | | } |
| | | 766 | | |
| | | 767 | | private static void Add( |
| | | 768 | | IDictionary<string, ApplicationAuditReconciliationCandidate> candidates, |
| | | 769 | | ApplicationAuditReconciliationCandidate candidate) |
| | | 770 | | { |
| | 21 | 771 | | candidates[candidate.FindingKey] = candidate; |
| | 21 | 772 | | } |
| | | 773 | | |
| | | 774 | | private static string NormalizeRequired(string value, int maximumLength, string parameterName) |
| | | 775 | | { |
| | 30 | 776 | | ArgumentException.ThrowIfNullOrWhiteSpace(value, parameterName); |
| | 29 | 777 | | string normalized = value.Trim(); |
| | 29 | 778 | | return normalized.Length <= maximumLength |
| | 29 | 779 | | ? normalized |
| | 29 | 780 | | : throw new ArgumentOutOfRangeException(parameterName, $"Value cannot exceed {maximumLength} characters."); |
| | | 781 | | } |
| | | 782 | | |
| | | 783 | | private static string? NormalizeOptional(string? value, int maximumLength) |
| | | 784 | | { |
| | 25 | 785 | | if (string.IsNullOrWhiteSpace(value)) |
| | | 786 | | { |
| | 16 | 787 | | return null; |
| | | 788 | | } |
| | | 789 | | |
| | 9 | 790 | | string normalized = value.Trim(); |
| | 9 | 791 | | return normalized.Length <= maximumLength |
| | 9 | 792 | | ? normalized |
| | 9 | 793 | | : normalized[..maximumLength]; |
| | | 794 | | } |
| | | 795 | | |
| | | 796 | | private DateTime UtcNow() |
| | | 797 | | { |
| | 23 | 798 | | return _timeProvider.GetUtcNow().UtcDateTime; |
| | | 799 | | } |
| | | 800 | | } |