summaryrefslogtreecommitdiff
path: root/src/mongo/db/s/migration_destination_manager.cpp
diff options
context:
space:
mode:
Diffstat (limited to 'src/mongo/db/s/migration_destination_manager.cpp')
-rw-r--r--src/mongo/db/s/migration_destination_manager.cpp40
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;
}