use std::io;
use std::path::{Path, PathBuf};
use std::process::Stdio;
use std::sync::Arc;
use std::time::Duration;
use futures_util::future::FutureExt;
use futures_util::stream::FuturesUnordered;
use futures_util::stream::StreamExt;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::process::{Child, ChildStderr, ChildStdin, ChildStdout};
use super::grep_engine;
use super::grep_engine::windows;
use super::scan;
use super::windows_line::{Executed, Written};
use super::{
CaptureBudget, Captured, KillOnDrop, Readers, RunOwner, SHELL_PIPE_READ_CAP, ShellPlatform,
ShellRunResult, Stream, Tree, Watchdog, build_program_command, build_shell_command,
};
#[derive(Debug)]
pub(super) struct Plan {
pub(super) steps: Vec<Step>,
pub(super) exe: PathBuf,
}
#[derive(Debug)]
pub(super) struct Step {
pub(super) run: Run,
pub(super) cwd: PathBuf,
pub(super) join: Join,
}
#[derive(Debug)]
pub(super) enum Run {
Shell { text: String },
Own {
args: Vec<String>,
redirects: Vec<String>,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(super) enum Join {
First,
And,
Or,
Always,
Pipe,
}
impl Join {
#[must_use]
pub(super) fn spelling(self) -> &'static str {
match self {
Join::First => "",
Join::And => "&&",
Join::Or => "||",
Join::Always => "&",
Join::Pipe => "|",
}
}
}
#[derive(Debug)]
pub(super) enum Member {
Shell {
text: String,
},
Own {
argv: Vec<String>,
redirects: Vec<String>,
cwd: PathBuf,
},
}
#[derive(Debug)]
pub(super) enum Refusal {
OwnImageInShell,
OwnImageSpelling(String),
Shape(String),
}
impl std::fmt::Display for Refusal {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Refusal::OwnImageInShell => write!(
f,
"the line also runs `mahbot` itself, which a shell does not wait for — \
such a call would come back truncated or empty"
),
Refusal::OwnImageSpelling(spelling) => write!(
f,
"a shell member runs `{spelling}`, which names `mahbot`'s file name but \
is not the program the runner spawns — and a shell does not wait for \
such an image either"
),
Refusal::Shape(reason) => write!(f, "{reason}"),
}
}
}
pub(super) fn build(
members: &[(Member, String)],
exe: &Path,
root_cwd: &Path,
) -> Result<Plan, Refusal> {
let mut steps: Vec<Step> = Vec::new();
let mut fold: Option<Fold> = None;
let mut tracked = root_cwd.to_path_buf();
let mut untracked_cd: Option<String> = None;
for (index, (member, _)) in members.iter().enumerate() {
let join = join_before(members, index)?;
let feeds_pipe = members[index].1 == "|";
if (matches!(member, Member::Own { .. }) || fold.is_none())
&& let Some(cause) = &untracked_cd
{
return Err(Refusal::Shape(cause.clone()));
}
let (args, redirects, cwd) = match member {
Member::Own {
argv,
redirects,
cwd,
} => (argv.clone(), redirects.clone(), cwd.clone()),
Member::Shell { text } => {
let words = member_words(text);
if let Some((args, redirects)) = own_call(&words, exe) {
if let Some(cause) = &untracked_cd {
return Err(Refusal::Shape(cause.clone()));
}
let cwd = tracked.clone();
(args, redirects, cwd)
} else {
if let Some(refusal) = own_image_reference(&words, text, exe) {
return Err(refusal);
}
if let Some(cause) = keyword_cd(&words, text) {
untracked_cd = Some(cause);
}
let fold_cwd = tracked.clone();
match tracked_cd(&words, text, &tracked, feeds_pipe || join == Join::Pipe) {
Ok(Some(moved)) => tracked = moved,
Ok(None) => {}
Err(cause) => untracked_cd = Some(cause),
}
if feeds_pipe {
flush_fold(&mut steps, &mut fold);
}
fold_push(&mut fold, text, join, fold_cwd);
if feeds_pipe || join == Join::Pipe {
flush_fold(&mut steps, &mut fold);
}
continue;
}
}
};
flush_fold(&mut steps, &mut fold);
redirects_supported(&redirects).map_err(Refusal::Shape)?;
tracked.clone_from(&cwd);
steps.push(Step {
run: Run::Own { args, redirects },
cwd,
join,
});
}
flush_fold(&mut steps, &mut fold);
if let Some(reason) = state_across_folds(&steps) {
return Err(Refusal::Shape(reason));
}
Ok(Plan {
steps,
exe: exe.to_path_buf(),
})
}
fn fold_push(fold: &mut Option<Fold>, text: &str, join: Join, start_cwd: PathBuf) {
match fold {
Some(open) => {
open.text.push(' ');
open.text.push_str(join.spelling());
open.text.push(' ');
open.text.push_str(text);
}
None => {
*fold = Some(Fold {
text: text.to_string(),
join,
cwd: start_cwd,
});
}
}
}
struct Fold {
text: String,
join: Join,
cwd: PathBuf,
}
fn flush_fold(steps: &mut Vec<Step>, fold: &mut Option<Fold>) {
if let Some(open) = fold.take() {
steps.push(Step {
run: Run::Shell { text: open.text },
cwd: open.cwd,
join: open.join,
});
}
}
fn join_before(members: &[(Member, String)], index: usize) -> Result<Join, Refusal> {
if index == 0 {
return Ok(Join::First);
}
match members[index - 1].1.as_str() {
"&&" => Ok(Join::And),
"||" => Ok(Join::Or),
"&" => Ok(Join::Always),
"|" => Ok(Join::Pipe),
other => Err(Refusal::Shape(format!(
"the rewritten line carries the connector `{other}`, which this \
platform's runner does not read"
))),
}
}
fn state_across_folds(steps: &[Step]) -> Option<String> {
for (index, step) in steps.iter().enumerate() {
let Run::Shell { text } = &step.run else {
continue;
};
if !steps[index + 1..]
.iter()
.any(|later| matches!(later.run, Run::Shell { .. }))
{
continue;
}
if let Some(verb) = state_verb(text) {
return Some(format!(
"the shell member `{verb}` changes state the line's later members would \
have seen, and the runner gives each member its own interpreter"
));
}
}
None
}
fn state_verb(text: &str) -> Option<String> {
const STATE_VERBS: &[&str] = &["set", "setlocal", "endlocal", "path", "pushd", "popd"];
for (member, _) in
windows::segment_command(text).expect("a fold is the members the segmenter read, rejoined")
{
let words = member_words(&member);
let Some(verb) = command_word_index(&words).map(|index| &words[index]) else {
continue;
};
if let Some(key) = windows::verb_key(command_word(&verb.value))
&& STATE_VERBS.contains(&key.as_str())
{
return Some(key);
}
}
None
}
fn redirects_supported(redirects: &[String]) -> Result<(), String> {
parse_redirects(redirects).map(|_| ())
}
#[derive(Debug)]
enum Redirect {
Stdin(String),
Stdout { target: String, append: bool },
Stderr { target: String, append: bool },
MergeIntoStdout,
}
fn parse_redirects(redirects: &[String]) -> Result<Vec<Redirect>, String> {
let mut parsed = Vec::new();
let mut words = redirects.iter();
while let Some(token) = words.next() {
let spelling = split_redirect(token).map_err(|why| unsupported(token, why))?;
let (stream, append, glued) = match spelling {
RedirectSpelling::MergeIntoStdout => {
parsed.push(Redirect::MergeIntoStdout);
continue;
}
RedirectSpelling::Operator {
stream,
append,
glued,
} => (stream, append, glued),
};
let target = if let Some(glued) = glued {
glued
} else {
let Some(word) = words.next() else {
return Err(format!(
"the runner cannot apply the redirection `{token}` to a member it \
spawns itself — its target word is missing"
));
};
word
};
if windows::has_percent_expansion(target) {
return Err(format!(
"the runner cannot apply the redirection `{token}` to a member it spawns \
itself — its target `{target}` carries a `%…%` pair, which cmd.exe \
would have expanded before the member ran"
));
}
let delivered = windows::unquote_word(target);
if windows::is_drive_relative(&delivered) {
return Err(format!(
"the runner cannot apply the redirection `{token}` to a member it spawns \
itself — its target `{delivered}` is drive-relative, naming another \
drive's current directory"
));
}
let target = target.to_string();
parsed.push(match stream {
Stream::Stdin => Redirect::Stdin(target),
Stream::Stdout => Redirect::Stdout { target, append },
Stream::Stderr => Redirect::Stderr { target, append },
});
}
Ok(parsed)
}
#[derive(Debug, Clone, Copy)]
enum RedirectSpelling<'a> {
Operator {
stream: Stream,
append: bool,
glued: Option<&'a str>,
},
MergeIntoStdout,
}
fn split_redirect(token: &str) -> Result<RedirectSpelling<'_>, Unsupported> {
if token == "2>&1" {
return Ok(RedirectSpelling::MergeIntoStdout);
}
if token.contains('&') {
return Err(Unsupported::Descriptor);
}
let (stream, append, tail) = match token.as_bytes() {
[b'1', b'>', b'>', ..] => (Stream::Stdout, true, &token[3..]),
[b'1', b'>', ..] => (Stream::Stdout, false, &token[2..]),
[b'2', b'>', b'>', ..] => (Stream::Stderr, true, &token[3..]),
[b'2', b'>', ..] => (Stream::Stderr, false, &token[2..]),
[b'>', b'>', ..] => (Stream::Stdout, true, &token[2..]),
[b'>', ..] => (Stream::Stdout, false, &token[1..]),
[b'<', ..] => (Stream::Stdin, false, &token[1..]),
_ => return Err(Unsupported::Spelling),
};
if tail.starts_with('>') {
return Err(Unsupported::Spelling);
}
Ok(RedirectSpelling::Operator {
stream,
append,
glued: (!tail.is_empty()).then_some(tail),
})
}
#[derive(Debug, Clone, Copy)]
enum Unsupported {
Descriptor,
Spelling,
}
#[must_use]
fn unsupported(token: &str, why: Unsupported) -> String {
match why {
Unsupported::Descriptor => format!(
"the runner cannot apply the redirection `{token}` to a member it spawns \
itself — the one descriptor merge it applies is `2>&1`, whose destination \
the member's stdout already has, and its combined `&>` is not applied either"
),
Unsupported::Spelling => format!(
"the runner cannot apply the redirection `{token}` to a member it spawns \
itself — it applies `2>&1`, `>`, `1>`, `>>`, `1>>`, `2>`, `2>>` and `<`, \
the target-taking ones with their target glued to the operator or as the \
next word"
),
}
}
impl Plan {
#[must_use]
pub(super) fn direct_own(&self) -> Option<(&[String], &[String], &Path)> {
let [step] = self.steps.as_slice() else {
return None;
};
match &step.run {
Run::Own { args, redirects } => Some((args, redirects, &step.cwd)),
Run::Shell { .. } => None,
}
}
}
pub(super) fn apply_redirects(
cmd: &mut tokio::process::Command,
redirects: &[String],
cwd: &Path,
) -> io::Result<AppliedRedirects> {
let parsed = parse_redirects(redirects).map_err(io::Error::other)?;
let mut applied = AppliedRedirects::default();
let mut stdout: Option<RedirectDest> = None;
for redirect in &parsed {
match redirect {
Redirect::Stdin(target) => {
cmd.stdin(stdin_stdio(target, cwd)?);
applied.streams.push(Stream::Stdin);
}
Redirect::Stdout { target, append } => {
let dest = RedirectDest::open(target, *append, cwd)?;
cmd.stdout(dest.stdio()?);
applied.streams.push(Stream::Stdout);
stdout = Some(dest);
}
Redirect::Stderr { target, append } => {
cmd.stderr(RedirectDest::open(target, *append, cwd)?.stdio()?);
applied.streams.push(Stream::Stderr);
}
Redirect::MergeIntoStdout => {
if let Some(dest) = &stdout {
cmd.stderr(dest.stdio()?);
applied.streams.push(Stream::Stderr);
} else {
applied.streams.retain(|stream| *stream != Stream::Stderr);
applied.merge_into_stdout = true;
}
}
}
}
Ok(applied)
}
#[derive(Debug, Default)]
pub(super) struct AppliedRedirects {
pub(super) streams: Vec<Stream>,
pub(super) merge_into_stdout: bool,
}
enum RedirectDest {
Null,
File(std::fs::File),
}
impl RedirectDest {
fn open(target: &str, append: bool, cwd: &Path) -> io::Result<Self> {
let target = windows::unquote_word(target);
if scan::is_null_device(&target) {
return Ok(RedirectDest::Null);
}
let path = cwd.join(&target);
Ok(RedirectDest::File(if append {
std::fs::OpenOptions::new()
.append(true)
.create(true)
.open(path)?
} else {
std::fs::File::create(path)?
}))
}
fn stdio(&self) -> io::Result<Stdio> {
match self {
RedirectDest::Null => Ok(Stdio::null()),
RedirectDest::File(file) => Ok(Stdio::from(file.try_clone()?)),
}
}
}
fn stdin_stdio(target: &str, cwd: &Path) -> io::Result<Stdio> {
let target = windows::unquote_word(target);
if scan::is_null_device(&target) {
return Ok(Stdio::null());
}
Ok(Stdio::from(std::fs::File::open(cwd.join(&target))?))
}
#[must_use]
fn own_image_in_text(text: &str, exe: &Path) -> bool {
let Some(words) = windows::tokenize(text) else {
return names_image(text.split_whitespace().next().unwrap_or(""), exe)
|| spelling_in_text(text, exe);
};
names_image_in_words(&words, exe)
}
#[must_use]
fn names_image_in_words(words: &[grep_engine::GrepWord], exe: &Path) -> bool {
let mut starts_command = true;
let mut index = 0;
while let Some(word) = words.get(index) {
if word.redirect {
index += if word.needs_target { 2 } else { 1 };
continue;
}
if starts_command {
if names_image(&word.value, exe) {
return true;
}
let launcher = windows::verb_key(command_word(&word.value))
.is_some_and(|key| LAUNCHER_VERBS.contains(&key.as_str()));
if launcher && launcher_hands_image(words, index, exe).is_some() {
return true;
}
}
starts_command = matches!(word.raw.as_str(), "&&" | "||" | "&" | "|")
|| (starts_command && is_decoration(&word.value));
index += 1;
}
false
}
fn own_image_reference(words: &[grep_engine::GrepWord], text: &str, exe: &Path) -> Option<Refusal> {
let first = command_word_index(words)?;
let word = &words[first];
if is_own_image_word(&word.value, exe) {
return Some(Refusal::OwnImageInShell);
}
if let Some(spelling) = names_image_file(&word.value, exe) {
return Some(Refusal::OwnImageSpelling(spelling.to_string()));
}
let launcher = windows::verb_key(command_word(&word.value))
.is_some_and(|key| LAUNCHER_VERBS.contains(&key.as_str()));
if launcher {
return match launcher_hands_image(words, first, exe) {
Some(HandsOn::ToCommand) => Some(Refusal::OwnImageInShell),
Some(HandsOn::Unreadable) => Some(Refusal::Shape(format!(
"the line's member `{text}` runs a launcher whose own command word this \
model cannot read while `mahbot` appears after it, and a call a shell \
might be starting is not one a shell waits for"
))),
None => None,
};
}
None
}
const LAUNCHER_VERBS: &[&str] = &["start", "call", "cmd", "powershell", "pwsh"];
fn launcher_hands_image(
words: &[grep_engine::GrepWord],
launcher: usize,
exe: &Path,
) -> Option<HandsOn> {
match launched_command(words, launcher) {
Some(at) => names_image(&words[at].value, exe).then_some(HandsOn::ToCommand),
None => words[launcher + 1..]
.iter()
.any(|word| names_image(&word.value, exe))
.then_some(HandsOn::Unreadable),
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum HandsOn {
ToCommand,
Unreadable,
}
fn launched_command(words: &[grep_engine::GrepWord], launcher: usize) -> Option<usize> {
let key = windows::verb_key(command_word(&words[launcher].value))?;
let rest = &words[launcher + 1..];
let offset = match key.as_str() {
"cmd" => {
let switch = rest.iter().position(|word| {
is_switch_word(&word.value)
&& matches!(word.value[1..].to_ascii_lowercase().as_str(), "c" | "k")
})?;
switch + 1
}
"start" => rest
.iter()
.position(|word| !is_switch_word(&word.value) && !word.raw.starts_with('"'))?,
"call" => 0,
"powershell" | "pwsh" => rest.iter().position(|word| !is_switch_word(&word.value))?,
_ => return None,
};
let at = launcher + 1 + offset;
(at < words.len()).then_some(at)
}
#[must_use]
fn is_switch_word(word: &str) -> bool {
word.starts_with(['/', '-']) && !word[1..].contains(['/', '\\'])
}
#[must_use]
fn names_image(word: &str, exe: &Path) -> bool {
is_own_image_word(word, exe) || names_image_file(word, exe).is_some()
}
#[must_use]
fn names_image_file<'a>(word: &'a str, exe: &Path) -> Option<&'a str> {
let name = image_file_name(exe)?;
let (_, file) = command_word(word).rsplit_once(['\\', '/'])?;
windows::verb_key(file)
.is_some_and(|key| Some(key) == windows::verb_key(name))
.then_some(word)
}
#[must_use]
fn is_own_image_word(word: &str, exe: &Path) -> bool {
if windows::same_spelling(Path::new(command_word(word)), exe) {
return true;
}
let Some(name) = image_file_name(exe) else {
return false;
};
windows::verb_key(command_word(word)).is_some_and(|key| Some(key) == windows::verb_key(name))
}
#[must_use]
fn command_word(word: &str) -> &str {
word.trim_start_matches(['@', '(']).trim_end_matches(')')
}
#[must_use]
fn command_word_index(words: &[grep_engine::GrepWord]) -> Option<usize> {
let mut index = 0;
while let Some(word) = words.get(index) {
if word.redirect {
index += if word.needs_target { 2 } else { 1 };
continue;
}
if is_decoration(&word.value) {
index += 1;
continue;
}
return Some(index);
}
None
}
#[must_use]
fn is_decoration(word: &str) -> bool {
!word.is_empty() && word.chars().all(|c| matches!(c, '(' | ')' | '@'))
}
#[must_use]
fn member_words(text: &str) -> Vec<grep_engine::GrepWord> {
windows::tokenize(text).expect("the segmenter's own text always tokenizes")
}
#[must_use]
fn image_file_name(exe: &Path) -> Option<&str> {
exe.to_str()?
.rsplit(['\\', '/'])
.next()
.filter(|name| !name.is_empty())
}
#[must_use]
fn spelling_in_text(text: &str, exe: &Path) -> bool {
let spelled = fold_path(&crate::util::strip_verbatim_prefix(exe).to_string_lossy());
!spelled.is_empty() && fold_path(text).contains(&spelled)
}
#[must_use]
fn fold_path(text: &str) -> String {
text.replace('/', "\\").to_lowercase()
}
pub(super) enum OwnImage {
None,
Direct(Plan),
Refused(String),
}
#[must_use]
fn own_image_plan(
executed: Executed<'_>,
written: Written<'_>,
exe: &Path,
cwd: &Path,
) -> OwnImage {
let Some(members) = windows::segment_command(executed.0) else {
return if own_image_in_text(executed.0, exe) || spelling_in_text(executed.0, exe) {
OwnImage::Refused(format!(
"the line names `mahbot` but cannot be read as cmd.exe's own chain of \
commands (`{}`), and a line the shell has to run cannot be \
waited for — run it as a line of its own",
written.0
))
} else {
OwnImage::None
};
};
if !members
.iter()
.any(|(text, _)| names_image_in_words(&member_words(text), exe))
{
return OwnImage::None;
}
let line: Vec<(Member, String)> = members
.into_iter()
.map(|(text, conn)| (Member::Shell { text }, conn))
.collect();
let root = grep_engine::canonical_or_lexical(cwd, ShellPlatform::Windows);
match build(&line, exe, &root) {
Ok(plan) => OwnImage::Direct(plan),
Err(refusal) => OwnImage::Refused(refusal.to_string()),
}
}
#[must_use]
pub(super) fn plan_for_command(
executed: Executed<'_>,
written: Written<'_>,
cwd: &Path,
) -> OwnImage {
if super::SHELL_PLATFORM != super::ShellPlatform::Windows {
return OwnImage::None;
}
let exe = match std::env::current_exe() {
Ok(exe) => exe,
Err(e) => {
tracing::debug!(
err = %e,
"no own executable path: a line naming this service's image runs as the shell runs it"
);
return OwnImage::None;
}
};
own_image_plan(executed, written, &exe, cwd)
}
fn own_call(words: &[grep_engine::GrepWord], exe: &Path) -> Option<(Vec<String>, Vec<String>)> {
let first = command_word_index(words)?;
if !is_own_image_word(&words[first].value, exe) {
return None;
}
let mut argv = Vec::new();
let mut redirects = Vec::new();
let mut index = 0;
while let Some(word) = words.get(index) {
if word.redirect {
redirects.push(word.raw.clone());
if word.needs_target
&& let Some(target) = words.get(index + 1)
{
redirects.push(target.raw.clone());
index += 1;
}
} else if index > first {
if windows::has_percent_expansion(&word.raw) {
return None;
}
argv.push(word.value.clone());
}
index += 1;
}
Some((argv, redirects))
}
fn tracked_cd(
words: &[grep_engine::GrepWord],
text: &str,
cwd: &Path,
piped: bool,
) -> Result<Option<PathBuf>, String> {
let platform = ShellPlatform::Windows;
let command = command_word_index(words).map(|at| (at, words[at].value.as_str()));
let (leads, verb) = match command {
Some((at, verb)) => (at == 0, verb),
None => (false, text.split_whitespace().next().unwrap_or("")),
};
if piped && names_cd_family(verb) {
return Err(untrackable_cd(
text,
"runs on one side of a pipe, which the shell gives an interpreter of its own",
));
}
if grep_engine::is_cd_segment(verb, platform) {
if !leads {
return Err(untrackable_cd(
text,
"cannot be read where the tracking starts",
));
}
return grep_engine::resolve_cd(text, cwd, Path::new(""), platform)
.map(Some)
.map_err(|_| untrackable_cd(text, "is a form the runner cannot track"));
}
if grep_engine::is_cd_spelling(verb, platform) {
return Err(untrackable_cd(
text,
"is a spelling of cmd.exe's own directory change that the runner cannot track",
));
}
Ok(None)
}
#[must_use]
fn names_cd_family(word: &str) -> bool {
let platform = ShellPlatform::Windows;
grep_engine::is_cd_segment(word, platform) || grep_engine::is_cd_spelling(word, platform)
}
fn keyword_cd(words: &[grep_engine::GrepWord], text: &str) -> Option<String> {
const KEYWORD_FORMS: &[&str] = &["if", "else", "for"];
let at = command_word_index(words)?;
let key = windows::verb_key(&words[at].value);
if !key
.as_deref()
.is_some_and(|key| KEYWORD_FORMS.contains(&key))
{
return None;
}
words[at + 1..]
.iter()
.any(|word| names_cd_family(command_word(&word.value)))
.then(|| {
format!(
"the line's member `{text}` runs a directory change inside one of the \
shell's own keyword forms, whose condition decides whether it runs at \
all, so the directory the line is left in is not something this model \
can read"
)
})
}
fn untrackable_cd(text: &str, why: &str) -> String {
format!(
"the line's `cd` member `{text}` {why}, so a step of its own after it cannot be \
run in the directory the shell would have used"
)
}
#[must_use]
pub(super) fn refusal_message(cause: &str) -> String {
format!(
"{}{cause}. This is not the command's output.",
super::REFUSAL_FRAME
)
}
#[must_use]
fn runs_after(join: Join, previous: Option<i32>) -> bool {
match join {
Join::And => previous == Some(0),
Join::Or => previous != Some(0),
Join::First | Join::Always | Join::Pipe => true,
}
}
#[must_use]
fn pipe_groups(steps: &[Step]) -> Vec<std::ops::Range<usize>> {
let mut groups = Vec::new();
let mut start = 0;
for (index, step) in steps.iter().enumerate() {
if index > 0 && step.join != Join::Pipe {
groups.push(start..index);
start = index;
}
}
if !steps.is_empty() {
groups.push(start..steps.len());
}
groups
}
struct MemberRun {
child: Child,
pid: u32,
status: Option<std::process::ExitStatus>,
}
struct CaptureBudgets {
stdout: CaptureBudget,
stderr: CaptureBudget,
}
impl CaptureBudgets {
fn new() -> Self {
Self {
stdout: CaptureBudget::new(SHELL_PIPE_READ_CAP),
stderr: CaptureBudget::new(SHELL_PIPE_READ_CAP),
}
}
}
pub(super) async fn run(
plan: &Plan,
timeout: Duration,
drain_limit: Duration,
memory_limit: Option<u64>,
owner: RunOwner,
) -> ShellRunResult {
let start = std::time::Instant::now();
let mut tree = Tree::new(owner);
let mut kill_guard = KillOnDrop::new(tree.clone());
let cancel = tokio_util::sync::CancellationToken::new();
let budgets = CaptureBudgets::new();
let mut readers = Readers::default();
let mut watchdog = Watchdog::new(timeout);
let mut previous: Option<i32> = None;
let mut reported: Option<std::process::ExitStatus> = None;
for group in pipe_groups(&plan.steps) {
if !runs_after(plan.steps[group.start].join, previous) {
continue;
}
let mut members =
match spawn_group(plan, group, &mut tree, &mut readers, &budgets, &cancel).await {
Ok(members) => members,
Err(e) => {
cancel.cancel();
kill_guard.disarm();
return ShellRunResult::SpawnFailed(e);
}
};
let named_pid = members.first().map(|member| member.pid);
match wait_group(&mut members, &mut watchdog, memory_limit).await {
GroupEnd::Reaped => {
let last = members.last().and_then(|member| member.status);
previous = last.and_then(|status| status.code());
reported = last;
}
GroupEnd::TimedOut => {
end_group(&tree, &mut members, &mut kill_guard).await;
let (stdout, stderr) = collect_stopped(&cancel, &mut readers).await;
return ShellRunResult::TimedOut {
stdout,
stderr,
pid: named_pid,
elapsed: start.elapsed(),
};
}
GroupEnd::MemoryExceeded { used, limit } => {
end_group(&tree, &mut members, &mut kill_guard).await;
let (stdout, stderr) = collect_stopped(&cancel, &mut readers).await;
return ShellRunResult::MemoryExceeded {
stdout,
stderr,
pid: named_pid,
elapsed: start.elapsed(),
used,
limit,
};
}
GroupEnd::Failed(e) => {
return ShellRunResult::SpawnFailed(e);
}
}
}
kill_guard.disarm();
let drained = readers.drain(drain_limit).await;
let Some(status) = reported else {
return ShellRunResult::SpawnFailed(io::Error::other("the plan has no step to run"));
};
if drained {
tree.retain_after_completion();
let (stdout, stderr) = readers.collect().await;
return ShellRunResult::Completed {
stdout,
stderr,
status,
elapsed: start.elapsed(),
};
}
tree.retain_after_completion();
let (stdout, stderr) = readers.snapshot_and_detach();
ShellRunResult::ExitedWithLeftovers {
stdout,
stderr,
status,
scope: None,
elapsed: start.elapsed(),
}
}
async fn spawn_group(
plan: &Plan,
group: std::ops::Range<usize>,
tree: &mut Tree,
readers: &mut Readers,
budgets: &CaptureBudgets,
cancel: &tokio_util::sync::CancellationToken,
) -> io::Result<Vec<MemberRun>> {
let mut members: Vec<MemberRun> = Vec::with_capacity(group.len());
let mut feed: Option<Pipe> = None;
for index in group.clone() {
let step = &plan.steps[index];
let mut cmd = match &step.run {
Run::Shell { text } => build_shell_command(text, &step.cwd),
Run::Own { args, .. } => build_program_command(&plan.exe, args, &step.cwd),
};
let applied = match &step.run {
Run::Own { redirects, .. } => match apply_redirects(&mut cmd, redirects, &step.cwd) {
Ok(applied) => applied,
Err(e) => {
end_group_members(tree, &mut members).await;
return Err(e);
}
},
Run::Shell { .. } => AppliedRedirects::default(),
};
let fed_by_pipe = feed.is_some();
if fed_by_pipe && !applied.streams.contains(&Stream::Stdin) {
cmd.stdin(Stdio::piped());
}
if !applied.streams.contains(&Stream::Stdout) {
cmd.stdout(Stdio::piped());
}
if !applied.streams.contains(&Stream::Stderr) {
cmd.stderr(Stdio::piped());
}
let mut child = match cmd.spawn() {
Ok(child) => child,
Err(e) => {
end_group_members(tree, &mut members).await;
return Err(e);
}
};
let pid = child
.id()
.expect("a freshly spawned child still has its pid");
tree.attach(pid);
match feed.take() {
None => {}
Some(pipe) => match child.stdin.take() {
Some(write_end) => connect_pipe(pipe, write_end),
None => drop(pipe),
},
}
let piped = index + 1 < group.end;
let mut stdout = child.stdout.take();
let mut stderr = child.stderr.take();
if piped {
feed = Some(Pipe {
stdout: stdout.take(),
stderr: if applied.merge_into_stdout {
stderr.take()
} else {
None
},
});
}
if let Some(pipe) = stdout {
readers.push(
Captured::Stdout,
pipe,
budgets.stdout.clone(),
cancel.clone(),
);
}
if let Some(pipe) = stderr {
let (stream, budget) = if applied.merge_into_stdout {
(Captured::Stdout, budgets.stdout.clone())
} else {
(Captured::Stderr, budgets.stderr.clone())
};
readers.push(stream, pipe, budget, cancel.clone());
}
members.push(MemberRun {
child,
pid,
status: None,
});
}
Ok(members)
}
struct Pipe {
stdout: Option<ChildStdout>,
stderr: Option<ChildStderr>,
}
fn connect_pipe(pipe: Pipe, write_end: ChildStdin) {
match (pipe.stdout, pipe.stderr) {
(Some(stdout), Some(stderr)) => {
let shared = SharedStdin::new(write_end);
spawn_merged_copy(stdout, shared.clone());
spawn_merged_copy(stderr, shared);
}
(Some(stdout), None) => spawn_pipe_copy(stdout, write_end),
(None, Some(stderr)) => spawn_pipe_copy(stderr, write_end),
(None, None) => drop(write_end),
}
}
#[derive(Clone)]
struct SharedStdin(Arc<tokio::sync::Mutex<ChildStdin>>);
impl SharedStdin {
fn new(write_end: ChildStdin) -> Self {
SharedStdin(Arc::new(tokio::sync::Mutex::new(write_end)))
}
}
fn spawn_pipe_copy(
mut from: impl tokio::io::AsyncRead + Unpin + Send + 'static,
mut to: ChildStdin,
) {
tokio::spawn(async move {
let _ = tokio::io::copy(&mut from, &mut to).await;
});
}
fn spawn_merged_copy(
mut from: impl tokio::io::AsyncRead + Unpin + Send + 'static,
to: SharedStdin,
) {
tokio::spawn(async move {
let mut buf = [0_u8; MERGED_COPY_CHUNK];
loop {
match from.read(&mut buf).await {
Ok(0) | Err(_) => break,
Ok(read) => {
let mut write_end = to.0.lock().await;
if write_end.write_all(&buf[..read]).await.is_err() {
break;
}
}
}
}
});
}
const MERGED_COPY_CHUNK: usize = 8 * 1024;
enum GroupEnd {
Reaped,
TimedOut,
MemoryExceeded { used: u64, limit: u64 },
Failed(io::Error),
}
async fn wait_group(
members: &mut [MemberRun],
watchdog: &mut Watchdog,
memory_limit: Option<u64>,
) -> GroupEnd {
debug_assert!(
members.iter().all(|member| member.status.is_none()),
"a member is waited for as soon as it is spawned, never twice"
);
let mut statuses: Vec<Option<std::process::ExitStatus>> = vec![None; members.len()];
let mut live: Vec<(usize, u32)> = members
.iter()
.enumerate()
.map(|(index, member)| (index, member.pid))
.collect();
let end = {
let mut pending: FuturesUnordered<_> = members
.iter_mut()
.enumerate()
.map(|(index, member)| member.child.wait().map(move |result| (index, result)))
.collect();
loop {
tokio::select! {
biased;
Some((index, result)) = pending.next() => {
match result {
Ok(status) => {
statuses[index] = Some(status);
live.retain(|(at, _)| *at != index);
}
Err(e) => break GroupEnd::Failed(e),
}
if statuses.iter().all(Option::is_some) {
break GroupEnd::Reaped;
}
}
() = &mut watchdog.deadline => break GroupEnd::TimedOut,
_ = watchdog.sample.tick() => {
let pids: Vec<u32> = live.iter().map(|(_, pid)| *pid).collect();
if let Some(limit) = memory_limit
&& let Some(sample) = watchdog.measure(&pids, limit)
{
break GroupEnd::MemoryExceeded {
used: sample.used,
limit: sample.limit,
};
}
}
}
}
};
match end {
GroupEnd::Reaped => {
for (member, status) in members.iter_mut().zip(statuses) {
member.status = status;
}
GroupEnd::Reaped
}
other => other,
}
}
async fn end_group(tree: &Tree, members: &mut [MemberRun], kill_guard: &mut KillOnDrop) {
end_group_members(tree, members).await;
kill_guard.disarm();
}
async fn end_group_members(tree: &Tree, members: &mut [MemberRun]) {
for member in members.iter_mut() {
let _ = member.child.start_kill();
}
tree.terminate();
for member in members.iter_mut() {
let _ = member.child.wait().await;
}
}
async fn collect_stopped(
cancel: &tokio_util::sync::CancellationToken,
readers: &mut Readers,
) -> (Vec<u8>, Vec<u8>) {
cancel.cancel();
readers.collect().await
}
#[cfg(test)]
mod tests {
use super::*;
use std::ops::Range;
const EXE: &str = r"C:\Program Files\MahBot\mahbot.exe";
const ROOT: &str = r"C:\ws";
fn shell(text: &str, conn: &str) -> (Member, String) {
(
Member::Shell {
text: text.to_string(),
},
conn.to_string(),
)
}
fn own(argv: &[&str], redirects: &[&str], cwd: &str, conn: &str) -> (Member, String) {
(
Member::Own {
argv: argv.iter().map(|word| (*word).to_string()).collect(),
redirects: redirects.iter().map(|word| (*word).to_string()).collect(),
cwd: PathBuf::from(cwd),
},
conn.to_string(),
)
}
fn served(cwd: &str, conn: &str) -> (Member, String) {
own(
&["__grep-engine", "--spec-file", r"C:\tmp\spec.json"],
&[],
cwd,
conn,
)
}
fn plan(members: &[(Member, String)]) -> Plan {
build_at(ROOT, members).expect("a decomposable line")
}
fn build_at(root: &str, members: &[(Member, String)]) -> Result<Plan, Refusal> {
build(members, Path::new(EXE), &tracked_root_of(root))
}
fn tracked_root_of(root: &str) -> PathBuf {
grep_engine::canonical_or_lexical(Path::new(root), ShellPlatform::Windows)
}
fn direct_plan(command: &str) -> Plan {
match own_image_plan(
Executed(command),
Written(command),
Path::new(EXE),
Path::new(ROOT),
) {
OwnImage::Direct(plan) => plan,
OwnImage::Refused(cause) => panic!("{command}: refused: {cause}"),
OwnImage::None => panic!("{command}: not read as an own-image line"),
}
}
fn refusal_cause(command: &str) -> String {
match own_image_plan(
Executed(command),
Written(command),
Path::new(EXE),
Path::new(ROOT),
) {
OwnImage::Refused(cause) => cause,
OwnImage::Direct(_) | OwnImage::None => {
panic!("{command}: expected a refusal")
}
}
}
fn shell_text(step: &Step) -> &str {
match &step.run {
Run::Shell { text } => text,
Run::Own { .. } => panic!("expected a shell step"),
}
}
fn step_cwds(plan: &Plan) -> Vec<&Path> {
plan.steps.iter().map(|step| step.cwd.as_path()).collect()
}
#[test]
fn a_lone_search_is_one_own_step() {
let plan = plan(&[served(ROOT, "")]);
assert_eq!(plan.steps.len(), 1);
assert_eq!(plan.steps[0].join, Join::First);
assert_eq!(plan.steps[0].cwd, PathBuf::from(ROOT));
let Run::Own { args, .. } = &plan.steps[0].run else {
panic!("the served member must be an own step");
};
assert_eq!(args[0], "__grep-engine");
}
#[test]
fn a_pipelined_search_is_one_pipe_group() {
let plan = plan(&[served(ROOT, "|"), shell("head -3", "")]);
assert_eq!(plan.steps.len(), 2);
assert_eq!(plan.steps[1].join, Join::Pipe);
assert_eq!(pipe_groups(&plan.steps), vec![0..2]);
}
#[test]
fn a_cd_before_a_search_is_a_fold_and_the_tracked_cwd() {
let plan = plan(&[shell("cd src", "&&"), served(r"C:\ws\src", "")]);
assert_eq!(plan.steps.len(), 2);
assert_eq!(shell_text(&plan.steps[0]), "cd src");
assert_eq!(plan.steps[0].cwd, tracked_root_of(ROOT));
assert_eq!(plan.steps[1].join, Join::And);
assert_eq!(plan.steps[1].cwd, PathBuf::from(r"C:\ws\src"));
}
#[test]
fn a_cd_search_and_tail_is_a_fold_then_a_pipe_group() {
let plan = plan(&[
shell("cd src", "&&"),
served(r"C:\ws\src", "|"),
shell("head -3", ""),
]);
assert_eq!(plan.steps.len(), 3);
assert_eq!(plan.steps[2].join, Join::Pipe);
assert_eq!(pipe_groups(&plan.steps), vec![0..1, 1..3]);
assert_eq!(shell_text(&plan.steps[0]), "cd src");
assert_eq!(plan.steps[2].cwd, plan.steps[1].cwd);
}
#[test]
fn a_member_that_feeds_a_pipe_is_a_step_of_its_own() {
let root = tracked_root_of(ROOT);
let plan = plan(&[
shell("cd src", "&&"),
shell("type f.txt", "|"),
served(r"C:\ws\src", ""),
]);
assert_eq!(plan.steps.len(), 3);
assert_eq!(pipe_groups(&plan.steps), vec![0..1, 1..3]);
assert_eq!(shell_text(&plan.steps[0]), "cd src");
assert_eq!(plan.steps[0].join, Join::First);
assert_eq!(shell_text(&plan.steps[1]), "type f.txt");
assert_eq!(plan.steps[1].join, Join::And);
assert!(matches!(plan.steps[2].run, Run::Own { .. }));
assert_eq!(plan.steps[2].join, Join::Pipe);
assert_eq!(
step_cwds(&plan),
vec![
root.as_path(),
root.join("src").as_path(),
Path::new(r"C:\ws\src")
]
);
}
#[test]
fn the_pipeline_consumer_is_a_step_of_its_own() {
let piped = plan(&[served(ROOT, "|"), shell("head -3", "&&"), shell("cat", "")]);
assert_eq!(piped.steps.len(), 3);
assert_eq!(pipe_groups(&piped.steps), vec![0..2, 2..3]);
assert!(matches!(piped.steps[0].run, Run::Own { .. }));
assert_eq!(shell_text(&piped.steps[1]), "head -3");
assert_eq!(piped.steps[1].join, Join::Pipe);
assert_eq!(shell_text(&piped.steps[2]), "cat");
assert_eq!(piped.steps[2].join, Join::And);
let tailed = plan(&[
served(ROOT, "|"),
shell("head -3", "|"),
shell("sort", "&&"),
shell("echo done", ""),
]);
assert_eq!(tailed.steps.len(), 4);
assert_eq!(pipe_groups(&tailed.steps), vec![0..3, 3..4]);
assert_eq!(tailed.steps[3].join, Join::And);
let bare = plan(&[shell("type a.txt", "|"), shell("sort", "")]);
assert_eq!(bare.steps.len(), 2);
assert_eq!(shell_text(&bare.steps[0]), "type a.txt");
assert_eq!(shell_text(&bare.steps[1]), "sort");
assert_eq!(pipe_groups(&bare.steps), vec![0..2]);
}
#[test]
fn two_searches_are_two_sequenced_own_steps() {
let plan = plan(&[served(ROOT, "&&"), served(ROOT, "")]);
assert_eq!(plan.steps.len(), 2);
assert_eq!(plan.steps[1].join, Join::And);
assert!(runs_after(plan.steps[1].join, Some(0)));
assert!(!runs_after(plan.steps[1].join, Some(1)));
}
#[test]
fn a_redirect_rides_its_own_step() {
let plan = plan(&[own(
&["__grep-engine", "--spec-file", r"C:\tmp\spec.json"],
&[">", "out.txt"],
ROOT,
"",
)]);
let Run::Own { redirects, .. } = &plan.steps[0].run else {
panic!("expected an own step");
};
assert_eq!(redirects, &[">".to_string(), "out.txt".to_string()]);
}
#[test]
fn a_fold_joins_its_members_with_their_own_connectors() {
let plan = plan(&[
shell("echo one", "&"),
shell("echo two", "&&"),
shell("echo three", ""),
]);
assert_eq!(plan.steps.len(), 1);
assert_eq!(
shell_text(&plan.steps[0]),
"echo one & echo two && echo three"
);
}
#[test]
fn an_own_step_breaks_the_fold_and_carries_the_cwd() {
let plan = plan(&[
shell("cd src", "&&"),
served(r"C:\ws\src", "&&"),
shell("type note.txt", ""),
]);
assert_eq!(plan.steps.len(), 3);
assert_eq!(shell_text(&plan.steps[0]), "cd src");
assert_eq!(shell_text(&plan.steps[2]), "type note.txt");
assert_eq!(plan.steps[2].cwd, PathBuf::from(r"C:\ws\src"));
}
#[test]
fn a_fold_member_that_calls_the_image_becomes_an_own_step() {
for (members, args) in [
(vec![shell("mahbot -V", "")], ["-V"].as_slice()),
(
vec![shell("echo hi", "&&"), shell("mahbot -V", "")],
["-V"].as_slice(),
),
(
vec![shell(r#""C:\Program Files\MahBot\mahbot.exe" debug"#, "")],
["debug"].as_slice(),
),
(vec![shell("@mahbot -V", "")], ["-V"].as_slice()),
(vec![shell("@ mahbot -V", "")], ["-V"].as_slice()),
(vec![shell("> out.txt mahbot -V", "")], ["-V"].as_slice()),
(
vec![
shell("grep -n x f.txt", "&&"),
shell("2> err.txt mahbot debug", ""),
],
["debug"].as_slice(),
),
] {
let plan = build_at(ROOT, &members).expect("a call of the image is the runner's");
let own: Vec<&Step> = plan
.steps
.iter()
.filter(|step| matches!(step.run, Run::Own { .. }))
.collect();
assert_eq!(own.len(), 1, "{members:?}");
let Run::Own { args: argv, .. } = &own[0].run else {
unreachable!("filtered to the own step");
};
assert_eq!(argv, args, "{members:?}");
}
let redirected = plan(&[shell("> out.txt mahbot -V", "")]);
let Run::Own { redirects, .. } = &redirected.steps[0].run else {
panic!("a call of the image is an own step");
};
assert_eq!(redirects, &[">".to_string(), "out.txt".to_string()]);
for text in [
"echo mahbot",
"type mahbot.exe",
"echo (mahbot -V)",
"echo @mahbot -V",
"echo hi > mahbot.exe",
"echo hi < mahbot -V.txt",
"cmd /c dir mahbot",
"cmd -c dir mahbot",
r#"cmd /c "dir mahbot""#,
r#"powershell -c "Get-Process mahbot""#,
r"start /d C:\ws notepad mahbot",
"mahbot.bat -V",
"mahbot.com -V",
] {
let plan = plan(&[shell(text, "&&"), served(ROOT, "")]);
assert_eq!(plan.steps.len(), 2, "{text}");
}
}
#[test]
fn a_fold_naming_the_image_in_a_shape_the_runner_cannot_run_refuses_the_line() {
for (members, spelling) in [
(vec![shell("start mahbot chrome", "")], false),
(vec![shell("call mahbot.exe -V", "")], false),
(vec![shell("cmd /c mahbot bench-openrouter", "")], false),
(vec![shell("cmd -c mahbot -V", "")], false),
(vec![shell(r".\mahbot.exe -V", "")], true),
(vec![shell(r"target\debug\mahbot.exe -V", "")], true),
(
vec![shell("echo hi", "&&"), shell(r".\mahbot.exe -V", "")],
true,
),
] {
let refusal = build_at(ROOT, &members).expect_err("a shape the runner cannot run");
match refusal {
Refusal::OwnImageSpelling(spelled) => {
assert!(
spelling,
"{members:?}: unexpected spelling cause: {spelled}"
);
}
Refusal::OwnImageInShell => {
assert!(!spelling, "{members:?}: expected the spelling cause");
}
Refusal::Shape(reason) => {
panic!("{members:?}: unexpected cause: {reason}");
}
}
}
for text in ["cmd /x dir mahbot", "cmd dir c mahbot", "cmd dir k mahbot"] {
let refusal = build_at(ROOT, &[shell(text, "")])
.expect_err("a launcher whose command word cannot be read");
let Refusal::Shape(reason) = refusal else {
panic!("{text}: expected the shape refusal: {refusal}");
};
assert!(
reason.contains("command word this model cannot read"),
"{text}: {reason}"
);
assert!(
!reason.contains("runs `mahbot` itself"),
"{text}: the cause must not claim what the reading did not see: {reason}"
);
}
}
#[test]
fn an_unsupported_redirect_refuses_the_line() {
let refusal = build_at(
ROOT,
&[own(
&["__grep-engine", "--spec-file", r"C:\tmp\spec.json"],
&[">&2"],
ROOT,
"",
)],
)
.expect_err("a dup the runner has no destination for is refused");
let Refusal::Shape(reason) = refusal else {
panic!("expected the redirect refusal");
};
assert!(reason.contains("`>&2`"), "{reason}");
assert!(
reason.contains("descriptor merge it applies is `2>&1`"),
"{reason}"
);
}
#[test]
fn a_state_changing_fold_before_another_fold_refuses_the_line() {
for member in ["set MAHBOT_TEST=1", "path C:\\other", "pushd src"] {
let refusal = build_at(
ROOT,
&[
shell(member, "&&"),
served(ROOT, "&&"),
shell("echo done", ""),
],
)
.expect_err("the state dies with the fold");
assert!(matches!(refusal, Refusal::Shape(_)), "{member}: {refusal}");
}
plan(&[shell("set MAHBOT_TEST=1", "&&"), served(ROOT, "")]);
plan(&[shell("cd src", "&&"), served(r"C:\ws\src", "")]);
plan(&[shell("pushd src", "&&"), served(r"C:\ws\src", "")]);
plan(&[
shell("color 0A", "&&"),
served(ROOT, "&&"),
shell("echo done", ""),
]);
}
#[test]
fn a_cd_the_tracking_cannot_follow_refuses_the_line() {
for member in [
"cd..",
r"cd\Users",
"cd/d C:\\ws",
"@cd C:\\ws",
"chdir..",
"> out.txt cd src",
"cd C:ws",
] {
let refusal = build_at(ROOT, &[shell(member, "&&"), served(ROOT, "")])
.expect_err("a directory change the tracking cannot follow");
let Refusal::Shape(reason) = refusal else {
panic!("{member}: expected the shape refusal: {refusal}");
};
assert!(reason.contains("`cd` member"), "{member}: {reason}");
assert!(reason.contains(member), "{member}: {reason}");
}
for members in [
vec![shell("cd src", "|"), served(ROOT, "")],
vec![shell("dir", "|"), shell("cd src", "&&"), served(ROOT, "")],
] {
let refusal = build_at(ROOT, &members).expect_err("a `cd` a pipeline owns");
let Refusal::Shape(reason) = refusal else {
panic!("expected the shape refusal: {refusal}");
};
assert!(reason.contains("side of a pipe"), "{reason}");
}
for members in [
vec![served(ROOT, "&&"), shell("cd..", "")],
vec![served(ROOT, "|"), shell("cd..", "")],
vec![shell("cd..", "&&"), shell("echo done", "")],
vec![
shell("cd..", "&&"),
shell("echo hi", "&&"),
shell("echo done", ""),
],
vec![shell("if exist src cd src", "&&"), shell("echo done", "")],
] {
assert!(!plan(&members).steps.is_empty(), "{members:?}");
}
for members in [
vec![shell("cd..", "&&"), shell("mahbot -V", "")],
vec![
shell("echo hi", "&&"),
shell("cd..", "&&"),
shell("mahbot -V", ""),
],
] {
assert!(build_at(ROOT, &members).is_err(), "{members:?}");
}
for member in [
"if exist src cd src",
"if not exist x cd sub",
"else cd sub",
"for %f in (*) do cd sub",
"if exist x (cd sub)",
] {
let refusal = build_at(ROOT, &[shell(member, "&&"), served(ROOT, "")])
.expect_err("a `cd` a keyword form owns");
let Refusal::Shape(reason) = refusal else {
panic!("{member}: expected the shape refusal: {refusal}");
};
assert!(reason.contains("keyword form"), "{member}: {reason}");
}
for (member, tracked) in [
("cd src", r"C:\ws\src"),
("chdir src", r"C:\ws\src"),
("CD \\ws", r"C:\ws"),
] {
let plan = plan(&[shell(member, "&&"), served(tracked, "")]);
assert_eq!(plan.steps[1].cwd, PathBuf::from(tracked), "{member}");
}
}
#[test]
fn redirects_supported_is_the_closed_set_it_documents() {
for supported in [
vec![">", "out.txt"],
vec!["1>", "out.txt"],
vec![">>", "out.txt"],
vec!["1>>", "out.txt"],
vec!["2>", "err.txt"],
vec!["2>>", "err.txt"],
vec!["<", "in.txt"],
vec![">", "nul"],
vec!["2>", "NUL"],
vec![">", "out.txt", "2>", "err.txt", "<", "in.txt"],
vec![">out.txt"],
vec!["1>>log.txt"],
vec!["2>err.txt"],
vec!["2>>err.txt"],
vec!["<in.txt"],
vec!["1>0"],
vec![">", r"C:\out.txt"],
vec![">", r#""C:\out.txt""#],
vec!["2>&1"],
vec![">", "out.txt", "2>&1"],
vec!["2>&1", ">", "out.txt"],
] {
let redirects: Vec<String> = supported.iter().map(|w| (*w).to_string()).collect();
assert_eq!(redirects_supported(&redirects), Ok(()), "{supported:?}");
}
for refused in [
vec!["1>&2"],
vec![">&2"],
vec![">>"],
vec!["<"],
vec!["10>", "out.txt"],
vec![">"],
vec!["<&2"],
vec!["<>file"],
vec![">", r"%TEMP%\out.txt"],
vec![r">%TEMP%\out.txt"],
vec![">", "C:out.txt"],
vec![">C:out.txt"],
vec![">", r#""C:out.txt""#],
vec![r#">"C:out.txt""#],
vec![">", r#""C:""out.txt""#],
vec!["<", "C:in.txt"],
vec!["3>&1"],
] {
let spellings: Vec<String> = refused.iter().map(|w| (*w).to_string()).collect();
let reason = redirects_supported(&spellings).expect_err("outside the set");
assert!(
reason.contains(&format!("`{}`", refused[0])),
"{refused:?}: {reason}"
);
}
}
#[test]
fn pipe_groups_split_on_every_join_but_a_pipe() {
let plan = plan(&[
shell("echo one", "&&"),
served(ROOT, "|"),
shell("head -1", "&&"),
served(ROOT, "||"),
shell("tail -1", ""),
]);
let groups: Vec<Range<usize>> = pipe_groups(&plan.steps);
assert_eq!(groups, vec![0..1, 1..3, 3..4, 4..5]);
assert_eq!(pipe_groups(&[]), Vec::<Range<usize>>::new());
}
#[test]
fn own_image_in_text_reads_the_command_position() {
let exe = Path::new(EXE);
for names in [
"mahbot -V",
"MAHBOT debug --db board",
r#""C:\Program Files\MahBot\mahbot.exe" chrome x"#,
r#""\\?\C:\Program Files\MahBot\mahbot.exe" -V"#,
r#""C:/Program Files/MahBot/mahbot.exe" -V"#,
"start mahbot -V",
"call mahbot.exe -V",
"cmd /c mahbot debug",
"powershell -c mahbot",
"pwsh mahbot.exe -V",
r".\mahbot.exe -V",
r"target\debug\mahbot.exe -V",
r"C:\somewhere\else\mahbot -V",
r"start .\mahbot.exe",
"cd sub && mahbot -V",
"type f.txt | mahbot -V",
"echo hi & mahbot -V",
"@mahbot -V",
"(mahbot -V)",
"( mahbot -V )",
"(cd x & mahbot -V)",
] {
assert!(own_image_in_text(names, exe), "{names}");
}
for text in [
"grep -rn mahbot .",
"grep -n mahbot.exe f.txt",
"echo mahbot",
"type mahbot.exe",
"copy mahbot.exe back\\",
"where mahbot",
r"C:\Program Files\MahBot\mahbot.exe chrome x",
r".\grep.exe -V",
r"target\debug\other.exe -V",
r#"echo "&&" mahbot"#,
"echo ( mahbot -V )",
"echo @ mahbot -V",
"cmd /c dir mahbot",
"cmd /c rd /s /q mahbot",
"cmd /d /c echo mahbot",
"powershell Get-Process mahbot",
"powershell -c Get-Process mahbot",
"start /wait notepad",
r#"start "" notepad mahbot"#,
] {
assert!(!own_image_in_text(text, exe), "{text}");
}
}
#[test]
fn own_image_plan_reads_a_lone_call() {
let exe = Path::new(EXE);
let cwd = Path::new(ROOT);
let OwnImage::Direct(direct) = own_image_plan(
Executed(r#""C:\Program Files\MahBot\mahbot.exe" debug --db board "select 1""#),
Written(r#""C:\Program Files\MahBot\mahbot.exe" debug --db board "select 1""#),
exe,
cwd,
) else {
panic!("a lone own-image call is run by the runner");
};
let Some((args, redirects, step_cwd)) = direct.direct_own() else {
panic!("a direct plan is one own step");
};
assert_eq!(args, &["debug", "--db", "board", "select 1"]);
assert!(redirects.is_empty());
assert_eq!(step_cwd, tracked_root_of(ROOT));
let OwnImage::Direct(redirected) = own_image_plan(
Executed(r#""C:\Program Files\MahBot\mahbot.exe" -V > out.txt"#),
Written(r#""C:\Program Files\MahBot\mahbot.exe" -V > out.txt"#),
exe,
cwd,
) else {
panic!("a redirected lone call is run by the runner");
};
let Some((args, redirects, _)) = redirected.direct_own() else {
panic!("a direct plan is one own step");
};
assert_eq!(args, &["-V"]);
assert_eq!(redirects, &[">".to_string(), "out.txt".to_string()]);
let OwnImage::Direct(glued) = own_image_plan(
Executed(r#""C:\Program Files\MahBot\mahbot.exe" -V >out.txt"#),
Written(r#""C:\Program Files\MahBot\mahbot.exe" -V >out.txt"#),
exe,
cwd,
) else {
panic!("a redirected lone call is run by the runner");
};
let Some((args, redirects, _)) = glued.direct_own() else {
panic!("a direct plan is one own step");
};
assert_eq!(args, &["-V"]);
assert_eq!(redirects, &[">out.txt".to_string()]);
}
#[test]
fn own_image_plan_decomposes_a_compound_line() {
let root = tracked_root_of(ROOT);
let sub = root.join("sub");
let plan = direct_plan(&format!(r#"cd sub && "{EXE}" --version && echo done"#));
assert_eq!(plan.steps.len(), 3);
assert_eq!(shell_text(&plan.steps[0]), "cd sub");
assert_eq!(shell_text(&plan.steps[2]), "echo done");
assert_eq!(plan.steps[2].join, Join::And);
assert_eq!(
step_cwds(&plan),
vec![root.as_path(), sub.as_path(), sub.as_path()]
);
let Run::Own { args, .. } = &plan.steps[1].run else {
panic!("the call of the image is an own step");
};
assert_eq!(args, &["--version"]);
let plan = direct_plan(&format!(r#""{EXE}" -V | head -3"#));
assert_eq!(plan.steps.len(), 2);
assert_eq!(plan.steps[1].join, Join::Pipe);
assert_eq!(pipe_groups(&plan.steps), vec![0..2]);
assert_eq!(shell_text(&plan.steps[1]), "head -3");
assert_eq!(step_cwds(&plan), vec![root.as_path(), root.as_path()]);
let plan = direct_plan(&format!(r#""{EXE}" debug --db board "SELECT 1" > out.txt"#));
let Some((args, redirects, _)) = plan.direct_own() else {
panic!("a lone call is one own step");
};
assert_eq!(args, &["debug", "--db", "board", "SELECT 1"]);
assert_eq!(redirects, &[">".to_string(), "out.txt".to_string()]);
let plan = direct_plan(&format!(r#""{EXE}" -V 2>&1 | head -2"#));
assert_eq!(plan.steps.len(), 2);
assert_eq!(plan.steps[1].join, Join::Pipe);
let Run::Own { args, redirects } = &plan.steps[0].run else {
panic!("the call of the image is an own step");
};
assert_eq!(args, &["-V"]);
assert_eq!(redirects, &["2>&1".to_string()]);
}
#[test]
fn an_unrepresentable_own_image_line_is_refused() {
for command in ["start mahbot -V", "cmd /c mahbot bench-openrouter"] {
let cause = refusal_cause(command);
assert!(cause.contains("runs `mahbot` itself"), "{command}: {cause}");
}
let cause = refusal_cause("mahbot -V >&2");
assert!(cause.contains("`>&2`"), "{cause}");
let cause = refusal_cause("cd /x && mahbot -V");
assert!(cause.contains("`cd` member"), "{cause}");
let cause = refusal_cause("mahbot -V %TEMP%");
assert!(cause.contains("runs `mahbot` itself"), "{cause}");
let cause = refusal_cause("mahbot -V (x)");
assert!(cause.contains("cannot be read"), "{cause}");
let cause = refusal_cause("cd sub && mahbot -V (x)");
assert!(cause.contains("cannot be read"), "{cause}");
let cause = refusal_cause(r"cd sub && C:\Program Files\MahBot\mahbot.exe -V (x)");
assert!(cause.contains("cannot be read"), "{cause}");
for command in ["@start mahbot -V", "@cmd /c mahbot -V"] {
let cause = refusal_cause(command);
assert!(cause.contains("runs `mahbot` itself"), "{command}: {cause}");
}
for command in [
"(mahbot -V)",
"( mahbot -V )",
"mahbot -V )",
"(cd x & mahbot -V)",
] {
let cause = refusal_cause(command);
assert!(cause.contains("cannot be read"), "{command}: {cause}");
}
}
#[test]
fn a_line_that_names_no_own_image_is_the_shells() {
let exe = Path::new(EXE);
let cwd = Path::new(ROOT);
for command in [
"grep -rn mahbot .",
"echo mahbot > out.txt",
"type mahbot.exe",
"echo hi && echo done",
"echo (x) mahbot",
"echo @ mahbot -V",
"for /f %i in ('mahbot -V') do echo x",
] {
assert!(
matches!(
own_image_plan(Executed(command), Written(command), exe, cwd),
OwnImage::None
),
"{command}"
);
}
if crate::tools::shell::SHELL_PLATFORM == ShellPlatform::Unix {
assert!(matches!(
plan_for_command(
Executed("mahbot -V && echo done"),
Written("mahbot -V && echo done"),
cwd
),
OwnImage::None
));
}
}
#[test]
fn the_cwd_a_fold_starts_in_is_the_members_own() {
let root = tracked_root_of(ROOT);
let plan = plan(&[
shell("cd src", "&&"),
served(r"C:\ws\src", "|"),
shell("head -3", ""),
]);
assert_eq!(
step_cwds(&plan),
vec![
root.as_path(),
Path::new(r"C:\ws\src"),
Path::new(r"C:\ws\src")
]
);
let generic = direct_plan(&format!(r#"cd src && "{EXE}" -V | head -3"#));
assert_eq!(generic.steps.len(), plan.steps.len());
assert_eq!(generic.steps[1].cwd, generic.steps[2].cwd);
assert!(generic.steps[2].cwd.ends_with("src"));
}
#[cfg(unix)]
fn shell_plan(steps: &[(&str, Join)]) -> Plan {
Plan {
steps: steps
.iter()
.map(|(text, join)| Step {
run: Run::Shell {
text: (*text).to_string(),
},
cwd: std::env::temp_dir(),
join: *join,
})
.collect(),
exe: PathBuf::from(EXE),
}
}
#[cfg(unix)]
async fn run_plan(plan: &Plan) -> ShellRunResult {
run(
plan,
Duration::from_secs(20),
Duration::from_secs(5),
None,
RunOwner::Agent,
)
.await
}
#[cfg(unix)]
#[tokio::test]
async fn a_group_runs_only_when_its_connector_says_so() {
let skipped = run_plan(&shell_plan(&[
("exit 3", Join::First),
("echo never", Join::And),
]))
.await;
let ShellRunResult::Completed { stdout, status, .. } = skipped else {
panic!("a completed run");
};
assert!(stdout.is_empty(), "the `&&` group must not run");
assert_eq!(status.code(), Some(3), "the failed group's status stays");
let ran = run_plan(&shell_plan(&[
("exit 3", Join::First),
("echo after-failure", Join::Or),
]))
.await;
let ShellRunResult::Completed { stdout, status, .. } = ran else {
panic!("a completed run");
};
assert_eq!(String::from_utf8_lossy(&stdout), "after-failure\n");
assert_eq!(status.code(), Some(0), "the last member's status");
}
#[cfg(unix)]
#[tokio::test]
async fn a_pipe_group_connects_its_members() {
let result = run_plan(&shell_plan(&[
("printf 'a\\nb\\nc\\n'", Join::First),
("head -2", Join::Pipe),
]))
.await;
let ShellRunResult::Completed { stdout, status, .. } = result else {
panic!("a completed run");
};
assert_eq!(String::from_utf8_lossy(&stdout), "a\nb\n");
assert_eq!(status.code(), Some(0));
}
#[cfg(unix)]
#[tokio::test]
async fn a_pipelines_input_is_its_heads_own_stdout() {
let result = run_plan(&shell_plan(&[
("echo hi", Join::First),
(r"printf 'x\ny\n'", Join::And),
("head -1", Join::Pipe),
]))
.await;
let ShellRunResult::Completed { stdout, .. } = result else {
panic!("a completed run");
};
assert_eq!(String::from_utf8_lossy(&stdout), "hi\nx\n");
}
#[cfg(unix)]
#[tokio::test]
async fn an_unsatisfied_head_takes_its_pipeline_with_it() {
let result = run_plan(&shell_plan(&[
("exit 3", Join::First),
("echo x", Join::And),
("cat", Join::Pipe),
]))
.await;
let ShellRunResult::Completed { stdout, status, .. } = result else {
panic!("a completed run");
};
assert!(stdout.is_empty(), "the `&&` group must not run");
assert_eq!(status.code(), Some(3), "the failed group's status stays");
}
#[cfg(unix)]
#[tokio::test]
async fn the_run_reports_the_last_members_status_and_concatenated_output() {
let result = run_plan(&shell_plan(&[
("echo one", Join::First),
("echo two; exit 4", Join::Always),
]))
.await;
let ShellRunResult::Completed { stdout, status, .. } = result else {
panic!("a completed run");
};
assert_eq!(String::from_utf8_lossy(&stdout), "one\ntwo\n");
assert_eq!(status.code(), Some(4));
}
#[cfg(unix)]
fn program_plan(steps: &[(&[&str], &[&str], Join)]) -> Plan {
Plan {
steps: steps
.iter()
.map(|(args, redirects, join)| Step {
run: Run::Own {
args: args.iter().map(|word| (*word).to_string()).collect(),
redirects: redirects.iter().map(|word| (*word).to_string()).collect(),
},
cwd: std::env::temp_dir(),
join: *join,
})
.collect(),
exe: PathBuf::from("/bin/sh"),
}
}
#[cfg(unix)]
#[tokio::test]
async fn a_merged_stderr_joins_the_captured_stdout() {
let result = run_plan(&program_plan(&[(
&["-c", "echo out; echo err >&2"],
&["2>&1"],
Join::First,
)]))
.await;
let ShellRunResult::Completed { stdout, stderr, .. } = result else {
panic!("a completed run");
};
assert_eq!(String::from_utf8_lossy(&stdout), "out\nerr\n");
assert!(
stderr.is_empty(),
"a merged stderr has no stderr channel: {}",
String::from_utf8_lossy(&stderr)
);
}
#[cfg(unix)]
#[tokio::test]
async fn a_merged_stderr_reaches_the_pipelines_consumer() {
let result = run_plan(&program_plan(&[
(&["-c", "echo out; echo err >&2"], &["2>&1"], Join::First),
(&["-c", "cat"], &[], Join::Pipe),
]))
.await;
let ShellRunResult::Completed { stdout, .. } = result else {
panic!("a completed run");
};
let captured = String::from_utf8_lossy(&stdout);
let mut lines: Vec<&str> = captured.lines().collect();
lines.sort_unstable();
assert_eq!(lines, ["err", "out"]);
}
#[cfg(unix)]
#[tokio::test]
async fn a_merge_after_a_redirect_takes_that_redirects_destination() {
let dir = std::env::temp_dir();
let target = dir.join(format!("plan-merge-{}.txt", std::process::id()));
let _ = std::fs::remove_file(&target);
let result = run_plan(&program_plan(&[(
&["-c", "echo out; echo err >&2"],
&[">", &target.to_string_lossy(), "2>&1"],
Join::First,
)]))
.await;
let ShellRunResult::Completed { stdout, .. } = result else {
panic!("a completed run");
};
assert!(
stdout.is_empty(),
"the member's own redirect took its streams: {}",
String::from_utf8_lossy(&stdout)
);
let written = std::fs::read_to_string(&target).expect("the redirect's file");
let mut lines: Vec<&str> = written.lines().collect();
lines.sort_unstable();
assert_eq!(lines, ["err", "out"]);
let _ = std::fs::remove_file(&target);
}
#[cfg(unix)]
#[tokio::test]
async fn a_merge_after_a_redirect_overrides_it() {
let dir = std::env::temp_dir();
let target = dir.join(format!("plan-override-{}.txt", std::process::id()));
let _ = std::fs::remove_file(&target);
let result = run_plan(&program_plan(&[(
&["-c", "echo out; echo err >&2"],
&["2>", &target.to_string_lossy(), "2>&1"],
Join::First,
)]))
.await;
let ShellRunResult::Completed { stdout, stderr, .. } = result else {
panic!("a completed run");
};
assert_eq!(String::from_utf8_lossy(&stdout), "out\nerr\n");
assert!(
stderr.is_empty(),
"the later merge left no stderr channel: {}",
String::from_utf8_lossy(&stderr)
);
assert_eq!(
std::fs::read_to_string(&target).expect("the redirect's file"),
""
);
let _ = std::fs::remove_file(&target);
}
#[cfg(unix)]
#[tokio::test]
async fn a_flooded_stdout_does_not_truncate_stderr() {
let result = run_plan(&program_plan(&[(
&["-c", "head -c 1048576 /dev/zero; echo marker >&2"],
&[],
Join::First,
)]))
.await;
let ShellRunResult::Completed { stdout, stderr, .. } = result else {
panic!("a completed run");
};
assert_eq!(stdout.len(), SHELL_PIPE_READ_CAP);
assert_eq!(String::from_utf8_lossy(&stderr), "marker\n");
}
#[cfg(unix)]
#[tokio::test]
async fn a_merge_before_a_redirect_still_feeds_the_pipe() {
let dir = std::env::temp_dir();
let target = dir.join(format!("plan-merge-pipe-{}.txt", std::process::id()));
let _ = std::fs::remove_file(&target);
let result = run_plan(&program_plan(&[
(
&["-c", "echo out; echo err >&2"],
&["2>&1", ">", &target.to_string_lossy()],
Join::First,
),
(&["-c", "cat"], &[], Join::Pipe),
]))
.await;
let ShellRunResult::Completed { stdout, .. } = result else {
panic!("a completed run");
};
assert_eq!(String::from_utf8_lossy(&stdout), "err\n");
assert_eq!(
std::fs::read_to_string(&target).expect("the redirect's file"),
"out\n"
);
let _ = std::fs::remove_file(&target);
}
#[cfg(unix)]
#[tokio::test]
async fn a_plan_past_its_deadline_is_stopped_and_named() {
let result = run(
&shell_plan(&[("echo started; sleep 30", Join::First)]),
Duration::from_secs(1),
Duration::from_secs(5),
None,
RunOwner::Agent,
)
.await;
let ShellRunResult::TimedOut {
stdout,
pid,
elapsed,
..
} = result
else {
panic!("expected a timeout, got {result:?}");
};
assert!(pid.is_some(), "the group the run waited for is named");
assert!(
elapsed < Duration::from_secs(20),
"the watchdog must not wait out the member: {elapsed:?}"
);
assert!(
String::from_utf8_lossy(&stdout).contains("started"),
"the member's output up to the stop is collected: {}",
String::from_utf8_lossy(&stdout)
);
}
#[cfg(unix)]
#[tokio::test]
async fn a_plan_over_the_memory_ceiling_is_stopped() {
const CEILING: u64 = 16 * 1024 * 1024;
const PAYLOAD: usize = 64 * 1024 * 1024;
let result = run(
&shell_plan(&[(
&format!("x=$(yes x | head -c {PAYLOAD}); sleep 30"),
Join::First,
)]),
Duration::from_secs(30),
Duration::from_secs(5),
Some(CEILING),
RunOwner::Agent,
)
.await;
let ShellRunResult::MemoryExceeded {
used,
limit,
elapsed,
..
} = result
else {
panic!("expected a memory failure, got {result:?}");
};
assert_eq!(limit, CEILING);
assert!(
used > limit,
"the tripping sample exceeds the ceiling: {used}"
);
assert!(
elapsed < Duration::from_secs(20),
"the watchdog must not wait out the member: {elapsed:?}"
);
}
#[cfg(unix)]
#[tokio::test]
async fn a_leftover_holder_past_the_drain_bound_keeps_its_status() {
let result = run(
&shell_plan(&[("echo out; sleep 5 &", Join::First)]),
Duration::from_secs(20),
Duration::from_millis(300),
None,
RunOwner::Agent,
)
.await;
let ShellRunResult::ExitedWithLeftovers {
stdout,
status,
scope,
..
} = result
else {
panic!("expected ExitedWithLeftovers, got {result:?}");
};
assert_eq!(status.code(), Some(0), "the member's own status");
assert!(
scope.is_none(),
"the leftover is not in one group the run can name"
);
assert!(
String::from_utf8_lossy(&stdout).contains("out"),
"what the member wrote before the bound is collected: {}",
String::from_utf8_lossy(&stdout)
);
}
#[cfg(unix)]
#[tokio::test]
async fn a_member_that_cannot_be_spawned_fails_the_run() {
let mut plan = shell_plan(&[("echo never", Join::First)]);
plan.steps[0].cwd = PathBuf::from("/nonexistent-plan-cwd");
let result = run_plan(&plan).await;
let ShellRunResult::SpawnFailed(e) = result else {
panic!("expected a spawn failure, got {result:?}");
};
assert!(!e.to_string().is_empty(), "the failure names its cause");
}
#[cfg(unix)]
#[tokio::test]
async fn a_redirect_that_cannot_be_opened_ends_the_group() {
let dir = std::env::temp_dir();
let marker = dir.join(format!("plan-cleanup-{}.txt", std::process::id()));
let _ = std::fs::remove_file(&marker);
let target = dir.join(format!("plan-missing-{}/out.txt", std::process::id()));
let result = run_plan(&program_plan(&[
(
&["-c", &format!("sleep 1; echo done > {}", marker.display())],
&[],
Join::First,
),
(
&["-c", "cat"],
&[">", &target.to_string_lossy()],
Join::Pipe,
),
]))
.await;
let ShellRunResult::SpawnFailed(_) = result else {
panic!("expected a spawn failure, got {result:?}");
};
tokio::time::sleep(Duration::from_millis(1500)).await;
assert!(
!marker.exists(),
"the group's earlier member must be killed, not left running"
);
}
}