From 70961a965d8257236ab586e636d25c49459c33a8 Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Tue, 18 Aug 2026 13:04:18 +0800 Subject: [PATCH 1/4] Pipe: respect OPC UA operation limits --- .../opcua/client/IoTDBOpcUaClient.java | 112 +++++++++++++++--- .../opcua/client/IoTDBOpcUaClientTest.java | 94 +++++++++++++++ 2 files changed, 190 insertions(+), 16 deletions(-) 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 83fc94e2ea2fd..238ac5afe31ff 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 @@ -64,6 +64,7 @@ import java.nio.file.Paths; import java.util.ArrayList; +import java.util.Arrays; import java.util.HashSet; import java.util.List; import java.util.Map; @@ -74,10 +75,14 @@ import static org.apache.iotdb.db.pipe.sink.protocol.opcua.server.OpcUaNameSpace.convertToOpcDataType; import static org.apache.iotdb.db.pipe.sink.protocol.opcua.server.OpcUaNameSpace.timestampToUtc; import static org.eclipse.milo.opcua.stack.core.StatusCodes.Bad_Timeout; +import static org.eclipse.milo.opcua.stack.core.types.enumerated.TimestampsToReturn.Neither; public class IoTDBOpcUaClient { private static final Logger LOGGER = LoggerFactory.getLogger(OpcUaNameSpace.class); + private static final int DEFAULT_MAX_NODES_PER_WRITE = 10_000; + private static final int DEFAULT_MAX_NODES_PER_NODE_MANAGEMENT = 250; + // Customized nodes private static final int NAME_SPACE_INDEX = 2; @@ -90,6 +95,8 @@ public class IoTDBOpcUaClient { private OpcUaClient client; private final boolean historizing; private ClientRunner runner; + private int maxNodesPerWrite = DEFAULT_MAX_NODES_PER_WRITE; + private int maxNodesPerNodeManagement = DEFAULT_MAX_NODES_PER_NODE_MANAGEMENT; public IoTDBOpcUaClient( final String nodeUrl, @@ -119,6 +126,62 @@ public void run(final OpcUaClient client) throws Exception { } break; } + updateOperationLimits(); + } + + private void updateOperationLimits() { + try { + final List operationLimits = + client + .readValuesAsync( + 0.0, + Neither, + Arrays.asList( + Identifiers.Server_ServerCapabilities_OperationLimits_MaxNodesPerWrite, + Identifiers + .Server_ServerCapabilities_OperationLimits_MaxNodesPerNodeManagement)) + .get(); + maxNodesPerWrite = getOperationLimit(operationLimits, 0, DEFAULT_MAX_NODES_PER_WRITE); + maxNodesPerNodeManagement = + getOperationLimit(operationLimits, 1, DEFAULT_MAX_NODES_PER_NODE_MANAGEMENT); + LOGGER.info( + "OPC UA server operation limits: maxNodesPerWrite={}, maxNodesPerNodeManagement={}", + maxNodesPerWrite, + maxNodesPerNodeManagement); + } catch (final InterruptedException e) { + Thread.currentThread().interrupt(); + LOGGER.warn( + "Interrupted while reading OPC UA server operation limits, use defaults: " + + "maxNodesPerWrite={}, maxNodesPerNodeManagement={}", + DEFAULT_MAX_NODES_PER_WRITE, + DEFAULT_MAX_NODES_PER_NODE_MANAGEMENT); + } catch (final Exception e) { + LOGGER.warn( + "Failed to read OPC UA server operation limits, use defaults: " + + "maxNodesPerWrite={}, maxNodesPerNodeManagement={}", + DEFAULT_MAX_NODES_PER_WRITE, + DEFAULT_MAX_NODES_PER_NODE_MANAGEMENT, + e); + } + } + + private static int getOperationLimit( + final List operationLimits, final int index, final int defaultValue) { + if (Objects.isNull(operationLimits) || operationLimits.size() <= index) { + return defaultValue; + } + + final DataValue dataValue = operationLimits.get(index); + if (Objects.isNull(dataValue) + || Objects.isNull(dataValue.getStatusCode()) + || !dataValue.getStatusCode().isGood() + || Objects.isNull(dataValue.getValue()) + || !(dataValue.getValue().getValue() instanceof Number)) { + return defaultValue; + } + + final long limit = ((Number) dataValue.getValue().getValue()).longValue(); + return limit == 0 ? Integer.MAX_VALUE : (int) Math.min(limit, Integer.MAX_VALUE); } // Only support tree model & client-server @@ -267,17 +330,23 @@ private void addMissingNodes(final List writeRequests) throws } } - 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) { - throw new PipeException( - DataNodePipeMessages.FAILED_TO_CREATE_NODES_AFTER_TRANSFER_DATA - + addStatus - + writeRequests - .get(0) - .getErrorString(new StatusCode(StatusCodes.Bad_NodeIdUnknown))); + for (int startIndex = 0; startIndex < nodesToAdd.size(); ) { + final int endIndex = + getBatchEndIndex(startIndex, nodesToAdd.size(), maxNodesPerNodeManagement); + final AddNodesResponse addStatus = + client.addNodesAsync(nodesToAdd.subList(startIndex, endIndex)).get(); + for (final AddNodesResult result : addStatus.getResults()) { + if (!result.getStatusCode().equals(StatusCode.GOOD) + && result.getStatusCode().getValue() != StatusCodes.Bad_NodeIdExists) { + throw new PipeException( + DataNodePipeMessages.FAILED_TO_CREATE_NODES_AFTER_TRANSFER_DATA + + addStatus + + writeRequests + .get(0) + .getErrorString(new StatusCode(StatusCodes.Bad_NodeIdUnknown))); + } } + startIndex = endIndex; } } @@ -294,13 +363,24 @@ private void validateRetriedWrites( private List writeValuesOnce(final List writeRequests) throws Exception { - final List nodeIds = new ArrayList<>(writeRequests.size()); - final List dataValues = new ArrayList<>(writeRequests.size()); - for (final OpcUaWriteRequest writeRequest : writeRequests) { - nodeIds.add(writeRequest.nodeId); - dataValues.add(writeRequest.dataValue); + final List writeStatuses = new ArrayList<>(writeRequests.size()); + for (int startIndex = 0; startIndex < writeRequests.size(); ) { + final int endIndex = getBatchEndIndex(startIndex, writeRequests.size(), maxNodesPerWrite); + final List nodeIds = new ArrayList<>(endIndex - startIndex); + final List dataValues = new ArrayList<>(endIndex - startIndex); + for (int i = startIndex; i < endIndex; ++i) { + nodeIds.add(writeRequests.get(i).nodeId); + dataValues.add(writeRequests.get(i).dataValue); + } + writeStatuses.addAll(client.writeValuesAsync(nodeIds, dataValues).get()); + startIndex = endIndex; } - return client.writeValuesAsync(nodeIds, dataValues).get(); + return writeStatuses; + } + + private static int getBatchEndIndex( + final int startIndex, final int totalSize, final int batchSize) { + return (int) Math.min((long) totalSize, (long) startIndex + batchSize); } private static final class OpcUaWriteRequest { 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 8f8333143e315..d327155f32d00 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 @@ -34,8 +34,12 @@ 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.DataValue; +import org.eclipse.milo.opcua.stack.core.types.builtin.ExpandedNodeId; import org.eclipse.milo.opcua.stack.core.types.builtin.NodeId; import org.eclipse.milo.opcua.stack.core.types.builtin.StatusCode; +import org.eclipse.milo.opcua.stack.core.types.builtin.Variant; +import org.eclipse.milo.opcua.stack.core.types.enumerated.TimestampsToReturn; import org.eclipse.milo.opcua.stack.core.types.structured.AddNodesItem; import org.eclipse.milo.opcua.stack.core.types.structured.AddNodesResponse; import org.eclipse.milo.opcua.stack.core.types.structured.AddNodesResult; @@ -52,6 +56,8 @@ import java.util.Map; import java.util.concurrent.CompletableFuture; +import static org.eclipse.milo.opcua.stack.core.types.builtin.unsigned.Unsigned.uint; + public class IoTDBOpcUaClientTest { @Test @@ -70,6 +76,30 @@ public void testTransferWritesAllMeasurementsInOneRequest() throws Exception { Mockito.argThat(listWithSize(2))); } + @Test + public void testTransferSplitsWritesAtServerLimit() throws Exception { + final OpcUaClient miloClient = Mockito.mock(OpcUaClient.class); + Mockito.when(miloClient.writeValuesAsync(Mockito.anyList(), Mockito.anyList())) + .thenAnswer( + invocation -> { + final int size = ((List) invocation.getArguments()[0]).size(); + return CompletableFuture.completedFuture(Collections.nCopies(size, StatusCode.GOOD)); + }); + final IoTDBOpcUaClient client = createClient(miloClient, 1, 250); + + client.transfer(createTablet(), createSink()); + + final InOrder inOrder = Mockito.inOrder(miloClient); + inOrder + .verify(miloClient) + .writeValuesAsync( + Mockito.argThat(nodeIds("root/db/d1/s1")), Mockito.argThat(listWithSize(1))); + inOrder + .verify(miloClient) + .writeValuesAsync( + Mockito.argThat(nodeIds("root/db/d1/s2")), Mockito.argThat(listWithSize(1))); + } + @Test public void testTransferLastValuesBatchesDevicesInOneRequest() throws Exception { final OpcUaClient miloClient = Mockito.mock(OpcUaClient.class); @@ -133,6 +163,56 @@ public void testTransferCreatesAndRetriesOnlyMissingNodes() throws Exception { Mockito.argThat(nodeIds("root/db/d1/s1")), Mockito.argThat(listWithSize(1))); } + @Test + public void testTransferSplitsMissingNodeCreationAtServerLimit() throws Exception { + final OpcUaClient miloClient = Mockito.mock(OpcUaClient.class); + Mockito.when(miloClient.writeValuesAsync(Mockito.anyList(), Mockito.anyList())) + .thenReturn( + CompletableFuture.completedFuture( + Arrays.asList( + new StatusCode(StatusCodes.Bad_NodeIdUnknown), + new StatusCode(StatusCodes.Bad_NodeIdUnknown)))) + .thenReturn( + CompletableFuture.completedFuture(Arrays.asList(StatusCode.GOOD, StatusCode.GOOD))); + + final AddNodesResponse addNodesResponse = Mockito.mock(AddNodesResponse.class); + 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.addNodesAsync(Mockito.anyList())) + .thenReturn(CompletableFuture.completedFuture(addNodesResponse)); + + final IoTDBOpcUaClient client = Mockito.spy(createClient(miloClient, 10_000, 1)); + final AddNodesItem firstNode = Mockito.mock(AddNodesItem.class); + final AddNodesItem secondNode = Mockito.mock(AddNodesItem.class); + final ExpandedNodeId firstNodeId = new NodeId(2, "root/db/d1/s1").expanded(); + final ExpandedNodeId secondNodeId = new NodeId(2, "root/db/d1/s2").expanded(); + Mockito.when(firstNode.getRequestedNewNodeId()).thenReturn(firstNodeId); + Mockito.when(secondNode.getRequestedNewNodeId()).thenReturn(secondNodeId); + Mockito.doAnswer( + invocation -> + Collections.singletonList( + "s1".equals(invocation.getArguments()[1]) ? firstNode : secondNode)) + .when(client) + .getNodesToAdd( + Mockito.any(String[].class), + Mockito.anyString(), + Mockito.any(NodeId.class), + Mockito.any()); + + client.transfer(createTablet(), createSink()); + + Mockito.verify(miloClient, Mockito.times(2)).addNodesAsync(Mockito.argThat(listWithSize(1))); + final InOrder inOrder = Mockito.inOrder(miloClient); + inOrder + .verify(miloClient) + .writeValuesAsync(Mockito.argThat(listWithSize(2)), Mockito.argThat(listWithSize(2))); + inOrder.verify(miloClient, Mockito.times(2)).addNodesAsync(Mockito.argThat(listWithSize(1))); + inOrder + .verify(miloClient) + .writeValuesAsync(Mockito.argThat(listWithSize(2)), Mockito.argThat(listWithSize(2))); + } + @Test public void testTransferFailsOnNonRecoverableStatus() throws Exception { final OpcUaClient miloClient = Mockito.mock(OpcUaClient.class); @@ -154,6 +234,12 @@ public void testTransferFailsOnNonRecoverableStatus() throws Exception { } private static IoTDBOpcUaClient createClient(final OpcUaClient miloClient) throws Exception { + return createClient(miloClient, 10_000, 250); + } + + private static IoTDBOpcUaClient createClient( + final OpcUaClient miloClient, final int maxNodesPerWrite, final int maxNodesPerNodeManagement) + throws Exception { final IoTDBOpcUaClient client = new IoTDBOpcUaClient( "opc.tcp://127.0.0.1:12686", SecurityPolicy.None, new AnonymousProvider(), false); @@ -163,6 +249,14 @@ private static IoTDBOpcUaClient createClient(final OpcUaClient miloClient) throw final CompletableFuture connectFuture = CompletableFuture.completedFuture(miloClient); Mockito.when(miloClient.connectAsync()).thenReturn(connectFuture); + Mockito.when( + miloClient.readValuesAsync( + Mockito.anyDouble(), Mockito.eq(TimestampsToReturn.Neither), Mockito.anyList())) + .thenReturn( + CompletableFuture.completedFuture( + Arrays.asList( + new DataValue(new Variant(uint(maxNodesPerWrite))), + new DataValue(new Variant(uint(maxNodesPerNodeManagement)))))); client.run(miloClient); return client; } From b67c2656640bc7f9ec237c7ea5eafbdae6da9d63 Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Tue, 18 Aug 2026 17:56:16 +0800 Subject: [PATCH 2/4] Pipe: remove OPC UA server operation limits --- .../opcua/server/OpcUaServerBuilder.java | 17 +++++++++++++ .../opcua/server/OpcUaServerBuilderTest.java | 25 +++++++++++++++++++ 2 files changed, 42 insertions(+) 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 5b7d820f8f7bb..09af7545b165c 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 @@ -26,6 +26,7 @@ import org.eclipse.milo.opcua.sdk.server.EndpointConfig; import org.eclipse.milo.opcua.sdk.server.OpcUaServer; import org.eclipse.milo.opcua.sdk.server.OpcUaServerConfig; +import org.eclipse.milo.opcua.sdk.server.OpcUaServerConfigLimits; 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; @@ -51,6 +52,7 @@ import org.eclipse.milo.opcua.stack.core.transport.TransportProfile; import org.eclipse.milo.opcua.stack.core.types.builtin.DateTime; import org.eclipse.milo.opcua.stack.core.types.builtin.LocalizedText; +import org.eclipse.milo.opcua.stack.core.types.builtin.unsigned.UInteger; 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; @@ -78,6 +80,7 @@ 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; +import static org.eclipse.milo.opcua.stack.core.types.builtin.unsigned.Unsigned.uint; /** * OPC UA Server builder for IoTDB to send data. The coding style referenced ExampleServer.java in @@ -87,6 +90,8 @@ public class OpcUaServerBuilder implements Closeable { private static final Logger LOGGER = LoggerFactory.getLogger(OpcUaServerBuilder.class); private static final String WILD_CARD_ADDRESS = "0.0.0.0"; + // Milo interprets zero as a zero-sized limit, so use the largest Java list size instead. + private static final int UNLIMITED_OPERATION_LIMIT = Integer.MAX_VALUE; private static final int BACKSLASH = 0x5c; private int tcpBindPort; @@ -310,6 +315,18 @@ protected X509Certificate[] createRsaSha256CertificateChain(final KeyPair keyPai .setSessionSecurityDiagnosticsAccessMode( SessionSecurityDiagnosticsAccessMode.RESTRICTED) .setProductUri("urn:apache:iotdb:opc-ua-server") + .setLimits( + new OpcUaServerConfigLimits() { + @Override + public UInteger getMaxNodesPerWrite() { + return uint(UNLIMITED_OPERATION_LIMIT); + } + + @Override + public UInteger getMaxNodesPerNodeManagement() { + return uint(UNLIMITED_OPERATION_LIMIT); + } + }) .build(); // Setup server to enable event posting 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 8c21c6632ce8f..3cb552acd17ab 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 @@ -145,6 +145,31 @@ public void testNewCertificateContainsAdvertisedHost() throws Exception { } } + @Test + public void testPipeOperationLimitsAreEffectivelyUnlimited() throws Exception { + final Path securityDir = temporaryFolder.newFolder("operation-limits-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)) { + final OpcUaServer server = builder.build(); + + Assert.assertEquals( + Integer.MAX_VALUE, server.getConfig().getLimits().getMaxNodesPerWrite().intValue()); + Assert.assertEquals( + Integer.MAX_VALUE, + server.getConfig().getLimits().getMaxNodesPerNodeManagement().intValue()); + } + } + @Test public void testRebuildWithChangedPassword() throws Exception { final Path securityDir = temporaryFolder.newFolder("changed-password-security").toPath(); From a5f2ef8da9eff6136f68c0d70e2d4da62a8f704f Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Tue, 18 Aug 2026 18:25:32 +0800 Subject: [PATCH 3/4] Revert "Pipe: remove OPC UA server operation limits" This reverts commit b67c2656640bc7f9ec237c7ea5eafbdae6da9d63. --- .../opcua/server/OpcUaServerBuilder.java | 17 ------------- .../opcua/server/OpcUaServerBuilderTest.java | 25 ------------------- 2 files changed, 42 deletions(-) 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 09af7545b165c..5b7d820f8f7bb 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 @@ -26,7 +26,6 @@ import org.eclipse.milo.opcua.sdk.server.EndpointConfig; import org.eclipse.milo.opcua.sdk.server.OpcUaServer; import org.eclipse.milo.opcua.sdk.server.OpcUaServerConfig; -import org.eclipse.milo.opcua.sdk.server.OpcUaServerConfigLimits; 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; @@ -52,7 +51,6 @@ import org.eclipse.milo.opcua.stack.core.transport.TransportProfile; import org.eclipse.milo.opcua.stack.core.types.builtin.DateTime; import org.eclipse.milo.opcua.stack.core.types.builtin.LocalizedText; -import org.eclipse.milo.opcua.stack.core.types.builtin.unsigned.UInteger; 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; @@ -80,7 +78,6 @@ 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; -import static org.eclipse.milo.opcua.stack.core.types.builtin.unsigned.Unsigned.uint; /** * OPC UA Server builder for IoTDB to send data. The coding style referenced ExampleServer.java in @@ -90,8 +87,6 @@ public class OpcUaServerBuilder implements Closeable { private static final Logger LOGGER = LoggerFactory.getLogger(OpcUaServerBuilder.class); private static final String WILD_CARD_ADDRESS = "0.0.0.0"; - // Milo interprets zero as a zero-sized limit, so use the largest Java list size instead. - private static final int UNLIMITED_OPERATION_LIMIT = Integer.MAX_VALUE; private static final int BACKSLASH = 0x5c; private int tcpBindPort; @@ -315,18 +310,6 @@ protected X509Certificate[] createRsaSha256CertificateChain(final KeyPair keyPai .setSessionSecurityDiagnosticsAccessMode( SessionSecurityDiagnosticsAccessMode.RESTRICTED) .setProductUri("urn:apache:iotdb:opc-ua-server") - .setLimits( - new OpcUaServerConfigLimits() { - @Override - public UInteger getMaxNodesPerWrite() { - return uint(UNLIMITED_OPERATION_LIMIT); - } - - @Override - public UInteger getMaxNodesPerNodeManagement() { - return uint(UNLIMITED_OPERATION_LIMIT); - } - }) .build(); // Setup server to enable event posting 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 3cb552acd17ab..8c21c6632ce8f 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 @@ -145,31 +145,6 @@ public void testNewCertificateContainsAdvertisedHost() throws Exception { } } - @Test - public void testPipeOperationLimitsAreEffectivelyUnlimited() throws Exception { - final Path securityDir = temporaryFolder.newFolder("operation-limits-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)) { - final OpcUaServer server = builder.build(); - - Assert.assertEquals( - Integer.MAX_VALUE, server.getConfig().getLimits().getMaxNodesPerWrite().intValue()); - Assert.assertEquals( - Integer.MAX_VALUE, - server.getConfig().getLimits().getMaxNodesPerNodeManagement().intValue()); - } - } - @Test public void testRebuildWithChangedPassword() throws Exception { final Path securityDir = temporaryFolder.newFolder("changed-password-security").toPath(); From 6dd7263cc41883ba9167a36fbeb04cb08102172b Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Fri, 21 Aug 2026 12:15:48 +0800 Subject: [PATCH 4/4] Localize OPC UA operation limit logs --- .../apache/iotdb/db/i18n/DataNodePipeMessages.java | 6 ++++++ .../apache/iotdb/db/i18n/DataNodePipeMessages.java | 6 ++++++ .../sink/protocol/opcua/client/IoTDBOpcUaClient.java | 11 ++++++----- 3 files changed, 18 insertions(+), 5 deletions(-) diff --git a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java index 7da8a5965794b..4df5d7b41479e 100644 --- a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java +++ b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java @@ -2590,4 +2590,10 @@ private DataNodePipeMessages() {} "Failed to release TsFile parser memory for Pipe {} (creation time {}) in DataRegion {} because no reservation exists."; public static final String LOG_PIPE_PROCESSOR_WORKER_ARG_HAS_BEEN_PROCESSING_THE_SAME_EVENT_FOR_ARG_MS_PIPE_ARG_DATAREGION_ARG_SUBTASK_ARG_EVENT_ARG_THREAD_STATE_ARG_STACK_ARG_63B40775 = "Pipe processor worker {} has been processing the same event for {} ms. Pipe: {}, DataRegion: {}, subtask: {}, event: {}, thread state: {}. Stack:{}"; + public static final String LOG_OPC_UA_SERVER_OPERATION_LIMITS_MAXNODESPERWRITE_ARG_MAXNODESPERNODEMANAGEMENT_ARG_5D2BCC90 = + "OPC UA server operation limits: maxNodesPerWrite={}, maxNodesPerNodeManagement={}"; + public static final String LOG_INTERRUPTED_WHILE_READING_OPC_UA_SERVER_OPERATION_LIMITS_USE_DEFAULTS_MAXNODESPERWRITE_ARG_MAXNODESPERNODEMANAGEMENT_ARG_357D46A4 = + "Interrupted while reading OPC UA server operation limits, use defaults: maxNodesPerWrite={}, maxNodesPerNodeManagement={}"; + public static final String LOG_FAILED_TO_READ_OPC_UA_SERVER_OPERATION_LIMITS_USE_DEFAULTS_MAXNODESPERWRITE_ARG_MAXNODESPERNODEMANAGEMENT_ARG_65460871 = + "Failed to read OPC UA server operation limits, use defaults: maxNodesPerWrite={}, maxNodesPerNodeManagement={}"; } diff --git a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java index 35c88b9abfe25..027e6b4c34c6d 100644 --- a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java +++ b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java @@ -2418,4 +2418,10 @@ private DataNodePipeMessages() {} "无法释放 Pipe {}(创建时间 {})在 DataRegion {} 中的 TsFile 解析器内存,因为不存在对应的预留。"; public static final String LOG_PIPE_PROCESSOR_WORKER_ARG_HAS_BEEN_PROCESSING_THE_SAME_EVENT_FOR_ARG_MS_PIPE_ARG_DATAREGION_ARG_SUBTASK_ARG_EVENT_ARG_THREAD_STATE_ARG_STACK_ARG_63B40775 = "Pipe processor worker {} 已连续处理同一 event {} ms。Pipe:{},DataRegion:{},subtask:{},event:{},线程状态:{}。栈:{}"; + public static final String LOG_OPC_UA_SERVER_OPERATION_LIMITS_MAXNODESPERWRITE_ARG_MAXNODESPERNODEMANAGEMENT_ARG_5D2BCC90 = + "OPC UA 服务器操作限制:maxNodesPerWrite={},maxNodesPerNodeManagement={}"; + public static final String LOG_INTERRUPTED_WHILE_READING_OPC_UA_SERVER_OPERATION_LIMITS_USE_DEFAULTS_MAXNODESPERWRITE_ARG_MAXNODESPERNODEMANAGEMENT_ARG_357D46A4 = + "读取 OPC UA 服务器操作限制时被中断,使用默认值:maxNodesPerWrite={},maxNodesPerNodeManagement={}"; + public static final String LOG_FAILED_TO_READ_OPC_UA_SERVER_OPERATION_LIMITS_USE_DEFAULTS_MAXNODESPERWRITE_ARG_MAXNODESPERNODEMANAGEMENT_ARG_65460871 = + "读取 OPC UA 服务器操作限制失败,使用默认值:maxNodesPerWrite={},maxNodesPerNodeManagement={}"; } 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 238ac5afe31ff..f00ec84504763 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 @@ -145,20 +145,21 @@ private void updateOperationLimits() { maxNodesPerNodeManagement = getOperationLimit(operationLimits, 1, DEFAULT_MAX_NODES_PER_NODE_MANAGEMENT); LOGGER.info( - "OPC UA server operation limits: maxNodesPerWrite={}, maxNodesPerNodeManagement={}", + DataNodePipeMessages + .LOG_OPC_UA_SERVER_OPERATION_LIMITS_MAXNODESPERWRITE_ARG_MAXNODESPERNODEMANAGEMENT_ARG_5D2BCC90, maxNodesPerWrite, maxNodesPerNodeManagement); } catch (final InterruptedException e) { Thread.currentThread().interrupt(); LOGGER.warn( - "Interrupted while reading OPC UA server operation limits, use defaults: " - + "maxNodesPerWrite={}, maxNodesPerNodeManagement={}", + DataNodePipeMessages + .LOG_INTERRUPTED_WHILE_READING_OPC_UA_SERVER_OPERATION_LIMITS_USE_DEFAULTS_MAXNODESPERWRITE_ARG_MAXNODESPERNODEMANAGEMENT_ARG_357D46A4, DEFAULT_MAX_NODES_PER_WRITE, DEFAULT_MAX_NODES_PER_NODE_MANAGEMENT); } catch (final Exception e) { LOGGER.warn( - "Failed to read OPC UA server operation limits, use defaults: " - + "maxNodesPerWrite={}, maxNodesPerNodeManagement={}", + DataNodePipeMessages + .LOG_FAILED_TO_READ_OPC_UA_SERVER_OPERATION_LIMITS_USE_DEFAULTS_MAXNODESPERWRITE_ARG_MAXNODESPERNODEMANAGEMENT_ARG_65460871, DEFAULT_MAX_NODES_PER_WRITE, DEFAULT_MAX_NODES_PER_NODE_MANAGEMENT, e);