From fa53133e9286588bb202be92043890be197cf3ab Mon Sep 17 00:00:00 2001 From: Ziya Suzen Date: Wed, 4 Mar 2026 13:25:57 +0000 Subject: [PATCH 1/7] Fix NuidTests flap --- tests/NATS.Client.CoreUnit.Tests/NuidTests.cs | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/tests/NATS.Client.CoreUnit.Tests/NuidTests.cs b/tests/NATS.Client.CoreUnit.Tests/NuidTests.cs index 7e1b1005b..8ffeaf5f8 100644 --- a/tests/NATS.Client.CoreUnit.Tests/NuidTests.cs +++ b/tests/NATS.Client.CoreUnit.Tests/NuidTests.cs @@ -137,7 +137,7 @@ public void GetNextNuid_PrefixRenewed_Char() }); executionThread.Start(); - executionThread.Join(1_000); + executionThread.Join(10_000); // Assert Assert.Equal(2, Interlocked.CompareExchange(ref _result, 0, 0)); @@ -175,7 +175,7 @@ public void InitAndWrite_Char() Volatile.Write(ref completedSuccessfully, didWrite && isMatch); }); t.Start(); - t.Join(1_000); + t.Join(10_000); Assert.True(completedSuccessfully); } @@ -202,7 +202,7 @@ public void DifferentThreads_DifferentPrefixes() threads.Add(t); } - threads.ForEach(t => t.Join(1_000)); + threads.ForEach(t => t.Join(10_000)); // Assert var uniquePrefixes = new HashSet(); From 5c74f8a53393cececa6e0ddaeab427de0a8481ef Mon Sep 17 00:00:00 2001 From: Ziya Suzen Date: Wed, 4 Mar 2026 13:35:41 +0000 Subject: [PATCH 2/7] Improve NuidTests stability --- tests/NATS.Client.CoreUnit.Tests/NuidTests.cs | 9 ++++++--- 1 file changed, 6 insertions(+), 3 deletions(-) diff --git a/tests/NATS.Client.CoreUnit.Tests/NuidTests.cs b/tests/NATS.Client.CoreUnit.Tests/NuidTests.cs index 8ffeaf5f8..f09fa29aa 100644 --- a/tests/NATS.Client.CoreUnit.Tests/NuidTests.cs +++ b/tests/NATS.Client.CoreUnit.Tests/NuidTests.cs @@ -137,7 +137,7 @@ public void GetNextNuid_PrefixRenewed_Char() }); executionThread.Start(); - executionThread.Join(10_000); + Assert.True(executionThread.Join(10_000), $"Thread {executionThread.ManagedThreadId} did not complete within the timeout."); // Assert Assert.Equal(2, Interlocked.CompareExchange(ref _result, 0, 0)); @@ -175,7 +175,7 @@ public void InitAndWrite_Char() Volatile.Write(ref completedSuccessfully, didWrite && isMatch); }); t.Start(); - t.Join(10_000); + Assert.True(t.Join(10_000), $"Thread {t.ManagedThreadId} did not complete within the timeout."); Assert.True(completedSuccessfully); } @@ -202,7 +202,10 @@ public void DifferentThreads_DifferentPrefixes() threads.Add(t); } - threads.ForEach(t => t.Join(10_000)); + foreach (var thread in threads) + { + Assert.True(thread.Join(10_000), $"Thread {thread.ManagedThreadId} did not complete within the timeout."); + } // Assert var uniquePrefixes = new HashSet(); From 8c5291194cd95a6137c08c958ea007602967b6c5 Mon Sep 17 00:00:00 2001 From: Ziya Suzen Date: Thu, 9 Apr 2026 12:19:58 +0100 Subject: [PATCH 3/7] tests: fix flaky test failures in CI Add MockServer ready signal and use ConnectRetryAsync for mock server tests to avoid race between acceptance loop startup and client connect. Add explicit ConnectRetryAsync to Request_reply_many_test_max_count. Increase Watcher_timeout_reconnect timeout from 30s to 60s to allow margin for idle timeout detection on slow CI. Skip first ping from max RTT assertion in SlowConsumer test since it catches the publish burst tail. Increase CommandTimeout in Buffer_high_pressure_pub_test from 5s to 30s for the 195MB high-pressure workload. --- tests/NATS.Client.Core2.Tests/BufferReferenceTest.cs | 1 + tests/NATS.Client.Core2.Tests/RequestReplyTest.cs | 2 ++ tests/NATS.Client.Core2.Tests/Utf8SubjectTest.cs | 10 +++++++--- tests/NATS.Client.JetStream.Tests/SlowConsumerTest.cs | 5 ++++- .../NatsKVWatcherTest.cs | 2 +- tests/NATS.Client.TestUtilities/MockServer.cs | 4 ++++ 6 files changed, 19 insertions(+), 5 deletions(-) diff --git a/tests/NATS.Client.Core2.Tests/BufferReferenceTest.cs b/tests/NATS.Client.Core2.Tests/BufferReferenceTest.cs index a34a26d7f..827c186c3 100644 --- a/tests/NATS.Client.Core2.Tests/BufferReferenceTest.cs +++ b/tests/NATS.Client.Core2.Tests/BufferReferenceTest.cs @@ -22,6 +22,7 @@ public async Task Buffer_high_pressure_pub_test() { Url = _server.Url, RequestReplyMode = NatsRequestReplyMode.Direct, + CommandTimeout = TimeSpan.FromSeconds(30), }); await nats.ConnectRetryAsync(); diff --git a/tests/NATS.Client.Core2.Tests/RequestReplyTest.cs b/tests/NATS.Client.Core2.Tests/RequestReplyTest.cs index c4d9f269b..848fef370 100644 --- a/tests/NATS.Client.Core2.Tests/RequestReplyTest.cs +++ b/tests/NATS.Client.Core2.Tests/RequestReplyTest.cs @@ -4,6 +4,7 @@ using System.Text.Json.Nodes; using NATS.Client.Core2.Tests; using NATS.Client.Core2.Tests.ExtraUtils.FrameworkPolyfillExtensions; +using NATS.Client.TestUtilities2; namespace NATS.Client.Core.Tests; @@ -233,6 +234,7 @@ public async Task Request_reply_many_test_start_up_timeout() public async Task Request_reply_many_test_max_count() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var sub = await nats.SubscribeCoreAsync("foo"); var reg = sub.Register(async msg => diff --git a/tests/NATS.Client.Core2.Tests/Utf8SubjectTest.cs b/tests/NATS.Client.Core2.Tests/Utf8SubjectTest.cs index 00704fee7..73211e873 100644 --- a/tests/NATS.Client.Core2.Tests/Utf8SubjectTest.cs +++ b/tests/NATS.Client.Core2.Tests/Utf8SubjectTest.cs @@ -1,6 +1,7 @@ using System.Text; using NATS.Client.Core2.Tests; using NATS.Client.TestUtilities; +using NATS.Client.TestUtilities2; namespace NATS.Client.Core.Tests; @@ -37,8 +38,9 @@ public async Task Utf8_subject_is_decoded_correctly() }, cancellationToken: cts.Token); + await server.Ready; await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); await foreach (var msg in nats.SubscribeAsync(">", cancellationToken: cts.Token)) { @@ -77,8 +79,9 @@ public async Task Emoji_subject_and_reply_to_are_decoded_correctly() }, cancellationToken: cts.Token); + await server.Ready; await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); await foreach (var msg in nats.SubscribeAsync(">", cancellationToken: cts.Token)) { @@ -123,8 +126,9 @@ public async Task Utf8_subject_with_hmsg_is_decoded_correctly_mock() }, cancellationToken: cts.Token); + await server.Ready; await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); await foreach (var msg in nats.SubscribeAsync(">", cancellationToken: cts.Token)) { diff --git a/tests/NATS.Client.JetStream.Tests/SlowConsumerTest.cs b/tests/NATS.Client.JetStream.Tests/SlowConsumerTest.cs index 14a9f391a..f8707eca0 100644 --- a/tests/NATS.Client.JetStream.Tests/SlowConsumerTest.cs +++ b/tests/NATS.Client.JetStream.Tests/SlowConsumerTest.cs @@ -276,7 +276,10 @@ public async Task JetStream_fetch_slow_consumer_should_not_block_connection() { var rtt = await pingTask; pingCount++; - if (rtt.TotalMilliseconds > maxPingRttMs) + + // Skip the first ping for max RTT calculation since it + // legitimately catches the tail of the 100-message publish burst. + if (i > 0 && rtt.TotalMilliseconds > maxPingRttMs) maxPingRttMs = rtt.TotalMilliseconds; _output.WriteLine($"Ping {i + 1}: RTT {rtt.TotalMilliseconds}ms"); } diff --git a/tests/NATS.Client.KeyValueStore.Tests/NatsKVWatcherTest.cs b/tests/NATS.Client.KeyValueStore.Tests/NatsKVWatcherTest.cs index 2589a8927..9d9301fb1 100644 --- a/tests/NATS.Client.KeyValueStore.Tests/NatsKVWatcherTest.cs +++ b/tests/NATS.Client.KeyValueStore.Tests/NatsKVWatcherTest.cs @@ -221,7 +221,7 @@ public async Task Watch_subset() [Fact] public async Task Watcher_timeout_reconnect() { - var timeout = TimeSpan.FromSeconds(30); + var timeout = TimeSpan.FromSeconds(60); var cts = new CancellationTokenSource(timeout); var cancellationToken = cts.Token; diff --git a/tests/NATS.Client.TestUtilities/MockServer.cs b/tests/NATS.Client.TestUtilities/MockServer.cs index 4bde4af8a..0583e6c69 100644 --- a/tests/NATS.Client.TestUtilities/MockServer.cs +++ b/tests/NATS.Client.TestUtilities/MockServer.cs @@ -17,6 +17,7 @@ public class MockServer : IAsyncDisposable private readonly List _clients = new(); private readonly Task _accept; private readonly CancellationTokenSource _cts; + private readonly TaskCompletionSource _ready = new(); private readonly bool _autoPong; public MockServer( @@ -37,6 +38,7 @@ public MockServer( _accept = Task.Run( async () => { + _ready.SetResult(); var n = 0; while (!cancellationToken.IsCancellationRequested) { @@ -151,6 +153,8 @@ public MockServer( public int Port { get; } + public Task Ready => _ready.Task; + public string Url => $"127.0.0.1:{Port}"; public async ValueTask DisposeAsync() From 616121953e30bd241ab6edeb14d1a9641beec8a6 Mon Sep 17 00:00:00 2001 From: Ziya Suzen Date: Thu, 9 Apr 2026 12:35:55 +0100 Subject: [PATCH 4/7] tests: add ready wait and connect retry to mock server tests Apply the same MockServer ready signal and ConnectRetryAsync pattern to ProtocolParserSizeCheckTest and PingCancellationTest to fix connection race on net481 under CI load. --- .../PingCancellationTest.cs | 13 +++++++++---- .../ProtocolParserSizeCheckTest.cs | 17 +++++++++++++---- 2 files changed, 22 insertions(+), 8 deletions(-) diff --git a/tests/NATS.Client.Core2.Tests/PingCancellationTest.cs b/tests/NATS.Client.Core2.Tests/PingCancellationTest.cs index 459c1eeac..4446404c7 100644 --- a/tests/NATS.Client.Core2.Tests/PingCancellationTest.cs +++ b/tests/NATS.Client.Core2.Tests/PingCancellationTest.cs @@ -1,4 +1,5 @@ using NATS.Client.TestUtilities; +using NATS.Client.TestUtilities2; using Synadia.Orbit.Testing.NatsServerProcessManager; namespace NATS.Client.Core.Tests; @@ -19,8 +20,9 @@ public async Task PingAsync_succeeds_with_mock_server() logger: m => _output.WriteLine(m), cancellationToken: cts.Token); + await server.Ready; await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); var rtt = await nats.PingAsync(cts.Token); rtt.Should().BeGreaterThan(TimeSpan.Zero); @@ -36,8 +38,9 @@ public async Task PingAsync_multiple_sequential_pings_succeed() logger: m => _output.WriteLine(m), cancellationToken: cts.Token); + await server.Ready; await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); for (var i = 0; i < 5; i++) { @@ -75,8 +78,9 @@ public async Task PingAsync_throws_when_cancelled_waiting_for_pong() autoPong: false, cancellationToken: cts.Token); + await server.Ready; await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); // This ping should time out because the server won't reply PONG using var pingCts = new CancellationTokenSource(TimeSpan.FromMilliseconds(500)); @@ -95,8 +99,9 @@ public async Task PingAsync_concurrent_pings_all_complete() logger: m => _output.WriteLine(m), cancellationToken: cts.Token); + await server.Ready; await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); // Fire multiple pings concurrently — exercises the pool and concurrent SetResult paths var tasks = new Task[10]; diff --git a/tests/NATS.Client.Core2.Tests/ProtocolParserSizeCheckTest.cs b/tests/NATS.Client.Core2.Tests/ProtocolParserSizeCheckTest.cs index 91fa59019..bd6d03a7e 100644 --- a/tests/NATS.Client.Core2.Tests/ProtocolParserSizeCheckTest.cs +++ b/tests/NATS.Client.Core2.Tests/ProtocolParserSizeCheckTest.cs @@ -3,6 +3,7 @@ using System.Text; using Microsoft.Extensions.Logging; using NATS.Client.TestUtilities; +using NATS.Client.TestUtilities2; namespace NATS.Client.Core.Tests; @@ -22,8 +23,9 @@ public async Task Msg_with_payload_exceeding_max_payload_does_not_oom() var logFactory = new InMemoryTestLoggerFactory(LogLevel.Error, m => output.WriteLine($"[LOG] {m.Message}")); await using var server = new FakeServer(output); + await server.Ready; await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, LoggerFactory = logFactory }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); await nats.SubscribeCoreAsync("foo", cancellationToken: cts.Token); @@ -49,8 +51,9 @@ public async Task Msg_with_negative_payload_length_does_not_oom() var logFactory = new InMemoryTestLoggerFactory(LogLevel.Error, m => output.WriteLine($"[LOG] {m.Message}")); await using var server = new FakeServer(output); + await server.Ready; await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, LoggerFactory = logFactory }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); await nats.SubscribeCoreAsync("foo", cancellationToken: cts.Token); @@ -76,8 +79,9 @@ public async Task Hmsg_with_total_less_than_headers_does_not_oom() var logFactory = new InMemoryTestLoggerFactory(LogLevel.Error, m => output.WriteLine($"[LOG] {m.Message}")); await using var server = new FakeServer(output); + await server.Ready; await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, LoggerFactory = logFactory }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); await nats.SubscribeCoreAsync("foo", cancellationToken: cts.Token); @@ -103,8 +107,9 @@ public async Task Hmsg_with_total_exceeding_max_payload_does_not_oom() var logFactory = new InMemoryTestLoggerFactory(LogLevel.Error, m => output.WriteLine($"[LOG] {m.Message}")); await using var server = new FakeServer(output); + await server.Ready; await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, LoggerFactory = logFactory }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); await nats.SubscribeCoreAsync("foo", cancellationToken: cts.Token); @@ -231,6 +236,7 @@ private sealed class FakeServer : IAsyncDisposable private readonly TcpListener _listener; private readonly CancellationTokenSource _cts = new(TimeSpan.FromSeconds(30)); private readonly TaskCompletionSource _accepted = new(); + private readonly TaskCompletionSource _ready = new(); private readonly Dictionary _subWaiters = new(); private TcpClient? _tcpClient; private StreamWriter? _writer; @@ -250,6 +256,8 @@ public FakeServer(ITestOutputHelper output, string info = "{\"max_payload\":1048 public int Port { get; } + public Task Ready => _ready.Task; + public string Url => $"127.0.0.1:{Port}"; /// @@ -309,6 +317,7 @@ public async ValueTask DisposeAsync() private async Task AcceptAndServeAsync(string info) { + _ready.SetResult(); _tcpClient = await _listener.AcceptTcpClientAsync(); var stream = _tcpClient.GetStream(); var encoding = Encoding.GetEncoding(28591); From 1512f1dd1d2e4594068c7c18cc08ef56480945de Mon Sep 17 00:00:00 2001 From: Ziya Suzen Date: Thu, 9 Apr 2026 13:14:05 +0100 Subject: [PATCH 5/7] tests: fix PingCancellation and SendBuffer flaky tests Use ConnectRetryAsync for real server PingCancellation tests. Add MockServer ready wait, ConnectRetryAsync, and increased CommandTimeout to SendBufferTest to handle 8MB publish under CI load on net481. --- tests/NATS.Client.Core2.Tests/PingCancellationTest.cs | 4 ++-- tests/NATS.Client.Core2.Tests/SendBufferTest.cs | 8 ++++++-- 2 files changed, 8 insertions(+), 4 deletions(-) diff --git a/tests/NATS.Client.Core2.Tests/PingCancellationTest.cs b/tests/NATS.Client.Core2.Tests/PingCancellationTest.cs index 4446404c7..81e8d823b 100644 --- a/tests/NATS.Client.Core2.Tests/PingCancellationTest.cs +++ b/tests/NATS.Client.Core2.Tests/PingCancellationTest.cs @@ -123,7 +123,7 @@ public async Task PingAsync_succeeds_against_real_server() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); var rtt = await nats.PingAsync(); rtt.Should().BeGreaterThan(TimeSpan.Zero); @@ -135,7 +135,7 @@ public async Task PingAsync_times_out_after_server_stopped() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); // Verify ping works while server is up var rtt = await nats.PingAsync(); diff --git a/tests/NATS.Client.Core2.Tests/SendBufferTest.cs b/tests/NATS.Client.Core2.Tests/SendBufferTest.cs index bdc5aa335..c995a720d 100644 --- a/tests/NATS.Client.Core2.Tests/SendBufferTest.cs +++ b/tests/NATS.Client.Core2.Tests/SendBufferTest.cs @@ -2,6 +2,7 @@ using System.Net.Sockets; using Microsoft.Extensions.Logging; using NATS.Client.TestUtilities; +using NATS.Client.TestUtilities2; #if !NET6_0_OR_GREATER using NATS.Client.Core.Internal.NetStandardExtensions; #endif @@ -40,8 +41,9 @@ void Log(string m) await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await server.Ready; Log($"[C] connect {server.Url}"); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); Log($"[C] ping"); var rtt = await nats.PingAsync(cts.Token); @@ -128,10 +130,12 @@ void Log(string m) { Url = server.Url, LoggerFactory = testLogger, + CommandTimeout = TimeSpan.FromSeconds(30), }); + await server.Ready; Log($"[C] connect {server.Url}"); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); Log($"[C] ping"); var rtt = await nats.PingAsync(cts.Token); From 7c4630b100b4d3a82ba51662571a0cc6b1c4353f Mon Sep 17 00:00:00 2001 From: Ziya Suzen Date: Thu, 9 Apr 2026 13:49:44 +0100 Subject: [PATCH 6/7] tests: use ConnectTimeout instead of retry for FakeServer tests FakeServer only accepts one client, so ConnectRetryAsync consumes the single accept slot on the first failed attempt, making all subsequent retries fail. Use increased ConnectTimeout (10s) with plain ConnectAsync instead. --- .../ProtocolParserSizeCheckTest.cs | 17 ++++++++--------- 1 file changed, 8 insertions(+), 9 deletions(-) diff --git a/tests/NATS.Client.Core2.Tests/ProtocolParserSizeCheckTest.cs b/tests/NATS.Client.Core2.Tests/ProtocolParserSizeCheckTest.cs index bd6d03a7e..24f45dc6e 100644 --- a/tests/NATS.Client.Core2.Tests/ProtocolParserSizeCheckTest.cs +++ b/tests/NATS.Client.Core2.Tests/ProtocolParserSizeCheckTest.cs @@ -3,7 +3,6 @@ using System.Text; using Microsoft.Extensions.Logging; using NATS.Client.TestUtilities; -using NATS.Client.TestUtilities2; namespace NATS.Client.Core.Tests; @@ -24,8 +23,8 @@ public async Task Msg_with_payload_exceeding_max_payload_does_not_oom() await using var server = new FakeServer(output); await server.Ready; - await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, LoggerFactory = logFactory }); - await nats.ConnectRetryAsync(); + await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, LoggerFactory = logFactory, ConnectTimeout = TimeSpan.FromSeconds(10) }); + await nats.ConnectAsync(); using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); await nats.SubscribeCoreAsync("foo", cancellationToken: cts.Token); @@ -52,8 +51,8 @@ public async Task Msg_with_negative_payload_length_does_not_oom() await using var server = new FakeServer(output); await server.Ready; - await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, LoggerFactory = logFactory }); - await nats.ConnectRetryAsync(); + await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, LoggerFactory = logFactory, ConnectTimeout = TimeSpan.FromSeconds(10) }); + await nats.ConnectAsync(); using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); await nats.SubscribeCoreAsync("foo", cancellationToken: cts.Token); @@ -80,8 +79,8 @@ public async Task Hmsg_with_total_less_than_headers_does_not_oom() await using var server = new FakeServer(output); await server.Ready; - await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, LoggerFactory = logFactory }); - await nats.ConnectRetryAsync(); + await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, LoggerFactory = logFactory, ConnectTimeout = TimeSpan.FromSeconds(10) }); + await nats.ConnectAsync(); using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); await nats.SubscribeCoreAsync("foo", cancellationToken: cts.Token); @@ -108,8 +107,8 @@ public async Task Hmsg_with_total_exceeding_max_payload_does_not_oom() await using var server = new FakeServer(output); await server.Ready; - await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, LoggerFactory = logFactory }); - await nats.ConnectRetryAsync(); + await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, LoggerFactory = logFactory, ConnectTimeout = TimeSpan.FromSeconds(10) }); + await nats.ConnectAsync(); using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); await nats.SubscribeCoreAsync("foo", cancellationToken: cts.Token); From 50d69daf179ea5562b3d80da4bddc6bb79ba8270 Mon Sep 17 00:00:00 2001 From: Ziya Suzen Date: Thu, 9 Apr 2026 14:04:20 +0100 Subject: [PATCH 7/7] tests: fix remaining ProtocolParserSizeCheck connection races Add ready wait and ConnectTimeout/ConnectRetryAsync to the two tests that were missed in the previous pass (Valid_msg_still_works and Valid_hmsg_still_works). --- .../NATS.Client.Core2.Tests/ProtocolParserSizeCheckTest.cs | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/tests/NATS.Client.Core2.Tests/ProtocolParserSizeCheckTest.cs b/tests/NATS.Client.Core2.Tests/ProtocolParserSizeCheckTest.cs index 24f45dc6e..099ba7400 100644 --- a/tests/NATS.Client.Core2.Tests/ProtocolParserSizeCheckTest.cs +++ b/tests/NATS.Client.Core2.Tests/ProtocolParserSizeCheckTest.cs @@ -3,6 +3,7 @@ using System.Text; using Microsoft.Extensions.Logging; using NATS.Client.TestUtilities; +using NATS.Client.TestUtilities2; namespace NATS.Client.Core.Tests; @@ -182,8 +183,9 @@ public async Task Valid_msg_still_works() logger: m => output.WriteLine(m), cancellationToken: cts.Token); + await server.Ready; await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); await foreach (var msg in nats.SubscribeAsync("foo", cancellationToken: cts.Token)) { @@ -201,7 +203,8 @@ public async Task Valid_hmsg_still_works() { await using var server = new FakeServer(output); - await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await server.Ready; + await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, ConnectTimeout = TimeSpan.FromSeconds(10) }); await nats.ConnectAsync(); using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10));