use std::sync::Arc;
use std::collections::HashMap;
use crate::arithmetic;
use crate::ast::{Arg, Command, Expr, PipelineStage, Redirect, RedirectKind, Value};
use crate::dispatch::{CommandDispatcher, PipelinePosition};
use crate::interpreter::{apply_output_format, ExecResult, OutputFormat, PathError};
use crate::tools::{global_flag_value_is_truthy, ExecContext, ToolArgs, ToolRegistry, ToolSchema};
use tokio::io::AsyncWriteExt;
use super::pipe_stream::pipe_stream_default;
use super::scatter::{
parse_gather_options, parse_scatter_options, ScatterGatherRunner,
};
fn has_json_flag(args: &[Arg]) -> bool {
let mut past_double_dash = false;
for arg in args {
match arg {
Arg::DoubleDash => past_double_dash = true,
_ if past_double_dash => {}
Arg::LongFlag(name) if name == "json" => return true,
Arg::Named { key, value } if key == "json" => {
let on = match value {
Expr::Literal(v) => global_flag_value_is_truthy(v),
_ => true,
};
if on {
return true;
}
}
_ => {}
}
}
false
}
fn finalize_scatter_gather_error(result: ExecResult, format: Option<OutputFormat>) -> ExecResult {
match format {
Some(format) => apply_output_format(result, format),
None => result,
}
}
pub(crate) async fn apply_redirects(
mut result: ExecResult,
redirects: &[Redirect],
ctx: &ExecContext,
dispatcher: &dyn CommandDispatcher,
) -> ExecResult {
for redir in redirects {
match redir.kind {
RedirectKind::MergeStderr => {
result.materialize();
if !result.err.is_empty() {
let err = std::mem::take(&mut result.err);
result.push_out(&err);
}
}
RedirectKind::MergeStdout => {
if result.is_bytes() {
return ExecResult::failure(
1,
"redirect: cannot merge binary stdout into stderr (1>&2) — \
redirect it to a file or pipe through base64/xxd",
);
}
result.materialize();
if !result.text_out().is_empty() {
let out = result.text_out().into_owned();
result.err.push_str(&out);
}
result.clear_stdout();
}
RedirectKind::StdoutOverwrite => {
let path = match eval_redirect_target(&redir.target, ctx, dispatcher).await {
Ok(p) => p,
Err(e) => return ExecResult::failure(1, format!("redirect: {e}")),
};
if let Some(bytes) = result.out_bytes() {
if let Err(e) = redirect_write(ctx, &path, bytes).await {
return ExecResult::failure(1, format!("redirect: {e}"));
}
} else if let Some(output) = result.take_output_for_stream() {
let mut buf = Vec::new();
if let Err(e) = output.write_canonical(&mut buf, None) {
return ExecResult::failure(1, format!("redirect: {e}"));
}
if let Err(e) = redirect_write(ctx, &path, &buf).await {
return ExecResult::failure(1, format!("redirect: {e}"));
}
} else if let Err(e) = redirect_write(ctx, &path, result.text_out().as_bytes()).await {
return ExecResult::failure(1, format!("redirect: {e}"));
}
result.clear_stdout();
}
RedirectKind::StdoutAppend => {
let path = match eval_redirect_target(&redir.target, ctx, dispatcher).await {
Ok(p) => p,
Err(e) => return ExecResult::failure(1, format!("redirect: {e}")),
};
if let Some(bytes) = result.out_bytes() {
if let Err(e) = redirect_append(ctx, &path, bytes).await {
return ExecResult::failure(1, format!("redirect: {e}"));
}
} else if let Some(output) = result.take_output_for_stream() {
let mut buf = Vec::new();
if let Err(e) = output.write_canonical(&mut buf, None) {
return ExecResult::failure(1, format!("redirect: {e}"));
}
if let Err(e) = redirect_append(ctx, &path, &buf).await {
return ExecResult::failure(1, format!("redirect: {e}"));
}
} else if let Err(e) = redirect_append(ctx, &path, result.text_out().as_bytes()).await {
return ExecResult::failure(1, format!("redirect: {e}"));
}
result.clear_stdout();
}
RedirectKind::Stderr => {
let path = match eval_redirect_target(&redir.target, ctx, dispatcher).await {
Ok(p) => p,
Err(e) => return ExecResult::failure(1, format!("redirect: {e}")),
};
if let Err(e) = redirect_write(ctx, &path, result.err.as_bytes()).await {
return ExecResult::failure(1, format!("redirect: {e}"));
}
result.err.clear();
}
RedirectKind::Both => {
let path = match eval_redirect_target(&redir.target, ctx, dispatcher).await {
Ok(p) => p,
Err(e) => return ExecResult::failure(1, format!("redirect: {e}")),
};
let mut combined: Vec<u8> = if let Some(b) = result.out_bytes() {
b.to_vec()
} else if let Some(output) = result.take_output_for_stream() {
let mut buf = Vec::new();
if let Err(e) = output.write_canonical(&mut buf, None) {
return ExecResult::failure(1, format!("redirect: {e}"));
}
buf
} else {
result.text_out().into_owned().into_bytes()
};
combined.extend_from_slice(result.err.as_bytes());
if let Err(e) = redirect_write(ctx, &path, &combined).await {
return ExecResult::failure(1, format!("redirect: {e}"));
}
result.clear_stdout();
result.err.clear();
}
RedirectKind::Stdin | RedirectKind::HereDoc(_) | RedirectKind::HereString => {}
}
}
result
}
async fn eval_redirect_target(
expr: &Expr,
ctx: &ExecContext,
dispatcher: &dyn CommandDispatcher,
) -> Result<String, String> {
let value = dispatcher
.eval_expr(expr, ctx)
.await
.map_err(|e| e.to_string())?;
if let Some(msg) = crate::interpreter::structured_boundary_error("a redirect target", &value) {
return Err(msg);
}
crate::interpreter::value_to_text_sink_named(&value, "a redirect target").map_err(|e| e.to_string())
}
async fn redirect_write(ctx: &ExecContext, path: &str, data: &[u8]) -> Result<(), String> {
use crate::backend::WriteMode;
let resolved = ctx.resolve_path(path);
ctx.backend.write(&resolved, data, WriteMode::Overwrite).await.map_err(|e| e.to_string())
}
async fn redirect_append(ctx: &ExecContext, path: &str, data: &[u8]) -> Result<(), String> {
let resolved = ctx.resolve_path(path);
ctx.backend.append(&resolved, data).await.map_err(|e| e.to_string())
}
async fn setup_stdin_redirects(
cmd: &Command,
ctx: &mut ExecContext,
dispatcher: &dyn CommandDispatcher,
) -> Result<(), String> {
use std::path::Path;
for redir in &cmd.redirects {
match &redir.kind {
RedirectKind::Stdin => {
let path = eval_redirect_target(&redir.target, ctx, dispatcher).await?;
let resolved = ctx.resolve_path(&path);
let data = ctx
.backend
.read(Path::new(&resolved), None)
.await
.map_err(|e| format!("redirect: {path}: {e}"))?;
ctx.set_stdin(data);
}
RedirectKind::HereDoc(_) => {
match &redir.target {
Expr::Literal(Value::String(content)) => {
ctx.set_stdin(content.clone());
}
expr => {
let body = eval_redirect_target(expr, ctx, dispatcher).await?;
ctx.set_stdin(body);
}
}
}
RedirectKind::HereString => {
let mut s = eval_redirect_target(&redir.target, ctx, dispatcher).await?;
s.push('\n');
ctx.set_stdin(s);
}
_ => {}
}
}
Ok(())
}
async fn setup_stdin_redirects_for(
stage: &PipelineStage,
ctx: &mut ExecContext,
dispatcher: &dyn CommandDispatcher,
) -> Result<(), String> {
match stage {
PipelineStage::Command(cmd) => setup_stdin_redirects(cmd, ctx, dispatcher).await,
PipelineStage::Compound(_) => Ok(()),
}
}
async fn dispatch_stage(
stage: &PipelineStage,
ctx: &mut ExecContext,
dispatcher: &dyn CommandDispatcher,
) -> anyhow::Result<ExecResult> {
match stage {
PipelineStage::Command(cmd) => dispatcher.dispatch(cmd, ctx).await,
PipelineStage::Compound(stmt) => dispatcher.dispatch_stmt(stmt, ctx).await,
}
}
#[derive(Clone)]
pub struct PipelineRunner {
tools: Arc<ToolRegistry>,
}
impl PipelineRunner {
pub fn new(tools: Arc<ToolRegistry>) -> Self {
Self { tools }
}
pub async fn run(
&self,
stages: &[PipelineStage],
ctx: &mut ExecContext,
dispatcher: &dyn CommandDispatcher,
) -> ExecResult {
if stages.is_empty() {
return ExecResult::success("");
}
if let Some((scatter_idx, gather_idx)) = find_scatter_gather(stages) {
let commands: Vec<Command> = match stages
.iter()
.map(|s| s.as_command().cloned())
.collect::<Option<Vec<_>>>()
{
Some(commands) => commands,
None => {
return ExecResult::failure(
2,
"scatter/gather cannot share a pipeline with an if/for/while/case \
stage. Run the compound on its own and pipe its output in.",
)
}
};
return self
.run_scatter_gather(&commands, scatter_idx, gather_idx, ctx, dispatcher)
.await;
}
self.run_stage_sequence(stages, ctx, dispatcher).await
}
pub async fn run_sequential(
&self,
commands: &[Command],
ctx: &mut ExecContext,
dispatcher: &dyn CommandDispatcher,
) -> ExecResult {
let stages: Vec<PipelineStage> = commands
.iter()
.cloned()
.map(PipelineStage::Command)
.collect();
self.run_stage_sequence(&stages, ctx, dispatcher).await
}
async fn run_stage_sequence(
&self,
stages: &[PipelineStage],
ctx: &mut ExecContext,
dispatcher: &dyn CommandDispatcher,
) -> ExecResult {
if stages.is_empty() {
return ExecResult::success("");
}
if stages.len() == 1 {
let result = self.run_single(&stages[0], ctx, None, dispatcher).await;
ctx.scope.set_pipestatus(&[result.code]);
return result;
}
self.run_pipeline(stages, ctx, dispatcher).await
}
async fn run_scatter_gather(
&self,
commands: &[Command],
scatter_idx: usize,
gather_idx: usize,
ctx: &mut ExecContext,
dispatcher: &dyn CommandDispatcher,
) -> ExecResult {
let pre_scatter = &commands[..scatter_idx];
let scatter_cmd = &commands[scatter_idx];
let parallel = &commands[scatter_idx + 1..gather_idx];
let gather_cmd = &commands[gather_idx];
let post_gather = &commands[gather_idx + 1..];
let format = (has_json_flag(&scatter_cmd.args) || has_json_flag(&gather_cmd.args))
.then_some(OutputFormat::Json);
let scatter_schema = self.tools.get("scatter").map(|t| t.schema());
let gather_schema = self.tools.get("gather").map(|t| t.schema());
let scatter_args = match build_tool_args(&scatter_cmd.args, ctx, scatter_schema.as_ref()).await {
Ok(args) => args,
Err(e) => {
return finalize_scatter_gather_error(
ExecResult::failure(1, format!("scatter: {e}")),
format,
)
}
};
let gather_args = match build_tool_args(&gather_cmd.args, ctx, gather_schema.as_ref()).await {
Ok(args) => args,
Err(e) => {
return finalize_scatter_gather_error(
ExecResult::failure(1, format!("gather: {e}")),
format,
)
}
};
let scatter_opts = match parse_scatter_options(&scatter_args) {
Ok(opts) => opts,
Err(e) => {
return finalize_scatter_gather_error(
ExecResult::failure(2, format!("scatter: {e}")),
format,
)
}
};
let gather_opts = match parse_gather_options(&gather_args) {
Ok(opts) => opts,
Err(e) => {
return finalize_scatter_gather_error(
ExecResult::failure(2, format!("gather: {e}")),
format,
)
}
};
let sequential_dispatcher: Arc<dyn CommandDispatcher> = dispatcher.fork_attached().await;
let runner = ScatterGatherRunner::new(self.tools.clone(), sequential_dispatcher);
runner
.run(
pre_scatter,
scatter_opts,
parallel,
gather_opts,
&gather_cmd.redirects,
post_gather,
ctx,
)
.await
}
async fn run_single(
&self,
stage: &PipelineStage,
ctx: &mut ExecContext,
stdin: Option<Vec<u8>>,
dispatcher: &dyn CommandDispatcher,
) -> ExecResult {
if let Err(e) = setup_stdin_redirects_for(stage, ctx, dispatcher).await {
return ExecResult::failure(1, e);
}
if let Some(input) = stdin {
ctx.set_stdin(input);
}
ctx.pipeline_position = PipelinePosition::Only;
let result = match dispatch_stage(stage, ctx, dispatcher).await {
Ok(result) => result,
Err(e) => ExecResult::failure(1, e.to_string()),
};
apply_redirects(result, stage.redirects(), ctx, dispatcher).await
}
async fn run_pipeline(
&self,
stages: &[PipelineStage],
ctx: &mut ExecContext,
dispatcher: &dyn CommandDispatcher,
) -> ExecResult {
let stage_count = stages.len();
let last_idx = stage_count - 1;
let mut pipe_writers: Vec<Option<super::pipe_stream::PipeWriter>> = Vec::new();
let mut pipe_readers: Vec<Option<super::pipe_stream::PipeReader>> = Vec::new();
for _ in 0..last_idx {
let (writer, reader) = pipe_stream_default();
pipe_writers.push(Some(writer));
pipe_readers.push(Some(reader));
}
let mut data_senders: Vec<Option<tokio::sync::oneshot::Sender<Option<Value>>>> = Vec::new();
let mut data_receivers: Vec<Option<tokio::sync::oneshot::Receiver<Option<Value>>>> = Vec::new();
for _ in 0..last_idx {
let (tx, rx) = tokio::sync::oneshot::channel();
data_senders.push(Some(tx));
data_receivers.push(Some(rx));
}
let mut handles: Vec<tokio::task::JoinHandle<(ExecResult, ExecContext)>> = Vec::with_capacity(stage_count);
let mut stage0_took_session_stdin = false;
let mut stage0_took_session_pipe_stdin = false;
for (i, stage) in stages.iter().enumerate() {
let mut stage_ctx = ctx.child_for_pipeline();
let stage = stage.clone();
let task_dispatcher: Arc<dyn CommandDispatcher> = dispatcher.fork_attached().await;
let stdin_setup = setup_stdin_redirects_for(&stage, &mut stage_ctx, dispatcher).await;
if i == 0 {
let redirect_set_stdin = stage_ctx.stdin.is_some();
stage0_took_session_stdin = !redirect_set_stdin;
if !redirect_set_stdin {
stage_ctx.stdin = ctx.stdin.take();
}
if stage_ctx.stdin_data.is_none() {
stage_ctx.stdin_data = ctx.stdin_data.take();
}
if !redirect_set_stdin && stage_ctx.pipe_stdin.is_none() {
stage_ctx.pipe_stdin = ctx.pipe_stdin.take();
stage0_took_session_pipe_stdin = true;
}
} else {
stage_ctx.pipe_stdin = pipe_readers[i - 1].take();
}
if i < last_idx {
stage_ctx.pipe_stdout = pipe_writers[i].take();
}
stage_ctx.pipeline_position = match i {
0 => PipelinePosition::First,
n if n == last_idx => PipelinePosition::Last,
_ => PipelinePosition::Middle,
};
let data_sender = if i < last_idx { data_senders[i].take() } else { None };
let data_receiver = if i > 0 { data_receivers[i - 1].take() } else { None };
let handle: tokio::task::JoinHandle<(ExecResult, ExecContext)> =
tokio::spawn(crate::telemetry::bind_current_context(async move {
if let Err(e) = stdin_setup {
return (ExecResult::failure(1, e), stage_ctx);
}
stage_ctx.stdin_data_rx = data_receiver;
let mut result = match dispatch_stage(&stage, &mut stage_ctx, &*task_dispatcher).await {
Ok(result) => result,
Err(e) => ExecResult::failure(1, e.to_string()),
};
result = apply_redirects(result, stage.redirects(), &stage_ctx, &*task_dispatcher).await;
if !result.err.is_empty() {
if let Some(ref stderr) = stage_ctx.stderr {
stderr.write_str(&result.err);
result.err.clear();
}
}
if let Some(tx) = data_sender {
let _ = tx.send(result.data.clone());
}
if let Some(mut pipe_out) = stage_ctx.pipe_stdout.take() {
let bytes: Vec<u8> = if let Some(b) = result.out_bytes() {
b.to_vec()
} else if let Some(output) = result.take_output_for_stream() {
let mut buf = Vec::new();
if output.write_canonical(&mut buf, None).is_err() {
buf = output.to_canonical_string().into_bytes();
}
buf
} else {
result.text_out().into_owned().into_bytes()
};
if !bytes.is_empty() {
let _ = pipe_out.write_all(&bytes).await;
let _ = pipe_out.shutdown().await;
}
}
(result, stage_ctx)
}));
handles.push(handle);
}
let mut last_result = ExecResult::success("");
let mut panics: Vec<String> = Vec::new();
let mut codes: Vec<i64> = vec![1; handles.len()];
for (i, handle) in handles.into_iter().enumerate() {
match handle.await {
Ok((result, mut stage_ctx)) => {
codes[i] = result.code;
if i == 0 && stage0_took_session_stdin {
ctx.stdin = stage_ctx.stdin.take();
}
if i == 0 && stage0_took_session_pipe_stdin {
ctx.pipe_stdin = stage_ctx.pipe_stdin.take();
}
if i == last_idx {
last_result = result;
ctx.scope = stage_ctx.scope;
ctx.cwd = stage_ctx.cwd;
ctx.prev_cwd = stage_ctx.prev_cwd;
ctx.aliases = stage_ctx.aliases;
}
}
Err(e) => {
panics.push(format!("stage {}: {}", i, e));
}
}
}
ctx.scope.set_pipestatus(&codes);
if !panics.is_empty() {
last_result = ExecResult::failure(
1,
format!("pipeline stage(s) panicked: {}", panics.join("; ")),
);
}
last_result
}
}
pub fn select_leaf<'a>(schema: &'a ToolSchema, args: &[Arg]) -> anyhow::Result<&'a ToolSchema> {
let root_lookup = schema_param_lookup(schema);
let is_root_value_flag = |name: &str| -> bool {
root_lookup.get(name).is_some_and(|(_, typ, ..)| !is_bool_type(typ))
};
let mut node = schema;
let mut skip_next_positional = false;
for arg in args {
match arg {
Arg::DoubleDash => break,
Arg::LongFlag(name) if is_root_value_flag(name) => skip_next_positional = true,
Arg::ShortFlag(name) if is_root_value_flag(name) => skip_next_positional = true,
Arg::Positional(expr) => {
if skip_next_positional {
skip_next_positional = false;
continue; }
if node.subcommands.is_empty() {
break; }
match classify_subcommand_positional(expr) {
SubcommandWord::Word(word) => {
match node.subcommands.iter().find(|c| c.matches_command(word)) {
Some(child) => node = child, None => break, }
}
SubcommandWord::OtherLiteral => break,
SubcommandWord::Computed(kind) => anyhow::bail!(
"{}: a subcommand name is required here, but got {kind}. \
Subcommands must be literal words — spell it out \
(e.g. `{} <subcommand> …`) or use the `--flag=value` form.",
node.name,
schema.name
),
}
}
_ => {}
}
}
Ok(node)
}
enum SubcommandWord<'a> {
Word(&'a str),
OtherLiteral,
Computed(&'static str),
}
fn classify_subcommand_positional(expr: &Expr) -> SubcommandWord<'_> {
match expr {
Expr::Literal(Value::String(s)) => SubcommandWord::Word(s),
Expr::Literal(_) => SubcommandWord::OtherLiteral,
Expr::CommandSubst(_) | Expr::Command(_) => SubcommandWord::Computed("a command substitution `$(…)`"),
Expr::VarRef(_)
| Expr::VarWithDefault { .. }
| Expr::VarLength(_)
| Expr::Positional(_)
| Expr::AllArgs
| Expr::ArgCount
| Expr::CurrentPid
| Expr::LastExitCode => SubcommandWord::Computed("a variable reference"),
Expr::Interpolated(_) | Expr::HereDocBody { .. } => SubcommandWord::Computed("an interpolated string"),
Expr::GlobPattern(_) => SubcommandWord::Computed("a glob pattern"),
Expr::Arithmetic(_) => SubcommandWord::Computed("an arithmetic expansion"),
_ => SubcommandWord::Computed("a value computed at runtime"),
}
}
pub fn schema_param_lookup(schema: &ToolSchema) -> HashMap<String, (&str, &str, usize, bool)> {
let mut map = HashMap::new();
for p in schema.params.iter().filter(|p| !p.positional) {
map.insert(p.name.clone(), (p.name.as_str(), p.param_type.as_str(), p.consumes, p.repeatable));
for alias in &p.aliases {
let stripped = alias.trim_start_matches('-');
map.insert(stripped.to_string(), (p.name.as_str(), p.param_type.as_str(), p.consumes, p.repeatable));
}
}
map
}
pub fn is_bool_type(param_type: &str) -> bool {
matches!(param_type.to_lowercase().as_str(), "bool" | "boolean")
}
struct SyncEvalSource<'a> {
ctx: &'a ExecContext,
}
#[async_trait::async_trait]
impl crate::kernel::ArgValueSource for SyncEvalSource<'_> {
async fn eval(&self, expr: &Expr) -> anyhow::Result<Option<Value>> {
eval_simple_expr(expr, self.ctx).map_err(|e| anyhow::anyhow!(e))
}
async fn expand_glob(&self, _pattern: &str) -> anyhow::Result<Option<Vec<String>>> {
Ok(None)
}
async fn home(&self) -> Option<String> {
None
}
}
pub async fn build_tool_args(
args: &[Arg],
ctx: &ExecContext,
schema: Option<&ToolSchema>,
) -> Result<ToolArgs, String> {
crate::kernel::bind_tool_args(args, schema, &SyncEvalSource { ctx })
.await
.map_err(|e| e.to_string())
}
pub(crate) fn eval_simple_expr(expr: &Expr, ctx: &ExecContext) -> Result<Option<Value>, String> {
match expr {
Expr::Literal(value) => Ok(Some(eval_literal(value, ctx))),
Expr::VarRef(path) => match ctx.scope.resolve_path(path) {
Ok(v) => Ok(Some(v)),
Err(PathError::UndefinedRoot(_)) if path.segments.len() <= 1 => Ok(None),
Err(PathError::UndefinedRoot(_)) => Err(format!(
"{}: undefined variable",
crate::interpreter::format_path(path)
)),
Err(PathError::Absence(msg)) | Err(PathError::Shape(msg)) => Err(msg),
},
Expr::Interpolated(parts) => Ok(Some(Value::String(eval_string_parts_sync(parts, ctx)?))),
Expr::VarLength(path) => {
crate::interpreter::resolve_length(&ctx.scope, path).map(|n| Some(Value::Int(n)))
}
Expr::VarWithDefault { path, default } => {
match crate::interpreter::resolve_default(&ctx.scope, path)? {
Some(value) => Ok(Some(value)),
None => Ok(Some(Value::String(eval_string_parts_sync(default, ctx)?))),
}
}
Expr::GlobPattern(s) => Ok(Some(Value::String(s.clone()))),
Expr::Arithmetic(expr_str) => arithmetic::eval_arithmetic(expr_str, &ctx.scope)
.map(|n| Some(Value::Int(n)))
.map_err(|e| format!("arithmetic error: {e}")),
Expr::HereDocBody { parts, strip_tabs } => {
let mut asm = crate::interpreter::HeredocAssembler::new(*strip_tabs);
for sp in parts {
match &sp.part {
crate::ast::StringPart::Literal(s) => asm.push_literal(s),
other => {
let s = eval_string_parts_sync(std::slice::from_ref(other), ctx)?;
asm.push_interpolated(&s);
}
}
}
Ok(Some(Value::String(asm.into_string())))
}
Expr::CommandSubst(_) | Expr::Command(_) => Err(
"command substitution `$(...)` is not supported in a scatter/gather flag value here; \
assign it to a variable first (e.g. `n=$(...); scatter --limit $n`)"
.to_string(),
),
_ => Ok(None), }
}
fn eval_literal(value: &Value, _ctx: &ExecContext) -> Value {
value.clone()
}
fn eval_string_parts_sync(parts: &[crate::ast::StringPart], ctx: &ExecContext) -> Result<String, String> {
let mut result = String::new();
for part in parts {
match part {
crate::ast::StringPart::Literal(s) => result.push_str(s),
crate::ast::StringPart::Var(path) => match ctx.scope.resolve_path(path) {
Ok(value) => result.push_str(
&crate::interpreter::value_to_text_sink(&value).map_err(|e| e.to_string())?,
),
Err(PathError::UndefinedRoot(_)) => {}
Err(PathError::Absence(msg)) | Err(PathError::Shape(msg)) => return Err(msg),
},
crate::ast::StringPart::VarWithDefault { path, default } => {
match crate::interpreter::resolve_default(&ctx.scope, path)? {
Some(value) => result.push_str(
&crate::interpreter::value_to_text_sink(&value).map_err(|e| e.to_string())?,
),
None => result.push_str(&eval_string_parts_sync(default, ctx)?),
}
}
crate::ast::StringPart::VarLength(path) => {
let len = crate::interpreter::resolve_length(&ctx.scope, path)?;
result.push_str(&len.to_string());
}
crate::ast::StringPart::Positional(n) => {
if let Some(s) = ctx.scope.get_positional(*n) {
result.push_str(s);
}
}
crate::ast::StringPart::AllArgs => {
result.push_str(&ctx.scope.all_args().join(" "));
}
crate::ast::StringPart::ArgCount => {
result.push_str(&ctx.scope.arg_count().to_string());
}
crate::ast::StringPart::Arithmetic(expr) => {
let value = arithmetic::eval_arithmetic(expr, &ctx.scope)
.map_err(|e| format!("arithmetic error: {e}"))?;
result.push_str(&value.to_string());
}
crate::ast::StringPart::CommandSubst(_) => {
return Err(
"command substitution `$(...)` is not supported inside a scatter/gather \
flag's interpolated value here; assign it to a variable first"
.to_string(),
);
}
crate::ast::StringPart::LastExitCode => {
result.push_str(&ctx.scope.last_result().code.to_string());
}
crate::ast::StringPart::CurrentPid => {
result.push_str(&ctx.scope.pid().to_string());
}
}
}
Ok(result)
}
fn find_scatter_gather(stages: &[PipelineStage]) -> Option<(usize, usize)> {
let named = |name: &str| {
stages
.iter()
.position(|s| s.as_command().is_some_and(|c| c.name == name))
};
let scatter_idx = named("scatter")?;
let gather_idx = named("gather")?;
if gather_idx > scatter_idx {
Some((scatter_idx, gather_idx))
} else {
None
}
}
#[cfg(test)]
mod select_leaf_tests {
use super::*;
use crate::tools::ParamSchema;
fn kj_schema() -> ToolSchema {
ToolSchema::new("kj", "kaijutsu")
.param(ParamSchema::new("confirm", "string"))
.param(ParamSchema::new("verbose", "bool"))
.subcommand(
ToolSchema::new("context", "context ops")
.with_command_aliases(["ctx"])
.subcommand(ToolSchema::new("list", "list").with_command_aliases(["ls"]))
.subcommand(
ToolSchema::new("create", "create").param(
ParamSchema::new("type", "string").with_aliases(["t"]),
),
),
)
}
fn word(s: &str) -> Arg {
Arg::Positional(Expr::Literal(Value::String(s.to_string())))
}
#[test]
fn flat_tool_returns_root() {
let schema = ToolSchema::new("cat", "concat")
.param(ParamSchema::required("path", "string", "f").positional());
let leaf = select_leaf(&schema, &[word("foo.txt")]).expect("flat ok");
assert_eq!(leaf.name, "cat");
}
#[test]
fn single_hop() {
let schema = kj_schema();
let leaf = select_leaf(&schema, &[word("context")]).expect("ok");
assert_eq!(leaf.name, "context");
}
#[test]
fn two_hops() {
let schema = kj_schema();
let leaf = select_leaf(&schema, &[word("context"), word("create")]).expect("ok");
assert_eq!(leaf.name, "create");
assert!(leaf.params.iter().any(|p| p.name == "type"), "leaf has --type");
}
#[test]
fn alias_hops_route() {
let schema = kj_schema();
let leaf = select_leaf(&schema, &[word("ctx"), word("ls")]).expect("ok");
assert_eq!(leaf.name, "list");
}
#[test]
fn unknown_subcommand_stops_at_current_node() {
let schema = kj_schema();
let leaf = select_leaf(&schema, &[word("context"), word("nonesuch")]).expect("ok");
assert_eq!(leaf.name, "context");
}
#[test]
fn root_bool_flag_before_path_does_not_disrupt_routing() {
let schema = kj_schema();
let args = vec![Arg::LongFlag("verbose".into()), word("context"), word("create")];
let leaf = select_leaf(&schema, &args).expect("ok");
assert_eq!(leaf.name, "create");
}
#[test]
fn root_value_flag_space_form_before_path_skips_its_value() {
let schema = kj_schema();
let args = vec![
Arg::LongFlag("confirm".into()),
word("token"),
word("context"),
word("create"),
];
let leaf = select_leaf(&schema, &args).expect("ok");
assert_eq!(leaf.name, "create");
}
#[test]
fn leaf_value_flag_after_path_routes_to_leaf() {
let schema = kj_schema();
let args = vec![
word("context"),
word("create"),
Arg::LongFlag("type".into()),
word("x"),
];
let leaf = select_leaf(&schema, &args).expect("ok");
assert_eq!(leaf.name, "create");
assert!(leaf.params.iter().any(|p| p.name == "type"));
}
#[test]
fn double_dash_stops_routing() {
let schema = kj_schema();
let leaf = select_leaf(&schema, &[Arg::DoubleDash, word("context")]).expect("ok");
assert_eq!(leaf.name, "kj");
}
#[test]
fn computed_subcommand_selector_errors() {
let schema = kj_schema();
let args = vec![Arg::Positional(Expr::CommandSubst(vec![
crate::ast::Stmt::Command(crate::ast::Command {
name: "echo".into(),
args: vec![],
redirects: vec![],
}),
]))];
let err = select_leaf(&schema, &args).expect_err("must error");
let msg = err.to_string();
assert!(msg.contains("subcommand name is required"), "got: {msg}");
assert!(msg.contains("command substitution"), "names the cause: {msg}");
}
#[test]
fn variable_subcommand_selector_errors() {
let schema = kj_schema();
let args = vec![Arg::Positional(Expr::VarRef(crate::ast::VarPath::simple("sub")))];
let err = select_leaf(&schema, &args).expect_err("must error");
assert!(err.to_string().contains("variable reference"), "got: {err}");
}
#[test]
fn computed_positional_after_leaf_is_fine() {
let schema = kj_schema();
let args = vec![
word("context"),
word("list"),
Arg::Positional(Expr::CommandSubst(vec![crate::ast::Stmt::Command(
crate::ast::Command { name: "echo".into(), args: vec![], redirects: vec![] },
)])),
];
let leaf = select_leaf(&schema, &args).expect("ok");
assert_eq!(leaf.name, "list");
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::dispatch::BackendDispatcher;
use crate::tools::register_builtins;
use crate::vfs::{Filesystem, MemoryFs, VfsRouter};
use std::path::Path;
async fn make_runner_and_ctx() -> (PipelineRunner, ExecContext, BackendDispatcher) {
let mut tools = ToolRegistry::new();
register_builtins(&mut tools);
let tools = Arc::new(tools);
let runner = PipelineRunner::new(tools.clone());
let dispatcher = BackendDispatcher::new(tools.clone());
let mut vfs = VfsRouter::new();
let mem = MemoryFs::new();
mem.write(Path::new("test.txt"), b"hello\nworld\nfoo").await.unwrap();
vfs.mount("/", mem);
let ctx = ExecContext::with_vfs_and_tools(Arc::new(vfs), tools);
(runner, ctx, dispatcher)
}
fn stages(commands: impl IntoIterator<Item = Command>) -> Vec<PipelineStage> {
commands.into_iter().map(PipelineStage::Command).collect()
}
fn make_cmd(name: &str, args: Vec<&str>) -> Command {
Command {
name: name.to_string(),
args: args.iter().map(|s| Arg::Positional(Expr::Literal(Value::String(s.to_string())))).collect(),
redirects: vec![],
}
}
#[tokio::test]
async fn test_single_command() {
let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
let cmd = make_cmd("echo", vec!["hello"]);
let result = runner.run(&stages([cmd]), &mut ctx, &dispatcher).await;
assert!(result.ok());
assert_eq!(result.text_out().trim(), "hello");
}
#[tokio::test]
async fn test_pipeline_echo_grep() {
let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
let echo_cmd = Command {
name: "echo".to_string(),
args: vec![Arg::Positional(Expr::Literal(Value::String("hello\nworld".to_string())))],
redirects: vec![],
};
let grep_cmd = Command {
name: "grep".to_string(),
args: vec![Arg::Positional(Expr::Literal(Value::String("world".to_string())))],
redirects: vec![],
};
let result = runner.run(&stages([echo_cmd, grep_cmd]), &mut ctx, &dispatcher).await;
assert!(result.ok());
assert_eq!(result.text_out().trim(), "world");
}
#[tokio::test]
async fn test_pipeline_cat_grep() {
let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
let cat_cmd = make_cmd("cat", vec!["/test.txt"]);
let grep_cmd = Command {
name: "grep".to_string(),
args: vec![Arg::Positional(Expr::Literal(Value::String("hello".to_string())))],
redirects: vec![],
};
let result = runner.run(&stages([cat_cmd, grep_cmd]), &mut ctx, &dispatcher).await;
assert!(result.ok());
assert!(result.text_out().contains("hello"));
}
#[tokio::test]
async fn test_command_not_found() {
let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
let cmd = make_cmd("nonexistent", vec![]);
let result = runner.run(&stages([cmd]), &mut ctx, &dispatcher).await;
assert!(!result.ok());
assert_eq!(result.code, 127);
#[cfg(feature = "subprocess")]
assert!(
result.err.contains("not found"),
"expected a not-found report, got {:?}",
result.err
);
#[cfg(not(feature = "subprocess"))]
assert!(
result.err.contains("external commands are"),
"expected the unavailable-externals report, got {:?}",
result.err
);
}
#[tokio::test]
async fn test_pipeline_continues_on_failure() {
let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
let cat_cmd = make_cmd("cat", vec!["/nonexistent"]);
let grep_cmd = Command {
name: "grep".to_string(),
args: vec![Arg::Positional(Expr::Literal(Value::String("hello".to_string())))],
redirects: vec![],
};
let result = runner.run(&stages([cat_cmd, grep_cmd]), &mut ctx, &dispatcher).await;
assert!(!result.ok());
}
#[tokio::test]
async fn test_pipeline_last_command_exit_code() {
let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
let echo_cmd = make_cmd("echo", vec!["hello"]);
let cat_cmd = make_cmd("cat", vec![]);
let result = runner.run(&stages([echo_cmd, cat_cmd]), &mut ctx, &dispatcher).await;
assert!(result.ok());
assert!(result.text_out().contains("hello"));
}
#[tokio::test]
async fn test_empty_pipeline() {
let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
let result = runner.run(&stages([]), &mut ctx, &dispatcher).await;
assert!(result.ok());
}
#[test]
fn test_find_scatter_gather_both_present() {
let commands = vec![
make_cmd("echo", vec!["a"]),
make_cmd("scatter", vec![]),
make_cmd("process", vec![]),
make_cmd("gather", vec![]),
];
let result = find_scatter_gather(&stages(commands));
assert_eq!(result, Some((1, 3)));
}
#[test]
fn test_find_scatter_gather_no_scatter() {
let commands = vec![
make_cmd("echo", vec!["a"]),
make_cmd("gather", vec![]),
];
let result = find_scatter_gather(&stages(commands));
assert!(result.is_none());
}
#[test]
fn test_find_scatter_gather_no_gather() {
let commands = vec![
make_cmd("echo", vec!["a"]),
make_cmd("scatter", vec![]),
];
let result = find_scatter_gather(&stages(commands));
assert!(result.is_none());
}
#[test]
fn test_find_scatter_gather_wrong_order() {
let commands = vec![
make_cmd("gather", vec![]),
make_cmd("scatter", vec![]),
];
let result = find_scatter_gather(&stages(commands));
assert!(result.is_none());
}
#[tokio::test]
async fn test_scatter_gather_simple() {
let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
let split_cmd = Command {
name: "split".to_string(),
args: vec![Arg::Positional(Expr::Literal(Value::String("a b c".to_string())))],
redirects: vec![],
};
let scatter_cmd = make_cmd("scatter", vec![]);
let process_cmd = Command {
name: "echo".to_string(),
args: vec![Arg::Positional(Expr::VarRef(crate::ast::VarPath::simple("ITEM")))],
redirects: vec![],
};
let gather_cmd = make_cmd("gather", vec![]);
let result = runner.run(&stages([split_cmd, scatter_cmd, process_cmd, gather_cmd]), &mut ctx, &dispatcher).await;
assert!(result.ok(), "scatter with structured data should succeed: {}", result.err);
assert!(result.text_out().contains("a"));
assert!(result.text_out().contains("b"));
assert!(result.text_out().contains("c"));
}
#[tokio::test]
async fn test_scatter_gather_empty_input() {
let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
let echo_cmd = Command {
name: "echo".to_string(),
args: vec![Arg::Positional(Expr::Literal(Value::String("".to_string())))],
redirects: vec![],
};
let scatter_cmd = make_cmd("scatter", vec![]);
let process_cmd = Command {
name: "echo".to_string(),
args: vec![Arg::Positional(Expr::VarRef(crate::ast::VarPath::simple("ITEM")))],
redirects: vec![],
};
let gather_cmd = make_cmd("gather", vec![]);
let result = runner.run(&stages([echo_cmd, scatter_cmd, process_cmd, gather_cmd]), &mut ctx, &dispatcher).await;
assert!(result.ok());
assert!(result.text_out().trim().is_empty());
}
#[tokio::test]
async fn test_scatter_gather_with_structured_stdin() {
let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
let data = Value::Json(serde_json::json!(["x", "y", "z"]));
ctx.set_stdin_with_data("x\ny\nz".to_string(), Some(data));
let scatter_cmd = make_cmd("scatter", vec![]);
let process_cmd = Command {
name: "echo".to_string(),
args: vec![Arg::Positional(Expr::VarRef(crate::ast::VarPath::simple("ITEM")))],
redirects: vec![],
};
let gather_cmd = make_cmd("gather", vec![]);
let result = runner.run(&stages([scatter_cmd, process_cmd, gather_cmd]), &mut ctx, &dispatcher).await;
assert!(result.ok(), "scatter with structured stdin should succeed: {}", result.err);
assert!(result.text_out().contains("x"));
assert!(result.text_out().contains("y"));
assert!(result.text_out().contains("z"));
}
#[tokio::test]
async fn test_scatter_gather_json_input() {
let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
let data = Value::Json(serde_json::json!(["one", "two", "three"]));
ctx.set_stdin_with_data(r#"["one", "two", "three"]"#.to_string(), Some(data));
let scatter_cmd = make_cmd("scatter", vec![]);
let process_cmd = Command {
name: "echo".to_string(),
args: vec![Arg::Positional(Expr::VarRef(crate::ast::VarPath::simple("ITEM")))],
redirects: vec![],
};
let gather_cmd = make_cmd("gather", vec![]);
let result = runner.run(&stages([scatter_cmd, process_cmd, gather_cmd]), &mut ctx, &dispatcher).await;
assert!(result.ok(), "scatter with JSON data should succeed: {}", result.err);
assert!(result.text_out().contains("one"));
assert!(result.text_out().contains("two"));
assert!(result.text_out().contains("three"));
}
#[tokio::test]
async fn test_scatter_gather_with_post_gather() {
let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
let split_cmd = Command {
name: "split".to_string(),
args: vec![Arg::Positional(Expr::Literal(Value::String("a b".to_string())))],
redirects: vec![],
};
let scatter_cmd = make_cmd("scatter", vec![]);
let process_cmd = Command {
name: "echo".to_string(),
args: vec![Arg::Positional(Expr::VarRef(crate::ast::VarPath::simple("ITEM")))],
redirects: vec![],
};
let gather_cmd = make_cmd("gather", vec![]);
let grep_cmd = Command {
name: "grep".to_string(),
args: vec![Arg::Positional(Expr::Literal(Value::String("a".to_string())))],
redirects: vec![],
};
let result = runner.run(&stages([split_cmd, scatter_cmd, process_cmd, gather_cmd, grep_cmd]), &mut ctx, &dispatcher).await;
assert!(result.ok(), "scatter with post_gather should succeed: {}", result.err);
assert!(result.text_out().contains("a"));
assert!(!result.text_out().contains("b"));
}
#[tokio::test]
async fn test_scatter_custom_var_name() {
let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
let data = Value::Json(serde_json::json!(["test1", "test2"]));
ctx.set_stdin_with_data("test1\ntest2".to_string(), Some(data));
let scatter_cmd = Command {
name: "scatter".to_string(),
args: vec![Arg::Named {
key: "as".to_string(),
value: Expr::Literal(Value::String("URL".to_string())),
}],
redirects: vec![],
};
let process_cmd = Command {
name: "echo".to_string(),
args: vec![Arg::Positional(Expr::VarRef(crate::ast::VarPath::simple("URL")))],
redirects: vec![],
};
let gather_cmd = make_cmd("gather", vec![]);
let result = runner.run(&stages([scatter_cmd, process_cmd, gather_cmd]), &mut ctx, &dispatcher).await;
assert!(result.ok(), "scatter with custom var should succeed: {}", result.err);
assert!(result.text_out().contains("test1"));
assert!(result.text_out().contains("test2"));
}
#[tokio::test]
async fn test_pipeline_routes_through_backend() {
use crate::backend::testing::MockBackend;
use std::sync::atomic::Ordering;
let (backend, call_count) = MockBackend::new();
let backend: std::sync::Arc<dyn crate::backend::KernelBackend> = std::sync::Arc::new(backend);
let mut ctx = crate::tools::ExecContext::with_backend(backend);
let tools = std::sync::Arc::new(ToolRegistry::new());
let runner = PipelineRunner::new(tools.clone());
let dispatcher = BackendDispatcher::new(tools);
let cmd = make_cmd("test-tool", vec!["arg1"]);
let result = runner.run(&stages([cmd]), &mut ctx, &dispatcher).await;
assert!(result.ok(), "Mock backend should return success");
assert_eq!(call_count.load(Ordering::SeqCst), 1, "call_tool should be invoked once");
assert!(result.text_out().contains("mock executed"), "Output should be from mock backend");
}
#[tokio::test]
async fn test_multi_command_pipeline_routes_through_backend() {
use crate::backend::testing::MockBackend;
use std::sync::atomic::Ordering;
let (backend, call_count) = MockBackend::new();
let backend: std::sync::Arc<dyn crate::backend::KernelBackend> = std::sync::Arc::new(backend);
let mut ctx = crate::tools::ExecContext::with_backend(backend);
let tools = std::sync::Arc::new(ToolRegistry::new());
let runner = PipelineRunner::new(tools.clone());
let dispatcher = BackendDispatcher::new(tools);
let cmd1 = make_cmd("tool1", vec![]);
let cmd2 = make_cmd("tool2", vec![]);
let cmd3 = make_cmd("tool3", vec![]);
let result = runner.run(&stages([cmd1, cmd2, cmd3]), &mut ctx, &dispatcher).await;
assert!(result.ok());
assert_eq!(call_count.load(Ordering::SeqCst), 3, "call_tool should be invoked for each command");
}
#[tokio::test]
async fn backend_dispatcher_scalar_data_matches_production_unwrap() {
use crate::backend::testing::MockBackend;
use crate::backend::ToolResult;
let (mock, _calls) = MockBackend::new();
let backend = mock.with_tool_result(|_name| Ok(ToolResult::with_data("", serde_json::json!(42))));
let backend: Arc<dyn crate::backend::KernelBackend> = Arc::new(backend);
let mut ctx = ExecContext::with_backend(backend);
let dispatcher = BackendDispatcher::new(Arc::new(ToolRegistry::new()));
let cmd = make_cmd("embedder_tool", vec![]);
let result = dispatcher.dispatch(&cmd, &mut ctx).await.expect("dispatch");
assert_eq!(
result.data,
Some(Value::Int(42)),
"a scalar ToolResult.data must unwrap to a native Value, matching \
the production From<ToolResult> path — not stay Value::Json(42)"
);
}
#[tokio::test]
async fn backend_dispatcher_envelope_shaped_data_stays_structured() {
use crate::backend::testing::MockBackend;
use crate::backend::ToolResult;
let envelope = kaish_types::bytes_to_envelope(&[1u8, 2, 3]);
let (mock, _calls) = MockBackend::new();
let backend = mock.with_tool_result(move |_name| Ok(ToolResult::with_data("", envelope.clone())));
let backend: Arc<dyn crate::backend::KernelBackend> = Arc::new(backend);
let mut ctx = ExecContext::with_backend(backend);
let dispatcher = BackendDispatcher::new(Arc::new(ToolRegistry::new()));
let cmd = make_cmd("embedder_tool", vec![]);
let result = dispatcher.dispatch(&cmd, &mut ctx).await.expect("dispatch");
assert!(
matches!(result.data, Some(Value::Json(_))),
"envelope-shaped external data must stay structured, not silently \
decode to Value::Bytes: got {:?}",
result.data
);
}
#[tokio::test]
async fn backend_dispatcher_preserves_did_spill_and_original_code() {
use crate::backend::testing::MockBackend;
use crate::backend::ToolResult;
let (mock, _calls) = MockBackend::new();
let backend = mock.with_tool_result(|_name| {
Ok(ToolResult::success("truncated...")
.with_did_spill(true)
.with_original_code(Some(0)))
});
let backend: Arc<dyn crate::backend::KernelBackend> = Arc::new(backend);
let mut ctx = ExecContext::with_backend(backend);
let dispatcher = BackendDispatcher::new(Arc::new(ToolRegistry::new()));
let cmd = make_cmd("embedder_tool", vec![]);
let result = dispatcher.dispatch(&cmd, &mut ctx).await.expect("dispatch");
assert!(result.did_spill, "did_spill must survive the backend seam");
assert_eq!(result.original_code, Some(0), "original_code must survive the backend seam");
}
use crate::tools::{ParamSchema, ToolSchema};
fn make_test_schema() -> ToolSchema {
ToolSchema::new("test-tool", "A test tool for schema-aware parsing")
.param(ParamSchema::required("query", "string", "Search query"))
.param(ParamSchema::optional("limit", "int", Value::Int(10), "Max results"))
.param(ParamSchema::optional("verbose", "bool", Value::Bool(false), "Verbose output"))
.param(ParamSchema::optional("output", "string", Value::String("stdout".into()), "Output destination"))
.with_positional_mapping()
}
fn make_minimal_ctx() -> ExecContext {
let mut vfs = VfsRouter::new();
vfs.mount("/", MemoryFs::new());
ExecContext::new(Arc::new(vfs))
}
fn test_dispatcher() -> BackendDispatcher {
BackendDispatcher::new(Arc::new(ToolRegistry::new()))
}
#[tokio::test]
async fn test_schema_aware_string_arg() {
let args = vec![
Arg::LongFlag("query".to_string()),
Arg::Positional(Expr::Literal(Value::String("test".to_string()))),
];
let schema = make_test_schema();
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert!(tool_args.flags.is_empty(), "No flags should be set");
assert!(tool_args.positional.is_empty(), "No positionals - consumed by --query");
assert_eq!(
tool_args.named.get("query"),
Some(&Value::String("test".to_string())),
"--query should consume 'test' as its value"
);
}
#[tokio::test]
async fn test_schema_aware_bool_flag() {
let args = vec![
Arg::LongFlag("verbose".to_string()),
];
let schema = make_test_schema();
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert!(tool_args.flags.contains("verbose"), "--verbose should be a flag");
assert!(tool_args.named.is_empty(), "No named args");
assert!(tool_args.positional.is_empty(), "No positionals");
}
#[tokio::test]
async fn test_schema_aware_mixed() {
let args = vec![
Arg::Positional(Expr::Literal(Value::String("file.txt".to_string()))),
Arg::LongFlag("output".to_string()),
Arg::Positional(Expr::Literal(Value::String("out.txt".to_string()))),
Arg::LongFlag("verbose".to_string()),
];
let schema = make_test_schema();
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert!(tool_args.positional.is_empty(), "file.txt consumed as query param");
assert_eq!(
tool_args.named.get("query"),
Some(&Value::String("file.txt".to_string()))
);
assert_eq!(
tool_args.named.get("output"),
Some(&Value::String("out.txt".to_string()))
);
assert!(tool_args.flags.contains("verbose"));
}
#[tokio::test]
async fn test_schema_aware_multiple_string_args() {
let args = vec![
Arg::LongFlag("query".to_string()),
Arg::Positional(Expr::Literal(Value::String("test".to_string()))),
Arg::LongFlag("output".to_string()),
Arg::Positional(Expr::Literal(Value::String("result.json".to_string()))),
Arg::LongFlag("verbose".to_string()),
Arg::LongFlag("limit".to_string()),
Arg::Positional(Expr::Literal(Value::Int(5))),
];
let schema = make_test_schema();
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert!(tool_args.positional.is_empty(), "All positionals consumed");
assert_eq!(
tool_args.named.get("query"),
Some(&Value::String("test".to_string()))
);
assert_eq!(
tool_args.named.get("output"),
Some(&Value::String("result.json".to_string()))
);
assert_eq!(
tool_args.named.get("limit"),
Some(&Value::Int(5))
);
assert!(tool_args.flags.contains("verbose"));
}
#[tokio::test]
async fn test_schema_aware_double_dash() {
let args = vec![
Arg::LongFlag("output".to_string()),
Arg::Positional(Expr::Literal(Value::String("out.txt".to_string()))),
Arg::DoubleDash,
Arg::Positional(Expr::Literal(Value::String("--this-is-data".to_string()))),
];
let schema = make_test_schema();
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert_eq!(
tool_args.named.get("output"),
Some(&Value::String("out.txt".to_string()))
);
assert_eq!(
tool_args.positional,
vec![Value::String("--this-is-data".to_string())]
);
}
#[tokio::test]
async fn word_assign_binary_value_is_loud_not_placeholder() {
let args = vec![Arg::WordAssign {
key: "if".to_string(),
value: Expr::Literal(Value::Bytes(vec![0xff, 0x00, 0xfe])),
}];
let ctx = make_minimal_ctx();
let err = build_tool_args(&args, &ctx, None).await.expect_err("binary WordAssign must error");
assert!(
err.contains("cannot be used as"),
"error should name the binary problem, got {err:?}"
);
}
#[tokio::test]
async fn word_assign_before_double_dash_binds_named_for_export() {
let args = vec![Arg::WordAssign {
key: "A".to_string(),
value: Expr::Literal(Value::String("1".to_string())),
}];
let schema = ToolSchema::new("export", "export");
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert_eq!(tool_args.named.get("A"), Some(&Value::String("1".to_string())));
assert!(tool_args.positional.is_empty());
}
#[tokio::test]
async fn word_assign_after_double_dash_is_positional_even_for_export() {
let args = vec![
Arg::DoubleDash,
Arg::WordAssign {
key: "A".to_string(),
value: Expr::Literal(Value::String("1".to_string())),
},
];
let schema = ToolSchema::new("export", "export");
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert!(
tool_args.named.is_empty(),
"past `--`, A=1 must not become a named assignment: {:?}",
tool_args.named
);
assert_eq!(tool_args.positional, vec![Value::String("A=1".to_string())]);
}
#[tokio::test]
async fn named_true_on_undeclared_flag_flagifies() {
let args = vec![Arg::Named {
key: "json".to_string(),
value: Expr::Literal(Value::Bool(true)),
}];
let schema = make_test_schema(); let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert!(tool_args.flags.contains("json"), "flags: {:?}", tool_args.flags);
assert!(!tool_args.named.contains_key("json"), "named: {:?}", tool_args.named);
}
#[tokio::test]
async fn named_false_on_undeclared_flag_drops() {
let args = vec![Arg::Named {
key: "json".to_string(),
value: Expr::Literal(Value::Bool(false)),
}];
let schema = make_test_schema();
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert!(!tool_args.flags.contains("json"));
assert!(!tool_args.named.contains_key("json"));
}
#[tokio::test]
async fn named_true_on_declared_bool_param_flagifies() {
let args = vec![Arg::Named {
key: "verbose".to_string(),
value: Expr::Literal(Value::Bool(true)),
}];
let schema = make_test_schema();
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert!(tool_args.flags.contains("verbose"));
assert!(!tool_args.named.contains_key("verbose"));
}
#[tokio::test]
async fn named_true_on_declared_value_flag_keeps_value() {
let args = vec![Arg::Named {
key: "output".to_string(),
value: Expr::Literal(Value::Bool(true)),
}];
let schema = make_test_schema();
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert_eq!(tool_args.named.get("output"), Some(&Value::Bool(true)));
assert!(!tool_args.flags.contains("output"));
}
#[tokio::test]
async fn named_true_with_no_schema_flagifies() {
let args = vec![Arg::Named {
key: "verbose".to_string(),
value: Expr::Literal(Value::Bool(true)),
}];
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, None).await.expect("build_tool_args");
assert!(tool_args.flags.contains("verbose"));
assert!(!tool_args.named.contains_key("verbose"));
}
#[tokio::test]
async fn test_no_schema_fallback() {
let args = vec![
Arg::LongFlag("query".to_string()),
Arg::Positional(Expr::Literal(Value::String("test".to_string()))),
];
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, None).await.expect("build_tool_args");
assert!(tool_args.flags.contains("query"), "--query should be a flag");
assert_eq!(
tool_args.positional,
vec![Value::String("test".to_string())],
"'test' should be a positional"
);
}
#[tokio::test]
async fn test_unknown_flag_ambiguous_space_value_now_errors_loud() {
let args = vec![
Arg::LongFlag("unknown".to_string()),
Arg::Positional(Expr::Literal(Value::String("value".to_string()))),
];
let schema = make_test_schema();
let ctx = make_minimal_ctx();
let err = build_tool_args(&args, &ctx, Some(&schema))
.await
.expect_err("an undeclared flag immediately before a positional must be ambiguous, not silently bool");
assert!(
err.contains("--unknown is not a declared flag"),
"got: {err}"
);
}
#[tokio::test]
async fn test_unknown_bool_flag_with_no_following_positional_is_fine() {
let args = vec![
Arg::Positional(Expr::Literal(Value::String("value".to_string()))),
Arg::LongFlag("unknown".to_string()),
];
let schema = make_test_schema();
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert!(tool_args.flags.contains("unknown"));
assert!(tool_args.positional.is_empty(), "value consumed as query param");
assert_eq!(
tool_args.named.get("query"),
Some(&Value::String("value".to_string()))
);
}
#[tokio::test]
async fn test_unknown_short_flag_ambiguous_space_value_now_errors_loud() {
let args = vec![
Arg::ShortFlag("t".to_string()),
Arg::Positional(Expr::Literal(Value::String("value".to_string()))),
];
let schema = make_test_schema();
let ctx = make_minimal_ctx();
let err = build_tool_args(&args, &ctx, Some(&schema))
.await
.expect_err("an undeclared short flag immediately before a positional must be ambiguous, not silently bool");
assert!(
err.contains("-t is not a declared flag"),
"got: {err}"
);
}
#[tokio::test]
async fn test_unknown_short_bool_flag_with_no_following_positional_is_fine() {
let args = vec![
Arg::Positional(Expr::Literal(Value::String("value".to_string()))),
Arg::ShortFlag("t".to_string()),
];
let schema = make_test_schema();
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert!(tool_args.flags.contains("t"));
assert!(tool_args.positional.is_empty(), "value consumed as query param");
assert_eq!(
tool_args.named.get("query"),
Some(&Value::String("value".to_string()))
);
}
#[tokio::test]
async fn test_unknown_short_flag_not_ambiguous_without_map_positionals() {
let args = vec![
Arg::ShortFlag("t".to_string()),
Arg::Positional(Expr::Literal(Value::String("value".to_string()))),
];
let schema = ToolSchema::new("test-tool", "A test tool")
.param(ParamSchema::required("query", "string", "Search query"));
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert!(tool_args.flags.contains("t"));
assert_eq!(tool_args.positional, vec![Value::String("value".to_string())]);
}
#[tokio::test]
async fn test_named_args_unchanged() {
let args = vec![
Arg::Named {
key: "query".to_string(),
value: Expr::Literal(Value::String("test".to_string())),
},
Arg::LongFlag("verbose".to_string()),
];
let schema = make_test_schema();
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert_eq!(
tool_args.named.get("query"),
Some(&Value::String("test".to_string()))
);
assert!(tool_args.flags.contains("verbose"));
}
#[tokio::test]
async fn test_short_flags_unchanged() {
let args = vec![
Arg::ShortFlag("la".to_string()),
Arg::Positional(Expr::Literal(Value::String("file.txt".to_string()))),
];
let schema = make_test_schema();
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert!(tool_args.flags.contains("l"));
assert!(tool_args.flags.contains("a"));
assert!(tool_args.positional.is_empty(), "file.txt consumed as query param");
assert_eq!(
tool_args.named.get("query"),
Some(&Value::String("file.txt".to_string()))
);
}
#[tokio::test]
async fn test_flag_at_end_no_value() {
let args = vec![
Arg::Positional(Expr::Literal(Value::String("file.txt".to_string()))),
Arg::LongFlag("output".to_string()),
];
let schema = make_test_schema();
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert!(tool_args.flags.contains("output"));
assert!(tool_args.positional.is_empty(), "file.txt consumed as query param");
assert_eq!(
tool_args.named.get("query"),
Some(&Value::String("file.txt".to_string()))
);
}
#[tokio::test]
async fn test_positional_skips_bool_params() {
let schema = ToolSchema::new("test", "")
.param(ParamSchema::required("query", "string", ""))
.param(ParamSchema::optional(
"verbose",
"bool",
Value::Bool(false),
"",
))
.param(ParamSchema::optional(
"output",
"string",
Value::Null,
"",
))
.with_positional_mapping();
let args = vec![
Arg::Positional(Expr::Literal(Value::String("val1".to_string()))),
Arg::Positional(Expr::Literal(Value::String("val2".to_string()))),
];
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert_eq!(
tool_args.named.get("query"),
Some(&Value::String("val1".to_string()))
);
assert_eq!(
tool_args.named.get("output"),
Some(&Value::String("val2".to_string()))
);
assert!(!tool_args.flags.contains("verbose"));
assert!(tool_args.positional.is_empty());
}
#[tokio::test]
async fn test_positionals_fill_available_slots() {
let args = vec![
Arg::Positional(Expr::Literal(Value::String("val1".to_string()))),
Arg::Positional(Expr::Literal(Value::String("val2".to_string()))),
Arg::Positional(Expr::Literal(Value::String("val3".to_string()))),
];
let schema = make_test_schema(); let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert_eq!(
tool_args.named.get("query"),
Some(&Value::String("val1".to_string()))
);
assert_eq!(
tool_args.named.get("limit"),
Some(&Value::String("val2".to_string()))
);
assert_eq!(
tool_args.named.get("output"),
Some(&Value::String("val3".to_string()))
);
assert!(tool_args.positional.is_empty());
}
#[tokio::test]
async fn test_truly_excess_positionals() {
let schema = ToolSchema::new("test", "")
.param(ParamSchema::required("name", "string", ""))
.with_positional_mapping();
let args = vec![
Arg::Positional(Expr::Literal(Value::String("first".to_string()))),
Arg::Positional(Expr::Literal(Value::String("second".to_string()))),
Arg::Positional(Expr::Literal(Value::String("third".to_string()))),
];
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert_eq!(
tool_args.named.get("name"),
Some(&Value::String("first".to_string()))
);
assert_eq!(
tool_args.positional,
vec![
Value::String("second".to_string()),
Value::String("third".to_string()),
]
);
}
#[tokio::test]
async fn test_double_dash_positional_not_mapped() {
let args = vec![
Arg::Positional(Expr::Literal(Value::String("val1".to_string()))),
Arg::DoubleDash,
Arg::Positional(Expr::Literal(Value::String("val2".to_string()))),
];
let schema = make_test_schema();
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert_eq!(
tool_args.named.get("query"),
Some(&Value::String("val1".to_string()))
);
assert_eq!(
tool_args.positional,
vec![Value::String("val2".to_string())]
);
}
#[tokio::test]
async fn test_all_params_filled_by_flags() {
let args = vec![
Arg::LongFlag("query".to_string()),
Arg::Positional(Expr::Literal(Value::String("search".to_string()))),
Arg::LongFlag("output".to_string()),
Arg::Positional(Expr::Literal(Value::String("out.txt".to_string()))),
Arg::LongFlag("verbose".to_string()),
];
let schema = make_test_schema();
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert_eq!(
tool_args.named.get("query"),
Some(&Value::String("search".to_string()))
);
assert_eq!(
tool_args.named.get("output"),
Some(&Value::String("out.txt".to_string()))
);
assert!(tool_args.flags.contains("verbose"));
assert!(tool_args.positional.is_empty());
}
#[tokio::test]
async fn test_mixed_flags_and_positional_fill() {
let args = vec![
Arg::LongFlag("output".to_string()),
Arg::Positional(Expr::Literal(Value::String("foo".to_string()))),
Arg::Positional(Expr::Literal(Value::String("val1".to_string()))),
];
let schema = make_test_schema();
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert_eq!(
tool_args.named.get("output"),
Some(&Value::String("foo".to_string()))
);
assert_eq!(
tool_args.named.get("query"),
Some(&Value::String("val1".to_string()))
);
assert!(tool_args.positional.is_empty());
}
#[tokio::test]
async fn test_alias_flag_prevents_mapping_overwrite() {
let schema = ToolSchema::new("test", "")
.param(ParamSchema::required("query", "string", "").with_aliases(["-q"]))
.param(ParamSchema::required("output", "string", ""))
.with_positional_mapping();
let args = vec![
Arg::ShortFlag("q".to_string()),
Arg::Positional(Expr::Literal(Value::String("search".to_string()))),
Arg::Positional(Expr::Literal(Value::String("out.txt".to_string()))),
];
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert_eq!(
tool_args.named.get("query"),
Some(&Value::String("search".to_string()))
);
assert_eq!(
tool_args.named.get("output"),
Some(&Value::String("out.txt".to_string()))
);
assert!(tool_args.positional.is_empty());
}
#[tokio::test]
async fn test_builtin_schema_no_positional_mapping() {
let schema = ToolSchema::new("echo", "")
.param(ParamSchema::optional("args", "any", Value::Null, ""))
.param(ParamSchema::optional("no_newline", "bool", Value::Bool(false), ""));
let args = vec![
Arg::Positional(Expr::Literal(Value::String("hello".to_string()))),
Arg::Positional(Expr::Literal(Value::String("world".to_string()))),
];
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert_eq!(
tool_args.positional,
vec![
Value::String("hello".to_string()),
Value::String("world".to_string()),
]
);
assert!(!tool_args.named.contains_key("args"));
}
#[tokio::test]
async fn test_short_flag_with_alias_consumes_value() {
let schema = ToolSchema::new("head", "Output first part of files")
.param(ParamSchema::optional("lines", "int", Value::Int(10), "Number of lines")
.with_aliases(["-n"]));
let args = vec![
Arg::ShortFlag("n".to_string()),
Arg::Positional(Expr::Literal(Value::Int(5))),
Arg::Positional(Expr::Literal(Value::String("/tmp/file.txt".to_string()))),
];
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert!(tool_args.flags.is_empty(), "no boolean flags: {:?}", tool_args.flags);
assert_eq!(tool_args.named.get("lines"), Some(&Value::Int(5)), "should resolve alias to canonical name");
assert_eq!(tool_args.positional, vec![Value::String("/tmp/file.txt".to_string())]);
}
#[tokio::test]
async fn test_glued_short_flag_value_now_binds() {
let schema = ToolSchema::new("cut", "")
.param(ParamSchema::optional("fields", "string", Value::Null, "").with_aliases(["-f"]));
let args = vec![Arg::ShortFlag("f1".to_string())];
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert!(tool_args.flags.is_empty(), "no bogus bool flags: {:?}", tool_args.flags);
assert_eq!(tool_args.named.get("fields"), Some(&Value::String("1".to_string())));
}
#[tokio::test]
async fn test_repeatable_flag_now_accumulates() {
let schema = ToolSchema::new("sed", "")
.param(ParamSchema::optional("expression", "string", Value::Null, "")
.with_aliases(["-e"])
.with_repeatable(true));
let args = vec![
Arg::ShortFlag("e".to_string()),
Arg::Positional(Expr::Literal(Value::String("A".to_string()))),
Arg::ShortFlag("e".to_string()),
Arg::Positional(Expr::Literal(Value::String("B".to_string()))),
];
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert_eq!(
tool_args.named.get("expression"),
Some(&Value::Json(serde_json::json!(["A", "B"]))),
"both occurrences must survive, not just the last: {:?}",
tool_args.named
);
}
#[tokio::test]
async fn test_multi_consume_flag_now_accumulates() {
let schema = ToolSchema::new("jq", "")
.param(ParamSchema::optional("arg", "any", Value::Null, "").consumes(2));
let args = vec![
Arg::LongFlag("arg".to_string()),
Arg::Positional(Expr::Literal(Value::String("name".to_string()))),
Arg::Positional(Expr::Literal(Value::String("val".to_string()))),
];
let ctx = make_minimal_ctx();
let tool_args = build_tool_args(&args, &ctx, Some(&schema)).await.expect("build_tool_args");
assert_eq!(
tool_args.named.get("arg"),
Some(&Value::Json(serde_json::json!([["name", "val"]]))),
"got: {:?}",
tool_args.named
);
assert!(tool_args.positional.is_empty());
}
#[tokio::test]
async fn test_merge_stderr_redirect() {
let result = ExecResult::from_output(0, "stdout content", "stderr content");
let redirects = vec![Redirect {
kind: RedirectKind::MergeStderr,
target: Expr::Literal(Value::Null),
}];
let ctx = make_minimal_ctx();
let result = apply_redirects(result, &redirects, &ctx, &test_dispatcher()).await;
assert_eq!(&*result.text_out(), "stdout contentstderr content");
assert!(result.err.is_empty());
}
#[tokio::test]
async fn test_merge_stderr_with_empty_stderr() {
let result = ExecResult::from_output(0, "stdout only", "");
let redirects = vec![Redirect {
kind: RedirectKind::MergeStderr,
target: Expr::Literal(Value::Null),
}];
let ctx = make_minimal_ctx();
let result = apply_redirects(result, &redirects, &ctx, &test_dispatcher()).await;
assert_eq!(&*result.text_out(), "stdout only");
assert!(result.err.is_empty());
}
#[tokio::test]
async fn test_merge_stderr_order_matters() {
let result = ExecResult::from_output(0, "stdout\n", "stderr\n");
let redirects = vec![Redirect {
kind: RedirectKind::MergeStderr,
target: Expr::Literal(Value::Null),
}];
let ctx = make_minimal_ctx();
let result = apply_redirects(result, &redirects, &ctx, &test_dispatcher()).await;
assert_eq!(&*result.text_out(), "stdout\nstderr\n");
assert!(result.err.is_empty());
}
#[tokio::test]
async fn test_redirect_with_command_execution() {
let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
let cmd = Command {
name: "echo".to_string(),
args: vec![Arg::Positional(Expr::Literal(Value::String("hello".to_string())))],
redirects: vec![Redirect {
kind: RedirectKind::MergeStderr,
target: Expr::Literal(Value::Null),
}],
};
let result = runner.run(&stages([cmd]), &mut ctx, &dispatcher).await;
assert!(result.ok());
assert!(result.text_out().contains("hello"));
}
#[tokio::test]
async fn test_merge_stderr_in_pipeline() {
let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
let echo_cmd = Command {
name: "echo".to_string(),
args: vec![Arg::Positional(Expr::Literal(Value::String("output".to_string())))],
redirects: vec![Redirect {
kind: RedirectKind::MergeStderr,
target: Expr::Literal(Value::Null),
}],
};
let grep_cmd = Command {
name: "grep".to_string(),
args: vec![Arg::Positional(Expr::Literal(Value::String("output".to_string())))],
redirects: vec![],
};
let result = runner.run(&stages([echo_cmd, grep_cmd]), &mut ctx, &dispatcher).await;
assert!(result.ok(), "result failed: code={}, err={}", result.code, result.err);
assert!(result.text_out().contains("output"));
}
fn big_table_output(rows: usize) -> crate::interpreter::OutputData {
use crate::interpreter::OutputNode;
let headers = vec!["id".to_string(), "name".to_string()];
let nodes: Vec<OutputNode> = (0..rows)
.map(|i| OutputNode::new(i.to_string()).with_cells(vec![format!("row-{i}")]))
.collect();
crate::interpreter::OutputData::table(headers, nodes)
}
#[tokio::test]
async fn test_both_redirect_streams_structured_output_to_file() {
let output = big_table_output(50);
let expected_stdout = output.to_canonical_string();
let mut result = ExecResult::with_output(output);
result.err = "warning: heads up\n".to_string();
let redirects = vec![Redirect {
kind: RedirectKind::Both,
target: Expr::Literal(Value::String("/out.txt".to_string())),
}];
let ctx = make_minimal_ctx();
let result = apply_redirects(result, &redirects, &ctx, &test_dispatcher()).await;
assert!(result.ok());
assert_eq!(&*result.text_out(), "");
assert!(result.err.is_empty());
assert!(!result.has_output());
let written = ctx.backend.read(Path::new("/out.txt"), None).await.expect("file written");
let written = String::from_utf8(written).expect("valid utf8");
assert_eq!(written, format!("{expected_stdout}warning: heads up\n"));
}
#[tokio::test]
async fn test_both_redirect_streams_large_structured_output_intact() {
let rows = 5_000;
let output = big_table_output(rows);
let expected_stdout = output.to_canonical_string();
let result = ExecResult::with_output(output);
let redirects = vec![Redirect {
kind: RedirectKind::Both,
target: Expr::Literal(Value::String("/big.txt".to_string())),
}];
let ctx = make_minimal_ctx();
let result = apply_redirects(result, &redirects, &ctx, &test_dispatcher()).await;
assert!(result.ok());
let written = ctx.backend.read(Path::new("/big.txt"), None).await.expect("file written");
let written = String::from_utf8(written).expect("valid utf8");
assert_eq!(written, expected_stdout);
assert!(written.contains("row-0"));
assert!(written.contains(&format!("row-{}", rows - 1)));
}
#[tokio::test]
async fn test_both_redirect_still_writes_binary_stdout_raw() {
let result = ExecResult::success_text_or_bytes(vec![0xff, 0x00, 0xfe, b'x']);
let redirects = vec![Redirect {
kind: RedirectKind::Both,
target: Expr::Literal(Value::String("/bin.out".to_string())),
}];
let ctx = make_minimal_ctx();
let result = apply_redirects(result, &redirects, &ctx, &test_dispatcher()).await;
assert!(result.ok());
let written = ctx.backend.read(Path::new("/bin.out"), None).await.expect("file written");
assert_eq!(written, vec![0xff, 0x00, 0xfe, b'x']);
}
}