diff --git a/src/test/consensus/ThreadedExtensions_test.cpp b/src/test/consensus/ThreadedExtensions_test.cpp index f362a82112..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 @@ -834,12 +837,141 @@ class ThreadedExtensions_test : public beast::unit_test::suite 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(); } }; diff --git a/src/xrpld/app/consensus/RCLConsensus.cpp b/src/xrpld/app/consensus/RCLConsensus.cpp index d2218b1a3d..0e7449d929 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, @@ -90,8 +91,26 @@ 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, std::unique_ptr&& feeVote, LedgerMaster& ledgerMaster, LocalTxs& localTxs, @@ -99,6 +118,7 @@ RCLConsensus::Adaptor::Adaptor( ValidatorKeys const& validatorKeys, beast::Journal journal) : app_(app) + , consensusMutex_(consensusMutex) , feeVote_(std::move(feeVote)) , ledgerMaster_(ledgerMaster) , localTxs_(localTxs) @@ -518,11 +538,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 +605,18 @@ 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)}; + if (beforeAcceptExtension_) + beforeAcceptExtension_(); + auto const liveBuild = [&] { + std::lock_guard lock{consensusMutex_}; + 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(); @@ -615,7 +641,14 @@ 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_}; + auto const holdStart = std::chrono::steady_clock::now(); + auto salt = ce().txnOrderingSalt(buildTxSetHash, buildSeq); + noteAcceptLock(holdStart); + return salt; + }(); + CanonicalTXSet retriableTxs{orderingSalt}; JLOG(j_.debug()) << "Building canonical tx set: " << retriableTxs.key(); @@ -641,17 +674,26 @@ 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_}; + auto const holdStart = std::chrono::steady_clock::now(); + if (insideAcceptExtension_) + insideAcceptExtension_(); + if (replayData) + { + ce().onReplayBuild(); + } + else if (ce().rngEnabled() || ce().exportEnabled()) + { + ce().onPreBuild(retriableTxs, buildSeq, buildTxSetHash); + } + else + { + ce().clearRngState(); + } + noteAcceptLock(holdStart); } //@@end accept-time-cleanup-disabled @@ -1120,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 0792ac23e0..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 @@ -69,6 +70,7 @@ class RCLConsensus class Adaptor { Application& app_; + std::recursive_mutex& consensusMutex_; std::unique_ptr feeVote_; LedgerMaster& ledgerMaster_; LocalTxs& localTxs_; @@ -114,6 +116,7 @@ class RCLConsensus Adaptor( Application& app, + std::recursive_mutex& consensusMutex, std::unique_ptr&& feeVote, LedgerMaster& ledgerMaster, LocalTxs& localTxs, @@ -202,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 @@ -210,12 +230,19 @@ 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; + 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. @@ -503,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; @@ -552,11 +601,15 @@ 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_; + std::function whileConsensusLocked_; + Adaptor adaptor_; std::unique_ptr> consensus_; beast::Journal const j_;