use std::fs;
use std::path::{Path, PathBuf};
use std::sync::Mutex;
use std::time::{Duration, Instant};
use crate::config::{FeffConfig, StageIsolation};
use crate::error::{Error, PipelineError};
use crate::output::{FeffOutputs, FeffTable, PathsDat};
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,
}
impl PipelineResult {
pub fn outputs(&self) -> Result<FeffOutputs, Error> {
FeffOutputs::discover(&self.work_dir)
}
pub fn read_xmu(&self) -> Result<FeffTable, Error> {
FeffTable::from_file(self.work_dir.join("xmu.dat"))
}
pub fn read_xmu_strict(&self) -> Result<FeffTable, Error> {
FeffTable::from_file_strict(self.work_dir.join("xmu.dat"))
}
pub fn read_chi(&self) -> Result<FeffTable, Error> {
FeffTable::from_file(self.work_dir.join("chi.dat"))
}
pub fn read_chi_strict(&self) -> Result<FeffTable, Error> {
FeffTable::from_file_strict(self.work_dir.join("chi.dat"))
}
pub fn read_eels(&self) -> Result<FeffTable, Error> {
FeffTable::from_file(self.work_dir.join("eels.dat"))
}
pub fn read_eels_strict(&self) -> Result<FeffTable, Error> {
FeffTable::from_file_strict(self.work_dir.join("eels.dat"))
}
pub fn read_ldos(&self, index: u32) -> Result<FeffTable, Error> {
FeffTable::from_file(self.work_dir.join(format!("ldos{index:02}.dat")))
}
pub fn read_ldos_strict(&self, index: u32) -> Result<FeffTable, Error> {
FeffTable::from_file_strict(self.work_dir.join(format!("ldos{index:02}.dat")))
}
pub fn read_paths(&self) -> Result<PathsDat, Error> {
PathsDat::from_file(self.work_dir.join("paths.dat"))
}
}
#[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,
self.config.stage_isolation,
)?;
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>,
isolation: StageIsolation,
) -> Result<(), Error> {
match isolation {
StageIsolation::Fork => run_stage_forked(stage, work_dir, timeout),
StageIsolation::Worker => run_stage_worker(stage, work_dir, timeout),
StageIsolation::InProcess => run_stage_in_process(stage, work_dir),
StageIsolation::Auto => {
if !fork_unsafe_host() {
run_stage_forked(stage, work_dir, timeout)
} else if crate::worker::installed() {
run_stage_worker(stage, work_dir, timeout)
} else {
Err(Error::Pipeline(PipelineError {
stage: stage.executable_name().to_string(),
exit_code: None,
stderr: "this process is fork-unsafe (GUI host): call \
feff10::worker::init() at the top of main(), or run the \
pipeline from a separate process"
.to_string(),
feff_error: None,
}))
}
}
}
}
fn run_stage_worker(
stage: Stage,
work_dir: &std::path::Path,
timeout: Option<Duration>,
) -> Result<(), Error> {
let exe = std::env::current_exe()?;
let mut child = std::process::Command::new(exe)
.env(crate::worker::ENV_STAGE, stage.executable_name())
.env(crate::worker::ENV_DIR, work_dir)
.spawn()?;
let start = Instant::now();
loop {
match child.try_wait()? {
Some(status) => {
if status.success() {
return Ok(());
}
return Err(Error::Pipeline(PipelineError {
stage: stage.executable_name().to_string(),
exit_code: status.code(),
stderr: format!("worker exited with {status}"),
feff_error: None,
}));
}
None => {
if let Some(t) = timeout
&& start.elapsed() > t
{
let _ = child.kill();
let _ = child.wait();
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(target_os = "macos")]
fn fork_unsafe_host() -> bool {
unsafe extern "C" {
fn _dyld_image_count() -> u32;
fn _dyld_get_image_name(image_index: u32) -> *const std::os::raw::c_char;
}
unsafe {
let count = _dyld_image_count();
for i in 0..count {
let name = _dyld_get_image_name(i);
if name.is_null() {
continue;
}
let name = std::ffi::CStr::from_ptr(name).to_bytes();
for marker in [
&b"AppKit.framework"[..],
&b"Metal.framework"[..],
&b"UIKit.framework"[..],
] {
if name.windows(marker.len()).any(|w| w == marker) {
return true;
}
}
}
}
false
}
#[cfg(all(unix, not(target_os = "macos")))]
fn fork_unsafe_host() -> bool {
false
}
fn run_stage_in_process(stage: Stage, work_dir: &std::path::Path) -> Result<(), Error> {
let _cwd_guard = CwdGuard::enter(work_dir)?;
let handle = std::thread::Builder::new()
.name(format!("feff-{}", stage.executable_name()))
.stack_size(64 * 1024 * 1024)
.spawn(move || unsafe { stage.call_ffi() })
.map_err(Error::Io)?;
handle.join().map_err(|_| {
Error::Pipeline(PipelineError {
stage: stage.executable_name().to_string(),
exit_code: None,
stderr: "stage panicked".to_string(),
feff_error: None,
})
})
}
#[cfg(unix)]
fn run_stage_forked(
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>,
isolation: StageIsolation,
) -> Result<(), Error> {
match isolation {
StageIsolation::Worker => run_stage_worker(stage, work_dir, timeout),
StageIsolation::Auto if crate::worker::installed() => {
run_stage_worker(stage, work_dir, timeout)
}
_ => run_stage_in_process(stage, work_dir),
}
}
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, PipelineResult};
#[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);
}
#[test]
fn pipeline_result_reads_common_outputs() {
let tmp = tempfile::tempdir().unwrap();
std::fs::write(tmp.path().join("xmu.dat"), "# header\n1.0 2.0 3.0 4.0\n").unwrap();
std::fs::write(
tmp.path().join("paths.dat"),
"PATH Rmax= 5.5\n\
1 1 1.0 index, nleg, degeneracy, r= 1.0\n\
x y z ipot label rleg beta eta\n\
0.0 0.0 0.0 0 'A' 1.0 180.0 0.0\n",
)
.unwrap();
let result = PipelineResult {
stages: vec![],
work_dir: tmp.path().to_path_buf(),
};
let xmu = result.read_xmu().unwrap();
assert_eq!(xmu.nrows(), 1);
assert_eq!(xmu.ncols(), 4);
let paths = result.read_paths().unwrap();
assert_eq!(paths.len(), 1);
let outputs = result.outputs().unwrap();
assert!(outputs.file("xmu.dat").is_some());
assert!(outputs.file("paths.dat").is_some());
}
}