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 {
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 {
fn drop(&mut self) {
let _ = fs::remove_dir_all(&self.path);
}
}
#[derive(Default)]
struct Collector {
payloads: Vec<Vec<u8>>,
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]
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]
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]
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)]
);
}