diff --git a/plc4j/drivers/opcua/src/main/java/org/apache/plc4x/java/opcua/protocol/OpcuaProtocolLogic.java b/plc4j/drivers/opcua/src/main/java/org/apache/plc4x/java/opcua/protocol/OpcuaProtocolLogic.java index c8f735fbff6..41d8c94cd70 100644 --- a/plc4j/drivers/opcua/src/main/java/org/apache/plc4x/java/opcua/protocol/OpcuaProtocolLogic.java +++ b/plc4j/drivers/opcua/src/main/java/org/apache/plc4x/java/opcua/protocol/OpcuaProtocolLogic.java @@ -811,6 +811,15 @@ public CompletableFuture subscribe(PlcSubscriptionReque long subscriptionId = response.getSubscriptionId(); OpcuaSubscriptionHandle handle = new OpcuaSubscriptionHandle(this, tm, conversation, subscriptionRequest, subscriptionId, cycleTime); + if (subscriptionRequest.getConsumer() != null) { + handle.register(subscriptionRequest.getConsumer()); + } + subscriptionRequest.getTagNames().forEach(tagName -> { + Consumer tagConsumer = subscriptionRequest.getTagConsumer(tagName); + if (tagConsumer != null) { + handle.registerTagConsumer(tagName, tagConsumer); + } + }); subscriptions.put(handle.getSubscriptionId(), handle); return handle; }) diff --git a/plc4j/drivers/opcua/src/main/java/org/apache/plc4x/java/opcua/protocol/OpcuaSubscriptionHandle.java b/plc4j/drivers/opcua/src/main/java/org/apache/plc4x/java/opcua/protocol/OpcuaSubscriptionHandle.java index 19e248a346f..dd159ff334d 100644 --- a/plc4j/drivers/opcua/src/main/java/org/apache/plc4x/java/opcua/protocol/OpcuaSubscriptionHandle.java +++ b/plc4j/drivers/opcua/src/main/java/org/apache/plc4x/java/opcua/protocol/OpcuaSubscriptionHandle.java @@ -64,6 +64,7 @@ public class OpcuaSubscriptionHandle extends DefaultPlcSubscriptionHandle { private final Logger logger = LoggerFactory.getLogger(OpcuaSubscriptionHandle.class); private final Set> consumers; + private final Map> tagConsumers; private final List tagNames; private final Conversation conversation; private final PlcSubscriptionRequest subscriptionRequest; @@ -83,6 +84,7 @@ public OpcuaSubscriptionHandle(OpcuaProtocolLogic plcSubscriber, RequestTransact super(plcSubscriber); this.tm = tm; this.consumers = new HashSet<>(); + this.tagConsumers = new HashMap<>(); this.subscriptionRequest = subscriptionRequest; this.tagNames = new ArrayList<>(subscriptionRequest.getTagNames()); this.conversation = conversation; @@ -292,8 +294,15 @@ private void onMonitoredValue(List values) { Map tagMap = new LinkedHashMap<>(); for (MonitoredItemNotification value : values) { String tagName = tagNames.get((int) value.getClientHandle() - 1); - tagMap.put(tagName, subscriptionRequest.getTag(tagName).getTag()); + PlcTag tag = subscriptionRequest.getTag(tagName).getTag(); + tagMap.put(tagName, tag); dataValues.add(value.getValue()); + Consumer tagConsumer = tagConsumers.get(tagName); + if (tagConsumer != null) { + Entry, Map>> mappedResponse = plcSubscriber.readResponse(Map.of(tagName, tag), List.of(value.getValue()), responseMetadata); + PlcSubscriptionEvent event = new DefaultPlcSubscriptionEvent(Instant.ofEpochMilli(receiveTs), mappedResponse.getValue(), mappedResponse.getKey()); + tagConsumer.accept(event); + } } Entry, Map>> mappedResponse = plcSubscriber.readResponse(tagMap, dataValues, responseMetadata); @@ -346,6 +355,12 @@ public PlcConsumerRegistration register(Consumer consumer) return new DefaultPlcConsumerRegistration(plcSubscriber, consumer, this); } + public PlcConsumerRegistration registerTagConsumer(String tagName, Consumer consumer) { + logger.info("Registering a new OPCUA subscription consumer for tag with name " + tagName); + tagConsumers.put(tagName, consumer); + return new DefaultPlcConsumerRegistration(plcSubscriber, consumer, this); + } + public Long getSubscriptionId() { return subscriptionId; }