use std::io;
use wkv::RangeIndexChunkedSerializer as Engine;
#[derive(Debug, thiserror::Error)]
pub enum ChunkStreamError {
#[error("RangeIndex chunk stream: {0}")]
InvalidState(String),
#[error(transparent)]
Io(#[from] io::Error),
}
pub const MIN_CHUNK_SIZE: usize = 8 + 4 + wkv::RANGE_INDEX_STUB_SIZE;
pub struct RangeIndexChunkedSerializer(pub(super) Engine);
impl RangeIndexChunkedSerializer {
pub fn new(key: &[u8], stub: &[u8], total_file_bytes: u64) -> Self {
Self(Engine::new(key, stub, total_file_bytes))
}
#[inline]
pub fn supply_file_data(&mut self, data: &[u8]) {
self.0.supply_file_data(data);
}
pub fn move_next(&mut self, destination: &mut [u8]) -> Result<usize, ChunkStreamError> {
if self.0.is_complete() {
return Err(ChunkStreamError::InvalidState(
"Serializer has already completed".to_string(),
));
}
self
.0
.move_next(destination)
.map_err(|e| ChunkStreamError::InvalidState(e.to_string()))
}
pub fn write_trailer(&self, target: &mut [u8], stub: &[u8]) -> Result<usize, ChunkStreamError> {
let trailer_len = 8 + 4 + stub.len();
if target.len() < trailer_len {
return Err(ChunkStreamError::InvalidState(format!(
"trailer needs {trailer_len} bytes, got {}",
target.len()
)));
}
Ok(trailer_len)
}
#[inline]
pub fn is_complete(&self) -> bool {
self.0.is_complete()
}
#[inline]
pub fn needs_file_data(&self) -> bool {
self.0.needs_file_data()
}
#[inline]
pub fn file_data_remaining(&self) -> u64 {
self.0.file_data_remaining()
}
#[inline]
pub fn total_file_bytes(&self) -> u64 {
self.0.total_file_bytes()
}
}
#[cfg(test)]
mod tests {
use super::*;
fn engine_serializer(key: &[u8], stub: &[u8], total: u64) -> Engine {
Engine::new(key, stub, total)
}
#[test]
fn min_chunk_size_matches_csharp_formula() {
assert_eq!(MIN_CHUNK_SIZE, 47);
assert_eq!(MIN_CHUNK_SIZE, 8 + 4 + wkv::RANGE_INDEX_STUB_SIZE);
}
#[test]
fn move_next_frames_key_header_first() {
let mut s = RangeIndexChunkedSerializer::new(b"key", &[0u8; 35], 0);
assert!(!s.is_complete());
let mut dest = [0u8; 64];
let n = s.move_next(&mut dest).unwrap();
assert_eq!(u32::from_le_bytes([dest[0], dest[1], dest[2], dest[3]]), 3);
assert!(n >= 4 + 3);
}
#[test]
fn move_next_zero_progress_on_undersized_destination() {
let mut s = RangeIndexChunkedSerializer::new(b"key", &[0u8; 35], 16);
let mut dest = [0u8; 2];
assert_eq!(s.move_next(&mut dest).unwrap(), 0);
}
#[test]
fn move_next_after_complete_is_contract_violation() {
let mut s = RangeIndexChunkedSerializer::new(b"", &[0u8; 35], 0);
let mut dest = [0u8; MIN_CHUNK_SIZE];
loop {
let n = s.move_next(&mut dest).unwrap();
if s.is_complete() {
break;
}
assert!(n > 0);
}
let err = s.move_next(&mut dest).unwrap_err();
assert!(err.to_string().contains("already completed"));
}
#[test]
fn supply_file_data_feeds_file_phase() {
let mut s = RangeIndexChunkedSerializer::new(b"k", &[0u8; 35], 10);
assert!(!s.needs_file_data()); let mut dest = vec![0u8; MIN_CHUNK_SIZE];
s.move_next(&mut dest).unwrap();
if !s.is_complete() && s.needs_file_data() {
s.supply_file_data(&[7u8; 10]);
assert_eq!(s.file_data_remaining(), 10);
}
assert_eq!(s.total_file_bytes(), 10);
let _ = engine_serializer(b"k", &[0u8; 35], 10);
}
#[test]
fn write_trailer_validates_capacity_without_touching_state() {
let s = RangeIndexChunkedSerializer::new(b"k", &[0u8; 35], 0);
let mut small = [0u8; 8 + 4 + 34];
assert!(s.write_trailer(&mut small, &[0u8; 35]).is_err());
let mut ok = [0u8; MIN_CHUNK_SIZE];
assert_eq!(
s.write_trailer(&mut ok, &[0u8; 35]).unwrap(),
MIN_CHUNK_SIZE
);
assert!(!s.is_complete());
}
}