use std::fs;
use std::path::PathBuf;
use std::sync::atomic::{AtomicU64, Ordering};
use pi_append_log::file::{DefaultFileLayout, FileAppendLogBuilder};
use pi_append_log::format::{BlockEncoder, DefaultBlockCodec};
use pi_append_log::{AppendLog, AppendLogBuilder, AppendLogVisitor, AppendOptions, ReadOrder};
static NEXT_ID: AtomicU64 = AtomicU64::new(5000);
struct TempDir(PathBuf);
impl TempDir {
fn new() -> Self {
let id = NEXT_ID.fetch_add(1, Ordering::Relaxed);
let path = std::env::temp_dir().join(format!(
"pi_append_log_detachable_append_{}_{}",
std::process::id(),
id
));
fs::create_dir(&path).expect("create temporary directory");
Self(path)
}
}
impl Drop for TempDir {
fn drop(&mut self) {
let _ = fs::remove_dir_all(&self.0);
}
}
#[derive(Default)]
struct Collector {
blocks: Vec<Vec<u8>>,
}
impl AppendLogVisitor for Collector {
fn visit(
&mut self,
block: &[u8],
_context: pi_append_log::BlockVisitContext,
) -> pi_result::Result<bool> {
self.blocks.push(block.to_vec());
Ok(false)
}
}
fn builder(path: PathBuf) -> FileAppendLogBuilder<DefaultFileLayout> {
FileAppendLogBuilder::new(path.clone(), DefaultFileLayout::new(path))
}
#[tokio::test]
async fn append_accepts_boxed_detachable_buffer_and_recovers_the_complete_block() {
let temp = TempDir::new();
let codec = DefaultBlockCodec::default();
let mut initial_visitor = Collector::default();
let log = builder(temp.0.clone())
.build(&codec, ReadOrder::Forward, &mut initial_visitor)
.await
.expect("initial build")
.storage;
let encoded = codec
.encode(1, 0, b"boxed detachable buffer")
.expect("encode block");
let expected = encoded.clone();
let boxed_block: Box<[u8]> = encoded.into_boxed_slice();
let active_size = log
.append(boxed_block, AppendOptions::default())
.await
.expect("boxed detachable buffer append must succeed");
assert_eq!(active_size, expected.len() as u64);
drop(log);
let mut restarted_visitor = Collector::default();
let restarted = builder(temp.0.clone())
.build(&codec, ReadOrder::Forward, &mut restarted_visitor)
.await
.expect("restart after boxed append");
assert_eq!(restarted_visitor.blocks, vec![expected]);
drop(restarted.storage);
}