From 051362c62c2803dfd1ba11f6d01b4c50a725c47c Mon Sep 17 00:00:00 2001 From: Mark Shabanov Date: Thu, 20 Aug 2026 14:27:25 +0300 Subject: [PATCH 01/12] Report degraded multi-region status when failover is stuck on an unavailable remote DC --- fdbcli/StatusCommand.cpp | 18 +- fdbclient/Schemas.cpp | 2 + fdbserver/SimulatedCluster.cpp | 28 +- .../clustercontroller/ClusterRecovery.cpp | 56 ++- fdbserver/clustercontroller/ClusterRecovery.h | 4 + fdbserver/clustercontroller/Status.cpp | 173 +++++++- fdbserver/core/ServerKnobs.cpp | 2 + fdbserver/core/include/fdbserver/core/Knobs.h | 6 + .../workloads/DegradedMultiRegionStatus.cpp | 409 ++++++++++++++++++ tests/CMakeLists.txt | 1 + tests/slow/DegradedMultiRegionStatus.toml | 37 ++ 11 files changed, 713 insertions(+), 23 deletions(-) create mode 100644 fdbserver/workloads/DegradedMultiRegionStatus.cpp create mode 100644 tests/slow/DegradedMultiRegionStatus.toml diff --git a/fdbcli/StatusCommand.cpp b/fdbcli/StatusCommand.cpp index f43d5e6e4f4..1cff7ae0213 100644 --- a/fdbcli/StatusCommand.cpp +++ b/fdbcli/StatusCommand.cpp @@ -714,10 +714,20 @@ void printStatus(StatusObjectReader statusObj, ASSERT_WE_THINK(availLoss == -1); const bool possiblyLosingData = logEpochsMayBeLosingData(statusObjCluster); if (possiblyLosingData) { - outputString += format( - "\n\n Warning: the database may have data loss and availability loss. Please " - "restart following tlog interfaces, otherwise storage servers may never be able " - "to catch up.\n"); + const std::string baseMessage = + "Please restart following tlog interfaces, otherwise storage servers " + "may never be able to catch up.\n"; + + bool degradedMultiRegion = false; + statusObjCluster.get("degraded_multi_region", degradedMultiRegion); + + const std::string header = + !degradedMultiRegion + ? "\n\n Warning: the database may have data loss and availability loss. " + : "\n\n Warning: one region is unavailable; committed data remains safe in " + "the surviving region. "; + + outputString += header + baseMessage; } else { outputString += format( "\n\n Warning: the database may have availability loss. The current log state " diff --git a/fdbclient/Schemas.cpp b/fdbclient/Schemas.cpp index df06fea0a5b..c90f571f04e 100644 --- a/fdbclient/Schemas.cpp +++ b/fdbclient/Schemas.cpp @@ -632,6 +632,7 @@ const KeyRef JSONSchemas::statusSchema = R"statusSchema( }, "active_tss_count":0, "degraded_processes":0, + "degraded_multi_region":true, "database_available":true, "database_lock_state": { "locked": true, @@ -1550,6 +1551,7 @@ file is writable and has not been overwritten externally." }, "maintenance_zone":"0ccb4e0fdbdb5583010f6b77d9d10ece", "maintenance_seconds_remaining":1.0, + "degraded_multi_region":true, "data":{ "least_operating_space_bytes_log_server":0, "average_partition_size_bytes":0, diff --git a/fdbserver/SimulatedCluster.cpp b/fdbserver/SimulatedCluster.cpp index 090b41777bb..3633b856015 100644 --- a/fdbserver/SimulatedCluster.cpp +++ b/fdbserver/SimulatedCluster.cpp @@ -481,6 +481,8 @@ class TestConfig : public BasicTestConfig { Optional generateFearless, buggify, faultInjection; Optional config; Optional remoteConfig; + // Satellite redundancy mode override applied to all regions (e.g. "one_satellite_single") + Optional satelliteRedundancyMode; bool randomlyRenameZoneId = false; bool simHTTPServerEnabled = true; @@ -548,6 +550,7 @@ class TestConfig : public BasicTestConfig { .add("storageEngineType", &storageEngineType) .add("config", &config) .add("remoteConfig", &remoteConfig) + .add("satelliteRedundancyMode", &satelliteRedundancyMode) .add("buggify", &buggify) .add("faultInjection", &faultInjection) .add("StderrSeverity", &stderrSeverity) @@ -1896,7 +1899,10 @@ void SimulationConfig::setRegions(const TestConfig& testConfig) { bool needsRemote = generateFearless; if (generateFearless) { - if (datacenters > 4) { + std::string satelliteRedundancyModeStr; + if (testConfig.satelliteRedundancyMode.present()) { + satelliteRedundancyModeStr = testConfig.satelliteRedundancyMode.get(); + } else if (datacenters > 4) { // FIXME: we cannot use one satellite replication with more than one satellite per region because // canKillProcesses does not respect usable_dcs int satellite_replication_type = deterministicRandom()->randomInt(0, 3); @@ -1912,14 +1918,12 @@ void SimulationConfig::setRegions(const TestConfig& testConfig) { } case 1: { CODE_PROBE(true, "Simulated cluster using two satellite fast redundancy mode"); - primaryObj["satellite_redundancy_mode"] = "two_satellite_fast"; - remoteObj["satellite_redundancy_mode"] = "two_satellite_fast"; + satelliteRedundancyModeStr = "two_satellite_fast"; break; } case 2: { CODE_PROBE(true, "Simulated cluster using two satellite safe redundancy mode"); - primaryObj["satellite_redundancy_mode"] = "two_satellite_safe"; - remoteObj["satellite_redundancy_mode"] = "two_satellite_safe"; + satelliteRedundancyModeStr = "two_satellite_safe"; break; } default: @@ -1939,20 +1943,17 @@ void SimulationConfig::setRegions(const TestConfig& testConfig) { } case 2: { CODE_PROBE(true, "Simulated cluster using single satellite redundancy mode"); - primaryObj["satellite_redundancy_mode"] = "one_satellite_single"; - remoteObj["satellite_redundancy_mode"] = "one_satellite_single"; + satelliteRedundancyModeStr = "one_satellite_single"; break; } case 3: { CODE_PROBE(true, "Simulated cluster using double satellite redundancy mode"); - primaryObj["satellite_redundancy_mode"] = "one_satellite_double"; - remoteObj["satellite_redundancy_mode"] = "one_satellite_double"; + satelliteRedundancyModeStr = "one_satellite_double"; break; } case 4: { CODE_PROBE(true, "Simulated cluster using triple satellite redundancy mode"); - primaryObj["satellite_redundancy_mode"] = "one_satellite_triple"; - remoteObj["satellite_redundancy_mode"] = "one_satellite_triple"; + satelliteRedundancyModeStr = "one_satellite_triple"; break; } default: @@ -1960,6 +1961,11 @@ void SimulationConfig::setRegions(const TestConfig& testConfig) { } } + if (!satelliteRedundancyModeStr.empty()) { + primaryObj["satellite_redundancy_mode"] = satelliteRedundancyModeStr; + remoteObj["satellite_redundancy_mode"] = satelliteRedundancyModeStr; + } + // Calculate the maximum satellite_logs we can support based on available machines bool useNormalDCsAsSatellites = datacenters > 4 && testConfig.minimumRegions < 2 && deterministicRandom()->random01() < 0.3; diff --git a/fdbserver/clustercontroller/ClusterRecovery.cpp b/fdbserver/clustercontroller/ClusterRecovery.cpp index d1de04e86dd..086af3f045f 100644 --- a/fdbserver/clustercontroller/ClusterRecovery.cpp +++ b/fdbserver/clustercontroller/ClusterRecovery.cpp @@ -514,6 +514,10 @@ Future trackTlogRecovery(Reference self, DBRecoveryCount recoverCount = self->cstate.myDBState.recoveryCount + 1; DatabaseConfiguration configuration = self->configuration; // self-configuration can be changed by configurationMonitor so we need a copy + // Start of the current remote-region stall, if any. Kept across loop iterations so that + // re-emissions of the event on unrelated core-state changes do not reset the reported + // stall duration; it grows monotonically until the log set is complete. + Optional remoteLogsMissingSince; while (true) { DBCoreState newState; self->logSystem->toCoreState(newState); @@ -581,6 +585,46 @@ Future trackTlogRecovery(Reference self, .trackLatest(self->clusterRecoveryStateEventHolder->trackingKey); } + // Surface "degraded multi-region" as a first-class signal: with usableRegions > 1, the remote region's + // log set has not been recruited (allLogs == false). Committed data is still safe in the surviving + // region, so this is distinct from data loss. Note that oldTLogData is deliberately not a + // discriminator: when the remote region is down, old log generations cannot be purged (finalUpdate + // requires allLogs), so oldTLogData stays non-empty precisely in this stalled state. The value is also + // (re)set when all logs are recruited or the cluster fully recovers, so the trackLatest event always + // reflects the current state. + // + // The remote log set is legitimately absent before recovery reaches accepting_commits even during a + // normal (non-region-loss) recovery, because remote tlog recruitment is part of the normal recovery + // path. Only begin tracking the stall once the cluster is accepting commits: at that point a missing + // remote log set is no longer a transient recruiting artifact but a genuine degraded multi-region + // condition. Without this gate, StallSeconds accumulates across the pre-accepting_commits recruiting + // window and a slow-but-healthy recovery trips degraded_multi_region once the stall exceeds + // DEGRADED_MULTI_REGION_MIN_STALL_SECONDS. + bool remoteRegionLogsMissing = + configuration.usableRegions > 1 && !allLogs && self->recoveryState >= RecoveryState::ACCEPTING_COMMITS; + if (remoteRegionLogsMissing) { + if (!remoteLogsMissingSince.present()) { + remoteLogsMissingSince = now(); + } + } else { + remoteLogsMissingSince = Optional(); + } + // StallSeconds is carried in the event itself rather than derived from the event's + // emission Time by status: the event is re-emitted on every core-state change, and + // deriving the duration from the latest emission would reset the stall counter while + // the stall is still ongoing. + double remoteRegionStallSeconds = + remoteLogsMissingSince.present() ? std::max(0.0, now() - remoteLogsMissingSince.get()) : 0.0; + TraceEvent( + getRecoveryEventName(ClusterRecoveryEventType::CLUSTER_RECOVERY_REMOTE_REGION_STALL_EVENT_NAME).c_str(), + self->dbgid) + .detail("RemoteRegionLogsMissing", remoteRegionLogsMissing) + .detail("AllLogs", allLogs) + .detail("OldTLogDataSize", newState.oldTLogData.size()) + .detail("UsableRegions", configuration.usableRegions) + .detail("StallSeconds", remoteRegionStallSeconds) + .trackLatest(self->clusterRecoveryRemoteRegionStallEventHolder->trackingKey); + self->registrationTrigger.trigger(); if (finalUpdate) { @@ -589,7 +633,15 @@ Future trackTlogRecovery(Reference self, co_return; } - co_await changed; + // While the remote log set is missing, keep re-emitting the stall event on a fixed cadence so + // StallSeconds keeps advancing even when no core-state change would otherwise wake this loop. + // Without this, an idle (commits-stalled) cluster would block forever in co_await changed and + // the sustained-stall threshold could never be crossed. + if (remoteRegionLogsMissing) { + co_await (changed || delay(SERVER_KNOBS->DEGRADED_MULTI_REGION_REFRESH_SECONDS)); + } else { + co_await changed; + } } } @@ -2099,6 +2151,8 @@ const std::string& getRecoveryEventName(ClusterRecoveryEventType type) { SERVER_KNOBS->CLUSTER_RECOVERY_EVENT_NAME_PREFIX + "RecoveryAvailable" }); recoveryEventNameMap.insert({ ClusterRecoveryEventType::CLUSTER_RECOVERY_METRICS_EVENT_NAME, SERVER_KNOBS->CLUSTER_RECOVERY_EVENT_NAME_PREFIX + "RecoveryMetrics" }); + recoveryEventNameMap.insert({ ClusterRecoveryEventType::CLUSTER_RECOVERY_REMOTE_REGION_STALL_EVENT_NAME, + SERVER_KNOBS->CLUSTER_RECOVERY_EVENT_NAME_PREFIX + "RecoveryRemoteRegionStall" }); } auto iter = recoveryEventNameMap.find(type); diff --git a/fdbserver/clustercontroller/ClusterRecovery.h b/fdbserver/clustercontroller/ClusterRecovery.h index ac6f4096619..c62f6aaca5e 100644 --- a/fdbserver/clustercontroller/ClusterRecovery.h +++ b/fdbserver/clustercontroller/ClusterRecovery.h @@ -55,6 +55,7 @@ enum ClusterRecoveryEventType { CLUSTER_RECOVERY_COMMIT_EVENT_NAME, CLUSTER_RECOVERY_AVAILABLE_EVENT_NAME, CLUSTER_RECOVERY_METRICS_EVENT_NAME, + CLUSTER_RECOVERY_REMOTE_REGION_STALL_EVENT_NAME, CLUSTER_RECOVERY_LAST // Always the last entry }; @@ -254,6 +255,7 @@ struct ClusterRecoveryData : NonCopyable, ReferenceCounted Reference clusterRecoveryGenerationsEventHolder; Reference clusterRecoveryDurationEventHolder; Reference clusterRecoveryAvailableEventHolder; + Reference clusterRecoveryRemoteRegionStallEventHolder; ClusterRecoveryData(ClusterControllerData* controllerData, Reference const> const& dbInfo, @@ -290,6 +292,8 @@ struct ClusterRecoveryData : NonCopyable, ReferenceCounted getRecoveryEventName(ClusterRecoveryEventType::CLUSTER_RECOVERY_DURATION_EVENT_NAME)); clusterRecoveryAvailableEventHolder = makeReference( getRecoveryEventName(ClusterRecoveryEventType::CLUSTER_RECOVERY_AVAILABLE_EVENT_NAME)); + clusterRecoveryRemoteRegionStallEventHolder = makeReference( + getRecoveryEventName(ClusterRecoveryEventType::CLUSTER_RECOVERY_REMOTE_REGION_STALL_EVENT_NAME)); logger = cc.traceCounters(getRecoveryEventName(ClusterRecoveryEventType::CLUSTER_RECOVERY_METRICS_EVENT_NAME), dbgid, diff --git a/fdbserver/clustercontroller/Status.cpp b/fdbserver/clustercontroller/Status.cpp index 71af46f3cf8..99a1877a1c4 100644 --- a/fdbserver/clustercontroller/Status.cpp +++ b/fdbserver/clustercontroller/Status.cpp @@ -1202,13 +1202,54 @@ static JsonBuilderObject clientStatusFetcher( return clientStatus; } +// Parse the recovery-provided remote-region-stall event into a boolean. The event is emitted by +// trackTlogRecovery when the remote region's log set could not be recruited (allLogs == false with +// usableRegions > 1). An absent or empty event (e.g. an older server that does not emit it yet) means the +// signal is not present, so the cluster is not flagged as degraded. +static bool parseRemoteRegionLogsMissing(const TraceEventFields& remoteRegionStallEvent) { + if (remoteRegionStallEvent.size() == 0) { + return false; + } + return remoteRegionStallEvent.getInt("RemoteRegionLogsMissing") != 0; +} + +// Parse the duration (in seconds) for which the remote-region-stall signal has been active. The recovery +// side carries the duration in the event itself ("StallSeconds") because the event is (re)emitted on every +// core-state change; deriving the duration from the latest emission time would reset the stall counter +// while the stall is still ongoing. An absent event, an event reporting that the logs are not missing, or +// any malformed value conservatively reports 0, which keeps the cluster from being flagged as degraded. +static double parseRemoteRegionStallSeconds(const TraceEventFields& remoteRegionStallEvent) { + try { + if (remoteRegionStallEvent.size() == 0) { + return 0.0; + } + if (remoteRegionStallEvent.getInt("RemoteRegionLogsMissing") == 0) { + return 0.0; + } + double stallSeconds = std::stod(remoteRegionStallEvent.getValue("StallSeconds")); + return stallSeconds > 0.0 ? stallSeconds : 0.0; + } catch (Error& e) { + if (e.code() == error_code_actor_cancelled) { + throw; + } + return 0.0; + } catch (std::exception&) { + // Malformed numeric fields fail safe to "no sustained stall". + return 0.0; + } +} + static AsyncResult recoveryStateStatusFetcher(Database cx, WorkerDetails ccWorker, WorkerDetails mWorker, int workerCount, std::set* incomplete_reasons, - int* statusCode) { + int* statusCode, + bool* remoteRegionLogsMissing, + double* remoteRegionStallSeconds) { JsonBuilderObject message; + *remoteRegionLogsMissing = false; + *remoteRegionStallSeconds = 0.0; Transaction tr(cx); try { Future mdActiveGensF = @@ -1223,10 +1264,24 @@ static AsyncResult recoveryStateStatusFetcher(Database cx, timeoutError(ccWorker.interf.eventLogRequest.getReply(EventLogRequest(StringRef( getRecoveryEventName(ClusterRecoveryEventType::CLUSTER_RECOVERY_AVAILABLE_EVENT_NAME)))), 1.0); + Future remoteRegionStallF = + timeoutError(ccWorker.interf.eventLogRequest.getReply(EventLogRequest(StringRef(getRecoveryEventName( + ClusterRecoveryEventType::CLUSTER_RECOVERY_REMOTE_REGION_STALL_EVENT_NAME)))), + 1.0); tr.setOption(FDBTransactionOptions::PRIORITY_SYSTEM_IMMEDIATE); Future> rvF = errorOr(timeoutError(tr.getReadVersion(), 1.0)); - co_await (success(mdActiveGensF) && success(mdF) && success(rvF) && success(mDBAvailableF)); + co_await (success(mdActiveGensF) && success(mdF) && success(rvF) && success(mDBAvailableF) && + success(remoteRegionStallF)); + + // Precisely report whether recovery is stalled because the remote region's log set could not be + // recruited (allLogs == false with usableRegions > 1). This replaces the coarser heuristic of matching + // on the accepting_commits recovery state, which would also flag a normal transient pass through + // accepting_commits as degraded. Additionally report how long the stall has been active, because the + // signal is also transiently true during a normal (non-initial) recovery of a healthy multi-region + // cluster until the remote epoch is recruited. + *remoteRegionLogsMissing = parseRemoteRegionLogsMissing(remoteRegionStallF.get()); + *remoteRegionStallSeconds = parseRemoteRegionStallSeconds(remoteRegionStallF.get()); const TraceEventFields& md = mdF.get(); int mStatusCode = md.getInt("StatusCode"); @@ -1805,6 +1860,7 @@ static JsonBuilderObject configurationFetcher(Optional co static AsyncResult dataStatusFetcher(WorkerDetails ddWorker, DatabaseConfiguration configuration, + bool degradedMultiRegion, int* minStorageReplicasRemaining) { JsonBuilderObject statusObjData; @@ -1831,7 +1887,9 @@ static AsyncResult dataStatusFetcher(WorkerDetails ddWorker, if (startingStats.size() && startingStats.getValue("State") != "Active") { JsonBuilderObject stateSectionObj; stateSectionObj["name"] = "initializing"; - stateSectionObj["description"] = "(Re)initializing automatic data distribution"; + stateSectionObj["description"] = degradedMultiRegion + ? "Degraded multiregional (Re)initializing automatic data distribution" + : "(Re)initializing automatic data distribution"; statusObjData["state"] = stateSectionObj; co_return statusObjData; } @@ -3070,9 +3128,17 @@ static AsyncResult clusterGetStatusImpl(Reference s // construct status information for cluster subsections int statusCode = (int)RecoveryStatus::END; + bool remoteRegionLogsMissing = false; + double remoteRegionStallSeconds = 0.0; std::vector>> clusterSubsectionFetchers; - clusterSubsectionFetchers.push_back(errorOr(recoveryStateStatusFetcher( - cx, ccWorker, mWorker, workers.size(), &status_incomplete_reasons, &statusCode))); + clusterSubsectionFetchers.push_back(errorOr(recoveryStateStatusFetcher(cx, + ccWorker, + mWorker, + workers.size(), + &status_incomplete_reasons, + &statusCode, + &remoteRegionLogsMissing, + &remoteRegionStallSeconds))); clusterSubsectionFetchers.push_back(errorOr(timeoutError(getIdmpKeyStatus(cx), 5.0))); clusterSubsectionFetchers.push_back(errorOr(versionEpochStatusFetcher(cx, &status_incomplete_reasons))); @@ -3210,9 +3276,22 @@ static AsyncResult clusterGetStatusImpl(Reference s int fullyReplicatedRegions = -1; // NOTE: here we should start all the transaction before wait in order to overlay latency Future> primaryDCFO = getActivePrimaryDC(cx, &fullyReplicatedRegions, &messages); + // The cluster is "degraded multi-region" when recovery explicitly reports that it is stalled because + // the remote region's log set could not be recruited (allLogs == false with usableRegions > 1) after + // accepting commits, and the stall has persisted for a while. Committed data is still safe in the + // surviving region (and its satellite), so this is a degraded-but-not-data-losing state. Requiring + // the recovery state to be at (or past) accepting_commits avoids the transient phases before the + // cluster is actually accepting commits. The sustained-stall requirement avoids the transient window + // during a normal (non-initial) recovery of a healthy multi-region cluster, where the remote log set + // is legitimately absent at accepting_commits until the remote epoch finishes recruiting. + bool degradedMultiRegion = + configuration.present() && statusCode >= RecoveryStatus::accepting_commits && + configuration.get().usableRegions > 1 && remoteRegionLogsMissing && + remoteRegionStallSeconds >= SERVER_KNOBS->DEGRADED_MULTI_REGION_MIN_STALL_SECONDS; + statusObj["degraded_multi_region"] = degradedMultiRegion; std::vector> statusSectionFetchers; statusSectionFetchers.push_back( - dataStatusFetcher(ddWorker, configuration.get(), &minStorageReplicasRemaining)); + dataStatusFetcher(ddWorker, configuration.get(), degradedMultiRegion, &minStorageReplicasRemaining)); statusSectionFetchers.push_back(workloadStatusFetcher( db, workers, mWorker, rkWorker, &qos, &dataOverlay, &status_incomplete_reasons, storageServerFuture)); statusSectionFetchers.push_back(layerStatusFetcher(cx, &messages, &status_incomplete_reasons)); @@ -3551,7 +3630,7 @@ StatusReply clusterGetFaultToleranceStatus(const std::string& statusStr) { std::string faultToleranceRelatedFields[] = { "fault_tolerance", "data", "logs", "maintenance_zone", "maintenance_seconds_remaining", "qos", - "recovery_state", "messages" + "recovery_state", "messages", "degraded_multi_region" }; JsonBuilderObject statusObj; @@ -3953,6 +4032,86 @@ TEST_CASE("/status/json/merging") { return Void(); } +// Verify the "degraded multi-region" signal parsing and computation, exercising the actual code paths used by +// recoveryStateStatusFetcher and clusterGetStatusImpl. A cluster is only considered degraded multi-region when +// recovery explicitly reports that the remote region's log set could not be recruited (allLogs == false with +// usableRegions > 1) and the stall has persisted long enough to rule out the transient window that occurs +// during a normal, non-initial recovery of a healthy multi-region cluster. +TEST_CASE("/fdbserver/clustercontroller/degradedMultiRegionComputation") { + // parseRemoteRegionLogsMissing: event carries the signal. + { + TraceEventFields remoteRegionStallEvent; + remoteRegionStallEvent.addField("RemoteRegionLogsMissing", "1"); + ASSERT(parseRemoteRegionLogsMissing(remoteRegionStallEvent)); + } + // parseRemoteRegionLogsMissing: event explicitly says logs are not missing (normal recovery progressing). + { + TraceEventFields remoteRegionStallEvent; + remoteRegionStallEvent.addField("RemoteRegionLogsMissing", "0"); + ASSERT(!parseRemoteRegionLogsMissing(remoteRegionStallEvent)); + } + // parseRemoteRegionLogsMissing: empty event (older server that does not emit the signal) means false. + { + TraceEventFields remoteRegionStallEvent; + ASSERT(!parseRemoteRegionLogsMissing(remoteRegionStallEvent)); + } + + // Real scenario from production logs: with 2 usable regions, the first region is down, recovery is stalled + // and the remote log set has not been recruited (AllLogs=0), while old log generations remain (OldTLogDataSize=2). + // This is exactly the "degraded but not data-losing" state the signal is meant to surface. + { + TraceEventFields remoteRegionStallEvent; + remoteRegionStallEvent.addField("RemoteRegionLogsMissing", "1"); + remoteRegionStallEvent.addField("AllLogs", "0"); + remoteRegionStallEvent.addField("OldTLogDataSize", "2"); + remoteRegionStallEvent.addField("UsableRegions", "2"); + ASSERT(parseRemoteRegionLogsMissing(remoteRegionStallEvent)); + } + + // parseRemoteRegionStallSeconds: logs not missing -> no active stall, regardless of the carried duration. + { + TraceEventFields remoteRegionStallEvent; + remoteRegionStallEvent.addField("StallSeconds", "3600.0"); + remoteRegionStallEvent.addField("RemoteRegionLogsMissing", "0"); + ASSERT_EQ(parseRemoteRegionStallSeconds(remoteRegionStallEvent), 0.0); + } + // parseRemoteRegionStallSeconds: empty event -> no active stall. + { + TraceEventFields remoteRegionStallEvent; + ASSERT_EQ(parseRemoteRegionStallSeconds(remoteRegionStallEvent), 0.0); + } + // parseRemoteRegionStallSeconds: sustained stall carries the monotonic duration reported by recovery. + { + const double threshold = SERVER_KNOBS->DEGRADED_MULTI_REGION_MIN_STALL_SECONDS; + TraceEventFields remoteRegionStallEvent; + remoteRegionStallEvent.addField("StallSeconds", format("%.6f", threshold + 5.0)); + remoteRegionStallEvent.addField("RemoteRegionLogsMissing", "1"); + ASSERT_GE(parseRemoteRegionStallSeconds(remoteRegionStallEvent), threshold); + } + { + const double threshold = SERVER_KNOBS->DEGRADED_MULTI_REGION_MIN_STALL_SECONDS; + TraceEventFields remoteRegionStallEvent; + remoteRegionStallEvent.addField("StallSeconds", "0.1"); + remoteRegionStallEvent.addField("RemoteRegionLogsMissing", "1"); + ASSERT_LT(parseRemoteRegionStallSeconds(remoteRegionStallEvent), threshold); + } + // parseRemoteRegionStallSeconds: malformed or negative StallSeconds fails safe to 0 (not degraded). + { + TraceEventFields remoteRegionStallEvent; + remoteRegionStallEvent.addField("StallSeconds", "garbage"); + remoteRegionStallEvent.addField("RemoteRegionLogsMissing", "1"); + ASSERT_EQ(parseRemoteRegionStallSeconds(remoteRegionStallEvent), 0.0); + } + { + TraceEventFields remoteRegionStallEvent; + remoteRegionStallEvent.addField("StallSeconds", "-5.0"); + remoteRegionStallEvent.addField("RemoteRegionLogsMissing", "1"); + ASSERT_EQ(parseRemoteRegionStallSeconds(remoteRegionStallEvent), 0.0); + } + + co_return; +} + // Test that clusterGetStatus returns partial results within the specified deadline // even when the database is unavailable and transactions hang forever. TEST_CASE("/fdbserver/clustercontroller/clusterGetStatusTimeout") { diff --git a/fdbserver/core/ServerKnobs.cpp b/fdbserver/core/ServerKnobs.cpp index 0bd7a4e8d55..b6e089fe0be 100644 --- a/fdbserver/core/ServerKnobs.cpp +++ b/fdbserver/core/ServerKnobs.cpp @@ -876,6 +876,8 @@ void ServerKnobs::initialize(Randomize randomize, ClientKnobs* clientKnobs, IsSi bool shortRecoveryDuration = randomize && buggify(); init( ENFORCED_MIN_RECOVERY_DURATION, 0.085 ); if( shortRecoveryDuration ) ENFORCED_MIN_RECOVERY_DURATION = 0.01; init( REQUIRED_MIN_RECOVERY_DURATION, 0.080 ); if( shortRecoveryDuration ) REQUIRED_MIN_RECOVERY_DURATION = 0.01; + init( DEGRADED_MULTI_REGION_MIN_STALL_SECONDS, 30.0 ); if( isSimulated ) DEGRADED_MULTI_REGION_MIN_STALL_SECONDS = 15.0; + init( DEGRADED_MULTI_REGION_REFRESH_SECONDS, 1.0 ); init( ALWAYS_CAUSAL_READ_RISKY, false ); init( MAX_COMMIT_UPDATES, 2000 ); if( randomize && buggify() ) MAX_COMMIT_UPDATES = 1; init( MAX_PROXY_COMPUTE, 2.0 ); diff --git a/fdbserver/core/include/fdbserver/core/Knobs.h b/fdbserver/core/include/fdbserver/core/Knobs.h index 2f65564bde0..b45b262d9b0 100644 --- a/fdbserver/core/include/fdbserver/core/Knobs.h +++ b/fdbserver/core/include/fdbserver/core/Knobs.h @@ -744,6 +744,12 @@ class SWIFT_CXX_IMMORTAL_SINGLETON_TYPE ServerKnobs : public KnobsImpl 1. +// +// Scenario mirroring production: kill ONLY the primary datacenter "0" of the first +// region, leaving its satellite "2" and the remote region "1"/"3" alive. The cluster +// fails over to the remote region, storage servers recover data from the surviving +// satellite "2", and recovery gets stuck at accepting_commits because the dead +// primary's log set cannot be recruited (allLogs == false => RemoteRegionLogsMissing=1). +// After the stall persists past DEGRADED_MULTI_REGION_MIN_STALL_SECONDS the test reads +// the status JSON and asserts: +// - cluster.degraded_multi_region == true +// - cluster.data.state.description contains "Degraded multiregional" +// +// Before killing the primary DC the test also exercises the false-positive regression +// from the remote log set being transiently absent at accepting_commits: it verifies +// (a) a fully healthy cluster is not flagged, and (b) a *normal* recovery of a healthy +// cluster (triggered by killing only the sequencer, no region loss) keeps +// degraded_multi_region == false throughout the recovery. +// +// Finally it resets usable_regions=1 so subsequent checks are not stuck. + +#include "fdbclient/NativeAPI.actor.h" +#include "fdbserver/core/TesterInterface.h" +#include "fdbserver/core/WorkerInterface.h" +#include "fdbserver/tester/workloads.h" +#include "fdbserver/core/FDBSimulationPolicy.h" +#include "fdbserver/core/RecoveryState.h" +#include "fdbserver/core/ServerDBInfo.h" +#include "fdbrpc/simulator.h" +#include "fdbclient/ManagementAPI.h" +#include "flow/CoroUtils.h" +#include "fdbclient/ReadYourWrites.h" +#include "fdbclient/json_spirit/json_spirit_value.h" + +struct DegradedMultiRegionStatusWorkload : TestWorkload { + static constexpr auto NAME = "DegradedMultiRegionStatus"; + bool enabled; + double testDuration; + bool testSucceeded; + + explicit DegradedMultiRegionStatusWorkload(WorkloadContext const& wcx) : TestWorkload(wcx) { + enabled = !clientId && g_network->isSimulated(); + testDuration = getOption(options, "testDuration"_sr, 120.0); + testSucceeded = false; + } + + void disableFailureInjectionWorkloads(std::set& out) const override { out.insert("all"); } + + Future setup(Database const& cx) override { + if (enabled) { + return _setup(cx); + } + return Void(); + } + Future start(Database const& cx) override { + if (enabled) { + return killPrimaryKeepSatellites(cx); + } + return Void(); + } + Future check(Database const& cx) override { return enabled ? testSucceeded : true; } + void getMetrics(std::vector& m) override {} + + // Wait until the cluster is fully recovered (stable starting point) before killing anything. + Future _setup(Database cx) { + double failedWait = 0.0; + while (dbInfo->get().recoveryState < RecoveryState::FULLY_RECOVERED) { + if (failedWait >= 300.0) { + TraceEvent(SevError, "DegradedMultiRegionStatus_FullRecoveryTimeout") + .detail("Elapsed", failedWait) + .detail("RecoveryState", dbInfo->get().recoveryState); + ASSERT(false); + } + co_await delay(1.0); + failedWait += 1.0; + } + TraceEvent("DegradedMultiRegionStatus_Setup").log(); + } + + // Read the degraded_multi_region flag from the status JSON. Returns false on any + // transient read/parse error (treated as "not degraded") so callers only fail on an + // explicit true. + Future readDegradedFlag(Database cx) { + ReadYourWritesTransaction tr(cx); + try { + tr.setOption(FDBTransactionOptions::READ_SYSTEM_KEYS); + Optional statusVal = co_await tr.get("\xff\xff/status/json"_sr); + if (statusVal.present()) { + json_spirit::mValue mv; + json_spirit::read_string(statusVal.get().toString(), mv); + auto& root = mv.get_obj(); + if (root.contains("cluster")) { + auto& clusterObj = root["cluster"].get_obj(); + if (clusterObj.contains("degraded_multi_region")) { + co_return clusterObj["degraded_multi_region"].get_bool(); + } + } + } + } catch (Error& e) { + if (e.code() == error_code_actor_cancelled) { + throw; + } + // Transient read errors: treat as not degraded. + } + co_return false; + } + + // Verify that a healthy (both regions up) cluster is not falsely flagged as + // degraded multi-region while idle. + Future verifyHealthyNotDegraded(Database cx) { + double tStart = now(); + while (true) { + bool degraded = co_await readDegradedFlag(cx); + if (degraded) { + TraceEvent(SevError, "DegradedMultiRegionStatus_HealthyFalsePositive") + .detail("Elapsed", now() - tStart) + .detail("Phase", "idle"); + co_return false; + } + if (now() - tStart > 10.0) { + co_return true; + } + co_await delay(1.0); + } + } + + // Trigger a *normal* recovery with no region loss by killing only the sequencer + // (master), then poll the status across the recovery window and assert that + // degraded_multi_region stays false the whole time. This exercises the exact + // regression where the remote log set is transiently absent at accepting_commits. + Future verifyNormalRecoveryNotDegraded(Database cx) { + ASSERT(g_network->isSimulated()); + + // Kill the current sequencer to force a recovery while keeping every region alive. + NetworkAddress masterAddr = dbInfo->get().master.address(); + ISimulator::ProcessInfo* masterProcess = g_simulator->getProcessByAddress(masterAddr); + if (masterProcess == nullptr) { + TraceEvent(SevError, "DegradedMultiRegionStatus_MasterProcessNotFound").detail("Address", masterAddr); + co_return false; + } + LifetimeToken preKillMasterLifetime = dbInfo->get().masterLifetime; + TraceEvent("DegradedMultiRegionStatus_KillSequencer").detail("Address", masterAddr); + g_simulator->killProcess(masterProcess, ISimulator::KillType::KillInstantly); + + // Wait until the kill is observed as a new recovery start. A changed ServerDBInfo + // alone is not a reliable recovery marker (broadcasts also happen for unrelated + // reasons), so require either a new masterLifetime or a drop from FULLY_RECOVERED. + // onChange() is awaited in a loop because any single change may not belong to + // this recovery; without this wait the polling below could pass on the stale + // fully-recovered state without ever observing the new accepting_commits window. + double observedAt = now(); + constexpr double observeTimeout = 60.0; + while (true) { + // isEqual is not const-qualified, so compare against a local copy. + LifetimeToken currentMasterLifetime = dbInfo->get().masterLifetime; + if (!currentMasterLifetime.isEqual(preKillMasterLifetime) || + dbInfo->get().recoveryState < RecoveryState::FULLY_RECOVERED) { + break; // recovery started: new lifetime or the state dropped + } + if (now() - observedAt > observeTimeout) { + TraceEvent(SevError, "DegradedMultiRegionStatus_NormalRecoveryNotObserved") + .detail("RecoveryState", dbInfo->get().recoveryState); + co_return false; + } + co_await race(dbInfo->onChange(), delay(1.0)); + } + + double tStart = now(); + double deadline = 120.0; + uint8_t cycles = 10; + while (true) { + bool degraded = co_await readDegradedFlag(cx); + if (degraded) { + TraceEvent(SevError, "DegradedMultiRegionStatus_NormalRecoveryFalsePositive") + .detail("Elapsed", now() - tStart) + .detail("RecoveryState", dbInfo->get().recoveryState); + co_return false; + } + if (dbInfo->get().recoveryState >= RecoveryState::ALL_LOGS_RECRUITED) { + TraceEvent("DegradedMultiRegionStatus_NormalRecoveryClean") + .detail("Elapsed", now() - tStart) + .detail("RecoveryState", dbInfo->get().recoveryState); + if (cycles == 0) { + co_return true; + } + cycles--; + } else { + cycles = 10; // reset the hold count if recovery is still in progress + } + if (now() - tStart > deadline) { + TraceEvent(SevError, "DegradedMultiRegionStatus_NormalRecoveryTimeout") + .detail("Elapsed", now() - tStart) + .detail("RecoveryState", dbInfo->get().recoveryState); + co_return false; + } + co_await delay(1.0); + } + } + + // Wait until status JSON reports the desired degraded multi-region state. + Future waitForDegradedStatus(Database cx) { + double tStart = now(); + Optional degradedSince; + while (true) { + ReadYourWritesTransaction tr(cx); + try { + tr.setOption(FDBTransactionOptions::READ_SYSTEM_KEYS); + Optional statusVal = co_await tr.get("\xff\xff/status/json"_sr); + if (statusVal.present()) { + json_spirit::mValue mv; + json_spirit::read_string(statusVal.get().toString(), mv); + auto& root = mv.get_obj(); + if (root.contains("cluster")) { + auto& clusterObj = root["cluster"].get_obj(); + bool degraded = false; + if (clusterObj.contains("degraded_multi_region")) { + degraded = clusterObj["degraded_multi_region"].get_bool(); + } + std::string dataStateDesc; + if (clusterObj.contains("data") && clusterObj["data"].get_obj().contains("state") && + clusterObj["data"].get_obj()["state"].get_obj().contains("description")) { + dataStateDesc = clusterObj["data"].get_obj()["state"].get_obj()["description"].get_str(); + } + + bool acceptingCommits = + clusterObj.contains("recovery_state") && + clusterObj["recovery_state"].get_obj().contains("name") && + clusterObj["recovery_state"].get_obj()["name"].get_str() == "accepting_commits"; + + const bool expectedState = degraded && acceptingCommits && + dataStateDesc.find("Degraded multiregional") != std::string::npos; + + std::string recoveryStateJson; + if (clusterObj.contains("recovery_state")) { + recoveryStateJson = json_spirit::write_string(clusterObj["recovery_state"], + json_spirit::Output_options::none); + } + TraceEvent("DegradedMultiRegionStatus_Probe") + .detail("Degraded", degraded) + .detail("DataStateDesc", dataStateDesc) + .detail("AcceptingCommits", acceptingCommits) + .detail("RecoveryStateJson", recoveryStateJson) + .detail("Elapsed", now() - tStart) + .detail("ExpectedState", expectedState); + if (expectedState) { + if (!degradedSince.present()) { + degradedSince = now(); + TraceEvent("DegradedMultiRegionStatus_DegradedHoldStarted") + .detail("Elapsed", now() - tStart); + } + + const double holdDuration = now() - degradedSince.get(); + if (holdDuration >= 10.0) { + TraceEvent("DegradedMultiRegionStatus_Success") + .detail("Elapsed", now() - tStart) + .detail("Degraded", degraded) + .detail("DataStateDesc", dataStateDesc); + printf("\n=== Degraded Multi-Region Status Found ===\n"); + printf( + "Warning: one region is unavailable; committed data remains safe in the surviving " + "region.\n"); + printf( + "Please restart following tlog interfaces, otherwise storage servers may never be " + "able to catch up.\n"); + printf("\nData:\n"); + printf(" Replication health - %s\n", dataStateDesc.c_str()); + printf("========================================\n\n"); + fflush(stdout); + co_return true; + } + } else { + if (degradedSince.present()) { + TraceEvent("DegradedMultiRegionStatus_DegradedHoldLost") + .detail("Elapsed", now() - tStart) + .detail("HoldDuration", now() - degradedSince.get()); + co_return false; + } + } + } + } + } catch (Error& e) { + if (e.code() == error_code_actor_cancelled) { + throw; + } + // Transient read errors: keep polling. + } + + if (now() - tStart > testDuration) { + TraceEvent(SevError, "DegradedMultiRegionStatus_Timeout").detail("Elapsed", now() - tStart); + co_return false; + } + co_await delay(5.0); + } + } + + Future killPrimaryKeepSatellites(Database cx) { + ASSERT(g_network->isSimulated()); + + // The cluster starts fully healthy (both regions up). Confirm it is not + // falsely marked degraded while idle. + testSucceeded = co_await verifyHealthyNotDegraded(cx); + if (!testSucceeded) { + co_return; + } + + // Trigger a normal recovery with no region loss and confirm the flag stays + // false throughout. This guards the transient accepting_commits window. + testSucceeded = co_await verifyNormalRecoveryNotDegraded(cx); + if (!testSucceeded) { + co_return; + } + + // Give the cluster a moment to settle back to fully recovered before the + // destructive kill below. + co_await _setup(cx); + + // The primary region of the first region is datacenter "0" in the fearless + // topology (with primary satellite "2"). Kill only "0"; keep "2", "1" and + // "3" alive so the remote region can recover data from the surviving + // satellite of the (dead) primary region. + LifetimeToken previousMasterLifetime = dbInfo->get().masterLifetime; + g_simulator->killDataCenter("0"_sr, ISimulator::KillType::KillInstantly, true); + TraceEvent("DegradedMultiRegionStatus_KilledPrimaryDC").log(); + + bool failoverReady = true; + + try { + co_await timeoutError(waitForPrimaryDC(cx, "1"_sr), 120.0); + } catch (Error& e) { + if (e.code() == error_code_actor_cancelled) { + throw; + } + TraceEvent(SevError, "DegradedMultiRegionStatus_WaitForPrimaryDCFailed").error(e); + failoverReady = false; + } + + if (failoverReady) { + double start = now(); + LifetimeToken currentMasterLifetime = dbInfo->get().masterLifetime; + while (currentMasterLifetime.isEqual(previousMasterLifetime)) { + if (now() - start > 120.0) { + TraceEvent(SevError, "DegradedMultiRegionStatus_MasterChangeTimeout").log(); + failoverReady = false; + break; + } + co_await dbInfo->onChange(); + currentMasterLifetime = dbInfo->get().masterLifetime; + } + } + + if (failoverReady) { + testSucceeded = co_await waitForDegradedStatus(cx); + } else { + testSucceeded = false; + } + + // Unstick the cluster regardless of the assertion outcome: recovery is parked + // at accepting_commits while region "0" is dead. A forced recovery in region + // "1" runs updateConfigForForcedRecovery, which forces usable_regions=1 and + // lets the cluster fully recover, so subsequent workloads/checks do not hang. + // This is cleanup, not part of the assertion: a cleanup failure is logged but + // does not overwrite the degraded-status result. + try { + co_await forceRecovery(cx->getConnectionRecord(), "1"_sr); + TraceEvent("DegradedMultiRegionStatus_Unstick").log(); + } catch (Error& e) { + if (e.code() == error_code_actor_cancelled) { + throw; + } + TraceEvent(SevWarnAlways, "DegradedMultiRegionStatus_UnstickFailed").error(e); + } + + // Safety net in case the forced recovery did not take effect. + try { + co_await ManagementAPI::changeConfig(cx.getReference(), "usable_regions=1", true); + TraceEvent("DegradedMultiRegionStatus_Reset").log(); + } catch (Error& e) { + if (e.code() == error_code_actor_cancelled) { + throw; + } + TraceEvent(SevWarnAlways, "DegradedMultiRegionStatus_ResetFailed").error(e); + } + } +}; + +WorkloadFactory DegradedMultiRegionStatusWorkloadFactory; diff --git a/tests/CMakeLists.txt b/tests/CMakeLists.txt index bf6b1785cac..a76b2cd2499 100644 --- a/tests/CMakeLists.txt +++ b/tests/CMakeLists.txt @@ -278,6 +278,7 @@ if(WITH_PYTHON) add_fdb_test(TEST_FILES slow/GcGenerations.toml) add_fdb_test(TEST_FILES fast/KillRegionCycle.toml) add_fdb_test(TEST_FILES rare/ClogRemoteTLog.toml) + add_fdb_test(TEST_FILES slow/DegradedMultiRegionStatus.toml) endif() if(WITH_ROCKSDB) diff --git a/tests/slow/DegradedMultiRegionStatus.toml b/tests/slow/DegradedMultiRegionStatus.toml new file mode 100644 index 00000000000..bed6b96bc26 --- /dev/null +++ b/tests/slow/DegradedMultiRegionStatus.toml @@ -0,0 +1,37 @@ +# Integration test for the "degraded multi-region" status signal. +# +# Simulated fearless topology (generateFearless=true) with two usable regions: +# - region 0 (primary): datacenter "0" (priority 2) + satellite "2" (priority 1) +# - region 1 (remote): datacenter "1" (priority 1) + satellite "3" (priority 1) +# - usable_regions=2, coordinators placed in non-primary DCs (minimumRegions=2) +# +# The workload kills ONLY datacenter "0" (the primary of the first region), leaving +# its satellite "2" and the remote region alive. The cluster fails over to region 1, +# recovers data from the surviving satellite "2", and recovery gets stuck at +# accepting_commits with the dead primary's log set not recruited +# (allLogs == false => RemoteRegionLogsMissing=1 => degraded_multi_region=true). + +[configuration] +generateFearless = true +minimumRegions = 2 +coordinators = 3 +datacenters = 4 +simHTTPServerEnabled = false +satelliteRedundancyMode = "one_satellite_single" +config = "single" + +[[test]] +testTitle = 'DegradedMultiRegionStatus' +clearAfterTest = false +runFailureWorkloads = false +runConsistencyCheck = true + + [[test.workload]] + testName = 'Cycle' + nodeCount = 100 + transactionsPerSecond = 100.0 + testDuration = 3.0 + + [[test.workload]] + testName = 'DegradedMultiRegionStatus' + testDuration = 200.0 From 013192d8dae2040930eecfb7d7adebd611bdbe9a Mon Sep 17 00:00:00 2001 From: Mark Shabanov Date: Thu, 27 Aug 2026 13:26:03 +0300 Subject: [PATCH 02/12] Apply clang-format to Status.cpp --- fdbserver/clustercontroller/Status.cpp | 13 +++++++++---- 1 file changed, 9 insertions(+), 4 deletions(-) diff --git a/fdbserver/clustercontroller/Status.cpp b/fdbserver/clustercontroller/Status.cpp index 99a1877a1c4..4fb4d10d086 100644 --- a/fdbserver/clustercontroller/Status.cpp +++ b/fdbserver/clustercontroller/Status.cpp @@ -3628,10 +3628,15 @@ StatusReply clusterGetFaultToleranceStatus(const std::string& statusStr) { json_spirit::mValue mv = readJSONStrictly(statusStr); JSONDoc jsonDoc(mv); - std::string faultToleranceRelatedFields[] = { - "fault_tolerance", "data", "logs", "maintenance_zone", "maintenance_seconds_remaining", "qos", - "recovery_state", "messages", "degraded_multi_region" - }; + std::string faultToleranceRelatedFields[] = { "fault_tolerance", + "data", + "logs", + "maintenance_zone", + "maintenance_seconds_remaining", + "qos", + "recovery_state", + "messages", + "degraded_multi_region" }; JsonBuilderObject statusObj; for (std::string& field : faultToleranceRelatedFields) { From 2a9819200cd1a24eeca76bed43ee7086eeb97d41 Mon Sep 17 00:00:00 2001 From: Mark Shabanov Date: Fri, 28 Aug 2026 17:00:20 +0300 Subject: [PATCH 03/12] Refresh remote-region stall event in a separate actor to avoid periodic cstate writes --- .../clustercontroller/ClusterRecovery.cpp | 111 ++++++++++++------ tests/slow/DegradedMultiRegionStatus.toml | 1 + 2 files changed, 76 insertions(+), 36 deletions(-) diff --git a/fdbserver/clustercontroller/ClusterRecovery.cpp b/fdbserver/clustercontroller/ClusterRecovery.cpp index 086af3f045f..c825a1eb72e 100644 --- a/fdbserver/clustercontroller/ClusterRecovery.cpp +++ b/fdbserver/clustercontroller/ClusterRecovery.cpp @@ -506,10 +506,58 @@ Future rejoinRequestHandler(Reference self) { } } +// Snapshot of the remote-region stall state shared between trackTlogRecovery (which updates it on +// every core-state change) and remoteRegionStallEventRefresher (which re-emits the trackLatest event +// on a fixed cadence while the stall persists). +struct RemoteRegionStallState : ReferenceCounted { + bool remoteRegionLogsMissing = false; + bool allLogs = false; + int oldTLogDataSize = 0; + int usableRegions = 1; + double missingSince = 0; // now() at stall start; valid when remoteRegionLogsMissing + std::string eventName; // stall event name, resolved once by trackTlogRecovery +}; + +// Emits the remote-region stall trackLatest event from the given snapshot. Called on core-state +// changes (where the signal can also be cleared) and on the refresher's cadence. +static void traceRemoteRegionStall(Reference self, + const RemoteRegionStallState* state, + bool remoteRegionLogsMissing, + double stallSeconds) { + TraceEvent(state->eventName.c_str(), self->dbgid) + .detail("RemoteRegionLogsMissing", remoteRegionLogsMissing) + .detail("AllLogs", state->allLogs) + .detail("OldTLogDataSize", state->oldTLogDataSize) + .detail("UsableRegions", state->usableRegions) + .detail("StallSeconds", stallSeconds) + .trackLatest(self->clusterRecoveryRemoteRegionStallEventHolder->trackingKey); +} + +// Re-emits the remote-region stall trackLatest event on a fixed cadence while the stall persists, +// keeping StallSeconds advancing on an otherwise idle cluster. Started lazily by trackTlogRecovery +// on the first observed stall; returns once the shared state says the stall is over. +Future remoteRegionStallEventRefresher(Reference self, + Reference state) { + while (state->remoteRegionLogsMissing) { + co_await delay(SERVER_KNOBS->DEGRADED_MULTI_REGION_REFRESH_SECONDS); + // trackTlogRecovery may have cleared the stall while the delay was pending. + if (!state->remoteRegionLogsMissing) { + co_return; + } + traceRemoteRegionStall(self, state.getPtr(), true, std::max(0.0, now() - state->missingSince)); + } +} + // Keeps the coordinated state (cstate) updated as the set of recruited tlogs change through recovery. Future trackTlogRecovery(Reference self, Reference>> oldLogSystems, Future minRecoveryDuration) { + // Shared with the refresher, which runs as a local future while a stall is active: when this + // function returns (final update), the refresher is cancelled automatically. + Reference stallState = makeReference(); + stallState->eventName = + getRecoveryEventName(ClusterRecoveryEventType::CLUSTER_RECOVERY_REMOTE_REGION_STALL_EVENT_NAME); + Future refresher; Future rejoinRequests = Never(); DBRecoveryCount recoverCount = self->cstate.myDBState.recoveryCount + 1; DatabaseConfiguration configuration = @@ -585,29 +633,36 @@ Future trackTlogRecovery(Reference self, .trackLatest(self->clusterRecoveryStateEventHolder->trackingKey); } - // Surface "degraded multi-region" as a first-class signal: with usableRegions > 1, the remote region's - // log set has not been recruited (allLogs == false). Committed data is still safe in the surviving - // region, so this is distinct from data loss. Note that oldTLogData is deliberately not a - // discriminator: when the remote region is down, old log generations cannot be purged (finalUpdate - // requires allLogs), so oldTLogData stays non-empty precisely in this stalled state. The value is also - // (re)set when all logs are recruited or the cluster fully recovers, so the trackLatest event always - // reflects the current state. + // "Degraded multi-region": with usableRegions > 1 the remote region's log set has not been + // recruited (allLogs == false). oldTLogData is deliberately not a discriminator: when the + // remote region is down, old generations cannot be purged (finalUpdate requires allLogs), + // so oldTLogData stays non-empty precisely in this stalled state. // - // The remote log set is legitimately absent before recovery reaches accepting_commits even during a - // normal (non-region-loss) recovery, because remote tlog recruitment is part of the normal recovery - // path. Only begin tracking the stall once the cluster is accepting commits: at that point a missing - // remote log set is no longer a transient recruiting artifact but a genuine degraded multi-region - // condition. Without this gate, StallSeconds accumulates across the pre-accepting_commits recruiting - // window and a slow-but-healthy recovery trips degraded_multi_region once the stall exceeds - // DEGRADED_MULTI_REGION_MIN_STALL_SECONDS. + // Gate on ACCEPTING_COMMITS: a missing remote log set is a transient recruiting artifact + // during normal recovery, so StallSeconds must not start accumulating before the cluster + // accepts commits, or a slow-but-healthy recovery trips degraded_multi_region. bool remoteRegionLogsMissing = configuration.usableRegions > 1 && !allLogs && self->recoveryState >= RecoveryState::ACCEPTING_COMMITS; + if (remoteRegionLogsMissing && !remoteLogsMissingSince.present()) { + remoteLogsMissingSince = now(); + } else if (!remoteRegionLogsMissing) { + remoteLogsMissingSince = Optional(); + } + // Publish the current snapshot for the refresher actor before starting it: the refresher reads + // remoteRegionLogsMissing synchronously on creation (before its first co_await), so publishing + // first lets it enter its loop immediately instead of deferring the first re-emission by a full + // loop iteration. + stallState->remoteRegionLogsMissing = remoteRegionLogsMissing; + stallState->allLogs = allLogs; + stallState->oldTLogDataSize = newState.oldTLogData.size(); + stallState->usableRegions = configuration.usableRegions; + if (remoteLogsMissingSince.present()) { + stallState->missingSince = remoteLogsMissingSince.get(); + } if (remoteRegionLogsMissing) { - if (!remoteLogsMissingSince.present()) { - remoteLogsMissingSince = now(); + if (!refresher.isValid() || refresher.isReady()) { + refresher = remoteRegionStallEventRefresher(self, stallState); } - } else { - remoteLogsMissingSince = Optional(); } // StallSeconds is carried in the event itself rather than derived from the event's // emission Time by status: the event is re-emitted on every core-state change, and @@ -615,15 +670,7 @@ Future trackTlogRecovery(Reference self, // the stall is still ongoing. double remoteRegionStallSeconds = remoteLogsMissingSince.present() ? std::max(0.0, now() - remoteLogsMissingSince.get()) : 0.0; - TraceEvent( - getRecoveryEventName(ClusterRecoveryEventType::CLUSTER_RECOVERY_REMOTE_REGION_STALL_EVENT_NAME).c_str(), - self->dbgid) - .detail("RemoteRegionLogsMissing", remoteRegionLogsMissing) - .detail("AllLogs", allLogs) - .detail("OldTLogDataSize", newState.oldTLogData.size()) - .detail("UsableRegions", configuration.usableRegions) - .detail("StallSeconds", remoteRegionStallSeconds) - .trackLatest(self->clusterRecoveryRemoteRegionStallEventHolder->trackingKey); + traceRemoteRegionStall(self, stallState.getPtr(), remoteRegionLogsMissing, remoteRegionStallSeconds); self->registrationTrigger.trigger(); @@ -633,15 +680,7 @@ Future trackTlogRecovery(Reference self, co_return; } - // While the remote log set is missing, keep re-emitting the stall event on a fixed cadence so - // StallSeconds keeps advancing even when no core-state change would otherwise wake this loop. - // Without this, an idle (commits-stalled) cluster would block forever in co_await changed and - // the sustained-stall threshold could never be crossed. - if (remoteRegionLogsMissing) { - co_await (changed || delay(SERVER_KNOBS->DEGRADED_MULTI_REGION_REFRESH_SECONDS)); - } else { - co_await changed; - } + co_await changed; } } diff --git a/tests/slow/DegradedMultiRegionStatus.toml b/tests/slow/DegradedMultiRegionStatus.toml index bed6b96bc26..f28d3aaf484 100644 --- a/tests/slow/DegradedMultiRegionStatus.toml +++ b/tests/slow/DegradedMultiRegionStatus.toml @@ -19,6 +19,7 @@ datacenters = 4 simHTTPServerEnabled = false satelliteRedundancyMode = "one_satellite_single" config = "single" +desiredTLogCount = 2 [[test]] testTitle = 'DegradedMultiRegionStatus' From 682f4667cd9aa0fbd81794a4c8b5e722e87c8a84 Mon Sep 17 00:00:00 2001 From: Mark Shabanov Date: Sat, 29 Aug 2026 22:53:16 +0300 Subject: [PATCH 04/12] Fix NormalRecoveryFalsePositive test for some seeds --- tests/slow/DegradedMultiRegionStatus.toml | 2 ++ 1 file changed, 2 insertions(+) diff --git a/tests/slow/DegradedMultiRegionStatus.toml b/tests/slow/DegradedMultiRegionStatus.toml index f28d3aaf484..e8a283b9f76 100644 --- a/tests/slow/DegradedMultiRegionStatus.toml +++ b/tests/slow/DegradedMultiRegionStatus.toml @@ -20,6 +20,8 @@ simHTTPServerEnabled = false satelliteRedundancyMode = "one_satellite_single" config = "single" desiredTLogCount = 2 +buggify = false +extraMachineCountDC = 2 [[test]] testTitle = 'DegradedMultiRegionStatus' From 74a72095ee41dc964974d3b6b9d57caa00704288 Mon Sep 17 00:00:00 2001 From: Mark Shabanov Date: Tue, 1 Sep 2026 00:31:57 +0300 Subject: [PATCH 05/12] fixed the included header --- fdbserver/workloads/DegradedMultiRegionStatus.cpp | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/fdbserver/workloads/DegradedMultiRegionStatus.cpp b/fdbserver/workloads/DegradedMultiRegionStatus.cpp index 2e0ebf75f83..82543d7261e 100644 --- a/fdbserver/workloads/DegradedMultiRegionStatus.cpp +++ b/fdbserver/workloads/DegradedMultiRegionStatus.cpp @@ -42,7 +42,7 @@ // // Finally it resets usable_regions=1 so subsequent checks are not stuck. -#include "fdbclient/NativeAPI.actor.h" +#include "fdbclient/NativeAPI.h" #include "fdbserver/core/TesterInterface.h" #include "fdbserver/core/WorkerInterface.h" #include "fdbserver/tester/workloads.h" From 06c45f9c1ec6bbb507f482410c261197ac2e9165 Mon Sep 17 00:00:00 2001 From: Mark Shabanov Date: Wed, 2 Sep 2026 15:31:49 +0300 Subject: [PATCH 06/12] Fix NormalRecoveryFalsePositive by using double redundancy --- tests/slow/DegradedMultiRegionStatus.toml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tests/slow/DegradedMultiRegionStatus.toml b/tests/slow/DegradedMultiRegionStatus.toml index e8a283b9f76..bb15b6148b7 100644 --- a/tests/slow/DegradedMultiRegionStatus.toml +++ b/tests/slow/DegradedMultiRegionStatus.toml @@ -18,7 +18,7 @@ coordinators = 3 datacenters = 4 simHTTPServerEnabled = false satelliteRedundancyMode = "one_satellite_single" -config = "single" +config = "double" desiredTLogCount = 2 buggify = false extraMachineCountDC = 2 From ba74433c46d2efd1513e95fc4259174ac04a34bf Mon Sep 17 00:00:00 2001 From: Mark Shabanov Date: Fri, 4 Sep 2026 15:07:41 +0300 Subject: [PATCH 07/12] Raise simulated DEGRADED_MULTI_REGION_MIN_STALL_SECONDS to avoid test false positive --- fdbserver/core/ServerKnobs.cpp | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/fdbserver/core/ServerKnobs.cpp b/fdbserver/core/ServerKnobs.cpp index b6e089fe0be..d22cc31e6b4 100644 --- a/fdbserver/core/ServerKnobs.cpp +++ b/fdbserver/core/ServerKnobs.cpp @@ -876,7 +876,7 @@ void ServerKnobs::initialize(Randomize randomize, ClientKnobs* clientKnobs, IsSi bool shortRecoveryDuration = randomize && buggify(); init( ENFORCED_MIN_RECOVERY_DURATION, 0.085 ); if( shortRecoveryDuration ) ENFORCED_MIN_RECOVERY_DURATION = 0.01; init( REQUIRED_MIN_RECOVERY_DURATION, 0.080 ); if( shortRecoveryDuration ) REQUIRED_MIN_RECOVERY_DURATION = 0.01; - init( DEGRADED_MULTI_REGION_MIN_STALL_SECONDS, 30.0 ); if( isSimulated ) DEGRADED_MULTI_REGION_MIN_STALL_SECONDS = 15.0; + init( DEGRADED_MULTI_REGION_MIN_STALL_SECONDS, 120.0 ); if( isSimulated ) DEGRADED_MULTI_REGION_MIN_STALL_SECONDS = 60.0; init( DEGRADED_MULTI_REGION_REFRESH_SECONDS, 1.0 ); init( ALWAYS_CAUSAL_READ_RISKY, false ); init( MAX_COMMIT_UPDATES, 2000 ); if( randomize && buggify() ) MAX_COMMIT_UPDATES = 1; From 59e5e4915f11d01273f5c0c5ef86ad72d6aa6d60 Mon Sep 17 00:00:00 2001 From: Mark Shabanov Date: Thu, 10 Sep 2026 17:10:46 +0300 Subject: [PATCH 08/12] Simplify degraded-multi-region stall signaling by deriving duration in status --- fdbcli/StatusCommand.cpp | 2 +- .../clustercontroller/ClusterRecovery.cpp | 92 ++++--------------- fdbserver/clustercontroller/Status.cpp | 32 ++++--- fdbserver/core/ServerKnobs.cpp | 1 - fdbserver/core/include/fdbserver/core/Knobs.h | 3 - .../workloads/DegradedMultiRegionStatus.cpp | 2 +- 6 files changed, 37 insertions(+), 95 deletions(-) diff --git a/fdbcli/StatusCommand.cpp b/fdbcli/StatusCommand.cpp index 1cff7ae0213..b907900cff2 100644 --- a/fdbcli/StatusCommand.cpp +++ b/fdbcli/StatusCommand.cpp @@ -724,7 +724,7 @@ void printStatus(StatusObjectReader statusObj, const std::string header = !degradedMultiRegion ? "\n\n Warning: the database may have data loss and availability loss. " - : "\n\n Warning: one region is unavailable; committed data remains safe in " + : "\n\n Warning: one region is unavailable; committed data is expected to remain safe in " "the surviving region. "; outputString += header + baseMessage; diff --git a/fdbserver/clustercontroller/ClusterRecovery.cpp b/fdbserver/clustercontroller/ClusterRecovery.cpp index c825a1eb72e..e0abcfa84a8 100644 --- a/fdbserver/clustercontroller/ClusterRecovery.cpp +++ b/fdbserver/clustercontroller/ClusterRecovery.cpp @@ -506,65 +506,17 @@ Future rejoinRequestHandler(Reference self) { } } -// Snapshot of the remote-region stall state shared between trackTlogRecovery (which updates it on -// every core-state change) and remoteRegionStallEventRefresher (which re-emits the trackLatest event -// on a fixed cadence while the stall persists). -struct RemoteRegionStallState : ReferenceCounted { - bool remoteRegionLogsMissing = false; - bool allLogs = false; - int oldTLogDataSize = 0; - int usableRegions = 1; - double missingSince = 0; // now() at stall start; valid when remoteRegionLogsMissing - std::string eventName; // stall event name, resolved once by trackTlogRecovery -}; - -// Emits the remote-region stall trackLatest event from the given snapshot. Called on core-state -// changes (where the signal can also be cleared) and on the refresher's cadence. -static void traceRemoteRegionStall(Reference self, - const RemoteRegionStallState* state, - bool remoteRegionLogsMissing, - double stallSeconds) { - TraceEvent(state->eventName.c_str(), self->dbgid) - .detail("RemoteRegionLogsMissing", remoteRegionLogsMissing) - .detail("AllLogs", state->allLogs) - .detail("OldTLogDataSize", state->oldTLogDataSize) - .detail("UsableRegions", state->usableRegions) - .detail("StallSeconds", stallSeconds) - .trackLatest(self->clusterRecoveryRemoteRegionStallEventHolder->trackingKey); -} - -// Re-emits the remote-region stall trackLatest event on a fixed cadence while the stall persists, -// keeping StallSeconds advancing on an otherwise idle cluster. Started lazily by trackTlogRecovery -// on the first observed stall; returns once the shared state says the stall is over. -Future remoteRegionStallEventRefresher(Reference self, - Reference state) { - while (state->remoteRegionLogsMissing) { - co_await delay(SERVER_KNOBS->DEGRADED_MULTI_REGION_REFRESH_SECONDS); - // trackTlogRecovery may have cleared the stall while the delay was pending. - if (!state->remoteRegionLogsMissing) { - co_return; - } - traceRemoteRegionStall(self, state.getPtr(), true, std::max(0.0, now() - state->missingSince)); - } -} - // Keeps the coordinated state (cstate) updated as the set of recruited tlogs change through recovery. Future trackTlogRecovery(Reference self, Reference>> oldLogSystems, Future minRecoveryDuration) { - // Shared with the refresher, which runs as a local future while a stall is active: when this - // function returns (final update), the refresher is cancelled automatically. - Reference stallState = makeReference(); - stallState->eventName = - getRecoveryEventName(ClusterRecoveryEventType::CLUSTER_RECOVERY_REMOTE_REGION_STALL_EVENT_NAME); - Future refresher; Future rejoinRequests = Never(); DBRecoveryCount recoverCount = self->cstate.myDBState.recoveryCount + 1; DatabaseConfiguration configuration = self->configuration; // self-configuration can be changed by configurationMonitor so we need a copy // Start of the current remote-region stall, if any. Kept across loop iterations so that - // re-emissions of the event on unrelated core-state changes do not reset the reported - // stall duration; it grows monotonically until the log set is complete. + // re-emissions of the event on unrelated core-state changes do not reset the reported stall + // start; status derives the elapsed duration from this timestamp. Optional remoteLogsMissingSince; while (true) { DBCoreState newState; @@ -639,7 +591,7 @@ Future trackTlogRecovery(Reference self, // so oldTLogData stays non-empty precisely in this stalled state. // // Gate on ACCEPTING_COMMITS: a missing remote log set is a transient recruiting artifact - // during normal recovery, so StallSeconds must not start accumulating before the cluster + // during normal recovery, so the stall must not start accumulating before the cluster // accepts commits, or a slow-but-healthy recovery trips degraded_multi_region. bool remoteRegionLogsMissing = configuration.usableRegions > 1 && !allLogs && self->recoveryState >= RecoveryState::ACCEPTING_COMMITS; @@ -648,29 +600,21 @@ Future trackTlogRecovery(Reference self, } else if (!remoteRegionLogsMissing) { remoteLogsMissingSince = Optional(); } - // Publish the current snapshot for the refresher actor before starting it: the refresher reads - // remoteRegionLogsMissing synchronously on creation (before its first co_await), so publishing - // first lets it enter its loop immediately instead of deferring the first re-emission by a full - // loop iteration. - stallState->remoteRegionLogsMissing = remoteRegionLogsMissing; - stallState->allLogs = allLogs; - stallState->oldTLogDataSize = newState.oldTLogData.size(); - stallState->usableRegions = configuration.usableRegions; - if (remoteLogsMissingSince.present()) { - stallState->missingSince = remoteLogsMissingSince.get(); - } - if (remoteRegionLogsMissing) { - if (!refresher.isValid() || refresher.isReady()) { - refresher = remoteRegionStallEventRefresher(self, stallState); - } - } - // StallSeconds is carried in the event itself rather than derived from the event's - // emission Time by status: the event is re-emitted on every core-state change, and - // deriving the duration from the latest emission would reset the stall counter while - // the stall is still ongoing. - double remoteRegionStallSeconds = - remoteLogsMissingSince.present() ? std::max(0.0, now() - remoteLogsMissingSince.get()) : 0.0; - traceRemoteRegionStall(self, stallState.getPtr(), remoteRegionLogsMissing, remoteRegionStallSeconds); + // Carry the absolute time the stall began rather than a pre-computed duration: status derives + // how long the stall has lasted when it reads this event, so the event does not have to be + // re-emitted periodically just to keep a duration field fresh. Format as a string because + // numeric trace fields are emitted with "%g" (six significant digits), which is far too coarse + // for an absolute timestamp. + double remoteRegionStallStartSeconds = remoteLogsMissingSince.present() ? remoteLogsMissingSince.get() : 0.0; + TraceEvent( + getRecoveryEventName(ClusterRecoveryEventType::CLUSTER_RECOVERY_REMOTE_REGION_STALL_EVENT_NAME).c_str(), + self->dbgid) + .detail("RemoteRegionLogsMissing", remoteRegionLogsMissing) + .detail("AllLogs", allLogs) + .detail("OldTLogDataSize", newState.oldTLogData.size()) + .detail("UsableRegions", configuration.usableRegions) + .detail("RemoteRegionStallStartSeconds", format("%.6f", remoteRegionStallStartSeconds)) + .trackLatest(self->clusterRecoveryRemoteRegionStallEventHolder->trackingKey); self->registrationTrigger.trigger(); diff --git a/fdbserver/clustercontroller/Status.cpp b/fdbserver/clustercontroller/Status.cpp index 4fb4d10d086..dda61293cff 100644 --- a/fdbserver/clustercontroller/Status.cpp +++ b/fdbserver/clustercontroller/Status.cpp @@ -1213,11 +1213,12 @@ static bool parseRemoteRegionLogsMissing(const TraceEventFields& remoteRegionSta return remoteRegionStallEvent.getInt("RemoteRegionLogsMissing") != 0; } -// Parse the duration (in seconds) for which the remote-region-stall signal has been active. The recovery -// side carries the duration in the event itself ("StallSeconds") because the event is (re)emitted on every -// core-state change; deriving the duration from the latest emission time would reset the stall counter -// while the stall is still ongoing. An absent event, an event reporting that the logs are not missing, or -// any malformed value conservatively reports 0, which keeps the cluster from being flagged as degraded. +// Parse how long (in seconds) the remote-region-stall signal has been active. The recovery side carries +// the absolute time the stall began ("RemoteRegionStallStartSeconds") rather than a pre-computed duration: +// the event is (re)emitted on core-state changes, so a duration field would go stale between emissions, +// whereas the start timestamp lets status compute the elapsed time from this event on every request. An +// absent event, an event reporting that the logs are not missing, or any malformed value conservatively +// reports 0, which keeps the cluster from being flagged as degraded. static double parseRemoteRegionStallSeconds(const TraceEventFields& remoteRegionStallEvent) { try { if (remoteRegionStallEvent.size() == 0) { @@ -1226,8 +1227,8 @@ static double parseRemoteRegionStallSeconds(const TraceEventFields& remoteRegion if (remoteRegionStallEvent.getInt("RemoteRegionLogsMissing") == 0) { return 0.0; } - double stallSeconds = std::stod(remoteRegionStallEvent.getValue("StallSeconds")); - return stallSeconds > 0.0 ? stallSeconds : 0.0; + double stallStartSeconds = remoteRegionStallEvent.getDouble("RemoteRegionStallStartSeconds"); + return stallStartSeconds > 0.0 ? std::max(0.0, now() - stallStartSeconds) : 0.0; } catch (Error& e) { if (e.code() == error_code_actor_cancelled) { throw; @@ -4073,10 +4074,10 @@ TEST_CASE("/fdbserver/clustercontroller/degradedMultiRegionComputation") { ASSERT(parseRemoteRegionLogsMissing(remoteRegionStallEvent)); } - // parseRemoteRegionStallSeconds: logs not missing -> no active stall, regardless of the carried duration. + // parseRemoteRegionStallSeconds: logs not missing -> no active stall, regardless of the carried start. { TraceEventFields remoteRegionStallEvent; - remoteRegionStallEvent.addField("StallSeconds", "3600.0"); + remoteRegionStallEvent.addField("RemoteRegionStallStartSeconds", format("%.6f", now() - 3600.0)); remoteRegionStallEvent.addField("RemoteRegionLogsMissing", "0"); ASSERT_EQ(parseRemoteRegionStallSeconds(remoteRegionStallEvent), 0.0); } @@ -4085,31 +4086,32 @@ TEST_CASE("/fdbserver/clustercontroller/degradedMultiRegionComputation") { TraceEventFields remoteRegionStallEvent; ASSERT_EQ(parseRemoteRegionStallSeconds(remoteRegionStallEvent), 0.0); } - // parseRemoteRegionStallSeconds: sustained stall carries the monotonic duration reported by recovery. + // parseRemoteRegionStallSeconds: sustained stall -> elapsed time since the reported start exceeds the + // threshold, so status computes a duration that keeps growing without the event being re-emitted. { const double threshold = SERVER_KNOBS->DEGRADED_MULTI_REGION_MIN_STALL_SECONDS; TraceEventFields remoteRegionStallEvent; - remoteRegionStallEvent.addField("StallSeconds", format("%.6f", threshold + 5.0)); + remoteRegionStallEvent.addField("RemoteRegionStallStartSeconds", format("%.6f", now() - (threshold + 5.0))); remoteRegionStallEvent.addField("RemoteRegionLogsMissing", "1"); ASSERT_GE(parseRemoteRegionStallSeconds(remoteRegionStallEvent), threshold); } { const double threshold = SERVER_KNOBS->DEGRADED_MULTI_REGION_MIN_STALL_SECONDS; TraceEventFields remoteRegionStallEvent; - remoteRegionStallEvent.addField("StallSeconds", "0.1"); + remoteRegionStallEvent.addField("RemoteRegionStallStartSeconds", format("%.6f", now() - 0.1)); remoteRegionStallEvent.addField("RemoteRegionLogsMissing", "1"); ASSERT_LT(parseRemoteRegionStallSeconds(remoteRegionStallEvent), threshold); } - // parseRemoteRegionStallSeconds: malformed or negative StallSeconds fails safe to 0 (not degraded). + // parseRemoteRegionStallSeconds: malformed or non-positive start fails safe to 0 (not degraded). { TraceEventFields remoteRegionStallEvent; - remoteRegionStallEvent.addField("StallSeconds", "garbage"); + remoteRegionStallEvent.addField("RemoteRegionStallStartSeconds", "garbage"); remoteRegionStallEvent.addField("RemoteRegionLogsMissing", "1"); ASSERT_EQ(parseRemoteRegionStallSeconds(remoteRegionStallEvent), 0.0); } { TraceEventFields remoteRegionStallEvent; - remoteRegionStallEvent.addField("StallSeconds", "-5.0"); + remoteRegionStallEvent.addField("RemoteRegionStallStartSeconds", "-5.0"); remoteRegionStallEvent.addField("RemoteRegionLogsMissing", "1"); ASSERT_EQ(parseRemoteRegionStallSeconds(remoteRegionStallEvent), 0.0); } diff --git a/fdbserver/core/ServerKnobs.cpp b/fdbserver/core/ServerKnobs.cpp index d22cc31e6b4..b33edfc8add 100644 --- a/fdbserver/core/ServerKnobs.cpp +++ b/fdbserver/core/ServerKnobs.cpp @@ -877,7 +877,6 @@ void ServerKnobs::initialize(Randomize randomize, ClientKnobs* clientKnobs, IsSi init( ENFORCED_MIN_RECOVERY_DURATION, 0.085 ); if( shortRecoveryDuration ) ENFORCED_MIN_RECOVERY_DURATION = 0.01; init( REQUIRED_MIN_RECOVERY_DURATION, 0.080 ); if( shortRecoveryDuration ) REQUIRED_MIN_RECOVERY_DURATION = 0.01; init( DEGRADED_MULTI_REGION_MIN_STALL_SECONDS, 120.0 ); if( isSimulated ) DEGRADED_MULTI_REGION_MIN_STALL_SECONDS = 60.0; - init( DEGRADED_MULTI_REGION_REFRESH_SECONDS, 1.0 ); init( ALWAYS_CAUSAL_READ_RISKY, false ); init( MAX_COMMIT_UPDATES, 2000 ); if( randomize && buggify() ) MAX_COMMIT_UPDATES = 1; init( MAX_PROXY_COMPUTE, 2.0 ); diff --git a/fdbserver/core/include/fdbserver/core/Knobs.h b/fdbserver/core/include/fdbserver/core/Knobs.h index b45b262d9b0..41944c0e556 100644 --- a/fdbserver/core/include/fdbserver/core/Knobs.h +++ b/fdbserver/core/include/fdbserver/core/Knobs.h @@ -747,9 +747,6 @@ class SWIFT_CXX_IMMORTAL_SINGLETON_TYPE ServerKnobs : public KnobsImpl Date: Thu, 10 Sep 2026 17:51:49 +0300 Subject: [PATCH 09/12] Add joshua test (cherry picked from src3 commit 8ff49e1dea) --- .github/workflows/build.yml | 98 +++++++++++++++++++ .github/workflows/run-tests.yml | 95 ++++++++++++++++++ .github/workflows/tests-only.yml | 19 ++++ CMakeLists.txt | 17 +++- .../for-github/calc-version-from-git.bash | 24 +++++ build-scripts/for-linux/build-on-linux.bash | 83 ++++++++++++++++ build-scripts/for-linux/find-swift.bash | 14 +++ build-scripts/for-linux/test-joshua.bash | 41 ++++++++ build-scripts/set-ver-prms.sh | 12 +++ flow/CMakeLists.txt | 2 +- 10 files changed, 402 insertions(+), 3 deletions(-) create mode 100644 .github/workflows/build.yml create mode 100644 .github/workflows/run-tests.yml create mode 100644 .github/workflows/tests-only.yml create mode 100755 build-scripts/for-github/calc-version-from-git.bash create mode 100755 build-scripts/for-linux/build-on-linux.bash create mode 100755 build-scripts/for-linux/find-swift.bash create mode 100755 build-scripts/for-linux/test-joshua.bash create mode 100755 build-scripts/set-ver-prms.sh diff --git a/.github/workflows/build.yml b/.github/workflows/build.yml new file mode 100644 index 00000000000..204ca03a89b --- /dev/null +++ b/.github/workflows/build.yml @@ -0,0 +1,98 @@ +--- +name: Build + +on: + push: + branches: ["ow-fork-*"] + tags: ["*-*ow"] + +jobs: + calc_ver: + # calculate versions from git tags + runs-on: ubuntu-latest + outputs: + project_ver: ${{steps.vers.outputs.project_ver}} + build_ver: ${{steps.vers.outputs.build_ver}} + full_ver: ${{steps.vers.outputs.full_ver}} + release_flag: ${{steps.vers.outputs.release_flag}} + + steps: + - name: Checkout + uses: actions/checkout@v4 + + - name: Calculate versions + id: vers + shell: bash + run: ${{github.workspace}}/build-scripts/for-github/calc-version-from-git.bash + + build: + needs: [calc_ver] + env: + # There are not enough free 19 Gb mounted in "/" to build and create rpm and deb packages. + # So, we use /mnt/bld which is mounted to a bigger space. In /mnt there is 66 Gb free space in github action server. + BLD_DIR: /mnt/bld + + + strategy: + matrix: + include: + - run_on: ubuntu-latest + for: linux + prepare: "debian-based" + build_on: linux + parallel: 5 + image: foundationdb-build:8.0.0-0.ow.build + owner: ${GITHUB_REPOSITORY_OWNER@L} + runs-on: ${{ matrix.run_on }} + steps: + - name: Checkout + uses: actions/checkout@v4 + + - name: Set building repo + run: | + echo "use_image=ghcr.io/${{matrix.owner}}/${{matrix.image}}" >> "$GITHUB_ENV" + + - name: Build + run: | + sudo mkdir -p ${{ env.BLD_DIR }} + sudo chmod 777 ${{ env.BLD_DIR }} + + echo "=== Disk space before build ===" + df -h + + podman run --rm \ + --name build \ + --mount=type=tmpfs,dst=/tmp \ + --mount=type=tmpfs,dst=/var/tmp \ + --security-opt label=disable \ + --mount=type=bind,src=${{github.workspace}},dst=/home/runner/src,readonly \ + --mount=type=bind,src=${{ env.BLD_DIR }},dst=/home/runner/bld \ + $use_image \ + /home/runner/src/build-scripts/for-${{ matrix.for }}/build-on-${{ matrix.build_on }}.bash \ + ${{needs.calc_ver.outputs.project_ver}} \ + ${{needs.calc_ver.outputs.build_ver}} \ + ${{needs.calc_ver.outputs.release_flag}} \ + ${{ matrix.parallel }} + + echo "=== Disk space after building and creating packages ===" + df -h + + # - name: Minimal tests + # working-directory: ${{ env.BLD_DIR }} + # shell: bash + # run: ctest --output-on-failure -V + + - name: Upload result + uses: nanoufo/action-upload-artifacts-and-release-assets@v2 + with: + path: | + ${{ env.BLD_DIR }}/linux/packages/*${{needs.calc_ver.outputs.full_ver}}* + if-no-files-found: error + + tests: + needs: [calc_ver, build] + uses: ./.github/workflows/run-tests.yml + with: + full_ver: ${{needs.calc_ver.outputs.full_ver}} + build_run_id: ${{ github.run_id }} + secrets: inherit diff --git a/.github/workflows/run-tests.yml b/.github/workflows/run-tests.yml new file mode 100644 index 00000000000..e0cdac8e81d --- /dev/null +++ b/.github/workflows/run-tests.yml @@ -0,0 +1,95 @@ +name: Run Tests (reusable) + +on: + workflow_call: + inputs: + full_ver: + description: 'Version of the build to test (e.g. 7.4.0-3.1.ow)' + required: true + type: string + build_run_id: + description: 'Run ID to download the correctness package from' + required: true + type: string + +jobs: + tests: + runs-on: ubuntu-latest + env: + JOSHUA_DB_VER: "7.1.57" + N_OF_TESTS: 500 # to fit in 360 minutes job run limit + JOSHUA_AGENT_TAG: "rockylinux9.6-20260309" + # parameter that controls the maximum lifetime of the Joshua agent (in seconds). + AGENT_TIMEOUT: 18000 + + steps: + - name: Set agent URL + run: | + echo "JOSHUA_AGENT_URL=ghcr.io/${GITHUB_REPOSITORY_OWNER,,}" >> $GITHUB_ENV + echo "Agent URL: ghcr.io/${GITHUB_REPOSITORY_OWNER,,}" + + - name: Checkout + uses: actions/checkout@v4 + with: + path: ${{github.workspace}}/src + + - name: Install dependencies + shell: bash + run: | + sudo apt-get update + sudo apt-get install -y sudo wget crudini git python3 python3-pip + sudo pip3 install wheel setuptools python-dateutil lxml boto3 + + - name: Install FoundationDb + shell: bash + run: | + mkdir deb + pushd deb + MY_ARCH=`dpkg-architecture -q DEB_BUILD_ARCH` + wget https://github.com/apple/foundationdb/releases/download/${{ env.JOSHUA_DB_VER }}/foundationdb-clients_${{ env.JOSHUA_DB_VER }}-1_${MY_ARCH}.deb https://github.com/apple/foundationdb/releases/download/${{ env.JOSHUA_DB_VER }}/foundationdb-server_${{ env.JOSHUA_DB_VER }}-1_${MY_ARCH}.deb + sudo apt-get install -y ./foundationdb-clients_${{ env.JOSHUA_DB_VER }}-1_${MY_ARCH}.deb ./foundationdb-server_${{ env.JOSHUA_DB_VER }}-1_${MY_ARCH}.deb + popd + sudo systemctl stop foundationdb + MY_IP=`hostname -I | awk '{print $1}'` + sudo sed -i s/127.0.0.1/$MY_IP/ /etc/foundationdb/fdb.cluster + sudo crudini --set /etc/foundationdb/foundationdb.conf fdbserver memory 4GiB + sudo systemctl start foundationdb + pip3 install 'foundationdb==${{ env.JOSHUA_DB_VER }}' + + - name: Download the correctness package + uses: actions/download-artifact@v4 + id: download_correctness + with: + name: correctness-${{ inputs.full_ver }}.tar.gz + run-id: ${{ inputs.build_run_id }} + github-token: ${{ secrets.GITHUB_TOKEN }} + + - name: Echo download path + run: echo ${{steps.download_correctness.outputs.download-path}} + + - name: Display structure of downloaded files + run: ls -R + working-directory: ${{github.workspace}} + + - name: Download joshua + shell: bash + run: | + git clone https://github.com/FoundationDB/fdb-joshua.git + + - name: run joshua-agent + shell: bash + run: | + podman pull ${{ env.JOSHUA_AGENT_URL }}/joshua-agent:${{ env.JOSHUA_AGENT_TAG }} + for i in 1 2 3 4; do + podman run -d \ + -v /etc/foundationdb:/etc/foundationdb \ + -e AGENT_TIMEOUT=${{ env.AGENT_TIMEOUT }} \ + joshua-agent:${{ env.JOSHUA_AGENT_TAG }} + done + + - name: run tests + shell: bash + working-directory: ${{github.workspace}}/fdb-joshua + run: | + podman ps + ${{github.workspace}}/src/build-scripts/for-linux/test-joshua.bash ${{github.workspace}}/correctness-${{ inputs.full_ver }}.tar.gz ${{env.N_OF_TESTS}} \ No newline at end of file diff --git a/.github/workflows/tests-only.yml b/.github/workflows/tests-only.yml new file mode 100644 index 00000000000..7feb20dec8a --- /dev/null +++ b/.github/workflows/tests-only.yml @@ -0,0 +1,19 @@ +name: Tests only + +on: + workflow_dispatch: + inputs: + full_ver: + description: 'Version of the build to test (e.g. 7.4.0-3.1.ow)' + required: true + build_run_id: + description: 'Run ID of the build workflow (find it in the URL of the build run)' + required: true + +jobs: + tests: + uses: ./.github/workflows/run-tests.yml + with: + full_ver: ${{ github.event.inputs.full_ver }} + build_run_id: ${{ github.event.inputs.build_run_id }} + secrets: inherit \ No newline at end of file diff --git a/CMakeLists.txt b/CMakeLists.txt index 184e28a9c4c..468b0660264 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -111,13 +111,26 @@ message(STATUS "Current git version ${CURRENT_GIT_VERSION}") ############################################################################### # Version information ############################################################################### +if(DEFINED VERSION AND NOT VERSION STREQUAL CMAKE_PROJECT_VERSION) + message(FATAL_ERROR "The version parameter ${VERSION} differs from the project version ${CMAKE_PROJECT_VERSION}") +endif() if(FDB_RELEASE_CANDIDATE) set(FDB_RELEASE_CANDIDATE_VERSION 1 CACHE STRING "release candidate version") set(FDB_VERSION ${PROJECT_VERSION}-rc${FDB_RELEASE_CANDIDATE_VERSION}) else() - set(FDB_VERSION ${PROJECT_VERSION}) + if(NOT DEFINED BUILD_VERSION) + if(NOT FDB_RELEASE) + if(CURRENT_GIT_VERSION) + set(git_string ".${CURRENT_GIT_VERSION}") + endif() + set(BUILD_VERSION "0${git_string}.PRERELEASE") + else() + set(BUILD_VERSION "1") + endif() + endif() + set(FDB_VERSION "${PROJECT_VERSION}-${BUILD_VERSION}") endif() if(NOT FDB_RELEASE) string(TIMESTAMP FDB_BUILD_TIMESTMAP %Y%m%d%H%M%S) @@ -129,7 +142,7 @@ if(NOT FDB_RELEASE) set(FDB_BUILDTIME_STRING ".${FDB_BUILDTIME}") set(PRERELEASE_TAG "prerelease") endif() -set(FDB_VERSION_PLAIN ${FDB_VERSION}) +set(FDB_VERSION_PLAIN ${PROJECT_VERSION}) string(REPLACE "." ";" FDB_VERSION_LIST ${FDB_VERSION_PLAIN}) list(GET FDB_VERSION_LIST 0 FDB_MAJOR) list(GET FDB_VERSION_LIST 1 FDB_MINOR) diff --git a/build-scripts/for-github/calc-version-from-git.bash b/build-scripts/for-github/calc-version-from-git.bash new file mode 100755 index 00000000000..4d2f525be9a --- /dev/null +++ b/build-scripts/for-github/calc-version-from-git.bash @@ -0,0 +1,24 @@ +#!/bin/bash + +# Calculate version numbers from git tags and set them to github output +# Assume that current directory is a github repository + +# Do not set github variables when running outside github +[[ -z "$GITHUB_OUTPUT" ]] && GITHUB_OUTPUT=/dev/null + +git fetch --prune --unshallow --tags --force +GIT_VERSION=`git describe --tags` +PROJECT_VERSION=`echo $GIT_VERSION | cut -d- -f1` +BUILD_VERSION=`echo $GIT_VERSION | cut -d- -f2-3 --output-delimiter=.` +GIT_CHANGE_NUM=`echo $GIT_VERSION | cut -d- -f3` +if [[ -n "$GIT_CHANGE_NUM" ]] || [[ "$BUILD_VERSION" < "1" ]]; then + RELEASE_FLAG=OFF +else + RELEASE_FLAG=ON +fi + +# Display versions and set Github environment +echo "project_ver=$PROJECT_VERSION" | tee -a $GITHUB_OUTPUT +echo "build_ver=$BUILD_VERSION" | tee -a $GITHUB_OUTPUT +echo "full_ver=$PROJECT_VERSION-$BUILD_VERSION" | tee -a $GITHUB_OUTPUT +echo "release_flag=$RELEASE_FLAG" | tee -a $GITHUB_OUTPUT diff --git a/build-scripts/for-linux/build-on-linux.bash b/build-scripts/for-linux/build-on-linux.bash new file mode 100755 index 00000000000..cc454523650 --- /dev/null +++ b/build-scripts/for-linux/build-on-linux.bash @@ -0,0 +1,83 @@ +#!/bin/bash + +# $1 - Version +# $2 - Build version +# $3 - Release flag +# $4 - Paralllel threads +# $5 - Source Dir. If not set then relative to the script dir + +get_oldest_java_path() +{ + if [ -x "$(command -v update-java-alternatives)" ] + then + update-java-alternatives -l | sort -Vk1 | head -n 1 | awk '{print $3}' + else + echo 'update-java-alternatives is not installed' >&2 + ls -d /etc/alternatives/java_sdk_*[0-9] | sort -V | head -n 1 + fi +} + +set -e + +BASE_DIR="$(readlink -f $(dirname $0))" +source $BASE_DIR/../set-ver-prms.sh "$1" "$2" +RELEASE_FLAG=${3:-OFF} +PARALLEL_PRMS="-j ${4:-$(nproc)}" +SRC_DIR=${5:-$(readlink -f $BASE_DIR/../..)} + +START_DIR=`pwd` +BUILD_DIR=$START_DIR/bld/linux + +mkdir -p $BUILD_DIR +pushd $BUILD_DIR + +rm -rf * +export LANG=C + +APP_PRMS="\ + $CMAKE_VERSION_PRMS \ + -DFDB_RELEASE=$RELEASE_FLAG \ + -DGENERATE_DEBUG_PACKAGES=OFF \ + -DUSE_LIBCXX=ON \ + -DUSE_LD=LLD \ + -DSTATIC_LINK_LIBCXX=ON \ + -DENABLE_SIMULATION_TESTS=ON" + +[ ! -e /usr/lib64/libcrypto.a -a -e /opt/openssl/lib/libcrypto.a ] && \ + APP_PRMS="$APP_PRMS -DOPENSSL_ROOT_DIR=/opt/openssl" + +# find swift +APP_PRMS="$APP_PRMS -DCMAKE_Swift_COMPILER=`$BASE_DIR/find-swift.bash`" + +# set oldest java +JAVA_OLDEST_PATH=$(get_oldest_java_path) + + +echo "env JAVA_HOME=$JAVA_OLDEST_PATH CC=clang CXX=clang++ cmake -G Ninja $APP_PRMS $SRC_DIR" +env JAVA_HOME=$JAVA_OLDEST_PATH CC=clang CXX=clang++ cmake -G Ninja $APP_PRMS $SRC_DIR + +echo "number of parallel jobs: [$PARALLEL_PRMS]" + +ninja $PARALLEL_PRMS -k 0 + +echo "=== Disk space after build ===" +df -h + +cpack -G RPM +for fn in packages/*.rpm; do echo "$fn:"; rpm -qpRv $fn; echo; done + +cpack -G DEB +for fn in packages/*.deb; do echo "$fn:"; dpkg -I $fn; done + +ninja $PARALLEL_PRMS -k 0 documentation/package_html + +# make all-binaries and libraries archives +# calculate filenames +# TO DO: calculate the processor architecture instead of hardcoding x86_64 +BINS_FILENAME=`ls -1 packages/foundationdb-docs-* | sed s/-docs/-bins/ | sed s/.tgz/.x86_64.tgz/` +LIBS_FILENAME=`ls -1 packages/foundationdb-docs-* | sed s/-docs/-libs/ | sed s/.tgz/.x86_64.tgz/` +tar -I pigz -cvf $BINS_FILENAME -C packages/bin . +tar -I pigz -cvf $LIBS_FILENAME -C packages/lib . + +popd + diff --git a/build-scripts/for-linux/find-swift.bash b/build-scripts/for-linux/find-swift.bash new file mode 100755 index 00000000000..bbbc50febfc --- /dev/null +++ b/build-scripts/for-linux/find-swift.bash @@ -0,0 +1,14 @@ +#!/bin/bash + +# find swift installation +# print the path to swift installed or return non zero rc + + +if ! which swiftc 2>/dev/null; then + WILDCARDS="/opt/swift*/usr/bin/swiftc" + if ls $WILDCARDS 2>/dev/null; then + ls -1d $WILDCARDS | tail -n 1 + else + return $? 2>/dev/null || exit $? + fi +fi diff --git a/build-scripts/for-linux/test-joshua.bash b/build-scripts/for-linux/test-joshua.bash new file mode 100755 index 00000000000..f7f6ff249d5 --- /dev/null +++ b/build-scripts/for-linux/test-joshua.bash @@ -0,0 +1,41 @@ +#!/bin/bash + +# $1 - the path to the correctness archive file, ex. correctness-7.3.49-2.ow.tar.gz +# $2 - number of tests + +python3 -m joshua.joshua start --tarball $1 --max-runs $2 + +python3 -m joshua.joshua tail | python3 -c " +import sys, re +sys.stdout = open(sys.stdout.fileno(), 'w', buffering=1) +from datetime import datetime +failed = False +count = 1 +for line in sys.stdin: + m = re.search(r'TestFile=\"([^\"]+)\".*?Ok=\"(\d+)\"', line) + if m: + timestamp = datetime.now().strftime('%Y-%m-%d %H:%M:%S') + print(f'[#{count} {timestamp}] TestFile=\"{m.group(1)}\" Ok=\"{m.group(2)}\"') + if m.group(2) == '0': + failed = True + else: + sys.stdout.write(line) + count += 1 +sys.exit(1 if failed else 0) +" & + +TAIL_PID=$! + +while kill -0 "$TAIL_PID" 2>/dev/null; do + sleep 30 + if ! podman ps -q 2>/dev/null | grep -q .; then + echo "ERROR: all joshua-agent containers have stopped unexpectedly" >&2 + kill "$TAIL_PID" 2>/dev/null + echo "Stopping joshua due to container failure..." + python3 -m joshua.joshua stop + exit 1 + fi +done + +wait "$TAIL_PID" +exit $? diff --git a/build-scripts/set-ver-prms.sh b/build-scripts/set-ver-prms.sh new file mode 100755 index 00000000000..656ab5b01c1 --- /dev/null +++ b/build-scripts/set-ver-prms.sh @@ -0,0 +1,12 @@ +#!/bin/bash + +# $1 - Project version +# $2 - Build version + +CMAKE_VERSION_PRMS= + +[[ -n "$1" ]] && CMAKE_VERSION_PRMS="$CMAKE_VERSION_PRMS -DVERSION=$1" +[[ -n "$2" ]] && CMAKE_VERSION_PRMS="$CMAKE_VERSION_PRMS -DBUILD_VERSION=$2" + +export CMAKE_VERSION_PRMS + diff --git a/flow/CMakeLists.txt b/flow/CMakeLists.txt index 63acc647bd1..0bff7d14a96 100644 --- a/flow/CMakeLists.txt +++ b/flow/CMakeLists.txt @@ -98,7 +98,7 @@ target_link_libraries(flowlinktest PRIVATE "$" add_flow_target(EXECUTABLE NAME flow_test SRCS FlowTest.cpp) register_fdb_unit_tests(flow_test) -target_link_libraries(flow_test PRIVATE "$" fdbrpc) +target_link_libraries(flow_test PRIVATE fdbrpc) set(IS_ARM_MAC NO) if(APPLE AND CMAKE_SYSTEM_PROCESSOR STREQUAL "arm64") From 8f1df8fa0bf92e665b64c79743875896c2e45f9b Mon Sep 17 00:00:00 2001 From: Mark Shabanov Date: Fri, 11 Sep 2026 10:22:37 +0300 Subject: [PATCH 10/12] clang-format fix --- fdbcli/StatusCommand.cpp | 3 ++- fdbserver/workloads/DegradedMultiRegionStatus.cpp | 6 +++--- 2 files changed, 5 insertions(+), 4 deletions(-) diff --git a/fdbcli/StatusCommand.cpp b/fdbcli/StatusCommand.cpp index b907900cff2..e1448c69a42 100644 --- a/fdbcli/StatusCommand.cpp +++ b/fdbcli/StatusCommand.cpp @@ -724,7 +724,8 @@ void printStatus(StatusObjectReader statusObj, const std::string header = !degradedMultiRegion ? "\n\n Warning: the database may have data loss and availability loss. " - : "\n\n Warning: one region is unavailable; committed data is expected to remain safe in " + : "\n\n Warning: one region is unavailable; committed data is expected to " + "remain safe in " "the surviving region. "; outputString += header + baseMessage; diff --git a/fdbserver/workloads/DegradedMultiRegionStatus.cpp b/fdbserver/workloads/DegradedMultiRegionStatus.cpp index 4c1d339548c..c6d7c973221 100644 --- a/fdbserver/workloads/DegradedMultiRegionStatus.cpp +++ b/fdbserver/workloads/DegradedMultiRegionStatus.cpp @@ -279,9 +279,9 @@ struct DegradedMultiRegionStatusWorkload : TestWorkload { .detail("Degraded", degraded) .detail("DataStateDesc", dataStateDesc); printf("\n=== Degraded Multi-Region Status Found ===\n"); - printf( - "Warning: one region is unavailable; committed data is expected to remain safe in the surviving " - "region.\n"); + printf("Warning: one region is unavailable; committed data is expected to remain safe " + "in the surviving " + "region.\n"); printf( "Please restart following tlog interfaces, otherwise storage servers may never be " "able to catch up.\n"); From cd9d3bf49d886d7cf816732a63e130c1f6b17fcb Mon Sep 17 00:00:00 2001 From: Mark Shabanov Date: Tue, 15 Sep 2026 14:12:33 +0300 Subject: [PATCH 11/12] Fix warning message in degraded_multi_region ternary --- fdbcli/StatusCommand.cpp | 15 +++++++-------- 1 file changed, 7 insertions(+), 8 deletions(-) diff --git a/fdbcli/StatusCommand.cpp b/fdbcli/StatusCommand.cpp index e1448c69a42..db15a9336c3 100644 --- a/fdbcli/StatusCommand.cpp +++ b/fdbcli/StatusCommand.cpp @@ -721,14 +721,13 @@ void printStatus(StatusObjectReader statusObj, bool degradedMultiRegion = false; statusObjCluster.get("degraded_multi_region", degradedMultiRegion); - const std::string header = - !degradedMultiRegion - ? "\n\n Warning: the database may have data loss and availability loss. " - : "\n\n Warning: one region is unavailable; committed data is expected to " - "remain safe in " - "the surviving region. "; - - outputString += header + baseMessage; + outputString += + "\n\n Warning: the database may have data loss and availability loss. "; + if (degradedMultiRegion) { + outputString += "One region is unavailable; committed data is expected to remain " + "safe in the surviving region. "; + } + outputString += baseMessage; } else { outputString += format( "\n\n Warning: the database may have availability loss. The current log state " From b113c179447208836cc93c79fbf3ab0fca88e889 Mon Sep 17 00:00:00 2001 From: Mark Shabanov Date: Wed, 30 Sep 2026 16:06:58 +0300 Subject: [PATCH 12/12] Do not warn about data loss when a lost log set keeps its satellite copy --- .../source/mr-status-json-schemas.rst.inc | 1 + fdbcli/StatusCommand.cpp | 65 ++++++++------ fdbcli/tests/fdbcli_tests.py | 88 +++++++++++++++++++ fdbserver/clustercontroller/Status.cpp | 44 +++++++++- .../workloads/DegradedMultiRegionStatus.cpp | 39 +++++--- 5 files changed, 197 insertions(+), 40 deletions(-) diff --git a/documentation/sphinx/source/mr-status-json-schemas.rst.inc b/documentation/sphinx/source/mr-status-json-schemas.rst.inc index a06bf5411bb..d040eae584d 100644 --- a/documentation/sphinx/source/mr-status-json-schemas.rst.inc +++ b/documentation/sphinx/source/mr-status-json-schemas.rst.inc @@ -650,6 +650,7 @@ }, "active_tss_count":0, "degraded_processes":0, + "degraded_multi_region":true, "database_available":true, "database_lock_state":{ "locked":true, diff --git a/fdbcli/StatusCommand.cpp b/fdbcli/StatusCommand.cpp index db15a9336c3..18e781bb77b 100644 --- a/fdbcli/StatusCommand.cpp +++ b/fdbcli/StatusCommand.cpp @@ -710,24 +710,19 @@ void printStatus(StatusObjectReader statusObj, outputString += format(" (%d without data loss)", dataLoss); } + // Both verdicts are evaluated before the branching below so that region availability can + // be reported independently of how the log state is classified. + const bool possiblyLosingData = logEpochsMayBeLosingData(statusObjCluster); + bool degradedMultiRegion = false; + statusObjCluster.get("degraded_multi_region", degradedMultiRegion); + if (dataLoss == -1) { ASSERT_WE_THINK(availLoss == -1); - const bool possiblyLosingData = logEpochsMayBeLosingData(statusObjCluster); if (possiblyLosingData) { - const std::string baseMessage = - "Please restart following tlog interfaces, otherwise storage servers " - "may never be able to catch up.\n"; - - bool degradedMultiRegion = false; - statusObjCluster.get("degraded_multi_region", degradedMultiRegion); - - outputString += - "\n\n Warning: the database may have data loss and availability loss. "; - if (degradedMultiRegion) { - outputString += "One region is unavailable; committed data is expected to remain " - "safe in the surviving region. "; - } - outputString += baseMessage; + outputString += format( + "\n\n Warning: the database may have data loss and availability loss. Please " + "restart following tlog interfaces, otherwise storage servers may never be able " + "to catch up.\n"); } else { outputString += format( "\n\n Warning: the database may have availability loss. The current log state " @@ -735,18 +730,14 @@ void printStatus(StatusObjectReader statusObj, } if (statusObjCluster.has("logs")) { for (StatusObjectReader logEpoch : statusObjCluster.last().get_array()) { - bool logEpochPossiblyLosingData; - if (logEpoch.get("possibly_losing_data", logEpochPossiblyLosingData) && - !logEpochPossiblyLosingData) { - continue; - } - // Current epoch doesn't have an end version. - int64_t epoch, beginVersion, endVersion = invalidVersion; - bool current; - logEpoch.get("epoch", epoch); - logEpoch.get("begin_version", beginVersion); - logEpoch.get("end_version", endVersion); - logEpoch.get("current", current); + // Unknown means "assume at risk": the field is absent on servers that predate it. + bool logEpochPossiblyLosingData = true; + const bool dataAtRisk = + !logEpoch.get("possibly_losing_data", logEpochPossiblyLosingData) || + logEpochPossiblyLosingData; + // Unavailable log interfaces are the ones that must come back for the cluster to + // finish recovering, so an epoch is reported when either its data is at risk or any + // of its log interfaces is unavailable. std::string missing_log_interfaces; if (logEpoch.has("log_interfaces")) { for (StatusObjectReader logInterface : logEpoch.last().get_array()) { @@ -759,6 +750,16 @@ void printStatus(StatusObjectReader statusObj, } } } + if (!dataAtRisk && missing_log_interfaces.empty()) { + continue; + } + // Current epoch doesn't have an end version. + int64_t epoch, beginVersion, endVersion = invalidVersion; + bool current; + logEpoch.get("epoch", epoch); + logEpoch.get("begin_version", beginVersion); + logEpoch.get("end_version", endVersion); + logEpoch.get("current", current); outputString += format( " %s log epoch: %lld begin: %lld end: %s, missing " "log interfaces(id,address): %s\n", @@ -770,6 +771,16 @@ void printStatus(StatusObjectReader statusObj, } } } + // Region availability is orthogonal to the data-loss verdict above: report it whenever the + // cluster says a region is unavailable, whichever way the log state was classified. The + // statement about committed data is only made when the log state does not indicate data loss. + if (degradedMultiRegion) { + outputString += possiblyLosingData + ? "\n One region is unavailable; data in the surviving region may be " + "incomplete.\n" + : "\n One region is unavailable; committed data is expected to remain " + "safe in the surviving region.\n"; + } } } diff --git a/fdbcli/tests/fdbcli_tests.py b/fdbcli/tests/fdbcli_tests.py index 68a7becf757..461fc97efff 100755 --- a/fdbcli/tests/fdbcli_tests.py +++ b/fdbcli/tests/fdbcli_tests.py @@ -458,6 +458,93 @@ def status_json_file_region_failover_message(): assert "may have data loss" not in stdout +def status_json_degraded_multi_region_message(): + # A multi-region cluster whose primary region is down while its satellite survives recovers without data + # loss: the log generation that lost its primary set still has an intact synchronous satellite copy, so + # the server reports possibly_losing_data=false for it. Status must warn about availability loss, explain + # that the surviving region is expected to hold all committed data, and must not claim data loss. + status_json = { + "client": { + "cluster_file": {"path": "fdb.cluster", "up_to_date": True}, + "coordinators": {"coordinators": [], "quorum_reachable": True}, + "database_status": {"available": True, "healthy": False}, + "messages": [], + "timestamp": 1417807090, + }, + "cluster": { + "configuration": { + "redundancy_mode": "double", + "storage_engine": "ssd-2", + "coordinators_count": 3, + "excluded_servers": [], + }, + "data": {"state": {"name": "healthy", "healthy": True}}, + "degraded_multi_region": True, + "fault_tolerance": { + "max_zone_failures_without_losing_availability": -1, + "max_zone_failures_without_losing_data": -1, + }, + "logs": [ + { + "epoch": 2, + "current": True, + "begin_version": 100, + "possibly_losing_data": False, + "log_interfaces": [], + }, + { + "epoch": 1, + "current": False, + "begin_version": 1, + "end_version": 100, + "possibly_losing_data": False, + "log_fault_tolerance": -1, + "satellite_log_replication_factor": 1, + "satellite_log_fault_tolerance": 0, + "log_interfaces": [ + { + "id": "aaaaaaaaaaaaaaaa", + "healthy": False, + "address": "1.1.1.1:4500", + } + ], + }, + ], + "machines": {}, + "processes": {}, + }, + } + + def render(status): + with tempfile.NamedTemporaryFile(mode="w", suffix=".json") as status_file: + json.dump(status, status_file) + status_file.flush() + result = subprocess.run( + [command_template[0], "--status-from-json", status_file.name], + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + env=fdbcli_env, + ) + assert result.returncode == 0, result.stderr.decode("utf-8") + return result.stdout.decode("utf-8") + + stdout = render(status_json) + assert "Warning: the database may have availability loss." in stdout + assert "may have data loss" not in stdout + assert ( + "One region is unavailable; committed data is expected to remain safe in the surviving region." + in stdout + ) + # The unavailable interfaces of the generation whose data is safe are still listed, so that the operator + # knows what has to come back for the cluster to finish recovering. + assert "missing log interfaces(id,address): aaaaaaaaaaaaaaaa,1.1.1.1:4500" in stdout + + # The region note is only printed when the cluster explicitly reports an unavailable region. + del status_json["cluster"]["degraded_multi_region"] + stdout = render(status_json) + assert "One region is unavailable" not in stdout + + @enable_logging() def consistencycheck(logger): consistency_check_on_output = "ConsistencyCheck is on" @@ -1053,6 +1140,7 @@ def tls_address_suffix(): integer_options() tls_address_suffix() status_json_file_region_failover_message() + status_json_degraded_multi_region_message() idempotency_ids() cdc_operator_commands() client_threads_per_version_env_ignored() diff --git a/fdbserver/clustercontroller/Status.cpp b/fdbserver/clustercontroller/Status.cpp index dda61293cff..1aaf12befa0 100644 --- a/fdbserver/clustercontroller/Status.cpp +++ b/fdbserver/clustercontroller/Status.cpp @@ -2485,6 +2485,23 @@ static AsyncResult clusterSummaryStatisticsFetcher( co_return statusObj; } +// A satellite log set carries a synchronous copy of its region's mutation stream, and a commit is acknowledged +// only after the satellite quorum has made it durable. A log generation whose local primary set has been lost +// is therefore still fully recoverable (its storage servers can still catch up) as long as its local satellite +// set still satisfies its replication policy, even though the fault-tolerance counters of the lost primary set +// are negative. This is deliberately conservative in the other direction: a generation with no satellite set, +// or with a satellite set that no longer satisfies its policy, stays marked as possibly losing data. +static bool computePossiblyLosingData(int minFaultTolerance, + Optional logFaultTolerance, + Optional satLogFaultTolerance) { + if (minFaultTolerance >= 0) { + return false; + } + const bool primarySetLost = logFaultTolerance.present() && logFaultTolerance.get() < 0; + const bool satelliteCopyIntact = satLogFaultTolerance.present() && satLogFaultTolerance.get() >= 0; + return !(primarySetLost && satelliteCopyIntact); +} + static JsonBuilderObject tlogFetcher(int* logFaultTolerance, const std::vector& tLogs, std::unordered_map const& address_workers) { @@ -2559,8 +2576,10 @@ static JsonBuilderObject tlogFetcher(int* logFaultTolerance, *logFaultTolerance = std::min(*logFaultTolerance, minFaultTolerance); statusObj["log_interfaces"] = logsObj; // We may lose logs in this log generation, storage servers may never be able to catch up this log - // generation. - statusObj["possibly_losing_data"] = minFaultTolerance < 0; + // generation, unless its synchronous satellite copy is still intact: in that case the committed mutations + // remain recoverable from the satellite. + statusObj["possibly_losing_data"] = + computePossiblyLosingData(minFaultTolerance, log_fault_tolerance, sat_log_fault_tolerance); if (sat_log_replication_factor.present()) statusObj["satellite_log_replication_factor"] = sat_log_replication_factor.get(); @@ -4116,6 +4135,27 @@ TEST_CASE("/fdbserver/clustercontroller/degradedMultiRegionComputation") { ASSERT_EQ(parseRemoteRegionStallSeconds(remoteRegionStallEvent), 0.0); } + // computePossiblyLosingData: the primary region's local log set is lost while its satellite set still + // satisfies its replication policy, so the committed mutations are recoverable from that synchronous copy. + // This is the satellite-survives-a-region-failure case, which must not be reported as data loss. + { + ASSERT(!computePossiblyLosingData(/*minFaultTolerance=*/-1, + /*logFaultTolerance=*/-1, + /*satLogFaultTolerance=*/0)); + } + // computePossiblyLosingData: the satellite set is degraded, missing, or there is no primary set to lose, + // so the negative fault tolerance does indicate that committed data may be lost. + { + ASSERT(computePossiblyLosingData(-1, -1, -1)); + ASSERT(computePossiblyLosingData(-1, -1, Optional())); + ASSERT(computePossiblyLosingData(-1, Optional(), 0)); + } + // computePossiblyLosingData: no local set is below its policy (including an epoch with no log sets). + { + ASSERT(!computePossiblyLosingData(0, 0, 0)); + ASSERT(!computePossiblyLosingData(0, Optional(), Optional())); + } + co_return; } diff --git a/fdbserver/workloads/DegradedMultiRegionStatus.cpp b/fdbserver/workloads/DegradedMultiRegionStatus.cpp index c6d7c973221..ee12653a824 100644 --- a/fdbserver/workloads/DegradedMultiRegionStatus.cpp +++ b/fdbserver/workloads/DegradedMultiRegionStatus.cpp @@ -33,6 +33,10 @@ // the status JSON and asserts: // - cluster.degraded_multi_region == true // - cluster.data.state.description contains "Degraded multiregional" +// - fault_tolerance.max_zone_failures_without_losing_data is still negative, while the +// generation that lost its primary log set (negative log_fault_tolerance) is NOT +// flagged as possibly losing data: its surviving satellite set still satisfies its +// replication policy, so committed mutations remain recoverable from it // // Before killing the primary DC the test also exercises the false-positive regression // from the remote log set being transiently absent at accepting_commits: it verifies @@ -278,17 +282,30 @@ struct DegradedMultiRegionStatusWorkload : TestWorkload { .detail("Elapsed", now() - tStart) .detail("Degraded", degraded) .detail("DataStateDesc", dataStateDesc); - printf("\n=== Degraded Multi-Region Status Found ===\n"); - printf("Warning: one region is unavailable; committed data is expected to remain safe " - "in the surviving " - "region.\n"); - printf( - "Please restart following tlog interfaces, otherwise storage servers may never be " - "able to catch up.\n"); - printf("\nData:\n"); - printf(" Replication health - %s\n", dataStateDesc.c_str()); - printf("========================================\n\n"); - fflush(stdout); + // The unavailable region forces the overall fault tolerance negative, so the log state + // is what decides whether committed data is reported as at risk. + ASSERT(clusterObj.contains("fault_tolerance") && + clusterObj["fault_tolerance"] + .get_obj()["max_zone_failures_without_losing_data"] + .get_int() == -1); + // A generation that lost its primary log set while keeping a satellite set that still + // satisfies its replication policy must not be flagged: the satellite holds a + // synchronous copy of the mutation stream, so its storage servers can still catch up. + ASSERT(clusterObj.contains("logs")); + bool sawGenerationWithLostPrimarySet = false; + for (auto& logEpoch : clusterObj["logs"].get_array()) { + auto& logEpochObj = logEpoch.get_obj(); + if (!logEpochObj.contains("log_fault_tolerance") || + logEpochObj["log_fault_tolerance"].get_int() >= 0 || + !logEpochObj.contains("satellite_log_fault_tolerance") || + logEpochObj["satellite_log_fault_tolerance"].get_int() < 0) { + continue; + } + sawGenerationWithLostPrimarySet = true; + ASSERT(logEpochObj.contains("possibly_losing_data") && + !logEpochObj["possibly_losing_data"].get_bool()); + } + ASSERT(sawGenerationWithLostPrimarySet); co_return true; } } else {