diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/FailedMessageRetryDataStore.cs b/src/ServiceControl.Persistence.EFCore/Implementation/FailedMessageRetryDataStore.cs index 4bef78cb47..4129ceeb0d 100644 --- a/src/ServiceControl.Persistence.EFCore/Implementation/FailedMessageRetryDataStore.cs +++ b/src/ServiceControl.Persistence.EFCore/Implementation/FailedMessageRetryDataStore.cs @@ -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) => @@ -52,31 +51,4 @@ public Task ProcessPendingRetries(DateTime periodFrom, DateTime periodTo, string await processCallback(uniqueMessageId, token); } }, cancellationToken); - - public async Task 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(); - } - } } \ No newline at end of file diff --git a/src/ServiceControl.Persistence.RavenDB/ErrorMessagesDataStore.cs b/src/ServiceControl.Persistence.RavenDB/ErrorMessagesDataStore.cs index fec09b8f8f..a381261bd6 100644 --- a/src/ServiceControl.Persistence.RavenDB/ErrorMessagesDataStore.cs +++ b/src/ServiceControl.Persistence.RavenDB/ErrorMessagesDataStore.cs @@ -29,7 +29,6 @@ class ErrorMessagesDataStore( IRavenSessionProvider sessionProvider, IRavenDocumentStoreProvider documentStoreProvider, - IBodyStorage bodyStorage, ExpirationManager expirationManager, ILogger logger) : IMessagesViewDataStore, IFailedMessageQueryDataStore, IFailedMessageLifecycleDataStore, IFailedMessageRetryDataStore @@ -507,37 +506,5 @@ public async Task GetRetryPendingMessages(DateTime from, DateTime to, } record struct FailedMessageProjection(string UniqueMessageId); - - public async Task 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(); - } - } } } diff --git a/src/ServiceControl.Persistence.Tests/EFCore/FailedMessageRetryBodyDataStoreTests.cs b/src/ServiceControl.Persistence.Tests/EFCore/FailedMessageRetryBodyDataStoreTests.cs deleted file mode 100644 index d964caedec..0000000000 --- a/src/ServiceControl.Persistence.Tests/EFCore/FailedMessageRetryBodyDataStoreTests.cs +++ /dev/null @@ -1,70 +0,0 @@ -namespace ServiceControl.Persistence.Tests; - -using System; -using System.Text; -using System.Threading.Tasks; -using NUnit.Framework; -using ServiceControl.MessageFailures; -using ServiceControl.Persistence.EFCore.Entities; - -class FailedMessageRetryBodyDataStoreTests : ErrorIngestionTestBase -{ - static readonly DateTime Noon = new(2026, 8, 1, 12, 0, 0, DateTimeKind.Utc); - - [Test] - public async Task GetFailedMessageBody_returns_inline_BodyText_bytes() - { - var id = await SeedFailedMessage(bodyText: "hello body"); - - var body = await FailedMessageRetryStore.GetFailedMessageBody(id.ToString()); - - Assert.That(body, Is.EqualTo(Encoding.UTF8.GetBytes("hello body"))); - } - - [Test] - public async Task GetFailedMessageBody_throws_when_the_body_is_unavailable() - { - var id = await SeedFailedMessage(bodyStoredExternally: true); - - Assert.ThrowsAsync(() => - FailedMessageRetryStore.GetFailedMessageBody(id.ToString())); - } - - [Test] - public async Task GetFailedMessageBody_returns_external_storage_body_when_BodyStoredExternally() - { - var expected = Encoding.UTF8.GetBytes("external body payload"); - var id = await SeedFailedMessage(bodyStoredExternally: true); - - await RecordedBodies.WriteBody(id.ToString(), expected, "text/plain"); - - var body = await FailedMessageRetryStore.GetFailedMessageBody(id.ToString()); - - Assert.That(body, Is.EqualTo(expected)); - } - - async Task SeedFailedMessage(string bodyText = null, bool bodyStoredExternally = false) - { - var id = Guid.NewGuid(); - - await Store(new FailedMessageEntity - { - UniqueMessageId = id, - Status = FailedMessageStatus.RetryIssued, - StatusChangedAt = Noon, - LastModified = Noon, - NumberOfProcessingAttempts = 1, - FirstTimeOfFailure = Noon, - LastTimeOfFailure = Noon, - LastAttemptedAt = Noon, - IsSystemMessage = false, - HeadersJson = "{}", - BodyText = bodyText, - BodyStoredExternally = bodyStoredExternally, - BodySize = bodyText is null ? 0 : Encoding.UTF8.GetByteCount(bodyText), - BodyContentType = bodyText is null ? null : "text/plain" - }); - - return id; - } -} diff --git a/src/ServiceControl.Persistence.Tests/Recoverability/FailedMessageRetryDataStoreTests.cs b/src/ServiceControl.Persistence.Tests/Recoverability/FailedMessageRetryDataStoreTests.cs index ea78b02c74..3808635944 100644 --- a/src/ServiceControl.Persistence.Tests/Recoverability/FailedMessageRetryDataStoreTests.cs +++ b/src/ServiceControl.Persistence.Tests/Recoverability/FailedMessageRetryDataStoreTests.cs @@ -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(() => - FailedMessageRetryStore.GetFailedMessageBody(Guid.NewGuid().ToString())); - } - async Task Insert(IngestedFailure failure, FailedMessageStatus status = FailedMessageStatus.Unresolved) { var message = failure.ToFailedMessage(status); diff --git a/src/ServiceControl.Persistence/IFailedMessageRetryDataStore.cs b/src/ServiceControl.Persistence/IFailedMessageRetryDataStore.cs index 8add2b6246..cb29ce3994 100644 --- a/src/ServiceControl.Persistence/IFailedMessageRetryDataStore.cs +++ b/src/ServiceControl.Persistence/IFailedMessageRetryDataStore.cs @@ -9,6 +9,5 @@ public interface IFailedMessageRetryDataStore Task ProcessPendingRetries(DateTime periodFrom, DateTime periodTo, string queueAddress, Func processCallback, CancellationToken cancellationToken = default); Task GetRetryPendingMessages(DateTime from, DateTime to, string queueAddress, CancellationToken cancellationToken = default); Task RemoveFailedMessageRetry(string uniqueMessageId, CancellationToken cancellationToken = default); - Task GetFailedMessageBody(string uniqueMessageId, CancellationToken cancellationToken = default); } } diff --git a/src/ServiceControl.UnitTests/MessageFailures/AsyncRangeAndQueueAuditTests.cs b/src/ServiceControl.UnitTests/MessageFailures/AsyncRangeAndQueueAuditTests.cs index edf12b4bc1..01fc61ccaf 100644 --- a/src/ServiceControl.UnitTests/MessageFailures/AsyncRangeAndQueueAuditTests.cs +++ b/src/ServiceControl.UnitTests/MessageFailures/AsyncRangeAndQueueAuditTests.cs @@ -175,6 +175,5 @@ internal sealed class StubErrorMessageDataStore : IFailedMessageQueryDataStore, public Task MarkAsResolved(string failedMessageId, CancellationToken cancellationToken = default) => throw new NotImplementedException(); public Task ProcessPendingRetries(DateTime periodFrom, DateTime periodTo, string queueAddress, Func processCallback, CancellationToken cancellationToken = default) => throw new NotImplementedException(); public Task RevertRetry(string messageUniqueId, CancellationToken cancellationToken = default) => throw new NotImplementedException(); - public Task GetFailedMessageBody(string uniqueMessageId, CancellationToken cancellationToken = default) => throw new NotImplementedException(); } }