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
126 changes: 87 additions & 39 deletions src/Opc.Ua.Server/Subscription/SessionPublishQueue.cs
Original file line number Diff line number Diff line change
Expand Up @@ -427,16 +427,15 @@ public void PublishCompleted(ISubscription subscription, bool moreNotifications)
{
lock (m_lock)
{
// Flag the subscription as available and let the selection policy decide
// which of the available subscriptions is handed to a waiting request.
queuedSubscription.Publishing = false;
queuedSubscription.ReadyToPublish = moreNotifications;
queuedSubscription.Timestamp = DateTime.UtcNow;

if (moreNotifications)
{
AssignSubscriptionToRequest(queuedSubscription);
}
else
{
queuedSubscription.ReadyToPublish = false;
queuedSubscription.Timestamp = DateTime.UtcNow;
AssignSubscriptionsToRequests();
}
}
}
Expand Down Expand Up @@ -484,6 +483,7 @@ internal IReadOnlyList<QueuedSubscription> CapturePublishTimerSnapshot()
internal void PublishTimerExpired(IReadOnlyList<QueuedSubscription> queuedSubscriptions)
{
var subscriptionsToDelete = new List<ISubscription>();
List<QueuedSubscription>? notifyingSubscriptions = null;

// check each available subscription.
for (int ii = 0; ii < queuedSubscriptions.Count; ii++)
Expand Down Expand Up @@ -521,16 +521,34 @@ internal void PublishTimerExpired(IReadOnlyList<QueuedSubscription> queuedSubscr
continue;
}

// assign subscription to request if one is available.
// collect the subscription, it is assigned to a request further below.
if (!subscription.Publishing)
{
lock (m_lock)
(notifyingSubscriptions ??= []).Add(subscription);
}
}

if (notifyingSubscriptions != null)
{
lock (m_lock)
{
// Flag every notifying subscription as available before any request is
// served. Assigning them one by one while iterating would hand the
// waiting requests out in the (unordered) iteration order of the
// subscription dictionary and bypass the priority and timestamp based
// selection policy applied by PublishAsync.
foreach (QueuedSubscription subscription in notifyingSubscriptions)
{
if (!subscription.Publishing)
if (subscription.Publishing || subscription.ReadyToPublish)
{
AssignSubscriptionToRequest(subscription);
continue;
}

subscription.ReadyToPublish = true;
subscription.Timestamp = DateTime.UtcNow;
}

AssignSubscriptionsToRequests();
}
}

Expand All @@ -556,47 +574,77 @@ private bool TryRemoveExact(QueuedSubscription queuedSubscription)
}

/// <summary>
/// Checks the state of the subscriptions.
/// Hands the subscriptions that are ready to publish to the waiting publish
/// requests. The subscriptions are selected with the same priority and timestamp
/// based policy as <see cref="PublishAsync"/>, so the order in which subscriptions
/// became ready does not determine which one is published first.
/// </summary>
private void AssignSubscriptionToRequest(QueuedSubscription subscription)
private void AssignSubscriptionsToRequests()
{
lock (m_lock)
while (m_queuedRequests.Count > 0)
{
// find a request.
while (m_queuedRequests.Count > 0)
QueuedSubscription? subscriptionToPublish = GetSubscriptionToPublish();

if (subscriptionToPublish == null)
{
QueuedPublishRequest request = m_queuedRequests.First!.Value;
m_queuedRequests.RemoveFirst();
break;
}

if (request.Tcs.Task.IsCompleted)
{
request.Dispose();
continue;
}
if (!TryAssignSubscriptionToRequest(subscriptionToPublish))
{
// no usable request left, keep the subscription available.
subscriptionToPublish.Publishing = false;
break;
}
}
}

// check secure channel.
if (!m_session.IsSecureChannelValid(request.SecureChannelId))
{
m_logger.PublishAbandonedBecauseTheSecureChannelChanged();
request.Tcs.TrySetException(new ServiceResultException(StatusCodes.BadSecureChannelIdInvalid));
request.Dispose();
continue;
}
/// <summary>
/// Completes the next usable publish request with the subscription. Returns false
/// if no usable request is queued, in which case the subscription stays available.
/// </summary>
private bool TryAssignSubscriptionToRequest(QueuedSubscription subscription)
{
// find a request.
while (m_queuedRequests.Count > 0)
{
QueuedPublishRequest request = m_queuedRequests.First!.Value;
m_queuedRequests.RemoveFirst();

if (request.Tcs.Task.IsCompleted)
{
request.Dispose();
continue;
}

// check secure channel.
if (!m_session.IsSecureChannelValid(request.SecureChannelId))
{
m_logger.PublishAbandonedBecauseTheSecureChannelChanged();
request.Tcs.TrySetException(new ServiceResultException(StatusCodes.BadSecureChannelIdInvalid));
request.Dispose();
continue;
}

m_logger.PUBLISHIdAssignedToSubscriptionSubscriptionId(
request.SecureChannelId,
subscription.Subscription.Id);
subscription.Publishing = true;

subscription.Publishing = true;
request.Tcs.TrySetResult(subscription.Subscription);
if (!request.Tcs.TrySetResult(subscription.Subscription))
{
// the request was cancelled or timed out in the meantime.
subscription.Publishing = false;
request.Dispose();
return;
continue;
}

// mark it as available.
subscription.ReadyToPublish = true;
subscription.Timestamp = DateTime.UtcNow;
m_logger.PUBLISHIdAssignedToSubscriptionSubscriptionId(
request.SecureChannelId,
subscription.Subscription.Id);

request.Dispose();
return true;
}

return false;
}

/// <summary>
Expand Down
118 changes: 118 additions & 0 deletions tests/Opc.Ua.Server.Tests/SessionPublishQueueTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -127,6 +127,124 @@ public async Task PublishAsync_ReturnsSubscriptionIfReadyAsync()
Assert.That(result, Is.SameAs(subMock.Object));
}

[Test]
[CancelAfter(10000)]
public async Task PublishTimerPreservesReadySubscriptionTimestampOrderAsync()
{
using var queue = new SessionPublishQueue(
m_serverMock.Object,
m_sessionMock.Object,
kMaxPublishRequests);

PublishingState publishingState = PublishingState.Idle;
var timerOrder = new List<ISubscription>();
var subscription1 = new Mock<ISubscription>();
subscription1.Setup(s => s.Id).Returns(1);
subscription1.Setup(s => s.Priority).Returns(1);
subscription1
.Setup(s => s.PublishTimerExpired())
.Callback(() => timerOrder.Add(subscription1.Object))
.Returns(() => publishingState);
queue.Add(subscription1.Object);

var subscription2 = new Mock<ISubscription>();
subscription2.Setup(s => s.Id).Returns(2);
subscription2.Setup(s => s.Priority).Returns(1);
subscription2
.Setup(s => s.PublishTimerExpired())
.Callback(() => timerOrder.Add(subscription2.Object))
.Returns(() => publishingState);
queue.Add(subscription2.Object);

queue.PublishTimerExpired();
Assert.That(timerOrder, Has.Count.EqualTo(2));

ISubscription newerSubscription = timerOrder[0];
ISubscription olderSubscription = timerOrder[1];
queue.PublishCompleted(olderSubscription, false);
DateTime timestampBoundary = DateTime.UtcNow;
Assert.That(
SpinWait.SpinUntil(
() => DateTime.UtcNow > timestampBoundary,
TimeSpan.FromSeconds(1)),
Is.True,
"The clock did not advance while preparing distinct subscription timestamps.");
queue.PublishCompleted(newerSubscription, false);

queue.Requeue(newerSubscription);
queue.Requeue(olderSubscription);
publishingState = PublishingState.NotificationsAvailable;
timerOrder.Clear();
queue.PublishTimerExpired();

Assert.That(timerOrder, Has.Count.EqualTo(2));
Assert.That(timerOrder[0], Is.SameAs(newerSubscription));
Assert.That(timerOrder[1], Is.SameAs(olderSubscription));
ISubscription result = await queue.PublishAsync(
"channel1",
DateTime.MaxValue,
false,
null,
CancellationToken.None).ConfigureAwait(false);
Assert.That(result, Is.SameAs(olderSubscription));
}

[Test]
[CancelAfter(10000)]
public async Task PublishTimerAssignsWaitingRequestToHighestPrioritySubscriptionAsync()
{
using var queue = new SessionPublishQueue(
m_serverMock.Object,
m_sessionMock.Object,
kMaxPublishRequests);

// The queued subscriptions are enumerated in an unspecified order, so the
// priority is derived from the observed order: the Subscription the publish
// timer visits first is the one with the lowest priority.
var timerOrder = new List<ISubscription>();
var subscription1 = new Mock<ISubscription>();
var subscription2 = new Mock<ISubscription>();

byte GetPriority(ISubscription subscription)
{
return timerOrder.Count > 0 && ReferenceEquals(timerOrder[0], subscription)
? (byte)1
: (byte)200;
}

subscription1.Setup(s => s.Id).Returns(1);
subscription1.Setup(s => s.Priority).Returns(() => GetPriority(subscription1.Object));
subscription1
.Setup(s => s.PublishTimerExpired())
.Callback(() => timerOrder.Add(subscription1.Object))
.Returns(PublishingState.NotificationsAvailable);
queue.Add(subscription1.Object);

subscription2.Setup(s => s.Id).Returns(2);
subscription2.Setup(s => s.Priority).Returns(() => GetPriority(subscription2.Object));
subscription2
.Setup(s => s.PublishTimerExpired())
.Callback(() => timerOrder.Add(subscription2.Object))
.Returns(PublishingState.NotificationsAvailable);
queue.Add(subscription2.Object);

Task<ISubscription> request = queue.PublishAsync(
"channel1",
DateTime.MaxValue,
false,
null,
CancellationToken.None);
Assert.That(request.IsCompleted, Is.False);

queue.PublishTimerExpired();

Assert.That(timerOrder, Has.Count.EqualTo(2));

ISubscription result = await request.ConfigureAwait(false);
Assert.That(result.Priority, Is.EqualTo(200));
Assert.That(result, Is.SameAs(timerOrder[1]));
}

[Test]
public void PublishAsync_WhenParked_NotifiesParkSinkOnce()
{
Expand Down
Loading