This is an automated email from the ASF dual-hosted git repository. dominikriemer pushed a commit to branch improve-opc-ua-subscription-mode-configuration in repository https://gitbox.apache.org/repos/asf/streampipes.git
commit 370d44ce56892e41f34ff9628c705970ebf6df71 Author: Dominik Riemer <[email protected]> AuthorDate: Mon Jun 29 22:32:25 2026 +0200 Add environment variable to disallow insecure OPC UA endpoints --- .../apache/streampipes/commons/constants/Envs.java | 1 + .../commons/environment/DefaultEnvironment.java | 5 ++ .../commons/environment/Environment.java | 2 + .../connectors/opcua/adapter/OpcUaNodeBrowser.java | 2 +- .../opcua/config/security/SecurityConfig.java | 21 +++++ .../documentation.md | 2 + .../documentation.md | 2 + .../opcua/config/security/SecurityConfigTest.java | 98 +++++++++++++++++++++- 8 files changed, 131 insertions(+), 2 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/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/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/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.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/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) + ); + } }
