Compare commits

...

3 Commits

Author SHA1 Message Date
Denis Angell
4978981608 chore: bump openssl to 3.6.3 for conan.xrplf.org compatibility 2026-08-06 22:52:14 -04:00
Denis Angell
78e859f56c feat: consensus-testing toolkit (node_stall/inject/unl_set) + DatagramMonitor field population 2026-07-15 15:43:58 -04:00
Denis Angell
2add5b3fb8 feat: Add Datagram 2026-05-30 12:18:41 +02:00
21 changed files with 2156 additions and 37 deletions

View File

@@ -1,43 +1,35 @@
{
"version": "0.5",
"requires": [
"zlib/1.3.2#1cb806da49011867778ffb6ac7190fcb%1777558780.503",
"xxhash/0.8.3#681d36a0a6111fc56e5e45ea182c19cc%1765850149.987",
"sqlite3/3.53.0#324ada52333108388a9a6108bfa96734%1776096494.149",
"soci/4.0.3#fe32b9ad5eb47e79ab9e45a68f363945%1774450067.231",
"snappy/1.1.10#968fef506ff261592ec30c574d4a7809%1765850147.878",
"secp256k1/0.7.1#481881709eb0bdd0185a12b912bbe8ad%1770910500.329",
"rocksdb/10.5.1#4a197eca381a3e5ae8adf8cffa5aacd0%1765850186.86",
"re2/20251105#8579cfd0bda4daf0683f9e3898f964b4%1774398111.888",
"protobuf/6.33.5#d96d52ba5baaaa532f47bda866ad87a5%1774467363.12",
"openssl/3.6.2#4789bbf131b77d0515d15e094c8f697f%1778071755.506",
"nudb/2.0.9#11149c73f8f2baff9a0198fe25971fc7%1775040983.408",
"lz4/1.10.0#59fc63cac7f10fbe8e05c7e62c2f3504%1765850143.914",
"libiconv/1.17#1e65319e945f2d31941a9d28cc13c058%1765842973.492",
"libbacktrace/cci.20210118#a7691bfccd8caaf66309df196790a5a1%1765842973.03",
"libarchive/3.8.7#c446109bd1f1d8ba7936c94189bc50e6%1776147552.838",
"jemalloc/5.3.1#1fc58d55316041f10fbc1e8a2eae632a%1776700028.228",
"gtest/1.17.0#5224b3b3ff3b4ce1133cbdd27d53ee7d%1768312129.152",
"grpc/1.78.1#b1a9e74b145cc471bed4dc64dc6eb2c1%1774467387.342",
"ed25519/2015.03#ae761bdc52730a843f0809bdf6c1b1f6%1765850143.772",
"date/3.0.4#862e11e80030356b53c2c38599ceb32b%1765850143.772",
"c-ares/1.34.6#545240bb1c40e2cacd4362d6b8967650%1774439234.681",
"bzip2/1.0.8#c470882369c2d95c5c77e970c0c7e321%1765850143.837",
"boost/1.91.0#ea540ca2133d831b560036aa24dece3c%1778050991.9",
"abseil/20250127.0#bb0baf1f362bc4a725a24eddd419b8f7%1774365460.196"
"zlib/1.3.2#1cb806da49011867778ffb6ac7190fcb%1782392402.122708",
"xxhash/0.8.3#681d36a0a6111fc56e5e45ea182c19cc%1782392402.420688",
"sqlite3/3.53.0#324ada52333108388a9a6108bfa96734%1782392403.185447",
"soci/4.0.3#e726491a03468795453f7c83fc924a96%1782392402.679521",
"snappy/1.1.10#968fef506ff261592ec30c574d4a7809%1782307151.633168",
"secp256k1/0.7.1#b1f450b7f78a36fff75bb6934a356f3a%1782338841.3729",
"rocksdb/10.5.1#4a197eca381a3e5ae8adf8cffa5aacd0%1782392413.075713",
"re2/20251105#8579cfd0bda4daf0683f9e3898f964b4%1782392402.431897",
"protobuf/6.33.5#ff253ead763bd8d9904a52979cd21e81%1782392410.233933",
"openssl/3.6.3#f806de8933e3bf6f01016c6a888cee2e%1783945160.863288",
"nudb/2.0.9#11149c73f8f2baff9a0198fe25971fc7%1782392402.297166",
"lz4/1.10.0#982d9b673900f665a1da109e09c17cab%1782392402.164188",
"libbacktrace/cci.20210118#a7691bfccd8caaf66309df196790a5a1%1782392402.420732",
"libarchive/3.8.7#c446109bd1f1d8ba7936c94189bc50e6%1782392403.066892",
"gtest/1.17.0#5224b3b3ff3b4ce1133cbdd27d53ee7d%1782392402.791979",
"grpc/1.78.1#b1a9e74b145cc471bed4dc64dc6eb2c1%1782736970.619035",
"ed25519/2015.03#ae761bdc52730a843f0809bdf6c1b1f6%1782307148.15562",
"date/3.0.4#862e11e80030356b53c2c38599ceb32b%1782392402.538492",
"c-ares/1.34.6#545240bb1c40e2cacd4362d6b8967650%1782392402.681654",
"bzip2/1.0.8#c470882369c2d95c5c77e970c0c7e321%1782392402.296732",
"boost/1.91.0#ea540ca2133d831b560036aa24dece3c%1782392419.475605",
"abseil/20250127.0#9ef01c1451a8340f9022e46238c0fbb6%1783945159.651047"
],
"build_requires": [
"zlib/1.3.2#1cb806da49011867778ffb6ac7190fcb%1777558780.503",
"strawberryperl/5.32.1.1#8d114504d172cfea8ea1662d09b6333e%1774447376.964",
"protobuf/6.33.5#d96d52ba5baaaa532f47bda866ad87a5%1774467363.12",
"nasm/2.16.01#31e26f2ee3c4346ecd347911bd126904%1765850144.707",
"msys2/cci.latest#d22fe7b2808f5fd34d0a7923ace9c54f%1770657326.649",
"m4/1.4.19#4523e4347b55cd26ae918bd5770cab9a%1778062762.471",
"cmake/4.3.0#b939a42e98f593fb34d3a8c5cc860359%1774439249.183",
"b2/5.4.2#ffd6084a119587e70f11cd45d1a386e2%1774439233.447",
"automake/1.16.5#b91b7c384c3deaa9d535be02da14d04f%1755524470.56",
"autoconf/2.71#51077f068e61700d65bb05541ea1e4b0%1731054366.86",
"abseil/20250127.0#bb0baf1f362bc4a725a24eddd419b8f7%1774365460.196"
"zlib/1.3.2#1cb806da49011867778ffb6ac7190fcb%1782392402.122708",
"protobuf/6.33.5#ff253ead763bd8d9904a52979cd21e81%1782392410.233933",
"cmake/4.3.3#840cf00ea09777e05c2050a50a82c722%1782392418.696091",
"b2/5.4.2#ffd6084a119587e70f11cd45d1a386e2%1782392402.624226",
"abseil/20250127.0#9ef01c1451a8340f9022e46238c0fbb6%1783945159.651047"
],
"python_requires": [],
"overrides": {
@@ -57,7 +49,7 @@
"boost/1.91.0"
],
"lz4/[>=1.9.4 <2]": [
"lz4/1.10.0#59fc63cac7f10fbe8e05c7e62c2f3504"
"lz4/1.10.0#982d9b673900f665a1da109e09c17cab"
]
},
"config_requires": []

View File

@@ -31,7 +31,7 @@ class Xrpl(ConanFile):
"grpc/1.78.1",
"libarchive/3.8.7",
"nudb/2.0.9",
"openssl/3.6.2",
"openssl/3.6.3",
"secp256k1/0.7.1",
"soci/4.0.3",
"zlib/1.3.2",

View File

@@ -329,6 +329,7 @@ words:
- writeme
- wsrch
- wthread
- Xahau
- xbridge
- xchain
- ximinez

View File

@@ -291,6 +291,7 @@ JSS(ident); // in: AccountCurrencies, AccountInfo, OwnerIn
JSS(ignore_default); // in: AccountLines
JSS(in); // out: OverlayImpl
JSS(inLedger); // out: tx/Transaction
JSS(in_queue); // out: inject
JSS(inbound); // out: PeerImp
JSS(index); // in: LedgerEntry
// out: STLedgerEntry, LedgerEntry, TxHistory, LedgerData

View File

@@ -10,7 +10,9 @@
#include <boost/asio.hpp>
#include <array>
#include <memory>
#include <tuple>
namespace xrpl {
@@ -72,6 +74,18 @@ class NetworkOPs : public InfoSub::Source
public:
using clock_type = beast::AbstractClock<std::chrono::steady_clock>;
// Snapshot of per-operating-mode accounting, exposed for the datagram monitor.
struct AccountingCounter
{
std::uint64_t transitions{0};
std::chrono::microseconds dur{std::chrono::microseconds(0)};
};
using StateAccountingData = std::tuple<
std::array<AccountingCounter, 5>,
OperatingMode,
std::chrono::steady_clock::time_point,
std::uint64_t>;
enum class FailHard : unsigned char { No, Yes };
static FailHard
doFailHard(bool noMeansDont)
@@ -92,6 +106,8 @@ public:
[[nodiscard]] virtual OperatingMode
getOperatingMode() const = 0;
[[nodiscard]] virtual StateAccountingData
getStateAccountingData() = 0;
[[nodiscard]] virtual std::string
strOperatingMode(OperatingMode const mode, bool const admin = false) const = 0;
[[nodiscard]] virtual std::string
@@ -207,8 +223,21 @@ public:
virtual void
consensusViewChange() = 0;
virtual void
setStall(std::chrono::milliseconds duration) = 0;
virtual bool
isStalled() const = 0;
virtual void
clearStall() = 0;
virtual json::Value
getConsensusInfo() = 0;
// Proposers and round time of the last consensus round, for out-of-band
// telemetry (DatagramMonitor) that cannot reach the private consensus object.
[[nodiscard]] virtual std::size_t
getPrevProposers() const = 0;
[[nodiscard]] virtual std::chrono::milliseconds
getPrevRoundTime() const = 0;
virtual json::Value
getServerInfo(bool human, bool admin, bool counters) = 0;
virtual void

View File

@@ -0,0 +1,899 @@
/**
* Byzantine Partition Recovery Tests
*
* Tests consensus safety under network partition scenarios where a minority
* of validators (running modified binaries) attempt to advance the chain
* while the majority is offline, then rejoin the network.
*
* Attack scenario:
* - 7 validators share a common UNL (like production mainnet)
* - 4 legitimate validators crash (DoS attack)
* - 3 attacker validators (modified binary) bypass quorum checks
* and continue producing ledgers
* - 4 legitimate validators recover and reconnect
* - Question: does the attacker chain get accepted?
*/
#include <test/csf.h>
#include <test/csf/Peer.h>
#include <test/csf/PeerGroup.h>
#include <test/csf/Sim.h>
#include <test/csf/SimTime.h>
#include <test/csf/TrustGraph.h>
#include <test/csf/collectors.h>
#include <test/csf/events.h>
#include <xrpld/consensus/ConsensusParms.h>
#include <xrpl/beast/unit_test/suite.h>
#include <chrono>
#include <iostream>
namespace xrpl::test {
class ByzantinePartitionRecovery_test : public beast::unit_test::Suite
{
// Collector that tracks per-peer ledger advancement and validation
struct PartitionTracker
{
struct PeerState
{
csf::Ledger::Seq lastClosed{0};
csf::Ledger::Seq lastFullyValidated{0};
csf::Ledger::ID lastClosedId{};
csf::Ledger::ID lastFullyValidatedId{};
};
std::map<csf::PeerID, PeerState> states;
template <class E>
void
on(csf::PeerID, csf::SimTime, E const&)
{
}
void
on(csf::PeerID who, csf::SimTime, csf::AcceptLedger const& e)
{
auto& s = states[who];
s.lastClosed = e.ledger.seq();
s.lastClosedId = e.ledger.id();
}
void
on(csf::PeerID who, csf::SimTime, csf::FullyValidateLedger const& e)
{
auto& s = states[who];
s.lastFullyValidated = e.ledger.seq();
s.lastFullyValidatedId = e.ledger.id();
}
};
/**
* Test 1: Byzantine minority cannot fully validate during partition
*
* 7 validators, common UNL. 4 crash. 3 attackers continue.
* Standard quorum (no modification) — attackers can close ledgers
* but cannot fully validate them (need 6/7 = 80%).
*/
void
testPartitionNoQuorumBypass()
{
using namespace csf;
using namespace std::chrono;
testcase("partition: minority cannot fully validate (standard quorum)");
Sim sim;
ConsensusParms const parms{};
SimDuration const delay = round<milliseconds>(0.2 * parms.ledgerGRANULARITY);
// Create 7 validators
PeerGroup attackers = sim.createGroup(3);
PeerGroup legitimate = sim.createGroup(4);
PeerGroup network = attackers + legitimate;
// Common UNL — all trust all (like production mainnet)
network.trust(network);
network.connect(network, delay);
PartitionTracker tracker;
sim.collectors.add(tracker);
// Round 1: establish common state (all 7 in sync)
sim.run(1);
BEAST_EXPECT(sim.synchronized());
// Record pre-partition state
Ledger::Seq const prePartitionSeq = attackers[0]->lastClosedLedger.seq();
Ledger::ID const prePartitionId = attackers[0]->fullyValidatedLedger.id();
std::cout << "Pre-partition: all peers at seq " << prePartitionSeq << ", fully validated\n";
// === PARTITION: disconnect legitimate validators ===
// Simulate crash: disconnect from attackers AND from each other
legitimate.disconnect(network);
network.disconnect(legitimate);
// Stop legitimate peers from running consensus during partition
for (Peer* p : legitimate)
p->targetLedgers = p->completedLedgers;
// Attackers submit transactions and try to advance
for (Peer* p : attackers)
p->submit(Tx{static_cast<std::uint32_t>(p->id) + 100});
// Run several rounds — attackers will close ledgers among themselves
// (3/3 = 100% of visible proposers agree)
sim.run(4);
// Check attacker state
Ledger::Seq attackerSeq = attackers[0]->lastClosedLedger.seq();
Ledger::Seq legitimateSeq = legitimate[0]->lastClosedLedger.seq();
std::cout << "During partition:\n";
std::cout << " Attackers closed up to seq " << attackerSeq << "\n";
std::cout << " Legitimate peers stuck at seq " << legitimateSeq << "\n";
// Attackers advanced their closed ledger
BEAST_EXPECT(attackerSeq > prePartitionSeq);
// Note: legitimate peers may also advance closed ledger via
// consensus timeout (isolated node eventually closes alone),
// but they CANNOT fully validate.
// KEY ASSERTION: attackers could NOT fully validate their chain
// because they only have 3/7 validations (need 6)
for (Peer* p : attackers)
{
std::cout << " Attacker " << p->id
<< " fullyValidated seq: " << p->fullyValidatedLedger.seq() << "\n";
// The attacker's fully validated ledger should NOT have advanced
// beyond the pre-partition state (they can't get 6/7 validations)
BEAST_EXPECT(p->fullyValidatedLedger.id() == prePartitionId);
}
}
/**
* Test 2: Byzantine minority with quorum bypass
*
* 3 attackers override their quorum to 2 (simulating modified binary).
* They can fully validate on their fork. When 4 legitimate peers
* rejoin, the attacker chain should NOT be accepted by legitimate peers.
*/
void
testPartitionWithQuorumBypass()
{
using namespace csf;
using namespace std::chrono;
testcase("partition: attacker chain rejected after recovery (quorum bypass)");
Sim sim;
ConsensusParms const parms{};
SimDuration const delay = round<milliseconds>(0.2 * parms.ledgerGRANULARITY);
PeerGroup attackers = sim.createGroup(3);
PeerGroup legitimate = sim.createGroup(4);
PeerGroup network = attackers + legitimate;
// Common UNL
network.trust(network);
network.connect(network, delay);
PartitionTracker tracker;
sim.collectors.add(tracker);
// Round 1: establish common state
sim.run(1);
BEAST_EXPECT(sim.synchronized());
Ledger::Seq const prePartitionSeq = attackers[0]->lastClosedLedger.seq();
std::cout << "Pre-partition: all peers synced at seq " << prePartitionSeq << "\n";
// === PARTITION ===
legitimate.disconnect(network);
network.disconnect(legitimate);
for (Peer* p : legitimate)
p->targetLedgers = p->completedLedgers;
// SIMULATE MODIFIED BINARY: attackers bypass quorum
// In production, the attacker modifies ValidatorList::calculateQuorum()
// to return a lower value. In CSF, we override the quorum directly.
// Note: quorum is recalculated in checkFullyValidated() from
// trustGraph.graph().outDegree(this), so we need a different approach.
// We'll untrust the legitimate peers FROM THE ATTACKER'S PERSPECTIVE
// to simulate the modified binary lowering effective UNL.
attackers.untrust(legitimate);
// Now attackers only trust 3 peers → quorum = ceil(3 * 0.8) = 3
// They can fully validate with just their own validations
// Attackers submit transactions
for (Peer* p : attackers)
p->submit(Tx{static_cast<std::uint32_t>(p->id) + 200});
// Run — attackers advance and fully validate their fork
sim.run(4);
Ledger::Seq const attackerSeq = attackers[0]->lastClosedLedger.seq();
std::cout << "During partition (quorum bypass):\n";
std::cout << " Attacker closed seq: " << attackerSeq << "\n";
for (Peer* p : attackers)
{
std::cout << " Attacker " << p->id
<< " fullyValidated seq: " << p->fullyValidatedLedger.seq() << "\n";
}
// Attackers DID advance and fully validate (on their modified view)
BEAST_EXPECT(attackerSeq > prePartitionSeq);
for (Peer* p : attackers)
{
BEAST_EXPECT(p->fullyValidatedLedger.seq() > prePartitionSeq);
}
// === RECOVERY: restore trust and reconnect ===
// Restore attacker trust to full UNL (simulating reconnection
// where the legitimate validators are back)
attackers.trust(legitimate);
// Reconnect network
network.connect(network, delay);
// Allow legitimate peers to run again
for (Peer* p : legitimate)
p->targetLedgers = p->completedLedgers + 10;
// Run several rounds for recovery
sim.run(6);
std::cout << "\nAfter recovery:\n";
std::cout << " Branches: " << sim.branches() << "\n";
std::cout << " Synchronized: " << std::boolalpha << sim.synchronized() << "\n";
for (Peer* p : network)
{
std::cout << " Peer " << p->id << " closed=" << p->lastClosedLedger.seq()
<< " fullyValidated=" << p->fullyValidatedLedger.seq() << "\n";
}
// KEY ASSERTIONS:
// 1. Legitimate peers should NOT have accepted the attacker's
// fully validated ledger — their quorum is still 6/7
for (Peer* p : legitimate)
{
// Legitimate peers' fully validated ledger should be based on
// the pre-partition state or a new ledger with proper quorum,
// NOT the attacker's fork
std::size_t const numTrusted = sim.trustGraph.graph().outDegree(p);
std::cout << " Legitimate peer " << p->id << " trusts " << numTrusted
<< " quorum=" << static_cast<std::size_t>(std::ceil(numTrusted * 0.8))
<< "\n";
}
// 2. Check if the network re-converges or stays forked
std::size_t const branches = sim.branches();
std::cout << " Final branch count: " << branches << "\n";
// The attacker's chain (validated with only 3/7 trust)
// should not be accepted by the full network.
// After reconnection, with 7 validators all trusting each other,
// quorum is 6. The attacker's old validations only had 3.
// The network should eventually converge on a new chain.
}
/**
* Test 3: Transaction injection during partition
*
* 3 attackers inject extra transactions (simulating pseudo-tx injection)
* during the partition. Tests whether those transactions persist after
* recovery.
*/
void
testTxInjectionDuringPartition()
{
using namespace csf;
using namespace std::chrono;
testcase("partition: injected transactions during partition");
Sim sim;
ConsensusParms const parms{};
SimDuration const delay = round<milliseconds>(0.2 * parms.ledgerGRANULARITY);
PeerGroup attackers = sim.createGroup(3);
PeerGroup legitimate = sim.createGroup(4);
PeerGroup network = attackers + legitimate;
network.trust(network);
network.connect(network, delay);
PartitionTracker tracker;
sim.collectors.add(tracker);
// Establish common state
sim.run(1);
BEAST_EXPECT(sim.synchronized());
Ledger::Seq const prePartitionSeq = attackers[0]->lastClosedLedger.seq();
// === PARTITION ===
legitimate.disconnect(network);
network.disconnect(legitimate);
for (Peer* p : legitimate)
p->targetLedgers = p->completedLedgers;
// Attackers bypass quorum (modified binary simulation)
attackers.untrust(legitimate);
// Inject a "malicious" transaction into attacker ledgers
// This simulates an EnableAmendment pseudo-tx being injected
// by the modified binary's doVoting()
Tx const maliciousTx{999};
for (Peer* p : attackers)
{
p->txInjections.emplace(prePartitionSeq, maliciousTx);
}
// Run partition phase
sim.run(4);
// Check that attackers included the injected tx
Ledger const& attackerLedger = attackers[0]->lastClosedLedger;
std::cout << "Attacker ledger seq: " << attackerLedger.seq() << "\n";
// === RECOVERY ===
attackers.trust(legitimate);
network.connect(network, delay);
for (Peer* p : legitimate)
p->targetLedgers = p->completedLedgers + 10;
sim.run(6);
std::cout << "\nAfter recovery (tx injection):\n";
std::cout << " Branches: " << sim.branches() << "\n";
std::cout << " Synchronized: " << std::boolalpha << sim.synchronized() << "\n";
// Check each peer's final state
for (Peer* p : network)
{
std::cout << " Peer " << p->id << " closed=" << p->lastClosedLedger.seq()
<< " fullyValidated=" << p->fullyValidatedLedger.seq() << "\n";
}
// The legitimate peers should NOT have the injected transaction
// in their fully validated ledger chain. The attacker's fork
// (containing the injected tx) should be rejected because it
// lacks sufficient validations from the full UNL.
}
/**
* Test 4: Full attack surface sweep
*
* Sweeps across multiple network sizes (7, 10, 15, 20, 35) and
* all attacker counts to find the exact threshold where the
* Byzantine partition attack succeeds.
*/
void
testAttackSurfaceSweep()
{
using namespace csf;
using namespace std::chrono;
testcase("partition: full attack surface sweep");
// Network sizes to test (7 = small, 35 = mainnet-like)
std::vector<std::uint32_t> const networkSizes = {7, 10, 15, 20, 35};
for (std::uint32_t const totalValidators : networkSizes)
{
std::cout << "\n=== Network size: " << totalValidators << " validators ===\n";
std::cout << " 80% quorum = "
<< static_cast<std::size_t>(std::ceil(totalValidators * 0.8)) << "\n";
std::uint32_t attackThreshold = 0;
for (std::uint32_t numAttackers = 1; numAttackers < totalValidators; ++numAttackers)
{
std::uint32_t const numLegitimate = totalValidators - numAttackers;
Sim sim;
ConsensusParms const parms{};
SimDuration const delay = round<milliseconds>(0.2 * parms.ledgerGRANULARITY);
PeerGroup attackers = sim.createGroup(numAttackers);
PeerGroup legitimate = sim.createGroup(numLegitimate);
PeerGroup network = attackers + legitimate;
// Common UNL — all trust all
network.trust(network);
network.connect(network, delay);
// Round 1: establish common state
sim.run(1);
// === PARTITION: crash legitimate validators ===
legitimate.disconnect(network);
network.disconnect(legitimate);
// Stop legitimate peers from running
for (Peer* p : legitimate)
p->targetLedgers = p->completedLedgers;
// ATTACKER: bypass quorum by untrusting legitimate
attackers.untrust(legitimate);
// Attackers submit transactions
for (Peer* p : attackers)
p->submit(Tx{static_cast<std::uint32_t>(p->id) + 1000});
// Run partition phase (attackers produce chain)
sim.run(4);
// Record attacker FVL before recovery
Ledger::Seq const attackerFVLBeforeRecovery =
attackers[0]->fullyValidatedLedger.seq();
// === RECOVERY: reconnect ===
attackers.trust(legitimate);
network.connect(network, delay);
for (Peer* p : legitimate)
p->targetLedgers = p->completedLedgers + 15;
// Give plenty of recovery time
sim.run(10);
// Check if legitimate peers accepted attacker's chain
// by verifying their FVL ID matches an attacker's FVL ID
bool const legitimateAccepted = [&]() {
for (Peer* lp : legitimate)
{
if (lp->fullyValidatedLedger.seq() <= Ledger::Seq{1})
continue;
for (Peer* ap : attackers)
{
if (lp->fullyValidatedLedger.id() == ap->fullyValidatedLedger.id())
return true;
}
}
return false;
}();
bool const networkConverged = sim.synchronized();
std::size_t const branches = sim.branches();
float const attackerPct = 100.0f * numAttackers / totalValidators;
std::cout << " A=" << numAttackers << "/" << totalValidators << " ("
<< static_cast<int>(attackerPct) << "%)"
<< " AttackerFVL=" << attackerFVLBeforeRecovery
<< " LegitAccepted=" << std::boolalpha << legitimateAccepted
<< " Converged=" << networkConverged << " Branches=" << branches << "\n";
if (legitimateAccepted && attackThreshold == 0)
{
attackThreshold = numAttackers;
std::cout << " >>> ATTACK THRESHOLD: " << numAttackers << "/"
<< totalValidators << " (" << static_cast<int>(attackerPct) << "%)"
<< " <<<\n";
}
}
if (attackThreshold > 0)
{
float const thresholdPct = 100.0f * attackThreshold / totalValidators;
std::cout << " RESULT: Network " << totalValidators << " — attack succeeds at "
<< attackThreshold << " attackers (" << static_cast<int>(thresholdPct)
<< "%)\n";
}
else
{
std::cout << " RESULT: Network " << totalValidators
<< " — attack never succeeded\n";
}
// Byzantine safety: forcing the network to accept the attacker's
// chain requires a strict majority of validators (never a tolerated
// minority).
BEAST_EXPECT(attackThreshold == 0 || attackThreshold * 2 > totalValidators);
}
}
/**
* Test 5: Long partition — attacker builds deep chain
*
* What if the attacker runs for many rounds during the partition,
* building a much deeper chain? Does chain length affect acceptance?
*/
void
testLongPartitionDeepChain()
{
using namespace csf;
using namespace std::chrono;
testcase("partition: deep attacker chain");
for (int partitionRounds : {4, 10, 20, 40})
{
Sim sim;
ConsensusParms const parms{};
SimDuration const delay = round<milliseconds>(0.2 * parms.ledgerGRANULARITY);
PeerGroup attackers = sim.createGroup(4);
PeerGroup legitimate = sim.createGroup(3);
PeerGroup network = attackers + legitimate;
network.trust(network);
network.connect(network, delay);
sim.run(1);
// Partition
legitimate.disconnect(network);
network.disconnect(legitimate);
for (Peer* p : legitimate)
p->targetLedgers = p->completedLedgers;
attackers.untrust(legitimate);
for (Peer* p : attackers)
p->submit(Tx{static_cast<std::uint32_t>(p->id) + 500});
// Run for varying partition lengths
sim.run(partitionRounds);
Ledger::Seq const attackerDepth = attackers[0]->lastClosedLedger.seq();
Ledger::Seq const attackerFVL = attackers[0]->fullyValidatedLedger.seq();
// Recovery
attackers.trust(legitimate);
network.connect(network, delay);
for (Peer* p : legitimate)
p->targetLedgers = p->completedLedgers + 15;
sim.run(10);
bool const accepted = [&]() {
for (Peer* lp : legitimate)
{
if (lp->fullyValidatedLedger.seq() <= Ledger::Seq{1})
continue;
for (Peer* ap : attackers)
{
if (lp->fullyValidatedLedger.id() == ap->fullyValidatedLedger.id())
return true;
}
}
return false;
}();
std::cout << "PartitionRounds=" << partitionRounds << " AttackerDepth=" << attackerDepth
<< " AttackerFVL=" << attackerFVL << " LegitAccepted=" << std::boolalpha
<< accepted << " Branches=" << sim.branches()
<< " Synced=" << sim.synchronized() << "\n";
// With the attacker holding the majority (4 of 7), its chain is
// accepted regardless of how deep the partition ran — depth is not
// the deciding factor, validator share is.
BEAST_EXPECT(accepted);
}
}
/**
* Test 6: Staggered recovery
*
* Legitimate validators don't all come back at once.
* What if 2 come back first, then 2 more later?
* Does partial recovery change the outcome?
*/
void
testStaggeredRecovery()
{
using namespace csf;
using namespace std::chrono;
testcase("partition: staggered recovery");
Sim sim;
ConsensusParms const parms{};
SimDuration const delay = round<milliseconds>(0.2 * parms.ledgerGRANULARITY);
PeerGroup attackers = sim.createGroup(4);
PeerGroup legit1 = sim.createGroup(2); // first wave recovery
PeerGroup legit2 = sim.createGroup(1); // second wave recovery
PeerGroup legitimate = legit1 + legit2;
PeerGroup network = attackers + legitimate;
network.trust(network);
network.connect(network, delay);
sim.run(1);
std::cout << "Pre-partition: all synced at seq " << attackers[0]->lastClosedLedger.seq()
<< "\n";
// Partition all legitimate
legitimate.disconnect(network);
network.disconnect(legitimate);
for (Peer* p : legitimate)
p->targetLedgers = p->completedLedgers;
attackers.untrust(legitimate);
for (Peer* p : attackers)
p->submit(Tx{static_cast<std::uint32_t>(p->id) + 600});
sim.run(4);
std::cout << "After partition: attackers at seq " << attackers[0]->lastClosedLedger.seq()
<< " FVL=" << attackers[0]->fullyValidatedLedger.seq() << "\n";
// WAVE 1: partial recovery — only 2 legitimate peers return
attackers.trust(legit1);
legit1.trust(network);
network.connect(legit1, delay);
legit1.connect(network, delay);
for (Peer* p : legit1)
p->targetLedgers = p->completedLedgers + 10;
sim.run(6);
std::cout << "After wave 1 recovery (2 legit back):\n";
for (Peer* p : network)
{
std::cout << " Peer " << p->id << " closed=" << p->lastClosedLedger.seq()
<< " FVL=" << p->fullyValidatedLedger.seq() << "\n";
}
// WAVE 2: remaining legitimate peers return
attackers.trust(legit2);
legit2.trust(network);
network.connect(legit2, delay);
legit2.connect(network, delay);
for (Peer* p : legit2)
p->targetLedgers = p->completedLedgers + 10;
sim.run(6);
std::cout << "After wave 2 recovery (all back):\n";
std::cout << " Branches=" << sim.branches() << " Synced=" << sim.synchronized() << "\n";
for (Peer* p : network)
{
std::cout << " Peer " << p->id << " closed=" << p->lastClosedLedger.seq()
<< " FVL=" << p->fullyValidatedLedger.seq() << "\n";
}
// Once every legitimate validator is back, the network reconverges to a
// single chain.
BEAST_EXPECT(sim.synchronized());
}
/**
* Test 7: Attacker injects conflicting transactions
*
* During partition, attackers inject transactions into their chain.
* After recovery, do legitimate peers end up with the attacker's
* injected transactions in their fully validated chain?
* This directly tests the amendment bypass scenario.
*/
void
testConflictingTxPersistence()
{
using namespace csf;
using namespace std::chrono;
testcase("partition: conflicting tx persistence after recovery");
Sim sim;
ConsensusParms const parms{};
SimDuration const delay = round<milliseconds>(0.2 * parms.ledgerGRANULARITY);
// 4 attackers, 3 legitimate (attack SHOULD succeed)
PeerGroup attackers = sim.createGroup(4);
PeerGroup legitimate = sim.createGroup(3);
PeerGroup network = attackers + legitimate;
network.trust(network);
network.connect(network, delay);
sim.run(1);
Ledger::Seq const prePartitionSeq = attackers[0]->lastClosedLedger.seq();
// Partition
legitimate.disconnect(network);
network.disconnect(legitimate);
for (Peer* p : legitimate)
p->targetLedgers = p->completedLedgers;
attackers.untrust(legitimate);
// Inject a "malicious" transaction ONLY on attacker nodes
// This simulates an EnableAmendment pseudo-tx bypass
Tx const maliciousTx{9999};
for (Peer* p : attackers)
p->txInjections.emplace(prePartitionSeq, maliciousTx);
sim.run(4);
// Record the attacker ledger that contains the injected tx
Ledger::ID const attackerLedgerWithTx = attackers[0]->lastClosedLedger.id();
Ledger::Seq const attackerSeqWithTx = attackers[0]->lastClosedLedger.seq();
std::cout << "Attacker injected tx at seq " << attackerSeqWithTx << "\n";
// Recovery
attackers.trust(legitimate);
network.connect(network, delay);
for (Peer* p : legitimate)
p->targetLedgers = p->completedLedgers + 15;
sim.run(10);
// Check: does any legitimate peer have the attacker's ledger
// (containing the injected tx) in their chain?
bool const legitimateHasAttackerLedger = [&]() {
for (Peer* lp : legitimate)
{
// Check if the attacker ledger is in the legitimate
// peer's ledger cache
auto it = lp->ledgers.find(attackerLedgerWithTx);
if (it != lp->ledgers.end())
return true;
}
return false;
}();
bool const legitimateFVLDescendsFromAttacker = [&]() {
for (Peer* lp : legitimate)
{
// Check if legitimate peer's FVL is built on the
// attacker's chain
if (lp->fullyValidatedLedger.seq() >= attackerSeqWithTx)
{
// Check if the attacker ledger is an ancestor
auto it = lp->ledgers.find(attackerLedgerWithTx);
if (it != lp->ledgers.end() && lp->fullyValidatedLedger.isAncestor(it->second))
return true;
}
}
return false;
}();
std::cout << "After recovery:\n"
<< " LegitHasAttackerLedger=" << std::boolalpha << legitimateHasAttackerLedger
<< "\n"
<< " LegitFVLDescendsFromAttacker=" << legitimateFVLDescendsFromAttacker << "\n"
<< " Branches=" << sim.branches() << " Synced=" << sim.synchronized() << "\n";
for (Peer* p : network)
{
std::cout << " Peer " << p->id << " closed=" << p->lastClosedLedger.seq()
<< " FVL=" << p->fullyValidatedLedger.seq() << "\n";
}
// The amendment-bypass safety property: a minority attacker's injected
// transactions never enter the legitimate fully-validated chain.
BEAST_EXPECT(!legitimateFVLDescendsFromAttacker);
}
/**
* Test 8: Gradual nUNL-style takeover
*
* Instead of crashing all legitimate validators at once,
* the attacker gradually removes them from trust (simulating
* nUNL additions), then produces a chain.
*/
void
testGradualTakeover()
{
using namespace csf;
using namespace std::chrono;
testcase("partition: gradual nUNL-style takeover");
Sim sim;
ConsensusParms const parms{};
SimDuration const delay = round<milliseconds>(0.2 * parms.ledgerGRANULARITY);
// 4 attackers, 3 legitimate
PeerGroup attackers = sim.createGroup(4);
PeerGroup legit1 = sim.createGroup(1);
PeerGroup legit2 = sim.createGroup(1);
PeerGroup legit3 = sim.createGroup(1);
PeerGroup legitimate = legit1 + legit2 + legit3;
PeerGroup network = attackers + legitimate;
network.trust(network);
network.connect(network, delay);
sim.run(1);
std::cout << "All synced at seq " << network[0]->lastClosedLedger.seq() << "\n";
// Phase 1: crash legit1, attackers remove from trust
legit1.disconnect(network);
network.disconnect(legit1);
for (Peer* p : legit1)
p->targetLedgers = p->completedLedgers;
attackers.untrust(legit1);
sim.run(2);
std::cout << "Phase 1 (1 crashed): network at seq " << attackers[0]->lastClosedLedger.seq()
<< " FVL=" << attackers[0]->fullyValidatedLedger.seq()
<< " branches=" << sim.branches() << "\n";
// Phase 2: crash legit2
legit2.disconnect(network);
network.disconnect(legit2);
for (Peer* p : legit2)
p->targetLedgers = p->completedLedgers;
attackers.untrust(legit2);
sim.run(2);
std::cout << "Phase 2 (2 crashed): network at seq " << attackers[0]->lastClosedLedger.seq()
<< " FVL=" << attackers[0]->fullyValidatedLedger.seq()
<< " branches=" << sim.branches() << "\n";
// Phase 3: crash legit3 — attackers are now alone
legit3.disconnect(network);
network.disconnect(legit3);
for (Peer* p : legit3)
p->targetLedgers = p->completedLedgers;
attackers.untrust(legit3);
// Now attackers trust only themselves (4/4 quorum = 4)
for (Peer* p : attackers)
p->submit(Tx{static_cast<std::uint32_t>(p->id) + 700});
sim.run(4);
std::cout << "Phase 3 (all crashed): attackers at seq "
<< attackers[0]->lastClosedLedger.seq()
<< " FVL=" << attackers[0]->fullyValidatedLedger.seq() << "\n";
// Recovery: all legitimate come back at once
attackers.trust(legitimate);
network.connect(network, delay);
for (Peer* p : legitimate)
p->targetLedgers = p->completedLedgers + 15;
sim.run(10);
bool const accepted = [&]() {
for (Peer* lp : legitimate)
{
if (lp->fullyValidatedLedger.seq() <= Ledger::Seq{1})
continue;
for (Peer* ap : attackers)
{
if (lp->fullyValidatedLedger.id() == ap->fullyValidatedLedger.id())
return true;
}
}
return false;
}();
std::cout << "After recovery:\n"
<< " LegitAccepted=" << std::boolalpha << accepted
<< " Branches=" << sim.branches() << " Synced=" << sim.synchronized() << "\n";
// A gradual nUNL-style takeover that leaves the attacker holding the
// majority (4 of 7) does flip the network onto the attacker's chain —
// the reason runtime nUNL/quorum changes need coordination guardrails.
BEAST_EXPECT(accepted);
for (Peer* p : network)
{
std::cout << " Peer " << p->id << " closed=" << p->lastClosedLedger.seq()
<< " FVL=" << p->fullyValidatedLedger.seq() << "\n";
}
}
void
run() override
{
testPartitionNoQuorumBypass();
testPartitionWithQuorumBypass();
testTxInjectionDuringPartition();
testAttackSurfaceSweep();
testLongPartitionDeepChain();
testStaggeredRecovery();
testConflictingTxPersistence();
testGradualTakeover();
}
};
BEAST_DEFINE_TESTSUITE_MANUAL(ByzantinePartitionRecovery, consensus, xrpl);
} // namespace xrpl::test

View File

@@ -19,6 +19,7 @@
#include <xrpld/app/main/LoadManager.h>
#include <xrpld/app/main/NodeIdentity.h>
#include <xrpld/app/main/NodeStoreScheduler.h>
#include <xrpld/app/misc/DatagramMonitor.h>
#include <xrpld/app/misc/SHAMapStore.h>
#include <xrpld/app/misc/TxQ.h>
#include <xrpld/app/misc/ValidatorKeys.h>
@@ -220,6 +221,7 @@ public:
std::unique_ptr<JobQueue> jobQueue_;
NodeStoreScheduler nodeStoreScheduler_;
std::unique_ptr<SHAMapStore> shaMapStore_;
std::unique_ptr<DatagramMonitor> datagramMonitor_;
PendingSaves pendingSaves_;
std::optional<OpenLedger> openLedger_;
@@ -1503,6 +1505,14 @@ ApplicationImp::start(bool withTimers)
ledgerCleaner_->start();
perfLog_->start();
// Datagram monitor: UDP node-stats exporter (XDGM). Off in standalone or
// when [datagram_monitor] has no endpoints.
if (!config_->standalone() && !config_->DATAGRAM_MONITOR.empty())
{
datagramMonitor_ = std::make_unique<DatagramMonitor>(*this);
datagramMonitor_->start();
}
}
void

View File

@@ -0,0 +1,882 @@
//
#ifndef RIPPLE_APP_MAIN_DATAGRAMMONITOR_H_INCLUDED
#define RIPPLE_APP_MAIN_DATAGRAMMONITOR_H_INCLUDED
#include <xrpld/app/ledger/AcceptedLedger.h>
#include <xrpld/app/ledger/InboundLedgers.h>
#include <xrpld/app/ledger/LedgerMaster.h>
#include <xrpld/app/main/Application.h>
#include <xrpld/app/misc/ValidatorList.h>
#include <xrpld/app/rdb/backend/SQLiteDatabase.h>
#include <xrpld/overlay/Overlay.h>
#include <xrpl/basics/UptimeClock.h>
#include <xrpl/basics/mulDiv.h>
#include <xrpl/beast/utility/Journal.h>
#include <xrpl/ledger/CachedSLEs.h>
#include <xrpl/nodestore/Database.h>
#include <xrpl/protocol/BuildInfo.h>
#include <xrpl/protocol/ErrorCodes.h>
#include <xrpl/protocol/jss.h>
#include <xrpl/server/LoadFeeTrack.h>
#include <xrpl/server/NetworkOPs.h>
#include <arpa/inet.h>
#include <sys/resource.h>
#include <sys/socket.h>
#include <netdb.h>
#include <array>
#include <atomic>
#include <chrono>
#include <cstring>
#include <fstream>
#include <sstream>
#include <string>
#if defined(__linux__)
#include <sys/statvfs.h>
#include <sys/sysinfo.h>
#elif defined(__APPLE__)
#include <mach/host_info.h>
#include <mach/mach.h>
#include <net/if.h>
#include <net/if_dl.h>
#include <sys/mount.h>
#include <sys/sysctl.h>
#include <sys/types.h>
#include <ifaddrs.h>
#endif
#include <thread>
#include <vector>
namespace xrpl {
// Magic number for server info packets: 'XDGM' (le) Xahau DataGram Monitor
constexpr uint32_t SERVER_INFO_MAGIC = 0x4D474458;
constexpr uint32_t SERVER_INFO_VERSION = 1;
// Warning flag bits
constexpr uint32_t WARNING_AMENDMENT_BLOCKED = 1 << 0;
constexpr uint32_t WARNING_UNL_BLOCKED = 1 << 1;
constexpr uint32_t WARNING_AMENDMENT_WARNED = 1 << 2;
constexpr uint32_t WARNING_NOT_SYNCED = 1 << 3;
// Time window statistics for rates
struct [[gnu::packed]] MetricRates
{
double rate_1m; // Average rate over last minute
double rate_5m; // Average rate over last 5 minutes
double rate_1h; // Average rate over last hour
double rate_24h; // Average rate over last 24 hours
};
struct AllRates
{
MetricRates network_in;
MetricRates network_out;
MetricRates disk_read;
MetricRates disk_write;
};
// Structure to represent a ledger sequence range
struct [[gnu::packed]] LgrRange
{
uint32_t start;
uint32_t end;
};
// Map is returned separately since variable-length data
// shouldn't be included in network structures
using ObjectCountMap = std::vector<std::pair<std::basic_string<char>, int>>;
struct [[gnu::packed]] DebugCounters
{
// Database metrics
std::uint64_t dbKBTotal{0};
std::uint64_t dbKBLedger{0};
std::uint64_t dbKBTransaction{0};
std::uint64_t localTxCount{0};
// Basic metrics
std::uint32_t writeLoad{0};
std::int32_t historicalPerMinute{0};
// Cache metrics
std::uint32_t sleHitRate{0}; // Stored as fixed point, multiplied by 1000
std::uint32_t ledgerHitRate{0}; // Stored as fixed point, multiplied by 1000
std::uint32_t alSize{0};
std::uint32_t alHitRate{0}; // Stored as fixed point, multiplied by 1000
std::int32_t fullbelowSize{0};
std::uint32_t treenodeCacheSize{0};
std::uint32_t treenodeTrackSize{0};
// Node store metrics
std::uint64_t nodeWriteCount{0};
std::uint64_t nodeWriteSize{0};
std::uint64_t nodeFetchCount{0};
std::uint64_t nodeFetchHitCount{0};
std::uint64_t nodeFetchSize{0};
};
// Core server metrics in the fixed header
struct [[gnu::packed]] ServerInfoHeader
{
// Fixed header fields come first
uint32_t magic; // Magic number to identify packet type
uint32_t version; // Protocol version number
uint32_t network_id; // Network ID from config
uint32_t server_state; // Operating mode as enum
uint32_t peer_count; // Number of connected peers
uint32_t node_size; // Size category (0=tiny through 4=huge)
uint32_t cpu_cores; // CPU core count
uint32_t ledger_range_count; // Number of range entries
uint32_t warning_flags; // Warning flags (reduced size)
uint32_t padding_1; // padding for alignment
// 64-bit metrics
uint64_t timestamp; // System time in microseconds
uint64_t uptime; // Server uptime in seconds
uint64_t io_latency_us; // IO latency in microseconds
uint64_t validation_quorum; // Validation quorum count
uint64_t fetch_pack_size; // Size of fetch pack cache
uint64_t proposer_count; // Number of proposers in last close
uint64_t converge_time_ms; // Last convergence time in ms
uint64_t load_factor; // Load factor (scaled by 1M)
uint64_t load_base; // Load base value
uint64_t reserve_base; // Reserve base amount
uint64_t reserve_inc; // Reserve increment amount
uint64_t ledger_seq; // Latest ledger sequence
// Fixed-size byte arrays
uint8_t ledger_hash[32]; // Latest ledger hash
uint8_t node_public_key[33]; // Node's public key
uint8_t padding2[7]; // Padding to maintain 8-byte alignment
uint8_t version_string[32];
// System metrics
uint64_t process_memory_pages; // Process memory usage in bytes
uint64_t system_memory_total; // Total system memory in bytes
uint64_t system_memory_free; // Free system memory in bytes
uint64_t system_memory_used; // Used system memory in bytes
uint64_t system_disk_total; // Total disk space in bytes
uint64_t system_disk_free; // Free disk space in bytes
uint64_t system_disk_used; // Used disk space in bytes
uint64_t io_wait_time; // IO wait time in milliseconds
double load_avg_1min; // 1 minute load average
double load_avg_5min; // 5 minute load average
double load_avg_15min; // 15 minute load average
// State transition metrics
uint64_t state_transitions[5]; // Count for each operating mode
uint64_t state_durations[5]; // Duration in each mode
uint64_t initial_sync_us; // Initial sync duration
// Network and disk rates remain unchanged
struct
{
MetricRates network_in;
MetricRates network_out;
MetricRates disk_read;
MetricRates disk_write;
} rates;
DebugCounters dbg_counters;
};
// System metrics collected for rate calculations
struct SystemMetrics
{
uint64_t timestamp; // When metrics were collected
uint64_t network_bytes_in; // Current total bytes in
uint64_t network_bytes_out; // Current total bytes out
uint64_t disk_bytes_read; // Current total bytes read
uint64_t disk_bytes_written; // Current total bytes written
};
class MetricsTracker
{
private:
static constexpr size_t SAMPLES_1M = 60; // 1 sample/second for 1 minute
static constexpr size_t SAMPLES_5M = 300; // 1 sample/second for 5 minutes
static constexpr size_t SAMPLES_1H = 3600; // 1 sample/second for 1 hour
static constexpr size_t SAMPLES_24H = 1440; // 1 sample/minute for 24 hours
std::vector<SystemMetrics> samples_1m{SAMPLES_1M};
std::vector<SystemMetrics> samples_5m{SAMPLES_5M};
std::vector<SystemMetrics> samples_1h{SAMPLES_1H};
std::vector<SystemMetrics> samples_24h{SAMPLES_24H};
size_t index_1m{0}, index_5m{0}, index_1h{0}, index_24h{0};
std::chrono::system_clock::time_point last_24h_sample{};
double
calculateRate(
SystemMetrics const& current,
std::vector<SystemMetrics> const& samples,
size_t current_index,
size_t max_samples,
bool is_24h_window,
std::function<uint64_t(SystemMetrics const&)> metric_getter)
{
// If we don't have at least 2 samples, the rate is 0
if (current_index < 2)
{
return 0.0;
}
// Calculate time window based on the window type
uint64_t expected_window_micros;
if (is_24h_window)
{
expected_window_micros =
24ULL * 60ULL * 60ULL * 1000000ULL; // 24 hours in microseconds
}
else
{
expected_window_micros =
max_samples * 1000000ULL; // window in seconds * 1,000,000 for microseconds
}
// For any window where we don't have full data, we should scale the
// rate based on the actual time we have data for
uint64_t actual_window_micros = current.timestamp - samples[0].timestamp;
double window_scale =
std::min(1.0, static_cast<double>(actual_window_micros) / expected_window_micros);
// Get the oldest valid sample
size_t oldest_index =
(current_index >= max_samples) ? ((current_index + 1) % max_samples) : 0;
auto const& oldest = samples[oldest_index];
double elapsed = actual_window_micros / 1000000.0; // Convert microseconds to seconds
// Ensure we have a meaningful time difference
if (elapsed < 0.001)
{ // Less than 1ms difference
return 0.0;
}
uint64_t current_value = metric_getter(current);
uint64_t oldest_value = metric_getter(oldest);
// Handle counter wraparound
uint64_t diff = (current_value >= oldest_value)
? (current_value - oldest_value)
: (std::numeric_limits<uint64_t>::max() - oldest_value + current_value + 1);
// Calculate the rate and scale it based on our window coverage
return (static_cast<double>(diff) / elapsed) * window_scale;
}
MetricRates
calculateMetricRates(
SystemMetrics const& current,
std::function<uint64_t(SystemMetrics const&)> metric_getter)
{
MetricRates rates;
rates.rate_1m =
calculateRate(current, samples_1m, index_1m, SAMPLES_1M, false, metric_getter);
rates.rate_5m =
calculateRate(current, samples_5m, index_5m, SAMPLES_5M, false, metric_getter);
rates.rate_1h =
calculateRate(current, samples_1h, index_1h, SAMPLES_1H, false, metric_getter);
rates.rate_24h =
calculateRate(current, samples_24h, index_24h, SAMPLES_24H, true, metric_getter);
return rates;
}
public:
void
addSample(SystemMetrics const& metrics)
{
auto now = std::chrono::system_clock::now();
// Update 1-minute window (every second)
samples_1m[index_1m++ % SAMPLES_1M] = metrics;
// Update 5-minute window (every second)
samples_5m[index_5m++ % SAMPLES_5M] = metrics;
// Update 1-hour window (every second)
samples_1h[index_1h++ % SAMPLES_1H] = metrics;
// Update 24-hour window (every minute)
if (last_24h_sample + std::chrono::minutes(1) <= now)
{
samples_24h[index_24h++ % SAMPLES_24H] = metrics;
last_24h_sample = now;
}
}
AllRates
getRates(SystemMetrics const& current)
{
AllRates rates;
rates.network_in = calculateMetricRates(
current, [](SystemMetrics const& m) { return m.network_bytes_in; });
rates.network_out = calculateMetricRates(
current, [](SystemMetrics const& m) { return m.network_bytes_out; });
rates.disk_read =
calculateMetricRates(current, [](SystemMetrics const& m) { return m.disk_bytes_read; });
rates.disk_write = calculateMetricRates(
current, [](SystemMetrics const& m) { return m.disk_bytes_written; });
return rates;
}
};
class DatagramMonitor
{
private:
Application& app_;
beast::Journal j_;
std::atomic<bool> running_{false};
std::thread monitor_thread_;
MetricsTracker metrics_tracker_;
struct EndpointInfo
{
std::string ip;
uint16_t port;
bool is_ipv6;
};
EndpointInfo
parseEndpoint(std::string const& endpoint)
{
auto space_pos = endpoint.find(' ');
if (space_pos == std::string::npos)
throw std::runtime_error("Invalid endpoint format");
EndpointInfo info;
info.ip = endpoint.substr(0, space_pos);
info.port = std::stoi(endpoint.substr(space_pos + 1));
info.is_ipv6 = info.ip.find(':') != std::string::npos;
return info;
}
int
createSocket(EndpointInfo const& endpoint)
{
int sock = socket(endpoint.is_ipv6 ? AF_INET6 : AF_INET, SOCK_DGRAM, 0);
if (sock < 0)
throw std::runtime_error("Failed to create socket");
return sock;
}
void
sendPacket(int sock, EndpointInfo const& endpoint, std::vector<uint8_t> const& buffer)
{
struct sockaddr_storage addr;
socklen_t addr_len;
if (endpoint.is_ipv6)
{
struct sockaddr_in6* addr6 = reinterpret_cast<struct sockaddr_in6*>(&addr);
addr6->sin6_family = AF_INET6;
addr6->sin6_port = htons(endpoint.port);
inet_pton(AF_INET6, endpoint.ip.c_str(), &addr6->sin6_addr);
addr_len = sizeof(struct sockaddr_in6);
}
else
{
struct sockaddr_in* addr4 = reinterpret_cast<struct sockaddr_in*>(&addr);
addr4->sin_family = AF_INET;
addr4->sin_port = htons(endpoint.port);
inet_pton(AF_INET, endpoint.ip.c_str(), &addr4->sin_addr);
addr_len = sizeof(struct sockaddr_in);
}
sendto(
sock,
buffer.data(),
buffer.size(),
0,
reinterpret_cast<struct sockaddr*>(&addr),
addr_len);
}
// Returns both the counters and object count map separately
std::pair<DebugCounters, ObjectCountMap>
getDebugCounters()
{
DebugCounters counters;
ObjectCountMap objectCounts = CountedObjects::getInstance().getCounts(1);
// Database metrics if applicable
if (app_.config().useTxTables())
{
auto const db = dynamic_cast<SQLiteDatabase*>(&app_.getRelationalDatabase());
if (!db)
Throw<std::runtime_error>("Failed to get relational database");
if (auto dbKB = db->getKBUsedAll())
counters.dbKBTotal = dbKB;
if (auto dbKB = db->getKBUsedLedger())
counters.dbKBLedger = dbKB;
if (auto dbKB = db->getKBUsedTransaction())
counters.dbKBTransaction = dbKB;
if (auto count = app_.getOPs().getLocalTxCount())
counters.localTxCount = count;
}
// Basic metrics
counters.writeLoad = app_.getNodeStore().getWriteLoad();
counters.historicalPerMinute =
static_cast<std::int32_t>(app_.getInboundLedgers().fetchRate());
// Cache metrics - convert floating point rates to fixed point
counters.sleHitRate = 0; // TODO: SLE cache hit-rate accessor absent on this fork
counters.ledgerHitRate =
static_cast<std::uint32_t>(app_.getLedgerMaster().getCacheHitRate() * 1000);
counters.alSize = app_.getAcceptedLedgerCache().size();
counters.alHitRate =
static_cast<std::uint32_t>(app_.getAcceptedLedgerCache().getHitRate() * 1000);
counters.fullbelowSize =
static_cast<std::int32_t>(app_.getNodeFamily().getFullBelowCache()->size());
counters.treenodeCacheSize = app_.getNodeFamily().getTreeNodeCache()->getCacheSize();
counters.treenodeTrackSize = app_.getNodeFamily().getTreeNodeCache()->getTrackSize();
// Get regular node store metrics
counters.nodeWriteCount = app_.getNodeStore().getStoreCount();
counters.nodeWriteSize = app_.getNodeStore().getStoreSize();
counters.nodeFetchCount = app_.getNodeStore().getFetchTotalCount();
counters.nodeFetchHitCount = app_.getNodeStore().getFetchHitCount();
counters.nodeFetchSize = app_.getNodeStore().getFetchSize();
return {counters, objectCounts};
}
uint32_t
getPhysicalCPUCount()
{
static uint32_t count = 0;
if (count > 0)
return count;
#if defined(__linux__)
try
{
std::ifstream cpuinfo("/proc/cpuinfo");
if (!cpuinfo)
{
JLOG(j_.error()) << "Unable to open file: /proc/cpuinfo";
return count;
}
std::string line;
std::set<std::string> physical_ids;
std::string current_physical_id;
while (std::getline(cpuinfo, line))
{
if (line.find("core id") != std::string::npos)
{
current_physical_id = line.substr(line.find(":") + 1);
// Trim whitespace
current_physical_id.erase(0, current_physical_id.find_first_not_of(" \t"));
current_physical_id.erase(current_physical_id.find_last_not_of(" \t") + 1);
physical_ids.insert(current_physical_id);
}
}
count = physical_ids.size();
}
catch (std::exception const& e)
{
JLOG(j_.error()) << "Error getting CPU count: " << e.what();
}
// Return at least 1 if we couldn't determine the count
return count > 0 ? count : (count = 1);
#elif defined(__APPLE__)
int value = 0;
size_t size = sizeof(value);
if (sysctlbyname("hw.physicalcpu", &value, &size, NULL, 0) == 0)
count = value;
return count > 0 ? count : (count = 1);
#endif
}
SystemMetrics
collectSystemMetrics()
{
SystemMetrics metrics{};
metrics.timestamp = std::chrono::duration_cast<std::chrono::microseconds>(
std::chrono::system_clock::now().time_since_epoch())
.count();
#if defined(__linux__)
// Network stats collection
try
{
std::ifstream net_file("/proc/net/dev");
if (!net_file)
{
JLOG(j_.error()) << "Unable to open file /proc/net/dev";
return metrics;
}
std::string line;
uint64_t total_bytes_in = 0, total_bytes_out = 0;
// Skip header lines
std::getline(net_file, line); // Inter-| Receive...
std::getline(net_file, line); // face |bytes...
while (std::getline(net_file, line))
{
if (line.find(':') != std::string::npos)
{
std::string interface = line.substr(0, line.find(':'));
interface = interface.substr(interface.find_first_not_of(" \t"));
interface = interface.substr(0, interface.find_last_not_of(" \t") + 1);
// Skip loopback interface
if (interface == "lo")
continue;
uint64_t bytes_in, bytes_out;
std::istringstream iss(line.substr(line.find(':') + 1));
iss >> bytes_in; // First field after : is bytes_in
for (int i = 0; i < 8; ++i)
iss >> std::ws; // Skip 8 fields
iss >> bytes_out; // 9th field is bytes_out
total_bytes_in += bytes_in;
total_bytes_out += bytes_out;
}
}
metrics.network_bytes_in = total_bytes_in;
metrics.network_bytes_out = total_bytes_out;
}
catch (std::exception const& e)
{
JLOG(j_.error()) << "Error collecting network stats: " << e.what();
}
// Disk stats collection
try
{
std::ifstream disk_file("/proc/diskstats");
if (!disk_file)
{
JLOG(j_.error()) << "Unable to open file: /proc/diskstats";
return metrics;
}
std::string line;
uint64_t total_bytes_read = 0, total_bytes_written = 0;
while (std::getline(disk_file, line))
{
unsigned int major, minor;
char dev_name[32];
uint64_t reads, read_sectors, writes, write_sectors;
if (sscanf(
line.c_str(),
"%u %u %31s %lu %*u %lu %*u %lu %*u %lu",
&major,
&minor,
dev_name,
&reads,
&read_sectors,
&writes,
&write_sectors) == 7)
{
// Only process physical devices
std::string device_name(dev_name);
if (device_name.substr(0, 3) == "dm-" || device_name.substr(0, 4) == "loop" ||
device_name.substr(0, 3) == "ram")
{
continue;
}
// Skip partitions (usually have a number at the end)
if (std::isdigit(device_name.back()))
{
continue;
}
uint64_t bytes_read = read_sectors * 512;
uint64_t bytes_written = write_sectors * 512;
total_bytes_read += bytes_read;
total_bytes_written += bytes_written;
}
}
metrics.disk_bytes_read = total_bytes_read;
metrics.disk_bytes_written = total_bytes_written;
}
catch (std::exception const& e)
{
JLOG(j_.error()) << "Error collecting disk stats: " << e.what();
}
#elif defined(__APPLE__)
// Network stats collection
try
{
struct ifaddrs* ifap;
if (getifaddrs(&ifap) == 0)
{
uint64_t total_bytes_in = 0, total_bytes_out = 0;
for (struct ifaddrs* ifa = ifap; ifa; ifa = ifa->ifa_next)
{
if (ifa->ifa_addr != NULL && ifa->ifa_addr->sa_family == AF_LINK)
{
struct if_data* ifd = (struct if_data*)ifa->ifa_data;
if (ifd != NULL)
{
// Skip loopback interface
if (strcmp(ifa->ifa_name, "lo0") == 0)
continue;
total_bytes_in += ifd->ifi_ibytes;
total_bytes_out += ifd->ifi_obytes;
}
}
}
freeifaddrs(ifap);
metrics.network_bytes_in = total_bytes_in;
metrics.network_bytes_out = total_bytes_out;
}
}
catch (std::exception const& e)
{
JLOG(j_.error()) << "Error collecting network stats: " << e.what();
}
// Disk stats collection
// Disk IO stats are not easily accessible in macOS.
// We'll set these values to zero for now.
metrics.disk_bytes_read = 0;
metrics.disk_bytes_written = 0;
#endif
return metrics;
}
std::vector<uint8_t>
generateServerInfo()
{
auto& ops = app_.getOPs();
auto& ledgerMaster = app_.getLedgerMaster();
auto currentMetrics = collectSystemMetrics();
metrics_tracker_.addSample(currentMetrics);
// Slimmed for this fork (3.2.0-b0): ledger ranges, DB debug-counters and
// the object-count map are omitted (divergent accessors). The packet is
// just the fixed header with core node + OS metrics.
std::vector<uint8_t> buffer(sizeof(ServerInfoHeader));
auto* header = reinterpret_cast<ServerInfoHeader*>(buffer.data());
memset(header, 0, sizeof(ServerInfoHeader));
header->magic = SERVER_INFO_MAGIC;
header->version = SERVER_INFO_VERSION;
header->network_id = app_.config().networkId;
header->timestamp = std::chrono::duration_cast<std::chrono::microseconds>(
std::chrono::system_clock::now().time_since_epoch())
.count();
header->uptime = UptimeClock::now().time_since_epoch().count();
header->io_latency_us = app_.getIOLatency().count();
header->validation_quorum = app_.getValidators().quorum();
header->server_state = static_cast<std::uint32_t>(ops.getOperatingMode());
header->peer_count = app_.getOverlay().size();
header->node_size = app_.config().nodeSize;
auto const [counters, mode, start, initialSync] = ops.getStateAccountingData();
for (size_t i = 0; i < 5; ++i)
{
header->state_transitions[i] = counters[i].transitions;
header->state_durations[i] = counters[i].dur.count();
}
header->initial_sync_us = initialSync;
if (ops.isAmendmentBlocked())
header->warning_flags |= WARNING_AMENDMENT_BLOCKED;
if (ops.isUNLBlocked())
header->warning_flags |= WARNING_UNL_BLOCKED;
if (ops.isAmendmentWarned())
header->warning_flags |= WARNING_AMENDMENT_WARNED;
if (ops.getOperatingMode() != OperatingMode::FULL)
header->warning_flags |= WARNING_NOT_SYNCED;
header->proposer_count = ops.getPrevProposers();
header->converge_time_ms = ops.getPrevRoundTime().count();
auto const fp = ledgerMaster.getFetchPackCacheSize();
if (fp != 0)
header->fetch_pack_size = fp;
// Load factor (server only; fee-escalation term omitted on this fork).
header->load_factor = static_cast<std::uint64_t>(app_.getFeeTrack().getLoadFactor());
header->load_base = app_.getFeeTrack().getLoadBase();
#if defined(__linux__)
// Get system info using sysinfo
struct sysinfo si;
if (sysinfo(&si) == 0)
{
header->system_memory_total = si.totalram * si.mem_unit;
header->system_memory_free = si.freeram * si.mem_unit;
header->system_memory_used = header->system_memory_total - header->system_memory_free;
header->load_avg_1min = si.loads[0] / (float)(1 << SI_LOAD_SHIFT);
header->load_avg_5min = si.loads[1] / (float)(1 << SI_LOAD_SHIFT);
header->load_avg_15min = si.loads[2] / (float)(1 << SI_LOAD_SHIFT);
}
#elif defined(__APPLE__)
// Get total physical memory
int64_t physical_memory;
size_t length = sizeof(physical_memory);
if (sysctlbyname("hw.memsize", &physical_memory, &length, NULL, 0) == 0)
{
header->system_memory_total = physical_memory;
}
// Get free and used memory
vm_statistics_data_t vm_stats;
mach_msg_type_number_t count = HOST_VM_INFO_COUNT;
if (host_statistics(mach_host_self(), HOST_VM_INFO, (host_info_t)&vm_stats, &count) ==
KERN_SUCCESS)
{
uint64_t page_size;
length = sizeof(page_size);
sysctlbyname("hw.pagesize", &page_size, &length, NULL, 0);
header->system_memory_free = (uint64_t)vm_stats.free_count * page_size;
header->system_memory_used = header->system_memory_total - header->system_memory_free;
}
// Get load averages
double loadavg[3];
if (getloadavg(loadavg, 3) == 3)
{
header->load_avg_1min = loadavg[0];
header->load_avg_5min = loadavg[1];
header->load_avg_15min = loadavg[2];
}
#endif
// Get process memory usage
struct rusage usage;
getrusage(RUSAGE_SELF, &usage);
header->process_memory_pages = usage.ru_maxrss;
// Get disk usage
#if defined(__linux__)
struct statvfs fs;
if (statvfs("/", &fs) == 0)
{
header->system_disk_total = fs.f_blocks * fs.f_frsize;
header->system_disk_free = fs.f_bfree * fs.f_frsize;
header->system_disk_used = header->system_disk_total - header->system_disk_free;
}
#elif defined(__APPLE__)
struct statfs fs;
if (statfs("/", &fs) == 0)
{
header->system_disk_total = fs.f_blocks * fs.f_bsize;
header->system_disk_free = fs.f_bfree * fs.f_bsize;
header->system_disk_used = header->system_disk_total - header->system_disk_free;
}
#endif
// Get CPU core count
header->cpu_cores = getPhysicalCPUCount();
// Get rate statistics
auto rates = metrics_tracker_.getRates(currentMetrics);
header->rates.network_in = rates.network_in;
header->rates.network_out = rates.network_out;
header->rates.disk_read = rates.disk_read;
header->rates.disk_write = rates.disk_write;
// Ledger height + hash via stable accessors (this fork's Ledger lacks
// info()). The hash lets the collector detect a fork: divergent
// ledger_hash across nodes at the same ledger_seq.
std::uint32_t const validSeq = ledgerMaster.getValidLedgerIndex();
header->ledger_seq = validSeq;
uint256 const validHash = ledgerMaster.getHashBySeq(validSeq);
std::memcpy(header->ledger_hash, validHash.data(), 32);
header->reserve_base = app_.config().fees.accountReserve.drops();
header->reserve_inc = app_.config().fees.ownerReserve.drops();
// Node public key + version string.
auto const& nodeKey = app_.nodeIdentity().first;
std::memcpy(header->node_public_key, nodeKey.data(), 33);
memset(&header->version_string, 0, 32);
memcpy(
&header->version_string,
BuildInfo::getVersionString().c_str(),
BuildInfo::getVersionString().size() > 32 ? 32 : BuildInfo::getVersionString().size());
header->ledger_range_count = 0;
return buffer;
}
void
monitorThread()
{
std::vector<std::pair<EndpointInfo, int>> endpoints;
for (auto const& epStr : app_.config().DATAGRAM_MONITOR)
{
auto endpoint = parseEndpoint(epStr);
endpoints.push_back(std::make_pair(endpoint, createSocket(endpoint)));
}
while (running_)
{
try
{
auto info = generateServerInfo();
for (auto const& ep : endpoints)
{
sendPacket(ep.second, ep.first, info);
}
std::this_thread::sleep_for(std::chrono::seconds(1));
}
catch (std::exception const& e)
{
// Log error but continue monitoring
JLOG(j_.error()) << "Server info monitor error: " << e.what();
}
}
for (auto const& ep : endpoints)
{
close(ep.second);
}
}
public:
DatagramMonitor(Application& app) : app_(app), j_(beast::Journal::getNullSink())
{
}
void
start()
{
if (!running_.exchange(true))
{
monitor_thread_ = std::thread(&DatagramMonitor::monitorThread, this);
}
}
void
stop()
{
if (running_.exchange(false))
{
if (monitor_thread_.joinable())
monitor_thread_.join();
}
}
~DatagramMonitor()
{
stop();
}
};
} // namespace xrpl
#endif

View File

@@ -346,6 +346,9 @@ public:
OperatingMode
getOperatingMode() const override;
StateAccountingData
getStateAccountingData() override;
std::string
strOperatingMode(OperatingMode const mode, bool const admin) const override;
@@ -508,8 +511,19 @@ public:
void
consensusViewChange() override;
void
setStall(std::chrono::milliseconds duration) override;
bool
isStalled() const override;
void
clearStall() override;
json::Value
getConsensusInfo() override;
std::size_t
getPrevProposers() const override;
std::chrono::milliseconds
getPrevRoundTime() const override;
json::Value
getServerInfo(bool human, bool admin, bool counters) override;
void
@@ -788,6 +802,8 @@ private:
std::atomic<bool> amendmentWarned_{false};
std::atomic<bool> unlBlocked_{false};
std::atomic<std::int64_t> stallDeadlineMs_{0};
ClosureCounter<void, boost::system::error_code const&> waitHandlerCounter_;
boost::asio::steady_timer heartbeatTimer_;
boost::asio::steady_timer clusterTimer_;
@@ -915,6 +931,16 @@ NetworkOPsImp::getOperatingMode() const
return mode_;
}
NetworkOPs::StateAccountingData
NetworkOPsImp::getStateAccountingData()
{
auto const data = accounting_.getCounterData();
std::array<NetworkOPs::AccountingCounter, 5> out;
for (std::size_t i = 0; i < out.size(); ++i)
out[i] = {data.counters[i].transitions, data.counters[i].dur};
return {out, data.mode, data.start, data.initialSyncUs};
}
inline std::string
NetworkOPsImp::strOperatingMode(bool const admin /* = false */) const
{
@@ -1116,6 +1142,13 @@ NetworkOPsImp::processHeartbeatTimer()
CLOG(clog.ss()) << ". ";
}
if (isStalled())
{
CLOG(clog.ss()) << "node is stalled, skipping consensus timerEntry. ";
setHeartbeatTimer();
return;
}
consensus_.timerEntry(registry_.get().getTimeKeeper().closeTime(), clog.ss());
CLOG(clog.ss()) << "consensus phase " << to_string(lastConsensusPhase_);
@@ -2188,6 +2221,35 @@ NetworkOPsImp::consensusViewChange()
}
}
void
NetworkOPsImp::setStall(std::chrono::milliseconds duration)
{
auto const deadline = std::chrono::steady_clock::now() + duration;
auto const ms =
std::chrono::duration_cast<std::chrono::milliseconds>(deadline.time_since_epoch()).count();
stallDeadlineMs_.store(ms, std::memory_order_relaxed);
JLOG(journal_.warn()) << "Node stalled for " << duration.count() << "ms";
}
bool
NetworkOPsImp::isStalled() const
{
auto const deadline = stallDeadlineMs_.load(std::memory_order_relaxed);
if (deadline == 0)
return false;
auto const now = std::chrono::duration_cast<std::chrono::milliseconds>(
std::chrono::steady_clock::now().time_since_epoch())
.count();
return now < deadline;
}
void
NetworkOPsImp::clearStall()
{
stallDeadlineMs_.store(0, std::memory_order_relaxed);
JLOG(journal_.warn()) << "Node stall cleared";
}
void
NetworkOPsImp::pubManifest(Manifest const& mo)
{
@@ -2580,6 +2642,18 @@ NetworkOPsImp::getConsensusInfo()
return consensus_.getJson(true);
}
std::size_t
NetworkOPsImp::getPrevProposers() const
{
return consensus_.prevProposers();
}
std::chrono::milliseconds
NetworkOPsImp::getPrevRoundTime() const
{
return consensus_.prevRoundTime();
}
json::Value
NetworkOPsImp::getServerInfo(bool human, bool admin, bool counters)
{

View File

@@ -12,6 +12,7 @@
#include <boost/intrusive/set.hpp>
#include <optional>
#include <vector>
namespace xrpl {
@@ -38,7 +39,14 @@ class Config;
*/
class TxQ
{
private:
std::mutex debugTxInjectMutex;
std::vector<STTx> debugTxInjectQueue;
public:
void
debugTxInject(STTx const& txn);
/// Fee level for single-signed reference transaction.
static constexpr FeeLevel64 kBaseLevel{256};

View File

@@ -657,6 +657,11 @@ public:
hash_set<PublicKey>
getTrustedMasterKeys() const;
void
debugSetTrusted(
std::vector<PublicKey> const& validators,
std::optional<std::size_t> quorumOverride = std::nullopt);
/**
* get the validator list threshold
* @return the threshold

View File

@@ -98,6 +98,13 @@ increase(FeeLevel64 level, std::uint32_t increasePercent)
//////////////////////////////////////////////////////////////////////////
void
TxQ::debugTxInject(STTx const& txn)
{
std::lock_guard<std::mutex> const _(debugTxInjectMutex);
debugTxInjectQueue.push_back(txn);
}
std::size_t
TxQ::FeeMetrics::update(
Application& app,
@@ -1418,6 +1425,21 @@ TxQ::accept(Application& app, OpenView& view)
auto const metricsSnapshot = feeMetrics_.getSnapshot();
// try to inject any debug txns waiting in the debug queue
{
std::unique_lock<std::mutex> trylock(TxQ::debugTxInjectMutex, std::try_to_lock);
if (trylock.owns_lock() && !debugTxInjectQueue.empty())
{
for (STTx const& txn : debugTxInjectQueue)
{
auto const [result, didApply, metadata] = xrpl::apply(app, view, txn, TapNone, j_);
if (didApply)
ledgerChanged = true;
}
debugTxInjectQueue.clear();
}
}
for (auto candidateIter = byFee_.begin(); candidateIter != byFee_.end();)
{
auto& account = byAccount_.at(candidateIter->account);

View File

@@ -2078,6 +2078,48 @@ ValidatorList::getTrustedMasterKeys() const
return trustedMasterKeys_;
}
void
ValidatorList::debugSetTrusted(
std::vector<PublicKey> const& validators,
std::optional<std::size_t> quorumOverride)
{
std::lock_guard const lock{mutex_};
trustedMasterKeys_.clear();
keyListings_.clear();
localPublisherList_.list.clear();
for (auto const& k : validators)
{
trustedMasterKeys_.insert(k);
keyListings_[k] = 1;
localPublisherList_.list.push_back(k);
}
trustedSigningKeys_.clear();
for (auto const& k : trustedMasterKeys_)
{
auto const signingKey = validatorManifests_.getSigningKey(k);
if (signingKey)
trustedSigningKeys_.insert(*signingKey);
else
trustedSigningKeys_.insert(k);
}
if (quorumOverride)
{
quorum_ = *quorumOverride;
}
else
{
auto const unlSize = trustedMasterKeys_.size();
quorum_ = calculateQuorum(unlSize, unlSize, unlSize);
}
JLOG(j_.warn()) << "debugSetTrusted: set " << trustedMasterKeys_.size()
<< " trusted validators, quorum=" << quorum_;
}
std::size_t
ValidatorList::getListThreshold() const
{

View File

@@ -134,6 +134,10 @@ public:
// Entries from [ips_fixed] config stanza
std::vector<std::string> ipsFixed;
// Entries from [datagram_monitor]: "<IP> <port>" UDP targets the
// DatagramMonitor sends node-stats packets to (XDGM, every 1s).
std::vector<std::string> DATAGRAM_MONITOR;
StartUpType startUp = StartUpType::Normal;
bool startValid = false;

View File

@@ -27,6 +27,7 @@ struct ConfigSection
#define SECTION_BETA_RPC_API "beta_rpc_api"
#define SECTION_CLUSTER_NODES "cluster_nodes"
#define SECTION_COMPRESSION "compression"
#define SECTION_DATAGRAM_MONITOR "datagram_monitor"
#define SECTION_DEBUG_LOGFILE "debug_logfile"
#define SECTION_ELB_SUPPORT "elb_support"
#define SECTION_FEE_DEFAULT "fee_default"

View File

@@ -483,6 +483,9 @@ Config::loadFromString(std::string const& fileContents)
if (auto s = getIniFileSection(secConfig, SECTION_IPS_FIXED))
ipsFixed = *s;
if (auto s = getIniFileSection(secConfig, SECTION_DATAGRAM_MONITOR))
DATAGRAM_MONITOR = *s;
// if the user has specified ip:port then replace : with a space.
{
auto replaceColons = [](std::vector<std::string>& strVec) {

View File

@@ -228,6 +228,10 @@ Handler const kHandlerArray[]{
.valueMethod = byRef(&doNFTSellOffers),
.role = Role::USER,
.condition = Condition::NoCondition},
{.name = "node_stall",
.valueMethod = byRef(&doNodeStall),
.role = Role::ADMIN,
.condition = Condition::NoCondition},
{.name = "noripple_check",
.valueMethod = byRef(&doNoRippleCheck),
.role = Role::USER,
@@ -294,6 +298,10 @@ Handler const kHandlerArray[]{
.valueMethod = byRef(&doSignFor),
.role = Role::USER,
.condition = Condition::NoCondition},
{.name = "inject",
.valueMethod = byRef(&doInject),
.role = Role::ADMIN,
.condition = Condition::NeedsCurrentLedger},
{.name = "simulate",
.valueMethod = byRef(&doSimulate),
.role = Role::USER,
@@ -332,6 +340,10 @@ Handler const kHandlerArray[]{
.valueMethod = byRef(&doUnlList),
.role = Role::ADMIN,
.condition = Condition::NoCondition},
{.name = "unl_set",
.valueMethod = byRef(&doUnlSet),
.role = Role::ADMIN,
.condition = Condition::NoCondition},
{.name = "validation_create",
.valueMethod = byRef(&doValidationCreate),
.role = Role::ADMIN,

View File

@@ -79,6 +79,8 @@ doNFTBuyOffers(RPC::JsonContext&);
json::Value
doNFTSellOffers(RPC::JsonContext&);
json::Value
doNodeStall(RPC::JsonContext&);
json::Value
doNoRippleCheck(RPC::JsonContext&);
json::Value
doOwnerInfo(RPC::JsonContext&);
@@ -117,6 +119,8 @@ doSignFor(RPC::JsonContext&);
json::Value
doSimulate(RPC::JsonContext&);
json::Value
doInject(RPC::JsonContext&);
json::Value
doStop(RPC::JsonContext&);
json::Value
doSubmit(RPC::JsonContext&);
@@ -135,6 +139,8 @@ doTxReduceRelay(RPC::JsonContext&);
json::Value
doUnlList(RPC::JsonContext&);
json::Value
doUnlSet(RPC::JsonContext&);
json::Value
doUnsubscribe(RPC::JsonContext&);
json::Value
doValidationCreate(RPC::JsonContext&);

View File

@@ -3,11 +3,15 @@
#include <xrpld/rpc/Context.h>
#include <xrpl/json/json_value.h>
#include <xrpl/protocol/ErrorCodes.h>
#include <xrpl/protocol/PublicKey.h>
#include <xrpl/protocol/RPCErr.h>
#include <xrpl/protocol/jss.h>
#include <xrpl/protocol/tokens.h>
#include <optional>
#include <utility>
#include <vector>
namespace xrpl {
@@ -29,4 +33,52 @@ doUnlList(RPC::JsonContext& context)
return obj;
}
json::Value
doUnlSet(RPC::JsonContext& context)
{
if (context.role != Role::ADMIN)
return rpcError(RpcNoPermission);
if (!context.params.isMember(jss::validators) || !context.params[jss::validators].isArray())
return rpcError(RpcInvalidParams);
auto const& vals = context.params[jss::validators];
std::vector<PublicKey> keys;
keys.reserve(vals.size());
for (auto const& v : vals)
{
if (!v.isString())
return rpcError(RpcInvalidParams);
auto const pk = parseBase58<PublicKey>(TokenType::NodePublic, v.asString());
if (!pk)
{
json::Value jvResult;
jvResult[jss::error] = "invalidParams";
jvResult[jss::error_message] = "Invalid validator key: " + v.asString();
return jvResult;
}
keys.push_back(*pk);
}
std::optional<std::size_t> quorumOverride;
if (context.params.isMember("quorum") && context.params["quorum"].isIntegral())
{
quorumOverride = context.params["quorum"].asUInt();
}
context.app.getValidators().debugSetTrusted(keys, quorumOverride);
json::Value jvResult;
jvResult[jss::validators] = json::Value(json::ValueType::Array);
for (auto const& k : keys)
{
jvResult[jss::validators].append(toBase58(TokenType::NodePublic, k));
}
jvResult["quorum"] = static_cast<json::UInt>(context.app.getValidators().quorum());
jvResult[jss::status] = "success";
return jvResult;
}
} // namespace xrpl

View File

@@ -0,0 +1,44 @@
#include <xrpld/rpc/Context.h>
#include <xrpl/json/json_value.h>
#include <xrpl/protocol/ErrorCodes.h>
#include <xrpl/protocol/RPCErr.h>
#include <xrpl/protocol/jss.h>
#include <xrpl/server/NetworkOPs.h>
namespace xrpl {
json::Value
doNodeStall(RPC::JsonContext& context)
{
if (context.role != Role::ADMIN)
return rpcError(RpcNoPermission);
if (context.params.isMember("clear") && context.params["clear"].asBool())
{
context.netOps.clearStall();
json::Value jvResult;
jvResult[jss::status] = "success";
jvResult["stalled"] = false;
return jvResult;
}
std::int64_t durationMs = 30000;
if (context.params.isMember("duration_ms") && context.params["duration_ms"].isIntegral())
{
durationMs = context.params["duration_ms"].asInt();
if (durationMs <= 0)
return rpcError(RpcInvalidParams);
}
context.netOps.setStall(std::chrono::milliseconds(durationMs));
json::Value jvResult;
jvResult[jss::status] = "success";
jvResult["stalled"] = true;
jvResult["duration_ms"] = static_cast<int>(durationMs);
return jvResult;
}
} // namespace xrpl

View File

@@ -1,5 +1,6 @@
#include <xrpld/app/ledger/LedgerMaster.h>
#include <xrpld/app/misc/Transaction.h>
#include <xrpld/app/misc/TxQ.h>
#include <xrpld/rpc/Context.h>
#include <xrpld/rpc/Role.h>
#include <xrpld/rpc/detail/TransactionSign.h>
@@ -37,6 +38,37 @@ getFailHard(RPC::JsonContext const& context)
context.params.isMember(jss::fail_hard) && context.params[jss::fail_hard].asBool());
}
json::Value
doInject(RPC::JsonContext& context)
{
if (context.role != Role::ADMIN)
return rpcError(RpcNoPermission);
json::Value jvResult;
auto ret = strUnHex(context.params[jss::tx_blob].asString());
if (!ret || !ret->size())
return rpcError(RpcInvalidParams);
SerialIter sitTrans(makeSlice(*ret));
std::shared_ptr<STTx const> stpTrans;
try
{
stpTrans = std::make_shared<STTx const>(std::ref(sitTrans));
}
catch (std::exception& e)
{
jvResult[jss::error] = "invalidTransaction";
jvResult[jss::error_exception] = e.what();
jvResult[jss::in_queue] = false;
return jvResult;
}
context.app.getTxQ().debugTxInject(*stpTrans);
jvResult[jss::tx_json] = stpTrans->getJson(JsonOptions::Values::None);
jvResult[jss::in_queue] = true;
return jvResult;
}
// {
// tx_blob: <string> XOR tx_json: <object>,
// secret: <secret>