From 9ef144066cc0b3f4ab47de2a1509a761a558c68a Mon Sep 17 00:00:00 2001 From: John Simons Date: Thu, 13 Aug 2026 18:57:18 +1000 Subject: [PATCH] Refactor and re-enable processing attempt acceptance tests Update the tests to focus on the most recent processing attempt, making them compatible with SQL and PostgreSQL persistence which store fewer historical attempts than RavenDB. The specific attempt trimming logic for RavenDB has been moved to a dedicated persistence test. --- ...eControl.AcceptanceTests.PostgreSql.csproj | 4 - ...ceControl.AcceptanceTests.SqlServer.csproj | 4 - .../When_a_messages_fails_multiple_times.cs | 39 ++++----- ..._errors_with_same_uniqueid_are_imported.cs | 83 ++++++++----------- .../ProcessingAttemptTrimmingTests.cs | 46 ++++++++++ 5 files changed, 98 insertions(+), 78 deletions(-) create mode 100644 src/ServiceControl.Persistence.Tests.RavenDB/Recoverability/ProcessingAttemptTrimmingTests.cs diff --git a/src/ServiceControl.AcceptanceTests.PostgreSql/ServiceControl.AcceptanceTests.PostgreSql.csproj b/src/ServiceControl.AcceptanceTests.PostgreSql/ServiceControl.AcceptanceTests.PostgreSql.csproj index c873aa8eb7..5d1bd8efeb 100644 --- a/src/ServiceControl.AcceptanceTests.PostgreSql/ServiceControl.AcceptanceTests.PostgreSql.csproj +++ b/src/ServiceControl.AcceptanceTests.PostgreSql/ServiceControl.AcceptanceTests.PostgreSql.csproj @@ -52,10 +52,6 @@ - - - - diff --git a/src/ServiceControl.AcceptanceTests.SqlServer/ServiceControl.AcceptanceTests.SqlServer.csproj b/src/ServiceControl.AcceptanceTests.SqlServer/ServiceControl.AcceptanceTests.SqlServer.csproj index 6677b0b3dd..d8ffdd841b 100644 --- a/src/ServiceControl.AcceptanceTests.SqlServer/ServiceControl.AcceptanceTests.SqlServer.csproj +++ b/src/ServiceControl.AcceptanceTests.SqlServer/ServiceControl.AcceptanceTests.SqlServer.csproj @@ -52,10 +52,6 @@ - - - - diff --git a/src/ServiceControl.AcceptanceTests/Recoverability/MessageFailures/When_a_messages_fails_multiple_times.cs b/src/ServiceControl.AcceptanceTests/Recoverability/MessageFailures/When_a_messages_fails_multiple_times.cs index 434f6c776d..1bb57db166 100644 --- a/src/ServiceControl.AcceptanceTests/Recoverability/MessageFailures/When_a_messages_fails_multiple_times.cs +++ b/src/ServiceControl.AcceptanceTests/Recoverability/MessageFailures/When_a_messages_fails_multiple_times.cs @@ -1,4 +1,4 @@ -namespace ServiceControl.AcceptanceTests.Recoverability +namespace ServiceControl.AcceptanceTests.Recoverability { using NServiceBus.AcceptanceTesting; using NServiceBus.Routing; @@ -17,15 +17,14 @@ class When_a_messages_fails_multiple_times : AcceptanceTest { const int NumberOfFailedAttempts = 20; - const int MaximalNumberOfStoredFailedAttempts = 10; - const string AttemptIdHeaderKey = "testing.failed_attempt_no"; + const string AttemptNumberHeaderKey = "testing.failed_attempt_no"; [Test] - public async Task Should_store_only_the_latest_processing_attempts() + public async Task Should_report_the_most_recent_attempt_last() { FailedMessage result = null; - var context = await Define() + await Define() .WithEndpoint() .Done(async c => { @@ -34,19 +33,23 @@ public async Task Should_store_only_the_latest_processing_attempts() return false; } - result = await this.TryGet($"/api/errors/{c.UniqueMessageId}"); + result = await this.TryGet( + $"/api/errors/{c.UniqueMessageId}", + m => LatestAttemptNumber(m) == NumberOfFailedAttempts.ToString()); - var failureTimes = result?.ProcessingAttempts.Select(pa => pa.Headers["NServiceBus.TimeOfFailure"]).ToArray() ?? []; - - return failureTimes.SequenceEqual([.. c.LatestFailureTimes]); + return result != null; }) .Run(); + + Assert.That(LatestAttemptNumber(result), Is.EqualTo(NumberOfFailedAttempts.ToString())); } + static string LatestAttemptNumber(FailedMessage message) => + message.ProcessingAttempts[^1].Headers.GetValueOrDefault(AttemptNumberHeaderKey); + class TestContext : ScenarioContext { public string UniqueMessageId { get; set; } - public List LatestFailureTimes { get; set; } = []; } class AnEndpoint : EndpointConfigurationBuilder @@ -66,30 +69,20 @@ protected override TransportOperations CreateMessage(TestContext context) var transportOperations = Enumerable.Range(0, NumberOfFailedAttempts) .Select(i => { - var timeOfFailure = DateTimeOffsetHelper.ToWireFormattedString(earliestTimeOfFailure.Add(TimeSpan.FromMinutes(i))); - var headers = new Dictionary { [Headers.MessageId] = messageId, [Headers.EnclosedMessageTypes] = typeof(MyMessage).FullName, ["NServiceBus.FailedQ"] = endpointName, ["$.diagnostics.hostid"] = Guid.NewGuid().ToString(), - ["NServiceBus.TimeOfFailure"] = timeOfFailure, - - [AttemptIdHeaderKey] = (i + 1).ToString() + ["NServiceBus.TimeOfFailure"] = DateTimeOffsetHelper.ToWireFormattedString(earliestTimeOfFailure.Add(TimeSpan.FromMinutes(i))), + [AttemptNumberHeaderKey] = (i + 1).ToString() }; - context.LatestFailureTimes.Add(timeOfFailure); - return new TransportOperation(new OutgoingMessage(messageId, headers, Array.Empty()), new UnicastAddressTag("error")); }) .ToArray(); - context.LatestFailureTimes = context.LatestFailureTimes - .Skip(context.LatestFailureTimes.Count - MaximalNumberOfStoredFailedAttempts) - .Take(MaximalNumberOfStoredFailedAttempts) - .ToList(); - return new TransportOperations(transportOperations); } } @@ -97,4 +90,4 @@ protected override TransportOperations CreateMessage(TestContext context) class MyMessage : ICommand; } } -} \ No newline at end of file +} diff --git a/src/ServiceControl.AcceptanceTests/Recoverability/MessageFailures/When_errors_with_same_uniqueid_are_imported.cs b/src/ServiceControl.AcceptanceTests/Recoverability/MessageFailures/When_errors_with_same_uniqueid_are_imported.cs index dc725480c0..af9b323433 100644 --- a/src/ServiceControl.AcceptanceTests/Recoverability/MessageFailures/When_errors_with_same_uniqueid_are_imported.cs +++ b/src/ServiceControl.AcceptanceTests/Recoverability/MessageFailures/When_errors_with_same_uniqueid_are_imported.cs @@ -1,9 +1,8 @@ -namespace ServiceControl.AcceptanceTests.Recoverability.MessageFailures +namespace ServiceControl.AcceptanceTests.Recoverability.MessageFailures { using System; using System.Collections.Concurrent; using System.Collections.Generic; - using System.Linq; using System.Threading.Tasks; using AcceptanceTesting; using AcceptanceTesting.EndpointTemplates; @@ -19,37 +18,32 @@ class When_errors_with_same_uniqueid_are_imported : AcceptanceTest { + const int NumberOfDuplicates = 10; + [Test] public async Task The_import_should_deduplicate_on_TimeOfFailure() { var criticalErrorExecuted = false; - SetSettings = settings => settings.MaximumConcurrencyLevel = 10; - CustomizeHostBuilder = builder => builder.Services.AddSingleton(); - CustomConfiguration = config => - { - config.DefineCriticalErrorAction((_, _) => + SetSettings = settings => settings.MaximumConcurrencyLevel = NumberOfDuplicates; + CustomizeHostBuilder = builder => builder.Services.AddSingleton(); + CustomConfiguration = config => config.DefineCriticalErrorAction((_, _) => { criticalErrorExecuted = true; return Task.CompletedTask; }); - }; FailedMessage failure = null; var context = await Define() .WithEndpoint() .Done(async c => { - if (c.UniqueId == null) + if (c.UniqueId == null || c.IngestedCount < NumberOfDuplicates) { return false; } - var result = await this.TryGet($"/api/errors/{c.UniqueId}", m => - { - Console.WriteLine("Processing attempts: " + m.ProcessingAttempts.Count); - return m.ProcessingAttempts.Count == 2; - }); + var result = await this.TryGet($"/api/errors/{c.UniqueId}"); failure = result; return criticalErrorExecuted || result; }) @@ -64,8 +58,8 @@ public async Task The_import_should_deduplicate_on_TimeOfFailure() var attempts = failure.ProcessingAttempts; using (Assert.EnterMultipleScope()) { - Assert.That(attempts, Has.Count.EqualTo(2)); - Assert.That(attempts.Select(a => a.AttemptedAt), Is.EquivalentTo(context.FailureTimes)); + Assert.That(attempts, Has.Count.EqualTo(1)); + Assert.That(attempts[^1].AttemptedAt, Is.EqualTo(context.FailureTime)); } } @@ -94,41 +88,34 @@ protected override TransportOperations CreateMessage(MyContext context) { var messageId = Guid.NewGuid().ToString(); context.UniqueId = DeterministicGuid.MakeId(messageId, "Error.SourceEndpoint").ToString(); - context.FailureTimes = new[] - { - new DateTime(2020, 09, 05, 13, 20, 00, 0, DateTimeKind.Utc), - new DateTime(2020, 09, 05, 12, 20, 00, 0, DateTimeKind.Utc), - }; + context.FailureTime = new DateTime(2020, 09, 05, 13, 20, 00, 0, DateTimeKind.Utc); - return new TransportOperations(GetMessages(context.UniqueId, context.FailureTimes).ToArray()); + return new TransportOperations([.. GetMessages(context.UniqueId, context.FailureTime)]); } - IEnumerable GetMessages(string uniqueId, DateTime[] failureTimes) + IEnumerable GetMessages(string uniqueId, DateTime failureTime) { - for (var failureNo = 0; failureNo < failureTimes.Length; failureNo++) + for (var i = 0; i < NumberOfDuplicates; i++) { - for (var i = 0; i < 5; i++) + var messageId = Guid.NewGuid().ToString(); + var headers = new Dictionary { - var messageId = Guid.NewGuid().ToString(); - var headers = new Dictionary - { - [Headers.MessageId] = messageId, - ["ServiceControl.Retry.UniqueMessageId"] = uniqueId, - [Headers.ProcessingEndpoint] = "Error.SourceEndpoint", - ["NServiceBus.ExceptionInfo.ExceptionType"] = typeof(Exception).FullName, - ["NServiceBus.ExceptionInfo.Message"] = "Bad thing happened", - ["NServiceBus.ExceptionInfo.InnerExceptionType"] = "System.Exception", - ["NServiceBus.ExceptionInfo.Source"] = "NServiceBus.Core", - ["NServiceBus.ExceptionInfo.StackTrace"] = string.Empty, - ["NServiceBus.FailedQ"] = "Error.SourceEndpoint", - ["NServiceBus.TimeOfFailure"] = DateTimeOffsetHelper.ToWireFormattedString(failureTimes[failureNo]), - ["Counter"] = i.ToString() - }; - - var outgoingMessage = new OutgoingMessage(messageId, headers, new byte[0]); - - yield return new TransportOperation(outgoingMessage, new UnicastAddressTag("error")); - } + [Headers.MessageId] = messageId, + ["ServiceControl.Retry.UniqueMessageId"] = uniqueId, + [Headers.ProcessingEndpoint] = "Error.SourceEndpoint", + ["NServiceBus.ExceptionInfo.ExceptionType"] = typeof(Exception).FullName, + ["NServiceBus.ExceptionInfo.Message"] = "Bad thing happened", + ["NServiceBus.ExceptionInfo.InnerExceptionType"] = "System.Exception", + ["NServiceBus.ExceptionInfo.Source"] = "NServiceBus.Core", + ["NServiceBus.ExceptionInfo.StackTrace"] = string.Empty, + ["NServiceBus.FailedQ"] = "Error.SourceEndpoint", + ["NServiceBus.TimeOfFailure"] = DateTimeOffsetHelper.ToWireFormattedString(failureTime), + ["Counter"] = i.ToString() + }; + + var outgoingMessage = new OutgoingMessage(messageId, headers, Array.Empty()); + + yield return new TransportOperation(outgoingMessage, new UnicastAddressTag("error")); } } } @@ -138,11 +125,13 @@ class MyContext : ScenarioContext { public string UniqueId { get; set; } - public DateTime[] FailureTimes { get; set; } + public DateTime FailureTime { get; set; } + + public int IngestedCount => receivedMessages.Count; public void OnMessage(string counter) => receivedMessages.AddOrUpdate(counter, true, (id, old) => true); readonly ConcurrentDictionary receivedMessages = new(); } } -} \ No newline at end of file +} diff --git a/src/ServiceControl.Persistence.Tests.RavenDB/Recoverability/ProcessingAttemptTrimmingTests.cs b/src/ServiceControl.Persistence.Tests.RavenDB/Recoverability/ProcessingAttemptTrimmingTests.cs new file mode 100644 index 0000000000..2ba7b07503 --- /dev/null +++ b/src/ServiceControl.Persistence.Tests.RavenDB/Recoverability/ProcessingAttemptTrimmingTests.cs @@ -0,0 +1,46 @@ +namespace ServiceControl.Persistence.Tests.RavenDB.Recoverability; + +using System.Linq; +using System.Threading.Tasks; +using NUnit.Framework; + +[TestFixture] +class ProcessingAttemptTrimmingTests : RavenPersistenceTestBase +{ + const int MaxStoredAttempts = 10; + const int IngestedAttempts = 15; + + [Test] + public async Task Only_the_latest_attempts_are_kept() + { + var failure = new IngestedFailure(); + + for (var i = 0; i < IngestedAttempts; i++) + { + await Ingest(failure.NextAttempt(failure.AttemptedAt.AddMinutes(i))); + } + + var message = await FailedMessageQueryStore.GetFailedMessage(failure.UniqueMessageIdString); + + using (Assert.EnterMultipleScope()) + { + Assert.That(message.ProcessingAttempts, Has.Count.EqualTo(MaxStoredAttempts)); + Assert.That( + message.ProcessingAttempts.Select(attempt => attempt.AttemptedAt), + Is.EqualTo(Enumerable + .Range(IngestedAttempts - MaxStoredAttempts, MaxStoredAttempts) + .Select(i => failure.AttemptedAt.AddMinutes(i)))); + } + } + + async Task Ingest(IngestedFailure failure) + { + await using (var unitOfWork = await UnitOfWorkFactory.StartNew()) + { + await unitOfWork.Recoverability.RecordFailedProcessingAttempt(failure.Context, failure.ProcessingAttempt, failure.Groups); + await unitOfWork.Complete(TestContext.CurrentContext.CancellationToken); + } + + await CompleteDatabaseOperation(); + } +}