use std::convert::TryInto;
use std::marker::PhantomData;
use std::path::{Path, PathBuf};
use std::collections::VecDeque;
use std::io::{Error, Result, ErrorKind};
use std::sync::{Arc,
atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering}};
use futures::future::{FutureExt, BoxFuture};
use async_lock::Mutex;
use bytes::BufMut;
use pi_guid::Guid;
use pi_hash::XHashMap;
use pi_async_rt::{lock::spin_lock::SpinLock,
rt::{AsyncRuntime, multi_thread::MultiTaskRuntime}};
use pi_async_transaction::AsyncCommitLog;
use crate::log_store::log_file::{PairLoader, LogMethod, LogFile, log_file_name_to_usize, PairLoaderExt};
const DEFAULT_COMMIT_LOG_FILE_SIZE: usize = 16 * 1024 * 1024 * 1024;
const DEFAULT_LOAD_BUFFER_LEN: u64 = 8192;
const DEFAULT_COMMIT_LOG_BLOCK_SIZE: usize = 8192;
const DEFAULT_DELAY_COMMIT_TIMEOUT: usize = 1;
const DEFAULT_COMMIT_LOG_FILE_MAX_LIMIT: u64 = 32 * 1024 * 1024;
const DEFAULT_COMMIT_LOG_COLLECT_INTERVAL: usize = 10 * 1000;
pub trait CommitLoggerExt: AsyncCommitLog {
fn start_replay_ext<B, F>(&self, callback: Arc<F>)
-> BoxFuture<'static, Result<(usize, usize)>>
where B: BufMut + AsRef<[u8]> + From<Vec<u8>> + Send + Sized + 'static,
F: Fn(Option<(Self::Cid, LogMethod, u64, B)>) -> Result<()> + Send + Sync + 'static;
}
pub struct CommitLoggerBuilder {
rt: MultiTaskRuntime<()>, path: PathBuf, log_block_limit: usize, delay_timeout: usize, log_file_limit: u64, collect_interval: usize, }
unsafe impl Send for CommitLoggerBuilder {}
unsafe impl Sync for CommitLoggerBuilder {}
impl CommitLoggerBuilder {
pub fn new<P: AsRef<Path>>(rt: MultiTaskRuntime<()>,
dir: P) -> Self {
CommitLoggerBuilder {
rt,
path: dir.as_ref().to_path_buf(),
log_block_limit: DEFAULT_COMMIT_LOG_BLOCK_SIZE,
delay_timeout: DEFAULT_DELAY_COMMIT_TIMEOUT,
log_file_limit: DEFAULT_COMMIT_LOG_FILE_MAX_LIMIT,
collect_interval: DEFAULT_COMMIT_LOG_COLLECT_INTERVAL,
}
}
pub fn log_block_limit(mut self, mut limit: usize) -> Self {
if limit < 2048 || limit > 32 * 1024 * 1024 {
limit = DEFAULT_COMMIT_LOG_BLOCK_SIZE
}
self.log_block_limit = limit;
self
}
pub fn delay_timeout(mut self, mut timeout: usize) -> Self {
if timeout < 1 || timeout > 10 {
timeout = DEFAULT_DELAY_COMMIT_TIMEOUT
}
self.delay_timeout = timeout;
self
}
pub fn log_file_limit(mut self, mut limit: u64) -> Self {
if limit < 2 * 1024 * 1024 || limit > 2 * 1024 * 1024 * 1024 {
limit = DEFAULT_COMMIT_LOG_FILE_MAX_LIMIT;
}
self.log_file_limit = limit;
self
}
pub fn collect_interval(mut self, mut interval: usize) -> Self {
if interval < 5 * 1000 || interval > 5 * 60 * 1000 {
interval = DEFAULT_COMMIT_LOG_COLLECT_INTERVAL;
}
self.collect_interval = interval;
self
}
pub async fn build(mut self) -> Result<CommitLogger> {
let file = LogFile::open(self.rt.clone(),
self.path.clone(),
self.log_block_limit,
DEFAULT_COMMIT_LOG_FILE_SIZE, None).await?;
let rt = self.rt;
let delay_timeout = self.delay_timeout;
let log_file_limit = self.log_file_limit;
let writed_size = AtomicU64::new(0); let check_point_counter = Arc::new(AtomicU64::new(0)); let check_point_path = Arc::new(file.writable_path().unwrap()); let writable = SpinLock::new((check_point_counter, check_point_path)); let only_reads = SpinLock::new(VecDeque::new());
let check_points = Mutex::new(XHashMap::default());
let is_replaying = AtomicBool::new(false); let replay_only_reads = SpinLock::new(VecDeque::new());
let replay_confirm_buf = SpinLock::new(VecDeque::new());
let commit_log_count = AtomicUsize::new(0);
let confirm_commited_count = AtomicUsize::new(0);
let inner = InnerCommitLogger {
rt: rt.clone(),
file,
delay_timeout,
log_file_limit,
writed_size,
writable,
only_reads,
check_points,
is_replaying,
replay_only_reads,
replay_confirm_buf,
commit_log_count,
confirm_commited_count,
};
let commit_logger = CommitLogger(Arc::new(inner));
let commit_logger_copy = commit_logger.clone();
let timeout = self.collect_interval;
let _ = rt.spawn(async move {
loop {
collect_commit_logger(&commit_logger_copy, timeout).await;
}
});
Ok(commit_logger)
}
}
#[derive(Clone)]
pub struct CommitLogger(Arc<InnerCommitLogger>);
unsafe impl Send for CommitLogger {}
unsafe impl Sync for CommitLogger {}
impl AsyncCommitLog for CommitLogger {
type C = usize;
type Cid = Guid;
fn append<B>(&self, commit_uid: Self::Cid, log: B) -> BoxFuture<'static, Result<Self::C>>
where B: BufMut + AsRef<[u8]> + Send + Sized + 'static {
let logger = self.clone();
async move {
if log.as_ref().len() == 0 {
return Ok(0);
}
let mut check_pointes_locked = logger.0.check_points.lock().await;
let log_handle = logger.0.file.append(LogMethod::PlainAppend,
commit_uid.0.to_le_bytes().as_ref(),
log.as_ref());
logger.0.writed_size.fetch_add(log.as_ref().len() as u64 + 16, Ordering::Relaxed);
logger.0.commit_log_count.fetch_add(1, Ordering::Relaxed);
let (counter, path) = &*logger.0.writable.lock();
counter.fetch_add(1, Ordering::AcqRel); check_pointes_locked.insert(commit_uid, (counter.clone(), path.clone()));
Ok(log_handle)
}.boxed()
}
fn flush(&self, log_handle: Self::C) -> BoxFuture<'static, Result<()>> {
let mut logger = self.clone();
async move {
logger.0.file.delay_commit_without_auto_split(log_handle,
false,
logger.0.delay_timeout).await
}.boxed()
}
fn confirm(&self, commit_uid: Self::Cid) -> BoxFuture<'static, Result<()>> {
if self.0.is_replaying.load(Ordering::Relaxed) {
return self.confirm_replay(commit_uid);
}
let logger = self.clone();
async move {
let mut check_pointes_locked = logger.0.check_points.lock().await;
if logger.0.writed_size.load(Ordering::Relaxed) >= logger.0.log_file_limit {
let _ = new_check_point(&logger, false).await;
}
if let Some((counter, check_point_path)) = check_pointes_locked.remove(&commit_uid) {
logger.0.confirm_commited_count.fetch_add(1, Ordering::Relaxed);
if counter.fetch_sub(1, Ordering::AcqRel) == 1 {
if check_point_path.as_ref() == logger.0.writable.lock().1.as_ref() {
let _ = new_check_point(&logger, true).await;
}
let mut swap = VecDeque::new();
{
let only_reads = &mut *logger.0.only_reads.lock();
for (path, is_finish_confirm) in only_reads.iter_mut() {
if check_point_path.as_ref() == path {
*is_finish_confirm = true; } else {
match path.metadata() {
Err(e) => {
warn!("Confirm commited transaction failed, path: {:?}, reason: {:?}", path, e);
},
Ok(meta) => {
if meta.len() == 0 {
*is_finish_confirm = true;
}
}
}
}
}
let mut prev = true; while let Some((path, is_finish_confirm)) = only_reads.pop_front() {
if prev && is_finish_confirm {
let _ = logger.0.file.readable_to_back(path).await?;
} else if !prev {
swap.push_back((path, is_finish_confirm));
} else {
swap.push_back((path, is_finish_confirm));
prev = false; }
}
}
*logger.0.only_reads.lock() = swap; }
}
Ok(())
}.boxed()
}
fn start_replay<B, F>(&self, mut callback: Arc<F>) -> BoxFuture<'static, Result<(usize, usize)>>
where B: BufMut + AsRef<[u8]> + From<Vec<u8>> + Send + Sized + 'static,
F: Fn(Self::Cid, B) -> Result<()> + Send + Sync + 'static {
self.0.is_replaying.store(true, Ordering::SeqCst); let commit_logger = self.clone();
async move {
if let Some(writable_path) = commit_logger.0.file.writable_path() {
match writable_path.metadata() {
Err(e) => {
return Err(Error::new(ErrorKind::Other, format!("Replay commit log failed, path: {:?}, reason: {:?}", writable_path, e)));
},
Ok(meta) => {
if meta.len() == 0 && commit_logger.0.file.readable_amount() == 0 {
return Ok((0, 0));
}
}
}
}
if let Err(e) = commit_logger.0.file.split().await {
return Err(Error::new(ErrorKind::Other, format!("Replay commit log failed, reason: {:?}", e)));
}
let mut invalid_only_read_paths = Vec::new(); let mut only_read_paths = commit_logger.0.file.all_readable_path();
for only_read_path in only_read_paths {
match only_read_path.metadata() {
Err(e) => {
return Err(Error::new(ErrorKind::Other, format!("Replay commit log failed, path: {:?}, reason: {:?}", only_read_path, e)));
},
Ok(meta) => {
if meta.len() == 0 {
invalid_only_read_paths.push(only_read_path);
continue;
}
}
}
commit_logger.0.replay_only_reads.lock().push_back(only_read_path);
}
if let Some(path) = commit_logger.0.replay_only_reads.lock().pop_front() {
*commit_logger.0.writable.lock() = (Arc::new(AtomicU64::new(0)), Arc::new(path));
}
let mut loader = CommitLoggerLoader {
logger: commit_logger.clone(),
buf: Vec::new(),
log_file: None,
callback,
result: Ok((0, 0)),
marker: PhantomData,
};
if let Err(e) = commit_logger.0.file.load_before(&mut loader,
None,
DEFAULT_LOAD_BUFFER_LEN,
true).await {
return Err(e);
}
for invalid_only_read_path in invalid_only_read_paths {
if let Err(e) = commit_logger
.0
.file.readable_to_back(invalid_only_read_path.clone())
.await {
return Err(Error::new(ErrorKind::Other, format!("Replay commit log failed, path: {:?}, reason: {:?}", invalid_only_read_path, e)));
}
}
loader.result()
}.boxed()
}
fn append_replay<B>(&self, commit_uid: Self::Cid, _log: B) -> BoxFuture<'static, Result<Self::C>>
where B: BufMut + AsRef<[u8]> + Send + Sized + 'static {
let logger = self.clone();
async move {
let mut check_pointes_locked = logger.0.check_points.lock().await;
let (counter, path) = &*logger.0.writable.lock();
counter.fetch_add(1, Ordering::AcqRel); check_pointes_locked.insert(commit_uid, (counter.clone(), path.clone()));
logger.0.commit_log_count.fetch_add(1, Ordering::Relaxed);
Ok(0)
}.boxed()
}
fn flush_replay(&self, _log_handle: Self::C) -> BoxFuture<'static, Result<()>> {
async move {
Ok(())
}.boxed()
}
fn confirm_replay(&self, commit_uid: Self::Cid) -> BoxFuture<'static, Result<()>> {
let logger = self.clone();
async move {
logger.0.replay_confirm_buf.lock().push_back(commit_uid);
Ok(())
}.boxed()
}
fn finish_replay(&self) -> BoxFuture<'static, Result<()>> {
let logger = self.clone();
async move {
logger.0.is_replaying.store(false, Ordering::SeqCst);
let replay_confirms = &mut *logger.0.replay_confirm_buf.lock();
while let Some(commit_uid) = replay_confirms.pop_front() {
let _ = logger.confirm(commit_uid).await?;
}
Ok(())
}.boxed()
}
fn check_point_of(&self, commit_uid: Self::Cid) -> BoxFuture<'static, Option<usize>> {
let logger = self.clone();
async move {
let check_point_path = if let Some((_counter, check_point_path)) = logger.0.check_points.lock().await.get(&commit_uid) {
check_point_path.as_ref().clone()
} else {
return None;
};
if let Some(file_name) = check_point_path.file_name() {
if let Some(file_name_str) = file_name.to_str() {
return log_file_name_to_usize(file_name_str);
}
}
None
}.boxed()
}
fn current_check_point(&self) -> BoxFuture<'static, usize> {
let logger = self.clone();
async move {
logger
.0
.file
.current_log_index()
}.boxed()
}
fn append_check_point(&self) -> BoxFuture<'static, Result<usize>> {
let logger = self.clone();
async move {
let _check_pointes_locked = logger.0.check_points.lock().await;
new_check_point(&logger, false).await
}.boxed()
}
fn waiting_confirm_count(&self) -> BoxFuture<'static, usize> {
let logger = self.clone();
async move {
logger
.0
.check_points
.lock()
.await
.len()
}.boxed()
}
fn append_total_count(&self) -> usize {
self
.0
.commit_log_count
.load(Ordering::Relaxed)
}
fn confirm_total_count(&self) -> usize {
self
.0
.confirm_commited_count
.load(Ordering::Relaxed)
}
}
async fn new_check_point(logger: &CommitLogger,
is_finish_confirm: bool) -> Result<usize> {
logger.0.file.commit_pending_block().await?;
let log_index = logger.0.file.split().await?;
let check_point_counter = Arc::new(AtomicU64::new(0)); let check_point_path = Arc::new(logger.0.file.writable_path().unwrap()); *logger.0.writable.lock() = (check_point_counter, check_point_path);
let only_read_path = logger.0.file.last_readable_path();
logger.0.only_reads.lock().push_back((only_read_path, is_finish_confirm));
logger.0.writed_size.store(0, Ordering::Relaxed);
Ok(log_index)
}
async fn collect_commit_logger(logger: &CommitLogger, timeout: usize) {
logger.0.rt.timeout(timeout).await;
if logger.0.is_replaying.load(Ordering::Relaxed) {
return;
}
let check_pointes_locked = logger.0.check_points.lock().await;
if logger.0.writed_size.load(Ordering::Relaxed) >= logger.0.log_file_limit {
new_check_point(&logger, false).await;
}
drop(check_pointes_locked); }
impl CommitLoggerExt for CommitLogger {
fn start_replay_ext<B, F>(&self, mut callback: Arc<F>)
-> BoxFuture<'static, Result<(usize, usize)>>
where B: BufMut + AsRef<[u8]> + From<Vec<u8>> + Send + Sized + 'static,
F: Fn(Option<(Self::Cid, LogMethod, u64, B)>) -> Result<()> + Send + Sync + 'static
{
self.0.is_replaying.store(true, Ordering::SeqCst); let commit_logger = self.clone();
async move {
if let Some(writable_path) = commit_logger.0.file.writable_path() {
match writable_path.metadata() {
Err(e) => {
return Err(Error::new(ErrorKind::Other,
format!("Replay commit log failed, path: {:?}, reason: {:?}",
writable_path,
e)));
},
Ok(meta) => {
if meta.len() == 0 && commit_logger.0.file.readable_amount() == 0 {
return Ok((0, 0));
}
}
}
}
if let Err(e) = commit_logger.0.file.split().await {
return Err(Error::new(ErrorKind::Other,
format!("Replay commit log failed, reason: {:?}",
e)));
}
let mut invalid_only_read_paths = Vec::new(); let mut only_read_paths = commit_logger.0.file.all_readable_path();
for only_read_path in only_read_paths {
match only_read_path.metadata() {
Err(e) => {
return Err(Error::new(ErrorKind::Other,
format!("Replay commit log failed, path: {:?}, reason: {:?}",
only_read_path,
e)));
},
Ok(meta) => {
if meta.len() == 0 {
invalid_only_read_paths.push(only_read_path);
continue;
}
}
}
commit_logger.0.replay_only_reads.lock().push_back(only_read_path);
}
if let Some(path) = commit_logger.0.replay_only_reads.lock().pop_front() {
*commit_logger.0.writable.lock() = (Arc::new(AtomicU64::new(0)), Arc::new(path));
}
let mut loader = CommitLoggerLoaderExt {
logger: commit_logger.clone(),
buf: Vec::new(),
log_file: None,
callback,
result: Ok((0, 0)),
marker: PhantomData,
};
if let Err(e) = commit_logger.0.file.load_before_with_payload_time(&mut loader,
None,
DEFAULT_LOAD_BUFFER_LEN,
true).await {
return Err(e);
}
for invalid_only_read_path in invalid_only_read_paths {
if let Err(e) = commit_logger
.0
.file.readable_to_back(invalid_only_read_path.clone())
.await {
return Err(Error::new(ErrorKind::Other,
format!("Replay commit log failed, path: {:?}, reason: {:?}",
invalid_only_read_path,
e)));
}
}
loader.result()
}.boxed()
}
}
struct InnerCommitLogger {
rt: MultiTaskRuntime<()>, file: LogFile, delay_timeout: usize, log_file_limit: u64, writed_size: AtomicU64, writable: SpinLock<(Arc<AtomicU64>, Arc<PathBuf>)>, only_reads: SpinLock<VecDeque<(PathBuf, bool)>>, check_points: Mutex<XHashMap<Guid, (Arc<AtomicU64>, Arc<PathBuf>)>>, is_replaying: AtomicBool, replay_only_reads: SpinLock<VecDeque<PathBuf>>, replay_confirm_buf: SpinLock<VecDeque<Guid>>, commit_log_count: AtomicUsize, confirm_commited_count: AtomicUsize, }
struct CommitLoggerLoader<
B: BufMut + AsRef<[u8]> + From<Vec<u8>> + Send + Sized + 'static,
F: Fn(Guid, B) -> Result<()> + Send + 'static,
> {
logger: CommitLogger, buf: Vec<(Guid, Vec<u8>)>, log_file: Option<PathBuf>, callback: Arc<F>, result: Result<(usize, usize)>, marker: PhantomData<B>,
}
impl<
B: BufMut + AsRef<[u8]> + From<Vec<u8>> + Send + Sized + 'static,
F: Fn(Guid, B) -> Result<()> + Send + 'static,
> PairLoader for CommitLoggerLoader<B, F> {
fn is_require(&self, _log_file: Option<&PathBuf>, _key: &Vec<u8>) -> bool {
true
}
fn load(&mut self,
log_file: Option<&PathBuf>,
_method: LogMethod,
key: Vec<u8>,
value: Option<Vec<u8>>) {
if self.result.is_err() {
return;
}
if let Some(log_file) = log_file {
if self.log_file.is_none() {
self.log_file = Some(log_file.clone());
}
if self.log_file.as_ref().unwrap() != log_file {
while let Some((commit_uid, log)) = self.buf.pop() {
if let Err(e) = (self.callback)(commit_uid.clone(), B::from(log)) {
self.result = Err(Error::new(ErrorKind::Other, format!("Replay commit log failed, commit_uid: {:?}, reason: {:?}", commit_uid, e)));
}
}
self.log_file = Some(log_file.clone());
next_check_point(&self.logger);
}
let uid = u128::from_le_bytes(key.try_into().unwrap());
let commit_uid = Guid(uid);
if let Some(log) = value {
if let Ok((log_count, bytes_count)) = self.result {
self.result = Ok((log_count + 1, bytes_count + 16 + log.len()));
}
self.buf.push((commit_uid, log));
}
}
}
}
fn next_check_point(logger: &CommitLogger) {
{
let (_, last_writable_path) = &*logger.0.writable.lock();
let only_read_path = last_writable_path.as_ref().clone();
logger.0.only_reads.lock().push_back((only_read_path, false));
}
if let Some(path) = logger.0.replay_only_reads.lock().pop_front() {
let check_point_counter = Arc::new(AtomicU64::new(0)); let check_point_path = Arc::new(path); *logger.0.writable.lock() = (check_point_counter, check_point_path);
} else {
let check_point_counter = Arc::new(AtomicU64::new(0)); let check_point_path = Arc::new(logger.0.file.writable_path().unwrap()); *logger.0.writable.lock() = (check_point_counter, check_point_path);
}
}
impl<
B: BufMut + AsRef<[u8]> + From<Vec<u8>> + Send + Sized + 'static,
F: Fn(Guid, B) -> Result<()> + Send + 'static,
> CommitLoggerLoader<B, F> {
pub fn result(mut self) -> Result<(usize, usize)> {
if self.buf.len() > 0 {
while let Some((commit_uid, log)) = self.buf.pop() {
if let Err(e) = (self.callback)(commit_uid.clone(), B::from(log)) {
self.result = Err(Error::new(ErrorKind::Other, format!("Replay commit log failed, commit_uid: {:?}, reason: {:?}", commit_uid, e)));
}
}
next_check_point(&self.logger);
}
self.result
}
}
#[cfg(test)]
mod checkpoint_rotation_tests {
use std::{
fs,
path::PathBuf,
sync::atomic::{AtomicU64 as TestAtomicU64, AtomicUsize as TestAtomicUsize},
time::{Duration, SystemTime, UNIX_EPOCH},
};
use crossbeam_channel::bounded;
use pi_async_rt::rt::{
multi_thread::{MultiTaskRuntime, MultiTaskRuntimeBuilder},
startup_global_time_loop,
AsyncRuntime,
};
use super::*;
const PHYSICAL_FILE_LIMIT: usize = 1024 * 1024;
const CURRENT_BLOCK_LIMIT: usize = 2 * 1024 * 1024;
const TEST_TIMEOUT: Duration = Duration::from_secs(30);
#[test]
fn test_commit_logger_checkpoint_crosses_logfile_limit_once() {
let _time_loop = startup_global_time_loop(1);
let rt = MultiTaskRuntimeBuilder::default()
.init_worker_size(1)
.build();
let root = unique_test_root();
fs::create_dir_all(&root)
.expect("creating checkpoint size-boundary test root must succeed");
let (sender, receiver) = bounded(1);
let test_rt = rt.clone();
let test_root = root.clone();
rt.spawn(async move {
let result = async {
verify_size_boundary(test_rt.clone(), test_root.join("commit-logger")).await?;
verify_public_delay_commit_boundary(test_rt, test_root.join("public-log-file")).await
}
.await;
let _ = sender.send(result);
})
.expect("spawning checkpoint size-boundary test must succeed");
receiver
.recv_timeout(TEST_TIMEOUT)
.expect("checkpoint size-boundary test must finish within 30 seconds")
.unwrap_or_else(|error| panic!("checkpoint size-boundary test failed: {error}"));
fs::remove_dir_all(&root)
.expect("cleaning checkpoint size-boundary test root must succeed");
}
async fn verify_size_boundary(
rt: MultiTaskRuntime<()>,
root: PathBuf,
) -> std::result::Result<(), String> {
let logger = build_small_physical_file_logger(rt, root.clone()).await?;
let old_path = logger
.0
.file
.writable_path()
.ok_or_else(|| "small-limit logger omitted initial writable file".to_owned())?;
let first_uid = Guid(0x7101);
let second_uid = Guid(0x7102);
let first_handle = logger
.append(first_uid.clone(), vec![0x61; 768 * 1024])
.await
.map_err(|error| format!("appending first WAL failed: {error}"))?;
logger
.flush(first_handle)
.await
.map_err(|error| format!("flushing first WAL failed: {error}"))?;
let old_len_before = logger.0.file.writable_size();
if old_len_before == 0 || old_len_before >= PHYSICAL_FILE_LIMIT {
return Err(format!(
"first WAL must leave the physical file below its limit: len={old_len_before}, limit={PHYSICAL_FILE_LIMIT}",
));
}
if logger.0.file.writable_path().as_ref() != Some(&old_path) {
return Err("CommitLogger flush must not auto-split the physical WAL".to_owned());
}
let crossing_payload_len = PHYSICAL_FILE_LIMIT - old_len_before + 64 * 1024;
let second_handle = logger
.append(second_uid.clone(), vec![0x71; crossing_payload_len])
.await
.map_err(|error| format!("appending threshold-crossing WAL failed: {error}"))?;
let new_checkpoint = logger
.append_check_point()
.await
.map_err(|error| format!("rotating threshold-crossing checkpoint failed: {error}"))?;
let new_path = logger
.0
.file
.writable_path()
.ok_or_else(|| "checkpoint rotation omitted new writable file".to_owned())?;
if new_path == old_path {
return Err("checkpoint rotation must replace the physical writable file".to_owned());
}
if logger.0.file.readable_amount() != 1 {
return Err(format!(
"threshold-crossing checkpoint must split exactly once: readable_amount={}",
logger.0.file.readable_amount(),
));
}
let parsed_checkpoint = new_path
.file_name()
.and_then(|name| name.to_str())
.and_then(log_file_name_to_usize)
.ok_or_else(|| format!("new checkpoint path is invalid: {new_path:?}"))?;
if parsed_checkpoint != new_checkpoint {
return Err(format!(
"returned checkpoint must identify the actual writable file: returned={new_checkpoint}, actual={parsed_checkpoint}",
));
}
let old_len_after = old_path
.metadata()
.map_err(|error| format!("reading old WAL metadata failed: {error}"))?
.len() as usize;
if old_len_after <= PHYSICAL_FILE_LIMIT {
return Err(format!(
"pending WAL must be written to the old file before its single split: len={old_len_after}",
));
}
let new_len = new_path
.metadata()
.map_err(|error| format!("reading new WAL metadata failed: {error}"))?
.len();
if new_len != 0 {
return Err(format!(
"new checkpoint must not contain pre-rotation WAL: len={new_len}",
));
}
logger
.flush(second_handle)
.await
.map_err(|error| format!("flushing helper-committed WAL failed: {error}"))?;
logger
.confirm(second_uid)
.await
.map_err(|error| format!("confirming second WAL failed: {error}"))?;
if !old_path.exists() {
return Err("old WAL must remain active until every registered transaction confirms".to_owned());
}
logger
.confirm(first_uid)
.await
.map_err(|error| format!("confirming first WAL failed: {error}"))?;
let mut backup_path = old_path.clone();
if !backup_path.set_extension("bak") {
return Err(format!("old WAL path cannot form a backup path: {old_path:?}"));
}
if old_path.exists() {
return Err("fully confirmed old WAL must no longer remain active".to_owned());
}
let backup_len = backup_path
.metadata()
.map_err(|error| format!("confirmed old WAL backup is missing: {error}"))?
.len() as usize;
if backup_len != old_len_after {
return Err(format!(
"confirmed backup must preserve the exact old WAL bytes: active={old_len_after}, backup={backup_len}",
));
}
if logger.waiting_confirm_count().await != 0 ||
logger.append_total_count() != 2 ||
logger.confirm_total_count() != 2 {
return Err(format!(
"checkpoint accounting did not close: waiting={}, appended={}, confirmed={}",
logger.waiting_confirm_count().await,
logger.append_total_count(),
logger.confirm_total_count(),
));
}
Ok(())
}
async fn verify_public_delay_commit_boundary(
rt: MultiTaskRuntime<()>,
path: PathBuf,
) -> std::result::Result<(), String> {
let file = LogFile::open(rt.clone(),
path,
CURRENT_BLOCK_LIMIT,
PHYSICAL_FILE_LIMIT,
None)
.await
.map_err(|error| format!("opening public LogFile boundary fixture failed: {error}"))?;
let old_path = file
.writable_path()
.ok_or_else(|| "public LogFile omitted initial writable path".to_owned())?;
let first_value = vec![0x81; PHYSICAL_FILE_LIMIT + 64 * 1024];
let first_handle = file.append(LogMethod::PlainAppend, b"first", &first_value);
if first_handle == 0 {
return Err("public LogFile append returned the reserved handle 0".to_owned());
}
file.delay_commit(first_handle, false, 1)
.await
.map_err(|error| format!("public threshold delay commit failed: {error}"))?;
let mut new_path = file
.writable_path()
.ok_or_else(|| "public threshold split omitted writable path".to_owned())?;
for _ in 0..100 {
if new_path != old_path && file.readable_amount() == 1 {
break;
}
rt.timeout(1).await;
new_path = file
.writable_path()
.ok_or_else(|| "public threshold split lost writable path".to_owned())?;
}
if new_path == old_path || file.readable_amount() != 1 {
return Err(format!(
"public delay commit must auto-split exactly once: old={old_path:?}, new={new_path:?}, readable={}",
file.readable_amount(),
));
}
let old_len = old_path
.metadata()
.map_err(|error| format!("reading public old WAL metadata failed: {error}"))?
.len() as usize;
if old_len <= PHYSICAL_FILE_LIMIT || file.writable_size() != 0 {
return Err(format!(
"public auto-split file sizes are invalid: old={old_len}, new={}",
file.writable_size(),
));
}
file.delay_commit(first_handle, false, 10)
.await
.map_err(|error| format!("repeating committed public handle failed: {error}"))?;
let second_handle = file.append(LogMethod::PlainAppend, b"second", b"value");
file.delay_commit(second_handle, false, 10)
.await
.map_err(|error| format!("public successor delay commit failed: {error}"))?;
if file.commited_uid() != second_handle || file.writable_path().as_ref() != Some(&new_path) {
return Err(format!(
"public successor did not close on the same writable file: committed={}, expected={}, writable={:?}, expected_path={new_path:?}",
file.commited_uid(),
second_handle,
file.writable_path(),
));
}
if file.writable_size() == 0 || file.readable_amount() != 1 {
return Err(format!(
"public successor produced an unexpected split or empty write: writable_len={}, readable={}",
file.writable_size(),
file.readable_amount(),
));
}
Ok(())
}
async fn build_small_physical_file_logger(
rt: MultiTaskRuntime<()>,
path: PathBuf,
) -> std::result::Result<CommitLogger, String> {
let file = LogFile::open(rt.clone(),
path,
CURRENT_BLOCK_LIMIT,
PHYSICAL_FILE_LIMIT,
None)
.await
.map_err(|error| format!("opening small-limit LogFile failed: {error}"))?;
let check_point_path = Arc::new(
file.writable_path()
.ok_or_else(|| "small-limit LogFile omitted writable path".to_owned())?,
);
let check_point_counter = Arc::new(TestAtomicU64::new(0));
Ok(CommitLogger(Arc::new(InnerCommitLogger {
rt,
file,
delay_timeout: 1,
log_file_limit: u64::MAX,
writed_size: TestAtomicU64::new(0),
writable: SpinLock::new((check_point_counter, check_point_path)),
only_reads: SpinLock::new(VecDeque::new()),
check_points: Mutex::new(XHashMap::default()),
is_replaying: AtomicBool::new(false),
replay_only_reads: SpinLock::new(VecDeque::new()),
replay_confirm_buf: SpinLock::new(VecDeque::new()),
commit_log_count: TestAtomicUsize::new(0),
confirm_commited_count: TestAtomicUsize::new(0),
})))
}
fn unique_test_root() -> PathBuf {
let nonce = SystemTime::now()
.duration_since(UNIX_EPOCH)
.expect("system time must be after UNIX_EPOCH")
.as_nanos();
std::env::temp_dir().join(format!(
"pi-store-checkpoint-size-boundary-{}-{nonce}",
std::process::id(),
))
}
}
struct CommitLoggerLoaderExt<
B: BufMut + AsRef<[u8]> + From<Vec<u8>> + Send + Sized + 'static,
F: Fn(Option<(Guid, LogMethod, u64, B)>) -> Result<()> + Send + 'static,
> {
logger: CommitLogger, buf: Vec<(Guid, LogMethod, u64, Vec<u8>)>, log_file: Option<PathBuf>, callback: Arc<F>, result: Result<(usize, usize)>, marker: PhantomData<B>,
}
impl<
B: BufMut + AsRef<[u8]> + From<Vec<u8>> + Send + Sized + 'static,
F: Fn(Option<(Guid, LogMethod, u64, B)>) -> Result<()> + Send + 'static,
> PairLoaderExt for CommitLoggerLoaderExt<B, F> {
fn is_require(&self,
_log_file: Option<&PathBuf>,
_payload_time: u64,
_key: &Vec<u8>) -> bool {
true
}
fn load(&mut self,
log_file: Option<&PathBuf>,
method: LogMethod,
payload_time: u64,
key: Vec<u8>,
value: Option<Vec<u8>>) {
if self.result.is_err() {
return;
}
if let Some(log_file) = log_file {
if self.log_file.is_none() {
self.log_file = Some(log_file.clone());
}
if self.log_file.as_ref().unwrap() != log_file {
while let Some((commit_uid, method, time, log)) = self.buf.pop() {
if let Err(e) = (self.callback)(Some((commit_uid.clone(), method, time, B::from(log)))) {
self.result = Err(Error::new(ErrorKind::Other, format!("Replay commit log failed, commit_uid: {:?}, reason: {:?}", commit_uid, e)));
}
}
self.log_file = Some(log_file.clone());
next_check_point(&self.logger);
}
let uid = u128::from_le_bytes(key.try_into().unwrap());
let commit_uid = Guid(uid);
if let Some(log) = value {
if let Ok((log_count, bytes_count)) = self.result {
self.result = Ok((log_count + 1, bytes_count + 16 + log.len()));
}
self.buf.push((commit_uid, method, payload_time, log));
}
}
}
}
impl<
B: BufMut + AsRef<[u8]> + From<Vec<u8>> + Send + Sized + 'static,
F: Fn(Option<(Guid, LogMethod, u64, B)>) -> Result<()> + Send + 'static,
> CommitLoggerLoaderExt<B, F> {
pub fn result(mut self) -> Result<(usize, usize)> {
if self.buf.len() > 0 {
while let Some((commit_uid, method, time, log)) = self.buf.pop() {
if let Err(e) = (self.callback)(Some((commit_uid.clone(), method, time, B::from(log)))) {
self.result = Err(Error::new(ErrorKind::Other,
format!("Replay commit log failed, commit_uid: {:?}, reason: {:?}",
commit_uid,
e)));
}
}
next_check_point(&self.logger);
}
(self.callback)(None);
self.result
}
}