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 76ee6f08165 Fix finite pipe auto-drop without DataRegions (#18540)
76ee6f08165 is described below

commit 76ee6f08165ee8fc05f9c410f8310b6717e8394c
Author: Zhenyu Luo <[email protected]>
AuthorDate: Mon Aug 31 11:27:02 2026 +0800

    Fix finite pipe auto-drop without DataRegions (#18540)
    
    Allow finite internal pipes to complete after an explicit empty DataRegion 
report while preserving realtime and external-source pipes.
---
 .../auto/enhanced/IoTDBPipeAutoDropIT.java         | 81 ++++++++++++++++++++++
 .../pipe/it/single/IoTDBPipePermissionIT.java      | 66 ++++++++++++++++++
 .../runtime/heartbeat/PipeHeartbeatParser.java     | 17 +++--
 .../runtime/heartbeat/PipeHeartbeatParserTest.java | 76 ++++++++++++++++++--
 4 files changed, 231 insertions(+), 9 deletions(-)

diff --git 
a/integration-test/src/test/java/org/apache/iotdb/pipe/it/dual/treemodel/auto/enhanced/IoTDBPipeAutoDropIT.java
 
b/integration-test/src/test/java/org/apache/iotdb/pipe/it/dual/treemodel/auto/enhanced/IoTDBPipeAutoDropIT.java
index 28928653ef9..2d1d7d51f33 100644
--- 
a/integration-test/src/test/java/org/apache/iotdb/pipe/it/dual/treemodel/auto/enhanced/IoTDBPipeAutoDropIT.java
+++ 
b/integration-test/src/test/java/org/apache/iotdb/pipe/it/dual/treemodel/auto/enhanced/IoTDBPipeAutoDropIT.java
@@ -137,6 +137,87 @@ public class IoTDBPipeAutoDropIT extends 
AbstractPipeDualTreeModelAutoIT {
     }
   }
 
+  @Test
+  public void testAutoDropFinitePipesWithoutDataRegion() throws Exception {
+    final DataNodeWrapper receiverDataNode = receiverEnv.getDataNodeWrapper(0);
+    final Map<String, String> sinkAttributes = new HashMap<>();
+    sinkAttributes.put("sink", "iotdb-thrift-sink");
+    sinkAttributes.put("sink.batch.enable", "false");
+    sinkAttributes.put("sink.ip", receiverDataNode.getIp());
+    sinkAttributes.put("sink.port", 
Integer.toString(receiverDataNode.getPort()));
+
+    try (final SyncConfigNodeIServiceClient client =
+        (SyncConfigNodeIServiceClient) 
senderEnv.getLeaderConfigNodeConnection()) {
+      final Map<String, String> querySourceAttributes = new HashMap<>();
+      querySourceAttributes.put("source.mode", "query");
+      querySourceAttributes.put("user", SessionConfig.DEFAULT_USER);
+      Assert.assertEquals(
+          TSStatusCode.SUCCESS_STATUS.getStatusCode(),
+          client
+              .createPipe(
+                  new TCreatePipeReq("query_pipe_without_data_region", 
sinkAttributes)
+                      .setExtractorAttributes(querySourceAttributes))
+              .getCode());
+
+      final Map<String, String> historySourceAttributes = new HashMap<>();
+      historySourceAttributes.put("source.realtime.enable", 
Boolean.FALSE.toString());
+      historySourceAttributes.put("user", SessionConfig.DEFAULT_USER);
+      Assert.assertEquals(
+          TSStatusCode.SUCCESS_STATUS.getStatusCode(),
+          client
+              .createPipe(
+                  new TCreatePipeReq("history_pipe_without_data_region", 
sinkAttributes)
+                      .setExtractorAttributes(historySourceAttributes))
+              .getCode());
+
+      final Map<String, String> realtimeSourceAttributes = new HashMap<>();
+      realtimeSourceAttributes.put("source.history.enable", 
Boolean.FALSE.toString());
+      realtimeSourceAttributes.put("user", SessionConfig.DEFAULT_USER);
+      Assert.assertEquals(
+          TSStatusCode.SUCCESS_STATUS.getStatusCode(),
+          client
+              .createPipe(
+                  new TCreatePipeReq("realtime_pipe_without_data_region", 
sinkAttributes)
+                      .setExtractorAttributes(realtimeSourceAttributes))
+              .getCode());
+
+      Assert.assertEquals(
+          TSStatusCode.SUCCESS_STATUS.getStatusCode(),
+          client.startPipe("query_pipe_without_data_region").getCode());
+      Assert.assertEquals(
+          TSStatusCode.SUCCESS_STATUS.getStatusCode(),
+          client.startPipe("history_pipe_without_data_region").getCode());
+      Assert.assertEquals(
+          TSStatusCode.SUCCESS_STATUS.getStatusCode(),
+          client.startPipe("realtime_pipe_without_data_region").getCode());
+
+      await()
+          .pollInSameThread()
+          .pollInterval(1L, TimeUnit.SECONDS)
+          .atMost(600, TimeUnit.SECONDS)
+          .untilAsserted(
+              () -> {
+                final List<TShowPipeInfo> pipeInfoList =
+                    client.showPipe(new 
TShowPipeReq().setUserName(SessionConfig.DEFAULT_USER))
+                        .pipeInfoList;
+                Assert.assertFalse(
+                    pipeInfoList.stream()
+                        .anyMatch(info -> 
info.getId().equals("query_pipe_without_data_region")));
+                Assert.assertFalse(
+                    pipeInfoList.stream()
+                        .anyMatch(info -> 
info.getId().equals("history_pipe_without_data_region")));
+                Assert.assertTrue(
+                    pipeInfoList.stream()
+                        .anyMatch(
+                            info -> 
info.getId().equals("realtime_pipe_without_data_region")));
+              });
+
+      Assert.assertEquals(
+          TSStatusCode.SUCCESS_STATUS.getStatusCode(),
+          client.dropPipe("realtime_pipe_without_data_region").getCode());
+    }
+  }
+
   @Test
   public void testAutoDropIgnoredUnmatchedDataRegions() throws Exception {
     final DataNodeWrapper receiverDataNode = receiverEnv.getDataNodeWrapper(0);
diff --git 
a/integration-test/src/test/java/org/apache/iotdb/pipe/it/single/IoTDBPipePermissionIT.java
 
b/integration-test/src/test/java/org/apache/iotdb/pipe/it/single/IoTDBPipePermissionIT.java
index 45fa762b02b..96870b0329b 100644
--- 
a/integration-test/src/test/java/org/apache/iotdb/pipe/it/single/IoTDBPipePermissionIT.java
+++ 
b/integration-test/src/test/java/org/apache/iotdb/pipe/it/single/IoTDBPipePermissionIT.java
@@ -37,7 +37,11 @@ import java.sql.ResultSet;
 import java.sql.SQLException;
 import java.sql.Statement;
 import java.util.Arrays;
+import java.util.HashSet;
+import java.util.Set;
+import java.util.concurrent.TimeUnit;
 
+import static org.awaitility.Awaitility.await;
 import static org.junit.Assert.fail;
 
 @RunWith(IoTDBTestRunner.class)
@@ -257,4 +261,66 @@ public class IoTDBPipePermissionIT extends 
AbstractPipeSingleIT {
     TableModelUtils.assertCountData("test", "test", 0, env);
     TableModelUtils.assertCountData("test", "test1", 100, env);
   }
+
+  @Test
+  public void testAuthenticatedRealtimePipesRemainWithoutDataRegion() throws 
Exception {
+    final Set<String> expectedPipeNames =
+        new HashSet<>(Arrays.asList("root_pipe", "source_pipe", "sink_pipe", 
"source_sink_pipe"));
+
+    try (final Connection connection = 
env.getConnection(BaseEnv.TABLE_SQL_DIALECT);
+        final Statement statement = connection.createStatement()) {
+      statement.execute("CREATE USER user_source 'paSs1234@56789'");
+      statement.execute("CREATE USER user_sink 'paSs1234@56789'");
+      statement.execute("CREATE DATABASE test_pipe_authentication");
+
+      statement.execute(
+          "create pipe root_pipe "
+              + "with source ('forwarding-pipe-requests'='false', "
+              + "'database-name'='test_pipe_authentication', 
'table-name'='table_0') "
+              + "with processor ('processor'='rename-database-processor', "
+              + "'new-db-name'='pipe_newDB1') "
+              + "with sink ('sink'='write-back-sink')");
+      statement.execute(
+          "create pipe source_pipe "
+              + "with source ('forwarding-pipe-requests'='false', "
+              + "'database-name'='test_pipe_authentication', 
'table-name'='table_0', "
+              + "'user'='user_source', 'password'='paSs1234@56789') "
+              + "with processor ('processor'='rename-database-processor', "
+              + "'new-db-name'='pipe_newDB2') "
+              + "with sink ('sink'='write-back-sink')");
+      statement.execute(
+          "create pipe sink_pipe "
+              + "with source ('forwarding-pipe-requests'='false', "
+              + "'database-name'='test_pipe_authentication', 
'table-name'='table_0') "
+              + "with processor ('processor'='rename-database-processor', "
+              + "'new-db-name'='pipe_newDB3') "
+              + "with sink ('sink'='write-back-sink', 'user'='user_sink', "
+              + "'password'='paSs1234@56789')");
+      statement.execute(
+          "create pipe source_sink_pipe "
+              + "with source ('forwarding-pipe-requests'='false', "
+              + "'database-name'='test_pipe_authentication', 
'table-name'='table_0', "
+              + "'user'='user_source', 'password'='paSs1234@56789') "
+              + "with processor ('processor'='rename-database-processor', "
+              + "'new-db-name'='pipe_newDB4') "
+              + "with sink ('sink'='write-back-sink', 'user'='user_sink', "
+              + "'password'='paSs1234@56789')");
+
+      await()
+          .pollInSameThread()
+          .pollInterval(1, TimeUnit.SECONDS)
+          .during(15, TimeUnit.SECONDS)
+          .atMost(30, TimeUnit.SECONDS)
+          .untilAsserted(
+              () -> {
+                final Set<String> actualPipeNames = new HashSet<>();
+                try (final ResultSet resultSet = statement.executeQuery("SHOW 
PIPES")) {
+                  while (resultSet.next()) {
+                    actualPipeNames.add(resultSet.getString("ID"));
+                  }
+                }
+                Assert.assertEquals(expectedPipeNames, actualPipeNames);
+              });
+    }
+  }
 }
diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatParser.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatParser.java
index db3d67e86c3..f702312e5dc 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatParser.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatParser.java
@@ -25,6 +25,7 @@ import org.apache.iotdb.commons.consensus.index.ProgressIndex;
 import org.apache.iotdb.commons.exception.pipe.PipeRuntimeCriticalException;
 import org.apache.iotdb.commons.exception.pipe.PipeRuntimeException;
 import 
org.apache.iotdb.commons.exception.pipe.PipeRuntimeSinkCriticalException;
+import org.apache.iotdb.commons.pipe.agent.task.PipeTaskAgent;
 import org.apache.iotdb.commons.pipe.agent.task.meta.PipeMeta;
 import org.apache.iotdb.commons.pipe.agent.task.meta.PipeRuntimeMeta;
 import org.apache.iotdb.commons.pipe.agent.task.meta.PipeStaticMeta;
@@ -178,10 +179,18 @@ public class PipeHeartbeatParser {
       final Set<Integer> requiredDataRegionIds =
           collectRequiredDataRegionIds(pipeMetaFromCoordinator);
 
-      // Remove completed pipes only when every required DataRegion has been 
reported complete.
-      // Relying on the region-level reports (instead of the DataNode-level 
boolean) prevents a
-      // leader-change / task-creation failure from being treated as a 
successful snapshot transfer.
-      if (!requiredDataRegionIds.isEmpty()
+      // A history-only internal pipe is finite and may be removed when all 
required DataRegions
+      // complete, or when CN determines that no DataRegion matched at 
creation time. An explicit
+      // region-level report proves that a DataNode has received the pipe 
meta, preventing an empty
+      // task map from completing the pipe before its initial push. Realtime 
and external-source
+      // pipes must remain alive because they may receive work in the future.
+      final boolean isFiniteInternalPipe =
+          !staticMeta.isSourceExternal()
+              && 
PipeTaskAgent.isHistoryOnlyPipe(staticMeta.getSourceParameters());
+      final boolean hasReliableDataRegionReport =
+          pipeHeartbeat.hasCompletedDataRegionReport(staticMeta);
+      if (isFiniteInternalPipe
+          && hasReliableDataRegionReport
           && 
temporaryMeta.getCompletedDataRegionIds().containsAll(requiredDataRegionIds)) {
         PipeLogger.log(
             LOGGER::info,
diff --git 
a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatParserTest.java
 
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatParserTest.java
index 95c29a56336..c36a09a7def 100644
--- 
a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatParserTest.java
+++ 
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/pipe/coordinator/runtime/heartbeat/PipeHeartbeatParserTest.java
@@ -369,7 +369,7 @@ public class PipeHeartbeatParserTest {
     
CommonDescriptor.getInstance().getConfig().setSeperatedPipeHeartbeatEnabled(false);
 
     final PipeTaskInfo pipeTaskInfo = new PipeTaskInfo();
-    final PipeMeta pipeMeta = createPipeMeta();
+    final PipeMeta pipeMeta = createHistoryOnlyPipeMeta(1);
     pipeTaskInfo.createPipe(
         new CreatePipePlanV2(pipeMeta.getStaticMeta(), 
pipeMeta.getRuntimeMeta()));
 
@@ -394,12 +394,67 @@ public class PipeHeartbeatParserTest {
     verify(context.procedureManager, 
never()).pipeHandleMetaChange(anyBoolean(), anyBoolean());
   }
 
+  @Test
+  public void 
testParseHeartbeatCompletesHistoryOnlyPipeWithoutRequiredDataRegion()
+      throws Exception {
+    
CommonDescriptor.getInstance().getConfig().setSeperatedPipeHeartbeatEnabled(false);
+
+    final Map<String, String> sourceAttributes = new HashMap<>();
+    sourceAttributes.put("source.realtime.enable", Boolean.FALSE.toString());
+    final PipeTaskInfo pipeTaskInfo = new PipeTaskInfo();
+    final PipeMeta pipeMeta = createPipeMeta(sourceAttributes);
+    pipeTaskInfo.createPipe(
+        new CreatePipePlanV2(pipeMeta.getStaticMeta(), 
pipeMeta.getRuntimeMeta()));
+
+    final ParserTestContext context = createParserTestContext(1, pipeTaskInfo);
+    context.parser.parseHeartbeat(
+        1, createPipeHeartbeatWithCompletedRegions(pipeMeta, false, 
Collections.emptyList()));
+
+    Assert.assertNull(pipeTaskInfo.getPipeMetaByPipeName("test_pipe"));
+    verify(context.procedureManager, times(1)).pipeHandleMetaChange(true, 
true);
+  }
+
+  @Test
+  public void testParseHeartbeatKeepsHistoryOnlyPipeWithoutDataRegionReport() 
throws Exception {
+    
CommonDescriptor.getInstance().getConfig().setSeperatedPipeHeartbeatEnabled(false);
+
+    final Map<String, String> sourceAttributes = new HashMap<>();
+    sourceAttributes.put("source.realtime.enable", Boolean.FALSE.toString());
+    final PipeTaskInfo pipeTaskInfo = new PipeTaskInfo();
+    final PipeMeta pipeMeta = createPipeMeta(sourceAttributes);
+    pipeTaskInfo.createPipe(
+        new CreatePipePlanV2(pipeMeta.getStaticMeta(), 
pipeMeta.getRuntimeMeta()));
+
+    final ParserTestContext context = createParserTestContext(1, pipeTaskInfo);
+    context.parser.parseHeartbeat(1, createPipeHeartbeat(pipeMeta, false));
+
+    Assert.assertNotNull(pipeTaskInfo.getPipeMetaByPipeName("test_pipe"));
+    verify(context.procedureManager, 
never()).pipeHandleMetaChange(anyBoolean(), anyBoolean());
+  }
+
+  @Test
+  public void testParseHeartbeatKeepsRealtimePipeWithoutRequiredDataRegion() 
throws Exception {
+    
CommonDescriptor.getInstance().getConfig().setSeperatedPipeHeartbeatEnabled(false);
+
+    final PipeTaskInfo pipeTaskInfo = new PipeTaskInfo();
+    final PipeMeta pipeMeta = createPipeMeta(Collections.emptyMap());
+    pipeTaskInfo.createPipe(
+        new CreatePipePlanV2(pipeMeta.getStaticMeta(), 
pipeMeta.getRuntimeMeta()));
+
+    final ParserTestContext context = createParserTestContext(1, pipeTaskInfo);
+    context.parser.parseHeartbeat(
+        1, createPipeHeartbeatWithCompletedRegions(pipeMeta, false, 
Collections.emptyList()));
+
+    Assert.assertNotNull(pipeTaskInfo.getPipeMetaByPipeName("test_pipe"));
+    verify(context.procedureManager, 
never()).pipeHandleMetaChange(anyBoolean(), anyBoolean());
+  }
+
   @Test
   public void testParseHeartbeatDoesNotTrustDataNodeBooleanForCompletion() 
throws Exception {
     
CommonDescriptor.getInstance().getConfig().setSeperatedPipeHeartbeatEnabled(false);
 
     final PipeTaskInfo pipeTaskInfo = new PipeTaskInfo();
-    final PipeMeta pipeMeta = createPipeMeta(1);
+    final PipeMeta pipeMeta = createHistoryOnlyPipeMeta(1);
     pipeTaskInfo.createPipe(
         new CreatePipePlanV2(pipeMeta.getStaticMeta(), 
pipeMeta.getRuntimeMeta()));
 
@@ -418,7 +473,7 @@ public class PipeHeartbeatParserTest {
     
CommonDescriptor.getInstance().getConfig().setSeperatedPipeHeartbeatEnabled(false);
 
     final PipeTaskInfo pipeTaskInfo = new PipeTaskInfo();
-    final PipeMeta pipeMeta = createPipeMeta(1, 2);
+    final PipeMeta pipeMeta = createHistoryOnlyPipeMeta(1, 2);
     
pipeMeta.getRuntimeMeta().getConsensusGroupId2TaskMetaMap().get(2).setLeaderNodeId(2);
     pipeTaskInfo.createPipe(
         new CreatePipePlanV2(pipeMeta.getStaticMeta(), 
pipeMeta.getRuntimeMeta()));
@@ -442,7 +497,7 @@ public class PipeHeartbeatParserTest {
     
CommonDescriptor.getInstance().getConfig().setSeperatedPipeHeartbeatEnabled(false);
 
     final PipeTaskInfo pipeTaskInfo = new PipeTaskInfo();
-    final PipeMeta pipeMeta = createPipeMeta(1, 2);
+    final PipeMeta pipeMeta = createHistoryOnlyPipeMeta(1, 2);
     
pipeMeta.getRuntimeMeta().getConsensusGroupId2TaskMetaMap().get(2).setLeaderNodeId(2);
     pipeTaskInfo.createPipe(
         new CreatePipePlanV2(pipeMeta.getStaticMeta(), 
pipeMeta.getRuntimeMeta()));
@@ -577,6 +632,17 @@ public class PipeHeartbeatParserTest {
   }
 
   private PipeMeta createPipeMeta(final int... regionIds) {
+    return createPipeMeta(Collections.emptyMap(), regionIds);
+  }
+
+  private PipeMeta createHistoryOnlyPipeMeta(final int... regionIds) {
+    final Map<String, String> sourceAttributes = new HashMap<>();
+    sourceAttributes.put("source.realtime.enable", Boolean.FALSE.toString());
+    return createPipeMeta(sourceAttributes, regionIds);
+  }
+
+  private PipeMeta createPipeMeta(
+      final Map<String, String> sourceAttributes, final int... regionIds) {
     final PipeRuntimeMeta pipeRuntimeMeta = new PipeRuntimeMeta();
     for (final int regionId : regionIds) {
       pipeRuntimeMeta
@@ -584,7 +650,7 @@ public class PipeHeartbeatParserTest {
           .put(regionId, new PipeTaskMeta(MinimumProgressIndex.INSTANCE, 1));
     }
     return new PipeMeta(
-        new PipeStaticMeta("test_pipe", 1L, new HashMap<>(), new HashMap<>(), 
new HashMap<>()),
+        new PipeStaticMeta("test_pipe", 1L, sourceAttributes, new HashMap<>(), 
new HashMap<>()),
         pipeRuntimeMeta);
   }
 

Reply via email to