2021-05-18 04:02:55 +08:00
|
|
|
/*
|
|
|
|
* PaxosConfigTransaction.actor.cpp
|
|
|
|
*
|
|
|
|
* This source file is part of the FoundationDB open source project
|
|
|
|
*
|
2022-03-22 04:36:23 +08:00
|
|
|
* Copyright 2013-2022 Apple Inc. and the FoundationDB project authors
|
2021-05-18 04:02:55 +08:00
|
|
|
*
|
|
|
|
* 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.
|
|
|
|
*/
|
|
|
|
|
2021-07-19 05:02:45 +08:00
|
|
|
#include "fdbclient/DatabaseContext.h"
|
2022-06-08 05:49:02 +08:00
|
|
|
#include "fdbclient/MonitorLeader.h"
|
2021-05-18 04:02:55 +08:00
|
|
|
#include "fdbclient/PaxosConfigTransaction.h"
|
|
|
|
#include "flow/actorcompiler.h" // must be last include
|
|
|
|
|
2022-02-23 02:40:36 +08:00
|
|
|
using ConfigTransactionInfo = ModelInterface<ConfigTransactionInterface>;
|
|
|
|
|
2021-08-21 06:53:13 +08:00
|
|
|
class CommitQuorum {
|
2021-08-26 07:28:21 +08:00
|
|
|
ActorCollection actors{ false };
|
2021-08-21 06:53:13 +08:00
|
|
|
std::vector<ConfigTransactionInterface> ctis;
|
|
|
|
size_t failed{ 0 };
|
|
|
|
size_t successful{ 0 };
|
|
|
|
size_t maybeCommitted{ 0 };
|
|
|
|
Promise<Void> result;
|
2021-08-26 06:21:36 +08:00
|
|
|
Standalone<VectorRef<ConfigMutationRef>> mutations;
|
|
|
|
ConfigCommitAnnotation annotation;
|
|
|
|
|
2022-09-14 06:43:56 +08:00
|
|
|
ConfigTransactionCommitRequest getCommitRequest(ConfigGeneration generation,
|
|
|
|
CoordinatorsHash coordinatorsHash) const {
|
2022-06-08 05:49:02 +08:00
|
|
|
return ConfigTransactionCommitRequest(coordinatorsHash, generation, mutations, annotation);
|
2021-08-26 06:21:36 +08:00
|
|
|
}
|
2021-08-21 06:53:13 +08:00
|
|
|
|
|
|
|
void updateResult() {
|
2021-08-26 07:28:21 +08:00
|
|
|
if (successful >= ctis.size() / 2 + 1 && result.canBeSet()) {
|
2021-08-21 06:53:13 +08:00
|
|
|
result.send(Void());
|
2021-08-26 07:28:21 +08:00
|
|
|
} else if (failed >= ctis.size() / 2 + 1 && result.canBeSet()) {
|
2021-10-19 03:02:39 +08:00
|
|
|
// Rollforwards could cause a version that didn't have quorum to
|
2021-10-20 08:36:43 +08:00
|
|
|
// commit, so send commit_unknown_result instead of commit_failed.
|
2022-05-26 03:08:30 +08:00
|
|
|
|
|
|
|
// Calling sendError could delete this
|
|
|
|
auto local = this->result;
|
|
|
|
local.sendError(commit_unknown_result());
|
2021-08-21 06:53:13 +08:00
|
|
|
} else {
|
|
|
|
// Check if it is possible to ever receive quorum agreement
|
|
|
|
auto totalRequestsOutstanding = ctis.size() - (failed + successful + maybeCommitted);
|
|
|
|
if ((failed + totalRequestsOutstanding < ctis.size() / 2 + 1) &&
|
2021-08-26 07:28:21 +08:00
|
|
|
(successful + totalRequestsOutstanding < ctis.size() / 2 + 1) && result.canBeSet()) {
|
2022-05-26 03:08:30 +08:00
|
|
|
// Calling sendError could delete this
|
|
|
|
auto local = this->result;
|
|
|
|
local.sendError(commit_unknown_result());
|
2021-08-21 06:53:13 +08:00
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2021-08-26 06:21:36 +08:00
|
|
|
ACTOR static Future<Void> addRequestActor(CommitQuorum* self,
|
|
|
|
ConfigGeneration generation,
|
2022-09-14 06:43:56 +08:00
|
|
|
CoordinatorsHash coordinatorsHash,
|
2021-08-26 06:21:36 +08:00
|
|
|
ConfigTransactionInterface cti) {
|
2021-08-21 06:53:13 +08:00
|
|
|
try {
|
2022-04-28 12:54:13 +08:00
|
|
|
if (cti.hostname.present()) {
|
2022-06-08 05:49:02 +08:00
|
|
|
wait(timeoutError(retryGetReplyFromHostname(self->getCommitRequest(generation, coordinatorsHash),
|
|
|
|
cti.hostname.get(),
|
|
|
|
WLTOKEN_CONFIGTXN_COMMIT),
|
2022-04-28 12:54:13 +08:00
|
|
|
CLIENT_KNOBS->COMMIT_QUORUM_TIMEOUT));
|
|
|
|
} else {
|
2022-06-08 05:49:02 +08:00
|
|
|
wait(timeoutError(cti.commit.getReply(self->getCommitRequest(generation, coordinatorsHash)),
|
2022-04-28 12:54:13 +08:00
|
|
|
CLIENT_KNOBS->COMMIT_QUORUM_TIMEOUT));
|
|
|
|
}
|
2021-08-21 06:53:13 +08:00
|
|
|
++self->successful;
|
|
|
|
} catch (Error& e) {
|
2022-02-02 14:27:12 +08:00
|
|
|
// self might be destroyed if this actor is cancelled
|
2021-10-12 03:17:09 +08:00
|
|
|
if (e.code() == error_code_actor_cancelled) {
|
|
|
|
throw;
|
|
|
|
}
|
|
|
|
|
2022-02-02 14:27:12 +08:00
|
|
|
if (e.code() == error_code_not_committed || e.code() == error_code_timed_out) {
|
2021-08-21 06:53:13 +08:00
|
|
|
++self->failed;
|
2021-08-27 15:44:12 +08:00
|
|
|
} else {
|
|
|
|
++self->maybeCommitted;
|
2021-08-21 06:53:13 +08:00
|
|
|
}
|
|
|
|
}
|
|
|
|
self->updateResult();
|
|
|
|
return Void();
|
|
|
|
}
|
|
|
|
|
|
|
|
public:
|
|
|
|
CommitQuorum() = default;
|
2021-08-26 07:28:21 +08:00
|
|
|
explicit CommitQuorum(std::vector<ConfigTransactionInterface> const& ctis) : ctis(ctis) {}
|
2021-08-26 06:21:36 +08:00
|
|
|
void set(KeyRef key, ValueRef value) {
|
|
|
|
if (key == configTransactionDescriptionKey) {
|
|
|
|
annotation.description = ValueRef(annotation.arena(), value);
|
|
|
|
} else {
|
2021-08-26 12:28:36 +08:00
|
|
|
mutations.push_back_deep(mutations.arena(),
|
|
|
|
IKnobCollection::createSetMutation(mutations.arena(), key, value));
|
2021-08-26 06:21:36 +08:00
|
|
|
}
|
|
|
|
}
|
|
|
|
void clear(KeyRef key) {
|
|
|
|
if (key == configTransactionDescriptionKey) {
|
|
|
|
annotation.description = ""_sr;
|
|
|
|
} else {
|
2021-08-26 12:28:36 +08:00
|
|
|
mutations.push_back_deep(mutations.arena(), IKnobCollection::createClearMutation(mutations.arena(), key));
|
2021-08-26 06:21:36 +08:00
|
|
|
}
|
|
|
|
}
|
|
|
|
void setTimestamp() { annotation.timestamp = now(); }
|
|
|
|
size_t expectedSize() const { return annotation.expectedSize() + mutations.expectedSize(); }
|
2022-09-14 06:43:56 +08:00
|
|
|
Future<Void> commit(ConfigGeneration generation, CoordinatorsHash coordinatorsHash) {
|
2021-08-21 06:53:13 +08:00
|
|
|
// Send commit message to all replicas, even those that did not return the used replica.
|
|
|
|
// This way, slow replicas are kept up date.
|
|
|
|
for (const auto& cti : ctis) {
|
2022-06-08 05:49:02 +08:00
|
|
|
actors.add(addRequestActor(this, generation, coordinatorsHash, cti));
|
2021-08-21 06:53:13 +08:00
|
|
|
}
|
|
|
|
return result.getFuture();
|
|
|
|
}
|
2022-05-26 03:08:30 +08:00
|
|
|
bool committed() const { return result.isSet() && !result.isError(); }
|
2021-08-21 06:53:13 +08:00
|
|
|
};
|
|
|
|
|
2021-08-08 09:35:12 +08:00
|
|
|
class GetGenerationQuorum {
|
2021-08-26 07:28:21 +08:00
|
|
|
ActorCollection actors{ false };
|
2022-09-14 06:43:56 +08:00
|
|
|
CoordinatorsHash coordinatorsHash{ 0 };
|
2021-08-21 01:12:07 +08:00
|
|
|
std::vector<ConfigTransactionInterface> ctis;
|
2021-08-08 10:42:41 +08:00
|
|
|
std::map<ConfigGeneration, std::vector<ConfigTransactionInterface>> seenGenerations;
|
2021-08-21 01:12:07 +08:00
|
|
|
Promise<ConfigGeneration> result;
|
2021-08-08 09:35:12 +08:00
|
|
|
size_t totalRepliesReceived{ 0 };
|
|
|
|
size_t maxAgreement{ 0 };
|
2022-06-08 05:49:02 +08:00
|
|
|
Future<Void> coordinatorsChangedFuture;
|
2021-08-08 09:35:12 +08:00
|
|
|
Optional<Version> lastSeenLiveVersion;
|
2021-08-21 01:12:07 +08:00
|
|
|
Future<ConfigGeneration> getGenerationFuture;
|
2021-08-08 09:35:12 +08:00
|
|
|
|
2021-08-08 10:42:41 +08:00
|
|
|
ACTOR static Future<Void> addRequestActor(GetGenerationQuorum* self, ConfigTransactionInterface cti) {
|
2022-02-02 14:27:12 +08:00
|
|
|
loop {
|
|
|
|
try {
|
2022-04-28 12:54:13 +08:00
|
|
|
state ConfigTransactionGetGenerationReply reply;
|
|
|
|
if (cti.hostname.present()) {
|
|
|
|
wait(timeoutError(store(reply,
|
|
|
|
retryGetReplyFromHostname(
|
2022-06-08 05:49:02 +08:00
|
|
|
ConfigTransactionGetGenerationRequest{ self->coordinatorsHash,
|
|
|
|
self->lastSeenLiveVersion },
|
2022-04-28 12:54:13 +08:00
|
|
|
cti.hostname.get(),
|
|
|
|
WLTOKEN_CONFIGTXN_GETGENERATION)),
|
|
|
|
CLIENT_KNOBS->GET_GENERATION_QUORUM_TIMEOUT));
|
|
|
|
} else {
|
|
|
|
wait(timeoutError(store(reply,
|
2022-06-08 05:49:02 +08:00
|
|
|
cti.getGeneration.getReply(ConfigTransactionGetGenerationRequest{
|
|
|
|
self->coordinatorsHash, self->lastSeenLiveVersion })),
|
2022-04-28 12:54:13 +08:00
|
|
|
CLIENT_KNOBS->GET_GENERATION_QUORUM_TIMEOUT));
|
|
|
|
}
|
2021-08-26 07:28:21 +08:00
|
|
|
|
2022-02-02 14:27:12 +08:00
|
|
|
++self->totalRepliesReceived;
|
|
|
|
auto gen = reply.generation;
|
|
|
|
self->lastSeenLiveVersion =
|
|
|
|
std::max(gen.liveVersion, self->lastSeenLiveVersion.orDefault(::invalidVersion));
|
|
|
|
auto& replicas = self->seenGenerations[gen];
|
|
|
|
replicas.push_back(cti);
|
|
|
|
self->maxAgreement = std::max(replicas.size(), self->maxAgreement);
|
2022-06-08 05:49:02 +08:00
|
|
|
// TraceEvent("ConfigTransactionGotGenerationReply")
|
|
|
|
// .detail("From", cti.getGeneration.getEndpoint().getPrimaryAddress())
|
|
|
|
// .detail("TotalRepliesReceived", self->totalRepliesReceived)
|
|
|
|
// .detail("ReplyGeneration", gen.toString())
|
|
|
|
// .detail("Replicas", replicas.size())
|
|
|
|
// .detail("Coordinators", self->ctis.size())
|
|
|
|
// .detail("MaxAgreement", self->maxAgreement)
|
|
|
|
// .detail("LastSeenLiveVersion", self->lastSeenLiveVersion);
|
2022-02-02 14:27:12 +08:00
|
|
|
if (replicas.size() >= self->ctis.size() / 2 + 1 && !self->result.isSet()) {
|
|
|
|
self->result.send(gen);
|
|
|
|
} else if (self->maxAgreement + (self->ctis.size() - self->totalRepliesReceived) <
|
|
|
|
(self->ctis.size() / 2 + 1)) {
|
|
|
|
if (!self->result.isError()) {
|
2022-05-26 03:08:30 +08:00
|
|
|
// Calling sendError could delete self
|
|
|
|
auto local = self->result;
|
|
|
|
local.sendError(failed_to_reach_quorum());
|
2022-02-02 14:27:12 +08:00
|
|
|
}
|
|
|
|
}
|
|
|
|
break;
|
|
|
|
} catch (Error& e) {
|
|
|
|
if (e.code() == error_code_broken_promise) {
|
|
|
|
continue;
|
|
|
|
} else if (e.code() == error_code_timed_out) {
|
|
|
|
++self->totalRepliesReceived;
|
|
|
|
if (self->totalRepliesReceived == self->ctis.size() && self->result.canBeSet() &&
|
|
|
|
!self->result.isError()) {
|
2022-05-26 03:08:30 +08:00
|
|
|
// Calling sendError could delete self
|
|
|
|
auto local = self->result;
|
|
|
|
local.sendError(failed_to_reach_quorum());
|
2022-02-02 14:27:12 +08:00
|
|
|
}
|
|
|
|
break;
|
|
|
|
} else {
|
|
|
|
throw;
|
|
|
|
}
|
2021-08-21 01:12:07 +08:00
|
|
|
}
|
2021-08-08 09:35:12 +08:00
|
|
|
}
|
|
|
|
return Void();
|
|
|
|
}
|
|
|
|
|
2021-08-21 01:12:07 +08:00
|
|
|
ACTOR static Future<ConfigGeneration> getGenerationActor(GetGenerationQuorum* self) {
|
2021-08-08 09:35:12 +08:00
|
|
|
state int retries = 0;
|
|
|
|
loop {
|
2021-08-21 01:12:07 +08:00
|
|
|
for (const auto& cti : self->ctis) {
|
2021-08-26 07:28:21 +08:00
|
|
|
self->actors.add(addRequestActor(self, cti));
|
2021-08-08 09:35:12 +08:00
|
|
|
}
|
|
|
|
try {
|
2021-08-26 07:28:21 +08:00
|
|
|
choose {
|
2022-12-09 01:26:45 +08:00
|
|
|
when(ConfigGeneration generation = wait(self->result.getFuture())) {
|
|
|
|
return generation;
|
|
|
|
}
|
|
|
|
when(wait(self->actors.getResult())) {
|
|
|
|
ASSERT(false);
|
|
|
|
}
|
2021-08-26 07:28:21 +08:00
|
|
|
}
|
2021-08-08 09:35:12 +08:00
|
|
|
} catch (Error& e) {
|
|
|
|
if (e.code() == error_code_failed_to_reach_quorum) {
|
2022-07-20 04:15:51 +08:00
|
|
|
CODE_PROBE(true, "Failed to reach quorum getting generation");
|
2022-06-08 05:49:02 +08:00
|
|
|
if (self->coordinatorsChangedFuture.isReady()) {
|
|
|
|
throw coordinators_changed();
|
|
|
|
}
|
2022-11-19 07:10:53 +08:00
|
|
|
if (deterministicRandom()->random01() < 0.95) {
|
|
|
|
// Add some random jitter to prevent clients from
|
|
|
|
// contending.
|
|
|
|
wait(delayJittered(std::clamp(
|
|
|
|
0.006 * (1 << std::min(retries, 30)), 0.0, CLIENT_KNOBS->TIMEOUT_RETRY_UPPER_BOUND)));
|
|
|
|
}
|
2022-06-08 05:49:02 +08:00
|
|
|
if (deterministicRandom()->random01() < 0.05) {
|
|
|
|
// Randomly inject a delay of at least the generation
|
|
|
|
// reply timeout, to try to prevent contention between
|
|
|
|
// clients.
|
|
|
|
wait(delay(CLIENT_KNOBS->GET_GENERATION_QUORUM_TIMEOUT *
|
|
|
|
(deterministicRandom()->random01() + 1.0)));
|
|
|
|
}
|
2021-08-08 09:35:12 +08:00
|
|
|
++retries;
|
2021-08-26 07:28:21 +08:00
|
|
|
self->actors.clear(false);
|
2022-02-02 14:27:12 +08:00
|
|
|
self->seenGenerations.clear();
|
2021-10-15 06:03:53 +08:00
|
|
|
self->result.reset();
|
|
|
|
self->totalRepliesReceived = 0;
|
|
|
|
self->maxAgreement = 0;
|
2021-08-08 09:35:12 +08:00
|
|
|
} else {
|
|
|
|
throw e;
|
|
|
|
}
|
|
|
|
}
|
2021-07-19 05:02:45 +08:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2021-08-21 01:12:07 +08:00
|
|
|
public:
|
|
|
|
GetGenerationQuorum() = default;
|
2022-09-14 06:43:56 +08:00
|
|
|
explicit GetGenerationQuorum(CoordinatorsHash coordinatorsHash,
|
2022-06-08 05:49:02 +08:00
|
|
|
std::vector<ConfigTransactionInterface> const& ctis,
|
|
|
|
Future<Void> coordinatorsChangedFuture,
|
2021-08-21 01:12:07 +08:00
|
|
|
Optional<Version> const& lastSeenLiveVersion = {})
|
2022-06-08 05:49:02 +08:00
|
|
|
: coordinatorsHash(coordinatorsHash), ctis(ctis), coordinatorsChangedFuture(coordinatorsChangedFuture),
|
|
|
|
lastSeenLiveVersion(lastSeenLiveVersion) {}
|
2021-08-21 01:12:07 +08:00
|
|
|
Future<ConfigGeneration> getGeneration() {
|
|
|
|
if (!getGenerationFuture.isValid()) {
|
|
|
|
getGenerationFuture = getGenerationActor(this);
|
2021-07-19 05:02:45 +08:00
|
|
|
}
|
2021-08-21 01:12:07 +08:00
|
|
|
return getGenerationFuture;
|
|
|
|
}
|
|
|
|
bool isReady() const {
|
|
|
|
return getGenerationFuture.isValid() && getGenerationFuture.isReady() && !getGenerationFuture.isError();
|
|
|
|
}
|
|
|
|
Optional<ConfigGeneration> getCachedGeneration() const {
|
|
|
|
return isReady() ? getGenerationFuture.get() : Optional<ConfigGeneration>{};
|
|
|
|
}
|
|
|
|
std::vector<ConfigTransactionInterface> getReadReplicas() const {
|
|
|
|
ASSERT(isReady());
|
|
|
|
return seenGenerations.at(getGenerationFuture.get());
|
|
|
|
}
|
|
|
|
Optional<Version> getLastSeenLiveVersion() const { return lastSeenLiveVersion; }
|
|
|
|
};
|
|
|
|
|
|
|
|
class PaxosConfigTransactionImpl {
|
2022-09-14 06:43:56 +08:00
|
|
|
CoordinatorsHash coordinatorsHash{ 0 };
|
2021-08-21 01:12:07 +08:00
|
|
|
std::vector<ConfigTransactionInterface> ctis;
|
2021-08-21 06:53:13 +08:00
|
|
|
GetGenerationQuorum getGenerationQuorum;
|
|
|
|
CommitQuorum commitQuorum;
|
2021-08-21 01:12:07 +08:00
|
|
|
int numRetries{ 0 };
|
|
|
|
Optional<UID> dID;
|
|
|
|
Database cx;
|
2022-06-08 05:49:02 +08:00
|
|
|
Future<Void> watchClusterFileFuture;
|
2021-08-21 01:12:07 +08:00
|
|
|
|
|
|
|
ACTOR static Future<Optional<Value>> get(PaxosConfigTransactionImpl* self, Key key) {
|
2022-02-09 08:28:30 +08:00
|
|
|
state ConfigKey configKey = ConfigKey::decodeKey(key);
|
2022-02-02 14:27:12 +08:00
|
|
|
loop {
|
|
|
|
try {
|
2022-04-28 12:54:13 +08:00
|
|
|
state ConfigGeneration generation = wait(self->getGenerationQuorum.getGeneration());
|
|
|
|
state std::vector<ConfigTransactionInterface> readReplicas =
|
|
|
|
self->getGenerationQuorum.getReadReplicas();
|
|
|
|
std::vector<Future<Void>> fs;
|
|
|
|
for (ConfigTransactionInterface& readReplica : readReplicas) {
|
|
|
|
if (readReplica.hostname.present()) {
|
|
|
|
fs.push_back(tryInitializeRequestStream(
|
|
|
|
&readReplica.get, readReplica.hostname.get(), WLTOKEN_CONFIGTXN_GET));
|
|
|
|
}
|
|
|
|
}
|
|
|
|
wait(waitForAll(fs));
|
|
|
|
state Reference<ConfigTransactionInfo> configNodes(new ConfigTransactionInfo(readReplicas));
|
2022-06-08 05:49:02 +08:00
|
|
|
ConfigTransactionGetReply reply = wait(timeoutError(
|
|
|
|
basicLoadBalance(configNodes,
|
|
|
|
&ConfigTransactionInterface::get,
|
|
|
|
ConfigTransactionGetRequest{ self->coordinatorsHash, generation, configKey }),
|
|
|
|
CLIENT_KNOBS->GET_KNOB_TIMEOUT));
|
2022-02-02 14:27:12 +08:00
|
|
|
if (reply.value.present()) {
|
|
|
|
return reply.value.get().toValue();
|
|
|
|
} else {
|
|
|
|
return Optional<Value>{};
|
|
|
|
}
|
|
|
|
} catch (Error& e) {
|
2022-06-08 05:49:02 +08:00
|
|
|
if (e.code() != error_code_timed_out && e.code() != error_code_broken_promise &&
|
|
|
|
e.code() != error_code_coordinators_changed) {
|
2022-02-02 14:27:12 +08:00
|
|
|
throw;
|
|
|
|
}
|
2022-02-09 08:28:30 +08:00
|
|
|
self->reset();
|
2022-02-02 14:27:12 +08:00
|
|
|
}
|
2021-07-19 05:02:45 +08:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2021-07-19 05:21:21 +08:00
|
|
|
ACTOR static Future<RangeResult> getConfigClasses(PaxosConfigTransactionImpl* self) {
|
2022-06-08 05:49:02 +08:00
|
|
|
loop {
|
|
|
|
try {
|
|
|
|
state ConfigGeneration generation = wait(self->getGenerationQuorum.getGeneration());
|
|
|
|
state std::vector<ConfigTransactionInterface> readReplicas =
|
|
|
|
self->getGenerationQuorum.getReadReplicas();
|
|
|
|
std::vector<Future<Void>> fs;
|
|
|
|
for (ConfigTransactionInterface& readReplica : readReplicas) {
|
|
|
|
if (readReplica.hostname.present()) {
|
|
|
|
fs.push_back(tryInitializeRequestStream(
|
|
|
|
&readReplica.getClasses, readReplica.hostname.get(), WLTOKEN_CONFIGTXN_GETCLASSES));
|
|
|
|
}
|
|
|
|
}
|
|
|
|
wait(waitForAll(fs));
|
|
|
|
state Reference<ConfigTransactionInfo> configNodes(new ConfigTransactionInfo(readReplicas));
|
2022-07-30 08:28:34 +08:00
|
|
|
ConfigTransactionGetConfigClassesReply reply = wait(
|
|
|
|
basicLoadBalance(configNodes,
|
|
|
|
&ConfigTransactionInterface::getClasses,
|
|
|
|
ConfigTransactionGetConfigClassesRequest{ self->coordinatorsHash, generation }));
|
2022-06-08 05:49:02 +08:00
|
|
|
RangeResult result;
|
|
|
|
result.reserve(result.arena(), reply.configClasses.size());
|
|
|
|
for (const auto& configClass : reply.configClasses) {
|
|
|
|
result.push_back_deep(result.arena(), KeyValueRef(configClass, ""_sr));
|
|
|
|
}
|
|
|
|
return result;
|
|
|
|
} catch (Error& e) {
|
|
|
|
if (e.code() != error_code_coordinators_changed) {
|
|
|
|
throw;
|
|
|
|
}
|
|
|
|
self->reset();
|
2022-04-28 12:54:13 +08:00
|
|
|
}
|
|
|
|
}
|
2021-07-19 05:21:21 +08:00
|
|
|
}
|
|
|
|
|
2021-07-19 05:26:15 +08:00
|
|
|
ACTOR static Future<RangeResult> getKnobs(PaxosConfigTransactionImpl* self, Optional<Key> configClass) {
|
2022-06-08 05:49:02 +08:00
|
|
|
loop {
|
|
|
|
try {
|
|
|
|
state ConfigGeneration generation = wait(self->getGenerationQuorum.getGeneration());
|
|
|
|
state std::vector<ConfigTransactionInterface> readReplicas =
|
|
|
|
self->getGenerationQuorum.getReadReplicas();
|
|
|
|
std::vector<Future<Void>> fs;
|
|
|
|
for (ConfigTransactionInterface& readReplica : readReplicas) {
|
|
|
|
if (readReplica.hostname.present()) {
|
|
|
|
fs.push_back(tryInitializeRequestStream(
|
|
|
|
&readReplica.getKnobs, readReplica.hostname.get(), WLTOKEN_CONFIGTXN_GETKNOBS));
|
|
|
|
}
|
|
|
|
}
|
|
|
|
wait(waitForAll(fs));
|
|
|
|
state Reference<ConfigTransactionInfo> configNodes(new ConfigTransactionInfo(readReplicas));
|
2022-07-30 08:28:34 +08:00
|
|
|
ConfigTransactionGetKnobsReply reply = wait(basicLoadBalance(
|
|
|
|
configNodes,
|
|
|
|
&ConfigTransactionInterface::getKnobs,
|
|
|
|
ConfigTransactionGetKnobsRequest{ self->coordinatorsHash, generation, configClass }));
|
2022-06-08 05:49:02 +08:00
|
|
|
RangeResult result;
|
|
|
|
result.reserve(result.arena(), reply.knobNames.size());
|
|
|
|
for (const auto& knobName : reply.knobNames) {
|
|
|
|
result.push_back_deep(result.arena(), KeyValueRef(knobName, ""_sr));
|
|
|
|
}
|
|
|
|
return result;
|
|
|
|
} catch (Error& e) {
|
|
|
|
if (e.code() != error_code_coordinators_changed) {
|
|
|
|
throw;
|
|
|
|
}
|
|
|
|
self->reset();
|
2022-04-28 12:54:13 +08:00
|
|
|
}
|
|
|
|
}
|
2021-07-19 05:21:21 +08:00
|
|
|
}
|
|
|
|
|
2021-07-19 05:43:58 +08:00
|
|
|
ACTOR static Future<Void> commit(PaxosConfigTransactionImpl* self) {
|
2022-06-08 05:49:02 +08:00
|
|
|
loop {
|
|
|
|
try {
|
|
|
|
ConfigGeneration generation = wait(self->getGenerationQuorum.getGeneration());
|
|
|
|
self->commitQuorum.setTimestamp();
|
|
|
|
wait(self->commitQuorum.commit(generation, self->coordinatorsHash));
|
|
|
|
return Void();
|
|
|
|
} catch (Error& e) {
|
|
|
|
if (e.code() != error_code_coordinators_changed) {
|
|
|
|
throw;
|
|
|
|
}
|
|
|
|
self->reset();
|
|
|
|
}
|
|
|
|
}
|
2021-07-19 05:43:58 +08:00
|
|
|
}
|
|
|
|
|
2021-08-28 06:06:33 +08:00
|
|
|
ACTOR static Future<Void> onError(PaxosConfigTransactionImpl* self, Error e) {
|
|
|
|
// TODO: Improve this:
|
|
|
|
TraceEvent("ConfigIncrementOnError").error(e).detail("NumRetries", self->numRetries);
|
|
|
|
if (e.code() == error_code_transaction_too_old || e.code() == error_code_not_committed) {
|
2022-02-09 08:28:30 +08:00
|
|
|
wait(delay(std::clamp((1 << self->numRetries++) * 0.01 * deterministicRandom()->random01(),
|
|
|
|
0.0,
|
|
|
|
CLIENT_KNOBS->TIMEOUT_RETRY_UPPER_BOUND)));
|
2021-08-28 06:06:33 +08:00
|
|
|
self->reset();
|
|
|
|
return Void();
|
|
|
|
}
|
|
|
|
throw e;
|
|
|
|
}
|
|
|
|
|
2022-06-08 05:49:02 +08:00
|
|
|
// Returns when the cluster interface updates with a new connection string.
|
|
|
|
ACTOR static Future<Void> watchClusterFile(Database cx) {
|
2022-08-17 08:21:49 +08:00
|
|
|
state Future<Void> leaderMonitor =
|
|
|
|
monitorLeader<ClusterInterface>(cx->getConnectionRecord(), cx->statusClusterInterface);
|
2022-06-08 05:49:02 +08:00
|
|
|
state std::string connectionString = cx->getConnectionRecord()->getConnectionString().toString();
|
|
|
|
|
|
|
|
loop {
|
2022-09-13 13:34:46 +08:00
|
|
|
wait(cx->statusClusterInterface->onChange());
|
|
|
|
if (cx->getConnectionRecord()->getConnectionString().toString() != connectionString) {
|
|
|
|
return Void();
|
2022-06-08 05:49:02 +08:00
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2021-07-19 05:02:45 +08:00
|
|
|
public:
|
|
|
|
Future<Version> getReadVersion() {
|
2021-08-21 01:12:07 +08:00
|
|
|
return map(getGenerationQuorum.getGeneration(), [](auto const& gen) { return gen.committedVersion; });
|
2021-07-19 05:02:45 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
Optional<Version> getCachedReadVersion() const {
|
2021-08-21 01:12:07 +08:00
|
|
|
auto gen = getGenerationQuorum.getCachedGeneration();
|
|
|
|
if (gen.present()) {
|
|
|
|
return gen.get().committedVersion;
|
2021-07-19 05:02:45 +08:00
|
|
|
} else {
|
|
|
|
return {};
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2021-08-08 10:42:41 +08:00
|
|
|
Version getCommittedVersion() const {
|
2021-08-21 06:53:13 +08:00
|
|
|
return commitQuorum.committed() ? getGenerationQuorum.getCachedGeneration().get().liveVersion
|
|
|
|
: ::invalidVersion;
|
2021-08-08 10:42:41 +08:00
|
|
|
}
|
2021-07-19 05:02:45 +08:00
|
|
|
|
2021-08-21 06:53:13 +08:00
|
|
|
int64_t getApproximateSize() const { return commitQuorum.expectedSize(); }
|
2021-07-19 05:02:45 +08:00
|
|
|
|
2021-08-21 06:53:13 +08:00
|
|
|
void set(KeyRef key, ValueRef value) { commitQuorum.set(key, value); }
|
2021-07-19 05:02:45 +08:00
|
|
|
|
2021-08-21 06:53:13 +08:00
|
|
|
void clear(KeyRef key) { commitQuorum.clear(key); }
|
2021-07-19 05:02:45 +08:00
|
|
|
|
|
|
|
Future<Optional<Value>> get(Key const& key) { return get(this, key); }
|
|
|
|
|
2021-07-19 05:21:21 +08:00
|
|
|
Future<RangeResult> getRange(KeyRangeRef keys) {
|
|
|
|
if (keys == configClassKeys) {
|
|
|
|
return getConfigClasses(this);
|
|
|
|
} else if (keys == globalConfigKnobKeys) {
|
|
|
|
return getKnobs(this, {});
|
|
|
|
} else if (configKnobKeys.contains(keys) && keys.singleKeyRange()) {
|
|
|
|
const auto configClass = keys.begin.removePrefix(configKnobKeys.begin);
|
|
|
|
return getKnobs(this, configClass);
|
|
|
|
} else {
|
|
|
|
throw invalid_config_db_range_read();
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2021-08-28 06:06:33 +08:00
|
|
|
Future<Void> onError(Error const& e) { return onError(this, e); }
|
2021-07-19 05:02:45 +08:00
|
|
|
|
|
|
|
void debugTransaction(UID dID) { this->dID = dID; }
|
|
|
|
|
|
|
|
void reset() {
|
2022-06-08 05:49:02 +08:00
|
|
|
ctis.clear();
|
|
|
|
// Re-read connection string. If the cluster file changed, this will
|
|
|
|
// return the updated value.
|
|
|
|
const ClusterConnectionString& cs = cx->getConnectionRecord()->getConnectionString();
|
|
|
|
ctis.reserve(cs.hostnames.size() + cs.coords.size());
|
|
|
|
for (const auto& h : cs.hostnames) {
|
|
|
|
ctis.emplace_back(h);
|
|
|
|
}
|
|
|
|
for (const auto& c : cs.coords) {
|
|
|
|
ctis.emplace_back(c);
|
|
|
|
}
|
|
|
|
coordinatorsHash = std::hash<std::string>()(cx->getConnectionRecord()->getConnectionString().toString());
|
2022-08-17 08:21:49 +08:00
|
|
|
if (!cx->statusLeaderMon.isValid() || cx->statusLeaderMon.isReady()) {
|
|
|
|
cx->statusClusterInterface = makeReference<AsyncVar<Optional<ClusterInterface>>>();
|
|
|
|
cx->statusLeaderMon = watchClusterFile(cx);
|
|
|
|
}
|
2022-06-08 05:49:02 +08:00
|
|
|
getGenerationQuorum = GetGenerationQuorum{
|
2022-08-17 08:21:49 +08:00
|
|
|
coordinatorsHash, ctis, cx->statusLeaderMon, getGenerationQuorum.getLastSeenLiveVersion()
|
2022-06-08 05:49:02 +08:00
|
|
|
};
|
2021-08-21 06:53:13 +08:00
|
|
|
commitQuorum = CommitQuorum{ ctis };
|
2021-07-19 05:02:45 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
void fullReset() {
|
|
|
|
numRetries = 0;
|
|
|
|
dID = {};
|
|
|
|
reset();
|
|
|
|
}
|
|
|
|
|
|
|
|
void checkDeferredError(Error const& deferredError) const {
|
|
|
|
if (deferredError.code() != invalid_error_code) {
|
|
|
|
throw deferredError;
|
|
|
|
}
|
|
|
|
if (cx.getPtr()) {
|
|
|
|
cx->checkDeferredError();
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2021-07-19 05:43:58 +08:00
|
|
|
Future<Void> commit() { return commit(this); }
|
|
|
|
|
2022-06-08 05:49:02 +08:00
|
|
|
PaxosConfigTransactionImpl(Database const& cx) : cx(cx) { reset(); }
|
2021-07-19 05:02:45 +08:00
|
|
|
|
2021-08-21 06:53:13 +08:00
|
|
|
PaxosConfigTransactionImpl(std::vector<ConfigTransactionInterface> const& ctis)
|
2022-06-08 05:49:02 +08:00
|
|
|
: ctis(ctis), getGenerationQuorum(0, ctis, Future<Void>()), commitQuorum(ctis) {}
|
2021-07-19 05:02:45 +08:00
|
|
|
};
|
2021-05-18 04:02:55 +08:00
|
|
|
|
|
|
|
Future<Version> PaxosConfigTransaction::getReadVersion() {
|
2021-08-03 03:32:11 +08:00
|
|
|
return impl->getReadVersion();
|
2021-05-18 04:02:55 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
Optional<Version> PaxosConfigTransaction::getCachedReadVersion() const {
|
2021-08-03 03:32:11 +08:00
|
|
|
return impl->getCachedReadVersion();
|
2021-05-18 04:02:55 +08:00
|
|
|
}
|
|
|
|
|
2021-07-19 05:02:45 +08:00
|
|
|
Future<Optional<Value>> PaxosConfigTransaction::get(Key const& key, Snapshot) {
|
2021-08-03 03:32:11 +08:00
|
|
|
return impl->get(key);
|
2021-05-18 04:02:55 +08:00
|
|
|
}
|
|
|
|
|
2021-07-19 05:21:21 +08:00
|
|
|
Future<RangeResult> PaxosConfigTransaction::getRange(KeySelector const& begin,
|
|
|
|
KeySelector const& end,
|
|
|
|
int limit,
|
|
|
|
Snapshot snapshot,
|
|
|
|
Reverse reverse) {
|
|
|
|
if (reverse) {
|
|
|
|
throw client_invalid_operation();
|
|
|
|
}
|
2021-08-03 03:32:11 +08:00
|
|
|
return impl->getRange(KeyRangeRef(begin.getKey(), end.getKey()));
|
2021-05-18 04:02:55 +08:00
|
|
|
}
|
|
|
|
|
2021-07-19 05:26:15 +08:00
|
|
|
Future<RangeResult> PaxosConfigTransaction::getRange(KeySelector begin,
|
|
|
|
KeySelector end,
|
|
|
|
GetRangeLimits limits,
|
|
|
|
Snapshot snapshot,
|
|
|
|
Reverse reverse) {
|
2021-07-19 05:21:21 +08:00
|
|
|
if (reverse) {
|
|
|
|
throw client_invalid_operation();
|
|
|
|
}
|
2021-08-03 03:32:11 +08:00
|
|
|
return impl->getRange(KeyRangeRef(begin.getKey(), end.getKey()));
|
2021-05-18 04:02:55 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
void PaxosConfigTransaction::set(KeyRef const& key, ValueRef const& value) {
|
2021-08-03 03:32:11 +08:00
|
|
|
return impl->set(key, value);
|
2021-05-18 04:02:55 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
void PaxosConfigTransaction::clear(KeyRef const& key) {
|
2021-08-03 03:32:11 +08:00
|
|
|
return impl->clear(key);
|
2021-05-18 04:02:55 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
Future<Void> PaxosConfigTransaction::commit() {
|
2021-08-03 03:32:11 +08:00
|
|
|
return impl->commit();
|
2021-05-18 04:02:55 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
Version PaxosConfigTransaction::getCommittedVersion() const {
|
2021-08-03 03:32:11 +08:00
|
|
|
return impl->getCommittedVersion();
|
2021-05-18 04:02:55 +08:00
|
|
|
}
|
|
|
|
|
2022-11-14 03:22:57 +08:00
|
|
|
double PaxosConfigTransaction::getTagThrottledDuration() const {
|
|
|
|
return 0.0;
|
|
|
|
}
|
|
|
|
|
2022-10-17 10:56:51 +08:00
|
|
|
int64_t PaxosConfigTransaction::getTotalCost() const {
|
|
|
|
return 0;
|
|
|
|
}
|
|
|
|
|
2021-05-18 04:02:55 +08:00
|
|
|
int64_t PaxosConfigTransaction::getApproximateSize() const {
|
2021-08-03 03:32:11 +08:00
|
|
|
return impl->getApproximateSize();
|
2021-05-18 04:02:55 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
void PaxosConfigTransaction::setOption(FDBTransactionOptions::Option option, Optional<StringRef> value) {
|
2021-07-19 05:02:45 +08:00
|
|
|
// TODO: Support using this option to determine atomicity
|
2021-05-18 04:02:55 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
Future<Void> PaxosConfigTransaction::onError(Error const& e) {
|
2021-08-03 03:32:11 +08:00
|
|
|
return impl->onError(e);
|
2021-05-18 04:02:55 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
void PaxosConfigTransaction::cancel() {
|
2021-07-19 05:02:45 +08:00
|
|
|
// TODO: Implement someday
|
|
|
|
throw client_invalid_operation();
|
2021-05-18 04:02:55 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
void PaxosConfigTransaction::reset() {
|
2021-08-03 03:32:11 +08:00
|
|
|
impl->reset();
|
2021-05-18 04:02:55 +08:00
|
|
|
}
|
|
|
|
|
2021-05-18 06:31:03 +08:00
|
|
|
void PaxosConfigTransaction::fullReset() {
|
2021-08-03 03:32:11 +08:00
|
|
|
impl->fullReset();
|
2021-05-18 06:31:03 +08:00
|
|
|
}
|
|
|
|
|
2021-05-18 04:02:55 +08:00
|
|
|
void PaxosConfigTransaction::debugTransaction(UID dID) {
|
2021-08-03 03:32:11 +08:00
|
|
|
impl->debugTransaction(dID);
|
2021-05-18 04:02:55 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
void PaxosConfigTransaction::checkDeferredError() const {
|
2021-08-03 03:32:11 +08:00
|
|
|
impl->checkDeferredError(deferredError);
|
2021-05-18 04:02:55 +08:00
|
|
|
}
|
|
|
|
|
2021-07-19 05:02:45 +08:00
|
|
|
PaxosConfigTransaction::PaxosConfigTransaction(std::vector<ConfigTransactionInterface> const& ctis)
|
2021-08-03 03:32:11 +08:00
|
|
|
: impl(PImpl<PaxosConfigTransactionImpl>::create(ctis)) {}
|
2021-07-19 05:02:45 +08:00
|
|
|
|
|
|
|
PaxosConfigTransaction::PaxosConfigTransaction() = default;
|
2021-05-18 04:02:55 +08:00
|
|
|
|
|
|
|
PaxosConfigTransaction::~PaxosConfigTransaction() = default;
|
2021-06-23 12:44:59 +08:00
|
|
|
|
2022-02-20 07:09:55 +08:00
|
|
|
void PaxosConfigTransaction::construct(Database const& cx) {
|
2021-08-03 03:32:11 +08:00
|
|
|
impl = PImpl<PaxosConfigTransactionImpl>::create(cx);
|
2021-06-30 01:29:33 +08:00
|
|
|
}
|