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 @@ -52,8 +52,7 @@
<Compile Remove="..\ServiceControl.AcceptanceTests\Recoverability\ExternalIntegration\When_encountered_an_error.cs" />
<Compile Remove="..\ServiceControl.AcceptanceTests\Recoverability\When_edited_message_fails_to_process.cs" />

<!-- EF body retrieval does not distinguish an empty body from unavailable content. -->
<Compile Remove="..\ServiceControl.AcceptanceTests\Recoverability\MessageFailures\When_a_retry_for_a_empty_body_message_is_successful.cs" />
<!-- Remaining EF retry scenarios are deferred for separate investigation. -->
<Compile Remove="..\ServiceControl.AcceptanceTests\Recoverability\When_a_message_is_retried.cs" />
<Compile Remove="..\ServiceControl.AcceptanceTests\Recoverability\When_a_message_is_retried_with_a_replyTo_header.cs" />
<Compile Remove="..\ServiceControl.AcceptanceTests\Recoverability\When_a_native_integration_message_is_retried.cs" />
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -52,8 +52,7 @@
<Compile Remove="..\ServiceControl.AcceptanceTests\Recoverability\ExternalIntegration\When_encountered_an_error.cs" />
<Compile Remove="..\ServiceControl.AcceptanceTests\Recoverability\When_edited_message_fails_to_process.cs" />

<!-- EF body retrieval does not distinguish an empty body from unavailable content. -->
<Compile Remove="..\ServiceControl.AcceptanceTests\Recoverability\MessageFailures\When_a_retry_for_a_empty_body_message_is_successful.cs" />
<!-- Remaining EF retry scenarios are deferred for separate investigation. -->
<Compile Remove="..\ServiceControl.AcceptanceTests\Recoverability\When_a_message_is_retried.cs" />
<Compile Remove="..\ServiceControl.AcceptanceTests\Recoverability\When_a_message_is_retried_with_a_replyTo_header.cs" />
<Compile Remove="..\ServiceControl.AcceptanceTests\Recoverability\When_a_native_integration_message_is_retried.cs" />
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -31,7 +31,7 @@ public async Task SubsequentBatchesShouldBeProcessed(CancellationToken cancellat

CustomizeHostBuilder = hostBuilder =>
{
hostBuilder.Services.AddSingleton<ReturnToSender>(provider => new FakeReturnToSender(provider.GetRequiredService<IFailedMessageRetryDataStore>(), provider.GetRequiredService<MyContext>()));
hostBuilder.Services.AddSingleton<ReturnToSender>(provider => new FakeReturnToSender(provider.GetRequiredService<IBodyStorage>(), provider.GetRequiredService<MyContext>()));
};

await Define<MyContext>()
Expand Down Expand Up @@ -148,8 +148,8 @@ public class MyContext : ScenarioContext

public class MessageThatWillFail : ICommand;

public class FakeReturnToSender(IFailedMessageRetryDataStore errorMessageStore, MyContext myContext)
: ReturnToSender(errorMessageStore, NullLogger<ReturnToSender>.Instance)
public class FakeReturnToSender(IBodyStorage bodyStorage, MyContext myContext)
: ReturnToSender(bodyStorage, NullLogger<ReturnToSender>.Instance)
{
public override Task HandleMessage(MessageContext message, IMessageDispatcher sender, string errorQueueTransportAddress, CancellationToken cancellationToken = default)
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,13 +19,13 @@ namespace ServiceControl.Persistence.EFCore.Implementation.BodyStorage;
/// </remarks>
public class BodyStorage(IServiceScopeFactory scopeFactory, IBodyStoragePersistence storagePersistence) : DataStoreBase(scopeFactory), IBodyStorage
{
public async Task<MessageBodyStreamResult?> TryFetch(string bodyId, CancellationToken cancellationToken = default)
public async Task<MessageBodyResult> 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.
Expand All @@ -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<BodyRow?> ResolveBody(ServiceControlDbContext dbContext, string bodyId, CancellationToken cancellationToken)
Expand All @@ -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);
Expand All @@ -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; }
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -55,18 +55,27 @@ public Task ProcessPendingRetries(DateTime periodFrom, DateTime periodTo, string

public async Task<byte[]> 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();
}
}
Expand Down
42 changes: 25 additions & 17 deletions src/ServiceControl.Persistence.RavenDB/ErrorMessagesDataStore.cs
Original file line number Diff line number Diff line change
Expand Up @@ -510,26 +510,34 @@ record struct FailedMessageProjection(string UniqueMessageId);

public async Task<byte[]> 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;
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@ class RavenAttachmentsBodyStorage(IRavenSessionProvider sessionProvider) : IBody
{
public const string AttachmentName = "body";

public async Task<MessageBodyStreamResult> TryFetch(string bodyId, CancellationToken cancellationToken = default)
public async Task<MessageBodyResult> TryFetch(string bodyId, CancellationToken cancellationToken = default)
{
using var session = await sessionProvider.OpenSession(cancellationToken: cancellationToken);

Expand All @@ -23,7 +23,7 @@ public async Task<MessageBodyStreamResult> 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;
}
Expand All @@ -42,28 +42,37 @@ public async Task<MessageBodyStreamResult> TryFetch(string bodyId, CancellationT
return await ResultForUniqueId(session, uniqueId, cancellationToken);
}

return null;
return MessageBodyResult.NotFound();
}

async Task<MessageBodyStreamResult> ResultForUniqueId(IAsyncDocumentSession session, string uniqueId, CancellationToken cancellationToken)
async Task<MessageBodyResult> ResultForUniqueId(IAsyncDocumentSession session, string uniqueId, CancellationToken cancellationToken)
{
var documentId = FailedMessageIdGenerator.MakeDocumentId(uniqueId);
var failedMessage = await session.LoadAsync<FailedMessage>(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));
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
using NUnit.Framework;
using ServiceControl.MessageFailures;
using ServiceControl.Operations;
using ServiceControl.Operations.BodyStorage;
using ServiceControl.Persistence.UnitOfWork;

[TestFixture]
Expand Down Expand Up @@ -81,14 +82,14 @@ async Task RunTest(Func<Dictionary<string, string>, 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));
Expand Down
26 changes: 13 additions & 13 deletions src/ServiceControl.Persistence.Tests/EFCore/BodyReadTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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("<order>1</order>"));
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("<order>1</order>"));
}
}

Expand All @@ -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));
}
}

Expand All @@ -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");
}
}

Expand All @@ -76,30 +76,30 @@ 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);

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]
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<MessageBodyStreamResult> Fetch(string bodyId)
async Task<MessageBodyResult> Fetch(string bodyId)
{
using var scope = ServiceProvider.CreateScope();
var bodyStorage = scope.ServiceProvider.GetRequiredService<IBodyStorage>();
Expand Down
Loading
Loading