use std::collections::BTreeMap;
use std::io;
use std::path::Path;
use std::sync::atomic::{AtomicU64, Ordering};
use crate::caps::{self, Caps};
use crate::error::{Error, Result};
use crate::journal::{self, FORMAT, Journal};
use crate::overlay::Overlay;
use crate::path::{RelPath, validate_user_path};
use crate::plan::{self, Plan};
use crate::progress::Record;
use crate::recover::{self, RecoveryReport, append_progress, barrier};
use crate::state::{Marker, TX_PREFIX, TxDir, private_dir};
use crate::vfs::{FileId, Kind, LockGuard, Vfs};
#[derive(Clone, Debug, Default)]
pub struct Options {
allow_untested_fs: bool,
}
impl Options {
pub fn new() -> Self {
Self::default()
}
pub(crate) fn allow_untested(&self) -> bool {
self.allow_untested_fs
}
pub fn allow_untested_fs(mut self, yes: bool) -> Self {
self.allow_untested_fs = yes;
self
}
}
pub struct Transaction {
vfs: Box<dyn Vfs>,
tx: TxDir,
overlay: Overlay,
blobs: BTreeMap<u32, FileId>,
next_blob: u32,
caps: Caps,
recovered: RecoveryReport,
finished: bool,
_lock: LockGuard,
}
impl std::fmt::Debug for Transaction {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Transaction")
.field("tx", &self.tx.name)
.finish_non_exhaustive()
}
}
fn tx_name() -> String {
static COUNTER: AtomicU64 = AtomicU64::new(0);
let nanos = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_nanos())
.unwrap_or(0);
format!(
"{TX_PREFIX}{nanos:x}-{:x}-{:x}",
std::process::id(),
COUNTER.fetch_add(1, Ordering::Relaxed)
)
}
pub(crate) fn ensure_private_dir(vfs: &dyn Vfs) -> Result<()> {
match vfs.stat(&private_dir())? {
Some(m) if m.kind == Kind::Dir => Ok(()),
Some(_) => Err(Error::UnsupportedFileType(private_dir().to_path_buf())),
None => {
vfs.mkdir(&private_dir())?;
vfs.sync_dir(&RelPath::root())?;
Ok(())
}
}
}
impl Transaction {
pub fn begin(root: impl AsRef<Path>) -> Result<Transaction> {
Self::begin_with(root, &Options::default())
}
pub fn begin_with(root: impl AsRef<Path>, opts: &Options) -> Result<Transaction> {
#[cfg(target_os = "linux")]
{
let fs = crate::sys::linux::LinuxFs::open(root.as_ref())?;
Self::begin_on(Box::new(fs), opts)
}
#[cfg(not(target_os = "linux"))]
{
let _ = (root, opts);
Err(Error::UnsupportedPlatform)
}
}
pub fn begin_on(vfs: Box<dyn Vfs>, opts: &Options) -> Result<Transaction> {
ensure_private_dir(&*vfs)?;
let lock = vfs.lock(true)?;
let caps = caps::probe(&*vfs, opts.allow_untested_fs)?;
let recovered = recover::recover_locked(&*vfs)?;
let tx = TxDir { name: tx_name() };
vfs.mkdir(&tx.dir())?;
vfs.mkdir(&tx.staged())?;
vfs.mkdir(&tx.backup())?;
let overlay = Overlay::new(vfs.root_id()?, caps.case_insensitive);
Ok(Transaction {
vfs,
tx,
overlay,
blobs: BTreeMap::new(),
next_blob: 0,
caps,
recovered,
finished: false,
_lock: lock,
})
}
pub fn recovered(&self) -> &RecoveryReport {
&self.recovered
}
pub fn case_insensitive(&self) -> bool {
self.caps.case_insensitive
}
fn blob_path(&self, n: u32) -> RelPath {
self.tx.staged().join(format!("b{n}"))
}
pub fn write(&mut self, path: impl AsRef<Path>, data: impl AsRef<[u8]>) -> Result<()> {
self.write_impl(path.as_ref(), data.as_ref(), None)
}
pub fn write_with_mode(
&mut self,
path: impl AsRef<Path>,
data: impl AsRef<[u8]>,
mode: u32,
) -> Result<()> {
self.write_impl(path.as_ref(), data.as_ref(), Some(mode & 0o7777))
}
fn write_impl(&mut self, path: &Path, data: &[u8], mode: Option<u32>) -> Result<()> {
let p = validate_user_path(path)?;
let inherited = self.overlay.check_write(&*self.vfs, &p)?;
let mode = mode.or(inherited);
let n = self.next_blob;
self.next_blob += 1;
let meta = self.vfs.create_file(&self.blob_path(n), data, mode)?;
self.blobs.insert(n, meta.id);
self.overlay.apply_write(&*self.vfs, &p, n, mode)
}
pub fn create_dir_all(&mut self, path: impl AsRef<Path>) -> Result<()> {
let p = validate_user_path(path.as_ref())?;
self.overlay.create_dir_all(&*self.vfs, &p)
}
pub fn rename(&mut self, from: impl AsRef<Path>, to: impl AsRef<Path>) -> Result<()> {
let from = validate_user_path(from.as_ref())?;
let to = validate_user_path(to.as_ref())?;
self.overlay.rename(&*self.vfs, &from, &to)
}
pub fn remove(&mut self, path: impl AsRef<Path>) -> Result<()> {
let p = validate_user_path(path.as_ref())?;
self.overlay.remove(&*self.vfs, &p)
}
pub fn remove_dir_all(&mut self, path: impl AsRef<Path>) -> Result<()> {
let p = validate_user_path(path.as_ref())?;
self.overlay.remove_dir_all(&*self.vfs, &p)
}
pub fn read(&self, path: impl AsRef<Path>) -> Result<Vec<u8>> {
let p = validate_user_path(path.as_ref())?;
self.overlay.read(&*self.vfs, &p, |n| self.blob_path(n))
}
pub fn exists(&self, path: impl AsRef<Path>) -> Result<bool> {
let p = validate_user_path(path.as_ref())?;
Ok(self.overlay.resolve(&*self.vfs, &p)?.is_some())
}
pub fn commit(mut self) -> Result<()> {
self.finished = true;
let vfs = &*self.vfs;
let tx = self.tx.clone();
let plan = match plan::compile(vfs, &self.overlay, &tx, &self.blobs) {
Ok(p) => p,
Err(e) => {
let _ = tx.gc(vfs);
return Err(e);
}
};
if plan.waves.is_empty() {
tx.gc(vfs)?;
return Ok(());
}
let j = Journal {
format: FORMAT,
tx: tx.name.clone(),
root: vfs.root_id()?,
plan,
};
if let Err(e) = prepare(vfs, &tx, &j) {
let _ = tx.gc(vfs);
return Err(e.into());
}
vfs.event("prepared");
if let Err(cause) = apply(vfs, &tx, &j.plan) {
return match recover::rollback(vfs, &tx, &j) {
Ok(()) => Err(cause.into()),
Err(rb) => Err(Error::RollbackFailed {
cause: Box::new(cause.into()),
rollback: Box::new(rb),
}),
};
}
if let Err(source) = tx.set_marker(vfs, Marker::Committed) {
return Err(Error::CommitOutcomeUnknown {
tx: tx.name.clone(),
source,
});
}
vfs.event("committed");
let _ = tx.gc(vfs);
Ok(())
}
}
fn prepare(vfs: &dyn Vfs, tx: &TxDir, j: &Journal) -> io::Result<()> {
for tok in &j.plan.tokens {
if tok.kind == Kind::File && tok.locs[0].path.is_private() {
vfs.sync_file(&tok.locs[0].path)?;
}
}
vfs.sync_dir(&tx.staged())?;
vfs.sync_dir(&tx.backup())?;
vfs.create_file(&tx.progress(), b"", None)?;
vfs.sync_file(&tx.progress())?;
vfs.sync_dir(&tx.dir())?;
vfs.sync_dir(&private_dir())?;
vfs.create_file(&tx.journal_tmp(), &journal::encode(j), None)?;
vfs.sync_file(&tx.journal_tmp())?;
vfs.rename_noreplace(&tx.journal_tmp(), None, &tx.journal(), None)?;
vfs.sync_dir(&tx.dir())
}
fn apply(vfs: &dyn Vfs, tx: &TxDir, plan: &Plan) -> io::Result<()> {
for (w, wave) in plan.waves.iter().enumerate() {
for s in wave {
let (src, dst) = (plan.src(*s), plan.dst(*s));
vfs.rename_noreplace(&src.path, Some(src.parent), &dst.path, Some(dst.parent))?;
}
barrier(vfs, plan, wave)?;
append_progress(vfs, tx, Record::ApplyDone(w as u32))?;
}
Ok(())
}
impl Drop for Transaction {
fn drop(&mut self) {
if !self.finished {
let _ = self.tx.gc(&*self.vfs);
}
}
}