From 99d5d440c37c9f198b397f6ecdcf24eee6c7a39a Mon Sep 17 00:00:00 2001 From: Marc Date: Fri, 31 Jul 2026 16:00:51 +0200 Subject: [PATCH 1/5] Make subscription transfer between sessions transactional Transfer now prepares every subscription before any of them moves, so a failure part way through leaves the source session exactly as it was rather than with a subset of its subscriptions already gone. Monitored item resend-data triggers captured during preparation are restored when a prepared transfer is rolled back. The session publish queue tracks the transfer claim so a subscription cannot be published by the source session once it has been prepared for transfer, and cannot be lost if the transfer is abandoned. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 9e6a5abf-3299-4cd1-9855-010fedbf0ad8 --- .../NodeManager/MasterNodeManager.cs | 11 + .../MonitoredItem/IMonitoredItem.cs | 12 + .../MonitoredItem/MonitoredItem.cs | 14 +- .../Subscription/SessionPublishQueue.cs | 254 ++++++++- .../Subscription/Subscription.cs | 440 ++++++++++++++- .../Subscription/SubscriptionManager.cs | 483 ++++++++++++----- .../MasterNodeManagerDeterministicTests.cs | 1 - .../Opc.Ua.Server.Tests/SubscriptionTests.cs | 501 ++++++++++++++++++ .../SessionPublishQueueConcurrencyTests.cs | 267 ++++++++++ 9 files changed, 1814 insertions(+), 169 deletions(-) create mode 100644 tests/Opc.Ua.Subscriptions.Tests/SessionPublishQueueConcurrencyTests.cs diff --git a/src/Opc.Ua.Server/NodeManager/MasterNodeManager.cs b/src/Opc.Ua.Server/NodeManager/MasterNodeManager.cs index 05a5fbacd7..f5ac5496c4 100644 --- a/src/Opc.Ua.Server/NodeManager/MasterNodeManager.cs +++ b/src/Opc.Ua.Server/NodeManager/MasterNodeManager.cs @@ -5515,6 +5515,7 @@ private async ValueTask var processedItems = new List(monitoredItems.Count); var originalErrors = new ServiceResult[errors.Count]; + var resendStates = new bool[monitoredItems.Count]; var effectiveTransferOptions = new MonitoredItemTransferOptions { DeferInitialValues = sendInitialValues || @@ -5526,6 +5527,7 @@ private async ValueTask { IMonitoredItem? monitoredItem = monitoredItems[ii]; originalErrors[ii] = errors[ii]; + resendStates[ii] = monitoredItem?.IsResendData ?? false; bool isDetached = monitoredItem is IDetachableMonitoredItem { IsDetached: true @@ -5581,6 +5583,7 @@ await owner.TransferMonitoredItemsAsync( monitoredItems, errors, originalErrors, + resendStates, participants, effectiveTransferOptions); try @@ -5617,6 +5620,7 @@ await failedTransaction.RollbackAsync(CancellationToken.None) monitoredItems, errors, originalErrors, + resendStates, participants, effectiveTransferOptions); } @@ -5677,6 +5681,7 @@ public MonitoredItemTransferTransaction( IList monitoredItems, IList errors, ServiceResult[] originalErrors, + bool[] resendStates, List participants, MonitoredItemTransferOptions transferOptions) { @@ -5685,6 +5690,7 @@ public MonitoredItemTransferTransaction( m_monitoredItems = monitoredItems; m_errors = errors; m_originalErrors = originalErrors; + m_resendStates = resendStates; m_participants = participants; m_transferOptions = transferOptions; } @@ -5759,6 +5765,10 @@ await participant.NodeManager.RollbackMonitoredItemsTransferAsync( for (int ii = 0; ii < m_monitoredItems.Count; ii++) { + if (m_monitoredItems[ii] is IMonitoredItemTransferState transferState) + { + transferState.RestoreResendDataTrigger(m_resendStates[ii]); + } m_errors[ii] = m_originalErrors[ii]; } @@ -5775,6 +5785,7 @@ await participant.NodeManager.RollbackMonitoredItemsTransferAsync( private readonly IList m_monitoredItems; private readonly IList m_errors; private readonly ServiceResult[] m_originalErrors; + private readonly bool[] m_resendStates; private readonly List m_participants; private readonly MonitoredItemTransferOptions m_transferOptions; private int m_state; diff --git a/src/Opc.Ua.Server/Subscription/MonitoredItem/IMonitoredItem.cs b/src/Opc.Ua.Server/Subscription/MonitoredItem/IMonitoredItem.cs index 3e3d2d8a5c..562b88f1d7 100644 --- a/src/Opc.Ua.Server/Subscription/MonitoredItem/IMonitoredItem.cs +++ b/src/Opc.Ua.Server/Subscription/MonitoredItem/IMonitoredItem.cs @@ -346,6 +346,18 @@ public interface ISampledDataChangeMonitoredItem : IDataChangeMonitoredItem2 void SetSamplingInterval(double samplingInterval); } + /// + /// Restores monitored item transient state when a prepared subscription transfer is rolled back. + /// + internal interface IMonitoredItemTransferState + { + /// + /// Restores the resend-data trigger to the value captured before transfer preparation. + /// + /// The original resend-data trigger state. + void RestoreResendDataTrigger(bool resendData); + } + /// /// Defines constants for the monitored item type. /// diff --git a/src/Opc.Ua.Server/Subscription/MonitoredItem/MonitoredItem.cs b/src/Opc.Ua.Server/Subscription/MonitoredItem/MonitoredItem.cs index 6cfb00ee21..ea8ee7f9a2 100644 --- a/src/Opc.Ua.Server/Subscription/MonitoredItem/MonitoredItem.cs +++ b/src/Opc.Ua.Server/Subscription/MonitoredItem/MonitoredItem.cs @@ -42,7 +42,8 @@ public class MonitoredItem : IEventMonitoredItem, ISampledDataChangeMonitoredItem, ITriggeredMonitoredItem, - IDetachableMonitoredItem + IDetachableMonitoredItem, + IMonitoredItemTransferState { /// /// Initializes the object with its node type. @@ -510,6 +511,14 @@ public void SetupResendDataTrigger() } } + void IMonitoredItemTransferState.RestoreResendDataTrigger(bool resendData) + { + lock (m_lock) + { + m_resendData = resendData; + } + } + /// /// Sets a flag indicating that the item has been triggered and should publish. /// @@ -1591,7 +1600,8 @@ public virtual bool Publish( else { // pull any unprocessed data. - if (m_calculator != null && m_calculator.HasEndTimePassed(DateTime.UtcNow)) + if (m_calculator != null && + m_calculator.HasEndTimePassed(DateTime.UtcNow)) { while (m_calculator.TryGetProcessedValue(false, out DataValue processedValue)) { diff --git a/src/Opc.Ua.Server/Subscription/SessionPublishQueue.cs b/src/Opc.Ua.Server/Subscription/SessionPublishQueue.cs index b858988a39..44d93c06e3 100644 --- a/src/Opc.Ua.Server/Subscription/SessionPublishQueue.cs +++ b/src/Opc.Ua.Server/Subscription/SessionPublishQueue.cs @@ -63,6 +63,7 @@ public SessionPublishQueue( m_session = session ?? throw new ArgumentNullException(nameof(session)); m_queuedRequests = new LinkedList(); m_queuedSubscriptions = new ConcurrentDictionary(); + m_transferClaims = []; m_maxRequestCount = maxPublishRequests; m_timeProvider = timeProvider ?? (server as ITimeProviderProvider)?.TimeProvider @@ -104,6 +105,7 @@ protected virtual void Dispose(bool disposing) } m_queuedSubscriptions.Clear(); + m_transferClaims.Clear(); } } } @@ -193,6 +195,7 @@ public IList Close() // clear the queue. m_queuedSubscriptions.Clear(); + m_transferClaims.Clear(); } foreach (ISubscription subscription in queuedSubscriptions) @@ -257,19 +260,94 @@ internal bool ContainsSubscription(ISubscription subscription) } /// - /// Removes the exact subscription entry so a transfer can claim it. + /// Claims and removes the exact subscription entry before transfer callbacks run. /// - internal bool TryRemoveForTransfer(ISubscription subscription) + internal bool TryClaimForTransfer( + Subscription subscription, + ISession sourceSession, + out SubscriptionTransferClaim? claim) { - if (!m_queuedSubscriptions.TryGetValue( - subscription.Id, - out QueuedSubscription? queuedSubscription) || - !ReferenceEquals(queuedSubscription.Subscription, subscription)) + lock (m_lock) { - return false; + claim = null; + if (!m_queuedSubscriptions.TryGetValue( + subscription.Id, + out QueuedSubscription? queuedSubscription) || + !ReferenceEquals(queuedSubscription.Subscription, subscription) || + m_transferClaims.ContainsKey(subscription.Id)) + { + return false; + } + + if (!subscription.TryBeginTransfer(sourceSession)) + { + return false; + } + if (!TryRemoveExact(queuedSubscription)) + { + subscription.AbortTransfer(sourceSession); + return false; + } + + claim = new SubscriptionTransferClaim(queuedSubscription); + m_transferClaims.Add(subscription.Id, claim); + return true; + } + } + + /// + /// Restores a previously claimed subscription entry when transfer preparation fails before ownership changes. + /// + /// The queue entry claim to return to active publishing. + /// true when the claim was current and the entry was restored. + internal bool RestoreTransferClaim(SubscriptionTransferClaim claim) + { + lock (m_lock) + { + uint subscriptionId = claim.Entry.Subscription.Id; + if (!m_transferClaims.TryGetValue( + subscriptionId, + out SubscriptionTransferClaim? currentClaim) || + !ReferenceEquals(currentClaim, claim) || + !m_queuedSubscriptions.TryAdd(subscriptionId, claim.Entry)) + { + return false; + } + + m_transferClaims.Remove(subscriptionId); + return true; + } + } + + /// + /// Removes a transfer claim after the destination session has accepted the subscription. + /// + /// The queue entry claim that completed. + internal void CompleteTransferClaim(SubscriptionTransferClaim claim) + { + lock (m_lock) + { + uint subscriptionId = claim.Entry.Subscription.Id; + if (m_transferClaims.TryGetValue( + subscriptionId, + out SubscriptionTransferClaim? currentClaim) && + ReferenceEquals(currentClaim, claim)) + { + m_transferClaims.Remove(subscriptionId); + } } + } - return TryRemoveExact(queuedSubscription); + internal bool TryRemoveForTransfer(ISubscription subscription) + { + lock (m_lock) + { + return m_queuedSubscriptions.TryGetValue( + subscription.Id, + out QueuedSubscription? queuedSubscription) && + ReferenceEquals(queuedSubscription.Subscription, subscription) && + TryRemoveExact(queuedSubscription); + } } /// @@ -380,12 +458,11 @@ public void Acknowledge( if ((context.DiagnosticsMask & DiagnosticsMasks.OperationAll) != 0) { - DiagnosticInfo? diagnosticInfo = ServerUtils - .CreateDiagnosticInfo( - m_server, - context, - result, - m_logger); + DiagnosticInfo? diagnosticInfo = ServerUtils.CreateDiagnosticInfo( + m_server, + context, + result, + m_logger); acknowledgeDiagnosticInfoList.Add(diagnosticInfo!); diagnosticsExist = true; } @@ -422,13 +499,25 @@ public void Acknowledge( /// public void PublishCompleted(ISubscription subscription, bool moreNotifications) { - if (m_queuedSubscriptions.TryGetValue(subscription.Id, - out QueuedSubscription? queuedSubscription)) + if (!m_queuedSubscriptions.TryGetValue( + subscription.Id, + out QueuedSubscription? queuedSubscription)) { lock (m_lock) { - queuedSubscription.Publishing = false; + PublishCompletedTransferClaimNoLock(subscription, moreNotifications); + } + return; + } + lock (m_lock) + { + if (m_queuedSubscriptions.TryGetValue( + subscription.Id, + out QueuedSubscription? currentSubscription) && + ReferenceEquals(currentSubscription, queuedSubscription)) + { + queuedSubscription.Publishing = false; if (moreNotifications) { AssignSubscriptionToRequest(queuedSubscription); @@ -438,7 +527,10 @@ public void PublishCompleted(ISubscription subscription, bool moreNotifications) queuedSubscription.ReadyToPublish = false; queuedSubscription.Timestamp = DateTime.UtcNow; } + return; } + + PublishCompletedTransferClaimNoLock(subscription, moreNotifications); } } @@ -447,13 +539,30 @@ public void PublishCompleted(ISubscription subscription, bool moreNotifications) /// public void Requeue(ISubscription subscription) { - if (m_queuedSubscriptions.TryGetValue(subscription.Id, out QueuedSubscription? queuedSubscription)) + if (!m_queuedSubscriptions.TryGetValue( + subscription.Id, + out QueuedSubscription? queuedSubscription)) { lock (m_lock) + { + RequeueTransferClaimNoLock(subscription); + } + return; + } + + lock (m_lock) + { + if (m_queuedSubscriptions.TryGetValue( + subscription.Id, + out QueuedSubscription? currentSubscription) && + ReferenceEquals(currentSubscription, queuedSubscription)) { queuedSubscription.Publishing = false; queuedSubscription.ReadyToPublish = true; + return; } + + RequeueTransferClaimNoLock(subscription); } } @@ -489,6 +598,11 @@ internal void PublishTimerExpired(IReadOnlyList queuedSubscr for (int ii = 0; ii < queuedSubscriptions.Count; ii++) { QueuedSubscription subscription = queuedSubscriptions[ii]; + if (!IsCurrentSubscription(subscription)) + { + continue; + } + PublishingState state = subscription.Subscription.PublishTimerExpired(); // check for expired subscription. @@ -511,7 +625,13 @@ internal void PublishTimerExpired(IReadOnlyList queuedSubscr // check if idle. if (state == PublishingState.Idle) { - subscription.ReadyToPublish = false; + lock (m_lock) + { + if (IsCurrentSubscriptionNoLock(subscription)) + { + subscription.ReadyToPublish = false; + } + } continue; } @@ -526,7 +646,8 @@ internal void PublishTimerExpired(IReadOnlyList queuedSubscr { lock (m_lock) { - if (!subscription.Publishing) + if (IsCurrentSubscriptionNoLock(subscription) && + !subscription.Publishing) { AssignSubscriptionToRequest(subscription); } @@ -543,7 +664,49 @@ internal void PublishTimerExpired(IReadOnlyList queuedSubscr /// internal bool TryRemoveForExpiration(QueuedSubscription queuedSubscription) { - return TryRemoveExact(queuedSubscription); + lock (m_lock) + { + return TryRemoveExact(queuedSubscription); + } + } + + private bool IsCurrentSubscription(QueuedSubscription queuedSubscription) + { + return IsCurrentSubscriptionNoLock(queuedSubscription); + } + + private void PublishCompletedTransferClaimNoLock( + ISubscription subscription, + bool moreNotifications) + { + if (m_transferClaims.TryGetValue( + subscription.Id, + out SubscriptionTransferClaim? transferClaim) && + ReferenceEquals(transferClaim.Entry.Subscription, subscription)) + { + transferClaim.Entry.Publishing = false; + transferClaim.Entry.ReadyToPublish = moreNotifications; + } + } + + private void RequeueTransferClaimNoLock(ISubscription subscription) + { + if (m_transferClaims.TryGetValue( + subscription.Id, + out SubscriptionTransferClaim? transferClaim) && + ReferenceEquals(transferClaim.Entry.Subscription, subscription)) + { + transferClaim.Entry.Publishing = false; + transferClaim.Entry.ReadyToPublish = true; + } + } + + private bool IsCurrentSubscriptionNoLock(QueuedSubscription queuedSubscription) + { + return m_queuedSubscriptions.TryGetValue( + queuedSubscription.Subscription.Id, + out QueuedSubscription? currentSubscription) && + ReferenceEquals(currentSubscription, queuedSubscription); } private bool TryRemoveExact(QueuedSubscription queuedSubscription) @@ -696,6 +859,10 @@ public void Dispose() /// internal sealed class QueuedSubscription { + /// + /// Initializes the queue entry for a subscription owned by this session. + /// + /// The subscription tracked by the publish queue. public QueuedSubscription(ISubscription subscription) { Subscription = subscription; @@ -703,12 +870,47 @@ public QueuedSubscription(ISubscription subscription) Timestamp = DateTime.UtcNow; } + /// + /// Gets the subscription associated with the queue entry. + /// public ISubscription Subscription { get; } + + /// + /// Gets or sets the UTC timestamp used for publish scheduling and timeout decisions. + /// public DateTime Timestamp { get; set; } + + /// + /// Gets or sets whether the subscription has notifications ready for a publish response. + /// public bool ReadyToPublish { get; set; } + + /// + /// Gets or sets whether the queue entry is currently assigned to an outstanding publish request. + /// public bool Publishing { get; set; } } + /// + /// Holds the exact queue entry removed while a subscription is being transferred to another session. + /// + internal sealed class SubscriptionTransferClaim + { + /// + /// Initializes a transfer claim for the removed queue entry. + /// + /// The queue entry held outside active publishing during transfer. + public SubscriptionTransferClaim(QueuedSubscription entry) + { + Entry = entry; + } + + /// + /// Gets the queue entry that must be restored or completed exactly once. + /// + public QueuedSubscription Entry { get; } + } + /// /// Dumps the current state of the session queue. /// @@ -773,6 +975,7 @@ internal void TraceState(string context, params object[] args) private readonly ISession m_session; private readonly LinkedList m_queuedRequests; private readonly ConcurrentDictionary m_queuedSubscriptions; + private readonly Dictionary m_transferClaims; private readonly int m_maxRequestCount; private readonly TimeProvider m_timeProvider; } @@ -782,10 +985,16 @@ internal void TraceState(string context, params object[] args) /// internal static partial class SessionPublishQueueLog { + /// + /// Logs that a publish request was abandoned because its secure channel no longer matches the queued request. + /// [LoggerMessage(EventId = ServerEventIds.SessionPublishQueue + 0, Level = LogLevel.Warning, Message = "Publish abandoned because the secure channel changed.")] public static partial void PublishAbandonedBecauseTheSecureChannelChanged(this ILogger logger); + /// + /// Logs the trace-level assignment of a queued publish request to a subscription. + /// [LoggerMessage(EventId = ServerEventIds.SessionPublishQueue + 1, Level = LogLevel.Trace, Message = "PUBLISH: #{Id} Assigned To Subscription({SubscriptionId}).")] public static partial void PUBLISHIdAssignedToSubscriptionSubscriptionId( @@ -793,6 +1002,9 @@ public static partial void PUBLISHIdAssignedToSubscriptionSubscriptionId( string id, uint subscriptionId); + /// + /// Logs a trace-level snapshot of the publish queue counters for diagnostics. + /// [LoggerMessage(EventId = ServerEventIds.SessionPublishQueue + 2, Level = LogLevel.Trace, Message = "PublishQueue {Context}, SessionId={SessionId}, SubscriptionCount={SubscriptionCount}, " + "RequestCount={RequestCount}, ReadyToPublishCount={ReadyToPublishCount}, " + diff --git a/src/Opc.Ua.Server/Subscription/Subscription.cs b/src/Opc.Ua.Server/Subscription/Subscription.cs index 94c4cad648..0019610048 100644 --- a/src/Opc.Ua.Server/Subscription/Subscription.cs +++ b/src/Opc.Ua.Server/Subscription/Subscription.cs @@ -483,7 +483,7 @@ monitoredItem.Value is not IDetachableMonitoredItem { IsDetached: true } && - AreSameNodeManager(monitoredItem.Value.NodeManager, nodeManager)); + IsOwnedBy(monitoredItem.Value, nodeManager)); } } @@ -503,7 +503,7 @@ monitoredItem is not IDetachableMonitoredItem { IsDetached: true } && - AreSameNodeManager(monitoredItem.NodeManager, nodeManager)) + IsOwnedBy(monitoredItem, nodeManager)) ]; } } @@ -560,12 +560,25 @@ bool ISubscriptionMonitoredItemLifecycle.ContainsMonitoredItem( } } - private static bool AreSameNodeManager( - IAsyncNodeManager first, - IAsyncNodeManager second) + private static bool IsOwnedBy( + IMonitoredItem monitoredItem, + IAsyncNodeManager nodeManager) { - return ReferenceEquals(first, second) || - ReferenceEquals(first.SyncNodeManager, second.SyncNodeManager); + IAsyncNodeManager? monitoredItemOwner = monitoredItem.NodeManager; + if (monitoredItemOwner is null) + { + return false; + } + if (ReferenceEquals(monitoredItemOwner, nodeManager)) + { + return true; + } + + INodeManager? monitoredItemSync = monitoredItemOwner.SyncNodeManager; + INodeManager? nodeManagerSync = nodeManager.SyncNodeManager; + return monitoredItemSync is not null && + nodeManagerSync is not null && + ReferenceEquals(monitoredItemSync, nodeManagerSync); } /// @@ -619,6 +632,13 @@ public PublishingState PublishTimerExpired() { lock (m_lock) { + // OPC 10000-4 ยง5.14.1.2, Table 79 handles TransferSubscriptions as a single transition + // that sets the new Session, returns the response and issues Good_SubscriptionTransferred. + if (m_transferInProgress) + { + return PublishingState.Idle; + } + long currentTime = m_timeProvider.GetTimestampMilliseconds(); // check if publish interval has elapsed. @@ -782,14 +802,13 @@ public async ValueTask TransferSessionAsync(OperationContext context, bool sendI errors.Add(null!); } - await m_server.NodeManager - .TransferMonitoredItemsAsync( - context, - sendInitialValues, - monitoredItems, - errors, - transferOptions: null, - cancellationToken) + await m_server.NodeManager.TransferMonitoredItemsAsync( + context, + sendInitialValues, + monitoredItems, + errors, + null, + cancellationToken) .ConfigureAwait(false); int badTransfers = 0; @@ -827,6 +846,7 @@ await m_server.NodeManager } } + /// public bool IsTransferIdentityCompatible(ISession targetSession) { @@ -860,6 +880,376 @@ public bool IsTransferIdentityCompatible(ISession targetSession) StringComparison.Ordinal); } + /// + /// Reserves the subscription for transfer while it is still owned by the expected source session. + /// + /// The session that currently owns the subscription. + /// true when the transfer reservation was acquired. + internal bool TryBeginTransfer(ISession? sourceSession) + { + lock (m_lock) + { + if (m_transferInProgress || + m_expired || + !ReferenceEquals(Session, sourceSession)) + { + return false; + } + + m_transferInProgress = true; + return true; + } + } + + /// + /// Prepares monitored item state for a session transfer without making the new owner visible yet. + /// + /// The operation context for the destination session. + /// The session that currently owns the subscription. + /// Whether initial values should be sent after transfer commits. + /// The token that aborts transfer preparation. + /// A prepared transfer that can be committed or rolled back by the caller. + /// + /// The subscription is no longer reserved by the source session. + /// + internal async ValueTask PrepareSessionTransferAsync( + OperationContext context, + ISession? sourceSession, + bool sendInitialValues, + CancellationToken cancellationToken) + { + List monitoredItems; + lock (m_lock) + { + if (!m_transferInProgress || + !ReferenceEquals(Session, sourceSession)) + { + throw new ServiceResultException( + StatusCodes.BadSubscriptionIdInvalid, + "Subscription source changed during transfer."); + } + monitoredItems = m_monitoredItems.Select(v => v.Value.Value).ToList(); + } + + var errors = new List(monitoredItems.Count); + for (int ii = 0; ii < monitoredItems.Count; ii++) + { + errors.Add(null!); + } + + OperationContext? sourceContext = null; + IMonitoredItemTransferTransaction? monitoredItemTransaction = null; + try + { + if (m_server.NodeManager is IMonitoredItemTransferCoordinator coordinator) + { + sourceContext = sourceSession != null + ? new OperationContext(sourceSession, context.DiagnosticsMask) + : new OperationContext( + new RequestHeader(), + null!, + RequestType.Unknown, + RequestLifetime.None); + monitoredItemTransaction = await coordinator.PrepareMonitoredItemsTransferAsync( + context, + sourceContext, + sendInitialValues, + monitoredItems, + errors, + new MonitoredItemTransferOptions + { + DeferInitialValues = sendInitialValues + }, + cancellationToken) + .ConfigureAwait(false); + } + else + { + var fallbackTransaction = new ResendStateTransferTransaction( + monitoredItems); + try + { + // The non-coordinator contract has no commit callback. + // Keep legacy eager resend semantics instead of + // requesting deferral that cannot be committed here. + await m_server.NodeManager.TransferMonitoredItemsAsync( + context, + sendInitialValues, + monitoredItems, + errors, + transferOptions: null, + cancellationToken) + .ConfigureAwait(false); + } + catch + { + await fallbackTransaction.RollbackAsync(CancellationToken.None) + .ConfigureAwait(false); + throw; + } + monitoredItemTransaction = fallbackTransaction; + } + } + catch + { + sourceContext?.Dispose(); + throw; + } + + int badTransfers = 0; + for (int ii = 0; ii < errors.Count; ii++) + { + if (ServiceResult.IsBad(errors[ii])) + { + badTransfers++; + } + } + if (badTransfers > 0) + { + m_logger.FailedToTransferCountMonitoredItems(badTransfers); + } + + return new PreparedSessionTransfer( + this, + sourceSession, + context.Session, + monitoredItemTransaction, + sourceContext); + } + + /// + /// Releases the transfer reservation after the destination session is already the owner. + /// + /// The session that must currently own the subscription. + /// Ownership changed before transfer completion. + internal void CompleteTransfer(ISession destinationSession) + { + lock (m_lock) + { + if (!m_transferInProgress || + !ReferenceEquals(Session, destinationSession)) + { + throw new ServiceResultException( + StatusCodes.BadSubscriptionIdInvalid, + "Subscription ownership changed while completing transfer."); + } + m_transferInProgress = false; + } + } + + /// + /// Releases a transfer reservation without changing ownership when preparation cannot continue. + /// + /// The source session that still owns the subscription. + internal void AbortTransfer(ISession? sourceSession) + { + lock (m_lock) + { + if (m_transferInProgress && + ReferenceEquals(Session, sourceSession)) + { + m_transferInProgress = false; + } + } + } + + /// + /// Represents a prepared subscription transfer whose ownership and monitored item effects can + /// still be committed or rolled back. + /// + internal sealed class PreparedSessionTransfer + { + /// + /// Initializes a prepared transfer with the source context and monitored item transaction to dispose. + /// + /// The subscription being transferred. + /// The session that owned the subscription when preparation started. + /// The session that will own the subscription after commit. + /// The prepared monitored item transfer, if one was needed. + /// The source operation context created for monitored item callbacks. + public PreparedSessionTransfer( + Subscription subscription, + ISession? sourceSession, + ISession destinationSession, + IMonitoredItemTransferTransaction? monitoredItemTransaction, + OperationContext? sourceContext) + { + m_subscription = subscription; + m_sourceSession = sourceSession; + m_destinationSession = destinationSession; + m_monitoredItemTransaction = monitoredItemTransaction; + m_sourceContext = sourceContext; + } + + /// + /// Moves subscription ownership and diagnostics to the destination session. + /// + /// + /// The subscription is no longer owned by the source session. + /// + public void CommitOwnership() + { + lock (m_subscription.m_lock) + { + if (!m_subscription.m_transferInProgress || + !ReferenceEquals(m_subscription.Session, m_sourceSession)) + { + throw new ServiceResultException( + StatusCodes.BadSubscriptionIdInvalid, + "Subscription ownership changed during transfer."); + } + m_subscription.Session = m_destinationSession; + } + + lock (m_subscription.DiagnosticsWriteLock) + { + m_subscription.Diagnostics.SessionId = m_destinationSession.Id; + } + } + + /// + /// Makes the prepared monitored item transfer effects visible after ownership commits. + /// + public void CommitMonitoredItemEffects() + { + m_monitoredItemTransaction?.Commit(); + } + + /// + /// Restores source ownership and monitored item state when a later transfer step fails. + /// + /// The token that aborts monitored item rollback. + /// A task that completes when rollback has finished. + /// One or more rollback steps failed. + public async ValueTask RollbackAsync(CancellationToken cancellationToken) + { + var rollbackErrors = new List(); + try + { + try + { + lock (m_subscription.m_lock) + { + if (ReferenceEquals(m_subscription.Session, m_destinationSession)) + { + m_subscription.Session = m_sourceSession!; + } + else if (!ReferenceEquals(m_subscription.Session, m_sourceSession)) + { + throw new ServiceResultException( + StatusCodes.BadSubscriptionIdInvalid, + "Subscription ownership changed while rolling back transfer."); + } + } + + lock (m_subscription.DiagnosticsWriteLock) + { + m_subscription.Diagnostics.SessionId = m_sourceSession?.Id ?? default; + } + } + catch (Exception error) + { + rollbackErrors.Add(error); + } + + if (m_monitoredItemTransaction != null) + { + try + { + await m_monitoredItemTransaction.RollbackAsync(cancellationToken) + .ConfigureAwait(false); + } + catch (Exception error) + { + rollbackErrors.Add(error); + } + } + } + finally + { + DisposeSourceContext(); + } + + if (rollbackErrors.Count > 0) + { + throw new AggregateException( + "The subscription transfer could not be fully rolled back.", + rollbackErrors); + } + } + + /// + /// Releases the source operation context after a successful transfer. + /// + public void Complete() + { + DisposeSourceContext(); + } + + private void DisposeSourceContext() + { + Interlocked.Exchange(ref m_sourceContext, null)?.Dispose(); + } + + private readonly Subscription m_subscription; + private readonly ISession? m_sourceSession; + private readonly ISession m_destinationSession; + private readonly IMonitoredItemTransferTransaction? m_monitoredItemTransaction; + private OperationContext? m_sourceContext; + } + + private sealed class ResendStateTransferTransaction : + IMonitoredItemTransferTransaction + { + /// + /// Captures resend-data state so the legacy transfer path can be rolled back. + /// + /// The monitored items whose resend state is captured. + public ResendStateTransferTransaction(IList monitoredItems) + { + m_monitoredItems = monitoredItems; + m_resendStates = new bool[monitoredItems.Count]; + for (int ii = 0; ii < monitoredItems.Count; ii++) + { + m_resendStates[ii] = monitoredItems[ii]?.IsResendData ?? false; + } + } + + /// + /// Completes the legacy transfer transaction; no deferred work is required. + /// + public void Commit() + { + } + + /// + /// Restores each monitored item's resend-data trigger to the captured value. + /// + /// Unused token; rollback is synchronous and idempotent. + /// A completed task. + public ValueTask RollbackAsync(CancellationToken cancellationToken) + { + if (Interlocked.Exchange(ref m_rolledBack, 1) != 0) + { + return default; + } + + for (int ii = 0; ii < m_monitoredItems.Count; ii++) + { + if (m_monitoredItems[ii] is IMonitoredItemTransferState transferState) + { + transferState.RestoreResendDataTrigger(m_resendStates[ii]); + } + } + return default; + } + + private readonly IList m_monitoredItems; + private readonly bool[] m_resendStates; + private int m_rolledBack; + + } + /// /// Restores ownership if a transfer failed after assigning its destination. /// @@ -2937,6 +3327,13 @@ private void VerifySession(OperationContext context) throw new ServiceResultException(StatusCodes.BadSubscriptionIdInvalid); } + if (m_transferInProgress) + { + throw new ServiceResultException( + StatusCodes.BadSubscriptionIdInvalid, + "Subscription transfer is in progress."); + } + if (!ReferenceEquals(context.Session, Session)) { throw new ServiceResultException( @@ -3075,6 +3472,7 @@ private void TraceState(LogLevel logLevel, TraceStateId id, string context) private readonly NodeId m_diagnosticsId; private bool m_refreshInProgress; private bool m_expired; + private bool m_transferInProgress; private readonly Dictionary> m_itemsToTrigger; private readonly bool m_supportsDurable; private readonly ILogger m_logger; @@ -3085,18 +3483,30 @@ private void TraceState(LogLevel logLevel, TraceStateId id, string context) /// internal static partial class SubscriptionLog { + /// + /// Logs that deleting monitored items for a subscription failed. + /// [LoggerMessage(EventId = ServerEventIds.Subscription + 0, Level = LogLevel.Error, Message = "Delete items for subscription failed.")] public static partial void DeleteItemsForSubscriptionFailed(this ILogger logger, Exception ex); + /// + /// Logs the number of monitored items that could not be transferred. + /// [LoggerMessage(EventId = ServerEventIds.Subscription + 1, Level = LogLevel.Trace, Message = "Failed to transfer {Count} Monitored Items")] public static partial void FailedToTransferCountMonitoredItems(this ILogger logger, int count); + /// + /// Logs an invariant violation where monitored items were queued without available notifications. + /// [LoggerMessage(EventId = ServerEventIds.Subscription + 2, Level = LogLevel.Error, Message = "Oops! MonitoredItems queued but no notifications available.")] public static partial void OopsMonitoredItemsQueuedButNoNotificationsAvailable(this ILogger logger); + /// + /// Logs that durable subscription setup was requested without a durable monitored item queue factory. + /// [LoggerMessage(EventId = ServerEventIds.Subscription + 3, Level = LogLevel.Error, Message = "SetSubscriptionDurable requested for subscription with id {SubscriptionId}, but no " + "IMonitoredItemQueueFactory that supports durable queues was registered")] diff --git a/src/Opc.Ua.Server/Subscription/SubscriptionManager.cs b/src/Opc.Ua.Server/Subscription/SubscriptionManager.cs index 7d99d21694..d3c8ff3f63 100644 --- a/src/Opc.Ua.Server/Subscription/SubscriptionManager.cs +++ b/src/Opc.Ua.Server/Subscription/SubscriptionManager.cs @@ -41,7 +41,7 @@ namespace Opc.Ua.Server /// /// A generic session manager object for a server. /// - public class SubscriptionManager : ISubscriptionManager + public class SubscriptionManager : ISubscriptionManager, IAsyncDisposable { /// /// Initializes the manager with its configuration. @@ -100,6 +100,7 @@ public SubscriptionManager( // create a event to signal shutdown. m_shutdownEvent = new ManualResetEvent(true); + m_workerShutdown = new CancellationTokenSource(); // create queue and event for condition refresh worker m_conditionRefreshEvent = new ManualResetEvent(false); @@ -115,47 +116,99 @@ public void Dispose() GC.SuppressFinalize(this); } + /// + /// Frees managed resources that require asynchronous shutdown. + /// + public async ValueTask DisposeAsync() + { + await DisposeAsyncCore().ConfigureAwait(false); + GC.SuppressFinalize(this); + } + /// /// An overrideable version of the Dispose. /// protected virtual void Dispose(bool disposing) { - if (disposing) + if (disposing && Interlocked.Exchange(ref m_disposed, 1) == 0) { - List? subscriptions = null; - List? publishQueues = null; + SignalConditionRefreshShutdown(); + CaptureManagedResources( + out List publishQueues, + out List subscriptions); + DisposeManagedResources(publishQueues, subscriptions); + DisposeSynchronizationResources(); + } + } - m_semaphoreSlim.Wait(); - try - { - publishQueues = [.. m_publishQueues.Values]; - m_publishQueues.Clear(); + /// + /// An overrideable version of the asynchronous dispose. + /// + protected virtual async ValueTask DisposeAsyncCore() + { + if (Interlocked.Exchange(ref m_disposed, 1) != 0) + { + return; + } - subscriptions = [.. m_subscriptions.Values]; - m_subscriptions.Clear(); - m_expiringSubscriptions.Clear(); - } - finally - { - m_semaphoreSlim.Release(); - } + SignalConditionRefreshShutdown(); - foreach (SessionPublishQueue publishQueue in publishQueues) - { - publishQueue?.Dispose(); - } + await Task.WhenAll( + m_publishSubscriptionsTask ?? Task.CompletedTask, + m_conditionRefreshTask ?? Task.CompletedTask) + .ConfigureAwait(false); - foreach (ISubscription subscription in subscriptions) - { - subscription?.Dispose(); - } + await m_semaphoreSlim.WaitAsync(CancellationToken.None).ConfigureAwait(false); + List publishQueues; + List subscriptions; + try + { + CaptureManagedResources(out publishQueues, out subscriptions); + } + finally + { + m_semaphoreSlim.Release(); + } + + DisposeManagedResources(publishQueues, subscriptions); + DisposeSynchronizationResources(); + } + + private void CaptureManagedResources( + out List publishQueues, + out List subscriptions) + { + publishQueues = [.. m_publishQueues.Values]; + m_publishQueues.Clear(); + + subscriptions = [.. m_subscriptions.Values]; + m_subscriptions.Clear(); + m_expiringSubscriptions.Clear(); + } + + private static void DisposeManagedResources( + List publishQueues, + List subscriptions) + { + foreach (SessionPublishQueue publishQueue in publishQueues) + { + publishQueue?.Dispose(); + } - m_shutdownEvent.Dispose(); - m_conditionRefreshEvent.Dispose(); - m_semaphoreSlim.Dispose(); + foreach (ISubscription subscription in subscriptions) + { + subscription?.Dispose(); } } + private void DisposeSynchronizationResources() + { + m_shutdownEvent.Dispose(); + m_conditionRefreshEvent.Dispose(); + m_workerShutdown.Dispose(); + m_semaphoreSlim.Dispose(); + } + /// /// Raised after a new subscription is created. /// @@ -251,23 +304,32 @@ public virtual async ValueTask StartupAsync(CancellationToken cancellationToken await RestoreSubscriptionsAsync(cancellationToken) .ConfigureAwait(false); + if (m_workerShutdown.IsCancellationRequested) + { + m_workerShutdown.Dispose(); + m_workerShutdown = new CancellationTokenSource(); + } + m_shutdownEvent.Reset(); + CancellationToken workerCancellationToken = m_workerShutdown.Token; - // TODO: Ensure shutdown awaits completion and a cancellation token is passed - _ = Task.Factory.StartNew( - () => PublishSubscriptionsAsync(m_publishingResolution), - default, - TaskCreationOptions.LongRunning | TaskCreationOptions.DenyChildAttach, - TaskScheduler.Default); + m_publishSubscriptionsTask = Task.Factory.StartNew( + () => PublishSubscriptionsAsync( + m_publishingResolution, + workerCancellationToken).AsTask(), + default, + TaskCreationOptions.LongRunning | TaskCreationOptions.DenyChildAttach, + TaskScheduler.Default) + .Unwrap(); m_conditionRefreshEvent.Reset(); - // TODO: Ensure shutdown awaits completion and a cancellation token is passed - _ = Task.Factory.StartNew( - ConditionRefreshWorkerAsync, - default, - TaskCreationOptions.LongRunning | TaskCreationOptions.DenyChildAttach, - TaskScheduler.Default); + m_conditionRefreshTask = Task.Factory.StartNew( + ConditionRefreshWorkerAsync, + default, + TaskCreationOptions.LongRunning | TaskCreationOptions.DenyChildAttach, + TaskScheduler.Default) + .Unwrap(); } finally { @@ -280,17 +342,30 @@ await RestoreSubscriptionsAsync(cancellationToken) /// public virtual async ValueTask ShutdownAsync(CancellationToken cancellationToken = default) { + Task publishSubscriptionsTask; + Task conditionRefreshTask; await m_semaphoreSlim.WaitAsync(cancellationToken).ConfigureAwait(false); try { - BeforeConditionRefreshShutdownSignalForTest?.Invoke(); - - // stop the publishing thread. - m_shutdownEvent.Set(); + // stop the publishing and condition refresh workers. + SignalConditionRefreshShutdown(); + publishSubscriptionsTask = + m_publishSubscriptionsTask ?? Task.CompletedTask; + conditionRefreshTask = + m_conditionRefreshTask ?? Task.CompletedTask; + } + finally + { + m_semaphoreSlim.Release(); + } - // trigger the condition refresh thread. - m_conditionRefreshEvent.Set(); + await Task.WhenAll(publishSubscriptionsTask, conditionRefreshTask) + .WaitAsync(cancellationToken) + .ConfigureAwait(false); + await m_semaphoreSlim.WaitAsync(cancellationToken).ConfigureAwait(false); + try + { // dispose of publish queues. foreach (SessionPublishQueue queue in m_publishQueues.Values) { @@ -1505,10 +1580,15 @@ public async ValueTask TransferSubscriptionsAsync } ISession ownerSession = null!; + var concreteSubscription = subscription as Subscription; SessionPublishQueue? sourcePublishQueue = null; + SessionPublishQueue.SubscriptionTransferClaim? sourceQueueClaim = null; bool sourceIsAbandoned = false; bool sourceRemoved = false; - bool transferCompleted = false; + bool transferStarted = false; + Subscription.PreparedSessionTransfer? preparedTransfer = null; + SessionPublishQueue? destinationPublishQueue = null; + bool destinationAdded = false; await m_semaphoreSlim.WaitAsync(cancellationToken).ConfigureAwait(false); try { @@ -1571,15 +1651,15 @@ is MessageSecurityMode.Sign continue; } - // Validate the exact current source while holding the same - // semaphore used by expiration claims. + // Claim the exact current source before any fallible monitored-item + // callback can run. Lock order is manager semaphore, then queue lock, + // then subscription lock; rollback follows the same order. if (ownerSession != null) { if (!m_publishQueues.TryGetValue( ownerSession.Id, out sourcePublishQueue) || - sourcePublishQueue == null || - !sourcePublishQueue.ContainsSubscription(subscription)) + sourcePublishQueue == null) { result.StatusCode = StatusCodes.BadSubscriptionIdInvalid; results.Add(result); @@ -1589,10 +1669,67 @@ is MessageSecurityMode.Sign } continue; } + + if (concreteSubscription != null) + { + if (!sourcePublishQueue.TryClaimForTransfer( + concreteSubscription, + ownerSession, + out sourceQueueClaim)) + { + result.StatusCode = StatusCodes.BadSubscriptionIdInvalid; + results.Add(result); + if ((context.DiagnosticsMask & DiagnosticsMasks.OperationAll) != 0) + { + diagnosticInfos.Add(null!); + } + continue; + } + sourceRemoved = true; + transferStarted = true; + } + else + { + sourceRemoved = sourcePublishQueue.TryRemoveForTransfer(subscription); + if (!sourceRemoved) + { + result.StatusCode = StatusCodes.BadSubscriptionIdInvalid; + results.Add(result); + if ((context.DiagnosticsMask & DiagnosticsMasks.OperationAll) != 0) + { + diagnosticInfos.Add(null!); + } + continue; + } + } } else if (ContainsAbandonedSubscription(subscription)) { sourceIsAbandoned = true; + if (concreteSubscription != null && + !concreteSubscription.TryBeginTransfer(null)) + { + result.StatusCode = StatusCodes.BadSubscriptionIdInvalid; + results.Add(result); + if ((context.DiagnosticsMask & DiagnosticsMasks.OperationAll) != 0) + { + diagnosticInfos.Add(null!); + } + continue; + } + transferStarted = concreteSubscription != null; + sourceRemoved = TryRemoveAbandonedSubscription(subscription); + if (!sourceRemoved) + { + concreteSubscription?.AbortTransfer(null); + result.StatusCode = StatusCodes.BadSubscriptionIdInvalid; + results.Add(result); + if ((context.DiagnosticsMask & DiagnosticsMasks.OperationAll) != 0) + { + diagnosticInfos.Add(null!); + } + continue; + } } else if (m_abandonedSubscriptions.ContainsKey(subscription.Id)) { @@ -1604,85 +1741,138 @@ is MessageSecurityMode.Sign } continue; } - - try + else if (concreteSubscription != null) { - // transfer session, add subscription to publish queue - await subscription.TransferSessionAsync( - context, - sendInitialValues, - cancellationToken) - .ConfigureAwait(false); - - if (sourcePublishQueue != null) + if (!concreteSubscription.TryBeginTransfer(null)) { - sourceRemoved = - sourcePublishQueue.TryRemoveForTransfer(subscription); + result.StatusCode = StatusCodes.BadSubscriptionIdInvalid; + results.Add(result); + if ((context.DiagnosticsMask & DiagnosticsMasks.OperationAll) != 0) + { + diagnosticInfos.Add(null!); + } + continue; } - else if (sourceIsAbandoned) + transferStarted = true; + } + + try + { + if (concreteSubscription != null) { - sourceRemoved = - TryRemoveAbandonedSubscription(subscription); + preparedTransfer = await concreteSubscription + .PrepareSessionTransferAsync( + context, + ownerSession, + sendInitialValues, + cancellationToken) + .ConfigureAwait(false); + preparedTransfer.CommitOwnership(); } - - if ((sourcePublishQueue != null || sourceIsAbandoned) && - !sourceRemoved) + else { - throw new ServiceResultException( - StatusCodes.BadSubscriptionIdInvalid, - "Subscription source changed during transfer."); + await subscription.TransferSessionAsync( + context, + sendInitialValues, + cancellationToken) + .ConfigureAwait(false); } // add to queue in new session, create queue if necessary if (!m_publishQueues.TryGetValue( context.SessionId, - out SessionPublishQueue? publishQueue) || - publishQueue == null) + out destinationPublishQueue) || + destinationPublishQueue == null) { m_publishQueues[context.SessionId] - = publishQueue = new SessionPublishQueue( + = destinationPublishQueue = new SessionPublishQueue( m_server, context.Session, m_maxPublishRequestCount, m_timeProvider); } - publishQueue.Add(subscription); - transferCompleted = true; + destinationPublishQueue.Add(subscription); + destinationAdded = true; + preparedTransfer?.CommitMonitoredItemEffects(); + if (concreteSubscription != null) + { + concreteSubscription.CompleteTransfer(context.Session); + } + if (sourceQueueClaim != null) + { + sourcePublishQueue!.CompleteTransferClaim(sourceQueueClaim); + } + preparedTransfer?.Complete(); } - finally + catch (Exception transferError) { - if (!transferCompleted) + var rollbackErrors = new List(); + if (destinationAdded && destinationPublishQueue != null) + { + destinationPublishQueue.TryRemoveForTransfer(subscription); + destinationPublishQueue.RemoveQueuedRequests(); + } + + if (preparedTransfer != null) { - bool ownershipRestored = - ReferenceEquals(subscription.Session, ownerSession); - if (!ownershipRestored && subscription is Subscription concrete) + try { - ownershipRestored = - concrete.TryRestoreSessionAfterFailedTransfer( - context.Session, - ownerSession); + await preparedTransfer.RollbackAsync(CancellationToken.None) + .ConfigureAwait(false); } - - if (ownershipRestored) + catch (Exception rollbackError) { - if (sourceRemoved && - sourcePublishQueue != null && - ownerSession != null && - m_publishQueues.TryGetValue( - ownerSession.Id, - out SessionPublishQueue? currentOwnerQueue) && - ReferenceEquals(currentOwnerQueue, sourcePublishQueue)) - { - sourcePublishQueue.Add(subscription); - } - else if (sourceRemoved && sourceIsAbandoned) - { - m_abandonedSubscriptions.TryAdd( - subscription.Id, - subscription); - } + rollbackErrors.Add(rollbackError); } } + else if (!ReferenceEquals(subscription.Session, ownerSession) && + concreteSubscription != null && + !concreteSubscription.TryRestoreSessionAfterFailedTransfer( + context.Session, + ownerSession)) + { + rollbackErrors.Add( + new ServiceResultException( + StatusCodes.BadSubscriptionIdInvalid, + "Subscription ownership could not be restored.")); + } + + if (sourceQueueClaim != null && + !sourcePublishQueue!.RestoreTransferClaim(sourceQueueClaim)) + { + rollbackErrors.Add( + new ServiceResultException( + StatusCodes.BadSubscriptionIdInvalid, + "Subscription source queue could not be restored.")); + } + else if (sourceQueueClaim == null && + sourceRemoved && + sourcePublishQueue != null) + { + sourcePublishQueue.Add(subscription); + } + else if (sourceRemoved && sourceIsAbandoned && + !m_abandonedSubscriptions.TryAdd( + subscription.Id, + subscription)) + { + rollbackErrors.Add( + new ServiceResultException( + StatusCodes.BadSubscriptionIdInvalid, + "Abandoned subscription source could not be restored.")); + } + + if (transferStarted) + { + concreteSubscription!.AbortTransfer(ownerSession); + } + + if (rollbackErrors.Count > 0) + { + rollbackErrors.Insert(0, transferError); + throw new AggregateException(rollbackErrors); + } + throw; } } finally @@ -2320,6 +2510,10 @@ private async ValueTask PublishSubscriptionsAsync(int sleepCycle, CancellationTo { m_logger.SubscriptionPublishTaskTaskIdX8ExitedNormally2(Task.CurrentId); } + catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) + { + m_logger.SubscriptionPublishTaskTaskIdX8ExitedNormally2(Task.CurrentId); + } catch (Exception e) { m_logger.SubscriptionPublishTaskTaskIdX8ExitedUnexpectedly(e, Task.CurrentId); @@ -2349,7 +2543,8 @@ internal void ProcessAbandonedPublishTimers( for (int ii = 0; ii < abandonedSubscriptions.Count; ii++) { ISubscription subscription = abandonedSubscriptions[ii]; - if (subscription.PublishTimerExpired() != PublishingState.Expired || + if (!ContainsAbandonedSubscription(subscription) || + subscription.PublishTimerExpired() != PublishingState.Expired || !TryClaimAbandonedSubscriptionExpiration(subscription)) { continue; @@ -2371,29 +2566,48 @@ private async Task ConditionRefreshWorkerAsync() try { m_logger.SubscriptionConditionRefreshTaskTaskIdX8Started(Task.CurrentId); + WaitHandle[] waitHandles = [m_conditionRefreshEvent, m_shutdownEvent]; while (true) { ConditionRefreshTask? conditionRefreshTask = null; + bool shutdown; lock (m_conditionRefreshLock) { - if (m_conditionRefreshQueue.Count > 0) + shutdown = m_shutdownEvent.WaitOne(0); + if (!shutdown && m_conditionRefreshQueue.Count > 0) { conditionRefreshTask = m_conditionRefreshQueue.Dequeue(); } - else + else if (!shutdown) { BeforeConditionRefreshResetForTest?.Invoke(); - m_conditionRefreshEvent.Reset(); + shutdown = m_shutdownEvent.WaitOne(0); + if (!shutdown) + { + m_conditionRefreshEvent.Reset(); + } } } + if (shutdown) + { + m_logger.SubscriptionConditionRefreshTaskTaskIdX8Exited(Task.CurrentId); + break; + } + if (conditionRefreshTask == null) { - m_conditionRefreshEvent.WaitOne(); + if (WaitHandle.WaitAny(waitHandles) == 1) + { + m_logger.SubscriptionConditionRefreshTaskTaskIdX8Exited(Task.CurrentId); + break; + } + continue; } - else if (conditionRefreshTask.MonitoredItemId == 0) + + if (conditionRefreshTask.MonitoredItemId == 0) { await DoConditionRefreshAsync(conditionRefreshTask.Subscription) .ConfigureAwait(false); @@ -2406,12 +2620,6 @@ await DoConditionRefresh2Async( .ConfigureAwait(false); } - // use shutdown event to end loop - if (m_shutdownEvent.WaitOne(0)) - { - m_logger.SubscriptionConditionRefreshTaskTaskIdX8Exited(Task.CurrentId); - break; - } } } catch (ObjectDisposedException) @@ -2424,6 +2632,17 @@ await DoConditionRefresh2Async( } } + private void SignalConditionRefreshShutdown() + { + BeforeConditionRefreshShutdownSignalForTest?.Invoke(); + lock (m_conditionRefreshLock) + { + m_shutdownEvent.Set(); + m_workerShutdown.Cancel(); + m_conditionRefreshEvent.Set(); + } + } + /// /// Cleanups the subscriptions. /// @@ -2537,16 +2756,11 @@ public override int GetHashCode() private readonly ManualResetEvent m_shutdownEvent; private readonly Queue m_conditionRefreshQueue; private readonly ManualResetEvent m_conditionRefreshEvent; - private readonly ISubscriptionStore m_subscriptionStore; - - private readonly Lock m_statusMessagesLock = new(); - private readonly Lock m_eventLock = new(); - private readonly Lock m_conditionRefreshLock = new(); - private event SubscriptionEventHandler? m_SubscriptionCreated; - private event SubscriptionEventHandler? m_SubscriptionDeleted; - + private CancellationTokenSource m_workerShutdown; + private Task? m_publishSubscriptionsTask; + private Task? m_conditionRefreshTask; + private int m_disposed; internal Action? BeforeConditionRefreshResetForTest { get; set; } - internal Action? BeforeConditionRefreshShutdownSignalForTest { get; set; } internal void WakeConditionRefreshWorkerForTest() @@ -2558,12 +2772,21 @@ internal void StartConditionRefreshWorkerForTest() { m_shutdownEvent.Reset(); m_conditionRefreshEvent.Reset(); - _ = Task.Factory.StartNew( - ConditionRefreshWorkerAsync, - default, - TaskCreationOptions.LongRunning | TaskCreationOptions.DenyChildAttach, - TaskScheduler.Default); + m_conditionRefreshTask = Task.Factory.StartNew( + ConditionRefreshWorkerAsync, + default, + TaskCreationOptions.LongRunning | TaskCreationOptions.DenyChildAttach, + TaskScheduler.Default) + .Unwrap(); } + + private readonly ISubscriptionStore m_subscriptionStore; + + private readonly Lock m_statusMessagesLock = new(); + private readonly Lock m_eventLock = new(); + private readonly Lock m_conditionRefreshLock = new(); + private event SubscriptionEventHandler? m_SubscriptionCreated; + private event SubscriptionEventHandler? m_SubscriptionDeleted; } /// diff --git a/tests/Opc.Ua.Server.Tests/NodeManager/MasterNodeManagerDeterministicTests.cs b/tests/Opc.Ua.Server.Tests/NodeManager/MasterNodeManagerDeterministicTests.cs index 63c84e033f..03d056ad7a 100644 --- a/tests/Opc.Ua.Server.Tests/NodeManager/MasterNodeManagerDeterministicTests.cs +++ b/tests/Opc.Ua.Server.Tests/NodeManager/MasterNodeManagerDeterministicTests.cs @@ -1237,7 +1237,6 @@ await sut.TransferMonitoredItemsAsync( [Test] - [Ignore("Resend trigger rollback is covered by the monitored-item transfer stack slice.")] public async Task TransferMonitoredItemsAsyncFailureRollsBackOwnerAndResendAsync() { var sourceSession = new Mock(); diff --git a/tests/Opc.Ua.Server.Tests/SubscriptionTests.cs b/tests/Opc.Ua.Server.Tests/SubscriptionTests.cs index 943dcfae14..b35eab819b 100644 --- a/tests/Opc.Ua.Server.Tests/SubscriptionTests.cs +++ b/tests/Opc.Ua.Server.Tests/SubscriptionTests.cs @@ -408,6 +408,27 @@ private static void SetPrivateField( field.SetValue(instance, value); } + private static void EnqueueConditionRefreshTask( + SubscriptionManager manager, + ISubscription subscription) + { + Type taskType = typeof(SubscriptionManager).GetNestedType( + "ConditionRefreshTask", + BindingFlags.NonPublic) + ?? throw new InvalidOperationException( + "ConditionRefreshTask type not found."); + object task = Activator.CreateInstance(taskType, subscription, 0u) + ?? throw new InvalidOperationException( + "ConditionRefreshTask could not be created."); + object queue = GetPrivateField( + manager, + "m_conditionRefreshQueue"); + MethodInfo enqueue = queue.GetType().GetMethod("Enqueue") + ?? throw new InvalidOperationException( + "ConditionRefreshTask queue Enqueue method not found."); + enqueue.Invoke(queue, [task]); + } + private static void ExpireOnNextPublishTimer(Subscription subscription) { uint maxLifetimeCount = GetPrivateField( @@ -933,6 +954,82 @@ public async Task CreateSubscriptionWithNotificationLimitsUsesEffectiveLimitAsyn Is.EqualTo((uint)expectedLimit)); } + [Test] + public async Task DisposeAsyncJoinsWorkersAndIsIdempotentAsync() + { + var timeProvider = new FakeTimeProvider( + new DateTimeOffset(2026, 1, 1, 0, 0, 0, TimeSpan.Zero)); + var configuration = new ApplicationConfiguration + { + ServerConfiguration = new ServerConfiguration + { + PublishingResolution = 1000 + } + }; + var manager = new SubscriptionManager( + m_serverMock.Object, + configuration, + timeProvider); + + await manager.StartupAsync().ConfigureAwait(false); + await Task.Yield(); + + await manager.DisposeAsync().ConfigureAwait(false); + await manager.DisposeAsync().ConfigureAwait(false); + + Assert.Multiple(() => + { + Assert.That( + GetPrivateField(manager, "m_publishSubscriptionsTask").IsCompleted, + Is.True); + Assert.That( + GetPrivateField(manager, "m_conditionRefreshTask").IsCompleted, + Is.True); + }); + } + + [Test] + public async Task DisposeAsyncWhileConditionRefreshIsRunningWaitsForWorkAsync() + { + var configuration = new ApplicationConfiguration + { + ServerConfiguration = new ServerConfiguration() + }; + var manager = new SubscriptionManager( + m_serverMock.Object, + configuration); + var refreshStarted = new TaskCompletionSource( + TaskCreationOptions.RunContinuationsAsynchronously); + var releaseRefresh = new TaskCompletionSource( + TaskCreationOptions.RunContinuationsAsynchronously); + bool refreshCompleted = false; + var subscription = new Mock(); + subscription + .Setup(s => s.ConditionRefreshAsync(It.IsAny())) + .Returns((CancellationToken _) => new ValueTask(WaitForReleaseAsync())); + + async Task WaitForReleaseAsync() + { + refreshStarted.SetResult(new object()); + await releaseRefresh.Task.ConfigureAwait(false); + refreshCompleted = true; + } + + EnqueueConditionRefreshTask(manager, subscription.Object); + manager.StartConditionRefreshWorkerForTest(); + await refreshStarted.Task.ConfigureAwait(false); + + Task disposeTask = manager.DisposeAsync().AsTask(); + + Assert.That(disposeTask.IsCompleted, Is.False); + + releaseRefresh.SetResult(new object()); + await disposeTask.ConfigureAwait(false); + await manager.DisposeAsync().ConfigureAwait(false); + + Assert.That(refreshCompleted, Is.True); + } + [TestCase(0, 0L, 0L)] [TestCase(0, 25L, 25L)] [TestCase(10, 0L, 10L)] @@ -1088,6 +1185,410 @@ public async Task TransferSerializesWithClosingSourceSessionAsync() Assert.DoesNotThrow(() => subscription.ResendData(destinationContext)); } + [Test] + public async Task TransferClaimsSourceBeforeCallbacksAndBlocksStalePublishAsync() + { + var fixture = await CreateTransferSubscriptionAsync().ConfigureAwait(false); + using SubscriptionManager manager = fixture.Manager; + Subscription subscription = fixture.Subscription; + SessionPublishQueue sourceQueue = GetPublishQueue( + manager, + fixture.SourceContext.SessionId); + SetPrivateField(subscription, "m_waitingForPublish", false); + SetPrivateField(subscription, "m_lifetimeCounter", 0u); + SetExpiryTime( + subscription, + TimeProvider.System.GetTimestampMilliseconds() - 100); + m_sessionMock + .Setup(session => session.IsSecureChannelValid("source-channel")) + .Returns(true); + IReadOnlyList staleSnapshot = + sourceQueue.CapturePublishTimerSnapshot(); + using var publishCancellation = new CancellationTokenSource(); + Task sourcePublish = sourceQueue.PublishAsync( + "source-channel", + DateTime.MaxValue, + requeue: false, + parkSink: null, + publishCancellation.Token); + var transferEntered = new TaskCompletionSource( + TaskCreationOptions.RunContinuationsAsynchronously); + var releaseTransfer = new TaskCompletionSource( + TaskCreationOptions.RunContinuationsAsynchronously); + m_nodeManagerMock + .Setup(nodeManager => nodeManager.TransferMonitoredItemsAsync( + It.IsAny(), + It.IsAny(), + It.IsAny>(), + It.IsAny>(), + It.IsAny(), + It.IsAny())) + .Callback(() => transferEntered.TrySetResult(true)) + .Returns(() => new ValueTask(releaseTransfer.Task)); + + Task transferTask = manager + .TransferSubscriptionsAsync( + fixture.DestinationContext, + [subscription.Id], + sendInitialValues: true) + .AsTask(); + Task entered = await Task.WhenAny( + transferEntered.Task, + Task.Delay(TimeSpan.FromSeconds(5))).ConfigureAwait(false); + if (!ReferenceEquals(entered, transferEntered.Task)) + { + releaseTransfer.TrySetResult(true); + Assert.Fail("Transfer did not reach the monitored-item callback barrier."); + } + + try + { + sourceQueue.PublishTimerExpired(staleSnapshot); + + Assert.Multiple(() => + { + Assert.That( + sourceQueue.CapturePublishTimerSnapshot(), + Is.Empty, + "The source entry must be claimed before callbacks run."); + Assert.That( + sourcePublish.IsCompleted, + Is.False, + "A stale source timer must not assign the claimed subscription."); + ServiceResultException error = Assert.Throws( + () => subscription.Publish( + fixture.SourceContext, + out _, + out _)); + Assert.That( + error.StatusCode, + Is.EqualTo(StatusCodes.BadSubscriptionIdInvalid)); + }); + } + finally + { + releaseTransfer.TrySetResult(true); + } + + TransferSubscriptionsResponse transferred = await transferTask.ConfigureAwait(false); + publishCancellation.Cancel(); + try + { + await sourcePublish.ConfigureAwait(false); + } + catch (Exception) + { + // The old session's parked Publish is cancelled or completed by + // the transferred-status notification, never with the subscription. + } + + Assert.Multiple(() => + { + Assert.That(transferred.Results, Has.Count.EqualTo(1)); + Assert.That(transferred.Results[0].StatusCode, Is.EqualTo(StatusCodes.Good)); + Assert.That( + subscription.Session, + Is.SameAs(fixture.DestinationSession.Object)); + }); + Assert.DoesNotThrow( + () => subscription.ResendData(fixture.DestinationContext)); + } + + [Test] + public async Task TransferInitialValueIsPublishedOnlyByDestinationAsync() + { + var fixture = await CreateTransferSubscriptionAsync().ConfigureAwait(false); + using SubscriptionManager manager = fixture.Manager; + Subscription subscription = fixture.Subscription; + using var queueFactory = new MonitoredItemQueueFactory(m_telemetry); + m_serverMock.Setup(server => server.MonitoredItemQueueFactory).Returns(queueFactory); + var itemOwner = new Mock(); + using var monitoredItem = new MonitoredItem( + m_serverMock.Object, + itemOwner.Object, + new object(), + subscription.Id, + id: 10, + new ReadValueId + { + NodeId = new NodeId("TransferValue", 2), + AttributeId = Attributes.Value + }, + DiagnosticsMasks.None, + TimestampsToReturn.Both, + MonitoringMode.Reporting, + clientHandle: 11, + originalFilter: null, + filterToUse: null, + range: null, + samplingInterval: 1000, + queueSize: 1, + discardOldest: true, + sourceSamplingInterval: 1000); + await RegisterMonitoredItemsAsync(subscription, monitoredItem).ConfigureAwait(false); + monitoredItem.QueueValue(new DataValue(new Variant(1234)), null); + var initialNotifications = new Queue(); + monitoredItem.Publish( + new OperationContext(monitoredItem), + initialNotifications, + new Queue(), + 1, + m_telemetry.CreateLogger()); + Assert.That(initialNotifications, Has.Count.EqualTo(1)); + + var transferEntered = new TaskCompletionSource( + TaskCreationOptions.RunContinuationsAsynchronously); + var releaseTransfer = new TaskCompletionSource( + TaskCreationOptions.RunContinuationsAsynchronously); + itemOwner + .Setup(nodeManager => nodeManager.TransferMonitoredItemsAsync( + It.IsAny(), + It.IsAny(), + It.IsAny>(), + It.IsAny>(), + It.IsAny>(), + It.IsAny(), + It.IsAny())) + .Returns< + OperationContext, + bool, + IList, + IList, + IList, + MonitoredItemTransferOptions, + CancellationToken>(async (_, sendInitialValues, monitoredItems, processedItems, errors, transferOptions, cancellationToken) => + { + for (int ii = 0; ii < monitoredItems.Count; ii++) + { + if (processedItems[ii]) + { + continue; + } + + processedItems[ii] = true; + errors[ii] = ServiceResult.Good; + if (sendInitialValues && !transferOptions.DeferInitialValues) + { + monitoredItems[ii].SetupResendDataTrigger(); + } + } + transferEntered.TrySetResult(true); + await releaseTransfer.Task.WaitAsync(cancellationToken).ConfigureAwait(false); + }); + itemOwner + .Setup(nodeManager => nodeManager.RollbackMonitoredItemsTransferAsync( + It.IsAny(), + It.IsAny>(), + It.IsAny>(), + It.IsAny>(), + It.IsAny(), + It.IsAny())) + .Returns(default(ValueTask)); + var configurationNodeManager = new Mock(); + configurationNodeManager + .SetupGet(nodeManager => nodeManager.NamespaceUris) + .Returns(System.Array.Empty()); + var coreNodeManager = new Mock(); + coreNodeManager + .SetupGet(nodeManager => nodeManager.NamespaceUris) + .Returns(System.Array.Empty()); + var factory = new Mock(); + factory + .Setup(nodeManagerFactory => nodeManagerFactory.CreateConfigurationNodeManager()) + .Returns(configurationNodeManager.Object); + factory + .Setup(nodeManagerFactory => nodeManagerFactory.CreateCoreNodeManager(It.IsAny())) + .Returns(coreNodeManager.Object); + m_serverMock.Setup(server => server.MainNodeManagerFactory).Returns(factory.Object); + using var masterNodeManager = new MasterNodeManager( + m_serverMock.Object, + new ApplicationConfiguration { ServerConfiguration = new ServerConfiguration() }, + null, + itemOwner.Object); + m_serverMock.Setup(server => server.NodeManager).Returns(masterNodeManager); + + Task transferTask = manager + .TransferSubscriptionsAsync( + fixture.DestinationContext, + [subscription.Id], + sendInitialValues: true) + .AsTask(); + Task entered = await Task.WhenAny( + transferEntered.Task, + Task.Delay(TimeSpan.FromSeconds(5))).ConfigureAwait(false); + if (!ReferenceEquals(entered, transferEntered.Task)) + { + releaseTransfer.TrySetResult(true); + Assert.Fail("Transfer did not reach the monitored-item callback barrier."); + } + + try + { + itemOwner.Verify( + nodeManager => nodeManager.TransferMonitoredItemsAsync( + It.IsAny(), + true, + It.IsAny>(), + It.IsAny>(), + It.IsAny>(), + It.Is(options => options.DeferInitialValues), + It.IsAny()), + Times.Once); + Assert.That(monitoredItem.IsResendData, Is.False); + ServiceResultException error = Assert.Throws( + () => subscription.Publish( + fixture.SourceContext, + out _, + out _)); + Assert.That(error.StatusCode, Is.EqualTo(StatusCodes.BadSubscriptionIdInvalid)); + } + finally + { + releaseTransfer.TrySetResult(true); + } + + TransferSubscriptionsResponse transferred = await transferTask.ConfigureAwait(false); + Assert.That(transferred.Results[0].StatusCode, Is.EqualTo(StatusCodes.Good)); + Assert.That(monitoredItem.IsResendData, Is.True); + + SessionPublishQueue destinationQueue = GetPublishQueue( + manager, + fixture.DestinationContext.SessionId); + SetPrivateField(subscription, "m_waitingForPublish", false); + SetPrivateField(subscription, "m_lifetimeCounter", 0u); + SetExpiryTime( + subscription, + TimeProvider.System.GetTimestampMilliseconds() - 100); + destinationQueue.PublishTimerExpired( + destinationQueue.CapturePublishTimerSnapshot()); + ISubscription readySubscription = await destinationQueue.PublishAsync( + "destination-channel", + DateTime.MaxValue, + requeue: false, + parkSink: null, + CancellationToken.None).ConfigureAwait(false); + NotificationMessage message = readySubscription.Publish( + fixture.DestinationContext, + out _, + out _); + + Assert.Multiple(() => + { + Assert.That(readySubscription, Is.SameAs(subscription)); + Assert.That(message, Is.Not.Null); + Assert.That(message.NotificationData, Is.Not.Empty); + Assert.That(monitoredItem.IsResendData, Is.False); + }); + } + + [Test] + public async Task TransferFallbackAppliesInitialValueExactlyOnceAsync() + { + var fixture = await CreateTransferSubscriptionAsync().ConfigureAwait(false); + using SubscriptionManager manager = fixture.Manager; + Subscription subscription = fixture.Subscription; + using var queueFactory = new MonitoredItemQueueFactory(m_telemetry); + m_serverMock.Setup(server => server.MonitoredItemQueueFactory).Returns(queueFactory); + var itemOwner = new Mock(); + using var monitoredItem = new MonitoredItem( + m_serverMock.Object, + itemOwner.Object, + new object(), + subscription.Id, + id: 11, + new ReadValueId + { + NodeId = new NodeId("FallbackTransferValue", 2), + AttributeId = Attributes.Value + }, + DiagnosticsMasks.None, + TimestampsToReturn.Both, + MonitoringMode.Reporting, + clientHandle: 12, + originalFilter: null, + filterToUse: null, + range: null, + samplingInterval: 1000, + queueSize: 1, + discardOldest: true, + sourceSamplingInterval: 1000); + await RegisterMonitoredItemsAsync(subscription, monitoredItem).ConfigureAwait(false); + monitoredItem.QueueValue(new DataValue(new Variant(5678)), null); + var initialNotifications = new Queue(); + monitoredItem.Publish( + new OperationContext(monitoredItem), + initialNotifications, + new Queue(), + 1, + m_telemetry.CreateLogger()); + Assert.That(initialNotifications, Has.Count.EqualTo(1)); + int setupResendCalls = 0; + m_nodeManagerMock + .Setup(nodeManager => nodeManager.TransferMonitoredItemsAsync( + It.IsAny(), + true, + It.Is>(items => items.Contains(monitoredItem)), + It.IsAny>(), + It.Is(options => options == null), + It.IsAny())) + .Callback, IList, MonitoredItemTransferOptions, CancellationToken>( + (_, sendInitialValues, monitoredItems, errors, transferOptions, _) => + { + Assert.That(transferOptions, Is.Null); + for (int ii = 0; ii < monitoredItems.Count; ii++) + { + errors[ii] = ServiceResult.Good; + if (sendInitialValues) + { + setupResendCalls++; + monitoredItems[ii].SetupResendDataTrigger(); + } + } + }) + .Returns(default(ValueTask)); + + TransferSubscriptionsResponse transferred = await manager + .TransferSubscriptionsAsync( + fixture.DestinationContext, + [subscription.Id], + sendInitialValues: true) + .ConfigureAwait(false); + + Assert.That(transferred.Results[0].StatusCode, Is.EqualTo(StatusCodes.Good)); + Assert.That(setupResendCalls, Is.EqualTo(1)); + SessionPublishQueue destinationQueue = GetPublishQueue( + manager, + fixture.DestinationContext.SessionId); + SetPrivateField(subscription, "m_waitingForPublish", false); + SetPrivateField(subscription, "m_lifetimeCounter", 0u); + SetExpiryTime( + subscription, + TimeProvider.System.GetTimestampMilliseconds() - 100); + destinationQueue.PublishTimerExpired( + destinationQueue.CapturePublishTimerSnapshot()); + ISubscription readySubscription = await destinationQueue.PublishAsync( + "destination-channel", + DateTime.MaxValue, + requeue: false, + parkSink: null, + CancellationToken.None).ConfigureAwait(false); + NotificationMessage firstMessage = readySubscription.Publish( + fixture.DestinationContext, + out _, + out _); + NotificationMessage secondMessage = readySubscription.Publish( + fixture.DestinationContext, + out _, + out _); + + Assert.Multiple(() => + { + Assert.That(firstMessage.NotificationData, Is.Not.Empty); + Assert.That(secondMessage?.NotificationData ?? [], Is.Empty); + Assert.That(monitoredItem.IsResendData, Is.False); + }); + } + [Test] public async Task SourcePublishTimerSnapshotDoesNotExpireTransferredSubscriptionAsync() { diff --git a/tests/Opc.Ua.Subscriptions.Tests/SessionPublishQueueConcurrencyTests.cs b/tests/Opc.Ua.Subscriptions.Tests/SessionPublishQueueConcurrencyTests.cs new file mode 100644 index 0000000000..5da528ff3e --- /dev/null +++ b/tests/Opc.Ua.Subscriptions.Tests/SessionPublishQueueConcurrencyTests.cs @@ -0,0 +1,267 @@ +/* ======================================================================== + * Copyright (c) 2005-2025 The OPC Foundation, Inc. All rights reserved. + * + * OPC Foundation MIT License 1.00 + * + * Permission is hereby granted, free of charge, to any person + * obtaining a copy of this software and associated documentation + * files (the "Software"), to deal in the Software without + * restriction, including without limitation the rights to use, + * copy, modify, merge, publish, distribute, sublicense, and/or sell + * copies of the Software, and to permit persons to whom the + * Software is furnished to do so, subject to the following + * conditions: + * + * The above copyright notice and this permission notice shall be + * included in all copies or substantial portions of the Software. + * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, + * EXPRESS OR IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES + * OF MERCHANTABILITY, FITNESS FOR A PARTICULAR PURPOSE AND + * NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR COPYRIGHT + * HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, + * WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING + * FROM, OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR + * OTHER DEALINGS IN THE SOFTWARE. + * + * The complete license agreement can be found here: + * http://opcfoundation.org/License/MIT/1.00/ + * ======================================================================*/ + +using System; +using System.Collections.Generic; +using System.Reflection; +using System.Runtime.CompilerServices; +using System.Threading; +using System.Threading.Tasks; +using Microsoft.Extensions.Logging.Abstractions; +using Moq; +using NUnit.Framework; +using Opc.Ua.Server; + +namespace Opc.Ua.Subscriptions.Tests +{ + /// + /// Covers queue races around publish, requeue, expiration and transfer claims. + /// + [TestFixture] + public class SessionPublishQueueConcurrencyTests + { + [Test] + public async Task ConcurrentPublishRequeueAndExpirationDoesNotDuplicateSubscriptionAsync() + { + SessionPublishQueue queue = CreateQueue(out Mock session); + Mock subscription = CreateSubscription(1, session.Object); + queue.Add(subscription.Object); + IReadOnlyList snapshot = queue.CapturePublishTimerSnapshot(); + + using var start = new Barrier(3); + using var cancellationTokenSource = new CancellationTokenSource(); + Task publishTask = Task.Run( + async () => + { + start.SignalAndWait(); + try + { + return await queue.PublishAsync( + "secure-channel", + DateTime.MaxValue, + requeue: false, + parkSink: null, + cancellationTokenSource.Token) + .ConfigureAwait(false); + } + catch (OperationCanceledException) + { + return null; + } + catch (ServiceResultException exception) + when (exception.StatusCode == StatusCodes.BadNoSubscription) + { + return null; + } + }); + Task requeueTask = Task.Run( + () => + { + start.SignalAndWait(); + queue.Requeue(subscription.Object); + }); + Task expirationTask = Task.Run( + () => + { + start.SignalAndWait(); + return queue.TryRemoveForExpiration(snapshot[0]); + }); + + await Task.WhenAll(requeueTask, expirationTask).ConfigureAwait(false); + cancellationTokenSource.Cancel(); + ISubscription? publishedSubscription = await publishTask.ConfigureAwait(false); + + Assert.That( + publishedSubscription == null || ReferenceEquals(publishedSubscription, subscription.Object), + Is.True); + + if (queue.ContainsSubscription(subscription.Object)) + { + queue.Requeue(subscription.Object); + Task firstPublish = queue.PublishAsync( + "secure-channel", + DateTime.MaxValue, + requeue: false, + parkSink: null, + CancellationToken.None); + Task secondPublish = queue.PublishAsync( + "secure-channel", + DateTime.MaxValue, + requeue: false, + parkSink: null, + CancellationToken.None); + + // Status == RanToCompletion rather than Task.IsCompletedSuccessfully, + // which is .NET 5+ only and would break the net48/net472 builds. + Assert.That(firstPublish.Status, Is.EqualTo(TaskStatus.RanToCompletion)); + Assert.That(await firstPublish.ConfigureAwait(false), Is.SameAs(subscription.Object)); + Assert.That(secondPublish.Status, Is.Not.EqualTo(TaskStatus.RanToCompletion)); + + queue.Remove(subscription.Object, removeQueuedRequests: true); + } + else + { + Assert.That(await expirationTask.ConfigureAwait(false), Is.True); + } + } + + [Test] + public async Task TransferClaimProtectsClaimedEntryFromExpirationAndPreservesRequeueAsync() + { + SessionPublishQueue queue = CreateQueue(out Mock session); + Subscription subscription = CreateConcreteSubscription(2, session.Object); + queue.Add(subscription); + IReadOnlyList snapshot = queue.CapturePublishTimerSnapshot(); + + bool claimed = queue.TryClaimForTransfer( + subscription, + session.Object, + out SessionPublishQueue.SubscriptionTransferClaim? claim); + + Assert.That(claimed, Is.True); + Assert.That(claim, Is.Not.Null); + + using var start = new Barrier(3); + Task expirationTask = Task.Run( + () => + { + start.SignalAndWait(); + return queue.TryRemoveForExpiration(snapshot[0]); + }); + Task secondClaimTask = Task.Run( + () => + { + start.SignalAndWait(); + return queue.TryClaimForTransfer(subscription, session.Object, out _); + }); + Task requeueTask = Task.Run( + () => + { + start.SignalAndWait(); + queue.Requeue(subscription); + }); + + await Task.WhenAll(expirationTask, secondClaimTask, requeueTask).ConfigureAwait(false); + + Assert.That(await expirationTask.ConfigureAwait(false), Is.False); + Assert.That(await secondClaimTask.ConfigureAwait(false), Is.False); + Assert.That(queue.RestoreTransferClaim(claim!), Is.True); + + Task publishTask = queue.PublishAsync( + "secure-channel", + DateTime.MaxValue, + requeue: false, + parkSink: null, + CancellationToken.None); + + Assert.That(publishTask.Status, Is.EqualTo(TaskStatus.RanToCompletion)); + Assert.That(await publishTask.ConfigureAwait(false), Is.SameAs(subscription)); + } + + private static SessionPublishQueue CreateQueue(out Mock session) + { + return CreateQueue(out session, out _); + } + + private static SessionPublishQueue CreateQueue( + out Mock session, + out Mock server) + { + var telemetry = new TestTelemetryContext(); + server = new Mock(MockBehavior.Loose); + server.SetupGet(value => value.Telemetry).Returns(telemetry); + + session = new Mock(MockBehavior.Loose); + session.SetupGet(value => value.Id).Returns(new NodeId(Guid.NewGuid())); + session.SetupGet(value => value.Identity).Returns(new UserIdentity()); + session.SetupGet(value => value.IdentityToken).Returns(CreateIdentityToken().Object); + session.SetupGet(value => value.SessionDiagnostics).Returns( + new SessionDiagnosticsDataType + { + ClientDescription = new ApplicationDescription + { + ApplicationUri = "urn:test" + } + }); + session.Setup(value => value.IsSecureChannelValid(It.IsAny())).Returns(true); + + return new SessionPublishQueue(server.Object, session.Object, 10, TimeProvider.System); + } + + private static Mock CreateIdentityToken() + { + var identityToken = new Mock(MockBehavior.Loose); + identityToken.SetupGet(value => value.TokenType).Returns(UserTokenType.Anonymous); + identityToken.SetupGet(value => value.Token).Returns(new AnonymousIdentityToken()); + return identityToken; + } + + private static Mock CreateSubscription(uint id, ISession session) + { + var subscription = new Mock(MockBehavior.Loose); + subscription.SetupGet(value => value.Id).Returns(id); + subscription.SetupGet(value => value.Session).Returns(session); + subscription.SetupGet(value => value.SessionId).Returns(session.Id); + subscription.SetupGet(value => value.Priority).Returns(0); + return subscription; + } + + private static Subscription CreateConcreteSubscription(uint id, ISession session) + { + // RuntimeHelpers.GetUninitializedObject is .NET 5+ only; net48/net472 have the + // equivalent on FormatterServices, which is obsolete on modern targets. +#if NET5_0_OR_GREATER + var subscription = (Subscription)RuntimeHelpers.GetUninitializedObject(typeof(Subscription)); +#else + var subscription = (Subscription)System.Runtime.Serialization.FormatterServices + .GetUninitializedObject(typeof(Subscription)); +#endif + SetField(subscription, "m_lock", new Lock()); + SetField(subscription, "k__BackingField", id); + SetField(subscription, "k__BackingField", session); + SetField(subscription, "k__BackingField", (byte)0); + return subscription; + } + + private static void SetField(Subscription subscription, string fieldName, TValue value) + { + typeof(Subscription) + .GetField(fieldName, BindingFlags.Instance | BindingFlags.NonPublic)! + .SetValue(subscription, value); + } + + private sealed class TestTelemetryContext : TelemetryContextBase + { + public TestTelemetryContext() + : base(NullLoggerFactory.Instance) + { + } + } + } +} From 41d353d9dd54f148925bf576add23054dcecc1f2 Mon Sep 17 00:00:00 2001 From: Marc Date: Fri, 31 Jul 2026 17:01:33 +0200 Subject: [PATCH 2/5] Add asynchronous server disposal for subscription shutdown ServerInternalData gains asynchronous disposal so the subscription manager, which now requires it, is shut down without blocking, and both dispose paths are guarded so a second call is a no-op. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 9e6a5abf-3299-4cd1-9855-010fedbf0ad8 --- .../Hosting/OpcUaServerHostedService.cs | 5 +- .../Server/ServerInternalData.cs | 50 +++++++++++++++++-- .../ServerInternalDataTests.cs | 43 ++++++++++++++++ 3 files changed, 92 insertions(+), 6 deletions(-) diff --git a/src/Opc.Ua.Server/Hosting/OpcUaServerHostedService.cs b/src/Opc.Ua.Server/Hosting/OpcUaServerHostedService.cs index e5b06ab7b1..d2f3de91d9 100644 --- a/src/Opc.Ua.Server/Hosting/OpcUaServerHostedService.cs +++ b/src/Opc.Ua.Server/Hosting/OpcUaServerHostedService.cs @@ -40,6 +40,7 @@ using Opc.Ua.Configuration; using Opc.Ua.Identity; using Opc.Ua.Schema; +using Opc.Ua.Security.Certificates; using Opc.Ua.Server.AliasNames; using Opc.Ua.Server.Historian; @@ -427,7 +428,7 @@ await certificateManager.UpdateAsync( configuration.SecurityConfiguration, configuration.ApplicationUri, ct).ConfigureAwait(false); - using var certificates = + using CertificateEntryCollection certificates = certificateManager.SnapshotApplicationCertificates(); return certificates.Count > 0; } @@ -514,7 +515,6 @@ public override void Dispose() m_server?.Dispose(); base.Dispose(); } - } /// @@ -547,5 +547,4 @@ public static partial void UserTokenPolicyTokenTypeIsConfiguredWithout( Message = "Error while stopping OPC UA server.")] public static partial void ErrorWhileStoppingOPCUAServer(this ILogger logger, Exception ex); } - } diff --git a/src/Opc.Ua.Server/Server/ServerInternalData.cs b/src/Opc.Ua.Server/Server/ServerInternalData.cs index 5d7d7da6f4..c4734b365f 100644 --- a/src/Opc.Ua.Server/Server/ServerInternalData.cs +++ b/src/Opc.Ua.Server/Server/ServerInternalData.cs @@ -66,7 +66,7 @@ public class ServerInternalData : Historian.IHistorianRegistryProvider, ITransportListenerRegistryProvider, IServerEndpointRegistryProvider, - + IAsyncDisposable, ITimeProviderProvider { /// @@ -138,13 +138,22 @@ public void Dispose() GC.SuppressFinalize(this); } + /// + /// Frees managed resources that require asynchronous shutdown. + /// + public async ValueTask DisposeAsync() + { + await DisposeAsyncCore().ConfigureAwait(false); + GC.SuppressFinalize(this); + } + /// /// An overrideable version of the Dispose. /// /// true to release both managed and unmanaged resources; false to release only unmanaged resources. protected virtual void Dispose(bool disposing) { - if (disposing) + if (disposing && Interlocked.Exchange(ref m_disposed, 1) == 0) { m_roleStateBinding?.Dispose(); m_roleStateBinding = null; @@ -163,6 +172,39 @@ protected virtual void Dispose(bool disposing) } } + /// + /// An overrideable version of asynchronous dispose. + /// + protected virtual async ValueTask DisposeAsyncCore() + { + if (Interlocked.Exchange(ref m_disposed, 1) != 0) + { + return; + } + + m_roleStateBinding?.Dispose(); + m_roleStateBinding = null; + (RoleManager as IDisposable)?.Dispose(); + ResourceManager?.Dispose(); + RequestManager?.Dispose(); + AggregateManager?.Dispose(); + ModellingRulesManager?.Dispose(); + ConformanceUnitsManager?.Dispose(); + (NodeManager as IDisposable)?.Dispose(); + SessionManager?.Dispose(); + if (SubscriptionManager is IAsyncDisposable asyncSubscriptionManager) + { + await asyncSubscriptionManager.DisposeAsync().ConfigureAwait(false); + } + else + { + SubscriptionManager?.Dispose(); + } + MonitoredItemQueueFactory?.Dispose(); + (AliasNameStoreRegistry as IDisposable)?.Dispose(); + (HistorianRegistry as IDisposable)?.Dispose(); + } + /// /// The server-wide registry of OPC UA Part 17 alias-name stores. /// Surfaces through the optional @@ -702,7 +744,8 @@ public async ValueTask CloseSessionAsync( { await NodeManager.SessionClosingAsync(context, sessionId, deleteSubscriptions, cancellationToken) .ConfigureAwait(false); - await SubscriptionManager.SessionClosingAsync(context, sessionId, deleteSubscriptions, cancellationToken) + await SubscriptionManager + .SessionClosingAsync(context, sessionId, deleteSubscriptions, cancellationToken) .ConfigureAwait(false); } finally @@ -1280,5 +1323,6 @@ private ServiceResult OnUpdateDiagnostics( private RoleStateBinding? m_roleStateBinding; private volatile IReadOnlyList? m_transportListeners; private ArrayOf m_serverEndpoints; + private int m_disposed; } } diff --git a/tests/Opc.Ua.Server.Tests/ServerInternalDataTests.cs b/tests/Opc.Ua.Server.Tests/ServerInternalDataTests.cs index b53d6f483c..5cc6db6258 100644 --- a/tests/Opc.Ua.Server.Tests/ServerInternalDataTests.cs +++ b/tests/Opc.Ua.Server.Tests/ServerInternalDataTests.cs @@ -30,6 +30,7 @@ using System; using System.Linq; using System.Threading; +using System.Threading.Tasks; using Moq; using NUnit.Framework; using Opc.Ua.Tests; @@ -543,5 +544,47 @@ public void ReportAuditEventDoesNothingWhenAuditingDisabled() using ServerInternalData data = CreateServerInternalData(); Assert.DoesNotThrow(() => data.ReportAuditEvent(data.DefaultSystemContext, null)); } + + [Test] + public async Task DisposeAsyncCompletesAsync() + { + ServerInternalData data = CreateServerInternalData(); + + await data.DisposeAsync().ConfigureAwait(false); + + Assert.That(data.RequestManager, Is.Null); + } + + [Test] + public async Task DisposeAsyncIsIdempotentAsync() + { + ServerInternalData data = CreateServerInternalData(); + + await data.DisposeAsync().ConfigureAwait(false); + + Assert.DoesNotThrowAsync(async () => await data.DisposeAsync().ConfigureAwait(false)); + } + + [Test] + public async Task DisposeAfterDisposeAsyncDoesNotDisposeTwiceAsync() + { + ServerInternalData data = CreateServerInternalData(); + + await data.DisposeAsync().ConfigureAwait(false); + + // The synchronous path shares the disposed guard with the asynchronous one, so a + // Dispose that follows DisposeAsync must be a no-op rather than releasing a second time. + Assert.DoesNotThrow(data.Dispose); + } + + [Test] + public void DisposeAsyncAfterDisposeDoesNotDisposeTwice() + { + ServerInternalData data = CreateServerInternalData(); + + data.Dispose(); + + Assert.DoesNotThrowAsync(async () => await data.DisposeAsync().ConfigureAwait(false)); + } } } From d67600072156176da3d2df850929dc0d48612aea Mon Sep 17 00:00:00 2001 From: Marc Date: Fri, 31 Jul 2026 21:38:40 +0200 Subject: [PATCH 3/5] Make subscription transfer rollback and disposal cleanup consistent Share ServerInternalData disposal cleanup so sync and async paths leave the same observable state, while async disposal still awaits async subscription managers when available. Protect synchronous SubscriptionManager disposal with the same worker shutdown and semaphore capture used by async disposal, and make failed transfer rollback preserve unrelated destination publish requests while clearing stale source claims. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 9e6a5abf-3299-4cd1-9855-010fedbf0ad8 --- .../Server/ServerInternalData.cs | 66 ++++++---- .../Subscription/SessionPublishQueue.cs | 5 +- .../Subscription/SubscriptionManager.cs | 59 +++++++-- .../ServerInternalDataTests.cs | 50 +++++++ .../Opc.Ua.Server.Tests/SubscriptionTests.cs | 122 +++++++++++++++--- 5 files changed, 248 insertions(+), 54 deletions(-) diff --git a/src/Opc.Ua.Server/Server/ServerInternalData.cs b/src/Opc.Ua.Server/Server/ServerInternalData.cs index aa43d90a62..3d28158020 100644 --- a/src/Opc.Ua.Server/Server/ServerInternalData.cs +++ b/src/Opc.Ua.Server/Server/ServerInternalData.cs @@ -155,20 +155,7 @@ protected virtual void Dispose(bool disposing) { if (disposing && Interlocked.Exchange(ref m_disposed, 1) == 0) { - m_roleStateBinding?.Dispose(); - m_roleStateBinding = null; - (RoleManager as IDisposable)?.Dispose(); - ResourceManager?.Dispose(); - RequestManager?.Dispose(); - AggregateManager?.Dispose(); - ModellingRulesManager?.Dispose(); - ConformanceUnitsManager?.Dispose(); - (NodeManager as IDisposable)?.Dispose(); - SessionManager?.Dispose(); - SubscriptionManager?.Dispose(); - MonitoredItemQueueFactory?.Dispose(); - (AliasNameStoreRegistry as IDisposable)?.Dispose(); - (HistorianRegistry as IDisposable)?.Dispose(); + DisposeManagedResources(); } } @@ -182,25 +169,60 @@ protected virtual async ValueTask DisposeAsyncCore() return; } + await DisposeManagedResourcesAsync().ConfigureAwait(false); + } + + private void DisposeManagedResources() + { + DisposeManagedResourcesBeforeSubscriptionManager(); + SubscriptionManager?.Dispose(); + DisposeManagedResourcesAfterSubscriptionManager(); + } + + private async ValueTask DisposeManagedResourcesAsync() + { + DisposeManagedResourcesBeforeSubscriptionManager(); + if (SubscriptionManager is IAsyncDisposable asyncSubscriptionManager) + { + await asyncSubscriptionManager.DisposeAsync().ConfigureAwait(false); + } + else + { + SubscriptionManager?.Dispose(); + } + DisposeManagedResourcesAfterSubscriptionManager(); + } + + private void DisposeManagedResourcesBeforeSubscriptionManager() + { m_roleStateBinding?.Dispose(); m_roleStateBinding = null; (RoleManager as IDisposable)?.Dispose(); + RoleManager = null!; ResourceManager?.Dispose(); + ResourceManager = null!; RequestManager?.Dispose(); + RequestManager = null!; AggregateManager?.Dispose(); + AggregateManager = null!; ModellingRulesManager?.Dispose(); + ModellingRulesManager = null!; ConformanceUnitsManager?.Dispose(); + ConformanceUnitsManager = null!; (NodeManager as IDisposable)?.Dispose(); + NodeManager = null!; + DiagnosticsNodeManager = null!; + ConfigurationNodeManager = null!; + CoreNodeManager = null!; SessionManager?.Dispose(); - if (SubscriptionManager is IAsyncDisposable asyncSubscriptionManager) - { - await asyncSubscriptionManager.DisposeAsync().ConfigureAwait(false); - } - else - { - SubscriptionManager?.Dispose(); - } + SessionManager = null!; + } + + private void DisposeManagedResourcesAfterSubscriptionManager() + { + SubscriptionManager = null!; MonitoredItemQueueFactory?.Dispose(); + MonitoredItemQueueFactory = null!; (AliasNameStoreRegistry as IDisposable)?.Dispose(); (HistorianRegistry as IDisposable)?.Dispose(); } diff --git a/src/Opc.Ua.Server/Subscription/SessionPublishQueue.cs b/src/Opc.Ua.Server/Subscription/SessionPublishQueue.cs index 316e0c0e81..929d387ab0 100644 --- a/src/Opc.Ua.Server/Subscription/SessionPublishQueue.cs +++ b/src/Opc.Ua.Server/Subscription/SessionPublishQueue.cs @@ -308,14 +308,13 @@ internal bool RestoreTransferClaim(SubscriptionTransferClaim claim) if (!m_transferClaims.TryGetValue( subscriptionId, out SubscriptionTransferClaim? currentClaim) || - !ReferenceEquals(currentClaim, claim) || - !m_queuedSubscriptions.TryAdd(subscriptionId, claim.Entry)) + !ReferenceEquals(currentClaim, claim)) { return false; } m_transferClaims.Remove(subscriptionId); - return true; + return m_queuedSubscriptions.TryAdd(subscriptionId, claim.Entry); } } diff --git a/src/Opc.Ua.Server/Subscription/SubscriptionManager.cs b/src/Opc.Ua.Server/Subscription/SubscriptionManager.cs index d3c8ff3f63..a69ab71bc4 100644 --- a/src/Opc.Ua.Server/Subscription/SubscriptionManager.cs +++ b/src/Opc.Ua.Server/Subscription/SubscriptionManager.cs @@ -133,11 +133,41 @@ protected virtual void Dispose(bool disposing) if (disposing && Interlocked.Exchange(ref m_disposed, 1) == 0) { SignalConditionRefreshShutdown(); + Task.WaitAll( + m_publishSubscriptionsTask ?? Task.CompletedTask, + m_conditionRefreshTask ?? Task.CompletedTask); + m_semaphoreSlim.Wait(); + List publishQueues; + List subscriptions; + try + { + CaptureManagedResources(out publishQueues, out subscriptions); + } + finally + { + m_semaphoreSlim.Release(); + } + + DisposeManagedResources(publishQueues, subscriptions); + DisposeSynchronizationResources(); + } + } + + private async ValueTask<( + List PublishQueues, + List Subscriptions)> CaptureManagedResourcesAsync() + { + await m_semaphoreSlim.WaitAsync(CancellationToken.None).ConfigureAwait(false); + try + { CaptureManagedResources( out List publishQueues, out List subscriptions); - DisposeManagedResources(publishQueues, subscriptions); - DisposeSynchronizationResources(); + return (publishQueues, subscriptions); + } + finally + { + m_semaphoreSlim.Release(); } } @@ -158,17 +188,9 @@ await Task.WhenAll( m_conditionRefreshTask ?? Task.CompletedTask) .ConfigureAwait(false); - await m_semaphoreSlim.WaitAsync(CancellationToken.None).ConfigureAwait(false); - List publishQueues; - List subscriptions; - try - { - CaptureManagedResources(out publishQueues, out subscriptions); - } - finally - { - m_semaphoreSlim.Release(); - } + (List publishQueues, List subscriptions) + = await CaptureManagedResourcesAsync() + .ConfigureAwait(false); DisposeManagedResources(publishQueues, subscriptions); DisposeSynchronizationResources(); @@ -1810,7 +1832,6 @@ await subscription.TransferSessionAsync( if (destinationAdded && destinationPublishQueue != null) { destinationPublishQueue.TryRemoveForTransfer(subscription); - destinationPublishQueue.RemoveQueuedRequests(); } if (preparedTransfer != null) @@ -2768,6 +2789,16 @@ internal void WakeConditionRefreshWorkerForTest() m_conditionRefreshEvent.Set(); } + internal void EnqueueConditionRefreshForTest( + ISubscription subscription, + uint monitoredItemId = 0) + { + lock (m_conditionRefreshLock) + { + m_conditionRefreshQueue.Enqueue(new ConditionRefreshTask(subscription, monitoredItemId)); + } + } + internal void StartConditionRefreshWorkerForTest() { m_shutdownEvent.Reset(); diff --git a/tests/Opc.Ua.Server.Tests/ServerInternalDataTests.cs b/tests/Opc.Ua.Server.Tests/ServerInternalDataTests.cs index 5cc6db6258..078bc4c0ac 100644 --- a/tests/Opc.Ua.Server.Tests/ServerInternalDataTests.cs +++ b/tests/Opc.Ua.Server.Tests/ServerInternalDataTests.cs @@ -555,6 +555,20 @@ public async Task DisposeAsyncCompletesAsync() Assert.That(data.RequestManager, Is.Null); } + [Test] + public async Task DisposeAndDisposeAsyncLeaveSameObservableStateAsync() + { + ServerInternalData syncData = CreateServerInternalData(); + ServerInternalData asyncData = CreateServerInternalData(); + ConfigureDisposableState(syncData); + ConfigureDisposableState(asyncData); + + syncData.Dispose(); + await asyncData.DisposeAsync().ConfigureAwait(false); + + Assert.That(CaptureDisposedState(asyncData), Is.EqualTo(CaptureDisposedState(syncData))); + } + [Test] public async Task DisposeAsyncIsIdempotentAsync() { @@ -586,5 +600,41 @@ public void DisposeAsyncAfterDisposeDoesNotDisposeTwice() Assert.DoesNotThrowAsync(async () => await data.DisposeAsync().ConfigureAwait(false)); } + + private static void ConfigureDisposableState(ServerInternalData data) + { + var mockNodeManager = new Mock(); + mockNodeManager.Setup(m => m.DiagnosticsNodeManager).Returns((IDiagnosticsNodeManager)null); + mockNodeManager.Setup(m => m.ConfigurationNodeManager).Returns((IConfigurationNodeManager)null); + mockNodeManager.Setup(m => m.CoreNodeManager).Returns((ICoreNodeManager)null); + data.SetNodeManager(mockNodeManager.Object); + + var mockSessionManager = new Mock(); + var mockSubscriptionManager = new Mock(); + mockSubscriptionManager + .As() + .Setup(manager => manager.DisposeAsync()) + .Returns(default(ValueTask)); + data.SetSessionManager(mockSessionManager.Object, mockSubscriptionManager.Object); + + data.SetMonitoredItemQueueFactory(new Mock().Object); + data.SetRoleManager(new Mock().Object); + } + + private static bool[] CaptureDisposedState(ServerInternalData data) + { + return + [ + data.RoleManager == null, + data.NodeManager == null, + data.DiagnosticsNodeManager == null, + data.ConfigurationNodeManager == null, + data.CoreNodeManager == null, + data.SessionManager == null, + data.SubscriptionManager == null, + data.MonitoredItemQueueFactory == null, + data.RequestManager == null + ]; + } } } diff --git a/tests/Opc.Ua.Server.Tests/SubscriptionTests.cs b/tests/Opc.Ua.Server.Tests/SubscriptionTests.cs index b35eab819b..eeff4e0c44 100644 --- a/tests/Opc.Ua.Server.Tests/SubscriptionTests.cs +++ b/tests/Opc.Ua.Server.Tests/SubscriptionTests.cs @@ -412,21 +412,7 @@ private static void EnqueueConditionRefreshTask( SubscriptionManager manager, ISubscription subscription) { - Type taskType = typeof(SubscriptionManager).GetNestedType( - "ConditionRefreshTask", - BindingFlags.NonPublic) - ?? throw new InvalidOperationException( - "ConditionRefreshTask type not found."); - object task = Activator.CreateInstance(taskType, subscription, 0u) - ?? throw new InvalidOperationException( - "ConditionRefreshTask could not be created."); - object queue = GetPrivateField( - manager, - "m_conditionRefreshQueue"); - MethodInfo enqueue = queue.GetType().GetMethod("Enqueue") - ?? throw new InvalidOperationException( - "ConditionRefreshTask queue Enqueue method not found."); - enqueue.Invoke(queue, [task]); + manager.EnqueueConditionRefreshForTest(subscription); } private static void ExpireOnNextPublishTimer(Subscription subscription) @@ -1030,6 +1016,46 @@ async Task WaitForReleaseAsync() Assert.That(refreshCompleted, Is.True); } + [Test] + public void RestoreTransferClaimRemovesCurrentClaimWhenRestoreEntryAlreadyExists() + { + using Subscription subscription = CreateSubscription(); + using var queue = new SessionPublishQueue( + m_serverMock.Object, + m_sessionMock.Object, + maxPublishRequests: 10); + queue.Add(subscription); + + Assert.That( + queue.TryClaimForTransfer( + subscription, + m_sessionMock.Object, + out SessionPublishQueue.SubscriptionTransferClaim claim), + Is.True); + Assert.That(claim, Is.Not.Null); + + var collidingSubscription = new Mock(); + collidingSubscription.SetupGet(sub => sub.Id).Returns(subscription.Id); + queue.Add(collidingSubscription.Object); + + Assert.That(queue.RestoreTransferClaim(claim!), Is.False); + + subscription.AbortTransfer(m_sessionMock.Object); + queue.Remove(collidingSubscription.Object, removeQueuedRequests: false); + queue.Add(subscription); + + Assert.That( + queue.TryClaimForTransfer( + subscription, + m_sessionMock.Object, + out SessionPublishQueue.SubscriptionTransferClaim retryClaim), + Is.True, + "The failed restore must not leave a stale claim that blocks future transfers."); + Assert.That(retryClaim, Is.Not.Null); + queue.CompleteTransferClaim(retryClaim!); + subscription.AbortTransfer(m_sessionMock.Object); + } + [TestCase(0, 0L, 0L)] [TestCase(0, 25L, 25L)] [TestCase(10, 0L, 10L)] @@ -1876,6 +1902,72 @@ await manager.SessionClosingAsync( () => subscription.ResendData(fixture.DestinationContext)); } + [Test] + public async Task FailedTransferPreservesDestinationQueuedPublishRequestsAsync() + { + var fixture = await CreateTransferSubscriptionAsync().ConfigureAwait(false); + using SubscriptionManager manager = fixture.Manager; + Subscription subscription = fixture.Subscription; + fixture.DestinationSession + .Setup(session => session.IsSecureChannelValid("destination-channel")) + .Returns(true); + CreateSubscriptionResponse destinationCreated = await manager + .CreateSubscriptionAsync( + fixture.DestinationContext, + requestedPublishingInterval: 1000, + requestedLifetimeCount: 30, + requestedMaxKeepAliveCount: 10, + maxNotificationsPerPublish: 0, + publishingEnabled: true, + priority: 0) + .ConfigureAwait(false); + Assert.That( + manager.TryGetSubscription( + destinationCreated.SubscriptionId, + out ISubscription destinationSubscription), + Is.True); + Assert.That(destinationSubscription, Is.Not.Null); + SessionPublishQueue destinationQueue = GetPublishQueue( + manager, + fixture.DestinationContext.SessionId); + Task destinationPublish = destinationQueue.PublishAsync( + "destination-channel", + DateTime.MaxValue, + requeue: false, + parkSink: null, + CancellationToken.None); + Assert.That(destinationPublish.IsCompleted, Is.False); + m_nodeManagerMock + .Setup(nodeManager => nodeManager.TransferMonitoredItemsAsync( + It.IsAny(), + It.IsAny(), + It.IsAny>(), + It.IsAny>(), + It.IsAny(), + It.IsAny())) + .Returns(new ValueTask(Task.FromException( + new ServiceResultException(StatusCodes.BadUnexpectedError)))); + + TransferSubscriptionsResponse failed = await manager + .TransferSubscriptionsAsync( + fixture.DestinationContext, + [subscription.Id], + sendInitialValues: false) + .ConfigureAwait(false); + + Assert.Multiple(() => + { + Assert.That(failed.Results, Has.Count.EqualTo(1)); + Assert.That(ServiceResult.IsBad(failed.Results[0].StatusCode), Is.True); + Assert.That(destinationPublish.IsCompleted, Is.False); + }); + + destinationQueue.PublishCompleted(destinationSubscription!, moreNotifications: true); + + ISubscription publishedSubscription = await destinationPublish.ConfigureAwait(false); + Assert.That(publishedSubscription, Is.SameAs(destinationSubscription!)); + } + private async Task RegisterMonitoredItemsAsync( Subscription subscription, params IMonitoredItem[] monitoredItems) From a2e80a7949c6a60ba72923ed4cdb5d769ce2b417 Mon Sep 17 00:00:00 2001 From: Marc Date: Sat, 1 Aug 2026 09:59:54 +0200 Subject: [PATCH 4/5] Cover the transactional subscription transfer failure paths The transfer rollback, abandoned-claim and publish-queue requeue paths carried no tests, which is where a partial transfer would actually be observable: a failure part way through must leave the source session holding every subscription it started with. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 9e6a5abf-3299-4cd1-9855-010fedbf0ad8 --- .../Opc.Ua.Server.Tests/SubscriptionTests.cs | 487 ++++++++++++++++++ 1 file changed, 487 insertions(+) diff --git a/tests/Opc.Ua.Server.Tests/SubscriptionTests.cs b/tests/Opc.Ua.Server.Tests/SubscriptionTests.cs index eeff4e0c44..b8931c6f48 100644 --- a/tests/Opc.Ua.Server.Tests/SubscriptionTests.cs +++ b/tests/Opc.Ua.Server.Tests/SubscriptionTests.cs @@ -1968,6 +1968,493 @@ public async Task FailedTransferPreservesDestinationQueuedPublishRequestsAsync() Assert.That(publishedSubscription, Is.SameAs(destinationSubscription!)); } + [Test] + public void TryClaimForTransferFailsWhenTheClaimingSessionIsNotTheOwner() + { + using Subscription subscription = CreateSubscription(); + using var queue = new SessionPublishQueue( + m_serverMock.Object, + m_sessionMock.Object, + maxPublishRequests: 10); + queue.Add(subscription); + var staleSession = new Mock(); + staleSession.Setup(session => session.Id).Returns(new NodeId(Guid.NewGuid())); + + bool claimed = queue.TryClaimForTransfer( + subscription, + staleSession.Object, + out SessionPublishQueue.SubscriptionTransferClaim claim); + + Assert.Multiple(() => + { + Assert.That(claimed, Is.False); + Assert.That(claim, Is.Null); + Assert.That( + queue.ContainsSubscription(subscription), + Is.True, + "A refused claim must leave the subscription publishable by its owner."); + }); + + Assert.That( + queue.TryClaimForTransfer( + subscription, + m_sessionMock.Object, + out SessionPublishQueue.SubscriptionTransferClaim ownerClaim), + Is.True, + "The refused claim must not block the real owner from starting a transfer."); + queue.CompleteTransferClaim(ownerClaim!); + subscription.AbortTransfer(m_sessionMock.Object); + } + + [Test] + public void RestoreTransferClaimFailsWhenTheClaimWasAlreadyRestored() + { + using Subscription subscription = CreateSubscription(); + using var queue = new SessionPublishQueue( + m_serverMock.Object, + m_sessionMock.Object, + maxPublishRequests: 10); + queue.Add(subscription); + Assert.That( + queue.TryClaimForTransfer( + subscription, + m_sessionMock.Object, + out SessionPublishQueue.SubscriptionTransferClaim claim), + Is.True); + + Assert.That(queue.RestoreTransferClaim(claim!), Is.True); + + Assert.Multiple(() => + { + Assert.That( + queue.RestoreTransferClaim(claim!), + Is.False, + "A claim may only be restored once, otherwise a later queue entry is overwritten."); + Assert.That(queue.ContainsSubscription(subscription), Is.True); + }); + subscription.AbortTransfer(m_sessionMock.Object); + } + + [Test] + public void TryRemoveForTransferRemovesOnlyTheExactQueuedEntry() + { + using Subscription subscription = CreateSubscription(); + using var queue = new SessionPublishQueue( + m_serverMock.Object, + m_sessionMock.Object, + maxPublishRequests: 10); + queue.Add(subscription); + var impostor = new Mock(); + impostor.Setup(candidate => candidate.Id).Returns(subscription.Id); + + Assert.Multiple(() => + { + Assert.That( + queue.TryRemoveForTransfer(impostor.Object), + Is.False, + "A different subscription instance with the same id must not remove the entry."); + Assert.That(queue.ContainsSubscription(subscription), Is.True); + }); + + Assert.That(queue.TryRemoveForTransfer(subscription), Is.True); + + Assert.Multiple(() => + { + Assert.That(queue.ContainsSubscription(subscription), Is.False); + Assert.That( + queue.TryRemoveForTransfer(subscription), + Is.False, + "A subscription that is no longer queued cannot be removed a second time."); + }); + } + + [Test] + public async Task PublishCompletedKeepsNotificationsOfAClaimedSubscriptionAsync() + { + using Subscription subscription = CreateSubscription(); + using var queue = new SessionPublishQueue( + m_serverMock.Object, + m_sessionMock.Object, + maxPublishRequests: 10); + queue.Add(subscription); + Assert.That( + queue.TryClaimForTransfer( + subscription, + m_sessionMock.Object, + out SessionPublishQueue.SubscriptionTransferClaim claim), + Is.True); + + queue.PublishCompleted(subscription, moreNotifications: true); + + Assert.Multiple(() => + { + Assert.That(claim!.Entry.Publishing, Is.False); + Assert.That( + claim.Entry.ReadyToPublish, + Is.True, + "A notification raised while the entry is claimed must survive an abandoned transfer."); + }); + + Assert.That(queue.RestoreTransferClaim(claim!), Is.True); + Task publish = queue.PublishAsync( + "channel1", + DateTime.MaxValue, + requeue: false, + parkSink: null, + CancellationToken.None); + + Assert.That(publish.Status, Is.EqualTo(TaskStatus.RanToCompletion)); + Assert.That(await publish.ConfigureAwait(false), Is.SameAs(subscription)); + subscription.AbortTransfer(m_sessionMock.Object); + } + + [Test] + public void TryBeginTransferFailsWhileATransferIsAlreadyInProgress() + { + using Subscription subscription = CreateSubscription(); + + Assert.That(subscription.TryBeginTransfer(m_sessionMock.Object), Is.True); + + Assert.That( + subscription.TryBeginTransfer(m_sessionMock.Object), + Is.False, + "Only one transfer may reserve a subscription at a time."); + + subscription.AbortTransfer(m_sessionMock.Object); + Assert.That( + subscription.TryBeginTransfer(m_sessionMock.Object), + Is.True, + "An aborted transfer must release the reservation."); + subscription.AbortTransfer(m_sessionMock.Object); + } + + [Test] + public void PublishTimerIsIdleWhileATransferIsInProgress() + { + using Subscription subscription = CreateSubscription(); + ExpireOnNextPublishTimer(subscription); + Assert.That(subscription.TryBeginTransfer(m_sessionMock.Object), Is.True); + + PublishingState state = subscription.PublishTimerExpired(); + + Assert.That( + state, + Is.EqualTo(PublishingState.Idle), + "A reserved subscription must not expire on the source session timer."); + subscription.AbortTransfer(m_sessionMock.Object); + } + + [Test] + public void PrepareSessionTransferAsyncRejectsAnUnreservedSubscription() + { + using Subscription subscription = CreateSubscription(); + var context = new OperationContext( + m_sessionMock.Object, + DiagnosticsMasks.None); + + ServiceResultException exception = Assert.ThrowsAsync( + async () => + { + await subscription + .PrepareSessionTransferAsync( + context, + m_sessionMock.Object, + sendInitialValues: false, + CancellationToken.None) + .ConfigureAwait(false); + }); + + Assert.That( + exception.StatusCode, + Is.EqualTo(StatusCodes.BadSubscriptionIdInvalid)); + } + + [Test] + public void CompleteTransferRejectsAnUnexpectedOwner() + { + using Subscription subscription = CreateSubscription(); + + ServiceResultException exception = Assert.Throws( + () => subscription.CompleteTransfer(m_sessionMock.Object)); + + Assert.That( + exception.StatusCode, + Is.EqualTo(StatusCodes.BadSubscriptionIdInvalid)); + } + + [Test] + public async Task TransferSessionAsyncMovesOwnershipToTheDestinationAsync() + { + using Subscription subscription = CreateSubscription(); + var destinationSession = new Mock(); + destinationSession.Setup(session => session.Id).Returns(new NodeId(Guid.NewGuid())); + SetSessionIdentity( + destinationSession, + new UserNameIdentityTokenHandler("transfer-user", [1, 2, 3])); + var context = new OperationContext( + destinationSession.Object, + DiagnosticsMasks.None); + + await subscription + .TransferSessionAsync(context, sendInitialValues: true, CancellationToken.None) + .ConfigureAwait(false); + + Assert.That(subscription.Session, Is.SameAs(destinationSession.Object)); + m_nodeManagerMock.Verify( + nodeManager => nodeManager.TransferMonitoredItemsAsync( + context, + true, + It.IsAny>(), + It.IsAny>(), + null, + It.IsAny()), + Times.Once); + } + + [Test] + public async Task TransferFailsWhenTheSourceQueueEntryIsAlreadyClaimedAsync() + { + var fixture = await CreateTransferSubscriptionAsync().ConfigureAwait(false); + using SubscriptionManager manager = fixture.Manager; + Subscription subscription = fixture.Subscription; + SessionPublishQueue sourceQueue = GetPublishQueue( + manager, + fixture.SourceContext.SessionId); + Assert.That( + sourceQueue.TryClaimForTransfer( + subscription, + m_sessionMock.Object, + out SessionPublishQueue.SubscriptionTransferClaim claim), + Is.True); + var diagnosticsContext = new OperationContext( + fixture.DestinationSession.Object, + DiagnosticsMasks.OperationAll); + + TransferSubscriptionsResponse failed = await manager + .TransferSubscriptionsAsync( + diagnosticsContext, + [subscription.Id], + sendInitialValues: false) + .ConfigureAwait(false); + + Assert.Multiple(() => + { + Assert.That(failed.Results, Has.Count.EqualTo(1)); + Assert.That( + failed.Results[0].StatusCode, + Is.EqualTo(StatusCodes.BadSubscriptionIdInvalid)); + Assert.That(failed.DiagnosticInfos, Has.Count.EqualTo(1)); + Assert.That( + subscription.Session, + Is.SameAs(m_sessionMock.Object), + "A refused claim must leave ownership with the source session."); + }); + + Assert.That(sourceQueue.RestoreTransferClaim(claim!), Is.True); + subscription.AbortTransfer(m_sessionMock.Object); + TransferSubscriptionsResponse retried = await manager + .TransferSubscriptionsAsync( + fixture.DestinationContext, + [subscription.Id], + sendInitialValues: false) + .ConfigureAwait(false); + + Assert.That(retried.Results[0].StatusCode, Is.EqualTo(StatusCodes.Good)); + } + + [Test] + public async Task TransferFailsWhenTheAbandonedSubscriptionIsAlreadyReservedAsync() + { + var fixture = await CreateTransferSubscriptionAsync().ConfigureAwait(false); + using SubscriptionManager manager = fixture.Manager; + Subscription subscription = fixture.Subscription; + m_sessionMock.SetupGet(session => session.IsClosing).Returns(true); + await manager.SessionClosingAsync( + fixture.SourceContext, + fixture.SourceContext.SessionId, + deleteSubscriptions: false, + CancellationToken.None) + .ConfigureAwait(false); + Assert.That(subscription.Session, Is.Null); + Assert.That(subscription.TryBeginTransfer(null), Is.True); + var diagnosticsContext = new OperationContext( + fixture.DestinationSession.Object, + DiagnosticsMasks.OperationAll); + + TransferSubscriptionsResponse failed = await manager + .TransferSubscriptionsAsync( + diagnosticsContext, + [subscription.Id], + sendInitialValues: false) + .ConfigureAwait(false); + + ConcurrentDictionary abandonedSubscriptions = + GetPrivateField>( + manager, + "m_abandonedSubscriptions"); + Assert.Multiple(() => + { + Assert.That(failed.Results, Has.Count.EqualTo(1)); + Assert.That( + failed.Results[0].StatusCode, + Is.EqualTo(StatusCodes.BadSubscriptionIdInvalid)); + Assert.That(failed.DiagnosticInfos, Has.Count.EqualTo(1)); + Assert.That( + abandonedSubscriptions.TryGetValue( + subscription.Id, + out ISubscription retainedSubscription), + Is.True, + "A refused reservation must leave the subscription abandoned, not lost."); + Assert.That(retainedSubscription, Is.SameAs(subscription)); + }); + + subscription.AbortTransfer(null); + TransferSubscriptionsResponse retried = await manager + .TransferSubscriptionsAsync( + fixture.DestinationContext, + [subscription.Id], + sendInitialValues: false) + .ConfigureAwait(false); + + Assert.That(retried.Results[0].StatusCode, Is.EqualTo(StatusCodes.Good)); + } + + [Test] + public async Task FailedOwnershipCommitRollsBackThePreparedTransferAsync() + { + var fixture = await CreateTransferSubscriptionAsync().ConfigureAwait(false); + using SubscriptionManager manager = fixture.Manager; + Subscription subscription = fixture.Subscription; + SessionPublishQueue sourceQueue = GetPublishQueue( + manager, + fixture.SourceContext.SessionId); + int transferCalls = 0; + m_nodeManagerMock + .Setup(nodeManager => nodeManager.TransferMonitoredItemsAsync( + It.IsAny(), + It.IsAny(), + It.IsAny>(), + It.IsAny>(), + It.IsAny(), + It.IsAny())) + .Callback(() => + { + // Drop the reservation while the monitored items are handed over so the + // ownership commit that follows preparation has to fail and roll back. + if (Interlocked.Increment(ref transferCalls) == 1) + { + subscription.AbortTransfer(m_sessionMock.Object); + } + }) + .Returns(default(ValueTask)); + + TransferSubscriptionsResponse failed = await manager + .TransferSubscriptionsAsync( + fixture.DestinationContext, + [subscription.Id], + sendInitialValues: false) + .ConfigureAwait(false); + + IReadOnlyList restoredSource = + sourceQueue.CapturePublishTimerSnapshot(); + Assert.Multiple(() => + { + Assert.That(failed.Results, Has.Count.EqualTo(1)); + Assert.That(ServiceResult.IsBad(failed.Results[0].StatusCode), Is.True); + Assert.That( + subscription.Session, + Is.SameAs(m_sessionMock.Object), + "A failed ownership commit must leave the source session as the owner."); + Assert.That(restoredSource, Has.Count.EqualTo(1)); + Assert.That(restoredSource[0].Subscription, Is.SameAs(subscription)); + }); + + TransferSubscriptionsResponse retried = await manager + .TransferSubscriptionsAsync( + fixture.DestinationContext, + [subscription.Id], + sendInitialValues: false) + .ConfigureAwait(false); + + Assert.Multiple(() => + { + Assert.That(retried.Results[0].StatusCode, Is.EqualTo(StatusCodes.Good)); + Assert.That( + subscription.Session, + Is.SameAs(fixture.DestinationSession.Object)); + }); + } + + [Test] + public async Task FailedSourceQueueRestoreAggregatesTheTransferErrorAsync() + { + var fixture = await CreateTransferSubscriptionAsync().ConfigureAwait(false); + using SubscriptionManager manager = fixture.Manager; + Subscription subscription = fixture.Subscription; + SessionPublishQueue sourceQueue = GetPublishQueue( + manager, + fixture.SourceContext.SessionId); + var collidingSubscription = new Mock(); + collidingSubscription.Setup(candidate => candidate.Id).Returns(subscription.Id); + int transferCalls = 0; + m_nodeManagerMock + .Setup(nodeManager => nodeManager.TransferMonitoredItemsAsync( + It.IsAny(), + It.IsAny(), + It.IsAny>(), + It.IsAny>(), + It.IsAny(), + It.IsAny())) + .Returns(() => + { + if (Interlocked.Increment(ref transferCalls) != 1) + { + return default; + } + + // Occupy the claimed slot so restoring the source queue entry has to + // fail, which is the only way the rollback itself can report an error. + sourceQueue.Add(collidingSubscription.Object); + return new ValueTask( + Task.FromException( + new ServiceResultException(StatusCodes.BadUnexpectedError))); + }); + + TransferSubscriptionsResponse failed = await manager + .TransferSubscriptionsAsync( + fixture.DestinationContext, + [subscription.Id], + sendInitialValues: false) + .ConfigureAwait(false); + + Assert.Multiple(() => + { + Assert.That(failed.Results, Has.Count.EqualTo(1)); + Assert.That(ServiceResult.IsBad(failed.Results[0].StatusCode), Is.True); + Assert.That( + subscription.Session, + Is.SameAs(m_sessionMock.Object), + "A rollback that cannot restore the queue must still keep source ownership."); + Assert.That( + sourceQueue.ContainsSubscription(subscription), + Is.False); + }); + + sourceQueue.Remove(collidingSubscription.Object, removeQueuedRequests: false); + sourceQueue.Add(subscription); + TransferSubscriptionsResponse retried = await manager + .TransferSubscriptionsAsync( + fixture.DestinationContext, + [subscription.Id], + sendInitialValues: false) + .ConfigureAwait(false); + + Assert.That( + retried.Results[0].StatusCode, + Is.EqualTo(StatusCodes.Good), + "A reported rollback failure must not leave a stale claim behind."); + } + private async Task RegisterMonitoredItemsAsync( Subscription subscription, params IMonitoredItem[] monitoredItems) From 1feef013604cc0051d23772b3d2cf6b575af941e Mon Sep 17 00:00:00 2001 From: Marc Date: Mon, 3 Aug 2026 17:37:26 +0200 Subject: [PATCH 5/5] Apply owner-approved sync disposal bridge ServerInternalData.Dispose now deliberately blocks on DisposeAsyncCore for the owner-approved sync-over-async exception, while preserving the shared disposed guard so Dispose after DisposeAsync remains a no-op. The disposal helpers were flattened into DisposeAsyncCore to keep ordering explicit around subscription manager disposal. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> Copilot-Session: 9e6a5abf-3299-4cd1-9855-010fedbf0ad8 --- .../Server/ServerInternalData.cs | 68 ++++++++-------- .../ServerInternalDataTests.cs | 81 ++++++++++++++++++- 2 files changed, 112 insertions(+), 37 deletions(-) diff --git a/src/Opc.Ua.Server/Server/ServerInternalData.cs b/src/Opc.Ua.Server/Server/ServerInternalData.cs index 3d28158020..11cd2034ca 100644 --- a/src/Opc.Ua.Server/Server/ServerInternalData.cs +++ b/src/Opc.Ua.Server/Server/ServerInternalData.cs @@ -130,8 +130,12 @@ public ServerInternalData( } /// - /// Frees any unmanaged resources. + /// Frees resources by running the asynchronous disposal core synchronously. /// + /// + /// Callers should prefer . If has already run, + /// this method is a no-op. + /// public void Dispose() { Dispose(true); @@ -139,8 +143,11 @@ public void Dispose() } /// - /// Frees managed resources that require asynchronous shutdown. + /// Frees resources asynchronously. /// + /// + /// This is the preferred disposal path. A subsequent call to is a no-op. + /// public async ValueTask DisposeAsync() { await DisposeAsyncCore().ConfigureAwait(false); @@ -148,20 +155,31 @@ public async ValueTask DisposeAsync() } /// - /// An overrideable version of the Dispose. + /// Runs the asynchronous disposal core synchronously when disposing managed resources. /// - /// true to release both managed and unmanaged resources; false to release only unmanaged resources. + /// + /// true to release managed resources; false to release only unmanaged resources. + /// + /// + /// If has already run, this method is a no-op. + /// protected virtual void Dispose(bool disposing) { - if (disposing && Interlocked.Exchange(ref m_disposed, 1) == 0) + if (disposing && Volatile.Read(ref m_disposed) == 0) { - DisposeManagedResources(); +#pragma warning disable CA2012 // Owner-approved sync dispose path blocks here; TODO: remove if contract changes. + DisposeAsyncCore().GetAwaiter().GetResult(); +#pragma warning restore CA2012 } } /// /// An overrideable version of asynchronous dispose. /// + /// + /// This method performs the full managed-resource cleanup. calls this method + /// synchronously when callers use the synchronous disposal path. + /// protected virtual async ValueTask DisposeAsyncCore() { if (Interlocked.Exchange(ref m_disposed, 1) != 0) @@ -169,32 +187,6 @@ protected virtual async ValueTask DisposeAsyncCore() return; } - await DisposeManagedResourcesAsync().ConfigureAwait(false); - } - - private void DisposeManagedResources() - { - DisposeManagedResourcesBeforeSubscriptionManager(); - SubscriptionManager?.Dispose(); - DisposeManagedResourcesAfterSubscriptionManager(); - } - - private async ValueTask DisposeManagedResourcesAsync() - { - DisposeManagedResourcesBeforeSubscriptionManager(); - if (SubscriptionManager is IAsyncDisposable asyncSubscriptionManager) - { - await asyncSubscriptionManager.DisposeAsync().ConfigureAwait(false); - } - else - { - SubscriptionManager?.Dispose(); - } - DisposeManagedResourcesAfterSubscriptionManager(); - } - - private void DisposeManagedResourcesBeforeSubscriptionManager() - { m_roleStateBinding?.Dispose(); m_roleStateBinding = null; (RoleManager as IDisposable)?.Dispose(); @@ -216,10 +208,14 @@ private void DisposeManagedResourcesBeforeSubscriptionManager() CoreNodeManager = null!; SessionManager?.Dispose(); SessionManager = null!; - } - - private void DisposeManagedResourcesAfterSubscriptionManager() - { + if (SubscriptionManager is IAsyncDisposable asyncSubscriptionManager) + { + await asyncSubscriptionManager.DisposeAsync().ConfigureAwait(false); + } + else + { + SubscriptionManager?.Dispose(); + } SubscriptionManager = null!; MonitoredItemQueueFactory?.Dispose(); MonitoredItemQueueFactory = null!; diff --git a/tests/Opc.Ua.Server.Tests/ServerInternalDataTests.cs b/tests/Opc.Ua.Server.Tests/ServerInternalDataTests.cs index 078bc4c0ac..c7d74ef15e 100644 --- a/tests/Opc.Ua.Server.Tests/ServerInternalDataTests.cs +++ b/tests/Opc.Ua.Server.Tests/ServerInternalDataTests.cs @@ -555,6 +555,20 @@ public async Task DisposeAsyncCompletesAsync() Assert.That(data.RequestManager, Is.Null); } + [Test] + public void DisposeReleasesManagedResourcesOnce() + { + ServerInternalData data = CreateServerInternalData(); + DisposalCounts counts = ConfigureCountingDisposableState(data); + + data.Dispose(); + + Assert.That(CaptureDisposedState(data), Is.All.True); + Assert.That(counts.Total, Is.EqualTo(5)); + Assert.That(counts.SubscriptionAsync, Is.EqualTo(1)); + Assert.That(counts.SubscriptionSync, Is.Zero); + } + [Test] public async Task DisposeAndDisposeAsyncLeaveSameObservableStateAsync() { @@ -583,22 +597,30 @@ public async Task DisposeAsyncIsIdempotentAsync() public async Task DisposeAfterDisposeAsyncDoesNotDisposeTwiceAsync() { ServerInternalData data = CreateServerInternalData(); + DisposalCounts counts = ConfigureCountingDisposableState(data); await data.DisposeAsync().ConfigureAwait(false); // The synchronous path shares the disposed guard with the asynchronous one, so a // Dispose that follows DisposeAsync must be a no-op rather than releasing a second time. Assert.DoesNotThrow(data.Dispose); + Assert.That(counts.Total, Is.EqualTo(5)); + Assert.That(counts.SubscriptionAsync, Is.EqualTo(1)); + Assert.That(counts.SubscriptionSync, Is.Zero); } [Test] - public void DisposeAsyncAfterDisposeDoesNotDisposeTwice() + public async Task DisposeAsyncAfterDisposeDoesNotDisposeTwiceAsync() { ServerInternalData data = CreateServerInternalData(); + DisposalCounts counts = ConfigureCountingDisposableState(data); data.Dispose(); Assert.DoesNotThrowAsync(async () => await data.DisposeAsync().ConfigureAwait(false)); + Assert.That(counts.Total, Is.EqualTo(5)); + Assert.That(counts.SubscriptionAsync, Is.EqualTo(1)); + Assert.That(counts.SubscriptionSync, Is.Zero); } private static void ConfigureDisposableState(ServerInternalData data) @@ -621,6 +643,40 @@ private static void ConfigureDisposableState(ServerInternalData data) data.SetRoleManager(new Mock().Object); } + private static DisposalCounts ConfigureCountingDisposableState(ServerInternalData data) + { + var counts = new DisposalCounts(); + + var mockNodeManager = new Mock(); + mockNodeManager.Setup(manager => manager.DiagnosticsNodeManager).Returns((IDiagnosticsNodeManager)null); + mockNodeManager.Setup(manager => manager.ConfigurationNodeManager).Returns((IConfigurationNodeManager)null); + mockNodeManager.Setup(manager => manager.CoreNodeManager).Returns((ICoreNodeManager)null); + mockNodeManager.As().Setup(manager => manager.Dispose()).Callback(() => counts.NodeManager++); + data.SetNodeManager(mockNodeManager.Object); + + var mockSessionManager = new Mock(); + mockSessionManager.Setup(manager => manager.Dispose()).Callback(() => counts.SessionManager++); + + var mockSubscriptionManager = new Mock(); + mockSubscriptionManager.Setup(manager => manager.Dispose()).Callback(() => counts.SubscriptionSync++); + mockSubscriptionManager + .As() + .Setup(manager => manager.DisposeAsync()) + .Callback(() => counts.SubscriptionAsync++) + .Returns(default(ValueTask)); + data.SetSessionManager(mockSessionManager.Object, mockSubscriptionManager.Object); + + var mockQueueFactory = new Mock(); + mockQueueFactory.Setup(factory => factory.Dispose()).Callback(() => counts.MonitoredItemQueueFactory++); + data.SetMonitoredItemQueueFactory(mockQueueFactory.Object); + + var mockRoleManager = new Mock(); + mockRoleManager.As().Setup(manager => manager.Dispose()).Callback(() => counts.RoleManager++); + data.SetRoleManager(mockRoleManager.Object); + + return counts; + } + private static bool[] CaptureDisposedState(ServerInternalData data) { return @@ -636,5 +692,28 @@ private static bool[] CaptureDisposedState(ServerInternalData data) data.RequestManager == null ]; } + + private sealed class DisposalCounts + { + public int Total => + RoleManager + + NodeManager + + SessionManager + + SubscriptionSync + + SubscriptionAsync + + MonitoredItemQueueFactory; + + public int RoleManager { get; set; } + + public int NodeManager { get; set; } + + public int SessionManager { get; set; } + + public int SubscriptionSync { get; set; } + + public int SubscriptionAsync { get; set; } + + public int MonitoredItemQueueFactory { get; set; } + } } }