use std::fs;
use std::path::{Path, PathBuf};
use std::sync::Mutex;
use std::time::{Duration, Instant};
use crate::config::FeffConfig;
use crate::error::{Error, PipelineError};
use crate::stage::Stage;
static FEFF_EXEC_LOCK: Mutex<()> = Mutex::new(());
#[derive(Debug)]
pub struct StageResult {
pub stage: Stage,
pub duration: Duration,
}
#[derive(Debug)]
pub struct PipelineResult {
pub stages: Vec<StageResult>,
pub work_dir: PathBuf,
}
#[derive(Debug)]
pub enum StageProgress {
Starting,
Finished { duration: Duration },
}
pub struct FeffPipeline {
config: FeffConfig,
}
impl FeffPipeline {
pub fn new(config: FeffConfig) -> Self {
Self { config }
}
pub fn run(&self) -> Result<PipelineResult, Error> {
self.run_with_progress(|_, _| {})
}
pub fn run_with_progress<F>(&self, mut callback: F) -> Result<PipelineResult, Error>
where
F: FnMut(Stage, StageProgress),
{
fs::create_dir_all(&self.config.work_dir)?;
let inp_path = self.config.work_dir.join("feff.inp");
let mut file = fs::File::create(&inp_path)?;
self.config.input.write_to(&mut file)?;
let mut stage_results = Vec::new();
let feff_error_path = self.config.work_dir.join(".feff.error");
let _ = fs::remove_file(&feff_error_path);
let _exec_lock = match FEFF_EXEC_LOCK.lock() {
Ok(lock) => lock,
Err(poisoned) => poisoned.into_inner(),
};
for &stage in &self.config.stages {
callback(stage, StageProgress::Starting);
let start = Instant::now();
run_stage_isolated(stage, &self.config.work_dir, self.config.stage_timeout)?;
let duration = start.elapsed();
callback(stage, StageProgress::Finished { duration });
let feff_error = fs::read_to_string(&feff_error_path)
.ok()
.and_then(|s| if s.trim().is_empty() { None } else { Some(s) });
if feff_error.is_some() {
return Err(Error::Pipeline(PipelineError {
stage: stage.executable_name().to_string(),
exit_code: None,
stderr: String::new(),
feff_error,
}));
}
stage_results.push(StageResult { stage, duration });
}
Ok(PipelineResult {
stages: stage_results,
work_dir: self.config.work_dir.clone(),
})
}
}
#[cfg(unix)]
fn run_stage_isolated(
stage: Stage,
work_dir: &std::path::Path,
timeout: Option<Duration>,
) -> Result<(), Error> {
let _cwd_guard = CwdGuard::enter(work_dir)?;
let pid = unsafe { libc::fork() };
match pid {
-1 => Err(Error::Io(std::io::Error::last_os_error())),
0 => {
unsafe { stage.call_ffi() };
unsafe { libc::_exit(0) };
}
child_pid => {
if timeout.is_some() {
wait_with_timeout(stage, child_pid, timeout)
} else {
wait_blocking(stage, child_pid)
}
}
}
}
#[cfg(unix)]
fn wait_blocking(stage: Stage, child_pid: libc::pid_t) -> Result<(), Error> {
let mut status: libc::c_int = 0;
loop {
let ret = unsafe { libc::waitpid(child_pid, &mut status, 0) };
if ret == -1 {
let err = std::io::Error::last_os_error();
if err.raw_os_error() == Some(libc::EINTR) {
continue;
}
return Err(Error::Io(err));
}
break;
}
check_child_status(stage, status)
}
#[cfg(unix)]
fn wait_with_timeout(
stage: Stage,
child_pid: libc::pid_t,
timeout: Option<Duration>,
) -> Result<(), Error> {
let start = Instant::now();
loop {
let mut status: libc::c_int = 0;
let ret = unsafe { libc::waitpid(child_pid, &mut status, libc::WNOHANG) };
if ret == -1 {
let err = std::io::Error::last_os_error();
if err.raw_os_error() == Some(libc::EINTR) {
continue;
}
return Err(Error::Io(err));
}
if ret == child_pid {
return check_child_status(stage, status);
}
if let Some(t) = timeout
&& start.elapsed() > t
{
unsafe { libc::kill(child_pid, libc::SIGKILL) };
unsafe { libc::waitpid(child_pid, std::ptr::null_mut(), 0) };
return Err(Error::Pipeline(PipelineError {
stage: stage.executable_name().to_string(),
exit_code: None,
stderr: format!("timed out after {}s", t.as_secs()),
feff_error: None,
}));
}
std::thread::sleep(Duration::from_millis(50));
}
}
#[cfg(unix)]
fn check_child_status(stage: Stage, status: libc::c_int) -> Result<(), Error> {
if libc::WIFEXITED(status) {
let exit_code = libc::WEXITSTATUS(status);
if exit_code == 0 {
Ok(())
} else {
Err(Error::Pipeline(PipelineError {
stage: stage.executable_name().to_string(),
exit_code: Some(exit_code),
stderr: String::new(),
feff_error: None,
}))
}
} else if libc::WIFSIGNALED(status) {
let signal = libc::WTERMSIG(status);
Err(Error::Pipeline(PipelineError {
stage: stage.executable_name().to_string(),
exit_code: None,
stderr: format!("killed by signal {signal}"),
feff_error: None,
}))
} else {
Err(Error::Pipeline(PipelineError {
stage: stage.executable_name().to_string(),
exit_code: None,
stderr: "unknown child status".to_string(),
feff_error: None,
}))
}
}
#[cfg(not(unix))]
fn run_stage_isolated(
stage: Stage,
work_dir: &std::path::Path,
_timeout: Option<Duration>,
) -> Result<(), Error> {
let _cwd_guard = CwdGuard::enter(work_dir)?;
unsafe { stage.call_ffi() };
Ok(())
}
struct CwdGuard {
old_dir: PathBuf,
}
impl CwdGuard {
fn enter(dir: &Path) -> Result<Self, Error> {
let old_dir = std::env::current_dir()?;
std::env::set_current_dir(dir)?;
Ok(Self { old_dir })
}
}
impl Drop for CwdGuard {
fn drop(&mut self) {
let _ = std::env::set_current_dir(&self.old_dir);
}
}
#[cfg(test)]
mod tests {
use super::CwdGuard;
#[test]
fn cwd_guard_restores_previous_directory() {
let original = std::env::current_dir().unwrap().canonicalize().unwrap();
let tmp = tempfile::tempdir().unwrap();
let tmp_canon = tmp.path().canonicalize().unwrap();
{
let _guard = CwdGuard::enter(tmp.path()).unwrap();
let now = std::env::current_dir().unwrap().canonicalize().unwrap();
assert_eq!(now, tmp_canon);
}
let restored = std::env::current_dir().unwrap().canonicalize().unwrap();
assert_eq!(restored, original);
}
}