From 241f7b0e330df68c43dd08d7209298448a2b13de Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Tue, 21 Jul 2026 10:43:31 +0800 Subject: [PATCH 1/2] [Pipe] Prevent unexpected OPC UA endpoint redirects --- .../iotdb/db/i18n/DataNodePipeMessages.java | 3 + .../iotdb/db/i18n/DataNodePipeMessages.java | 3 + .../pipe/sink/protocol/opcua/OpcUaSink.java | 22 ++- .../protocol/opcua/client/ClientRunner.java | 89 ++++++++++- .../opcua/client/IoTDBOpcUaClient.java | 12 +- .../opcua/client/ClientRunnerTest.java | 147 ++++++++++++++++++ .../config/constant/PipeSinkConstant.java | 5 + 7 files changed, 269 insertions(+), 12 deletions(-) create mode 100644 iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/client/ClientRunnerTest.java diff --git a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java index 40b499410aaec..eef4e4d71813b 100644 --- a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java +++ b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java @@ -945,6 +945,9 @@ public final class DataNodePipeMessages { "Network failed to receive tsFile %s, status: %s"; public static final String SECURITY_DIR = "security dir: {}"; public static final String SECURITY_PKI_DIR = "security pki dir: {}"; + public static final String + LOG_OPC_UA_ENDPOINT_SELECTED_CONFIGURED_ARG_ADVERTISED_ARG_EFFECTIVE_ARG_ALLOWENDPOINTREDIRECT_ARG_4FE076CB = + "OPC UA endpoint selected: configured={}, advertised={}, effective={}, allowEndpointRedirect={}."; public static final String SSL_TRUST_STORE_PAIR_REQUIRED_WHEN_SSL_ENABLED = "When %s or %s is true, specify a complete trust-store pair under the same " + "alias: %s and %s, %s and %s, or %s and %s"; diff --git a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java index 85a8ae289cf4f..21512774cf290 100644 --- a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java +++ b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java @@ -892,6 +892,9 @@ public final class DataNodePipeMessages { "网络接收 TsFile %s 失败,状态:%s"; public static final String SECURITY_DIR = "security 目录:{}"; public static final String SECURITY_PKI_DIR = "security pki 目录:{}"; + public static final String + LOG_OPC_UA_ENDPOINT_SELECTED_CONFIGURED_ARG_ADVERTISED_ARG_EFFECTIVE_ARG_ALLOWENDPOINTREDIRECT_ARG_4FE076CB = + "已选择 OPC UA endpoint:configured={},advertised={},effective={},allowEndpointRedirect={}。"; public static final String SSL_TRUST_STORE_PAIR_REQUIRED_WHEN_SSL_ENABLED = "当 %s 或 %s 为 true 时,请在同一别名下指定完整的 trust-store 参数对:%s 和 %s、%s 和 %s,或 %s 和 %s"; public static final String SSL_KEY_STORE_PATH_AND_PASSWORD_MUST_BE_SPECIFIED_TOGETHER = diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/OpcUaSink.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/OpcUaSink.java index ea32d6263aaed..a264beed09aea 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/OpcUaSink.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/OpcUaSink.java @@ -70,6 +70,8 @@ import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CONNECTOR_IOTDB_USERNAME_KEY; import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CONNECTOR_IOTDB_USER_DEFAULT_VALUE; import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CONNECTOR_IOTDB_USER_KEY; +import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CONNECTOR_OPC_UA_ALLOW_ENDPOINT_REDIRECT_DEFAULT_VALUE; +import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CONNECTOR_OPC_UA_ALLOW_ENDPOINT_REDIRECT_KEY; import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CONNECTOR_OPC_UA_DEBOUNCE_TIME_MS_DEFAULT_VALUE; import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CONNECTOR_OPC_UA_DEBOUNCE_TIME_MS_KEY; import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CONNECTOR_OPC_UA_DEFAULT_QUALITY_BAD_VALUE; @@ -112,6 +114,7 @@ import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.SINK_IOTDB_PASSWORD_KEY; import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.SINK_IOTDB_USERNAME_KEY; import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.SINK_IOTDB_USER_KEY; +import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.SINK_OPC_UA_ALLOW_ENDPOINT_REDIRECT_KEY; import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.SINK_OPC_UA_DEBOUNCE_TIME_MS_KEY; import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.SINK_OPC_UA_DEFAULT_QUALITY_KEY; import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.SINK_OPC_UA_ENABLE_ANONYMOUS_ACCESS_KEY; @@ -399,6 +402,12 @@ private void customizeClient(final PipeParameters parameters) { parameters.getLongOrDefault( Arrays.asList(CONNECTOR_OPC_UA_TIMEOUT_SECONDS_KEY, SINK_OPC_UA_TIMEOUT_SECONDS_KEY), CONNECTOR_OPC_UA_TIMEOUT_SECONDS_DEFAULT_VALUE); + final boolean allowEndpointRedirect = + parameters.getBooleanOrDefault( + Arrays.asList( + CONNECTOR_OPC_UA_ALLOW_ENDPOINT_REDIRECT_KEY, + SINK_OPC_UA_ALLOW_ENDPOINT_REDIRECT_KEY), + CONNECTOR_OPC_UA_ALLOW_ENDPOINT_REDIRECT_DEFAULT_VALUE); synchronized (CLIENT_KEY_TO_REFERENCE_COUNT_AND_CLIENT_MAP) { client = @@ -418,11 +427,20 @@ private void customizeClient(final PipeParameters parameters) { SINK_OPC_UA_HISTORIZING_KEY), CONNECTOR_OPC_UA_HISTORIZING_DEFAULT_VALUE)); final ClientRunner runner = - new ClientRunner(result, securityDir, password, userName, timeoutSeconds); + new ClientRunner( + result, + securityDir, + password, + userName, + timeoutSeconds, + allowEndpointRedirect); runner.run(); return new Pair<>(new AtomicInteger(0), result); } - oldValue.getRight().checkEquals(userName, password, securityDir, policy); + oldValue + .getRight() + .checkEquals( + userName, password, securityDir, policy, allowEndpointRedirect); return oldValue; }) .getRight(); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/client/ClientRunner.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/client/ClientRunner.java index dc68e0d2b969b..6cbc5c3d61df7 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/client/ClientRunner.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/client/ClientRunner.java @@ -27,7 +27,10 @@ import org.eclipse.milo.opcua.stack.client.security.DefaultClientCertificateValidator; import org.eclipse.milo.opcua.stack.core.security.DefaultTrustListManager; import org.eclipse.milo.opcua.stack.core.security.SecurityPolicy; +import org.eclipse.milo.opcua.stack.core.transport.TransportProfile; import org.eclipse.milo.opcua.stack.core.types.builtin.LocalizedText; +import org.eclipse.milo.opcua.stack.core.types.structured.EndpointDescription; +import org.eclipse.milo.opcua.stack.core.util.EndpointUtil; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -37,7 +40,10 @@ import java.nio.file.Path; import java.nio.file.Paths; import java.security.Security; +import java.util.List; +import java.util.Locale; import java.util.Objects; +import java.util.Optional; import static org.eclipse.milo.opcua.stack.core.types.builtin.unsigned.Unsigned.uint; @@ -54,6 +60,7 @@ public class ClientRunner { private final Path securityDir; private final String password; private final long timeoutSeconds; + private final boolean allowEndpointRedirect; // For conflict checking private final String user; @@ -64,11 +71,22 @@ public ClientRunner( final String password, final String user, final long timeoutSeconds) { + this(configurableUaClient, securityDir, password, user, timeoutSeconds, false); + } + + public ClientRunner( + final IoTDBOpcUaClient configurableUaClient, + final String securityDir, + final String password, + final String user, + final long timeoutSeconds, + final boolean allowEndpointRedirect) { this.configurableUaClient = configurableUaClient; this.securityDir = Paths.get(securityDir); this.password = password; this.user = user; this.timeoutSeconds = timeoutSeconds; + this.allowEndpointRedirect = allowEndpointRedirect; configurableUaClient.setRunner(this); } @@ -93,7 +111,12 @@ private OpcUaClient createClient() throws Exception { return OpcUaClient.create( configurableUaClient.getNodeUrl(), - endpoints -> endpoints.stream().filter(configurableUaClient.endpointFilter()).findFirst(), + endpoints -> + selectEndpoint( + endpoints, + configurableUaClient.getNodeUrl(), + configurableUaClient.getSecurityPolicy(), + allowEndpointRedirect), configBuilder -> configBuilder .setApplicationName(LocalizedText.english("Apache IoTDB OPC UA client")) @@ -109,6 +132,66 @@ private OpcUaClient createClient() throws Exception { .build()); } + static Optional selectEndpoint( + final List endpoints, + final String configuredNodeUrl, + final SecurityPolicy securityPolicy, + final boolean allowEndpointRedirect) { + final String configuredScheme = normalizeScheme(EndpointUtil.getScheme(configuredNodeUrl)); + + return endpoints.stream() + .filter(endpoint -> securityPolicy.getUri().equals(endpoint.getSecurityPolicyUri())) + .filter(endpoint -> matchesConfiguredTransport(endpoint, configuredScheme)) + .findFirst() + .map( + advertisedEndpoint -> { + final EndpointDescription effectiveEndpoint = + allowEndpointRedirect + ? advertisedEndpoint + : EndpointUtil.updateUrl( + advertisedEndpoint, + EndpointUtil.getHost(configuredNodeUrl), + EndpointUtil.getPort(configuredNodeUrl)); + logger.info( + DataNodePipeMessages + .LOG_OPC_UA_ENDPOINT_SELECTED_CONFIGURED_ARG_ADVERTISED_ARG_EFFECTIVE_ARG_ALLOWENDPOINTREDIRECT_ARG_4FE076CB, + configuredNodeUrl, + advertisedEndpoint.getEndpointUrl(), + effectiveEndpoint.getEndpointUrl(), + allowEndpointRedirect); + return effectiveEndpoint; + }); + } + + private static boolean matchesConfiguredTransport( + final EndpointDescription endpoint, final String configuredScheme) { + if (Objects.isNull(configuredScheme) + || !Objects.equals( + configuredScheme, normalizeScheme(EndpointUtil.getScheme(endpoint.getEndpointUrl())))) { + return false; + } + + final String transportProfileUri = endpoint.getTransportProfileUri(); + if (Objects.isNull(transportProfileUri)) { + return true; + } + try { + return Objects.equals( + configuredScheme, + normalizeScheme(TransportProfile.fromUri(transportProfileUri).getScheme())); + } catch (final IllegalArgumentException ignored) { + // Preserve compatibility with servers that advertise a custom transport profile. The + // endpoint URL scheme still has to match the configured URL. + return true; + } + } + + private static String normalizeScheme(final String scheme) { + return Objects.nonNull(scheme) && scheme.equalsIgnoreCase("opc.https") + ? "https" + : Objects.isNull(scheme) ? null : scheme.toLowerCase(Locale.ROOT); + } + public void run() { try { final OpcUaClient client = createClient(); @@ -143,7 +226,8 @@ void checkEquals( final String user, final String password, final Path securityDir, - final SecurityPolicy securityPolicy) { + final SecurityPolicy securityPolicy, + final boolean allowEndpointRedirect) { checkEquals("user", this.user, user); checkEquals("password", this.password, password); checkEquals( @@ -151,6 +235,7 @@ void checkEquals( FileSystems.getDefault().getPath(this.securityDir.toAbsolutePath().toString()), FileSystems.getDefault().getPath(securityDir.toAbsolutePath().toString())); checkEquals("securityPolicy", configurableUaClient.getSecurityPolicy(), securityPolicy); + checkEquals("allow endpoint redirect", this.allowEndpointRedirect, allowEndpointRedirect); } private void checkEquals(final String attrName, Object thisAttr, Object thatAttr) { diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/client/IoTDBOpcUaClient.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/client/IoTDBOpcUaClient.java index 62e489b2e1f1b..3dd2141fad52a 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/client/IoTDBOpcUaClient.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/client/IoTDBOpcUaClient.java @@ -52,7 +52,6 @@ import org.eclipse.milo.opcua.stack.core.types.structured.AddNodesItem; import org.eclipse.milo.opcua.stack.core.types.structured.AddNodesResponse; import org.eclipse.milo.opcua.stack.core.types.structured.AddNodesResult; -import org.eclipse.milo.opcua.stack.core.types.structured.EndpointDescription; import org.eclipse.milo.opcua.stack.core.types.structured.ObjectAttributes; import org.eclipse.milo.opcua.stack.core.types.structured.VariableAttributes; import org.slf4j.Logger; @@ -65,7 +64,6 @@ import java.util.List; import java.util.Objects; import java.util.concurrent.ExecutionException; -import java.util.function.Predicate; import static org.apache.iotdb.db.pipe.sink.protocol.opcua.server.OpcUaNameSpace.convertToOpcDataType; import static org.apache.iotdb.db.pipe.sink.protocol.opcua.server.OpcUaNameSpace.timestampToUtc; @@ -307,10 +305,6 @@ String getNodeUrl() { return nodeUrl; } - Predicate endpointFilter() { - return e -> getSecurityPolicy().getUri().equals(e.getSecurityPolicyUri()); - } - SecurityPolicy getSecurityPolicy() { return securityPolicy; } @@ -365,7 +359,9 @@ public void checkEquals( final String user, final String password, final String securityDir, - final SecurityPolicy securityPolicy) { - runner.checkEquals(user, password, Paths.get(securityDir), securityPolicy); + final SecurityPolicy securityPolicy, + final boolean allowEndpointRedirect) { + runner.checkEquals( + user, password, Paths.get(securityDir), securityPolicy, allowEndpointRedirect); } } diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/client/ClientRunnerTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/client/ClientRunnerTest.java new file mode 100644 index 0000000000000..c760f661a711e --- /dev/null +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/client/ClientRunnerTest.java @@ -0,0 +1,147 @@ +/* + * 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.iotdb.db.pipe.sink.protocol.opcua.client; + +import org.apache.iotdb.pipe.api.exception.PipeException; + +import org.eclipse.milo.opcua.sdk.client.api.identity.AnonymousProvider; +import org.eclipse.milo.opcua.stack.core.Stack; +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.builtin.unsigned.Unsigned; +import org.eclipse.milo.opcua.stack.core.types.enumerated.MessageSecurityMode; +import org.eclipse.milo.opcua.stack.core.types.structured.EndpointDescription; +import org.eclipse.milo.opcua.stack.core.types.structured.UserTokenPolicy; +import org.junit.Assert; +import org.junit.Test; + +import java.nio.file.Paths; +import java.util.Arrays; + +public class ClientRunnerTest { + + private static final String CONFIGURED_NODE_URL = "opc.tcp://10.60.80.65:12686/iotdb"; + private static final SecurityPolicy SECURITY_POLICY = SecurityPolicy.Basic256Sha256; + + @Test + public void testConfiguredHostAndPortAreUsedByDefault() { + final EndpointDescription advertisedEndpoint = + createEndpoint( + "opc.tcp://fwq03-15:4840/server-path", + SECURITY_POLICY, + Stack.TCP_UASC_UABINARY_TRANSPORT_URI); + + final EndpointDescription effectiveEndpoint = + ClientRunner.selectEndpoint( + Arrays.asList(advertisedEndpoint), CONFIGURED_NODE_URL, SECURITY_POLICY, false) + .orElseThrow(AssertionError::new); + + Assert.assertEquals( + "opc.tcp://10.60.80.65:12686/server-path", effectiveEndpoint.getEndpointUrl()); + Assert.assertNotSame(advertisedEndpoint, effectiveEndpoint); + Assert.assertSame( + advertisedEndpoint.getServerCertificate(), effectiveEndpoint.getServerCertificate()); + Assert.assertEquals(advertisedEndpoint.getSecurityMode(), effectiveEndpoint.getSecurityMode()); + Assert.assertEquals( + advertisedEndpoint.getSecurityPolicyUri(), effectiveEndpoint.getSecurityPolicyUri()); + Assert.assertSame( + advertisedEndpoint.getUserIdentityTokens(), effectiveEndpoint.getUserIdentityTokens()); + Assert.assertEquals( + advertisedEndpoint.getTransportProfileUri(), effectiveEndpoint.getTransportProfileUri()); + } + + @Test + public void testAdvertisedEndpointIsUsedWhenRedirectIsAllowed() { + final EndpointDescription advertisedEndpoint = + createEndpoint( + "opc.tcp://fwq03-15:4840/server-path", + SECURITY_POLICY, + Stack.TCP_UASC_UABINARY_TRANSPORT_URI); + + final EndpointDescription effectiveEndpoint = + ClientRunner.selectEndpoint( + Arrays.asList(advertisedEndpoint), CONFIGURED_NODE_URL, SECURITY_POLICY, true) + .orElseThrow(AssertionError::new); + + Assert.assertSame(advertisedEndpoint, effectiveEndpoint); + } + + @Test + public void testEndpointSelectionMatchesConfiguredTransport() { + final EndpointDescription wrongSchemeEndpoint = + createEndpoint( + "https://wrong-scheme:12686/iotdb", + SECURITY_POLICY, + Stack.HTTPS_UABINARY_TRANSPORT_URI); + final EndpointDescription wrongTransportEndpoint = + createEndpoint( + "opc.tcp://wrong-transport:12686/iotdb", + SECURITY_POLICY, + Stack.HTTPS_UABINARY_TRANSPORT_URI); + final EndpointDescription matchingEndpoint = + createEndpoint( + "opc.tcp://matching:12686/iotdb", + SECURITY_POLICY, + Stack.TCP_UASC_UABINARY_TRANSPORT_URI); + + final EndpointDescription selectedEndpoint = + ClientRunner.selectEndpoint( + Arrays.asList(wrongSchemeEndpoint, wrongTransportEndpoint, matchingEndpoint), + CONFIGURED_NODE_URL, + SECURITY_POLICY, + true) + .orElseThrow(AssertionError::new); + + Assert.assertSame(matchingEndpoint, selectedEndpoint); + } + + @Test + public void testAllowEndpointRedirectParticipatesInConflictDetection() { + final String securityDir = "target/opcua-client-runner-test"; + final IoTDBOpcUaClient client = + new IoTDBOpcUaClient( + CONFIGURED_NODE_URL, SECURITY_POLICY, AnonymousProvider.INSTANCE, false); + final ClientRunner runner = new ClientRunner(client, securityDir, "password", null, 10, false); + + runner.checkEquals(null, "password", Paths.get(securityDir), SECURITY_POLICY, false); + final PipeException exception = + Assert.assertThrows( + PipeException.class, + () -> + runner.checkEquals( + null, "password", Paths.get(securityDir), SECURITY_POLICY, true)); + Assert.assertTrue(exception.getMessage().contains("allow endpoint redirect")); + } + + private static EndpointDescription createEndpoint( + final String endpointUrl, + final SecurityPolicy securityPolicy, + final String transportProfileUri) { + return new EndpointDescription( + endpointUrl, + null, + ByteString.of(new byte[] {1, 2, 3}), + MessageSecurityMode.SignAndEncrypt, + securityPolicy.getUri(), + new UserTokenPolicy[0], + transportProfileUri, + Unsigned.ubyte(1)); + } +} diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/constant/PipeSinkConstant.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/constant/PipeSinkConstant.java index 25eb4c8bb21ce..a6dedfad08150 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/constant/PipeSinkConstant.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/constant/PipeSinkConstant.java @@ -245,6 +245,11 @@ private static String getDefaultConnectorOrSinkName(final PipeParameters paramet public static final String CONNECTOR_OPC_UA_NODE_URL_KEY = "connector.opcua.node-url"; public static final String SINK_OPC_UA_NODE_URL_KEY = "sink.opcua.node-url"; + public static final String CONNECTOR_OPC_UA_ALLOW_ENDPOINT_REDIRECT_KEY = + "connector.opcua.allow-endpoint-redirect"; + public static final String SINK_OPC_UA_ALLOW_ENDPOINT_REDIRECT_KEY = + "sink.opcua.allow-endpoint-redirect"; + public static final boolean CONNECTOR_OPC_UA_ALLOW_ENDPOINT_REDIRECT_DEFAULT_VALUE = false; public static final String CONNECTOR_OPC_UA_SECURITY_POLICY_KEY = "connector.opcua.security-policy"; From df7c05857cff338509f650a6f0a84d17dbf8d28d Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Tue, 21 Jul 2026 12:14:28 +0800 Subject: [PATCH 2/2] [Pipe] Add advertised host for internal OPC UA server --- .../iotdb/db/i18n/DataNodePipeMessages.java | 8 + .../iotdb/db/i18n/DataNodePipeMessages.java | 7 + .../pipe/sink/protocol/opcua/OpcUaSink.java | 7 + .../opcua/server/OpcUaKeyStoreLoader.java | 21 ++- .../protocol/opcua/server/OpcUaNameSpace.java | 2 + .../opcua/server/OpcUaServerBuilder.java | 143 +++++++++++++-- .../opcua/server/OpcUaServerBuilderTest.java | 164 ++++++++++++++++++ .../config/constant/PipeSinkConstant.java | 4 + 8 files changed, 334 insertions(+), 22 deletions(-) create mode 100644 iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaServerBuilderTest.java diff --git a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java index eef4e4d71813b..f13bca9a4a8a6 100644 --- a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java +++ b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java @@ -1538,6 +1538,14 @@ public final class DataNodePipeMessages { public static final String UNABLE_CREATE_SECURITY_DIR = "Unable to create security dir: "; public static final String OPC_UA_SECURITY_DIR = "Security dir: {}"; public static final String OPC_UA_SECURITY_PKI_DIR = "Security pki dir: {}"; + public static final String + EXCEPTION_THE_ADVERTISED_HOST_MUST_BE_A_HOSTNAME_OR_IP_ADDRESS_WITHOUT_A_SCHEME_PORT_OR_PATH_6857C67A = + "The advertised host must be a hostname or IP address without a scheme, port, or path."; + public static final String + LOG_ADVERTISED_HOST_ARG_IS_NOT_PRESENT_IN_THE_LOADED_OPC_UA_SERVER_CERTIFICATE_SUBJECT_ALTERNATIVE_NAMES_SECURED_CLIENTS_MAY_REJECT_IT_REPLACE_OR_REGENERATE_THE_CERTIFICATE_AND_ESTABLISH_TRUST_AGAIN_912358AF = + "Advertised host {} is not present in the loaded OPC UA server certificate subject " + + "alternative names. Secured clients may reject it; replace or regenerate the " + + "certificate and establish trust again."; // --------------------------------------------------------------------------- // pipe – PipeDataNodePluginAgent diff --git a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java index 21512774cf290..17dd4c6d836f7 100644 --- a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java +++ b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java @@ -1442,6 +1442,13 @@ public final class DataNodePipeMessages { "安全目录:{}"; public static final String OPC_UA_SECURITY_PKI_DIR = "安全 PKI 目录:{}"; + public static final String + EXCEPTION_THE_ADVERTISED_HOST_MUST_BE_A_HOSTNAME_OR_IP_ADDRESS_WITHOUT_A_SCHEME_PORT_OR_PATH_6857C67A = + "advertised host 必须是不带 scheme、port 或 path 的 hostname 或 IP 地址。"; + public static final String + LOG_ADVERTISED_HOST_ARG_IS_NOT_PRESENT_IN_THE_LOADED_OPC_UA_SERVER_CERTIFICATE_SUBJECT_ALTERNATIVE_NAMES_SECURED_CLIENTS_MAY_REJECT_IT_REPLACE_OR_REGENERATE_THE_CERTIFICATE_AND_ESTABLISH_TRUST_AGAIN_912358AF = + "advertised host {} 不在已加载的 OPC UA server 证书 subject alternative names 中。安全客户端可能拒绝该证书;" + + "请替换或重新生成证书并重新建立信任。"; // --------------------------------------------------------------------------- // pipe – PipeDataNodePluginAgent diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/OpcUaSink.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/OpcUaSink.java index a264beed09aea..2eb4e7e81361c 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/OpcUaSink.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/OpcUaSink.java @@ -70,6 +70,7 @@ import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CONNECTOR_IOTDB_USERNAME_KEY; import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CONNECTOR_IOTDB_USER_DEFAULT_VALUE; import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CONNECTOR_IOTDB_USER_KEY; +import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CONNECTOR_OPC_UA_ADVERTISED_HOST_KEY; import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CONNECTOR_OPC_UA_ALLOW_ENDPOINT_REDIRECT_DEFAULT_VALUE; import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CONNECTOR_OPC_UA_ALLOW_ENDPOINT_REDIRECT_KEY; import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.CONNECTOR_OPC_UA_DEBOUNCE_TIME_MS_DEFAULT_VALUE; @@ -114,6 +115,7 @@ import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.SINK_IOTDB_PASSWORD_KEY; import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.SINK_IOTDB_USERNAME_KEY; import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.SINK_IOTDB_USER_KEY; +import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.SINK_OPC_UA_ADVERTISED_HOST_KEY; import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.SINK_OPC_UA_ALLOW_ENDPOINT_REDIRECT_KEY; import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.SINK_OPC_UA_DEBOUNCE_TIME_MS_KEY; import static org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant.SINK_OPC_UA_DEFAULT_QUALITY_KEY; @@ -271,6 +273,9 @@ private void customizeServer(final PipeParameters parameters) { parameters.getIntOrDefault( Arrays.asList(CONNECTOR_OPC_UA_HTTPS_BIND_PORT_KEY, SINK_OPC_UA_HTTPS_BIND_PORT_KEY), CONNECTOR_OPC_UA_HTTPS_BIND_PORT_DEFAULT_VALUE); + final String advertisedHost = + parameters.getStringByKeys( + CONNECTOR_OPC_UA_ADVERTISED_HOST_KEY, SINK_OPC_UA_ADVERTISED_HOST_KEY); final String user = parameters.getStringOrDefault( @@ -333,6 +338,7 @@ private void customizeServer(final PipeParameters parameters) { new OpcUaServerBuilder() .setTcpBindPort(tcpBindPort) .setHttpsBindPort(httpsBindPort) + .setAdvertisedHost(advertisedHost) .setUser(user) .setPassword(password) .setSecurityDir(securityDir) @@ -348,6 +354,7 @@ private void customizeServer(final PipeParameters parameters) { oldValue .getRight() .checkEquals( + advertisedHost, user, password, securityDir, diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaKeyStoreLoader.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaKeyStoreLoader.java index 26fa208443091..2e27ffc22777c 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaKeyStoreLoader.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaKeyStoreLoader.java @@ -22,7 +22,7 @@ import org.apache.iotdb.commons.utils.FileUtils; import org.apache.iotdb.db.i18n.DataNodePipeMessages; -import com.google.common.collect.Sets; +import com.google.common.net.InetAddresses; import org.eclipse.milo.opcua.sdk.server.util.HostnameUtil; import org.eclipse.milo.opcua.stack.core.util.SelfSignedCertificateBuilder; import org.eclipse.milo.opcua.stack.core.util.SelfSignedCertificateGenerator; @@ -41,22 +41,21 @@ import java.security.PrivateKey; import java.security.PublicKey; import java.security.cert.X509Certificate; +import java.util.LinkedHashSet; import java.util.Set; import java.util.UUID; -import java.util.regex.Pattern; class OpcUaKeyStoreLoader { private static final Logger LOGGER = LoggerFactory.getLogger(OpcUaKeyStoreLoader.class); - private static final Pattern IP_ADDR_PATTERN = - Pattern.compile("^(([01]?\\d\\d?|2[0-4]\\d|25[0-5])\\.){3}([01]?\\d\\d?|2[0-4]\\d|25[0-5])$"); - private static final String SERVER_ALIAS = "server-ai"; private X509Certificate serverCertificate; private KeyPair serverKeyPair; - OpcUaKeyStoreLoader load(final Path baseDir, final char[] password) throws Exception { + OpcUaKeyStoreLoader load( + final Path baseDir, final char[] password, final Set advertisedHostnames) + throws Exception { final KeyStore keyStore = KeyStore.getInstance("PKCS12"); final File serverKeyStore = baseDir.resolve("iotdb-server.pfx").toFile(); @@ -90,14 +89,14 @@ OpcUaKeyStoreLoader load(final Path baseDir, final char[] password) throws Excep .setApplicationUri(applicationUri); // Get as many hostnames and IP addresses as we can list in the certificate. - final Set hostnames = - Sets.union( - Sets.newHashSet(HostnameUtil.getHostname()), - HostnameUtil.getHostnames("0.0.0.0", false)); + final Set hostnames = new LinkedHashSet<>(); + hostnames.add(HostnameUtil.getHostname()); + hostnames.addAll(HostnameUtil.getHostnames("0.0.0.0", false)); + hostnames.addAll(advertisedHostnames); hostnames.forEach( hostname -> { - if (IP_ADDR_PATTERN.matcher(hostname).matches()) { + if (InetAddresses.isInetAddress(hostname)) { builder.addIpAddress(hostname); } else { builder.addDnsName(hostname); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaNameSpace.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaNameSpace.java index f95e8ed245b3e..86c52149719a2 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaNameSpace.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaNameSpace.java @@ -703,6 +703,7 @@ public void onMonitoringModeChanged(final List monitoredItems) { /////////////////////////////// Conflict detection /////////////////////////////// public void checkEquals( + final String advertisedHost, final String user, final String password, final String securityDir, @@ -710,6 +711,7 @@ public void checkEquals( final Set securityPolicies, final long debounceTimeMs) { builder.checkEquals( + advertisedHost, user, password, Paths.get(securityDir), diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaServerBuilder.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaServerBuilder.java index 6ac50c959f5f2..687c153351950 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaServerBuilder.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaServerBuilder.java @@ -22,6 +22,7 @@ import org.apache.iotdb.db.i18n.DataNodePipeMessages; import org.apache.iotdb.pipe.api.exception.PipeException; +import com.google.common.net.InetAddresses; import org.eclipse.milo.opcua.sdk.server.OpcUaServer; import org.eclipse.milo.opcua.sdk.server.api.config.OpcUaServerConfig; import org.eclipse.milo.opcua.sdk.server.identity.CompositeValidator; @@ -78,9 +79,11 @@ public class OpcUaServerBuilder implements Closeable { private static final Logger LOGGER = LoggerFactory.getLogger(OpcUaServerBuilder.class); private static final String WILD_CARD_ADDRESS = "0.0.0.0"; + private static final int BACKSLASH = 0x5c; private int tcpBindPort; private int httpsBindPort; + private String advertisedHost; private String user; private String password; private Path securityDir; @@ -99,6 +102,54 @@ public OpcUaServerBuilder setHttpsBindPort(final int httpsBindPort) { return this; } + /** Configures the host published in endpoint URLs without changing the wildcard bind address. */ + public OpcUaServerBuilder setAdvertisedHost(final String advertisedHost) { + this.advertisedHost = normalizeAdvertisedHost(advertisedHost); + return this; + } + + private static String normalizeAdvertisedHost(final String advertisedHost) { + if (Objects.isNull(advertisedHost)) { + return null; + } + + String normalizedAdvertisedHost = advertisedHost.trim(); + if (normalizedAdvertisedHost.isEmpty()) { + throw invalidAdvertisedHost(); + } + final boolean bracketed = + normalizedAdvertisedHost.startsWith("[") || normalizedAdvertisedHost.endsWith("]"); + if (bracketed) { + if (!normalizedAdvertisedHost.startsWith("[") || !normalizedAdvertisedHost.endsWith("]")) { + throw invalidAdvertisedHost(); + } + normalizedAdvertisedHost = + normalizedAdvertisedHost.substring(1, normalizedAdvertisedHost.length() - 1); + } + + final boolean isIpAddress = InetAddresses.isInetAddress(normalizedAdvertisedHost); + if (normalizedAdvertisedHost.isEmpty() + || (bracketed && !isIpAddress) + || normalizedAdvertisedHost.chars().anyMatch(Character::isWhitespace) + || normalizedAdvertisedHost.contains("/") + || normalizedAdvertisedHost.indexOf(BACKSLASH) >= 0 + || normalizedAdvertisedHost.contains("?") + || normalizedAdvertisedHost.contains("#") + || normalizedAdvertisedHost.contains("@") + || normalizedAdvertisedHost.contains("[") + || normalizedAdvertisedHost.contains("]") + || (!isIpAddress && normalizedAdvertisedHost.contains(":"))) { + throw invalidAdvertisedHost(); + } + return normalizedAdvertisedHost; + } + + private static IllegalArgumentException invalidAdvertisedHost() { + return new IllegalArgumentException( + DataNodePipeMessages + .EXCEPTION_THE_ADVERTISED_HOST_MUST_BE_A_HOSTNAME_OR_IP_ADDRESS_WITHOUT_A_SCHEME_PORT_OR_PATH_6857C67A); + } + public OpcUaServerBuilder setUser(final String user) { this.user = user; return this; @@ -147,8 +198,10 @@ public OpcUaServer build() throws Exception { LoggerFactory.getLogger(OpcUaServerBuilder.class) .info(DataNodePipeMessages.OPC_UA_SECURITY_PKI_DIR, pkiDir.getAbsolutePath()); + final Set endpointHostnames = getEndpointHostnames(); + final Set certificateHostnames = getCertificateHostnames(endpointHostnames); final OpcUaKeyStoreLoader loader = - new OpcUaKeyStoreLoader().load(securityDir, password.toCharArray()); + new OpcUaKeyStoreLoader().load(securityDir, password.toCharArray(), certificateHostnames); final DefaultCertificateManager certificateManager = new DefaultCertificateManager(loader.getServerKeyPair(), loader.getServerCertificate()); @@ -165,8 +218,15 @@ public OpcUaServer build() throws Exception { final SelfSignedHttpsCertificateBuilder httpsCertificateBuilder = new SelfSignedHttpsCertificateBuilder(httpsKeyPair); - httpsCertificateBuilder.setCommonName(HostnameUtil.getHostname()); - HostnameUtil.getHostnames(WILD_CARD_ADDRESS).forEach(httpsCertificateBuilder::addDnsName); + httpsCertificateBuilder.setCommonName(certificateHostnames.iterator().next()); + certificateHostnames.forEach( + hostname -> { + if (InetAddresses.isInetAddress(hostname)) { + httpsCertificateBuilder.addIpAddress(hostname); + } else { + httpsCertificateBuilder.addDnsName(hostname); + } + }); final X509Certificate httpsCertificate = httpsCertificateBuilder.build(); final DefaultServerCertificateValidator certificateValidator = @@ -190,6 +250,13 @@ public OpcUaServer build() throws Exception { StatusCodes.Bad_ConfigurationError, DataNodePipeMessages.NO_CERTIFICATE_FOUND)); + if (Objects.nonNull(advertisedHost) && !isAdvertisedHostInCertificate(certificate)) { + LOGGER.warn( + DataNodePipeMessages + .LOG_ADVERTISED_HOST_ARG_IS_NOT_PRESENT_IN_THE_LOADED_OPC_UA_SERVER_CERTIFICATE_SUBJECT_ALTERNATIVE_NAMES_SECURED_CLIENTS_MAY_REJECT_IT_REPLACE_OR_REGENERATE_THE_CERTIFICATE_AND_ESTABLISH_TRUST_AGAIN_912358AF, + advertisedHost); + } + final String applicationUri = CertificateUtil.getSanUri(certificate) .orElseThrow( @@ -199,7 +266,7 @@ public OpcUaServer build() throws Exception { DataNodePipeMessages.CERTIFICATE_MISSING_APPLICATION_URI)); final Set endpointConfigurations = - createEndpointConfigurations(certificate, tcpBindPort, httpsBindPort); + createEndpointConfigurations(certificate, tcpBindPort, httpsBindPort, endpointHostnames); serverConfig = OpcUaServerConfig.builder() @@ -233,19 +300,71 @@ public OpcUaServer build() throws Exception { return server; } - private Set createEndpointConfigurations( - final X509Certificate certificate, final int tcpBindPort, final int httpsBindPort) { + private Set getEndpointHostnames() { + if (Objects.nonNull(advertisedHost)) { + final Set hostnames = new LinkedHashSet<>(); + hostnames.add(toEndpointHostname(advertisedHost)); + return hostnames; + } + final Set hostnames = new LinkedHashSet<>(); + hostnames.add(toEndpointHostname(HostnameUtil.getHostname())); + HostnameUtil.getHostnames(WILD_CARD_ADDRESS).stream() + .map(OpcUaServerBuilder::toEndpointHostname) + .forEach(hostnames::add); + return hostnames; + } + + private static Set getCertificateHostnames(final Set endpointHostnames) { + final Set certificateHostnames = new LinkedHashSet<>(); + endpointHostnames.stream() + .map(OpcUaServerBuilder::removeIpv6Brackets) + .forEach(certificateHostnames::add); + return certificateHostnames; + } + + private static String toEndpointHostname(final String hostname) { + return InetAddresses.isInetAddress(hostname) && hostname.indexOf(':') >= 0 + ? '[' + hostname + ']' + : hostname; + } + + private static String removeIpv6Brackets(final String hostname) { + return hostname.startsWith("[") && hostname.endsWith("]") + ? hostname.substring(1, hostname.length() - 1) + : hostname; + } + + private boolean isAdvertisedHostInCertificate(final X509Certificate certificate) { + if (InetAddresses.isInetAddress(advertisedHost)) { + return CertificateUtil.getSanIpAddresses(certificate).stream() + .filter(InetAddresses::isInetAddress) + .map(InetAddresses::forString) + .anyMatch(InetAddresses.forString(advertisedHost)::equals); + } + return CertificateUtil.getSanDnsNames(certificate).stream() + .anyMatch(hostname -> hostname.equalsIgnoreCase(advertisedHost)); + } + + Set createEndpointConfigurations( + final X509Certificate certificate, + final int tcpBindPort, + final int httpsBindPort, + final Set hostnames) { final Set endpointConfigurations = new LinkedHashSet<>(); + final Set effectiveHostnames = new LinkedHashSet<>(); + if (Objects.nonNull(advertisedHost)) { + effectiveHostnames.add(toEndpointHostname(advertisedHost)); + } else { + hostnames.stream() + .map(OpcUaServerBuilder::toEndpointHostname) + .forEach(effectiveHostnames::add); + } final List bindAddresses = newArrayList(); bindAddresses.add(WILD_CARD_ADDRESS); - final Set hostnames = new LinkedHashSet<>(); - hostnames.add(HostnameUtil.getHostname()); - hostnames.addAll(HostnameUtil.getHostnames(WILD_CARD_ADDRESS)); - for (final String bindAddress : bindAddresses) { - for (final String hostname : hostnames) { + for (final String hostname : effectiveHostnames) { final EndpointConfiguration.Builder builder = EndpointConfiguration.newBuilder() .setBindAddress(bindAddress) @@ -322,12 +441,14 @@ private EndpointConfiguration buildHttpsEndpoint( /////////////////////////////// Conflict detection /////////////////////////////// void checkEquals( + final String advertisedHost, final String user, final String password, final Path securityDir, final boolean enableAnonymousAccess, final Set securityPolicies, final long debounceTimeMs) { + checkEquals("advertised host", this.advertisedHost, normalizeAdvertisedHost(advertisedHost)); checkEquals("user", this.user, user); checkEquals("password", this.password, password); checkEquals( diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaServerBuilderTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaServerBuilderTest.java new file mode 100644 index 0000000000000..97938f56d09fe --- /dev/null +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaServerBuilderTest.java @@ -0,0 +1,164 @@ +/* + * 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.iotdb.db.pipe.sink.protocol.opcua.server; + +import org.apache.iotdb.pipe.api.exception.PipeException; + +import org.eclipse.milo.opcua.sdk.server.OpcUaServer; +import org.eclipse.milo.opcua.stack.core.security.SecurityPolicy; +import org.eclipse.milo.opcua.stack.core.util.CertificateUtil; +import org.eclipse.milo.opcua.stack.server.EndpointConfiguration; +import org.junit.Assert; +import org.junit.Rule; +import org.junit.Test; +import org.junit.rules.TemporaryFolder; + +import java.nio.file.Path; +import java.util.Arrays; +import java.util.Collections; +import java.util.LinkedHashSet; +import java.util.Set; +import java.util.stream.Collectors; + +public class OpcUaServerBuilderTest { + + @Rule public final TemporaryFolder temporaryFolder = new TemporaryFolder(); + + @Test + public void testDetectedHostsArePublishedByDefault() { + final Set securityPolicies = Collections.singleton(SecurityPolicy.None); + final Set detectedHostnames = + new LinkedHashSet<>(Arrays.asList("opc-server", "10.60.80.65")); + final OpcUaServerBuilder builder = + new OpcUaServerBuilder().setSecurityPolicies(securityPolicies); + + final Set endpoints = + builder.createEndpointConfigurations(null, 12686, 8443, detectedHostnames); + + for (final String hostname : detectedHostnames) { + Assert.assertEquals( + 2, + endpoints.stream() + .filter(endpoint -> endpoint.getPath().equals("/iotdb")) + .filter(endpoint -> endpoint.getHostname().equals(hostname)) + .filter(endpoint -> endpoint.getSecurityPolicy() == SecurityPolicy.None) + .count()); + } + Assert.assertEquals(Collections.singleton(SecurityPolicy.None), securityPolicies); + } + + @Test + public void testOnlyExplicitAdvertisedHostIsPublished() { + final Set detectedHostnames = + new LinkedHashSet<>(Arrays.asList("opc-server", "10.60.80.65")); + final OpcUaServerBuilder builder = + new OpcUaServerBuilder() + .setAdvertisedHost("opc.example.com") + .setSecurityPolicies(Collections.singleton(SecurityPolicy.None)); + + final Set endpoints = + builder.createEndpointConfigurations(null, 12686, 8443, detectedHostnames); + + Assert.assertEquals( + Collections.singleton("opc.example.com"), + endpoints.stream().map(EndpointConfiguration::getHostname).collect(Collectors.toSet())); + Assert.assertTrue( + endpoints.stream().allMatch(endpoint -> "0.0.0.0".equals(endpoint.getBindAddress()))); + } + + @Test + public void testIpv6AdvertisedHostIsNormalizedAndEndpointUrlIsRejected() { + final OpcUaServerBuilder builder = + new OpcUaServerBuilder() + .setAdvertisedHost("[2001:db8::1]") + .setSecurityPolicies(Collections.singleton(SecurityPolicy.None)); + + final Set endpoints = + builder.createEndpointConfigurations( + null, 12686, 8443, Collections.singleton("opc-server")); + + Assert.assertEquals( + Collections.singleton("[2001:db8::1]"), + endpoints.stream().map(EndpointConfiguration::getHostname).collect(Collectors.toSet())); + Assert.assertTrue( + endpoints.stream() + .map(EndpointConfiguration::getEndpointUrl) + .allMatch(endpointUrl -> endpointUrl.contains("://[2001:db8::1]:"))); + Assert.assertThrows( + IllegalArgumentException.class, + () -> builder.setAdvertisedHost("opc.tcp://opc.example.com:12686/iotdb")); + } + + @Test + public void testNewCertificateContainsAdvertisedHost() throws Exception { + final String advertisedHost = "opc.example.com"; + final Path securityDir = temporaryFolder.newFolder("security").toPath(); + + try (final OpcUaServerBuilder builder = + new OpcUaServerBuilder() + .setTcpBindPort(12686) + .setHttpsBindPort(8443) + .setAdvertisedHost(advertisedHost) + .setUser("root") + .setPassword("root") + .setSecurityDir(securityDir.toString()) + .setEnableAnonymousAccess(true) + .setSecurityPolicies(Collections.singleton(SecurityPolicy.None)) + .setDebounceTimeMs(50)) { + final OpcUaServer server = builder.build(); + final Set endpoints = server.getConfig().getEndpoints(); + + Assert.assertEquals( + Collections.singleton(advertisedHost), + endpoints.stream().map(EndpointConfiguration::getHostname).collect(Collectors.toSet())); + Assert.assertTrue( + endpoints.stream() + .map(EndpointConfiguration::getCertificate) + .allMatch( + certificate -> + CertificateUtil.getSanDnsNames(certificate).contains(advertisedHost))); + } + } + + @Test + public void testAdvertisedHostParticipatesInConflictDetection() { + final Path securityDir = temporaryFolder.getRoot().toPath(); + final Set securityPolicies = Collections.singleton(SecurityPolicy.None); + final OpcUaServerBuilder builder = + new OpcUaServerBuilder() + .setAdvertisedHost("opc.example.com") + .setUser("root") + .setPassword("root") + .setSecurityDir(securityDir.toString()) + .setEnableAnonymousAccess(true) + .setSecurityPolicies(securityPolicies) + .setDebounceTimeMs(50); + + builder.checkEquals("opc.example.com", "root", "root", securityDir, true, securityPolicies, 50); + final PipeException exception = + Assert.assertThrows( + PipeException.class, + () -> + builder.checkEquals( + "other.example.com", "root", "root", securityDir, true, securityPolicies, 50)); + + Assert.assertTrue(exception.getMessage().contains("advertised host")); + } +} diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/constant/PipeSinkConstant.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/constant/PipeSinkConstant.java index a6dedfad08150..2db1bd0c06b49 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/constant/PipeSinkConstant.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/constant/PipeSinkConstant.java @@ -207,6 +207,10 @@ private static String getDefaultConnectorOrSinkName(final PipeParameters paramet public static final String SINK_OPC_UA_HTTPS_BIND_PORT_KEY = "sink.opcua.https.port"; public static final int CONNECTOR_OPC_UA_HTTPS_BIND_PORT_DEFAULT_VALUE = 8443; + public static final String CONNECTOR_OPC_UA_ADVERTISED_HOST_KEY = + "connector.opcua.advertised-host"; + public static final String SINK_OPC_UA_ADVERTISED_HOST_KEY = "sink.opcua.advertised-host"; + public static final String CONNECTOR_OPC_UA_SECURITY_DIR_KEY = "connector.opcua.security.dir"; public static final String SINK_OPC_UA_SECURITY_DIR_KEY = "sink.opcua.security.dir"; public static final String CONNECTOR_OPC_UA_SECURITY_DIR_DEFAULT_VALUE =