From 23062b892e555a3a6b31870771ba56ab4f0722f1 Mon Sep 17 00:00:00 2001 From: Dan Lambright Date: Fri, 15 Oct 2021 16:05:18 -0400 Subject: [PATCH 1/2] Calculate tpcv on resolvers --- fdbserver/CommitProxyServer.actor.cpp | 54 +++------------------------ fdbserver/MasterInterface.h | 6 +-- fdbserver/Resolver.actor.cpp | 32 ++++++++++++++++ fdbserver/ResolverInterface.h | 8 ++++ fdbserver/TLogServer.actor.cpp | 3 -- fdbserver/masterserver.actor.cpp | 41 -------------------- 6 files changed, 48 insertions(+), 96 deletions(-) diff --git a/fdbserver/CommitProxyServer.actor.cpp b/fdbserver/CommitProxyServer.actor.cpp index 6d29158b31..fcbf85ba7b 100644 --- a/fdbserver/CommitProxyServer.actor.cpp +++ b/fdbserver/CommitProxyServer.actor.cpp @@ -536,7 +536,7 @@ struct CommitBatchContext { double commitStartTime; - std::unordered_map tpcvMap; // obtained from sequencer + std::unordered_map tpcvMap; // obtained from resolver std::set writtenTLogs; // the set of tlog locations written to in the mutation. std::set writtenTags; // final set tags written to in the batch std::set writtenTagsPreResolution; // tags written to in the batch not including any changes from the resolver. @@ -733,12 +733,6 @@ ACTOR Future preresolutionProcessing(CommitBatchContext* self) { if (SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST) { self->writtenTagsPreResolution = self->getWrittenTagsPreResolution(); - if (self->hasMetadataMutation) { - int numLogs = pProxyCommitData->db->get().logSystemConfig.numLogs(); - for (int i = 0; i < numLogs; i++) { - self->writtenTLogs.insert(i); - } - } } GetCommitVersionRequest req(span.context, pProxyCommitData->commitVersionRequestNumber++, @@ -749,10 +743,6 @@ ACTOR Future preresolutionProcessing(CommitBatchContext* self) { GetCommitVersionReply versionReply = wait(brokenPromiseToNever( pProxyCommitData->master.getCommitVersion.getReply(req, TaskPriority::ProxyMasterVersionReply))); - if (SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST) { - self->tpcvMap = versionReply.tpcvMap; - } - pProxyCommitData->mostRecentProcessedRequestNumber = versionReply.requestNum; pProxyCommitData->stats.txnCommitVersionAssigned += trs.size(); @@ -807,6 +797,7 @@ ACTOR Future getResolution(CommitBatchContext* self) { std::vector> replies; for (int r = 0; r < pProxyCommitData->resolvers.size(); r++) { requests.requests[r].debugID = self->debugID; + requests.requests[r].writtenTags = self->writtenTagsPreResolution; replies.push_back(trackResolutionMetrics(pProxyCommitData->stats.resolverDist[r], brokenPromiseToNever(pProxyCommitData->resolvers[r].resolve.getReply( requests.requests[r], TaskPriority::ProxyResolverReply)))); @@ -994,6 +985,10 @@ ACTOR Future applyMetadataToCommittedTransactions(CommitBatchContext* self // in the same set of transactions as this proxy. ResolveTransactionBatchReply& reply = self->resolution[0]; self->toCommit.setMutations(reply.privateMutationCount, reply.privateMutations); + if (SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST) { + // TraceEvent("ResolverReturn").detail("ReturnTags",reply.writtenTags).detail("TPCVsize",reply.tpcvMap.size()).detail("ReqTags",self->writtenTagsPreResolution); + self->tpcvMap = reply.tpcvMap; + } } self->lockedKey = pProxyCommitData->txnStateStore->readValue(databaseLockedKey).get(); @@ -1018,16 +1013,6 @@ ACTOR Future applyMetadataToCommittedTransactions(CommitBatchContext* self return Void(); } -ACTOR Future getTPCV(CommitBatchContext* self, std::set writtenTLogs) { - state ProxyCommitData* const pProxyCommitData = self->pProxyCommitData; - GetTLogPrevCommitVersionReply rep = - wait(brokenPromiseToNever(pProxyCommitData->master.getTLogPrevCommitVersion.getReply( - GetTLogPrevCommitVersionRequest(writtenTLogs, self->commitVersion, self->prevVersion)))); - // TraceEvent("GetTLogPrevCommitVersionRequest"); - self->tpcvMap.insert(rep.tpcvMap.begin(), rep.tpcvMap.end()); - return Void(); -} - /// This second pass through committed transactions assigns the actual mutations to the appropriate storage servers' /// tags ACTOR Future assignMutationsToStorageServers(CommitBatchContext* self) { @@ -1242,14 +1227,6 @@ ACTOR Future postResolution(CommitBatchContext* self) { wait(Future(Never())); } - if (SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST && self->metadataMutationFromProxy && - !self->hasMetadataMutation) { - // TraceEvent("Abort metadataMutationFromProxy"); - for (int transactionNum = 0; transactionNum < trs.size(); transactionNum++) { - self->committed[transactionNum] = ConflictBatch::TransactionConflict; - } - } - // First pass wait(applyMetadataToCommittedTransactions(self)); @@ -1258,25 +1235,6 @@ ACTOR Future postResolution(CommitBatchContext* self) { self->toCommit.saveTags(self->writtenTags); - if (self->writtenTags.size() && !self->metadataMutationFromProxy && - SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST) { - // confirm all serialized tags are sent to a tLog for which the previous commit version was obtained. - std::set postResolutionTLogs; - self->toCommit.getLocations(self->writtenTags, postResolutionTLogs); - for (auto& t : postResolutionTLogs) { - if (self->writtenTLogs.find(t) == self->writtenTLogs.end()) { - TraceEvent(SevError, "TagHasNoPCV", pProxyCommitData->dbgid) - .detail("tagsBeforeResolution", self->writtenTagsPreResolution) - .detail("tagsAfterResolution", self->writtenTags) - .detail("numTLogsBeforeResolution", self->writtenTLogs.size()) - .detail("numTLogsAfterResolution", postResolutionTLogs.size()) - .detail("hasMetadataMutation", self->hasMetadataMutation) - .detail("metadataMutationFromProxy", self->metadataMutationFromProxy); - ASSERT(false); - } - } - } - // Serialize and backup the mutations as a single mutation if ((pProxyCommitData->vecBackupKeys.size() > 1) && self->logRangeMutations.size()) { wait(addBackupMutations(pProxyCommitData, diff --git a/fdbserver/MasterInterface.h b/fdbserver/MasterInterface.h index 21772295b7..f3956df044 100644 --- a/fdbserver/MasterInterface.h +++ b/fdbserver/MasterInterface.h @@ -156,7 +156,6 @@ struct GetCommitVersionReply { Version version; Version prevVersion; uint64_t requestNum; - std::unordered_map tpcvMap; GetCommitVersionReply() : resolverChangesVersion(0), version(0), prevVersion(0), requestNum(0) {} explicit GetCommitVersionReply(Version version, Version prevVersion, uint64_t requestNum) @@ -164,7 +163,7 @@ struct GetCommitVersionReply { template void serialize(Ar& ar) { - serializer(ar, resolverChanges, resolverChangesVersion, version, prevVersion, requestNum, tpcvMap); + serializer(ar, resolverChanges, resolverChangesVersion, version, prevVersion, requestNum); } }; @@ -194,11 +193,10 @@ struct GetCommitVersionRequest { struct GetTLogPrevCommitVersionReply { constexpr static FileIdentifier file_identifier = 16683183; - std::unordered_map tpcvMap; GetTLogPrevCommitVersionReply() {} template void serialize(Ar& ar) { - serializer(ar, tpcvMap); + serializer(ar); } }; diff --git a/fdbserver/Resolver.actor.cpp b/fdbserver/Resolver.actor.cpp index 4cab849cbd..39cf745784 100644 --- a/fdbserver/Resolver.actor.cpp +++ b/fdbserver/Resolver.actor.cpp @@ -78,6 +78,9 @@ struct Resolver : ReferenceCounted { Version debugMinRecentStateVersion = 0; + // The previous commit versions per tlog + std::vector tpcvVector; + CounterCollection cc; Counter resolveBatchIn; Counter resolveBatchStart; @@ -94,6 +97,7 @@ struct Resolver : ReferenceCounted { Counter resolveBatchOut; Counter metricsRequests; Counter splitRequests; + int numLogs; Future logger; @@ -192,6 +196,7 @@ ACTOR Future resolveBatch(Reference self, ResolveTransactionBatc g_traceBatch.addEvent("CommitDebug", debugID.get().first(), "Resolver.resolveBatch.AfterOrderer"); ResolveTransactionBatchReply& reply = proxyInfo.outstandingBatches[req.version]; + reply.writtenTags = req.writtenTags; std::vector commitList; std::vector tooOldList; @@ -281,6 +286,10 @@ ACTOR Future resolveBatch(Reference self, ResolveTransactionBatc reply.privateMutations.push_back(reply.arena, mutations); reply.arena.dependsOn(mutations.arena()); } + if (SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST) { + // merge mutation tags with sent client tags + toCommit.saveTags(reply.writtenTags); + } reply.privateMutationCount = toCommit.getMutationCount(); } @@ -340,6 +349,23 @@ ACTOR Future resolveBatch(Reference self, ResolveTransactionBatc } } + if (SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST) { + std::set writtenTLogs; + if (reply.privateMutationCount) { + for (int i = 0; i < self->numLogs; i++) { + writtenTLogs.insert(i); + } + } else { + toCommit.getLocations(reply.writtenTags, writtenTLogs); + } + if (self->tpcvVector[0] == invalidVersion) { + std::fill(self->tpcvVector.begin(), self->tpcvVector.end(), req.prevVersion); + } + for (uint16_t tLog : writtenTLogs) { + reply.tpcvMap[tLog] = self->tpcvVector[tLog]; + self->tpcvVector[tLog] = req.version; + } + } self->version.set(req.version); bool breachedLimit = self->totalStateBytes.get() <= SERVER_KNOBS->RESOLVER_STATE_MEMORY_LIMIT && self->totalStateBytes.get() + stateBytes > SERVER_KNOBS->RESOLVER_STATE_MEMORY_LIMIT; @@ -568,6 +594,12 @@ ACTOR Future resolverCore(ResolverInterface resolver, // This has to be declared after the self->txnStateStore get initialized transactionStateResolveContext = TransactionStateResolveContext(self, &addActor); + + if (SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST) { + self->numLogs = db->get().logSystemConfig.numLogs(); + self->tpcvVector.resize(1 + self->numLogs, 0); + std::fill(self->tpcvVector.begin(), self->tpcvVector.end(), invalidVersion); + } } loop choose { diff --git a/fdbserver/ResolverInterface.h b/fdbserver/ResolverInterface.h index 31501f0205..da41a62298 100644 --- a/fdbserver/ResolverInterface.h +++ b/fdbserver/ResolverInterface.h @@ -96,6 +96,9 @@ struct ResolveTransactionBatchReply { VectorRef privateMutations; uint32_t privateMutationCount; + std::unordered_map tpcvMap; + std::set writtenTags; + template void serialize(Archive& ar) { serializer(ar, @@ -105,6 +108,8 @@ struct ResolveTransactionBatchReply { conflictingKeyRangeMap, privateMutations, privateMutationCount, + tpcvMap, + writtenTags, arena); } }; @@ -123,6 +128,8 @@ struct ResolveTransactionBatchRequest { ReplyPromise reply; Optional debugID; + std::set writtenTags; + template void serialize(Archive& ar) { serializer(ar, @@ -134,6 +141,7 @@ struct ResolveTransactionBatchRequest { reply, arena, debugID, + writtenTags, spanContext); } }; diff --git a/fdbserver/TLogServer.actor.cpp b/fdbserver/TLogServer.actor.cpp index f1fe9b55f7..3a6d691d04 100644 --- a/fdbserver/TLogServer.actor.cpp +++ b/fdbserver/TLogServer.actor.cpp @@ -2148,9 +2148,7 @@ ACTOR Future tLogCommit(TLogData* self, } logData->minKnownCommittedVersion = std::max(logData->minKnownCommittedVersion, req.minKnownCommittedVersion); - wait(logData->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())) { wait(delay(0, g_network->getCurrentTask())); @@ -2198,7 +2196,6 @@ ACTOR Future tLogCommit(TLogData* self, 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 logData->version.set(req.version); diff --git a/fdbserver/masterserver.actor.cpp b/fdbserver/masterserver.actor.cpp index a0afb4e694..a5481be7ca 100644 --- a/fdbserver/masterserver.actor.cpp +++ b/fdbserver/masterserver.actor.cpp @@ -250,8 +250,6 @@ struct MasterData : NonCopyable, ReferenceCounted { // up-to-date in the presence of key range splits/merges. VersionVector ssVersionVector; - // The previous commit versions per tlog - std::vector tpcvVector; CounterCollection cc; Counter changeCoordinatorsRequests; Counter getCommitVersionRequests; @@ -1188,9 +1186,6 @@ ACTOR Future getVersion(Reference self, GetCommitVersionReques self->lastVersionTime = now(); self->version = self->recoveryTransactionVersion; rep.prevVersion = self->lastEpochEnd; - if (SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST) { - std::fill(self->tpcvVector.begin(), self->tpcvVector.end(), self->lastEpochEnd); - } } else { double t1 = now(); @@ -1227,13 +1222,6 @@ ACTOR Future getVersion(Reference self, GetCommitVersionReques proxyItr->second.replies[req.requestNum] = rep; ASSERT(rep.prevVersion >= 0); - if (SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST) { - for (uint16_t tLog : req.writtenTLogs) { - rep.tpcvMap[tLog] = self->tpcvVector[tLog]; - self->tpcvVector[tLog] = rep.version; - } - } - req.reply.send(rep); ASSERT(proxyItr->second.latestRequestNum.get() == req.requestNum - 1); @@ -1281,24 +1269,6 @@ ACTOR Future waitForPrev(Reference self, ReportRawCommittedVer return Void(); } -ACTOR Future waitForTLogPrev(Reference self, GetTLogPrevCommitVersionRequest req) { - // TraceEvent("WaitForTLogPrev").detail("Prev",req.prev).detail("CommitVersion",req.commitVersion).detail("PrevTLogVersion",self->prevTLogVersion.get()); - if (self->prevTLogVersion.get() != invalidVersion) { - wait(self->prevTLogVersion.whenAtLeast(req.prev)); - } else { - std::fill(self->tpcvVector.begin(), self->tpcvVector.end(), self->lastEpochEnd); - } - - GetTLogPrevCommitVersionReply reply; - for (uint16_t tLog : req.writtenTLogs) { - reply.tpcvMap[tLog] = self->tpcvVector[tLog]; - self->tpcvVector[tLog] = req.commitVersion; - } - self->prevTLogVersion.set(req.commitVersion); - req.reply.send(reply); - return Void(); -} - ACTOR Future serveLiveCommittedVersion(Reference self) { loop { choose { @@ -1332,10 +1302,6 @@ ACTOR Future serveLiveCommittedVersion(Reference self) { req.reply.send(Void()); } } - when(GetTLogPrevCommitVersionRequest req = - waitNext(self->myInterface.getTLogPrevCommitVersion.getFuture())) { - self->addActor.send(waitForTLogPrev(self, req)); - } } } } @@ -1973,10 +1939,6 @@ ACTOR Future masterCore(Reference self) { tr.read_snapshot = self->recoveryTransactionVersion; // lastEpochEnd would make more sense, but isn't in the initial // window of the resolver(s) - // resize the TPCV vector to the number of tlogs and initialize to first transaction - int numLogs = self->dbInfo->get().logSystemConfig.numLogs(); - self->tpcvVector.resize(1 + numLogs, 0); - TraceEvent("MasterRecoveryCommit", self->dbgid).log(); state Future> recoveryCommit = self->commitProxies[0].commit.tryGetReply(recoveryCommitRequest); self->addActor.send(self->logSystem->onError()); @@ -2116,9 +2078,6 @@ ACTOR Future masterServer(MasterInterface mi, wait(delay(5)); throw worker_removed(); } - // resize the TPCV vector to the number of tlogs and initialize to first transaction - int numLogs = self->dbInfo->get().logSystemConfig.numLogs(); - self->tpcvVector.resize(1 + numLogs, 0); } when(BackupWorkerDoneRequest req = waitNext(mi.notifyBackupWorkerDone.getFuture())) { if (self->logSystem.isValid() && self->logSystem->removeBackupWorker(req)) { From eb814ce07021f6b4785bcb69d29b2839efe208a8 Mon Sep 17 00:00:00 2001 From: Dan Lambright Date: Mon, 18 Oct 2021 10:23:08 -0400 Subject: [PATCH 2/2] Remove dead code, check replies match for all resolvers --- fdbclient/ServerKnobs.cpp | 2 +- fdbserver/CommitProxyServer.actor.cpp | 20 ++++---------------- fdbserver/MasterInterface.h | 27 +++------------------------ 3 files changed, 8 insertions(+), 41 deletions(-) diff --git a/fdbclient/ServerKnobs.cpp b/fdbclient/ServerKnobs.cpp index b4c0ef707c..314b1e73a0 100644 --- a/fdbclient/ServerKnobs.cpp +++ b/fdbclient/ServerKnobs.cpp @@ -37,7 +37,7 @@ void ServerKnobs::initialize(Randomize randomize, ClientKnobs* clientKnobs, IsSi init( MAX_WRITE_TRANSACTION_LIFE_VERSIONS, 5 * VERSIONS_PER_SECOND ); if (randomize && BUGGIFY) MAX_WRITE_TRANSACTION_LIFE_VERSIONS=std::max(1, 1 * VERSIONS_PER_SECOND); init( MAX_COMMIT_BATCH_INTERVAL, 2.0 ); if( randomize && BUGGIFY ) MAX_COMMIT_BATCH_INTERVAL = 0.5; // Each commit proxy generates a CommitTransactionBatchRequest at least this often, so that versions always advance smoothly MAX_COMMIT_BATCH_INTERVAL = std::min(MAX_COMMIT_BATCH_INTERVAL, MAX_READ_TRANSACTION_LIFE_VERSIONS/double(2*VERSIONS_PER_SECOND)); // Ensure that the proxy commits 2 times every MAX_READ_TRANSACTION_LIFE_VERSIONS, otherwise the master will not give out versions fast enough - init( ENABLE_VERSION_VECTOR, true ); + init( ENABLE_VERSION_VECTOR, true ); init( ENABLE_VERSION_VECTOR_TLOG_UNICAST, false ); // TLogs diff --git a/fdbserver/CommitProxyServer.actor.cpp b/fdbserver/CommitProxyServer.actor.cpp index fcbf85ba7b..4b80b287cf 100644 --- a/fdbserver/CommitProxyServer.actor.cpp +++ b/fdbserver/CommitProxyServer.actor.cpp @@ -537,11 +537,8 @@ struct CommitBatchContext { double commitStartTime; std::unordered_map tpcvMap; // obtained from resolver - std::set writtenTLogs; // the set of tlog locations written to in the mutation. std::set writtenTags; // final set tags written to in the batch std::set writtenTagsPreResolution; // tags written to in the batch not including any changes from the resolver. - bool hasMetadataMutation = false; - bool metadataMutationFromProxy = false; CommitBatchContext(ProxyCommitData*, const std::vector*, const int); @@ -564,9 +561,7 @@ std::set CommitBatchContext::getWrittenTagsPreResolution() { if (isSingleKeyMutation((MutationRef::Type)m.type)) { auto& tags = pProxyCommitData->tagsForKey(m.param1); transactionTags.insert(tags.begin(), tags.end()); - toCommit.getLocations(tags, writtenTLogs); if (pProxyCommitData->cacheInfo[m.param1]) { - toCommit.getLocations(cacheVector, writtenTLogs); transactionTags.insert(cacheTag); } } else if (m.type == MutationRef::ClearRange) { @@ -579,7 +574,6 @@ std::set CommitBatchContext::getWrittenTagsPreResolution() { ranges.begin().value().populateTags(); filteredTags.insert(ranges.begin().value().tags.begin(), ranges.begin().value().tags.end()); transactionTags.insert(ranges.begin().value().tags.begin(), ranges.begin().value().tags.end()); - toCommit.getLocations(filteredTags, writtenTLogs); } else { std::set allSources; for (auto r : ranges) { @@ -587,17 +581,14 @@ std::set CommitBatchContext::getWrittenTagsPreResolution() { allSources.insert(r.value().tags.begin(), r.value().tags.end()); transactionTags.insert(r.value().tags.begin(), r.value().tags.end()); } - toCommit.getLocations(allSources, writtenTLogs); } if (pProxyCommitData->needsCacheTag(clearRange)) { - toCommit.getLocations(cacheVector, writtenTLogs); transactionTags.insert(cacheTag); } } else { UNREACHABLE(); } } - hasMetadataMutation = containsMetadataMutation(trs[transactionNum].transaction.mutations); } return transactionTags; @@ -737,8 +728,7 @@ ACTOR Future preresolutionProcessing(CommitBatchContext* self) { GetCommitVersionRequest req(span.context, pProxyCommitData->commitVersionRequestNumber++, pProxyCommitData->mostRecentProcessedRequestNumber, - pProxyCommitData->dbgid, - self->writtenTLogs); + pProxyCommitData->dbgid); state double beforeGettingCommitVersion = now(); GetCommitVersionReply versionReply = wait(brokenPromiseToNever( pProxyCommitData->master.getCommitVersion.getReply(req, TaskPriority::ProxyMasterVersionReply))); @@ -839,6 +829,9 @@ void assertResolutionStateMutationsSizeConsistent(const std::vectorENABLE_VERSION_VECTOR_TLOG_UNICAST) { + ASSERT_EQ(resolution[0].tpcvMap.size(), resolution[r].tpcvMap.size()); + } for (int s = 0; s < resolution[r].stateMutations.size(); s++) { ASSERT(resolution[r].stateMutations[s].size() == resolution[0].stateMutations[s].size()); } @@ -872,11 +865,6 @@ void applyMetadataEffect(CommitBatchContext* self) { self->forceRecovery, /* popVersion= */ 0, /* initialCommit */ false); - if (SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST && - containsMetadataMutation( - self->resolution[0].stateMutations[versionIndex][transactionIndex].mutations)) { - self->metadataMutationFromProxy = true; - } } if (self->resolution[0].stateMutations[versionIndex][transactionIndex].mutations.size() && self->firstStateMutations) { diff --git a/fdbserver/MasterInterface.h b/fdbserver/MasterInterface.h index f3956df044..069d2c3a17 100644 --- a/fdbserver/MasterInterface.h +++ b/fdbserver/MasterInterface.h @@ -44,7 +44,6 @@ struct MasterInterface { RequestStream getLiveCommittedVersion; // Report a proxy's committed version. RequestStream reportLiveCommittedVersion; - RequestStream getTLogPrevCommitVersion; NetworkAddress address() const { return changeCoordinators.getEndpoint().getPrimaryAddress(); } NetworkAddressList addresses() const { return changeCoordinators.getEndpoint().addresses; } @@ -68,8 +67,6 @@ struct MasterInterface { RequestStream(waitFailure.getEndpoint().getAdjustedEndpoint(5)); reportLiveCommittedVersion = RequestStream( waitFailure.getEndpoint().getAdjustedEndpoint(6)); - getTLogPrevCommitVersion = - RequestStream(waitFailure.getEndpoint().getAdjustedEndpoint(7)); } } @@ -82,7 +79,6 @@ struct MasterInterface { streams.push_back(notifyBackupWorkerDone.getReceiver()); streams.push_back(getLiveCommittedVersion.getReceiver(TaskPriority::GetLiveCommittedVersion)); streams.push_back(reportLiveCommittedVersion.getReceiver(TaskPriority::ReportLiveCommittedVersion)); - streams.push_back(getTLogPrevCommitVersion.getReceiver(TaskPriority::GetTLogPrevCommitVersion)); FlowTransport::transport().addEndpoints(streams); } }; @@ -173,21 +169,19 @@ struct GetCommitVersionRequest { uint64_t requestNum; uint64_t mostRecentProcessedRequestNum; UID requestingProxy; - std::set writtenTLogs; ReplyPromise reply; GetCommitVersionRequest() {} GetCommitVersionRequest(SpanID spanContext, uint64_t requestNum, uint64_t mostRecentProcessedRequestNum, - UID requestingProxy, - std::set& writtenTLogs) + UID requestingProxy) : spanContext(spanContext), requestNum(requestNum), mostRecentProcessedRequestNum(mostRecentProcessedRequestNum), - requestingProxy(requestingProxy), writtenTLogs(writtenTLogs) {} + requestingProxy(requestingProxy) {} template void serialize(Ar& ar) { - serializer(ar, requestNum, mostRecentProcessedRequestNum, requestingProxy, writtenTLogs, reply, spanContext); + serializer(ar, requestNum, mostRecentProcessedRequestNum, requestingProxy, reply, spanContext); } }; @@ -200,21 +194,6 @@ struct GetTLogPrevCommitVersionReply { } }; -struct GetTLogPrevCommitVersionRequest { - constexpr static FileIdentifier file_identifier = 16683184; - std::set writtenTLogs; - Version commitVersion; - Version prev; - ReplyPromise reply; - GetTLogPrevCommitVersionRequest() {} - GetTLogPrevCommitVersionRequest(std::set& writtenTLogs, Version commitVersion, Version prev) - : writtenTLogs(writtenTLogs), commitVersion(commitVersion), prev(prev) {} - template - void serialize(Ar& ar) { - serializer(ar, writtenTLogs, commitVersion, prev, reply); - } -}; - struct ReportRawCommittedVersionRequest { constexpr static FileIdentifier file_identifier = 1853148; Version version;