From ee4710cdeeda50db75746c79228242580ea83872 Mon Sep 17 00:00:00 2001 From: pzajaczkowski Date: Wed, 15 Jan 2025 19:03:38 +0100 Subject: [PATCH 1/6] Add Lame Duck Mode event handler --- src/NATS.Client.Core/INatsConnection.cs | 5 ++++ src/NATS.Client.Core/NatsConnection.cs | 24 +++++++++++++++++-- .../NatsJsContextFactoryTest.cs | 2 ++ 3 files changed, 29 insertions(+), 2 deletions(-) diff --git a/src/NATS.Client.Core/INatsConnection.cs b/src/NATS.Client.Core/INatsConnection.cs index 403eecde2..d893ca4ac 100644 --- a/src/NATS.Client.Core/INatsConnection.cs +++ b/src/NATS.Client.Core/INatsConnection.cs @@ -25,6 +25,11 @@ public interface INatsConnection : INatsClient /// event AsyncEventHandler? MessageDropped; + /// + /// Event that is raised when server goes into Lame Duck Mode. + /// + public event AsyncEventHandler? LameDuckModeActivated; + /// /// Server information received from the NATS server. /// diff --git a/src/NATS.Client.Core/NatsConnection.cs b/src/NATS.Client.Core/NatsConnection.cs index 73ef8c9d9..71ec05f2a 100644 --- a/src/NATS.Client.Core/NatsConnection.cs +++ b/src/NATS.Client.Core/NatsConnection.cs @@ -24,14 +24,13 @@ internal enum NatsEvent ConnectionDisconnected, ReconnectFailed, MessageDropped, + LameDuckModeActivated, } public partial class NatsConnection : INatsConnection { #pragma warning disable SA1401 internal readonly ConnectionStatsCounter Counter; // allow to call from external sources - internal volatile ServerInfo? WritableServerInfo; - #pragma warning restore SA1401 private readonly object _gate = new object(); private readonly ILogger _logger; @@ -44,6 +43,7 @@ public partial class NatsConnection : INatsConnection private readonly ClientOpts _clientOpts; private readonly SubscriptionManager _subscriptionManager; + private ServerInfo? _writableServerInfo; private int _pongCount; private int _connectionState; private int _isDisposed; @@ -109,6 +109,8 @@ public NatsConnection(NatsOpts opts) public event AsyncEventHandler? MessageDropped; + public event AsyncEventHandler? LameDuckModeActivated; + public INatsConnection Connection => this; public NatsOpts Opts { get; } @@ -134,6 +136,21 @@ private set public Func>? OnSocketAvailableAsync { get; set; } + internal ServerInfo? WritableServerInfo + { + get => Interlocked.CompareExchange(ref _writableServerInfo, null, null); + set + { + var current = Interlocked.CompareExchange(ref _writableServerInfo, null, null); + if (current?.LameDuckMode == false && value?.LameDuckMode == true) + { + _eventChannel.Writer.TryWrite((NatsEvent.LameDuckModeActivated, new NatsEventArgs(string.Empty))); + } + + Interlocked.Exchange(ref _writableServerInfo, value); + } + } + internal bool IsDisposed { get => Interlocked.CompareExchange(ref _isDisposed, 0, 0) == 1; @@ -762,6 +779,9 @@ private async Task PublishEventsAsync() case NatsEvent.MessageDropped when MessageDropped != null && args is NatsMessageDroppedEventArgs error: await MessageDropped.InvokeAsync(this, error).ConfigureAwait(false); break; + case NatsEvent.LameDuckModeActivated when LameDuckModeActivated != null: + await LameDuckModeActivated.InvokeAsync(this, args).ConfigureAwait(false); + break; } } } diff --git a/tests/NATS.Client.JetStream.Tests/NatsJsContextFactoryTest.cs b/tests/NATS.Client.JetStream.Tests/NatsJsContextFactoryTest.cs index 19bfc1e2d..d2d476f19 100644 --- a/tests/NATS.Client.JetStream.Tests/NatsJsContextFactoryTest.cs +++ b/tests/NATS.Client.JetStream.Tests/NatsJsContextFactoryTest.cs @@ -92,6 +92,8 @@ public class MockConnection : INatsConnection public event AsyncEventHandler? ReconnectFailed; public event AsyncEventHandler? MessageDropped; + + public event AsyncEventHandler? LameDuckModeActivated; #pragma warning restore CS0067 public INatsServerInfo? ServerInfo { get; } = null; From 15f7ac285223bd2ee06afba8857f27053f44ccf6 Mon Sep 17 00:00:00 2001 From: pzajaczkowski Date: Wed, 15 Jan 2025 19:03:51 +0100 Subject: [PATCH 2/6] Add test for ldm --- .../NatsConnectionTest.cs | 26 +++++++++++++++++ tests/NATS.Client.TestUtilities/NatsServer.cs | 28 +++++++++++++++++++ 2 files changed, 54 insertions(+) diff --git a/tests/NATS.Client.Core.Tests/NatsConnectionTest.cs b/tests/NATS.Client.Core.Tests/NatsConnectionTest.cs index da4a5c405..6cdb585f9 100644 --- a/tests/NATS.Client.Core.Tests/NatsConnectionTest.cs +++ b/tests/NATS.Client.Core.Tests/NatsConnectionTest.cs @@ -505,6 +505,32 @@ public async Task ReconnectOnOpenConnection_ShouldDisconnectAndOpenNewConnection openedCount.ShouldBe(1); disconnectedCount.ShouldBe(1); } + + [Fact] + public async Task LameDuckModeActivated_EventHandlerShouldBeInvokedWhenInfoWithLDMRecievied() + { + // Arrange + await using var server = NatsServer.Start(_output, _transportType); + await using var connection = server.CreateClientConnection(); + await connection.ConnectAsync(); + + var invocationCount = 0; + var ldmSignal = new WaitSignal(); + + connection.LameDuckModeActivated += (_, _) => + { + Interlocked.Increment(ref invocationCount); + ldmSignal.Pulse(); + return default; + }; + + // Act + server.SignalLameDuckMode(); + await ldmSignal; + + // Assert + invocationCount.ShouldBe(1); + } } [JsonSerializable(typeof(SampleClass))] diff --git a/tests/NATS.Client.TestUtilities/NatsServer.cs b/tests/NATS.Client.TestUtilities/NatsServer.cs index 43bd0b2f5..8cabe8f99 100644 --- a/tests/NATS.Client.TestUtilities/NatsServer.cs +++ b/tests/NATS.Client.TestUtilities/NatsServer.cs @@ -438,6 +438,34 @@ public void LogMessage(string categoryName, LogLevel logLevel, EventId e } } + public void SignalLameDuckMode() + { + if (ServerProcess == null || ServerProcess.HasExited) + { + throw new Exception("Cannot signal LDM, server process is not running."); + } + + var signalProcess = new Process + { + StartInfo = new ProcessStartInfo + { + FileName = NatsServerPath, + Arguments = $"--signal ldm={ServerProcess.Id}", + RedirectStandardOutput = true, + UseShellExecute = false, + }, + }; + + signalProcess.Start(); + signalProcess.WaitForExit(); + + if (signalProcess.ExitCode != 0) + { + var error = signalProcess.StandardError.ReadToEnd(); + throw new Exception($"Failed to signal lame duck mode: {error}"); + } + } + private static (string configFileName, string config, string cmd) GetCmd(NatsServerOpts opts) { var configFileName = Path.GetTempFileName(); From 8c631a8b5d35cee0524b98785212e445ba8a00ea Mon Sep 17 00:00:00 2001 From: pzajaczkowski Date: Mon, 20 Jan 2025 22:44:47 +0100 Subject: [PATCH 3/6] Change `LameDuckModeActivated` EventArgs to include `Uri` and change LDM state checking --- src/NATS.Client.Core/INatsConnection.cs | 2 +- src/NATS.Client.Core/NatsConnection.cs | 11 +++++------ src/NATS.Client.Core/NatsEventArgs.cs | 8 ++++++++ .../NatsJsContextFactoryTest.cs | 2 +- 4 files changed, 15 insertions(+), 8 deletions(-) diff --git a/src/NATS.Client.Core/INatsConnection.cs b/src/NATS.Client.Core/INatsConnection.cs index d893ca4ac..a2b75f4ed 100644 --- a/src/NATS.Client.Core/INatsConnection.cs +++ b/src/NATS.Client.Core/INatsConnection.cs @@ -28,7 +28,7 @@ public interface INatsConnection : INatsClient /// /// Event that is raised when server goes into Lame Duck Mode. /// - public event AsyncEventHandler? LameDuckModeActivated; + public event AsyncEventHandler? LameDuckModeActivated; /// /// Server information received from the NATS server. diff --git a/src/NATS.Client.Core/NatsConnection.cs b/src/NATS.Client.Core/NatsConnection.cs index 71ec05f2a..797e8dcd5 100644 --- a/src/NATS.Client.Core/NatsConnection.cs +++ b/src/NATS.Client.Core/NatsConnection.cs @@ -109,7 +109,7 @@ public NatsConnection(NatsOpts opts) public event AsyncEventHandler? MessageDropped; - public event AsyncEventHandler? LameDuckModeActivated; + public event AsyncEventHandler? LameDuckModeActivated; public INatsConnection Connection => this; @@ -141,10 +141,9 @@ internal ServerInfo? WritableServerInfo get => Interlocked.CompareExchange(ref _writableServerInfo, null, null); set { - var current = Interlocked.CompareExchange(ref _writableServerInfo, null, null); - if (current?.LameDuckMode == false && value?.LameDuckMode == true) + if (value?.LameDuckMode == true) { - _eventChannel.Writer.TryWrite((NatsEvent.LameDuckModeActivated, new NatsEventArgs(string.Empty))); + _eventChannel.Writer.TryWrite((NatsEvent.LameDuckModeActivated, new NatsLameDuckModeActivatedEventArgs(_currentConnectUri!.Uri))); } Interlocked.Exchange(ref _writableServerInfo, value); @@ -779,8 +778,8 @@ private async Task PublishEventsAsync() case NatsEvent.MessageDropped when MessageDropped != null && args is NatsMessageDroppedEventArgs error: await MessageDropped.InvokeAsync(this, error).ConfigureAwait(false); break; - case NatsEvent.LameDuckModeActivated when LameDuckModeActivated != null: - await LameDuckModeActivated.InvokeAsync(this, args).ConfigureAwait(false); + case NatsEvent.LameDuckModeActivated when LameDuckModeActivated != null && args is NatsLameDuckModeActivatedEventArgs uri: + await LameDuckModeActivated.InvokeAsync(this, uri).ConfigureAwait(false); break; } } diff --git a/src/NATS.Client.Core/NatsEventArgs.cs b/src/NATS.Client.Core/NatsEventArgs.cs index 254930288..265da8441 100644 --- a/src/NATS.Client.Core/NatsEventArgs.cs +++ b/src/NATS.Client.Core/NatsEventArgs.cs @@ -33,3 +33,11 @@ public NatsMessageDroppedEventArgs(NatsSubBase subscription, int pending, string public object? Data { get; } } + +public class NatsLameDuckModeActivatedEventArgs : NatsEventArgs +{ + public NatsLameDuckModeActivatedEventArgs(Uri uri) + : base("Lame duck mode activated") => Uri = uri; + + public Uri Uri { get; } +} diff --git a/tests/NATS.Client.JetStream.Tests/NatsJsContextFactoryTest.cs b/tests/NATS.Client.JetStream.Tests/NatsJsContextFactoryTest.cs index d2d476f19..615412237 100644 --- a/tests/NATS.Client.JetStream.Tests/NatsJsContextFactoryTest.cs +++ b/tests/NATS.Client.JetStream.Tests/NatsJsContextFactoryTest.cs @@ -93,7 +93,7 @@ public class MockConnection : INatsConnection public event AsyncEventHandler? MessageDropped; - public event AsyncEventHandler? LameDuckModeActivated; + public event AsyncEventHandler? LameDuckModeActivated; #pragma warning restore CS0067 public INatsServerInfo? ServerInfo { get; } = null; From 5d344bf14e470e9a8e5eda470750bddc007f60a3 Mon Sep 17 00:00:00 2001 From: pzajaczkowski Date: Tue, 21 Jan 2025 01:29:13 +0100 Subject: [PATCH 4/6] Change LDM test to use SystemEvents to invoke LDM mode --- .../NatsConnectionTest.cs | 39 ++++++++++++++++--- tests/NATS.Client.TestUtilities/NatsServer.cs | 28 ------------- 2 files changed, 34 insertions(+), 33 deletions(-) diff --git a/tests/NATS.Client.Core.Tests/NatsConnectionTest.cs b/tests/NATS.Client.Core.Tests/NatsConnectionTest.cs index 6cdb585f9..211abe97b 100644 --- a/tests/NATS.Client.Core.Tests/NatsConnectionTest.cs +++ b/tests/NATS.Client.Core.Tests/NatsConnectionTest.cs @@ -507,11 +507,36 @@ public async Task ReconnectOnOpenConnection_ShouldDisconnectAndOpenNewConnection } [Fact] - public async Task LameDuckModeActivated_EventHandlerShouldBeInvokedWhenInfoWithLDMRecievied() + public async Task LameDuckModeActivated_EventHandlerShouldBeInvokedWhenInfoWithLDMReceived() { - // Arrange - await using var server = NatsServer.Start(_output, _transportType); - await using var connection = server.CreateClientConnection(); + await using var natsServer = NatsServer.Start( + _output, + new NatsServerOptsBuilder() + .AddServerConfigText(""" + accounts: { + SYS: { + users: [ + {user: "sys", password: "password"} + ] + }, + } + + system_account: SYS + """) + .UseTransport(_transportType) + .Build()); + + var natsOpts = new NatsOpts + { + Url = natsServer.ClientUrl, + AuthOpts = new NatsAuthOpts + { + Username = "sys", + Password = "password", + }, + }; + + await using var connection = natsServer.CreateClientConnection(natsOpts); await connection.ConnectAsync(); var invocationCount = 0; @@ -524,8 +549,12 @@ public async Task LameDuckModeActivated_EventHandlerShouldBeInvokedWhenInfoWithL return default; }; + var subject = $"$SYS.REQ.SERVER.{connection.ServerInfo!.Id}.LDM"; + // Act - server.SignalLameDuckMode(); + var response = await connection.RequestAsync( + subject: subject, + data: "{}"); await ldmSignal; // Assert diff --git a/tests/NATS.Client.TestUtilities/NatsServer.cs b/tests/NATS.Client.TestUtilities/NatsServer.cs index 8cabe8f99..43bd0b2f5 100644 --- a/tests/NATS.Client.TestUtilities/NatsServer.cs +++ b/tests/NATS.Client.TestUtilities/NatsServer.cs @@ -438,34 +438,6 @@ public void LogMessage(string categoryName, LogLevel logLevel, EventId e } } - public void SignalLameDuckMode() - { - if (ServerProcess == null || ServerProcess.HasExited) - { - throw new Exception("Cannot signal LDM, server process is not running."); - } - - var signalProcess = new Process - { - StartInfo = new ProcessStartInfo - { - FileName = NatsServerPath, - Arguments = $"--signal ldm={ServerProcess.Id}", - RedirectStandardOutput = true, - UseShellExecute = false, - }, - }; - - signalProcess.Start(); - signalProcess.WaitForExit(); - - if (signalProcess.ExitCode != 0) - { - var error = signalProcess.StandardError.ReadToEnd(); - throw new Exception($"Failed to signal lame duck mode: {error}"); - } - } - private static (string configFileName, string config, string cmd) GetCmd(NatsServerOpts opts) { var configFileName = Path.GetTempFileName(); From 54741473a3bf5000222d8908a5467851cf003a39 Mon Sep 17 00:00:00 2001 From: pzajaczkowski Date: Tue, 21 Jan 2025 15:45:22 +0100 Subject: [PATCH 5/6] Fix `LameDuckModeActivated` event test --- tests/NATS.Client.Core.Tests/NatsConnectionTest.cs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/tests/NATS.Client.Core.Tests/NatsConnectionTest.cs b/tests/NATS.Client.Core.Tests/NatsConnectionTest.cs index 211abe97b..d5dd827fa 100644 --- a/tests/NATS.Client.Core.Tests/NatsConnectionTest.cs +++ b/tests/NATS.Client.Core.Tests/NatsConnectionTest.cs @@ -552,9 +552,9 @@ public async Task LameDuckModeActivated_EventHandlerShouldBeInvokedWhenInfoWithL var subject = $"$SYS.REQ.SERVER.{connection.ServerInfo!.Id}.LDM"; // Act - var response = await connection.RequestAsync( + await connection.RequestAsync( subject: subject, - data: "{}"); + data: $$"""{"cid":{{connection.ServerInfo!.ClientId}}}"""); await ldmSignal; // Assert From ea50c59a9ccc0562ca0c53410a2497a90b2cddcf Mon Sep 17 00:00:00 2001 From: Ziya Suzen Date: Tue, 21 Jan 2025 15:53:04 +0000 Subject: [PATCH 6/6] Restrict test to >=2.10 The LDM request feature test uses added to server version 2.10. --- tests/NATS.Client.Core.Tests/NatsConnectionTest.cs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tests/NATS.Client.Core.Tests/NatsConnectionTest.cs b/tests/NATS.Client.Core.Tests/NatsConnectionTest.cs index d5dd827fa..23ec3797f 100644 --- a/tests/NATS.Client.Core.Tests/NatsConnectionTest.cs +++ b/tests/NATS.Client.Core.Tests/NatsConnectionTest.cs @@ -506,7 +506,7 @@ public async Task ReconnectOnOpenConnection_ShouldDisconnectAndOpenNewConnection disconnectedCount.ShouldBe(1); } - [Fact] + [SkipIfNatsServer(versionEarlierThan: "2.10")] public async Task LameDuckModeActivated_EventHandlerShouldBeInvokedWhenInfoWithLDMReceived() { await using var natsServer = NatsServer.Start(