Skip to content
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
70 changes: 60 additions & 10 deletions fdbclient/StatusClient.actor.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
*/

#include "flow/flow.h"
#include "fdbclient/ClusterConnectionMemoryRecord.h"
#include "fdbclient/CoordinationInterface.h"
#include "fdbclient/MonitorLeader.h"
#include "fdbclient/ClusterInterface.h"
Expand Down Expand Up @@ -305,7 +306,8 @@ void JSONDoc::mergeValueInto(json_spirit::mValue& dst, const json_spirit::mValue
// Will not throw, will just return non-present Optional if error
ACTOR Future<Optional<StatusObject>> clientCoordinatorsStatusFetcher(Reference<IClusterConnectionRecord> connRecord,
bool* quorum_reachable,
int* coordinatorsFaultTolerance) {
int* coordinatorsFaultTolerance,
double statusDeadline) {
try {
state ClientCoordinators coord(connRecord);
state StatusObject statusObj;
Expand Down Expand Up @@ -340,7 +342,7 @@ ACTOR Future<Optional<StatusObject>> clientCoordinatorsStatusFetcher(Reference<I

wait(smartQuorum(leaderServers, leaderServers.size() / 2 + 1, 1.5) &&
smartQuorum(coordProtocols, coordProtocols.size() / 2 + 1, 1.5) ||
delay(2.0));
delay(std::min(2.0, std::max(0.0, statusDeadline - now()))));

statusObj["quorum_reachable"] = *quorum_reachable =
quorum(leaderServers, leaderServers.size() / 2 + 1).isReady();
Expand Down Expand Up @@ -380,11 +382,12 @@ ACTOR Future<Optional<StatusObject>> clientCoordinatorsStatusFetcher(Reference<I
ACTOR Future<StatusObject> clientStatusFetcher(Reference<IClusterConnectionRecord> connRecord,
StatusArray* messages,
bool* quorum_reachable,
int* coordinatorsFaultTolerance) {
int* coordinatorsFaultTolerance,
double statusDeadline) {
state StatusObject statusObj;

state Optional<StatusObject> coordsStatusObj =
wait(clientCoordinatorsStatusFetcher(connRecord, quorum_reachable, coordinatorsFaultTolerance));
wait(clientCoordinatorsStatusFetcher(connRecord, quorum_reachable, coordinatorsFaultTolerance, statusDeadline));
state bool contentsUpToDate = wait(connRecord->upToDate());

if (coordsStatusObj.present()) {
Expand Down Expand Up @@ -422,9 +425,10 @@ ACTOR Future<StatusObject> clientStatusFetcher(Reference<IClusterConnectionRecor
// Cluster section of json output
ACTOR Future<Optional<StatusObject>> clusterStatusFetcher(ClusterInterface cI,
StatusArray* messages,
std::string statusField) {
std::string statusField,
double timeoutSeconds) {
state StatusRequest req(statusField);
state Future<Void> clusterTimeout = delay(CLIENT_KNOBS->STATUS_TIMEOUT);
state Future<Void> clusterTimeout = delay(timeoutSeconds);
state Optional<StatusObject> oStatusObj;

wait(delay(0.0)); // make sure the cluster controller is marked as not failed
Expand Down Expand Up @@ -462,6 +466,46 @@ ACTOR Future<Optional<StatusObject>> clusterStatusFetcher(ClusterInterface cI,
return oStatusObj;
}

TEST_CASE("/fdbclient/status/overallTimeoutBudget") {
state ClusterInterface clusterInterface;
state StatusArray messages;
state double overallTimeout = 0.2;
state double coordinatorProbeDuration = 0.1;
state double startTime = now();
state double statusDeadline = startTime + overallTimeout;

wait(delay(coordinatorProbeDuration));
state Optional<StatusObject> status =
wait(clusterStatusFetcher(clusterInterface, &messages, "", std::max(0.0, statusDeadline - now())));

double elapsed = now() - startTime;
ASSERT(!status.present());
ASSERT_EQ(messages.size(), 1);
ASSERT_EQ(messages.front().get_obj().at("name").get_str(), "status_incomplete_timeout");
ASSERT_GE(elapsed, overallTimeout);
ASSERT_LT(elapsed, overallTimeout + 0.05);
return Void();
}

TEST_CASE("/fdbclient/status/clientCoordinatorTimeoutBudget") {
state Reference<IClusterConnectionRecord> connRecord =
makeReference<ClusterConnectionMemoryRecord>(ClusterConnectionString("test:test@127.0.0.1:1"));
state bool quorumReachable = false;
state int coordinatorsFaultTolerance = 0;
state double timeout = 0.2;
state double startTime = now();

state Optional<StatusObject> status = wait(clientCoordinatorsStatusFetcher(
connRecord, &quorumReachable, &coordinatorsFaultTolerance, startTime + timeout));

double elapsed = now() - startTime;
ASSERT(status.present());
ASSERT(!quorumReachable);
ASSERT_GE(elapsed, timeout);
ASSERT_LT(elapsed, timeout + 0.05);
return Void();
}

// Create and return a database_status section.
// Will not throw, will not return an empty section.
StatusObject getClientDatabaseStatus(StatusObjectReader client, StatusObjectReader cluster) {
Expand Down Expand Up @@ -516,6 +560,9 @@ ACTOR Future<StatusObject> statusFetcherImpl(Reference<IClusterConnectionRecord>
if (!g_network)
throw network_not_setup();

// Keep the overall status request within STATUS_TIMEOUT, including the client-side coordinator probe that runs
// before we ask the cluster controller for its status payload.
state double statusDeadline = now() + CLIENT_KNOBS->STATUS_TIMEOUT;
state StatusObject statusObj;
state StatusObject statusObjClient;
state StatusArray clientMessages;
Expand All @@ -527,8 +574,8 @@ ACTOR Future<StatusObject> statusFetcherImpl(Reference<IClusterConnectionRecord>
try {
state int64_t clientTime = g_network->timer();

StatusObject _statusObjClient =
wait(clientStatusFetcher(connRecord, &clientMessages, &quorum_reachable, &coordinatorsFaultTolerance));
StatusObject _statusObjClient = wait(clientStatusFetcher(
connRecord, &clientMessages, &quorum_reachable, &coordinatorsFaultTolerance, statusDeadline));
statusObjClient = _statusObjClient;

if (clientTime != -1)
Expand All @@ -549,12 +596,15 @@ ACTOR Future<StatusObject> statusFetcherImpl(Reference<IClusterConnectionRecord>
if (quorum_reachable) {
try {

state Future<Void> interfaceTimeout = delay(2.0);
state Future<Void> interfaceTimeout = delay(std::min(2.0, std::max(0.0, statusDeadline - now())));

loop {
if (clusterInterface->get().present()) {
Optional<StatusObject> _statusObjCluster =
wait(clusterStatusFetcher(clusterInterface->get().get(), &clientMessages, statusField));
wait(clusterStatusFetcher(clusterInterface->get().get(),
&clientMessages,
statusField,
std::max(0.0, statusDeadline - now())));
if (_statusObjCluster.present()) {
statusObjCluster = _statusObjCluster.get();
// TODO: this is a temporary fix, getting the number of available coordinators should move to
Expand Down