use crate::error::{Error, Result};
use crate::storage::serialization::vint::{encode_signed, encode_unsigned};
use crate::storage::write_engine::mutation::DecoratedKey;
use std::io::Write;
use std::path::PathBuf;
const INDEX_SINK_BUFFER_BYTES: usize = 512 * 1024;
pub const COLUMN_INDEX_SIZE_BYTES: u64 = 64 * 1024;
pub const INDEX_INFO_WIDTH_BASE: u64 = 64 * 1024;
#[derive(Debug, Clone, Default)]
pub struct PromotedIndexBlock {
pub first_name: Vec<u8>,
pub last_name: Vec<u8>,
pub offset: u64,
pub width: u64,
pub oss50_separator: Option<Vec<u8>>,
}
#[derive(Debug)]
pub struct IndexWriter {
buffer: Vec<u8>,
sink: Option<std::io::BufWriter<std::fs::File>>,
index_path: Option<PathBuf>,
counting: bool,
position: u64,
entry_count: usize,
}
#[derive(Debug, Clone, Copy)]
pub struct IndexEntryInfo {
pub index_offset: u64,
pub entry_size: usize,
}
impl IndexWriter {
pub fn new() -> Self {
Self {
buffer: Vec::new(),
sink: None,
index_path: None,
counting: false,
position: 0,
entry_count: 0,
}
}
pub fn counting() -> Self {
Self {
buffer: Vec::new(),
sink: None,
index_path: None,
counting: true,
position: 0,
entry_count: 0,
}
}
pub fn with_sink(index_path: PathBuf) -> Self {
Self {
buffer: Vec::new(),
sink: None,
index_path: Some(index_path),
counting: false,
position: 0,
entry_count: 0,
}
}
fn ensure_sink(&mut self) -> Result<()> {
if self.sink.is_some() {
return Ok(());
}
if let Some(path) = self.index_path.clone() {
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent)?;
}
let file = std::fs::File::create(&path)?;
self.sink = Some(std::io::BufWriter::with_capacity(
INDEX_SINK_BUFFER_BYTES,
file,
));
}
Ok(())
}
fn flush_entry(&mut self) -> Result<()> {
if self.counting {
self.position += self.buffer.len() as u64;
self.buffer.clear();
return Ok(());
}
if self.index_path.is_none() {
return Ok(());
}
self.ensure_sink()?;
if let Some(sink) = self.sink.as_mut() {
sink.write_all(&self.buffer)?;
}
self.position += self.buffer.len() as u64;
self.buffer.clear();
Ok(())
}
pub fn add_partition(
&mut self,
key: &DecoratedKey,
data_offset: u64,
) -> Result<IndexEntryInfo> {
self.add_partition_with_promoted(key, data_offset, &[])
}
pub fn add_partition_with_promoted(
&mut self,
key: &DecoratedKey,
data_offset: u64,
blocks: &[PromotedIndexBlock],
) -> Result<IndexEntryInfo> {
let index_offset = self.position + self.buffer.len() as u64;
let entry_size = self.write_entry(key, data_offset, blocks)?;
self.flush_entry()?;
self.entry_count += 1;
Ok(IndexEntryInfo {
index_offset,
entry_size,
})
}
fn write_entry(
&mut self,
key: &DecoratedKey,
data_offset: u64,
blocks: &[PromotedIndexBlock],
) -> Result<usize> {
let start_len = self.buffer.len();
let key_len = key.key.len() as u16;
self.buffer.extend_from_slice(&key_len.to_be_bytes());
self.buffer.extend_from_slice(&key.key);
encode_unsigned(data_offset, &mut self.buffer);
if blocks.len() >= 2 {
let promoted_payload = serialize_promoted_index(blocks, key.key.len());
encode_unsigned(promoted_payload.len() as u64, &mut self.buffer);
self.buffer.extend_from_slice(&promoted_payload);
} else {
encode_unsigned(0, &mut self.buffer);
}
let bytes_written = self.buffer.len() - start_len;
Ok(bytes_written)
}
pub fn finish(self) -> Result<Vec<u8>> {
if self.index_path.is_some() {
return Err(Error::InvalidInput(
"IndexWriter::finish() called on a streaming writer; use finish_streaming()"
.to_string(),
));
}
Ok(self.buffer)
}
pub fn finish_streaming(mut self) -> Result<u64> {
if self.index_path.is_none() {
return Err(Error::InvalidInput(
"finish_streaming() called on an in-memory IndexWriter".to_string(),
));
}
self.flush_entry()?;
if let Some(mut sink) = self.sink.take() {
sink.flush()?;
}
Ok(self.position)
}
pub fn entry_count(&self) -> usize {
self.entry_count
}
pub fn buffered_len(&self) -> usize {
self.buffer.len()
}
}
impl Default for IndexWriter {
fn default() -> Self {
Self::new()
}
}
const NB_DELETION_TIME_LIVE_SIZE: usize = 12;
fn serialize_promoted_index(blocks: &[PromotedIndexBlock], raw_key_len: usize) -> Vec<u8> {
debug_assert!(blocks.len() >= 2, "caller must ensure at least 2 blocks");
let header_length: u64 = (2 + raw_key_len + NB_DELETION_TIME_LIVE_SIZE) as u64;
let mut deletion_time_bytes = [0u8; NB_DELETION_TIME_LIVE_SIZE];
deletion_time_bytes[0..4].copy_from_slice(&i32::MAX.to_be_bytes());
deletion_time_bytes[4..12].copy_from_slice(&i64::MIN.to_be_bytes());
let mut index_info_bytes: Vec<u8> = Vec::new();
let mut block_start_offsets: Vec<u32> = Vec::with_capacity(blocks.len());
for block in blocks {
let block_offset_in_info = index_info_bytes.len() as u32;
block_start_offsets.push(block_offset_in_info);
serialize_index_info(&mut index_info_bytes, block);
}
let mut payload: Vec<u8> = Vec::new();
encode_unsigned(header_length, &mut payload);
payload.extend_from_slice(&deletion_time_bytes);
encode_unsigned(blocks.len() as u64, &mut payload);
payload.extend_from_slice(&index_info_bytes);
for &offset in &block_start_offsets {
payload.extend_from_slice(&(offset as i32).to_be_bytes());
}
payload
}
fn serialize_index_info(buf: &mut Vec<u8>, block: &PromotedIndexBlock) {
buf.extend_from_slice(&block.first_name);
buf.extend_from_slice(&block.last_name);
encode_unsigned(block.offset, buf);
let width_delta = (block.width as i64) - (INDEX_INFO_WIDTH_BASE as i64);
encode_signed(width_delta, buf);
buf.push(0x00u8);
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_index_writer_new() {
let writer = IndexWriter::new();
assert_eq!(writer.entry_count(), 0);
}
#[test]
fn test_add_single_partition_int_key() {
let mut writer = IndexWriter::new();
let key = DecoratedKey::new(12345, vec![0x00, 0x00, 0x00, 0x2A]);
let info = writer.add_partition(&key, 0).unwrap();
assert_eq!(writer.entry_count(), 1);
assert_eq!(info.index_offset, 0);
assert_eq!(info.entry_size, 8);
}
#[test]
fn test_add_single_partition_uuid_key() {
let mut writer = IndexWriter::new();
let key = DecoratedKey::new(12345, vec![0xBB; 16]);
let info = writer.add_partition(&key, 0).unwrap();
assert_eq!(writer.entry_count(), 1);
assert_eq!(info.entry_size, 20);
}
#[test]
fn test_raw_key_bytes_written() {
let mut writer = IndexWriter::new();
let pk_bytes = vec![0x00, 0x00, 0x00, 0x2A];
let key = DecoratedKey::new(12345, pk_bytes.clone());
writer.add_partition(&key, 0).unwrap();
let bytes = writer.finish().unwrap();
assert_eq!(&bytes[0..2], &[0x00, 0x04], "Key length should be 4");
assert_eq!(&bytes[2..6], &pk_bytes, "Should be raw key bytes");
assert_eq!(bytes[6], 0x00, "Offset should be 0");
assert_eq!(bytes[7], 0x00, "Promoted size should be 0");
}
#[test]
fn test_uuid_key_raw_bytes() {
let mut writer = IndexWriter::new();
let pk_bytes = vec![
0x55, 0x0e, 0x84, 0x00, 0xe2, 0x9b, 0x41, 0xd4, 0xa7, 0x16, 0x44, 0x66, 0x55, 0x44,
0x00, 0x00,
];
let key = DecoratedKey::new(12345, pk_bytes.clone());
writer.add_partition(&key, 0).unwrap();
let bytes = writer.finish().unwrap();
assert_eq!(&bytes[0..2], &[0x00, 0x10], "Key length should be 16");
assert_eq!(&bytes[2..18], &pk_bytes, "Should be raw UUID bytes");
}
#[test]
fn test_add_multiple_partitions() {
let mut writer = IndexWriter::new();
let key1 = DecoratedKey::new(100, vec![0x00, 0x00, 0x00, 0x01]);
let key2 = DecoratedKey::new(200, vec![0x00, 0x00, 0x00, 0x02]);
let key3 = DecoratedKey::new(300, vec![0x00, 0x00, 0x00, 0x03]);
let info1 = writer.add_partition(&key1, 0).unwrap();
let info2 = writer.add_partition(&key2, 150).unwrap();
let info3 = writer.add_partition(&key3, 300).unwrap();
assert_eq!(writer.entry_count(), 3);
assert_eq!(info1.index_offset, 0);
assert_eq!(info2.index_offset, info1.entry_size as u64);
assert_eq!(
info3.index_offset,
(info1.entry_size + info2.entry_size) as u64
);
}
#[test]
fn test_finish_multiple_entries() {
let mut writer = IndexWriter::new();
let key1 = DecoratedKey::new(100, vec![0x00, 0x00, 0x00, 0x01]);
let key2 = DecoratedKey::new(200, vec![0x00, 0x00, 0x00, 0x02]);
writer.add_partition(&key1, 0).unwrap();
writer.add_partition(&key2, 150).unwrap();
let bytes = writer.finish().unwrap();
assert_eq!(bytes.len(), 17);
assert_eq!(&bytes[0..2], &[0x00, 0x04]);
assert_eq!(&bytes[8..10], &[0x00, 0x04]);
}
#[test]
fn test_position_encoding() {
let mut writer = IndexWriter::new();
let key = DecoratedKey::new(12345, vec![0x00, 0x00, 0x00, 0x2A]);
writer.add_partition(&key, 127).unwrap();
let bytes = writer.finish().unwrap();
assert_eq!(bytes[6], 0x7F);
assert_eq!(bytes[7], 0x00); }
#[test]
fn test_position_encoding_large_offset() {
let mut writer = IndexWriter::new();
let key = DecoratedKey::new(12345, vec![0x00, 0x00, 0x00, 0x2A]);
writer.add_partition(&key, 12381).unwrap();
let bytes = writer.finish().unwrap();
assert_eq!(bytes[6], 0xB0);
assert_eq!(bytes[7], 0x5D);
assert_eq!(bytes[8], 0x00);
assert_eq!(bytes.len(), 9);
}
#[test]
fn test_variable_key_sizes() {
let mut writer = IndexWriter::new();
let key = DecoratedKey::new(100, vec![0x42]);
let info = writer.add_partition(&key, 0).unwrap();
assert_eq!(info.entry_size, 5);
let mut writer = IndexWriter::new();
let key = DecoratedKey::new(100, vec![0x00, 0x00, 0x00, 0x2A]);
let info = writer.add_partition(&key, 0).unwrap();
assert_eq!(info.entry_size, 8);
let mut writer = IndexWriter::new();
let key = DecoratedKey::new(100, vec![0; 8]);
let info = writer.add_partition(&key, 0).unwrap();
assert_eq!(info.entry_size, 12);
let mut writer = IndexWriter::new();
let key = DecoratedKey::new(100, vec![0; 16]);
let info = writer.add_partition(&key, 0).unwrap();
assert_eq!(info.entry_size, 20); }
#[test]
fn test_empty_index() {
let writer = IndexWriter::new();
let bytes = writer.finish().unwrap();
assert_eq!(bytes.len(), 0);
}
#[test]
fn test_token_order_preservation() {
let mut writer = IndexWriter::new();
let key1 = DecoratedKey::new(100, vec![0x01]);
let key2 = DecoratedKey::new(200, vec![0x02]);
let key3 = DecoratedKey::new(300, vec![0x03]);
writer.add_partition(&key1, 0).unwrap();
writer.add_partition(&key2, 100).unwrap();
writer.add_partition(&key3, 200).unwrap();
let bytes = writer.finish().unwrap();
assert_eq!(bytes.len(), 16);
assert_eq!(&bytes[0..2], &[0x00, 0x01]);
assert_eq!(&bytes[5..7], &[0x00, 0x01]);
assert_eq!(&bytes[10..12], &[0x00, 0x01]);
}
#[test]
fn test_vint_encoding_boundaries() {
let key = DecoratedKey::new(12345, vec![0x00, 0x00, 0x00, 0x2A]);
let mut writer = IndexWriter::new();
writer.add_partition(&key, 127).unwrap();
assert_eq!(writer.finish().unwrap().len(), 8);
let mut writer = IndexWriter::new();
writer.add_partition(&key, 128).unwrap();
assert_eq!(writer.finish().unwrap().len(), 9);
let mut writer = IndexWriter::new();
writer.add_partition(&key, 16383).unwrap();
assert_eq!(writer.finish().unwrap().len(), 9);
let mut writer = IndexWriter::new();
writer.add_partition(&key, 16384).unwrap();
assert_eq!(writer.finish().unwrap().len(), 10); }
#[test]
fn test_index_offset_tracking() {
let mut writer = IndexWriter::new();
let key1 = DecoratedKey::new(100, vec![0x01, 0x02, 0x03, 0x04]);
let info1 = writer.add_partition(&key1, 0).unwrap();
let key2 = DecoratedKey::new(200, vec![0x05, 0x06]);
let info2 = writer.add_partition(&key2, 127).unwrap();
let key3 = DecoratedKey::new(300, vec![0x07]);
let info3 = writer.add_partition(&key3, 12381).unwrap();
assert_eq!(info1.index_offset, 0);
assert_eq!(info1.entry_size, 8, "Entry 1: 2 + 4 + 1 + 1 = 8");
assert_eq!(info2.index_offset, 8);
assert_eq!(info2.entry_size, 6, "Entry 2: 2 + 2 + 1 + 1 = 6");
assert_eq!(info3.index_offset, 14);
assert_eq!(info3.entry_size, 6, "Entry 3: 2 + 1 + 2 + 1 = 6");
let bytes = writer.finish().unwrap();
assert_eq!(
bytes.len(),
info1.entry_size + info2.entry_size + info3.entry_size,
"Total size matches sum of entry sizes"
);
}
#[test]
fn test_realistic_scenario() {
let mut writer = IndexWriter::new();
let key1 = DecoratedKey::new(-5000000000, vec![0x00, 0x00, 0x03, 0xE9]);
writer.add_partition(&key1, 0).unwrap();
let key2 = DecoratedKey::new(-2000000000, vec![0x00, 0x00, 0x03, 0xEA]);
writer.add_partition(&key2, 250).unwrap();
let key3 = DecoratedKey::new(3000000000, vec![0x00, 0x00, 0x03, 0xEB]);
writer.add_partition(&key3, 500).unwrap();
assert_eq!(writer.entry_count(), 3);
let bytes = writer.finish().unwrap();
assert_eq!(bytes.len(), 26);
}
#[test]
fn test_single_block_no_promoted_index() {
let mut writer = IndexWriter::new();
let key = DecoratedKey::new(100, vec![0x01, 0x02, 0x03, 0x04]);
let block = PromotedIndexBlock {
first_name: vec![0x00], last_name: vec![0x00],
offset: 0,
width: 70_000,
oss50_separator: None,
};
let info = writer
.add_partition_with_promoted(&key, 0, &[block])
.unwrap();
let bytes = writer.finish().unwrap();
assert_eq!(
bytes[info.entry_size - 1],
0x00,
"promoted size should be 0"
);
assert_eq!(info.entry_size, 8);
}
#[test]
fn test_two_blocks_emits_promoted_index() {
let mut writer = IndexWriter::new();
let key = DecoratedKey::new(100, vec![0x01, 0x02, 0x03, 0x04]);
let ck_prefix = vec![0x00u8];
let block1 = PromotedIndexBlock {
first_name: ck_prefix.clone(),
last_name: ck_prefix.clone(),
offset: 0,
width: 65_536, oss50_separator: None,
};
let block2 = PromotedIndexBlock {
first_name: ck_prefix.clone(),
last_name: ck_prefix.clone(),
offset: 65_536,
width: 65_536,
oss50_separator: None,
};
let info = writer
.add_partition_with_promoted(&key, 0, &[block1, block2])
.unwrap();
let bytes = writer.finish().unwrap();
let promoted_size_offset = 2 + 4 + 1; assert!(
bytes[promoted_size_offset] > 0,
"promoted_index_size must be > 0 for 2 blocks"
);
assert!(
info.entry_size > 8,
"Wide partition entry must be larger than small"
);
}
#[test]
fn test_three_blocks_promoted_index_size_matches_payload() {
let key = DecoratedKey::new(42, vec![0xAA, 0xBB]);
let ck_prefix = vec![0x00u8, 0x61, 0x62];
let make_block = |off: u64, w: u64| PromotedIndexBlock {
first_name: ck_prefix.clone(),
last_name: ck_prefix.clone(),
offset: off,
width: w,
oss50_separator: None,
};
let blocks = vec![
make_block(0, 70_000),
make_block(70_000, 68_000),
make_block(138_000, 65_000),
];
let mut writer = IndexWriter::new();
let info = writer
.add_partition_with_promoted(&key, 1234, &blocks)
.unwrap();
let bytes = writer.finish().unwrap();
assert_eq!(bytes.len(), info.entry_size);
let pos_vint_len = 2; let promoted_size_vint_start = 2 + 2 + pos_vint_len;
let promoted_size = parse_vint_simple(&bytes[promoted_size_vint_start..]);
let (promoted_size_value, promoted_size_vint_bytes) = promoted_size;
assert!(promoted_size_value > 0, "3 blocks → promoted index present");
let payload_start = promoted_size_vint_start + promoted_size_vint_bytes;
let payload_end = payload_start + promoted_size_value as usize;
assert_eq!(
payload_end,
bytes.len(),
"entry size should exactly cover key + vints + promoted payload"
);
let payload = &bytes[payload_start..payload_end];
let (header_len, hl_bytes) = parse_vint_simple(payload);
assert_eq!(
header_len, 16,
"headerLength = 2 (key_len_prefix) + 2 (key bytes) + 12 (NB DeletionTime) = 16"
);
let dt_start = hl_bytes;
let dt_end = dt_start + NB_DELETION_TIME_LIVE_SIZE;
let ldt = i32::from_be_bytes(payload[dt_start..dt_start + 4].try_into().unwrap());
let mfda = i64::from_be_bytes(payload[dt_start + 4..dt_end].try_into().unwrap());
assert_eq!(ldt, i32::MAX, "NB LIVE DeletionTime: ldt = i32::MAX");
assert_eq!(mfda, i64::MIN, "NB LIVE DeletionTime: mfda = i64::MIN");
let (count, _) = parse_vint_simple(&payload[dt_end..]);
assert_eq!(count, 3, "Three IndexInfo blocks");
}
#[test]
fn test_block_boundary_math_threshold_crossing() {
assert_eq!(COLUMN_INDEX_SIZE_BYTES, 64 * 1024);
let block_at_threshold = PromotedIndexBlock {
first_name: vec![0x00],
last_name: vec![0x00],
offset: COLUMN_INDEX_SIZE_BYTES,
width: COLUMN_INDEX_SIZE_BYTES,
oss50_separator: None,
};
let block_below = PromotedIndexBlock {
first_name: vec![0x00],
last_name: vec![0x00],
offset: 0,
width: COLUMN_INDEX_SIZE_BYTES,
oss50_separator: None,
};
let raw_key_len = 4usize;
let blocks = vec![block_below, block_at_threshold];
let payload = serialize_promoted_index(&blocks, raw_key_len);
assert!(!payload.is_empty());
let (header_len, hl_bytes) = parse_vint_simple(&payload);
assert_eq!(
header_len, 18,
"headerLength for a 4-byte key = 2 + 4 + 12 = 18"
);
let dt_skip = NB_DELETION_TIME_LIVE_SIZE;
let (count, _) = parse_vint_simple(&payload[hl_bytes + dt_skip..]);
assert_eq!(count, 2);
}
#[test]
fn test_width_delta_encoding_exact_threshold() {
let block = PromotedIndexBlock {
first_name: vec![0x00],
last_name: vec![0x00],
offset: 0,
width: INDEX_INFO_WIDTH_BASE, oss50_separator: None,
};
let block2 = PromotedIndexBlock {
first_name: vec![0x00],
last_name: vec![0x00],
offset: INDEX_INFO_WIDTH_BASE,
width: INDEX_INFO_WIDTH_BASE,
oss50_separator: None,
};
let mut info_bytes: Vec<u8> = Vec::new();
serialize_index_info(&mut info_bytes, &block);
assert_eq!(info_bytes, vec![0x00, 0x00, 0x00, 0x00, 0x00]);
let payload = serialize_promoted_index(&[block, block2], 4);
assert!(!payload.is_empty());
}
#[test]
fn test_clustering_prefix_preserved_verbatim() {
let first_name = vec![0x00u8, 0x01, 0x02]; let last_name = vec![0x00u8, 0xFF, 0xFE];
let block1 = PromotedIndexBlock {
first_name: first_name.clone(),
last_name: last_name.clone(),
offset: 0,
width: 100_000,
oss50_separator: None,
};
let mut info_bytes: Vec<u8> = Vec::new();
serialize_index_info(&mut info_bytes, &block1);
assert_eq!(&info_bytes[0..3], &first_name[..]);
assert_eq!(&info_bytes[3..6], &last_name[..]);
}
#[test]
fn test_streaming_byte_identical_to_in_memory() {
let partitions: Vec<(i64, Vec<u8>, u64)> = vec![
(100, vec![0x00, 0x00, 0x00, 0x01], 0),
(200, vec![0x00, 0x00, 0x00, 0x02], 256),
(300, vec![0x00, 0x00, 0x00, 0x03], 512),
(400, vec![0x00, 0x00, 0x00, 0x04], 12381), (500, vec![0xBB; 16], 1_000_000_000), ];
let mut mem_writer = IndexWriter::new();
let mut mem_infos = Vec::new();
for (token, key_bytes, offset) in &partitions {
let key = DecoratedKey::new(*token, key_bytes.clone());
let info = mem_writer.add_partition(&key, *offset).unwrap();
mem_infos.push(info);
}
let mem_bytes = mem_writer.finish().unwrap();
let tmp = tempfile::NamedTempFile::new().unwrap();
let path = tmp.path().to_path_buf();
let mut stream_writer = IndexWriter::with_sink(path.clone());
let mut stream_infos = Vec::new();
for (token, key_bytes, offset) in &partitions {
let key = DecoratedKey::new(*token, key_bytes.clone());
let info = stream_writer.add_partition(&key, *offset).unwrap();
stream_infos.push(info);
}
let total_bytes = stream_writer.finish_streaming().unwrap();
let stream_bytes = std::fs::read(&path).unwrap();
assert_eq!(
mem_bytes, stream_bytes,
"Streaming Index.db output must be byte-identical to in-memory output"
);
assert_eq!(
total_bytes as usize,
stream_bytes.len(),
"finish_streaming() must return exact byte count"
);
assert_eq!(
mem_infos.len(),
stream_infos.len(),
"Entry count must match"
);
for (i, (m, s)) in mem_infos.iter().zip(stream_infos.iter()).enumerate() {
assert_eq!(
m.index_offset, s.index_offset,
"Entry {i}: index_offset must match between modes"
);
assert_eq!(
m.entry_size, s.entry_size,
"Entry {i}: entry_size must match between modes"
);
}
}
#[test]
fn test_streaming_byte_identical_with_promoted_index() {
let ck_prefix = vec![0x00u8, 0x61, 0x62]; let make_block = |off: u64, w: u64| PromotedIndexBlock {
first_name: ck_prefix.clone(),
last_name: ck_prefix.clone(),
offset: off,
width: w,
oss50_separator: None,
};
let blocks = vec![
make_block(0, 70_000),
make_block(70_000, 68_000),
make_block(138_000, 65_000),
];
let key = DecoratedKey::new(42, vec![0xAA, 0xBB]);
let mut mem_writer = IndexWriter::new();
let mem_info = mem_writer
.add_partition_with_promoted(&key, 1234, &blocks)
.unwrap();
let mem_bytes = mem_writer.finish().unwrap();
let tmp = tempfile::NamedTempFile::new().unwrap();
let path = tmp.path().to_path_buf();
let mut stream_writer = IndexWriter::with_sink(path.clone());
let stream_info = stream_writer
.add_partition_with_promoted(&key, 1234, &blocks)
.unwrap();
let total_bytes = stream_writer.finish_streaming().unwrap();
let stream_bytes = std::fs::read(&path).unwrap();
assert_eq!(
mem_bytes, stream_bytes,
"Streaming wide-partition Index.db must be byte-identical to in-memory"
);
assert_eq!(total_bytes as usize, stream_bytes.len());
assert_eq!(mem_info.index_offset, stream_info.index_offset);
assert_eq!(mem_info.entry_size, stream_info.entry_size);
}
#[test]
fn test_streaming_entry_count() {
let tmp = tempfile::NamedTempFile::new().unwrap();
let mut writer = IndexWriter::with_sink(tmp.path().to_path_buf());
assert_eq!(writer.entry_count(), 0);
for i in 0..5u64 {
let key = DecoratedKey::new(i as i64 * 100, vec![i as u8]);
writer.add_partition(&key, i * 64).unwrap();
}
assert_eq!(writer.entry_count(), 5);
let _ = writer.finish_streaming().unwrap();
}
#[test]
fn test_streaming_index_offset_tracking() {
let tmp = tempfile::NamedTempFile::new().unwrap();
let mut writer = IndexWriter::with_sink(tmp.path().to_path_buf());
let key1 = DecoratedKey::new(100, vec![0x01, 0x02, 0x03, 0x04]);
let info1 = writer.add_partition(&key1, 0).unwrap();
let key2 = DecoratedKey::new(200, vec![0x05, 0x06]);
let info2 = writer.add_partition(&key2, 127).unwrap();
let key3 = DecoratedKey::new(300, vec![0x07]);
let info3 = writer.add_partition(&key3, 12381).unwrap();
assert_eq!(info1.index_offset, 0);
assert_eq!(info1.entry_size, 8);
assert_eq!(info2.index_offset, 8);
assert_eq!(info2.entry_size, 6);
assert_eq!(info3.index_offset, 14);
assert_eq!(info3.entry_size, 6);
let total = writer.finish_streaming().unwrap();
assert_eq!(total, 20, "Total bytes: 8 + 6 + 6 = 20");
}
#[test]
fn test_finish_on_streaming_writer_errors() {
let tmp = tempfile::NamedTempFile::new().unwrap();
let writer = IndexWriter::with_sink(tmp.path().to_path_buf());
let result = writer.finish();
assert!(result.is_err(), "finish() on streaming writer must fail");
}
#[test]
fn test_finish_streaming_on_in_memory_writer_errors() {
let writer = IndexWriter::new();
let result = writer.finish_streaming();
assert!(
result.is_err(),
"finish_streaming() on in-memory writer must fail"
);
}
#[test]
fn test_counting_mode_does_not_retain_entry_bytes() {
let mut counting = IndexWriter::counting();
let mut in_memory = IndexWriter::new();
for i in 0..1000u64 {
let key = DecoratedKey::new(i as i64, vec![(i & 0xff) as u8; 16]);
counting.add_partition(&key, i * 64).unwrap();
in_memory.add_partition(&key, i * 64).unwrap();
}
assert_eq!(
counting.buffered_len(),
0,
"counting mode must not accumulate any Index.db entry bytes"
);
assert!(
in_memory.buffered_len() > 1000,
"in-memory mode retains all entries (sanity check the contrast)"
);
assert_eq!(counting.entry_count(), 1000);
}
#[test]
fn test_counting_mode_offsets_match_in_memory() {
let partitions: Vec<(i64, Vec<u8>, u64)> = vec![
(100, vec![0x01, 0x02, 0x03, 0x04], 0),
(200, vec![0x05, 0x06], 127),
(300, vec![0x07], 12381),
(400, vec![0xBB; 16], 1_000_000_000),
];
let mut counting = IndexWriter::counting();
let mut in_memory = IndexWriter::new();
for (token, key_bytes, offset) in &partitions {
let key = DecoratedKey::new(*token, key_bytes.clone());
let c = counting.add_partition(&key, *offset).unwrap();
let m = in_memory.add_partition(&key, *offset).unwrap();
assert_eq!(
c.index_offset, m.index_offset,
"counting index_offset must match in-memory"
);
assert_eq!(
c.entry_size, m.entry_size,
"counting entry_size must match in-memory"
);
assert_eq!(counting.buffered_len(), 0);
}
}
#[test]
fn test_counting_mode_with_promoted_index_offsets_match() {
let ck_prefix = vec![0x00u8, 0x61, 0x62];
let make_block = |off: u64, w: u64| PromotedIndexBlock {
first_name: ck_prefix.clone(),
last_name: ck_prefix.clone(),
offset: off,
width: w,
oss50_separator: None,
};
let blocks = vec![
make_block(0, 70_000),
make_block(70_000, 68_000),
make_block(138_000, 65_000),
];
let key = DecoratedKey::new(42, vec![0xAA, 0xBB]);
let mut counting = IndexWriter::counting();
let c = counting
.add_partition_with_promoted(&key, 1234, &blocks)
.unwrap();
let mut in_memory = IndexWriter::new();
let m = in_memory
.add_partition_with_promoted(&key, 1234, &blocks)
.unwrap();
assert_eq!(c.index_offset, m.index_offset);
assert_eq!(c.entry_size, m.entry_size);
assert_eq!(
counting.buffered_len(),
0,
"counting mode must not retain the promoted payload"
);
}
fn parse_vint_simple(buf: &[u8]) -> (u64, usize) {
if buf.is_empty() {
return (0, 0);
}
let first = buf[0];
if first <= 0x7F {
return (first as u64, 1);
}
let leading = first.leading_ones() as usize;
let extra_bits = 8 - leading - 1; let mask = (1u8 << extra_bits).wrapping_sub(1);
let mut value = (first & mask) as u64;
for b in buf.iter().take(leading + 1).skip(1) {
value = (value << 8) | *b as u64;
}
(value, 1 + leading)
}
}