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 @@ -4,10 +4,9 @@ namespace ServiceControl.Persistence.EFCore.Implementation;
using Microsoft.EntityFrameworkCore;
using Microsoft.Extensions.DependencyInjection;
using ServiceControl.MessageFailures;
using ServiceControl.Operations.BodyStorage;
using ServiceControl.Persistence.EFCore.Infrastructure;

public class FailedMessageRetryDataStore(IServiceScopeFactory scopeFactory, IBodyStorage bodyStorage)
public class FailedMessageRetryDataStore(IServiceScopeFactory scopeFactory)
: DataStoreBase(scopeFactory), IFailedMessageRetryDataStore
{
public Task RemoveFailedMessageRetry(string uniqueMessageId, CancellationToken cancellationToken = default) =>
Expand Down Expand Up @@ -52,31 +51,4 @@ public Task ProcessPendingRetries(DateTime periodFrom, DateTime periodTo, string
await processCallback(uniqueMessageId, token);
}
}, cancellationToken);

public async Task<byte[]> GetFailedMessageBody(string uniqueMessageId, CancellationToken cancellationToken = default)
{
var result = await bodyStorage.TryFetch(uniqueMessageId, cancellationToken);

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

if (result.State == MessageBodyState.Empty)
{
return [];
}

await using (result.Content.Stream)
{
using var memoryStream = new MemoryStream();
await result.Content.Stream.CopyToAsync(memoryStream, cancellationToken);
return memoryStream.ToArray();
}
}
}
33 changes: 0 additions & 33 deletions src/ServiceControl.Persistence.RavenDB/ErrorMessagesDataStore.cs
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,6 @@
class ErrorMessagesDataStore(
IRavenSessionProvider sessionProvider,
IRavenDocumentStoreProvider documentStoreProvider,
IBodyStorage bodyStorage,
ExpirationManager expirationManager,
ILogger<ErrorMessagesDataStore> logger)
: IMessagesViewDataStore, IFailedMessageQueryDataStore, IFailedMessageLifecycleDataStore, IFailedMessageRetryDataStore
Expand Down Expand Up @@ -507,37 +506,5 @@ public async Task<string[]> GetRetryPendingMessages(DateTime from, DateTime to,
}

record struct FailedMessageProjection(string UniqueMessageId);

public async Task<byte[]> GetFailedMessageBody(string uniqueMessageId, CancellationToken cancellationToken = default)
{
var result = await bodyStorage.TryFetch(uniqueMessageId, cancellationToken);

if (result.State == MessageBodyState.NotFound)
{
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();
}
}
}
}

This file was deleted.

Original file line number Diff line number Diff line change
Expand Up @@ -131,13 +131,6 @@ await FailedMessageRetryStore.ProcessPendingRetries(WindowStart, WindowEnd, null
Assert.That(captured, Is.EquivalentTo(new[] { shipping, billing }).IgnoreCase);
}

[Test]
public void GetFailedMessageBody_throws_for_a_nonexistent_message()
{
Assert.ThrowsAsync<InvalidOperationException>(() =>
FailedMessageRetryStore.GetFailedMessageBody(Guid.NewGuid().ToString()));
}

async Task<string> Insert(IngestedFailure failure, FailedMessageStatus status = FailedMessageStatus.Unresolved)
{
var message = failure.ToFailedMessage(status);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,5 @@ public interface IFailedMessageRetryDataStore
Task ProcessPendingRetries(DateTime periodFrom, DateTime periodTo, string queueAddress, Func<string, CancellationToken, Task> processCallback, CancellationToken cancellationToken = default);
Task<string[]> GetRetryPendingMessages(DateTime from, DateTime to, string queueAddress, CancellationToken cancellationToken = default);
Task RemoveFailedMessageRetry(string uniqueMessageId, CancellationToken cancellationToken = default);
Task<byte[]> GetFailedMessageBody(string uniqueMessageId, CancellationToken cancellationToken = default);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -175,6 +175,5 @@ internal sealed class StubErrorMessageDataStore : IFailedMessageQueryDataStore,
public Task<bool> MarkAsResolved(string failedMessageId, CancellationToken cancellationToken = default) => throw new NotImplementedException();
public Task ProcessPendingRetries(DateTime periodFrom, DateTime periodTo, string queueAddress, Func<string, CancellationToken, Task> processCallback, CancellationToken cancellationToken = default) => throw new NotImplementedException();
public Task RevertRetry(string messageUniqueId, CancellationToken cancellationToken = default) => throw new NotImplementedException();
public Task<byte[]> GetFailedMessageBody(string uniqueMessageId, CancellationToken cancellationToken = default) => throw new NotImplementedException();
}
}
Loading