diff --git a/tests/NATS.Client.Core.Tests/LowLevelApiTest.cs b/tests/NATS.Client.Core.Tests/LowLevelApiTest.cs index 7b259630f..e796eb518 100644 --- a/tests/NATS.Client.Core.Tests/LowLevelApiTest.cs +++ b/tests/NATS.Client.Core.Tests/LowLevelApiTest.cs @@ -1,5 +1,6 @@ using System.Buffers; using System.Text; +using NATS.Client.TestUtilities2; using Synadia.Orbit.Testing.NatsServerProcessManager; namespace NATS.Client.Core.Tests; @@ -15,6 +16,7 @@ public async Task Sub_custom_builder_test() { await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var subject = "foo.*"; var builder = new NatsSubCustomTestBuilder(_output); diff --git a/tests/NATS.Client.Core.Tests/ProtocolTest.cs b/tests/NATS.Client.Core.Tests/ProtocolTest.cs index 464d3068b..4f8c430ee 100644 --- a/tests/NATS.Client.Core.Tests/ProtocolTest.cs +++ b/tests/NATS.Client.Core.Tests/ProtocolTest.cs @@ -1,5 +1,6 @@ using Microsoft.Extensions.Logging; using NATS.Client.TestUtilities; +using NATS.Client.TestUtilities2; using Synadia.Orbit.Testing.NatsServerProcessManager; namespace NATS.Client.Core.Tests; @@ -16,6 +17,7 @@ public async Task Protocol_parser_under_load(int size) var logger = new InMemoryTestLoggerFactory(LogLevel.Error); var opts = new NatsOpts { Url = server.Url, LoggerFactory = logger }; var nats = new NatsConnection(opts); + await nats.ConnectRetryAsync(); using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(120)); diff --git a/tests/NATS.Client.Core.Tests/SubscriptionTest.cs b/tests/NATS.Client.Core.Tests/SubscriptionTest.cs index 2124f9822..ef79ee316 100644 --- a/tests/NATS.Client.Core.Tests/SubscriptionTest.cs +++ b/tests/NATS.Client.Core.Tests/SubscriptionTest.cs @@ -1,4 +1,5 @@ using System.Net; +using NATS.Client.TestUtilities2; using Synadia.Orbit.Testing.NatsServerProcessManager; namespace NATS.Client.Core.Tests; @@ -35,6 +36,7 @@ public async Task Subscription_periodic_cleanup_test() // SharedInbox so the connection doesn't open an inbox subscription at connect, // which would make the SUB frame count below never settle at 1. var nats = new NatsConnection(new NatsOpts { Url = $"nats://127.0.0.1:{proxy.Port}", SubscriptionCleanUpInterval = TimeSpan.FromSeconds(1), RequestReplyMode = NatsRequestReplyMode.SharedInbox }); + await nats.ConnectRetryAsync(); async Task Isolator() { @@ -75,6 +77,7 @@ public async Task Subscription_cleanup_on_message_receive_test() // SharedInbox so the connection doesn't open an inbox subscription at connect, // which would make the SUB frame count below never settle at 1. var nats = new NatsConnection(new NatsOpts { Url = $"nats://127.0.0.1:{proxy.Port}", SubscriptionCleanUpInterval = TimeSpan.MaxValue, RequestReplyMode = NatsRequestReplyMode.SharedInbox }); + await nats.ConnectRetryAsync(); async Task Isolator() { @@ -108,6 +111,7 @@ public async Task Auto_unsubscribe_on_max_messages_with_inbox_subscription_test( { await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var subject = nats.NewInbox(); await using var sub1 = await nats.SubscribeCoreAsync(subject, opts: new NatsSubOpts { MaxMsgs = 1 }); @@ -147,6 +151,7 @@ public async Task Auto_unsubscribe_on_max_messages_test() { await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); const string subject = "foo1"; const int maxMsgs = 99; var opts = new NatsSubOpts { MaxMsgs = maxMsgs }; @@ -177,6 +182,7 @@ public async Task Auto_unsubscribe_on_timeout_test() { await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); const string subject = "foo2"; var opts = new NatsSubOpts { Timeout = TimeSpan.FromSeconds(1) }; @@ -200,6 +206,7 @@ public async Task Auto_unsubscribe_on_idle_timeout_test() { await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); const string subject = "foo3"; var opts = new NatsSubOpts { IdleTimeout = TimeSpan.FromSeconds(3) }; @@ -231,6 +238,7 @@ public async Task Subscribe_with_no_timeout_sentinels_does_not_throw() { await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); const string subject = "foo-max"; var opts = new NatsSubOpts @@ -266,6 +274,7 @@ public async Task Manual_unsubscribe_test() { await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); const string subject = "foo4"; await using var sub = await nats.SubscribeCoreAsync(subject); @@ -352,6 +361,7 @@ public async Task Serialization_exceptions() { await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); diff --git a/tests/NATS.Client.Core2.Tests/ErrorHandlerTest.cs b/tests/NATS.Client.Core2.Tests/ErrorHandlerTest.cs index ef9a19f6c..d9a041b51 100644 --- a/tests/NATS.Client.Core2.Tests/ErrorHandlerTest.cs +++ b/tests/NATS.Client.Core2.Tests/ErrorHandlerTest.cs @@ -1,6 +1,7 @@ using Microsoft.Extensions.Logging; using NATS.Client.Core.Tests; using NATS.Client.TestUtilities; +using NATS.Client.TestUtilities2; namespace NATS.Client.Core2.Tests; @@ -29,6 +30,7 @@ public async Task Handle_permissions_violation() LoggerFactory = logger, AuthOpts = new NatsAuthOpts { Username = "u" }, }); + await nats.ConnectRetryAsync(); var errors = new List(); diff --git a/tests/NATS.Client.Core2.Tests/JsonSerializerTests.cs b/tests/NATS.Client.Core2.Tests/JsonSerializerTests.cs index f23480240..54c175d0b 100644 --- a/tests/NATS.Client.Core2.Tests/JsonSerializerTests.cs +++ b/tests/NATS.Client.Core2.Tests/JsonSerializerTests.cs @@ -1,5 +1,6 @@ using NATS.Client.Core2.Tests; using NATS.Client.Serializers.Json; +using NATS.Client.TestUtilities2; namespace NATS.Client.Core.Tests; @@ -24,6 +25,7 @@ public async Task Serialize_any_type() SerializerRegistry = NatsJsonSerializerRegistry.Default, }; await using var nats = new NatsConnection(natsOpts); + await nats.ConnectRetryAsync(); // in local runs server start is taking too long when running // the whole suite. @@ -39,6 +41,7 @@ public async Task Serialize_any_type() // Default serializer won't work with random types await using var nats1 = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats1.ConnectRetryAsync(); var exception = await Assert.ThrowsAsync(() => nats1.PublishAsync( subject: "would.not.work", diff --git a/tests/NATS.Client.Core2.Tests/MessageInterfaceTest.cs b/tests/NATS.Client.Core2.Tests/MessageInterfaceTest.cs index 526500545..0acb019fe 100644 --- a/tests/NATS.Client.Core2.Tests/MessageInterfaceTest.cs +++ b/tests/NATS.Client.Core2.Tests/MessageInterfaceTest.cs @@ -1,4 +1,5 @@ using NATS.Client.Core2.Tests; +using NATS.Client.TestUtilities2; using Synadia.Orbit.Testing.NatsServerProcessManager; namespace NATS.Client.Core.Tests; @@ -17,6 +18,7 @@ public MessageInterfaceTest(NatsServerFixture server) public async Task Sub_custom_builder_test() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var sync = 0; var sub = Task.Run(async () => diff --git a/tests/NATS.Client.Core2.Tests/ProtocolParserSizeCheckTest.cs b/tests/NATS.Client.Core2.Tests/ProtocolParserSizeCheckTest.cs index 1d26bb219..2ce1e6fc8 100644 --- a/tests/NATS.Client.Core2.Tests/ProtocolParserSizeCheckTest.cs +++ b/tests/NATS.Client.Core2.Tests/ProtocolParserSizeCheckTest.cs @@ -25,7 +25,7 @@ public async Task Msg_with_payload_exceeding_max_payload_does_not_oom() await server.Ready; await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, LoggerFactory = logFactory, ConnectTimeout = TimeSpan.FromSeconds(30) }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); await nats.SubscribeCoreAsync("foo", cancellationToken: cts.Token); @@ -53,7 +53,7 @@ public async Task Msg_with_negative_payload_length_does_not_oom() await server.Ready; await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, LoggerFactory = logFactory, ConnectTimeout = TimeSpan.FromSeconds(10) }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); await nats.SubscribeCoreAsync("foo", cancellationToken: cts.Token); @@ -81,7 +81,7 @@ public async Task Hmsg_with_total_less_than_headers_does_not_oom() await server.Ready; await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, LoggerFactory = logFactory, ConnectTimeout = TimeSpan.FromSeconds(10) }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); await nats.SubscribeCoreAsync("foo", cancellationToken: cts.Token); @@ -109,7 +109,7 @@ public async Task Hmsg_with_total_exceeding_max_payload_does_not_oom() await server.Ready; await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, LoggerFactory = logFactory, ConnectTimeout = TimeSpan.FromSeconds(10) }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); await nats.SubscribeCoreAsync("foo", cancellationToken: cts.Token); @@ -138,7 +138,7 @@ public async Task Protocol_violation_exits_read_loop_cleanly() // MaxReconnectRetry=0 because FakeServer only accepts one connection await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, LoggerFactory = logFactory, MaxReconnectRetry = 0 }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); await nats.SubscribeCoreAsync("foo", cancellationToken: cts.Token); @@ -208,7 +208,7 @@ public async Task Valid_hmsg_still_works() // SharedInbox so the connection doesn't open an inbox subscription at connect; // the test injects an HMSG with sid 1, which must map to the "foo" subscription. await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, ConnectTimeout = TimeSpan.FromSeconds(10), RequestReplyMode = NatsRequestReplyMode.SharedInbox }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); var sub = await nats.SubscribeCoreAsync("foo", cancellationToken: cts.Token); diff --git a/tests/NATS.Client.Core2.Tests/ProtocolTest.cs b/tests/NATS.Client.Core2.Tests/ProtocolTest.cs index 05b6e0747..757423e4d 100644 --- a/tests/NATS.Client.Core2.Tests/ProtocolTest.cs +++ b/tests/NATS.Client.Core2.Tests/ProtocolTest.cs @@ -2,6 +2,7 @@ using System.Text; using NATS.Client.Core2.Tests; using NATS.Client.Core2.Tests.ExtraUtils.FrameworkPolyfillExtensions; +using NATS.Client.TestUtilities2; namespace NATS.Client.Core.Tests; @@ -21,11 +22,13 @@ public ProtocolTest(ITestOutputHelper output, NatsServerFixture server) public async Task Subscription_with_same_subject() { var nats1 = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats1.ConnectRetryAsync(); var proxy = new NatsProxy(_server.Port); // SharedInbox so the connection doesn't open an inbox subscription at connect, // which would add an extra SUB frame to the proxy capture asserted below. var nats2 = new NatsConnection(new NatsOpts { Url = $"nats://127.0.0.1:{proxy.Port}", ConnectTimeout = TimeSpan.FromSeconds(10), RequestReplyMode = NatsRequestReplyMode.SharedInbox }); + await nats2.ConnectRetryAsync(); var sub1 = await nats2.SubscribeCoreAsync("foo.bar"); var sub2 = await nats2.SubscribeCoreAsync("foo.bar"); @@ -119,6 +122,7 @@ public async Task Subscription_queue_group() // SharedInbox so the connection doesn't open an inbox subscription at connect, // which would shift the SUB frames asserted by index below. var nats = new NatsConnection(new NatsOpts { Url = $"nats://127.0.0.1:{proxy.Port}", ConnectTimeout = TimeSpan.FromSeconds(10), RequestReplyMode = NatsRequestReplyMode.SharedInbox }); + await nats.ConnectRetryAsync(); var subject = $"{_server.GetNextId()}.foo"; await using var sub1 = await nats.SubscribeCoreAsync(subject, queueGroup: "group1"); @@ -151,6 +155,7 @@ void Log(string text) var proxy = new NatsProxy(_server.Port); var nats = new NatsConnection(new NatsOpts { Url = $"nats://127.0.0.1:{proxy.Port}", ConnectTimeout = TimeSpan.FromSeconds(10) }); + await nats.ConnectRetryAsync(); var prefix = $"{_server.GetNextId()}.foo"; @@ -223,6 +228,7 @@ void Log(string text) // SharedInbox so the connection doesn't open an inbox subscription at connect, // which would consume sid 1 and break the sid sequence asserted below. var nats = new NatsConnection(new NatsOpts { Url = $"nats://127.0.0.1:{proxy.Port}", ConnectTimeout = TimeSpan.FromSeconds(10), RequestReplyMode = NatsRequestReplyMode.SharedInbox }); + await nats.ConnectRetryAsync(); var sid = 0; Log("### Auto-unsubscribe after consuming max-msgs"); @@ -350,6 +356,7 @@ public async Task Reconnect_with_sub_and_additional_commands() // SharedInbox so the connection doesn't open an inbox subscription at connect, // which would add an extra SUB frame to the proxy capture asserted below. var nats = new NatsConnection(new NatsOpts { Url = $"nats://127.0.0.1:{proxy.Port}", ConnectTimeout = TimeSpan.FromSeconds(10), RequestReplyMode = NatsRequestReplyMode.SharedInbox }); + await nats.ConnectRetryAsync(); var subject = $"{_server.GetNextId()}.foo"; var cmdSubject = $"{_server.GetNextId()}.bar"; @@ -406,6 +413,7 @@ await Retry.Until( public async Task Proactively_reject_payloads_over_the_threshold_set_by_server() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); diff --git a/tests/NATS.Client.Core2.Tests/RequestReplyTest.cs b/tests/NATS.Client.Core2.Tests/RequestReplyTest.cs index 7c0cfae0e..c228c469e 100644 --- a/tests/NATS.Client.Core2.Tests/RequestReplyTest.cs +++ b/tests/NATS.Client.Core2.Tests/RequestReplyTest.cs @@ -24,6 +24,7 @@ public RequestReplyTest(ITestOutputHelper output, NatsServerFixture server) public async Task Simple_request_reply_test() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); const string subject = "foo"; var cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); @@ -58,6 +59,7 @@ public async Task Request_reply_command_timeout_test() Url = _server.Url, RequestTimeout = TimeSpan.FromSeconds(1), }); + await nats.ConnectRetryAsync(); var sub = await nats.SubscribeCoreAsync("foo"); var reg = sub.Register(async msg => @@ -76,6 +78,7 @@ await Assert.ThrowsAsync(async () => // Cancellation token usage { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var cts = new CancellationTokenSource(TimeSpan.FromSeconds(1)); var sub = await nats.SubscribeCoreAsync("foo"); @@ -103,6 +106,7 @@ public async Task Request_reply_no_responders_default_throws_test() // Default mode (Direct transport, not set explicitly) throws on no-responders, // preserving the pre-3.x default behavior. await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); await Assert.ThrowsAsync(async () => await nats.RequestAsync(Guid.NewGuid().ToString(), 0)); } @@ -110,6 +114,7 @@ public async Task Request_reply_no_responders_default_throws_test() public async Task Request_reply_no_responders_shared_inbox_throws_test() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url, RequestReplyMode = NatsRequestReplyMode.SharedInbox }); + await nats.ConnectRetryAsync(); await Assert.ThrowsAsync(async () => await nats.RequestAsync(Guid.NewGuid().ToString(), 0)); } @@ -119,6 +124,7 @@ public async Task Request_reply_no_responders_explicit_direct_returns_message_te // Explicitly selecting Direct preserves the pre-3.x Direct behavior: the 503 sentinel // comes back as a message with HasNoResponders set instead of throwing. await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url, RequestReplyMode = NatsRequestReplyMode.Direct }); + await nats.ConnectRetryAsync(); var reply = await nats.RequestAsync(Guid.NewGuid().ToString(), 0); Assert.True(reply.HasNoResponders); } @@ -130,6 +136,7 @@ public async Task Request_reply_no_responders_per_call_suppressed_test(NatsReque { // Per-call ThrowIfNoResponders=false returns the sentinel as a message in either mode. await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var reply = await nats.RequestAsync(Guid.NewGuid().ToString(), 0, replyOpts: new NatsSubOpts { ThrowIfNoResponders = false }); Assert.True(reply.HasNoResponders); } @@ -139,6 +146,7 @@ public async Task Request_reply_no_responders_per_call_throw_overrides_explicit_ { // Per-call ThrowIfNoResponders=true forces a throw even when Direct was selected explicitly. await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url, RequestReplyMode = NatsRequestReplyMode.Direct }); + await nats.ConnectRetryAsync(); await Assert.ThrowsAsync(async () => await nats.RequestAsync(Guid.NewGuid().ToString(), 0, replyOpts: new NatsSubOpts { ThrowIfNoResponders = true })); } @@ -149,6 +157,7 @@ public async Task Request_reply_many_test() const int msgs = 10; 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 => @@ -160,6 +169,7 @@ public async Task Request_reply_many_test() await msg.ReplyAsync(null); // stop iteration with a sentinel }); + await nats.PingAsync(); var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); const int data = 2; @@ -180,6 +190,7 @@ public async Task Request_reply_many_test() public async Task Request_reply_many_test_overall_timeout() { 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 => @@ -187,6 +198,7 @@ public async Task Request_reply_many_test_overall_timeout() await msg.ReplyAsync(msg.Data * 2); await msg.ReplyAsync(msg.Data * 3); }); + await nats.PingAsync(); var results = new[] { 8, 12 }; var count = 0; @@ -210,6 +222,7 @@ public async Task Request_reply_many_test_overall_timeout() public async Task Request_reply_many_test_idle_timeout() { 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 => @@ -220,6 +233,7 @@ public async Task Request_reply_many_test_idle_timeout() await Task.Delay(5_000); await msg.ReplyAsync(msg.Data * 4); }); + await nats.PingAsync(); var results = new[] { 6, 9 }; var count = 0; @@ -243,6 +257,7 @@ public async Task Request_reply_many_test_idle_timeout() public async Task Request_reply_many_test_start_up_timeout() { 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 => @@ -250,6 +265,7 @@ public async Task Request_reply_many_test_start_up_timeout() await Task.Delay(2_000); await msg.ReplyAsync(msg.Data * 2); }); + await nats.PingAsync(); var count = 0; var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); @@ -282,6 +298,7 @@ public async Task Request_reply_many_test_max_count() await msg.ReplyAsync(msg.Data * 4); await msg.ReplyAsync(null); // sentinel }); + await nats.PingAsync(); var results = new[] { 2, 3, 4 }; var count = 0; @@ -305,6 +322,7 @@ public async Task Request_reply_many_test_max_count() public async Task Request_reply_many_test_sentinel() { 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 => @@ -314,6 +332,7 @@ public async Task Request_reply_many_test_sentinel() await msg.ReplyAsync(msg.Data * 4); await msg.ReplyAsync(null); // sentinel }); + await nats.PingAsync(); var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); var results = new[] { 2, 3, 4 }; @@ -338,6 +357,7 @@ static string ToStr(ReadOnlyMemory input) } await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); @@ -372,7 +392,7 @@ public async Task Request_reply_many_multiple_with_timeout_test() await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); // connect to avoid race to subscribe and publish - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); const string subject = "foo"; var cts = new CancellationTokenSource(TimeSpan.FromSeconds(60)); @@ -436,6 +456,7 @@ public async Task Request_reply_many_multiple_with_timeout_test() public async Task Simple_empty_request_reply_test() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); const string subject = "foo"; var cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); @@ -461,6 +482,7 @@ public async Task Simple_empty_request_reply_test() public async Task Direct_request_reply_test() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url, RequestReplyMode = NatsRequestReplyMode.Direct }); + await nats.ConnectRetryAsync(); var reply = await nats.RequestAsync("$JS.API.INFO", cancellationToken: default); // reply-to should be inbox with id @@ -481,6 +503,7 @@ public async Task Default_mode_is_direct_test() // Default RequestReplyMode is Direct, so the reply-to inbox carries a numeric id // e.g. _INBOX.Hu5HPpWesrJhvQq2NG3YJ6.1 await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); Assert.Equal(NatsRequestReplyMode.Direct, nats.Opts.RequestReplyMode); var reply = await nats.RequestAsync("$JS.API.INFO", cancellationToken: default); @@ -499,6 +522,7 @@ public async Task Default_mode_is_direct_test() public async Task SharedInbox_request_reply_test() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url, RequestReplyMode = NatsRequestReplyMode.SharedInbox }); + await nats.ConnectRetryAsync(); var reply = await nats.RequestAsync("$JS.API.INFO", cancellationToken: default); // reply-to should be inbox @@ -517,6 +541,7 @@ public async Task SharedInbox_request_reply_test() public async Task Request_reply_no_timeout_sentinels_do_not_throw(NatsRequestReplyMode mode) { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var subject = $"foo.{Guid.NewGuid():N}"; await using var sub = await nats.SubscribeCoreAsync(subject); @@ -546,6 +571,7 @@ public async Task Request_reply_no_timeout_sentinels_do_not_throw(NatsRequestRep public async Task Request_reply_out_of_range_timeout_throws(NatsRequestReplyMode mode) { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var replyOpts = new NatsSubOpts { Timeout = TimeSpan.FromDays(50) }; diff --git a/tests/NATS.Client.Core2.Tests/SerializerTest.cs b/tests/NATS.Client.Core2.Tests/SerializerTest.cs index 89fedfaa8..99732ed94 100644 --- a/tests/NATS.Client.Core2.Tests/SerializerTest.cs +++ b/tests/NATS.Client.Core2.Tests/SerializerTest.cs @@ -5,6 +5,7 @@ using System.Text.Json.Serialization.Metadata; using NATS.Client.Core2.Tests; using NATS.Client.Serializers.Json; +using NATS.Client.TestUtilities2; // ReSharper disable RedundantTypeArgumentsOfMethod // ReSharper disable ReturnTypeCanBeNotNullable @@ -25,7 +26,7 @@ public async Task Serializer_exceptions() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); await Assert.ThrowsAsync(() => nats.PublishAsync( @@ -51,7 +52,7 @@ public async Task NatsMemoryOwner_empty_payload_should_not_throw() using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); var cancellationToken = cts.Token; - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); var sub = await nats.SubscribeCoreAsync>("foo", cancellationToken: cancellationToken); await nats.PingAsync(cancellationToken); @@ -242,7 +243,7 @@ public async Task Deserialize_with_empty_should_still_go_through_the_deserialize using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); var cancellationToken = cts.Token; - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); var serializer = new TestSerializerWithEmpty(); var sub = await nats.SubscribeCoreAsync("foo", serializer: serializer, cancellationToken: cancellationToken); @@ -271,7 +272,7 @@ public async Task Deserialize_chained_with_empty() using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); var cancellationToken = cts.Token; - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); var sub = await nats.SubscribeCoreAsync("foo", cancellationToken: cancellationToken); @@ -307,7 +308,7 @@ public async Task Deserialize_using_json_stream_serializer_registry() using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); var cancellationToken = cts.Token; - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); var sub1 = await nats.SubscribeCoreAsync($"{prefix}.1", cancellationToken: cancellationToken); var sub2 = await nats.SubscribeCoreAsync($"{prefix}.2", cancellationToken: cancellationToken); @@ -336,7 +337,7 @@ public async Task Serializer_can_mutate_headers_during_serialization() using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); var cancellationToken = cts.Token; - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); var subject = _server.GetNextId(); var sub = await nats.SubscribeCoreAsync(subject, cancellationToken: cancellationToken); diff --git a/tests/NATS.Client.Core2.Tests/SlowConsumerTest.cs b/tests/NATS.Client.Core2.Tests/SlowConsumerTest.cs index ecb5f88cd..0ac10dd66 100644 --- a/tests/NATS.Client.Core2.Tests/SlowConsumerTest.cs +++ b/tests/NATS.Client.Core2.Tests/SlowConsumerTest.cs @@ -1,4 +1,5 @@ using NATS.Client.Core2.Tests; +using NATS.Client.TestUtilities2; namespace NATS.Client.Core.Tests; @@ -22,6 +23,7 @@ public async Task Slow_consumer() Url = _server.Url, SubPendingChannelCapacity = 3, }); + await nats.ConnectRetryAsync(); var lost = 0; nats.MessageDropped += (_, e) => @@ -108,6 +110,7 @@ public async Task SlowConsumerDetected_fires_once_per_episode() Url = _server.Url, SubPendingChannelCapacity = 3, }); + await nats.ConnectRetryAsync(); var droppedCount = 0; var slowConsumerCount = 0; @@ -206,6 +209,7 @@ public async Task SlowConsumerDetected_fires_again_after_recovery() Url = _server.Url, SubPendingChannelCapacity = 3, }); + await nats.ConnectRetryAsync(); var droppedCount = 0; var slowConsumerCount = 0; @@ -317,7 +321,7 @@ await Retry.Until( // Wait for more drops and the second slow consumer event await Retry.Until( "episode 2 drops", - () => Volatile.Read(ref droppedCount) > droppedAfterEpisode1 && Volatile.Read(ref slowConsumerCount) >= 2); + () => Volatile.Read(ref droppedCount) > droppedAfterEpisode1 && Volatile.Read(ref slowConsumerCount) > slowConsumerAfterEpisode1); var droppedAfterEpisode2 = Volatile.Read(ref droppedCount); var slowConsumerAfterEpisode2 = Volatile.Read(ref slowConsumerCount); diff --git a/tests/NATS.Client.Core2.Tests/SubscriptionDrainTest.cs b/tests/NATS.Client.Core2.Tests/SubscriptionDrainTest.cs index 56b8d63d0..1c8ba3e3a 100644 --- a/tests/NATS.Client.Core2.Tests/SubscriptionDrainTest.cs +++ b/tests/NATS.Client.Core2.Tests/SubscriptionDrainTest.cs @@ -1,4 +1,5 @@ using NATS.Client.Core2.Tests; +using NATS.Client.TestUtilities2; namespace NATS.Client.Core.Tests; @@ -16,6 +17,7 @@ public SubscriptionDrainTest(NatsServerFixture server) public async Task Drain_preserves_buffered_messages_and_completes_channel() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); var cancellationToken = cts.Token; @@ -44,6 +46,7 @@ public async Task Drain_preserves_buffered_messages_and_completes_channel() public async Task Drain_delivers_messages_still_in_flight_after_unsub() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); var cancellationToken = cts.Token; @@ -76,6 +79,7 @@ public async Task Drain_delivers_messages_still_in_flight_after_unsub() public async Task Drain_stops_new_messages_but_connection_stays_usable() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); var cancellationToken = cts.Token; @@ -109,6 +113,7 @@ public async Task Drain_stops_new_messages_but_connection_stays_usable() public async Task Drain_is_idempotent_and_dispose_after_drain_is_safe() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); var cancellationToken = cts.Token; diff --git a/tests/NATS.Client.Core2.Tests/TlsPreferTest.cs b/tests/NATS.Client.Core2.Tests/TlsPreferTest.cs index 7fa4c73d1..bfb46242f 100644 --- a/tests/NATS.Client.Core2.Tests/TlsPreferTest.cs +++ b/tests/NATS.Client.Core2.Tests/TlsPreferTest.cs @@ -2,6 +2,7 @@ using System.Net.Sockets; using System.Security.Authentication; using System.Text; +using NATS.Client.TestUtilities2; using Synadia.Orbit.Testing.GoHarness; using Synadia.Orbit.Testing.NatsServerProcessManager; @@ -86,7 +87,7 @@ public async Task Prefer_mode_sends_credentials_in_plaintext_when_info_has_no_tl ConnectTimeout = TimeSpan.FromSeconds(10), }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); await serverTask; listener.Stop(); @@ -371,7 +372,7 @@ public async Task Disable_mode_skips_tls_when_server_advertises_tls_available() ConnectTimeout = TimeSpan.FromSeconds(10), }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); await serverTask; listener.Stop(); @@ -419,7 +420,7 @@ public async Task Real_server_tls_available_prefer_upgrades_disable_stays_plaint TlsOpts = new NatsTlsOpts { Mode = TlsMode.Disable }, }); - await natsDisable.ConnectAsync(); + await natsDisable.ConnectRetryAsync(); await natsDisable.PingAsync(cts.Token); output.WriteLine("Disable mode: connected in plaintext"); } diff --git a/tests/NATS.Client.Core2.Tests/Utf8SubjectTest.cs b/tests/NATS.Client.Core2.Tests/Utf8SubjectTest.cs index c42472fbc..086e20282 100644 --- a/tests/NATS.Client.Core2.Tests/Utf8SubjectTest.cs +++ b/tests/NATS.Client.Core2.Tests/Utf8SubjectTest.cs @@ -151,6 +151,7 @@ public class Utf8SubjectServerTest public async Task Utf8_subject_pub_sub_with_real_server() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var subject = "test.café.🔥"; var sync = 0; @@ -188,6 +189,7 @@ await Retry.Until( public async Task Utf8_subject_request_reply_with_real_server() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var subject = "svc.ñoño.日本語"; var sync = 0; diff --git a/tests/NATS.Client.JetStream.Tests/CancellationTokenTests.cs b/tests/NATS.Client.JetStream.Tests/CancellationTokenTests.cs index 5b4a25426..22dcda2c5 100644 --- a/tests/NATS.Client.JetStream.Tests/CancellationTokenTests.cs +++ b/tests/NATS.Client.JetStream.Tests/CancellationTokenTests.cs @@ -12,6 +12,7 @@ public async Task FetchAsync_with_cancelled_token_throws_immediately() { using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var prefix = server.GetNextId(); var js = new NatsJSContext(nats); await js.CreateStreamAsync($"{prefix}s1", [$"{prefix}s1.*"], cts.Token); @@ -43,6 +44,7 @@ public async Task FetchNoWaitAsync_with_cancelled_token_throws_immediately() { using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var prefix = server.GetNextId(); var js = new NatsJSContext(nats); await js.CreateStreamAsync($"{prefix}s1", [$"{prefix}s1.*"], cts.Token); @@ -74,6 +76,7 @@ public async Task NextAsync_with_cancelled_token_throws_immediately() { using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, ConnectTimeout = TimeSpan.FromSeconds(30) }); + await nats.ConnectRetryAsync(); var prefix = server.GetNextId(); var js = new NatsJSContext(nats); await js.CreateStreamAsync($"{prefix}s1", [$"{prefix}s1.*"], cts.Token); @@ -100,6 +103,7 @@ public async Task ConsumeAsync_with_cancelled_token_throws_immediately() { using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var prefix = server.GetNextId(); var js = new NatsJSContext(nats); await js.CreateStreamAsync($"{prefix}s1", [$"{prefix}s1.*"], cts.Token); diff --git a/tests/NATS.Client.JetStream.Tests/ConsumeResponseTest.cs b/tests/NATS.Client.JetStream.Tests/ConsumeResponseTest.cs index f22eff883..5258277e7 100644 --- a/tests/NATS.Client.JetStream.Tests/ConsumeResponseTest.cs +++ b/tests/NATS.Client.JetStream.Tests/ConsumeResponseTest.cs @@ -1,5 +1,6 @@ using System.Text; using NATS.Client.TestUtilities; +using NATS.Client.TestUtilities2; using NATS.Net; namespace NATS.Client.JetStream.Tests; @@ -38,6 +39,7 @@ public async Task Consume_response(NatsRequestReplyMode mode) var cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); await using var nats = new NatsConnection(new NatsOpts { Url = ms.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var js = nats.CreateJetStreamContext(); var consumer = await js.GetConsumerAsync("x", "x", cts.Token); diff --git a/tests/NATS.Client.JetStream.Tests/ConsumerConsumeTest.cs b/tests/NATS.Client.JetStream.Tests/ConsumerConsumeTest.cs index 1b2fbc7a7..7c7aa874e 100644 --- a/tests/NATS.Client.JetStream.Tests/ConsumerConsumeTest.cs +++ b/tests/NATS.Client.JetStream.Tests/ConsumerConsumeTest.cs @@ -282,6 +282,7 @@ public async Task Consume_dispose_test() Url = _server.Url, DrainSubscriptionsOnDispose = true, }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var cts = new CancellationTokenSource(TimeSpan.FromSeconds(60)); @@ -380,6 +381,7 @@ public async Task Consume_connection_dispose_drains_buffered_messages() await using (var setup = new NatsConnection(new NatsOpts { Url = _server.Url })) { + await setup.ConnectRetryAsync(); var setupJs = new NatsJSContext(setup); var stream = await setupJs.CreateStreamAsync(new StreamConfig(streamName, [subject]), cts.Token); @@ -395,6 +397,7 @@ public async Task Consume_connection_dispose_drains_buffered_messages() DrainSubscriptionsOnDispose = true, ConsumerDrainOnDisposeTimeout = TimeSpan.FromSeconds(30), }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var consumer = await js.GetConsumerAsync(streamName, consumerName, cts.Token); @@ -423,6 +426,7 @@ public async Task Consume_connection_dispose_drains_buffered_messages() var consumed = await consumeTask; await using var check = new NatsConnection(new NatsOpts { Url = _server.Url }); + await check.ConnectRetryAsync(); var checkJs = new NatsJSContext(check); var info = (await checkJs.GetConsumerAsync(streamName, consumerName, cts.Token)).Info; @@ -834,7 +838,7 @@ public async Task Consume_503_counter_resets_on_success_test() { 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 cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); @@ -878,7 +882,7 @@ public async Task Consume_ephemeral_consumer_deleted_on_server_terminates_with_4 { 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 cts = new CancellationTokenSource(TimeSpan.FromSeconds(60)); @@ -986,7 +990,7 @@ public async Task Consume_slow_message_processing_does_not_prevent_consumer_dele { 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 cts = new CancellationTokenSource(TimeSpan.FromSeconds(60)); diff --git a/tests/NATS.Client.JetStream.Tests/ConsumerFetchTest.cs b/tests/NATS.Client.JetStream.Tests/ConsumerFetchTest.cs index c67d7f124..73220bb4a 100644 --- a/tests/NATS.Client.JetStream.Tests/ConsumerFetchTest.cs +++ b/tests/NATS.Client.JetStream.Tests/ConsumerFetchTest.cs @@ -1,6 +1,7 @@ using NATS.Client.Core.Tests; using NATS.Client.Core2.Tests; using NATS.Client.JetStream.Models; +using NATS.Client.TestUtilities2; namespace NATS.Client.JetStream.Tests; @@ -23,6 +24,7 @@ public async Task Fetch_test(NatsRequestReplyMode mode) { var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContext(nats); await js.CreateStreamAsync($"{prefix}s1", new[] { $"{prefix}s1.*" }, cts.Token); @@ -55,6 +57,7 @@ public async Task FetchNoWait_test(NatsRequestReplyMode mode) { var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContext(nats); await js.CreateStreamAsync($"{prefix}s1", new[] { $"{prefix}s1.*" }, cts.Token); @@ -93,6 +96,7 @@ public async Task Fetch_dispose_test(NatsRequestReplyMode mode) RequestReplyMode = mode, DrainSubscriptionsOnDispose = true, }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var cts = new CancellationTokenSource(TimeSpan.FromSeconds(60)); @@ -190,6 +194,7 @@ public async Task Fetch_connection_dispose_drains_buffered_messages() await using (var setup = new NatsConnection(new NatsOpts { Url = _server.Url })) { + await setup.ConnectRetryAsync(); var setupJs = new NatsJSContext(setup); var stream = await setupJs.CreateStreamAsync(new StreamConfig(streamName, [subject]), cts.Token); @@ -205,6 +210,7 @@ public async Task Fetch_connection_dispose_drains_buffered_messages() DrainSubscriptionsOnDispose = true, ConsumerDrainOnDisposeTimeout = TimeSpan.FromSeconds(30), }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var consumer = await js.GetConsumerAsync(streamName, consumerName, cts.Token); @@ -234,6 +240,7 @@ public async Task Fetch_connection_dispose_drains_buffered_messages() var consumed = await fetchTask; await using var check = new NatsConnection(new NatsOpts { Url = _server.Url }); + await check.ConnectRetryAsync(); var checkJs = new NatsJSContext(check); var info = (await checkJs.GetConsumerAsync(streamName, consumerName, cts.Token)).Info; diff --git a/tests/NATS.Client.JetStream.Tests/ConsumerNotificationTest.cs b/tests/NATS.Client.JetStream.Tests/ConsumerNotificationTest.cs index 01606aa58..bafdbdd49 100644 --- a/tests/NATS.Client.JetStream.Tests/ConsumerNotificationTest.cs +++ b/tests/NATS.Client.JetStream.Tests/ConsumerNotificationTest.cs @@ -3,6 +3,7 @@ using NATS.Client.Core2.Tests; using NATS.Client.JetStream.Models; using NATS.Client.TestUtilities; +using NATS.Client.TestUtilities2; using Synadia.Orbit.Testing.NatsServerProcessManager; namespace NATS.Client.JetStream.Tests; @@ -24,6 +25,7 @@ public async Task Non_terminal_errors_sent_as_notifications() { await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); @@ -97,6 +99,7 @@ public async Task Non_terminal_errors_sent_as_notifications() public async Task Exceeded_max_errors() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContext(nats); diff --git a/tests/NATS.Client.JetStream.Tests/ConsumerSetupTest.cs b/tests/NATS.Client.JetStream.Tests/ConsumerSetupTest.cs index 1132d27c8..17766c76c 100644 --- a/tests/NATS.Client.JetStream.Tests/ConsumerSetupTest.cs +++ b/tests/NATS.Client.JetStream.Tests/ConsumerSetupTest.cs @@ -3,6 +3,7 @@ using NATS.Client.Core2.Tests; using NATS.Client.JetStream.Models; using NATS.Client.TestUtilities; +using NATS.Client.TestUtilities2; using Synadia.Orbit.Testing.NatsServerProcessManager; namespace NATS.Client.JetStream.Tests; @@ -25,6 +26,7 @@ public ConsumerSetupTest(ITestOutputHelper output, NatsServerFixture server) public async Task Create_push_consumer(NatsRequestReplyMode mode) { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContext(nats); @@ -57,6 +59,7 @@ await js.CreateOrUpdateConsumerAsync( public async Task Create_paused_consumer() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContextFactory().CreateContext(nats); @@ -90,6 +93,7 @@ await js.CreateOrUpdateConsumerAsync( public async Task Consumer_config(NatsRequestReplyMode mode) { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContextFactory().CreateContext(nats); diff --git a/tests/NATS.Client.JetStream.Tests/CounterTest.cs b/tests/NATS.Client.JetStream.Tests/CounterTest.cs index 735176996..fc4c788e6 100644 --- a/tests/NATS.Client.JetStream.Tests/CounterTest.cs +++ b/tests/NATS.Client.JetStream.Tests/CounterTest.cs @@ -1,6 +1,7 @@ using NATS.Client.Core2.Tests; using NATS.Client.JetStream.Models; using NATS.Client.TestUtilities; +using NATS.Client.TestUtilities2; namespace NATS.Client.JetStream.Tests; @@ -20,6 +21,7 @@ public CounterTest(ITestOutputHelper output, NatsServerFixture server) public async Task Counter_functionality_test() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContext(nats); @@ -94,6 +96,7 @@ public async Task Counter_functionality_test() public async Task Counter_without_AllowMsgCounter_should_return_error() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContext(nats); diff --git a/tests/NATS.Client.JetStream.Tests/DirectGetTest.cs b/tests/NATS.Client.JetStream.Tests/DirectGetTest.cs index 96e6f7eaf..0e36c15b3 100644 --- a/tests/NATS.Client.JetStream.Tests/DirectGetTest.cs +++ b/tests/NATS.Client.JetStream.Tests/DirectGetTest.cs @@ -1,6 +1,7 @@ using NATS.Client.Core2.Tests; using NATS.Client.JetStream.Models; using NATS.Client.TestUtilities; +using NATS.Client.TestUtilities2; namespace NATS.Client.JetStream.Tests; @@ -22,6 +23,7 @@ public async Task Direct_get_returns_503_no_responders() // https://github.com/nats-io/nats-server/commit/ce309b79d99552996e18dce47dc04bdc730c0d84 // When we fail to deliver a message through a service import respond with no responders. await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContext(nats); @@ -42,6 +44,7 @@ await Assert.ThrowsAsync(async () => await s1.GetDire public async Task Direct_get_returns_no_reply() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContext(nats); diff --git a/tests/NATS.Client.JetStream.Tests/DoubleAckTest.cs b/tests/NATS.Client.JetStream.Tests/DoubleAckTest.cs index 0a5f4de23..e20359c37 100644 --- a/tests/NATS.Client.JetStream.Tests/DoubleAckTest.cs +++ b/tests/NATS.Client.JetStream.Tests/DoubleAckTest.cs @@ -1,4 +1,5 @@ using NATS.Client.Core2.Tests; +using NATS.Client.TestUtilities2; using Synadia.Orbit.Testing.NatsServerProcessManager; namespace NATS.Client.JetStream.Tests; @@ -20,6 +21,7 @@ public async Task Fetch_should_not_block_socket() { var cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContext(nats); diff --git a/tests/NATS.Client.JetStream.Tests/ErrorHandlerTest.cs b/tests/NATS.Client.JetStream.Tests/ErrorHandlerTest.cs index b9c350012..944fa3f47 100644 --- a/tests/NATS.Client.JetStream.Tests/ErrorHandlerTest.cs +++ b/tests/NATS.Client.JetStream.Tests/ErrorHandlerTest.cs @@ -24,6 +24,7 @@ public async Task Consumer_fetch_error_handling() { var proxy = new NatsProxy(_server.Port); await using var nats = new NatsConnection(new NatsOpts { Url = $"nats://127.0.0.1:{proxy.Port}", ConnectTimeout = TimeSpan.FromSeconds(10) }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContext(nats); @@ -97,6 +98,7 @@ public async Task Consumer_consume_handling() { var proxy = new NatsProxy(_server.Port); await using var nats = new NatsConnection(new NatsOpts { Url = $"nats://127.0.0.1:{proxy.Port}", ConnectTimeout = TimeSpan.FromSeconds(10) }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContext(nats); @@ -187,6 +189,7 @@ public async Task Ordered_consumer_fetch_error_handling() { var proxy = new NatsProxy(_server.Port); await using var nats = new NatsConnection(new NatsOpts { Url = $"nats://127.0.0.1:{proxy.Port}", ConnectTimeout = TimeSpan.FromSeconds(10) }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContext(nats); @@ -357,6 +360,7 @@ public async Task Exception_propagation_handling() { var proxy = new NatsProxy(_server.Port); await using var nats = new NatsConnection(new NatsOpts { Url = $"nats://127.0.0.1:{proxy.Port}", ConnectTimeout = TimeSpan.FromSeconds(10) }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContext(nats); diff --git a/tests/NATS.Client.JetStream.Tests/JetStreamApiSerializerTest.cs b/tests/NATS.Client.JetStream.Tests/JetStreamApiSerializerTest.cs index 58c0cdb2b..dcdd23f19 100644 --- a/tests/NATS.Client.JetStream.Tests/JetStreamApiSerializerTest.cs +++ b/tests/NATS.Client.JetStream.Tests/JetStreamApiSerializerTest.cs @@ -4,6 +4,7 @@ using NATS.Client.Core2.Tests; using NATS.Client.JetStream.Internal; using NATS.Client.JetStream.Models; +using NATS.Client.TestUtilities2; using JsonSerializer = System.Text.Json.JsonSerializer; namespace NATS.Client.JetStream.Tests; @@ -24,6 +25,7 @@ public JetStreamApiSerializerTest(ITestOutputHelper output, NatsServerFixture se public async Task Should_respect_buffers_lifecycle() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContext(nats); var apiSubject = $"{prefix}.js.fake.api"; diff --git a/tests/NATS.Client.JetStream.Tests/ListTests.cs b/tests/NATS.Client.JetStream.Tests/ListTests.cs index 6dfa36b8c..62ed37ded 100644 --- a/tests/NATS.Client.JetStream.Tests/ListTests.cs +++ b/tests/NATS.Client.JetStream.Tests/ListTests.cs @@ -1,6 +1,7 @@ using NATS.Client.Core.Tests; using NATS.Client.Core2.Tests; using NATS.Client.JetStream.Models; +using NATS.Client.TestUtilities2; using Synadia.Orbit.Testing.NatsServerProcessManager; namespace NATS.Client.JetStream.Tests; @@ -21,6 +22,7 @@ public ListTests(ITestOutputHelper output, NatsServerFixture server) public async Task List_streams() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url, RequestTimeout = TimeSpan.FromSeconds(5) }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId() + "-"; _output.WriteLine($"prefix: {prefix}"); @@ -95,6 +97,7 @@ public async Task List_streams() public async Task List_consumers() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url, RequestTimeout = TimeSpan.FromSeconds(5) }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId() + "-"; _output.WriteLine($"prefix: {prefix}"); var js = new NatsJSContext(nats); diff --git a/tests/NATS.Client.JetStream.Tests/NatsJsContextFactoryTest.cs b/tests/NATS.Client.JetStream.Tests/NatsJsContextFactoryTest.cs index 2f353dd71..f86ba1c9a 100644 --- a/tests/NATS.Client.JetStream.Tests/NatsJsContextFactoryTest.cs +++ b/tests/NATS.Client.JetStream.Tests/NatsJsContextFactoryTest.cs @@ -2,6 +2,7 @@ using System.Text; using System.Threading.Channels; using NATS.Client.Core.Tests; +using NATS.Client.TestUtilities2; using Synadia.Orbit.Testing.NatsServerProcessManager; namespace NATS.Client.JetStream.Tests; @@ -18,6 +19,7 @@ public async Task Create_Context_Test() // Arrange await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, RequestTimeout = TimeSpan.FromSeconds(10) }); + await nats.ConnectRetryAsync(); var factory = new NatsJSContextFactory(); // Act @@ -33,6 +35,7 @@ public async Task Create_Context_WithOpts_Test() // Arrange await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, RequestTimeout = TimeSpan.FromSeconds(10) }); + await nats.ConnectRetryAsync(); var factory = new NatsJSContextFactory(); var opts = new NatsJSOpts(nats.Opts); diff --git a/tests/NATS.Client.JetStream.Tests/OrderedConsumerTest.cs b/tests/NATS.Client.JetStream.Tests/OrderedConsumerTest.cs index dd240f95b..1dd5cf636 100644 --- a/tests/NATS.Client.JetStream.Tests/OrderedConsumerTest.cs +++ b/tests/NATS.Client.JetStream.Tests/OrderedConsumerTest.cs @@ -510,6 +510,7 @@ public async Task OrderedConsume_connection_dispose_drains_buffered_messages() await using (var setup = new NatsConnection(new NatsOpts { Url = _server.Url })) { + await setup.ConnectRetryAsync(); var setupJs = new NatsJSContext(setup); await setupJs.CreateStreamAsync(new StreamConfig(streamName, [subject]), cts.Token); @@ -523,6 +524,7 @@ public async Task OrderedConsume_connection_dispose_drains_buffered_messages() DrainSubscriptionsOnDispose = true, ConsumerDrainOnDisposeTimeout = TimeSpan.FromSeconds(30), }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var stream = await js.GetStreamAsync(streamName, cancellationToken: cts.Token); var consumer = await stream.CreateOrderedConsumerAsync(cancellationToken: cts.Token); diff --git a/tests/NATS.Client.JetStream.Tests/PinnedClientTest.cs b/tests/NATS.Client.JetStream.Tests/PinnedClientTest.cs index 73c1c2d89..8e8e447e1 100644 --- a/tests/NATS.Client.JetStream.Tests/PinnedClientTest.cs +++ b/tests/NATS.Client.JetStream.Tests/PinnedClientTest.cs @@ -27,6 +27,7 @@ public PinnedClientTest(ITestOutputHelper output, NatsServerFixture server) public async Task Pinned_client_basic_flow_with_consume() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var prefix = _server.GetNextId(); @@ -95,6 +96,7 @@ public async Task Pinned_client_basic_flow_with_consume() public async Task Unpin_allows_other_consumer_to_receive() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var prefix = _server.GetNextId(); @@ -257,6 +259,7 @@ await SendHMsgAsync( using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); await using var nats = new NatsConnection(new NatsOpts { Url = ms.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var js = nats.CreateJetStreamContext(); var consumer = (NatsJSConsumer)await js.GetConsumerAsync("x", "x", cts.Token); @@ -376,6 +379,7 @@ public async Task Pinned_client_not_allowed_with_next() public async Task Invalid_priority_group_returns_error() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var prefix = _server.GetNextId(); @@ -410,6 +414,7 @@ public async Task Invalid_priority_group_returns_error() public async Task Context_unpin_consumer_async() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var prefix = _server.GetNextId(); @@ -474,6 +479,7 @@ public async Task Context_unpin_consumer_async() public async Task Consumer_info_shows_priority_groups_state() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var prefix = _server.GetNextId(); @@ -538,6 +544,7 @@ public async Task Pin_id_from_headers_should_use_last_value_when_multiple_header using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); await using var nats = new NatsConnection(new NatsOpts { Url = ms.Url, ConnectTimeout = TimeSpan.FromSeconds(30) }); + await nats.ConnectRetryAsync(); var js = nats.CreateJetStreamContext(); var consumer = (NatsJSConsumer)await js.GetConsumerAsync("x", "x", cts.Token); var headers = new NatsHeaders @@ -578,6 +585,7 @@ public async Task Queued_consume_pull_request_should_use_latest_pin_id_when_sent using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(30)); await using var nats = new NatsConnection(new NatsOpts { Url = ms.Url }); + await nats.ConnectRetryAsync(); var js = nats.CreateJetStreamContext(); var consumer = (NatsJSConsumer)await js.GetConsumerAsync("x", "x", cts.Token); diff --git a/tests/NATS.Client.JetStream.Tests/PriorityGroupTest.cs b/tests/NATS.Client.JetStream.Tests/PriorityGroupTest.cs index f07f3cd69..d436cd01f 100644 --- a/tests/NATS.Client.JetStream.Tests/PriorityGroupTest.cs +++ b/tests/NATS.Client.JetStream.Tests/PriorityGroupTest.cs @@ -2,6 +2,7 @@ using NATS.Client.Core2.Tests; using NATS.Client.JetStream.Models; using NATS.Client.TestUtilities; +using NATS.Client.TestUtilities2; namespace NATS.Client.JetStream.Tests; @@ -21,6 +22,7 @@ public PriorityGroupTest(ITestOutputHelper output, NatsServerFixture server) public async Task Next_from_overflow_group() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var prefix = _server.GetNextId(); @@ -64,6 +66,7 @@ public async Task Next_from_overflow_group() public async Task Fetch_from_overflow_group() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var prefix = _server.GetNextId(); @@ -110,6 +113,7 @@ public async Task Fetch_from_overflow_group() public async Task Consume_from_overflow_group() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var prefix = _server.GetNextId(); @@ -156,6 +160,7 @@ public async Task Consume_from_overflow_group() public async Task Fetch_from_prioritized_group_with_priority() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var prefix = _server.GetNextId(); @@ -218,6 +223,7 @@ public async Task Fetch_from_prioritized_group_with_priority() public async Task Consumer_with_prioritized_policy_validation() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var prefix = _server.GetNextId(); diff --git a/tests/NATS.Client.JetStream.Tests/PublishConcurrentTests.cs b/tests/NATS.Client.JetStream.Tests/PublishConcurrentTests.cs index cc2dfe86c..757b1ac5b 100644 --- a/tests/NATS.Client.JetStream.Tests/PublishConcurrentTests.cs +++ b/tests/NATS.Client.JetStream.Tests/PublishConcurrentTests.cs @@ -1,6 +1,7 @@ using System.Diagnostics; using NATS.Client.Core.Tests; using NATS.Client.Core2.Tests; +using NATS.Client.TestUtilities2; using Synadia.Orbit.Testing.NatsServerProcessManager; namespace NATS.Client.JetStream.Tests; @@ -21,6 +22,7 @@ public PublishConcurrentTests(ITestOutputHelper output, NatsServerFixture server public async Task Publish_concurrently() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContext(nats); diff --git a/tests/NATS.Client.JetStream.Tests/PublishRetryTest.cs b/tests/NATS.Client.JetStream.Tests/PublishRetryTest.cs index 21ce4d10a..506f89063 100644 --- a/tests/NATS.Client.JetStream.Tests/PublishRetryTest.cs +++ b/tests/NATS.Client.JetStream.Tests/PublishRetryTest.cs @@ -24,6 +24,7 @@ public async Task Publish_without_telemetry(NatsRequestReplyMode mode) { var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); // Without telemetry diff --git a/tests/NATS.Client.JetStream.Tests/PublishTest.cs b/tests/NATS.Client.JetStream.Tests/PublishTest.cs index 865b6681c..fe727db24 100644 --- a/tests/NATS.Client.JetStream.Tests/PublishTest.cs +++ b/tests/NATS.Client.JetStream.Tests/PublishTest.cs @@ -343,6 +343,7 @@ public async Task Publish_retry_on_503(NatsRequestReplyMode mode) LoggerFactory = logger, RequestReplyMode = mode, }); + await nats.ConnectRetryAsync(); var cts = new CancellationTokenSource(TimeSpan.FromSeconds(60)); var js = new NatsJSContext(nats); @@ -379,6 +380,7 @@ public async Task Publish_no_responders(NatsRequestReplyMode mode) { await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var js = nats.CreateJetStreamContext(); var result = await js.TryPublishAsync("foo", 1); Assert.IsType(result.Error); diff --git a/tests/NATS.Client.JetStream.Tests/SlowConsumerTest.cs b/tests/NATS.Client.JetStream.Tests/SlowConsumerTest.cs index cea079d80..df218311a 100644 --- a/tests/NATS.Client.JetStream.Tests/SlowConsumerTest.cs +++ b/tests/NATS.Client.JetStream.Tests/SlowConsumerTest.cs @@ -1,6 +1,7 @@ using NATS.Client.Core.Tests; using NATS.Client.Core2.Tests; using NATS.Client.JetStream.Models; +using NATS.Client.TestUtilities2; namespace NATS.Client.JetStream.Tests; @@ -28,7 +29,7 @@ public async Task JetStream_consume_slow_consumer_should_not_block_connection() Url = _server.Url, SubPendingChannelCapacity = 10, // Small capacity to trigger drops quickly }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContext(nats); @@ -208,7 +209,7 @@ public async Task JetStream_fetch_slow_consumer_should_not_block_connection() Url = _server.Url, SubPendingChannelCapacity = 10, }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContext(nats); diff --git a/tests/NATS.Client.JetStream.Tests/StreamDomainAdjustmentTest.cs b/tests/NATS.Client.JetStream.Tests/StreamDomainAdjustmentTest.cs index 5da808331..d1e33e30f 100644 --- a/tests/NATS.Client.JetStream.Tests/StreamDomainAdjustmentTest.cs +++ b/tests/NATS.Client.JetStream.Tests/StreamDomainAdjustmentTest.cs @@ -38,6 +38,7 @@ public async Task Stream_operations_should_convert_domain_to_external_api() }); await using var nats = new NatsConnection(new NatsOpts { Url = ms.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); // Configure a stream with Source that has Domain set diff --git a/tests/NATS.Client.KeyValueStore.Tests/GetKeysTest.cs b/tests/NATS.Client.KeyValueStore.Tests/GetKeysTest.cs index d450184cc..687055a6c 100644 --- a/tests/NATS.Client.KeyValueStore.Tests/GetKeysTest.cs +++ b/tests/NATS.Client.KeyValueStore.Tests/GetKeysTest.cs @@ -1,5 +1,6 @@ using NATS.Client.Core.Tests; using NATS.Client.TestUtilities; +using NATS.Client.TestUtilities2; using Synadia.Orbit.Testing.NatsServerProcessManager; namespace NATS.Client.KeyValueStore.Tests; @@ -17,6 +18,7 @@ public async Task Get_keys_should_not_hang_when_there_are_deleted_keys() await using var server = await NatsServerProcess.StartAsync(); await using var nats1 = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats1.ConnectRetryAsync(); var js1 = new NatsJSContext(nats1); var kv1 = new NatsKVContext(js1); var store1 = await kv1.CreateStoreAsync(config, cancellationToken: cancellationToken); @@ -59,6 +61,7 @@ public async Task Get_keys_when_empty() await using var server = await NatsServerProcess.StartAsync(); await using var nats1 = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats1.ConnectRetryAsync(); var js1 = new NatsJSContext(nats1); var kv1 = new NatsKVContext(js1); var store1 = await kv1.CreateStoreAsync(config, cancellationToken: cancellationToken); @@ -94,6 +97,7 @@ public async Task Get_filtered_keys() await using var server = await NatsServerProcess.StartAsync(); await using var nats1 = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats1.ConnectRetryAsync(); var js1 = new NatsJSContext(nats1); var kv1 = new NatsKVContext(js1); var store1 = await kv1.CreateStoreAsync(config, cancellationToken: cancellationToken); diff --git a/tests/NATS.Client.KeyValueStore.Tests/KeyValueContextTest.cs b/tests/NATS.Client.KeyValueStore.Tests/KeyValueContextTest.cs index edf420b7b..bd836eaa4 100644 --- a/tests/NATS.Client.KeyValueStore.Tests/KeyValueContextTest.cs +++ b/tests/NATS.Client.KeyValueStore.Tests/KeyValueContextTest.cs @@ -1,4 +1,5 @@ using NATS.Client.Core.Tests; +using NATS.Client.TestUtilities2; using Synadia.Orbit.Testing.NatsServerProcessManager; namespace NATS.Client.KeyValueStore.Tests; @@ -14,6 +15,7 @@ public async Task Create_store_test() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -33,6 +35,7 @@ public async Task Update_store_test() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -58,6 +61,7 @@ public async Task Create_store_via_create_or_update_store_test() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -81,6 +85,7 @@ public async Task Update_store_via_create_or_update_store_test() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -106,6 +111,7 @@ public async Task Duplicate_window_capped_by_max_age_on_create(int maxAgeSeconds await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -125,6 +131,7 @@ public async Task Duplicate_window_recapped_on_update() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -148,6 +155,7 @@ public async Task Delete_store_test() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -170,6 +178,7 @@ public async Task Get_bucket_names_test() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -199,6 +208,7 @@ public async Task Get_statuses_test() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); diff --git a/tests/NATS.Client.KeyValueStore.Tests/KeyValueStoreTest.cs b/tests/NATS.Client.KeyValueStore.Tests/KeyValueStoreTest.cs index d8eddf3b8..7d3e4fed1 100644 --- a/tests/NATS.Client.KeyValueStore.Tests/KeyValueStoreTest.cs +++ b/tests/NATS.Client.KeyValueStore.Tests/KeyValueStoreTest.cs @@ -2,6 +2,7 @@ using NATS.Client.Core.Tests; using NATS.Client.JetStream.Models; using NATS.Client.TestUtilities; +using NATS.Client.TestUtilities2; using Synadia.Orbit.Testing.NatsServerProcessManager; namespace NATS.Client.KeyValueStore.Tests; @@ -19,6 +20,7 @@ public async Task Simple_create_put_get_test(NatsRequestReplyMode mode) { await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -40,6 +42,7 @@ public async Task Handle_non_direct_gets(NatsRequestReplyMode mode) { await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -83,6 +86,7 @@ public async Task Get_keys(NatsRequestReplyMode mode) await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -122,6 +126,7 @@ public async Task Get_key_revisions(NatsRequestReplyMode mode) await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -182,6 +187,7 @@ public async Task Delete_and_purge(NatsRequestReplyMode mode) await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -310,6 +316,7 @@ public async Task Purge_deletes(NatsRequestReplyMode mode) await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -394,6 +401,7 @@ public async Task Update_with_revisions(NatsRequestReplyMode mode) await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -434,6 +442,7 @@ public async Task Create(NatsRequestReplyMode mode) await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -486,6 +495,7 @@ public async Task TestMessageTTLApiNotSupportedupport() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -505,6 +515,7 @@ public async Task TestMessageTTL(NatsRequestReplyMode mode) await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -567,6 +578,7 @@ public async Task TestTTLMessageWhenTTLDisabledOnStream(NatsRequestReplyMode mod await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -586,6 +598,7 @@ public async Task SetsSubjectDeleteMarkerTTL(NatsRequestReplyMode mode) await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -602,6 +615,7 @@ public async Task SubjectDeleteMarkerTTL_enabled_removals_should_be_interpreted_ { await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -656,6 +670,7 @@ public async Task TestMessageNeverExpire(NatsRequestReplyMode mode) await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -703,6 +718,7 @@ public async Task History(NatsRequestReplyMode mode) await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -749,6 +765,7 @@ public async Task Status(NatsRequestReplyMode mode) await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -787,6 +804,7 @@ public async Task Compressed_storage(NatsRequestReplyMode mode) { await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -818,6 +836,7 @@ public async Task Validate_keys(NatsRequestReplyMode mode) { await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -883,6 +902,7 @@ public async Task TestDirectMessageRepublishedSubject(NatsRequestReplyMode mode) await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -916,6 +936,7 @@ public async Task Test_CombinedSources(NatsRequestReplyMode mode) { await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -965,6 +986,7 @@ public async Task Try_Create(NatsRequestReplyMode mode) { await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -1001,6 +1023,7 @@ public async Task Try_Delete(NatsRequestReplyMode mode) { await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); @@ -1028,6 +1051,7 @@ public async Task Try_Update(NatsRequestReplyMode mode) { await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, RequestReplyMode = mode }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var kv = new NatsKVContext(js); diff --git a/tests/NATS.Client.KeyValueStore.Tests/NatsKVContextFactoryTest.cs b/tests/NATS.Client.KeyValueStore.Tests/NatsKVContextFactoryTest.cs index 65fe6f7bf..a6683d0ef 100644 --- a/tests/NATS.Client.KeyValueStore.Tests/NatsKVContextFactoryTest.cs +++ b/tests/NATS.Client.KeyValueStore.Tests/NatsKVContextFactoryTest.cs @@ -1,5 +1,6 @@ using NATS.Client.Core.Tests; using NATS.Client.JetStream.Models; +using NATS.Client.TestUtilities2; using Synadia.Orbit.Testing.NatsServerProcessManager; namespace NATS.Client.KeyValueStore.Tests; @@ -16,6 +17,7 @@ public async Task Create_Context_Test() // Arrange await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, RequestTimeout = TimeSpan.FromSeconds(10) }); + await nats.ConnectRetryAsync(); var jsFactory = new NatsJSContextFactory(); var jsContext = jsFactory.CreateContext(nats); var factory = new NatsKVContextFactory(); diff --git a/tests/NATS.Client.KeyValueStore.Tests/NatsKVWatcherTest.cs b/tests/NATS.Client.KeyValueStore.Tests/NatsKVWatcherTest.cs index 9d9301fb1..3d4cbbc9f 100644 --- a/tests/NATS.Client.KeyValueStore.Tests/NatsKVWatcherTest.cs +++ b/tests/NATS.Client.KeyValueStore.Tests/NatsKVWatcherTest.cs @@ -40,6 +40,7 @@ public async Task Watcher_reconnect_with_history() var proxy = new NatsProxy(_server.Port); await using var nats2 = new NatsConnection(new NatsOpts { Url = $"nats://127.0.0.1:{proxy.Port}" }); + await nats2.ConnectRetryAsync(); await nats1.ConnectRetryAsync(); var js2 = new NatsJSContext(nats2); var kv2 = new NatsKVContext(js2); @@ -118,6 +119,7 @@ public async Task Watch_all() var cancellationToken = cts.Token; await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContext(nats); @@ -162,6 +164,7 @@ public async Task Watch_subset() var cancellationToken = cts.Token; await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContext(nats); @@ -328,6 +331,7 @@ await Retry.Until( public async Task Watch_push_consumer_flow_control() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContext(nats); @@ -368,6 +372,7 @@ public async Task Watch_empty_bucket_for_end_of_data() var cancellationToken = cts.Token; await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContext(nats); @@ -409,6 +414,7 @@ public async Task Watch_empty_bucket_for_end_of_data() public async Task Serialization_errors() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContext(nats); @@ -435,6 +441,7 @@ public async Task Serialization_errors() public async Task Watch_with_empty_filter() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContext(nats); @@ -456,6 +463,7 @@ await Assert.ThrowsAsync(async () => public async Task Watch_with_multiple_filter_on_old_server() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContext(nats); @@ -485,6 +493,7 @@ public async Task Watch_with_multiple_filter_on_old_server() public async Task Watch_resume_at_revision() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var bucket = $"{prefix}Watch_resume_at_revision"; @@ -568,6 +577,7 @@ public async Task Watch_resume_at_revision() public async Task Validate_watch_options() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var bucket = prefix + nameof(Validate_watch_options); @@ -643,6 +653,7 @@ await Assert.ThrowsAsync(async () => public async Task ReadAfterDelete() { await using var nats = new NatsConnection(new NatsOpts { Url = _server.Url }); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContext(nats); diff --git a/tests/NATS.Client.KeyValueStore.Tests/SlowConsumerTest.cs b/tests/NATS.Client.KeyValueStore.Tests/SlowConsumerTest.cs index 23ee6aa29..76a44c576 100644 --- a/tests/NATS.Client.KeyValueStore.Tests/SlowConsumerTest.cs +++ b/tests/NATS.Client.KeyValueStore.Tests/SlowConsumerTest.cs @@ -1,5 +1,6 @@ using NATS.Client.Core.Tests; using NATS.Client.Core2.Tests; +using NATS.Client.TestUtilities2; namespace NATS.Client.KeyValueStore.Tests; @@ -27,7 +28,7 @@ public async Task KV_watch_slow_consumer_should_not_block_connection() Url = _server.Url, SubPendingChannelCapacity = 10, // Small capacity to trigger drops quickly }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); var prefix = _server.GetNextId(); var js = new NatsJSContext(nats); diff --git a/tests/NATS.Client.ObjectStore.Tests/CompatTest.cs b/tests/NATS.Client.ObjectStore.Tests/CompatTest.cs index 3ad3f6198..e791e1cf0 100644 --- a/tests/NATS.Client.ObjectStore.Tests/CompatTest.cs +++ b/tests/NATS.Client.ObjectStore.Tests/CompatTest.cs @@ -2,6 +2,7 @@ using System.Text.Json.Nodes; using NATS.Client.Core.Tests; using NATS.Client.ObjectStore.Models; +using NATS.Client.TestUtilities2; using Synadia.Orbit.Testing.NatsServerProcessManager; namespace NATS.Client.ObjectStore.Tests; @@ -13,6 +14,7 @@ public async Task Headers_serialization() { await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var ob = new NatsObjContext(js); diff --git a/tests/NATS.Client.ObjectStore.Tests/ObjectStoreTest.cs b/tests/NATS.Client.ObjectStore.Tests/ObjectStoreTest.cs index 2fd5c31cf..33a835284 100644 --- a/tests/NATS.Client.ObjectStore.Tests/ObjectStoreTest.cs +++ b/tests/NATS.Client.ObjectStore.Tests/ObjectStoreTest.cs @@ -8,6 +8,7 @@ using NATS.Client.ObjectStore.Models; using NATS.Client.Serializers.Json; using NATS.Client.TestUtilities; +using NATS.Client.TestUtilities2; using Synadia.Orbit.Testing.NatsServerProcessManager; namespace NATS.Client.ObjectStore.Tests; @@ -26,6 +27,7 @@ public async Task Create_delete_object_store() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var ob = new NatsObjContext(js); @@ -53,6 +55,7 @@ public async Task Put_chunks() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var ob = new NatsObjContext(js); @@ -135,6 +138,7 @@ public async Task Get_chunks() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var ob = new NatsObjContext(js); @@ -189,6 +193,7 @@ public async Task Delete_object() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var ob = new NatsObjContext(js); @@ -227,6 +232,7 @@ public async Task Put_and_get_large_file() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var obj = new NatsObjContext(js); @@ -258,6 +264,7 @@ public async Task Add_link() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var obj = new NatsObjContext(js); @@ -305,6 +312,7 @@ public async Task Seal_and_get_status() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var obj = new NatsObjContext(js); @@ -344,6 +352,7 @@ public async Task List() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var obj = new NatsObjContext(js); @@ -402,6 +411,7 @@ public async Task List_empty_store_for_end_of_data() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var obj = new NatsObjContext(js); @@ -441,6 +451,7 @@ public async Task Compressed_storage() { await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var obj = new NatsObjContext(js); @@ -473,6 +484,7 @@ public async Task Put_get_serialization_when_default_serializer_is_not_used() Url = server.Url, SerializerRegistry = NatsJsonSerializerRegistry.Default, }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var ob = new NatsObjContext(js); @@ -508,6 +520,7 @@ public async Task Put_with_activity() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var obj = new NatsObjContext(js); @@ -540,6 +553,7 @@ public async Task Put_multiple_times_with_activity() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var obj = new NatsObjContext(js); @@ -564,6 +578,7 @@ public async Task Rename_object_should_perge_old_named_meta() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var ob = new NatsObjContext(js); @@ -632,6 +647,7 @@ public async Task Metadata_field_types_match_spec() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var obj = new NatsObjContext(js); @@ -673,6 +689,7 @@ public async Task Get_does_not_leak_connection_opened_handlers() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var ob = new NatsObjContext(js); diff --git a/tests/NATS.Client.ObjectStore.Tests/SlowConsumerTest.cs b/tests/NATS.Client.ObjectStore.Tests/SlowConsumerTest.cs index a1d3bc271..f9d2dc92d 100644 --- a/tests/NATS.Client.ObjectStore.Tests/SlowConsumerTest.cs +++ b/tests/NATS.Client.ObjectStore.Tests/SlowConsumerTest.cs @@ -1,5 +1,6 @@ using NATS.Client.Core.Tests; using NATS.Client.ObjectStore.Models; +using NATS.Client.TestUtilities2; using Synadia.Orbit.Testing.NatsServerProcessManager; namespace NATS.Client.ObjectStore.Tests; @@ -24,7 +25,7 @@ public async Task ObjectStore_get_slow_consumer_should_not_block_connection() Url = server.Url, SubPendingChannelCapacity = 10, }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var ob = new NatsObjContext(js); diff --git a/tests/NATS.Client.ObjectStore.Tests/WatcherTest.cs b/tests/NATS.Client.ObjectStore.Tests/WatcherTest.cs index ea7bdf726..b0962125f 100644 --- a/tests/NATS.Client.ObjectStore.Tests/WatcherTest.cs +++ b/tests/NATS.Client.ObjectStore.Tests/WatcherTest.cs @@ -1,4 +1,5 @@ using NATS.Client.Core.Tests; +using NATS.Client.TestUtilities2; using Synadia.Orbit.Testing.NatsServerProcessManager; namespace NATS.Client.ObjectStore.Tests; @@ -10,6 +11,7 @@ public async Task Watcher_test() { await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var js = new NatsJSContext(nats); var ob = new NatsObjContext(js); diff --git a/tests/NATS.Client.Services.Tests/ServicesSerializationTest.cs b/tests/NATS.Client.Services.Tests/ServicesSerializationTest.cs index d05d8afd7..52a853905 100644 --- a/tests/NATS.Client.Services.Tests/ServicesSerializationTest.cs +++ b/tests/NATS.Client.Services.Tests/ServicesSerializationTest.cs @@ -5,6 +5,7 @@ using NATS.Client.Serializers.Json; using NATS.Client.Services.Internal; using NATS.Client.Services.Models; +using NATS.Client.TestUtilities2; using Synadia.Orbit.Testing.NatsServerProcessManager; namespace NATS.Client.Services.Tests; @@ -22,6 +23,7 @@ public async Task Service_info_and_stat_request_serialization() // Set serializer registry to use anything but a raw bytes (NatsMemory in this case) serializer await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, SerializerRegistry = NatsJsonSerializerRegistry.Default }); + await nats.ConnectRetryAsync(); var svc = new NatsSvcContext(nats); @@ -49,6 +51,7 @@ public async Task Service_message_serialization() { await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url, SerializerRegistry = NatsJsonSerializerRegistry.Default }); + await nats.ConnectRetryAsync(); var svc = new NatsSvcContext(nats); diff --git a/tests/NATS.Client.Services.Tests/ServicesTests.cs b/tests/NATS.Client.Services.Tests/ServicesTests.cs index 9eeade7a6..2fb930f8a 100644 --- a/tests/NATS.Client.Services.Tests/ServicesTests.cs +++ b/tests/NATS.Client.Services.Tests/ServicesTests.cs @@ -22,6 +22,7 @@ public async Task Add_service_listeners_ping_info_and_stats() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var svc = new NatsSvcContext(nats); await using var s1 = await svc.AddServiceAsync("s1", "1.0.0", cancellationToken: cancellationToken); @@ -57,6 +58,7 @@ public async Task Add_end_point() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var svc = new NatsSvcContext(nats); await using var s1 = await svc.AddServiceAsync("s1", "1.0.0", cancellationToken: cancellationToken); @@ -139,6 +141,7 @@ public async Task Add_groups_metadata_and_stats() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var svc = new NatsSvcContext(nats); await using var s1 = await svc.AddServiceAsync("s1", "1.0.0", cancellationToken: cancellationToken); @@ -240,6 +243,7 @@ public async Task Add_multiple_service_listeners_ping_info_and_stats() { await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var svc = new NatsSvcContext(nats); var cts = new CancellationTokenSource(TimeSpan.FromSeconds(10)); @@ -276,6 +280,7 @@ public async Task Pass_headers_to_request_and_in_response() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var svc = new NatsSvcContext(nats); await using var s1 = await svc.AddServiceAsync("s1", "1.0.0", cancellationToken: cancellationToken); @@ -334,6 +339,7 @@ public async Task Service_started_time() await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var svc = new NatsSvcContext(nats); await using var s1 = await svc.AddServiceAsync("s1", "1.0.0", cancellationToken: cancellationToken); @@ -358,6 +364,7 @@ public async Task Service_ids_unique() { await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var svc = new NatsSvcContext(nats); var cts = new CancellationTokenSource(TimeSpan.FromSeconds(60)); @@ -392,6 +399,7 @@ public async Task Service_without_queue_group() await using var server = await NatsServerProcess.StartAsync(); var proxy = new NatsProxy(server.Port); await using var nats = new NatsConnection(new NatsOpts { Url = $"127.0.0.1:{proxy.Port}" }); + await nats.ConnectRetryAsync(); var svc = new NatsSvcContext(nats); var cts = new CancellationTokenSource(TimeSpan.FromSeconds(60)); @@ -532,6 +540,7 @@ public async Task AddEndpoint_SubjectWithWhitespace_ThrowsWhenValidationEnabled( await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var svc = new NatsSvcContext(nats); await using var s1 = await svc.AddServiceAsync("s1", "1.0.0", cancellationToken: cancellationToken); @@ -557,6 +566,7 @@ public async Task AddEndpoint_QueueGroupWithWhitespace_ThrowsWhenValidationEnabl await using var server = await NatsServerProcess.StartAsync(); await using var nats = new NatsConnection(new NatsOpts { Url = server.Url }); + await nats.ConnectRetryAsync(); var svc = new NatsSvcContext(nats); await using var s1 = await svc.AddServiceAsync( diff --git a/tests/NATS.Client.Services.Tests/SlowConsumerTest.cs b/tests/NATS.Client.Services.Tests/SlowConsumerTest.cs index 74d2dc5e5..41d4c7263 100644 --- a/tests/NATS.Client.Services.Tests/SlowConsumerTest.cs +++ b/tests/NATS.Client.Services.Tests/SlowConsumerTest.cs @@ -1,4 +1,5 @@ using NATS.Client.Core.Tests; +using NATS.Client.TestUtilities2; using Synadia.Orbit.Testing.NatsServerProcessManager; namespace NATS.Client.Services.Tests; @@ -23,7 +24,7 @@ public async Task Service_slow_handler_should_not_block_connection() Url = server.Url, SubPendingChannelCapacity = 10, }); - await nats.ConnectAsync(); + await nats.ConnectRetryAsync(); var svc = new NatsSvcContext(nats);