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 @@ -47,6 +47,9 @@ protected static void RegisterDataStores(IServiceCollection services, EFPersiste
services.AddSingleton<RetentionSweeper>();
services.AddHostedService(sp => sp.GetRequiredService<RetentionSweeper>());
services.AddSingleton<IRetentionSweeper>(sp => sp.GetRequiredService<RetentionSweeper>());

services.AddSingleton<RetentionSweepCustomCheck.State>();
services.AddCustomCheck<RetentionSweepCustomCheck>();
}

services.AddSingleton<OperationsManager>();
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,60 @@
namespace ServiceControl.Persistence.EFCore.Infrastructure;

using System.Threading;
using System.Threading.Tasks;
using ServiceControl.Persistence.EFCore.Infrastructure.Metrics;
using NServiceBus.CustomChecks;

// The standard custom-check state-transition pipeline deduplicates repeated results.
class RetentionSweepCustomCheck(RetentionSweepCustomCheck.State state)
: CustomCheck("ServiceControl Retention", "ServiceControl Health", TimeSpan.FromMinutes(1))
{
public override Task<CheckResult> PerformCheck(CancellationToken cancellationToken = default)
{
var failures = state.GetFailures();
var failureSummary = string.Join("; ", failures.Select(failure => $"{failure.Key}: {failure.Value}"));

return Task.FromResult(state.ConsecutiveFailedSweeps < 3
? CheckResult.Pass
: CheckResult.Failed($"Retention processing has failures. Last failure per entity: {failureSummary}. See https://docs.particular.net/servicecontrol/troubleshooting for guidance on resolving the issue."));

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We may need to adjust ServicePulse to render these links correctly

}

internal class State
{
readonly Dictionary<RetentionEntity, string> failures = [];

public int ConsecutiveFailedSweeps { get; private set; }

public void Clear(RetentionEntity entity)
{
lock (failures)
{
failures.Remove(entity);
}
}

public void ReportError(RetentionEntity entity, string reason)
{
lock (failures)
{
failures[entity] = reason;
}
}

public void SweepComplete()
{
lock (failures)
{
ConsecutiveFailedSweeps = failures.Count == 0 ? 0 : ConsecutiveFailedSweeps + 1;
}
}

internal KeyValuePair<RetentionEntity, string>[] GetFailures()
{
lock (failures)
{
return failures.ToArray();
}
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -17,12 +17,13 @@ namespace ServiceControl.Persistence.EFCore.Infrastructure;
//
// A manual sweep can be triggered via the API (see IRetentionSweeper / IRetentionApi) with
// caller-supplied cutoffs.
public class RetentionSweeper(
class RetentionSweeper(
ILogger<RetentionSweeper> logger,
TimeProvider timeProvider,
IServiceScopeFactory serviceScopeFactory,
IBodyStoragePersistence bodyStorage,
RetentionMetrics metrics,
RetentionSweepCustomCheck.State retentionState,
EFPersisterSettings settings,
IHostApplicationLifetime hostApplicationLifetime) : BackgroundService, IRetentionSweeper
{
Expand Down Expand Up @@ -155,6 +156,7 @@ async Task SweepBody(DateTime? errorCutoff, DateTime? eventsCutoff, bool pace, C
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();
}

// Each pass is isolated so one failing kind of row does not stop the others from being
Expand All @@ -166,14 +168,15 @@ async Task RunPass(RetentionEntity entity, Func<CancellationToken, Task> pass, C
try
{
await pass(cancellationToken);

cycle.Complete();
retentionState.Clear(entity);
}
#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)
{
logger.LogError(ex, "Error during the {RetentionEntity} retention pass", entity);
retentionState.ReportError(entity, ex.Message);
}
#pragma warning restore PS0019
}
Expand Down
Original file line number Diff line number Diff line change
@@ -1 +1,2 @@
ServiceControl Health: ServiceControl Retention
Storage space: ServiceControl body storage
Original file line number Diff line number Diff line change
@@ -1 +1,2 @@
ServiceControl Health: ServiceControl Retention
Storage space: ServiceControl body storage
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ namespace ServiceControl.Persistence.Tests;
using System.IO;
using System.Linq;
using System.Threading.Tasks;
using EFCore.Infrastructure;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging.Abstractions;
using NServiceBus.CustomChecks;
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,150 @@
namespace ServiceControl.Persistence.Tests;

using System;
using System.Linq;
using System.Threading.Tasks;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Logging.Abstractions;
using NServiceBus.CustomChecks;
using NUnit.Framework;
using ServiceControl.Contracts.CustomChecks;
using ServiceControl.CustomChecks;
using ServiceControl.Infrastructure.DomainEvents;
using ServiceControl.Operations;
using ServiceControl.Persistence.EFCore.Infrastructure;
using ServiceControl.Persistence.EFCore.Infrastructure.Metrics;

class RetentionSweepCustomCheckTests : PersistenceTestBase
{
RetentionSweeper Sweeper => ServiceProvider.GetRequiredService<RetentionSweeper>();
RetentionSweepCustomCheck.State State => ServiceProvider.GetRequiredService<RetentionSweepCustomCheck.State>();

RetentionSweepCustomCheck Check =>
ServiceProvider.GetServices<ICustomCheck>().OfType<RetentionSweepCustomCheck>().Single();

[Test]
public async Task Check_initially_passes()
{
var result = await Check.PerformCheck();
Assert.That(result, Is.EqualTo(CheckResult.Pass));
}

[Test]
public async Task Check_remains_passing_until_three_consecutive_failed_sweeps()
{
State.ReportError(RetentionEntity.FailedMessages, "db timeout");

State.SweepComplete();
Assert.That(await Check.PerformCheck(), Is.EqualTo(CheckResult.Pass));
State.SweepComplete();
Assert.That(await Check.PerformCheck(), Is.EqualTo(CheckResult.Pass));
State.SweepComplete();

var result = await Check.PerformCheck();
using (Assert.EnterMultipleScope())
{
Assert.That(result.HasFailed, Is.True);
Assert.That(result.FailureReason, Does.Contain("FailedMessages"));
Assert.That(result.FailureReason, Does.Contain("db timeout"));
Assert.That(result.FailureReason, Does.Contain("https://docs.particular.net/servicecontrol/troubleshooting"));
}
}

[Test]
public async Task Successful_sweep_resets_the_consecutive_failure_count()
{
State.ReportError(RetentionEntity.FailedMessages, "db timeout");
CompleteThreeSweeps();

await Sweeper.SweepNow();
State.ReportError(RetentionEntity.FailedMessages, "another timeout");
State.SweepComplete();

var result = await Check.PerformCheck();
Assert.That(result, Is.EqualTo(CheckResult.Pass));
}

[Test]
public async Task Check_does_not_expire_with_time()
{
State.ReportError(RetentionEntity.FailedMessages, "db timeout");
CompleteThreeSweeps();
AdvanceClock(TimeSpan.FromHours(2));

var result = await Check.PerformCheck();
Assert.That(result.HasFailed, Is.True);
}

[Test]
public async Task Multiple_failing_entities_are_represented()
{
State.ReportError(RetentionEntity.FailedMessages, "body delete failed");
State.ReportError(RetentionEntity.EventLog, "batch delete timeout");
CompleteThreeSweeps();

var result = await Check.PerformCheck();

using (Assert.EnterMultipleScope())
{
Assert.That(result.HasFailed, Is.True);
Assert.That(result.FailureReason, Does.Contain("FailedMessages"));
Assert.That(result.FailureReason, Does.Contain("EventLog"));
Assert.That(result.FailureReason, Does.Contain("body delete failed"));
Assert.That(result.FailureReason, Does.Contain("batch delete timeout"));
}
}

[Test]
public async Task Clearing_one_retention_type_keeps_other_failures()
{
State.ReportError(RetentionEntity.FailedMessages, "body delete failed");
State.ReportError(RetentionEntity.EventLog, "batch delete timeout");
CompleteThreeSweeps();

State.Clear(RetentionEntity.FailedMessages);

var result = await Check.PerformCheck();
using (Assert.EnterMultipleScope())
{
Assert.That(result.FailureReason, Does.Not.Contain("FailedMessages"));
Assert.That(result.FailureReason, Does.Contain("EventLog"));
}
}

[Test]
public async Task Repeated_failed_results_raise_one_state_change_event()
{
State.ReportError(RetentionEntity.FailedMessages, "db timeout");
CompleteThreeSweeps();

var result = await Check.PerformCheck();
var domainEvents = (FakeDomainEvents)ServiceProvider.GetRequiredService<IDomainEvents>();
var processor = new CustomCheckResultProcessor(domainEvents, CustomChecks, NullLogger<CustomCheckResultProcessor>.Instance);
var detail = new CustomCheckDetail
{
Category = "ServiceControl Health",
CustomCheckId = "ServiceControl Retention",
HasFailed = result.HasFailed,
FailureReason = result.FailureReason,
ReportedAt = Now,
OriginatingEndpoint = new EndpointDetails
{
Host = "localhost",
HostId = Guid.NewGuid(),
Name = "ServiceControl"
}
};

await processor.ProcessResult(detail);
await processor.ProcessResult(detail);

Assert.That(domainEvents.RaisedEvents.OfType<CustomCheckFailed>().Count(), Is.EqualTo(1));
}

void CompleteThreeSweeps()
{
State.SweepComplete();
State.SweepComplete();
State.SweepComplete();
}
}
19 changes: 19 additions & 0 deletions src/ServiceControl.Persistence.Tests/EFCore/RetentionSweepTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ namespace ServiceControl.Persistence.Tests;
using ServiceControl.EventLog;
using ServiceControl.MessageFailures;
using ServiceControl.Persistence.EFCore.Entities;
using ServiceControl.Persistence.EFCore.Infrastructure;
using ServiceControl.Persistence.EFCore.Infrastructure.Metrics;
using ServiceControl.Persistence.Infrastructure;

Expand Down Expand Up @@ -308,6 +309,24 @@ public async Task A_failing_pass_does_not_stop_the_others()
}
}

[Test]
public async Task A_failing_pass_is_tracked()
{
// Subtracting this from the clock cannot be represented, so the failed messages pass throws
// before it reaches the database.
EFSettings.ErrorRetentionPeriod = TimeSpan.FromDays(1_000_000);

await RunRetentionSweep();

var failure = ServiceProvider.GetRequiredService<RetentionSweepCustomCheck.State>().GetFailures().Single();

using (Assert.EnterMultipleScope())
{
Assert.That(failure.Key, Is.EqualTo(RetentionEntity.FailedMessages));
Assert.That(failure.Value, Is.Not.Empty);
}
}

RecordedRetentionMetrics ListenToRetentionMetrics() => new(ServiceProvider.GetRequiredService<IMeterFactory>());

async Task<string> SeedGroup(Guid uniqueMessageId)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ public async Task SetUp()
}

var hostBuilder = Host.CreateApplicationBuilder();
hostBuilder.Services.AddMetrics();

LoggerUtil.ActiveLoggers = Loggers.Test;
hostBuilder.Logging.ConfigureLogging(LogLevel.Information);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@ public static class InternalCustomCheckClassification
"Error Database Search Engine", // RavenDB persister
"ServiceControl body storage", // EF Core persisters
"Dead Letter Queue", // ASBS / IBMMQ / MSMQ
"ServiceControl Retention", // EF Core persisters

// ----- Audit instance (forwarded to the primary via ReportCustomCheckResult) -----
"Audit Message Ingestion",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ public void Internal_checks_are_flagged_internal()
[TestCase("Error Message Ingestion")]
[TestCase("Dead Letter Queue")]
[TestCase("ServiceControl body storage")]
[TestCase("ServiceControl Retention")]
[TestCase("Audit Message Ingestion Process")]
public void Every_shipped_check_is_internal(string id)
{
Expand Down
6 changes: 5 additions & 1 deletion src/ServiceControl/App.config
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,9 @@ These settings are only here so that we can debug ServiceControl while developin
<!--<add key="ServiceControl/TransportType" value="RabbitMQ.QuorumConventionalRouting" />-->
<!--<add key="ServiceControl/TransportType" value="SQLServer" />-->

<add key="ServiceControl/PersistenceType" value="RavenDB" />
<add key="ServiceControl/PersistenceType" value="RavenDB" />
<!--<add key="ServiceControl/PersistenceType" value="SQLServer" />-->
<!--<add key="ServiceControl/PersistenceType" value="PostgreSQL" />-->

<!-- options are any comma separated combination of NLog,Seq,Otlp -->
<add key="ServiceControl/LoggingProviders" value="NLog,Seq"/>
Expand Down Expand Up @@ -84,6 +86,8 @@ These settings are only here so that we can debug ServiceControl while developin
<!-- Learning -->
<!-- If the LearningTransport connectionString is empty, it will default to the solution directory when running/debugging from the IDE -->
<add name="NServiceBus/Transport" connectionString="" />



<!--Amazon SQS -->
<!--<add name="NServiceBus/Transport" connectionString="Region=;QueueNamePrefix=;TopicNamePrefix=;AccessKeyId=;SecretAccessKey=;S3BucketForLargeMessages=;S3KeyPrefix=" />-->
Expand Down
Loading