From 915d37cc31c2588802695e81dfc84e3514d60cbc Mon Sep 17 00:00:00 2001 From: Nicholas Dudfield Date: Wed, 23 Sep 2026 09:17:16 +0700 Subject: [PATCH 1/3] fix(consensus): hold the consensus mutex around accept-time extension prep doAccept reacquires the consensus mutex for the live-build view, the ordering salt, and the replay or pre-build block. Ledger building, open-ledger locks, and endConsensus stay outside it. --- src/xrpld/app/consensus/RCLConsensus.cpp | 52 +++++++++++++++--------- src/xrpld/app/consensus/RCLConsensus.h | 18 ++++---- 2 files changed, 43 insertions(+), 27 deletions(-) diff --git a/src/xrpld/app/consensus/RCLConsensus.cpp b/src/xrpld/app/consensus/RCLConsensus.cpp index 92377b6292..b456d238d9 100644 --- a/src/xrpld/app/consensus/RCLConsensus.cpp +++ b/src/xrpld/app/consensus/RCLConsensus.cpp @@ -77,6 +77,7 @@ RCLConsensus::RCLConsensus( beast::Journal journal) : adaptor_( app, + mutex_, std::move(feeVote), ledgerMaster, localTxs, @@ -92,6 +93,7 @@ RCLConsensus::~RCLConsensus() = default; RCLConsensus::Adaptor::Adaptor( Application& app, + std::recursive_mutex& consensusMutex, std::unique_ptr&& feeVote, LedgerMaster& ledgerMaster, LocalTxs& localTxs, @@ -99,6 +101,7 @@ RCLConsensus::Adaptor::Adaptor( ValidatorKeys const& validatorKeys, beast::Journal journal) : app_(app) + , consensusMutex_(consensusMutex) , feeVote_(std::move(feeVote)) , ledgerMaster_(ledgerMaster) , localTxs_(localTxs) @@ -518,11 +521,9 @@ RCLConsensus::Adaptor::onAccept( "acceptLedger", [=, this, cj = std::move(consensusJson)]() mutable { //@@start do-accept-freeze-contract - // Note that no lock is held or acquired during this job. - // This is because generic Consensus guarantees that once a ledger - // is accepted, the consensus results and capture by reference state - // will not change until startRound is called (which happens via - // endConsensus). + // The job locks only for extension preparation, not ledger + // building. The accepted result must remain frozen until + // startRound, reached through endConsensus after doAccept returns. //@@end do-accept-freeze-contract RclConsensusLogger clog("onAccept", validating, j_); this->doAccept( @@ -587,10 +588,12 @@ RCLConsensus::Adaptor::doAccept( // influence fallback entropy, transaction ordering, or ledger state. auto replayData = ledgerMaster_.releaseReplay(); auto const consensusTxSetHash = result.txns.id(); - auto const liveBuild = replayData - ? std::optional{} - : std::optional{ - ce().makeLiveBuildTxSet(result.txns)}; + auto const liveBuild = [&] { + std::lock_guard lock{consensusMutex_}; + return replayData ? std::optional{} + : std::optional{ + ce().makeLiveBuildTxSet(result.txns)}; + }(); auto const& buildTxs = liveBuild ? liveBuild->txns : result.txns; auto const buildTxSetHash = buildTxs.id(); @@ -615,7 +618,11 @@ RCLConsensus::Adaptor::doAccept( // FIXME: Use a std::vector and a custom sorter instead of CanonicalTXSet? //@@start txn-ordering-salt-build-inputs auto const buildSeq = prevLedger.seq() + 1; - CanonicalTXSet retriableTxs{ce().txnOrderingSalt(buildTxSetHash, buildSeq)}; + auto const orderingSalt = [&] { + std::lock_guard lock{consensusMutex_}; + return ce().txnOrderingSalt(buildTxSetHash, buildSeq); + }(); + CanonicalTXSet retriableTxs{orderingSalt}; JLOG(j_.debug()) << "Building canonical tx set: " << retriableTxs.key(); @@ -641,17 +648,22 @@ RCLConsensus::Adaptor::doAccept( // Export witness injection are independently gated inside onPreBuild; // export-only rounds still need this hook even when RNG is off. //@@start accept-time-cleanup-disabled - if (replayData) { - ce().onReplayBuild(); - } - else if (ce().rngEnabled() || ce().exportEnabled()) - { - ce().onPreBuild(retriableTxs, buildSeq, buildTxSetHash); - } - else - { - ce().clearRngState(); + // Match consensus-side readers/writers. Never extend this scope across + // buildLCL, the open-ledger locks, or endConsensus. + std::lock_guard lock{consensusMutex_}; + if (replayData) + { + ce().onReplayBuild(); + } + else if (ce().rngEnabled() || ce().exportEnabled()) + { + ce().onPreBuild(retriableTxs, buildSeq, buildTxSetHash); + } + else + { + ce().clearRngState(); + } } //@@end accept-time-cleanup-disabled diff --git a/src/xrpld/app/consensus/RCLConsensus.h b/src/xrpld/app/consensus/RCLConsensus.h index 0792ac23e0..76095dc756 100644 --- a/src/xrpld/app/consensus/RCLConsensus.h +++ b/src/xrpld/app/consensus/RCLConsensus.h @@ -69,6 +69,7 @@ class RCLConsensus class Adaptor { Application& app_; + std::recursive_mutex& consensusMutex_; std::unique_ptr feeVote_; LedgerMaster& ledgerMaster_; LocalTxs& localTxs_; @@ -114,6 +115,7 @@ class RCLConsensus Adaptor( Application& app, + std::recursive_mutex& consensusMutex, std::unique_ptr&& feeVote, LedgerMaster& ledgerMaster, LocalTxs& localTxs, @@ -210,10 +212,10 @@ class RCLConsensus // Consensus methods and since RCLConsensus::consensus_ should // only be accessed under lock, these will only be called under lock. // - // In general, the idea is that there is only ONE thread that is running - // consensus code at anytime. The only special case is the dispatched - // onAccept call, which does not take a lock and relies on Consensus not - // changing state until a future call to startRound. + // Normally only one thread runs consensus code at a time. The + // dispatched accept job builds the ledger outside the lock, but + // reacquires it for extension preparation. The accepted result must + // still remain unchanged until a future call to startRound. friend class Consensus; /** Attempt to acquire a specific ledger. @@ -552,9 +554,11 @@ public: } private: - // Since Consensus does not provide intrinsic thread-safety, this mutex - // guards all calls to consensus_. adaptor_ uses atomics internally - // to allow concurrent access of its data members that have getters. + // Guards mutable consensus state and accept-job round-state access. + // Atomic-only status polls and the extension's independently synchronized + // cross-thread APIs are exempt. Constructed before adaptor_. + // Lock order: C before LedgerMaster, never the reverse; C before busyMu_ + // before the collector. mutable std::recursive_mutex mutex_; Adaptor adaptor_; From 87e3a2cbd6daf78f2d4395b460528ab011a85f5d Mon Sep 17 00:00:00 2001 From: Nicholas Dudfield Date: Wed, 23 Sep 2026 09:50:48 +0700 Subject: [PATCH 2/3] test(seam): accept-lock probe and hold-time counter Hooks and the max-hold counter sit on top of the mutex fix so a suite can show the accept block waiting and report how long the lock is held. --- src/xrpld/app/consensus/RCLConsensus.cpp | 40 +++++++++++++++++-- src/xrpld/app/consensus/RCLConsensus.h | 49 ++++++++++++++++++++++++ 2 files changed, 85 insertions(+), 4 deletions(-) diff --git a/src/xrpld/app/consensus/RCLConsensus.cpp b/src/xrpld/app/consensus/RCLConsensus.cpp index b456d238d9..7ebccecaa5 100644 --- a/src/xrpld/app/consensus/RCLConsensus.cpp +++ b/src/xrpld/app/consensus/RCLConsensus.cpp @@ -91,6 +91,23 @@ RCLConsensus::RCLConsensus( RCLConsensus::~RCLConsensus() = default; +void +RCLConsensus::Adaptor::noteAcceptLock( + std::chrono::steady_clock::time_point start) +{ + auto const ns = std::chrono::duration_cast( + std::chrono::steady_clock::now() - start) + .count(); + if (ns <= 0) + return; + auto cur = maxAcceptLockNs_.load(std::memory_order_relaxed); + while (static_cast(ns) > cur && + !maxAcceptLockNs_.compare_exchange_weak( + cur, static_cast(ns), std::memory_order_relaxed)) + { + } +} + RCLConsensus::Adaptor::Adaptor( Application& app, std::recursive_mutex& consensusMutex, @@ -588,11 +605,17 @@ RCLConsensus::Adaptor::doAccept( // influence fallback entropy, transaction ordering, or ledger state. auto replayData = ledgerMaster_.releaseReplay(); auto const consensusTxSetHash = result.txns.id(); + if (beforeAcceptExtension_) + beforeAcceptExtension_(); auto const liveBuild = [&] { std::lock_guard lock{consensusMutex_}; - return replayData ? std::optional{} - : std::optional{ - ce().makeLiveBuildTxSet(result.txns)}; + auto const holdStart = std::chrono::steady_clock::now(); + auto built = replayData + ? std::optional{} + : std::optional{ + ce().makeLiveBuildTxSet(result.txns)}; + noteAcceptLock(holdStart); + return built; }(); auto const& buildTxs = liveBuild ? liveBuild->txns : result.txns; auto const buildTxSetHash = buildTxs.id(); @@ -620,7 +643,10 @@ RCLConsensus::Adaptor::doAccept( auto const buildSeq = prevLedger.seq() + 1; auto const orderingSalt = [&] { std::lock_guard lock{consensusMutex_}; - return ce().txnOrderingSalt(buildTxSetHash, buildSeq); + auto const holdStart = std::chrono::steady_clock::now(); + auto salt = ce().txnOrderingSalt(buildTxSetHash, buildSeq); + noteAcceptLock(holdStart); + return salt; }(); CanonicalTXSet retriableTxs{orderingSalt}; @@ -652,6 +678,9 @@ RCLConsensus::Adaptor::doAccept( // Match consensus-side readers/writers. Never extend this scope across // buildLCL, the open-ledger locks, or endConsensus. std::lock_guard lock{consensusMutex_}; + auto const holdStart = std::chrono::steady_clock::now(); + if (insideAcceptExtension_) + insideAcceptExtension_(); if (replayData) { ce().onReplayBuild(); @@ -664,6 +693,7 @@ RCLConsensus::Adaptor::doAccept( { ce().clearRngState(); } + noteAcceptLock(holdStart); } //@@end accept-time-cleanup-disabled @@ -1132,6 +1162,8 @@ RCLConsensus::timerEntry( { std::lock_guard _{mutex_}; consensus_->timerEntry(now, clog); + if (whileConsensusLocked_) + whileConsensusLocked_(); } catch (SHAMapMissingNode const& mn) { diff --git a/src/xrpld/app/consensus/RCLConsensus.h b/src/xrpld/app/consensus/RCLConsensus.h index 76095dc756..d49ecbaddb 100644 --- a/src/xrpld/app/consensus/RCLConsensus.h +++ b/src/xrpld/app/consensus/RCLConsensus.h @@ -40,6 +40,7 @@ #include #include #include +#include #include #include #include @@ -204,6 +205,23 @@ class RCLConsensus ConsensusExtensions const& ce() const; + // Test seam. Empty unless a suite is showing that the accept job's + // extension block waits while consensus holds this mutex. + void + setAcceptExtensionProbe( + std::function beforeLock, + std::function insideLock) + { + beforeAcceptExtension_ = std::move(beforeLock); + insideAcceptExtension_ = std::move(insideLock); + } + + std::uint64_t + maxAcceptLockHoldNs() const + { + return maxAcceptLockNs_.load(std::memory_order_relaxed); + } + private: //--------------------------------------------------------------------- // The following members implement the generic Consensus requirements @@ -218,6 +236,13 @@ class RCLConsensus // still remain unchanged until a future call to startRound. friend class Consensus; + std::function beforeAcceptExtension_; + std::function insideAcceptExtension_; + std::atomic maxAcceptLockNs_{0}; + + void + noteAcceptLock(std::chrono::steady_clock::time_point start); + /** Attempt to acquire a specific ledger. If not available, asynchronously acquires from the network. @@ -505,6 +530,28 @@ public: bool extensionsBusy() const; + // Test seams for the accept-path lock. Production leaves them empty. + void + setWhileConsensusLocked(std::function hook) + { + whileConsensusLocked_ = std::move(hook); + } + + void + setAcceptExtensionProbe( + std::function beforeLock, + std::function insideLock) + { + adaptor_.setAcceptExtensionProbe( + std::move(beforeLock), std::move(insideLock)); + } + + std::uint64_t + maxAcceptLockHoldNs() const + { + return adaptor_.maxAcceptLockHoldNs(); + } + //! @see Consensus::getJson Json::Value getJson(bool full) const; @@ -561,6 +608,8 @@ private: // before the collector. mutable std::recursive_mutex mutex_; + std::function whileConsensusLocked_; + Adaptor adaptor_; std::unique_ptr> consensus_; beast::Journal const j_; From edd5ee0f8b4cebff7788ee630946a4cb2759dd1f Mon Sep 17 00:00:00 2001 From: Nicholas Dudfield Date: Wed, 23 Sep 2026 09:17:24 +0700 Subject: [PATCH 3/3] test(consensus): show the accept extension block waits on the consensus lock A heartbeat-side hook holds the mutex while a captured accept job reaches it, and standalone simulate reenters the same lock. --- .../consensus/ThreadedExtensions_test.cpp | 186 +++++++++++++++--- 1 file changed, 161 insertions(+), 25 deletions(-) diff --git a/src/test/consensus/ThreadedExtensions_test.cpp b/src/test/consensus/ThreadedExtensions_test.cpp index fa974f8c7c..21f042389b 100644 --- a/src/test/consensus/ThreadedExtensions_test.cpp +++ b/src/test/consensus/ThreadedExtensions_test.cpp @@ -2,6 +2,7 @@ #include #include +#include #include #include #include @@ -24,8 +25,10 @@ #include #include #include +#include #include #include +#include #include #include #include @@ -88,7 +91,9 @@ class ThreadedExtensions_test : public beast::unit_test::suite { bool found = false; forEachItem( - ledger, keylet::pendingExports(), [&](std::shared_ptr const& sle) { + ledger, + keylet::pendingExports(), + [&](std::shared_ptr const& sle) { if (sle && sle->key() == key) found = true; }); @@ -224,7 +229,8 @@ class ThreadedExtensions_test : public beast::unit_test::suite std::vector keys; std::vector unl; - for (auto const* name : {"thread-val-0", "thread-val-1", "thread-val-2"}) + for (auto const* name : + {"thread-val-0", "thread-val-1", "thread-val-2"}) { keys.push_back(ValidatorKey::fromPassphrase(name)); unl.push_back(keys.back().pubKey); @@ -323,11 +329,12 @@ class ThreadedExtensions_test : public beast::unit_test::suite << " delay=" << counts.delayCalls.load() << std::endl; BEAST_EXPECT(false); }; - auto const submitOk = [&](std::size_t node, Json::Value tx, jtx::Account const& signer) { - tx[jss::NetworkID] = networkID; - auto const txn = net.submit(node, std::move(tx), signer); - return txn && txn->getResult() == tesSUCCESS; - }; + auto const submitOk = + [&](std::size_t node, Json::Value tx, jtx::Account const& signer) { + tx[jss::NetworkID] = networkID; + auto const txn = net.submit(node, std::move(tx), signer); + return txn && txn->getResult() == tesSUCCESS; + }; if (!submitOk( 0, jtx::pay(jtx::Account::master, owner, jtx::XRP(20'000)), @@ -387,8 +394,7 @@ class ThreadedExtensions_test : public beast::unit_test::suite parent) : nullptr; log << " threaded view diag fromReport=" - << (view && view->fromUNLReport ? 1 : 0) - << " masters=" + << (view && view->fromUNLReport ? 1 : 0) << " masters=" << (view ? view->orderedOriginalMasterKeys.size() : 0) << std::endl; return fail("validator view"); @@ -472,10 +478,11 @@ class ThreadedExtensions_test : public beast::unit_test::suite }; std::vector origins; auto const historyFrom = net.minValidated(); - auto const sendIntent = [&](std::uint32_t ticket) -> std::optional { + auto const sendIntent = + [&](std::uint32_t ticket) -> std::optional { auto const open = net[0].app().openLedger().current()->seq(); - Json::Value tx = intent( - ticket, open + ExportLimits::maxAdmissionWindowLedgers); + Json::Value tx = + intent(ticket, open + ExportLimits::maxAdmissionWindowLedgers); tx[jss::NetworkID] = networkID; auto const txn = net.submit(0, std::move(tx), owner); if (!txn || txn->getResult() != tesSUCCESS) @@ -544,8 +551,8 @@ class ThreadedExtensions_test : public beast::unit_test::suite PeerFaultConfig cfg; cfg.sendDropPctX100 = 10000; if (category >= 0) - cfg.messageCategories = std::set{ - static_cast(category)}; + cfg.messageCategories = + std::set{static_cast(category)}; net[2].app().getRuntimeConfig().setPeerDefaults(cfg); } { @@ -632,7 +639,8 @@ class ThreadedExtensions_test : public beast::unit_test::suite if (!beat()) return; if (validatorsAdvanced() < downSeq + 2) - return fail("validators did not advance while the observer was down"); + return fail( + "validators did not advance while the observer was down"); if (!net.restartNode(observer).isUp()) return fail("observer did not restart"); if (!relink(observer)) @@ -784,8 +792,9 @@ class ThreadedExtensions_test : public beast::unit_test::suite tx->getFieldH256(sfTransactionHash) == origin) return fail("witnessed more than once"); } - log << " threaded " << schedule << " origin=" << to_string(origin) - << " witnessed:" << seqW << std::endl; + log << " threaded " << schedule + << " origin=" << to_string(origin) << " witnessed:" << seqW + << std::endl; continue; } for (std::uint32_t n = 0; n < nNodes; ++n) @@ -805,10 +814,10 @@ class ThreadedExtensions_test : public beast::unit_test::suite auto const inDir = latch && pendingDirContains(*ledger, latch->key()); char const* const shape = !latch ? "missing" - : hasSig ? "signed" - : hasNode && inDir ? "pending" - : !hasNode && !inDir ? "pruned" - : "inconsistent"; + : hasSig ? "signed" + : hasNode && inDir ? "pending" + : !hasNode && !inDir ? "pruned" + : "inconsistent"; log << " threaded " << schedule << " origin=" << to_string(origin) << " expired shape=" << shape << std::endl; if (!latch) @@ -821,21 +830,148 @@ class ThreadedExtensions_test : public beast::unit_test::suite : "directory holds an unlinked latch"); } - log << " threaded " << schedule - << " origins=" << origins.size() + log << " threaded " << schedule << " origins=" << origins.size() << " directDropped=" << droppedCalls - droppedQueued - << " delayed=" << delayed - << " observerDownFrom=" << downSeq + << " delayed=" << delayed << " observerDownFrom=" << downSeq << " pre=" << preSeq << " end=" << agreed << std::endl; BEAST_EXPECT(origins.size() == 8); BEAST_EXPECT(droppedCalls > droppedQueued); BEAST_EXPECT(delayed > 0); + log << " accept lock hold ns=" + << net[0].app().getOPs().getConsensus().maxAcceptLockHoldNs() + << std::endl; + } + + void + testSimulateReenters() + { + testcase("simulate force-accept reenters the consensus lock"); + using namespace jtx; + Env env{*this, envconfig()}; + auto& consensus = env.app().getOPs().getConsensus(); + consensus.simulate( + env.app().timeKeeper().closeTime(), std::chrono::milliseconds{1}); + BEAST_EXPECT(consensus.maxAcceptLockHoldNs() > 0); + } + + void + testAcceptLockWaits() + { + testcase("accept extension block waits while consensus holds the lock"); + using namespace std::chrono_literals; + + std::mutex mu; + std::condition_variable cv; + std::function pending; + bool havePending = false; + std::atomic capture{false}; + bool reaching = false; + std::atomic entered{false}; + std::atomic sawWait{false}; + std::mutex threadMu; + std::thread acceptThread; + + JobQueue::DispatchHook hook = [&](JobType type, + std::string const&, + JobQueue::JobFunction const& func) { + if (capture.load(std::memory_order_acquire) && type == jtACCEPT) + { + std::lock_guard lock(mu); + pending = func; + havePending = true; + return JobQueue::JobDisposition::claimedQueued; + } + return JobQueue::JobDisposition::pass; + }; + + MultiNode net(*this, /*virtualClock=*/true, /*stepping=*/false); + std::vector keys; + std::vector unl; + for (std::size_t i = 0; i < 3; ++i) + { + keys.push_back(ValidatorKey::fromPassphrase( + "accept-lock-" + std::to_string(i))); + unl.push_back(keys.back().pubKey); + } + auto const factory = [](Application& app) { + return std::unique_ptr(std::make_unique(app)); + }; + auto const configHook = [](Config& cfg) { + cfg.NETWORK_ID = networkID; + cfg.features.insert(featureConsensusEntropy); + cfg.features.insert(featureExport); + }; + for (std::size_t i = 0; i < keys.size(); ++i) + { + net.add( + TrustConfig{keys[i].seed, unl}, + factory, + i == 0 ? hook : JobQueue::DispatchHook{}, + /*bindServerListeners=*/false, + configHook); + } + if (!BEAST_EXPECT(net.allUp())) + return; + for (std::size_t i = 0; i < 3; ++i) + for (std::size_t j = i + 1; j < 3; ++j) + if (!BEAST_EXPECT(net.simConnect(i, j) != nullptr)) + return; + if (!BEAST_EXPECT(net.waitForPeers(2, 20s))) + return; + + auto& consensus = net[0].app().getOPs().getConsensus(); + consensus.setAcceptExtensionProbe( + [&] { + std::lock_guard lock(mu); + reaching = true; + cv.notify_all(); + }, + [&] { entered.store(true, std::memory_order_release); }); + consensus.setWhileConsensusLocked([&] { + if (sawWait.load(std::memory_order_acquire)) + return; + std::function job; + { + std::lock_guard lock(mu); + if (!havePending) + return; + job = std::move(pending); + havePending = false; + reaching = false; + } + std::thread worker(std::move(job)); + { + std::unique_lock lock(mu); + cv.wait(lock, [&] { return reaching; }); + } + sawWait.store( + !entered.load(std::memory_order_acquire), + std::memory_order_release); + std::lock_guard lock(threadMu); + acceptThread = std::move(worker); + }); + + capture.store(true, std::memory_order_release); + for (int i = 0; i < 40 && !sawWait.load(std::memory_order_acquire); ++i) + (void)net.threadedTick(1s); + capture.store(false, std::memory_order_release); + { + std::lock_guard lock(threadMu); + if (acceptThread.joinable()) + acceptThread.join(); + } + BEAST_EXPECT(sawWait.load(std::memory_order_acquire)); + BEAST_EXPECT(entered.load(std::memory_order_acquire)); + log << " accept lock hold ns=" << consensus.maxAcceptLockHoldNs() + << std::endl; } public: void run() override { + testSimulateReenters(); + testAcceptLockWaits(); realThreads(); } };