diff --git a/src/NATS.Client.JetStream/Internal/NatsJSConsume.cs b/src/NATS.Client.JetStream/Internal/NatsJSConsume.cs index 310966c5c..273290325 100644 --- a/src/NATS.Client.JetStream/Internal/NatsJSConsume.cs +++ b/src/NATS.Client.JetStream/Internal/NatsJSConsume.cs @@ -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( diff --git a/src/NATS.Client.JetStream/Internal/NatsJSExtensionsInternal.cs b/src/NATS.Client.JetStream/Internal/NatsJSExtensionsInternal.cs index 582a21767..0d137113b 100644 --- a/src/NATS.Client.JetStream/Internal/NatsJSExtensionsInternal.cs +++ b/src/NATS.Client.JetStream/Internal/NatsJSExtensionsInternal.cs @@ -26,9 +26,8 @@ public static void HandlePinIdMismatch(NatsJSConsumer? jsConsumer, NatsJSNotific /// The consumer to set the pin ID on. 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); diff --git a/tests/NATS.Client.JetStream.Tests/PinnedClientTest.cs b/tests/NATS.Client.JetStream.Tests/PinnedClientTest.cs index 1e7863e17..8407e5c71 100644 --- a/tests/NATS.Client.JetStream.Tests/PinnedClientTest.cs +++ b/tests/NATS.Client.JetStream.Tests/PinnedClientTest.cs @@ -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; @@ -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() {