use std::process::ExitStatus;
use std::sync::Arc;
use std::sync::Mutex;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use anyhow::{Result, bail};
use oxdock_parser::{Arg, AssertTarget, Step, StepKind, Value, guard_option_allows};
use oxdock_process::{BackgroundHandle, CommandStdin, ProcessManager};
fn exit_status_from_code(code: i32) -> ExitStatus {
#[cfg(unix)]
{
use std::os::unix::process::ExitStatusExt;
ExitStatus::from_raw(code << 8)
}
#[cfg(windows)]
{
use std::os::windows::process::ExitStatusExt;
ExitStatus::from_raw(code as u32)
}
}
use super::capture::SpillBuffer;
use super::handlers;
use super::io::{ExactCapture, SlidingWindow, StreamHandle};
use super::state::{ExecState, TaskEntry, TaskPhase};
pub(super) struct ThreadJoinHandle {
join: Option<std::thread::JoinHandle<Result<()>>>,
cancel_token: Arc<AtomicBool>,
active_process: Arc<Mutex<Option<Box<dyn BackgroundHandle>>>>,
thread_error: Option<anyhow::Error>,
}
impl ThreadJoinHandle {
pub(super) fn new(
join: std::thread::JoinHandle<Result<()>>,
cancel_token: Arc<AtomicBool>,
active_process: Arc<Mutex<Option<Box<dyn BackgroundHandle>>>>,
) -> Self {
Self {
join: Some(join),
cancel_token,
active_process,
thread_error: None,
}
}
fn reap(&mut self) {
if self.join.is_none() {
return;
}
let handle = self.join.take().unwrap();
match handle.join() {
Ok(Ok(())) => {}
Ok(Err(e)) => {
self.thread_error = Some(e);
}
Err(panic) => {
let msg = if let Some(s) = panic.downcast_ref::<&str>() {
s.to_string()
} else if let Some(s) = panic.downcast_ref::<String>() {
s.clone()
} else {
"thread panicked".to_string()
};
self.thread_error = Some(anyhow::anyhow!("{msg}"));
}
}
}
}
impl BackgroundHandle for ThreadJoinHandle {
fn try_wait(&mut self) -> Result<Option<ExitStatus>> {
if let Some(join) = &self.join {
if join.is_finished() {
self.reap();
} else {
return Ok(None);
}
}
if let Some(ref err) = self.thread_error {
Err(anyhow::anyhow!("{err}"))
} else {
Ok(Some(exit_status_from_code(0)))
}
}
fn kill(&mut self) -> Result<()> {
self.cancel_token.store(true, Ordering::SeqCst);
if let Ok(mut guard) = self.active_process.lock()
&& let Some(ref mut proc) = *guard
{
let _ = proc.kill();
}
self.reap();
Ok(())
}
fn wait(&mut self) -> Result<ExitStatus> {
self.reap();
if let Some(ref err) = self.thread_error {
Err(anyhow::anyhow!("{err}"))
} else {
Ok(exit_status_from_code(0))
}
}
}
impl Drop for ThreadJoinHandle {
fn drop(&mut self) {
let _ = self.kill();
}
}
static ASSERT_GENERATION: AtomicUsize = AtomicUsize::new(0);
#[derive(Debug)]
pub(super) enum Flow {
Done,
Break { idx: usize },
Continue { idx: usize },
Return { idx: usize, value: Value },
}
pub(super) fn allocate_assert_generation() -> usize {
ASSERT_GENERATION.fetch_add(1, Ordering::Relaxed)
}
#[derive(Clone, Copy, PartialEq, Eq)]
pub(super) enum AssertStream {
Stdout,
Stderr,
}
fn extract_stream_needle(kind: &StepKind) -> Option<(AssertStream, &Arg)> {
let step = match kind {
StepKind::WithIo { cmd, .. } => cmd.as_ref(),
other => other,
};
match step {
StepKind::AssertContains { haystack, needle } => match haystack {
AssertTarget::Stdout => Some((AssertStream::Stdout, needle)),
AssertTarget::Stderr => Some((AssertStream::Stderr, needle)),
_ => None,
},
_ => None,
}
}
fn needs_exact_stdout(kind: &StepKind) -> bool {
let step = match kind {
StepKind::WithIo { cmd, .. } => cmd.as_ref(),
other => other,
};
matches!(
step,
StepKind::AssertEq {
actual: AssertTarget::Stdout,
..
}
)
}
pub(super) enum ResolvedAssertTarget {
Value(Value),
Stdout,
Stderr,
Pipe(Vec<u8>),
}
pub(super) fn resolve_assert_target<P: ProcessManager>(
target: &AssertTarget,
cx: &mut StepCtx<'_, P>,
) -> Result<ResolvedAssertTarget> {
match target {
AssertTarget::Value(arg) => Ok(ResolvedAssertTarget::Value(
super::args::evaluate_assert_operand(arg, cx)?,
)),
AssertTarget::Stdout => Ok(ResolvedAssertTarget::Stdout),
AssertTarget::Stderr => Ok(ResolvedAssertTarget::Stderr),
AssertTarget::Pipe(name) => {
let bytes = cx.state.io.peek_pipe_content(name).map_err(|e| {
anyhow::anyhow!("step pipe assertion cannot read pipe {name:?}: {e}")
})?;
Ok(ResolvedAssertTarget::Pipe(bytes))
}
}
}
pub(super) fn pre_register_assertions<P: ProcessManager>(
state: &mut ExecState<P>,
steps: &[Step],
generation: usize,
) -> Result<()> {
let mut windows = match state.assert_windows.lock() {
Ok(guard) => guard,
Err(_) => bail!("assert_windows poisoned"),
};
let mut stderr_windows = match state.assert_windows_stderr.lock() {
Ok(guard) => guard,
Err(_) => bail!("assert_windows_stderr poisoned"),
};
let mut exact = match state.exact_stdout.lock() {
Ok(guard) => guard,
Err(_) => bail!("exact_stdout poisoned"),
};
for (idx, step) in steps.iter().enumerate() {
if let Some((stream, arg)) = extract_stream_needle(&step.kind) {
let resolved = super::args::resolve_arg_state(arg, state)?;
let map = match stream {
AssertStream::Stdout => &mut windows,
AssertStream::Stderr => &mut stderr_windows,
};
map.insert((generation, idx), SlidingWindow::new(resolved.into_bytes()));
}
if needs_exact_stdout(&step.kind) {
exact.entry(generation).or_insert_with(ExactCapture::new);
}
}
Ok(())
}
#[allow(clippy::collapsible_if)]
pub(super) fn sync_iteration_assert_needles<P: ProcessManager>(
state: &ExecState<P>,
steps: &[Step],
generation: usize,
) -> Result<()> {
let mut windows = match state.assert_windows.lock() {
Ok(guard) => guard,
Err(_) => bail!("assert_windows poisoned"),
};
let mut stderr_windows = match state.assert_windows_stderr.lock() {
Ok(guard) => guard,
Err(_) => bail!("assert_windows_stderr poisoned"),
};
for (idx, step) in steps.iter().enumerate() {
if let Some((stream, arg)) = extract_stream_needle(&step.kind) {
let map = match stream {
AssertStream::Stdout => &mut windows,
AssertStream::Stderr => &mut stderr_windows,
};
if let Some(w) = map.get_mut(&(generation, idx)) {
let resolved = super::args::resolve_arg_state(arg, state)?;
w.update_needle(resolved.into_bytes());
}
}
}
Ok(())
}
pub struct StepCtx<'a, P: ProcessManager> {
pub(super) state: &'a mut ExecState<P>,
pub(super) process: &'a mut P,
pub(super) stdin: CommandStdin,
pub(super) expose_stdin: bool,
pub(super) out: Option<StreamHandle>,
pub(super) err: Option<StreamHandle>,
}
#[allow(clippy::too_many_arguments)]
pub(super) fn execute_steps<P: ProcessManager>(
state: &mut ExecState<P>,
process: &mut P,
steps: &[Step],
stdin: CommandStdin,
expose_stdin: bool,
out: Option<StreamHandle>,
err: Option<StreamHandle>,
wait_at_end: bool,
) -> Result<Flow> {
let generation = allocate_assert_generation();
let flow = execute_steps_inner(
state,
process,
generation,
steps,
stdin,
expose_stdin,
out,
err,
wait_at_end,
)?;
let mut windows = match state.assert_windows.lock() {
Ok(guard) => guard,
Err(_) => bail!("assert_windows poisoned"),
};
windows.retain(|(g, _), _| *g != generation);
let mut stderr_windows = match state.assert_windows_stderr.lock() {
Ok(guard) => guard,
Err(_) => bail!("assert_windows_stderr poisoned"),
};
stderr_windows.retain(|(g, _), _| *g != generation);
let mut exact = match state.exact_stdout.lock() {
Ok(guard) => guard,
Err(_) => bail!("exact_stdout poisoned"),
};
exact.retain(|g, _| *g != generation);
Ok(flow)
}
#[allow(clippy::too_many_arguments)]
pub(super) fn execute_single_step_with_generation<P: ProcessManager>(
state: &mut ExecState<P>,
process: &mut P,
cmd: &StepKind,
generation: usize,
idx: usize,
stdin: CommandStdin,
expose_stdin: bool,
out: Option<StreamHandle>,
err: Option<StreamHandle>,
) -> Result<Flow> {
let mut cx = StepCtx {
state,
process,
stdin,
expose_stdin,
out,
err,
};
match cmd {
StepKind::FuncDef { .. }
| StepKind::Call { .. }
| StepKind::Return { .. }
| StepKind::While { .. }
| StepKind::Break
| StepKind::Continue
| StepKind::For { .. }
| StepKind::If { .. }
| StepKind::Timeout { .. }
| StepKind::WithIo { .. }
| StepKind::AssignCapture { .. } => {
return dispatch_flow_step(cmd, &mut cx, generation, idx);
}
_ => {}
}
match cmd {
StepKind::Run(arg) => {
let cmd = super::args::resolve_arg(arg, &mut cx)?;
let cmd = super::args::expand_dsl_vars(&cmd, cx.state);
handlers::run(&mut cx, idx, &cmd)
}
StepKind::RunExec { argv } => {
let resolved = handlers::resolve_run_exec_argv(argv, &mut cx)?;
handlers::run_argv(&mut cx, idx, &resolved)
}
StepKind::Echo(arg) => {
let msg = super::args::resolve_arg(arg, &mut cx)?;
handlers::echo(&mut cx, &msg)
}
StepKind::AsyncBlock { .. } => handlers::dispatch_async_block(cmd, &mut cx),
StepKind::Workdir(arg) => {
let path = super::args::resolve_arg(arg, &mut cx)?;
handlers::workdir(&mut cx, idx, &path)
}
StepKind::Workspace(target) => handlers::workspace(&mut cx, target),
StepKind::Env { key, value } => {
let resolved = super::args::resolve_arg(value, &mut cx)?;
handlers::env(&mut cx, key, &resolved)
}
StepKind::InheritEnv { keys } => {
handlers::inherit_env(&mut cx, keys)?;
sync_iteration_assert_needles(
cx.state,
&[Step {
guard: None,
kind: cmd.clone(),
scope_enter: 0,
scope_exit: 0,
}],
generation,
)?;
Ok(())
}
StepKind::Copy {
from_current_workspace,
from,
to,
} => {
let from_resolved = super::args::resolve_arg(from, &mut cx)?;
let to_resolved = super::args::resolve_arg(to, &mut cx)?;
handlers::copy(
&mut cx,
idx,
*from_current_workspace,
&from_resolved,
&to_resolved,
)
}
StepKind::CopyGit {
rev,
from,
to,
include_dirty,
} => {
let rev_resolved = super::args::resolve_arg(rev, &mut cx)?;
let from_resolved = super::args::resolve_arg(from, &mut cx)?;
let to_resolved = super::args::resolve_arg(to, &mut cx)?;
handlers::copy_git(
&mut cx,
idx,
&rev_resolved,
&from_resolved,
&to_resolved,
*include_dirty,
)
}
StepKind::HashSha256 { path } => {
let path_resolved = super::args::resolve_arg(path, &mut cx)?;
handlers::hash_sha256(&mut cx, idx, &path_resolved)
}
StepKind::Symlink { from, to } => {
let from_resolved = super::args::resolve_arg(from, &mut cx)?;
let to_resolved = super::args::resolve_arg(to, &mut cx)?;
handlers::symlink(&mut cx, idx, &from_resolved, &to_resolved)
}
StepKind::Mkdir(arg) => {
let path = super::args::resolve_arg(arg, &mut cx)?;
handlers::mkdir(&mut cx, idx, &path)
}
StepKind::Ls(arg) => {
let resolved = super::args::resolve_arg_opt(arg, &mut cx)?;
handlers::ls(&mut cx, idx, &resolved)
}
StepKind::Cwd => handlers::cwd(&mut cx, idx),
StepKind::Read(arg) => {
let resolved = super::args::resolve_arg_opt(arg, &mut cx)?;
handlers::read(&mut cx, idx, &resolved)
}
StepKind::ReadLine { var } => handlers::read_line(&mut cx, idx, var),
StepKind::Write { path, contents } => {
let path_resolved = super::args::resolve_arg(path, &mut cx)?;
let contents_resolved = super::args::resolve_arg_opt(contents, &mut cx)?;
handlers::write(&mut cx, idx, &path_resolved, contents_resolved.as_deref())
}
StepKind::Append { path, contents } => {
let path_resolved = super::args::resolve_arg(path, &mut cx)?;
let contents_resolved = super::args::resolve_arg_opt(contents, &mut cx)?;
handlers::append(&mut cx, idx, &path_resolved, contents_resolved.as_deref())
}
StepKind::Expand { path, overrides } => {
let path_resolved = super::args::resolve_arg_opt(path, &mut cx)?;
let overrides_resolved = super::args::resolve_overrides(overrides, &mut cx)?;
handlers::replace(&mut cx, idx, &path_resolved, &overrides_resolved)
}
StepKind::AssertEq {
hash,
actual,
expected,
} => {
let target = resolve_assert_target(actual, &mut cx)?;
let expected_resolved = match expected {
Some(e) => Some(super::args::evaluate_assert_operand(e, &mut cx)?),
None => None,
};
handlers::assert_eq(
&mut cx,
idx,
generation,
idx,
hash,
&target,
expected_resolved.as_ref(),
)
}
StepKind::AssertContains { haystack, needle } => {
let target = resolve_assert_target(haystack, &mut cx)?;
handlers::assert_contains(&mut cx, idx, generation, idx, &target, needle)
}
StepKind::WithIoBlock { .. } => {
bail!("WITH_IO block should have been expanded during parsing")
}
StepKind::Exit(code) => {
let code = super::args::resolve_arg_as_int(code, &mut cx)?;
handlers::exit(&mut cx, code)
}
StepKind::Assign {
var,
decl_type,
expr,
} => handlers::assign(&mut cx, var, *decl_type, expr),
StepKind::Set { var, expr } => handlers::set_var_value(&mut cx, var, expr),
StepKind::AssignAsync {
var,
decl_type,
body,
} => handlers::dispatch_assign_async(var, *decl_type, body, &mut cx),
StepKind::Await { var } => handlers::dispatch_await(var, &mut cx),
StepKind::AwaitCapture {
out_var,
out_type,
task_var,
} => handlers::dispatch_await_capture(out_var, *out_type, task_var, &mut cx),
StepKind::Cancel { var } => handlers::dispatch_cancel(var, &mut cx),
StepKind::Sleep { duration } => {
let duration = super::args::resolve_arg_as_duration(duration, &mut cx)?;
handlers::sleep(&mut cx, idx, &duration)
}
StepKind::FuncDef { .. }
| StepKind::Call { .. }
| StepKind::Return { .. }
| StepKind::While { .. }
| StepKind::Break
| StepKind::Continue
| StepKind::For { .. }
| StepKind::If { .. }
| StepKind::Timeout { .. }
| StepKind::WithIo { .. }
| StepKind::AssignCapture { .. } => {
unreachable!("compound steps dispatch before this match")
}
}?;
Ok(Flow::Done)
}
#[allow(clippy::too_many_arguments)]
fn execute_steps_inner<P: ProcessManager>(
state: &mut ExecState<P>,
process: &mut P,
generation: usize,
steps: &[Step],
stdin: CommandStdin,
expose_stdin: bool,
out: Option<StreamHandle>,
err: Option<StreamHandle>,
wait_at_end: bool,
) -> Result<Flow> {
pre_register_assertions(state, steps, generation)?;
for (idx, step) in steps.iter().enumerate() {
if state.cancel_token.load(Ordering::SeqCst) {
bail!("ASYNC task cancelled");
}
if step.scope_enter > 0 {
for _ in 0..step.scope_enter {
state.push_scope();
}
}
let should_run = guard_option_allows(step.guard.as_ref(), &state.envs);
let flow_result: Result<Flow> = if !should_run {
Ok(Flow::Done)
} else {
let mut cx = StepCtx {
state,
process,
stdin: stdin.clone(),
expose_stdin,
out: out.clone(),
err: err.clone(),
};
let flow_result: Result<Flow> = match &step.kind {
StepKind::FuncDef { .. }
| StepKind::Call { .. }
| StepKind::Return { .. }
| StepKind::While { .. }
| StepKind::Break
| StepKind::Continue
| StepKind::For { .. }
| StepKind::If { .. }
| StepKind::Timeout { .. }
| StepKind::WithIo { .. }
| StepKind::AssignCapture { .. } => {
dispatch_flow_step(&step.kind, &mut cx, generation, idx)
}
_ => {
match &step.kind {
StepKind::InheritEnv { keys } => {
handlers::inherit_env(&mut cx, keys)?;
sync_iteration_assert_needles(cx.state, steps, generation)?;
Ok(())
}
StepKind::Workdir(arg) => {
let path = super::args::resolve_arg(arg, &mut cx)?;
handlers::workdir(&mut cx, idx, &path)
}
StepKind::Workspace(target) => handlers::workspace(&mut cx, target),
StepKind::Env { key, value } => {
let resolved = super::args::resolve_arg(value, &mut cx)?;
handlers::env(&mut cx, key, &resolved)?;
sync_iteration_assert_needles(cx.state, steps, generation)?;
Ok(())
}
StepKind::Run(arg) => {
let cmd = super::args::resolve_arg(arg, &mut cx)?;
let cmd = super::args::expand_dsl_vars(&cmd, cx.state);
handlers::run(&mut cx, idx, &cmd)
}
StepKind::RunExec { argv } => {
let resolved = handlers::resolve_run_exec_argv(argv, &mut cx)?;
handlers::run_argv(&mut cx, idx, &resolved)
}
StepKind::Echo(arg) => {
let msg = super::args::resolve_arg(arg, &mut cx)?;
handlers::echo(&mut cx, &msg)
}
StepKind::AsyncBlock { .. } => {
handlers::dispatch_async_block(&step.kind, &mut cx)
}
StepKind::Copy {
from_current_workspace,
from,
to,
} => {
let from_resolved = super::args::resolve_arg(from, &mut cx)?;
let to_resolved = super::args::resolve_arg(to, &mut cx)?;
handlers::copy(
&mut cx,
idx,
*from_current_workspace,
&from_resolved,
&to_resolved,
)
}
StepKind::CopyGit {
rev,
from,
to,
include_dirty,
} => {
let rev_resolved = super::args::resolve_arg(rev, &mut cx)?;
let from_resolved = super::args::resolve_arg(from, &mut cx)?;
let to_resolved = super::args::resolve_arg(to, &mut cx)?;
handlers::copy_git(
&mut cx,
idx,
&rev_resolved,
&from_resolved,
&to_resolved,
*include_dirty,
)
}
StepKind::HashSha256 { path } => {
let path_resolved = super::args::resolve_arg(path, &mut cx)?;
handlers::hash_sha256(&mut cx, idx, &path_resolved)
}
StepKind::Symlink { from, to } => {
let from_resolved = super::args::resolve_arg(from, &mut cx)?;
let to_resolved = super::args::resolve_arg(to, &mut cx)?;
handlers::symlink(&mut cx, idx, &from_resolved, &to_resolved)
}
StepKind::Mkdir(arg) => {
let path = super::args::resolve_arg(arg, &mut cx)?;
handlers::mkdir(&mut cx, idx, &path)
}
StepKind::Ls(arg) => {
let resolved = super::args::resolve_arg_opt(arg, &mut cx)?;
handlers::ls(&mut cx, idx, &resolved)
}
StepKind::Cwd => handlers::cwd(&mut cx, idx),
StepKind::Read(arg) => {
let resolved = super::args::resolve_arg_opt(arg, &mut cx)?;
handlers::read(&mut cx, idx, &resolved)
}
StepKind::ReadLine { var } => handlers::read_line(&mut cx, idx, var),
StepKind::Write { path, contents } => {
let path_resolved = super::args::resolve_arg(path, &mut cx)?;
let contents_resolved =
super::args::resolve_arg_opt(contents, &mut cx)?;
handlers::write(
&mut cx,
idx,
&path_resolved,
contents_resolved.as_deref(),
)
}
StepKind::Append { path, contents } => {
let path_resolved = super::args::resolve_arg(path, &mut cx)?;
let contents_resolved =
super::args::resolve_arg_opt(contents, &mut cx)?;
handlers::append(
&mut cx,
idx,
&path_resolved,
contents_resolved.as_deref(),
)
}
StepKind::Expand { path, overrides } => {
let path_resolved = super::args::resolve_arg_opt(path, &mut cx)?;
let overrides_resolved =
super::args::resolve_overrides(overrides, &mut cx)?;
handlers::replace(&mut cx, idx, &path_resolved, &overrides_resolved)
}
StepKind::AssertEq {
hash,
actual,
expected,
} => {
let target = resolve_assert_target(actual, &mut cx)?;
let expected_resolved = match expected {
Some(e) => Some(super::args::evaluate_assert_operand(e, &mut cx)?),
None => None,
};
handlers::assert_eq(
&mut cx,
idx,
generation,
idx,
hash,
&target,
expected_resolved.as_ref(),
)
}
StepKind::AssertContains { haystack, needle } => {
let target = resolve_assert_target(haystack, &mut cx)?;
handlers::assert_contains(
&mut cx, idx, generation, idx, &target, needle,
)
}
StepKind::WithIoBlock { .. } => {
bail!("WITH_IO block should have been expanded during parsing")
}
StepKind::Exit(code) => {
let code = super::args::resolve_arg_as_int(code, &mut cx)?;
handlers::exit(&mut cx, code)
}
StepKind::Assign {
var,
decl_type,
expr,
} => handlers::assign(&mut cx, var, *decl_type, expr),
StepKind::Set { var, expr } => handlers::set_var_value(&mut cx, var, expr),
StepKind::AssignAsync {
var,
decl_type,
body,
} => handlers::dispatch_assign_async(var, *decl_type, body, &mut cx),
StepKind::Await { var } => handlers::dispatch_await(var, &mut cx),
StepKind::AwaitCapture {
out_var,
out_type,
task_var,
} => {
handlers::dispatch_await_capture(out_var, *out_type, task_var, &mut cx)
}
StepKind::Cancel { var } => handlers::dispatch_cancel(var, &mut cx),
StepKind::Sleep { duration } => {
let duration = super::args::resolve_arg_as_duration(duration, &mut cx)?;
handlers::sleep(&mut cx, idx, &duration)
}
StepKind::FuncDef { .. }
| StepKind::Call { .. }
| StepKind::Return { .. }
| StepKind::While { .. }
| StepKind::Break
| StepKind::Continue
| StepKind::For { .. }
| StepKind::If { .. }
| StepKind::Timeout { .. }
| StepKind::WithIo { .. }
| StepKind::AssignCapture { .. } => {
unreachable!("compound steps dispatch in the outer match")
}
}?;
Ok(Flow::Done)
}
};
flow_result
};
let restore_result = restore_scopes(state, step.scope_exit);
let expiry_drained = if let Some(expiry) = state.keeper_expiry.as_mut() {
expiry.expire_step(steps, idx)
} else {
false
};
if expiry_drained {
state.keeper_expiry = None;
}
let flow = flow_result?;
restore_result?;
match flow {
Flow::Done => {}
Flow::Break { .. } | Flow::Continue { .. } | Flow::Return { .. } => {
return Ok(flow);
}
}
}
let reap_named = !state.inside_async;
let has_bg = !state.bg_children.is_empty();
let named_pending = |state: &ExecState<P>| {
reap_named
&& state
.named_tasks
.lock()
.unwrap_or_else(|e| e.into_inner())
.values()
.any(|entry| !entry.state.lock().unwrap_or_else(|e| e.into_inner()).reaped)
};
let has_named = named_pending(state);
if wait_at_end && (has_bg || has_named) {
loop {
let mut failed_status: Option<anyhow::Error> = None;
if failed_status.is_none() && state.cancel_token.load(Ordering::SeqCst) {
failed_status = Some(anyhow::anyhow!("ASYNC task cancelled"));
}
let mut i = 0;
while i < state.bg_children.len() {
match state.bg_children[i].try_wait() {
Ok(Some(status)) => {
if !status.success() && failed_status.is_none() {
failed_status =
Some(anyhow::anyhow!("ASYNC process exited with status {status}"));
break;
}
state.bg_children.swap_remove(i);
}
Ok(None) => {
i += 1;
}
Err(e) => {
if failed_status.is_none() {
failed_status = Some(e);
}
break;
}
}
}
if failed_status.is_none() && reap_named {
let entries: Vec<(u64, Arc<TaskEntry>)> = {
let named = state.named_tasks.lock().unwrap_or_else(|e| e.into_inner());
named
.iter()
.map(|(id, entry)| (*id, Arc::clone(entry)))
.collect()
};
for (id, entry) in &entries {
enum Poll {
Pending,
CompletedOk { sink: Option<Arc<SpillBuffer>> },
CompletedErr(anyhow::Error),
}
let poll = {
let mut guard = entry.state.lock().unwrap_or_else(|e| e.into_inner());
match guard.phase {
TaskPhase::Running | TaskPhase::Awaiting => {
let take_sink = matches!(guard.phase, TaskPhase::Running);
match guard.handle.as_mut() {
Some(handle) => match handle.try_wait() {
Ok(Some(status)) => {
let _ = guard.handle.take();
guard.phase = TaskPhase::Completed;
let sink =
if take_sink { guard.sink.take() } else { None };
if status.success() {
Poll::CompletedOk { sink }
} else {
Poll::CompletedErr(anyhow::anyhow!(
"named ASYNC task {id} exited with status {status}"
))
}
}
Ok(None) => Poll::Pending,
Err(e) => {
let _ = guard.handle.take();
guard.phase = TaskPhase::Completed;
Poll::CompletedErr(e)
}
},
None => Poll::Pending,
}
}
TaskPhase::Cancelled | TaskPhase::Completed => Poll::Pending,
}
};
match poll {
Poll::Pending => {}
Poll::CompletedOk { sink } => {
entry.finish_teardown();
if let Some(sink) = sink
&& let Err(e) = forward_task_sink(&sink, &out, *id)
{
if failed_status.is_none() {
failed_status = Some(e);
}
break;
}
}
Poll::CompletedErr(e) => {
entry.finish_teardown();
if failed_status.is_none() {
failed_status = Some(e);
}
break;
}
}
}
}
if let Some(err) = failed_status {
for survivor in state.bg_children.iter_mut() {
let _ = survivor.kill();
}
if reap_named {
let entries: Vec<Arc<TaskEntry>> = {
let named = state.named_tasks.lock().unwrap_or_else(|e| e.into_inner());
named.values().cloned().collect()
};
let mut to_kill: Vec<(Arc<TaskEntry>, Box<dyn BackgroundHandle>)> = Vec::new();
for entry in &entries {
let mut guard = entry.state.lock().unwrap_or_else(|e| e.into_inner());
match guard.phase {
TaskPhase::Running | TaskPhase::Awaiting => {
guard.phase = TaskPhase::Cancelled;
if let Some(handle) = guard.handle.take() {
to_kill.push((Arc::clone(entry), handle));
}
}
TaskPhase::Cancelled | TaskPhase::Completed => {}
}
}
for (entry, mut handle) in to_kill {
let _ = handle.kill();
entry.finish_teardown();
}
}
state.bg_children.clear();
return Err(err);
}
let bg_empty = state.bg_children.is_empty();
if bg_empty && !named_pending(state) {
return Ok(Flow::Done);
}
if reap_named {
let unreaped: Vec<Arc<TaskEntry>> = {
let named = state.named_tasks.lock().unwrap_or_else(|e| e.into_inner());
named
.values()
.filter(|entry| {
let guard = entry.state.lock().unwrap_or_else(|e| e.into_inner());
matches!(guard.phase, TaskPhase::Cancelled) && !guard.reaped
})
.cloned()
.collect()
};
for entry in &unreaped {
entry.wait_reaped();
}
}
std::thread::sleep(std::time::Duration::from_millis(10));
}
}
Ok(Flow::Done)
}
fn forward_task_sink(sink: &Arc<SpillBuffer>, out: &Option<StreamHandle>, id: u64) -> Result<()> {
let bytes = sink
.drain_bytes()
.map_err(|e| anyhow::anyhow!("named ASYNC task {id} output drain failed: {e}"))?;
if !bytes.is_empty() {
super::io::write_stdout(out.clone(), |writer| {
writer
.write_all(&bytes)
.map_err(|e| anyhow::anyhow!("named ASYNC task {id} output forward failed: {e}"))?;
Ok(())
})?;
}
Ok(())
}
fn restore_scopes<P: ProcessManager>(state: &mut ExecState<P>, count: usize) -> Result<()> {
for _ in 0..count {
state.pop_scope()?;
}
Ok(())
}
#[allow(clippy::too_many_arguments)]
pub(super) fn execute_scoped_steps<P: ProcessManager>(
state: &mut ExecState<P>,
process: &mut P,
steps: &[Step],
stdin: CommandStdin,
expose_stdin: bool,
out: Option<StreamHandle>,
err: Option<StreamHandle>,
wait_at_end: bool,
) -> Result<Flow> {
state.push_scope();
let res = execute_steps(
state,
process,
steps,
stdin,
expose_stdin,
out,
err,
wait_at_end,
);
let pop_res = state.pop_scope();
match (res, pop_res) {
(Ok(flow), Ok(())) => Ok(flow),
(Err(e), _) => Err(e),
(Ok(_), Err(e)) => Err(e),
}
}
fn dispatch_flow_step<P: ProcessManager>(
cmd: &StepKind,
cx: &mut StepCtx<'_, P>,
generation: usize,
idx: usize,
) -> Result<Flow> {
match cmd {
StepKind::FuncDef { name, params, body } => {
handlers::define_func(cx, name, params, body)?;
Ok(Flow::Done)
}
StepKind::Call { name, args } => {
let _ = handlers::call_func_value(cx, idx, name, args)?;
Ok(Flow::Done)
}
StepKind::Return { expr } => handlers::handle_return(cx, idx, expr),
StepKind::While { cond, body } => handlers::while_loop(cx, idx, cond, body),
StepKind::Break => Ok(Flow::Break { idx }),
StepKind::Continue => Ok(Flow::Continue { idx }),
StepKind::For {
key_var,
key_type,
var,
var_type,
in_expr,
body,
} => handlers::for_loop(
cx,
key_var.as_deref(),
*key_type,
var,
*var_type,
in_expr,
body,
),
StepKind::If {
cond,
then_body,
else_ifs,
else_body,
} => handlers::if_then(cx, cond, then_body, else_ifs, else_body),
StepKind::Timeout { duration, body } => {
let duration = super::args::resolve_arg_as_duration(duration, cx)?;
handlers::timeout(cx, idx, &duration, body)
}
StepKind::WithIo { bindings, cmd } => handlers::with_io(cx, generation, idx, bindings, cmd),
StepKind::AssignCapture {
var,
decl_type,
cmd,
} => handlers::assign_capture(cx, generation, idx, var, *decl_type, cmd),
_ => {
unreachable!("dispatch_flow_step handles only compound steps")
}
}
}