summaryrefslogtreecommitdiff
path: root/src/mongo/s/client/shard_registry.cpp
diff options
context:
space:
mode:
authorLucas de Castro Borges <lucas@gnuabordo.com.br>2025-02-14 14:26:38 -0300
committerLucas de Castro Borges <lucas@gnuabordo.com.br>2025-02-14 14:26:38 -0300
commit294bc6ecabf14c09c9bc8644704921dcf97cb44e (patch)
tree279b1e0bab53901a1647ac63c1c724f0f789a663 /src/mongo/s/client/shard_registry.cpp
parent70be7c27a251621187a1de533462ae2bb1e3bd39 (diff)
parent1e917fd798aa25b7066d4b414b51184f13d5a092 (diff)
Update upstream source from tag 'upstream/6.0.10'debian/6.0.10-1
Update to upstream version '6.0.10' with Debian dir 2d176fa254eee97b139f712fec5709641335a8c3
Diffstat (limited to 'src/mongo/s/client/shard_registry.cpp')
-rw-r--r--src/mongo/s/client/shard_registry.cpp97
1 files changed, 71 insertions, 26 deletions
diff --git a/src/mongo/s/client/shard_registry.cpp b/src/mongo/s/client/shard_registry.cpp
index d722a9116ec..0b1691af82b 100644
--- a/src/mongo/s/client/shard_registry.cpp
+++ b/src/mongo/s/client/shard_registry.cpp
@@ -130,12 +130,8 @@ 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,
- fetchedFromConfigServers] = [&]()
- -> std::tuple<ShardRegistryData, Timestamp, Increment, ShardRegistryData::ShardMap, bool> {
+ auto [returnData, returnTopologyTime, returnForceReloadIncrement, removedShards] =
+ [&]() -> std::tuple<ShardRegistryData, Timestamp, Increment, ShardRegistryData::ShardMap> {
if (timeInStore.topologyTime > cachedData.getTime().topologyTime ||
timeInStore.forceReloadIncrement > cachedData.getTime().forceReloadIncrement) {
auto [reloadedData, maxTopologyTime] =
@@ -144,14 +140,12 @@ ShardRegistry::Cache::LookupResult ShardRegistry::_lookup(OperationContext* opCt
auto [mergedData, removedShards] =
ShardRegistryData::mergeExisting(*cachedData, reloadedData);
- return {
- mergedData, maxTopologyTime, timeInStore.forceReloadIncrement, removedShards, true};
+ return {mergedData, maxTopologyTime, timeInStore.forceReloadIncrement, removedShards};
} else {
return {*cachedData,
cachedData.getTime().topologyTime,
cachedData.getTime().forceReloadIncrement,
- {},
- false};
+ {}};
}
}();
@@ -186,11 +180,6 @@ 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,
@@ -218,9 +207,9 @@ void ShardRegistry::startupPeriodicReloader(OperationContext* opCtx) {
AsyncTry([this] {
LOGV2_DEBUG(22726, 1, "Reloading shardRegistry");
- return _reloadInternal();
+ return _reloadAsyncNoRetry();
})
- .until([](auto sw) {
+ .until([](auto&& sw) {
if (!sw.isOK()) {
LOGV2(22727,
"Error running periodic reload of shard registry",
@@ -232,7 +221,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",
@@ -295,6 +284,49 @@ 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()) {
@@ -380,8 +412,18 @@ std::unique_ptr<Shard> ShardRegistry::createConnection(const ConnectionString& c
return _shardFactory->createUniqueShard(ShardId("<unnamed>"), connStr);
}
-bool ShardRegistry::isUp() const {
- return _isUp.load();
+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;
}
void ShardRegistry::toBSON(BSONObjBuilder* result) const {
@@ -400,23 +442,26 @@ 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.
- _reloadInternal().get(opCtx);
+ return _reloadAsyncNoRetry();
} else {
- AsyncTry([=]() mutable { return _reloadInternal(); })
+ return AsyncTry([=]() mutable { return _reloadAsyncNoRetry(); })
.until([](auto sw) mutable {
return sw.getStatus() != ErrorCodes::ReadConcernMajorityNotAvailableYet;
})
.withBackoffBetweenIterations(kExponentialBackoff)
- .on(Grid::get(opCtx)->getExecutorPool()->getFixedExecutor(),
+ .on(Grid::get(getGlobalServiceContext())->getExecutorPool()->getFixedExecutor(),
CancellationToken::uncancelable())
- .semi()
- .get(opCtx);
+ .share();
}
}
-SharedSemiFuture<ShardRegistry::Cache::ValueHandle> ShardRegistry::_reloadInternal() {
+SharedSemiFuture<ShardRegistry::Cache::ValueHandle> ShardRegistry::_reloadAsyncNoRetry() {
// Make the next acquire do a lookup.
auto value = _forceReloadIncrement.addAndFetch(1);
LOGV2_DEBUG(4620253, 2, "Forcing ShardRegistry reload", "newForceReloadIncrement"_attr = value);