diff --git a/clients/da-vinci-client/src/main/java/com/linkedin/davinci/stats/AggVersionedStorageEngineStats.java b/clients/da-vinci-client/src/main/java/com/linkedin/davinci/stats/AggVersionedStorageEngineStats.java index 6692fb885ae..d059e887099 100644 --- a/clients/da-vinci-client/src/main/java/com/linkedin/davinci/stats/AggVersionedStorageEngineStats.java +++ b/clients/da-vinci-client/src/main/java/com/linkedin/davinci/stats/AggVersionedStorageEngineStats.java @@ -93,6 +93,24 @@ public void setStorageEngine(String topicName, StorageEngine storageEngine) { } } + public void unsetStorageEngine(String topicName) { + if (!Version.isVersionTopicOrStreamReprocessingTopic(topicName)) { + LOGGER.warn("Invalid topic name: {}", topicName); + return; + } + String storeName = Version.parseStoreFromKafkaTopicName(topicName); + int version = Version.parseVersionFromKafkaTopicName(topicName); + try { + getStats(storeName, version).setStorageEngine(null); + otelStatsMap.computeIfPresent(storeName, (k, stats) -> { + stats.onVersionRemoved(version); + return stats; + }); + } catch (Exception e) { + LOGGER.warn("Failed to unset StorageEngine for store: {}, version: {}", storeName, version, e); + } + } + public void recordRocksDBOpenFailure(String topicName) { if (!Version.isVersionTopicOrStreamReprocessingTopic(topicName)) { LOGGER.warn("Invalid topic name: {}", topicName); diff --git a/clients/da-vinci-client/src/main/java/com/linkedin/davinci/stats/StorageEngineOtelMetricEntity.java b/clients/da-vinci-client/src/main/java/com/linkedin/davinci/stats/StorageEngineOtelMetricEntity.java index eb0b5fd2164..7feac1573ca 100644 --- a/clients/da-vinci-client/src/main/java/com/linkedin/davinci/stats/StorageEngineOtelMetricEntity.java +++ b/clients/da-vinci-client/src/main/java/com/linkedin/davinci/stats/StorageEngineOtelMetricEntity.java @@ -17,10 +17,11 @@ /** * OTel metric entities for storage engine statistics. * - *

Consolidates 4 Tehuti AsyncGauge sensors into 3 OTel metrics: + *

Consolidates storage-engine Tehuti sensors and adds bounded role-level visibility: *

@@ -28,10 +29,16 @@ public enum StorageEngineOtelMetricEntity implements ModuleMetricEntityInterface { DISK_USAGE( "ingestion.disk.used", MetricType.ASYNC_GAUGE, MetricUnit.BYTES, - "Disk usage in bytes by record type (data or replication metadata)", + "Total disk usage in bytes across all versions with each role, by record type", setOf(VENICE_CLUSTER_NAME, VENICE_STORE_NAME, VENICE_VERSION_ROLE, VENICE_RECORD_TYPE) ), + VERSION_COUNT( + "ingestion.disk.version_count", MetricType.ASYNC_GAUGE, MetricUnit.NUMBER, + "Number of locally loaded storage-engine versions by role", + setOf(VENICE_CLUSTER_NAME, VENICE_STORE_NAME, VENICE_VERSION_ROLE) + ), + ROCKSDB_OPEN_FAILURE_COUNT( "rocksdb.open.failure_count", MetricType.COUNTER, MetricUnit.NUMBER, "Count of RocksDB open failures; VERSION_ROLE reflects the version's role at failure time, not its current role", diff --git a/clients/da-vinci-client/src/main/java/com/linkedin/davinci/stats/StorageEngineOtelStats.java b/clients/da-vinci-client/src/main/java/com/linkedin/davinci/stats/StorageEngineOtelStats.java index b5cd40e1da1..2a5d586a2d2 100644 --- a/clients/da-vinci-client/src/main/java/com/linkedin/davinci/stats/StorageEngineOtelStats.java +++ b/clients/da-vinci-client/src/main/java/com/linkedin/davinci/stats/StorageEngineOtelStats.java @@ -5,6 +5,7 @@ import static com.linkedin.davinci.stats.StorageEngineOtelMetricEntity.DISK_USAGE; import static com.linkedin.davinci.stats.StorageEngineOtelMetricEntity.KEY_COUNT_ESTIMATE; import static com.linkedin.davinci.stats.StorageEngineOtelMetricEntity.ROCKSDB_OPEN_FAILURE_COUNT; +import static com.linkedin.davinci.stats.StorageEngineOtelMetricEntity.VERSION_COUNT; import static com.linkedin.venice.meta.Store.NON_EXISTING_VERSION; import com.linkedin.davinci.stats.AggVersionedStorageEngineStats.StorageEngineStatsWrapper; @@ -26,10 +27,11 @@ /** * Per-store OTel stats for storage engine metrics. * - *

Holds 3 OTel metrics: + *

Holds 4 OTel metrics: *

@@ -61,6 +63,9 @@ public class StorageEngineOtelStats implements Closeable { /** Key count ASYNC_GAUGE with VersionRole dimension */ private final AsyncMetricEntityStateOneEnum keyCountMetric; + /** Number of local storage engines with each VersionRole */ + private final AsyncMetricEntityStateOneEnum versionCountMetric; + /** RocksDB open failure COUNTER with VersionRole dimension */ private final MetricEntityStateOneEnum openFailureMetric; @@ -78,9 +83,8 @@ public StorageEngineOtelStats(MetricsRepository metricsRepository, String storeN Map baseDimensionsMap = otelSetup.getBaseDimensionsMap(); /* - * Two-callback contract: the liveStateResolver returns the wrapper or null (null -> dormant, - * no emission); the valueResolver reads the metric value from that wrapper. The null return - * is the liveness signal, enforced by the API. + * The live-state resolver returns the bounded per-version map when a role is present, or null + * when dormant. The value resolver aggregates that role without creating per-version attributes. */ this.diskUsageMetrics = AsyncMetricEntityStateTwoEnums.create( DISK_USAGE.getMetricEntity(), @@ -88,8 +92,8 @@ public StorageEngineOtelStats(MetricsRepository metricsRepository, String storeN baseDimensionsMap, VeniceRecordType.class, VersionRole.class, - (recordType, role) -> getWrapperForRole(role), - (wrapper, recordType, role) -> diskUsage(wrapper, recordType)); + (recordType, role) -> hasVersionForRole(role) ? wrappersByVersion : null, + (wrappers, recordType, role) -> diskUsageForRole(wrappers, role, recordType)); this.keyCountMetric = AsyncMetricEntityStateOneEnum.create( KEY_COUNT_ESTIMATE.getMetricEntity(), @@ -99,12 +103,21 @@ public StorageEngineOtelStats(MetricsRepository metricsRepository, String storeN role -> getWrapperForRole(role), (wrapper, role) -> wrapper.getKeyCountEstimate()); + this.versionCountMetric = AsyncMetricEntityStateOneEnum.create( + VERSION_COUNT.getMetricEntity(), + otelRepository, + baseDimensionsMap, + VersionRole.class, + role -> wrappersByVersion.isEmpty() ? null : wrappersByVersion, + (wrappers, role) -> countVersionsForRole(wrappers, role)); + // RocksDB open failure count: COUNTER with VersionRole dimension this.openFailureMetric = MetricEntityStateOneEnum .create(ROCKSDB_OPEN_FAILURE_COUNT.getMetricEntity(), otelRepository, baseDimensionsMap, VersionRole.class); } else { this.diskUsageMetrics = null; this.keyCountMetric = null; + this.versionCountMetric = null; this.openFailureMetric = null; } } @@ -175,6 +188,33 @@ private StorageEngineStatsWrapper getWrapperForRole(VersionRole role) { return wrappersByVersion.get(version); } + private boolean hasVersionForRole(VersionRole role) { + return countVersionsForRole(wrappersByVersion, role) > 0; + } + + private long countVersionsForRole(Map wrappers, VersionRole role) { + VersionInfo snapshot = versionInfo; + return wrappers.keySet().stream().filter(version -> classifyVersion(version, snapshot) == role).count(); + } + + private long diskUsageForRole( + Map wrappers, + VersionRole role, + VeniceRecordType recordType) { + VersionInfo snapshot = versionInfo; + if (role == VersionRole.BACKUP) { + return wrappers.entrySet() + .stream() + .filter(entry -> classifyVersion(entry.getKey(), snapshot) == VersionRole.BACKUP) + .mapToLong(entry -> diskUsage(entry.getValue(), recordType)) + .sum(); + } + + int version = getVersionForRole(role, snapshot, wrappers.keySet()); + StorageEngineStatsWrapper wrapper = wrappers.get(version); + return wrapper == null ? 0 : diskUsage(wrapper, recordType); + } + /** Reads disk usage (data or RMD) from a resolved wrapper. */ private static long diskUsage(StorageEngineStatsWrapper wrapper, VeniceRecordType recordType) { switch (recordType) { @@ -187,14 +227,15 @@ private static long diskUsage(StorageEngineStatsWrapper wrapper, VeniceRecordTyp } } - /** - * Clears internal wrapper references. On subsequent collections each async-gauge's - * {@code liveStateResolver} will return {@code null} for every role and no data points will be - * emitted. The SDK instruments themselves are NOT deregistered — they remain registered and - * are polled until the SDK is shut down. - */ + /** Stops observable callbacks and clears all storage-engine references. */ @Override public void close() { wrappersByVersion.clear(); + versionInfo = VersionInfo.NON_EXISTING; + if (emitOtelMetrics) { + diskUsageMetrics.close(); + keyCountMetric.close(); + versionCountMetric.close(); + } } } diff --git a/clients/da-vinci-client/src/main/java/com/linkedin/davinci/storage/StorageService.java b/clients/da-vinci-client/src/main/java/com/linkedin/davinci/storage/StorageService.java index a423372ec62..da64baa94ce 100644 --- a/clients/da-vinci-client/src/main/java/com/linkedin/davinci/storage/StorageService.java +++ b/clients/da-vinci-client/src/main/java/com/linkedin/davinci/storage/StorageService.java @@ -526,6 +526,7 @@ public synchronized void removeStorageEngine(String kafkaTopic) { LOGGER.warn("Storage engine {} does not exist, ignoring remove request.", kafkaTopic); return; } + aggVersionedStorageEngineStats.unsetStorageEngine(kafkaTopic); storageEngine.drop(); VeniceStoreVersionConfig storeConfig = configLoader.getStoreConfig(kafkaTopic); @@ -552,6 +553,7 @@ public synchronized void closeStorageEngine(String kafkaTopic) { LOGGER.warn("Storage engine {} does not exist, ignoring close request.", kafkaTopic); return; } + aggVersionedStorageEngineStats.unsetStorageEngine(kafkaTopic); storageEngine.close(); VeniceStoreVersionConfig storeConfig = configLoader.getStoreConfig(kafkaTopic); diff --git a/clients/da-vinci-client/src/test/java/com/linkedin/davinci/stats/AggVersionedStorageEngineStatsTest.java b/clients/da-vinci-client/src/test/java/com/linkedin/davinci/stats/AggVersionedStorageEngineStatsTest.java index e2b79865ea3..c5d7d09569d 100644 --- a/clients/da-vinci-client/src/test/java/com/linkedin/davinci/stats/AggVersionedStorageEngineStatsTest.java +++ b/clients/da-vinci-client/src/test/java/com/linkedin/davinci/stats/AggVersionedStorageEngineStatsTest.java @@ -84,6 +84,31 @@ public void testSetStorageEngineWiresOtelStats() { } } + @Test + public void testUnsetStorageEngineRemovesAndReopenRestoresOtelStats() { + InMemoryMetricReader reader = InMemoryMetricReader.create(); + try (VeniceMetricsRepository veniceRepo = createOtelMetricsRepository(reader)) { + OtelTestContext ctx = createOtelTestContext(veniceRepo); + StorageEngineStats mockEngineStats = mock(StorageEngineStats.class); + doReturn(5000L).when(mockEngineStats).getStoreSizeInBytes(); + StorageEngine mockEngine = mock(StorageEngine.class); + doReturn(mockEngineStats).when(mockEngine).getStats(); + Attributes attrs = buildDiskUsageDataAttrs(ctx.clusterName, ctx.storeName, VersionRole.BACKUP); + String diskMetric = StorageEngineOtelMetricEntity.DISK_USAGE.getMetricEntity().getMetricName(); + + ctx.stats.setStorageEngine(ctx.topicName, mockEngine); + OpenTelemetryDataTestUtils.validateLongPointDataFromGauge(reader, 5000, attrs, diskMetric, OTEL_PREFIX); + + ctx.stats.unsetStorageEngine(ctx.topicName); + LongPointData pointAfterClose = OpenTelemetryDataTestUtils + .getLongPointDataFromGaugeIfPresent(reader.collectAllMetrics(), diskMetric, OTEL_PREFIX, attrs); + Assert.assertNull(pointAfterClose, "Closed local engine must not emit a version metric"); + + ctx.stats.setStorageEngine(ctx.topicName, mockEngine); + OpenTelemetryDataTestUtils.validateLongPointDataFromGauge(reader, 5000, attrs, diskMetric, OTEL_PREFIX); + } + } + @Test public void testRecordRocksDBOpenFailureWiresBothTehutiAndOtel() { InMemoryMetricReader reader = InMemoryMetricReader.create(); diff --git a/clients/da-vinci-client/src/test/java/com/linkedin/davinci/stats/ServerMetricEntityTest.java b/clients/da-vinci-client/src/test/java/com/linkedin/davinci/stats/ServerMetricEntityTest.java index 8fee846e59d..2c1ae4978be 100644 --- a/clients/da-vinci-client/src/test/java/com/linkedin/davinci/stats/ServerMetricEntityTest.java +++ b/clients/da-vinci-client/src/test/java/com/linkedin/davinci/stats/ServerMetricEntityTest.java @@ -22,7 +22,7 @@ public class ServerMetricEntityTest { @Test public void testServerMetricEntitiesCount() { - assertEquals(SERVER_METRIC_ENTITIES.size(), 185, "Expected 185 unique metric entities"); + assertEquals(SERVER_METRIC_ENTITIES.size(), 186, "Expected 186 unique metric entities"); } /** diff --git a/clients/da-vinci-client/src/test/java/com/linkedin/davinci/stats/StorageEngineOtelMetricEntityTest.java b/clients/da-vinci-client/src/test/java/com/linkedin/davinci/stats/StorageEngineOtelMetricEntityTest.java index d3a819f5428..e281176b649 100644 --- a/clients/da-vinci-client/src/test/java/com/linkedin/davinci/stats/StorageEngineOtelMetricEntityTest.java +++ b/clients/da-vinci-client/src/test/java/com/linkedin/davinci/stats/StorageEngineOtelMetricEntityTest.java @@ -29,8 +29,16 @@ private static Map expec "ingestion.disk.used", MetricType.ASYNC_GAUGE, MetricUnit.BYTES, - "Disk usage in bytes by record type (data or replication metadata)", + "Total disk usage in bytes across all versions with each role, by record type", setOf(VENICE_CLUSTER_NAME, VENICE_STORE_NAME, VENICE_VERSION_ROLE, VENICE_RECORD_TYPE))); + map.put( + StorageEngineOtelMetricEntity.VERSION_COUNT, + new MetricEntityExpectation( + "ingestion.disk.version_count", + MetricType.ASYNC_GAUGE, + MetricUnit.NUMBER, + "Number of locally loaded storage-engine versions by role", + setOf(VENICE_CLUSTER_NAME, VENICE_STORE_NAME, VENICE_VERSION_ROLE))); map.put( StorageEngineOtelMetricEntity.ROCKSDB_OPEN_FAILURE_COUNT, new MetricEntityExpectation( diff --git a/clients/da-vinci-client/src/test/java/com/linkedin/davinci/stats/StorageEngineOtelStatsTest.java b/clients/da-vinci-client/src/test/java/com/linkedin/davinci/stats/StorageEngineOtelStatsTest.java index aa35a118fde..5ed869b80bb 100644 --- a/clients/da-vinci-client/src/test/java/com/linkedin/davinci/stats/StorageEngineOtelStatsTest.java +++ b/clients/da-vinci-client/src/test/java/com/linkedin/davinci/stats/StorageEngineOtelStatsTest.java @@ -32,6 +32,8 @@ public class StorageEngineOtelStatsTest { private static final String DISK_USAGE_METRIC = StorageEngineOtelMetricEntity.DISK_USAGE.getMetricEntity().getMetricName(); + private static final String VERSION_COUNT_METRIC = + StorageEngineOtelMetricEntity.VERSION_COUNT.getMetricEntity().getMetricName(); private static final String KEY_COUNT_METRIC = StorageEngineOtelMetricEntity.KEY_COUNT_ESTIMATE.getMetricEntity().getMetricName(); private static final String OPEN_FAILURE_METRIC = @@ -144,17 +146,23 @@ public void testDiskUsageRmdBackupVersion() { } @Test - public void testDiskUsageBackupSelectsSmallestVersion() { - // current=1, future=2. Versions 3 and 5 are both backups — should select version 3 (smallest) + public void testDiskUsageAggregatesAllBackupVersions() { + // current=1, future=2. Versions 3 and 5 are both backups. stats.setStatsWrapper(3, new MockWrapper(3000, 300, 30)); stats.setStatsWrapper(5, new MockWrapper(5000, 500, 50)); OpenTelemetryDataTestUtils.validateLongPointDataFromGauge( inMemoryMetricReader, - 3000, + 8000, buildDiskUsageAttributes(VersionRole.BACKUP, VeniceRecordType.DATA), DISK_USAGE_METRIC, METRIC_PREFIX); + OpenTelemetryDataTestUtils.validateLongPointDataFromGauge( + inMemoryMetricReader, + 800, + buildDiskUsageAttributes(VersionRole.BACKUP, VeniceRecordType.REPLICATION_METADATA), + DISK_USAGE_METRIC, + METRIC_PREFIX); } @Test @@ -326,6 +334,29 @@ public void testLiveValueUpdatesAfterVersionInfoChange() { METRIC_PREFIX); } + @Test + public void testVersionCountTracksAddRemoveAndRoleChanges() { + stats.setStatsWrapper(1, new MockWrapper(1000, 100, 10)); + stats.setStatsWrapper(2, new MockWrapper(2000, 200, 20)); + stats.setStatsWrapper(3, new MockWrapper(3000, 300, 30)); + stats.setStatsWrapper(5, new MockWrapper(5000, 500, 50)); + + assertVersionCount(VersionRole.CURRENT, 1); + assertVersionCount(VersionRole.FUTURE, 1); + assertVersionCount(VersionRole.BACKUP, 2); + + stats.onVersionRemoved(3); + assertVersionCount(VersionRole.BACKUP, 1); + + stats.updateVersionInfo(5, 2); + assertVersionCount(VersionRole.CURRENT, 1); + assertVersionCount(VersionRole.FUTURE, 1); + assertVersionCount(VersionRole.BACKUP, 1); + + stats.onVersionRemoved(1); + assertVersionCount(VersionRole.BACKUP, 0); + } + @Test public void testRemoveVersionClearsWrapper() { AggVersionedStorageEngineStats.StorageEngineStatsWrapper wrapper = new MockWrapper(1000, 100, 10); @@ -394,6 +425,15 @@ private static Attributes buildVersionRoleAttributes(VersionRole role) { .build(); } + private void assertVersionCount(VersionRole role, long expectedCount) { + OpenTelemetryDataTestUtils.validateLongPointDataFromGauge( + inMemoryMetricReader, + expectedCount, + buildVersionRoleAttributes(role), + VERSION_COUNT_METRIC, + METRIC_PREFIX); + } + @Test public void testVersionRoleEnumCount() { // getVersionForRole returns NON_EXISTING_VERSION for unknown roles — but a new VersionRole diff --git a/clients/da-vinci-client/src/test/java/com/linkedin/davinci/storage/StorageServiceTest.java b/clients/da-vinci-client/src/test/java/com/linkedin/davinci/storage/StorageServiceTest.java index 8832dd3015b..8234d9422e9 100644 --- a/clients/da-vinci-client/src/test/java/com/linkedin/davinci/storage/StorageServiceTest.java +++ b/clients/da-vinci-client/src/test/java/com/linkedin/davinci/storage/StorageServiceTest.java @@ -16,6 +16,7 @@ import com.linkedin.davinci.config.VeniceStoreVersionConfig; import com.linkedin.davinci.stats.AggVersionedStorageEngineStats; import com.linkedin.davinci.stats.RocksDBMemoryStats; +import com.linkedin.davinci.store.DelegatingStorageEngine; import com.linkedin.davinci.store.StorageEngine; import com.linkedin.davinci.store.StorageEngineFactory; import com.linkedin.venice.exceptions.VeniceNoStoreException; @@ -76,6 +77,47 @@ public void testDeleteStorageEngineOnRocksDBError() { verify(factory, times(2)).removeStorageEngine(storageEngineName); } + @Test + public void testRemoveAndCloseStorageEngineDetachStats() { + VeniceConfigLoader configLoader = mock(VeniceConfigLoader.class); + VeniceServerConfig serverConfig = mock(VeniceServerConfig.class); + when(serverConfig.getDataBasePath()).thenReturn("/tmp"); + when(configLoader.getVeniceServerConfig()).thenReturn(serverConfig); + VeniceStoreVersionConfig storeConfig = mock(VeniceStoreVersionConfig.class); + when(storeConfig.getStorePersistenceType()).thenReturn(PersistenceType.BLACK_HOLE); + when(configLoader.getStoreConfig(storageEngineName)).thenReturn(storeConfig); + + AggVersionedStorageEngineStats storageEngineStats = mock(AggVersionedStorageEngineStats.class); + StorageEngineFactory factory = mock(StorageEngineFactory.class); + Map factories = new HashMap<>(); + factories.put(PersistenceType.BLACK_HOLE, factory); + StorageService storageService = new StorageService( + configLoader, + storageEngineStats, + mock(RocksDBMemoryStats.class), + mock(InternalAvroSpecificSerializer.class), + mock(InternalAvroSpecificSerializer.class), + storeRepository, + false, + false, + ignored -> true, + Optional.of(factories)); + + DelegatingStorageEngine storageEngine = mock(DelegatingStorageEngine.class); + when(storageEngine.getStoreVersionName()).thenReturn(storageEngineName); + when(storageEngine.getType()).thenReturn(PersistenceType.BLACK_HOLE); + storageService.getStorageEngineRepository().addLocalStorageEngine(storageEngine); + + storageService.removeStorageEngine(storageEngineName); + verify(storageEngineStats).unsetStorageEngine(storageEngineName); + verify(storageEngine).drop(); + + storageService.getStorageEngineRepository().addLocalStorageEngine(storageEngine); + storageService.closeStorageEngine(storageEngineName); + verify(storageEngineStats, times(2)).unsetStorageEngine(storageEngineName); + verify(storageEngine).close(); + } + @Test public void testGetStoreAndUserPartitionsMapping() { VeniceConfigLoader configLoader = mock(VeniceConfigLoader.class); diff --git a/docs/operations/data-management/version-lifecycle.md b/docs/operations/data-management/version-lifecycle.md new file mode 100644 index 00000000000..0a3db0db84e --- /dev/null +++ b/docs/operations/data-management/version-lifecycle.md @@ -0,0 +1,99 @@ +# Version Lifecycle and Deletion + +Venice protects versions that may still serve reads or complete an in-flight push, while cleaning +terminal non-current copies by count and time. A version's status alone is not enough to decide deletion: +the controller also considers whether it is current, the newest completed bootstrap (future), a +backup, within configured retention, or part of a store migration. + +## Deletion decision table + +The count-based sweep is implemented by `Store.retrieveVersionsToDelete`. The time-based sweep is +implemented by `StoreBackupVersionCleanupService`. Deferred version swap terminal transitions +directly delete bootstrap-complete, non-current child copies via `ControllerClient.deleteOldVersion`, +so they cannot remain invisible to both sweeps. + +| Initial status | Role / initial condition | Count retention | Time / safety gate | Trigger | Deletion decision | +| --- | --- | --- | --- | --- | --- | +| Any | Current version | Any | Any | Any cleanup | **KEEP** | +| `NOT_CREATED` or `CREATED` | Non-current version at or above current | Not eligible | Not below current | Any cleanup | **KEEP** | +| `NOT_CREATED` or `CREATED` | Below current before standard cleanup is ready | Not eligible | Minimum/default retention gate not met | Time sweep | **DEFER** | +| `NOT_CREATED` or `CREATED` | Below current after standard cleanup is ready | Not eligible | Minimum delay met with multiple eligible backups, or retention expired | Time sweep | **DELETE** | +| `STARTED` | Newest version (active ingestion) | Protected | Not below current | Any cleanup | **KEEP** | +| `STARTED` | Older version at or above current while store is migrating | Protected during migration | Not below current | Any cleanup | **KEEP** | +| `STARTED` | Below current while migrating, before standard cleanup is ready | Protected during migration | Minimum/default retention gate not met | Time sweep | **DEFER** | +| `STARTED` | Below current while migrating, after standard cleanup is ready | Protected from count only | Minimum delay met with multiple eligible backups, or retention expired | Time sweep | **DELETE** | +| `STARTED` | Older non-current version, store is not migrating | Not protected | Not eligible | Count sweep | **DELETE** | +| `PUSHED` | Non-current version without a terminal parent decision | Excluded because it may be an active deferred or concurrent swap candidate | Not eligible | Count or time cleanup | **KEEP** | +| `PUSHED` or `ONLINE` | Still current in a child after parent `ERROR` or `ROLLED_BACK` | Protected while serving; reconciliation remains incomplete | Not eligible | Parent terminal reconciliation | **DEFER: keep and retry** | +| `PUSHED` or `ONLINE` | Still current in a child after parent `PARTIALLY_ONLINE` | Serving the intentional partial state | Not eligible | Parent terminal reconciliation | **KEEP** | +| `PUSHED` or non-current `ONLINE` | Parent deferred swap becomes `ERROR`, `PARTIALLY_ONLINE`, or `ROLLED_BACK` | Removed from count sweep | None — immediately deleted | Parent terminal transition | **DELETE** | +| `ONLINE` | Backup within the configured preserved count | Within limit | Not considered | Count sweep | **KEEP** | +| `ONLINE` | Backup beyond the configured preserved count | Exceeds limit | Not considered | Count sweep | **DELETE** | +| `ONLINE` | Backup considered by retention cleanup | Not considered | Before minimum retention | Time sweep | **DEFER** | +| `ONLINE` | Backup considered by retention cleanup | Not considered | Routers or servers do not agree on current version | Time sweep | **DEFER** | +| `ONLINE` | Eligible old backup | Not considered | Retention expired and metadata validation passes | Time sweep | **DELETE** | +| `ERROR` or `KILLED` | Non-current version | Not protected | Minimum cleanup safety gate where applicable | Immediate or periodic cleanup | **DELETE** | +| `ROLLED_BACK` | Non-current version | Excluded from count sweep | Before rolled-back retention expires | Time sweep | **DEFER** | +| `ROLLED_BACK` | Non-current version | Excluded from count sweep | Rolled-back retention expired | Time sweep | **DELETE** | + +Count-based and time-based cleanup are independent triggers. The first applicable trigger may delete +an eligible backup, but neither trigger may delete the current version. `PUSHED` versions are never +deleted from status and version number alone because multiple deferred or concurrent swaps can be +active. A terminal parent decision explicitly deletes bootstrap-complete non-current child copies +via per-region `ControllerClient.deleteOldVersion`. The controller scans every deferred terminal +parent version, including versions superseded by a newer push, and retries while a child is +unreachable, missing the target metadata, in progress, or unexpectedly still current after parent +`ERROR` or `ROLLED_BACK`. + +## State machine + +```mermaid +stateDiagram-v2 + [*] --> NOT_CREATED + NOT_CREATED --> CREATED: metadata initialized + CREATED --> STARTED: ingestion starts + STARTED --> PUSHED: bootstrap completes + STARTED --> ERROR: push fails + STARTED --> KILLED: push is killed + PUSHED --> ONLINE: version swap succeeds + PUSHED --> DELETED: deferred swap terminates while non-current (parent terminal sweep) + ONLINE --> ROLLED_BACK: rollback or abandoned non-current copy + + NOT_CREATED --> DELETED: stale below-current metadata after time gate + CREATED --> DELETED: stale below-current metadata after time gate + STARTED --> DELETED: stale backup via count or time + ONLINE --> DELETED: backup exceeds count or time retention + ERROR --> DELETED: immediate or periodic cleanup + KILLED --> DELETED: immediate or periodic cleanup + ROLLED_BACK --> DELETED: rolled-back retention expires + + note right of STARTED + Newest active ingestion: KEEP + Migrating copy at/above current: KEEP + Below current: time policy still applies + end note + + note right of PUSHED + Before terminal parent decision: KEEP + Current after ERROR/ROLLED_BACK: retry + Current after PARTIALLY_ONLINE: KEEP + Parent terminal + non-current: DELETE immediately + end note + + note right of ONLINE + Current: always KEEP + Backup: count/time policy + end note + + note right of ROLLED_BACK + Before retention expiry: DEFER + end note +``` + +## Backup observability + +The `ingestion.disk.used` gauge sums disk usage across every local version with the same role, so the +`backup` series represents cumulative backup disk rather than one representative version. The +`ingestion.disk.version_count` gauge reports the number of locally loaded storage engines by +`current`, `future`, and `backup` role. Both metrics use bounded role dimensions rather than a +per-version dimension. diff --git a/docs/operations/index.md b/docs/operations/index.md index 901c0a7faf1..2f275e87887 100644 --- a/docs/operations/index.md +++ b/docs/operations/index.md @@ -6,6 +6,8 @@ This guide covers operational tasks for Venice administrators and operators. - [Repush](data-management/repush.md) - Re-ingest data from source of truth to repair inconsistencies or apply schema changes +- [Version Lifecycle and Deletion](data-management/version-lifecycle.md) - Version states, retention gates, and deletion + decisions - [TTL](data-management/ttl.md) - Configure time-to-live to automatically expire old records - [System Stores](data-management/system-stores.md) - Internal stores used by Venice for metadata and coordination diff --git a/internal/venice-client-common/src/main/java/com/linkedin/venice/stats/metrics/AsyncMetricEntityStateOneEnum.java b/internal/venice-client-common/src/main/java/com/linkedin/venice/stats/metrics/AsyncMetricEntityStateOneEnum.java index 52e002bb0bb..ede5f075280 100644 --- a/internal/venice-client-common/src/main/java/com/linkedin/venice/stats/metrics/AsyncMetricEntityStateOneEnum.java +++ b/internal/venice-client-common/src/main/java/com/linkedin/venice/stats/metrics/AsyncMetricEntityStateOneEnum.java @@ -6,6 +6,8 @@ import com.linkedin.venice.stats.metrics.AsyncMetricResolvers.LiveStateResolverOneEnum; import com.linkedin.venice.stats.metrics.AsyncMetricResolvers.ValueResolverOneEnum; import io.opentelemetry.api.common.Attributes; +import io.opentelemetry.api.metrics.ObservableDoubleGauge; +import io.opentelemetry.api.metrics.ObservableLongGauge; import java.util.EnumMap; import java.util.Map; import java.util.function.ObjDoubleConsumer; @@ -36,7 +38,7 @@ * cost is {@code O(|E|)} {@code liveStateResolver} calls plus one {@code measurement.record(...)} * per emitted combo. */ -public class AsyncMetricEntityStateOneEnum & VeniceDimensionInterface> { +public class AsyncMetricEntityStateOneEnum & VeniceDimensionInterface> implements AutoCloseable { private final boolean emitOpenTelemetryMetrics; /** Precomputed per-enum attributes; {@code null} when OTel is disabled. */ private final EnumMap attributesByEnum; @@ -177,4 +179,13 @@ public EnumMap getAttributesByEnum() { public Object getInstrument() { return instrument; } + + @Override + public void close() { + if (instrument instanceof ObservableLongGauge) { + ((ObservableLongGauge) instrument).close(); + } else if (instrument instanceof ObservableDoubleGauge) { + ((ObservableDoubleGauge) instrument).close(); + } + } } diff --git a/internal/venice-client-common/src/main/java/com/linkedin/venice/stats/metrics/AsyncMetricEntityStateTwoEnums.java b/internal/venice-client-common/src/main/java/com/linkedin/venice/stats/metrics/AsyncMetricEntityStateTwoEnums.java index 31bd8eed731..0f67e2114cc 100644 --- a/internal/venice-client-common/src/main/java/com/linkedin/venice/stats/metrics/AsyncMetricEntityStateTwoEnums.java +++ b/internal/venice-client-common/src/main/java/com/linkedin/venice/stats/metrics/AsyncMetricEntityStateTwoEnums.java @@ -6,6 +6,8 @@ import com.linkedin.venice.stats.metrics.AsyncMetricResolvers.LiveStateResolverTwoEnums; import com.linkedin.venice.stats.metrics.AsyncMetricResolvers.ValueResolverTwoEnums; import io.opentelemetry.api.common.Attributes; +import io.opentelemetry.api.metrics.ObservableDoubleGauge; +import io.opentelemetry.api.metrics.ObservableLongGauge; import java.util.EnumMap; import java.util.Map; import java.util.function.ObjDoubleConsumer; @@ -37,7 +39,8 @@ * is {@code O(|E1| × |E2|)} {@code liveStateResolver} calls plus one * {@code measurement.record(...)} per emitted pair. */ -public class AsyncMetricEntityStateTwoEnums & VeniceDimensionInterface, E2 extends Enum & VeniceDimensionInterface> { +public class AsyncMetricEntityStateTwoEnums & VeniceDimensionInterface, E2 extends Enum & VeniceDimensionInterface> + implements AutoCloseable { private final boolean emitOpenTelemetryMetrics; /** Precomputed per-pair attributes; {@code null} when OTel is disabled. */ private final EnumMap> attributesByEnum; @@ -187,4 +190,13 @@ public EnumMap> getAttributesByEnum() { public Object getInstrument() { return instrument; } + + @Override + public void close() { + if (instrument instanceof ObservableLongGauge) { + ((ObservableLongGauge) instrument).close(); + } else if (instrument instanceof ObservableDoubleGauge) { + ((ObservableDoubleGauge) instrument).close(); + } + } } diff --git a/internal/venice-client-common/src/test/java/com/linkedin/venice/stats/metrics/AsyncMetricEntityStateOneEnumTest.java b/internal/venice-client-common/src/test/java/com/linkedin/venice/stats/metrics/AsyncMetricEntityStateOneEnumTest.java index aafbafb4be3..9c6bdf82d49 100644 --- a/internal/venice-client-common/src/test/java/com/linkedin/venice/stats/metrics/AsyncMetricEntityStateOneEnumTest.java +++ b/internal/venice-client-common/src/test/java/com/linkedin/venice/stats/metrics/AsyncMetricEntityStateOneEnumTest.java @@ -19,7 +19,9 @@ import com.linkedin.venice.stats.metrics.AsyncMetricResolvers.LiveStateResolverOneEnum; import com.linkedin.venice.stats.metrics.AsyncMetricResolvers.ValueResolverOneEnum; import io.opentelemetry.api.common.Attributes; +import io.opentelemetry.api.metrics.ObservableDoubleGauge; import io.opentelemetry.api.metrics.ObservableDoubleMeasurement; +import io.opentelemetry.api.metrics.ObservableLongGauge; import io.opentelemetry.api.metrics.ObservableLongMeasurement; import java.util.EnumMap; import java.util.EnumSet; @@ -84,9 +86,57 @@ public void testCreateRegistersExactlyOneObservableGauge() { } } + @Test + public void testCloseUnregistersObservableGauge() { + ObservableLongGauge gauge = mock(ObservableLongGauge.class); + when(mockOtelRepository.registerObservableLongGauge(eq(mockMetricEntity), any())).thenReturn(gauge); + AsyncMetricEntityStateOneEnum metricState = AsyncMetricEntityStateOneEnum.create( + mockMetricEntity, + mockOtelRepository, + baseDimensionsMap, + DimensionEnum1.class, + e -> e, + (state, e) -> 1L); + + metricState.close(); + + verify(gauge).close(); + } + + @Test + public void testCloseUnregistersObservableDoubleGauge() { + when(mockMetricEntity.getMetricType()).thenReturn(MetricType.ASYNC_DOUBLE_GAUGE); + ObservableDoubleGauge gauge = mock(ObservableDoubleGauge.class); + when(mockOtelRepository.registerObservableDoubleGauge(eq(mockMetricEntity), any())).thenReturn(gauge); + AsyncMetricEntityStateOneEnum metricState = AsyncMetricEntityStateOneEnum.create( + mockMetricEntity, + mockOtelRepository, + baseDimensionsMap, + DimensionEnum1.class, + e -> e, + (state, e) -> 1L); + + metricState.close(); + + verify(gauge).close(); + } + + @Test + public void testCloseOnDisabledInstanceIsNoOp() { + AsyncMetricEntityStateOneEnum metricState = AsyncMetricEntityStateOneEnum.create( + mockMetricEntity, + null /* OTel disabled */, + baseDimensionsMap, + DimensionEnum1.class, + e -> e, + (state, e) -> 1L); + + // Must not throw when instrument is null. + metricState.close(); + } + @Test public void testCallbackEmitsOnlyWhenLiveStateResolverReturnsNonNull() { - // liveStateResolver returns state for DIMENSION_ONE only; DIMENSION_TWO is dormant. LiveStateResolverOneEnum liveStateResolver = e -> e == DimensionEnum1.DIMENSION_ONE ? "live" : null; ValueResolverOneEnum valueResolver = (state, e) -> 42L; diff --git a/internal/venice-client-common/src/test/java/com/linkedin/venice/stats/metrics/AsyncMetricEntityStateTwoEnumsTest.java b/internal/venice-client-common/src/test/java/com/linkedin/venice/stats/metrics/AsyncMetricEntityStateTwoEnumsTest.java index b8c70b57286..04a3060ee11 100644 --- a/internal/venice-client-common/src/test/java/com/linkedin/venice/stats/metrics/AsyncMetricEntityStateTwoEnumsTest.java +++ b/internal/venice-client-common/src/test/java/com/linkedin/venice/stats/metrics/AsyncMetricEntityStateTwoEnumsTest.java @@ -20,7 +20,9 @@ import com.linkedin.venice.stats.metrics.AsyncMetricResolvers.LiveStateResolverTwoEnums; import com.linkedin.venice.stats.metrics.AsyncMetricResolvers.ValueResolverTwoEnums; import io.opentelemetry.api.common.Attributes; +import io.opentelemetry.api.metrics.ObservableDoubleGauge; import io.opentelemetry.api.metrics.ObservableDoubleMeasurement; +import io.opentelemetry.api.metrics.ObservableLongGauge; import io.opentelemetry.api.metrics.ObservableLongMeasurement; import java.util.HashMap; import java.util.HashSet; @@ -94,6 +96,58 @@ public void testCreateRegistersExactlyOneObservableGauge() { assertEquals(leafCount, DimensionEnum1.values().length * DimensionEnum2.values().length); } + @Test + public void testCloseUnregistersObservableGauge() { + ObservableLongGauge gauge = mock(ObservableLongGauge.class); + when(mockOtelRepository.registerObservableLongGauge(eq(mockMetricEntity), any())).thenReturn(gauge); + AsyncMetricEntityStateTwoEnums metricState = AsyncMetricEntityStateTwoEnums.create( + mockMetricEntity, + mockOtelRepository, + baseDimensionsMap, + DimensionEnum1.class, + DimensionEnum2.class, + (e1, e2) -> e1, + (state, e1, e2) -> 1L); + + metricState.close(); + + verify(gauge).close(); + } + + @Test + public void testCloseUnregistersObservableDoubleGauge() { + when(mockMetricEntity.getMetricType()).thenReturn(MetricType.ASYNC_DOUBLE_GAUGE); + ObservableDoubleGauge gauge = mock(ObservableDoubleGauge.class); + when(mockOtelRepository.registerObservableDoubleGauge(eq(mockMetricEntity), any())).thenReturn(gauge); + AsyncMetricEntityStateTwoEnums metricState = AsyncMetricEntityStateTwoEnums.create( + mockMetricEntity, + mockOtelRepository, + baseDimensionsMap, + DimensionEnum1.class, + DimensionEnum2.class, + (e1, e2) -> e1, + (state, e1, e2) -> 1L); + + metricState.close(); + + verify(gauge).close(); + } + + @Test + public void testCloseOnDisabledInstanceIsNoOp() { + AsyncMetricEntityStateTwoEnums metricState = AsyncMetricEntityStateTwoEnums.create( + mockMetricEntity, + null /* OTel disabled */, + baseDimensionsMap, + DimensionEnum1.class, + DimensionEnum2.class, + (e1, e2) -> e1, + (state, e1, e2) -> 1L); + + // Must not throw when instrument is null. + metricState.close(); + } + @Test public void testCallbackEmitsOnlyWhenLiveStateResolverReturnsNonNull() { LiveStateResolverTwoEnums liveStateResolver = (e1, e2) -> { diff --git a/internal/venice-common/src/main/java/com/linkedin/venice/meta/AbstractStore.java b/internal/venice-common/src/main/java/com/linkedin/venice/meta/AbstractStore.java index 44c494e28a9..691d60235c9 100644 --- a/internal/venice-common/src/main/java/com/linkedin/venice/meta/AbstractStore.java +++ b/internal/venice-common/src/main/java/com/linkedin/venice/meta/AbstractStore.java @@ -343,6 +343,7 @@ public static List computeVersionsToDelete( * b) ERROR version (ideally should not be there as AbstractPushmonitor#handleErrorPush deletes those) * c) STARTED versions if its not the last one and the store is not migrating. * d) KILLED versions by {@link org.apache.kafka.clients.admin.Admin#killOfflinePush} api. + * e) ROLLED_BACK versions after their dedicated time-based retention expires. */ // current version need not be the largest version, preseve it before finding other versions > current version for (int i = lastElementIndex; i >= 0; i--) { @@ -353,26 +354,27 @@ public static List computeVersionsToDelete( for (int i = lastElementIndex; i >= 0; i--) { Version version = versions.get(i); + VersionStatus status = version.getStatus(); if (version.getNumber() == currentVersion) { // currentVersion is always preserved continue; } - if (VersionStatus.isVersionRolledBack(version.getStatus())) { + if (VersionStatus.isVersionRolledBack(status)) { // ROLLED_BACK versions are retained and cleaned up by StoreBackupVersionCleanupService // with a dedicated retention period. Skip them here so the retention window is honored. // Note: if the backup version retention-based cleanup service is disabled for the cluster, // ROLLED_BACK versions will not be automatically deleted by this path. The admin tool's // deleteOldVersion can still be used for manual cleanup in that case. continue; - } else if (VersionStatus.canDelete(version.getStatus())) { // ERROR and KILLED versions are always deleted + } else if (VersionStatus.canDelete(status)) { // ERROR and KILLED versions are always deleted versionsToDelete.add(version); - } else if (VersionStatus.ONLINE.equals(version.getStatus())) { + } else if (VersionStatus.preserveLastFew(status)) { if (curNumVersionsToPreserve > 0) { // keep the minimum number of version to preserve curNumVersionsToPreserve--; } else { versionsToDelete.add(version); } - } else if (VersionStatus.STARTED.equals(version.getStatus()) && (i != lastElementIndex) && !isMigrating) { + } else if (VersionStatus.STARTED.equals(status) && (i != lastElementIndex) && !isMigrating) { // For the non-last started version, if it's not the current version(STARTED version should not be the current // version, just prevent some edge cases here.), we should delete it only if the store is not migrating // as during store migration are there are concurrent pushes with STARTED version. @@ -380,8 +382,8 @@ public static List computeVersionsToDelete( // version status properly. versionsToDelete.add(version); } - // TODO here we don't deal with the PUSHED version, just keep all of them, need to consider collect them too in - // the future. + // PUSHED versions can still be active deferred or concurrent swap candidates. They require an + // explicit terminal transition before cleanup and are never deleted by count alone. } return versionsToDelete; } diff --git a/mkdocs.yml b/mkdocs.yml index a10502b2951..a2758c5f3cc 100644 --- a/mkdocs.yml +++ b/mkdocs.yml @@ -121,6 +121,7 @@ nav: - operations/index.md - Data Management: - Repush: operations/data-management/repush.md + - Version Lifecycle and Deletion: operations/data-management/version-lifecycle.md - TTL: operations/data-management/ttl.md - System Stores: operations/data-management/system-stores.md - Alerting: diff --git a/services/venice-controller/src/main/java/com/linkedin/venice/controller/DeferredVersionSwapService.java b/services/venice-controller/src/main/java/com/linkedin/venice/controller/DeferredVersionSwapService.java index 5f87b36c3f7..561f17174ac 100644 --- a/services/venice-controller/src/main/java/com/linkedin/venice/controller/DeferredVersionSwapService.java +++ b/services/venice-controller/src/main/java/com/linkedin/venice/controller/DeferredVersionSwapService.java @@ -85,8 +85,12 @@ public class DeferredVersionSwapService extends AbstractVeniceService { private static final Set VERSION_SWAP_COMPLETION_STATUSES = Utils.setOf(ONLINE, PARTIALLY_ONLINE, ERROR); private static final Set TERMINAL_PUSH_VERSION_STATUSES = Utils.setOf(ONLINE); + private static final Set ABANDONED_VERSION_STATUSES = + Utils.setOf(ERROR, PARTIALLY_ONLINE, VersionStatus.ROLLED_BACK); private Cache storeWaitTimeCacheForSequentialRollout = Caffeine.newBuilder().expireAfterWrite(1, TimeUnit.HOURS).build(); + private final Cache terminalVersionReconciliationCache = + Caffeine.newBuilder().maximumSize(100_000).expireAfterWrite(1, TimeUnit.DAYS).build(); private static final int CONTROLLER_CLIENT_REQUEST_TIMEOUT = 1 * Time.MS_PER_SECOND; private static final int LOG_LATENCY_THRESHOLD = 5 * Time.MS_PER_SECOND; private final Map clusterToExecutorMap = new ConcurrentHashMap<>(); @@ -421,6 +425,20 @@ private boolean isPushInTerminalState( logMessageIfNotRedundant(message); return true; } + break; + case ERROR: + case PARTIALLY_ONLINE: + case ROLLED_BACK: + String reconciliationKey = getVersionProcessingKey(clusterName, storeName, targetVersionNum); + if (terminalVersionReconciliationCache.getIfPresent(reconciliationKey) == null + && reconcileAbandonedVersionInChildRegions( + clusterName, + storeName, + targetVersionNum, + targetVersion.getStatus())) { + terminalVersionReconciliationCache.put(reconciliationKey, true); + } + return false; } logMessageIfNotRedundant( @@ -899,7 +917,7 @@ private void handleFailedRollForward( return v + 1; }); - if (attemptedRetries == MAX_ROLL_FORWARD_RETRY_LIMIT) { + if (attemptedRetries >= MAX_ROLL_FORWARD_RETRY_LIMIT) { deferredVersionSwapStats.recordDeferredVersionSwapFailedRollForwardMetric(clusterName, parentStore.getName()); updateStore(clusterName, parentStore.getName(), PARTIALLY_ONLINE, targetVersionNum); failedRollforwardRetryCountMap.remove(kafkaTopicName); @@ -948,6 +966,56 @@ private void finishProcessingStore(String kafkaTopicName) { storesBeingProcessed.remove(kafkaTopicName); } + private static String getVersionProcessingKey(String clusterName, String storeName, int versionNumber) { + return clusterName + ":" + Version.composeKafkaTopic(storeName, versionNumber); + } + + private static boolean isAbandonedDeferredVersion(Version version) { + return version != null && version.isVersionSwapDeferred() + && ABANDONED_VERSION_STATUSES.contains(version.getStatus()); + } + + private void submitTerminalVersionReconciliationTasks( + String clusterName, + Store parentStore, + ThreadPoolExecutor clusterExecutorService) { + for (Version version: parentStore.getVersions()) { + if (!isAbandonedDeferredVersion(version)) { + continue; + } + + int versionNumber = version.getNumber(); + VersionStatus parentStatus = version.getStatus(); + String processingKey = getVersionProcessingKey(clusterName, parentStore.getName(), versionNumber); + if (terminalVersionReconciliationCache.getIfPresent(processingKey) != null + || !tryStartProcessingStore(processingKey)) { + continue; + } + + clusterExecutorService.submit(() -> { + try { + if (reconcileAbandonedVersionInChildRegions( + clusterName, + parentStore.getName(), + versionNumber, + parentStatus)) { + terminalVersionReconciliationCache.put(processingKey, true); + } + } catch (Exception e) { + LOGGER.warn( + "Caught exception while reconciling terminal version: {} for store: {} in cluster: {}", + versionNumber, + parentStore.getName(), + clusterName, + e); + deferredVersionSwapStats.recordDeferredVersionSwapExceptionMetric(clusterName); + } finally { + finishProcessingStore(processingKey); + } + }); + } + } + private Runnable getRunnableForDeferredVersionSwap() { return () -> { LogContext.setLogContext(veniceControllerMultiClusterConfig.getLogContext()); @@ -982,6 +1050,11 @@ private Runnable getRunnableForDeferredVersionSwap() { // Filter out stores that aren't doing a target region push w/ deferred swap List eligibleStoresToProcess = new ArrayList<>(); for (Store parentStore: parentStores) { + submitTerminalVersionReconciliationTasks(cluster, parentStore, clusterExecutorService); + Version latestVersion = parentStore.getVersion(parentStore.getLargestUsedVersionNumber()); + if (isAbandonedDeferredVersion(latestVersion)) { + continue; + } if (!isTargetRegionPushWithDeferredSwapEnabled(parentStore)) { continue; } @@ -994,9 +1067,9 @@ private Runnable getRunnableForDeferredVersionSwap() { Version targetVersion = parentStore.getVersion(parentStore.getLargestUsedVersionNumber()); // Check if store is already being processed - String kafkaTopicName = - Version.composeKafkaTopic(parentStore.getName(), parentStore.getLargestUsedVersionNumber()); - if (!tryStartProcessingStore(kafkaTopicName)) { + String processingKey = + getVersionProcessingKey(cluster, parentStore.getName(), parentStore.getLargestUsedVersionNumber()); + if (!tryStartProcessingStore(processingKey)) { String message = "Skipping store " + parentStore.getName() + " as it's already being processed"; logMessageIfNotRedundant(message); continue; @@ -1015,8 +1088,16 @@ private Runnable getRunnableForDeferredVersionSwap() { } else { performParallelRollForward(cluster, parentStore, childControllerClientMap, targetVersion); } + } catch (Exception e) { + LOGGER.warn( + "Caught exception while processing deferred version swap for store: {} in cluster: {}", + parentStore.getName(), + cluster, + e); + deferredVersionSwapStats.recordDeferredVersionSwapExceptionMetric(cluster); } finally { - finishProcessingStore(Version.composeKafkaTopic(parentStore.getName(), targetVersion.getNumber())); + finishProcessingStore( + getVersionProcessingKey(cluster, parentStore.getName(), targetVersion.getNumber())); } }); } @@ -1251,36 +1332,116 @@ private void performParallelRollForward( public void updateStore(String clusterName, String storeName, VersionStatus status, int targetVersionNum) { HelixVeniceClusterResources resources = veniceParentHelixAdmin.getVeniceHelixAdmin().getHelixVeniceClusterResources(clusterName); - try (AutoCloseableLock ignore = resources.getClusterLockManager().createStoreWriteLock(storeName)) { - ReadWriteStoreRepository repository = resources.getStoreMetadataRepository(); - Store store = repository.getStore(storeName); - LOGGER.info( - "Updating store: {} version: {} from status {} to status {}", - storeName, - targetVersionNum, - store.getVersionStatus(targetVersionNum), - status); - store.updateVersionStatus(targetVersionNum, status); - if (status == ONLINE || status == PARTIALLY_ONLINE) { - store.setCurrentVersion(targetVersionNum); - - // For jobs that stop polling early or for pushes that don't poll (empty push), we need to truncate the parent - // VT here to unblock the next push - String kafkaTopicName = Version.composeKafkaTopic(storeName, targetVersionNum); - ConcurrentPushDetectionStrategy strategy = - veniceControllerMultiClusterConfig.getControllerConfig(clusterName).getConcurrentPushDetectionStrategy(); - // skip truncating if the topic was not created based on ConcurrentPushDetectionStrategy - if (strategy.isTopicWriteNeeded() && !veniceParentHelixAdmin.isTopicTruncated(kafkaTopicName)) { - LOGGER.info("Truncating parent VT for {}", kafkaTopicName); - veniceParentHelixAdmin.truncateKafkaTopic(Version.composeKafkaTopic(storeName, targetVersionNum)); + try { + try (AutoCloseableLock ignore = resources.getClusterLockManager().createStoreWriteLock(storeName)) { + ReadWriteStoreRepository repository = resources.getStoreMetadataRepository(); + Store store = repository.getStore(storeName); + LOGGER.info( + "Updating store: {} version: {} from status {} to status {}", + storeName, + targetVersionNum, + store.getVersionStatus(targetVersionNum), + status); + store.updateVersionStatus(targetVersionNum, status); + if (status == ONLINE || status == PARTIALLY_ONLINE) { + store.setCurrentVersion(targetVersionNum); + + // For jobs that stop polling early or for pushes that don't poll (empty push), we need to truncate the parent + // VT here to unblock the next push + String kafkaTopicName = Version.composeKafkaTopic(storeName, targetVersionNum); + ConcurrentPushDetectionStrategy strategy = + veniceControllerMultiClusterConfig.getControllerConfig(clusterName).getConcurrentPushDetectionStrategy(); + // skip truncating if the topic was not created based on ConcurrentPushDetectionStrategy + if (strategy.isTopicWriteNeeded() && !veniceParentHelixAdmin.isTopicTruncated(kafkaTopicName)) { + LOGGER.info("Truncating parent VT for {}", kafkaTopicName); + veniceParentHelixAdmin.truncateKafkaTopic(Version.composeKafkaTopic(storeName, targetVersionNum)); + } + } + repository.updateStore(store); + } + + if (status == ERROR || status == PARTIALLY_ONLINE || status == VersionStatus.ROLLED_BACK) { + String reconciliationKey = getVersionProcessingKey(clusterName, storeName, targetVersionNum); + if (reconcileAbandonedVersionInChildRegions(clusterName, storeName, targetVersionNum, status)) { + terminalVersionReconciliationCache.put(reconciliationKey, true); } } - repository.updateStore(store); } catch (Exception e) { LOGGER.warn("Failed to execute updateStore for store: {} in cluster: {}", storeName, clusterName, e); } } + private boolean reconcileAbandonedVersionInChildRegions( + String clusterName, + String storeName, + int targetVersionNum, + VersionStatus parentStatus) { + Map controllerClientMap = + veniceParentHelixAdmin.getVeniceHelixAdmin().getControllerClientMap(clusterName); + if (controllerClientMap == null || controllerClientMap.isEmpty()) { + throw new VeniceException( + "Cannot reconcile abandoned version " + targetVersionNum + " for store " + storeName + " in cluster " + + clusterName + ": no child regions are configured"); + } + + Set regionsToReconcile = new HashSet<>(); + boolean allRegionsTerminal = true; + for (String region: controllerClientMap.keySet()) { + StoreResponse storeResponse = getStoreForRegion(clusterName, region, storeName); + if (storeResponse == null || storeResponse.getStore() == null) { + allRegionsTerminal = false; + continue; + } + + StoreInfo childStore = storeResponse.getStore(); + Version childVersion = getVersionFromStoreInRegion(region, storeName, targetVersionNum, storeResponse); + if (childVersion == null) { + allRegionsTerminal = false; + continue; + } + if (VersionStatus.isVersionRolledBack(childVersion.getStatus()) + || VersionStatus.canDelete(childVersion.getStatus())) { + continue; + } + if (childStore.getCurrentVersion() == targetVersionNum) { + if (parentStatus != PARTIALLY_ONLINE) { + allRegionsTerminal = false; + } + continue; + } + + if (VersionStatus.isBootstrapCompleted(childVersion.getStatus())) { + regionsToReconcile.add(region); + } else { + allRegionsTerminal = false; + } + } + + if (!regionsToReconcile.isEmpty()) { + for (String region: regionsToReconcile) { + try { + controllerClientMap.get(region).deleteOldVersion(storeName, targetVersionNum); + LOGGER.info( + "Deleted abandoned deferred version {} for store {} in region {} (parent status: {})", + targetVersionNum, + storeName, + region, + parentStatus); + } catch (Exception e) { + LOGGER.warn( + "Failed to delete abandoned deferred version {} for store {} in region {} (parent status: {})", + targetVersionNum, + storeName, + region, + parentStatus, + e); + allRegionsTerminal = false; + } + } + } + return allRegionsTerminal; + } + private void markTargetRegionPromoted(String clusterName, String storeName, int targetVersionNum) { try { LOGGER.info("Marking targetRegionPromoted=true for store: {} version: {}", storeName, targetVersionNum); diff --git a/services/venice-controller/src/main/java/com/linkedin/venice/controller/StoreBackupVersionCleanupService.java b/services/venice-controller/src/main/java/com/linkedin/venice/controller/StoreBackupVersionCleanupService.java index 6a9941eda9b..d9bc875bb16 100644 --- a/services/venice-controller/src/main/java/com/linkedin/venice/controller/StoreBackupVersionCleanupService.java +++ b/services/venice-controller/src/main/java/com/linkedin/venice/controller/StoreBackupVersionCleanupService.java @@ -156,10 +156,10 @@ protected static boolean whetherStoreReadyToBeCleanup( long minCleanupDelayMs) { List versions = store.getVersions(); - // Regardless of retention, if there are more than 1 non-rolled-back versions strictly below the current version, - // we should clean up. ROLLED_BACK versions are excluded since they have their own retention-based cleanup path. + // Regardless of retention, clean up when more than one standard-cleanup-eligible version is below current. + // PUSHED versions remain protected until an explicit terminal decision, and ROLLED_BACK has its own retention path. if (versions.stream() - .filter(v -> v.getNumber() < currentVersion && !VersionStatus.isVersionRolledBack(v.getStatus())) + .filter(v -> v.getNumber() < currentVersion && isEligibleForStandardBackupCleanup(v)) .count() > 1) { return true; } @@ -333,7 +333,7 @@ protected boolean cleanupBackupVersion(Store store, String clusterName) { HashSet 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 {