mirror of
https://github.com/XRPLF/rippled.git
synced 2026-08-24 07:40:53 +00:00
Compare commits
3 Commits
dangell7/m
...
dangell7/d
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
4978981608 | ||
|
|
78e859f56c | ||
|
|
2add5b3fb8 |
64
conan.lock
64
conan.lock
@@ -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": []
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -329,6 +329,7 @@ words:
|
||||
- writeme
|
||||
- wsrch
|
||||
- wthread
|
||||
- Xahau
|
||||
- xbridge
|
||||
- xchain
|
||||
- ximinez
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
899
src/test/consensus/ByzantinePartitionRecovery_test.cpp
Normal file
899
src/test/consensus/ByzantinePartitionRecovery_test.cpp
Normal 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
|
||||
@@ -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
|
||||
|
||||
882
src/xrpld/app/misc/DatagramMonitor.h
Normal file
882
src/xrpld/app/misc/DatagramMonitor.h
Normal 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
|
||||
@@ -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)
|
||||
{
|
||||
|
||||
@@ -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};
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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
|
||||
{
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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&);
|
||||
|
||||
@@ -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
|
||||
|
||||
44
src/xrpld/rpc/handlers/admin/server_control/NodeStall.cpp
Normal file
44
src/xrpld/rpc/handlers/admin/server_control/NodeStall.cpp
Normal 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
|
||||
@@ -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>
|
||||
|
||||
Reference in New Issue
Block a user