implemented setPeekCursor

removed oldTLogServer
first compiling version
This commit is contained in:
Evan Tschannen 2017-07-10 17:41:32 -07:00
parent 979ebcef6c
commit 81ae263ad9
15 changed files with 510 additions and 1719 deletions

View File

@ -130,16 +130,12 @@ static void applyMetadataMutations(UID const& dbgid, Arena &arena, VectorRef<Mut
if (logSystem && m.param1.startsWith( excludedServersPrefix )) {
// If one of our existing tLogs is now excluded, we have to die and recover
auto addr = decodeExcludedServersKey(m.param1);
for( auto tl : logSystem->getLogSystemConfig().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) {

View File

@ -636,14 +636,15 @@ std::vector<std::pair<WorkerInterface, ProcessClass>> getWorkersForTlogsAcrossDa
if(oldMasterFit < newMasterFit) return false;
//FIXME: implement for remote logs and log routers
std::vector<ProcessClass> 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;

View File

@ -105,66 +105,24 @@ struct DBCoreState {
template <class Archive>
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);
}
};

View File

@ -32,6 +32,109 @@
struct DBCoreState;
template <class Collection>
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<Reference<AsyncVar<OptionalInterface<TLogInterface>>>> logServers;
std::vector<Reference<AsyncVar<OptionalInterface<TLogInterface>>>> logRouters;
int32_t tLogWriteAntiQuorum;
int32_t tLogReplicationFactor;
std::vector< LocalityData > tLogLocalities; // Stores the localities of the log servers
IRepPolicyRef tLogPolicy;
LocalitySetRef logServerSet;
std::vector<int> logIndexArray;
std::map<int,LocalityEntry> 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<int>* logServerMap;
logServerSet = LocalitySetRef(new LocalityMap<int>());
logServerMap = (LocalityMap<int>*) 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<WorkerInterface> const& workers ) {
LocalityMap<int>* logServerMap;
logServerSet = LocalitySetRef(new LocalityMap<int>());
logServerMap = (LocalityMap<int>*) 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<Tag> const& tags, std::vector<int>& 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<int>* logServerMap = (LocalityMap<int>*) 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<LocalityEntry> alsoServers, resultEntries;
std::vector<int> newLocations;
};
struct ILogSystem {
// Represents a particular (possibly provisional) epoch of the log subsystem
@ -165,9 +268,8 @@ struct ILogSystem {
};
struct MergedPeekCursor : IPeekCursor, ReferenceCounted<MergedPeekCursor> {
LocalityGroup localityGroup;
std::vector< std::pair<LogMessageVersion, int> > sortedVersions;
vector< Reference<IPeekCursor> > serverCursors;
std::vector< std::pair<LogMessageVersion, int> > sortedVersions;
Tag tag;
int bestServer, currentCursor, readQuorum;
Optional<LogMessageVersion> nextVersion;
@ -176,11 +278,10 @@ struct ILogSystem {
UID randomID;
int tLogReplicationFactor;
IRepPolicyRef tLogPolicy;
std::vector< LocalityData > tLogLocalities;
MergedPeekCursor( std::vector<Reference<AsyncVar<OptionalInterface<TLogInterface>>>> 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<Reference<AsyncVar<OptionalInterface<TLogInterface>>>> const& logServers, int bestServer, int readQuorum, Tag tag, Version begin, Version end, bool parallelGetMore );
MergedPeekCursor( vector< Reference<IPeekCursor> > const& serverCursors, LogMessageVersion const& messageVersion, int bestServer, int readQuorum, Optional<LogMessageVersion> nextVersion, std::vector< LocalityData > const& tLogLocalities, IRepPolicyRef const tLogPolicy, int tLogReplicationFactor );
MergedPeekCursor( vector< Reference<IPeekCursor> > const& serverCursors, LogMessageVersion const& messageVersion, int bestServer, int readQuorum, Optional<LogMessageVersion> 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<SetPeekCursor> {
std::vector<LogSet> logSets;
std::vector< std::vector< Reference<IPeekCursor> > > serverCursors;
Tag tag;
int bestSet, bestServer, currentSet, currentCursor;
LocalityGroup localityGroup;
std::vector< std::pair<LogMessageVersion, int> > sortedVersions;
Optional<LogMessageVersion> nextVersion;
LogMessageVersion messageVersion;
bool hasNextMessage;
bool useBestSet;
UID randomID;
SetPeekCursor( std::vector<LogSet> const& logSets, int bestSet, int bestServer, Tag tag, Version begin, Version end, bool parallelGetMore );
virtual Reference<IPeekCursor> 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<Tag> getTags();
virtual void advanceTo(LogMessageVersion n);
virtual Future<Void> getMore();
virtual Future<Void> onFailed();
virtual bool isActive();
virtual LogMessageVersion version();
virtual Version popped();
virtual void addref() {
ReferenceCounted<SetPeekCursor>::addref();
}
virtual void delref() {
ReferenceCounted<SetPeekCursor>::delref();
}
};
struct MultiCursor : IPeekCursor, ReferenceCounted<MultiCursor> {
std::vector<Reference<IPeekCursor>> cursors;
std::vector<LogMessageVersion> epochEnds;

View File

@ -57,7 +57,7 @@ protected:
struct TLogSet {
std::vector<OptionalInterface<TLogInterface>> tLogs;
std::vector<TLogInterface> logRouters;
std::vector<OptionalInterface<TLogInterface>> logRouters;
int32_t tLogWriteAntiQuorum, tLogReplicationFactor;
std::vector< LocalityData > tLogLocalities; // Stores the localities of the log servers
IRepPolicyRef tLogPolicy;

View File

@ -237,19 +237,17 @@ LogMessageVersion ILogSystem::ServerPeekCursor::version() { return messageVersio
Version ILogSystem::ServerPeekCursor::popped() { return poppedVersion; }
ILogSystem::MergedPeekCursor::MergedPeekCursor( std::vector<Reference<AsyncVar<OptionalInterface<TLogInterface>>>> 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<Reference<AsyncVar<OptionalInterface<TLogInterface>>>> 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<ILogSystem::ServerPeekCursor> 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<ILogSystem::IPeekCursor> > const& serverCursors, LogMessageVersion const& messageVersion, int bestServer, int readQuorum, Optional<LogMessageVersion> 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<ILogSystem::IPeekCursor> > const& serverCursors, LogMessageVersion const& messageVersion, int bestServer, int readQuorum, Optional<LogMessageVersion> 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::IPeekCursor> ILogSystem::MergedPeekCursor::cloneNoMore() {
for( auto it : serverCursors ) {
cursors.push_back(it->cloneNoMore());
}
return Reference<ILogSystem::MergedPeekCursor>( new ILogSystem::MergedPeekCursor( cursors, messageVersion, bestServer, readQuorum, nextVersion, tLogLocalities, tLogPolicy, tLogReplicationFactor ) );
return Reference<ILogSystem::MergedPeekCursor>( 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<LogMessageVersion, int>(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<LogSet> 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<ILogSystem::ServerPeekCursor> cursor( new ILogSystem::ServerPeekCursor( logSets[i].logServers[j], tag, begin, end, true, parallelGetMore ) );
serverCursors[i].push_back( cursor );
}
maxServers = std::max<int>(maxServers, serverCursors[i].size());
}
sortedVersions.resize(maxServers);
}
Reference<ILogSystem::IPeekCursor> 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<LogMessageVersion, int>(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<Tag> 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<Void> 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<Future<Void>> 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<Future<Void>> 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<Void> 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<Void> 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<Reference<IPeekCursor>> cursors, std::vector<LogMessageVersion> epochEnds ) : cursors(cursors), epochEnds(epochEnds), poppedVersion(0) {}
Reference<ILogSystem::IPeekCursor> ILogSystem::MultiCursor::cloneNoMore() {

File diff suppressed because it is too large Load Diff

View File

@ -1447,26 +1447,34 @@ static StatusArray oldTlogFetcher(int* oldLogFaultTolerance, Reference<AsyncVar<
if(db->get().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);
}
}

View File

@ -213,7 +213,6 @@ struct TLogData : NonCopyable {
WorkerCache<TLogInterface> tlogCache;
Future<Void> updatePersist; //SOMEDAY: integrate the recovery and update storage so that only one of them is committing to persistant data.
Future<Void> oldLogServer;
PromiseStream<Future<Void>> 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<Void> 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 ) {
logData->addActor = PromiseStream<Future<Void>>(); //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<Void> checkEmptyQueue(TLogData* self) {
}
}
ACTOR Future<Void> restorePersistentState( TLogData* self, LocalityData locality, Promise<Void> oldLog, PromiseStream<InitializeTLogRequest> tlogRequests ) {
ACTOR Future<Void> restorePersistentState( TLogData* self, LocalityData locality, PromiseStream<InitializeTLogRequest> tlogRequests ) {
state double startt = now();
state Reference<LogData> logData;
state KeyRange tagKeys;
@ -1455,28 +1454,7 @@ ACTOR Future<Void> restorePersistentState( TLogData* self, LocalityData locality
state std::vector<Future<ErrorOr<Void>>> 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<Void> tLogStart( TLogData* self, InitializeTLogRequest req, Localit
}
// New tLog (if !recoverFrom.size()) or restore from network
ACTOR Future<Void> tLog( IKeyValueStore* persistentData, IDiskQueue* persistentQueue, Reference<AsyncVar<ServerDBInfo>> db, LocalityData locality, PromiseStream<InitializeTLogRequest> tlogRequests, UID tlogId, bool restoreFromDisk, Promise<Void> oldLog )
ACTOR Future<Void> tLog( IKeyValueStore* persistentData, IDiskQueue* persistentQueue, Reference<AsyncVar<ServerDBInfo>> db, LocalityData locality, PromiseStream<InitializeTLogRequest> tlogRequests, UID tlogId, bool restoreFromDisk )
{
state TLogData self( tlogId, persistentData, persistentQueue, db );
state Future<Void> error = actorCollection( self.sharedActors.getFuture() );
@ -1866,7 +1844,7 @@ ACTOR Future<Void> 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) );
}

View File

@ -29,12 +29,6 @@
#include "fdbrpc/Replication.h"
#include "fdbrpc/ReplicationUtils.h"
template <class Collection>
void uniquify( Collection& c ) {
std::sort(c.begin(), c.end());
c.resize( std::unique(c.begin(), c.end()) - c.begin() );
}
ACTOR static Future<Void> reportTLogCommitErrors( Future<Void> commitReply, UID debugID ) {
try {
Void _ = wait(commitReply);
@ -48,103 +42,6 @@ ACTOR static Future<Void> reportTLogCommitErrors( Future<Void> commitReply, UID
}
}
class LogSet {
public:
std::vector<Reference<AsyncVar<OptionalInterface<TLogInterface>>>> logServers;
std::vector<Reference<AsyncVar<OptionalInterface<TLogInterface>>>> logRouters;
int32_t tLogWriteAntiQuorum;
int32_t tLogReplicationFactor;
std::vector< LocalityData > tLogLocalities; // Stores the localities of the log servers
IRepPolicyRef tLogPolicy;
LocalitySetRef logServerSet;
std::vector<int> logIndexArray;
std::map<int,LocalityEntry> 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<int>* logServerMap;
logServerSet = LocalitySetRef(new LocalityMap<int>());
logServerMap = (LocalityMap<int>*) 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<WorkerInterface> const& workers ) {
LocalityMap<int>* logServerMap;
logServerSet = LocalitySetRef(new LocalityMap<int>());
logServerMap = (LocalityMap<int>*) 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<Tag> const& tags, std::vector<int>& 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<int>* logServerMap = (LocalityMap<int>*) 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<LocalityEntry> alsoServers, resultEntries;
std::vector<int> newLocations;
};
struct OldLogData {
std::vector<LogSet> tLogs;
Version epochEnd;
@ -213,11 +110,11 @@ struct TagPartitionedLogSystem : ILogSystem, ReferenceCounted<TagPartitionedLogS
for( auto& it : lsConf.tLogs ) {
LogSet logSet;
for( auto & log : it.tLogs) {
for( auto& log : it.tLogs) {
logSet.logServers.push_back( Reference<AsyncVar<OptionalInterface<TLogInterface>>>( new AsyncVar<OptionalInterface<TLogInterface>>( log ) ) );
}
for( auto & log : it.logRouters) {
logSet.logRouters.push_back( Reference<AsyncVar<TLogInterface>>( new AsyncVar<TLogInterface>( log ) ) );
for( auto& log : it.logRouters) {
logSet.logRouters.push_back( Reference<AsyncVar<OptionalInterface<TLogInterface>>>( new AsyncVar<OptionalInterface<TLogInterface>>( log ) ) );
}
logSet.tLogWriteAntiQuorum = it.tLogWriteAntiQuorum;
logSet.tLogReplicationFactor = it.tLogReplicationFactor;
@ -238,7 +135,7 @@ struct TagPartitionedLogSystem : ILogSystem, ReferenceCounted<TagPartitionedLogS
logSet.logServers.push_back( Reference<AsyncVar<OptionalInterface<TLogInterface>>>( new AsyncVar<OptionalInterface<TLogInterface>>( log ) ) );
}
for( auto & log : it.logRouters) {
logSet.logRouters.push_back( Reference<AsyncVar<TLogInterface>>( new AsyncVar<TLogInterface>( log ) ) );
logSet.logRouters.push_back( Reference<AsyncVar<OptionalInterface<TLogInterface>>>( new AsyncVar<OptionalInterface<TLogInterface>>( log ) ) );
}
logSet.tLogWriteAntiQuorum = it.tLogWriteAntiQuorum;
logSet.tLogReplicationFactor = it.tLogReplicationFactor;
@ -268,7 +165,7 @@ struct TagPartitionedLogSystem : ILogSystem, ReferenceCounted<TagPartitionedLogS
logSet.logServers.push_back( Reference<AsyncVar<OptionalInterface<TLogInterface>>>( new AsyncVar<OptionalInterface<TLogInterface>>( log ) ) );
}
for( auto & log : it.logRouters) {
logSet.logRouters.push_back( Reference<AsyncVar<TLogInterface>>( new AsyncVar<TLogInterface>( log ) ) );
logSet.logRouters.push_back( Reference<AsyncVar<OptionalInterface<TLogInterface>>>( new AsyncVar<OptionalInterface<TLogInterface>>( log ) ) );
}
logSet.tLogWriteAntiQuorum = it.tLogWriteAntiQuorum;
logSet.tLogReplicationFactor = it.tLogReplicationFactor;
@ -290,7 +187,7 @@ struct TagPartitionedLogSystem : ILogSystem, ReferenceCounted<TagPartitionedLogS
logSet.logServers.push_back( Reference<AsyncVar<OptionalInterface<TLogInterface>>>( new AsyncVar<OptionalInterface<TLogInterface>>( log ) ) );
}
for( auto & log : it.logRouters) {
logSet.logRouters.push_back( Reference<AsyncVar<TLogInterface>>( new AsyncVar<TLogInterface>( log ) ) );
logSet.logRouters.push_back( Reference<AsyncVar<OptionalInterface<TLogInterface>>>( new AsyncVar<OptionalInterface<TLogInterface>>( log ) ) );
}
logSet.tLogWriteAntiQuorum = it.tLogWriteAntiQuorum;
logSet.tLogReplicationFactor = it.tLogReplicationFactor;
@ -410,40 +307,41 @@ struct TagPartitionedLogSystem : ILogSystem, ReferenceCounted<TagPartitionedLogS
}
virtual Reference<IPeekCursor> peek( Version begin, Tag tag, bool parallelGetMore ) {
if(oldLogData.size() == 0 || begin >= oldLogData[0].epochEnd) {
return Reference<ILogSystem::MergedPeekCursor>( 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<ILogSystem::MergedPeekCursor>( new ILogSystem::MergedPeekCursor( tLogs[1].logRouters, -1, (int)tLogs[1].logRouters.size(), tag, begin, getPeekEnd(), false ) );
} else {
std::vector< Reference<ILogSystem::IPeekCursor> > cursors;
std::vector< LogMessageVersion > epochEnds;
cursors.push_back( Reference<ILogSystem::MergedPeekCursor>( 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<ILogSystem::MergedPeekCursor>( 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<ILogSystem::SetPeekCursor>( new ILogSystem::SetPeekCursor( tLogs, 1, tLogs[1].logServers.size() ? tLogs[1].bestLocationFor( tag ) : -1, tag, begin, getPeekEnd(), parallelGetMore ) );
} else {
std::vector< Reference<ILogSystem::IPeekCursor> > cursors;
std::vector< LogMessageVersion > epochEnds;
cursors.push_back( Reference<ILogSystem::SetPeekCursor>( 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<ILogSystem::SetPeekCursor>( 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<ILogSystem::MultiCursor>( new ILogSystem::MultiCursor(cursors, epochEnds) );
return Reference<ILogSystem::MultiCursor>( new ILogSystem::MultiCursor(cursors, epochEnds) );
}
}
}
virtual Reference<IPeekCursor> peekSingle( Version begin, Tag tag ) {
ASSERT(tag < SERVER_KNOBS->MAX_TAG);
if(oldLogData.size() == 0 || begin >= oldLogData[0].epochEnd) {
return Reference<ILogSystem::ServerPeekCursor>( new ILogSystem::ServerPeekCursor( logServers.size() ?
logServers[bestLocationFor( tag )] :
return Reference<ILogSystem::ServerPeekCursor>( new ILogSystem::ServerPeekCursor( tLogs[1].logServers.size() ?
tLogs[1].logServers[tLogs[1].bestLocationFor( tag )] :
Reference<AsyncVar<OptionalInterface<TLogInterface>>>(), tag, begin, getPeekEnd(), false, false ) );
} else {
TEST(true); //peekSingle used during non-copying tlog recovery
std::vector< Reference<ILogSystem::IPeekCursor> > cursors;
std::vector< LogMessageVersion > epochEnds;
cursors.push_back( Reference<ILogSystem::ServerPeekCursor>( new ILogSystem::ServerPeekCursor( logServers.size() ?
logServers[bestLocationFor( tag )] :
cursors.push_back( Reference<ILogSystem::ServerPeekCursor>( new ILogSystem::ServerPeekCursor( tLogs[1].logServers.size() ?
tLogs[1].logServers[tLogs[1].bestLocationFor( tag )] :
Reference<AsyncVar<OptionalInterface<TLogInterface>>>(), tag, oldLogData[0].epochEnd, getPeekEnd(), false, false) ) );
for(int i = 0; i < oldLogData.size() && begin < oldLogData[i].epochEnd; i++) {
cursors.push_back( Reference<ILogSystem::MergedPeekCursor>( 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<ILogSystem::SetPeekCursor>( 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<TagPartitionedLogS
return waitForAll(quorumResults);
}
virtual Future<Reference<ILogSystem>> newEpoch( vector<WorkerInterface> availableLogServers, DatabaseConfiguration const& config, LogEpoch recoveryCount ) {
virtual Future<Reference<ILogSystem>> newEpoch( vector<WorkerInterface> availableLogServers, vector<WorkerInterface> availableRemoteLogServers, vector<WorkerInterface> 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<TagPartitionedLogSystem>::addRef(this), availableLogServers, config, recoveryCount );
return newEpoch( Reference<TagPartitionedLogSystem>::addRef(this), availableLogServers, availableRemoteLogServers, availableLogRouters, config, recoveryCount );
}
virtual LogSystemConfig getLogSystemConfig() {
@ -533,7 +431,7 @@ struct TagPartitionedLogSystem : ILogSystem, ReferenceCounted<TagPartitionedLogS
}
for( int i = 0; i < t.logRouters.size(); i++ ) {
log.logRouters.push_back(t.logRouters[i]->get().interf());
log.logRouters.push_back(t.logRouters[i]->get());
}
logSystemConfig.tLogs.push_back(log);
@ -557,7 +455,7 @@ struct TagPartitionedLogSystem : ILogSystem, ReferenceCounted<TagPartitionedLogS
}
for( int i = 0; i < t.logRouters.size(); i++ ) {
log.logRouters.push_back(t.logRouters[i]->get().interf());
log.logRouters.push_back(t.logRouters[i]->get());
}
logSystemConfig.oldTLogs[i].tLogs.push_back(log);

View File

@ -304,7 +304,7 @@ Future<Void> storageServer(
std::string const& folder ); // changes pssi->id() to be the recovered ID
Future<Void> masterServer( MasterInterface const& mi, Reference<AsyncVar<ServerDBInfo>> const& db, class ServerCoordinators const&, LifetimeToken const& lifetime );
Future<Void> masterProxyServer(MasterProxyInterface const& proxy, InitializeMasterProxyRequest const& req, Reference<AsyncVar<ServerDBInfo>> const& db);
Future<Void> tLog( class IKeyValueStore* const& persistentData, class IDiskQueue* const& persistentQueue, Reference<AsyncVar<ServerDBInfo>> const& db, LocalityData const& locality, PromiseStream<InitializeTLogRequest> const& tlogRequests, UID const& tlogId, bool const& restoreFromDisk, Promise<Void> const& oldLog ); // changes tli->id() to be the recovered ID
Future<Void> tLog( class IKeyValueStore* const& persistentData, class IDiskQueue* const& persistentQueue, Reference<AsyncVar<ServerDBInfo>> const& db, LocalityData const& locality, PromiseStream<InitializeTLogRequest> const& tlogRequests, UID const& tlogId, bool const& restoreFromDisk ); // changes tli->id() to be the recovered ID
Future<Void> debugQueryServer( DebugQueryRequest const& req );
Future<Void> monitorServerDBInfo( Reference<AsyncVar<Optional<ClusterControllerFullInterface>>> const& ccInterface, Reference<ClusterConnectionFile> const&, LocalityData const&, Reference<AsyncVar<ServerDBInfo>> const& dbInfo );
Future<Void> resolver( ResolverInterface const& proxy, InitializeResolverRequest const&, Reference<AsyncVar<ServerDBInfo>> const& db );

View File

@ -53,7 +53,6 @@
<ActorCompiler Include="Resolver.actor.cpp" />
<ActorCompiler Include="LogSystemDiskQueueAdapter.actor.cpp" />
<ActorCompiler Include="LogSystemPeekCursor.actor.cpp" />
<ActorCompiler Include="OldTLogServer.actor.cpp" />
<ActorCompiler Include="LogRouter.actor.cpp" />
<ClCompile Include="SkipList.cpp" />
<ActorCompiler Include="WaitFailure.actor.cpp" />

View File

@ -246,7 +246,6 @@
<ActorCompiler Include="workloads\AtomicRestore.actor.cpp">
<Filter>workloads</Filter>
</ActorCompiler>
<ActorCompiler Include="OldTLogServer.actor.cpp" />
<ActorCompiler Include="LogRouter.actor.cpp" />
<ActorCompiler Include="workloads\SlowTaskWorkload.actor.cpp">
<Filter>workloads</Filter>

View File

@ -149,7 +149,7 @@ ACTOR Future<Void> writeTransitionMasterState( Reference<MasterData> 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<Void> writeRecoveredMasterState( Reference<MasterData> 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<Void> updateLogsValue( Reference<MasterData> self, Database cx ) {
Optional<Standalone<StringRef>> value = wait( tr.get(logsKey) );
ASSERT(value.present());
auto logConf = self->logSystem->getLogSystemConfig();
std::vector<OptionalInterface<TLogInterface>> 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;
}

View File

@ -596,13 +596,9 @@ ACTOR Future<Void> workerServer( Reference<ClusterConnectionFile> connFile, Refe
details["StorageEngine"] = s.storeType.toString();
startRole( s.storeID, interf.id(), "SharedTLog", details, "Restored" );
Promise<Void> oldLog;
Future<Void> tl = tLog( kv, queue, dbInfo, locality, tlog.isReady() ? tlogRequests : PromiseStream<InitializeTLogRequest>(), s.storeID, true, oldLog );
Future<Void> tl = tLog( kv, queue, dbInfo, locality, tlog.isReady() ? tlogRequests : PromiseStream<InitializeTLogRequest>(), 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<Void> workerServer( Reference<ClusterConnectionFile> 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<Void>() );
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 ) );