diff --git a/src/ServiceControl.AcceptanceTests.PostgreSql/ServiceControl.AcceptanceTests.PostgreSql.csproj b/src/ServiceControl.AcceptanceTests.PostgreSql/ServiceControl.AcceptanceTests.PostgreSql.csproj
index d5ac47492d..74ed777f0b 100644
--- a/src/ServiceControl.AcceptanceTests.PostgreSql/ServiceControl.AcceptanceTests.PostgreSql.csproj
+++ b/src/ServiceControl.AcceptanceTests.PostgreSql/ServiceControl.AcceptanceTests.PostgreSql.csproj
@@ -52,8 +52,7 @@
-
-
+
diff --git a/src/ServiceControl.AcceptanceTests.SqlServer/ServiceControl.AcceptanceTests.SqlServer.csproj b/src/ServiceControl.AcceptanceTests.SqlServer/ServiceControl.AcceptanceTests.SqlServer.csproj
index cecd9a36e8..5dda932f0e 100644
--- a/src/ServiceControl.AcceptanceTests.SqlServer/ServiceControl.AcceptanceTests.SqlServer.csproj
+++ b/src/ServiceControl.AcceptanceTests.SqlServer/ServiceControl.AcceptanceTests.SqlServer.csproj
@@ -52,8 +52,7 @@
-
-
+
diff --git a/src/ServiceControl.AcceptanceTests/Recoverability/MessageFailures/When_a_retry_fails_to_be_sent.cs b/src/ServiceControl.AcceptanceTests/Recoverability/MessageFailures/When_a_retry_fails_to_be_sent.cs
index e367de15cc..bbeea18019 100644
--- a/src/ServiceControl.AcceptanceTests/Recoverability/MessageFailures/When_a_retry_fails_to_be_sent.cs
+++ b/src/ServiceControl.AcceptanceTests/Recoverability/MessageFailures/When_a_retry_fails_to_be_sent.cs
@@ -17,7 +17,7 @@
using NUnit.Framework;
using ServiceControl.MessageFailures;
using ServiceControl.MessageFailures.Api;
- using ServiceControl.Persistence;
+ using ServiceControl.Operations.BodyStorage;
using ServiceControl.Recoverability;
using TestSupport;
@@ -31,7 +31,7 @@ public async Task SubsequentBatchesShouldBeProcessed(CancellationToken cancellat
CustomizeHostBuilder = hostBuilder =>
{
- hostBuilder.Services.AddSingleton(provider => new FakeReturnToSender(provider.GetRequiredService(), provider.GetRequiredService()));
+ hostBuilder.Services.AddSingleton(provider => new FakeReturnToSender(provider.GetRequiredService(), provider.GetRequiredService()));
};
await Define()
@@ -148,8 +148,8 @@ public class MyContext : ScenarioContext
public class MessageThatWillFail : ICommand;
- public class FakeReturnToSender(IFailedMessageRetryDataStore errorMessageStore, MyContext myContext)
- : ReturnToSender(errorMessageStore, NullLogger.Instance)
+ public class FakeReturnToSender(IBodyStorage bodyStorage, MyContext myContext)
+ : ReturnToSender(bodyStorage, NullLogger.Instance)
{
public override Task HandleMessage(MessageContext message, IMessageDispatcher sender, string errorQueueTransportAddress, CancellationToken cancellationToken = default)
{
diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/BodyStorage/BodyStorage.cs b/src/ServiceControl.Persistence.EFCore/Implementation/BodyStorage/BodyStorage.cs
index 233771e24e..cdfdda6b4e 100644
--- a/src/ServiceControl.Persistence.EFCore/Implementation/BodyStorage/BodyStorage.cs
+++ b/src/ServiceControl.Persistence.EFCore/Implementation/BodyStorage/BodyStorage.cs
@@ -19,13 +19,13 @@ namespace ServiceControl.Persistence.EFCore.Implementation.BodyStorage;
///
public class BodyStorage(IServiceScopeFactory scopeFactory, IBodyStoragePersistence storagePersistence) : DataStoreBase(scopeFactory), IBodyStorage
{
- public async Task TryFetch(string bodyId, CancellationToken cancellationToken = default)
+ public async Task TryFetch(string bodyId, CancellationToken cancellationToken = default)
{
var row = await ExecuteWithDbContext((dbContext, token) => ResolveBody(dbContext, bodyId, token), cancellationToken);
if (row == null)
{
- return null; // No such message: the API turns this into a 404.
+ return MessageBodyResult.NotFound();
}
// Bodies are immutable per message, so the id is a stable ETag.
@@ -35,33 +35,42 @@ public class BodyStorage(IServiceScopeFactory scopeFactory, IBodyStoragePersiste
{
var external = await storagePersistence.ReadBody(uniqueMessageId, cancellationToken);
- return external == null
- ? new MessageBodyStreamResult { HasResult = false }
- : new MessageBodyStreamResult
- {
- HasResult = true,
- Stream = external.Stream,
- ContentType = external.ContentType,
- BodySize = external.BodySize,
- Etag = uniqueMessageId
- };
+ if (external == null)
+ {
+ return MessageBodyResult.Unavailable();
+ }
+
+ if (external.BodySize == 0)
+ {
+ await external.Stream.DisposeAsync();
+ return MessageBodyResult.Empty();
+ }
+
+ return MessageBodyResult.Available(new MessageBodyStreamContent(external.Stream, external.ContentType, external.BodySize, uniqueMessageId));
}
if (row.BodyText != null)
{
var bytes = Encoding.UTF8.GetBytes(row.BodyText);
- return new MessageBodyStreamResult
+ if (bytes.Length == 0)
{
- HasResult = true,
- Stream = new MemoryStream(bytes, writable: false),
- ContentType = row.BodyContentType,
- BodySize = bytes.Length,
- Etag = uniqueMessageId
- };
+ return MessageBodyResult.Empty();
+ }
+
+ return MessageBodyResult.Available(new MessageBodyStreamContent(
+ new MemoryStream(bytes, writable: false),
+ row.BodyContentType,
+ bytes.Length,
+ uniqueMessageId));
+ }
+
+ if (row.BodySize == 0)
+ {
+ return MessageBodyResult.Empty();
}
- return new MessageBodyStreamResult { HasResult = false }; // Message exists but carries no body.
+ return MessageBodyResult.Unavailable();
}
static async Task ResolveBody(ServiceControlDbContext dbContext, string bodyId, CancellationToken cancellationToken)
@@ -88,6 +97,7 @@ public class BodyStorage(IServiceScopeFactory scopeFactory, IBodyStoragePersiste
UniqueMessageId = message.UniqueMessageId,
BodyText = message.BodyText,
BodyStoredExternally = message.BodyStoredExternally,
+ BodySize = message.BodySize,
BodyContentType = message.BodyContentType
})
.FirstOrDefaultAsync(cancellationToken);
@@ -97,6 +107,7 @@ sealed class BodyRow
public Guid UniqueMessageId { get; init; }
public string? BodyText { get; init; }
public bool BodyStoredExternally { get; init; }
+ public int BodySize { get; init; }
public string? BodyContentType { get; init; }
}
}
diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/FailedMessageRetryDataStore.cs b/src/ServiceControl.Persistence.EFCore/Implementation/FailedMessageRetryDataStore.cs
index a2850f9ddd..4bef78cb47 100644
--- a/src/ServiceControl.Persistence.EFCore/Implementation/FailedMessageRetryDataStore.cs
+++ b/src/ServiceControl.Persistence.EFCore/Implementation/FailedMessageRetryDataStore.cs
@@ -55,18 +55,27 @@ public Task ProcessPendingRetries(DateTime periodFrom, DateTime periodTo, string
public async Task GetFailedMessageBody(string uniqueMessageId, CancellationToken cancellationToken = default)
{
- var result = await bodyStorage.TryFetch(uniqueMessageId, cancellationToken)
- ?? throw new InvalidOperationException("IBodyStorage.TryFetch result cannot be null");
+ var result = await bodyStorage.TryFetch(uniqueMessageId, cancellationToken);
- if (!result.HasResult)
+ if (result.State == MessageBodyState.NotFound)
+ {
+ throw new InvalidOperationException("IBodyStorage.TryFetch result cannot be null");
+ }
+
+ if (result.State == MessageBodyState.Unavailable)
{
throw new InvalidOperationException("IBodyStorage.TryFetch did not return a body");
}
- await using (result.Stream)
+ if (result.State == MessageBodyState.Empty)
+ {
+ return [];
+ }
+
+ await using (result.Content.Stream)
{
using var memoryStream = new MemoryStream();
- await result.Stream.CopyToAsync(memoryStream, cancellationToken);
+ await result.Content.Stream.CopyToAsync(memoryStream, cancellationToken);
return memoryStream.ToArray();
}
}
diff --git a/src/ServiceControl.Persistence.RavenDB/ErrorMessagesDataStore.cs b/src/ServiceControl.Persistence.RavenDB/ErrorMessagesDataStore.cs
index 2e5bbce988..fec09b8f8f 100644
--- a/src/ServiceControl.Persistence.RavenDB/ErrorMessagesDataStore.cs
+++ b/src/ServiceControl.Persistence.RavenDB/ErrorMessagesDataStore.cs
@@ -510,26 +510,34 @@ record struct FailedMessageProjection(string UniqueMessageId);
public async Task GetFailedMessageBody(string uniqueMessageId, CancellationToken cancellationToken = default)
{
- byte[] body = null;
- var result = await bodyStorage.TryFetch(uniqueMessageId, cancellationToken)
- ?? throw new InvalidOperationException("IBodyStorage.TryFetch result cannot be null");
+ var result = await bodyStorage.TryFetch(uniqueMessageId, cancellationToken);
- if (result.HasResult)
+ if (result.State == MessageBodyState.NotFound)
{
- await using (result.Stream) // Not strictly required for MemoryStream but might be different behavior in future .NET versions
- {
- // Unfortunately we can't use the buffer manager here yet because core doesn't allow to set the length property so usage of GetBuffer is not possible
- // furthermore call ToArray would neglect many of the benefits of the recyclable stream
- // RavenDB always returns a memory stream in ver. 3.5 so there is no need to pretend we need to do buffered reads since the memory is anyway fully allocated already
- // this assumption might change when we stop supporting RavenDB 3.5 but right now this is the most memory efficient way to do things
- // https://github.com/microsoft/Microsoft.IO.RecyclableMemoryStream#getbuffer-and-toarray
- using var memoryStream = new MemoryStream();
- await result.Stream.CopyToAsync(memoryStream, cancellationToken);
-
- body = memoryStream.ToArray();
- }
+ throw new InvalidOperationException("IBodyStorage.TryFetch result cannot be null");
+ }
+
+ if (result.State == MessageBodyState.Unavailable)
+ {
+ throw new InvalidOperationException("IBodyStorage.TryFetch result cannot be null");
+ }
+
+ if (result.State == MessageBodyState.Empty)
+ {
+ return [];
+ }
+
+ await using (result.Content.Stream) // Not strictly required for MemoryStream but might be different behavior in future .NET versions
+ {
+ // Unfortunately we can't use the buffer manager here yet because core doesn't allow to set the length property so usage of GetBuffer is not possible
+ // furthermore call ToArray would neglect many of the benefits of the recyclable stream
+ // RavenDB always returns a memory stream in ver. 3.5 so there is no need to pretend we need to do buffered reads since the memory is anyway fully allocated already
+ // this assumption might change when we stop supporting RavenDB 3.5 but right now this is the most memory efficient way to do things
+ // https://github.com/microsoft/Microsoft.IO.RecyclableMemoryStream#getbuffer-and-toarray
+ using var memoryStream = new MemoryStream();
+ await result.Content.Stream.CopyToAsync(memoryStream, cancellationToken);
+ return memoryStream.ToArray();
}
- return body;
}
}
}
diff --git a/src/ServiceControl.Persistence.RavenDB/RavenAttachmentsBodyStorage.cs b/src/ServiceControl.Persistence.RavenDB/RavenAttachmentsBodyStorage.cs
index 64cc215ebc..69a498fac0 100644
--- a/src/ServiceControl.Persistence.RavenDB/RavenAttachmentsBodyStorage.cs
+++ b/src/ServiceControl.Persistence.RavenDB/RavenAttachmentsBodyStorage.cs
@@ -14,7 +14,7 @@ class RavenAttachmentsBodyStorage(IRavenSessionProvider sessionProvider) : IBody
{
public const string AttachmentName = "body";
- public async Task TryFetch(string bodyId, CancellationToken cancellationToken = default)
+ public async Task TryFetch(string bodyId, CancellationToken cancellationToken = default)
{
using var session = await sessionProvider.OpenSession(cancellationToken: cancellationToken);
@@ -23,7 +23,7 @@ public async Task TryFetch(string bodyId, CancellationT
if (Guid.TryParse(bodyId, out _))
{
var result = await ResultForUniqueId(session, bodyId, cancellationToken);
- if (result != null)
+ if (result.State != MessageBodyState.NotFound)
{
return result;
}
@@ -42,28 +42,37 @@ public async Task TryFetch(string bodyId, CancellationT
return await ResultForUniqueId(session, uniqueId, cancellationToken);
}
- return null;
+ return MessageBodyResult.NotFound();
}
- async Task ResultForUniqueId(IAsyncDocumentSession session, string uniqueId, CancellationToken cancellationToken)
+ async Task ResultForUniqueId(IAsyncDocumentSession session, string uniqueId, CancellationToken cancellationToken)
{
var documentId = FailedMessageIdGenerator.MakeDocumentId(uniqueId);
+ var failedMessage = await session.LoadAsync(documentId, cancellationToken);
+
+ if (failedMessage == null)
+ {
+ return MessageBodyResult.NotFound();
+ }
var result = await session.Advanced.Attachments.GetAsync(documentId, AttachmentName, cancellationToken);
if (result == null)
{
- return null;
+ return MessageBodyResult.Unavailable();
}
- return new MessageBodyStreamResult
+ if (result.Details.Size == 0)
{
- HasResult = true,
- Stream = result.Stream,
- ContentType = result.Details.ContentType,
- BodySize = (int)result.Details.Size,
- Etag = result.Details.ChangeVector
- };
+ await result.Stream.DisposeAsync();
+ return MessageBodyResult.Empty();
+ }
+
+ return MessageBodyResult.Available(new MessageBodyStreamContent(
+ result.Stream,
+ result.Details.ContentType,
+ (int)result.Details.Size,
+ result.Details.ChangeVector));
}
}
}
\ No newline at end of file
diff --git a/src/ServiceControl.Persistence.Tests/BodyStorage/AttachmentsBodyStorageTests.cs b/src/ServiceControl.Persistence.Tests/BodyStorage/AttachmentsBodyStorageTests.cs
index 8fd23948d9..d4053ac891 100644
--- a/src/ServiceControl.Persistence.Tests/BodyStorage/AttachmentsBodyStorageTests.cs
+++ b/src/ServiceControl.Persistence.Tests/BodyStorage/AttachmentsBodyStorageTests.cs
@@ -10,6 +10,7 @@
using NUnit.Framework;
using ServiceControl.MessageFailures;
using ServiceControl.Operations;
+ using ServiceControl.Operations.BodyStorage;
using ServiceControl.Persistence.UnitOfWork;
[TestFixture]
@@ -81,14 +82,14 @@ async Task RunTest(Func, string> getIdToQuery)
Assert.That(retrieved, Is.Not.Null);
using (Assert.EnterMultipleScope())
{
- Assert.That(retrieved.HasResult, Is.True);
- Assert.That(retrieved.ContentType, Is.EqualTo(contentType));
+ Assert.That(retrieved.State, Is.EqualTo(MessageBodyState.Available));
+ Assert.That(retrieved.Content.ContentType, Is.EqualTo(contentType));
}
- var buffer = new byte[retrieved.BodySize];
- await using (retrieved.Stream)
+ var buffer = new byte[retrieved.Content.BodySize];
+ await using (retrieved.Content.Stream)
{
- retrieved.Stream.ReadExactly(buffer);
+ retrieved.Content.Stream.ReadExactly(buffer);
}
Assert.That(buffer, Is.EqualTo(body));
diff --git a/src/ServiceControl.Persistence.Tests/EFCore/BodyReadTests.cs b/src/ServiceControl.Persistence.Tests/EFCore/BodyReadTests.cs
index 3223dc8323..7a2996c262 100644
--- a/src/ServiceControl.Persistence.Tests/EFCore/BodyReadTests.cs
+++ b/src/ServiceControl.Persistence.Tests/EFCore/BodyReadTests.cs
@@ -26,9 +26,9 @@ public async Task Fetches_an_inline_text_body()
Assert.That(result, Is.Not.Null);
using (Assert.EnterMultipleScope())
{
- Assert.That(result.HasResult, Is.True);
- Assert.That(result.ContentType, Is.EqualTo("text/xml"));
- Assert.That(Encoding.UTF8.GetString(ReadAll(result.Stream)), Is.EqualTo("1"));
+ Assert.That(result.State, Is.EqualTo(MessageBodyState.Available));
+ Assert.That(result.Content.ContentType, Is.EqualTo("text/xml"));
+ Assert.That(Encoding.UTF8.GetString(ReadAll(result.Content.Stream)), Is.EqualTo("1"));
}
}
@@ -44,9 +44,9 @@ public async Task Fetches_an_external_binary_body()
Assert.That(result, Is.Not.Null);
using (Assert.EnterMultipleScope())
{
- Assert.That(result.HasResult, Is.True);
- Assert.That(result.ContentType, Is.EqualTo("application/octet-stream"));
- Assert.That(ReadAll(result.Stream), Is.EqualTo(body));
+ Assert.That(result.State, Is.EqualTo(MessageBodyState.Available));
+ Assert.That(result.Content.ContentType, Is.EqualTo("application/octet-stream"));
+ Assert.That(ReadAll(result.Content.Stream), Is.EqualTo(body));
}
}
@@ -62,8 +62,8 @@ public async Task Fetches_the_whole_body_for_large_text_not_the_inline_prefix()
Assert.That(result, Is.Not.Null);
using (Assert.EnterMultipleScope())
{
- Assert.That(result.HasResult, Is.True);
- Assert.That(ReadAll(result.Stream), Is.EqualTo(body), "external storage is authoritative, not the inline search prefix");
+ Assert.That(result.State, Is.EqualTo(MessageBodyState.Available));
+ Assert.That(ReadAll(result.Content.Stream), Is.EqualTo(body), "external storage is authoritative, not the inline search prefix");
}
}
@@ -76,11 +76,11 @@ public async Task Fetches_by_message_id()
var result = await Fetch(failure.MessageId);
Assert.That(result, Is.Not.Null);
- Assert.That(result.HasResult, Is.True);
+ Assert.That(result.State, Is.EqualTo(MessageBodyState.Available));
}
[Test]
- public async Task Reports_no_body_for_an_empty_body()
+ public async Task Reports_an_empty_body()
{
var failure = new IngestedFailure { Body = [] };
await Ingest(failure);
@@ -88,7 +88,7 @@ public async Task Reports_no_body_for_an_empty_body()
var result = await Fetch(failure.UniqueMessageIdString);
Assert.That(result, Is.Not.Null);
- Assert.That(result.HasResult, Is.False);
+ Assert.That(result.State, Is.EqualTo(MessageBodyState.Empty));
}
[Test]
@@ -96,10 +96,10 @@ public async Task Returns_null_for_an_unknown_message()
{
var result = await Fetch(Guid.NewGuid().ToString());
- Assert.That(result, Is.Null);
+ Assert.That(result.State, Is.EqualTo(MessageBodyState.NotFound));
}
- async Task Fetch(string bodyId)
+ async Task Fetch(string bodyId)
{
using var scope = ServiceProvider.CreateScope();
var bodyStorage = scope.ServiceProvider.GetRequiredService();
diff --git a/src/ServiceControl.Persistence.Tests/EFCore/FailedMessageRetryBodyDataStoreTests.cs b/src/ServiceControl.Persistence.Tests/EFCore/FailedMessageRetryBodyDataStoreTests.cs
index 0a46b887cc..d964caedec 100644
--- a/src/ServiceControl.Persistence.Tests/EFCore/FailedMessageRetryBodyDataStoreTests.cs
+++ b/src/ServiceControl.Persistence.Tests/EFCore/FailedMessageRetryBodyDataStoreTests.cs
@@ -24,7 +24,7 @@ public async Task GetFailedMessageBody_returns_inline_BodyText_bytes()
[Test]
public async Task GetFailedMessageBody_throws_when_the_body_is_unavailable()
{
- var id = await SeedFailedMessage();
+ var id = await SeedFailedMessage(bodyStoredExternally: true);
Assert.ThrowsAsync(() =>
FailedMessageRetryStore.GetFailedMessageBody(id.ToString()));
@@ -33,8 +33,8 @@ public async Task GetFailedMessageBody_throws_when_the_body_is_unavailable()
[Test]
public async Task GetFailedMessageBody_returns_external_storage_body_when_BodyStoredExternally()
{
- var id = await SeedFailedMessage(bodyStoredExternally: true);
var expected = Encoding.UTF8.GetBytes("external body payload");
+ var id = await SeedFailedMessage(bodyStoredExternally: true);
await RecordedBodies.WriteBody(id.ToString(), expected, "text/plain");
diff --git a/src/ServiceControl.Persistence.Tests/Recoverability/ReturnToSenderDequeuerTests.cs b/src/ServiceControl.Persistence.Tests/Recoverability/ReturnToSenderDequeuerTests.cs
index 7adc218149..bc366dfcb4 100644
--- a/src/ServiceControl.Persistence.Tests/Recoverability/ReturnToSenderDequeuerTests.cs
+++ b/src/ServiceControl.Persistence.Tests/Recoverability/ReturnToSenderDequeuerTests.cs
@@ -2,6 +2,7 @@
{
using System;
using System.Collections.Generic;
+ using System.IO;
using System.Linq;
using System.Text;
using System.Threading;
@@ -18,6 +19,7 @@
using Persistence.Infrastructure;
using ServiceControl.CompositeViews.Messages;
using ServiceControl.Operations;
+ using ServiceControl.Operations.BodyStorage;
using ServiceControl.Recoverability;
[TestFixture]
@@ -67,11 +69,47 @@ public async Task It_fetches_the_body_from_storage_if_provided()
};
var message = CreateMessage(Guid.NewGuid().ToString(), headers);
- await new ReturnToSender(new FakeErrorMessageDataStore(), NullLogger.Instance).HandleMessage(message, sender, "error");
+ await new ReturnToSender(new FakeBodyStorage(), NullLogger.Instance).HandleMessage(message, sender, "error");
Assert.That(Encoding.UTF8.GetString(sender.Message.Body.ToArray()), Is.EqualTo("MessageBodyId"));
}
+ [Test]
+ public async Task It_sends_an_empty_body_when_storage_reports_empty()
+ {
+ var sender = new FakeSender();
+ var headers = new Dictionary
+ {
+ ["ServiceControl.TargetEndpointAddress"] = "TargetEndpoint",
+ ["ServiceControl.Retry.Attempt.MessageId"] = "MessageBodyId",
+ ["ServiceControl.Retry.UniqueMessageId"] = "MessageBodyId"
+ };
+ var message = CreateMessage(Guid.NewGuid().ToString(), headers);
+
+ await new ReturnToSender(new FakeBodyStorage(MessageBodyState.Empty), NullLogger.Instance).HandleMessage(message, sender, "error");
+
+ Assert.That(sender.Message.Body.ToArray(), Is.Empty);
+ }
+
+ [TestCase(MessageBodyState.NotFound)]
+ [TestCase(MessageBodyState.Unavailable)]
+ public void It_does_not_send_when_the_body_cannot_be_retrieved(MessageBodyState state)
+ {
+ var sender = new FakeSender();
+ var headers = new Dictionary
+ {
+ ["ServiceControl.TargetEndpointAddress"] = "TargetEndpoint",
+ ["ServiceControl.Retry.Attempt.MessageId"] = "MessageBodyId",
+ ["ServiceControl.Retry.UniqueMessageId"] = "MessageBodyId"
+ };
+ var message = CreateMessage(Guid.NewGuid().ToString(), headers);
+
+ Assert.ThrowsAsync(() =>
+ new ReturnToSender(new FakeBodyStorage(state), NullLogger.Instance).HandleMessage(message, sender, "error"));
+
+ Assert.That(sender.Message, Is.Null);
+ }
+
[Test]
public async Task It_uses_retry_to_if_provided()
{
@@ -167,15 +205,19 @@ public Task Dispatch(TransportOperations outgoingMessages, TransportTransaction
}
}
- class FakeErrorMessageDataStore : IFailedMessageRetryDataStore
+ class FakeBodyStorage(MessageBodyState state = MessageBodyState.Available) : IBodyStorage
{
- public Task ProcessPendingRetries(DateTime periodFrom, DateTime periodTo, string queueAddress, Func processCallback, CancellationToken cancellationToken = default) => throw new NotImplementedException();
-
- public Task GetRetryPendingMessages(DateTime from, DateTime to, string queueAddress, CancellationToken cancellationToken = default) => throw new NotImplementedException();
-
- public Task RemoveFailedMessageRetry(string uniqueMessageId, CancellationToken cancellationToken = default) => throw new NotImplementedException();
-
- public Task GetFailedMessageBody(string bodyId, CancellationToken cancellationToken = default) => Task.FromResult(Encoding.UTF8.GetBytes(bodyId));
+ public Task TryFetch(string bodyId, CancellationToken cancellationToken = default) =>
+ Task.FromResult(state switch
+ {
+ MessageBodyState.NotFound => MessageBodyResult.NotFound(),
+ MessageBodyState.Empty => MessageBodyResult.Empty(),
+ MessageBodyState.Unavailable => MessageBodyResult.Unavailable(),
+ MessageBodyState.Available => MessageBodyResult.Available(Content(Encoding.UTF8.GetBytes(bodyId))),
+ _ => throw new ArgumentOutOfRangeException(nameof(state), state, null)
+ });
+
+ static MessageBodyStreamContent Content(byte[] body) => new(new MemoryStream(body), "text/plain", body.Length, "etag");
}
}
}
\ No newline at end of file
diff --git a/src/ServiceControl.Persistence.Tests/RetryStateTests.cs b/src/ServiceControl.Persistence.Tests/RetryStateTests.cs
index 34b3cabb1d..5cdac0108c 100644
--- a/src/ServiceControl.Persistence.Tests/RetryStateTests.cs
+++ b/src/ServiceControl.Persistence.Tests/RetryStateTests.cs
@@ -71,7 +71,7 @@ public async Task When_the_dequeuer_is_created_then_the_error_address_is_cached(
var transportCustomization = new TestTransportCustomization { TransportInfrastructure = transportInfrastructure };
- var testReturnToSenderDequeuer = new TestReturnToSenderDequeuer(new ReturnToSender(FailedMessageRetryStore, NullLogger.Instance), FailedMessageLifecycleStore, domainEvents, "TestEndpoint",
+ var testReturnToSenderDequeuer = new TestReturnToSenderDequeuer(new ReturnToSender(BodyStorage, NullLogger.Instance), FailedMessageLifecycleStore, domainEvents, "TestEndpoint",
errorQueueNameCache, transportCustomization);
await testReturnToSenderDequeuer.StartAsync(new CancellationToken());
@@ -93,7 +93,7 @@ public async Task When_a_group_is_prepared_with_three_batches_and_SC_is_restarte
MessageRedirectsDataStore,
domainEvents,
new TestReturnToSenderDequeuer(
- new ReturnToSender(FailedMessageRetryStore, NullLogger.Instance),
+ new ReturnToSender(BodyStorage, NullLogger.Instance),
FailedMessageLifecycleStore,
domainEvents,
"TestEndpoint",
@@ -121,7 +121,7 @@ public async Task When_a_group_is_prepared_with_three_batches_and_SC_is_restarte
MessageRedirectsDataStore,
domainEvents,
new TestReturnToSenderDequeuer(
- new ReturnToSender(FailedMessageRetryStore, NullLogger.Instance),
+ new ReturnToSender(BodyStorage, NullLogger.Instance),
FailedMessageLifecycleStore,
domainEvents,
"TestEndpoint",
@@ -148,7 +148,7 @@ public async Task When_a_group_is_forwarded_the_status_is_Completed()
var sender = new TestSender();
- var returnToSender = new TestReturnToSenderDequeuer(new ReturnToSender(FailedMessageRetryStore, NullLogger.Instance), FailedMessageLifecycleStore, domainEvents, "TestEndpoint", new ErrorQueueNameCache(), new TestTransportCustomization());
+ 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);
await processor.ProcessBatches(); // mark ready
@@ -213,7 +213,7 @@ public async Task When_there_is_one_poison_message_it_is_removed_from_batch_and_
}
};
- var returnToSender = new TestReturnToSenderDequeuer(new ReturnToSender(FailedMessageRetryStore, NullLogger.Instance), FailedMessageLifecycleStore, domainEvents, "TestEndpoint", new ErrorQueueNameCache(), new TestTransportCustomization());
+ 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);
bool c;
@@ -250,7 +250,7 @@ public async Task When_a_group_has_one_batch_out_of_two_forwarded_the_status_is_
await CreateAFailedMessageAndMarkAsPartOfRetryBatch(retryManager, "Test-group", true, 1001);
- var returnToSender = new ReturnToSender(FailedMessageRetryStore, NullLogger.Instance);
+ var returnToSender = new ReturnToSender(BodyStorage, NullLogger.Instance);
var sender = new TestSender();
@@ -300,7 +300,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(FailedMessageRetryStore, NullLogger.Instance), FailedMessageLifecycleStore, domainEvents, "TestEndpoint", new ErrorQueueNameCache(), new TestTransportCustomization());
+ 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);
await processor.ProcessBatches(); // stage
@@ -327,7 +327,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(FailedMessageRetryStore, NullLogger.Instance), FailedMessageLifecycleStore, domainEvents, "TestEndpoint", new ErrorQueueNameCache(), new TestTransportCustomization());
+ 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);
await processor.ProcessBatches(); // stage (emits per-message audit)
@@ -362,7 +362,7 @@ RetryProcessor CreateProcessor(IDomainEvents domainEvents, TestSender sender) =>
new(RetryStagingStore,
MessageRedirectsDataStore,
domainEvents,
- new TestReturnToSenderDequeuer(new ReturnToSender(FailedMessageRetryStore, NullLogger.Instance), FailedMessageLifecycleStore, domainEvents, "TestEndpoint", new ErrorQueueNameCache(), new TestTransportCustomization()),
+ new TestReturnToSenderDequeuer(new ReturnToSender(BodyStorage, NullLogger.Instance), FailedMessageLifecycleStore, domainEvents, "TestEndpoint", new ErrorQueueNameCache(), new TestTransportCustomization()),
new RetryingManager(domainEvents, NullLogger.Instance),
new Lazy(() => sender),
new RecordingMessageActionAuditLog(),
diff --git a/src/ServiceControl.Persistence/IBodyStorage.cs b/src/ServiceControl.Persistence/IBodyStorage.cs
index 4becd1ea2e..483e347206 100644
--- a/src/ServiceControl.Persistence/IBodyStorage.cs
+++ b/src/ServiceControl.Persistence/IBodyStorage.cs
@@ -1,20 +1,51 @@
namespace ServiceControl.Operations.BodyStorage
{
+ using System;
using System.IO;
using System.Threading;
using System.Threading.Tasks;
public interface IBodyStorage
{
- Task TryFetch(string bodyId, CancellationToken cancellationToken = default);
+ Task TryFetch(string bodyId, CancellationToken cancellationToken = default);
}
- public class MessageBodyStreamResult
+ public enum MessageBodyState
{
- public bool HasResult;
- public Stream Stream; // Intentional, other streams could require a context
- public string ContentType;
- public int BodySize;
- public string Etag;
+ NotFound,
+ Empty,
+ Unavailable,
+ Available
}
+
+ public sealed class MessageBodyResult
+ {
+ MessageBodyResult(MessageBodyState state, MessageBodyStreamContent content = null)
+ {
+ State = state;
+ ContentValue = content;
+ }
+
+ public MessageBodyState State { get; }
+
+ public MessageBodyStreamContent Content => State == MessageBodyState.Available
+ ? ContentValue
+ : throw new InvalidOperationException($"Body content is not available when the state is {State}.");
+
+ public static MessageBodyResult NotFound() => new(MessageBodyState.NotFound);
+
+ public static MessageBodyResult Empty() => new(MessageBodyState.Empty);
+
+ public static MessageBodyResult Unavailable() => new(MessageBodyState.Unavailable);
+
+ public static MessageBodyResult Available(MessageBodyStreamContent content)
+ {
+ ArgumentNullException.ThrowIfNull(content);
+ return new MessageBodyResult(MessageBodyState.Available, content);
+ }
+
+ MessageBodyStreamContent ContentValue { get; }
+ }
+
+ public sealed record MessageBodyStreamContent(Stream Stream, string ContentType, int BodySize, string Etag);
}
\ No newline at end of file
diff --git a/src/ServiceControl.UnitTests/BodyStorage/MessageBodyResultTests.cs b/src/ServiceControl.UnitTests/BodyStorage/MessageBodyResultTests.cs
new file mode 100644
index 0000000000..5c2b04f97a
--- /dev/null
+++ b/src/ServiceControl.UnitTests/BodyStorage/MessageBodyResultTests.cs
@@ -0,0 +1,40 @@
+namespace ServiceControl.UnitTests.BodyStorage;
+
+using System;
+using System.IO;
+using NUnit.Framework;
+using ServiceControl.Operations.BodyStorage;
+
+[TestFixture]
+public class MessageBodyResultTests
+{
+ [TestCase(MessageBodyState.NotFound)]
+ [TestCase(MessageBodyState.Empty)]
+ [TestCase(MessageBodyState.Unavailable)]
+ public void Body_is_not_accessible_without_content(MessageBodyState state)
+ {
+ var result = state switch
+ {
+ MessageBodyState.NotFound => MessageBodyResult.NotFound(),
+ MessageBodyState.Empty => MessageBodyResult.Empty(),
+ MessageBodyState.Unavailable => MessageBodyResult.Unavailable(),
+ MessageBodyState.Available => throw new ArgumentOutOfRangeException(nameof(state), state, null),
+ _ => throw new ArgumentOutOfRangeException(nameof(state), state, null)
+ };
+
+ Assert.Throws(() => _ = result.Content);
+ }
+
+ [Test]
+ public void Body_is_accessible_when_available()
+ {
+ var content = new MessageBodyStreamContent(Stream.Null, "text/plain", 1, "etag");
+ var result = MessageBodyResult.Available(content);
+
+ Assert.That(result.Content, Is.SameAs(content));
+ }
+
+ [Test]
+ public void Available_rejects_null() =>
+ Assert.Throws(() => MessageBodyResult.Available(null));
+}
\ No newline at end of file
diff --git a/src/ServiceControl/CompositeViews/Messages/GetMessagesController.cs b/src/ServiceControl/CompositeViews/Messages/GetMessagesController.cs
index a7d22eb027..e4de1674d7 100644
--- a/src/ServiceControl/CompositeViews/Messages/GetMessagesController.cs
+++ b/src/ServiceControl/CompositeViews/Messages/GetMessagesController.cs
@@ -90,18 +90,18 @@ public async Task Get(string id, [FromQuery(Name = "instance_id")
{
var result = await bodyStorage.TryFetch(id, cancellationToken);
- if (result == null)
+ if (result.State == MessageBodyState.NotFound)
{
return NotFound();
}
- if (!result.HasResult)
+ if (result.State is MessageBodyState.Empty or MessageBodyState.Unavailable)
{
return NoContent();
}
- Response.WithEtag(result.Etag);
- return File(result.Stream, result.ContentType ?? "text/*");
+ Response.WithEtag(result.Content.Etag);
+ return File(result.Content.Stream, result.Content.ContentType ?? "text/*");
}
var remote = settings.RemoteInstances.SingleOrDefault(r => r.InstanceId == instanceId);
diff --git a/src/ServiceControl/Recoverability/Retrying/Infrastructure/ReturnToSender.cs b/src/ServiceControl/Recoverability/Retrying/Infrastructure/ReturnToSender.cs
index 38737f4de1..e1cbd44e36 100644
--- a/src/ServiceControl/Recoverability/Retrying/Infrastructure/ReturnToSender.cs
+++ b/src/ServiceControl/Recoverability/Retrying/Infrastructure/ReturnToSender.cs
@@ -2,14 +2,15 @@ namespace ServiceControl.Recoverability
{
using System;
using System.Collections.Generic;
+ using System.IO;
using System.Threading;
using System.Threading.Tasks;
using Microsoft.Extensions.Logging;
using NServiceBus.Routing;
using NServiceBus.Transport;
- using ServiceControl.Persistence;
+ using Operations.BodyStorage;
- class ReturnToSender(IFailedMessageRetryDataStore errorMessageStore, ILogger logger)
+ class ReturnToSender(IBodyStorage bodyStorage, ILogger logger)
{
public virtual async Task HandleMessage(MessageContext message, IMessageDispatcher sender, string errorQueueTransportAddress, CancellationToken cancellationToken = default)
{
@@ -23,9 +24,9 @@ public virtual async Task HandleMessage(MessageContext message, IMessageDispatch
logger.LogDebug("{MessageId}: Retrieving message body", messageId);
- if (outgoingHeaders.TryGetValue("ServiceControl.Retry.Attempt.MessageId", out var attemptMessageId))
+ if (outgoingHeaders.ContainsKey("ServiceControl.Retry.Attempt.MessageId"))
{
- body = await FetchFromFailedMessage(outgoingHeaders, messageId, attemptMessageId, cancellationToken);
+ body = await FetchFromFailedMessage(outgoingHeaders, messageId, cancellationToken);
outgoingHeaders.Remove("ServiceControl.Retry.Attempt.MessageId");
}
else
@@ -56,21 +57,30 @@ public virtual async Task HandleMessage(MessageContext message, IMessageDispatch
logger.LogDebug("{MessageId}: Forwarded message to {RetryTo}", messageId, retryTo);
}
- async Task FetchFromFailedMessage(Dictionary outgoingHeaders, string messageId, string attemptMessageId, CancellationToken cancellationToken)
+ async Task FetchFromFailedMessage(Dictionary outgoingHeaders, string messageId, CancellationToken cancellationToken)
{
var uniqueMessageId = outgoingHeaders["ServiceControl.Retry.UniqueMessageId"];
- byte[] body = await errorMessageStore.GetFailedMessageBody(uniqueMessageId, cancellationToken);
+ var result = await bodyStorage.TryFetch(uniqueMessageId, cancellationToken);
- if (body == null)
+ if (result.State is MessageBodyState.NotFound or MessageBodyState.Unavailable)
{
- logger.LogWarning("{MessageId}: Message Body not found in failed message with unique id {UniqueMessageId} for attempt Id {AttemptMessageId}", messageId, uniqueMessageId, attemptMessageId);
+ throw new InvalidOperationException($"Cannot retry failed message {uniqueMessageId} because its body state is {result.State}.");
}
- else
+
+ if (result.State == MessageBodyState.Empty)
{
- logger.LogDebug("{MessageId}: Body size: {MessageLength} bytes retrieved from failed message attachment", messageId, body.LongLength);
+ return EmptyBody;
}
- return body;
+ await using (result.Content.Stream)
+ {
+ using var memoryStream = new MemoryStream();
+ await result.Content.Stream.CopyToAsync(memoryStream, cancellationToken);
+ var body = memoryStream.ToArray();
+
+ logger.LogDebug("{MessageId}: Body size: {MessageLength} bytes retrieved from failed message attachment", messageId, body.LongLength);
+ return body;
+ }
}
static readonly byte[] EmptyBody = Array.Empty();