diff --git a/internal/venice-test-common/src/integrationTest/java/com/linkedin/venice/controller/StoreUpdateHandlerIntegrationTest.java b/internal/venice-test-common/src/integrationTest/java/com/linkedin/venice/controller/StoreUpdateHandlerIntegrationTest.java new file mode 100644 index 00000000000..4eea99339ff --- /dev/null +++ b/internal/venice-test-common/src/integrationTest/java/com/linkedin/venice/controller/StoreUpdateHandlerIntegrationTest.java @@ -0,0 +1,197 @@ +package com.linkedin.venice.controller; + +import static com.linkedin.venice.ConfigKeys.CONTROLLER_AUTO_MATERIALIZE_DAVINCI_PUSH_STATUS_SYSTEM_STORE; +import static com.linkedin.venice.ConfigKeys.CONTROLLER_AUTO_MATERIALIZE_META_SYSTEM_STORE; +import static com.linkedin.venice.controllerapi.ControllerApiConstants.READ_QUOTA_IN_CU; + +import com.linkedin.venice.controllerapi.ControllerClient; +import com.linkedin.venice.controllerapi.ControllerResponse; +import com.linkedin.venice.controllerapi.NewStoreResponse; +import com.linkedin.venice.controllerapi.UpdateStoreQueryParams; +import com.linkedin.venice.exceptions.VeniceRetriableException; +import com.linkedin.venice.integration.utils.ServiceFactory; +import com.linkedin.venice.integration.utils.VeniceControllerWrapper; +import com.linkedin.venice.integration.utils.VeniceMultiRegionClusterCreateOptions; +import com.linkedin.venice.integration.utils.VeniceTwoLayerMultiRegionMultiClusterWrapper; +import com.linkedin.venice.meta.Store; +import com.linkedin.venice.meta.StoreInfo; +import com.linkedin.venice.utils.TestUtils; +import com.linkedin.venice.utils.Time; +import com.linkedin.venice.utils.Utils; +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; +import java.util.Properties; +import java.util.Set; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; +import org.testng.Assert; +import org.testng.annotations.Test; + + +public class StoreUpdateHandlerIntegrationTest { + private static final long TEST_TIMEOUT_MS = 2 * Time.MS_PER_MINUTE; + private static final long UPDATED_READ_QUOTA = 1234; + + @Test(timeOut = TEST_TIMEOUT_MS) + public void testStoreUpdateHandlerRetriesWithFinalReadOnlySnapshot() throws InterruptedException { + String storeName = Utils.getUniqueString("store-update-handler"); + String originalOwner = "test-owner"; + RetryingStoreUpdateHandler storeUpdateHandler = new RetryingStoreUpdateHandler(storeName, UPDATED_READ_QUOTA); + AtomicInteger childHandlerInvocationCount = new AtomicInteger(); + + Properties parentControllerProperties = new Properties(); + parentControllerProperties.setProperty(CONTROLLER_AUTO_MATERIALIZE_META_SYSTEM_STORE, Boolean.FALSE.toString()); + parentControllerProperties + .setProperty(CONTROLLER_AUTO_MATERIALIZE_DAVINCI_PUSH_STATUS_SYSTEM_STORE, Boolean.FALSE.toString()); + parentControllerProperties.put(VeniceControllerWrapper.STORE_UPDATE_HANDLER, storeUpdateHandler); + Properties childControllerProperties = new Properties(); + childControllerProperties.setProperty(CONTROLLER_AUTO_MATERIALIZE_META_SYSTEM_STORE, Boolean.FALSE.toString()); + childControllerProperties + .setProperty(CONTROLLER_AUTO_MATERIALIZE_DAVINCI_PUSH_STATUS_SYSTEM_STORE, Boolean.FALSE.toString()); + childControllerProperties.put( + VeniceControllerWrapper.STORE_UPDATE_HANDLER, + (StoreUpdateHandler) (clusterName, store, updatedConfigs) -> childHandlerInvocationCount.incrementAndGet()); + + VeniceMultiRegionClusterCreateOptions options = + new VeniceMultiRegionClusterCreateOptions.Builder().numberOfRegions(1) + .numberOfClusters(1) + .numberOfParentControllers(1) + .numberOfChildControllers(1) + .numberOfServers(0) + .numberOfRouters(0) + .replicationFactor(1) + .parentControllerProperties(parentControllerProperties) + .childControllerProperties(childControllerProperties) + .build(); + + try (VeniceTwoLayerMultiRegionMultiClusterWrapper venice = + ServiceFactory.getVeniceTwoLayerMultiRegionMultiClusterWrapper(options)) { + String clusterName = venice.getClusterNames()[0]; + String childControllerUrl = venice.getChildRegions().get(0).getControllerConnectString(); + try ( + ControllerClient parentControllerClient = + new ControllerClient(clusterName, venice.getControllerConnectString()); + ControllerClient childControllerClient = new ControllerClient(clusterName, childControllerUrl)) { + NewStoreResponse newStoreResponse = + parentControllerClient.createNewStore(storeName, originalOwner, "\"string\"", "\"string\""); + Assert.assertFalse(newStoreResponse.isError(), newStoreResponse.getError()); + TestUtils.waitForNonDeterministicAssertion( + 30, + TimeUnit.SECONDS, + () -> Assert.assertFalse(childControllerClient.getStore(storeName).isError())); + + ControllerResponse updateStoreResponse = parentControllerClient + .updateStore(storeName, new UpdateStoreQueryParams().setReadQuotaInCU(UPDATED_READ_QUOTA)); + Assert.assertFalse(updateStoreResponse.isError(), updateStoreResponse.getError()); + Assert.assertTrue( + storeUpdateHandler.awaitSuccessfulInvocation(30, TimeUnit.SECONDS), + "The store update handler did not succeed after its first-attempt failure"); + + Store callbackStore = storeUpdateHandler.getLatestStore(); + Assert.assertTrue(storeUpdateHandler.getInvocationCount() >= 2); + Assert.assertEquals(storeUpdateHandler.getLatestClusterName(), clusterName); + Assert.assertEquals(callbackStore.getName(), storeName); + Assert.assertEquals(callbackStore.getOwner(), originalOwner); + Assert.assertEquals(callbackStore.getReadQuotaInCU(), UPDATED_READ_QUOTA); + Assert.assertTrue(storeUpdateHandler.receivedOnlyReadOnlyStores()); + Assert.assertTrue(storeUpdateHandler.receivedOnlyImmutableUpdatedConfigs()); + Assert.assertTrue( + storeUpdateHandler.getReceivedUpdatedConfigs() + .stream() + .allMatch(updatedConfigs -> updatedConfigs.equals(Collections.singleton(READ_QUOTA_IN_CU)))); + + String barrierOwner = "owner-after-update"; + ControllerResponse setOwnerResponse = parentControllerClient.setStoreOwner(storeName, barrierOwner); + Assert.assertFalse(setOwnerResponse.isError(), setOwnerResponse.getError()); + TestUtils.waitForNonDeterministicAssertion(30, TimeUnit.SECONDS, () -> { + StoreInfo childStore = childControllerClient.getStore(storeName).getStore(); + Assert.assertEquals(childStore.getOwner(), barrierOwner); + Assert.assertEquals(childStore.getReadQuotaInCU(), UPDATED_READ_QUOTA); + }); + + StoreInfo parentStore = parentControllerClient.getStore(storeName).getStore(); + Assert.assertEquals(parentStore.getReadQuotaInCU(), UPDATED_READ_QUOTA); + Assert.assertTrue(storeUpdateHandler.getInvocationCount() >= 2); + Assert.assertEquals(childHandlerInvocationCount.get(), 0); + } + } + } + + private static final class RetryingStoreUpdateHandler implements StoreUpdateHandler { + private final String targetStoreName; + private final long targetReadQuota; + private final AtomicInteger invocationCount = new AtomicInteger(); + private final AtomicReference latestClusterName = new AtomicReference<>(); + private final AtomicReference latestStore = new AtomicReference<>(); + private final AtomicBoolean receivedOnlyReadOnlyStores = new AtomicBoolean(true); + private final AtomicBoolean receivedOnlyImmutableUpdatedConfigs = new AtomicBoolean(true); + private final CopyOnWriteArrayList> receivedUpdatedConfigs = new CopyOnWriteArrayList<>(); + private final CountDownLatch successfulInvocation = new CountDownLatch(1); + + private RetryingStoreUpdateHandler(String targetStoreName, long targetReadQuota) { + this.targetStoreName = targetStoreName; + this.targetReadQuota = targetReadQuota; + } + + @Override + public void handleStoreUpdate(String clusterName, Store store, Set updatedConfigs) { + if (!targetStoreName.equals(store.getName()) || store.getReadQuotaInCU() != targetReadQuota) { + return; + } + + latestClusterName.set(clusterName); + latestStore.set(store); + receivedUpdatedConfigs.add(updatedConfigs); + try { + store.setOwner("unexpected-mutation"); + receivedOnlyReadOnlyStores.set(false); + } catch (UnsupportedOperationException expected) { + // Expected for callback snapshots. + } + try { + updatedConfigs.add("unexpected-config"); + receivedOnlyImmutableUpdatedConfigs.set(false); + } catch (UnsupportedOperationException expected) { + // Expected for callback config sets. + } + + if (invocationCount.incrementAndGet() == 1) { + throw new VeniceRetriableException("Expected first-attempt store update handler failure"); + } + successfulInvocation.countDown(); + } + + private boolean awaitSuccessfulInvocation(long timeout, TimeUnit unit) throws InterruptedException { + return successfulInvocation.await(timeout, unit); + } + + private int getInvocationCount() { + return invocationCount.get(); + } + + private String getLatestClusterName() { + return latestClusterName.get(); + } + + private Store getLatestStore() { + return latestStore.get(); + } + + private boolean receivedOnlyReadOnlyStores() { + return receivedOnlyReadOnlyStores.get(); + } + + private boolean receivedOnlyImmutableUpdatedConfigs() { + return receivedOnlyImmutableUpdatedConfigs.get(); + } + + private List> getReceivedUpdatedConfigs() { + return new ArrayList<>(receivedUpdatedConfigs); + } + } +} diff --git a/internal/venice-test-common/src/integrationTest/java/com/linkedin/venice/integration/utils/VeniceControllerWrapper.java b/internal/venice-test-common/src/integrationTest/java/com/linkedin/venice/integration/utils/VeniceControllerWrapper.java index c37c8e69a1d..5c1e8469508 100644 --- a/internal/venice-test-common/src/integrationTest/java/com/linkedin/venice/integration/utils/VeniceControllerWrapper.java +++ b/internal/venice-test-common/src/integrationTest/java/com/linkedin/venice/integration/utils/VeniceControllerWrapper.java @@ -68,6 +68,7 @@ import com.linkedin.venice.acl.VeniceComponent; import com.linkedin.venice.client.store.ClientConfig; import com.linkedin.venice.controller.Admin; +import com.linkedin.venice.controller.StoreUpdateHandler; import com.linkedin.venice.controller.VeniceController; import com.linkedin.venice.controller.VeniceControllerContext; import com.linkedin.venice.controller.VeniceHelixAdmin; @@ -117,6 +118,7 @@ public class VeniceControllerWrapper extends ProcessWrapper { public static final String PARENT_D2_SERVICE_NAME = "ParentController"; public static final String SUPERSET_SCHEMA_GENERATOR = "SupersetSchemaGenerator"; + public static final String STORE_UPDATE_HANDLER = "StoreUpdateHandler"; public static final double DEFAULT_STORAGE_ENGINE_OVERHEAD_RATIO = 0.85d; @@ -421,6 +423,11 @@ static StatefulServiceProvider generateService(VeniceCo if (passedSupersetSchemaGenerator instanceof SupersetSchemaGenerator) { supersetSchemaGenerator = Optional.of((SupersetSchemaGenerator) passedSupersetSchemaGenerator); } + Optional storeUpdateHandler = Optional.empty(); + Object passedStoreUpdateHandler = options.getExtraProperties().get(STORE_UPDATE_HANDLER); + if (passedStoreUpdateHandler instanceof StoreUpdateHandler) { + storeUpdateHandler = Optional.of((StoreUpdateHandler) passedStoreUpdateHandler); + } Map d2Clients = options.getD2Clients(); VeniceControllerContext ctx = new VeniceControllerContext.Builder().setPropertiesList(propertiesList) .setMetricsRepository(metricsRepository) @@ -431,6 +438,7 @@ static StatefulServiceProvider generateService(VeniceCo .setRouterClientConfig(consumerClientConfig.orElse(null)) .setExternalSupersetSchemaGenerator(supersetSchemaGenerator.orElse(null)) .setAccessController(options.getDynamicAccessController()) + .setStoreUpdateHandler(storeUpdateHandler.orElse(null)) .build(); VeniceController veniceController = new VeniceController(ctx); return new VeniceControllerWrapper( diff --git a/services/venice-controller/src/main/java/com/linkedin/venice/controller/StoreUpdateHandler.java b/services/venice-controller/src/main/java/com/linkedin/venice/controller/StoreUpdateHandler.java new file mode 100644 index 00000000000..e31c7a2f3fc --- /dev/null +++ b/services/venice-controller/src/main/java/com/linkedin/venice/controller/StoreUpdateHandler.java @@ -0,0 +1,37 @@ +package com.linkedin.venice.controller; + +import com.linkedin.venice.meta.Store; +import java.util.Set; + + +/** + * Handles a successful store update consumed by a parent controller. + */ +@FunctionalInterface +public interface StoreUpdateHandler { + StoreUpdateHandler NO_OP = new StoreUpdateHandler() { + @Override + public void handleStoreUpdate(String clusterName, Store store, Set updatedConfigs) { + } + + @Override + public boolean isNoOp() { + return true; + } + }; + + /** + * @param clusterName the cluster containing the updated store + * @param store a read-only snapshot of the final store state + * @param updatedConfigs the immutable set of config keys copied from the durable UPDATE_STORE message; the set is + * stable across retries of the same admin operation + */ + void handleStoreUpdate(String clusterName, Store store, Set updatedConfigs); + + /** + * @return whether callback-specific work, including fetching the final store snapshot, should be skipped + */ + default boolean isNoOp() { + return false; + } +} diff --git a/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceController.java b/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceController.java index 68676953c58..c5292f85c74 100644 --- a/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceController.java +++ b/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceController.java @@ -156,6 +156,7 @@ public static List> getMetricEntity private final Optional> versionLifecycleEventListeners; private final Optional> valueSchemaCreatedListeners; private final Optional externalETLService; + private final StoreUpdateHandler storeUpdateHandler; /** * Allocates a new {@code VeniceController} object. @@ -226,6 +227,7 @@ public VeniceController(VeniceControllerContext ctx) { this.versionLifecycleEventListeners = Optional.ofNullable(ctx.getVersionLifecycleEventListeners()); this.valueSchemaCreatedListeners = Optional.ofNullable(ctx.getValueSchemaCreatedListeners()); this.externalETLService = Optional.ofNullable(ctx.getExternalETLService()); + this.storeUpdateHandler = ctx.getStoreUpdateHandler(); this.controllerService = createControllerService(); this.adminServer = createAdminServer(false); this.secureAdminServer = sslEnabled ? createAdminServer(true) : null; @@ -262,7 +264,8 @@ private VeniceControllerService createControllerService() { pubSubPositionTypeRegistry, versionLifecycleEventListeners, valueSchemaCreatedListeners, - externalETLService); + externalETLService, + storeUpdateHandler); Admin admin = veniceControllerService.getVeniceHelixAdmin(); if (multiClusterConfigs.isParent() && !(admin instanceof VeniceParentHelixAdmin)) { throw new VeniceException( diff --git a/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceControllerContext.java b/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceControllerContext.java index 7a6a0107b98..beecfc9085f 100644 --- a/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceControllerContext.java +++ b/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceControllerContext.java @@ -40,6 +40,7 @@ public class VeniceControllerContext { private List versionLifecycleEventListeners; private List valueSchemaCreatedListeners; private ExternalETLService externalETLService; + private StoreUpdateHandler storeUpdateHandler; public List getPropertiesList() { return propertiesList; @@ -97,6 +98,10 @@ public ExternalETLService getExternalETLService() { return externalETLService; } + public StoreUpdateHandler getStoreUpdateHandler() { + return storeUpdateHandler; + } + public VeniceControllerContext(Builder builder) { this.propertiesList = builder.propertiesList; this.metricsRepository = builder.metricsRepository; @@ -112,6 +117,8 @@ public VeniceControllerContext(Builder builder) { this.versionLifecycleEventListeners = builder.versionLifecycleEventListeners; this.valueSchemaCreatedListeners = builder.valueSchemaCreatedListeners; this.externalETLService = builder.externalETLService; + this.storeUpdateHandler = + builder.storeUpdateHandler == null ? StoreUpdateHandler.NO_OP : builder.storeUpdateHandler; } public static class Builder { @@ -132,6 +139,7 @@ public static class Builder { private List versionLifecycleEventListeners; private List valueSchemaCreatedListeners; private ExternalETLService externalETLService; + private StoreUpdateHandler storeUpdateHandler; public Builder setPropertiesList(List propertiesList) { this.propertiesList = propertiesList; @@ -207,6 +215,11 @@ public Builder setExternalETLService(ExternalETLService externalETLService) { return this; } + public Builder setStoreUpdateHandler(StoreUpdateHandler storeUpdateHandler) { + this.storeUpdateHandler = storeUpdateHandler; + return this; + } + private void addDefaultValues() { if (metricsRepository == null && !isMetricsRepositorySet) { @@ -223,6 +236,9 @@ private void addDefaultValues() { if (serviceDiscoveryAnnouncers == null && !isServiceDiscoveryAnnouncerSet) { serviceDiscoveryAnnouncers = Collections.emptyList(); } + if (storeUpdateHandler == null) { + storeUpdateHandler = StoreUpdateHandler.NO_OP; + } } public VeniceControllerContext build() { diff --git a/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceControllerService.java b/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceControllerService.java index 1890c36c126..d75a8006d09 100644 --- a/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceControllerService.java +++ b/services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceControllerService.java @@ -73,6 +73,46 @@ public VeniceControllerService( Optional> versionLifecycleEventListeners, Optional> valueSchemaCreatedListeners, Optional externalETLService) { + this( + multiClusterConfigs, + metricsRepository, + sslEnabled, + sslConfig, + accessController, + authorizerService, + d2Client, + d2Clients, + routerClientConfig, + icProvider, + externalSupersetSchemaGenerator, + pubSubTopicRepository, + pubSubClientsFactory, + pubSubPositionTypeRegistry, + versionLifecycleEventListeners, + valueSchemaCreatedListeners, + externalETLService, + StoreUpdateHandler.NO_OP); + } + + public VeniceControllerService( + VeniceControllerMultiClusterConfig multiClusterConfigs, + MetricsRepository metricsRepository, + boolean sslEnabled, + Optional sslConfig, + Optional accessController, + Optional authorizerService, + D2Client d2Client, + Map d2Clients, + Optional routerClientConfig, + Optional icProvider, + Optional externalSupersetSchemaGenerator, + PubSubTopicRepository pubSubTopicRepository, + PubSubClientsFactory pubSubClientsFactory, + PubSubPositionTypeRegistry pubSubPositionTypeRegistry, + Optional> versionLifecycleEventListeners, + Optional> valueSchemaCreatedListeners, + Optional externalETLService, + StoreUpdateHandler storeUpdateHandler) { this.multiClusterConfigs = multiClusterConfigs; DelegatingClusterLeaderInitializationRoutine initRoutineForPushJobDetailsSystemStore = @@ -207,7 +247,8 @@ public VeniceControllerService( metricsRepository, pubSubClientsFactory.getConsumerAdapterFactory(), pubSubTopicRepository, - pubSubMessageDeserializer); + pubSubMessageDeserializer, + storeUpdateHandler); this.consumerServicesByClusters.put(cluster, adminConsumerService); this.admin.setAdminConsumerService(cluster, adminConsumerService); diff --git a/services/venice-controller/src/main/java/com/linkedin/venice/controller/kafka/consumer/AdminConsumerService.java b/services/venice-controller/src/main/java/com/linkedin/venice/controller/kafka/consumer/AdminConsumerService.java index ef85a827c03..812f1b2586f 100644 --- a/services/venice-controller/src/main/java/com/linkedin/venice/controller/kafka/consumer/AdminConsumerService.java +++ b/services/venice-controller/src/main/java/com/linkedin/venice/controller/kafka/consumer/AdminConsumerService.java @@ -3,6 +3,7 @@ import static com.linkedin.venice.pubsub.PubSubUtil.getPubSubPositionWireFormat; import com.linkedin.venice.annotation.VisibleForTesting; +import com.linkedin.venice.controller.StoreUpdateHandler; import com.linkedin.venice.controller.VeniceControllerClusterConfig; import com.linkedin.venice.controller.VeniceHelixAdmin; import com.linkedin.venice.controller.ZkAdminTopicMetadataAccessor; @@ -50,6 +51,7 @@ public class AdminConsumerService extends AbstractVeniceService { private final PubSubMessageDeserializer pubSubMessageDeserializer; private final LogContext logContext; private final PubSubPositionDeserializer pubSubPositionDeserializer; + private final StoreUpdateHandler storeUpdateHandler; public AdminConsumerService( VeniceHelixAdmin admin, @@ -58,6 +60,24 @@ public AdminConsumerService( PubSubConsumerAdapterFactory consumerFactory, PubSubTopicRepository pubSubTopicRepository, PubSubMessageDeserializer pubSubMessageDeserializer) { + this( + admin, + config, + metricsRepository, + consumerFactory, + pubSubTopicRepository, + pubSubMessageDeserializer, + StoreUpdateHandler.NO_OP); + } + + public AdminConsumerService( + VeniceHelixAdmin admin, + VeniceControllerClusterConfig config, + MetricsRepository metricsRepository, + PubSubConsumerAdapterFactory consumerFactory, + PubSubTopicRepository pubSubTopicRepository, + PubSubMessageDeserializer pubSubMessageDeserializer, + StoreUpdateHandler storeUpdateHandler) { this.config = config; this.logContext = config.getLogContext(); this.admin = admin; @@ -77,6 +97,7 @@ public AdminConsumerService( this.consumerFactory = consumerFactory; this.pubSubPositionDeserializer = config.getPubSubPositionDeserializer(); this.threadFactory = new DaemonThreadFactory("AdminConsumerService-" + config.getClusterName(), logContext); + this.storeUpdateHandler = storeUpdateHandler; } @Override @@ -118,7 +139,8 @@ private AdminConsumptionTask getAdminConsumptionTaskForCluster(String clusterNam config.getAdminConsumptionCycleTimeoutMs(), config.getAdminConsumptionMaxWorkerThreadPoolSize(), pubSubTopicRepository, - config.getRegionName()); + config.getRegionName(), + storeUpdateHandler); } /** diff --git a/services/venice-controller/src/main/java/com/linkedin/venice/controller/kafka/consumer/AdminConsumptionTask.java b/services/venice-controller/src/main/java/com/linkedin/venice/controller/kafka/consumer/AdminConsumptionTask.java index 918d0e67714..5460ed1e326 100644 --- a/services/venice-controller/src/main/java/com/linkedin/venice/controller/kafka/consumer/AdminConsumptionTask.java +++ b/services/venice-controller/src/main/java/com/linkedin/venice/controller/kafka/consumer/AdminConsumptionTask.java @@ -5,6 +5,7 @@ import com.linkedin.venice.common.VeniceSystemStoreType; import com.linkedin.venice.controller.AdminTopicMetadataAccessor; import com.linkedin.venice.controller.ExecutionIdAccessor; +import com.linkedin.venice.controller.StoreUpdateHandler; import com.linkedin.venice.controller.VeniceHelixAdmin; import com.linkedin.venice.controller.kafka.AdminTopicUtils; import com.linkedin.venice.controller.kafka.protocol.admin.AdminOperation; @@ -262,6 +263,7 @@ public ExecutorService getExecutorService() { * The local region name of the controller. */ private final String regionName; + private final StoreUpdateHandler storeUpdateHandler; public AdminConsumptionTask( String clusterName, @@ -279,6 +281,42 @@ public AdminConsumptionTask( int maxWorkerThreadPoolSize, PubSubTopicRepository pubSubTopicRepository, String regionName) { + this( + clusterName, + consumer, + remoteConsumptionEnabled, + remoteKafkaServerUrl, + admin, + adminTopicMetadataAccessor, + executionIdAccessor, + isParentController, + stats, + adminTopicReplicationFactor, + minInSyncReplicas, + processingCycleTimeoutInMs, + maxWorkerThreadPoolSize, + pubSubTopicRepository, + regionName, + StoreUpdateHandler.NO_OP); + } + + public AdminConsumptionTask( + String clusterName, + PubSubConsumerAdapter consumer, + boolean remoteConsumptionEnabled, + Optional remoteKafkaServerUrl, + VeniceHelixAdmin admin, + AdminTopicMetadataAccessor adminTopicMetadataAccessor, + ExecutionIdAccessor executionIdAccessor, + boolean isParentController, + AdminConsumptionStats stats, + int adminTopicReplicationFactor, + Optional minInSyncReplicas, + long processingCycleTimeoutInMs, + int maxWorkerThreadPoolSize, + PubSubTopicRepository pubSubTopicRepository, + String regionName, + StoreUpdateHandler storeUpdateHandler) { this.clusterName = clusterName; this.pubSubAdminTopic = pubSubTopicRepository.getTopic(AdminTopicUtils.getTopicNameFromClusterName(clusterName)); this.adminTopicPartition = new PubSubTopicPartitionImpl(pubSubAdminTopic, AdminTopicUtils.ADMIN_TOPIC_PARTITION_ID); @@ -312,6 +350,7 @@ public AdminConsumptionTask( new DaemonThreadFactory(String.format("Venice-Admin-Execution-Task-%s", clusterName), admin.getLogContext())); this.undelegatedRecords = new LinkedList<>(); this.regionName = regionName; + this.storeUpdateHandler = storeUpdateHandler; this.storeRetryCountMap = new ConcurrentHashMap<>(); if (remoteConsumptionEnabled) { @@ -569,7 +608,8 @@ private void executeMessagesAndCollectResults() throws InterruptedException { isParentController, stats, regionName, - inflightThreadsByStore); + inflightThreadsByStore, + storeUpdateHandler); // Check if there is previously created scheduled task still occupying one thread from the pool. if (storesWithScheduledTask.add(storeName)) { // Log the store name and the position of the task being added into the task list diff --git a/services/venice-controller/src/main/java/com/linkedin/venice/controller/kafka/consumer/AdminExecutionTask.java b/services/venice-controller/src/main/java/com/linkedin/venice/controller/kafka/consumer/AdminExecutionTask.java index 6a40dedec59..8a55fca4d39 100644 --- a/services/venice-controller/src/main/java/com/linkedin/venice/controller/kafka/consumer/AdminExecutionTask.java +++ b/services/venice-controller/src/main/java/com/linkedin/venice/controller/kafka/consumer/AdminExecutionTask.java @@ -7,6 +7,7 @@ import com.linkedin.venice.common.VeniceSystemStoreUtils; import com.linkedin.venice.compression.CompressionStrategy; import com.linkedin.venice.controller.ExecutionIdAccessor; +import com.linkedin.venice.controller.StoreUpdateHandler; import com.linkedin.venice.controller.VeniceHelixAdmin; import com.linkedin.venice.controller.kafka.protocol.admin.AbortMigration; import com.linkedin.venice.controller.kafka.protocol.admin.AddVersion; @@ -53,6 +54,7 @@ import com.linkedin.venice.meta.IngestionPauseMode; import com.linkedin.venice.meta.LifecycleHooksRecord; import com.linkedin.venice.meta.LifecycleHooksRecordImpl; +import com.linkedin.venice.meta.ReadOnlyStore; import com.linkedin.venice.meta.StorageMode; import com.linkedin.venice.meta.Store; import com.linkedin.venice.meta.VeniceETLStrategy; @@ -63,6 +65,7 @@ import java.util.ArrayList; import java.util.Collections; import java.util.HashSet; +import java.util.LinkedHashSet; import java.util.List; import java.util.Map; import java.util.Queue; @@ -95,6 +98,7 @@ public class AdminExecutionTask implements Callable { private final long lastPersistedExecutionId; private final ConcurrentHashMap inflightThreadsByStore; + private final StoreUpdateHandler storeUpdateHandler; AdminExecutionTask( Logger LOGGER, @@ -109,6 +113,36 @@ public class AdminExecutionTask implements Callable { AdminConsumptionStats stats, String regionName, ConcurrentHashMap inflightThreadsByStore) { + this( + LOGGER, + clusterName, + storeName, + lastSucceededExecutionIdMap, + lastPersistedExecutionId, + internalTopic, + admin, + executionIdAccessor, + isParentController, + stats, + regionName, + inflightThreadsByStore, + StoreUpdateHandler.NO_OP); + } + + AdminExecutionTask( + Logger LOGGER, + String clusterName, + String storeName, + ConcurrentHashMap lastSucceededExecutionIdMap, + long lastPersistedExecutionId, + Queue internalTopic, + VeniceHelixAdmin admin, + ExecutionIdAccessor executionIdAccessor, + boolean isParentController, + AdminConsumptionStats stats, + String regionName, + ConcurrentHashMap inflightThreadsByStore, + StoreUpdateHandler storeUpdateHandler) { this.LOGGER = LOGGER; this.clusterName = clusterName; this.storeName = storeName; @@ -121,6 +155,7 @@ public class AdminExecutionTask implements Callable { this.stats = stats; this.regionName = regionName; this.inflightThreadsByStore = inflightThreadsByStore; + this.storeUpdateHandler = storeUpdateHandler; } @Override @@ -250,6 +285,8 @@ private void processMessage(AdminOperation adminOperation) { lastSucceededExecutionId); return; } + boolean storeUpdated = false; + Set updatedConfigs = Collections.emptySet(); try { switch (AdminMessageType.valueOf(adminOperation)) { case STORE_CREATION: @@ -286,7 +323,10 @@ private void processMessage(AdminOperation adminOperation) { handleSetStorePartitionCount((SetStorePartitionCount) adminOperation.payloadUnion); break; case UPDATE_STORE: - handleSetStore((UpdateStore) adminOperation.payloadUnion); + UpdateStore updateStore = (UpdateStore) adminOperation.payloadUnion; + updatedConfigs = extractUpdatedConfigs(updateStore); + handleSetStore(updateStore); + storeUpdated = true; break; case DELETE_STORE: handleDeleteStore((DeleteStore) adminOperation.payloadUnion); @@ -354,10 +394,30 @@ private void processMessage(AdminOperation adminOperation) { AdminMessageType.valueOf(adminOperation), e.getMessage()); } + if (storeUpdated && isParentController && !storeUpdateHandler.isNoOp()) { + Store finalStore = admin.getStore(clusterName, storeName).cloneStore(); + // Invoke before advancing checkpoints so callback failures leave the admin operation eligible for retry. + storeUpdateHandler.handleStoreUpdate(clusterName, new ReadOnlyStore(finalStore), updatedConfigs); + } executionIdAccessor.updateLastSucceededExecutionIdMap(clusterName, storeName, adminOperation.executionId); lastSucceededExecutionIdMap.put(storeName, adminOperation.executionId); } + /** + * Copies the config keys from the durable UPDATE_STORE message so callbacks cannot mutate the Avro collection. + * The durable message is retried unchanged, so the returned set is deterministic and stable across retry attempts. + */ + private Set extractUpdatedConfigs(UpdateStore updateStore) { + if (updateStore.updatedConfigsList == null || updateStore.updatedConfigsList.isEmpty()) { + return Collections.emptySet(); + } + Set updatedConfigs = new LinkedHashSet<>(); + for (CharSequence config: updateStore.updatedConfigsList) { + updatedConfigs.add(config.toString()); + } + return Collections.unmodifiableSet(updatedConfigs); + } + private void handleStoreCreation(StoreCreation message) { String clusterName = message.clusterName.toString(); String storeName = message.storeName.toString(); diff --git a/services/venice-controller/src/test/java/com/linkedin/venice/controller/VeniceControllerContextTest.java b/services/venice-controller/src/test/java/com/linkedin/venice/controller/VeniceControllerContextTest.java index 8e6272eb6b9..85c0e79d351 100644 --- a/services/venice-controller/src/test/java/com/linkedin/venice/controller/VeniceControllerContextTest.java +++ b/services/venice-controller/src/test/java/com/linkedin/venice/controller/VeniceControllerContextTest.java @@ -4,6 +4,7 @@ import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertNotNull; import static org.testng.Assert.assertNull; +import static org.testng.Assert.assertSame; import com.linkedin.d2.balancer.D2Client; import com.linkedin.venice.acl.DynamicAccessController; @@ -24,6 +25,7 @@ public void testVeniceServerContextCanSetDefaults() { assertNotNull(veniceControllerContext.getMetricsRepository()); assertNotNull(veniceControllerContext.getServiceDiscoveryAnnouncers()); assertEquals(veniceControllerContext.getServiceDiscoveryAnnouncers(), Collections.emptyList()); + assertSame(veniceControllerContext.getStoreUpdateHandler(), StoreUpdateHandler.NO_OP); } @Test @@ -36,6 +38,7 @@ public void testVeniceServerContextCanSetValues() { ClientConfig routerClientConfig = mock(ClientConfig.class); ICProvider icProvider = mock(ICProvider.class); SupersetSchemaGenerator externalSupersetSchemaGenerator = mock(SupersetSchemaGenerator.class); + StoreUpdateHandler storeUpdateHandler = mock(StoreUpdateHandler.class); VeniceControllerContext veniceControllerContext = new VeniceControllerContext.Builder().setPropertiesList(propertiesList) @@ -45,6 +48,7 @@ public void testVeniceServerContextCanSetValues() { .setRouterClientConfig(routerClientConfig) .setIcProvider(icProvider) .setExternalSupersetSchemaGenerator(externalSupersetSchemaGenerator) + .setStoreUpdateHandler(storeUpdateHandler) .setMetricsRepository(null) .setServiceDiscoveryAnnouncers(null) .build(); @@ -57,5 +61,6 @@ public void testVeniceServerContextCanSetValues() { assertEquals(veniceControllerContext.getD2Client(), d2Client); assertEquals(veniceControllerContext.getRouterClientConfig(), routerClientConfig); assertEquals(veniceControllerContext.getIcProvider(), icProvider); + assertSame(veniceControllerContext.getStoreUpdateHandler(), storeUpdateHandler); } } diff --git a/services/venice-controller/src/test/java/com/linkedin/venice/controller/kafka/consumer/AdminExecutionTaskTest.java b/services/venice-controller/src/test/java/com/linkedin/venice/controller/kafka/consumer/AdminExecutionTaskTest.java index fb2870b36c4..f0d286a74e0 100644 --- a/services/venice-controller/src/test/java/com/linkedin/venice/controller/kafka/consumer/AdminExecutionTaskTest.java +++ b/services/venice-controller/src/test/java/com/linkedin/venice/controller/kafka/consumer/AdminExecutionTaskTest.java @@ -1,11 +1,16 @@ package com.linkedin.venice.controller.kafka.consumer; +import static com.linkedin.venice.controllerapi.ControllerApiConstants.READ_QUOTA_IN_CU; +import static com.linkedin.venice.controllerapi.ControllerApiConstants.THROUGHPUT_QUOTA_IN_BYTES; +import static com.linkedin.venice.controllerapi.ControllerApiConstants.THROUGHPUT_QUOTA_IN_RECORDS; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyInt; import static org.mockito.ArgumentMatchers.anyLong; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.atLeastOnce; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.inOrder; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.never; import static org.mockito.Mockito.verify; @@ -13,9 +18,11 @@ import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertNotNull; import static org.testng.Assert.assertNull; +import static org.testng.Assert.assertThrows; import static org.testng.Assert.assertTrue; import com.linkedin.venice.controller.ExecutionIdAccessor; +import com.linkedin.venice.controller.StoreUpdateHandler; import com.linkedin.venice.controller.VeniceHelixAdmin; import com.linkedin.venice.controller.kafka.protocol.admin.AddVersion; import com.linkedin.venice.controller.kafka.protocol.admin.AdminOperation; @@ -25,12 +32,17 @@ import com.linkedin.venice.controller.kafka.protocol.enums.SchemaType; import com.linkedin.venice.controller.stats.AdminConsumptionStats; import com.linkedin.venice.controllerapi.UpdateStoreQueryParams; +import com.linkedin.venice.exceptions.VeniceUnsupportedOperationException; +import com.linkedin.venice.meta.Store; import com.linkedin.venice.meta.Version; import com.linkedin.venice.pubsub.api.PubSubPosition; import com.linkedin.venice.pubsub.mock.InMemoryPubSubPosition; import java.util.Arrays; +import java.util.Collections; +import java.util.List; import java.util.Optional; import java.util.Queue; +import java.util.Set; import java.util.concurrent.CancellationException; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentLinkedQueue; @@ -40,6 +52,7 @@ import java.util.concurrent.atomic.AtomicReference; import org.apache.logging.log4j.Logger; import org.mockito.ArgumentCaptor; +import org.mockito.InOrder; import org.testng.annotations.BeforeMethod; import org.testng.annotations.Test; @@ -611,6 +624,172 @@ public void testHandleSetStore_ThroughputQuota_PropagatedToParams() { "updateStore must be called with throughput quota values propagated from the UpdateStore message"); } + @Test + public void testParentStoreUpdateHandlerReceivesFinalStoreBeforeCheckpoint() { + when(mockAdmin.isLeaderControllerFor(clusterName)).thenReturn(true); + Store mutableStore = mock(Store.class); + Store finalStoreSnapshot = mock(Store.class); + when(mockAdmin.getStore(clusterName, storeName)).thenReturn(mutableStore); + when(mutableStore.cloneStore()).thenReturn(finalStoreSnapshot); + when(finalStoreSnapshot.getName()).thenReturn(storeName); + when(finalStoreSnapshot.getOwner()).thenReturn("final-owner"); + StoreUpdateHandler storeUpdateHandler = mock(StoreUpdateHandler.class); + doAnswer(invocation -> { + assertNull(lastSucceededExecutionIdMap.get(storeName)); + return null; + }).when(storeUpdateHandler).handleStoreUpdate(eq(clusterName), any(Store.class), any()); + + Queue queue = new ConcurrentLinkedQueue<>(); + queue.add(createUpdateStoreWrapper(1L, false)); + + AdminExecutionTask task = new AdminExecutionTask( + mockLogger, + clusterName, + storeName, + lastSucceededExecutionIdMap, + lastPersistedExecutionId, + queue, + mockAdmin, + mockExecutionIdAccessor, + true, + mockStats, + regionName, + inflightThreadsByStore, + storeUpdateHandler); + + task.call(); + + ArgumentCaptor storeCaptor = ArgumentCaptor.forClass(Store.class); + ArgumentCaptor> updatedConfigsCaptor = ArgumentCaptor.forClass(Set.class); + InOrder inOrder = inOrder(mockAdmin, mutableStore, storeUpdateHandler, mockExecutionIdAccessor); + inOrder.verify(mockAdmin).updateStore(eq(clusterName), eq(storeName), any(UpdateStoreQueryParams.class)); + inOrder.verify(mockAdmin).getStore(clusterName, storeName); + inOrder.verify(mutableStore).cloneStore(); + inOrder.verify(storeUpdateHandler) + .handleStoreUpdate(eq(clusterName), storeCaptor.capture(), updatedConfigsCaptor.capture()); + inOrder.verify(mockExecutionIdAccessor).updateLastSucceededExecutionIdMap(clusterName, storeName, 1L); + + Store callbackStore = storeCaptor.getValue(); + assertEquals(callbackStore.getName(), storeName); + assertEquals(callbackStore.getOwner(), "final-owner"); + assertThrows(UnsupportedOperationException.class, () -> callbackStore.setOwner("new-owner")); + Set updatedConfigs = updatedConfigsCaptor.getValue(); + assertEquals(updatedConfigs, Collections.singleton(READ_QUOTA_IN_CU)); + assertThrows(UnsupportedOperationException.class, () -> updatedConfigs.add("another-config")); + assertEquals(lastSucceededExecutionIdMap.get(storeName), Long.valueOf(1L)); + } + + @Test + public void testParentStoreUpdateWithDefaultNoOpHandlerDoesNotFetchStore() { + when(mockAdmin.isLeaderControllerFor(clusterName)).thenReturn(true); + Queue queue = new ConcurrentLinkedQueue<>(); + queue.add(createUpdateStoreWrapper(1L, false)); + + AdminExecutionTask task = new AdminExecutionTask( + mockLogger, + clusterName, + storeName, + lastSucceededExecutionIdMap, + lastPersistedExecutionId, + queue, + mockAdmin, + mockExecutionIdAccessor, + true, + mockStats, + regionName, + inflightThreadsByStore); + + task.call(); + + verify(mockAdmin, never()).getStore(anyString(), anyString()); + verify(mockExecutionIdAccessor).updateLastSucceededExecutionIdMap(clusterName, storeName, 1L); + assertEquals(lastSucceededExecutionIdMap.get(storeName), Long.valueOf(1L)); + assertTrue(queue.isEmpty()); + } + + @Test + public void testChildControllerDoesNotInvokeStoreUpdateHandlerOrFetchStore() { + when(mockAdmin.isLeaderControllerFor(clusterName)).thenReturn(true); + StoreUpdateHandler storeUpdateHandler = mock(StoreUpdateHandler.class); + + Queue queue = new ConcurrentLinkedQueue<>(); + queue.add(createUpdateStoreWrapper(1L, false)); + + AdminExecutionTask task = new AdminExecutionTask( + mockLogger, + clusterName, + storeName, + lastSucceededExecutionIdMap, + lastPersistedExecutionId, + queue, + mockAdmin, + mockExecutionIdAccessor, + false, + mockStats, + regionName, + inflightThreadsByStore, + storeUpdateHandler); + + task.call(); + + verify(storeUpdateHandler, never()).handleStoreUpdate(anyString(), any(Store.class), any()); + verify(mockAdmin, never()).getStore(anyString(), anyString()); + verify(mockExecutionIdAccessor).updateLastSucceededExecutionIdMap(clusterName, storeName, 1L); + } + + @Test + public void testStoreUpdateHandlerFailureLeavesExecutionIdUnadvanced() { + when(mockAdmin.isLeaderControllerFor(clusterName)).thenReturn(true); + Store mutableStore = mock(Store.class); + Store finalStoreSnapshot = mock(Store.class); + when(mockAdmin.getStore(clusterName, storeName)).thenReturn(mutableStore); + when(mutableStore.cloneStore()).thenReturn(finalStoreSnapshot); + StoreUpdateHandler storeUpdateHandler = mock(StoreUpdateHandler.class); + AtomicInteger handlerInvocationCount = new AtomicInteger(); + List> receivedUpdatedConfigs = new java.util.concurrent.CopyOnWriteArrayList<>(); + doAnswer(invocation -> { + receivedUpdatedConfigs.add(invocation.getArgument(2)); + if (handlerInvocationCount.incrementAndGet() == 1) { + throw new VeniceUnsupportedOperationException("store update callback"); + } + return null; + }).when(storeUpdateHandler).handleStoreUpdate(eq(clusterName), any(Store.class), any()); + + Queue queue = new ConcurrentLinkedQueue<>(); + queue.add(createUpdateStoreWrapper(1L, false)); + + AdminExecutionTask task = new AdminExecutionTask( + mockLogger, + clusterName, + storeName, + lastSucceededExecutionIdMap, + lastPersistedExecutionId, + queue, + mockAdmin, + mockExecutionIdAccessor, + true, + mockStats, + regionName, + inflightThreadsByStore, + storeUpdateHandler); + + assertThrows(VeniceUnsupportedOperationException.class, task::call); + + verify(mockExecutionIdAccessor, never()).updateLastSucceededExecutionIdMap(anyString(), anyString(), anyLong()); + assertNull(lastSucceededExecutionIdMap.get(storeName)); + assertEquals(queue.size(), 1); + + task.call(); + + assertEquals(handlerInvocationCount.get(), 2); + assertEquals( + receivedUpdatedConfigs, + Arrays.asList(Collections.singleton(READ_QUOTA_IN_CU), Collections.singleton(READ_QUOTA_IN_CU))); + assertThrows(UnsupportedOperationException.class, () -> receivedUpdatedConfigs.get(0).add("another-config")); + verify(mockExecutionIdAccessor).updateLastSucceededExecutionIdMap(clusterName, storeName, 1L); + assertEquals(lastSucceededExecutionIdMap.get(storeName), Long.valueOf(1L)); + } + private AdminOperationWrapper createUpdateStoreWrapper(long executionId, boolean targetRegionPromoted) { return createUpdateStoreWrapper(executionId, targetRegionPromoted, -1L, -1L); } @@ -668,6 +847,13 @@ private AdminOperationWrapper createUpdateStoreWrapper( updateStore.flinkVeniceViewsEnabled = false; updateStore.unusedSchemaDeletionEnabled = false; updateStore.updatedConfigsList = new java.util.ArrayList<>(); + updateStore.updatedConfigsList.add(READ_QUOTA_IN_CU); + if (throughputQuotaInBytes >= 0) { + updateStore.updatedConfigsList.add(THROUGHPUT_QUOTA_IN_BYTES); + } + if (throughputQuotaInRecords >= 0) { + updateStore.updatedConfigsList.add(THROUGHPUT_QUOTA_IN_RECORDS); + } updateStore.replicateAllConfigs = true; updateStore.storeLifecycleHooks = new java.util.ArrayList<>(); updateStore.keyUrnCompressionEnabled = false;