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
1 change: 1 addition & 0 deletions tests/NATS.Client.Core2.Tests/BufferReferenceTest.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
17 changes: 11 additions & 6 deletions tests/NATS.Client.Core2.Tests/PingCancellationTest.cs
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
using NATS.Client.TestUtilities;
using NATS.Client.TestUtilities2;
using Synadia.Orbit.Testing.NatsServerProcessManager;

namespace NATS.Client.Core.Tests;
Expand All @@ -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);
Expand All @@ -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++)
{
Expand Down Expand Up @@ -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));
Expand All @@ -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<TimeSpan>[10];
Expand All @@ -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);
Expand All @@ -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();
Expand Down
23 changes: 17 additions & 6 deletions tests/NATS.Client.Core2.Tests/ProtocolParserSizeCheckTest.cs
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
using System.Text;
using Microsoft.Extensions.Logging;
using NATS.Client.TestUtilities;
using NATS.Client.TestUtilities2;

namespace NATS.Client.Core.Tests;

Expand All @@ -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));
Expand All @@ -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));
Expand All @@ -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));
Expand All @@ -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));
Expand Down Expand Up @@ -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<string>("foo", cancellationToken: cts.Token))
{
Expand All @@ -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));
Expand Down Expand Up @@ -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<string, TaskCompletionSource> _subWaiters = new();
private TcpClient? _tcpClient;
private StreamWriter? _writer;
Expand All @@ -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}";

/// <summary>
Expand Down Expand Up @@ -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);
Expand Down
2 changes: 2 additions & 0 deletions tests/NATS.Client.Core2.Tests/RequestReplyTest.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -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<int>("foo");
var reg = sub.Register(async msg =>
Expand Down
8 changes: 6 additions & 2 deletions tests/NATS.Client.Core2.Tests/SendBufferTest.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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);
Expand Down
10 changes: 7 additions & 3 deletions tests/NATS.Client.Core2.Tests/Utf8SubjectTest.cs
Original file line number Diff line number Diff line change
@@ -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;

Expand Down Expand Up @@ -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<string>(">", cancellationToken: cts.Token))
{
Expand Down Expand Up @@ -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<string>(">", cancellationToken: cts.Token))
{
Expand Down Expand Up @@ -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<string>(">", cancellationToken: cts.Token))
{
Expand Down
5 changes: 4 additions & 1 deletion tests/NATS.Client.JetStream.Tests/SlowConsumerTest.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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");
}
Expand Down
2 changes: 1 addition & 1 deletion tests/NATS.Client.KeyValueStore.Tests/NatsKVWatcherTest.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
4 changes: 4 additions & 0 deletions tests/NATS.Client.TestUtilities/MockServer.cs
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ public class MockServer : IAsyncDisposable
private readonly List<Task> _clients = new();
private readonly Task _accept;
private readonly CancellationTokenSource _cts;
private readonly TaskCompletionSource _ready = new();
private readonly bool _autoPong;

public MockServer(
Expand All @@ -37,6 +38,7 @@ public MockServer(
_accept = Task.Run(
async () =>
{
_ready.SetResult();
var n = 0;
while (!cancellationToken.IsCancellationRequested)
{
Expand Down Expand Up @@ -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()
Expand Down
Loading