diff --git a/changelog/unreleased/SOLR-13136.yml b/changelog/unreleased/SOLR-13136.yml new file mode 100644 index 000000000000..74d88dc8e7b2 --- /dev/null +++ b/changelog/unreleased/SOLR-13136.yml @@ -0,0 +1,7 @@ +title: Create new shards in CONSTRUCTION state so queries don't fail during shard creation +type: fixed +authors: + - name: Nick Shanin +links: + - name: SOLR-13136 + url: https://issues.apache.org/jira/browse/SOLR-13136 diff --git a/solr/core/src/java/org/apache/solr/cloud/api/collections/CreateShardCmd.java b/solr/core/src/java/org/apache/solr/cloud/api/collections/CreateShardCmd.java index 7c0e16336038..3c0a68a64eed 100644 --- a/solr/core/src/java/org/apache/solr/cloud/api/collections/CreateShardCmd.java +++ b/solr/core/src/java/org/apache/solr/cloud/api/collections/CreateShardCmd.java @@ -20,20 +20,37 @@ import static org.apache.solr.common.cloud.ZkStateReader.SHARD_ID_PROP; import static org.apache.solr.common.params.CollectionAdminParams.FOLLOW_ALIASES; import static org.apache.solr.common.params.CollectionParams.CollectionAction.CREATESHARD; +import static org.apache.solr.common.params.CommonAdminParams.TIMEOUT; import java.lang.invoke.MethodHandles; +import java.util.HashMap; +import java.util.List; import java.util.Map; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.stream.Collectors; +import org.apache.solr.cloud.ActiveReplicaWatcher; import org.apache.solr.cloud.DistributedClusterStateUpdater; +import org.apache.solr.cloud.Overseer; +import org.apache.solr.cloud.overseer.OverseerAction; +import org.apache.solr.common.SolrCloseableLatch; import org.apache.solr.common.SolrException; import org.apache.solr.common.cloud.ClusterState; import org.apache.solr.common.cloud.DocCollection; +import org.apache.solr.common.cloud.Replica; import org.apache.solr.common.cloud.ReplicaCount; +import org.apache.solr.common.cloud.Slice; import org.apache.solr.common.cloud.ZkNodeProps; +import org.apache.solr.common.cloud.ZkStateReader; import org.apache.solr.common.params.CollectionParams; import org.apache.solr.common.params.CommonAdminParams; +import org.apache.solr.common.params.CoreAdminParams; +import org.apache.solr.common.params.ModifiableSolrParams; import org.apache.solr.common.util.NamedList; import org.apache.solr.common.util.SimpleOrderedMap; import org.apache.solr.common.util.Utils; +import org.apache.solr.handler.component.ShardHandler; +import org.apache.zookeeper.KeeperException; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -67,6 +84,7 @@ public void call(AdminCmdContext adminCmdContext, ZkNodeProps message, NamedList } ClusterState clusterState = adminCmdContext.getClusterState(); DocCollection collection = clusterState.getCollection(collectionName); + boolean sliceAlreadyExists = collection.getSlice(sliceName) != null; ReplicaCount numReplicas = ReplicaCount.fromMessage(message, collection, 1); if (!numReplicas.hasLeaderReplica()) { @@ -115,6 +133,7 @@ public void call(AdminCmdContext adminCmdContext, ZkNodeProps message, NamedList CollectionHandlingUtils.addPropertyParams(message, addReplicasProps); final NamedList addResult = new NamedList<>(); + int timeout = message.getInt(TIMEOUT, 10 * 60); // 10 minutes try { new AddReplicaCmd(ccc) .addReplica( @@ -148,18 +167,162 @@ public void call(AdminCmdContext adminCmdContext, ZkNodeProps message, NamedList success.addAll(addResultSuccess); } }); - } catch (Assign.AssignmentException e) { - // clean up the slice that we created - new DeleteShardCmd(ccc) - .call( - adminCmdContext - .subRequestContext(CollectionParams.CollectionAction.DELETESHARD, async) - .withClusterState(clusterState), - new ZkNodeProps(COLLECTION_PROP, collectionName, SHARD_ID_PROP, sliceName), - results); + } catch (Exception e) { + if (!sliceAlreadyExists) { + // Don't leave a half-created shard stuck in CONSTRUCTION; remove it so the create can be + // retried cleanly. DeleteShardCmd needs a fresh cluster state here: the snapshot from + // waitForNewShard predates AddReplicaCmd, so it shows a slice with no replicas and the + // delete would leave the added replicas' cores behind. + try { + new DeleteShardCmd(ccc) + .call( + adminCmdContext + .subRequestContext(CollectionParams.CollectionAction.DELETESHARD, async) + .withClusterState(ccc.getZkStateReader().getClusterState()), + new ZkNodeProps(COLLECTION_PROP, collectionName, SHARD_ID_PROP, sliceName), + results); + } catch (Exception cleanupEx) { + log.warn("Failed to delete shard {} after failed create", sliceName, cleanupEx); + } + } throw e; } + if (!sliceAlreadyExists) { + // The new slice is in CONSTRUCTION state, so queries skip it until it is activated. Only + // activate it once its replicas are active and any buffered updates are applied. If that + // fails, delete the half-created slice instead of activating a shard that may not be + // ready, mirroring the AddReplica failure handling above, and report the failure so the + // create can be retried. + try { + waitForShardReplicasActive(collectionName, sliceName, numReplicas.total(), timeout); + applyBufferedUpdatesOnLeader(adminCmdContext, collectionName, sliceName); + } catch (Exception e) { + try { + new DeleteShardCmd(ccc) + .call( + adminCmdContext + .subRequestContext(CollectionParams.CollectionAction.DELETESHARD, async) + .withClusterState(ccc.getZkStateReader().getClusterState()), + new ZkNodeProps(COLLECTION_PROP, collectionName, SHARD_ID_PROP, sliceName), + results); + } catch (Exception cleanupEx) { + log.warn("Failed to delete shard {} after failed create", sliceName, cleanupEx); + } + throw e; + } + activateShard(collectionName, sliceName); + } + log.info("Finished create command on all shards for collection: {}", collectionName); } + + /** + * Cores of a CONSTRUCTION shard buffer their updates, which a split lifts for its sub-shards but + * nothing else would for a created shard, so ask the leader to apply them. + */ + private void applyBufferedUpdatesOnLeader( + AdminCmdContext adminCmdContext, String collectionName, String sliceName) { + Slice slice = + ccc.getZkStateReader().getClusterState().getCollection(collectionName).getSlice(sliceName); + Replica leader = slice == null ? null : slice.getLeader(); + if (leader == null) { + log.warn("No leader for new shard {} of collection {}", sliceName, collectionName); + return; + } + ModifiableSolrParams params = new ModifiableSolrParams(); + params.set( + CoreAdminParams.ACTION, CoreAdminParams.CoreAdminAction.REQUESTAPPLYUPDATES.toString()); + params.set(CoreAdminParams.NAME, leader.getCoreName()); + + ShardHandler shardHandler = ccc.newShardHandler(); + CollectionHandlingUtils.ShardRequestTracker tracker = + CollectionHandlingUtils.asyncRequestTracker(adminCmdContext, ccc); + tracker.sendShardRequest(leader, params, shardHandler); + + // The core refuses the request when it is not buffering, which is fine, so only log a failure. + NamedList applyResults = new NamedList<>(); + tracker.processResponses(applyResults, shardHandler, false, null); + if (applyResults.get("failure") != null) { + log.warn( + "Leader {} of new shard {} did not apply buffered updates: {}", + leader.getName(), + sliceName, + applyResults.get("failure")); + } + } + + /** Flip a newly created shard from CONSTRUCTION to ACTIVE so it becomes visible to queries. */ + private void activateShard(String collectionName, String sliceName) + throws KeeperException, InterruptedException { + Map activateProps = new HashMap<>(); + activateProps.put(Overseer.QUEUE_OPERATION, OverseerAction.UPDATESHARDSTATE.toLower()); + activateProps.put(COLLECTION_PROP, collectionName); + activateProps.put(sliceName, Slice.State.ACTIVE.toString()); + ZkNodeProps activateMsg = new ZkNodeProps(activateProps); + if (ccc.getDistributedClusterStateUpdater().isDistributedStateUpdate()) { + ccc.getDistributedClusterStateUpdater() + .doSingleStateUpdate( + DistributedClusterStateUpdater.MutatingCommand.SliceUpdateShardState, + activateMsg, + ccc.getSolrCloudManager(), + ccc.getZkStateReader()); + } else { + ccc.offerStateUpdate(activateMsg); + } + } + + /** Wait until a shard has the expected number of replicas and all of them are ACTIVE. */ + private void waitForShardReplicasActive( + String collectionName, String sliceName, int expectedReplicas, int timeout) + throws InterruptedException { + ZkStateReader zkStateReader = ccc.getZkStateReader(); + try { + zkStateReader.waitForState( + collectionName, + timeout, + TimeUnit.SECONDS, + collection -> { + Slice newSlice = collection == null ? null : collection.getSlice(sliceName); + return newSlice != null && newSlice.getReplicas().size() >= expectedReplicas; + }); + } catch (TimeoutException e) { + throw new SolrException( + SolrException.ErrorCode.SERVER_ERROR, + "Timeout waiting " + timeout + " seconds for the replicas of shard " + sliceName, + e); + } + Slice slice = zkStateReader.getClusterState().getCollection(collectionName).getSlice(sliceName); + if (slice == null) { + throw new SolrException( + SolrException.ErrorCode.SERVER_ERROR, + "Newly created shard " + sliceName + " not visible in cluster state"); + } + List coreNames = + slice.getReplicas().stream().map(Replica::getCoreName).collect(Collectors.toList()); + if (coreNames.isEmpty()) { + log.warn( + "No replicas found for new shard {} of collection {}; activating without waiting", + sliceName, + collectionName); + return; + } + SolrCloseableLatch latch = + new SolrCloseableLatch(coreNames.size(), ccc.getCloseableToLatchOn()); + ActiveReplicaWatcher watcher = new ActiveReplicaWatcher(collectionName, null, coreNames, latch); + try { + zkStateReader.registerCollectionStateWatcher(collectionName, watcher); + if (!latch.await(timeout, TimeUnit.SECONDS)) { + throw new SolrException( + SolrException.ErrorCode.SERVER_ERROR, + "Timeout waiting " + + timeout + + " seconds for replicas of shard " + + sliceName + + " to become active."); + } + } finally { + zkStateReader.removeCollectionStateWatcher(collectionName, watcher); + } + } } diff --git a/solr/core/src/java/org/apache/solr/cloud/overseer/CollectionMutator.java b/solr/core/src/java/org/apache/solr/cloud/overseer/CollectionMutator.java index ba8f563ca3c0..3d3be8e0a87e 100644 --- a/solr/core/src/java/org/apache/solr/cloud/overseer/CollectionMutator.java +++ b/solr/core/src/java/org/apache/solr/cloud/overseer/CollectionMutator.java @@ -67,6 +67,10 @@ public ZkWriteCommand createShard(final ClusterState clusterState, ZkNodeProps m Map sliceProps = new HashMap<>(); String shardRange = message.getStr(ZkStateReader.SHARD_RANGE_PROP); String shardState = message.getStr(ZkStateReader.SHARD_STATE_PROP); + if (shardState == null) { + // A brand-new slice has no replicas yet; keep it out of query routing until they are added + shardState = Slice.State.CONSTRUCTION.toString(); + } String shardParent = message.getStr(ZkStateReader.SHARD_PARENT_PROP); String shardParentZkSession = message.getStr("shard_parent_zk_session"); String shardParentNode = message.getStr("shard_parent_node"); diff --git a/solr/core/src/java/org/apache/solr/handler/component/ShardResponse.java b/solr/core/src/java/org/apache/solr/handler/component/ShardResponse.java index 289a6e941829..fca9292599ea 100644 --- a/solr/core/src/java/org/apache/solr/handler/component/ShardResponse.java +++ b/solr/core/src/java/org/apache/solr/handler/component/ShardResponse.java @@ -78,7 +78,7 @@ void setShard(String shard) { this.shard = shard; } - void setException(Throwable exception) { + public void setException(Throwable exception) { this.exception = exception; } diff --git a/solr/core/src/test/org/apache/solr/cloud/CreateShardConstructionTest.java b/solr/core/src/test/org/apache/solr/cloud/CreateShardConstructionTest.java new file mode 100644 index 000000000000..207a23af9def --- /dev/null +++ b/solr/core/src/test/org/apache/solr/cloud/CreateShardConstructionTest.java @@ -0,0 +1,132 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.solr.cloud; + +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; +import org.apache.solr.client.solrj.request.CollectionAdminRequest; +import org.apache.solr.client.solrj.request.UpdateRequest; +import org.apache.solr.common.cloud.CollectionStateWatcher; +import org.apache.solr.common.cloud.Replica; +import org.apache.solr.common.cloud.Slice; +import org.apache.solr.common.cloud.ZkStateReader; +import org.apache.solr.common.params.ModifiableSolrParams; +import org.apache.solr.common.params.ShardParams; +import org.junit.BeforeClass; +import org.junit.Test; + +/** + * Verifies that a shard created via the Collections API stays in CONSTRUCTION state until its + * replicas are ACTIVE, so queries never route to a half-created shard (SOLR-13136). + */ +public class CreateShardConstructionTest extends SolrCloudTestCase { + + @BeforeClass + public static void setupCluster() throws Exception { + configureCluster(2).addConfig("conf", configset("cloud-minimal")).configure(); + } + + @Test + public void testNewShardStaysConstructionUntilReplicasActive() throws Exception { + String collection = "shard_construction_test"; + CollectionAdminRequest.createCollectionWithImplicitRouter(collection, "conf", "shard1", 1) + .processAndWait(cluster.getSolrClient(), 60); + cluster.waitForActiveCollection(collection, 1, 1); + + // Record every state the new shard is observed in; a ZK watcher fires on each state change + // so the brief CONSTRUCTION window can't be missed by polling. + Set observedStates = ConcurrentHashMap.newKeySet(); + CountDownLatch activeLatch = new CountDownLatch(1); + CollectionStateWatcher watcher = + (liveNodes, collectionState) -> { + if (collectionState != null) { + Slice slice = collectionState.getSlice("shard2"); + if (slice != null) { + observedStates.add(slice.getState()); + if (slice.getState() == Slice.State.ACTIVE + && slice.getReplicas().stream() + .allMatch(r -> r.getState() == Replica.State.ACTIVE)) { + activeLatch.countDown(); + } + } + } + return false; + }; + ZkStateReader zkStateReader = ZkStateReader.from(cluster.getSolrClient()); + zkStateReader.registerCollectionStateWatcher(collection, watcher); + try { + AtomicReference createError = new AtomicReference<>(); + Thread createThread = + new Thread( + () -> { + try { + CollectionAdminRequest.createShard(collection, "shard2") + .process(cluster.getSolrClient()); + } catch (Exception e) { + createError.set(e); + } + }); + createThread.start(); + + // Query while the shard is under construction; it must not fail with + // "no servers hosting shard". + ModifiableSolrParams queryParams = params("q", "*:*", "rows", "0"); + Exception queryError = null; + long deadlineNanos = System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(120000); + while (!activeLatch.await(200, TimeUnit.MILLISECONDS)) { + if (System.nanoTime() > deadlineNanos) { + break; + } + try { + cluster.getSolrClient().query(collection, queryParams); + } catch (Exception e) { + if (e.getMessage() != null && e.getMessage().contains("no servers hosting shard")) { + queryError = e; + break; + } + } + } + createThread.join(180000); + + assertNull("createShard failed: " + createError.get(), createError.get()); + assertNull("query failed during shard creation: " + queryError, queryError); + assertTrue( + "new shard was never observed in CONSTRUCTION state (saw: " + observedStates + ")", + observedStates.contains(Slice.State.CONSTRUCTION)); + assertTrue("new shard never reached ACTIVE", activeLatch.getCount() == 0); + + // updates to the new shard must be applied and searchable once it is active + UpdateRequest update = new UpdateRequest(); + update.add("id", "doc1"); + update.setParam(ShardParams._ROUTE_, "shard2"); + update.commit(cluster.getSolrClient(), collection); + assertEquals( + 1, + cluster + .getSolrClient() + .query(collection, params("q", "*:*", ShardParams._ROUTE_, "shard2")) + .getResults() + .getNumFound()); + } finally { + zkStateReader.removeCollectionStateWatcher(collection, watcher); + CollectionAdminRequest.deleteCollection(collection).process(cluster.getSolrClient()); + } + } +} diff --git a/solr/core/src/test/org/apache/solr/cloud/api/collections/CreateShardCmdTest.java b/solr/core/src/test/org/apache/solr/cloud/api/collections/CreateShardCmdTest.java new file mode 100644 index 000000000000..03c7d520aa93 --- /dev/null +++ b/solr/core/src/test/org/apache/solr/cloud/api/collections/CreateShardCmdTest.java @@ -0,0 +1,924 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.solr.cloud.api.collections; + +import static org.apache.solr.SolrTestCaseJ4.assumeWorkingMockito; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +import java.time.Instant; +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Queue; +import java.util.Set; +import java.util.concurrent.AbstractExecutorService; +import java.util.concurrent.ConcurrentLinkedQueue; +import java.util.concurrent.RejectedExecutionException; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicReference; +import java.util.function.Predicate; +import org.apache.solr.SolrTestCase; +import org.apache.solr.client.solrj.cloud.DistribStateManager; +import org.apache.solr.client.solrj.cloud.SolrCloudManager; +import org.apache.solr.client.solrj.impl.ClusterStateProvider; +import org.apache.solr.client.solrj.response.QueryResponse; +import org.apache.solr.cloud.DistributedClusterStateUpdater; +import org.apache.solr.cloud.ZkController; +import org.apache.solr.cluster.placement.PlacementPluginFactory; +import org.apache.solr.cluster.placement.plugins.SimplePlacementFactory; +import org.apache.solr.common.MapWriter; +import org.apache.solr.common.SolrCloseable; +import org.apache.solr.common.SolrException; +import org.apache.solr.common.cloud.ClusterState; +import org.apache.solr.common.cloud.CollectionStateWatcher; +import org.apache.solr.common.cloud.DocCollection; +import org.apache.solr.common.cloud.SolrZkClient; +import org.apache.solr.common.cloud.ZkNodeProps; +import org.apache.solr.common.cloud.ZkStateReader; +import org.apache.solr.common.params.CollectionParams.CollectionAction; +import org.apache.solr.common.params.CoreAdminParams; +import org.apache.solr.common.params.ModifiableSolrParams; +import org.apache.solr.common.util.NamedList; +import org.apache.solr.common.util.TimeSource; +import org.apache.solr.common.util.Utils; +import org.apache.solr.core.CoreContainer; +import org.apache.solr.handler.component.ShardHandler; +import org.apache.solr.handler.component.ShardHandlerFactory; +import org.apache.solr.handler.component.ShardRequest; +import org.apache.solr.handler.component.ShardResponse; +import org.junit.Before; +import org.junit.Test; +import org.mockito.stubbing.Answer; + +/** + * Unit tests for {@link CreateShardCmd} failure handling. The cluster collaborators are mocked so + * the replica wait can be made to time out deterministically: the replica is present in the cluster + * state (so the wait's first phase passes) but no state watcher ever reports it active, so the + * wait's latch times out. + */ +public class CreateShardCmdTest extends SolrTestCase { + + private static final String COLLECTION = "testcoll"; + private static final String SHARD = "shard2"; + private static final String NODE1 = "node1:8983_solr"; + private static final String NODE2 = "node2:8983_solr"; + + @Before + public void setUpMocks() { + assumeWorkingMockito(); + } + + @Test + @SuppressWarnings({"unchecked", "rawtypes"}) + public void testReplicaWaitFailureDeletesShardInsteadOfActivating() throws Exception { + // Cluster states in the ZooKeeper JSON shape, parsed by the production factory. + // Before the create: the collection has one unrelated slice, no shard2. + DocCollection before = + ClusterState.collectionFromObjects( + COLLECTION, + new HashMap<>( + Map.of( + "configName", + "conf1", + "shards", + Map.of( + "shard1", + Map.of( + "state", + "active", + "range", + "80000000-ffffffff", + "replicas", + Map.of())))), + 1, + Instant.now(), + null); + // Right after the shard state update lands: shard2 exists, still no replicas. + DocCollection shardCreated = + ClusterState.collectionFromObjects( + COLLECTION, + new HashMap<>( + Map.of( + "configName", + "conf1", + "shards", + Map.of( + "shard1", + Map.of( + "state", "active", "range", "80000000-7fffffff", "replicas", Map.of()), + SHARD, + Map.of( + "state", + "construction", + "range", + "80000000-ffffffff", + "replicas", + Map.of())))), + 2, + Instant.now(), + null); + // After AddReplicaCmd: shard2 has one replica, registered but never active. + DocCollection replicaAdded = + ClusterState.collectionFromObjects( + COLLECTION, + new HashMap<>( + Map.of( + "configName", + "conf1", + "shards", + Map.of( + "shard1", + Map.of( + "state", "active", "range", "80000000-7fffffff", "replicas", Map.of()), + SHARD, + Map.of( + "state", + "construction", + "range", + "80000000-ffffffff", + "replicas", + Map.of( + "core_node1", + Map.of( + "core", + "testcoll_shard2_replica_n1", + "node_name", + NODE1, + "base_url", + "http://node1:8983/solr", + "state", + "down", + "leader", + "true", + "type", + "NRT")))))), + 3, + Instant.now(), + null); + Set liveNodes = Set.of(NODE1, NODE2); + ClusterState beforeState = new ClusterState(liveNodes, Map.of(COLLECTION, before)); + ClusterState replicaState = new ClusterState(liveNodes, Map.of(COLLECTION, replicaAdded)); + + ZkStateReader zkStateReader = mock(ZkStateReader.class); + // Simulate the cluster state advancing as the state updates recorded below are offered: + // the createshard update brings the slice into being, the addreplica update adds the + // replica, the deletecore update (DeleteReplicaCmd removing the replica from the state) + // drops it again, and the deleteshard update removes the slice. getClusterState reports + // the current simulated state, and waitForState answers its predicate against it, timing + // out like the real implementation when the predicate does not hold. In particular the + // replica never becomes active (no watcher fires), and the added replica is still in the + // state after its core is unloaded, so DeleteReplicaCmd has to remove it itself. + AtomicReference currentDoc = new AtomicReference<>(shardCreated); + when(zkStateReader.getClusterState()) + .thenAnswer( + invocation -> new ClusterState(liveNodes, Map.of(COLLECTION, currentDoc.get()))); + doAnswer( + invocation -> { + Predicate predicate = invocation.getArgument(3); + DocCollection current = currentDoc.get(); + if (predicate.test(current)) { + return current; + } + throw new TimeoutException("simulated state never satisfied the wait predicate"); + }) + .when(zkStateReader) + .waitForState(anyString(), anyLong(), any(TimeUnit.class), any(Predicate.class)); + // DeleteShardCmd cleans up shard metadata in ZooKeeper after deleting a shard. + when(zkStateReader.getZkClient()).thenReturn(mock(SolrZkClient.class)); + + ClusterStateProvider stateProvider = mock(ClusterStateProvider.class); + when(stateProvider.getClusterState()).thenReturn(replicaState); + when(stateProvider.getLiveNodes()).thenReturn(liveNodes); + SolrCloudManager cloudManager = mock(SolrCloudManager.class); + when(cloudManager.getClusterStateProvider()).thenReturn(stateProvider); + when(cloudManager.getClusterState()).thenReturn(replicaState); + when(cloudManager.getTimeSource()).thenReturn(new TimeSource.NanoTimeSource()); + when(cloudManager.getDistribStateManager()).thenReturn(mock(DistribStateManager.class)); + + // The shard handler accepts the CREATE request and reports success for it, once. + ShardHandler shardHandler = mock(ShardHandler.class); + ShardHandlerFactory shardHandlerFactory = mock(ShardHandlerFactory.class); + when(shardHandler.getShardHandlerFactory()).thenReturn(shardHandlerFactory); + when(shardHandlerFactory.getShardHandler()).thenReturn(shardHandler); + AtomicReference submitted = new AtomicReference<>(); + AtomicBoolean responded = new AtomicBoolean(); + doAnswer( + invocation -> { + submitted.set(invocation.getArgument(0)); + // Each submitted request gets its own one-shot success response, so later + // requests (for example the cleanup's replica unload) can also complete. + responded.set(false); + return null; + }) + .when(shardHandler) + .submit(any(ShardRequest.class), any(), any(ModifiableSolrParams.class)); + Answer takeAnswer = + invocation -> { + if (submitted.get() == null || !responded.compareAndSet(false, true)) { + return null; + } + ShardResponse response = new ShardResponse(); + response.setShardRequest(submitted.get()); + QueryResponse queryResponse = new QueryResponse(); + queryResponse.setResponse( + new NamedList<>(Map.of("responseHeader", new NamedList<>(Map.of("status", 0))))); + response.setSolrResponse(queryResponse); + return response; + }; + when(shardHandler.takeCompletedOrError()).thenAnswer(takeAnswer); + when(shardHandler.takeCompletedIncludingErrors()).thenAnswer(takeAnswer); + + DistributedClusterStateUpdater stateUpdater = mock(DistributedClusterStateUpdater.class); + when(stateUpdater.isDistributedStateUpdate()).thenReturn(false); + + CoreContainer coreContainer = mock(CoreContainer.class); + PlacementPluginFactory placementFactory = mock(PlacementPluginFactory.class); + when(placementFactory.createPluginInstance()) + .thenReturn(new SimplePlacementFactory().createPluginInstance()); + when(coreContainer.getPlacementPluginFactory()).thenReturn(placementFactory); + ZkController zkController = mock(ZkController.class); + when(coreContainer.getZkController()).thenReturn(zkController); + when(zkController.getNodeName()).thenReturn(NODE1); + + CollectionCommandContext ccc = mock(CollectionCommandContext.class); + when(ccc.isDistributedCollectionAPI()).thenReturn(false); + when(ccc.newShardHandler()).thenReturn(shardHandler); + when(ccc.getSolrCloudManager()).thenReturn(cloudManager); + when(ccc.getZkStateReader()).thenReturn(zkStateReader); + when(ccc.getDistributedClusterStateUpdater()).thenReturn(stateUpdater); + when(ccc.getCoreContainer()).thenReturn(coreContainer); + when(ccc.getAdminPath()).thenReturn("/admin/collections"); + when(ccc.getCloseableToLatchOn()).thenReturn(mock(SolrCloseable.class)); + // DeleteReplicaCmd unloads the core on the command context's executor; run it inline so + // no thread outlives the test. + when(ccc.getExecutorService()).thenReturn(new SameThreadExecutorService()); + // Record every cluster state update the command offers, and advance the simulated + // cluster state (see the ZkStateReader stub above) as each update lands. + List> offeredUpdates = new ArrayList<>(); + doAnswer( + invocation -> { + Object update = invocation.getArgument(0); + Map recorded; + if (update instanceof ZkNodeProps zkNodeProps) { + recorded = new HashMap<>(zkNodeProps.getProperties()); + } else { + recorded = + new HashMap<>((Map) Utils.fromJSON(Utils.toJSON(update))); + } + offeredUpdates.add(recorded); + switch (String.valueOf(recorded.get("operation"))) { + case "createshard" -> currentDoc.set(shardCreated); + case "addreplica" -> currentDoc.set(replicaAdded); + case "deletecore" -> currentDoc.set(shardCreated); + case "deleteshard" -> currentDoc.set(before); + default -> {} + } + return null; + }) + .when(ccc) + .offerStateUpdate(any(MapWriter.class)); + + AdminCmdContext adminCmdContext = + new AdminCmdContext(CollectionAction.CREATESHARD).withClusterState(beforeState); + ZkNodeProps message = + new ZkNodeProps( + Map.of( + "collection", + COLLECTION, + "shard", + SHARD, + "timeout", + 1, + "replicationFactor", + 1, + "createNodeSet", + NODE1)); + + Exception thrown = null; + try { + new CreateShardCmd(ccc).call(adminCmdContext, message, new NamedList<>()); + } catch (Exception e) { + thrown = e; + } + assertNotNull("expected the create to fail in the replica wait", thrown); + if (!String.valueOf(thrown.getMessage()).contains("replicas of shard")) { + throw new AssertionError("expected the replica wait timeout, got: " + thrown, thrown); + } + + List operations = new ArrayList<>(); + for (Map update : offeredUpdates) { + operations.add(String.valueOf(update.get("operation"))); + } + assertTrue( + "expected a deleteshard state update, offered: " + offeredUpdates, + operations.contains("deleteshard")); + assertTrue( + "expected the cleanup to also delete the replica that AddReplicaCmd added, offered: " + + offeredUpdates, + operations.contains("deletecore")); + for (Map update : offeredUpdates) { + if ("updateshardstate".equalsIgnoreCase(String.valueOf(update.get("operation")))) { + fail("the failed shard must not be activated, offered: " + update); + } + } + } + + /** + * Drives the first cleanup block, the one around adding the replicas. The replica is registered + * in the cluster state first and creating its core fails afterwards, so {@code AddReplicaCmd} + * throws with the replica already in place. The create must then delete the half-created shard, + * including that replica, and the caller must see the original add failure rather than a cleanup + * result. + */ + @Test + @SuppressWarnings({"unchecked", "rawtypes"}) + public void testAddReplicaFailureDeletesShardAndReplica() throws Exception { + // Cluster states in the ZooKeeper JSON shape, parsed by the production factory. + // Before the create: the collection has one unrelated slice, no shard2. + DocCollection before = + ClusterState.collectionFromObjects( + COLLECTION, + new HashMap<>( + Map.of( + "configName", + "conf1", + "shards", + Map.of( + "shard1", + Map.of( + "state", + "active", + "range", + "80000000-ffffffff", + "replicas", + Map.of())))), + 1, + Instant.now(), + null); + // Right after the shard state update lands: shard2 exists, still no replicas. + DocCollection shardCreated = + ClusterState.collectionFromObjects( + COLLECTION, + new HashMap<>( + Map.of( + "configName", + "conf1", + "shards", + Map.of( + "shard1", + Map.of( + "state", "active", "range", "80000000-7fffffff", "replicas", Map.of()), + SHARD, + Map.of( + "state", + "construction", + "range", + "80000000-ffffffff", + "replicas", + Map.of())))), + 2, + Instant.now(), + null); + // After the replica is registered: shard2 has one replica, not yet active. Creating its + // core is what fails below, so this is the state the cleanup re-reads. + DocCollection replicaAdded = + ClusterState.collectionFromObjects( + COLLECTION, + new HashMap<>( + Map.of( + "configName", + "conf1", + "shards", + Map.of( + "shard1", + Map.of( + "state", "active", "range", "80000000-7fffffff", "replicas", Map.of()), + SHARD, + Map.of( + "state", + "construction", + "range", + "80000000-ffffffff", + "replicas", + Map.of( + "core_node1", + Map.of( + "core", + "testcoll_shard2_replica_n1", + "node_name", + NODE1, + "base_url", + "http://node1:8983/solr", + "state", + "down", + "leader", + "true", + "type", + "NRT")))))), + 3, + Instant.now(), + null); + Set liveNodes = Set.of(NODE1, NODE2); + ClusterState beforeState = new ClusterState(liveNodes, Map.of(COLLECTION, before)); + ClusterState replicaState = new ClusterState(liveNodes, Map.of(COLLECTION, replicaAdded)); + + ZkStateReader zkStateReader = mock(ZkStateReader.class); + // Simulate the cluster state advancing as the state updates recorded below are offered: + // the createshard update brings the slice into being, the addreplica update registers the + // replica, the deletecore update (DeleteReplicaCmd removing the replica from the state) + // drops it again, and the deleteshard update removes the slice. getClusterState reports + // the current simulated state, and waitForState answers its predicate against it, timing + // out like the real implementation when the predicate does not hold. + AtomicReference currentDoc = new AtomicReference<>(shardCreated); + when(zkStateReader.getClusterState()) + .thenAnswer( + invocation -> new ClusterState(liveNodes, Map.of(COLLECTION, currentDoc.get()))); + doAnswer( + invocation -> { + Predicate predicate = invocation.getArgument(3); + DocCollection current = currentDoc.get(); + if (predicate.test(current)) { + return current; + } + throw new TimeoutException("simulated state never satisfied the wait predicate"); + }) + .when(zkStateReader) + .waitForState(anyString(), anyLong(), any(TimeUnit.class), any(Predicate.class)); + // DeleteShardCmd cleans up shard metadata in ZooKeeper after deleting a shard. + when(zkStateReader.getZkClient()).thenReturn(mock(SolrZkClient.class)); + + ClusterStateProvider stateProvider = mock(ClusterStateProvider.class); + when(stateProvider.getClusterState()).thenReturn(replicaState); + when(stateProvider.getLiveNodes()).thenReturn(liveNodes); + SolrCloudManager cloudManager = mock(SolrCloudManager.class); + when(cloudManager.getClusterStateProvider()).thenReturn(stateProvider); + when(cloudManager.getClusterState()).thenReturn(replicaState); + when(cloudManager.getTimeSource()).thenReturn(new TimeSource.NanoTimeSource()); + when(cloudManager.getDistribStateManager()).thenReturn(mock(DistribStateManager.class)); + + // The CREATE request for the replica's core fails; every other request (in particular + // the cleanup's UNLOAD of that same core) succeeds. + ShardHandler shardHandler = mock(ShardHandler.class); + ShardHandlerFactory shardHandlerFactory = mock(ShardHandlerFactory.class); + when(shardHandler.getShardHandlerFactory()).thenReturn(shardHandlerFactory); + when(shardHandlerFactory.getShardHandler()).thenReturn(shardHandler); + Queue submittedRequests = new ConcurrentLinkedQueue<>(); + Queue submittedParams = new ConcurrentLinkedQueue<>(); + doAnswer( + invocation -> { + submittedRequests.add(invocation.getArgument(0)); + submittedParams.add(invocation.getArgument(2)); + return null; + }) + .when(shardHandler) + .submit(any(ShardRequest.class), any(), any(ModifiableSolrParams.class)); + Answer takeAnswer = + invocation -> { + ShardRequest sreq = submittedRequests.poll(); + ModifiableSolrParams params = submittedParams.poll(); + if (sreq == null) { + return null; + } + ShardResponse response = new ShardResponse(); + response.setShardRequest(sreq); + if (params != null + && CoreAdminParams.CoreAdminAction.CREATE + .toString() + .equals(params.get(CoreAdminParams.ACTION))) { + response.setException( + new SolrException( + SolrException.ErrorCode.SERVER_ERROR, "simulated core creation failure")); + } else { + QueryResponse queryResponse = new QueryResponse(); + queryResponse.setResponse( + new NamedList<>(Map.of("responseHeader", new NamedList<>(Map.of("status", 0))))); + response.setSolrResponse(queryResponse); + } + return response; + }; + when(shardHandler.takeCompletedOrError()).thenAnswer(takeAnswer); + when(shardHandler.takeCompletedIncludingErrors()).thenAnswer(takeAnswer); + + DistributedClusterStateUpdater stateUpdater = mock(DistributedClusterStateUpdater.class); + when(stateUpdater.isDistributedStateUpdate()).thenReturn(false); + + CoreContainer coreContainer = mock(CoreContainer.class); + PlacementPluginFactory placementFactory = mock(PlacementPluginFactory.class); + when(placementFactory.createPluginInstance()) + .thenReturn(new SimplePlacementFactory().createPluginInstance()); + when(coreContainer.getPlacementPluginFactory()).thenReturn(placementFactory); + ZkController zkController = mock(ZkController.class); + when(coreContainer.getZkController()).thenReturn(zkController); + when(zkController.getNodeName()).thenReturn(NODE1); + + CollectionCommandContext ccc = mock(CollectionCommandContext.class); + when(ccc.isDistributedCollectionAPI()).thenReturn(false); + when(ccc.newShardHandler()).thenReturn(shardHandler); + when(ccc.getSolrCloudManager()).thenReturn(cloudManager); + when(ccc.getZkStateReader()).thenReturn(zkStateReader); + when(ccc.getDistributedClusterStateUpdater()).thenReturn(stateUpdater); + when(ccc.getCoreContainer()).thenReturn(coreContainer); + when(ccc.getAdminPath()).thenReturn("/admin/collections"); + when(ccc.getCloseableToLatchOn()).thenReturn(mock(SolrCloseable.class)); + // DeleteReplicaCmd unloads the core on the command context's executor; run it inline so + // no thread outlives the test. + when(ccc.getExecutorService()).thenReturn(new SameThreadExecutorService()); + // Record every cluster state update the command offers, and advance the simulated + // cluster state (see the ZkStateReader stub above) as each update lands. + List> offeredUpdates = new ArrayList<>(); + doAnswer( + invocation -> { + Object update = invocation.getArgument(0); + Map recorded; + if (update instanceof ZkNodeProps zkNodeProps) { + recorded = new HashMap<>(zkNodeProps.getProperties()); + } else { + recorded = + new HashMap<>((Map) Utils.fromJSON(Utils.toJSON(update))); + } + offeredUpdates.add(recorded); + switch (String.valueOf(recorded.get("operation"))) { + case "createshard" -> currentDoc.set(shardCreated); + case "addreplica" -> currentDoc.set(replicaAdded); + case "deletecore" -> currentDoc.set(shardCreated); + case "deleteshard" -> currentDoc.set(before); + default -> {} + } + return null; + }) + .when(ccc) + .offerStateUpdate(any(MapWriter.class)); + + AdminCmdContext adminCmdContext = + new AdminCmdContext(CollectionAction.CREATESHARD).withClusterState(beforeState); + ZkNodeProps message = + new ZkNodeProps( + Map.of( + "collection", + COLLECTION, + "shard", + SHARD, + "timeout", + 1, + "replicationFactor", + 1, + "createNodeSet", + NODE1)); + + Exception thrown = null; + try { + new CreateShardCmd(ccc).call(adminCmdContext, message, new NamedList<>()); + } catch (Exception e) { + thrown = e; + } + assertNotNull("expected the create to fail when creating the replica's core fails", thrown); + if (!String.valueOf(thrown.getMessage()).contains("ADDREPLICA failed to create replica")) { + throw new AssertionError("expected the add replica failure, got: " + thrown, thrown); + } + + List operations = new ArrayList<>(); + for (Map update : offeredUpdates) { + operations.add(String.valueOf(update.get("operation"))); + } + assertTrue( + "expected the cleanup to delete the replica that was registered before the failure," + + " offered: " + + offeredUpdates, + operations.contains("deletecore")); + for (Map update : offeredUpdates) { + if ("deletecore".equals(String.valueOf(update.get("operation")))) { + assertEquals( + "the cleanup must delete the replica that AddReplicaCmd registered", + "testcoll_shard2_replica_n1", + String.valueOf(update.get("core"))); + } + if ("updateshardstate".equalsIgnoreCase(String.valueOf(update.get("operation")))) { + fail("the failed shard must not be activated, offered: " + update); + } + } + assertTrue( + "expected a deleteshard state update, offered: " + offeredUpdates, + operations.contains("deleteshard")); + assertNull("expected no slice left behind after the cleanup", currentDoc.get().getSlice(SHARD)); + } + + /** + * Pins the current buffered-updates policy: when the leader answers the REQUESTAPPLYUPDATES + * request with a failure, the create logs it and still activates the shard. A core refuses that + * request when it is not buffering, which is expected for a plain create, so the failure alone + * must not delete the shard or fail the create. Whether a failure should instead abort the + * creation is an open question on the PR; if the policy changes, this test changes with it. + */ + @Test + @SuppressWarnings({"unchecked", "rawtypes"}) + public void testBufferedUpdatesFailureDoesNotFailCreate() throws Exception { + DocCollection before = + ClusterState.collectionFromObjects( + COLLECTION, + new HashMap<>( + Map.of( + "configName", + "conf1", + "shards", + Map.of( + "shard1", + Map.of( + "state", + "active", + "range", + "80000000-ffffffff", + "replicas", + Map.of())))), + 1, + Instant.now(), + null); + DocCollection shardCreated = + ClusterState.collectionFromObjects( + COLLECTION, + new HashMap<>( + Map.of( + "configName", + "conf1", + "shards", + Map.of( + "shard1", + Map.of( + "state", "active", "range", "80000000-7fffffff", "replicas", Map.of()), + SHARD, + Map.of( + "state", + "construction", + "range", + "80000000-ffffffff", + "replicas", + Map.of())))), + 2, + Instant.now(), + null); + // shard2 has one replica, already active and marked as the leader. + DocCollection replicaActive = + ClusterState.collectionFromObjects( + COLLECTION, + new HashMap<>( + Map.of( + "configName", + "conf1", + "shards", + Map.of( + "shard1", + Map.of( + "state", "active", "range", "80000000-7fffffff", "replicas", Map.of()), + SHARD, + Map.of( + "state", + "construction", + "range", + "80000000-ffffffff", + "replicas", + Map.of( + "core_node1", + Map.of( + "core", + "testcoll_shard2_replica_n1", + "node_name", + NODE1, + "base_url", + "http://node1:8983/solr", + "state", + "active", + "leader", + "true", + "type", + "NRT")))))), + 3, + Instant.now(), + null); + Set liveNodes = Set.of(NODE1, NODE2); + ClusterState beforeState = new ClusterState(liveNodes, Map.of(COLLECTION, before)); + ClusterState createdState = new ClusterState(liveNodes, Map.of(COLLECTION, shardCreated)); + ClusterState activeState = new ClusterState(liveNodes, Map.of(COLLECTION, replicaActive)); + + ZkStateReader zkStateReader = mock(ZkStateReader.class); + when(zkStateReader.getClusterState()).thenReturn(createdState, activeState); + doAnswer( + invocation -> { + Predicate predicate = invocation.getArgument(3); + predicate.test(replicaActive); + return replicaActive; + }) + .when(zkStateReader) + .waitForState(anyString(), anyLong(), any(TimeUnit.class), any(Predicate.class)); + // The replica is already active, so notify the active-replica watcher immediately. + doAnswer( + invocation -> { + CollectionStateWatcher watcher = invocation.getArgument(1); + watcher.onStateChanged(liveNodes, replicaActive); + return null; + }) + .when(zkStateReader) + .registerCollectionStateWatcher(anyString(), any(CollectionStateWatcher.class)); + + ClusterStateProvider stateProvider = mock(ClusterStateProvider.class); + when(stateProvider.getClusterState()).thenReturn(activeState); + when(stateProvider.getLiveNodes()).thenReturn(liveNodes); + SolrCloudManager cloudManager = mock(SolrCloudManager.class); + when(cloudManager.getClusterStateProvider()).thenReturn(stateProvider); + when(cloudManager.getClusterState()).thenReturn(activeState); + when(cloudManager.getTimeSource()).thenReturn(new TimeSource.NanoTimeSource()); + when(cloudManager.getDistribStateManager()).thenReturn(mock(DistribStateManager.class)); + + // The CREATE request succeeds; the REQUESTAPPLYUPDATES request to the leader fails, as it + // does when the core is not buffering updates. + ShardHandler shardHandler = mock(ShardHandler.class); + ShardHandlerFactory shardHandlerFactory = mock(ShardHandlerFactory.class); + when(shardHandler.getShardHandlerFactory()).thenReturn(shardHandlerFactory); + when(shardHandlerFactory.getShardHandler()).thenReturn(shardHandler); + Queue submittedRequests = new ConcurrentLinkedQueue<>(); + Queue submittedParams = new ConcurrentLinkedQueue<>(); + AtomicBoolean applyUpdatesSubmitted = new AtomicBoolean(); + doAnswer( + invocation -> { + submittedRequests.add(invocation.getArgument(0)); + ModifiableSolrParams params = invocation.getArgument(2); + submittedParams.add(params); + if (CoreAdminParams.CoreAdminAction.REQUESTAPPLYUPDATES + .toString() + .equals(params.get(CoreAdminParams.ACTION))) { + applyUpdatesSubmitted.set(true); + } + return null; + }) + .when(shardHandler) + .submit(any(ShardRequest.class), any(), any(ModifiableSolrParams.class)); + Answer takeAnswer = + invocation -> { + ShardRequest sreq = submittedRequests.poll(); + ModifiableSolrParams params = submittedParams.poll(); + if (sreq == null) { + return null; + } + ShardResponse response = new ShardResponse(); + response.setShardRequest(sreq); + if (params != null + && CoreAdminParams.CoreAdminAction.REQUESTAPPLYUPDATES + .toString() + .equals(params.get(CoreAdminParams.ACTION))) { + response.setException( + new SolrException( + SolrException.ErrorCode.SERVER_ERROR, "Core is not buffering updates")); + } else { + QueryResponse queryResponse = new QueryResponse(); + queryResponse.setResponse( + new NamedList<>(Map.of("responseHeader", new NamedList<>(Map.of("status", 0))))); + response.setSolrResponse(queryResponse); + } + return response; + }; + when(shardHandler.takeCompletedOrError()).thenAnswer(takeAnswer); + when(shardHandler.takeCompletedIncludingErrors()).thenAnswer(takeAnswer); + + DistributedClusterStateUpdater stateUpdater = mock(DistributedClusterStateUpdater.class); + when(stateUpdater.isDistributedStateUpdate()).thenReturn(false); + + CoreContainer coreContainer = mock(CoreContainer.class); + PlacementPluginFactory placementFactory = mock(PlacementPluginFactory.class); + when(placementFactory.createPluginInstance()) + .thenReturn(new SimplePlacementFactory().createPluginInstance()); + when(coreContainer.getPlacementPluginFactory()).thenReturn(placementFactory); + ZkController zkController = mock(ZkController.class); + when(coreContainer.getZkController()).thenReturn(zkController); + when(zkController.getNodeName()).thenReturn(NODE1); + + CollectionCommandContext ccc = mock(CollectionCommandContext.class); + when(ccc.isDistributedCollectionAPI()).thenReturn(false); + when(ccc.newShardHandler()).thenReturn(shardHandler); + when(ccc.getSolrCloudManager()).thenReturn(cloudManager); + when(ccc.getZkStateReader()).thenReturn(zkStateReader); + when(ccc.getDistributedClusterStateUpdater()).thenReturn(stateUpdater); + when(ccc.getCoreContainer()).thenReturn(coreContainer); + when(ccc.getAdminPath()).thenReturn("/admin/collections"); + when(ccc.getCloseableToLatchOn()).thenReturn(mock(SolrCloseable.class)); + List> offeredUpdates = new ArrayList<>(); + doAnswer( + invocation -> { + Object update = invocation.getArgument(0); + if (update instanceof ZkNodeProps zkNodeProps) { + offeredUpdates.add(new HashMap<>(zkNodeProps.getProperties())); + } else { + offeredUpdates.add( + new HashMap<>((Map) Utils.fromJSON(Utils.toJSON(update)))); + } + return null; + }) + .when(ccc) + .offerStateUpdate(any(MapWriter.class)); + + AdminCmdContext adminCmdContext = + new AdminCmdContext(CollectionAction.CREATESHARD).withClusterState(beforeState); + ZkNodeProps message = + new ZkNodeProps( + Map.of( + "collection", + COLLECTION, + "shard", + SHARD, + "timeout", + 5, + "replicationFactor", + 1, + "createNodeSet", + NODE1)); + + Exception thrown = null; + try { + new CreateShardCmd(ccc).call(adminCmdContext, message, new NamedList<>()); + } catch (Exception e) { + thrown = e; + } + assertNull("the create must not fail when only applying buffered updates fails", thrown); + assertTrue( + "expected the leader to be asked to apply buffered updates", applyUpdatesSubmitted.get()); + + List operations = new ArrayList<>(); + for (Map update : offeredUpdates) { + operations.add(String.valueOf(update.get("operation"))); + } + assertFalse( + "the shard must not be deleted over a buffered-updates failure, offered: " + offeredUpdates, + operations.contains("deleteshard")); + boolean activated = + offeredUpdates.stream() + .anyMatch( + update -> + "updateshardstate".equalsIgnoreCase(String.valueOf(update.get("operation"))) + && "active".equalsIgnoreCase(String.valueOf(update.get(SHARD)))); + assertTrue( + "expected the shard to be activated despite the buffered-updates failure, offered: " + + offeredUpdates, + activated); + } + + /** + * An {@link AbstractExecutorService} that runs every task on the submitting thread, so tests + * using it leave no pool thread behind. + */ + private static final class SameThreadExecutorService extends AbstractExecutorService { + private boolean shutdown; + + @Override + public void shutdown() { + shutdown = true; + } + + @Override + public List shutdownNow() { + shutdown = true; + return List.of(); + } + + @Override + public boolean isShutdown() { + return shutdown; + } + + @Override + public boolean isTerminated() { + return shutdown; + } + + @Override + public boolean awaitTermination(long timeout, TimeUnit unit) { + return true; + } + + @Override + public void execute(Runnable command) { + if (shutdown) { + throw new RejectedExecutionException("executor is shut down"); + } + command.run(); + } + } +}