use crate::{error::Error, storage::sstable::bti::node::BtiResult};
use std::io::{Read, Seek, SeekFrom};
use super::node_decode::parse_bti_node;
use super::partitions::{payload_start_in_node, sized_ints_read_from_slice};
use super::traversal::{dfs_collect_in_order, load_bti_trie_via_footer};
#[allow(unused_imports)] use super::partitions::BtiPartitionLocation;
pub const FLAG_OPEN_MARKER: u8 = 0x8;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct BtiRowIndexEntry {
pub data_offset: u64,
pub open_marker: Option<(i32, i64)>,
}
fn read_unsigned_vint_from_slice(data: &[u8]) -> BtiResult<(u64, usize)> {
if data.is_empty() {
return Err(Error::Parse(
"Rows.db payload: unexpected end of data reading unsigned vint".to_string(),
));
}
let first = data[0];
let extra_bytes = first.leading_ones() as usize;
if extra_bytes > 8 {
return Err(Error::Parse(format!(
"Rows.db payload: invalid unsigned vint first byte 0x{first:02x}"
)));
}
let total = extra_bytes + 1;
if data.len() < total {
return Err(Error::Parse(format!(
"Rows.db payload: unsigned vint needs {total} bytes, have {}",
data.len()
)));
}
let mut value: u64 = if extra_bytes >= 8 {
0
} else {
let data_bits = 8 - extra_bytes - 1;
let mask = if data_bits == 0 {
0
} else {
(1u16 << data_bits) - 1
};
(first as u16 & mask) as u64
};
for &b in &data[1..total] {
value = (value << 8) | (b as u64);
}
Ok((value, total))
}
fn read_signed_vint_from_slice(data: &[u8]) -> BtiResult<(i64, usize)> {
let (u, n) = read_unsigned_vint_from_slice(data)?;
let value = ((u >> 1) as i64) ^ -((u & 1) as i64);
Ok((value, n))
}
#[doc(hidden)]
pub fn read_unsigned_vint_from_slice_for_test(data: &[u8]) -> BtiResult<(u64, usize)> {
read_unsigned_vint_from_slice(data)
}
#[doc(hidden)]
pub fn read_signed_vint_from_slice_for_test(data: &[u8]) -> BtiResult<(i64, usize)> {
read_signed_vint_from_slice(data)
}
const DA_DELETION_TIME_LIVE_SENTINEL: u8 = 0x80;
const DA_DELETION_TIME_BODY_LEN: usize = 12;
fn decode_da_deletion_time(data: &[u8], start: usize) -> BtiResult<(Option<(i32, i64)>, usize)> {
if start >= data.len() {
return Err(Error::Parse(format!(
"DA DeletionTime: start {start} beyond buffer size {}",
data.len()
)));
}
if data[start] == DA_DELETION_TIME_LIVE_SENTINEL {
return Ok((None, 1));
}
if start + DA_DELETION_TIME_BODY_LEN > data.len() {
return Err(Error::Parse(format!(
"DA DeletionTime: non-live value needs {DA_DELETION_TIME_BODY_LEN} bytes, have {}",
data.len().saturating_sub(start)
)));
}
let b = &data[start..start + DA_DELETION_TIME_BODY_LEN];
let marked_for_delete_at = i64::from_be_bytes([b[0], b[1], b[2], b[3], b[4], b[5], b[6], b[7]]);
let local_deletion_time = u32::from_be_bytes([b[8], b[9], b[10], b[11]]) as i32;
Ok((
Some((local_deletion_time, marked_for_delete_at)),
DA_DELETION_TIME_BODY_LEN,
))
}
pub fn decode_bti_row_payload(
trie_data: &[u8],
payload_start: usize,
payload_bits: u8,
) -> BtiResult<BtiRowIndexEntry> {
if payload_start > trie_data.len() {
return Err(Error::Parse(format!(
"Rows.db payload start {payload_start} beyond trie size {}",
trie_data.len()
)));
}
let offset_bytes = (payload_bits & !FLAG_OPEN_MARKER) as usize;
if offset_bytes == 0 || offset_bytes > 7 {
return Err(Error::Parse(format!(
"Rows.db payload: invalid SizedInts byte count {offset_bytes} \
(payload_bits=0x{payload_bits:02x}); expected 1..=7"
)));
}
if payload_start + offset_bytes > trie_data.len() {
return Err(Error::Parse(format!(
"Rows.db payload: SizedInts offset needs {offset_bytes} bytes, have {}",
trie_data.len().saturating_sub(payload_start)
)));
}
let raw = sized_ints_read_from_slice(&trie_data[payload_start..payload_start + offset_bytes])?;
let data_offset = raw as u64;
let open_marker = if payload_bits & FLAG_OPEN_MARKER != 0 {
let dt_start = payload_start + offset_bytes;
let (deletion, _consumed) = decode_da_deletion_time(trie_data, dt_start)?;
deletion
} else {
None
};
Ok(BtiRowIndexEntry {
data_offset,
open_marker,
})
}
fn read_row_node_payload(
trie_data: &[u8],
node_offset: usize,
) -> BtiResult<Option<BtiRowIndexEntry>> {
if node_offset >= trie_data.len() {
return Err(Error::Parse(format!(
"Rows.db payload read: node_offset {node_offset} out of bounds"
)));
}
let header_byte = trie_data[node_offset];
let ordinal = (header_byte >> 4) & 0x0F;
let payload_flags = header_byte & 0x0F;
if ordinal == 1 || ordinal == 3 {
return Ok(None);
}
if ordinal == 0 {
if payload_flags == 0 {
return Err(Error::Parse(
"Rows.db PayloadOnly node has zero payload_flags".to_string(),
));
}
let payload_start = node_offset + 1;
Ok(Some(decode_bti_row_payload(
trie_data,
payload_start,
payload_flags,
)?))
} else if payload_flags != 0 {
let node = parse_bti_node(&trie_data[node_offset..], node_offset as u64)?;
let payload_start = payload_start_in_node(&node, trie_data, node_offset)?;
Ok(Some(decode_bti_row_payload(
trie_data,
payload_start,
payload_flags,
)?))
} else {
Ok(None)
}
}
pub(crate) fn dfs_collect_row_entries(
trie_data: &[u8],
root_offset: usize,
) -> BtiResult<Vec<(Vec<u8>, BtiRowIndexEntry)>> {
dfs_collect_in_order(trie_data, root_offset, |data, off| {
read_row_node_payload(data, off)
})
}
pub fn iterate_rows_in_bti_trie(
trie_data: &[u8],
root_offset: usize,
) -> BtiResult<Vec<(Vec<u8>, BtiRowIndexEntry)>> {
dfs_collect_row_entries(trie_data, root_offset)
}
pub type BtiRowIndexEntryWithKey = (Vec<u8>, BtiRowIndexEntry);
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct BtiRowIndexHeader {
pub data_position: u64,
pub trie_root: usize,
pub block_count: u32,
pub partition_deletion: Option<(i32, i64)>,
}
pub fn resolve_rows_db_entry(rows_db: &[u8], rows_offset: usize) -> BtiResult<BtiRowIndexHeader> {
if rows_offset + 2 > rows_db.len() {
return Err(Error::Parse(format!(
"Rows.db entry: rows_offset {rows_offset} + 2 (key length) exceeds file size {}",
rows_db.len()
)));
}
let key_length = u16::from_be_bytes([rows_db[rows_offset], rows_db[rows_offset + 1]]) as usize;
let entry_start = rows_offset + 2 + key_length;
if entry_start > rows_db.len() {
return Err(Error::Parse(format!(
"Rows.db entry: key length {key_length} at offset {rows_offset} overruns file size {}",
rows_db.len()
)));
}
let base = rows_offset + key_length;
let mut cur = entry_start;
let (data_position, n) = read_unsigned_vint_from_slice(&rows_db[cur..])?;
cur += n;
let (root_delta, n) = read_signed_vint_from_slice(&rows_db[cur..])?;
cur += n;
let trie_root_signed = root_delta + base as i64;
if trie_root_signed < 0 || (trie_root_signed as usize) >= rows_db.len() {
return Err(Error::Parse(format!(
"Rows.db entry: recovered trie root {trie_root_signed} out of bounds \
(base={base}, delta={root_delta}, file size={})",
rows_db.len()
)));
}
let trie_root = trie_root_signed as usize;
let (block_count_u64, n) = read_unsigned_vint_from_slice(&rows_db[cur..])?;
cur += n;
let block_count = u32::try_from(block_count_u64).map_err(|_| {
Error::Parse(format!(
"Rows.db entry: implausible block count {block_count_u64}"
))
})?;
let partition_deletion = match decode_da_deletion_time(rows_db, cur) {
Ok((deletion, _consumed)) => deletion,
Err(_) => None,
};
Ok(BtiRowIndexHeader {
data_position,
trie_root,
block_count,
partition_deletion,
})
}
pub fn select_row_index_blocks_for_range(
entries: &[(Vec<u8>, BtiRowIndexEntry)],
start: &[u8],
end: &[u8],
) -> Vec<BtiRowIndexEntry> {
if start > end || entries.is_empty() {
return Vec::new();
}
let mut out = Vec::new();
for (i, (sep_i, block)) in entries.iter().enumerate() {
let next_is_greater_than_start = match entries.get(i + 1) {
Some((sep_next, _)) => sep_next.as_slice() > start,
None => true, };
let overlaps = sep_i.as_slice() <= end && next_is_greater_than_start;
if overlaps {
out.push(block.clone());
}
}
out
}
pub fn iterate_rows_for_partition(
rows_db: &[u8],
rows_offset: usize,
) -> BtiResult<(BtiRowIndexHeader, Vec<BtiRowIndexEntryWithKey>)> {
let header = resolve_rows_db_entry(rows_db, rows_offset)?;
let entries = iterate_rows_in_bti_trie(rows_db, header.trie_root)?;
Ok((header, entries))
}
pub fn iterate_rows_in_bti_file<R: Read + Seek>(
reader: &mut R,
) -> BtiResult<Vec<(Vec<u8>, BtiRowIndexEntry)>> {
let file_size = reader.seek(SeekFrom::End(0))?;
if file_size < 8 {
return Ok(Vec::new());
}
let (trie_data, root_offset) = load_bti_trie_via_footer(reader)?;
iterate_rows_in_bti_trie(&trie_data, root_offset)
}
#[cfg(test)]
mod tests {
use super::*;
fn dense16_node(payload_flags: u8, start: u8, deltas: &[u16]) -> Vec<u8> {
let len = deltas.len() as u8;
let mut v = vec![0xB0 | (payload_flags & 0x0F), start, len - 1];
for &d in deltas {
v.extend_from_slice(&d.to_be_bytes());
}
v
}
fn row_leaf_no_marker(pos: u8) -> Vec<u8> {
assert!(pos <= 127, "use a 1-byte unsigned vint position");
vec![0x01, pos] }
fn make_rows_trie_three(
(k1, p1): (u8, u8),
(k2, p2): (u8, u8),
(k3, p3): (u8, u8),
) -> (Vec<u8>, usize) {
let mut trie = Vec::new();
let o1 = trie.len() as u64; trie.extend_from_slice(&row_leaf_no_marker(p1));
let o2 = trie.len() as u64; trie.extend_from_slice(&row_leaf_no_marker(p2));
let o3 = trie.len() as u64; trie.extend_from_slice(&row_leaf_no_marker(p3));
let root = trie.len() as u64; trie.push(0x50); trie.push(0x03); trie.push(k1);
trie.push(k2);
trie.push(k3);
trie.push((root - o1) as u8);
trie.push((root - o2) as u8);
trie.push((root - o3) as u8);
(trie, root as usize)
}
#[test]
fn dfs_rows_yields_byte_order_with_row_payloads() {
let (trie, root) = make_rows_trie_three((0x10, 5), (0x20, 17), (0x30, 99));
let entries = dfs_collect_row_entries(&trie, root).unwrap();
assert_eq!(
entries,
vec![
(
vec![0x10],
BtiRowIndexEntry {
data_offset: 5,
open_marker: None
}
),
(
vec![0x20],
BtiRowIndexEntry {
data_offset: 17,
open_marker: None
}
),
(
vec![0x30],
BtiRowIndexEntry {
data_offset: 99,
open_marker: None
}
),
],
"Rows.db DFS must yield byte-ordered keys with decoded Data.db positions"
);
}
#[test]
fn dfs_dense_emits_offset_zero_child_and_skips_gap() {
let mut trie = Vec::new();
trie.extend_from_slice(&row_leaf_no_marker(5)); trie.extend_from_slice(&row_leaf_no_marker(9)); let root = trie.len() as u64;
let deltas = [root as u16, 0x0000, (root - 2) as u16];
trie.extend(dense16_node(0, 0x10, &deltas));
let entries = dfs_collect_row_entries(&trie, root as usize).unwrap();
assert_eq!(
entries,
vec![
(
vec![0x10],
BtiRowIndexEntry {
data_offset: 5,
open_marker: None
}
),
(
vec![0x12],
BtiRowIndexEntry {
data_offset: 9,
open_marker: None
}
),
],
"DFS must emit the real child at absolute offset 0 and skip the \
no-transition gap (0x11)"
);
}
#[test]
fn decode_row_payload_open_marker_modern() {
let mut data = vec![0x07u8]; data.extend_from_slice(&567890i64.to_be_bytes());
data.extend_from_slice(&1234u32.to_be_bytes());
let entry = decode_bti_row_payload(&data, 0, 0x9).unwrap();
assert_eq!(
entry,
BtiRowIndexEntry {
data_offset: 7,
open_marker: Some((1234, 567890)),
}
);
}
#[test]
fn decode_row_payload_open_marker_live_sentinel() {
let data = vec![0x07u8, 0x80u8];
let entry = decode_bti_row_payload(&data, 0, 0x9).unwrap();
assert_eq!(
entry,
BtiRowIndexEntry {
data_offset: 7,
open_marker: None,
}
);
}
#[test]
fn da_deletion_time_decoder() {
assert_eq!(decode_da_deletion_time(&[0x80], 0).unwrap(), (None, 1));
assert_eq!(
decode_da_deletion_time(&[0x00, 0x80, 0xFF], 1).unwrap(),
(None, 1)
);
let mut buf = Vec::new();
buf.extend_from_slice(&987_654_321_000i64.to_be_bytes()); buf.extend_from_slice(&1_700_000_000u32.to_be_bytes()); let (del, n) = decode_da_deletion_time(&buf, 0).unwrap();
assert_eq!(n, 12);
assert_eq!(del, Some((1_700_000_000i32, 987_654_321_000i64)));
assert_ne!(buf[0], 0x80);
assert!(decode_da_deletion_time(&[0x00, 0x01, 0x02], 0).is_err());
assert!(decode_da_deletion_time(&[0x00], 5).is_err());
}
#[test]
fn read_unsigned_vint_multibyte() {
let (v, n) = read_unsigned_vint_from_slice(&[0x81, 0x2C]).unwrap();
assert_eq!((v, n), (300, 2));
let (v, n) = read_unsigned_vint_from_slice(&[0x7F]).unwrap();
assert_eq!((v, n), (127, 1));
}
#[test]
fn range_filter_subset_and_empty_and_reversed() {
let (trie, root) = make_rows_trie_three((0x10, 5), (0x20, 17), (0x30, 99));
let all = dfs_collect_row_entries(&trie, root).unwrap();
let filter = |lo: &[u8], hi: &[u8]| -> Vec<u64> {
if lo > hi {
return Vec::new();
}
all.iter()
.filter(|(k, _)| k.as_slice() >= lo && k.as_slice() <= hi)
.map(|(_, e)| e.data_offset)
.collect()
};
assert_eq!(filter(&[0x10], &[0x20]), vec![5, 17]); assert_eq!(filter(&[0x20], &[0x30]), vec![17, 99]); assert_eq!(filter(&[0x00], &[0x0F]), Vec::<u64>::new()); assert_eq!(filter(&[0x31], &[0xFF]), Vec::<u64>::new()); assert_eq!(filter(&[0x30], &[0x10]), Vec::<u64>::new()); assert_eq!(filter(&[0x10], &[0x30]), vec![5, 17, 99]); }
#[test]
fn select_blocks_separator_semantics() {
let entries = vec![
(
vec![0x10u8],
BtiRowIndexEntry {
data_offset: 5,
open_marker: None,
},
),
(
vec![0x20u8],
BtiRowIndexEntry {
data_offset: 17,
open_marker: None,
},
),
(
vec![0x30u8],
BtiRowIndexEntry {
data_offset: 99,
open_marker: None,
},
),
];
let offs = |start: &[u8], end: &[u8]| -> Vec<u64> {
select_row_index_blocks_for_range(&entries, start, end)
.into_iter()
.map(|b| b.data_offset)
.collect()
};
assert_eq!(
offs(&[0x18], &[0x18]),
vec![5],
"floor block must be selected"
);
assert_eq!(offs(&[0x18], &[0x28]), vec![5, 17]);
assert_eq!(offs(&[0x20], &[0x2F]), vec![17]);
assert_eq!(offs(&[0x40], &[0x50]), vec![99]);
assert_eq!(offs(&[0x00], &[0x0F]), Vec::<u64>::new());
assert_eq!(offs(&[0x00], &[0xFF]), vec![5, 17, 99]);
assert_eq!(offs(&[0x30], &[0x10]), Vec::<u64>::new());
assert!(select_row_index_blocks_for_range(&[], &[0x00], &[0xFF]).is_empty());
}
#[test]
fn resolve_rows_db_entry_recovers_root_and_metadata() {
let mut buf = vec![0xEEu8; 4]; let rows_offset = buf.len();
buf.extend_from_slice(&4u16.to_be_bytes());
buf.extend_from_slice(&[0x00, 0x00, 0x00, 0x07]);
let base = rows_offset + 4;
buf.push(123);
let root_delta: i64 = 2 - base as i64;
let zig = ((root_delta << 1) ^ (root_delta >> 63)) as u64;
assert!(zig < 128, "test setup expects a 1-byte vint");
buf.push(zig as u8);
buf.push(38);
buf.extend_from_slice(&17i64.to_be_bytes());
buf.extend_from_slice(&9u32.to_be_bytes());
let header = resolve_rows_db_entry(&buf, rows_offset).unwrap();
assert_eq!(header.data_position, 123);
assert_eq!(
header.trie_root, 2,
"trie root = rootΔ + (RowsOffset + keylen)"
);
assert_eq!(header.block_count, 38);
assert_eq!(header.partition_deletion, Some((9, 17)));
assert!(resolve_rows_db_entry(&buf, buf.len() + 10).is_err());
}
#[test]
fn resolve_rows_db_entry_live_partition_deletion() {
let mut buf = vec![0xEEu8; 4];
let rows_offset = buf.len();
buf.extend_from_slice(&4u16.to_be_bytes());
buf.extend_from_slice(&[0x00, 0x00, 0x00, 0x07]);
let base = rows_offset + 4;
buf.push(123); let root_delta: i64 = 2 - base as i64;
let zig = ((root_delta << 1) ^ (root_delta >> 63)) as u64;
buf.push(zig as u8); buf.push(38); buf.push(0x80);
let header = resolve_rows_db_entry(&buf, rows_offset).unwrap();
assert_eq!(header.block_count, 38);
assert_eq!(
header.partition_deletion, None,
"0x80 live sentinel must decode to no partition deletion"
);
}
#[test]
fn read_signed_vint_zigzag() {
let (v, n) = read_signed_vint_from_slice(&[0x13]).unwrap();
assert_eq!((v, n), (-10, 1));
assert_eq!(read_signed_vint_from_slice(&[0x00]).unwrap(), (0, 1));
assert_eq!(read_signed_vint_from_slice(&[0x7E]).unwrap(), (63, 1));
}
#[test]
fn decode_row_payload_sizedints_two_bytes() {
let data = vec![0x40u8, 0x80];
let entry = decode_bti_row_payload(&data, 0, 0x2).unwrap();
assert_eq!(
entry,
BtiRowIndexEntry {
data_offset: 16512,
open_marker: None,
}
);
}
#[test]
fn iterate_rows_in_bti_trie_empty_and_oob_root() {
let err = iterate_rows_in_bti_trie(&[], 0);
assert!(err.is_err(), "empty Rows.db trie must error, not panic");
let (trie, _root) = make_rows_trie_three((0x10, 5), (0x20, 17), (0x30, 99));
let err = iterate_rows_in_bti_trie(&trie, trie.len() + 100);
assert!(err.is_err(), "out-of-bounds root must error, not panic");
let (trie, root) = make_rows_trie_three((0x10, 5), (0x20, 17), (0x30, 99));
let entries = iterate_rows_in_bti_trie(&trie, root).unwrap();
assert_eq!(entries.len(), 3);
}
#[test]
fn range_filter_prefix_relationship_is_sound() {
let mut trie = Vec::new();
trie.extend_from_slice(&row_leaf_no_marker(2));
let k_off = trie.len() as u64; trie.push(0x21); trie.push(0x20); trie.push(k_off as u8); trie.push(0x01);
let root = trie.len() as u64; trie.push(0x20); trie.push(0x10); trie.push((root - k_off) as u8);
let all = iterate_rows_in_bti_trie(&trie, root as usize).unwrap();
assert_eq!(
all,
vec![
(
vec![0x10],
BtiRowIndexEntry {
data_offset: 1,
open_marker: None
}
),
(
vec![0x10, 0x20],
BtiRowIndexEntry {
data_offset: 2,
open_marker: None
}
),
],
"K (internal payload) must sort before its descendant K2"
);
let filter = |lo: &[u8], hi: &[u8]| -> Vec<u64> {
if lo > hi {
return Vec::new();
}
all.iter()
.filter(|(k, _)| k.as_slice() >= lo && k.as_slice() <= hi)
.map(|(_, e)| e.data_offset)
.collect()
};
let k = [0x10u8];
let k2 = [0x10u8, 0x20u8];
assert_eq!(
filter(&k, &k),
vec![1],
"[K..=K] must include K and exclude the longer K2"
);
assert_eq!(
filter(&k, &k2),
vec![1, 2],
"[K..=K2] must include both K and K2"
);
assert_eq!(
filter(&[0x10, 0x00], &[0x10, 0x10]),
Vec::<u64>::new(),
"a range strictly between K and K2 must exclude both"
);
}
}