Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
namespace ServiceControl.AcceptanceTests.RavenDB.Recoverability.MessageFailures
{
using ServiceControl.Infrastructure;
using System.Threading;
using System.Threading.Tasks;
using Infrastructure.WebApi;
Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
namespace ServiceControl.AcceptanceTests.RavenDB.Recoverability.MessageFailures
{
using ServiceControl.Infrastructure;
using System.Threading;
using System.Threading.Tasks;
using Infrastructure.WebApi;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,4 +23,9 @@
<PackageReference Include="NUnit3TestAdapter" />
</ItemGroup>

<ItemGroup>
<!-- InMemory does not version list queries, so its list endpoints carry no ETag -->
<Compile Remove="WebApi\When_a_list_endpoint_is_revalidated.cs" />
</ItemGroup>

</Project>
Original file line number Diff line number Diff line change
@@ -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<MyContext>()
.WithEndpoint<Receiver>(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<MessagesView>("/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<MyContext>()
.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<MyContext>()
.WithEndpoint<Receiver>(b => b.When(bus => bus.SendLocal(new MyMessage { Payload = "PAYLOAD" })))
.Done(async c =>
{
if (c.MessageId == null)
{
return false;
}

MessagesView audited = await this.TryGetSingle<MessagesView>("/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<Answer> 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<HttpResponseMessage> 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<DefaultServerWithAudit>();

[Handler]
public class MyMessageHandler(MyContext testContext, IReadOnlySettings settings) : IHandleMessages<MyMessage>
{
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; }
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -58,9 +58,61 @@ await Define<MyContext>()
.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<MyContext>()
.WithEndpoint<Receiver>(b => b.When(bus => bus.SendLocal(new MyMessage { Payload = "PAYLOAD" })))
.Done(async c =>
{
if (c.MessageId == null)
{
return false;
}

MessagesView audited = await this.TryGetSingle<MessagesView>("/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<DefaultServerWithAudit>();
Expand Down Expand Up @@ -89,4 +141,4 @@ public class MyContext : ScenarioContext
public string EndpointNameOfReceivingEndpoint { get; set; }
}
}
}
}
Original file line number Diff line number Diff line change
@@ -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
{
Expand Down Expand Up @@ -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)
{
Expand All @@ -44,25 +48,17 @@ public Task Store(string bodyId, string contentType, int bodySize, Stream bodySt
return Task.CompletedTask;
}

public async Task<StreamResult> TryFetch(string bodyId, CancellationToken cancellationToken = default)
public Task<MessageBodyView> 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
Expand All @@ -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; }
}
}
}
Loading
Loading