This is an automated email from the ASF dual-hosted git repository.
jt2594838 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/master by this push:
new 02f43e73f47 Upgrade Eclipse Milo to 1.1.6 (#18431)
02f43e73f47 is described below
commit 02f43e73f4705186210c2ec9d7e07bb16827b2a5
Author: Caideyipi <[email protected]>
AuthorDate: Tue Aug 11 14:15:58 2026 +0800
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
---
dependencies.json | 17 +-
example/pipe-opc-ua-sink/pom.xml | 4 +-
.../java/org/apache/iotdb/opcua/ClientExample.java | 4 +-
.../apache/iotdb/opcua/ClientExampleRunner.java | 41 +++-
.../java/org/apache/iotdb/opcua/ClientTest.java | 54 ++---
.../iotdb/pipe/it/single/IoTDBPipeOPCUAIT.java | 33 +--
iotdb-core/datanode/pom.xml | 14 +-
.../db/pipe/sink/protocol/opcua/OpcUaSink.java | 6 +-
.../sink/protocol/opcua/client/ClientRunner.java | 46 +++-
.../protocol/opcua/client/IoTDBOpcUaClient.java | 24 +-
.../sink/protocol/opcua/server/OpcUaNameSpace.java | 10 +-
.../protocol/opcua/server/OpcUaServerBuilder.java | 246 ++++++++++-----------
.../protocol/opcua/client/ClientRunnerTest.java | 2 +-
.../opcua/client/IoTDBOpcUaClientTest.java | 27 +--
.../opcua/server/OpcUaServerBuilderTest.java | 104 +++++++--
.../server/OpcUaTcpNoneCompatibilityTest.java | 132 +++++++++++
pom.xml | 35 +--
17 files changed, 492 insertions(+), 307 deletions(-)
diff --git a/dependencies.json b/dependencies.json
index 4bee8c84542..ac20a3f7586 100644
--- a/dependencies.json
+++ b/dependencies.json
@@ -28,7 +28,6 @@
"com.google.j2objc:j2objc-annotations",
"com.h2database:h2-mvstore",
"com.sun.activation:jakarta.activation",
- "com.sun.istack:istack-commons-runtime",
"com.zaxxer:HikariCP",
"commons-cli:commons-cli",
"commons-codec:commons-codec",
@@ -116,14 +115,11 @@
"org.eclipse.jetty:jetty-session",
"org.eclipse.jetty:jetty-util",
"org.eclipse.jetty.ee10:jetty-ee10-servlet",
- "org.eclipse.milo:bsd-core",
- "org.eclipse.milo:bsd-generator",
- "org.eclipse.milo:sdk-client",
- "org.eclipse.milo:sdk-core",
- "org.eclipse.milo:sdk-server",
- "org.eclipse.milo:stack-client",
- "org.eclipse.milo:stack-core",
- "org.eclipse.milo:stack-server",
+ "org.eclipse.milo:milo-sdk-client",
+ "org.eclipse.milo:milo-sdk-core",
+ "org.eclipse.milo:milo-sdk-server",
+ "org.eclipse.milo:milo-stack-core",
+ "org.eclipse.milo:milo-transport",
"org.fusesource.hawtbuf:hawtbuf",
"org.fusesource.hawtdispatch:hawtdispatch",
"org.fusesource.hawtdispatch:hawtdispatch-transport",
@@ -133,8 +129,6 @@
"org.glassfish.hk2:hk2-utils",
"org.glassfish.hk2:osgi-resource-locator",
"org.glassfish.hk2.external:aopalliance-repackaged",
- "org.glassfish.jaxb:jaxb-runtime",
- "org.glassfish.jaxb:txw2",
"org.glassfish.jersey.containers:jersey-container-servlet-core",
"org.glassfish.jersey.core:jersey-client",
"org.glassfish.jersey.core:jersey-common",
@@ -145,6 +139,7 @@
"org.java-websocket:Java-WebSocket",
"org.javassist:javassist",
"org.jline:jline",
+ "org.jspecify:jspecify",
"org.jvnet.mimepull:mimepull",
"org.latencyutils:LatencyUtils",
"org.ops4j.pax.jdbc:pax-jdbc-common",
diff --git a/example/pipe-opc-ua-sink/pom.xml b/example/pipe-opc-ua-sink/pom.xml
index 916e5262c62..0627a854096 100644
--- a/example/pipe-opc-ua-sink/pom.xml
+++ b/example/pipe-opc-ua-sink/pom.xml
@@ -31,12 +31,12 @@
<dependencies>
<dependency>
<groupId>org.eclipse.milo</groupId>
- <artifactId>sdk-client</artifactId>
+ <artifactId>milo-sdk-client</artifactId>
<version>${milo.version}</version>
</dependency>
<dependency>
<groupId>org.eclipse.milo</groupId>
- <artifactId>sdk-server</artifactId>
+ <artifactId>milo-sdk-server</artifactId>
</dependency>
</dependencies>
<profiles>
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 6b7f6997763..dba324c755d 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 1fe49a75007..60b1f8b2ee6 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 @@ package org.apache.iotdb.opcua;
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<OpcUaClient> future = new
CompletableFuture<>();
private final ClientExample clientExample;
+ private FileBasedTrustListManager trustListManager;
public ClientExampleRunner(ClientExample clientExample) {
this.clientExample = clientExample;
@@ -61,21 +63,24 @@ public class ClientExampleRunner {
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 @@ public class ClientExampleRunner {
.setCertificateChain(loader.getClientCertificateChain())
.setCertificateValidator(certificateValidator)
.setIdentityProvider(clientExample.getIdentityProvider())
- .setRequestTimeout(uint(5000))
- .build());
+ .setRequestTimeout(uint(5000)));
}
public void run() {
@@ -100,11 +104,13 @@ public class ClientExampleRunner {
}
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 class ClientExampleRunner {
} catch (Throwable t) {
System.out.println("Error getting client: {}" + t.getMessage());
+ closeTrustListManager();
future.completeExceptionally(t);
try {
@@ -144,4 +151,16 @@ public class ClientExampleRunner {
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 cc09a7dc7bb..6db181f56c4 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 class ClientTest implements ClientExample {
new ClientExampleRunner(example).run();
}
- private final AtomicLong clientHandles = new AtomicLong(1L);
-
@Override
public void run(OpcUaClient client, CompletableFuture<OpcUaClient> 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 class ClientTest implements ClientExample {
},
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<UaMonitoredItem> 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 9dfa3bc012c..4054f7fae5f 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
@@ -35,9 +35,9 @@ import org.apache.iotdb.rpc.TSStatusCode;
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;
@@ -127,13 +127,10 @@ public class IoTDBPipeOPCUAIT extends
AbstractPipeSingleIT {
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;
}
@@ -194,26 +191,19 @@ public class IoTDBPipeOPCUAIT extends
AbstractPipeSingleIT {
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());
@@ -232,10 +222,7 @@ public class IoTDBPipeOPCUAIT extends AbstractPipeSingleIT
{
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());
@@ -247,7 +234,7 @@ public class IoTDBPipeOPCUAIT extends AbstractPipeSingleIT {
}
}
- 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 045c9f3c2d6..b615061fc74 100644
--- a/iotdb-core/datanode/pom.xml
+++ b/iotdb-core/datanode/pom.xml
@@ -187,23 +187,19 @@
</dependency>
<dependency>
<groupId>org.eclipse.milo</groupId>
- <artifactId>stack-core</artifactId>
+ <artifactId>milo-stack-core</artifactId>
</dependency>
<dependency>
<groupId>org.eclipse.milo</groupId>
- <artifactId>sdk-core</artifactId>
+ <artifactId>milo-sdk-core</artifactId>
</dependency>
<dependency>
<groupId>org.eclipse.milo</groupId>
- <artifactId>stack-server</artifactId>
+ <artifactId>milo-transport</artifactId>
</dependency>
<dependency>
<groupId>org.eclipse.milo</groupId>
- <artifactId>stack-client</artifactId>
- </dependency>
- <dependency>
- <groupId>org.eclipse.milo</groupId>
- <artifactId>sdk-client</artifactId>
+ <artifactId>milo-sdk-client</artifactId>
</dependency>
<dependency>
<groupId>org.bouncycastle</groupId>
@@ -223,7 +219,7 @@
</dependency>
<dependency>
<groupId>org.eclipse.milo</groupId>
- <artifactId>sdk-server</artifactId>
+ <artifactId>milo-sdk-server</artifactId>
</dependency>
<dependency>
<groupId>commons-cli</groupId>
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 2eb4e7e8136..1263603bd49 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
@@ -44,9 +44,9 @@ import org.apache.iotdb.pipe.api.exception.PipeException;
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 6cbc5c3d61d..64757596d16 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
@@ -24,8 +24,9 @@ import org.apache.iotdb.pipe.api.exception.PipeException;
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;
@@ -34,7 +35,8 @@ import org.eclipse.milo.opcua.stack.core.util.EndpointUtil;
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;
@@ -47,7 +49,7 @@ import java.util.Optional;
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);
@@ -61,6 +63,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;
@@ -96,18 +99,20 @@ public class ClientRunner {
throw new Exception(DataNodePipeMessages.UNABLE_TO_CREATE_SECURITY_DIR +
securityDir);
}
- final File pkiDir = securityDir.resolve("pki").toFile();
+ final Path pkiDir = securityDir.resolve("pki");
logger.info(DataNodePipeMessages.SECURITY_DIR,
securityDir.toAbsolutePath());
- logger.info(DataNodePipeMessages.SECURITY_PKI_DIR,
pkiDir.getAbsolutePath());
+ logger.info(DataNodePipeMessages.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(),
@@ -117,6 +122,7 @@ public class ClientRunner {
configurableUaClient.getNodeUrl(),
configurableUaClient.getSecurityPolicy(),
allowEndpointRedirect),
+ transportBuilder ->
transportBuilder.setConnectTimeout(uint(timeoutSeconds * 1000L)),
configBuilder ->
configBuilder
.setApplicationName(LocalizedText.english("Apache IoTDB OPC UA
client"))
@@ -127,9 +133,7 @@ public class ClientRunner {
.setCertificateValidator(certificateValidator)
.setIdentityProvider(configurableUaClient.getIdentityProvider())
.setRequestTimeout(uint(timeoutSeconds * 1000L))
- .setConnectTimeout(uint(timeoutSeconds * 1000L))
- .setMaxResponseMessageSize(uint(0))
- .build());
+ .setMaxResponseMessageSize(uint(0)));
}
static Optional<EndpointDescription> selectEndpoint(
@@ -207,6 +211,7 @@ public class ClientRunner {
e);
}
} catch (final Exception e) {
+ closeOnFailure(e);
throw new PipeException(
String.format(
DataNodePipeMessages.ERROR_GETTING_OPC_CLIENT_FMT,
@@ -216,6 +221,25 @@ public class ClientRunner {
}
}
+ 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 c6ebadedb89..8c6c30dfeba 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
@@ -31,7 +31,7 @@ import org.apache.tsfile.enums.TSDataType;
import org.apache.tsfile.write.record.Tablet;
import org.apache.tsfile.write.schema.IMeasurementSchema;
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;
@@ -104,7 +104,7 @@ public class IoTDBOpcUaClient {
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) {
@@ -234,7 +234,7 @@ public class IoTDBOpcUaClient {
}
}
- 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) {
@@ -267,7 +267,7 @@ public class IoTDBOpcUaClient {
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 {
@@ -341,7 +341,7 @@ public class IoTDBOpcUaClient {
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
@@ -356,7 +356,7 @@ public class IoTDBOpcUaClient {
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;
}
@@ -371,7 +371,7 @@ public class IoTDBOpcUaClient {
new QualifiedName(NAME_SPACE_INDEX, measurementName),
NodeClass.Variable,
ExtensionObject.encode(
- client.getStaticSerializationContext(),
+ client.getStaticEncodingContext(),
createMeasurementAttributes(measurementName, opcDataType,
initialValue)),
Identifiers.BaseDataVariableType.expanded()));
@@ -379,8 +379,14 @@ public class IoTDBOpcUaClient {
}
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 86c52149719..1d6262c2c6a 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
@@ -40,11 +40,11 @@ import org.apache.tsfile.write.schema.IMeasurementSchema;
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;
@@ -537,7 +537,7 @@ public class OpcUaNameSpace extends
ManagedNamespaceWithLifecycle {
}
// 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 687c1533519..5b7d820f8f7 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
@@ -23,19 +23,30 @@ 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.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;
@@ -43,15 +54,12 @@ import
org.eclipse.milo.opcua.stack.core.types.builtin.LocalizedText;
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;
@@ -59,16 +67,16 @@ import java.nio.file.Path;
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;
/**
@@ -89,7 +97,7 @@ public class OpcUaServerBuilder implements Closeable {
private Path securityDir;
private boolean enableAnonymousAccess;
private Set<SecurityPolicy> securityPolicies;
- private DefaultTrustListManager trustListManager;
+ private FileBasedTrustListManager trustListManager;
private long debounceTimeMs;
public OpcUaServerBuilder setTcpBindPort(final int tcpBindPort) {
@@ -191,59 +199,75 @@ public class OpcUaServerBuilder implements Closeable {
throw new PipeException(DataNodePipeMessages.UNABLE_CREATE_SECURITY_DIR
+ securityDir);
}
- final File pkiDir = securityDir.resolve("pki").toFile();
+ final Path pkiDir = securityDir.resolve("pki");
LoggerFactory.getLogger(OpcUaServerBuilder.class)
.info(DataNodePipeMessages.OPC_UA_SECURITY_DIR,
securityDir.toAbsolutePath());
LoggerFactory.getLogger(OpcUaServerBuilder.class)
- .info(DataNodePipeMessages.OPC_UA_SECURITY_PKI_DIR,
pkiDir.getAbsolutePath());
+ .info(DataNodePipeMessages.OPC_UA_SECURITY_PKI_DIR,
pkiDir.toAbsolutePath());
final Set<String> endpointHostnames = getEndpointHostnames();
final Set<String> 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(
DataNodePipeMessages.CERTIFICATE_DIRECTORY_IS_PLEASE_MOVE_CERTIFICATES_FROM,
- 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<IdentityValidator> 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(
@@ -265,10 +289,10 @@ public class OpcUaServerBuilder implements Closeable {
StatusCodes.Bad_ConfigurationError,
DataNodePipeMessages.CERTIFICATE_MISSING_APPLICATION_URI));
- final Set<EndpointConfiguration> endpointConfigurations =
- createEndpointConfigurations(certificate, tcpBindPort, httpsBindPort,
endpointHostnames);
+ final Set<EndpointConfig> endpointConfigurations =
+ createEndpointConfigurations(certificate, tcpBindPort,
endpointHostnames);
- serverConfig =
+ final OpcUaServerConfig serverConfig =
OpcUaServerConfig.builder()
.setApplicationUri(applicationUri)
.setApplicationName(LocalizedText.english("Apache IoTDB OPC UA
server"))
@@ -282,16 +306,18 @@ public class OpcUaServerBuilder implements Closeable {
"",
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) {
@@ -345,12 +371,9 @@ public class OpcUaServerBuilder implements Closeable {
.anyMatch(hostname -> hostname.equalsIgnoreCase(advertisedHost));
}
- Set<EndpointConfiguration> createEndpointConfigurations(
- final X509Certificate certificate,
- final int tcpBindPort,
- final int httpsBindPort,
- final Set<String> hostnames) {
- final Set<EndpointConfiguration> endpointConfigurations = new
LinkedHashSet<>();
+ Set<EndpointConfig> createEndpointConfigurations(
+ final X509Certificate certificate, final int tcpBindPort, final
Set<String> hostnames) {
+ final Set<EndpointConfig> endpointConfigurations = new LinkedHashSet<>();
final Set<String> effectiveHostnames = new LinkedHashSet<>();
if (Objects.nonNull(advertisedHost)) {
effectiveHostnames.add(toEndpointHostname(advertisedHost));
@@ -360,84 +383,61 @@ public class OpcUaServerBuilder implements Closeable {
.forEach(effectiveHostnames::add);
}
- final List<String> 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<SecurityPolicy> 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<SecurityPolicy> 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 c760f661a71..9ccffbfc450 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 @@ 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.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 9fc418766cb..5cb881939d4 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.enums.TSDataType;
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 class IoTDBOpcUaClientTest {
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 class IoTDBOpcUaClientTest {
@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 class IoTDBOpcUaClientTest {
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 class IoTDBOpcUaClientTest {
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 class IoTDBOpcUaClientTest {
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 @@ public class IoTDBOpcUaClientTest {
final ClientRunner runner = Mockito.mock(ClientRunner.class);
Mockito.when(runner.getTimeoutSeconds()).thenReturn(1L);
client.setRunner(runner);
- final CompletableFuture<UaClient> connectFuture =
CompletableFuture.completedFuture(miloClient);
- Mockito.when(miloClient.connect()).thenReturn(connectFuture);
+ final CompletableFuture<OpcUaClient> 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 97938f56d09..8c21c6632ce 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 @@ 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.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 class OpcUaServerBuilderTest {
final OpcUaServerBuilder builder =
new OpcUaServerBuilder().setSecurityPolicies(securityPolicies);
- final Set<EndpointConfiguration> endpoints =
- builder.createEndpointConfigurations(null, 12686, 8443,
detectedHostnames);
+ final Set<EndpointConfig> 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 class OpcUaServerBuilderTest {
.setAdvertisedHost("opc.example.com")
.setSecurityPolicies(Collections.singleton(SecurityPolicy.None));
- final Set<EndpointConfiguration> endpoints =
- builder.createEndpointConfigurations(null, 12686, 8443,
detectedHostnames);
+ final Set<EndpointConfig> 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 class OpcUaServerBuilderTest {
.setAdvertisedHost("[2001:db8::1]")
.setSecurityPolicies(Collections.singleton(SecurityPolicy.None));
- final Set<EndpointConfiguration> endpoints =
- builder.createEndpointConfigurations(
- null, 12686, 8443, Collections.singleton("opc-server"));
+ final Set<EndpointConfig> 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 class OpcUaServerBuilderTest {
.setSecurityPolicies(Collections.singleton(SecurityPolicy.None))
.setDebounceTimeMs(50)) {
final OpcUaServer server = builder.build();
- final Set<EndpointConfiguration> endpoints =
server.getConfig().getEndpoints();
+ final Set<EndpointConfig> 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 00000000000..fabee999e2e
--- /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<EndpointDescription> selectTcpNoneEndpoint(
+ final List<EndpointDescription> 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 43ec486debb..b4ab0ad3133 100644
--- a/pom.xml
+++ b/pom.xml
@@ -109,7 +109,7 @@
<maven.compiler.source>17</maven.compiler.source>
<maven.compiler.target>17</maven.compiler.target>
<micrometer.version>1.11.4</micrometer.version>
- <milo.version>0.6.14</milo.version>
+ <milo.version>1.1.6</milo.version>
<!-- It seems that powermock is having issues with the newest mockito
versions -->
<mockito.version>2.23.4</mockito.version>
<!--mockito.version>4.11.0</mockito.version-->
@@ -383,39 +383,22 @@
</dependency>
<dependency>
<groupId>org.eclipse.milo</groupId>
- <artifactId>stack-core</artifactId>
+ <artifactId>milo-stack-core</artifactId>
<version>${milo.version}</version>
- <exclusions>
- <exclusion>
- <groupId>com.sun.activation</groupId>
- <artifactId>jakarta.activation</artifactId>
- </exclusion>
- </exclusions>
- </dependency>
- <dependency>
- <groupId>org.eclipse.milo</groupId>
- <artifactId>sdk-core</artifactId>
- <version>${milo.version}</version>
- <exclusions>
- <exclusion>
- <groupId>com.sun.activation</groupId>
- <artifactId>jakarta.activation</artifactId>
- </exclusion>
- </exclusions>
</dependency>
<dependency>
<groupId>org.eclipse.milo</groupId>
- <artifactId>stack-server</artifactId>
+ <artifactId>milo-sdk-core</artifactId>
<version>${milo.version}</version>
</dependency>
<dependency>
<groupId>org.eclipse.milo</groupId>
- <artifactId>stack-client</artifactId>
+ <artifactId>milo-transport</artifactId>
<version>${milo.version}</version>
</dependency>
<dependency>
<groupId>org.eclipse.milo</groupId>
- <artifactId>sdk-client</artifactId>
+ <artifactId>milo-sdk-client</artifactId>
<version>${milo.version}</version>
</dependency>
<!-- TODO: Deprecated: Use Airline 2 or Picocli instead -->
@@ -432,14 +415,8 @@
</dependency>
<dependency>
<groupId>org.eclipse.milo</groupId>
- <artifactId>sdk-server</artifactId>
+ <artifactId>milo-sdk-server</artifactId>
<version>${milo.version}</version>
- <exclusions>
- <exclusion>
- <groupId>com.sun.activation</groupId>
- <artifactId>jakarta.activation</artifactId>
- </exclusion>
- </exclusions>
</dependency>
<dependency>
<groupId>org.reflections</groupId>