Skip to content

Commit 54b946c

Browse files
Implement compaction
Add methods to borrow and steal internal memory. Move compression enum
1 parent 3f8fc66 commit 54b946c

13 files changed

Lines changed: 300 additions & 84 deletions

src/amulet/rocksdb/__init__.pyi

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@ import typing
77
from . import _rocksdb, _version
88

99
__all__: list[str] = [
10+
"CompactRangeOptions",
1011
"CompressionType",
1112
"Options",
1213
"ReadOptions",
@@ -16,6 +17,9 @@ __all__: list[str] = [
1617
"compiler_config",
1718
]
1819

20+
class CompactRangeOptions:
21+
def __init__(self) -> None: ...
22+
1923
class CompressionType:
2024
"""
2125
Members:
@@ -86,6 +90,7 @@ class RocksDB:
8690
options: Options,
8791
read_options: ReadOptions,
8892
write_options: WriteOptions,
93+
compact_range_options: CompactRangeOptions,
8994
) -> None:
9095
"""
9196
Construct a new :class:`RocksDB` instance from the database at the given path.
@@ -96,6 +101,7 @@ class RocksDB:
96101
:param options: The RocksDB Options object.
97102
:param read_options: The RocksDB ReadOptions object.
98103
:param write_options: The RocksDB WriteOptions object.
104+
:param compact_range_options: The RocksDB CompactRangeOptions object.
99105
:raises: RocksDBException if an error occured.
100106
"""
101107

@@ -107,6 +113,16 @@ class RocksDB:
107113
If needed, an external lock must be used to ensure that no other threads are accessing the database.
108114
"""
109115

116+
def compact(self) -> None:
117+
"""
118+
Remove deleted entries from the database to reduce its size.
119+
"""
120+
121+
def compact_range(self, arg0: str | None, arg1: str | None) -> None:
122+
"""
123+
Remove deleted entries from the database to reduce its size.
124+
"""
125+
110126
def delete(self, key: bytes) -> None:
111127
"""
112128
Delete a key from the database.

src/amulet/rocksdb/_rocksdb.py.cpp

Lines changed: 33 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@
2121
#include <amulet/pybind11_extensions/iterator.hpp>
2222

2323
// #include <amulet/rocksdb.hpp>
24+
#include <amulet/rocksdb/compact_range_options.hpp>
2425
#include <amulet/rocksdb/db.hpp>
2526
#include <amulet/rocksdb/options.hpp>
2627
#include <amulet/rocksdb/read_options.hpp>
@@ -243,6 +244,10 @@ void init_module(py::module m)
243244
&Amulet::RocksDB::WriteOptions::get_disable_wal,
244245
&Amulet::RocksDB::WriteOptions::set_disable_wal);
245246

247+
py::classh<Amulet::RocksDB::CompactRangeOptions> CompactRangeOptions(m, "CompactRangeOptions");
248+
CompactRangeOptions.def(
249+
py::init());
250+
246251
// py::classh<Amulet::RocksDBIterator> RocksDBIterator(m, "RocksDBIterator", py::release_gil_before_calling_cpp_dtor());
247252
// RocksDBIterator.def(
248253
// "valid",
@@ -364,15 +369,26 @@ void init_module(py::module m)
364369
":param compression_type: The compression type to use. (Default ZStandardCompression)\n"
365370
":raises: RocksDBException if an error occured."));
366371
RocksDB.def(
367-
py::init<
368-
std::filesystem::path,
369-
const Amulet::RocksDB::Options&,
370-
const Amulet::RocksDB::ReadOptions&,
371-
const Amulet::RocksDB::WriteOptions&>(),
372+
py::init([](
373+
std::filesystem::path path,
374+
Amulet::RocksDB::Options& options,
375+
Amulet::RocksDB::ReadOptions& read_options,
376+
Amulet::RocksDB::WriteOptions& write_options,
377+
Amulet::RocksDB::CompactRangeOptions& compact_range_options
378+
) {
379+
return std::make_unique<Amulet::RocksDB::RocksDB>(
380+
std::move(path),
381+
std::move(options),
382+
std::move(read_options),
383+
std::move(write_options),
384+
std::move(compact_range_options)
385+
);
386+
}),
372387
py::arg("path"),
373388
py::arg("options"),
374389
py::arg("read_options"),
375390
py::arg("write_options"),
391+
py::arg("compact_range_options"),
376392
py::doc(
377393
"Construct a new :class:`RocksDB` instance from the database at the given path.\n"
378394
"\n"
@@ -382,6 +398,7 @@ void init_module(py::module m)
382398
":param options: The RocksDB Options object.\n"
383399
":param read_options: The RocksDB ReadOptions object.\n"
384400
":param write_options: The RocksDB WriteOptions object.\n"
401+
":param compact_range_options: The RocksDB CompactRangeOptions object.\n"
385402
":raises: RocksDBException if an error occured."));
386403

387404
RocksDB.def(
@@ -393,16 +410,17 @@ void init_module(py::module m)
393410
"If needed, an external lock must be used to ensure that no other threads are accessing the database."),
394411
py::call_guard<py::gil_scoped_release>());
395412

396-
// RocksDB.def(
397-
// "compact",
398-
// [](Amulet::RocksDB& self) {
399-
// if (!self) {
400-
// throw std::runtime_error("The RocksDB database has been closed.");
401-
// }
402-
// //self->CompactRange(nullptr, nullptr);
403-
// },
404-
// py::doc("Remove deleted entries from the database to reduce its size."),
405-
// py::call_guard<py::gil_scoped_release>());
413+
RocksDB.def(
414+
"compact_range",
415+
&Amulet::RocksDB::RocksDB::compact_range,
416+
py::doc("Remove deleted entries from the database to reduce its size."),
417+
py::call_guard<py::gil_scoped_release>());
418+
419+
RocksDB.def(
420+
"compact",
421+
&Amulet::RocksDB::RocksDB::compact,
422+
py::doc("Remove deleted entries from the database to reduce its size."),
423+
py::call_guard<py::gil_scoped_release>());
406424

407425
auto put = [](Amulet::RocksDB::RocksDB& self, py::bytes key, py::bytes value) {
408426
std::string_view key_view = key;
Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,43 @@
1+
#include <rocksdb/options.h>
2+
3+
#include "compact_range_options.hpp"
4+
5+
namespace Amulet {
6+
namespace RocksDB {
7+
8+
CompactRangeOptions::CompactRangeOptions()
9+
: _impl(new ROCKSDB_NAMESPACE::CompactRangeOptions())
10+
{
11+
}
12+
13+
CompactRangeOptions::~CompactRangeOptions()
14+
{
15+
delete _impl;
16+
}
17+
18+
CompactRangeOptions::CompactRangeOptions(CompactRangeOptions&& other)
19+
{
20+
_impl = other.steal();
21+
}
22+
23+
CompactRangeOptions& CompactRangeOptions::operator=(CompactRangeOptions&& other)
24+
{
25+
delete _impl;
26+
_impl = other.steal();
27+
return *this;
28+
}
29+
30+
ROCKSDB_NAMESPACE::CompactRangeOptions* CompactRangeOptions::steal()
31+
{
32+
auto* impl = _impl;
33+
_impl = new ROCKSDB_NAMESPACE::CompactRangeOptions();
34+
return impl;
35+
}
36+
37+
ROCKSDB_NAMESPACE::CompactRangeOptions* CompactRangeOptions::borrow()
38+
{
39+
return _impl;
40+
}
41+
42+
} // namespace RocksDB
43+
} // namespace Amulet
Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,37 @@
1+
#pragma once
2+
3+
#include "export.hpp"
4+
5+
namespace ROCKSDB_NAMESPACE {
6+
struct CompactRangeOptions;
7+
}
8+
9+
namespace Amulet {
10+
namespace RocksDB {
11+
class RocksDB;
12+
13+
class AMULET_ROCKSDB_EXPORT CompactRangeOptions {
14+
private:
15+
ROCKSDB_NAMESPACE::CompactRangeOptions* _impl;
16+
17+
public:
18+
CompactRangeOptions();
19+
~CompactRangeOptions();
20+
21+
CompactRangeOptions(const CompactRangeOptions&) = delete;
22+
CompactRangeOptions& operator=(const CompactRangeOptions&) = delete;
23+
CompactRangeOptions(CompactRangeOptions&&);
24+
CompactRangeOptions& operator=(CompactRangeOptions&&);
25+
26+
// Steal the internal pointer
27+
// It is your responsibiliy to delete the pointer
28+
// This allocates a new internal pointer
29+
ROCKSDB_NAMESPACE::CompactRangeOptions* steal();
30+
31+
// Borrow the internal pointer
32+
// The pointer is still managed by this object
33+
ROCKSDB_NAMESPACE::CompactRangeOptions* borrow();
34+
};
35+
36+
} // namespace RocksDB
37+
} // namespace Amulet

src/amulet/rocksdb/compression.hpp

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,10 @@
1+
#pragma once
2+
3+
namespace Amulet {
4+
namespace RocksDB {
5+
enum class CompressionType : unsigned char {
6+
NoCompression = 0x00,
7+
ZStandardCompression = 0x07,
8+
};
9+
} // namespace RocksDB
10+
} // namespace Amulet

src/amulet/rocksdb/db.cpp

Lines changed: 59 additions & 32 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@
55

66
#include <rocksdb/db.h>
77

8+
#include "compact_range_options.hpp"
89
#include "db.hpp"
910
#include "options.hpp"
1011
#include "read_options.hpp"
@@ -17,32 +18,33 @@ namespace RocksDB {
1718
ROCKSDB_NAMESPACE::DB* db;
1819
ROCKSDB_NAMESPACE::ReadOptions read_options;
1920
ROCKSDB_NAMESPACE::WriteOptions write_options;
20-
// ROCKSDB_NAMESPACE::CompactRangeOptions compact_range_options;
21+
ROCKSDB_NAMESPACE::CompactRangeOptions compact_range_options;
2122
};
2223

2324
RocksDB::RocksDB(
2425
std::filesystem::path path,
25-
const Options& options,
26-
const ReadOptions& read_options,
27-
const WriteOptions& write_options)
28-
// CompactRangeOptions& compact_range_options)
26+
Options&& options,
27+
ReadOptions&& read_options,
28+
WriteOptions&& write_options,
29+
CompactRangeOptions&& compact_range_options)
2930
{
3031
// Expand dots and symbolic links
3132
path = std::filesystem::absolute(path);
3233

3334
std::unique_ptr<ROCKSDB_NAMESPACE::DB> db = nullptr;
34-
auto status = ROCKSDB_NAMESPACE::DB::Open(*options._impl, path.string(), &db);
35+
auto status = ROCKSDB_NAMESPACE::DB::Open(*options.borrow(), path.string(), &db);
3536
if (status.IsCorruption()) {
36-
ROCKSDB_NAMESPACE::RepairDB(path.string(), *options._impl);
37-
status = ROCKSDB_NAMESPACE::DB::Open(*options._impl, path.string(), &db);
37+
ROCKSDB_NAMESPACE::RepairDB(path.string(), *options.borrow());
38+
status = ROCKSDB_NAMESPACE::DB::Open(*options.borrow(), path.string(), &db);
3839
}
3940
if (!status.ok()) {
4041
throw RocksDBException(status.ToString());
4142
}
4243
_impl = new RocksDBImpl {
4344
db.release(),
44-
*read_options._impl,
45-
*write_options._impl,
45+
*read_options.steal(),
46+
*write_options.steal(),
47+
*compact_range_options.steal()
4648
};
4749
}
4850

@@ -55,11 +57,12 @@ namespace RocksDB {
5557
[create_if_missing, compression_type]() {
5658
Options options;
5759
options.set_create_if_missing(create_if_missing);
58-
options._impl->compression = static_cast<ROCKSDB_NAMESPACE::CompressionType>(compression_type);
60+
options.set_compression_type(compression_type);
5961
return options;
6062
}(),
6163
ReadOptions(),
62-
WriteOptions())
64+
WriteOptions(),
65+
CompactRangeOptions())
6366
{
6467
Options options;
6568
ReadOptions read_options;
@@ -166,27 +169,51 @@ namespace RocksDB {
166169
}
167170
}
168171

169-
// void RocksDB::compact_range(std::optional<std::string_view> begin, std::optional<std::string_view> end)
170-
//{
171-
// if (!_impl) {
172-
// throw std::runtime_error("RocksDB has been closed");
173-
// }
174-
// _impl->db->CompactRange(
175-
// _impl->compact_range_options,
176-
// begin ? *begin : nullptr,
177-
// end ? *end : nullptr);
178-
// }
172+
void RocksDB::compact_range(std::optional<std::string_view> begin, std::optional<std::string_view> end)
173+
{
174+
if (!_impl) {
175+
throw std::runtime_error("RocksDB has been closed");
176+
}
177+
if (begin) {
178+
ROCKSDB_NAMESPACE::Slice begin_slice = *begin;
179+
if (end) {
180+
ROCKSDB_NAMESPACE::Slice end_slice = *end;
181+
_impl->db->CompactRange(
182+
_impl->compact_range_options,
183+
&begin_slice,
184+
&end_slice);
185+
} else {
186+
_impl->db->CompactRange(
187+
_impl->compact_range_options,
188+
&begin_slice,
189+
nullptr);
190+
}
191+
} else {
192+
if (end) {
193+
ROCKSDB_NAMESPACE::Slice end_slice = *end;
194+
_impl->db->CompactRange(
195+
_impl->compact_range_options,
196+
nullptr,
197+
&end_slice);
198+
} else {
199+
_impl->db->CompactRange(
200+
_impl->compact_range_options,
201+
nullptr,
202+
nullptr);
203+
}
204+
}
205+
}
179206

180-
// void RocksDB::compact()
181-
//{
182-
// if (!_impl) {
183-
// throw std::runtime_error("RocksDB has been closed");
184-
// }
185-
// _impl->db->CompactRange(
186-
// _impl->compact_range_options,
187-
// nullptr,
188-
// nullptr);
189-
// }
207+
void RocksDB::compact()
208+
{
209+
if (!_impl) {
210+
throw std::runtime_error("RocksDB has been closed");
211+
}
212+
_impl->db->CompactRange(
213+
_impl->compact_range_options,
214+
nullptr,
215+
nullptr);
216+
}
190217

191218
} // namespace RocksDB
192219
} // namespace Amulet

0 commit comments

Comments
 (0)