Add .cluster.configuration status json field

This commit is contained in:
sfc-gh-tclinkenbeard 2021-05-18 10:47:16 -07:00
parent bb0838676b
commit fcc6efd3b1
9 changed files with 85 additions and 48 deletions

View File

@ -2718,7 +2718,8 @@ ACTOR Future<Void> timeKeeper(ClusterControllerData* self) {
ACTOR Future<Void> statusServer(FutureStream<StatusRequest> requests,
ClusterControllerData* self,
ServerCoordinators coordinators) {
ServerCoordinators coordinators,
ConfigBroadcaster const* configBroadcaster) {
// Seconds since the END of the last GetStatus executed
state double last_request_time = 0.0;
@ -2785,7 +2786,8 @@ ACTOR Future<Void> statusServer(FutureStream<StatusRequest> requests,
&self->db.clientStatus,
coordinators,
incompatibleConnections,
self->datacenterVersionDifference)));
self->datacenterVersionDifference,
configBroadcaster)));
if (result.isError() && result.getError().code() == error_code_actor_cancelled)
throw result.getError();
@ -3417,17 +3419,21 @@ ACTOR Future<Void> dbInfoUpdater(ClusterControllerData* self) {
ACTOR Future<Void> clusterControllerCore(ClusterControllerFullInterface interf,
Future<Void> leaderFail,
ServerCoordinators coordinators,
LocalityData locality) {
LocalityData locality,
Optional<bool> useTestConfigDB) {
state ClusterControllerData self(interf, locality, coordinators);
state ConfigBroadcaster configBroadcaster(coordinators);
state ConfigBroadcaster configBroadcaster(coordinators, useTestConfigDB);
state Future<Void> coordinationPingDelay = delay(SERVER_KNOBS->WORKER_COORDINATION_PING_DELAY);
state uint64_t step = 0;
state Future<ErrorOr<Void>> error = errorOr(actorCollection(self.addActor.getFuture()));
self.addActor.send(configBroadcaster.serve(self.db.serverInfo->get().configBroadcaster));
if (useTestConfigDB.present()) {
self.addActor.send(configBroadcaster.serve(self.db.serverInfo->get().configBroadcaster));
}
self.addActor.send(clusterWatchDatabase(&self, &self.db)); // Start the master database
self.addActor.send(self.updateWorkerList.init(self.db.db));
self.addActor.send(statusServer(interf.clientInterface.databaseStatus.getFuture(), &self, coordinators));
self.addActor.send(
statusServer(interf.clientInterface.databaseStatus.getFuture(), &self, coordinators, &configBroadcaster));
self.addActor.send(timeKeeper(&self));
self.addActor.send(monitorProcessClasses(&self));
self.addActor.send(monitorServerInfoConfig(&self.db));
@ -3544,7 +3550,8 @@ ACTOR Future<Void> clusterController(ServerCoordinators coordinators,
Reference<AsyncVar<Optional<ClusterControllerFullInterface>>> currentCC,
bool hasConnected,
Reference<AsyncVar<ClusterControllerPriorityInfo>> asyncPriorityInfo,
LocalityData locality) {
LocalityData locality,
Optional<bool> useTestConfigDB) {
loop {
state ClusterControllerFullInterface cci;
state bool inRole = false;
@ -3571,7 +3578,7 @@ ACTOR Future<Void> clusterController(ServerCoordinators coordinators,
startRole(Role::CLUSTER_CONTROLLER, cci.id(), UID());
inRole = true;
wait(clusterControllerCore(cci, leaderFail, coordinators, locality));
wait(clusterControllerCore(cci, leaderFail, coordinators, locality, useTestConfigDB));
}
} catch (Error& e) {
if (inRole)
@ -3594,13 +3601,15 @@ ACTOR Future<Void> clusterController(Reference<ClusterConnectionFile> connFile,
Reference<AsyncVar<Optional<ClusterControllerFullInterface>>> currentCC,
Reference<AsyncVar<ClusterControllerPriorityInfo>> asyncPriorityInfo,
Future<Void> recoveredDiskFiles,
LocalityData locality) {
LocalityData locality,
Optional<bool> useTestConfigDB) {
wait(recoveredDiskFiles);
state bool hasConnected = false;
loop {
try {
ServerCoordinators coordinators(connFile);
wait(clusterController(coordinators, currentCC, hasConnected, asyncPriorityInfo, locality));
wait(
clusterController(coordinators, currentCC, hasConnected, asyncPriorityInfo, locality, useTestConfigDB));
} catch (Error& e) {
if (e.code() != error_code_coordinators_changed)
throw; // Expected to terminate fdbserver

View File

@ -18,7 +18,6 @@
* limitations under the License.
*/
#include "fdbclient/JsonBuilder.h"
#include "fdbserver/ConfigBroadcaster.h"
#include "fdbserver/IConfigConsumer.h"
#include "flow/actorcompiler.h" // must be last include
@ -29,15 +28,6 @@ bool matchesConfigClass(Optional<ConfigClassSet> const& configClassSet, Optional
return !configClassSet.present() || !configClass.present() || configClassSet.get().contains(configClass.get());
}
template <class T>
std::string configClassToString(Optional<T> const& configClass) {
if (configClass.present()) {
return configClass.get().toString();
} else {
return "<global>";
}
}
} // namespace
class ConfigBroadcasterImpl {
@ -120,6 +110,15 @@ class ConfigBroadcasterImpl {
}
}
ConfigBroadcasterImpl()
: id(deterministicRandom()->randomUniqueID()), lastCompactedVersion(0), mostRecentVersion(0),
cc("ConfigBroadcaster"), compactRequest("CompactRequest", cc),
successfulChangeRequest("SuccessfulChangeRequest", cc), failedChangeRequest("FailedChangeRequest", cc),
snapshotRequest("SnapshotRequest", cc) {
logger = traceCounters(
"ConfigBroadcasterMetrics", id, SERVER_KNOBS->WORKER_LOGGING_INTERVAL, &cc, "ConfigBroadcasterMetrics");
}
public:
Future<Void> serve(ConfigBroadcaster* self, ConfigFollowerInterface const& cfi) { return serve(self, this, cfi); }
@ -145,18 +144,25 @@ public:
return Void();
}
template <class ConfigSource>
ConfigBroadcasterImpl(ConfigSource const& configSource)
: id(deterministicRandom()->randomUniqueID()), lastCompactedVersion(0), mostRecentVersion(0),
cc("ConfigBroadcaster"), compactRequest("CompactRequest", cc),
successfulChangeRequest("SuccessfulChangeRequest", cc), failedChangeRequest("FailedChangeRequest", cc),
snapshotRequest("SnapshotRequest", cc) {
logger = traceCounters(
"ConfigBroadcasterMetrics", id, SERVER_KNOBS->WORKER_LOGGING_INTERVAL, &cc, "ConfigBroadcasterMetrics");
ConfigBroadcasterImpl(ConfigFollowerInterface const& configSource) : ConfigBroadcasterImpl() {
consumer = IConfigConsumer::createSimple(configSource, 0.5, Optional<double>{});
TraceEvent(SevDebug, "BroadcasterStartingConsumer", id).detail("Consumer", consumer->getID());
}
ConfigBroadcasterImpl(ServerCoordinators const& configSource, Optional<bool> useTestConfigDB)
: ConfigBroadcasterImpl() {
if (useTestConfigDB.present()) {
if (useTestConfigDB.get()) {
consumer = IConfigConsumer::createSimple(configSource, 0.5, Optional<double>{});
} else {
consumer = IConfigConsumer::createPaxos(configSource, 0.5, Optional<double>{});
}
TraceEvent(SevDebug, "BroadcasterStartingConsumer", id)
.detail("Consumer", consumer->getID())
.detail("UsingSimpleConsumer", useTestConfigDB.get());
}
}
JsonBuilderObject getStatus() const {
JsonBuilderObject result;
JsonBuilderArray mutationsArray;
@ -165,10 +171,9 @@ public:
mutationObject["version"] = versionedMutation.version;
const auto& mutation = versionedMutation.mutation;
mutationObject["description"] = mutation.getDescription();
mutationObject["config_class"] = configClassToString(mutation.getConfigClass());
mutationObject["config_class"] = mutation.getConfigClass().orDefault("<global>"_sr);
mutationObject["knob_name"] = mutation.getKnobName();
mutationObject["knob_value"] =
mutation.getValue().present() ? mutation.getValue().get().toString() : "<cleared>";
mutationObject["knob_value"] = mutation.getValue().orDefault("<cleared>"_sr);
mutationObject["timestamp"] = mutation.getTimestamp();
mutationsArray.push_back(std::move(mutationObject));
}
@ -183,7 +188,7 @@ public:
for (const auto& [knobName, knobValue] : kvs) {
kvsObject[knobName] = knobValue;
}
snapshotObject[configClassToString(configClass)] = std::move(kvsObject);
snapshotObject[configClass.orDefault("<global>"_sr)] = std::move(kvsObject);
}
result["snapshot"] = std::move(snapshotObject);
result["last_compacted_version"] = lastCompactedVersion;
@ -197,8 +202,8 @@ public:
ConfigBroadcaster::ConfigBroadcaster(ConfigFollowerInterface const& cfi)
: impl(std::make_unique<ConfigBroadcasterImpl>(cfi)) {}
ConfigBroadcaster::ConfigBroadcaster(ServerCoordinators const& coordinators)
: impl(std::make_unique<ConfigBroadcasterImpl>(coordinators)) {}
ConfigBroadcaster::ConfigBroadcaster(ServerCoordinators const& coordinators, Optional<bool> useTestConfigDB)
: impl(std::make_unique<ConfigBroadcasterImpl>(coordinators, useTestConfigDB)) {}
ConfigBroadcaster::ConfigBroadcaster(ConfigBroadcaster&&) = default;
@ -223,3 +228,7 @@ Future<Void> ConfigBroadcaster::setSnapshot(std::map<ConfigKey, Value>&& snapsho
UID ConfigBroadcaster::getID() const {
return impl->getID();
}
JsonBuilderObject ConfigBroadcaster::getStatus() const {
return impl->getStatus();
}

View File

@ -21,6 +21,7 @@
#pragma once
#include "fdbclient/CoordinationInterface.h"
#include "fdbclient/JsonBuilder.h"
#include "fdbserver/CoordinationInterface.h"
#include "fdbserver/ConfigFollowerInterface.h"
#include "flow/flow.h"
@ -31,7 +32,7 @@ class ConfigBroadcaster {
public:
explicit ConfigBroadcaster(ConfigFollowerInterface const&);
explicit ConfigBroadcaster(ServerCoordinators const&);
explicit ConfigBroadcaster(ServerCoordinators const&, Optional<bool> useTestConfigDB);
ConfigBroadcaster(ConfigBroadcaster&&);
ConfigBroadcaster& operator=(ConfigBroadcaster&&);
~ConfigBroadcaster();
@ -40,4 +41,5 @@ public:
Version mostRecentVersion);
Future<Void> setSnapshot(std::map<ConfigKey, Value>&& snapshot, Version snapshotVersion);
UID getID() const;
JsonBuilderObject getStatus() const;
};

View File

@ -229,7 +229,7 @@ ACTOR Future<ISimulator::KillType> simulatedFDBDRebooter(Reference<ClusterConnec
whitelistBinPaths,
"",
{},
{}));
{ true }));
}
if (runBackupAgents != AgentNone) {
futures.push_back(runBackup(connFile));

View File

@ -2643,7 +2643,8 @@ ACTOR Future<StatusReply> clusterGetStatus(
std::map<NetworkAddress, std::pair<double, OpenDatabaseRequest>>* clientStatus,
ServerCoordinators coordinators,
std::vector<NetworkAddress> incompatibleConnections,
Version datacenterVersionDifference) {
Version datacenterVersionDifference,
ConfigBroadcaster const* configBroadcaster) {
state double tStart = timer();
state JsonBuilderArray messages;
@ -2874,6 +2875,8 @@ ACTOR Future<StatusReply> clusterGetStatus(
statusObj["workload"] = workerStatuses[1];
statusObj["layers"] = workerStatuses[2];
// TODO: Read from coordinators for more up-to-date config database status?
statusObj["configuration_database"] = configBroadcaster->getStatus();
// Add qos section if it was populated
if (!qos.empty())

View File

@ -23,6 +23,7 @@
#pragma once
#include "fdbrpc/fdbrpc.h"
#include "fdbserver/ConfigBroadcaster.h"
#include "fdbserver/WorkerInterface.actor.h"
#include "fdbserver/MasterInterface.h"
#include "fdbclient/ClusterInterface.h"
@ -42,6 +43,7 @@ Future<StatusReply> clusterGetStatus(
std::map<NetworkAddress, std::pair<double, OpenDatabaseRequest>>* const& clientStatus,
ServerCoordinators const& coordinators,
std::vector<NetworkAddress> const& incompatibleConnections,
Version const& datacenterVersionDifference);
Version const& datacenterVersionDifference,
ConfigBroadcaster const* const& conifgBroadcaster);
#endif

View File

@ -832,7 +832,8 @@ ACTOR Future<Void> clusterController(Reference<ClusterConnectionFile> ccf,
Reference<AsyncVar<Optional<ClusterControllerFullInterface>>> currentCC,
Reference<AsyncVar<ClusterControllerPriorityInfo>> asyncPriorityInfo,
Future<Void> recoveredDiskFiles,
LocalityData locality);
LocalityData locality,
Optional<bool> useTestConfigDB);
// These servers are started by workerServer
class IKeyValueStore;

View File

@ -973,7 +973,7 @@ struct CLIOptions {
const char* blobCredsFromENV = nullptr;
std::string configPath;
Optional<bool> useTestConfigDB{ false };
Optional<bool> useTestConfigDB{ true };
Reference<ClusterConnectionFile> connectionFile;
Standalone<StringRef> machineId;

View File

@ -1963,7 +1963,8 @@ ACTOR Future<Void> monitorLeaderRemotelyWithDelayedCandidacy(
Reference<AsyncVar<ClusterControllerPriorityInfo>> asyncPriorityInfo,
Future<Void> recoveredDiskFiles,
LocalityData locality,
Reference<AsyncVar<ServerDBInfo>> dbInfo) {
Reference<AsyncVar<ServerDBInfo>> dbInfo,
Optional<bool> useTestConfigDB) {
state Future<Void> monitor = monitorLeaderRemotely(connFile, currentCC);
state Future<Void> timeout;
@ -1989,7 +1990,8 @@ ACTOR Future<Void> monitorLeaderRemotelyWithDelayedCandidacy(
: Never())) {}
when(wait(timeout.isValid() ? timeout : Never())) {
monitor.cancel();
wait(clusterController(connFile, currentCC, asyncPriorityInfo, recoveredDiskFiles, locality));
wait(clusterController(
connFile, currentCC, asyncPriorityInfo, recoveredDiskFiles, locality, useTestConfigDB));
return Void();
}
}
@ -2022,7 +2024,11 @@ ACTOR Future<Void> fdbd(Reference<ClusterConnectionFile> connFile,
state vector<Future<Void>> actors;
state Promise<Void> recoveredDiskFiles;
state LocalConfiguration localConfig(configPath, manualKnobOverrides);
wait(localConfig.initialize(dataFolder, deterministicRandom()->randomUniqueID()));
if (useTestConfigDB.present()) {
// TODO: Shouldn't block here
wait(localConfig.initialize(dataFolder, deterministicRandom()->randomUniqueID()));
}
actors.push_back(serveProtocolInfo());
@ -2069,13 +2075,18 @@ ACTOR Future<Void> fdbd(Reference<ClusterConnectionFile> connFile,
actors.push_back(reportErrors(monitorLeader(connFile, cc), "ClusterController"));
} else if (processClass.machineClassFitness(ProcessClass::ClusterController) == ProcessClass::WorstFit &&
SERVER_KNOBS->MAX_DELAY_CC_WORST_FIT_CANDIDACY_SECONDS > 0) {
actors.push_back(
reportErrors(monitorLeaderRemotelyWithDelayedCandidacy(
connFile, cc, asyncPriorityInfo, recoveredDiskFiles.getFuture(), localities, dbInfo),
"ClusterController"));
actors.push_back(reportErrors(monitorLeaderRemotelyWithDelayedCandidacy(connFile,
cc,
asyncPriorityInfo,
recoveredDiskFiles.getFuture(),
localities,
dbInfo,
useTestConfigDB),
"ClusterController"));
} else {
actors.push_back(reportErrors(
clusterController(connFile, cc, asyncPriorityInfo, recoveredDiskFiles.getFuture(), localities),
clusterController(
connFile, cc, asyncPriorityInfo, recoveredDiskFiles.getFuture(), localities, useTestConfigDB),
"ClusterController"));
}
actors.push_back(reportErrors(extractClusterInterface(cc, ci), "ExtractClusterInterface"));