use std::{io::Read, path::PathBuf};
use wkv::RangeIndexMigrationReader as Engine;
use super::range_index_chunked_serializer::{ChunkStreamError, MIN_CHUNK_SIZE};
pub const DEFAULT_FILE_READ_BUFFER_SIZE: usize = 1 << 20;
pub struct RangeIndexMigrationReader<R: Read>(Engine<R>);
impl<R: Read> RangeIndexMigrationReader<R> {
pub fn new(
serializer: super::range_index_chunked_serializer::RangeIndexChunkedSerializer,
reader: R,
temp_file_path: Option<PathBuf>,
read_buffer_size: usize,
) -> Result<Self, ChunkStreamError> {
if read_buffer_size == 0 {
return Err(ChunkStreamError::InvalidState(
"readBufferSize must be positive".to_string(),
));
}
Ok(Self(
Engine::new(serializer.0, reader, temp_file_path, read_buffer_size)
.map_err(|e| ChunkStreamError::InvalidState(e.to_string()))?,
))
}
pub fn read_next_chunk_async(
&mut self,
destination: &mut [u8],
) -> Result<usize, ChunkStreamError> {
self.validate_destination(destination.len())?;
self
.0
.read_next_chunk(destination)
.map_err(|e| ChunkStreamError::InvalidState(e.to_string()))
}
pub fn read_next_chunk(&mut self, destination: &mut [u8]) -> Result<usize, ChunkStreamError> {
self.read_next_chunk_async(destination)
}
pub fn validate_destination(&self, length: usize) -> Result<(), ChunkStreamError> {
if length < MIN_CHUNK_SIZE {
return Err(ChunkStreamError::InvalidState(format!(
"destination must be at least {MIN_CHUNK_SIZE} bytes (the trailer size) so the stream can complete, got {length}"
)));
}
Ok(())
}
pub fn supply_file_data_or_throw(
&self,
bytes_read: usize,
file_data_remaining: u64,
) -> Result<(), ChunkStreamError> {
if bytes_read == 0 && file_data_remaining > 0 {
return Err(ChunkStreamError::InvalidState(format!(
"RangeIndex file truncated: {file_data_remaining} bytes remaining"
)));
}
Ok(())
}
#[inline]
pub fn is_complete(&self) -> bool {
self.0.is_complete()
}
#[inline]
pub fn total_file_bytes(&self) -> u64 {
self.0.total_file_bytes()
}
pub fn dispose(&mut self) {
self.0.dispose();
}
}
#[cfg(test)]
mod tests {
use std::{fs, io};
use tempfile::tempdir;
use super::*;
use crate::resp::rangeindex::range_index_chunked_serializer::RangeIndexChunkedSerializer;
fn reader_for<'a>(
file: &'a [u8],
key: &[u8],
stub: &[u8],
temp: Option<PathBuf>,
) -> RangeIndexMigrationReader<&'a [u8]> {
let serializer = RangeIndexChunkedSerializer::new(key, stub, file.len() as u64);
RangeIndexMigrationReader::new(serializer, file, temp, DEFAULT_FILE_READ_BUFFER_SIZE).unwrap()
}
#[test]
fn default_read_buffer_is_one_mib() {
assert_eq!(DEFAULT_FILE_READ_BUFFER_SIZE, 1 << 20);
}
#[test]
fn drives_stream_to_completion_in_destination_sized_chunks() {
let stub = [0x11u8; 35];
let file: Vec<u8> = (0..2000u32).map(|i| i as u8).collect();
let mut r = reader_for(&file, b"idx", &stub, None);
assert_eq!(r.total_file_bytes(), 2000);
let mut out = Vec::new();
let mut buf = vec![0u8; MIN_CHUNK_SIZE + 3];
while !r.is_complete() {
let written = r.read_next_chunk(&mut buf).unwrap();
assert!(written > 0, "incomplete stream must make progress");
out.extend_from_slice(&buf[..written]);
}
let key_len = u32::from_le_bytes([out[0], out[1], out[2], out[3]]) as usize;
assert_eq!(&out[4..4 + key_len], b"idx");
assert_eq!(&out[out.len() - 35..], &stub);
r.dispose();
}
#[test]
fn validate_destination_rejects_undersized_buffer() {
let r = reader_for(b"data", b"k", &[0u8; 35], None);
assert!(r.validate_destination(MIN_CHUNK_SIZE - 1).is_err());
assert!(r.validate_destination(MIN_CHUNK_SIZE).is_ok());
let mut r2 = reader_for(b"data", b"k", &[0u8; 35], None);
let mut small = vec![0u8; MIN_CHUNK_SIZE - 1];
assert!(r2.read_next_chunk(&mut small).is_err());
}
#[test]
fn truncated_file_fails_supply() {
let src: &[u8] = &[7u8; 10];
let serializer = RangeIndexChunkedSerializer::new(b"k", &[0u8; 35], 100);
let mut r =
RangeIndexMigrationReader::new(serializer, src, None, DEFAULT_FILE_READ_BUFFER_SIZE).unwrap();
let mut buf = vec![0u8; 4096];
let err = loop {
match r.read_next_chunk(&mut buf) {
Ok(_) => continue,
Err(e) => break e,
}
};
assert!(err.to_string().contains("truncated"));
assert!(r.supply_file_data_or_throw(0, 90).is_err());
assert!(r.supply_file_data_or_throw(16, 90).is_ok());
assert!(r.supply_file_data_or_throw(0, 0).is_ok());
}
#[test]
fn dispose_deletes_owned_temp_snapshot() {
let dir = tempdir().unwrap();
let temp = dir.path().join("snapshot.bftree");
fs::write(&temp, b"payload").unwrap();
let mut r = reader_for(b"abc", b"k", &[0u8; 35], Some(temp.clone()));
assert!(!r.is_complete());
r.dispose();
assert!(!temp.exists(), "owned temp snapshot must be deleted");
r.dispose();
}
#[test]
fn zero_read_buffer_rejected() {
let serializer = RangeIndexChunkedSerializer::new(b"k", &[0u8; 35], 0);
let err = RangeIndexMigrationReader::new(serializer, io::empty(), None, 0)
.err()
.expect("zero buffer must be rejected");
assert!(err.to_string().contains("positive"));
}
}