summaryrefslogtreecommitdiff
path: root/src/mongo/s/client/shard_registry.cpp
diff options
context:
space:
mode:
Diffstat (limited to 'src/mongo/s/client/shard_registry.cpp')
-rw-r--r--src/mongo/s/client/shard_registry.cpp99
1 files changed, 27 insertions, 72 deletions
diff --git a/src/mongo/s/client/shard_registry.cpp b/src/mongo/s/client/shard_registry.cpp
index dfb3932f3c0..d722a9116ec 100644
--- a/src/mongo/s/client/shard_registry.cpp
+++ b/src/mongo/s/client/shard_registry.cpp
@@ -130,8 +130,12 @@ ShardRegistry::Cache::LookupResult ShardRegistry::_lookup(OperationContext* opCt
// Check if we need to refresh from the configsvrs. If so, then do that and get the results,
// otherwise (this is a lookup only to incorporate updated connection strings from the RSM),
// then get the equivalent values from the previously cached data.
- auto [returnData, returnTopologyTime, returnForceReloadIncrement, removedShards] =
- [&]() -> std::tuple<ShardRegistryData, Timestamp, Increment, ShardRegistryData::ShardMap> {
+ auto [returnData,
+ returnTopologyTime,
+ returnForceReloadIncrement,
+ removedShards,
+ fetchedFromConfigServers] = [&]()
+ -> std::tuple<ShardRegistryData, Timestamp, Increment, ShardRegistryData::ShardMap, bool> {
if (timeInStore.topologyTime > cachedData.getTime().topologyTime ||
timeInStore.forceReloadIncrement > cachedData.getTime().forceReloadIncrement) {
auto [reloadedData, maxTopologyTime] =
@@ -140,12 +144,14 @@ ShardRegistry::Cache::LookupResult ShardRegistry::_lookup(OperationContext* opCt
auto [mergedData, removedShards] =
ShardRegistryData::mergeExisting(*cachedData, reloadedData);
- return {mergedData, maxTopologyTime, timeInStore.forceReloadIncrement, removedShards};
+ return {
+ mergedData, maxTopologyTime, timeInStore.forceReloadIncrement, removedShards, true};
} else {
return {*cachedData,
cachedData.getTime().topologyTime,
cachedData.getTime().forceReloadIncrement,
- {}};
+ {},
+ false};
}
}();
@@ -180,6 +186,11 @@ ShardRegistry::Cache::LookupResult ShardRegistry::_lookup(OperationContext* opCt
}
}
+ // The registry is "up" once there has been a successful lookup from the config servers.
+ if (fetchedFromConfigServers) {
+ _isUp.store(true);
+ }
+
Time returnTime{returnTopologyTime, rsmIncrementForConnStrings, returnForceReloadIncrement};
LOGV2_DEBUG(4620251,
2,
@@ -207,9 +218,9 @@ void ShardRegistry::startupPeriodicReloader(OperationContext* opCtx) {
AsyncTry([this] {
LOGV2_DEBUG(22726, 1, "Reloading shardRegistry");
- return _reloadAsyncNoRetry();
+ return _reloadInternal();
})
- .until([](auto&& sw) {
+ .until([](auto sw) {
if (!sw.isOK()) {
LOGV2(22727,
"Error running periodic reload of shard registry",
@@ -221,7 +232,7 @@ void ShardRegistry::startupPeriodicReloader(OperationContext* opCtx) {
})
.withDelayBetweenIterations(kRefreshPeriod) // This call is optional.
.on(_executor, CancellationToken::uncancelable())
- .getAsync([](auto&& sw) {
+ .getAsync([](auto sw) {
LOGV2_DEBUG(22725,
1,
"Exiting periodic shard registry reloader",
@@ -284,49 +295,6 @@ StatusWith<std::shared_ptr<Shard>> ShardRegistry::getShard(OperationContext* opC
return {ErrorCodes::ShardNotFound, str::stream() << "Shard " << shardId << " not found"};
}
-SemiFuture<std::shared_ptr<Shard>> ShardRegistry::getShard(ExecutorPtr executor,
- const ShardId& shardId) noexcept {
-
- // Fetch the shard registry data associated to the latest known topology time
- return _getDataAsync()
- .thenRunOn(executor)
- .then([this, executor, shardId](auto&& cachedData) {
- // First check if this is a non config shard lookup
- if (auto shard = cachedData->findShard(shardId)) {
- return SemiFuture<std::shared_ptr<Shard>>::makeReady(std::move(shard));
- }
-
- // then check if this is a config shard (this call is blocking in any case)
- {
- stdx::lock_guard<Latch> lk(_mutex);
- if (auto shard = _configShardData.findShard(shardId)) {
- return SemiFuture<std::shared_ptr<Shard>>::makeReady(std::move(shard));
- }
- }
-
- // If the shard was not found, force reload the shard regitry data and try again.
- //
- // This is to cover the following scenario:
- // 1. Primary of the replicaset fetch the list of shards and store it on disk
- // 2. Primary crash before the latest VectorClock topology time is majority written to
- // disk
- // 3. A new primary with a stale ShardRegistry is elected and read the set of shards
- // from disk and calls ShardRegistry::getShard
-
- return _reloadAsync()
- .thenRunOn(executor)
- .then([this, executor, shardId](auto&& cachedData) -> std::shared_ptr<Shard> {
- auto shard = cachedData->findShard(shardId);
- uassert(ErrorCodes::ShardNotFound,
- str::stream() << "Shard " << shardId << " not found",
- shard);
- return shard;
- })
- .semi();
- })
- .semi();
-}
-
std::vector<ShardId> ShardRegistry::getAllShardIds(OperationContext* opCtx) {
auto shardIds = _getData(opCtx)->getAllShardIds();
if (shardIds.empty()) {
@@ -412,18 +380,8 @@ std::unique_ptr<Shard> ShardRegistry::createConnection(const ConnectionString& c
return _shardFactory->createUniqueShard(ShardId("<unnamed>"), connStr);
}
-bool ShardRegistry::isUp() {
- if (_isUp.load())
- return true;
-
- // Before the first lookup is completed, the latest cached value is either empty or it is
- // associated to the default constructed time
- const auto latestCached = _cache->peekLatestCached(_kSingleton);
- if (latestCached && latestCached.getTime() != Time()) {
- _isUp.store(true);
- return true;
- }
- return false;
+bool ShardRegistry::isUp() const {
+ return _isUp.load();
}
void ShardRegistry::toBSON(BSONObjBuilder* result) const {
@@ -442,26 +400,23 @@ void ShardRegistry::toBSON(BSONObjBuilder* result) const {
}
void ShardRegistry::reload(OperationContext* opCtx) {
- _reloadAsync().get(opCtx);
-}
-
-SharedSemiFuture<ShardRegistry::Cache::ValueHandle> ShardRegistry::_reloadAsync() {
if (MONGO_unlikely(TestingProctor::instance().isEnabled())) {
// Some unit tests don't support running the reload's AsyncTry on the fixed executor.
- return _reloadAsyncNoRetry();
+ _reloadInternal().get(opCtx);
} else {
- return AsyncTry([=]() mutable { return _reloadAsyncNoRetry(); })
+ AsyncTry([=]() mutable { return _reloadInternal(); })
.until([](auto sw) mutable {
return sw.getStatus() != ErrorCodes::ReadConcernMajorityNotAvailableYet;
})
.withBackoffBetweenIterations(kExponentialBackoff)
- .on(Grid::get(getGlobalServiceContext())->getExecutorPool()->getFixedExecutor(),
+ .on(Grid::get(opCtx)->getExecutorPool()->getFixedExecutor(),
CancellationToken::uncancelable())
- .share();
+ .semi()
+ .get(opCtx);
}
}
-SharedSemiFuture<ShardRegistry::Cache::ValueHandle> ShardRegistry::_reloadAsyncNoRetry() {
+SharedSemiFuture<ShardRegistry::Cache::ValueHandle> ShardRegistry::_reloadInternal() {
// Make the next acquire do a lookup.
auto value = _forceReloadIncrement.addAndFetch(1);
LOGV2_DEBUG(4620253, 2, "Forcing ShardRegistry reload", "newForceReloadIncrement"_attr = value);
@@ -614,7 +569,7 @@ std::pair<ShardRegistryData, Timestamp> ShardRegistryData::createFromCatalogClie
OperationContext* opCtx, ShardFactory* shardFactory) {
auto const catalogClient = Grid::get(opCtx)->catalogClient();
- auto readConcern = repl::ReadConcernLevel::kSnapshotReadConcern;
+ auto readConcern = repl::ReadConcernLevel::kMajorityReadConcern;
// ShardRemote requires a majority read. We can only allow a non-majority read if we are a
// config server.