This is an automated email from the ASF dual-hosted git repository.

dominikriemer pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/streampipes.git


The following commit(s) were added to refs/heads/dev by this push:
     new 80cd72594f feat: Extend OPC UA subscription mode configuration (#4659)
80cd72594f is described below

commit 80cd72594f816f25fbad9a607cb7f58451b411a1
Author: Dominik Riemer <[email protected]>
AuthorDate: Thu Jul 2 16:03:25 2026 +0200

    feat: Extend OPC UA subscription mode configuration (#4659)
---
 .../apache/streampipes/commons/constants/Envs.java |   1 +
 .../commons/environment/DefaultEnvironment.java    |   5 +
 .../commons/environment/Environment.java           |   2 +
 .../opcua/OpcUaConnectorsModuleExport.java         |   2 +
 .../connectors/opcua/adapter/OpcUaAdapter.java     |  88 ++++++---
 .../connectors/opcua/adapter/OpcUaNodeBrowser.java |   2 +-
 .../opcua/client/ConnectedOpcUaClient.java         |  31 ++--
 .../opcua/config/OpcUaAdapterConfig.java           |  41 +++++
 .../opcua/config/SharedUserConfiguration.java      |  61 +++++-
 .../opcua/config/SpOpcUaConfigExtractor.java       |  54 ++++++
 .../opcua/config/security/SecurityConfig.java      |  21 +++
 .../opcua/migration/OpcUaAdapterMigrationV7.java   |  59 ++++++
 .../connectors/opcua/utils/OpcUaLabels.java        |   4 +
 .../documentation.md                               |   2 +
 .../strings.en                                     |  15 ++
 .../documentation.md                               |   2 +
 .../adapter/OpcUaAdapterSubscriptionTest.java      | 204 +++++++++++++++++++++
 .../opcua/config/SpOpcUaConfigExtractorTest.java   |  48 +++++
 .../opcua/config/security/SecurityConfigTest.java  |  98 +++++++++-
 19 files changed, 695 insertions(+), 45 deletions(-)

diff --git 
a/streampipes-commons/src/main/java/org/apache/streampipes/commons/constants/Envs.java
 
b/streampipes-commons/src/main/java/org/apache/streampipes/commons/constants/Envs.java
index 3cb48ac964..d73cdfbda0 100644
--- 
a/streampipes-commons/src/main/java/org/apache/streampipes/commons/constants/Envs.java
+++ 
b/streampipes-commons/src/main/java/org/apache/streampipes/commons/constants/Envs.java
@@ -138,6 +138,7 @@ public enum Envs {
   SP_OPCUA_KEYSTORE_PASSWORD("SP_OPCUA_KEYSTORE_PASSWORD", "password"),
   SP_OPCUA_KEYSTORE_TYPE("SP_OPCUA_KEYSTORE_TYPE", "PKCS12"),
   SP_OPCUA_KEYSTORE_ALIAS("SP_OPCUA_KEYSTORE_ALIAS", "apache-streampipes"),
+  SP_OPCUA_DISALLOW_INSECURE_ENDPOINTS("SP_OPCUA_DISALLOW_INSECURE_ENDPOINTS", 
"false"),
   SP_OPCUA_APPLICATION_URI(
       "SP_OPCUA_APPLICATION_URI", "urn:org:apache:streampipes:opcua:client"
   ),
diff --git 
a/streampipes-commons/src/main/java/org/apache/streampipes/commons/environment/DefaultEnvironment.java
 
b/streampipes-commons/src/main/java/org/apache/streampipes/commons/environment/DefaultEnvironment.java
index e0e00d694c..09f3877eb7 100644
--- 
a/streampipes-commons/src/main/java/org/apache/streampipes/commons/environment/DefaultEnvironment.java
+++ 
b/streampipes-commons/src/main/java/org/apache/streampipes/commons/environment/DefaultEnvironment.java
@@ -377,6 +377,11 @@ public class DefaultEnvironment implements Environment {
     return new StringEnvironmentVariable(Envs.SP_OPCUA_KEYSTORE_ALIAS);
   }
 
+  @Override
+  public BooleanEnvironmentVariable getOpcUaDisallowInsecureEndpoints() {
+    return new 
BooleanEnvironmentVariable(Envs.SP_OPCUA_DISALLOW_INSECURE_ENDPOINTS);
+  }
+
   @Override
   public IntEnvironmentVariable getOpcUaMinPullIntervalMs() {
     return new IntEnvironmentVariable(Envs.SP_OPCUA_MIN_PULL_INTERVAL_MS);
diff --git 
a/streampipes-commons/src/main/java/org/apache/streampipes/commons/environment/Environment.java
 
b/streampipes-commons/src/main/java/org/apache/streampipes/commons/environment/Environment.java
index c1d50e3057..2460a2d10b 100644
--- 
a/streampipes-commons/src/main/java/org/apache/streampipes/commons/environment/Environment.java
+++ 
b/streampipes-commons/src/main/java/org/apache/streampipes/commons/environment/Environment.java
@@ -182,6 +182,8 @@ public interface Environment {
 
   StringEnvironmentVariable getOpcUaKeystoreAlias();
 
+  BooleanEnvironmentVariable getOpcUaDisallowInsecureEndpoints();
+
   IntEnvironmentVariable getOpcUaMinPullIntervalMs();
 
   StringEnvironmentVariable getKeystoreFilename();
diff --git 
a/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/OpcUaConnectorsModuleExport.java
 
b/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/OpcUaConnectorsModuleExport.java
index 0c7bf2088f..9ea53e469b 100644
--- 
a/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/OpcUaConnectorsModuleExport.java
+++ 
b/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/OpcUaConnectorsModuleExport.java
@@ -31,6 +31,7 @@ import 
org.apache.streampipes.extensions.connectors.opcua.migration.OpcUaAdapter
 import 
org.apache.streampipes.extensions.connectors.opcua.migration.OpcUaAdapterMigrationV4;
 import 
org.apache.streampipes.extensions.connectors.opcua.migration.OpcUaAdapterMigrationV5;
 import 
org.apache.streampipes.extensions.connectors.opcua.migration.OpcUaAdapterMigrationV6;
+import 
org.apache.streampipes.extensions.connectors.opcua.migration.OpcUaAdapterMigrationV7;
 import 
org.apache.streampipes.extensions.connectors.opcua.migration.OpcUaSinkMigrationV1;
 import 
org.apache.streampipes.extensions.connectors.opcua.migration.OpcUaSinkMigrationV2;
 import org.apache.streampipes.extensions.connectors.opcua.sink.OpcUaSink;
@@ -69,6 +70,7 @@ public class OpcUaConnectorsModuleExport implements 
IExtensionModuleExport {
         new OpcUaAdapterMigrationV4(),
         new OpcUaAdapterMigrationV5(),
         new OpcUaAdapterMigrationV6(),
+        new OpcUaAdapterMigrationV7(),
         new OpcUaSinkMigrationV1(),
         new OpcUaSinkMigrationV2()
     );
diff --git 
a/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/adapter/OpcUaAdapter.java
 
b/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/adapter/OpcUaAdapter.java
index 33571ec394..20799195a6 100644
--- 
a/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/adapter/OpcUaAdapter.java
+++ 
b/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/adapter/OpcUaAdapter.java
@@ -84,14 +84,14 @@ public class OpcUaAdapter implements StreamPipesAdapter, 
IPullAdapter, SupportsR
   private int numberOfEventProperties = 0;
 
   /**
-   * This variable is used to map the node ids during the subscription to the 
labels of the nodes
+   * This variable is used to map the node ids during the subscription to the 
selected nodes.
    */
-  private final Map<String, String> nodeIdToLabelMapping;
+  private final Map<String, OpcUaNode> nodeIdToNodeMapping;
 
   public OpcUaAdapter(OpcUaClientProvider clientProvider) {
     this.clientProvider = clientProvider;
     this.event = new HashMap<>();
-    this.nodeIdToLabelMapping = new HashMap<>();
+    this.nodeIdToNodeMapping = new HashMap<>();
   }
 
   private void prepareAdapter() throws AdapterException {
@@ -130,14 +130,10 @@ public class OpcUaAdapter implements StreamPipesAdapter, 
IPullAdapter, SupportsR
         }
 
         this.pullingIntervalMilliSeconds = effectiveIntervalMs;
-      } else {
-        var allNodeIds = this.allNodes.stream()
-            .map(node -> node.nodeInfo().getNodeId()).toList();
-        this.connectedClient.createListSubscription(allNodeIds, this);
       }
 
-      this.allNodes.forEach(node -> this.nodeIdToLabelMapping
-          .put(node.nodeInfo().getNodeId().toString(), 
node.nodeInfo().getDisplayName()));
+      this.allNodes.forEach(node -> this.nodeIdToNodeMapping
+          .put(node.nodeInfo().getNodeId().toString(), node));
 
     } catch (Exception e) {
       var errorMessage = buildStartupErrorMessage(e);
@@ -274,29 +270,61 @@ public class OpcUaAdapter implements StreamPipesAdapter, 
IPullAdapter, SupportsR
   public void onSubscriptionValue(OpcUaMonitoredItem item,
                                   DataValue value) {
 
-    String key = 
this.nodeIdToLabelMapping.get(item.getReadValueId().getNodeId().toString());
-
-    var currNode = this.allNodes.stream()
-        .filter(node -> key.equals(node.nodeInfo().getDisplayName()))
-        .findFirst()
-        .orElse(null);
+    var nodeId = item.getReadValueId().getNodeId();
+    var currNode = this.nodeIdToNodeMapping.get(nodeId.toString());
 
     if (currNode != null) {
-      currNode.addToEvent(connectedClient.getClient(), event, 
value.getValue());
-      // ensure that event is complete and all opc ua subscriptions 
transmitted at least one value
-      if (event.size() >= numberOfEventProperties) {
-        Map<String, Object> newEvent = new HashMap<>();
-        // deep copy of event to prevent preprocessor error
-        for (String k : event.keySet()) {
-          newEvent.put(k, event.get(k));
+      if (value == null) {
+        LOG.warn("Received null DataValue for OPC UA subscription node label: 
{}, node id: {}",
+            currNode.nodeInfo().getDisplayName(),
+            nodeId);
+        return;
+      }
+
+      var status = value.getStatusCode();
+      if (status != null && status.isGood()) {
+        synchronized (event) {
+          currNode.addToEvent(connectedClient.getClient(), event, 
value.getValue());
+          publishSubscriptionEventIfComplete();
+        }
+      } else if (shouldSkipEvent(true)) {
+        LOG.warn("Received status code {} for OPC UA subscription node label: 
{}, node id: {}",
+            status,
+            currNode.nodeInfo().getDisplayName(),
+            nodeId);
+      } else {
+        synchronized (event) {
+          event.remove(currNode.nodeInfo().getDesiredName(""));
+          publishSubscriptionEvent();
         }
-        collector.collect(newEvent);
       }
     } else {
       LOG.error("No event is produced, because subscription item {} could not 
be found within all nodes", item);
     }
   }
 
+  private void publishSubscriptionEvent() {
+    if (!event.isEmpty()) {
+      collector.collect(copySubscriptionEvent());
+    }
+  }
+
+  private void publishSubscriptionEventIfComplete() {
+    // ensure that event is complete and all opc ua subscriptions transmitted 
at least one value
+    if (event.size() >= numberOfEventProperties) {
+      collector.collect(copySubscriptionEvent());
+    }
+  }
+
+  private Map<String, Object> copySubscriptionEvent() {
+    Map<String, Object> newEvent = new HashMap<>();
+    // deep copy of event to prevent preprocessor error
+    for (String k : event.keySet()) {
+      newEvent.put(k, event.get(k));
+    }
+    return newEvent;
+  }
+
   @Override
   public PollingSettings getPollingInterval() {
     return PollingSettings.from(TimeUnit.MILLISECONDS, 
this.pullingIntervalMilliSeconds);
@@ -320,6 +348,15 @@ public class OpcUaAdapter implements StreamPipesAdapter, 
IPullAdapter, SupportsR
     if (this.opcUaAdapterConfig.inPullMode()) {
       this.pullAdapterScheduler = new PullAdapterScheduler();
       this.pullAdapterScheduler.schedule(this, 
extractor.getAdapterDescription().getElementId());
+    } else {
+      var allNodeIds = this.allNodes.stream()
+          .map(node -> node.nodeInfo().getNodeId()).toList();
+      try {
+        this.connectedClient.createListSubscription(allNodeIds, 
this.opcUaAdapterConfig, this);
+      } catch (Exception e) {
+        var errorMessage = buildStartupErrorMessage(e);
+        throw new AdapterException(errorMessage, e);
+      }
     }
   }
 
@@ -341,14 +378,15 @@ public class OpcUaAdapter implements StreamPipesAdapter, 
IPullAdapter, SupportsR
 
   @Override
   public IAdapterConfiguration declareConfig() {
-    var builder = AdapterConfigurationBuilder.create(ID, 6, () -> new 
OpcUaAdapter(clientProvider))
+    var builder = AdapterConfigurationBuilder.create(ID, 7, () -> new 
OpcUaAdapter(clientProvider))
         .withAssets(ExtensionAssetType.DOCUMENTATION, ExtensionAssetType.ICON)
         .withLocales(Locales.EN)
         .requiredAlternatives(Labels.withId(ADAPTER_TYPE),
             Alternatives.from(Labels.withId(PULL_MODE),
                 SharedUserConfiguration.getPullModeGroup()
             ),
-            Alternatives.from(Labels.withId(SUBSCRIPTION_MODE)));
+            Alternatives.from(Labels.withId(SUBSCRIPTION_MODE),
+                SharedUserConfiguration.getSubscriptionModeGroup()));
     SharedUserConfiguration.appendSharedOpcUaConfig(builder, true);
     
builder.requiredStaticProperty(SharedUserConfiguration.makeNamingStrategyOption());
     return builder.buildConfiguration();
diff --git 
a/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/adapter/OpcUaNodeBrowser.java
 
b/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/adapter/OpcUaNodeBrowser.java
index 5e366af76e..cdaa3ecbe1 100644
--- 
a/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/adapter/OpcUaNodeBrowser.java
+++ 
b/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/adapter/OpcUaNodeBrowser.java
@@ -104,7 +104,7 @@ public class OpcUaNodeBrowser {
       );
     }
 
-    LOG.info(
+    LOG.debug(
         "Using node of type {}",
         node.getNodeClass()
             .toString()
diff --git 
a/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/client/ConnectedOpcUaClient.java
 
b/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/client/ConnectedOpcUaClient.java
index ddb27e5bd3..1441695cbe 100644
--- 
a/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/client/ConnectedOpcUaClient.java
+++ 
b/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/client/ConnectedOpcUaClient.java
@@ -19,6 +19,7 @@
 package org.apache.streampipes.extensions.connectors.opcua.client;
 
 import org.apache.streampipes.extensions.connectors.opcua.adapter.OpcUaAdapter;
+import 
org.apache.streampipes.extensions.connectors.opcua.config.OpcUaAdapterConfig;
 
 import org.eclipse.milo.opcua.sdk.client.OpcUaClient;
 import 
org.eclipse.milo.opcua.sdk.client.subscriptions.MonitoredItemServiceOperationResult;
@@ -27,7 +28,6 @@ import 
org.eclipse.milo.opcua.sdk.client.subscriptions.OpcUaSubscription;
 import org.eclipse.milo.opcua.stack.core.UaException;
 import org.eclipse.milo.opcua.stack.core.types.builtin.NodeId;
 import org.eclipse.milo.opcua.stack.core.types.builtin.StatusCode;
-import org.eclipse.milo.opcua.stack.core.types.enumerated.TimestampsToReturn;
 import org.jspecify.annotations.NonNull;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -49,32 +49,28 @@ public class ConnectedOpcUaClient {
   /***
    * Register subscriptions for given OPC UA nodes
    * @param nodes List of {@link 
org.eclipse.milo.opcua.stack.core.types.builtin.NodeId}
+   * @param config subscription configuration
    * @param opcUaAdapter current instance of {@link OpcUaAdapter}
    * @throws Exception
    */
   public void createListSubscription(List<NodeId> nodes,
+                                     OpcUaAdapterConfig config,
                                      OpcUaAdapter opcUaAdapter) throws 
Exception {
-    initSubscription(nodes, opcUaAdapter);
+    initSubscription(nodes, config, opcUaAdapter);
   }
 
 
   public void initSubscription(List<NodeId> nodes,
+                               OpcUaAdapterConfig config,
                                OpcUaAdapter opcUaAdapter) throws Exception {
-    var subscription = getOpcUaSubscription(nodes, opcUaAdapter);
-
-    for (NodeId node : nodes) {
-      var value = this.client.readValue(0, TimestampsToReturn.Both, node);
-      if (value == null || value.getValue().getValue() == null) {
-        LOG.error("Node has no value");
-      }
-    }
+    var subscription = getOpcUaSubscription(nodes, config, opcUaAdapter);
 
     List<OpcUaMonitoredItem> items = new ArrayList<>();
     for (NodeId node : nodes) {
       var item = OpcUaMonitoredItem.newDataItem(node);
-      item.setSamplingInterval(1000.0);
-      item.setQueueSize(uint(10));
-      item.setDiscardOldest(true);
+      item.setSamplingInterval((double) 
config.getSubscriptionSamplingIntervalMs());
+      item.setQueueSize(uint(config.getSubscriptionQueueSize()));
+      item.setDiscardOldest(config.isSubscriptionDiscardOldest());
       item.setDataValueListener(opcUaAdapter::onSubscriptionValue);
       items.add(item);
     }
@@ -95,14 +91,15 @@ public class ConnectedOpcUaClient {
   }
 
   private @NonNull OpcUaSubscription getOpcUaSubscription(List<NodeId> nodes,
+                                                          OpcUaAdapterConfig 
config,
                                                           OpcUaAdapter 
opcUaAdapter) throws UaException {
-    OpcUaSubscription subscription = createManagedSubscription();
+    OpcUaSubscription subscription = createManagedSubscription(config);
     subscription.setSubscriptionListener(new 
OpcUaSubscription.SubscriptionListener() {
       @Override
       public void onTransferFailed(OpcUaSubscription subscription, StatusCode 
statusCode) {
         LOG.warn("Transfer for subscriptionId={} failed: {}", 
subscription.getSubscriptionId(), statusCode);
         try {
-          initSubscription(nodes, opcUaAdapter);
+          initSubscription(nodes, config, opcUaAdapter);
         } catch (Exception e) {
           LOG.error("Re-creating the subscription failed", e);
         }
@@ -111,8 +108,8 @@ public class ConnectedOpcUaClient {
     return subscription;
   }
 
-  private OpcUaSubscription createManagedSubscription() throws UaException {
-    var subscription = new OpcUaSubscription(this.client, 1000.0);
+  private OpcUaSubscription createManagedSubscription(OpcUaAdapterConfig 
config) throws UaException {
+    var subscription = new OpcUaSubscription(this.client, (double) 
config.getSubscriptionPublishingIntervalMs());
     subscription.create();
     return subscription;
   }
diff --git 
a/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/config/OpcUaAdapterConfig.java
 
b/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/config/OpcUaAdapterConfig.java
index a0bac48efd..14b5683f33 100644
--- 
a/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/config/OpcUaAdapterConfig.java
+++ 
b/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/config/OpcUaAdapterConfig.java
@@ -22,9 +22,18 @@ import 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaNamingStrat
 
 public class OpcUaAdapterConfig extends OpcUaConfig {
 
+  public static final int DEFAULT_SUBSCRIPTION_PUBLISHING_INTERVAL_MS = 1000;
+  public static final int DEFAULT_SUBSCRIPTION_SAMPLING_INTERVAL_MS = 1000;
+  public static final int DEFAULT_SUBSCRIPTION_QUEUE_SIZE = 10;
+  public static final boolean DEFAULT_SUBSCRIPTION_DISCARD_OLDEST = true;
+
   private Integer pullIntervalMilliSeconds;
   private String incompleteEventStrategy;
   private OpcUaNamingStrategy namingStrategy;
+  private int subscriptionPublishingIntervalMs = 
DEFAULT_SUBSCRIPTION_PUBLISHING_INTERVAL_MS;
+  private int subscriptionSamplingIntervalMs = 
DEFAULT_SUBSCRIPTION_SAMPLING_INTERVAL_MS;
+  private int subscriptionQueueSize = DEFAULT_SUBSCRIPTION_QUEUE_SIZE;
+  private boolean subscriptionDiscardOldest = 
DEFAULT_SUBSCRIPTION_DISCARD_OLDEST;
 
   public Integer getPullIntervalMilliSeconds() {
     return pullIntervalMilliSeconds;
@@ -54,4 +63,36 @@ public class OpcUaAdapterConfig extends OpcUaConfig {
     this.namingStrategy = namingStrategy;
   }
 
+  public int getSubscriptionPublishingIntervalMs() {
+    return subscriptionPublishingIntervalMs;
+  }
+
+  public void setSubscriptionPublishingIntervalMs(int 
subscriptionPublishingIntervalMs) {
+    this.subscriptionPublishingIntervalMs = subscriptionPublishingIntervalMs;
+  }
+
+  public int getSubscriptionSamplingIntervalMs() {
+    return subscriptionSamplingIntervalMs;
+  }
+
+  public void setSubscriptionSamplingIntervalMs(int 
subscriptionSamplingIntervalMs) {
+    this.subscriptionSamplingIntervalMs = subscriptionSamplingIntervalMs;
+  }
+
+  public int getSubscriptionQueueSize() {
+    return subscriptionQueueSize;
+  }
+
+  public void setSubscriptionQueueSize(int subscriptionQueueSize) {
+    this.subscriptionQueueSize = subscriptionQueueSize;
+  }
+
+  public boolean isSubscriptionDiscardOldest() {
+    return subscriptionDiscardOldest;
+  }
+
+  public void setSubscriptionDiscardOldest(boolean subscriptionDiscardOldest) {
+    this.subscriptionDiscardOldest = subscriptionDiscardOldest;
+  }
+
 }
diff --git 
a/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/config/SharedUserConfiguration.java
 
b/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/config/SharedUserConfiguration.java
index dbf25d83c6..e939b641ff 100644
--- 
a/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/config/SharedUserConfiguration.java
+++ 
b/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/config/SharedUserConfiguration.java
@@ -21,8 +21,10 @@ package 
org.apache.streampipes.extensions.connectors.opcua.config;
 import org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels;
 import 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaNamingStrategy;
 import org.apache.streampipes.extensions.connectors.opcua.utils.SecurityUtils;
+import org.apache.streampipes.model.staticproperty.FreeTextStaticProperty;
 import org.apache.streampipes.model.staticproperty.OneOfStaticProperty;
 import org.apache.streampipes.model.staticproperty.Option;
+import org.apache.streampipes.model.staticproperty.PropertyValueSpecification;
 import org.apache.streampipes.model.staticproperty.StaticPropertyGroup;
 import org.apache.streampipes.sdk.StaticProperties;
 import 
org.apache.streampipes.sdk.builder.AbstractConfigurablePipelineElementBuilder;
@@ -43,6 +45,10 @@ import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabe
 import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.OPC_URL;
 import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.PASSWORD;
 import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.PULLING_INTERVAL;
+import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.SUBSCRIPTION_DISCARD_OLDEST;
+import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.SUBSCRIPTION_PUBLISHING_INTERVAL;
+import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.SUBSCRIPTION_QUEUE_SIZE;
+import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.SUBSCRIPTION_SAMPLING_INTERVAL;
 import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.USERNAME;
 import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.USERNAME_GROUP;
 
@@ -51,6 +57,9 @@ public class SharedUserConfiguration {
   public static final String INCOMPLETE_EVENT_HANDLING_KEY = 
"incomplete-event-handling";
   public static final String INCOMPLETE_OPTION_IGNORE = "ignore-event";
   public static final String INCOMPLETE_OPTION_SEND = "send-event";
+  public static final String SUBSCRIPTION_GROUP = "subscription-mode-group";
+  public static final String SUBSCRIPTION_DISCARD_OLDEST_TRUE = "true";
+  public static final String SUBSCRIPTION_DISCARD_OLDEST_FALSE = "false";
 
   public static final String SECURITY_MODE = "securityMode";
   public static final String SECURITY_POLICY = "securityPolicy";
@@ -146,14 +155,64 @@ public class SharedUserConfiguration {
     return group;
   }
 
+  public static StaticPropertyGroup getSubscriptionModeGroup() {
+    var group = StaticProperties.group(
+        Labels.withId(SUBSCRIPTION_GROUP),
+        false,
+        positiveIntegerProperty(
+            SUBSCRIPTION_PUBLISHING_INTERVAL,
+            OpcUaAdapterConfig.DEFAULT_SUBSCRIPTION_PUBLISHING_INTERVAL_MS,
+            100,
+            3600000),
+        positiveIntegerProperty(
+            SUBSCRIPTION_SAMPLING_INTERVAL,
+            OpcUaAdapterConfig.DEFAULT_SUBSCRIPTION_SAMPLING_INTERVAL_MS,
+            100,
+            3600000),
+        positiveIntegerProperty(
+            SUBSCRIPTION_QUEUE_SIZE,
+            OpcUaAdapterConfig.DEFAULT_SUBSCRIPTION_QUEUE_SIZE,
+            1,
+            100),
+        getSubscriptionDiscardOldestConfig(),
+        getIncompleteEventConfig()
+    );
+    group.setHorizontalRendering(false);
+    return group;
+  }
+
+  private static FreeTextStaticProperty positiveIntegerProperty(OpcUaLabels 
label,
+                                                                int 
defaultValue,
+                                                                int min,
+                                                                int max) {
+    var property = 
StaticProperties.integerFreeTextProperty(Labels.withId(label), defaultValue);
+    property.setValueSpecification(new PropertyValueSpecification(min, max, 
1));
+    return property;
+  }
+
+  public static OneOfStaticProperty getSubscriptionDiscardOldestConfig() {
+    var property = StaticProperties.singleValueSelection(
+        Labels.withId(SUBSCRIPTION_DISCARD_OLDEST),
+        List.of(
+            new Option("Discard oldest", SUBSCRIPTION_DISCARD_OLDEST_TRUE),
+            new Option("Discard newest", SUBSCRIPTION_DISCARD_OLDEST_FALSE)
+        )
+    );
+    
property.getOptions().get(0).setSelected(OpcUaAdapterConfig.DEFAULT_SUBSCRIPTION_DISCARD_OLDEST);
+    
property.getOptions().get(1).setSelected(!OpcUaAdapterConfig.DEFAULT_SUBSCRIPTION_DISCARD_OLDEST);
+    return property;
+  }
+
   public static OneOfStaticProperty getIncompleteEventConfig() {
-    return StaticProperties.singleValueSelection(
+    var property = StaticProperties.singleValueSelection(
         Labels.withId(INCOMPLETE_EVENT_HANDLING_KEY),
         List.of(
             new Option("Ignore (only complete messages are sent)", 
INCOMPLETE_OPTION_IGNORE),
             new Option("Send (incomplete messages are sent)", 
INCOMPLETE_OPTION_SEND)
         )
     );
+    property.getOptions().get(0).setSelected(true);
+    return property;
   }
 
   public static List<String> getDependsOn(boolean adapterConfig) {
diff --git 
a/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/config/SpOpcUaConfigExtractor.java
 
b/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/config/SpOpcUaConfigExtractor.java
index fae0dfcde8..6c7c8b3491 100644
--- 
a/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/config/SpOpcUaConfigExtractor.java
+++ 
b/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/config/SpOpcUaConfigExtractor.java
@@ -47,6 +47,10 @@ import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabe
 import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.PASSWORD;
 import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.PULLING_INTERVAL;
 import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.PULL_MODE;
+import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.SUBSCRIPTION_DISCARD_OLDEST;
+import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.SUBSCRIPTION_PUBLISHING_INTERVAL;
+import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.SUBSCRIPTION_QUEUE_SIZE;
+import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.SUBSCRIPTION_SAMPLING_INTERVAL;
 import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.USERNAME;
 import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.USERNAME_GROUP;
 
@@ -83,6 +87,34 @@ public class SpOpcUaConfigExtractor {
 
       config.setPullIntervalMilliSeconds(pullIntervalSeconds);
       config.setIncompleteEventStrategy(incompleteEventStrategy);
+    } else {
+      config.setSubscriptionPublishingIntervalMs(extractOptionalInteger(
+          extractor,
+          SUBSCRIPTION_PUBLISHING_INTERVAL.name(),
+          OpcUaAdapterConfig.DEFAULT_SUBSCRIPTION_PUBLISHING_INTERVAL_MS
+      ));
+      config.setSubscriptionSamplingIntervalMs(extractOptionalInteger(
+          extractor,
+          SUBSCRIPTION_SAMPLING_INTERVAL.name(),
+          OpcUaAdapterConfig.DEFAULT_SUBSCRIPTION_SAMPLING_INTERVAL_MS
+      ));
+      config.setSubscriptionQueueSize(extractOptionalInteger(
+          extractor,
+          SUBSCRIPTION_QUEUE_SIZE.name(),
+          OpcUaAdapterConfig.DEFAULT_SUBSCRIPTION_QUEUE_SIZE
+      ));
+      var discardOldest = extractOptionalSelectedInternalName(
+          extractor,
+          SUBSCRIPTION_DISCARD_OLDEST.name(),
+          
String.valueOf(OpcUaAdapterConfig.DEFAULT_SUBSCRIPTION_DISCARD_OLDEST)
+      );
+      config.setSubscriptionDiscardOldest(Boolean.parseBoolean(discardOldest));
+      var incompleteEventStrategy = extractOptionalSelectedInternalName(
+          extractor,
+          SharedUserConfiguration.INCOMPLETE_EVENT_HANDLING_KEY,
+          SharedUserConfiguration.INCOMPLETE_OPTION_IGNORE
+      );
+      config.setIncompleteEventStrategy(incompleteEventStrategy);
     }
 
     var namingStrategySelection = extractor.selectedSingleValueInternalName(
@@ -94,6 +126,28 @@ public class SpOpcUaConfigExtractor {
     return config;
   }
 
+  private static Integer extractOptionalInteger(IStaticPropertyExtractor 
extractor,
+                                                String internalName,
+                                                Integer defaultValue) {
+    try {
+      var value = extractor.singleValueParameter(internalName, Integer.class);
+      return value != null ? value : defaultValue;
+    } catch (RuntimeException e) {
+      return defaultValue;
+    }
+  }
+
+  private static String 
extractOptionalSelectedInternalName(IStaticPropertyExtractor extractor,
+                                                           String internalName,
+                                                           String 
defaultValue) {
+    try {
+      var value = extractor.selectedSingleValueInternalName(internalName, 
String.class);
+      return value != null ? value : defaultValue;
+    } catch (RuntimeException e) {
+      return defaultValue;
+    }
+  }
+
   public static OpcUaConfig extractSinkConfig(IParameterExtractor extractor,
                                               IStreamPipesClient 
streamPipesClient) {
     return extractSharedConfig(extractor, new OpcUaConfig(), 
streamPipesClient);
diff --git 
a/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/config/security/SecurityConfig.java
 
b/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/config/security/SecurityConfig.java
index 2d4b433ba0..ccd64b78ea 100644
--- 
a/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/config/security/SecurityConfig.java
+++ 
b/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/config/security/SecurityConfig.java
@@ -48,19 +48,40 @@ public class SecurityConfig {
   private final MessageSecurityMode securityMode;
   private final SecurityPolicy securityPolicy;
   private final IStreamPipesClient streamPipesClient;
+  private final boolean disallowInsecureEndpoints;
 
   public SecurityConfig(MessageSecurityMode securityMode,
                         SecurityPolicy securityPolicy,
                         IStreamPipesClient streamPipesClient) {
+    this(
+        securityMode,
+        securityPolicy,
+        streamPipesClient,
+        
Environments.getEnvironment().getOpcUaDisallowInsecureEndpoints().getValueOrDefault()
+    );
+  }
+
+  SecurityConfig(MessageSecurityMode securityMode,
+                 SecurityPolicy securityPolicy,
+                 IStreamPipesClient streamPipesClient,
+                 boolean disallowInsecureEndpoints) {
     this.securityMode = securityMode;
     this.securityPolicy = securityPolicy;
     this.streamPipesClient = streamPipesClient;
+    this.disallowInsecureEndpoints = disallowInsecureEndpoints;
   }
 
   public void configureSecurityPolicy(OpcUaConfig config,
                                       List<EndpointDescription> endpoints,
                                       OpcUaClientConfigBuilder builder)
       throws SpConfigurationException, URISyntaxException {
+    if (disallowInsecureEndpoints && (securityMode == MessageSecurityMode.None 
|| securityPolicy == SecurityPolicy.None)) {
+      throw new SpConfigurationException(
+          "OPC UA connections with security mode None or security policy None 
are disabled by "
+              + "SP_OPCUA_DISALLOW_INSECURE_ENDPOINTS"
+      );
+    }
+
     URI configuredServerUri = new 
URI(config.getOpcServerURL()).parseServerAuthority();
 
     EndpointDescription tmpEndpoint = endpoints
diff --git 
a/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/migration/OpcUaAdapterMigrationV7.java
 
b/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/migration/OpcUaAdapterMigrationV7.java
new file mode 100644
index 0000000000..91f5acbe0f
--- /dev/null
+++ 
b/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/migration/OpcUaAdapterMigrationV7.java
@@ -0,0 +1,59 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ *
+ */
+
+package org.apache.streampipes.extensions.connectors.opcua.migration;
+
+import 
org.apache.streampipes.extensions.api.extractor.IStaticPropertyExtractor;
+import org.apache.streampipes.extensions.api.migration.IAdapterMigrator;
+import org.apache.streampipes.extensions.connectors.opcua.adapter.OpcUaAdapter;
+import 
org.apache.streampipes.extensions.connectors.opcua.config.SharedUserConfiguration;
+import org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels;
+import org.apache.streampipes.model.connect.adapter.AdapterDescription;
+import org.apache.streampipes.model.extensions.svcdiscovery.SpServiceTagPrefix;
+import org.apache.streampipes.model.migration.MigrationResult;
+import org.apache.streampipes.model.migration.ModelMigratorConfig;
+import org.apache.streampipes.model.staticproperty.StaticPropertyAlternatives;
+
+public class OpcUaAdapterMigrationV7 implements IAdapterMigrator {
+
+  @Override
+  public ModelMigratorConfig config() {
+    return new ModelMigratorConfig(
+        OpcUaAdapter.ID,
+        SpServiceTagPrefix.ADAPTER,
+        6,
+        7
+    );
+  }
+
+  @Override
+  public MigrationResult<AdapterDescription> migrate(AdapterDescription 
element,
+                                                     IStaticPropertyExtractor 
extractor) throws RuntimeException {
+    element.getConfig().forEach(sp -> {
+      if 
(sp.getInternalName().equalsIgnoreCase(OpcUaLabels.ADAPTER_TYPE.name())) {
+        var alternatives = (StaticPropertyAlternatives) sp;
+        alternatives.getAlternatives().stream()
+            .filter(alternative -> 
alternative.getInternalName().equals(OpcUaLabels.SUBSCRIPTION_MODE.name()))
+            .findFirst()
+            .ifPresent(alternative ->
+                
alternative.setStaticProperty(SharedUserConfiguration.getSubscriptionModeGroup()));
+      }
+    });
+    return MigrationResult.success(element);
+  }
+}
diff --git 
a/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/utils/OpcUaLabels.java
 
b/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/utils/OpcUaLabels.java
index 98ce898dff..725ea32d59 100644
--- 
a/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/utils/OpcUaLabels.java
+++ 
b/streampipes-extensions/streampipes-connectors-opcua/src/main/java/org/apache/streampipes/extensions/connectors/opcua/utils/OpcUaLabels.java
@@ -41,6 +41,10 @@ public enum OpcUaLabels {
   ADAPTER_TYPE,
   PULL_MODE,
   SUBSCRIPTION_MODE,
+  SUBSCRIPTION_PUBLISHING_INTERVAL,
+  SUBSCRIPTION_SAMPLING_INTERVAL,
+  SUBSCRIPTION_QUEUE_SIZE,
+  SUBSCRIPTION_DISCARD_OLDEST,
   NAMING_STRATEGY,
   MAPPING_PROPERY;
 }
diff --git 
a/streampipes-extensions/streampipes-connectors-opcua/src/main/resources/org.apache.streampipes.connect.iiot.adapters.opcua/documentation.md
 
b/streampipes-extensions/streampipes-connectors-opcua/src/main/resources/org.apache.streampipes.connect.iiot.adapters.opcua/documentation.md
index 24ec39a874..0a37586a06 100644
--- 
a/streampipes-extensions/streampipes-connectors-opcua/src/main/resources/org.apache.streampipes.connect.iiot.adapters.opcua/documentation.md
+++ 
b/streampipes-extensions/streampipes-connectors-opcua/src/main/resources/org.apache.streampipes.connect.iiot.adapters.opcua/documentation.md
@@ -39,6 +39,8 @@ The following environment variables control that location and 
certificate identi
 * SP_OPCUA_KEYSTORE_FILE the keystore file to create or reuse (e.g., 
keystore.pfx, must be of type PKCS12)
 * SP_OPCUA_KEYSTORE_PASSWORD the password to the keystore
 * SP_OPCUA_APPLICATION_URI the application URI used by the client to identify 
itself
+* SP_OPCUA_DISALLOW_INSECURE_ENDPOINTS set to `true` to reject connections 
that use security mode `None`
+  or security policy `None`
 
 Certificate requirements:
 
diff --git 
a/streampipes-extensions/streampipes-connectors-opcua/src/main/resources/org.apache.streampipes.connect.iiot.adapters.opcua/strings.en
 
b/streampipes-extensions/streampipes-connectors-opcua/src/main/resources/org.apache.streampipes.connect.iiot.adapters.opcua/strings.en
index 853b9f11ea..026dc4e889 100644
--- 
a/streampipes-extensions/streampipes-connectors-opcua/src/main/resources/org.apache.streampipes.connect.iiot.adapters.opcua/strings.en
+++ 
b/streampipes-extensions/streampipes-connectors-opcua/src/main/resources/org.apache.streampipes.connect.iiot.adapters.opcua/strings.en
@@ -80,6 +80,21 @@ PULL_MODE.description=
 SUBSCRIPTION_MODE.title=Subscription mode
 SUBSCRIPTION_MODE.description=
 
+subscription-mode-group.title=Subscription Settings
+subscription-mode-group.description=
+
+SUBSCRIPTION_PUBLISHING_INTERVAL.title=Publishing Interval
+SUBSCRIPTION_PUBLISHING_INTERVAL.description=Requested OPC UA subscription 
publishing interval, in milliseconds
+
+SUBSCRIPTION_SAMPLING_INTERVAL.title=Sampling Interval
+SUBSCRIPTION_SAMPLING_INTERVAL.description=Requested monitored item sampling 
interval, in milliseconds
+
+SUBSCRIPTION_QUEUE_SIZE.title=Queue Size
+SUBSCRIPTION_QUEUE_SIZE.description=Requested queue size for each monitored 
item
+
+SUBSCRIPTION_DISCARD_OLDEST.title=Queue Overflow Handling
+SUBSCRIPTION_DISCARD_OLDEST.description=Select which queued value is discarded 
when the monitored item queue is full
+
 incomplete-event-handling.title=Incomplete Events
 incomplete-event-handling.description=Select how events with missing values 
(e.g., due to bad status codes) are handled.
 
diff --git 
a/streampipes-extensions/streampipes-connectors-opcua/src/main/resources/org.apache.streampipes.sinks.databases.jvm.opcua/documentation.md
 
b/streampipes-extensions/streampipes-connectors-opcua/src/main/resources/org.apache.streampipes.sinks.databases.jvm.opcua/documentation.md
index 7032e2115a..7db32fea80 100644
--- 
a/streampipes-extensions/streampipes-connectors-opcua/src/main/resources/org.apache.streampipes.sinks.databases.jvm.opcua/documentation.md
+++ 
b/streampipes-extensions/streampipes-connectors-opcua/src/main/resources/org.apache.streampipes.sinks.databases.jvm.opcua/documentation.md
@@ -39,6 +39,8 @@ The following environment variables control that location and 
certificate identi
 * SP_OPCUA_KEYSTORE_FILE the keystore file to create or reuse (e.g., 
keystore.pfx, must be of type PKCS12)
 * SP_OPCUA_KEYSTORE_PASSWORD the password to the keystore
 * SP_OPCUA_APPLICATION_URI the application URI used by the client to identify 
itself
+* SP_OPCUA_DISALLOW_INSECURE_ENDPOINTS set to `true` to reject connections 
that use security mode `None`
+  or security policy `None`
 
 Certificate requirements:
 
diff --git 
a/streampipes-extensions/streampipes-connectors-opcua/src/test/java/org/apache/streampipes/extensions/connectors/opcua/adapter/OpcUaAdapterSubscriptionTest.java
 
b/streampipes-extensions/streampipes-connectors-opcua/src/test/java/org/apache/streampipes/extensions/connectors/opcua/adapter/OpcUaAdapterSubscriptionTest.java
new file mode 100644
index 0000000000..449c1f4d5e
--- /dev/null
+++ 
b/streampipes-extensions/streampipes-connectors-opcua/src/test/java/org/apache/streampipes/extensions/connectors/opcua/adapter/OpcUaAdapterSubscriptionTest.java
@@ -0,0 +1,204 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ *
+ */
+
+package org.apache.streampipes.extensions.connectors.opcua.adapter;
+
+import org.apache.streampipes.extensions.api.connect.IEventCollector;
+import 
org.apache.streampipes.extensions.connectors.opcua.client.ConnectedOpcUaClient;
+import 
org.apache.streampipes.extensions.connectors.opcua.client.OpcUaClientProvider;
+import 
org.apache.streampipes.extensions.connectors.opcua.config.OpcUaAdapterConfig;
+import 
org.apache.streampipes.extensions.connectors.opcua.config.SharedUserConfiguration;
+import 
org.apache.streampipes.extensions.connectors.opcua.model.node.BasicVariableNodeInfo;
+import org.apache.streampipes.extensions.connectors.opcua.model.node.OpcUaNode;
+import org.apache.streampipes.model.connect.guess.FieldStatusInfo;
+
+import org.eclipse.milo.opcua.sdk.client.OpcUaClient;
+import org.eclipse.milo.opcua.sdk.client.subscriptions.OpcUaMonitoredItem;
+import org.eclipse.milo.opcua.stack.core.StatusCodes;
+import org.eclipse.milo.opcua.stack.core.types.builtin.DataValue;
+import org.eclipse.milo.opcua.stack.core.types.builtin.NodeId;
+import org.eclipse.milo.opcua.stack.core.types.builtin.StatusCode;
+import org.eclipse.milo.opcua.stack.core.types.builtin.Variant;
+import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
+
+import java.lang.reflect.Field;
+import java.util.Map;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+class OpcUaAdapterSubscriptionTest {
+
+  @Test
+  @SuppressWarnings("unchecked")
+  void shouldRouteSubscriptionValuesByNodeIdForDuplicateDisplayNames() throws 
Exception {
+    var firstNodeId = new NodeId(2, "LineA.Temperature");
+    var secondNodeId = new NodeId(2, "LineB.Temperature");
+    var firstNode = new TestNode(firstNodeId, "Temperature", "lineA");
+    var secondNode = new TestNode(secondNodeId, "Temperature", "lineB");
+    var collector = mock(IEventCollector.class);
+    var adapter = makeAdapter(collector, firstNode, secondNode);
+
+    adapter.onSubscriptionValue(
+        OpcUaMonitoredItem.newDataItem(firstNodeId),
+        DataValue.valueOnly(new Variant(101))
+    );
+    adapter.onSubscriptionValue(
+        OpcUaMonitoredItem.newDataItem(secondNodeId),
+        DataValue.valueOnly(new Variant(202))
+    );
+
+    ArgumentCaptor<Map<String, Object>> eventCaptor = 
ArgumentCaptor.forClass(Map.class);
+    verify(collector).collect(eventCaptor.capture());
+    assertEquals(101, eventCaptor.getValue().get("lineA"));
+    assertEquals(202, eventCaptor.getValue().get("lineB"));
+  }
+
+  @Test
+  void 
shouldNotPublishBadStatusSubscriptionValuesWhenIncompleteEventsAreIgnored() 
throws Exception {
+    var nodeId = new NodeId(2, "IntermittentValue");
+    var node = new TestNode(nodeId, "IntermittentValue", "intermittent");
+    var collector = mock(IEventCollector.class);
+    var adapter = makeAdapter(collector, node);
+
+    adapter.onSubscriptionValue(
+        OpcUaMonitoredItem.newDataItem(nodeId),
+        new DataValue(new Variant(1), new 
StatusCode(StatusCodes.Bad_NotReadable))
+    );
+
+    verify(collector, never()).collect(org.mockito.ArgumentMatchers.any());
+  }
+
+  @Test
+  @SuppressWarnings("unchecked")
+  void shouldRemoveStaleSubscriptionValueBeforeSendingIncompleteEvent() throws 
Exception {
+    var firstNodeId = new NodeId(2, "StableCounter");
+    var secondNodeId = new NodeId(2, "IntermittentValue");
+    var firstNode = new TestNode(firstNodeId, "StableCounter", "stable");
+    var secondNode = new TestNode(secondNodeId, "IntermittentValue", 
"intermittent");
+    var collector = mock(IEventCollector.class);
+    var adapter = makeAdapter(collector, 
SharedUserConfiguration.INCOMPLETE_OPTION_SEND, firstNode, secondNode);
+
+    adapter.onSubscriptionValue(
+        OpcUaMonitoredItem.newDataItem(firstNodeId),
+        DataValue.valueOnly(new Variant(1))
+    );
+    adapter.onSubscriptionValue(
+        OpcUaMonitoredItem.newDataItem(secondNodeId),
+        DataValue.valueOnly(new Variant(2))
+    );
+    adapter.onSubscriptionValue(
+        OpcUaMonitoredItem.newDataItem(secondNodeId),
+        new DataValue(new Variant(3), new 
StatusCode(StatusCodes.Bad_NotReadable))
+    );
+
+    ArgumentCaptor<Map<String, Object>> eventCaptor = 
ArgumentCaptor.forClass(Map.class);
+    verify(collector, 
org.mockito.Mockito.times(2)).collect(eventCaptor.capture());
+
+    var incompleteEvent = eventCaptor.getAllValues().get(1);
+    assertEquals(1, incompleteEvent.get("stable"));
+    assertFalse(incompleteEvent.containsKey("intermittent"));
+  }
+
+  private OpcUaAdapter makeAdapter(IEventCollector collector,
+                                   TestNode... nodes) throws Exception {
+    return makeAdapter(collector, 
SharedUserConfiguration.INCOMPLETE_OPTION_IGNORE, nodes);
+  }
+
+  private OpcUaAdapter makeAdapter(IEventCollector collector,
+                                   String incompleteEventStrategy,
+                                   TestNode... nodes) throws Exception {
+    var adapter = new OpcUaAdapter(mock(OpcUaClientProvider.class));
+    var nodeIdToNodeMapping = getNodeIdToNodeMapping(adapter);
+    for (var node : nodes) {
+      nodeIdToNodeMapping.put(node.nodeInfo().getNodeId().toString(), node);
+    }
+
+    var config = new OpcUaAdapterConfig();
+    config.setIncompleteEventStrategy(incompleteEventStrategy);
+
+    setField(adapter, "collector", collector);
+    setField(adapter, "connectedClient", new 
ConnectedOpcUaClient(mock(OpcUaClient.class)));
+    setField(adapter, "numberOfEventProperties", nodes.length);
+    setField(adapter, "opcUaAdapterConfig", config);
+    return adapter;
+  }
+
+  @SuppressWarnings("unchecked")
+  private Map<String, OpcUaNode> getNodeIdToNodeMapping(OpcUaAdapter adapter) 
throws Exception {
+    Field field = adapter.getClass().getDeclaredField("nodeIdToNodeMapping");
+    field.setAccessible(true);
+    return (Map<String, OpcUaNode>) field.get(adapter);
+  }
+
+  private void setField(Object target,
+                        String fieldName,
+                        Object value) throws Exception {
+    Field field = target.getClass().getDeclaredField(fieldName);
+    field.setAccessible(true);
+    field.set(target, value);
+  }
+
+  private static class TestNode implements OpcUaNode {
+
+    private final BasicVariableNodeInfo nodeInfo;
+    private final String eventPropertyName;
+
+    TestNode(NodeId nodeId,
+             String displayName,
+             String eventPropertyName) {
+      this.nodeInfo = mock(BasicVariableNodeInfo.class);
+      this.eventPropertyName = eventPropertyName;
+
+      when(this.nodeInfo.getNodeId()).thenReturn(nodeId);
+      when(this.nodeInfo.getDisplayName()).thenReturn(displayName);
+      when(this.nodeInfo.getDesiredName("")).thenReturn(eventPropertyName);
+    }
+
+    @Override
+    public BasicVariableNodeInfo nodeInfo() {
+      return nodeInfo;
+    }
+
+    @Override
+    public int getNumberOfEventProperties(OpcUaClient client) {
+      return 1;
+    }
+
+    @Override
+    public void addToEvent(OpcUaClient client,
+                           Map<String, Object> event,
+                           Variant variant) {
+      event.put(eventPropertyName, variant.getValue());
+    }
+
+    @Override
+    public void addToEventPreview(OpcUaClient client,
+                                  Map<String, Object> eventPreview,
+                                  Map<String, FieldStatusInfo> 
fieldStatusInfos,
+                                  Variant variant,
+                                  FieldStatusInfo fieldStatusInfo) {
+      eventPreview.put(eventPropertyName, variant.getValue());
+    }
+  }
+}
diff --git 
a/streampipes-extensions/streampipes-connectors-opcua/src/test/java/org/apache/streampipes/extensions/connectors/opcua/config/SpOpcUaConfigExtractorTest.java
 
b/streampipes-extensions/streampipes-connectors-opcua/src/test/java/org/apache/streampipes/extensions/connectors/opcua/config/SpOpcUaConfigExtractorTest.java
index 33b80be2cf..ee49eabbc5 100644
--- 
a/streampipes-extensions/streampipes-connectors-opcua/src/test/java/org/apache/streampipes/extensions/connectors/opcua/config/SpOpcUaConfigExtractorTest.java
+++ 
b/streampipes-extensions/streampipes-connectors-opcua/src/test/java/org/apache/streampipes/extensions/connectors/opcua/config/SpOpcUaConfigExtractorTest.java
@@ -36,7 +36,13 @@ import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabe
 import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.OPC_URL;
 import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.PULLING_INTERVAL;
 import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.PULL_MODE;
+import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.SUBSCRIPTION_DISCARD_OLDEST;
+import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.SUBSCRIPTION_MODE;
+import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.SUBSCRIPTION_PUBLISHING_INTERVAL;
+import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.SUBSCRIPTION_QUEUE_SIZE;
+import static 
org.apache.streampipes.extensions.connectors.opcua.utils.OpcUaLabels.SUBSCRIPTION_SAMPLING_INTERVAL;
 import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.when;
 
@@ -73,4 +79,46 @@ class SpOpcUaConfigExtractorTest {
 
     assertEquals(List.of("ns=2;s=Demo.DataTypeTest.ExtensionObject"), 
config.getSelectedNodeNames());
   }
+
+  @Test
+  void shouldExtractSubscriptionSettings() {
+    IStaticPropertyExtractor extractor = mock(IStaticPropertyExtractor.class);
+
+    when(extractor.selectedAlternativeInternalId(ADAPTER_TYPE.name()))
+        .thenReturn(SUBSCRIPTION_MODE.name());
+    when(extractor.selectedAlternativeInternalId(OPC_HOST_OR_URL.name()))
+        .thenReturn(OPC_URL.name());
+    
when(extractor.selectedAlternativeInternalId(SharedUserConfiguration.USER_AUTHENTICATION))
+        .thenReturn(SharedUserConfiguration.USER_AUTHENTICATION_ANONYMOUS);
+    when(extractor.selectedTreeNodesInternalNames(AVAILABLE_NODES.name(), 
String.class))
+        .thenReturn(List.of("ns=2;s=Demo.Subscription.HighFrequencyCounter"));
+    
when(extractor.selectedSingleValueInternalName(SharedUserConfiguration.SECURITY_MODE,
 String.class))
+        .thenReturn(MessageSecurityMode.None.name());
+    
when(extractor.selectedSingleValue(SharedUserConfiguration.SECURITY_POLICY, 
String.class))
+        .thenReturn(SecurityPolicy.None.name());
+    when(extractor.singleValueParameter(OPC_SERVER_URL.name(), String.class))
+        .thenReturn("opc.tcp://localhost:4840/milo");
+    
when(extractor.singleValueParameter(SUBSCRIPTION_PUBLISHING_INTERVAL.name(), 
Integer.class))
+        .thenReturn(250);
+    when(extractor.singleValueParameter(SUBSCRIPTION_SAMPLING_INTERVAL.name(), 
Integer.class))
+        .thenReturn(100);
+    when(extractor.singleValueParameter(SUBSCRIPTION_QUEUE_SIZE.name(), 
Integer.class))
+        .thenReturn(25);
+    
when(extractor.selectedSingleValueInternalName(SUBSCRIPTION_DISCARD_OLDEST.name(),
 String.class))
+        .thenReturn(SharedUserConfiguration.SUBSCRIPTION_DISCARD_OLDEST_FALSE);
+    when(extractor.selectedSingleValueInternalName(
+        SharedUserConfiguration.INCOMPLETE_EVENT_HANDLING_KEY,
+        String.class
+    )).thenReturn(SharedUserConfiguration.INCOMPLETE_OPTION_SEND);
+    
when(extractor.selectedSingleValueInternalName(OpcUaLabels.NAMING_STRATEGY.name(),
 String.class))
+        .thenReturn(OpcUaNamingStrategy.DISPLAY_NAME.name());
+
+    var config = SpOpcUaConfigExtractor.extractAdapterConfig(extractor, 
mock(IStreamPipesClient.class));
+
+    assertEquals(250, config.getSubscriptionPublishingIntervalMs());
+    assertEquals(100, config.getSubscriptionSamplingIntervalMs());
+    assertEquals(25, config.getSubscriptionQueueSize());
+    assertFalse(config.isSubscriptionDiscardOldest());
+    assertEquals(SharedUserConfiguration.INCOMPLETE_OPTION_SEND, 
config.getIncompleteEventStrategy());
+  }
 }
diff --git 
a/streampipes-extensions/streampipes-connectors-opcua/src/test/java/org/apache/streampipes/extensions/connectors/opcua/config/security/SecurityConfigTest.java
 
b/streampipes-extensions/streampipes-connectors-opcua/src/test/java/org/apache/streampipes/extensions/connectors/opcua/config/security/SecurityConfigTest.java
index 020f368d00..ebfa2d0fda 100644
--- 
a/streampipes-extensions/streampipes-connectors-opcua/src/test/java/org/apache/streampipes/extensions/connectors/opcua/config/security/SecurityConfigTest.java
+++ 
b/streampipes-extensions/streampipes-connectors-opcua/src/test/java/org/apache/streampipes/extensions/connectors/opcua/config/security/SecurityConfigTest.java
@@ -18,6 +18,10 @@
 
 package org.apache.streampipes.extensions.connectors.opcua.config.security;
 
+import org.apache.streampipes.commons.exceptions.SpConfigurationException;
+import org.apache.streampipes.extensions.connectors.opcua.config.OpcUaConfig;
+
+import org.eclipse.milo.opcua.sdk.client.OpcUaClientConfigBuilder;
 import org.eclipse.milo.opcua.stack.core.security.SecurityPolicy;
 import org.eclipse.milo.opcua.stack.core.types.builtin.ByteString;
 import org.eclipse.milo.opcua.stack.core.types.enumerated.ApplicationType;
@@ -28,9 +32,13 @@ import 
org.eclipse.milo.opcua.stack.core.types.structured.UserTokenPolicy;
 import org.junit.jupiter.api.Test;
 
 import java.net.URI;
+import java.util.List;
 
 import static 
org.eclipse.milo.opcua.stack.core.types.builtin.unsigned.Unsigned.ubyte;
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
 import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
 
 class SecurityConfigTest {
 
@@ -39,7 +47,8 @@ class SecurityConfigTest {
     var securityConfig = new SecurityConfig(
         MessageSecurityMode.None,
         SecurityPolicy.None,
-        null
+        null,
+        false
     );
 
     var endpoint = new EndpointDescription(
@@ -68,4 +77,91 @@ class SecurityConfigTest {
 
     assertEquals("opc.tcp://127.0.0.1:32791/milo", updated.getEndpointUrl());
   }
+
+  @Test
+  void configureSecurityPolicyRejectsNoneSecurityModeWhenDisallowed() {
+    var securityConfig = new SecurityConfig(
+        MessageSecurityMode.None,
+        SecurityPolicy.Basic256Sha256,
+        null,
+        true
+    );
+
+    var exception = assertThrows(
+        SpConfigurationException.class,
+        () -> securityConfig.configureSecurityPolicy(
+            makeConfig(),
+            List.of(makeNoneEndpoint()),
+            new OpcUaClientConfigBuilder()
+        )
+    );
+
+    
assertTrue(exception.getMessage().contains("SP_OPCUA_DISALLOW_INSECURE_ENDPOINTS"));
+  }
+
+  @Test
+  void configureSecurityPolicyRejectsNoneSecurityPolicyWhenDisallowed() {
+    for (var securityMode : List.of(MessageSecurityMode.Sign, 
MessageSecurityMode.SignAndEncrypt)) {
+      var securityConfig = new SecurityConfig(
+          securityMode,
+          SecurityPolicy.None,
+          null,
+          true
+      );
+
+      var exception = assertThrows(
+          SpConfigurationException.class,
+          () -> securityConfig.configureSecurityPolicy(
+              makeConfig(),
+              List.of(makeNoneEndpoint()),
+              new OpcUaClientConfigBuilder()
+          )
+      );
+
+      
assertTrue(exception.getMessage().contains("SP_OPCUA_DISALLOW_INSECURE_ENDPOINTS"));
+    }
+  }
+
+  @Test
+  void configureSecurityPolicyAllowsNoneNoneByDefault() {
+    var securityConfig = new SecurityConfig(
+        MessageSecurityMode.None,
+        SecurityPolicy.None,
+        null,
+        false
+    );
+
+    assertDoesNotThrow(() -> securityConfig.configureSecurityPolicy(
+        makeConfig(),
+        List.of(makeNoneEndpoint()),
+        new OpcUaClientConfigBuilder()
+    ));
+  }
+
+  private OpcUaConfig makeConfig() {
+    var config = new OpcUaConfig();
+    config.setOpcServerURL("opc.tcp://127.0.0.1:4840/milo");
+    return config;
+  }
+
+  private EndpointDescription makeNoneEndpoint() {
+    return new EndpointDescription(
+        "opc.tcp://localhost:4840/milo",
+        new ApplicationDescription(
+            "urn:test",
+            "urn:test:product",
+            null,
+            ApplicationType.Server,
+            null,
+            null,
+            null
+        ),
+        ByteString.NULL_VALUE,
+        MessageSecurityMode.None,
+        SecurityPolicy.None.getUri(),
+        new UserTokenPolicy[0],
+        null,
+        ubyte(0)
+    );
+  }
 }

Reply via email to