Skip to content
Merged
6 changes: 5 additions & 1 deletion src/NATS.Client.Core/Internal/Telemetry.cs
Original file line number Diff line number Diff line change
Expand Up @@ -87,7 +87,11 @@ internal static void AddTraceContextHeaders(Activity? activity, ref NatsHeaders?
return;
}

headers[fieldName] = fieldValue;
// There are cases where headers reused internally (e.g. JetStream publish retry)
// there may also be cases where application can reuse the same header
// in which case we should still be able to overwrite headers with telemetry fields
// even though headers would be set to readonly before being passed down in publish methods.
headers.SetOverrideReadOnly(fieldName, fieldValue);
});
}

Expand Down
1 change: 1 addition & 0 deletions src/NATS.Client.Core/NatsConnection.Publish.cs
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ public ValueTask PublishAsync(string subject, NatsHeaders? headers = default, st
if (Telemetry.HasListeners())
{
using var activity = Telemetry.StartSendActivity($"{SpanDestinationName(subject)} {Telemetry.Constants.PublishActivityName}", this, subject, replyTo);
Telemetry.AddTraceContextHeaders(activity, ref headers);

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

this was a bug fix

try
{
headers?.SetReadOnly();
Expand Down
18 changes: 18 additions & 0 deletions src/NATS.Client.Core/NatsHeaders.cs
Original file line number Diff line number Diff line change
Expand Up @@ -423,6 +423,24 @@ IEnumerator IEnumerable.GetEnumerator()

internal void SetReadOnly() => Interlocked.Exchange(ref _readonly, 1);

internal void SetOverrideReadOnly(string key, StringValues value)
{
if (key == null)
{
throw new ArgumentNullException(nameof(key));
}

if (value.Count == 0)
{
Store?.Remove(key);
}
else
{
EnsureStore(1);
Store[key] = value;
}
}

private void ThrowIfReadOnly()
{
if (IsReadOnly)
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
authorization: {
users: [
{ user:u, permissions: { publish:x } }
]
}
16 changes: 14 additions & 2 deletions tests/NATS.Client.Core2.Tests/NatsServerFixture.cs
Original file line number Diff line number Diff line change
Expand Up @@ -8,8 +8,13 @@ public class NatsServerFixture : IDisposable
private int _next;

public NatsServerFixture()
: this(null)
{
Server = NatsServerProcess.Start();
}

protected NatsServerFixture(string? config)
{
Server = NatsServerProcess.Start(config: config);
}

public NatsServerProcess Server { get; }
Expand All @@ -27,6 +32,13 @@ public void Dispose()
}

[CollectionDefinition("nats-server")]
public class DatabaseCollection : ICollectionFixture<NatsServerFixture>
public class NatsServerCollection : ICollectionFixture<NatsServerFixture>
{
}

public class NatsServerRestrictedUserFixture() : NatsServerFixture("resources/configs/restricted-user.conf");

[CollectionDefinition("nats-server-restricted-user")]
public class NatsServerRestrictedUserCollection : ICollectionFixture<NatsServerRestrictedUserFixture>
{
}
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,13 @@
<ProjectReference Include="..\NATS.Client.TestUtilities\NATS.Client.TestUtilities.csproj" />
</ItemGroup>

<ItemGroup>
<Content Include="..\NATS.Client.Core.Tests\resources\**\*">
<Link>resources\%(RecursiveDir)%(Filename)%(Extension)</Link>
<CopyToOutputDirectory>PreserveNewest</CopyToOutputDirectory>
</Content>
</ItemGroup>

<ItemGroup>
<Compile Include="..\NATS.Client.Core2.Tests\NatsServerFixture.cs" />
</ItemGroup>
Expand Down
152 changes: 152 additions & 0 deletions tests/NATS.Client.JetStream.Tests/PublishRetryTest.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,152 @@
using System.Diagnostics;
using NATS.Client.Core2.Tests;
using NATS.Net;

namespace NATS.Client.JetStream.Tests;

[Collection("nats-server-restricted-user")]
public class PublishRetryTest
{
private readonly ITestOutputHelper _output;
private readonly NatsServerRestrictedUserFixture _server;

public PublishRetryTest(ITestOutputHelper output, NatsServerRestrictedUserFixture server)
{
_output = output;
_server = server;
}

[Fact]
public async Task Publish_with_or_without_telemetry()
{
var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10));
await using var nats = new NatsConnection(new NatsOpts
{
Url = _server.Url,
AuthOpts = new NatsAuthOpts { Username = "u" },
});
var prefix = _server.GetNextId();

// With telemetry
{
using var activityListener = new ActivityListener { ShouldListenTo = _ => true, Sample = (ref ActivityCreationOptions<ActivityContext> _) => ActivitySamplingResult.AllData, };
ActivitySource.AddActivityListener(activityListener);
var headers = new NatsHeaders();
await nats.PublishAsync(prefix, headers: headers, cancellationToken: cts.Token);
Assert.Contains("traceparent", headers);
}

// Without telemetry
{
var headers = new NatsHeaders();
await nats.PublishAsync(prefix, headers: headers, cancellationToken: cts.Token);
Assert.DoesNotContain("traceparent", headers);
Assert.Empty(headers);
}

// With telemetry and data
{
using var activityListener = new ActivityListener { ShouldListenTo = _ => true, Sample = (ref ActivityCreationOptions<ActivityContext> _) => ActivitySamplingResult.AllData, };
ActivitySource.AddActivityListener(activityListener);
var headers = new NatsHeaders();
await nats.PublishAsync(prefix, data: 1, headers: headers, cancellationToken: cts.Token);
Assert.Contains("traceparent", headers);
}

// Without telemetry and data
{
var headers = new NatsHeaders();
await nats.PublishAsync(prefix, data: 1, headers: headers, cancellationToken: cts.Token);
Assert.DoesNotContain("traceparent", headers);
Assert.Empty(headers);
}
}

[Fact]
public async Task Multiple_publish_with_same_headers_when_telemetry_on_should_not_throw_header_readonly_exception()
{
using var activityListener = new ActivityListener
{
ShouldListenTo = _ => true,
Sample = (ref ActivityCreationOptions<ActivityContext> _) => ActivitySamplingResult.AllData,
};
ActivitySource.AddActivityListener(activityListener);

var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10));
await using var nats = new NatsConnection(new NatsOpts
{
Url = _server.Url,
AuthOpts = new NatsAuthOpts { Username = "u" },
});
var prefix = _server.GetNextId();

// Publish with no data implemented in a different method
{
var headers = new NatsHeaders();
await nats.PublishAsync(prefix, headers: headers, cancellationToken: cts.Token);
await nats.PublishAsync(prefix, headers: headers, cancellationToken: cts.Token);
Assert.Contains("traceparent", headers);
}

// Also try publish-with-data method
{
var headers = new NatsHeaders();
await nats.PublishAsync(prefix, data: 1, headers: headers, cancellationToken: cts.Token);
await nats.PublishAsync(prefix, data: 1, headers: headers, cancellationToken: cts.Token);
Assert.Contains("traceparent", headers);
}
}

[Fact]
public async Task Retry_with_telemetry_on_should_not_throw_header_readonly_exception()
{
using var activityListener = new ActivityListener
{
ShouldListenTo = _ => true,
Sample = (ref ActivityCreationOptions<ActivityContext> _) => ActivitySamplingResult.AllData,
};
ActivitySource.AddActivityListener(activityListener);

var cts = new CancellationTokenSource(TimeSpan.FromSeconds(30));
await using var nats = new NatsConnection(new NatsOpts
{
Url = _server.Url,
AuthOpts = new NatsAuthOpts { Username = "u" },
RequestTimeout = TimeSpan.FromSeconds(.5),
});
var prefix = _server.GetNextId();
var js = nats.CreateJetStreamContext();

// Because of the permission error (associated user should only have permission to subject 'x' to publish)
// JetStream publish internally retry publishing again which in turn would use the same header.
// Telemetry altering headers should not cause 'readonly' exceptions.
{
var headers = new NatsHeaders();
var exception = await Assert.ThrowsAnyAsync<Exception>(async () =>
{
await js.PublishAsync(
subject: $"{prefix}.foo",
data: "order 1",
headers: headers,
cancellationToken: cts.Token);
});
Assert.IsNotType<InvalidOperationException>(exception);
Assert.DoesNotMatch("response headers cannot be modified", exception.Message);
Assert.Contains("traceparent", headers);
}

// Also check for specific exception
{
var headers = new NatsHeaders();
await Assert.ThrowsAsync<NatsJSPublishNoResponseException>(async () =>
{
await js.PublishAsync(
subject: $"{prefix}.foo",
data: "order 1",
headers: headers,
cancellationToken: cts.Token);
});
Assert.Contains("traceparent", headers);
}
}
}