From 333f2e48655168b2bf1b25a8d07f68917057ac15 Mon Sep 17 00:00:00 2001 From: Alvin Moore Date: Wed, 31 May 2017 14:21:50 -0700 Subject: [PATCH 1/2] Added connection fault disabler within setup of backup submission. It should be reviewed to determine the amount of time to wait before disabling --- fdbserver/workloads/AtomicSwitchover.actor.cpp | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/fdbserver/workloads/AtomicSwitchover.actor.cpp b/fdbserver/workloads/AtomicSwitchover.actor.cpp index 0ef27bd83d..9134303c0f 100644 --- a/fdbserver/workloads/AtomicSwitchover.actor.cpp +++ b/fdbserver/workloads/AtomicSwitchover.actor.cpp @@ -56,6 +56,7 @@ struct AtomicSwitchoverWorkload : TestWorkload { ACTOR static Future _setup(Database cx, AtomicSwitchoverWorkload* self) { state DatabaseBackupAgent backupAgent(cx); + state Future disabler = disableConnectionFailuresAfter(300, "atomicSwitchover"); try { TraceEvent("AS_Submit1"); Void _ = wait( backupAgent.submitBackup(self->extraDB, BackupAgentBase::getDefaultTag(), self->backupRanges, false, StringRef(), StringRef(), true) ); @@ -96,7 +97,7 @@ struct AtomicSwitchoverWorkload : TestWorkload { auto src = srcFuture.get().begin(); auto bkp = bkpFuture.get().begin(); - + while (src != srcFuture.get().end() && bkp != bkpFuture.get().end()) { KeyRef bkpKey = bkp->key.substr(backupPrefix.size()); if (src->key != bkpKey && src->value != bkp->value) { @@ -156,7 +157,7 @@ struct AtomicSwitchoverWorkload : TestWorkload { state Future switch2After = delay(self->switch2After); state Future stopAfter = delay(self->stopAfter); state Future disabler = disableConnectionFailuresAfter(300, "atomicSwitchover"); - + TraceEvent("AS_Wait1"); int _ = wait( backupAgent.waitBackup(self->extraDB, BackupAgentBase::getDefaultTag(), false) ); TraceEvent("AS_Ready1"); @@ -186,4 +187,4 @@ struct AtomicSwitchoverWorkload : TestWorkload { } }; -WorkloadFactory AtomicSwitchoverWorkloadFactory("AtomicSwitchover"); \ No newline at end of file +WorkloadFactory AtomicSwitchoverWorkloadFactory("AtomicSwitchover"); From 1626e16377caaf406c6647c33326ecdc5dc38670 Mon Sep 17 00:00:00 2001 From: Evan Tschannen Date: Wed, 31 May 2017 16:23:37 -0700 Subject: [PATCH 2/2] Merge branch 'release-4.6' into release-5.0 --- fdbclient/DatabaseBackupAgent.actor.cpp | 7 ++-- fdbserver/masterserver.actor.cpp | 8 +++-- fdbserver/workloads/ApiWorkload.actor.cpp | 10 +++--- fdbserver/workloads/ApiWorkload.h | 35 +++++++++++++++---- fdbserver/workloads/ChangeConfig.actor.cpp | 24 ++++++++++--- fdbserver/workloads/WriteDuringRead.actor.cpp | 5 +++ tests/slow/ApiCorrectnessSwitchover.txt | 30 ++++++++++++++++ 7 files changed, 99 insertions(+), 20 deletions(-) create mode 100644 tests/slow/ApiCorrectnessSwitchover.txt diff --git a/fdbclient/DatabaseBackupAgent.actor.cpp b/fdbclient/DatabaseBackupAgent.actor.cpp index 21ed3b556d..827eb33624 100755 --- a/fdbclient/DatabaseBackupAgent.actor.cpp +++ b/fdbclient/DatabaseBackupAgent.actor.cpp @@ -1375,6 +1375,7 @@ public: } ACTOR static Future atomicSwitchover(DatabaseBackupAgent* backupAgent, Database dest, Key tagName, Standalone> backupRanges, Key addPrefix, Key removePrefix) { + state DatabaseBackupAgent drAgent(dest); state UID destlogUid = wait(backupAgent->getLogUid(dest, tagName)); state int status = wait(backupAgent->getStateValue(dest, destlogUid)); @@ -1385,7 +1386,7 @@ public: state UID logUid = g_random->randomUniqueID(); state Key logUidValue = BinaryWriter::toValue(logUid, Unversioned()); - state UID logUidCurrent = wait(backupAgent->getLogUid(backupAgent->taskBucket->src, tagName)); + state UID logUidCurrent = wait(drAgent.getLogUid(backupAgent->taskBucket->src, tagName)); if (logUidCurrent.isValid()) { logUid = logUidCurrent; @@ -1446,7 +1447,7 @@ public: TraceEvent("DBA_switchover_stopped"); try { - Void _ = wait( backupAgent->submitBackup(backupAgent->taskBucket->src, tagName, backupRanges, false, addPrefix, removePrefix, true, true) ); + Void _ = wait( drAgent.submitBackup(backupAgent->taskBucket->src, tagName, backupRanges, false, addPrefix, removePrefix, true, true) ); } catch( Error &e ) { if( e.code() != error_code_backup_duplicate ) throw; @@ -1454,7 +1455,7 @@ public: TraceEvent("DBA_switchover_submitted"); - int _ = wait( backupAgent->waitSubmitted(backupAgent->taskBucket->src, tagName) ); + int _ = wait( drAgent.waitSubmitted(backupAgent->taskBucket->src, tagName) ); TraceEvent("DBA_switchover_started"); diff --git a/fdbserver/masterserver.actor.cpp b/fdbserver/masterserver.actor.cpp index cb90ab83e4..622b5fb59e 100644 --- a/fdbserver/masterserver.actor.cpp +++ b/fdbserver/masterserver.actor.cpp @@ -549,10 +549,14 @@ ACTOR Future readTransactionSystemState( Reference self, Refer // Recover version info self->lastEpochEnd = oldLogSystem->getEnd() - 1; - if (self->lastEpochEnd == 0) + if (self->lastEpochEnd == 0) { self->recoveryTransactionVersion = 1; - else + } else { self->recoveryTransactionVersion = self->lastEpochEnd + SERVER_KNOBS->MAX_VERSIONS_IN_FLIGHT; + if(BUGGIFY) { + self->recoveryTransactionVersion += g_random->randomInt64(0, 1e6*SERVER_KNOBS->VERSIONS_PER_SECOND); + } + } TraceEvent("MasterRecovering", self->dbgid).detail("lastEpochEnd", self->lastEpochEnd).detail("recoveryTransactionVersion", self->recoveryTransactionVersion); diff --git a/fdbserver/workloads/ApiWorkload.actor.cpp b/fdbserver/workloads/ApiWorkload.actor.cpp index a11345d79d..d8af3de22c 100644 --- a/fdbserver/workloads/ApiWorkload.actor.cpp +++ b/fdbserver/workloads/ApiWorkload.actor.cpp @@ -47,7 +47,7 @@ Future ApiWorkload::clearKeyspace() { ACTOR Future setup(Database cx, ApiWorkload *self) { state Future disabler = disableConnectionFailuresAfter(300, "ApiWorkload"); - self->transactionFactory = Reference(new TransactionFactory, const Database>(cx)); + self->transactionFactory = Reference(new TransactionFactory, const Database>(cx, cx, false)); //Clear keyspace before running Void _ = wait(timeoutError(self->clearKeyspace(), 600)); @@ -270,25 +270,25 @@ ACTOR Future chooseTransactionFactory(Database cx, std::vectorclientPrefixInt); - self->transactionFactory = Reference(new TransactionFactory, const Database>(cx)); + self->transactionFactory = Reference(new TransactionFactory, const Database>(cx, self->extraDB, self->useExtraDB)); } else if(transactionType == READ_YOUR_WRITES) { printf("client %d: Running ReadYourWrites Transactions\n", self->clientPrefixInt); - self->transactionFactory = Reference(new TransactionFactory, const Database>(cx)); + self->transactionFactory = Reference(new TransactionFactory, const Database>(cx, self->extraDB, self->useExtraDB)); } else if(transactionType == THREAD_SAFE) { printf("client %d: Running ThreadSafe Transactions\n", self->clientPrefixInt); Reference dbHandle = wait(unsafeThreadFutureToFuture(ThreadSafeDatabase::createFromExistingDatabase(cx))); - self->transactionFactory = Reference(new TransactionFactory>(dbHandle)); + self->transactionFactory = Reference(new TransactionFactory>(dbHandle, dbHandle, false)); } else if(transactionType == MULTI_VERSION) { printf("client %d: Running Multi-Version Transactions\n", self->clientPrefixInt); Reference threadSafeHandle = wait(unsafeThreadFutureToFuture(ThreadSafeDatabase::createFromExistingDatabase(cx))); Reference dbHandle = MultiVersionDatabase::debugCreateFromExistingDatabase(threadSafeHandle); - self->transactionFactory = Reference(new TransactionFactory>(dbHandle)); + self->transactionFactory = Reference(new TransactionFactory>(dbHandle, dbHandle, false)); } return Void(); diff --git a/fdbserver/workloads/ApiWorkload.h b/fdbserver/workloads/ApiWorkload.h index a1ce135f13..c36eb0bbdd 100644 --- a/fdbserver/workloads/ApiWorkload.h +++ b/fdbserver/workloads/ApiWorkload.h @@ -85,9 +85,16 @@ struct TransactionWrapper : public ReferenceCounted { //A wrapper class for flow based transactions (NativeAPI, ReadYourWrites) template struct FlowTransactionWrapper : public TransactionWrapper { - + Database cx; + Database extraDB; + bool useExtraDB; T transaction; - FlowTransactionWrapper(Database cx) : transaction(cx) { } + T lastTransaction; + FlowTransactionWrapper(Database cx, Database extraDB, bool useExtraDB) : cx(cx), extraDB(extraDB), useExtraDB(useExtraDB), transaction(cx) { + if(useExtraDB && g_random->random01() < 0.5) { + transaction = T(extraDB); + } + } virtual ~FlowTransactionWrapper() { } //Sets a key-value pair in the database @@ -132,7 +139,12 @@ struct FlowTransactionWrapper : public TransactionWrapper { //Processes transaction error conditions Future onError(Error const& e) { - return transaction.onError(e); + Future returnVal = transaction.onError(e); + if( useExtraDB ) { + lastTransaction = std::move(transaction); + transaction = T( g_random->random01() < 0.5 ? extraDB : cx ); + } + return returnVal; } //Gets the read version of a transaction @@ -160,7 +172,7 @@ struct ThreadTransactionWrapper : public TransactionWrapper { Reference transaction; - ThreadTransactionWrapper(Reference db) : transaction(db->createTransaction()) { } + ThreadTransactionWrapper(Reference db, Reference extraDB, bool useExtraDB) : transaction(db->createTransaction()) { } virtual ~ThreadTransactionWrapper() { } //Sets a key-value pair in the database @@ -237,17 +249,21 @@ struct TransactionFactory : public TransactionFactoryInterface { //The database used to create transaction (of type Database, Reference, etc.) DB dbHandle; + DB extraDbHandle; + bool useExtraDB; - TransactionFactory(DB dbHandle) : dbHandle(dbHandle) { } + TransactionFactory(DB dbHandle, DB extraDbHandle, bool useExtraDB) : dbHandle(dbHandle), extraDbHandle(extraDbHandle), useExtraDB(useExtraDB) { } virtual ~TransactionFactory() { } //Creates a new transaction Reference createTransaction() { - return Reference(new T(dbHandle)); + return Reference(new T(dbHandle, extraDbHandle, useExtraDB)); } }; struct ApiWorkload : TestWorkload { + bool useExtraDB; + Database extraDB; ApiWorkload(WorkloadContext const& wcx, int maxClients = -1) : TestWorkload(wcx), success(true), transactionFactory(NULL), maxClients(maxClients) { clientPrefixInt = getOption(options, LiteralStringRef("clientId"), clientId); @@ -262,6 +278,13 @@ struct ApiWorkload : TestWorkload { maxLongKeyLength = getOption(options, LiteralStringRef("maxLongKeyLength"), 128); minValueLength = getOption(options, LiteralStringRef("minValueLength"), 1); maxValueLength = getOption(options, LiteralStringRef("maxValueLength"), 10000); + + useExtraDB = g_simulator.extraDB != NULL; + if(useExtraDB) { + Reference extraFile(new ClusterConnectionFile(*g_simulator.extraDB)); + Reference extraCluster = Cluster::createCluster(extraFile, -1); + extraDB = extraCluster->createDatabase(LiteralStringRef("DB")).get(); + } } Future setup(Database const& cx); diff --git a/fdbserver/workloads/ChangeConfig.actor.cpp b/fdbserver/workloads/ChangeConfig.actor.cpp index dbcbb2830f..91471d6da4 100644 --- a/fdbserver/workloads/ChangeConfig.actor.cpp +++ b/fdbserver/workloads/ChangeConfig.actor.cpp @@ -54,15 +54,13 @@ struct ChangeConfigWorkload : TestWorkload { virtual void getMetrics( vector& m ) {} - ACTOR Future ChangeConfigClient( Database cx, ChangeConfigWorkload *self) { - state Future disabler = disableConnectionFailuresAfter(300, "ChangeConfig"); - Void _ = wait( delay( self->minDelayBeforeChange + g_random->random01() * ( self->maxDelayBeforeChange - self->minDelayBeforeChange ) ) ); - + ACTOR Future extraDatabaseConfigure(ChangeConfigWorkload *self) { if (g_network->isSimulated() && g_simulator.extraDB) { Reference extraFile(new ClusterConnectionFile(*g_simulator.extraDB)); Reference cluster = Cluster::createCluster(extraFile, -1); state Database extraDB = cluster->createDatabase(LiteralStringRef("DB")).get(); + Void _ = wait(delay(5*g_random->random01())); if (self->configMode.size()) ConfigurationResult::Type _ = wait(changeConfig(extraDB, self->configMode)); if (self->networkAddresses.size()) { @@ -71,6 +69,19 @@ struct ChangeConfigWorkload : TestWorkload { else CoordinatorsResult::Type _ = wait(changeQuorum(extraDB, specifiedQuorumChange(NetworkAddress::parseList(self->networkAddresses)))); } + Void _ = wait(delay(5*g_random->random01())); + } + return Void(); + } + + ACTOR Future ChangeConfigClient( Database cx, ChangeConfigWorkload *self) { + state Future disabler = disableConnectionFailuresAfter(300, "ChangeConfig"); + Void _ = wait( delay( self->minDelayBeforeChange + g_random->random01() * ( self->maxDelayBeforeChange - self->minDelayBeforeChange ) ) ); + + state bool extraConfigureBefore = g_random->random01() < 0.5; + + if(extraConfigureBefore) { + Void _ = wait( self->extraDatabaseConfigure(self) ); } if( self->configMode.size() ) @@ -81,6 +92,11 @@ struct ChangeConfigWorkload : TestWorkload { else CoordinatorsResult::Type _ = wait( changeQuorum( cx, specifiedQuorumChange(NetworkAddress::parseList( self->networkAddresses )) ) ); } + + if(!extraConfigureBefore) { + Void _ = wait( self->extraDatabaseConfigure(self) ); + } + return Void(); } }; diff --git a/fdbserver/workloads/WriteDuringRead.actor.cpp b/fdbserver/workloads/WriteDuringRead.actor.cpp index 49e1beb8b1..6b2a5ec507 100644 --- a/fdbserver/workloads/WriteDuringRead.actor.cpp +++ b/fdbserver/workloads/WriteDuringRead.actor.cpp @@ -808,6 +808,11 @@ struct WriteDuringReadWorkload : TestWorkload { self->changeCount.insert( allKeys, 0 ); doingCommit = false; //TraceEvent("WDRError").error(e, true); + if(e.code() == error_code_database_locked) { + self->memoryDatabase = self->lastCommittedDatabase; + self->addedConflicts.insert(allKeys, false); + return Void(); + } if( e.code() == error_code_not_committed || e.code() == error_code_commit_unknown_result || e.code() == error_code_transaction_too_large || e.code() == error_code_key_too_large || e.code() == error_code_value_too_large || cancelled ) throw not_committed(); try { diff --git a/tests/slow/ApiCorrectnessSwitchover.txt b/tests/slow/ApiCorrectnessSwitchover.txt new file mode 100644 index 0000000000..6a9c452b8d --- /dev/null +++ b/tests/slow/ApiCorrectnessSwitchover.txt @@ -0,0 +1,30 @@ +testTitle=ApiCorrectnessTest +testName=ApiCorrectness +runSetup=true +clearAfterTest=true +numKeys=5000 +onlyLowerCase=true +shortKeysRatio=0.5 +minShortKeyLength=1 +maxShortKeyLength=3 +minLongKeyLength=1 +maxLongKeyLength=128 +minValueLength=1 +maxValueLength=1000 +numGets=1000 +numGetRanges=100 +numGetRangeSelectors=100 +numGetKeys=100 +numClears=100 +numClearRanges=10 +maxTransactionBytes=500000 +randomTestDuration=60 +timeout=2100 + +testName=AtomicSwitchover +switch1After=10.0 +switch2After=20.0 +stopAfter=130.0 +clearAfterTest=false +simBackupAgents=BackupToDB +extraDB=2 \ No newline at end of file