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]
async fn append_reports_accumulated_active_size_without_automatic_rotation() {
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
.expect("build")
.storage;
let block = codec.encode(1, 0, b"active size").unwrap();
let first_size = log
.append(block.clone(), AppendOptions::default())
.await
.expect("first append");
assert_eq!(first_size, block.len() as u64);
let second_size = log
.append(block, AppendOptions::default())
.await
.expect("second append");
assert_eq!(second_size, first_size * 2);
assert!(temp.0.join("00000001").exists());
}
#[tokio::test]
async fn durable_append_is_recovered_after_restart() {
let temp = TempDir::new();
let codec = DefaultBlockCodec::default();
let mut initial_visitor = pi_append_log_test_visitor::Collector::default();
let log = builder(temp.0.clone())
.build(&codec, ReadOrder::Forward, &mut initial_visitor)
.await
.expect("build")
.storage;
let block = codec.encode(1, 0, b"durable").unwrap();
log.append(block.clone(), AppendOptions::default())
.await
.expect("durable append");
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, vec![block]);
drop(restarted.storage);
}
#[tokio::test]
async fn returned_closed_handle_archives_exact_structure_and_rejects_other_namespace() {
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
.expect("append before rotate");
let closed = first.rotate().await.unwrap().expect("non-empty closed");
assert!(other.archive(closed.clone()).await.is_err());
first
.archive(closed)
.await
.expect("archive returned handle");
assert!(first_temp.0.join("00000001.archive").exists());
assert!(!second_temp.0.join("00000001.archive").exists());
}
#[cfg(windows)]
#[tokio::test]
async fn windows_root_alias_can_archive_rotated_closed_handle() {
let temp = TempDir::new();
let codec = DefaultBlockCodec::default();
let mut first_visitor = pi_append_log_test_visitor::Collector::default();
let first = builder(temp.0.clone())
.build(&codec, ReadOrder::Forward, &mut first_visitor)
.await
.expect("build first log")
.storage;
let alias = PathBuf::from(temp.0.to_string_lossy().to_ascii_uppercase());
let mut second_visitor = pi_append_log_test_visitor::Collector::default();
let second = builder(alias.clone())
.build(&codec, ReadOrder::Forward, &mut second_visitor)
.await
.expect("build aliased log")
.storage;
first
.append(
codec.encode(1, 0, b"windows alias").unwrap(),
AppendOptions::default(),
)
.await
.expect("append before rotate");
let closed = first
.rotate()
.await
.expect("rotate")
.expect("non-empty closed");
second
.archive(closed)
.await
.expect("archive through root alias");
assert!(temp.0.join("00000001.archive").exists());
}
#[tokio::test]
async fn repeated_archive_of_same_closed_handle_is_idempotent() {
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
.expect("build")
.storage;
log.append(
codec.encode(1, 0, b"one").unwrap(),
AppendOptions::default(),
)
.await
.expect("append before rotate");
let closed = log.rotate().await.unwrap().expect("non-empty closed");
log.archive(closed.clone()).await.expect("first archive");
log.archive(closed).await.expect("repeated archive");
assert!(temp.0.join("00000001.archive").exists());
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
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 {
#[derive(Default)]
pub struct Collector {
pub blocks: Vec<Vec<u8>>,
}
impl pi_append_log::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)
}
}
}