xrpld
Loading...
Searching...
No Matches
RocksDBFactory.cpp
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>
21
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>
34
35#include <atomic>
36#include <cstddef>
37#include <filesystem>
38#include <format>
39#include <functional>
40#include <memory>
41#include <stdexcept>
42#include <string>
43
44namespace xrpl::node_store {
45
46class RocksDBEnv : public rocksdb::EnvWrapper
47{
48public:
49 RocksDBEnv() : EnvWrapper(rocksdb::Env::Default())
50 {
51 }
52
53 struct ThreadParams
54 {
55 ThreadParams(void (*f)(void*), void* a) : f(f), a(a)
56 {
57 }
58
59 void (*f)(void*);
60 void* a;
61 };
62
63 static void
64 threadEntry(void* ptr)
65 {
66 ThreadParams const* const p(reinterpret_cast<ThreadParams*>(ptr));
67 auto const f = p->f;
68
69 void* a(p->a);
70 delete p;
71
72 static std::atomic<std::size_t> kN;
73 std::size_t const id(++kN);
75
76 f(a);
77 }
78
79 void
80 StartThread(void (*f)(void*), void* a) override
81 {
82 auto* const p = new ThreadParams(f, a);
83 EnvWrapper::StartThread(&RocksDBEnv::threadEntry, p);
84 }
85};
86
87//------------------------------------------------------------------------------
88
89class RocksDBBackend : public Backend, public BatchWriter::Callback
90{
91private:
92 std::atomic<bool> deletePath_;
93
94public:
95 beast::Journal journal;
96 size_t const keyBytes;
97 BatchWriter batch;
98 std::string name;
99 std::unique_ptr<rocksdb::DB> db;
100 int fdMinRequired = 2048;
101 rocksdb::Options options;
102
103 RocksDBBackend(
104 int keyBytes,
105 Section const& keyValues,
106 Scheduler& scheduler,
107 beast::Journal journal,
108 RocksDBEnv* env)
109 : deletePath_(false), journal(journal), keyBytes(keyBytes), batch(*this, scheduler)
110 {
111 if (!getIfExists(keyValues, Keys::kPath, name))
112 Throw<std::runtime_error>("Missing path in RocksDBFactory backend");
113
114 rocksdb::BlockBasedTableOptions tableOptions;
115 options.env = env;
116
117 bool const hardSet =
118 keyValues.exists(Keys::kHardSet) && get<bool>(keyValues, Keys::kHardSet);
119
120 if (keyValues.exists(Keys::kCacheMb))
121 {
122 auto size = get<int>(keyValues, Keys::kCacheMb);
123
124 if (!hardSet && size == 256)
125 size = 1024;
126
127 tableOptions.block_cache = rocksdb::NewLRUCache(megabytes(size));
128 }
129
130 if (auto const v = get<int>(keyValues, Keys::kFilterBits))
131 {
132 bool const filterBlocks = !keyValues.exists(Keys::kFilterFull) ||
133 (get<int>(keyValues, Keys::kFilterFull) == 0);
134 tableOptions.filter_policy.reset(rocksdb::NewBloomFilterPolicy(v, filterBlocks));
135 }
136
137 if (getIfExists(keyValues, Keys::kOpenFiles, options.max_open_files))
138 {
139 if (!hardSet && options.max_open_files == 2000)
140 options.max_open_files = 8000;
141
142 fdMinRequired = options.max_open_files + 128;
143 }
144
145 if (keyValues.exists(Keys::kFileSizeMb))
146 {
147 auto fileSizeMb = get<int>(keyValues, Keys::kFileSizeMb);
148
149 if (!hardSet && fileSizeMb == 8)
150 fileSizeMb = 256;
151
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;
155 }
156
157 getIfExists(keyValues, Keys::kFileSizeMult, options.target_file_size_multiplier);
158
159 if (keyValues.exists(Keys::kBgThreads))
160 {
161 options.env->SetBackgroundThreads(
162 get<int>(keyValues, Keys::kBgThreads), rocksdb::Env::LOW);
163 }
164
165 if (keyValues.exists(Keys::kHighThreads))
166 {
167 auto const highThreads = get<int>(keyValues, Keys::kHighThreads);
168 options.env->SetBackgroundThreads(highThreads, rocksdb::Env::HIGH);
169
170 // If we have high-priority threads, presumably we want to
171 // use them for background flushes
172 if (highThreads > 0)
173 options.max_background_flushes = highThreads;
174 }
175
176 options.compression = rocksdb::kSnappyCompression;
177
178 getIfExists(keyValues, Keys::kBlockSize, tableOptions.block_size);
179
180 if (keyValues.exists(Keys::kUniversalCompaction) &&
181 (get<int>(keyValues, Keys::kUniversalCompaction) != 0))
182 {
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;
187 }
188
189 if (keyValues.exists(Keys::kBbtOptions))
190 {
191 rocksdb::ConfigOptions const configOptions;
192 auto const s = rocksdb::GetBlockBasedTableOptionsFromString(
193 configOptions, tableOptions, get(keyValues, Keys::kBbtOptions), &tableOptions);
194 if (!s.ok())
195 {
196 Throw<std::runtime_error>(
197 std::format("Unable to set RocksDB bbt_options: {}", s.ToString()));
198 }
199 }
200
201 options.table_factory.reset(NewBlockBasedTableFactory(tableOptions));
202
203 if (keyValues.exists(Keys::kOptions))
204 {
205 auto const s =
206 rocksdb::GetOptionsFromString(options, get(keyValues, Keys::kOptions), &options);
207 if (!s.ok())
208 {
209 Throw<std::runtime_error>(
210 std::format("Unable to set RocksDB options: {}", s.ToString()));
211 }
212 }
213
214 std::string s1, s2;
215 rocksdb::GetStringFromDBOptions(&s1, options, "; ");
216 rocksdb::GetStringFromColumnFamilyOptions(&s2, options, "; ");
217 JLOG(journal.debug()) << "RocksDB DBOptions: " << s1;
218 JLOG(journal.debug()) << "RocksDB CFOptions: " << s2;
219 }
220
221 ~RocksDBBackend() override
222 {
223 close();
224 }
225
226 void
227 open(bool createIfMissing) override
228 {
229 if (db)
230 {
231 // LCOV_EXCL_START
232 UNREACHABLE(
233 "xrpl::node_store::RocksDBBackend::open : database is already "
234 "open");
235 JLOG(journal.error()) << "database is already open";
236 return;
237 // LCOV_EXCL_STOP
238 }
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))
243 {
244 Throw<std::runtime_error>(
245 std::format("Unable to open/create RocksDB: {}", status.ToString()));
246 }
247 db.reset(localDb);
248 }
249
250 bool
251 isOpen() override
252 {
253 return static_cast<bool>(db);
254 }
255
256 void
257 close() override
258 {
259 if (db)
260 {
261 db.reset();
262 if (deletePath_)
263 {
264 std::filesystem::path const dir = name;
266 }
267 }
268 }
269
271 getName() override
272 {
273 return name;
274 }
275
276 //--------------------------------------------------------------------------
277
278 Status
279 fetch(UInt256 const& hash, std::shared_ptr<NodeObject>* pObject) override
280 {
281 XRPL_ASSERT(db, "xrpl::node_store::RocksDBBackend::fetch : non-null database");
282 pObject->reset();
283
284 Status status = Status::Ok;
285
286 rocksdb::ReadOptions const options;
287 rocksdb::Slice const slice(reinterpret_cast<char const*>(hash.data()), keyBytes);
288
289 std::string string;
290
291 rocksdb::Status const getStatus = db->Get(options, slice, &string);
292
293 if (getStatus.ok())
294 {
295 DecodedBlob decoded(hash.data(), string.data(), string.size());
296
297 if (decoded.wasOk())
298 {
299 *pObject = decoded.createObject();
300 }
301 else
302 {
303 // Decoding failed, probably corrupted!
304 //
305 status = Status::DataCorrupt;
306 }
307 }
308 else
309 {
310 if (getStatus.IsCorruption())
311 {
312 status = Status::DataCorrupt;
313 }
314 else if (getStatus.IsNotFound())
315 {
316 status = Status::NotFound;
317 }
318 else
319 {
320 status = static_cast<Status>(
321 static_cast<int>(Status::CustomCode) + unsafeCast<int>(getStatus.code()));
322
323 JLOG(journal.error()) << getStatus.ToString();
324 }
325 }
326
327 return status;
328 }
329
330 void
331 store(std::shared_ptr<NodeObject> const& object) override
332 {
333 batch.store(object);
334 }
335
336 void
337 storeBatch(Batch const& batch) override
338 {
339 XRPL_ASSERT(
340 db,
341 "xrpl::node_store::RocksDBBackend::storeBatch : non-null "
342 "database");
343 rocksdb::WriteBatch wb;
344
345 for (auto const& e : batch)
346 {
347 EncodedBlob const encoded(e);
348
349 wb.Put(
350 rocksdb::Slice(reinterpret_cast<char const*>(encoded.getKey()), keyBytes),
351 rocksdb::Slice(
352 reinterpret_cast<char const*>(encoded.getData()), encoded.getSize()));
353 }
354
355 rocksdb::WriteOptions const options;
356
357 auto ret = db->Write(options, &wb);
358
359 if (!ret.ok())
360 Throw<std::runtime_error>(std::format("storeBatch failed: {}", ret.ToString()));
361 }
362
363 void
364 sync() override
365 {
366 }
367
368 void
369 forEach(std::function<void(std::shared_ptr<NodeObject>)> f) override
370 {
371 XRPL_ASSERT(db, "xrpl::node_store::RocksDBBackend::forEach : non-null database");
372 rocksdb::ReadOptions const options;
373
374 std::unique_ptr<rocksdb::Iterator> it(db->NewIterator(options));
375
376 for (it->SeekToFirst(); it->Valid(); it->Next())
377 {
378 if (it->key().size() == keyBytes)
379 {
380 DecodedBlob decoded(it->key().data(), it->value().data(), it->value().size());
381
382 if (decoded.wasOk())
383 {
384 f(decoded.createObject());
385 }
386 else
387 {
388 // Uh oh, corrupted data!
389 JLOG(journal.fatal()) << "Corrupt NodeObject #" << it->key().ToString(true);
390 }
391 }
392 else
393 {
394 // VFALCO NOTE What does it mean to find an
395 // incorrectly sized key? Corruption?
396 JLOG(journal.fatal()) << "Bad key size = " << it->key().size();
397 }
398 }
399 }
400
401 int
402 getWriteLoad() override
403 {
404 return batch.getWriteLoad();
405 }
406
407 void
408 setDeletePath() override
409 {
410 deletePath_ = true;
411 }
412
413 //--------------------------------------------------------------------------
414
415 void
416 writeBatch(Batch const& batch) override
417 {
418 storeBatch(batch);
419 }
420
424 [[nodiscard]] int
425 fdRequired() const override
426 {
427 return fdMinRequired;
428 }
429};
430
431//------------------------------------------------------------------------------
432
433class RocksDBFactory : public Factory
434{
435private:
436 Manager& manager_;
437
438public:
439 RocksDBEnv env;
440
441 RocksDBFactory(Manager& manager) : manager_(manager)
442 {
443 manager_.insert(*this);
444 }
445
446 [[nodiscard]] std::string
447 getName() const override
448 {
449 return "RocksDB";
450 }
451
452 std::unique_ptr<Backend>
453 createInstance(
454 size_t keyBytes,
455 Section const& keyValues,
456 std::size_t,
457 Scheduler& scheduler,
458 beast::Journal journal) override
459 {
460 return std::make_unique<RocksDBBackend>(keyBytes, keyValues, scheduler, journal, &env);
461 }
462};
463
464void
465registerRocksDBFactory(Manager& manager)
466{
467 static RocksDBFactory const kInstance{manager};
468}
469
470} // namespace xrpl::node_store
471
472#endif
Stream fatal() const
Definition Journal.h:368
Stream error() const
Definition Journal.h:362
Stream debug() const
Definition Journal.h:344
A backend used for the NodeStore.
Definition Backend.h:29
T format(T... args)
T make_unique(T... args)
void setCurrentThreadName(std::string_view newThreadName)
Changes the name of the caller thread.
void storeBatch(Backend &backend, Batch const &batch)
Definition TestBase.h:85
Status
Return codes from Backend operations.
bool getIfExists(Section const &section, std::string const &name, T &v)
void open(soci::session &s, BasicConfig const &config, std::string const &dbName)
Open a soci session.
Definition SociDB.cpp:91
T remove_all(T... args)
T reset(T... args)
This callback does the actual writing.
Definition BatchWriter.h:30
T to_string(T... args)