pi_append_log 0.3.0

Storage-agnostic append-only block log traits, codec, layout, and file backend
use std::fs::{self, OpenOptions};
use std::io::Write;
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(0);

struct TempDir {
    // 测试专用临时目录;Drop 时递归删除,避免测试文件污染系统临时目录。
    path: 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_file_test_{}_{}",
            std::process::id(),
            id
        ));
        fs::create_dir(&path).expect("create temporary directory");
        Self { path }
    }
}

impl Drop for TempDir {
    // 测试结束后清理目录;Windows 下由测试显式关闭文件句柄后再执行清理。
    fn drop(&mut self) {
        let _ = fs::remove_dir_all(&self.path);
    }
}

#[derive(Default)]
struct Collector {
    // 保存 Visitor 收到的完整 envelope,便于验证读取数量和原始字节顺序。
    payloads: Vec<Vec<u8>>,
    // 记录每个 block 后是否跨越物理结构边界。
    boundaries: Vec<(bool, bool)>,
}

impl AppendLogVisitor for Collector {
    fn visit(
        &mut self,
        block: &[u8],
        context: pi_append_log::BlockVisitContext,
    ) -> pi_result::Result<bool> {
        self.payloads.push(block.to_vec());
        self.boundaries
            .push((context.is_first_in_structure, context.is_last_in_structure));
        Ok(false)
    }
}

fn builder(root: PathBuf) -> FileAppendLogBuilder<DefaultFileLayout> {
    FileAppendLogBuilder::new(root.clone(), DefaultFileLayout::new(root))
}

#[tokio::test]
// 验证:轮转后生成 Closed;active 尾部半写可在重启时截断;未归档 Closed 被返回;恢复后仍可追加并归档。
async fn file_log_recovers_tail_and_returns_unarchived_closed_handles() {
    let temp = TempDir::new();
    let codec = DefaultBlockCodec::default();
    let mut visitor = Collector::default();
    let result = builder(temp.path.clone())
        .build(&codec, ReadOrder::Forward, &mut visitor)
        .await
        .expect("initial build");
    let log = result.storage;
    let first = codec.encode(1, 0, b"first").expect("encode first");
    let second = codec.encode(2, 0, b"second").expect("encode second");
    let first_size = log
        .append(first.clone(), AppendOptions::default())
        .await
        .expect("append first");
    assert_eq!(first_size, first.len() as u64);
    let second_size = log
        .append(second.clone(), AppendOptions::default())
        .await
        .expect("append second");
    assert_eq!(second_size, (first.len() + second.len()) as u64);
    let closed = log
        .rotate()
        .await
        .expect("rotate")
        .expect("non-empty closed");
    drop(log);

    assert!(
        temp.path.join("00000001").exists(),
        "rotate must not rename the old segment"
    );
    assert!(
        temp.path.join("00000002").exists(),
        "rotate must create the next segment"
    );
    let active = temp.path.join("00000002");
    let mut file = OpenOptions::new()
        .append(true)
        .open(&active)
        .expect("open active");
    file.write_all(&second[..second.len() - 3])
        .expect("write torn tail");
    drop(file);

    let mut recovered_visitor = Collector::default();
    let recovered = builder(temp.path.clone())
        .build(&codec, ReadOrder::Forward, &mut recovered_visitor)
        .await
        .expect("recover active tail");
    assert_eq!(recovered.recovered_closed.len(), 1);
    assert_eq!(recovered_visitor.payloads.len(), 2);
    assert_eq!(
        recovered_visitor.boundaries,
        vec![(true, false), (false, true)]
    );
    assert!(
        recovered
            .storage
            .append(
                codec.encode(3, 0, b"third").unwrap(),
                AppendOptions::default()
            )
            .await
            .is_ok()
    );

    recovered
        .storage
        .archive(closed)
        .await
        .expect("archive closed");
    assert!(temp.path.join("00000001.archive").exists());
    assert!(
        !temp.path.join("00000001").exists(),
        "archiving removes the segment path"
    );
}

#[tokio::test]
async fn recovery_allows_structure_id_gaps() {
    let temp = TempDir::new();
    let codec = DefaultBlockCodec::default();
    fs::write(
        temp.path.join("00000001"),
        codec
            .encode(1, 0, b"closed-one")
            .expect("encode closed segment"),
    )
    .expect("write closed segment");
    fs::write(
        temp.path.join("00000003"),
        codec
            .encode(2, 0, b"active-three")
            .expect("encode active segment"),
    )
    .expect("write active segment");

    let mut visitor = Collector::default();
    let recovered = builder(temp.path.clone())
        .build(&codec, ReadOrder::Forward, &mut visitor)
        .await
        .expect("build must accept non-contiguous structure ids");

    assert_eq!(recovered.recovered_closed.len(), 1);
    assert!(temp.path.join("00000001").exists());
    assert!(temp.path.join("00000003").exists());
    assert!(!temp.path.join("00000002").exists());
    recovered
        .storage
        .append(
            codec
                .encode(3, 0, b"append-to-three")
                .expect("encode append"),
            AppendOptions::default(),
        )
        .await
        .expect("append must target maximum structure id");
    assert!(!temp.path.join("00000002").exists());
}

#[tokio::test]
// 验证空 active 执行 rotate 返回 None,不创建无意义的空结构。
async fn empty_rotation_returns_none_and_archive_is_idempotent() {
    let temp = TempDir::new();
    let codec = DefaultBlockCodec::default();
    let mut visitor = Collector::default();
    let result = builder(temp.path.clone())
        .build(&codec, ReadOrder::Forward, &mut visitor)
        .await
        .expect("build");
    assert!(
        result
            .storage
            .rotate()
            .await
            .expect("empty rotate")
            .is_none()
    );
}

#[tokio::test]
// 验证 Backward 读取虽然反转 block 访问顺序,但 is_first/is_last 仍描述物理结构中的真实位置。
async fn backward_read_preserves_physical_structure_boundaries() {
    let temp = TempDir::new();
    let codec = DefaultBlockCodec::default();
    let mut initial_visitor = Collector::default();
    let result = builder(temp.path.clone())
        .build(&codec, ReadOrder::Forward, &mut initial_visitor)
        .await
        .expect("build");
    result
        .storage
        .append(
            codec.encode(1, 0, b"first").unwrap(),
            AppendOptions::default(),
        )
        .await
        .unwrap();
    result
        .storage
        .append(
            codec.encode(2, 0, b"second").unwrap(),
            AppendOptions::default(),
        )
        .await
        .unwrap();
    result.storage.rotate().await.unwrap();
    drop(result.storage);

    let mut backward_visitor = Collector::default();
    builder(temp.path.clone())
        .build(&codec, ReadOrder::Backward, &mut backward_visitor)
        .await
        .expect("backward build");
    assert_eq!(
        backward_visitor.boundaries,
        vec![(false, true), (true, false)]
    );
}