From 714cd994355c3ab4b9f231fc2eefca82be9fe89d Mon Sep 17 00:00:00 2001 From: Anatoly Kochnev Date: Fri, 24 Apr 2026 15:21:06 +0500 Subject: [PATCH] [Client] Fixed excessive tasks spawning during session connection loss Added checks for Session.KeepAliveStopped before trying to execute Publish for subscriptions. --- Libraries/Opc.Ua.Client/Session/Session.cs | 6 ++ .../Subscription/Subscription.cs | 16 ++- .../Subscription/SubscriptionUnitTests.cs | 97 +++++++++++++++---- 3 files changed, 97 insertions(+), 22 deletions(-) diff --git a/Libraries/Opc.Ua.Client/Session/Session.cs b/Libraries/Opc.Ua.Client/Session/Session.cs index f510d97e00..0d3fe10a64 100644 --- a/Libraries/Opc.Ua.Client/Session/Session.cs +++ b/Libraries/Opc.Ua.Client/Session/Session.cs @@ -3554,6 +3554,12 @@ public bool BeginPublish(int timeout) return false; } + if (KeepAliveStopped) + { + m_logger.LogWarning("Publish skipped due to session lost connection. Last successfull keepalive: {LastKeepAlive}", LastKeepAliveTime); + return false; + } + // get event handler to modify ack list PublishSequenceNumbersToAcknowledgeEventHandler? callback = m_PublishSequenceNumbersToAcknowledge; diff --git a/Libraries/Opc.Ua.Client/Subscription/Subscription.cs b/Libraries/Opc.Ua.Client/Subscription/Subscription.cs index b5a2e8f699..42f26daed1 100644 --- a/Libraries/Opc.Ua.Client/Subscription/Subscription.cs +++ b/Libraries/Opc.Ua.Client/Subscription/Subscription.cs @@ -44,7 +44,6 @@ namespace Opc.Ua.Client /// public class Subscription : ISnapshotRestore, IDisposable, ICloneable { - private const int kMinKeepAliveTimerInterval = 1000; private const int kKeepAliveTimerMargin = 1000; private const int kRepublishMessageExpiredTimeout = 10000; @@ -53,6 +52,11 @@ public class Subscription : ISnapshotRestore, IDisposable, IC /// public const int RepublishMessageTimeout = 2500; + /// + /// Minimum keep alive interval + /// + public const int MinKeepAliveTimerInterval = 1000; + /// /// Create subscription /// @@ -2104,7 +2108,9 @@ private void HandleOnKeepAliveStopped() if (session != null && session.Connected && - !session.Reconnecting) + !session.Reconnecting && + !session.KeepAliveStopped + ) { TraceState("PUBLISHING STOPPED"); @@ -2198,7 +2204,7 @@ private int BeginPublishTimeout() { return Math.Max( Math.Min(m_keepAliveInterval * 3, int.MaxValue), - kMinKeepAliveTimerInterval); + MinKeepAliveTimerInterval); } /// @@ -2312,12 +2318,12 @@ private int CalculateKeepAliveInterval() { int keepAliveInterval = (int) Math.Min(CurrentPublishingInterval * (CurrentKeepAliveCount + 1), int.MaxValue); - if (keepAliveInterval < kMinKeepAliveTimerInterval) + if (keepAliveInterval < MinKeepAliveTimerInterval) { keepAliveInterval = (int)Math.Min( PublishingInterval * (KeepAliveCount + 1), int.MaxValue); - keepAliveInterval = Math.Max(kMinKeepAliveTimerInterval, keepAliveInterval); + keepAliveInterval = Math.Max(MinKeepAliveTimerInterval, keepAliveInterval); } return keepAliveInterval; } diff --git a/Tests/Opc.Ua.Client.Tests/Subscription/SubscriptionUnitTests.cs b/Tests/Opc.Ua.Client.Tests/Subscription/SubscriptionUnitTests.cs index 783558df32..2dbc432c3f 100644 --- a/Tests/Opc.Ua.Client.Tests/Subscription/SubscriptionUnitTests.cs +++ b/Tests/Opc.Ua.Client.Tests/Subscription/SubscriptionUnitTests.cs @@ -42,6 +42,22 @@ namespace Opc.Ua.Client.Tests [Parallelizable] public class SubscriptionUnitTests { + private const PublishStateChangedMask kSessionNotConnected = PublishStateChangedMask.Stopped | PublishStateChangedMask.SessionNotConnected; + + public record KeepAliveTestDataProvider(PublishStateChangedMask ExpectedPublishState) : IFormattable + { + public bool SessionConnected { get; init; } + public bool SessionReconnecting { get; init; } + public bool SessionKeepAliveStopped { get; init; } + public string ToString(string format, IFormatProvider formatProvider) + { + return $"Connected={SessionConnected}, " + + $"reconnecting={SessionReconnecting}, " + + $"keepAlive={SessionKeepAliveStopped}. " + + $"Expected status:{ExpectedPublishState}"; + } + } + private sealed class SubscriptionContainer : IDisposable { private readonly CancellationTokenRegistration m_tokedCancellation; @@ -73,28 +89,31 @@ private static Task AwaitForRepublishTimeout(CancellationToken ct) return Task.Delay(Subscription.RepublishMessageTimeout + 100, ct); } - private static ISession BuildSessionMock(Func republishHandler) + private static ISession BuildSessionMock(Func republishHandler = null, Action> setup = null) { uint subscriptionIdSeed = 0u; var session = new Mock(); - session - .Setup(x => x.RepublishAsync(It.IsAny(), It.IsAny(), It.IsAny())) - .ReturnsAsync(( - subscriptionId, - sequenceNumber, - ct) => - { - if (subscriptionId > subscriptionIdSeed) - { - return (true, StatusCodes.BadSubscriptionIdInvalid); - } - if (republishHandler(subscriptionId, sequenceNumber)) + if (republishHandler is not null) + { + session + .Setup(x => x.RepublishAsync(It.IsAny(), It.IsAny(), It.IsAny())) + .ReturnsAsync(( + subscriptionId, + sequenceNumber, + ct) => { - return (true, ServiceResult.Good); - } - return (true, StatusCodes.BadMessageNotAvailable); - }); + if (subscriptionId > subscriptionIdSeed) + { + return (true, StatusCodes.BadSubscriptionIdInvalid); + } + if (republishHandler(subscriptionId, sequenceNumber)) + { + return (true, ServiceResult.Good); + } + return (true, StatusCodes.BadMessageNotAvailable); + }); + } session .Setup(x => x .CreateSubscriptionAsync( @@ -149,6 +168,7 @@ private static ISession BuildSessionMock(Func republishHandler StatusCodes.Good)], DiagnosticInfos = [.. subscriptionIds.ConvertAll(_ => new DiagnosticInfo())] }); + setup?.Invoke(session); return session.Object; } @@ -350,5 +370,48 @@ public async Task WillRepublishIfMissedMessagesInBetweenOfPublishesAsync( Assert.That(subscription.Notifications, Is.EquivalentTo(messages.Skip(1))); } } + + [DatapointSource] + public IEnumerable SubscriptionKeepAliveValues() + { + yield return new(PublishStateChangedMask.Stopped) { SessionConnected = true, SessionReconnecting = false, SessionKeepAliveStopped = false }; + + yield return new(kSessionNotConnected) { SessionConnected = false, SessionReconnecting = false, SessionKeepAliveStopped = false }; + yield return new(kSessionNotConnected) { SessionConnected = false, SessionReconnecting = true, SessionKeepAliveStopped = false }; + yield return new(kSessionNotConnected) { SessionConnected = false, SessionReconnecting = false, SessionKeepAliveStopped = true }; + yield return new(kSessionNotConnected) { SessionConnected = false, SessionReconnecting = true, SessionKeepAliveStopped = true }; + + yield return new(kSessionNotConnected) { SessionConnected = true, SessionReconnecting = true, SessionKeepAliveStopped = false }; + yield return new(kSessionNotConnected) { SessionConnected = true, SessionReconnecting = true, SessionKeepAliveStopped = true }; + + yield return new(kSessionNotConnected) { SessionConnected = true, SessionReconnecting = false, SessionKeepAliveStopped = true }; + } + + [Theory] + [CancelAfter(Subscription.MinKeepAliveTimerInterval * 10)] + public async Task RespectsStateOfSessionDuringKeepAliveCalls(KeepAliveTestDataProvider testData, CancellationToken ct) + { + var keepAliveCompleted = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + void KeepAliveHasTriggered(Subscription x, PublishStateChangedEventArgs y) => keepAliveCompleted.TrySetResult(y.Status); + ISession session = BuildSessionMock( + setup: mock => + { + mock.Setup(x => x.Connected).Returns(testData.SessionConnected); + mock.Setup(x => x.Reconnecting).Returns(testData.SessionReconnecting); + mock.Setup(x => x.KeepAliveStopped).Returns(testData.SessionKeepAliveStopped); + }); + + using var subscription = new Subscription( + NUnitTelemetryContext.Create(), + new() { PublishingEnabled = true }) + { + Session = session + }; + subscription.PublishStatusChanged += KeepAliveHasTriggered; + await subscription.CreateAsync(ct).ConfigureAwait(false); + await Task.WhenAny(keepAliveCompleted.Task, Task.Delay(-1, ct)).ConfigureAwait(false); + + Assert.That(keepAliveCompleted.Task.Result, Is.EqualTo(testData.ExpectedPublishState)); + } } }