Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,15 @@ public async Task ProcessFailedMessages(Func<FailedTransportMessage, Func<Cancel
public Task SaveFailedAuditImport(FailedAuditImport message, CancellationToken cancellationToken = default)
{
cancellationToken.ThrowIfCancellationRequested();
dataStore.failedAuditImports.Add(message);
var existing = message.Id == null ? -1 : dataStore.failedAuditImports.FindIndex(x => x.Id == message.Id);
if (existing >= 0)
{
dataStore.failedAuditImports[existing] = message;
}
else
{
dataStore.failedAuditImports.Add(message);
}
return Task.CompletedTask;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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))

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is strange. Why would an ID ever include the "FailedAuditImports/" prefix at all when it reaches this method?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The code is replicating what we were doing in error, where there are places that records get roundtripped.
In this case I'm pretty sure it's always unset on the way in.

I did wonder whether it would be better to either:

  • Verify that and enforce it being unset on save; or
  • Check if the "id" is even used outside of raven and remove it from the contract (I think we did this in at least one place)

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

From our discussion, this is how the error instance does it and it can stay.
It has to be returned for now so deletions work

{
message.Id = $"FailedAuditImports/{message.Id}";
}
await session.StoreAsync(message, token: cancellationToken);
await session.SaveChangesAsync(cancellationToken);
}
Expand Down
Original file line number Diff line number Diff line change
@@ -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;

Expand All @@ -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<string, string>
{
{ 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<string, string> { ["ServiceControl.Retry.UniqueMessageId"] = retryId.ToString() };

var fallback = FailedAuditImport.DeriveKey(new Dictionary<string, string>(), "native-1");
Assert.That(FailedAuditImport.DeriveKey(headers, "native-1"), Is.EqualTo(retryId));
Assert.That(FailedAuditImport.DeriveKey(new Dictionary<string, string>(), "native-1"), Is.EqualTo(fallback));
Assert.That(FailedAuditImport.DeriveKey(new Dictionary<string, string>(), "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()
{
Expand Down
21 changes: 21 additions & 0 deletions src/ServiceControl.Audit.Persistence/FailedAuditImport.cs
Original file line number Diff line number Diff line change
@@ -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<string, string> 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);
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Loading