use super::lock::{self, ExecutorLock};
use std::ffi::OsStr;
use std::fs::File;
use std::io;
use std::os::fd::{FromRawFd, IntoRawFd, RawFd};
use std::os::unix::fs::MetadataExt;
use std::path::{Path, PathBuf};
use std::process::Command;
pub const LOCK_FD_ENV: &str = "LERNIE_LOCK_FD";
#[derive(Debug, thiserror::Error)]
pub enum AdoptError {
#[error("{LOCK_FD_ENV}={0:?} is not a valid fd number")]
Parse(String),
#[error("fstat adopted fd {fd}: {source}")]
Fstat {
fd: RawFd,
#[source]
source: io::Error,
},
#[error("stat inbox {inbox}: {source}")]
InboxStat {
inbox: PathBuf,
#[source]
source: io::Error,
},
#[error("adopted fd {fd} is not the inbox directory {inbox} (device/inode mismatch)")]
Mismatch { fd: RawFd, inbox: PathBuf },
#[error("restore close-on-exec on adopted fd {fd}: {source}")]
Cloexec {
fd: RawFd,
#[source]
source: io::Error,
},
#[error("re-assert flock on adopted fd {fd}: {source}")]
Flock {
fd: RawFd,
#[source]
source: io::Error,
},
}
#[derive(Debug, thiserror::Error)]
pub enum LeaseError {
#[error(transparent)]
Adopt(#[from] AdoptError),
#[error("acquire executor lock: {0}")]
Acquire(#[source] io::Error),
}
pub fn take_lease(
lease_env: Option<&OsStr>,
inbox_dir: &Path,
) -> Result<Option<ExecutorLock>, LeaseError> {
match lease_env {
Some(v) => Ok(adopt(v, inbox_dir)?),
None => lock::try_acquire(inbox_dir).map_err(LeaseError::Acquire),
}
}
pub fn adopt(env_val: &OsStr, inbox_dir: &Path) -> Result<Option<ExecutorLock>, AdoptError> {
let fd: RawFd = env_val
.to_str()
.and_then(|s| s.parse().ok())
.ok_or_else(|| AdoptError::Parse(env_val.to_string_lossy().into_owned()))?;
let file = unsafe { File::from_raw_fd(fd) };
if let Err(e) = validate(&file, fd, inbox_dir) {
let _ = file.into_raw_fd();
return Err(e);
}
interpret_cloexec(fd, set_fd_flags(fd, libc::FD_CLOEXEC))?;
interpret_flock(fd, lock::lock_or_none(file))
}
fn interpret_cloexec(fd: RawFd, r: io::Result<()>) -> Result<(), AdoptError> {
match r {
Ok(()) => Ok(()),
Err(source) => Err(AdoptError::Cloexec { fd, source }),
}
}
fn interpret_flock(
fd: RawFd,
r: io::Result<Option<ExecutorLock>>,
) -> Result<Option<ExecutorLock>, AdoptError> {
match r {
Ok(lock) => Ok(lock),
Err(source) => Err(AdoptError::Flock { fd, source }),
}
}
fn validate(file: &File, fd: RawFd, inbox_dir: &Path) -> Result<(), AdoptError> {
let fd_meta = file
.metadata()
.map_err(|source| AdoptError::Fstat { fd, source })?;
let dir_meta = std::fs::metadata(inbox_dir).map_err(|source| AdoptError::InboxStat {
inbox: inbox_dir.to_path_buf(),
source,
})?;
if fd_meta.dev() != dir_meta.dev() || fd_meta.ino() != dir_meta.ino() {
return Err(AdoptError::Mismatch {
fd,
inbox: inbox_dir.to_path_buf(),
});
}
Ok(())
}
pub fn successor_command(
exe: &Path,
workspace: &Path,
agent_id: &str,
lease: ExecutorLock,
) -> io::Result<Command> {
let fd = lease.as_raw_fd();
set_fd_flags(fd, 0)?;
std::mem::forget(lease);
let mut cmd = Command::new(exe);
cmd.arg("advance")
.arg(workspace)
.arg(agent_id)
.env(LOCK_FD_ENV, fd.to_string());
Ok(cmd)
}
fn set_fd_flags(fd: RawFd, flags: libc::c_int) -> io::Result<()> {
let ret = unsafe { libc::fcntl(fd, libc::F_SETFD, flags) };
if ret == -1 {
Err(io::Error::last_os_error())
} else {
Ok(())
}
}
#[cfg(test)]
mod tests;