Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -94,6 +94,7 @@ await Define<Context>()
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");
}
}

Expand Down
8 changes: 8 additions & 0 deletions src/ServiceControl.Api/Contracts/RetentionPurgeOutcome.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
namespace ServiceControl.Api.Contracts;

public enum RetentionPurgeOutcome
{
Succeeded,
Failed,
Cancelled
}
11 changes: 11 additions & 0 deletions src/ServiceControl.Api/Contracts/RetentionPurgeStatusResponse.cs
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,17 @@ public class RetentionPurgeStatusResponse

public DateTime? LastEventsCutoff { get; set; }

/// <summary>
/// How the most recent purge ended. Null while it is still running or before any has run.
/// </summary>
public RetentionPurgeOutcome? LastOutcome { get; set; }

/// <summary>
/// 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.
/// </summary>
public string LastError { get; set; }

/// <summary>Present only on the 501 Not Implemented response.</summary>
public string Reason { get; set; }
}
Original file line number Diff line number Diff line change
Expand Up @@ -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)
{
Expand All @@ -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.
Expand All @@ -88,80 +82,108 @@ 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()
{
using var cancellation = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken, hostApplicationLifetime.ApplicationStopping);
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<string?> passErrors)
{
var errors = passErrors.OfType<string>().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<string?> 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<CancellationToken, Task> pass, CancellationToken cancellationToken)
// Returns what went wrong for the status to report, or null when the pass succeeded.
async Task<string?> RunPass(RetentionEntity entity, Func<CancellationToken, Task> pass, CancellationToken cancellationToken)
{
using var cycle = metrics.BeginCycle(entity, cancellationToken);

Expand All @@ -170,17 +192,21 @@ async Task RunPass(RetentionEntity entity, Func<CancellationToken, Task> 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.
Expand All @@ -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<ServiceControlDbContext>();

Expand Down Expand Up @@ -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<ServiceControlDbContext>();

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,11 @@ class InMemoryBodyStoragePersistence : IBodyStoragePersistence

public HashSet<string> FailDeleteFor { get; } = [];

/// <summary>
/// Extra lines added to the simulated delete failure, the way a storage SDK appends response detail.
/// </summary>
public string DeleteFailureDetail { get; set; }

public IReadOnlyList<StoredBody> Written
{
get
Expand Down Expand Up @@ -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)
Expand Down
Loading
Loading