pi_append_log 0.1.0

Storage-agnostic append-only block log traits, codec, layout, and file backend
Documentation
use std::fs;
use std::path::PathBuf;
use std::sync::Arc;
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, AppendOptions, ReadOrder};

static NEXT_ID: AtomicU64 = AtomicU64::new(1000);

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_runtime_{}_{}",
            std::process::id(),
            id
        ));
        fs::create_dir(&path).expect("create temp directory");
        Self(path)
    }
}

impl Drop for TempDir {
    // 测试结束后删除临时文件。
    fn drop(&mut self) {
        let _ = fs::remove_dir_all(&self.0);
    }
}

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

#[tokio::test]
// 验证 append 返回活动文件大小,外部可据此判断阈值且日志不会自动轮转。
async fn append_reports_active_size_for_external_rotation_policy_and_durable_data_survives_restart()
{
    let temp = TempDir::new();
    let codec = DefaultBlockCodec::default();
    let mut visitor = pi_append_log_test_visitor::Collector::default();
    let result = builder(temp.0.clone())
        .build(&codec, ReadOrder::Forward, &mut visitor)
        .await
        .expect("build");
    let block = codec.encode(1, 0, b"durable").unwrap();
    let log = result.storage;
    let first_size = log
        .append(block.clone(), AppendOptions::default())
        .await
        .unwrap();
    assert_eq!(first_size, block.len() as u64);
    assert!(first_size < 70);
    let second_size = log.append(block, AppendOptions::default()).await.unwrap();
    assert_eq!(second_size, first_size * 2);
    assert!(second_size >= 70);
    assert!(temp.0.join("00000001.active").exists());
    drop(log);

    let mut restarted_visitor = pi_append_log_test_visitor::Collector::default();
    let restarted = builder(temp.0.clone())
        .build(&codec, ReadOrder::Forward, &mut restarted_visitor)
        .await
        .expect("restart");
    assert_eq!(restarted_visitor.blocks.len(), 2);
    drop(restarted);
}

#[tokio::test]
// 验证 rotate 返回的 Closed 只归档对应结构,重复 archive 幂等成功,不能影响其他日志目录。
async fn explicit_rotation_and_archive_target_the_returned_closed_handle() {
    let first_temp = TempDir::new();
    let second_temp = TempDir::new();
    let codec = DefaultBlockCodec::default();
    let mut visitor = pi_append_log_test_visitor::Collector::default();
    let first = builder(first_temp.0.clone())
        .build(&codec, ReadOrder::Forward, &mut visitor)
        .await
        .unwrap()
        .storage;
    let mut other_visitor = pi_append_log_test_visitor::Collector::default();
    let other = builder(second_temp.0.clone())
        .build(&codec, ReadOrder::Forward, &mut other_visitor)
        .await
        .unwrap()
        .storage;
    first
        .append(
            codec.encode(1, 0, b"one").unwrap(),
            AppendOptions::default(),
        )
        .await
        .unwrap();
    let closed = first.rotate().await.unwrap().unwrap();
    assert!(first.archive(closed.clone()).await.is_ok());
    assert!(first.archive(closed).await.is_ok());
    assert!(!second_temp.0.join("00000001.archive").exists());
    let _ = other;
}

#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
// 验证多个 Tokio 任务并发 append 时,两个完整 envelope 都可恢复,说明 block 没有发生字节交错。
async fn concurrent_appends_produce_complete_blocks_without_interleaving() {
    let temp = TempDir::new();
    let codec = DefaultBlockCodec::default();
    let mut visitor = pi_append_log_test_visitor::Collector::default();
    let log = builder(temp.0.clone())
        .build(&codec, ReadOrder::Forward, &mut visitor)
        .await
        .unwrap()
        .storage;
    let shared = Arc::new(log);
    let left = Arc::clone(&shared);
    let right = Arc::clone(&shared);
    let left_block = codec.encode(1, 0, &vec![1; 1024]).unwrap();
    let right_block = codec.encode(2, 0, &vec![2; 1024]).unwrap();
    let first_size = left_block.len() as u64;
    let final_size = first_size + right_block.len() as u64;
    let (left_result, right_result) = tokio::join!(
        left.append(left_block, AppendOptions::default()),
        right.append(right_block, AppendOptions::default()),
    );
    let mut sizes = [left_result.unwrap(), right_result.unwrap()];
    sizes.sort_unstable();
    assert_eq!(sizes, [first_size, final_size]);
    drop(shared);
    let mut restarted_visitor = pi_append_log_test_visitor::Collector::default();
    let restarted = builder(temp.0.clone())
        .build(&codec, ReadOrder::Forward, &mut restarted_visitor)
        .await
        .unwrap();
    assert_eq!(restarted_visitor.blocks.len(), 2);
    drop(restarted);
}

mod pi_append_log_test_visitor {
    // 运行期测试用 Visitor:复制收到的借用 block,仅用于在回调返回后断言恢复数量。
    #[derive(Default)]
    pub struct Collector {
        // 保存 Visitor 回调次数对应的完整 block 副本。
        pub blocks: Vec<Vec<u8>>,
    }

    impl pi_append_log::AppendLogVisitor for Collector {
        fn visit(
            &mut self,
            block: &[u8],
            _context: pi_append_log::BlockVisitContext,
        ) -> std::io::Result<bool> {
            self.blocks.push(block.to_vec());
            Ok(false)
        }
    }
}