From 9e8c43801bd46e9058997806348b262a3c4b709a Mon Sep 17 00:00:00 2001 From: Ziya Suzen Date: Thu, 1 May 2025 14:27:12 +0100 Subject: [PATCH] Fix subscriptions dispose exception handling --- src/NATS.Client.Core/NatsSubBase.cs | 4 ++-- .../Internal/NatsJSConsume.cs | 20 +++++++++++------ .../Internal/NatsJSFetch.cs | 22 ++++++++++++------- .../Internal/NatsJSOrderedConsume.cs | 15 ++++++++----- 4 files changed, 39 insertions(+), 22 deletions(-) diff --git a/src/NATS.Client.Core/NatsSubBase.cs b/src/NATS.Client.Core/NatsSubBase.cs index 2737cb6b7..521266ff0 100644 --- a/src/NATS.Client.Core/NatsSubBase.cs +++ b/src/NATS.Client.Core/NatsSubBase.cs @@ -239,6 +239,8 @@ public virtual ValueTask DisposeAsync() _idleTimeoutTimer?.Dispose(); _startUpTimeoutTimer?.Dispose(); + _tokenRegistration.Dispose(); + if (Exception != null) { if (Exception is NatsSubException { Exception: not null } nse) @@ -249,8 +251,6 @@ public virtual ValueTask DisposeAsync() throw Exception; } - _tokenRegistration.Dispose(); - return unsubscribeAsync; } diff --git a/src/NATS.Client.JetStream/Internal/NatsJSConsume.cs b/src/NATS.Client.JetStream/Internal/NatsJSConsume.cs index 29d696538..0b8124341 100644 --- a/src/NATS.Client.JetStream/Internal/NatsJSConsume.cs +++ b/src/NATS.Client.JetStream/Internal/NatsJSConsume.cs @@ -190,16 +190,22 @@ public void Delivered(int msgSize) public override async ValueTask DisposeAsync() { Interlocked.Exchange(ref _disposed, 1); - await base.DisposeAsync().ConfigureAwait(false); - await _pullTask.ConfigureAwait(false); + try + { + await base.DisposeAsync().ConfigureAwait(false); + } + finally + { + await _pullTask.ConfigureAwait(false); #if NETSTANDARD2_0 - _timer.Dispose(); + _timer.Dispose(); #else - await _timer.DisposeAsync().ConfigureAwait(false); + await _timer.DisposeAsync().ConfigureAwait(false); #endif - if (_notificationChannel != null) - { - await _notificationChannel.DisposeAsync(); + if (_notificationChannel != null) + { + await _notificationChannel.DisposeAsync(); + } } } diff --git a/src/NATS.Client.JetStream/Internal/NatsJSFetch.cs b/src/NATS.Client.JetStream/Internal/NatsJSFetch.cs index 98e218360..1d30663bf 100644 --- a/src/NATS.Client.JetStream/Internal/NatsJSFetch.cs +++ b/src/NATS.Client.JetStream/Internal/NatsJSFetch.cs @@ -160,17 +160,23 @@ public void ResetHeartbeatTimer() public override async ValueTask DisposeAsync() { Interlocked.Exchange(ref _disposed, 1); - await base.DisposeAsync().ConfigureAwait(false); + try + { + await base.DisposeAsync().ConfigureAwait(false); + } + finally + { #if NETSTANDARD2_0 - _hbTimer.Dispose(); - _expiresTimer.Dispose(); + _hbTimer.Dispose(); + _expiresTimer.Dispose(); #else - await _hbTimer.DisposeAsync().ConfigureAwait(false); - await _expiresTimer.DisposeAsync().ConfigureAwait(false); + await _hbTimer.DisposeAsync().ConfigureAwait(false); + await _expiresTimer.DisposeAsync().ConfigureAwait(false); #endif - if (_notificationChannel != null) - { - await _notificationChannel.DisposeAsync(); + if (_notificationChannel != null) + { + await _notificationChannel.DisposeAsync(); + } } } diff --git a/src/NATS.Client.JetStream/Internal/NatsJSOrderedConsume.cs b/src/NATS.Client.JetStream/Internal/NatsJSOrderedConsume.cs index 0969e4bd5..806f6d83a 100644 --- a/src/NATS.Client.JetStream/Internal/NatsJSOrderedConsume.cs +++ b/src/NATS.Client.JetStream/Internal/NatsJSOrderedConsume.cs @@ -136,14 +136,19 @@ public override async ValueTask DisposeAsync() Interlocked.Exchange(ref _disposed, 1); _context.Connection.ConnectionDisconnected -= ConnectionOnConnectionDisconnected; - - await base.DisposeAsync().ConfigureAwait(false); - await _pullTask.ConfigureAwait(false); + try + { + await base.DisposeAsync().ConfigureAwait(false); + } + finally + { + await _pullTask.ConfigureAwait(false); #if NETSTANDARD2_0 - _timer.Dispose(); + _timer.Dispose(); #else - await _timer.DisposeAsync().ConfigureAwait(false); + await _timer.DisposeAsync().ConfigureAwait(false); #endif + } } internal override ValueTask WriteReconnectCommandsAsync(CommandWriter commandWriter, int sid)