Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
16 commits
Select commit Hold shift + click to select a range
99d5d44
Make subscription transfer between sessions transactional
marcschier Jul 31, 2026
41d353d
Add asynchronous server disposal for subscription shutdown
marcschier Jul 31, 2026
098eb48
Merge branch 'marcschier/wot-05-lifecycle' into marcschier/wot-07-sub…
marcschier Jul 31, 2026
d676000
Make subscription transfer rollback and disposal cleanup consistent
marcschier Jul 31, 2026
a2e80a7
Cover the transactional subscription transfer failure paths
marcschier Aug 1, 2026
4b77f5a
Merge branch 'marcschier/wot-05-lifecycle' into marcschier/wot-07-sub…
marcschier Aug 1, 2026
a287a39
Merge branch 'marcschier/wot-05-lifecycle' into marcschier/wot-07-sub…
marcschier Aug 1, 2026
eb6e24e
Merge branch 'marcschier/wot-05-lifecycle' into marcschier/wot-07-sub…
marcschier Aug 1, 2026
8110f7a
Merge remote-tracking branch 'origin/marcschier/wot-05-lifecycle' int…
marcschier Aug 2, 2026
c360578
Merge remote-tracking branch 'origin/marcschier/wot-05-lifecycle' int…
marcschier Aug 3, 2026
2cdeae6
Merge remote-tracking branch 'origin/marcschier/wot-05-lifecycle' int…
marcschier Aug 3, 2026
8623b6a
Merge remote-tracking branch 'origin/marcschier/wot-05-lifecycle' int…
marcschier Aug 3, 2026
f70dd4e
Merge remote-tracking branch 'origin/marcschier/wot-05-lifecycle' int…
marcschier Aug 3, 2026
1feef01
Apply owner-approved sync disposal bridge
marcschier Aug 3, 2026
d47e8aa
Merge remote-tracking branch 'origin/marcschier/wot-05-lifecycle' int…
marcschier Aug 3, 2026
5260eda
Merge remote-tracking branch 'origin/marcschier/wot-05-lifecycle' int…
marcschier Aug 3, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 2 additions & 3 deletions src/Opc.Ua.Server/Hosting/OpcUaServerHostedService.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -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;
}
Expand Down Expand Up @@ -514,7 +515,6 @@ public override void Dispose()
m_server?.Dispose();
base.Dispose();
}

}

/// <summary>
Expand Down Expand Up @@ -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);
}

}
11 changes: 11 additions & 0 deletions src/Opc.Ua.Server/NodeManager/MasterNodeManager.cs
Original file line number Diff line number Diff line change
Expand Up @@ -5522,6 +5522,7 @@ private async ValueTask<IMonitoredItemTransferTransaction>

var processedItems = new List<bool>(monitoredItems.Count);
var originalErrors = new ServiceResult[errors.Count];
var resendStates = new bool[monitoredItems.Count];
var effectiveTransferOptions = new MonitoredItemTransferOptions
{
DeferInitialValues = sendInitialValues ||
Expand All @@ -5533,6 +5534,7 @@ private async ValueTask<IMonitoredItemTransferTransaction>
{
IMonitoredItem? monitoredItem = monitoredItems[ii];
originalErrors[ii] = errors[ii];
resendStates[ii] = monitoredItem?.IsResendData ?? false;
bool isDetached = monitoredItem is IDetachableMonitoredItem
{
IsDetached: true
Expand Down Expand Up @@ -5588,6 +5590,7 @@ await owner.TransferMonitoredItemsAsync(
monitoredItems,
errors,
originalErrors,
resendStates,
participants,
effectiveTransferOptions);
try
Expand Down Expand Up @@ -5624,6 +5627,7 @@ await failedTransaction.RollbackAsync(CancellationToken.None)
monitoredItems,
errors,
originalErrors,
resendStates,
participants,
effectiveTransferOptions);
}
Expand Down Expand Up @@ -5684,6 +5688,7 @@ public MonitoredItemTransferTransaction(
IList<IMonitoredItem> monitoredItems,
IList<ServiceResult> errors,
ServiceResult[] originalErrors,
bool[] resendStates,
List<MonitoredItemTransferParticipant> participants,
MonitoredItemTransferOptions transferOptions)
{
Expand All @@ -5692,6 +5697,7 @@ public MonitoredItemTransferTransaction(
m_monitoredItems = monitoredItems;
m_errors = errors;
m_originalErrors = originalErrors;
m_resendStates = resendStates;
m_participants = participants;
m_transferOptions = transferOptions;
}
Expand Down Expand Up @@ -5769,6 +5775,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];
}

Expand All @@ -5785,6 +5795,7 @@ await participant.NodeManager.RollbackMonitoredItemsTransferAsync(
private readonly IList<IMonitoredItem> m_monitoredItems;
private readonly IList<ServiceResult> m_errors;
private readonly ServiceResult[] m_originalErrors;
private readonly bool[] m_resendStates;
private readonly List<MonitoredItemTransferParticipant> m_participants;
private readonly MonitoredItemTransferOptions m_transferOptions;
private int m_state;
Expand Down
100 changes: 81 additions & 19 deletions src/Opc.Ua.Server/Server/ServerInternalData.cs
Original file line number Diff line number Diff line change
Expand Up @@ -66,7 +66,7 @@ public class ServerInternalData :
Historian.IHistorianRegistryProvider,
ITransportListenerRegistryProvider,
IServerEndpointRegistryProvider,

IAsyncDisposable,
ITimeProviderProvider
{
/// <summary>
Expand Down Expand Up @@ -130,37 +130,97 @@ public ServerInternalData(
}

/// <summary>
/// Frees any unmanaged resources.
/// Frees resources by running the asynchronous disposal core synchronously.
/// </summary>
/// <remarks>
/// Callers should prefer <see cref="DisposeAsync"/>. If <see cref="DisposeAsync"/> has already run,
/// this method is a no-op.
/// </remarks>
public void Dispose()
{
Dispose(true);
GC.SuppressFinalize(this);
}

/// <summary>
/// An overrideable version of the Dispose.
/// Frees resources asynchronously.
/// </summary>
/// <param name="disposing"><c>true</c> to release both managed and unmanaged resources; <c>false</c> to release only unmanaged resources.</param>
/// <remarks>
/// This is the preferred disposal path. A subsequent call to <see cref="Dispose()"/> is a no-op.
/// </remarks>
public async ValueTask DisposeAsync()
{
await DisposeAsyncCore().ConfigureAwait(false);
GC.SuppressFinalize(this);
}

/// <summary>
/// Runs the asynchronous disposal core synchronously when disposing managed resources.
/// </summary>
/// <param name="disposing">
/// <c>true</c> to release managed resources; <c>false</c> to release only unmanaged resources.
/// </param>
/// <remarks>
/// If <see cref="DisposeAsyncCore"/> has already run, this method is a no-op.
/// </remarks>
protected virtual void Dispose(bool disposing)
{
if (disposing)
if (disposing && Volatile.Read(ref m_disposed) == 0)
{
#pragma warning disable CA2012 // Owner-approved sync dispose path blocks here; TODO: remove if contract changes.
DisposeAsyncCore().GetAwaiter().GetResult();
#pragma warning restore CA2012
}
}

/// <summary>
/// An overrideable version of asynchronous dispose.
/// </summary>
/// <remarks>
/// This method performs the full managed-resource cleanup. <see cref="Dispose()"/> calls this method
/// synchronously when callers use the synchronous disposal path.
/// </remarks>
protected virtual async ValueTask DisposeAsyncCore()
{
if (Interlocked.Exchange(ref m_disposed, 1) != 0)
{
return;
}

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();
SessionManager = null!;
if (SubscriptionManager is IAsyncDisposable asyncSubscriptionManager)
{
await asyncSubscriptionManager.DisposeAsync().ConfigureAwait(false);
}
else
{
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();
}
SubscriptionManager = null!;
MonitoredItemQueueFactory?.Dispose();
MonitoredItemQueueFactory = null!;
(AliasNameStoreRegistry as IDisposable)?.Dispose();
(HistorianRegistry as IDisposable)?.Dispose();
}
Comment thread
marcschier marked this conversation as resolved.

/// <summary>
Expand Down Expand Up @@ -702,7 +762,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
Expand Down Expand Up @@ -1289,5 +1350,6 @@ private ServiceResult OnUpdateDiagnostics(
private RoleStateBinding? m_roleStateBinding;
private volatile IReadOnlyList<ITransportListener>? m_transportListeners;
private ArrayOf<EndpointDescription> m_serverEndpoints;
private int m_disposed;
}
}
12 changes: 12 additions & 0 deletions src/Opc.Ua.Server/Subscription/MonitoredItem/IMonitoredItem.cs
Original file line number Diff line number Diff line change
Expand Up @@ -346,6 +346,18 @@ public interface ISampledDataChangeMonitoredItem : IDataChangeMonitoredItem2
void SetSamplingInterval(double samplingInterval);
}

/// <summary>
/// Restores monitored item transient state when a prepared subscription transfer is rolled back.
/// </summary>
internal interface IMonitoredItemTransferState
{
/// <summary>
/// Restores the resend-data trigger to the value captured before transfer preparation.
/// </summary>
/// <param name="resendData">The original resend-data trigger state.</param>
void RestoreResendDataTrigger(bool resendData);
}

/// <summary>
/// Defines constants for the monitored item type.
/// </summary>
Expand Down
14 changes: 12 additions & 2 deletions src/Opc.Ua.Server/Subscription/MonitoredItem/MonitoredItem.cs
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,8 @@ public class MonitoredItem :
IEventMonitoredItem,
ISampledDataChangeMonitoredItem,
ITriggeredMonitoredItem,
IDetachableMonitoredItem
IDetachableMonitoredItem,
IMonitoredItemTransferState
{
/// <summary>
/// Initializes the object with its node type.
Expand Down Expand Up @@ -510,6 +511,14 @@ public void SetupResendDataTrigger()
}
}

void IMonitoredItemTransferState.RestoreResendDataTrigger(bool resendData)
{
lock (m_lock)
{
m_resendData = resendData;
}
}

/// <summary>
/// Sets a flag indicating that the item has been triggered and should publish.
/// </summary>
Expand Down Expand Up @@ -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))
{
Expand Down
Loading
Loading