Merge branch 'develop' into ximinez/online-delete-gaps

This commit is contained in:
Ed Hennis
2026-07-27 17:00:51 -04:00
committed by GitHub
109 changed files with 2602 additions and 1895 deletions

View File

@@ -17,7 +17,7 @@
#include <utility>
#include <vector>
namespace xrpl::NodeStore {
namespace xrpl::node_store {
namespace {
constexpr std::size_t kPoolSizes[] = {1000, 10000, 100000};
@@ -326,4 +326,4 @@ registerStoreBatch(BackendConfig const& bc)
}();
} // namespace
} // namespace xrpl::NodeStore
} // namespace xrpl::node_store

View File

@@ -18,7 +18,7 @@
#include <utility>
#include <vector>
namespace xrpl::NodeStore {
namespace xrpl::node_store {
namespace {
// Number of distinct objects pre-generated per run.
@@ -240,4 +240,4 @@ registerWorkload(BackendConfig const& bc, Workload const& w)
}();
} // namespace
} // namespace xrpl::NodeStore
} // namespace xrpl::node_store

View File

@@ -32,7 +32,7 @@
// Shared helpers for the NodeStore benchmarks.
//
namespace xrpl::NodeStore {
namespace xrpl::node_store {
// Fill `bytes` of memory at `buffer` with random bits drawn from `g`.
template <class Generator>
@@ -315,4 +315,4 @@ backendConfigs()
return kConfigs;
}
} // namespace xrpl::NodeStore
} // namespace xrpl::node_store

View File

@@ -15,7 +15,7 @@
namespace beast {
std::string
printIdentifiers(SemanticVersion::identifier_list const& list)
printIdentifiers(SemanticVersion::IdentifierList const& list)
{
std::string ret;
@@ -115,7 +115,7 @@ extractIdentifier(std::string& value, bool allowLeadingZeroes, std::string& inpu
bool
extractIdentifiers(
SemanticVersion::identifier_list& identifiers,
SemanticVersion::IdentifierList& identifiers,
bool allowLeadingZeroes,
std::string& input)
{

View File

@@ -11,7 +11,7 @@
#include <mutex>
#include <vector>
namespace xrpl::NodeStore {
namespace xrpl::node_store {
BatchWriter::BatchWriter(Callback& callback, Scheduler& scheduler)
: callback_(callback), scheduler_(scheduler)
@@ -72,7 +72,7 @@ BatchWriter::writeBatch()
writeSet_.swap(set);
XRPL_ASSERT(
writeSet_.empty(), "xrpl::NodeStore::BatchWriter::writeBatch : writes not set");
writeSet_.empty(), "xrpl::node_store::BatchWriter::writeBatch : writes not set");
writeLoad_ = set.size();
if (set.empty())
@@ -107,4 +107,4 @@ BatchWriter::waitForWriting()
writeCondition_.wait(sl);
}
} // namespace xrpl::NodeStore
} // namespace xrpl::node_store

View File

@@ -30,7 +30,7 @@
#include <thread>
#include <utility>
namespace xrpl::NodeStore {
namespace xrpl::node_store {
Database::Database(
Scheduler& scheduler,
@@ -43,7 +43,7 @@ Database::Database(
, requestBundle_(get<int>(config, Keys::kRqBundle, 4))
, readThreads_(std::max(1, readThreads))
{
XRPL_ASSERT(readThreads, "xrpl::NodeStore::Database::Database : nonzero threads input");
XRPL_ASSERT(readThreads, "xrpl::node_store::Database::Database : nonzero threads input");
if (earliestLedgerSeq_ < 1)
Throw<std::runtime_error>("Invalid earliest_seq");
@@ -89,7 +89,7 @@ Database::Database(
{
XRPL_ASSERT(
!it->second.empty(),
"xrpl::NodeStore::Database::Database : non-empty "
"xrpl::node_store::Database::Database : non-empty "
"data");
auto const& hash = it->first;
@@ -164,7 +164,7 @@ Database::stop()
{
XRPL_ASSERT(
steady_clock::now() - start < 30s,
"xrpl::NodeStore::Database::stop : maximum stop duration");
"xrpl::node_store::Database::stop : maximum stop duration");
std::this_thread::yield();
}
@@ -213,7 +213,7 @@ Database::importInternal(Backend& dstBackend, Database& srcDB)
};
srcDB.forEach([&](std::shared_ptr<NodeObject> nodeObject) {
XRPL_ASSERT(nodeObject, "xrpl::NodeStore::Database::importInternal : non-null node");
XRPL_ASSERT(nodeObject, "xrpl::node_store::Database::importInternal : non-null node");
if (!nodeObject) // This should never happen
return;
@@ -257,7 +257,7 @@ Database::fetchNodeObject(
void
Database::getCountsJson(json::Value& obj)
{
XRPL_ASSERT(obj.isObject(), "xrpl::NodeStore::Database::getCountsJson : valid input type");
XRPL_ASSERT(obj.isObject(), "xrpl::node_store::Database::getCountsJson : valid input type");
{
std::unique_lock<std::mutex> const lock(readLock_);
@@ -276,4 +276,4 @@ Database::getCountsJson(json::Value& obj)
obj[jss::node_reads_duration_us] = std::to_string(fetchDurationUs_);
}
} // namespace xrpl::NodeStore
} // namespace xrpl::node_store

View File

@@ -15,7 +15,7 @@
#include <memory>
#include <utility>
namespace xrpl::NodeStore {
namespace xrpl::node_store {
void
DatabaseNodeImp::store(NodeObjectType type, Blob&& data, uint256 const& hash, std::uint32_t)
@@ -125,4 +125,4 @@ DatabaseNodeImp::fetchNodeObject(
return nodeObject;
}
} // namespace xrpl::NodeStore
} // namespace xrpl::node_store

View File

@@ -22,7 +22,7 @@
#include <string>
#include <utility>
namespace xrpl::NodeStore {
namespace xrpl::node_store {
DatabaseRotatingImp::DatabaseRotatingImp(
Scheduler& scheduler,
@@ -43,7 +43,7 @@ DatabaseRotatingImp::DatabaseRotatingImp(
void
DatabaseRotatingImp::rotate(
std::unique_ptr<NodeStore::Backend>&& newBackend,
std::unique_ptr<node_store::Backend>&& newBackend,
std::function<void(std::string const& writableName, std::string const& archiveName)> const& f)
{
// Pass these two names to the callback function
@@ -52,7 +52,7 @@ DatabaseRotatingImp::rotate(
// Hold on to current archive backend pointer until after the
// callback finishes. Only then will the archive directory be
// deleted.
std::shared_ptr<NodeStore::Backend> oldArchiveBackend;
std::shared_ptr<node_store::Backend> oldArchiveBackend;
std::uint64_t copyForwards = 0;
{
std::scoped_lock const lock(mutex_);
@@ -232,4 +232,4 @@ DatabaseRotatingImp::forEach(std::function<void(std::shared_ptr<NodeObject>)> f)
archive->forEach(f);
}
} // namespace xrpl::NodeStore
} // namespace xrpl::node_store

View File

@@ -10,7 +10,7 @@
#include <memory>
#include <utility>
namespace xrpl::NodeStore {
namespace xrpl::node_store {
DecodedBlob::DecodedBlob(void const* key, void const* value, int valueBytes) : key_(key)
{
@@ -55,7 +55,7 @@ DecodedBlob::DecodedBlob(void const* key, void const* value, int valueBytes) : k
std::shared_ptr<NodeObject>
DecodedBlob::createObject()
{
XRPL_ASSERT(success_, "xrpl::NodeStore::DecodedBlob::createObject : valid object type");
XRPL_ASSERT(success_, "xrpl::node_store::DecodedBlob::createObject : valid object type");
std::shared_ptr<NodeObject> object;
@@ -69,4 +69,4 @@ DecodedBlob::createObject()
return object;
}
} // namespace xrpl::NodeStore
} // namespace xrpl::node_store

View File

@@ -3,7 +3,7 @@
#include <xrpl/nodestore/Scheduler.h>
#include <xrpl/nodestore/Task.h>
namespace xrpl::NodeStore {
namespace xrpl::node_store {
void
DummyScheduler::scheduleTask(Task& task)
@@ -22,4 +22,4 @@ DummyScheduler::onBatchWrite(BatchWriteReport const& report)
{
}
} // namespace xrpl::NodeStore
} // namespace xrpl::node_store

View File

@@ -22,7 +22,7 @@
#include <string>
#include <utility>
namespace xrpl::NodeStore {
namespace xrpl::node_store {
ManagerImp&
ManagerImp::instance()
@@ -112,7 +112,7 @@ ManagerImp::erase(Factory& factory)
std::scoped_lock const _(mutex_);
auto const iter =
std::ranges::find_if(list_, [&factory](Factory* other) { return other == &factory; });
XRPL_ASSERT(iter != list_.end(), "xrpl::NodeStore::ManagerImp::erase : valid input");
XRPL_ASSERT(iter != list_.end(), "xrpl::node_store::ManagerImp::erase : valid input");
list_.erase(iter);
}
@@ -135,4 +135,4 @@ Manager::instance()
return ManagerImp::instance();
}
} // namespace xrpl::NodeStore
} // namespace xrpl::node_store

View File

@@ -24,7 +24,7 @@
#include <tuple>
#include <utility>
namespace xrpl::NodeStore {
namespace xrpl::node_store {
struct MemoryDB
{
@@ -132,7 +132,7 @@ public:
Status
fetch(uint256 const& hash, std::shared_ptr<NodeObject>* pObject) override
{
XRPL_ASSERT(db_, "xrpl::NodeStore::MemoryBackend::fetch : non-null database");
XRPL_ASSERT(db_, "xrpl::node_store::MemoryBackend::fetch : non-null database");
std::scoped_lock const _(db_->mutex);
@@ -149,7 +149,7 @@ public:
void
store(std::shared_ptr<NodeObject> const& object) override
{
XRPL_ASSERT(db_, "xrpl::NodeStore::MemoryBackend::store : non-null database");
XRPL_ASSERT(db_, "xrpl::node_store::MemoryBackend::store : non-null database");
std::scoped_lock const _(db_->mutex);
db_->table.emplace(object->getHash(), object);
}
@@ -169,7 +169,7 @@ public:
void
forEach(std::function<void(std::shared_ptr<NodeObject>)> f) override
{
XRPL_ASSERT(db_, "xrpl::NodeStore::MemoryBackend::forEach : non-null database");
XRPL_ASSERT(db_, "xrpl::node_store::MemoryBackend::forEach : non-null database");
for (auto const& e : db_->table)
f(e.second);
}
@@ -216,4 +216,4 @@ MemoryFactory::createInstance(
return std::make_unique<MemoryBackend>(keyBytes, keyValues, journal);
}
} // namespace xrpl::NodeStore
} // namespace xrpl::node_store

View File

@@ -44,7 +44,7 @@
#include <string>
#include <utility>
namespace xrpl::NodeStore {
namespace xrpl::node_store {
class NuDBBackend : public Backend
{
@@ -136,7 +136,7 @@ public:
{
// LCOV_EXCL_START
UNREACHABLE(
"xrpl::NodeStore::NuDBBackend::open : database is already "
"xrpl::node_store::NuDBBackend::open : database is already "
"open");
JLOG(j.error()) << "database is already open";
return;
@@ -441,4 +441,4 @@ registerNuDBFactory(Manager& manager)
static NuDBFactory const kInstance{manager};
}
} // namespace xrpl::NodeStore
} // namespace xrpl::node_store

View File

@@ -13,7 +13,7 @@
#include <memory>
#include <string>
namespace xrpl::NodeStore {
namespace xrpl::node_store {
class NullBackend : public Backend
{
@@ -125,4 +125,4 @@ registerNullFactory(Manager& manager)
static NullFactory const kInstance{manager};
}
} // namespace xrpl::NodeStore
} // namespace xrpl::node_store

View File

@@ -42,7 +42,7 @@
#include <stdexcept>
#include <string>
namespace xrpl::NodeStore {
namespace xrpl::node_store {
class RocksDBEnv : public rocksdb::EnvWrapper
{
@@ -231,7 +231,7 @@ public:
{
// LCOV_EXCL_START
UNREACHABLE(
"xrpl::NodeStore::RocksDBBackend::open : database is already "
"xrpl::node_store::RocksDBBackend::open : database is already "
"open");
JLOG(journal.error()) << "database is already open";
return;
@@ -279,7 +279,7 @@ public:
Status
fetch(uint256 const& hash, std::shared_ptr<NodeObject>* pObject) override
{
XRPL_ASSERT(db, "xrpl::NodeStore::RocksDBBackend::fetch : non-null database");
XRPL_ASSERT(db, "xrpl::node_store::RocksDBBackend::fetch : non-null database");
pObject->reset();
Status status = Status::Ok;
@@ -339,7 +339,7 @@ public:
{
XRPL_ASSERT(
db,
"xrpl::NodeStore::RocksDBBackend::storeBatch : non-null "
"xrpl::node_store::RocksDBBackend::storeBatch : non-null "
"database");
rocksdb::WriteBatch wb;
@@ -369,7 +369,7 @@ public:
void
forEach(std::function<void(std::shared_ptr<NodeObject>)> f) override
{
XRPL_ASSERT(db, "xrpl::NodeStore::RocksDBBackend::forEach : non-null database");
XRPL_ASSERT(db, "xrpl::node_store::RocksDBBackend::forEach : non-null database");
rocksdb::ReadOptions const options;
std::unique_ptr<rocksdb::Iterator> it(db->NewIterator(options));
@@ -468,6 +468,6 @@ registerRocksDBFactory(Manager& manager)
static RocksDBFactory const kInstance{manager};
}
} // namespace xrpl::NodeStore
} // namespace xrpl::node_store
#endif

View File

@@ -491,7 +491,7 @@ public:
lastRotated = ledgerSeq - 1;
}
std::unique_ptr<NodeStore::Backend>
std::unique_ptr<node_store::Backend>
makeBackendRotating(jtx::Env& env, NodeStoreScheduler& scheduler, std::string path)
{
Section section{env.app().config().section(Sections::kNodeDatabase)};
@@ -502,7 +502,7 @@ public:
newPath = path;
section.set(Keys::kPath, newPath.string());
auto backend{NodeStore::Manager::instance().makeBackend(
auto backend{node_store::Manager::instance().makeBackend(
section,
megabytes(env.app().config().getValueFor(SizedItem::BurstSize, std::nullopt)),
scheduler,
@@ -551,7 +551,7 @@ public:
auto archiveBackend = makeBackendRotating(env, scheduler, archiveDb);
static constexpr int kReadThreads = 4;
auto dbr = std::make_unique<NodeStore::DatabaseRotatingImp>(
auto dbr = std::make_unique<node_store::DatabaseRotatingImp>(
scheduler,
kReadThreads,
std::move(writableBackend),

View File

@@ -1,266 +0,0 @@
#include <xrpl/beast/core/SemanticVersion.h>
#include <xrpl/beast/unit_test/suite.h>
#include <string>
namespace beast {
class SemanticVersion_test : public unit_test::Suite
{
using identifier_list = SemanticVersion::identifier_list;
public:
void
checkPass(std::string const& input, bool shouldPass = true)
{
SemanticVersion v;
if (shouldPass)
{
BEAST_EXPECT(v.parse(input));
BEAST_EXPECT(v.print() == input);
}
else
{
BEAST_EXPECT(!v.parse(input));
}
}
void
checkFail(std::string const& input)
{
checkPass(input, false);
}
// check input and input with appended metadata
void
checkMeta(std::string const& input, bool shouldPass)
{
checkPass(input, shouldPass);
checkPass(input + "+a", shouldPass);
checkPass(input + "+1", shouldPass);
checkPass(input + "+a.b", shouldPass);
checkPass(input + "+ab.cd", shouldPass);
checkFail(input + "!");
checkFail(input + "+");
checkFail(input + "++");
checkFail(input + "+!");
checkFail(input + "+.");
checkFail(input + "+a.!");
}
void
checkMetaFail(std::string const& input)
{
checkMeta(input, false);
}
// check input, input with appended release data,
// input with appended metadata, and input with both
// appended release data and appended metadata
//
void
checkRelease(std::string const& input, bool shouldPass = true)
{
checkMeta(input, shouldPass);
checkMeta(input + "-1", shouldPass);
checkMeta(input + "-a", shouldPass);
checkMeta(input + "-a1", shouldPass);
checkMeta(input + "-a1.b1", shouldPass);
checkMeta(input + "-ab.cd", shouldPass);
checkMeta(input + "--", shouldPass);
checkMetaFail(input + "+");
checkMetaFail(input + "!");
checkMetaFail(input + "-");
checkMetaFail(input + "-!");
checkMetaFail(input + "-.");
checkMetaFail(input + "-a.!");
checkMetaFail(input + "-0.a");
}
// Checks the major.minor.version string alone and with all
// possible combinations of release identifiers and metadata.
//
void
check(std::string const& input, bool shouldPass = true)
{
checkRelease(input, shouldPass);
}
void
negcheck(std::string const& input)
{
check(input, false);
}
void
testParse()
{
testcase("parsing");
check("0.0.0");
check("1.2.3");
check("2147483647.2147483647.2147483647"); // max int
// negative values
negcheck("-1.2.3");
negcheck("1.-2.3");
negcheck("1.2.-3");
// missing parts
negcheck("");
negcheck("1");
negcheck("1.");
negcheck("1.2");
negcheck("1.2.");
negcheck(".2.3");
// whitespace
negcheck(" 1.2.3");
negcheck("1 .2.3");
negcheck("1.2 .3");
negcheck("1.2.3 ");
// leading zeroes
negcheck("01.2.3");
negcheck("1.02.3");
negcheck("1.2.03");
}
static identifier_list
ids()
{
return identifier_list();
}
static identifier_list
ids(std::string const& s1)
{
identifier_list v;
v.push_back(s1);
return v;
}
static identifier_list
ids(std::string const& s1, std::string const& s2)
{
identifier_list v;
v.push_back(s1);
v.push_back(s2);
return v;
}
static identifier_list
ids(std::string const& s1, std::string const& s2, std::string const& s3)
{
identifier_list v;
v.push_back(s1);
v.push_back(s2);
v.push_back(s3);
return v;
}
// Checks the decomposition of the input into appropriate values
void
checkValues(
std::string const& input,
int majorVersion,
int minorVersion,
int patchVersion,
identifier_list const& preReleaseIdentifiers = identifier_list(),
identifier_list const& metaData = identifier_list())
{
SemanticVersion v;
BEAST_EXPECT(v.parse(input));
BEAST_EXPECT(v.majorVersion == majorVersion);
BEAST_EXPECT(v.minorVersion == minorVersion);
BEAST_EXPECT(v.patchVersion == patchVersion);
BEAST_EXPECT(v.preReleaseIdentifiers == preReleaseIdentifiers);
BEAST_EXPECT(v.metaData == metaData);
}
void
testValues()
{
testcase("values");
checkValues("0.1.2", 0, 1, 2);
checkValues("1.2.3", 1, 2, 3);
checkValues("1.2.3-rc1", 1, 2, 3, ids("rc1"));
checkValues("1.2.3-rc1.debug", 1, 2, 3, ids("rc1", "debug"));
checkValues("1.2.3-rc1.debug.asm", 1, 2, 3, ids("rc1", "debug", "asm"));
checkValues("1.2.3+full", 1, 2, 3, ids(), ids("full"));
checkValues("1.2.3+full.prod", 1, 2, 3, ids(), ids("full", "prod"));
checkValues("1.2.3+full.prod.x86", 1, 2, 3, ids(), ids("full", "prod", "x86"));
checkValues(
"1.2.3-rc1.debug.asm+full.prod.x86",
1,
2,
3,
ids("rc1", "debug", "asm"),
ids("full", "prod", "x86"));
}
// makes sure the left version is less than the right
void
checkLessInternal(std::string const& lhs, std::string const& rhs)
{
SemanticVersion left;
SemanticVersion right;
BEAST_EXPECT(left.parse(lhs));
BEAST_EXPECT(right.parse(rhs));
BEAST_EXPECT(compare(left, left) == 0);
BEAST_EXPECT(compare(right, right) == 0);
BEAST_EXPECT(compare(left, right) < 0);
BEAST_EXPECT(compare(right, left) > 0);
BEAST_EXPECT(left < right);
BEAST_EXPECT(right > left);
BEAST_EXPECT(left == left);
BEAST_EXPECT(right == right);
}
void
checkLess(std::string const& lhs, std::string const& rhs)
{
checkLessInternal(lhs, rhs);
checkLessInternal(lhs + "+meta", rhs);
checkLessInternal(lhs, rhs + "+meta");
checkLessInternal(lhs + "+meta", rhs + "+meta");
}
void
testCompare()
{
testcase("comparisons");
checkLess("1.0.0-alpha", "1.0.0-alpha.1");
checkLess("1.0.0-alpha.1", "1.0.0-alpha.beta");
checkLess("1.0.0-alpha.beta", "1.0.0-beta");
checkLess("1.0.0-beta", "1.0.0-beta.2");
checkLess("1.0.0-beta.2", "1.0.0-beta.11");
checkLess("1.0.0-beta.11", "1.0.0-rc.1");
checkLess("1.0.0-rc.1", "1.0.0");
checkLess("0.9.9", "1.0.0");
}
void
run() override
{
testParse();
testValues();
testCompare();
}
};
BEAST_DEFINE_TESTSUITE(SemanticVersion, beast, beast);
} // namespace beast

View File

@@ -1,115 +0,0 @@
#include <xrpl/beast/unit_test/suite.h>
#include <xrpl/beast/utility/Zero.h>
namespace beast {
struct AdlTester
{
};
int
signum(AdlTester)
{
return 0;
}
namespace inner_adl_test {
struct AdlTester2
{
};
int
signum(AdlTester2)
{
return 0;
}
} // namespace inner_adl_test
class Zero_test : public beast::unit_test::Suite
{
private:
struct IntegerWrapper
{
int value;
IntegerWrapper(int v) : value(v)
{
}
[[nodiscard]] int
signum() const
{
return value;
}
};
public:
void
expectSame(bool result, bool correct, char const* message)
{
expect(result == correct, message);
}
void
testLhsZero(IntegerWrapper x)
{
expectSame(x >= kZero, x.signum() >= 0, "lhs greater-than-or-equal-to");
expectSame(x > kZero, x.signum() > 0, "lhs greater than");
expectSame(x == kZero, x.signum() == 0, "lhs equal to");
expectSame(x != kZero, x.signum() != 0, "lhs not equal to");
expectSame(x < kZero, x.signum() < 0, "lhs less than");
expectSame(x <= kZero, x.signum() <= 0, "lhs less-than-or-equal-to");
}
void
testLhsZero()
{
testcase("lhs zero");
testLhsZero(-7);
testLhsZero(0);
testLhsZero(32);
}
void
testRhsZero(IntegerWrapper x)
{
expectSame(kZero >= x, 0 >= x.signum(), "rhs greater-than-or-equal-to");
expectSame(kZero > x, 0 > x.signum(), "rhs greater than");
expectSame(kZero == x, 0 == x.signum(), "rhs equal to");
expectSame(kZero != x, 0 != x.signum(), "rhs not equal to");
expectSame(kZero < x, 0 < x.signum(), "rhs less than");
expectSame(kZero <= x, 0 <= x.signum(), "rhs less-than-or-equal-to");
}
void
testRhsZero()
{
testcase("rhs zero");
testRhsZero(-4);
testRhsZero(0);
testRhsZero(64);
}
void
testAdl()
{
expect(AdlTester{} == kZero, "ADL failure!");
expect(inner_adl_test::AdlTester2{} == kZero, "ADL failure!");
}
void
run() override
{
testLhsZero();
testRhsZero();
testAdl();
}
};
BEAST_DEFINE_TESTSUITE(Zero, beast, beast);
} // namespace beast

View File

@@ -40,6 +40,17 @@ public:
*/
[[nodiscard]] virtual unsigned
version() const = 0;
/**
* Close the client's connection to the server.
*
* Releases the connection the client holds against the server's per-port
* connection limit. After this call the client must not be used to
* invoke() again. Tests use this to deterministically free the slot
* rather than waiting for the server's idle timeout to drop it.
*/
virtual void
disconnect() = 0;
};
} // namespace xrpl::test

View File

@@ -482,11 +482,52 @@ public:
app().getNumberOfThreads() == 1,
"syncClose() is only useful on an application with a single thread");
auto const result = close();
auto serverBarrier = std::make_shared<std::promise<void>>();
auto future = serverBarrier->get_future();
boost::asio::post(app().getIOContext(), [serverBarrier]() { serverBarrier->set_value(); });
auto const status = future.wait_for(timeout);
return result && status == std::future_status::ready;
return result && drainServerIo(timeout);
}
/**
* Disconnect the Env's built-in client and wait for the server to
* register the dropped connection.
*
* Env holds one persistent client connection to the server's RPC port for
* its whole lifetime (see client()), and that connection counts against
* the port's connection limit. Tests that need a known starting occupancy
* can call this to deterministically release that slot instead of waiting
* out the server's localhost idle timeout.
*
* The server decrements its per-port connection count in the peer's
* destructor, which runs when the io_context processes the end-of-stream
* on the closed socket. After closing the client this drains the server's
* io_context twice: the first barrier guarantees the reactor has reaped
* the closed socket and queued the peer's teardown, and the second
* guarantees that teardown (and therefore the count decrement) has run.
*
* This is only sound when the server uses a single io_context thread, so
* that draining establishes ordering against the teardown - configure the
* Env with singleThreadIo() (as syncClose() also requires). Like
* syncClose(), it relies on loopback teardown latency being negligible.
*
* @param timeout Maximum time to wait for each barrier task to execute
* @return true if both barriers executed within timeout, false otherwise
*/
[[nodiscard]] bool
disconnectClient(std::chrono::steady_clock::duration timeout = std::chrono::seconds{1})
{
XRPL_ASSERT(
app().getNumberOfThreads() == 1,
"disconnectClient() is only useful on an application with a single "
"thread");
bundle_.client->disconnect();
// Drain the server's single io thread twice: the first barrier flushes
// the reactor's reap of the closed socket (queuing the peer teardown),
// the second flushes that teardown - and therefore the connection-count
// decrement. Both run unconditionally so a timed-out first drain does
// not short-circuit the second.
bool const reaped = drainServerIo(timeout);
bool const toreDown = drainServerIo(timeout);
return reaped && toreDown;
}
/**
@@ -846,6 +887,25 @@ public:
}
private:
/**
* Drain the (single) server io_context thread once.
*
* Posts a barrier task to the server's io_context and blocks until it
* runs, so every task queued before it has been processed. Only meaningful
* with a single io thread (see syncClose()/disconnectClient()).
*
* @param timeout Maximum time to wait for the barrier task to execute
* @return true if the barrier ran within timeout, false otherwise
*/
[[nodiscard]] bool
drainServerIo(std::chrono::steady_clock::duration timeout)
{
auto barrier = std::make_shared<std::promise<void>>();
auto future = barrier->get_future();
boost::asio::post(app().getIOContext(), [barrier]() { barrier->set_value(); });
return future.wait_for(timeout) == std::future_status::ready;
}
void
fund(bool setDefaultRipple, STAmount const& amount, Account const& account);

View File

@@ -13,8 +13,9 @@ class WSClient_test : public beast::unit_test::Suite
{
public:
void
run() override
testSmoke()
{
testcase("smoke");
using namespace jtx;
Env env(*this);
auto wsc = makeWSClient(env.app().config());
@@ -28,6 +29,47 @@ public:
auto jv = wsc->getMsg(std::chrono::seconds(1));
pass();
}
void
testGracefulDisconnect()
{
testcase("graceful disconnect");
using namespace jtx;
using namespace std::chrono;
Env env(*this);
auto wsc = makeWSClient(env.app().config());
// Put real traffic on the connection before closing it.
json::Value stream;
stream["streams"] = json::ValueType::Array;
stream["streams"].append("ledger");
auto const sub = wsc->invoke("subscribe", stream);
BEAST_EXPECT(sub.isMember("result") || sub.isMember("status"));
// disconnect() performs a graceful WebSocket closing handshake and
// blocks until the server acknowledges. On loopback that completes in
// well under its internal 1s timeout; only a broken async_close/ack
// coordination would fall through to the force-close path at ~1s. A
// generous bound keeps this from flaking under load while still
// catching that regression.
auto const start = steady_clock::now();
wsc->disconnect();
auto const elapsed = duration_cast<milliseconds>(steady_clock::now() - start);
BEAST_EXPECT(elapsed < milliseconds{750});
// disconnect() must be idempotent: a second call (and the subsequent
// destructor) must not hang, double-close, or crash.
wsc->disconnect();
pass();
}
void
run() override
{
testSmoke();
testGracefulDisconnect();
}
};
BEAST_DEFINE_TESTSUITE(WSClient, jtx, xrpl);

View File

@@ -14,18 +14,23 @@
#include <xrpl/server/Port.h>
#include <boost/asio/buffer.hpp>
#include <boost/asio/error.hpp>
#include <boost/asio/io_context.hpp>
#include <boost/asio/ip/address_v4.hpp>
#include <boost/asio/ip/address_v6.hpp>
#include <boost/asio/ip/tcp.hpp>
#include <boost/beast/core/multi_buffer.hpp>
#include <boost/beast/http/dynamic_body.hpp>
#include <boost/beast/http/error.hpp>
#include <boost/beast/http/message.hpp>
#include <boost/beast/http/read.hpp>
#include <boost/beast/http/string_body.hpp>
#include <boost/beast/http/verb.hpp>
#include <boost/beast/http/write.hpp>
#include <boost/system/system_error.hpp>
#include <algorithm>
#include <array>
#include <iostream>
#include <memory>
#include <sstream>
@@ -84,6 +89,40 @@ class JSONRPCClient : public AbstractClient
boost::beast::multi_buffer bout_;
unsigned rpcVersion_;
bool disconnected_ = false;
// Errors that mean the persistent keep-alive connection was dropped by the
// server (rather than a genuine protocol failure), so the request can be
// safely retried on a fresh connection.
static bool
droppedConnection(boost::system::error_code const& ec)
{
namespace error = boost::asio::error;
static auto const kDroppedConnectionErrors = std::to_array<boost::system::error_code>({
boost::beast::http::error::end_of_stream,
error::eof,
error::connection_reset,
error::connection_aborted,
error::broken_pipe,
error::not_connected,
});
return std::ranges::any_of(
kDroppedConnectionErrors,
[&ec](boost::system::error_code const& e) { return ec == e; });
}
// Tear down and re-establish the socket to ep_, discarding any buffered
// bytes left over from the dropped connection.
void
reconnect()
{
boost::system::error_code ec;
stream_.close(ec);
bin_.clear();
stream_.connect(ep_);
}
public:
explicit JSONRPCClient(Config const& cfg, unsigned rpcVersion)
: ep_(getEndpoint(cfg)), stream_(ios_), rpcVersion_(rpcVersion)
@@ -91,12 +130,10 @@ public:
stream_.connect(ep_);
}
/*
Return value is an Object type with up to three keys:
status
error
result
*/
// Return value is an Object type with up to three keys:
// status
// error
// result
json::Value
invoke(std::string const& cmd, json::Value const& params) override
{
@@ -104,6 +141,13 @@ public:
using namespace boost::asio;
using namespace std::string_literals;
// Once disconnect() has released the slot, the client must not be
// reused (see AbstractClient::disconnect). Refuse rather than let the
// failed write/read below trip the reconnect path and silently
// re-consume a connection slot, which would defeat disconnectClient().
if (disconnected_)
Throw<std::logic_error>("JSONRPCClient::invoke called after disconnect()");
request<string_body> req;
req.method(boost::beast::http::verb::post);
req.target("/");
@@ -131,10 +175,29 @@ public:
req.body() = to_string(jr);
}
req.prepare_payload();
write(stream_, req);
// The client keeps a single keep-alive connection for its whole
// lifetime, but the server drops idle localhost connections after a few
// seconds (BaseHTTPPeer::kTimeoutSecondsLocal). If a slow gap between
// requests let the server close the socket, the write/read here fails
// with end_of_stream; reconnect and retry the request exactly once.
response<dynamic_body> res;
read(stream_, bin_, res);
auto writeAndRead = [&] {
write(stream_, req);
read(stream_, bin_, res);
};
try
{
writeAndRead();
}
catch (boost::system::system_error const& e)
{
if (!droppedConnection(e.code()))
throw;
reconnect();
res = {};
writeAndRead();
}
json::Reader jr;
json::Value jv;
@@ -151,6 +214,19 @@ public:
{
return rpcVersion_;
}
void
disconnect() override
{
if (disconnected_)
return;
disconnected_ = true;
boost::system::error_code ec;
stream_.shutdown(boost::asio::ip::tcp::socket::shutdown_both, ec);
stream_.close(ec);
}
};
std::unique_ptr<AbstractClient>

View File

@@ -2,6 +2,7 @@
#include <xrpld/core/Config.h>
#include <xrpl/basics/Mutex.hpp>
#include <xrpl/basics/contract.h>
#include <xrpl/config/BasicConfig.h>
#include <xrpl/config/Constants.h>
@@ -112,10 +113,11 @@ class WSClientImpl : public WSClient
bool peerClosed_ = false;
// synchronize destructor
bool b0_ = false;
std::mutex m0_;
std::condition_variable cv0_;
// disconnect() waits on this until the read loop ends (for any reason:
// the server acknowledged our close, or a timeout force-closed the socket).
static constexpr auto kDisconnectTimeout = std::chrono::seconds{1};
xrpl::Mutex<bool> readEnded_;
std::condition_variable readEndCv_;
// synchronize message queue
std::mutex m_;
@@ -127,23 +129,26 @@ class WSClientImpl : public WSClient
void
cleanup()
{
boost::asio::post(ios_, boost::asio::bind_executor(strand_, [this] {
if (!peerClosed_)
{
ws_.async_close(
{}, boost::asio::bind_executor(strand_, [&](error_code) {
try
{
stream_.cancel();
}
// NOLINTNEXTLINE(bugprone-empty-catch)
catch (boost::system::system_error const&)
{
// ignored
}
}));
}
}));
boost::asio::post(
ios_, //
boost::asio::bind_executor(strand_, [this] {
if (!peerClosed_)
{
ws_.async_close(
{}, //
boost::asio::bind_executor(strand_, [&](error_code) {
try
{
stream_.cancel();
}
// NOLINTNEXTLINE(bugprone-empty-catch)
catch (boost::system::system_error const&)
{
// ignored
}
}));
}
}));
work_ = std::nullopt;
thread_.join();
}
@@ -289,6 +294,44 @@ public:
return rpcVersion_;
}
void
disconnect() override
{
// Perform a graceful WebSocket closing handshake and block until the
// read loop ends, so the server observes a clean close (not a RST) and
// has finished tearing the connection down by the time we return.
// If the server already closed, the wait below returns immediately.
boost::asio::post(
ios_,
boost::asio::bind_executor(
strand_, //
[this] {
if (!peerClosed_)
{
ws_.async_close(
boost::beast::websocket::close_code::normal,
boost::asio::bind_executor(strand_, [](error_code) {}));
}
}));
auto lock = readEnded_.lock<std::unique_lock>();
readEndCv_.wait_for(lock, kDisconnectTimeout, [&lock] { return *lock; });
// On timeout (server gone or not replying) force the socket closed so
// the outstanding read ends and the worker thread can later be joined.
if (!*lock)
{
boost::asio::post(
ios_,
boost::asio::bind_executor(
strand_, //
[this] {
boost::system::error_code ec;
stream_.close(ec);
}));
}
}
private:
void
onReadMsg(error_code const& ec)
@@ -297,33 +340,31 @@ private:
{
if (ec == boost::beast::websocket::error::closed)
peerClosed_ = true;
*readEnded_.lock() = true;
readEndCv_.notify_all();
return;
}
json::Value jv;
json::Reader jr;
jr.parse(bufferString(rb_.data()), jv);
rb_.consume(rb_.size());
auto m = std::make_shared<Msg>(std::move(jv));
{
std::scoped_lock const lock(m_);
msgs_.push_front(m);
cv_.notify_all();
}
ws_.async_read(
rb_, boost::asio::bind_executor(strand_, [this](error_code const& ec, std::size_t) {
onReadMsg(ec);
}));
}
// Called when the read op terminates
void
onReadDone()
{
std::scoped_lock const lock(m0_);
b0_ = true;
cv0_.notify_all();
}
};
std::unique_ptr<WSClient>

View File

@@ -1,110 +0,0 @@
#include <test/nodestore/TestBase.h>
#include <test/unit_test/SuiteJournal.h>
#include <xrpl/basics/ByteUtilities.h>
#include <xrpl/beast/unit_test/suite.h>
#include <xrpl/beast/utility/Journal.h>
#include <xrpl/beast/utility/temp_dir.h>
#include <xrpl/beast/xor_shift_engine.h>
#include <xrpl/config/BasicConfig.h>
#include <xrpl/config/Constants.h>
#include <xrpl/nodestore/Backend.h>
#include <xrpl/nodestore/DummyScheduler.h>
#include <xrpl/nodestore/Manager.h>
#include <xrpl/nodestore/Types.h>
#include <algorithm>
#include <cstdint>
#include <memory>
#include <string>
namespace xrpl::NodeStore {
// Tests the Backend interface
//
class Backend_test : public TestBase
{
public:
void
testBackend(std::string const& type, std::uint64_t const seedValue, int numObjsToTest = 2000)
{
DummyScheduler scheduler;
testcase("Backend type=" + type);
Section params;
beast::TempDir const tempDir;
params.set(Keys::kType, type);
params.set(Keys::kPath, tempDir.path());
beast::xor_shift_engine rng(seedValue);
// Create a batch
auto batch = createPredictableBatch(numObjsToTest, rng());
using beast::Severity;
test::SuiteJournal journal("Backend_test", *this);
{
// Open the backend
std::unique_ptr<Backend> backend =
Manager::instance().makeBackend(params, megabytes(4), scheduler, journal);
backend->open();
// Write the batch
storeBatch(*backend, batch);
{
// Read it back in
Batch copy;
fetchCopyOfBatch(*backend, &copy, batch);
BEAST_EXPECT(areBatchesEqual(batch, copy));
}
{
// Reorder and read the copy again
std::shuffle(batch.begin(), batch.end(), rng);
Batch copy;
fetchCopyOfBatch(*backend, &copy, batch);
BEAST_EXPECT(areBatchesEqual(batch, copy));
}
}
{
// Re-open the backend
std::unique_ptr<Backend> backend =
Manager::instance().makeBackend(params, megabytes(4), scheduler, journal);
backend->open();
// Read it back in
Batch copy;
fetchCopyOfBatch(*backend, &copy, batch);
// Canonicalize the source and destination batches
std::ranges::sort(batch, LessThan{});
std::ranges::sort(copy, LessThan{});
BEAST_EXPECT(areBatchesEqual(batch, copy));
}
}
//--------------------------------------------------------------------------
void
run() override
{
std::uint64_t const seedValue = 50;
testBackend("nudb", seedValue);
#if XRPL_ROCKSDB_AVAILABLE
testBackend("rocksdb", seedValue);
#endif
#ifdef XRPL_ENABLE_SQLITE_BACKEND_TESTS
testBackend("sqlite", seedValue);
#endif
}
};
BEAST_DEFINE_TESTSUITE(Backend, nodestore, xrpl);
} // namespace xrpl::NodeStore

View File

@@ -1,73 +0,0 @@
#include <test/nodestore/TestBase.h>
#include <xrpl/beast/unit_test/suite.h>
#include <xrpl/nodestore/NodeObject.h>
#include <xrpl/nodestore/detail/DecodedBlob.h>
#include <xrpl/nodestore/detail/EncodedBlob.h>
#include <cstdint>
#include <memory>
namespace xrpl::NodeStore {
// Tests predictable batches, and NodeObject blob encoding
//
class NodeStoreBasic_test : public TestBase
{
public:
// Make sure predictable object generation works!
void
testBatches(std::uint64_t const seedValue)
{
testcase("batch");
auto batch1 = createPredictableBatch(kNumObjectsToTest, seedValue);
auto batch2 = createPredictableBatch(kNumObjectsToTest, seedValue);
BEAST_EXPECT(areBatchesEqual(batch1, batch2));
auto batch3 = createPredictableBatch(kNumObjectsToTest, seedValue + 1);
BEAST_EXPECT(!areBatchesEqual(batch1, batch3));
}
// Checks encoding/decoding blobs
void
testBlobs(std::uint64_t const seedValue)
{
testcase("encoding");
auto batch = createPredictableBatch(kNumObjectsToTest, seedValue);
for (auto const& expected : batch)
{
EncodedBlob const encoded(expected);
DecodedBlob decoded(encoded.getKey(), encoded.getData(), encoded.getSize());
BEAST_EXPECT(decoded.wasOk());
if (decoded.wasOk())
{
std::shared_ptr<NodeObject> const object(decoded.createObject());
BEAST_EXPECT(isSame(expected, object));
}
}
}
void
run() override
{
std::uint64_t const seedValue = 50;
testBatches(seedValue);
testBlobs(seedValue);
}
};
BEAST_DEFINE_TESTSUITE(NodeStoreBasic, nodestore, xrpl);
} // namespace xrpl::NodeStore

View File

@@ -1,44 +1,21 @@
#include <test/jtx/CheckMessageLogs.h>
#include <test/jtx/Env.h>
#include <test/jtx/envconfig.h>
#include <test/nodestore/TestBase.h>
#include <test/unit_test/SuiteJournal.h>
#include <xrpld/core/Config.h>
#include <xrpl/basics/ByteUtilities.h>
#include <xrpl/beast/unit_test/suite.h>
#include <xrpl/beast/utility/Journal.h>
#include <xrpl/beast/utility/temp_dir.h>
#include <xrpl/beast/xor_shift_engine.h>
#include <xrpl/config/BasicConfig.h>
#include <xrpl/config/Constants.h>
#include <xrpl/nodestore/Database.h>
#include <xrpl/nodestore/DummyScheduler.h>
#include <xrpl/nodestore/Manager.h>
#include <xrpl/nodestore/Types.h>
#include <xrpl/protocol/SystemParameters.h>
#include <xrpl/rdb/DatabaseCon.h>
#include <algorithm>
#include <cstdint>
#include <cstring>
#include <memory>
#include <stdexcept>
#include <string>
#include <utility>
namespace xrpl::NodeStore {
namespace xrpl::node_store {
class Database_test : public TestBase
class DatabaseConfig_test : public beast::unit_test::Suite
{
test::SuiteJournal journal_;
public:
Database_test() : journal_("Database_test", *this)
{
}
void
testConfig()
{
@@ -73,8 +50,8 @@ public:
Env env = [&]() {
auto p = test::jtx::envconfig();
{
auto& section = p->section(Sections::kSqlite);
section.set(Keys::kSafetyLevel, "high");
auto& section = p->section("sqlite");
section.set("safety_level", "high");
}
p->ledgerHistory = 100'000'000;
@@ -102,8 +79,8 @@ public:
Env env = [&]() {
auto p = test::jtx::envconfig();
{
auto& section = p->section(Sections::kSqlite);
section.set(Keys::kSafetyLevel, "low");
auto& section = p->section("sqlite");
section.set("safety_level", "low");
}
p->ledgerHistory = 100'000'000;
@@ -131,10 +108,10 @@ public:
Env env = [&]() {
auto p = test::jtx::envconfig();
{
auto& section = p->section(Sections::kSqlite);
section.set(Keys::kJournalMode, "off");
section.set(Keys::kSynchronous, "extra");
section.set(Keys::kTempStore, "default");
auto& section = p->section("sqlite");
section.set("journal_mode", "off");
section.set("synchronous", "extra");
section.set("temp_store", "default");
}
return Env(
@@ -145,7 +122,7 @@ public:
}();
// No warning, even though higher risk settings were used because
// LEDGER_HISTORY is small
// ledgerHistory is small
BEAST_EXPECT(!found);
auto const s = setupDatabaseCon(env.app().config());
if (BEAST_EXPECT(s.globalPragma->size() == 3))
@@ -163,10 +140,10 @@ public:
Env env = [&]() {
auto p = test::jtx::envconfig();
{
auto& section = p->section(Sections::kSqlite);
section.set(Keys::kJournalMode, "off");
section.set(Keys::kSynchronous, "extra");
section.set(Keys::kTempStore, "default");
auto& section = p->section("sqlite");
section.set("journal_mode", "off");
section.set("synchronous", "extra");
section.set("temp_store", "default");
}
p->ledgerHistory = 50'000'000;
@@ -178,7 +155,7 @@ public:
}();
// No warning, even though higher risk settings were used because
// LEDGER_HISTORY is small
// ledgerHistory is small
BEAST_EXPECT(found);
auto const s = setupDatabaseCon(env.app().config());
if (BEAST_EXPECT(s.globalPragma->size() == 3))
@@ -199,11 +176,11 @@ public:
auto p = test::jtx::envconfig();
{
auto& section = p->section(Sections::kSqlite);
section.set(Keys::kSafetyLevel, "low");
section.set(Keys::kJournalMode, "off");
section.set(Keys::kSynchronous, "extra");
section.set(Keys::kTempStore, "default");
auto& section = p->section("sqlite");
section.set("safety_level", "low");
section.set("journal_mode", "off");
section.set("synchronous", "extra");
section.set("temp_store", "default");
}
try
@@ -230,9 +207,9 @@ public:
auto p = test::jtx::envconfig();
{
auto& section = p->section(Sections::kSqlite);
section.set(Keys::kSafetyLevel, "high");
section.set(Keys::kJournalMode, "off");
auto& section = p->section("sqlite");
section.set("safety_level", "high");
section.set("journal_mode", "off");
}
try
@@ -259,9 +236,9 @@ public:
auto p = test::jtx::envconfig();
{
auto& section = p->section(Sections::kSqlite);
section.set(Keys::kSafetyLevel, "low");
section.set(Keys::kSynchronous, "extra");
auto& section = p->section("sqlite");
section.set("safety_level", "low");
section.set("synchronous", "extra");
}
try
@@ -288,9 +265,9 @@ public:
auto p = test::jtx::envconfig();
{
auto& section = p->section(Sections::kSqlite);
section.set(Keys::kSafetyLevel, "high");
section.set(Keys::kTempStore, "default");
auto& section = p->section("sqlite");
section.set("safety_level", "high");
section.set("temp_store", "default");
}
try
@@ -317,8 +294,8 @@ public:
auto p = test::jtx::envconfig();
{
auto& section = p->section(Sections::kSqlite);
section.set(Keys::kSafetyLevel, "slow");
auto& section = p->section("sqlite");
section.set("safety_level", "slow");
}
try
@@ -345,8 +322,8 @@ public:
auto p = test::jtx::envconfig();
{
auto& section = p->section(Sections::kSqlite);
section.set(Keys::kJournalMode, "fast");
auto& section = p->section("sqlite");
section.set("journal_mode", "fast");
}
try
@@ -373,8 +350,8 @@ public:
auto p = test::jtx::envconfig();
{
auto& section = p->section(Sections::kSqlite);
section.set(Keys::kSynchronous, "instant");
auto& section = p->section("sqlite");
section.set("synchronous", "instant");
}
try
@@ -401,8 +378,8 @@ public:
auto p = test::jtx::envconfig();
{
auto& section = p->section(Sections::kSqlite);
section.set(Keys::kTempStore, "network");
auto& section = p->section("sqlite");
section.set("temp_store", "network");
}
try
@@ -436,9 +413,9 @@ public:
Env env = [&]() {
auto p = test::jtx::envconfig();
{
auto& section = p->section(Sections::kSqlite);
section.set(Keys::kPageSize, "512");
section.set(Keys::kJournalSizeLimit, "2582080");
auto& section = p->section("sqlite");
section.set("page_size", "512");
section.set("journal_size_limit", "2582080");
}
return Env(*this, std::move(p));
}();
@@ -457,8 +434,8 @@ public:
bool found = false;
auto p = test::jtx::envconfig();
{
auto& section = p->section(Sections::kSqlite);
section.set(Keys::kPageSize, "256");
auto& section = p->section("sqlite");
section.set("page_size", "256");
}
try
{
@@ -480,8 +457,8 @@ public:
bool found = false;
auto p = test::jtx::envconfig();
{
auto& section = p->section(Sections::kSqlite);
section.set(Keys::kPageSize, "131072");
auto& section = p->section("sqlite");
section.set("page_size", "131072");
}
try
{
@@ -503,8 +480,8 @@ public:
bool found = false;
auto p = test::jtx::envconfig();
{
auto& section = p->section(Sections::kSqlite);
section.set(Keys::kPageSize, "513");
auto& section = p->section("sqlite");
section.set("page_size", "513");
}
try
{
@@ -522,208 +499,13 @@ public:
}
}
//--------------------------------------------------------------------------
void
testImport(
std::string const& destBackendType,
std::string const& srcBackendType,
std::int64_t seedValue)
{
DummyScheduler scheduler;
beast::TempDir const nodeDb;
Section srcParams;
srcParams.set(Keys::kType, srcBackendType);
srcParams.set(Keys::kPath, nodeDb.path());
// Create a batch
auto batch = createPredictableBatch(kNumObjectsToTest, seedValue);
// Write to source db
{
std::unique_ptr<Database> src =
Manager::instance().makeDatabase(megabytes(4), scheduler, 2, srcParams, journal_);
storeBatch(*src, batch);
}
Batch copy;
{
// Re-open the db
std::unique_ptr<Database> src =
Manager::instance().makeDatabase(megabytes(4), scheduler, 2, srcParams, journal_);
// Set up the destination database
beast::TempDir const destDb;
Section destParams;
destParams.set(Keys::kType, destBackendType);
destParams.set(Keys::kPath, destDb.path());
std::unique_ptr<Database> dest =
Manager::instance().makeDatabase(megabytes(4), scheduler, 2, destParams, journal_);
testcase("import into '" + destBackendType + "' from '" + srcBackendType + "'");
// Do the import
dest->importDatabase(*src);
// Get the results of the import
fetchCopyOfBatch(*dest, &copy, batch);
}
// Canonicalize the source and destination batches
std::ranges::sort(batch, LessThan{});
std::ranges::sort(copy, LessThan{});
BEAST_EXPECT(areBatchesEqual(batch, copy));
}
//--------------------------------------------------------------------------
void
testNodeStore(
std::string const& type,
bool const testPersistence,
std::int64_t const seedValue,
int numObjsToTest = 2000)
{
DummyScheduler scheduler;
std::string const s = "NodeStore backend '" + type + "'";
testcase(s);
beast::TempDir const nodeDb;
Section nodeParams;
nodeParams.set(Keys::kType, type);
nodeParams.set(Keys::kPath, nodeDb.path());
beast::xor_shift_engine rng(seedValue);
// Create a batch
auto batch = createPredictableBatch(numObjsToTest, rng());
{
// Open the database
std::unique_ptr<Database> db =
Manager::instance().makeDatabase(megabytes(4), scheduler, 2, nodeParams, journal_);
// Write the batch
storeBatch(*db, batch);
{
// Read it back in
Batch copy;
fetchCopyOfBatch(*db, &copy, batch);
BEAST_EXPECT(areBatchesEqual(batch, copy));
}
{
// Reorder and read the copy again
std::shuffle(batch.begin(), batch.end(), rng);
Batch copy;
fetchCopyOfBatch(*db, &copy, batch);
BEAST_EXPECT(areBatchesEqual(batch, copy));
}
}
if (testPersistence)
{
// Re-open the database without the ephemeral DB
std::unique_ptr<Database> db =
Manager::instance().makeDatabase(megabytes(4), scheduler, 2, nodeParams, journal_);
// Read it back in
Batch copy;
fetchCopyOfBatch(*db, &copy, batch);
// Canonicalize the source and destination batches
std::ranges::sort(batch, LessThan{});
std::ranges::sort(copy, LessThan{});
BEAST_EXPECT(areBatchesEqual(batch, copy));
}
if (type == "memory")
{
// Verify default earliest ledger sequence
{
std::unique_ptr<Database> db = Manager::instance().makeDatabase(
megabytes(4), scheduler, 2, nodeParams, journal_);
BEAST_EXPECT(db->earliestLedgerSeq() == kXrpLedgerEarliestSeq);
}
// Set an invalid earliest ledger sequence
try
{
nodeParams.set(Keys::kEarliestSeq, "0");
std::unique_ptr<Database> const db = Manager::instance().makeDatabase(
megabytes(4), scheduler, 2, nodeParams, journal_);
}
catch (std::runtime_error const& e)
{
BEAST_EXPECT(std::strcmp(e.what(), "Invalid earliest_seq") == 0);
}
{
// Set a valid earliest ledger sequence
nodeParams.set(Keys::kEarliestSeq, "1");
std::unique_ptr<Database> db = Manager::instance().makeDatabase(
megabytes(4), scheduler, 2, nodeParams, journal_);
// Verify database uses the earliest ledger sequence setting
BEAST_EXPECT(db->earliestLedgerSeq() == 1);
}
// Create another database that attempts to set the value again
try
{
// Set to default earliest ledger sequence
nodeParams.set(Keys::kEarliestSeq, std::to_string(kXrpLedgerEarliestSeq));
std::unique_ptr<Database> const db2 = Manager::instance().makeDatabase(
megabytes(4), scheduler, 2, nodeParams, journal_);
}
catch (std::runtime_error const& e)
{
BEAST_EXPECT(std::strcmp(e.what(), "earliest_seq set more than once") == 0);
}
}
}
//--------------------------------------------------------------------------
void
run() override
{
std::int64_t const seedValue = 50;
testConfig();
testNodeStore("memory", false, seedValue);
// Persistent backend tests
{
testNodeStore("nudb", true, seedValue);
#if XRPL_ROCKSDB_AVAILABLE
testNodeStore("rocksdb", true, seedValue);
#endif
}
// Import tests
{
testImport("nudb", "nudb", seedValue);
#if XRPL_ROCKSDB_AVAILABLE
testImport("rocksdb", "rocksdb", seedValue);
#endif
#if XRPL_ENABLE_SQLITE_BACKEND_TESTS
testImport("sqlite", "sqlite", seedValue);
#endif
}
}
};
BEAST_DEFINE_TESTSUITE(Database, nodestore, xrpl);
BEAST_DEFINE_TESTSUITE(DatabaseConfig, nodestore, xrpl);
} // namespace xrpl::NodeStore
} // namespace xrpl::node_store

View File

@@ -1,443 +0,0 @@
#include <test/nodestore/TestBase.h>
#include <test/unit_test/SuiteJournal.h>
#include <xrpl/basics/ByteUtilities.h>
#include <xrpl/basics/Number.h>
#include <xrpl/beast/unit_test/suite.h>
#include <xrpl/beast/utility/Journal.h>
#include <xrpl/beast/utility/temp_dir.h>
#include <xrpl/config/BasicConfig.h>
#include <xrpl/config/Constants.h>
#include <xrpl/nodestore/DummyScheduler.h>
#include <xrpl/nodestore/Manager.h>
#include <xrpl/nodestore/Types.h>
#include <cstddef>
#include <exception>
#include <memory>
#include <sstream>
#include <string>
#include <utility>
#include <vector>
namespace xrpl::NodeStore {
class NuDBFactory_test : public TestBase
{
private:
// Helper function to create a Section with specified parameters
static Section
createSection(std::string const& path, std::string const& blockSize = "")
{
Section params;
params.set(Keys::kType, "nudb");
params.set(Keys::kPath, path);
if (!blockSize.empty())
params.set(Keys::kNudbBlockSize, blockSize);
return params;
}
// Helper function to create a backend and test basic functionality
bool
testBackendFunctionality(Section const& params, std::size_t expectedBlocksize)
{
try
{
DummyScheduler scheduler;
test::SuiteJournal journal("NuDBFactory_test", *this);
auto backend =
Manager::instance().makeBackend(params, megabytes(4), scheduler, journal);
if (!BEAST_EXPECT(backend))
return false;
if (!BEAST_EXPECT(backend->getBlockSize() == expectedBlocksize))
return false;
backend->open();
if (!BEAST_EXPECT(backend->isOpen()))
return false;
// Test basic store/fetch functionality
auto batch = createPredictableBatch(10, 12345);
storeBatch(*backend, batch);
Batch copy;
fetchCopyOfBatch(*backend, &copy, batch);
backend->close();
return areBatchesEqual(batch, copy);
}
catch (...)
{
return false;
}
}
// Helper function to test log messages
void
testLogMessage(Section const& params, beast::Severity level, std::string const& expectedMessage)
{
test::StreamSink sink(level);
beast::Journal const journal(sink);
DummyScheduler scheduler;
auto backend = Manager::instance().makeBackend(params, megabytes(4), scheduler, journal);
std::string const logOutput = sink.messages().str();
BEAST_EXPECT(logOutput.contains(expectedMessage));
}
// Helper function to test power of two validation
void
testPowerOfTwoValidation(std::string const& size, bool shouldWork)
{
beast::TempDir const tempDir;
auto params = createSection(tempDir.path(), size);
test::StreamSink sink(beast::Severity::Warning);
beast::Journal const journal(sink);
DummyScheduler scheduler;
auto backend = Manager::instance().makeBackend(params, megabytes(4), scheduler, journal);
std::string const logOutput = sink.messages().str();
bool const hasWarning = logOutput.contains("Invalid nudb_block_size");
BEAST_EXPECT(hasWarning == !shouldWork);
}
public:
void
testDefaultBlockSize()
{
testcase("Default block size (no nudb_block_size specified)");
beast::TempDir const tempDir;
auto params = createSection(tempDir.path());
// Should work with default 4096 block size
BEAST_EXPECT(testBackendFunctionality(params, 4096));
}
void
testValidBlockSizes()
{
testcase("Valid block sizes");
std::vector<std::size_t> const validSizes = {4096, 8192, 16384, 32768};
for (auto const& size : validSizes)
{
beast::TempDir const tempDir;
auto params = createSection(tempDir.path(), to_string(size));
BEAST_EXPECT(testBackendFunctionality(params, size));
}
// Empty value is ignored by the config parser, so uses the
// default
beast::TempDir const tempDir;
auto params = createSection(tempDir.path(), "");
BEAST_EXPECT(testBackendFunctionality(params, 4096));
}
void
testInvalidBlockSizes()
{
testcase("Invalid block sizes");
std::vector<std::string> const invalidSizes = {
"2048", // Too small
"1024", // Too small
"65536", // Too large
"131072", // Too large
"5000", // Not power of 2
"6000", // Not power of 2
"10000", // Not power of 2
"0", // Zero
"-1", // Negative
"abc", // Non-numeric
"4k", // Invalid format
"4096.5" // Decimal
};
for (auto const& size : invalidSizes)
{
beast::TempDir const tempDir;
auto params = createSection(tempDir.path(), size);
// Fails
BEAST_EXPECT(!testBackendFunctionality(params, 4096));
}
// Test whitespace cases separately since lexical_cast may handle them
std::vector<std::string> const whitespaceInvalidSizes = {
"4096 ", // Trailing space - might be handled by lexical_cast
" 4096" // Leading space - might be handled by lexical_cast
};
for (auto const& size : whitespaceInvalidSizes)
{
beast::TempDir const tempDir;
auto params = createSection(tempDir.path(), size);
// Fails
BEAST_EXPECT(!testBackendFunctionality(params, 4096));
}
}
void
testLogMessages()
{
testcase("Log message verification");
// Test valid custom block size logging
{
beast::TempDir const tempDir;
auto params = createSection(tempDir.path(), "8192");
testLogMessage(params, beast::Severity::Info, "Using custom NuDB block size: 8192");
}
// Test invalid block size failure
{
beast::TempDir const tempDir;
auto params = createSection(tempDir.path(), "5000");
test::StreamSink sink(beast::Severity::Warning);
beast::Journal const journal(sink);
DummyScheduler scheduler;
try
{
auto backend =
Manager::instance().makeBackend(params, megabytes(4), scheduler, journal);
fail();
}
catch (std::exception const& e)
{
std::string const logOutput{e.what()};
BEAST_EXPECT(logOutput.contains("Invalid nudb_block_size: 5000"));
BEAST_EXPECT(logOutput.contains("Must be power of 2 between 4096 and 32768"));
}
}
// Test non-numeric value failure
{
beast::TempDir const tempDir;
auto params = createSection(tempDir.path(), "invalid");
test::StreamSink sink(beast::Severity::Warning);
beast::Journal const journal(sink);
DummyScheduler scheduler;
try
{
auto backend =
Manager::instance().makeBackend(params, megabytes(4), scheduler, journal);
fail();
}
catch (std::exception const& e)
{
std::string const logOutput{e.what()};
BEAST_EXPECT(logOutput.contains("Invalid nudb_block_size value: invalid"));
}
}
}
void
testPowerOfTwoValidation()
{
testcase("Power of 2 validation logic");
// Test edge cases around valid range
std::vector<std::pair<std::string, bool>> const testCases = {
{"4095", false}, // Just below minimum
{"4096", true}, // Minimum valid
{"4097", false}, // Just above minimum, not power of 2
{"8192", true}, // Valid power of 2
{"8193", false}, // Just above valid power of 2
{"16384", true}, // Valid power of 2
{"32768", true}, // Maximum valid
{"32769", false}, // Just above maximum
{"65536", false} // Power of 2 but too large
};
for (auto const& [size, shouldWork] : testCases)
{
beast::TempDir const tempDir;
auto params = createSection(tempDir.path(), size);
// We test the validation logic by catching exceptions for invalid
// values
test::StreamSink sink(beast::Severity::Warning);
beast::Journal const journal(sink);
DummyScheduler scheduler;
try
{
auto backend =
Manager::instance().makeBackend(params, megabytes(4), scheduler, journal);
BEAST_EXPECT(shouldWork);
}
catch (std::exception const& e)
{
std::string const logOutput{e.what()};
BEAST_EXPECT(logOutput.contains("Invalid nudb_block_size"));
}
}
}
void
testBothConstructorVariants()
{
testcase("Both constructor variants work with custom block size");
beast::TempDir const tempDir;
auto params = createSection(tempDir.path(), "16384");
DummyScheduler scheduler;
test::SuiteJournal journal("NuDBFactory_test", *this);
// Test first constructor (without nudb::context)
{
auto backend1 =
Manager::instance().makeBackend(params, megabytes(4), scheduler, journal);
BEAST_EXPECT(backend1 != nullptr);
BEAST_EXPECT(testBackendFunctionality(params, 16384));
}
// Test second constructor (with nudb::context)
// Note: This would require access to nudb::context, which might not be
// easily testable without more complex setup. For now, we test that
// the factory can create backends with the first constructor.
}
void
testConfigurationParsing()
{
testcase("Configuration parsing edge cases");
// Test that whitespace is handled correctly
std::vector<std::string> const validFormats = {
"8192" // Basic valid format
};
// Test whitespace handling separately since lexical_cast behavior may
// vary
std::vector<std::string> const whitespaceFormats = {
" 8192", // Leading space - may or may not be handled by
// lexical_cast
"8192 " // Trailing space - may or may not be handled by
// lexical_cast
};
// Test basic valid format
for (auto const& format : validFormats)
{
beast::TempDir const tempDir;
auto params = createSection(tempDir.path(), format);
test::StreamSink sink(beast::Severity::Info);
beast::Journal const journal(sink);
DummyScheduler scheduler;
auto backend =
Manager::instance().makeBackend(params, megabytes(4), scheduler, journal);
// Should log success message for valid values
std::string const logOutput = sink.messages().str();
bool const hasSuccessMessage = logOutput.contains("Using custom NuDB block size");
BEAST_EXPECT(hasSuccessMessage);
}
// Test whitespace formats - these should work if lexical_cast handles
// them
for (auto const& format : whitespaceFormats)
{
beast::TempDir const tempDir;
auto params = createSection(tempDir.path(), format);
// Use a lower threshold to capture both info and warning messages
test::StreamSink sink(beast::Severity::Debug);
beast::Journal const journal(sink);
DummyScheduler scheduler;
try
{
auto backend =
Manager::instance().makeBackend(params, megabytes(4), scheduler, journal);
fail();
}
catch (...)
{
// Fails
BEAST_EXPECT(!testBackendFunctionality(params, 8192));
}
}
}
void
testDataPersistence()
{
testcase("Data persistence with different block sizes");
std::vector<std::string> const blockSizes = {"4096", "8192", "16384", "32768"};
for (auto const& size : blockSizes)
{
beast::TempDir const tempDir;
auto params = createSection(tempDir.path(), size);
DummyScheduler scheduler;
test::SuiteJournal journal("NuDBFactory_test", *this);
// Create test data
auto batch = createPredictableBatch(50, 54321);
// Store data
{
auto backend =
Manager::instance().makeBackend(params, megabytes(4), scheduler, journal);
backend->open();
storeBatch(*backend, batch);
backend->close();
}
// Retrieve data in new backend instance
{
auto backend =
Manager::instance().makeBackend(params, megabytes(4), scheduler, journal);
backend->open();
Batch copy;
fetchCopyOfBatch(*backend, &copy, batch);
BEAST_EXPECT(areBatchesEqual(batch, copy));
backend->close();
}
}
}
void
run() override
{
testDefaultBlockSize();
testValidBlockSizes();
testInvalidBlockSizes();
testLogMessages();
testPowerOfTwoValidation();
testBothConstructorVariants();
testConfigurationParsing();
testDataPersistence();
}
};
BEAST_DEFINE_TESTSUITE(NuDBFactory, xrpl_core, xrpl);
} // namespace xrpl::NodeStore

View File

@@ -1,202 +0,0 @@
#pragma once
#include <xrpl/basics/Blob.h>
#include <xrpl/basics/base_uint.h>
#include <xrpl/basics/random.h>
#include <xrpl/beast/unit_test/suite.h>
#include <xrpl/beast/utility/rngfill.h>
#include <xrpl/beast/xor_shift_engine.h>
#include <xrpl/nodestore/Backend.h>
#include <xrpl/nodestore/Database.h>
#include <xrpl/nodestore/NodeObject.h>
#include <xrpl/nodestore/Types.h>
#include <boost/algorithm/string.hpp>
#include <cstddef>
#include <cstdint>
#include <memory>
#include <utility>
namespace xrpl::NodeStore {
/**
* Binary function that satisfies the strict-weak-ordering requirement.
*
* This compares the hashes of both objects and returns true if
* the first hash is considered to go before the second.
*
* @see std::sort
*/
struct LessThan
{
bool
operator()(std::shared_ptr<NodeObject> const& lhs, std::shared_ptr<NodeObject> const& rhs)
const noexcept
{
return lhs->getHash() < rhs->getHash();
}
};
/**
* Returns `true` if objects are identical.
*/
inline bool
isSame(std::shared_ptr<NodeObject> const& lhs, std::shared_ptr<NodeObject> const& rhs)
{
return (lhs->getType() == rhs->getType()) && (lhs->getHash() == rhs->getHash()) &&
(lhs->getData() == rhs->getData());
}
// Some common code for the unit tests
//
class TestBase : public beast::unit_test::Suite
{
public:
// Tunable parameters
//
static std::size_t const kMinPayloadBytes = 1;
static std::size_t const kMaxPayloadBytes = 2000;
static int const kNumObjectsToTest = 2000;
public:
// Create a predictable batch of objects
static Batch
createPredictableBatch(int numObjects, std::uint64_t seed)
{
Batch batch;
batch.reserve(numObjects);
beast::xor_shift_engine rng(seed);
for (int i = 0; i < numObjects; ++i)
{
NodeObjectType const type = [&] {
switch (randInt(rng, 3))
{
case 0:
return NodeObjectType::Ledger;
case 1:
return NodeObjectType::AccountNode;
case 2:
return NodeObjectType::TransactionNode;
case 3:
default:
return NodeObjectType::Unknown;
}
}();
uint256 hash;
beast::rngfill(hash.begin(), hash.size(), rng);
Blob blob(randInt(rng, kMinPayloadBytes, kMaxPayloadBytes));
beast::rngfill(blob.data(), blob.size(), rng);
batch.push_back(NodeObject::createObject(type, std::move(blob), hash));
}
return batch;
}
// Compare two batches for equality
static bool
areBatchesEqual(Batch const& lhs, Batch const& rhs)
{
bool result = true;
if (lhs.size() == rhs.size())
{
for (int i = 0; i < lhs.size(); ++i)
{
if (!isSame(lhs[i], rhs[i]))
{
result = false;
break;
}
}
}
else
{
result = false;
}
return result;
}
// Store a batch in a backend
static void
storeBatch(Backend& backend, Batch const& batch)
{
for (auto const& object : batch)
{
backend.store(object);
}
}
// Get a copy of a batch in a backend
void
fetchCopyOfBatch(Backend& backend, Batch* pCopy, Batch const& batch)
{
pCopy->clear();
pCopy->reserve(batch.size());
for (auto const& expected : batch)
{
std::shared_ptr<NodeObject> object;
Status const status = backend.fetch(expected->getHash(), &object);
BEAST_EXPECT(status == Status::Ok);
if (status == Status::Ok)
{
BEAST_EXPECT(object != nullptr);
pCopy->push_back(object);
}
}
}
void
fetchMissing(Backend& backend, Batch const& batch)
{
for (auto const& expected : batch)
{
std::shared_ptr<NodeObject> object;
Status const status = backend.fetch(expected->getHash(), &object);
BEAST_EXPECT(status == Status::NotFound);
}
}
// Store all objects in a batch
static void
storeBatch(Database& db, Batch const& batch)
{
for (auto const& object : batch)
{
Blob data(object->getData());
db.store(object->getType(), std::move(data), object->getHash(), db.earliestLedgerSeq());
}
}
// Fetch all the hashes in one batch, into another batch.
static void
fetchCopyOfBatch(Database& db, Batch* pCopy, Batch const& batch)
{
pCopy->clear();
pCopy->reserve(batch.size());
for (auto const& expected : batch)
{
std::shared_ptr<NodeObject> const object = db.fetchNodeObject(expected->getHash(), 0);
if (object != nullptr)
pCopy->push_back(object);
}
}
};
} // namespace xrpl::NodeStore

View File

@@ -195,7 +195,7 @@ fmtdur(std::chrono::duration<Period, Rep> const& d)
} // namespace detail
namespace NodeStore {
namespace node_store {
//------------------------------------------------------------------------------
@@ -552,5 +552,5 @@ BEAST_DEFINE_TESTSUITE_MANUAL(import, nodestore, xrpl);
//------------------------------------------------------------------------------
} // namespace NodeStore
} // namespace node_store
} // namespace xrpl

View File

@@ -1,57 +0,0 @@
#include <xrpl/beast/unit_test/suite.h>
#include <xrpl/nodestore/detail/varint.h>
#include <array>
#include <cstddef>
#include <cstdint>
#include <vector>
namespace xrpl::NodeStore::tests {
class varint_test : public beast::unit_test::Suite
{
public:
void
testVarints(std::vector<std::size_t> vv)
{
testcase("encode, decode");
for (auto const v : vv)
{
std::array<std::uint8_t, varint_traits<std::size_t>::kMax> vi{};
auto const n0 = writeVarint(vi.data(), v);
expect(n0 > 0, "write error");
expect(n0 == sizeVarint(v), "size error");
std::size_t v1 = 0;
auto const n1 = readVarint(vi.data(), n0, v1);
expect(n1 == n0, "read error");
expect(v == v1, "wrong value");
}
}
void
run() override
{
testVarints(
{0,
1,
2,
126,
127,
128,
253,
254,
255,
16127,
16128,
16129,
0xff,
0xffff,
0xffffffff,
0xffffffffffffUL,
0xffffffffffffffffUL});
}
};
BEAST_DEFINE_TESTSUITE(varint, nodestore, xrpl);
} // namespace xrpl::NodeStore::tests

View File

@@ -1,41 +0,0 @@
#include <xrpl/beast/unit_test/suite.h>
#include <xrpl/protocol/ApiVersion.h>
namespace xrpl::test {
struct ApiVersion_test : beast::unit_test::Suite
{
void
run() override
{
{
testcase("API versions invariants");
static_assert(RPC::kApiMinimumSupportedVersion <= RPC::kApiMaximumSupportedVersion);
static_assert(RPC::kApiMinimumSupportedVersion <= RPC::kApiMaximumValidVersion);
static_assert(RPC::kApiMaximumSupportedVersion <= RPC::kApiMaximumValidVersion);
static_assert(RPC::kApiBetaVersion <= RPC::kApiMaximumValidVersion);
BEAST_EXPECT(true);
}
{
// Update when we change versions
testcase("API versions");
static_assert(RPC::kApiMinimumSupportedVersion >= 1);
static_assert(RPC::kApiMinimumSupportedVersion < 2);
static_assert(RPC::kApiMaximumSupportedVersion >= 2);
static_assert(RPC::kApiMaximumSupportedVersion < 3);
static_assert(RPC::kApiMaximumValidVersion >= 3);
static_assert(RPC::kApiMaximumValidVersion < 4);
static_assert(RPC::kApiBetaVersion >= 3);
static_assert(RPC::kApiBetaVersion < 4);
BEAST_EXPECT(true);
}
}
};
BEAST_DEFINE_TESTSUITE(ApiVersion, protocol, xrpl);
} // namespace xrpl::test

View File

@@ -1,52 +0,0 @@
#include <xrpl/beast/unit_test/suite.h>
#include <xrpl/protocol/Serializer.h>
#include <cstdint>
#include <initializer_list>
#include <limits>
namespace xrpl {
struct Serializer_test : public beast::unit_test::Suite
{
void
run() override
{
{
std::initializer_list<std::int32_t> const values = {
std::numeric_limits<std::int32_t>::min(),
-1,
0,
1,
std::numeric_limits<std::int32_t>::max()};
for (std::int32_t const value : values)
{
Serializer s;
s.add32(value);
BEAST_EXPECT(s.size() == 4);
SerialIter sit(s.slice());
BEAST_EXPECT(sit.geti32() == value);
}
}
{
std::initializer_list<std::int64_t> const values = {
std::numeric_limits<std::int64_t>::min(),
-1,
0,
1,
std::numeric_limits<std::int64_t>::max()};
for (std::int64_t const value : values)
{
Serializer s;
s.add64(value);
BEAST_EXPECT(s.size() == 8);
SerialIter sit(s.slice());
BEAST_EXPECT(sit.geti64() == value);
}
}
}
};
BEAST_DEFINE_TESTSUITE(Serializer, protocol, xrpl);
} // namespace xrpl

View File

@@ -558,10 +558,12 @@ class ServerStatus_test : public beast::unit_test::Suite, public beast::test::En
using namespace test::jtx;
using namespace boost::asio;
using namespace boost::beast::http;
Env env{*this, envconfig([&](std::unique_ptr<Config> cfg) {
// Run the server with a single io thread so disconnectClient() below
// can deterministically drain the server's io_context (see its docs).
Env env{*this, singleThreadIo(envconfig([&](std::unique_ptr<Config> cfg) {
(*cfg)[Sections::kPortRpc].set(Keys::kLimit, std::to_string(limit));
return cfg;
})};
}))};
auto const section = env.app().config().section(Sections::kPortRpc);
// NOLINTBEGIN(bugprone-unchecked-optional-access)
@@ -580,16 +582,27 @@ class ServerStatus_test : public beast::unit_test::Suite, public beast::test::En
BEAST_EXPECT(!ec);
std::vector<std::pair<ip::tcp::socket, boost::beast::multi_buffer>> clients;
int connectionCount{1}; // starts at 1 because the Env already has one
// for JSONRPCCLient
// for nonzero limits, go one past the limit, although failures happen
// at the limit, so this really leads to the last two clients failing.
// for zero limit, pick an arbitrary nonzero number of clients - all
// should connect fine.
// Env owns a persistent JSON-RPC HTTP client connection to port_rpc as
// part of startup, which counts against this port's connection limit.
// This test wants a known starting occupancy of zero, so for nonzero
// limits it deterministically drops that hidden client and waits for
// the server to register the disconnect before opening its own clients.
//
// Starting from zero is important because the port limit rejects once
// the incremented connection count reaches the configured limit. With a
// zero baseline and N = limit + 1 test-owned clients, exactly the last
// two requests should be rejected.
if (limit != 0)
BEAST_EXPECT(env.disconnectClient());
// For nonzero limits, go one past the limit. The port rejects at the
// limit, not only above it, so this yields the last two clients
// failing. For zero limit, pick an arbitrary nonzero number of clients
// and expect them all to succeed.
int const testTo = (limit == 0) ? 50 : limit + 1;
while (connectionCount < testTo)
while (static_cast<int>(clients.size()) < testTo)
{
clients.emplace_back(ip::tcp::socket{ios}, boost::beast::multi_buffer{});
async_connect(clients.back().first, it, yield[ec]);
@@ -597,19 +610,24 @@ class ServerStatus_test : public beast::unit_test::Suite, public beast::test::En
auto req = makeHTTPRequest(ip, port, to_string(jr), {});
async_write(clients.back().first, req, yield[ec]);
BEAST_EXPECT(!ec);
++connectionCount;
}
int readCount = 0;
int successfulReads = 0;
for (auto& [soc, buf] : clients)
{
boost::beast::http::response<boost::beast::http::string_body> resp;
async_read(soc, buf, resp, yield[ec]);
++readCount;
// expect the reads to fail for the clients that connected at or
// above the limit. If limit is 0, all reads should succeed
BEAST_EXPECT((limit == 0 || readCount < limit - 1) ? (!ec) : bool(ec));
if (!ec)
++successfulReads;
}
// This test cares about the exact number of accepted requests, not which
// specific client observed the rejection. With a zero baseline (the
// hidden Env client dropped above), the server accepts until the
// connection count reaches the limit: all clients for limit 0, else
// limit - 1 of the limit + 1 clients (the last two are rejected).
int const expectedReads = (limit == 0) ? static_cast<int>(clients.size()) : limit - 1;
BEAST_EXPECT(successfulReads == expectedReads);
}
void

View File

@@ -27,10 +27,13 @@ target_link_libraries(xrpl_tests PRIVATE GTest::gtest GTest::gmock xrpl.libxrpl)
# supported on Windows.
set(test_modules
basics
beast
consensus
crypto
json
nodestore
peerfinder
protocol
resource
shamap
tx

View File

@@ -0,0 +1,333 @@
#include <xrpl/beast/core/SemanticVersion.h>
#include <gtest/gtest.h>
#include <algorithm>
#include <array>
#include <cctype>
#include <locale>
#include <string>
#include <string_view>
#include <vector>
namespace beast {
namespace {
using IdentifierList = SemanticVersion::IdentifierList;
// Version strings are not valid C++ identifiers, so squash their punctuation to
// turn one into a gtest parameter name.
std::string
identifierFor(std::string_view version)
{
std::string name{version};
std::ranges::replace_if(
name, [](char c) { return !std::isalnum(c, std::locale::classic()); }, '_');
if (!name.empty() && std::isdigit(name.front(), std::locale::classic()))
name.insert(0, "v_");
return name;
}
// Pre-release and metadata suffixes, each applied to a "major.minor.patch" base.
// The valid ones leave a well-formed base well-formed; the invalid ones make any
// base malformed.
constexpr auto kValidPreRelease =
std::to_array<std::string_view>({"", "-1", "-a", "-a1", "-a1.b1", "-ab.cd", "--"});
constexpr auto kInvalidPreRelease =
std::to_array<std::string_view>({"+", "!", "-", "-!", "-.", "-a.!", "-0.a"});
constexpr auto kValidMetaData = std::to_array<std::string_view>({"", "+a", "+1", "+a.b", "+ab.cd"});
constexpr auto kInvalidMetaData =
std::to_array<std::string_view>({"!", "+", "++", "+!", "+.", "+a.!"});
// Assembles base + preRelease + metaData and checks whether it parses. A version
// we accept must also round-trip through print().
void
expectParse(
std::string_view base,
std::string_view preRelease,
std::string_view metaData,
bool shouldPass)
{
auto const input = std::string{base}.append(preRelease).append(metaData);
SCOPED_TRACE(::testing::Message() << '"' << input << '"');
SemanticVersion v;
if (shouldPass)
{
EXPECT_TRUE(v.parse(input));
EXPECT_EQ(v.print(), input);
}
else
{
EXPECT_FALSE(v.parse(input));
}
}
struct ParseCase
{
std::string_view testName;
std::string_view base;
bool shouldPass;
};
std::string
parseCaseName(::testing::TestParamInfo<ParseCase> const& info)
{
return std::string{info.param.testName};
}
constexpr auto kParseCases = std::to_array<ParseCase>({
{.testName = "zeroes", .base = "0.0.0", .shouldPass = true},
{.testName = "simple", .base = "1.2.3", .shouldPass = true},
{.testName = "max_int", .base = "2147483647.2147483647.2147483647", .shouldPass = true},
// negative values
{.testName = "negative_major", .base = "-1.2.3", .shouldPass = false},
{.testName = "negative_minor", .base = "1.-2.3", .shouldPass = false},
{.testName = "negative_patch", .base = "1.2.-3", .shouldPass = false},
// missing parts
{.testName = "empty", .base = "", .shouldPass = false},
{.testName = "major_only", .base = "1", .shouldPass = false},
{.testName = "major_then_dot", .base = "1.", .shouldPass = false},
{.testName = "major_and_minor", .base = "1.2", .shouldPass = false},
{.testName = "major_minor_then_dot", .base = "1.2.", .shouldPass = false},
{.testName = "missing_major", .base = ".2.3", .shouldPass = false},
// whitespace
{.testName = "leading_space", .base = " 1.2.3", .shouldPass = false},
{.testName = "space_after_major", .base = "1 .2.3", .shouldPass = false},
{.testName = "space_after_minor", .base = "1.2 .3", .shouldPass = false},
{.testName = "trailing_space", .base = "1.2.3 ", .shouldPass = false},
// leading zeroes
{.testName = "leading_zero_in_major", .base = "01.2.3", .shouldPass = false},
{.testName = "leading_zero_in_minor", .base = "1.02.3", .shouldPass = false},
{.testName = "leading_zero_in_patch", .base = "1.2.03", .shouldPass = false},
});
struct ValuesCase
{
std::string_view testName;
std::string_view input;
int majorVersion;
int minorVersion;
int patchVersion;
IdentifierList preReleaseIdentifiers{}; // NOLINT(readability-redundant-member-init)
IdentifierList metaData{}; // NOLINT(readability-redundant-member-init)
};
std::string
valuesCaseName(::testing::TestParamInfo<ValuesCase> const& info)
{
return std::string{info.param.testName};
}
std::vector<ValuesCase> const kValuesCases{
{
.testName = "zero_major",
.input = "0.1.2",
.majorVersion = 0,
.minorVersion = 1,
.patchVersion = 2,
},
{
.testName = "simple",
.input = "1.2.3",
.majorVersion = 1,
.minorVersion = 2,
.patchVersion = 3,
},
{
.testName = "one_pre_release_identifier",
.input = "1.2.3-rc1",
.majorVersion = 1,
.minorVersion = 2,
.patchVersion = 3,
.preReleaseIdentifiers = {"rc1"},
},
{
.testName = "two_pre_release_identifiers",
.input = "1.2.3-rc1.debug",
.majorVersion = 1,
.minorVersion = 2,
.patchVersion = 3,
.preReleaseIdentifiers = {"rc1", "debug"},
},
{
.testName = "three_pre_release_identifiers",
.input = "1.2.3-rc1.debug.asm",
.majorVersion = 1,
.minorVersion = 2,
.patchVersion = 3,
.preReleaseIdentifiers = {"rc1", "debug", "asm"},
},
{
.testName = "one_metadata_identifier",
.input = "1.2.3+full",
.majorVersion = 1,
.minorVersion = 2,
.patchVersion = 3,
.metaData = {"full"},
},
{
.testName = "two_metadata_identifiers",
.input = "1.2.3+full.prod",
.majorVersion = 1,
.minorVersion = 2,
.patchVersion = 3,
.metaData = {"full", "prod"},
},
{
.testName = "three_metadata_identifiers",
.input = "1.2.3+full.prod.x86",
.majorVersion = 1,
.minorVersion = 2,
.patchVersion = 3,
.metaData = {"full", "prod", "x86"},
},
{
.testName = "pre_release_and_metadata",
.input = "1.2.3-rc1.debug.asm+full.prod.x86",
.majorVersion = 1,
.minorVersion = 2,
.patchVersion = 3,
.preReleaseIdentifiers = {"rc1", "debug", "asm"},
.metaData = {"full", "prod", "x86"},
},
};
struct OrderCase
{
std::string_view lesser;
std::string_view greater;
};
std::string
orderCaseName(::testing::TestParamInfo<OrderCase> const& info)
{
return identifierFor(info.param.lesser) + "_below_" + identifierFor(info.param.greater);
}
constexpr auto kOrderCases = std::to_array<OrderCase>({
{.lesser = "1.0.0-alpha", .greater = "1.0.0-alpha.1"},
{.lesser = "1.0.0-alpha.1", .greater = "1.0.0-alpha.beta"},
{.lesser = "1.0.0-alpha.beta", .greater = "1.0.0-beta"},
{.lesser = "1.0.0-beta", .greater = "1.0.0-beta.2"},
{.lesser = "1.0.0-beta.2", .greater = "1.0.0-beta.11"},
{.lesser = "1.0.0-beta.11", .greater = "1.0.0-rc.1"},
{.lesser = "1.0.0-rc.1", .greater = "1.0.0"},
{.lesser = "0.9.9", .greater = "1.0.0"},
});
} // namespace
class SemanticVersionParse : public ::testing::TestWithParam<ParseCase>
{
};
// Exercises the base string on its own and with every combination of appended
// pre-release identifiers and metadata.
TEST_P(SemanticVersionParse, pre_release_and_metadata_combinations)
{
auto const& [testName, base, shouldPass] = GetParam();
for (auto const preRelease : kValidPreRelease)
{
for (auto const metaData : kValidMetaData)
expectParse(base, preRelease, metaData, shouldPass);
for (auto const metaData : kInvalidMetaData)
expectParse(base, preRelease, metaData, false);
}
// A malformed pre-release section poisons the whole string, whatever
// metadata follows it.
for (auto const preRelease : kInvalidPreRelease)
{
for (auto const metaData : kValidMetaData)
expectParse(base, preRelease, metaData, false);
for (auto const metaData : kInvalidMetaData)
expectParse(base, preRelease, metaData, false);
}
}
INSTANTIATE_TEST_SUITE_P(
Inputs,
SemanticVersionParse,
::testing::ValuesIn(kParseCases),
parseCaseName);
class SemanticVersionValues : public ::testing::TestWithParam<ValuesCase>
{
};
TEST_P(SemanticVersionValues, decomposes_into_components)
{
auto const& expected = GetParam();
SemanticVersion v;
EXPECT_TRUE(v.parse(expected.input));
EXPECT_EQ(v.majorVersion, expected.majorVersion);
EXPECT_EQ(v.minorVersion, expected.minorVersion);
EXPECT_EQ(v.patchVersion, expected.patchVersion);
EXPECT_EQ(v.preReleaseIdentifiers, expected.preReleaseIdentifiers);
EXPECT_EQ(v.metaData, expected.metaData);
}
INSTANTIATE_TEST_SUITE_P(
Inputs,
SemanticVersionValues,
::testing::ValuesIn(kValuesCases),
valuesCaseName);
class SemanticVersionOrder : public ::testing::TestWithParam<OrderCase>
{
};
TEST_P(SemanticVersionOrder, lesser_precedes_greater)
{
auto const& [lesser, greater] = GetParam();
// Metadata takes no part in precedence, so attaching it to either side must
// leave the ordering untouched.
static constexpr auto kMetaData = std::to_array<std::string_view>({"", "+meta"});
for (auto const lesserMetaData : kMetaData)
{
for (auto const greaterMetaData : kMetaData)
{
auto const lesserInput = std::string{lesser}.append(lesserMetaData);
auto const greaterInput = std::string{greater}.append(greaterMetaData);
SCOPED_TRACE(
::testing::Message() << '"' << lesserInput << "\" < \"" << greaterInput << '"');
SemanticVersion lesserVersion;
SemanticVersion greaterVersion;
EXPECT_TRUE(lesserVersion.parse(lesserInput));
EXPECT_TRUE(greaterVersion.parse(greaterInput));
EXPECT_EQ(compare(lesserVersion, lesserVersion), 0);
EXPECT_EQ(compare(greaterVersion, greaterVersion), 0);
EXPECT_LT(compare(lesserVersion, greaterVersion), 0);
EXPECT_GT(compare(greaterVersion, lesserVersion), 0);
EXPECT_LT(lesserVersion, greaterVersion);
EXPECT_GT(greaterVersion, lesserVersion);
EXPECT_EQ(lesserVersion, lesserVersion);
EXPECT_EQ(greaterVersion, greaterVersion);
}
}
}
INSTANTIATE_TEST_SUITE_P(
Pairs,
SemanticVersionOrder,
::testing::ValuesIn(kOrderCases),
orderCaseName);
} // namespace beast

View File

@@ -0,0 +1,92 @@
#include <xrpl/beast/utility/Zero.h>
#include <gtest/gtest.h>
namespace beast {
struct AdlTester
{
};
int
signum(AdlTester)
{
return 0;
}
namespace inner_adl_test {
struct AdlTester2
{
};
int
signum(AdlTester2)
{
return 0;
}
} // namespace inner_adl_test
namespace {
struct IntegerWrapper
{
int value;
IntegerWrapper(int v) : value(v)
{
}
[[nodiscard]] int
signum() const
{
return value;
}
};
void
testLhsZero(IntegerWrapper x)
{
EXPECT_EQ(x >= kZero, x.signum() >= 0);
EXPECT_EQ(x > kZero, x.signum() > 0);
EXPECT_EQ(x == kZero, x.signum() == 0);
EXPECT_EQ(x != kZero, x.signum() != 0);
EXPECT_EQ(x < kZero, x.signum() < 0);
EXPECT_EQ(x <= kZero, x.signum() <= 0);
}
void
testRhsZero(IntegerWrapper x)
{
EXPECT_EQ(kZero >= x, 0 >= x.signum());
EXPECT_EQ(kZero > x, 0 > x.signum());
EXPECT_EQ(kZero == x, 0 == x.signum());
EXPECT_EQ(kZero != x, 0 != x.signum());
EXPECT_EQ(kZero < x, 0 < x.signum());
EXPECT_EQ(kZero <= x, 0 <= x.signum());
}
} // namespace
TEST(Zero, lhs)
{
testLhsZero(-7);
testLhsZero(0);
testLhsZero(32);
}
TEST(Zero, rhs)
{
testRhsZero(-4);
testRhsZero(0);
testRhsZero(64);
}
TEST(Zero, adl)
{
EXPECT_TRUE(AdlTester{} == kZero);
EXPECT_TRUE(inner_adl_test::AdlTester2{} == kZero);
}
} // namespace beast

View File

@@ -0,0 +1,50 @@
#pragma once
#include <xrpl/beast/utility/Journal.h>
#include <mutex>
#include <sstream>
#include <string>
namespace xrpl::test {
class CaptureSink : public beast::Journal::Sink
{
mutable std::mutex mutex_;
std::stringstream strm_;
public:
explicit CaptureSink(beast::Severity threshold = beast::Severity::Debug)
: Sink{threshold, false}
{
}
void
write(beast::Severity level, std::string const& text) override
{
if (level < threshold())
return;
writeAlways(level, text);
}
void
writeAlways(beast::Severity /*level*/, std::string const& text) override
{
// Journal sinks may be written to concurrently (e.g. from a backend's background workers),
// so serialize access to strm_. write() funnels into writeAlways(), so the lock lives here
// only: locking in both would self-deadlock on this non-recursive mutex.
std::scoped_lock const lock(mutex_);
strm_ << text << '\n';
}
[[nodiscard]] std::string
messages() const
{
// Returns a snapshot of the captured output. Takes the lock so the read is safe even if a
// writer is still active.
std::scoped_lock const lock(mutex_);
return strm_.str();
}
};
} // namespace xrpl::test

View File

@@ -29,11 +29,11 @@ namespace xrpl::test {
class TestFamily : public Family
{
private:
std::unique_ptr<NodeStore::Database> db_;
std::unique_ptr<node_store::Database> db_;
TestStopwatch clock_;
std::shared_ptr<FullBelowCache> fbCache_;
std::shared_ptr<TreeNodeCache> tnCache_;
NodeStore::DummyScheduler scheduler_;
node_store::DummyScheduler scheduler_;
beast::Journal j_;
public:
@@ -51,16 +51,16 @@ public:
Section config;
config.set(Keys::kType, "memory");
config.set(Keys::kPath, "TestFamily");
db_ = NodeStore::Manager::instance().makeDatabase(megabytes(4), scheduler_, 1, config, j);
db_ = node_store::Manager::instance().makeDatabase(megabytes(4), scheduler_, 1, config, j);
}
NodeStore::Database&
node_store::Database&
db() override
{
return *db_;
}
[[nodiscard]] NodeStore::Database const&
[[nodiscard]] node_store::Database const&
db() const override
{
return *db_;

View File

@@ -220,7 +220,7 @@ public:
}
// Storage services
NodeStore::Database&
node_store::Database&
getNodeStore() override
{
throw std::logic_error("TestServiceRegistry::getNodeStore() not implemented");

View File

@@ -0,0 +1,182 @@
#include <xrpl/nodestore/Backend.h>
#include <xrpl/basics/ByteUtilities.h>
#include <xrpl/beast/utility/Journal.h>
#include <xrpl/beast/utility/temp_dir.h>
#include <xrpl/beast/xor_shift_engine.h>
#include <xrpl/config/BasicConfig.h>
#include <xrpl/nodestore/DummyScheduler.h>
#include <xrpl/nodestore/Manager.h>
#include <xrpl/nodestore/NodeObject.h>
#include <xrpl/nodestore/Types.h>
#include <gtest/gtest.h>
#include <helpers/TestSink.h>
#include <nodestore/TestBase.h>
#include <algorithm>
#include <atomic>
#include <cstddef>
#include <memory>
#include <ranges>
#include <string>
#include <thread>
#include <vector>
namespace xrpl::node_store {
namespace {
std::vector<std::string>
backendTypes()
{
std::vector<std::string> types{"nudb"};
#if XRPL_ROCKSDB_AVAILABLE
types.emplace_back("rocksdb");
#endif
#ifdef XRPL_ENABLE_SQLITE_BACKEND_TESTS
types.emplace_back("sqlite");
#endif
return types;
}
// Run work(i) for every i in [0, n) spread across numThreads threads, handing
// out indices via a shared atomic counter (mirrors the old Timing_test
// parallel-for so the N items are partitioned, not duplicated).
template <class Work>
void
parallelFor(std::size_t n, std::size_t numThreads, Work work)
{
std::atomic<std::size_t> next{0};
auto const runner = [&] {
for (std::size_t i = next++; i < n; i = next++)
work(i);
};
auto threads = std::views::iota(std::size_t{0}, numThreads) |
std::views::transform([&](std::size_t) { return std::thread{runner}; }) |
std::ranges::to<std::vector>();
std::ranges::for_each(threads, &std::thread::join);
}
} // namespace
class BackendTypeTest : public ::testing::TestWithParam<std::string>
{
protected:
void
SetUp() override
{
params_.set("type", GetParam());
params_.set("path", tempDir_.path());
beast::xor_shift_engine rng(kSeedValue);
batch_ = createPredictableBatch(kNumObjects, rng());
}
std::unique_ptr<Backend>
makeOpenBackend()
{
auto backend = Manager::instance().makeBackend(params_, megabytes(4), scheduler_, journal_);
backend->open();
return backend;
}
DummyScheduler scheduler_;
beast::TempDir const tempDir_;
beast::Journal const journal_{TestSink::instance()};
Section params_;
Batch batch_;
};
TEST_P(BackendTypeTest, store_and_fetch)
{
auto backend = makeOpenBackend();
storeBatch(*backend, batch_);
{
SCOPED_TRACE("read in original order");
auto const copy = fetchCopyOfBatch(*backend, batch_);
EXPECT_EQ(batch_, copy);
}
{
SCOPED_TRACE("read in shuffled order");
beast::xor_shift_engine rng(kSeedValue);
std::shuffle(batch_.begin(), batch_.end(), rng);
auto const copy = fetchCopyOfBatch(*backend, batch_);
EXPECT_EQ(batch_, copy);
}
}
TEST_P(BackendTypeTest, persists_after_reopen)
{
{
auto backend = makeOpenBackend();
storeBatch(*backend, batch_);
}
// re-open a fresh backend instance over the same path
auto backend = makeOpenBackend();
auto copy = fetchCopyOfBatch(*backend, batch_);
std::ranges::sort(batch_, LessThan{});
std::ranges::sort(copy, LessThan{});
EXPECT_EQ(batch_, copy);
}
// missing-key path. Replaces the correctness half of Timing_test::doMissing
// (and the missing branch of doMixed): every fetch on an empty backend must
// report Status::NotFound.
TEST_P(BackendTypeTest, fetch_missing)
{
auto backend = makeOpenBackend();
// deliberately do NOT store batch_ — every key must be absent
fetchMissing(*backend, batch_);
}
// concurrent store/fetch correctness. Replaces the correctness half of the
// multi-threaded Timing_test workloads (which only ran manually, never in CI):
// many threads store disjoint objects, then many threads fetch and verify each
// round-trips. Doubles as a thread-safety smoke test for the backend.
TEST_P(BackendTypeTest, concurrent_store_and_fetch)
{
// The SQLite backend is not designed for concurrent writers (and the old
// Timing_test only exercised nudb/rocksdb under threads).
if (GetParam() == "sqlite")
GTEST_SKIP() << "sqlite backend is not exercised under concurrency";
for (auto const numThreads : {4uz, 8uz})
{
SCOPED_TRACE("threads=" + std::to_string(numThreads));
auto backend = makeOpenBackend();
// concurrent stores of disjoint objects
parallelFor(batch_.size(), numThreads, [&](std::size_t i) { backend->store(batch_[i]); });
// concurrent fetches, each verifying its object round-trips. Worker
// threads only touch an atomic counter; the EXPECT runs on the main
// thread after join to avoid relying on cross-thread assertion support.
std::atomic<std::size_t> mismatches{0};
parallelFor(batch_.size(), numThreads, [&](std::size_t i) {
std::shared_ptr<NodeObject> result;
if (backend->fetch(batch_[i]->getHash(), &result) != Status::Ok || !result ||
!isSame(result, batch_[i]))
{
++mismatches;
}
});
EXPECT_EQ(mismatches.load(), 0u);
backend->close();
}
}
INSTANTIATE_TEST_SUITE_P(
BackendTypes,
BackendTypeTest,
::testing::ValuesIn(backendTypes()),
[](::testing::TestParamInfo<std::string> const& info) { return info.param; });
} // namespace xrpl::node_store

View File

@@ -0,0 +1,41 @@
#include <xrpl/nodestore/NodeObject.h>
#include <xrpl/nodestore/detail/DecodedBlob.h>
#include <xrpl/nodestore/detail/EncodedBlob.h>
#include <gtest/gtest.h>
#include <nodestore/TestBase.h>
#include <cstddef>
#include <memory>
#include <string>
namespace xrpl::node_store {
TEST(NodeStoreBasics, predictable_batches)
{
auto const batch1 = createPredictableBatch(kNumObjectsToTest, kSeedValue);
auto const batch2 = createPredictableBatch(kNumObjectsToTest, kSeedValue);
EXPECT_EQ(batch1, batch2);
auto const batch3 = createPredictableBatch(kNumObjectsToTest, kSeedValue + 1);
EXPECT_NE(batch1, batch3);
}
TEST(NodeStoreBasics, blob_encoding)
{
auto const batch = createPredictableBatch(kNumObjectsToTest, kSeedValue);
for (std::size_t i = 0; i < batch.size(); ++i)
{
SCOPED_TRACE("blob index=" + std::to_string(i));
EncodedBlob const encoded(batch[i]);
DecodedBlob decoded(encoded.getKey(), encoded.getData(), encoded.getSize());
EXPECT_TRUE(decoded.wasOk());
if (decoded.wasOk())
{
std::shared_ptr<NodeObject> const object(decoded.createObject());
EXPECT_TRUE(isSame(batch[i], object));
}
}
}
} // namespace xrpl::node_store

View File

@@ -0,0 +1,146 @@
#include <xrpl/nodestore/detail/codec.h>
#include <xrpl/nodestore/NodeObject.h>
#include <xrpl/protocol/HashPrefix.h>
#include <gtest/gtest.h>
#include <nudb/detail/buffer.hpp>
#include <nudb/detail/stream.hpp>
#include <array>
#include <cstddef>
#include <cstdint>
#include <cstring>
#include <utility>
#include <vector>
using namespace xrpl;
using namespace xrpl::node_store;
namespace {
// v1 inner-node layout: 16 hashes of 32 bytes each
constexpr std::size_t kHashCount = 16;
constexpr std::size_t kHashSize = 32;
std::vector<std::uint8_t>
makeInnerNode(std::size_t nonEmptySlots)
{
using namespace nudb::detail;
static constexpr std::size_t kInnerNodeSize = 525;
std::array<std::uint8_t, kHashCount * kHashSize> hashes{};
for (auto slot = 0uz; slot < nonEmptySlots; ++slot)
{
for (auto byte = 0uz; byte < kHashSize; ++byte)
{
std::size_t const offset = (slot * kHashSize) + byte;
hashes[offset] = static_cast<std::uint8_t>((offset % 255) + 1);
}
}
std::vector<std::uint8_t> blob(kInnerNodeSize);
ostream os(blob.data(), blob.size());
write<std::uint32_t>(os, 0); // index
write<std::uint32_t>(os, 0); // unused
write<std::uint8_t>(os, static_cast<std::uint8_t>(NodeObjectType::Unknown));
write<std::uint32_t>(os, static_cast<std::uint32_t>(HashPrefix::InnerNode));
write(os, hashes.data(), hashes.size());
return blob;
}
std::uint8_t
codecType(std::pair<void const*, std::size_t> const& compressed)
{
return static_cast<std::uint8_t const*>(compressed.first)[0];
}
} // namespace
// All 16 hash slots populated - "full v1 inner node"
TEST(Codec, inner_node_full_roundtrip)
{
static constexpr std::uint8_t kTypeInnerNodeFull = 3;
auto const blob = makeInnerNode(kHashCount);
nudb::detail::buffer compressBuf;
auto const compressed = nodeobjectCompress(blob.data(), blob.size(), compressBuf);
EXPECT_EQ(codecType(compressed), kTypeInnerNodeFull);
EXPECT_EQ(compressed.second, sizeVarint(kTypeInnerNodeFull) + (kHashCount * kHashSize));
nudb::detail::buffer decompressBuf;
auto const restored = nodeobjectDecompress(compressed.first, compressed.second, decompressBuf);
EXPECT_EQ(restored.second, blob.size());
EXPECT_EQ(std::memcmp(restored.first, blob.data(), blob.size()), 0);
}
// Some hash slots empty - "compressed v1 inner node"
TEST(Codec, inner_node_compressed_roundtrip)
{
static constexpr std::uint8_t kTypeInnerNodeCompressed = 2;
static constexpr std::size_t kNonEmpty = 5;
auto const blob = makeInnerNode(kNonEmpty);
nudb::detail::buffer compressBuf;
auto const compressed = nodeobjectCompress(blob.data(), blob.size(), compressBuf);
EXPECT_EQ(codecType(compressed), kTypeInnerNodeCompressed);
EXPECT_EQ(
compressed.second,
sizeVarint(kTypeInnerNodeCompressed) + sizeof(std::uint16_t) + (kNonEmpty * kHashSize));
EXPECT_LT(compressed.second, blob.size());
nudb::detail::buffer decompressBuf;
auto const restored = nodeobjectDecompress(compressed.first, compressed.second, decompressBuf);
EXPECT_EQ(restored.second, blob.size());
EXPECT_EQ(std::memcmp(restored.first, blob.data(), blob.size()), 0);
}
// Anything that is not a v1 inner node - lz4 compressed
TEST(Codec, lz4_roundtrip)
{
// A payload that is deliberately not a v1 inner node (any size other than 525), filled with a
// short repeating pattern so lz4 actually shrinks it.
static constexpr std::size_t kNonInnerNodeSize = 1000;
static constexpr std::size_t kBytePatternPeriod = 7;
static constexpr std::uint8_t kTypeLz4 = 1;
std::vector<std::uint8_t> blob(kNonInnerNodeSize);
for (auto i = 0uz; i < blob.size(); ++i)
blob[i] = static_cast<std::uint8_t>(i % kBytePatternPeriod);
nudb::detail::buffer compressBuf;
auto const compressed = nodeobjectCompress(blob.data(), blob.size(), compressBuf);
EXPECT_EQ(codecType(compressed), kTypeLz4);
nudb::detail::buffer decompressBuf;
auto const restored = nodeobjectDecompress(compressed.first, compressed.second, decompressBuf);
EXPECT_EQ(restored.second, blob.size());
EXPECT_EQ(std::memcmp(restored.first, blob.data(), blob.size()), 0);
}
// An uncompressed blob is never produced by the compressor but must still decode: leading varint 0
// followed by the raw payload.
TEST(Codec, uncompressed_passthrough)
{
static constexpr std::uint8_t kTypeUncompressed = 0;
static constexpr auto payload = std::to_array<std::uint8_t>({0xde, 0xad, 0xbe, 0xef, 0x2a});
std::vector<std::uint8_t> blob;
blob.push_back(kTypeUncompressed); // leading varint type tag
blob.insert(blob.end(), payload.begin(), payload.end());
nudb::detail::buffer decompressBuf;
auto const restored = nodeobjectDecompress(blob.data(), blob.size(), decompressBuf);
EXPECT_EQ(restored.second, payload.size());
EXPECT_EQ(std::memcmp(restored.first, payload.data(), payload.size()), 0);
}

View File

@@ -0,0 +1,248 @@
#include <xrpl/nodestore/Database.h>
#include <xrpl/basics/ByteUtilities.h>
#include <xrpl/beast/utility/Journal.h>
#include <xrpl/beast/utility/temp_dir.h>
#include <xrpl/beast/xor_shift_engine.h>
#include <xrpl/config/BasicConfig.h>
#include <xrpl/nodestore/DummyScheduler.h>
#include <xrpl/nodestore/Manager.h>
#include <xrpl/nodestore/Types.h>
#include <xrpl/protocol/SystemParameters.h>
#include <gtest/gtest.h>
#include <helpers/TestSink.h>
#include <nodestore/TestBase.h>
#include <algorithm>
#include <memory>
#include <stdexcept>
#include <string>
#include <vector>
namespace xrpl::node_store {
namespace {
std::vector<std::string>
allBackends()
{
std::vector<std::string> types{"memory", "nudb"};
#if XRPL_ROCKSDB_AVAILABLE
types.emplace_back("rocksdb");
#endif
return types;
}
std::vector<std::string>
persistentBackends()
{
std::vector<std::string> types{"nudb"};
#if XRPL_ROCKSDB_AVAILABLE
types.emplace_back("rocksdb");
#endif
return types;
}
std::vector<std::string>
importBackends()
{
std::vector<std::string> types{"nudb"};
#if XRPL_ROCKSDB_AVAILABLE
types.emplace_back("rocksdb");
#endif
#ifdef XRPL_ENABLE_SQLITE_BACKEND_TESTS
types.emplace_back("sqlite");
#endif
return types;
}
} // namespace
// Shared setup for the parameterized Database tests: builds the node params,
// journal and a predictable batch per test, mirroring Backend.cpp's fixture.
class NodeStoreDatabaseTestBase : public ::testing::TestWithParam<std::string>
{
protected:
void
SetUp() override
{
nodeParams_.set("type", GetParam());
nodeParams_.set("path", nodeDb_.path());
beast::xor_shift_engine rng(kSeedValue);
batch_ = createPredictableBatch(kNumObjects, rng());
}
std::unique_ptr<Database>
makeDatabase()
{
return Manager::instance().makeDatabase(megabytes(4), scheduler_, 2, nodeParams_, journal_);
}
DummyScheduler scheduler_;
beast::TempDir const nodeDb_;
beast::Journal const journal_{TestSink::instance()};
Section nodeParams_;
Batch batch_;
};
class NodeStoreDatabaseTest : public NodeStoreDatabaseTestBase
{
};
class NodeStoreDatabasePersistenceTest : public NodeStoreDatabaseTestBase
{
};
TEST_P(NodeStoreDatabaseTest, store_and_fetch)
{
auto db = makeDatabase();
storeBatch(*db, batch_);
{
SCOPED_TRACE("read in original order");
auto const copy = fetchCopyOfBatch(*db, batch_);
EXPECT_EQ(batch_, copy);
}
{
SCOPED_TRACE("read in shuffled order");
beast::xor_shift_engine rng(kSeedValue);
std::shuffle(batch_.begin(), batch_.end(), rng);
auto const copy = fetchCopyOfBatch(*db, batch_);
EXPECT_EQ(batch_, copy);
}
}
TEST_P(NodeStoreDatabasePersistenceTest, round_trip)
{
{
auto db = makeDatabase();
storeBatch(*db, batch_);
}
// re-open without the ephemeral db
auto db = makeDatabase();
auto copy = fetchCopyOfBatch(*db, batch_);
std::ranges::sort(batch_, LessThan{});
std::ranges::sort(copy, LessThan{});
EXPECT_EQ(batch_, copy);
}
// missing-key path at the Database layer. Mirrors Backend's fetch_missing —
// fetching keys that were never stored must return nullptr (NotFound).
TEST_P(NodeStoreDatabaseTest, fetch_missing)
{
auto db = makeDatabase();
// never store: every key must be absent
fetchMissing(*db, batch_);
}
INSTANTIATE_TEST_SUITE_P(
NodeStoreBackends,
NodeStoreDatabaseTest,
::testing::ValuesIn(allBackends()),
[](::testing::TestParamInfo<std::string> const& info) { return info.param; });
INSTANTIATE_TEST_SUITE_P(
PersistentBackends,
NodeStoreDatabasePersistenceTest,
::testing::ValuesIn(persistentBackends()),
[](::testing::TestParamInfo<std::string> const& info) { return info.param; });
TEST(NodeStoreDatabase, memory_earliest_seq)
{
DummyScheduler scheduler;
beast::TempDir const nodeDb;
Section nodeParams;
nodeParams.set("type", "memory");
nodeParams.set("path", nodeDb.path());
beast::Journal const journal(TestSink::instance());
// default earliest ledger sequence
{
auto db = Manager::instance().makeDatabase(megabytes(4), scheduler, 2, nodeParams, journal);
EXPECT_EQ(db->earliestLedgerSeq(), kXrpLedgerEarliestSeq);
}
// invalid earliest_seq value
{
nodeParams.set("earliest_seq", "0");
try
{
auto db =
Manager::instance().makeDatabase(megabytes(4), scheduler, 2, nodeParams, journal);
FAIL() << "expected runtime_error for earliest_seq=0";
}
catch (std::runtime_error const& e)
{
EXPECT_STREQ(e.what(), "Invalid earliest_seq");
}
}
// valid earliest_seq value
{
nodeParams.set("earliest_seq", "1");
auto db = Manager::instance().makeDatabase(megabytes(4), scheduler, 2, nodeParams, journal);
EXPECT_EQ(db->earliestLedgerSeq(), 1u);
}
}
class DatabaseImportTest : public ::testing::TestWithParam<std::string>
{
};
TEST_P(DatabaseImportTest, same_backend)
{
auto const type = GetParam();
DummyScheduler scheduler;
beast::Journal const journal(TestSink::instance());
beast::TempDir const srcDir;
Section srcParams;
srcParams.set("type", type);
srcParams.set("path", srcDir.path());
auto batch = createPredictableBatch(kNumObjects, kSeedValue);
// write to source db
{
auto src = Manager::instance().makeDatabase(megabytes(4), scheduler, 2, srcParams, journal);
storeBatch(*src, batch);
}
Batch copy;
{
// re-open source and import into a fresh destination
auto src = Manager::instance().makeDatabase(megabytes(4), scheduler, 2, srcParams, journal);
beast::TempDir const destDir;
Section destParams;
destParams.set("type", type);
destParams.set("path", destDir.path());
auto dest =
Manager::instance().makeDatabase(megabytes(4), scheduler, 2, destParams, journal);
dest->importDatabase(*src);
copy = fetchCopyOfBatch(*dest, batch);
}
std::ranges::sort(batch, LessThan{});
std::ranges::sort(copy, LessThan{});
EXPECT_EQ(batch, copy);
}
INSTANTIATE_TEST_SUITE_P(
ImportBackends,
DatabaseImportTest,
::testing::ValuesIn(importBackends()),
[](::testing::TestParamInfo<std::string> const& info) { return info.param; });
} // namespace xrpl::node_store

View File

@@ -0,0 +1,297 @@
#include <xrpl/basics/ByteUtilities.h>
#include <xrpl/beast/utility/Journal.h>
#include <xrpl/beast/utility/temp_dir.h>
#include <xrpl/config/BasicConfig.h>
#include <xrpl/nodestore/DummyScheduler.h>
#include <xrpl/nodestore/Manager.h>
#include <gtest/gtest.h>
#include <helpers/CaptureSink.h>
#include <helpers/TestSink.h>
#include <nodestore/TestBase.h>
#include <array>
#include <cstddef>
#include <exception>
#include <memory>
#include <string>
#include <utility>
#include <vector>
namespace xrpl::node_store {
namespace {
Section
makeSection(std::string const& path, std::string const& blockSize = "")
{
Section params;
params.set("type", "nudb");
params.set("path", path);
if (!blockSize.empty())
params.set("nudb_block_size", blockSize);
return params;
}
void
runRoundTrip(Section const& params, std::size_t expectedBlocksize)
{
DummyScheduler scheduler;
beast::Journal const journal(TestSink::instance());
auto backend = Manager::instance().makeBackend(params, megabytes(4), scheduler, journal);
ASSERT_TRUE(backend);
ASSERT_EQ(backend->getBlockSize(), expectedBlocksize);
backend->open();
ASSERT_TRUE(backend->isOpen());
auto const batch = createPredictableBatch(10, 12345);
storeBatch(*backend, batch);
auto const copy = fetchCopyOfBatch(*backend, batch);
backend->close();
EXPECT_EQ(batch, copy);
}
} // namespace
TEST(NuDBFactory, default_block_size)
{
beast::TempDir const tempDir;
auto const params = makeSection(tempDir.path());
ASSERT_NO_FATAL_FAILURE(runRoundTrip(params, 4096));
}
TEST(NuDBFactory, valid_block_sizes)
{
auto const kValidSizes = std::to_array<std::size_t>({4096, 8192, 16384, 32768});
for (auto const size : kValidSizes)
{
SCOPED_TRACE("size=" + std::to_string(size));
beast::TempDir const tempDir;
auto const params = makeSection(tempDir.path(), std::to_string(size));
ASSERT_NO_FATAL_FAILURE(runRoundTrip(params, size));
}
// empty value is ignored by config parser; default (4096) is used
{
beast::TempDir const tempDir;
auto const params = makeSection(tempDir.path(), "");
ASSERT_NO_FATAL_FAILURE(runRoundTrip(params, 4096));
}
}
TEST(NuDBFactory, invalid_block_sizes)
{
std::vector<std::string> const kInvalidSizes = {
"2048", // too small
"1024", // too small
"65536", // too large
"131072", // too large
"5000", // not power of 2
"6000", // not power of 2
"10000", // not power of 2
"0", // zero
"-1", // negative
"abc", // non-numeric
"4k", // invalid format
"4096.5"}; // decimal
for (auto const& size : kInvalidSizes)
{
SCOPED_TRACE("size='" + size + "'");
beast::TempDir const tempDir;
auto const params = makeSection(tempDir.path(), size);
EXPECT_THROW(runRoundTrip(params, 4096), std::exception);
}
// whitespace handling — lexical_cast may or may not strip; treat as invalid
std::vector<std::string> const kWhitespaceSizes = {"4096 ", " 4096"};
for (auto const& size : kWhitespaceSizes)
{
SCOPED_TRACE("size='" + size + "'");
beast::TempDir const tempDir;
auto const params = makeSection(tempDir.path(), size);
EXPECT_THROW(runRoundTrip(params, 4096), std::exception);
}
}
TEST(NuDBFactory, log_messages)
{
// valid custom block size emits info log
{
beast::TempDir const tempDir;
auto const params = makeSection(tempDir.path(), "8192");
test::CaptureSink sink(beast::Severity::Info);
beast::Journal const journal(sink);
DummyScheduler scheduler;
[[maybe_unused]] auto backend =
Manager::instance().makeBackend(params, megabytes(4), scheduler, journal);
EXPECT_TRUE(sink.messages().contains("Using custom NuDB block size: 8192"));
}
// invalid block size throws with informative message
{
beast::TempDir const tempDir;
auto const params = makeSection(tempDir.path(), "5000");
test::CaptureSink sink(beast::Severity::Warning);
beast::Journal const journal(sink);
DummyScheduler scheduler;
try
{
auto backend =
Manager::instance().makeBackend(params, megabytes(4), scheduler, journal);
FAIL() << "expected exception for invalid block size 5000";
}
catch (std::exception const& e)
{
std::string const what{e.what()};
EXPECT_TRUE(what.contains("Invalid nudb_block_size: 5000"));
EXPECT_TRUE(what.contains("Must be power of 2 between 4096 and 32768"));
}
}
// non-numeric value throws
{
beast::TempDir const tempDir;
auto const params = makeSection(tempDir.path(), "invalid");
test::CaptureSink sink(beast::Severity::Warning);
beast::Journal const journal(sink);
DummyScheduler scheduler;
try
{
auto backend =
Manager::instance().makeBackend(params, megabytes(4), scheduler, journal);
FAIL() << "expected exception for non-numeric block size";
}
catch (std::exception const& e)
{
std::string const what{e.what()};
EXPECT_TRUE(what.contains("Invalid nudb_block_size value: invalid"));
}
}
}
TEST(NuDBFactory, power_of_two_validation)
{
std::vector<std::pair<std::string, bool>> const kCASES = {
{"4095", false}, // just below minimum
{"4096", true}, // minimum valid
{"4097", false}, // not power of 2
{"8192", true}, // valid power of 2
{"8193", false}, // not power of 2
{"16384", true}, // valid power of 2
{"32768", true}, // maximum valid
{"32769", false}, // just above maximum
{"65536", false}}; // power of 2 but too large
for (auto const& [size, shouldWork] : kCASES)
{
SCOPED_TRACE("size=" + size + " shouldWork=" + (shouldWork ? "true" : "false"));
beast::TempDir const tempDir;
auto const params = makeSection(tempDir.path(), size);
test::CaptureSink sink(beast::Severity::Warning);
beast::Journal const journal(sink);
DummyScheduler scheduler;
try
{
auto backend =
Manager::instance().makeBackend(params, megabytes(4), scheduler, journal);
EXPECT_TRUE(shouldWork);
}
catch (std::exception const& e)
{
// A throw is only expected for sizes that should NOT work; if a
// valid size throws, fail here instead of silently matching the
// message below (which would mask the regression).
EXPECT_FALSE(shouldWork);
std::string const what{e.what()};
EXPECT_TRUE(what.contains("Invalid nudb_block_size"));
}
}
}
TEST(NuDBFactory, both_constructor_variants)
{
beast::TempDir const tempDir;
auto const params = makeSection(tempDir.path(), "16384");
DummyScheduler scheduler;
beast::Journal const journal(TestSink::instance());
auto backend1 = Manager::instance().makeBackend(params, megabytes(4), scheduler, journal);
EXPECT_NE(backend1, nullptr);
ASSERT_NO_FATAL_FAILURE(runRoundTrip(params, 16384));
// Test second constructor (with nudb::context)
// Note: This would require access to nudb::context, which might not be
// easily testable without more complex setup. For now, we test that
// the factory can create backends with the first constructor.
}
TEST(NuDBFactory, configuration_parsing)
{
// basic valid format emits success log
{
beast::TempDir const tempDir;
auto const params = makeSection(tempDir.path(), "8192");
test::CaptureSink sink(beast::Severity::Info);
beast::Journal const journal(sink);
DummyScheduler scheduler;
[[maybe_unused]] auto backend =
Manager::instance().makeBackend(params, megabytes(4), scheduler, journal);
EXPECT_TRUE(sink.messages().contains("Using custom NuDB block size"));
}
// Test whitespace handling separately since lexical_cast behavior may vary
std::vector<std::string> const kWhitespaceFormats = {" 8192", "8192 "};
for (auto const& format : kWhitespaceFormats)
{
SCOPED_TRACE("format='" + format + "'");
beast::TempDir const tempDir;
auto const params = makeSection(tempDir.path(), format);
test::CaptureSink sink(beast::Severity::Debug);
beast::Journal const journal(sink);
DummyScheduler scheduler;
EXPECT_ANY_THROW(Manager::instance().makeBackend(params, megabytes(4), scheduler, journal));
}
}
TEST(NuDBFactory, data_persistence)
{
std::vector<std::string> const kBlockSizes = {"4096", "8192", "16384", "32768"};
for (auto const& size : kBlockSizes)
{
SCOPED_TRACE("size=" + size);
beast::TempDir const tempDir;
auto const params = makeSection(tempDir.path(), size);
DummyScheduler scheduler;
beast::Journal const journal(TestSink::instance());
// Create test data
auto const batch = createPredictableBatch(50, 54321);
// Store data
{
auto backend =
Manager::instance().makeBackend(params, megabytes(4), scheduler, journal);
backend->open();
storeBatch(*backend, batch);
backend->close();
}
// Retrieve data in new backend instance
{
auto backend =
Manager::instance().makeBackend(params, megabytes(4), scheduler, journal);
backend->open();
auto const copy = fetchCopyOfBatch(*backend, batch);
EXPECT_EQ(batch, copy);
backend->close();
}
}
}
} // namespace xrpl::node_store

View File

@@ -0,0 +1,169 @@
#pragma once
#include <xrpl/basics/Blob.h>
#include <xrpl/basics/base_uint.h>
#include <xrpl/basics/random.h>
#include <xrpl/beast/utility/rngfill.h>
#include <xrpl/beast/xor_shift_engine.h>
#include <xrpl/nodestore/Backend.h>
#include <xrpl/nodestore/Database.h>
#include <xrpl/nodestore/NodeObject.h>
#include <xrpl/nodestore/Types.h>
#include <gtest/gtest.h>
#include <algorithm>
#include <cstddef>
#include <cstdint>
#include <memory>
#include <string>
#include <utility>
namespace xrpl::node_store {
constexpr std::size_t kMinPayloadBytes = 1;
constexpr std::size_t kMaxPayloadBytes = 2000;
constexpr int kNumObjectsToTest = 2000;
constexpr int kNumObjects = 2000;
constexpr std::uint64_t kSeedValue = 50;
struct LessThan
{
bool
operator()(std::shared_ptr<NodeObject> const& lhs, std::shared_ptr<NodeObject> const& rhs)
const noexcept
{
return lhs->getHash() < rhs->getHash();
}
};
[[nodiscard]] inline bool
isSame(std::shared_ptr<NodeObject> const& lhs, std::shared_ptr<NodeObject> const& rhs)
{
return (lhs->getType() == rhs->getType()) && (lhs->getHash() == rhs->getHash()) &&
(lhs->getData() == rhs->getData());
}
[[nodiscard]] inline Batch
createPredictableBatch(std::size_t numObjects, std::uint64_t seed)
{
Batch batch;
batch.reserve(numObjects);
beast::xor_shift_engine rng(seed);
for (auto i = 0uz; i < numObjects; ++i)
{
NodeObjectType const type = [&] {
switch (randInt(rng, 3))
{
case 0:
return NodeObjectType::Ledger;
case 1:
return NodeObjectType::AccountNode;
case 2:
return NodeObjectType::TransactionNode;
case 3:
default:
return NodeObjectType::Unknown;
}
}();
uint256 hash;
beast::rngfill(hash.begin(), hash.size(), rng);
Blob blob(randInt(rng, kMinPayloadBytes, kMaxPayloadBytes));
beast::rngfill(blob.data(), blob.size(), rng);
batch.emplace_back(NodeObject::createObject(type, std::move(blob), hash));
}
return batch;
}
inline void
storeBatch(Backend& backend, Batch const& batch)
{
for (auto const& obj : batch)
backend.store(obj);
}
[[nodiscard]] inline Batch
fetchCopyOfBatch(Backend& backend, Batch const& batch)
{
Batch copy;
copy.reserve(batch.size());
for (auto i = 0uz; i < batch.size(); ++i)
{
SCOPED_TRACE("fetchCopyOfBatch index=" + std::to_string(i));
std::shared_ptr<NodeObject> object;
Status const status = backend.fetch(batch[i]->getHash(), &object);
EXPECT_EQ(status, Status::Ok);
if (status == Status::Ok)
{
EXPECT_NE(object, nullptr);
copy.emplace_back(object);
}
}
return copy;
}
inline void
fetchMissing(Backend& backend, Batch const& batch)
{
for (auto i = 0uz; i < batch.size(); ++i)
{
SCOPED_TRACE("fetchMissing index=" + std::to_string(i));
std::shared_ptr<NodeObject> object;
Status const status = backend.fetch(batch[i]->getHash(), &object);
EXPECT_EQ(status, Status::NotFound);
}
}
inline void
storeBatch(Database& db, Batch const& batch)
{
for (auto const& obj : batch)
{
Blob data(obj->getData());
db.store(obj->getType(), std::move(data), obj->getHash(), db.earliestLedgerSeq());
}
}
[[nodiscard]] inline Batch
fetchCopyOfBatch(Database& db, Batch const& batch)
{
Batch copy;
copy.reserve(batch.size());
for (auto const& obj : batch)
{
std::shared_ptr<NodeObject> const result = db.fetchNodeObject(obj->getHash(), 0);
if (result != nullptr)
copy.emplace_back(result);
}
return copy;
}
inline void
fetchMissing(Database& db, Batch const& batch)
{
for (auto i = 0uz; i < batch.size(); ++i)
{
SCOPED_TRACE("fetchMissing(Database) index=" + std::to_string(i));
EXPECT_EQ(db.fetchNodeObject(batch[i]->getHash(), 0), nullptr);
}
}
} // namespace xrpl::node_store
namespace xrpl {
[[nodiscard]] inline bool
operator==(node_store::Batch const& lhs, node_store::Batch const& rhs)
{
return std::ranges::equal(lhs, rhs, node_store::isSame);
}
} // namespace xrpl

View File

@@ -0,0 +1,46 @@
#include <xrpl/nodestore/detail/Varint.h>
#include <gtest/gtest.h>
#include <array>
#include <cstddef>
#include <cstdint>
#include <limits>
using namespace xrpl::node_store;
TEST(Varint, encode_decode)
{
static constexpr auto kValues = std::to_array<std::size_t>({
0,
1,
2,
126,
127,
128,
253,
254,
255,
16127,
16128,
16129,
0xff,
0xffff,
0xffffffff,
0xffffffffffffUL,
std::numeric_limits<std::size_t>::max(),
});
for (auto const value : kValues)
{
std::array<std::uint8_t, VarintTraits<std::size_t>::kMax> buffer{};
auto const bytesWritten = writeVarint(buffer.data(), value);
EXPECT_GT(bytesWritten, 0u);
EXPECT_EQ(bytesWritten, sizeVarint(value));
std::size_t decoded = 0;
auto const bytesRead = readVarint(buffer.data(), bytesWritten, decoded);
EXPECT_EQ(bytesRead, bytesWritten);
EXPECT_EQ(value, decoded);
}
}

View File

@@ -0,0 +1,26 @@
#include <xrpl/protocol/ApiVersion.h>
#include <gtest/gtest.h>
using namespace xrpl;
TEST(ApiVersion, invariants)
{
static_assert(RPC::kApiMinimumSupportedVersion <= RPC::kApiMaximumSupportedVersion);
static_assert(RPC::kApiMinimumSupportedVersion <= RPC::kApiMaximumValidVersion);
static_assert(RPC::kApiMaximumSupportedVersion <= RPC::kApiMaximumValidVersion);
static_assert(RPC::kApiBetaVersion <= RPC::kApiMaximumValidVersion);
}
// Update when we change versions
TEST(ApiVersion, versions)
{
static_assert(RPC::kApiMinimumSupportedVersion >= 1);
static_assert(RPC::kApiMinimumSupportedVersion < 2);
static_assert(RPC::kApiMaximumSupportedVersion >= 2);
static_assert(RPC::kApiMaximumSupportedVersion < 3);
static_assert(RPC::kApiMaximumValidVersion >= 3);
static_assert(RPC::kApiMaximumValidVersion < 4);
static_assert(RPC::kApiBetaVersion >= 3);
static_assert(RPC::kApiBetaVersion < 4);
}

View File

@@ -0,0 +1,49 @@
#include <xrpl/protocol/Serializer.h>
#include <gtest/gtest.h>
#include <array>
#include <cstdint>
#include <limits>
using namespace xrpl;
TEST(Serializer, add32_roundtrip)
{
static constexpr auto kValues = std::to_array<std::int32_t>({
std::numeric_limits<std::int32_t>::min(),
-1,
0,
1,
std::numeric_limits<std::int32_t>::max(),
});
for (std::int32_t const value : kValues)
{
Serializer s;
s.add32(value);
EXPECT_EQ(s.size(), 4);
SerialIter sit(s.slice());
EXPECT_EQ(sit.geti32(), value);
}
}
TEST(Serializer, add64_roundtrip)
{
static constexpr auto kValues = std::to_array<std::int64_t>({
std::numeric_limits<std::int64_t>::min(),
-1,
0,
1,
std::numeric_limits<std::int64_t>::max(),
});
for (std::int64_t const value : kValues)
{
Serializer s;
s.add64(value);
EXPECT_EQ(s.size(), 8);
SerialIter sit(s.slice());
EXPECT_EQ(sit.geti64(), value);
}
}

View File

@@ -24,13 +24,13 @@ namespace xrpl::tests {
class TestNodeFamily : public Family
{
private:
std::unique_ptr<NodeStore::Database> db_;
std::unique_ptr<node_store::Database> db_;
std::shared_ptr<FullBelowCache> fbCache_;
std::shared_ptr<TreeNodeCache> tnCache_;
TestStopwatch clock_;
NodeStore::DummyScheduler scheduler_;
node_store::DummyScheduler scheduler_;
beast::Journal const j_;
@@ -49,17 +49,17 @@ public:
Section testSection;
testSection.set(Keys::kType, "memory");
testSection.set(Keys::kPath, "SHAMap_test");
db_ = NodeStore::Manager::instance().makeDatabase(
db_ = node_store::Manager::instance().makeDatabase(
megabytes(4), scheduler_, 1, testSection, j);
}
NodeStore::Database&
node_store::Database&
db() override
{
return *db_;
}
[[nodiscard]] NodeStore::Database const&
[[nodiscard]] node_store::Database const&
db() const override
{
return *db_;

View File

@@ -18,7 +18,7 @@ namespace xrpl {
class AccountStateSF : public SHAMapSyncFilter
{
public:
AccountStateSF(NodeStore::Database& db, AbstractFetchPackContainer& fp) : db_(db), fp_(fp)
AccountStateSF(node_store::Database& db, AbstractFetchPackContainer& fp) : db_(db), fp_(fp)
{
}
@@ -34,7 +34,7 @@ public:
getNode(SHAMapHash const& nodeHash) const override;
private:
NodeStore::Database& db_;
node_store::Database& db_;
AbstractFetchPackContainer& fp_;
};

View File

@@ -136,7 +136,7 @@ private:
addPeers();
void
tryDB(NodeStore::Database& srcDB);
tryDB(node_store::Database& srcDB);
void
done();

View File

@@ -18,7 +18,7 @@ namespace xrpl {
class TransactionStateSF : public SHAMapSyncFilter
{
public:
TransactionStateSF(NodeStore::Database& db, AbstractFetchPackContainer& fp) : db_(db), fp_(fp)
TransactionStateSF(node_store::Database& db, AbstractFetchPackContainer& fp) : db_(db), fp_(fp)
{
}
@@ -34,7 +34,7 @@ public:
getNode(SHAMapHash const& nodeHash) const override;
private:
NodeStore::Database& db_;
node_store::Database& db_;
AbstractFetchPackContainer& fp_;
};

View File

@@ -224,7 +224,7 @@ InboundLedger::neededStateHashes(int max, SHAMapSyncFilter const* filter) const
// See how much of the ledger data is stored locally
// Data found in a fetch pack will be stored
void
InboundLedger::tryDB(NodeStore::Database& srcDB)
InboundLedger::tryDB(node_store::Database& srcDB)
{
if (!haveHeader_)
{

View File

@@ -654,7 +654,7 @@ LedgerMaster::tryFill(std::shared_ptr<Ledger const> ledger)
std::uint32_t minHas = seq;
std::uint32_t maxHas = seq;
NodeStore::Database& nodeStore{app_.getNodeStore()};
node_store::Database& nodeStore{app_.getNodeStore()};
while (!app_.getJobQueue().isStopping() && seq > 0)
{
{

View File

@@ -232,7 +232,7 @@ public:
std::unique_ptr<Resource::Manager> resourceManager_;
std::unique_ptr<NodeStore::Database> nodeStore_;
std::unique_ptr<node_store::Database> nodeStore_;
NodeFamily nodeFamily_;
std::unique_ptr<OrderBookDB> orderBookDB_;
std::unique_ptr<PathRequestManager> pathRequestManager_;
@@ -655,7 +655,7 @@ public:
return tempNodeCache_;
}
NodeStore::Database&
node_store::Database&
getNodeStore() override
{
return *nodeStore_;
@@ -860,9 +860,9 @@ public:
if (config_->doImport)
{
auto j = logs_->journal("NodeObject");
NodeStore::DummyScheduler dummyScheduler;
std::unique_ptr<NodeStore::Database> source =
NodeStore::Manager::instance().makeDatabase(
node_store::DummyScheduler dummyScheduler;
std::unique_ptr<node_store::Database> source =
node_store::Manager::instance().makeDatabase(
megabytes(config_->getValueFor(SizedItem::BurstSize, std::nullopt)),
dummyScheduler,
0,

View File

@@ -12,7 +12,7 @@ NodeStoreScheduler::NodeStoreScheduler(JobQueue& jobQueue) : jobQueue_(jobQueue)
}
void
NodeStoreScheduler::scheduleTask(NodeStore::Task& task)
NodeStoreScheduler::scheduleTask(node_store::Task& task)
{
if (jobQueue_.isStopped())
return;
@@ -26,19 +26,19 @@ NodeStoreScheduler::scheduleTask(NodeStore::Task& task)
}
void
NodeStoreScheduler::onFetch(NodeStore::FetchReport const& report)
NodeStoreScheduler::onFetch(node_store::FetchReport const& report)
{
if (jobQueue_.isStopped())
return;
jobQueue_.addLoadEvents(
report.fetchType == NodeStore::FetchType::Async ? JtNsAsyncRead : JtNsSyncRead,
report.fetchType == node_store::FetchType::Async ? JtNsAsyncRead : JtNsSyncRead,
1,
report.elapsed);
}
void
NodeStoreScheduler::onBatchWrite(NodeStore::BatchWriteReport const& report)
NodeStoreScheduler::onBatchWrite(node_store::BatchWriteReport const& report)
{
if (jobQueue_.isStopped())
return;

View File

@@ -7,19 +7,19 @@
namespace xrpl {
/**
* A NodeStore::Scheduler which uses the JobQueue.
* A node_store::Scheduler which uses the JobQueue.
*/
class NodeStoreScheduler : public NodeStore::Scheduler
class NodeStoreScheduler : public node_store::Scheduler
{
public:
explicit NodeStoreScheduler(JobQueue& jobQueue);
void
scheduleTask(NodeStore::Task& task) override;
scheduleTask(node_store::Task& task) override;
void
onFetch(NodeStore::FetchReport const& report) override;
onFetch(node_store::FetchReport const& report) override;
void
onBatchWrite(NodeStore::BatchWriteReport const& report) override;
onBatchWrite(node_store::BatchWriteReport const& report) override;
private:
JobQueue& jobQueue_;

View File

@@ -44,7 +44,7 @@ public:
[[nodiscard]] virtual std::uint32_t
clampFetchDepth(std::uint32_t fetchDepth) const = 0;
virtual std::unique_ptr<NodeStore::Database>
virtual std::unique_ptr<node_store::Database>
makeNodeStore(int readThreads) = 0;
/**
@@ -102,5 +102,5 @@ public:
//------------------------------------------------------------------------------
std::unique_ptr<SHAMapStore>
makeSHAMapStore(Application& app, NodeStore::Scheduler& scheduler, beast::Journal journal);
makeSHAMapStore(Application& app, node_store::Scheduler& scheduler, beast::Journal journal);
} // namespace xrpl

View File

@@ -97,7 +97,7 @@ SHAMapStoreImp::SavedStateDB::setLastRotated(LedgerIndex seq)
SHAMapStoreImp::SHAMapStoreImp(
Application& app,
NodeStore::Scheduler& scheduler,
node_store::Scheduler& scheduler,
beast::Journal journal)
: app_(app)
, scheduler_(scheduler)
@@ -168,7 +168,7 @@ SHAMapStoreImp::SHAMapStoreImp(
}
}
std::unique_ptr<NodeStore::Database>
std::unique_ptr<node_store::Database>
SHAMapStoreImp::makeNodeStore(int readThreads)
{
auto nscfg = app_.config().section(Sections::kNodeDatabase);
@@ -188,7 +188,7 @@ SHAMapStoreImp::makeNodeStore(int readThreads)
std::to_string(app_.config().getValueFor(SizedItem::TreeCacheAge, std::nullopt)));
}
std::unique_ptr<NodeStore::Database> db;
std::unique_ptr<node_store::Database> db;
if (deleteInterval_ != 0u)
{
@@ -204,7 +204,7 @@ SHAMapStoreImp::makeNodeStore(int readThreads)
// Create NodeStore with two backends to allow online deletion of
// data
auto dbr = std::make_unique<NodeStore::DatabaseRotatingImp>(
auto dbr = std::make_unique<node_store::DatabaseRotatingImp>(
scheduler_,
readThreads,
std::move(writableBackend),
@@ -213,11 +213,11 @@ SHAMapStoreImp::makeNodeStore(int readThreads)
app_.getJournal(kNodeStoreName));
fdRequired_ += dbr->fdRequired();
dbRotating_ = dbr.get();
db.reset(dynamic_cast<NodeStore::Database*>(dbr.release()));
db.reset(dynamic_cast<node_store::Database*>(dbr.release()));
}
else
{
db = NodeStore::Manager::instance().makeDatabase(
db = node_store::Manager::instance().makeDatabase(
megabytes(app_.config().getValueFor(SizedItem::BurstSize, std::nullopt)),
scheduler_,
readThreads,
@@ -267,7 +267,7 @@ SHAMapStoreImp::copyNode(std::uint64_t& nodeCount, SHAMapTreeNode const& node)
{
// Copy a single record from node to dbRotating_
auto obj = dbRotating_->fetchNodeObject(
node.getHash().asUInt256(), 0, NodeStore::FetchType::Synchronous, true);
node.getHash().asUInt256(), 0, node_store::FetchType::Synchronous, true);
if (!obj)
{
XRPL_ASSERT(node.cowid() == 0, "SHAMapStoreImp::copyNode : rescued node must be clean");
@@ -395,7 +395,7 @@ SHAMapStoreImp::run()
// exception) also clear the flag.
struct RotationExposureGuard
{
NodeStore::DatabaseRotating& db;
node_store::DatabaseRotating& db;
~RotationExposureGuard()
{
db.setRotationInFlight(false);
@@ -539,7 +539,7 @@ SHAMapStoreImp::dbPaths()
boost::filesystem::remove_all(p);
}
std::unique_ptr<NodeStore::Backend>
std::unique_ptr<node_store::Backend>
SHAMapStoreImp::makeBackendRotating(std::string path)
{
Section section{app_.config().section(Sections::kNodeDatabase)};
@@ -558,7 +558,7 @@ SHAMapStoreImp::makeBackendRotating(std::string path)
}
section.set(Keys::kPath, newPath.string());
auto backend{NodeStore::Manager::instance().makeBackend(
auto backend{node_store::Manager::instance().makeBackend(
section,
megabytes(app_.config().getValueFor(SizedItem::BurstSize, std::nullopt)),
scheduler_,
@@ -775,7 +775,7 @@ SHAMapStoreImp::minimumOnline() const
//------------------------------------------------------------------------------
std::unique_ptr<SHAMapStore>
makeSHAMapStore(Application& app, NodeStore::Scheduler& scheduler, beast::Journal journal)
makeSHAMapStore(Application& app, node_store::Scheduler& scheduler, beast::Journal journal)
{
return std::make_unique<SHAMapStoreImp>(app, scheduler, journal);
}

View File

@@ -81,9 +81,9 @@ private:
// minimum ledger to maintain online.
std::atomic<LedgerIndex> minimumOnline_;
NodeStore::Scheduler& scheduler_;
node_store::Scheduler& scheduler_;
beast::Journal const journal_;
NodeStore::DatabaseRotating* dbRotating_ = nullptr;
node_store::DatabaseRotating* dbRotating_ = nullptr;
SavedStateDB stateDb_;
std::thread thread_;
bool stop_ = false;
@@ -124,7 +124,7 @@ private:
static constexpr auto kNodeStoreName = "NodeStore";
public:
SHAMapStoreImp(Application& app, NodeStore::Scheduler& scheduler, beast::Journal journal);
SHAMapStoreImp(Application& app, node_store::Scheduler& scheduler, beast::Journal journal);
std::uint32_t
clampFetchDepth(std::uint32_t fetchDepth) const override
@@ -132,7 +132,7 @@ public:
return (deleteInterval_ != 0u) ? std::min(fetchDepth, deleteInterval_) : fetchDepth;
}
std::unique_ptr<NodeStore::Database>
std::unique_ptr<node_store::Database>
makeNodeStore(int readThreads) override;
LedgerIndex
@@ -185,7 +185,7 @@ private:
void
dbPaths();
std::unique_ptr<NodeStore::Backend>
std::unique_ptr<node_store::Backend>
makeBackendRotating(std::string path = std::string());
template <class CacheInstance>
@@ -196,7 +196,7 @@ private:
for (auto const& key : cache.getKeys())
{
dbRotating_->fetchNodeObject(key, 0, NodeStore::FetchType::Synchronous, true);
dbRotating_->fetchNodeObject(key, 0, node_store::FetchType::Synchronous, true);
if (!(++check % checkHealthInterval_) && healthWait() == HealthResult::Stopping)
return true;
}

View File

@@ -33,13 +33,13 @@ public:
NodeFamily(Application& app, CollectorManager& cm);
NodeStore::Database&
node_store::Database&
db() override
{
return db_;
}
[[nodiscard]] NodeStore::Database const&
[[nodiscard]] node_store::Database const&
db() const override
{
return db_;
@@ -80,7 +80,7 @@ public:
private:
Application& app_;
NodeStore::Database& db_;
node_store::Database& db_;
beast::Journal const j_;
std::shared_ptr<FullBelowCache> fbCache_;