diff options
Diffstat (limited to 'src/mongo/db/s/sharding_ddl_util.cpp')
| -rw-r--r-- | src/mongo/db/s/sharding_ddl_util.cpp | 208 |
1 files changed, 25 insertions, 183 deletions
diff --git a/src/mongo/db/s/sharding_ddl_util.cpp b/src/mongo/db/s/sharding_ddl_util.cpp index ba13de4abbd..d516fd5e668 100644 --- a/src/mongo/db/s/sharding_ddl_util.cpp +++ b/src/mongo/db/s/sharding_ddl_util.cpp @@ -33,7 +33,6 @@ #include "mongo/db/catalog/collection_catalog.h" #include "mongo/db/commands/feature_compatibility_version.h" -#include "mongo/db/concurrency/exception_util.h" #include "mongo/db/db_raii.h" #include "mongo/db/dbdirectclient.h" #include "mongo/db/repl/repl_client_info.h" @@ -77,14 +76,11 @@ void updateTags(OperationContext* opCtx, }()}); return updateOp; }()); + request.setWriteConcern(writeConcern.toBSON()); auto configShard = Grid::get(opCtx)->shardRegistry()->getConfigShard(); - auto response = - configShard->runBatchWriteCommand(opCtx, - Milliseconds::max(), - request, - writeConcern, - Shard::RetryPolicy::kIdempotentOrCursorInvalidated); + auto response = configShard->runBatchWriteCommand( + opCtx, Milliseconds::max(), request, Shard::RetryPolicy::kIdempotentOrCursorInvalidated); uassertStatusOK(response.toStatus()); } @@ -108,13 +104,11 @@ void deleteChunks(OperationContext* opCtx, return deleteOp; }()); + request.setWriteConcern(writeConcern.toBSON()); + auto configShard = Grid::get(opCtx)->shardRegistry()->getConfigShard(); - auto response = - configShard->runBatchWriteCommand(opCtx, - Milliseconds::max(), - request, - writeConcern, - Shard::RetryPolicy::kIdempotentOrCursorInvalidated); + auto response = configShard->runBatchWriteCommand( + opCtx, Milliseconds::max(), request, Shard::RetryPolicy::kIdempotentOrCursorInvalidated); uassertStatusOK(response.toStatus()); } @@ -179,93 +173,6 @@ void setAllowMigrations(OperationContext* opCtx, } } - -// Check that the collection UUID is the same in every shard knowing the collection -void checkCollectionUUIDConsistencyAcrossShards( - OperationContext* opCtx, - const NamespaceString& nss, - const UUID& collectionUuid, - const std::vector<mongo::ShardId>& shardIds, - std::shared_ptr<executor::ScopedTaskExecutor> executor) { - const BSONObj filterObj = BSON("name" << nss.coll()); - BSONObj cmdObj = BSON("listCollections" << 1 << "filter" << filterObj); - - auto responses = sharding_ddl_util::sendAuthenticatedCommandToShards( - opCtx, nss.db().toString(), cmdObj, shardIds, **executor); - - struct MismatchedShard { - std::string shardId; - std::string uuid; - }; - - std::vector<MismatchedShard> mismatches; - - for (auto cmdResponse : responses) { - auto responseData = uassertStatusOK(cmdResponse.swResponse); - auto collectionVector = responseData.data.firstElement()["firstBatch"].Array(); - auto shardId = cmdResponse.shardId; - - if (collectionVector.empty()) { - // Collection does not exist on the shard - continue; - } - - auto bsonCollectionUuid = collectionVector.front()["info"]["uuid"]; - if (collectionUuid.data() != bsonCollectionUuid.uuid()) { - mismatches.push_back({shardId.toString(), bsonCollectionUuid.toString()}); - } - } - - if (!mismatches.empty()) { - std::stringstream errorMessage; - errorMessage << "The collection " << nss.toString() - << " with expected UUID: " << collectionUuid.toString() - << " has different UUIDs on the following shards: ["; - - for (auto mismatch : mismatches) { - errorMessage << "{ " << mismatch.shardId << ":" << mismatch.uuid << " },"; - } - errorMessage << "]"; - uasserted(ErrorCodes::InvalidUUID, errorMessage.str()); - } -} - - -// Check the collection does not exist in any shard when `dropTarget` is set to false -void checkTargetCollectionDoesNotExistInCluster( - OperationContext* opCtx, - const NamespaceString& toNss, - const std::vector<mongo::ShardId>& shardIds, - std::shared_ptr<executor::ScopedTaskExecutor> executor) { - const BSONObj filterObj = BSON("name" << toNss.coll()); - BSONObj cmdObj = BSON("listCollections" << 1 << "filter" << filterObj); - - auto responses = sharding_ddl_util::sendAuthenticatedCommandToShards( - opCtx, toNss.db(), cmdObj, shardIds, **executor); - - std::vector<std::string> shardsContainingTargetCollection; - for (auto cmdResponse : responses) { - uassertStatusOK(cmdResponse.swResponse); - auto responseData = uassertStatusOK(cmdResponse.swResponse); - auto collectionVector = responseData.data.firstElement()["firstBatch"].Array(); - - if (!collectionVector.empty()) { - shardsContainingTargetCollection.push_back(cmdResponse.shardId.toString()); - } - } - - if (!shardsContainingTargetCollection.empty()) { - std::stringstream errorMessage; - errorMessage << "The collection " << toNss.toString() - << " already exists in the following shards: ["; - std::move(shardsContainingTargetCollection.begin(), - shardsContainingTargetCollection.end(), - std::ostream_iterator<std::string>(errorMessage, ", ")); - errorMessage << "]"; - uasserted(ErrorCodes::NamespaceExists, errorMessage.str()); - } -} - } // namespace void linearizeCSRSReads(OperationContext* opCtx) { @@ -285,20 +192,27 @@ std::vector<AsyncRequestsSender::Response> sendAuthenticatedCommandToShards( StringData dbName, const BSONObj& command, const std::vector<ShardId>& shardIds, - const std::shared_ptr<executor::TaskExecutor>& executor, - const bool throwOnError) { + const std::shared_ptr<executor::TaskExecutor>& executor) { + // TODO SERVER-57519: remove the following scope + { + // Ensure ShardRegistry is initialized before using the AsyncRequestsSender that relies on + // unsafe functions (SERVER-57280) + auto shardRegistry = Grid::get(opCtx)->shardRegistry(); + if (!shardRegistry->isUp()) { + shardRegistry->reload(opCtx); + } + } // The AsyncRequestsSender ignore impersonation metadata so we need to manually attach them to // the command BSONObjBuilder bob(command); rpc::writeAuthDataToImpersonatedUserMetadata(opCtx, &bob); - if (serverGlobalParams.featureCompatibility.isVersionInitialized() && - gFeatureFlagUserWriteBlocking.isEnabled(serverGlobalParams.featureCompatibility)) { + if (gFeatureFlagUserWriteBlocking.isEnabled(serverGlobalParams.featureCompatibility)) { WriteBlockBypass::get(opCtx).writeAsMetadata(&bob); } auto authenticatedCommand = bob.obj(); return sharding_util::sendCommandToShards( - opCtx, dbName, authenticatedCommand, shardIds, executor, throwOnError); + opCtx, dbName, authenticatedCommand, shardIds, executor); } void removeTagsMetadataFromConfig(OperationContext* opCtx, @@ -341,13 +255,11 @@ void removeTagsMetadataFromConfig_notIdempotent(OperationContext* opCtx, return deleteOp; }()); + request.setWriteConcern(writeConcern.toBSON()); + auto configShard = Grid::get(opCtx)->shardRegistry()->getConfigShard(); - auto response = - configShard->runBatchWriteCommand(opCtx, - Milliseconds::max(), - request, - writeConcern, - Shard::RetryPolicy::kIdempotentOrCursorInvalidated); + auto response = configShard->runBatchWriteCommand( + opCtx, Milliseconds::max(), request, Shard::RetryPolicy::kIdempotentOrCursorInvalidated); uassertStatusOK(response.toStatus()); } @@ -440,24 +352,6 @@ void shardedRenameMetadata(OperationContext* opCtx, opCtx, CollectionType::ConfigNS, fromCollType.toBSON(), writeConcern)); } -void checkCatalogConsistencyAcrossShardsForRename( - OperationContext* opCtx, - const NamespaceString& fromNss, - const NamespaceString& toNss, - const bool dropTarget, - std::shared_ptr<executor::ScopedTaskExecutor> executor) { - - auto participants = Grid::get(opCtx)->shardRegistry()->getAllShardIds(opCtx); - - auto sourceCollUuid = *getCollectionUUID(opCtx, fromNss); - checkCollectionUUIDConsistencyAcrossShards( - opCtx, fromNss, sourceCollUuid, participants, executor); - - if (!dropTarget) { - checkTargetCollectionDoesNotExistInCluster(opCtx, toNss, participants, executor); - } -} - void checkRenamePreconditions(OperationContext* opCtx, bool sourceIsSharded, const NamespaceString& toNss, @@ -557,26 +451,6 @@ void resumeMigrations(OperationContext* opCtx, setAllowMigrations(opCtx, nss, expectedCollectionUUID, true); } -bool checkAllowMigrations(OperationContext* opCtx, const NamespaceString& nss) { - auto collDoc = - uassertStatusOK(Grid::get(opCtx)->shardRegistry()->getConfigShard()->exhaustiveFindOnConfig( - opCtx, - ReadPreferenceSetting(ReadPreference::PrimaryOnly, TagSet{}), - repl::ReadConcernLevel::kMajorityReadConcern, - CollectionType::ConfigNS, - BSON(CollectionType::kNssFieldName << nss.ns()), - BSONObj(), - 1)) - .docs; - - uassert(ErrorCodes::NamespaceNotFound, - str::stream() << "collection " << nss.ns() << " not found", - !collDoc.empty()); - - auto coll = CollectionType(collDoc[0]); - return coll.getAllowMigrations(); -} - boost::optional<UUID> getCollectionUUID(OperationContext* opCtx, const NamespaceString& nss, bool allowViews) { @@ -625,11 +499,8 @@ void sendDropCollectionParticipantCommandToShards(OperationContext* opCtx, const NamespaceString& nss, const std::vector<ShardId>& shardIds, std::shared_ptr<executor::TaskExecutor> executor, - const OperationSessionInfo& osi, - bool fromMigrate) { - ShardsvrDropCollectionParticipant dropCollectionParticipant(nss); - dropCollectionParticipant.setFromMigrate(fromMigrate); - + const OperationSessionInfo& osi) { + const ShardsvrDropCollectionParticipant dropCollectionParticipant(nss); const auto cmdObj = CommandHelpers::appendMajorityWriteConcern(dropCollectionParticipant.toBSON({})); @@ -645,34 +516,5 @@ void sendDropCollectionParticipantCommandToShards(OperationContext* opCtx, } } -BSONObj getCriticalSectionReasonForRename(const NamespaceString& from, const NamespaceString& to) { - return BSON("command" - << "rename" - << "from" << from.toString() << "to" << to.toString()); -} - -void ensureCollectionDroppedNoChangeEvent(OperationContext* opCtx, - const NamespaceString& nss, - const boost::optional<UUID>& uuid) { - invariant(!opCtx->lockState()->isLocked()); - invariant(!opCtx->lockState()->inAWriteUnitOfWork()); - - writeConflictRetry(opCtx, - "mongo::sharding_ddl_util::ensureCollectionDroppedNoChangeEvent", - nss.toString(), - [&] { - AutoGetCollection coll(opCtx, nss, MODE_X); - if (!coll || (uuid && coll->uuid() != uuid)) { - // If the collection doesn't exist or exists with a different UUID, - // then the requested collection has been dropped already. - return; - } - - WriteUnitOfWork wuow(opCtx); - uassertStatusOK(coll.getDb()->dropCollectionEvenIfSystem( - opCtx, nss, {} /* dropOpTime */, true /* markFromMigrate */)); - wuow.commit(); - }); -} } // namespace sharding_ddl_util } // namespace mongo |
