use crate::error::{Error, Result};
use crate::storage::sstable::bti::encode_partition_key_for_bti_trie;
use crate::storage::sstable::bti::parser::FLAG_HAS_HASH_BYTE;
use crate::util::cassandra_murmur3::cassandra_partition_filter_hash_lower_bits;
use std::collections::BTreeMap;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PartitionPayload {
DataOffset(u64),
RowsOffset(u64),
}
#[derive(Debug, Clone)]
pub struct PartitionTrieEntry {
key: [u8; 9],
hash_byte: u8,
payload: PartitionPayload,
}
#[derive(Debug, Default)]
pub struct PartitionsTrieWriter {
entries: Vec<PartitionTrieEntry>,
first_raw: Option<([u8; 9], Vec<u8>)>,
last_raw: Option<([u8; 9], Vec<u8>)>,
}
impl PartitionsTrieWriter {
pub fn new() -> Self {
Self {
entries: Vec::new(),
first_raw: None,
last_raw: None,
}
}
pub fn len(&self) -> usize {
self.entries.len()
}
pub fn is_empty(&self) -> bool {
self.entries.is_empty()
}
#[cfg(test)]
pub(crate) fn retained_raw_key_bytes(&self) -> usize {
self.first_raw.as_ref().map_or(0, |(_, k)| k.len())
+ self.last_raw.as_ref().map_or(0, |(_, k)| k.len())
}
pub fn add_partition(&mut self, raw_key_bytes: &[u8], data_offset: u64) {
self.add_partition_with_payload(raw_key_bytes, PartitionPayload::DataOffset(data_offset));
}
pub fn add_partition_with_payload(&mut self, raw_key_bytes: &[u8], payload: PartitionPayload) {
let key = encode_partition_key_for_bti_trie(raw_key_bytes);
let hash_byte = filter_hash_byte(raw_key_bytes);
match &self.first_raw {
Some((min_key, _)) if key >= *min_key => {}
_ => self.first_raw = Some((key, raw_key_bytes.to_vec())),
}
match &self.last_raw {
Some((max_key, _)) if key <= *max_key => {}
_ => self.last_raw = Some((key, raw_key_bytes.to_vec())),
}
self.entries.push(PartitionTrieEntry {
key,
hash_byte,
payload,
});
}
pub fn finish(self) -> Result<Vec<u8>> {
if self.entries.is_empty() {
return Ok(Vec::new());
}
let mut entries = self.entries;
entries.sort_by(|a, b| a.key.cmp(&b.key));
for w in entries.windows(2) {
if w[0].key == w[1].key {
return Err(Error::InvalidInput(
"duplicate partition trie key (token collision) in Partitions.db".to_string(),
));
}
}
let first_key = self.first_raw.map(|(_, k)| k).unwrap_or_default();
let last_key = self.last_raw.map(|(_, k)| k).unwrap_or_default();
let key_count = entries.len() as i64;
let root = build_trie(&entries);
let mut buf = Vec::new();
let root_offset = write_node(&root, &mut buf)?;
let first_pos = buf.len() as i64;
write_with_short_length(&mut buf, &first_key)?;
write_with_short_length(&mut buf, &last_key)?;
buf.extend_from_slice(&first_pos.to_be_bytes());
buf.extend_from_slice(&key_count.to_be_bytes());
buf.extend_from_slice(&(root_offset as i64).to_be_bytes());
Ok(buf)
}
}
fn write_with_short_length(buf: &mut Vec<u8>, bytes: &[u8]) -> Result<()> {
if bytes.len() > u16::MAX as usize {
return Err(Error::InvalidInput(format!(
"partition key too long for Partitions.db short-length prefix: {} bytes",
bytes.len()
)));
}
buf.extend_from_slice(&(bytes.len() as u16).to_be_bytes());
buf.extend_from_slice(bytes);
Ok(())
}
fn filter_hash_byte(raw_key_bytes: &[u8]) -> u8 {
cassandra_partition_filter_hash_lower_bits(raw_key_bytes) as u8
}
enum TrieBuildNode {
Leaf {
hash_byte: u8,
payload: PartitionPayload,
},
Internal {
children: BTreeMap<u8, TrieBuildNode>,
},
}
fn build_trie(entries: &[PartitionTrieEntry]) -> TrieBuildNode {
let mut root = TrieBuildNode::Internal {
children: BTreeMap::new(),
};
for entry in entries {
insert(&mut root, &entry.key, entry.hash_byte, entry.payload);
}
root
}
fn insert(node: &mut TrieBuildNode, key: &[u8], hash_byte: u8, payload: PartitionPayload) {
match node {
TrieBuildNode::Internal { children } => {
if key.is_empty() {
children
.entry(0)
.or_insert(TrieBuildNode::Leaf { hash_byte, payload });
return;
}
let first = key[0];
let rest = &key[1..];
if rest.is_empty() {
children.insert(first, TrieBuildNode::Leaf { hash_byte, payload });
} else {
let child = children
.entry(first)
.or_insert_with(|| TrieBuildNode::Internal {
children: BTreeMap::new(),
});
insert(child, rest, hash_byte, payload);
}
}
TrieBuildNode::Leaf { .. } => {
}
}
}
#[cfg(test)]
fn serialize_trie(root: &TrieBuildNode) -> Result<Vec<u8>> {
let mut buf = Vec::new();
let root_offset = write_node(root, &mut buf)?;
buf.extend_from_slice(&(root_offset as u64).to_be_bytes());
Ok(buf)
}
fn write_node(node: &TrieBuildNode, buf: &mut Vec<u8>) -> Result<usize> {
match node {
TrieBuildNode::Leaf { hash_byte, payload } => write_leaf(*hash_byte, *payload, buf),
TrieBuildNode::Internal { children } => {
let mut child_offsets: Vec<(u8, usize)> = Vec::with_capacity(children.len());
for (&byte, child) in children.iter() {
let off = write_node(child, buf)?;
child_offsets.push((byte, off));
}
if child_offsets.len() == 256 {
write_dense(&child_offsets, buf)
} else {
write_sparse(&child_offsets, buf)
}
}
}
}
fn write_leaf(hash_byte: u8, payload: PartitionPayload, buf: &mut Vec<u8>) -> Result<usize> {
let position: i64 = match payload {
PartitionPayload::DataOffset(data_offset) => {
if data_offset > i64::MAX as u64 {
return Err(Error::InvalidInput(format!(
"Data.db offset {data_offset} too large to encode as a signed BTI position"
)));
}
!(data_offset as i64)
}
PartitionPayload::RowsOffset(rows_offset) => {
if rows_offset > i64::MAX as u64 {
return Err(Error::InvalidInput(format!(
"Rows.db offset {rows_offset} too large to encode as a signed BTI position"
)));
}
rows_offset as i64
}
};
let position_bytes = sized_ints_non_zero_size(position);
debug_assert!((1..=8).contains(&position_bytes));
let payload_bits = FLAG_HAS_HASH_BYTE + (position_bytes as u8 - 1);
debug_assert!(payload_bits <= 16);
let offset = buf.len();
let header = payload_bits & 0x0F;
buf.push(header);
buf.push(hash_byte);
write_sized_int_be(buf, position, position_bytes);
Ok(offset)
}
fn write_sparse(child_offsets: &[(u8, usize)], buf: &mut Vec<u8>) -> Result<usize> {
let count = child_offsets.len();
if count == 0 {
return Err(Error::InvalidInput(
"internal BTI trie node has no children".to_string(),
));
}
if count > 255 {
return Err(Error::InvalidInput(format!(
"BTI Sparse node fan-out {count} exceeds 255; Dense node required"
)));
}
let node_offset = buf.len();
let max_delta = child_offsets
.iter()
.map(|(_, child_off)| node_offset - child_off)
.max()
.unwrap_or(0);
let (ordinal, ptr_bytes) = sparse_ordinal_for_delta(max_delta as u64)?;
let header = ordinal << 4;
buf.push(header);
buf.push(count as u8);
for (byte, _) in child_offsets {
buf.push(*byte);
}
for (_, child_off) in child_offsets {
let delta = (node_offset - child_off) as u64;
write_be_unsigned(buf, delta, ptr_bytes);
}
Ok(node_offset)
}
fn write_dense(child_offsets: &[(u8, usize)], buf: &mut Vec<u8>) -> Result<usize> {
debug_assert_eq!(child_offsets.len(), 256, "write_dense expects full fan-out");
if child_offsets.len() != 256 {
return Err(Error::InvalidInput(format!(
"write_dense requires a full 256-child node, got {}",
child_offsets.len()
)));
}
let start_byte = child_offsets[0].0; let node_offset = buf.len();
let max_delta = child_offsets
.iter()
.map(|(_, child_off)| node_offset - child_off)
.max()
.unwrap_or(0);
let (ordinal, ptr_bytes) = dense_ordinal_for_delta(max_delta as u64)?;
let header = ordinal << 4;
buf.push(header);
buf.push(start_byte);
buf.push((child_offsets.len() - 1) as u8);
for (_, child_off) in child_offsets {
let delta = (node_offset - child_off) as u64;
write_be_unsigned(buf, delta, ptr_bytes);
}
Ok(node_offset)
}
fn dense_ordinal_for_delta(max_delta: u64) -> Result<(u8, usize)> {
if max_delta <= 0xFFFF {
Ok((11, 2))
} else if max_delta <= 0xFF_FFFF {
Ok((12, 3))
} else if max_delta <= 0xFFFF_FFFF {
Ok((13, 4))
} else if max_delta <= 0xFF_FFFF_FFFF {
Ok((14, 5))
} else {
Ok((15, 8))
}
}
fn sparse_ordinal_for_delta(max_delta: u64) -> Result<(u8, usize)> {
if max_delta <= 0xFF {
Ok((5, 1))
} else if max_delta <= 0xFFFF {
Ok((7, 2))
} else if max_delta <= 0xFF_FFFF {
Ok((8, 3))
} else if max_delta <= 0xFF_FFFF_FFFF {
Ok((9, 5))
} else {
Err(Error::InvalidInput(format!(
"BTI Partitions.db trie too large: backward delta {max_delta} exceeds 40 bits"
)))
}
}
fn sized_ints_non_zero_size(value: i64) -> usize {
let abs_value = if value < 0 { !value } else { value } as u64;
if abs_value == 0 {
return 1;
}
let significant_bits = 64 - abs_value.leading_zeros() as usize;
(significant_bits + 1).div_ceil(8).clamp(1, 8)
}
fn write_sized_int_be(buf: &mut Vec<u8>, value: i64, bytes: usize) {
let raw = value as u64;
write_be_unsigned(buf, raw, bytes);
}
fn write_be_unsigned(buf: &mut Vec<u8>, value: u64, bytes: usize) {
let all = value.to_be_bytes();
buf.extend_from_slice(&all[8 - bytes..]);
}
#[derive(Debug, Clone)]
pub struct RowIndexBlock {
pub separator_key: Vec<u8>,
pub block_offset: u64,
pub open_marker: Option<(i32, i64)>,
}
#[derive(Debug, Clone)]
struct PartitionRowIndex {
partition_key: Vec<u8>,
data_position: u64,
blocks: Vec<RowIndexBlock>,
partition_deletion: Option<(i32, i64)>,
}
#[derive(Debug, Default)]
pub struct RowsTrieWriter {
partitions: Vec<PartitionRowIndex>,
}
impl RowsTrieWriter {
pub fn new() -> Self {
Self {
partitions: Vec::new(),
}
}
pub fn len(&self) -> usize {
self.partitions.len()
}
pub fn is_empty(&self) -> bool {
self.partitions.is_empty()
}
pub fn add_partition_row_index(
&mut self,
partition_key: &[u8],
data_position: u64,
blocks: Vec<RowIndexBlock>,
partition_deletion: Option<(i32, i64)>,
) {
self.partitions.push(PartitionRowIndex {
partition_key: partition_key.to_vec(),
data_position,
blocks,
partition_deletion,
});
}
pub fn finish(self) -> Result<(Vec<u8>, Vec<u64>)> {
let mut buf = Vec::new();
let mut rows_offsets = Vec::with_capacity(self.partitions.len());
for part in &self.partitions {
if part.blocks.is_empty() {
return Err(Error::InvalidInput(
"Rows.db: wide partition must have at least one row-index block".to_string(),
));
}
for w in part.blocks.windows(2) {
if w[0].separator_key >= w[1].separator_key {
return Err(Error::InvalidInput(
"Rows.db: row-index separators must be strictly ascending".to_string(),
));
}
}
let root = build_row_trie(&part.blocks);
let trie_root = write_row_node(&root, &mut buf)?;
let rows_offset = buf.len();
write_trie_index_entry(
&mut buf,
&part.partition_key,
part.data_position,
trie_root,
part.blocks.len() as u64,
part.partition_deletion,
)?;
rows_offsets.push(rows_offset as u64);
}
Ok((buf, rows_offsets))
}
}
enum RowTrieBuildNode {
Leaf {
block_offset: u64,
open_marker: Option<(i32, i64)>,
},
Internal {
children: BTreeMap<u8, RowTrieBuildNode>,
},
}
fn build_row_trie(blocks: &[RowIndexBlock]) -> RowTrieBuildNode {
let mut root = RowTrieBuildNode::Internal {
children: BTreeMap::new(),
};
for b in blocks {
insert_row(&mut root, &b.separator_key, b.block_offset, b.open_marker);
}
root
}
fn insert_row(
node: &mut RowTrieBuildNode,
key: &[u8],
block_offset: u64,
open_marker: Option<(i32, i64)>,
) {
match node {
RowTrieBuildNode::Internal { children } => {
if key.is_empty() {
children.entry(0).or_insert(RowTrieBuildNode::Leaf {
block_offset,
open_marker,
});
return;
}
let first = key[0];
let rest = &key[1..];
if rest.is_empty() {
children.insert(
first,
RowTrieBuildNode::Leaf {
block_offset,
open_marker,
},
);
} else {
let child = children
.entry(first)
.or_insert_with(|| RowTrieBuildNode::Internal {
children: BTreeMap::new(),
});
insert_row(child, rest, block_offset, open_marker);
}
}
RowTrieBuildNode::Leaf { .. } => {
}
}
}
fn write_row_node(node: &RowTrieBuildNode, buf: &mut Vec<u8>) -> Result<usize> {
match node {
RowTrieBuildNode::Leaf {
block_offset,
open_marker,
} => write_row_leaf(*block_offset, *open_marker, buf),
RowTrieBuildNode::Internal { children } => {
let mut child_offsets: Vec<(u8, usize)> = Vec::with_capacity(children.len());
for (&byte, child) in children.iter() {
let off = write_row_node(child, buf)?;
child_offsets.push((byte, off));
}
if child_offsets.len() == 256 {
write_dense(&child_offsets, buf)
} else {
write_sparse(&child_offsets, buf)
}
}
}
}
fn write_row_leaf(
block_offset: u64,
open_marker: Option<(i32, i64)>,
buf: &mut Vec<u8>,
) -> Result<usize> {
if block_offset > i64::MAX as u64 {
return Err(Error::InvalidInput(format!(
"Rows.db block offset {block_offset} too large to encode as SizedInts"
)));
}
let offset_bytes = sized_ints_non_zero_size(block_offset as i64);
if !(1..=7).contains(&offset_bytes) {
return Err(Error::InvalidInput(format!(
"Rows.db block offset {block_offset} needs {offset_bytes} SizedInts bytes; \
expected 1..=7"
)));
}
let mut payload_bits = offset_bytes as u8;
if open_marker.is_some() {
payload_bits |= crate::storage::sstable::bti::parser::FLAG_OPEN_MARKER;
}
let offset = buf.len();
buf.push(payload_bits & 0x0F);
write_sized_int_be(buf, block_offset as i64, offset_bytes);
if let Some((ldt, mfda)) = open_marker {
write_da_deletion_time(buf, Some((ldt, mfda)));
}
Ok(offset)
}
fn write_trie_index_entry(
buf: &mut Vec<u8>,
partition_key: &[u8],
data_position: u64,
trie_root: usize,
block_count: u64,
partition_deletion: Option<(i32, i64)>,
) -> Result<()> {
let entry_start = buf.len();
let key_length = partition_key.len();
let key_length_u16 = u16::try_from(key_length).map_err(|_| {
Error::InvalidInput(format!(
"Rows.db TrieIndexEntry: partition key length {key_length} exceeds u16"
))
})?;
buf.extend_from_slice(&key_length_u16.to_be_bytes());
buf.extend_from_slice(partition_key);
write_unsigned_vint(buf, data_position);
let base = entry_start + key_length;
let root_delta = trie_root as i64 - base as i64;
write_signed_vint(buf, root_delta);
write_unsigned_vint(buf, block_count);
write_da_deletion_time(buf, partition_deletion);
Ok(())
}
fn write_unsigned_vint(buf: &mut Vec<u8>, value: u64) {
let extra_bytes = if value == 0 {
0
} else {
let significant_bits = 64 - value.leading_zeros() as usize;
let mut n = 0usize;
while n < 8 && (7 - n) + 8 * n < significant_bits {
n += 1;
}
n
};
if extra_bytes == 0 {
buf.push(value as u8);
return;
}
if extra_bytes >= 8 {
buf.push(0xFF);
buf.extend_from_slice(&value.to_be_bytes());
return;
}
let data_bits_first = 7 - extra_bytes;
let leading_ones: u8 = (!0u8) << (8 - extra_bytes);
let total_bytes = extra_bytes + 1;
let mut bytes = value.to_be_bytes().to_vec();
let tail = bytes.split_off(8 - total_bytes);
let first = leading_ones | (tail[0] & ((1u8 << data_bits_first) - 1));
buf.push(first);
buf.extend_from_slice(&tail[1..]);
}
fn write_signed_vint(buf: &mut Vec<u8>, value: i64) {
let zigzag = ((value << 1) ^ (value >> 63)) as u64;
write_unsigned_vint(buf, zigzag);
}
fn write_da_deletion_time(buf: &mut Vec<u8>, deletion: Option<(i32, i64)>) {
match deletion {
None => buf.push(0x80),
Some((local_deletion_time, marked_for_delete_at)) => {
buf.extend_from_slice(&marked_for_delete_at.to_be_bytes());
buf.extend_from_slice(&(local_deletion_time as u32).to_be_bytes());
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::storage::sstable::bti::sized_ints;
use crate::storage::sstable::bti::{lookup_raw_key_in_bti_partitions_db, BtiPartitionLocation};
use std::io::Cursor;
#[test]
fn sized_int_size_matches_reader() {
let values = [
0i64,
1,
-1,
127,
-128,
128,
-129,
255,
-256,
32767,
-32768,
32768,
-32769,
i64::MAX,
i64::MIN,
!0i64,
!63i64,
!125i64,
!1000i64,
!1_000_000i64,
!300_000_000_000i64,
];
for v in values {
assert_eq!(
sized_ints_non_zero_size(v),
sized_ints::non_zero_size(v),
"size mismatch for {v}"
);
}
}
#[test]
fn sized_int_write_read_roundtrip() {
for v in [0i64, !0i64, !63i64, !125i64, !1_000_000i64, i64::MIN] {
let n = sized_ints_non_zero_size(v);
let mut buf = Vec::new();
write_sized_int_be(&mut buf, v, n);
assert_eq!(buf.len(), n);
let mut cur = Cursor::new(buf);
let got = sized_ints::read(&mut cur, n).unwrap();
assert_eq!(got, v, "SizedInt roundtrip failed for {v}");
}
}
fn assert_roundtrip(keys_and_offsets: &[(Vec<u8>, u64)]) {
let mut w = PartitionsTrieWriter::new();
for (k, off) in keys_and_offsets {
w.add_partition(k, *off);
}
let bytes = w.finish().expect("finish trie");
assert!(bytes.len() >= 8, "trie must include 8-byte footer");
for (k, expected) in keys_and_offsets {
let mut cur = Cursor::new(bytes.clone());
let loc = lookup_raw_key_in_bti_partitions_db(&mut cur, k)
.expect("lookup")
.unwrap_or_else(|| panic!("key {k:?} not found in written trie"));
match loc {
BtiPartitionLocation::DataOffset(got) => assert_eq!(
got, *expected,
"key {k:?}: expected DataOffset({expected}) got DataOffset({got})"
),
BtiPartitionLocation::RowsOffset(r) => {
panic!("key {k:?}: phase-1 writer must emit DataOffset, got RowsOffset({r})")
}
}
}
}
#[test]
fn canonical_filter_hash_byte_matches_real_bti_fixture() {
let vectors: [(Vec<u8>, u8); 3] = [
(vec![0x22u8; 16], 0x24),
(vec![0x11u8; 16], 0x22),
(vec![0x33u8; 16], 0xf4),
];
for (raw_key, expected) in vectors {
let got = filter_hash_byte(&raw_key);
assert_eq!(
got, expected,
"canonical hash byte mismatch for key {raw_key:02x?}: \
expected 0x{expected:02x} (from real da-2-bti-Partitions.db), got 0x{got:02x}"
);
}
for (raw_key, expected) in [(vec![0x22u8; 16], 0x90u8), (vec![0x11u8; 16], 0xbc)] {
let token = crate::util::cassandra_murmur3::cassandra_murmur3_token(&raw_key);
let placeholder = (((token as u64) ^ 0x8000_0000_0000_0000u64) >> 56) as u8;
assert_eq!(placeholder, expected, "placeholder reference value drifted");
assert_ne!(
placeholder,
filter_hash_byte(&raw_key),
"canonical hash byte must differ from the old token-derived placeholder"
);
}
}
#[test]
fn written_leaf_hash_byte_is_canonical() {
let raw_key = vec![0x22u8; 16];
let mut w = PartitionsTrieWriter::new();
w.add_partition(&raw_key, 0);
let bytes = w.finish().expect("finish trie");
assert_eq!(
bytes[0], 0x08,
"expected PayloadOnly leaf with payloadBits=8"
);
assert_eq!(
bytes[1],
filter_hash_byte(&raw_key),
"serialized leaf hash byte must be the canonical value"
);
assert_eq!(bytes[1], 0x24, "canonical hash byte for UUID 2222… is 0x24");
}
#[test]
fn empty_trie_is_empty_bytes() {
let w = PartitionsTrieWriter::new();
assert!(w.finish().unwrap().is_empty());
}
#[test]
fn single_partition_roundtrip() {
assert_roundtrip(&[(vec![0x11u8; 16], 0)]);
}
#[test]
fn three_uuid_partitions_roundtrip() {
assert_roundtrip(&[
(vec![0x11u8; 16], 63),
(vec![0x22u8; 16], 0),
(vec![0x33u8; 16], 125),
]);
}
#[test]
fn interior_raw_keys_are_not_retained() {
let mut w = PartitionsTrieWriter::new();
for i in 0u64..1000 {
let mut k = vec![0u8; 16];
k[0..8].copy_from_slice(&i.to_be_bytes());
w.add_partition(&k, i);
}
assert!(
w.retained_raw_key_bytes() <= 32,
"retained {} raw-key bytes; expected <= 32 (first+last only)",
w.retained_raw_key_bytes()
);
}
#[test]
fn boundary_raw_keys_match_sorted_first_last() {
let raws: Vec<Vec<u8>> = [0x33u8, 0x11, 0x88, 0x22, 0x55, 0x44]
.iter()
.map(|b| vec![*b; 16])
.collect();
let mut w = PartitionsTrieWriter::new();
for (i, r) in raws.iter().enumerate() {
w.add_partition(r, i as u64);
}
let mut sorted: Vec<(_, Vec<u8>)> = raws
.iter()
.map(|r| (encode_partition_key_for_bti_trie(r), r.clone()))
.collect();
sorted.sort_by(|a, b| a.0.cmp(&b.0));
let expect_first = sorted.first().map(|(_, r)| r.clone()).unwrap();
let expect_last = sorted.last().map(|(_, r)| r.clone()).unwrap();
assert_eq!(
w.first_raw.as_ref().map(|(_, k)| k.clone()),
Some(expect_first),
"first boundary raw key must be the min-trie-key partition"
);
assert_eq!(
w.last_raw.as_ref().map(|(_, k)| k.clone()),
Some(expect_last),
"last boundary raw key must be the max-trie-key partition"
);
}
#[test]
fn large_offsets_roundtrip() {
assert_roundtrip(&[
(vec![0xA1u8; 16], 1_000_000),
(vec![0xB2u8; 16], 300_000_000_000),
(vec![0xC3u8; 16], 5),
]);
}
#[test]
fn many_partitions_roundtrip() {
let mut data = Vec::new();
for i in 0u64..200 {
let mut key = vec![0u8; 16];
key[0..8].copy_from_slice(&i.to_be_bytes());
key[8..16].copy_from_slice(&(i.wrapping_mul(2654435761)).to_be_bytes());
data.push((key, i * 37));
}
assert_roundtrip(&data);
}
#[test]
fn duplicate_key_is_rejected() {
let mut w = PartitionsTrieWriter::new();
w.add_partition(&[0x55u8; 16], 0);
w.add_partition(&[0x55u8; 16], 100);
assert!(w.finish().is_err());
}
#[test]
fn full_256_fanout_internal_node_serializes_and_roundtrips() {
use std::collections::BTreeMap;
let mut inner_children: BTreeMap<u8, TrieBuildNode> = BTreeMap::new();
for b in 0u16..=255 {
inner_children.insert(
b as u8,
TrieBuildNode::Leaf {
hash_byte: b as u8,
payload: PartitionPayload::DataOffset((b as u64) * 17),
},
);
}
let inner = TrieBuildNode::Internal {
children: inner_children,
};
let mut root_children: BTreeMap<u8, TrieBuildNode> = BTreeMap::new();
root_children.insert(0xFF, inner);
let root = TrieBuildNode::Internal {
children: root_children,
};
let bytes = serialize_trie(&root).expect("256-fan-out node must serialize");
for b in 0u16..=255 {
let key = [0xFFu8, b as u8];
let loc = lookup_key_in_trie(&bytes, &key)
.unwrap_or_else(|| panic!("byte {b} not found in 256-fan-out trie"));
assert_eq!(
loc,
(b as u64) * 17,
"byte {b}: wrong Data.db offset resolved"
);
}
}
fn lookup_key_in_trie(bytes: &[u8], key: &[u8]) -> Option<u64> {
use crate::storage::sstable::bti::lookup_partition_in_bti_file;
let mut cur = Cursor::new(bytes.to_vec());
match lookup_partition_in_bti_file(&mut cur, key).ok()?? {
BtiPartitionLocation::DataOffset(o) => Some(o),
BtiPartitionLocation::RowsOffset(_) => None,
}
}
#[test]
fn unsigned_vint_roundtrips_through_reader() {
use crate::storage::sstable::bti::parser::read_unsigned_vint_from_slice_for_test as read_u;
let values: [u64; 20] = [
0,
1,
63,
64,
127,
128,
255,
256,
16_383,
16_384,
65_535,
65_536,
1_000_000,
300_000_000_000,
(1u64 << 35) - 1,
1u64 << 35,
(1u64 << 49) - 1,
1u64 << 49,
u64::MAX - 1,
u64::MAX,
];
for v in values {
let mut buf = Vec::new();
write_unsigned_vint(&mut buf, v);
let (got, n) = read_u(&buf).expect("read");
assert_eq!(got, v, "unsigned vint roundtrip failed for {v}: {buf:02x?}");
assert_eq!(n, buf.len(), "consumed all bytes for {v}");
}
}
#[test]
fn signed_vint_roundtrips_through_reader() {
use crate::storage::sstable::bti::parser::read_signed_vint_from_slice_for_test as read_s;
for v in [
0i64,
1,
-1,
10,
-10,
127,
-128,
1000,
-1000,
i32::MIN as i64,
i64::MAX,
i64::MIN,
] {
let mut buf = Vec::new();
write_signed_vint(&mut buf, v);
let (got, _n) = read_s(&buf).expect("read");
assert_eq!(got, v, "signed vint roundtrip failed for {v}: {buf:02x?}");
}
}
#[test]
fn rows_db_single_wide_partition_roundtrips() {
use crate::storage::sstable::bti::{iterate_rows_in_bti_trie, resolve_rows_db_entry};
let sep = |ck: i32| ((ck as u32) ^ 0x8000_0000).to_be_bytes().to_vec();
let blocks = vec![
RowIndexBlock {
separator_key: sep(8),
block_offset: 16_512,
open_marker: None,
},
RowIndexBlock {
separator_key: sep(16),
block_offset: 33_024,
open_marker: None,
},
RowIndexBlock {
separator_key: sep(24),
block_offset: 49_536,
open_marker: None,
},
];
let raw_pk = 1i32.to_be_bytes().to_vec();
let mut w = RowsTrieWriter::new();
w.add_partition_row_index(&raw_pk, 0, blocks.clone(), None);
let (rows_db, offsets) = w.finish().expect("finish Rows.db");
assert_eq!(offsets.len(), 1);
let rows_offset = offsets[0] as usize;
assert!(
iterate_rows_in_bti_trie(&rows_db, rows_offset).is_err(),
"RowsOffset is a TrieIndexEntry, not a trie root"
);
let header = resolve_rows_db_entry(&rows_db, rows_offset).expect("resolve entry");
assert_eq!(header.data_position, 0);
assert_eq!(header.block_count, blocks.len() as u32);
assert_eq!(header.partition_deletion, None, "LIVE sentinel → None");
let entries =
iterate_rows_in_bti_trie(&rows_db, header.trie_root).expect("traverse from root");
assert_eq!(entries.len(), blocks.len());
for (got, expected) in entries.iter().zip(blocks.iter()) {
assert_eq!(got.0, expected.separator_key, "separator key");
assert_eq!(got.1.data_offset, expected.block_offset, "block offset");
assert_eq!(got.1.open_marker, None);
}
}
#[test]
fn rows_db_multiple_wide_partitions_roundtrip() {
use crate::storage::sstable::bti::{iterate_rows_in_bti_trie, resolve_rows_db_entry};
let sep = |ck: i32| ((ck as u32) ^ 0x8000_0000).to_be_bytes().to_vec();
let mk = |base: u64| {
vec![
RowIndexBlock {
separator_key: sep(8),
block_offset: base + 16_512,
open_marker: None,
},
RowIndexBlock {
separator_key: sep(16),
block_offset: base + 33_024,
open_marker: None,
},
]
};
let mut w = RowsTrieWriter::new();
w.add_partition_row_index(&1i32.to_be_bytes(), 0, mk(0), None);
w.add_partition_row_index(&2i32.to_be_bytes(), 700_000, mk(0), None);
w.add_partition_row_index(&3i32.to_be_bytes(), 1_400_000, mk(0), None);
let (rows_db, offsets) = w.finish().expect("finish");
assert_eq!(offsets.len(), 3);
let data_positions = [0u64, 700_000, 1_400_000];
for (i, &ro) in offsets.iter().enumerate() {
let header = resolve_rows_db_entry(&rows_db, ro as usize).expect("resolve");
assert_eq!(
header.data_position, data_positions[i],
"partition {i} data position"
);
assert_eq!(header.block_count, 2);
let entries = iterate_rows_in_bti_trie(&rows_db, header.trie_root).expect("traverse");
assert_eq!(entries.len(), 2, "partition {i} block count");
assert_eq!(entries[0].0, sep(8));
assert_eq!(entries[1].0, sep(16));
}
}
#[test]
fn rows_db_empty_writer_is_zero_bytes() {
let w = RowsTrieWriter::new();
let (bytes, offsets) = w.finish().expect("finish empty");
assert!(bytes.is_empty(), "empty Rows.db must be 0 bytes");
assert!(offsets.is_empty());
}
#[test]
fn rows_db_open_marker_roundtrips() {
use crate::storage::sstable::bti::{iterate_rows_in_bti_trie, resolve_rows_db_entry};
let sep = |ck: i32| ((ck as u32) ^ 0x8000_0000).to_be_bytes().to_vec();
let blocks = vec![
RowIndexBlock {
separator_key: sep(8),
block_offset: 16_512,
open_marker: Some((1_700_000_000, 1_700_000_000_000_000)),
},
RowIndexBlock {
separator_key: sep(16),
block_offset: 33_024,
open_marker: None,
},
];
let mut w = RowsTrieWriter::new();
w.add_partition_row_index(&7i32.to_be_bytes(), 0, blocks.clone(), None);
let (rows_db, offsets) = w.finish().expect("finish");
let header = resolve_rows_db_entry(&rows_db, offsets[0] as usize).expect("resolve");
let entries = iterate_rows_in_bti_trie(&rows_db, header.trie_root).expect("traverse");
assert_eq!(
entries[0].1.open_marker,
Some((1_700_000_000, 1_700_000_000_000_000))
);
assert_eq!(entries[1].1.open_marker, None);
}
#[test]
fn rows_db_partition_deletion_roundtrips() {
use crate::storage::sstable::bti::resolve_rows_db_entry;
let sep = |ck: i32| ((ck as u32) ^ 0x8000_0000).to_be_bytes().to_vec();
let blocks = vec![
RowIndexBlock {
separator_key: sep(8),
block_offset: 16_512,
open_marker: None,
},
RowIndexBlock {
separator_key: sep(16),
block_offset: 33_024,
open_marker: None,
},
];
let mut w = RowsTrieWriter::new();
w.add_partition_row_index(&9i32.to_be_bytes(), 0, blocks, Some((1234, 5678)));
let (rows_db, offsets) = w.finish().expect("finish");
let header = resolve_rows_db_entry(&rows_db, offsets[0] as usize).expect("resolve");
assert_eq!(header.partition_deletion, Some((1234, 5678)));
}
#[test]
fn rows_db_rejects_non_ascending_separators() {
let sep = |ck: i32| ((ck as u32) ^ 0x8000_0000).to_be_bytes().to_vec();
let blocks = vec![
RowIndexBlock {
separator_key: sep(16),
block_offset: 16_512,
open_marker: None,
},
RowIndexBlock {
separator_key: sep(8),
block_offset: 33_024,
open_marker: None,
},
];
let mut w = RowsTrieWriter::new();
w.add_partition_row_index(&1i32.to_be_bytes(), 0, blocks, None);
assert!(w.finish().is_err());
}
#[test]
fn partition_leaf_rows_offset_roundtrips() {
let raw_key = vec![0x11u8; 16];
let mut w = PartitionsTrieWriter::new();
w.add_partition_with_payload(&raw_key, PartitionPayload::RowsOffset(242));
let bytes = w.finish().expect("finish");
let mut cur = Cursor::new(bytes);
let loc = lookup_raw_key_in_bti_partitions_db(&mut cur, &raw_key)
.expect("lookup")
.expect("found");
match loc {
BtiPartitionLocation::RowsOffset(o) => assert_eq!(o, 242),
BtiPartitionLocation::DataOffset(o) => {
panic!("expected RowsOffset(242), got DataOffset({o})")
}
}
}
}