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 6559aa52ba9 Pipe: Handle hybrid meta progress indexes (#18331)
6559aa52ba9 is described below

commit 6559aa52ba9fd09dae82be45060ddab81839730b
Author: Caideyipi <[email protected]>
AuthorDate: Fri Jul 31 10:23:59 2026 +0800

    Pipe: Handle hybrid meta progress indexes (#18331)
    
    * Pipe: Handle hybrid meta progress indexes
    
    * Pipe: Handle hybrid meta progress indexes
    
    * Fix getter
    
    * Pipe: Delegate progress index extraction to subclasses
---
 .../pipe/agent/task/PipeConfigNodeTaskAgent.java   |  12 +-
 .../apache/iotdb/db/i18n/DataNodePipeMessages.java |   2 +
 .../apache/iotdb/db/i18n/DataNodePipeMessages.java |   2 +
 .../db/pipe/agent/task/PipeDataNodeTaskAgent.java  |  10 +-
 .../processor/aggregate/AggregateProcessor.java    |  16 +-
 .../twostage/plugin/TwoStageCountProcessor.java    |  29 +++-
 ...istoricalDataRegionTsFileAndDeletionSource.java |  17 +-
 .../plugin/TwoStageCountProcessorTest.java         |  55 +++++++
 .../TsFileResourceProgressIndexTest.java           |   9 ++
 .../apache/iotdb/commons/i18n/PipeMessages.java    |   2 +
 .../apache/iotdb/commons/i18n/PipeMessages.java    |   2 +
 .../commons/consensus/index/ProgressIndex.java     |  10 ++
 .../consensus/index/impl/HybridProgressIndex.java  |  25 +++
 .../consensus/index/impl/IoTProgressIndex.java     |   9 ++
 .../consensus/index/impl/MetaProgressIndex.java    |   9 ++
 .../consensus/index/impl/MinimumProgressIndex.java |   9 ++
 .../consensus/index/impl/RecoverProgressIndex.java |   9 ++
 .../consensus/index/impl/SimpleProgressIndex.java  |   9 ++
 .../consensus/index/impl/StateProgressIndex.java   |   9 ++
 .../index/impl/TimePartitionProgressIndex.java     |   9 ++
 .../index/impl/TimeWindowStateProgressIndex.java   |  11 +-
 .../pipe/source/IoTDBNonDataRegionSource.java      |  45 ++++--
 .../commons/consensus/index/ProgressIndexTest.java |  82 ++++++++++
 .../pipe/source/IoTDBNonDataRegionSourceTest.java  | 172 +++++++++++++++++++++
 24 files changed, 525 insertions(+), 39 deletions(-)

diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/agent/task/PipeConfigNodeTaskAgent.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/agent/task/PipeConfigNodeTaskAgent.java
index 7f244962847..8668debe222 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/agent/task/PipeConfigNodeTaskAgent.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/agent/task/PipeConfigNodeTaskAgent.java
@@ -52,6 +52,7 @@ import java.util.Collection;
 import java.util.Collections;
 import java.util.List;
 import java.util.Map;
+import java.util.Objects;
 import java.util.Optional;
 import java.util.concurrent.atomic.AtomicLong;
 import java.util.function.Consumer;
@@ -182,10 +183,13 @@ public class PipeConfigNodeTaskAgent extends 
PipeTaskAgent {
         }
 
         final ProgressIndex progressIndex = 
groupId2TaskMetaMap.get(regionId).getProgressIndex();
-        if (progressIndex instanceof MetaProgressIndex) {
-          if (((MetaProgressIndex) progressIndex).getIndex() + 1
-              < listeningQueueNewFirstIndex.get()) {
-            listeningQueueNewFirstIndex.set(((MetaProgressIndex) 
progressIndex).getIndex() + 1);
+        final Optional<MetaProgressIndex> metaProgressIndex =
+            Objects.isNull(progressIndex)
+                ? Optional.empty()
+                : 
progressIndex.getProgressIndexByType(MetaProgressIndex.class);
+        if (metaProgressIndex.isPresent()) {
+          if (metaProgressIndex.get().getIndex() + 1 < 
listeningQueueNewFirstIndex.get()) {
+            listeningQueueNewFirstIndex.set(metaProgressIndex.get().getIndex() 
+ 1);
           }
         } else {
           // Do not clear "minimumProgressIndex"s related queues to avoid 
clearing
diff --git 
a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java
 
b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java
index b75097a50cd..abda0e8031f 100644
--- 
a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java
+++ 
b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodePipeMessages.java
@@ -434,6 +434,8 @@ public final class DataNodePipeMessages {
       "Failed to send request {} (watermark = {}) to {}";
   public static final String 
FAILED_TO_TRIGGER_COMBINE_WATERMARK_COUNT_PROGRESSINDEX =
       "Failed to trigger combine. watermark={}, count={}, progressIndex={}";
+  public static final String 
EXCEPTION_FAILED_TO_INITIALIZE_STATEPROGRESSINDEX_FROM_PROGRESS_INDEX_ARG_E95617F9
 =
+      "Failed to initialize StateProgressIndex from progress index %s.";
   public static final String FAILURE_OCCURRED_WHEN_TRYING_TO_COMMIT_PROGRESS =
       "Failure occurred when trying to commit progress index. timestamp={}, 
count={}, "
           + "progressIndex={}";
diff --git 
a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java
 
b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java
index f03d649ac81..11604aef115 100644
--- 
a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java
+++ 
b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodePipeMessages.java
@@ -412,6 +412,8 @@ public final class DataNodePipeMessages {
       "发送 request {}(watermark = {})到 {} 失败";
   public static final String 
FAILED_TO_TRIGGER_COMBINE_WATERMARK_COUNT_PROGRESSINDEX =
       "触发合并失败。watermark={}, count={}, progressIndex={}";
+  public static final String 
EXCEPTION_FAILED_TO_INITIALIZE_STATEPROGRESSINDEX_FROM_PROGRESS_INDEX_ARG_E95617F9
 =
+      "无法从进度索引 %s 初始化 StateProgressIndex。";
   public static final String FAILURE_OCCURRED_WHEN_TRYING_TO_COMMIT_PROGRESS =
       "尝试提交进度索引时发生失败。timestamp={}, count={}, "
           + "progressIndex={}";
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 dce7ee5f6d9..3adad062db7 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
@@ -240,11 +240,15 @@ public class PipeDataNodeTaskAgent extends PipeTaskAgent {
         }
 
         final ProgressIndex progressIndex = pipeTaskMeta.getProgressIndex();
-        if (progressIndex instanceof MetaProgressIndex) {
-          if (((MetaProgressIndex) progressIndex).getIndex() + 1
+        final Optional<MetaProgressIndex> metaProgressIndex =
+            Objects.isNull(progressIndex)
+                ? Optional.empty()
+                : 
progressIndex.getProgressIndexByType(MetaProgressIndex.class);
+        if (metaProgressIndex.isPresent()) {
+          if (metaProgressIndex.get().getIndex() + 1
               < schemaRegionId2ListeningQueueNewFirstIndex.getOrDefault(id, 
Long.MAX_VALUE)) {
             schemaRegionId2ListeningQueueNewFirstIndex.put(
-                id, ((MetaProgressIndex) progressIndex).getIndex() + 1);
+                id, metaProgressIndex.get().getIndex() + 1);
           }
         } else {
           // Do not clear "minimumProgressIndex"s related queues to avoid 
clearing
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/processor/aggregate/AggregateProcessor.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/processor/aggregate/AggregateProcessor.java
index 2224fc12fcc..d9140659e26 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/processor/aggregate/AggregateProcessor.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/processor/aggregate/AggregateProcessor.java
@@ -351,19 +351,15 @@ public class AggregateProcessor implements PipeProcessor {
 
     // Restore window state
     final ProgressIndex index = pipeTaskMeta.getProgressIndex();
-    if (index == MinimumProgressIndex.INSTANCE) {
+    if (Objects.isNull(index) || index == MinimumProgressIndex.INSTANCE) {
       return;
     }
-    if (!(index instanceof TimeWindowStateProgressIndex)) {
-      throw new PipeException(
-          String.format(
-              DataNodePipeMessages
-                  
.PIPE_EXCEPTION_THE_AGGREGATE_PROCESSOR_DOES_NOT_SUPPORT_PROGRESSINDEXTYPE_35351D27,
-              index.getType()));
-    }
-
     final TimeWindowStateProgressIndex timeWindowStateProgressIndex =
-        (TimeWindowStateProgressIndex) index;
+        
index.getProgressIndexByType(TimeWindowStateProgressIndex.class).orElse(null);
+    // A pipe altered from another processor may not have window state yet.
+    if (Objects.isNull(timeWindowStateProgressIndex)) {
+      return;
+    }
     for (final Map.Entry<String, Pair<Long, ByteBuffer>> entry :
         
timeWindowStateProgressIndex.getTimeSeries2TimestampWindowBufferPairMap().entrySet())
 {
       final AtomicReference<TimeSeriesRuntimeState> stateReference =
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/processor/twostage/plugin/TwoStageCountProcessor.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/processor/twostage/plugin/TwoStageCountProcessor.java
index 447cf0bbaf5..f937989dfb9 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/processor/twostage/plugin/TwoStageCountProcessor.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/processor/twostage/plugin/TwoStageCountProcessor.java
@@ -156,13 +156,7 @@ public class TwoStageCountProcessor implements 
PipeProcessor {
     outputSeries = parseOutputSeries(parameters);
 
     if (Objects.nonNull(pipeTaskMeta) && 
Objects.nonNull(pipeTaskMeta.getProgressIndex())) {
-      if (pipeTaskMeta.getProgressIndex() instanceof MinimumProgressIndex) {
-        pipeTaskMeta.updateProgressIndex(
-            new StateProgressIndex(Long.MIN_VALUE, new HashMap<>(), 
MinimumProgressIndex.INSTANCE));
-      }
-
-      final StateProgressIndex stateProgressIndex =
-          (StateProgressIndex) pipeTaskMeta.getProgressIndex();
+      final StateProgressIndex stateProgressIndex = 
initializeStateProgressIndex(pipeTaskMeta);
       localCommitProgressIndex.set(stateProgressIndex.getInnerProgressIndex());
       final Binary localCountState = 
stateProgressIndex.getState().get(LOCAL_COUNT_STATE_KEY);
       localCount.set(
@@ -191,6 +185,27 @@ public class TwoStageCountProcessor implements 
PipeProcessor {
         parameters.getStringByKeys(PROCESSOR_OUTPUT_SERIES_KEY, 
_PROCESSOR_OUTPUT_SERIES_KEY));
   }
 
+  static StateProgressIndex initializeStateProgressIndex(final PipeTaskMeta 
pipeTaskMeta) {
+    final ProgressIndex progressIndex = pipeTaskMeta.getProgressIndex();
+    if (progressIndex instanceof StateProgressIndex stateProgressIndex) {
+      return stateProgressIndex;
+    }
+
+    final ProgressIndex updatedProgressIndex =
+        pipeTaskMeta.updateProgressIndex(
+            new StateProgressIndex(
+                Long.MIN_VALUE, Collections.emptyMap(), 
MinimumProgressIndex.INSTANCE));
+    return updatedProgressIndex
+        .getProgressIndexByType(StateProgressIndex.class)
+        .orElseThrow(
+            () ->
+                new PipeException(
+                    String.format(
+                        DataNodePipeMessages
+                            
.EXCEPTION_FAILED_TO_INITIALIZE_STATEPROGRESSINDEX_FROM_PROGRESS_INDEX_ARG_E95617F9,
+                        updatedProgressIndex)));
+  }
+
   @Override
   public void process(TabletInsertionEvent tabletInsertionEvent, 
EventCollector eventCollector)
       throws Exception {
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/historical/PipeHistoricalDataRegionTsFileAndDeletionSource.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/historical/PipeHistoricalDataRegionTsFileAndDeletionSource.java
index b5ab9e0c65e..bfd705f7419 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/historical/PipeHistoricalDataRegionTsFileAndDeletionSource.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/historical/PipeHistoricalDataRegionTsFileAndDeletionSource.java
@@ -696,7 +696,7 @@ public class PipeHistoricalDataRegionTsFileAndDeletionSource
 
     resourceList.sort(
         (o1, o2) ->
-            startIndex instanceof TimeWindowStateProgressIndex
+            Objects.nonNull(getTimeWindowStateProgressIndex(startIndex))
                 ? Long.compare(o1.getFileStartTime(), o2.getFileStartTime())
                 : comparePersistentResourcesByProgressIndex(o1, o2));
   }
@@ -898,11 +898,11 @@ public class 
PipeHistoricalDataRegionTsFileAndDeletionSource
   }
 
   private boolean mayTsFileContainUnprocessedData(final TsFileResource 
resource) {
-    final ProgressIndex innerStartIndex = getInnerProgressIndex(startIndex);
-    if (innerStartIndex instanceof TimeWindowStateProgressIndex) {
+    final TimeWindowStateProgressIndex timeWindowStateProgressIndex =
+        getTimeWindowStateProgressIndex(startIndex);
+    if (Objects.nonNull(timeWindowStateProgressIndex)) {
       // The resource is closed thus the TsFileResource#getFileEndTime() is 
safe to use
-      return ((TimeWindowStateProgressIndex) innerStartIndex).getMinTime()
-          <= resource.getFileEndTime();
+      return timeWindowStateProgressIndex.getMinTime() <= 
resource.getFileEndTime();
     }
 
     if (pipeName.startsWith(PipeStaticMeta.CONSENSUS_PIPE_PREFIX)) {
@@ -941,6 +941,13 @@ public class 
PipeHistoricalDataRegionTsFileAndDeletionSource
         : Objects.isNull(progressIndex) ? MinimumProgressIndex.INSTANCE : 
progressIndex;
   }
 
+  private TimeWindowStateProgressIndex getTimeWindowStateProgressIndex(
+      final ProgressIndex progressIndex) {
+    return Objects.isNull(progressIndex)
+        ? null
+        : 
progressIndex.getProgressIndexByType(TimeWindowStateProgressIndex.class).orElse(null);
+  }
+
   private boolean isProgressIndexCoveredByTimePartitionProgressIndex(
       final PersistentResource resource,
       final ProgressIndex progressIndex,
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/processor/twostage/plugin/TwoStageCountProcessorTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/processor/twostage/plugin/TwoStageCountProcessorTest.java
index 66db9ccdde8..2343fd3a61d 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/processor/twostage/plugin/TwoStageCountProcessorTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/processor/twostage/plugin/TwoStageCountProcessorTest.java
@@ -19,7 +19,14 @@
 
 package org.apache.iotdb.db.pipe.processor.twostage.plugin;
 
+import org.apache.iotdb.commons.consensus.index.ProgressIndex;
+import org.apache.iotdb.commons.consensus.index.impl.HybridProgressIndex;
+import org.apache.iotdb.commons.consensus.index.impl.MetaProgressIndex;
+import org.apache.iotdb.commons.consensus.index.impl.SimpleProgressIndex;
+import org.apache.iotdb.commons.consensus.index.impl.StateProgressIndex;
+import 
org.apache.iotdb.commons.consensus.index.impl.TimeWindowStateProgressIndex;
 import org.apache.iotdb.commons.path.PartialPath;
+import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTaskMeta;
 import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameters;
 
 import org.junit.Assert;
@@ -37,6 +44,54 @@ public class TwoStageCountProcessorTest {
         "root.db.d.s2", parseOutputSeries("processor.output-series", 
"root.db.d.s2").getFullPath());
   }
 
+  @Test
+  public void testInitializeStateProgressIndexFromHybridProgressIndex() {
+    final MetaProgressIndex metaProgressIndex = new MetaProgressIndex(10L);
+    final SimpleProgressIndex simpleProgressIndex = new SimpleProgressIndex(1, 
2L);
+    final ProgressIndex hybridProgressIndex =
+        new HybridProgressIndex(metaProgressIndex)
+            .updateToMinimumEqualOrIsAfterProgressIndex(simpleProgressIndex);
+    final PipeTaskMeta pipeTaskMeta = new PipeTaskMeta(hybridProgressIndex, 0);
+
+    final StateProgressIndex stateProgressIndex =
+        TwoStageCountProcessor.initializeStateProgressIndex(pipeTaskMeta);
+
+    Assert.assertSame(stateProgressIndex, pipeTaskMeta.getProgressIndex());
+    Assert.assertEquals(
+        metaProgressIndex,
+        
stateProgressIndex.getProgressIndexByType(MetaProgressIndex.class).orElse(null));
+    Assert.assertEquals(
+        simpleProgressIndex,
+        
stateProgressIndex.getProgressIndexByType(SimpleProgressIndex.class).orElse(null));
+  }
+
+  @Test
+  public void 
testInitializeStateProgressIndexFromTimeWindowStateProgressIndex() {
+    final TimeWindowStateProgressIndex timeWindowStateProgressIndex =
+        new TimeWindowStateProgressIndex(Collections.emptyMap());
+    final PipeTaskMeta pipeTaskMeta = new 
PipeTaskMeta(timeWindowStateProgressIndex, 0);
+
+    final StateProgressIndex stateProgressIndex =
+        TwoStageCountProcessor.initializeStateProgressIndex(pipeTaskMeta);
+    Assert.assertEquals(
+        timeWindowStateProgressIndex,
+        
stateProgressIndex.getProgressIndexByType(TimeWindowStateProgressIndex.class).orElse(null));
+
+    final SimpleProgressIndex simpleProgressIndex = new SimpleProgressIndex(1, 
2L);
+    final ProgressIndex updatedProgressIndex =
+        pipeTaskMeta.updateProgressIndex(
+            new StateProgressIndex(1L, Collections.emptyMap(), 
simpleProgressIndex));
+    Assert.assertTrue(updatedProgressIndex instanceof StateProgressIndex);
+    Assert.assertEquals(
+        timeWindowStateProgressIndex,
+        updatedProgressIndex
+            .getProgressIndexByType(TimeWindowStateProgressIndex.class)
+            .orElse(null));
+    Assert.assertEquals(
+        simpleProgressIndex,
+        
updatedProgressIndex.getProgressIndexByType(SimpleProgressIndex.class).orElse(null));
+  }
+
   private PartialPath parseOutputSeries(final String key, final String value) 
throws Exception {
     return TwoStageCountProcessor.parseOutputSeries(
         new PipeParameters(Collections.singletonMap(key, value)));
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/TsFileResourceProgressIndexTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/TsFileResourceProgressIndexTest.java
index 774f6560999..ff7ca591b16 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/TsFileResourceProgressIndexTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/TsFileResourceProgressIndexTest.java
@@ -54,6 +54,7 @@ import java.util.Collections;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
+import java.util.Optional;
 import java.util.Random;
 import java.util.stream.IntStream;
 
@@ -217,6 +218,14 @@ public class TsFileResourceProgressIndexTest {
       throw new UnsupportedOperationException("method not implemented.");
     }
 
+    @Override
+    public <T extends ProgressIndex> Optional<T> getProgressIndexByType(
+        final Class<T> progressIndexClass) {
+      return progressIndexClass.isInstance(this)
+          ? Optional.of(progressIndexClass.cast(this))
+          : Optional.empty();
+    }
+
     @Override
     public TotalOrderSumTuple getTotalOrderSumTuple() {
       return new TotalOrderSumTuple((long) val);
diff --git 
a/iotdb-core/node-commons/src/main/i18n/en/org/apache/iotdb/commons/i18n/PipeMessages.java
 
b/iotdb-core/node-commons/src/main/i18n/en/org/apache/iotdb/commons/i18n/PipeMessages.java
index 742239bc1ca..010cf08f196 100644
--- 
a/iotdb-core/node-commons/src/main/i18n/en/org/apache/iotdb/commons/i18n/PipeMessages.java
+++ 
b/iotdb-core/node-commons/src/main/i18n/en/org/apache/iotdb/commons/i18n/PipeMessages.java
@@ -957,6 +957,8 @@ public final class PipeMessages {
   public static final String 
LOG_ORIGIN_REQUEST_TYPE_MISMATCH_EXPECTED_ARG_ACTUAL_ARG_D96D10AE = "Origin 
request type mismatch: expected {}, actual {}";
   public static final String 
LOG_ORIGIN_BODY_SIZE_MISMATCH_EXPECTED_ARG_ACTUAL_ARG_5D410B75 = "Origin body 
size mismatch: expected {}, actual {}";
   public static final String 
LOG_INVALID_SLICE_INDEX_EXPECTED_ARG_ACTUAL_ARG_2AC41628 = "Invalid slice 
index: expected {}, actual {}";
+  public static final String 
LOG_PIPE_ARG_ARG_ENCOUNTERED_AN_UNEXPECTED_HYBRIDPROGRESSINDEX_IN_ARG_PROGRESS_INDEX_ARG_7C578B17
 =
+      "Pipe {}@{} encountered an unexpected HybridProgressIndex in {}. 
Progress index: {}.";
   public static final String 
EXCEPTION_DECOMPRESSED_LENGTH_SHOULD_BETWEEN_0_ARG_BUT_GOT_ARG_488B3073 = 
"Decompressed length should be between 0 and %d, but got %d.";
   public static final String EXCEPTION_COMMA_50AD1C01 = ", ";
   public static final String 
MESSAGE_NO_DATAPARTITIONTABLE_GENERATION_TASK_FOUND_4414BE55 = "No 
DataPartitionTable generation task found";
diff --git 
a/iotdb-core/node-commons/src/main/i18n/zh/org/apache/iotdb/commons/i18n/PipeMessages.java
 
b/iotdb-core/node-commons/src/main/i18n/zh/org/apache/iotdb/commons/i18n/PipeMessages.java
index 0a8c0badb35..8ea19a43f89 100644
--- 
a/iotdb-core/node-commons/src/main/i18n/zh/org/apache/iotdb/commons/i18n/PipeMessages.java
+++ 
b/iotdb-core/node-commons/src/main/i18n/zh/org/apache/iotdb/commons/i18n/PipeMessages.java
@@ -923,6 +923,8 @@ public final class PipeMessages {
   public static final String 
LOG_ORIGIN_REQUEST_TYPE_MISMATCH_EXPECTED_ARG_ACTUAL_ARG_D96D10AE = "Origin 
request type 不匹配:期望 {},实际 {}";
   public static final String 
LOG_ORIGIN_BODY_SIZE_MISMATCH_EXPECTED_ARG_ACTUAL_ARG_5D410B75 = "Origin body 
size 不匹配:期望 {},实际 {}";
   public static final String 
LOG_INVALID_SLICE_INDEX_EXPECTED_ARG_ACTUAL_ARG_2AC41628 = "无效的 slice index:期望 
{},实际 {}";
+  public static final String 
LOG_PIPE_ARG_ARG_ENCOUNTERED_AN_UNEXPECTED_HYBRIDPROGRESSINDEX_IN_ARG_PROGRESS_INDEX_ARG_7C578B17
 =
+      "Pipe {}@{} 在 {} 中遇到了非预期的 HybridProgressIndex。进度索引:{}。";
   public static final String 
EXCEPTION_DECOMPRESSED_LENGTH_SHOULD_BETWEEN_0_ARG_BUT_GOT_ARG_488B3073 = 
"解压后长度应介于 0 和 %d 之间,但实际为 %d。";
   public static final String EXCEPTION_COMMA_50AD1C01 = ", ";
   public static final String 
MESSAGE_NO_DATAPARTITIONTABLE_GENERATION_TASK_FOUND_4414BE55 = "未找到 
DataPartitionTable 生成任务";
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/ProgressIndex.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/ProgressIndex.java
index 979eee0c8db..83fdac10b5c 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/ProgressIndex.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/ProgressIndex.java
@@ -34,6 +34,7 @@ import java.io.OutputStream;
 import java.nio.ByteBuffer;
 import java.util.List;
 import java.util.Map;
+import java.util.Optional;
 import java.util.concurrent.locks.ReentrantReadWriteLock;
 import java.util.stream.Collectors;
 import java.util.stream.IntStream;
@@ -158,6 +159,15 @@ public abstract class ProgressIndex implements Accountable 
{
    */
   public abstract ProgressIndexType getType();
 
+  /**
+   * Extracts a progress index of the given type from this progress index.
+   *
+   * <p>{@link StateProgressIndex} and {@link HybridProgressIndex} are 
recursively unwrapped because
+   * they may contain progress indexes from other causal chains.
+   */
+  public abstract <T extends ProgressIndex> Optional<T> getProgressIndexByType(
+      Class<T> progressIndexClass);
+
   /**
    * Get the sum of the tuples of each total order relation of the {@link 
ProgressIndex}, which is
    * used for topological sorting of the {@link ProgressIndex}.
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/HybridProgressIndex.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/HybridProgressIndex.java
index 2c8895532dc..724a1500911 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/HybridProgressIndex.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/HybridProgressIndex.java
@@ -36,6 +36,7 @@ import java.util.Collections;
 import java.util.HashMap;
 import java.util.Map;
 import java.util.Objects;
+import java.util.Optional;
 import java.util.concurrent.locks.ReentrantReadWriteLock;
 import java.util.stream.Collectors;
 
@@ -225,6 +226,30 @@ public class HybridProgressIndex extends ProgressIndex {
     return ProgressIndexType.HYBRID_PROGRESS_INDEX;
   }
 
+  @Override
+  public <T extends ProgressIndex> Optional<T> getProgressIndexByType(
+      final Class<T> progressIndexClass) {
+    if (progressIndexClass.isInstance(this)) {
+      return Optional.of(progressIndexClass.cast(this));
+    }
+
+    final Map<Short, ProgressIndex> type2Index = getType2Index();
+    // Prefer a direct component over one nested in another composite progress 
index.
+    for (final ProgressIndex progressIndex : type2Index.values()) {
+      if (progressIndexClass.isInstance(progressIndex)) {
+        return Optional.of(progressIndexClass.cast(progressIndex));
+      }
+    }
+    for (final ProgressIndex progressIndex : type2Index.values()) {
+      final Optional<T> extractedProgressIndex =
+          progressIndex.getProgressIndexByType(progressIndexClass);
+      if (extractedProgressIndex.isPresent()) {
+        return extractedProgressIndex;
+      }
+    }
+    return Optional.empty();
+  }
+
   @Override
   public TotalOrderSumTuple getTotalOrderSumTuple() {
     lock.readLock().lock();
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/IoTProgressIndex.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/IoTProgressIndex.java
index 8f6a24845aa..2de6e4cb135 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/IoTProgressIndex.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/IoTProgressIndex.java
@@ -35,6 +35,7 @@ import java.util.Collections;
 import java.util.HashMap;
 import java.util.Map;
 import java.util.Objects;
+import java.util.Optional;
 import java.util.concurrent.locks.ReentrantReadWriteLock;
 
 public class IoTProgressIndex extends ProgressIndex {
@@ -197,6 +198,14 @@ public class IoTProgressIndex extends ProgressIndex {
     return ProgressIndexType.IOT_PROGRESS_INDEX;
   }
 
+  @Override
+  public <T extends ProgressIndex> Optional<T> getProgressIndexByType(
+      final Class<T> progressIndexClass) {
+    return progressIndexClass.isInstance(this)
+        ? Optional.of(progressIndexClass.cast(this))
+        : Optional.empty();
+  }
+
   @Override
   public TotalOrderSumTuple getTotalOrderSumTuple() {
     lock.readLock().lock();
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/MetaProgressIndex.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/MetaProgressIndex.java
index 75322152d45..80a31a7e4f3 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/MetaProgressIndex.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/MetaProgressIndex.java
@@ -32,6 +32,7 @@ import java.io.InputStream;
 import java.io.OutputStream;
 import java.nio.ByteBuffer;
 import java.util.Objects;
+import java.util.Optional;
 import java.util.concurrent.locks.ReentrantReadWriteLock;
 
 public class MetaProgressIndex extends ProgressIndex {
@@ -153,6 +154,14 @@ public class MetaProgressIndex extends ProgressIndex {
     return ProgressIndexType.META_PROGRESS_INDEX;
   }
 
+  @Override
+  public <T extends ProgressIndex> Optional<T> getProgressIndexByType(
+      final Class<T> progressIndexClass) {
+    return progressIndexClass.isInstance(this)
+        ? Optional.of(progressIndexClass.cast(this))
+        : Optional.empty();
+  }
+
   @Override
   public TotalOrderSumTuple getTotalOrderSumTuple() {
     lock.readLock().lock();
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/MinimumProgressIndex.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/MinimumProgressIndex.java
index e22f82c9fbb..b34e97e3493 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/MinimumProgressIndex.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/MinimumProgressIndex.java
@@ -28,6 +28,7 @@ import java.io.IOException;
 import java.io.InputStream;
 import java.io.OutputStream;
 import java.nio.ByteBuffer;
+import java.util.Optional;
 
 public class MinimumProgressIndex extends ProgressIndex {
 
@@ -82,6 +83,14 @@ public class MinimumProgressIndex extends ProgressIndex {
     return ProgressIndexType.MINIMUM_PROGRESS_INDEX;
   }
 
+  @Override
+  public <T extends ProgressIndex> Optional<T> getProgressIndexByType(
+      final Class<T> progressIndexClass) {
+    return progressIndexClass.isInstance(this)
+        ? Optional.of(progressIndexClass.cast(this))
+        : Optional.empty();
+  }
+
   @Override
   public TotalOrderSumTuple getTotalOrderSumTuple() {
     return TOTAL_ORDER_SUM_TUPLE;
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/RecoverProgressIndex.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/RecoverProgressIndex.java
index 5756594abeb..f42ad50f63a 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/RecoverProgressIndex.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/RecoverProgressIndex.java
@@ -36,6 +36,7 @@ import java.util.Collections;
 import java.util.HashMap;
 import java.util.Map;
 import java.util.Objects;
+import java.util.Optional;
 import java.util.concurrent.locks.ReentrantReadWriteLock;
 import java.util.stream.Collectors;
 
@@ -202,6 +203,14 @@ public class RecoverProgressIndex extends ProgressIndex {
     return ProgressIndexType.RECOVER_PROGRESS_INDEX;
   }
 
+  @Override
+  public <T extends ProgressIndex> Optional<T> getProgressIndexByType(
+      final Class<T> progressIndexClass) {
+    return progressIndexClass.isInstance(this)
+        ? Optional.of(progressIndexClass.cast(this))
+        : Optional.empty();
+  }
+
   @Override
   public TotalOrderSumTuple getTotalOrderSumTuple() {
     lock.readLock().lock();
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/SimpleProgressIndex.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/SimpleProgressIndex.java
index 162dc9f128d..be9a0c5ddee 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/SimpleProgressIndex.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/SimpleProgressIndex.java
@@ -32,6 +32,7 @@ import java.io.InputStream;
 import java.io.OutputStream;
 import java.nio.ByteBuffer;
 import java.util.Objects;
+import java.util.Optional;
 import java.util.concurrent.locks.ReentrantReadWriteLock;
 
 public class SimpleProgressIndex extends ProgressIndex {
@@ -179,6 +180,14 @@ public class SimpleProgressIndex extends ProgressIndex {
     return ProgressIndexType.SIMPLE_PROGRESS_INDEX;
   }
 
+  @Override
+  public <T extends ProgressIndex> Optional<T> getProgressIndexByType(
+      final Class<T> progressIndexClass) {
+    return progressIndexClass.isInstance(this)
+        ? Optional.of(progressIndexClass.cast(this))
+        : Optional.empty();
+  }
+
   @Override
   public TotalOrderSumTuple getTotalOrderSumTuple() {
     return new TotalOrderSumTuple(memTableFlushOrderId, (long) rebootTimes);
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/StateProgressIndex.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/StateProgressIndex.java
index b6c51665e19..ed2b4d91145 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/StateProgressIndex.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/StateProgressIndex.java
@@ -36,6 +36,7 @@ import java.nio.ByteBuffer;
 import java.util.HashMap;
 import java.util.Map;
 import java.util.Objects;
+import java.util.Optional;
 import java.util.concurrent.locks.ReentrantReadWriteLock;
 
 /**
@@ -196,6 +197,14 @@ public class StateProgressIndex extends ProgressIndex {
     return ProgressIndexType.STATE_PROGRESS_INDEX;
   }
 
+  @Override
+  public <T extends ProgressIndex> Optional<T> getProgressIndexByType(
+      final Class<T> progressIndexClass) {
+    return progressIndexClass.isInstance(this)
+        ? Optional.of(progressIndexClass.cast(this))
+        : getInnerProgressIndex().getProgressIndexByType(progressIndexClass);
+  }
+
   @Override
   public TotalOrderSumTuple getTotalOrderSumTuple() {
     return innerProgressIndex.getTotalOrderSumTuple();
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/TimePartitionProgressIndex.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/TimePartitionProgressIndex.java
index fa1bb84fff6..05ce7c86718 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/TimePartitionProgressIndex.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/TimePartitionProgressIndex.java
@@ -37,6 +37,7 @@ import java.util.Collections;
 import java.util.HashMap;
 import java.util.Map;
 import java.util.Objects;
+import java.util.Optional;
 import java.util.concurrent.locks.ReentrantReadWriteLock;
 import java.util.stream.Collectors;
 
@@ -233,6 +234,14 @@ public class TimePartitionProgressIndex extends 
ProgressIndex {
     return ProgressIndexType.TIME_PARTITION_PROGRESS_INDEX;
   }
 
+  @Override
+  public <T extends ProgressIndex> Optional<T> getProgressIndexByType(
+      final Class<T> progressIndexClass) {
+    return progressIndexClass.isInstance(this)
+        ? Optional.of(progressIndexClass.cast(this))
+        : Optional.empty();
+  }
+
   @Override
   public TotalOrderSumTuple getTotalOrderSumTuple() {
     lock.readLock().lock();
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/TimeWindowStateProgressIndex.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/TimeWindowStateProgressIndex.java
index 7e959cf14b8..56293fd5fc0 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/TimeWindowStateProgressIndex.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/TimeWindowStateProgressIndex.java
@@ -39,6 +39,7 @@ import java.util.Collections;
 import java.util.HashMap;
 import java.util.Map;
 import java.util.Objects;
+import java.util.Optional;
 import java.util.concurrent.locks.ReentrantReadWriteLock;
 import java.util.stream.Collectors;
 
@@ -210,7 +211,7 @@ public class TimeWindowStateProgressIndex extends 
ProgressIndex {
     lock.writeLock().lock();
     try {
       if (!(progressIndex instanceof TimeWindowStateProgressIndex)) {
-        return this;
+        return ProgressIndex.blendProgressIndex(this, progressIndex);
       }
 
       final TimeWindowStateProgressIndex thisTimeWindowStateProgressIndex = 
this;
@@ -240,6 +241,14 @@ public class TimeWindowStateProgressIndex extends 
ProgressIndex {
     return ProgressIndexType.TIME_WINDOW_STATE_PROGRESS_INDEX;
   }
 
+  @Override
+  public <T extends ProgressIndex> Optional<T> getProgressIndexByType(
+      final Class<T> progressIndexClass) {
+    return progressIndexClass.isInstance(this)
+        ? Optional.of(progressIndexClass.cast(this))
+        : Optional.empty();
+  }
+
   @Override
   public TotalOrderSumTuple getTotalOrderSumTuple() {
     throw new UnsupportedOperationException(
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/source/IoTDBNonDataRegionSource.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/source/IoTDBNonDataRegionSource.java
index 286f373c8a7..d4d55ba7dac 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/source/IoTDBNonDataRegionSource.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/source/IoTDBNonDataRegionSource.java
@@ -20,8 +20,8 @@
 package org.apache.iotdb.commons.pipe.source;
 
 import org.apache.iotdb.commons.consensus.index.ProgressIndex;
+import org.apache.iotdb.commons.consensus.index.impl.HybridProgressIndex;
 import org.apache.iotdb.commons.consensus.index.impl.MetaProgressIndex;
-import org.apache.iotdb.commons.consensus.index.impl.MinimumProgressIndex;
 import org.apache.iotdb.commons.exception.IllegalPathException;
 import org.apache.iotdb.commons.exception.auth.AccessDeniedException;
 import org.apache.iotdb.commons.i18n.PipeMessages;
@@ -42,6 +42,8 @@ import org.apache.iotdb.pipe.api.event.Event;
 import org.apache.iotdb.pipe.api.exception.PipeException;
 
 import org.apache.tsfile.utils.Pair;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
 
 import java.io.IOException;
 import java.util.LinkedList;
@@ -54,6 +56,8 @@ import java.util.concurrent.atomic.AtomicBoolean;
 @TableModel
 public abstract class IoTDBNonDataRegionSource extends IoTDBSource {
 
+  private static final Logger LOGGER = 
LoggerFactory.getLogger(IoTDBNonDataRegionSource.class);
+
   protected IoTDBTreePatternOperations treePattern;
   protected TablePattern tablePattern;
 
@@ -68,6 +72,7 @@ public abstract class IoTDBNonDataRegionSource extends 
IoTDBSource {
   // If the extractor is closed, it should not be started again. This is to 
avoid the case that
   // the extractor is closed and then be reused by processor.
   protected final AtomicBoolean hasBeenClosed = new AtomicBoolean(false);
+  private final AtomicBoolean hasWarnedUnexpectedHybridProgressIndex = new 
AtomicBoolean(false);
 
   protected PipeWritePlanEvent lastEvent = null;
 
@@ -100,14 +105,15 @@ public abstract class IoTDBNonDataRegionSource extends 
IoTDBSource {
     }
 
     final ProgressIndex progressIndex = pipeTaskMeta.getProgressIndex();
+    warnIfUnexpectedHybridProgressIndex(progressIndex);
+    final MetaProgressIndex metaProgressIndex = 
extractMetaProgressIndex(progressIndex);
     final long nextIndex =
-        progressIndex instanceof MinimumProgressIndex
+        Objects.isNull(metaProgressIndex)
                 // If the index is invalid, the queue is seen as cleared 
before and thus
                 // needs snapshot re-transferring
-                || !getListeningQueue()
-                    .isGivenNextIndexValid(((MetaProgressIndex) 
progressIndex).getIndex() + 1)
+                || 
!getListeningQueue().isGivenNextIndexValid(metaProgressIndex.getIndex() + 1)
             ? getNextIndexAfterSnapshot()
-            : ((MetaProgressIndex) progressIndex).getIndex() + 1;
+            : metaProgressIndex.getIndex() + 1;
     iterator = getListeningQueue().newIterator(nextIndex);
     super.start();
   }
@@ -296,10 +302,31 @@ public abstract class IoTDBNonDataRegionSource extends 
IoTDBSource {
     if (Objects.isNull(pipeTaskMeta)) {
       return 0L;
     }
-    return !(pipeTaskMeta.getProgressIndex() instanceof MinimumProgressIndex)
-        ? getListeningQueue().getTailIndex()
-            - ((MetaProgressIndex) pipeTaskMeta.getProgressIndex()).getIndex()
-            - 1
+    final ProgressIndex progressIndex = pipeTaskMeta.getProgressIndex();
+    warnIfUnexpectedHybridProgressIndex(progressIndex);
+    final MetaProgressIndex metaProgressIndex = 
extractMetaProgressIndex(progressIndex);
+    return Objects.nonNull(metaProgressIndex)
+        ? getListeningQueue().getTailIndex() - metaProgressIndex.getIndex() - 1
         : getListeningQueue().getSize() + historicalEventsCount;
   }
+
+  private static MetaProgressIndex extractMetaProgressIndex(final 
ProgressIndex progressIndex) {
+    return Objects.isNull(progressIndex)
+        ? null
+        : 
progressIndex.getProgressIndexByType(MetaProgressIndex.class).orElse(null);
+  }
+
+  private void warnIfUnexpectedHybridProgressIndex(final ProgressIndex 
progressIndex) {
+    if (Objects.nonNull(progressIndex)
+        && 
progressIndex.getProgressIndexByType(HybridProgressIndex.class).isPresent()
+        && hasWarnedUnexpectedHybridProgressIndex.compareAndSet(false, true)) {
+      LOGGER.warn(
+          PipeMessages
+              
.LOG_PIPE_ARG_ARG_ENCOUNTERED_AN_UNEXPECTED_HYBRIDPROGRESSINDEX_IN_ARG_PROGRESS_INDEX_ARG_7C578B17,
+          pipeName,
+          creationTime,
+          getClass().getSimpleName(),
+          progressIndex);
+    }
+  }
 }
diff --git 
a/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/consensus/index/ProgressIndexTest.java
 
b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/consensus/index/ProgressIndexTest.java
new file mode 100644
index 00000000000..07ee3ef268e
--- /dev/null
+++ 
b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/consensus/index/ProgressIndexTest.java
@@ -0,0 +1,82 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.commons.consensus.index;
+
+import org.apache.iotdb.commons.consensus.index.impl.HybridProgressIndex;
+import org.apache.iotdb.commons.consensus.index.impl.MetaProgressIndex;
+import org.apache.iotdb.commons.consensus.index.impl.SimpleProgressIndex;
+import org.apache.iotdb.commons.consensus.index.impl.StateProgressIndex;
+import 
org.apache.iotdb.commons.consensus.index.impl.TimeWindowStateProgressIndex;
+
+import org.junit.Assert;
+import org.junit.Test;
+
+import java.util.Collections;
+
+public class ProgressIndexTest {
+
+  @Test
+  public void testGetProgressIndexByTypeFromStateWrappedHybridProgressIndex() {
+    final MetaProgressIndex metaProgressIndex = new MetaProgressIndex(10L);
+    final SimpleProgressIndex simpleProgressIndex = new SimpleProgressIndex(1, 
2L);
+    final ProgressIndex hybridProgressIndex =
+        new HybridProgressIndex(metaProgressIndex)
+            .updateToMinimumEqualOrIsAfterProgressIndex(simpleProgressIndex);
+    final StateProgressIndex stateProgressIndex =
+        new StateProgressIndex(1L, Collections.emptyMap(), 
hybridProgressIndex);
+
+    Assert.assertEquals(
+        metaProgressIndex,
+        
stateProgressIndex.getProgressIndexByType(MetaProgressIndex.class).orElse(null));
+    Assert.assertEquals(
+        simpleProgressIndex,
+        
stateProgressIndex.getProgressIndexByType(SimpleProgressIndex.class).orElse(null));
+    Assert.assertSame(
+        hybridProgressIndex,
+        
stateProgressIndex.getProgressIndexByType(HybridProgressIndex.class).orElse(null));
+    Assert.assertFalse(
+        
stateProgressIndex.getProgressIndexByType(TimeWindowStateProgressIndex.class).isPresent());
+  }
+
+  @Test
+  public void 
testTimeWindowStateProgressIndexBlendsWithOtherProgressIndexTypes() {
+    final TimeWindowStateProgressIndex timeWindowStateProgressIndex =
+        new TimeWindowStateProgressIndex(Collections.emptyMap());
+    final SimpleProgressIndex simpleProgressIndex = new SimpleProgressIndex(1, 
2L);
+
+    final ProgressIndex blendedProgressIndex =
+        
timeWindowStateProgressIndex.updateToMinimumEqualOrIsAfterProgressIndex(
+            simpleProgressIndex);
+    Assert.assertTrue(blendedProgressIndex instanceof HybridProgressIndex);
+    Assert.assertEquals(
+        timeWindowStateProgressIndex,
+        blendedProgressIndex
+            .getProgressIndexByType(TimeWindowStateProgressIndex.class)
+            .orElse(null));
+    Assert.assertEquals(
+        simpleProgressIndex,
+        
blendedProgressIndex.getProgressIndexByType(SimpleProgressIndex.class).orElse(null));
+
+    final ProgressIndex reverseBlendedProgressIndex =
+        simpleProgressIndex.updateToMinimumEqualOrIsAfterProgressIndex(
+            timeWindowStateProgressIndex);
+    Assert.assertEquals(blendedProgressIndex, reverseBlendedProgressIndex);
+  }
+}
diff --git 
a/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/pipe/source/IoTDBNonDataRegionSourceTest.java
 
b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/pipe/source/IoTDBNonDataRegionSourceTest.java
new file mode 100644
index 00000000000..3bfbd6dcefd
--- /dev/null
+++ 
b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/pipe/source/IoTDBNonDataRegionSourceTest.java
@@ -0,0 +1,172 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.commons.pipe.source;
+
+import org.apache.iotdb.commons.consensus.index.ProgressIndex;
+import org.apache.iotdb.commons.consensus.index.impl.HybridProgressIndex;
+import org.apache.iotdb.commons.consensus.index.impl.MetaProgressIndex;
+import org.apache.iotdb.commons.consensus.index.impl.SimpleProgressIndex;
+import org.apache.iotdb.commons.consensus.index.impl.StateProgressIndex;
+import org.apache.iotdb.commons.exception.IllegalPathException;
+import org.apache.iotdb.commons.exception.auth.AccessDeniedException;
+import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTaskMeta;
+import 
org.apache.iotdb.commons.pipe.datastructure.queue.listening.AbstractPipeListeningQueue;
+import org.apache.iotdb.commons.pipe.event.PipeSnapshotEvent;
+import org.apache.iotdb.commons.pipe.event.PipeWritePlanEvent;
+
+import org.junit.Assert;
+import org.junit.Test;
+import org.mockito.Mockito;
+
+import javax.annotation.Nonnull;
+
+import java.io.IOException;
+import java.util.Collections;
+import java.util.Optional;
+
+public class IoTDBNonDataRegionSourceTest {
+
+  @Test
+  public void testStartWithStateWrappedHybridProgressIndex() throws Exception {
+    final AbstractPipeListeningQueue listeningQueue =
+        Mockito.mock(AbstractPipeListeningQueue.class);
+    Mockito.when(listeningQueue.isGivenNextIndexValid(11L)).thenReturn(true);
+
+    final ProgressIndex hybridProgressIndex =
+        new HybridProgressIndex(new MetaProgressIndex(10L))
+            .updateToMinimumEqualOrIsAfterProgressIndex(new 
SimpleProgressIndex(1, 2L));
+    final StateProgressIndex stateProgressIndex =
+        new StateProgressIndex(1L, Collections.emptyMap(), 
hybridProgressIndex);
+    final TestNonDataRegionSource source =
+        new TestNonDataRegionSource(listeningQueue, new 
PipeTaskMeta(stateProgressIndex, 0));
+
+    source.start();
+
+    Mockito.verify(listeningQueue).newIterator(11L);
+  }
+
+  @Test
+  public void testGetUnTransferredEventCountWithHybridProgressIndex() {
+    final AbstractPipeListeningQueue listeningQueue =
+        Mockito.mock(AbstractPipeListeningQueue.class);
+    Mockito.when(listeningQueue.getTailIndex()).thenReturn(20L);
+
+    final ProgressIndex hybridProgressIndex =
+        new HybridProgressIndex(new MetaProgressIndex(10L))
+            .updateToMinimumEqualOrIsAfterProgressIndex(new 
SimpleProgressIndex(1, 2L));
+    final TestNonDataRegionSource source =
+        new TestNonDataRegionSource(listeningQueue, new 
PipeTaskMeta(hybridProgressIndex, 0));
+
+    Assert.assertEquals(9L, source.getUnTransferredEventCount());
+  }
+
+  @Test
+  public void 
testGetUnTransferredEventCountWithHybridProgressIndexWithoutMetaIndex() {
+    final AbstractPipeListeningQueue listeningQueue =
+        Mockito.mock(AbstractPipeListeningQueue.class);
+    Mockito.when(listeningQueue.getSize()).thenReturn(7L);
+
+    final TestNonDataRegionSource source =
+        new TestNonDataRegionSource(
+            listeningQueue,
+            new PipeTaskMeta(new HybridProgressIndex(new 
SimpleProgressIndex(1, 2L)), 0));
+
+    Assert.assertEquals(7L, source.getUnTransferredEventCount());
+  }
+
+  private static final class TestNonDataRegionSource extends 
IoTDBNonDataRegionSource {
+
+    private final AbstractPipeListeningQueue listeningQueue;
+
+    private TestNonDataRegionSource(
+        final AbstractPipeListeningQueue listeningQueue, final PipeTaskMeta 
pipeTaskMeta) {
+      this.listeningQueue = listeningQueue;
+      this.pipeTaskMeta = pipeTaskMeta;
+    }
+
+    @Override
+    protected AbstractPipeListeningQueue getListeningQueue() {
+      return listeningQueue;
+    }
+
+    @Override
+    protected boolean needTransferSnapshot() {
+      return false;
+    }
+
+    @Override
+    protected void triggerSnapshot() {
+      // Do nothing
+    }
+
+    @Override
+    protected long getMaxBlockingTimeMs() {
+      return 0L;
+    }
+
+    @Override
+    protected boolean canSkipSnapshotPrivilegeCheck(final PipeSnapshotEvent 
event) {
+      return false;
+    }
+
+    @Override
+    protected void initSnapshotGenerator(final PipeSnapshotEvent event)
+        throws IOException, IllegalPathException {
+      // Do nothing
+    }
+
+    @Override
+    protected boolean hasNextEventInCurrentSnapshot() {
+      return false;
+    }
+
+    @Override
+    protected PipeWritePlanEvent getNextEventInCurrentSnapshot() {
+      return null;
+    }
+
+    @Override
+    protected Optional<PipeWritePlanEvent> trimRealtimeEventByPrivilege(
+        final PipeWritePlanEvent event) throws AccessDeniedException {
+      return Optional.of(event);
+    }
+
+    @Override
+    protected Optional<PipeWritePlanEvent> trimRealtimeEventByPipePattern(
+        final PipeWritePlanEvent event) {
+      return Optional.of(event);
+    }
+
+    @Override
+    protected boolean isTypeListened(final PipeWritePlanEvent event) {
+      return true;
+    }
+
+    @Override
+    protected void confineHistoricalEventTransferTypes(final PipeSnapshotEvent 
event) {
+      // Do nothing
+    }
+
+    @Override
+    protected void login(final @Nonnull String password) {
+      // Do nothing
+    }
+  }
+}

Reply via email to