diff --git a/CMakeLists.txt b/CMakeLists.txt index a417f9a..0df2727 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -410,6 +410,8 @@ endif() add_library(goblin_core src/auth.cpp + src/exanic_fabric_model.cpp + src/exanic_pubsub_plane.cpp src/command.cpp src/pubsub.cpp src/luau_script.cpp @@ -1123,6 +1125,20 @@ if(GOBLIN_CORE_BUILD_TESTS) CXX_STANDARD 23 CXX_STANDARD_REQUIRED ON CXX_EXTENSIONS OFF) add_test(NAME goblin_core_pubsub COMMAND goblin_core_pubsub_test) + add_executable(goblin_core_exanic_fabric_test tests/exanic_fabric_test.cpp) + target_link_libraries(goblin_core_exanic_fabric_test PRIVATE goblin_core) + set_target_properties(goblin_core_exanic_fabric_test PROPERTIES + CXX_STANDARD 23 CXX_STANDARD_REQUIRED ON CXX_EXTENSIONS OFF) + add_test(NAME goblin_core_exanic_fabric + COMMAND goblin_core_exanic_fabric_test + ${CMAKE_CURRENT_BINARY_DIR}/exanic_fabric_vectors.txt) + + add_executable(goblin_core_exanic_plane_test tests/exanic_plane_test.cpp) + target_link_libraries(goblin_core_exanic_plane_test PRIVATE goblin_core) + set_target_properties(goblin_core_exanic_plane_test PROPERTIES + CXX_STANDARD 23 CXX_STANDARD_REQUIRED ON CXX_EXTENSIONS OFF) + add_test(NAME goblin_core_exanic_plane COMMAND goblin_core_exanic_plane_test) + if(GOBLIN_HAS_KAFKA_BUILD) add_executable(goblin_core_kafka_ingest_test tests/kafka_ingest_test.cpp) target_link_libraries(goblin_core_kafka_ingest_test PRIVATE goblin_core) diff --git a/README.md b/README.md index c3ec041..d264cc6 100644 --- a/README.md +++ b/README.md @@ -603,6 +603,12 @@ threshold to `ceil(leaf_capacity ** k)`, with `k` in `[0, 1]` and a default of leaf when its threshold is reached. See [fixed-width packed sorted sets](docs/packed-zsets.md). +`--packed-zset-score-rle` compresses runs of four or more equal scores in packed +B+ tree leaf bases. It applies to all packed member/score types and defaults to +off; `--no-packed-zset-score-rle` disables it explicitly. Dirty tails and logical +leaf capacities are preserved. Snapshots use the receiving server's setting. +`GOBLIN.MEMORY` reports the setting, compressed leaf count, and used score bytes. + `--zset-implementation standard|packed-int32-float32|packed-int32-float64|packed-int64-float32|packed-int64-float64|packed-uuid-float32|packed-uuid-float64` selects the representation created by unqualified sorted-set commands. It defaults to `standard`. Existing and restored keys retain their representation; diff --git a/docs/packed-zsets.md b/docs/packed-zsets.md index 2e3addc..48d6e50 100644 --- a/docs/packed-zsets.md +++ b/docs/packed-zsets.md @@ -123,6 +123,48 @@ all dirty leaves to compact, while `GOBLIN.MEMORY key` reports the representation, exponent, per-leaf threshold, and aggregate sorted/dirty entry counts. +## Optional score run-length encoding + +Enable sorted-base score compression with: + +```text +goblin-core --zset-implementation packed-int32-float32 --packed-zset-score-rle +``` + +The default is off; `--no-packed-zset-score-rle` explicitly disables it. The +option applies to all six packed representations, including keys created by +representation-qualified commands, copies, aggregate stores, and snapshot +restores. Ordinary sorted sets retain their existing storage. Like the merge +exponent, this is server policy: snapshots retain the member/score type and +logical contents, and use the receiving server's compression setting. + +Each compressed-layout leaf stores its sorted member IDs separately from its +score stream. Four or more adjacent equal scores become `[NaN, count, score]`. +All three fields occupy one score-width word; the count is a binary `uint32` +for FLOAT32 or `uint64` for FLOAT64. Shorter runs remain literal scores. NaN +markers are internal and never appear in replies. Dirty tails keep their +ordinary tuples and NaN deletion markers in separate storage. + +Logical leaf capacities and merge thresholds stay the same. A checkpoint every +32 base positions bounds score decoding for binary-search probes; sequential +reads decode each run once. Compaction rebuilds one leaf of logical entries; +redistribution handles at most two leaves. All scratch space is bounded by the +fixed leaf capacity. Runs end at leaf boundaries, and rank counts and +invalidation bits still count members. + +Score buffers grow before a mutation is committed and can shrink as scores +become more repetitive. If memory pressure prevents optional redistribution +or shrinking, the existing leaves remain valid and maintenance can be retried +by later writes. The separately allocated buffers and checkpoints add overhead +for diverse scores; this option is intended for workloads with many ties and +adds encoding work during compaction. + +`GOBLIN.MEMORY key` additionally reports `score_rle` (0 or 1), +`compressed_leaf_count` (leaves currently containing encoded runs), and +`sorted_score_bytes` (used score-stream bytes, excluding reserved capacity, +members, checkpoints, and dirty records). `total_allocated_bytes` includes all +of those allocations and is the appropriate counter for comparing footprints. + Native snapshots serialize canonical logical entries, not append history, and restore the same packed representation. Representation-qualified names are part of the RESP command surface. The unqualified typed SBE sorted-set templates also diff --git a/include/goblin/core/packed_bplus_tree.hpp b/include/goblin/core/packed_bplus_tree.hpp index f352477..1ff5e3b 100644 --- a/include/goblin/core/packed_bplus_tree.hpp +++ b/include/goblin/core/packed_bplus_tree.hpp @@ -17,6 +17,7 @@ #include #include "goblin/core/memory_limit.hpp" +#include "goblin/core/packed_score_runs.hpp" namespace goblin::core::detail { @@ -25,7 +26,7 @@ namespace goblin::core::detail { // leaf overwrites its dirty record; reaching the local threshold compacts only // that leaf. Branch and leaf references are 32-bit arena indices, never heap // pointers. -template +template class PackedBPlusTree { public: using Key = typename Traits::Key; @@ -77,11 +78,35 @@ class PackedBPlusTree { std::size_t bytes = leaves_.capacity() * sizeof(Leaf) + branches_.capacity() * sizeof(Branch); for (const auto& leaf : leaves_) { - bytes += leaf.entries.capacity() * sizeof(Entry); + if constexpr (ScoreRle) bytes += leaf.entries.allocated_bytes(); + else bytes += leaf.entries.capacity() * sizeof(Entry); } return bytes; } + [[nodiscard]] static constexpr bool score_rle_enabled() noexcept { + return ScoreRle; + } + [[nodiscard]] std::size_t compressed_leaf_count() const noexcept { + std::size_t count = 0; + if constexpr (ScoreRle) { + for (const auto& leaf : leaves_) { + count += leaf.active && leaf.entries.encoded; + } + } + return count; + } + [[nodiscard]] std::size_t sorted_score_bytes() const noexcept { + if constexpr (!ScoreRle) return sorted_count_ * sizeof(Score); + else { + std::size_t bytes = 0; + for (const auto& leaf : leaves_) { + if (leaf.active) bytes += leaf.entries.scores.size() * sizeof(Score); + } + return bytes; + } + } + void insert(Entry entry) { if (root_ == kNull) { const auto leaf_id = allocate_leaf(); @@ -106,6 +131,7 @@ class PackedBPlusTree { } auto& leaf = leaves_[path.leaf]; + prepare_dirty(leaf); ++leaf.live_count; ++size_; set_dirty(leaf, entry.key, entry.score); @@ -122,6 +148,7 @@ class PackedBPlusTree { auto new_path = locate(new_entry); if (old_path.leaf == new_path.leaf) { auto& leaf = leaves_[old_path.leaf]; + prepare_dirty(leaf); set_dirty(leaf, new_entry.key, new_entry.score, old_entry); expand_fence(old_path, new_entry); maybe_compact_leaf(old_path.leaf); @@ -136,6 +163,8 @@ class PackedBPlusTree { auto& old_leaf = leaves_[old_path.leaf]; auto& new_leaf = leaves_[new_path.leaf]; + prepare_dirty(old_leaf); + prepare_dirty(new_leaf); --old_leaf.live_count; ++new_leaf.live_count; set_dirty(old_leaf, old_entry.key, std::nullopt, old_entry); @@ -152,6 +181,7 @@ class PackedBPlusTree { assert(root_ != kNull && size_ != 0); const auto path = locate(old_entry); auto& leaf = leaves_[path.leaf]; + prepare_dirty(leaf); --leaf.live_count; --size_; set_dirty(leaf, old_entry.key, std::nullopt, old_entry); @@ -303,22 +333,26 @@ class PackedBPlusTree { if (leaf_id >= leaves_.size()) return false; const auto& leaf = leaves_[leaf_id]; if (!leaf.active || leaf.prev != previous || !leaf.has_fence || - leaf.sorted_size > leaf.entries.size() || - leaf.entries.size() > kLeafCapacity + merge_threshold_) { + leaf.sorted_size + tail_size(leaf) > kLeafCapacity + merge_threshold_) { return false; } + if constexpr (ScoreRle) { + if (!leaf.entries.check_invariants(leaf.sorted_size) || + tail_size(leaf) >= merge_threshold_) return false; + } else if (leaf.sorted_size > leaf.entries.size()) return false; if (previous_fence && !less(*previous_fence, leaf.fence)) return false; previous_fence = leaf.fence; for (std::size_t i = 1; i < leaf.sorted_size; ++i) { - if (less(leaf.entries[i], leaf.entries[i - 1])) return false; + if (less(base_entry(leaf, i), base_entry(leaf, i - 1))) return false; } for (std::size_t i = 0; i < kLeafCapacity; ++i) { const bool shadowed = i < leaf.sorted_size && - tail_contains(leaf, leaf.entries[i].key); + tail_contains(leaf, base_entry(leaf, i).key); if (base_invalidated(leaf, i) != shadowed) return false; } - for (std::size_t i = leaf.sorted_size + 1; i < leaf.entries.size(); ++i) { - if (!Traits::less(leaf.entries[i - 1].key, leaf.entries[i].key)) { + const auto dirty_begin = tail_begin(leaf); + for (std::size_t i = 1; i < tail_size(leaf); ++i) { + if (!Traits::less(dirty_begin[i - 1].key, dirty_begin[i].key)) { return false; } } @@ -334,7 +368,7 @@ class PackedBPlusTree { if (!local_ok || local_live != leaf.live_count) return false; live_seen += local_live; sorted_seen += leaf.sorted_size; - dirty_seen += leaf.entries.size() - leaf.sorted_size; + dirty_seen += tail_size(leaf); previous = leaf_id; ++leaves_seen; } @@ -358,8 +392,9 @@ class PackedBPlusTree { std::numeric_limits::max(); static constexpr std::size_t kMaxHeight = 16; + using RleEntries = PackedScoreRuns; struct Leaf { - std::vector entries; + std::conditional_t> entries; // Base tuples stay unchanged until compaction, so their slots are stable. // One bit replaces a dirty-tail key search per base tuple on merges/reads. std::array invalidated{}; @@ -441,8 +476,9 @@ class PackedBPlusTree { [[nodiscard]] std::uint32_t allocate_leaf() { Leaf prepared; - reserve_memory_vector(prepared.entries, - kLeafCapacity + merge_threshold_); + if constexpr (ScoreRle) prepared.entries.reserve(merge_threshold_); + else reserve_memory_vector(prepared.entries, + kLeafCapacity + merge_threshold_); std::uint32_t id = kNull; if (free_leaf_ != kNull) { id = free_leaf_; @@ -566,29 +602,41 @@ class PackedBPlusTree { return {.leaf = node, .offset = rank}; } - [[nodiscard]] static auto tail_lower_bound(Leaf& leaf, const Key& key) { - return std::lower_bound( - leaf.entries.begin() + static_cast(leaf.sorted_size), - leaf.entries.end(), key, - [](const Entry& entry, const Key& candidate) { - return key_less(entry, candidate); - }); + [[nodiscard]] static auto tail_begin(auto& leaf) { + if constexpr (ScoreRle) return leaf.entries.dirty.begin(); + else return leaf.entries.begin() + + static_cast(leaf.sorted_size); } - - [[nodiscard]] static auto tail_lower_bound(const Leaf& leaf, - const Key& key) { - return std::lower_bound( - leaf.entries.begin() + static_cast(leaf.sorted_size), - leaf.entries.end(), key, + [[nodiscard]] static auto tail_end(auto& leaf) { + if constexpr (ScoreRle) return leaf.entries.dirty.end(); + else return leaf.entries.end(); + } + [[nodiscard]] static std::size_t tail_size(const Leaf& leaf) noexcept { + return static_cast(tail_end(leaf) - tail_begin(leaf)); + } + [[nodiscard]] static Entry base_entry(const Leaf& leaf, + std::size_t slot) noexcept { + if constexpr (ScoreRle) return leaf.entries.entry(slot); + else return leaf.entries[slot]; + } + [[nodiscard]] static auto tail_lower_bound(auto& leaf, const Key& key) { + return std::lower_bound(tail_begin(leaf), tail_end(leaf), key, [](const Entry& entry, const Key& candidate) { return key_less(entry, candidate); }); } - [[nodiscard]] static bool tail_contains(const Leaf& leaf, const Key& key) noexcept { const auto found = tail_lower_bound(leaf, key); - return found != leaf.entries.end() && found->key == key; + return found != tail_end(leaf) && found->key == key; + } + + void prepare_dirty(Leaf& leaf) { + if constexpr (ScoreRle) { + if (tail_size(leaf) + 1 >= merge_threshold_) { + leaf.entries.prepare_merge(tail_size(leaf) + 1); + } + } } [[nodiscard]] static bool base_invalidated(const Leaf& leaf, @@ -601,34 +649,50 @@ class PackedBPlusTree { std::optional retired = std::nullopt) { auto found = tail_lower_bound(leaf, key); const auto stored = score.value_or(std::numeric_limits::quiet_NaN()); - if (found != leaf.entries.end() && found->key == key) { + if (found != tail_end(leaf) && found->key == key) { found->score = stored; return; } if (retired) { // Only the first dirty record invalidates a base slot. Repeated updates, // deletes and reinsertions retain that bit until the leaf is compacted. - const auto base_end = leaf.entries.begin() + - static_cast(leaf.sorted_size); - const auto base = std::lower_bound(leaf.entries.begin(), base_end, *retired, - [](const Entry& a, const Entry& b) { - return less(a, b); - }); - assert(base != base_end && equivalent(*base, *retired)); - const auto slot = static_cast(base - leaf.entries.begin()); + std::size_t slot; + if constexpr (ScoreRle) { + std::size_t first = 0; + std::size_t end = leaf.sorted_size; + while (first < end) { + const auto middle = first + (end - first) / 2; + if (less(base_entry(leaf, middle), *retired)) first = middle + 1; + else end = middle; + } + slot = first; + } else { + const auto base_end = leaf.entries.begin() + + static_cast(leaf.sorted_size); + const auto base = std::lower_bound(leaf.entries.begin(), base_end, *retired, + [](const Entry& a, const Entry& b) { return less(a, b); }); + slot = static_cast(base - leaf.entries.begin()); + } + assert(slot < leaf.sorted_size && + equivalent(base_entry(leaf, slot), *retired)); leaf.invalidated[slot / 64] |= std::uint64_t{1} << (slot % 64); } - assert(leaf.entries.size() < leaf.entries.capacity()); - leaf.entries.insert(found, Entry{stored, key}); + if constexpr (ScoreRle) { + assert(leaf.entries.dirty.size() < leaf.entries.dirty.capacity()); + leaf.entries.dirty.insert(found, Entry{stored, key}); + } else { + assert(leaf.entries.size() < leaf.entries.capacity()); + leaf.entries.insert(found, Entry{stored, key}); + } ++dirty_count_; } [[nodiscard]] static std::size_t collect_additions( const Leaf& leaf, std::array& additions) { std::size_t addition_count = 0; - for (std::size_t i = leaf.sorted_size; i < leaf.entries.size(); ++i) { - if (!std::isnan(leaf.entries[i].score)) { - std::construct_at(&additions[addition_count++].entry, leaf.entries[i]); + for (auto it = tail_begin(leaf); it != tail_end(leaf); ++it) { + if (!std::isnan(it->score)) { + std::construct_at(&additions[addition_count++].entry, *it); } } std::sort(additions.begin(), additions.begin() + @@ -645,6 +709,13 @@ class PackedBPlusTree { std::array additions; const auto addition_count = collect_additions(leaf, additions); + struct RawCursor { + const std::vector& entries; + Entry entry(std::size_t slot) const { return entries[slot]; } + }; + using Cursor = std::conditional_t; + Cursor cursor{leaf.entries}; std::size_t base = 0; std::size_t addition = 0; const auto next_base = [&]() { @@ -655,12 +726,13 @@ class PackedBPlusTree { }; next_base(); while (base < leaf.sorted_size || addition < addition_count) { + const auto candidate = base < leaf.sorted_size ? cursor.entry(base) : Entry{}; const bool use_base = addition == addition_count || - (base < leaf.sorted_size && - less(leaf.entries[base], additions[addition].entry)); + (base < leaf.sorted_size && less(candidate, additions[addition].entry)); if (use_base) { - if (!fn(leaf.entries[base++])) return false; + ++base; + if (!fn(candidate)) return false; next_base(); } else { if (!fn(additions[addition++].entry)) return false; @@ -671,92 +743,177 @@ class PackedBPlusTree { void maybe_compact_leaf(std::uint32_t leaf_id) { const auto& leaf = leaves_[leaf_id]; - if (leaf.entries.size() - leaf.sorted_size >= merge_threshold_) { + if (tail_size(leaf) >= merge_threshold_) { compact_leaf(leaf_id); } } - void compact_leaf(std::uint32_t leaf_id) { - auto& leaf = leaves_[leaf_id]; - const auto tail_size = leaf.entries.size() - leaf.sorted_size; - if (tail_size == 0) return; + void assign_rle_base(Leaf& leaf, std::span values) noexcept { + sorted_count_ = sorted_count_ - leaf.sorted_size + values.size(); + dirty_count_ -= tail_size(leaf); + leaf.entries.assign(values); + leaf.sorted_size = values.size(); + leaf.invalidated.fill(0); + } - std::array additions; - const auto addition_count = collect_additions(leaf, additions); + void compact_rle_leaf(std::uint32_t leaf_id) { + auto& leaf = leaves_[leaf_id]; + if (tail_size(leaf) != 0) { + std::array values; + std::size_t count = 0; + for_each_leaf_entry(leaf_id, [&](const Entry& entry) { + assert(count < values.size()); + values[count++] = entry; + return true; + }); + const std::span live(values.data(), count); + leaf.entries.reserve_words(RleEntries::word_count(live)); + assign_rle_base(leaf, live); + assert(count == leaf.live_count); + } + leaf.entries.maybe_shrink(merge_threshold_); + } - const auto old_sorted = leaf.sorted_size; - std::size_t prefix = 0; - std::size_t source = 0; - const auto copy_run = [&](std::size_t end) { - const auto length = end - source; - if (length != 0 && prefix != source) { - std::memmove(leaf.entries.data() + prefix, leaf.entries.data() + source, - length * sizeof(Entry)); - } - prefix += length; - }; - // Visit only the holes. The unchanged prefix needs no writes, and each - // surviving run after it is compacted with one overlapping block move. - for (std::size_t word = 0; word < leaf.invalidated.size(); ++word) { - auto holes = leaf.invalidated[word]; - while (holes != 0) { - const auto hole = word * 64 + std::countr_zero(holes); - assert(hole < old_sorted); - copy_run(hole); - source = hole + 1; - holes &= holes - 1; - } + void rebalance_rle_leaves(std::uint32_t left_id, std::uint32_t right_id) { + auto& left = leaves_[left_id]; + auto& right = leaves_[right_id]; + const auto left_old = left.live_count; + const auto right_old = right.live_count; + const auto combined = left_old + right_old; + const auto left_new = combined > kLeafCapacity ? combined / 2 : combined; + const auto right_new = combined - left_new; + std::array values; + std::size_t count = 0; + for (const auto id : {left_id, right_id}) { + for_each_leaf_entry(id, [&](const Entry& entry) { + values[count++] = entry; + return true; + }); + } + assert(count == combined); + const std::span all(values.data(), count); + try { + left.entries.reserve_words(RleEntries::word_count(all.first(left_new))); + right.entries.reserve_words(RleEntries::word_count(all.subspan(left_new))); + } catch (const std::bad_alloc&) { + // This is optional maintenance after a committed deletion/rescore. + // Both old leaves, their dirty records, and all counts remain valid. + return; + } + const auto left_path = locate(left.fence); + const auto right_path = locate(right.fence); + assign_rle_base(left, all.first(left_new)); + assign_rle_base(right, all.subspan(left_new)); + left.live_count = left_new; + right.live_count = right_new; + if (right_new != 0) { + left.fence = values[left_new - 1]; + refresh_path(left_path, static_cast(left_new) - + static_cast(left_old)); + refresh_path(right_path, static_cast(right_new) - + static_cast(right_old)); + } else { + left.fence = right.fence; + left.next = right.next; + if (right.next != kNull) leaves_[right.next].prev = left_id; + if (last_leaf_ == right_id) last_leaf_ = left_id; + refresh_path(left_path, static_cast(right_old)); + remove_child(right_path, right_old); + release_leaf(right_id); + --active_leaf_count_; + } + left.entries.maybe_shrink(merge_threshold_); + right.entries.maybe_shrink(merge_threshold_); + if (left.live_count < kLeafCapacity / 4 && active_leaf_count_ > 1) { + maybe_rebalance_leaf(left_id); } - copy_run(old_sorted); - const auto final_size = prefix + addition_count; - leaf.entries.resize(final_size); - - if (addition_count != 0) { - if (prefix == 0 || less(leaf.entries[prefix - 1], additions[0].entry)) { - // An ordered append (including an empty base) never moves the base. - for (std::size_t i = 0; i < addition_count; ++i) { - leaf.entries[prefix + i] = additions[i].entry; + } + + void compact_leaf(std::uint32_t leaf_id) { + if constexpr (ScoreRle) { + compact_rle_leaf(leaf_id); + } else { + auto& leaf = leaves_[leaf_id]; + const auto tail_size = leaf.entries.size() - leaf.sorted_size; + if (tail_size == 0) return; + + std::array additions; + const auto addition_count = collect_additions(leaf, additions); + + const auto old_sorted = leaf.sorted_size; + std::size_t prefix = 0; + std::size_t source = 0; + const auto copy_run = [&](std::size_t end) { + const auto length = end - source; + if (length != 0 && prefix != source) { + std::memmove(leaf.entries.data() + prefix, leaf.entries.data() + source, + length * sizeof(Entry)); } - } else { - std::size_t left = prefix; - std::size_t right = addition_count; - std::size_t output = final_size; - while (right != 0) { - const auto& addition = additions[right - 1].entry; - if (left != 0 && less(addition, leaf.entries[left - 1])) { - // Find the run from its right edge. Exponential probes make short - // runs cheap, while long runs still use logarithmic comparisons. - auto start = left - 1; - std::size_t stride = 1; - while (start != 0) { - const auto probe = start > stride ? start - stride : 0; - if (!less(addition, leaf.entries[probe])) { - const auto first = std::upper_bound( - leaf.entries.begin() + static_cast(probe + 1), - leaf.entries.begin() + static_cast(start), - addition, - [](const Entry& a, const Entry& b) { return less(a, b); }); - start = static_cast(first - leaf.entries.begin()); - break; + prefix += length; + }; + // Visit only the holes. The unchanged prefix needs no writes, and each + // surviving run after it is compacted with one overlapping block move. + for (std::size_t word = 0; word < leaf.invalidated.size(); ++word) { + auto holes = leaf.invalidated[word]; + while (holes != 0) { + const auto hole = word * 64 + std::countr_zero(holes); + assert(hole < old_sorted); + copy_run(hole); + source = hole + 1; + holes &= holes - 1; + } + } + copy_run(old_sorted); + const auto final_size = prefix + addition_count; + leaf.entries.resize(final_size); + + if (addition_count != 0) { + if (prefix == 0 || less(leaf.entries[prefix - 1], additions[0].entry)) { + // An ordered append (including an empty base) never moves the base. + for (std::size_t i = 0; i < addition_count; ++i) { + leaf.entries[prefix + i] = additions[i].entry; + } + } else { + std::size_t left = prefix; + std::size_t right = addition_count; + std::size_t output = final_size; + while (right != 0) { + const auto& addition = additions[right - 1].entry; + if (left != 0 && less(addition, leaf.entries[left - 1])) { + // Find the run from its right edge. Exponential probes make short + // runs cheap, while long runs still use logarithmic comparisons. + auto start = left - 1; + std::size_t stride = 1; + while (start != 0) { + const auto probe = start > stride ? start - stride : 0; + if (!less(addition, leaf.entries[probe])) { + const auto first = std::upper_bound( + leaf.entries.begin() + static_cast(probe + 1), + leaf.entries.begin() + static_cast(start), + addition, + [](const Entry& a, const Entry& b) { return less(a, b); }); + start = static_cast(first - leaf.entries.begin()); + break; + } + start = probe; + stride *= 2; } - start = probe; - stride *= 2; + const auto length = left - start; + output -= length; + std::memmove(leaf.entries.data() + output, leaf.entries.data() + start, + length * sizeof(Entry)); + left = start; } - const auto length = left - start; - output -= length; - std::memmove(leaf.entries.data() + output, leaf.entries.data() + start, - length * sizeof(Entry)); - left = start; + leaf.entries[--output] = additions[--right].entry; } - leaf.entries[--output] = additions[--right].entry; } } + assert(final_size == leaf.live_count); + sorted_count_ = sorted_count_ - old_sorted + final_size; + dirty_count_ -= tail_size; + leaf.sorted_size = final_size; + leaf.invalidated.fill(0); } - assert(final_size == leaf.live_count); - sorted_count_ = sorted_count_ - old_sorted + final_size; - dirty_count_ -= tail_size; - leaf.sorted_size = final_size; - leaf.invalidated.fill(0); } void expand_fence(const Path& path, const Entry& entry) noexcept { @@ -791,8 +948,7 @@ class PackedBPlusTree { void split_leaf(const Path& path) { compact_leaf(path.leaf); - auto& current = leaves_[path.leaf]; - assert(current.live_count == kLeafCapacity); + assert(leaves_[path.leaf].live_count == kLeafCapacity); // All allocations precede changes to the live directory. reserve_memory_vector(branches_, branches_.size() + height_ + 1); @@ -800,17 +956,34 @@ class PackedBPlusTree { auto& leaf = leaves_[path.leaf]; auto& right = leaves_[right_id]; const auto split = leaf.sorted_size / 2; - right.entries.insert( - right.entries.end(), - leaf.entries.begin() + static_cast(split), - leaf.entries.end()); - leaf.entries.resize(split); - right.sorted_size = right.entries.size(); + const auto old_fence = leaf.fence; + if constexpr (ScoreRle) { + std::array values; + for (std::size_t i = 0; i < leaf.sorted_size; ++i) { + values[i] = base_entry(leaf, i); + } + const std::span all(values.data(), leaf.sorted_size); + try { + leaf.entries.reserve_words(RleEntries::word_count(all.first(split))); + right.entries.reserve_words(RleEntries::word_count(all.subspan(split))); + } catch (...) { + release_leaf(right_id); + throw; + } + assign_rle_base(right, all.subspan(split)); + assign_rle_base(leaf, all.first(split)); + } else { + right.entries.insert( + right.entries.end(), + leaf.entries.begin() + static_cast(split), + leaf.entries.end()); + leaf.entries.resize(split); + right.sorted_size = right.entries.size(); + leaf.sorted_size = split; + } right.live_count = right.sorted_size; - leaf.sorted_size = split; leaf.live_count = split; - const auto old_fence = leaf.fence; - leaf.fence = leaf.entries.back(); + leaf.fence = base_entry(leaf, split - 1); right.fence = old_fence; right.has_fence = true; right.prev = path.leaf; @@ -899,60 +1072,64 @@ class PackedBPlusTree { right_id = leaf_id; } if (left_id == kNull || right_id == kNull) return; - compact_leaf(left_id); - compact_leaf(right_id); - - auto left_path = locate(leaves_[left_id].fence); - auto right_path = locate(leaves_[right_id].fence); - assert(left_path.leaf == left_id && right_path.leaf == right_id); - auto& left = leaves_[left_id]; - auto& right = leaves_[right_id]; - const auto left_old = left.live_count; - const auto right_old = right.live_count; - const auto combined = left_old + right_old; - - if (combined > kLeafCapacity) { - const auto left_new = combined / 2; - const auto right_new = combined - left_new; - // Compacted neighbors are already globally ordered. Transfer only the - // boundary run; the destination's reserved leaf capacity is sufficient. - assert(left_new <= left.entries.capacity()); - assert(right_new <= right.entries.capacity()); - if (left_old < left_new) { - const auto boundary = right.entries.begin() + - static_cast(left_new - left_old); - left.entries.insert(left.entries.end(), right.entries.begin(), boundary); - right.entries.erase(right.entries.begin(), boundary); - } else if (left_old > left_new) { - const auto boundary = left.entries.begin() + - static_cast(left_new); - right.entries.insert(right.entries.begin(), boundary, left.entries.end()); - left.entries.resize(left_new); + if constexpr (ScoreRle) { + rebalance_rle_leaves(left_id, right_id); + } else { + compact_leaf(left_id); + compact_leaf(right_id); + + auto left_path = locate(leaves_[left_id].fence); + auto right_path = locate(leaves_[right_id].fence); + assert(left_path.leaf == left_id && right_path.leaf == right_id); + auto& left = leaves_[left_id]; + auto& right = leaves_[right_id]; + const auto left_old = left.live_count; + const auto right_old = right.live_count; + const auto combined = left_old + right_old; + + if (combined > kLeafCapacity) { + const auto left_new = combined / 2; + const auto right_new = combined - left_new; + // Compacted neighbors are already globally ordered. Transfer only the + // boundary run; the destination's reserved leaf capacity is sufficient. + assert(left_new <= left.entries.capacity()); + assert(right_new <= right.entries.capacity()); + if (left_old < left_new) { + const auto boundary = right.entries.begin() + + static_cast(left_new - left_old); + left.entries.insert(left.entries.end(), right.entries.begin(), boundary); + right.entries.erase(right.entries.begin(), boundary); + } else if (left_old > left_new) { + const auto boundary = left.entries.begin() + + static_cast(left_new); + right.entries.insert(right.entries.begin(), boundary, left.entries.end()); + left.entries.resize(left_new); + } + left.sorted_size = left.live_count = left_new; + right.sorted_size = right.live_count = right_new; + left.fence = left.entries.back(); + refresh_path(left_path, static_cast(left_new) - + static_cast(left_old)); + refresh_path(right_path, static_cast(right_new) - + static_cast(right_old)); + return; } - left.sorted_size = left.live_count = left_new; - right.sorted_size = right.live_count = right_new; - left.fence = left.entries.back(); - refresh_path(left_path, static_cast(left_new) - - static_cast(left_old)); - refresh_path(right_path, static_cast(right_new) - - static_cast(right_old)); - return; - } - - left.entries.insert(left.entries.end(), right.entries.begin(), - right.entries.end()); - left.sorted_size = left.live_count = combined; - left.fence = right.fence; - left.next = right.next; - if (right.next != kNull) leaves_[right.next].prev = left_id; - if (last_leaf_ == right_id) last_leaf_ = left_id; - refresh_path(left_path, static_cast(right_old)); - remove_child(right_path, right_old); - release_leaf(right_id); - --active_leaf_count_; - if (left.live_count < kLeafCapacity / 4 && active_leaf_count_ > 1) { - maybe_rebalance_leaf(left_id); + left.entries.insert(left.entries.end(), right.entries.begin(), + right.entries.end()); + left.sorted_size = left.live_count = combined; + left.fence = right.fence; + left.next = right.next; + if (right.next != kNull) leaves_[right.next].prev = left_id; + if (last_leaf_ == right_id) last_leaf_ = left_id; + refresh_path(left_path, static_cast(right_old)); + remove_child(right_path, right_old); + release_leaf(right_id); + --active_leaf_count_; + + if (left.live_count < kLeafCapacity / 4 && active_leaf_count_ > 1) { + maybe_rebalance_leaf(left_id); + } } } diff --git a/include/goblin/core/packed_zset.hpp b/include/goblin/core/packed_zset.hpp index 503342d..648dd04 100644 --- a/include/goblin/core/packed_zset.hpp +++ b/include/goblin/core/packed_zset.hpp @@ -350,11 +350,11 @@ template return (bits & kSign) != 0 ? ~bits : bits ^ kSign; } -template +template class PackedZSetIndex { public: using Key = typename Traits::Key; - using Order = PackedBPlusTree; + using Order = PackedBPlusTree; using Entry = typename Order::Entry; explicit PackedZSetIndex( @@ -392,6 +392,15 @@ class PackedZSetIndex { [[nodiscard]] std::size_t tree_height() const noexcept { return order_.tree_height(); } + [[nodiscard]] static constexpr bool score_rle_enabled() noexcept { + return ScoreRle; + } + [[nodiscard]] std::size_t compressed_leaf_count() const noexcept { + return order_.compressed_leaf_count(); + } + [[nodiscard]] std::size_t sorted_score_bytes() const noexcept { + return order_.sorted_score_bytes(); + } [[nodiscard]] double merge_exponent() const noexcept { return merge_exponent_; } @@ -625,6 +634,17 @@ using PackedInt64Float64 = using PackedUuidFloat32 = PackedZSetIndex; using PackedUuidFloat64 = PackedZSetIndex; +using PackedRleInt32Float32 = + PackedZSetIndex, float, true>; +using PackedRleInt32Float64 = + PackedZSetIndex, double, true>; +using PackedRleInt64Float32 = + PackedZSetIndex, float, true>; +using PackedRleInt64Float64 = + PackedZSetIndex, double, true>; +using PackedRleUuidFloat32 = PackedZSetIndex; +using PackedRleUuidFloat64 = PackedZSetIndex; + } // namespace detail // Runtime wrapper around the six compile-time layouts. Keyspace stores this as @@ -634,7 +654,8 @@ class PackedZSet { public: explicit PackedZSet( PackedZSetKind kind, - double merge_exponent = kDefaultPackedZSetMergeExponent) + double merge_exponent = kDefaultPackedZSetMergeExponent, + bool score_rle = false) : kind_(kind) { if (!valid_packed_zset_merge_exponent(merge_exponent)) { throw std::invalid_argument( @@ -642,22 +663,28 @@ class PackedZSet { } switch (kind) { case PackedZSetKind::Int32Float32: - impl_.emplace(merge_exponent); + if (score_rle) impl_.emplace(merge_exponent); + else impl_.emplace(merge_exponent); break; case PackedZSetKind::Int32Float64: - impl_.emplace(merge_exponent); + if (score_rle) impl_.emplace(merge_exponent); + else impl_.emplace(merge_exponent); break; case PackedZSetKind::Int64Float32: - impl_.emplace(merge_exponent); + if (score_rle) impl_.emplace(merge_exponent); + else impl_.emplace(merge_exponent); break; case PackedZSetKind::Int64Float64: - impl_.emplace(merge_exponent); + if (score_rle) impl_.emplace(merge_exponent); + else impl_.emplace(merge_exponent); break; case PackedZSetKind::UuidFloat32: - impl_.emplace(merge_exponent); + if (score_rle) impl_.emplace(merge_exponent); + else impl_.emplace(merge_exponent); break; case PackedZSetKind::UuidFloat64: - impl_.emplace(merge_exponent); + if (score_rle) impl_.emplace(merge_exponent); + else impl_.emplace(merge_exponent); break; } } @@ -704,6 +731,18 @@ class PackedZSet { return std::visit([](const auto& value) { return value.tree_height(); }, impl_); } + [[nodiscard]] bool score_rle_enabled() const noexcept { + return std::visit([](const auto& value) { return value.score_rle_enabled(); }, + impl_); + } + [[nodiscard]] std::size_t compressed_leaf_count() const noexcept { + return std::visit([](const auto& value) { return value.compressed_leaf_count(); }, + impl_); + } + [[nodiscard]] std::size_t sorted_score_bytes() const noexcept { + return std::visit([](const auto& value) { return value.sorted_score_bytes(); }, + impl_); + } [[nodiscard]] double merge_exponent() const noexcept { return std::visit([](const auto& value) { return value.merge_exponent(); }, impl_); @@ -863,7 +902,10 @@ class PackedZSet { using Implementation = std::variant; + detail::PackedUuidFloat32, detail::PackedUuidFloat64, + detail::PackedRleInt32Float32, detail::PackedRleInt32Float64, + detail::PackedRleInt64Float32, detail::PackedRleInt64Float64, + detail::PackedRleUuidFloat32, detail::PackedRleUuidFloat64>; template [[nodiscard]] static std::optional parse_for_index( diff --git a/include/goblin/core/store.hpp b/include/goblin/core/store.hpp index eb928ca..da429c5 100644 --- a/include/goblin/core/store.hpp +++ b/include/goblin/core/store.hpp @@ -163,6 +163,9 @@ struct PackedZSetMemoryStats { std::size_t tree_height{0}; double merge_exponent{kDefaultPackedZSetMergeExponent}; std::size_t merge_threshold{0}; + bool score_rle{false}; + std::size_t compressed_leaf_count{0}; + std::size_t sorted_score_bytes{0}; std::size_t total_allocated_bytes{0}; }; @@ -1556,6 +1559,8 @@ struct StoreOptions { // Zero merges every mutation; one permits a tail as large as the leaf's sorted // capacity. Reads reconcile only visited leaves and never trigger maintenance. double packed_zset_merge_exponent{kDefaultPackedZSetMergeExponent}; + // Encode equal sorted-base scores in runs; snapshots use the receiver policy. + bool packed_zset_score_rle{false}; // Unqualified zset commands create this representation. A live key remains // pinned to the representation that created or restored it. ZSetImplementation zset_implementation{ZSetImplementation::Standard}; diff --git a/src/command.cpp b/src/command.cpp index bf148e3..11e7c35 100644 --- a/src/command.cpp +++ b/src/command.cpp @@ -568,6 +568,9 @@ void append_hello_response(std::string& out, resp::Version version, fields.emplace_back("merge_exponent"); fields.push_back(format_score(stats.merge_exponent)); add("merge_threshold", stats.merge_threshold); + add("score_rle", stats.score_rle); + add("compressed_leaf_count", stats.compressed_leaf_count); + add("sorted_score_bytes", stats.sorted_score_bytes); add("total_allocated_bytes", stats.total_allocated_bytes); return fields; } diff --git a/src/main.cpp b/src/main.cpp index 75b927a..d8e07ca 100644 --- a/src/main.cpp +++ b/src/main.cpp @@ -643,6 +643,7 @@ void print_usage(std::string_view program) { " packed-int64-float64|packed-uuid-float32|\n" " packed-uuid-float64] (default: standard)\n" << " [--packed-zset-merge-exponent K] (0 <= K <= 1)\n" + << " [--packed-zset-score-rle|--no-packed-zset-score-rle]\n" << " [--block-shrink on|off]\n" << " [--zset-chunk-bytes BYTES] [--hash-chunk-bytes BYTES]\n" << " [--hash-compaction-knapsack|--no-hash-compaction-knapsack]\n" @@ -1553,6 +1554,12 @@ int main(int argc, char** argv) { continue; } + if (arg == "--packed-zset-score-rle" || + arg == "--no-packed-zset-score-rle") { + store_options.packed_zset_score_rle = arg == "--packed-zset-score-rle"; + continue; + } + if (arg == "--packed-zset-merge-exponent") { if (i + 1 >= argc) { print_usage(argv[0]); diff --git a/src/store.cpp b/src/store.cpp index 3c17140..61cded3 100644 --- a/src/store.cpp +++ b/src/store.cpp @@ -1300,7 +1300,8 @@ PackedZSetAddResult Store::packed_zadd( // Mutate a detached object first. Invalid members/scores and XX-only misses // therefore never leave an empty key behind. - PackedZSet prepared(kind, options_.packed_zset_merge_exponent); + PackedZSet prepared(kind, options_.packed_zset_merge_exponent, + options_.packed_zset_score_rle); auto result = prepared.add(items, options); if (!result.invalid_member && !result.invalid_score && !prepared.empty()) { (void)keyspace_.place_loaded_packed_zset(key, std::move(prepared)); @@ -1441,7 +1442,8 @@ ZStoreResult Store::packed_zunionstore( } if (invalid_score) return {.invalid_score = true}; - PackedZSet merged(kind, options_.packed_zset_merge_exponent); + PackedZSet merged(kind, options_.packed_zset_merge_exponent, + options_.packed_zset_score_rle); std::vector items; reserve_memory_vector(items, totals.size()); totals.for_each([&items](const auto& entry) { @@ -1486,7 +1488,8 @@ ZStoreResult Store::packed_zinterstore( } } - PackedZSet intersection(kind, options_.packed_zset_merge_exponent); + PackedZSet intersection(kind, options_.packed_zset_merge_exponent, + options_.packed_zset_score_rle); bool invalid_score = false; bool invalid_member = false; if (!sources.empty()) { @@ -1570,6 +1573,9 @@ std::optional Store::packed_zset_memory_stats( .tree_height = zset->tree_height(), .merge_exponent = zset->merge_exponent(), .merge_threshold = zset->merge_threshold(), + .score_rle = zset->score_rle_enabled(), + .compressed_leaf_count = zset->compressed_leaf_count(), + .sorted_score_bytes = zset->sorted_score_bytes(), .total_allocated_bytes = zset->allocated_bytes(), }; } @@ -3151,7 +3157,8 @@ CopyResult Store::copy(std::string_view source, std::string_view destination, } case KeyType::PackedZset: { const auto* original = find_packed_zset(source); - PackedZSet clone(original->kind(), options_.packed_zset_merge_exponent); + PackedZSet clone(original->kind(), options_.packed_zset_merge_exponent, + options_.packed_zset_score_rle); const auto entries = original->range_by_rank(0, -1); std::vector items; reserve_memory_vector(items, entries.size()); @@ -4224,7 +4231,8 @@ SnapshotLoadStats Store::load_native(std::istream& in) { items.push_back({score, member}); } PackedZSet zset(static_cast(encoded_kind), - options_.packed_zset_merge_exponent); + options_.packed_zset_merge_exponent, + options_.packed_zset_score_rle); const auto result = zset.add(items); if (result.invalid_member || result.invalid_score || zset.size() != count) { diff --git a/tests/packed_bplus_tree_test.cpp b/tests/packed_bplus_tree_test.cpp index d294cec..6d515e0 100644 --- a/tests/packed_bplus_tree_test.cpp +++ b/tests/packed_bplus_tree_test.cpp @@ -25,9 +25,9 @@ typename Traits::Key key_for(std::size_t id) { } } -template +template void redistribute_full_neighbor() { - using Tree = goblin::core::detail::PackedBPlusTree; + using Tree = goblin::core::detail::PackedBPlusTree; const auto capacity = Tree::kLeafCapacity; const auto count = capacity + capacity / 2; for (const bool reverse : {false, true}) { @@ -44,9 +44,10 @@ void redistribute_full_neighbor() { tree.erase({static_cast(id), key_for(id)}); } // An underfull leaf plus a full neighbor cannot coalesce: both directions - // must redistribute, preserving the mapping without another allocation. + // must redistribute. The raw layout does so without another allocation; + // RLE may grow a score stream when the neighbor has more diverse scores. assert(tree.leaf_count() == 2 && tree.check_invariants()); - assert(tree.allocated_bytes() == allocated); + if constexpr (!Rle) assert(tree.allocated_bytes() == allocated); const auto entries = tree.range_by_rank(0, -1, false); assert(entries.size() == count - capacity / 4 - 1); for (std::size_t i = 0; i < entries.size(); ++i) { @@ -57,9 +58,9 @@ void redistribute_full_neighbor() { } } -template +template void exercise(double exponent) { - using Tree = goblin::core::detail::PackedBPlusTree; + using Tree = goblin::core::detail::PackedBPlusTree; using Entry = typename Tree::Entry; Tree tree(exponent); const auto capacity = Tree::kLeafCapacity; @@ -185,9 +186,9 @@ void exercise(double exponent) { verify(); } -template +template void deep_rescore_routing() { - using Tree = goblin::core::detail::PackedBPlusTree; + using Tree = goblin::core::detail::PackedBPlusTree; using Entry = typename Tree::Entry; const auto capacity = Tree::kLeafCapacity; const auto count = capacity * 40 + 11; @@ -259,28 +260,156 @@ void deep_rescore_routing() { verify(); } +template +void score_run_boundaries() { + using Traits = goblin::core::detail::PackedIntegerTraits; + using Tree = goblin::core::detail::PackedBPlusTree; + using Raw = goblin::core::detail::PackedBPlusTree; + Tree tree(0.0); + for (int i = 0; i < 3; ++i) tree.insert({Score{7}, i}); + assert(tree.sorted_score_bytes() == 3 * sizeof(Score)); + assert(tree.compressed_leaf_count() == 0); + tree.insert({Score{7}, 3}); + assert(tree.sorted_score_bytes() == 3 * sizeof(Score)); + assert(tree.compressed_leaf_count() == 1); + tree.erase({Score{7}, 1}); + assert(tree.sorted_score_bytes() == 3 * sizeof(Score)); + assert(tree.compressed_leaf_count() == 0 && tree.check_invariants()); + + Tree tied(0.5); + Raw raw(0.5); + const auto count = Tree::kLeafCapacity * 3; + for (std::size_t i = 0; i < count; ++i) { + tied.insert({Score{1}, static_cast(i)}); + raw.insert({Score{1}, static_cast(i)}); + } + tied.force_merge(); + raw.force_merge(); + assert(tied.check_invariants()); + assert(tied.compressed_leaf_count() == tied.leaf_count()); + assert(tied.sorted_score_bytes() == tied.leaf_count() * 3 * sizeof(Score)); + assert(tied.allocated_bytes() < raw.allocated_bytes()); + // Mutate a copy, including dirty slots and score-stream growth. + auto copy = tied; + for (std::size_t i = 0; i < count; ++i) { + copy.replace({Score{1}, static_cast(i)}, + {static_cast(i + 2), static_cast(i)}); + } + copy.force_merge(); + assert(copy.check_invariants() && tied.check_invariants()); + assert(copy.compressed_leaf_count() == 0); + assert(copy.sorted_score_bytes() == count * sizeof(Score)); +} + +template +void score_run_allocation_failures() { + using namespace goblin::core; + using Traits = detail::PackedIntegerTraits; + using Tree = detail::PackedBPlusTree; + const auto capacity = Tree::kLeafCapacity; + MemoryCeiling deny_growth(1); + deny_growth.bind(nullptr, [](const void*) noexcept { return std::size_t{1}; }); + + // Break long runs into distinct scores until compaction needs more storage. + // Exercise both a single leaf and moves between separate source/dest leaves. + for (const bool cross_leaf : {false, true}) { + Tree tree(0.5); + const auto count = cross_leaf ? capacity * 2 + capacity / 4 : capacity / 2; + for (std::size_t i = 0; i < count; ++i) { + tree.insert({Score{1}, static_cast(i)}); + } + tree.force_merge(); + std::size_t changed = 0; + bool rejected = false; + for (; changed < capacity / 2; ++changed) { + MemoryCeilingScope scope(&deny_growth); + try { + tree.replace({Score{1}, static_cast(changed)}, + {static_cast(changed + 2), + static_cast(changed)}); + } catch (const MaxMemoryExceeded&) { + rejected = true; + break; + } + } + assert(rejected && changed != 0 && tree.check_invariants()); + assert(tree.size() == count); + for (const auto& entry : tree.range_by_rank(0, -1, false)) { + const auto id = static_cast(entry.key); + assert(entry.score == (id < changed ? static_cast(id + 2) : Score{1})); + } + // A refused write remains retryable with the original old tuple. + tree.replace({Score{1}, static_cast(changed)}, + {static_cast(changed + 2), static_cast(changed)}); + tree.force_merge(); + assert(tree.check_invariants()); + } + + // Redistribution from a diverse neighbor would expand a compressed leaf. + // Denying that optional growth must not reject the already committed erase. + Tree tree(0.5); + const auto count = capacity + capacity / 2; + for (std::size_t i = 0; i < count; ++i) { + tree.insert({i < capacity / 2 ? Score{1} : static_cast(i + 2), + static_cast(i)}); + } + tree.force_merge(); + assert(tree.leaf_count() == 2); + { + MemoryCeilingScope scope(&deny_growth); + for (std::size_t i = 0; i <= capacity / 4; ++i) { + tree.erase({Score{1}, static_cast(i)}); + } + } + assert(tree.size() == count - capacity / 4 - 1 && tree.check_invariants()); + tree.erase({Score{1}, static_cast(capacity / 4 + 1)}); + assert(tree.check_invariants()); +} + } // namespace int main() { using namespace goblin::core::detail; + score_run_boundaries(); + score_run_boundaries(); + score_run_allocation_failures(); + score_run_allocation_failures(); redistribute_full_neighbor, float>(); + redistribute_full_neighbor, float, true>(); redistribute_full_neighbor, double>(); + redistribute_full_neighbor, double, true>(); redistribute_full_neighbor, float>(); + redistribute_full_neighbor, float, true>(); redistribute_full_neighbor, double>(); + redistribute_full_neighbor, double, true>(); redistribute_full_neighbor(); + redistribute_full_neighbor(); redistribute_full_neighbor(); + redistribute_full_neighbor(); deep_rescore_routing, float>(); + deep_rescore_routing, float, true>(); deep_rescore_routing, double>(); + deep_rescore_routing, double, true>(); deep_rescore_routing, float>(); + deep_rescore_routing, float, true>(); deep_rescore_routing, double>(); + deep_rescore_routing, double, true>(); deep_rescore_routing(); + deep_rescore_routing(); deep_rescore_routing(); + deep_rescore_routing(); for (const auto exponent : {0.0, 0.5, 1.0}) { exercise, float>(exponent); + exercise, float, true>(exponent); exercise, double>(exponent); + exercise, double, true>(exponent); exercise, float>(exponent); + exercise, float, true>(exponent); exercise, double>(exponent); + exercise, double, true>(exponent); exercise(exponent); + exercise(exponent); exercise(exponent); + exercise(exponent); } } diff --git a/tests/packed_zset_test.cpp b/tests/packed_zset_test.cpp index 7a646f9..c8e23ea 100644 --- a/tests/packed_zset_test.cpp +++ b/tests/packed_zset_test.cpp @@ -64,11 +64,11 @@ void test_integer_hash_distribution() { } } -template +template void test_update_slots_and_allocation_failures() { using namespace goblin::core; using Key = typename Traits::Key; - using Index = detail::PackedZSetIndex; + using Index = detail::PackedZSetIndex; const auto key_for = [](std::size_t id) { if constexpr (std::is_integral_v) { return static_cast(id); @@ -188,7 +188,7 @@ void test_leaf_dirty_dedup_and_in_place_merge() { } } -void test_randomized_append_log_against_reference() { +void test_randomized_append_log_against_reference(bool rle = false) { constexpr std::array cases{ std::pair{PackedZSetKind::Int32Float32, 0.0}, std::pair{PackedZSetKind::Int32Float32, 0.5}, @@ -198,7 +198,7 @@ void test_randomized_append_log_against_reference() { std::pair{PackedZSetKind::Int32Float64, 1.0}, }; for (const auto [kind, merge_exponent] : cases) { - PackedZSet zset(kind, merge_exponent); + PackedZSet zset(kind, merge_exponent, rle); std::map reference; std::uint64_t random = 0x8c3c'010c'cb47'563dULL; const auto next = [&random] { @@ -397,6 +397,58 @@ void test_configurable_merge_exponent() { assert(loaded->unsorted_entries == 0); } +void test_score_rle_store_policy() { + StoreOptions options; + options.packed_zset_score_rle = true; + options.zset_implementation = goblin::core::ZSetImplementation::PackedInt32Float32; + Store store(options); + std::vector members; + std::vector items; + members.reserve(1200); + items.reserve(1200); + for (int i = 0; i < 1200; ++i) { + members.push_back(std::to_string(i)); + items.push_back({static_cast(i / 100), members.back()}); + } + assert(store.packed_zadd("runs", PackedZSetKind::Int32Float32, items).added == 1200); + assert(run(store, {"COPY", "runs", "copy"}) == ":1\r\n"); + assert(run(store, {"ZUNIONSTORE", "union", "2", "runs", "copy"}) == ":1200\r\n"); + assert(run(store, {"ZINTERSTORE", "intersection", "2", "runs", "copy"}) == ":1200\r\n"); + assert(run(store, {"ZADD", "ordinary", "1", "1", "1", "2", "1", "3", "1", "4"}) == ":4\r\n"); + for (const auto key : {"runs", "copy", "union", "intersection", "ordinary"}) { + assert(run(store, {"GOBLIN.OPTIMIZE", key}).front() == ':'); + const auto stats = store.packed_zset_memory_stats(key); + assert(stats && stats->score_rle && stats->compressed_leaf_count != 0); + assert(stats->sorted_score_bytes < stats->sorted_entries * sizeof(float)); + const auto reply = run(store, {"GOBLIN.MEMORY", key}); + assert(reply.find("score_rle") != std::string::npos); + assert(reply.find("compressed_leaf_count") != std::string::npos); + assert(reply.find("sorted_score_bytes") != std::string::npos); + } + const auto expected = run(store, {"ZRANGE", "runs", "0", "-1", "WITHSCORES"}); + assert(run(store, {"ZRANGE", "copy", "0", "-1", "WITHSCORES"}) == expected); + std::stringstream snapshot; + store.save(snapshot, false); + for (const bool enabled : {false, true}) { + StoreOptions receiver; + receiver.packed_zset_score_rle = enabled; + Store restored(receiver); + snapshot.clear(); + snapshot.seekg(0); + assert(restored.load(snapshot).keys == 5); + const auto stats = restored.packed_zset_memory_stats("runs"); + assert(stats && stats->score_rle == enabled); + assert(run(restored, {"ZRANGE", "runs", "0", "-1", "WITHSCORES"}) == expected); + // Also cover a snapshot emitted by the ordinary layout loading into RLE. + std::stringstream again; + restored.save(again, false); + Store compressed(options); + assert(compressed.load(again).keys == 5); + assert(compressed.packed_zset_memory_stats("runs")->score_rle); + assert(run(compressed, {"ZRANGE", "runs", "0", "-1", "WITHSCORES"}) == expected); + } +} + void test_all_command_prefixes() { const std::vector prefixes{ "GOBLIN.PACKED_INT32_FLOAT32", @@ -715,9 +767,11 @@ void test_snapshot_copy_and_optimize() { "$19\r\n9223372036854775807\r\n$3\r\n3.5\r\n"); } -void test_replication_round_trip() { +void test_replication_round_trip(bool rle = false) { Store source; - Store target; + StoreOptions receiver_options; + receiver_options.packed_zset_score_rle = rle; + Store target(receiver_options); std::vector fields{ "GOBLIN.PACKED_INT32_FLOAT64.ZADD", "replicated", "2.5", "17"}; auto parsed = goblin::core::parse_command(fields); @@ -740,6 +794,7 @@ void test_replication_round_trip() { assert(goblin::core::apply_firehose_batch(target, batch, error)); assert(target.packed_zset_kind("replicated") == PackedZSetKind::Int32Float64); + assert(target.packed_zset_memory_stats("replicated")->score_rle == rle); assert(run(target, {"GOBLIN.PACKED_INT32_FLOAT64.ZSCORE", "replicated", "17"}) == "$3\r\n2.5\r\n"); @@ -748,7 +803,7 @@ void test_replication_round_trip() { default_options.zset_implementation = ZSetImplementation::PackedInt64Float32; Store default_source(default_options); - Store default_target; + Store default_target(receiver_options); std::vector ordinary_fields{ "ZADD", "ordinary-replicated", "0.1", "91"}; auto ordinary = goblin::core::parse_command(ordinary_fields); @@ -770,6 +825,7 @@ void test_replication_round_trip() { error)); assert(default_target.packed_zset_kind("ordinary-replicated") == PackedZSetKind::Int64Float32); + assert(default_target.packed_zset_memory_stats("ordinary-replicated")->score_rle == rle); assert(run(default_target, {"ZSCORE", "ordinary-replicated", "91"}) == "$19\r\n0.10000000149011612\r\n"); @@ -782,15 +838,23 @@ int main() { test_integer_hash_distribution(); using namespace goblin::core::detail; test_update_slots_and_allocation_failures, float>(); + test_update_slots_and_allocation_failures, float, true>(); test_update_slots_and_allocation_failures, double>(); + test_update_slots_and_allocation_failures, double, true>(); test_update_slots_and_allocation_failures, float>(); + test_update_slots_and_allocation_failures, float, true>(); test_update_slots_and_allocation_failures, double>(); + test_update_slots_and_allocation_failures, double, true>(); test_update_slots_and_allocation_failures(); + test_update_slots_and_allocation_failures(); test_update_slots_and_allocation_failures(); + test_update_slots_and_allocation_failures(); test_leaf_dirty_dedup_and_in_place_merge(); test_randomized_append_log_against_reference(); + test_randomized_append_log_against_reference(true); test_multilevel_tree_and_cross_leaf_churn(); test_configurable_merge_exponent(); + test_score_rle_store_policy(); test_all_command_prefixes(); test_uuid_commands_and_representation_gate(); test_float32_scores_and_full_surface(); @@ -798,4 +862,5 @@ int main() { test_default_implementation_selector(); test_snapshot_copy_and_optimize(); test_replication_round_trip(); + test_replication_round_trip(true); }