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

luwei16 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/master by this push:
     new 3dbe4d3ca53 [fix](cloud) Invalidate version caches on visible commit 
retries (#67813)
3dbe4d3ca53 is described below

commit 3dbe4d3ca53307f73db96eda23d4efc9f02963c5
Author: Luwei <[email protected]>
AuthorDate: Wed Sep 16 12:42:59 2026 +0800

    [fix](cloud) Invalidate version caches on visible commit retries (#67813)
    
    ### What problem does this PR solve?
    
    Issue Number: close #67099
    
    Related PR: None
    
    Problem Summary: After a cloud commit response is lost, a retry can
    return a VISIBLE transaction without version results. FE can then keep
    using stale partition versions and the table version used to validate
    SQL caches.
    
    Invalidate the affected table and partition version caches, including
    the MoW already-visible shortcut, and notify other FEs through the
    existing version synchronization RPC. Subsequent reads fetch fresh
    versions from MS. Commit processing performs no extra MS version query,
    so a version-service outage does not turn an already committed
    transaction into a commit error or require replaying its callbacks.
    
    Publish batch cache updates under the tables' version write locks.
    Invalidation epochs prevent old in-flight MS responses and delayed
    commit notifications from restoring stale caches. The synchronization
    RPC uses table version -1 for invalidation; older FEs ignore it and
    retain their existing periodic synchronization. Peer notification keeps
    its existing asynchronous behavior and configuration.
    
    ### Release note
    
    Prevent stale reads and false commit failures after cloud commit
    retries.
    
    ### Check List (For Author)
    
    - Test: Unit Test
    - 58 FE tests passed in both ordinary and JaCoCo coverage modes via
    run-fe-ut.sh (CloudGlobalTransactionMgrTest, CloudPartitionTest,
    OlapTableTest, VersionHelperTest); FE Checkstyle: 0 violations; PR diff
    whitespace check passed.
    - Tests cover invalidation and batch-publication locks, table/partition
    cache recovery, old RPC completions, successful commit callbacks during
    a version outage, and FE notification serialization/receiver handling.
        - No live cluster deployment or SQL regression run.
    - Behavior changed: Yes. VISIBLE retries without versions invalidate
    local and peer caches; subsequent reads retrieve versions without
    changing the successful commit result.
    - Does this need documentation: No
---
 .../java/org/apache/doris/catalog/OlapTable.java   |  20 +-
 .../cloud/catalog/CloudFEVersionSynchronizer.java  |  59 +++
 .../apache/doris/cloud/catalog/CloudPartition.java |  46 ++-
 .../transaction/CloudGlobalTransactionMgr.java     |  11 +
 .../transaction/CloudGlobalTransactionMgrTest.java | 398 ++++++++++++++++++++-
 5 files changed, 520 insertions(+), 14 deletions(-)

diff --git a/fe/fe-core/src/main/java/org/apache/doris/catalog/OlapTable.java 
b/fe/fe-core/src/main/java/org/apache/doris/catalog/OlapTable.java
index 17cac68dc68..4645a7796c9 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/catalog/OlapTable.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/catalog/OlapTable.java
@@ -248,6 +248,10 @@ public class OlapTable extends Table implements 
MTMVRelatedTableIf, GsonPostProc
     private volatile long lastTableVersionCachedTimeMs = 0;
     private volatile long cachedTableVersion = -1;
 
+    // Commit notifications cannot clear invalidation: their versions may 
precede the missing commit result.
+    private final AtomicLong tableVersionCacheEpoch = new AtomicLong();
+    private final AtomicLong refreshedTableVersionCacheEpoch = new 
AtomicLong();
+
     private ReadWriteLock versionLock = Config.isCloudMode() ? new 
ReentrantReadWriteLock(true) : null;
 
     public OlapTable() {
@@ -3591,7 +3595,7 @@ public class OlapTable extends Table implements 
MTMVRelatedTableIf, GsonPostProc
     @VisibleForTesting
     protected boolean isCachedTableVersionExpired() {
         // -1 means no cache yet, need to fetch from MS
-        if (cachedTableVersion == -1) {
+        if (cachedTableVersion == -1 || tableVersionCacheEpoch.get() != 
refreshedTableVersionCacheEpoch.get()) {
             return true;
         }
         ConnectContext ctx = ConnectContext.get();
@@ -3605,13 +3609,18 @@ public class OlapTable extends Table implements 
MTMVRelatedTableIf, GsonPostProc
 
     public boolean isCachedTableVersionExpired(long expirationMs) {
         // -1 means no cache yet, need to fetch from MS
-        if (cachedTableVersion == -1 || expirationMs <= 0) {
+        if (cachedTableVersion == -1 || expirationMs <= 0
+                || tableVersionCacheEpoch.get() != 
refreshedTableVersionCacheEpoch.get()) {
             return true;
         }
         return System.currentTimeMillis() - lastTableVersionCachedTimeMs > 
expirationMs;
     }
 
-    public void setCachedTableVersion(long version) {
+    public void invalidateCachedTableVersion() {
+        tableVersionCacheEpoch.incrementAndGet();
+    }
+
+    public synchronized void setCachedTableVersion(long version) {
         if (version >= cachedTableVersion) {
             cachedTableVersion = version;
             lastTableVersionCachedTimeMs = System.currentTimeMillis();
@@ -3632,6 +3641,7 @@ public class OlapTable extends Table implements 
MTMVRelatedTableIf, GsonPostProc
             return getCachedTableVersion();
         }
 
+        long cacheEpoch = tableVersionCacheEpoch.get();
         // get version rpc
         Cloud.GetVersionRequest request = Cloud.GetVersionRequest.newBuilder()
                 .setRequestIp(FrontendOptions.getLocalHostAddressCached())
@@ -3658,6 +3668,7 @@ public class OlapTable extends Table implements 
MTMVRelatedTableIf, GsonPostProc
             }
             // update cache
             setCachedTableVersion(version);
+            refreshedTableVersionCacheEpoch.accumulateAndGet(cacheEpoch, 
Math::max);
             return version;
         } catch (RpcException e) {
             LOG.warn("get version from meta service failed", e);
@@ -3722,9 +3733,11 @@ public class OlapTable extends Table implements 
MTMVRelatedTableIf, GsonPostProc
     private static List<Long> getVisibleVersionInBatchFromMs(List<OlapTable> 
tables) {
         List<Long> dbIds = new ArrayList<>(tables.size());
         List<Long> tableIds = new ArrayList<>(tables.size());
+        List<Long> cacheEpochs = new ArrayList<>(tables.size());
         for (OlapTable table : tables) {
             dbIds.add(table.getDatabase().getId());
             tableIds.add(table.getId());
+            cacheEpochs.add(table.tableVersionCacheEpoch.get());
         }
 
         List<Long> versions = getVisibleVersionFromMeta(dbIds, tableIds);
@@ -3733,6 +3746,7 @@ public class OlapTable extends Table implements 
MTMVRelatedTableIf, GsonPostProc
         Preconditions.checkState(tables.size() == versions.size());
         for (int i = 0; i < tables.size(); i++) {
             tables.get(i).setCachedTableVersion(versions.get(i));
+            
tables.get(i).refreshedTableVersionCacheEpoch.accumulateAndGet(cacheEpochs.get(i),
 Math::max);
         }
         return versions;
     }
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudFEVersionSynchronizer.java
 
b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudFEVersionSynchronizer.java
index 5b8f5507a34..6ac2abbbace 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudFEVersionSynchronizer.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudFEVersionSynchronizer.java
@@ -49,6 +49,9 @@ import java.util.stream.Collectors;
 
 public class CloudFEVersionSynchronizer {
     private static final Logger LOG = 
LogManager.getLogger(CloudFEVersionSynchronizer.class);
+    // Reuse the version RPC for whole-table invalidation when a commit 
response has no versions.
+    // Older FEs ignore this value and continue relying on their periodic 
version synchronization.
+    private static final long INVALIDATE_VERSION_CACHE = -1;
 
     private static final ExecutorService SYNC_VERSION_THREAD_POOL = 
Executors.newFixedThreadPool(
             Config.cloud_sync_version_task_threads_num,
@@ -57,6 +60,54 @@ public class CloudFEVersionSynchronizer {
     public CloudFEVersionSynchronizer() {
     }
 
+    public void invalidateVersionCaches(long dbId, List<Long> tableIds) {
+        Database db = Env.getCurrentInternalCatalog().getDbNullable(dbId);
+        if (db == null) {
+            return; // The database may have been dropped since the original 
commit.
+        }
+        List<OlapTable> tables = new ArrayList<>();
+        for (long tableId : tableIds) {
+            Table table = db.getTableNullable(tableId);
+            if (table != null && table.isManagedTable()) {
+                tables.add((OlapTable) table);
+            }
+        }
+        invalidateVersionCaches(tables);
+        List<Pair<OlapTable, Long>> tableVersions = tables.stream()
+                .map(table -> Pair.of(table, 
INVALIDATE_VERSION_CACHE)).collect(Collectors.toList());
+        pushVersionAsync(dbId, tableVersions, Collections.emptyMap());
+    }
+
+    private void invalidateVersionCaches(List<OlapTable> tables) {
+        tables.sort(Comparator.comparingLong(OlapTable::getId));
+        for (OlapTable table : tables) {
+            table.readLock();
+        }
+        try {
+            for (OlapTable table : tables) {
+                table.versionWriteLock();
+            }
+            try {
+                for (OlapTable table : tables) {
+                    table.invalidateCachedTableVersion();
+                    for (Partition partition : table.getAllPartitions()) {
+                        ((CloudPartition) 
partition).invalidateCachedVisibleVersion();
+                    }
+                }
+            } finally {
+                for (int i = tables.size() - 1; i >= 0; i--) {
+                    tables.get(i).versionWriteUnlock();
+                }
+            }
+        } finally {
+            for (int i = tables.size() - 1; i >= 0; i--) {
+                tables.get(i).readUnlock();
+            }
+        }
+        LOG.info("Invalidated cloud version caches for tables {}",
+                
tables.stream().map(OlapTable::getId).collect(Collectors.toList()));
+    }
+
     // master FE send sync version rpc to other FEs
     public void pushVersionAsync(long dbId, OlapTable table, long version) {
         if (!Config.cloud_enable_version_syncer) {
@@ -154,17 +205,25 @@ public class CloudFEVersionSynchronizer {
         }
         // only update table version
         if (request.getPartitionVersionInfos().isEmpty()) {
+            List<OlapTable> invalidatedTables = new ArrayList<>();
             request.getTableVersionInfos().forEach(tableVersionInfo -> {
                 Table table = 
db.getTableNullable(tableVersionInfo.getTableId());
                 if (table == null || !table.isManagedTable()) {
                     return;
                 }
                 OlapTable olapTable = (OlapTable) table;
+                if (tableVersionInfo.getVersion() == INVALIDATE_VERSION_CACHE) 
{
+                    invalidatedTables.add(olapTable);
+                    return;
+                }
                 olapTable.setCachedTableVersion(tableVersionInfo.getVersion());
                 if (LOG.isDebugEnabled()) {
                     LOG.debug("Update tableId: {}, version: {}", 
olapTable.getId(), tableVersionInfo.getVersion());
                 }
             });
+            if (!invalidatedTables.isEmpty()) {
+                invalidateVersionCaches(invalidatedTables);
+            }
             return;
         }
         // update partition and table version
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudPartition.java 
b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudPartition.java
index 86c3b72409c..07c2a6d504e 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudPartition.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/CloudPartition.java
@@ -48,6 +48,7 @@ import java.util.Comparator;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
+import java.util.concurrent.atomic.AtomicLong;
 import java.util.concurrent.locks.ReentrantLock;
 import java.util.stream.Collectors;
 
@@ -66,6 +67,10 @@ public class CloudPartition extends Partition {
     // This value is set when get the version from meta-service, 0 means 
version is not cached yet
     private volatile long lastVersionCachedTimeMs = 0;
 
+    // Only an MS read started after invalidation may make this cache valid 
again.
+    private final AtomicLong versionCacheEpoch = new AtomicLong();
+    private final AtomicLong refreshedVersionCacheEpoch = new AtomicLong();
+
     private ReentrantLock lock = new ReentrantLock(true);
 
     public CloudPartition(long id, String name, MaterializedIndex baseIndex,
@@ -127,9 +132,13 @@ public class CloudPartition extends Partition {
         return super.getVisibleVersion();
     }
 
+    public void invalidateCachedVisibleVersion() {
+        versionCacheEpoch.incrementAndGet();
+    }
+
     @VisibleForTesting
     protected boolean isCachedVersionExpired() {
-        if (lastVersionCachedTimeMs == 0) {
+        if (lastVersionCachedTimeMs == 0 || versionCacheEpoch.get() != 
refreshedVersionCacheEpoch.get()) {
             return true;
         }
         ConnectContext ctx = ConnectContext.get();
@@ -155,6 +164,7 @@ public class CloudPartition extends Partition {
     }
 
     private long getVisibleVersionFromMs(boolean waitForPendingTxns) {
+        long cacheEpoch = versionCacheEpoch.get();
         if (LOG.isDebugEnabled()) {
             LOG.debug("getVisibleVersionFromMs use CloudPartition {}, 
waitForPendingTxns: {}",
                     super.getName(), waitForPendingTxns);
@@ -183,6 +193,7 @@ public class CloudPartition extends Partition {
                             version, tso, super.getId());
                 }
                 setCachedVisibleVersion(version, mTime, tso);
+                refreshedVersionCacheEpoch.accumulateAndGet(cacheEpoch, 
Math::max);
                 return version;
             } else {
                 assert resp.getStatus().getCode() == 
MetaServiceCode.VERSION_NOT_FOUND;
@@ -193,6 +204,7 @@ public class CloudPartition extends Partition {
                 LOG.debug("get version from meta service, version: {}, 
partition: {}", version, super.getId());
             }
             setCachedVisibleVersion(version, mTime);
+            refreshedVersionCacheEpoch.accumulateAndGet(cacheEpoch, Math::max);
             return version;
         } catch (RpcException e) {
             throw new RuntimeException("get version from meta service failed");
@@ -254,10 +266,12 @@ public class CloudPartition extends Partition {
         List<Long> partitionIds = new ArrayList<>();
         List<Long> versionUpdateTimesMs = new ArrayList<>();
         List<Long> commitTsos = new ArrayList<>();
+        List<Long> cacheEpochs = new ArrayList<>();
         for (CloudPartition partition : partitions) {
             dbIds.add(partition.getDbId());
             tableIds.add(partition.getTableId());
             partitionIds.add(partition.getId());
+            cacheEpochs.add(partition.versionCacheEpoch.get());
         }
 
         List<Long> versions = getSnapshotVisibleVersion(
@@ -266,15 +280,27 @@ public class CloudPartition extends Partition {
         // Cache visible version, see hasData() for details.
         int size = versions.size();
         boolean hasCommitTsos = commitTsos.size() == size;
-        for (int i = 0; i < size; ++i) {
-            Long version = versions.get(i);
-            if (version > Partition.PARTITION_INIT_VERSION) {
-                // For compatibility, the existing partitions may not have 
mtime
-                long mTime = versions.size() == versionUpdateTimesMs.size() ? 
versionUpdateTimesMs.get(i) : 0;
-                long tso = hasCommitTsos ? commitTsos.get(i) : -1;
-                partitions.get(i).setCachedVisibleVersion(versions.get(i), 
mTime, tso);
-            } else { // No data has been written to this partition
-                
partitions.get(i).setCachedVisibleVersion(Partition.PARTITION_INIT_VERSION, 
System.currentTimeMillis());
+        List<OlapTable> tables = getTables(partitions);
+        for (OlapTable table : tables) {
+            table.versionWriteLock();
+        }
+        try {
+            for (int i = 0; i < size; ++i) {
+                Long version = versions.get(i);
+                if (version > Partition.PARTITION_INIT_VERSION) {
+                    // For compatibility, the existing partitions may not have 
mtime
+                    long mTime = versions.size() == 
versionUpdateTimesMs.size() ? versionUpdateTimesMs.get(i) : 0;
+                    long tso = hasCommitTsos ? commitTsos.get(i) : -1;
+                    partitions.get(i).setCachedVisibleVersion(versions.get(i), 
mTime, tso);
+                } else { // No data has been written to this partition
+                    
partitions.get(i).setCachedVisibleVersion(Partition.PARTITION_INIT_VERSION,
+                            System.currentTimeMillis());
+                }
+                
partitions.get(i).refreshedVersionCacheEpoch.accumulateAndGet(cacheEpochs.get(i),
 Math::max);
+            }
+        } finally {
+            for (int i = tables.size() - 1; i >= 0; i--) {
+                tables.get(i).versionWriteUnlock();
             }
         }
 
diff --git 
a/fe/fe-core/src/main/java/org/apache/doris/cloud/transaction/CloudGlobalTransactionMgr.java
 
b/fe/fe-core/src/main/java/org/apache/doris/cloud/transaction/CloudGlobalTransactionMgr.java
index 52b7a84bea2..0aa654c80fc 100644
--- 
a/fe/fe-core/src/main/java/org/apache/doris/cloud/transaction/CloudGlobalTransactionMgr.java
+++ 
b/fe/fe-core/src/main/java/org/apache/doris/cloud/transaction/CloudGlobalTransactionMgr.java
@@ -491,6 +491,10 @@ public class CloudGlobalTransactionMgr implements 
GlobalTransactionMgrIface {
                         + "] is already aborted. abort reason: " + 
transactionState.getReason());
             } else if (transactionState.getTransactionStatus() == 
TransactionStatus.COMMITTED
                     || transactionState.getTransactionStatus() == 
TransactionStatus.VISIBLE) {
+                if (transactionState.getTransactionStatus() == 
TransactionStatus.VISIBLE) {
+                    ((CloudEnv) 
Env.getCurrentEnv()).getCloudFEVersionSynchronizer()
+                            .invalidateVersionCaches(dbId, 
transactionState.getTableIdList());
+                }
                 LOG.info("txn={}, status={} not need to calculate delete 
bitmap again, just return ",
                         transactionId,
                         transactionState.getTransactionStatus().toString());
@@ -597,6 +601,13 @@ public class CloudGlobalTransactionMgr implements 
GlobalTransactionMgrIface {
                 ? commitTxnResponse.getTxnInfo().getCommitTso() : -1;
         int totalPartitionNum = commitTxnResponse.getPartitionIdsList().size();
         if (totalPartitionNum == 0 && 
commitTxnResponse.getTableStatsList().isEmpty()) {
+            // A visible retry can contain only txn_info. Invalidate locally 
without another commit-time RPC;
+            // subsequent reads fetch fresh versions, including the table 
version used to validate SQL caches.
+            if (commitTxnResponse.getTxnInfo().getStatus() == 
TxnStatusPB.TXN_STATUS_VISIBLE
+                    && !(commitTxnResponse.getIsLazyCommit() && 
commitTxnResponse.getIsLazyCommitIncomplete())) {
+                ((CloudEnv) 
Env.getCurrentEnv()).getCloudFEVersionSynchronizer()
+                        .invalidateVersionCaches(dbId, 
commitTxnResponse.getTxnInfo().getTableIdsList());
+            }
             return Collections.emptyMap();
         }
         Env env = Env.getCurrentEnv();
diff --git 
a/fe/fe-core/src/test/java/org/apache/doris/cloud/transaction/CloudGlobalTransactionMgrTest.java
 
b/fe/fe-core/src/test/java/org/apache/doris/cloud/transaction/CloudGlobalTransactionMgrTest.java
index 3d200ce80d3..b2927e26eba 100644
--- 
a/fe/fe-core/src/test/java/org/apache/doris/cloud/transaction/CloudGlobalTransactionMgrTest.java
+++ 
b/fe/fe-core/src/test/java/org/apache/doris/cloud/transaction/CloudGlobalTransactionMgrTest.java
@@ -21,9 +21,17 @@ import org.apache.doris.catalog.CatalogTestUtil;
 import org.apache.doris.catalog.Env;
 import org.apache.doris.catalog.FakeEditLog;
 import org.apache.doris.catalog.FakeEnv;
+import org.apache.doris.catalog.KeysType;
+import org.apache.doris.catalog.MaterializedIndex;
+import org.apache.doris.catalog.OlapTable;
+import org.apache.doris.catalog.PartitionInfo;
+import org.apache.doris.catalog.RandomDistributionInfo;
 import org.apache.doris.catalog.Table;
 import org.apache.doris.catalog.stream.CloudOlapTableStreamUpdate;
 import org.apache.doris.catalog.stream.TableStreamUpdateInfo;
+import org.apache.doris.cloud.catalog.CloudEnv;
+import org.apache.doris.cloud.catalog.CloudFEVersionSynchronizer;
+import org.apache.doris.cloud.catalog.CloudPartition;
 import org.apache.doris.cloud.proto.Cloud;
 import org.apache.doris.cloud.proto.Cloud.AbortTxnResponse;
 import org.apache.doris.cloud.proto.Cloud.BeginTxnResponse;
@@ -33,16 +41,32 @@ import 
org.apache.doris.cloud.proto.Cloud.GetCurrentMaxTxnResponse;
 import org.apache.doris.cloud.proto.Cloud.MetaServiceCode;
 import org.apache.doris.cloud.proto.Cloud.TxnInfoPB;
 import org.apache.doris.cloud.rpc.MetaServiceProxy;
+import org.apache.doris.cloud.rpc.VersionHelper;
+import org.apache.doris.common.ClientPool;
 import org.apache.doris.common.Config;
 import org.apache.doris.common.DuplicatedRequestException;
 import org.apache.doris.common.FeMetaVersion;
+import org.apache.doris.common.GenericPool;
 import org.apache.doris.common.LabelAlreadyUsedException;
 import org.apache.doris.common.UserException;
+import org.apache.doris.common.jmockit.Deencapsulation;
 import org.apache.doris.load.routineload.RLTaskTxnCommitAttachment;
+import org.apache.doris.qe.ConnectContext;
+import org.apache.doris.qe.SessionVariable;
+import org.apache.doris.rpc.RpcException;
+import org.apache.doris.service.FrontendServiceImpl;
+import org.apache.doris.system.Frontend;
+import org.apache.doris.system.SystemInfoService.HostInfo;
+import org.apache.doris.thrift.FrontendService;
+import org.apache.doris.thrift.TCloudVersionInfo;
+import org.apache.doris.thrift.TFrontendSyncCloudVersionRequest;
+import org.apache.doris.thrift.TStatus;
+import org.apache.doris.thrift.TStatusCode;
 import org.apache.doris.thrift.TTabletCommitInfo;
 import org.apache.doris.transaction.GlobalTransactionMgrIface;
 import org.apache.doris.transaction.TabletCommitInfo;
 import org.apache.doris.transaction.TransactionState;
+import org.apache.doris.transaction.TransactionStatus;
 import org.apache.doris.transaction.TxnStateChangeCallback;
 
 import com.google.common.collect.Lists;
@@ -50,6 +74,7 @@ import org.junit.jupiter.api.AfterEach;
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
+import org.mockito.AdditionalAnswers;
 import org.mockito.ArgumentCaptor;
 import org.mockito.MockedStatic;
 import org.mockito.Mockito;
@@ -57,8 +82,15 @@ import org.mockito.Mockito;
 import java.lang.reflect.Method;
 import java.util.List;
 import java.util.Map;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
 import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.concurrent.atomic.AtomicLong;
+import java.util.concurrent.atomic.AtomicReference;
 
 public class CloudGlobalTransactionMgrTest {
 
@@ -77,13 +109,19 @@ public class CloudGlobalTransactionMgrTest {
         Config.meta_service_endpoint = "127.0.0.1:20121";
         fakeEditLog = new FakeEditLog();
         fakeEnv = new FakeEnv();
-        masterEnv = CatalogTestUtil.createTestCatalog();
+        Env catalog = CatalogTestUtil.createTestCatalog();
+        masterEnv = Mockito.mock(CloudEnv.class, 
AdditionalAnswers.delegatesTo(catalog));
+        Mockito.doReturn(new CloudFEVersionSynchronizer()).when((CloudEnv) 
masterEnv).getCloudFEVersionSynchronizer();
+        // Env.getCurrentGlobalTransactionMgr reads the field directly rather 
than calling the delegated getter.
+        Deencapsulation.setField(masterEnv, "globalTransactionMgr", 
catalog.getGlobalTransactionMgr());
+        FakeEnv.setEnv(masterEnv);
         FakeEnv.setMetaVersion(FeMetaVersion.VERSION_CURRENT);
         masterTransMgr = masterEnv.getGlobalTransactionMgr();
     }
 
     @AfterEach
     public void tearDown() {
+        ConnectContext.remove();
         if (fakeEnv != null) {
             fakeEnv.close();
         }
@@ -641,6 +679,364 @@ public class CloudGlobalTransactionMgrTest {
         }
     }
 
+    @Test
+    public void testVisibleRetryInvalidatesTableAndPartitionVersions() throws 
Exception {
+        useVersionCaches();
+        CloudPartition first = addCloudPartition(1000);
+        CloudPartition second = addCloudPartition(2000);
+        OlapTable firstTable = getCloudTable(first);
+        OlapTable secondTable = getCloudTable(second);
+        CommitTxnResponse response = visibleRetry(List.of(first.getTableId(), 
second.getTableId()));
+        try (MockedStatic<VersionHelper> versions = mockVersionHelper()) {
+            versions.when(() -> 
VersionHelper.getVersionFromMeta(Mockito.any())).thenReturn(partitionVersion(4));
+            masterTransMgr.afterCommitTxnResp(response, null, List.of());
+            versions.verifyNoInteractions();
+            Assertions.assertEquals(2, first.getCachedVisibleVersion());
+            Assertions.assertEquals(2, firstTable.getCachedTableVersion());
+            Assertions.assertEquals(4, first.getVisibleVersion());
+            Assertions.assertEquals(List.of(4L), 
CloudPartition.getSnapshotVisibleVersion(List.of(second)));
+            Assertions.assertEquals(4, firstTable.getVisibleVersion());
+            Assertions.assertEquals(List.of(4L), 
OlapTable.getVisibleVersionInBatch(List.of(secondTable)));
+            Assertions.assertEquals(4, first.getCachedVisibleVersion());
+            Assertions.assertEquals(4, second.getCachedVisibleVersion());
+            Assertions.assertEquals(40, first.getVisibleVersionTime());
+            Assertions.assertEquals(400, first.getTso());
+            // Current versions must not become the original transaction's BE 
promotion outcome.
+            Assertions.assertEquals(0, response.getVersionsCount());
+            Assertions.assertEquals(4, first.getVisibleVersion());
+            Assertions.assertEquals(List.of(4L), 
CloudPartition.getSnapshotVisibleVersion(List.of(second)));
+            Assertions.assertEquals(4, firstTable.getVisibleVersion());
+            Assertions.assertEquals(List.of(4L), 
OlapTable.getVisibleVersionInBatch(List.of(secondTable)));
+            versions.verify(() -> 
VersionHelper.getVersionFromMeta(Mockito.any()), Mockito.times(4));
+        }
+    }
+
+    @Test
+    public void testIncompleteCommitDoesNotRefreshPartitionVersions() throws 
Exception {
+        CloudPartition partition = addCloudPartition(1000);
+        CommitTxnResponse response = 
visibleRetry(List.of(partition.getTableId()));
+        try (MockedStatic<VersionHelper> versions = mockVersionHelper()) {
+            
masterTransMgr.afterCommitTxnResp(response.toBuilder().setTxnInfo(response.getTxnInfo().toBuilder()
+                    
.setStatus(Cloud.TxnStatusPB.TXN_STATUS_COMMITTED)).build(), null, List.of());
+            
masterTransMgr.afterCommitTxnResp(response.toBuilder().setIsLazyCommit(true)
+                    .setIsLazyCommitIncomplete(true).build(), null, List.of());
+            Assertions.assertEquals(2, partition.getCachedVisibleVersion());
+            versions.verifyNoInteractions();
+        }
+    }
+
+    @Test
+    public void testVisibleRetrySkipsDroppedTable() throws Exception {
+        try (MockedStatic<VersionHelper> versions = mockVersionHelper()) {
+            masterTransMgr.afterCommitTxnResp(visibleRetry(List.of(1000L)), 
null, List.of());
+            versions.verifyNoInteractions();
+        }
+    }
+
+    @Test
+    public void testMowVisiblePrecheckInvalidatesVersions() throws Exception {
+        useVersionCaches();
+        CloudPartition partition = addCloudPartition(1000);
+        TxnInfoPB txnInfo = 
visibleRetry(List.of(partition.getTableId())).getTxnInfo();
+        MetaServiceProxy proxy = Mockito.mock(MetaServiceProxy.class);
+        
Mockito.when(proxy.getTxn(Mockito.any())).thenReturn(Cloud.GetTxnResponse.newBuilder()
+                
.setStatus(Cloud.MetaServiceResponseStatus.newBuilder().setCode(MetaServiceCode.OK))
+                .setTxnInfo(txnInfo).build());
+        try (MockedStatic<MetaServiceProxy> proxyMock = 
Mockito.mockStatic(MetaServiceProxy.class);
+                MockedStatic<VersionHelper> versions = mockVersionHelper()) {
+            proxyMock.when(MetaServiceProxy::getInstance).thenReturn(proxy);
+            versions.when(() -> 
VersionHelper.getVersionFromMeta(Mockito.any())).thenReturn(partitionVersion(4));
+            Method check = 
CloudGlobalTransactionMgr.class.getDeclaredMethod("checkTransactionStateBeforeCommit",
+                    long.class, long.class);
+            check.setAccessible(true);
+            Assertions.assertEquals(false, check.invoke(masterTransMgr, 
txnInfo.getDbId(), txnInfo.getTxnId()));
+            versions.verifyNoInteractions();
+            Assertions.assertEquals(4, partition.getVisibleVersion());
+            Assertions.assertEquals(4, 
getCloudTable(partition).getVisibleVersion());
+            Assertions.assertEquals(4, partition.getCachedVisibleVersion());
+        }
+    }
+
+    @Test
+    public void testRefreshFailurePreservesCommitCallbacks() throws Exception {
+        useVersionCaches();
+        CloudPartition partition = addCloudPartition(1000);
+        CommitTxnResponse retry = 
visibleRetry(List.of(partition.getTableId()));
+        CommitTxnResponse response = retry.toBuilder()
+                
.setTxnInfo(retry.getTxnInfo().toBuilder().setListenerId(42)).build();
+        TxnStateChangeCallback callback = 
Mockito.mock(TxnStateChangeCallback.class);
+        Mockito.when(callback.getId()).thenReturn(42L);
+        masterTransMgr.getCallbackFactory().addCallback(callback);
+        Mockito.clearInvocations(callback);
+        MetaServiceProxy proxy = Mockito.mock(MetaServiceProxy.class);
+        Mockito.when(proxy.commitTxn(Mockito.any())).thenReturn(response);
+        Table table = 
masterEnv.getInternalCatalog().getDbOrMetaException(CatalogTestUtil.testDbId1)
+                .getTableOrMetaException(partition.getTableId());
+        try (MockedStatic<MetaServiceProxy> proxyMock = 
Mockito.mockStatic(MetaServiceProxy.class);
+                MockedStatic<VersionHelper> versions = mockVersionHelper()) {
+            proxyMock.when(MetaServiceProxy::getInstance).thenReturn(proxy);
+            versions.when(() -> 
VersionHelper.getVersionFromMeta(Mockito.any()))
+                    .thenThrow(new RpcException("MS", "unavailable"));
+            Assertions.assertDoesNotThrow(() -> 
masterTransMgr.commitTransactionWithoutLock(CatalogTestUtil.testDbId1,
+                    List.of(table), response.getTxnInfo().getTxnId(), null, 
null));
+            Assertions.assertEquals(2, partition.getCachedVisibleVersion());
+            versions.verifyNoInteractions();
+            Mockito.verify(callback).afterCommitted(Mockito.argThat(state ->
+                    state.getTransactionStatus() == 
TransactionStatus.VISIBLE), Mockito.eq(true));
+            Mockito.verify(callback).afterVisible(Mockito.argThat(state ->
+                    state.getTransactionStatus() == 
TransactionStatus.VISIBLE), Mockito.eq(true));
+            // A version outage fails the following read, not the already 
committed transaction.
+            Assertions.assertThrows(RuntimeException.class, 
partition::getVisibleVersion);
+            Assertions.assertThrows(RpcException.class, ((OlapTable) 
table)::getVisibleVersion);
+            versions.when(() -> 
VersionHelper.getVersionFromMeta(Mockito.any())).thenReturn(partitionVersion(4));
+            Assertions.assertEquals(4, partition.getVisibleVersion());
+            Assertions.assertEquals(4, ((OlapTable) 
table).getVisibleVersion());
+            Mockito.verifyNoMoreInteractions(callback);
+        }
+    }
+
+    @Test
+    public void testSinglePartitionReadCannotClearConcurrentInvalidation() 
throws Exception {
+        checkReadCannotClearConcurrentInvalidation(false, false);
+    }
+
+    @Test
+    public void testBatchPartitionReadCannotClearConcurrentInvalidation() 
throws Exception {
+        checkReadCannotClearConcurrentInvalidation(false, true);
+    }
+
+    @Test
+    public void testSingleTableReadCannotClearConcurrentInvalidation() throws 
Exception {
+        checkReadCannotClearConcurrentInvalidation(true, false);
+    }
+
+    @Test
+    public void testBatchTableReadCannotClearConcurrentInvalidation() throws 
Exception {
+        checkReadCannotClearConcurrentInvalidation(true, true);
+    }
+
+    private void checkReadCannotClearConcurrentInvalidation(boolean 
tableVersion, boolean batch) throws Exception {
+        useVersionCaches();
+        CloudPartition partition = addCloudPartition(1000);
+        OlapTable table = getCloudTable(partition);
+        
ConnectContext.get().getSessionVariable().cloudPartitionVersionCacheTtlMs = 0;
+        ConnectContext.get().getSessionVariable().cloudTableVersionCacheTtlMs 
= 0;
+        try (MockedStatic<VersionHelper> versions = mockVersionHelper()) {
+            versions.when(() -> 
VersionHelper.getVersionFromMeta(Mockito.any())).thenAnswer(invocation -> {
+                // The MS snapshot predates the commit, but its reply arrives 
after cache invalidation.
+                
masterTransMgr.afterCommitTxnResp(visibleRetry(List.of(table.getId())), null, 
List.of());
+                useVersionCaches();
+                return partitionVersion(2);
+            });
+            if (tableVersion) {
+                Assertions.assertEquals(2L, batch
+                        ? 
OlapTable.getVisibleVersionInBatch(List.of(table)).get(0) : 
table.getVisibleVersion());
+            } else {
+                Assertions.assertEquals(2L, batch
+                        ? 
CloudPartition.getSnapshotVisibleVersion(List.of(partition)).get(0)
+                        : partition.getVisibleVersion());
+            }
+            // A delayed ordinary commit notification cannot repair an unknown 
version either.
+            table.setCachedTableVersion(3);
+            partition.setCachedVisibleVersion(3, 30);
+            versions.when(() -> 
VersionHelper.getVersionFromMeta(Mockito.any())).thenReturn(partitionVersion(4));
+            if (tableVersion) {
+                Assertions.assertEquals(4, table.getVisibleVersion());
+                Assertions.assertEquals(4, table.getVisibleVersion());
+            } else {
+                Assertions.assertEquals(4, partition.getVisibleVersion());
+                Assertions.assertEquals(4, partition.getVisibleVersion());
+            }
+            versions.verify(() -> 
VersionHelper.getVersionFromMeta(Mockito.any()), Mockito.times(2));
+        }
+    }
+
+    @Test
+    public void testInvalidationWaitsForVersionSnapshotReaders() throws 
Exception {
+        CloudPartition partition = addCloudPartition(1000);
+        OlapTable table = getCloudTable(partition);
+        ExecutorService executor = Executors.newSingleThreadExecutor();
+        CountDownLatch started = new CountDownLatch(1);
+        table.versionReadLock();
+        CompletableFuture<Void> invalidation;
+        try {
+            invalidation = CompletableFuture.runAsync(() -> {
+                try (FakeEnv ignored = new FakeEnv()) {
+                    started.countDown();
+                    ((CloudEnv) masterEnv).getCloudFEVersionSynchronizer()
+                            
.invalidateVersionCaches(CatalogTestUtil.testDbId1, List.of(table.getId()));
+                }
+            }, executor);
+            Assertions.assertTrue(started.await(5, TimeUnit.SECONDS));
+            Assertions.assertThrows(TimeoutException.class, () -> 
invalidation.get(200, TimeUnit.MILLISECONDS));
+        } finally {
+            table.versionReadUnlock();
+            executor.shutdown();
+        }
+        invalidation.get(5, TimeUnit.SECONDS);
+        Assertions.assertTrue(executor.awaitTermination(5, TimeUnit.SECONDS));
+    }
+
+    @Test
+    public void testBatchRefreshInstallsVersionsUnderWriteLock() throws 
Exception {
+        CloudPartition first = addCloudPartition(1000);
+        OlapTable table = getCloudTable(first);
+        CloudPartition second = Mockito.spy(new CloudPartition(1002, "p2", new 
MaterializedIndex(),
+                new RandomDistributionInfo(1), CatalogTestUtil.testDbId1, 
table.getId()));
+        table.addPartition(second);
+        ExecutorService executor = Executors.newSingleThreadExecutor();
+        AtomicReference<CompletableFuture<Void>> snapshotReader = new 
AtomicReference<>();
+        try (MockedStatic<VersionHelper> versions = mockVersionHelper()) {
+            versions.when(() -> 
VersionHelper.getVersionFromMeta(Mockito.any())).thenReturn(
+                    
partitionVersion(4).toBuilder().addVersions(4).addVersionUpdateTimeMs(40).addCommitTsos(400).build());
+            Mockito.doAnswer(invocation -> {
+                CountDownLatch started = new CountDownLatch(1);
+                CompletableFuture<Void> reader = CompletableFuture.runAsync(() 
-> {
+                    started.countDown();
+                    table.versionReadLock();
+                    try {
+                        Assertions.assertEquals(4, 
first.getCachedVisibleVersion());
+                        Assertions.assertEquals(4, 
second.getCachedVisibleVersion());
+                    } finally {
+                        table.versionReadUnlock();
+                    }
+                }, executor);
+                snapshotReader.set(reader);
+                Assertions.assertTrue(started.await(5, TimeUnit.SECONDS));
+                Assertions.assertThrows(TimeoutException.class, () -> 
reader.get(200, TimeUnit.MILLISECONDS));
+                return invocation.callRealMethod();
+            }).when(second).setCachedVisibleVersion(4, 40, 400);
+            CloudPartition.getSnapshotVisibleVersionFromMs(List.of(first, 
second), false);
+            snapshotReader.get().get(5, TimeUnit.SECONDS);
+        } finally {
+            executor.shutdown();
+            Assertions.assertTrue(executor.awaitTermination(5, 
TimeUnit.SECONDS));
+        }
+    }
+
+    @Test
+    public void testFollowerInvalidatesTableAndAllPartitions() throws 
Exception {
+        useVersionCaches();
+        CloudPartition partition = addCloudPartition(1000);
+        OlapTable table = getCloudTable(partition);
+        CloudPartition second = new CloudPartition(1002, "p2", new 
MaterializedIndex(),
+                new RandomDistributionInfo(1), CatalogTestUtil.testDbId1, 
table.getId());
+        second.setCachedVisibleVersion(2, 20);
+        table.addPartition(second);
+        TFrontendSyncCloudVersionRequest request = new 
TFrontendSyncCloudVersionRequest()
+                
.setDbId(CatalogTestUtil.testDbId1).setPartitionVersionInfos(List.of())
+                .setTableVersionInfos(List.of(new 
TCloudVersionInfo().setTableId(table.getId()).setVersion(-1)));
+        Mockito.doReturn(false).when(masterEnv).isMaster();
+        FrontendServiceImpl service = new FrontendServiceImpl(null);
+        try (MockedStatic<VersionHelper> versions = mockVersionHelper()) {
+            Assertions.assertEquals(TStatusCode.OK, 
service.syncCloudVersion(request).getStatusCode());
+            versions.verifyNoInteractions();
+            // Even a delayed regular version push must leave these caches 
invalid.
+            service.syncCloudVersion(new TFrontendSyncCloudVersionRequest()
+                    
.setDbId(CatalogTestUtil.testDbId1).setPartitionVersionInfos(List.of())
+                    .setTableVersionInfos(List.of(new 
TCloudVersionInfo().setTableId(table.getId()).setVersion(3))));
+            versions.when(() -> 
VersionHelper.getVersionFromMeta(Mockito.any())).thenReturn(partitionVersion(4));
+            Assertions.assertEquals(4, table.getVisibleVersion());
+            Assertions.assertEquals(4, partition.getVisibleVersion());
+            Assertions.assertEquals(4, second.getVisibleVersion());
+            versions.verify(() -> 
VersionHelper.getVersionFromMeta(Mockito.any()), Mockito.times(3));
+        }
+    }
+
+    @Test
+    @SuppressWarnings("unchecked")
+    public void testInvalidationIsPushedThroughVersionRpc() throws Exception {
+        CloudPartition partition = addCloudPartition(1000);
+        boolean syncEnabled = Config.cloud_enable_version_syncer;
+        GenericPool<FrontendService.Client> originalPool = 
ClientPool.frontendVersionPool;
+        GenericPool<FrontendService.Client> pool = 
Mockito.mock(GenericPool.class);
+        FrontendService.Client client = 
Mockito.mock(FrontendService.Client.class);
+        Mockito.when(pool.borrowObject(Mockito.any())).thenReturn(client);
+        CompletableFuture<TFrontendSyncCloudVersionRequest> sent = new 
CompletableFuture<>();
+        CountDownLatch returned = new CountDownLatch(1);
+        
Mockito.when(client.syncCloudVersion(Mockito.any())).thenAnswer(invocation -> {
+            sent.complete(invocation.getArgument(0));
+            return new TStatus(TStatusCode.OK);
+        });
+        Mockito.doAnswer(invocation -> {
+            returned.countDown();
+            return null;
+        }).when(pool).returnObject(Mockito.any(), Mockito.eq(client));
+        CloudEnv sender = (CloudEnv) masterEnv;
+        Frontend follower = Mockito.mock(Frontend.class);
+        Mockito.when(follower.isAlive()).thenReturn(true);
+        Mockito.when(follower.getHost()).thenReturn("127.0.0.2");
+        Mockito.when(follower.getRpcPort()).thenReturn(9020);
+        Mockito.doReturn(List.of(follower)).when(sender).getFrontends(null);
+        Mockito.doReturn(new HostInfo("127.0.0.1", 
9010)).when(sender).getSelfNode();
+        try {
+            FakeEnv.setEnv(sender);
+            Config.cloud_enable_version_syncer = true;
+            ClientPool.frontendVersionPool = pool;
+            
masterTransMgr.afterCommitTxnResp(visibleRetry(List.of(partition.getTableId())),
 null, List.of());
+            TFrontendSyncCloudVersionRequest request = sent.get(5, 
TimeUnit.SECONDS);
+            Assertions.assertTrue(returned.await(5, TimeUnit.SECONDS));
+            // Check the actual Thrift payload as well as the sender's Java 
objects.
+            TFrontendSyncCloudVersionRequest decoded = new 
TFrontendSyncCloudVersionRequest();
+            new org.apache.thrift.TDeserializer().deserialize(decoded,
+                    new org.apache.thrift.TSerializer().serialize(request));
+            Assertions.assertEquals(CatalogTestUtil.testDbId1, 
decoded.getDbId());
+            
Assertions.assertTrue(decoded.getPartitionVersionInfos().isEmpty());
+            Assertions.assertEquals(1, decoded.getTableVersionInfosSize());
+            Assertions.assertEquals(partition.getTableId(), 
decoded.getTableVersionInfos().get(0).getTableId());
+            Assertions.assertEquals(-1, 
decoded.getTableVersionInfos().get(0).getVersion());
+        } finally {
+            FakeEnv.setEnv(masterEnv);
+            ClientPool.frontendVersionPool = originalPool;
+            Config.cloud_enable_version_syncer = syncEnabled;
+        }
+    }
+
+    private CloudPartition addCloudPartition(long tableId) throws Exception {
+        OlapTable table = new OlapTable(tableId, "version_cache_" + tableId, 
List.of(), KeysType.DUP_KEYS,
+                new PartitionInfo(), new RandomDistributionInfo(1));
+        CloudPartition partition = new CloudPartition(tableId + 1, "p1", new 
MaterializedIndex(),
+                new RandomDistributionInfo(1), CatalogTestUtil.testDbId1, 
tableId);
+        partition.setCachedVisibleVersion(2, 20, 200);
+        table.addPartition(partition);
+        table.setCachedTableVersion(2);
+        
masterEnv.getInternalCatalog().getDbOrMetaException(CatalogTestUtil.testDbId1).registerTable(table);
+        return partition;
+    }
+
+    private MockedStatic<VersionHelper> mockVersionHelper() {
+        MockedStatic<VersionHelper> versions = 
Mockito.mockStatic(VersionHelper.class);
+        // Batch reads pass an explicit retry limit; share the response stub 
with single-version reads.
+        versions.when(() -> VersionHelper.getVersionFromMeta(Mockito.any(), 
Mockito.anyInt()))
+                .thenAnswer(invocation -> 
VersionHelper.getVersionFromMeta(invocation.getArgument(0)));
+        return versions;
+    }
+
+    private CommitTxnResponse visibleRetry(List<Long> tableIds) {
+        return CommitTxnResponse.newBuilder()
+                
.setStatus(Cloud.MetaServiceResponseStatus.newBuilder().setCode(MetaServiceCode.OK))
+                
.setTxnInfo(buildTxnInfo(100).toBuilder().clearTableIds().addAllTableIds(tableIds)
+                        
.setStatus(Cloud.TxnStatusPB.TXN_STATUS_VISIBLE).setCommitTso(300)).build();
+    }
+
+    private Cloud.GetVersionResponse partitionVersion(long version) {
+        return 
Cloud.GetVersionResponse.newBuilder().setVersion(version).addVersions(version)
+                .addVersionUpdateTimeMs(version * 10).addCommitTsos(version * 
100).build();
+    }
+
+    private OlapTable getCloudTable(CloudPartition partition) throws Exception 
{
+        return (OlapTable) 
masterEnv.getInternalCatalog().getDbOrMetaException(CatalogTestUtil.testDbId1)
+                .getTableOrMetaException(partition.getTableId());
+    }
+
+    private void useVersionCaches() {
+        ConnectContext context = new ConnectContext();
+        context.setSessionVariable(new SessionVariable());
+        context.getSessionVariable().cloudPartitionVersionCacheTtlMs = 
Long.MAX_VALUE;
+        context.getSessionVariable().cloudTableVersionCacheTtlMs = 
Long.MAX_VALUE;
+        context.setThreadLocalInfo();
+    }
+
     private TxnInfoPB buildTxnInfo(long transactionId) {
         return TxnInfoPB.newBuilder()
                 .setDbId(CatalogTestUtil.testDbId1)


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to