#![allow(clippy::new_ret_no_self, clippy::unused_unit)]
use std::ffi::OsString;
use std::path::PathBuf;
use std::sync::Arc;
use tokio::sync::Mutex;
use id_effect::kernel::Effect;
use id_effect::{Env, Needs, ProviderError, ProviderSpec};
use crate::error::ProcessError;
#[derive(Clone, Debug)]
pub struct CommandSpec {
pub program: OsString,
pub args: Vec<OsString>,
pub current_dir: Option<PathBuf>,
}
impl CommandSpec {
#[inline]
pub fn new(program: impl Into<OsString>) -> Self {
Self {
program: program.into(),
args: Vec::new(),
current_dir: None,
}
}
#[inline]
pub fn arg(mut self, a: impl Into<OsString>) -> Self {
self.args.push(a.into());
self
}
#[inline]
pub fn dir(mut self, d: impl Into<PathBuf>) -> Self {
self.current_dir = Some(d.into());
self
}
}
#[derive(Clone)]
pub struct ChildHandle {
inner: Arc<Mutex<tokio::process::Child>>,
}
pub trait ProcessRuntime: Send + Sync + 'static {
fn spawn(&self, cmd: CommandSpec) -> Effect<ChildHandle, ProcessError, ()>;
fn spawn_wait(&self, cmd: CommandSpec) -> Effect<std::process::ExitStatus, ProcessError, ()>;
}
#[derive(Clone, Copy, Debug, Default)]
pub struct TokioProcessRuntime;
impl ProcessRuntime for TokioProcessRuntime {
fn spawn(&self, cmd: CommandSpec) -> Effect<ChildHandle, ProcessError, ()> {
Effect::new_async(move |_r: &mut ()| {
Box::pin(async move {
let mut c = tokio::process::Command::new(&cmd.program);
c.args(&cmd.args);
if let Some(dir) = &cmd.current_dir {
c.current_dir(dir);
}
let child = c.spawn().map_err(ProcessError::from)?;
Ok(ChildHandle {
inner: Arc::new(Mutex::new(child)),
})
})
})
}
fn spawn_wait(&self, cmd: CommandSpec) -> Effect<std::process::ExitStatus, ProcessError, ()> {
let this = *self;
Effect::new_async(move |_r: &mut ()| {
Box::pin(async move {
let handle = this.spawn(cmd).run(&mut ()).await?;
child_wait(handle).run(&mut ()).await
})
})
}
}
pub type ProcessRuntimeService = Arc<dyn ProcessRuntime>;
#[derive(::id_effect::ProviderSpecDerive)]
#[provides(ProcessRuntimeService)]
pub struct TokioProcessRuntimeProvider;
impl TokioProcessRuntimeProvider {
fn new() -> ProcessRuntimeService {
Arc::new(TokioProcessRuntime)
}
}
#[inline]
pub fn spawn<R>(cmd: CommandSpec) -> Effect<ChildHandle, ProcessError, R>
where
R: Needs<ProcessRuntimeService> + 'static,
{
Effect::new_async(move |r: &mut R| {
let rt = r.need().clone();
let inner = rt.spawn(cmd);
Box::pin(async move { inner.run(&mut ()).await })
})
}
#[inline]
pub fn spawn_wait<R>(cmd: CommandSpec) -> Effect<std::process::ExitStatus, ProcessError, R>
where
R: Needs<ProcessRuntimeService> + 'static,
{
Effect::new_async(move |r: &mut R| {
let rt = r.need().clone();
let inner = rt.spawn_wait(cmd);
Box::pin(async move { inner.run(&mut ()).await })
})
}
#[inline]
pub fn child_wait<R>(handle: ChildHandle) -> Effect<std::process::ExitStatus, ProcessError, R>
where
R: 'static,
{
Effect::new_async(move |_r: &mut R| {
Box::pin(async move {
let mut guard = handle.inner.lock().await;
let status = guard.wait().await.map_err(ProcessError::from)?;
Ok(status)
})
})
}
#[inline]
pub fn child_kill<R>(handle: ChildHandle) -> Effect<(), ProcessError, R>
where
R: 'static,
{
Effect::new_async(move |_r: &mut R| {
Box::pin(async move {
let mut guard = handle.inner.lock().await;
guard.kill().await.map_err(ProcessError::from)?;
Ok(())
})
})
}
#[cfg(test)]
mod tests {
use super::*;
use std::ffi::OsString;
mod command_spec {
use super::*;
#[test]
fn new_sets_program_with_no_args_or_dir() {
let c = CommandSpec::new("prog");
assert_eq!(c.program, OsString::from("prog"));
assert!(c.args.is_empty());
assert!(c.current_dir.is_none());
}
#[test]
fn arg_appends_in_order() {
let c = CommandSpec::new("p").arg("a").arg("b");
assert_eq!(c.args.len(), 2);
assert_eq!(c.args[0], OsString::from("a"));
assert_eq!(c.args[1], OsString::from("b"));
}
#[test]
fn dir_sets_working_directory() {
let d = PathBuf::from("/tmp/work");
let c = CommandSpec::new("p").dir(d.clone());
assert_eq!(c.current_dir, Some(d));
}
}
}