use std::{
fs,
path::{Path, PathBuf},
};
use wkv::RangeIndexChunkedDeserializer as Engine;
use super::range_index_chunked_serializer::ChunkStreamError;
pub struct RangeIndexChunkedDeserializer(pub(super) Engine);
impl RangeIndexChunkedDeserializer {
pub fn new(temp_path: impl Into<PathBuf>) -> Result<Self, ChunkStreamError> {
Ok(Self(Engine::new(temp_path).map_err(|e| {
ChunkStreamError::InvalidState(e.to_string())
})?))
}
pub fn process_chunk(&mut self, data: &[u8]) -> Result<bool, ChunkStreamError> {
self
.0
.process_chunk(data)
.map_err(|e| ChunkStreamError::InvalidState(e.to_string()))
}
pub fn parse_trailer(&self) -> Result<usize, ChunkStreamError> {
if !self.0.is_complete() {
return Err(ChunkStreamError::InvalidState(
"trailer not parsed: stream is not complete".to_string(),
));
}
Ok(self.0.stub().len())
}
pub fn write_file_bytes(&self) -> u64 {
fs::metadata(self.temp_path()).map(|m| m.len()).unwrap_or(0)
}
pub fn close_stream(&self) -> bool {
self.temp_path().exists()
}
#[inline]
pub fn is_complete(&self) -> bool {
self.0.is_complete()
}
#[inline]
pub fn has_error(&self) -> bool {
self.0.has_error()
}
#[inline]
pub fn take_error(&mut self) -> Option<String> {
self.0.take_error().map(|e| e.to_string())
}
#[inline]
pub fn key(&self) -> &[u8] {
self.0.key()
}
#[inline]
pub fn stub(&self) -> &[u8] {
self.0.stub()
}
#[inline]
pub fn temp_path(&self) -> &Path {
self.0.temp_path()
}
pub fn dispose(&mut self) {
self.0.dispose();
}
}
#[cfg(test)]
mod tests {
use std::io::Write;
use tempfile::tempdir;
use super::*;
fn frame_stream(key: &[u8], stub: &[u8], file: &[u8], chunk_size: usize) -> Vec<Vec<u8>> {
let mut serializer = wkv::RangeIndexChunkedSerializer::new(key, stub, file.len() as u64);
let mut chunks = Vec::new();
let mut supplied = 0usize;
let mut dest = vec![0u8; chunk_size];
while !serializer.is_complete() {
if serializer.needs_file_data() && !file.is_empty() {
let n = (file.len() - supplied).min(5);
serializer.supply_file_data(&file[supplied..supplied + n]);
supplied += n;
}
dest.fill(0);
let written = serializer.move_next(&mut dest).unwrap();
assert!(written > 0);
chunks.push(dest[..written].to_vec());
}
chunks
}
#[test]
fn reassembles_multi_chunk_stream() {
let dir = tempdir().unwrap();
let temp = dir.path().join("reassembly.bftree");
let key = b"idx-key";
let stub = [0xABu8; 35];
let file: Vec<u8> = (0..1000u32).map(|i| i as u8).collect();
let chunks = frame_stream(key, &stub, &file, 64);
assert!(chunks.len() > 3, "expected genuinely multi-chunk stream");
let mut d = RangeIndexChunkedDeserializer::new(&temp).unwrap();
for chunk in &chunks {
assert!(d.process_chunk(chunk).unwrap(), "chunk rejected");
}
assert!(d.is_complete());
assert!(!d.has_error());
assert_eq!(d.key(), &key[..]);
assert_eq!(d.stub(), &stub);
assert_eq!(d.parse_trailer().unwrap(), 35);
assert!(d.close_stream());
assert_eq!(fs::read(&temp).unwrap(), file);
d.dispose();
assert!(!temp.exists(), "dispose removes temp file");
}
#[test]
fn file_bytes_written_tracks_progress() {
let dir = tempdir().unwrap();
let temp = dir.path().join("progress.bftree");
let stub = [1u8; 35];
let file = vec![9u8; 300];
let chunks = frame_stream(b"k", &stub, &file, 47);
let mut d = RangeIndexChunkedDeserializer::new(&temp).unwrap();
assert!(!d.close_stream());
assert_eq!(d.write_file_bytes(), 0);
let mut last = 0u64;
for chunk in chunks.iter().take(chunks.len() - 1) {
d.process_chunk(chunk).unwrap();
let written = d.write_file_bytes();
assert!(written >= last);
last = written;
}
d.process_chunk(chunks.last().unwrap()).unwrap();
assert!(d.is_complete());
assert_eq!(d.write_file_bytes(), 300);
assert!(d.close_stream());
d.dispose();
}
#[test]
fn empty_chunks_are_noops() {
let dir = tempdir().unwrap();
let mut d = RangeIndexChunkedDeserializer::new(dir.path().join("noop.bftree")).unwrap();
assert!(d.process_chunk(&[]).unwrap());
assert!(!d.is_complete());
assert!(!d.has_error());
d.dispose();
}
#[test]
fn corrupt_checksum_is_detected() {
let dir = tempdir().unwrap();
let mut chunks = frame_stream(b"k", &[2u8; 35], &[3u8; 64], 47);
let last = chunks.last_mut().unwrap();
last[0] ^= 0xFF;
let mut d = RangeIndexChunkedDeserializer::new(dir.path().join("bad.bftree")).unwrap();
for chunk in chunks.iter().take(chunks.len() - 1) {
assert!(d.process_chunk(chunk).unwrap());
}
assert!(!d.process_chunk(&chunks[chunks.len() - 1]).unwrap());
assert!(d.has_error());
assert!(d.take_error().is_some());
assert!(!d.is_complete());
d.dispose();
}
#[test]
fn feeding_after_terminal_state_is_rejected() {
let dir = tempdir().unwrap();
let chunks = frame_stream(b"k", &[5u8; 35], &[6u8; 16], 47);
let mut d = RangeIndexChunkedDeserializer::new(dir.path().join("done.bftree")).unwrap();
for chunk in &chunks {
d.process_chunk(chunk).unwrap();
}
assert!(d.is_complete());
assert!(!d.process_chunk(b"trailing").unwrap());
assert!(d.is_complete());
d.dispose();
assert!(!d.process_chunk(&chunks[0]).unwrap());
}
#[test]
fn std_file_write_replay_guard() {
let dir = tempdir().unwrap();
let p = dir.path().join("w.bftree");
{
let mut f = fs::File::create(&p).unwrap();
f.write_all(&[1, 2, 3]).unwrap();
}
assert_eq!(fs::read(&p).unwrap(), vec![1, 2, 3]);
}
}