use std::collections::VecDeque;
use std::fs;
use std::future::Future;
use std::path::PathBuf;
use std::pin::Pin;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use std::task::{Context, Poll, Wake, Waker};
use std::time::Duration;
use futures_core::Stream;
use pi_append_log::file::{DefaultFileLayout, FileAppendLogBuilder};
use pi_append_log::format::{BlockEncoder, DefaultBlockCodec};
use pi_append_log::{AppendLog, AppendLogBuilder, AppendLogVisitor, AppendOptions, ReadOrder};
use pi_result::{ClassifyErrorKind, ErrorKind};
static NEXT_ID: AtomicU64 = AtomicU64::new(9000);
const TEST_TIMEOUT: Duration = Duration::from_secs(5);
type Chunk = Box<[u8]>;
type ChunkResult = pi_result::Result<Chunk>;
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_stream_{}_{}",
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.0);
}
}
#[derive(Default)]
struct Collector {
blocks: Vec<Vec<u8>>,
}
impl 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)
}
}
struct PendingFirstChunkStream {
entered: Arc<tokio::sync::Notify>,
}
impl Stream for PendingFirstChunkStream {
type Item = ChunkResult;
fn poll_next(self: Pin<&mut Self>, _context: &mut Context<'_>) -> Poll<Option<Self::Item>> {
self.entered.notify_one();
Poll::Pending
}
}
struct ChunkStream {
items: VecDeque<ChunkResult>,
}
impl ChunkStream {
fn new(items: impl IntoIterator<Item = ChunkResult>) -> Self {
Self {
items: items.into_iter().collect(),
}
}
}
impl Stream for ChunkStream {
type Item = ChunkResult;
fn poll_next(mut self: Pin<&mut Self>, _context: &mut Context<'_>) -> Poll<Option<Self::Item>> {
Poll::Ready(self.items.pop_front())
}
}
struct StreamGate {
between_chunks: AtomicBool,
released: AtomicBool,
stream_waker: Mutex<Option<Waker>>,
reached_between_chunks: tokio::sync::Notify,
}
impl StreamGate {
fn release(&self) {
self.released.store(true, Ordering::SeqCst);
if let Some(waker) = self.stream_waker.lock().expect("lock stream waker").take() {
waker.wake();
}
}
}
struct GatedChunkStream {
first: Option<Chunk>,
second: Option<Chunk>,
gate: Arc<StreamGate>,
}
impl Stream for GatedChunkStream {
type Item = ChunkResult;
fn poll_next(mut self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Option<Self::Item>> {
if let Some(first) = self.first.take() {
return Poll::Ready(Some(Ok(first)));
}
if !self.gate.released.load(Ordering::SeqCst) {
if !self.gate.between_chunks.swap(true, Ordering::SeqCst) {
self.gate.reached_between_chunks.notify_one();
}
*self.gate.stream_waker.lock().expect("lock stream waker") =
Some(context.waker().clone());
return Poll::Pending;
}
Poll::Ready(self.second.take().map(Ok))
}
}
struct ObservableWake {
notified: tokio::sync::Notify,
}
impl Wake for ObservableWake {
fn wake(self: Arc<Self>) {
self.notified.notify_one();
}
}
fn builder(path: PathBuf) -> FileAppendLogBuilder<DefaultFileLayout> {
FileAppendLogBuilder::new(path.clone(), DefaultFileLayout::new(path))
}
fn split_block(block: &[u8], first_end: usize, second_end: usize) -> [Chunk; 3] {
[
block[..first_end].to_vec().into_boxed_slice(),
block[first_end..second_end].to_vec().into_boxed_slice(),
block[second_end..].to_vec().into_boxed_slice(),
]
}
#[tokio::test]
async fn append_stream_combines_chunks_into_one_recoverable_block_and_returns_active_size() {
let temp = TempDir::new();
let codec = DefaultBlockCodec::default();
let mut initial_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"one logical block from several chunks")
.expect("encode block");
let chunks = split_block(&block, 7, block.len() - 5);
let active_size = log
.append_stream(
ChunkStream::new(chunks.into_iter().map(Ok)),
AppendOptions::default(),
)
.await
.expect("stream append must succeed");
assert_eq!(active_size, block.len() as u64);
drop(log);
let mut restarted_visitor = Collector::default();
let restarted = builder(temp.0.clone())
.build(&codec, ReadOrder::Forward, &mut restarted_visitor)
.await
.expect("restart after stream append");
assert_eq!(restarted_visitor.blocks, vec![block]);
drop(restarted.storage);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn ordinary_append_cannot_interleave_between_stream_chunks() {
let temp = TempDir::new();
let codec = DefaultBlockCodec::default();
let mut initial_visitor = Collector::default();
let log = Arc::new(
builder(temp.0.clone())
.build(&codec, ReadOrder::Forward, &mut initial_visitor)
.await
.expect("build")
.storage,
);
let streamed_block = codec
.encode(1, 0, b"streamed block")
.expect("encode streamed block");
let ordinary_block = codec
.encode(2, 0, b"ordinary block")
.expect("encode ordinary block");
let split = streamed_block.len() / 2;
let gate = Arc::new(StreamGate {
between_chunks: AtomicBool::new(false),
released: AtomicBool::new(false),
stream_waker: Mutex::new(None),
reached_between_chunks: tokio::sync::Notify::new(),
});
let stream = GatedChunkStream {
first: Some(streamed_block[..split].to_vec().into_boxed_slice()),
second: Some(streamed_block[split..].to_vec().into_boxed_slice()),
gate: Arc::clone(&gate),
};
let stream_log = Arc::clone(&log);
let stream_future = async move {
stream_log
.append_stream(stream, AppendOptions { durable: false })
.await
};
tokio::pin!(stream_future);
let wake = Arc::new(ObservableWake {
notified: tokio::sync::Notify::new(),
});
let waker = Waker::from(Arc::clone(&wake));
let mut context = Context::from_waker(&waker);
loop {
match Future::poll(stream_future.as_mut(), &mut context) {
Poll::Ready(result) => {
gate.release();
assert!(
result.is_ok(),
"append_stream completed before the controlled chunk boundary: {result:?}"
);
panic!("append_stream must remain pending between controlled chunks");
}
Poll::Pending if gate.between_chunks.load(Ordering::SeqCst) => break,
Poll::Pending => {
tokio::time::timeout(TEST_TIMEOUT, wake.notified.notified())
.await
.expect("append_stream stopped waking before polling its stream");
}
}
}
let append_log = Arc::clone(&log);
let ordinary_append_input = ordinary_block.clone();
let append_future = async move {
append_log
.append(ordinary_append_input, AppendOptions { durable: false })
.await
};
tokio::pin!(append_future);
assert!(matches!(
Future::poll(append_future.as_mut(), &mut context),
Poll::Pending
));
gate.release();
tokio::time::timeout(TEST_TIMEOUT, stream_future.as_mut())
.await
.expect("append_stream did not finish after release")
.expect("append_stream failed after release");
tokio::time::timeout(TEST_TIMEOUT, append_future.as_mut())
.await
.expect("ordinary append did not finish after stream")
.expect("ordinary append failed");
drop(stream_future);
drop(append_future);
drop(log);
let mut restarted_visitor = Collector::default();
let restarted = builder(temp.0.clone())
.build(&codec, ReadOrder::Forward, &mut restarted_visitor)
.await
.expect("restart after stream and ordinary append");
assert_eq!(
restarted_visitor.blocks,
vec![streamed_block, ordinary_block]
);
drop(restarted.storage);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn rotate_cannot_interleave_between_stream_chunks() {
let temp = TempDir::new();
let codec = DefaultBlockCodec::default();
let mut initial_visitor = Collector::default();
let log = Arc::new(
builder(temp.0.clone())
.build(&codec, ReadOrder::Forward, &mut initial_visitor)
.await
.expect("build")
.storage,
);
let streamed_block = codec
.encode(1, 0, b"streamed block")
.expect("encode streamed block");
let split = streamed_block.len() / 2;
let gate = Arc::new(StreamGate {
between_chunks: AtomicBool::new(false),
released: AtomicBool::new(false),
stream_waker: Mutex::new(None),
reached_between_chunks: tokio::sync::Notify::new(),
});
let stream = GatedChunkStream {
first: Some(streamed_block[..split].to_vec().into_boxed_slice()),
second: Some(streamed_block[split..].to_vec().into_boxed_slice()),
gate: Arc::clone(&gate),
};
let stream_log = Arc::clone(&log);
let stream_future = async move {
stream_log
.append_stream(stream, AppendOptions { durable: false })
.await
};
tokio::pin!(stream_future);
let wake = Arc::new(ObservableWake {
notified: tokio::sync::Notify::new(),
});
let waker = Waker::from(Arc::clone(&wake));
let mut context = Context::from_waker(&waker);
loop {
match Future::poll(stream_future.as_mut(), &mut context) {
Poll::Ready(result) => {
gate.release();
assert!(
result.is_ok(),
"append_stream completed before the controlled chunk boundary: {result:?}"
);
panic!("append_stream must remain pending between controlled chunks");
}
Poll::Pending if gate.between_chunks.load(Ordering::SeqCst) => break,
Poll::Pending => {
tokio::time::timeout(TEST_TIMEOUT, wake.notified.notified())
.await
.expect("append_stream stopped waking before polling its stream");
}
}
}
let rotate_log = Arc::clone(&log);
let rotate_future = async move { rotate_log.rotate().await };
tokio::pin!(rotate_future);
assert!(matches!(
Future::poll(rotate_future.as_mut(), &mut context),
Poll::Pending
));
gate.release();
tokio::time::timeout(TEST_TIMEOUT, stream_future.as_mut())
.await
.expect("append_stream did not finish after release")
.expect("append_stream failed after release");
tokio::time::timeout(TEST_TIMEOUT, rotate_future.as_mut())
.await
.expect("rotate did not finish after stream")
.expect("rotate failed")
.expect("streamed block must produce a closed structure");
drop(stream_future);
drop(rotate_future);
drop(log);
assert!(temp.0.join("00000001").exists());
assert_eq!(fs::metadata(temp.0.join("00000002")).unwrap().len(), 0);
let active_count = fs::read_dir(&temp.0)
.unwrap()
.filter_map(Result::ok)
.filter(|entry| entry.path().extension().is_none())
.count();
assert_eq!(active_count, 2);
let mut restarted_visitor = Collector::default();
let restarted = builder(temp.0.clone())
.build(&codec, ReadOrder::Forward, &mut restarted_visitor)
.await
.expect("restart after stream and rotate");
assert_eq!(restarted_visitor.blocks, vec![streamed_block]);
assert_eq!(restarted.recovered_closed.len(), 1);
drop(restarted.storage);
}
#[tokio::test]
async fn empty_stream_is_invalid_input_and_leaves_active_unchanged() {
let temp = TempDir::new();
let codec = DefaultBlockCodec::default();
let mut visitor = Collector::default();
let log = builder(temp.0.clone())
.build(&codec, ReadOrder::Forward, &mut visitor)
.await
.expect("build")
.storage;
let active = temp.0.join("00000001");
let before = fs::metadata(&active).unwrap().len();
let error = log
.append_stream(ChunkStream::new([]), AppendOptions::default())
.await
.expect_err("empty stream must fail");
assert_eq!(error.classify_error_kind(), ErrorKind::InvalidInput);
assert_eq!(fs::metadata(active).unwrap().len(), before);
}
#[tokio::test]
async fn all_empty_chunks_are_invalid_input_and_leave_active_unchanged() {
let temp = TempDir::new();
let codec = DefaultBlockCodec::default();
let mut visitor = Collector::default();
let log = builder(temp.0.clone())
.build(&codec, ReadOrder::Forward, &mut visitor)
.await
.expect("build")
.storage;
let active = temp.0.join("00000001");
let before = fs::metadata(&active).unwrap().len();
let empty_chunks = [
Ok(Vec::<u8>::new().into_boxed_slice()),
Ok(Vec::<u8>::new().into_boxed_slice()),
];
let error = log
.append_stream(ChunkStream::new(empty_chunks), AppendOptions::default())
.await
.expect_err("all-empty stream must fail");
assert_eq!(error.classify_error_kind(), ErrorKind::InvalidInput);
assert_eq!(fs::metadata(active).unwrap().len(), before);
}
#[tokio::test]
async fn stream_error_after_prefix_preserves_prefix_and_returns_dependency_error() {
let temp = TempDir::new();
let codec = DefaultBlockCodec::default();
let mut 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"partial stream failure")
.expect("encode block");
let prefix = block[..block.len() / 2].to_vec().into_boxed_slice();
let expected_prefix = prefix.to_vec();
let dependency_error = pi_result::error_stack::Report::new(ErrorKind::Dependency);
let stream = ChunkStream::new([Ok(prefix), Err(dependency_error)]);
let error = log
.append_stream(stream, AppendOptions { durable: false })
.await
.expect_err("stream dependency failure must propagate");
assert_eq!(error.classify_error_kind(), ErrorKind::Dependency);
assert_eq!(fs::read(temp.0.join("00000001")).unwrap(), expected_prefix);
assert!(temp.0.join("00000001").exists());
}
#[tokio::test]
async fn stream_error_after_prefix_invalidates_instance() {
let temp = TempDir::new();
let codec = DefaultBlockCodec::default();
let mut visitor = Collector::default();
let log = builder(temp.0.clone())
.build(&codec, ReadOrder::Forward, &mut visitor)
.await
.expect("build")
.storage;
let partial_block = codec
.encode(1, 0, b"partial stream invalidates instance")
.expect("encode partial block");
let prefix = partial_block[..partial_block.len() / 2]
.to_vec()
.into_boxed_slice();
let expected_prefix = prefix.to_vec();
let dependency_error = pi_result::error_stack::Report::new(ErrorKind::Dependency);
let error = log
.append_stream(
ChunkStream::new([Ok(prefix), Err(dependency_error)]),
AppendOptions { durable: false },
)
.await
.expect_err("stream dependency failure must propagate");
assert_eq!(error.classify_error_kind(), ErrorKind::Dependency);
assert_eq!(fs::read(temp.0.join("00000001")).unwrap(), expected_prefix);
let complete_block = codec
.encode(2, 0, b"must not follow a partial prefix")
.expect("encode complete block");
let followup = log
.append(complete_block, AppendOptions { durable: false })
.await;
let followup_error = followup.expect_err("partial stream failure must invalidate the instance");
assert_eq!(
followup_error.classify_error_kind(),
ErrorKind::InvalidState
);
assert_eq!(fs::read(temp.0.join("00000001")).unwrap(), expected_prefix);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn stream_cancellation_after_prefix_invalidates_instance() {
let temp = TempDir::new();
let codec = DefaultBlockCodec::default();
let mut visitor = Collector::default();
let log = Arc::new(
builder(temp.0.clone())
.build(&codec, ReadOrder::Forward, &mut visitor)
.await
.expect("build")
.storage,
);
let streamed_block = codec
.encode(1, 0, b"cancelled between stream chunks")
.expect("encode streamed block");
let split = streamed_block.len() / 2;
let expected_prefix = streamed_block[..split].to_vec();
let gate = Arc::new(StreamGate {
between_chunks: AtomicBool::new(false),
released: AtomicBool::new(false),
stream_waker: Mutex::new(None),
reached_between_chunks: tokio::sync::Notify::new(),
});
let stream = GatedChunkStream {
first: Some(expected_prefix.clone().into_boxed_slice()),
second: Some(streamed_block[split..].to_vec().into_boxed_slice()),
gate: Arc::clone(&gate),
};
let stream_log = Arc::clone(&log);
let append_task = tokio::spawn(async move {
stream_log
.append_stream(stream, AppendOptions { durable: false })
.await
});
tokio::time::timeout(TEST_TIMEOUT, gate.reached_between_chunks.notified())
.await
.expect("append_stream never reached the controlled chunk boundary");
assert!(gate.between_chunks.load(Ordering::SeqCst));
assert_eq!(fs::read(temp.0.join("00000001")).unwrap(), expected_prefix);
append_task.abort();
let join_error = append_task
.await
.expect_err("aborted append_stream task must report cancellation");
assert!(join_error.is_cancelled());
let complete_block = codec
.encode(2, 0, b"must not follow cancelled stream prefix")
.expect("encode complete block");
let followup = log
.append(complete_block, AppendOptions { durable: false })
.await;
let followup_error =
followup.expect_err("cancelled partial stream must invalidate the instance");
assert_eq!(
followup_error.classify_error_kind(),
ErrorKind::InvalidState
);
assert_eq!(fs::read(temp.0.join("00000001")).unwrap(), expected_prefix);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn stream_cancellation_before_first_chunk_keeps_shared_namespace_available() {
let temp = TempDir::new();
let codec = DefaultBlockCodec::default();
let mut first_visitor = Collector::default();
let first = Arc::new(
builder(temp.0.clone())
.build(&codec, ReadOrder::Forward, &mut first_visitor)
.await
.expect("build first instance")
.storage,
);
let mut second_visitor = Collector::default();
let second = builder(temp.0.clone())
.build(&codec, ReadOrder::Forward, &mut second_visitor)
.await
.expect("build second instance")
.storage;
let entered = Arc::new(tokio::sync::Notify::new());
let stream = PendingFirstChunkStream {
entered: Arc::clone(&entered),
};
let first_log = Arc::clone(&first);
let append_task = tokio::spawn(async move {
first_log
.append_stream(stream, AppendOptions { durable: false })
.await
});
tokio::time::timeout(TEST_TIMEOUT, entered.notified())
.await
.expect("append_stream never polled its first chunk");
assert_eq!(
fs::metadata(temp.0.join("00000001"))
.expect("initial active file")
.len(),
0
);
append_task.abort();
let join_error = append_task
.await
.expect_err("aborted append_stream task must report cancellation");
assert!(join_error.is_cancelled());
let block = codec
.encode(1, 0, b"must remain appendable after pre-chunk cancellation")
.expect("encode followup block");
let active_size = second
.append(block.clone(), AppendOptions { durable: false })
.await
.expect("cancellation before the first chunk must leave the namespace available");
assert!(active_size > 0);
assert_eq!(active_size, block.len() as u64);
assert_eq!(
fs::read(temp.0.join("00000001")).expect("active file after followup append"),
block
);
assert!(temp.0.join("00000001").exists());
}