diff --git a/fdbserver/ClusterController.actor.cpp b/fdbserver/ClusterController.actor.cpp index cc0c7e6f77..2c8e871b38 100644 --- a/fdbserver/ClusterController.actor.cpp +++ b/fdbserver/ClusterController.actor.cpp @@ -2718,7 +2718,8 @@ ACTOR Future timeKeeper(ClusterControllerData* self) { ACTOR Future statusServer(FutureStream 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 statusServer(FutureStream 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 dbInfoUpdater(ClusterControllerData* self) { ACTOR Future clusterControllerCore(ClusterControllerFullInterface interf, Future leaderFail, ServerCoordinators coordinators, - LocalityData locality) { + LocalityData locality, + Optional useTestConfigDB) { state ClusterControllerData self(interf, locality, coordinators); - state ConfigBroadcaster configBroadcaster(coordinators); + state ConfigBroadcaster configBroadcaster(coordinators, useTestConfigDB); state Future coordinationPingDelay = delay(SERVER_KNOBS->WORKER_COORDINATION_PING_DELAY); state uint64_t step = 0; state Future> 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 clusterController(ServerCoordinators coordinators, Reference>> currentCC, bool hasConnected, Reference> asyncPriorityInfo, - LocalityData locality) { + LocalityData locality, + Optional useTestConfigDB) { loop { state ClusterControllerFullInterface cci; state bool inRole = false; @@ -3571,7 +3578,7 @@ ACTOR Future 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 clusterController(Reference connFile, Reference>> currentCC, Reference> asyncPriorityInfo, Future recoveredDiskFiles, - LocalityData locality) { + LocalityData locality, + Optional 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 diff --git a/fdbserver/ConfigBroadcaster.actor.cpp b/fdbserver/ConfigBroadcaster.actor.cpp index cd45f01f8c..8ec43f511e 100644 --- a/fdbserver/ConfigBroadcaster.actor.cpp +++ b/fdbserver/ConfigBroadcaster.actor.cpp @@ -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 const& configClassSet, Optional return !configClassSet.present() || !configClass.present() || configClassSet.get().contains(configClass.get()); } -template -std::string configClassToString(Optional const& configClass) { - if (configClass.present()) { - return configClass.get().toString(); - } else { - return ""; - } -} - } // 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 serve(ConfigBroadcaster* self, ConfigFollowerInterface const& cfi) { return serve(self, this, cfi); } @@ -145,18 +144,25 @@ public: return Void(); } - template - 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{}); TraceEvent(SevDebug, "BroadcasterStartingConsumer", id).detail("Consumer", consumer->getID()); } + ConfigBroadcasterImpl(ServerCoordinators const& configSource, Optional useTestConfigDB) + : ConfigBroadcasterImpl() { + if (useTestConfigDB.present()) { + if (useTestConfigDB.get()) { + consumer = IConfigConsumer::createSimple(configSource, 0.5, Optional{}); + } else { + consumer = IConfigConsumer::createPaxos(configSource, 0.5, Optional{}); + } + 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(""_sr); mutationObject["knob_name"] = mutation.getKnobName(); - mutationObject["knob_value"] = - mutation.getValue().present() ? mutation.getValue().get().toString() : ""; + mutationObject["knob_value"] = mutation.getValue().orDefault(""_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(""_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(cfi)) {} -ConfigBroadcaster::ConfigBroadcaster(ServerCoordinators const& coordinators) - : impl(std::make_unique(coordinators)) {} +ConfigBroadcaster::ConfigBroadcaster(ServerCoordinators const& coordinators, Optional useTestConfigDB) + : impl(std::make_unique(coordinators, useTestConfigDB)) {} ConfigBroadcaster::ConfigBroadcaster(ConfigBroadcaster&&) = default; @@ -223,3 +228,7 @@ Future ConfigBroadcaster::setSnapshot(std::map&& snapsho UID ConfigBroadcaster::getID() const { return impl->getID(); } + +JsonBuilderObject ConfigBroadcaster::getStatus() const { + return impl->getStatus(); +} diff --git a/fdbserver/ConfigBroadcaster.h b/fdbserver/ConfigBroadcaster.h index dd34bfda6c..ed780606ce 100644 --- a/fdbserver/ConfigBroadcaster.h +++ b/fdbserver/ConfigBroadcaster.h @@ -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 useTestConfigDB); ConfigBroadcaster(ConfigBroadcaster&&); ConfigBroadcaster& operator=(ConfigBroadcaster&&); ~ConfigBroadcaster(); @@ -40,4 +41,5 @@ public: Version mostRecentVersion); Future setSnapshot(std::map&& snapshot, Version snapshotVersion); UID getID() const; + JsonBuilderObject getStatus() const; }; diff --git a/fdbserver/SimulatedCluster.actor.cpp b/fdbserver/SimulatedCluster.actor.cpp index d30d42773f..6cfb4ecda4 100644 --- a/fdbserver/SimulatedCluster.actor.cpp +++ b/fdbserver/SimulatedCluster.actor.cpp @@ -229,7 +229,7 @@ ACTOR Future simulatedFDBDRebooter(Reference clusterGetStatus( std::map>* clientStatus, ServerCoordinators coordinators, std::vector incompatibleConnections, - Version datacenterVersionDifference) { + Version datacenterVersionDifference, + ConfigBroadcaster const* configBroadcaster) { state double tStart = timer(); state JsonBuilderArray messages; @@ -2874,6 +2875,8 @@ ACTOR Future 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()) diff --git a/fdbserver/Status.h b/fdbserver/Status.h index 5cb9a74de1..3cfb019a8e 100644 --- a/fdbserver/Status.h +++ b/fdbserver/Status.h @@ -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 clusterGetStatus( std::map>* const& clientStatus, ServerCoordinators const& coordinators, std::vector const& incompatibleConnections, - Version const& datacenterVersionDifference); + Version const& datacenterVersionDifference, + ConfigBroadcaster const* const& conifgBroadcaster); #endif diff --git a/fdbserver/WorkerInterface.actor.h b/fdbserver/WorkerInterface.actor.h index 58e9810d82..5c6b9a343d 100644 --- a/fdbserver/WorkerInterface.actor.h +++ b/fdbserver/WorkerInterface.actor.h @@ -832,7 +832,8 @@ ACTOR Future clusterController(Reference ccf, Reference>> currentCC, Reference> asyncPriorityInfo, Future recoveredDiskFiles, - LocalityData locality); + LocalityData locality, + Optional useTestConfigDB); // These servers are started by workerServer class IKeyValueStore; diff --git a/fdbserver/fdbserver.actor.cpp b/fdbserver/fdbserver.actor.cpp index 9ad65669f0..35551602f4 100644 --- a/fdbserver/fdbserver.actor.cpp +++ b/fdbserver/fdbserver.actor.cpp @@ -973,7 +973,7 @@ struct CLIOptions { const char* blobCredsFromENV = nullptr; std::string configPath; - Optional useTestConfigDB{ false }; + Optional useTestConfigDB{ true }; Reference connectionFile; Standalone machineId; diff --git a/fdbserver/worker.actor.cpp b/fdbserver/worker.actor.cpp index 4f041781f6..ba636825de 100644 --- a/fdbserver/worker.actor.cpp +++ b/fdbserver/worker.actor.cpp @@ -1963,7 +1963,8 @@ ACTOR Future monitorLeaderRemotelyWithDelayedCandidacy( Reference> asyncPriorityInfo, Future recoveredDiskFiles, LocalityData locality, - Reference> dbInfo) { + Reference> dbInfo, + Optional useTestConfigDB) { state Future monitor = monitorLeaderRemotely(connFile, currentCC); state Future timeout; @@ -1989,7 +1990,8 @@ ACTOR Future 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 fdbd(Reference connFile, state vector> actors; state Promise 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 fdbd(Reference 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"));