use super::file_lock::{self, FileLockError};
use crate::model::command_custody::{CommandCustodyRecord, CommandCustodyRecordError};
use std::{
fs::{self, File},
io,
path::{Path, PathBuf},
process::{Child, Command},
thread,
time::{Duration, Instant},
};
use thiserror::Error;
const QUIESCENCE_GRACE: Duration = Duration::from_millis(250);
pub const COMMAND_CUSTODY_DESCRIPTOR_ENV: &str = "IC_BACKUP_COMMAND_CUSTODY_FD";
#[derive(Debug)]
pub struct CommandLifetimeLock {
file: File,
path: PathBuf,
record: CommandCustodyRecord,
dispatched: bool,
}
impl CommandLifetimeLock {
pub fn acquire(
journal: &Path,
operation_sequence: u64,
) -> Result<Self, CommandLifetimeLockError> {
let journal = resolve_journal(journal)?;
let path = lock_path(&journal, operation_sequence);
let file = file_lock::acquire(&path).map_err(|error| project_error(&path, error))?;
let (device, inode) = file_identity(&file.metadata()?)?;
let record = CommandCustodyRecord::new(journal, operation_sequence, device, inode)?;
Ok(Self {
file,
path,
record,
dispatched: false,
})
}
#[must_use]
pub fn record(&self) -> &CommandCustodyRecord {
&self.record
}
#[must_use]
pub fn path(&self) -> &Path {
&self.path
}
pub fn spawn(&mut self, mut command: Command) -> Result<Child, CommandLifetimeLockError> {
if self.dispatched {
return Err(CommandLifetimeLockError::AlreadyDispatched {
path: self.path.clone(),
});
}
require_identity(&self.record, &fs::symlink_metadata(&self.path)?, &self.path)?;
self.dispatched = true;
#[cfg(unix)]
{
use command_fds::CommandFdExt;
use std::os::fd::AsRawFd;
let descriptor =
rustix::io::fcntl_dupfd_cloexec(&self.file, 3).map_err(io::Error::from)?;
command.env(
COMMAND_CUSTODY_DESCRIPTOR_ENV,
descriptor.as_raw_fd().to_string(),
);
command.preserved_fds(vec![descriptor]);
let child = command.spawn();
drop(command);
child.map_err(CommandLifetimeLockError::from)
}
#[cfg(not(unix))]
{
let _ = command;
Err(io::Error::from(io::ErrorKind::Unsupported).into())
}
}
pub fn finish(self) -> Result<CommandQuiescenceGuard, CommandLifetimeLockError> {
let Self { file, record, .. } = self;
drop(file);
let deadline = Instant::now() + QUIESCENCE_GRACE;
loop {
match CommandQuiescenceGuard::acquire(&record) {
Err(CommandLifetimeLockError::InFlight { .. }) if Instant::now() < deadline => {
thread::sleep(Duration::from_millis(5));
}
result => return result,
}
}
}
}
#[derive(Debug)]
pub struct CommandQuiescenceGuard {
file: File,
record: CommandCustodyRecord,
}
impl CommandQuiescenceGuard {
pub fn acquire(expected: &CommandCustodyRecord) -> Result<Self, CommandLifetimeLockError> {
let path = lock_path(expected.journal(), expected.operation_sequence());
let file =
file_lock::acquire_existing(&path).map_err(|error| project_error(&path, error))?;
require_identity(expected, &file.metadata()?, &path)?;
Ok(Self {
file,
record: expected.clone(),
})
}
#[must_use]
pub fn record(&self) -> &CommandCustodyRecord {
&self.record
}
}
impl Drop for CommandQuiescenceGuard {
fn drop(&mut self) {
#[cfg(unix)]
file_lock::unlock(&self.file);
}
}
fn resolve_journal(path: &Path) -> Result<PathBuf, CommandLifetimeLockError> {
let name = path
.file_name()
.ok_or_else(|| CommandLifetimeLockError::InvalidJournal {
path: path.to_path_buf(),
})?;
let parent = path
.parent()
.filter(|parent| !parent.as_os_str().is_empty())
.unwrap_or_else(|| Path::new("."));
let parent = parent.canonicalize()?;
if !parent.is_dir() {
return Err(io::Error::from(io::ErrorKind::NotADirectory).into());
}
let journal = parent.join(name);
match fs::symlink_metadata(&journal) {
Ok(metadata) if !metadata.is_file() => {
return Err(CommandLifetimeLockError::InvalidJournal { path: journal });
}
Err(error) if error.kind() != io::ErrorKind::NotFound => return Err(error.into()),
_ => {}
}
Ok(journal)
}
fn lock_path(journal: &Path, operation_sequence: u64) -> PathBuf {
let mut path = journal.as_os_str().to_os_string();
path.push(format!(".command-{operation_sequence}.lock"));
PathBuf::from(path)
}
fn file_identity(metadata: &fs::Metadata) -> io::Result<(u64, u64)> {
if !metadata.is_file() {
return Err(io::Error::from(io::ErrorKind::InvalidInput));
}
#[cfg(unix)]
{
use std::os::unix::fs::MetadataExt;
Ok((metadata.dev(), metadata.ino()))
}
#[cfg(not(unix))]
{
let _ = metadata;
Err(io::Error::from(io::ErrorKind::Unsupported))
}
}
fn require_identity(
expected: &CommandCustodyRecord,
metadata: &fs::Metadata,
path: &Path,
) -> Result<(), CommandLifetimeLockError> {
if !metadata.is_file() {
return Err(CommandLifetimeLockError::IdentityChanged {
path: path.to_path_buf(),
});
}
let (device, inode) = file_identity(metadata)?;
if expected.matches_file(device, inode) {
return Ok(());
}
Err(CommandLifetimeLockError::IdentityChanged {
path: path.to_path_buf(),
})
}
fn project_error(path: &Path, error: FileLockError) -> CommandLifetimeLockError {
match error {
FileLockError::Locked => CommandLifetimeLockError::InFlight {
path: path.to_path_buf(),
},
FileLockError::UnsafeEntry { kind } => CommandLifetimeLockError::UnsafeEntry {
path: path.to_path_buf(),
kind,
},
FileLockError::Io(error) => CommandLifetimeLockError::Io(error),
}
}
#[derive(Debug, Error)]
pub enum CommandLifetimeLockError {
#[error("command custody remains in flight: {path:?}")]
InFlight {
path: PathBuf,
},
#[error("command custody identity changed: {path:?}")]
IdentityChanged {
path: PathBuf,
},
#[error("command spawn already attempted: {path:?}")]
AlreadyDispatched {
path: PathBuf,
},
#[error("unsafe command custody entry at {path:?}: {kind}")]
UnsafeEntry {
path: PathBuf,
kind: String,
},
#[error("invalid command journal location: {path:?}")]
InvalidJournal {
path: PathBuf,
},
#[error(transparent)]
Record(#[from] CommandCustodyRecordError),
#[error(transparent)]
Io(#[from] io::Error),
}
#[cfg(all(test, unix))]
mod tests;