From 5185ed4231ed5cbab21f44bdc061a72fb93ed07f Mon Sep 17 00:00:00 2001
From: Caideyipi <87789683+Caideyipi@users.noreply.github.com>
Date: Tue, 11 Aug 2026 14:15:58 +0800
Subject: [PATCH] Upgrade Eclipse Milo to 1.1.6 (#18431)
* Upgrade Eclipse Milo to 1.1.6
* Add OPC UA TCP None compatibility test
* Update dependency manifest for Milo 1.1.6
* CI fix
---
example/pipe-opc-ua-sink/pom.xml | 4 +-
.../org/apache/iotdb/opcua/ClientExample.java | 4 +-
.../iotdb/opcua/ClientExampleRunner.java | 41 ++-
.../org/apache/iotdb/opcua/ClientTest.java | 54 +---
.../pipe/it/single/IoTDBPipeOPCUAIT.java | 33 +--
iotdb-core/datanode/pom.xml | 14 +-
.../pipe/sink/protocol/opcua/OpcUaSink.java | 6 +-
.../protocol/opcua/client/ClientRunner.java | 46 +++-
.../opcua/client/IoTDBOpcUaClient.java | 24 +-
.../protocol/opcua/server/OpcUaNameSpace.java | 10 +-
.../opcua/server/OpcUaServerBuilder.java | 246 +++++++++---------
.../opcua/client/ClientRunnerTest.java | 2 +-
.../opcua/client/IoTDBOpcUaClientTest.java | 27 +-
.../opcua/server/OpcUaServerBuilderTest.java | 104 ++++++--
.../server/OpcUaTcpNoneCompatibilityTest.java | 132 ++++++++++
pom.xml | 41 +--
16 files changed, 489 insertions(+), 299 deletions(-)
create mode 100644 iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaTcpNoneCompatibilityTest.java
diff --git a/example/pipe-opc-ua-sink/pom.xml b/example/pipe-opc-ua-sink/pom.xml
index f9564ba01dc3..b3f3d4e99d7d 100644
--- a/example/pipe-opc-ua-sink/pom.xml
+++ b/example/pipe-opc-ua-sink/pom.xml
@@ -31,12 +31,12 @@
org.eclipse.milo
- sdk-client
+ milo-sdk-client
${milo.version}
org.eclipse.milo
- sdk-server
+ milo-sdk-server
diff --git a/example/pipe-opc-ua-sink/src/main/java/org/apache/iotdb/opcua/ClientExample.java b/example/pipe-opc-ua-sink/src/main/java/org/apache/iotdb/opcua/ClientExample.java
index 6b7f6997763b..dba324c755d4 100644
--- a/example/pipe-opc-ua-sink/src/main/java/org/apache/iotdb/opcua/ClientExample.java
+++ b/example/pipe-opc-ua-sink/src/main/java/org/apache/iotdb/opcua/ClientExample.java
@@ -20,8 +20,8 @@
package org.apache.iotdb.opcua;
import org.eclipse.milo.opcua.sdk.client.OpcUaClient;
-import org.eclipse.milo.opcua.sdk.client.api.identity.AnonymousProvider;
-import org.eclipse.milo.opcua.sdk.client.api.identity.IdentityProvider;
+import org.eclipse.milo.opcua.sdk.client.identity.AnonymousProvider;
+import org.eclipse.milo.opcua.sdk.client.identity.IdentityProvider;
import org.eclipse.milo.opcua.stack.core.security.SecurityPolicy;
import org.eclipse.milo.opcua.stack.core.types.structured.EndpointDescription;
diff --git a/example/pipe-opc-ua-sink/src/main/java/org/apache/iotdb/opcua/ClientExampleRunner.java b/example/pipe-opc-ua-sink/src/main/java/org/apache/iotdb/opcua/ClientExampleRunner.java
index 1fe49a75007a..60b1f8b2ee68 100644
--- a/example/pipe-opc-ua-sink/src/main/java/org/apache/iotdb/opcua/ClientExampleRunner.java
+++ b/example/pipe-opc-ua-sink/src/main/java/org/apache/iotdb/opcua/ClientExampleRunner.java
@@ -21,13 +21,14 @@
import org.bouncycastle.jce.provider.BouncyCastleProvider;
import org.eclipse.milo.opcua.sdk.client.OpcUaClient;
-import org.eclipse.milo.opcua.stack.client.security.DefaultClientCertificateValidator;
import org.eclipse.milo.opcua.stack.core.Stack;
-import org.eclipse.milo.opcua.stack.core.security.DefaultTrustListManager;
+import org.eclipse.milo.opcua.stack.core.security.DefaultClientCertificateValidator;
+import org.eclipse.milo.opcua.stack.core.security.FileBasedCertificateQuarantine;
+import org.eclipse.milo.opcua.stack.core.security.FileBasedTrustListManager;
import org.eclipse.milo.opcua.stack.core.types.builtin.LocalizedText;
import org.slf4j.LoggerFactory;
-import java.io.File;
+import java.io.IOException;
import java.nio.file.Files;
import java.nio.file.Path;
import java.nio.file.Paths;
@@ -48,6 +49,7 @@ public class ClientExampleRunner {
private final CompletableFuture future = new CompletableFuture<>();
private final ClientExample clientExample;
+ private FileBasedTrustListManager trustListManager;
public ClientExampleRunner(ClientExample clientExample) {
this.clientExample = clientExample;
@@ -61,21 +63,24 @@ private OpcUaClient createClient() throws Exception {
throw new Exception("unable to create security dir: " + securityTempDir);
}
- final File pkiDir = securityTempDir.resolve("pki").toFile();
+ final Path pkiDir = securityTempDir.resolve("pki");
System.out.println("security dir: " + securityTempDir.toAbsolutePath());
- LoggerFactory.getLogger(getClass()).info("security pki dir: {}", pkiDir.getAbsolutePath());
+ LoggerFactory.getLogger(getClass()).info("security pki dir: {}", pkiDir.toAbsolutePath());
final IoTDBKeyStoreLoaderClient loader = new IoTDBKeyStoreLoaderClient().load(securityTempDir);
- final DefaultTrustListManager trustListManager = new DefaultTrustListManager(pkiDir);
+ trustListManager = FileBasedTrustListManager.createAndInitialize(pkiDir);
+ final FileBasedCertificateQuarantine certificateQuarantine =
+ FileBasedCertificateQuarantine.create(pkiDir.resolve("rejected").resolve("certs"));
final DefaultClientCertificateValidator certificateValidator =
- new DefaultClientCertificateValidator(trustListManager);
+ new DefaultClientCertificateValidator(trustListManager, certificateQuarantine);
return OpcUaClient.create(
clientExample.getEndpointUrl(),
endpoints -> endpoints.stream().filter(clientExample.endpointFilter()).findFirst(),
+ transportBuilder -> {},
configBuilder ->
configBuilder
.setApplicationName(LocalizedText.english("eclipse milo opc-ua client"))
@@ -85,8 +90,7 @@ private OpcUaClient createClient() throws Exception {
.setCertificateChain(loader.getClientCertificateChain())
.setCertificateValidator(certificateValidator)
.setIdentityProvider(clientExample.getIdentityProvider())
- .setRequestTimeout(uint(5000))
- .build());
+ .setRequestTimeout(uint(5000)));
}
public void run() {
@@ -100,11 +104,13 @@ public void run() {
}
try {
- client.disconnect().get();
- Stack.releaseSharedResources();
+ client.disconnectAsync().get();
} catch (InterruptedException | ExecutionException e) {
Thread.currentThread().interrupt();
System.out.println("Error disconnecting: {}" + e.getMessage());
+ } finally {
+ closeTrustListManager();
+ Stack.releaseSharedResources();
}
try {
@@ -126,6 +132,7 @@ public void run() {
} catch (Throwable t) {
System.out.println("Error getting client: {}" + t.getMessage());
+ closeTrustListManager();
future.completeExceptionally(t);
try {
@@ -144,4 +151,16 @@ public void run() {
e.printStackTrace();
}
}
+
+ private void closeTrustListManager() {
+ if (trustListManager != null) {
+ try {
+ trustListManager.close();
+ } catch (IOException e) {
+ e.printStackTrace();
+ } finally {
+ trustListManager = null;
+ }
+ }
+ }
}
diff --git a/example/pipe-opc-ua-sink/src/main/java/org/apache/iotdb/opcua/ClientTest.java b/example/pipe-opc-ua-sink/src/main/java/org/apache/iotdb/opcua/ClientTest.java
index cc09a7dc7bb0..6db181f56c4e 100644
--- a/example/pipe-opc-ua-sink/src/main/java/org/apache/iotdb/opcua/ClientTest.java
+++ b/example/pipe-opc-ua-sink/src/main/java/org/apache/iotdb/opcua/ClientTest.java
@@ -20,27 +20,18 @@
package org.apache.iotdb.opcua;
import org.eclipse.milo.opcua.sdk.client.OpcUaClient;
-import org.eclipse.milo.opcua.sdk.client.api.subscriptions.UaMonitoredItem;
-import org.eclipse.milo.opcua.sdk.client.api.subscriptions.UaSubscription;
+import org.eclipse.milo.opcua.sdk.client.subscriptions.OpcUaMonitoredItem;
+import org.eclipse.milo.opcua.sdk.client.subscriptions.OpcUaSubscription;
import org.eclipse.milo.opcua.stack.core.AttributeId;
import org.eclipse.milo.opcua.stack.core.Identifiers;
-import org.eclipse.milo.opcua.stack.core.types.builtin.ExtensionObject;
import org.eclipse.milo.opcua.stack.core.types.builtin.QualifiedName;
-import org.eclipse.milo.opcua.stack.core.types.builtin.unsigned.UInteger;
import org.eclipse.milo.opcua.stack.core.types.enumerated.MonitoringMode;
-import org.eclipse.milo.opcua.stack.core.types.enumerated.TimestampsToReturn;
import org.eclipse.milo.opcua.stack.core.types.structured.ContentFilter;
import org.eclipse.milo.opcua.stack.core.types.structured.EventFilter;
-import org.eclipse.milo.opcua.stack.core.types.structured.MonitoredItemCreateRequest;
-import org.eclipse.milo.opcua.stack.core.types.structured.MonitoringParameters;
import org.eclipse.milo.opcua.stack.core.types.structured.ReadValueId;
import org.eclipse.milo.opcua.stack.core.types.structured.SimpleAttributeOperand;
-import java.util.Collections;
-import java.util.List;
import java.util.concurrent.CompletableFuture;
-import java.util.concurrent.atomic.AtomicInteger;
-import java.util.concurrent.atomic.AtomicLong;
import static org.eclipse.milo.opcua.stack.core.types.builtin.unsigned.Unsigned.uint;
@@ -52,24 +43,19 @@ public static void main(String[] args) {
new ClientExampleRunner(example).run();
}
- private final AtomicLong clientHandles = new AtomicLong(1L);
-
@Override
public void run(OpcUaClient client, CompletableFuture future) throws Exception {
// synchronous connect
- client.connect().get();
+ client.connect();
// create a subscription and a monitored item
- final UaSubscription subscription =
- client.getSubscriptionManager().createSubscription(200.0).get();
+ final OpcUaSubscription subscription = new OpcUaSubscription(client, 200.0);
+ subscription.create();
final ReadValueId readValueId =
new ReadValueId(
Identifiers.Server, AttributeId.EventNotifier.uid(), null, QualifiedName.NULL_VALUE);
- // client handle must be unique per item
- final UInteger clientHandle = uint(clientHandles.getAndIncrement());
-
final EventFilter eventFilter =
new EventFilter(
new SimpleAttributeOperand[] {
@@ -96,30 +82,18 @@ public void run(OpcUaClient client, CompletableFuture future) throw
},
new ContentFilter(null));
- final MonitoringParameters parameters =
- new MonitoringParameters(
- clientHandle,
- 0.0,
- ExtensionObject.encode(client.getStaticSerializationContext(), eventFilter),
- uint(10000),
- true);
-
- final MonitoredItemCreateRequest request =
- new MonitoredItemCreateRequest(readValueId, MonitoringMode.Reporting, parameters);
-
- final List items =
- subscription
- .createMonitoredItems(TimestampsToReturn.Both, Collections.singletonList(request))
- .get();
+ final OpcUaMonitoredItem monitoredItem =
+ new OpcUaMonitoredItem(readValueId, MonitoringMode.Reporting);
+ monitoredItem.setSamplingInterval(0.0);
+ monitoredItem.setFilter(eventFilter);
+ monitoredItem.setQueueSize(uint(10000));
+ monitoredItem.setDiscardOldest(true);
+ subscription.addMonitoredItem(monitoredItem);
+ subscription.synchronizeMonitoredItems();
// do something with the value updates
- final UaMonitoredItem monitoredItem = items.get(0);
-
- final AtomicInteger eventCount = new AtomicInteger(0);
-
- monitoredItem.setEventConsumer(
+ monitoredItem.setEventValueListener(
(item, vs) -> {
- eventCount.incrementAndGet();
System.out.println("Event Received from " + item.getReadValueId().getNodeId());
for (int i = 0; i < vs.length; i++) {
diff --git a/integration-test/src/test/java/org/apache/iotdb/pipe/it/single/IoTDBPipeOPCUAIT.java b/integration-test/src/test/java/org/apache/iotdb/pipe/it/single/IoTDBPipeOPCUAIT.java
index 33ee72a61e8c..089887729b4a 100644
--- a/integration-test/src/test/java/org/apache/iotdb/pipe/it/single/IoTDBPipeOPCUAIT.java
+++ b/integration-test/src/test/java/org/apache/iotdb/pipe/it/single/IoTDBPipeOPCUAIT.java
@@ -33,9 +33,9 @@
import org.apache.tsfile.common.conf.TSFileConfig;
import org.eclipse.milo.opcua.sdk.client.OpcUaClient;
-import org.eclipse.milo.opcua.sdk.client.api.identity.AnonymousProvider;
-import org.eclipse.milo.opcua.sdk.client.api.identity.IdentityProvider;
-import org.eclipse.milo.opcua.sdk.client.api.identity.UsernameProvider;
+import org.eclipse.milo.opcua.sdk.client.identity.AnonymousProvider;
+import org.eclipse.milo.opcua.sdk.client.identity.IdentityProvider;
+import org.eclipse.milo.opcua.sdk.client.identity.UsernameProvider;
import org.eclipse.milo.opcua.stack.core.security.SecurityPolicy;
import org.eclipse.milo.opcua.stack.core.types.builtin.DataValue;
import org.eclipse.milo.opcua.stack.core.types.builtin.DateTime;
@@ -107,13 +107,10 @@ public void testOPCUAServerSink() throws Exception {
throw e;
}
}
- value =
- opcUaClient
- .readValue(0, TimestampsToReturn.Both, new NodeId(2, "root/db/d1/`1`"))
- .get();
+ value = opcUaClient.readValue(0, TimestampsToReturn.Both, new NodeId(2, "root/db/d1/`1`"));
Assert.assertEquals(new Variant(1.0), value.getValue());
Assert.assertEquals(new DateTime(timestampToUtc(1)), value.getSourceTime());
- opcUaClient.disconnect().get();
+ opcUaClient.disconnect();
break;
}
@@ -174,26 +171,19 @@ public void testOPCUAServerSink() throws Exception {
long startTime = System.currentTimeMillis();
while (true) {
try {
- value =
- opcUaClient
- .readValue(0, TimestampsToReturn.Both, new NodeId(2, "root/db/`123`"))
- .get();
+ value = opcUaClient.readValue(0, TimestampsToReturn.Both, new NodeId(2, "root/db/`123`"));
Assert.assertEquals(new Variant(1.0), value.getValue());
Assert.assertEquals(StatusCode.BAD, value.getStatusCode());
Assert.assertEquals(new DateTime(timestampToUtc(1)), value.getSourceTime());
value =
- opcUaClient
- .readValue(0, TimestampsToReturn.Both, new NodeId(2, "root/db/`1231`"))
- .get();
+ opcUaClient.readValue(0, TimestampsToReturn.Both, new NodeId(2, "root/db/`1231`"));
Assert.assertEquals(new Variant(1.0), value.getValue());
Assert.assertEquals(StatusCode.BAD, value.getStatusCode());
Assert.assertEquals(new DateTime(timestampToUtc(1)), value.getSourceTime());
value =
- opcUaClient
- .readValue(0, TimestampsToReturn.Both, new NodeId(2, "root/db/`1232`"))
- .get();
+ opcUaClient.readValue(0, TimestampsToReturn.Both, new NodeId(2, "root/db/`1232`"));
Assert.assertEquals(new Variant(1.0), value.getValue());
Assert.assertEquals(StatusCode.BAD, value.getStatusCode());
Assert.assertEquals(new DateTime(timestampToUtc(1)), value.getSourceTime());
@@ -212,10 +202,7 @@ public void testOPCUAServerSink() throws Exception {
startTime = System.currentTimeMillis();
while (true) {
try {
- value =
- opcUaClient
- .readValue(0, TimestampsToReturn.Both, new NodeId(2, "root/db/`123`"))
- .get();
+ value = opcUaClient.readValue(0, TimestampsToReturn.Both, new NodeId(2, "root/db/`123`"));
Assert.assertEquals(new DateTime(timestampToUtc(2)), value.getSourceTime());
Assert.assertEquals(new Variant(2.0), value.getValue());
Assert.assertEquals(StatusCode.UNCERTAIN, value.getStatusCode());
@@ -227,7 +214,7 @@ public void testOPCUAServerSink() throws Exception {
}
}
- opcUaClient.disconnect().get();
+ opcUaClient.disconnect();
Assert.assertEquals(
TSStatusCode.SUCCESS_STATUS.getStatusCode(), client.dropPipe("testPipe").getCode());
diff --git a/iotdb-core/datanode/pom.xml b/iotdb-core/datanode/pom.xml
index 299e50f1be93..52b259e4fe46 100644
--- a/iotdb-core/datanode/pom.xml
+++ b/iotdb-core/datanode/pom.xml
@@ -180,11 +180,11 @@
org.eclipse.milo
- stack-core
+ milo-stack-core
org.eclipse.milo
- sdk-core
+ milo-sdk-core
commons-io
@@ -200,7 +200,7 @@
org.eclipse.milo
- stack-server
+ milo-transport
org.checkerframework
@@ -208,11 +208,7 @@
org.eclipse.milo
- stack-client
-
-
- org.eclipse.milo
- sdk-client
+ milo-sdk-client
org.bouncycastle
@@ -252,7 +248,7 @@
org.eclipse.milo
- sdk-server
+ milo-sdk-server
commons-cli
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 d6cf17c07adf..926c8c136139 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
@@ -37,9 +37,9 @@
import org.apache.tsfile.common.conf.TSFileConfig;
import org.apache.tsfile.utils.Pair;
import org.apache.tsfile.write.record.Tablet;
-import org.eclipse.milo.opcua.sdk.client.api.identity.AnonymousProvider;
-import org.eclipse.milo.opcua.sdk.client.api.identity.IdentityProvider;
-import org.eclipse.milo.opcua.sdk.client.api.identity.UsernameProvider;
+import org.eclipse.milo.opcua.sdk.client.identity.AnonymousProvider;
+import org.eclipse.milo.opcua.sdk.client.identity.IdentityProvider;
+import org.eclipse.milo.opcua.sdk.client.identity.UsernameProvider;
import org.eclipse.milo.opcua.sdk.server.OpcUaServer;
import org.eclipse.milo.opcua.stack.core.security.SecurityPolicy;
import org.eclipse.milo.opcua.stack.core.types.builtin.StatusCode;
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 5ed25f47935f..d90bd33a49e0 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
@@ -23,8 +23,9 @@
import org.bouncycastle.jce.provider.BouncyCastleProvider;
import org.eclipse.milo.opcua.sdk.client.OpcUaClient;
-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.DefaultClientCertificateValidator;
+import org.eclipse.milo.opcua.stack.core.security.FileBasedCertificateQuarantine;
+import org.eclipse.milo.opcua.stack.core.security.FileBasedTrustListManager;
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;
@@ -33,7 +34,8 @@
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
-import java.io.File;
+import java.io.Closeable;
+import java.io.IOException;
import java.nio.file.FileSystems;
import java.nio.file.Files;
import java.nio.file.Path;
@@ -46,7 +48,7 @@
import static org.eclipse.milo.opcua.stack.core.types.builtin.unsigned.Unsigned.uint;
-public class ClientRunner {
+public class ClientRunner implements Closeable {
private static final Logger logger = LoggerFactory.getLogger(ClientRunner.class);
@@ -60,6 +62,7 @@ public class ClientRunner {
private final String password;
private final long timeoutSeconds;
private final boolean allowEndpointRedirect;
+ private FileBasedTrustListManager trustListManager;
// For conflict checking
private final String user;
@@ -95,18 +98,20 @@ private OpcUaClient createClient() throws Exception {
throw new Exception("unable to create security dir: " + securityDir);
}
- final File pkiDir = securityDir.resolve("pki").toFile();
+ final Path pkiDir = securityDir.resolve("pki");
logger.info("security dir: {}", securityDir.toAbsolutePath());
- logger.info("security pki dir: {}", pkiDir.getAbsolutePath());
+ logger.info("security pki dir: {}", pkiDir.toAbsolutePath());
final IoTDBKeyStoreLoaderClient loader =
new IoTDBKeyStoreLoaderClient().load(securityDir, password.toCharArray());
- final DefaultTrustListManager trustListManager = new DefaultTrustListManager(pkiDir);
+ trustListManager = FileBasedTrustListManager.createAndInitialize(pkiDir);
+ final FileBasedCertificateQuarantine certificateQuarantine =
+ FileBasedCertificateQuarantine.create(pkiDir.resolve("rejected").resolve("certs"));
final DefaultClientCertificateValidator certificateValidator =
- new DefaultClientCertificateValidator(trustListManager);
+ new DefaultClientCertificateValidator(trustListManager, certificateQuarantine);
return OpcUaClient.create(
configurableUaClient.getNodeUrl(),
@@ -116,6 +121,7 @@ private OpcUaClient createClient() throws Exception {
configurableUaClient.getNodeUrl(),
configurableUaClient.getSecurityPolicy(),
allowEndpointRedirect),
+ transportBuilder -> transportBuilder.setConnectTimeout(uint(timeoutSeconds * 1000L)),
configBuilder ->
configBuilder
.setApplicationName(LocalizedText.english("Apache IoTDB OPC UA client"))
@@ -126,9 +132,7 @@ private OpcUaClient createClient() throws Exception {
.setCertificateValidator(certificateValidator)
.setIdentityProvider(configurableUaClient.getIdentityProvider())
.setRequestTimeout(uint(timeoutSeconds * 1000L))
- .setConnectTimeout(uint(timeoutSeconds * 1000L))
- .setMaxResponseMessageSize(uint(0))
- .build());
+ .setMaxResponseMessageSize(uint(0)));
}
static Optional selectEndpoint(
@@ -201,11 +205,31 @@ public void run() {
"Error running opc client: " + e.getClass().getSimpleName() + ": " + e.getMessage());
}
} catch (final Exception e) {
+ closeOnFailure(e);
throw new PipeException(
"Error getting opc client: " + e.getClass().getSimpleName() + ": " + e.getMessage());
}
}
+ private void closeOnFailure(final Exception failure) {
+ try {
+ close();
+ } catch (final IOException closeException) {
+ failure.addSuppressed(closeException);
+ }
+ }
+
+ @Override
+ public void close() throws IOException {
+ if (trustListManager != null) {
+ try {
+ trustListManager.close();
+ } finally {
+ trustListManager = null;
+ }
+ }
+ }
+
long getTimeoutSeconds() {
return timeoutSeconds;
}
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 cf8ab9001fb3..afc3ccb2e654 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
@@ -30,7 +30,7 @@
import org.apache.tsfile.write.record.Tablet;
import org.apache.tsfile.write.schema.MeasurementSchema;
import org.eclipse.milo.opcua.sdk.client.OpcUaClient;
-import org.eclipse.milo.opcua.sdk.client.api.identity.IdentityProvider;
+import org.eclipse.milo.opcua.sdk.client.identity.IdentityProvider;
import org.eclipse.milo.opcua.sdk.core.AccessLevel;
import org.eclipse.milo.opcua.sdk.core.ValueRanks;
import org.eclipse.milo.opcua.stack.core.Identifiers;
@@ -103,7 +103,7 @@ public void run(final OpcUaClient client) throws Exception {
long startTime = System.currentTimeMillis();
while (System.currentTimeMillis() - startTime < runner.getTimeoutSeconds() * 1000L) {
try {
- client.connect().get();
+ client.connectAsync().get();
} catch (final ExecutionException e) {
if (e.getCause() instanceof UaException
&& ((UaException) e.getCause()).getStatusCode().getValue() == Bad_Timeout) {
@@ -232,7 +232,7 @@ private void addMissingNodes(final List writeRequests) throws
}
}
- final AddNodesResponse addStatus = client.addNodes(nodesToAdd).get();
+ final AddNodesResponse addStatus = client.addNodesAsync(nodesToAdd).get();
for (final AddNodesResult result : addStatus.getResults()) {
if (!result.getStatusCode().equals(StatusCode.GOOD)
&& result.getStatusCode().getValue() != StatusCodes.Bad_NodeIdExists) {
@@ -265,7 +265,7 @@ private List writeValuesOnce(final List writeRequ
nodeIds.add(writeRequest.nodeId);
dataValues.add(writeRequest.dataValue);
}
- return client.writeValues(nodeIds, dataValues).get();
+ return client.writeValuesAsync(nodeIds, dataValues).get();
}
private static final class OpcUaWriteRequest {
@@ -339,7 +339,7 @@ public List getNodesToAdd(
new QualifiedName(NAME_SPACE_INDEX, segments[0]),
NodeClass.Object,
ExtensionObject.encode(
- client.getStaticSerializationContext(), createFolderAttributes(segments[0])),
+ client.getStaticEncodingContext(), createFolderAttributes(segments[0])),
Identifiers.FolderType.expanded()));
// segments.length >= 3
@@ -354,7 +354,7 @@ public List getNodesToAdd(
new QualifiedName(NAME_SPACE_INDEX, segments[i]),
NodeClass.Object,
ExtensionObject.encode(
- client.getStaticSerializationContext(), createFolderAttributes(segments[i])),
+ client.getStaticEncodingContext(), createFolderAttributes(segments[i])),
Identifiers.FolderType.expanded()));
curNodeId = nextId;
}
@@ -369,7 +369,7 @@ public List getNodesToAdd(
new QualifiedName(NAME_SPACE_INDEX, measurementName),
NodeClass.Variable,
ExtensionObject.encode(
- client.getStaticSerializationContext(),
+ client.getStaticEncodingContext(),
createMeasurementAttributes(measurementName, opcDataType, initialValue)),
Identifiers.BaseDataVariableType.expanded()));
@@ -377,8 +377,14 @@ public List getNodesToAdd(
}
public void disconnect() throws Exception {
- if (Objects.nonNull(client)) {
- client.disconnect().get();
+ try {
+ if (Objects.nonNull(client)) {
+ client.disconnectAsync().get();
+ }
+ } finally {
+ if (Objects.nonNull(runner)) {
+ runner.close();
+ }
}
}
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 64fd9fce3107..b3d078498252 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
@@ -37,11 +37,11 @@
import org.eclipse.milo.opcua.sdk.core.AccessLevel;
import org.eclipse.milo.opcua.sdk.core.Reference;
import org.eclipse.milo.opcua.sdk.server.Lifecycle;
+import org.eclipse.milo.opcua.sdk.server.ManagedNamespaceWithLifecycle;
import org.eclipse.milo.opcua.sdk.server.OpcUaServer;
-import org.eclipse.milo.opcua.sdk.server.api.DataItem;
-import org.eclipse.milo.opcua.sdk.server.api.ManagedNamespaceWithLifecycle;
-import org.eclipse.milo.opcua.sdk.server.api.MonitoredItem;
-import org.eclipse.milo.opcua.sdk.server.model.nodes.objects.BaseEventTypeNode;
+import org.eclipse.milo.opcua.sdk.server.items.DataItem;
+import org.eclipse.milo.opcua.sdk.server.items.MonitoredItem;
+import org.eclipse.milo.opcua.sdk.server.model.objects.BaseEventTypeNode;
import org.eclipse.milo.opcua.sdk.server.nodes.UaFolderNode;
import org.eclipse.milo.opcua.sdk.server.nodes.UaNode;
import org.eclipse.milo.opcua.sdk.server.nodes.UaVariableNode;
@@ -439,7 +439,7 @@ private void transferTabletForPubSubModel(final Tablet tablet) throws UaExceptio
}
// Send the event
- getServer().getEventBus().post(eventNode);
+ getServer().getEventNotifier().fire(eventNode);
}
}
eventNode.delete();
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 adc505b87168..8f9e23c7a02b 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,19 +22,30 @@
import org.apache.iotdb.pipe.api.exception.PipeException;
import com.google.common.net.InetAddresses;
+import org.eclipse.milo.opcua.sdk.server.EndpointConfig;
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.OpcUaServerConfig;
+import org.eclipse.milo.opcua.sdk.server.diagnostics.SessionSecurityDiagnosticsAccessMode;
+import org.eclipse.milo.opcua.sdk.server.identity.AnonymousIdentityValidator;
import org.eclipse.milo.opcua.sdk.server.identity.CompositeValidator;
+import org.eclipse.milo.opcua.sdk.server.identity.IdentityValidator;
import org.eclipse.milo.opcua.sdk.server.identity.UsernameIdentityValidator;
import org.eclipse.milo.opcua.sdk.server.identity.X509IdentityValidator;
-import org.eclipse.milo.opcua.sdk.server.model.nodes.objects.ServerTypeNode;
+import org.eclipse.milo.opcua.sdk.server.model.objects.ServerTypeNode;
import org.eclipse.milo.opcua.sdk.server.nodes.UaNode;
import org.eclipse.milo.opcua.sdk.server.util.HostnameUtil;
import org.eclipse.milo.opcua.stack.core.Identifiers;
+import org.eclipse.milo.opcua.stack.core.NodeIds;
import org.eclipse.milo.opcua.stack.core.StatusCodes;
+import org.eclipse.milo.opcua.stack.core.UaException;
import org.eclipse.milo.opcua.stack.core.UaRuntimeException;
+import org.eclipse.milo.opcua.stack.core.security.DefaultApplicationGroup;
import org.eclipse.milo.opcua.stack.core.security.DefaultCertificateManager;
-import org.eclipse.milo.opcua.stack.core.security.DefaultTrustListManager;
+import org.eclipse.milo.opcua.stack.core.security.DefaultServerCertificateValidator;
+import org.eclipse.milo.opcua.stack.core.security.FileBasedCertificateQuarantine;
+import org.eclipse.milo.opcua.stack.core.security.FileBasedTrustListManager;
+import org.eclipse.milo.opcua.stack.core.security.MemoryCertificateStore;
+import org.eclipse.milo.opcua.stack.core.security.RsaSha256CertificateFactory;
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.DateTime;
@@ -42,15 +53,12 @@
import org.eclipse.milo.opcua.stack.core.types.enumerated.MessageSecurityMode;
import org.eclipse.milo.opcua.stack.core.types.structured.BuildInfo;
import org.eclipse.milo.opcua.stack.core.util.CertificateUtil;
-import org.eclipse.milo.opcua.stack.core.util.SelfSignedCertificateGenerator;
-import org.eclipse.milo.opcua.stack.core.util.SelfSignedHttpsCertificateBuilder;
-import org.eclipse.milo.opcua.stack.server.EndpointConfiguration;
-import org.eclipse.milo.opcua.stack.server.security.DefaultServerCertificateValidator;
+import org.eclipse.milo.opcua.stack.transport.server.tcp.OpcTcpServerTransport;
+import org.eclipse.milo.opcua.stack.transport.server.tcp.OpcTcpServerTransportConfig;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.io.Closeable;
-import java.io.File;
import java.io.IOException;
import java.nio.file.FileSystems;
import java.nio.file.Files;
@@ -58,16 +66,16 @@
import java.nio.file.Paths;
import java.security.KeyPair;
import java.security.cert.X509Certificate;
+import java.util.ArrayList;
import java.util.HashSet;
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Objects;
import java.util.Set;
-import static com.google.common.collect.Lists.newArrayList;
-import static org.eclipse.milo.opcua.sdk.server.api.config.OpcUaServerConfig.USER_TOKEN_POLICY_ANONYMOUS;
-import static org.eclipse.milo.opcua.sdk.server.api.config.OpcUaServerConfig.USER_TOKEN_POLICY_USERNAME;
-import static org.eclipse.milo.opcua.sdk.server.api.config.OpcUaServerConfig.USER_TOKEN_POLICY_X509;
+import static org.eclipse.milo.opcua.sdk.server.OpcUaServerConfig.USER_TOKEN_POLICY_ANONYMOUS;
+import static org.eclipse.milo.opcua.sdk.server.OpcUaServerConfig.USER_TOKEN_POLICY_USERNAME;
+import static org.eclipse.milo.opcua.sdk.server.OpcUaServerConfig.USER_TOKEN_POLICY_X509;
import static org.eclipse.milo.opcua.stack.core.types.builtin.unsigned.Unsigned.ubyte;
/**
@@ -88,7 +96,7 @@ public class OpcUaServerBuilder implements Closeable {
private Path securityDir;
private boolean enableAnonymousAccess;
private Set securityPolicies;
- private DefaultTrustListManager trustListManager;
+ private FileBasedTrustListManager trustListManager;
private long debounceTimeMs;
public OpcUaServerBuilder setTcpBindPort(final int tcpBindPort) {
@@ -189,59 +197,75 @@ public OpcUaServer build() throws Exception {
throw new PipeException("Unable to create security dir: " + securityDir);
}
- final File pkiDir = securityDir.resolve("pki").toFile();
+ final Path pkiDir = securityDir.resolve("pki");
LoggerFactory.getLogger(OpcUaServerBuilder.class)
.info("Security dir: {}", securityDir.toAbsolutePath());
LoggerFactory.getLogger(OpcUaServerBuilder.class)
- .info("Security pki dir: {}", pkiDir.getAbsolutePath());
+ .info("Security pki dir: {}", pkiDir.toAbsolutePath());
final Set endpointHostnames = getEndpointHostnames();
final Set certificateHostnames = getCertificateHostnames(endpointHostnames);
final OpcUaKeyStoreLoader loader =
new OpcUaKeyStoreLoader().load(securityDir, password.toCharArray(), certificateHostnames);
- final DefaultCertificateManager certificateManager =
- new DefaultCertificateManager(loader.getServerKeyPair(), loader.getServerCertificate());
-
- final OpcUaServerConfig serverConfig;
+ trustListManager = FileBasedTrustListManager.createAndInitialize(pkiDir);
+ final FileBasedCertificateQuarantine certificateQuarantine =
+ FileBasedCertificateQuarantine.create(pkiDir.resolve("rejected").resolve("certs"));
+ // OpcUaKeyStoreLoader already persists the application key pair in iotdb-server.pfx. Keeping
+ // Milo's application-group store in memory avoids a second password-protected stale copy.
+ final MemoryCertificateStore certificateStore = new MemoryCertificateStore();
+ final RsaSha256CertificateFactory certificateFactory =
+ new RsaSha256CertificateFactory() {
+ @Override
+ protected KeyPair createRsaSha256KeyPair() {
+ return loader.getServerKeyPair();
+ }
- trustListManager = new DefaultTrustListManager(pkiDir);
+ @Override
+ protected X509Certificate[] createRsaSha256CertificateChain(final KeyPair keyPair) {
+ return new X509Certificate[] {loader.getServerCertificate()};
+ }
+ };
+ final DefaultServerCertificateValidator certificateValidator =
+ new DefaultServerCertificateValidator(trustListManager, certificateQuarantine);
+ final DefaultApplicationGroup applicationGroup =
+ DefaultApplicationGroup.createAndInitialize(
+ trustListManager, certificateStore, certificateFactory, certificateValidator);
+ final DefaultCertificateManager certificateManager =
+ new DefaultCertificateManager(certificateQuarantine, applicationGroup);
LOGGER.info(
"Certificate directory is: {}, Please move certificates from the reject dir to the trusted directory to allow encrypted access",
- pkiDir.getAbsolutePath());
-
- final KeyPair httpsKeyPair = SelfSignedCertificateGenerator.generateRsaKeyPair(2048);
-
- final SelfSignedHttpsCertificateBuilder httpsCertificateBuilder =
- new SelfSignedHttpsCertificateBuilder(httpsKeyPair);
- 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 =
- new DefaultServerCertificateValidator(trustListManager);
+ pkiDir.toAbsolutePath());
- final UsernameIdentityValidator identityValidator =
+ final UsernameIdentityValidator usernameIdentityValidator =
new UsernameIdentityValidator(
- enableAnonymousAccess,
authChallenge ->
- authChallenge.getUsername().equals(user)
- && authChallenge.getPassword().equals(password));
-
- final X509IdentityValidator x509IdentityValidator = new X509IdentityValidator(c -> true);
+ Objects.equals(authChallenge.getUsername(), user)
+ && Objects.equals(authChallenge.getPassword(), password));
+ final X509IdentityValidator x509IdentityValidator =
+ new X509IdentityValidator(
+ userCertificate -> {
+ try {
+ certificateValidator.validateCertificateChain(List.of(userCertificate), null, null);
+ return true;
+ } catch (final UaException ignored) {
+ return false;
+ }
+ });
+ final List identityValidators = new ArrayList<>();
+ if (enableAnonymousAccess) {
+ identityValidators.add(AnonymousIdentityValidator.INSTANCE);
+ }
+ identityValidators.add(usernameIdentityValidator);
+ identityValidators.add(x509IdentityValidator);
final X509Certificate certificate =
- certificateManager.getCertificates().stream()
- .findFirst()
+ applicationGroup
+ .getCertificateChain(NodeIds.RsaSha256ApplicationCertificateType)
+ .filter(certificateChain -> certificateChain.length > 0)
+ .map(certificateChain -> certificateChain[0])
.orElseThrow(
() ->
new UaRuntimeException(
@@ -261,10 +285,10 @@ public OpcUaServer build() throws Exception {
StatusCodes.Bad_ConfigurationError,
"Certificate is missing the application URI"));
- final Set endpointConfigurations =
- createEndpointConfigurations(certificate, tcpBindPort, httpsBindPort, endpointHostnames);
+ final Set endpointConfigurations =
+ createEndpointConfigurations(certificate, tcpBindPort, endpointHostnames);
- serverConfig =
+ final OpcUaServerConfig serverConfig =
OpcUaServerConfig.builder()
.setApplicationUri(applicationUri)
.setApplicationName(LocalizedText.english("Apache IoTDB OPC UA server"))
@@ -278,16 +302,18 @@ public OpcUaServer build() throws Exception {
"",
DateTime.now()))
.setCertificateManager(certificateManager)
- .setTrustListManager(trustListManager)
- .setCertificateValidator(certificateValidator)
- .setHttpsKeyPair(httpsKeyPair)
- .setHttpsCertificateChain(new X509Certificate[] {httpsCertificate})
- .setIdentityValidator(new CompositeValidator(identityValidator, x509IdentityValidator))
+ .setIdentityValidator(new CompositeValidator(identityValidators))
+ .setSessionSecurityDiagnosticsAccessMode(
+ SessionSecurityDiagnosticsAccessMode.RESTRICTED)
.setProductUri("urn:apache:iotdb:opc-ua-server")
.build();
// Setup server to enable event posting
- final OpcUaServer server = new OpcUaServer(serverConfig);
+ final OpcTcpServerTransportConfig transportConfig =
+ OpcTcpServerTransportConfig.newBuilder().build();
+ final OpcUaServer server =
+ new OpcUaServer(
+ serverConfig, transportProfile -> new OpcTcpServerTransport(transportConfig));
final UaNode serverNode =
server.getAddressSpaceManager().getManagedNode(Identifiers.Server).orElse(null);
if (serverNode instanceof ServerTypeNode) {
@@ -341,12 +367,9 @@ private boolean isAdvertisedHostInCertificate(final X509Certificate certificate)
.anyMatch(hostname -> hostname.equalsIgnoreCase(advertisedHost));
}
- Set createEndpointConfigurations(
- final X509Certificate certificate,
- final int tcpBindPort,
- final int httpsBindPort,
- final Set hostnames) {
- final Set endpointConfigurations = new LinkedHashSet<>();
+ Set createEndpointConfigurations(
+ final X509Certificate certificate, final int tcpBindPort, final Set hostnames) {
+ final Set endpointConfigurations = new LinkedHashSet<>();
final Set effectiveHostnames = new LinkedHashSet<>();
if (Objects.nonNull(advertisedHost)) {
effectiveHostnames.add(toEndpointHostname(advertisedHost));
@@ -356,84 +379,61 @@ Set createEndpointConfigurations(
.forEach(effectiveHostnames::add);
}
- final List bindAddresses = newArrayList();
- bindAddresses.add(WILD_CARD_ADDRESS);
-
- for (final String bindAddress : bindAddresses) {
- for (final String hostname : effectiveHostnames) {
- final EndpointConfiguration.Builder builder =
- EndpointConfiguration.newBuilder()
- .setBindAddress(bindAddress)
- .setHostname(hostname)
- .setPath("/iotdb")
- .setCertificate(certificate)
- .addTokenPolicies(
- USER_TOKEN_POLICY_ANONYMOUS,
- USER_TOKEN_POLICY_USERNAME,
- USER_TOKEN_POLICY_X509);
-
- final Set securityPolicySet = new HashSet<>(securityPolicies);
- if (securityPolicySet.contains(SecurityPolicy.None)) {
- final EndpointConfiguration.Builder noSecurityBuilder =
- builder
- .copy()
- .setSecurityPolicy(SecurityPolicy.None)
- .setSecurityMode(MessageSecurityMode.None);
-
- endpointConfigurations.add(buildTcpEndpoint(noSecurityBuilder, tcpBindPort));
- endpointConfigurations.add(buildHttpsEndpoint(noSecurityBuilder, httpsBindPort));
- securityPolicySet.remove(SecurityPolicy.None);
- }
-
- for (final SecurityPolicy securityPolicy : securityPolicySet) {
- endpointConfigurations.add(
- buildTcpEndpoint(
- builder
- .copy()
- .setSecurityPolicy(securityPolicy)
- .setSecurityMode(MessageSecurityMode.SignAndEncrypt),
- tcpBindPort));
-
- endpointConfigurations.add(
- buildHttpsEndpoint(
- builder
- .copy()
- .setSecurityPolicy(securityPolicy)
- .setSecurityMode(MessageSecurityMode.Sign),
- httpsBindPort));
- }
-
- final EndpointConfiguration.Builder discoveryBuilder =
+ for (final String hostname : effectiveHostnames) {
+ final EndpointConfig.Builder builder =
+ EndpointConfig.newBuilder()
+ .setBindAddress(WILD_CARD_ADDRESS)
+ .setHostname(hostname)
+ .setPath("/iotdb")
+ .setCertificate(certificate);
+ if (enableAnonymousAccess) {
+ builder.addTokenPolicy(USER_TOKEN_POLICY_ANONYMOUS);
+ }
+ builder.addTokenPolicies(USER_TOKEN_POLICY_USERNAME, USER_TOKEN_POLICY_X509);
+
+ final Set securityPolicySet = new HashSet<>(securityPolicies);
+ if (securityPolicySet.contains(SecurityPolicy.None)) {
+ final EndpointConfig.Builder noSecurityBuilder =
builder
.copy()
- .setPath("/iotdb/discovery")
.setSecurityPolicy(SecurityPolicy.None)
.setSecurityMode(MessageSecurityMode.None);
- endpointConfigurations.add(buildTcpEndpoint(discoveryBuilder, tcpBindPort));
- endpointConfigurations.add(buildHttpsEndpoint(discoveryBuilder, httpsBindPort));
+ endpointConfigurations.add(buildTcpEndpoint(noSecurityBuilder, tcpBindPort));
+ securityPolicySet.remove(SecurityPolicy.None);
}
+
+ for (final SecurityPolicy securityPolicy : securityPolicySet) {
+ endpointConfigurations.add(
+ buildTcpEndpoint(
+ builder
+ .copy()
+ .setSecurityPolicy(securityPolicy)
+ .setSecurityMode(MessageSecurityMode.SignAndEncrypt),
+ tcpBindPort));
+ }
+
+ final EndpointConfig.Builder discoveryBuilder =
+ builder
+ .copy()
+ .setPath("/iotdb/discovery")
+ .setSecurityPolicy(SecurityPolicy.None)
+ .setSecurityMode(MessageSecurityMode.None);
+
+ endpointConfigurations.add(buildTcpEndpoint(discoveryBuilder, tcpBindPort));
}
return endpointConfigurations;
}
- private EndpointConfiguration buildTcpEndpoint(
- final EndpointConfiguration.Builder base, final int tcpBindPort) {
+ private EndpointConfig buildTcpEndpoint(
+ final EndpointConfig.Builder base, final int tcpBindPort) {
return base.copy()
.setTransportProfile(TransportProfile.TCP_UASC_UABINARY)
.setBindPort(tcpBindPort)
.build();
}
- private EndpointConfiguration buildHttpsEndpoint(
- final EndpointConfiguration.Builder base, final int httpsBindPort) {
- return base.copy()
- .setTransportProfile(TransportProfile.HTTPS_UABINARY)
- .setBindPort(httpsBindPort)
- .build();
- }
-
/////////////////////////////// Conflict detection ///////////////////////////////
void checkEquals(
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
index c760f661a711..9ccffbfc4507 100644
--- 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
@@ -21,7 +21,7 @@
import org.apache.iotdb.pipe.api.exception.PipeException;
-import org.eclipse.milo.opcua.sdk.client.api.identity.AnonymousProvider;
+import org.eclipse.milo.opcua.sdk.client.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;
diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/client/IoTDBOpcUaClientTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/client/IoTDBOpcUaClientTest.java
index 98f8cfaa88d3..e4ab28692c8f 100644
--- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/client/IoTDBOpcUaClientTest.java
+++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/client/IoTDBOpcUaClientTest.java
@@ -26,8 +26,7 @@
import org.apache.tsfile.write.record.Tablet;
import org.apache.tsfile.write.schema.MeasurementSchema;
import org.eclipse.milo.opcua.sdk.client.OpcUaClient;
-import org.eclipse.milo.opcua.sdk.client.api.UaClient;
-import org.eclipse.milo.opcua.sdk.client.api.identity.AnonymousProvider;
+import org.eclipse.milo.opcua.sdk.client.identity.AnonymousProvider;
import org.eclipse.milo.opcua.stack.core.StatusCodes;
import org.eclipse.milo.opcua.stack.core.security.SecurityPolicy;
import org.eclipse.milo.opcua.stack.core.types.builtin.NodeId;
@@ -51,7 +50,7 @@ public class IoTDBOpcUaClientTest {
@Test
public void testTransferWritesAllMeasurementsInOneRequest() throws Exception {
final OpcUaClient miloClient = Mockito.mock(OpcUaClient.class);
- Mockito.when(miloClient.writeValues(Mockito.anyList(), Mockito.anyList()))
+ Mockito.when(miloClient.writeValuesAsync(Mockito.anyList(), Mockito.anyList()))
.thenReturn(
CompletableFuture.completedFuture(Arrays.asList(StatusCode.GOOD, StatusCode.GOOD)));
final IoTDBOpcUaClient client = createClient(miloClient);
@@ -59,7 +58,7 @@ public void testTransferWritesAllMeasurementsInOneRequest() throws Exception {
client.transfer(createTablet(), createSink());
Mockito.verify(miloClient)
- .writeValues(
+ .writeValuesAsync(
Mockito.argThat(nodeIds("root/db/d1/s1", "root/db/d1/s2")),
Mockito.argThat(listWithSize(2)));
}
@@ -67,7 +66,7 @@ public void testTransferWritesAllMeasurementsInOneRequest() throws Exception {
@Test
public void testTransferCreatesAndRetriesOnlyMissingNodes() throws Exception {
final OpcUaClient miloClient = Mockito.mock(OpcUaClient.class);
- Mockito.when(miloClient.writeValues(Mockito.anyList(), Mockito.anyList()))
+ Mockito.when(miloClient.writeValuesAsync(Mockito.anyList(), Mockito.anyList()))
.thenReturn(
CompletableFuture.completedFuture(
Arrays.asList(new StatusCode(StatusCodes.Bad_NodeIdUnknown), StatusCode.GOOD)))
@@ -77,7 +76,7 @@ public void testTransferCreatesAndRetriesOnlyMissingNodes() throws Exception {
final AddNodesResult addNodesResult = Mockito.mock(AddNodesResult.class);
Mockito.when(addNodesResult.getStatusCode()).thenReturn(StatusCode.GOOD);
Mockito.when(addNodesResponse.getResults()).thenReturn(new AddNodesResult[] {addNodesResult});
- Mockito.when(miloClient.addNodes(Mockito.anyList()))
+ Mockito.when(miloClient.addNodesAsync(Mockito.anyList()))
.thenReturn(CompletableFuture.completedFuture(addNodesResponse));
final IoTDBOpcUaClient client = Mockito.spy(createClient(miloClient));
@@ -95,17 +94,18 @@ public void testTransferCreatesAndRetriesOnlyMissingNodes() throws Exception {
final InOrder inOrder = Mockito.inOrder(miloClient);
inOrder
.verify(miloClient)
- .writeValues(Mockito.argThat(listWithSize(2)), Mockito.argThat(listWithSize(2)));
- inOrder.verify(miloClient).addNodes(Mockito.argThat(listWithSize(1)));
+ .writeValuesAsync(Mockito.argThat(listWithSize(2)), Mockito.argThat(listWithSize(2)));
+ inOrder.verify(miloClient).addNodesAsync(Mockito.argThat(listWithSize(1)));
inOrder
.verify(miloClient)
- .writeValues(Mockito.argThat(nodeIds("root/db/d1/s1")), Mockito.argThat(listWithSize(1)));
+ .writeValuesAsync(
+ Mockito.argThat(nodeIds("root/db/d1/s1")), Mockito.argThat(listWithSize(1)));
}
@Test
public void testTransferFailsOnNonRecoverableStatus() throws Exception {
final OpcUaClient miloClient = Mockito.mock(OpcUaClient.class);
- Mockito.when(miloClient.writeValues(Mockito.anyList(), Mockito.anyList()))
+ Mockito.when(miloClient.writeValuesAsync(Mockito.anyList(), Mockito.anyList()))
.thenReturn(
CompletableFuture.completedFuture(
Arrays.asList(new StatusCode(StatusCodes.Bad_NotWritable), StatusCode.GOOD)));
@@ -119,7 +119,7 @@ public void testTransferFailsOnNonRecoverableStatus() throws Exception {
Assert.assertTrue(e.getMessage().contains("Bad_NotWritable"));
}
- Mockito.verify(miloClient, Mockito.never()).addNodes(Mockito.anyList());
+ Mockito.verify(miloClient, Mockito.never()).addNodesAsync(Mockito.anyList());
}
private static IoTDBOpcUaClient createClient(final OpcUaClient miloClient) throws Exception {
@@ -129,8 +129,9 @@ private static IoTDBOpcUaClient createClient(final OpcUaClient miloClient) throw
final ClientRunner runner = Mockito.mock(ClientRunner.class);
Mockito.when(runner.getTimeoutSeconds()).thenReturn(1L);
client.setRunner(runner);
- final CompletableFuture connectFuture = CompletableFuture.completedFuture(miloClient);
- Mockito.when(miloClient.connect()).thenReturn(connectFuture);
+ final CompletableFuture connectFuture =
+ CompletableFuture.completedFuture(miloClient);
+ Mockito.when(miloClient.connectAsync()).thenReturn(connectFuture);
client.run(miloClient);
return client;
}
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
index 97938f56d09f..8c21c6632ce8 100644
--- 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
@@ -21,10 +21,13 @@
import org.apache.iotdb.pipe.api.exception.PipeException;
+import org.eclipse.milo.opcua.sdk.server.EndpointConfig;
import org.eclipse.milo.opcua.sdk.server.OpcUaServer;
+import org.eclipse.milo.opcua.sdk.server.diagnostics.SessionSecurityDiagnosticsAccessMode;
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.enumerated.UserTokenType;
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;
@@ -49,18 +52,24 @@ public void testDetectedHostsArePublishedByDefault() {
final OpcUaServerBuilder builder =
new OpcUaServerBuilder().setSecurityPolicies(securityPolicies);
- final Set endpoints =
- builder.createEndpointConfigurations(null, 12686, 8443, detectedHostnames);
+ final Set endpoints =
+ builder.createEndpointConfigurations(null, 12686, detectedHostnames);
for (final String hostname : detectedHostnames) {
Assert.assertEquals(
- 2,
+ 1,
endpoints.stream()
.filter(endpoint -> endpoint.getPath().equals("/iotdb"))
.filter(endpoint -> endpoint.getHostname().equals(hostname))
.filter(endpoint -> endpoint.getSecurityPolicy() == SecurityPolicy.None)
.count());
}
+ Assert.assertTrue(
+ endpoints.stream()
+ .allMatch(
+ endpoint ->
+ endpoint.getTransportProfile() == TransportProfile.TCP_UASC_UABINARY
+ && endpoint.getBindPort() == 12686));
Assert.assertEquals(Collections.singleton(SecurityPolicy.None), securityPolicies);
}
@@ -73,12 +82,12 @@ public void testOnlyExplicitAdvertisedHostIsPublished() {
.setAdvertisedHost("opc.example.com")
.setSecurityPolicies(Collections.singleton(SecurityPolicy.None));
- final Set endpoints =
- builder.createEndpointConfigurations(null, 12686, 8443, detectedHostnames);
+ final Set endpoints =
+ builder.createEndpointConfigurations(null, 12686, detectedHostnames);
Assert.assertEquals(
Collections.singleton("opc.example.com"),
- endpoints.stream().map(EndpointConfiguration::getHostname).collect(Collectors.toSet()));
+ endpoints.stream().map(EndpointConfig::getHostname).collect(Collectors.toSet()));
Assert.assertTrue(
endpoints.stream().allMatch(endpoint -> "0.0.0.0".equals(endpoint.getBindAddress())));
}
@@ -90,16 +99,15 @@ public void testIpv6AdvertisedHostIsNormalizedAndEndpointUrlIsRejected() {
.setAdvertisedHost("[2001:db8::1]")
.setSecurityPolicies(Collections.singleton(SecurityPolicy.None));
- final Set endpoints =
- builder.createEndpointConfigurations(
- null, 12686, 8443, Collections.singleton("opc-server"));
+ final Set endpoints =
+ builder.createEndpointConfigurations(null, 12686, Collections.singleton("opc-server"));
Assert.assertEquals(
Collections.singleton("[2001:db8::1]"),
- endpoints.stream().map(EndpointConfiguration::getHostname).collect(Collectors.toSet()));
+ endpoints.stream().map(EndpointConfig::getHostname).collect(Collectors.toSet()));
Assert.assertTrue(
endpoints.stream()
- .map(EndpointConfiguration::getEndpointUrl)
+ .map(EndpointConfig::getEndpointUrl)
.allMatch(endpointUrl -> endpointUrl.contains("://[2001:db8::1]:")));
Assert.assertThrows(
IllegalArgumentException.class,
@@ -123,20 +131,86 @@ public void testNewCertificateContainsAdvertisedHost() throws Exception {
.setSecurityPolicies(Collections.singleton(SecurityPolicy.None))
.setDebounceTimeMs(50)) {
final OpcUaServer server = builder.build();
- final Set endpoints = server.getConfig().getEndpoints();
+ final Set endpoints = server.getConfig().getEndpoints();
Assert.assertEquals(
Collections.singleton(advertisedHost),
- endpoints.stream().map(EndpointConfiguration::getHostname).collect(Collectors.toSet()));
+ endpoints.stream().map(EndpointConfig::getHostname).collect(Collectors.toSet()));
Assert.assertTrue(
endpoints.stream()
- .map(EndpointConfiguration::getCertificate)
+ .map(EndpointConfig::getCertificate)
.allMatch(
certificate ->
CertificateUtil.getSanDnsNames(certificate).contains(advertisedHost)));
}
}
+ @Test
+ public void testRebuildWithChangedPassword() throws Exception {
+ final Path securityDir = temporaryFolder.newFolder("changed-password-security").toPath();
+
+ try (final OpcUaServerBuilder builder =
+ new OpcUaServerBuilder()
+ .setTcpBindPort(12686)
+ .setHttpsBindPort(8443)
+ .setAdvertisedHost("127.0.0.1")
+ .setUser("root")
+ .setPassword("root")
+ .setSecurityDir(securityDir.toString())
+ .setEnableAnonymousAccess(true)
+ .setSecurityPolicies(Collections.singleton(SecurityPolicy.None))
+ .setDebounceTimeMs(50)) {
+ builder.build();
+ }
+
+ try (final OpcUaServerBuilder builder =
+ new OpcUaServerBuilder()
+ .setTcpBindPort(12686)
+ .setHttpsBindPort(8443)
+ .setAdvertisedHost("127.0.0.1")
+ .setUser("root")
+ .setPassword("changed")
+ .setSecurityDir(securityDir.toString())
+ .setEnableAnonymousAccess(true)
+ .setSecurityPolicies(Collections.singleton(SecurityPolicy.None))
+ .setDebounceTimeMs(50)) {
+ builder.build();
+ }
+ }
+
+ @Test
+ public void testAnonymousAccessCanBeDisabledAndDiagnosticsStayRestricted() throws Exception {
+ final Path securityDir = temporaryFolder.newFolder("restricted-security").toPath();
+
+ try (final OpcUaServerBuilder builder =
+ new OpcUaServerBuilder()
+ .setTcpBindPort(12686)
+ .setHttpsBindPort(8443)
+ .setAdvertisedHost("127.0.0.1")
+ .setUser("root")
+ .setPassword("root")
+ .setSecurityDir(securityDir.toString())
+ .setEnableAnonymousAccess(false)
+ .setSecurityPolicies(Collections.singleton(SecurityPolicy.None))
+ .setDebounceTimeMs(50)) {
+ final OpcUaServer server = builder.build();
+
+ Assert.assertEquals(
+ SessionSecurityDiagnosticsAccessMode.RESTRICTED,
+ server.getConfig().getSessionSecurityDiagnosticsAccessMode());
+ Assert.assertFalse(
+ server
+ .getConfig()
+ .getIdentityValidator()
+ .getSupportedTokenTypes()
+ .contains(UserTokenType.Anonymous));
+ Assert.assertTrue(
+ server.getConfig().getEndpoints().stream()
+ .flatMap(endpoint -> endpoint.getTokenPolicies().stream())
+ .noneMatch(policy -> policy.getTokenType() == UserTokenType.Anonymous));
+ }
+ }
+
@Test
public void testAdvertisedHostParticipatesInConflictDetection() {
final Path securityDir = temporaryFolder.getRoot().toPath();
diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaTcpNoneCompatibilityTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaTcpNoneCompatibilityTest.java
new file mode 100644
index 000000000000..fabee999e2ed
--- /dev/null
+++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/sink/protocol/opcua/server/OpcUaTcpNoneCompatibilityTest.java
@@ -0,0 +1,132 @@
+/*
+ * 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.eclipse.milo.opcua.sdk.client.OpcUaClient;
+import org.eclipse.milo.opcua.sdk.client.identity.AnonymousProvider;
+import org.eclipse.milo.opcua.sdk.client.identity.IdentityProvider;
+import org.eclipse.milo.opcua.sdk.client.identity.UsernameProvider;
+import org.eclipse.milo.opcua.sdk.server.OpcUaServer;
+import org.eclipse.milo.opcua.stack.core.NodeIds;
+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.DataValue;
+import org.eclipse.milo.opcua.stack.core.types.enumerated.MessageSecurityMode;
+import org.eclipse.milo.opcua.stack.core.types.enumerated.TimestampsToReturn;
+import org.eclipse.milo.opcua.stack.core.types.structured.EndpointDescription;
+import org.junit.Assert;
+import org.junit.Rule;
+import org.junit.Test;
+import org.junit.rules.TemporaryFolder;
+
+import java.net.ServerSocket;
+import java.nio.file.Path;
+import java.util.Collections;
+import java.util.List;
+import java.util.Optional;
+import java.util.concurrent.TimeUnit;
+
+import static org.eclipse.milo.opcua.stack.core.types.builtin.unsigned.Unsigned.uint;
+
+public class OpcUaTcpNoneCompatibilityTest {
+
+ private static final long TIMEOUT_SECONDS = 15;
+
+ @Rule public final TemporaryFolder temporaryFolder = new TemporaryFolder();
+
+ @Test
+ public void testTcpNoneSupportsAnonymousAndUsernameSessions() throws Exception {
+ final int tcpBindPort = findAvailablePort();
+ final Path securityDir = temporaryFolder.newFolder("tcp-none-security").toPath();
+ OpcUaServer server = null;
+
+ try (final OpcUaServerBuilder builder =
+ new OpcUaServerBuilder()
+ .setTcpBindPort(tcpBindPort)
+ .setHttpsBindPort(tcpBindPort == 65535 ? 65534 : tcpBindPort + 1)
+ .setAdvertisedHost("127.0.0.1")
+ .setUser("root")
+ .setPassword("root")
+ .setSecurityDir(securityDir.toString())
+ .setEnableAnonymousAccess(true)
+ .setSecurityPolicies(Collections.singleton(SecurityPolicy.None))
+ .setDebounceTimeMs(50)) {
+ server = builder.build();
+ server.startup().get(TIMEOUT_SECONDS, TimeUnit.SECONDS);
+
+ final String endpointUrl = "opc.tcp://127.0.0.1:" + tcpBindPort + "/iotdb";
+ assertCanReadServerState(endpointUrl, AnonymousProvider.INSTANCE);
+ assertCanReadServerState(endpointUrl, new UsernameProvider("root", "root"));
+ } finally {
+ if (server != null) {
+ server.shutdown().get(TIMEOUT_SECONDS, TimeUnit.SECONDS);
+ }
+ Stack.releaseSharedResources();
+ }
+ }
+
+ private static void assertCanReadServerState(
+ final String endpointUrl, final IdentityProvider identityProvider) throws Exception {
+ final OpcUaClient client =
+ OpcUaClient.create(
+ endpointUrl,
+ OpcUaTcpNoneCompatibilityTest::selectTcpNoneEndpoint,
+ transportBuilder -> transportBuilder.setConnectTimeout(uint(TIMEOUT_SECONDS * 1000L)),
+ configBuilder ->
+ configBuilder
+ .setIdentityProvider(identityProvider)
+ .setRequestTimeout(uint(TIMEOUT_SECONDS * 1000L)));
+
+ try {
+ client.connectAsync().get(TIMEOUT_SECONDS, TimeUnit.SECONDS);
+
+ final EndpointDescription selectedEndpoint = client.getConfig().getEndpoint();
+ Assert.assertEquals(MessageSecurityMode.None, selectedEndpoint.getSecurityMode());
+ Assert.assertEquals(SecurityPolicy.None.getUri(), selectedEndpoint.getSecurityPolicyUri());
+ Assert.assertEquals(
+ Stack.TCP_UASC_UABINARY_TRANSPORT_URI, selectedEndpoint.getTransportProfileUri());
+
+ final DataValue serverState =
+ client.readValue(0.0, TimestampsToReturn.Neither, NodeIds.Server_ServerStatus_State);
+ Assert.assertNotNull(serverState.getValue().getValue());
+ Assert.assertFalse(serverState.getStatusCode().isBad());
+ } finally {
+ client.disconnectAsync().get(TIMEOUT_SECONDS, TimeUnit.SECONDS);
+ }
+ }
+
+ private static Optional selectTcpNoneEndpoint(
+ final List endpoints) {
+ return endpoints.stream()
+ .filter(endpoint -> endpoint.getEndpointUrl().endsWith("/iotdb"))
+ .filter(endpoint -> endpoint.getSecurityMode() == MessageSecurityMode.None)
+ .filter(endpoint -> SecurityPolicy.None.getUri().equals(endpoint.getSecurityPolicyUri()))
+ .filter(
+ endpoint ->
+ Stack.TCP_UASC_UABINARY_TRANSPORT_URI.equals(endpoint.getTransportProfileUri()))
+ .findFirst();
+ }
+
+ private static int findAvailablePort() throws Exception {
+ try (final ServerSocket socket = new ServerSocket(0)) {
+ return socket.getLocalPort();
+ }
+ }
+}
diff --git a/pom.xml b/pom.xml
index 11ffc5335b29..5107e2a9a656 100644
--- a/pom.xml
+++ b/pom.xml
@@ -120,7 +120,7 @@
1.8
1.8
1.11.4
- 0.6.14
+ 1.1.6
2.23.4
@@ -371,44 +371,27 @@
org.eclipse.milo
- stack-core
+ milo-stack-core
${milo.version}
-
-
- com.sun.activation
- jakarta.activation
-
-
- org.eclipse.milo
- sdk-core
- ${milo.version}
-
-
- com.sun.activation
- jakarta.activation
-
-
+ org.checkerframework
+ checker-qual
+ ${checker-qual.version}
org.eclipse.milo
- stack-server
+ milo-sdk-core
${milo.version}
-
- org.checkerframework
- checker-qual
- ${checker-qual.version}
-
org.eclipse.milo
- stack-client
+ milo-transport
${milo.version}
org.eclipse.milo
- sdk-client
+ milo-sdk-client
${milo.version}
@@ -430,14 +413,8 @@
org.eclipse.milo
- sdk-server
+ milo-sdk-server
${milo.version}
-
-
- com.sun.activation
- jakarta.activation
-
-
org.reflections