diff --git a/fdbserver/ApplyMetadataMutation.h b/fdbserver/ApplyMetadataMutation.h index 4107b95357..1fc6a29154 100644 --- a/fdbserver/ApplyMetadataMutation.h +++ b/fdbserver/ApplyMetadataMutation.h @@ -130,16 +130,12 @@ static void applyMetadataMutations(UID const& dbgid, Arena &arena, VectorRefgetLogSystemConfig().tLogs ) { - if(!tl.present() || addr.excludes(tl.interf().commit.getEndpoint().address)) { - TraceEvent("MutationRequiresRestart", dbgid).detail("M", m.toString()).detail("PrevValue", t.present() ? printable(t.get()) : "(none)").detail("toCommit", toCommit!=NULL).detail("addr", addr.toString()); - if(confChange) *confChange = true; - } - } - for( auto tl : logSystem->getLogSystemConfig().remoteTLogs ) { - if(!tl.present() || addr.excludes(tl.interf().commit.getEndpoint().address)) { - TraceEvent("MutationRequiresRestart", dbgid).detail("M", m.toString()).detail("PrevValue", t.present() ? printable(t.get()) : "(none)").detail("toCommit", toCommit!=NULL).detail("addr", addr.toString()); - if(confChange) *confChange = true; + for( auto& logs : logSystem->getLogSystemConfig().tLogs ) { + for( auto& tl : logs.tLogs ) { + if(!tl.present() || addr.excludes(tl.interf().commit.getEndpoint().address)) { + TraceEvent("MutationRequiresRestart", dbgid).detail("M", m.toString()).detail("PrevValue", t.present() ? printable(t.get()) : "(none)").detail("toCommit", toCommit!=NULL).detail("addr", addr.toString()); + if(confChange) *confChange = true; + } } } } else if(m.param1 != excludedServersVersionKey) { diff --git a/fdbserver/ClusterController.actor.cpp b/fdbserver/ClusterController.actor.cpp index f700162e31..f3040b56ea 100644 --- a/fdbserver/ClusterController.actor.cpp +++ b/fdbserver/ClusterController.actor.cpp @@ -636,14 +636,15 @@ std::vector> getWorkersForTlogsAcrossDa if(oldMasterFit < newMasterFit) return false; + //FIXME: implement for remote logs and log routers std::vector tlogProcessClasses; - for(auto& it : dbi.logSystemConfig.tLogs ) { + for(auto& it : dbi.logSystemConfig.tLogs[0].tLogs ) { auto tlogWorker = id_worker.find(it.interf().locality.processId()); if ( tlogWorker == id_worker.end() ) return false; tlogProcessClasses.push_back(tlogWorker->second.processClass); } - AcrossDatacenterFitness oldAcrossFit(dbi.logSystemConfig.tLogs, tlogProcessClasses); + AcrossDatacenterFitness oldAcrossFit(dbi.logSystemConfig.tLogs[0].tLogs, tlogProcessClasses); AcrossDatacenterFitness newAcrossFit(getWorkersForTlogsAcrossDatacenters(db.config, id_used, true)); if(oldAcrossFit < newAcrossFit) return false; diff --git a/fdbserver/DBCoreState.h b/fdbserver/DBCoreState.h index eae61a34c4..12e236f5ea 100644 --- a/fdbserver/DBCoreState.h +++ b/fdbserver/DBCoreState.h @@ -105,66 +105,24 @@ struct DBCoreState { template void serialize(Archive& ar) { - ASSERT( ar.protocolVersion() >= 0x0FDB00A320050001LL ); + ASSERT( ar.protocolVersion() >= 0x0FDB00A460010001LL); if( ar.protocolVersion() >= 0x0FDB00A560010001LL) { ar & tLogs & oldTLogData & recoveryCount & logSystemType; } else if(ar.isDeserializing) { tLogs.push_back(CoreTLogSet()); ar & tLogs[0].tLogs & tLogs[0].tLogWriteAntiQuorum & recoveryCount & tLogs[0].tLogReplicationFactor & logSystemType; - if( ar.protocolVersion() >= 0x0FDB00A460010001LL) { - uint64_t tLocalitySize = (uint64_t)tLogLocalities.size(); - ar & oldTLogData & tLogs[0].tLogPolicy & tLocalitySize; - if (ar.isDeserializing) { - tLogs[0].tLogLocalities.reserve(tLocalitySize); - for (size_t i = 0; i < tLocalitySize; i++) { - LocalityData locality; - ar & locality; - tLogs[0].tLogLocalities.push_back(locality); - } - } - } - else { - oldTLogData.clear(); - oldTLogData.push_back(OldTLogCoreData()); - oldTLogData.tLogs.push_back(CoreTLogSet()); - ar & oldTLogData[0].tLogs[0].tLogs & oldTLogData[0].tLogs[0].epochEnd & oldTLogData[0].tLogs[0].tLogWriteAntiQuorum & oldTLogData[0].tLogs[0].tLogReplicationFactor; - tLogs[0].tLogPolicy = IRepPolicyRef(new PolicyAcross(tLogs[0].tLogReplicationFactor, "zoneid", IRepPolicyRef(new PolicyOne()))); - if(!oldTLogData[0].tLogs[0].tLogs.size()) { - oldTLogData.pop_back(); - } - else { - for(int i = 0; i < oldTLogData.tLogs[0].size(); i++ ) { - oldTLogData[i].tLogs[0].tLogPolicy = IRepPolicyRef(new PolicyAcross(oldTLogData[i].tLogs[0].tLogReplicationFactor, "zoneid", IRepPolicyRef(new PolicyOne()))); - if (oldTLogData[i].tLogs[0].tLogs.size()) - { - oldTLogData[i].tLogs[0].tLogLocalities.reserve(oldTLogData[i].tLogs[0].tLogs.size()); - for (auto& tLog : oldTLogData[i].tLogs[0].tLogs) { - LocalityData locality; - locality.set(LocalityData::keyZoneId, g_random->randomUniqueID().toString()); - locality.set(LocalityData::keyDataHallId, LiteralStringRef("0")); - oldTLogData[i].tLogs[0].tLogLocalities.push_back(locality); - } - } - } - } - tLogs[0].tLogLocalities.reserve(tLogs[0].tLogs.size()); - for (auto& tLog : tLogs[0].tLogs) { + uint64_t tLocalitySize = (uint64_t)tLogs[0].tLogLocalities.size(); + ar & oldTLogData & tLogs[0].tLogPolicy & tLocalitySize; + if (ar.isDeserializing) { + tLogs[0].tLogLocalities.reserve(tLocalitySize); + for (size_t i = 0; i < tLocalitySize; i++) { LocalityData locality; - locality.set(LocalityData::keyZoneId, g_random->randomUniqueID().toString()); - locality.set(LocalityData::keyDataHallId, LiteralStringRef("0")); + ar & locality; tLogs[0].tLogLocalities.push_back(locality); } } } - - TraceEvent("CoreStateSerialize").detail("AntiQuorum", tLogWriteAntiQuorum) - .detail("logSystemType", logSystemType).detail("recoveryCount", recoveryCount) - .detail("tLogReplicationFactor", tLogReplicationFactor) - .detail("tLogPolicy", (tLogPolicy.getPtr()) ? tLogPolicy->info() : "[unset]") - .detail("logs", describe(tLogs)).detail("procotol", ar.protocolVersion()) - .detail("oldTLogData", oldTLogData.size()) - .detail("deserializing", ar.isDeserializing); } }; diff --git a/fdbserver/LogSystem.h b/fdbserver/LogSystem.h index a7e6e96109..ee8722a462 100644 --- a/fdbserver/LogSystem.h +++ b/fdbserver/LogSystem.h @@ -32,6 +32,109 @@ struct DBCoreState; +template +void uniquify( Collection& c ) { + std::sort(c.begin(), c.end()); + c.resize( std::unique(c.begin(), c.end()) - c.begin() ); +} + +class LogSet { +public: + std::vector>>> logServers; + std::vector>>> logRouters; + int32_t tLogWriteAntiQuorum; + int32_t tLogReplicationFactor; + std::vector< LocalityData > tLogLocalities; // Stores the localities of the log servers + IRepPolicyRef tLogPolicy; + LocalitySetRef logServerSet; + std::vector logIndexArray; + std::map logEntryMap; + bool isLocal; + bool hasBest; + + LogSet() : tLogWriteAntiQuorum(0), tLogReplicationFactor(0), isLocal(true), hasBest(true) {} + + int bestLocationFor( Tag tag ) { + return hasBest ? tag % logServers.size() : invalidTag; + } + + void updateLocalitySet() { + LocalityMap* logServerMap; + logServerSet = LocalitySetRef(new LocalityMap()); + logServerMap = (LocalityMap*) logServerSet.getPtr(); + + logEntryMap.clear(); + logIndexArray.clear(); + logIndexArray.reserve(logServers.size()); + + for( int i = 0; i < logServers.size(); i++ ) { + if (logServers[i]->get().present()) { + logIndexArray.push_back(i); + ASSERT(logEntryMap.find(i) == logEntryMap.end()); + logEntryMap[logIndexArray.back()] = logServerMap->add(logServers[i]->get().interf().locality, &logIndexArray.back()); + } + } + } + + void updateLocalitySet( vector const& workers ) { + LocalityMap* logServerMap; + + logServerSet = LocalitySetRef(new LocalityMap()); + logServerMap = (LocalityMap*) logServerSet.getPtr(); + + logEntryMap.clear(); + logIndexArray.clear(); + logIndexArray.reserve(workers.size()); + + for( int i = 0; i < workers.size(); i++ ) { + ASSERT(logEntryMap.find(i) == logEntryMap.end()); + logIndexArray.push_back(i); + logEntryMap[logIndexArray.back()] = logServerMap->add(workers[i].locality, &logIndexArray.back()); + } + } + + void getPushLocations( std::vector const& tags, std::vector& locations, int locationOffset ) { + newLocations.clear(); + alsoServers.clear(); + resultEntries.clear(); + + if(hasBest) { + for(auto& t : tags) { + newLocations.push_back(bestLocationFor(t)); + } + } + + uniquify( newLocations ); + + if (newLocations.size()) + alsoServers.reserve(newLocations.size()); + + // Convert locations to the also servers + for (auto location : newLocations) { + ASSERT(logEntryMap[location]._id == location); + locations.push_back(locationOffset + location); + alsoServers.push_back(logEntryMap[location]); + } + + // Run the policy, assert if unable to satify + bool result = logServerSet->selectReplicas(tLogPolicy, alsoServers, resultEntries); + ASSERT(result); + + // Add the new servers to the location array + LocalityMap* logServerMap = (LocalityMap*) logServerSet.getPtr(); + for (auto entry : resultEntries) { + locations.push_back(locationOffset + *logServerMap->getObject(entry)); + } + //TraceEvent("getPushLocations").detail("Policy", tLogPolicy->info()) + // .detail("Results", locations.size()).detail("Selection", logServerSet->size()) + // .detail("Included", alsoServers.size()).detail("Duration", timer() - t); + } + +private: + std::vector alsoServers, resultEntries; + std::vector newLocations; +}; + struct ILogSystem { // Represents a particular (possibly provisional) epoch of the log subsystem @@ -165,9 +268,8 @@ struct ILogSystem { }; struct MergedPeekCursor : IPeekCursor, ReferenceCounted { - LocalityGroup localityGroup; - std::vector< std::pair > sortedVersions; vector< Reference > serverCursors; + std::vector< std::pair > sortedVersions; Tag tag; int bestServer, currentCursor, readQuorum; Optional nextVersion; @@ -176,11 +278,10 @@ struct ILogSystem { UID randomID; int tLogReplicationFactor; IRepPolicyRef tLogPolicy; - std::vector< LocalityData > tLogLocalities; - MergedPeekCursor( std::vector>>> const& logServers, int bestServer, int readQuorum, Tag tag, Version begin, Version end, bool parallelGetMore, std::vector< LocalityData > const& tLogLocalities, IRepPolicyRef const tLogPolicy, int tLogReplicationFactor ); + MergedPeekCursor( std::vector>>> const& logServers, int bestServer, int readQuorum, Tag tag, Version begin, Version end, bool parallelGetMore ); - MergedPeekCursor( vector< Reference > const& serverCursors, LogMessageVersion const& messageVersion, int bestServer, int readQuorum, Optional nextVersion, std::vector< LocalityData > const& tLogLocalities, IRepPolicyRef const tLogPolicy, int tLogReplicationFactor ); + MergedPeekCursor( vector< Reference > const& serverCursors, LogMessageVersion const& messageVersion, int bestServer, int readQuorum, Optional nextVersion ); // if server_cursors[c]->hasMessage(), then nextSequence <= server_cursors[c]->sequence() and there are no messages known to that server with sequences in [nextSequence,server_cursors[c]->sequence()) @@ -194,7 +295,7 @@ struct ILogSystem { void calcHasMessage(); - void updateMessage(bool usePolicy); + void updateMessage(); virtual bool hasMessage(); @@ -225,6 +326,62 @@ struct ILogSystem { } }; + struct SetPeekCursor : IPeekCursor, ReferenceCounted { + std::vector logSets; + std::vector< std::vector< Reference > > serverCursors; + Tag tag; + int bestSet, bestServer, currentSet, currentCursor; + LocalityGroup localityGroup; + std::vector< std::pair > sortedVersions; + Optional nextVersion; + LogMessageVersion messageVersion; + bool hasNextMessage; + bool useBestSet; + UID randomID; + + SetPeekCursor( std::vector const& logSets, int bestSet, int bestServer, Tag tag, Version begin, Version end, bool parallelGetMore ); + + virtual Reference cloneNoMore(); + + virtual void setProtocolVersion( uint64_t version ); + + virtual Arena& arena(); + + virtual ArenaReader* reader(); + + void calcHasMessage(); + + void updateMessage(int logIdx, bool usePolicy); + + virtual bool hasMessage(); + + virtual void nextMessage(); + + virtual StringRef getMessage(); + + virtual std::vector getTags(); + + virtual void advanceTo(LogMessageVersion n); + + virtual Future getMore(); + + virtual Future onFailed(); + + virtual bool isActive(); + + virtual LogMessageVersion version(); + + virtual Version popped(); + + virtual void addref() { + ReferenceCounted::addref(); + } + + virtual void delref() { + ReferenceCounted::delref(); + } + }; + struct MultiCursor : IPeekCursor, ReferenceCounted { std::vector> cursors; std::vector epochEnds; diff --git a/fdbserver/LogSystemConfig.h b/fdbserver/LogSystemConfig.h index dccfc0ca01..46f8ea6fb4 100644 --- a/fdbserver/LogSystemConfig.h +++ b/fdbserver/LogSystemConfig.h @@ -57,7 +57,7 @@ protected: struct TLogSet { std::vector> tLogs; - std::vector logRouters; + std::vector> logRouters; int32_t tLogWriteAntiQuorum, tLogReplicationFactor; std::vector< LocalityData > tLogLocalities; // Stores the localities of the log servers IRepPolicyRef tLogPolicy; diff --git a/fdbserver/LogSystemPeekCursor.actor.cpp b/fdbserver/LogSystemPeekCursor.actor.cpp index 5f56fe141e..d97143e378 100644 --- a/fdbserver/LogSystemPeekCursor.actor.cpp +++ b/fdbserver/LogSystemPeekCursor.actor.cpp @@ -237,19 +237,17 @@ LogMessageVersion ILogSystem::ServerPeekCursor::version() { return messageVersio Version ILogSystem::ServerPeekCursor::popped() { return poppedVersion; } -ILogSystem::MergedPeekCursor::MergedPeekCursor( std::vector>>> const& logServers, int bestServer, int readQuorum, Tag tag, Version begin, Version end, bool parallelGetMore, std::vector< LocalityData > const& tLogLocalities, IRepPolicyRef const tLogPolicy, int tLogReplicationFactor ) - : bestServer(bestServer), readQuorum(readQuorum), tag(tag), currentCursor(0), hasNextMessage(false), messageVersion(begin), randomID(g_random->randomUniqueID()), tLogLocalities(tLogLocalities), tLogPolicy(tLogPolicy), tLogReplicationFactor(tLogReplicationFactor) { +ILogSystem::MergedPeekCursor::MergedPeekCursor( std::vector>>> const& logServers, int bestServer, int readQuorum, Tag tag, Version begin, Version end, bool parallelGetMore ) + : bestServer(bestServer), readQuorum(readQuorum), tag(tag), currentCursor(0), hasNextMessage(false), messageVersion(begin), randomID(g_random->randomUniqueID()) { for( int i = 0; i < logServers.size(); i++ ) { Reference cursor( new ILogSystem::ServerPeekCursor( logServers[i], tag, begin, end, bestServer >= 0, parallelGetMore ) ); //TraceEvent("MPC_starting", randomID).detail("cursor", cursor->randomID).detail("end", end); serverCursors.push_back( cursor ); } - sortedVersions.resize(serverCursors.size()); } -ILogSystem::MergedPeekCursor::MergedPeekCursor( vector< Reference > const& serverCursors, LogMessageVersion const& messageVersion, int bestServer, int readQuorum, Optional nextVersion, std::vector< LocalityData > const& tLogLocalities, IRepPolicyRef const tLogPolicy, int tLogReplicationFactor ) - : serverCursors(serverCursors), bestServer(bestServer), readQuorum(readQuorum), currentCursor(0), hasNextMessage(false), messageVersion(messageVersion), nextVersion(nextVersion), randomID(g_random->randomUniqueID()), tLogLocalities(tLogLocalities), tLogPolicy(tLogPolicy), tLogReplicationFactor(tLogReplicationFactor) { - sortedVersions.resize(serverCursors.size()); +ILogSystem::MergedPeekCursor::MergedPeekCursor( vector< Reference > const& serverCursors, LogMessageVersion const& messageVersion, int bestServer, int readQuorum, Optional nextVersion ) + : serverCursors(serverCursors), bestServer(bestServer), readQuorum(readQuorum), currentCursor(0), hasNextMessage(false), messageVersion(messageVersion), nextVersion(nextVersion), randomID(g_random->randomUniqueID()) { calcHasMessage(); } @@ -258,7 +256,7 @@ Reference ILogSystem::MergedPeekCursor::cloneNoMore() { for( auto it : serverCursors ) { cursors.push_back(it->cloneNoMore()); } - return Reference( new ILogSystem::MergedPeekCursor( cursors, messageVersion, bestServer, readQuorum, nextVersion, tLogLocalities, tLogPolicy, tLogReplicationFactor ) ); + return Reference( new ILogSystem::MergedPeekCursor( cursors, messageVersion, bestServer, readQuorum, nextVersion ) ); } void ILogSystem::MergedPeekCursor::setProtocolVersion( uint64_t version ) { @@ -292,14 +290,10 @@ void ILogSystem::MergedPeekCursor::calcHasMessage() { } hasNextMessage = false; - updateMessage(false); // Use Quorum logic - - if(!hasNextMessage) { - updateMessage(true); - } + updateMessage(); } -void ILogSystem::MergedPeekCursor::updateMessage(bool usePolicy) { +void ILogSystem::MergedPeekCursor::updateMessage() { loop { bool advancedPast = false; sortedVersions.clear(); @@ -309,24 +303,8 @@ void ILogSystem::MergedPeekCursor::updateMessage(bool usePolicy) { sortedVersions.push_back(std::pair(serverCursor->version(), i)); } - if(usePolicy) { - ASSERT(tLogPolicy); - localityGroup.clear(); - std::sort(sortedVersions.begin(), sortedVersions.end()); - - for(auto sortedVersion : sortedVersions) { - auto& locality = tLogLocalities[sortedVersion.second]; - localityGroup.add(locality); - - if( localityGroup.size() >= tLogReplicationFactor && localityGroup.validate(tLogPolicy) ) { - messageVersion = sortedVersion.first; - break; - } - } - } else { - std::nth_element(sortedVersions.begin(), sortedVersions.end()-readQuorum, sortedVersions.end()); - messageVersion = sortedVersions[sortedVersions.size()-readQuorum].first; - } + std::nth_element(sortedVersions.begin(), sortedVersions.end()-readQuorum, sortedVersions.end()); + messageVersion = sortedVersions[sortedVersions.size()-readQuorum].first; for(int i = 0; i < serverCursors.size(); i++) { auto& c = serverCursors[i]; @@ -433,6 +411,252 @@ Version ILogSystem::MergedPeekCursor::popped() { return poppedVersion; } +ILogSystem::SetPeekCursor::SetPeekCursor( std::vector const& logSets, int bestSet, int bestServer, Tag tag, Version begin, Version end, bool parallelGetMore ) + : logSets(logSets), bestSet(bestSet), bestServer(bestServer), tag(tag), currentCursor(0), currentSet(bestSet), hasNextMessage(false), messageVersion(begin), useBestSet(true), randomID(g_random->randomUniqueID()) { + serverCursors.resize(logSets.size()); + int maxServers = 0; + for( int i = 0; i < logSets.size(); i++ ) { + for( int j = 0; j < logSets[i].logServers.size(); j++) { + Reference cursor( new ILogSystem::ServerPeekCursor( logSets[i].logServers[j], tag, begin, end, true, parallelGetMore ) ); + serverCursors[i].push_back( cursor ); + } + maxServers = std::max(maxServers, serverCursors[i].size()); + } + sortedVersions.resize(maxServers); +} + +Reference ILogSystem::SetPeekCursor::cloneNoMore() { + ASSERT(false); //not implemented + throw internal_error(); +} + +void ILogSystem::SetPeekCursor::setProtocolVersion( uint64_t version ) { + for( auto& cursors : serverCursors ) { + for( auto& it : cursors ) { + if( it->hasMessage() ) { + it->setProtocolVersion( version ); + } + } + } +} + +Arena& ILogSystem::SetPeekCursor::arena() { return serverCursors[currentSet][currentCursor]->arena(); } + +ArenaReader* ILogSystem::SetPeekCursor::reader() { return serverCursors[currentSet][currentCursor]->reader(); } + + +void ILogSystem::SetPeekCursor::calcHasMessage() { + if(nextVersion.present()) serverCursors[bestSet][bestServer]->advanceTo( nextVersion.get() ); + if( serverCursors[bestSet][bestServer]->hasMessage() ) { + messageVersion = serverCursors[bestSet][bestServer]->version(); + currentSet = bestSet; + currentCursor = bestServer; + hasNextMessage = true; + + for (auto& cursors : serverCursors) { + for(auto& c : cursors) { + c->advanceTo(messageVersion); + } + } + + return; + } + + auto bestVersion = serverCursors[bestSet][bestServer]->version(); + for (auto& cursors : serverCursors) { + for (auto& c : cursors) { + c->advanceTo(bestVersion); + } + } + + if(useBestSet) { + hasNextMessage = false; + updateMessage(bestSet, false); // Use Quorum logic + + if(!hasNextMessage) { + updateMessage(bestSet, true); + } + } else { + for(int i = 0; i < logSets.size() && !hasNextMessage; i++) { + if(i != bestSet) { + updateMessage(i, false); // Use Quorum logic + } + } + + for(int i = 0; i < logSets.size() && !hasNextMessage; i++) { + if(i != bestSet) { + updateMessage(i, true); + } + } + } +} + +void ILogSystem::SetPeekCursor::updateMessage(int logIdx, bool usePolicy) { + loop { + bool advancedPast = false; + sortedVersions.clear(); + for(int i = 0; i < serverCursors[logIdx].size(); i++) { + auto& serverCursor = serverCursors[i]; + if (nextVersion.present()) serverCursor[logIdx]->advanceTo(nextVersion.get()); + sortedVersions.push_back(std::pair(serverCursor[logIdx]->version(), i)); + } + + if(usePolicy) { + localityGroup.clear(); + std::sort(sortedVersions.begin(), sortedVersions.end()); + + for(auto sortedVersion : sortedVersions) { + auto& locality = logSets[logIdx].tLogLocalities[sortedVersion.second]; + localityGroup.add(locality); + + if( localityGroup.size() >= logSets[logIdx].tLogReplicationFactor && localityGroup.validate(logSets[logIdx].tLogPolicy) ) { + messageVersion = sortedVersion.first; + break; + } + } + } else { + //(int)oldLogData[i].logServers.size() + 1 - oldLogData[i].tLogReplicationFactor + std::nth_element(sortedVersions.begin(), sortedVersions.end()-(logSets[logIdx].logServers.size()+1-logSets[logIdx].tLogReplicationFactor), sortedVersions.end()); + messageVersion = sortedVersions[sortedVersions.size()-(logSets[logIdx].logServers.size()+1-logSets[logIdx].tLogReplicationFactor)].first; + } + + for(int i = 0; i < serverCursors[logIdx].size(); i++) { + auto& c = serverCursors[logIdx][i]; + auto start = c->version(); + c->advanceTo(messageVersion); + if( start < messageVersion && messageVersion < c->version() ) { + advancedPast = true; + TEST(true); //Merge peek cursor advanced past desired sequence + } + } + + if(!advancedPast) + break; + } + + for(int i = 0; i < serverCursors[logIdx].size(); i++) { + auto& c = serverCursors[logIdx][i]; + ASSERT_WE_THINK( !c->hasMessage() || c->version() >= messageVersion ); // Seems like the loop above makes this unconditionally true + if (c->version() == messageVersion && c->hasMessage()) { + hasNextMessage = true; + currentSet = logIdx; + currentCursor = i; + break; + } + } +} + +bool ILogSystem::SetPeekCursor::hasMessage() { + return hasNextMessage; +} + +void ILogSystem::SetPeekCursor::nextMessage() { + nextVersion = version(); + nextVersion.get().sub++; + serverCursors[currentSet][currentCursor]->nextMessage(); + calcHasMessage(); + ASSERT(hasMessage() || !version().sub); +} + +StringRef ILogSystem::SetPeekCursor::getMessage() { return serverCursors[currentSet][currentCursor]->getMessage(); } + +std::vector ILogSystem::SetPeekCursor::getTags() { + return serverCursors[currentSet][currentCursor]->getTags(); +} + +void ILogSystem::SetPeekCursor::advanceTo(LogMessageVersion n) { + for( auto& cursors : serverCursors ) { + for (auto& c : cursors) { + c->advanceTo(n); + } + } + calcHasMessage(); +} + +ACTOR Future setPeekGetMore(ILogSystem::SetPeekCursor* self, LogMessageVersion startVersion) { + loop { + //TraceEvent("LPC_getMoreA", self->randomID).detail("start", startVersion.toString()); + if(self->bestServer >= 0 && self->bestSet >= 0 && self->serverCursors[self->bestSet][self->bestServer]->isActive()) { + ASSERT(!self->serverCursors[self->bestSet][self->bestServer]->hasMessage()); + Void _ = wait( self->serverCursors[self->bestSet][self->bestServer]->getMore() || self->serverCursors[self->bestSet][self->bestServer]->onFailed() ); + self->useBestSet = true; + } else { + bool bestSetValid = self->bestSet >= 0; + if(bestSetValid) { + self->localityGroup.clear(); + for( int i = 0; i < self->serverCursors[self->bestSet].size(); i++) { + if(!self->serverCursors[self->bestSet][i]->isActive()) { + self->localityGroup.add(self->logSets[self->bestSet].tLogLocalities[i]); + } + } + bestSetValid = self->localityGroup.size() < self->logSets[self->bestSet].tLogReplicationFactor || !self->localityGroup.validate(self->logSets[self->bestSet].tLogPolicy); + } + if(bestSetValid) { + vector> q; + for (auto& c : self->serverCursors[self->bestSet]) { + if (!c->hasMessage()) { + q.push_back(c->getMore()); + } + } + Void _ = wait(quorum(q, 1)); + self->useBestSet = true; + } else { + //FIXME: this will peeking way too many cursors when satellites exist, and does not need to peek bestSet cursors since we cannot get anymore data from them + vector> q; + for(auto& cursors : self->serverCursors) { + for (auto& c :cursors) { + if (!c->hasMessage()) { + q.push_back(c->getMore()); + } + } + } + Void _ = wait(quorum(q, 1)); + self->useBestSet = false; + } + } + self->calcHasMessage(); + //TraceEvent("LPC_getMoreB", self->randomID).detail("hasMessage", self->hasMessage()).detail("start", startVersion.toString()).detail("seq", self->version().toString()); + if (self->hasMessage() || self->version() > startVersion) + return Void(); + } +} + +Future ILogSystem::SetPeekCursor::getMore() { + auto startVersion = version(); + calcHasMessage(); + if( hasMessage() ) + return Void(); + if (nextVersion.present()) + advanceTo(nextVersion.get()); + ASSERT(!hasMessage()); + if (version() > startVersion) + return Void(); + + return setPeekGetMore(this, startVersion); +} + +Future ILogSystem::SetPeekCursor::onFailed() { + ASSERT(false); + return Never(); +} + +bool ILogSystem::SetPeekCursor::isActive() { + ASSERT(false); + return false; +} + +LogMessageVersion ILogSystem::SetPeekCursor::version() { return messageVersion; } + +Version ILogSystem::SetPeekCursor::popped() { + Version poppedVersion = 0; + for (auto& cursors : serverCursors) { + for(auto& c : cursors) { + poppedVersion = std::max(poppedVersion, c->popped()); + } + } + return poppedVersion; +} + ILogSystem::MultiCursor::MultiCursor( std::vector> cursors, std::vector epochEnds ) : cursors(cursors), epochEnds(epochEnds), poppedVersion(0) {} Reference ILogSystem::MultiCursor::cloneNoMore() { diff --git a/fdbserver/OldTLogServer.actor.cpp b/fdbserver/OldTLogServer.actor.cpp deleted file mode 100644 index c34296edb7..0000000000 --- a/fdbserver/OldTLogServer.actor.cpp +++ /dev/null @@ -1,1428 +0,0 @@ -/* - * OldTLogServer.actor.cpp - * - * This source file is part of the FoundationDB open source project - * - * Copyright 2013-2018 Apple Inc. and the FoundationDB project authors - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -#include "flow/actorcompiler.h" -#include "flow/Hash3.h" -#include "flow/Stats.h" -#include "flow/UnitTest.h" -#include "fdbclient/NativeAPI.h" -#include "fdbclient/KeyRangeMap.h" -#include "fdbclient/SystemData.h" -#include "WorkerInterface.h" -#include "TLogInterface.h" -#include "flow/Notified.h" -#include "Knobs.h" -#include "IKeyValueStore.h" -#include "flow/ActorCollection.h" -#include "fdbrpc/FailureMonitor.h" -#include "IDiskQueue.h" -#include "fdbrpc/sim_validation.h" -#include "ServerDBInfo.h" -#include "LogSystem.h" -#include "WaitFailure.h" - -using std::pair; -using std::make_pair; -using std::min; -using std::max; - -namespace oldTLog { - struct TLogQueueEntryRef { - Version version; - Version knownCommittedVersion; - StringRef messages; - VectorRef< TagMessagesRef > tags; - - TLogQueueEntryRef() : version(0), knownCommittedVersion(0) {} - TLogQueueEntryRef(Arena &a, TLogQueueEntryRef const &from) - : version(from.version), knownCommittedVersion(from.knownCommittedVersion), messages(a, from.messages), tags(a, from.tags) { - } - - template - void serialize(Ar& ar) { - if( ar.protocolVersion() >= 0x0FDB00A460010001) { - ar & version & messages & tags & knownCommittedVersion; - } else if(ar.isDeserializing) { - ar & version & messages & tags; - knownCommittedVersion = 0; - } - } - size_t expectedSize() const { - return messages.expectedSize() + tags.expectedSize(); - } - }; - - typedef Standalone TLogQueueEntry; - - struct TLogQueue : public IClosable { - public: - TLogQueue( IDiskQueue* queue, UID dbgid ) : queue(queue), debugNextReadVersion(1), dbgid(dbgid) {} - - // Each packet in the queue is - // uint32_t payloadSize - // uint8_t payload[payloadSize] (begins with uint64_t protocolVersion via IncludeVersion) - // uint8_t validFlag - - // TLogQueue is a durable queue of TLogQueueEntry objects with an interface similar to IDiskQueue - - // TLogQueue pushes (but not commits) are atomic - after commit fails to return, a prefix of entire calls to push are durable. This is - // implemented on top of the weaker guarantee of IDiskQueue::commit (that a prefix of bytes is durable) using validFlag and by - // padding any incomplete packet with zeros after recovery. - - // Before calling push, pop, or commit, the user must call readNext() until it throws - // end_of_stream(). It may not be called again thereafter. - Future readNext() { - return readNext( this ); - } - - void push( TLogQueueEntryRef const& qe ) { - ASSERT( version_location.empty() || version_location.lastItem()->key < qe.version ); - BinaryWriter wr( Unversioned() ); // outer framing is not versioned - wr << uint32_t(0); - IncludeVersion().write(wr); // payload is versioned - wr << qe; - wr << uint8_t(1); - *(uint32_t*)wr.getData() = wr.getLength() - sizeof(uint32_t) - sizeof(uint8_t); - auto loc = queue->push( wr.toStringRef() ); - //TraceEvent("TLogQueueVersionWritten", dbgid).detail("Size", wr.getLength() - sizeof(uint32_t) - sizeof(uint8_t)).detail("Loc", loc); - version_location[qe.version] = loc; - } - void pop( Version upTo ) { - // Keep only the given and all subsequent version numbers - // Find the first version >= upTo - auto v = version_location.lower_bound(upTo); - if (v == version_location.begin()) return; - - if(v == version_location.end()) { - v = version_location.lastItem(); - } - else { - v.decrementNonEnd(); - } - - queue->pop( v->value ); - version_location.erase( version_location.begin(), v ); // ... and then we erase that previous version and all prior versions - } - Future commit() { return queue->commit(); } - - // Implements IClosable - virtual Future getError() { return queue->getError(); } - virtual Future onClosed() { return queue->onClosed(); } - virtual void dispose() { queue->dispose(); delete this; } - virtual void close() { queue->close(); delete this; } - - private: - IDiskQueue* queue; - Map version_location; // For the version of each entry that was push()ed, the end location of the serialized bytes - Version debugNextReadVersion; - UID dbgid; - - ACTOR static Future readNext( TLogQueue* self ) { - state TLogQueueEntry result; - state int zeroFillSize = 0; - - loop { - Standalone h = wait( self->queue->readNext( sizeof(uint32_t) ) ); - if (h.size() != sizeof(uint32_t)) { - if (h.size()) { - TEST( true ); // Zero fill within size field - int payloadSize = 0; - memcpy(&payloadSize, h.begin(), h.size()); - zeroFillSize = sizeof(uint32_t)-h.size(); // zero fill the size itself - zeroFillSize += payloadSize+1; // and then the contents and valid flag - } - break; - } - - state uint32_t payloadSize = *(uint32_t*)h.begin(); - ASSERT( payloadSize < (100<<20) ); - - Standalone e = wait( self->queue->readNext( payloadSize+1 ) ); - if (e.size() != payloadSize+1) { - TEST( true ); // Zero fill within payload - zeroFillSize = payloadSize+1 - e.size(); - break; - } - - if (e[payloadSize]) { - Arena a = e.arena(); - ArenaReader ar( a, e.substr(0, payloadSize), IncludeVersion() ); - ar >> result; - ASSERT( result.version >= self->debugNextReadVersion ); - self->debugNextReadVersion = result.version + 1; - self->version_location[result.version] = self->queue->getNextReadLocation(); - return result; - } - } - if (zeroFillSize) { - TEST( true ); // Fixing a partial commit at the end of the tlog queue - for(int i=0; iqueue->push( StringRef((const uint8_t*)"",1) ); - } - throw end_of_stream(); - } - }; - - struct TLogData : NonCopyable { - struct TagData { - std::deque> version_messages; - bool nothing_persistent; // true means tag is *known* to have no messages in persistentData. false means nothing. - bool popped_recently; // `popped` has changed since last updatePersistentData - Version popped; // see popped version tracking contract below - bool update_version_sizes; - - TagData( Version popped, bool nothing_persistent, bool popped_recently, Tag tag ) : nothing_persistent(nothing_persistent), popped(popped), popped_recently(popped_recently), update_version_sizes(tag != txsTag) {} - - TagData(TagData&& r) noexcept(true) : version_messages(std::move(r.version_messages)), nothing_persistent(r.nothing_persistent), popped_recently(r.popped_recently), popped(r.popped), update_version_sizes(r.update_version_sizes) {} - void operator= (TagData&& r) noexcept(true) { - version_messages = std::move(r.version_messages); - nothing_persistent = r.nothing_persistent; - popped_recently = r.popped_recently; - popped = r.popped; - update_version_sizes = r.update_version_sizes; - } - - // Erase messages not needed to update *from* versions >= before (thus, messages with toversion <= before) - ACTOR Future eraseMessagesBefore( TagData *self, Version before, Counter* bytesErased, TLogData *tlogData, int taskID ) { - while(!self->version_messages.empty() && self->version_messages.front().first < before) { - Version version = self->version_messages.front().first; - std::pair &sizes = tlogData->version_sizes[version]; - int64_t messagesErased = 0; - - while(!self->version_messages.empty() && self->version_messages.front().first == version) { - auto const& m = self->version_messages.front(); - ++messagesErased; - - if(self->update_version_sizes) { - sizes.first -= m.second.expectedSize(); - } - - self->version_messages.pop_front(); - } - - *bytesErased += (messagesErased * sizeof(std::pair) * SERVER_KNOBS->VERSION_MESSAGES_OVERHEAD_FACTOR_1024THS) >> 10; - Void _ = wait(yield(taskID)); - } - - return Void(); - } - - Future eraseMessagesBefore(Version before, Counter* bytesErased, TLogData *tlogData, int taskID) { - return eraseMessagesBefore(this, before, bytesErased, tlogData, taskID); - } - }; - - /* - Popped version tracking contract needed by log system to implement ILogCursor::popped(): - - - Log server tracks for each (possible) tag a popped_version - Impl: TagData::popped (in memory) and persistTagPoppedKeys (in persistentData) - - popped_version(tag) is <= the maximum version for which log server (or a predecessor) is ever asked to pop the tag - Impl: Only increased by tLogPop() in response to either a pop request or recovery from a predecessor - - popped_version(tag) is > the maximum version for which log server is unable to peek messages due to previous pops (on this server or a predecessor) - Impl: Increased by tLogPop() atomically with erasing messages from memory; persisted by updatePersistentData() atomically with erasing messages from store; messages are not erased from queue where popped_version is not persisted - - LockTLogReply returns all tags which either have messages, or which have nonzero popped_versions - Impl: tag_data is present for all such tags - - peek(tag, v) returns the popped_version for tag if that is greater than v - Impl: Check tag_data->popped (after all waits) - */ - - struct peekTrackerData { - std::map> sequence_version; - double lastUpdate; - }; - - std::map peekTracker; - - UID dbgid; - bool coreStarted; - bool stopped; - DBRecoveryCount recoveryCount; - - IKeyValueStore* persistentData; - IDiskQueue* rawPersistentQueue; - TLogQueue *persistentQueue; - VersionMetricHandle persistentDataVersion, persistentDataDurableVersion; // The last version number in the portion of the log (written|durable) to persistentData - NotifiedVersion version, queueCommittedVersion, queueCommitEnd; - Version queueCommitBegin, queueCommittingVersion; - int64_t diskQueueCommitBytes; - AsyncVar largeDiskQueueCommitBytes; //becomes true when diskQueueCommitBytes is greater than MAX_QUEUE_COMMIT_BYTES - Version prevVersion, knownCommittedVersion; - - Deque>>> messageBlocks; - Map< Tag, TagData > tag_data; - - Map> version_sizes; - - int64_t instanceID; - CounterCollection cc; - Counter bytesInput; - Counter bytesDurable; - - Reference> dbInfo; - Future updatePersist; //SOMEDAY: integrate the recovery and update storage so that only one of them is committing to persistant data. - PromiseStream> addActor; - - TLogData(UID dbgid, IKeyValueStore* persistentData, IDiskQueue * persistentQueue, Reference> const& dbInfo) - : dbgid(dbgid), - persistentData(persistentData), rawPersistentQueue(persistentQueue), persistentQueue(new TLogQueue(persistentQueue, dbgid)), - prevVersion(0), knownCommittedVersion(0), - dbInfo(dbInfo), - updatePersist(Void()), - instanceID(g_random->randomUniqueID().first()), - cc("TLog", dbgid.toString()), - bytesInput("bytesInput", cc), - bytesDurable("bytesDurable", cc), - // These are initialized differently on init() or recovery - recoveryCount(), coreStarted(false), stopped(false), queueCommitBegin(0), queueCommitEnd(0), diskQueueCommitBytes(0), largeDiskQueueCommitBytes(false), queueCommittingVersion(0) - { - persistentDataVersion.init(LiteralStringRef("TLog.PersistentDataVersion"), cc.id); - persistentDataDurableVersion.init(LiteralStringRef("TLog.PersistentDataDurableVersion"), cc.id); - version.initMetric(LiteralStringRef("TLog.Version"), cc.id); - queueCommittedVersion.initMetric(LiteralStringRef("TLog.QueueCommittedVersion"), cc.id); - - specialCounter(cc, "version", [this](){ return this->version.get(); }); - specialCounter(cc, "kvstoreBytesUsed", [this](){ return this->persistentData->getStorageBytes().used; }); - specialCounter(cc, "kvstoreBytesFree", [this](){ return this->persistentData->getStorageBytes().free; }); - specialCounter(cc, "kvstoreBytesAvailable", [this](){ return this->persistentData->getStorageBytes().available; }); - specialCounter(cc, "kvstoreBytesTotal", [this](){ return this->persistentData->getStorageBytes().total; }); - specialCounter(cc, "queueDiskBytesUsed", [this](){ return this->rawPersistentQueue->getStorageBytes().used; }); - specialCounter(cc, "queueDiskBytesFree", [this](){ return this->rawPersistentQueue->getStorageBytes().free; }); - specialCounter(cc, "queueDiskBytesAvailable", [this](){ return this->rawPersistentQueue->getStorageBytes().available; }); - specialCounter(cc, "queueDiskBytesTotal", [this](){ return this->rawPersistentQueue->getStorageBytes().total; }); - } - - LogEpoch epoch() const { return recoveryCount; } - }; - - ACTOR Future tLogLock( TLogData* self, ReplyPromise< TLogLockResult > reply ) { - state Version stopVersion = self->version.get(); - - TEST(true); // TLog stopped by recovering master - TEST( self->stopped ); - TEST( !self->stopped ); - - TraceEvent("TLogStop", self->dbgid).detail("Ver", stopVersion).detail("isStopped", self->stopped); - - self->stopped = true; - - // Lock once the current version has been committed - Void _ = wait( self->queueCommittedVersion.whenAtLeast( stopVersion ) ); - - ASSERT(stopVersion == self->version.get()); - - TLogLockResult result; - result.end = stopVersion; - result.knownCommittedVersion = self->knownCommittedVersion; - for( auto & tag : self->tag_data ) - result.tags.push_back( tag.key ); - - reply.send( result ); - return Void(); - } - - KeyRange prefixRange( KeyRef prefix ) { - Key end = prefix; - UNSTOPPABLE_ASSERT( end.size() && end.end()[-1] != 0xFF ); - ++const_cast( end.end()[-1] ); - return KeyRangeRef( prefix, end ); - } - - ////// Persistence format (for self->persistentData) - - // Immutable keys - static const KeyValueRef persistFormat( LiteralStringRef( "Format" ), LiteralStringRef("FoundationDB/LogServer/2/2") ); - static const KeyRangeRef persistFormatReadableRange( LiteralStringRef("FoundationDB/LogServer/2/2"), LiteralStringRef("FoundationDB/LogServer/2/3") ); - static const KeyRef persistID = LiteralStringRef( "ID" ); - static const KeyRef persistRecoveryCountKey = LiteralStringRef("DbRecoveryCount"); - - // Updated on updatePersistentData() - static const KeyRef persistCurrentVersionKey = LiteralStringRef("version"); - static const KeyRange persistTagMessagesKeys = prefixRange(LiteralStringRef("TagMsg/")); - static const KeyRange persistTagPoppedKeys = prefixRange(LiteralStringRef("TagPop/")); - - // Only present during network recovery process - static const KeyValueRef persistRecoveryInProgress( LiteralStringRef("RecoveryInProgress"), LiteralStringRef("1") ); - - static Key persistTagMessagesKey( Tag tag, Version version ) { - BinaryWriter wr( Unversioned() ); - wr.serializeBytes(persistTagMessagesKeys.begin); - wr << tag; - wr << bigEndian64( version ); - return wr.toStringRef(); - } - - static Key persistTagPoppedKey( Tag tag ) { - BinaryWriter wr(Unversioned()); - wr.serializeBytes( persistTagPoppedKeys.begin ); - wr << tag; - return wr.toStringRef(); - } - - static Value persistTagPoppedValue( Version popped ) { - return BinaryWriter::toValue( popped, Unversioned() ); - } - - static Tag decodeTagPoppedKey( KeyRef key ) { - Tag s; - BinaryReader rd( key.removePrefix(persistTagPoppedKeys.begin), Unversioned() ); - rd >> s; - return s; - } - - static Version decodeTagPoppedValue( ValueRef value ) { - return BinaryReader::fromStringRef( value, Unversioned() ); - } - - static StringRef stripTagMessagesKey( StringRef key ) { - return key.substr( sizeof(Tag) + persistTagMessagesKeys.begin.size() ); - } - - static Version decodeTagMessagesKey( StringRef key ) { - return bigEndian64( BinaryReader::fromStringRef( stripTagMessagesKey(key), Unversioned() ) ); - } - - static Standalone decodeTagMessagesKeyTag( StringRef key ) { - key = key.removePrefix( persistTagMessagesKeys.begin ); // \x00\xff - BinaryWriter wr( Unversioned() ); - for(auto c = key.begin(); c != key.end(); ++c) { - if (*c) - wr << *c; - else { - ASSERT( c+1 != key.end() ); - if (c[1] == 0xff) { - wr << uint8_t(0); - c++; - } else if (c[1] == 0) - break; - else - throw internal_error(); - } - } - return wr.toStringRef(); - } - - void validate( TLogData* self, bool force = false ) { - } - - void updatePersistentPopped( TLogData* self, Tag tag, TLogData::TagData& data ) { - if (!data.popped_recently) return; - self->persistentData->set(KeyValueRef( persistTagPoppedKey(tag), persistTagPoppedValue(data.popped) )); - data.popped_recently = false; - - if (data.nothing_persistent) return; - - self->persistentData->clear( KeyRangeRef( - persistTagMessagesKey( tag, Version(0) ), - persistTagMessagesKey( tag, data.popped ) ) ); - if (data.popped > self->persistentDataVersion) - data.nothing_persistent = true; - //TraceEvent("TLogPopWrite", self->dbgid).detail("Tag", tag).detail("To", data.popped); - } - - ACTOR Future updatePersistentData( TLogData* self, Version newPersistentDataVersion ) { - // PERSIST: Changes self->persistentDataVersion and writes and commits the relevant changes - ASSERT( newPersistentDataVersion <= self->version.get() ); - ASSERT( newPersistentDataVersion <= self->queueCommittedVersion.get() ); - ASSERT( newPersistentDataVersion > self->persistentDataVersion ); - ASSERT( self->persistentDataVersion == self->persistentDataDurableVersion ); - - //TraceEvent("updatePersistentData", self->dbgid).detail("seq", newPersistentDataSeq); - - state bool anyData = false; - state Map::iterator tag; - // For all existing tags - for(tag = self->tag_data.begin(); tag != self->tag_data.end(); ++tag) { - state Version currentVersion = 0; - // Clear recently popped versions from persistentData if necessary - updatePersistentPopped( self, tag->key, tag->value ); - // Transfer unpopped messages with version numbers less than newPersistentDataVersion to persistentData - state std::deque>::iterator msg = tag->value.version_messages.begin(); - while(msg != tag->value.version_messages.end() && msg->first <= newPersistentDataVersion) { - currentVersion = msg->first; - anyData = true; - tag->value.nothing_persistent = false; - BinaryWriter wr( Unversioned() ); - - for(; msg != tag->value.version_messages.end() && msg->first == currentVersion; ++msg) - wr << msg->second.toStringRef(); - - self->persistentData->set( KeyValueRef( persistTagMessagesKey( tag->key, currentVersion ), wr.toStringRef() ) ); - - Future f = yield(TaskUpdateStorage); - if(!f.isReady()) { - Void _ = wait(f); - msg = std::upper_bound(tag->value.version_messages.begin(), tag->value.version_messages.end(), std::make_pair(currentVersion, LengthPrefixedStringRef()), CompareFirst>()); - } - } - - Void _ = wait(yield(TaskUpdateStorage)); - } - - self->persistentData->set( KeyValueRef( persistCurrentVersionKey, BinaryWriter::toValue(newPersistentDataVersion, Unversioned()) ) ); - self->persistentDataVersion = newPersistentDataVersion; - - Void _ = wait( self->persistentData->commit() ); // SOMEDAY: This seems to be running pretty often, should we slow it down??? - Void _ = wait( delay(0, TaskUpdateStorage) ); - - // Now that the changes we made to persistentData are durable, erase the data we moved from memory and the queue, increase bytesDurable accordingly, and update persistentDataDurableVersion. - - TEST(anyData); // TLog moved data to persistentData - self->persistentDataDurableVersion = newPersistentDataVersion; - - for(tag = self->tag_data.begin(); tag != self->tag_data.end(); ++tag) { - Void _ = wait(tag->value.eraseMessagesBefore( newPersistentDataVersion+1, &self->bytesDurable, self, TaskUpdateStorage )); - Void _ = wait(yield(TaskUpdateStorage)); - } - - self->version_sizes.erase(self->version_sizes.begin(), self->version_sizes.lower_bound(self->persistentDataDurableVersion)); - - Void _ = wait(yield(TaskUpdateStorage)); - - while(!self->messageBlocks.empty() && self->messageBlocks.front().first <= newPersistentDataVersion) { - self->bytesDurable += self->messageBlocks.front().second.size() * SERVER_KNOBS->TLOG_MESSAGE_BLOCK_OVERHEAD_FACTOR; - self->messageBlocks.pop_front(); - Void _ = wait(yield(TaskUpdateStorage)); - } - - ASSERT(self->bytesDurable.getValue() <= self->bytesInput.getValue()); - - if( self->queueCommitEnd.get() > 0 ) - self->persistentQueue->pop( newPersistentDataVersion+1 ); // SOMEDAY: this can cause a slow task (~0.5ms), presumably from erasing too many versions. Should we limit the number of versions cleared at a time? - - return Void(); - } - - // This function (and updatePersistentData, which is called by this function) run at a low priority and can soak up all CPU resources. - // For this reason, they employ aggressive use of yields to avoid causing slow tasks that could introduce latencies for more important - // work (e.g. commits). - ACTOR Future updateStorage( TLogData* self ) { - Void _ = wait(delay(0, TaskUpdateStorage)); - loop { - state Version prevVersion = 0; - state Version nextVersion = 0; - state int totalSize = 0; - - state Map>::iterator sizeItr = self->version_sizes.begin(); - while( totalSize < SERVER_KNOBS->UPDATE_STORAGE_BYTE_LIMIT && sizeItr != self->version_sizes.end() - && (self->bytesInput.getValue() - self->bytesDurable.getValue() - totalSize >= SERVER_KNOBS->TLOG_SPILL_THRESHOLD || sizeItr->value.first == 0) ) - { - Void _ = wait( yield(TaskUpdateStorage) ); - - ++sizeItr; - nextVersion = sizeItr == self->version_sizes.end() ? self->version.get() : sizeItr->key; - - state Map::iterator tag; - for(tag = self->tag_data.begin(); tag != self->tag_data.end(); ++tag) { - auto it = std::lower_bound(tag->value.version_messages.begin(), tag->value.version_messages.end(), std::make_pair(prevVersion, LengthPrefixedStringRef()), CompareFirst>()); - for(; it != tag->value.version_messages.end() && it->first < nextVersion; ++it) { - totalSize += it->second.expectedSize(); - } - - Void _ = wait(yield(TaskUpdateStorage)); - } - - prevVersion = nextVersion; - } - - nextVersion = std::max(nextVersion, self->persistentDataVersion); - - TraceEvent("UpdateStorageVer", self->dbgid).detail("nextVersion", nextVersion).detail("persistentDataVersion", self->persistentDataVersion).detail("totalSize", totalSize); - - Void _ = wait( self->queueCommittedVersion.whenAtLeast( nextVersion ) ); - Void _ = wait( delay(0, TaskUpdateStorage) ); - - if (nextVersion > self->persistentDataVersion) { - self->updatePersist = updatePersistentData(self, nextVersion); - Void _ = wait( self->updatePersist ); - } - - if( totalSize < SERVER_KNOBS->UPDATE_STORAGE_BYTE_LIMIT ) { - Void _ = wait( delay(BUGGIFY ? SERVER_KNOBS->BUGGIFY_TLOG_STORAGE_MIN_UPDATE_INTERVAL : SERVER_KNOBS->TLOG_STORAGE_MIN_UPDATE_INTERVAL, TaskUpdateStorage) ); - } - else { - //recovery wants to commit to persistant data when updatePersistentData is not active, this delay ensures that immediately after - //updatePersist returns another one has not been started yet. - Void _ = wait( delay(0.0, TaskUpdateStorage) ); - } - } - } - - void commitMessages( TLogData* self, Version version, Arena arena, StringRef messages, VectorRef< TagMessagesRef > tags) { - // SOMEDAY: This method of copying messages is reasonably memory efficient, but it's still a lot of bytes copied. Find a - // way to do the memory allocation right as we receive the messages in the network layer. - - int64_t addedBytes = 0; - int64_t expectedBytes = 0; - - if(!messages.size()) { - return; - } - - StringRef messages1; // the first block of messages, if they aren't all stored contiguously. otherwise empty - - // Grab the last block in the blocks list so we can share its arena - // We pop all of the elements of it to create a "fresh" vector that starts at the end of the previous vector - Standalone> block; - if(self->messageBlocks.empty()) { - block = Standalone>(); - block.reserve(block.arena(), std::max(SERVER_KNOBS->TLOG_MESSAGE_BLOCK_BYTES, messages.size())); - } - else { - block = self->messageBlocks.back().second; - } - - block.pop_front(block.size()); - - // If the current batch of messages doesn't fit entirely in the remainder of the last block in the list - if(messages.size() + block.size() > block.capacity()) { - // Find how many messages will fit - LengthPrefixedStringRef r((uint32_t*)messages.begin()); - uint8_t const* end = messages.begin() + block.capacity() - block.size(); - while(r.toStringRef().end() <= end) { - r = LengthPrefixedStringRef( (uint32_t*)r.toStringRef().end() ); - } - - // Fill up the rest of this block - int bytes = (uint8_t*)r.getLengthPtr()-messages.begin(); - if (bytes) { - TEST(true); // Splitting commit messages across multiple blocks - messages1 = StringRef(block.end(), bytes); - block.append(block.arena(), messages.begin(), bytes); - self->messageBlocks.push_back( std::make_pair(version, block) ); - addedBytes += int64_t(block.size()) * SERVER_KNOBS->TLOG_MESSAGE_BLOCK_OVERHEAD_FACTOR; - messages = messages.substr(bytes); - } - - // Make a new block - block = Standalone>(); - block.reserve(block.arena(), std::max(SERVER_KNOBS->TLOG_MESSAGE_BLOCK_BYTES, messages.size())); - } - - // Copy messages into block - ASSERT(messages.size() <= block.capacity() - block.size()); - block.append(block.arena(), messages.begin(), messages.size()); - self->messageBlocks.push_back( std::make_pair(version, block) ); - addedBytes += int64_t(block.size()) * SERVER_KNOBS->TLOG_MESSAGE_BLOCK_OVERHEAD_FACTOR; - messages = StringRef(block.end()-messages.size(), messages.size()); - - for(auto tag = tags.begin(); tag != tags.end(); ++tag) { - int64_t tagMessages = 0; - - auto tsm = self->tag_data.find(tag->tag); - if (tsm == self->tag_data.end()) { - tsm = self->tag_data.insert( mapPair(std::move(Tag(tag->tag)), TLogData::TagData(Version(0), true, true, tag->tag) ), false ); - } - - if (version >= tsm->value.popped) { - for(int m = 0; m < tag->messageOffsets.size(); ++m) { - int offs = tag->messageOffsets[m]; - uint8_t const* p = offs < messages1.size() ? messages1.begin() + offs : messages.begin() + offs - messages1.size(); - tsm->value.version_messages.push_back(std::make_pair(version, LengthPrefixedStringRef((uint32_t*)p))); - if(tsm->value.version_messages.back().second.expectedSize() > SERVER_KNOBS->MAX_MESSAGE_SIZE) { - TraceEvent(SevWarnAlways, "LargeMessage").detail("Size", tsm->value.version_messages.back().second.expectedSize()); - } - if (tag->tag != txsTag) - expectedBytes += tsm->value.version_messages.back().second.expectedSize(); - - ++tagMessages; - } - } - - // The factor of VERSION_MESSAGES_OVERHEAD is intended to be an overestimate of the actual memory used to store this data in a std::deque. - // In practice, this number is probably something like 528/512 ~= 1.03, but this could vary based on the implementation. - // There will also be a fixed overhead per std::deque, but its size should be trivial relative to the size of the TLog - // queue and can be thought of as increasing the capacity of the queue slightly. - addedBytes += (tagMessages * sizeof(std::pair) * SERVER_KNOBS->VERSION_MESSAGES_OVERHEAD_FACTOR_1024THS) >> 10; - } - - self->version_sizes[version] = make_pair(expectedBytes, expectedBytes); - self->bytesInput += addedBytes; - - //TraceEvent("TLogPushed", self->dbgid).detail("Bytes", addedBytes).detail("MessageBytes", messages.size()).detail("Tags", tags.size()).detail("expectedBytes", expectedBytes).detail("mCount", mCount).detail("tCount", tCount); - } - - Version poppedVersion( TLogData* self, Tag tag) { - auto mapIt = self->tag_data.find(tag); - if (mapIt == self->tag_data.end()) - return Version(0); - return mapIt->value.popped; - } - - std::deque> & get_version_messages( TLogData* self, Tag tag ) { - auto mapIt = self->tag_data.find(tag); - if (mapIt == self->tag_data.end()) { - static std::deque> empty; - return empty; - } - return mapIt->value.version_messages; - }; - - ACTOR Future tLogPop( TLogData* self, TLogPopRequest req ) { - auto ti = self->tag_data.find(req.tag); - if (ti == self->tag_data.end()) { - ti = self->tag_data.insert( mapPair(std::move(Tag(req.tag)), TLogData::TagData(req.to, true, true, req.tag)) ); - } else if (req.to > ti->value.popped) { - ti->value.popped = req.to; - ti->value.popped_recently = true; - //if (to.epoch == self->epoch()) - if ( req.to > self->persistentDataDurableVersion ) - Void _ = wait(ti->value.eraseMessagesBefore( req.to, &self->bytesDurable, self, TaskTLogPop )); - //TraceEvent("TLogPop", self->dbgid).detail("Tag", req.tag).detail("To", req.to); - } - - req.reply.send(Void()); - return Void(); - } - - void peekMessagesFromMemory( TLogData* self, TLogPeekRequest const& req, BinaryWriter& messages, Version& endVersion ) { - ASSERT( !messages.getLength() ); - - auto& deque = get_version_messages(self, req.tag); - //TraceEvent("tLogPeekMem", self->dbgid).detail("Tag", printable(req.tag1)).detail("pDS", self->persistentDataSequence).detail("pDDS", self->persistentDataDurableSequence).detail("Oldest", map1.empty() ? 0 : map1.begin()->key ).detail("OldestMsgCount", map1.empty() ? 0 : map1.begin()->value.size()); - - Version begin = std::max( req.begin, self->persistentDataDurableVersion+1 ); - auto it = std::lower_bound(deque.begin(), deque.end(), std::make_pair(begin, LengthPrefixedStringRef()), CompareFirst>()); - - Version currentVersion = -1; - for(; it != deque.end(); ++it) { - if(it->first != currentVersion) { - if (messages.getLength() >= SERVER_KNOBS->DESIRED_TOTAL_BYTES) { - endVersion = it->first; - //TraceEvent("tLogPeekMessagesReached2", self->dbgid); - break; - } - - currentVersion = it->first; - messages << int32_t(-1) << currentVersion; - } - - messages << it->second.toStringRef(); - } - } - - ACTOR Future tLogPeekMessages( TLogData* self, TLogPeekRequest req ) { - state BinaryWriter messages(Unversioned()); - state BinaryWriter messages2(Unversioned()); - state int sequence = -1; - state UID peekId; - - if(req.sequence.present()) { - try { - peekId = req.sequence.get().first; - sequence = req.sequence.get().second; - if(sequence > 0) { - auto& trackerData = self->peekTracker[peekId]; - trackerData.lastUpdate = now(); - Version ver = wait(trackerData.sequence_version[sequence].getFuture()); - req.begin = ver; - Void _ = wait(yield()); - } - } catch( Error &e ) { - if(e.code() == error_code_timed_out) { - req.reply.sendError(timed_out()); - return Void(); - } else { - throw; - } - } - } - - if( req.returnIfBlocked && self->version.get() < req.begin ) { - req.reply.sendError(end_of_stream()); - return Void(); - } - - //TraceEvent("tLogPeekMessages0", self->dbgid).detail("reqBeginEpoch", req.begin.epoch).detail("reqBeginSeq", req.begin.sequence).detail("epoch", self->epoch()).detail("persistentDataSeq", self->persistentDataSequence).detail("Tag1", printable(req.tag1)).detail("Tag2", printable(req.tag2)); - // Wait until we have something to return that the caller doesn't already have - if( self->version.get() < req.begin ) { - Void _ = wait( self->version.whenAtLeast( req.begin ) ); - Void _ = wait( delay(SERVER_KNOBS->TLOG_PEEK_DELAY, g_network->getCurrentTask()) ); - } - - state Version endVersion = self->version.get() + 1; - - //grab messages from disk - //TraceEvent("tLogPeekMessages", self->dbgid).detail("reqBeginEpoch", req.begin.epoch).detail("reqBeginSeq", req.begin.sequence).detail("epoch", self->epoch()).detail("persistentDataSeq", self->persistentDataSequence).detail("Tag1", printable(req.tag1)).detail("Tag2", printable(req.tag2)); - if( req.begin <= self->persistentDataDurableVersion ) { - // Just in case the durable version changes while we are waiting for the read, we grab this data from memory. We may or may not actually send it depending on - // whether we get enough data from disk. - // SOMEDAY: Only do this if an initial attempt to read from disk results in insufficient data and the required data is no longer in memory - // SOMEDAY: Should we only send part of the messages we collected, to actually limit the size of the result? - - peekMessagesFromMemory( self, req, messages2, endVersion ); - - Standalone> kvs = wait( - self->persistentData->readRange(KeyRangeRef( - persistTagMessagesKey(req.tag, req.begin), - persistTagMessagesKey(req.tag, self->persistentDataDurableVersion + 1)), SERVER_KNOBS->DESIRED_TOTAL_BYTES, SERVER_KNOBS->DESIRED_TOTAL_BYTES)); - - //TraceEvent("TLogPeekResults", self->dbgid).detail("ForAddress", req.reply.getEndpoint().address).detail("Tag1Results", s1).detail("Tag2Results", s2).detail("Tag1ResultsLim", kv1.size()).detail("Tag2ResultsLim", kv2.size()).detail("Tag1ResultsLast", kv1.size() ? printable(kv1[0].key) : "").detail("Tag2ResultsLast", kv2.size() ? printable(kv2[0].key) : "").detail("Limited", limited).detail("NextEpoch", next_pos.epoch).detail("NextSeq", next_pos.sequence).detail("NowEpoch", self->epoch()).detail("NowSeq", self->sequence.getNextSequence()); - - for (auto &kv : kvs) { - auto ver = decodeTagMessagesKey(kv.key); - messages << int32_t(-1) << ver; - messages.serializeBytes(kv.value); - } - - if (kvs.expectedSize() >= SERVER_KNOBS->DESIRED_TOTAL_BYTES) - endVersion = decodeTagMessagesKey(kvs.end()[-1].key) + 1; - else - messages.serializeBytes( messages2.toStringRef() ); - } else { - peekMessagesFromMemory( self, req, messages, endVersion ); - //TraceEvent("TLogPeekResults", self->dbgid).detail("ForAddress", req.reply.getEndpoint().address).detail("MessageBytes", messages.getLength()).detail("NextEpoch", next_pos.epoch).detail("NextSeq", next_pos.sequence).detail("NowSeq", self->sequence.getNextSequence()); - } - - Version poppedVer = poppedVersion(self, req.tag); - - TLogPeekReply reply; - reply.maxKnownVersion = self->version.get(); - if(poppedVer > req.begin) { - reply.popped = poppedVer; - reply.end = poppedVer; - } else { - reply.messages = messages.toStringRef(); - reply.end = endVersion; - } - //TraceEvent("TlogPeek", self->dbgid).detail("endVer", reply.end).detail("msgBytes", reply.messages.expectedSize()).detail("ForAddress", req.reply.getEndpoint().address); - - if(req.sequence.present()) { - auto& trackerData = self->peekTracker[peekId]; - trackerData.lastUpdate = now(); - auto& sequenceData = trackerData.sequence_version[sequence+1]; - if(sequenceData.isSet()) { - if(sequenceData.getFuture().get() != reply.end) { - TEST(true); //tlog peek second attempt ended at a different version - req.reply.sendError(timed_out()); - return Void(); - } - } else { - sequenceData.send(reply.end); - } - } - - req.reply.send( reply ); - return Void(); - } - - ACTOR Future doQueueCommit( TLogData* self ) { - state Version ver = self->version.get(); - state Version commitNumber = self->queueCommitBegin+1; - self->queueCommitBegin = commitNumber; - self->queueCommittingVersion = ver; - - Future c = self->persistentQueue->commit(); - self->diskQueueCommitBytes = 0; - self->largeDiskQueueCommitBytes.set(false); - - Void _ = wait(c); - Void _ = wait(self->queueCommitEnd.whenAtLeast(commitNumber-1)); - - //Calling check_yield instead of yield to avoid a destruction ordering problem in simulation - if(g_network->check_yield(g_network->getCurrentTask())) { - Void _ = wait(delay(0, g_network->getCurrentTask())); - } - - ASSERT( ver > self->queueCommittedVersion.get() ); - - self->queueCommittedVersion.set(ver); - self->queueCommitEnd.set(commitNumber); - - TraceEvent("TLogCommitDurable", self->dbgid).detail("Version", ver); - - return Void(); - } - - ACTOR Future commitQueue( TLogData* self ) { - loop { - Void _ = wait( self->version.whenAtLeast( std::max(self->queueCommittingVersion, self->queueCommittedVersion.get()) + 1 ) ); - while( self->queueCommitBegin != self->queueCommitEnd.get() && !self->largeDiskQueueCommitBytes.get() ) { - Void _ = wait( self->queueCommitEnd.whenAtLeast(self->queueCommitBegin) || self->largeDiskQueueCommitBytes.onChange() ); - } - self->addActor.send(doQueueCommit(self)); - } - } - - ACTOR Future tLogCommit( - TLogData* self, - TLogCommitRequest req, - PromiseStream warningCollectorInput ) { - - state Optional tlogDebugID; - if(req.debugID.present()) - { - tlogDebugID = g_nondeterministic_random->randomUniqueID(); - g_traceBatch.addAttach("CommitAttachID", req.debugID.get().first(), tlogDebugID.get().first()); - g_traceBatch.addEvent("CommitDebug", tlogDebugID.get().first(), "TLog.tLogCommit.BeforeWaitForVersion"); - } - - self->knownCommittedVersion = std::max(self->knownCommittedVersion, req.knownCommittedVersion); - - Void _ = wait( self->version.whenAtLeast( req.prevVersion ) ); - - //Calling check_yield instead of yield to avoid a destruction ordering problem in simulation - if(g_network->check_yield(g_network->getCurrentTask())) { - Void _ = wait(delay(0, g_network->getCurrentTask())); - } - - if(self->stopped) { - req.reply.sendError( tlog_stopped() ); - return Void(); - } - - if (self->version.get() == req.prevVersion) { // Not a duplicate (check relies on no waiting between here and self->version.set() below!) - if(req.debugID.present()) - g_traceBatch.addEvent("CommitDebug", tlogDebugID.get().first(), "TLog.tLogCommit.Before"); - - TraceEvent("TLogCommit", self->dbgid).detail("Version", req.version); - commitMessages(self, req.version, req.arena, req.messages, req.tags); - - // Log the changes to the persistent queue, to be committed by commitQueue() - TLogQueueEntryRef qe; - qe.version = req.version; - qe.knownCommittedVersion = req.knownCommittedVersion; - qe.messages = req.messages; - qe.tags = req.tags; - self->persistentQueue->push( qe ); - - self->diskQueueCommitBytes += qe.expectedSize(); - if( self->diskQueueCommitBytes > SERVER_KNOBS->MAX_QUEUE_COMMIT_BYTES ) { - self->largeDiskQueueCommitBytes.set(true); - } - - // Notifies the commitQueue actor to commit persistentQueue, and also unblocks tLogPeekMessages actors - self->prevVersion = self->version.get(); - self->version.set( req.version ); - - if(req.debugID.present()) - g_traceBatch.addEvent("CommitDebug", tlogDebugID.get().first(), "TLog.tLogCommit.AfterTLogCommit"); - } - // Send replies only once all prior messages have been received and committed. - Void _ = wait( timeoutWarning( self->queueCommittedVersion.whenAtLeast( req.version ), 0.1, warningCollectorInput ) ); - - if(req.debugID.present()) - g_traceBatch.addEvent("CommitDebug", tlogDebugID.get().first(), "TLog.tLogCommit.After"); - - req.reply.send( Void() ); - return Void(); - } - - ACTOR Future initPersistentState( TLogData* self ) { - // PERSIST: Initial setup of persistentData for a brand new tLog for a new database - IKeyValueStore *storage = self->persistentData; - storage->set( persistFormat ); - storage->set( KeyValueRef( persistID, BinaryWriter::toValue( self->dbgid, Unversioned() ) ) ); - storage->set( KeyValueRef( persistCurrentVersionKey, BinaryWriter::toValue(self->version.get(), Unversioned()) ) ); - storage->set( KeyValueRef( persistRecoveryCountKey, BinaryWriter::toValue(self->recoveryCount, Unversioned()) ) ); - - TraceEvent("TLogInitCommit", self->dbgid).detail("Version", self->version.get()); - Void _ = wait( storage->commit() ); - return Void(); - } - - ACTOR Future restorePersistentState( TLogData* self, Promise outRecoveryCount, bool processQueue, TLogInterface myInterface ) { - state double startt = now(); - // PERSIST: Read basic state from persistentData; replay persistentQueue but don't erase it - IKeyValueStore *storage = self->persistentData; - - TraceEvent("TLogRestorePersistentState", self->dbgid).detail("pq", processQueue); - - state Future> fFormat = storage->readValue(persistFormat.key); - state Future> fID = storage->readValue(persistID); - state Future> fVer = storage->readValue(persistCurrentVersionKey); - state Future> fRecoveryCount = storage->readValue(persistRecoveryCountKey); - state Future> fRecoveryInProgress = storage->readValue( persistRecoveryInProgress.key ); - - // FIXME: metadata in queue? - - Void _ = wait( waitForAll( (vector>>(), fFormat, fID, fVer, fRecoveryCount, fRecoveryInProgress) ) ); - - if (fFormat.get().present() && !persistFormatReadableRange.contains( fFormat.get().get() )) { - TraceEvent(SevError, "UnsupportedDBFormat", self->dbgid).detail("Format", printable(fFormat.get().get())).detail("Expected", persistFormat.value.toString()); - throw worker_recovery_failed(); - } - - if (fRecoveryInProgress.get().present()) { - TEST(true); // We must have rebooted during network recovery; the master recovery that depended on us will fail and we can permanently delete our (incomplete) storage - TraceEvent("RestartedDuringNetworkRecovery", self->dbgid); - throw worker_removed(); - } - - if (!fFormat.get().present()) { - Standalone> v = wait( self->persistentData->readRange( KeyRangeRef(StringRef(), LiteralStringRef("\xff")), 1 ) ); - if (!v.size()) { - TEST(true); // The DB is completely empty, so it was never initialized. Delete it. - throw worker_removed(); - } else { - // This should never happen - TraceEvent(SevError, "NoDBFormatKey", self->dbgid).detail("FirstKey", printable(v[0].key)); - ASSERT( false ); - throw worker_recovery_failed(); - } - } - - - ASSERT( self->dbgid == BinaryReader::fromStringRef(fID.get().get(), Unversioned()) ); - - Version ver = BinaryReader::fromStringRef( fVer.get().get(), Unversioned() ); - self->persistentDataVersion = ver; - self->persistentDataDurableVersion = ver; - self->version.set( ver ); - - TraceEvent("TLogRestorePersistentStateVer", self->dbgid).detail("ver", self->version.get()); - - self->recoveryCount = BinaryReader::fromStringRef( fRecoveryCount.get().get(), Unversioned() ); - - outRecoveryCount.send( self->recoveryCount ); // This might cancel this actor (if the recovery count is ancient) and destroy self - Void _ = wait(Future(Void())); // ... so check for cancellation - - // Restore popped keys. Pop operations that took place after the last (committed) updatePersistentDataVersion might be lost, but - // that is fine because we will get the corresponding data back, too. - state KeyRange tagKeys = persistTagPoppedKeys; - loop { - Standalone> data = wait( self->persistentData->readRange( tagKeys, BUGGIFY ? 3 : 1<<30, 1<<20 ) ); - if (!data.size()) break; - ((KeyRangeRef&)tagKeys) = KeyRangeRef( keyAfter(data.back().key, tagKeys.arena()), tagKeys.end ); - - for(auto &kv : data) { - Tag tag = decodeTagPoppedKey(kv.key); - Version popped = decodeTagPoppedValue(kv.value); - TraceEvent("TLogRestorePop", self->dbgid).detail("Tag", tag).detail("To", popped); - ASSERT( self->tag_data.find(tag) == self->tag_data.end() ); - self->tag_data.insert( mapPair( std::move(Tag(tag)), TLogData::TagData( popped, false, false, tag )) ); - } - } - - // PERSIST: Apply changes from queue - if (processQueue) { - state Version lastVer = 0; - state double recoverMemoryLimit = SERVER_KNOBS->TARGET_BYTES_PER_TLOG + SERVER_KNOBS->SPRING_BYTES_TLOG; - if (BUGGIFY) recoverMemoryLimit = SERVER_KNOBS->BUGGIFY_RECOVER_MEMORY_LIMIT; - try { - loop { - TLogQueueEntry qe = wait( self->persistentQueue->readNext() ); - //TraceEvent("TLogRecoveredQE", self->dbgid).detail("ver", qe.version).detail("MessageBytes", qe.messages.size()).detail("Tags", qe.tags.size()) - // .detail("Tag0", qe.tags.size() ? qe.tags[0].tag : invalidTag); - - ASSERT( qe.version > lastVer ); - lastVer = qe.version; - self->knownCommittedVersion = std::max(self->knownCommittedVersion, qe.knownCommittedVersion); - if( qe.version > self->version.get() ) { - commitMessages(self, qe.version, qe.arena(), qe.messages, qe.tags); - self->version.set( qe.version ); - self->queueCommittedVersion.set( qe.version ); - - if (self->bytesInput.getValue() - self->bytesDurable.getValue() > recoverMemoryLimit) { - TEST(true); // Flush excess data during TLog queue recovery - TraceEvent("FlushLargeQueueDuringRecovery", self->dbgid).detail("BytesInput", self->bytesInput.getValue()).detail("BytesDurable", self->bytesDurable.getValue()).detail("Version", self->version.get()).detail("PVer", self->persistentDataVersion); - - while(self->persistentDataDurableVersion != self->version.get()) { - Version nextVersion; - int totalSize = 0; - std::vector>::iterator, std::deque>::iterator>> iters; - for(auto tag = self->tag_data.begin(); tag != self->tag_data.end(); ++tag) - iters.push_back(std::make_pair(tag->value.version_messages.begin(), tag->value.version_messages.end())); - - while( totalSize < SERVER_KNOBS->UPDATE_STORAGE_BYTE_LIMIT ) { - nextVersion = self->version.get(); - for( auto &it : iters ) - if(it.first != it.second) - nextVersion = std::min( nextVersion, it.first->first + 1 ); - - if(nextVersion == self->version.get()) - break; - - for( auto &it : iters ) { - while (it.first != it.second && it.first->first < nextVersion) { - totalSize += it.first->second.expectedSize(); - ++it.first; - } - } - } - - Void _ = wait( updatePersistentData(self, nextVersion ) ); - } - } - } - } - } catch (Error& e) { - if (e.code() != error_code_end_of_stream) throw; - } - } - - TraceEvent("TLogRestorePersistentStateDone", self->dbgid) - .detail("pq", processQueue).detail("version", self->version.get()).detail("durableVer", self->persistentDataDurableVersion) - .detail("Took", now()-startt); - TEST( now()-startt >= 1.0 ); // TLog recovery took more than 1 second - TEST( processQueue ); // TLog recovered from disk queue - - return Void(); - } - - void getQueuingMetrics( TLogData* self, TLogQueuingMetricsRequest const& req ) { - TLogQueuingMetricsReply reply; - reply.localTime = now(); - reply.instanceID = self->instanceID; - reply.bytesInput = self->bytesInput.getValue(); - reply.bytesDurable = self->bytesDurable.getValue(); - reply.storageBytes = self->persistentData->getStorageBytes(); - reply.v = self->prevVersion; - req.reply.send( reply ); - } - - ACTOR Future respondToRecovered( TLogInterface tli, Future recovery ) { - Void _ = wait( recovery ); - - loop { - TLogRecoveryFinishedRequest req = waitNext( tli.recoveryFinished.getFuture() ); - req.reply.send(Void()); - } - } - - ACTOR Future cleanupPeekTrackers( TLogData* self ) { - loop { - double minExpireTime = SERVER_KNOBS->PEEK_TRACKER_EXPIRATION_TIME; - auto it = self->peekTracker.begin(); - while(it != self->peekTracker.end()) { - double expireTime = SERVER_KNOBS->PEEK_TRACKER_EXPIRATION_TIME - now()-it->second.lastUpdate; - if(expireTime < 1.0e-6) { - for(auto seq : it->second.sequence_version) { - if(!seq.second.isSet()) { - seq.second.sendError(timed_out()); - } - } - it = self->peekTracker.erase(it); - } else { - minExpireTime = std::min(minExpireTime, expireTime); - ++it; - } - } - - Void _ = wait( delay(minExpireTime) ); - } - } - - ACTOR Future serveTLogInterface( TLogData* self, TLogInterface tli, PromiseStream warningCollectorInput ) { - loop choose { - when( TLogPeekRequest req = waitNext( tli.peekMessages.getFuture() ) ) { - self->addActor.send( tLogPeekMessages( self, req ) ); - } - when( TLogPopRequest req = waitNext( tli.popMessages.getFuture() ) ) { - self->addActor.send( tLogPop( self, req ) ); - } - when( TLogCommitRequest req = waitNext( tli.commit.getFuture() ) ) { - TEST(self->stopped); // TLogCommitRequest while stopped - if (!self->stopped) - self->addActor.send( tLogCommit( self, req, warningCollectorInput ) ); - else - req.reply.sendError( tlog_stopped() ); - } - when( ReplyPromise< TLogLockResult > reply = waitNext( tli.lock.getFuture() ) ) { - self->addActor.send( tLogLock(self, reply) ); - } - when (TLogQueuingMetricsRequest req = waitNext(tli.getQueuingMetrics.getFuture())) { - getQueuingMetrics(self, req); - } - when (TLogConfirmRunningRequest req = waitNext(tli.confirmRunning.getFuture())){ - if (req.debugID.present() ) { - UID tlogDebugID = g_nondeterministic_random->randomUniqueID(); - g_traceBatch.addAttach("TransactionAttachID", req.debugID.get().first(), tlogDebugID.first()); - g_traceBatch.addEvent("TransactionDebug", tlogDebugID.first(), "TLogServer.TLogConfirmRunningRequest"); - } - if (!self->stopped) - req.reply.send(Void()); - else - req.reply.sendError( tlog_stopped() ); - } - } - } - - ACTOR Future tLogCore( TLogData* self, TLogInterface tli, Future recovery ) { - state PromiseStream warningCollectorInput; - state Future warningCollector = timeoutWarningCollector( warningCollectorInput.getFuture(), 1.0, "TLogQueueCommitSlow", self->dbgid ); - state Future error = actorCollection( self->addActor.getFuture() ); - - self->addActor.send( updateStorage(self) ); - self->addActor.send( commitQueue(self) ); - self->addActor.send( waitFailureServer(tli.waitFailure.getFuture()) ); - self->addActor.send( respondToRecovered(tli, recovery) ); - self->addActor.send( traceCounters("TLogMetrics", self->dbgid, SERVER_KNOBS->STORAGE_LOGGING_DELAY, &self->cc, self->dbgid.toString() + "/TLogMetrics")); - self->addActor.send( cleanupPeekTrackers(self) ); - - if( recovery.isValid() && !recovery.isReady()) { - self->addActor.send( recovery ); - } - - self->coreStarted = true; - - Void _ = wait( serveTLogInterface(self, tli, warningCollectorInput) || error ); - throw internal_error(); - }; - - ACTOR Future checkEmptyQueue(TLogData* self) { - TraceEvent("TLogCheckEmptyQueueBegin", self->dbgid); - try { - TLogQueueEntry r = wait( self->persistentQueue->readNext() ); - throw internal_error(); - } catch (Error& e) { - if (e.code() != error_code_end_of_stream) throw; - TraceEvent("TLogCheckEmptyQueueEnd", self->dbgid); - return Void(); - } - } - - ACTOR Future recoverTagFromLogSystem( TLogData* self, Version beginVersion, Version endVersion, Tag tag, Reference> uncommittedBytes, Reference>> logSystem ) { - state Future dbInfoChange = Void(); - state Reference r; - state Version tagAt = beginVersion; - state Version tagPopped = 0; - state Version lastVer = 0; - - TraceEvent("LogRecoveringTagBegin", self->dbgid).detail("Tag", tag).detail("recoverAt", endVersion); - - while (tagAt <= endVersion) { - loop { - choose { - when(Void _ = wait( r ? r->getMore() : Never() ) ) { - break; - } - when( Void _ = wait( dbInfoChange ) ) { - if(r) tagPopped = std::max(tagPopped, r->popped()); - if( logSystem->get() ) - r = logSystem->get()->peek( tagAt, tag ); - else - r = Reference(); - dbInfoChange = logSystem->onChange(); - } - } - } - - TraceEvent("LogRecoveringTagResults", self->dbgid).detail("Tag", tag); - - Version ver = 0; - BinaryWriter wr( Unversioned() ); - int writtenBytes = 0; - while (true) { - bool foundMessage = r->hasMessage(); - //TraceEvent("LogRecoveringMsg").detail("Tag", tag).detail("foundMessage", foundMessage).detail("ver", r->version().toString()); - if (!foundMessage || r->version().version != ver) { - ASSERT(r->version().version > lastVer); - if (ver) { - //TraceEvent("LogRecoveringTagVersion", self->dbgid).detail("Tag", tag).detail("Ver", ver).detail("Bytes", wr.getLength()); - writtenBytes += 100 + wr.getLength(); - self->persistentData->set( KeyValueRef( persistTagMessagesKey( tag, ver ), wr.toStringRef() ) ); - } - lastVer = ver; - ver = r->version().version; - wr = BinaryWriter( Unversioned() ); - if (!foundMessage || ver > endVersion) - break; - } - - // FIXME: This logic duplicates stuff in LogPushData::addMessage(), and really would be better in PeekResults or somewhere else. Also unnecessary copying. - StringRef msg = r->getMessage(); - wr << uint32_t( msg.size() + sizeof(uint32_t) ) << r->version().sub; - wr.serializeBytes( msg ); - r->nextMessage(); - } - - tagAt = r->version().version; - - if(writtenBytes) - uncommittedBytes->set(uncommittedBytes->get() + writtenBytes); - - while(uncommittedBytes->get() >= SERVER_KNOBS->LARGE_TLOG_COMMIT_BYTES) { - Void _ = wait(uncommittedBytes->onChange()); - } - } - if(r) tagPopped = std::max(tagPopped, r->popped()); - - auto tsm = self->tag_data.find(tag); - if (tsm == self->tag_data.end()) { - self->tag_data.insert( mapPair(std::move(Tag(tag)), TLogData::TagData(tagPopped, false, true, tag)) ); - } - - Void _ = wait(tLogPop( self, TLogPopRequest(tagPopped, tag) )); - - updatePersistentPopped( self, tag, self->tag_data.find(tag)->value ); - return Void(); - } - - ACTOR Future updateLogSystem(TLogData* self, LogSystemConfig recoverFrom, Reference>> logSystem) { - loop { - TraceEvent("TLogUpdate", self->dbgid).detail("recoverFrom", recoverFrom.toString()).detail("dbInfo", self->dbInfo->get().logSystemConfig.toString()); - if( self->dbInfo->get().logSystemConfig.isEqualIds(recoverFrom) ) { - logSystem->set(ILogSystem::fromLogSystemConfig( self->dbgid, self->dbInfo->get().myLocality, self->dbInfo->get().logSystemConfig )); - } else if( self->dbInfo->get().logSystemConfig.isNextGenerationOf(recoverFrom) && std::count( self->dbInfo->get().logSystemConfig.tLogs.begin(), self->dbInfo->get().logSystemConfig.tLogs.end(), self->dbgid ) ) { - logSystem->set(ILogSystem::fromOldLogSystemConfig( self->dbgid, self->dbInfo->get().myLocality, self->dbInfo->get().logSystemConfig )); - } else { - logSystem->set(Reference()); - } - Void _ = wait( self->dbInfo->onChange() ); - } - } - - ACTOR Future recoverFromLogSystem( TLogData* self, LogSystemConfig recoverFrom, Version recoverAt, Version knownCommittedVersion, std::vector recoverTags, Promise copyComplete ) { - state Future committing = Void(); - state double lastCommitT = now(); - state Reference> uncommittedBytes = Reference>(new AsyncVar()); - state std::vector> recoverFutures; - state Reference>> logSystem = Reference>>(new AsyncVar>()); - state Future updater = updateLogSystem(self, recoverFrom, logSystem); - - for(auto tag : recoverTags ) - recoverFutures.push_back(recoverTagFromLogSystem(self, knownCommittedVersion, recoverAt, tag, uncommittedBytes, logSystem)); - - state Future copyDone = waitForAll(recoverFutures); - state Future recoveryDone = Never(); - state Future commitTimeout = delay(SERVER_KNOBS->LONG_TLOG_COMMIT_TIME); - - loop { - choose { - when(Void _ = wait(copyDone)) { - recoverFutures.clear(); - for(auto tag : recoverTags ) - recoverFutures.push_back(recoverTagFromLogSystem(self, 0, knownCommittedVersion, tag, uncommittedBytes, logSystem)); - copyDone = Never(); - recoveryDone = waitForAll(recoverFutures); - - Void __ = wait( committing ); - Void __ = wait( self->updatePersist ); - committing = self->persistentData->commit(); - commitTimeout = delay(SERVER_KNOBS->LONG_TLOG_COMMIT_TIME); - uncommittedBytes->set(0); - Void __ = wait( committing ); - TraceEvent("TLogCommitCopyData", self->dbgid); - - if(!copyComplete.isSet()) - copyComplete.send(Void()); - } - when(Void _ = wait(recoveryDone)) { break; } - when(Void _ = wait(commitTimeout)) { - TEST(true); // We need to commit occasionally if this process is long to avoid running out of memory. - // We let one, but not more, commits pipeline with the network transfer - Void __ = wait( committing ); - Void __ = wait( self->updatePersist ); - committing = self->persistentData->commit(); - commitTimeout = delay(SERVER_KNOBS->LONG_TLOG_COMMIT_TIME); - uncommittedBytes->set(0); - TraceEvent("TLogCommitRecoveryData", self->dbgid).detail("MemoryUsage", DEBUG_DETERMINISM ? 0 : getMemoryUsage()); - } - when(Void _ = wait(uncommittedBytes->onChange())) { - if(uncommittedBytes->get() >= SERVER_KNOBS->LARGE_TLOG_COMMIT_BYTES) - commitTimeout = Void(); - } - } - } - - Void _ = wait( committing ); - Void _ = wait( self->updatePersist ); - Void _ = wait( self->persistentData->commit() ); - - TraceEvent("TLogRecoveryComplete", self->dbgid).detail("Locality", self->dbInfo->get().myLocality.toString()); - TEST(true); // tLog restore from old log system completed - - return Void(); - } - - ACTOR Future tLogStart( TLogData* self, LogSystemConfig recoverFrom, Version recoverAt, Version knownCommittedVersion, std::vector recoverTags, bool recoverFromDisk, - TLogInterface tli, ReplyPromise outInterface, Promise outRecoveryCount ) { - state Future recovery = Void(); - if (recoverFrom.logSystemType == 1) { - ASSERT(false); - } else if (recoverFrom.logSystemType == 2) { - Void _ = wait( checkEmptyQueue(self) ); - - self->persistentDataVersion = recoverAt; - self->persistentDataDurableVersion = recoverAt; // durable is a white lie until initPersistentState() commits the store - self->queueCommittedVersion.set( recoverAt ); - self->version.set( recoverAt ); - - Void _ = wait( initPersistentState( self ) ); - - state Promise copyComplete; - recovery = recoverFromLogSystem( self, recoverFrom, recoverAt, knownCommittedVersion, recoverTags, copyComplete ); - Void _ = wait(copyComplete.getFuture()); - } else if (recoverFromDisk) { - Void _ = wait( restorePersistentState( self, outRecoveryCount, true, tli ) ); - TEST(true); // tLog restore from disk completed - } else { - // Brand new tlog, initialization has already been done by caller - Void _ = wait( checkEmptyQueue(self) ); - Void _ = wait( initPersistentState( self ) ); - } - - TraceEvent("TLogReady", self->dbgid); - - validate(self); - //dump(self); - - outInterface.send( tli ); - - Void _ = wait( tLogCore( self, tli, recovery ) ); - throw internal_error(); // tLogCore() shouldn't return without an error - } - - ACTOR Future rejoinMasters( TLogData* self, TLogInterface tli, Future fRecoveryCount ) { - state DBRecoveryCount recoveryCount = wait( fRecoveryCount ); - state UID lastMasterID(0,0); - loop { - auto const& inf = self->dbInfo->get(); - bool isDisplaced = inf.recoveryCount >= recoveryCount && inf.recoveryState != 0 && - !std::count( inf.logSystemConfig.tLogs.begin(), inf.logSystemConfig.tLogs.end(), tli.id() ) && - !std::count( inf.priorCommittedLogServers.begin(), inf.priorCommittedLogServers.end(), tli.id() ); - for(int i = 0; i < inf.logSystemConfig.oldTLogs.size() && isDisplaced; i++) { - isDisplaced = !std::count( inf.logSystemConfig.oldTLogs[i].tLogs.begin(), inf.logSystemConfig.oldTLogs[i].tLogs.end(), tli.id() ); - } - if ( isDisplaced ) - { - TraceEvent("TLogDisplaced", tli.id()).detail("Reason", "DBInfoDoesNotContain"); - if (BUGGIFY) Void _ = wait( delay( SERVER_KNOBS->BUGGIFY_WORKER_REMOVED_MAX_LAG * g_random->random01() ) ); - throw worker_removed(); - } - - if (self->dbInfo->get().master.id() != lastMasterID) { - // The TLogRejoinRequest is needed to establish communications with a new master, which doesn't have our TLogInterface - TLogRejoinRequest req; - req.myInterface = tli; - TraceEvent("TLogRejoining", self->dbgid).detail("Master", self->dbInfo->get().master.id()); - choose { - when ( bool success = wait( brokenPromiseToNever( self->dbInfo->get().master.tlogRejoin.getReply( req ) ) ) ) { - if (success) - lastMasterID = self->dbInfo->get().master.id(); - } - when ( Void _ = wait( self->dbInfo->onChange() ) ) { } - } - } else - Void _ = wait( self->dbInfo->onChange() ); - } - } - - // Restore from disk - ACTOR Future tLog( IKeyValueStore* persistentData, IDiskQueue* persistentQueue, TLogInterface tli, Reference> db ) { - state TLogData self( tli.id(), persistentData, persistentQueue, db ); - state Promise recoveryCount; - state Future removed = rejoinMasters(&self, tli, recoveryCount.getFuture()); - - Void _ = wait( tLogStart( &self, LogSystemConfig(), Version(0), Version(0), std::vector(), true, tli, ReplyPromise(), recoveryCount ) || removed ); - throw internal_error(); // tLogStart doesn't return without an error - } -} \ No newline at end of file diff --git a/fdbserver/Status.actor.cpp b/fdbserver/Status.actor.cpp index 40ec4f403f..67e4d24528 100644 --- a/fdbserver/Status.actor.cpp +++ b/fdbserver/Status.actor.cpp @@ -1447,26 +1447,34 @@ static StatusArray oldTlogFetcher(int* oldLogFaultTolerance, Referenceget().recoveryState == RecoveryState::FULLY_RECOVERED) { for(auto it : db->get().logSystemConfig.oldTLogs) { StatusObject statusObj; - int failedLogs = 0; StatusArray logsObj; - for(auto log : it.tLogs) { - StatusObject logObj; - bool failed = !log.present() || !address_workers.count(log.interf().address()); - logObj["id"] = log.id().shortString(); - logObj["healthy"] = !failed; - if(log.present()) { - logObj["address"] = log.interf().address().toString(); + int maxFaultTolerance = 0; + + for(int i = 0; i < it.tLogs.size(); i++) { + int failedLogs = 0; + for(auto& log : it.tLogs[i].tLogs) { + StatusObject logObj; + bool failed = !log.present() || !address_workers.count(log.interf().address()); + logObj["id"] = log.id().shortString(); + logObj["healthy"] = !failed; + if(log.present()) { + logObj["address"] = log.interf().address().toString(); + } + logsObj.push_back(logObj); + if(failed) { + failedLogs++; + } } - logsObj.push_back(logObj); - if(failed) { - failedLogs++; + maxFaultTolerance = std::max(maxFaultTolerance, it.tLogs[i].tLogReplicationFactor - 1 - it.tLogs[i].tLogWriteAntiQuorum - failedLogs); + //FIXME: add information for remote and satellites + if(i==0) { + statusObj["log_replication_factor"] = it.tLogs[i].tLogReplicationFactor; + statusObj["log_write_anti_quorum"] = it.tLogs[i].tLogWriteAntiQuorum; + statusObj["log_fault_tolerance"] = it.tLogs[i].tLogReplicationFactor - 1 - it.tLogs[i].tLogWriteAntiQuorum - failedLogs; } } - *oldLogFaultTolerance = std::min(*oldLogFaultTolerance, it.tLogReplicationFactor - 1 - it.tLogWriteAntiQuorum - failedLogs); + *oldLogFaultTolerance = std::min(*oldLogFaultTolerance, maxFaultTolerance); statusObj["logs"] = logsObj; - statusObj["log_replication_factor"] = it.tLogReplicationFactor; - statusObj["log_write_anti_quorum"] = it.tLogWriteAntiQuorum; - statusObj["log_fault_tolerance"] = it.tLogReplicationFactor - 1 - it.tLogWriteAntiQuorum - failedLogs; oldTlogsArray.push_back(statusObj); } } diff --git a/fdbserver/TLogServer.actor.cpp b/fdbserver/TLogServer.actor.cpp index 73d0f9e0ee..df6bcfbab9 100644 --- a/fdbserver/TLogServer.actor.cpp +++ b/fdbserver/TLogServer.actor.cpp @@ -213,7 +213,6 @@ struct TLogData : NonCopyable { WorkerCache tlogCache; Future updatePersist; //SOMEDAY: integrate the recovery and update storage so that only one of them is committing to persistant data. - Future oldLogServer; PromiseStream> sharedActors; @@ -395,7 +394,7 @@ KeyRange prefixRange( KeyRef prefix ) { // Immutable keys static const KeyValueRef persistFormat( LiteralStringRef( "Format" ), LiteralStringRef("FoundationDB/LogServer/2/4") ); -static const KeyRangeRef persistFormatReadableRange( LiteralStringRef("FoundationDB/LogServer/2/2"), LiteralStringRef("FoundationDB/LogServer/2/5") ); +static const KeyRangeRef persistFormatReadableRange( LiteralStringRef("FoundationDB/LogServer/2/3"), LiteralStringRef("FoundationDB/LogServer/2/5") ); static const KeyRangeRef persistRecoveryCountKeys = KeyRangeRef( LiteralStringRef( "DbRecoveryCount/" ), LiteralStringRef( "DbRecoveryCount0" ) ); // Updated on updatePersistentData() @@ -1195,8 +1194,8 @@ ACTOR Future serveTLogInterface( TLogData* self, TLogInterface tli, Refere dbInfoChange = self->dbInfo->onChange(); bool found = false; if(self->dbInfo->get().recoveryState >= RecoveryState::FULLY_RECOVERED) { - for(auto& log : self->dbInfo->get().logSystemConfig.tLogs) { - if( std::count( self->dbInfo->get().logSystemConfig.tLogs.begin(), self->dbInfo->get().logSystemConfig.tLogs.end(), logData->logId ) ) { + for(auto& logs : self->dbInfo->get().logSystemConfig.tLogs) { + if( std::count( logs.tLogs.begin(), logs.tLogs.end(), logData->logId ) ) { found = true; break; } @@ -1249,7 +1248,7 @@ void removeLog( TLogData* self, Reference logData ) { logData->addActor = PromiseStream>(); //there could be items still in the promise stream if one of the actors threw an error immediately self->id_data.erase(logData->logId); - if(self->id_data.size() || (self->oldLogServer.isValid() && !self->oldLogServer.isReady())) { + if(self->id_data.size()) { return; } else { throw worker_removed(); @@ -1417,7 +1416,7 @@ ACTOR Future checkEmptyQueue(TLogData* self) { } } -ACTOR Future restorePersistentState( TLogData* self, LocalityData locality, Promise oldLog, PromiseStream tlogRequests ) { +ACTOR Future restorePersistentState( TLogData* self, LocalityData locality, PromiseStream tlogRequests ) { state double startt = now(); state Reference logData; state KeyRange tagKeys; @@ -1455,28 +1454,7 @@ ACTOR Future restorePersistentState( TLogData* self, LocalityData locality state std::vector>> removed; state int persistentDataFormat = 0; - if(fFormat.get().get() == LiteralStringRef("FoundationDB/LogServer/2/2")) { - TLogInterface recruited; - recruited.uniqueID = self->dbgid; - recruited.locality = locality; - recruited.initEndpoints(); - - DUMPTOKEN( recruited.peekMessages ); - DUMPTOKEN( recruited.popMessages ); - DUMPTOKEN( recruited.commit ); - DUMPTOKEN( recruited.lock ); - DUMPTOKEN( recruited.getQueuingMetrics ); - DUMPTOKEN( recruited.confirmRunning ); - - //FIXME: need for upgrades from 4.X to 5.0, remove once this upgrade path is no longer needed - oldLog.send(Void()); - while(!tlogRequests.isEmpty()) { - tlogRequests.getFuture().pop().reply.sendError(recruitment_failed()); - } - - Void _ = wait( oldTLog::tLog(self->persistentData, self->rawPersistentQueue, recruited, self->dbInfo) ); - throw internal_error(); - } else if(fFormat.get().get() >= LiteralStringRef("FoundationDB/LogServer/2/4")) { + if(fFormat.get().get() >= LiteralStringRef("FoundationDB/LogServer/2/4")) { persistentDataFormat = 1; } @@ -1857,7 +1835,7 @@ ACTOR Future tLogStart( TLogData* self, InitializeTLogRequest req, Localit } // New tLog (if !recoverFrom.size()) or restore from network -ACTOR Future tLog( IKeyValueStore* persistentData, IDiskQueue* persistentQueue, Reference> db, LocalityData locality, PromiseStream tlogRequests, UID tlogId, bool restoreFromDisk, Promise oldLog ) +ACTOR Future tLog( IKeyValueStore* persistentData, IDiskQueue* persistentQueue, Reference> db, LocalityData locality, PromiseStream tlogRequests, UID tlogId, bool restoreFromDisk ) { state TLogData self( tlogId, persistentData, persistentQueue, db ); state Future error = actorCollection( self.sharedActors.getFuture() ); @@ -1866,7 +1844,7 @@ ACTOR Future tLog( IKeyValueStore* persistentData, IDiskQueue* persistentQ try { if(restoreFromDisk) { - Void _ = wait( restorePersistentState( &self, locality, oldLog, tlogRequests ) ); + Void _ = wait( restorePersistentState( &self, locality, tlogRequests ) ); } else { Void _ = wait( checkEmptyQueue(&self) ); } diff --git a/fdbserver/TagPartitionedLogSystem.actor.cpp b/fdbserver/TagPartitionedLogSystem.actor.cpp index 9aab8461a7..4b8d3df68c 100644 --- a/fdbserver/TagPartitionedLogSystem.actor.cpp +++ b/fdbserver/TagPartitionedLogSystem.actor.cpp @@ -29,12 +29,6 @@ #include "fdbrpc/Replication.h" #include "fdbrpc/ReplicationUtils.h" -template -void uniquify( Collection& c ) { - std::sort(c.begin(), c.end()); - c.resize( std::unique(c.begin(), c.end()) - c.begin() ); -} - ACTOR static Future reportTLogCommitErrors( Future commitReply, UID debugID ) { try { Void _ = wait(commitReply); @@ -48,103 +42,6 @@ ACTOR static Future reportTLogCommitErrors( Future commitReply, UID } } -class LogSet { -public: - std::vector>>> logServers; - std::vector>>> logRouters; - int32_t tLogWriteAntiQuorum; - int32_t tLogReplicationFactor; - std::vector< LocalityData > tLogLocalities; // Stores the localities of the log servers - IRepPolicyRef tLogPolicy; - LocalitySetRef logServerSet; - std::vector logIndexArray; - std::map logEntryMap; - bool isLocal; - bool hasBest; - - LogSet() : tLogWriteAntiQuorum(0), tLogReplicationFactor(0), isLocal(true), hasBest(true) {} - - int bestLocationFor( Tag tag ) { - return hasBest ? tag % logServers.size() : invalidTag; - } - - void updateLocalitySet() { - LocalityMap* logServerMap; - logServerSet = LocalitySetRef(new LocalityMap()); - logServerMap = (LocalityMap*) logServerSet.getPtr(); - - logEntryMap.clear(); - logIndexArray.clear(); - logIndexArray.reserve(logServers.size()); - - for( int i = 0; i < logServers.size(); i++ ) { - if (logServers[i]->get().present()) { - logIndexArray.push_back(i); - ASSERT(logEntryMap.find(i) == logEntryMap.end()); - logEntryMap[logIndexArray.back()] = logServerMap->add(logServers[i]->get().interf().locality, &logIndexArray.back()); - } - } - } - - void updateLocalitySet( vector const& workers ) { - LocalityMap* logServerMap; - - logServerSet = LocalitySetRef(new LocalityMap()); - logServerMap = (LocalityMap*) logServerSet.getPtr(); - - logEntryMap.clear(); - logIndexArray.clear(); - logIndexArray.reserve(workers.size()); - - for( int i = 0; i < workers.size(); i++ ) { - ASSERT(logEntryMap.find(i) == logEntryMap.end()); - logIndexArray.push_back(i); - logEntryMap[logIndexArray.back()] = logServerMap->add(workers[i].locality, &logIndexArray.back()); - } - } - - void getPushLocations( std::vector const& tags, std::vector& locations, int locationOffset ) { - newLocations.clear(); - alsoServers.clear(); - resultEntries.clear(); - - if(hasBest) { - for(auto& t : tags) { - newLocations.push_back(bestLocationFor(t)); - } - } - - uniquify( newLocations ); - - if (newLocations.size()) - alsoServers.reserve(newLocations.size()); - - // Convert locations to the also servers - for (auto location : newLocations) { - ASSERT(logEntryMap[location]._id == location); - locations.push_back(locationOffset + location); - alsoServers.push_back(logEntryMap[location]); - } - - // Run the policy, assert if unable to satify - bool result = logServerSet->selectReplicas(tLogPolicy, alsoServers, resultEntries); - ASSERT(result); - - // Add the new servers to the location array - LocalityMap* logServerMap = (LocalityMap*) logServerSet.getPtr(); - for (auto entry : resultEntries) { - locations.push_back(locationOffset + *logServerMap->getObject(entry)); - } - //TraceEvent("getPushLocations").detail("Policy", tLogPolicy->info()) - // .detail("Results", locations.size()).detail("Selection", logServerSet->size()) - // .detail("Included", alsoServers.size()).detail("Duration", timer() - t); - } - -private: - std::vector alsoServers, resultEntries; - std::vector newLocations; -}; - struct OldLogData { std::vector tLogs; Version epochEnd; @@ -213,11 +110,11 @@ struct TagPartitionedLogSystem : ILogSystem, ReferenceCounted>>( new AsyncVar>( log ) ) ); } - for( auto & log : it.logRouters) { - logSet.logRouters.push_back( Reference>( new AsyncVar( log ) ) ); + for( auto& log : it.logRouters) { + logSet.logRouters.push_back( Reference>>( new AsyncVar>( log ) ) ); } logSet.tLogWriteAntiQuorum = it.tLogWriteAntiQuorum; logSet.tLogReplicationFactor = it.tLogReplicationFactor; @@ -238,7 +135,7 @@ struct TagPartitionedLogSystem : ILogSystem, ReferenceCounted>>( new AsyncVar>( log ) ) ); } for( auto & log : it.logRouters) { - logSet.logRouters.push_back( Reference>( new AsyncVar( log ) ) ); + logSet.logRouters.push_back( Reference>>( new AsyncVar>( log ) ) ); } logSet.tLogWriteAntiQuorum = it.tLogWriteAntiQuorum; logSet.tLogReplicationFactor = it.tLogReplicationFactor; @@ -268,7 +165,7 @@ struct TagPartitionedLogSystem : ILogSystem, ReferenceCounted>>( new AsyncVar>( log ) ) ); } for( auto & log : it.logRouters) { - logSet.logRouters.push_back( Reference>( new AsyncVar( log ) ) ); + logSet.logRouters.push_back( Reference>>( new AsyncVar>( log ) ) ); } logSet.tLogWriteAntiQuorum = it.tLogWriteAntiQuorum; logSet.tLogReplicationFactor = it.tLogReplicationFactor; @@ -290,7 +187,7 @@ struct TagPartitionedLogSystem : ILogSystem, ReferenceCounted>>( new AsyncVar>( log ) ) ); } for( auto & log : it.logRouters) { - logSet.logRouters.push_back( Reference>( new AsyncVar( log ) ) ); + logSet.logRouters.push_back( Reference>>( new AsyncVar>( log ) ) ); } logSet.tLogWriteAntiQuorum = it.tLogWriteAntiQuorum; logSet.tLogReplicationFactor = it.tLogReplicationFactor; @@ -410,40 +307,41 @@ struct TagPartitionedLogSystem : ILogSystem, ReferenceCounted peek( Version begin, Tag tag, bool parallelGetMore ) { - if(oldLogData.size() == 0 || begin >= oldLogData[0].epochEnd) { - return Reference( new ILogSystem::MergedPeekCursor( logServers, logServers.size() ? bestLocationFor( tag ) : -1, - (int)logServers.size() + 1 - tLogReplicationFactor, tag, begin, getPeekEnd(), parallelGetMore, tLogLocalities, tLogPolicy, tLogReplicationFactor)); + if(tag >= SERVER_KNOBS->MAX_TAG) { + //FIXME: non-static logRouters + return Reference( new ILogSystem::MergedPeekCursor( tLogs[1].logRouters, -1, (int)tLogs[1].logRouters.size(), tag, begin, getPeekEnd(), false ) ); } else { - std::vector< Reference > cursors; - std::vector< LogMessageVersion > epochEnds; - cursors.push_back( Reference( new ILogSystem::MergedPeekCursor( logServers, logServers.size() ? bestLocationFor( tag ) : -1, - (int)logServers.size() + 1 - tLogReplicationFactor, tag, oldLogData[0].epochEnd, getPeekEnd(), parallelGetMore, tLogLocalities, tLogPolicy, tLogReplicationFactor)) ); - for(int i = 0; i < oldLogData.size() && begin < oldLogData[i].epochEnd; i++) { - cursors.push_back( Reference( new ILogSystem::MergedPeekCursor( oldLogData[i].logServers, oldLogData[i].logServers.size() ? oldBestLocationFor( tag, i ) : -1, - (int)oldLogData[i].logServers.size() + 1 - oldLogData[i].tLogReplicationFactor, tag, i+1 == oldLogData.size() ? begin : std::max(oldLogData[i+1].epochEnd, begin), oldLogData[i].epochEnd, parallelGetMore, oldLogData[i].tLogLocalities, oldLogData[i].tLogPolicy, oldLogData[i].tLogReplicationFactor)) ); - epochEnds.push_back(LogMessageVersion(oldLogData[i].epochEnd)); - } + if(oldLogData.size() == 0 || begin >= oldLogData[0].epochEnd) { + return Reference( new ILogSystem::SetPeekCursor( tLogs, 1, tLogs[1].logServers.size() ? tLogs[1].bestLocationFor( tag ) : -1, tag, begin, getPeekEnd(), parallelGetMore ) ); + } else { + std::vector< Reference > cursors; + std::vector< LogMessageVersion > epochEnds; + cursors.push_back( Reference( new ILogSystem::SetPeekCursor( tLogs, 1, tLogs[1].logServers.size() ? tLogs[1].bestLocationFor( tag ) : -1, tag, oldLogData[0].epochEnd, getPeekEnd(), parallelGetMore)) ); + for(int i = 0; i < oldLogData.size() && begin < oldLogData[i].epochEnd; i++) { + cursors.push_back( Reference( new ILogSystem::SetPeekCursor( oldLogData[i].tLogs, 1, oldLogData[i].tLogs[1].logServers.size() ? oldLogData[i].tLogs[1].bestLocationFor( tag ) : -1, tag, i+1 == oldLogData.size() ? begin : std::max(oldLogData[i+1].epochEnd, begin), oldLogData[i].epochEnd, parallelGetMore)) ); + epochEnds.push_back(LogMessageVersion(oldLogData[i].epochEnd)); + } - return Reference( new ILogSystem::MultiCursor(cursors, epochEnds) ); + return Reference( new ILogSystem::MultiCursor(cursors, epochEnds) ); + } } } virtual Reference peekSingle( Version begin, Tag tag ) { + ASSERT(tag < SERVER_KNOBS->MAX_TAG); if(oldLogData.size() == 0 || begin >= oldLogData[0].epochEnd) { - return Reference( new ILogSystem::ServerPeekCursor( logServers.size() ? - logServers[bestLocationFor( tag )] : + return Reference( new ILogSystem::ServerPeekCursor( tLogs[1].logServers.size() ? + tLogs[1].logServers[tLogs[1].bestLocationFor( tag )] : Reference>>(), tag, begin, getPeekEnd(), false, false ) ); } else { TEST(true); //peekSingle used during non-copying tlog recovery std::vector< Reference > cursors; std::vector< LogMessageVersion > epochEnds; - cursors.push_back( Reference( new ILogSystem::ServerPeekCursor( logServers.size() ? - logServers[bestLocationFor( tag )] : + cursors.push_back( Reference( new ILogSystem::ServerPeekCursor( tLogs[1].logServers.size() ? + tLogs[1].logServers[tLogs[1].bestLocationFor( tag )] : Reference>>(), tag, oldLogData[0].epochEnd, getPeekEnd(), false, false) ) ); for(int i = 0; i < oldLogData.size() && begin < oldLogData[i].epochEnd; i++) { - cursors.push_back( Reference( new ILogSystem::MergedPeekCursor( oldLogData[i].logServers, oldLogData[i].logServers.size() ? oldBestLocationFor( tag, i ) : -1, - (int)oldLogData[i].logServers.size() + 1 - oldLogData[i].tLogReplicationFactor, tag, i+1 == oldLogData.size() ? begin : std::max(oldLogData[i+1].epochEnd, begin), oldLogData[i].epochEnd, false, - oldLogData[i].tLogLocalities, oldLogData[i].tLogPolicy, oldLogData[i].tLogReplicationFactor)) ); + cursors.push_back( Reference( new ILogSystem::SetPeekCursor( oldLogData[i].tLogs, 1, oldLogData[i].tLogs[1].logServers.size() ? oldLogData[i].tLogs[1].bestLocationFor( tag ) : -1, tag, i+1 == oldLogData.size() ? begin : std::max(oldLogData[i+1].epochEnd, begin), oldLogData[i].epochEnd, false)) ); epochEnds.push_back(LogMessageVersion(oldLogData[i].epochEnd)); } @@ -509,10 +407,10 @@ struct TagPartitionedLogSystem : ILogSystem, ReferenceCounted> newEpoch( vector availableLogServers, DatabaseConfiguration const& config, LogEpoch recoveryCount ) { + virtual Future> newEpoch( vector availableLogServers, vector availableRemoteLogServers, vector availableLogRouters, DatabaseConfiguration const& config, LogEpoch recoveryCount ) { // Call only after end_epoch() has successfully completed. Returns a new epoch immediately following this one. The new epoch // is only provisional until the caller updates the coordinated DBCoreState - return newEpoch( Reference::addRef(this), availableLogServers, config, recoveryCount ); + return newEpoch( Reference::addRef(this), availableLogServers, availableRemoteLogServers, availableLogRouters, config, recoveryCount ); } virtual LogSystemConfig getLogSystemConfig() { @@ -533,7 +431,7 @@ struct TagPartitionedLogSystem : ILogSystem, ReferenceCountedget().interf()); + log.logRouters.push_back(t.logRouters[i]->get()); } logSystemConfig.tLogs.push_back(log); @@ -557,7 +455,7 @@ struct TagPartitionedLogSystem : ILogSystem, ReferenceCountedget().interf()); + log.logRouters.push_back(t.logRouters[i]->get()); } logSystemConfig.oldTLogs[i].tLogs.push_back(log); diff --git a/fdbserver/WorkerInterface.h b/fdbserver/WorkerInterface.h index 3634b02d69..de022d5843 100644 --- a/fdbserver/WorkerInterface.h +++ b/fdbserver/WorkerInterface.h @@ -304,7 +304,7 @@ Future storageServer( std::string const& folder ); // changes pssi->id() to be the recovered ID Future masterServer( MasterInterface const& mi, Reference> const& db, class ServerCoordinators const&, LifetimeToken const& lifetime ); Future masterProxyServer(MasterProxyInterface const& proxy, InitializeMasterProxyRequest const& req, Reference> const& db); -Future tLog( class IKeyValueStore* const& persistentData, class IDiskQueue* const& persistentQueue, Reference> const& db, LocalityData const& locality, PromiseStream const& tlogRequests, UID const& tlogId, bool const& restoreFromDisk, Promise const& oldLog ); // changes tli->id() to be the recovered ID +Future tLog( class IKeyValueStore* const& persistentData, class IDiskQueue* const& persistentQueue, Reference> const& db, LocalityData const& locality, PromiseStream const& tlogRequests, UID const& tlogId, bool const& restoreFromDisk ); // changes tli->id() to be the recovered ID Future debugQueryServer( DebugQueryRequest const& req ); Future monitorServerDBInfo( Reference>> const& ccInterface, Reference const&, LocalityData const&, Reference> const& dbInfo ); Future resolver( ResolverInterface const& proxy, InitializeResolverRequest const&, Reference> const& db ); diff --git a/fdbserver/fdbserver.vcxproj b/fdbserver/fdbserver.vcxproj index cd76556f17..489c678f6e 100644 --- a/fdbserver/fdbserver.vcxproj +++ b/fdbserver/fdbserver.vcxproj @@ -53,7 +53,6 @@ - diff --git a/fdbserver/fdbserver.vcxproj.filters b/fdbserver/fdbserver.vcxproj.filters index 0a8d61451e..67e43e0cea 100644 --- a/fdbserver/fdbserver.vcxproj.filters +++ b/fdbserver/fdbserver.vcxproj.filters @@ -246,7 +246,6 @@ workloads - workloads diff --git a/fdbserver/masterserver.actor.cpp b/fdbserver/masterserver.actor.cpp index 2440ccff9e..4f3698bebd 100644 --- a/fdbserver/masterserver.actor.cpp +++ b/fdbserver/masterserver.actor.cpp @@ -149,7 +149,7 @@ ACTOR Future writeTransitionMasterState( Reference self, bool self->logSystem->toCoreState( newState ); newState.recoveryCount = self->prevDBState.recoveryCount + 1; - ASSERT( newState.tLogWriteAntiQuorum == self->configuration.tLogWriteAntiQuorum && newState.tLogReplicationFactor == self->configuration.tLogReplicationFactor ); + ASSERT( newState.tLogs[0].tLogWriteAntiQuorum == self->configuration.tLogWriteAntiQuorum && newState.tLogs[0].tLogReplicationFactor == self->configuration.tLogReplicationFactor ); try { Void _ = wait( self->cstate2.setExclusive( BinaryWriter::toValue(newState, IncludeVersion()) ) ); @@ -181,7 +181,7 @@ ACTOR Future writeRecoveredMasterState( Reference self ) { state DBCoreState newState = self->myDBState.get(); self->logSystem->toCoreState( newState ); - ASSERT( newState.tLogWriteAntiQuorum == self->configuration.tLogWriteAntiQuorum && newState.tLogReplicationFactor == self->configuration.tLogReplicationFactor ); + ASSERT( newState.tLogs[0].tLogWriteAntiQuorum == self->configuration.tLogWriteAntiQuorum && newState.tLogs[0].tLogReplicationFactor == self->configuration.tLogReplicationFactor ); try { Void _ = wait( self->cstate3.setExclusive( BinaryWriter::toValue(newState, IncludeVersion()) ) ); @@ -340,13 +340,18 @@ ACTOR Future updateLogsValue( Reference self, Database cx ) { Optional> value = wait( tr.get(logsKey) ); ASSERT(value.present()); - auto logConf = self->logSystem->getLogSystemConfig(); + std::vector> logConf; auto logs = decodeLogsValue(value.get()); + for(auto& log : self->logSystem->getLogSystemConfig().tLogs) { + for(auto& tl : log.tLogs) { + logConf.push_back(tl); + } + } - bool match = (logs.first.size() == logConf.tLogs.size()); + bool match = (logs.first.size() == logConf.size()); if(match) { for(int i = 0; i < logs.first.size(); i++) { - if(logs.first[i].first != logConf.tLogs[i].id()) { + if(logs.first[i].first != logConf[i].id()) { match = false; break; } diff --git a/fdbserver/worker.actor.cpp b/fdbserver/worker.actor.cpp index 1cf0e21103..3784ff7641 100644 --- a/fdbserver/worker.actor.cpp +++ b/fdbserver/worker.actor.cpp @@ -596,13 +596,9 @@ ACTOR Future workerServer( Reference connFile, Refe details["StorageEngine"] = s.storeType.toString(); startRole( s.storeID, interf.id(), "SharedTLog", details, "Restored" ); - Promise oldLog; - Future tl = tLog( kv, queue, dbInfo, locality, tlog.isReady() ? tlogRequests : PromiseStream(), s.storeID, true, oldLog ); + Future tl = tLog( kv, queue, dbInfo, locality, tlog.isReady() ? tlogRequests : PromiseStream(), s.storeID, true ); tl = handleIOErrors( tl, kv, s.storeID ); tl = handleIOErrors( tl, queue, s.storeID ); - if(tlog.isReady()) { - tlog = oldLog.getFuture() || tl; - } errorForwarders.add( forwardError( errors, "SharedTLog", s.storeID, tl ) ); } } @@ -669,7 +665,7 @@ ACTOR Future workerServer( Reference connFile, Refe IDiskQueue* queue = openDiskQueue( joinPath( folder, fileLogQueuePrefix.toString() + logId.toString() + "-" ), logId ); filesClosed.add( data->onClosed() ); filesClosed.add( queue->onClosed() ); - tlog = tLog( data, queue, dbInfo, locality, tlogRequests, logId, false, Promise() ); + tlog = tLog( data, queue, dbInfo, locality, tlogRequests, logId, false ); tlog = handleIOErrors( tlog, data, logId ); tlog = handleIOErrors( tlog, queue, logId ); errorForwarders.add( forwardError( errors, "SharedTLog", logId, tlog ) );