foundationdb/fdbserver/commitproxy/CommitProxyServer.cpp

3276 lines
134 KiB
C++

/*
* CommitProxyServer.cpp
*
* This source file is part of the FoundationDB open source project
*
* Copyright 2013-2026 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 <algorithm>
#include <string_view>
#include <tuple>
#include <variant>
#include "fdbclient/AccumulativeChecksum.h"
#include "fdbclient/Atomic.h"
#include "fdbclient/BackupAgent.h"
#include "fdbclient/BuildIdempotencyIdMutations.h"
#include "fdbclient/CommitTransaction.h"
#include "fdbclient/DatabaseContext.h"
#include "fdbclient/FDBTypes.h"
#include "fdbclient/IdempotencyId.h"
#include "fdbclient/Knobs.h"
#include "fdbclient/CommitProxyInterface.h"
#include "fdbclient/NativeAPI.actor.h"
#include "fdbclient/SystemData.h"
#include "fdbclient/Tracing.h"
#include "fdbclient/TransactionLineage.h"
#include "fdbrpc/sim_validation.h"
#include "fdbserver/core/AccumulativeChecksumUtil.h"
#include "fdbserver/logsystem/ApplyMetadataMutation.h"
#include "fdbserver/core/ConflictBatch.h"
#include "fdbserver/core/DataDistributorInterface.h"
#include "fdbserver/kvstore/IKeyValueStore.h"
#include "fdbserver/core/Knobs.h"
#include "fdbserver/logsystem/LogSystem.h"
#include "fdbserver/kvstore/FDBExecHelper.h"
#include "fdbserver/logsystem/LogSystemFactory.h"
#include "fdbserver/logsystem/LogSystemDiskQueueAdapter.h"
#include "fdbserver/core/MasterInterface.h"
#include "fdbserver/core/MutationTracking.h"
#include "ProxyCommitData.h"
#include "fdbserver/core/RatekeeperInterface.h"
#include "fdbserver/core/RecoveryState.h"
#include "fdbserver/core/ServerDBInfo.h"
#include "fdbserver/core/WaitFailure.h"
#include "fdbserver/commitproxy/CommitProxyServer.h"
#include "fdbserver/core/WorkerInterface.h"
#include "flow/ActorCollection.h"
#include "flow/CodeProbe.h"
#include "flow/CoroUtils.h"
#include "flow/EncryptUtils.h"
#include "flow/Error.h"
#include "flow/IRandom.h"
#include "flow/Knobs.h"
#include "flow/Trace.h"
#include "flow/UnitTest.h"
#include "flow/network.h"
using WriteMutationRefVar = std::variant<MutationRef, VectorRef<MutationRef>>;
struct ResolutionRequestBuilder {
const ProxyCommitData* self;
// One request per resolver.
std::vector<ResolveTransactionBatchRequest> requests;
// Txn i to resolvers that have i'th data sent
std::vector<std::vector<int>> transactionResolverMap;
std::vector<CommitTransactionRef*> outTr;
// Used to report conflicting keys, the format is
// [CommitTransactionRef_Index][Resolver_Index][Read_Conflict_Range_Index_on_Resolver]
// -> read_conflict_range's original index in the commitTransactionRef
std::vector<std::vector<std::vector<int>>> txReadConflictRangeIndexMap;
ResolutionRequestBuilder(ProxyCommitData* self,
Version version,
Version prevVersion,
Version lastReceivedVersion,
Version lastShardMove,
Span& parentSpan)
: self(self), requests(self->resolvers.size()) {
for (auto& req : requests) {
req.spanContext = parentSpan.context;
req.prevVersion = prevVersion;
req.version = version;
req.lastReceivedVersion = lastReceivedVersion;
req.lastShardMove = lastShardMove;
}
}
CommitTransactionRef& getOutTransaction(int resolver, Version read_snapshot) {
CommitTransactionRef*& out = outTr[resolver];
if (!out) {
ResolveTransactionBatchRequest& request = requests[resolver];
request.transactions.resize(request.arena, request.transactions.size() + 1);
out = &request.transactions.back();
out->read_snapshot = read_snapshot;
}
return *out;
}
std::vector<int> getResolversForRange(const KeyRangeRef& range, Optional<Version> readSnapshot) const {
std::vector<int> resolvers;
resolvers.reserve(self->resolvers.size());
std::vector<unsigned char> seen(self->resolvers.size(), 0);
for (auto& intersectingRange : self->keyResolvers.intersectingRanges(range)) {
auto& versionResolvers = intersectingRange.value();
if (readSnapshot.present()) {
for (int i = versionResolvers.size() - 1; i >= 0; --i) {
const int resolver = versionResolvers[i].second;
if (!seen[resolver]) {
seen[resolver] = 1;
resolvers.push_back(resolver);
}
if (versionResolvers[i].first < readSnapshot.get()) {
break;
}
}
} else if (!versionResolvers.empty()) {
const int resolver = versionResolvers.back().second;
if (!seen[resolver]) {
seen[resolver] = 1;
resolvers.push_back(resolver);
}
}
}
if (SERVER_KNOBS->PROXY_USE_RESOLVER_PRIVATE_MUTATIONS && systemKeys.intersects(range)) {
resolvers.clear();
for (int resolver = 0; resolver < self->resolvers.size(); ++resolver) {
resolvers.push_back(resolver);
}
}
ASSERT(!resolvers.empty());
return resolvers;
}
// Returns a read conflict index map: [resolver_index][read_conflict_range_index_on_the_resolver]
// -> read_conflict_range's original index
std::vector<std::vector<int>> addReadConflictRanges(CommitTransactionRef& trIn) {
std::vector<std::vector<int>> rCRIndexMap(requests.size());
for (int idx = 0; idx < trIn.read_conflict_ranges.size(); ++idx) {
const auto& r = trIn.read_conflict_ranges[idx];
for (int resolver : getResolversForRange(r, trIn.read_snapshot)) {
getOutTransaction(resolver, trIn.read_snapshot)
.read_conflict_ranges.push_back(requests[resolver].arena, r);
rCRIndexMap[resolver].push_back(idx);
}
}
return rCRIndexMap;
}
void addWriteConflictRanges(CommitTransactionRef& trIn) {
for (auto& r : trIn.write_conflict_ranges) {
for (int resolver : getResolversForRange(r, Optional<Version>())) {
getOutTransaction(resolver, trIn.read_snapshot)
.write_conflict_ranges.push_back(requests[resolver].arena, r);
}
}
}
void addTransaction(CommitTransactionRequest& trRequest, Version ver, int transactionNumberInBatch) {
auto& trIn = trRequest.transaction;
// SOMEDAY: There are a couple of unnecessary O( # resolvers ) steps here
outTr.assign(requests.size(), nullptr);
ASSERT(transactionNumberInBatch >= 0 && transactionNumberInBatch < 32768);
bool isTXNStateTransaction = false;
for (auto& m : trIn.mutations) {
DEBUG_MUTATION("AddTr", ver, m, self->dbgid).detail("Idx", transactionNumberInBatch);
if (m.type == MutationRef::SetVersionstampedKey) {
transformVersionstampMutation(m, &MutationRef::param1, requests[0].version, transactionNumberInBatch);
trIn.write_conflict_ranges.push_back(requests[0].arena, singleKeyRange(m.param1, requests[0].arena));
} else if (m.type == MutationRef::SetVersionstampedValue) {
transformVersionstampMutation(m, &MutationRef::param2, requests[0].version, transactionNumberInBatch);
}
if (isMetadataMutation(m)) {
isTXNStateTransaction = true;
auto& tr = getOutTransaction(0, trIn.read_snapshot);
tr.mutations.push_back(requests[0].arena, m);
tr.lock_aware = trRequest.isLockAware();
}
}
if (isTXNStateTransaction && !trRequest.isLockAware()) {
// This mitigates https://github.com/apple/foundationdb/issues/3647. Since this transaction is not lock
// aware, if this transaction got a read version then \xff/dbLocked must not have been set at this
// transaction's read snapshot. If that changes by commit time, then it won't commit on any proxy because of
// a conflict. A client could set a read version manually so this isn't totally bulletproof.
trIn.read_conflict_ranges.push_back(trRequest.arena, KeyRangeRef(databaseLockedKey, databaseLockedKeyEnd));
}
std::vector<std::vector<int>> rCRIndexMap = addReadConflictRanges(trIn);
txReadConflictRangeIndexMap.push_back(std::move(rCRIndexMap));
addWriteConflictRanges(trIn);
if (isTXNStateTransaction) {
for (int r = 0; r < requests.size(); r++) {
int transactionNumberInRequest =
&getOutTransaction(r, trIn.read_snapshot) - requests[r].transactions.begin();
requests[r].txnStateTransactions.push_back(requests[r].arena, transactionNumberInRequest);
}
// Note only Resolver 0 got the correct spanContext, which means
// the reply from Resolver 0 has the right one back.
auto& tr = getOutTransaction(0, trIn.read_snapshot);
tr.spanContext = trRequest.spanContext;
}
std::vector<int> resolversUsed;
for (int r = 0; r < outTr.size(); r++) {
if (outTr[r]) {
resolversUsed.push_back(r);
outTr[r]->report_conflicting_keys = trIn.report_conflicting_keys;
}
}
transactionResolverMap.emplace_back(std::move(resolversUsed));
}
};
Future<Void> commitBatcher(ProxyCommitData* commitData,
PromiseStream<std::pair<std::vector<CommitTransactionRequest>, int>> out,
FutureStream<CommitTransactionRequest> in,
int desiredBytes,
int64_t memBytesLimit) {
co_await delayJittered(commitData->commitBatchInterval, TaskPriority::ProxyCommitBatcher);
double lastBatch = 0;
while (true) {
Future<Void> timeout;
std::vector<CommitTransactionRequest> batch;
int batchBytes = 0;
auto flushBatch = [&](ProxyStats::CommitBatchFlushReason reason) {
commitData->stats.recordCommitBatchFlush(reason);
out.send({ std::move(batch), batchBytes });
lastBatch = now();
};
// TODO: Enable this assertion (currently failing with gcc)
// static_assert(std::is_nothrow_move_constructible_v<CommitTransactionRequest>);
if (SERVER_KNOBS->MAX_COMMIT_BATCH_INTERVAL <= 0) {
timeout = Never();
} else {
timeout = delayJittered(SERVER_KNOBS->MAX_COMMIT_BATCH_INTERVAL, TaskPriority::ProxyCommitBatcher);
}
while (!timeout.isReady() &&
!(batch.size() == SERVER_KNOBS->COMMIT_TRANSACTION_BATCH_COUNT_MAX || batchBytes >= desiredBytes)) {
auto res = co_await race(in, timeout, commitData->triggerCommit.onChange());
if (res.index() == 0) {
CommitTransactionRequest req = std::get<0>(std::move(res));
// WARNING: this code is run at a high priority, so it needs to do as little work as possible
int bytes = getBytes(req);
// Drop requests if memory is under severe pressure
if (commitData->commitBatchesMemBytesCount + bytes > memBytesLimit) {
++commitData->stats.txnCommitErrors;
req.reply.sendError(commit_proxy_memory_limit_exceeded());
TraceEvent(SevWarnAlways, "ProxyCommitBatchMemoryThresholdExceeded")
.suppressFor(60)
.detail("MemBytesCount", commitData->commitBatchesMemBytesCount)
.detail("MemLimit", memBytesLimit);
continue;
}
commitData->stats.transactionSizeDist->sample(bytes);
if (bytes > FLOW_KNOBS->PACKET_WARNING) {
TraceEvent(SevWarn, "LargeTransaction")
.suppressFor(1.0)
.detail("Size", bytes)
.detail("Client", req.reply.getEndpoint().getPrimaryAddress());
}
++commitData->stats.txnCommitIn;
commitData->stats.uniqueClients.insert(req.reply.getEndpoint().getPrimaryAddress());
if (req.debugID.present()) {
g_traceBatch.addEvent("CommitDebug", req.debugID.get().first(), "CommitProxyServer.batcher");
}
if (batch.empty()) {
if (now() - lastBatch > commitData->commitBatchInterval) {
timeout = delayJittered(SERVER_KNOBS->COMMIT_TRANSACTION_BATCH_INTERVAL_FROM_IDLE,
TaskPriority::ProxyCommitBatcher);
} else {
timeout = delayJittered(commitData->commitBatchInterval - (now() - lastBatch),
TaskPriority::ProxyCommitBatcher);
}
}
if ((batchBytes + bytes > CLIENT_KNOBS->TRANSACTION_SIZE_LIMIT || req.firstInBatch()) &&
!batch.empty()) {
auto reason = batchBytes + bytes > CLIENT_KNOBS->TRANSACTION_SIZE_LIMIT
? ProxyStats::CommitBatchFlushReason::TRANSACTION_SIZE_LIMIT
: ProxyStats::CommitBatchFlushReason::FIRST_IN_BATCH;
commitData->triggerCommit.set(false);
flushBatch(reason);
timeout = delayJittered(commitData->commitBatchInterval, TaskPriority::ProxyCommitBatcher);
batch.clear();
batchBytes = 0;
}
batch.push_back(req);
batchBytes += bytes;
commitData->commitBatchesMemBytesCount += bytes;
} else if (res.index() == 1) {
} else if (res.index() == 2) {
ASSERT(commitData->triggerCommit.get());
double commitTime = lastBatch + SERVER_KNOBS->COMMIT_TRIGGER_DELAY;
if (now() > commitTime) {
break;
}
timeout = timeout || delayJittered(commitTime - now(), TaskPriority::ProxyCommitBatcher);
} else {
UNREACHABLE();
}
}
auto reason = batch.size() == SERVER_KNOBS->COMMIT_TRANSACTION_BATCH_COUNT_MAX
? ProxyStats::CommitBatchFlushReason::COUNT_LIMIT
: batchBytes >= desiredBytes ? ProxyStats::CommitBatchFlushReason::BYTE_LIMIT
: ProxyStats::CommitBatchFlushReason::TIMEOUT;
commitData->triggerCommit.set(false);
flushBatch(reason);
}
}
void createWhitelistBinPathVec(const std::string& binPath, std::vector<Standalone<StringRef>>& binPathVec) {
TraceEvent(SevDebug, "BinPathConverter").detail("Input", binPath);
StringRef input(binPath);
while (!input.empty()) {
StringRef token = input.eat(","_sr);
if (!token.empty()) {
const uint8_t* ptr = token.begin();
while (ptr != token.end() && *ptr == ' ') {
ptr++;
}
if (ptr != token.end()) {
Standalone<StringRef> newElement(token.substr(ptr - token.begin()));
TraceEvent(SevDebug, "BinPathItem").detail("Element", newElement);
binPathVec.push_back(newElement);
}
}
}
return;
}
bool isWhitelisted(const std::vector<Standalone<StringRef>>& binPathVec, StringRef binPath) {
TraceEvent("BinPath").detail("Value", binPath);
for (const auto& item : binPathVec) {
TraceEvent("Element").detail("Value", item);
}
return std::find(binPathVec.begin(), binPathVec.end(), binPath) != binPathVec.end();
}
Future<Void> addBackupMutations(ProxyCommitData* self,
const std::map<Key, MutationListRef>* logRangeMutations,
LogPushData* toCommit,
Version commitVersion,
double* computeDuration,
double* computeStart) {
const int32_t version = commitVersion / CLIENT_KNOBS->LOG_RANGE_BLOCK_SIZE;
int yieldBytes = 0;
toCommit->addTransactionInfo(SpanContext());
// Serialize the log range mutations within the map
for (auto logRangeMutation = logRangeMutations->cbegin(); logRangeMutation != logRangeMutations->cend();
++logRangeMutation) {
// FIXME: this is re-implementing the serialize function of MutationListRef in order to have a yield
// this is 0x0FDB00A200090001
BinaryWriter valueWriter(IncludeVersion(ProtocolVersion::withBackupMutations()));
valueWriter << logRangeMutation->second.totalSize(); // this is int32 by default
MutationListRef::Blob* blobIter = logRangeMutation->second.blob_begin;
while (blobIter) {
if (yieldBytes > SERVER_KNOBS->DESIRED_TOTAL_BYTES) {
yieldBytes = 0;
if (g_network->check_yield(TaskPriority::ProxyCommitYield1)) {
*computeDuration += g_network->timer_monotonic() - *computeStart;
co_await delay(0, TaskPriority::ProxyCommitYield1);
*computeStart = g_network->timer_monotonic();
}
}
valueWriter.serializeBytes(blobIter->data);
yieldBytes += blobIter->data.size();
blobIter = blobIter->next;
}
Key val = valueWriter.toValue();
BinaryWriter wr(Unversioned()); // backupName/hash/commitVersion/part, so wr is param1
// Serialize the log destination
wr.serializeBytes(logRangeMutation->first);
// Write the log keys and version information
wr << (uint8_t)hashlittle(&version, sizeof(version), 0);
wr << bigEndian64(commitVersion);
uint32_t* partBuffer = nullptr;
for (int part = 0; part * CLIENT_KNOBS->MUTATION_BLOCK_SIZE < val.size(); part++) {
MutationRef backupMutation;
backupMutation.type = MutationRef::SetValue;
// Assign the second parameter as the part
// Define the mutation type and and location
backupMutation.param2 = getBackupValue(val, part);
Key key = getBackupKey(wr, &partBuffer, part); // holds the memory for backupMutation
backupMutation.param1 = key;
ASSERT(backupMutation.param1.startsWith(
logRangeMutation->first)); // We are writing into the configured destination
auto& tags = self->tagsForKey(backupMutation.param1);
toCommit->addTags(tags);
if (self->acsBuilder != nullptr) {
updateMutationWithAcsAndAddMutationToAcsBuilder(
self->acsBuilder,
backupMutation,
tags,
getCommitProxyAccumulativeChecksumIndex(self->commitProxyIndex),
self->epoch,
commitVersion,
self->dbgid);
}
toCommit->writeTypedMessage(backupMutation);
// if (DEBUG_MUTATION("BackupProxyCommit", commitVersion, backupMutation)) {
// TraceEvent("BackupProxyCommitTo", self->dbgid).detail("To",
// describe(tags)).detail("BackupMutation", backupMutation.toString())
// .detail("BackupMutationSize", val.size()).detail("Version", commitVersion).detail("DestPath",
// logRangeMutation.first) .detail("PartIndex", part).detail("PartIndexEndian",
// bigEndian32(part)).detail("PartData", backupMutation.param1);
// }
}
}
}
static Future<Void> releaseResolvingAfterImpl(ProxyCommitData* self,
Future<Void> releaseDelay,
int64_t localBatchNumber) {
co_await releaseDelay;
ASSERT(self->latestLocalCommitBatchResolving.get() == localBatchNumber - 1);
self->latestLocalCommitBatchResolving.set(localBatchNumber);
co_return;
}
Future<Void> releaseResolvingAfter(ProxyCommitData* self, Future<Void> releaseDelay, int64_t localBatchNumber) {
if (releaseDelay.isReady()) {
releaseDelay.get();
ASSERT(self->latestLocalCommitBatchResolving.get() == localBatchNumber - 1);
self->latestLocalCommitBatchResolving.set(localBatchNumber);
return Void();
}
return releaseResolvingAfterImpl(self, releaseDelay, localBatchNumber);
}
static Future<ResolveTransactionBatchReply> trackResolutionMetrics(Reference<Histogram> dist,
Future<ResolveTransactionBatchReply> in) {
double startTime = g_network->timer_monotonic();
ResolveTransactionBatchReply reply = co_await in;
dist->sampleSeconds(g_network->timer_monotonic() - startTime);
co_return reply;
}
namespace CommitBatch {
constexpr const std::string_view UNSET = std::string_view();
constexpr const std::string_view INITIALIZE = "initialize"sv;
constexpr const std::string_view PRE_RESOLUTION = "preResolution"sv;
constexpr const std::string_view RESOLUTION = "resolution"sv;
constexpr const std::string_view POST_RESOLUTION = "postResolution"sv;
constexpr const std::string_view TRANSACTION_LOGGING = "transactionLogging"sv;
constexpr const std::string_view REPLY = "reply"sv;
constexpr const std::string_view COMPLETE = "complete"sv;
struct CommitBatchContext {
using StoreCommit_t = std::vector<std::pair<Future<LogSystemDiskQueueAdapter::CommitMessage>, Future<Void>>>;
ProxyCommitData* const pProxyCommitData;
std::vector<CommitTransactionRequest> trs;
const int currentBatchMemBytesCount;
double startTime;
// The current stage of batch commit
std::string_view stage = UNSET;
Optional<UID> debugID;
bool forceRecovery = false;
bool rejected = false; // If rejected due to long queue length
int64_t localBatchNumber;
LogPushData toCommit;
int batchOperations = 0;
Span span;
int64_t batchBytes = 0;
int latencyBucket = 0;
Version commitVersion;
Version prevVersion;
int64_t maxTransactionBytes;
std::vector<std::vector<int>> transactionResolverMap;
std::vector<std::vector<std::vector<int>>> txReadConflictRangeIndexMap;
Future<Void> releaseDelay;
Future<Void> releaseFuture;
std::vector<ResolveTransactionBatchReply> resolution;
double computeStart;
double computeDuration = 0;
Arena arena;
/// true if the batch is the 1st batch for this proxy, additional metadata
/// processing is involved for this batch.
bool isMyFirstBatch;
bool firstStateMutations;
Optional<Value> previousCoordinators;
StoreCommit_t storeCommits;
std::vector<uint8_t> committed;
Optional<Key> lockedKey;
bool locked;
int commitCount = 0;
std::vector<int> nextTr;
bool lockedAfter;
Optional<Value> metadataVersionAfter;
int mutationCount = 0;
int mutationBytes = 0;
std::map<Key, MutationListRef> logRangeMutations;
Arena logRangeMutationsArena;
int transactionNum = 0;
int yieldBytes = 0;
LogSystemDiskQueueAdapter::CommitMessage msg;
Future<Version> loggingComplete;
double commitStartTime;
std::unordered_map<uint16_t, Version> tpcvMap; // obtained from resolver
std::set<Tag> writtenTags; // final set tags written to in the batch
std::set<Tag> writtenTagsPreResolution; // tags written to in the batch not including any changes from the resolver.
IdempotencyIdKVBuilder idempotencyKVBuilder;
CommitBatchContext(ProxyCommitData*, const std::vector<CommitTransactionRequest>*, const int);
void setupTraceBatch();
std::set<Tag> getWrittenTagsPreResolution();
void checkHotShards();
bool rangeLockEnabled();
Version lastShardMove;
private:
void evaluateBatchSize();
};
bool CommitBatchContext::rangeLockEnabled() {
return pProxyCommitData->rangeLockEnabled();
}
void CommitBatchContext::checkHotShards() {
// removed expired hot shards
for (auto it = pProxyCommitData->hotShards.begin(); it != pProxyCommitData->hotShards.end();) {
if (now() > it->second) {
it = pProxyCommitData->hotShards.erase(it);
} else {
++it;
}
}
if (pProxyCommitData->hotShards.empty()) {
return;
}
auto trsBegin = trs.begin();
std::vector<size_t> transactionsToRemove;
for (int transactionNum = 0; transactionNum < trs.size(); transactionNum++) {
VectorRef<MutationRef>* pMutations = &trs[transactionNum].transaction.mutations;
bool abortTransaction = false;
for (int mutationNum = 0; mutationNum < pMutations->size(); mutationNum++) {
auto& m = (*pMutations)[mutationNum];
if (isSingleKeyMutation((MutationRef::Type)m.type)) {
for (const auto& shard : pProxyCommitData->hotShards) {
if (shard.first.contains(KeyRef(m.param1))) {
abortTransaction = true;
break;
}
}
} else if (m.type == MutationRef::ClearRange) {
for (const auto& shard : pProxyCommitData->hotShards) {
if (shard.first.intersects(KeyRangeRef(m.param1, m.param2))) {
abortTransaction = true;
break;
}
}
} else {
UNREACHABLE();
}
}
if (abortTransaction) {
trs[transactionNum].reply.sendError(transaction_throttled_hot_shard());
transactionsToRemove.push_back(transactionNum);
}
}
// Remove transactions marked for removal in reverse order to avoid shifting indices
for (auto it = transactionsToRemove.rbegin(); it != transactionsToRemove.rend(); ++it) {
trs.erase(trsBegin + *it);
}
committed.resize(trs.size());
return;
}
// Check whether the mutation intersects any legal backup ranges
// If so, it will be clamped to the intersecting range(s) later
inline bool shouldBackup(MutationRef const& m) {
if (normalKeys.contains(m.param1) || m.param1 == metadataVersionKey) {
return true;
} else if (m.type != MutationRef::Type::ClearRange) {
return systemBackupMutationMask().rangeContaining(m.param1).value();
} else {
for (auto& r : systemBackupMutationMask().intersectingRanges(KeyRangeRef(m.param1, m.param2))) {
if (r->value()) {
return true;
}
}
}
return false;
}
// Find the set of logs the batch is sent to. An empty set indicates it cannot be
// determined. In version vector, this means the batch should be sent to all logs.
std::set<Tag> CommitBatchContext::getWrittenTagsPreResolution() {
std::set<Tag> transactionTags;
lastShardMove = pProxyCommitData->lastShardMove;
if (pProxyCommitData->txnStateStore->getReplaceContent()) {
return std::set<Tag>();
}
if (!pProxyCommitData->idempotencyClears.empty()) {
return std::set<Tag>();
}
for (int transactionNum = 0; transactionNum < trs.size(); transactionNum++) {
int mutationNum = 0;
VectorRef<MutationRef>* pMutations = &trs[transactionNum].transaction.mutations;
if (trs[transactionNum].idempotencyId.valid()) {
return std::set<Tag>();
}
for (; mutationNum < pMutations->size(); mutationNum++) {
auto& m = (*pMutations)[mutationNum];
// disable version vector's effect if any mutation in the batch is backed up.
// TODO: make backup work with version vector.
if (pProxyCommitData->vecBackupKeys.size() > 1 && shouldBackup(m)) {
return std::set<Tag>();
}
if (isSingleKeyMutation((MutationRef::Type)m.type)) {
auto& tags = pProxyCommitData->tagsForKey(m.param1);
transactionTags.insert(tags.begin(), tags.end());
if (!pProxyCommitData->cdcRouting.empty()) {
const auto& cdcTags = pProxyCommitData->cdcRouting.tagsForKey(m.param1);
transactionTags.insert(cdcTags.begin(), cdcTags.end());
}
} else if (m.type == MutationRef::ClearRange) {
auto range = pProxyCommitData->keyInfo.rangeContaining(m.param1);
if (range.end() >= m.param2) {
range.value().populateTags();
transactionTags.insert(range.value().tags.begin(), range.value().tags.end());
} else {
std::set<Tag> allSources;
while (range.begin() < m.param2) {
range.value().populateTags();
allSources.insert(range.value().tags.begin(), range.value().tags.end());
transactionTags.insert(range.value().tags.begin(), range.value().tags.end());
++range;
}
}
KeyRangeRef clearRange(KeyRangeRef(m.param1, m.param2));
if (!pProxyCommitData->cdcRouting.empty()) {
const auto cdcTags = pProxyCommitData->cdcRouting.tagsForRange(clearRange);
transactionTags.insert(cdcTags.begin(), cdcTags.end());
}
} else {
UNREACHABLE();
}
}
}
if (toCommit.getLogRouterTags()) {
toCommit.storeRandomRouterTag();
transactionTags.insert(toCommit.savedRandomRouterTag.get());
}
return transactionTags;
}
CommitBatchContext::CommitBatchContext(ProxyCommitData* const pProxyCommitData_,
const std::vector<CommitTransactionRequest>* trs_,
const int currentBatchMemBytesCount)
: pProxyCommitData(pProxyCommitData_), trs(std::move(*const_cast<std::vector<CommitTransactionRequest>*>(trs_))),
currentBatchMemBytesCount(currentBatchMemBytesCount), startTime(g_network->now()),
localBatchNumber(++pProxyCommitData->localCommitBatchesStarted),
toCommit(pProxyCommitData->logSystem, pProxyCommitData->localTLogCount), span("MP:commitBatch"_loc),
committed(trs.size()), lastShardMove(invalidVersion) {
evaluateBatchSize();
if (batchOperations != 0) {
latencyBucket =
std::min<int>(SERVER_KNOBS->PROXY_COMPUTE_BUCKETS - 1,
SERVER_KNOBS->PROXY_COMPUTE_BUCKETS * batchBytes /
(batchOperations * (CLIENT_KNOBS->VALUE_SIZE_LIMIT + CLIENT_KNOBS->KEY_SIZE_LIMIT)));
}
// since we are using just the former to limit the number of versions actually in flight!
ASSERT(SERVER_KNOBS->MAX_READ_TRANSACTION_LIFE_VERSIONS <= SERVER_KNOBS->MAX_VERSIONS_IN_FLIGHT);
}
void CommitBatchContext::setupTraceBatch() {
for (const auto& tr : trs) {
if (tr.debugID.present()) {
if (!debugID.present()) {
debugID = nondeterministicRandom()->randomUniqueID();
}
g_traceBatch.addAttach("CommitAttachID", tr.debugID.get().first(), debugID.get().first());
}
span.addLink(tr.spanContext);
}
if (debugID.present()) {
g_traceBatch.addEvent("CommitDebug", debugID.get().first(), "CommitProxyServer.commitBatch.Before");
}
}
void CommitBatchContext::evaluateBatchSize() {
for (const auto& tr : trs) {
const auto& mutations = tr.transaction.mutations;
batchOperations += mutations.size();
batchBytes += mutations.expectedSize();
}
}
// Try to identify recovery transaction and backup's apply mutations (blind writes).
// Both cannot be rejected and are approximated by looking at first mutation
// starting with 0xff.
bool canReject(const std::vector<CommitTransactionRequest>& trs) {
for (const auto& tr : trs) {
if (tr.transaction.mutations.empty())
continue;
if (tr.transaction.mutations[0].param1.startsWith("\xff"_sr) || tr.transaction.read_conflict_ranges.empty()) {
return false;
}
}
return true;
}
double computeReleaseDelay(CommitBatchContext* self, double latencyBucket) {
return std::min(SERVER_KNOBS->MAX_PROXY_COMPUTE,
self->batchOperations * self->pProxyCommitData->commitComputePerOperation[latencyBucket]);
}
Future<Void> preresolutionProcessing(CommitBatchContext* self) {
ProxyCommitData* const pProxyCommitData = self->pProxyCommitData;
std::vector<CommitTransactionRequest>& trs = self->trs;
const int64_t localBatchNumber = self->localBatchNumber;
const int latencyBucket = self->latencyBucket;
const Optional<UID>& debugID = self->debugID;
Span span("MP:preresolutionProcessing"_loc, self->span.context);
double startTime = g_network->timer_monotonic();
if (self->localBatchNumber - self->pProxyCommitData->latestLocalCommitBatchResolving.get() >
SERVER_KNOBS->RESET_MASTER_BATCHES &&
now() - self->pProxyCommitData->lastMasterReset > SERVER_KNOBS->RESET_MASTER_DELAY) {
TraceEvent(SevWarnAlways, "ResetMasterNetwork", self->pProxyCommitData->dbgid)
.detail("CurrentBatch", self->localBatchNumber)
.detail("InProcessBatch", self->pProxyCommitData->latestLocalCommitBatchResolving.get());
FlowTransport::transport().resetConnection(self->pProxyCommitData->master.address());
self->pProxyCommitData->lastMasterReset = now();
}
// Pre-resolution the commits
CODE_PROBE(pProxyCommitData->latestLocalCommitBatchResolving.get() < localBatchNumber - 1, "Wait for local batch");
co_await pProxyCommitData->latestLocalCommitBatchResolving.whenAtLeast(localBatchNumber - 1);
double queuingDelay = g_network->timer_monotonic() - startTime;
pProxyCommitData->stats.computeLatency.addMeasurement(queuingDelay);
pProxyCommitData->stats.commitBatchQueuingDist->sampleSeconds(queuingDelay);
if ((queuingDelay > (double)SERVER_KNOBS->MAX_READ_TRANSACTION_LIFE_VERSIONS / SERVER_KNOBS->VERSIONS_PER_SECOND ||
(g_network->isSimulated() && buggify(0.01))) &&
SERVER_KNOBS->PROXY_REJECT_BATCH_QUEUED_TOO_LONG && canReject(trs)) {
// Disabled for the recovery transaction. otherwise, recovery can't finish and keeps doing more recoveries.
CODE_PROBE(true, "Reject transactions in the batch");
TraceEvent(g_network->isSimulated() ? SevInfo : SevWarnAlways, "ProxyReject", pProxyCommitData->dbgid)
.suppressFor(0.1)
.detail("QDelay", queuingDelay)
.detail("Transactions", trs.size())
.detail("BatchNumber", localBatchNumber);
ASSERT(pProxyCommitData->latestLocalCommitBatchResolving.get() == localBatchNumber - 1);
pProxyCommitData->latestLocalCommitBatchResolving.set(localBatchNumber);
co_await pProxyCommitData->latestLocalCommitBatchLogging.whenAtLeast(localBatchNumber - 1);
ASSERT(pProxyCommitData->latestLocalCommitBatchLogging.get() == localBatchNumber - 1);
pProxyCommitData->latestLocalCommitBatchLogging.set(localBatchNumber);
for (const auto& tr : trs) {
tr.reply.sendError(transaction_too_old());
}
++pProxyCommitData->stats.commitBatchOut;
pProxyCommitData->stats.txnCommitOut += trs.size();
pProxyCommitData->stats.txnRejectedForQueuedTooLong += trs.size();
self->rejected = true;
co_return;
}
self->releaseDelay = delay(computeReleaseDelay(self, latencyBucket), TaskPriority::ProxyMasterVersionReply);
if (debugID.present()) {
g_traceBatch.addEvent(
"CommitDebug", debugID.get().first(), "CommitProxyServer.commitBatch.GettingCommitVersion");
}
if (SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST) {
self->writtenTagsPreResolution = self->getWrittenTagsPreResolution();
}
if (SERVER_KNOBS->HOT_SHARD_THROTTLING_ENABLED && !pProxyCommitData->hotShards.empty()) {
self->checkHotShards();
}
GetCommitVersionRequest req(span.context,
pProxyCommitData->commitVersionRequestNumber++,
pProxyCommitData->mostRecentProcessedRequestNumber,
pProxyCommitData->dbgid);
double beforeGettingCommitVersion = g_network->timer_monotonic();
GetCommitVersionReply versionReply = co_await brokenPromiseToNever(
pProxyCommitData->master.getCommitVersion.getReply(req, TaskPriority::ProxyMasterVersionReply));
pProxyCommitData->mostRecentProcessedRequestNumber = versionReply.requestNum;
pProxyCommitData->stats.txnCommitVersionAssigned += trs.size();
pProxyCommitData->stats.lastCommitVersionAssigned = versionReply.version;
pProxyCommitData->stats.getCommitVersionDist->sampleSeconds(g_network->timer_monotonic() -
beforeGettingCommitVersion);
self->commitVersion = versionReply.version;
self->prevVersion = versionReply.prevVersion;
// TraceEvent("CPGetVersion", pProxyCommitData->dbgid).detail("Master", pProxyCommitData->master.id().toString()).detail("CommitVersion", self->commitVersion).detail("PrvVersion", self->prevVersion);
for (auto it : versionReply.resolverChanges) {
auto rs = pProxyCommitData->keyResolvers.modify(it.range);
for (auto r = rs.begin(); r != rs.end(); ++r)
r->value().emplace_back(versionReply.resolverChangesVersion, it.dest);
}
//TraceEvent("ProxyGotVer", pProxyContext->dbgid).detail("Commit", commitVersion).detail("Prev", prevVersion);
if (debugID.present()) {
g_traceBatch.addEvent("CommitDebug", debugID.get().first(), "CommitProxyServer.commitBatch.GotCommitVersion");
}
}
Future<Void> getResolution(CommitBatchContext* self) {
double resolutionStart = g_network->timer_monotonic();
// Sending these requests is the fuzzy border between phase 1 and phase 2; it could conceivably overlap with
// resolution processing but is still using CPU
ProxyCommitData* pProxyCommitData = self->pProxyCommitData;
std::vector<CommitTransactionRequest>& trs = self->trs;
Span span("MP:getResolution"_loc, self->span.context);
ResolutionRequestBuilder requests(pProxyCommitData,
self->commitVersion,
self->prevVersion,
pProxyCommitData->version.get(),
self->lastShardMove,
span);
int conflictRangeCount = 0;
self->maxTransactionBytes = 0;
for (int t = 0; t < trs.size(); t++) {
requests.addTransaction(trs[t], self->commitVersion, t);
conflictRangeCount +=
trs[t].transaction.read_conflict_ranges.size() + trs[t].transaction.write_conflict_ranges.size();
//TraceEvent("MPTransactionDump", self->dbgid).detail("Snapshot", trs[t].transaction.read_snapshot);
// for(auto& m : trs[t].transaction.mutations)
self->maxTransactionBytes = std::max<int64_t>(self->maxTransactionBytes, trs[t].transaction.expectedSize());
// TraceEvent("MPTransactionsDump", self->dbgid).detail("Mutation", m.toString());
}
pProxyCommitData->stats.conflictRanges += conflictRangeCount;
for (int r = 1; r < pProxyCommitData->resolvers.size(); r++)
ASSERT(requests.requests[r].txnStateTransactions.size() == requests.requests[0].txnStateTransactions.size());
pProxyCommitData->stats.txnCommitResolving += trs.size();
std::vector<Future<ResolveTransactionBatchReply>> replies;
Future<ResolveTransactionBatchReply> singleResolverReply;
double singleResolverStart = 0;
if (pProxyCommitData->resolvers.size() == 1) {
requests.requests[0].debugID = self->debugID;
requests.requests[0].writtenTags = self->writtenTagsPreResolution;
singleResolverStart = g_network->timer_monotonic();
singleResolverReply = brokenPromiseToNever(
pProxyCommitData->resolvers[0].resolve.getReply(requests.requests[0], TaskPriority::ProxyResolverReply));
} else {
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))));
}
}
self->transactionResolverMap.swap(requests.transactionResolverMap);
// Used to report conflicting keys
self->txReadConflictRangeIndexMap.swap(requests.txReadConflictRangeIndexMap);
self->releaseFuture = releaseResolvingAfter(pProxyCommitData, self->releaseDelay, self->localBatchNumber);
if (self->localBatchNumber - self->pProxyCommitData->latestLocalCommitBatchLogging.get() >
SERVER_KNOBS->RESET_RESOLVER_BATCHES &&
now() - self->pProxyCommitData->lastResolverReset > SERVER_KNOBS->RESET_RESOLVER_DELAY) {
for (int r = 0; r < self->pProxyCommitData->resolvers.size(); r++) {
TraceEvent(SevWarnAlways, "ResetResolverNetwork", self->pProxyCommitData->dbgid)
.detail("PeerAddr", self->pProxyCommitData->resolvers[r].address())
.detail("PeerAddress", self->pProxyCommitData->resolvers[r].address())
.detail("CurrentBatch", self->localBatchNumber)
.detail("InProcessBatch", self->pProxyCommitData->latestLocalCommitBatchLogging.get());
FlowTransport::transport().resetConnection(self->pProxyCommitData->resolvers[r].address());
}
self->pProxyCommitData->lastResolverReset = now();
}
// Wait for the final resolution
if (pProxyCommitData->resolvers.size() == 1) {
ResolveTransactionBatchReply resolutionResp = co_await singleResolverReply;
pProxyCommitData->stats.resolverDist[0]->sampleSeconds(g_network->timer_monotonic() - singleResolverStart);
self->resolution.clear();
self->resolution.push_back(std::move(resolutionResp));
} else {
std::vector<ResolveTransactionBatchReply> resolutionResp = co_await getAll(replies);
self->resolution = std::move(resolutionResp);
}
self->pProxyCommitData->stats.resolutionDist->sampleSeconds(g_network->timer_monotonic() - resolutionStart);
if (self->debugID.present()) {
g_traceBatch.addEvent(
"CommitDebug", self->debugID.get().first(), "CommitProxyServer.commitBatch.AfterResolution");
}
}
void assertResolutionStateMutationsSizeConsistent(const std::vector<ResolveTransactionBatchReply>& resolution) {
for (int r = 1; r < resolution.size(); r++) {
ASSERT(resolution[r].stateMutations.size() == resolution[0].stateMutations.size());
for (int s = 0; s < resolution[r].stateMutations.size(); s++) {
ASSERT(resolution[r].stateMutations[s].size() == resolution[0].stateMutations[s].size());
}
}
}
// If the splitMutations is not empty, which means some clear range in mutations are split into multiple clear range
// ops. Modify mutations by replace the old clear range with the split clear ranges
void replaceRawClearRanges(Arena& arena,
VectorRef<MutationRef>& mutations,
std::vector<std::pair<int, std::vector<MutationRef>>>& splitMutations,
size_t totalSize,
Optional<UID> debugId = Optional<UID>()) {
if (splitMutations.empty())
return;
int i = mutations.size() - 1;
mutations.resize(arena, totalSize);
// place from back
int curr = totalSize - 1;
for (; i >= 0; --i) {
if (splitMutations.empty()) {
ASSERT_EQ(curr, i);
break;
}
if (splitMutations.back().first == i) {
ASSERT_EQ(mutations[i].type, MutationRef::ClearRange);
// TODO(gglass): legacy comment below references tenant. Possibly some
// opportunity for simplification here. Legacy comment:
// replace with tenant aligned mutations
auto& currMutations = splitMutations.back().second;
while (!currMutations.empty()) {
mutations[curr] = currMutations.back();
currMutations.pop_back();
curr--;
}
splitMutations.pop_back();
} else {
ASSERT_GT(curr, i);
mutations[curr] = mutations[i];
curr--;
}
}
ASSERT_EQ(splitMutations.size(), 0);
}
// Acknowledge transaction state store commits.
// Note: This acknowledgement will cause the transaction state store's popped version ("poppedUpTo", that's
// maintained in LogSystemDiskQueueAdapter) to get updated.
void acknowledgeTransactionStateStoreCommits(CommitBatchContext* self) {
for (auto& p : self->storeCommits) {
ASSERT(!p.second.isReady());
p.first.get().acknowledge.send(Void());
ASSERT(p.second.isReady());
}
}
// Compute and apply "metadata" effects of each other proxy's most recent batch
void applyMetadataEffect(CommitBatchContext* self) {
bool initialState = self->isMyFirstBatch;
self->firstStateMutations = self->isMyFirstBatch;
for (int versionIndex = 0; versionIndex < self->resolution[0].stateMutations.size(); versionIndex++) {
// pProxyCommitData->logAdapter->setNextVersion( ??? ); << Ideally we would be telling the log adapter that the
// pushes in this commit will be in the version at which these state mutations were committed by another proxy,
// but at present we don't have that information here. So the disk queue may be unnecessarily conservative
// about popping.
for (int transactionIndex = 0;
transactionIndex < self->resolution[0].stateMutations[versionIndex].size() && !self->forceRecovery;
transactionIndex++) {
bool committed = true;
for (int resolver = 0; resolver < self->resolution.size(); resolver++) {
committed =
committed && self->resolution[resolver].stateMutations[versionIndex][transactionIndex].committed;
}
if (committed) {
applyMetadataMutations(SpanContext(),
self->pProxyCommitData->getApplyMetadataProxyContext(),
self->arena,
self->pProxyCommitData->logSystemConsumer,
self->resolution[0].stateMutations[versionIndex][transactionIndex].mutations,
/* pToCommit= */ nullptr,
self->forceRecovery,
/* version= */ self->commitVersion,
/* popVersion= */ 0,
/* initialCommit */ false,
/* provisionalCommitProxy */ self->pProxyCommitData->provisional);
}
if (!self->resolution[0].stateMutations[versionIndex][transactionIndex].mutations.empty() &&
self->firstStateMutations) {
ASSERT(committed);
self->firstStateMutations = false;
self->forceRecovery = false;
}
}
// These changes to txnStateStore will be committed by the other proxy, so we simply discard the commit message
auto fcm = self->pProxyCommitData->logAdapter->getCommitMessage();
self->storeCommits.emplace_back(fcm, self->pProxyCommitData->txnStateStore->commit());
if (initialState) {
initialState = false;
self->forceRecovery = false;
self->pProxyCommitData->txnStateStore->resyncLog();
acknowledgeTransactionStateStoreCommits(self);
self->storeCommits.clear();
}
}
}
/// Determine which transactions actually committed (conservatively) by combining results from the resolvers
void determineCommittedTransactions(CommitBatchContext* self) {
auto pProxyCommitData = self->pProxyCommitData;
const auto& trs = self->trs;
ASSERT(self->transactionResolverMap.size() == self->committed.size());
// For each commitTransactionRef, it is only sent to resolvers specified in transactionResolverMap
// Thus, we use this nextTr to track the correct transaction index on each resolver.
self->nextTr.resize(self->resolution.size());
for (int t = 0; t < trs.size(); t++) {
uint8_t commit = ConflictBatchStatus::TransactionCommitted;
for (int r : self->transactionResolverMap[t]) {
commit = std::min(self->resolution[r].committed[self->nextTr[r]++], commit);
}
self->committed[t] = commit;
}
for (int r = 0; r < self->resolution.size(); r++)
ASSERT(self->nextTr[r] == self->resolution[r].committed.size());
pProxyCommitData->logAdapter->setNextVersion(self->commitVersion);
self->lockedKey = pProxyCommitData->txnStateStore->readValue(databaseLockedKey).get();
self->locked = self->lockedKey.present() && !self->lockedKey.get().empty();
const Optional<Value> mustContainSystemKey =
pProxyCommitData->txnStateStore->readValue(mustContainSystemMutationsKey).get();
if (mustContainSystemKey.present() && !mustContainSystemKey.get().empty()) {
for (int t = 0; t < trs.size(); t++) {
if (self->committed[t] == ConflictBatchStatus::TransactionCommitted) {
bool foundSystem = false;
for (auto& m : trs[t].transaction.mutations) {
if ((m.type == MutationRef::ClearRange ? m.param2 : m.param1) >= nonMetadataSystemKeys.end) {
foundSystem = true;
break;
}
}
if (!foundSystem) {
self->committed[t] = ConflictBatchStatus::TransactionConflict;
}
}
}
}
}
// This first pass through committed transactions deals with "metadata" effects (modifications of txnStateStore, changes
// to storage servers' responsibilities)
Future<Void> changeCoordinatorsForMetadata(CommitBatchContext* self) {
ProxyCommitData* const pProxyCommitData = self->pProxyCommitData;
co_await brokenPromiseToNever(pProxyCommitData->db->get().clusterInterface.changeCoordinators.getReply(
ChangeCoordinatorsRequest(pProxyCommitData->txnStateStore->readValue(coordinatorsKey).get().get(),
self->pProxyCommitData->master.id())));
ASSERT(false); // ChangeCoordinatorsRequest should always throw
}
Future<Void> applyMetadataToCommittedTransactions(CommitBatchContext* self) {
ProxyCommitData* const pProxyCommitData = self->pProxyCommitData;
auto& trs = self->trs;
int t;
for (t = 0; t < trs.size() && !self->forceRecovery; t++) {
if (self->committed[t] == ConflictBatchStatus::TransactionCommitted &&
(!self->locked || trs[t].isLockAware())) {
self->commitCount++;
applyMetadataMutations(trs[t].spanContext,
pProxyCommitData->getApplyMetadataProxyContext(),
self->arena,
pProxyCommitData->logSystemConsumer,
trs[t].transaction.mutations,
SERVER_KNOBS->PROXY_USE_RESOLVER_PRIVATE_MUTATIONS ? nullptr : &self->toCommit,
self->forceRecovery,
self->commitVersion,
self->commitVersion + 1,
/* initialCommit= */ false,
/* provisionalCommitProxy */ self->pProxyCommitData->provisional);
}
if (self->firstStateMutations) {
ASSERT(self->committed[t] == ConflictBatchStatus::TransactionCommitted);
self->firstStateMutations = false;
self->forceRecovery = false;
}
}
if (self->forceRecovery) {
for (; t < trs.size(); t++)
self->committed[t] = ConflictBatchStatus::TransactionConflict;
TraceEvent(SevWarn, "RestartingTxnSubsystem", pProxyCommitData->dbgid).detail("Stage", "AwaitCommit");
}
if (SERVER_KNOBS->PROXY_USE_RESOLVER_PRIVATE_MUTATIONS) {
// Resolver also calculates forceRecovery and only applies metadata mutations
// 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->pProxyCommitData->lastShardMove = reply.lastShardMove;
// extract push locations from tpcv
std::vector<int> fromLocations;
fromLocations.reserve(reply.tpcvMap.size());
for (const auto& pair : self->tpcvMap) {
fromLocations.push_back(pair.first);
}
// save push locations for each tag
self->toCommit.setPushLocationsForTags(fromLocations);
}
self->toCommit.addWrittenTags(reply.writtenTags);
}
self->lockedKey = pProxyCommitData->txnStateStore->readValue(databaseLockedKey).get();
self->lockedAfter = self->lockedKey.present() && !self->lockedKey.get().empty();
self->metadataVersionAfter = pProxyCommitData->txnStateStore->readValue(metadataVersionKey).get();
auto fcm = pProxyCommitData->logAdapter->getCommitMessage();
self->storeCommits.emplace_back(fcm, pProxyCommitData->txnStateStore->commit());
pProxyCommitData->version.set(self->commitVersion);
if (!pProxyCommitData->validState.isSet())
pProxyCommitData->validState.send(Void());
ASSERT(self->commitVersion);
if (!self->isMyFirstBatch &&
pProxyCommitData->txnStateStore->readValue(coordinatorsKey).get().get() != self->previousCoordinators.get()) {
return changeCoordinatorsForMetadata(self);
}
return Void();
}
WriteMutationRefVar writeMutation(CommitBatchContext* self, const MutationRef* mutation) {
self->toCommit.writeTypedMessage(*mutation);
return std::variant<MutationRef, VectorRef<MutationRef>>{ *mutation };
}
void pushToBackupMutations(CommitBatchContext* self,
ProxyCommitData* const pProxyCommitData,
Arena& arena,
MutationRef const& m,
MutationRef const& writtenMutation) {
if (m.type != MutationRef::Type::ClearRange) {
// Add the mutation to the relevant backup tag
for (const auto& backupName : pProxyCommitData->vecBackupKeys[m.param1]) {
self->logRangeMutations[backupName].push_back_deep(self->logRangeMutationsArena, writtenMutation);
}
} else {
KeyRangeRef mutationRange(m.param1, m.param2);
KeyRangeRef intersectionRange;
// Identify and add the intersecting ranges of the mutation to the array of mutations to serialize
for (auto backupRange : pProxyCommitData->vecBackupKeys.intersectingRanges(mutationRange)) {
// Get the backup sub range
const auto& backupSubrange = backupRange.range();
// Determine the intersecting range
intersectionRange = mutationRange & backupSubrange;
// Create the custom mutation for the specific backup tag
MutationRef backupMutation(MutationRef::Type::ClearRange, intersectionRange.begin, intersectionRange.end);
// Add the mutation to the relevant backup tag
for (const auto& backupName : backupRange.value()) {
self->logRangeMutations[backupName].push_back_deep(self->logRangeMutationsArena, backupMutation);
}
}
}
}
void addAccumulativeChecksumMutations(CommitBatchContext* self) {
ASSERT(self->pProxyCommitData->acsBuilder != nullptr);
const uint16_t acsIndex = getCommitProxyAccumulativeChecksumIndex(self->pProxyCommitData->commitProxyIndex);
for (const auto& [tag, acsState] : self->pProxyCommitData->acsBuilder->getAcsTable()) {
ASSERT(tagSupportAccumulativeChecksum(tag));
ASSERT(acsState.version <= self->commitVersion);
if (acsState.version < self->commitVersion) {
// Have not updated in the current commit batch
// So, need not send acs mutation for this tag
continue;
}
ASSERT(acsState.epoch == self->pProxyCommitData->epoch);
MutationRef acsMutation;
acsMutation.type = MutationRef::SetValue;
acsMutation.param1 = accumulativeChecksumKey; // private mutation
AccumulativeChecksumState acsToSend(acsIndex, acsState.acs, self->commitVersion, self->pProxyCommitData->epoch);
Value acsValue = accumulativeChecksumValue(acsToSend);
acsMutation.param2 = acsValue;
acsMutation.setAccumulativeChecksumIndex(acsIndex);
if (CLIENT_KNOBS->ENABLE_ACCUMULATIVE_CHECKSUM_LOGGING) {
TraceEvent(SevInfo, "AcsBuilderIssueAccumulativeChecksumMutation", self->pProxyCommitData->dbgid)
.detail("AcsTag", tag)
.detail("AcsIndex", acsIndex)
.detail("AcsToSend", acsToSend.toString())
.detail("Mutation", acsMutation)
.detail("Version", self->commitVersion)
.detail("CommitProxyIndex", self->pProxyCommitData->commitProxyIndex);
}
DEBUG_MUTATION("ProxyCommit", self->commitVersion, acsMutation, self->pProxyCommitData->dbgid);
self->toCommit.addTag(tag);
self->toCommit.writeTypedMessage(acsMutation);
}
}
void rejectMutationsForReadLockOnRange(CommitBatchContext* self) {
ASSERT(self->rangeLockEnabled());
ProxyCommitData* const pProxyCommitData = self->pProxyCommitData;
ASSERT(pProxyCommitData->rangeLock != nullptr);
// Fast path: no exclusive locks held -> nothing to check, skip the
// per-mutation loop entirely. Steady state when no bulkload is running.
if (!pProxyCommitData->rangeLock->anyExclusiveLockHeld()) {
++pProxyCommitData->stats.rangeLockFastPath;
return;
}
++pProxyCommitData->stats.rangeLockSlowPath;
std::vector<CommitTransactionRequest>& trs = self->trs;
for (int i = self->transactionNum; i < trs.size(); i++) {
if (self->committed[i] != ConflictBatchStatus::TransactionCommitted) {
continue;
} else if (trs[i].isLockAware()) {
continue; // rangeLock is transparent to lock-aware transactions
}
VectorRef<MutationRef>* pMutations = &trs[i].transaction.mutations;
for (int j = 0; j < pMutations->size(); j++) {
MutationRef m = (*pMutations)[j];
KeyRange rangeToCheck;
if (isSingleKeyMutation((MutationRef::Type)m.type)) {
rangeToCheck = singleKeyRange(m.param1);
} else if (m.type == MutationRef::ClearRange) {
rangeToCheck = KeyRangeRef(m.param1, m.param2);
}
bool shouldReject = pProxyCommitData->rangeLock->isLocked(rangeToCheck);
if (shouldReject) {
self->committed[i] = ConflictBatchStatus::TransactionLockReject;
trs[i].reply.sendError(transaction_rejected_range_locked());
break;
}
}
}
}
/// This second pass through committed transactions assigns the actual mutations to the appropriate storage servers'
/// tags
Future<Void> assignMutationsToStorageServers(CommitBatchContext* self) {
ProxyCommitData* const pProxyCommitData = self->pProxyCommitData;
std::vector<CommitTransactionRequest>& trs = self->trs;
for (; self->transactionNum < trs.size(); self->transactionNum++) {
if (!(self->committed[self->transactionNum] == ConflictBatchStatus::TransactionCommitted &&
(!self->locked || trs[self->transactionNum].isLockAware()))) {
continue;
}
bool checkSample = trs[self->transactionNum].commitCostEstimation.present();
Optional<ClientTrCommitCostEstimation>* trCost = &trs[self->transactionNum].commitCostEstimation;
int mutationNum = 0;
VectorRef<MutationRef>* pMutations = &trs[self->transactionNum].transaction.mutations;
self->toCommit.addTransactionInfo(trs[self->transactionNum].spanContext);
for (; mutationNum < pMutations->size(); mutationNum++) {
if (self->yieldBytes > SERVER_KNOBS->DESIRED_TOTAL_BYTES) {
self->yieldBytes = 0;
if (g_network->check_yield(TaskPriority::ProxyCommitYield1)) {
self->computeDuration += g_network->timer_monotonic() - self->computeStart;
co_await delay(0, TaskPriority::ProxyCommitYield1);
self->computeStart = g_network->timer_monotonic();
}
}
MutationRef m = (*pMutations)[mutationNum];
Arena arena;
MutationRef writtenMutation;
self->mutationCount++;
self->mutationBytes += m.expectedSize();
self->yieldBytes += m.expectedSize();
// Determine the set of tags (responsible storage servers) for the mutation, splitting it
// if necessary. Serialize (splits of) the mutation into the message buffer and add the tags.
if (isSingleKeyMutation((MutationRef::Type)m.type)) {
auto& tags = pProxyCommitData->tagsForKey(m.param1);
// sample single key mutation based on cost
// the expectation of sampling is every COMMIT_SAMPLE_COST sample once
if (checkSample) {
double totalCosts = trCost->get().writeCosts;
double cost = getWriteOperationCost(m.expectedSize());
double mul = std::max(1.0, totalCosts / std::max(1.0, (double)CLIENT_KNOBS->COMMIT_SAMPLE_COST));
ASSERT(totalCosts > 0);
double prob = mul * cost / totalCosts;
if (deterministicRandom()->random01() < prob) {
const auto& storageServers = pProxyCommitData->keyInfo[m.param1].src_info;
for (const auto& ssInfo : storageServers) {
auto id = ssInfo->interf.id();
// scale cost
cost = cost < CLIENT_KNOBS->COMMIT_SAMPLE_COST ? CLIENT_KNOBS->COMMIT_SAMPLE_COST : cost;
pProxyCommitData->updateSSTagCost(
id, trs[self->transactionNum].tagSet.get(), m, cost / storageServers.size());
}
}
}
DEBUG_MUTATION("ProxyCommit", self->commitVersion, m, pProxyCommitData->dbgid).detail("To", tags);
self->toCommit.addTags(tags);
if (!pProxyCommitData->cdcRouting.empty()) {
self->toCommit.addTags(pProxyCommitData->cdcRouting.tagsForKey(m.param1));
}
if (pProxyCommitData->acsBuilder != nullptr) {
updateMutationWithAcsAndAddMutationToAcsBuilder(
pProxyCommitData->acsBuilder,
m,
tags,
getCommitProxyAccumulativeChecksumIndex(pProxyCommitData->commitProxyIndex),
pProxyCommitData->epoch,
self->commitVersion,
pProxyCommitData->dbgid);
}
WriteMutationRefVar var = writeMutation(self, &m);
// FIXME: Remove assert once ClearRange RAW_ACCESS usecase handling is done
ASSERT(std::holds_alternative<MutationRef>(var));
writtenMutation = std::get<MutationRef>(var);
} else if (m.type == MutationRef::ClearRange) {
auto range = pProxyCommitData->keyInfo.rangeContaining(m.param1);
if (range.end() >= m.param2) {
// Fast path
DEBUG_MUTATION("ProxyCommit", self->commitVersion, m, pProxyCommitData->dbgid)
.detail("To", range.value().tags);
range.value().populateTags();
self->toCommit.addTags(range.value().tags);
if (pProxyCommitData->acsBuilder != nullptr) {
updateMutationWithAcsAndAddMutationToAcsBuilder(
pProxyCommitData->acsBuilder,
m,
range.value().tags,
getCommitProxyAccumulativeChecksumIndex(pProxyCommitData->commitProxyIndex),
pProxyCommitData->epoch,
self->commitVersion,
pProxyCommitData->dbgid);
}
// check whether clear is sampled
if (checkSample && !trCost->get().clearIdxCosts.empty() &&
trCost->get().clearIdxCosts[0].first == mutationNum) {
auto const& ssInfos = range.value().src_info;
for (auto const& ssInfo : ssInfos) {
auto id = ssInfo->interf.id();
pProxyCommitData->updateSSTagCost(id,
trs[self->transactionNum].tagSet.get(),
m,
trCost->get().clearIdxCosts[0].second / ssInfos.size());
}
trCost->get().clearIdxCosts.pop_front();
}
} else {
CODE_PROBE(true, "A clear range extends past a shard boundary");
std::set<Tag> allSources;
while (range.begin() < m.param2) {
range.value().populateTags();
allSources.insert(range.value().tags.begin(), range.value().tags.end());
// check whether clear is sampled
if (checkSample && !trCost->get().clearIdxCosts.empty() &&
trCost->get().clearIdxCosts[0].first == mutationNum) {
auto const& ssInfos = range.value().src_info;
for (auto const& ssInfo : ssInfos) {
auto id = ssInfo->interf.id();
pProxyCommitData->updateSSTagCost(id,
trs[self->transactionNum].tagSet.get(),
m,
trCost->get().clearIdxCosts[0].second /
ssInfos.size());
}
trCost->get().clearIdxCosts.pop_front();
}
++range;
}
DEBUG_MUTATION("ProxyCommit", self->commitVersion, m)
.detail("Dbgid", pProxyCommitData->dbgid)
.detail("To", allSources);
self->toCommit.addTags(allSources);
if (self->pProxyCommitData->acsBuilder != nullptr) {
updateMutationWithAcsAndAddMutationToAcsBuilder(
pProxyCommitData->acsBuilder,
m,
allSources,
getCommitProxyAccumulativeChecksumIndex(pProxyCommitData->commitProxyIndex),
pProxyCommitData->epoch,
self->commitVersion,
pProxyCommitData->dbgid);
}
}
KeyRangeRef clearRange(KeyRangeRef(m.param1, m.param2));
if (!pProxyCommitData->cdcRouting.empty()) {
self->toCommit.addTags(pProxyCommitData->cdcRouting.tagsForRange(clearRange));
}
WriteMutationRefVar var = writeMutation(self, &m);
// FIXME: Remove assert once ClearRange RAW_ACCESS usecase handling is done
ASSERT(std::holds_alternative<MutationRef>(var));
writtenMutation = std::get<MutationRef>(var);
} else if (m.type == MutationRef::NoOp) {
// TODO(gglass): what is the deal with MutationRef::NoOp? Is it needed?
// This used to be the following:
// ASSERT_EQ(pProxyCommitData->getTenantMode(), TenantMode::REQUIRED);
ASSERT(false);
continue;
} else {
UNREACHABLE();
}
DisabledTraceEvent(SevDebug, "BeforeBackup", pProxyCommitData->dbgid)
.detail("M1", m.param1)
.detail("M2", m.param2)
.detail("MT", getTypeString(m.type))
.detail("VecBackupKeys", pProxyCommitData->vecBackupKeys.size())
.detail("ShouldBackup", shouldBackup(m));
if (pProxyCommitData->vecBackupKeys.size() <= 1 || !shouldBackup(m)) {
continue;
}
pushToBackupMutations(self, pProxyCommitData, arena, m, writtenMutation);
}
if (checkSample) {
self->pProxyCommitData->stats.txnExpensiveClearCostEstCount +=
trs[self->transactionNum].commitCostEstimation.get().expensiveCostEstCount;
}
}
}
Future<Void> postResolution(CommitBatchContext* self) {
double postResolutionStart = g_network->timer_monotonic();
ProxyCommitData* const pProxyCommitData = self->pProxyCommitData;
std::vector<CommitTransactionRequest>& trs = self->trs;
const int64_t localBatchNumber = self->localBatchNumber;
const Optional<UID>& debugID = self->debugID;
Span span("MP:postResolution"_loc, self->span.context);
bool queuedCommits = pProxyCommitData->latestLocalCommitBatchLogging.get() < localBatchNumber - 1;
CODE_PROBE(queuedCommits, "Queuing post-resolution commit processing");
co_await pProxyCommitData->latestLocalCommitBatchLogging.whenAtLeast(localBatchNumber - 1);
double postResolutionQueuing = g_network->timer_monotonic();
pProxyCommitData->stats.postResolutionDist->sampleSeconds(postResolutionQueuing - postResolutionStart);
co_await yield(TaskPriority::ProxyCommitYield1);
self->computeStart = g_network->timer_monotonic();
pProxyCommitData->stats.txnCommitResolved += trs.size();
if (debugID.present()) {
g_traceBatch.addEvent(
"CommitDebug", debugID.get().first(), "CommitProxyServer.commitBatch.ProcessingMutations");
}
self->isMyFirstBatch = !pProxyCommitData->version.get();
self->previousCoordinators = pProxyCommitData->txnStateStore->readValue(coordinatorsKey).get();
assertResolutionStateMutationsSizeConsistent(self->resolution);
applyMetadataEffect(self);
if (debugID.present()) {
g_traceBatch.addEvent(
"CommitDebug", debugID.get().first(), "CommitProxyServer.commitBatch.ApplyMetadataEffect");
}
determineCommittedTransactions(self);
if (debugID.present()) {
g_traceBatch.addEvent(
"CommitDebug", debugID.get().first(), "CommitProxyServer.commitBatch.DetermineCommittedTransactions");
}
if (self->forceRecovery) {
co_await Future<Void>(Never());
}
// First pass
co_await applyMetadataToCommittedTransactions(self);
if (debugID.present()) {
g_traceBatch.addEvent(
"CommitDebug", debugID.get().first(), "CommitProxyServer.commitBatch.ApplyMetadataToCommittedTxn");
}
// After applyed metadata change, this commit proxy has the latest view of locked ranges.
// If a transaction has any mutation accessing to the locked range, reject the transaction with
// error_code_transaction_rejected_range_locked
if (self->rangeLockEnabled()) {
rejectMutationsForReadLockOnRange(self);
}
// Second pass
co_await assignMutationsToStorageServers(self);
if (debugID.present()) {
g_traceBatch.addEvent("CommitDebug", debugID.get().first(), "CommitProxyServer.commitBatch.AssignMutationToSS");
}
// Serialize and backup the mutations as a single mutation
if ((pProxyCommitData->vecBackupKeys.size() > 1) && !self->logRangeMutations.empty()) {
co_await addBackupMutations(pProxyCommitData,
&self->logRangeMutations,
&self->toCommit,
self->commitVersion,
&self->computeDuration,
&self->computeStart);
}
// When version vector is enabled, idempotency entries should only be created or cleared
// if the operation was detected at pre resolution time. This ensures that the
// operation is broadcast to all logs, and does not lead to logs being included
// that are not part of the expected tag set (tpcv).
buildIdempotencyIdMutations(self->trs,
self->idempotencyKVBuilder,
self->commitVersion,
self->committed,
ConflictBatchStatus::TransactionCommitted,
self->locked,
[&](const KeyValue& kv) {
MutationRef idempotencyIdSet;
idempotencyIdSet.type = MutationRef::Type::SetValue;
idempotencyIdSet.param1 = kv.key;
idempotencyIdSet.param2 = kv.value;
auto& tags = pProxyCommitData->tagsForKey(kv.key);
ASSERT(!SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST ||
pProxyCommitData->db->get().logSystemConfig.numLogs() ==
self->tpcvMap.size());
self->toCommit.addTags(tags);
if (pProxyCommitData->acsBuilder != nullptr) {
updateMutationWithAcsAndAddMutationToAcsBuilder(
pProxyCommitData->acsBuilder,
idempotencyIdSet,
tags,
getCommitProxyAccumulativeChecksumIndex(pProxyCommitData->commitProxyIndex),
pProxyCommitData->epoch,
self->commitVersion,
pProxyCommitData->dbgid);
}
self->toCommit.writeTypedMessage(idempotencyIdSet);
});
if (!SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST ||
pProxyCommitData->db->get().logSystemConfig.numLogs() == self->tpcvMap.size()) {
for (int i = 0; i < pProxyCommitData->idempotencyClears.size(); i++) {
auto& tags = pProxyCommitData->tagsForKey(pProxyCommitData->idempotencyClears[i].param1);
self->toCommit.addTags(tags);
if (pProxyCommitData->acsBuilder != nullptr) {
updateMutationWithAcsAndAddMutationToAcsBuilder(
pProxyCommitData->acsBuilder,
pProxyCommitData->idempotencyClears[i],
tags,
getCommitProxyAccumulativeChecksumIndex(pProxyCommitData->commitProxyIndex),
pProxyCommitData->epoch,
self->commitVersion,
pProxyCommitData->dbgid);
}
WriteMutationRefVar var = writeMutation(self, &pProxyCommitData->idempotencyClears[i]);
ASSERT(std::holds_alternative<MutationRef>(var));
}
pProxyCommitData->idempotencyClears = Standalone<VectorRef<MutationRef>>();
}
self->toCommit.saveTags(self->writtenTags);
pProxyCommitData->stats.mutations += self->mutationCount;
pProxyCommitData->stats.mutationBytes += self->mutationBytes;
// Storage servers mustn't make durable versions which are not fully committed (because then they are impossible
// to roll back) We prevent this by limiting the number of versions which are semi-committed but not fully
// committed to be less than the MVCC window
if (pProxyCommitData->committedVersion.get() <
self->commitVersion - SERVER_KNOBS->MAX_READ_TRANSACTION_LIFE_VERSIONS) {
self->computeDuration += g_network->timer_monotonic() - self->computeStart;
Span waitVersionSpan;
while (pProxyCommitData->committedVersion.get() <
self->commitVersion - SERVER_KNOBS->MAX_READ_TRANSACTION_LIFE_VERSIONS) {
// This should be *extremely* rare in the real world, but knob buggification should make it happen in
// simulation
CODE_PROBE(true, "Semi-committed pipeline limited by MVCC window");
//TraceEvent("ProxyWaitingForCommitted", pProxyCommitData->dbgid).detail("CommittedVersion", pProxyCommitData->committedVersion.get()).detail("NeedToCommit", commitVersion);
waitVersionSpan = Span("MP:overMaxReadTransactionLifeVersions"_loc, span.context);
// @todo probably there is no need to get the (entire) version vector from the sequencer
// in this case, and if so, consider adding a flag to the request to tell the sequencer
// to not send the version vector information.
auto res =
co_await race(pProxyCommitData->committedVersion.whenAtLeast(
self->commitVersion - SERVER_KNOBS->MAX_READ_TRANSACTION_LIFE_VERSIONS),
pProxyCommitData->cx->onProxiesChanged(),
pProxyCommitData->master.getLiveCommittedVersion.getReply(
GetRawCommittedVersionRequest(waitVersionSpan.context, debugID, invalidVersion),
TaskPriority::GetLiveCommittedVersionReply));
if (res.index() == 0) {
co_await yield();
break;
} else if (res.index() == 1) {
} else if (res.index() == 2) {
GetRawCommittedVersionReply v = std::get<2>(std::move(res));
if (v.version > pProxyCommitData->committedVersion.get()) {
pProxyCommitData->locked = v.locked;
pProxyCommitData->metadataVersion = v.metadataVersion;
pProxyCommitData->committedVersion.set(v.version);
}
if (pProxyCommitData->committedVersion.get() <
self->commitVersion - SERVER_KNOBS->MAX_READ_TRANSACTION_LIFE_VERSIONS) {
co_await delay(SERVER_KNOBS->PROXY_SPIN_DELAY);
}
} else {
UNREACHABLE();
}
}
waitVersionSpan = Span{};
self->computeStart = g_network->timer_monotonic();
}
self->msg = self->storeCommits.back().first.get();
if (self->debugID.present())
g_traceBatch.addEvent(
"CommitDebug", self->debugID.get().first(), "CommitProxyServer.commitBatch.AfterStoreCommits");
// txnState (transaction subsystem state) tag: message extracted from log adapter
bool firstMessage = true;
for (auto m : self->msg.messages) {
if (firstMessage) {
ASSERT(!SERVER_KNOBS->ENABLE_VERSION_VECTOR ||
pProxyCommitData->db->get().logSystemConfig.numLogs() == self->tpcvMap.size());
self->toCommit.addTxsTag();
}
self->toCommit.writeMessage(StringRef(m.begin(), m.size()), !firstMessage);
firstMessage = false;
}
if (self->prevVersion && self->commitVersion - self->prevVersion < SERVER_KNOBS->MAX_VERSIONS_IN_FLIGHT / 2)
debug_advanceMaxCommittedVersion(UID(), self->commitVersion); //< Is this valid?
// TraceEvent("ProxyPush", pProxyCommitData->dbgid)
// .detail("PrevVersion", self->prevVersion)
// .detail("Version", self->commitVersion)
// .detail("TransactionsSubmitted", trs.size())
// .detail("TransactionsCommitted", self->commitCount)
// .detail("TxsPopTo", self->msg.popTo);
if (self->prevVersion && self->commitVersion - self->prevVersion < SERVER_KNOBS->MAX_VERSIONS_IN_FLIGHT / 2)
debug_advanceMaxCommittedVersion(UID(), self->commitVersion);
self->commitStartTime = now();
pProxyCommitData->lastStartCommit = self->commitStartTime;
Optional<std::unordered_map<uint16_t, Version>> tpcvMap = Optional<std::unordered_map<uint16_t, Version>>();
if (SERVER_KNOBS->ENABLE_VERSION_VECTOR) {
tpcvMap = self->tpcvMap;
}
if (self->pProxyCommitData->acsBuilder != nullptr) {
// Issue acs mutation at the end of this commit batch
addAccumulativeChecksumMutations(self);
}
const auto versionSet = LogPushVersionSet{ self->prevVersion,
self->commitVersion,
pProxyCommitData->committedVersion.get(),
pProxyCommitData->minKnownCommittedVersion };
self->loggingComplete =
pProxyCommitData->logSystem->push(versionSet, self->toCommit, span.context, self->debugID, tpcvMap);
float ratio = self->toCommit.getEmptyMessageRatio();
pProxyCommitData->stats.commitBatchingEmptyMessageRatio.addMeasurement(ratio);
if (!self->forceRecovery) {
ASSERT(pProxyCommitData->latestLocalCommitBatchLogging.get() == self->localBatchNumber - 1);
pProxyCommitData->latestLocalCommitBatchLogging.set(self->localBatchNumber);
}
self->computeDuration += g_network->timer_monotonic() - self->computeStart;
if (self->batchOperations > 0) {
double estimatedDelay = computeReleaseDelay(self, self->latencyBucket);
double computePerOperation =
std::min(SERVER_KNOBS->MAX_COMPUTE_PER_OPERATION, self->computeDuration / self->batchOperations);
if (computePerOperation <= pProxyCommitData->commitComputePerOperation[self->latencyBucket]) {
pProxyCommitData->commitComputePerOperation[self->latencyBucket] = computePerOperation;
} else {
pProxyCommitData->commitComputePerOperation[self->latencyBucket] =
SERVER_KNOBS->PROXY_COMPUTE_GROWTH_RATE * computePerOperation +
((1.0 - SERVER_KNOBS->PROXY_COMPUTE_GROWTH_RATE) *
pProxyCommitData->commitComputePerOperation[self->latencyBucket]);
}
pProxyCommitData->stats.maxComputeNS =
std::max<int64_t>(pProxyCommitData->stats.maxComputeNS,
1e9 * pProxyCommitData->commitComputePerOperation[self->latencyBucket]);
pProxyCommitData->stats.minComputeNS =
std::min<int64_t>(pProxyCommitData->stats.minComputeNS,
1e9 * pProxyCommitData->commitComputePerOperation[self->latencyBucket]);
if (estimatedDelay >= SERVER_KNOBS->MAX_COMPUTE_DURATION_LOG_CUTOFF ||
self->computeDuration >= SERVER_KNOBS->MAX_COMPUTE_DURATION_LOG_CUTOFF) {
TraceEvent(SevInfo, "LongComputeDuration", pProxyCommitData->dbgid)
.suppressFor(10.0)
.detail("EstimatedComputeDuration", estimatedDelay)
.detail("ComputeDuration", self->computeDuration)
.detail("ComputePerOperation", computePerOperation)
.detail("LatencyBucket", self->latencyBucket)
.detail("UpdatedComputePerOperationEstimate",
pProxyCommitData->commitComputePerOperation[self->latencyBucket])
.detail("BatchBytes", self->batchBytes)
.detail("BatchOperations", self->batchOperations);
}
}
pProxyCommitData->stats.processingMutationDist->sampleSeconds(g_network->timer_monotonic() - postResolutionQueuing);
}
Future<Void> transactionLogging(CommitBatchContext* self) {
double tLoggingStart = g_network->timer_monotonic();
ProxyCommitData* const pProxyCommitData = self->pProxyCommitData;
Span span("MP:transactionLogging"_loc, self->span.context);
try {
auto res = co_await race(self->loggingComplete,
pProxyCommitData->committedVersion.whenAtLeast(self->commitVersion + 1));
if (res.index() == 0) {
Version ver = std::get<0>(res);
if (!SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST) {
pProxyCommitData->minKnownCommittedVersion = std::max(pProxyCommitData->minKnownCommittedVersion, ver);
}
}
} catch (Error& e) {
if (e.code() == error_code_broken_promise) {
throw tlog_failed();
}
throw;
}
pProxyCommitData->lastCommitLatency = now() - self->commitStartTime;
pProxyCommitData->lastCommitTime = std::max(pProxyCommitData->lastCommitTime.get(), self->commitStartTime);
co_await yield(TaskPriority::ProxyCommitYield2);
if (pProxyCommitData->popRemoteTxs &&
self->msg.popTo > (!pProxyCommitData->txsPopVersions.empty() ? pProxyCommitData->txsPopVersions.back().second
: pProxyCommitData->lastTxsPop)) {
if (pProxyCommitData->txsPopVersions.size() >= SERVER_KNOBS->MAX_TXS_POP_VERSION_HISTORY) {
TraceEvent(SevWarnAlways, "DiscardingTxsPopHistory").suppressFor(1.0);
pProxyCommitData->txsPopVersions.pop_front();
}
pProxyCommitData->txsPopVersions.emplace_back(self->commitVersion, self->msg.popTo);
}
pProxyCommitData->logSystemConsumer->popTxs(self->msg.popTo);
pProxyCommitData->stats.tlogLoggingDist->sampleSeconds(g_network->timer_monotonic() - tLoggingStart);
}
Future<Void> reply(CommitBatchContext* self) {
double replyStart = g_network->timer_monotonic();
ProxyCommitData* const pProxyCommitData = self->pProxyCommitData;
Span span("MP:reply"_loc, self->span.context);
const Optional<UID>& debugID = self->debugID;
if (!SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST) {
// Version vector/unicast is disabled: Logging completed, so the current version (and all versions prior to
// the current versions) can be treated as commited.
// Advance min committed version.
if (self->prevVersion && self->commitVersion - self->prevVersion < SERVER_KNOBS->MAX_VERSIONS_IN_FLIGHT / 2) {
//TraceEvent("CPAdvanceMinVersion", self->pProxyCommitData->dbgid).detail("PrvVersion", self->prevVersion).detail("CommitVersion", self->commitVersion).detail("Master", self->pProxyCommitData->master.id().toString()).detail("TxSize", self->trs.size());
debug_advanceMinCommittedVersion(UID(), self->commitVersion);
}
// Acknowledge transaction state store commits.
acknowledgeTransactionStateStoreCommits(self);
}
// TraceEvent("ProxyPushed", pProxyCommitData->dbgid)
// .detail("PrevVersion", self->prevVersion)
// .detail("Version", self->commitVersion);
if (debugID.present())
g_traceBatch.addEvent("CommitDebug", debugID.get().first(), "CommitProxyServer.commitBatch.AfterLogPush");
// After logging finishes, we report the commit version to master so that every other proxy can get the most
// up-to-date live committed version. We also maintain the invariant that master's committed version >=
// self->committedVersion by reporting commit version first before updating self->committedVersion. Otherwise, a
// client may get a commit version that the master is not aware of, and next GRV request may get a version less
// than self->committedVersion.
CODE_PROBE(pProxyCommitData->committedVersion.get() > self->commitVersion,
"later version was reported committed first");
if (self->commitVersion >= pProxyCommitData->committedVersion.get()) {
Optional<std::set<Tag>> writtenTags;
if (SERVER_KNOBS->ENABLE_VERSION_VECTOR) {
writtenTags = self->writtenTags;
}
co_await pProxyCommitData->master.reportLiveCommittedVersion.getReply(
ReportRawCommittedVersionRequest(self->commitVersion,
self->lockedAfter,
self->metadataVersionAfter,
pProxyCommitData->minKnownCommittedVersion,
self->prevVersion,
writtenTags),
TaskPriority::ProxyMasterVersionReply);
}
if (debugID.present()) {
g_traceBatch.addEvent(
"CommitDebug", debugID.get().first(), "CommitProxyServer.commitBatch.AfterReportRawCommittedVersion");
}
if (SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST) {
// Version vector/unicast is enabled: Received a reply from the sequencer, so we can treat the
// current version (and all versions prior to the current version) as committed.
// Advance min committed version now.
if (self->prevVersion && self->commitVersion - self->prevVersion < SERVER_KNOBS->MAX_VERSIONS_IN_FLIGHT / 2) {
//TraceEvent("CPAdvanceMinVersion", self->pProxyCommitData->dbgid).detail("PrvVersion", self->prevVersion).detail("CommitVersion", self->commitVersion).detail("Master", self->pProxyCommitData->master.id().toString()).detail("TxSize", self->trs.size());
debug_advanceMinCommittedVersion(UID(), self->commitVersion);
}
// Acknowledge transaction state store commits.
acknowledgeTransactionStateStoreCommits(self);
}
if (self->commitVersion > pProxyCommitData->committedVersion.get()) {
pProxyCommitData->locked = self->lockedAfter;
pProxyCommitData->metadataVersion = self->metadataVersionAfter;
pProxyCommitData->committedVersion.set(self->commitVersion);
if (SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST) {
ASSERT(self->loggingComplete.isReady());
pProxyCommitData->minKnownCommittedVersion =
std::max(pProxyCommitData->minKnownCommittedVersion, self->loggingComplete.get());
}
}
if (self->forceRecovery) {
TraceEvent(SevWarn, "RestartingTxnSubsystem", pProxyCommitData->dbgid).detail("Stage", "ProxyShutdown");
throw worker_removed();
}
// Send replies to clients
// TODO: should be timer_monotonic(), but gets compared to request time, which uses g_network->timer().
double endTime = g_network->timer();
// Reset all to zero, used to track the correct index of each commitTransacitonRef on each resolver
std::fill(self->nextTr.begin(), self->nextTr.end(), 0);
std::unordered_map<uint8_t, int16_t> idCountsForKey;
for (int t = 0; t < self->trs.size(); t++) {
auto& tr = self->trs[t];
if (self->committed[t] == ConflictBatchStatus::TransactionCommitted && (!self->locked || tr.isLockAware())) {
ASSERT_WE_THINK(self->commitVersion != invalidVersion);
if (self->trs[t].idempotencyId.valid()) {
idCountsForKey[uint8_t(t >> 8)] += 1;
}
tr.reply.send(CommitID(self->commitVersion, t, self->metadataVersionAfter));
} else if (self->committed[t] == ConflictBatchStatus::TransactionTooOld) {
tr.reply.sendError(transaction_too_old());
} else if (self->committed[t] == ConflictBatchStatus::TransactionLockReject) {
// We already sent the error
ASSERT(tr.reply.isSet());
} else {
// If enable the option to report conflicting keys from resolvers, we send back all keyranges' indices
// through CommitID
if (tr.transaction.report_conflicting_keys) {
Standalone<VectorRef<int>> conflictingKRIndices;
for (int resolverInd : self->transactionResolverMap[t]) {
auto const& cKRs =
self->resolution[resolverInd]
.conflictingKeyRangeMap[self->nextTr[resolverInd]]; // nextTr[resolverInd] -> index of
// this trs[t] on the resolver
for (auto const& rCRIndex : cKRs) {
// read_conflict_range can change when sent to resolvers, mapping the index from
// resolver-side to original index in commitTransactionRef
conflictingKRIndices.push_back(conflictingKRIndices.arena(),
self->txReadConflictRangeIndexMap[t][resolverInd][rCRIndex]);
}
}
// At least one keyRange index should be returned
ASSERT(!conflictingKRIndices.empty());
tr.reply.send(CommitID(
invalidVersion, t, Optional<Value>(), Optional<Standalone<VectorRef<int>>>(conflictingKRIndices)));
} else {
tr.reply.sendError(not_committed());
}
}
// Update corresponding transaction indices on each resolver
for (int resolverInd : self->transactionResolverMap[t])
self->nextTr[resolverInd]++;
// TODO: filter if pipelined with large commit
const double duration = endTime - tr.requestTime();
pProxyCommitData->stats.commitLatencySample.addMeasurement(duration);
if (pProxyCommitData->latencyBandConfig.present()) {
bool filter = self->maxTransactionBytes >
pProxyCommitData->latencyBandConfig.get().commitConfig.maxCommitBytes.orDefault(
std::numeric_limits<int>::max());
pProxyCommitData->stats.commitLatencyBands.addMeasurement(duration, 1, Filtered(filter));
}
}
for (auto [highOrderBatchIndex, count] : idCountsForKey) {
pProxyCommitData->expectedIdempotencyIdCountForKey.send(
ExpectedIdempotencyIdCountForKey{ self->commitVersion, count, highOrderBatchIndex });
}
++pProxyCommitData->stats.commitBatchOut;
pProxyCommitData->stats.txnCommitOut += self->trs.size();
pProxyCommitData->stats.txnConflicts += self->trs.size() - self->commitCount;
pProxyCommitData->stats.txnCommitOutSuccess += self->commitCount;
if (now() - pProxyCommitData->lastCoalesceTime > SERVER_KNOBS->RESOLVER_COALESCE_TIME) {
pProxyCommitData->lastCoalesceTime = now();
int lastSize = pProxyCommitData->keyResolvers.size();
auto rs = pProxyCommitData->keyResolvers.ranges();
Version oldestVersion = self->prevVersion - SERVER_KNOBS->MAX_WRITE_TRANSACTION_LIFE_VERSIONS;
for (auto r = rs.begin(); r != rs.end(); ++r) {
while (r->value().size() > 1 && r->value()[1].first < oldestVersion)
r->value().pop_front();
if (!r->value().empty() && r->value().front().first < oldestVersion)
r->value().front().first = 0;
}
if (SERVER_KNOBS->PROXY_USE_RESOLVER_PRIVATE_MUTATIONS) {
// Only normal key space, because \xff key space is processed by all resolvers.
pProxyCommitData->keyResolvers.coalesce(normalKeys);
auto& versions = pProxyCommitData->systemKeyVersions;
while (versions.size() > 1 && versions[1] < oldestVersion) {
versions.pop_front();
}
if (!versions.empty() && versions[0] < oldestVersion) {
versions[0] = 0;
}
} else {
pProxyCommitData->keyResolvers.coalesce(allKeys);
}
if (pProxyCommitData->keyResolvers.size() != lastSize)
TraceEvent("KeyResolverSize", pProxyCommitData->dbgid)
.detail("Size", pProxyCommitData->keyResolvers.size());
}
// Dynamic batching for commits
double target_latency =
(now() - self->startTime) * SERVER_KNOBS->COMMIT_TRANSACTION_BATCH_INTERVAL_LATENCY_FRACTION;
pProxyCommitData->commitBatchInterval =
std::max(SERVER_KNOBS->COMMIT_TRANSACTION_BATCH_INTERVAL_MIN,
std::min(SERVER_KNOBS->COMMIT_TRANSACTION_BATCH_INTERVAL_MAX,
target_latency * SERVER_KNOBS->COMMIT_TRANSACTION_BATCH_INTERVAL_SMOOTHER_ALPHA +
pProxyCommitData->commitBatchInterval *
(1 - SERVER_KNOBS->COMMIT_TRANSACTION_BATCH_INTERVAL_SMOOTHER_ALPHA)));
pProxyCommitData->stats.commitBatchingWindowSize.addMeasurement(pProxyCommitData->commitBatchInterval);
pProxyCommitData->commitBatchesMemBytesCount -= self->currentBatchMemBytesCount;
ASSERT_ABORT(pProxyCommitData->commitBatchesMemBytesCount >= 0);
co_await self->releaseFuture;
pProxyCommitData->stats.replyCommitDist->sampleSeconds(g_network->timer_monotonic() - replyStart);
}
// Commit one batch of transactions trs
Future<Void> commitBatchImpl(CommitBatchContext* pContext) {
// WARNING: this code is run at a high priority (until the first delay(0)), so it needs to do as little work as
// possible
pContext->stage = INITIALIZE;
getCurrentLineage()->modify(&TransactionLineage::operation) = TransactionLineage::Operation::Commit;
// Active load balancing runs at a very high priority (to obtain accurate estimate of memory used by commit batches)
// so we need to downgrade here
co_await delay(0, TaskPriority::ProxyCommit);
pContext->pProxyCommitData->lastVersionTime = pContext->startTime;
++pContext->pProxyCommitData->stats.commitBatchIn;
pContext->setupTraceBatch();
/////// Phase 1: Pre-resolution processing (CPU bound except waiting for a version # which is separately pipelined
/// and *should* be available by now (unless empty commit); ordered; currently atomic but could yield)
pContext->stage = PRE_RESOLUTION;
co_await CommitBatch::preresolutionProcessing(pContext);
if (pContext->rejected) {
pContext->pProxyCommitData->commitBatchesMemBytesCount -= pContext->currentBatchMemBytesCount;
co_return;
}
/////// Phase 2: Resolution (waiting on the network; pipelined)
pContext->stage = RESOLUTION;
co_await CommitBatch::getResolution(pContext);
////// Phase 3: Post-resolution processing (CPU bound except for very rare situations; ordered; currently atomic but
/// doesn't need to be)
pContext->stage = POST_RESOLUTION;
co_await CommitBatch::postResolution(pContext);
/////// Phase 4: Logging (network bound; pipelined up to MAX_READ_TRANSACTION_LIFE_VERSIONS (limited by loop above))
pContext->stage = TRANSACTION_LOGGING;
co_await CommitBatch::transactionLogging(pContext);
/////// Phase 5: Replies (CPU bound; no particular order required, though ordered execution would be best for
/// latency)
pContext->stage = REPLY;
co_await CommitBatch::reply(pContext);
pContext->stage = COMPLETE;
}
} // namespace CommitBatch
Future<Void> commitBatch(ProxyCommitData* pCommitData,
std::vector<CommitTransactionRequest>* trs,
int currentBatchMemBytesCount) {
CommitBatch::CommitBatchContext context(pCommitData, trs, currentBatchMemBytesCount);
Future<Void> commit = CommitBatch::commitBatchImpl(&context);
try {
co_await timeoutError(commit, SERVER_KNOBS->COMMIT_PROXY_LIVENESS_TIMEOUT);
} catch (Error& err) {
if (err.code() == error_code_actor_cancelled) {
throw;
}
TraceEvent(SevInfo, "CommitBatchFailed", pCommitData->dbgid)
.detail("Stage", context.stage)
.detail("ErrorCode", err.code());
throw failed_to_progress();
}
}
// Add tss mapping data to the reply, if any of the included storage servers have a TSS pair
void maybeAddTssMapping(GetKeyServerLocationsReply& reply,
ProxyCommitData* commitData,
std::unordered_set<UID>& included,
UID ssId) {
if (!included.contains(ssId)) {
auto mappingItr = commitData->tssMapping.find(ssId);
if (mappingItr != commitData->tssMapping.end()) {
reply.resultsTssMapping.push_back(*mappingItr);
}
included.insert(ssId);
}
}
void addTagMapping(GetKeyServerLocationsReply& reply, ProxyCommitData* commitData) {
for (const auto& [_, shard] : reply.results) {
for (auto& ssi : shard) {
auto iter = commitData->storageCache.find(ssi.id());
ASSERT_WE_THINK(iter != commitData->storageCache.end());
reply.resultsTagMapping.emplace_back(ssi.id(), iter->second->tag);
}
}
}
static Future<Void> doKeyServerLocationRequest(GetKeyServerLocationsRequest req, ProxyCommitData* commitData) {
// We can't respond to these requests until we have valid txnStateStore
getCurrentLineage()->modify(&TransactionLineage::operation) = TransactionLineage::Operation::GetKeyServersLocations;
getCurrentLineage()->modify(&TransactionLineage::txID) = req.spanContext.traceID;
co_await commitData->validState.getFuture();
co_await delay(0, TaskPriority::DefaultEndpoint);
std::unordered_set<UID> tssMappingsIncluded;
GetKeyServerLocationsReply rep;
if (!req.end.present()) {
auto r = req.reverse ? commitData->keyInfo.rangeContainingKeyBefore(req.begin)
: commitData->keyInfo.rangeContaining(req.begin);
std::vector<StorageServerInterface> ssis;
ssis.reserve(r.value().src_info.size());
for (auto& it : r.value().src_info) {
ssis.push_back(it->interf);
maybeAddTssMapping(rep, commitData, tssMappingsIncluded, it->interf.id());
}
rep.results.emplace_back(r.range(), ssis);
} else if (!req.reverse) {
int count = 0;
for (auto r = commitData->keyInfo.rangeContaining(req.begin);
r != commitData->keyInfo.ranges().end() && count < req.limit && r.begin() < req.end.get();
++r) {
std::vector<StorageServerInterface> ssis;
ssis.reserve(r.value().src_info.size());
for (auto& it : r.value().src_info) {
ssis.push_back(it->interf);
maybeAddTssMapping(rep, commitData, tssMappingsIncluded, it->interf.id());
}
rep.results.emplace_back(r.range(), ssis);
count++;
}
} else {
int count = 0;
auto r = commitData->keyInfo.rangeContainingKeyBefore(req.end.get());
while (count < req.limit && req.begin < r.end()) {
std::vector<StorageServerInterface> ssis;
ssis.reserve(r.value().src_info.size());
for (auto& it : r.value().src_info) {
ssis.push_back(it->interf);
maybeAddTssMapping(rep, commitData, tssMappingsIncluded, it->interf.id());
}
rep.results.emplace_back(r.range(), ssis);
if (r == commitData->keyInfo.ranges().begin()) {
break;
}
count++;
--r;
}
}
addTagMapping(rep, commitData);
req.reply.send(rep);
++commitData->stats.keyServerLocationOut;
}
static Future<Void> readRequestServer(CommitProxyInterface proxy,
PromiseStream<Future<Void>> addActor,
ProxyCommitData* commitData) {
while (true) {
GetKeyServerLocationsRequest req = co_await proxy.getKeyServersLocations.getFuture();
// WARNING: this code is run at a high priority, so it needs to do as little work as possible
if (req.limit != CLIENT_KNOBS->STORAGE_METRICS_SHARD_LIMIT && // Always do data distribution requests
(commitData->stats.keyServerLocationIn.getValue() - commitData->stats.keyServerLocationOut.getValue() >
SERVER_KNOBS->KEY_LOCATION_MAX_QUEUE_SIZE ||
(g_network->isSimulated() && buggify(0.001)))) {
++commitData->stats.keyServerLocationErrors;
req.reply.sendError(commit_proxy_memory_limit_exceeded());
TraceEvent(SevWarnAlways, "ProxyLocationRequestThresholdExceeded").suppressFor(60);
} else {
++commitData->stats.keyServerLocationIn;
addActor.send(doKeyServerLocationRequest(req, commitData));
}
}
}
static Future<Void> rejoinServer(CommitProxyInterface proxy, ProxyCommitData* commitData) {
// We can't respond to these requests until we have valid txnStateStore
co_await commitData->validState.getFuture();
TraceEvent("ProxyReadyForReads", proxy.id()).log();
while (true) {
GetStorageServerRejoinInfoRequest req = co_await proxy.getStorageServerRejoinInfo.getFuture();
if (commitData->txnStateStore->readValue(serverListKeyFor(req.id)).get().present()) {
GetStorageServerRejoinInfoReply rep;
rep.version = commitData->version.get();
rep.tag = decodeServerTagValue(commitData->txnStateStore->readValue(serverTagKeyFor(req.id)).get().get());
RangeResult history = commitData->txnStateStore->readRange(serverTagHistoryRangeFor(req.id)).get();
for (int i = history.size() - 1; i >= 0; i--) {
rep.history.push_back(
std::make_pair(decodeServerTagHistoryKey(history[i].key), decodeServerTagValue(history[i].value)));
}
auto localityKey = commitData->txnStateStore->readValue(tagLocalityListKeyFor(req.dcId)).get();
rep.newLocality = false;
if (localityKey.present()) {
int8_t locality = decodeTagLocalityListValue(localityKey.get());
if (locality != rep.tag.locality) {
TraceEvent(SevWarnAlways, "SSRejoinedWithChangedLocality")
.detail("Tag", rep.tag.toString())
.detail("DcId", req.dcId)
.detail("NewLocality", locality);
} else if (locality != rep.tag.locality) {
uint16_t tagId = 0;
std::vector<uint16_t> usedTags;
auto tagKeys = commitData->txnStateStore->readRange(serverTagKeys).get();
for (auto& kv : tagKeys) {
Tag t = decodeServerTagValue(kv.value);
if (t.locality == locality) {
usedTags.push_back(t.id);
}
}
auto historyKeys = commitData->txnStateStore->readRange(serverTagHistoryKeys).get();
for (auto& kv : historyKeys) {
Tag t = decodeServerTagValue(kv.value);
if (t.locality == locality) {
usedTags.push_back(t.id);
}
}
std::sort(usedTags.begin(), usedTags.end());
int usedIdx = 0;
for (; !usedTags.empty() && tagId <= usedTags.end()[-1]; tagId++) {
if (tagId < usedTags[usedIdx]) {
break;
} else {
usedIdx++;
}
}
rep.newTag = Tag(locality, tagId);
}
} else {
ASSERT_WE_THINK(rep.tag.locality != tagLocalityUpgraded);
TraceEvent(SevWarnAlways, "SSRejoinedWithUnknownLocality")
.detail("Tag", rep.tag.toString())
.detail("DcId", req.dcId);
}
req.reply.send(rep);
} else {
req.reply.sendError(worker_removed());
}
}
}
Future<Void> ddMetricsRequestServer(CommitProxyInterface proxy, Reference<AsyncVar<ServerDBInfo> const> db) {
while (true) {
GetDDMetricsRequest req = co_await proxy.getDDMetrics.getFuture();
if (!db->get().distributor.present()) {
req.reply.sendError(dd_not_found());
continue;
}
ErrorOr<GetDataDistributorMetricsReply> reply =
co_await errorOr(db->get().distributor.get().dataDistributorMetrics.getReply(
GetDataDistributorMetricsRequest(req.keys, req.shardLimit)));
if (reply.isError()) {
req.reply.sendError(reply.getError());
} else {
GetDDMetricsReply newReply;
newReply.storageMetricsList = reply.get().storageMetricsList;
req.reply.send(newReply);
}
}
}
Future<Void> monitorRemoteCommitted(ProxyCommitData* self) {
while (true) {
co_await delay(0); // allow this actor to be cancelled if we are removed after db changes.
Optional<std::vector<OptionalInterface<TLogInterface>>> remoteLogs;
if (self->db->get().recoveryState >= RecoveryState::ALL_LOGS_RECRUITED) {
for (auto& logSet : self->db->get().logSystemConfig.tLogs) {
if (!logSet.isLocal) {
remoteLogs = logSet.tLogs;
for (auto& tLog : logSet.tLogs) {
if (!tLog.present()) {
remoteLogs = Optional<std::vector<OptionalInterface<TLogInterface>>>();
break;
}
}
break;
}
}
}
if (!remoteLogs.present()) {
co_await self->db->onChange();
continue;
}
self->popRemoteTxs = true;
Future<Void> onChange = self->db->onChange();
while (true) {
std::vector<Future<TLogQueuingMetricsReply>> replies;
for (auto& it : remoteLogs.get()) {
replies.push_back(
brokenPromiseToNever(it.interf().getQueuingMetrics.getReply(TLogQueuingMetricsRequest())));
}
co_await (waitForAll(replies) || onChange);
if (onChange.isReady()) {
break;
}
// FIXME: use the configuration to calculate a more precise minimum recovery version.
Version minVersion = std::numeric_limits<Version>::max();
for (auto& it : replies) {
minVersion = std::min(minVersion, it.get().v);
}
while (!self->txsPopVersions.empty() && self->txsPopVersions.front().first <= minVersion) {
self->lastTxsPop = self->txsPopVersions.front().second;
self->logSystemConsumer->popTxs(self->txsPopVersions.front().second, tagLocalityRemoteLog);
self->txsPopVersions.pop_front();
}
co_await (delay(SERVER_KNOBS->UPDATE_REMOTE_LOG_VERSION_INTERVAL) || onChange);
if (onChange.isReady()) {
break;
}
}
}
}
Future<Void> proxySnapCreate(ProxySnapRequest snapReq, ProxyCommitData* commitData) {
TraceEvent("SnapCommitProxy_SnapReqEnter")
.detail("SnapPayload", snapReq.snapPayload)
.detail("SnapUID", snapReq.snapUID);
try {
// whitelist check
ExecCmdValueString execArg(snapReq.snapPayload);
StringRef binPath = execArg.getBinaryPath();
if (!isWhitelisted(commitData->whitelistedBinPathVec, binPath)) {
TraceEvent("SnapCommitProxy_WhiteListCheckFailed")
.detail("SnapPayload", snapReq.snapPayload)
.detail("SnapUID", snapReq.snapUID);
throw snap_path_not_whitelisted();
}
// db fully recovered check
if (commitData->db->get().recoveryState != RecoveryState::FULLY_RECOVERED) {
// Cluster is not fully recovered and needs TLogs
// from previous generation for full recovery.
// Currently, snapshot of old tlog generation is not
// supported and hence failing the snapshot request until
// cluster is fully_recovered.
TraceEvent("SnapCommitProxy_ClusterNotFullyRecovered")
.detail("SnapPayload", snapReq.snapPayload)
.detail("SnapUID", snapReq.snapUID);
throw snap_not_fully_recovered_unsupported();
}
auto result = commitData->txnStateStore->readValue("log_anti_quorum"_sr.withPrefix(configKeysPrefix)).get();
int logAntiQuorum = 0;
if (result.present()) {
logAntiQuorum = atoi(result.get().toString().c_str());
}
// FIXME: logAntiQuorum not supported, remove it later,
// In version2, we probably don't need this limitation, but this needs to be tested.
if (logAntiQuorum > 0) {
TraceEvent("SnapCommitProxy_LogAntiQuorumNotSupported")
.detail("SnapPayload", snapReq.snapPayload)
.detail("SnapUID", snapReq.snapUID);
throw snap_log_anti_quorum_unsupported();
}
int snapReqRetry = 0;
double snapRetryBackoff = FLOW_KNOBS->PREVENT_FAST_SPIN_DELAY;
while (true) {
// send a snap request to DD
if (!commitData->db->get().distributor.present()) {
TraceEvent(SevWarnAlways, "DataDistributorNotPresent").detail("Operation", "SnapRequest");
throw dd_not_found();
}
Error err;
try {
Future<ErrorOr<Void>> ddSnapReq =
commitData->db->get().distributor.get().distributorSnapReq.tryGetReply(
DistributorSnapRequest(snapReq.snapPayload, snapReq.snapUID));
co_await throwErrorOr(ddSnapReq);
break;
} catch (Error& e) {
err = e;
}
TraceEvent("SnapCommitProxy_DDSnapResponseError")
.errorUnsuppressed(err)
.detail("SnapPayload", snapReq.snapPayload)
.detail("SnapUID", snapReq.snapUID)
.detail("Retry", snapReqRetry);
// Retry if we have network issues
if (err.code() != error_code_request_maybe_delivered ||
++snapReqRetry > SERVER_KNOBS->SNAP_NETWORK_FAILURE_RETRY_LIMIT)
throw err;
co_await delay(snapRetryBackoff);
snapRetryBackoff = snapRetryBackoff * 2; // exponential backoff
}
snapReq.reply.send(Void());
} catch (Error& e) {
TraceEvent("SnapCommitProxy_SnapReqError")
.errorUnsuppressed(e)
.detail("SnapPayload", snapReq.snapPayload)
.detail("SnapUID", snapReq.snapUID);
if (e.code() != error_code_operation_cancelled) {
snapReq.reply.sendError(e);
} else {
throw e;
}
}
TraceEvent("SnapCommitProxy_SnapReqExit")
.detail("SnapPayload", snapReq.snapPayload)
.detail("SnapUID", snapReq.snapUID);
}
Future<Void> proxyCheckSafeExclusion(Reference<AsyncVar<ServerDBInfo> const> db, ExclusionSafetyCheckRequest req) {
TraceEvent("SafetyCheckCommitProxyBegin").log();
ExclusionSafetyCheckReply reply(false);
if (!db->get().distributor.present()) {
TraceEvent(SevWarnAlways, "DataDistributorNotPresent").detail("Operation", "ExclusionSafetyCheck");
req.reply.send(reply);
co_return;
}
try {
Future<ErrorOr<DistributorExclusionSafetyCheckReply>> ddSafeFuture =
db->get().distributor.get().distributorExclCheckReq.tryGetReply(
DistributorExclusionSafetyCheckRequest(req.exclusions));
DistributorExclusionSafetyCheckReply ddReply = co_await throwErrorOr(ddSafeFuture);
reply.safe = ddReply.safe;
} catch (Error& e) {
TraceEvent("SafetyCheckCommitProxyResponseError").error(e);
if (e.code() != error_code_operation_cancelled) {
req.reply.sendError(e);
co_return;
} else {
throw e;
}
}
TraceEvent("SafetyCheckCommitProxyFinish").log();
req.reply.send(reply);
}
class TxnTagCommitCostReporter {
UID myID;
Reference<AsyncVar<ServerDBInfo> const> db;
UIDTransactionTagMap<TransactionCommitCostEstimation>* ssTrTagCommitCost;
Future<Void> traceRatekeeperChanges() {
while (true) {
co_await db->onChange();
if (db->get().ratekeeper.present()) {
TraceEvent("ProxyRatekeeperChanged", myID).detail("RKID", db->get().ratekeeper.get().id());
} else {
TraceEvent("ProxyRatekeeperDied", myID).log();
}
}
}
Future<Void> reportCommitCosts() {
bool sendImmediately = db->get().ratekeeper.present();
while (true) {
if (!sendImmediately) {
co_await db->onChange();
}
sendImmediately = false;
if (!db->get().ratekeeper.present()) {
continue;
}
auto reply = brokenPromiseToNever(db->get().ratekeeper.get().reportCommitCostEstimation.getReply(
ReportCommitCostEstimationRequest(std::move(*ssTrTagCommitCost))));
ssTrTagCommitCost->clear();
auto dbChanged = db->onChange();
auto replyOrChange = co_await race(reply, dbChanged);
if (replyOrChange.index() == 1) {
sendImmediately = true;
continue;
}
co_await race(delay(SERVER_KNOBS->REPORT_TRANSACTION_COST_ESTIMATION_DELAY), dbChanged);
sendImmediately = true;
}
}
public:
TxnTagCommitCostReporter(UID myID,
Reference<AsyncVar<ServerDBInfo> const> db,
UIDTransactionTagMap<TransactionCommitCostEstimation>* ssTrTagCommitCost)
: myID(myID), db(db), ssTrTagCommitCost(ssTrTagCommitCost) {}
Future<Void> run() { co_await race(traceRatekeeperChanges(), reportCommitCosts()); }
};
namespace {
struct ExpireServerEntry {
int64_t timeReceived;
int expectedCount = 0;
int receivedCount = 0;
bool initialized = false;
};
struct IdempotencyKey {
Version version;
uint8_t highOrderBatchIndex;
bool operator==(const IdempotencyKey& other) const {
return version == other.version && highOrderBatchIndex == other.highOrderBatchIndex;
}
};
} // namespace
namespace std {
template <>
struct hash<IdempotencyKey> {
std::size_t operator()(const IdempotencyKey& key) const {
std::size_t seed = 0;
boost::hash_combine(seed, std::hash<Version>{}(key.version));
boost::hash_combine(seed, std::hash<uint8_t>{}(key.highOrderBatchIndex));
return seed;
}
};
} // namespace std
namespace {
class IdempotencyIdsExpireServer {
PublicRequestStream<ExpireIdempotencyIdRequest> expireIdempotencyId;
PromiseStream<ExpectedIdempotencyIdCountForKey> expectedIdempotencyIdCountForKey;
Standalone<VectorRef<MutationRef>>* idempotencyClears;
std::unordered_map<IdempotencyKey, ExpireServerEntry> idStatus;
void updateStatus(IdempotencyKey key, ExpireServerEntry* status) {
if (status->initialized) {
if (status->receivedCount == status->expectedCount) {
auto keyRange =
makeIdempotencySingleKeyRange(idempotencyClears->arena(), key.version, key.highOrderBatchIndex);
idempotencyClears->push_back(idempotencyClears->arena(),
MutationRef(MutationRef::ClearRange, keyRange.begin, keyRange.end));
idStatus.erase(key);
}
} else {
status->timeReceived = now();
status->initialized = true;
}
}
Future<Void> receiveExpireRequests() {
while (true) {
ExpireIdempotencyIdRequest req = co_await expireIdempotencyId.getFuture();
IdempotencyKey key{ req.commitVersion, req.batchIndexHighByte };
ExpireServerEntry* status = &idStatus[key];
status->receivedCount += 1;
CODE_PROBE(status->expectedCount == 0, "ExpireIdempotencyIdRequest received before count is known");
if (status->expectedCount > 0) {
ASSERT_LE(status->receivedCount, status->expectedCount);
}
updateStatus(key, status);
}
}
Future<Void> receiveExpectedCounts() {
while (true) {
ExpectedIdempotencyIdCountForKey req = co_await expectedIdempotencyIdCountForKey.getFuture();
IdempotencyKey key{ req.commitVersion, req.batchIndexHighByte };
ExpireServerEntry* status = &idStatus[key];
ASSERT_EQ(status->expectedCount, 0);
status->expectedCount = req.idempotencyIdCount;
updateStatus(key, status);
}
}
Future<Void> purgeOldEntries() {
while (true) {
co_await delay(SERVER_KNOBS->IDEMPOTENCY_ID_IN_MEMORY_LIFETIME);
int64_t purgeBefore = now() - SERVER_KNOBS->IDEMPOTENCY_ID_IN_MEMORY_LIFETIME;
std::vector<IdempotencyKey> keys;
keys.reserve(idStatus.size());
for (const auto& entry : idStatus) {
keys.push_back(entry.first);
}
std::vector<IdempotencyKey> keysToErase;
for (const auto& key : keys) {
co_await yield();
auto status = idStatus.find(key);
if (status != idStatus.end() && status->second.timeReceived < purgeBefore) {
keysToErase.push_back(key);
}
}
for (const auto& key : keysToErase) {
idStatus.erase(key);
}
}
}
public:
IdempotencyIdsExpireServer(PublicRequestStream<ExpireIdempotencyIdRequest> expireIdempotencyId,
PromiseStream<ExpectedIdempotencyIdCountForKey> expectedIdempotencyIdCountForKey,
Standalone<VectorRef<MutationRef>>* idempotencyClears)
: expireIdempotencyId(expireIdempotencyId), expectedIdempotencyIdCountForKey(expectedIdempotencyIdCountForKey),
idempotencyClears(idempotencyClears) {}
Future<Void> run() { co_await race(receiveExpireRequests(), receiveExpectedCounts(), purgeOldEntries()); }
};
TEST_CASE("/fdbserver/commitproxy/IdempotencyIdsExpireServer/ExpireBeforeExpectedCount") {
constexpr Version version = 100;
constexpr uint8_t firstBatch = 1;
constexpr uint8_t secondBatch = 2;
PublicRequestStream<ExpireIdempotencyIdRequest> expireRequests;
PromiseStream<ExpectedIdempotencyIdCountForKey> expectedCounts;
Standalone<VectorRef<MutationRef>> clears;
IdempotencyIdsExpireServer expireServer(expireRequests, expectedCounts, &clears);
Future<Void> server = expireServer.run();
// Expire requests can arrive before their expected counts; neither batch may be cleared yet.
expireRequests.send(ExpireIdempotencyIdRequest(version, firstBatch));
expireRequests.send(ExpireIdempotencyIdRequest(version, firstBatch));
expireRequests.send(ExpireIdempotencyIdRequest(version, secondBatch));
co_await yield();
ASSERT(!server.isReady());
ASSERT(clears.empty());
// The first batch is complete once counts arrive, while the second is still one request short.
expectedCounts.send(ExpectedIdempotencyIdCountForKey(version, 2, firstBatch));
expectedCounts.send(ExpectedIdempotencyIdCountForKey(version, 2, secondBatch));
co_await yield();
ASSERT(!server.isReady());
ASSERT_EQ(clears.size(), 1);
Arena expectedArena;
auto firstRange = makeIdempotencySingleKeyRange(expectedArena, version, firstBatch);
ASSERT_EQ(clears[0].type, MutationRef::ClearRange);
ASSERT(clears[0].param1 == firstRange.begin);
ASSERT(clears[0].param2 == firstRange.end);
// Completing the second batch must emit its distinct clear range.
expireRequests.send(ExpireIdempotencyIdRequest(version, secondBatch));
co_await yield();
ASSERT(!server.isReady());
ASSERT_EQ(clears.size(), 2);
auto secondRange = makeIdempotencySingleKeyRange(expectedArena, version, secondBatch);
ASSERT_EQ(clears[1].type, MutationRef::ClearRange);
ASSERT(clears[1].param1 == secondRange.begin);
ASSERT(clears[1].param2 == secondRange.end);
server.cancel();
}
struct TransactionStateResolveContext {
// Maximum sequence for txnStateRequest, this is defined when the request last flag is set.
Sequence maxSequence = std::numeric_limits<Sequence>::max();
// Flags marks received transaction state requests, we only process the transaction request when *all* requests are
// received.
std::unordered_set<Sequence> receivedSequences;
ProxyCommitData* pCommitData = nullptr;
// Pointer to transaction state store, shortcut for commitData.txnStateStore
IKeyValueStore* pTxnStateStore = nullptr;
Future<Void> txnRecovery;
// Actor streams
PromiseStream<Future<Void>>* pActors = nullptr;
// Flag reports if the transaction state request is complete. This request should only happen during recover, i.e.
// once per commit proxy.
bool processed = false;
TransactionStateResolveContext() = default;
TransactionStateResolveContext(ProxyCommitData* pCommitData_, PromiseStream<Future<Void>>* pActors_)
: pCommitData(pCommitData_), pTxnStateStore(pCommitData_->txnStateStore), pActors(pActors_) {
ASSERT(pTxnStateStore != nullptr);
}
};
Future<Void> processCompleteTransactionStateRequest(TransactionStateResolveContext* pContext) {
KeyRange txnKeys = allKeys;
std::map<Tag, UID> tag_uid;
RangeResult UIDtoTagMap = pContext->pTxnStateStore->readRange(serverTagKeys).get();
for (const KeyValueRef& kv : UIDtoTagMap) {
tag_uid[decodeServerTagValue(kv.value)] = decodeServerTagKey(kv.key);
}
while (true) {
co_await yield();
RangeResult data =
pContext->pTxnStateStore
->readRange(txnKeys, SERVER_KNOBS->BUGGIFIED_ROW_LIMIT, SERVER_KNOBS->APPLY_MUTATION_BYTES)
.get();
if (data.empty())
break;
((KeyRangeRef&)txnKeys) = KeyRangeRef(keyAfter(data.back().key, txnKeys.arena()), txnKeys.end);
Standalone<VectorRef<MutationRef>> mutations;
std::vector<std::pair<MapPair<Key, ServerCacheInfo>, int>> keyInfoData;
std::vector<UID> src, dest;
ServerCacheInfo info;
auto updateTagInfo = [pContext = pContext](const std::vector<UID>& uids,
std::vector<Tag>& tags,
std::vector<Reference<StorageInfo>>& storageInfoItems) {
for (const auto& id : uids) {
auto storageInfo = getStorageInfo(id, &pContext->pCommitData->storageCache, pContext->pTxnStateStore);
ASSERT(storageInfo->tag != invalidTag);
tags.push_back(storageInfo->tag);
storageInfoItems.push_back(storageInfo);
}
};
for (auto& kv : data) {
if (kv.key.startsWith(keyServersPrefix)) {
KeyRef k = kv.key.removePrefix(keyServersPrefix);
if (k == allKeys.end) {
continue;
}
decodeKeyServersValue(tag_uid, kv.value, src, dest);
info.tags.clear();
info.src_info.clear();
updateTagInfo(src, info.tags, info.src_info);
info.dest_info.clear();
updateTagInfo(dest, info.tags, info.dest_info);
uniquify(info.tags);
keyInfoData.emplace_back(MapPair<Key, ServerCacheInfo>(k, info), 1);
} else if (kv.key.startsWith(rangeLockPrefix)) {
if (pContext->pCommitData->rangeLockEnabled()) {
ASSERT(pContext->pCommitData->rangeLock != nullptr);
Key keyInsert = kv.key.removePrefix(rangeLockPrefix);
pContext->pCommitData->rangeLock->initKeyPoint(keyInsert, kv.value);
}
} else {
mutations.emplace_back(mutations.arena(), MutationRef::SetValue, kv.key, kv.value);
continue;
}
}
// insert keyTag data separately from metadata mutations so that we can do one bulk insert which
// avoids a lot of map lookups.
pContext->pCommitData->keyInfo.rawInsert(keyInfoData);
Arena arena;
bool confChanges;
applyMetadataMutations(SpanContext(),
pContext->pCommitData->getApplyMetadataProxyContext(),
arena,
Reference<LogSystemConsumer>(),
mutations,
/* pToCommit= */ nullptr,
confChanges,
/* version= */ 0,
/* popVersion= */ 0,
/* initialCommit= */ true,
/* provisionalCommitProxy */ pContext->pCommitData->provisional);
}
auto lockedKey = pContext->pTxnStateStore->readValue(databaseLockedKey).get();
pContext->pCommitData->locked = lockedKey.present() && !lockedKey.get().empty();
pContext->pCommitData->metadataVersion = pContext->pTxnStateStore->readValue(metadataVersionKey).get();
pContext->pCommitData->cdcRouting.reload(pContext->pTxnStateStore);
pContext->pTxnStateStore->enableSnapshot();
}
Future<Void> processTransactionStateRequestPart(TransactionStateResolveContext* pContext, TxnStateRequest request) {
ASSERT(pContext->pCommitData != nullptr);
ASSERT(pContext->pActors != nullptr);
if (pContext->receivedSequences.contains(request.sequence)) {
if (pContext->receivedSequences.size() == pContext->maxSequence) {
co_await pContext->txnRecovery;
}
// This part is already received. Still we will re-broadcast it to other CommitProxies
pContext->pActors->send(broadcastTxnRequest(request, SERVER_KNOBS->TXN_STATE_SEND_AMOUNT, true));
co_await yield();
co_return;
}
if (request.last) {
// This is the last piece of subsequence, yet other pieces might still on the way.
pContext->maxSequence = request.sequence + 1;
}
pContext->receivedSequences.insert(request.sequence);
// Although we may receive the CommitTransactionRequest for the recovery transaction before all of the
// TxnStateRequest, we will not get a resolution result from any resolver until the master has submitted its initial
// (sequence 0) resolution request, which it doesn't do until we have acknowledged all TxnStateRequests
ASSERT(!pContext->pCommitData->validState.isSet());
for (auto& kv : request.data) {
pContext->pTxnStateStore->set(kv, &request.arena);
}
pContext->pTxnStateStore->commit(true);
if (pContext->receivedSequences.size() == pContext->maxSequence) {
// Received all components of the txnStateRequest
ASSERT(!pContext->processed);
pContext->txnRecovery = processCompleteTransactionStateRequest(pContext);
co_await pContext->txnRecovery;
pContext->processed = true;
}
pContext->pActors->send(broadcastTxnRequest(request, SERVER_KNOBS->TXN_STATE_SEND_AMOUNT, true));
co_await yield();
}
} // anonymous namespace
//
// Metrics related to the commit proxy are logged on a five second interval in
// the `ProxyMetrics` trace. However, it can be hard to determine workload
// burstiness when looking at such a large time range. This function adds much
// more frequent logging for certain metrics to provide fine-grained insight
// into workload patterns. The metrics logged by this function break down into
// two categories:
//
// * existing counters reported by `ProxyMetrics`
// * new counters that are only reported by this function
//
// Neither is implemented optimally, but the data collected should be helpful
// in identifying workload patterns on the server.
//
// Metrics reporting by this function can be disabled by setting the
// `BURSTINESS_METRICS_ENABLED` knob to false. The reporting interval can be
// adjusted by modifying the knob `BURSTINESS_METRICS_LOG_INTERVAL`.
//
Future<Void> logDetailedMetrics(ProxyCommitData* commitData) {
while (true) {
if (!SERVER_KNOBS->BURSTINESS_METRICS_ENABLED) {
co_return;
}
double startTime = now();
int64_t commitBatchInBaseline = commitData->stats.commitBatchIn.getValue();
int64_t commitBatchFlushByteLimitBaseline = commitData->stats.commitBatchFlushByteLimit.getValue();
int64_t commitBatchFlushCountLimitBaseline = commitData->stats.commitBatchFlushCountLimit.getValue();
int64_t commitBatchFlushTimeoutBaseline = commitData->stats.commitBatchFlushTimeout.getValue();
int64_t commitBatchFlushFirstInBatchBaseline = commitData->stats.commitBatchFlushFirstInBatch.getValue();
int64_t commitBatchFlushTransactionSizeLimitBaseline =
commitData->stats.commitBatchFlushTransactionSizeLimit.getValue();
int64_t txnCommitInBaseline = commitData->stats.txnCommitIn.getValue();
int64_t mutationsBaseline = commitData->stats.mutations.getValue();
int64_t mutationBytesBaseline = commitData->stats.mutationBytes.getValue();
co_await delay(SERVER_KNOBS->BURSTINESS_METRICS_LOG_INTERVAL);
int64_t commitBatchInReal = commitData->stats.commitBatchIn.getValue();
int64_t commitBatchFlushByteLimitReal = commitData->stats.commitBatchFlushByteLimit.getValue();
int64_t commitBatchFlushCountLimitReal = commitData->stats.commitBatchFlushCountLimit.getValue();
int64_t commitBatchFlushTimeoutReal = commitData->stats.commitBatchFlushTimeout.getValue();
int64_t commitBatchFlushFirstInBatchReal = commitData->stats.commitBatchFlushFirstInBatch.getValue();
int64_t commitBatchFlushTransactionSizeLimitReal =
commitData->stats.commitBatchFlushTransactionSizeLimit.getValue();
int64_t txnCommitInReal = commitData->stats.txnCommitIn.getValue();
int64_t mutationsReal = commitData->stats.mutations.getValue();
int64_t mutationBytesReal = commitData->stats.mutationBytes.getValue();
// Don't log anything if any of the counters got reset during the wait
// interval. Assume that typically all the counters get reset at once.
if (commitBatchInReal < commitBatchInBaseline ||
commitBatchFlushByteLimitReal < commitBatchFlushByteLimitBaseline ||
commitBatchFlushCountLimitReal < commitBatchFlushCountLimitBaseline ||
commitBatchFlushTimeoutReal < commitBatchFlushTimeoutBaseline ||
commitBatchFlushFirstInBatchReal < commitBatchFlushFirstInBatchBaseline ||
commitBatchFlushTransactionSizeLimitReal < commitBatchFlushTransactionSizeLimitBaseline ||
txnCommitInReal < txnCommitInBaseline || mutationsReal < mutationsBaseline ||
mutationBytesReal < mutationBytesBaseline) {
continue;
}
TraceEvent("ProxyDetailedMetrics")
.detail("Elapsed", now() - startTime)
.detail("CommitBatchIn", commitBatchInReal - commitBatchInBaseline)
.detail("CommitBatchFlushByteLimit", commitBatchFlushByteLimitReal - commitBatchFlushByteLimitBaseline)
.detail("CommitBatchFlushCountLimit", commitBatchFlushCountLimitReal - commitBatchFlushCountLimitBaseline)
.detail("CommitBatchFlushTimeout", commitBatchFlushTimeoutReal - commitBatchFlushTimeoutBaseline)
.detail("CommitBatchFlushFirstInBatch",
commitBatchFlushFirstInBatchReal - commitBatchFlushFirstInBatchBaseline)
.detail("CommitBatchFlushTransactionSizeLimit",
commitBatchFlushTransactionSizeLimitReal - commitBatchFlushTransactionSizeLimitBaseline)
.detail("TxnCommitIn", txnCommitInReal - txnCommitInBaseline)
.detail("Mutations", mutationsReal - mutationsBaseline)
.detail("MutationBytes", mutationBytesReal - mutationBytesBaseline)
.detail("UniqueClients", commitData->stats.getSizeAndResetUniqueClients());
}
}
class CommitProxyServerCore {
CommitProxyInterface proxy;
LifetimeToken masterLifetime;
Reference<AsyncVar<ServerDBInfo> const> db;
bool firstProxy;
std::string whitelistBinPaths;
ProxyCommitData commitData;
PromiseStream<std::pair<std::vector<CommitTransactionRequest>, int>> batchedCommits;
Future<Void> commitBatcherActor;
Future<Void> lastCommitComplete = Void();
TxnTagCommitCostReporter txnTagCommitCostReporter;
IdempotencyIdsExpireServer idempotencyIdsExpireServer;
PromiseStream<Future<Void>> addActor;
Future<Void> onError;
Future<Void> dbInfoChange;
TransactionStateResolveContext transactionStateResolveContext;
Future<Void> serveDbInfoChanges() {
while (true) {
co_await dbInfoChange;
dbInfoChange = commitData.db->onChange();
if (masterLifetime.isEqual(commitData.db->get().masterLifetime) &&
commitData.db->get().recoveryState >= RecoveryState::RECOVERY_TRANSACTION) {
commitData.logSystem = makeLogSystemFromServerDBInfo(proxy.id(), commitData.db->get(), false, addActor);
commitData.logSystemConsumer = commitData.logSystem->makeConsumer();
for (auto it : commitData.tag_popped) {
commitData.logSystemConsumer->pop(it.second, it.first);
}
commitData.logSystemConsumer->popTxs(commitData.lastTxsPop, tagLocalityRemoteLog);
}
commitData.updateLatencyBandConfig(commitData.db->get().latencyBandConfig);
}
}
Future<Void> serveActorErrors() { co_await onError; }
Future<Void> serveBatchedCommits() {
while (true) {
std::pair<std::vector<CommitTransactionRequest>, int> batchedRequests = co_await batchedCommits.getFuture();
// WARNING: this code is run at a high priority, so it needs to do as little work as possible
/*
TraceEvent("CommitProxyCTR", proxy.id())
.detail("CommitTransactions", trs.size())
.detail("TransactionRate", transactionRate)
.detail("TransactionQueue", transactionQueue.size())
.detail("ReleasedTransactionCount", transactionCount);
TraceEvent("CommitProxyCore", commitData.dbgid)
.detail("TxSize", trs.size())
.detail("MasterLifetime", masterLifetime.toString())
.detail("DbMasterLifetime", commitData.db->get().masterLifetime.toString())
.detail("RecoveryState", commitData.db->get().recoveryState)
.detail("CCInf", commitData.db->get().clusterInterface.id().toString());
*/
const std::vector<CommitTransactionRequest>& trs = batchedRequests.first;
const int batchBytes = batchedRequests.second;
if (!trs.empty() ||
(commitData.db->get().recoveryState >= RecoveryState::ACCEPTING_COMMITS &&
masterLifetime.isEqual(commitData.db->get().masterLifetime) && lastCommitComplete.isReady())) {
lastCommitComplete =
commitBatch(&commitData,
const_cast<std::vector<CommitTransactionRequest>*>(&batchedRequests.first),
batchBytes);
addActor.send(lastCommitComplete);
}
}
}
Future<Void> serveProxySnapRequests() {
while (true) {
ProxySnapRequest snapReq = co_await proxy.proxySnapReq.getFuture();
TraceEvent(SevDebug, "SnapMasterEnqueue").log();
addActor.send(proxySnapCreate(snapReq, &commitData));
}
}
Future<Void> serveExclusionSafetyCheckRequests() {
while (true) {
ExclusionSafetyCheckRequest exclCheckReq = co_await proxy.exclusionSafetyCheckReq.getFuture();
addActor.send(proxyCheckSafeExclusion(db, exclCheckReq));
}
}
Future<Void> serveTxnStateRequests() {
while (true) {
TxnStateRequest request = co_await proxy.txnState.getFuture();
addActor.send(processTransactionStateRequestPart(&transactionStateResolveContext, request));
}
}
Future<Void> serveSetThrottledShardRequests() {
while (true) {
SetThrottledShardRequest request = co_await proxy.setThrottledShard.getFuture();
for (auto& shard : request.throttledShards) {
auto it = commitData.hotShards.begin();
for (; it != commitData.hotShards.end(); ++it) {
if (it->first == shard) {
it->second = request.expirationTime;
break;
}
}
if (it == commitData.hotShards.end()) {
commitData.hotShards.emplace_back(std::make_pair(shard, request.expirationTime));
}
}
// TraceEvent(SevDebug, "ReceivedSetThrottledShards").detail("NumHotShards", commitData.hotShards.size());
}
}
public:
CommitProxyServerCore(CommitProxyInterface proxy,
MasterInterface master,
LifetimeToken masterLifetime,
Reference<AsyncVar<ServerDBInfo> const> db,
LogEpoch epoch,
Version recoveryTransactionVersion,
bool firstProxy,
std::string whitelistBinPaths,
bool provisional,
uint16_t commitProxyIndex)
: proxy(proxy), masterLifetime(masterLifetime), db(db), firstProxy(firstProxy),
whitelistBinPaths(std::move(whitelistBinPaths)), commitData(proxy.id(),
master,
recoveryTransactionVersion,
proxy.commit,
db,
firstProxy,
provisional,
commitProxyIndex,
epoch),
txnTagCommitCostReporter(proxy.id(), db, &commitData.ssTrTagCommitCost),
idempotencyIdsExpireServer(proxy.expireIdempotencyId,
commitData.expectedIdempotencyIdCountForKey,
&commitData.idempotencyClears),
onError(transformError(actorCollection(addActor.getFuture()), broken_promise(), tlog_failed())) {}
Future<Void> run() {
addActor.send(waitFailureServer(proxy.waitFailure.getFuture()));
addActor.send(traceRole(Role::COMMIT_PROXY, proxy.id()));
//TraceEvent("CommitProxyInit1", proxy.id());
// Wait until we can load the "real" logsystem, since we don't support switching them currently
while (!(masterLifetime.isEqual(commitData.db->get().masterLifetime) &&
commitData.db->get().recoveryState >= RecoveryState::RECOVERY_TRANSACTION)) {
//TraceEvent("ProxyInit2", proxy.id()).detail("LSEpoch", db->get().logSystemConfig.epoch).detail("Need", epoch);
co_await commitData.db->onChange();
}
dbInfoChange = commitData.db->onChange();
//TraceEvent("ProxyInit3", proxy.id());
commitData.resolvers = commitData.db->get().resolvers;
commitData.localTLogCount = commitData.db->get().logSystemConfig.numLogs();
ASSERT(!commitData.resolvers.empty());
for (int i = 0; i < commitData.resolvers.size(); ++i) {
commitData.stats.resolverDist.push_back(
Histogram::getHistogram("CommitProxy"_sr,
"ToResolver_" + commitData.resolvers[i].id().toString(),
Histogram::Unit::milliseconds));
}
// Initialize keyResolvers map
auto rs =
commitData.keyResolvers.modify(SERVER_KNOBS->PROXY_USE_RESOLVER_PRIVATE_MUTATIONS ? normalKeys : allKeys);
for (auto r = rs.begin(); r != rs.end(); ++r)
r->value().emplace_back(0, 0);
commitData.systemKeyVersions.push_back(0);
commitData.logSystem = makeLogSystemFromServerDBInfo(proxy.id(), commitData.db->get(), false, addActor);
commitData.logSystemConsumer = commitData.logSystem->makeConsumer();
commitData.logAdapter =
new LogSystemDiskQueueAdapter(commitData.logSystem, Reference<AsyncVar<PeekTxsInfo>>(), 1, false);
commitData.txnStateStore = keyValueStoreLogSystem(commitData.logAdapter,
commitData.db,
proxy.id(),
2e9,
DisableSnapshot::True,
ReplaceContent::True,
ExactRecovery::True);
createWhitelistBinPathVec(whitelistBinPaths, commitData.whitelistedBinPathVec);
commitData.updateLatencyBandConfig(commitData.db->get().latencyBandConfig);
// ((SERVER_MEM_LIMIT * COMMIT_BATCHES_MEM_FRACTION_OF_TOTAL) / COMMIT_BATCHES_MEM_TO_TOTAL_MEM_SCALE_FACTOR) is
// only a approximate formula for limiting the memory used. COMMIT_BATCHES_MEM_TO_TOTAL_MEM_SCALE_FACTOR is an
// estimate based on experiments and not an accurate one.
int64_t commitBatchesMemoryLimit = SERVER_KNOBS->COMMIT_BATCHES_MEM_BYTES_HARD_LIMIT;
if (SERVER_KNOBS->SERVER_MEM_LIMIT > 0) {
commitBatchesMemoryLimit =
std::min(commitBatchesMemoryLimit,
static_cast<int64_t>(
(SERVER_KNOBS->SERVER_MEM_LIMIT * SERVER_KNOBS->COMMIT_BATCHES_MEM_FRACTION_OF_TOTAL) /
SERVER_KNOBS->COMMIT_BATCHES_MEM_TO_TOTAL_MEM_SCALE_FACTOR));
}
TraceEvent(SevInfo, "CommitBatchesMemoryLimit").detail("BytesLimit", commitBatchesMemoryLimit);
// Initialize RangeLock
if (commitData.rangeLockEnabled()) {
commitData.rangeLock = std::make_shared<RangeLock>(&commitData);
TraceEvent(SevInfo, "CommitProxyRangeLockEnabled", commitData.dbgid);
}
addActor.send(monitorRemoteCommitted(&commitData));
addActor.send(readRequestServer(proxy, addActor, &commitData));
addActor.send(rejoinServer(proxy, &commitData));
addActor.send(ddMetricsRequestServer(proxy, db));
addActor.send(txnTagCommitCostReporter.run());
addActor.send(logDetailedMetrics(&commitData));
auto openDb = openDBOnServer(db);
if (firstProxy) {
addActor.send(recurringAsync(
[openDb = openDb]() {
return cleanIdempotencyIds(openDb, SERVER_KNOBS->IDEMPOTENCY_IDS_MIN_AGE_SECONDS);
},
SERVER_KNOBS->IDEMPOTENCY_IDS_CLEANER_POLLING_INTERVAL,
true,
SERVER_KNOBS->IDEMPOTENCY_IDS_CLEANER_POLLING_INTERVAL));
}
addActor.send(idempotencyIdsExpireServer.run());
// wait for txnStateStore recovery
co_await commitData.txnStateStore->readValue(StringRef());
int commitBatchByteLimit =
(int)std::min<double>(SERVER_KNOBS->COMMIT_TRANSACTION_BATCH_BYTES_MAX,
std::max<double>(SERVER_KNOBS->COMMIT_TRANSACTION_BATCH_BYTES_MIN,
SERVER_KNOBS->COMMIT_TRANSACTION_BATCH_BYTES_SCALE_BASE *
pow(commitData.db->get().client.commitProxies.size(),
SERVER_KNOBS->COMMIT_TRANSACTION_BATCH_BYTES_SCALE_POWER)));
commitBatcherActor = commitBatcher(
&commitData, batchedCommits, proxy.commit.getFuture(), commitBatchByteLimit, commitBatchesMemoryLimit);
// This has to be initialized after the commitData.txnStateStore gets initialized.
transactionStateResolveContext = TransactionStateResolveContext(&commitData, &addActor);
co_await race(serveDbInfoChanges(),
serveActorErrors(),
serveBatchedCommits(),
serveProxySnapRequests(),
serveExclusionSafetyCheckRequests(),
serveTxnStateRequests(),
serveSetThrottledShardRequests());
}
};
// only update the local Db info if the CP is not removed
Future<Void> updateLocalDbInfo(Reference<AsyncVar<ServerDBInfo> const> in,
Reference<AsyncVar<ServerDBInfo>> out,
uint64_t recoveryCount,
CommitProxyInterface myInterface) {
// whether this CP already receive the db info including itself
bool firstValidDbInfo = false;
while (true) {
bool isIncluded =
std::count(in->get().client.commitProxies.begin(), in->get().client.commitProxies.end(), myInterface);
if (in->get().recoveryCount >= recoveryCount && !isIncluded) {
throw worker_removed();
}
if (isIncluded) {
firstValidDbInfo = true;
}
// only update the db info if this is the current CP, or before we received first one including current CP.
// Several db infos at the beginning just contain the provisional CP
if (isIncluded || !firstValidDbInfo) {
DisabledTraceEvent("UpdateLocalDbInfo", myInterface.id())
.detail("Provisional", myInterface.provisional)
.detail("Included", isIncluded)
.detail("FirstValid", firstValidDbInfo)
.detail("ReceivedRC", in->get().recoveryCount)
.detail("RecoveryCount", recoveryCount);
if (in->get().recoveryCount >= out->get().recoveryCount) {
out->set(in->get());
}
}
co_await in->onChange();
}
}
Future<Void> commitProxyServer(CommitProxyInterface proxy,
InitializeCommitProxyRequest req,
Reference<AsyncVar<ServerDBInfo> const> db,
std::string whitelistBinPaths) {
try {
auto localDb = makeReference<AsyncVar<ServerDBInfo>>();
CommitProxyServerCore core(proxy,
req.master,
req.masterLifetime,
localDb,
req.recoveryCount,
req.recoveryTransactionVersion,
req.firstProxy,
std::move(whitelistBinPaths),
proxy.provisional,
req.commitProxyIndex);
co_await race(core.run(), updateLocalDbInfo(db, localDb, req.recoveryCount, proxy));
} catch (Error& e) {
Severity sev = e.code() == error_code_failed_to_progress ? SevWarnAlways : SevInfo;
TraceEvent(sev, "CommitProxyTerminated", proxy.id()).errorUnsuppressed(e);
if (e.code() != error_code_worker_removed && e.code() != error_code_tlog_stopped &&
e.code() != error_code_tlog_failed && e.code() != error_code_coordinators_changed &&
e.code() != error_code_coordinated_state_conflict && e.code() != error_code_new_coordinators_timed_out &&
e.code() != error_code_failed_to_progress) {
throw;
}
CODE_PROBE(e.code() == error_code_failed_to_progress, "Commit proxy failed to progress");
}
}
void forceLinkCommitProxyTests() {}