diff options
Diffstat (limited to 'src/mongo/db/s/migration_destination_manager.cpp')
| -rw-r--r-- | src/mongo/db/s/migration_destination_manager.cpp | 40 |
1 files changed, 34 insertions, 6 deletions
diff --git a/src/mongo/db/s/migration_destination_manager.cpp b/src/mongo/db/s/migration_destination_manager.cpp index 2e748cf3ca9..1c02208ee27 100644 --- a/src/mongo/db/s/migration_destination_manager.cpp +++ b/src/mongo/db/s/migration_destination_manager.cpp @@ -148,10 +148,14 @@ bool opReplicatedEnough(OperationContext* txn, const repl::OpTime& lastOpApplied, const WriteConcernOptions& writeConcern) { WriteConcernResult writeConcernResult; + writeConcernResult.wTimedOut = false; - Status waitForMajorityWriteConcernStatus = + Status majorityStatus = waitForWriteConcern(txn, lastOpApplied, kMajorityWriteConcern, &writeConcernResult); - if (!waitForMajorityWriteConcernStatus.isOK()) { + if (!majorityStatus.isOK()) { + if (!writeConcernResult.wTimedOut) { + uassertStatusOK(majorityStatus); + } return false; } @@ -159,13 +163,16 @@ bool opReplicatedEnough(OperationContext* txn, // write concerns in case the user's write concern is stronger than majority WriteConcernOptions userWriteConcern(writeConcern); userWriteConcern.wTimeout = -1; + writeConcernResult.wTimedOut = false; - Status waitForUserWriteConcernStatus = + Status userStatus = waitForWriteConcern(txn, lastOpApplied, userWriteConcern, &writeConcernResult); - if (!waitForUserWriteConcernStatus.isOK()) { + if (!userStatus.isOK()) { + if (!writeConcernResult.wTimedOut) { + uassertStatusOK(userStatus); + } return false; } - return true; } @@ -220,6 +227,7 @@ MigrationDestinationManager::State MigrationDestinationManager::getState() const void MigrationDestinationManager::setState(State newState) { stdx::lock_guard<stdx::mutex> sl(_mutex); _state = newState; + _stateChangedCV.notify_all(); } bool MigrationDestinationManager::isActive() const { @@ -231,7 +239,21 @@ bool MigrationDestinationManager::_isActive_inlock() const { return _sessionId.is_initialized(); } -void MigrationDestinationManager::report(BSONObjBuilder& b) { +void MigrationDestinationManager::report(BSONObjBuilder& b, + OperationContext* opCtx, + bool waitForSteadyOrDone) { + if (waitForSteadyOrDone) { + stdx::unique_lock<stdx::mutex> lock(_mutex); + try { + opCtx->waitForConditionOrInterruptFor(_stateChangedCV, lock, Seconds(1), [&]() -> bool { + return _state != READY && _state != CLONE && _state != CATCHUP; + }); + } catch (...) { + // Ignoring this error because this is an optional parameter and we catch timeout + // exceptions later. + } + b.append("waited", true); + } stdx::lock_guard<stdx::mutex> sl(_mutex); b.appendBool("active", _sessionId.is_initialized()); @@ -286,6 +308,7 @@ Status MigrationDestinationManager::start(const NamespaceString& nss, invariant(!_scopedRegisterReceiveChunk); _state = READY; + _stateChangedCV.notify_all(); _errmsg = ""; _nss = nss; @@ -336,6 +359,7 @@ bool MigrationDestinationManager::abort(const MigrationSessionId& sessionId) { } _state = ABORT; + _stateChangedCV.notify_all(); _errmsg = "aborted"; return true; @@ -344,6 +368,7 @@ bool MigrationDestinationManager::abort(const MigrationSessionId& sessionId) { void MigrationDestinationManager::abortWithoutSessionIdCheck() { stdx::lock_guard<stdx::mutex> sl(_mutex); _state = ABORT; + _stateChangedCV.notify_all(); _errmsg = "aborted without session id check"; } @@ -368,6 +393,7 @@ bool MigrationDestinationManager::startCommit(const MigrationSessionId& sessionI } _state = COMMIT_START; + _stateChangedCV.notify_all(); const auto deadline = Date_t::now() + Seconds(30); @@ -375,6 +401,8 @@ bool MigrationDestinationManager::startCommit(const MigrationSessionId& sessionI if (stdx::cv_status::timeout == _isActiveCV.wait_until(lock, deadline.toSystemTimePoint())) { _state = FAIL; + _stateChangedCV.notify_all(); + log() << "startCommit never finished!" << migrateLog; return false; } |
