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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 17 additions & 1 deletion Libraries/Opc.Ua.Client/Subscription/MonitoredItem.cs
Original file line number Diff line number Diff line change
Expand Up @@ -136,6 +136,8 @@ public virtual void Restore(MonitoredItemState state)
State = state;
ClientHandle = state.ClientId;
ServerId = state.ServerId;
TriggeringItemId = state.TriggeringItemId;
TriggeredItems = state.TriggeredItems != null ? new UInt32Collection(state.TriggeredItems) : null;
}

/// <inheritdoc/>
Expand All @@ -144,7 +146,9 @@ public virtual void Snapshot(out MonitoredItemState state)
state = new MonitoredItemState(State)
{
ServerId = Status.Id,
ClientId = ClientHandle
ClientId = ClientHandle,
TriggeringItemId = TriggeringItemId,
TriggeredItems = TriggeredItems != null ? new UInt32Collection(TriggeredItems) : null
};
}

Expand Down Expand Up @@ -1070,6 +1074,18 @@ private static EventFilter GetDefaultEventFilter()
private MonitoredItemEventCache? m_eventCache;
private IEncodeable? m_lastNotification;
private event MonitoredItemNotificationEventHandler? m_Notification;

/// <summary>
/// Server-side identifier of the triggering item if this monitored item
/// is triggered by another item. 0 indicates this item is not triggered.
/// </summary>
internal uint TriggeringItemId { get; set; }

/// <summary>
/// Collection of server-side identifiers of monitored items that are
/// triggered by this item. Null if this item does not trigger any other items.
/// </summary>
internal UInt32Collection? TriggeredItems { get; set; }
}

/// <summary>
Expand Down
16 changes: 16 additions & 0 deletions Libraries/Opc.Ua.Client/Subscription/MonitoredItemState.cs
Original file line number Diff line number Diff line change
Expand Up @@ -79,6 +79,22 @@ public MonitoredItemState(MonitoredItemOptions options)
/// </summary>
[DataMember(Order = 15)]
public DateTime Timestamp { get; set; } = DateTime.UtcNow;

/// <summary>
/// Server-side identifier of the triggering item if this monitored item
/// is triggered by another item. 0 indicates this item is not triggered.
/// Used to restore triggering links after session reconnect.
/// </summary>
[DataMember(Order = 16)]
public uint TriggeringItemId { get; init; }

/// <summary>
/// Collection of server-side identifiers of monitored items that are
/// triggered by this item. Empty or null if this item does not trigger
/// any other items. Used to restore triggering links after session reconnect.
/// </summary>
[DataMember(Order = 17)]
public UInt32Collection? TriggeredItems { get; init; }
}

/// <summary>
Expand Down
201 changes: 201 additions & 0 deletions Libraries/Opc.Ua.Client/Subscription/Subscription.cs
Original file line number Diff line number Diff line change
Expand Up @@ -1037,6 +1037,9 @@ public async Task<IList<MonitoredItem>> CreateItemsAsync(CancellationToken ct =
m_changeMask |= SubscriptionChangeMask.ItemsCreated;
ChangesCompleted();

// Restore triggering relationships after items are created
await RestoreTriggeringAsync(ct).ConfigureAwait(false);

// return the list of items affected by the change.
return itemsToCreate;
}
Expand Down Expand Up @@ -1136,6 +1139,82 @@ public async Task<IList<MonitoredItem>> DeleteItemsAsync(
return itemsToDelete;
}

/// <summary>
/// Restores triggering relationships for monitored items that were
/// configured with triggers before reconnection.
/// </summary>
private async Task RestoreTriggeringAsync(CancellationToken ct = default)
{
VerifySessionAndSubscriptionState(true);

// Build triggering groups outside of lock to avoid await in lock
Dictionary<uint, List<uint>> triggeringGroups;
lock (m_cache)
{
// Group monitored items by their triggering item
triggeringGroups = new Dictionary<uint, List<uint>>();
foreach (MonitoredItem item in m_monitoredItems.Values)
{
if (item.TriggeredItems != null && item.TriggeredItems.Count > 0)
{
// This item triggers other items
var triggeredServerIds = new List<uint>();
foreach (uint triggeredClientHandle in item.TriggeredItems)
{
// Find the monitored item by client handle
if (m_monitoredItems.TryGetValue(triggeredClientHandle, out MonitoredItem? triggeredItem) &&
triggeredItem.Status.Created)
{
triggeredServerIds.Add(triggeredItem.Status.Id);
}
}

if (triggeredServerIds.Count > 0)
{
if (!triggeringGroups.TryGetValue(item.Status.Id, out List<uint>? list))
{
list = [];
triggeringGroups[item.Status.Id] = list;
}
list.AddRange(triggeredServerIds);
}
}
}
}

// Call SetTriggering for each triggering item
foreach (var kvp in triggeringGroups)
{
uint triggeringItemId = kvp.Key;
var linksToAdd = new UInt32Collection(kvp.Value);

try
{
await Session.SetTriggeringAsync(
null,
Id,
triggeringItemId,
linksToAdd,
null,
ct).ConfigureAwait(false);

m_logger.LogInformation(
"Restored {Count} triggering links for MonitoredItem {TriggeringItemId} in Subscription {SubscriptionId}",
linksToAdd.Count,
triggeringItemId,
Id);
}
catch (Exception ex)
{
m_logger.LogError(
ex,
"Failed to restore triggering links for MonitoredItem {TriggeringItemId} in Subscription {SubscriptionId}",
triggeringItemId,
Id);
}
}
}

/// <summary>
/// Set monitoring mode of items.
/// </summary>
Expand Down Expand Up @@ -1197,6 +1276,125 @@ public async Task<IList<MonitoredItem>> DeleteItemsAsync(
return errors;
}

/// <summary>
/// Sets the triggering relationships for a monitored item in this subscription
/// and tracks them for automatic restoration after reconnection.
/// </summary>
/// <param name="triggeringItem">The monitored item that will trigger other items.</param>
/// <param name="linksToAdd">Monitored items to be reported when the triggering item changes.</param>
/// <param name="linksToRemove">Monitored items to stop reporting when the triggering item changes.</param>
/// <param name="ct">Cancellation token.</param>
/// <returns>The response from the server.</returns>
/// <exception cref="ArgumentNullException">Thrown when triggeringItem is null.</exception>
/// <exception cref="ServiceResultException">Thrown when the operation fails.</exception>
public async Task<SetTriggeringResponse> SetTriggeringAsync(
MonitoredItem triggeringItem,
IList<MonitoredItem>? linksToAdd,
IList<MonitoredItem>? linksToRemove,
CancellationToken ct = default)
{
if (triggeringItem == null)
{
throw new ArgumentNullException(nameof(triggeringItem));
}

using Activity? activity = m_telemetry.StartActivity();
VerifySessionAndSubscriptionState(true);

if (!triggeringItem.Status.Created)
{
throw new ServiceResultException(
StatusCodes.BadInvalidState,
"Triggering item has not been created on the server.");
}

// Convert monitored items to server IDs
var serverIdsToAdd = new UInt32Collection();
var clientHandlesToAdd = new UInt32Collection();
if (linksToAdd != null)
{
foreach (MonitoredItem item in linksToAdd)
{
if (!item.Status.Created)
{
throw new ServiceResultException(
StatusCodes.BadInvalidState,
$"Monitored item '{item.DisplayName}' has not been created on the server.");
}
serverIdsToAdd.Add(item.Status.Id);
clientHandlesToAdd.Add(item.ClientHandle);
}
}

var serverIdsToRemove = new UInt32Collection();
var clientHandlesToRemove = new UInt32Collection();
if (linksToRemove != null)
{
foreach (MonitoredItem item in linksToRemove)
{
if (!item.Status.Created)
{
throw new ServiceResultException(
StatusCodes.BadInvalidState,
$"Monitored item '{item.DisplayName}' has not been created on the server.");
}
serverIdsToRemove.Add(item.Status.Id);
clientHandlesToRemove.Add(item.ClientHandle);
}
}

// Call the Session SetTriggering method
SetTriggeringResponse response = await Session.SetTriggeringAsync(
null,
Id,
triggeringItem.Status.Id,
serverIdsToAdd,
serverIdsToRemove,
ct).ConfigureAwait(false);

// Update the triggering relationships for automatic restoration
lock (m_cache)
{
// Initialize the triggered items collection if needed
triggeringItem.TriggeredItems ??= new UInt32Collection();

// Add new links
if (clientHandlesToAdd.Count > 0)
{
foreach (uint clientHandle in clientHandlesToAdd)
{
if (!triggeringItem.TriggeredItems.Contains(clientHandle))
{
triggeringItem.TriggeredItems.Add(clientHandle);
}

// Update the triggered item to remember its triggering item
if (m_monitoredItems.TryGetValue(clientHandle, out MonitoredItem? triggeredItem))
{
triggeredItem.TriggeringItemId = triggeringItem.Status.Id;
}
}
}

// Remove links
if (clientHandlesToRemove.Count > 0)
{
foreach (uint clientHandle in clientHandlesToRemove)
{
triggeringItem.TriggeredItems.Remove(clientHandle);

// Clear the triggering item reference
if (m_monitoredItems.TryGetValue(clientHandle, out MonitoredItem? triggeredItem))
{
triggeredItem.TriggeringItemId = 0;
}
}
}
}

return response;
}

/// <summary>
/// Tells the server to refresh all conditions being monitored by the subscription.
/// </summary>
Expand Down Expand Up @@ -1349,6 +1547,9 @@ await session.RemoveSubscriptionsAsync(subscriptionsToRemove, ct)
m_changeMask |= SubscriptionChangeMask.Transferred;
ChangesCompleted();

// Restore triggering relationships after subscription transfer
await RestoreTriggeringAsync(ct).ConfigureAwait(false);

StartKeepAliveTimer();

TraceState("TRANSFERRED ASYNC");
Expand Down
94 changes: 94 additions & 0 deletions Tests/Opc.Ua.Client.Tests/SubscriptionTest.cs
Original file line number Diff line number Diff line change
Expand Up @@ -1462,5 +1462,99 @@ private void DeferSubscriptionAcknowledge(
e.DeferredAcknowledgementsToSend.Clear();
e.AcknowledgementsToSend.Clear();
}

[Test]
[Order(900)]
public async Task SetTriggeringTrackingAsync()
{
// Create a subscription
var subscription = new Subscription(Session.DefaultSubscription)
{
PublishingEnabled = true,
PublishingInterval = 1000,
KeepAliveCount = 10,
LifetimeCount = 100,
MaxNotificationsPerPublish = 1000,
Priority = 100
};

Session.AddSubscription(subscription);
await subscription.CreateAsync(CancellationToken.None).ConfigureAwait(false);
Assert.That(subscription.Created, Is.True);

// Create monitored items
var triggeringItem = new MonitoredItem(subscription.DefaultItem)
{
StartNodeId = VariableIds.Server_ServerStatus_CurrentTime,
AttributeId = Attributes.Value,
MonitoringMode = MonitoringMode.Reporting,
SamplingInterval = 0,
QueueSize = 0,
DiscardOldest = true
};

var triggeredItem1 = new MonitoredItem(subscription.DefaultItem)
{
StartNodeId = VariableIds.Server_ServerStatus_State,
AttributeId = Attributes.Value,
MonitoringMode = MonitoringMode.Sampling,
SamplingInterval = 0,
QueueSize = 0,
DiscardOldest = true
};

var triggeredItem2 = new MonitoredItem(subscription.DefaultItem)
{
StartNodeId = VariableIds.Server_ServerStatus_BuildInfo,
AttributeId = Attributes.Value,
MonitoringMode = MonitoringMode.Sampling,
SamplingInterval = 0,
QueueSize = 0,
DiscardOldest = true
};

subscription.AddItem(triggeringItem);
subscription.AddItem(triggeredItem1);
subscription.AddItem(triggeredItem2);

// Create the items
await subscription.ApplyChangesAsync(CancellationToken.None).ConfigureAwait(false);

Assert.That(triggeringItem.Created, Is.True);
Assert.That(triggeredItem1.Created, Is.True);
Assert.That(triggeredItem2.Created, Is.True);

// Set up triggering relationship using the new method
var linksToAdd = new List<MonitoredItem> { triggeredItem1, triggeredItem2 };
SetTriggeringResponse response = await subscription.SetTriggeringAsync(
triggeringItem,
linksToAdd,
null,
CancellationToken.None).ConfigureAwait(false);

Assert.That(response, Is.Not.Null);

// Verify the triggering relationships are tracked
Assert.That(triggeringItem.TriggeredItems, Is.Not.Null);
Assert.That(triggeringItem.TriggeredItems.Count, Is.EqualTo(2));
Assert.That(triggeringItem.TriggeredItems, Does.Contain(triggeredItem1.ClientHandle));
Assert.That(triggeringItem.TriggeredItems, Does.Contain(triggeredItem2.ClientHandle));

Assert.That(triggeredItem1.TriggeringItemId, Is.EqualTo(triggeringItem.Status.Id));
Assert.That(triggeredItem2.TriggeringItemId, Is.EqualTo(triggeringItem.Status.Id));

// Snapshot the subscription state
subscription.Snapshot(out SubscriptionState state);

// Verify that the triggering relationships are persisted
MonitoredItemState triggeringItemState = state.MonitoredItems
.FirstOrDefault(m => m.ClientId == triggeringItem.ClientHandle);
Assert.That(triggeringItemState, Is.Not.Null);
Assert.That(triggeringItemState.TriggeredItems, Is.Not.Null);
Assert.That(triggeringItemState.TriggeredItems.Count, Is.EqualTo(2));

// Clean up
await subscription.DeleteAsync(true, CancellationToken.None).ConfigureAwait(false);
}
}
}
Loading