1#if XRPL_ROCKSDB_AVAILABLE
2#include <xrpl/basics/ByteUtilities.h>
3#include <xrpl/basics/Log.h>
4#include <xrpl/basics/base_uint.h>
5#include <xrpl/basics/contract.h>
6#include <xrpl/basics/safe_cast.h>
7#include <xrpl/beast/core/CurrentThreadName.h>
8#include <xrpl/beast/utility/Journal.h>
9#include <xrpl/beast/utility/instrumentation.h>
10#include <xrpl/config/BasicConfig.h>
11#include <xrpl/config/Constants.h>
12#include <xrpl/nodestore/Backend.h>
13#include <xrpl/nodestore/Factory.h>
14#include <xrpl/nodestore/Manager.h>
15#include <xrpl/nodestore/NodeObject.h>
16#include <xrpl/nodestore/Scheduler.h>
17#include <xrpl/nodestore/Types.h>
18#include <xrpl/nodestore/detail/BatchWriter.h>
19#include <xrpl/nodestore/detail/DecodedBlob.h>
20#include <xrpl/nodestore/detail/EncodedBlob.h>
22#include <rocksdb/advanced_options.h>
23#include <rocksdb/cache.h>
24#include <rocksdb/compression_type.h>
25#include <rocksdb/convenience.h>
26#include <rocksdb/db.h>
27#include <rocksdb/env.h>
28#include <rocksdb/filter_policy.h>
29#include <rocksdb/iterator.h>
30#include <rocksdb/options.h>
31#include <rocksdb/slice.h>
32#include <rocksdb/table.h>
33#include <rocksdb/write_batch.h>
46class RocksDBEnv :
public rocksdb::EnvWrapper
49 RocksDBEnv() : EnvWrapper(rocksdb::Env::Default())
55 ThreadParams(
void (*f)(
void*),
void* a) : f(f), a(a)
64 threadEntry(
void* ptr)
66 ThreadParams
const*
const p(
reinterpret_cast<ThreadParams*
>(ptr));
72 static std::atomic<std::size_t> kN;
73 std::size_t
const id(++kN);
80 StartThread(
void (*f)(
void*),
void* a)
override
82 auto*
const p =
new ThreadParams(f, a);
83 EnvWrapper::StartThread(&RocksDBEnv::threadEntry, p);
92 std::atomic<bool> deletePath_;
95 beast::Journal journal;
96 size_t const keyBytes;
99 std::unique_ptr<rocksdb::DB> db;
100 int fdMinRequired = 2048;
101 rocksdb::Options options;
105 Section
const& keyValues,
106 Scheduler& scheduler,
107 beast::Journal journal,
109 : deletePath_(false), journal(journal), keyBytes(keyBytes), batch(*this, scheduler)
112 Throw<std::runtime_error>(
"Missing path in RocksDBFactory backend");
114 rocksdb::BlockBasedTableOptions tableOptions;
118 keyValues.exists(Keys::kHardSet) && get<bool>(keyValues, Keys::kHardSet);
120 if (keyValues.exists(Keys::kCacheMb))
122 auto size = get<int>(keyValues, Keys::kCacheMb);
124 if (!hardSet && size == 256)
127 tableOptions.block_cache = rocksdb::NewLRUCache(megabytes(size));
130 if (
auto const v = get<int>(keyValues, Keys::kFilterBits))
132 bool const filterBlocks = !keyValues.exists(Keys::kFilterFull) ||
133 (get<int>(keyValues, Keys::kFilterFull) == 0);
134 tableOptions.filter_policy.reset(rocksdb::NewBloomFilterPolicy(v, filterBlocks));
137 if (
getIfExists(keyValues, Keys::kOpenFiles, options.max_open_files))
139 if (!hardSet && options.max_open_files == 2000)
140 options.max_open_files = 8000;
142 fdMinRequired = options.max_open_files + 128;
145 if (keyValues.exists(Keys::kFileSizeMb))
147 auto fileSizeMb = get<int>(keyValues, Keys::kFileSizeMb);
149 if (!hardSet && fileSizeMb == 8)
152 options.target_file_size_base = megabytes(fileSizeMb);
153 options.max_bytes_for_level_base = 5 * options.target_file_size_base;
154 options.write_buffer_size = 2 * options.target_file_size_base;
157 getIfExists(keyValues, Keys::kFileSizeMult, options.target_file_size_multiplier);
159 if (keyValues.exists(Keys::kBgThreads))
161 options.env->SetBackgroundThreads(
162 get<int>(keyValues, Keys::kBgThreads), rocksdb::Env::LOW);
165 if (keyValues.exists(Keys::kHighThreads))
167 auto const highThreads = get<int>(keyValues, Keys::kHighThreads);
168 options.env->SetBackgroundThreads(highThreads, rocksdb::Env::HIGH);
173 options.max_background_flushes = highThreads;
176 options.compression = rocksdb::kSnappyCompression;
178 getIfExists(keyValues, Keys::kBlockSize, tableOptions.block_size);
180 if (keyValues.exists(Keys::kUniversalCompaction) &&
181 (get<int>(keyValues, Keys::kUniversalCompaction) != 0))
183 options.compaction_style = rocksdb::kCompactionStyleUniversal;
184 options.min_write_buffer_number_to_merge = 2;
185 options.max_write_buffer_number = 6;
186 options.write_buffer_size = 6 * options.target_file_size_base;
189 if (keyValues.exists(Keys::kBbtOptions))
191 rocksdb::ConfigOptions const configOptions;
192 auto const s = rocksdb::GetBlockBasedTableOptionsFromString(
193 configOptions, tableOptions, get(keyValues, Keys::kBbtOptions), &tableOptions);
196 Throw<std::runtime_error>(
197 std::format(
"Unable to set RocksDB bbt_options: {}", s.ToString()));
201 options.table_factory.reset(NewBlockBasedTableFactory(tableOptions));
203 if (keyValues.exists(Keys::kOptions))
206 rocksdb::GetOptionsFromString(options, get(keyValues, Keys::kOptions), &options);
209 Throw<std::runtime_error>(
210 std::format(
"Unable to set RocksDB options: {}", s.ToString()));
215 rocksdb::GetStringFromDBOptions(&s1, options,
"; ");
216 rocksdb::GetStringFromColumnFamilyOptions(&s2, options,
"; ");
217 JLOG(journal.
debug()) <<
"RocksDB DBOptions: " << s1;
218 JLOG(journal.
debug()) <<
"RocksDB CFOptions: " << s2;
221 ~RocksDBBackend()
override
227 open(
bool createIfMissing)
override
233 "xrpl::node_store::RocksDBBackend::open : database is already "
235 JLOG(journal.
error()) <<
"database is already open";
239 rocksdb::DB* localDb =
nullptr;
240 options.create_if_missing = createIfMissing;
241 rocksdb::Status
const status = rocksdb::DB::Open(options, name, &localDb);
242 if (!
status.ok() || (localDb ==
nullptr))
244 Throw<std::runtime_error>(
253 return static_cast<bool>(db);
281 XRPL_ASSERT(db,
"xrpl::node_store::RocksDBBackend::fetch : non-null database");
286 rocksdb::ReadOptions
const options;
287 rocksdb::Slice
const slice(
reinterpret_cast<char const*
>(hash.data()), keyBytes);
291 rocksdb::Status
const getStatus = db->Get(options, slice, &
string);
295 DecodedBlob decoded(hash.data(),
string.data(),
string.size());
299 *pObject = decoded.createObject();
305 status = Status::DataCorrupt;
310 if (getStatus.IsCorruption())
312 status = Status::DataCorrupt;
314 else if (getStatus.IsNotFound())
316 status = Status::NotFound;
321 static_cast<int>(Status::CustomCode) + unsafeCast<int>(getStatus.code()));
323 JLOG(journal.
error()) << getStatus.ToString();
341 "xrpl::node_store::RocksDBBackend::storeBatch : non-null "
343 rocksdb::WriteBatch wb;
345 for (
auto const& e : batch)
347 EncodedBlob
const encoded(e);
350 rocksdb::Slice(
reinterpret_cast<char const*
>(encoded.getKey()), keyBytes),
352 reinterpret_cast<char const*
>(encoded.getData()), encoded.getSize()));
355 rocksdb::WriteOptions
const options;
357 auto ret = db->Write(options, &wb);
360 Throw<std::runtime_error>(
std::format(
"storeBatch failed: {}", ret.ToString()));
371 XRPL_ASSERT(db,
"xrpl::node_store::RocksDBBackend::forEach : non-null database");
372 rocksdb::ReadOptions
const options;
376 for (it->SeekToFirst(); it->Valid(); it->Next())
378 if (it->key().size() == keyBytes)
380 DecodedBlob decoded(it->key().data(), it->value().data(), it->value().size());
384 f(decoded.createObject());
389 JLOG(journal.
fatal()) <<
"Corrupt NodeObject #" << it->key().ToString(
true);
396 JLOG(journal.
fatal()) <<
"Bad key size = " << it->key().size();
402 getWriteLoad()
override
404 return batch.getWriteLoad();
408 setDeletePath()
override
416 writeBatch(Batch
const& batch)
override
425 fdRequired()
const override
427 return fdMinRequired;
433class RocksDBFactory :
public Factory
441 RocksDBFactory(Manager& manager) : manager_(manager)
443 manager_.insert(*
this);
446 [[nodiscard]] std::string
447 getName()
const override
452 std::unique_ptr<Backend>
455 Section
const& keyValues,
457 Scheduler& scheduler,
458 beast::Journal journal)
override
465registerRocksDBFactory(Manager& manager)
467 static RocksDBFactory
const kInstance{manager};
A backend used for the NodeStore.
void setCurrentThreadName(std::string_view newThreadName)
Changes the name of the caller thread.
void storeBatch(Backend &backend, Batch const &batch)
Status
Return codes from Backend operations.
bool getIfExists(Section const §ion, std::string const &name, T &v)
void open(soci::session &s, BasicConfig const &config, std::string const &dbName)
Open a soci session.
This callback does the actual writing.