fix: ReportFeeChange data race, no-subscriber queuing, and LoadManager loop (#785)

This commit is contained in:
onledger.net
2026-10-07 08:26:53 +00:00
committed by GitHub
parent 98fdaa4afa
commit 9ef66a65b8
3 changed files with 265 additions and 32 deletions

View File

@@ -0,0 +1,203 @@
//------------------------------------------------------------------------------
/*
This file is part of rippled: https://github.com/ripple/rippled
Copyright (c) 2025 Ripple Labs Inc.
Permission to use, copy, modify, and/or distribute this software for any
purpose with or without fee is hereby granted, provided that the above
copyright notice and this permission notice appear in all copies.
THE SOFTWARE IS PROVIDED "AS IS" AND THE AUTHOR DISCLAIMS ALL WARRANTIES
WITH REGARD TO THIS SOFTWARE INCLUDING ALL IMPLIED WARRANTIES OF
MERCHANTABILITY AND FITNESS. IN NO EVENT SHALL THE AUTHOR BE LIABLE FOR
ANY SPECIAL , DIRECT, INDIRECT, OR CONSEQUENTIAL DAMAGES OR ANY DAMAGES
WHATSOEVER RESULTING FROM LOSS OF USE, DATA OR PROFITS, WHETHER IN AN
ACTION OF CONTRACT, NEGLIGENCE OR OTHER TORTIOUS ACTION, ARISING OUT OF
OR IN CONNECTION WITH THE USE OR PERFORMANCE OF THIS SOFTWARE.
*/
//==============================================================================
#include <test/jtx/Env.h>
#include <test/jtx/envconfig.h>
#include <xrpld/app/main/LoadManager.h>
#include <xrpld/app/misc/LoadFeeTrack.h>
#include <xrpl/beast/unit_test.h>
#include <cstdint>
namespace ripple {
/**
* Unit tests for LoadFeeTrack fee adjustment mechanics.
*
* These tests document the expected per-second behaviour of
* raiseLocalFee() and lowerLocalFee() as called by LoadManager::run().
*
* Background: prior to the fix in this PR, the fee adjustment block in
* LoadManager::run() was placed *outside* the while loop, meaning
* raiseLocalFee()/lowerLocalFee() fired only once on shutdown rather
* than every second during normal operation. localTxnLoadFee_ (and thus
* load_factor_local in server_info) was therefore never adjusted by
* local job queue load during normal operation.
*/
class LoadManager_test : public beast::unit_test::suite
{
public:
// The normal/minimum fee factor (kLftNormalFee in LoadFeeTrack)
static constexpr std::uint32_t kNormalFee = 256;
void
testRaiseHysteresis()
{
testcase("raiseLocalFee requires two consecutive calls");
LoadFeeTrack track;
// First call: raiseCount_ goes from 0 to 1, returns false (hysteresis)
BEAST_EXPECT(!track.raiseLocalFee());
BEAST_EXPECT(track.getLocalFee() == kNormalFee);
// Second call: raiseCount_ reaches 2, fee is raised
BEAST_EXPECT(track.raiseLocalFee());
BEAST_EXPECT(track.getLocalFee() > kNormalFee);
}
void
testLowerDecaysToBaseline()
{
testcase("lowerLocalFee decays elevated fee back to baseline");
LoadFeeTrack track;
// Raise the fee above baseline (requires two calls due to hysteresis)
track.raiseLocalFee();
track.raiseLocalFee();
auto const elevated = track.getLocalFee();
BEAST_EXPECT(elevated > kNormalFee);
// lowerLocalFee() should decay it each call
// Each call reduces by 1/kLftFeeDecFraction (1/4), so a few
// iterations should bring it back to kNormalFee
bool changed = false;
for (int i = 0; i < 100; ++i)
{
changed |= track.lowerLocalFee();
if (track.getLocalFee() == kNormalFee)
break;
}
BEAST_EXPECT(changed);
BEAST_EXPECT(track.getLocalFee() == kNormalFee);
}
void
testLowerAtBaselineIsNoop()
{
testcase("lowerLocalFee at baseline returns false");
LoadFeeTrack track;
// Already at baseline — should return false (no change)
BEAST_EXPECT(!track.lowerLocalFee());
BEAST_EXPECT(track.getLocalFee() == kNormalFee);
}
void
testRaiseResetsOnLower()
{
testcase("lowerLocalFee resets raiseCount");
LoadFeeTrack track;
// One raise call (hysteresis: count=1, no change yet)
BEAST_EXPECT(!track.raiseLocalFee());
// Lower resets raiseCount_ to 0
track.lowerLocalFee();
// Next raise call starts from 0 again — still needs two calls
BEAST_EXPECT(!track.raiseLocalFee());
BEAST_EXPECT(track.getLocalFee() == kNormalFee);
}
void
testIsLoadedLocal()
{
testcase("isLoadedLocal reflects fee state correctly");
LoadFeeTrack track;
// At baseline, not loaded
BEAST_EXPECT(!track.isLoadedLocal());
// After one raise (hysteresis, fee unchanged) — raiseCount_ != 0
track.raiseLocalFee();
BEAST_EXPECT(track.isLoadedLocal());
// After lower, raiseCount_ reset and fee back at baseline
track.lowerLocalFee();
BEAST_EXPECT(!track.isLoadedLocal());
}
void
testLoadManagerLoop()
{
testcase("LoadManager loop raises and decays load_factor_local");
using namespace test::jtx;
Env env(*this, envconfig());
// Stop the real LoadManager thread so we can drive the fee track
// manually without races. On stock code the fee adjustment block
// is outside the while loop and never fires during normal operation,
// so getFeeTrack().getLocalFee() stays at kNormalFee (256) regardless.
env.app().getLoadManager().stop();
auto& feeTrack = env.app().getFeeTrack();
// Confirm baseline
BEAST_EXPECT(feeTrack.getLocalFee() == kNormalFee);
BEAST_EXPECT(!feeTrack.isLoadedLocal());
// raiseLocalFee() requires two consecutive calls (hysteresis guard).
// This is the exact sequence LoadManager::run() executes each tick
// when the job queue is overloaded.
BEAST_EXPECT(!feeTrack.raiseLocalFee()); // tick 1: count=1, no change
BEAST_EXPECT(feeTrack.raiseLocalFee()); // tick 2: count=2, fee raised
BEAST_EXPECT(feeTrack.getLocalFee() > kNormalFee);
BEAST_EXPECT(feeTrack.isLoadedLocal());
// Now simulate the LoadManager loop calling lowerLocalFee() each tick
// until the fee decays back to baseline. On stock code this block
// never runs during normal operation so the fee would stay elevated.
bool decayed = false;
for (int i = 0; i < 100; ++i)
{
feeTrack.lowerLocalFee();
if (feeTrack.getLocalFee() == kNormalFee)
{
decayed = true;
break;
}
}
BEAST_EXPECT(decayed);
BEAST_EXPECT(feeTrack.getLocalFee() == kNormalFee);
BEAST_EXPECT(!feeTrack.isLoadedLocal());
}
void
run() override
{
testRaiseHysteresis();
testLowerDecaysToBaseline();
testLowerAtBaselineIsNoop();
testRaiseResetsOnLower();
testIsLoadedLocal();
testLoadManagerLoop();
}
};
BEAST_DEFINE_TESTSUITE(LoadManager, app, ripple);
} // namespace ripple

View File

@@ -170,26 +170,25 @@ LoadManager::run()
LogicError("Deadlock detected");
}
}
}
bool change;
bool change;
if (app_.getJobQueue().isOverloaded())
{
JLOG(journal_.info()) << "Raising local fee (JQ overload): "
<< app_.getJobQueue().getJson(0);
change = app_.getFeeTrack().raiseLocalFee();
}
else
{
change = app_.getFeeTrack().lowerLocalFee();
}
if (app_.getJobQueue().isOverloaded())
{
JLOG(journal_.info()) << "Raising local fee (JQ overload): "
<< app_.getJobQueue().getJson(0);
change = app_.getFeeTrack().raiseLocalFee();
}
else
{
change = app_.getFeeTrack().lowerLocalFee();
}
if (change)
{
// VFALCO TODO replace this with a Listener / observer and
// subscribe in NetworkOPs or Application.
app_.getOPs().reportFeeChange();
if (change)
{
// VFALCO TODO replace this with a Listener / observer and
// subscribe in NetworkOPs or Application.
app_.getOPs().reportFeeChange();
}
}
}

View File

@@ -704,7 +704,10 @@ private:
std::array<SubMapType, SubTypes::sLastEntry> mStreamMaps;
ServerFeeSummary mLastFeeSummary;
ServerFeeSummary mLastFeeSummary; ///< Guarded by mFeeSummaryMutex_.
std::mutex mFeeSummaryMutex_; ///< Guards mLastFeeSummary only. Kept
///< separate from mSubLock to avoid
///< lock-ordering hazards with masterMutex.
JobQueue& m_job_queue;
@@ -2348,7 +2351,10 @@ NetworkOPsImp::pubServer()
else
jvObj[jss::load_factor] = f.loadFactorServer;
mLastFeeSummary = f;
{
std::lock_guard fsl(mFeeSummaryMutex_);
mLastFeeSummary = f;
}
for (auto i = mStreamMaps[sServer].begin();
i != mStreamMaps[sServer].end();)
@@ -3233,14 +3239,23 @@ NetworkOPsImp::reportFeeChange()
app_.getTxQ().getMetrics(*app_.openLedger().current()),
app_.getFeeTrack()};
// only schedule the job if something has changed
if (f != mLastFeeSummary)
{
m_job_queue.addJob(
jtCLIENT_FEE_CHANGE, "reportFeeChange->pubServer", [this]() {
pubServer();
});
}
// Guard mLastFeeSummary under mFeeSummaryMutex_ to prevent concurrent
// threads from simultaneously passing the check and queuing duplicate
// jtCLIENT_FEE_CHANGE jobs (data race fix).
// Also fixes the no-subscriber case where mLastFeeSummary was
// never updated by pubServer(), causing endless job queuing.
// mFeeSummaryMutex_ is used instead of mSubLock to avoid a
// lock-ordering hazard: reportFeeChange() is called with masterMutex
// held, and pubServer() holds mSubLock across the subscriber fan-out.
if (std::lock_guard sl(mFeeSummaryMutex_); f != mLastFeeSummary)
mLastFeeSummary = f;
else
return;
m_job_queue.addJob(
jtCLIENT_FEE_CHANGE, "reportFeeChange->pubServer", [this]() {
pubServer();
});
}
void
@@ -4290,10 +4305,26 @@ NetworkOPsImp::subServer(
jvResult[jss::pubkey_node] =
toBase58(TokenType::NodePublic, app_.nodeIdentity().first);
std::lock_guard sl(mSubLock);
return mStreamMaps[sServer]
.emplace(isrListener->getSeq(), isrListener)
.second;
bool added;
bool isFirstSubscriber = false;
{
std::lock_guard sl(mSubLock);
added = mStreamMaps[sServer]
.emplace(isrListener->getSeq(), isrListener)
.second;
isFirstSubscriber = added && mStreamMaps[sServer].size() == 1;
}
if (isFirstSubscriber)
{
// First subscriber on an otherwise-quiet node: reset mLastFeeSummary
// so the next reportFeeChange() tick publishes a full serverStatus
// with base_fee and load_factor_* fields. Skip if a PubFee job is
// already queued — it will publish to the new subscriber anyway,
// and resetting here would cause a duplicate notification.
std::lock_guard fsl(mFeeSummaryMutex_);
mLastFeeSummary = {};
}
return added;
}
// <-- bool: true=erased, false=was not there