diff --git a/docs/telemetry.md b/docs/telemetry.md index ff2bc9939d..e55820a22c 100644 --- a/docs/telemetry.md +++ b/docs/telemetry.md @@ -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`. diff --git a/src/Directory.Packages.props b/src/Directory.Packages.props index c4ee893c22..bc9c338e44 100644 --- a/src/Directory.Packages.props +++ b/src/Directory.Packages.props @@ -67,6 +67,9 @@ + + + diff --git a/src/ServiceControl.Audit/Auditing/Metrics/IngestionMetrics.cs b/src/ServiceControl.Audit/Auditing/Metrics/IngestionMetrics.cs index 038fefb7df..cdd6788eab 100644 --- a/src/ServiceControl.Audit/Auditing/Metrics/IngestionMetrics.cs +++ b/src/ServiceControl.Audit/Auditing/Metrics/IngestionMetrics.cs @@ -1,4 +1,4 @@ -namespace ServiceControl.Audit.Auditing.Metrics; +namespace ServiceControl.Audit.Auditing.Metrics; using System.Collections.Generic; using System.Diagnostics; @@ -23,7 +23,7 @@ public IngestionMetrics(IMeterFactory meterFactory) batchDuration = meter.CreateHistogram(BatchDurationInstrumentName, unit: "seconds", "Message batch processing duration in seconds"); ingestionDuration = meter.CreateHistogram(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($"{InstrumentPrefix}.failures_total", description: "Audit ingestion failure count"); } @@ -64,10 +64,6 @@ void RecordBatchOutcome(bool success) long consecutiveBatchFailures; readonly Histogram batchDuration; -#pragma warning disable IDE0052 - // this can be changed to Gauge once we can use the latest version of System.Diagnostics.DiagnosticSource - readonly ObservableGauge consecutiveBatchFailureGauge; -#pragma warning restore IDE0052 readonly Histogram ingestionDuration; readonly Counter failureCounter; diff --git a/src/ServiceControl.Persistence.EFCore/Abstractions/BasePersistence.cs b/src/ServiceControl.Persistence.EFCore/Abstractions/BasePersistence.cs index 6988f622d5..e88ea72bfc 100644 --- a/src/ServiceControl.Persistence.EFCore/Abstractions/BasePersistence.cs +++ b/src/ServiceControl.Persistence.EFCore/Abstractions/BasePersistence.cs @@ -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 { @@ -43,6 +44,7 @@ protected static void RegisterDataStores(IServiceCollection services, EFPersiste } services.AddSingleton(); + services.AddSingleton(); services.AddSingleton(); services.AddSingleton(); services.AddSingleton(); diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/EFCoreArchivingManager.cs b/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/EFCoreArchivingManager.cs index 35c6576c34..0602bfe2e0 100644 --- a/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/EFCoreArchivingManager.cs +++ b/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/EFCoreArchivingManager.cs @@ -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; /// /// EFCore equivalent of the RavenDB . Wraps the shared /// singleton to manage in-memory archive progress state. /// -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; } diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/EFCoreUnarchivingManager.cs b/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/EFCoreUnarchivingManager.cs index e9c41c9c37..276ae07788 100644 --- a/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/EFCoreUnarchivingManager.cs +++ b/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/EFCoreUnarchivingManager.cs @@ -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; /// /// EFCore equivalent of the RavenDB . Wraps the shared /// singleton to manage in-memory unarchive progress state. /// -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; } diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/MessageArchiver.cs b/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/MessageArchiver.cs index 58ca73ded2..579abd6211 100644 --- a/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/MessageArchiver.cs +++ b/src/ServiceControl.Persistence.EFCore/Implementation/Recoverability/MessageArchiver.cs @@ -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 { @@ -19,6 +20,7 @@ public MessageArchiver( IDomainEvents domainEvents, IMessageActionAuditLog auditLog, TimeProvider timeProvider, + ArchiveMetrics metrics, ILogger logger ) { @@ -29,8 +31,8 @@ ILogger 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) diff --git a/src/ServiceControl.Persistence.EFCore/Infrastructure/Metrics/RetentionMetrics.cs b/src/ServiceControl.Persistence.EFCore/Infrastructure/Metrics/RetentionMetrics.cs index 08f1f97f64..700a4358d6 100644 --- a/src/ServiceControl.Persistence.EFCore/Infrastructure/Metrics/RetentionMetrics.cs +++ b/src/ServiceControl.Persistence.EFCore/Infrastructure/Metrics/RetentionMetrics.cs @@ -16,7 +16,7 @@ public RetentionMetrics(IMeterFactory meterFactory) { var meter = meterFactory.Create(MeterName, MeterVersion); - cycleDuration = meter.CreateHistogram( + cycleDuration = meter.CreateHistogram( CycleDurationInstrumentName, unit: "seconds", description: "Retention sweep pass duration in seconds", @@ -26,7 +26,7 @@ public RetentionMetrics(IMeterFactory meterFactory) advice: new InstrumentAdvice { HistogramBucketBoundaries = [0.1, 0.5, 1, 5, 15, 60, 300, 900] }); rowsDeleted = meter.CreateCounter(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); @@ -74,9 +74,6 @@ IEnumerable> ObserveConsecutiveFailures() readonly Histogram cycleDuration; readonly Counter rowsDeleted; -#pragma warning disable IDE0052 - readonly ObservableGauge consecutiveFailureGauge; -#pragma warning restore IDE0052 static readonly TagList[] EntityTags = [ diff --git a/src/ServiceControl.Persistence.Tests.RavenDB/Archiving/ArchiveGroupPerMessageAuditTests.cs b/src/ServiceControl.Persistence.Tests.RavenDB/Archiving/ArchiveGroupPerMessageAuditTests.cs index 55fcc5f300..16e76129c1 100644 --- a/src/ServiceControl.Persistence.Tests.RavenDB/Archiving/ArchiveGroupPerMessageAuditTests.cs +++ b/src/ServiceControl.Persistence.Tests.RavenDB/Archiving/ArchiveGroupPerMessageAuditTests.cs @@ -21,6 +21,7 @@ public ArchiveGroupPerMessageAuditTests() => { services.AddSingleton(); services.AddSingleton(); + services.AddSingleton(TestRetryMetrics.Create()); services.AddSingleton(audit); }; diff --git a/src/ServiceControl.Persistence.Tests.RavenDB/Archiving/ArchiveGroupTests.cs b/src/ServiceControl.Persistence.Tests.RavenDB/Archiving/ArchiveGroupTests.cs index 8e1546f180..d7736132e0 100644 --- a/src/ServiceControl.Persistence.Tests.RavenDB/Archiving/ArchiveGroupTests.cs +++ b/src/ServiceControl.Persistence.Tests.RavenDB/Archiving/ArchiveGroupTests.cs @@ -15,6 +15,7 @@ public ArchiveGroupTests() => { services.AddSingleton(); services.AddSingleton(); + services.AddSingleton(TestRetryMetrics.Create()); }; [Test] diff --git a/src/ServiceControl.Persistence.Tests.RavenDB/Unarchiving/UnarchiveGroupTests.cs b/src/ServiceControl.Persistence.Tests.RavenDB/Unarchiving/UnarchiveGroupTests.cs index c40fd55ca3..9d3e8c06cf 100644 --- a/src/ServiceControl.Persistence.Tests.RavenDB/Unarchiving/UnarchiveGroupTests.cs +++ b/src/ServiceControl.Persistence.Tests.RavenDB/Unarchiving/UnarchiveGroupTests.cs @@ -15,6 +15,7 @@ public UnarchiveGroupTests() => { services.AddSingleton(); services.AddSingleton(); + services.AddSingleton(TestRetryMetrics.Create()); }; [Test] diff --git a/src/ServiceControl.Persistence.Tests/EFCore/ArchiveMetricsTests.cs b/src/ServiceControl.Persistence.Tests/EFCore/ArchiveMetricsTests.cs new file mode 100644 index 0000000000..37317fc884 --- /dev/null +++ b/src/ServiceControl.Persistence.Tests/EFCore/ArchiveMetricsTests.cs @@ -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; + +/// +/// Instrument names and tag values are what dashboards and alerts are built on, so they are a +/// published contract and not an implementation detail. +/// +[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(); + + 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(); + + ServiceProvider provider; + FakeTimeProvider fakeTime; +} diff --git a/src/ServiceControl.Persistence.Tests/EFCore/RecordedArchiveMetrics.cs b/src/ServiceControl.Persistence.Tests/EFCore/RecordedArchiveMetrics.cs new file mode 100644 index 0000000000..7a1a15dd9a --- /dev/null +++ b/src/ServiceControl.Persistence.Tests/EFCore/RecordedArchiveMetrics.cs @@ -0,0 +1,85 @@ +namespace ServiceControl.Persistence.Tests; + +using System; +using System.Collections.Generic; +using System.Diagnostics.Metrics; +using System.Linq; +using ServiceControl.Recoverability.Archiving.Metrics; + +/// +/// Collects everything the archive instruments record, for the meter belonging to one factory. +/// Every fixture in the run shares the meter name, so the factory is what tells these instruments +/// apart from the ones another test left behind. +/// +sealed class RecordedArchiveMetrics : IDisposable +{ + public RecordedArchiveMetrics(IMeterFactory meterFactory) + { + listener = new MeterListener + { + InstrumentPublished = (instrument, activeListener) => + { + if (instrument.Meter.Name == ArchiveMetrics.MeterName && ReferenceEquals(instrument.Meter.Scope, meterFactory)) + { + activeListener.EnableMeasurementEvents(instrument); + } + } + }; + + listener.SetMeasurementEventCallback((instrument, measurement, tags, _) => Add(instrument, measurement, tags)); + listener.SetMeasurementEventCallback((instrument, measurement, tags, _) => Add(instrument, measurement, tags)); + listener.Start(); + } + + public IReadOnlyList Of(string instrumentName, ArchiveOperationKind kind) + { + lock (measurements) + { + return + [ + .. measurements.Where(measurement => + measurement.InstrumentName == instrumentName && + Equals(measurement.Tags["archive.operation"], KindName(kind))) + ]; + } + } + + public IReadOnlyList BatchDurations(ArchiveOperationKind kind) => Of(ArchiveMetrics.BatchDurationInstrumentName, kind); + + public IReadOnlyList OperationDurations(ArchiveOperationKind kind) => Of(ArchiveMetrics.OperationDurationInstrumentName, kind); + + public double Messages(ArchiveOperationKind kind) => + Of(ArchiveMetrics.MessagesInstrumentName, kind).Sum(measurement => measurement.Value); + + public double InProgress(ArchiveOperationKind kind, string state) + { + lock (measurements) + { + measurements.RemoveAll(measurement => measurement.InstrumentName == ArchiveMetrics.OperationsInProgressInstrumentName); + } + + listener.RecordObservableInstruments(); + + return Of(ArchiveMetrics.OperationsInProgressInstrumentName, kind) + .Single(measurement => Equals(measurement.Tags["archive.state"], state)).Value; + } + + public void Dispose() => listener.Dispose(); + + static string KindName(ArchiveOperationKind kind) => kind == ArchiveOperationKind.Archive ? "archive" : "unarchive"; + + void Add(Instrument instrument, double value, ReadOnlySpan> tags) + { + var recorded = new Recorded(instrument.Name, value, tags.ToArray().ToDictionary(tag => tag.Key, tag => tag.Value)); + + lock (measurements) + { + measurements.Add(recorded); + } + } + + readonly List measurements = []; + readonly MeterListener listener; + + public sealed record Recorded(string InstrumentName, double Value, IReadOnlyDictionary Tags); +} diff --git a/src/ServiceControl.Persistence.Tests/Recoverability/TestRetryMetrics.cs b/src/ServiceControl.Persistence.Tests/Recoverability/TestRetryMetrics.cs new file mode 100644 index 0000000000..42a97162d9 --- /dev/null +++ b/src/ServiceControl.Persistence.Tests/Recoverability/TestRetryMetrics.cs @@ -0,0 +1,20 @@ +namespace ServiceControl.Persistence.Tests +{ + using System; + using System.Diagnostics.Metrics; + using ServiceControl.Recoverability.Retrying.Metrics; + + static class TestRetryMetrics + { + public static RetryMetrics Create() => new(new TestMeterFactory(), TimeProvider.System); + } + + sealed class TestMeterFactory : IMeterFactory + { + public Meter Create(MeterOptions options) => new(options.Name, options.Version, options.Tags, scope: this); + + public void Dispose() + { + } + } +} diff --git a/src/ServiceControl.Persistence.Tests/RetryStateTests.cs b/src/ServiceControl.Persistence.Tests/RetryStateTests.cs index 82bbef5e16..ea748e65d5 100644 --- a/src/ServiceControl.Persistence.Tests/RetryStateTests.cs +++ b/src/ServiceControl.Persistence.Tests/RetryStateTests.cs @@ -30,7 +30,7 @@ class RetryStateTests : PersistenceTestBase public async Task When_a_group_is_processed_it_is_set_to_the_Preparing_state() { var domainEvents = new FakeDomainEvents(); - var retryManager = new RetryingManager(domainEvents, NullLogger.Instance); + var retryManager = new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance); await CreateAFailedMessageAndMarkAsPartOfRetryBatch(retryManager, "Test-group", true, 1); var status = retryManager.GetStatusForRetryOperation("Test-group", RetryType.FailureGroup); @@ -42,7 +42,7 @@ public async Task When_a_group_is_processed_it_is_set_to_the_Preparing_state() public async Task When_a_group_is_prepared_and_SC_is_started_the_group_is_marked_as_failed() { var domainEvents = new FakeDomainEvents(); - var retryManager = new RetryingManager(domainEvents, NullLogger.Instance); + var retryManager = new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance); await CreateAFailedMessageAndMarkAsPartOfRetryBatch(retryManager, "Test-group", false, 1); @@ -83,7 +83,7 @@ public async Task When_the_dequeuer_is_created_then_the_error_address_is_cached( public async Task When_a_group_is_prepared_with_three_batches_and_SC_is_restarted_while_the_first_group_is_being_forwarded_then_the_count_still_matches() { var domainEvents = new FakeDomainEvents(); - var retryManager = new RetryingManager(domainEvents, NullLogger.Instance); + var retryManager = new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance); await CreateAFailedMessageAndMarkAsPartOfRetryBatch(retryManager, "Test-group", true, 2001); @@ -100,7 +100,7 @@ public async Task When_a_group_is_prepared_with_three_batches_and_SC_is_restarte new ErrorQueueNameCache(), new TestTransportCustomization()), retryManager, - new Lazy(() => sender), + TestRetryMetrics.Create(), new Lazy(() => sender), new RecordingMessageActionAuditLog(), NullLogger.Instance); @@ -110,7 +110,7 @@ public async Task When_a_group_is_prepared_with_three_batches_and_SC_is_restarte await processor.ProcessBatches(); // mark ready // Simulate SC restart - retryManager = new RetryingManager(domainEvents, NullLogger.Instance); + retryManager = new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance); var documentManager = new CustomRetryDocumentManager(false, RetryBatchStore, retryManager); @@ -128,7 +128,7 @@ public async Task When_a_group_is_prepared_with_three_batches_and_SC_is_restarte new ErrorQueueNameCache(), new TestTransportCustomization()), retryManager, - new Lazy(() => sender), + TestRetryMetrics.Create(), new Lazy(() => sender), new RecordingMessageActionAuditLog(), NullLogger.Instance); @@ -142,14 +142,14 @@ public async Task When_a_group_is_prepared_with_three_batches_and_SC_is_restarte public async Task When_a_group_is_forwarded_the_status_is_Completed() { var domainEvents = new FakeDomainEvents(); - var retryManager = new RetryingManager(domainEvents, NullLogger.Instance); + var retryManager = new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance); await CreateAFailedMessageAndMarkAsPartOfRetryBatch(retryManager, "Test-group", true, 1); var sender = new TestSender(); var returnToSender = new TestReturnToSenderDequeuer(new ReturnToSender(BodyStorage, NullLogger.Instance), FailedMessageLifecycleStore, domainEvents, "TestEndpoint", new ErrorQueueNameCache(), new TestTransportCustomization()); - var processor = new RetryProcessor(RetryStagingStore, MessageRedirectsDataStore, domainEvents, returnToSender, retryManager, new Lazy(() => sender), new RecordingMessageActionAuditLog(), NullLogger.Instance); + var processor = new RetryProcessor(RetryStagingStore, MessageRedirectsDataStore, domainEvents, returnToSender, retryManager, TestRetryMetrics.Create(), new Lazy(() => sender), new RecordingMessageActionAuditLog(), NullLogger.Instance); await processor.ProcessBatches(); // mark ready await processor.ProcessBatches(); @@ -197,7 +197,7 @@ public async Task When_the_batch_being_forwarded_is_gone_the_forwarding_pointer_ public async Task When_there_is_one_poison_message_it_is_removed_from_batch_and_the_status_is_Complete() { var domainEvents = new FakeDomainEvents(); - var retryManager = new RetryingManager(domainEvents, NullLogger.Instance); + var retryManager = new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance); await CreateAFailedMessageAndMarkAsPartOfRetryBatch(retryManager, "Test-group", true, "A", "B", "C"); @@ -214,7 +214,7 @@ public async Task When_there_is_one_poison_message_it_is_removed_from_batch_and_ }; var returnToSender = new TestReturnToSenderDequeuer(new ReturnToSender(BodyStorage, NullLogger.Instance), FailedMessageLifecycleStore, domainEvents, "TestEndpoint", new ErrorQueueNameCache(), new TestTransportCustomization()); - var processor = new RetryProcessor(RetryStagingStore, MessageRedirectsDataStore, domainEvents, returnToSender, retryManager, new Lazy(() => sender), new RecordingMessageActionAuditLog(), NullLogger.Instance); + var processor = new RetryProcessor(RetryStagingStore, MessageRedirectsDataStore, domainEvents, returnToSender, retryManager, TestRetryMetrics.Create(), new Lazy(() => sender), new RecordingMessageActionAuditLog(), NullLogger.Instance); bool c; do @@ -246,7 +246,7 @@ public async Task When_there_is_one_poison_message_it_is_removed_from_batch_and_ public async Task When_a_group_has_one_batch_out_of_two_forwarded_the_status_is_Forwarding() { var domainEvents = new FakeDomainEvents(); - var retryManager = new RetryingManager(domainEvents, NullLogger.Instance); + var retryManager = new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance); await CreateAFailedMessageAndMarkAsPartOfRetryBatch(retryManager, "Test-group", true, 1001); @@ -254,7 +254,7 @@ public async Task When_a_group_has_one_batch_out_of_two_forwarded_the_status_is_ var sender = new TestSender(); - var processor = new RetryProcessor(RetryStagingStore, MessageRedirectsDataStore, domainEvents, new TestReturnToSenderDequeuer(returnToSender, FailedMessageLifecycleStore, domainEvents, "TestEndpoint", new ErrorQueueNameCache(), new TestTransportCustomization()), retryManager, new Lazy(() => sender), new RecordingMessageActionAuditLog(), NullLogger.Instance); + var processor = new RetryProcessor(RetryStagingStore, MessageRedirectsDataStore, domainEvents, new TestReturnToSenderDequeuer(returnToSender, FailedMessageLifecycleStore, domainEvents, "TestEndpoint", new ErrorQueueNameCache(), new TestTransportCustomization()), retryManager, TestRetryMetrics.Create(), new Lazy(() => sender), new RecordingMessageActionAuditLog(), NullLogger.Instance); await CompleteDatabaseOperation(); @@ -269,7 +269,7 @@ public async Task When_a_group_has_one_batch_out_of_two_forwarded_the_status_is_ public async Task When_a_selection_is_staged_each_message_is_audited_as_a_batch() { var domainEvents = new FakeDomainEvents(); - var retryManager = new RetryingManager(domainEvents, NullLogger.Instance); + var retryManager = new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance); var user = new AuditUser("alice-sub", "Alice"); const string operationId = "op-sel"; var ids = new[] { "A", "B" }; @@ -301,7 +301,7 @@ public async Task When_a_selection_is_staged_each_message_is_audited_as_a_batch( var audit = new RecordingMessageActionAuditLog(); var sender = new TestSender(); var returnToSender = new TestReturnToSenderDequeuer(new ReturnToSender(BodyStorage, NullLogger.Instance), FailedMessageLifecycleStore, domainEvents, "TestEndpoint", new ErrorQueueNameCache(), new TestTransportCustomization()); - var processor = new RetryProcessor(RetryStagingStore, MessageRedirectsDataStore, domainEvents, returnToSender, retryManager, new Lazy(() => sender), audit, NullLogger.Instance); + var processor = new RetryProcessor(RetryStagingStore, MessageRedirectsDataStore, domainEvents, returnToSender, retryManager, TestRetryMetrics.Create(), new Lazy(() => sender), audit, NullLogger.Instance); await processor.ProcessBatches(); // stage await processor.ProcessBatches(); // forward @@ -319,7 +319,7 @@ public async Task When_a_selection_is_staged_each_message_is_audited_as_a_batch( public async Task When_a_group_is_staged_each_message_is_audited_with_the_initiating_user() { var domainEvents = new FakeDomainEvents(); - var retryManager = new RetryingManager(domainEvents, NullLogger.Instance); + var retryManager = new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance); var user = new AuditUser("alice-sub", "Alice"); const string operationId = "op-abc"; @@ -328,7 +328,7 @@ public async Task When_a_group_is_staged_each_message_is_audited_with_the_initia var audit = new RecordingMessageActionAuditLog(); var sender = new TestSender(); var returnToSender = new TestReturnToSenderDequeuer(new ReturnToSender(BodyStorage, NullLogger.Instance), FailedMessageLifecycleStore, domainEvents, "TestEndpoint", new ErrorQueueNameCache(), new TestTransportCustomization()); - var processor = new RetryProcessor(RetryStagingStore, MessageRedirectsDataStore, domainEvents, returnToSender, retryManager, new Lazy(() => sender), audit, NullLogger.Instance); + var processor = new RetryProcessor(RetryStagingStore, MessageRedirectsDataStore, domainEvents, returnToSender, retryManager, TestRetryMetrics.Create(), new Lazy(() => sender), audit, NullLogger.Instance); await processor.ProcessBatches(); // stage (emits per-message audit) await processor.ProcessBatches(); // forward @@ -363,8 +363,8 @@ RetryProcessor CreateProcessor(IDomainEvents domainEvents, TestSender sender) => MessageRedirectsDataStore, domainEvents, new TestReturnToSenderDequeuer(new ReturnToSender(BodyStorage, NullLogger.Instance), FailedMessageLifecycleStore, domainEvents, "TestEndpoint", new ErrorQueueNameCache(), new TestTransportCustomization()), - new RetryingManager(domainEvents, NullLogger.Instance), - new Lazy(() => sender), + new RetryingManager(domainEvents, TestRetryMetrics.Create(), NullLogger.Instance), + TestRetryMetrics.Create(), new Lazy(() => sender), new RecordingMessageActionAuditLog(), NullLogger.Instance); @@ -426,7 +426,7 @@ async Task CreateAFailedMessageAndMarkAsPartOfRetryBatch(RetryingManager retryMa class CustomRetriesGateway : RetriesGateway { public CustomRetriesGateway(bool progressToStaged, IRetryBatchStore store, RetryingManager retryManager) - : base(store, retryManager, NullLogger.Instance) + : base(store, retryManager, TestRetryMetrics.Create(), NullLogger.Instance) { this.progressToStaged = progressToStaged; } diff --git a/src/ServiceControl.Persistence/Recoverability/Archiving/InMemoryArchive.cs b/src/ServiceControl.Persistence/Recoverability/Archiving/InMemoryArchive.cs index 2453e6e6c6..90f933d506 100644 --- a/src/ServiceControl.Persistence/Recoverability/Archiving/InMemoryArchive.cs +++ b/src/ServiceControl.Persistence/Recoverability/Archiving/InMemoryArchive.cs @@ -4,14 +4,16 @@ using System.Threading; using System.Threading.Tasks; using Infrastructure.DomainEvents; + using ServiceControl.Recoverability.Archiving.Metrics; public class InMemoryArchive // in memory { - public InMemoryArchive(string requestId, ArchiveType archiveType, IDomainEvents domainEvents) + public InMemoryArchive(string requestId, ArchiveType archiveType, IDomainEvents domainEvents, ArchiveMetrics? metrics = null) { RequestId = requestId; ArchiveType = archiveType; this.domainEvents = domainEvents; + operationMetrics = metrics?.CreateOperation(ArchiveOperationKind.Archive); } public int TotalNumberOfMessages { get; set; } @@ -45,6 +47,7 @@ public Task Start(CancellationToken cancellationToken = default) { ArchiveState = ArchiveState.ArchiveStarted; CompletionTime = null; + operationMetrics?.Started(); return domainEvents.Raise(new ArchiveOperationStarting { @@ -61,6 +64,7 @@ public Task BatchArchived(int numberOfMessagesArchivedInBatch, CancellationToken NumberOfMessagesArchived += numberOfMessagesArchivedInBatch; CurrentBatch++; Last = DateTime.UtcNow; + operationMetrics?.BatchCompleted(numberOfMessagesArchivedInBatch); return domainEvents.Raise(new ArchiveOperationBatchCompleted { @@ -77,6 +81,7 @@ public Task FinalizeArchive(CancellationToken cancellationToken = default) ArchiveState = ArchiveState.ArchiveFinalizing; NumberOfMessagesArchived = TotalNumberOfMessages; Last = DateTime.UtcNow; + operationMetrics?.Finalizing(); return domainEvents.Raise(new ArchiveOperationFinalizing { @@ -94,6 +99,7 @@ public Task Complete(CancellationToken cancellationToken = default) NumberOfMessagesArchived = TotalNumberOfMessages; CompletionTime = DateTime.UtcNow; Last = DateTime.UtcNow; + operationMetrics?.Completed(); return domainEvents.Raise(new ArchiveOperationCompleted { @@ -113,5 +119,6 @@ public bool NeedsAcknowledgement() } IDomainEvents domainEvents; + readonly ArchiveOperationMetrics? operationMetrics; } } \ No newline at end of file diff --git a/src/ServiceControl.Persistence/Recoverability/Archiving/InMemoryUnarchive.cs b/src/ServiceControl.Persistence/Recoverability/Archiving/InMemoryUnarchive.cs index 0816fd1856..bbba1d009b 100644 --- a/src/ServiceControl.Persistence/Recoverability/Archiving/InMemoryUnarchive.cs +++ b/src/ServiceControl.Persistence/Recoverability/Archiving/InMemoryUnarchive.cs @@ -4,14 +4,16 @@ using System.Threading; using System.Threading.Tasks; using Infrastructure.DomainEvents; + using ServiceControl.Recoverability.Archiving.Metrics; public class InMemoryUnarchive // in memory { - public InMemoryUnarchive(string requestId, ArchiveType archiveType, IDomainEvents domainEvents) + public InMemoryUnarchive(string requestId, ArchiveType archiveType, IDomainEvents domainEvents, ArchiveMetrics? metrics = null) { RequestId = requestId; ArchiveType = archiveType; this.domainEvents = domainEvents; + operationMetrics = metrics?.CreateOperation(ArchiveOperationKind.Unarchive); } public int TotalNumberOfMessages { get; set; } @@ -45,6 +47,7 @@ public Task Start(CancellationToken cancellationToken = default) { ArchiveState = ArchiveState.ArchiveStarted; CompletionTime = null; + operationMetrics?.Started(); return domainEvents.Raise(new UnarchiveOperationStarting { @@ -61,6 +64,7 @@ public Task BatchUnarchived(int numberOfMessagesUnarchivedInBatch, CancellationT NumberOfMessagesUnarchived += numberOfMessagesUnarchivedInBatch; CurrentBatch++; Last = DateTime.UtcNow; + operationMetrics?.BatchCompleted(numberOfMessagesUnarchivedInBatch); return domainEvents.Raise(new UnarchiveOperationBatchCompleted { @@ -77,6 +81,7 @@ public Task FinalizeUnarchive(CancellationToken cancellationToken = default) ArchiveState = ArchiveState.ArchiveFinalizing; NumberOfMessagesUnarchived = TotalNumberOfMessages; Last = DateTime.UtcNow; + operationMetrics?.Finalizing(); return domainEvents.Raise(new UnarchiveOperationFinalizing { @@ -94,6 +99,7 @@ public Task Complete(CancellationToken cancellationToken = default) NumberOfMessagesUnarchived = TotalNumberOfMessages; CompletionTime = DateTime.UtcNow; Last = DateTime.UtcNow; + operationMetrics?.Completed(); return domainEvents.Raise(new UnarchiveOperationCompleted { @@ -113,5 +119,6 @@ internal bool NeedsAcknowledgement() } IDomainEvents domainEvents; + readonly ArchiveOperationMetrics? operationMetrics; } } \ No newline at end of file diff --git a/src/ServiceControl.Persistence/Recoverability/Archiving/Metrics/ArchiveMetrics.cs b/src/ServiceControl.Persistence/Recoverability/Archiving/Metrics/ArchiveMetrics.cs new file mode 100644 index 0000000000..3ddd925186 --- /dev/null +++ b/src/ServiceControl.Persistence/Recoverability/Archiving/Metrics/ArchiveMetrics.cs @@ -0,0 +1,111 @@ +namespace ServiceControl.Recoverability.Archiving.Metrics; + +using System; +using System.Collections.Generic; +using System.Diagnostics; +using System.Diagnostics.Metrics; +using System.Threading; +using ServiceControl.Infrastructure; + +/// +/// The sc.archive.* instruments. +/// +public class ArchiveMetrics +{ + public const string MeterName = ServiceControlMeters.Error; + + public static readonly string OperationDurationInstrumentName = $"{InstrumentPrefix}.operation_duration_seconds"; + public static readonly string BatchDurationInstrumentName = $"{InstrumentPrefix}.batch_duration_seconds"; + public static readonly string MessagesInstrumentName = $"{InstrumentPrefix}.messages_total"; + public static readonly string OperationsInProgressInstrumentName = $"{InstrumentPrefix}.operations_in_progress"; + + public ArchiveMetrics(IMeterFactory meterFactory, TimeProvider timeProvider) + { + this.timeProvider = timeProvider; + var meter = meterFactory.Create(MeterName, MeterVersion); + + operationDuration = meter.CreateHistogram( + OperationDurationInstrumentName, + unit: "seconds", + description: "Group archive or unarchive operation duration in seconds", + tags: null, + advice: new InstrumentAdvice { HistogramBucketBoundaries = [0.5, 1, 5, 15, 60, 300, 900, 3600] }); + + batchDuration = meter.CreateHistogram( + BatchDurationInstrumentName, + unit: "seconds", + description: "Archive batch duration in seconds", + tags: null, + advice: new InstrumentAdvice { HistogramBucketBoundaries = [0.05, 0.1, 0.5, 1, 5, 15, 60] }); + + messages = meter.CreateCounter(MessagesInstrumentName, description: "Messages archived and unarchived"); + meter.CreateObservableGauge(OperationsInProgressInstrumentName, ObserveInProgress, description: "Archive operations currently in progress"); + } + + public ArchiveOperationMetrics CreateOperation(ArchiveOperationKind kind) => new(this, kind); + + internal long GetTimestamp() => timeProvider.GetTimestamp(); + + internal void RecordBatch(ArchiveOperationKind kind, long sinceTimestamp, int messagesInBatch) + { + batchDuration.Record(timeProvider.GetElapsedTime(sinceTimestamp).TotalSeconds, KindTags[(int)kind]); + messages.Add(messagesInBatch, KindTags[(int)kind]); + } + + internal void RecordOperationCompleted(ArchiveOperationKind kind, long startTimestamp) => + operationDuration.Record(timeProvider.GetElapsedTime(startTimestamp).TotalSeconds, KindTags[(int)kind]); + + internal void RecordStateTransition(ArchiveOperationKind kind, ArchiveState? from, ArchiveState to) + { + if (from is { } previous && TryStateBucket(previous, out var fromBucket)) + { + Interlocked.Decrement(ref inProgress[Index(kind, fromBucket)]); + } + + if (TryStateBucket(to, out var toBucket)) + { + Interlocked.Increment(ref inProgress[Index(kind, toBucket)]); + } + } + + IEnumerable> ObserveInProgress() + { + for (var kind = 0; kind < KindTags.Length; kind++) + { + for (var state = 0; state < StateNames.Length; state++) + { + var tags = KindTags[kind]; + tags.Add("archive.state", StateNames[state]); + + yield return new Measurement(Volatile.Read(ref inProgress[Index((ArchiveOperationKind)kind, state)]), tags); + } + } + } + + static bool TryStateBucket(ArchiveState state, out int bucket) + { + bucket = (int)state; + return state != ArchiveState.ArchiveCompleted; + } + + static int Index(ArchiveOperationKind kind, int stateBucket) => ((int)kind * StateNames.Length) + stateBucket; + + // Indexed by the ArchiveState enum values; Completed is never in progress. + static readonly string[] StateNames = ["started", "progressing", "finalizing"]; + + static readonly TagList[] KindTags = + [ + new() { { "archive.operation", "archive" } }, + new() { { "archive.operation", "unarchive" } } + ]; + + readonly long[] inProgress = new long[KindTags.Length * StateNames.Length]; + + readonly TimeProvider timeProvider; + readonly Histogram operationDuration; + readonly Histogram batchDuration; + readonly Counter messages; + + const string MeterVersion = "0.1.0"; + const string InstrumentPrefix = "sc.archive"; +} diff --git a/src/ServiceControl.Persistence/Recoverability/Archiving/Metrics/ArchiveOperationKind.cs b/src/ServiceControl.Persistence/Recoverability/Archiving/Metrics/ArchiveOperationKind.cs new file mode 100644 index 0000000000..48d577ff09 --- /dev/null +++ b/src/ServiceControl.Persistence/Recoverability/Archiving/Metrics/ArchiveOperationKind.cs @@ -0,0 +1,7 @@ +namespace ServiceControl.Recoverability.Archiving.Metrics; + +public enum ArchiveOperationKind +{ + Archive, + Unarchive +} diff --git a/src/ServiceControl.Persistence/Recoverability/Archiving/Metrics/ArchiveOperationMetrics.cs b/src/ServiceControl.Persistence/Recoverability/Archiving/Metrics/ArchiveOperationMetrics.cs new file mode 100644 index 0000000000..4bdf2410d4 --- /dev/null +++ b/src/ServiceControl.Persistence/Recoverability/Archiving/Metrics/ArchiveOperationMetrics.cs @@ -0,0 +1,50 @@ +namespace ServiceControl.Recoverability.Archiving.Metrics; + +public sealed class ArchiveOperationMetrics +{ + internal ArchiveOperationMetrics(ArchiveMetrics metrics, ArchiveOperationKind kind) + { + this.metrics = metrics; + this.kind = kind; + } + + public void Started() + { + Transition(ArchiveState.ArchiveStarted); + operationStartTimestamp = metrics.GetTimestamp(); + lastBatchTimestamp = operationStartTimestamp; + } + + public void BatchCompleted(int messagesInBatch) + { + Transition(ArchiveState.ArchiveProgressing); + metrics.RecordBatch(kind, lastBatchTimestamp, messagesInBatch); + lastBatchTimestamp = metrics.GetTimestamp(); + } + + public void Finalizing() => Transition(ArchiveState.ArchiveFinalizing); + + public void Completed() + { + var started = trackedState is not null and not ArchiveState.ArchiveCompleted; + Transition(ArchiveState.ArchiveCompleted); + + if (started) + { + metrics.RecordOperationCompleted(kind, operationStartTimestamp); + } + } + + void Transition(ArchiveState to) + { + metrics.RecordStateTransition(kind, trackedState, to); + trackedState = to; + } + + ArchiveState? trackedState; + long operationStartTimestamp; + long lastBatchTimestamp; + + readonly ArchiveMetrics metrics; + readonly ArchiveOperationKind kind; +} diff --git a/src/ServiceControl.UnitTests/Recoverability/FailureGroupsRetryControllerAuditTests.cs b/src/ServiceControl.UnitTests/Recoverability/FailureGroupsRetryControllerAuditTests.cs index e332202078..af91ad48ab 100644 --- a/src/ServiceControl.UnitTests/Recoverability/FailureGroupsRetryControllerAuditTests.cs +++ b/src/ServiceControl.UnitTests/Recoverability/FailureGroupsRetryControllerAuditTests.cs @@ -1,4 +1,4 @@ -#nullable enable +#nullable enable namespace ServiceControl.UnitTests.Recoverability; using System; @@ -22,7 +22,7 @@ public async Task Emits_group_retry_operation_entry() var session = new TestableMessageSession(); var audit = new RecordingMessageActionAuditLog(); var user = new AuditUser("alice-sub", "Alice"); - var retryingManager = new RetryingManager(new FakeDomainEvents(), NullLogger.Instance); + var retryingManager = new RetryingManager(new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance); var controller = new FailureGroupsRetryController(session, retryingManager, new StubCurrentUserAccessor(user), audit); await controller.ArchiveGroupErrors("group-42"); @@ -41,7 +41,7 @@ public async Task Group_retry_skipped_as_already_in_progress_is_not_audited() { var session = new TestableMessageSession(); var audit = new RecordingMessageActionAuditLog(); - var retryingManager = new RetryingManager(new FakeDomainEvents(), NullLogger.Instance); + var retryingManager = new RetryingManager(new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance); await retryingManager.Preparing("group-42", RetryType.FailureGroup, totalNumberOfMessages: 10); var controller = new FailureGroupsRetryController(session, retryingManager, new StubCurrentUserAccessor(new AuditUser("alice-sub", "Alice")), audit); diff --git a/src/ServiceControl.UnitTests/Recoverability/Metrics/RecordedRetryMetrics.cs b/src/ServiceControl.UnitTests/Recoverability/Metrics/RecordedRetryMetrics.cs new file mode 100644 index 0000000000..222d17b700 --- /dev/null +++ b/src/ServiceControl.UnitTests/Recoverability/Metrics/RecordedRetryMetrics.cs @@ -0,0 +1,70 @@ +namespace ServiceControl.UnitTests.Recoverability.Metrics +{ + using System; + using System.Collections.Generic; + using System.Diagnostics.Metrics; + using System.Linq; + using ServiceControl.Recoverability.Retrying.Metrics; + + /// + /// Collects everything the retry instruments record, for the meter belonging to one factory. + /// Every fixture in the run shares the meter name, so the factory is what tells these + /// instruments apart from the ones another test left behind. + /// + sealed class RecordedRetryMetrics : IDisposable + { + public RecordedRetryMetrics(IMeterFactory meterFactory) + { + listener = new MeterListener + { + InstrumentPublished = (instrument, activeListener) => + { + if (instrument.Meter.Name == RetryMetrics.MeterName && ReferenceEquals(instrument.Meter.Scope, meterFactory)) + { + activeListener.EnableMeasurementEvents(instrument); + } + } + }; + + listener.SetMeasurementEventCallback((instrument, measurement, tags, _) => Add(instrument, measurement, tags)); + listener.SetMeasurementEventCallback((instrument, measurement, tags, _) => Add(instrument, measurement, tags)); + listener.Start(); + } + + public IReadOnlyList Of(string instrumentName) + { + lock (measurements) + { + return [.. measurements.Where(measurement => measurement.InstrumentName == instrumentName)]; + } + } + + public IReadOnlyList Observe(string instrumentName) + { + lock (measurements) + { + measurements.RemoveAll(measurement => measurement.InstrumentName == instrumentName); + } + + listener.RecordObservableInstruments(); + return Of(instrumentName); + } + + public void Dispose() => listener.Dispose(); + + void Add(Instrument instrument, double value, ReadOnlySpan> tags) + { + var recorded = new Recorded(instrument.Name, value, tags.ToArray().ToDictionary(tag => tag.Key, tag => tag.Value)); + + lock (measurements) + { + measurements.Add(recorded); + } + } + + readonly List measurements = []; + readonly MeterListener listener; + + public sealed record Recorded(string InstrumentName, double Value, IReadOnlyDictionary Tags); + } +} diff --git a/src/ServiceControl.UnitTests/Recoverability/Metrics/RetryMetricsTests.cs b/src/ServiceControl.UnitTests/Recoverability/Metrics/RetryMetricsTests.cs new file mode 100644 index 0000000000..eda89f6fb0 --- /dev/null +++ b/src/ServiceControl.UnitTests/Recoverability/Metrics/RetryMetricsTests.cs @@ -0,0 +1,222 @@ +namespace ServiceControl.UnitTests.Recoverability.Metrics +{ + using System; + using System.Collections.Generic; + using System.Diagnostics.Metrics; + using System.Linq; + using System.Threading; + using Microsoft.Extensions.DependencyInjection; + using NUnit.Framework; + using ServiceControl.Persistence; + using ServiceControl.Recoverability; + using ServiceControl.Recoverability.Retrying.Metrics; + + /// + /// Instrument names and tag values are what dashboards and alerts are built on, so they are a + /// published contract and not an implementation detail. + /// + [TestFixture] + class RetryMetricsTests + { + [SetUp] + public void CreateMeterFactory() => provider = new ServiceCollection().AddMetrics().BuildServiceProvider(); + + [TearDown] + public void DisposeMeterFactory() => provider.Dispose(); + + [Test] + public void The_meter_publishes_the_instruments_it_is_named_for() + { + var published = new List(); + + using var listener = new MeterListener + { + InstrumentPublished = (instrument, _) => + { + if (BelongsToThisTest(instrument)) + { + published.Add(instrument.Name); + } + } + }; + + listener.Start(); + + var metrics = new RetryMetrics(MeterFactory, TimeProvider.System); + metrics.ObserveOperationsInProgress(() => []); + metrics.ObservePendingBulkRequests(() => 0); + + Assert.That(published.Order(), Is.EqualTo(new[] + { + "sc.retry.forward_duration_seconds", + "sc.retry.messages_total", + "sc.retry.operation_duration_seconds", + "sc.retry.operations_in_progress", + "sc.retry.pending_bulk_requests", + "sc.retry.prepare_duration_seconds", + "sc.retry.stage_duration_seconds" + })); + } + + [Test] + public void Every_retry_type_maps_to_its_own_tag_value() + { + var metrics = new RetryMetrics(MeterFactory, TimeProvider.System); + using var recorded = new RecordedRetryMetrics(MeterFactory); + + foreach (var retryType in Enum.GetValues()) + { + metrics.RecordMessages(retryType, RetryMessageOutcome.Staged, 1); + } + + var tagValues = recorded.Of(RetryMetrics.MessagesInstrumentName).Select(measurement => measurement.Tags["retry.type"]); + + Assert.That(tagValues, Is.EquivalentTo(new[] { "unknown", "single", "group", "batch", "endpoint", "all", "queue" })); + } + + [Test] + public void Operation_completion_is_tagged_with_its_result() + { + var metrics = new RetryMetrics(MeterFactory, TimeProvider.System); + using var recorded = new RecordedRetryMetrics(MeterFactory); + + metrics.RecordOperationCompleted(RetryType.FailureGroup, metrics.GetTimestamp(), failed: false); + metrics.RecordOperationCompleted(RetryType.FailureGroup, metrics.GetTimestamp(), failed: true); + + var durations = recorded.Of(RetryMetrics.OperationDurationInstrumentName); + + Assert.That(durations, Has.Count.EqualTo(2)); + Assert.That(durations[0].Tags["result"], Is.EqualTo("success")); + Assert.That(durations[0].Tags["retry.type"], Is.EqualTo("group")); + Assert.That(durations[1].Tags["result"], Is.EqualTo("failed")); + } + + [Test] + public void A_scope_records_the_outcome_it_was_given() + { + var metrics = new RetryMetrics(MeterFactory, TimeProvider.System); + using var recorded = new RecordedRetryMetrics(MeterFactory); + + using (var staging = metrics.BeginStaging(RetryType.FailureGroup)) + { + staging.Complete(); + } + + using (var staging = metrics.BeginStaging(RetryType.FailureGroup)) + { + staging.Empty(); + } + + using (metrics.BeginStaging(RetryType.FailureGroup)) + { + } + + var results = recorded.Of(RetryMetrics.StageDurationInstrumentName).Select(measurement => measurement.Tags["result"]); + + Assert.That(results, Is.EqualTo(new[] { "success", "empty", "failed" })); + } + + [Test] + public void A_scope_cut_short_by_shutdown_is_recorded_as_cancelled() + { + var metrics = new RetryMetrics(MeterFactory, TimeProvider.System); + using var recorded = new RecordedRetryMetrics(MeterFactory); + using var shutdown = new CancellationTokenSource(); + + using (metrics.BeginPreparation(RetryType.All, shutdown.Token)) + { + shutdown.Cancel(); + } + + Assert.That(recorded.Of(RetryMetrics.PrepareDurationInstrumentName).Single().Tags["result"], Is.EqualTo("cancelled")); + } + + [Test] + public void A_scope_that_finished_before_shutdown_is_still_a_success() + { + var metrics = new RetryMetrics(MeterFactory, TimeProvider.System); + using var recorded = new RecordedRetryMetrics(MeterFactory); + using var shutdown = new CancellationTokenSource(); + + using (var preparation = metrics.BeginPreparation(RetryType.All, shutdown.Token)) + { + preparation.Complete(); + shutdown.Cancel(); + } + + Assert.That(recorded.Of(RetryMetrics.PrepareDurationInstrumentName).Single().Tags["result"], Is.EqualTo("success")); + } + + [Test] + public void Forwarding_is_tagged_with_its_mode() + { + var metrics = new RetryMetrics(MeterFactory, TimeProvider.System); + using var recorded = new RecordedRetryMetrics(MeterFactory); + + using (var forwarding = metrics.BeginForwarding(RetryType.FailureGroup, recoveringFromPrematureShutdown: false)) + { + forwarding.Complete(); + } + + using (metrics.BeginForwarding(RetryType.FailureGroup, recoveringFromPrematureShutdown: true)) + { + } + + var forwarded = recorded.Of(RetryMetrics.ForwardDurationInstrumentName); + + Assert.That(forwarded.Select(measurement => measurement.Tags["mode"]), Is.EqualTo(new[] { "counting", "timeout" })); + Assert.That(forwarded.Select(measurement => measurement.Tags["result"]), Is.EqualTo(new[] { "success", "failed" })); + } + + [Test] + public void The_gauge_counts_operations_by_type_and_state_and_excludes_completed() + { + var metrics = new RetryMetrics(MeterFactory, TimeProvider.System); + using var recorded = new RecordedRetryMetrics(MeterFactory); + + metrics.ObserveOperationsInProgress(() => + [ + (RetryType.FailureGroup, RetryState.Preparing), + (RetryType.FailureGroup, RetryState.Preparing), + (RetryType.SingleMessage, RetryState.Forwarding), + (RetryType.All, RetryState.Waiting), + (RetryType.All, RetryState.Completed) + ]); + + var observed = recorded.Observe(RetryMetrics.OperationsInProgressInstrumentName) + .ToDictionary(measurement => (measurement.Tags["retry.type"], measurement.Tags["retry.state"]), measurement => measurement.Value); + + Assert.That(observed, Is.EqualTo(new Dictionary<(object, object), double> + { + { ("group", "preparing"), 2 }, + { ("single", "forwarding"), 1 }, + { ("all", "waiting"), 1 } + })); + } + + [Test] + public void The_bulk_request_gauge_reads_the_queue_depth() + { + var metrics = new RetryMetrics(MeterFactory, TimeProvider.System); + using var recorded = new RecordedRetryMetrics(MeterFactory); + + var depth = 3; + metrics.ObservePendingBulkRequests(() => depth); + + Assert.That(recorded.Observe(RetryMetrics.PendingBulkRequestsInstrumentName).Single().Value, Is.EqualTo(3)); + + depth = 0; + + Assert.That(recorded.Observe(RetryMetrics.PendingBulkRequestsInstrumentName).Single().Value, Is.Zero); + } + + // Every fixture in the run shares the meter name, so the factory is what tells the instruments + // created here apart from the ones another test left behind. + bool BelongsToThisTest(Instrument instrument) => + instrument.Meter.Name == RetryMetrics.MeterName && ReferenceEquals(instrument.Meter.Scope, MeterFactory); + + IMeterFactory MeterFactory => provider.GetRequiredService(); + + ServiceProvider provider; + } +} diff --git a/src/ServiceControl.UnitTests/Recoverability/RetryOperationTests.cs b/src/ServiceControl.UnitTests/Recoverability/RetryOperationTests.cs index 0254f5a6dd..95c99d9428 100644 --- a/src/ServiceControl.UnitTests/Recoverability/RetryOperationTests.cs +++ b/src/ServiceControl.UnitTests/Recoverability/RetryOperationTests.cs @@ -6,6 +6,7 @@ using NUnit.Framework; using ServiceControl.Persistence; using ServiceControl.Recoverability; + using ServiceControl.UnitTests.Recoverability; [TestFixture] public class RetryOperationTests @@ -13,7 +14,7 @@ public class RetryOperationTests [Test] public async Task Wait_should_set_wait_state() { - var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), NullLogger.Instance); + var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance); await summary.Wait(DateTime.UtcNow, "FailureGroup1"); using (Assert.EnterMultipleScope()) { @@ -29,7 +30,7 @@ public async Task Wait_should_set_wait_state() [Test] public void Fail_should_set_failed() { - var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), NullLogger.Instance); + var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance); summary.Fail(); Assert.That(summary.Failed, Is.True); } @@ -37,7 +38,7 @@ public void Fail_should_set_failed() [Test] public async Task Prepare_should_set_prepare_state() { - var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), NullLogger.Instance); + var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance); await summary.Prepare(1000); using (Assert.EnterMultipleScope()) { @@ -50,7 +51,7 @@ public async Task Prepare_should_set_prepare_state() [Test] public async Task Prepared_batch_should_set_prepare_state() { - var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), NullLogger.Instance); + var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance); await summary.Prepare(1000); await summary.PrepareBatch(1000); using (Assert.EnterMultipleScope()) @@ -64,7 +65,7 @@ public async Task Prepared_batch_should_set_prepare_state() [Test] public async Task Forwarding_should_set_forwarding_state() { - var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), NullLogger.Instance); + var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance); await summary.Prepare(1000); await summary.PrepareBatch(1000); await summary.Forwarding(); @@ -80,7 +81,7 @@ public async Task Forwarding_should_set_forwarding_state() [Test] public async Task Batch_forwarded_should_set_forwarding_state() { - var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), NullLogger.Instance); + var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance); await summary.Prepare(1000); await summary.PrepareBatch(1000); await summary.Forwarding(); @@ -98,7 +99,7 @@ public async Task Batch_forwarded_should_set_forwarding_state() public async Task Should_raise_domain_events() { var domainEvents = new FakeDomainEvents(); - var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, domainEvents, NullLogger.Instance); + var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, domainEvents, TestRetryMetrics.Create(), NullLogger.Instance); await summary.Prepare(1000); await summary.PrepareBatch(1000); await summary.Forwarding(); @@ -117,7 +118,7 @@ public async Task Should_raise_domain_events() [Test] public async Task Batch_forwarded_all_forwarded_should_set_completed_state() { - var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), NullLogger.Instance); + var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance); await summary.Prepare(1000); await summary.PrepareBatch(1000); await summary.Forwarding(); @@ -134,7 +135,7 @@ public async Task Batch_forwarded_all_forwarded_should_set_completed_state() [Test] public async Task Skip_should_set_update_skipped_messages() { - var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), NullLogger.Instance); + var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance); await summary.Wait(DateTime.UtcNow); await summary.Prepare(2000); await summary.PrepareBatch(1000); @@ -150,7 +151,7 @@ public async Task Skip_should_set_update_skipped_messages() [Test] public async Task Skip_should_complete_when_all_skipped() { - var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), NullLogger.Instance); + var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance); await summary.Wait(DateTime.UtcNow); await summary.Prepare(1000); await summary.PrepareBatch(1000); @@ -166,7 +167,7 @@ public async Task Skip_should_complete_when_all_skipped() [Test] public async Task Skip_and_forward_combination_should_complete_when_done() { - var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), NullLogger.Instance); + var summary = new InMemoryRetry("abc123", RetryType.FailureGroup, new FakeDomainEvents(), TestRetryMetrics.Create(), NullLogger.Instance); await summary.Wait(DateTime.UtcNow); await summary.Prepare(2000); await summary.PrepareBatch(1000); diff --git a/src/ServiceControl.UnitTests/Recoverability/TestRetryMetrics.cs b/src/ServiceControl.UnitTests/Recoverability/TestRetryMetrics.cs new file mode 100644 index 0000000000..8b8c2b6f71 --- /dev/null +++ b/src/ServiceControl.UnitTests/Recoverability/TestRetryMetrics.cs @@ -0,0 +1,20 @@ +namespace ServiceControl.UnitTests.Recoverability +{ + using System; + using System.Diagnostics.Metrics; + using ServiceControl.Recoverability.Retrying.Metrics; + + static class TestRetryMetrics + { + public static RetryMetrics Create() => new(new TestMeterFactory(), TimeProvider.System); + } + + sealed class TestMeterFactory : IMeterFactory + { + public Meter Create(MeterOptions options) => new(options.Name, options.Version, options.Tags, scope: this); + + public void Dispose() + { + } + } +} diff --git a/src/ServiceControl/HostApplicationBuilderExtensions.cs b/src/ServiceControl/HostApplicationBuilderExtensions.cs index 109e37b8a7..b511f16ac2 100644 --- a/src/ServiceControl/HostApplicationBuilderExtensions.cs +++ b/src/ServiceControl/HostApplicationBuilderExtensions.cs @@ -15,11 +15,13 @@ using global::ServiceControl.Infrastructure.WebApi; using global::ServiceControl.Notifications.Email; using global::ServiceControl.Operations.Metrics; + using global::ServiceControl.Recoverability.Retrying.Metrics; using global::ServiceControl.Persistence; using global::ServiceControl.Transports; using Licensing; using Microsoft.AspNetCore.HttpLogging; using Microsoft.Extensions.DependencyInjection; + using Microsoft.Extensions.DependencyInjection.Extensions; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Hosting.WindowsServices; using Microsoft.Extensions.Logging; @@ -99,7 +101,7 @@ public static void AddServiceControl(this IHostApplicationBuilder hostBuilder, S services.AddPersistence(settings); services.AddMetrics(settings.PrintMetrics); - hostBuilder.AddIngestionMetrics(settings); + hostBuilder.AddTelemetry(settings); services.AddServiceControlHealthChecks(); if (settings.ErrorIngestionOnly) @@ -149,9 +151,11 @@ public static void AddServiceControlInstallers(this IHostApplicationBuilder host persistence.AddInstaller(hostApplicationBuilder.Services); } - public static void AddIngestionMetrics(this IHostApplicationBuilder hostBuilder, Settings settings) + public static void AddTelemetry(this IHostApplicationBuilder hostBuilder, Settings settings) { + hostBuilder.Services.TryAddSingleton(TimeProvider.System); hostBuilder.Services.AddSingleton(); + hostBuilder.Services.AddSingleton(); var otlpEndpoint = OtlpEndpoint.Read(hostBuilder.Configuration); @@ -168,6 +172,9 @@ public static void AddIngestionMetrics(this IHostApplicationBuilder hostBuilder, .WithMetrics(metrics => { metrics.AddIngestionMetrics(); + metrics.AddAspNetCoreInstrumentation(); + metrics.AddHttpClientInstrumentation(); + metrics.AddRuntimeInstrumentation(); metrics.AddOtlpExporter(); if (Debugger.IsAttached) diff --git a/src/ServiceControl/Operations/Metrics/IngestionMetrics.cs b/src/ServiceControl/Operations/Metrics/IngestionMetrics.cs index e395919797..ffbc5d63f7 100644 --- a/src/ServiceControl/Operations/Metrics/IngestionMetrics.cs +++ b/src/ServiceControl/Operations/Metrics/IngestionMetrics.cs @@ -1,4 +1,4 @@ -namespace ServiceControl.Operations.Metrics; +namespace ServiceControl.Operations.Metrics; using System.Collections.Generic; using System.Diagnostics; @@ -23,7 +23,7 @@ public IngestionMetrics(IMeterFactory meterFactory) batchDuration = meter.CreateHistogram(BatchDurationInstrumentName, unit: "seconds", description: "Message batch processing duration in seconds"); ingestionDuration = meter.CreateHistogram(MessageDurationInstrumentName, unit: "seconds", description: "Error message processing duration in seconds"); storageDuration = meter.CreateHistogram(StorageDurationInstrumentName, unit: "seconds", description: "Error ingestion batch storage write duration in seconds"); - consecutiveBatchFailureGauge = meter.CreateObservableGauge($"{InstrumentPrefix}.consecutive_batch_failures_total", () => Volatile.Read(ref consecutiveBatchFailures), description: "Consecutive error ingestion batch failures"); + meter.CreateObservableGauge($"{InstrumentPrefix}.consecutive_batch_failures_total", () => Volatile.Read(ref consecutiveBatchFailures), description: "Consecutive error ingestion batch failures"); failureCounter = meter.CreateCounter($"{InstrumentPrefix}.failures_total", description: "Error ingestion failure count"); } @@ -67,10 +67,6 @@ void RecordBatchOutcome(bool success) long consecutiveBatchFailures; readonly Histogram batchDuration; -#pragma warning disable IDE0052 - // this can be changed to Gauge once we can use the latest version of System.Diagnostics.DiagnosticSource - readonly ObservableGauge consecutiveBatchFailureGauge; -#pragma warning restore IDE0052 readonly Histogram ingestionDuration; readonly Histogram storageDuration; readonly Counter failureCounter; diff --git a/src/ServiceControl/Recoverability/Retrying/InMemoryRetry.cs b/src/ServiceControl/Recoverability/Retrying/InMemoryRetry.cs index 9e3160c9e8..98d203bd8c 100644 --- a/src/ServiceControl/Recoverability/Retrying/InMemoryRetry.cs +++ b/src/ServiceControl/Recoverability/Retrying/InMemoryRetry.cs @@ -6,18 +6,22 @@ using Infrastructure.DomainEvents; using Microsoft.Extensions.Logging; using ServiceControl.Persistence; + using ServiceControl.Recoverability.Retrying.Metrics; public class InMemoryRetry { - public InMemoryRetry(string requestId, RetryType retryType, IDomainEvents domainEvents, ILogger logger) + public InMemoryRetry(string requestId, RetryType retryType, IDomainEvents domainEvents, RetryMetrics metrics, ILogger logger) { RequestId = requestId; - this.retryType = retryType; + RetryType = retryType; this.domainEvents = domainEvents; + this.metrics = metrics; this.logger = logger; + operationStartTimestamp = metrics.GetTimestamp(); } public string RequestId { get; } + public RetryType RetryType { get; } public int TotalNumberOfMessages { get; private set; } public int NumberOfMessagesPrepared { get; private set; } public int NumberOfMessagesForwarded { get; private set; } @@ -39,6 +43,7 @@ public static string MakeOperationId(string requestId, RetryType retryType) public Task Wait(DateTime started, string originator = null, string classifier = null, DateTime? last = null, CancellationToken cancellationToken = default) { RetryState = RetryState.Waiting; + operationStartTimestamp = metrics.GetTimestamp(); NumberOfMessagesPrepared = 0; NumberOfMessagesForwarded = 0; TotalNumberOfMessages = 0; @@ -53,7 +58,7 @@ public Task Wait(DateTime started, string originator = null, string classifier = return domainEvents.Raise(new RetryOperationWaiting { RequestId = RequestId, - RetryType = retryType, + RetryType = RetryType, Progress = GetProgress(), StartTime = Started }, cancellationToken); @@ -66,6 +71,12 @@ public void Fail() public Task Prepare(int totalNumberOfMessages, CancellationToken cancellationToken = default) { + // A completed operation being prepared again is a new run that never went through Wait. + if (RetryState == RetryState.Completed) + { + operationStartTimestamp = metrics.GetTimestamp(); + } + RetryState = RetryState.Preparing; TotalNumberOfMessages = totalNumberOfMessages; NumberOfMessagesForwarded = 0; @@ -74,7 +85,7 @@ public Task Prepare(int totalNumberOfMessages, CancellationToken cancellationTok return domainEvents.Raise(new RetryOperationPreparing { RequestId = RequestId, - RetryType = retryType, + RetryType = RetryType, TotalNumberOfMessages = TotalNumberOfMessages, Progress = GetProgress(), IsFailed = Failed, @@ -89,7 +100,7 @@ public Task PrepareBatch(int numberOfMessagesPrepared, CancellationToken cancell return domainEvents.Raise(new RetryOperationPreparing { RequestId = RequestId, - RetryType = retryType, + RetryType = RetryType, TotalNumberOfMessages = TotalNumberOfMessages, Progress = GetProgress(), IsFailed = Failed, @@ -114,7 +125,7 @@ public Task Forwarding(CancellationToken cancellationToken = default) return domainEvents.Raise(new RetryOperationForwarding { RequestId = RequestId, - RetryType = retryType, + RetryType = RetryType, TotalNumberOfMessages = TotalNumberOfMessages, Progress = GetProgress(), IsFailed = Failed, @@ -125,11 +136,12 @@ public Task Forwarding(CancellationToken cancellationToken = default) public async Task BatchForwarded(int numberOfMessagesForwarded, CancellationToken cancellationToken = default) { NumberOfMessagesForwarded += numberOfMessagesForwarded; + metrics.RecordMessages(RetryType, RetryMessageOutcome.Forwarded, numberOfMessagesForwarded); await domainEvents.Raise(new RetryMessagesForwarded { RequestId = RequestId, - RetryType = retryType, + RetryType = RetryType, TotalNumberOfMessages = TotalNumberOfMessages, Progress = GetProgress(), IsFailed = Failed, @@ -142,6 +154,7 @@ await domainEvents.Raise(new RetryMessagesForwarded public Task Skip(int numberOfMessagesSkipped, CancellationToken cancellationToken = default) { NumberOfMessagesSkipped += numberOfMessagesSkipped; + metrics.RecordMessages(RetryType, RetryMessageOutcome.Skipped, numberOfMessagesSkipped); return CheckForCompletion(cancellationToken); } @@ -154,11 +167,12 @@ async Task CheckForCompletion(CancellationToken cancellationToken) RetryState = RetryState.Completed; CompletionTime = DateTime.UtcNow; + metrics.RecordOperationCompleted(RetryType, operationStartTimestamp, Failed); await domainEvents.Raise(new RetryOperationCompleted { RequestId = RequestId, - RetryType = retryType, + RetryType = RetryType, Failed = Failed, Progress = GetProgress(), StartTime = Started, @@ -169,7 +183,7 @@ await domainEvents.Raise(new RetryOperationCompleted Classifier = Classifier }, cancellationToken); - if (retryType == RetryType.FailureGroup) + if (RetryType == RetryType.FailureGroup) { await domainEvents.Raise(new MessagesSubmittedForRetry { @@ -206,9 +220,11 @@ public bool IsInProgress() return RetryState is not RetryState.Completed and not RetryState.Waiting; } - readonly RetryType retryType; + + long operationStartTimestamp; IDomainEvents domainEvents; + readonly RetryMetrics metrics; readonly ILogger logger; } } \ No newline at end of file diff --git a/src/ServiceControl/Recoverability/Retrying/Metrics/RetryDurationScope.cs b/src/ServiceControl/Recoverability/Retrying/Metrics/RetryDurationScope.cs new file mode 100644 index 0000000000..4ff10e7b29 --- /dev/null +++ b/src/ServiceControl/Recoverability/Retrying/Metrics/RetryDurationScope.cs @@ -0,0 +1,61 @@ +namespace ServiceControl.Recoverability.Retrying.Metrics; + +using System; +using System.Diagnostics; +using System.Diagnostics.Metrics; +using System.Threading; + +/// +/// One timed stretch of the retry pipeline. A stretch that finished counts as its outcome even if +/// shutdown has since been requested; one that shutdown cut short is recorded as cancelled rather +/// than as a failure. +/// +public sealed class RetryDurationScope : IDisposable +{ + internal RetryDurationScope(Histogram histogram, TimeProvider timeProvider, TagList tags, CancellationToken cancellationToken) + { + this.histogram = histogram; + this.timeProvider = timeProvider; + this.tags = tags; + this.cancellationToken = cancellationToken; + start = timeProvider.GetTimestamp(); + } + + public void Complete() => Finish(ScopeOutcome.Success); + + public void Empty() => Finish(ScopeOutcome.Empty); + + public void Dispose() + { + var result = outcome ?? (cancellationToken.IsCancellationRequested ? ScopeOutcome.Cancelled : ScopeOutcome.Failed); + tags.Add("result", ResultTag(result)); + + histogram.Record(timeProvider.GetElapsedTime(start, end == 0 ? timeProvider.GetTimestamp() : end).TotalSeconds, tags); + } + + // The work is over once it finishes, so the clock stops here rather than wherever the scope + // happens to be disposed. + void Finish(ScopeOutcome result) + { + end = timeProvider.GetTimestamp(); + outcome = result; + } + + static string ResultTag(ScopeOutcome outcome) => outcome switch + { + ScopeOutcome.Success => "success", + ScopeOutcome.Empty => "empty", + ScopeOutcome.Cancelled => "cancelled", + ScopeOutcome.Failed => "failed", + _ => "failed" + }; + + ScopeOutcome? outcome; + long end; + TagList tags; + + readonly Histogram histogram; + readonly TimeProvider timeProvider; + readonly CancellationToken cancellationToken; + readonly long start; +} diff --git a/src/ServiceControl/Recoverability/Retrying/Metrics/RetryMessageOutcome.cs b/src/ServiceControl/Recoverability/Retrying/Metrics/RetryMessageOutcome.cs new file mode 100644 index 0000000000..c53182a74a --- /dev/null +++ b/src/ServiceControl/Recoverability/Retrying/Metrics/RetryMessageOutcome.cs @@ -0,0 +1,10 @@ +namespace ServiceControl.Recoverability.Retrying.Metrics; + +public enum RetryMessageOutcome +{ + Staged, + Forwarded, + Skipped, + StagingRetried, + Abandoned +} diff --git a/src/ServiceControl/Recoverability/Retrying/Metrics/RetryMetrics.cs b/src/ServiceControl/Recoverability/Retrying/Metrics/RetryMetrics.cs new file mode 100644 index 0000000000..510dcfbd9c --- /dev/null +++ b/src/ServiceControl/Recoverability/Retrying/Metrics/RetryMetrics.cs @@ -0,0 +1,129 @@ +namespace ServiceControl.Recoverability.Retrying.Metrics; + +using System; +using System.Collections.Generic; +using System.Diagnostics; +using System.Diagnostics.Metrics; +using System.Linq; +using System.Threading; +using ServiceControl.Infrastructure; +using ServiceControl.Persistence; + +public class RetryMetrics +{ + public const string MeterName = ServiceControlMeters.Error; + + public static readonly string OperationDurationInstrumentName = $"{InstrumentPrefix}.operation_duration_seconds"; + public static readonly string PrepareDurationInstrumentName = $"{InstrumentPrefix}.prepare_duration_seconds"; + public static readonly string StageDurationInstrumentName = $"{InstrumentPrefix}.stage_duration_seconds"; + public static readonly string ForwardDurationInstrumentName = $"{InstrumentPrefix}.forward_duration_seconds"; + public static readonly string MessagesInstrumentName = $"{InstrumentPrefix}.messages_total"; + public static readonly string OperationsInProgressInstrumentName = $"{InstrumentPrefix}.operations_in_progress"; + public static readonly string PendingBulkRequestsInstrumentName = $"{InstrumentPrefix}.pending_bulk_requests"; + + public RetryMetrics(IMeterFactory meterFactory, TimeProvider timeProvider) + { + this.timeProvider = timeProvider; + meter = meterFactory.Create(MeterName, MeterVersion); + + operationDuration = meter.CreateHistogram( + OperationDurationInstrumentName, + unit: "seconds", + description: "Retry operation duration from request to completion in seconds", + tags: null, + advice: new InstrumentAdvice { HistogramBucketBoundaries = [1, 5, 15, 60, 300, 900, 1800, 3600] }); + + prepareDuration = meter.CreateHistogram( + PrepareDurationInstrumentName, + unit: "seconds", + description: "Retry preparation duration in seconds, covering the store scan and batch creation", + tags: null, + advice: new InstrumentAdvice { HistogramBucketBoundaries = [0.1, 0.5, 1, 5, 15, 60, 300, 900] }); + + stageDuration = meter.CreateHistogram( + StageDurationInstrumentName, + unit: "seconds", + description: "Retry batch staging duration in seconds", + tags: null, + advice: new InstrumentAdvice { HistogramBucketBoundaries = [0.05, 0.1, 0.5, 1, 5, 15, 60] }); + + forwardDuration = meter.CreateHistogram( + ForwardDurationInstrumentName, + unit: "seconds", + description: "Retry batch forwarding duration in seconds", + tags: null, + advice: new InstrumentAdvice { HistogramBucketBoundaries = [0.5, 1, 5, 15, 45, 60, 300, 900] }); + + messages = meter.CreateCounter(MessagesInstrumentName, description: "Messages moved through the retry pipeline"); + } + + public long GetTimestamp() => timeProvider.GetTimestamp(); + + public void RecordOperationCompleted(RetryType retryType, long startTimestamp, bool failed) + { + var tags = RetryTypeTags(retryType); + tags.Add("result", failed ? "failed" : "success"); + + operationDuration.Record(timeProvider.GetElapsedTime(startTimestamp).TotalSeconds, tags); + } + + public void RecordMessages(RetryType retryType, RetryMessageOutcome outcome, long count) + { + var tags = RetryTypeTags(retryType); + tags.Add("result", OutcomeNames[(int)outcome]); + + messages.Add(count, tags); + } + + public RetryDurationScope BeginPreparation(RetryType retryType, CancellationToken cancellationToken = default) => + new(prepareDuration, timeProvider, RetryTypeTags(retryType), cancellationToken); + + public RetryDurationScope BeginStaging(RetryType retryType, CancellationToken cancellationToken = default) => + new(stageDuration, timeProvider, RetryTypeTags(retryType), cancellationToken); + + public RetryDurationScope BeginForwarding(RetryType retryType, bool recoveringFromPrematureShutdown, CancellationToken cancellationToken = default) + { + var tags = RetryTypeTags(retryType); + tags.Add("mode", recoveringFromPrematureShutdown ? "timeout" : "counting"); + + return new RetryDurationScope(forwardDuration, timeProvider, tags, cancellationToken); + } + + public void ObserveOperationsInProgress(Func> operations) => + meter.CreateObservableGauge(OperationsInProgressInstrumentName, () => MeasureInProgress(operations()), description: "Retry operations currently in progress"); + + public void ObservePendingBulkRequests(Func queueDepth) => + meter.CreateObservableGauge(PendingBulkRequestsInstrumentName, () => (long)queueDepth(), description: "Bulk retry requests queued for preparation"); + + static IEnumerable> MeasureInProgress(IEnumerable<(RetryType RetryType, RetryState RetryState)> operations) + { + foreach (var group in operations.Where(operation => operation.RetryState != RetryState.Completed).GroupBy(operation => operation)) + { + var tags = RetryTypeTags(group.Key.RetryType); + tags.Add("retry.state", StateNames[(int)group.Key.RetryState]); + + yield return new Measurement(group.Count(), tags); + } + } + + static TagList RetryTypeTags(RetryType retryType) => new() { { "retry.type", RetryTypeNames[(int)retryType] } }; + + // Indexed by the RetryType enum values. + static readonly string[] RetryTypeNames = ["unknown", "single", "group", "batch", "endpoint", "all", "queue"]; + + // Indexed by the RetryState enum values; Completed is never emitted. + static readonly string[] StateNames = ["waiting", "preparing", "forwarding", "completed"]; + + static readonly string[] OutcomeNames = ["staged", "forwarded", "skipped", "staging_retried", "abandoned"]; + + readonly Meter meter; + readonly TimeProvider timeProvider; + readonly Histogram operationDuration; + readonly Histogram prepareDuration; + readonly Histogram stageDuration; + readonly Histogram forwardDuration; + readonly Counter messages; + + const string MeterVersion = "0.1.0"; + const string InstrumentPrefix = "sc.retry"; +} diff --git a/src/ServiceControl/Recoverability/Retrying/Metrics/ScopeOutcome.cs b/src/ServiceControl/Recoverability/Retrying/Metrics/ScopeOutcome.cs new file mode 100644 index 0000000000..d84b7b9b45 --- /dev/null +++ b/src/ServiceControl/Recoverability/Retrying/Metrics/ScopeOutcome.cs @@ -0,0 +1,9 @@ +namespace ServiceControl.Recoverability.Retrying.Metrics; + +enum ScopeOutcome +{ + Success, + Empty, + Cancelled, + Failed +} diff --git a/src/ServiceControl/Recoverability/Retrying/RetriesGateway.cs b/src/ServiceControl/Recoverability/Retrying/RetriesGateway.cs index dbb53b04c5..3b183defc0 100644 --- a/src/ServiceControl/Recoverability/Retrying/RetriesGateway.cs +++ b/src/ServiceControl/Recoverability/Retrying/RetriesGateway.cs @@ -11,14 +11,18 @@ namespace ServiceControl.Recoverability using MessageFailures; using Microsoft.Extensions.Logging; using ServiceControl.Persistence; + using ServiceControl.Recoverability.Retrying.Metrics; class RetriesGateway { - public RetriesGateway(IRetryBatchStore store, RetryingManager operationManager, ILogger logger) + public RetriesGateway(IRetryBatchStore store, RetryingManager operationManager, RetryMetrics metrics, ILogger logger) { this.store = store; this.operationManager = operationManager; + this.metrics = metrics; this.logger = logger; + + metrics.ObservePendingBulkRequests(() => bulkRequests.Count); } public async Task StartRetryForSingleMessage(string uniqueMessageId, AuditUser? initiatedBy = null, string operationId = null, CancellationToken cancellationToken = default) @@ -29,9 +33,13 @@ public async Task StartRetryForSingleMessage(string uniqueMessageId, AuditUser? var retryType = RetryType.SingleMessage; var numberOfMessages = 1; + using var preparation = metrics.BeginPreparation(retryType, cancellationToken); + await operationManager.Preparing(requestId, retryType, numberOfMessages, cancellationToken); await AssignMessagesToBatch(requestId, retryType, new[] { uniqueMessageId }, DateTime.UtcNow, cancellationToken, initiatedBy: initiatedBy, operationId: operationId); await operationManager.PreparedBatch(requestId, retryType, numberOfMessages, cancellationToken); + + preparation.Complete(); } public async Task StartRetryForMessageSelection(string[] uniqueMessageIds, AuditUser? initiatedBy = null, string operationId = null, CancellationToken cancellationToken = default) @@ -42,9 +50,13 @@ public async Task StartRetryForMessageSelection(string[] uniqueMessageIds, Audit var retryType = RetryType.MultipleMessages; var numberOfMessages = uniqueMessageIds.Length; + using var preparation = metrics.BeginPreparation(retryType, cancellationToken); + await operationManager.Preparing(requestId, retryType, numberOfMessages, cancellationToken); await AssignMessagesToBatch(requestId, retryType, uniqueMessageIds, DateTime.UtcNow, cancellationToken, initiatedBy: initiatedBy, operationId: operationId); await operationManager.PreparedBatch(requestId, retryType, numberOfMessages, cancellationToken); + + preparation.Complete(); } async Task AssignMessagesToBatch(string requestId, RetryType retryType, string[] messageIds, DateTime startTime, CancellationToken cancellationToken, DateTime? last = null, string originator = null, string batchName = null, string classifier = null, AuditUser? initiatedBy = null, string operationId = null) @@ -85,6 +97,8 @@ public async Task ProcessNextBulkRetry(CancellationToken cancellationToken async Task ProcessRequest(BulkRetryRequest request, CancellationToken cancellationToken) { + using var preparation = metrics.BeginPreparation(request.RetryType, cancellationToken); + var (batches, latestAttempt) = await request.GetRequestedBatches(store, cancellationToken); var totalMessages = batches.Sum(b => b.Length); @@ -102,6 +116,8 @@ async Task ProcessRequest(BulkRetryRequest request, CancellationToken cancellati await operationManager.PreparedBatch(request.RequestId, request.RetryType, numberOfMessagesAdded, cancellationToken); } } + + preparation.Complete(); } static string GetBatchName(int pageNum, int totalPages, string context) @@ -143,6 +159,7 @@ public void EnqueueRetryForFailureGroup(RetryForFailureGroup item) readonly IRetryBatchStore store; readonly RetryingManager operationManager; + readonly RetryMetrics metrics; readonly ConcurrentQueue bulkRequests = new ConcurrentQueue(); const int BatchSize = 1000; diff --git a/src/ServiceControl/Recoverability/Retrying/RetryProcessor.cs b/src/ServiceControl/Recoverability/Retrying/RetryProcessor.cs index 663d70c3be..3f3e11836a 100644 --- a/src/ServiceControl/Recoverability/Retrying/RetryProcessor.cs +++ b/src/ServiceControl/Recoverability/Retrying/RetryProcessor.cs @@ -14,6 +14,7 @@ namespace ServiceControl.Recoverability using NServiceBus.Transport; using Persistence.MessageRedirects; using ServiceControl.Persistence; + using ServiceControl.Recoverability.Retrying.Metrics; class RetryProcessor { @@ -23,6 +24,7 @@ public RetryProcessor( IDomainEvents domainEvents, ReturnToSenderDequeuer returnToSender, RetryingManager retryingManager, + RetryMetrics metrics, Lazy messageDispatcher, IMessageActionAuditLog auditLog, ILogger logger) @@ -31,6 +33,7 @@ public RetryProcessor( this.redirectsStore = redirectsStore; this.returnToSender = returnToSender; this.retryingManager = retryingManager; + this.metrics = metrics; this.domainEvents = domainEvents; this.messageDispatcher = messageDispatcher; this.auditLog = auditLog; @@ -118,6 +121,8 @@ async Task ForwardCurrentBatch(CancellationToken cancellationToken) async Task Forward(RetryBatch forwardingBatch, CancellationToken cancellationToken) { + using var forwarding = metrics.BeginForwarding(forwardingBatch.RetryType, isRecoveringFromPrematureShutdown, cancellationToken); + var messageCount = forwardingBatch.MessageCount; await retryingManager.Forwarding(forwardingBatch.RequestId, forwardingBatch.RetryType, cancellationToken); @@ -144,6 +149,7 @@ async Task Forward(RetryBatch forwardingBatch, CancellationToken cancellationTok } logger.LogInformation("Done forwarding batch {ForwardingBatchId}", forwardingBatch.Id); + forwarding.Complete(); } static Predicate IsPartOfStagedBatch(string stagingId) @@ -157,6 +163,8 @@ static Predicate IsPartOfStagedBatch(string stagingId) async Task Stage(RetryBatch stagingBatch, CancellationToken cancellationToken) { + using var staging = metrics.BeginStaging(stagingBatch.RetryType, cancellationToken); + var stagingId = Guid.NewGuid().ToString(); var messagesToStage = await store.GetMessagesToStage(stagingBatch.Id, cancellationToken); @@ -165,6 +173,7 @@ async Task Stage(RetryBatch stagingBatch, CancellationToken cancellationTok { logger.LogInformation("Retry batch {RetryBatchId} cancelled as it has no messages left to stage", stagingBatch.Id); await store.DiscardBatch(stagingBatch.Id, cancellationToken); + staging.Empty(); return 0; } @@ -180,7 +189,9 @@ async Task Stage(RetryBatch stagingBatch, CancellationToken cancellationTok transportOperations[current++] = ToTransportOperation(messageToStage, stagingId); } - await TryDispatch(stagingBatch.Id, transportOperations, messagesToStage, stageAttemptsById, previousAttemptFailed, cancellationToken); + await TryDispatch(stagingBatch.Id, stagingBatch.RetryType, transportOperations, messagesToStage, stageAttemptsById, previousAttemptFailed, cancellationToken); + + metrics.RecordMessages(stagingBatch.RetryType, RetryMessageOutcome.Staged, messagesToStage.Length); AuditStagedMessages(stagingBatch, messagesToStage); @@ -198,6 +209,7 @@ await domainEvents.Raise(new MessagesSubmittedForRetry await store.MarkBatchAsForwarding(stagingBatch.Id, stagingId, [.. stageAttemptsById.Keys], cancellationToken); logger.LogInformation("Retry batch {RetryBatchId} staged with Staging Id {StagingId} and {RetryFailureCount} matching failure retries", stagingBatch.Id, stagingId, messagesToStage.Length); + staging.Complete(); return messagesToStage.Length; } @@ -234,24 +246,24 @@ void AuditStagedMessages(RetryBatch stagingBatch, IReadOnlyCollection messages, + Task TryDispatch(string batchId, RetryType retryType, TransportOperation[] transportOperations, IReadOnlyCollection messages, IReadOnlyDictionary stageAttemptsById, bool previousAttemptFailed, CancellationToken cancellationToken) { - return previousAttemptFailed ? ConcurrentDispatchToTransport(transportOperations, stageAttemptsById, cancellationToken) : - BatchDispatchToTransport(batchId, transportOperations, messages, cancellationToken); + return previousAttemptFailed ? ConcurrentDispatchToTransport(retryType, transportOperations, stageAttemptsById, cancellationToken) : + BatchDispatchToTransport(batchId, retryType, transportOperations, messages, cancellationToken); } - Task ConcurrentDispatchToTransport(IReadOnlyCollection transportOperations, IReadOnlyDictionary stageAttemptsById, CancellationToken cancellationToken) + Task ConcurrentDispatchToTransport(RetryType retryType, IReadOnlyCollection transportOperations, IReadOnlyDictionary stageAttemptsById, CancellationToken cancellationToken) { var tasks = new List(transportOperations.Count); foreach (var transportOperation in transportOperations) { - tasks.Add(TryStageMessage(transportOperation, stageAttemptsById, cancellationToken)); + tasks.Add(TryStageMessage(transportOperation, retryType, stageAttemptsById, cancellationToken)); } return Task.WhenAll(tasks); } - async Task BatchDispatchToTransport(string batchId, TransportOperation[] transportOperations, IReadOnlyCollection messages, CancellationToken cancellationToken) + async Task BatchDispatchToTransport(string batchId, RetryType retryType, TransportOperation[] transportOperations, IReadOnlyCollection messages, CancellationToken cancellationToken) { try { @@ -267,11 +279,13 @@ async Task BatchDispatchToTransport(string batchId, TransportOperation[] transpo await store.RecordStagingFailure([.. messages.Select(message => message.UniqueMessageId)], cancellationToken); + metrics.RecordMessages(retryType, RetryMessageOutcome.StagingRetried, messages.Count); + throw new RetryStagingException(e); } } - async Task TryStageMessage(TransportOperation transportOperation, IReadOnlyDictionary stageAttemptsById, CancellationToken cancellationToken) + async Task TryStageMessage(TransportOperation transportOperation, RetryType retryType, IReadOnlyDictionary stageAttemptsById, CancellationToken cancellationToken) { var uniqueMessageId = transportOperation.Message.Headers["ServiceControl.Retry.UniqueMessageId"]; @@ -292,6 +306,8 @@ async Task TryStageMessage(TransportOperation transportOperation, IReadOnlyDicti logger.LogWarning(e, "Attempt {StagingRetryAttempt} of {StagingRetryLimit} to stage a retry message {RetryMessageId} failed", incrementedAttempts, MaxStagingAttempts, uniqueMessageId); await store.IncrementStagingAttempts(uniqueMessageId, cancellationToken); + + metrics.RecordMessages(retryType, RetryMessageOutcome.StagingRetried, 1); } else { @@ -303,6 +319,8 @@ await domainEvents.Raise(new MessageFailedInStaging { UniqueMessageId = uniqueMessageId }, cancellationToken); + + metrics.RecordMessages(retryType, RetryMessageOutcome.Abandoned, 1); } throw new RetryStagingException(e); @@ -338,6 +356,7 @@ TransportOperation ToTransportOperation(StagingMessage message, string stagingId readonly IMessageRedirectsDataStore redirectsStore; readonly ReturnToSenderDequeuer returnToSender; readonly RetryingManager retryingManager; + readonly RetryMetrics metrics; readonly Lazy messageDispatcher; readonly IMessageActionAuditLog auditLog; IReadOnlyList redirects; diff --git a/src/ServiceControl/Recoverability/Retrying/RetryingManager.cs b/src/ServiceControl/Recoverability/Retrying/RetryingManager.cs index 3b9c91ac18..e766a6f982 100644 --- a/src/ServiceControl/Recoverability/Retrying/RetryingManager.cs +++ b/src/ServiceControl/Recoverability/Retrying/RetryingManager.cs @@ -8,13 +8,17 @@ using Infrastructure.DomainEvents; using Microsoft.Extensions.Logging; using ServiceControl.Persistence; + using ServiceControl.Recoverability.Retrying.Metrics; public class RetryingManager { - public RetryingManager(IDomainEvents domainEvents, ILogger logger) + public RetryingManager(IDomainEvents domainEvents, RetryMetrics metrics, ILogger logger) { this.domainEvents = domainEvents; + this.metrics = metrics; this.logger = logger; + + metrics.ObserveOperationsInProgress(() => retryOperations.Values.Select(operation => (operation.RetryType, operation.RetryState))); } public Task Wait(string requestId, RetryType retryType, DateTime started, string originator = null, string classifier = null, DateTime? last = null, CancellationToken cancellationToken = default) @@ -93,7 +97,7 @@ InMemoryRetry GetOrCreate(RetryType retryType, string requestId) ArgumentException.ThrowIfNullOrWhiteSpace(requestId); var key = InMemoryRetry.MakeOperationId(requestId, retryType); - return retryOperations.GetOrAdd(key, _ => new InMemoryRetry(requestId, retryType, domainEvents, logger)); + return retryOperations.GetOrAdd(key, _ => new InMemoryRetry(requestId, retryType, domainEvents, metrics, logger)); } public InMemoryRetry GetStatusForRetryOperation(string requestId, RetryType retryType) @@ -104,6 +108,7 @@ public InMemoryRetry GetStatusForRetryOperation(string requestId, RetryType retr } IDomainEvents domainEvents; + readonly RetryMetrics metrics; readonly ILogger logger; ConcurrentDictionary retryOperations = new ConcurrentDictionary(); } diff --git a/src/ServiceControl/ServiceControl.csproj b/src/ServiceControl/ServiceControl.csproj index 4401575e13..7afd8f88cf 100644 --- a/src/ServiceControl/ServiceControl.csproj +++ b/src/ServiceControl/ServiceControl.csproj @@ -38,6 +38,9 @@ + + +