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 09252d40064 feat(load): reject managed TsFile directories (#18271)
09252d40064 is described below

commit 09252d40064eefbdc6df01a0ac5d2868fc5bcfd3
Author: Zhenyu Luo <[email protected]>
AuthorDate: Mon Aug 3 08:51:34 2026 +0800

    feat(load): reject managed TsFile directories (#18271)
    
    * feat(load): reject managed TsFile directories
    
    * perf(load): cache internal TsFile directories
    
    * fix(load): reject all local data directories
    
    * test(load): cover pipe receiver outside data directory
    
    * fix(load): reject IoTDB data directory
    
    * fix(pipe): preserve load source during type conversion
    
    * fix(pipe): preserve load source during scheduler retry
    
    * fix(pipe): preserve load source during active load
    
    * fix(pipe): keep active load source validation
    
    * style(load): replace FQCN with LoadTsFileStatement import
    
    Co-authored-by: Cursor <[email protected]>
    
    ---------
    
    Co-authored-by: Cursor <[email protected]>
---
 .../org/apache/iotdb/db/it/IoTDBLoadTsFileIT.java  | 22 +++++++
 .../iotdb/db/i18n/DataNodeQueryMessages.java       |  3 +
 .../iotdb/db/i18n/DataNodeQueryMessages.java       |  3 +
 .../java/org/apache/iotdb/db/conf/IoTDBConfig.java | 22 +++++++
 .../protocol/legacy/loader/TsFileLoader.java       |  2 +-
 .../protocol/thrift/IoTDBDataNodeReceiver.java     |  2 +-
 .../plan/analyze/load/LoadTsFileAnalyzer.java      | 22 +++++--
 .../plan/relational/sql/ast/LoadTsFile.java        | 22 +++++--
 .../plan/scheduler/load/LoadTsFileScheduler.java   |  9 ++-
 .../plan/statement/crud/LoadTsFileStatement.java   | 67 +++++++++++++++++++---
 .../statement/crud/LoadTsFileStatementTest.java    | 60 +++++++++++++++++++
 11 files changed, 209 insertions(+), 25 deletions(-)

diff --git 
a/integration-test/src/test/java/org/apache/iotdb/db/it/IoTDBLoadTsFileIT.java 
b/integration-test/src/test/java/org/apache/iotdb/db/it/IoTDBLoadTsFileIT.java
index 8a0800771cf..fa5e0ecc313 100644
--- 
a/integration-test/src/test/java/org/apache/iotdb/db/it/IoTDBLoadTsFileIT.java
+++ 
b/integration-test/src/test/java/org/apache/iotdb/db/it/IoTDBLoadTsFileIT.java
@@ -890,6 +890,28 @@ public class IoTDBLoadTsFileIT {
     }
   }
 
+  @Test
+  public void 
testLoadDataNodeInternalDataDirectoryIsRejectedWithoutLeakingPath() throws 
Exception {
+    final DataNodeWrapper dataNodeWrapper = 
EnvFactory.getEnv().getDataNodeWrapper(0);
+    final File dataDir = new File(dataNodeWrapper.getDataPath());
+
+    try (final Connection connection =
+            
EnvFactory.getEnv().getConnectionWithSpecifiedDataNode(dataNodeWrapper);
+        final Statement statement = connection.createStatement()) {
+      try {
+        statement.execute(String.format("load \"%s\"", 
dataDir.getAbsolutePath()));
+        Assert.fail("Expected LOAD from the DataNode internal data directory 
to be rejected.");
+      } catch (final SQLException e) {
+        Assert.assertTrue(
+            e.getMessage(),
+            e.getMessage()
+                .contains(
+                    "Cannot load files because the specified directory 
contains IoTDB data."));
+        Assert.assertFalse(e.getMessage(), 
e.getMessage().contains(dataDir.getAbsolutePath()));
+      }
+    }
+  }
+
   @Test
   public void testLoadWithMods() throws Exception {
     final long writtenPoint1;
diff --git 
a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java
 
b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java
index e8e1954751b..1c7c83061e3 100644
--- 
a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java
+++ 
b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java
@@ -3037,6 +3037,9 @@ public final class DataNodeQueryMessages {
       "Can not find %s on this machine, notice that load can only handle files 
on this machine.";
   public static final String 
QUERY_EXCEPTION_LOAD_TSFILE_SOURCE_PATH_S_IS_OUTSIDE_ALLOWED_DIRECTORIES_85A6019F
 =
       "Load TsFile source path %s is outside allowed directories %s.";
+  public static final String
+      
QUERY_EXCEPTION_CANNOT_LOAD_FILES_BECAUSE_SPECIFIED_DIRECTORY_CONTAINS_IOTDB_DATA_B0A1B93D
 =
+          "Cannot load files because the specified directory contains IoTDB 
data.";
   public static final String 
QUERY_EXCEPTION_FAILED_TO_RESOLVE_CANONICAL_PATH_FOR_LOAD_TSFILE_SOURCE_09CC9AC6
 =
       "Failed to resolve canonical path for Load TsFile source %s: %s";
   public static final String 
QUERY_EXCEPTION_DATA_TYPE_IS_NOT_CONSISTENT_INPUT_S_REGISTERED_S_AE9DBDC0 =
diff --git 
a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java
 
b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java
index c06b5332f7b..ba502738b83 100644
--- 
a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java
+++ 
b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeQueryMessages.java
@@ -3654,6 +3654,9 @@ public final class DataNodeQueryMessages {
   public static final String 
QUERY_EXCEPTION_LOAD_TSFILE_SOURCE_PATH_S_IS_OUTSIDE_ALLOWED_DIRECTORIES_85A6019F
 =
 
       "加载 TsFile 的源路径 %s 位于允许目录 %s 之外。";
+  public static final String
+      
QUERY_EXCEPTION_CANNOT_LOAD_FILES_BECAUSE_SPECIFIED_DIRECTORY_CONTAINS_IOTDB_DATA_B0A1B93D
 =
+          "指定目录包含 IoTDB 数据,无法加载文件。";
   public static final String 
QUERY_EXCEPTION_FAILED_TO_RESOLVE_CANONICAL_PATH_FOR_LOAD_TSFILE_SOURCE_09CC9AC6
 =
       "无法解析 load TsFile source %s 的 canonical path:%s";
   public static final String 
QUERY_EXCEPTION_DATA_TYPE_IS_NOT_CONSISTENT_INPUT_S_REGISTERED_S_AE9DBDC0 =
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java
index cf901ebd2db..19c3cc6217f 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java
@@ -310,6 +310,8 @@ public class IoTDBConfig {
 
   private CanonicalPaths loadTsFileAllowedDirCanonicalPaths = 
canonicalPaths(loadTsFileAllowedDirs);
 
+  private volatile CanonicalPaths internalDataDirCanonicalPaths = new 
CanonicalPaths(new Path[0]);
+
   private boolean loadTsFileSourcePathCheckEnabled = false;
 
   /** Strategy of multiple directories. */
@@ -1435,6 +1437,7 @@ public class IoTDBConfig {
     queryDir = addDataHomeDir(queryDir);
     sortTmpDir = addDataHomeDir(sortTmpDir);
     formulateDataDirs(tierDataDirs);
+    formulateInternalDataDirs(tierDataDirs);
   }
 
   private void formulateDataDirs(String[][] tierDataDirs) {
@@ -1486,6 +1489,7 @@ public class IoTDBConfig {
       }
     }
     this.tierDataDirs = newTierDataDirs;
+    formulateInternalDataDirs(newTierDataDirs);
     reloadSystemMetrics();
   }
 
@@ -1562,6 +1566,10 @@ public class IoTDBConfig {
         .toArray(String[]::new);
   }
 
+  public Path[] getInternalDataDirCanonicalPaths() throws 
FileNotFoundException {
+    return internalDataDirCanonicalPaths.getPaths();
+  }
+
   public String[][] getTierDataDirs() {
     return tierDataDirs;
   }
@@ -1570,6 +1578,7 @@ public class IoTDBConfig {
   public void setTierDataDirs(String[][] tierDataDirs) {
     formulateDataDirs(tierDataDirs);
     this.tierDataDirs = tierDataDirs;
+    formulateInternalDataDirs(tierDataDirs);
     // TODO(szywilliam): rewrite the logic here when ratis supports complete 
snapshot semantic
     setRatisDataRegionSnapshotDir(
         tierDataDirs[0][0] + File.separator + 
IoTDBConstant.SNAPSHOT_FOLDER_NAME);
@@ -1669,6 +1678,19 @@ public class IoTDBConfig {
     this.loadTsFileDirCanonicalPaths = canonicalPaths(newLoadTsFileDirs);
   }
 
+  private void formulateInternalDataDirs(final String[][] tierDataDirs) {
+    final List<String> internalDataDirs = new ArrayList<>();
+    internalDataDirs.add(addDataHomeDir("data"));
+    for (final String[] tierDataDir : tierDataDirs) {
+      for (final String dataDir : tierDataDir) {
+        if (FSUtils.isLocal(dataDir)) {
+          internalDataDirs.add(dataDir);
+        }
+      }
+    }
+    internalDataDirCanonicalPaths = 
canonicalPaths(internalDataDirs.toArray(new String[0]));
+  }
+
   private static CanonicalPaths canonicalPaths(final String[] dirs) {
     final Path[] paths = new Path[dirs.length];
     for (int i = 0; i < dirs.length; i++) {
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/legacy/loader/TsFileLoader.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/legacy/loader/TsFileLoader.java
index bf6f717be01..b91547137fb 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/legacy/loader/TsFileLoader.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/legacy/loader/TsFileLoader.java
@@ -53,7 +53,7 @@ public class TsFileLoader implements ILoader {
   @Override
   public void load(final SessionInfo sessionInfo) {
     try {
-      LoadTsFileStatement statement = 
LoadTsFileStatement.createUnchecked(tsFile.getAbsolutePath());
+      LoadTsFileStatement statement = 
LoadTsFileStatement.createForPipe(tsFile.getAbsolutePath());
       statement.setDeleteAfterLoad(true);
       statement.setConvertOnTypeMismatch(true);
       statement.setDatabaseLevel(parseSgLevel());
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/thrift/IoTDBDataNodeReceiver.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/thrift/IoTDBDataNodeReceiver.java
index 7424f20bb8a..8257dadfc99 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/thrift/IoTDBDataNodeReceiver.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/thrift/IoTDBDataNodeReceiver.java
@@ -714,7 +714,7 @@ public class IoTDBDataNodeReceiver extends 
IoTDBFileReceiver {
       final boolean shouldConvertDataTypeOnTypeMismatch,
       final boolean shouldWaitForSchemaBeforeLoad)
       throws FileNotFoundException {
-    final LoadTsFileStatement statement = 
LoadTsFileStatement.createUnchecked(fileAbsolutePath);
+    final LoadTsFileStatement statement = 
LoadTsFileStatement.createForPipe(fileAbsolutePath);
     statement.setDeleteAfterLoad(true);
     statement.setConvertOnTypeMismatch(shouldConvertDataTypeOnTypeMismatch);
     statement.setVerifySchema(
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileAnalyzer.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileAnalyzer.java
index 34827de0c63..68b23ae5d00 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileAnalyzer.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/analyze/load/LoadTsFileAnalyzer.java
@@ -507,15 +507,20 @@ public class LoadTsFileAnalyzer implements AutoCloseable {
           isTableModelTsFile.get(i)
               ? loadTsFileDataTypeConverter
                   .convertForTableModel(
-                      LoadTsFile.createUnchecked(
-                              null, tsFiles.get(i).getPath(), 
Collections.emptyMap())
+                      (isGeneratedByPipe
+                              ? LoadTsFile.createForPipe(
+                                  null, tsFiles.get(i).getPath(), 
Collections.emptyMap())
+                              : LoadTsFile.createUnchecked(
+                                  null, tsFiles.get(i).getPath(), 
Collections.emptyMap()))
                           .setDatabase(databaseForTableData)
                           .setDeleteAfterLoad(isDeleteAfterLoad)
                           .setConvertOnTypeMismatch(isConvertOnTypeMismatch))
                   .orElse(null)
               : loadTsFileDataTypeConverter
                   .convertForTreeModel(
-                      
LoadTsFileStatement.createUnchecked(tsFiles.get(i).getPath())
+                      (isGeneratedByPipe
+                              ? 
LoadTsFileStatement.createForPipe(tsFiles.get(i).getPath())
+                              : 
LoadTsFileStatement.createUnchecked(tsFiles.get(i).getPath()))
                           .setDeleteAfterLoad(isDeleteAfterLoad)
                           .setConvertOnTypeMismatch(isConvertOnTypeMismatch))
                   .orElse(null);
@@ -800,15 +805,20 @@ public class LoadTsFileAnalyzer implements AutoCloseable {
             isTableModelTsFile.get(i)
                 ? loadTsFileDataTypeConverter
                     .convertForTableModel(
-                        LoadTsFile.createUnchecked(
-                                null, tsFiles.get(i).getPath(), 
Collections.emptyMap())
+                        (isGeneratedByPipe
+                                ? LoadTsFile.createForPipe(
+                                    null, tsFiles.get(i).getPath(), 
Collections.emptyMap())
+                                : LoadTsFile.createUnchecked(
+                                    null, tsFiles.get(i).getPath(), 
Collections.emptyMap()))
                             .setDatabase(databaseForTableData)
                             .setDeleteAfterLoad(isDeleteAfterLoad)
                             .setConvertOnTypeMismatch(isConvertOnTypeMismatch))
                     .orElse(null)
                 : loadTsFileDataTypeConverter
                     .convertForTreeModel(
-                        
LoadTsFileStatement.createUnchecked(tsFiles.get(i).getPath())
+                        (isGeneratedByPipe
+                                ? 
LoadTsFileStatement.createForPipe(tsFiles.get(i).getPath())
+                                : 
LoadTsFileStatement.createUnchecked(tsFiles.get(i).getPath()))
                             .setDeleteAfterLoad(isDeleteAfterLoad)
                             .setConvertOnTypeMismatch(isConvertOnTypeMismatch))
                     .orElse(null);
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/sql/ast/LoadTsFile.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/sql/ast/LoadTsFile.java
index 17f95824929..8e46f7cd72a 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/sql/ast/LoadTsFile.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/relational/sql/ast/LoadTsFile.java
@@ -27,6 +27,7 @@ import 
org.apache.iotdb.commons.queryengine.plan.relational.sql.ast.NodeLocation
 import org.apache.iotdb.commons.queryengine.plan.relational.sql.ast.Statement;
 import org.apache.iotdb.db.conf.IoTDBDescriptor;
 import org.apache.iotdb.db.i18n.DataNodeQueryMessages;
+import org.apache.iotdb.db.queryengine.plan.statement.crud.LoadTsFileStatement;
 import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
 import org.apache.iotdb.db.storageengine.load.config.LoadTsFileConfigurator;
 
@@ -74,19 +75,25 @@ public class LoadTsFile extends Statement {
   private boolean needDecode4TimeColumn;
 
   public LoadTsFile(NodeLocation location, String filePath, Map<String, 
String> loadAttributes) {
-    this(location, filePath, loadAttributes, true);
+    this(location, filePath, loadAttributes, true, true);
   }
 
   public static LoadTsFile createUnchecked(
       NodeLocation location, String filePath, Map<String, String> 
loadAttributes) {
-    return new LoadTsFile(location, filePath, loadAttributes, false);
+    return new LoadTsFile(location, filePath, loadAttributes, false, true);
+  }
+
+  public static LoadTsFile createForPipe(
+      NodeLocation location, String filePath, Map<String, String> 
loadAttributes) {
+    return new LoadTsFile(location, filePath, loadAttributes, false, false);
   }
 
   private LoadTsFile(
       NodeLocation location,
       String filePath,
       Map<String, String> loadAttributes,
-      boolean validateSourcePath) {
+      boolean validateSourcePath,
+      boolean validateInternalDataDir) {
     super(location);
     this.filePath =
         requireNonNull(filePath, 
DataNodeQueryMessages.EXCEPTION_FILEPATH_IS_NULL_84CE8A66);
@@ -105,8 +112,9 @@ public class LoadTsFile extends Statement {
 
     try {
       this.tsFiles =
-          
org.apache.iotdb.db.queryengine.plan.statement.crud.LoadTsFileStatement.processTsFile(
-              new File(filePath), validateSourcePath);
+          validateInternalDataDir
+              ? LoadTsFileStatement.processTsFile(new File(filePath), 
validateSourcePath)
+              : LoadTsFileStatement.processTsFileForPipe(new File(filePath));
       this.resources = new ArrayList<>();
       this.writePointCountList = new ArrayList<>();
       this.isTableModel = new 
ArrayList<>(Collections.nCopies(this.tsFiles.size(), true));
@@ -306,7 +314,9 @@ public class LoadTsFile extends Statement {
       final Map<String, String> properties = this.loadAttributes;
 
       final LoadTsFile subStatement =
-          LoadTsFile.createUnchecked(getLocation().orElse(null), filePath, 
properties);
+          isGeneratedByPipe
+              ? LoadTsFile.createForPipe(getLocation().orElse(null), filePath, 
properties)
+              : LoadTsFile.createUnchecked(getLocation().orElse(null), 
filePath, properties);
 
       // Copy all configuration properties
       subStatement.databaseLevel = this.databaseLevel;
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileScheduler.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileScheduler.java
index ef2e948c0e6..da143c6b8cb 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileScheduler.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/scheduler/load/LoadTsFileScheduler.java
@@ -625,7 +625,10 @@ public class LoadTsFileScheduler implements IScheduler {
             failedNode.isTableModel()
                 ? loadTsFileDataTypeConverter
                     .convertForTableModel(
-                        LoadTsFile.createUnchecked(null, filePath, 
Collections.emptyMap())
+                        (isGeneratedByPipe
+                                ? LoadTsFile.createForPipe(null, filePath, 
Collections.emptyMap())
+                                : LoadTsFile.createUnchecked(
+                                    null, filePath, Collections.emptyMap()))
                             .setDatabase(failedNode.getDatabase())
                             .setDeleteAfterLoad(failedNode.isDeleteAfterLoad())
                             .setConvertOnTypeMismatch(true))
@@ -684,7 +687,9 @@ public class LoadTsFileScheduler implements IScheduler {
       final String filePath, final boolean deleteAfterLoad, final String 
database)
       throws FileNotFoundException {
     final LoadTsFileStatement statement =
-        LoadTsFileStatement.createUnchecked(filePath)
+        (isGeneratedByPipe
+                ? LoadTsFileStatement.createForPipe(filePath)
+                : LoadTsFileStatement.createUnchecked(filePath))
             .setDeleteAfterLoad(deleteAfterLoad)
             .setConvertOnTypeMismatch(true);
     if (database != null) {
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatement.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatement.java
index ca823fcb557..0b0dea76311 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatement.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatement.java
@@ -80,14 +80,19 @@ public class LoadTsFileStatement extends Statement {
   private boolean needDecode4TimeColumn;
 
   public LoadTsFileStatement(String filePath) throws FileNotFoundException {
-    this(filePath, true);
+    this(filePath, true, true);
   }
 
   public static LoadTsFileStatement createUnchecked(String filePath) throws 
FileNotFoundException {
-    return new LoadTsFileStatement(filePath, false);
+    return new LoadTsFileStatement(filePath, false, true);
   }
 
-  private LoadTsFileStatement(String filePath, boolean validateSourcePath)
+  public static LoadTsFileStatement createForPipe(String filePath) throws 
FileNotFoundException {
+    return new LoadTsFileStatement(filePath, false, false);
+  }
+
+  private LoadTsFileStatement(
+      String filePath, boolean validateSourcePath, boolean 
validateInternalDataDir)
       throws FileNotFoundException {
     this.file = new File(filePath).getAbsoluteFile();
     this.databaseLevel = 
IoTDBDescriptor.getInstance().getConfig().getDefaultDatabaseLevel();
@@ -99,7 +104,7 @@ public class LoadTsFileStatement extends Statement {
         
IoTDBDescriptor.getInstance().getConfig().getLoadTabletConversionThresholdBytes();
     this.autoCreateDatabase = 
IoTDBDescriptor.getInstance().getConfig().isAutoCreateSchemaEnabled();
 
-    this.tsFiles = processTsFile(file, validateSourcePath);
+    this.tsFiles = processTsFile(file, validateSourcePath, 
validateInternalDataDir);
     this.resources = new ArrayList<>();
     this.writePointCountList = new ArrayList<>();
     this.isTableModel = new 
ArrayList<>(Collections.nCopies(this.tsFiles.size(), false));
@@ -107,11 +112,26 @@ public class LoadTsFileStatement extends Statement {
   }
 
   public static List<File> processTsFile(final File file) throws 
FileNotFoundException {
-    return processTsFile(file, true);
+    return processTsFile(file, true, true);
   }
 
   public static List<File> processTsFile(final File file, final boolean 
validateSourcePath)
       throws FileNotFoundException {
+    return processTsFile(file, validateSourcePath, true);
+  }
+
+  public static List<File> processTsFileForPipe(final File file) throws 
FileNotFoundException {
+    return processTsFile(file, false, false);
+  }
+
+  private static List<File> processTsFile(
+      final File file, final boolean validateSourcePath, final boolean 
validateInternalDataDir)
+      throws FileNotFoundException {
+    final Path[] internalDataDirCanonicalPaths =
+        
IoTDBDescriptor.getInstance().getConfig().getInternalDataDirCanonicalPaths();
+    if (validateInternalDataDir) {
+      validateNotLoadingInternalTsFile(file, internalDataDirCanonicalPaths);
+    }
     if (validateSourcePath) {
       validateLoadSourcePath(file);
     }
@@ -127,7 +147,9 @@ public class LoadTsFileStatement extends Statement {
                     
.QUERY_EXCEPTION_CAN_NOT_FIND_S_ON_THIS_MACHINE_NOTICE_THAT_LOAD_CAN_ONLY_B7886C0E,
                 file.getPath()));
       }
-      tsFiles.addAll(findAllTsFile(file, validateSourcePath));
+      tsFiles.addAll(
+          findAllTsFile(
+              file, validateSourcePath, validateInternalDataDir, 
internalDataDirCanonicalPaths));
     }
     sortTsFiles(tsFiles);
     return tsFiles;
@@ -150,7 +172,11 @@ public class LoadTsFileStatement extends Statement {
     this.statementType = StatementType.MULTI_BATCH_INSERT;
   }
 
-  private static List<File> findAllTsFile(File file, boolean 
validateSourcePath)
+  private static List<File> findAllTsFile(
+      File file,
+      boolean validateSourcePath,
+      boolean validateInternalDataDir,
+      Path[] internalDataDirCanonicalPaths)
       throws FileNotFoundException {
     final File[] files = file.listFiles();
     if (files == null) {
@@ -159,13 +185,21 @@ public class LoadTsFileStatement extends Statement {
 
     final List<File> tsFiles = new ArrayList<>();
     for (File nowFile : files) {
+      if (validateInternalDataDir) {
+        validateNotLoadingInternalTsFile(nowFile, 
internalDataDirCanonicalPaths);
+      }
       if (validateSourcePath) {
         validateLoadSourcePath(nowFile);
       }
       if (nowFile.getName().endsWith(TsFileConstant.TSFILE_SUFFIX)) {
         tsFiles.add(nowFile);
       } else if (nowFile.isDirectory()) {
-        tsFiles.addAll(findAllTsFile(nowFile, validateSourcePath));
+        tsFiles.addAll(
+            findAllTsFile(
+                nowFile,
+                validateSourcePath,
+                validateInternalDataDir,
+                internalDataDirCanonicalPaths));
       }
     }
     return tsFiles;
@@ -199,6 +233,19 @@ public class LoadTsFileStatement extends Statement {
             Arrays.toString(allowedDirs)));
   }
 
+  private static void validateNotLoadingInternalTsFile(
+      final File file, final Path[] internalDataDirCanonicalPaths) throws 
FileNotFoundException {
+    final Path sourcePath = canonicalPath(file);
+    for (final Path internalDataDirCanonicalPath : 
internalDataDirCanonicalPaths) {
+      if (sourcePath.startsWith(internalDataDirCanonicalPath)
+          || internalDataDirCanonicalPath.startsWith(sourcePath)) {
+        throw new FileNotFoundException(
+            DataNodeQueryMessages
+                
.QUERY_EXCEPTION_CANNOT_LOAD_FILES_BECAUSE_SPECIFIED_DIRECTORY_CONTAINS_IOTDB_DATA_B0A1B93D);
+      }
+    }
+  }
+
   private static Path canonicalPath(final File file) throws 
FileNotFoundException {
     try {
       return file.getCanonicalFile().toPath();
@@ -498,7 +545,9 @@ public class LoadTsFileStatement extends Statement {
       loadAttributes.put(PIPE_GENERATED_KEY, String.valueOf(true));
     }
 
-    return LoadTsFile.createUnchecked(null, file.getAbsolutePath(), 
loadAttributes);
+    return isGeneratedByPipe
+        ? LoadTsFile.createForPipe(null, file.getAbsolutePath(), 
loadAttributes)
+        : LoadTsFile.createUnchecked(null, file.getAbsolutePath(), 
loadAttributes);
   }
 
   @Override
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatementTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatementTest.java
index bfebf51d281..fac3227f4b5 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatementTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/queryengine/plan/statement/crud/LoadTsFileStatementTest.java
@@ -109,6 +109,66 @@ public class LoadTsFileStatementTest {
     }
   }
 
+  @Test
+  public void testLoadInternalTsFileIsRejectedWithoutLeakingPath() throws 
Exception {
+    final IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig();
+    final String[][] originalTierDataDirs = config.getTierDataDirs();
+    final boolean originalCheckEnabled = 
config.isLoadTsFileSourcePathCheckEnabled();
+    final Path dataNodeDir = Files.createTempDirectory("load-tsfile-datanode");
+    final Path dataDir = dataNodeDir.resolve("data");
+    final Path internalTsFile =
+        
Files.createDirectories(dataDir.resolve("pipe-hardlink")).resolve("a.tsfile");
+    Files.createFile(internalTsFile);
+
+    try {
+      config.setTierDataDirs(new String[][] {{dataDir.toString()}});
+      config.setLoadTsFileSourcePathCheckEnabled(false);
+
+      try {
+        new LoadTsFileStatement(internalTsFile.toString());
+        Assert.fail("Expected internal IoTDB data directory to be rejected.");
+      } catch (final FileNotFoundException e) {
+        Assert.assertEquals(
+            "Cannot load files because the specified directory contains IoTDB 
data.",
+            e.getMessage());
+        Assert.assertFalse(e.getMessage().contains(dataDir.toString()));
+        Assert.assertFalse(e.getMessage().contains(internalTsFile.toString()));
+      }
+
+      Assert.assertEquals(
+          1, 
LoadTsFileStatement.createForPipe(internalTsFile.toString()).getTsFiles().size());
+    } finally {
+      config.setTierDataDirs(originalTierDataDirs);
+      config.setLoadTsFileSourcePathCheckEnabled(originalCheckEnabled);
+      deleteRecursively(dataNodeDir);
+    }
+  }
+
+  @Test
+  public void testLoadPipeReceiverTsFileOutsideDataDirIsAllowed() throws 
Exception {
+    final IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig();
+    final String[][] originalTierDataDirs = config.getTierDataDirs();
+    final boolean originalCheckEnabled = 
config.isLoadTsFileSourcePathCheckEnabled();
+    final Path dataNodeDir = Files.createTempDirectory("load-tsfile-datanode");
+    final Path dataDir = dataNodeDir.resolve("data");
+    final Path pipeReceiverDir =
+        
Files.createDirectories(dataNodeDir.resolve("system").resolve("pipe").resolve("receiver"));
+    final Path pipeReceiverTsFile = 
Files.createFile(pipeReceiverDir.resolve("a.tsfile"));
+
+    try {
+      config.setTierDataDirs(new String[][] {{dataDir.toString()}});
+      config.setLoadTsFileSourcePathCheckEnabled(false);
+
+      final LoadTsFileStatement statement = new 
LoadTsFileStatement(pipeReceiverTsFile.toString());
+      Assert.assertEquals(1, statement.getTsFiles().size());
+      Assert.assertEquals(pipeReceiverTsFile.toFile(), 
statement.getTsFiles().get(0));
+    } finally {
+      config.setTierDataDirs(originalTierDataDirs);
+      config.setLoadTsFileSourcePathCheckEnabled(originalCheckEnabled);
+      deleteRecursively(dataNodeDir);
+    }
+  }
+
   private static void assertLoadSourcePathRejected(final Path sourcePath) {
     try {
       new LoadTsFileStatement(sourcePath.toString());

Reply via email to