Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions Libraries/Opc.Ua.Client/Session/Session.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
16 changes: 11 additions & 5 deletions Libraries/Opc.Ua.Client/Subscription/Subscription.cs
Original file line number Diff line number Diff line change
Expand Up @@ -44,7 +44,6 @@ namespace Opc.Ua.Client
/// </summary>
public class Subscription : ISnapshotRestore<SubscriptionState>, IDisposable, ICloneable
{
private const int kMinKeepAliveTimerInterval = 1000;
private const int kKeepAliveTimerMargin = 1000;
private const int kRepublishMessageExpiredTimeout = 10000;

Expand All @@ -53,6 +52,11 @@ public class Subscription : ISnapshotRestore<SubscriptionState>, IDisposable, IC
/// </summary>
public const int RepublishMessageTimeout = 2500;

/// <summary>
/// Minimum keep alive interval
/// </summary>
public const int MinKeepAliveTimerInterval = 1000;

/// <summary>
/// Create subscription
/// </summary>
Expand Down Expand Up @@ -2104,7 +2108,9 @@ private void HandleOnKeepAliveStopped()

if (session != null &&
session.Connected &&
!session.Reconnecting)
!session.Reconnecting &&
!session.KeepAliveStopped
)
{
TraceState("PUBLISHING STOPPED");

Expand Down Expand Up @@ -2198,7 +2204,7 @@ private int BeginPublishTimeout()
{
return Math.Max(
Math.Min(m_keepAliveInterval * 3, int.MaxValue),
kMinKeepAliveTimerInterval);
MinKeepAliveTimerInterval);
}

/// <summary>
Expand Down Expand Up @@ -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;
}
Expand Down
97 changes: 80 additions & 17 deletions Tests/Opc.Ua.Client.Tests/Subscription/SubscriptionUnitTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -73,28 +89,31 @@ private static Task AwaitForRepublishTimeout(CancellationToken ct)
return Task.Delay(Subscription.RepublishMessageTimeout + 100, ct);
}

private static ISession BuildSessionMock(Func<uint, uint, bool> republishHandler)
private static ISession BuildSessionMock(Func<uint, uint, bool> republishHandler = null, Action<Mock<ISession>> setup = null)
{
uint subscriptionIdSeed = 0u;

var session = new Mock<ISession>();
session
.Setup(x => x.RepublishAsync(It.IsAny<uint>(), It.IsAny<uint>(), It.IsAny<CancellationToken>()))
.ReturnsAsync<uint, uint, CancellationToken, ISession, (bool, ServiceResult)>((
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<uint>(), It.IsAny<uint>(), It.IsAny<CancellationToken>()))
.ReturnsAsync<uint, uint, CancellationToken, ISession, (bool, ServiceResult)>((
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(
Expand Down Expand Up @@ -149,6 +168,7 @@ private static ISession BuildSessionMock(Func<uint, uint, bool> republishHandler
StatusCodes.Good)],
DiagnosticInfos = [.. subscriptionIds.ConvertAll(_ => new DiagnosticInfo())]
});
setup?.Invoke(session);
return session.Object;
}

Expand Down Expand Up @@ -350,5 +370,48 @@ public async Task WillRepublishIfMissedMessagesInBetweenOfPublishesAsync(
Assert.That(subscription.Notifications, Is.EquivalentTo(messages.Skip(1)));
}
}

[DatapointSource]
public IEnumerable<KeepAliveTestDataProvider> 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<PublishStateChangedMask>(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));
}
}
}
Loading