This is an automated email from the ASF dual-hosted git repository.
rong 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 457848c39fb [IOTDB-6220] Pipe: Added connector endpoint check logic to
avoid self-transmission (#11417)
457848c39fb is described below
commit 457848c39fb1758100b430d63794f15fb2b650fe
Author: Caideyipi <[email protected]>
AuthorDate: Tue Oct 31 13:35:59 2023 +0800
[IOTDB-6220] Pipe: Added connector endpoint check logic to avoid
self-transmission (#11417)
---
.../db/pipe/connector/protocol/IoTDBConnector.java | 22 ++++---
.../protocol/airgap/IoTDBAirGapConnector.java | 28 ++++++++
.../protocol/legacy/IoTDBLegacyPipeConnector.java | 67 +++++++++++++++----
.../thrift/sync/IoTDBThriftSyncConnector.java | 22 +++++++
.../iotdb/db/pipe/connector/PipeConnectorTest.java | 76 ++++++++++++++++++----
5 files changed, 181 insertions(+), 34 deletions(-)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/IoTDBConnector.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/IoTDBConnector.java
index 3a2bae03068..5f295b550ff 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/IoTDBConnector.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/IoTDBConnector.java
@@ -81,6 +81,18 @@ public abstract class IoTDBConnector implements
PipeConnector {
@Override
public void customize(PipeParameters parameters,
PipeConnectorRuntimeConfiguration configuration)
throws Exception {
+ nodeUrls.clear();
+ nodeUrls.addAll(parseNodeUrls(parameters));
+ LOGGER.info("IoTDBConnector nodeUrls: {}", nodeUrls);
+
+ isTabletBatchModeEnabled =
+ parameters.getBooleanOrDefault(
+ Arrays.asList(CONNECTOR_IOTDB_BATCH_MODE_ENABLE_KEY,
SINK_IOTDB_BATCH_MODE_ENABLE_KEY),
+ CONNECTOR_IOTDB_BATCH_MODE_ENABLE_DEFAULT_VALUE);
+ LOGGER.info("IoTDBConnector isTabletBatchModeEnabled: {}",
isTabletBatchModeEnabled);
+ }
+
+ protected Set<TEndPoint> parseNodeUrls(PipeParameters parameters) {
final Set<TEndPoint> givenNodeUrls = new HashSet<>(nodeUrls);
if (parameters.hasAttribute(CONNECTOR_IOTDB_IP_KEY)
@@ -110,14 +122,6 @@ public abstract class IoTDBConnector implements
PipeConnector {
Arrays.asList(parameters.getString(SINK_IOTDB_NODE_URLS_KEY).split(","))));
}
- nodeUrls.clear();
- nodeUrls.addAll(givenNodeUrls);
- LOGGER.info("IoTDBConnector nodeUrls: {}", nodeUrls);
-
- isTabletBatchModeEnabled =
- parameters.getBooleanOrDefault(
- Arrays.asList(CONNECTOR_IOTDB_BATCH_MODE_ENABLE_KEY,
SINK_IOTDB_BATCH_MODE_ENABLE_KEY),
- CONNECTOR_IOTDB_BATCH_MODE_ENABLE_DEFAULT_VALUE);
- LOGGER.info("IoTDBConnector isTabletBatchModeEnabled: {}",
isTabletBatchModeEnabled);
+ return givenNodeUrls;
}
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/airgap/IoTDBAirGapConnector.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/airgap/IoTDBAirGapConnector.java
index f7c7d26f94c..427db8030c7 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/airgap/IoTDBAirGapConnector.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/airgap/IoTDBAirGapConnector.java
@@ -19,8 +19,11 @@
package org.apache.iotdb.db.pipe.connector.protocol.airgap;
+import org.apache.iotdb.common.rpc.thrift.TEndPoint;
import org.apache.iotdb.commons.conf.CommonDescriptor;
import org.apache.iotdb.commons.pipe.config.PipeConfig;
+import org.apache.iotdb.db.conf.IoTDBConfig;
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
import
org.apache.iotdb.db.pipe.connector.payload.airgap.AirGapELanguageConstant;
import org.apache.iotdb.db.pipe.connector.payload.airgap.AirGapOneByteResponse;
import
org.apache.iotdb.db.pipe.connector.payload.evolvable.request.PipeTransferFilePieceReq;
@@ -37,6 +40,7 @@ import
org.apache.iotdb.db.pipe.event.common.tablet.PipeRawTabletInsertionEvent;
import org.apache.iotdb.db.pipe.event.common.tsfile.PipeTsFileInsertionEvent;
import
org.apache.iotdb.db.storageengine.dataregion.wal.exception.WALPipeException;
import
org.apache.iotdb.pipe.api.customizer.configuration.PipeConnectorRuntimeConfiguration;
+import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameterValidator;
import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameters;
import org.apache.iotdb.pipe.api.event.Event;
import org.apache.iotdb.pipe.api.event.dml.insertion.TabletInsertionEvent;
@@ -57,6 +61,7 @@ import java.net.Socket;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
+import java.util.Set;
import java.util.zip.CRC32;
import static org.apache.iotdb.commons.utils.BasicStructureSerDeUtil.LONG_LEN;
@@ -81,6 +86,29 @@ public class IoTDBAirGapConnector extends IoTDBConnector {
private long currentClientIndex = 0;
+ @Override
+ public void validate(PipeParameterValidator validator) throws Exception {
+ super.validate(validator);
+ final IoTDBConfig ioTDBConfig = IoTDBDescriptor.getInstance().getConfig();
+ final PipeConfig pipeConfig = PipeConfig.getInstance();
+ Set<TEndPoint> givenNodeUrls = parseNodeUrls(validator.getParameters());
+
+ validator.validate(
+ empty ->
+ !(pipeConfig.getPipeAirGapReceiverEnabled()
+ && (givenNodeUrls.contains(
+ new TEndPoint(
+ ioTDBConfig.getRpcAddress(),
pipeConfig.getPipeAirGapReceiverPort())))
+ || givenNodeUrls.contains(
+ new TEndPoint("127.0.0.1",
pipeConfig.getPipeAirGapReceiverPort()))
+ || givenNodeUrls.contains(
+ new TEndPoint("0.0.0.0",
pipeConfig.getPipeAirGapReceiverPort()))),
+ String.format(
+ "One of the endpoints %s of the receivers is pointing back to the
air gap receiver %s on sender itself",
+ givenNodeUrls,
+ new TEndPoint(ioTDBConfig.getRpcAddress(),
pipeConfig.getPipeAirGapReceiverPort())));
+ }
+
@Override
public void customize(PipeParameters parameters,
PipeConnectorRuntimeConfiguration configuration)
throws Exception {
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/legacy/IoTDBLegacyPipeConnector.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/legacy/IoTDBLegacyPipeConnector.java
index af607cdb06d..f202fca7f23 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/legacy/IoTDBLegacyPipeConnector.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/legacy/IoTDBLegacyPipeConnector.java
@@ -19,6 +19,7 @@
package org.apache.iotdb.db.pipe.connector.protocol.legacy;
+import org.apache.iotdb.common.rpc.thrift.TEndPoint;
import org.apache.iotdb.common.rpc.thrift.TSStatus;
import org.apache.iotdb.commons.client.property.ThriftClientProperty;
import org.apache.iotdb.commons.conf.CommonConfig;
@@ -26,6 +27,8 @@ import org.apache.iotdb.commons.conf.CommonDescriptor;
import org.apache.iotdb.commons.conf.IoTDBConstant;
import org.apache.iotdb.commons.exception.pipe.PipeRuntimeCriticalException;
import org.apache.iotdb.commons.pipe.config.PipeConfig;
+import org.apache.iotdb.db.conf.IoTDBConfig;
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
import org.apache.iotdb.db.pipe.connector.payload.legacy.TsFilePipeData;
import
org.apache.iotdb.db.pipe.connector.protocol.thrift.sync.IoTDBThriftSyncConnectorClient;
import org.apache.iotdb.db.pipe.event.common.heartbeat.PipeHeartbeatEvent;
@@ -59,6 +62,8 @@ import java.io.IOException;
import java.io.RandomAccessFile;
import java.nio.ByteBuffer;
import java.util.Arrays;
+import java.util.HashSet;
+import java.util.Set;
import static
org.apache.iotdb.db.pipe.config.constant.PipeConnectorConstant.CONNECTOR_IOTDB_IP_KEY;
import static
org.apache.iotdb.db.pipe.config.constant.PipeConnectorConstant.CONNECTOR_IOTDB_PASSWORD_DEFAULT_VALUE;
@@ -98,19 +103,55 @@ public class IoTDBLegacyPipeConnector implements
PipeConnector {
@Override
public void validate(PipeParameterValidator validator) throws Exception {
final PipeParameters parameters = validator.getParameters();
- validator.validate(
- args ->
- ((boolean) args[0] && (boolean) args[1]) || ((boolean) args[2] &&
(boolean) args[3]),
- String.format(
- "Either %s:%s or %s:%s must be specified",
- CONNECTOR_IOTDB_IP_KEY,
- CONNECTOR_IOTDB_PORT_KEY,
- SINK_IOTDB_IP_KEY,
- SINK_IOTDB_PORT_KEY),
- parameters.hasAttribute(CONNECTOR_IOTDB_IP_KEY),
- parameters.hasAttribute(CONNECTOR_IOTDB_PORT_KEY),
- parameters.hasAttribute(SINK_IOTDB_IP_KEY),
- parameters.hasAttribute(SINK_IOTDB_PORT_KEY));
+ final IoTDBConfig ioTDBConfig = IoTDBDescriptor.getInstance().getConfig();
+ Set<TEndPoint> givenNodeUrls = parseNodeUrls(validator.getParameters());
+
+ validator
+ .validate(
+ args ->
+ ((boolean) args[0] && (boolean) args[1])
+ || ((boolean) args[2] && (boolean) args[3]),
+ String.format(
+ "Either %s:%s or %s:%s must be specified",
+ CONNECTOR_IOTDB_IP_KEY,
+ CONNECTOR_IOTDB_PORT_KEY,
+ SINK_IOTDB_IP_KEY,
+ SINK_IOTDB_PORT_KEY),
+ parameters.hasAttribute(CONNECTOR_IOTDB_IP_KEY),
+ parameters.hasAttribute(CONNECTOR_IOTDB_PORT_KEY),
+ parameters.hasAttribute(SINK_IOTDB_IP_KEY),
+ parameters.hasAttribute(SINK_IOTDB_PORT_KEY))
+ .validate(
+ empty ->
+ !(givenNodeUrls.contains(
+ new TEndPoint(ioTDBConfig.getRpcAddress(),
ioTDBConfig.getRpcPort()))
+ || givenNodeUrls.contains(new TEndPoint("127.0.0.1",
ioTDBConfig.getRpcPort()))
+ || givenNodeUrls.contains(new TEndPoint("0.0.0.0",
ioTDBConfig.getRpcPort()))),
+ String.format(
+ "One of the endpoints %s of the receivers is pointing back to
the legacy receiver on sender %s itself",
+ givenNodeUrls,
+ new TEndPoint(ioTDBConfig.getRpcAddress(),
ioTDBConfig.getRpcPort())));
+ }
+
+ private Set<TEndPoint> parseNodeUrls(PipeParameters parameters) {
+ final Set<TEndPoint> givenNodeUrls = new HashSet<>();
+
+ if (parameters.hasAttribute(CONNECTOR_IOTDB_IP_KEY)
+ && parameters.hasAttribute(CONNECTOR_IOTDB_PORT_KEY)) {
+ givenNodeUrls.add(
+ new TEndPoint(
+ parameters.getString(CONNECTOR_IOTDB_IP_KEY),
+ parameters.getInt(CONNECTOR_IOTDB_PORT_KEY)));
+ }
+
+ if (parameters.hasAttribute(SINK_IOTDB_IP_KEY)
+ && parameters.hasAttribute(SINK_IOTDB_PORT_KEY)) {
+ givenNodeUrls.add(
+ new TEndPoint(
+ parameters.getString(SINK_IOTDB_IP_KEY),
parameters.getInt(SINK_IOTDB_PORT_KEY)));
+ }
+
+ return givenNodeUrls;
}
@Override
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/sync/IoTDBThriftSyncConnector.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/sync/IoTDBThriftSyncConnector.java
index dd74f2dd7f5..bfa2297086d 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/sync/IoTDBThriftSyncConnector.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/sync/IoTDBThriftSyncConnector.java
@@ -19,9 +19,12 @@
package org.apache.iotdb.db.pipe.connector.protocol.thrift.sync;
+import org.apache.iotdb.common.rpc.thrift.TEndPoint;
import org.apache.iotdb.commons.client.property.ThriftClientProperty;
import org.apache.iotdb.commons.conf.CommonDescriptor;
import org.apache.iotdb.commons.pipe.config.PipeConfig;
+import org.apache.iotdb.db.conf.IoTDBConfig;
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
import
org.apache.iotdb.db.pipe.connector.payload.evolvable.builder.IoTDBThriftSyncPipeTransferBatchReqBuilder;
import
org.apache.iotdb.db.pipe.connector.payload.evolvable.reponse.PipeTransferFilePieceResp;
import
org.apache.iotdb.db.pipe.connector.payload.evolvable.request.PipeTransferFilePieceReq;
@@ -39,6 +42,7 @@ import
org.apache.iotdb.db.pipe.event.common.tablet.PipeRawTabletInsertionEvent;
import org.apache.iotdb.db.pipe.event.common.tsfile.PipeTsFileInsertionEvent;
import
org.apache.iotdb.db.storageengine.dataregion.wal.exception.WALPipeException;
import
org.apache.iotdb.pipe.api.customizer.configuration.PipeConnectorRuntimeConfiguration;
+import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameterValidator;
import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameters;
import org.apache.iotdb.pipe.api.event.Event;
import org.apache.iotdb.pipe.api.event.dml.insertion.TabletInsertionEvent;
@@ -58,6 +62,7 @@ import java.io.RandomAccessFile;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
+import java.util.Set;
public class IoTDBThriftSyncConnector extends IoTDBConnector {
@@ -76,6 +81,23 @@ public class IoTDBThriftSyncConnector extends IoTDBConnector
{
// Do nothing
}
+ @Override
+ public void validate(PipeParameterValidator validator) throws Exception {
+ super.validate(validator);
+ final IoTDBConfig ioTDBConfig = IoTDBDescriptor.getInstance().getConfig();
+ Set<TEndPoint> givenNodeUrls = parseNodeUrls(validator.getParameters());
+
+ validator.validate(
+ empty ->
+ !(givenNodeUrls.contains(
+ new TEndPoint(ioTDBConfig.getRpcAddress(),
ioTDBConfig.getRpcPort()))
+ || givenNodeUrls.contains(new TEndPoint("127.0.0.1",
ioTDBConfig.getRpcPort()))
+ || givenNodeUrls.contains(new TEndPoint("0.0.0.0",
ioTDBConfig.getRpcPort()))),
+ String.format(
+ "One of the endpoints %s of the receivers is pointing back to the
thrift receiver %s on sender itself",
+ givenNodeUrls, new TEndPoint(ioTDBConfig.getRpcAddress(),
ioTDBConfig.getRpcPort())));
+ }
+
@Override
public void customize(PipeParameters parameters,
PipeConnectorRuntimeConfiguration configuration)
throws Exception {
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/connector/PipeConnectorTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/connector/PipeConnectorTest.java
index 5e1192c9dab..d5c84227abd 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/connector/PipeConnectorTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/connector/PipeConnectorTest.java
@@ -26,6 +26,7 @@ import
org.apache.iotdb.db.pipe.connector.protocol.thrift.async.IoTDBThriftAsync
import
org.apache.iotdb.db.pipe.connector.protocol.thrift.sync.IoTDBThriftSyncConnector;
import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameterValidator;
import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameters;
+import org.apache.iotdb.pipe.api.exception.PipeParameterNotValidException;
import org.junit.Assert;
import org.junit.Test;
@@ -33,10 +34,10 @@ import org.junit.Test;
import java.util.HashMap;
public class PipeConnectorTest {
- @Test
- public void testIoTDBLegacyPipeConnector() {
- IoTDBLegacyPipeConnector connector = new IoTDBLegacyPipeConnector();
- try {
+
+ @Test(expected = PipeParameterNotValidException.class)
+ public void testIoTDBLegacyPipeConnectorToSelf() throws Exception {
+ try (IoTDBLegacyPipeConnector connector = new IoTDBLegacyPipeConnector()) {
connector.validate(
new PipeParameterValidator(
new PipeParameters(
@@ -49,15 +50,32 @@ public class PipeConnectorTest {
put(PipeConnectorConstant.CONNECTOR_IOTDB_PORT_KEY,
"6667");
}
})));
+ }
+ }
+
+ @Test
+ public void testIoTDBLegacyPipeConnectorToOthers() {
+ try (IoTDBLegacyPipeConnector connector = new IoTDBLegacyPipeConnector()) {
+ connector.validate(
+ new PipeParameterValidator(
+ new PipeParameters(
+ new HashMap<String, String>() {
+ {
+ put(
+ PipeConnectorConstant.CONNECTOR_KEY,
+
BuiltinPipePlugin.IOTDB_LEGACY_PIPE_CONNECTOR.getPipePluginName());
+ put(PipeConnectorConstant.CONNECTOR_IOTDB_IP_KEY,
"127.0.0.1");
+ put(PipeConnectorConstant.CONNECTOR_IOTDB_PORT_KEY,
"6668");
+ }
+ })));
} catch (Exception e) {
Assert.fail();
}
}
- @Test
- public void testIoTDBThriftSyncConnector() {
- IoTDBThriftSyncConnector connector = new IoTDBThriftSyncConnector();
- try {
+ @Test(expected = PipeParameterNotValidException.class)
+ public void testIoTDBThriftSyncConnectorToSelf() throws Exception {
+ try (IoTDBThriftSyncConnector connector = new IoTDBThriftSyncConnector()) {
connector.validate(
new PipeParameterValidator(
new PipeParameters(
@@ -70,15 +88,32 @@ public class PipeConnectorTest {
put(PipeConnectorConstant.CONNECTOR_IOTDB_PORT_KEY,
"6667");
}
})));
+ }
+ }
+
+ @Test
+ public void testIoTDBThriftSyncConnectorToOthers() {
+ try (IoTDBThriftSyncConnector connector = new IoTDBThriftSyncConnector()) {
+ connector.validate(
+ new PipeParameterValidator(
+ new PipeParameters(
+ new HashMap<String, String>() {
+ {
+ put(
+ PipeConnectorConstant.CONNECTOR_KEY,
+
BuiltinPipePlugin.IOTDB_THRIFT_SYNC_CONNECTOR.getPipePluginName());
+ put(PipeConnectorConstant.CONNECTOR_IOTDB_IP_KEY,
"127.0.0.1");
+ put(PipeConnectorConstant.CONNECTOR_IOTDB_PORT_KEY,
"6668");
+ }
+ })));
} catch (Exception e) {
Assert.fail();
}
}
- @Test
- public void testIoTDBThriftAsyncConnector() {
- IoTDBThriftAsyncConnector connector = new IoTDBThriftAsyncConnector();
- try {
+ @Test(expected = PipeParameterNotValidException.class)
+ public void testIoTDBThriftAsyncConnectorToSelf() throws Exception {
+ try (IoTDBThriftAsyncConnector connector = new
IoTDBThriftAsyncConnector()) {
connector.validate(
new PipeParameterValidator(
new PipeParameters(
@@ -90,6 +125,23 @@ public class PipeConnectorTest {
put(PipeConnectorConstant.CONNECTOR_IOTDB_NODE_URLS_KEY,
"127.0.0.1:6667");
}
})));
+ }
+ }
+
+ @Test
+ public void testIoTDBThriftAsyncConnectorToOthers() {
+ try (IoTDBThriftAsyncConnector connector = new
IoTDBThriftAsyncConnector()) {
+ connector.validate(
+ new PipeParameterValidator(
+ new PipeParameters(
+ new HashMap<String, String>() {
+ {
+ put(
+ PipeConnectorConstant.CONNECTOR_KEY,
+
BuiltinPipePlugin.IOTDB_THRIFT_ASYNC_CONNECTOR.getPipePluginName());
+ put(PipeConnectorConstant.CONNECTOR_IOTDB_NODE_URLS_KEY,
"127.0.0.1:6668");
+ }
+ })));
} catch (Exception e) {
Assert.fail();
}