diff options
Diffstat (limited to 'src/mongo/db/s/resharding/resharding_recipient_service_test.cpp')
| -rw-r--r-- | src/mongo/db/s/resharding/resharding_recipient_service_test.cpp | 150 |
1 files changed, 94 insertions, 56 deletions
diff --git a/src/mongo/db/s/resharding/resharding_recipient_service_test.cpp b/src/mongo/db/s/resharding/resharding_recipient_service_test.cpp index 43ccab32bec..0b6c078603d 100644 --- a/src/mongo/db/s/resharding/resharding_recipient_service_test.cpp +++ b/src/mongo/db/s/resharding/resharding_recipient_service_test.cpp @@ -48,10 +48,13 @@ #include "mongo/db/s/resharding/resharding_recipient_service.h" #include "mongo/db/s/resharding/resharding_recipient_service_external_state.h" #include "mongo/db/s/resharding/resharding_service_test_helpers.h" +#include "mongo/db/s/sharding_ddl_util.h" #include "mongo/logv2/log.h" #include "mongo/unittest/death_test.h" +#include "mongo/util/clock_source_mock.h" #include "mongo/util/fail_point.h" + namespace mongo { namespace { @@ -208,6 +211,8 @@ public: */ class ReshardingRecipientServiceTest : public repl::PrimaryOnlyServiceMongoDTest { public: + ReshardingRecipientServiceTest() : PrimaryOnlyServiceMongoDTest(Options{}.useMockClock(true)) {} + using RecipientStateMachine = ReshardingRecipientService::RecipientStateMachine; std::unique_ptr<repl::PrimaryOnlyService> makeService(ServiceContext* serviceContext) override { @@ -222,7 +227,6 @@ public: repl::DropPendingCollectionReaper::set( serviceContext, std::make_unique<repl::DropPendingCollectionReaper>(storageMock.get())); repl::StorageInterface::set(serviceContext, std::move(storageMock)); - _controller = std::make_shared<RecipientStateTransitionController>(); _opObserverRegistry->addObserver(std::make_unique<RecipientOpObserverForTest>(_controller)); } @@ -263,7 +267,8 @@ public: const ReshardingRecipientDocument& recipientDoc) { CollectionOptions options; options.uuid = recipientDoc.getSourceUUID(); - resharding::data_copy::ensureCollectionDropped(opCtx, recipientDoc.getSourceNss()); + mongo::sharding_ddl_util::ensureCollectionDroppedNoChangeEvent(opCtx, + recipientDoc.getSourceNss()); resharding::data_copy::ensureCollectionExists(opCtx, recipientDoc.getSourceNss(), options); } @@ -496,69 +501,87 @@ DEATH_TEST_REGEX_F(ReshardingRecipientServiceTest, CommitFn, "4457001.*tripwire" TEST_F(ReshardingRecipientServiceTest, DropsTemporaryReshardingCollectionOnAbort) { auto metrics = ReshardingRecipientServiceTest::metrics(); for (bool isAlsoDonor : {false, true}) { - LOGV2(5551107, - "Running case", - "test"_attr = _agent.getTestName(), - "isAlsoDonor"_attr = isAlsoDonor); + for (bool waitForMetricsInitialized : {false, true}) { + LOGV2(5551107, + "Running case", + "test"_attr = _agent.getTestName(), + "isAlsoDonor"_attr = isAlsoDonor, + "waitForMetricsInitialized"_attr = waitForMetricsInitialized); + + boost::optional<PauseDuringStateTransitions> stateTransitionGuard; + if (waitForMetricsInitialized) { + std::vector<RecipientStateEnum> recipientStates{ + RecipientStateEnum::kDone, RecipientStateEnum::kCreatingCollection}; + stateTransitionGuard.emplace(controller(), recipientStates); + } else { + stateTransitionGuard.emplace(controller(), RecipientStateEnum::kDone); + } - boost::optional<PauseDuringStateTransitions> doneTransitionGuard; - doneTransitionGuard.emplace(controller(), RecipientStateEnum::kDone); + auto doc = makeStateDocument(isAlsoDonor); + auto instanceId = BSON(ReshardingRecipientDocument::kReshardingUUIDFieldName + << doc.getReshardingUUID()); - auto doc = makeStateDocument(isAlsoDonor); - auto instanceId = - BSON(ReshardingRecipientDocument::kReshardingUUIDFieldName << doc.getReshardingUUID()); + auto opCtx = makeOperationContext(); - auto opCtx = makeOperationContext(); + if (isAlsoDonor) { + // If the recipient is also a donor, the original collection should already exist on + // this shard. + createSourceCollection(opCtx.get(), doc); + } - if (isAlsoDonor) { - // If the recipient is also a donor, the original collection should already exist on - // this shard. - createSourceCollection(opCtx.get(), doc); - } + RecipientStateMachine::insertStateDocument(opCtx.get(), doc); + auto recipient = + RecipientStateMachine::getOrCreate(opCtx.get(), _service, doc.toBSON()); - RecipientStateMachine::insertStateDocument(opCtx.get(), doc); - auto recipient = RecipientStateMachine::getOrCreate(opCtx.get(), _service, doc.toBSON()); + notifyToStartCloning(opCtx.get(), *recipient, doc); + if (waitForMetricsInitialized) { + // Waiting for the metrics to be initialized here causes the second abort to occur + // before the metrics are initialized after the step up, thereby testing a different + // code path. + stateTransitionGuard->wait(RecipientStateEnum::kCreatingCollection); + } - notifyToStartCloning(opCtx.get(), *recipient, doc); - recipient->abort(false); + recipient->abort(false); - doneTransitionGuard->wait(RecipientStateEnum::kDone); - stepDown(); + stateTransitionGuard->wait(RecipientStateEnum::kDone); - ASSERT_EQ(recipient->getCompletionFuture().getNoThrow(), - ErrorCodes::InterruptedDueToReplStateChange); + stepDown(); - recipient.reset(); - stepUp(opCtx.get()); + ASSERT_EQ(recipient->getCompletionFuture().getNoThrow(), + ErrorCodes::InterruptedDueToReplStateChange); - auto maybeRecipient = RecipientStateMachine::lookup(opCtx.get(), _service, instanceId); - ASSERT_TRUE(bool(maybeRecipient)); - recipient = *maybeRecipient; + recipient.reset(); + stepUp(opCtx.get()); - doneTransitionGuard.reset(); - recipient->abort(false); + auto maybeRecipient = RecipientStateMachine::lookup(opCtx.get(), _service, instanceId); + ASSERT_TRUE(bool(maybeRecipient)); + recipient = *maybeRecipient; - ASSERT_OK(recipient->getCompletionFuture().getNoThrow()); - checkStateDocumentRemoved(opCtx.get()); + stateTransitionGuard.reset(); + recipient->abort(false); - if (isAlsoDonor) { - // Verify original collection still exists after aborting. - AutoGetCollection coll(opCtx.get(), doc.getSourceNss(), MODE_IS); - ASSERT_TRUE(bool(coll)); - ASSERT_EQ(coll->uuid(), doc.getSourceUUID()); - } + ASSERT_OK(recipient->getCompletionFuture().getNoThrow()); + checkStateDocumentRemoved(opCtx.get()); - // Verify the temporary collection no longer exists. - { - AutoGetCollection coll(opCtx.get(), doc.getTempReshardingNss(), MODE_IS); - ASSERT_FALSE(bool(coll)); + if (isAlsoDonor) { + // Verify original collection still exists after aborting. + AutoGetCollection coll(opCtx.get(), doc.getSourceNss(), MODE_IS); + ASSERT_TRUE(bool(coll)); + ASSERT_EQ(coll->uuid(), doc.getSourceUUID()); + } + + // Verify the temporary collection no longer exists. + { + AutoGetCollection coll(opCtx.get(), doc.getTempReshardingNss(), MODE_IS); + ASSERT_FALSE(bool(coll)); + } } } BSONObjBuilder result; metrics->serializeCumulativeOpMetrics(&result); - ASSERT_EQ(result.obj().getField("countReshardingFailures").numberLong(), 2); + ASSERT_LESS_THAN_OR_EQUALS(result.obj().getField("countReshardingFailures").numberLong(), 4); } TEST_F(ReshardingRecipientServiceTest, RenamesTemporaryReshardingCollectionWhenDone) { @@ -816,22 +839,33 @@ TEST_F(ReshardingRecipientServiceTest, RestoreMetricsAfterStepUp) { } // Step down before the transition to state can complete. stateTransitionsGuard.wait(state); - if (state == RecipientStateEnum::kStrictConsistency) { - auto currOp = recipient - ->reportForCurrentOp( - MongoProcessInterface::CurrentOpConnectionsMode::kExcludeIdle, - MongoProcessInterface::CurrentOpSessionsMode::kExcludeIdle) - .get(); + + dynamic_cast<ClockSourceMock*>(getServiceContext()->getFastClockSource()) + ->advance(Seconds(1)); + auto currOp = + recipient + ->reportForCurrentOp(MongoProcessInterface::CurrentOpConnectionsMode::kExcludeIdle, + MongoProcessInterface::CurrentOpSessionsMode::kExcludeIdle) + .get(); + + + if (state == RecipientStateEnum::kApplying) { + ASSERT_EQ(currOp.getField("totalApplyTimeElapsedSecs").Long(), 0); + ASSERT_EQ(currOp.getStringField("recipientState"), + RecipientState_serializer(RecipientStateEnum::kCloning)); + ASSERT_GT(currOp.getField("totalCopyTimeElapsedSecs").Long(), 0); + ASSERT_GT(currOp.getField("approxBytesToCopy").Long(), 0); + + } else if (state == RecipientStateEnum::kStrictConsistency) { ASSERT_EQ(currOp.getField("documentsCopied").Long(), 1L); ASSERT_EQ(currOp.getField("bytesCopied").Long(), (long)reshardedDoc.objsize()); ASSERT_EQ(currOp.getStringField("recipientState"), RecipientState_serializer(RecipientStateEnum::kApplying)); + ASSERT_GT(currOp.getField("totalCopyTimeElapsedSecs").Long(), 0); + ASSERT_GT(currOp.getField("totalApplyTimeElapsedSecs").Long(), 0); + ASSERT_GT(currOp.getField("approxBytesToCopy").Long(), 0); + } else if (state == RecipientStateEnum::kDone) { - auto currOp = recipient - ->reportForCurrentOp( - MongoProcessInterface::CurrentOpConnectionsMode::kExcludeIdle, - MongoProcessInterface::CurrentOpSessionsMode::kExcludeIdle) - .get(); ASSERT_EQ(currOp.getField("documentsCopied").Long(), 1L); ASSERT_EQ(currOp.getField("bytesCopied").Long(), (long)reshardedDoc.objsize()); ASSERT_EQ(currOp.getField("oplogEntriesFetched").Long(), @@ -840,7 +874,11 @@ TEST_F(ReshardingRecipientServiceTest, RestoreMetricsAfterStepUp) { oplogEntriesAppliedOnEachDonor * doc.getDonorShards().size()); ASSERT_EQ(currOp.getStringField("recipientState"), RecipientState_serializer(RecipientStateEnum::kStrictConsistency)); + ASSERT_GT(currOp.getField("totalCopyTimeElapsedSecs").Long(), 0); + ASSERT_GT(currOp.getField("totalApplyTimeElapsedSecs").Long(), 0); + ASSERT_GT(currOp.getField("approxBytesToCopy").Long(), 0); } + stepDown(); ASSERT_EQ(recipient->getCompletionFuture().getNoThrow(), |
