merge fix-accept-consensus-lock: consensus mutex around accept-time extension prep, probe seams, waits-on-lock case

This commit is contained in:
Nicholas Dudfield
2026-09-23 09:53:03 +07:00
3 changed files with 256 additions and 27 deletions

View File

@@ -2,6 +2,7 @@
#include <test/jtx/pay.h>
#include <xrpld/app/consensus/ConsensusExtensions.h>
#include <xrpld/app/consensus/RCLConsensus.h>
#include <xrpld/app/ledger/Ledger.h>
#include <xrpld/app/ledger/LedgerMaster.h>
#include <xrpld/app/misc/RuntimeConfig.h>
@@ -24,8 +25,10 @@
#include <array>
#include <atomic>
#include <chrono>
#include <condition_variable>
#include <cstdint>
#include <limits>
#include <mutex>
#include <optional>
#include <set>
#include <string>
@@ -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<void()> pending;
bool havePending = false;
std::atomic<bool> capture{false};
bool reaching = false;
std::atomic<bool> entered{false};
std::atomic<bool> 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<ValidatorKey> keys;
std::vector<std::string> 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<Overlay>(std::make_unique<SimOverlay>(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<void()> 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();
}
};

View File

@@ -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::nanoseconds>(
std::chrono::steady_clock::now() - start)
.count();
if (ns <= 0)
return;
auto cur = maxAcceptLockNs_.load(std::memory_order_relaxed);
while (static_cast<std::uint64_t>(ns) > cur &&
!maxAcceptLockNs_.compare_exchange_weak(
cur, static_cast<std::uint64_t>(ns), std::memory_order_relaxed))
{
}
}
RCLConsensus::Adaptor::Adaptor(
Application& app,
std::recursive_mutex& consensusMutex,
std::unique_ptr<FeeVote>&& 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<ConsensusExtensions::LiveBuildTxSet>{}
: std::optional<ConsensusExtensions::LiveBuildTxSet>{
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<ConsensusExtensions::LiveBuildTxSet>{}
: std::optional<ConsensusExtensions::LiveBuildTxSet>{
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)
{

View File

@@ -40,6 +40,7 @@
#include <xrpl/protocol/STValidation.h>
#include <atomic>
#include <chrono>
#include <functional>
#include <memory>
#include <mutex>
#include <set>
@@ -69,6 +70,7 @@ class RCLConsensus
class Adaptor
{
Application& app_;
std::recursive_mutex& consensusMutex_;
std::unique_ptr<FeeVote> feeVote_;
LedgerMaster& ledgerMaster_;
LocalTxs& localTxs_;
@@ -114,6 +116,7 @@ class RCLConsensus
Adaptor(
Application& app,
std::recursive_mutex& consensusMutex,
std::unique_ptr<FeeVote>&& 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<void()> beforeLock,
std::function<void()> 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<Adaptor> 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<Adaptor>;
std::function<void()> beforeAcceptExtension_;
std::function<void()> insideAcceptExtension_;
std::atomic<std::uint64_t> 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<void()> hook)
{
whileConsensusLocked_ = std::move(hook);
}
void
setAcceptExtensionProbe(
std::function<void()> beforeLock,
std::function<void()> 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<void()> whileConsensusLocked_;
Adaptor adaptor_;
std::unique_ptr<Consensus<Adaptor>> consensus_;
beast::Journal const j_;