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 ae163537785 [Subscription] Protect retained WAL by committed progress 
(#18399)
ae163537785 is described below

commit ae163537785586b43c0549f35410f1d79219bcca
Author: Caideyipi <[email protected]>
AuthorDate: Tue Aug 25 09:33:22 2026 +0800

    [Subscription] Protect retained WAL by committed progress (#18399)
    
    * [Subscription] Protect retained WAL by committed progress
    
    * fix(subscription): cache WAL retention boundary
    
    * Keep subscription queue registration API compatible
---
 .../consensus/iot/IoTConsensusServerImpl.java      |  14 +-
 .../subscription/SubscriptionQueueRegistry.java    |  55 ++++-
 .../SubscriptionWalRetentionCalculator.java        |  12 +-
 .../SubscriptionWalRetentionCalculatorTest.java    | 101 ++++++++++
 .../consensus/ConsensusPrefetchingQueue.java       | 222 +++++++++++++++++++--
 .../consensus/ConsensusPrefetchingQueueTest.java   |  46 +++++
 6 files changed, 429 insertions(+), 21 deletions(-)

diff --git 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java
 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java
index 338bf30c1d9..e074e7204ee 100644
--- 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java
+++ 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java
@@ -110,6 +110,7 @@ import java.util.concurrent.atomic.AtomicLong;
 import java.util.concurrent.locks.Condition;
 import java.util.concurrent.locks.Lock;
 import java.util.concurrent.locks.ReentrantLock;
+import java.util.function.LongSupplier;
 import java.util.regex.Pattern;
 
 import static 
org.apache.iotdb.commons.utils.FileUtils.humanReadableByteCountSI;
@@ -1144,8 +1145,10 @@ public class IoTConsensusServerImpl {
    */
   public void registerSubscriptionQueue(
       final BlockingQueue<IndexedConsensusRequest> queue,
-      final SubscriptionWalRetentionPolicy retentionPolicy) {
-    subscriptionQueueRegistry.register(queue, retentionPolicy);
+      final SubscriptionWalRetentionPolicy retentionPolicy,
+      final LongSupplier committedRetainedMinVersionIdSupplier) {
+    subscriptionQueueRegistry.register(
+        queue, retentionPolicy, committedRetainedMinVersionIdSupplier);
     // Immediately re-evaluate the safe delete index with new subscription 
awareness
     checkAndUpdateSafeDeletedSearchIndex();
     logger.info(
@@ -1288,8 +1291,8 @@ public class IoTConsensusServerImpl {
   }
 
   /**
-   * Computes and updates the safe-to-delete WAL search index based on 
replication progress and
-   * subscription WAL retention policy.
+   * Computes and updates the safe-to-delete WAL search index based on 
replication progress,
+   * subscription WAL retention policy, and consumer-group committed progress.
    *
    * <p>Because multiple subscription topics share one region WAL, the 
effective per-region
    * retention policy is the most conservative policy across all active 
subscription queues on this
@@ -1316,7 +1319,8 @@ public class IoTConsensusServerImpl {
 
     final SubscriptionRetentionBound subscriptionRetentionBound =
         subscriptionWalRetentionCalculator.calculate(
-            subscriptionQueueRegistry.getRetentionPolicies());
+            subscriptionQueueRegistry.getRetentionPolicies(),
+            subscriptionQueueRegistry.getCommittedRetainedMinVersionIds());
 
     consensusReqReader.setSafelyDeletedSearchIndex(
         Math.min(replicationIndex, 
subscriptionRetentionBound.getSafelyDeletedSearchIndex()));
diff --git 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionQueueRegistry.java
 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionQueueRegistry.java
index a0a9422e127..8d7c27f8cf2 100644
--- 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionQueueRegistry.java
+++ 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionQueueRegistry.java
@@ -33,6 +33,7 @@ import java.util.concurrent.BlockingQueue;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicLong;
+import java.util.function.LongSupplier;
 
 public class SubscriptionQueueRegistry {
 
@@ -41,7 +42,7 @@ public class SubscriptionQueueRegistry {
   private static final long QUEUE_FULL_LOG_INTERVAL_MS = 
TimeUnit.SECONDS.toMillis(10);
 
   private final String consensusGroupId;
-  private final Map<BlockingQueue<IndexedConsensusRequest>, 
SubscriptionWalRetentionPolicy> queues =
+  private final Map<BlockingQueue<IndexedConsensusRequest>, 
SubscriptionQueueRegistration> queues =
       new ConcurrentHashMap<>();
   private final AtomicLong droppedEntries = new AtomicLong();
   private final AtomicLong lastDropLogTimeMs = new AtomicLong();
@@ -50,10 +51,26 @@ public class SubscriptionQueueRegistry {
     this.consensusGroupId = consensusGroupId;
   }
 
+  /**
+   * Registers a queue without a committed-progress constraint.
+   *
+   * <p>This overload keeps callers using the original queue-registration API 
source-compatible.
+   * {@link Long#MAX_VALUE} is the neutral value for the retention calculator 
and therefore does not
+   * add an extra WAL-retention constraint.
+   */
   public synchronized void register(
       final BlockingQueue<IndexedConsensusRequest> queue,
       final SubscriptionWalRetentionPolicy retentionPolicy) {
-    queues.put(queue, retentionPolicy);
+    register(queue, retentionPolicy, () -> Long.MAX_VALUE);
+  }
+
+  public synchronized void register(
+      final BlockingQueue<IndexedConsensusRequest> queue,
+      final SubscriptionWalRetentionPolicy retentionPolicy,
+      final LongSupplier committedRetainedMinVersionIdSupplier) {
+    queues.put(
+        queue,
+        new SubscriptionQueueRegistration(retentionPolicy, 
committedRetainedMinVersionIdSupplier));
   }
 
   // Shares the monitor with offer() so unregister() is a real stop-receiving 
barrier.
@@ -70,7 +87,26 @@ public class SubscriptionQueueRegistry {
   }
 
   public synchronized Collection<SubscriptionWalRetentionPolicy> 
getRetentionPolicies() {
-    return new ArrayList<>(queues.values());
+    final Collection<SubscriptionWalRetentionPolicy> retentionPolicies = new 
ArrayList<>();
+    for (final SubscriptionQueueRegistration registration : queues.values()) {
+      retentionPolicies.add(registration.retentionPolicy);
+    }
+    return retentionPolicies;
+  }
+
+  public Collection<Long> getCommittedRetainedMinVersionIds() {
+    final Collection<LongSupplier> suppliers = new ArrayList<>();
+    synchronized (this) {
+      for (final SubscriptionQueueRegistration registration : queues.values()) 
{
+        suppliers.add(registration.committedRetainedMinVersionIdSupplier);
+      }
+    }
+
+    final Collection<Long> committedRetainedMinVersionIds = new ArrayList<>();
+    for (final LongSupplier supplier : suppliers) {
+      committedRetainedMinVersionIds.add(supplier.getAsLong());
+    }
+    return committedRetainedMinVersionIds;
   }
 
   public synchronized boolean offer(final IndexedConsensusRequest 
indexedConsensusRequest) {
@@ -132,4 +168,17 @@ public class SubscriptionQueueRegistry {
     }
     return offeredToAnyQueue;
   }
+
+  private static final class SubscriptionQueueRegistration {
+
+    private final SubscriptionWalRetentionPolicy retentionPolicy;
+    private final LongSupplier committedRetainedMinVersionIdSupplier;
+
+    private SubscriptionQueueRegistration(
+        final SubscriptionWalRetentionPolicy retentionPolicy,
+        final LongSupplier committedRetainedMinVersionIdSupplier) {
+      this.retentionPolicy = retentionPolicy;
+      this.committedRetainedMinVersionIdSupplier = 
committedRetainedMinVersionIdSupplier;
+    }
+  }
 }
diff --git 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionWalRetentionCalculator.java
 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionWalRetentionCalculator.java
index 3bd2be60f37..ddddf7943de 100644
--- 
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionWalRetentionCalculator.java
+++ 
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionWalRetentionCalculator.java
@@ -84,7 +84,8 @@ public class SubscriptionWalRetentionCalculator {
   }
 
   public SubscriptionRetentionBound calculate(
-      final Collection<SubscriptionWalRetentionPolicy> retentionPolicies) {
+      final Collection<SubscriptionWalRetentionPolicy> retentionPolicies,
+      final Collection<Long> committedRetainedMinVersionIds) {
     SubscriptionRetentionBound mergedBound = 
SubscriptionRetentionBound.noConstraint();
     for (final SubscriptionWalRetentionPolicy policy : retentionPolicies) {
       // For each topic, data can be deleted once either its size retention or 
its time retention
@@ -95,6 +96,15 @@ public class SubscriptionWalRetentionCalculator {
               
.mergeDeleteEither(buildTimeRetentionBound(policy.getRetentionMs()));
       mergedBound = mergedBound.mergeDeleteOnlyIfBoth(perQueueBound);
     }
+    for (final long committedRetainedMinVersionId : 
committedRetainedMinVersionIds) {
+      // Topic retention is a historical replay window, while committed 
progress protects data
+      // that a consumer group has not acknowledged yet. A WAL file can be 
reclaimed only when
+      // both constraints allow it, so keep the more conservative (smaller) 
file-version bound.
+      mergedBound =
+          mergedBound.mergeDeleteOnlyIfBoth(
+              SubscriptionRetentionBound.of(
+                  Long.MAX_VALUE, Math.max(0L, 
committedRetainedMinVersionId)));
+    }
     return mergedBound;
   }
 
diff --git 
a/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionWalRetentionCalculatorTest.java
 
b/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionWalRetentionCalculatorTest.java
new file mode 100644
index 00000000000..6815e8faf23
--- /dev/null
+++ 
b/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionWalRetentionCalculatorTest.java
@@ -0,0 +1,101 @@
+/*
+ * 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.consensus.iot.subscription;
+
+import org.apache.iotdb.consensus.iot.SubscriptionWalRetentionPolicy;
+import org.apache.iotdb.consensus.iot.log.ConsensusReqReader;
+import 
org.apache.iotdb.consensus.iot.subscription.SubscriptionWalRetentionCalculator.SubscriptionRetentionBound;
+
+import org.apache.tsfile.utils.Pair;
+import org.junit.Test;
+
+import java.util.Arrays;
+import java.util.Collections;
+
+import static org.junit.Assert.assertEquals;
+
+public class SubscriptionWalRetentionCalculatorTest {
+
+  @Test
+  public void testCommittedProgressKeepsMoreWalThanTopicRetention() {
+    final SubscriptionWalRetentionCalculator calculator =
+        new SubscriptionWalRetentionCalculator(new 
RetentionTestConsensusReqReader());
+    final SubscriptionRetentionBound bound =
+        calculator.calculate(
+            Collections.singletonList(
+                new SubscriptionWalRetentionPolicy(
+                    "topic", 100L, SubscriptionWalRetentionPolicy.UNBOUNDED)),
+            Arrays.asList(8L, 4L));
+
+    assertEquals(100L, bound.getSafelyDeletedSearchIndex());
+    assertEquals(4L, bound.getRetainedMinVersionId());
+  }
+
+  @Test
+  public void testTopicRetentionKeepsMoreWalThanCommittedProgress() {
+    final SubscriptionWalRetentionCalculator calculator =
+        new SubscriptionWalRetentionCalculator(new 
RetentionTestConsensusReqReader());
+    final SubscriptionRetentionBound bound =
+        calculator.calculate(
+            Collections.singletonList(
+                new SubscriptionWalRetentionPolicy(
+                    "topic", 100L, SubscriptionWalRetentionPolicy.UNBOUNDED)),
+            Collections.singletonList(20L));
+
+    assertEquals(100L, bound.getSafelyDeletedSearchIndex());
+    assertEquals(10L, bound.getRetainedMinVersionId());
+  }
+
+  private static final class RetentionTestConsensusReqReader implements 
ConsensusReqReader {
+
+    @Override
+    public void setSafelyDeletedSearchIndex(final long 
safelyDeletedSearchIndex) {}
+
+    @Override
+    public ReqIterator getReqIterator(final long startIndex) {
+      throw new UnsupportedOperationException();
+    }
+
+    @Override
+    public long getCurrentSearchIndex() {
+      return 0L;
+    }
+
+    @Override
+    public long getCurrentWALFileVersion() {
+      return 0L;
+    }
+
+    @Override
+    public long getTotalSize() {
+      return 200L;
+    }
+
+    @Override
+    public long getRegionDiskUsage() {
+      return 200L;
+    }
+
+    @Override
+    public Pair<Long, Long> getDeletionBoundToFreeAtLeast(final long 
bytesToFree) {
+      return new Pair<>(100L, 10L);
+    }
+  }
+}
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java
index 6bdbdbc642e..a856ac7b942 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java
@@ -111,6 +111,21 @@ public class ConsensusPrefetchingQueue {
 
   private final SubscriptionWalRetentionPolicy retentionPolicy;
 
+  /**
+   * Earliest WAL file version still needed by this consumer group's committed 
progress. Topic
+   * retention and this boundary are merged conservatively by IoTConsensus 
before deletion.
+   */
+  private volatile long committedRetainedMinVersionId = 0L;
+
+  private final Object committedRetentionLock = new Object();
+
+  private RegionProgress lastCommittedProgressForRetention;
+
+  private File lastRetainedWalFileForRetention;
+
+  private final Map<Long, WalFileCommitRequirement> walFileCommitRequirements =
+      new ConcurrentHashMap<>();
+
   private final WakeableIndexedConsensusQueue pendingEntries;
 
   private static final int PENDING_QUEUE_CAPACITY = 4096;
@@ -536,7 +551,8 @@ public class ConsensusPrefetchingQueue {
     this.pendingEntries =
         new WakeableIndexedConsensusQueue(
             PENDING_QUEUE_CAPACITY, this::requestPrefetch, 
this::canAcceptRealtimeEntry);
-    serverImpl.registerSubscriptionQueue(pendingEntries, retentionPolicy);
+    serverImpl.registerSubscriptionQueue(
+        pendingEntries, retentionPolicy, 
this::getCommittedRetainedMinVersionId);
 
     LOGGER.info(
         DataNodePipeMessages
@@ -790,6 +806,7 @@ public class ConsensusPrefetchingQueue {
     this.observedSeekGeneration = seekGeneration.get();
     discardBatch(this.lingerBatch);
     resetBatchWriterProgress();
+    refreshCommittedWalRetentionBoundAndNotify();
 
     LOGGER.info(
         DataNodePipeMessages
@@ -1104,7 +1121,7 @@ public class ConsensusPrefetchingQueue {
             compareWriterProgress(candidate, existing) > 0 ? candidate : 
existing);
   }
 
-  private int compareWriterProgress(
+  private static int compareWriterProgress(
       final WriterProgress leftProgress, final WriterProgress rightProgress) {
     int cmp = Long.compare(leftProgress.getPhysicalTime(), 
rightProgress.getPhysicalTime());
     if (cmp != 0) {
@@ -2094,8 +2111,7 @@ public class ConsensusPrefetchingQueue {
     commitManager.recordMapping(
         consumerGroupId, topicName, consensusGroupId, writerId, 
writerProgress);
     if (tablets.isEmpty()) {
-      return commitManager.commit(
-          consumerGroupId, topicName, consensusGroupId, writerId, 
writerProgress);
+      return commitAndRefreshWalRetention(writerId, writerProgress);
     }
 
     // nextOffset <= 0 means all tablets delivered in single batch
@@ -2304,8 +2320,7 @@ public class ConsensusPrefetchingQueue {
       boolean committed = false;
       try {
         committed =
-            commitManager.commitWithoutOutstanding(
-                consumerGroupId, topicName, consensusGroupId, commitWriterId, 
commitWriterProgress);
+            commitWithoutOutstandingAndRefreshWalRetention(commitWriterId, 
commitWriterProgress);
       } finally {
         if (!committed
             && Objects.nonNull(event)
@@ -2710,12 +2725,7 @@ public class ConsensusPrefetchingQueue {
           }
 
           final boolean committed =
-              commitManager.commit(
-                  consumerGroupId,
-                  topicName,
-                  consensusGroupId,
-                  commitWriterId,
-                  commitWriterProgress);
+              commitAndRefreshWalRetention(commitWriterId, 
commitWriterProgress);
           if (!committed) {
             if (!silent) {
               LOGGER.warn(
@@ -3230,6 +3240,7 @@ public class ConsensusPrefetchingQueue {
     // entry so seek/rebind resumes from the intended frontier.
     commitManager.resetState(
         consumerGroupId, topicName, consensusGroupId, 
request.committedRegionProgress);
+    refreshCommittedWalRetentionBoundAndNotify();
 
     LOGGER.info(
         DataNodePipeMessages
@@ -3243,6 +3254,122 @@ public class ConsensusPrefetchingQueue {
         seekGeneration.get());
   }
 
+  private boolean commitAndRefreshWalRetention(
+      final WriterId writerId, final WriterProgress writerProgress) {
+    final boolean committed =
+        commitManager.commit(
+            consumerGroupId, topicName, consensusGroupId, writerId, 
writerProgress);
+    if (committed) {
+      refreshCommittedWalRetentionBoundAndNotify();
+    }
+    return committed;
+  }
+
+  private boolean commitWithoutOutstandingAndRefreshWalRetention(
+      final WriterId writerId, final WriterProgress writerProgress) {
+    final boolean committed =
+        commitManager.commitWithoutOutstanding(
+            consumerGroupId, topicName, consensusGroupId, writerId, 
writerProgress);
+    if (committed) {
+      refreshCommittedWalRetentionBoundAndNotify();
+    }
+    return committed;
+  }
+
+  private long getCommittedRetainedMinVersionId() {
+    refreshCommittedWalRetentionBound();
+    return committedRetainedMinVersionId;
+  }
+
+  private void refreshCommittedWalRetentionBoundAndNotify() {
+    if (refreshCommittedWalRetentionBound()) {
+      serverImpl.checkAndUpdateSafeDeletedSearchIndex();
+    }
+  }
+
+  private boolean refreshCommittedWalRetentionBound() {
+    final RegionProgress committedRegionProgress =
+        commitManager.getCommittedRegionProgress(consumerGroupId, topicName, 
consensusGroupId);
+
+    synchronized (committedRetentionLock) {
+      if (Objects.equals(lastCommittedProgressForRetention, 
committedRegionProgress)
+          && Objects.nonNull(lastRetainedWalFileForRetention)
+          && lastRetainedWalFileForRetention.exists()) {
+        return false;
+      }
+
+      final CommittedWalRetentionBound newRetentionBound =
+          computeCommittedRetainedMinVersionId(committedRegionProgress);
+      final long newRetainedMinVersionId = 
newRetentionBound.retainedMinVersionId;
+      final boolean changed = committedRetainedMinVersionId != 
newRetainedMinVersionId;
+      committedRetainedMinVersionId = newRetainedMinVersionId;
+      lastCommittedProgressForRetention = committedRegionProgress;
+      lastRetainedWalFileForRetention = newRetentionBound.retainedWalFile;
+      walFileCommitRequirements.keySet().removeIf(versionId -> versionId < 
newRetainedMinVersionId);
+      return changed;
+    }
+  }
+
+  private CommittedWalRetentionBound computeCommittedRetainedMinVersionId(
+      final RegionProgress committedRegionProgress) {
+    if (!(consensusReqReader instanceof WALNode)) {
+      return new CommittedWalRetentionBound(0L, null);
+    }
+
+    final WALNode walNode = (WALNode) consensusReqReader;
+    final long currentWalVersion = walNode.getCurrentWALFileVersion();
+    final File[] walFiles = 
WALFileUtils.listAllWALFiles(walNode.getLogDirectory());
+    if (Objects.isNull(walFiles) || walFiles.length == 0) {
+      return new CommittedWalRetentionBound(Math.max(0L, currentWalVersion), 
null);
+    }
+
+    WALFileUtils.ascSortByVersionId(walFiles);
+    for (final File walFile : walFiles) {
+      final long versionId = WALFileUtils.parseVersionId(walFile.getName());
+      if (versionId >= currentWalVersion) {
+        return new CommittedWalRetentionBound(Math.max(0L, currentWalVersion), 
walFile);
+      }
+      if (ProgressWALIterator.isHeaderOnlyWalFile(walFile)) {
+        continue;
+      }
+
+      WalFileCommitRequirement requirement = 
walFileCommitRequirements.get(versionId);
+      if (Objects.isNull(requirement)) {
+        try (final ProgressWALReader reader = openProgressWALReader(walFile)) {
+          requirement =
+              WalFileCommitRequirement.fromMetadata(
+                  consensusGroupId.toString(), reader.getMetaData());
+          walFileCommitRequirements.put(versionId, requirement);
+        } catch (final IOException e) {
+          LOGGER.warn(
+              DataNodePipeMessages
+                  
.PIPE_LOG_CONSENSUSPREFETCHINGQUEUE_FAILED_TO_READ_WAL_METADATA_FROM_A2ED50D1,
+              this,
+              walFile,
+              e);
+          return new CommittedWalRetentionBound(versionId, walFile);
+        }
+      }
+
+      if (!requirement.isCoveredBy(committedRegionProgress)) {
+        return new CommittedWalRetentionBound(versionId, walFile);
+      }
+    }
+    return new CommittedWalRetentionBound(Math.max(0L, currentWalVersion), 
null);
+  }
+
+  private static final class CommittedWalRetentionBound {
+
+    private final long retainedMinVersionId;
+    private final File retainedWalFile;
+
+    private CommittedWalRetentionBound(
+        final long retainedMinVersionId, final File retainedWalFile) {
+      this.retainedMinVersionId = retainedMinVersionId;
+      this.retainedWalFile = retainedWalFile;
+    }
+  }
+
   public RegionProgress computeTailRegionProgress() {
     if (!(consensusReqReader instanceof WALNode)) {
       return new RegionProgress(Collections.emptyMap());
@@ -3307,6 +3434,77 @@ public class ConsensusPrefetchingQueue {
     }
   }
 
+  static final class WalFileCommitRequirement {
+
+    private final Map<WriterId, WriterProgress> requiredWriterProgress;
+    private final boolean containsUnsupportedProgress;
+
+    private WalFileCommitRequirement(
+        final Map<WriterId, WriterProgress> requiredWriterProgress,
+        final boolean containsUnsupportedProgress) {
+      this.requiredWriterProgress = requiredWriterProgress;
+      this.containsUnsupportedProgress = containsUnsupportedProgress;
+    }
+
+    static WalFileCommitRequirement fromMetadata(
+        final String regionId, final WALMetaData metadata) {
+      if (Objects.isNull(metadata)) {
+        return new WalFileCommitRequirement(Collections.emptyMap(), true);
+      }
+
+      final List<Integer> buffersSize = metadata.getBuffersSize();
+      final List<Long> physicalTimes = metadata.getPhysicalTimes();
+      final List<Short> nodeIds = metadata.getNodeIds();
+      final List<Long> localSeqs = metadata.getLocalSeqs();
+      if (physicalTimes.size() < buffersSize.size()
+          || nodeIds.size() < buffersSize.size()
+          || localSeqs.size() < buffersSize.size()) {
+        return new WalFileCommitRequirement(Collections.emptyMap(), true);
+      }
+
+      final Map<WriterId, WriterProgress> requiredWriterProgress = new 
LinkedHashMap<>();
+      for (int i = 0; i < buffersSize.size(); i++) {
+        final int writerNodeId = nodeIds.get(i);
+        final long physicalTime = physicalTimes.get(i);
+        final long localSeq = localSeqs.get(i);
+        if (writerNodeId < 0 && physicalTime == 0L && localSeq < 0L) {
+          // Non-search WAL entries do not participate in subscription 
progress.
+          continue;
+        }
+        if (writerNodeId < 0 || physicalTime < 0L || localSeq < 0L) {
+          // Legacy or incomplete writer metadata cannot be compared safely 
with RegionProgress.
+          return new WalFileCommitRequirement(Collections.emptyMap(), true);
+        }
+
+        final WriterId writerId = new WriterId(regionId, writerNodeId);
+        final WriterProgress candidateProgress = new 
WriterProgress(physicalTime, localSeq);
+        requiredWriterProgress.merge(
+            writerId,
+            candidateProgress,
+            (currentProgress, candidate) ->
+                compareWriterProgress(candidate, currentProgress) > 0
+                    ? candidate
+                    : currentProgress);
+      }
+      return new WalFileCommitRequirement(requiredWriterProgress, false);
+    }
+
+    boolean isCoveredBy(final RegionProgress committedRegionProgress) {
+      if (containsUnsupportedProgress || 
Objects.isNull(committedRegionProgress)) {
+        return false;
+      }
+      for (final Map.Entry<WriterId, WriterProgress> entry : 
requiredWriterProgress.entrySet()) {
+        final WriterProgress committedWriterProgress =
+            committedRegionProgress.getWriterPositions().get(entry.getKey());
+        if (Objects.isNull(committedWriterProgress)
+            || compareWriterProgress(committedWriterProgress, 
entry.getValue()) < 0) {
+          return false;
+        }
+      }
+      return true;
+    }
+  }
+
   private void mergeTailProgress(
       final Map<WriterId, WriterProgress> tailProgressByWriter, final 
WALMetaData metadata) {
     if (Objects.isNull(metadata)) {
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java
index da73a003808..a2907745179 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java
@@ -106,6 +106,52 @@ public class ConsensusPrefetchingQueueTest {
                 .getModifiers()));
   }
 
+  @Test
+  public void testWalFileCommitRequirementUsesPerWriterMaximum() {
+    final String regionId = "DataRegion[1]";
+    final WALMetaData metadata = new WALMetaData();
+    metadata.add(1, -1L, 0L, 0L, -1, -1L);
+    metadata.add(1, 1L, 0L, 100L, 7, 1L);
+    metadata.add(1, 2L, 0L, 100L, 7, 2L);
+    metadata.add(1, -1L, 0L, 200L, 8, 1L);
+
+    final ConsensusPrefetchingQueue.WalFileCommitRequirement requirement =
+        
ConsensusPrefetchingQueue.WalFileCommitRequirement.fromMetadata(regionId, 
metadata);
+
+    assertFalse(
+        requirement.isCoveredBy(
+            new RegionProgress(
+                Collections.singletonMap(
+                    new WriterId(regionId, 7), new WriterProgress(100L, 
2L)))));
+    assertFalse(
+        requirement.isCoveredBy(
+            new RegionProgress(
+                java.util.Map.of(
+                    new WriterId(regionId, 7),
+                    new WriterProgress(100L, 1L),
+                    new WriterId(regionId, 8),
+                    new WriterProgress(200L, 1L)))));
+    assertTrue(
+        requirement.isCoveredBy(
+            new RegionProgress(
+                java.util.Map.of(
+                    new WriterId(regionId, 7),
+                    new WriterProgress(100L, 2L),
+                    new WriterId(regionId, 8),
+                    new WriterProgress(200L, 1L)))));
+  }
+
+  @Test
+  public void testWalFileCommitRequirementRejectsUnsupportedWriterMetadata() {
+    final WALMetaData metadata = new WALMetaData();
+    metadata.add(1, 1L, 0L, 0L, -1, 1L);
+
+    final ConsensusPrefetchingQueue.WalFileCommitRequirement requirement =
+        
ConsensusPrefetchingQueue.WalFileCommitRequirement.fromMetadata("DataRegion[1]",
 metadata);
+
+    assertFalse(requirement.isCoveredBy(new 
RegionProgress(Collections.emptyMap())));
+  }
+
   @Test
   public void testReplayStartPreservesUncoveredFollowerEntries() throws 
Exception {
     final String originalSystemDir = 
IoTDBDescriptor.getInstance().getConfig().getSystemDir();

Reply via email to