From de03132412799e3feac30d3f5c64f5187cfeb112 Mon Sep 17 00:00:00 2001 From: Warwick Schroeder Date: Wed, 23 Sep 2026 12:29:04 +0800 Subject: [PATCH] Add retention sweep outcome tracking and reporting --- ...hen_triggering_a_manual_retention_purge.cs | 1 + .../Contracts/RetentionPurgeOutcome.cs | 8 + .../Contracts/RetentionPurgeStatusResponse.cs | 11 + .../Infrastructure/RetentionSweeper.cs | 146 +++++---- .../EFCore/InMemoryBodyStoragePersistence.cs | 8 +- .../EFCore/RetentionSweepTests.cs | 294 +++++++++++++++++- .../RetentionSweepCurrentStatus.cs | 9 +- .../RetentionSweepOutcome.cs | 20 ++ .../Infrastructure/RetentionApiTests.cs | 49 +++ .../Infrastructure/Api/RetentionApi.cs | 11 +- 10 files changed, 486 insertions(+), 71 deletions(-) create mode 100644 src/ServiceControl.Api/Contracts/RetentionPurgeOutcome.cs create mode 100644 src/ServiceControl.Persistence/RetentionSweepOutcome.cs create mode 100644 src/ServiceControl.UnitTests/Infrastructure/RetentionApiTests.cs diff --git a/src/ServiceControl.AcceptanceTests/WebApi/When_triggering_a_manual_retention_purge.cs b/src/ServiceControl.AcceptanceTests/WebApi/When_triggering_a_manual_retention_purge.cs index da514f32db..91660ca073 100644 --- a/src/ServiceControl.AcceptanceTests/WebApi/When_triggering_a_manual_retention_purge.cs +++ b/src/ServiceControl.AcceptanceTests/WebApi/When_triggering_a_manual_retention_purge.cs @@ -94,6 +94,7 @@ await Define() Assert.That(completion, Is.Not.Null, "the background purge must complete"); Assert.That(completion.IsRunning, Is.False); Assert.That(completion.LastFinishedAt, Is.Not.Null, "a completed run records its finish time"); + Assert.That(completion.LastOutcome, Is.EqualTo(RetentionPurgeOutcome.Succeeded), "a run with no failed pass reports success"); } } diff --git a/src/ServiceControl.Api/Contracts/RetentionPurgeOutcome.cs b/src/ServiceControl.Api/Contracts/RetentionPurgeOutcome.cs new file mode 100644 index 0000000000..366daea99b --- /dev/null +++ b/src/ServiceControl.Api/Contracts/RetentionPurgeOutcome.cs @@ -0,0 +1,8 @@ +namespace ServiceControl.Api.Contracts; + +public enum RetentionPurgeOutcome +{ + Succeeded, + Failed, + Cancelled +} \ No newline at end of file diff --git a/src/ServiceControl.Api/Contracts/RetentionPurgeStatusResponse.cs b/src/ServiceControl.Api/Contracts/RetentionPurgeStatusResponse.cs index 716c666fdf..bbec15cb5f 100644 --- a/src/ServiceControl.Api/Contracts/RetentionPurgeStatusResponse.cs +++ b/src/ServiceControl.Api/Contracts/RetentionPurgeStatusResponse.cs @@ -18,6 +18,17 @@ public class RetentionPurgeStatusResponse public DateTime? LastEventsCutoff { get; set; } + /// + /// How the most recent purge ended. Null while it is still running or before any has run. + /// + public RetentionPurgeOutcome? LastOutcome { get; set; } + + /// + /// What went wrong in each pass that failed, or null when none did. + /// A cancelled purge keeps the errors from before the cancellation. The full detail is in the log. + /// + public string LastError { get; set; } + /// Present only on the 501 Not Implemented response. public string Reason { get; set; } } \ No newline at end of file diff --git a/src/ServiceControl.Persistence.EFCore/Infrastructure/RetentionSweeper.cs b/src/ServiceControl.Persistence.EFCore/Infrastructure/RetentionSweeper.cs index cdda882ea1..ad5082de5c 100644 --- a/src/ServiceControl.Persistence.EFCore/Infrastructure/RetentionSweeper.cs +++ b/src/ServiceControl.Persistence.EFCore/Infrastructure/RetentionSweeper.cs @@ -36,14 +36,9 @@ class RetentionSweeper( // never overlap. Precedent: ExternalIntegrationRequestsDataStore.drainLock. readonly SemaphoreSlim sweepLock = new(1, 1); - // Status snapshot for GET /api/maintenance/retention/purge/status polling. Volatile reads/writes are - // sufficient here: the fields are written under sweepLock (or once at start) and read - // lock-free for status reporting, which only needs an eventually-consistent snapshot. - volatile bool isRunning; - DateTime? lastStartedAt; - DateTime? lastFinishedAt; - DateTime? lastErrorCutoff; - DateTime? lastEventsCutoff; + // Read without the lock by status polling, so it is replaced whole rather than edited and a poll + // never sees half of one run and half of another. + volatile RetentionSweepCurrentStatus status = new(false, null, null, null, null, null, null); protected override async Task ExecuteAsync(CancellationToken cancellationToken = default) { @@ -61,19 +56,18 @@ protected override async Task ExecuteAsync(CancellationToken cancellationToken = { await Sweep(errorCutoff: null, eventsCutoff: null, pace: true, cancellationToken); } -#pragma warning disable PS0019 // The filter already excludes OperationCanceledException, so - // cancellation propagates; PS0019 only recognises a cancellationToken guard. - catch (Exception ex) when (ex is not OperationCanceledException) +#pragma warning disable PS0019 // Filtered on the token alone because SqlClient reports a cancelled command as a SqlException. + catch (Exception ex) when (!cancellationToken.IsCancellationRequested) { logger.LogError(ex, "Error during retention sweep"); } -#pragma warning restore PS0019 } while (await timer.WaitForNextTickAsync(cancellationToken)); } - catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) + catch (Exception) when (cancellationToken.IsCancellationRequested) { logger.LogInformation("Stopping retention sweep"); } +#pragma warning restore PS0019 } // Runs a full sweep immediately, bypassing the timer and the inter-batch pause. @@ -88,20 +82,16 @@ public ManualSweepAttempt TryStartManualSweep(DateTime? errorCutoff, DateTime? e // already running (holding the lock) if (!sweepLock.Wait(0, cancellationToken)) { - return new ManualSweepAttempt(RetentionSweepStatus.AlreadyRunning, lastStartedAt, errorCutoff, eventsCutoff); + return new ManualSweepAttempt(RetentionSweepStatus.AlreadyRunning, status.LastStartedAt, errorCutoff, eventsCutoff); } - // Lock acquired on this thread. The background task owns it from here and releases it when - // the sweep body completes (SemaphoreSlim is not thread-affine, so releasing from the - // background thread is safe). isRunning is set now so a concurrent manual call sees it. - isRunning = true; - lastStartedAt = timeProvider.GetUtcNow().UtcDateTime; - lastErrorCutoff = errorCutoff; - lastEventsCutoff = eventsCutoff; + // The background task releases the lock when the sweep ends, which SemaphoreSlim allows from + // another thread. The run is marked started first so a status poll sees it straight away. + var running = RecordStart(errorCutoff, eventsCutoff); _ = SweepWithoutAcquiringLock(); - return new ManualSweepAttempt(RetentionSweepStatus.Started, lastStartedAt, errorCutoff, eventsCutoff); + return new ManualSweepAttempt(RetentionSweepStatus.Started, running.LastStartedAt, errorCutoff, eventsCutoff); async Task SweepWithoutAcquiringLock() { @@ -109,59 +99,91 @@ async Task SweepWithoutAcquiringLock() try { // if the caller doesn't hand over a real cancellation token then use the application lifetime. - await SweepBody(errorCutoff, eventsCutoff, false, cancellation.Token); - lastFinishedAt = timeProvider.GetUtcNow().UtcDateTime; + await SweepBody(running, errorCutoff, eventsCutoff, false, cancellation.Token); } - catch (OperationCanceledException) when (cancellation.Token.IsCancellationRequested) +#pragma warning disable PS0019 // Filtered on the token alone because SqlClient reports a cancelled command as a SqlException. + catch (Exception) when (cancellation.Token.IsCancellationRequested) { - //smother this exception, cancelling + // Cancelled by the caller or by shutdown, which is not an error. The status records it. } catch (Exception e) { logger.LogError(e, "Error during retention sweep"); } - finally - { - isRunning = false; - sweepLock.Release(); - } +#pragma warning restore PS0019 } } - public RetentionSweepCurrentStatus GetStatus() => new(isRunning, lastStartedAt, lastFinishedAt, lastErrorCutoff, lastEventsCutoff); + public RetentionSweepCurrentStatus GetStatus() => status; async Task Sweep(DateTime? errorCutoff, DateTime? eventsCutoff, bool pace, CancellationToken cancellationToken) { await sweepLock.WaitAsync(cancellationToken); - isRunning = true; - lastStartedAt = timeProvider.GetUtcNow().UtcDateTime; - lastErrorCutoff = errorCutoff; - lastEventsCutoff = eventsCutoff; + var running = RecordStart(errorCutoff, eventsCutoff); + await SweepBody(running, errorCutoff, eventsCutoff, pace, cancellationToken); + } + + // A new snapshot leaves out the previous run's finish time and outcome so they are never read as this run's. + RetentionSweepCurrentStatus RecordStart(DateTime? errorCutoff, DateTime? eventsCutoff) + { + var running = new RetentionSweepCurrentStatus(true, timeProvider.GetUtcNow().UtcDateTime, null, errorCutoff, eventsCutoff, null, null); + status = running; + return running; + } + + // The lock is released before the run is shown as finished, so a caller who sees it finished can start + // the next one. The compare-exchange keeps a run that has already started in that gap from being overwritten. + void Finish(RetentionSweepCurrentStatus running, RetentionSweepOutcome outcome, List passErrors) + { + var errors = passErrors.OfType().ToList(); + + var finished = running with + { + IsRunning = false, + LastFinishedAt = timeProvider.GetUtcNow().UtcDateTime, + LastOutcome = outcome, + LastError = errors.Count == 0 ? null : string.Join("; ", errors) + }; + + sweepLock.Release(); + Interlocked.CompareExchange(ref status, finished, running); + } + + // The three sub-sweeps, shared by the hourly Sweep path and the manual background path. The caller + // holds the lock, and every way out of here releases it. + async Task SweepBody(RetentionSweepCurrentStatus running, DateTime? errorCutoff, DateTime? eventsCutoff, bool pace, CancellationToken cancellationToken) + { + List passErrors = []; +#pragma warning disable PS0021 // The pass lambdas' token is the cancellationToken that RunPass hands back, so there is only one token. +#pragma warning disable PS0019 // Filtered on the token alone because SqlClient reports a cancelled command as a SqlException. try { - await SweepBody(errorCutoff, eventsCutoff, pace, cancellationToken); - lastFinishedAt = timeProvider.GetUtcNow().UtcDateTime; + passErrors.Add(await RunPass(RetentionEntity.FailedMessages, token => SweepFailedMessages(pace, errorCutoff, token), cancellationToken)); + passErrors.Add(await RunPass(RetentionEntity.EventLog, token => SweepEventLogItems(pace, eventsCutoff, token), cancellationToken)); + passErrors.Add(await RunPass(RetentionEntity.GroupComments, SweepOrphanedGroupComments, cancellationToken)); + retentionState.SweepComplete(); + + Finish(running, passErrors.Any(error => error is not null) ? RetentionSweepOutcome.Failed : RetentionSweepOutcome.Succeeded, passErrors); } - finally + catch (Exception) when (cancellationToken.IsCancellationRequested) { - isRunning = false; - sweepLock.Release(); + Finish(running, RetentionSweepOutcome.Cancelled, passErrors); + throw; } - } - - // The three sub-sweeps, isolated from lock management so both the locked Sweep path and the - // manual background path (which already holds the lock) share one implementation. - async Task SweepBody(DateTime? errorCutoff, DateTime? eventsCutoff, bool pace, CancellationToken cancellationToken) - { - await RunPass(RetentionEntity.FailedMessages, token => SweepFailedMessages(pace, errorCutoff, token), cancellationToken); - await RunPass(RetentionEntity.EventLog, token => SweepEventLogItems(pace, eventsCutoff, token), cancellationToken); - await RunPass(RetentionEntity.GroupComments, SweepOrphanedGroupComments, cancellationToken); - retentionState.SweepComplete(); + catch (Exception ex) + { + passErrors.Add(FirstLine(ex.Message)); + Finish(running, RetentionSweepOutcome.Failed, passErrors); + throw; + } +#pragma warning restore PS0019 +#pragma warning restore PS0021 } // Each pass is isolated so one failing kind of row does not stop the others from being // reclaimed, and so the metrics report an outcome for every pass on every run. - async Task RunPass(RetentionEntity entity, Func pass, CancellationToken cancellationToken) + // Returns what went wrong for the status to report, or null when the pass succeeded. + async Task RunPass(RetentionEntity entity, Func pass, CancellationToken cancellationToken) { using var cycle = metrics.BeginCycle(entity, cancellationToken); @@ -170,17 +192,21 @@ async Task RunPass(RetentionEntity entity, Func pass, C await pass(cancellationToken); cycle.Complete(); retentionState.Clear(entity); + return null; } -#pragma warning disable PS0019 // The filter already excludes OperationCanceledException, so - // cancellation propagates; PS0019 only recognises a cancellationToken guard. - catch (Exception ex) when (ex is not OperationCanceledException) +#pragma warning disable PS0019 // Filtered on the token alone because SqlClient reports a cancelled command as a SqlException. + catch (Exception ex) when (!cancellationToken.IsCancellationRequested) { logger.LogError(ex, "Error during the {RetentionEntity} retention pass", entity); retentionState.ReportError(entity, ex.Message); + return $"{entity}: {FirstLine(ex.Message)}"; } #pragma warning restore PS0019 } + // Storage SDK messages can carry many lines of response detail, and the full text is already in the log. + static string FirstLine(string message) => message.Split('\n', 2)[0].TrimEnd('\r'); + // Once the last message of a group has been swept the group cannot be displayed at all, so its // comment is unreachable. Leaving it behind would both accumulate invisible rows and, because // group ids are deterministic, reattach a stale comment if the same failure ever recurs. @@ -202,8 +228,11 @@ async Task SweepEventLogItems(bool pace, DateTime? eventsCutoff, CancellationTok { var cutoff = eventsCutoff ?? (timeProvider.GetUtcNow().UtcDateTime - settings.EventsRetentionPeriod); - while (!cancellationToken.IsCancellationRequested) + while (true) { + // Throws rather than leaving the loop, so a pass cut short between batches is not reported as finished. + cancellationToken.ThrowIfCancellationRequested(); + using var scope = serviceScopeFactory.CreateScope(); var dbContext = scope.ServiceProvider.GetRequiredService(); @@ -231,8 +260,11 @@ async Task SweepFailedMessages(bool pace, DateTime? errorCutoff, CancellationTok { var cutoff = errorCutoff ?? (timeProvider.GetUtcNow().UtcDateTime - settings.ErrorRetentionPeriod); - while (!cancellationToken.IsCancellationRequested) + while (true) { + // Throws rather than leaving the loop, so a pass cut short between batches is not reported as finished. + cancellationToken.ThrowIfCancellationRequested(); + using var scope = serviceScopeFactory.CreateScope(); var dbContext = scope.ServiceProvider.GetRequiredService(); diff --git a/src/ServiceControl.Persistence.Tests/EFCore/InMemoryBodyStoragePersistence.cs b/src/ServiceControl.Persistence.Tests/EFCore/InMemoryBodyStoragePersistence.cs index 4c0d30b2c2..e03d5c130b 100644 --- a/src/ServiceControl.Persistence.Tests/EFCore/InMemoryBodyStoragePersistence.cs +++ b/src/ServiceControl.Persistence.Tests/EFCore/InMemoryBodyStoragePersistence.cs @@ -16,6 +16,11 @@ class InMemoryBodyStoragePersistence : IBodyStoragePersistence public HashSet FailDeleteFor { get; } = []; + /// + /// Extra lines added to the simulated delete failure, the way a storage SDK appends response detail. + /// + public string DeleteFailureDetail { get; set; } + public IReadOnlyList Written { get @@ -83,7 +88,8 @@ public Task DeleteBodyIfExists(string bodyId, CancellationToken cancellationToke { if (FailDeleteFor.Contains(bodyId)) { - throw new InvalidOperationException($"Simulated body storage failure for {bodyId}"); + var detail = DeleteFailureDetail is null ? "" : Environment.NewLine + DeleteFailureDetail; + throw new InvalidOperationException($"Simulated body storage failure for {bodyId}{detail}"); } lock (gate) diff --git a/src/ServiceControl.Persistence.Tests/EFCore/RetentionSweepTests.cs b/src/ServiceControl.Persistence.Tests/EFCore/RetentionSweepTests.cs index 7d7fbcdb91..3a73a8350d 100644 --- a/src/ServiceControl.Persistence.Tests/EFCore/RetentionSweepTests.cs +++ b/src/ServiceControl.Persistence.Tests/EFCore/RetentionSweepTests.cs @@ -2,9 +2,13 @@ namespace ServiceControl.Persistence.Tests; using System; using System.Collections.Generic; +using System.Data.Common; using System.Diagnostics.Metrics; using System.Linq; +using System.Threading; using System.Threading.Tasks; +using Microsoft.EntityFrameworkCore; +using Microsoft.EntityFrameworkCore.Diagnostics; using Microsoft.Extensions.DependencyInjection; using NUnit.Framework; using ServiceControl.EventLog; @@ -16,6 +20,20 @@ namespace ServiceControl.Persistence.Tests; class RetentionSweepTests : ErrorIngestionTestBase { + readonly CancelAfterBatchDelete cancelAfterBatchDelete = new(); + readonly CancelDuringNextQuery cancelDuringNextQuery = new(); + + public RetentionSweepTests() + { + var registerBodyStorage = RegisterServices; + RegisterServices = services => + { + registerBodyStorage(services); + PersistenceTestsContext.InterceptDatabaseCommands(services, cancelAfterBatchDelete); + PersistenceTestsContext.InterceptDatabaseCommands(services, cancelDuringNextQuery); + }; + } + [SetUp] public void SetRetention() => EFSettings.ErrorRetentionPeriod = TimeSpan.FromDays(30); @@ -547,6 +565,125 @@ public async Task Manual_sweep_runs_in_the_background_and_reports_status() Assert.That(status.LastStartedAt, Is.Not.Null); Assert.That(status.LastFinishedAt, Is.Not.Null); Assert.That(status.LastErrorCutoff, Is.Not.Null); + Assert.That(status.LastOutcome, Is.EqualTo(RetentionSweepOutcome.Succeeded)); + Assert.That(status.LastError, Is.Null); + } + } + + [Test] + public async Task Manual_sweep_with_failing_passes_reports_failed_with_every_error() + { + // Subtracting these from the clock cannot be represented, so both passes throw. + EFSettings.ErrorRetentionPeriod = TimeSpan.FromDays(1_000_000); + EFSettings.EventsRetentionPeriod = TimeSpan.FromDays(1_000_000); + + var sweeper = GetSweeper(); + sweeper.TryStartManualSweep(null, null); + + await WaitForManualSweepToFinish(); + + var status = sweeper.GetStatus(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(status.LastOutcome, Is.EqualTo(RetentionSweepOutcome.Failed)); + Assert.That(status.LastError, Does.Contain(nameof(RetentionEntity.FailedMessages))); + Assert.That(status.LastError, Does.Contain(nameof(RetentionEntity.EventLog))); + Assert.That(status.LastFinishedAt, Is.Not.Null); + } + } + + [Test] + public async Task Scheduled_sweep_with_a_failing_pass_reports_failed() + { + EFSettings.ErrorRetentionPeriod = TimeSpan.FromDays(1_000_000); + + await RunRetentionSweep(); + + var status = GetSweeper().GetStatus(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(status.LastOutcome, Is.EqualTo(RetentionSweepOutcome.Failed)); + Assert.That(status.LastFinishedAt, Is.Not.Null); + } + } + + [Test] + public async Task A_cancelled_sweep_keeps_the_errors_of_passes_that_failed_first() + { + EFSettings.ErrorRetentionPeriod = TimeSpan.FromDays(1_000_000); + + using var cancellation = new CancellationTokenSource(); + var sweeper = GetSweeper(); + + // The failed messages pass throws before any database call, so by the time the start returns + // the event log pass is waiting on the database and the cancellation lands there. + sweeper.TryStartManualSweep(null, null, cancellation.Token); + await cancellation.CancelAsync(); + + await WaitForManualSweepToFinish(); + + var status = sweeper.GetStatus(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(status.LastOutcome, Is.EqualTo(RetentionSweepOutcome.Cancelled)); + Assert.That(status.LastError, Does.Contain(nameof(RetentionEntity.FailedMessages))); + } + } + + [Test] + public async Task Manual_sweep_cancelled_part_way_reports_cancelled() + { + using var cancellation = new CancellationTokenSource(); + var sweeper = GetSweeper(); + + // The background sweep runs on this thread until its first database call, so it is still + // running when the token is cancelled. + sweeper.TryStartManualSweep(null, null, cancellation.Token); + await cancellation.CancelAsync(); + + await WaitForManualSweepToFinish(); + + var status = sweeper.GetStatus(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(status.LastOutcome, Is.EqualTo(RetentionSweepOutcome.Cancelled)); + Assert.That(status.LastError, Is.Null); + Assert.That(status.LastFinishedAt, Is.Not.Null); + } + } + + [Test] + public async Task Starting_a_sweep_clears_the_previous_run_s_outcome() + { + EFSettings.ErrorRetentionPeriod = TimeSpan.FromDays(1_000_000); + + var sweeper = GetSweeper(); + sweeper.TryStartManualSweep(null, null); + await WaitForManualSweepToFinish(); + + var previous = sweeper.GetStatus(); + + EFSettings.ErrorRetentionPeriod = TimeSpan.FromDays(30); + + var second = sweeper.TryStartManualSweep(null, null); + var running = sweeper.GetStatus(); + + await WaitForManualSweepToFinish(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(second.Outcome, Is.EqualTo(RetentionSweepStatus.Started)); + Assert.That(previous.LastOutcome, Is.EqualTo(RetentionSweepOutcome.Failed)); + Assert.That(previous.LastError, Is.Not.Null); + Assert.That(previous.LastFinishedAt, Is.Not.Null); + Assert.That(running.IsRunning, Is.True); + Assert.That(running.LastOutcome, Is.Null); + Assert.That(running.LastError, Is.Null); + Assert.That(running.LastFinishedAt, Is.Null); } } @@ -556,6 +693,116 @@ public async Task A_second_manual_sweep_is_refused_while_one_is_running() // Seed enough rows to force multiple delete batches so the first sweep is still running when the // second, synchronous call is made. The single-flight lock is held from the moment the first call // returns Started until the background body completes. + await SeedMoreThanOneBatchOfExpiredMessages(); + + var sweeper = GetSweeper(); + var first = sweeper.TryStartManualSweep(Now.AddDays(-30), null); + // Immediately request a second sweep on the same thread while the first is still deleting. + var second = sweeper.TryStartManualSweep(Now.AddDays(-30), null); + + await WaitForManualSweepToFinish(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(first.Outcome, Is.EqualTo(RetentionSweepStatus.Started), + "the first call should start the sweep"); + Assert.That(second.Outcome, Is.EqualTo(RetentionSweepStatus.AlreadyRunning), + "a second sweep must not run in parallel with the first"); + } + } + + [Test] + public async Task A_sweep_cancelled_between_batches_is_recorded_as_cancelled() + { + await SeedMoreThanOneBatchOfExpiredMessages(); + + using var cancellation = new CancellationTokenSource(); + cancelAfterBatchDelete.Cancellation = cancellation; + using var recorded = ListenToRetentionMetrics(); + + var sweeper = GetSweeper(); + sweeper.TryStartManualSweep(Now.AddDays(-30), null, cancellation.Token); + + await WaitForManualSweepToFinish(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(recorded.Cycles(RetentionEntity.FailedMessages).Select(cycle => cycle.Result), Is.EqualTo(new[] { "cancelled" })); + Assert.That(recorded.Cycles(RetentionEntity.EventLog), Is.Empty, "the sweep must stop at the pass that was cut short"); + Assert.That(recorded.Cycles(RetentionEntity.GroupComments), Is.Empty, "the sweep must stop at the pass that was cut short"); + Assert.That(sweeper.GetStatus().LastOutcome, Is.EqualTo(RetentionSweepOutcome.Cancelled)); + Assert.That(await Query(dbContext => dbContext.FailedMessages.CountAsync()), Is.EqualTo(500), "exactly one batch is deleted before the cancellation"); + } + } + + [Test] + public async Task An_event_log_pass_cancelled_between_batches_is_recorded_as_cancelled() + { + EFSettings.EventsRetentionPeriod = TimeSpan.FromDays(14); + await Store([.. Enumerable.Range(0, 1500).Select(i => EventLogRow($"expired-{i}", Now.AddDays(-15)))]); + + // There are no expired failed messages, so the first delete the sweep runs is an event log batch. + using var cancellation = new CancellationTokenSource(); + cancelAfterBatchDelete.Cancellation = cancellation; + using var recorded = ListenToRetentionMetrics(); + + var sweeper = GetSweeper(); + sweeper.TryStartManualSweep(null, null, cancellation.Token); + + await WaitForManualSweepToFinish(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(recorded.Cycles(RetentionEntity.EventLog).Select(cycle => cycle.Result), Is.EqualTo(new[] { "cancelled" })); + Assert.That(recorded.Cycles(RetentionEntity.GroupComments), Is.Empty, "the sweep must stop at the pass that was cut short"); + Assert.That(sweeper.GetStatus().LastOutcome, Is.EqualTo(RetentionSweepOutcome.Cancelled)); + Assert.That(await Query(dbContext => dbContext.EventLogItems.CountAsync()), Is.EqualTo(500), "exactly one batch is deleted before the cancellation"); + } + } + + [TearDown] + public void DisarmInterceptors() + { + cancelAfterBatchDelete.Cancellation = null; + cancelDuringNextQuery.Arm(null, null); + } + + [Test] + public async Task A_sweep_cancelled_while_a_query_runs_reports_cancelled_without_an_error() + { + await SeedFailedMessage(FailedMessageStatus.Resolved, Now.AddDays(-31)); + + using var cancellation = new CancellationTokenSource(); + cancelDuringNextQuery.Arm(cancellation, PersistenceTestsContext.SqlToDelayFor(TimeSpan.FromSeconds(20))); + + var sweeper = GetSweeper(); + sweeper.TryStartManualSweep(null, null, cancellation.Token); + + await WaitForManualSweepToFinish(); + + var status = sweeper.GetStatus(); + + using (Assert.EnterMultipleScope()) + { + Assert.That(status.LastOutcome, Is.EqualTo(RetentionSweepOutcome.Cancelled)); + Assert.That(status.LastError, Is.Null, "a cancelled query is not a failed pass"); + } + } + + [Test] + public async Task A_failed_pass_reports_only_the_first_line_of_its_error() + { + var unluckyBody = await SeedFailedMessage(FailedMessageStatus.Resolved, Now.AddDays(-31), bodyStoredExternally: true); + RecordedBodies.FailDeleteFor.Add(unluckyBody.ToString()); + RecordedBodies.DeleteFailureDetail = "AuthorizationFailure"; + + await RunRetentionSweep(); + + Assert.That(GetSweeper().GetStatus().LastError, Is.EqualTo($"FailedMessages: Simulated body storage failure for {unluckyBody}")); + } + + async Task SeedMoreThanOneBatchOfExpiredMessages() + { var rows = new List(); for (var i = 0; i < 1500; i++) { @@ -578,20 +825,47 @@ public async Task A_second_manual_sweep_is_refused_while_one_is_running() } await Store([.. rows]); + } - var sweeper = GetSweeper(); - var first = sweeper.TryStartManualSweep(Now.AddDays(-30), null); - // Immediately request a second sweep on the same thread while the first is still deleting. - var second = sweeper.TryStartManualSweep(Now.AddDays(-30), null); + // Slows the next query down on the server and cancels while it runs, where SqlClient reports the + // cancellation as a SqlException rather than an OperationCanceledException. + class CancelDuringNextQuery : DbCommandInterceptor + { + CancellationTokenSource cancellation; + string delaySql; - await WaitForManualSweepToFinish(); + public void Arm(CancellationTokenSource cancellation, string delaySql) + { + this.cancellation = cancellation; + this.delaySql = delaySql; + } - using (Assert.EnterMultipleScope()) + public override ValueTask> ReaderExecutingAsync(DbCommand command, CommandEventData eventData, InterceptionResult result, CancellationToken cancellationToken = default) { - Assert.That(first.Outcome, Is.EqualTo(RetentionSweepStatus.Started), - "the first call should start the sweep"); - Assert.That(second.Outcome, Is.EqualTo(RetentionSweepStatus.AlreadyRunning), - "a second sweep must not run in parallel with the first"); + if (Interlocked.Exchange(ref cancellation, null) is { } armed) + { + command.CommandText = delaySql + Environment.NewLine + command.CommandText; + armed.CancelAfter(TimeSpan.FromMilliseconds(500)); + } + + return base.ReaderExecutingAsync(command, eventData, result, cancellationToken); + } + } + + // Cancels right after a delete command finishes, which is when a batch loop checks the token itself + // rather than inside an awaited database call. + class CancelAfterBatchDelete : DbCommandInterceptor + { + public CancellationTokenSource Cancellation { get; set; } + + public override async ValueTask NonQueryExecutedAsync(DbCommand command, CommandExecutedEventData eventData, int result, CancellationToken cancellationToken = default) + { + if (Cancellation is { IsCancellationRequested: false } armed) + { + await armed.CancelAsync(); + } + + return result; } } } diff --git a/src/ServiceControl.Persistence/RetentionSweepCurrentStatus.cs b/src/ServiceControl.Persistence/RetentionSweepCurrentStatus.cs index 2bc32143ea..9a5da4a2bf 100644 --- a/src/ServiceControl.Persistence/RetentionSweepCurrentStatus.cs +++ b/src/ServiceControl.Persistence/RetentionSweepCurrentStatus.cs @@ -2,10 +2,15 @@ namespace ServiceControl.Persistence; using System; -/// A point-in-time snapshot of sweep execution state. +/// +/// A point-in-time snapshot of sweep execution state. It describes the most recent sweep, whether +/// the hourly timer or a manual request started it. +/// public sealed record RetentionSweepCurrentStatus( bool IsRunning, DateTime? LastStartedAt, DateTime? LastFinishedAt, DateTime? LastErrorCutoff, - DateTime? LastEventsCutoff); \ No newline at end of file + DateTime? LastEventsCutoff, + RetentionSweepOutcome? LastOutcome, + string? LastError); \ No newline at end of file diff --git a/src/ServiceControl.Persistence/RetentionSweepOutcome.cs b/src/ServiceControl.Persistence/RetentionSweepOutcome.cs new file mode 100644 index 0000000000..23154897ab --- /dev/null +++ b/src/ServiceControl.Persistence/RetentionSweepOutcome.cs @@ -0,0 +1,20 @@ +namespace ServiceControl.Persistence; + +/// +/// How a sweep that has ended turned out. +/// +public enum RetentionSweepOutcome +{ + /// + /// Every pass ran to the end without an error. + /// + Succeeded, + /// + /// At least one pass threw. The other passes still ran. + /// + Failed, + /// + /// The sweep was cancelled before it finished, for example by shutdown. + /// + Cancelled +} \ No newline at end of file diff --git a/src/ServiceControl.UnitTests/Infrastructure/RetentionApiTests.cs b/src/ServiceControl.UnitTests/Infrastructure/RetentionApiTests.cs new file mode 100644 index 0000000000..9d239f829d --- /dev/null +++ b/src/ServiceControl.UnitTests/Infrastructure/RetentionApiTests.cs @@ -0,0 +1,49 @@ +namespace ServiceControl.UnitTests.Infrastructure; + +using System; +using System.Threading; +using System.Threading.Tasks; +using Microsoft.Extensions.DependencyInjection; +using NUnit.Framework; +using ServiceControl.Infrastructure.Api; +using ServiceControl.Persistence; + +[TestFixture] +public class RetentionApiTests +{ + [Test] + public async Task Every_sweep_outcome_is_reported_under_the_same_name([Values] RetentionSweepOutcome outcome) + { + var response = await GetStatus(outcome, "FailedMessages: boom"); + + using (Assert.EnterMultipleScope()) + { + Assert.That(response.LastOutcome?.ToString(), Is.EqualTo(outcome.ToString())); + Assert.That(response.LastError, Is.EqualTo("FailedMessages: boom")); + } + } + + [Test] + public async Task A_running_sweep_reports_no_outcome() + { + var response = await GetStatus(null, null); + + Assert.That(response.LastOutcome, Is.Null); + } + + static Task GetStatus(RetentionSweepOutcome? outcome, string error) + { + var status = new RetentionSweepCurrentStatus(outcome is null, DateTime.UtcNow, null, null, null, outcome, error); + var services = new ServiceCollection().AddSingleton(new StubSweeper(status)).BuildServiceProvider(); + + return new RetentionApi(services).GetStatus(); + } + + class StubSweeper(RetentionSweepCurrentStatus status) : IRetentionSweeper + { + public ManualSweepAttempt TryStartManualSweep(DateTime? errorCutoff, DateTime? eventsCutoff, CancellationToken cancellationToken = default) => + throw new NotSupportedException(); + + public RetentionSweepCurrentStatus GetStatus() => status; + } +} diff --git a/src/ServiceControl/Infrastructure/Api/RetentionApi.cs b/src/ServiceControl/Infrastructure/Api/RetentionApi.cs index 44d3c6fa71..1a64dd6a7b 100644 --- a/src/ServiceControl/Infrastructure/Api/RetentionApi.cs +++ b/src/ServiceControl/Infrastructure/Api/RetentionApi.cs @@ -63,7 +63,16 @@ public Task GetStatus(CancellationToken cancellati LastStartedAt = status.LastStartedAt, LastFinishedAt = status.LastFinishedAt, LastErrorCutoff = status.LastErrorCutoff, - LastEventsCutoff = status.LastEventsCutoff + LastEventsCutoff = status.LastEventsCutoff, + LastOutcome = status.LastOutcome switch + { + null => null, + RetentionSweepOutcome.Succeeded => RetentionPurgeOutcome.Succeeded, + RetentionSweepOutcome.Failed => RetentionPurgeOutcome.Failed, + RetentionSweepOutcome.Cancelled => RetentionPurgeOutcome.Cancelled, + _ => throw new ArgumentOutOfRangeException(nameof(status), status.LastOutcome, "Unknown retention sweep outcome.") + }, + LastError = status.LastError }); }