Skip to content
30 changes: 30 additions & 0 deletions docs/telemetry.md
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,36 @@ What the shapes mean when tuning:

Example Grafana dashboard - https://github.com/andreasohlund/Docker/blob/main/otel-monitoring/grafana-platform-template.json

## Retention

Only the SQL persisters report retention. On RavenDB the expiry is metadata on each document and the deletes happen inside the Raven server, so there is no sweep of ours to measure.

Meter `Particular.ServiceControl`, the same meter the instance uses for its ingestion. The prefix carries no instance segment, so when the audit SQL persister reports its own retention it uses these same names, and `job` (or `exported_job`) separates the two.

The sweep runs hourly and makes one pass per kind of row. Each pass is measured on its own, tagged `retention.entity`: `failed_messages`, `event_log` or `group_comments`.

- `sc.retention.cycle_duration_seconds` - Retention sweep pass duration in seconds
- `retention.entity` - Which pass this was
- `result` - The outcome: `success`, `failed`, or `cancelled` if shutdown cut the pass short
- `sc.retention.rows_deleted_total` - Rows deleted by the retention sweep
- `retention.entity` - Which pass deleted them
- `sc.retention.consecutive_failures_total` - Consecutive failures of that pass
- `retention.entity` - Which pass is failing

### Reading the retention metrics

- Rows reclaimed per hour: `sum(rate(sc_retention_rows_deleted_total[1h])) by (exported_job,retention_entity)`
- Sweep duration: `histogram_quantile(0.9,sum(rate(sc_retention_cycle_duration_seconds_bucket[6h])) by (le,exported_job,retention_entity))`
- Retention is broken: `max(sc_retention_consecutive_failures_total) by (exported_job,retention_entity) > 2`

What the shapes mean:

- `consecutive_failures_total` above zero is the signal that rows are no longer being reclaimed. Nothing else in the product reports this, and the database grows without bound while it lasts.
- `rows_deleted_total` flat at zero across a long window is only healthy if the instance is also not ingesting. Deletion stopping while ingestion continues means the retention window is not being enforced.
- A cycle duration climbing towards the hourly interval means the sweep is no longer keeping up with the arrival rate, and each run starts further behind than the last.
- A body store that refuses a delete fails the `failed_messages` pass rather than orphaning the body, so an expired credential or a changed permission shows as a climbing gauge instead of storage that quietly keeps growing. Those rows stay until a sweep can delete the body and the row together.
- Each pass is isolated, so one kind of row failing to be reclaimed does not stop the others. Every pass reports a result on every run, and the gauge is per pass, so alert across all entities rather than on any one of them.

## Monitoring

No telemetry is currently available.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -7,11 +7,12 @@ namespace ServiceControl.Audit.Auditing.Metrics;
using EndpointPlugin.Messages.SagaState;
using NServiceBus;
using NServiceBus.Transport;
using ServiceControl.Infrastructure;
using ServiceControl.Infrastructure.Ingestion.Metrics;

public class IngestionMetrics
{
public const string MeterName = "Particular.ServiceControl.Audit";
public const string MeterName = ServiceControlMeters.Audit;

public static readonly string BatchDurationInstrumentName = $"{InstrumentPrefix}.batch_duration_seconds";
public static readonly string MessageDurationInstrumentName = $"{InstrumentPrefix}.message_duration_seconds";
Expand Down
12 changes: 12 additions & 0 deletions src/ServiceControl.Infrastructure/ServiceControlMeters.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,12 @@
namespace ServiceControl.Infrastructure;

/// <summary>
/// The meters each instance publishes on. Shared because persisters publish onto the meter their
/// host has already registered with the exporter, and the two assemblies cannot reference each
/// other.
/// </summary>
public static class ServiceControlMeters
{
public const string Error = "Particular.ServiceControl";
public const string Audit = "Particular.ServiceControl.Audit";
}
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ namespace ServiceControl.Persistence.EFCore.Abstractions;
using ServiceControl.Persistence.EFCore.Implementation.Recoverability;
using ServiceControl.Persistence.EFCore.Implementation.UnitOfWork;
using ServiceControl.Persistence.EFCore.Infrastructure;
using ServiceControl.Persistence.EFCore.Infrastructure.Metrics;
using ServiceControl.Persistence.MessageRedirects;
using ServiceControl.Persistence.Recoverability;
using ServiceControl.Persistence.UnitOfWork;
Expand All @@ -37,6 +38,7 @@ protected static void RegisterDataStores(IServiceCollection services, EFPersiste

if (settings.RunRetentionSweep)
{
services.AddSingleton<RetentionMetrics>();
services.AddHostedService<RetentionSweeper>();
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,6 @@ public async Task WriteBody(string bodyId, ReadOnlyMemory<byte> body, string con
}
}

public Task DeleteBody(string bodyId, CancellationToken cancellationToken = default) =>
public Task DeleteBodyIfExists(string bodyId, CancellationToken cancellationToken = default) =>
container.GetBlobClient(bodyId).DeleteIfExistsAsync(cancellationToken: cancellationToken);
}
Original file line number Diff line number Diff line change
Expand Up @@ -135,7 +135,7 @@ public async Task WriteBody(string bodyId, ReadOnlyMemory<byte> body, string con
}
}

public Task DeleteBody(string bodyId, CancellationToken cancellationToken = default)
public Task DeleteBodyIfExists(string bodyId, CancellationToken cancellationToken = default)
{
TryDelete(GetBodyFilePath(bodyId));
return Task.CompletedTask;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -93,7 +93,7 @@ public async Task WriteBody(string bodyId, ReadOnlyMemory<byte> body, string con
}
}

public Task DeleteBody(string bodyId, CancellationToken cancellationToken = default) =>
public Task DeleteBodyIfExists(string bodyId, CancellationToken cancellationToken = default) =>
client.DeleteObjectAsync(bucketName, Key(bodyId), cancellationToken);

async Task<bool> Exists(string key, CancellationToken cancellationToken)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -189,7 +189,7 @@ async Task DeleteExternalBody(Guid uniqueMessageId, CancellationToken cancellati
{
try
{
await bodyStorage.DeleteBody(FailedErrorImportEntity.ExternalBodyId(uniqueMessageId), cancellationToken);
await bodyStorage.DeleteBodyIfExists(FailedErrorImportEntity.ExternalBodyId(uniqueMessageId), cancellationToken);
}
#pragma warning disable PS0019 // The filter already excludes OperationCanceledException, so cancellation
// propagates; PS0019 only recognises a cancellationToken.IsCancellationRequested guard.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,5 +15,10 @@ public interface IBodyStoragePersistence
{
Task WriteBody(string bodyId, ReadOnlyMemory<byte> body, string contentType, CancellationToken cancellationToken = default);
Task<MessageBodyFileResult?> ReadBody(string bodyId, CancellationToken cancellationToken = default);
Task DeleteBody(string bodyId, CancellationToken cancellationToken = default);

/// <summary>
/// Implementations throw only when the store itself fails, never because the body was already
/// gone: callers delete the body before the row that names it.
/// </summary>
Task DeleteBodyIfExists(string bodyId, CancellationToken cancellationToken = default);
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
namespace ServiceControl.Persistence.EFCore.Infrastructure.Metrics;

enum CycleOutcome
{
Success,
Cancelled,
Failed
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
namespace ServiceControl.Persistence.EFCore.Infrastructure.Metrics;

using System.Diagnostics;

/// <summary>
/// One pass of the retention sweep. A pass that finished counts as a success even if shutdown has
/// since been requested; one that shutdown cut short is recorded as cancelled rather than as a
/// failure.
/// </summary>
public sealed class RetentionCycleMetrics : IDisposable
{
internal RetentionCycleMetrics(RetentionMetrics metrics, RetentionEntity entity, CancellationToken cancellationToken)
{
this.metrics = metrics;
this.entity = entity;
this.cancellationToken = cancellationToken;
}

// The pass is over once it completes, so the clock stops here rather than wherever the scope
// happens to be disposed.
public void Complete()
{
stopwatch.Stop();
completed = true;
}

public void Dispose() => metrics.RecordCycle(entity, stopwatch.Elapsed, Outcome);

CycleOutcome Outcome =>
completed ? CycleOutcome.Success
: cancellationToken.IsCancellationRequested ? CycleOutcome.Cancelled
: CycleOutcome.Failed;

bool completed;

readonly RetentionMetrics metrics;
readonly RetentionEntity entity;
readonly CancellationToken cancellationToken;
readonly Stopwatch stopwatch = Stopwatch.StartNew();
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
namespace ServiceControl.Persistence.EFCore.Infrastructure.Metrics;

public enum RetentionEntity
{
FailedMessages,
EventLog,
GroupComments
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,90 @@
namespace ServiceControl.Persistence.EFCore.Infrastructure.Metrics;

using System.Diagnostics;
using System.Diagnostics.Metrics;
using ServiceControl.Infrastructure;

public class RetentionMetrics
{
public const string MeterName = ServiceControlMeters.Error;

public static readonly string CycleDurationInstrumentName = $"{InstrumentPrefix}.cycle_duration_seconds";
public static readonly string RowsDeletedInstrumentName = $"{InstrumentPrefix}.rows_deleted_total";
public static readonly string ConsecutiveFailuresInstrumentName = $"{InstrumentPrefix}.consecutive_failures_total";

public RetentionMetrics(IMeterFactory meterFactory)
{
var meter = meterFactory.Create(MeterName, MeterVersion);

cycleDuration = meter.CreateHistogram<double>(
CycleDurationInstrumentName,
unit: "seconds",
description: "Retention sweep pass duration in seconds",
tags: null,
// A sweep pass is sub-second when it is keeping up and minutes long when it is working
// through a backlog, so the default boundaries resolve neither end.
advice: new InstrumentAdvice<double> { HistogramBucketBoundaries = [0.1, 0.5, 1, 5, 15, 60, 300, 900] });

rowsDeleted = meter.CreateCounter<long>(RowsDeletedInstrumentName, description: "Rows deleted by the retention sweep");
consecutiveFailureGauge = meter.CreateObservableGauge(ConsecutiveFailuresInstrumentName, ObserveConsecutiveFailures, description: "Consecutive retention sweep failures");
}

public RetentionCycleMetrics BeginCycle(RetentionEntity entity, CancellationToken cancellationToken = default) => new(this, entity, cancellationToken);

public void RecordRowsDeleted(RetentionEntity entity, int rows) => rowsDeleted.Add(rows, EntityTags[(int)entity]);

internal void RecordCycle(RetentionEntity entity, TimeSpan elapsed, CycleOutcome outcome)
{
var tags = EntityTags[(int)entity];
tags.Add("result", ResultTag(outcome));

cycleDuration.Record(elapsed.TotalSeconds, tags);

// A pass cut short by shutdown neither proves the sweep healthy nor faulty, so it leaves
// the gauge where it was.
if (outcome == CycleOutcome.Success)
{
Interlocked.Exchange(ref consecutiveFailures[(int)entity], 0);
}
else if (outcome == CycleOutcome.Failed)
{
Interlocked.Increment(ref consecutiveFailures[(int)entity]);
}
}

IEnumerable<Measurement<long>> ObserveConsecutiveFailures()
{
for (var entity = 0; entity < consecutiveFailures.Length; entity++)
{
yield return new Measurement<long>(Volatile.Read(ref consecutiveFailures[entity]), EntityTags[entity]);
}
}

static TagList EntityTag(string entity) => new() { { "retention.entity", entity } };

static string ResultTag(CycleOutcome outcome) => outcome switch
{
CycleOutcome.Success => "success",
CycleOutcome.Cancelled => "cancelled",
CycleOutcome.Failed => "failed",
_ => throw new ArgumentOutOfRangeException(nameof(outcome))
};

readonly long[] consecutiveFailures = new long[EntityTags.Length];

readonly Histogram<double> cycleDuration;
readonly Counter<long> rowsDeleted;
#pragma warning disable IDE0052
readonly ObservableGauge<long> consecutiveFailureGauge;
#pragma warning restore IDE0052

static readonly TagList[] EntityTags =
[
EntityTag("failed_messages"),
EntityTag("event_log"),
EntityTag("group_comments")
];

const string MeterVersion = "0.1.0";
const string InstrumentPrefix = "sc.retention";
}
Loading