Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 19 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,25 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

## [Unreleased]

### Performance
- **SCAN hot-plane pages are now true O(COUNT) — no full-table walk per
page (#368).** The SCAN cursor hash is now the DashTable's own
fixed-seed key hash truncated to its top 48 bits, which makes cursor
ranges line up exactly with the extendible-hashing directory (indexed
by top hash bits): segments are range-partitioned in hash space, so a
page visits only the segments covering hashes at or after the cursor
and stops as soon as COUNT entries are collected
(`DashTable::hash_page` → `Database::scan_hot_page`). Per-page hot
cost is now independent of keyspace size (previously every page walked
and hashed all n live entries). Splits, merges, and directory doubling
between pages remain safe by construction — the cursor is a position
in hash space, and structural churn only changes which segment covers
that position, never the set of keys at or above it. The cold plane
keeps its filtered in-RAM index walk (bounded by spilled-key count;
ordered cold-side paging remains a follow-up in #368). SCAN cursors
from before this change are invalidated (cursors are documented as
ephemeral; restart scans at 0).

### Fixed
- **SCAN now honors the Redis stable-key guarantee under churn (#368).**
The cursor was a positional index into a keyspace snapshot re-collected
Expand Down
89 changes: 64 additions & 25 deletions src/command/key.rs
Original file line number Diff line number Diff line change
Expand Up @@ -970,6 +970,12 @@ pub fn scan(db: &mut Database, args: &[Frame]) -> Frame {
Some(c) => c,
None => return Frame::Error(Bytes::from_static(b"ERR invalid cursor")),
};
// Clamp to the 48-bit hash space. Legitimate resumed cursors always
// fit (multi-shard composites are unpacked by `coordinate_scan` before
// reaching here); an out-of-range client cursor would otherwise filter
// out every key (h48 < 2^48) and falsely report "scan complete" on a
// non-empty keyspace.
let cursor = cursor & 0x0000_FFFF_FFFF_FFFF;

// Parse optional arguments
let mut match_pattern: Option<&[u8]> = None;
Expand Down Expand Up @@ -1012,21 +1018,21 @@ pub fn scan(db: &mut Database, args: &[Frame]) -> Frame {
scan_core(db, cursor, count, match_pattern, type_filter, now_ms)
}

/// Stable 48-bit key hash for SCAN cursors (FNV-1a 64 truncated).
/// Stable 48-bit key hash for SCAN cursors.
///
/// The cursor is a POSITION IN HASH SPACE, so this hash must be stable
/// across pages and across the two planes — a fixed-seed hash, never the
/// tables' randomized hashers. 48 bits because the multi-shard composite
/// cursor reserves the upper 16 bits for the shard index
/// (`coordinate_scan`).
/// across pages and across the two planes — a fixed-seed hash. This is the
/// hot table's own fixed-seed key hash truncated to its TOP 48 bits, which
/// makes the cursor range-compatible with the DashTable's extendible-
/// hashing directory (indexed by top hash bits): `Database::scan_hot_page`
/// walks only the segments covering `[cursor << 16, ..)` — O(COUNT) per
/// page instead of a full-table walk (#368). The cold plane computes the
/// same value from the key bytes, keeping both planes in one hash order.
/// 48 bits because the multi-shard composite cursor reserves the upper
/// 16 bits for the shard index (`coordinate_scan`).
#[inline]
fn scan_hash48(key: &[u8]) -> u64 {
let mut h: u64 = 0xcbf2_9ce4_8422_2325;
for &b in key {
h ^= b as u64;
h = h.wrapping_mul(0x0000_0100_0000_01b3);
}
h & 0x0000_FFFF_FFFF_FFFF
crate::storage::dashtable::hash_key(key) >> 16
}

/// Shared SCAN page walk (#368): hash-ordered iteration with a real
Expand All @@ -1037,11 +1043,13 @@ fn scan_hash48(key: &[u8]) -> u64 {
/// unvisited hash. A key's hash never changes, so inserts/deletes/spills/
/// promotions between pages cannot displace another key's position —
/// Redis's contract ("a key present for the entire scan is returned at
/// least once"; here exactly once) holds under churn. Per page this does
/// ONE walk over both planes with a bounded `count`-min selection heap —
/// no full-keyspace sort and no second lookup pass (the old design paid
/// `collect + sort + exists() + get_if_alive()` over every key on every
/// page).
/// least once"; here exactly once) holds under churn. Per page the hot
/// plane is a true O(COUNT) segment-range walk (`Database::scan_hot_page`
/// visits only DashTable segments covering hashes ≥ cursor, #368); the
/// cold plane is still a filtered in-RAM index walk. Both feed one bounded
/// `count`-min selection heap — no full-keyspace sort and no second
/// lookup pass (the old design paid `collect + sort + exists() +
/// get_if_alive()` over every key on every page).
///
/// Page-boundary rule: a full page never advances the cursor past a hash
/// value whose key group might be only partially selected, so hash
Expand All @@ -1059,8 +1067,15 @@ fn scan_core(
) -> Frame {
use std::collections::BinaryHeap;

// Hot plane: O(COUNT) hash-range page — only the DashTable segments
// covering hashes ≥ cursor are visited; entries arrive live-filtered
// with their hash48 precomputed. When `more`, the page holds ≥ count
// entries, all below anything in unvisited segments, so it always
// contains the count smallest hot candidates.
let (hot_page, _hot_more) = db.scan_hot_page(cursor, count, now_ms);

// Bounded max-heap selection: the `count` smallest (hash, key) pairs
// at or after the cursor, in one pass over both planes.
// at or after the cursor, across both planes.
let mut heap: BinaryHeap<(u64, CompactKey, bool)> = BinaryHeap::with_capacity(count + 1);
{
let mut consider = |h: u64, key_bytes: &[u8], is_cold: bool| {
Expand All @@ -1076,8 +1091,8 @@ fn scan_core(
}
}
};
for key in db.iter_live_keys(now_ms) {
consider(scan_hash48(key.as_bytes()), key.as_bytes(), false);
for (h, key) in &hot_page {
consider(*h, key.as_bytes(), false);
}
// Cold plane: spilled keys with no live hot shadow (pure in-RAM
// index probe — no disk I/O). Partitioned from the hot walk, so no
Expand All @@ -1096,13 +1111,13 @@ fn scan_core(
let h_first = selected.first().map(|e| e.0).unwrap_or(h_last);
if h_first == h_last {
// Whole page one hash value: emit the ENTIRE equal-hash group
// (second filtered pass; unreachable in practice at 48 bits)
// so the cursor may step past it.
// (unreachable in practice at 48 bits) so the cursor may step
// past it. Hot members all live in the already-fetched page:
// equal h48 ⇒ same DashTable segment, and pages are
// whole-segment granular.
let mut extra: Vec<(u64, CompactKey, bool)> = Vec::new();
for key in db.iter_live_keys(now_ms) {
if scan_hash48(key.as_bytes()) == h_last
&& !selected.iter().any(|(_, k, _)| k == key)
{
for (h, key) in &hot_page {
if *h == h_last && !selected.iter().any(|(_, k, _)| k == key) {
extra.push((h_last, key.clone(), false));
}
}
Expand Down Expand Up @@ -1297,6 +1312,12 @@ pub fn scan_readonly(db: &Database, args: &[Frame], now_ms: u64) -> Frame {
Some(c) => c,
None => return Frame::Error(Bytes::from_static(b"ERR invalid cursor")),
};
// Clamp to the 48-bit hash space. Legitimate resumed cursors always
// fit (multi-shard composites are unpacked by `coordinate_scan` before
// reaching here); an out-of-range client cursor would otherwise filter
// out every key (h48 < 2^48) and falsely report "scan complete" on a
// non-empty keyspace.
let cursor = cursor & 0x0000_FFFF_FFFF_FFFF;

let mut match_pattern: Option<&[u8]> = None;
let mut count: usize = 10;
Expand Down Expand Up @@ -1551,6 +1572,24 @@ mod tests {
assert_eq!(seen.len(), 57, "every live key exactly once");
}

/// An out-of-range client cursor (bits above the 48-bit hash space,
/// e.g. a composite cursor replayed against a single-shard server)
/// must be clamped, not silently filter out every key and report a
/// false "scan complete" on a non-empty keyspace.
#[test]
fn scan_out_of_range_cursor_clamps_to_hash_space() {
let mut db = Database::new();
for i in 0..20 {
db.set(
Bytes::from(format!("k:{i}")),
Entry::new_string(Bytes::from_static(b"v")),
);
}
// 2^63: garbage in the shard-index bits, zero in the low 48.
let (_, keys) = scan_page(&mut db, 1u64 << 63, 50);
assert_eq!(keys.len(), 20, "clamped cursor must scan from position 0");
}

/// MATCH + COUNT paging still terminates and honors the filter.
#[test]
fn scan_match_filter_pages_terminate() {
Expand Down
172 changes: 172 additions & 0 deletions src/storage/dashtable/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -276,6 +276,17 @@ impl<V> DashTable<CompactKey, V> {
self.split_count
}

/// Global directory depth (log2 of directory size).
///
/// SCAN's 48-bit cursor mapping (`Database::scan_hot_page`) relies on
/// `depth <= 48` so that equal-top-48-bit hashes always route to the
/// same segment; a depth beyond 48 would need a 2^48-entry directory,
/// unreachable on real hardware.
#[inline]
pub fn directory_depth(&self) -> u32 {
self.depth
}

/// Resident bytes used by the DashTable structural overhead (segments +
/// directory + index map). Does NOT include per-entry key/value data --
/// that is tracked separately by `Database::used_memory`.
Expand Down Expand Up @@ -528,6 +539,68 @@ impl<V> DashTable<CompactKey, V> {
Iter::new(self.segments.collect_refs(), self.len)
}

/// Hash-ordered page collection for SCAN (#368 O(COUNT) walk).
///
/// Because the extendible-hashing directory is indexed by the hash's
/// TOP `global_depth` bits (`segment_index`), ascending directory order
/// is ascending hash order, and directory entry `d` covers exactly the
/// hash range `[d << (64-D), (d+1) << (64-D))`. Segments are therefore
/// range-partitioned: every entry of the segment at a lower directory
/// index hashes below every entry at a higher one. This walk starts at
/// the segment covering `from_hash` and visits segments in ascending
/// range order, stopping as soon as `want` qualifying entries are
/// collected — later segments can only contain larger hashes, so the
/// result is complete without touching the rest of the table.
///
/// A segment with `local_depth < global_depth` occupies a CONTIGUOUS
/// run of directory slots, so alias-dedup is a consecutive store-index
/// comparison.
///
/// Returns entries with `hash_key(key) >= from_hash` passing `alive`,
/// sorted ascending by `(hash, key)`, plus `true` if the walk stopped
/// with unvisited segments remaining (i.e. more entries may exist).
///
/// Split/merge/directory-doubling between calls is safe by
/// construction: the caller's cursor is a position in hash space, and
/// structural churn only changes WHICH segment covers that position,
/// never the set of keys at or above it.
pub fn hash_page<F: Fn(&CompactKey, &V) -> bool>(
&self,
from_hash: u64,
want: usize,
alive: F,
) -> (Vec<(u64, CompactKey)>, bool) {
let mut out: Vec<(u64, CompactKey)> = Vec::with_capacity(want.min(1024));
let start = segment_index(from_hash, self.depth);
let mut last_store_idx = usize::MAX;
let mut dir_idx = start;
while dir_idx < self.directory.len() {
let store_idx = self.directory[dir_idx];
if store_idx != last_store_idx {
last_store_idx = store_idx;
if out.len() >= want {
// Enough collected and at least one unvisited segment
// remains; everything in it hashes above what we have.
return (Self::finish_page(out), true);
}
let seg = self.segments.get(store_idx);
for (k, v) in seg.iter_occupied() {
let h = hash_key(k.as_ref());
if h >= from_hash && alive(k, v) {
out.push((h, k.clone()));
}
Comment on lines +581 to +591

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🚀 Performance & Scalability | 🟠 Major | 🏗️ Heavy lift

Preserve O(COUNT) when alive rejects entries.

The stop condition counts only qualifying entries. An expired-heavy table can therefore visit every segment before producing a small page, leaving SCAN page cost O(keyspace). Bound traversal independently of alive, return a continuation hash, and add an all-rejected/mostly-expired regression test.

Also applies to: 727-750

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@src/storage/dashtable/mod.rs` around lines 570 - 580, Bound the segment
traversal independently of the alive(k, v) filter so expired or rejected entries
cannot force scanning the entire keyspace for one page. In the scan path around
finish_page and the segment iteration, stop after examining the requested COUNT
of candidate entries, preserve the last examined hash as the continuation hash,
and return it even when no entries qualify. Add a regression test covering
all-rejected and mostly-expired tables while preserving normal pagination
behavior.

}
}
dir_idx += 1;
}
(Self::finish_page(out), false)
}

fn finish_page(mut out: Vec<(u64, CompactKey)>) -> Vec<(u64, CompactKey)> {
out.sort_unstable_by(|a, b| (a.0, a.1.as_bytes()).cmp(&(b.0, b.1.as_bytes())));
out
}

/// Return a mutable iterator over `(&Bytes, &mut V)` pairs.
pub fn iter_mut(&mut self) -> IterMut<'_, CompactKey, V> {
let total = self.len;
Expand Down Expand Up @@ -653,6 +726,105 @@ mod tests {
}
}

#[test]
fn hash_page_empty_table_is_terminal() {
let table: DashTable<CompactKey, String> = DashTable::new();
let (page, more) = table.hash_page(0, 16, |_, _| true);
assert!(page.is_empty());
assert!(!more);
}

#[test]
fn hash_page_alive_filter_and_more_flag() {
let mut table: DashTable<CompactKey, String> = DashTable::new();
for i in 0..2000u32 {
table.insert(CompactKey::from(format!("af_{i}")), test_value(i));
}
// Filter out every odd value: only evens may appear.
let (page, _) = table.hash_page(0, 200, |_, v| {
let n: u32 = v.trim_start_matches("value_").parse().unwrap();
n % 2 == 0
});
assert!(!page.is_empty());
for (_, k) in &page {
let n: u32 = std::str::from_utf8(k.as_ref())
.unwrap()
.trim_start_matches("af_")
.parse()
.unwrap();
assert_eq!(n % 2, 0, "alive filter leaked odd key {n}");
}
// A want larger than the table drains everything in one page.
let (all, more) = table.hash_page(0, usize::MAX, |_, _| true);
assert_eq!(all.len(), 2000);
assert!(!more, "full drain must report no further segments");
}

#[test]
fn hash_page_drains_in_hash_order_under_split_churn() {
// 4000 keys force many segment splits + directory doublings; paging
// with concurrent inserts between pages exercises the split-safety
// claim (cursor is a hash-space position, not a structure position).
let mut table: DashTable<CompactKey, String> = DashTable::new();
let mut original: Vec<String> = Vec::with_capacity(4000);
for i in 0..4000u32 {
let k = format!("hp_{i}");
table.insert(CompactKey::from(k.clone()), test_value(i));
original.push(k);
}

let mut seen: std::collections::HashSet<Vec<u8>> = std::collections::HashSet::new();
let mut cursor = 0u64;
let mut churn = 0u32;
loop {
let (page, more) = table.hash_page(cursor, 64, |_, _| true);
let mut prev: Option<(u64, &[u8])> = None;
for (h, k) in &page {
assert!(*h >= cursor, "entry below cursor");
assert_eq!(*h, hash_key(k.as_ref()), "stale hash in page");
if let Some((ph, pk)) = prev {
assert!(
(ph, pk) < (*h, k.as_ref()),
"page not ascending by (hash, key)"
);
}
prev = Some((*h, k.as_ref()));
}
if page.is_empty() {
assert!(!more, "empty page must be terminal");
break;
}
for (_, k) in &page {
assert!(
seen.insert(k.as_ref().to_vec()),
"duplicate key across pages: {:?}",
String::from_utf8_lossy(k.as_ref())
);
}
#[allow(clippy::unwrap_used)] // page verified non-empty above
let last = page.last().unwrap().0;
cursor = last + 1;
if !more {
break;
}
// Structural churn between pages: force splits mid-walk.
for _ in 0..50 {
table.insert(
CompactKey::from(format!("churn_{churn}")),
test_value(churn),
);
churn += 1;
}
}

for k in &original {
assert!(
seen.contains(k.as_bytes()),
"stable key {k} lost during split churn"
);
}
}

#[test]
fn test_new_empty() {
let table: DashTable<CompactKey, String> = DashTable::new();
Expand Down
Loading
Loading