diff options
Diffstat (limited to 'src/mongo/db/serverless')
6 files changed, 64 insertions, 28 deletions
diff --git a/src/mongo/db/serverless/SConscript b/src/mongo/db/serverless/SConscript index a3151a9da15..ad92d014d37 100644 --- a/src/mongo/db/serverless/SConscript +++ b/src/mongo/db/serverless/SConscript @@ -51,7 +51,6 @@ env.Library( 'shard_split_utils.cpp' ], LIBDEPS_PRIVATE=[ - '$BUILD_DIR/mongo/db/concurrency/exception_util', '$BUILD_DIR/mongo/db/dbhelpers', '$BUILD_DIR/mongo/db/repl/replica_set_messages', 'shard_split_state_machine', @@ -71,7 +70,6 @@ env.Library( ], LIBDEPS_PRIVATE=[ '$BUILD_DIR/mongo/db/catalog/local_oplog_info', - '$BUILD_DIR/mongo/db/concurrency/exception_util', '$BUILD_DIR/mongo/db/db_raii', '$BUILD_DIR/mongo/db/dbhelpers', '$BUILD_DIR/mongo/db/namespace_string', diff --git a/src/mongo/db/serverless/shard_split_commands.cpp b/src/mongo/db/serverless/shard_split_commands.cpp index 34e44a72a30..caf3a44cfde 100644 --- a/src/mongo/db/serverless/shard_split_commands.cpp +++ b/src/mongo/db/serverless/shard_split_commands.cpp @@ -219,7 +219,7 @@ public: auto splitService = repl::PrimaryOnlyServiceRegistry::get(opCtx->getServiceContext()) ->lookupServiceByName(ShardSplitDonorService::kServiceName); - auto [optionalDonor, _] = ShardSplitDonorService::DonorStateMachine::lookup( + auto optionalDonor = ShardSplitDonorService::DonorStateMachine::lookup( opCtx, splitService, BSON("_id" << cmd.getMigrationId())); uassert(ErrorCodes::NoSuchTenantMigration, diff --git a/src/mongo/db/serverless/shard_split_donor_op_observer.h b/src/mongo/db/serverless/shard_split_donor_op_observer.h index d8b607db6df..55796479ffe 100644 --- a/src/mongo/db/serverless/shard_split_donor_op_observer.h +++ b/src/mongo/db/serverless/shard_split_donor_op_observer.h @@ -204,10 +204,6 @@ public: size_t numberOfPrePostImagesToWrite, Date_t wallClockTime) final {} - void onTransactionPrepareNonPrimary(OperationContext* opCtx, - const std::vector<repl::OplogEntry>& statements, - const repl::OpTime& prepareOpTime) final {} - void onTransactionAbort(OperationContext* opCtx, boost::optional<OplogSlot> abortOplogEntryOpTime) final {} diff --git a/src/mongo/db/serverless/shard_split_donor_service.cpp b/src/mongo/db/serverless/shard_split_donor_service.cpp index 7b15d250450..633fc00e974 100644 --- a/src/mongo/db/serverless/shard_split_donor_service.cpp +++ b/src/mongo/db/serverless/shard_split_donor_service.cpp @@ -33,7 +33,7 @@ #include "mongo/db/serverless/shard_split_donor_service.h" #include "mongo/client/streamable_replica_set_monitor.h" #include "mongo/db/catalog_raii.h" -#include "mongo/db/concurrency/exception_util.h" +#include "mongo/db/concurrency/write_conflict_exception.h" #include "mongo/db/dbdirectclient.h" #include "mongo/db/dbhelpers.h" #include "mongo/db/repl/repl_client_info.h" diff --git a/src/mongo/db/serverless/shard_split_donor_service_test.cpp b/src/mongo/db/serverless/shard_split_donor_service_test.cpp index d48274ac15e..2d44ded9b7a 100644 --- a/src/mongo/db/serverless/shard_split_donor_service_test.cpp +++ b/src/mongo/db/serverless/shard_split_donor_service_test.cpp @@ -324,7 +324,7 @@ void processHelloRequest(executor::NetworkInterfaceMock* net, MockReplicaSet* re auto noi = net->getNextReadyRequest(); auto request = noi->getRequest(); - assertRemoteCommandNameEquals("hello", request); + assertRemoteCommandNameEquals("isMaster", request); auto requestHost = request.target.toString(); const auto node = replSet->getNode(requestHost); if (node->isRunning()) { @@ -343,6 +343,55 @@ void waitForHelloRequest(executor::NetworkInterfaceMock* net) { } } +TEST_F(ShardSplitDonorServiceTest, BasicShardSplitDonorServiceInstanceCreation) { + auto opCtx = makeOperationContext(); + test::shard_split::ScopedTenantAccessBlocker scopedTenants(_tenantIds, opCtx.get()); + test::shard_split::reconfigToAddRecipientNodes( + getServiceContext(), _recipientTagName, _replSet.getHosts(), _recipientSet.getHosts()); + + ShardSplitDonorService::DonorStateMachine::setSplitAcceptanceTaskExecutor_forTest(_executor); + _skipAcceptanceFP.reset(); + + // Create and start the instance. + auto serviceInstance = ShardSplitDonorService::DonorStateMachine::getOrCreate( + opCtx.get(), _service, defaultStateDocument().toBSON()); + ASSERT(serviceInstance.get()); + ASSERT_EQ(_uuid, serviceInstance->getId()); + + // Wait for monitors to start, and enqueue successfull hello responses + _net->enterNetwork(); + waitForHelloRequest(_net); + processHelloRequest(_net, &_recipientSet); + waitForHelloRequest(_net); + processHelloRequest(_net, &_recipientSet); + waitForHelloRequest(_net); + processHelloRequest(_net, &_recipientSet); + _net->runReadyNetworkOperations(); + _net->exitNetwork(); + + auto decisionFuture = serviceInstance->decisionFuture(); + decisionFuture.wait(); + + ASSERT_TRUE(hasActiveSplitForTenants(opCtx.get(), _tenantIds)); + + auto result = decisionFuture.get(); + ASSERT(!result.abortReason); + ASSERT_EQ(result.state, mongo::ShardSplitDonorStateEnum::kCommitted); + + BSONObj splitConfigBson = mockReplSetReconfigCmd.getLatestConfig(); + ASSERT_TRUE(splitConfigBson.hasField("replSetReconfig")); + auto splitConfig = repl::ReplSetConfig::parse(splitConfigBson["replSetReconfig"].Obj()); + ASSERT(splitConfig.isSplitConfig()); + + serviceInstance->tryForget(); + + auto completionFuture = serviceInstance->completionFuture(); + completionFuture.wait(); + + ASSERT_OK(serviceInstance->completionFuture().getNoThrow()); + ASSERT_TRUE(serviceInstance->isGarbageCollectable()); +} + TEST_F(ShardSplitDonorServiceTest, ShardSplitDonorServiceTimeout) { FailPointEnableBlock fp("pauseShardSplitAfterBlocking"); @@ -548,8 +597,7 @@ TEST(RecipientAcceptSplitListenerTest, FutureReady) { ASSERT_FALSE(listener.getFuture().isReady()); listener.onServerHeartbeatSucceededEvent( - donor.getHosts().front(), - BSON("setName" << donor.getSetName() << "isWritablePrimary" << true)); + donor.getHosts().front(), BSON("setName" << donor.getSetName() << "ismaster" << true)); ASSERT_TRUE(listener.getFuture().isReady()); } @@ -570,8 +618,7 @@ TEST(RecipientAcceptSplitListenerTest, FutureReadyNameChange) { listener.onServerHeartbeatSucceededEvent(host, BSON("setName" << donor.getSetName())); } listener.onServerHeartbeatSucceededEvent( - donor.getHosts().front(), - BSON("setName" << donor.getSetName() << "isWritablePrimary" << true)); + donor.getHosts().front(), BSON("setName" << donor.getSetName() << "ismaster" << true)); ASSERT_TRUE(listener.getFuture().isReady()); } @@ -593,7 +640,7 @@ TEST(RecipientAcceptSplitListenerTest, FutureNotReadyMissingNodes) { ASSERT_FALSE(listener.getFuture().isReady()); donor.setPrimary(donor.getHosts()[0].host()); listener.onServerHeartbeatSucceededEvent( - donor.getHosts()[0], BSON("setName" << donor.getSetName() << "isWritablePrimary" << true)); + donor.getHosts()[0], BSON("setName" << donor.getSetName() << "ismaster" << true)); ASSERT_TRUE(listener.getFuture().isReady()); } @@ -660,10 +707,9 @@ TEST_F(ShardSplitDonorServiceTest, ResumeAfterStepdownTest) { // verify that the state document exists ASSERT_OK(getStateDocument(opCtx.get(), _uuid).getStatus()); - auto [donor, isPausedOrShutdown] = ShardSplitDonorService::DonorStateMachine::lookup( + auto donor = ShardSplitDonorService::DonorStateMachine::lookup( opCtx.get(), _service, BSON("_id" << _uuid)); - ASSERT_TRUE(donor); - ASSERT_FALSE(isPausedOrShutdown); + ASSERT(donor); fp.reset(); @@ -743,12 +789,9 @@ TEST_F(ShardSplitRecipientCleanupTest, ShardSplitRecipientCleanup) { auto splitService = repl::PrimaryOnlyServiceRegistry::get(opCtx->getServiceContext()) ->lookupServiceByName(ShardSplitDonorService::kServiceName); - auto [optionalDonor, isPausedOrShutdown] = - ShardSplitDonorService::DonorStateMachine::lookup( - opCtx.get(), splitService, BSON("_id" << _uuid)); + auto optionalDonor = ShardSplitDonorService::DonorStateMachine::lookup( + opCtx.get(), splitService, BSON("_id" << _uuid)); - ASSERT_TRUE(optionalDonor); - ASSERT_FALSE(isPausedOrShutdown); ASSERT_TRUE(hasActiveSplitForTenants(opCtx.get(), _tenantIds)); ASSERT_TRUE(optionalDonor); @@ -815,14 +858,13 @@ TEST_F(ShardSplitStepUpWithCommitted, StepUpWithkCommitted) { auto splitService = repl::PrimaryOnlyServiceRegistry::get(opCtx->getServiceContext()) ->lookupServiceByName(ShardSplitDonorService::kServiceName); - auto [optionalDonor, isPausedOrShutdown] = ShardSplitDonorService::DonorStateMachine::lookup( + auto optionalInstance = ShardSplitDonorService::DonorStateMachine::lookup( opCtx.get(), splitService, BSON("_id" << _uuid)); - ASSERT_TRUE(optionalDonor); - ASSERT_FALSE(isPausedOrShutdown); + ASSERT(optionalInstance); _pauseBeforeRecipientCleanupFp.reset(); - auto serviceInstance = optionalDonor->get(); + auto serviceInstance = optionalInstance->get(); auto result = serviceInstance->decisionFuture().get(); diff --git a/src/mongo/db/serverless/shard_split_utils.cpp b/src/mongo/db/serverless/shard_split_utils.cpp index 2380f9bb427..0af9ccc6195 100644 --- a/src/mongo/db/serverless/shard_split_utils.cpp +++ b/src/mongo/db/serverless/shard_split_utils.cpp @@ -29,8 +29,8 @@ #include "mongo/db/serverless/shard_split_utils.h" #include "mongo/db/catalog_raii.h" -#include "mongo/db/concurrency/exception_util.h" #include "mongo/db/concurrency/lock_manager_defs.h" +#include "mongo/db/concurrency/write_conflict_exception.h" #include "mongo/db/db_raii.h" #include "mongo/db/dbhelpers.h" #include "mongo/db/ops/delete.h" @@ -279,7 +279,7 @@ void RecipientAcceptSplitListener::onServerHeartbeatSucceededEvent(const HostAnd _reportedSetNames[hostAndPort] = reply["setName"].str(); - if (!_hasPrimary && reply["isWritablePrimary"].booleanSafe()) { + if (!_hasPrimary && reply["ismaster"].booleanSafe()) { _hasPrimary = true; } |
