From 4d3550c4198bb8092ac99e7a33d169caa508b827 Mon Sep 17 00:00:00 2001 From: David Turner Date: Mon, 19 Mar 2018 15:09:22 +0000 Subject: [PATCH 01/23] Avoid loading shard metadata while closing If `ShardStateMetaData.FORMAT.loadLatestState` is called while a shard is closing, the shard metadata directory may be deleted after its existence has been checked but before the Lucene `Directory` has been created. When the `Directory` is created, the just-deleted directory is brought back into existence. There are three places where `loadLatestState` is called in a manner that leaves it open to this race. This change ensures that these calls occur either under a `ShardLock` or else while holding a reference to the existing `Store`. In either case, this protects the shard metadata directory from concurrent deletion. Cf #19338, #21463, #25335 and https://issues.apache.org/jira/browse/LUCENE-7375 --- ...ransportNodesListGatewayStartedShards.java | 31 +++++++++++++++++-- .../TransportNodesListShardStoreMetaData.java | 6 +++- 2 files changed, 33 insertions(+), 4 deletions(-) diff --git a/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java b/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java index 11df875d4dd99..58135e71c26b4 100644 --- a/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java +++ b/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java @@ -42,7 +42,10 @@ import org.elasticsearch.common.settings.Settings; import org.elasticsearch.common.xcontent.NamedXContentRegistry; import org.elasticsearch.env.NodeEnvironment; +import org.elasticsearch.env.ShardLock; +import org.elasticsearch.index.IndexService; import org.elasticsearch.index.IndexSettings; +import org.elasticsearch.index.shard.IndexShard; import org.elasticsearch.index.shard.ShardId; import org.elasticsearch.index.shard.ShardPath; import org.elasticsearch.index.shard.ShardStateMetaData; @@ -53,6 +56,7 @@ import java.io.IOException; import java.util.List; +import java.util.concurrent.TimeUnit; /** * This transport action is used to fetch the shard version from each node during primary allocation in {@link GatewayAllocator}. @@ -112,13 +116,32 @@ protected NodesGatewayStartedShards newResponse(Request request, return new NodesGatewayStartedShards(clusterService.getClusterName(), responses, failures); } + private AutoCloseable shardStoreReference(ShardId shardId) { + final IndexService indexService = indicesService.indexService(shardId.getIndex()); + if (indexService != null) { + final IndexShard indexShard = indexService.getShardOrNull(shardId.getId()); + if (indexShard != null) { + final Store store = indexShard.store(); + if (store.tryIncRef()) { + return store::decRef; + } + } + } + + return nodeEnv.shardLock(shardId, TimeUnit.SECONDS.toMillis(5)); + } + @Override protected NodeGatewayStartedShards nodeOperation(NodeRequest request) { try { final ShardId shardId = request.getShardId(); logger.trace("{} loading local shard state info", shardId); - ShardStateMetaData shardStateMetaData = ShardStateMetaData.FORMAT.loadLatestState(logger, NamedXContentRegistry.EMPTY, - nodeEnv.availableShardPaths(request.shardId)); + + ShardStateMetaData shardStateMetaData; + try (AutoCloseable ignored = shardStoreReference(shardId)) { + shardStateMetaData = ShardStateMetaData.FORMAT.loadLatestState(logger, NamedXContentRegistry.EMPTY, + nodeEnv.availableShardPaths(request.shardId)); + } if (shardStateMetaData != null) { IndexMetaData metaData = clusterService.state().metaData().index(shardId.getIndex()); if (metaData == null) { @@ -139,7 +162,9 @@ protected NodeGatewayStartedShards nodeOperation(NodeRequest request) { ShardPath shardPath = null; try { IndexSettings indexSettings = new IndexSettings(metaData, settings); - shardPath = ShardPath.loadShardPath(logger, nodeEnv, shardId, indexSettings); + try (ShardLock ignored = nodeEnv.shardLock(shardId, TimeUnit.SECONDS.toMillis(5))) { + shardPath = ShardPath.loadShardPath(logger, nodeEnv, shardId, indexSettings); + } if (shardPath == null) { throw new IllegalStateException(shardId + " no shard path found"); } diff --git a/server/src/main/java/org/elasticsearch/indices/store/TransportNodesListShardStoreMetaData.java b/server/src/main/java/org/elasticsearch/indices/store/TransportNodesListShardStoreMetaData.java index 404a19b0ab354..d91bba2ada929 100644 --- a/server/src/main/java/org/elasticsearch/indices/store/TransportNodesListShardStoreMetaData.java +++ b/server/src/main/java/org/elasticsearch/indices/store/TransportNodesListShardStoreMetaData.java @@ -41,6 +41,7 @@ import org.elasticsearch.common.unit.TimeValue; import org.elasticsearch.common.xcontent.NamedXContentRegistry; import org.elasticsearch.env.NodeEnvironment; +import org.elasticsearch.env.ShardLock; import org.elasticsearch.gateway.AsyncShardFetch; import org.elasticsearch.index.IndexService; import org.elasticsearch.index.IndexSettings; @@ -139,7 +140,10 @@ private StoreFilesMetaData listStoreMetaData(ShardId shardId) throws IOException return new StoreFilesMetaData(shardId, Store.MetadataSnapshot.EMPTY); } final IndexSettings indexSettings = indexService != null ? indexService.getIndexSettings() : new IndexSettings(metaData, settings); - final ShardPath shardPath = ShardPath.loadShardPath(logger, nodeEnv, shardId, indexSettings); + final ShardPath shardPath; + try (ShardLock ignored = nodeEnv.shardLock(shardId, TimeUnit.SECONDS.toMillis(5))) { + shardPath = ShardPath.loadShardPath(logger, nodeEnv, shardId, indexSettings); + } if (shardPath == null) { return new StoreFilesMetaData(shardId, Store.MetadataSnapshot.EMPTY); } From 487e7852a6a99b581f578e36c8f1586c5f36c1dd Mon Sep 17 00:00:00 2001 From: David Turner Date: Tue, 27 Mar 2018 13:46:34 +0100 Subject: [PATCH 02/23] Use IndexShard's mutex to protect the call to loadLatestState --- ...ransportNodesListGatewayStartedShards.java | 30 +++++++++---------- .../elasticsearch/index/shard/IndexShard.java | 13 ++++++++ 2 files changed, 27 insertions(+), 16 deletions(-) diff --git a/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java b/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java index 01ec7111efc58..1b7edeb05a127 100644 --- a/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java +++ b/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java @@ -20,6 +20,7 @@ package org.elasticsearch.gateway; import org.apache.logging.log4j.message.ParameterizedMessage; +import org.apache.lucene.store.AlreadyClosedException; import org.elasticsearch.ElasticsearchException; import org.elasticsearch.Version; import org.elasticsearch.action.ActionListener; @@ -42,7 +43,6 @@ import org.elasticsearch.common.xcontent.NamedXContentRegistry; import org.elasticsearch.env.NodeEnvironment; import org.elasticsearch.env.ShardLock; -import org.elasticsearch.index.IndexService; import org.elasticsearch.index.IndexSettings; import org.elasticsearch.index.shard.IndexShard; import org.elasticsearch.index.shard.ShardId; @@ -115,19 +115,21 @@ protected NodesGatewayStartedShards newResponse(Request request, return new NodesGatewayStartedShards(clusterService.getClusterName(), responses, failures); } - private AutoCloseable shardStoreReference(ShardId shardId) { - final IndexService indexService = indicesService.indexService(shardId.getIndex()); - if (indexService != null) { - final IndexShard indexShard = indexService.getShardOrNull(shardId.getId()); - if (indexShard != null) { - final Store store = indexShard.store(); - if (store.tryIncRef()) { - return store::decRef; - } + private ShardStateMetaData safelyLoadLatestState(ShardId shardId) throws Exception { + final IndexShard indexShard = indicesService.getShardOrNull(shardId); + if (indexShard != null) { + try { + return indexShard.loadShardStateMetaDataIfOpen(NamedXContentRegistry.EMPTY, nodeEnv.availableShardPaths(shardId)); + } + catch (AlreadyClosedException ignored) { + // Ok, fall through to trying to get the shard lock ourselves. } } - return nodeEnv.shardLock(shardId, TimeUnit.SECONDS.toMillis(5)); + try (ShardLock ignored = nodeEnv.shardLock(shardId, TimeUnit.SECONDS.toMillis(5))) { + return ShardStateMetaData.FORMAT.loadLatestState(logger, NamedXContentRegistry.EMPTY, + nodeEnv.availableShardPaths(shardId)); + } } @Override @@ -136,11 +138,7 @@ protected NodeGatewayStartedShards nodeOperation(NodeRequest request) { final ShardId shardId = request.getShardId(); logger.trace("{} loading local shard state info", shardId); - ShardStateMetaData shardStateMetaData; - try (AutoCloseable ignored = shardStoreReference(shardId)) { - shardStateMetaData = ShardStateMetaData.FORMAT.loadLatestState(logger, NamedXContentRegistry.EMPTY, - nodeEnv.availableShardPaths(request.shardId)); - } + ShardStateMetaData shardStateMetaData = safelyLoadLatestState(shardId); if (shardStateMetaData != null) { IndexMetaData metaData = clusterService.state().metaData().index(shardId.getIndex()); if (metaData == null) { diff --git a/server/src/main/java/org/elasticsearch/index/shard/IndexShard.java b/server/src/main/java/org/elasticsearch/index/shard/IndexShard.java index 30f813e86e234..c99e0fe498896 100644 --- a/server/src/main/java/org/elasticsearch/index/shard/IndexShard.java +++ b/server/src/main/java/org/elasticsearch/index/shard/IndexShard.java @@ -64,6 +64,7 @@ import org.elasticsearch.common.util.BigArrays; import org.elasticsearch.common.util.concurrent.AbstractRunnable; import org.elasticsearch.common.util.concurrent.AsyncIOProcessor; +import org.elasticsearch.common.xcontent.NamedXContentRegistry; import org.elasticsearch.common.xcontent.XContentHelper; import org.elasticsearch.core.internal.io.IOUtils; import org.elasticsearch.index.Index; @@ -139,6 +140,7 @@ import java.io.PrintStream; import java.nio.channels.ClosedByInterruptException; import java.nio.charset.StandardCharsets; +import java.nio.file.Path; import java.util.ArrayList; import java.util.Collections; import java.util.EnumSet; @@ -2067,6 +2069,17 @@ public void startRecovery(RecoveryState recoveryState, PeerRecoveryTargetService } } + public ShardStateMetaData loadShardStateMetaDataIfOpen(NamedXContentRegistry namedXContentRegistry, Path[] dataLocations) + throws IOException { + synchronized (mutex) { + if (state == IndexShardState.CLOSED) { + throw new AlreadyClosedException(shardId + " can't load shard state metadata - shard is closed"); + } + + return ShardStateMetaData.FORMAT.loadLatestState(logger, namedXContentRegistry, dataLocations); + } + } + class ShardEventListener implements Engine.EventListener { private final CopyOnWriteArrayList> delegates = new CopyOnWriteArrayList<>(); From fb5f8c967ccbc750f86ae43685babd165fd1703c Mon Sep 17 00:00:00 2001 From: David Turner Date: Wed, 4 Apr 2018 16:46:23 +0100 Subject: [PATCH 03/23] Add test --- .../gateway/RecoveryFromGatewayIT.java | 75 +++++++++++++++++++ 1 file changed, 75 insertions(+) diff --git a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java index 154d702e7fb77..f50ee79fd886a 100644 --- a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java +++ b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java @@ -20,6 +20,7 @@ package org.elasticsearch.gateway; import com.carrotsearch.hppc.cursors.ObjectCursor; +import org.elasticsearch.action.admin.indices.get.GetIndexResponse; import org.elasticsearch.action.admin.indices.recovery.RecoveryResponse; import org.elasticsearch.action.admin.indices.stats.IndexStats; import org.elasticsearch.action.admin.indices.stats.ShardStats; @@ -46,6 +47,7 @@ import org.elasticsearch.test.ESIntegTestCase.Scope; import org.elasticsearch.test.InternalTestCluster; import org.elasticsearch.test.InternalTestCluster.RestartCallback; +import org.elasticsearch.test.junit.annotations.TestLogging; import org.elasticsearch.test.store.MockFSIndexStore; import java.nio.file.DirectoryStream; @@ -58,11 +60,13 @@ import java.util.HashSet; import java.util.Map; import java.util.Set; +import java.util.concurrent.ExecutionException; import java.util.stream.IntStream; import static org.elasticsearch.cluster.metadata.IndexMetaData.SETTING_NUMBER_OF_REPLICAS; import static org.elasticsearch.cluster.metadata.IndexMetaData.SETTING_NUMBER_OF_SHARDS; import static org.elasticsearch.common.xcontent.XContentFactory.jsonBuilder; +import static org.elasticsearch.index.query.QueryBuilders.existsQuery; import static org.elasticsearch.index.query.QueryBuilders.matchAllQuery; import static org.elasticsearch.index.query.QueryBuilders.termQuery; import static org.elasticsearch.test.hamcrest.ElasticsearchAssertions.assertAcked; @@ -567,4 +571,75 @@ public Settings onNodeStopped(String nodeName) throws Exception { // start another node so cluster consistency checks won't time out due to the lack of state internalCluster().startNode(); } + + @TestLogging("org.elasticsearch.env.NodeEnvironment:TRACE,org.elasticsearch.gateway.TransportNodesListGatewayStartedShards:TRACE,org.elasticsearch.gateway.MetaDataStateFormat:TRACE,org.elasticsearch.index.shard.IndexShard:TRACE") + public void testLoadLatestStateOnClosingShard() throws Exception { + + final String nodeName = internalCluster().startNode(); + DiscoveryNode node = internalCluster().getInstance(ClusterService.class, nodeName).localNode(); + + assertAcked(prepareCreate("test").setSettings(Settings.builder() + .put(SETTING_NUMBER_OF_SHARDS, 1).put(SETTING_NUMBER_OF_REPLICAS, 0))); + final ShardId shardId = new ShardId(resolveIndex("test"), 0); + + try (AutoCloseable ignored = new ConcurrentShardLister(shardId, node, 4).start()) { + boolean indexExists; + do { + logger.info("--> deleting index"); + assertAcked(client().admin().indices().prepareDelete("test")); + logger.info("--> checking whether index still exists"); + GetIndexResponse getIndexResponse = client().admin().indices().prepareGetIndex().setIndices("_all").get(); + indexExists = getIndexResponse.getSettings().containsKey("test"); + if (indexExists) { + logger.info("--> index still exists, retrying"); + } + } while (indexExists); + logger.info("--> index no longer exists, exiting"); + } + logger.info("--> test done"); + } + + private class ConcurrentShardLister { + private final ShardId shardId; + private final DiscoveryNode node; + private volatile boolean shouldExit = false; + private Exception listerException; + private final int concurrentListingThreadCount; + + ConcurrentShardLister(ShardId shardId, DiscoveryNode node, int concurrentListingThreadCount) { + this.shardId = shardId; + this.node = node; + this.concurrentListingThreadCount = concurrentListingThreadCount; + } + + public AutoCloseable start() { + final Thread[] listingThreads = new Thread[concurrentListingThreadCount]; + for (int i = 0; i < concurrentListingThreadCount; i++) { + listingThreads[i] = new Thread(() -> { + for (int iterations = 0; iterations < 100 && shouldExit == false; iterations++) { + try { + internalCluster().getInstance(TransportNodesListGatewayStartedShards.class) + .execute(new TransportNodesListGatewayStartedShards.Request(shardId, new DiscoveryNode[]{node})) + .get(); + } catch (InterruptedException | ExecutionException exception) { + shouldExit = true; + listerException = exception; + } + } + }); + listingThreads[i].setName("ConcurrentShardLister[" + i + "]"); + listingThreads[i].start(); + } + + return () -> { + shouldExit = true; + for (int i = 0; i < concurrentListingThreadCount; i++) { + listingThreads[i].join(); + } + if (listerException != null) { + throw listerException; + } + }; + } + } } From fed85bfcb2df4efdf0487d175e40370e6fa90764 Mon Sep 17 00:00:00 2001 From: David Turner Date: Thu, 5 Apr 2018 08:07:04 +0100 Subject: [PATCH 04/23] Much simpler test --- .../gateway/RecoveryFromGatewayIT.java | 74 ++++--------------- 1 file changed, 15 insertions(+), 59 deletions(-) diff --git a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java index f50ee79fd886a..5c39280857712 100644 --- a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java +++ b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java @@ -60,6 +60,7 @@ import java.util.HashSet; import java.util.Map; import java.util.Set; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutionException; import java.util.stream.IntStream; @@ -573,8 +574,7 @@ public Settings onNodeStopped(String nodeName) throws Exception { } @TestLogging("org.elasticsearch.env.NodeEnvironment:TRACE,org.elasticsearch.gateway.TransportNodesListGatewayStartedShards:TRACE,org.elasticsearch.gateway.MetaDataStateFormat:TRACE,org.elasticsearch.index.shard.IndexShard:TRACE") - public void testLoadLatestStateOnClosingShard() throws Exception { - + public void testLoadLatestStateWhileClosingShardDoesNotResurrectMetadataDirectory() throws Exception { final String nodeName = internalCluster().startNode(); DiscoveryNode node = internalCluster().getInstance(ClusterService.class, nodeName).localNode(); @@ -582,64 +582,20 @@ public void testLoadLatestStateOnClosingShard() throws Exception { .put(SETTING_NUMBER_OF_SHARDS, 1).put(SETTING_NUMBER_OF_REPLICAS, 0))); final ShardId shardId = new ShardId(resolveIndex("test"), 0); - try (AutoCloseable ignored = new ConcurrentShardLister(shardId, node, 4).start()) { - boolean indexExists; - do { - logger.info("--> deleting index"); - assertAcked(client().admin().indices().prepareDelete("test")); - logger.info("--> checking whether index still exists"); - GetIndexResponse getIndexResponse = client().admin().indices().prepareGetIndex().setIndices("_all").get(); - indexExists = getIndexResponse.getSettings().containsKey("test"); - if (indexExists) { - logger.info("--> index still exists, retrying"); - } - } while (indexExists); - logger.info("--> index no longer exists, exiting"); - } - logger.info("--> test done"); - } - - private class ConcurrentShardLister { - private final ShardId shardId; - private final DiscoveryNode node; - private volatile boolean shouldExit = false; - private Exception listerException; - private final int concurrentListingThreadCount; - - ConcurrentShardLister(ShardId shardId, DiscoveryNode node, int concurrentListingThreadCount) { - this.shardId = shardId; - this.node = node; - this.concurrentListingThreadCount = concurrentListingThreadCount; - } - - public AutoCloseable start() { - final Thread[] listingThreads = new Thread[concurrentListingThreadCount]; - for (int i = 0; i < concurrentListingThreadCount; i++) { - listingThreads[i] = new Thread(() -> { - for (int iterations = 0; iterations < 100 && shouldExit == false; iterations++) { - try { - internalCluster().getInstance(TransportNodesListGatewayStartedShards.class) - .execute(new TransportNodesListGatewayStartedShards.Request(shardId, new DiscoveryNode[]{node})) - .get(); - } catch (InterruptedException | ExecutionException exception) { - shouldExit = true; - listerException = exception; - } - } - }); - listingThreads[i].setName("ConcurrentShardLister[" + i + "]"); - listingThreads[i].start(); + Thread listingThread = new Thread(() -> { + try { + internalCluster().getInstance(TransportNodesListGatewayStartedShards.class) + .execute(new TransportNodesListGatewayStartedShards.Request(shardId, new DiscoveryNode[]{node})) + .get(); + } catch (InterruptedException | ExecutionException exception) { + // don't care if this fails } + }); + listingThread.start(); + assertAcked(client().admin().indices().prepareDelete("test")); + // Asserts that the shard is really deleted - at one time, the concurrent TransportNodesListGatewayStartedShards request + // might have prevented this. - return () -> { - shouldExit = true; - for (int i = 0; i < concurrentListingThreadCount; i++) { - listingThreads[i].join(); - } - if (listerException != null) { - throw listerException; - } - }; - } + listingThread.join(); } } From 0ce802111e043f1ad729e5a2e785d5a75554cbbe Mon Sep 17 00:00:00 2001 From: David Turner Date: Thu, 5 Apr 2018 08:08:32 +0100 Subject: [PATCH 05/23] Tidy imports --- .../java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java | 3 --- 1 file changed, 3 deletions(-) diff --git a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java index 5c39280857712..7a9d54ee4a370 100644 --- a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java +++ b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java @@ -20,7 +20,6 @@ package org.elasticsearch.gateway; import com.carrotsearch.hppc.cursors.ObjectCursor; -import org.elasticsearch.action.admin.indices.get.GetIndexResponse; import org.elasticsearch.action.admin.indices.recovery.RecoveryResponse; import org.elasticsearch.action.admin.indices.stats.IndexStats; import org.elasticsearch.action.admin.indices.stats.ShardStats; @@ -60,14 +59,12 @@ import java.util.HashSet; import java.util.Map; import java.util.Set; -import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutionException; import java.util.stream.IntStream; import static org.elasticsearch.cluster.metadata.IndexMetaData.SETTING_NUMBER_OF_REPLICAS; import static org.elasticsearch.cluster.metadata.IndexMetaData.SETTING_NUMBER_OF_SHARDS; import static org.elasticsearch.common.xcontent.XContentFactory.jsonBuilder; -import static org.elasticsearch.index.query.QueryBuilders.existsQuery; import static org.elasticsearch.index.query.QueryBuilders.matchAllQuery; import static org.elasticsearch.index.query.QueryBuilders.termQuery; import static org.elasticsearch.test.hamcrest.ElasticsearchAssertions.assertAcked; From b654d9de900c7774646dcf70e6950355a66618fa Mon Sep 17 00:00:00 2001 From: David Turner Date: Thu, 5 Apr 2018 08:13:34 +0100 Subject: [PATCH 06/23] Try repeats --- .../gateway/RecoveryFromGatewayIT.java | 14 ++++++++------ 1 file changed, 8 insertions(+), 6 deletions(-) diff --git a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java index 7a9d54ee4a370..ff04fe65cc76c 100644 --- a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java +++ b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java @@ -580,12 +580,14 @@ public void testLoadLatestStateWhileClosingShardDoesNotResurrectMetadataDirector final ShardId shardId = new ShardId(resolveIndex("test"), 0); Thread listingThread = new Thread(() -> { - try { - internalCluster().getInstance(TransportNodesListGatewayStartedShards.class) - .execute(new TransportNodesListGatewayStartedShards.Request(shardId, new DiscoveryNode[]{node})) - .get(); - } catch (InterruptedException | ExecutionException exception) { - // don't care if this fails + for (int i = 0; i < 10; i++) { + try { + internalCluster().getInstance(TransportNodesListGatewayStartedShards.class) + .execute(new TransportNodesListGatewayStartedShards.Request(shardId, new DiscoveryNode[]{node})) + .get(); + } catch (InterruptedException | ExecutionException exception) { + // don't care if this fails + } } }); listingThread.start(); From e9ef547f5dde79235648de80d6ca386280fdc065 Mon Sep 17 00:00:00 2001 From: David Turner Date: Thu, 5 Apr 2018 09:57:12 +0200 Subject: [PATCH 07/23] Try without a loop --- .../gateway/RecoveryFromGatewayIT.java | 14 ++++++-------- 1 file changed, 6 insertions(+), 8 deletions(-) diff --git a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java index ff04fe65cc76c..7a9d54ee4a370 100644 --- a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java +++ b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java @@ -580,14 +580,12 @@ public void testLoadLatestStateWhileClosingShardDoesNotResurrectMetadataDirector final ShardId shardId = new ShardId(resolveIndex("test"), 0); Thread listingThread = new Thread(() -> { - for (int i = 0; i < 10; i++) { - try { - internalCluster().getInstance(TransportNodesListGatewayStartedShards.class) - .execute(new TransportNodesListGatewayStartedShards.Request(shardId, new DiscoveryNode[]{node})) - .get(); - } catch (InterruptedException | ExecutionException exception) { - // don't care if this fails - } + try { + internalCluster().getInstance(TransportNodesListGatewayStartedShards.class) + .execute(new TransportNodesListGatewayStartedShards.Request(shardId, new DiscoveryNode[]{node})) + .get(); + } catch (InterruptedException | ExecutionException exception) { + // don't care if this fails } }); listingThread.start(); From 7641ac251aa64084e3f7e6dc9b03466e89bc3070 Mon Sep 17 00:00:00 2001 From: David Turner Date: Thu, 5 Apr 2018 09:09:37 +0100 Subject: [PATCH 08/23] Multiple threads, one request each, synchronised starts --- .../gateway/RecoveryFromGatewayIT.java | 40 ++++++++++++++----- 1 file changed, 29 insertions(+), 11 deletions(-) diff --git a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java index 7a9d54ee4a370..20d0f39034308 100644 --- a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java +++ b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java @@ -59,6 +59,7 @@ import java.util.HashSet; import java.util.Map; import java.util.Set; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutionException; import java.util.stream.IntStream; @@ -579,20 +580,37 @@ public void testLoadLatestStateWhileClosingShardDoesNotResurrectMetadataDirector .put(SETTING_NUMBER_OF_SHARDS, 1).put(SETTING_NUMBER_OF_REPLICAS, 0))); final ShardId shardId = new ShardId(resolveIndex("test"), 0); - Thread listingThread = new Thread(() -> { - try { - internalCluster().getInstance(TransportNodesListGatewayStartedShards.class) - .execute(new TransportNodesListGatewayStartedShards.Request(shardId, new DiscoveryNode[]{node})) - .get(); - } catch (InterruptedException | ExecutionException exception) { - // don't care if this fails - } - }); - listingThread.start(); + final int listingThreadCount = 4; + final CountDownLatch countDownLatch = new CountDownLatch(listingThreadCount + 1); + + Thread listingThreads[] = new Thread[listingThreadCount]; + for (int threadIndex = 0; threadIndex < listingThreadCount; threadIndex++) { + listingThreads[threadIndex] = new Thread(() -> { + try { + countDownLatch.countDown(); + countDownLatch.await(); + internalCluster().getInstance(TransportNodesListGatewayStartedShards.class) + .execute(new TransportNodesListGatewayStartedShards.Request(shardId, new DiscoveryNode[]{node})) + .get(); + } catch (InterruptedException | ExecutionException exception) { + // don't care if this fails + } + }); + } + + for (final Thread listingThread : listingThreads) { + listingThread.start(); + } + + countDownLatch.countDown(); + countDownLatch.await(); + assertAcked(client().admin().indices().prepareDelete("test")); // Asserts that the shard is really deleted - at one time, the concurrent TransportNodesListGatewayStartedShards request // might have prevented this. - listingThread.join(); + for (final Thread listingThread : listingThreads) { + listingThread.join(); + } } } From 62ae05c3c2f03d23f0b8ef7c6ed66e9c9a6c0425 Mon Sep 17 00:00:00 2001 From: David Turner Date: Thu, 5 Apr 2018 09:24:46 +0100 Subject: [PATCH 09/23] Try doing listing and deletions all in the background, with multiple deletions --- .../gateway/RecoveryFromGatewayIT.java | 39 ++++++++++--------- 1 file changed, 21 insertions(+), 18 deletions(-) diff --git a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java index 20d0f39034308..829c0162b8782 100644 --- a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java +++ b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java @@ -580,36 +580,39 @@ public void testLoadLatestStateWhileClosingShardDoesNotResurrectMetadataDirector .put(SETTING_NUMBER_OF_SHARDS, 1).put(SETTING_NUMBER_OF_REPLICAS, 0))); final ShardId shardId = new ShardId(resolveIndex("test"), 0); - final int listingThreadCount = 4; - final CountDownLatch countDownLatch = new CountDownLatch(listingThreadCount + 1); - - Thread listingThreads[] = new Thread[listingThreadCount]; - for (int threadIndex = 0; threadIndex < listingThreadCount; threadIndex++) { - listingThreads[threadIndex] = new Thread(() -> { + final int listingThreadCount = 2; + final int deletingThreadCount = 2; + final CountDownLatch countDownLatch = new CountDownLatch(listingThreadCount + deletingThreadCount); + + Thread threads[] = new Thread[listingThreadCount + deletingThreadCount]; + for (int threadIndex = 0; threadIndex < listingThreadCount + deletingThreadCount; threadIndex++) { + final boolean isListingThread = threadIndex < listingThreadCount; + threads[threadIndex] = new Thread(() -> { try { countDownLatch.countDown(); countDownLatch.await(); - internalCluster().getInstance(TransportNodesListGatewayStartedShards.class) - .execute(new TransportNodesListGatewayStartedShards.Request(shardId, new DiscoveryNode[]{node})) - .get(); + + if (isListingThread) { + internalCluster().getInstance(TransportNodesListGatewayStartedShards.class) + .execute(new TransportNodesListGatewayStartedShards.Request(shardId, new DiscoveryNode[]{node})) + .get(); + } else { + assertAcked(client().admin().indices().prepareDelete("test")); + } } catch (InterruptedException | ExecutionException exception) { // don't care if this fails } - }); + }, (isListingThread ? "Listing" : "Deleting") + "[" + threadIndex + "]"); } - for (final Thread listingThread : listingThreads) { + for (final Thread listingThread : threads) { listingThread.start(); } - countDownLatch.countDown(); - countDownLatch.await(); - - assertAcked(client().admin().indices().prepareDelete("test")); - // Asserts that the shard is really deleted - at one time, the concurrent TransportNodesListGatewayStartedShards request - // might have prevented this. + // Deleting an index asserts that it really is gone from disk - at one time, the concurrent + // TransportNodesListGatewayStartedShards requests sometimes prevented this. - for (final Thread listingThread : listingThreads) { + for (final Thread listingThread : threads) { listingThread.join(); } } From ed371747bc9db13428e372ea008e83081eb87609 Mon Sep 17 00:00:00 2001 From: David Turner Date: Thu, 5 Apr 2018 09:29:42 +0100 Subject: [PATCH 10/23] Also don't care if the index is not found --- .../java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java index 829c0162b8782..6d9745e0057a4 100644 --- a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java +++ b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java @@ -34,6 +34,7 @@ import org.elasticsearch.env.Environment; import org.elasticsearch.env.NodeEnvironment; import org.elasticsearch.index.Index; +import org.elasticsearch.index.IndexNotFoundException; import org.elasticsearch.index.IndexSettings; import org.elasticsearch.index.engine.Engine; import org.elasticsearch.index.query.QueryBuilders; @@ -599,7 +600,7 @@ public void testLoadLatestStateWhileClosingShardDoesNotResurrectMetadataDirector } else { assertAcked(client().admin().indices().prepareDelete("test")); } - } catch (InterruptedException | ExecutionException exception) { + } catch (InterruptedException | ExecutionException | IndexNotFoundException ignored) { // don't care if this fails } }, (isListingThread ? "Listing" : "Deleting") + "[" + threadIndex + "]"); From be76bb0c6e90c075ed334007a4ad4fac3c47d494 Mon Sep 17 00:00:00 2001 From: David Turner Date: Thu, 5 Apr 2018 09:52:33 +0100 Subject: [PATCH 11/23] Add comments --- .../elasticsearch/gateway/RecoveryFromGatewayIT.java | 11 ++++++++--- 1 file changed, 8 insertions(+), 3 deletions(-) diff --git a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java index 6d9745e0057a4..ffc6847933389 100644 --- a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java +++ b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java @@ -572,8 +572,13 @@ public Settings onNodeStopped(String nodeName) throws Exception { internalCluster().startNode(); } - @TestLogging("org.elasticsearch.env.NodeEnvironment:TRACE,org.elasticsearch.gateway.TransportNodesListGatewayStartedShards:TRACE,org.elasticsearch.gateway.MetaDataStateFormat:TRACE,org.elasticsearch.index.shard.IndexShard:TRACE") public void testLoadLatestStateWhileClosingShardDoesNotResurrectMetadataDirectory() throws Exception { + + // This test pertains to a race condition in which deleting a shard concurrently with a TransportNodesListGatewayStartedShards + // request could resurrect the shard's metadata folder after it was deleted. Here we try and recreate the race, but it is quite + // delicate so this test does not always fail. Experimentation showed that setting the thread counts as below would yield a failure + // after a reasonable number of iterations: running repeatedly with -Dtests.iters=1000 saw 6 failures out of 10 runs. + final String nodeName = internalCluster().startNode(); DiscoveryNode node = internalCluster().getInstance(ClusterService.class, nodeName).localNode(); @@ -583,6 +588,7 @@ public void testLoadLatestStateWhileClosingShardDoesNotResurrectMetadataDirector final int listingThreadCount = 2; final int deletingThreadCount = 2; + final CountDownLatch countDownLatch = new CountDownLatch(listingThreadCount + deletingThreadCount); Thread threads[] = new Thread[listingThreadCount + deletingThreadCount]; @@ -610,8 +616,7 @@ public void testLoadLatestStateWhileClosingShardDoesNotResurrectMetadataDirector listingThread.start(); } - // Deleting an index asserts that it really is gone from disk - at one time, the concurrent - // TransportNodesListGatewayStartedShards requests sometimes prevented this. + // Deleting an index asserts that it really is gone from disk, so no other assertions are necessary here. for (final Thread listingThread : threads) { listingThread.join(); From d534ffd14b4c2a083b39edf58e17ce467785a4db Mon Sep 17 00:00:00 2001 From: David Turner Date: Thu, 5 Apr 2018 10:23:04 +0100 Subject: [PATCH 12/23] Reinstate logging --- .../java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java | 2 ++ 1 file changed, 2 insertions(+) diff --git a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java index ffc6847933389..3fb66b775ea37 100644 --- a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java +++ b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java @@ -572,6 +572,8 @@ public Settings onNodeStopped(String nodeName) throws Exception { internalCluster().startNode(); } + @TestLogging("org.elasticsearch.env.NodeEnvironment:TRACE,org.elasticsearch.gateway.TransportNodesListGatewayStartedShards:TRACE," + + "org.elasticsearch.gateway.MetaDataStateFormat:TRACE,org.elasticsearch.index.shard.IndexShard:TRACE") public void testLoadLatestStateWhileClosingShardDoesNotResurrectMetadataDirectory() throws Exception { // This test pertains to a race condition in which deleting a shard concurrently with a TransportNodesListGatewayStartedShards From f89d9b16ed766773cf8cdff460beb88ef6d2e8e2 Mon Sep 17 00:00:00 2001 From: David Turner Date: Thu, 5 Apr 2018 11:03:03 +0100 Subject: [PATCH 13/23] Add note about logging making failures more common --- .../org/elasticsearch/gateway/RecoveryFromGatewayIT.java | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java index 3fb66b775ea37..70728a1c1373f 100644 --- a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java +++ b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java @@ -572,14 +572,17 @@ public Settings onNodeStopped(String nodeName) throws Exception { internalCluster().startNode(); } - @TestLogging("org.elasticsearch.env.NodeEnvironment:TRACE,org.elasticsearch.gateway.TransportNodesListGatewayStartedShards:TRACE," + - "org.elasticsearch.gateway.MetaDataStateFormat:TRACE,org.elasticsearch.index.shard.IndexShard:TRACE") public void testLoadLatestStateWhileClosingShardDoesNotResurrectMetadataDirectory() throws Exception { // This test pertains to a race condition in which deleting a shard concurrently with a TransportNodesListGatewayStartedShards // request could resurrect the shard's metadata folder after it was deleted. Here we try and recreate the race, but it is quite // delicate so this test does not always fail. Experimentation showed that setting the thread counts as below would yield a failure // after a reasonable number of iterations: running repeatedly with -Dtests.iters=1000 saw 6 failures out of 10 runs. + // + // NB this experiment was run with + // @TestLogging("org.elasticsearch.env.NodeEnvironment:TRACE,org.elasticsearch.gateway.MetaDataStateFormat:TRACE," + + // "org.elasticsearch.gateway.TransportNodesListGatewayStartedShards:TRACE,org.elasticsearch.index.shard.IndexShard:TRACE") + // but with less verbose logging the failures seem rarer. final String nodeName = internalCluster().startNode(); DiscoveryNode node = internalCluster().getInstance(ClusterService.class, nodeName).localNode(); From 819d274ecf51006df138b0525fc98f0c53966b0e Mon Sep 17 00:00:00 2001 From: David Turner Date: Thu, 5 Apr 2018 11:06:44 +0100 Subject: [PATCH 14/23] NOCOMMIT reinstate logging for repro test --- .../java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java | 2 ++ 1 file changed, 2 insertions(+) diff --git a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java index 70728a1c1373f..80a9c5b899e94 100644 --- a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java +++ b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java @@ -572,6 +572,8 @@ public Settings onNodeStopped(String nodeName) throws Exception { internalCluster().startNode(); } + @TestLogging("org.elasticsearch.env.NodeEnvironment:TRACE,org.elasticsearch.gateway.MetaDataStateFormat:TRACE," + + "org.elasticsearch.gateway.TransportNodesListGatewayStartedShards:TRACE,org.elasticsearch.index.shard.IndexShard:TRACE") public void testLoadLatestStateWhileClosingShardDoesNotResurrectMetadataDirectory() throws Exception { // This test pertains to a race condition in which deleting a shard concurrently with a TransportNodesListGatewayStartedShards From fd9dc31f6d651f56d89119175f29b0d3b30fb665 Mon Sep 17 00:00:00 2001 From: David Turner Date: Thu, 5 Apr 2018 15:08:49 +0100 Subject: [PATCH 15/23] Revert "NOCOMMIT reinstate logging for repro test" This reverts commit 819d274ecf51006df138b0525fc98f0c53966b0e. --- .../java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java | 2 -- 1 file changed, 2 deletions(-) diff --git a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java index 80a9c5b899e94..70728a1c1373f 100644 --- a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java +++ b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java @@ -572,8 +572,6 @@ public Settings onNodeStopped(String nodeName) throws Exception { internalCluster().startNode(); } - @TestLogging("org.elasticsearch.env.NodeEnvironment:TRACE,org.elasticsearch.gateway.MetaDataStateFormat:TRACE," + - "org.elasticsearch.gateway.TransportNodesListGatewayStartedShards:TRACE,org.elasticsearch.index.shard.IndexShard:TRACE") public void testLoadLatestStateWhileClosingShardDoesNotResurrectMetadataDirectory() throws Exception { // This test pertains to a race condition in which deleting a shard concurrently with a TransportNodesListGatewayStartedShards From 3eff6c9c3881cd296e057521f52067fe5a8e1c3b Mon Sep 17 00:00:00 2001 From: David Turner Date: Fri, 18 May 2018 12:15:10 +0100 Subject: [PATCH 16/23] Can construct ShardStateMetaData from an IndexShard directly, no need to load it --- .../TransportNodesListGatewayStartedShards.java | 8 +------- .../org/elasticsearch/index/shard/IndexShard.java | 11 ++--------- 2 files changed, 3 insertions(+), 16 deletions(-) diff --git a/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java b/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java index 1b7edeb05a127..c3c34cb9ae42d 100644 --- a/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java +++ b/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java @@ -20,7 +20,6 @@ package org.elasticsearch.gateway; import org.apache.logging.log4j.message.ParameterizedMessage; -import org.apache.lucene.store.AlreadyClosedException; import org.elasticsearch.ElasticsearchException; import org.elasticsearch.Version; import org.elasticsearch.action.ActionListener; @@ -118,12 +117,7 @@ protected NodesGatewayStartedShards newResponse(Request request, private ShardStateMetaData safelyLoadLatestState(ShardId shardId) throws Exception { final IndexShard indexShard = indicesService.getShardOrNull(shardId); if (indexShard != null) { - try { - return indexShard.loadShardStateMetaDataIfOpen(NamedXContentRegistry.EMPTY, nodeEnv.availableShardPaths(shardId)); - } - catch (AlreadyClosedException ignored) { - // Ok, fall through to trying to get the shard lock ourselves. - } + return indexShard.getShardStateMetaData(); } try (ShardLock ignored = nodeEnv.shardLock(shardId, TimeUnit.SECONDS.toMillis(5))) { diff --git a/server/src/main/java/org/elasticsearch/index/shard/IndexShard.java b/server/src/main/java/org/elasticsearch/index/shard/IndexShard.java index 0536bcc476651..8364d0018fa2c 100644 --- a/server/src/main/java/org/elasticsearch/index/shard/IndexShard.java +++ b/server/src/main/java/org/elasticsearch/index/shard/IndexShard.java @@ -64,7 +64,6 @@ import org.elasticsearch.common.util.BigArrays; import org.elasticsearch.common.util.concurrent.AbstractRunnable; import org.elasticsearch.common.util.concurrent.AsyncIOProcessor; -import org.elasticsearch.common.xcontent.NamedXContentRegistry; import org.elasticsearch.common.xcontent.XContentHelper; import org.elasticsearch.core.internal.io.IOUtils; import org.elasticsearch.index.Index; @@ -139,7 +138,6 @@ import java.io.PrintStream; import java.nio.channels.ClosedByInterruptException; import java.nio.charset.StandardCharsets; -import java.nio.file.Path; import java.util.ArrayList; import java.util.Collections; import java.util.EnumSet; @@ -2067,14 +2065,9 @@ public void startRecovery(RecoveryState recoveryState, PeerRecoveryTargetService } } - public ShardStateMetaData loadShardStateMetaDataIfOpen(NamedXContentRegistry namedXContentRegistry, Path[] dataLocations) - throws IOException { + public ShardStateMetaData getShardStateMetaData() { synchronized (mutex) { - if (state == IndexShardState.CLOSED) { - throw new AlreadyClosedException(shardId + " can't load shard state metadata - shard is closed"); - } - - return ShardStateMetaData.FORMAT.loadLatestState(logger, namedXContentRegistry, dataLocations); + return new ShardStateMetaData(shardRouting.primary(), indexSettings.getUUID(), shardRouting.allocationId()); } } From 7f835ccb4608ba568ff089a3c9a84279844be9d9 Mon Sep 17 00:00:00 2001 From: David Turner Date: Fri, 18 May 2018 13:04:22 +0100 Subject: [PATCH 17/23] Inline method and avoid calling indicesService.getShardOrNull(shardId) twice --- ...ransportNodesListGatewayStartedShards.java | 116 +++++++++--------- 1 file changed, 59 insertions(+), 57 deletions(-) diff --git a/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java b/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java index c3c34cb9ae42d..8a354accdde2e 100644 --- a/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java +++ b/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java @@ -114,74 +114,76 @@ protected NodesGatewayStartedShards newResponse(Request request, return new NodesGatewayStartedShards(clusterService.getClusterName(), responses, failures); } - private ShardStateMetaData safelyLoadLatestState(ShardId shardId) throws Exception { - final IndexShard indexShard = indicesService.getShardOrNull(shardId); - if (indexShard != null) { - return indexShard.getShardStateMetaData(); - } - - try (ShardLock ignored = nodeEnv.shardLock(shardId, TimeUnit.SECONDS.toMillis(5))) { - return ShardStateMetaData.FORMAT.loadLatestState(logger, NamedXContentRegistry.EMPTY, - nodeEnv.availableShardPaths(shardId)); - } - } - @Override protected NodeGatewayStartedShards nodeOperation(NodeRequest request) { try { final ShardId shardId = request.getShardId(); logger.trace("{} loading local shard state info", shardId); - ShardStateMetaData shardStateMetaData = safelyLoadLatestState(shardId); - if (shardStateMetaData != null) { - IndexMetaData metaData = clusterService.state().metaData().index(shardId.getIndex()); - if (metaData == null) { - // we may send this requests while processing the cluster state that recovered the index - // sometimes the request comes in before the local node processed that cluster state - // in such cases we can load it from disk - metaData = IndexMetaData.FORMAT.loadLatestState(logger, NamedXContentRegistry.EMPTY, - nodeEnv.indexPaths(shardId.getIndex())); - } - if (metaData == null) { - ElasticsearchException e = new ElasticsearchException("failed to find local IndexMetaData"); - e.setShard(request.shardId); - throw e; - } + final IndexShard indexShard = indicesService.getShardOrNull(shardId); + if (indexShard != null) { + final ShardStateMetaData shardStateMetaData = indexShard.getShardStateMetaData(); + final String allocationId = shardStateMetaData.allocationId != null ? + shardStateMetaData.allocationId.getId() : null; + logger.debug("{} shard state info found: [{}]", shardId, shardStateMetaData); + return new NodeGatewayStartedShards(clusterService.localNode(), allocationId, shardStateMetaData.primary); + } - if (indicesService.getShardOrNull(shardId) == null) { - // we don't have an open shard on the store, validate the files on disk are openable - ShardPath shardPath = null; - try { - IndexSettings indexSettings = new IndexSettings(metaData, settings); - try (ShardLock ignored = nodeEnv.shardLock(shardId, TimeUnit.SECONDS.toMillis(5))) { - shardPath = ShardPath.loadShardPath(logger, nodeEnv, shardId, indexSettings); - } - if (shardPath == null) { - throw new IllegalStateException(shardId + " no shard path found"); - } - Store.tryOpenIndex(shardPath.resolveIndex(), shardId, nodeEnv::shardLock, logger); - } catch (Exception exception) { - final ShardPath finalShardPath = shardPath; - logger.trace(() -> new ParameterizedMessage( - "{} can't open index for shard [{}] in path [{}]", - shardId, - shardStateMetaData, - (finalShardPath != null) ? finalShardPath.resolveIndex() : ""), - exception); - String allocationId = shardStateMetaData.allocationId != null ? - shardStateMetaData.allocationId.getId() : null; - return new NodeGatewayStartedShards(clusterService.localNode(), allocationId, shardStateMetaData.primary, - exception); - } - } + final ShardStateMetaData shardStateMetaData; + try (ShardLock ignored = nodeEnv.shardLock(shardId, TimeUnit.SECONDS.toMillis(5))) { + shardStateMetaData = ShardStateMetaData.FORMAT.loadLatestState(logger, NamedXContentRegistry.EMPTY, + nodeEnv.availableShardPaths(shardId)); + } - logger.debug("{} shard state info found: [{}]", shardId, shardStateMetaData); + if (shardStateMetaData == null) { + logger.trace("{} no local shard info found", shardId); + return new NodeGatewayStartedShards(clusterService.localNode(), null, false); + } + + IndexMetaData metaData = clusterService.state().metaData().index(shardId.getIndex()); + if (metaData == null) { + // we may send this requests while processing the cluster state that recovered the index + // sometimes the request comes in before the local node processed that cluster state + // in such cases we can load it from disk + metaData = IndexMetaData.FORMAT.loadLatestState(logger, NamedXContentRegistry.EMPTY, + nodeEnv.indexPaths(shardId.getIndex())); + } + if (metaData == null) { + ElasticsearchException e = new ElasticsearchException("failed to find local IndexMetaData"); + e.setShard(request.shardId); + throw e; + } + + // we don't have an open shard on the store, validate the files on disk are openable + ShardPath shardPath = null; + try { + IndexSettings indexSettings = new IndexSettings(metaData, settings); + try (ShardLock ignored = nodeEnv.shardLock(shardId, TimeUnit.SECONDS.toMillis(5))) { + shardPath = ShardPath.loadShardPath(logger, nodeEnv, shardId, indexSettings); + } + if (shardPath == null) { + throw new IllegalStateException(shardId + " no shard path found"); + } + Store.tryOpenIndex(shardPath.resolveIndex(), shardId, nodeEnv::shardLock, logger); + } catch (Exception exception) { + final ShardPath finalShardPath = shardPath; + logger.trace(() -> new ParameterizedMessage( + "{} can't open index for shard [{}] in path [{}]", + shardId, + shardStateMetaData, + (finalShardPath != null) ? finalShardPath.resolveIndex() : ""), + exception); String allocationId = shardStateMetaData.allocationId != null ? shardStateMetaData.allocationId.getId() : null; - return new NodeGatewayStartedShards(clusterService.localNode(), allocationId, shardStateMetaData.primary); + return new NodeGatewayStartedShards(clusterService.localNode(), allocationId, shardStateMetaData.primary, + exception); } - logger.trace("{} no local shard info found", shardId); - return new NodeGatewayStartedShards(clusterService.localNode(), null, false); + + logger.debug("{} shard state info found: [{}]", shardId, shardStateMetaData); + String allocationId = shardStateMetaData.allocationId != null ? + shardStateMetaData.allocationId.getId() : null; + return new NodeGatewayStartedShards(clusterService.localNode(), allocationId, shardStateMetaData.primary); + } catch (Exception e) { throw new ElasticsearchException("failed to load started shards", e); } From 48f6d46d0a7754801de1deb13cae4b63d8849c4b Mon Sep 17 00:00:00 2001 From: David Turner Date: Fri, 18 May 2018 13:28:12 +0100 Subject: [PATCH 18/23] Assert that directory is really deleted --- .../gateway/RecoveryFromGatewayIT.java | 21 ++++++++++++++----- 1 file changed, 16 insertions(+), 5 deletions(-) diff --git a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java index 70728a1c1373f..845a9e7dda144 100644 --- a/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java +++ b/server/src/test/java/org/elasticsearch/gateway/RecoveryFromGatewayIT.java @@ -50,6 +50,7 @@ import org.elasticsearch.test.junit.annotations.TestLogging; import org.elasticsearch.test.store.MockFSIndexStore; +import java.io.File; import java.nio.file.DirectoryStream; import java.nio.file.Files; import java.nio.file.Path; @@ -617,14 +618,24 @@ public void testLoadLatestStateWhileClosingShardDoesNotResurrectMetadataDirector }, (isListingThread ? "Listing" : "Deleting") + "[" + threadIndex + "]"); } - for (final Thread listingThread : threads) { - listingThread.start(); + NodeEnvironment nodeEnvironment = internalCluster().getInstance(NodeEnvironment.class, nodeName); + + boolean directoryExists = false; + for (Path path : nodeEnvironment.availableShardPaths(shardId)) { + directoryExists = directoryExists || Files.exists(path); + } + assertTrue(directoryExists); + + for (final Thread thread : threads) { + thread.start(); } - // Deleting an index asserts that it really is gone from disk, so no other assertions are necessary here. + for (final Thread thread : threads) { + thread.join(); + } - for (final Thread listingThread : threads) { - listingThread.join(); + for (Path path : nodeEnvironment.availableShardPaths(shardId)) { + assertFalse(path + " should not exist", Files.exists(path)); } } } From 7e58bc67f02e7a210c127cf41ef03498ab87b705 Mon Sep 17 00:00:00 2001 From: David Turner Date: Thu, 24 May 2018 12:13:53 +0100 Subject: [PATCH 19/23] debug -> trace --- .../gateway/TransportNodesListGatewayStartedShards.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java b/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java index 8a354accdde2e..48a35621be84d 100644 --- a/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java +++ b/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java @@ -125,7 +125,7 @@ protected NodeGatewayStartedShards nodeOperation(NodeRequest request) { final ShardStateMetaData shardStateMetaData = indexShard.getShardStateMetaData(); final String allocationId = shardStateMetaData.allocationId != null ? shardStateMetaData.allocationId.getId() : null; - logger.debug("{} shard state info found: [{}]", shardId, shardStateMetaData); + logger.trace("{} shard state info found: [{}]", shardId, shardStateMetaData); return new NodeGatewayStartedShards(clusterService.localNode(), allocationId, shardStateMetaData.primary); } From 1d4e0447651c6c14e8774c2cc93f459686121e78 Mon Sep 17 00:00:00 2001 From: David Turner Date: Wed, 30 May 2018 14:33:53 +0100 Subject: [PATCH 20/23] No need for mutex, just read shardRouting once --- .../main/java/org/elasticsearch/index/shard/IndexShard.java | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/server/src/main/java/org/elasticsearch/index/shard/IndexShard.java b/server/src/main/java/org/elasticsearch/index/shard/IndexShard.java index 8364d0018fa2c..5d346d544c768 100644 --- a/server/src/main/java/org/elasticsearch/index/shard/IndexShard.java +++ b/server/src/main/java/org/elasticsearch/index/shard/IndexShard.java @@ -2066,9 +2066,8 @@ public void startRecovery(RecoveryState recoveryState, PeerRecoveryTargetService } public ShardStateMetaData getShardStateMetaData() { - synchronized (mutex) { - return new ShardStateMetaData(shardRouting.primary(), indexSettings.getUUID(), shardRouting.allocationId()); - } + final ShardRouting shardRouting = this.shardRouting; + return new ShardStateMetaData(shardRouting.primary(), indexSettings.getUUID(), shardRouting.allocationId()); } /** From 61b4e4eed7ba6cf57e10fa3e029649bddbaed090 Mon Sep 17 00:00:00 2001 From: David Turner Date: Wed, 30 May 2018 14:45:49 +0100 Subject: [PATCH 21/23] Only obtain lock once in TransportNodesListShardStoreMetaData --- .../org/elasticsearch/index/store/Store.java | 12 +++----- .../TransportNodesListShardStoreMetaData.java | 29 +++++++++++-------- 2 files changed, 21 insertions(+), 20 deletions(-) diff --git a/server/src/main/java/org/elasticsearch/index/store/Store.java b/server/src/main/java/org/elasticsearch/index/store/Store.java index ccaae9d5f79df..d70fd7a990138 100644 --- a/server/src/main/java/org/elasticsearch/index/store/Store.java +++ b/server/src/main/java/org/elasticsearch/index/store/Store.java @@ -241,7 +241,6 @@ final void ensureOpen() { * Note that this method requires the caller verify it has the right to access the store and * no concurrent file changes are happening. If in doubt, you probably want to use one of the following: * - * {@link #readMetadataSnapshot(Path, ShardId, NodeEnvironment.ShardLocker, Logger)} to read a meta data while locking * {@link IndexShard#snapshotStoreMetadata()} to safely read from an existing shard * {@link IndexShard#acquireLastIndexCommit(boolean)} to get an {@link IndexCommit} which is safe to use but has to be freed * @param commit the index commit to read the snapshot from or null if the latest snapshot should be read from the @@ -265,7 +264,6 @@ public MetadataSnapshot getMetadata(IndexCommit commit) throws IOException { * Note that this method requires the caller verify it has the right to access the store and * no concurrent file changes are happening. If in doubt, you probably want to use one of the following: * - * {@link #readMetadataSnapshot(Path, ShardId, NodeEnvironment.ShardLocker, Logger)} to read a meta data while locking * {@link IndexShard#snapshotStoreMetadata()} to safely read from an existing shard * {@link IndexShard#acquireLastIndexCommit(boolean)} to get an {@link IndexCommit} which is safe to use but has to be freed * @@ -456,18 +454,16 @@ private void closeInternal() { * * @throws IOException if the index we try to read is corrupted */ - public static MetadataSnapshot readMetadataSnapshot(Path indexLocation, ShardId shardId, NodeEnvironment.ShardLocker shardLocker, - Logger logger) throws IOException { - try (ShardLock lock = shardLocker.lock(shardId, TimeUnit.SECONDS.toMillis(5)); - Directory dir = new SimpleFSDirectory(indexLocation)) { + public static MetadataSnapshot readMetadataSnapshot(Path indexLocation, ShardId shardId, Logger logger, ShardLock shardLock) + throws IOException { + assert shardLock.isOpen(); + try (Directory dir = new SimpleFSDirectory(indexLocation)) { failIfCorrupted(dir, shardId); return new MetadataSnapshot(null, dir, logger); } catch (IndexNotFoundException ex) { // that's fine - happens all the time no need to log } catch (FileNotFoundException | NoSuchFileException ex) { logger.info("Failed to open / find files while reading metadata snapshot"); - } catch (ShardLockObtainFailedException ex) { - logger.info(() -> new ParameterizedMessage("{}: failed to obtain shard lock", shardId), ex); } return MetadataSnapshot.EMPTY; } diff --git a/server/src/main/java/org/elasticsearch/indices/store/TransportNodesListShardStoreMetaData.java b/server/src/main/java/org/elasticsearch/indices/store/TransportNodesListShardStoreMetaData.java index d91bba2ada929..fa5ce6e6533da 100644 --- a/server/src/main/java/org/elasticsearch/indices/store/TransportNodesListShardStoreMetaData.java +++ b/server/src/main/java/org/elasticsearch/indices/store/TransportNodesListShardStoreMetaData.java @@ -19,6 +19,7 @@ package org.elasticsearch.indices.store; +import org.apache.logging.log4j.message.ParameterizedMessage; import org.elasticsearch.ElasticsearchException; import org.elasticsearch.action.ActionListener; import org.elasticsearch.action.FailedNodeException; @@ -42,6 +43,7 @@ import org.elasticsearch.common.xcontent.NamedXContentRegistry; import org.elasticsearch.env.NodeEnvironment; import org.elasticsearch.env.ShardLock; +import org.elasticsearch.env.ShardLockObtainFailedException; import org.elasticsearch.gateway.AsyncShardFetch; import org.elasticsearch.index.IndexService; import org.elasticsearch.index.IndexSettings; @@ -139,20 +141,23 @@ private StoreFilesMetaData listStoreMetaData(ShardId shardId) throws IOException logger.trace("{} node doesn't have meta data for the requests index, responding with empty", shardId); return new StoreFilesMetaData(shardId, Store.MetadataSnapshot.EMPTY); } - final IndexSettings indexSettings = indexService != null ? indexService.getIndexSettings() : new IndexSettings(metaData, settings); - final ShardPath shardPath; - try (ShardLock ignored = nodeEnv.shardLock(shardId, TimeUnit.SECONDS.toMillis(5))) { - shardPath = ShardPath.loadShardPath(logger, nodeEnv, shardId, indexSettings); - } - if (shardPath == null) { - return new StoreFilesMetaData(shardId, Store.MetadataSnapshot.EMPTY); - } + final IndexSettings indexSettings + = indexService != null ? indexService.getIndexSettings() : new IndexSettings(metaData, settings); + // note that this may fail if it can't get access to the shard lock. Since we check above there is an active shard, this means: - // 1) a shard is being constructed, which means the master will not use a copy of this replica - // 2) A shard is shutting down and has not cleared it's content within lock timeout. In this case the master may not + // 1) a shard is being constructed, which means the master will not use a copy of this replica. + // 2) a shard is shutting down and has not cleared its content within the lock timeout. In this case the master may not // reuse local resources. - return new StoreFilesMetaData(shardId, Store.readMetadataSnapshot(shardPath.resolveIndex(), shardId, - nodeEnv::shardLock, logger)); + try (ShardLock shardLock = nodeEnv.shardLock(shardId, TimeUnit.SECONDS.toMillis(5))) { + final ShardPath shardPath = ShardPath.loadShardPath(logger, nodeEnv, shardId, indexSettings); + if (shardPath != null) { + return new StoreFilesMetaData(shardId, + Store.readMetadataSnapshot(shardPath.resolveIndex(), shardId, logger, shardLock)); + } + } catch (ShardLockObtainFailedException ex) { + logger.info(() -> new ParameterizedMessage("{}: failed to obtain shard lock", shardId), ex); + } + return new StoreFilesMetaData(shardId, Store.MetadataSnapshot.EMPTY); } finally { TimeValue took = new TimeValue(System.nanoTime() - startTimeNS, TimeUnit.NANOSECONDS); if (exists) { From 8f1a5e2c2dcbd47a1929a05b3f1f432e9eed2489 Mon Sep 17 00:00:00 2001 From: David Turner Date: Wed, 30 May 2018 16:09:55 +0100 Subject: [PATCH 22/23] Only obtain lock once in TransportNodesListGatewayStartedShards --- .../java/org/elasticsearch/env/NodeEnvironment.java | 8 -------- .../TransportNodesListGatewayStartedShards.java | 10 +++++----- .../java/org/elasticsearch/index/store/Store.java | 12 +++++------- .../elasticsearch/index/shard/IndexShardTests.java | 3 +-- .../org/elasticsearch/index/store/StoreTests.java | 6 +++--- 5 files changed, 14 insertions(+), 25 deletions(-) diff --git a/server/src/main/java/org/elasticsearch/env/NodeEnvironment.java b/server/src/main/java/org/elasticsearch/env/NodeEnvironment.java index 87874bd45000c..f97f34670f19a 100644 --- a/server/src/main/java/org/elasticsearch/env/NodeEnvironment.java +++ b/server/src/main/java/org/elasticsearch/env/NodeEnvironment.java @@ -605,14 +605,6 @@ protected void closeInternal() { }; } - /** - * A functional interface that people can use to reference {@link #shardLock(ShardId, long)} - */ - @FunctionalInterface - public interface ShardLocker { - ShardLock lock(ShardId shardId, long lockTimeoutMS) throws ShardLockObtainFailedException; - } - /** * Returns all currently lock shards. * diff --git a/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java b/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java index 48a35621be84d..c514c75884554 100644 --- a/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java +++ b/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java @@ -158,13 +158,13 @@ protected NodeGatewayStartedShards nodeOperation(NodeRequest request) { ShardPath shardPath = null; try { IndexSettings indexSettings = new IndexSettings(metaData, settings); - try (ShardLock ignored = nodeEnv.shardLock(shardId, TimeUnit.SECONDS.toMillis(5))) { + try (ShardLock shardLock = nodeEnv.shardLock(shardId, TimeUnit.SECONDS.toMillis(5))) { shardPath = ShardPath.loadShardPath(logger, nodeEnv, shardId, indexSettings); + if (shardPath == null) { + throw new IllegalStateException(shardId + " no shard path found"); + } + Store.tryOpenIndex(shardPath.resolveIndex(), shardId, shardLock, logger); } - if (shardPath == null) { - throw new IllegalStateException(shardId + " no shard path found"); - } - Store.tryOpenIndex(shardPath.resolveIndex(), shardId, nodeEnv::shardLock, logger); } catch (Exception exception) { final ShardPath finalShardPath = shardPath; logger.trace(() -> new ParameterizedMessage( diff --git a/server/src/main/java/org/elasticsearch/index/store/Store.java b/server/src/main/java/org/elasticsearch/index/store/Store.java index d70fd7a990138..5a46e1b9cbdb9 100644 --- a/server/src/main/java/org/elasticsearch/index/store/Store.java +++ b/server/src/main/java/org/elasticsearch/index/store/Store.java @@ -72,9 +72,7 @@ import org.elasticsearch.common.util.concurrent.RefCounted; import org.elasticsearch.common.util.iterable.Iterables; import org.elasticsearch.core.internal.io.IOUtils; -import org.elasticsearch.env.NodeEnvironment; import org.elasticsearch.env.ShardLock; -import org.elasticsearch.env.ShardLockObtainFailedException; import org.elasticsearch.index.IndexSettings; import org.elasticsearch.index.engine.CombinedDeletionPolicy; import org.elasticsearch.index.engine.Engine; @@ -473,9 +471,9 @@ public static MetadataSnapshot readMetadataSnapshot(Path indexLocation, ShardId * can be successfully opened. This includes reading the segment infos and possible * corruption markers. */ - public static boolean canOpenIndex(Logger logger, Path indexLocation, ShardId shardId, NodeEnvironment.ShardLocker shardLocker) throws IOException { + public static boolean canOpenIndex(Logger logger, Path indexLocation, ShardId shardId, ShardLock shardLock) { try { - tryOpenIndex(indexLocation, shardId, shardLocker, logger); + tryOpenIndex(indexLocation, shardId, shardLock, logger); } catch (Exception ex) { logger.trace(() -> new ParameterizedMessage("Can't open index for path [{}]", indexLocation), ex); return false; @@ -488,9 +486,9 @@ public static boolean canOpenIndex(Logger logger, Path indexLocation, ShardId sh * segment infos and possible corruption markers. If the index can not * be opened, an exception is thrown */ - public static void tryOpenIndex(Path indexLocation, ShardId shardId, NodeEnvironment.ShardLocker shardLocker, Logger logger) throws IOException, ShardLockObtainFailedException { - try (ShardLock lock = shardLocker.lock(shardId, TimeUnit.SECONDS.toMillis(5)); - Directory dir = new SimpleFSDirectory(indexLocation)) { + public static void tryOpenIndex(Path indexLocation, ShardId shardId, ShardLock shardLock, Logger logger) throws IOException { + assert shardLock.isOpen(); + try (Directory dir = new SimpleFSDirectory(indexLocation)) { failIfCorrupted(dir, shardId); SegmentInfos segInfo = Lucene.readSegmentInfos(dir); logger.trace("{} loaded segment info [{}]", shardId, segInfo); diff --git a/server/src/test/java/org/elasticsearch/index/shard/IndexShardTests.java b/server/src/test/java/org/elasticsearch/index/shard/IndexShardTests.java index 027b595ee761f..b1b804853d0f0 100644 --- a/server/src/test/java/org/elasticsearch/index/shard/IndexShardTests.java +++ b/server/src/test/java/org/elasticsearch/index/shard/IndexShardTests.java @@ -234,8 +234,7 @@ public void testFailShard() throws Exception { assertEquals(shardStateMetaData, getShardStateMetadata(shard)); // but index can't be opened for a failed shard assertThat("store index should be corrupted", Store.canOpenIndex(logger, shardPath.resolveIndex(), shard.shardId(), - (shardId, lockTimeoutMS) -> new DummyShardLock(shardId)), - equalTo(false)); + new DummyShardLock(shard.shardId())), equalTo(false)); } ShardStateMetaData getShardStateMetadata(IndexShard shard) { diff --git a/server/src/test/java/org/elasticsearch/index/store/StoreTests.java b/server/src/test/java/org/elasticsearch/index/store/StoreTests.java index 9352d978e6e46..63eaa88af9705 100644 --- a/server/src/test/java/org/elasticsearch/index/store/StoreTests.java +++ b/server/src/test/java/org/elasticsearch/index/store/StoreTests.java @@ -948,14 +948,14 @@ public void testCanOpenIndex() throws IOException { IndexWriterConfig iwc = newIndexWriterConfig(); Path tempDir = createTempDir(); final BaseDirectoryWrapper dir = newFSDirectory(tempDir); - assertFalse(Store.canOpenIndex(logger, tempDir, shardId, (id, l) -> new DummyShardLock(id))); + assertFalse(Store.canOpenIndex(logger, tempDir, shardId, new DummyShardLock(shardId))); IndexWriter writer = new IndexWriter(dir, iwc); Document doc = new Document(); doc.add(new StringField("id", "1", random().nextBoolean() ? Field.Store.YES : Field.Store.NO)); writer.addDocument(doc); writer.commit(); writer.close(); - assertTrue(Store.canOpenIndex(logger, tempDir, shardId, (id, l) -> new DummyShardLock(id))); + assertTrue(Store.canOpenIndex(logger, tempDir, shardId, new DummyShardLock(shardId))); DirectoryService directoryService = new DirectoryService(shardId, INDEX_SETTINGS) { @@ -966,7 +966,7 @@ public Directory newDirectory() throws IOException { }; Store store = new Store(shardId, INDEX_SETTINGS, directoryService, new DummyShardLock(shardId)); store.markStoreCorrupted(new CorruptIndexException("foo", "bar")); - assertFalse(Store.canOpenIndex(logger, tempDir, shardId, (id, l) -> new DummyShardLock(id))); + assertFalse(Store.canOpenIndex(logger, tempDir, shardId, new DummyShardLock(shardId))); store.close(); } From 78c0526c96e4a8589b28850f4404474c7d05edf8 Mon Sep 17 00:00:00 2001 From: David Turner Date: Thu, 31 May 2018 11:32:22 +0100 Subject: [PATCH 23/23] Only get lock once in TransportNodesListGatewayStartedShards, and return a null result if it fails --- ...ransportNodesListGatewayStartedShards.java | 85 ++++++++++--------- 1 file changed, 43 insertions(+), 42 deletions(-) diff --git a/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java b/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java index c514c75884554..695d2d1d4adac 100644 --- a/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java +++ b/server/src/main/java/org/elasticsearch/gateway/TransportNodesListGatewayStartedShards.java @@ -42,6 +42,7 @@ import org.elasticsearch.common.xcontent.NamedXContentRegistry; import org.elasticsearch.env.NodeEnvironment; import org.elasticsearch.env.ShardLock; +import org.elasticsearch.env.ShardLockObtainFailedException; import org.elasticsearch.index.IndexSettings; import org.elasticsearch.index.shard.IndexShard; import org.elasticsearch.index.shard.ShardId; @@ -129,60 +130,60 @@ protected NodeGatewayStartedShards nodeOperation(NodeRequest request) { return new NodeGatewayStartedShards(clusterService.localNode(), allocationId, shardStateMetaData.primary); } - final ShardStateMetaData shardStateMetaData; - try (ShardLock ignored = nodeEnv.shardLock(shardId, TimeUnit.SECONDS.toMillis(5))) { - shardStateMetaData = ShardStateMetaData.FORMAT.loadLatestState(logger, NamedXContentRegistry.EMPTY, - nodeEnv.availableShardPaths(shardId)); - } + try (ShardLock shardLock = nodeEnv.shardLock(shardId, TimeUnit.SECONDS.toMillis(5))) { + final ShardStateMetaData shardStateMetaData + = ShardStateMetaData.FORMAT.loadLatestState(logger, NamedXContentRegistry.EMPTY, nodeEnv.availableShardPaths(shardId)); - if (shardStateMetaData == null) { - logger.trace("{} no local shard info found", shardId); - return new NodeGatewayStartedShards(clusterService.localNode(), null, false); - } + if (shardStateMetaData == null) { + logger.trace("{} no local shard info found", shardId); + return new NodeGatewayStartedShards(clusterService.localNode(), null, false); + } - IndexMetaData metaData = clusterService.state().metaData().index(shardId.getIndex()); - if (metaData == null) { - // we may send this requests while processing the cluster state that recovered the index - // sometimes the request comes in before the local node processed that cluster state - // in such cases we can load it from disk - metaData = IndexMetaData.FORMAT.loadLatestState(logger, NamedXContentRegistry.EMPTY, - nodeEnv.indexPaths(shardId.getIndex())); - } - if (metaData == null) { - ElasticsearchException e = new ElasticsearchException("failed to find local IndexMetaData"); - e.setShard(request.shardId); - throw e; - } + IndexMetaData metaData = clusterService.state().metaData().index(shardId.getIndex()); + if (metaData == null) { + // we may send this requests while processing the cluster state that recovered the index + // sometimes the request comes in before the local node processed that cluster state + // in such cases we can load it from disk + metaData = IndexMetaData.FORMAT.loadLatestState(logger, NamedXContentRegistry.EMPTY, + nodeEnv.indexPaths(shardId.getIndex())); + } + if (metaData == null) { + ElasticsearchException e = new ElasticsearchException("failed to find local IndexMetaData"); + e.setShard(request.shardId); + throw e; + } - // we don't have an open shard on the store, validate the files on disk are openable - ShardPath shardPath = null; - try { - IndexSettings indexSettings = new IndexSettings(metaData, settings); - try (ShardLock shardLock = nodeEnv.shardLock(shardId, TimeUnit.SECONDS.toMillis(5))) { + // we don't have an open shard on the store, validate the files on disk are openable + ShardPath shardPath = null; + try { + IndexSettings indexSettings = new IndexSettings(metaData, settings); shardPath = ShardPath.loadShardPath(logger, nodeEnv, shardId, indexSettings); if (shardPath == null) { throw new IllegalStateException(shardId + " no shard path found"); } Store.tryOpenIndex(shardPath.resolveIndex(), shardId, shardLock, logger); + } catch (Exception exception) { + final ShardPath finalShardPath = shardPath; + logger.trace(() -> new ParameterizedMessage( + "{} can't open index for shard [{}] in path [{}]", + shardId, + shardStateMetaData, + (finalShardPath != null) ? finalShardPath.resolveIndex() : ""), + exception); + String allocationId = shardStateMetaData.allocationId != null ? + shardStateMetaData.allocationId.getId() : null; + return new NodeGatewayStartedShards(clusterService.localNode(), allocationId, shardStateMetaData.primary, + exception); } - } catch (Exception exception) { - final ShardPath finalShardPath = shardPath; - logger.trace(() -> new ParameterizedMessage( - "{} can't open index for shard [{}] in path [{}]", - shardId, - shardStateMetaData, - (finalShardPath != null) ? finalShardPath.resolveIndex() : ""), - exception); + + logger.debug("{} shard state info found: [{}]", shardId, shardStateMetaData); String allocationId = shardStateMetaData.allocationId != null ? shardStateMetaData.allocationId.getId() : null; - return new NodeGatewayStartedShards(clusterService.localNode(), allocationId, shardStateMetaData.primary, - exception); - } + return new NodeGatewayStartedShards(clusterService.localNode(), allocationId, shardStateMetaData.primary); - logger.debug("{} shard state info found: [{}]", shardId, shardStateMetaData); - String allocationId = shardStateMetaData.allocationId != null ? - shardStateMetaData.allocationId.getId() : null; - return new NodeGatewayStartedShards(clusterService.localNode(), allocationId, shardStateMetaData.primary); + } catch (ShardLockObtainFailedException e) { + return new NodeGatewayStartedShards(clusterService.localNode(), null, false, e); + } } catch (Exception e) { throw new ElasticsearchException("failed to load started shards", e);