repushChainVersions = new HashSet<>(); // all versions repushed into the current version
readyToBeRemovedVersions = versions.stream()
- .filter(v -> !VersionStatus.isVersionRolledBack(v.getStatus())) // rolled-back handled separately
+ .filter(StoreBackupVersionCleanupService::isEligibleForStandardBackupCleanup)
.sorted((v1, v2) -> Integer.compare(v2.getNumber(), v1.getNumber())) // sort in descending order
.filter(v -> {
// always delete past default retention and less than current version
@@ -355,7 +355,7 @@ protected boolean cleanupBackupVersion(Store store, String clusterName) {
if (isCurrentVersionRepushed && readyToBeRemovedVersions.isEmpty()) {
for (Version v: versions) {
if (v.getNumber() < currentVersion && v.getRepushSourceVersion() > NON_EXISTING_VERSION
- && !VersionStatus.isVersionRolledBack(v.getStatus())) {
+ && isEligibleForStandardBackupCleanup(v)) {
readyToBeRemovedVersions.add(v);
}
}
@@ -409,15 +409,18 @@ protected boolean cleanupBackupVersion(Store store, String clusterName) {
return true;
}
+ private static boolean isEligibleForStandardBackupCleanup(Version version) {
+ VersionStatus status = version.getStatus();
+ return status != VersionStatus.PUSHED && !VersionStatus.isVersionRolledBack(status);
+ }
+
/**
* Deletes ROLLED_BACK versions whose retention period has expired.
*
- * The retention clock is anchored to {@link Store#getLatestVersionPromoteToCurrentTimestamp()},
- * which is set when a version becomes current — including the rollback-target version being
- * re-promoted. If a subsequent push completes before the retention expires, the timestamp resets
- * and the ROLLED_BACK version survives longer than the configured retention. This is intentionally
- * conservative: a per-version {@code rolledBackTimestamp} would give exact retention but requires
- * a Version schema change.
+ *
The retention clock is anchored to {@link Store#getLatestVersionPromoteToCurrentTimestamp()}.
+ * Marking a child version ROLLED_BACK resets this store-level timestamp, and a later promotion may
+ * reset it again. This is intentionally conservative: exact independent retention for multiple
+ * rolled-back versions would require a per-version rollback timestamp.
*/
boolean cleanupRolledBackVersions(Store store, String clusterName, List versions) {
long rolledBackVersionRetentionMs = getRolledBackVersionRetentionMs(clusterName);
diff --git a/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestDeferredVersionSwapServiceWithSequentialRollout.java b/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestDeferredVersionSwapServiceWithSequentialRollout.java
index 87bd8720453..dad9c24d45a 100644
--- a/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestDeferredVersionSwapServiceWithSequentialRollout.java
+++ b/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestDeferredVersionSwapServiceWithSequentialRollout.java
@@ -5,9 +5,11 @@
import static org.mockito.ArgumentMatchers.anyInt;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.Mockito.anyDouble;
+import static org.mockito.Mockito.atLeast;
import static org.mockito.Mockito.atLeastOnce;
import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.inOrder;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.times;
@@ -39,6 +41,7 @@
import java.time.ZoneOffset;
import java.util.ArrayList;
import java.util.Arrays;
+import java.util.Collections;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
@@ -48,8 +51,10 @@
import java.util.Set;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.TimeUnit;
+import org.mockito.InOrder;
import org.testng.Assert;
import org.testng.annotations.BeforeMethod;
+import org.testng.annotations.DataProvider;
import org.testng.annotations.Test;
@@ -438,6 +443,9 @@ public void testSequentialRolloutFailurePath() throws Exception {
TestUtils.waitForNonDeterministicAssertion(5, TimeUnit.SECONDS, () -> {
// Verify error recording was called due to the failure
verify(store, atLeastOnce()).updateVersionStatus(2, VersionStatus.PARTIALLY_ONLINE);
+ Map childControllers = veniceHelixAdmin.getControllerClientMap(clusterName);
+ verify(childControllers.get(region2), atLeastOnce()).deleteOldVersion(storeName, versionTwo);
+ verify(childControllers.get(region3), atLeastOnce()).deleteOldVersion(storeName, versionTwo);
verify(admin, never()).rollForwardToFutureVersion(clusterName, storeName, region3);
verify(admin, never()).truncateKafkaTopic(anyString());
});
@@ -513,6 +521,295 @@ public void testSequentialRolloutVersionValidationFails() throws Exception {
verify(stats, never()).recordDeferredVersionSwapExceptionMetric(anyString());
}
+ /**
+ * When post-version-swap validation returns ROLLBACK, the target region(s) are rolled back and
+ * bootstrap-complete non-current child copies are reconciled. A child that still reports the failed
+ * version as current remains protected and keeps the reconciliation retriable.
+ */
+ @Test
+ public void testSequentialRolloutPostSwapValidationRollbackMarksNonTargetRegions() throws Exception {
+ String storeName = "testStore";
+ Store store = mockStore(versionOne, versionTwo, storeName);
+
+ // Lifecycle hook that returns ROLLBACK during post-version-swap validation.
+ List lifecycleHooks = new ArrayList<>();
+ Map params = new HashMap<>();
+ params.put("outcome", StoreVersionLifecycleEventOutcome.ROLLBACK.toString());
+ lifecycleHooks.add(new LifecycleHooksRecordImpl(MockStoreLifecycleHooks.class.getName(), params));
+ doReturn(lifecycleHooks).when(store).getStoreLifecycleHooks();
+
+ // Wire the hooks cache to instantiate the mock hook.
+ StoreLifecycleHooksCache hooksCache = mock(StoreLifecycleHooksCache.class);
+ doReturn(new MockStoreLifecycleHooks(new VeniceProperties(new Properties()))).when(hooksCache)
+ .getOrInstantiateHook(MockStoreLifecycleHooks.class.getName());
+ doReturn(hooksCache).when(veniceHelixAdmin).getStoreLifecycleHooksCache();
+
+ List storeList = new ArrayList<>();
+ storeList.add(store);
+ doReturn(storeList).when(admin).getAllStores(clusterName);
+ doReturn(true).when(admin).isLeaderControllerFor(clusterName);
+
+ Version versionOneImpl = new VersionImpl(storeName, versionOne);
+ Version versionTwoImpl = new VersionImpl(storeName, versionTwo);
+ versionTwoImpl.setStatus(VersionStatus.PUSHED);
+ List versionList = new ArrayList<>();
+ versionList.add(versionOneImpl);
+ versionList.add(versionTwoImpl);
+ StoreResponse storeResponse = getStoreResponse(versionList);
+
+ Map controllerClientMap = mockControllerClients(versionList);
+ for (Map.Entry entry: controllerClientMap.entrySet()) {
+ doReturn(storeResponse).when(entry.getValue()).getStore(anyString(), anyInt());
+ }
+
+ doReturn(store).when(repository).getStore(storeName);
+
+ Long time = LocalDateTime.now().toEpochSecond(ZoneOffset.UTC);
+ Admin.OfflinePushStatusInfo completedPush = getOfflinePushStatusInfo(
+ ExecutionStatus.COMPLETED.toString(),
+ ExecutionStatus.COMPLETED.toString(),
+ ExecutionStatus.COMPLETED.toString(),
+ time - TimeUnit.MINUTES.toSeconds(90),
+ time - TimeUnit.MINUTES.toSeconds(30),
+ time - TimeUnit.MINUTES.toSeconds(30));
+ String kafkaTopicName = Version.composeKafkaTopic(storeName, versionTwo);
+ doReturn(completedPush).when(admin).getOffLinePushStatus(clusterName, kafkaTopicName);
+
+ DeferredVersionSwapService deferredVersionSwapService =
+ new DeferredVersionSwapService(admin, veniceControllerMultiClusterConfig, stats, metricsRepository);
+ deferredVersionSwapService.startInner();
+
+ TestUtils.waitForNonDeterministicAssertion(5, TimeUnit.SECONDS, () -> {
+ // Target region (region1, the prior rolled-forward region) is rolled back.
+ verify(admin, atLeastOnce()).rollbackToBackupVersion(clusterName, storeName, region1);
+ // Non-current child regions have their orphaned version deleted.
+ Map childControllers = veniceHelixAdmin.getControllerClientMap(clusterName);
+ verify(childControllers.get(region2), atLeastOnce()).deleteOldVersion(storeName, versionTwo);
+ verify(childControllers.get(region3), atLeastOnce()).deleteOldVersion(storeName, versionTwo);
+ // No roll forward should happen for region2 after a ROLLBACK.
+ verify(admin, never()).rollForwardToFutureVersion(clusterName, storeName, region2);
+ });
+ }
+
+ @DataProvider(name = "abandonedParentStatuses")
+ public Object[][] abandonedParentStatuses() {
+ return new Object[][] { { VersionStatus.ERROR }, { VersionStatus.PARTIALLY_ONLINE },
+ { VersionStatus.ROLLED_BACK } };
+ }
+
+ @Test(dataProvider = "abandonedParentStatuses")
+ public void testAbandonedParentStatusReconcilesEligibleChildren(VersionStatus parentStatus) {
+ String storeName = "testStore";
+ Store store = mockStore(versionOne, versionTwo, storeName);
+ doReturn(store).when(repository).getStore(storeName);
+ doReturn(ConcurrentPushDetectionStrategy.PARENT_VERSION_STATUS_ONLY).when(clusterConfig)
+ .getConcurrentPushDetectionStrategy();
+
+ DeferredVersionSwapService deferredVersionSwapService =
+ new DeferredVersionSwapService(admin, veniceControllerMultiClusterConfig, stats, metricsRepository);
+ deferredVersionSwapService.updateStore(clusterName, storeName, parentStatus, versionTwo);
+
+ Map childControllers = veniceHelixAdmin.getControllerClientMap(clusterName);
+ InOrder inOrder = inOrder(store, childControllers.get(region2));
+ inOrder.verify(store).updateVersionStatus(versionTwo, parentStatus);
+ inOrder.verify(childControllers.get(region2)).deleteOldVersion(storeName, versionTwo);
+ verify(childControllers.get(region3)).deleteOldVersion(storeName, versionTwo);
+ verify(childControllers.get(region1), never()).deleteOldVersion(anyString(), anyInt());
+ }
+
+ @Test
+ public void testAbandonedParentStatusPersistsWhenChildReconciliationFails() {
+ String storeName = "testStore";
+ Store store = mockStore(versionOne, versionTwo, storeName);
+ doReturn(store).when(repository).getStore(storeName);
+
+ Map childControllers = veniceHelixAdmin.getControllerClientMap(clusterName);
+ doThrow(new VeniceException("child reconciliation failed")).when(childControllers.get(region2))
+ .deleteOldVersion(anyString(), anyInt());
+ doThrow(new VeniceException("child reconciliation failed")).when(childControllers.get(region3))
+ .deleteOldVersion(anyString(), anyInt());
+
+ DeferredVersionSwapService deferredVersionSwapService =
+ new DeferredVersionSwapService(admin, veniceControllerMultiClusterConfig, stats, metricsRepository);
+ deferredVersionSwapService.updateStore(clusterName, storeName, VersionStatus.ERROR, versionTwo);
+
+ verify(store).updateVersionStatus(versionTwo, VersionStatus.ERROR);
+ verify(childControllers.get(region2)).deleteOldVersion(storeName, versionTwo);
+ }
+
+ @Test
+ public void testTerminalParentStatusRepairsStrandedChildVersions() throws Exception {
+ String storeName = "testStore";
+ Store store = mockStore(versionOne, versionTwo, storeName);
+ Version targetVersion = store.getVersion(versionTwo);
+ doReturn(VersionStatus.ERROR).when(targetVersion).getStatus();
+ doReturn(VersionStatus.ERROR).when(store).getVersionStatus(versionTwo);
+ doReturn(Arrays.asList(targetVersion)).when(store).getVersions();
+ doReturn(Arrays.asList(store)).when(admin).getAllStores(clusterName);
+ doReturn(true).when(admin).isLeaderControllerFor(clusterName);
+
+ DeferredVersionSwapService deferredVersionSwapService =
+ new DeferredVersionSwapService(admin, veniceControllerMultiClusterConfig, stats, metricsRepository);
+ deferredVersionSwapService.startInner();
+
+ TestUtils.waitForNonDeterministicAssertion(5, TimeUnit.SECONDS, () -> {
+ Map childControllers = veniceHelixAdmin.getControllerClientMap(clusterName);
+ verify(childControllers.get(region2), atLeastOnce()).deleteOldVersion(storeName, versionTwo);
+ verify(childControllers.get(region3), atLeastOnce()).deleteOldVersion(storeName, versionTwo);
+ });
+ }
+
+ @Test
+ public void testTerminalParentStatusRetriesUntilInProgressChildCompletes() throws Exception {
+ String storeName = "testStore";
+ Store store = mockStore(versionOne, versionTwo, storeName);
+ Version targetVersion = store.getVersion(versionTwo);
+ doReturn(VersionStatus.ERROR).when(targetVersion).getStatus();
+ doReturn(VersionStatus.ERROR).when(store).getVersionStatus(versionTwo);
+ doReturn(Arrays.asList(targetVersion)).when(store).getVersions();
+ doReturn(Arrays.asList(store)).when(admin).getAllStores(clusterName);
+ doReturn(true).when(admin).isLeaderControllerFor(clusterName);
+
+ Version startedVersion = new VersionImpl(storeName, versionTwo);
+ startedVersion.setStatus(VersionStatus.STARTED);
+ Version pushedVersion = new VersionImpl(storeName, versionTwo);
+ pushedVersion.setStatus(VersionStatus.PUSHED);
+ Store regionStore = mockRegionalStore(versionOne, versionTwo, storeName);
+ StoreInfo startedStoreInfo = StoreInfo.fromStore(regionStore);
+ startedStoreInfo.setVersions(Arrays.asList(startedVersion));
+ StoreResponse startedResponse = new StoreResponse();
+ startedResponse.setStore(startedStoreInfo);
+ StoreInfo pushedStoreInfo = StoreInfo.fromStore(regionStore);
+ pushedStoreInfo.setVersions(Arrays.asList(pushedVersion));
+ StoreResponse pushedResponse = new StoreResponse();
+ pushedResponse.setStore(pushedStoreInfo);
+
+ Map childControllers = veniceHelixAdmin.getControllerClientMap(clusterName);
+ doReturn(startedResponse, pushedResponse).when(childControllers.get(region2))
+ .getStore(storeName, controllerTimeout);
+
+ DeferredVersionSwapService deferredVersionSwapService =
+ new DeferredVersionSwapService(admin, veniceControllerMultiClusterConfig, stats, metricsRepository);
+ deferredVersionSwapService.startInner();
+
+ TestUtils.waitForNonDeterministicAssertion(5, TimeUnit.SECONDS, () -> {
+ verify(childControllers.get(region2), atLeast(2)).getStore(storeName, controllerTimeout);
+ verify(childControllers.get(region2), atLeastOnce()).deleteOldVersion(storeName, versionTwo);
+ });
+ }
+
+ @Test
+ public void testTerminalParentStatusRetriesWhenChildVersionAppearsLater() throws Exception {
+ String storeName = "testStore";
+ Store store = mockStore(versionOne, versionTwo, storeName);
+ Version targetVersion = store.getVersion(versionTwo);
+ doReturn(VersionStatus.ERROR).when(targetVersion).getStatus();
+ doReturn(VersionStatus.ERROR).when(store).getVersionStatus(versionTwo);
+ doReturn(Arrays.asList(targetVersion)).when(store).getVersions();
+ doReturn(Arrays.asList(store)).when(admin).getAllStores(clusterName);
+ doReturn(true).when(admin).isLeaderControllerFor(clusterName);
+
+ Store regionStore = mockRegionalStore(versionOne, versionTwo, storeName);
+ StoreInfo missingStoreInfo = StoreInfo.fromStore(regionStore);
+ missingStoreInfo.setVersions(Collections.emptyList());
+ StoreResponse missingResponse = new StoreResponse();
+ missingResponse.setStore(missingStoreInfo);
+ Version pushedVersion = new VersionImpl(storeName, versionTwo);
+ pushedVersion.setStatus(VersionStatus.PUSHED);
+ StoreInfo pushedStoreInfo = StoreInfo.fromStore(regionStore);
+ pushedStoreInfo.setVersions(Arrays.asList(pushedVersion));
+ StoreResponse pushedResponse = new StoreResponse();
+ pushedResponse.setStore(pushedStoreInfo);
+
+ Map childControllers = veniceHelixAdmin.getControllerClientMap(clusterName);
+ doReturn(missingResponse, pushedResponse).when(childControllers.get(region2))
+ .getStore(storeName, controllerTimeout);
+
+ DeferredVersionSwapService deferredVersionSwapService =
+ new DeferredVersionSwapService(admin, veniceControllerMultiClusterConfig, stats, metricsRepository);
+ deferredVersionSwapService.startInner();
+
+ TestUtils.waitForNonDeterministicAssertion(5, TimeUnit.SECONDS, () -> {
+ verify(childControllers.get(region2), atLeast(2)).getStore(storeName, controllerTimeout);
+ verify(childControllers.get(region2), atLeastOnce()).deleteOldVersion(storeName, versionTwo);
+ });
+ }
+
+ @Test
+ public void testTerminalParentStatusRetriesCurrentTargetAfterItBecomesNonCurrent() throws Exception {
+ String storeName = "testStore";
+ Store store = mockStore(versionOne, versionTwo, storeName);
+ Version targetVersion = store.getVersion(versionTwo);
+ doReturn(VersionStatus.ERROR).when(targetVersion).getStatus();
+ doReturn(VersionStatus.ERROR).when(store).getVersionStatus(versionTwo);
+ doReturn(Arrays.asList(targetVersion)).when(store).getVersions();
+ doReturn(Arrays.asList(store)).when(admin).getAllStores(clusterName);
+ doReturn(true).when(admin).isLeaderControllerFor(clusterName);
+
+ Version pushedVersion = new VersionImpl(storeName, versionTwo);
+ pushedVersion.setStatus(VersionStatus.PUSHED);
+ StoreInfo currentTargetStoreInfo = StoreInfo.fromStore(mockRegionalStore(versionTwo, versionTwo, storeName));
+ currentTargetStoreInfo.setVersions(Arrays.asList(pushedVersion));
+ StoreResponse currentTargetResponse = new StoreResponse();
+ currentTargetResponse.setStore(currentTargetStoreInfo);
+ StoreInfo nonCurrentTargetStoreInfo = StoreInfo.fromStore(mockRegionalStore(versionOne, versionTwo, storeName));
+ nonCurrentTargetStoreInfo.setVersions(Arrays.asList(pushedVersion));
+ StoreResponse nonCurrentTargetResponse = new StoreResponse();
+ nonCurrentTargetResponse.setStore(nonCurrentTargetStoreInfo);
+
+ Map childControllers = veniceHelixAdmin.getControllerClientMap(clusterName);
+ doReturn(currentTargetResponse, nonCurrentTargetResponse).when(childControllers.get(region1))
+ .getStore(storeName, controllerTimeout);
+
+ DeferredVersionSwapService deferredVersionSwapService =
+ new DeferredVersionSwapService(admin, veniceControllerMultiClusterConfig, stats, metricsRepository);
+ deferredVersionSwapService.startInner();
+
+ TestUtils.waitForNonDeterministicAssertion(5, TimeUnit.SECONDS, () -> {
+ verify(childControllers.get(region1), atLeast(2)).getStore(storeName, controllerTimeout);
+ // After region1 becomes non-current, all 3 regions should have deleteOldVersion called.
+ verify(childControllers.get(region1), atLeastOnce()).deleteOldVersion(storeName, versionTwo);
+ verify(childControllers.get(region2), atLeastOnce()).deleteOldVersion(storeName, versionTwo);
+ verify(childControllers.get(region3), atLeastOnce()).deleteOldVersion(storeName, versionTwo);
+ });
+ }
+
+ @Test
+ public void testTerminalReconciliationIncludesSupersededParentVersions() throws Exception {
+ String storeName = "testStore";
+ int latestVersionNumber = 3;
+ Store store = mockStore(versionOne, latestVersionNumber, storeName);
+ Version latestVersion = store.getVersion(latestVersionNumber);
+ Version supersededVersion = mock(Version.class);
+ doReturn(versionTwo).when(supersededVersion).getNumber();
+ doReturn(VersionStatus.ERROR).when(supersededVersion).getStatus();
+ doReturn(true).when(supersededVersion).isVersionSwapDeferred();
+ doReturn(Arrays.asList(supersededVersion, latestVersion)).when(store).getVersions();
+ doReturn(Arrays.asList(store)).when(admin).getAllStores(clusterName);
+ doReturn(true).when(admin).isLeaderControllerFor(clusterName);
+
+ Version pushedVersion = new VersionImpl(storeName, versionTwo);
+ pushedVersion.setStatus(VersionStatus.PUSHED);
+ StoreInfo childStoreInfo = StoreInfo.fromStore(mockRegionalStore(versionOne, versionTwo, storeName));
+ childStoreInfo.setVersions(Arrays.asList(pushedVersion));
+ StoreResponse childStoreResponse = new StoreResponse();
+ childStoreResponse.setStore(childStoreInfo);
+ for (ControllerClient childController: veniceHelixAdmin.getControllerClientMap(clusterName).values()) {
+ doReturn(childStoreResponse).when(childController).getStore(storeName, controllerTimeout);
+ }
+
+ DeferredVersionSwapService deferredVersionSwapService =
+ new DeferredVersionSwapService(admin, veniceControllerMultiClusterConfig, stats, metricsRepository);
+ deferredVersionSwapService.startInner();
+
+ Map childControllers = veniceHelixAdmin.getControllerClientMap(clusterName);
+ TestUtils.waitForNonDeterministicAssertion(5, TimeUnit.SECONDS, () -> {
+ verify(childControllers.get(region1), atLeastOnce()).deleteOldVersion(storeName, versionTwo);
+ verify(childControllers.get(region2), atLeastOnce()).deleteOldVersion(storeName, versionTwo);
+ verify(childControllers.get(region3), atLeastOnce()).deleteOldVersion(storeName, versionTwo);
+ });
+ }
+
/**
* When the last region in rollout order is ONLINE,
* parent version is marked as ONLINE
diff --git a/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestStoreBackupVersionCleanupService.java b/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestStoreBackupVersionCleanupService.java
index 2c6f6365261..33f07cb3eff 100644
--- a/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestStoreBackupVersionCleanupService.java
+++ b/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestStoreBackupVersionCleanupService.java
@@ -822,6 +822,55 @@ public void testRegularPushBackupNotDeletedEarlyWhenRepushChainBroken() {
verify(admin, atLeast(1)).deleteOldVersionInStore(CLUSTER_NAME, storeExpired.getName(), 10);
}
+ @Test
+ public void testPushedVersionIsNotDeletedAfterDefaultRetention() {
+ Map versions = new HashMap<>();
+ versions.put(1, VersionStatus.PUSHED);
+ versions.put(2, VersionStatus.ONLINE);
+ long expiredRetention = mockTime.getMilliseconds() - DEFAULT_RETENTION_MS - 1;
+ Store store = mockStore(-1, expiredRetention, versions, 2);
+ Version currentVersion = store.getVersion(2);
+ doReturn(-1).when(currentVersion).getRepushSourceVersion();
+ doReturn(currentVersion).when(store).getVersionOrThrow(2);
+
+ Assert.assertFalse(service.cleanupBackupVersion(store, CLUSTER_NAME));
+ verify(admin, never()).deleteOldVersionInStore(CLUSTER_NAME, store.getName(), 1);
+ }
+
+ @Test
+ public void testPushedVersionIsNotDeletedAsRepushSource() {
+ Map versions = new HashMap<>();
+ versions.put(2, VersionStatus.PUSHED);
+ versions.put(3, VersionStatus.ONLINE);
+ long pastMinimumRetention = mockTime.getMilliseconds() - TimeUnit.HOURS.toMillis(2);
+ Store store = mockStore(-1, pastMinimumRetention, versions, 3);
+ Version pushedVersion = store.getVersion(2);
+ doReturn(1).when(pushedVersion).getRepushSourceVersion();
+ Version currentVersion = store.getVersion(3);
+ doReturn(2).when(currentVersion).getRepushSourceVersion();
+ doReturn(currentVersion).when(store).getVersionOrThrow(3);
+
+ Assert.assertFalse(service.cleanupBackupVersion(store, CLUSTER_NAME));
+ verify(admin, never()).deleteOldVersionInStore(CLUSTER_NAME, store.getName(), 2);
+ }
+
+ @Test
+ public void testPushedVersionDoesNotMakeOnlineBackupEligibleEarly() {
+ Map versions = new HashMap<>();
+ versions.put(1, VersionStatus.ONLINE);
+ versions.put(2, VersionStatus.PUSHED);
+ versions.put(3, VersionStatus.ONLINE);
+ long pastMinimumRetention = mockTime.getMilliseconds() - TimeUnit.HOURS.toMillis(2);
+ Store store = mockStore(-1, pastMinimumRetention, versions, 3);
+ Version currentVersion = store.getVersion(3);
+ doReturn(-1).when(currentVersion).getRepushSourceVersion();
+ doReturn(currentVersion).when(store).getVersionOrThrow(3);
+
+ Assert.assertFalse(service.cleanupBackupVersion(store, CLUSTER_NAME));
+ verify(admin, never()).deleteOldVersionInStore(CLUSTER_NAME, store.getName(), 1);
+ verify(admin, never()).deleteOldVersionInStore(CLUSTER_NAME, store.getName(), 2);
+ }
+
@org.testng.annotations.DataProvider(name = "rolledBackVersionCleanupParams")
public Object[][] rolledBackVersionCleanupParams() {
// { hoursElapsed, extraVersions (status map beyond v1 ONLINE + v2 ROLLED_BACK), expectCleanup,
diff --git a/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestVeniceHelixAdmin.java b/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestVeniceHelixAdmin.java
index d117d7c5493..594d23503de 100644
--- a/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestVeniceHelixAdmin.java
+++ b/services/venice-controller/src/test/java/com/linkedin/venice/controller/TestVeniceHelixAdmin.java
@@ -1294,7 +1294,7 @@ public void testRollForwardNoFutureVersions() {
/**
* isPartitionReadyToServe=>true: Future version exists and partitions are ready → success
- * isPartitionReadyToServe=>false: Future version exists but partitions aren’t ready → exception
+ * isPartitionReadyToServe=>false: Future version exists but partitions aren't ready → exception
*/
@Test(dataProvider = "True-and-False", dataProviderClass = DataProviderUtils.class)
public void testRollForwardPartitionNotReady(boolean isPartitionReadyToServe) throws Exception {