use command_stream::commands::{CommandContext, VirtualCommandRegistry};
use command_stream::zx;
use command_stream::{
cmd, create, exec, run, run_sync, set_shell_option, unset_shell_option, AnsiUtils,
CommandResult, EventData, EventType, OutputChunk, Pipeline, ProcessRunner, RunOptions,
StdinOption, StreamEmitter, StreamingRunner,
};
use serde::Serialize;
use serde_json::{json, Value};
use std::collections::HashMap;
use std::error::Error;
use std::pin::Pin;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
const PARITY_START: &str = "<<<PARITY_JSON";
const PARITY_END: &str = "PARITY_JSON>>>";
#[derive(Serialize)]
struct Observation {
label: &'static str,
value: Value,
}
type ExampleResult = Result<Vec<Observation>, Box<dyn Error>>;
fn observation(label: &'static str, value: impl Serialize) -> Observation {
Observation {
label,
value: serde_json::to_value(value).expect("example observations are serializable"),
}
}
fn quiet_options() -> RunOptions {
RunOptions {
mirror: false,
..RunOptions::default()
}
}
async fn quiet(command: &str) -> command_stream::Result<CommandResult> {
exec(command, quiet_options()).await
}
async fn await_result() -> ExampleResult {
let result = quiet("echo hello").await?;
Ok(vec![
observation("stdout", result.stdout),
observation("stderr", result.stderr),
observation("exit code", result.code),
])
}
async fn result_text() -> ExampleResult {
let result = quiet("echo hello").await?;
Ok(vec![observation("text output", result.stdout)])
}
async fn sync_execution() -> ExampleResult {
let result = tokio::task::spawn_blocking(|| run_sync("echo synchronous")).await??;
Ok(vec![observation("stdout", result.stdout)])
}
async fn exit_codes() -> ExampleResult {
let result = quiet("false").await?;
let checked = result.clone().error_for_status().unwrap_err();
Ok(vec![
observation("result code", result.code),
observation("checked error code", checked.code()),
])
}
async fn options() -> ExampleResult {
let directory = tempfile::tempdir()?;
let mut env = HashMap::new();
env.insert(
"COMMAND_STREAM_DEMO".to_string(),
"from-options".to_string(),
);
let result = exec(
"cat",
RunOptions {
mirror: false,
cwd: Some(directory.path().to_path_buf()),
env: Some(env),
stdin: StdinOption::Content("from-stdin\n".to_string()),
..RunOptions::default()
},
)
.await?;
Ok(vec![observation("stdin and cwd options", result.stdout)])
}
async fn function_api() -> ExampleResult {
let simple = run("echo run").await?;
let configured = exec("echo exec", quiet_options()).await?;
let mut runner = create("echo create", quiet_options());
let created = runner.run().await?;
Ok(vec![observation(
"run, exec and create",
[
simple.stdout.trim(),
configured.stdout.trim(),
created.stdout.trim(),
],
)])
}
async fn cancellation() -> ExampleResult {
let child_command = if cfg!(windows) {
"ping -n 31 127.0.0.1"
} else {
"/bin/sleep 30"
};
let mut runner = ProcessRunner::new(child_command, quiet_options());
runner.start().await?;
let (child_available, child_pid_available) = match runner.child() {
Some(mut child) => {
let pid_available = child.pid().is_some();
child.kill_with("SIGTERM")?;
(true, pid_available)
}
None => (false, false),
};
let _ = runner.run().await?;
let mut stream = StreamingRunner::new("sleep 30").stream();
let started = stream.wait_for_pid().await.is_some();
stream.kill();
let mut exit_code = 0;
while let Some(chunk) = stream.next().await {
if let OutputChunk::Exit(code) = chunk {
exit_code = code;
}
}
Ok(vec![
observation("child handle available after start", child_available),
observation("child pid available", child_pid_available),
observation("stream started", started),
observation("cancelled stream exit is non-zero", exit_code != 0),
])
}
async fn async_iteration() -> ExampleResult {
let mut stream = StreamingRunner::new("printf 'one\\ntwo\\n'").stream();
let mut stdout = Vec::new();
let mut exit_code = None;
while let Some(chunk) = stream.next().await {
match chunk {
OutputChunk::Stdout(data) => stdout.extend(data),
OutputChunk::Stderr(_) => {}
OutputChunk::Exit(code) => exit_code = Some(code),
}
}
Ok(vec![
observation("collected chunks", String::from_utf8(stdout)?),
observation("exit code", exit_code),
])
}
async fn events() -> ExampleResult {
let emitter = StreamEmitter::new();
let count = Arc::new(AtomicUsize::new(0));
let listener_count = Arc::clone(&count);
emitter
.on(EventType::Stdout, move |_| {
listener_count.fetch_add(1, Ordering::SeqCst);
})
.await;
emitter
.emit(EventType::Stdout, EventData::String("hello".to_string()))
.await;
Ok(vec![observation(
"stdout events",
count.load(Ordering::SeqCst),
)])
}
async fn stdin_streaming() -> ExampleResult {
let mut runner = ProcessRunner::new(
"cat",
RunOptions {
mirror: false,
stdin: StdinOption::Pipe,
..RunOptions::default()
},
);
runner.start().await?;
runner.write_stdin("first line\n").await?;
runner.write_stdin("second line\n").await?;
runner.close_stdin().await?;
let result = runner.run().await?;
Ok(vec![observation("what cat echoed back", result.stdout)])
}
async fn buffers_strings() -> ExampleResult {
let result = quiet("printf bytes").await?;
Ok(vec![
observation("string", &result.stdout),
observation("bytes", result.stdout.as_bytes()),
])
}
async fn mirror_capture() -> ExampleResult {
let captured = quiet("echo captured").await?;
let uncaptured = exec(
"true",
RunOptions {
mirror: false,
capture: false,
..RunOptions::default()
},
)
.await?;
Ok(vec![
observation("captured output", captured.stdout),
observation("capture can be disabled", uncaptured.stdout.is_empty()),
])
}
async fn builtin_catalog() -> ExampleResult {
let registry = VirtualCommandRegistry::with_builtins();
let mut commands = registry.list();
commands.sort_unstable();
Ok(vec![
observation("available built-ins", &commands),
observation("number of built-ins", commands.len()),
])
}
async fn builtin_filesystem() -> ExampleResult {
let directory = tempfile::tempdir()?;
let options = RunOptions {
mirror: false,
cwd: Some(directory.path().to_path_buf()),
..RunOptions::default()
};
exec("mkdir demo", options.clone()).await?;
exec("touch demo/file.txt", options.clone()).await?;
let listed = exec("ls demo", options.clone()).await?;
exec("rm -r demo", options).await?;
Ok(vec![observation("created and listed", listed.stdout)])
}
async fn builtin_text() -> ExampleResult {
let sequence = quiet("seq 1 3").await?;
let basename = quiet("basename /tmp/example.txt").await?;
Ok(vec![
observation("sequence", sequence.stdout),
observation("basename", basename.stdout),
])
}
async fn builtin_environment() -> ExampleResult {
let mut env = HashMap::new();
env.insert("COMMAND_STREAM_DEMO".to_string(), "visible".to_string());
let result = exec(
"env",
RunOptions {
mirror: false,
env: Some(env),
..RunOptions::default()
},
)
.await?;
Ok(vec![observation(
"configured environment visible",
result.stdout.contains("COMMAND_STREAM_DEMO=visible"),
)])
}
fn greet_handler(
context: CommandContext,
) -> Pin<Box<dyn std::future::Future<Output = CommandResult> + Send>> {
Box::pin(async move { CommandResult::success(format!("Hello, {}!\n", context.args.join(" "))) })
}
async fn virtual_commands() -> ExampleResult {
let mut registry = VirtualCommandRegistry::new();
registry.register("greet", greet_handler);
let handler = registry.get("greet").expect("registered handler");
let result = handler(CommandContext::new(vec!["Rust".to_string()])).await;
let removed = registry.unregister("greet");
Ok(vec![
observation("custom command output", result.stdout),
observation("unregistered again", removed),
])
}
async fn virtual_context() -> ExampleResult {
let mut context = CommandContext::new(vec!["one".to_string(), "two".to_string()]);
context.stdin = Some("piped\n".to_string());
context.cwd = Some(std::env::temp_dir());
context.env = Some(HashMap::from([("DEMO".to_string(), "value".to_string())]));
Ok(vec![observation(
"handler context",
json!({
"args": context.args,
"stdin": context.stdin,
"has_cwd": context.cwd.is_some(),
"env_value": context.env.and_then(|env| env.get("DEMO").cloned()),
}),
)])
}
fn streaming_handler(
context: CommandContext,
) -> Pin<Box<dyn std::future::Future<Output = CommandResult> + Send>> {
Box::pin(async move {
if let Some(output) = context.output_tx {
let _ = output
.send(command_stream::StreamChunk::Stdout("one\n".to_string()))
.await;
let _ = output
.send(command_stream::StreamChunk::Stdout("two\n".to_string()))
.await;
}
CommandResult::success("one\ntwo\n")
})
}
async fn virtual_streaming() -> ExampleResult {
let (sender, mut receiver) = tokio::sync::mpsc::channel(4);
let mut context = CommandContext::new(Vec::new());
context.output_tx = Some(sender);
let result = streaming_handler(context).await;
let mut chunks = Vec::new();
while let Ok(chunk) = receiver.try_recv() {
if let command_stream::StreamChunk::Stdout(text) = chunk {
chunks.push(text);
}
}
Ok(vec![
observation("chunks", chunks),
observation("collected output", result.stdout),
])
}
async fn pipelines() -> ExampleResult {
let result = Pipeline::new()
.add("printf 'hello\\nworld\\n'")
.add("grep world")
.mirror_output(false)
.run()
.await?;
Ok(vec![observation("pipeline output", result.stdout)])
}
async fn redirection() -> ExampleResult {
let directory = tempfile::tempdir()?;
let file = directory.path().join("output.txt");
let result = exec(
"echo redirected > output.txt",
RunOptions {
mirror: false,
cwd: Some(directory.path().to_path_buf()),
..RunOptions::default()
},
)
.await?;
Ok(vec![
observation("exit code", result.code),
observation("file contents", std::fs::read_to_string(file)?),
])
}
async fn sequences() -> ExampleResult {
let result = quiet("false || echo fallback; echo next").await?;
Ok(vec![observation("sequence output", result.stdout)])
}
async fn interpolation() -> ExampleResult {
let value = "hello from Rust";
let result = cmd!("echo {}", value).await?;
Ok(vec![observation("macro interpolation", result.stdout)])
}
async fn shell_settings() -> ExampleResult {
set_shell_option("pipefail").await;
let with_pipefail = Pipeline::new().add("false").add("true").run().await?;
unset_shell_option("pipefail").await;
let without_pipefail = Pipeline::new().add("false").add("true").run().await?;
Ok(vec![
observation("with pipefail", with_pipefail.code),
observation("without pipefail", without_pipefail.code),
])
}
async fn ansi_utils() -> ExampleResult {
Ok(vec![observation(
"stripped output",
AnsiUtils::strip_all("\u{1b}[31mred\u{1b}[0m"),
)])
}
async fn zx_compat() -> ExampleResult {
let words = "hello world";
let greeting = zx!("echo {}", words).await?;
let rejected = zx!("exit 2").await.unwrap_err();
let tolerated = zx!(zx::Shell::new().nothrow(true), "exit 3").await?;
let sorted = zx!("printf 'b\\na\\n'").pipe(zx!("sort")).await?;
let dir = std::env::temp_dir().canonicalize()?;
let inside = zx::within(async {
zx::configure(|options| options.cwd = Some(dir.clone()));
zx!("pwd").await
})
.await?;
Ok(vec![
observation("interpolation", greeting.stdout),
observation("rejected exit code", rejected.exit_code),
observation("nothrow exit code", tolerated.exit_code),
observation("pipe", sorted.lines()),
observation("within cwd", inside.stdout.trim() == dir.to_string_lossy()),
observation("cwd restored", zx::current_options().cwd.is_none()),
])
}
async fn execa_compat() -> ExampleResult {
use command_stream::execa::{Execa, Options};
let isolated = command_stream::execa::execa("node", ["--version"]).await?;
let general = command_stream::execa("node", ["--version"]).await?;
let argv = command_stream::execa(
"node",
[
"-e",
"process.stdout.write(process.argv.slice(1).join('|'))",
"hello world",
"$HOME",
],
)
.await?;
let api = Execa::new(Options {
strip_final_newline: false,
..Options::default()
});
let newline = api.command("node", ["-e", "console.log('hello')"]).await?;
let tolerated = api
.command("node", ["-e", "process.exit(3)"])
.reject(false)
.await?;
let input = command_stream::execa("node", ["-e", "process.stdin.pipe(process.stdout)"])
.input(b"input".to_vec())
.await?;
Ok(vec![
observation(
"isolated and general API",
isolated.stdout == general.stdout,
),
observation("exact argv", argv.text()),
observation("preserved newline", newline.text()),
observation("tolerated exit code", tolerated.exit_code),
observation("binary stdin", input.text()),
])
}
async fn shelljs_compat() -> ExampleResult {
let directory = tempfile::tempdir()?;
std::fs::write(directory.path().join("file with spaces"), "z\na\na\nb\n")?;
let mut shell = command_stream::shelljs::ShellJs::new();
shell.config.silent = true;
shell
.cd(&[directory.path().to_str().ok_or("non-UTF8 path")?])
.await?;
Ok(vec![
observation("general API", true),
observation(
"separate arguments",
shell.echo(&["hello", "two words"]).await?.stdout,
),
observation(
"head",
shell.head(&["-n", "2", "file with spaces"]).await?.stdout,
),
observation(
"tail",
shell.tail(&["-n", "2", "file with spaces"]).await?.stdout,
),
observation("missing file code", shell.cat(&["missing"]).await?.code),
])
}
async fn native_text() -> ExampleResult {
let options = RunOptions {
stdin: StdinOption::Content("b\nb\na\n".into()),
..quiet_options()
};
let mut observations = Vec::new();
for (label, command) in [
("head", "head -n 2"),
("tail", "tail -n 1"),
("sort", "sort -u"),
("uniq", "uniq -cd"),
("zero lines", "tail -n 0"),
] {
observations.push(observation(
label,
ProcessRunner::new(command, options.clone())
.run()
.await?
.stdout,
));
}
Ok(observations)
}
async fn execute(id: &str) -> ExampleResult {
match id {
"await-result" => await_result().await,
"result-text" => result_text().await,
"sync-execution" => sync_execution().await,
"exit-codes" => exit_codes().await,
"options" => options().await,
"function-api" => function_api().await,
"cancellation" => cancellation().await,
"async-iteration" => async_iteration().await,
"events" => events().await,
"stdin-streaming" => stdin_streaming().await,
"buffers-strings" => buffers_strings().await,
"mirror-capture" => mirror_capture().await,
"builtin-catalog" => builtin_catalog().await,
"builtin-filesystem" => builtin_filesystem().await,
"builtin-text" => builtin_text().await,
"builtin-environment" => builtin_environment().await,
"virtual-commands" => virtual_commands().await,
"virtual-context" => virtual_context().await,
"virtual-streaming" => virtual_streaming().await,
"pipelines" => pipelines().await,
"redirection" => redirection().await,
"sequences" => sequences().await,
"interpolation" => interpolation().await,
"shell-settings" => shell_settings().await,
"ansi-utils" => ansi_utils().await,
"zx-compat" => zx_compat().await,
"execa-compat" => execa_compat().await,
"shelljs-compat" => shelljs_compat().await,
"native-text" => native_text().await,
_ => Err(format!("unknown feature: {id}").into()),
}
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn Error>> {
let id = std::env::args().nth(1).ok_or("pass a feature id")?;
let observations = execute(&id).await?;
println!("# {id} — Rust");
for item in &observations {
println!("{}: {}", item.label, item.value);
}
println!("{PARITY_START}");
println!(
"{}",
serde_json::to_string(&json!({
"id": id,
"language": "rust",
"observations": observations,
"failure": Value::Null,
}))?
);
println!("{PARITY_END}");
Ok(())
}