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
11 changes: 10 additions & 1 deletion src/NATS.Client.JetStream/Internal/NatsJSConsume.cs
Original file line number Diff line number Diff line change
Expand Up @@ -173,7 +173,16 @@ public ValueTask CallMsgNextAsync(string origin, ConsumerGetnextRequest request,

if (_debug)
{
_logger.LogDebug(NatsJSLogEvents.PullRequest, "Sending pull request for {Origin} {Msgs}, {Bytes}", origin, request.Batch, request.MaxBytes);
_logger.LogDebug(
NatsJSLogEvents.PullRequest,
"Sending pull request for {Origin} {Msgs}, {Bytes}, pinId={PinId}, expires={Expires}, idleHeartbeat={IdleHeartbeat}, group={Group}",
origin,
request.Batch,
request.MaxBytes,
request.Id,
request.Expires,
request.IdleHeartbeat,
request.Group);
}

return Connection.PublishAsync(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,9 +26,8 @@ public static void HandlePinIdMismatch(NatsJSConsumer? jsConsumer, NatsJSNotific
/// <param name="jsConsumer">The consumer to set the pin ID on.</param>
public static void TrySetPinIdFromHeaders(NatsHeaders? headers, NatsJSConsumer? jsConsumer)
{
if (jsConsumer != null && headers != null && headers.TryGetValue(NatsPinIdHeader, out var pinIdValues))
if (jsConsumer != null && headers != null && headers.TryGetLastValue(NatsPinIdHeader, out var pinId))
{
var pinId = pinIdValues.ToString();
if (!string.IsNullOrEmpty(pinId))
{
jsConsumer.SetPinId(pinId);
Expand Down
30 changes: 30 additions & 0 deletions tests/NATS.Client.JetStream.Tests/PinnedClientTest.cs
Original file line number Diff line number Diff line change
@@ -1,6 +1,9 @@
using System.Collections.Concurrent;
using Microsoft.Extensions.Primitives;
using NATS.Client.Core;
using NATS.Client.Core.Tests;
using NATS.Client.Core2.Tests;
using NATS.Client.JetStream.Internal;
using NATS.Client.JetStream.Models;
using NATS.Client.TestUtilities;
using NATS.Client.TestUtilities2;
Expand Down Expand Up @@ -520,6 +523,33 @@ public async Task Consumer_info_shows_priority_groups_state()

public class PinnedClientMockServerTest
{
[Fact]
public async Task Pin_id_from_headers_should_use_last_value_when_multiple_headers_present()
{
await using var ms = new MockServer((_, cmd) =>
{
if (cmd.Name == "PUB" && cmd.Subject.Contains("CONSUMER.INFO", StringComparison.Ordinal))
{
cmd.Reply(payload: """{"stream_name":"x","name":"x"}""");
}

return Task.CompletedTask;
});

using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(30));
await using var nats = new NatsConnection(new NatsOpts { Url = ms.Url });
var js = nats.CreateJetStreamContext();
var consumer = (NatsJSConsumer)await js.GetConsumerAsync("x", "x", cts.Token);
var headers = new NatsHeaders
{
{ "Nats-Pin-Id", new StringValues(["pin-stale", "pin-current"]) },
};

NatsJSExtensionsInternal.TrySetPinIdFromHeaders(headers, consumer);

Assert.Equal("pin-current", consumer.GetPinId());
}

[Fact]
public async Task Queued_consume_pull_request_should_use_latest_pin_id_when_sent()
{
Expand Down
Loading