use crate::cap::{Auth, Cap};
use crate::error::{BotError, DispatchCertainty};
use crate::verb;
#[cfg(feature = "process")]
use crate::rt::process::{Frames, ProcessRun, ProcessRunError, ProcessSpec};
#[cfg(feature = "process")]
use crate::rt::supervise::Supervisor;
#[cfg(feature = "process")]
use crate::rt::sync::Semaphore;
#[cfg(feature = "process")]
use std::num::NonZeroUsize;
#[cfg(feature = "process")]
use std::time::Duration;
#[cfg(feature = "process")]
pub const DEFAULT_CAPTURE_LIMIT: NonZeroUsize = match NonZeroUsize::new(64 * 1024) {
Some(limit) => limit,
None => NonZeroUsize::MIN,
};
#[cfg(feature = "process")]
pub const DEFAULT_DEADLINE: Duration = Duration::from_secs(30);
#[cfg(feature = "process")]
pub const DEFAULT_MAX_CONCURRENT: NonZeroUsize = match NonZeroUsize::new(16) {
Some(limit) => limit,
None => NonZeroUsize::MIN,
};
#[cfg(feature = "process")]
#[derive(Debug)]
pub struct Process {
spec: ProcessSpec,
caps: Vec<Cap>,
slots: Semaphore,
frame_stdout: Option<NonZeroUsize>,
}
#[cfg(not(feature = "process"))]
#[derive(Debug)]
pub struct Process {
command: String,
caps: Vec<Cap>,
}
#[derive(PartialEq, Debug, Clone)]
#[non_exhaustive]
pub struct ProcessState {
pub running: bool,
pub exit_code: Option<i32>,
stdout: String,
stderr: String,
pub stdout_truncated: bool,
pub stderr_truncated: bool,
pub stdout_total_bytes: u64,
pub stderr_total_bytes: u64,
pub deadline_fired: bool,
pub signal: Option<i32>,
#[cfg(feature = "process")]
stdout_frames: Option<Frames>,
}
impl ProcessState {
#[must_use]
pub fn stdout(&self) -> &str {
&self.stdout
}
#[must_use]
pub fn stderr(&self) -> &str {
&self.stderr
}
#[cfg(feature = "process")]
#[must_use]
pub fn stdout_frames(&self) -> Option<&Frames> {
self.stdout_frames.as_ref()
}
}
#[cfg(feature = "process")]
impl Process {
pub fn new(command: impl Into<String>) -> Self {
let mut spec = ProcessSpec::new(command.into());
spec.capture_stdout(DEFAULT_CAPTURE_LIMIT);
spec.capture_stderr(DEFAULT_CAPTURE_LIMIT);
spec.deadline(DEFAULT_DEADLINE);
Self::from_spec(spec)
}
#[must_use]
pub fn from_spec(spec: ProcessSpec) -> Self {
Self {
spec,
caps: vec![Cap::sys()],
slots: Semaphore::new(DEFAULT_MAX_CONCURRENT.get()),
frame_stdout: None,
}
}
#[must_use]
pub fn frame_stdout(mut self, ceiling: NonZeroUsize) -> Self {
self.frame_stdout = Some(ceiling);
self
}
#[must_use]
pub fn max_concurrent(mut self, limit: NonZeroUsize) -> Self {
self.slots = Semaphore::new(limit.get());
self
}
async fn run_once(&self) -> Result<ProcessState, BotError> {
let Ok(_slot) = self.slots.acquire().await else {
let refusal = Err(BotError::DomainError {
domain: String::from("sys::process"),
certainty: DispatchCertainty::Refused,
cause: String::from("the process slots were closed before the process started"),
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "run_once: returning an error to the caller");
return refusal;
};
let mut supervisor = Supervisor::new(1);
match supervisor.run_process(&self.spec).await {
Ok(run) => Ok(ProcessState::from_run(&run, self.frame_stdout)),
Err(ProcessRunError::Refused) => Err(BotError::DomainError {
domain: String::from("sys::process"),
certainty: DispatchCertainty::Refused,
cause: String::from("the supervisor was cancelled before the process started"),
}),
Err(ProcessRunError::NotStarted { source }) => Err(BotError::DomainError {
domain: String::from("sys::process"),
certainty: DispatchCertainty::Refused,
cause: format!("the process did not start: {source}"),
}),
Err(ProcessRunError::AfterStart { source }) => Err(BotError::EffectIndeterminate {
domain: String::from("sys::process"),
cause: format!("the process started but its outcome is unknown: {source}"),
}),
}
}
}
#[cfg(not(feature = "process"))]
impl Process {
pub fn new(command: impl Into<String>) -> Self {
Self {
command: command.into(),
caps: vec![Cap::sys()],
}
}
}
#[cfg(feature = "process")]
impl ProcessState {
fn from_run(run: &ProcessRun, frame_stdout: Option<NonZeroUsize>) -> Self {
let stdout = run.stdout();
let stderr = run.stderr();
#[cfg(unix)]
let signal = run.signal();
#[cfg(not(unix))]
let signal: Option<i32> = None;
Self {
running: false,
exit_code: run.exit_code(),
stdout: String::from_utf8_lossy(stdout.bytes()).into_owned(),
stderr: String::from_utf8_lossy(stderr.bytes()).into_owned(),
stdout_truncated: stdout.truncated(),
stderr_truncated: stderr.truncated(),
stdout_total_bytes: stdout.total_bytes(),
stderr_total_bytes: stderr.total_bytes(),
deadline_fired: run.deadline_fired(),
signal,
stdout_frames: frame_stdout.map(|ceiling| stdout.frames(ceiling.get())),
}
}
}
impl Process {
async fn dispatch(&self, verb: &'static str) -> Result<ProcessState, BotError> {
#[cfg(feature = "process")]
{
let _ = verb;
self.run_once().await
}
#[cfg(not(feature = "process"))]
{
Err(BotError::DomainError {
domain: String::from("sys::process"),
certainty: DispatchCertainty::Refused,
cause: format!("{verb} {:?} — binding required", self.command),
})
}
}
}
impl verb::Observe for Process {
type Output = ProcessState;
fn required_caps(&self) -> &[Cap] {
&self.caps
}
async fn poll(&self, call: (Auth, ())) -> Result<ProcessState, BotError> {
call.0.check(self.required_caps())?;
self.dispatch("polling").await
}
fn domain_id(&self) -> &str {
"sys::process"
}
}
impl verb::Execute for Process {
type Input = ();
type Output = ProcessState;
fn required_caps(&self) -> &[Cap] {
&self.caps
}
async fn execute_action(&self, call: (Auth, &())) -> Result<ProcessState, BotError> {
call.0.check(self.required_caps())?;
self.dispatch("executing").await
}
fn domain_id(&self) -> &str {
"sys::process"
}
}
impl verb::Query for Process {
type Input = ();
type Output = ProcessState;
fn required_caps(&self) -> &[Cap] {
&self.caps
}
async fn query(&self, call: (Auth, &())) -> Result<ProcessState, BotError> {
call.0.check(self.required_caps())?;
self.dispatch("querying").await
}
fn domain_id(&self) -> &str {
"sys::process"
}
}