< Summary

Information
Class: ProjectTemplate.Infrastructure.Data.Auditing.ApplicationAuditReconciler
Assembly: ProjectTemplate.Infrastructure
File(s): /home/runner/work/NetCoreApplicationTemplate/NetCoreApplicationTemplate/src/ProjectTemplate.Infrastructure/Data/Auditing/ApplicationAuditReconciler.cs
Line coverage
90%
Covered lines: 467
Uncovered lines: 51
Coverable lines: 518
Total lines: 800
Line coverage: 90.1%
Branch coverage
72%
Covered branches: 84
Total branches: 116
Branch coverage: 72.4%
Method coverage

Feature is only available for sponsors

Upgrade to PRO version

Metrics

MethodBranch coverage Crap Score Cyclomatic complexity Line coverage
.ctor(...)50%1212100%
ReconcileAsync()100%22100%
GetSummaryAsync()100%22100%
QueryFindingsAsync()100%88100%
RecordRemediationAsync()100%22100%
RecordRemediationCoreAsync()62.5%88100%
BuildCandidates(...)83.33%443075%
AddDeliveryCandidate(...)56.25%621643.75%
PersistCandidatesAsync()100%22100%
PersistCandidatesCoreAsync()87.5%8881.33%
InsertAsync(...)66.66%66100%
EscapeFormatBraces(...)100%11100%
EnsureSingleRowWritten(...)100%22100%
GetSummaryCoreAsync()100%11100%
HasMalformedCorrelation(...)50%1010100%
HasMultipleValues(...)100%11100%
IsWhitespaceOnly(...)50%22100%
ToReceipt(...)100%11100%
Candidate(...)100%11100%
Add(...)100%11100%
NormalizeRequired(...)50%22100%
NormalizeOptional(...)75%44100%
UtcNow()100%11100%

File(s)

/home/runner/work/NetCoreApplicationTemplate/NetCoreApplicationTemplate/src/ProjectTemplate.Infrastructure/Data/Auditing/ApplicationAuditReconciler.cs

#LineLine coverage
 1using System.Runtime.CompilerServices;
 2using System.Security.Cryptography;
 3using System.Text;
 4using Microsoft.EntityFrameworkCore;
 5using Microsoft.EntityFrameworkCore.Infrastructure;
 6using Microsoft.EntityFrameworkCore.Metadata;
 7using Microsoft.EntityFrameworkCore.Storage;
 8using Microsoft.Extensions.Options;
 9using ProjectTemplate.Infrastructure.Data.Entities;
 10
 11namespace ProjectTemplate.Infrastructure.Data.Auditing;
 12
 1913public sealed class ApplicationAuditReconciler(
 1914    ApplicationDbContext dbContext,
 1915    IApplicationMutationManifestVerifier manifestVerifier,
 1916    IOptions<ApplicationAuditReconciliationOptions> options,
 1917    ApplicationAuditReconciliationMetrics metrics,
 1918    TimeProvider timeProvider)
 19    : IApplicationAuditReconciler
 20{
 1921    private readonly ApplicationDbContext _dbContext =
 1922        dbContext ?? throw new ArgumentNullException(nameof(dbContext));
 1923    private readonly IApplicationMutationManifestVerifier _manifestVerifier =
 1924        manifestVerifier ?? throw new ArgumentNullException(nameof(manifestVerifier));
 1925    private readonly ApplicationAuditReconciliationOptions _options =
 1926        options?.Value ?? throw new ArgumentNullException(nameof(options));
 1927    private readonly ApplicationAuditReconciliationMetrics _metrics =
 1928        metrics ?? throw new ArgumentNullException(nameof(metrics));
 1929    private readonly TimeProvider _timeProvider =
 1930        timeProvider ?? throw new ArgumentNullException(nameof(timeProvider));
 31
 32    public async Task<ApplicationAuditReconciliationSummary> ReconcileAsync(
 33        CancellationToken cancellationToken = default)
 34    {
 2035        cancellationToken.ThrowIfCancellationRequested();
 2036        if (!_options.Enabled)
 37        {
 138            return ApplicationAuditReconciliationMetrics.DisabledSummary;
 39        }
 40
 1941        DateTime now = UtcNow();
 1942        List<string> batchIds = await _dbContext.AuditRecords
 1943            .AsNoTracking()
 1944            .Where(record => record.MutationBatchId != string.Empty)
 1945            .GroupBy(record => record.MutationBatchId)
 1946            .OrderByDescending(group => group.Max(record => record.ModifiedOnUtc))
 1947            .Select(group => group.Key)
 1948            .Take(_options.MaximumBatchesPerRun)
 1949            .ToListAsync(cancellationToken)
 1950            .ConfigureAwait(false);
 51
 1952        List<AuditRecord> auditRecords = await _dbContext.AuditRecords
 1953            .AsNoTracking()
 1954            .Where(record => batchIds.Contains(record.MutationBatchId))
 1955            .ToListAsync(cancellationToken)
 1956            .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.
 1960        List<AuditRecord> malformedRecords = await _dbContext.AuditRecords
 1961            .AsNoTracking()
 1962            .Where(record => record.MutationBatchId == string.Empty)
 1963            .OrderByDescending(record => record.ModifiedOnUtc)
 1964            .ThenByDescending(record => record.Id)
 1965            .Take(_options.MaximumMalformedRecordsPerRun)
 1966            .ToListAsync(cancellationToken)
 1967            .ConfigureAwait(false);
 1968        auditRecords.AddRange(malformedRecords);
 69
 1970        List<ApplicationAuditCompletionOutboxEntry> completionEntries = await _dbContext
 1971            .ApplicationAuditCompletionOutboxEntries
 1972            .AsNoTracking()
 1973            .Where(entry => batchIds.Contains(entry.MutationBatchId) ||
 1974                entry.Status != ApplicationAuditCompletionOutboxStatuses.Delivered)
 1975            .OrderByDescending(entry => entry.CreatedUtc)
 1976            .Take(_options.MaximumBatchesPerRun * 2)
 1977            .ToListAsync(cancellationToken)
 1978            .ConfigureAwait(false);
 79
 1980        List<ApplicationAuditReconciliationCandidate> candidates = BuildCandidates(
 1981            auditRecords,
 1982            completionEntries,
 1983            now);
 84
 1985        await PersistCandidatesAsync(candidates, batchIds, completionEntries, now, cancellationToken)
 1986            .ConfigureAwait(false);
 87
 1788        ApplicationAuditReconciliationSummary summary = await GetSummaryCoreAsync(now, cancellationToken)
 1789            .ConfigureAwait(false);
 1790        _metrics.Update(summary);
 1791        return summary;
 1892    }
 93
 94    public async Task<ApplicationAuditReconciliationSummary> GetSummaryAsync(
 95        CancellationToken cancellationToken = default)
 96    {
 297        cancellationToken.ThrowIfCancellationRequested();
 298        if (!_options.Enabled)
 99        {
 1100            return ApplicationAuditReconciliationMetrics.DisabledSummary;
 101        }
 102
 1103        ApplicationAuditReconciliationSummary summary = await GetSummaryCoreAsync(
 1104                _metrics.LastRunUtc,
 1105                cancellationToken)
 1106            .ConfigureAwait(false);
 1107        _metrics.Update(summary);
 1108        return summary;
 2109    }
 110
 111    public async Task<IReadOnlyList<ApplicationAuditReconciliationFindingItem>> QueryFindingsAsync(
 112        ApplicationAuditReconciliationQuery request,
 113        CancellationToken cancellationToken = default)
 114    {
 1115        ArgumentNullException.ThrowIfNull(request);
 1116        cancellationToken.ThrowIfCancellationRequested();
 117
 1118        IQueryable<ApplicationAuditReconciliationFinding> query = _dbContext
 1119            .ApplicationAuditReconciliationFindings
 1120            .AsNoTracking();
 121
 1122        if (!string.IsNullOrWhiteSpace(request.ReasonCode))
 123        {
 1124            string reasonCode = request.ReasonCode.Trim();
 1125            query = query.Where(finding => finding.ReasonCode == reasonCode);
 126        }
 127
 1128        if (!string.IsNullOrWhiteSpace(request.Severity))
 129        {
 1130            string severity = request.Severity.Trim();
 1131            query = query.Where(finding => finding.Severity == severity);
 132        }
 133
 1134        if (!string.IsNullOrWhiteSpace(request.MutationBatchId))
 135        {
 1136            string batchId = request.MutationBatchId.Trim();
 1137            query = query.Where(finding => finding.MutationBatchId == batchId);
 138        }
 139
 1140        if (!string.IsNullOrWhiteSpace(request.RemediationStatus))
 141        {
 1142            string status = request.RemediationStatus.Trim();
 1143            query = query.Where(finding => finding.RemediationStatus == status);
 144        }
 145
 1146        int maximumResults = Math.Clamp(request.MaximumResults, 1, 500);
 1147        return await query
 1148            .OrderByDescending(finding => finding.LastObservedUtc)
 1149            .Take(maximumResults)
 1150            .Select(finding => new ApplicationAuditReconciliationFindingItem(
 1151                finding.Id,
 1152                finding.SchemaVersion,
 1153                finding.FindingKey,
 1154                finding.ReasonCode,
 1155                finding.Severity,
 1156                finding.MutationBatchId,
 1157                finding.Destination,
 1158                finding.Guidance,
 1159                finding.RemediationStatus,
 1160                finding.FirstObservedUtc,
 1161                finding.LastObservedUtc,
 1162                finding.ResolvedUtc))
 1163            .ToListAsync(cancellationToken)
 1164            .ConfigureAwait(false);
 1165    }
 166
 167    public async Task<ApplicationAuditReconciliationRemediationItem> RecordRemediationAsync(
 168        Guid findingId,
 169        ApplicationAuditReconciliationRemediationRequest request,
 170        CancellationToken cancellationToken = default)
 171    {
 5172        ArgumentNullException.ThrowIfNull(request);
 5173        cancellationToken.ThrowIfCancellationRequested();
 174
 5175        string actionCode = NormalizeRequired(request.ActionCode, 64, nameof(request.ActionCode));
 4176        string actorId = NormalizeRequired(request.ActorId, 256, nameof(request.ActorId));
 4177        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.
 4182        if (_dbContext.Database.CurrentTransaction is not null)
 183        {
 1184            return await RecordRemediationCoreAsync(
 1185                    findingId,
 1186                    actionCode,
 1187                    actorId,
 1188                    evidenceReference,
 1189                    request.ResolveFinding,
 1190                    cancellationToken)
 1191                .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.
 3196        IExecutionStrategy executionStrategy = _dbContext.Database.CreateExecutionStrategy();
 197
 3198        return await executionStrategy.ExecuteAsync(
 3199                async strategyCancellationToken =>
 3200                {
 3201                    await using IDbContextTransaction transaction = await _dbContext.Database
 3202                        .BeginTransactionAsync(strategyCancellationToken)
 3203                        .ConfigureAwait(false);
 3204
 3205                    ApplicationAuditReconciliationRemediationItem remediation = await RecordRemediationCoreAsync(
 3206                            findingId,
 3207                            actionCode,
 3208                            actorId,
 3209                            evidenceReference,
 3210                            request.ResolveFinding,
 3211                            strategyCancellationToken)
 3212                        .ConfigureAwait(false);
 3213
 3214                    await transaction.CommitAsync(strategyCancellationToken).ConfigureAwait(false);
 3215                    return remediation;
 3216                },
 3217                cancellationToken)
 3218            .ConfigureAwait(false);
 2219    }
 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    {
 4229        ApplicationAuditReconciliationFinding finding = await _dbContext
 4230            .ApplicationAuditReconciliationFindings
 4231            .AsNoTracking()
 4232            .SingleOrDefaultAsync(item => item.Id == findingId, cancellationToken)
 4233            .ConfigureAwait(false)
 4234            ?? throw new KeyNotFoundException($"Audit reconciliation finding '{findingId}' was not found.");
 235
 4236        DateTime now = UtcNow();
 4237        string remediationStatus = resolveFinding
 4238            ? ApplicationAuditReconciliationRemediationStatuses.Resolved
 4239            : ApplicationAuditReconciliationRemediationStatuses.Acknowledged;
 4240        DateTime? resolvedUtc = resolveFinding ? now : null;
 4241        string expectedFindingConcurrencyStamp = finding.ConcurrencyStamp;
 4242        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.
 4246        int updatedFindingCount = await _dbContext.ApplicationAuditReconciliationFindings
 4247            .Where(item => item.Id == findingId && item.ConcurrencyStamp == expectedFindingConcurrencyStamp)
 4248            .ExecuteUpdateAsync(
 4249                setters => setters
 4250                    .SetProperty(item => item.RemediationStatus, remediationStatus)
 4251                    .SetProperty(item => item.ResolvedUtc, resolvedUtc)
 4252                    .SetProperty(item => item.ConcurrencyStamp, nextFindingConcurrencyStamp),
 4253                cancellationToken)
 4254            .ConfigureAwait(false);
 255
 4256        if (updatedFindingCount != 1)
 257        {
 1258            throw new DbUpdateConcurrencyException(
 1259                $"Audit reconciliation finding '{findingId}' was modified after it was read. " +
 1260                "No remediation was recorded. Reload the finding, verify its current state, and retry.");
 261        }
 262
 3263        var remediationId = Guid.NewGuid();
 3264        string remediationConcurrencyStamp = Guid.NewGuid().ToString("N");
 265
 3266        _ = await InsertAsync<ApplicationAuditReconciliationRemediation>(
 3267                [
 3268                    (nameof(ApplicationAuditReconciliationRemediation.Id), remediationId),
 3269                    (nameof(ApplicationAuditReconciliationRemediation.FindingId), findingId),
 3270                    (nameof(ApplicationAuditReconciliationRemediation.MutationBatchId), finding.MutationBatchId),
 3271                    (nameof(ApplicationAuditReconciliationRemediation.ActionCode), actionCode),
 3272                    (nameof(ApplicationAuditReconciliationRemediation.ActorId), actorId),
 3273                    (nameof(ApplicationAuditReconciliationRemediation.EvidenceReference), evidenceReference),
 3274                    (nameof(ApplicationAuditReconciliationRemediation.RecordedUtc), now),
 3275                    (nameof(ApplicationAuditReconciliationRemediation.ConcurrencyStamp), remediationConcurrencyStamp)
 3276                ],
 3277                uniqueProperty: null,
 3278                cancellationToken)
 3279            .ConfigureAwait(false);
 280
 2281        return new(
 2282            remediationId,
 2283            findingId,
 2284            finding.MutationBatchId,
 2285            actionCode,
 2286            actorId,
 2287            evidenceReference,
 2288            now);
 2289    }
 290
 291    private List<ApplicationAuditReconciliationCandidate> BuildCandidates(
 292        IReadOnlyCollection<AuditRecord> auditRecords,
 293        IReadOnlyCollection<ApplicationAuditCompletionOutboxEntry> completionEntries,
 294        DateTime now)
 295    {
 19296        var candidates = new Dictionary<string, ApplicationAuditReconciliationCandidate>(StringComparer.Ordinal);
 19297        ILookup<string, AuditRecord> auditBatches = auditRecords
 19298            .Where(record => !string.IsNullOrWhiteSpace(record.MutationBatchId))
 19299            .ToLookup(record => record.MutationBatchId, StringComparer.Ordinal);
 19300        ILookup<string, ApplicationAuditCompletionOutboxEntry> completionBatches = completionEntries
 19301            .Where(entry => !string.IsNullOrWhiteSpace(entry.MutationBatchId))
 19302            .ToLookup(entry => entry.MutationBatchId, StringComparer.Ordinal);
 303
 42304        foreach (AuditRecord malformed in auditRecords.Where(record => string.IsNullOrWhiteSpace(record.MutationBatchId)
 305        {
 2306            Add(candidates, Candidate(
 2307                ApplicationAuditReconciliationReasonCodes.MalformedCorrelation,
 2308                ApplicationAuditReconciliationSeverities.Error,
 2309                $"missing-{malformed.Id:N}",
 2310                null,
 2311                "Retain the row, investigate the originating save path, and append remediation evidence."));
 312        }
 313
 76314        foreach (IGrouping<string, AuditRecord> batch in auditBatches)
 315        {
 19316            List<AuditRecord> records = [.. batch];
 19317            List<ApplicationAuditCompletionOutboxEntry> completions = [.. completionBatches[batch.Key]];
 19318            DateTime newestRecordUtc = records.Max(record => record.ModifiedOnUtc);
 319
 19320            if (completions.Count == 0 && newestRecordUtc <= now - _options.CompletionGracePeriod)
 321            {
 14322                Add(candidates, Candidate(
 14323                    ApplicationAuditReconciliationReasonCodes.MissingCompletion,
 14324                    ApplicationAuditReconciliationSeverities.Critical,
 14325                    batch.Key,
 14326                    null,
 14327                    "Preserve the audit batch, investigate transaction completion, and append an operator remediation re
 328            }
 329
 50330            foreach (ApplicationAuditCompletionOutboxEntry completion in completions)
 331            {
 6332                if (completion.AuditRecordCount != records.Count)
 333                {
 1334                    Add(candidates, Candidate(
 1335                        ApplicationAuditReconciliationReasonCodes.AuditRecordCountMismatch,
 1336                        ApplicationAuditReconciliationSeverities.Critical,
 1337                        batch.Key,
 1338                        completion.Destination,
 1339                        "Do not rewrite audit rows; compare retained evidence with the originating transaction and docum
 340                }
 341
 6342                ApplicationMutationAuditReceipt receipt = ToReceipt(completion);
 6343                if (!_manifestVerifier.Verify(receipt, records))
 344                {
 2345                    Add(candidates, Candidate(
 2346                        ApplicationAuditReconciliationReasonCodes.ManifestVerificationFailed,
 2347                        ApplicationAuditReconciliationSeverities.Critical,
 2348                        batch.Key,
 2349                        completion.Destination,
 2350                        "Quarantine downstream use of the batch, preserve all records, and investigate unauthorized or i
 351                }
 352            }
 353
 19354            if (records.Any(record => record.State == "Added" && string.IsNullOrWhiteSpace(record.KeyValues)))
 355            {
 0356                Add(candidates, Candidate(
 0357                    ApplicationAuditReconciliationReasonCodes.IncompleteGeneratedValues,
 0358                    ApplicationAuditReconciliationSeverities.Error,
 0359                    batch.Key,
 0360                    null,
 0361                    "Verify generated keys in the business database and append remediation evidence without modifying th
 362            }
 363
 19364            if (HasMalformedCorrelation(records))
 365            {
 0366                Add(candidates, Candidate(
 0367                    ApplicationAuditReconciliationReasonCodes.MalformedCorrelation,
 0368                    ApplicationAuditReconciliationSeverities.Warning,
 0369                    batch.Key,
 0370                    null,
 0371                    "Investigate inconsistent correlation metadata and preserve the original records as evidence."));
 372            }
 373        }
 374
 48375        foreach (IGrouping<string, ApplicationAuditCompletionOutboxEntry> batch in completionBatches)
 376        {
 5377            if (!auditBatches.Contains(batch.Key))
 378            {
 0379                foreach (ApplicationAuditCompletionOutboxEntry completion in batch)
 380                {
 0381                    Add(candidates, Candidate(
 0382                        ApplicationAuditReconciliationReasonCodes.MissingAuditBatch,
 0383                        ApplicationAuditReconciliationSeverities.Critical,
 0384                        batch.Key,
 0385                        completion.Destination,
 0386                        "Preserve the completion record and investigate missing or externally stored audit evidence."));
 387                }
 388            }
 389        }
 390
 29391        foreach (IGrouping<(string MutationBatchId, string Destination), ApplicationAuditCompletionOutboxEntry> duplicat
 19392            completionEntries.GroupBy(entry => (entry.MutationBatchId, entry.Destination)))
 393        {
 5394            if (duplicate.Count() > 1)
 395            {
 1396                Add(candidates, Candidate(
 1397                    ApplicationAuditReconciliationReasonCodes.DuplicateCompletion,
 1398                    ApplicationAuditReconciliationSeverities.Critical,
 1399                    duplicate.Key.MutationBatchId,
 1400                    duplicate.Key.Destination,
 1401                    "Preserve all records, stop dispatch for the destination, and investigate uniqueness or migration dr
 402            }
 403        }
 404
 50405        foreach (ApplicationAuditCompletionOutboxEntry entry in completionEntries)
 406        {
 6407            AddDeliveryCandidate(candidates, entry, now);
 408        }
 409
 19410        return [.. candidates.Values];
 411    }
 412
 413    private void AddDeliveryCandidate(
 414        IDictionary<string, ApplicationAuditReconciliationCandidate> candidates,
 415        ApplicationAuditCompletionOutboxEntry entry,
 416        DateTime now)
 417    {
 6418        if ((entry.Status == ApplicationAuditCompletionOutboxStatuses.Pending ||
 6419             entry.Status == ApplicationAuditCompletionOutboxStatuses.Deferred) &&
 6420            entry.CreatedUtc <= now - _options.StalePendingThreshold)
 421        {
 1422            Add(candidates, Candidate(
 1423                ApplicationAuditReconciliationReasonCodes.StalePending,
 1424                ApplicationAuditReconciliationSeverities.Warning,
 1425                entry.MutationBatchId,
 1426                entry.Destination,
 1427                "Verify dispatcher availability and destination registration before retrying delivery."));
 428        }
 5429        else if (entry.Status == ApplicationAuditCompletionOutboxStatuses.RetryableFailure &&
 5430                 entry.NextAttemptUtc <= now - _options.StaleRetryReadyThreshold)
 431        {
 0432            Add(candidates, Candidate(
 0433                ApplicationAuditReconciliationReasonCodes.StaleRetryReady,
 0434                ApplicationAuditReconciliationSeverities.Error,
 0435                entry.MutationBatchId,
 0436                entry.Destination,
 0437                "Inspect destination availability and retry policy; preserve prior attempt diagnostics."));
 438        }
 5439        else if (entry.Status == ApplicationAuditCompletionOutboxStatuses.Failed)
 440        {
 0441            Add(candidates, Candidate(
 0442                ApplicationAuditReconciliationReasonCodes.DeliveryFailed,
 0443                ApplicationAuditReconciliationSeverities.Error,
 0444                entry.MutationBatchId,
 0445                entry.Destination,
 0446                "Investigate the terminal delivery failure and append operator remediation evidence."));
 447        }
 5448        else if (entry.Status == ApplicationAuditCompletionOutboxStatuses.DeadLettered)
 449        {
 0450            Add(candidates, Candidate(
 0451                ApplicationAuditReconciliationReasonCodes.DeadLettered,
 0452                ApplicationAuditReconciliationSeverities.Critical,
 0453                entry.MutationBatchId,
 0454                entry.Destination,
 0455                "Review the dead letter, preserve diagnostics, correct the destination, and explicitly requeue only unde
 456        }
 5457    }
 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.
 19468        if (_dbContext.Database.CurrentTransaction is not null)
 469        {
 1470            await PersistCandidatesCoreAsync(candidates, auditBatchIds, completionEntries, now, cancellationToken)
 1471                .ConfigureAwait(false);
 1472            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.
 18478        IExecutionStrategy executionStrategy = _dbContext.Database.CreateExecutionStrategy();
 479
 18480        await executionStrategy.ExecuteAsync(
 18481                async strategyCancellationToken =>
 18482                {
 18483                    await using IDbContextTransaction transaction = await _dbContext.Database
 18484                        .BeginTransactionAsync(strategyCancellationToken)
 18485                        .ConfigureAwait(false);
 18486
 18487                    await PersistCandidatesCoreAsync(
 18488                            candidates,
 18489                            auditBatchIds,
 18490                            completionEntries,
 18491                            now,
 18492                            strategyCancellationToken)
 18493                        .ConfigureAwait(false);
 18494
 18495                    await transaction.CommitAsync(strategyCancellationToken).ConfigureAwait(false);
 18496                },
 18497                cancellationToken)
 18498            .ConfigureAwait(false);
 17499    }
 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    {
 19508        string[] keys = [.. candidates.Select(candidate => candidate.FindingKey)];
 19509        string[] scopeBatchIds = [.. auditBatchIds
 19510            .Concat(completionEntries.Select(entry => entry.MutationBatchId))
 19511            .Where(batchId => !string.IsNullOrWhiteSpace(batchId))
 19512            .Distinct(StringComparer.Ordinal)];
 513
 514        // Read the active and the resolvable findings in one round trip. Each row's stamp guards its update.
 19515        List<ApplicationAuditReconciliationFinding> existing = await _dbContext
 19516            .ApplicationAuditReconciliationFindings
 19517            .AsNoTracking()
 19518            .Where(finding => keys.Contains(finding.FindingKey) ||
 19519                (scopeBatchIds.Contains(finding.MutationBatchId) &&
 19520                    finding.RemediationStatus != ApplicationAuditReconciliationRemediationStatuses.Resolved))
 19521            .ToListAsync(cancellationToken)
 19522            .ConfigureAwait(false);
 19523        var existingByKey = existing
 19524            .ToDictionary(finding => finding.FindingKey, StringComparer.Ordinal);
 19525        var activeKeys = new HashSet<string>(keys, StringComparer.Ordinal);
 19526        var scopeBatchIdSet = new HashSet<string>(scopeBatchIds, StringComparer.Ordinal);
 527
 76528        foreach (ApplicationAuditReconciliationCandidate candidate in candidates)
 529        {
 20530            if (existingByKey.TryGetValue(candidate.FindingKey, out ApplicationAuditReconciliationFinding? finding))
 531            {
 2532                Guid findingId = finding.Id;
 2533                string expectedConcurrencyStamp = finding.ConcurrencyStamp;
 2534                string nextConcurrencyStamp = Guid.NewGuid().ToString("N");
 2535                int updatedCount = await _dbContext.ApplicationAuditReconciliationFindings
 2536                    .Where(item => item.Id == findingId && item.ConcurrencyStamp == expectedConcurrencyStamp)
 2537                    .ExecuteUpdateAsync(
 2538                        setters => setters
 2539                            .SetProperty(item => item.Severity, candidate.Severity)
 2540                            .SetProperty(item => item.Guidance, candidate.Guidance)
 2541                            .SetProperty(item => item.LastObservedUtc, now)
 2542                            .SetProperty(item => item.RemediationStatus, ApplicationAuditReconciliationRemediationStatus
 2543                            .SetProperty(item => item.ResolvedUtc, (DateTime?)null)
 2544                            .SetProperty(item => item.ConcurrencyStamp, nextConcurrencyStamp),
 2545                        cancellationToken)
 2546                    .ConfigureAwait(false);
 547
 2548                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.
 18554                int insertedCount = await InsertAsync<ApplicationAuditReconciliationFinding>(
 18555                        [
 18556                            (nameof(ApplicationAuditReconciliationFinding.Id), Guid.NewGuid()),
 18557                            (nameof(ApplicationAuditReconciliationFinding.SchemaVersion), ApplicationAuditReconciliation
 18558                            (nameof(ApplicationAuditReconciliationFinding.FindingKey), candidate.FindingKey),
 18559                            (nameof(ApplicationAuditReconciliationFinding.ReasonCode), candidate.ReasonCode),
 18560                            (nameof(ApplicationAuditReconciliationFinding.Severity), candidate.Severity),
 18561                            (nameof(ApplicationAuditReconciliationFinding.MutationBatchId), candidate.MutationBatchId),
 18562                            (nameof(ApplicationAuditReconciliationFinding.Destination), candidate.Destination),
 18563                            (nameof(ApplicationAuditReconciliationFinding.Guidance), candidate.Guidance),
 18564                            (nameof(ApplicationAuditReconciliationFinding.RemediationStatus), ApplicationAuditReconcilia
 18565                            (nameof(ApplicationAuditReconciliationFinding.FirstObservedUtc), now),
 18566                            (nameof(ApplicationAuditReconciliationFinding.LastObservedUtc), now),
 18567                            (nameof(ApplicationAuditReconciliationFinding.ResolvedUtc), null),
 18568                            (nameof(ApplicationAuditReconciliationFinding.ConcurrencyStamp), Guid.NewGuid().ToString("N"
 18569                        ],
 18570                        uniqueProperty: nameof(ApplicationAuditReconciliationFinding.FindingKey),
 18571                        cancellationToken)
 18572                    .ConfigureAwait(false);
 573
 18574                EnsureSingleRowWritten(insertedCount, candidate.FindingKey);
 575            }
 18576        }
 577
 34578        foreach (ApplicationAuditReconciliationFinding finding in existing.Where(finding =>
 17579                     !activeKeys.Contains(finding.FindingKey) &&
 17580                     scopeBatchIdSet.Contains(finding.MutationBatchId) &&
 17581                     finding.RemediationStatus != ApplicationAuditReconciliationRemediationStatuses.Resolved))
 582        {
 0583            Guid findingId = finding.Id;
 0584            string expectedConcurrencyStamp = finding.ConcurrencyStamp;
 0585            string nextConcurrencyStamp = Guid.NewGuid().ToString("N");
 0586            int resolvedCount = await _dbContext.ApplicationAuditReconciliationFindings
 0587                .Where(item => item.Id == findingId && item.ConcurrencyStamp == expectedConcurrencyStamp)
 0588                .ExecuteUpdateAsync(
 0589                    setters => setters
 0590                        .SetProperty(item => item.RemediationStatus, ApplicationAuditReconciliationRemediationStatuses.R
 0591                        .SetProperty(item => item.ResolvedUtc, (DateTime?)now)
 0592                        .SetProperty(item => item.ConcurrencyStamp, nextConcurrencyStamp),
 0593                    cancellationToken)
 0594                .ConfigureAwait(false);
 595
 0596            EnsureSingleRowWritten(resolvedCount, finding.FindingKey);
 0597        }
 17598    }
 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    {
 21611        IEntityType entityType = _dbContext.Model.FindEntityType(typeof(TEntity))
 21612            ?? throw new InvalidOperationException($"'{typeof(TEntity).Name}' is not part of the EF Core model.");
 21613        string tableName = entityType.GetTableName()
 21614            ?? throw new InvalidOperationException($"'{typeof(TEntity).Name}' is not mapped to a table.");
 21615        var table = StoreObjectIdentifier.Table(tableName, entityType.GetSchema());
 21616        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
 21626        string delimitedTable = EscapeFormatBraces(sqlGenerationHelper.DelimitIdentifier(table.Name, table.Schema));
 21627        var parameters = values.Select(value => value.Value).ToList();
 21628        StringBuilder sql = new StringBuilder()
 21629            .Append("INSERT INTO ").Append(delimitedTable)
 21630            .Append(" (").AppendJoin(", ", values.Select(value => Column(value.Property))).Append(") SELECT ")
 21631            .AppendJoin(", ", values.Select((_, index) => $"{{{index}}}"));
 632
 21633        if (uniqueProperty is not null)
 634        {
 18635            object? uniqueValue = values.Single(value => value.Property == uniqueProperty).Value;
 18636            _ = sql.Append(" WHERE NOT EXISTS (SELECT 1 FROM ").Append(delimitedTable)
 18637                .Append(" WHERE ").Append(Column(uniqueProperty)).Append(" = {").Append(parameters.Count).Append("})");
 18638            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.
 21643        return _dbContext.Database.ExecuteSqlAsync(
 21644            FormattableStringFactory.Create(sql.ToString(), [.. parameters]),
 21645            cancellationToken);
 646    }
 647
 648    private static string EscapeFormatBraces(string identifier)
 649    {
 297650        return identifier.Replace("{", "{{", StringComparison.Ordinal).Replace("}", "}}", StringComparison.Ordinal);
 651    }
 652
 653    private static void EnsureSingleRowWritten(int affectedCount, string findingKey)
 654    {
 20655        if (affectedCount != 1)
 656        {
 2657            throw new DbUpdateConcurrencyException(
 2658                $"Audit reconciliation finding '{findingKey}' was modified by another writer after it was read. " +
 2659                "No reconciliation changes were committed; the next reconciliation run re-evaluates the finding.");
 660        }
 18661    }
 662
 663    private async Task<ApplicationAuditReconciliationSummary> GetSummaryCoreAsync(
 664        DateTime? lastRunUtc,
 665        CancellationToken cancellationToken)
 666    {
 18667        IQueryable<ApplicationAuditReconciliationFinding> open = _dbContext
 18668            .ApplicationAuditReconciliationFindings
 18669            .AsNoTracking()
 18670            .Where(finding => finding.RemediationStatus != ApplicationAuditReconciliationRemediationStatuses.Resolved);
 671
 18672        long openCount = await open.LongCountAsync(cancellationToken).ConfigureAwait(false);
 18673        long errorCount = await open.LongCountAsync(
 18674            finding => finding.Severity == ApplicationAuditReconciliationSeverities.Error,
 18675            cancellationToken).ConfigureAwait(false);
 18676        long criticalCount = await open.LongCountAsync(
 18677            finding => finding.Severity == ApplicationAuditReconciliationSeverities.Critical,
 18678            cancellationToken).ConfigureAwait(false);
 18679        long manifestFailures = await open.LongCountAsync(
 18680            finding => finding.ReasonCode == ApplicationAuditReconciliationReasonCodes.ManifestVerificationFailed,
 18681            cancellationToken).ConfigureAwait(false);
 18682        long missingCompletion = await open.LongCountAsync(
 18683            finding => finding.ReasonCode == ApplicationAuditReconciliationReasonCodes.MissingCompletion,
 18684            cancellationToken).ConfigureAwait(false);
 18685        long staleDelivery = await open.LongCountAsync(
 18686            finding => finding.ReasonCode == ApplicationAuditReconciliationReasonCodes.StalePending ||
 18687                finding.ReasonCode == ApplicationAuditReconciliationReasonCodes.StaleRetryReady ||
 18688                finding.ReasonCode == ApplicationAuditReconciliationReasonCodes.DeliveryFailed,
 18689            cancellationToken).ConfigureAwait(false);
 18690        long deadLetters = await open.LongCountAsync(
 18691            finding => finding.ReasonCode == ApplicationAuditReconciliationReasonCodes.DeadLettered,
 18692            cancellationToken).ConfigureAwait(false);
 693
 18694        return new(
 18695            true,
 18696            lastRunUtc,
 18697            openCount,
 18698            errorCount,
 18699            criticalCount,
 18700            manifestFailures,
 18701            missingCompletion,
 18702            staleDelivery,
 18703            deadLetters);
 18704    }
 705
 706    private static bool HasMalformedCorrelation(IReadOnlyCollection<AuditRecord> records)
 707    {
 19708        return records.Any(record =>
 19709                IsWhitespaceOnly(record.OperationExecutionId) ||
 19710                IsWhitespaceOnly(record.ExecutionAttemptId) ||
 19711                IsWhitespaceOnly(record.DecisionAuditRecordId) ||
 19712                IsWhitespaceOnly(record.CorrelationId) ||
 19713                IsWhitespaceOnly(record.TraceId)) ||
 19714            HasMultipleValues(records.Select(record => record.OperationExecutionId)) ||
 19715            HasMultipleValues(records.Select(record => record.ExecutionAttemptId)) ||
 19716            HasMultipleValues(records.Select(record => record.DecisionAuditRecordId)) ||
 19717            HasMultipleValues(records.Select(record => record.CorrelationId)) ||
 19718            HasMultipleValues(records.Select(record => record.TraceId));
 719    }
 720
 721    private static bool HasMultipleValues(IEnumerable<string?> values)
 722    {
 95723        return values
 95724            .Where(value => !string.IsNullOrWhiteSpace(value))
 95725            .Select(value => value!.Trim())
 95726            .Distinct(StringComparer.Ordinal)
 95727            .Skip(1)
 95728            .Any();
 729    }
 730
 731    private static bool IsWhitespaceOnly(string? value)
 732    {
 95733        return value is not null && string.IsNullOrWhiteSpace(value);
 734    }
 735
 736    private static ApplicationMutationAuditReceipt ToReceipt(ApplicationAuditCompletionOutboxEntry entry)
 737    {
 6738        return new(
 6739            entry.MutationBatchId,
 6740            entry.AuditRecordCount,
 6741            entry.PersistenceOutcome,
 6742            new DateTimeOffset(DateTime.SpecifyKind(entry.ReceiptCompletedUtc, DateTimeKind.Utc)),
 6743            entry.MutationManifestHash,
 6744            entry.MutationManifestAlgorithm,
 6745            entry.MutationManifestSchemaVersion,
 6746            entry.OperationExecutionId,
 6747            entry.ExecutionAttemptId,
 6748            entry.DecisionAuditRecordId,
 6749            entry.CorrelationId,
 6750            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    {
 21760        string normalizedBatchId = NormalizeRequired(batchId, 64, nameof(batchId));
 21761        string? normalizedDestination = NormalizeOptional(destination, 128);
 21762        string identity = $"{reasonCode}\n{normalizedBatchId}\n{normalizedDestination}";
 21763        string key = Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(identity)));
 21764        return new(key, reasonCode, severity, normalizedBatchId, normalizedDestination, guidance);
 765    }
 766
 767    private static void Add(
 768        IDictionary<string, ApplicationAuditReconciliationCandidate> candidates,
 769        ApplicationAuditReconciliationCandidate candidate)
 770    {
 21771        candidates[candidate.FindingKey] = candidate;
 21772    }
 773
 774    private static string NormalizeRequired(string value, int maximumLength, string parameterName)
 775    {
 30776        ArgumentException.ThrowIfNullOrWhiteSpace(value, parameterName);
 29777        string normalized = value.Trim();
 29778        return normalized.Length <= maximumLength
 29779            ? normalized
 29780            : throw new ArgumentOutOfRangeException(parameterName, $"Value cannot exceed {maximumLength} characters.");
 781    }
 782
 783    private static string? NormalizeOptional(string? value, int maximumLength)
 784    {
 25785        if (string.IsNullOrWhiteSpace(value))
 786        {
 16787            return null;
 788        }
 789
 9790        string normalized = value.Trim();
 9791        return normalized.Length <= maximumLength
 9792            ? normalized
 9793            : normalized[..maximumLength];
 794    }
 795
 796    private DateTime UtcNow()
 797    {
 23798        return _timeProvider.GetUtcNow().UtcDateTime;
 799    }
 800}

Methods/Properties

.ctor(ProjectTemplate.Infrastructure.Data.ApplicationDbContext,ProjectTemplate.Infrastructure.Data.Auditing.IApplicationMutationManifestVerifier,Microsoft.Extensions.Options.IOptions`1<ProjectTemplate.Infrastructure.Data.Auditing.ApplicationAuditReconciliationOptions>,ProjectTemplate.Infrastructure.Data.Auditing.ApplicationAuditReconciliationMetrics,System.TimeProvider)
ReconcileAsync()
GetSummaryAsync()
QueryFindingsAsync()
RecordRemediationAsync()
RecordRemediationCoreAsync()
BuildCandidates(System.Collections.Generic.IReadOnlyCollection`1<ProjectTemplate.Infrastructure.Data.Entities.AuditRecord>,System.Collections.Generic.IReadOnlyCollection`1<ProjectTemplate.Infrastructure.Data.Entities.ApplicationAuditCompletionOutboxEntry>,System.DateTime)
AddDeliveryCandidate(System.Collections.Generic.IDictionary`2<System.String,ProjectTemplate.Infrastructure.Data.Auditing.ApplicationAuditReconciliationCandidate>,ProjectTemplate.Infrastructure.Data.Entities.ApplicationAuditCompletionOutboxEntry,System.DateTime)
PersistCandidatesAsync()
PersistCandidatesCoreAsync()
InsertAsync(System.Collections.Generic.IReadOnlyList`1<System.ValueTuple`2<System.String,System.Object>>,System.String,System.Threading.CancellationToken)
EscapeFormatBraces(System.String)
EnsureSingleRowWritten(System.Int32,System.String)
GetSummaryCoreAsync()
HasMalformedCorrelation(System.Collections.Generic.IReadOnlyCollection`1<ProjectTemplate.Infrastructure.Data.Entities.AuditRecord>)
HasMultipleValues(System.Collections.Generic.IEnumerable`1<System.String>)
IsWhitespaceOnly(System.String)
ToReceipt(ProjectTemplate.Infrastructure.Data.Entities.ApplicationAuditCompletionOutboxEntry)
Candidate(System.String,System.String,System.String,System.String,System.String)
Add(System.Collections.Generic.IDictionary`2<System.String,ProjectTemplate.Infrastructure.Data.Auditing.ApplicationAuditReconciliationCandidate>,ProjectTemplate.Infrastructure.Data.Auditing.ApplicationAuditReconciliationCandidate)
NormalizeRequired(System.String,System.Int32,System.String)
NormalizeOptional(System.String,System.Int32)
UtcNow()