diff --git a/src/ServiceControl.AcceptanceTests.RavenDB/Recoverability/MessageFailures/FailedErrorsController.cs b/src/ServiceControl.AcceptanceTests.RavenDB/Recoverability/MessageFailures/FailedErrorsController.cs index 8f6053a386..a12643282f 100644 --- a/src/ServiceControl.AcceptanceTests.RavenDB/Recoverability/MessageFailures/FailedErrorsController.cs +++ b/src/ServiceControl.AcceptanceTests.RavenDB/Recoverability/MessageFailures/FailedErrorsController.cs @@ -1,5 +1,6 @@ namespace ServiceControl.AcceptanceTests.RavenDB.Recoverability.MessageFailures { + using ServiceControl.Infrastructure; using System.Threading; using System.Threading.Tasks; using Infrastructure.WebApi; diff --git a/src/ServiceControl.AcceptanceTests.RavenDB/Recoverability/MessageFailures/FailedMessageRetriesController.cs b/src/ServiceControl.AcceptanceTests.RavenDB/Recoverability/MessageFailures/FailedMessageRetriesController.cs index 27ab2d4c4e..67ff1ae51f 100644 --- a/src/ServiceControl.AcceptanceTests.RavenDB/Recoverability/MessageFailures/FailedMessageRetriesController.cs +++ b/src/ServiceControl.AcceptanceTests.RavenDB/Recoverability/MessageFailures/FailedMessageRetriesController.cs @@ -1,5 +1,6 @@ namespace ServiceControl.AcceptanceTests.RavenDB.Recoverability.MessageFailures { + using ServiceControl.Infrastructure; using System.Threading; using System.Threading.Tasks; using Infrastructure.WebApi; diff --git a/src/ServiceControl.Audit.AcceptanceTests/ServiceControl.Audit.AcceptanceTests.csproj b/src/ServiceControl.Audit.AcceptanceTests/ServiceControl.Audit.AcceptanceTests.csproj index afab38a784..a36762973f 100644 --- a/src/ServiceControl.Audit.AcceptanceTests/ServiceControl.Audit.AcceptanceTests.csproj +++ b/src/ServiceControl.Audit.AcceptanceTests/ServiceControl.Audit.AcceptanceTests.csproj @@ -23,4 +23,9 @@ + + + + + \ No newline at end of file diff --git a/src/ServiceControl.Audit.AcceptanceTests/WebApi/When_a_list_endpoint_is_revalidated.cs b/src/ServiceControl.Audit.AcceptanceTests/WebApi/When_a_list_endpoint_is_revalidated.cs new file mode 100644 index 0000000000..2cac3878ab --- /dev/null +++ b/src/ServiceControl.Audit.AcceptanceTests/WebApi/When_a_list_endpoint_is_revalidated.cs @@ -0,0 +1,157 @@ +namespace ServiceControl.Audit.AcceptanceTests.WebApi +{ + using System.Net; + using System.Net.Http; + using System.Threading.Tasks; + using AcceptanceTesting; + using AcceptanceTesting.EndpointTemplates; + using Audit.Auditing.MessagesView; + using NServiceBus; + using NServiceBus.AcceptanceTesting; + using NServiceBus.Settings; + using NUnit.Framework; + + class When_a_list_endpoint_is_revalidated : AcceptanceTest + { + [Test] + public async Task Should_answer_not_modified_for_a_list_endpoint() + { + Answer issued = null; + Answer repeated = null; + + await Define() + .WithEndpoint(b => b.When(bus => bus.SendLocal(new MyMessage { Payload = "PAYLOAD" }))) + .Done(async c => + { + if (c.MessageId == null) + { + return false; + } + + // Wait for the message to appear + MessagesView audited = await this.TryGetSingle("/api/messages?include_system_messages=false&sort=id", m => m.MessageId == c.MessageId); + + if (audited == null) + { + return false; + } + + issued = await Ask("GET", "/api/messages?include_system_messages=false&sort=id", ifNoneMatch: null); + + repeated = await Ask("GET", "/api/messages?include_system_messages=false&sort=id", issued.Etag); + + return true; + }) + .Run(); + + Assert.That(issued.Etag, Is.Not.Null, "the list response carried no ETag, so there is nothing for a client to revalidate against"); + Assert.That(repeated.Status, Is.EqualTo(HttpStatusCode.NotModified), $"the list was sent again to a client that already held {issued.Etag}"); + Assert.That(repeated.TotalCount, Is.Not.Null.And.EqualTo(issued.TotalCount), "the 304 did not carry its Total-Count through"); + } + + [Test] + public async Task Should_answer_not_modified_for_an_empty_list_endpoint() + { + Answer issued = null; + Answer repeated = null; + + await Define() + .Done(async c => + { + issued = await Ask("GET", "/api/conversations/no-such-conversation", ifNoneMatch: null); + + repeated = await Ask("GET", "/api/conversations/no-such-conversation", issued.Etag); + + return true; + }) + .Run(); + + Assert.That(issued.Etag, Is.Not.Null, "an empty result should still carry an ETag of its own"); + Assert.That(repeated.Status, Is.EqualTo(HttpStatusCode.NotModified), $"the empty list was sent again to a client that already held {issued.Etag}"); + } + + [Test] + public async Task Audit_counts_should_emit_no_etag() + { + Answer answer = null; + + await Define() + .WithEndpoint(b => b.When(bus => bus.SendLocal(new MyMessage { Payload = "PAYLOAD" }))) + .Done(async c => + { + if (c.MessageId == null) + { + return false; + } + + MessagesView audited = await this.TryGetSingle("/api/messages?include_system_messages=false&sort=id", m => m.MessageId == c.MessageId); + + if (audited == null) + { + return false; + } + + answer = await Ask("GET", $"/api/endpoints/{c.EndpointNameOfReceivingEndpoint}/audit-count", ifNoneMatch: null); + + return true; + }) + .Run(); + + Assert.That(answer.Etag, Is.Null, "audit counts were not previously versioned and should not gain an ETag"); + } + + async Task Ask(string method, string url, string ifNoneMatch) + { + using var response = await Send(method, url, ifNoneMatch); + + return new Answer( + response.StatusCode, + Header(response, "ETag"), + Header(response, "Total-Count")); + } + + static string Header(HttpResponseMessage response, string name) => + response.Headers.TryGetValues(name, out var values) ? string.Join(string.Empty, values) : null; + + Task Send(string method, string url, string ifNoneMatch) + { + var request = new HttpRequestMessage(new HttpMethod(method), url); + + if (ifNoneMatch != null) + { + request.Headers.TryAddWithoutValidation("If-None-Match", ifNoneMatch); + } + + return HttpClient.SendAsync(request); + } + + record Answer(HttpStatusCode Status, string Etag, string TotalCount); + + public class Receiver : EndpointConfigurationBuilder + { + public Receiver() => EndpointSetup(); + + [Handler] + public class MyMessageHandler(MyContext testContext, IReadOnlySettings settings) : IHandleMessages + { + public Task Handle(MyMessage message, IMessageHandlerContext context) + { + testContext.EndpointNameOfReceivingEndpoint = settings.EndpointName(); + testContext.MessageId = context.MessageId; + return Task.CompletedTask; + } + } + } + + public class MyMessage : ICommand + { + public string Payload { get; set; } + } + + public class MyContext : ScenarioContext + { + public string MessageId { get; set; } + public string EndpointNameOfReceivingEndpoint { get; set; } + } + } +} \ No newline at end of file diff --git a/src/ServiceControl.Audit.AcceptanceTests/WebApi/When_a_message_body_is_requested_twice.cs b/src/ServiceControl.Audit.AcceptanceTests/WebApi/When_a_message_body_is_requested_twice.cs index 50fb1e4bfd..ed05bc95e6 100644 --- a/src/ServiceControl.Audit.AcceptanceTests/WebApi/When_a_message_body_is_requested_twice.cs +++ b/src/ServiceControl.Audit.AcceptanceTests/WebApi/When_a_message_body_is_requested_twice.cs @@ -58,9 +58,61 @@ await Define() .Run(); Assert.That(issued, Is.Not.Null, "the body response carried no validator, so a client can never revalidate it"); + Assert.That(issued, Does.StartWith("W/"), "the body ETag must be weak-tagged (W/...) so it carries clients through the upgrade"); Assert.That(repeated, Is.EqualTo(HttpStatusCode.NotModified), $"the body was sent again to a client that already held {issued}"); } + [Test] + public async Task Should_answer_not_modified_for_an_old_strong_tag() + { + string issued = null; + HttpStatusCode? repeated = null; + + await Define() + .WithEndpoint(b => b.When(bus => bus.SendLocal(new MyMessage { Payload = "PAYLOAD" }))) + .Done(async c => + { + if (c.MessageId == null) + { + return false; + } + + MessagesView audited = await this.TryGetSingle("/api/messages?include_system_messages=false&sort=id", m => m.MessageId == c.MessageId); + + if (audited == null) + { + return false; + } + + var url = $"/api{audited.BodyUrl}"; + + using var first = await this.GetRaw(url); + + if (!first.Headers.TryGetValues("ETag", out var values)) + { + return false; + } + + issued = string.Join(string.Empty, values); + + // Strip the W/ prefix to simulate an old strong tag from a pre-upgrade client + var strongTag = issued.StartsWith("W/") ? issued[2..] : issued; + + var request = new HttpRequestMessage(HttpMethod.Get, url); + request.Headers.TryAddWithoutValidation("If-None-Match", strongTag); + + using var second = await HttpClient.SendAsync(request); + repeated = second.StatusCode; + + return true; + }) + .Run(); + + Assert.That(issued, Is.Not.Null, "the body response carried no validator, so a client can never revalidate it"); + Assert.That(repeated, Is.EqualTo(HttpStatusCode.NotModified), + "weak comparison ignores the W/ prefix, so an old strong tag must still match the new weak tag"); + } + public class Receiver : EndpointConfigurationBuilder { public Receiver() => EndpointSetup(); @@ -89,4 +141,4 @@ public class MyContext : ScenarioContext public string EndpointNameOfReceivingEndpoint { get; set; } } } -} +} \ No newline at end of file diff --git a/src/ServiceControl.Audit.Persistence.InMemory/InMemoryAttachmentsBodyStorage.cs b/src/ServiceControl.Audit.Persistence.InMemory/InMemoryAttachmentsBodyStorage.cs index 9fbf8f5d34..c25dc228cc 100644 --- a/src/ServiceControl.Audit.Persistence.InMemory/InMemoryAttachmentsBodyStorage.cs +++ b/src/ServiceControl.Audit.Persistence.InMemory/InMemoryAttachmentsBodyStorage.cs @@ -1,11 +1,14 @@ namespace ServiceControl.Audit.Persistence.InMemory { + using System; using System.Collections.Generic; using System.IO; using System.Linq; using System.Threading; using System.Threading.Tasks; using ServiceControl.Audit.Auditing.BodyStorage; + using ServiceControl.Audit.Persistence; + using ServiceControl.Infrastructure; class InMemoryAttachmentsBodyStorage : IBodyStorage { @@ -35,6 +38,7 @@ public Task Store(string bodyId, string contentType, int bodySize, Stream bodySt } messageBody.ContentType = contentType; + messageBody.Version = DataVersion.FromToken(Guid.NewGuid().ToString("N")); if (needToAdd) { @@ -44,25 +48,17 @@ public Task Store(string bodyId, string contentType, int bodySize, Stream bodySt return Task.CompletedTask; } - public async Task TryFetch(string bodyId, CancellationToken cancellationToken = default) + public Task TryFetch(string bodyId, CancellationToken cancellationToken = default) { var messageBody = messageBodies.FirstOrDefault(w => w.BodyId == bodyId); - return await Task.FromResult(messageBody == null - ? new StreamResult - { - HasResult = false, - Stream = null - } - : new StreamResult - { - HasResult = true, - Stream = new MemoryStream(messageBody.Content), - ContentType = messageBody.ContentType, - BodySize = messageBody.BodySize, - // Bodies are immutable per message, so the id is a stable validator - Etag = bodyId - }); + return Task.FromResult(messageBody == null + ? MessageBodyView.NotFound() + : MessageBodyView.FromStream( + new MemoryStream(messageBody.Content), + messageBody.ContentType, + messageBody.BodySize, + messageBody.Version)); } class MessageBody @@ -71,6 +67,7 @@ class MessageBody public string ContentType { get; set; } public int BodySize { get; set; } public byte[] Content { get; set; } + public DataVersion Version { get; set; } } } } \ No newline at end of file diff --git a/src/ServiceControl.Audit.Persistence.InMemory/InMemoryAuditDataStore.cs b/src/ServiceControl.Audit.Persistence.InMemory/InMemoryAuditDataStore.cs index 6e4b5ed3c0..a5695b45f4 100644 --- a/src/ServiceControl.Audit.Persistence.InMemory/InMemoryAuditDataStore.cs +++ b/src/ServiceControl.Audit.Persistence.InMemory/InMemoryAuditDataStore.cs @@ -9,6 +9,8 @@ using ServiceControl.Audit.Auditing.BodyStorage; using ServiceControl.Audit.Auditing.MessagesView; using ServiceControl.Audit.Infrastructure; + using ServiceControl.Audit.Persistence; + using ServiceControl.Infrastructure; using ServiceControl.SagaAudit; class InMemoryAuditDataStore : IAuditDataStore @@ -34,7 +36,7 @@ public async Task> QuerySagaHistoryById(Guid input, Can return await Task.FromResult(QueryResult.Empty()); } - return await Task.FromResult(new QueryResult(sagaHistory, new QueryStatsInfo(string.Empty, 1))); + return await Task.FromResult(new QueryResult(sagaHistory, new QueryStatsInfo(DataVersion.None, 1))); } public async Task>> GetMessages(bool includeSystemMessages, PagingInfo pagingInfo, SortInfo sortInfo, DateTimeRange timeSentRange, CancellationToken cancellationToken = default) @@ -45,7 +47,7 @@ public async Task>> GetMessages(bool includeSyst (timeSentRange == null || !timeSentRange.To.HasValue || w.TimeSent <= timeSentRange.To.Value)) .ToList(); - return await Task.FromResult(new QueryResult>(matched, new QueryStatsInfo(string.Empty, matched.Count))); + return await Task.FromResult(new QueryResult>(matched, new QueryStatsInfo(DataVersion.None, matched.Count))); } public async Task>> QueryMessages(string keyword, PagingInfo pagingInfo, SortInfo sortInfo, DateTimeRange timeSentRange, CancellationToken cancellationToken = default) @@ -57,7 +59,7 @@ public async Task>> QueryMessages(string keyword (timeSentRange == null || !timeSentRange.From.HasValue || w.TimeSent >= timeSentRange.From.Value) && (timeSentRange == null || !timeSentRange.To.HasValue || w.TimeSent <= timeSentRange.To.Value)) .ToList(); - return await Task.FromResult(new QueryResult>(matched, new QueryStatsInfo(string.Empty, matched.Count()))); + return await Task.FromResult(new QueryResult>(matched, new QueryStatsInfo(DataVersion.None, matched.Count()))); } public async Task>> QueryMessagesByReceivingEndpointAndKeyword(string endpoint, string keyword, PagingInfo pagingInfo, SortInfo sortInfo, DateTimeRange timeSentRange, CancellationToken cancellationToken = default) @@ -68,7 +70,7 @@ public async Task>> QueryMessagesByReceivingEndp (timeSentRange == null || !timeSentRange.From.HasValue || w.TimeSent >= timeSentRange.From.Value) && (timeSentRange == null || !timeSentRange.To.HasValue || w.TimeSent <= timeSentRange.To.Value)) .ToList(); - return await Task.FromResult(new QueryResult>(matched, new QueryStatsInfo(string.Empty, matched.Count))); + return await Task.FromResult(new QueryResult>(matched, new QueryStatsInfo(DataVersion.None, matched.Count))); } public async Task>> QueryMessagesByReceivingEndpoint(bool includeSystemMessages, string endpointName, PagingInfo pagingInfo, SortInfo sortInfo, DateTimeRange timeSentRange, CancellationToken cancellationToken = default) @@ -77,13 +79,13 @@ public async Task>> QueryMessagesByReceivingEndp (timeSentRange == null || !timeSentRange.From.HasValue || w.TimeSent >= timeSentRange.From.Value) && (timeSentRange == null || !timeSentRange.To.HasValue || w.TimeSent <= timeSentRange.To.Value)) .ToList(); - return await Task.FromResult(new QueryResult>(matched, new QueryStatsInfo(string.Empty, matched.Count))); + return await Task.FromResult(new QueryResult>(matched, new QueryStatsInfo(DataVersion.None, matched.Count))); } public async Task>> QueryMessagesByConversationId(string conversationId, PagingInfo pagingInfo, SortInfo sortInfo, CancellationToken cancellationToken = default) { var matched = messageViews.Where(w => w.ConversationId == conversationId).ToList(); - return await Task.FromResult(new QueryResult>(matched, new QueryStatsInfo(string.Empty, matched.Count))); + return await Task.FromResult(new QueryResult>(matched, new QueryStatsInfo(DataVersion.None, matched.Count))); } public async Task GetMessageBody(string messageId, CancellationToken cancellationToken = default) @@ -127,21 +129,10 @@ IList GetMessageIdsMatchingQuery(string keyword) .Select(pm => pm.MessageMetadata["MessageId"] as string) .ToList(); } + async Task GetMessageBodyFromAttachments(string messageId, CancellationToken cancellationToken) { - var fromBodyStorage = await bodyStorage.TryFetch(messageId, cancellationToken); - - if (fromBodyStorage.HasResult) - { - return MessageBodyView.FromStream( - fromBodyStorage.Stream, - fromBodyStorage.ContentType, - fromBodyStorage.BodySize, - fromBodyStorage.Etag - ); - } - - return MessageBodyView.NotFound(); + return await bodyStorage.TryFetch(messageId, cancellationToken); } Task GetMessageBodyFromMetadata(string messageId, CancellationToken cancellationToken) @@ -168,7 +159,7 @@ Task GetMessageBodyFromMetadata(string messageId, CancellationT return Task.FromResult(MessageBodyView.NotFound()); } - return Task.FromResult(MessageBodyView.FromString(body, contentType, bodySize, messageId)); + return Task.FromResult(MessageBodyView.FromString(body, contentType, bodySize, DataVersion.FromToken(messageId))); } public Task>> QueryAuditCounts(string endpointName, CancellationToken cancellationToken = default) diff --git a/src/ServiceControl.Audit.Persistence.InMemory/ServiceControl.Audit.Persistence.InMemory.csproj b/src/ServiceControl.Audit.Persistence.InMemory/ServiceControl.Audit.Persistence.InMemory.csproj index f4f86456e4..8184641fa3 100644 --- a/src/ServiceControl.Audit.Persistence.InMemory/ServiceControl.Audit.Persistence.InMemory.csproj +++ b/src/ServiceControl.Audit.Persistence.InMemory/ServiceControl.Audit.Persistence.InMemory.csproj @@ -10,6 +10,7 @@ + diff --git a/src/ServiceControl.Audit.Persistence.RavenDB/Extensions/RavenQueryStatisticsExtensions.cs b/src/ServiceControl.Audit.Persistence.RavenDB/Extensions/RavenQueryStatisticsExtensions.cs index 56629a265c..b0d02eb695 100644 --- a/src/ServiceControl.Audit.Persistence.RavenDB/Extensions/RavenQueryStatisticsExtensions.cs +++ b/src/ServiceControl.Audit.Persistence.RavenDB/Extensions/RavenQueryStatisticsExtensions.cs @@ -1,12 +1,14 @@ -namespace ServiceControl.Audit.Persistence.RavenDB.Extensions +namespace ServiceControl.Audit.Persistence.RavenDB.Extensions { using System.Globalization; using Auditing.MessagesView; using Raven.Client.Documents.Session; + using ServiceControl.Infrastructure; static class RavenQueryStatisticsExtensions { public static QueryStatsInfo ToQueryStatsInfo(this QueryStatistics stats) => - new(stats.ResultEtag?.ToString(CultureInfo.InvariantCulture) ?? string.Empty, stats.TotalResults); + new(stats.ResultEtag is { } resultEtag ? DataVersion.FromToken(resultEtag) : DataVersion.None, + stats.TotalResults); } -} +} \ No newline at end of file diff --git a/src/ServiceControl.Audit.Persistence.RavenDB/RavenAttachmentsBodyStorage.cs b/src/ServiceControl.Audit.Persistence.RavenDB/RavenAttachmentsBodyStorage.cs index 441212614c..86dd91867c 100644 --- a/src/ServiceControl.Audit.Persistence.RavenDB/RavenAttachmentsBodyStorage.cs +++ b/src/ServiceControl.Audit.Persistence.RavenDB/RavenAttachmentsBodyStorage.cs @@ -5,6 +5,7 @@ using System.Threading.Tasks; using Auditing.BodyStorage; using Raven.Client.Documents.BulkInsert; + using ServiceControl.Infrastructure; class RavenAttachmentsBodyStorage( IRavenSessionProvider sessionProvider, @@ -23,24 +24,21 @@ public Task Store(string bodyId, string contentType, int bodySize, Stream bodySt .StoreAsync("body", bodyStream, contentType, cancellationToken); } - 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); var result = await session.Advanced.Attachments.GetAsync($"MessageBodies/{bodyId}", "body", cancellationToken); if (result == null) { - return new StreamResult { HasResult = false }; + return MessageBodyView.NotFound(); } - return new StreamResult - { - HasResult = true, - Stream = result.Stream, - BodySize = (int)result.Details.Size, - ContentType = result.Details.ContentType, - Etag = result.Details.ChangeVector - }; + return MessageBodyView.FromStream( + result.Stream, + result.Details.ContentType, + (int)result.Details.Size, + DataVersion.FromToken(result.Details.ChangeVector)); } } } \ No newline at end of file diff --git a/src/ServiceControl.Audit.Persistence.RavenDB/RavenAuditDataStore.cs b/src/ServiceControl.Audit.Persistence.RavenDB/RavenAuditDataStore.cs index 18e2aa7706..3f7c6ba091 100644 --- a/src/ServiceControl.Audit.Persistence.RavenDB/RavenAuditDataStore.cs +++ b/src/ServiceControl.Audit.Persistence.RavenDB/RavenAuditDataStore.cs @@ -124,7 +124,7 @@ public async Task GetMessageBody(string messageId, Cancellation result.Stream, result.Details.ContentType, (int)result.Details.Size, - result.Details.ChangeVector + DataVersion.FromToken(result.Details.ChangeVector) ); } diff --git a/src/ServiceControl.Audit.Persistence.Tests/AuditTests.cs b/src/ServiceControl.Audit.Persistence.Tests/AuditTests.cs index b9a1e8b04c..05ee2795c1 100644 --- a/src/ServiceControl.Audit.Persistence.Tests/AuditTests.cs +++ b/src/ServiceControl.Audit.Persistence.Tests/AuditTests.cs @@ -102,7 +102,7 @@ public async Task Can_roundtrip_message_body() Assert.That(retrievedMessage.Found, Is.True); Assert.That(retrievedMessage.HasContent, Is.True); Assert.That(retrievedMessage.ContentLength, Is.EqualTo(body.Length)); - Assert.That(retrievedMessage.ETag, Is.Not.Null.And.Not.Empty); + Assert.That(retrievedMessage.Version.HasValue, Is.True); Assert.That(retrievedMessage.StreamContent, Is.Not.Null); Assert.That(retrievedMessage.ContentType, Is.EqualTo(expectedContentType)); } @@ -134,8 +134,8 @@ public async Task Message_body_validator_is_stable_across_reads() var first = await DataStore.GetMessageBody(bodyId, TestContext.CurrentContext.CancellationToken); var second = await DataStore.GetMessageBody(bodyId, TestContext.CurrentContext.CancellationToken); - Assert.That(first.ETag, Is.Not.Null.And.Not.Empty, "a body with no validator cannot be revalidated, so conditional GET is dead on it"); - Assert.That(second.ETag, Is.EqualTo(first.ETag), "the body did not change, so neither may its validator"); + Assert.That(first.Version.HasValue, Is.True, "a body with no validator cannot be revalidated, so conditional GET is dead on it"); + Assert.That(second.Version, Is.EqualTo(first.Version), "the body did not change, so neither may its validator"); } [Test] diff --git a/src/ServiceControl.Audit.Persistence.Tests/BodyStorageTests.cs b/src/ServiceControl.Audit.Persistence.Tests/BodyStorageTests.cs index 2f63b14c4c..b7e33b9461 100644 --- a/src/ServiceControl.Audit.Persistence.Tests/BodyStorageTests.cs +++ b/src/ServiceControl.Audit.Persistence.Tests/BodyStorageTests.cs @@ -13,7 +13,7 @@ public async Task Handles_no_results_gracefully() var nonExistentBodyId = Guid.NewGuid().ToString(); var result = await BodyStorage.TryFetch(nonExistentBodyId, TestContext.CurrentContext.CancellationToken); - Assert.That(result.HasResult, Is.False); + Assert.That(result.Found, Is.False); } } } \ No newline at end of file diff --git a/src/ServiceControl.Audit.Persistence.Tests/InMemory/BodyVersionTests.cs b/src/ServiceControl.Audit.Persistence.Tests/InMemory/BodyVersionTests.cs new file mode 100644 index 0000000000..34dfef456a --- /dev/null +++ b/src/ServiceControl.Audit.Persistence.Tests/InMemory/BodyVersionTests.cs @@ -0,0 +1,43 @@ +namespace ServiceControl.Audit.Persistence.Tests +{ + using System; + using System.IO; + using System.Threading.Tasks; + using NUnit.Framework; + using ServiceControl.Audit.Infrastructure; + + [TestFixture] + class BodyVersionTests : PersistenceTestFixture + { + [Test] + public async Task In_memory_list_queries_remain_unversioned() + { + var result = await MessagesViewStore.GetMessages(false, new PagingInfo(), new SortInfo("message_id", "asc")); + + Assert.That(result.QueryStats.Version.HasValue, Is.False); + } + + [Test] + public async Task Replacing_a_body_changes_its_validator_without_changing_reads() + { + var bodyId = Guid.NewGuid().ToString(); + await BodyStorage.Store(bodyId, "text/plain", 1, new MemoryStream([1])); + var first = await BodyStorage.TryFetch(bodyId); + var repeated = await BodyStorage.TryFetch(bodyId); + + await BodyStorage.Store(bodyId, "text/plain", 1, new MemoryStream([2])); + var replaced = await BodyStorage.TryFetch(bodyId); + + using (Assert.EnterMultipleScope()) + { + Assert.That(first.Version.HasValue, Is.True); + Assert.That(repeated.Version, Is.EqualTo(first.Version)); + Assert.That(replaced.Version, Is.Not.EqualTo(first.Version)); + } + + first.StreamContent.Dispose(); + repeated.StreamContent.Dispose(); + replaced.StreamContent.Dispose(); + } + } +} diff --git a/src/ServiceControl.Audit.Persistence/IBodyStorage.cs b/src/ServiceControl.Audit.Persistence/IBodyStorage.cs index 98b6c42e84..2ea1d4bcbc 100644 --- a/src/ServiceControl.Audit.Persistence/IBodyStorage.cs +++ b/src/ServiceControl.Audit.Persistence/IBodyStorage.cs @@ -1,21 +1,13 @@ -namespace ServiceControl.Audit.Auditing.BodyStorage +namespace ServiceControl.Audit.Auditing.BodyStorage { using System.IO; using System.Threading; using System.Threading.Tasks; + using ServiceControl.Audit.Persistence; public interface IBodyStorage { Task Store(string bodyId, string contentType, int bodySize, Stream bodyStream, CancellationToken cancellationToken = default); - Task TryFetch(string bodyId, CancellationToken cancellationToken = default); - } - - public class StreamResult - { - public bool HasResult; - public Stream Stream; - public string ContentType; - public int BodySize; - public string Etag; + Task TryFetch(string bodyId, CancellationToken cancellationToken = default); } } \ No newline at end of file diff --git a/src/ServiceControl.Audit.Persistence/Infrastructure/QueryStatsInfo.cs b/src/ServiceControl.Audit.Persistence/Infrastructure/QueryStatsInfo.cs index 7e21326bd6..f44d64806d 100644 --- a/src/ServiceControl.Audit.Persistence/Infrastructure/QueryStatsInfo.cs +++ b/src/ServiceControl.Audit.Persistence/Infrastructure/QueryStatsInfo.cs @@ -1,18 +1,20 @@ namespace ServiceControl.Audit.Auditing.MessagesView { - public struct QueryStatsInfo + using ServiceControl.Infrastructure; + + public readonly struct QueryStatsInfo { - public readonly string ETag; + public readonly DataVersion Version; public readonly long TotalCount; public readonly long HighestTotalCountOfAllTheInstances; - public QueryStatsInfo(string eTag, long totalCount, long? highestTotalCountOfAllTheInstances = null) + public QueryStatsInfo(DataVersion version, long totalCount, long? highestTotalCountOfAllTheInstances = null) { - ETag = eTag; + Version = version; TotalCount = totalCount; HighestTotalCountOfAllTheInstances = highestTotalCountOfAllTheInstances ?? totalCount; } - public static readonly QueryStatsInfo Zero = new QueryStatsInfo(string.Empty, 0); + public static readonly QueryStatsInfo Zero = new(DataVersion.None, 0); } } \ No newline at end of file diff --git a/src/ServiceControl.Audit.Persistence/MessageBodyView.cs b/src/ServiceControl.Audit.Persistence/MessageBodyView.cs index 52991bdec8..357a56dd7d 100644 --- a/src/ServiceControl.Audit.Persistence/MessageBodyView.cs +++ b/src/ServiceControl.Audit.Persistence/MessageBodyView.cs @@ -1,6 +1,7 @@ -namespace ServiceControl.Audit.Persistence +namespace ServiceControl.Audit.Persistence { using System.IO; + using ServiceControl.Infrastructure; public class MessageBodyView { @@ -10,13 +11,13 @@ public class MessageBodyView public string StringContent { get; private set; } public string ContentType { get; private set; } public int ContentLength { get; private set; } - public string ETag { get; private set; } + public DataVersion Version { get; private set; } public static MessageBodyView NotFound() => new MessageBodyView { Found = false, HasContent = false }; public static MessageBodyView NoContent() => new MessageBodyView { Found = true, HasContent = false }; - public static MessageBodyView FromString(string content, string contentType, int contentLength, string etag) + public static MessageBodyView FromString(string content, string contentType, int contentLength, DataVersion version) => new MessageBodyView { Found = true, @@ -24,10 +25,10 @@ public static MessageBodyView FromString(string content, string contentType, int StringContent = content, ContentType = contentType, ContentLength = contentLength, - ETag = etag + Version = version }; - public static MessageBodyView FromStream(Stream content, string contentType, int contentLength, string etag) + public static MessageBodyView FromStream(Stream content, string contentType, int contentLength, DataVersion version) => new MessageBodyView { Found = true, @@ -35,7 +36,7 @@ public static MessageBodyView FromStream(Stream content, string contentType, int StreamContent = content, ContentType = contentType, ContentLength = contentLength, - ETag = etag + Version = version }; } } \ No newline at end of file diff --git a/src/ServiceControl.Audit.UnitTests/BodyStorage/BodyStorageEnricherTests.cs b/src/ServiceControl.Audit.UnitTests/BodyStorage/BodyStorageEnricherTests.cs index d559fdcf4e..2ba04f03b6 100644 --- a/src/ServiceControl.Audit.UnitTests/BodyStorage/BodyStorageEnricherTests.cs +++ b/src/ServiceControl.Audit.UnitTests/BodyStorage/BodyStorageEnricherTests.cs @@ -268,7 +268,7 @@ public Task Store(string bodyId, string contentType, int bodySize, Stream bodySt return Task.CompletedTask; } - public Task TryFetch(string bodyId, CancellationToken cancellationToken = default) + public Task TryFetch(string bodyId, CancellationToken cancellationToken = default) { throw new NotImplementedException(); } diff --git a/src/ServiceControl.Audit.UnitTests/Infrastructure/WebApi/ConditionalGetTests.cs b/src/ServiceControl.Audit.UnitTests/Infrastructure/WebApi/ConditionalGetTests.cs index 621bfd9426..457628445b 100644 --- a/src/ServiceControl.Audit.UnitTests/Infrastructure/WebApi/ConditionalGetTests.cs +++ b/src/ServiceControl.Audit.UnitTests/Infrastructure/WebApi/ConditionalGetTests.cs @@ -1,13 +1,21 @@ namespace ServiceControl.Audit.UnitTests.Infrastructure.WebApi { + using System; + using System.Collections.Generic; + using System.IO; + using System.Threading.Tasks; using System.Net; + using Audit.Auditing.MessagesView; + using Audit.Infrastructure; using Audit.Infrastructure.WebApi; using Microsoft.AspNetCore.Http; + using Microsoft.AspNetCore.Http.Features; using Microsoft.AspNetCore.Mvc; using Microsoft.AspNetCore.Mvc.Abstractions; using Microsoft.AspNetCore.Mvc.Filters; using Microsoft.AspNetCore.Routing; using NUnit.Framework; + using ServiceControl.Infrastructure; [TestFixture] public class ConditionalGetTests @@ -17,7 +25,7 @@ public void Repeating_a_request_with_the_etag_just_issued_is_not_modified() { var httpContext = new DefaultHttpContext(); - httpContext.Response.WithEtag("4611686018427387904"); + httpContext.Response.WithEtag(DataVersion.FromToken("4611686018427387904")); httpContext.Request.Headers.IfNoneMatch = httpContext.Response.Headers.ETag; var context = ResultExecuting(httpContext); @@ -34,7 +42,7 @@ public void A_different_etag_still_returns_the_payload() { var httpContext = new DefaultHttpContext(); - httpContext.Response.WithEtag("4611686018427387904"); + httpContext.Response.WithEtag(DataVersion.FromToken("4611686018427387904")); httpContext.Request.Headers.IfNoneMatch = "\"something-else\""; var context = ResultExecuting(httpContext); @@ -49,42 +57,144 @@ public void The_emitted_etag_is_a_well_formed_entity_tag() { var httpContext = new DefaultHttpContext(); - httpContext.Response.WithEtag("4611686018427387904"); + httpContext.Response.WithEtag(DataVersion.FromToken("4611686018427387904")); Assert.That(httpContext.Response.GetTypedHeaders().ETag, Is.Not.Null, "an ETag that cannot be parsed as an entity-tag disables conditional GET without any error"); } [Test] - public void A_deterministic_etag_is_a_well_formed_entity_tag() + public void Every_emitted_etag_is_marked_weak() { var httpContext = new DefaultHttpContext(); - httpContext.Response.WithDeterministicEtag("any-non-empty-payload-signature"); + httpContext.Response.WithEtag(DataVersion.FromToken("4611686018427387904")); - Assert.That(httpContext.Response.GetTypedHeaders().ETag, Is.Not.Null); + using (Assert.EnterMultipleScope()) + { + Assert.That(httpContext.Response.Headers.ETag.ToString(), Is.EqualTo("W/\"4611686018427387904\"")); + Assert.That(httpContext.Response.GetTypedHeaders().ETag.IsWeak, Is.True); + } } [Test] - public void The_emitted_etag_quotes_the_value_without_altering_it() + public void An_absent_data_version_emits_no_header() { var httpContext = new DefaultHttpContext(); - httpContext.Response.WithEtag("4611686018427387904"); + httpContext.Response.WithEtag(DataVersion.None); - Assert.That(httpContext.Response.Headers.ETag.ToString(), Is.EqualTo("\"4611686018427387904\"")); + Assert.That(httpContext.Response.Headers.ContainsKey("ETag"), Is.False, + "an empty entity-tag is well formed, so it would match itself and answer 304 for unrelated payloads"); } - [TestCase(null)] - [TestCase("")] - public void A_call_site_with_nothing_to_validate_emits_no_etag_header(string value) + [Test] + public void An_unmarked_validator_from_an_older_client_still_matches() { var httpContext = new DefaultHttpContext(); - httpContext.Response.WithEtag(value); + httpContext.Response.WithEtag(DataVersion.FromToken("4611686018427387904")); + httpContext.Request.Headers.IfNoneMatch = "\"4611686018427387904\""; - Assert.That(httpContext.Response.Headers.ContainsKey("ETag"), Is.False, - "an empty entity-tag is well formed, so it would match itself and answer 304 for unrelated payloads"); + var context = ResultExecuting(httpContext); + + new NotModifiedStatusHttpHandler().OnResultExecuting(context); + + Assert.That(context.Result, Is.InstanceOf(), + "weak comparison ignores the W/ prefix, which is what carries a client through the upgrade"); + } + + [Test] + public void A_wildcard_precondition_is_not_modified_when_a_representation_exists() + { + var httpContext = new DefaultHttpContext(); + + httpContext.Response.WithEtag(DataVersion.FromToken("4611686018427387904")); + httpContext.Request.Headers.IfNoneMatch = "*"; + + var context = ResultExecuting(httpContext); + + new NotModifiedStatusHttpHandler().OnResultExecuting(context); + + Assert.That(context.Result, Is.InstanceOf(), + "RFC 9110: * matches whenever a current representation exists"); + } + + [Test] + public void A_wildcard_precondition_is_ignored_when_there_is_no_validator() + { + var httpContext = new DefaultHttpContext(); + + httpContext.Request.Headers.IfNoneMatch = "*"; + + var context = ResultExecuting(httpContext); + + new NotModifiedStatusHttpHandler().OnResultExecuting(context); + + Assert.That(context.Result, Is.InstanceOf(), + "an endpoint that publishes no validator has nothing for a client to have cached"); + } + + [Test] + public void A_failed_result_cannot_be_replaced_by_not_modified() + { + var httpContext = new DefaultHttpContext(); + httpContext.Response.WithEtag(DataVersion.FromToken("version")); + httpContext.Request.Headers.IfNoneMatch = "*"; + var context = ResultExecuting(httpContext); + context.Result = new NotFoundResult(); + + new NotModifiedStatusHttpHandler().OnResultExecuting(context); + + Assert.That(context.Result, Is.InstanceOf()); + } + + [Test] + public async Task A_file_stream_is_disposed_when_replaced() + { + var httpContext = new DefaultHttpContext(); + var responseFeature = new CompletingResponseFeature(); + httpContext.Features.Set(responseFeature); + httpContext.Response.WithEtag(DataVersion.FromToken("version")); + httpContext.Request.Headers.IfNoneMatch = "*"; + var stream = new MemoryStream(); + var context = ResultExecuting(httpContext); + context.Result = new FileStreamResult(stream, "application/octet-stream"); + + new NotModifiedStatusHttpHandler().OnResultExecuting(context); + + Assert.That(context.Result, Is.InstanceOf()); + await responseFeature.Complete(); + Assert.That(stream.CanRead, Is.False); + } + + [Test] + public void A_paged_endpoint_emits_the_store_version_rather_than_a_hash_of_it() + { + var httpContext = new DefaultHttpContext(); + var version = DataVersion.FromToken("4611686018427387904"); + + httpContext.Response.WithQueryStatsAndPagingInfo( + new QueryStatsInfo(version, totalCount: 1), + new PagingInfo()); + + Assert.That(httpContext.Response.Headers.ETag.ToString(), Does.Contain(version.ToString())); + } + + class CompletingResponseFeature : IHttpResponseFeature + { + Func onCompleted; + object onCompletedState; + public int StatusCode { get; set; } = 200; + public string ReasonPhrase { get; set; } + public IHeaderDictionary Headers { get; set; } = new HeaderDictionary(); + public Stream Body { get; set; } = new MemoryStream(); + public bool HasStarted => false; +#pragma warning disable PS0013 // ASP.NET Core interface signatures cannot accept a cancellation token + public void OnStarting(Func callback, object state) { } + public void OnCompleted(Func callback, object state) => (onCompleted, onCompletedState) = (callback, state); +#pragma warning restore PS0013 + public Task Complete() => onCompleted(onCompletedState); } static ResultExecutingContext ResultExecuting(HttpContext httpContext) => @@ -94,4 +204,4 @@ static ResultExecutingContext ResultExecuting(HttpContext httpContext) => new OkObjectResult(new object()), controller: null); } -} +} \ No newline at end of file diff --git a/src/ServiceControl.Audit/Auditing/MessagesView/GetMessages2Controller.cs b/src/ServiceControl.Audit/Auditing/MessagesView/GetMessages2Controller.cs index 682e618eb9..b7a7fd9b48 100644 --- a/src/ServiceControl.Audit/Auditing/MessagesView/GetMessages2Controller.cs +++ b/src/ServiceControl.Audit/Auditing/MessagesView/GetMessages2Controller.cs @@ -54,7 +54,7 @@ public async Task> GetAllMessages( } Response.WithTotalCount(result.QueryStats.TotalCount); - Response.WithEtag(result.QueryStats.ETag); + Response.WithEtag(result.QueryStats.Version); return result.Results; } diff --git a/src/ServiceControl.Audit/Auditing/MessagesView/GetMessagesController.cs b/src/ServiceControl.Audit/Auditing/MessagesView/GetMessagesController.cs index 66031fcd9c..fa17d96f1d 100644 --- a/src/ServiceControl.Audit/Auditing/MessagesView/GetMessagesController.cs +++ b/src/ServiceControl.Audit/Auditing/MessagesView/GetMessagesController.cs @@ -68,7 +68,7 @@ public async Task Get(string id, CancellationToken cancellationTo throw new Exception($"Metadata for message '{id}' indicated that a body was present but no content could be found in storage"); } - Response.WithEtag(result.ETag); + Response.WithEtag(result.Version); var contentType = result.ContentType ?? "text/*"; return result.StringContent != null ? Content(result.StringContent, contentType) : File(result.StreamContent, contentType); } diff --git a/src/ServiceControl.Audit/Infrastructure/WebApi/HttpResponseExtensions.cs b/src/ServiceControl.Audit/Infrastructure/WebApi/HttpResponseExtensions.cs index 532273710a..92e2d7cc3b 100644 --- a/src/ServiceControl.Audit/Infrastructure/WebApi/HttpResponseExtensions.cs +++ b/src/ServiceControl.Audit/Infrastructure/WebApi/HttpResponseExtensions.cs @@ -9,40 +9,28 @@ namespace ServiceControl.Audit.Infrastructure.WebApi using Microsoft.AspNetCore.WebUtilities; using Microsoft.Extensions.Primitives; using Persistence.Infrastructure; + using ServiceControl.Infrastructure; static class HttpResponseExtensions { public static void WithTotalCount(this HttpResponse response, long totalCount) => response.WithHeader("Total-Count", totalCount.ToString(CultureInfo.InvariantCulture)); - public static void WithEtag(this HttpResponse response, StringValues value) + public static void WithEtag(this HttpResponse response, DataVersion version) { - var validator = value.ToString(); - - if (string.IsNullOrEmpty(validator)) + if (!version.HasValue) { return; } - // RFC 9110 requires an entity-tag to be a quoted string. Unquoted, EntityTagHeaderValue - // cannot parse it and NotModifiedStatusHttpHandler never matches a client's If-None-Match. - response.Headers.ETag = $"\"{validator}\""; + // Quotes are required by RFC 9110. Without them EntityTagHeaderValue cannot parse the tag and + // NotModifiedStatusHttpHandler never matches a client's If-None-Match. + response.Headers.ETag = $"W/\"{version}\""; } public static void WithQueryStatsInfo(this HttpResponse response, QueryStatsInfo queryStatsInfo) { response.WithTotalCount(queryStatsInfo.TotalCount); - response.WithEtag(queryStatsInfo.ETag); - } - - public static void WithDeterministicEtag(this HttpResponse response, string data) - { - if (string.IsNullOrEmpty(data)) - { - return; - } - - var guid = DeterministicGuid.MakeId(data); - response.WithEtag(guid.ToString()); + response.WithEtag(queryStatsInfo.Version); } static void WithHeader(this HttpResponse response, string header, StringValues value) => response.Headers.Append(header, value); @@ -104,7 +92,7 @@ static void AddLink(ICollection links, int page, string rel, string uriP public static void WithQueryStatsAndPagingInfo(this HttpResponse response, QueryStatsInfo queryStats, PagingInfo pagingInfo) { response.WithPagingLinksAndTotalCount(pagingInfo, queryStats.TotalCount, queryStats.HighestTotalCountOfAllTheInstances); - response.WithDeterministicEtag(queryStats.ETag); + response.WithEtag(queryStats.Version); } public static void WithPagingLinksAndTotalCount(this HttpResponse response, diff --git a/src/ServiceControl.Audit/Infrastructure/WebApi/NotModifiedStatusHttpHandler.cs b/src/ServiceControl.Audit/Infrastructure/WebApi/NotModifiedStatusHttpHandler.cs index 6c67704dac..d140a572e0 100644 --- a/src/ServiceControl.Audit/Infrastructure/WebApi/NotModifiedStatusHttpHandler.cs +++ b/src/ServiceControl.Audit/Infrastructure/WebApi/NotModifiedStatusHttpHandler.cs @@ -1,24 +1,43 @@ namespace ServiceControl.Audit.Infrastructure.WebApi { using System; + using System.Linq; using System.Net; using Microsoft.AspNetCore.Http; using Microsoft.AspNetCore.Http.Headers; using Microsoft.AspNetCore.Mvc; using Microsoft.AspNetCore.Mvc.Filters; + using Microsoft.AspNetCore.Mvc.Infrastructure; class NotModifiedStatusHttpHandler : IResultFilter { - static bool IfNoneMatch(RequestHeaders requestHeaders, ResponseHeaders responseHeaders) => - responseHeaders.ETag != null && requestHeaders.IfNoneMatch.Contains(responseHeaders.ETag); + static bool IfNoneMatch(RequestHeaders requestHeaders, ResponseHeaders responseHeaders) + { + var current = responseHeaders.ETag; + + if (current is null) + { + return false; + } + + // EntityTagHeaderValue.Equals compares strength as well as tag, and its own documentation + // says not to use it for this. RFC 9110 requires If-None-Match to use weak comparison. + return requestHeaders.IfNoneMatch.Any(candidate => + candidate.Tag.Equals("*", StringComparison.Ordinal) || candidate.Compare(current, useStrongComparison: false)); + } static bool IfNotModifiedSince(DateTimeOffset? ifModifiedSince, DateTimeOffset? lastModified) => lastModified <= ifModifiedSince; public void OnResultExecuting(ResultExecutingContext context) { - var statusCode = context.HttpContext.Response.StatusCode; - if (statusCode is < 200 or > 299) + if (context.HttpContext.Response.HasStarted) + { + // In forwarding scenarios we don't want to alter headers set by other instances + return; + } + + if (!IsSuccess(context)) { return; } @@ -31,10 +50,23 @@ public void OnResultExecuting(ResultExecutingContext context) if (ifNoneMatch || ifNotModifiedSince) { + // The replaced result never executes, so whatever it owned would never be released. + if (context.Result is FileStreamResult file) + { + context.HttpContext.Response.RegisterForDisposeAsync(file.FileStream); + } + context.Result = new StatusCodeResult((int)HttpStatusCode.NotModified); } } + // Response.StatusCode is still whatever the pipeline defaulted to, because the result that + // would set it has not executed yet. + static bool IsSuccess(ResultExecutingContext context) => + context.Result is IStatusCodeActionResult { StatusCode: { } statusCode } + ? statusCode is >= 200 and <= 299 + : context.HttpContext.Response.StatusCode is >= 200 and <= 299; + public void OnResultExecuted(ResultExecutedContext context) { // NOP diff --git a/src/ServiceControl.Persistence/Infrastructure/DataVersion.cs b/src/ServiceControl.Infrastructure/DataVersion.cs similarity index 96% rename from src/ServiceControl.Persistence/Infrastructure/DataVersion.cs rename to src/ServiceControl.Infrastructure/DataVersion.cs index d51bbd8d99..bfc83bda7e 100644 --- a/src/ServiceControl.Persistence/Infrastructure/DataVersion.cs +++ b/src/ServiceControl.Infrastructure/DataVersion.cs @@ -1,4 +1,5 @@ -namespace ServiceControl.Persistence.Infrastructure +#nullable enable +namespace ServiceControl.Infrastructure { using System; using System.Collections.Generic; @@ -67,10 +68,6 @@ public static DataVersion OverRows((string Name, object? Value)[]? summary return Compose([.. terms]); } - public static DataVersion OverRows((string Name, object? Value)[]? summary, IEnumerable rows) - where TRow : IVersionedRow => - OverRows(summary, rows, row => row.GetVersionFields()); - /// /// One version for a result gathered from several instances. Missing anywhere means missing overall. /// Keyed on the instance, so a validator moving from one instance to another still moves the @@ -187,4 +184,4 @@ static string Prefixed(string value) => _ => throw new ArgumentException($"A version term cannot be built from {value.GetType()}.", nameof(value)) }; } -} +} \ No newline at end of file diff --git a/src/ServiceControl.Persistence.EFCore/Implementation/BodyStorage/BodyStorage.cs b/src/ServiceControl.Persistence.EFCore/Implementation/BodyStorage/BodyStorage.cs index 02c93e381d..01b40b4b8a 100644 --- a/src/ServiceControl.Persistence.EFCore/Implementation/BodyStorage/BodyStorage.cs +++ b/src/ServiceControl.Persistence.EFCore/Implementation/BodyStorage/BodyStorage.cs @@ -1,5 +1,6 @@ namespace ServiceControl.Persistence.EFCore.Implementation.BodyStorage; +using ServiceControl.Infrastructure; using System.Linq.Expressions; using System.Text; using Microsoft.EntityFrameworkCore; diff --git a/src/ServiceControl.Persistence.EFCore/Infrastructure/QueryStatsInfoExtensions.cs b/src/ServiceControl.Persistence.EFCore/Infrastructure/QueryStatsInfoExtensions.cs index 3e0b8e30f8..4578a1a1c3 100644 --- a/src/ServiceControl.Persistence.EFCore/Infrastructure/QueryStatsInfoExtensions.cs +++ b/src/ServiceControl.Persistence.EFCore/Infrastructure/QueryStatsInfoExtensions.cs @@ -1,5 +1,6 @@ namespace ServiceControl.Persistence.EFCore.Infrastructure; +using ServiceControl.Infrastructure; using Microsoft.EntityFrameworkCore; using ServiceControl.Persistence.EFCore.Entities; using ServiceControl.Persistence.Infrastructure; @@ -26,6 +27,7 @@ public static async Task ToCountQueryStatsInfo(this IQu public static QueryStatsInfo ToQueryStatsInfo(this RetryHistory history) => new QueryStatsInfo(DataVersion.OverRows( [("historic", history.HistoricOperations.Count), ("unacknowledged", history.UnacknowledgedOperations.Count)], - history.HistoricOperations.Concat(history.UnacknowledgedOperations)), + history.HistoricOperations.Concat(history.UnacknowledgedOperations), + row => row.GetVersionFields()), history.HistoricOperations.Count); } diff --git a/src/ServiceControl.Persistence.RavenDB/RavenAttachmentsBodyStorage.cs b/src/ServiceControl.Persistence.RavenDB/RavenAttachmentsBodyStorage.cs index 98772b9496..11e1a956ce 100644 --- a/src/ServiceControl.Persistence.RavenDB/RavenAttachmentsBodyStorage.cs +++ b/src/ServiceControl.Persistence.RavenDB/RavenAttachmentsBodyStorage.cs @@ -1,5 +1,6 @@ namespace ServiceControl.Operations.BodyStorage.RavenAttachments { + using ServiceControl.Infrastructure; using System; using System.Linq; using System.Threading; diff --git a/src/ServiceControl.Persistence.RavenDB/RavenQueryStatisticsExtensions.cs b/src/ServiceControl.Persistence.RavenDB/RavenQueryStatisticsExtensions.cs index 7e58d3ecdb..37a85be551 100644 --- a/src/ServiceControl.Persistence.RavenDB/RavenQueryStatisticsExtensions.cs +++ b/src/ServiceControl.Persistence.RavenDB/RavenQueryStatisticsExtensions.cs @@ -1,5 +1,6 @@ namespace ServiceControl.Persistence { + using ServiceControl.Infrastructure; using Raven.Client.Documents.Session; using ServiceControl.Persistence.Infrastructure; diff --git a/src/ServiceControl.Persistence.RavenDB/Recoverability/RetryHistoryDataStore.cs b/src/ServiceControl.Persistence.RavenDB/Recoverability/RetryHistoryDataStore.cs index e8aae034e3..93f55c7ef1 100644 --- a/src/ServiceControl.Persistence.RavenDB/Recoverability/RetryHistoryDataStore.cs +++ b/src/ServiceControl.Persistence.RavenDB/Recoverability/RetryHistoryDataStore.cs @@ -1,5 +1,6 @@ namespace ServiceControl.Persistence.RavenDB.Recoverability { + using ServiceControl.Infrastructure; using System; using System.Threading; using System.Threading.Tasks; diff --git a/src/ServiceControl.Persistence.Tests/BodyStorage/BodyVersionTests.cs b/src/ServiceControl.Persistence.Tests/BodyStorage/BodyVersionTests.cs index 34ab459369..0c9312f596 100644 --- a/src/ServiceControl.Persistence.Tests/BodyStorage/BodyVersionTests.cs +++ b/src/ServiceControl.Persistence.Tests/BodyStorage/BodyVersionTests.cs @@ -5,6 +5,7 @@ namespace ServiceControl.Persistence.Tests; using System.Text; using System.Threading.Tasks; using NUnit.Framework; +using ServiceControl.Infrastructure; using ServiceControl.Operations.BodyStorage; using ServiceControl.Persistence.Infrastructure; diff --git a/src/ServiceControl.Persistence.Tests/EventLogDataStoreTests.cs b/src/ServiceControl.Persistence.Tests/EventLogDataStoreTests.cs index 02c9481eda..ea8847c4fa 100644 --- a/src/ServiceControl.Persistence.Tests/EventLogDataStoreTests.cs +++ b/src/ServiceControl.Persistence.Tests/EventLogDataStoreTests.cs @@ -5,6 +5,7 @@ namespace ServiceControl.Persistence.Tests; using System.Threading.Tasks; using NUnit.Framework; using ServiceControl.EventLog; +using ServiceControl.Infrastructure; using ServiceControl.Persistence.Infrastructure; class EventLogDataStoreTests : PersistenceTestBase diff --git a/src/ServiceControl.Persistence.Tests/Recoverability/ReturnToSenderDequeuerTests.cs b/src/ServiceControl.Persistence.Tests/Recoverability/ReturnToSenderDequeuerTests.cs index d6e0473429..819d134c1b 100644 --- a/src/ServiceControl.Persistence.Tests/Recoverability/ReturnToSenderDequeuerTests.cs +++ b/src/ServiceControl.Persistence.Tests/Recoverability/ReturnToSenderDequeuerTests.cs @@ -18,6 +18,7 @@ using Persistence; using Persistence.Infrastructure; using ServiceControl.CompositeViews.Messages; + using ServiceControl.Infrastructure; using ServiceControl.Operations; using ServiceControl.Operations.BodyStorage; using ServiceControl.Recoverability; diff --git a/src/ServiceControl.Persistence.Tests/VersionAssert.cs b/src/ServiceControl.Persistence.Tests/VersionAssert.cs index 25edc5bc5c..7db0333ee7 100644 --- a/src/ServiceControl.Persistence.Tests/VersionAssert.cs +++ b/src/ServiceControl.Persistence.Tests/VersionAssert.cs @@ -1,7 +1,7 @@ namespace ServiceControl.Persistence.Tests; using NUnit.Framework; -using ServiceControl.Persistence.Infrastructure; +using ServiceControl.Infrastructure; static class VersionAssert { diff --git a/src/ServiceControl.Persistence/IBodyStorage.cs b/src/ServiceControl.Persistence/IBodyStorage.cs index 187c3e8dc9..8f2b739737 100644 --- a/src/ServiceControl.Persistence/IBodyStorage.cs +++ b/src/ServiceControl.Persistence/IBodyStorage.cs @@ -1,5 +1,6 @@ namespace ServiceControl.Operations.BodyStorage { + using ServiceControl.Infrastructure; using System; using System.IO; using System.Threading; diff --git a/src/ServiceControl.Persistence/Infrastructure/QueryStatsInfo.cs b/src/ServiceControl.Persistence/Infrastructure/QueryStatsInfo.cs index 7aa2f0e1f7..ec6277d5ab 100644 --- a/src/ServiceControl.Persistence/Infrastructure/QueryStatsInfo.cs +++ b/src/ServiceControl.Persistence/Infrastructure/QueryStatsInfo.cs @@ -1,5 +1,7 @@ namespace ServiceControl.Persistence.Infrastructure { + using ServiceControl.Infrastructure; + public readonly struct QueryStatsInfo { public readonly DataVersion Version; diff --git a/src/ServiceControl.Persistence/Infrastructure/VersionedRowExtensions.cs b/src/ServiceControl.Persistence/Infrastructure/VersionedRowExtensions.cs index 3b093c4c96..55736c7338 100644 --- a/src/ServiceControl.Persistence/Infrastructure/VersionedRowExtensions.cs +++ b/src/ServiceControl.Persistence/Infrastructure/VersionedRowExtensions.cs @@ -1,6 +1,7 @@ namespace ServiceControl.Persistence.Infrastructure { using System.Collections.Generic; + using ServiceControl.Infrastructure; public static class VersionedRowExtensions { @@ -8,6 +9,6 @@ public static class VersionedRowExtensions /// Creates a QueryStatsInfo from a collection of rows, and versions it over those rows. /// public static QueryStatsInfo ToQueryStatsInfo(this IReadOnlyCollection rows, string name, long totalCount) - where TRow : IVersionedRow => new(DataVersion.OverRows([(name, totalCount)], rows), totalCount); + where TRow : IVersionedRow => new(DataVersion.OverRows([(name, totalCount)], rows, row => row.GetVersionFields()), totalCount); } } diff --git a/src/ServiceControl.UnitTests/BodyStorage/MessageBodyResultTests.cs b/src/ServiceControl.UnitTests/BodyStorage/MessageBodyResultTests.cs index 39adb49c2d..4937df3f7a 100644 --- a/src/ServiceControl.UnitTests/BodyStorage/MessageBodyResultTests.cs +++ b/src/ServiceControl.UnitTests/BodyStorage/MessageBodyResultTests.cs @@ -5,6 +5,7 @@ namespace ServiceControl.UnitTests.BodyStorage; using NUnit.Framework; using ServiceControl.Operations.BodyStorage; using ServiceControl.Persistence.Infrastructure; +using ServiceControl.Infrastructure; [TestFixture] public class MessageBodyResultTests diff --git a/src/ServiceControl.UnitTests/HeaderAssertions.cs b/src/ServiceControl.UnitTests/HeaderAssertions.cs index 1622ea5d79..2925edd445 100644 --- a/src/ServiceControl.UnitTests/HeaderAssertions.cs +++ b/src/ServiceControl.UnitTests/HeaderAssertions.cs @@ -1,5 +1,6 @@ namespace ServiceControl.UnitTests { + using ServiceControl.Infrastructure; using System.Collections.Generic; using NUnit.Framework; using ServiceControl.Persistence.Infrastructure; diff --git a/src/ServiceControl.UnitTests/Infrastructure/DataVersionTests.cs b/src/ServiceControl.UnitTests/Infrastructure/DataVersionTests.cs index 1498b200c8..6eb24ea88a 100644 --- a/src/ServiceControl.UnitTests/Infrastructure/DataVersionTests.cs +++ b/src/ServiceControl.UnitTests/Infrastructure/DataVersionTests.cs @@ -1,5 +1,7 @@ namespace ServiceControl.UnitTests.Infrastructure; +using ServiceControl.Infrastructure; + using System; using NUnit.Framework; using ServiceControl.Persistence.Infrastructure; diff --git a/src/ServiceControl.UnitTests/Infrastructure/VersionedRowTests.cs b/src/ServiceControl.UnitTests/Infrastructure/VersionedRowTests.cs index 7262271fe9..a4a090cf00 100644 --- a/src/ServiceControl.UnitTests/Infrastructure/VersionedRowTests.cs +++ b/src/ServiceControl.UnitTests/Infrastructure/VersionedRowTests.cs @@ -1,6 +1,7 @@ #nullable enable namespace ServiceControl.UnitTests.Infrastructure { + using ServiceControl.Infrastructure; using System; using System.Collections.Generic; using System.Linq; @@ -75,7 +76,7 @@ public void Changing_any_rendered_field_moves_the_version(Type type) } } - static DataVersion VersionOf(IVersionedRow row) => DataVersion.OverRows([("rows", 1)], [row]); + static DataVersion VersionOf(IVersionedRow row) => DataVersion.OverRows([("rows", 1)], [row], r => r.GetVersionFields()); static PropertyInfo[] Rendered(Type type) => [.. type.GetProperties(BindingFlags.Public | BindingFlags.Instance).Where(property => property.CanWrite)]; diff --git a/src/ServiceControl.UnitTests/Infrastructure/WebApi/ConditionalGetTests.cs b/src/ServiceControl.UnitTests/Infrastructure/WebApi/ConditionalGetTests.cs index 3c804c7f70..c3e1084380 100644 --- a/src/ServiceControl.UnitTests/Infrastructure/WebApi/ConditionalGetTests.cs +++ b/src/ServiceControl.UnitTests/Infrastructure/WebApi/ConditionalGetTests.cs @@ -1,5 +1,7 @@ namespace ServiceControl.UnitTests.Infrastructure.WebApi; +using ServiceControl.Infrastructure; + using System.Collections.Generic; using System.Net; using Microsoft.AspNetCore.Http; diff --git a/src/ServiceControl.UnitTests/Recoverability/MessageRedirectVersionTests.cs b/src/ServiceControl.UnitTests/Recoverability/MessageRedirectVersionTests.cs index 09feafe2a9..b6b5ae45d1 100644 --- a/src/ServiceControl.UnitTests/Recoverability/MessageRedirectVersionTests.cs +++ b/src/ServiceControl.UnitTests/Recoverability/MessageRedirectVersionTests.cs @@ -1,5 +1,6 @@ namespace ServiceControl.UnitTests.Operations { + using ServiceControl.Infrastructure; using System; using System.Collections.Generic; using NUnit.Framework; diff --git a/src/ServiceControl.UnitTests/Recoverability/RetryGroupVersionTests.cs b/src/ServiceControl.UnitTests/Recoverability/RetryGroupVersionTests.cs index 68acf0a198..950931bf82 100644 --- a/src/ServiceControl.UnitTests/Recoverability/RetryGroupVersionTests.cs +++ b/src/ServiceControl.UnitTests/Recoverability/RetryGroupVersionTests.cs @@ -1,5 +1,6 @@ namespace ServiceControl.UnitTests.Operations { + using ServiceControl.Infrastructure; using System; using System.Collections.Generic; using NUnit.Framework; diff --git a/src/ServiceControl.UnitTests/ScatterGather/IncompleteResultsTests.cs b/src/ServiceControl.UnitTests/ScatterGather/IncompleteResultsTests.cs index 85d76224cb..744fa7bf19 100644 --- a/src/ServiceControl.UnitTests/ScatterGather/IncompleteResultsTests.cs +++ b/src/ServiceControl.UnitTests/ScatterGather/IncompleteResultsTests.cs @@ -17,6 +17,7 @@ namespace ServiceControl.UnitTests.ScatterGather; using ServiceControl.Infrastructure.Api; using ServiceControl.Infrastructure.WebApi; using ServiceControl.Persistence.Infrastructure; +using ServiceControl.Infrastructure; /// /// A timed-out or failed instance is not an instance with no data: the composite keeps what the others diff --git a/src/ServiceControl.UnitTests/ScatterGather/MessageView_ScatterGatherTest.cs b/src/ServiceControl.UnitTests/ScatterGather/MessageView_ScatterGatherTest.cs index 8fc1cf9225..3a8e127173 100644 --- a/src/ServiceControl.UnitTests/ScatterGather/MessageView_ScatterGatherTest.cs +++ b/src/ServiceControl.UnitTests/ScatterGather/MessageView_ScatterGatherTest.cs @@ -1,5 +1,6 @@ namespace ServiceControl.UnitTests.ScatterGather { + using ServiceControl.Infrastructure; using System; using System.Collections.Generic; using System.Linq; diff --git a/src/ServiceControl.UnitTests/ScatterGather/MessagesView_ScatterGather_DataFromBothInstances.cs b/src/ServiceControl.UnitTests/ScatterGather/MessagesView_ScatterGather_DataFromBothInstances.cs index 69ac74e3b5..6dd84fabc8 100644 --- a/src/ServiceControl.UnitTests/ScatterGather/MessagesView_ScatterGather_DataFromBothInstances.cs +++ b/src/ServiceControl.UnitTests/ScatterGather/MessagesView_ScatterGather_DataFromBothInstances.cs @@ -1,5 +1,6 @@ namespace ServiceControl.UnitTests.ScatterGather { + using ServiceControl.Infrastructure; using System; using System.Collections.Generic; using System.Linq; diff --git a/src/ServiceControl.UnitTests/ScatterGather/ScatterGatherVersionTests.cs b/src/ServiceControl.UnitTests/ScatterGather/ScatterGatherVersionTests.cs index 22eebc1e3f..9f239c581d 100644 --- a/src/ServiceControl.UnitTests/ScatterGather/ScatterGatherVersionTests.cs +++ b/src/ServiceControl.UnitTests/ScatterGather/ScatterGatherVersionTests.cs @@ -1,5 +1,6 @@ namespace ServiceControl.UnitTests.ScatterGather { + using ServiceControl.Infrastructure; using System.Collections.Generic; using System.Linq; using System.Threading; diff --git a/src/ServiceControl/CompositeViews/Messages/ScatterGatherApi.cs b/src/ServiceControl/CompositeViews/Messages/ScatterGatherApi.cs index 1af5732b9f..da91191eb2 100644 --- a/src/ServiceControl/CompositeViews/Messages/ScatterGatherApi.cs +++ b/src/ServiceControl/CompositeViews/Messages/ScatterGatherApi.cs @@ -1,5 +1,6 @@ namespace ServiceControl.CompositeViews.Messages { + using ServiceControl.Infrastructure; using System; using System.Collections.Generic; using System.Linq; diff --git a/src/ServiceControl/Infrastructure/WebApi/HttpResponseExtensions.cs b/src/ServiceControl/Infrastructure/WebApi/HttpResponseExtensions.cs index cf44cf3f0d..988359f144 100644 --- a/src/ServiceControl/Infrastructure/WebApi/HttpResponseExtensions.cs +++ b/src/ServiceControl/Infrastructure/WebApi/HttpResponseExtensions.cs @@ -1,5 +1,6 @@ namespace ServiceControl.Infrastructure.WebApi { + using ServiceControl.Infrastructure; using System; using System.Collections.Generic; using System.Globalization; diff --git a/src/ServiceControl/MessageRedirects/Api/MessageRedirectsController.cs b/src/ServiceControl/MessageRedirects/Api/MessageRedirectsController.cs index fc52d6fbc0..50b8b3e10c 100644 --- a/src/ServiceControl/MessageRedirects/Api/MessageRedirectsController.cs +++ b/src/ServiceControl/MessageRedirects/Api/MessageRedirectsController.cs @@ -1,5 +1,6 @@ namespace ServiceControl.MessageRedirects.Api { + using ServiceControl.Infrastructure; using System; using System.Collections.Generic; using System.Linq; @@ -174,7 +175,7 @@ public async Task CountRedirects(CancellationToken cancellationToken = default) { var redirects = await store.GetRedirects(cancellationToken); - Response.WithEtag(DataVersion.OverRows([("redirects", redirects.Count)], redirects)); + Response.WithEtag(DataVersion.OverRows([("redirects", redirects.Count)], redirects, r => ((IVersionedRow)r).GetVersionFields())); Response.WithTotalCount(redirects.Count); } @@ -190,7 +191,7 @@ public async Task> Redirects(string sort, stri .Paging(pagingInfo) .ToList(); - Response.WithQueryStatsAndPagingInfo(new QueryStatsInfo(DataVersion.OverRows([("redirects", redirects.Count)], page), redirects.Count), pagingInfo); + Response.WithQueryStatsAndPagingInfo(new QueryStatsInfo(DataVersion.OverRows([("redirects", redirects.Count)], page, r => ((IVersionedRow)r).GetVersionFields()), redirects.Count), pagingInfo); return page.Select(r => new RedirectsQueryResult ( diff --git a/src/ServiceControl/Monitoring/Web/EndpointsMonitoringController.cs b/src/ServiceControl/Monitoring/Web/EndpointsMonitoringController.cs index 8a4a72c49f..362fdb0bc6 100644 --- a/src/ServiceControl/Monitoring/Web/EndpointsMonitoringController.cs +++ b/src/ServiceControl/Monitoring/Web/EndpointsMonitoringController.cs @@ -1,5 +1,6 @@ namespace ServiceControl.Monitoring { + using ServiceControl.Infrastructure; using System; using System.Collections.Generic; using System.Threading; diff --git a/src/ServiceControl/Recoverability/API/FailureGroupsController.cs b/src/ServiceControl/Recoverability/API/FailureGroupsController.cs index bf105537a1..bfdc7918a8 100644 --- a/src/ServiceControl/Recoverability/API/FailureGroupsController.cs +++ b/src/ServiceControl/Recoverability/API/FailureGroupsController.cs @@ -1,5 +1,6 @@ namespace ServiceControl.Recoverability.API { + using ServiceControl.Infrastructure; using System.Collections.Generic; using System.Linq; using System.Threading;