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();
     }

Reply via email to