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_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]
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)]
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,
) -> std::io::Result<bool> {
self.blocks.push(block.to_vec());
Ok(false)
}
}
}