From 298ef092fefcec56be656fa891ea8d29ad5b9634 Mon Sep 17 00:00:00 2001 From: Roman Ettlinger Date: Sun, 16 Aug 2026 11:00:22 +0200 Subject: [PATCH] Treat zero MaxNotificationsPerPublish as unlimited (OPC 10000-4) (#4047) Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> --- .../Subscription/Subscription.cs | 12 +- .../Subscription/SubscriptionManager.cs | 11 +- .../Opc.Ua.Server.Tests/SubscriptionTests.cs | 370 +++++++++++++++++- 3 files changed, 384 insertions(+), 9 deletions(-) diff --git a/Libraries/Opc.Ua.Server/Subscription/Subscription.cs b/Libraries/Opc.Ua.Server/Subscription/Subscription.cs index f385569eec..e706f842a8 100644 --- a/Libraries/Opc.Ua.Server/Subscription/Subscription.cs +++ b/Libraries/Opc.Ua.Server/Subscription/Subscription.cs @@ -1133,12 +1133,16 @@ private NotificationMessage ConstructMessage( Diagnostics.NextSequenceNumber = m_sequenceNumber; } + uint notificationLimit = m_maxNotificationsPerPublish == 0 + ? uint.MaxValue + : m_maxNotificationsPerPublish; + // add events. - if (events.Count > 0 && notificationCount < m_maxNotificationsPerPublish) + if (events.Count > 0 && notificationCount < notificationLimit) { var notification = new EventNotificationList(); - while (events.Count > 0 && notificationCount < m_maxNotificationsPerPublish) + while (events.Count > 0 && notificationCount < notificationLimit) { notification.Events.Add(events.Dequeue()); notificationCount++; @@ -1148,7 +1152,7 @@ private NotificationMessage ConstructMessage( } // add datachanges (space permitting). - if (datachanges.Count > 0 && notificationCount < m_maxNotificationsPerPublish) + if (datachanges.Count > 0 && notificationCount < notificationLimit) { bool diagnosticsExist = false; var notification = new DataChangeNotification @@ -1157,7 +1161,7 @@ private NotificationMessage ConstructMessage( DiagnosticInfos = new DiagnosticInfoCollection(datachanges.Count) }; - while (datachanges.Count > 0 && notificationCount < m_maxNotificationsPerPublish) + while (datachanges.Count > 0 && notificationCount < notificationLimit) { MonitoredItemNotification datachange = datachanges.Dequeue(); notification.MonitoredItems.Add(datachange); diff --git a/Libraries/Opc.Ua.Server/Subscription/SubscriptionManager.cs b/Libraries/Opc.Ua.Server/Subscription/SubscriptionManager.cs index 77367ba65a..6c5547ffed 100644 --- a/Libraries/Opc.Ua.Server/Subscription/SubscriptionManager.cs +++ b/Libraries/Opc.Ua.Server/Subscription/SubscriptionManager.cs @@ -1976,13 +1976,18 @@ protected virtual uint CalculateLifetimeCount( /// protected virtual uint CalculateMaxNotificationsPerPublish(uint maxNotificationsPerPublish) { - if (maxNotificationsPerPublish == 0 || - maxNotificationsPerPublish > m_maxNotificationsPerPublish) + if (maxNotificationsPerPublish == 0) { return m_maxNotificationsPerPublish; } - return maxNotificationsPerPublish; + if (m_maxNotificationsPerPublish == 0 || + maxNotificationsPerPublish <= m_maxNotificationsPerPublish) + { + return maxNotificationsPerPublish; + } + + return m_maxNotificationsPerPublish; } /// diff --git a/Tests/Opc.Ua.Server.Tests/SubscriptionTests.cs b/Tests/Opc.Ua.Server.Tests/SubscriptionTests.cs index ee5a89f110..cb79c8c582 100644 --- a/Tests/Opc.Ua.Server.Tests/SubscriptionTests.cs +++ b/Tests/Opc.Ua.Server.Tests/SubscriptionTests.cs @@ -1,6 +1,9 @@ using System; using System.Collections.Generic; +using System.Linq; using System.Reflection; +using System.Threading; +using System.Threading.Tasks; using Moq; using NUnit.Framework; using Opc.Ua.Tests; @@ -15,6 +18,7 @@ public class SubscriptionTests private Mock m_serverMock; private Mock m_sessionMock; private DiagnosticsNodeManager m_diagnosticsNodeManager; + private Mock m_nodeManagerMock; private Mock m_queueFactoryMock; private ITelemetryContext m_telemetry; @@ -28,6 +32,8 @@ public void SetUp() m_serverMock.Setup(s => s.Telemetry).Returns(m_telemetry); m_serverMock.Setup(s => s.MonitoredItemQueueFactory).Returns(m_queueFactoryMock.Object); + m_serverMock.Setup(s => s.DiagnosticsWriteLock).Returns(new object()); + m_serverMock.Setup(s => s.ServerDiagnostics).Returns(new ServerDiagnosticsSummaryDataType()); var namespaceUris = new NamespaceTable(); m_serverMock.Setup(s => s.NamespaceUris).Returns(namespaceUris); @@ -46,10 +52,28 @@ public void SetUp() false); m_serverMock.Setup(s => s.DiagnosticsNodeManager).Returns(m_diagnosticsNodeManager); + var mainNodeManagerFactoryMock = new Mock(); + var coreNodeManagerMock = new Mock(); + var configNodeManagerMock = new Mock(); + mainNodeManagerFactoryMock.Setup(f => f.CreateConfigurationNodeManager()).Returns(configNodeManagerMock.Object); + mainNodeManagerFactoryMock.Setup(f => f.CreateCoreNodeManager(It.IsAny())).Returns(coreNodeManagerMock.Object); + m_serverMock.Setup(s => s.MainNodeManagerFactory).Returns(mainNodeManagerFactoryMock.Object); + + m_nodeManagerMock = new Mock( + m_serverMock.Object, + new ApplicationConfiguration { ServerConfiguration = new ServerConfiguration() }, + null, + Array.Empty()); + m_serverMock.Setup(s => s.NodeManager).Returns(m_nodeManagerMock.Object); + m_sessionMock.Setup(s => s.Id).Returns(new NodeId(Guid.NewGuid())); + m_sessionMock.Setup(s => s.DiagnosticsLock).Returns(new object()); + m_sessionMock.Setup(s => s.SessionDiagnostics).Returns(new SessionDiagnosticsDataType()); } - private Subscription CreateSubscription(double publishingInterval = 1000) + private Subscription CreateSubscription( + double publishingInterval = 1000, + uint maxNotificationsPerPublish = 0) { return new Subscription( m_serverMock.Object, @@ -58,7 +82,7 @@ private Subscription CreateSubscription(double publishingInterval = 1000) publishingInterval: publishingInterval, maxLifetimeCount: 10, maxKeepAliveCount: 5, - maxNotificationsPerPublish: 0, + maxNotificationsPerPublish: maxNotificationsPerPublish, priority: 0, publishingEnabled: true, maxMessageCount: 10); @@ -328,5 +352,347 @@ public void Publish_MultipleTimes_WithMaxMessageCount() Assert.That(moreNotifications2, Is.True); Assert.That(moreNotifications3, Is.False); } + + [Test] + public async Task PublishWithZeroLimitDrainsEventAndDataChangeQueuesAsync() + { + using Subscription subscription = CreateSubscription( + publishingInterval: 100, + maxNotificationsPerPublish: 0); + var publishLimits = new List(); + Mock eventItem = CreateEventMonitoredItem( + id: 1, + notificationCount: 2, + publishLimits); + Mock dataChangeItem = CreateDataChangeMonitoredItem( + id: 2, + notificationCount: 2, + publishLimits); + await RegisterMonitoredItemsAsync( + subscription, + eventItem.Object, + dataChangeItem.Object).ConfigureAwait(false); + + SetExpiryTime(subscription, HiResClock.TickCount64 - 100); + Assert.That( + subscription.PublishTimerExpired(), + Is.EqualTo(PublishingState.NotificationsAvailable)); + + var context = new OperationContext(m_sessionMock.Object, new DiagnosticsMasks()); + NotificationMessage message = subscription.Publish( + context, + out _, + out bool moreNotifications); + + Assert.That(message.NotificationData, Has.Count.EqualTo(2)); + var eventNotification = (EventNotificationList)ExtensionObject.ToEncodeable( + message.NotificationData[0]); + var dataChangeNotification = (DataChangeNotification)ExtensionObject.ToEncodeable( + message.NotificationData[1]); + Assert.Multiple(() => + { + Assert.That(moreNotifications, Is.False); + Assert.That(eventNotification.Events, Has.Count.EqualTo(2)); + Assert.That(dataChangeNotification.MonitoredItems, Has.Count.EqualTo(2)); + Assert.That(dataChangeNotification.DiagnosticInfos, Has.Count.EqualTo(2)); + Assert.That( + publishLimits, + Is.EqualTo(new[] { uint.MaxValue, uint.MaxValue })); + }); + } + + [Test] + public async Task PublishWithFiniteLimitBuildsAndDrainsQueuedMessagesAsync() + { + using Subscription subscription = CreateSubscription( + publishingInterval: 100, + maxNotificationsPerPublish: 2); + var publishLimits = new List(); + Mock firstEventItem = CreateEventMonitoredItem( + id: 1, + notificationCount: 3, + publishLimits); + Mock secondEventItem = CreateEventMonitoredItem( + id: 2, + notificationCount: 1, + publishLimits); + Mock dataChangeItem = CreateDataChangeMonitoredItem( + id: 3, + notificationCount: 2, + publishLimits); + await RegisterMonitoredItemsAsync( + subscription, + firstEventItem.Object, + secondEventItem.Object, + dataChangeItem.Object).ConfigureAwait(false); + + SetExpiryTime(subscription, HiResClock.TickCount64 - 100); + Assert.That( + subscription.PublishTimerExpired(), + Is.EqualTo(PublishingState.NotificationsAvailable)); + + var context = new OperationContext(m_sessionMock.Object, new DiagnosticsMasks()); + NotificationMessage firstMessage = subscription.Publish( + context, + out _, + out bool moreAfterFirst); + NotificationMessage secondMessage = subscription.Publish( + context, + out _, + out bool moreAfterSecond); + NotificationMessage thirdMessage = subscription.Publish( + context, + out _, + out bool moreAfterThird); + + Assert.Multiple(() => + { + Assert.That(firstMessage.NotificationData, Has.Count.EqualTo(1)); + Assert.That(secondMessage.NotificationData, Has.Count.EqualTo(1)); + Assert.That(thirdMessage.NotificationData, Has.Count.EqualTo(1)); + }); + var firstEvents = (EventNotificationList)ExtensionObject.ToEncodeable( + firstMessage.NotificationData[0]); + var secondEvents = (EventNotificationList)ExtensionObject.ToEncodeable( + secondMessage.NotificationData[0]); + var dataChanges = (DataChangeNotification)ExtensionObject.ToEncodeable( + thirdMessage.NotificationData[0]); + Assert.Multiple(() => + { + Assert.That(moreAfterFirst, Is.True); + Assert.That(moreAfterSecond, Is.True); + Assert.That(moreAfterThird, Is.False); + Assert.That(firstEvents.Events, Has.Count.EqualTo(2)); + Assert.That(secondEvents.Events, Has.Count.EqualTo(2)); + Assert.That(dataChanges.MonitoredItems, Has.Count.EqualTo(2)); + Assert.That(publishLimits, Is.EqualTo(new uint[] { 6, 6, 6 })); + }); + } + + [TestCase(0, 0L, 0L)] + [TestCase(0, 25L, 25L)] + [TestCase(10, 0L, 10L)] + [TestCase(10, 10L, 10L)] + [TestCase(10, 25L, 10L)] + [TestCase(0, uint.MaxValue, uint.MaxValue)] + [TestCase(10, uint.MaxValue, 10L)] + [TestCase(int.MaxValue, uint.MaxValue, int.MaxValue)] + public async Task CreateSubscriptionWithNotificationLimitsUsesEffectiveLimitAsync( + int serverLimit, + long requestedLimit, + long expectedLimit) + { + var configuration = new ApplicationConfiguration + { + ServerConfiguration = new ServerConfiguration + { + MaxNotificationsPerPublish = serverLimit + } + }; + using var manager = new SubscriptionManager( + m_serverMock.Object, + configuration); + + var context = new OperationContext(m_sessionMock.Object, new DiagnosticsMasks()); + CreateSubscriptionResponse response = await manager.CreateSubscriptionAsync( + context, + requestedPublishingInterval: 1000, + requestedLifetimeCount: 30, + requestedMaxKeepAliveCount: 10, + maxNotificationsPerPublish: (uint)requestedLimit, + publishingEnabled: true, + priority: 0).ConfigureAwait(false); + + ISubscription subscription = manager.GetSubscriptions() + .FirstOrDefault(s => s.Id == response.SubscriptionId); + Assert.That(subscription, Is.Not.Null); + Assert.That( + subscription.Diagnostics.MaxNotificationsPerPublish, + Is.EqualTo((uint)expectedLimit)); + } + + [TestCase(0, 0L, 0L)] + [TestCase(0, 25L, 25L)] + [TestCase(10, 0L, 10L)] + [TestCase(10, 10L, 10L)] + [TestCase(10, 25L, 10L)] + [TestCase(0, uint.MaxValue, uint.MaxValue)] + [TestCase(10, uint.MaxValue, 10L)] + [TestCase(int.MaxValue, uint.MaxValue, int.MaxValue)] + public async Task ModifySubscriptionWithNotificationLimitsUsesEffectiveLimitAsync( + int serverLimit, + long requestedLimit, + long expectedLimit) + { + var configuration = new ApplicationConfiguration + { + ServerConfiguration = new ServerConfiguration + { + MaxNotificationsPerPublish = serverLimit + } + }; + using var manager = new SubscriptionManager( + m_serverMock.Object, + configuration); + + var context = new OperationContext(m_sessionMock.Object, new DiagnosticsMasks()); + CreateSubscriptionResponse response = await manager.CreateSubscriptionAsync( + context, + requestedPublishingInterval: 1000, + requestedLifetimeCount: 30, + requestedMaxKeepAliveCount: 10, + maxNotificationsPerPublish: 1, + publishingEnabled: true, + priority: 0).ConfigureAwait(false); + + manager.ModifySubscription( + context, + response.SubscriptionId, + requestedPublishingInterval: 1000, + requestedLifetimeCount: 30, + requestedMaxKeepAliveCount: 10, + maxNotificationsPerPublish: (uint)requestedLimit, + priority: 0, + revisedPublishingInterval: out _, + revisedLifetimeCount: out _, + revisedMaxKeepAliveCount: out _); + + ISubscription subscription = manager.GetSubscriptions() + .FirstOrDefault(s => s.Id == response.SubscriptionId); + Assert.That(subscription, Is.Not.Null); + Assert.That( + subscription.Diagnostics.MaxNotificationsPerPublish, + Is.EqualTo((uint)expectedLimit)); + } + + private async Task RegisterMonitoredItemsAsync( + Subscription subscription, + params IMonitoredItem[] monitoredItems) + { + m_nodeManagerMock + .Setup(n => n.CreateMonitoredItemsAsync( + It.IsAny(), + subscription.Id, + It.IsAny(), + It.IsAny(), + It.IsAny>(), + It.IsAny>(), + It.IsAny>(), + It.IsAny>(), + It.IsAny(), + It.IsAny())) + .Callback< + OperationContext, + uint, + double, + TimestampsToReturn, + IList, + IList, + IList, + IList, + bool, + CancellationToken>( + (_, _, _, _, _, _, _, createdItems, _, _) => + { + for (int ii = 0; ii < monitoredItems.Length; ii++) + { + createdItems[ii] = monitoredItems[ii]; + } + }) + .Returns(default(ValueTask)); + + var requests = new MonitoredItemCreateRequest[monitoredItems.Length]; + for (int ii = 0; ii < monitoredItems.Length; ii++) + { + requests[ii] = new MonitoredItemCreateRequest + { + MonitoringMode = MonitoringMode.Reporting + }; + } + var itemsToCreate = new MonitoredItemCreateRequestCollection(requests); + + var context = new OperationContext(m_sessionMock.Object, new DiagnosticsMasks()); + CreateMonitoredItemsResponse response = await subscription.CreateMonitoredItemsAsync( + context, + TimestampsToReturn.Both, + itemsToCreate).ConfigureAwait(false); + + Assert.That(response.Results, Has.Count.EqualTo(monitoredItems.Length)); + } + + private static Mock CreateEventMonitoredItem( + uint id, + int notificationCount, + List publishLimits) + { + var item = new Mock(); + var createResult = new MonitoredItemCreateResult + { + StatusCode = StatusCodes.Good, + RevisedSamplingInterval = 0, + RevisedQueueSize = (uint)notificationCount + }; + item.SetupGet(i => i.Id).Returns(id); + item.SetupGet(i => i.IsReadyToPublish).Returns(true); + item.SetupGet(i => i.MonitoredItemType).Returns(MonitoredItemTypeMask.Events); + item.Setup(i => i.GetCreateResult(out createResult)).Returns(ServiceResult.Good); + item.Setup(i => i.Publish( + It.IsAny(), + It.IsAny>(), + It.IsAny())) + .Returns, uint>( + (_, notifications, maxNotificationsPerPublish) => + { + publishLimits.Add(maxNotificationsPerPublish); + for (int ii = 0; ii < notificationCount; ii++) + { + notifications.Enqueue(new EventFieldList()); + } + return false; + }); + return item; + } + + private static Mock CreateDataChangeMonitoredItem( + uint id, + int notificationCount, + List publishLimits) + { + var item = new Mock(); + var createResult = new MonitoredItemCreateResult + { + StatusCode = StatusCodes.Good, + RevisedSamplingInterval = 0, + RevisedQueueSize = (uint)notificationCount + }; + item.SetupGet(i => i.Id).Returns(id); + item.SetupGet(i => i.IsReadyToPublish).Returns(true); + item.SetupGet(i => i.MonitoredItemType).Returns(MonitoredItemTypeMask.DataChange); + item.Setup(i => i.GetCreateResult(out createResult)).Returns(ServiceResult.Good); + item.Setup(i => i.Publish( + It.IsAny(), + It.IsAny>(), + It.IsAny>(), + It.IsAny(), + It.IsAny())) + .Returns< + OperationContext, + Queue, + Queue, + uint, + Microsoft.Extensions.Logging.ILogger>( + (_, notifications, diagnostics, maxNotificationsPerPublish, _) => + { + publishLimits.Add(maxNotificationsPerPublish); + for (int ii = 0; ii < notificationCount; ii++) + { + notifications.Enqueue( + new MonitoredItemNotification { Value = new DataValue(ii) }); + diagnostics.Enqueue(new DiagnosticInfo()); + } + return false; + }); + return item; + } } }