pi_append_log 0.3.0

Storage-agnostic append-only block log traits, codec, layout, and file backend
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);
}