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 b9683c0c106 fix(load): validate TsFile before async move (#18531)
b9683c0c106 is described below

commit b9683c0c10664d90ae85102deaeb554a9aab0806
Author: Zhenyu Luo <[email protected]>
AuthorDate: Mon Aug 31 11:22:57 2026 +0800

    fix(load): validate TsFile before async move (#18531)
    
    * fix(load): validate TsFile before async move
    
    * fix(load): report invalid tsfile in async pipe path
---
 .../apache/iotdb/db/it/IoTDBLoadTsFileAuthIT.java  | 29 +++++++++++
 .../plan/statement/crud/LoadTsFileStatement.java   |  2 +-
 .../iotdb/db/storageengine/load/util/LoadUtil.java | 38 ++++++++++++++
 .../db/storageengine/load/util/LoadUtilTest.java   | 59 ++++++++++++++++++++++
 4 files changed, 127 insertions(+), 1 deletion(-)

diff --git 
a/integration-test/src/test/java/org/apache/iotdb/db/it/IoTDBLoadTsFileAuthIT.java
 
b/integration-test/src/test/java/org/apache/iotdb/db/it/IoTDBLoadTsFileAuthIT.java
index d842941e0cc..0424eb189cb 100644
--- 
a/integration-test/src/test/java/org/apache/iotdb/db/it/IoTDBLoadTsFileAuthIT.java
+++ 
b/integration-test/src/test/java/org/apache/iotdb/db/it/IoTDBLoadTsFileAuthIT.java
@@ -42,6 +42,7 @@ import org.junit.experimental.categories.Category;
 import org.junit.runner.RunWith;
 
 import java.io.File;
+import java.nio.charset.StandardCharsets;
 import java.nio.file.Files;
 import java.sql.Connection;
 import java.sql.ResultSet;
@@ -180,6 +181,34 @@ public class IoTDBLoadTsFileAuthIT {
     assertCountEventually(10, TimeUnit.SECONDS.toMillis(60));
   }
 
+  @Test
+  public void testAsyncLoadShouldNotMoveFileWithoutTsFileSuffix() throws 
Exception {
+    assertAsyncLoadDoesNotMoveInvalidTsFile(new File(tmpDir, 
"not-a-tsfile.txt"), "Can not find");
+  }
+
+  @Test
+  public void testAsyncLoadShouldNotMoveInvalidTsFile() throws Exception {
+    assertAsyncLoadDoesNotMoveInvalidTsFile(new File(tmpDir, 
"invalid.tsfile"), "Loading file");
+  }
+
+  private static void assertAsyncLoadDoesNotMoveInvalidTsFile(
+      final File sourceFile, final String expectedErrorMessage) throws 
Exception {
+    final String originalContent = "ordinary file content";
+    Files.write(sourceFile.toPath(), 
originalContent.getBytes(StandardCharsets.UTF_8));
+
+    assertNonQueryTestFail(
+        String.format(
+            "load \"%s\" with ('async'='true', 'on-success'='delete')",
+            sourceFile.getAbsolutePath()),
+        expectedErrorMessage);
+
+    Assert.assertTrue(
+        "Non-TsFile source must remain in its original location", 
sourceFile.isFile());
+    Assert.assertEquals(
+        originalContent,
+        new String(Files.readAllBytes(sourceFile.toPath()), 
StandardCharsets.UTF_8));
+  }
+
   private static void prepareSchemaAndTsFile(final File tsFile) throws 
Exception {
     prepareSchema(MEASUREMENT.getType());
     generateTsFile(tsFile);
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 ca0a0a3e2bf..cacbe36cdbb 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
@@ -145,7 +145,7 @@ public class LoadTsFileStatement extends Statement {
     }
 
     final List<File> tsFiles = new ArrayList<>();
-    if (file.isFile()) {
+    if (file.isFile() && 
file.getName().endsWith(TsFileConstant.TSFILE_SUFFIX)) {
       tsFiles.add(file);
     } else {
       if (file.listFiles() == null) {
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/util/LoadUtil.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/util/LoadUtil.java
index ea1440bd316..8d653002dbc 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/util/LoadUtil.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/util/LoadUtil.java
@@ -26,6 +26,7 @@ import org.apache.iotdb.commons.utils.FileUtils;
 import org.apache.iotdb.commons.utils.RetryUtils;
 import org.apache.iotdb.db.auth.AuthorityChecker;
 import org.apache.iotdb.db.conf.IoTDBDescriptor;
+import org.apache.iotdb.db.i18n.DataNodeQueryMessages;
 import org.apache.iotdb.db.i18n.StorageEngineMessages;
 import org.apache.iotdb.db.protocol.session.IClientSession;
 import org.apache.iotdb.db.protocol.session.SessionManager;
@@ -35,7 +36,9 @@ import 
org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
 import org.apache.iotdb.db.storageengine.load.active.ActiveLoadPathHelper;
 import org.apache.iotdb.db.storageengine.load.disk.ILoadDiskSelector;
 
+import org.apache.tsfile.common.conf.TSFileConfig;
 import org.apache.tsfile.common.constant.TsFileConstant;
+import org.apache.tsfile.read.TsFileSequenceReader;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -68,6 +71,13 @@ public class LoadUtil {
     }
 
     try {
+      // Validate the complete batch before moving any source file. Otherwise 
a later malformed
+      // file could leave earlier files queued in the active-load directory.
+      for (final File file : tsFiles) {
+        if (file != null && !isValidTsFile(file)) {
+          return false;
+        }
+      }
       for (File file : tsFiles) {
         if (!loadTsFilesToActiveDir(loadAttributes, file, isDeleteAfterLoad)) {
           return false;
@@ -119,6 +129,11 @@ public class LoadUtil {
       return true;
     }
 
+    // Validate before moving the source so ordinary or malformed files remain 
in place.
+    if (!isValidTsFile(file)) {
+      return false;
+    }
+
     final File targetFilePath;
     try {
       targetFilePath =
@@ -148,6 +163,19 @@ public class LoadUtil {
     return true;
   }
 
+  private static boolean isValidTsFile(final File file) {
+    if (!file.isFile() || 
!file.getName().endsWith(TsFileConstant.TSFILE_SUFFIX)) {
+      return false;
+    }
+    try (final TsFileSequenceReader reader =
+        new TsFileSequenceReader(file.getAbsolutePath(), false)) {
+      return TSFileConfig.MAGIC_STRING.equals(reader.readHeadMagic())
+          && TSFileConfig.MAGIC_STRING.equals(reader.readTailMagic());
+    } catch (Exception e) {
+      return false;
+    }
+  }
+
   private static Map<String, String> appendCurrentUserIfAbsent(
       final Map<String, String> loadAttributes) {
     final Map<String, String> attributes =
@@ -196,6 +224,16 @@ public class LoadUtil {
     for (final String file : files) {
       sourceFiles.add(new File(file));
     }
+    // The main TsFile must be valid before any TsFile or sidecar is 
transferred. Pipe acknowledges
+    // this method immediately, so deferring validation to the active loader 
is too late.
+    for (final File sourceFile : sourceFiles) {
+      if (isTsFile(sourceFile) && !isValidTsFile(sourceFile)) {
+        throw new IOException(
+            String.format(
+                
DataNodeQueryMessages.THE_FILE_S_IS_NOT_A_VALID_TSFILE_PLEASE_CHECK_THE_INPUT_FILE,
+                sourceFile.getAbsolutePath()));
+      }
+    }
     sourceFiles.sort(Comparator.comparing(LoadUtil::isTsFile));
     transferFilesToActiveDir(
         targetDir,
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/util/LoadUtilTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/util/LoadUtilTest.java
index bffa2c5578c..5531734e105 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/util/LoadUtilTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/util/LoadUtilTest.java
@@ -19,11 +19,15 @@
 
 package org.apache.iotdb.db.storageengine.load.util;
 
+import org.apache.iotdb.db.conf.IoTDBConfig;
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
+import org.apache.iotdb.db.i18n.DataNodeQueryMessages;
 import 
org.apache.iotdb.db.storageengine.dataregion.modification.ModificationFile;
 import 
org.apache.iotdb.db.storageengine.dataregion.modification.v1.ModificationFileV1;
 import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
 import org.apache.iotdb.db.storageengine.load.active.ActiveLoadPathHelper;
 
+import org.apache.tsfile.write.writer.TsFileIOWriter;
 import org.junit.After;
 import org.junit.Assert;
 import org.junit.Before;
@@ -40,6 +44,8 @@ import java.util.Set;
 
 public class LoadUtilTest {
 
+  private final IoTDBConfig config = IoTDBDescriptor.getInstance().getConfig();
+  private String[] originalListeningDirs;
   private File tempDir;
   private File sourceDir;
   private File targetDir;
@@ -51,13 +57,58 @@ public class LoadUtilTest {
     targetDir = new File(tempDir, "target");
     Assert.assertTrue(sourceDir.mkdirs());
     Assert.assertTrue(targetDir.mkdirs());
+    originalListeningDirs = config.getLoadActiveListeningDirs();
+    config.setLoadActiveListeningDirs(new String[] 
{targetDir.getAbsolutePath()});
+    LoadUtil.updateLoadDiskSelector();
   }
 
   @After
   public void tearDown() {
+    config.setLoadActiveListeningDirs(originalListeningDirs);
+    LoadUtil.updateLoadDiskSelector();
     deleteRecursively(tempDir);
   }
 
+  @Test
+  public void testAsyncLoadValidatesCompleteBatchBeforeTransfer() throws 
Exception {
+    final File validTsFile = createCompletedTsFile("valid.tsfile");
+    final File invalidTsFile = new File(sourceDir, "invalid.tsfile");
+    Files.write(invalidTsFile.toPath(), 
"invalid".getBytes(StandardCharsets.UTF_8));
+
+    Assert.assertFalse(
+        LoadUtil.loadTsFileAsyncToActiveDir(Arrays.asList(validTsFile, 
invalidTsFile), null, true));
+    Assert.assertTrue(validTsFile.exists());
+    Assert.assertTrue(invalidTsFile.exists());
+    Assert.assertEquals(0, targetDir.listFiles().length);
+  }
+
+  @Test
+  public void testPipeAsyncLoadReportsInvalidTsFileAndKeepsSources() throws 
Exception {
+    final File invalidTsFile = new File(sourceDir, "invalid.tsfile");
+    final File resourceFile =
+        new File(invalidTsFile.getAbsolutePath() + 
TsFileResource.RESOURCE_SUFFIX);
+    Files.write(invalidTsFile.toPath(), 
"invalid".getBytes(StandardCharsets.UTF_8));
+    Files.write(resourceFile.toPath(), 
"resource".getBytes(StandardCharsets.UTF_8));
+
+    try {
+      LoadUtil.loadFilesToActiveDir(
+          null,
+          Arrays.asList(resourceFile.getAbsolutePath(), 
invalidTsFile.getAbsolutePath()),
+          true);
+      Assert.fail("Expected invalid TsFile error");
+    } catch (final IOException e) {
+      Assert.assertEquals(
+          String.format(
+              
DataNodeQueryMessages.THE_FILE_S_IS_NOT_A_VALID_TSFILE_PLEASE_CHECK_THE_INPUT_FILE,
+              invalidTsFile.getAbsolutePath()),
+          e.getMessage());
+    }
+
+    Assert.assertTrue(invalidTsFile.exists());
+    Assert.assertTrue(resourceFile.exists());
+    Assert.assertEquals(0, targetDir.listFiles().length);
+  }
+
   @Test
   public void 
testTransferFilesKeepsSameNamedGroupsIsolatedAndDeletesSourcesAfterHandoff()
       throws Exception {
@@ -224,6 +275,14 @@ public class LoadUtilTest {
     return sourceFiles;
   }
 
+  private File createCompletedTsFile(final String fileName) throws Exception {
+    final File tsFile = new File(sourceDir, fileName);
+    try (final TsFileIOWriter writer = new TsFileIOWriter(tsFile)) {
+      writer.endFile();
+    }
+    return tsFile;
+  }
+
   private static void deleteRecursively(final File file) {
     if (file == null || !file.exists()) {
       return;

Reply via email to