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