mod format;
mod lock;
mod platform;
mod worker;
mod workspace;
use std::fs;
use std::path::{Path, PathBuf};
use std::sync::{Arc, Condvar, Mutex};
use std::thread::JoinHandle;
use lora_wal::{WalError, WalMirror};
use self::format::write_archive_atomic;
use self::lock::ArchiveLock;
use self::platform::sync_dir;
use self::worker::spawn_archive_worker;
use self::workspace::{cleanup_stale_temp_paths, make_work_dir, prepare_work_dir};
pub(crate) struct WalArchive {
archive_path: PathBuf,
work_dir: PathBuf,
max_archive_bytes: u64,
state: Arc<(Mutex<ArchiveState>, Condvar)>,
write_lock: Arc<Mutex<()>>,
worker: Option<JoinHandle<()>>,
_archive_lock: ArchiveLock,
}
#[derive(Debug, Default)]
struct ArchiveState {
dirty: bool,
force: bool,
shutdown: bool,
failure: Option<String>,
}
impl WalArchive {
pub fn open(archive_path: PathBuf, max_archive_bytes: u64) -> Result<Self, WalError> {
if archive_path.is_dir() {
return Err(WalError::Malformed(format!(
"database archive path is a directory: {}",
archive_path.display()
)));
}
if let Some(parent) = archive_path.parent() {
fs::create_dir_all(parent)?;
}
let archive_lock = ArchiveLock::acquire(&archive_path)?;
cleanup_stale_temp_paths(&archive_path)?;
let work_dir = make_work_dir(&archive_path);
prepare_work_dir(&archive_path, &work_dir, max_archive_bytes)?;
let state = Arc::new((Mutex::new(ArchiveState::default()), Condvar::new()));
let write_lock = Arc::new(Mutex::new(()));
let worker = Some(spawn_archive_worker(
state.clone(),
write_lock.clone(),
work_dir.clone(),
archive_path.clone(),
max_archive_bytes,
));
Ok(Self {
archive_path,
work_dir,
max_archive_bytes,
state,
write_lock,
worker,
_archive_lock: archive_lock,
})
}
pub fn work_dir(&self) -> &Path {
&self.work_dir
}
}
impl WalMirror for WalArchive {
fn persist(&self, wal_dir: &Path) -> Result<(), WalError> {
if wal_dir != self.work_dir {
return Err(WalError::Malformed(format!(
"archive mirror received unexpected WAL dir: {}",
wal_dir.display()
)));
}
let (lock, cv) = &*self.state;
let mut state = lock.lock().unwrap();
if let Some(failure) = &state.failure {
return Err(WalError::Malformed(format!(
"database archive writer failed: {failure}"
)));
}
state.dirty = true;
cv.notify_one();
Ok(())
}
fn persist_force(&self, wal_dir: &Path) -> Result<(), WalError> {
if wal_dir != self.work_dir {
return Err(WalError::Malformed(format!(
"archive mirror received unexpected WAL dir: {}",
wal_dir.display()
)));
}
{
let (lock, _) = &*self.state;
let state = lock.lock().unwrap();
if let Some(failure) = &state.failure {
return Err(WalError::Malformed(format!(
"database archive writer failed: {failure}"
)));
}
}
let _write_guard = self.write_lock.lock().unwrap();
{
let (lock, _) = &*self.state;
let state = lock.lock().unwrap();
if let Some(failure) = &state.failure {
return Err(WalError::Malformed(format!(
"database archive writer failed: {failure}"
)));
}
}
let result =
write_archive_atomic(&self.work_dir, &self.archive_path, self.max_archive_bytes);
let (lock, _) = &*self.state;
let mut state = lock.lock().unwrap();
match result {
Ok(()) => {
state.dirty = false;
state.force = false;
Ok(())
}
Err(err) => {
state.failure = Some(err.to_string());
Err(err)
}
}
}
}
impl Drop for WalArchive {
fn drop(&mut self) {
{
let (lock, cv) = &*self.state;
let mut state = lock.lock().unwrap();
state.dirty = true;
state.shutdown = true;
state.force = true;
cv.notify_one();
}
let mut shutdown_cleanly = true;
if let Some(worker) = self.worker.take() {
shutdown_cleanly = worker.join().is_ok();
}
{
let (lock, _) = &*self.state;
let state = lock.lock().unwrap();
shutdown_cleanly &= state.failure.is_none();
}
if shutdown_cleanly {
let _ = fs::remove_dir_all(&self.work_dir);
if let Some(parent) = self.work_dir.parent() {
let _ = sync_dir(parent);
}
}
}
}