diff --git a/src/ServiceControl.Audit.Persistence.InMemory/InMemoryFailedAuditStorage.cs b/src/ServiceControl.Audit.Persistence.InMemory/InMemoryFailedAuditStorage.cs index 652b79983e..e6dbf35e02 100644 --- a/src/ServiceControl.Audit.Persistence.InMemory/InMemoryFailedAuditStorage.cs +++ b/src/ServiceControl.Audit.Persistence.InMemory/InMemoryFailedAuditStorage.cs @@ -25,7 +25,15 @@ public async Task ProcessFailedMessages(Func x.Id == message.Id); + if (existing >= 0) + { + dataStore.failedAuditImports[existing] = message; + } + else + { + dataStore.failedAuditImports.Add(message); + } return Task.CompletedTask; } diff --git a/src/ServiceControl.Audit.Persistence.RavenDB/RavenFailedAuditStorage.cs b/src/ServiceControl.Audit.Persistence.RavenDB/RavenFailedAuditStorage.cs index 183907a9ed..56573c3bb5 100644 --- a/src/ServiceControl.Audit.Persistence.RavenDB/RavenFailedAuditStorage.cs +++ b/src/ServiceControl.Audit.Persistence.RavenDB/RavenFailedAuditStorage.cs @@ -15,6 +15,10 @@ class RavenFailedAuditStorage(IRavenSessionProvider sessionProvider) : IFailedAu public async Task SaveFailedAuditImport(FailedAuditImport message, CancellationToken cancellationToken = default) { using var session = await sessionProvider.OpenSession(cancellationToken: cancellationToken); + if (message.Id != null && !message.Id.StartsWith("FailedAuditImports/", StringComparison.OrdinalIgnoreCase)) + { + message.Id = $"FailedAuditImports/{message.Id}"; + } await session.StoreAsync(message, token: cancellationToken); await session.SaveChangesAsync(cancellationToken); } diff --git a/src/ServiceControl.Audit.Persistence.Tests/FailedAuditStorageTests.cs b/src/ServiceControl.Audit.Persistence.Tests/FailedAuditStorageTests.cs index c450ff6e25..473fcf6e59 100644 --- a/src/ServiceControl.Audit.Persistence.Tests/FailedAuditStorageTests.cs +++ b/src/ServiceControl.Audit.Persistence.Tests/FailedAuditStorageTests.cs @@ -1,8 +1,10 @@ namespace ServiceControl.Audit.Persistence.Tests { using System; + using System.Collections.Generic; using System.Threading; using System.Threading.Tasks; + using NServiceBus; using NUnit.Framework; using ServiceControl.Audit.Auditing; @@ -21,6 +23,35 @@ public async Task Should_store_failures() Assert.That(numFailures, Is.EqualTo(1)); } + [Test] + public async Task Repeated_failure_of_the_same_message_stores_one_document() + { + var headers = new Dictionary + { + { Headers.MessageId, "message-1" }, + { Headers.ProcessingEndpoint, "Sales" } + }; + await FailedAuditStorage.SaveFailedAuditImport(new FailedAuditImport { Id = FailedAuditImport.DeriveKey(headers, "native-1").ToString() }); + await FailedAuditStorage.SaveFailedAuditImport(new FailedAuditImport { Id = FailedAuditImport.DeriveKey(headers, "native-1").ToString() }); + await configuration.CompleteDBOperation(); + + Assert.That(await FailedAuditStorage.GetFailedAuditsCount(), Is.EqualTo(1)); + } + + [Test] + public void DeriveKey_uses_retry_id_or_falls_back_for_malformed_headers() + { + var retryId = Guid.NewGuid(); + var headers = new Dictionary { ["ServiceControl.Retry.UniqueMessageId"] = retryId.ToString() }; + + var fallback = FailedAuditImport.DeriveKey(new Dictionary(), "native-1"); + Assert.That(FailedAuditImport.DeriveKey(headers, "native-1"), Is.EqualTo(retryId)); + Assert.That(FailedAuditImport.DeriveKey(new Dictionary(), "native-1"), Is.EqualTo(fallback)); + Assert.That(FailedAuditImport.DeriveKey(new Dictionary(), "native-2"), Is.Not.EqualTo(fallback)); + headers["ServiceControl.Retry.UniqueMessageId"] = "not-a-guid"; + Assert.That(FailedAuditImport.DeriveKey(headers, "native-1"), Is.EqualTo(fallback)); + } + [Test] public async Task Should_be_able_to_process_failures() { diff --git a/src/ServiceControl.Audit.Persistence/FailedAuditImport.cs b/src/ServiceControl.Audit.Persistence/FailedAuditImport.cs index da916151e3..bcf824281d 100644 --- a/src/ServiceControl.Audit.Persistence/FailedAuditImport.cs +++ b/src/ServiceControl.Audit.Persistence/FailedAuditImport.cs @@ -1,9 +1,30 @@ namespace ServiceControl.Audit.Auditing { + using System; + using System.Collections.Generic; + using ServiceControl.Audit.Persistence.Infrastructure; + public class FailedAuditImport { public string Id { get; set; } public FailedTransportMessage Message { get; set; } public string ExceptionInfo { get; set; } + + public static Guid DeriveKey(IReadOnlyDictionary headers, string nativeMessageId) + { + try + { + if (Guid.TryParse(headers.UniqueId(), out var uniqueMessageId)) + { + return uniqueMessageId; + } + } + catch (Exception) + { + // Malformed messages may have no processing endpoint; the transport ID is always available. + } + + return DeterministicGuid.MakeId(nativeMessageId); + } } } \ No newline at end of file diff --git a/src/ServiceControl.Audit/Auditing/AuditIngestionFaultPolicy.cs b/src/ServiceControl.Audit/Auditing/AuditIngestionFaultPolicy.cs index b9d323747f..54c1d6221d 100644 --- a/src/ServiceControl.Audit/Auditing/AuditIngestionFaultPolicy.cs +++ b/src/ServiceControl.Audit/Auditing/AuditIngestionFaultPolicy.cs @@ -56,7 +56,7 @@ async Task StoreFailedMessageDocument(ErrorContext errorContext, CancellationTok { var failure = new FailedAuditImport { - Id = Guid.NewGuid().ToString(), + Id = FailedAuditImport.DeriveKey(errorContext.Headers, errorContext.MessageId).ToString(), Message = new FailedTransportMessage { Id = errorContext.MessageId,