Skip to content

Commit f17ea66

Browse files
authored
Add custom check for failing retention (#5903)
* Add custom check for failing retention * Better naming for custom check * Reset dev config * Add targeted interface for faking so custom check doesn't need a real retention sweeper in tests * Move the retention custom check state out of the sweeper * Add link to troubleshooting page, keep the check silent until it fails 3 times * Check consecutive failures in custom check instead of baking it into state
1 parent 6618615 commit f17ea66

12 files changed

Lines changed: 248 additions & 3 deletions

File tree

‎src/ServiceControl.Persistence.EFCore/Abstractions/BasePersistence.cs‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -47,6 +47,9 @@ protected static void RegisterDataStores(IServiceCollection services, EFPersiste
4747
services.AddSingleton<RetentionSweeper>();
4848
services.AddHostedService(sp => sp.GetRequiredService<RetentionSweeper>());
4949
services.AddSingleton<IRetentionSweeper>(sp => sp.GetRequiredService<RetentionSweeper>());
50+
51+
services.AddSingleton<RetentionSweepCustomCheck.State>();
52+
services.AddCustomCheck<RetentionSweepCustomCheck>();
5053
}
5154

5255
services.AddSingleton<OperationsManager>();
Lines changed: 60 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,60 @@
1+
namespace ServiceControl.Persistence.EFCore.Infrastructure;
2+
3+
using System.Threading;
4+
using System.Threading.Tasks;
5+
using ServiceControl.Persistence.EFCore.Infrastructure.Metrics;
6+
using NServiceBus.CustomChecks;
7+
8+
// The standard custom-check state-transition pipeline deduplicates repeated results.
9+
class RetentionSweepCustomCheck(RetentionSweepCustomCheck.State state)
10+
: CustomCheck("ServiceControl Retention", "ServiceControl Health", TimeSpan.FromMinutes(1))
11+
{
12+
public override Task<CheckResult> PerformCheck(CancellationToken cancellationToken = default)
13+
{
14+
var failures = state.GetFailures();
15+
var failureSummary = string.Join("; ", failures.Select(failure => $"{failure.Key}: {failure.Value}"));
16+
17+
return Task.FromResult(state.ConsecutiveFailedSweeps < 3
18+
? CheckResult.Pass
19+
: CheckResult.Failed($"Retention processing has failures. Last failure per entity: {failureSummary}. See https://docs.particular.net/servicecontrol/troubleshooting for guidance on resolving the issue."));
20+
}
21+
22+
internal class State
23+
{
24+
readonly Dictionary<RetentionEntity, string> failures = [];
25+
26+
public int ConsecutiveFailedSweeps { get; private set; }
27+
28+
public void Clear(RetentionEntity entity)
29+
{
30+
lock (failures)
31+
{
32+
failures.Remove(entity);
33+
}
34+
}
35+
36+
public void ReportError(RetentionEntity entity, string reason)
37+
{
38+
lock (failures)
39+
{
40+
failures[entity] = reason;
41+
}
42+
}
43+
44+
public void SweepComplete()
45+
{
46+
lock (failures)
47+
{
48+
ConsecutiveFailedSweeps = failures.Count == 0 ? 0 : ConsecutiveFailedSweeps + 1;
49+
}
50+
}
51+
52+
internal KeyValuePair<RetentionEntity, string>[] GetFailures()
53+
{
54+
lock (failures)
55+
{
56+
return failures.ToArray();
57+
}
58+
}
59+
}
60+
}

‎src/ServiceControl.Persistence.EFCore/Infrastructure/RetentionSweeper.cs‎

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,12 +17,13 @@ namespace ServiceControl.Persistence.EFCore.Infrastructure;
1717
//
1818
// A manual sweep can be triggered via the API (see IRetentionSweeper / IRetentionApi) with
1919
// caller-supplied cutoffs.
20-
public class RetentionSweeper(
20+
class RetentionSweeper(
2121
ILogger<RetentionSweeper> logger,
2222
TimeProvider timeProvider,
2323
IServiceScopeFactory serviceScopeFactory,
2424
IBodyStoragePersistence bodyStorage,
2525
RetentionMetrics metrics,
26+
RetentionSweepCustomCheck.State retentionState,
2627
EFPersisterSettings settings,
2728
IHostApplicationLifetime hostApplicationLifetime) : BackgroundService, IRetentionSweeper
2829
{
@@ -155,6 +156,7 @@ async Task SweepBody(DateTime? errorCutoff, DateTime? eventsCutoff, bool pace, C
155156
await RunPass(RetentionEntity.FailedMessages, token => SweepFailedMessages(pace, errorCutoff, token), cancellationToken);
156157
await RunPass(RetentionEntity.EventLog, token => SweepEventLogItems(pace, eventsCutoff, token), cancellationToken);
157158
await RunPass(RetentionEntity.GroupComments, SweepOrphanedGroupComments, cancellationToken);
159+
retentionState.SweepComplete();
158160
}
159161

160162
// Each pass is isolated so one failing kind of row does not stop the others from being
@@ -166,14 +168,15 @@ async Task RunPass(RetentionEntity entity, Func<CancellationToken, Task> pass, C
166168
try
167169
{
168170
await pass(cancellationToken);
169-
170171
cycle.Complete();
172+
retentionState.Clear(entity);
171173
}
172174
#pragma warning disable PS0019 // The filter already excludes OperationCanceledException, so
173175
// cancellation propagates; PS0019 only recognises a cancellationToken guard.
174176
catch (Exception ex) when (ex is not OperationCanceledException)
175177
{
176178
logger.LogError(ex, "Error during the {RetentionEntity} retention pass", entity);
179+
retentionState.ReportError(entity, ex.Message);
177180
}
178181
#pragma warning restore PS0019
179182
}
Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1 +1,2 @@
1+
ServiceControl Health: ServiceControl Retention
12
Storage space: ServiceControl body storage
Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1 +1,2 @@
1+
ServiceControl Health: ServiceControl Retention
12
Storage space: ServiceControl body storage

‎src/ServiceControl.Persistence.Tests/EFCore/FileSystemBodyStorageCustomCheckTests.cs‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@ namespace ServiceControl.Persistence.Tests;
44
using System.IO;
55
using System.Linq;
66
using System.Threading.Tasks;
7+
using EFCore.Infrastructure;
78
using Microsoft.Extensions.DependencyInjection;
89
using Microsoft.Extensions.Logging.Abstractions;
910
using NServiceBus.CustomChecks;
Lines changed: 150 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,150 @@
1+
namespace ServiceControl.Persistence.Tests;
2+
3+
using System;
4+
using System.Linq;
5+
using System.Threading.Tasks;
6+
using Microsoft.Extensions.DependencyInjection;
7+
using Microsoft.Extensions.Logging.Abstractions;
8+
using NServiceBus.CustomChecks;
9+
using NUnit.Framework;
10+
using ServiceControl.Contracts.CustomChecks;
11+
using ServiceControl.CustomChecks;
12+
using ServiceControl.Infrastructure.DomainEvents;
13+
using ServiceControl.Operations;
14+
using ServiceControl.Persistence.EFCore.Infrastructure;
15+
using ServiceControl.Persistence.EFCore.Infrastructure.Metrics;
16+
17+
class RetentionSweepCustomCheckTests : PersistenceTestBase
18+
{
19+
RetentionSweeper Sweeper => ServiceProvider.GetRequiredService<RetentionSweeper>();
20+
RetentionSweepCustomCheck.State State => ServiceProvider.GetRequiredService<RetentionSweepCustomCheck.State>();
21+
22+
RetentionSweepCustomCheck Check =>
23+
ServiceProvider.GetServices<ICustomCheck>().OfType<RetentionSweepCustomCheck>().Single();
24+
25+
[Test]
26+
public async Task Check_initially_passes()
27+
{
28+
var result = await Check.PerformCheck();
29+
Assert.That(result, Is.EqualTo(CheckResult.Pass));
30+
}
31+
32+
[Test]
33+
public async Task Check_remains_passing_until_three_consecutive_failed_sweeps()
34+
{
35+
State.ReportError(RetentionEntity.FailedMessages, "db timeout");
36+
37+
State.SweepComplete();
38+
Assert.That(await Check.PerformCheck(), Is.EqualTo(CheckResult.Pass));
39+
State.SweepComplete();
40+
Assert.That(await Check.PerformCheck(), Is.EqualTo(CheckResult.Pass));
41+
State.SweepComplete();
42+
43+
var result = await Check.PerformCheck();
44+
using (Assert.EnterMultipleScope())
45+
{
46+
Assert.That(result.HasFailed, Is.True);
47+
Assert.That(result.FailureReason, Does.Contain("FailedMessages"));
48+
Assert.That(result.FailureReason, Does.Contain("db timeout"));
49+
Assert.That(result.FailureReason, Does.Contain("https://docs.particular.net/servicecontrol/troubleshooting"));
50+
}
51+
}
52+
53+
[Test]
54+
public async Task Successful_sweep_resets_the_consecutive_failure_count()
55+
{
56+
State.ReportError(RetentionEntity.FailedMessages, "db timeout");
57+
CompleteThreeSweeps();
58+
59+
await Sweeper.SweepNow();
60+
State.ReportError(RetentionEntity.FailedMessages, "another timeout");
61+
State.SweepComplete();
62+
63+
var result = await Check.PerformCheck();
64+
Assert.That(result, Is.EqualTo(CheckResult.Pass));
65+
}
66+
67+
[Test]
68+
public async Task Check_does_not_expire_with_time()
69+
{
70+
State.ReportError(RetentionEntity.FailedMessages, "db timeout");
71+
CompleteThreeSweeps();
72+
AdvanceClock(TimeSpan.FromHours(2));
73+
74+
var result = await Check.PerformCheck();
75+
Assert.That(result.HasFailed, Is.True);
76+
}
77+
78+
[Test]
79+
public async Task Multiple_failing_entities_are_represented()
80+
{
81+
State.ReportError(RetentionEntity.FailedMessages, "body delete failed");
82+
State.ReportError(RetentionEntity.EventLog, "batch delete timeout");
83+
CompleteThreeSweeps();
84+
85+
var result = await Check.PerformCheck();
86+
87+
using (Assert.EnterMultipleScope())
88+
{
89+
Assert.That(result.HasFailed, Is.True);
90+
Assert.That(result.FailureReason, Does.Contain("FailedMessages"));
91+
Assert.That(result.FailureReason, Does.Contain("EventLog"));
92+
Assert.That(result.FailureReason, Does.Contain("body delete failed"));
93+
Assert.That(result.FailureReason, Does.Contain("batch delete timeout"));
94+
}
95+
}
96+
97+
[Test]
98+
public async Task Clearing_one_retention_type_keeps_other_failures()
99+
{
100+
State.ReportError(RetentionEntity.FailedMessages, "body delete failed");
101+
State.ReportError(RetentionEntity.EventLog, "batch delete timeout");
102+
CompleteThreeSweeps();
103+
104+
State.Clear(RetentionEntity.FailedMessages);
105+
106+
var result = await Check.PerformCheck();
107+
using (Assert.EnterMultipleScope())
108+
{
109+
Assert.That(result.FailureReason, Does.Not.Contain("FailedMessages"));
110+
Assert.That(result.FailureReason, Does.Contain("EventLog"));
111+
}
112+
}
113+
114+
[Test]
115+
public async Task Repeated_failed_results_raise_one_state_change_event()
116+
{
117+
State.ReportError(RetentionEntity.FailedMessages, "db timeout");
118+
CompleteThreeSweeps();
119+
120+
var result = await Check.PerformCheck();
121+
var domainEvents = (FakeDomainEvents)ServiceProvider.GetRequiredService<IDomainEvents>();
122+
var processor = new CustomCheckResultProcessor(domainEvents, CustomChecks, NullLogger<CustomCheckResultProcessor>.Instance);
123+
var detail = new CustomCheckDetail
124+
{
125+
Category = "ServiceControl Health",
126+
CustomCheckId = "ServiceControl Retention",
127+
HasFailed = result.HasFailed,
128+
FailureReason = result.FailureReason,
129+
ReportedAt = Now,
130+
OriginatingEndpoint = new EndpointDetails
131+
{
132+
Host = "localhost",
133+
HostId = Guid.NewGuid(),
134+
Name = "ServiceControl"
135+
}
136+
};
137+
138+
await processor.ProcessResult(detail);
139+
await processor.ProcessResult(detail);
140+
141+
Assert.That(domainEvents.RaisedEvents.OfType<CustomCheckFailed>().Count(), Is.EqualTo(1));
142+
}
143+
144+
void CompleteThreeSweeps()
145+
{
146+
State.SweepComplete();
147+
State.SweepComplete();
148+
State.SweepComplete();
149+
}
150+
}

‎src/ServiceControl.Persistence.Tests/EFCore/RetentionSweepTests.cs‎

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ namespace ServiceControl.Persistence.Tests;
1010
using ServiceControl.EventLog;
1111
using ServiceControl.MessageFailures;
1212
using ServiceControl.Persistence.EFCore.Entities;
13+
using ServiceControl.Persistence.EFCore.Infrastructure;
1314
using ServiceControl.Persistence.EFCore.Infrastructure.Metrics;
1415
using ServiceControl.Persistence.Infrastructure;
1516

@@ -308,6 +309,24 @@ public async Task A_failing_pass_does_not_stop_the_others()
308309
}
309310
}
310311

312+
[Test]
313+
public async Task A_failing_pass_is_tracked()
314+
{
315+
// Subtracting this from the clock cannot be represented, so the failed messages pass throws
316+
// before it reaches the database.
317+
EFSettings.ErrorRetentionPeriod = TimeSpan.FromDays(1_000_000);
318+
319+
await RunRetentionSweep();
320+
321+
var failure = ServiceProvider.GetRequiredService<RetentionSweepCustomCheck.State>().GetFailures().Single();
322+
323+
using (Assert.EnterMultipleScope())
324+
{
325+
Assert.That(failure.Key, Is.EqualTo(RetentionEntity.FailedMessages));
326+
Assert.That(failure.Value, Is.Not.Empty);
327+
}
328+
}
329+
311330
RecordedRetentionMetrics ListenToRetentionMetrics() => new(ServiceProvider.GetRequiredService<IMeterFactory>());
312331

313332
async Task<string> SeedGroup(Guid uniqueMessageId)

‎src/ServiceControl.Persistence.Tests/PersistenceTestBase.cs‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,7 @@ public async Task SetUp()
3434
}
3535

3636
var hostBuilder = Host.CreateApplicationBuilder();
37+
hostBuilder.Services.AddMetrics();
3738

3839
LoggerUtil.ActiveLoggers = Loggers.Test;
3940
hostBuilder.Logging.ConfigureLogging(LogLevel.Information);

‎src/ServiceControl.Persistence/InternalCustomCheckClassification.cs‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,7 @@ public static class InternalCustomCheckClassification
4343
"Error Database Search Engine", // RavenDB persister
4444
"ServiceControl body storage", // EF Core persisters
4545
"Dead Letter Queue", // ASBS / IBMMQ / MSMQ
46+
"ServiceControl Retention", // EF Core persisters
4647

4748
// ----- Audit instance (forwarded to the primary via ReportCustomCheckResult) -----
4849
"Audit Message Ingestion",

0 commit comments

Comments
 (0)