This is an automated email from the ASF dual-hosted git repository.

Caideyipi 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 cec4f9b111b Fix audit database filtering for consensus pipes (#18476)
cec4f9b111b is described below

commit cec4f9b111bccf5302f7ac1e375c7401f425ef8c
Author: Caideyipi <[email protected]>
AuthorDate: Mon Aug 17 09:47:54 2026 +0800

    Fix audit database filtering for consensus pipes (#18476)
    
    * Fix audit database filtering for consensus pipes
    
    * Keep audit data excluded from user pipes
---
 .../impl/pipe/AbstractOperatePipeProcedureV2.java  |  5 ++++-
 .../db/pipe/agent/task/PipeDataNodeTaskAgent.java  | 25 +++++++++++++++-------
 .../agent/task/builder/PipeDataNodeBuilder.java    |  2 +-
 .../dataregion/DataRegionListeningFilter.java      | 13 +++++++----
 .../pipe/agent/task/PipeDataNodeTaskAgentTest.java | 23 ++++++++++++++++++++
 .../dataregion/DataRegionListeningFilterTest.java  | 25 +++++++++++++++++-----
 6 files changed, 74 insertions(+), 19 deletions(-)

diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/AbstractOperatePipeProcedureV2.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/AbstractOperatePipeProcedureV2.java
index 88cd54dbe28..0710fe345b8 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/AbstractOperatePipeProcedureV2.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/pipe/AbstractOperatePipeProcedureV2.java
@@ -908,7 +908,10 @@ public abstract class AbstractOperatePipeProcedureV2
 
               try {
                 return !DataRegionListeningFilter.shouldDatabaseBeListened(
-                    copiedPipeMeta.getStaticMeta().getSourceParameters(), 
isTableModel, database);
+                    copiedPipeMeta.getStaticMeta().getSourceParameters(),
+                    isTableModel,
+                    database,
+                    copiedPipeMeta.getStaticMeta().getPipeType());
               } catch (final Exception e) {
                 return false;
               }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java
index c47c54f675a..5241ea0c916 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java
@@ -161,7 +161,7 @@ public class PipeDataNodeTaskAgent extends PipeTaskAgent {
       final boolean needConstructDataRegionTask =
           
StorageEngine.getInstance().getAllDataRegionIds().contains(dataRegionId)
               && DataRegionListeningFilter.shouldDataRegionBeListened(
-                  sourceParameters, dataRegionId);
+                  sourceParameters, dataRegionId, 
pipeStaticMeta.getPipeType());
       final boolean needConstructSchemaRegionTask =
           SchemaEngine.getInstance()
                   .getAllSchemaRegionIds()
@@ -861,12 +861,21 @@ public class PipeDataNodeTaskAgent extends PipeTaskAgent {
         throw new PipeException(DataNodePipeMessages.PIPE_META_NOT_FOUND + 
pipeName);
       }
 
-      return pipeMetaKeeper
-          .getPipeMeta(pipeName)
-          .getRuntimeMeta()
-          .getConsensusGroupId2TaskMetaMap()
-          .get(consensusGroupId)
-          .getProgressIndex();
+      final PipeTaskMeta pipeTaskMeta =
+          pipeMetaKeeper
+              .getPipeMeta(pipeName)
+              .getRuntimeMeta()
+              .getConsensusGroupId2TaskMetaMap()
+              .get(consensusGroupId);
+      if (pipeTaskMeta == null) {
+        throw new PipeException(
+            String.format(
+                DataNodePipeMessages
+                    
.PIPE_EXCEPTION_FAILED_TO_GET_PIPE_TASK_PROGRESS_INDEX_WITH_PIPE_NAME_S_CFE9DE7C,
+                pipeName,
+                consensusGroupId));
+      }
+      return pipeTaskMeta.getProgressIndex();
     } finally {
       releaseReadLock();
     }
@@ -951,7 +960,7 @@ public class PipeDataNodeTaskAgent extends PipeTaskAgent {
         final boolean needConstructDataRegionTask =
             dataRegionIds.contains(dataRegionId)
                 && DataRegionListeningFilter.shouldDataRegionBeListened(
-                    sourceParameters, dataRegionId);
+                    sourceParameters, dataRegionId, 
pipeStaticMeta.getPipeType());
         final boolean needConstructSchemaRegionTask =
             schemaRegionIds.contains(new SchemaRegionId(consensusGroupId))
                 && SchemaRegionListeningFilter.shouldSchemaRegionBeListened(
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/builder/PipeDataNodeBuilder.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/builder/PipeDataNodeBuilder.java
index 46a10135d88..2896ff6039a 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/builder/PipeDataNodeBuilder.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/builder/PipeDataNodeBuilder.java
@@ -68,7 +68,7 @@ public class PipeDataNodeBuilder {
         final boolean needConstructDataRegionTask =
             dataRegionIds.contains(dataRegionId)
                 && DataRegionListeningFilter.shouldDataRegionBeListened(
-                    sourceParameters, dataRegionId);
+                    sourceParameters, dataRegionId, 
pipeStaticMeta.getPipeType());
         final boolean needConstructSchemaRegionTask =
             schemaRegionIds.contains(new SchemaRegionId(consensusGroupId))
                 && SchemaRegionListeningFilter.shouldSchemaRegionBeListened(
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/DataRegionListeningFilter.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/DataRegionListeningFilter.java
index ad4c1fddff3..bbaf9ff9745 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/DataRegionListeningFilter.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/DataRegionListeningFilter.java
@@ -23,6 +23,7 @@ import org.apache.iotdb.commons.consensus.DataRegionId;
 import org.apache.iotdb.commons.exception.IllegalPathException;
 import org.apache.iotdb.commons.path.PartialPath;
 import org.apache.iotdb.commons.pipe.agent.task.PipeTask;
+import org.apache.iotdb.commons.pipe.agent.task.meta.PipeType;
 import org.apache.iotdb.commons.pipe.datastructure.pattern.TablePattern;
 import org.apache.iotdb.commons.pipe.datastructure.pattern.TreePattern;
 import org.apache.iotdb.db.storageengine.StorageEngine;
@@ -58,9 +59,12 @@ public class DataRegionListeningFilter {
   }
 
   public static boolean shouldDatabaseBeListened(
-      final PipeParameters parameters, final boolean isTableModel, final 
String databaseRawName)
+      final PipeParameters parameters,
+      final boolean isTableModel,
+      final String databaseRawName,
+      final PipeType pipeType)
       throws IllegalPathException {
-    if (isAuditDatabase(databaseRawName)) {
+    if (!PipeType.CONSENSUS.equals(pipeType) && 
isAuditDatabase(databaseRawName)) {
       return false;
     }
 
@@ -90,7 +94,8 @@ public class DataRegionListeningFilter {
   }
 
   public static boolean shouldDataRegionBeListened(
-      PipeParameters parameters, DataRegionId dataRegionId) throws 
IllegalPathException {
+      final PipeParameters parameters, final DataRegionId dataRegionId, final 
PipeType pipeType)
+      throws IllegalPathException {
     final Pair<Boolean, Boolean> insertionDeletionListeningOptionPair =
         parseInsertionDeletionListeningOptionPair(parameters);
     final boolean hasSpecificListeningOption =
@@ -106,7 +111,7 @@ public class DataRegionListeningFilter {
     }
 
     final String databaseRawName = dataRegion.getDatabaseName();
-    if (isAuditDatabase(databaseRawName)) {
+    if (!PipeType.CONSENSUS.equals(pipeType) && 
isAuditDatabase(databaseRawName)) {
       return false;
     }
 
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgentTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgentTest.java
index 3dd93e87789..018ec371d91 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgentTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgentTest.java
@@ -23,7 +23,9 @@ import org.apache.iotdb.commons.conf.CommonDescriptor;
 import org.apache.iotdb.commons.consensus.index.ProgressIndex;
 import org.apache.iotdb.commons.consensus.index.impl.MinimumProgressIndex;
 import org.apache.iotdb.commons.consensus.index.impl.SimpleProgressIndex;
+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.PipeMetaKeeper;
 import org.apache.iotdb.commons.pipe.agent.task.meta.PipeRuntimeMeta;
 import org.apache.iotdb.commons.pipe.agent.task.meta.PipeStaticMeta;
 import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTaskMeta;
@@ -34,6 +36,7 @@ import org.apache.iotdb.pipe.api.exception.PipeException;
 import org.junit.Assert;
 import org.junit.Test;
 
+import java.lang.reflect.Field;
 import java.util.HashMap;
 import java.util.Map;
 import java.util.concurrent.ConcurrentHashMap;
@@ -44,6 +47,26 @@ public class PipeDataNodeTaskAgentTest {
   private static final int LOCAL_NODE_ID = 1;
   private static final int REGION_ID = 7;
 
+  @Test
+  public void testGetPipeTaskProgressIndexReportsMissingTaskMeta() throws 
Exception {
+    final PipeDataNodeTaskAgent taskAgent = new PipeDataNodeTaskAgent();
+    final Field pipeMetaKeeperField = 
PipeTaskAgent.class.getDeclaredField("pipeMetaKeeper");
+    pipeMetaKeeperField.setAccessible(true);
+    final PipeMetaKeeper pipeMetaKeeper = (PipeMetaKeeper) 
pipeMetaKeeperField.get(taskAgent);
+
+    final String pipeName = PipeStaticMeta.CONSENSUS_PIPE_PREFIX + 
"DataRegion[7]_1_2";
+    pipeMetaKeeper.addPipeMeta(
+        new PipeMeta(
+            new PipeStaticMeta(pipeName, 1L, new HashMap<>(), new HashMap<>(), 
new HashMap<>()),
+            new PipeRuntimeMeta()));
+
+    final PipeException exception =
+        Assert.assertThrows(
+            PipeException.class, () -> 
taskAgent.getPipeTaskProgressIndex(pipeName, REGION_ID));
+    Assert.assertTrue(exception.getMessage().contains(pipeName));
+    
Assert.assertTrue(exception.getMessage().contains(String.valueOf(REGION_ID)));
+  }
+
   @Test
   public void testCreateMemoryCheckStillRunsWhenNoPipeTasksNeedToBeCreated() 
throws Exception {
     final boolean originalPipeEnableMemoryCheck =
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/DataRegionListeningFilterTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/DataRegionListeningFilterTest.java
index 4901afcf912..4257b032ef4 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/DataRegionListeningFilterTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/DataRegionListeningFilterTest.java
@@ -19,6 +19,7 @@
 
 package org.apache.iotdb.db.pipe.source.dataregion;
 
+import org.apache.iotdb.commons.pipe.agent.task.meta.PipeType;
 import org.apache.iotdb.commons.pipe.config.constant.SystemConstant;
 import org.apache.iotdb.commons.subscription.meta.topic.TopicMeta;
 import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameters;
@@ -34,17 +35,31 @@ import static org.junit.Assert.assertTrue;
 public class DataRegionListeningFilterTest {
 
   @Test
-  public void testAuditDatabaseIsNeverListened() throws Exception {
+  public void testAuditDatabaseIsOnlyListenedByConsensusPipes() throws 
Exception {
     final Map<String, String> topicAttributes = new HashMap<>();
     topicAttributes.put(SystemConstant.SQL_DIALECT_KEY, 
SystemConstant.SQL_DIALECT_TABLE_VALUE);
     final PipeParameters parameters =
         new PipeParameters(
             new TopicMeta("topic", 1, 
topicAttributes).generateExtractorAttributes("root"));
 
-    assertFalse(DataRegionListeningFilter.shouldDatabaseBeListened(parameters, 
true, "__audit"));
     assertFalse(
-        DataRegionListeningFilter.shouldDatabaseBeListened(parameters, false, 
"root.__audit"));
-    assertFalse(DataRegionListeningFilter.shouldDatabaseBeListened(parameters, 
true, "__AUDIT"));
-    assertTrue(DataRegionListeningFilter.shouldDatabaseBeListened(parameters, 
true, "user_db"));
+        DataRegionListeningFilter.shouldDatabaseBeListened(
+            parameters, true, "__audit", PipeType.SUBSCRIPTION));
+    assertFalse(
+        DataRegionListeningFilter.shouldDatabaseBeListened(
+            parameters, false, "root.__audit", PipeType.SUBSCRIPTION));
+    assertFalse(
+        DataRegionListeningFilter.shouldDatabaseBeListened(
+            parameters, true, "__AUDIT", PipeType.SUBSCRIPTION));
+    assertTrue(
+        DataRegionListeningFilter.shouldDatabaseBeListened(
+            parameters, true, "user_db", PipeType.SUBSCRIPTION));
+
+    assertTrue(
+        DataRegionListeningFilter.shouldDatabaseBeListened(
+            parameters, true, "__audit", PipeType.CONSENSUS));
+    assertFalse(
+        DataRegionListeningFilter.shouldDatabaseBeListened(
+            parameters, true, "__audit", PipeType.USER));
   }
 }

Reply via email to