1use std::sync::Arc;
9
10use std::collections::HashMap;
11
12use crate::arithmetic;
13use crate::ast::{Arg, Command, Expr, Redirect, RedirectKind, Value};
14use crate::dispatch::{CommandDispatcher, PipelinePosition};
15use crate::interpreter::{ExecResult, PathError};
16use crate::tools::{ExecContext, ToolArgs, ToolRegistry, ToolSchema};
17use tokio::io::AsyncWriteExt;
18
19use super::pipe_stream::pipe_stream_default;
20use super::scatter::{
21 parse_gather_options, parse_scatter_options, ScatterGatherRunner,
22};
23
24pub(crate) async fn apply_redirects(
30 mut result: ExecResult,
31 redirects: &[Redirect],
32 ctx: &ExecContext,
33 dispatcher: &dyn CommandDispatcher,
34) -> ExecResult {
35 for redir in redirects {
40 match redir.kind {
41 RedirectKind::MergeStderr => {
42 result.materialize();
45 if !result.err.is_empty() {
46 let err = std::mem::take(&mut result.err);
47 result.push_out(&err);
48 }
49 }
50 RedirectKind::MergeStdout => {
51 if result.is_bytes() {
55 return ExecResult::failure(
56 1,
57 "redirect: cannot merge binary stdout into stderr (1>&2) — \
58 redirect it to a file or pipe through base64/xxd",
59 );
60 }
61 result.materialize();
62 if !result.text_out().is_empty() {
63 let out = result.text_out().into_owned();
64 result.err.push_str(&out);
65 }
66 result.clear_stdout();
73 }
74 RedirectKind::StdoutOverwrite => {
75 let path = match eval_redirect_target(&redir.target, ctx, dispatcher).await {
76 Ok(p) => p,
77 Err(e) => return ExecResult::failure(1, format!("redirect: {e}")),
78 };
79 if let Some(bytes) = result.out_bytes() {
81 if let Err(e) = redirect_write(ctx, &path, bytes).await {
82 return ExecResult::failure(1, format!("redirect: {e}"));
83 }
84 } else if let Some(output) = result.take_output_for_stream() {
85 let mut buf = Vec::new();
87 if let Err(e) = output.write_canonical(&mut buf, None) {
88 return ExecResult::failure(1, format!("redirect: {e}"));
89 }
90 if let Err(e) = redirect_write(ctx, &path, &buf).await {
91 return ExecResult::failure(1, format!("redirect: {e}"));
92 }
93 } else if let Err(e) = redirect_write(ctx, &path, result.text_out().as_bytes()).await {
94 return ExecResult::failure(1, format!("redirect: {e}"));
95 }
96 result.clear_stdout();
98 }
99 RedirectKind::StdoutAppend => {
100 let path = match eval_redirect_target(&redir.target, ctx, dispatcher).await {
101 Ok(p) => p,
102 Err(e) => return ExecResult::failure(1, format!("redirect: {e}")),
103 };
104 if let Some(bytes) = result.out_bytes() {
106 if let Err(e) = redirect_append(ctx, &path, bytes).await {
107 return ExecResult::failure(1, format!("redirect: {e}"));
108 }
109 } else if let Some(output) = result.take_output_for_stream() {
110 let mut buf = Vec::new();
112 if let Err(e) = output.write_canonical(&mut buf, None) {
113 return ExecResult::failure(1, format!("redirect: {e}"));
114 }
115 if let Err(e) = redirect_append(ctx, &path, &buf).await {
116 return ExecResult::failure(1, format!("redirect: {e}"));
117 }
118 } else if let Err(e) = redirect_append(ctx, &path, result.text_out().as_bytes()).await {
119 return ExecResult::failure(1, format!("redirect: {e}"));
120 }
121 result.clear_stdout();
123 }
124 RedirectKind::Stderr => {
125 let path = match eval_redirect_target(&redir.target, ctx, dispatcher).await {
126 Ok(p) => p,
127 Err(e) => return ExecResult::failure(1, format!("redirect: {e}")),
128 };
129 if let Err(e) = redirect_write(ctx, &path, result.err.as_bytes()).await {
130 return ExecResult::failure(1, format!("redirect: {e}"));
131 }
132 result.err.clear();
133 }
134 RedirectKind::Both => {
135 let path = match eval_redirect_target(&redir.target, ctx, dispatcher).await {
136 Ok(p) => p,
137 Err(e) => return ExecResult::failure(1, format!("redirect: {e}")),
138 };
139 let mut combined: Vec<u8> = if let Some(b) = result.out_bytes() {
147 b.to_vec()
148 } else if let Some(output) = result.take_output_for_stream() {
149 let mut buf = Vec::new();
150 if let Err(e) = output.write_canonical(&mut buf, None) {
151 return ExecResult::failure(1, format!("redirect: {e}"));
152 }
153 buf
154 } else {
155 result.text_out().into_owned().into_bytes()
156 };
157 combined.extend_from_slice(result.err.as_bytes());
158 if let Err(e) = redirect_write(ctx, &path, &combined).await {
159 return ExecResult::failure(1, format!("redirect: {e}"));
160 }
161 result.clear_stdout();
163 result.err.clear();
164 }
165 RedirectKind::Stdin | RedirectKind::HereDoc | RedirectKind::HereString => {}
167 }
168 }
169 result.materialize();
174 result
175}
176
177async fn eval_redirect_target(
189 expr: &Expr,
190 ctx: &ExecContext,
191 dispatcher: &dyn CommandDispatcher,
192) -> Result<String, String> {
193 let value = dispatcher
194 .eval_expr(expr, ctx)
195 .await
196 .map_err(|e| e.to_string())?;
197 if let Some(msg) = crate::interpreter::structured_boundary_error("a redirect target", &value) {
200 return Err(msg);
201 }
202 crate::interpreter::value_to_text_sink_named(&value, "a redirect target").map_err(|e| e.to_string())
205}
206
207async fn redirect_write(ctx: &ExecContext, path: &str, data: &[u8]) -> Result<(), String> {
214 use crate::backend::WriteMode;
215 let resolved = ctx.resolve_path(path);
216 ctx.backend.write(&resolved, data, WriteMode::Overwrite).await.map_err(|e| e.to_string())
217}
218
219async fn redirect_append(ctx: &ExecContext, path: &str, data: &[u8]) -> Result<(), String> {
223 let resolved = ctx.resolve_path(path);
224 ctx.backend.append(&resolved, data).await.map_err(|e| e.to_string())
225}
226
227async fn setup_stdin_redirects(
235 cmd: &Command,
236 ctx: &mut ExecContext,
237 dispatcher: &dyn CommandDispatcher,
238) -> Result<(), String> {
239 use std::path::Path;
240 for redir in &cmd.redirects {
241 match &redir.kind {
242 RedirectKind::Stdin => {
243 let path = eval_redirect_target(&redir.target, ctx, dispatcher).await?;
244 let resolved = ctx.resolve_path(&path);
245 let data = ctx
246 .backend
247 .read(Path::new(&resolved), None)
248 .await
249 .map_err(|e| format!("redirect: {path}: {e}"))?;
250 let content = String::from_utf8(data)
251 .map_err(|_| format!("redirect: {path}: invalid UTF-8"))?;
252 ctx.set_stdin(content);
253 }
254 RedirectKind::HereDoc => {
255 match &redir.target {
256 Expr::Literal(Value::String(content)) => {
257 ctx.set_stdin(content.clone());
258 }
259 expr => {
262 let body = eval_redirect_target(expr, ctx, dispatcher).await?;
263 ctx.set_stdin(body);
264 }
265 }
266 }
267 RedirectKind::HereString => {
268 let mut s = eval_redirect_target(&redir.target, ctx, dispatcher).await?;
271 s.push('\n');
272 ctx.set_stdin(s);
273 }
274 _ => {}
275 }
276 }
277 Ok(())
278}
279
280#[derive(Clone)]
282pub struct PipelineRunner {
283 tools: Arc<ToolRegistry>,
284}
285
286impl PipelineRunner {
287 pub fn new(tools: Arc<ToolRegistry>) -> Self {
289 Self { tools }
290 }
291
292 pub async fn run(
302 &self,
303 commands: &[Command],
304 ctx: &mut ExecContext,
305 dispatcher: &dyn CommandDispatcher,
306 ) -> ExecResult {
307 if commands.is_empty() {
308 return ExecResult::success("");
309 }
310
311 if let Some((scatter_idx, gather_idx)) = find_scatter_gather(commands) {
313 return self.run_scatter_gather(commands, scatter_idx, gather_idx, ctx, dispatcher).await;
314 }
315
316 self.run_sequential(commands, ctx, dispatcher).await
317 }
318
319 pub async fn run_sequential(
324 &self,
325 commands: &[Command],
326 ctx: &mut ExecContext,
327 dispatcher: &dyn CommandDispatcher,
328 ) -> ExecResult {
329 if commands.is_empty() {
330 return ExecResult::success("");
331 }
332
333 if commands.len() == 1 {
334 return self.run_single(&commands[0], ctx, None, dispatcher).await;
336 }
337
338 self.run_pipeline(commands, ctx, dispatcher).await
340 }
341
342 async fn run_scatter_gather(
344 &self,
345 commands: &[Command],
346 scatter_idx: usize,
347 gather_idx: usize,
348 ctx: &mut ExecContext,
349 dispatcher: &dyn CommandDispatcher,
350 ) -> ExecResult {
351 let pre_scatter = &commands[..scatter_idx];
353 let scatter_cmd = &commands[scatter_idx];
354 let parallel = &commands[scatter_idx + 1..gather_idx];
355 let gather_cmd = &commands[gather_idx];
356 let post_gather = &commands[gather_idx + 1..];
357
358 let scatter_schema = self.tools.get("scatter").map(|t| t.schema());
365 let gather_schema = self.tools.get("gather").map(|t| t.schema());
366 let scatter_args = match build_tool_args(&scatter_cmd.args, ctx, scatter_schema.as_ref()) {
367 Ok(args) => args,
368 Err(e) => return ExecResult::failure(1, format!("scatter: {e}")),
369 };
370 let gather_args = match build_tool_args(&gather_cmd.args, ctx, gather_schema.as_ref()) {
371 Ok(args) => args,
372 Err(e) => return ExecResult::failure(1, format!("gather: {e}")),
373 };
374 let scatter_opts = match parse_scatter_options(&scatter_args) {
375 Ok(opts) => opts,
376 Err(e) => return ExecResult::failure(2, format!("scatter: {e}")),
377 };
378 let gather_opts = match parse_gather_options(&gather_args) {
379 Ok(opts) => opts,
380 Err(e) => return ExecResult::failure(2, format!("gather: {e}")),
381 };
382
383 let sequential_dispatcher: Arc<dyn CommandDispatcher> = dispatcher.fork_attached().await;
388
389 let runner = ScatterGatherRunner::new(self.tools.clone(), sequential_dispatcher);
390 runner
391 .run(
392 pre_scatter,
393 scatter_opts,
394 parallel,
395 gather_opts,
396 &gather_cmd.redirects,
397 post_gather,
398 ctx,
399 )
400 .await
401 }
402
403 async fn run_single(
408 &self,
409 cmd: &Command,
410 ctx: &mut ExecContext,
411 stdin: Option<String>,
412 dispatcher: &dyn CommandDispatcher,
413 ) -> ExecResult {
414 if let Err(e) = setup_stdin_redirects(cmd, ctx, dispatcher).await {
416 return ExecResult::failure(1, e);
417 }
418
419 if let Some(input) = stdin {
421 ctx.set_stdin(input);
422 }
423
424 ctx.pipeline_position = PipelinePosition::Only;
426
427 let result = match dispatcher.dispatch(cmd, ctx).await {
429 Ok(result) => result,
430 Err(e) => ExecResult::failure(1, e.to_string()),
431 };
432
433 apply_redirects(result, &cmd.redirects, ctx, dispatcher).await
435 }
436
437 async fn run_pipeline(
447 &self,
448 commands: &[Command],
449 ctx: &mut ExecContext,
450 dispatcher: &dyn CommandDispatcher,
451 ) -> ExecResult {
452 let stage_count = commands.len();
453 let last_idx = stage_count - 1;
454
455 let mut pipe_writers: Vec<Option<super::pipe_stream::PipeWriter>> = Vec::new();
457 let mut pipe_readers: Vec<Option<super::pipe_stream::PipeReader>> = Vec::new();
458
459 for _ in 0..last_idx {
460 let (writer, reader) = pipe_stream_default();
461 pipe_writers.push(Some(writer));
462 pipe_readers.push(Some(reader));
463 }
464
465 let mut data_senders: Vec<Option<tokio::sync::oneshot::Sender<Option<Value>>>> = Vec::new();
467 let mut data_receivers: Vec<Option<tokio::sync::oneshot::Receiver<Option<Value>>>> = Vec::new();
468
469 for _ in 0..last_idx {
470 let (tx, rx) = tokio::sync::oneshot::channel();
471 data_senders.push(Some(tx));
472 data_receivers.push(Some(rx));
473 }
474
475 let mut handles: Vec<tokio::task::JoinHandle<(ExecResult, ExecContext)>> = Vec::with_capacity(stage_count);
476
477 for (i, cmd) in commands.iter().enumerate() {
478 let mut stage_ctx = ctx.child_for_pipeline();
479 let cmd = cmd.clone();
480
481 let task_dispatcher: Arc<dyn CommandDispatcher> = dispatcher.fork_attached().await;
486
487 let stdin_setup = setup_stdin_redirects(&cmd, &mut stage_ctx, dispatcher).await;
491
492 if i == 0 {
494 if stage_ctx.stdin.is_none() {
497 stage_ctx.stdin = ctx.stdin.take();
498 }
499 if stage_ctx.stdin_data.is_none() {
500 stage_ctx.stdin_data = ctx.stdin_data.take();
501 }
502 if stage_ctx.stdin.is_none() && stage_ctx.pipe_stdin.is_none() {
506 stage_ctx.pipe_stdin = ctx.pipe_stdin.take();
507 }
508 } else {
509 stage_ctx.pipe_stdin = pipe_readers[i - 1].take();
511 }
513
514 if i < last_idx {
516 stage_ctx.pipe_stdout = pipe_writers[i].take();
517 }
518
519 stage_ctx.pipeline_position = match i {
521 0 => PipelinePosition::First,
522 n if n == last_idx => PipelinePosition::Last,
523 _ => PipelinePosition::Middle,
524 };
525
526 let data_sender = if i < last_idx { data_senders[i].take() } else { None };
527 let data_receiver = if i > 0 { data_receivers[i - 1].take() } else { None };
528
529 let handle: tokio::task::JoinHandle<(ExecResult, ExecContext)> =
532 tokio::spawn(crate::telemetry::bind_current_context(async move {
533 if let Err(e) = stdin_setup {
535 return (ExecResult::failure(1, e), stage_ctx);
536 }
537
538 stage_ctx.stdin_data_rx = data_receiver;
546
547 let mut result = match task_dispatcher.dispatch(&cmd, &mut stage_ctx).await {
549 Ok(result) => result,
550 Err(e) => ExecResult::failure(1, e.to_string()),
551 };
552
553 result = apply_redirects(result, &cmd.redirects, &stage_ctx, &*task_dispatcher).await;
558
559 if !result.err.is_empty() {
565 if let Some(ref stderr) = stage_ctx.stderr {
566 stderr.write_str(&result.err);
567 result.err.clear();
568 }
569 }
570
571 if let Some(tx) = data_sender {
577 let _ = tx.send(result.data.clone());
578 }
579
580 if let Some(mut pipe_out) = stage_ctx.pipe_stdout.take() {
583 let bytes: Vec<u8> = if let Some(b) = result.out_bytes() {
591 b.to_vec()
592 } else if let Some(output) = result.take_output_for_stream() {
593 let mut buf = Vec::new();
594 if output.write_canonical(&mut buf, None).is_err() {
600 buf = output.to_canonical_string().into_bytes();
601 }
602 buf
603 } else {
604 result.text_out().into_owned().into_bytes()
605 };
606 if !bytes.is_empty() {
607 let _ = pipe_out.write_all(&bytes).await;
609 let _ = pipe_out.shutdown().await;
610 }
611 }
613
614 (result, stage_ctx)
615 }));
616
617 handles.push(handle);
618 }
619
620 let mut last_result = ExecResult::success("");
625 let mut panics: Vec<String> = Vec::new();
626 let mut gated: Option<ExecResult> = None;
634 for (i, handle) in handles.into_iter().enumerate() {
635 match handle.await {
636 Ok((result, stage_ctx)) => {
637 if result.latch.is_some() {
638 gated.get_or_insert_with(|| result.clone());
639 }
640 if i == last_idx {
641 last_result = result;
642 ctx.scope = stage_ctx.scope;
644 ctx.cwd = stage_ctx.cwd;
645 ctx.prev_cwd = stage_ctx.prev_cwd;
646 ctx.aliases = stage_ctx.aliases;
647 }
648 }
649 Err(e) => {
650 panics.push(format!("stage {}: {}", i, e));
651 }
652 }
653 }
654
655 if let Some(gate) = gated {
659 last_result = gate;
660 }
661
662 if !panics.is_empty() {
668 last_result = ExecResult::failure(
669 1,
670 format!("pipeline stage(s) panicked: {}", panics.join("; ")),
671 );
672 }
673
674 last_result
675 }
676}
677
678pub fn select_leaf<'a>(schema: &'a ToolSchema, args: &[Arg]) -> anyhow::Result<&'a ToolSchema> {
718 let root_lookup = schema_param_lookup(schema);
721 let is_root_value_flag = |name: &str| -> bool {
722 root_lookup.get(name).is_some_and(|(_, typ, ..)| !is_bool_type(typ))
723 };
724
725 let mut node = schema;
726 let mut skip_next_positional = false;
727 for arg in args {
728 match arg {
729 Arg::DoubleDash => break,
731 Arg::LongFlag(name) if is_root_value_flag(name) => skip_next_positional = true,
734 Arg::ShortFlag(name) if is_root_value_flag(name) => skip_next_positional = true,
735 Arg::Positional(expr) => {
736 if skip_next_positional {
737 skip_next_positional = false;
738 continue; }
740 if node.subcommands.is_empty() {
741 break; }
743 match classify_subcommand_positional(expr) {
744 SubcommandWord::Word(word) => {
745 match node.subcommands.iter().find(|c| c.matches_command(word)) {
746 Some(child) => node = child, None => break, }
749 }
750 SubcommandWord::OtherLiteral => break,
754 SubcommandWord::Computed(kind) => anyhow::bail!(
755 "{}: a subcommand name is required here, but got {kind}. \
756 Subcommands must be literal words — spell it out \
757 (e.g. `{} <subcommand> …`) or use the `--flag=value` form.",
758 node.name,
759 schema.name
760 ),
761 }
762 }
763 _ => {}
765 }
766 }
767 Ok(node)
768}
769
770enum SubcommandWord<'a> {
772 Word(&'a str),
774 OtherLiteral,
776 Computed(&'static str),
778}
779
780fn classify_subcommand_positional(expr: &Expr) -> SubcommandWord<'_> {
781 match expr {
782 Expr::Literal(Value::String(s)) => SubcommandWord::Word(s),
783 Expr::Literal(_) => SubcommandWord::OtherLiteral,
784 Expr::CommandSubst(_) | Expr::Command(_) => SubcommandWord::Computed("a command substitution `$(…)`"),
785 Expr::VarRef(_)
786 | Expr::VarWithDefault { .. }
787 | Expr::VarLength(_)
788 | Expr::Positional(_)
789 | Expr::AllArgs
790 | Expr::ArgCount
791 | Expr::CurrentPid
792 | Expr::LastExitCode => SubcommandWord::Computed("a variable reference"),
793 Expr::Interpolated(_) | Expr::HereDocBody { .. } => SubcommandWord::Computed("an interpolated string"),
794 Expr::GlobPattern(_) => SubcommandWord::Computed("a glob pattern"),
795 Expr::Arithmetic(_) => SubcommandWord::Computed("an arithmetic expansion"),
796 _ => SubcommandWord::Computed("a value computed at runtime"),
797 }
798}
799
800pub fn schema_param_lookup(schema: &ToolSchema) -> HashMap<String, (&str, &str, usize, bool)> {
801 let mut map = HashMap::new();
802 for p in schema.params.iter().filter(|p| !p.positional) {
803 map.insert(p.name.clone(), (p.name.as_str(), p.param_type.as_str(), p.consumes, p.repeatable));
804 for alias in &p.aliases {
805 let stripped = alias.trim_start_matches('-');
806 map.insert(stripped.to_string(), (p.name.as_str(), p.param_type.as_str(), p.consumes, p.repeatable));
807 }
808 }
809 map
810}
811
812pub fn is_bool_type(param_type: &str) -> bool {
814 matches!(param_type.to_lowercase().as_str(), "bool" | "boolean")
815}
816
817pub fn build_tool_args(
830 args: &[Arg],
831 ctx: &ExecContext,
832 schema: Option<&ToolSchema>,
833) -> Result<ToolArgs, String> {
834 let mut tool_args = ToolArgs::new();
835 let param_lookup = schema.map(schema_param_lookup).unwrap_or_default();
836 let accepts_word_assign = schema
837 .map(|s| crate::tools::accepts_word_assign(s.name.as_str()))
838 .unwrap_or(false);
839
840 let mut consumed_positionals: std::collections::HashSet<usize> = std::collections::HashSet::new();
842 let mut past_double_dash = false;
843
844 let mut positional_indices: Vec<(usize, &Expr)> = Vec::new();
846 for (i, arg) in args.iter().enumerate() {
847 if let Arg::Positional(expr) = arg {
848 positional_indices.push((i, expr));
849 }
850 }
851
852 let mut i = 0;
854 while i < args.len() {
855 let arg = &args[i];
856
857 match arg {
858 Arg::DoubleDash => {
859 past_double_dash = true;
860 }
861 Arg::Positional(expr) => {
862 if !consumed_positionals.contains(&i)
864 && let Some(value) = eval_simple_expr(expr, ctx)?
865 {
866 tool_args.positional.push(value);
867 }
868 }
869 Arg::Named { key, value } => {
870 if let Some(val) = eval_simple_expr(value, ctx)? {
871 tool_args.named.insert(key.clone(), val);
872 }
873 }
874 Arg::WordAssign { key, value } => {
875 if let Some(val) = eval_simple_expr(value, ctx)? {
876 if accepts_word_assign {
877 tool_args.named.insert(key.clone(), val);
878 } else {
879 let val_str = crate::interpreter::value_to_text_sink_named(
882 &val,
883 "a key=value argument",
884 )
885 .map_err(|e| e.to_string())?;
886 tool_args.positional.push(Value::String(format!("{key}={val_str}")));
887 }
888 }
889 }
890 Arg::ShortFlag(name) => {
891 if past_double_dash {
892 tool_args.positional.push(Value::String(format!("-{name}")));
893 } else if name.len() == 1 {
894 let flag_name = name.as_str();
897 let lookup = param_lookup.get(flag_name);
898 let is_bool = lookup
899 .map(|(_, typ, ..)| is_bool_type(typ))
900 .unwrap_or(true);
901
902 if is_bool {
903 tool_args.flags.insert(flag_name.to_string());
904 } else {
905 let canonical = lookup.map(|(n, ..)| *n).unwrap_or(flag_name);
907 let next_positional = positional_indices
908 .iter()
909 .find(|(idx, _)| *idx > i && !consumed_positionals.contains(idx));
910
911 if let Some((pos_idx, expr)) = next_positional {
912 if let Some(value) = eval_simple_expr(expr, ctx)? {
913 tool_args.named.insert(canonical.to_string(), value);
914 consumed_positionals.insert(*pos_idx);
915 } else {
916 tool_args.flags.insert(flag_name.to_string());
917 }
918 } else {
919 tool_args.flags.insert(flag_name.to_string());
920 }
921 }
922 } else if let Some(&(canonical, typ, ..)) = param_lookup.get(name.as_str()) {
923 if is_bool_type(typ) {
925 tool_args.flags.insert(canonical.to_string());
926 } else {
927 let next_positional = positional_indices
928 .iter()
929 .find(|(idx, _)| *idx > i && !consumed_positionals.contains(idx));
930 if let Some((pos_idx, expr)) = next_positional {
931 if let Some(value) = eval_simple_expr(expr, ctx)? {
932 tool_args.named.insert(canonical.to_string(), value);
933 consumed_positionals.insert(*pos_idx);
934 } else {
935 tool_args.flags.insert(name.clone());
936 }
937 } else {
938 tool_args.flags.insert(name.clone());
939 }
940 }
941 } else {
942 for c in name.chars() {
944 tool_args.flags.insert(c.to_string());
945 }
946 }
947 }
948 Arg::LongFlag(name) => {
949 if past_double_dash {
950 tool_args.positional.push(Value::String(format!("--{name}")));
951 } else {
952 let lookup = param_lookup.get(name.as_str());
954 let is_bool = lookup
955 .map(|(_, typ, ..)| is_bool_type(typ))
956 .unwrap_or(true); if is_bool {
959 tool_args.flags.insert(name.clone());
960 } else {
961 let canonical = lookup.map(|(n, ..)| *n).unwrap_or(name.as_str());
970 let next_positional = positional_indices
971 .iter()
972 .find(|(idx, _)| *idx > i && !consumed_positionals.contains(idx));
973
974 if let Some((pos_idx, expr)) = next_positional {
975 if let Some(value) = eval_simple_expr(expr, ctx)? {
976 tool_args.named.insert(canonical.to_string(), value);
977 consumed_positionals.insert(*pos_idx);
978 } else {
979 tool_args.flags.insert(name.clone());
980 }
981 } else {
982 tool_args.flags.insert(name.clone());
983 }
984 }
985 }
986 }
987 }
988 i += 1;
989 }
990
991 if let Some(schema) = schema.filter(|s| s.map_positionals) {
996 let pre_dash_count = if past_double_dash {
998 let dash_pos = args.iter().position(|a| matches!(a, Arg::DoubleDash)).unwrap_or(args.len());
1000 positional_indices.iter()
1002 .filter(|(idx, _)| *idx < dash_pos && !consumed_positionals.contains(idx))
1003 .count()
1004 } else {
1005 tool_args.positional.len()
1006 };
1007
1008 let mut remaining = Vec::new();
1009 let mut positional_iter = tool_args.positional.drain(..).enumerate();
1010
1011 for param in &schema.params {
1012 if tool_args.named.contains_key(¶m.name) || tool_args.flags.contains(¶m.name) {
1013 continue; }
1015 if is_bool_type(¶m.param_type) {
1016 continue; }
1018 loop {
1020 match positional_iter.next() {
1021 Some((idx, val)) if idx < pre_dash_count => {
1022 tool_args.named.insert(param.name.clone(), val);
1023 break;
1024 }
1025 Some((_, val)) => {
1026 remaining.push(val); }
1028 None => break,
1029 }
1030 }
1031 }
1032
1033 remaining.extend(positional_iter.map(|(_, v)| v));
1035 tool_args.positional = remaining;
1036 }
1037
1038 Ok(tool_args)
1039}
1040
1041pub(crate) fn eval_simple_expr(expr: &Expr, ctx: &ExecContext) -> Result<Option<Value>, String> {
1054 match expr {
1055 Expr::Literal(value) => Ok(Some(eval_literal(value, ctx))),
1056 Expr::VarRef(path) => match ctx.scope.resolve_path(path) {
1057 Ok(v) => Ok(Some(v)),
1058 Err(PathError::UndefinedRoot(_)) if path.segments.len() <= 1 => Ok(None),
1065 Err(PathError::UndefinedRoot(_)) => Err(format!(
1066 "{}: undefined variable",
1067 crate::interpreter::format_path(path)
1068 )),
1069 Err(PathError::Absence(msg)) | Err(PathError::Shape(msg)) => Err(msg),
1072 },
1073 Expr::Interpolated(parts) => Ok(Some(Value::String(eval_string_parts_sync(parts, ctx)?))),
1074 Expr::VarLength(path) => {
1079 crate::interpreter::resolve_length(&ctx.scope, path).map(|n| Some(Value::Int(n)))
1080 }
1081 Expr::VarWithDefault { path, default } => {
1082 match crate::interpreter::resolve_default(&ctx.scope, path)? {
1083 Some(value) => Ok(Some(value)),
1084 None => Ok(Some(Value::String(eval_string_parts_sync(default, ctx)?))),
1085 }
1086 }
1087 Expr::GlobPattern(s) => Ok(Some(Value::String(s.clone()))),
1088 Expr::HereDocBody { parts, strip_tabs } => {
1089 let mut asm = crate::interpreter::HeredocAssembler::new(*strip_tabs);
1093 for sp in parts {
1094 match &sp.part {
1095 crate::ast::StringPart::Literal(s) => asm.push_literal(s),
1096 other => {
1097 let s = eval_string_parts_sync(std::slice::from_ref(other), ctx)?;
1098 asm.push_interpolated(&s);
1099 }
1100 }
1101 }
1102 Ok(Some(Value::String(asm.into_string())))
1103 }
1104 Expr::CommandSubst(_) | Expr::Command(_) => Err(
1111 "command substitution `$(...)` is not supported in a scatter/gather flag value here; \
1112 assign it to a variable first (e.g. `n=$(...); scatter --limit $n`)"
1113 .to_string(),
1114 ),
1115 _ => Ok(None), }
1117}
1118
1119fn eval_literal(value: &Value, _ctx: &ExecContext) -> Value {
1121 value.clone()
1122}
1123
1124fn eval_string_parts_sync(parts: &[crate::ast::StringPart], ctx: &ExecContext) -> Result<String, String> {
1132 let mut result = String::new();
1133 for part in parts {
1134 match part {
1135 crate::ast::StringPart::Literal(s) => result.push_str(s),
1136 crate::ast::StringPart::Var(path) => match ctx.scope.resolve_path(path) {
1137 Ok(value) => result.push_str(
1141 &crate::interpreter::value_to_text_sink(&value).map_err(|e| e.to_string())?,
1142 ),
1143 Err(PathError::UndefinedRoot(_)) => {}
1151 Err(PathError::Absence(msg)) | Err(PathError::Shape(msg)) => return Err(msg),
1152 },
1153 crate::ast::StringPart::VarWithDefault { path, default } => {
1154 match crate::interpreter::resolve_default(&ctx.scope, path)? {
1155 Some(value) => result.push_str(
1156 &crate::interpreter::value_to_text_sink(&value).map_err(|e| e.to_string())?,
1157 ),
1158 None => result.push_str(&eval_string_parts_sync(default, ctx)?),
1159 }
1160 }
1161 crate::ast::StringPart::VarLength(path) => {
1162 let len = crate::interpreter::resolve_length(&ctx.scope, path)?;
1167 result.push_str(&len.to_string());
1168 }
1169 crate::ast::StringPart::Positional(n) => {
1170 if let Some(s) = ctx.scope.get_positional(*n) {
1171 result.push_str(s);
1172 }
1173 }
1174 crate::ast::StringPart::AllArgs => {
1175 result.push_str(&ctx.scope.all_args().join(" "));
1176 }
1177 crate::ast::StringPart::ArgCount => {
1178 result.push_str(&ctx.scope.arg_count().to_string());
1179 }
1180 crate::ast::StringPart::Arithmetic(expr) => {
1181 if let Ok(value) = arithmetic::eval_arithmetic(expr, &ctx.scope) {
1185 result.push_str(&value.to_string());
1186 }
1187 }
1188 crate::ast::StringPart::CommandSubst(_) => {
1189 return Err(
1195 "command substitution `$(...)` is not supported inside a scatter/gather \
1196 flag's interpolated value here; assign it to a variable first"
1197 .to_string(),
1198 );
1199 }
1200 crate::ast::StringPart::LastExitCode => {
1201 result.push_str(&ctx.scope.last_result().code.to_string());
1202 }
1203 crate::ast::StringPart::CurrentPid => {
1204 result.push_str(&ctx.scope.pid().to_string());
1205 }
1206 }
1207 }
1208 Ok(result)
1209}
1210
1211fn find_scatter_gather(commands: &[Command]) -> Option<(usize, usize)> {
1216 let scatter_idx = commands.iter().position(|c| c.name == "scatter")?;
1217 let gather_idx = commands.iter().position(|c| c.name == "gather")?;
1218
1219 if gather_idx > scatter_idx {
1221 Some((scatter_idx, gather_idx))
1222 } else {
1223 None
1224 }
1225}
1226
1227#[cfg(test)]
1228mod select_leaf_tests {
1229 use super::*;
1230 use crate::tools::ParamSchema;
1231
1232 fn kj_schema() -> ToolSchema {
1237 ToolSchema::new("kj", "kaijutsu")
1238 .param(ParamSchema::new("confirm", "string"))
1239 .param(ParamSchema::new("verbose", "bool"))
1240 .subcommand(
1241 ToolSchema::new("context", "context ops")
1242 .with_command_aliases(["ctx"])
1243 .subcommand(ToolSchema::new("list", "list").with_command_aliases(["ls"]))
1244 .subcommand(
1245 ToolSchema::new("create", "create").param(
1246 ParamSchema::new("type", "string").with_aliases(["t"]),
1247 ),
1248 ),
1249 )
1250 }
1251
1252 fn word(s: &str) -> Arg {
1253 Arg::Positional(Expr::Literal(Value::String(s.to_string())))
1254 }
1255
1256 #[test]
1257 fn flat_tool_returns_root() {
1258 let schema = ToolSchema::new("cat", "concat")
1259 .param(ParamSchema::required("path", "string", "f").positional());
1260 let leaf = select_leaf(&schema, &[word("foo.txt")]).expect("flat ok");
1261 assert_eq!(leaf.name, "cat");
1262 }
1263
1264 #[test]
1265 fn single_hop() {
1266 let schema = kj_schema();
1267 let leaf = select_leaf(&schema, &[word("context")]).expect("ok");
1268 assert_eq!(leaf.name, "context");
1269 }
1270
1271 #[test]
1272 fn two_hops() {
1273 let schema = kj_schema();
1274 let leaf = select_leaf(&schema, &[word("context"), word("create")]).expect("ok");
1275 assert_eq!(leaf.name, "create");
1276 assert!(leaf.params.iter().any(|p| p.name == "type"), "leaf has --type");
1277 }
1278
1279 #[test]
1280 fn alias_hops_route() {
1281 let schema = kj_schema();
1282 let leaf = select_leaf(&schema, &[word("ctx"), word("ls")]).expect("ok");
1284 assert_eq!(leaf.name, "list");
1285 }
1286
1287 #[test]
1288 fn unknown_subcommand_stops_at_current_node() {
1289 let schema = kj_schema();
1290 let leaf = select_leaf(&schema, &[word("context"), word("nonesuch")]).expect("ok");
1293 assert_eq!(leaf.name, "context");
1294 }
1295
1296 #[test]
1297 fn root_bool_flag_before_path_does_not_disrupt_routing() {
1298 let schema = kj_schema();
1299 let args = vec![Arg::LongFlag("verbose".into()), word("context"), word("create")];
1302 let leaf = select_leaf(&schema, &args).expect("ok");
1303 assert_eq!(leaf.name, "create");
1304 }
1305
1306 #[test]
1307 fn root_value_flag_space_form_before_path_skips_its_value() {
1308 let schema = kj_schema();
1309 let args = vec![
1312 Arg::LongFlag("confirm".into()),
1313 word("nonce"),
1314 word("context"),
1315 word("create"),
1316 ];
1317 let leaf = select_leaf(&schema, &args).expect("ok");
1318 assert_eq!(leaf.name, "create");
1319 }
1320
1321 #[test]
1322 fn leaf_value_flag_after_path_routes_to_leaf() {
1323 let schema = kj_schema();
1324 let args = vec![
1327 word("context"),
1328 word("create"),
1329 Arg::LongFlag("type".into()),
1330 word("x"),
1331 ];
1332 let leaf = select_leaf(&schema, &args).expect("ok");
1333 assert_eq!(leaf.name, "create");
1334 assert!(leaf.params.iter().any(|p| p.name == "type"));
1335 }
1336
1337 #[test]
1338 fn double_dash_stops_routing() {
1339 let schema = kj_schema();
1340 let leaf = select_leaf(&schema, &[Arg::DoubleDash, word("context")]).expect("ok");
1342 assert_eq!(leaf.name, "kj");
1343 }
1344
1345 #[test]
1346 fn computed_subcommand_selector_errors() {
1347 let schema = kj_schema();
1348 let args = vec![Arg::Positional(Expr::CommandSubst(vec![
1351 crate::ast::Stmt::Command(crate::ast::Command {
1352 name: "echo".into(),
1353 args: vec![],
1354 redirects: vec![],
1355 }),
1356 ]))];
1357 let err = select_leaf(&schema, &args).expect_err("must error");
1358 let msg = err.to_string();
1359 assert!(msg.contains("subcommand name is required"), "got: {msg}");
1360 assert!(msg.contains("command substitution"), "names the cause: {msg}");
1361 }
1362
1363 #[test]
1364 fn variable_subcommand_selector_errors() {
1365 let schema = kj_schema();
1366 let args = vec![Arg::Positional(Expr::VarRef(crate::ast::VarPath::simple("sub")))];
1367 let err = select_leaf(&schema, &args).expect_err("must error");
1368 assert!(err.to_string().contains("variable reference"), "got: {err}");
1369 }
1370
1371 #[test]
1372 fn computed_positional_after_leaf_is_fine() {
1373 let schema = kj_schema();
1374 let args = vec![
1377 word("context"),
1378 word("list"),
1379 Arg::Positional(Expr::CommandSubst(vec![crate::ast::Stmt::Command(
1380 crate::ast::Command { name: "echo".into(), args: vec![], redirects: vec![] },
1381 )])),
1382 ];
1383 let leaf = select_leaf(&schema, &args).expect("ok");
1384 assert_eq!(leaf.name, "list");
1385 }
1386}
1387
1388#[cfg(test)]
1389mod tests {
1390 use super::*;
1391 use crate::dispatch::BackendDispatcher;
1392 use crate::tools::register_builtins;
1393 use crate::vfs::{Filesystem, MemoryFs, VfsRouter};
1394 use std::path::Path;
1395
1396 async fn make_runner_and_ctx() -> (PipelineRunner, ExecContext, BackendDispatcher) {
1397 let mut tools = ToolRegistry::new();
1398 register_builtins(&mut tools);
1399 let tools = Arc::new(tools);
1400 let runner = PipelineRunner::new(tools.clone());
1401 let dispatcher = BackendDispatcher::new(tools.clone());
1402
1403 let mut vfs = VfsRouter::new();
1404 let mem = MemoryFs::new();
1405 mem.write(Path::new("test.txt"), b"hello\nworld\nfoo").await.unwrap();
1406 vfs.mount("/", mem);
1407 let ctx = ExecContext::with_vfs_and_tools(Arc::new(vfs), tools);
1408
1409 (runner, ctx, dispatcher)
1410 }
1411
1412 fn make_cmd(name: &str, args: Vec<&str>) -> Command {
1413 Command {
1414 name: name.to_string(),
1415 args: args.iter().map(|s| Arg::Positional(Expr::Literal(Value::String(s.to_string())))).collect(),
1416 redirects: vec![],
1417 }
1418 }
1419
1420 #[tokio::test]
1421 async fn test_single_command() {
1422 let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
1423 let cmd = make_cmd("echo", vec!["hello"]);
1424
1425 let result = runner.run(&[cmd], &mut ctx, &dispatcher).await;
1426 assert!(result.ok());
1427 assert_eq!(result.text_out().trim(), "hello");
1428 }
1429
1430 #[tokio::test]
1431 async fn test_pipeline_echo_grep() {
1432 let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
1433
1434 let echo_cmd = Command {
1436 name: "echo".to_string(),
1437 args: vec![Arg::Positional(Expr::Literal(Value::String("hello\nworld".to_string())))],
1438 redirects: vec![],
1439 };
1440 let grep_cmd = Command {
1441 name: "grep".to_string(),
1442 args: vec![Arg::Positional(Expr::Literal(Value::String("world".to_string())))],
1443 redirects: vec![],
1444 };
1445
1446 let result = runner.run(&[echo_cmd, grep_cmd], &mut ctx, &dispatcher).await;
1447 assert!(result.ok());
1448 assert_eq!(result.text_out().trim(), "world");
1449 }
1450
1451 #[tokio::test]
1452 async fn test_pipeline_cat_grep() {
1453 let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
1454
1455 let cat_cmd = make_cmd("cat", vec!["/test.txt"]);
1457 let grep_cmd = Command {
1458 name: "grep".to_string(),
1459 args: vec![Arg::Positional(Expr::Literal(Value::String("hello".to_string())))],
1460 redirects: vec![],
1461 };
1462
1463 let result = runner.run(&[cat_cmd, grep_cmd], &mut ctx, &dispatcher).await;
1464 assert!(result.ok());
1465 assert!(result.text_out().contains("hello"));
1466 }
1467
1468 #[tokio::test]
1469 async fn test_command_not_found() {
1470 let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
1471 let cmd = make_cmd("nonexistent", vec![]);
1472
1473 let result = runner.run(&[cmd], &mut ctx, &dispatcher).await;
1474 assert!(!result.ok());
1475 assert_eq!(result.code, 127);
1476 assert!(result.err.contains("not found"));
1477 }
1478
1479 #[tokio::test]
1480 async fn test_pipeline_continues_on_failure() {
1481 let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
1484
1485 let cat_cmd = make_cmd("cat", vec!["/nonexistent"]);
1488 let grep_cmd = Command {
1489 name: "grep".to_string(),
1490 args: vec![Arg::Positional(Expr::Literal(Value::String("hello".to_string())))],
1491 redirects: vec![],
1492 };
1493
1494 let result = runner.run(&[cat_cmd, grep_cmd], &mut ctx, &dispatcher).await;
1495 assert!(!result.ok());
1497 }
1498
1499 #[tokio::test]
1500 async fn test_pipeline_last_command_exit_code() {
1501 let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
1503
1504 let echo_cmd = make_cmd("echo", vec!["hello"]);
1505 let cat_cmd = make_cmd("cat", vec![]);
1506
1507 let result = runner.run(&[echo_cmd, cat_cmd], &mut ctx, &dispatcher).await;
1508 assert!(result.ok());
1509 assert!(result.text_out().contains("hello"));
1510 }
1511
1512 #[tokio::test]
1513 async fn test_empty_pipeline() {
1514 let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
1515 let result = runner.run(&[], &mut ctx, &dispatcher).await;
1516 assert!(result.ok());
1517 }
1518
1519 #[test]
1522 fn test_find_scatter_gather_both_present() {
1523 let commands = vec![
1524 make_cmd("echo", vec!["a"]),
1525 make_cmd("scatter", vec![]),
1526 make_cmd("process", vec![]),
1527 make_cmd("gather", vec![]),
1528 ];
1529 let result = find_scatter_gather(&commands);
1530 assert_eq!(result, Some((1, 3)));
1531 }
1532
1533 #[test]
1534 fn test_find_scatter_gather_no_scatter() {
1535 let commands = vec![
1536 make_cmd("echo", vec!["a"]),
1537 make_cmd("gather", vec![]),
1538 ];
1539 let result = find_scatter_gather(&commands);
1540 assert!(result.is_none());
1541 }
1542
1543 #[test]
1544 fn test_find_scatter_gather_no_gather() {
1545 let commands = vec![
1546 make_cmd("echo", vec!["a"]),
1547 make_cmd("scatter", vec![]),
1548 ];
1549 let result = find_scatter_gather(&commands);
1550 assert!(result.is_none());
1551 }
1552
1553 #[test]
1554 fn test_find_scatter_gather_wrong_order() {
1555 let commands = vec![
1556 make_cmd("gather", vec![]),
1557 make_cmd("scatter", vec![]),
1558 ];
1559 let result = find_scatter_gather(&commands);
1560 assert!(result.is_none());
1561 }
1562
1563 #[tokio::test]
1564 async fn test_scatter_gather_simple() {
1565 let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
1566
1567 let split_cmd = Command {
1569 name: "split".to_string(),
1570 args: vec![Arg::Positional(Expr::Literal(Value::String("a b c".to_string())))],
1571 redirects: vec![],
1572 };
1573 let scatter_cmd = make_cmd("scatter", vec![]);
1574 let process_cmd = Command {
1575 name: "echo".to_string(),
1576 args: vec![Arg::Positional(Expr::VarRef(crate::ast::VarPath::simple("ITEM")))],
1577 redirects: vec![],
1578 };
1579 let gather_cmd = make_cmd("gather", vec![]);
1580
1581 let result = runner.run(&[split_cmd, scatter_cmd, process_cmd, gather_cmd], &mut ctx, &dispatcher).await;
1582 assert!(result.ok(), "scatter with structured data should succeed: {}", result.err);
1583 assert!(result.text_out().contains("a"));
1585 assert!(result.text_out().contains("b"));
1586 assert!(result.text_out().contains("c"));
1587 }
1588
1589 #[tokio::test]
1590 async fn test_scatter_gather_empty_input() {
1591 let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
1592
1593 let echo_cmd = Command {
1595 name: "echo".to_string(),
1596 args: vec![Arg::Positional(Expr::Literal(Value::String("".to_string())))],
1597 redirects: vec![],
1598 };
1599 let scatter_cmd = make_cmd("scatter", vec![]);
1600 let process_cmd = Command {
1601 name: "echo".to_string(),
1602 args: vec![Arg::Positional(Expr::VarRef(crate::ast::VarPath::simple("ITEM")))],
1603 redirects: vec![],
1604 };
1605 let gather_cmd = make_cmd("gather", vec![]);
1606
1607 let result = runner.run(&[echo_cmd, scatter_cmd, process_cmd, gather_cmd], &mut ctx, &dispatcher).await;
1608 assert!(result.ok());
1609 assert!(result.text_out().trim().is_empty());
1610 }
1611
1612 #[tokio::test]
1613 async fn test_scatter_gather_with_structured_stdin() {
1614 let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
1615
1616 let data = Value::Json(serde_json::json!(["x", "y", "z"]));
1618 ctx.set_stdin_with_data("x\ny\nz".to_string(), Some(data));
1619
1620 let scatter_cmd = make_cmd("scatter", vec![]);
1621 let process_cmd = Command {
1622 name: "echo".to_string(),
1623 args: vec![Arg::Positional(Expr::VarRef(crate::ast::VarPath::simple("ITEM")))],
1624 redirects: vec![],
1625 };
1626 let gather_cmd = make_cmd("gather", vec![]);
1627
1628 let result = runner.run(&[scatter_cmd, process_cmd, gather_cmd], &mut ctx, &dispatcher).await;
1629 assert!(result.ok(), "scatter with structured stdin should succeed: {}", result.err);
1630 assert!(result.text_out().contains("x"));
1631 assert!(result.text_out().contains("y"));
1632 assert!(result.text_out().contains("z"));
1633 }
1634
1635 #[tokio::test]
1636 async fn test_scatter_gather_json_input() {
1637 let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
1638
1639 let data = Value::Json(serde_json::json!(["one", "two", "three"]));
1641 ctx.set_stdin_with_data(r#"["one", "two", "three"]"#.to_string(), Some(data));
1642
1643 let scatter_cmd = make_cmd("scatter", vec![]);
1644 let process_cmd = Command {
1645 name: "echo".to_string(),
1646 args: vec![Arg::Positional(Expr::VarRef(crate::ast::VarPath::simple("ITEM")))],
1647 redirects: vec![],
1648 };
1649 let gather_cmd = make_cmd("gather", vec![]);
1650
1651 let result = runner.run(&[scatter_cmd, process_cmd, gather_cmd], &mut ctx, &dispatcher).await;
1652 assert!(result.ok(), "scatter with JSON data should succeed: {}", result.err);
1653 assert!(result.text_out().contains("one"));
1654 assert!(result.text_out().contains("two"));
1655 assert!(result.text_out().contains("three"));
1656 }
1657
1658 #[tokio::test]
1659 async fn test_scatter_gather_with_post_gather() {
1660 let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
1661
1662 let split_cmd = Command {
1664 name: "split".to_string(),
1665 args: vec![Arg::Positional(Expr::Literal(Value::String("a b".to_string())))],
1666 redirects: vec![],
1667 };
1668 let scatter_cmd = make_cmd("scatter", vec![]);
1669 let process_cmd = Command {
1670 name: "echo".to_string(),
1671 args: vec![Arg::Positional(Expr::VarRef(crate::ast::VarPath::simple("ITEM")))],
1672 redirects: vec![],
1673 };
1674 let gather_cmd = make_cmd("gather", vec![]);
1675 let grep_cmd = Command {
1676 name: "grep".to_string(),
1677 args: vec![Arg::Positional(Expr::Literal(Value::String("a".to_string())))],
1678 redirects: vec![],
1679 };
1680
1681 let result = runner.run(&[split_cmd, scatter_cmd, process_cmd, gather_cmd, grep_cmd], &mut ctx, &dispatcher).await;
1682 assert!(result.ok(), "scatter with post_gather should succeed: {}", result.err);
1683 assert!(result.text_out().contains("a"));
1684 assert!(!result.text_out().contains("b"));
1685 }
1686
1687 #[tokio::test]
1688 async fn test_scatter_custom_var_name() {
1689 let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
1690
1691 let data = Value::Json(serde_json::json!(["test1", "test2"]));
1693 ctx.set_stdin_with_data("test1\ntest2".to_string(), Some(data));
1694
1695 let scatter_cmd = Command {
1697 name: "scatter".to_string(),
1698 args: vec![Arg::Named {
1699 key: "as".to_string(),
1700 value: Expr::Literal(Value::String("URL".to_string())),
1701 }],
1702 redirects: vec![],
1703 };
1704 let process_cmd = Command {
1705 name: "echo".to_string(),
1706 args: vec![Arg::Positional(Expr::VarRef(crate::ast::VarPath::simple("URL")))],
1707 redirects: vec![],
1708 };
1709 let gather_cmd = make_cmd("gather", vec![]);
1710
1711 let result = runner.run(&[scatter_cmd, process_cmd, gather_cmd], &mut ctx, &dispatcher).await;
1712 assert!(result.ok(), "scatter with custom var should succeed: {}", result.err);
1713 assert!(result.text_out().contains("test1"));
1714 assert!(result.text_out().contains("test2"));
1715 }
1716
1717 #[tokio::test]
1720 async fn test_pipeline_routes_through_backend() {
1721 use crate::backend::testing::MockBackend;
1722 use std::sync::atomic::Ordering;
1723
1724 let (backend, call_count) = MockBackend::new();
1726 let backend: std::sync::Arc<dyn crate::backend::KernelBackend> = std::sync::Arc::new(backend);
1727
1728 let mut ctx = crate::tools::ExecContext::with_backend(backend);
1730
1731 let tools = std::sync::Arc::new(ToolRegistry::new());
1733 let runner = PipelineRunner::new(tools.clone());
1734 let dispatcher = BackendDispatcher::new(tools);
1735
1736 let cmd = make_cmd("test-tool", vec!["arg1"]);
1738 let result = runner.run(&[cmd], &mut ctx, &dispatcher).await;
1739
1740 assert!(result.ok(), "Mock backend should return success");
1741 assert_eq!(call_count.load(Ordering::SeqCst), 1, "call_tool should be invoked once");
1742 assert!(result.text_out().contains("mock executed"), "Output should be from mock backend");
1743 }
1744
1745 #[tokio::test]
1746 async fn test_multi_command_pipeline_routes_through_backend() {
1747 use crate::backend::testing::MockBackend;
1748 use std::sync::atomic::Ordering;
1749
1750 let (backend, call_count) = MockBackend::new();
1751 let backend: std::sync::Arc<dyn crate::backend::KernelBackend> = std::sync::Arc::new(backend);
1752 let mut ctx = crate::tools::ExecContext::with_backend(backend);
1753
1754 let tools = std::sync::Arc::new(ToolRegistry::new());
1755 let runner = PipelineRunner::new(tools.clone());
1756 let dispatcher = BackendDispatcher::new(tools);
1757
1758 let cmd1 = make_cmd("tool1", vec![]);
1760 let cmd2 = make_cmd("tool2", vec![]);
1761 let cmd3 = make_cmd("tool3", vec![]);
1762
1763 let result = runner.run(&[cmd1, cmd2, cmd3], &mut ctx, &dispatcher).await;
1764
1765 assert!(result.ok());
1766 assert_eq!(call_count.load(Ordering::SeqCst), 3, "call_tool should be invoked for each command");
1767 }
1768
1769 #[tokio::test]
1776 async fn backend_dispatcher_scalar_data_matches_production_unwrap() {
1777 use crate::backend::testing::MockBackend;
1778 use crate::backend::ToolResult;
1779
1780 let (mock, _calls) = MockBackend::new();
1781 let backend = mock.with_tool_result(|_name| Ok(ToolResult::with_data("", serde_json::json!(42))));
1782 let backend: Arc<dyn crate::backend::KernelBackend> = Arc::new(backend);
1783 let mut ctx = ExecContext::with_backend(backend);
1784
1785 let dispatcher = BackendDispatcher::new(Arc::new(ToolRegistry::new()));
1786 let cmd = make_cmd("embedder_tool", vec![]);
1787
1788 let result = dispatcher.dispatch(&cmd, &mut ctx).await.expect("dispatch");
1789 assert_eq!(
1790 result.data,
1791 Some(Value::Int(42)),
1792 "a scalar ToolResult.data must unwrap to a native Value, matching \
1793 the production From<ToolResult> path — not stay Value::Json(42)"
1794 );
1795 }
1796
1797 #[tokio::test]
1803 async fn backend_dispatcher_envelope_shaped_data_stays_structured() {
1804 use crate::backend::testing::MockBackend;
1805 use crate::backend::ToolResult;
1806
1807 let envelope = kaish_types::bytes_to_envelope(&[1u8, 2, 3]);
1808 let (mock, _calls) = MockBackend::new();
1809 let backend = mock.with_tool_result(move |_name| Ok(ToolResult::with_data("", envelope.clone())));
1810 let backend: Arc<dyn crate::backend::KernelBackend> = Arc::new(backend);
1811 let mut ctx = ExecContext::with_backend(backend);
1812
1813 let dispatcher = BackendDispatcher::new(Arc::new(ToolRegistry::new()));
1814 let cmd = make_cmd("embedder_tool", vec![]);
1815
1816 let result = dispatcher.dispatch(&cmd, &mut ctx).await.expect("dispatch");
1817 assert!(
1818 matches!(result.data, Some(Value::Json(_))),
1819 "envelope-shaped external data must stay structured, not silently \
1820 decode to Value::Bytes: got {:?}",
1821 result.data
1822 );
1823 }
1824
1825 #[tokio::test]
1830 async fn backend_dispatcher_preserves_did_spill_and_original_code() {
1831 use crate::backend::testing::MockBackend;
1832 use crate::backend::ToolResult;
1833
1834 let (mock, _calls) = MockBackend::new();
1835 let backend = mock.with_tool_result(|_name| {
1836 Ok(ToolResult::success("truncated...")
1837 .with_did_spill(true)
1838 .with_original_code(Some(0)))
1839 });
1840 let backend: Arc<dyn crate::backend::KernelBackend> = Arc::new(backend);
1841 let mut ctx = ExecContext::with_backend(backend);
1842
1843 let dispatcher = BackendDispatcher::new(Arc::new(ToolRegistry::new()));
1844 let cmd = make_cmd("embedder_tool", vec![]);
1845
1846 let result = dispatcher.dispatch(&cmd, &mut ctx).await.expect("dispatch");
1847 assert!(result.did_spill, "did_spill must survive the backend seam");
1848 assert_eq!(result.original_code, Some(0), "original_code must survive the backend seam");
1849 }
1850
1851 use crate::tools::{ParamSchema, ToolSchema};
1854
1855 fn make_test_schema() -> ToolSchema {
1856 ToolSchema::new("test-tool", "A test tool for schema-aware parsing")
1857 .param(ParamSchema::required("query", "string", "Search query"))
1858 .param(ParamSchema::optional("limit", "int", Value::Int(10), "Max results"))
1859 .param(ParamSchema::optional("verbose", "bool", Value::Bool(false), "Verbose output"))
1860 .param(ParamSchema::optional("output", "string", Value::String("stdout".into()), "Output destination"))
1861 .with_positional_mapping()
1862 }
1863
1864 fn make_minimal_ctx() -> ExecContext {
1865 let mut vfs = VfsRouter::new();
1866 vfs.mount("/", MemoryFs::new());
1867 ExecContext::new(Arc::new(vfs))
1868 }
1869
1870 fn test_dispatcher() -> BackendDispatcher {
1874 BackendDispatcher::new(Arc::new(ToolRegistry::new()))
1875 }
1876
1877 #[test]
1878 fn test_schema_aware_string_arg() {
1879 let args = vec![
1881 Arg::LongFlag("query".to_string()),
1882 Arg::Positional(Expr::Literal(Value::String("test".to_string()))),
1883 ];
1884 let schema = make_test_schema();
1885 let ctx = make_minimal_ctx();
1886
1887 let tool_args = build_tool_args(&args, &ctx, Some(&schema)).expect("build_tool_args");
1888
1889 assert!(tool_args.flags.is_empty(), "No flags should be set");
1890 assert!(tool_args.positional.is_empty(), "No positionals - consumed by --query");
1891 assert_eq!(
1892 tool_args.named.get("query"),
1893 Some(&Value::String("test".to_string())),
1894 "--query should consume 'test' as its value"
1895 );
1896 }
1897
1898 #[test]
1899 fn test_schema_aware_bool_flag() {
1900 let args = vec![
1902 Arg::LongFlag("verbose".to_string()),
1903 ];
1904 let schema = make_test_schema();
1905 let ctx = make_minimal_ctx();
1906
1907 let tool_args = build_tool_args(&args, &ctx, Some(&schema)).expect("build_tool_args");
1908
1909 assert!(tool_args.flags.contains("verbose"), "--verbose should be a flag");
1910 assert!(tool_args.named.is_empty(), "No named args");
1911 assert!(tool_args.positional.is_empty(), "No positionals");
1912 }
1913
1914 #[test]
1915 fn test_schema_aware_mixed() {
1916 let args = vec![
1919 Arg::Positional(Expr::Literal(Value::String("file.txt".to_string()))),
1920 Arg::LongFlag("output".to_string()),
1921 Arg::Positional(Expr::Literal(Value::String("out.txt".to_string()))),
1922 Arg::LongFlag("verbose".to_string()),
1923 ];
1924 let schema = make_test_schema();
1925 let ctx = make_minimal_ctx();
1926
1927 let tool_args = build_tool_args(&args, &ctx, Some(&schema)).expect("build_tool_args");
1928
1929 assert!(tool_args.positional.is_empty(), "file.txt consumed as query param");
1930 assert_eq!(
1931 tool_args.named.get("query"),
1932 Some(&Value::String("file.txt".to_string()))
1933 );
1934 assert_eq!(
1935 tool_args.named.get("output"),
1936 Some(&Value::String("out.txt".to_string()))
1937 );
1938 assert!(tool_args.flags.contains("verbose"));
1939 }
1940
1941 #[test]
1942 fn test_schema_aware_multiple_string_args() {
1943 let args = vec![
1945 Arg::LongFlag("query".to_string()),
1946 Arg::Positional(Expr::Literal(Value::String("test".to_string()))),
1947 Arg::LongFlag("output".to_string()),
1948 Arg::Positional(Expr::Literal(Value::String("result.json".to_string()))),
1949 Arg::LongFlag("verbose".to_string()),
1950 Arg::LongFlag("limit".to_string()),
1951 Arg::Positional(Expr::Literal(Value::Int(5))),
1952 ];
1953 let schema = make_test_schema();
1954 let ctx = make_minimal_ctx();
1955
1956 let tool_args = build_tool_args(&args, &ctx, Some(&schema)).expect("build_tool_args");
1957
1958 assert!(tool_args.positional.is_empty(), "All positionals consumed");
1959 assert_eq!(
1960 tool_args.named.get("query"),
1961 Some(&Value::String("test".to_string()))
1962 );
1963 assert_eq!(
1964 tool_args.named.get("output"),
1965 Some(&Value::String("result.json".to_string()))
1966 );
1967 assert_eq!(
1968 tool_args.named.get("limit"),
1969 Some(&Value::Int(5))
1970 );
1971 assert!(tool_args.flags.contains("verbose"));
1972 }
1973
1974 #[test]
1975 fn test_schema_aware_double_dash() {
1976 let args = vec![
1979 Arg::LongFlag("output".to_string()),
1980 Arg::Positional(Expr::Literal(Value::String("out.txt".to_string()))),
1981 Arg::DoubleDash,
1982 Arg::Positional(Expr::Literal(Value::String("--this-is-data".to_string()))),
1983 ];
1984 let schema = make_test_schema();
1985 let ctx = make_minimal_ctx();
1986
1987 let tool_args = build_tool_args(&args, &ctx, Some(&schema)).expect("build_tool_args");
1988
1989 assert_eq!(
1990 tool_args.named.get("output"),
1991 Some(&Value::String("out.txt".to_string()))
1992 );
1993 assert_eq!(
1995 tool_args.positional,
1996 vec![Value::String("--this-is-data".to_string())]
1997 );
1998 }
1999
2000 #[test]
2006 fn word_assign_binary_value_is_loud_not_placeholder() {
2007 let args = vec![Arg::WordAssign {
2008 key: "if".to_string(),
2009 value: Expr::Literal(Value::Bytes(vec![0xff, 0x00, 0xfe])),
2010 }];
2011 let ctx = make_minimal_ctx();
2012
2013 let err = build_tool_args(&args, &ctx, None).expect_err("binary WordAssign must error");
2016 assert!(
2017 err.contains("cannot be used as"),
2018 "error should name the binary problem, got {err:?}"
2019 );
2020 }
2021
2022 #[test]
2023 fn test_no_schema_fallback() {
2024 let args = vec![
2026 Arg::LongFlag("query".to_string()),
2027 Arg::Positional(Expr::Literal(Value::String("test".to_string()))),
2028 ];
2029 let ctx = make_minimal_ctx();
2030
2031 let tool_args = build_tool_args(&args, &ctx, None).expect("build_tool_args");
2032
2033 assert!(tool_args.flags.contains("query"), "--query should be a flag");
2035 assert_eq!(
2036 tool_args.positional,
2037 vec![Value::String("test".to_string())],
2038 "'test' should be a positional"
2039 );
2040 }
2041
2042 #[test]
2043 fn test_unknown_flag_in_schema() {
2044 let args = vec![
2046 Arg::LongFlag("unknown".to_string()),
2047 Arg::Positional(Expr::Literal(Value::String("value".to_string()))),
2048 ];
2049 let schema = make_test_schema();
2050 let ctx = make_minimal_ctx();
2051
2052 let tool_args = build_tool_args(&args, &ctx, Some(&schema)).expect("build_tool_args");
2053
2054 assert!(tool_args.flags.contains("unknown"));
2055 assert!(tool_args.positional.is_empty(), "value consumed as query param");
2056 assert_eq!(
2057 tool_args.named.get("query"),
2058 Some(&Value::String("value".to_string()))
2059 );
2060 }
2061
2062 #[test]
2063 fn test_named_args_unchanged() {
2064 let args = vec![
2066 Arg::Named {
2067 key: "query".to_string(),
2068 value: Expr::Literal(Value::String("test".to_string())),
2069 },
2070 Arg::LongFlag("verbose".to_string()),
2071 ];
2072 let schema = make_test_schema();
2073 let ctx = make_minimal_ctx();
2074
2075 let tool_args = build_tool_args(&args, &ctx, Some(&schema)).expect("build_tool_args");
2076
2077 assert_eq!(
2078 tool_args.named.get("query"),
2079 Some(&Value::String("test".to_string()))
2080 );
2081 assert!(tool_args.flags.contains("verbose"));
2082 }
2083
2084 #[test]
2085 fn test_short_flags_unchanged() {
2086 let args = vec![
2088 Arg::ShortFlag("la".to_string()),
2089 Arg::Positional(Expr::Literal(Value::String("file.txt".to_string()))),
2090 ];
2091 let schema = make_test_schema();
2092 let ctx = make_minimal_ctx();
2093
2094 let tool_args = build_tool_args(&args, &ctx, Some(&schema)).expect("build_tool_args");
2095
2096 assert!(tool_args.flags.contains("l"));
2097 assert!(tool_args.flags.contains("a"));
2098 assert!(tool_args.positional.is_empty(), "file.txt consumed as query param");
2099 assert_eq!(
2100 tool_args.named.get("query"),
2101 Some(&Value::String("file.txt".to_string()))
2102 );
2103 }
2104
2105 #[test]
2106 fn test_flag_at_end_no_value() {
2107 let args = vec![
2110 Arg::Positional(Expr::Literal(Value::String("file.txt".to_string()))),
2111 Arg::LongFlag("output".to_string()),
2112 ];
2113 let schema = make_test_schema();
2114 let ctx = make_minimal_ctx();
2115
2116 let tool_args = build_tool_args(&args, &ctx, Some(&schema)).expect("build_tool_args");
2117
2118 assert!(tool_args.flags.contains("output"));
2120 assert!(tool_args.positional.is_empty(), "file.txt consumed as query param");
2121 assert_eq!(
2122 tool_args.named.get("query"),
2123 Some(&Value::String("file.txt".to_string()))
2124 );
2125 }
2126
2127 #[test]
2128 fn test_positional_skips_bool_params() {
2129 let schema = ToolSchema::new("test", "")
2133 .param(ParamSchema::required("query", "string", ""))
2134 .param(ParamSchema::optional(
2135 "verbose",
2136 "bool",
2137 Value::Bool(false),
2138 "",
2139 ))
2140 .param(ParamSchema::optional(
2141 "output",
2142 "string",
2143 Value::Null,
2144 "",
2145 ))
2146 .with_positional_mapping();
2147 let args = vec![
2148 Arg::Positional(Expr::Literal(Value::String("val1".to_string()))),
2149 Arg::Positional(Expr::Literal(Value::String("val2".to_string()))),
2150 ];
2151 let ctx = make_minimal_ctx();
2152
2153 let tool_args = build_tool_args(&args, &ctx, Some(&schema)).expect("build_tool_args");
2154
2155 assert_eq!(
2156 tool_args.named.get("query"),
2157 Some(&Value::String("val1".to_string()))
2158 );
2159 assert_eq!(
2160 tool_args.named.get("output"),
2161 Some(&Value::String("val2".to_string()))
2162 );
2163 assert!(!tool_args.flags.contains("verbose"));
2164 assert!(tool_args.positional.is_empty());
2165 }
2166
2167 #[test]
2168 fn test_positionals_fill_available_slots() {
2169 let args = vec![
2172 Arg::Positional(Expr::Literal(Value::String("val1".to_string()))),
2173 Arg::Positional(Expr::Literal(Value::String("val2".to_string()))),
2174 Arg::Positional(Expr::Literal(Value::String("val3".to_string()))),
2175 ];
2176 let schema = make_test_schema(); let ctx = make_minimal_ctx();
2178
2179 let tool_args = build_tool_args(&args, &ctx, Some(&schema)).expect("build_tool_args");
2180
2181 assert_eq!(
2184 tool_args.named.get("query"),
2185 Some(&Value::String("val1".to_string()))
2186 );
2187 assert_eq!(
2188 tool_args.named.get("limit"),
2189 Some(&Value::String("val2".to_string()))
2190 );
2191 assert_eq!(
2192 tool_args.named.get("output"),
2193 Some(&Value::String("val3".to_string()))
2194 );
2195 assert!(tool_args.positional.is_empty());
2196 }
2197
2198 #[test]
2199 fn test_truly_excess_positionals() {
2200 let schema = ToolSchema::new("test", "")
2202 .param(ParamSchema::required("name", "string", ""))
2203 .with_positional_mapping();
2204 let args = vec![
2205 Arg::Positional(Expr::Literal(Value::String("first".to_string()))),
2206 Arg::Positional(Expr::Literal(Value::String("second".to_string()))),
2207 Arg::Positional(Expr::Literal(Value::String("third".to_string()))),
2208 ];
2209 let ctx = make_minimal_ctx();
2210
2211 let tool_args = build_tool_args(&args, &ctx, Some(&schema)).expect("build_tool_args");
2212
2213 assert_eq!(
2214 tool_args.named.get("name"),
2215 Some(&Value::String("first".to_string()))
2216 );
2217 assert_eq!(
2218 tool_args.positional,
2219 vec![
2220 Value::String("second".to_string()),
2221 Value::String("third".to_string()),
2222 ]
2223 );
2224 }
2225
2226 #[test]
2227 fn test_double_dash_positional_not_mapped() {
2228 let args = vec![
2230 Arg::Positional(Expr::Literal(Value::String("val1".to_string()))),
2231 Arg::DoubleDash,
2232 Arg::Positional(Expr::Literal(Value::String("val2".to_string()))),
2233 ];
2234 let schema = make_test_schema();
2235 let ctx = make_minimal_ctx();
2236
2237 let tool_args = build_tool_args(&args, &ctx, Some(&schema)).expect("build_tool_args");
2238
2239 assert_eq!(
2240 tool_args.named.get("query"),
2241 Some(&Value::String("val1".to_string()))
2242 );
2243 assert_eq!(
2245 tool_args.positional,
2246 vec![Value::String("val2".to_string())]
2247 );
2248 }
2249
2250 #[test]
2251 fn test_all_params_filled_by_flags() {
2252 let args = vec![
2254 Arg::LongFlag("query".to_string()),
2255 Arg::Positional(Expr::Literal(Value::String("search".to_string()))),
2256 Arg::LongFlag("output".to_string()),
2257 Arg::Positional(Expr::Literal(Value::String("out.txt".to_string()))),
2258 Arg::LongFlag("verbose".to_string()),
2259 ];
2260 let schema = make_test_schema();
2261 let ctx = make_minimal_ctx();
2262
2263 let tool_args = build_tool_args(&args, &ctx, Some(&schema)).expect("build_tool_args");
2264
2265 assert_eq!(
2266 tool_args.named.get("query"),
2267 Some(&Value::String("search".to_string()))
2268 );
2269 assert_eq!(
2270 tool_args.named.get("output"),
2271 Some(&Value::String("out.txt".to_string()))
2272 );
2273 assert!(tool_args.flags.contains("verbose"));
2274 assert!(tool_args.positional.is_empty());
2275 }
2276
2277 #[test]
2278 fn test_mixed_flags_and_positional_fill() {
2279 let args = vec![
2281 Arg::LongFlag("output".to_string()),
2282 Arg::Positional(Expr::Literal(Value::String("foo".to_string()))),
2283 Arg::Positional(Expr::Literal(Value::String("val1".to_string()))),
2284 ];
2285 let schema = make_test_schema();
2286 let ctx = make_minimal_ctx();
2287
2288 let tool_args = build_tool_args(&args, &ctx, Some(&schema)).expect("build_tool_args");
2289
2290 assert_eq!(
2291 tool_args.named.get("output"),
2292 Some(&Value::String("foo".to_string()))
2293 );
2294 assert_eq!(
2295 tool_args.named.get("query"),
2296 Some(&Value::String("val1".to_string()))
2297 );
2298 assert!(tool_args.positional.is_empty());
2299 }
2300
2301 #[test]
2302 fn test_alias_flag_prevents_mapping_overwrite() {
2303 let schema = ToolSchema::new("test", "")
2305 .param(ParamSchema::required("query", "string", "").with_aliases(["-q"]))
2306 .param(ParamSchema::required("output", "string", ""))
2307 .with_positional_mapping();
2308 let args = vec![
2309 Arg::ShortFlag("q".to_string()),
2310 Arg::Positional(Expr::Literal(Value::String("search".to_string()))),
2311 Arg::Positional(Expr::Literal(Value::String("out.txt".to_string()))),
2312 ];
2313 let ctx = make_minimal_ctx();
2314
2315 let tool_args = build_tool_args(&args, &ctx, Some(&schema)).expect("build_tool_args");
2316
2317 assert_eq!(
2318 tool_args.named.get("query"),
2319 Some(&Value::String("search".to_string()))
2320 );
2321 assert_eq!(
2322 tool_args.named.get("output"),
2323 Some(&Value::String("out.txt".to_string()))
2324 );
2325 assert!(tool_args.positional.is_empty());
2326 }
2327
2328 #[test]
2329 fn test_builtin_schema_no_positional_mapping() {
2330 let schema = ToolSchema::new("echo", "")
2332 .param(ParamSchema::optional("args", "any", Value::Null, ""))
2333 .param(ParamSchema::optional("no_newline", "bool", Value::Bool(false), ""));
2334 let args = vec![
2336 Arg::Positional(Expr::Literal(Value::String("hello".to_string()))),
2337 Arg::Positional(Expr::Literal(Value::String("world".to_string()))),
2338 ];
2339 let ctx = make_minimal_ctx();
2340
2341 let tool_args = build_tool_args(&args, &ctx, Some(&schema)).expect("build_tool_args");
2342
2343 assert_eq!(
2345 tool_args.positional,
2346 vec![
2347 Value::String("hello".to_string()),
2348 Value::String("world".to_string()),
2349 ]
2350 );
2351 assert!(!tool_args.named.contains_key("args"));
2352 }
2353
2354 #[test]
2355 fn test_short_flag_with_alias_consumes_value() {
2356 let schema = ToolSchema::new("head", "Output first part of files")
2359 .param(ParamSchema::optional("lines", "int", Value::Int(10), "Number of lines")
2360 .with_aliases(["-n"]));
2361 let args = vec![
2362 Arg::ShortFlag("n".to_string()),
2363 Arg::Positional(Expr::Literal(Value::Int(5))),
2364 Arg::Positional(Expr::Literal(Value::String("/tmp/file.txt".to_string()))),
2365 ];
2366 let ctx = make_minimal_ctx();
2367
2368 let tool_args = build_tool_args(&args, &ctx, Some(&schema)).expect("build_tool_args");
2369
2370 assert!(tool_args.flags.is_empty(), "no boolean flags: {:?}", tool_args.flags);
2371 assert_eq!(tool_args.named.get("lines"), Some(&Value::Int(5)), "should resolve alias to canonical name");
2372 assert_eq!(tool_args.positional, vec![Value::String("/tmp/file.txt".to_string())]);
2373 }
2374
2375 #[tokio::test]
2378 async fn test_merge_stderr_redirect() {
2379 let result = ExecResult::from_output(0, "stdout content", "stderr content");
2381
2382 let redirects = vec![Redirect {
2383 kind: RedirectKind::MergeStderr,
2384 target: Expr::Literal(Value::Null),
2385 }];
2386
2387 let ctx = make_minimal_ctx();
2388 let result = apply_redirects(result, &redirects, &ctx, &test_dispatcher()).await;
2389
2390 assert_eq!(&*result.text_out(), "stdout contentstderr content");
2391 assert!(result.err.is_empty());
2392 }
2393
2394 #[tokio::test]
2395 async fn test_merge_stderr_with_empty_stderr() {
2396 let result = ExecResult::from_output(0, "stdout only", "");
2398
2399 let redirects = vec![Redirect {
2400 kind: RedirectKind::MergeStderr,
2401 target: Expr::Literal(Value::Null),
2402 }];
2403
2404 let ctx = make_minimal_ctx();
2405 let result = apply_redirects(result, &redirects, &ctx, &test_dispatcher()).await;
2406
2407 assert_eq!(&*result.text_out(), "stdout only");
2408 assert!(result.err.is_empty());
2409 }
2410
2411 #[tokio::test]
2412 async fn test_merge_stderr_order_matters() {
2413 let result = ExecResult::from_output(0, "stdout\n", "stderr\n");
2418
2419 let redirects = vec![Redirect {
2421 kind: RedirectKind::MergeStderr,
2422 target: Expr::Literal(Value::Null),
2423 }];
2424
2425 let ctx = make_minimal_ctx();
2426 let result = apply_redirects(result, &redirects, &ctx, &test_dispatcher()).await;
2427
2428 assert_eq!(&*result.text_out(), "stdout\nstderr\n");
2429 assert!(result.err.is_empty());
2430 }
2431
2432 #[tokio::test]
2433 async fn test_redirect_with_command_execution() {
2434 let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
2435
2436 let cmd = Command {
2438 name: "echo".to_string(),
2439 args: vec![Arg::Positional(Expr::Literal(Value::String("hello".to_string())))],
2440 redirects: vec![Redirect {
2441 kind: RedirectKind::MergeStderr,
2442 target: Expr::Literal(Value::Null),
2443 }],
2444 };
2445
2446 let result = runner.run(&[cmd], &mut ctx, &dispatcher).await;
2447 assert!(result.ok());
2448 assert!(result.text_out().contains("hello"));
2450 }
2451
2452 #[tokio::test]
2453 async fn test_merge_stderr_in_pipeline() {
2454 let (runner, mut ctx, dispatcher) = make_runner_and_ctx().await;
2455
2456 let echo_cmd = Command {
2459 name: "echo".to_string(),
2460 args: vec![Arg::Positional(Expr::Literal(Value::String("output".to_string())))],
2461 redirects: vec![Redirect {
2462 kind: RedirectKind::MergeStderr,
2463 target: Expr::Literal(Value::Null),
2464 }],
2465 };
2466 let grep_cmd = Command {
2467 name: "grep".to_string(),
2468 args: vec![Arg::Positional(Expr::Literal(Value::String("output".to_string())))],
2469 redirects: vec![],
2470 };
2471
2472 let result = runner.run(&[echo_cmd, grep_cmd], &mut ctx, &dispatcher).await;
2473 assert!(result.ok(), "result failed: code={}, err={}", result.code, result.err);
2474 assert!(result.text_out().contains("output"));
2475 }
2476
2477 fn big_table_output(rows: usize) -> crate::interpreter::OutputData {
2489 use crate::interpreter::OutputNode;
2490 let headers = vec!["id".to_string(), "name".to_string()];
2491 let nodes: Vec<OutputNode> = (0..rows)
2492 .map(|i| OutputNode::new(i.to_string()).with_cells(vec![format!("row-{i}")]))
2493 .collect();
2494 crate::interpreter::OutputData::table(headers, nodes)
2495 }
2496
2497 #[tokio::test]
2498 async fn test_both_redirect_streams_structured_output_to_file() {
2499 let output = big_table_output(50);
2504 let expected_stdout = output.to_canonical_string();
2505 let mut result = ExecResult::with_output(output);
2506 result.err = "warning: heads up\n".to_string();
2507
2508 let redirects = vec![Redirect {
2509 kind: RedirectKind::Both,
2510 target: Expr::Literal(Value::String("/out.txt".to_string())),
2511 }];
2512 let ctx = make_minimal_ctx();
2513 let result = apply_redirects(result, &redirects, &ctx, &test_dispatcher()).await;
2514
2515 assert!(result.ok());
2518 assert_eq!(&*result.text_out(), "");
2519 assert!(result.err.is_empty());
2520 assert!(!result.has_output());
2521
2522 let written = ctx.backend.read(Path::new("/out.txt"), None).await.expect("file written");
2523 let written = String::from_utf8(written).expect("valid utf8");
2524 assert_eq!(written, format!("{expected_stdout}warning: heads up\n"));
2528 }
2529
2530 #[tokio::test]
2531 async fn test_both_redirect_streams_large_structured_output_intact() {
2532 let rows = 5_000;
2537 let output = big_table_output(rows);
2538 let expected_stdout = output.to_canonical_string();
2539 let result = ExecResult::with_output(output);
2540
2541 let redirects = vec![Redirect {
2542 kind: RedirectKind::Both,
2543 target: Expr::Literal(Value::String("/big.txt".to_string())),
2544 }];
2545 let ctx = make_minimal_ctx();
2546 let result = apply_redirects(result, &redirects, &ctx, &test_dispatcher()).await;
2547 assert!(result.ok());
2548
2549 let written = ctx.backend.read(Path::new("/big.txt"), None).await.expect("file written");
2550 let written = String::from_utf8(written).expect("valid utf8");
2551 assert_eq!(written, expected_stdout);
2552 assert!(written.contains("row-0"));
2553 assert!(written.contains(&format!("row-{}", rows - 1)));
2554 }
2555
2556 #[tokio::test]
2557 async fn test_both_redirect_still_writes_binary_stdout_raw() {
2558 let result = ExecResult::success_text_or_bytes(vec![0xff, 0x00, 0xfe, b'x']);
2562 let redirects = vec![Redirect {
2563 kind: RedirectKind::Both,
2564 target: Expr::Literal(Value::String("/bin.out".to_string())),
2565 }];
2566 let ctx = make_minimal_ctx();
2567 let result = apply_redirects(result, &redirects, &ctx, &test_dispatcher()).await;
2568 assert!(result.ok());
2569
2570 let written = ctx.backend.read(Path::new("/bin.out"), None).await.expect("file written");
2571 assert_eq!(written, vec![0xff, 0x00, 0xfe, b'x']);
2572 }
2573}