diff --git a/changelog/unreleased/SOLR-18391-collection-creation-failure-cleanup.yml b/changelog/unreleased/SOLR-18391-collection-creation-failure-cleanup.yml new file mode 100644 index 000000000000..c5d6300d02e3 --- /dev/null +++ b/changelog/unreleased/SOLR-18391-collection-creation-failure-cleanup.yml @@ -0,0 +1,8 @@ +# See https://github.com/apache/solr/blob/main/dev-docs/changelog.adoc +title: A collection creation that fails part way no longer leaves a half-created collection behind, and the error returned to the client names the cause; a create whose final alias write fails now removes the completed collection instead of leaving it behind +type: fixed +authors: + - name: Nick Shanin +links: + - name: SOLR-18391 + url: https://issues.apache.org/jira/browse/SOLR-18391 diff --git a/solr/core/src/java/org/apache/solr/cloud/api/collections/CreateCollectionCmd.java b/solr/core/src/java/org/apache/solr/cloud/api/collections/CreateCollectionCmd.java index 14291d7b9411..a1f9e9dd9be0 100644 --- a/solr/core/src/java/org/apache/solr/cloud/api/collections/CreateCollectionCmd.java +++ b/solr/core/src/java/org/apache/solr/cloud/api/collections/CreateCollectionCmd.java @@ -21,7 +21,6 @@ import static org.apache.solr.common.params.CollectionAdminParams.COLL_CONF; import static org.apache.solr.common.params.CollectionParams.CollectionAction.ADDREPLICA; import static org.apache.solr.common.params.CollectionParams.CollectionAction.CREATE; -import static org.apache.solr.common.params.CollectionParams.CollectionAction.DELETE; import static org.apache.solr.common.params.CommonAdminParams.ASYNC; import static org.apache.solr.common.params.CommonAdminParams.WAIT_FOR_FINAL_STATE; import static org.apache.solr.common.params.CommonParams.NAME; @@ -80,6 +79,7 @@ import org.apache.solr.common.params.ModifiableSolrParams; import org.apache.solr.common.util.EnvUtils; import org.apache.solr.common.util.NamedList; +import org.apache.solr.common.util.RetryUtil; import org.apache.solr.common.util.SimpleOrderedMap; import org.apache.solr.common.util.Utils; import org.apache.solr.core.ConfigSetService; @@ -98,6 +98,11 @@ public class CreateCollectionCmd implements CollApiCmds.CollectionApiCommand { public static final String PRS_DEFAULT_PROP = "solr.cloud.prs.enabled"; + // The alias write at the end of a create is retried a few times on a ZooKeeper error: the + // collection itself is complete by then, and a transient error should not fail the create. + private static final int ALIAS_CREATION_ATTEMPTS = 3; + private static final long ALIAS_CREATION_RETRY_PAUSE_MS = 200; + public CreateCollectionCmd(CollectionCommandContext ccc) { this.ccc = ccc; } @@ -152,6 +157,9 @@ public void call(AdminCmdContext adminCmdContext, ZkNodeProps message, NamedList DocCollection newColl = null; final String collectionPath = DocCollection.getCollectionPath(collectionName); + // True once this command may have written the collection state. Any later failure then deletes + // the collection again, whatever the failure is, so that a failed create leaves nothing behind. + boolean stateWritten = false; try { ZkStateReader zkStateReader = ccc.getZkStateReader(); @@ -201,6 +209,7 @@ public void call(AdminCmdContext adminCmdContext, ZkNodeProps message, NamedList new ClusterStateMutator(ccc.getSolrCloudManager()).createCollection(clusterState, m); byte[] data = Utils.toJSON(Map.of(collectionName, command.collection)); ccc.getZkStateReader().getZkClient().create(collectionPath, data, CreateMode.PERSISTENT); + stateWritten = true; clusterState = clusterState.copyWith(collectionName, command.collection); newColl = command.collection; ccc.submitIntraProcessMessage(new RefreshCollectionMessage(collectionName)); @@ -216,6 +225,17 @@ public void call(AdminCmdContext adminCmdContext, ZkNodeProps message, NamedList e); } } else { + // The check at the top of this command reads this node's cluster state view, which can + // lag behind ZooKeeper. The state update below creates the collection's state.json, so + // if one is already in ZooKeeper it belongs to an existing collection, and a failed + // create must not delete a collection this command did not create. This command holds + // the collection lock, so state found here is not its own. + if (zkStateReader.getZkClient().exists(collectionPath)) { + throw new SolrException( + SolrException.ErrorCode.BAD_REQUEST, "collection already exists: " + collectionName); + } + // set before submitting: a submission that fails may still have been applied + stateWritten = true; if (ccc.getDistributedClusterStateUpdater().isDistributedStateUpdate()) { // The message has been crafted by CollectionsHandler.CollectionOperation.CREATE_OP and // defines the QUEUE_OPERATION to be CollectionParams.CollectionAction.CREATE. @@ -243,28 +263,22 @@ public void call(AdminCmdContext adminCmdContext, ZkNodeProps message, NamedList // refresh cluster state (value read below comes from Zookeeper watch firing following the // update done previously, be it by Overseer or by this thread when updates are distributed) clusterState = ccc.getSolrCloudManager().getClusterState(); + // The wait above saw the collection through a watch. This is a second, separate read, and + // while the watch is being released the collection can be absent from it for a moment. + // Replica assignment reads this cluster state, so make sure it holds what the wait saw. + if (!clusterState.hasCollection(collectionName)) { + clusterState = clusterState.copyWith(collectionName, newColl); + } } - final List replicaPositions; - try { - replicaPositions = - buildReplicaPositions( - ccc.getCoreContainer(), - ccc.getSolrCloudManager(), - clusterState, - message, - shardNames, - numReplicas); - } catch (Assign.AssignmentException e) { - ZkNodeProps deleteMessage = new ZkNodeProps("name", collectionName); - new DeleteCollectionCmd(ccc) - .call( - adminCmdContext.subRequestContext(DELETE).withClusterState(clusterState), - deleteMessage, - results); - // unwrap the exception - throw new SolrException(ErrorCode.BAD_REQUEST, e.getMessage(), e.getCause()); - } + final List replicaPositions = + buildReplicaPositions( + ccc.getCoreContainer(), + ccc.getSolrCloudManager(), + clusterState, + message, + shardNames, + numReplicas); if (replicaPositions.isEmpty()) { log.debug("Finished create command for collection: {}", collectionName); @@ -428,6 +442,12 @@ public void call(AdminCmdContext adminCmdContext, ZkNodeProps message, NamedList boolean failure = results.get("failure") != null && ((SimpleOrderedMap) results.get("failure")).size() > 0; + String failureDetail = null; + if (failure) { + // Name the first failure only. The map holds one entry per failed core, and dumping all + // of them puts every node's error text, URLs and paths included, in the client message. + failureDetail = String.valueOf(((SimpleOrderedMap) results.get("failure")).getVal(0)); + } if (isPRS) { TimeOut timeout = new TimeOut( @@ -446,18 +466,20 @@ public void call(AdminCmdContext adminCmdContext, ZkNodeProps message, NamedList // we have successfully found all replicas to be ACTIVE } else { failure = true; + if (failureDetail == null) { + failureDetail = "not all replicas became active"; + } } } if (failure) { - // Let's cleanup as we hit an exception - // We shouldn't be passing 'results' here for the cleanup as the response would then contain - // 'success' element, which may be interpreted by the user as a positive ack - CollectionHandlingUtils.cleanupCollection( - adminCmdContext, collectionName, new NamedList<>(), ccc); - log.info("Cleaned up artifacts for failed create collection for [{}]", collectionName); + // the collection is cleaned up where this exception is caught, below throw new SolrException( ErrorCode.BAD_REQUEST, - "Underlying core creation failed while creating collection: " + collectionName); + "Underlying core creation failed while creating collection: " + + collectionName + + " (" + + failureDetail + + ")"); } else { ccc.submitIntraProcessMessage(new RefreshCollectionMessage(collectionName)); log.debug("Finished create command on all shards for collection: {}", collectionName); @@ -480,15 +502,94 @@ public void call(AdminCmdContext adminCmdContext, ZkNodeProps message, NamedList // create an alias pointing to the new collection, if different from the collectionName if (!alias.equals(collectionName)) { - ccc.getZkStateReader() - .aliasesManager - .applyModificationAndExportToZk(a -> a.cloneWithCollectionAlias(alias, collectionName)); + runWithBoundedRetries( + () -> + ccc.getZkStateReader() + .aliasesManager + .applyModificationAndExportToZk( + a -> a.cloneWithCollectionAlias(alias, collectionName)), + ALIAS_CREATION_ATTEMPTS, + ALIAS_CREATION_RETRY_PAUSE_MS); } - } catch (SolrException ex) { - throw ex; } catch (Exception ex) { - throw new SolrException(SolrException.ErrorCode.SERVER_ERROR, null, ex); + if (ex instanceof InterruptedException) { + Thread.currentThread().interrupt(); + } + // An interrupted thread cannot complete the ZooKeeper calls the cleanup needs + if (stateWritten && !Thread.currentThread().isInterrupted()) { + cleanupFailedCreate(adminCmdContext, collectionName, ex); + } + if (ex instanceof SolrException solrException) { + throw solrException; + } + if (ex instanceof Assign.AssignmentException) { + // unwrap the exception + throw new SolrException(ErrorCode.BAD_REQUEST, ex.getMessage(), ex.getCause()); + } + throw new SolrException( + ErrorCode.SERVER_ERROR, "Could not create collection " + collectionName + ": " + ex, ex); + } + } + + /** + * Runs {@code op}, retrying it if it fails with a {@link ZooKeeperException}, for up to {@code + * maxAttempts} attempts in total and pausing {@code pauseMillis} between attempts. Any other + * failure propagates at once, and so does a ZooKeeper failure on the last attempt. An interrupt + * is never retried: the interrupt flag is restored and the {@link InterruptedException} + * propagates. + */ + static void runWithBoundedRetries(RetryUtil.RetryCmd op, int maxAttempts, long pauseMillis) + throws Exception { + int attempt = 0; + while (true) { + attempt++; + try { + op.execute(); + return; + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw e; + } catch (ZooKeeperException e) { + if (attempt >= maxAttempts) { + throw e; + } + log.warn( + "Attempt {} of {} failed with a ZooKeeper error; retrying after a pause", + attempt, + maxAttempts, + e); + try { + Thread.sleep(pauseMillis); + } catch (InterruptedException interruptedDuringPause) { + Thread.currentThread().interrupt(); + throw interruptedDuringPause; + } + } + } + } + + /** + * Deletes what a failed create left behind. A failure of the cleanup itself is logged and added + * to {@code cause} as suppressed, so that the client still gets the original error. + */ + private void cleanupFailedCreate( + AdminCmdContext adminCmdContext, String collectionName, Exception cause) { + try { + // We shouldn't be passing 'results' here for the cleanup as the response would then contain + // 'success' element, which may be interpreted by the user as a positive ack + CollectionHandlingUtils.cleanupCollection( + adminCmdContext, collectionName, new NamedList<>(), ccc); + log.info("Cleaned up artifacts for failed create collection for [{}]", collectionName); + } catch (Exception cleanupFailure) { + if (cleanupFailure instanceof InterruptedException) { + Thread.currentThread().interrupt(); + } + log.error( + "Could not clean up after failed create collection for [{}]", + collectionName, + cleanupFailure); + cause.addSuppressed(cleanupFailure); } } diff --git a/solr/core/src/java/org/apache/solr/cluster/placement/impl/PlacementPluginAssignStrategy.java b/solr/core/src/java/org/apache/solr/cluster/placement/impl/PlacementPluginAssignStrategy.java index 13e2a0b53a4c..c3a10e6eddd5 100644 --- a/solr/core/src/java/org/apache/solr/cluster/placement/impl/PlacementPluginAssignStrategy.java +++ b/solr/core/src/java/org/apache/solr/cluster/placement/impl/PlacementPluginAssignStrategy.java @@ -27,6 +27,7 @@ import org.apache.solr.cloud.api.collections.Assign; import org.apache.solr.cluster.Node; import org.apache.solr.cluster.Replica.ReplicaType; +import org.apache.solr.cluster.SolrCollection; import org.apache.solr.cluster.placement.BalanceRequest; import org.apache.solr.cluster.placement.DeleteCollectionRequest; import org.apache.solr.cluster.placement.DeleteReplicasRequest; @@ -70,7 +71,7 @@ public List assign( placementRequests.add( PlacementRequestImpl.toPlacementRequest( placementContext.getCluster(), - placementContext.getCluster().getCollection(assignRequest.collectionName), + getCollection(placementContext, assignRequest.collectionName), assignRequest)); } @@ -171,6 +172,20 @@ public void verifyDeleteReplicas( } } + /** + * Looks up the collection replicas are requested for. It can be missing from the cluster state + * this request sees, for example when it was deleted concurrently. + */ + private static SolrCollection getCollection( + PlacementContext placementContext, String collectionName) throws IOException { + SolrCollection collection = placementContext.getCluster().getCollection(collectionName); + if (collection == null) { + throw new Assign.AssignmentException( + "Collection " + collectionName + " not found in cluster state; cannot assign replicas"); + } + return collection; + } + /** Very minimal placement logic for System collections */ private static List computeSystemCollectionPositions( PlacementContext placementContext, Assign.AssignRequest assignRequest) throws IOException { @@ -186,7 +201,7 @@ private static List computeSystemCollectionPositions( } PlacementRequestImpl request = new PlacementRequestImpl( - placementContext.getCluster().getCollection(assignRequest.collectionName), + getCollection(placementContext, assignRequest.collectionName), new HashSet<>(assignRequest.shardNames), nodes, assignRequest.numReplicas); diff --git a/solr/core/src/test/org/apache/solr/cloud/CreateCollectionCleanupTest.java b/solr/core/src/test/org/apache/solr/cloud/CreateCollectionCleanupTest.java index d5c342e36295..5e2b47c8f972 100644 --- a/solr/core/src/test/org/apache/solr/cloud/CreateCollectionCleanupTest.java +++ b/solr/core/src/test/org/apache/solr/cloud/CreateCollectionCleanupTest.java @@ -17,18 +17,50 @@ package org.apache.solr.cloud; +import static org.hamcrest.CoreMatchers.containsString; import static org.hamcrest.CoreMatchers.hasItem; import static org.hamcrest.CoreMatchers.is; import static org.hamcrest.CoreMatchers.not; import java.nio.file.Files; import java.nio.file.Path; +import java.util.Collection; +import java.util.HashMap; +import java.util.List; +import java.util.Map; import java.util.Properties; +import java.util.concurrent.ExecutorService; import org.apache.solr.client.solrj.RemoteSolrException; import org.apache.solr.client.solrj.impl.CloudSolrClient; import org.apache.solr.client.solrj.request.CollectionAdminRequest; import org.apache.solr.client.solrj.response.RequestStatusState; +import org.apache.solr.cloud.api.collections.AdminCmdContext; +import org.apache.solr.cloud.api.collections.CollectionCommandContext; +import org.apache.solr.cloud.api.collections.CreateCollectionCmd; +import org.apache.solr.cloud.api.collections.DistributedCollectionCommandContext; +import org.apache.solr.cluster.placement.BalancePlan; +import org.apache.solr.cluster.placement.BalanceRequest; +import org.apache.solr.cluster.placement.PlacementContext; +import org.apache.solr.cluster.placement.PlacementException; +import org.apache.solr.cluster.placement.PlacementPlan; +import org.apache.solr.cluster.placement.PlacementPlugin; +import org.apache.solr.cluster.placement.PlacementPluginFactory; +import org.apache.solr.cluster.placement.PlacementRequest; +import org.apache.solr.cluster.placement.impl.DelegatingPlacementPluginFactory; +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.ZkNodeProps; +import org.apache.solr.common.cloud.ZkStateReader; +import org.apache.solr.common.params.CollectionAdminParams; +import org.apache.solr.common.params.CollectionParams; import org.apache.solr.common.params.CoreAdminParams; +import org.apache.solr.common.util.ExecutorUtil; +import org.apache.solr.common.util.NamedList; +import org.apache.solr.common.util.SolrNamedThreadFactory; +import org.apache.solr.core.CoreContainer; +import org.apache.solr.embedded.JettySolrRunner; +import org.junit.After; import org.junit.BeforeClass; import org.junit.Test; @@ -59,6 +91,8 @@ public class CreateCollectionCleanupTest extends SolrCloudTestCase { + " \n" + "\n"; + private static final String PLACEMENT_FAILURE = "simulated placement failure"; + @BeforeClass public static void createCluster() throws Exception { configureCluster(1) @@ -68,6 +102,11 @@ public static void createCluster() throws Exception { .configure(); } + @After + public void restoreDefaultPlacement() { + setPlacementPluginFactory(null); + } + @Test public void testCreateCollectionCleanup() throws Exception { final CloudSolrClient cloudClient = cluster.getSolrClient(); @@ -128,4 +167,161 @@ public void testAsyncCreateCollectionCleanup() throws Exception { CollectionAdminRequest.listCollections(cloudClient), not(hasItem(collectionName))); } + + /** A placement plugin failure that is not a {@link PlacementException}, such as a plugin bug. */ + @Test + public void testCleanupAfterUnexpectedPlacementFailure() throws Exception { + final CloudSolrClient cloudClient = cluster.getSolrClient(); + String collectionName = "foo3"; + setPlacementPluginFactory(new FailingPlacementFactory(false)); + + CollectionAdminRequest.Create create = + CollectionAdminRequest.createCollection(collectionName, "conf1", 1, 1) + .setPerReplicaState(random().nextBoolean()); + RemoteSolrException e = + expectThrows(RemoteSolrException.class, () -> create.process(cloudClient)); + assertThat(e.getMessage(), containsString(PLACEMENT_FAILURE)); + assertNoCollection(collectionName); + + // nothing is left that would get in the way of creating the collection again + setPlacementPluginFactory(null); + CollectionAdminRequest.createCollection(collectionName, "conf1", 1, 1).process(cloudClient); + cluster.waitForActiveCollection(collectionName, 1, 1); + CollectionAdminRequest.deleteCollection(collectionName).process(cloudClient); + } + + /** A placement plugin that rejects the request the way plugins are expected to. */ + @Test + public void testCleanupAfterPlacementException() throws Exception { + final CloudSolrClient cloudClient = cluster.getSolrClient(); + String collectionName = "foo4"; + setPlacementPluginFactory(new FailingPlacementFactory(true)); + + CollectionAdminRequest.Create create = + CollectionAdminRequest.createCollection(collectionName, "conf1", 1, 1) + .setPerReplicaState(random().nextBoolean()); + RemoteSolrException e = + expectThrows(RemoteSolrException.class, () -> create.process(cloudClient)); + assertThat(e.getMessage(), containsString(PLACEMENT_FAILURE)); + assertNoCollection(collectionName); + } + + /** + * A create that names a collection which already exists must fail with BAD_REQUEST and must not + * delete that collection. The command is invoked directly with a cluster state that does not + * contain the collection, the view a node hands to the command when its local state lags behind + * ZooKeeper. The up-front check passes on that view, so only the state.json check in ZooKeeper + * keeps the create from going ahead and, on a later failure, deleting a collection it did not + * create. + */ + @Test + public void testCreateDoesNotDeleteExistingCollectionOnStaleView() throws Exception { + final CloudSolrClient cloudClient = cluster.getSolrClient(); + String collectionName = "existingColl"; + CollectionAdminRequest.createCollection(collectionName, "conf1", 1, 1).process(cloudClient); + cluster.waitForActiveCollection(collectionName, 1, 1); + try { + String collectionPath = DocCollection.getCollectionPath(collectionName); + byte[] stateBefore = zkClient().getData(collectionPath, null, null); + assertNotNull(stateBefore); + + // the stale view: same live nodes, but the collection is not in it + ClusterState currentState = cloudClient.getClusterState(); + Map otherCollections = new HashMap<>(); + for (String name : currentState.getCollectionNames()) { + if (!name.equals(collectionName)) { + otherCollections.put(name, currentState.getCollectionRef(name)); + } + } + ClusterState staleState = new ClusterState(otherCollections, currentState.getLiveNodes()); + + CoreContainer coreContainer = cluster.getJettySolrRunners().get(0).getCoreContainer(); + ExecutorService executor = + ExecutorUtil.newMDCAwareSingleThreadExecutor( + new SolrNamedThreadFactory("createCollectionCleanupTest")); + try { + CollectionCommandContext ccc = + new DistributedCollectionCommandContext(coreContainer, executor); + AdminCmdContext adminCmdContext = + new AdminCmdContext(CollectionParams.CollectionAction.CREATE) + .withClusterState(staleState); + ZkNodeProps message = + new ZkNodeProps( + "name", + collectionName, + CollectionAdminParams.COLL_CONF, + "conf1", + "numShards", + "1", + "nrtReplicas", + "1", + DocCollection.CollectionStateProps.PER_REPLICA_STATE, + "false"); + SolrException e = + expectThrows( + SolrException.class, + () -> + new CreateCollectionCmd(ccc).call(adminCmdContext, message, new NamedList<>())); + + // the existing collection is untouched: same state.json, still listed + assertArrayEquals(stateBefore, zkClient().getData(collectionPath, null, null)); + assertThat(CollectionAdminRequest.listCollections(cloudClient), hasItem(collectionName)); + assertEquals(SolrException.ErrorCode.BAD_REQUEST.code, e.code()); + assertThat(e.getMessage(), containsString("collection already exists: " + collectionName)); + } finally { + executor.shutdown(); + } + } finally { + CollectionAdminRequest.deleteCollection(collectionName).process(cloudClient); + } + } + + private static void assertNoCollection(String collectionName) throws Exception { + assertThat( + "Failed collection is still in the clusterstate: " + + cluster.getSolrClient().getClusterState().getCollectionOrNull(collectionName), + CollectionAdminRequest.listCollections(cluster.getSolrClient()), + not(hasItem(collectionName))); + assertFalse( + "Failed collection still has a node in ZooKeeper", + zkClient().exists(ZkStateReader.COLLECTIONS_ZKNODE + "/" + collectionName)); + } + + /** Sets the placement plugin factory of the node; {@code null} restores the default. */ + private static void setPlacementPluginFactory(PlacementPluginFactory factory) { + for (JettySolrRunner jetty : cluster.getJettySolrRunners()) { + ((DelegatingPlacementPluginFactory) jetty.getCoreContainer().getPlacementPluginFactory()) + .setDelegate(factory); + } + } + + private static class FailingPlacementFactory + implements PlacementPluginFactory { + private final boolean placementException; + + FailingPlacementFactory(boolean placementException) { + this.placementException = placementException; + } + + @Override + public PlacementPlugin createPluginInstance() { + return new PlacementPlugin() { + @Override + public List computePlacements( + Collection placementRequests, PlacementContext placementContext) + throws PlacementException { + if (placementException) { + throw new PlacementException(PLACEMENT_FAILURE); + } + throw new IllegalStateException(PLACEMENT_FAILURE); + } + + @Override + public BalancePlan computeBalancing( + BalanceRequest balanceRequest, PlacementContext placementContext) { + throw new UnsupportedOperationException(); + } + }; + } + } } diff --git a/solr/core/src/test/org/apache/solr/cloud/api/collections/CreateCollectionCmdRetryTest.java b/solr/core/src/test/org/apache/solr/cloud/api/collections/CreateCollectionCmdRetryTest.java new file mode 100644 index 000000000000..5aa0ef2b3300 --- /dev/null +++ b/solr/core/src/test/org/apache/solr/cloud/api/collections/CreateCollectionCmdRetryTest.java @@ -0,0 +1,94 @@ +/* + * 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 java.util.concurrent.atomic.AtomicInteger; +import org.apache.solr.SolrTestCase; +import org.apache.solr.common.SolrException; +import org.apache.solr.common.SolrException.ErrorCode; +import org.apache.solr.common.cloud.ZooKeeperException; + +/** Unit tests for the bounded retry {@link CreateCollectionCmd} puts around the alias write. */ +public class CreateCollectionCmdRetryTest extends SolrTestCase { + + public void testSucceedsAfterOneZkFailure() throws Exception { + AtomicInteger attempts = new AtomicInteger(); + CreateCollectionCmd.runWithBoundedRetries( + () -> { + if (attempts.incrementAndGet() < 2) { + throw new ZooKeeperException(ErrorCode.SERVER_ERROR, "transient"); + } + }, + 3, + 1); + assertEquals(2, attempts.get()); + } + + public void testZkFailurePropagatesAfterAllAttempts() { + AtomicInteger attempts = new AtomicInteger(); + ZooKeeperException thrown = + expectThrows( + ZooKeeperException.class, + () -> + CreateCollectionCmd.runWithBoundedRetries( + () -> { + attempts.incrementAndGet(); + throw new ZooKeeperException(ErrorCode.SERVER_ERROR, "still down"); + }, + 3, + 1)); + assertEquals("still down", thrown.getMessage()); + assertEquals(3, attempts.get()); + } + + public void testNonZkFailureIsNotRetried() { + AtomicInteger attempts = new AtomicInteger(); + expectThrows( + SolrException.class, + () -> + CreateCollectionCmd.runWithBoundedRetries( + () -> { + attempts.incrementAndGet(); + throw new SolrException(ErrorCode.SERVER_ERROR, "not a zk error"); + }, + 3, + 1)); + assertEquals(1, attempts.get()); + } + + public void testInterruptDuringPauseStopsRetrying() { + AtomicInteger attempts = new AtomicInteger(); + Thread.currentThread().interrupt(); + try { + expectThrows( + InterruptedException.class, + () -> + CreateCollectionCmd.runWithBoundedRetries( + () -> { + attempts.incrementAndGet(); + throw new ZooKeeperException(ErrorCode.SERVER_ERROR, "transient"); + }, + 3, + 60_000)); + assertEquals(1, attempts.get()); + assertTrue(Thread.currentThread().isInterrupted()); + } finally { + // do not leak the interrupt flag into the rest of the suite + Thread.interrupted(); + } + } +} diff --git a/solr/core/src/test/org/apache/solr/cluster/placement/impl/PlacementPluginIntegrationTest.java b/solr/core/src/test/org/apache/solr/cluster/placement/impl/PlacementPluginIntegrationTest.java index 2fb85423006f..bc09502bf9f9 100644 --- a/solr/core/src/test/org/apache/solr/cluster/placement/impl/PlacementPluginIntegrationTest.java +++ b/solr/core/src/test/org/apache/solr/cluster/placement/impl/PlacementPluginIntegrationTest.java @@ -22,6 +22,7 @@ import java.util.Arrays; import java.util.HashMap; import java.util.HashSet; +import java.util.List; import java.util.Map; import java.util.Optional; import java.util.Set; @@ -38,6 +39,7 @@ import org.apache.solr.client.solrj.response.V2Response; import org.apache.solr.cloud.MiniSolrCloudCluster; import org.apache.solr.cloud.SolrCloudTestCase; +import org.apache.solr.cloud.api.collections.Assign; import org.apache.solr.cluster.Cluster; import org.apache.solr.cluster.Node; import org.apache.solr.cluster.SolrCollection; @@ -55,6 +57,8 @@ import org.apache.solr.cluster.placement.plugins.SimplePlacementFactory; 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.util.RetryUtil; import org.apache.solr.core.CoreContainer; import org.apache.solr.util.LogLevel; @@ -438,6 +442,23 @@ public void testNodeTypeIntegration() throws Exception { System.clearProperty(AffinityPlacementConfig.NODE_TYPE_SYSPROP); } + /** Replica assignment for a collection that does not exist, as a concurrent delete can cause. */ + @Test + public void testAssignForMissingCollection() { + Assign.AssignRequest request = + new Assign.AssignRequestBuilder() + .forCollection("missingCollection") + .forShard(List.of("shard1")) + .assignReplicas(ReplicaCount.of(Replica.Type.NRT, 1)) + .build(); + + Assign.AssignmentException e = + expectThrows( + Assign.AssignmentException.class, + () -> Assign.createAssignStrategy(cc).assign(cloudManager, request)); + assertTrue(e.getMessage(), e.getMessage().contains("missingCollection")); + } + @Test public void testAttributeFetcherImpl() throws Exception { CollectionAdminResponse rsp =