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
39 changes: 39 additions & 0 deletions docs/telemetry.md
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,45 @@ Meter `Particular.ServiceControl`.

`ServiceControl/PrintMetrics` predates this and no longer has anything to print.

### Retry

The retry pipeline runs in four stages: a bulk request is queued, preparation scans the store and cuts batches of 1000, staging dispatches a batch to the staging queue, and forwarding returns the staged messages to their senders. Each stage gets its own duration histogram, and the whole operation gets one end-to-end histogram on top.

Every instrument carries `retry.type`, one of `all`, `endpoint`, `group`, `queue`, `batch` or `single`.

- `sc.retry.operation_duration_seconds` - The whole retry operation, from the request arriving to the last message forwarded or skipped
- `result` - `success` or `failed`
- `sc.retry.prepare_duration_seconds` - Preparation, covering the store scan and batch creation. This is the stage that grows with the size of the error store.
- `result` - `success`, `failed`, or `cancelled` if shutdown cut the preparation short
- `sc.retry.stage_duration_seconds` - Staging one batch to the staging queue
- `result` - `success`, `failed`, `empty` for a batch that had no messages left and was discarded, or `cancelled` if shutdown cut the staging short
- `sc.retry.forward_duration_seconds` - Forwarding one batch back to the senders
- `result` - `success`, `failed`, or `cancelled` if shutdown cut the forwarding short
- `mode` - `counting`, or `timeout` when recovering from a premature shutdown. Timeout mode only ends on the forwarder's 45 second idle timer, so its distribution has a floor at that value.
- `sc.retry.messages_total` - Messages moved through the pipeline
- `result` - `staged`, `forwarded`, `skipped`, `staging_retried`, or `abandoned` for a message that hit the staging retry limit and was dropped from its batch. `abandoned` is the one to alert on: it is a message the user asked to retry that will not be retried.
- `sc.retry.operations_in_progress` - Retry operations currently in progress
- `retry.state` - `waiting`, `preparing` or `forwarding`
- `sc.retry.pending_bulk_requests` - Bulk retry requests queued behind each other, drained one per five second tick

A retry that hangs never records a duration, so on the histograms alone a stuck operation looks identical to no traffic. `operations_in_progress` holding a non-zero value while the duration histograms stay flat is the stuck-operation signal.

### Archive

Group archive and unarchive run as a loop over batches of 1000 until the group is drained. These instruments are emitted only when the instance runs on the SQL Server or PostgreSQL persistence; on RavenDB the family does not exist.

Every instrument carries `archive.operation`, either `archive` or `unarchive`.

- `sc.archive.operation_duration_seconds` - The whole group operation
- `sc.archive.batch_duration_seconds` - One batch of the loop, measured as the wall time between successive batch completions
- `sc.archive.messages_total` - Messages archived and unarchived
- `sc.archive.operations_in_progress` - Operations currently in progress
- `archive.state` - `started`, `progressing` or `finalizing`

### Host

With an OTLP endpoint configured the error instance also exports the standard [ASP.NET Core](https://learn.microsoft.com/en-us/dotnet/core/diagnostics/built-in-metrics-aspnetcore) instruments (`http.server.request.duration` per route, which covers the read APIs), the [HTTP client](https://learn.microsoft.com/en-us/dotnet/core/diagnostics/built-in-metrics-system-net) instruments (`http.client.request.duration`, which covers the scatter-gather calls to remote instances), and the [runtime](https://learn.microsoft.com/en-us/dotnet/core/diagnostics/built-in-metrics-runtime) instruments (GC, thread pool, exceptions).

## Audit

Meter `Particular.ServiceControl.Audit`.
Expand Down
3 changes: 3 additions & 0 deletions src/Directory.Packages.props
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,9 @@
<PackageVersion Include="OpenTelemetry.Exporter.Console" Version="1.18.0" />
<PackageVersion Include="OpenTelemetry.Exporter.OpenTelemetryProtocol" Version="1.18.0" />
<PackageVersion Include="OpenTelemetry.Extensions.Hosting" Version="1.18.0" />
<PackageVersion Include="OpenTelemetry.Instrumentation.AspNetCore" Version="1.18.0" />
<PackageVersion Include="OpenTelemetry.Instrumentation.Http" Version="1.18.0" />
<PackageVersion Include="OpenTelemetry.Instrumentation.Runtime" Version="1.18.0" />
<PackageVersion Include="Particular.Approvals" Version="2.0.1" />
<PackageVersion Include="Particular.LicensingComponent.Report" Version="1.2.0" />
<PackageVersion Include="Particular.Licensing.Sources" Version="7.4.0" />
Expand Down
8 changes: 2 additions & 6 deletions src/ServiceControl.Audit/Auditing/Metrics/IngestionMetrics.cs
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
namespace ServiceControl.Audit.Auditing.Metrics;
namespace ServiceControl.Audit.Auditing.Metrics;

using System.Collections.Generic;
using System.Diagnostics;
Expand All @@ -23,7 +23,7 @@ public IngestionMetrics(IMeterFactory meterFactory)

batchDuration = meter.CreateHistogram<double>(BatchDurationInstrumentName, unit: "seconds", "Message batch processing duration in seconds");
ingestionDuration = meter.CreateHistogram<double>(MessageDurationInstrumentName, unit: "seconds", description: "Audit message processing duration in seconds");
consecutiveBatchFailureGauge = meter.CreateObservableGauge($"{InstrumentPrefix}.consecutive_batch_failures_total", () => Volatile.Read(ref consecutiveBatchFailures), description: "Consecutive audit ingestion batch failures");
meter.CreateObservableGauge($"{InstrumentPrefix}.consecutive_batch_failures_total", () => Volatile.Read(ref consecutiveBatchFailures), description: "Consecutive audit ingestion batch failures");
failureCounter = meter.CreateCounter<long>($"{InstrumentPrefix}.failures_total", description: "Audit ingestion failure count");
}

Expand Down Expand Up @@ -64,10 +64,6 @@ void RecordBatchOutcome(bool success)
long consecutiveBatchFailures;

readonly Histogram<double> batchDuration;
#pragma warning disable IDE0052
// this can be changed to Gauge<T> once we can use the latest version of System.Diagnostics.DiagnosticSource
readonly ObservableGauge<long> consecutiveBatchFailureGauge;
#pragma warning restore IDE0052
readonly Histogram<double> ingestionDuration;
readonly Counter<long> failureCounter;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ namespace ServiceControl.Persistence.EFCore.Abstractions;
using ServiceControl.Persistence.Recoverability;
using ServiceControl.Persistence.UnitOfWork;
using ServiceControl.Recoverability;
using ServiceControl.Recoverability.Archiving.Metrics;

public abstract class BasePersistence
{
Expand All @@ -43,6 +44,7 @@ protected static void RegisterDataStores(IServiceCollection services, EFPersiste
}

services.AddSingleton<OperationsManager>();
services.AddSingleton<ArchiveMetrics>();
services.AddSingleton<IArchiveMessages, MessageArchiver>();
services.AddSingleton<ICustomChecksDataStore, CustomCheckDataStore>();
services.AddSingleton<IMessagesViewDataStore, MessagesViewDataStore>();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3,19 +3,20 @@ namespace ServiceControl.Persistence.EFCore.Implementation;
using ServiceControl.Infrastructure.DomainEvents;
using ServiceControl.Persistence.EFCore.Entities;
using ServiceControl.Recoverability;
using ServiceControl.Recoverability.Archiving.Metrics;

/// <summary>
/// EFCore equivalent of the RavenDB <see cref="ArchivingManager"/>. Wraps the shared
/// <see cref="OperationsManager"/> singleton to manage in-memory archive progress state.
/// </summary>
class EFCoreArchivingManager(IDomainEvents domainEvents, OperationsManager operationsManager)
class EFCoreArchivingManager(IDomainEvents domainEvents, OperationsManager operationsManager, ArchiveMetrics metrics)
{
InMemoryArchive GetOrCreate(ArchiveType archiveType, string requestId)
{
var id = InMemoryArchive.MakeId(requestId, archiveType);
if (!operationsManager.ArchiveOperations.TryGetValue(id, out var summary))
{
summary = new InMemoryArchive(requestId, archiveType, domainEvents);
summary = new InMemoryArchive(requestId, archiveType, domainEvents, metrics);
operationsManager.ArchiveOperations[id] = summary;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3,19 +3,20 @@ namespace ServiceControl.Persistence.EFCore.Implementation;
using ServiceControl.Infrastructure.DomainEvents;
using ServiceControl.Persistence.EFCore.Entities;
using ServiceControl.Recoverability;
using ServiceControl.Recoverability.Archiving.Metrics;

/// <summary>
/// EFCore equivalent of the RavenDB <see cref="UnarchivingManager"/>. Wraps the shared
/// <see cref="OperationsManager"/> singleton to manage in-memory unarchive progress state.
/// </summary>
class EFCoreUnarchivingManager(IDomainEvents domainEvents, OperationsManager operationsManager)
class EFCoreUnarchivingManager(IDomainEvents domainEvents, OperationsManager operationsManager, ArchiveMetrics metrics)
{
InMemoryUnarchive GetOrCreate(ArchiveType archiveType, string requestId)
{
var id = InMemoryUnarchive.MakeId(requestId, archiveType);
if (!operationsManager.UnarchiveOperations.TryGetValue(id, out var summary))
{
summary = new InMemoryUnarchive(requestId, archiveType, domainEvents);
summary = new InMemoryUnarchive(requestId, archiveType, domainEvents, metrics);
operationsManager.UnarchiveOperations[id] = summary;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ namespace ServiceControl.Persistence.EFCore.Implementation.Recoverability;
using ServiceControl.Persistence.EFCore.Entities;
using ServiceControl.Persistence.Recoverability;
using ServiceControl.Recoverability;
using ServiceControl.Recoverability.Archiving.Metrics;

public class MessageArchiver : IArchiveMessages
{
Expand All @@ -19,6 +20,7 @@ public MessageArchiver(
IDomainEvents domainEvents,
IMessageActionAuditLog auditLog,
TimeProvider timeProvider,
ArchiveMetrics metrics,
ILogger<MessageArchiver> logger
)
{
Expand All @@ -29,8 +31,8 @@ ILogger<MessageArchiver> logger
this.logger = logger;
this.domainEvents = domainEvents;

archivingManager = new EFCoreArchivingManager(domainEvents, operationsManager);
unarchivingManager = new EFCoreUnarchivingManager(domainEvents, operationsManager);
archivingManager = new EFCoreArchivingManager(domainEvents, operationsManager, metrics);
unarchivingManager = new EFCoreUnarchivingManager(domainEvents, operationsManager, metrics);
}

public async Task ArchiveAllInGroup(string groupId, AuditUser? initiatedBy = null, string? operationId = null, CancellationToken cancellationToken = default)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@ public RetentionMetrics(IMeterFactory meterFactory)
{
var meter = meterFactory.Create(MeterName, MeterVersion);

cycleDuration = meter.CreateHistogram<double>(
cycleDuration = meter.CreateHistogram(
CycleDurationInstrumentName,
unit: "seconds",
description: "Retention sweep pass duration in seconds",
Expand All @@ -26,7 +26,7 @@ public RetentionMetrics(IMeterFactory meterFactory)
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");
meter.CreateObservableGauge(ConsecutiveFailuresInstrumentName, ObserveConsecutiveFailures, description: "Consecutive retention sweep failures");
}

public RetentionCycleMetrics BeginCycle(RetentionEntity entity, CancellationToken cancellationToken = default) => new(this, entity, cancellationToken);
Expand Down Expand Up @@ -74,9 +74,6 @@ IEnumerable<Measurement<long>> ObserveConsecutiveFailures()

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

static readonly TagList[] EntityTags =
[
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ public ArchiveGroupPerMessageAuditTests() =>
{
services.AddSingleton<ArchiveAllInGroupHandler>();
services.AddSingleton<RetryingManager>();
services.AddSingleton(TestRetryMetrics.Create());
services.AddSingleton<IMessageActionAuditLog>(audit);
};

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ public ArchiveGroupTests() =>
{
services.AddSingleton<ArchiveAllInGroupHandler>();
services.AddSingleton<RetryingManager>();
services.AddSingleton(TestRetryMetrics.Create());
};

[Test]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ public UnarchiveGroupTests() =>
{
services.AddSingleton<UnarchiveAllInGroupHandler>();
services.AddSingleton<RetryingManager>();
services.AddSingleton(TestRetryMetrics.Create());
};

[Test]
Expand Down
170 changes: 170 additions & 0 deletions src/ServiceControl.Persistence.Tests/EFCore/ArchiveMetricsTests.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,170 @@
namespace ServiceControl.Persistence.Tests;

using System;
using System.Collections.Generic;
using System.Diagnostics.Metrics;
using System.Linq;
using System.Threading.Tasks;
using Microsoft.Extensions.DependencyInjection;
using Microsoft.Extensions.Time.Testing;
using NUnit.Framework;
using ServiceControl.Recoverability;
using ServiceControl.Recoverability.Archiving.Metrics;

/// <summary>
/// Instrument names and tag values are what dashboards and alerts are built on, so they are a
/// published contract and not an implementation detail.
/// </summary>
[TestFixture]
class ArchiveMetricsTests
{
[SetUp]
public void CreateMeterFactory()
{
provider = new ServiceCollection().AddMetrics().BuildServiceProvider();
fakeTime = new FakeTimeProvider(DateTimeOffset.UtcNow);
}

[TearDown]
public void DisposeMeterFactory() => provider.Dispose();

[Test]
public void The_meter_publishes_the_instruments_it_is_named_for()
{
var published = new List<string>();

using var listener = new MeterListener
{
InstrumentPublished = (instrument, _) =>
{
if (instrument.Meter.Name == ArchiveMetrics.MeterName && ReferenceEquals(instrument.Meter.Scope, MeterFactory))
{
published.Add(instrument.Name);
}
}
};

listener.Start();

_ = new ArchiveMetrics(MeterFactory, fakeTime);

Assert.That(published.Order(), Is.EqualTo(new[]
{
"sc.archive.batch_duration_seconds",
"sc.archive.messages_total",
"sc.archive.operation_duration_seconds",
"sc.archive.operations_in_progress"
}));
}

[Test]
public async Task An_archive_operation_records_batch_gaps_messages_and_total_duration()
{
var metrics = new ArchiveMetrics(MeterFactory, fakeTime);
using var recorded = new RecordedArchiveMetrics(MeterFactory);
var archive = new InMemoryArchive("group-1", ArchiveType.FailureGroup, new FakeDomainEvents(), metrics) { TotalNumberOfMessages = 1500 };

await archive.Start();
Assert.That(recorded.InProgress(ArchiveOperationKind.Archive, "started"), Is.EqualTo(1));

fakeTime.Advance(TimeSpan.FromSeconds(2));
await archive.BatchArchived(1000);
Assert.That(recorded.InProgress(ArchiveOperationKind.Archive, "started"), Is.Zero);
Assert.That(recorded.InProgress(ArchiveOperationKind.Archive, "progressing"), Is.EqualTo(1));

fakeTime.Advance(TimeSpan.FromSeconds(1));
await archive.BatchArchived(500);

await archive.FinalizeArchive();
Assert.That(recorded.InProgress(ArchiveOperationKind.Archive, "finalizing"), Is.EqualTo(1));

await archive.Complete();

using (Assert.EnterMultipleScope())
{
Assert.That(recorded.BatchDurations(ArchiveOperationKind.Archive).Select(batch => batch.Value), Is.EqualTo(new[] { 2.0, 1.0 }));
Assert.That(recorded.Messages(ArchiveOperationKind.Archive), Is.EqualTo(1500));
Assert.That(recorded.OperationDurations(ArchiveOperationKind.Archive).Single().Value, Is.EqualTo(3.0));
Assert.That(recorded.InProgress(ArchiveOperationKind.Archive, "started"), Is.Zero);
Assert.That(recorded.InProgress(ArchiveOperationKind.Archive, "progressing"), Is.Zero);
Assert.That(recorded.InProgress(ArchiveOperationKind.Archive, "finalizing"), Is.Zero);
}
}

[Test]
public async Task An_unarchive_operation_is_tagged_unarchive()
{
var metrics = new ArchiveMetrics(MeterFactory, fakeTime);
using var recorded = new RecordedArchiveMetrics(MeterFactory);
var unarchive = new InMemoryUnarchive("group-1", ArchiveType.FailureGroup, new FakeDomainEvents(), metrics) { TotalNumberOfMessages = 200 };

await unarchive.Start();
fakeTime.Advance(TimeSpan.FromSeconds(5));
await unarchive.BatchUnarchived(200);
await unarchive.FinalizeUnarchive();
await unarchive.Complete();

using (Assert.EnterMultipleScope())
{
Assert.That(recorded.BatchDurations(ArchiveOperationKind.Unarchive).Single().Value, Is.EqualTo(5.0));
Assert.That(recorded.Messages(ArchiveOperationKind.Unarchive), Is.EqualTo(200));
Assert.That(recorded.OperationDurations(ArchiveOperationKind.Unarchive).Single().Value, Is.EqualTo(5.0));
Assert.That(recorded.BatchDurations(ArchiveOperationKind.Archive), Is.Empty);
}
}

[Test]
public async Task An_operation_that_never_completes_stays_visible_on_the_gauge()
{
var metrics = new ArchiveMetrics(MeterFactory, fakeTime);
using var recorded = new RecordedArchiveMetrics(MeterFactory);
var archive = new InMemoryArchive("group-stuck", ArchiveType.FailureGroup, new FakeDomainEvents(), metrics) { TotalNumberOfMessages = 2000 };

await archive.Start();
await archive.BatchArchived(1000);

using (Assert.EnterMultipleScope())
{
Assert.That(recorded.InProgress(ArchiveOperationKind.Archive, "progressing"), Is.EqualTo(1));
Assert.That(recorded.OperationDurations(ArchiveOperationKind.Archive), Is.Empty);
}
}

[Test]
public async Task A_restarted_operation_counts_once_on_the_gauge()
{
var metrics = new ArchiveMetrics(MeterFactory, fakeTime);
using var recorded = new RecordedArchiveMetrics(MeterFactory);
var archive = new InMemoryArchive("group-1", ArchiveType.FailureGroup, new FakeDomainEvents(), metrics) { TotalNumberOfMessages = 10 };

await archive.Start();
await archive.BatchArchived(10);
await archive.Complete();

await archive.Start();

using (Assert.EnterMultipleScope())
{
Assert.That(recorded.InProgress(ArchiveOperationKind.Archive, "started"), Is.EqualTo(1));
Assert.That(recorded.InProgress(ArchiveOperationKind.Archive, "progressing"), Is.Zero);
}
}

[Test]
public async Task Without_metrics_the_state_machine_runs_unchanged()
{
var archive = new InMemoryArchive("group-1", ArchiveType.FailureGroup, new FakeDomainEvents()) { TotalNumberOfMessages = 10 };

await archive.Start();
await archive.BatchArchived(10);
await archive.FinalizeArchive();
await archive.Complete();

Assert.That(archive.ArchiveState, Is.EqualTo(ArchiveState.ArchiveCompleted));
}

IMeterFactory MeterFactory => provider.GetRequiredService<IMeterFactory>();

ServiceProvider provider;
FakeTimeProvider fakeTime;
}
Loading
Loading