pub mod corpus;
pub mod registry;
use std::fs::File;
use std::future::Future;
use std::ops::ControlFlow;
use std::path::{Path, PathBuf};
use std::time::Duration;
use tokio_util::sync::CancellationToken;
const CHECKPOINT_LOCK_FILE: &str = ".bridge-checkpoint.lock";
#[doc(hidden)]
pub fn acquire_checkpoint_lock(dir: &Path, prefix: &str) -> Result<File, String> {
std::fs::create_dir_all(dir).map_err(|error| {
format!(
"create {prefix} checkpoint directory {}: {error}",
dir.display()
)
})?;
let lock_path = dir.join(CHECKPOINT_LOCK_FILE);
let lock = std::fs::OpenOptions::new()
.create(true)
.truncate(false)
.read(true)
.write(true)
.open(&lock_path)
.map_err(|error| format!("open {prefix} lock {}: {error}", lock_path.display()))?;
lock.lock()
.map_err(|error| format!("acquire {prefix} lock {}: {error}", lock_path.display()))?;
Ok(lock)
}
#[doc(hidden)]
pub async fn acquire_checkpoint_lock_async(
dir: PathBuf,
prefix: &'static str,
) -> Result<File, String> {
tokio::task::spawn_blocking(move || acquire_checkpoint_lock(&dir, prefix))
.await
.map_err(|error| format!("{prefix} lock task failed: {error}"))?
}
#[doc(hidden)]
pub async fn rotation_watch_loop<F, Fut>(
interval: Duration,
shutdown: CancellationToken,
mut tick: F,
) where
F: FnMut() -> Fut,
Fut: Future<Output = ControlFlow<()>>,
{
let start = tokio::time::Instant::now() + interval;
let mut ticks = tokio::time::interval_at(start, interval);
ticks.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop {
tokio::select! {
_ = shutdown.cancelled() => break,
_ = ticks.tick() => {}
}
if tick().await.is_break() {
break;
}
}
}