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/PingCancellationTest.cs b/tests/NATS.Client.Core2.Tests/PingCancellationTest.cs index 459c1eeac..81e8d823b 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]; @@ -118,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); @@ -130,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/ProtocolParserSizeCheckTest.cs b/tests/NATS.Client.Core2.Tests/ProtocolParserSizeCheckTest.cs index 91fa59019..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; @@ -22,7 +23,8 @@ 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 using var nats = new NatsConnection(new NatsOpts { Url = server.Url, LoggerFactory = logFactory }); + await server.Ready; + 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)); @@ -49,7 +51,8 @@ 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 using var nats = new NatsConnection(new NatsOpts { Url = server.Url, LoggerFactory = logFactory }); + await server.Ready; + 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)); @@ -76,7 +79,8 @@ 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 using var nats = new NatsConnection(new NatsOpts { Url = server.Url, LoggerFactory = logFactory }); + await server.Ready; + 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)); @@ -103,7 +107,8 @@ 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 using var nats = new NatsConnection(new NatsOpts { Url = server.Url, LoggerFactory = logFactory }); + await server.Ready; + 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)); @@ -178,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)) { @@ -197,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)); @@ -231,6 +238,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 +258,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 +319,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); 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/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); 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()