From 4018ec6da14d78d141d8edd2d00ef0f22c0c0c9e Mon Sep 17 00:00:00 2001 From: Dave Nagoda Date: Thu, 1 Oct 2026 12:00:32 -0700 Subject: [PATCH 1/5] Add batch mode for JSON Lines inputs `--batch` runs one Function against many inputs in one process. It reads JSON Lines from `--input` or stdin and writes one JSON record per input line to stdout. The Function module, its provider, and the optional schema and query are loaded once, so each input pays only for the run. - Each record has the 1-based input `line`, so records match inputs even when blank lines are skipped. - Records are written with serde, so errors with quotes or newlines stay valid JSON Lines. - By default the batch stops at the first failed input. `--batch-continue-on-error` runs all inputs. The exit code is non-zero if any input failed. - `--batch-full-output` writes the full run result for each input. - A summary goes to stderr. `BluejaySchemaAnalyzer::with_analyzer` parses the schema and query once and computes the scale factor for each input. --- README.md | 66 +++++++ src/batch.rs | 244 +++++++++++++++++++++++ src/bluejay_schema_analyzer.rs | 64 ++++-- src/main.rs | 51 ++++- tests/integration_tests.rs | 344 +++++++++++++++++++++++++++++++++ 5 files changed, 753 insertions(+), 16 deletions(-) create mode 100644 src/batch.rs diff --git a/README.md b/README.md index 211eefca..7878be07 100644 --- a/README.md +++ b/README.md @@ -23,6 +23,72 @@ contains binaries that can be run on your computer. To see the list of possible commands and arguments, run `function-runner --help`. +### Batch mode + +Use `--batch` to run a Function against many inputs in one process. The +Function is loaded and compiled once, so each input costs only the run itself. + +The input is [JSON Lines](https://jsonlines.org/): one JSON input per line, +from `--input` or stdin. Blank lines are skipped. + +```sh +function-runner -f function.wasm --batch -i inputs.jsonl > results.jsonl +``` + +For each input, batch mode writes one JSON record on one line to stdout. +`line` is the 1-based line number of the input, so you can match each record +to its input. + +- An input that ran: + `{"line":1,"success":true,"instructions":5069,"memory_usage":1088,"logs":"","output":{...}}`. + `success` is `false` if the Function failed. If the output is not valid + JSON, `output` is `null` and `output_error` gives the reason. +- An input that could not run, for example invalid JSON: + `{"line":2,"success":false,"error":"Invalid input JSON: ..."}`. If the + Function itself cannot run, for example because it imports both WASI and a + provider that does not allow WASI, each input gets an error record with that + reason. + +Each record is one complete JSON line. If a result cannot be serialized, the +record for that input is an error record, never partial output. + +A summary goes to stderr, for example +`Batch complete: 3 inputs processed, 2 successful, 1 failed`. + +Batch options: + +- `--batch-continue-on-error`: run all inputs even if some fail. Without it, + the batch stops after the first failed input. +- `--batch-full-output`: write the full result for each input, the same fields + as `--json` plus `line`. + +The exit code is `0` only if every input succeeds. `--schema-path` and +`--query-path` work in batch mode; the schema and query are parsed once. +Profiling is not available in batch mode. + +## Library usage + +To compute scale factors for many inputs, use +`bluejay_schema_analyzer::BluejaySchemaAnalyzer::with_analyzer`. It parses and +validates the schema and query once, then calls your closure with an `analyze` +function that returns the scale factor for one input: + +```rust +use function_runner::bluejay_schema_analyzer::BluejaySchemaAnalyzer; + +let scale_factors = BluejaySchemaAnalyzer::with_analyzer( + &schema, + Some("schema.graphql"), + &query, + Some("input.graphql"), + |analyze| inputs.iter().map(|input| analyze(input)).collect::>>(), +)??; +``` + +The outer `Result` holds schema and query errors. The inner `Result` holds +analysis errors for an input. The `test_with_analyzer_analyzes_many_inputs` +test in `src/bluejay_schema_analyzer.rs` runs this pattern. + ## Development Building requires a rust toolchain of `1.66.0` to `1.67.0`. `cargo install --path . --locked` will build diff --git a/src/batch.rs b/src/batch.rs new file mode 100644 index 00000000..649a9c5b --- /dev/null +++ b/src/batch.rs @@ -0,0 +1,244 @@ +use std::io::{BufRead, Write}; +use std::path::Path; + +use anyhow::{anyhow, bail, Result}; +use function_runner::{ + bluejay_schema_analyzer::BluejaySchemaAnalyzer, + engine::{run, FunctionRunParams}, + function_run_result::FunctionRunResult, + BytesContainer, BytesContainerType, Codec, +}; +use serde::Serialize; +use wasmtime::Module; + +type ScaleFactorFn<'a> = dyn Fn(&serde_json::Value) -> Result + 'a; + +pub(crate) struct ScaleLimitsSource<'a> { + pub schema: &'a str, + pub schema_path: Option<&'a str>, + pub query: &'a str, + pub query_path: Option<&'a str>, +} + +pub(crate) struct BatchOptions<'a> { + pub function_path: &'a Path, + pub export: &'a str, + pub codec: Codec, + pub default_scale_factor: f64, + pub continue_on_error: bool, + pub full_output: bool, +} + +#[derive(Serialize)] +struct MinimalRecord<'a> { + line: usize, + success: bool, + instructions: u64, + memory_usage: u64, + logs: &'a str, + output: Option<&'a serde_json::Value>, + #[serde(skip_serializing_if = "Option::is_none")] + output_error: Option<&'a str>, +} + +#[derive(Serialize)] +struct FullRecord<'a> { + line: usize, + #[serde(flatten)] + result: &'a FunctionRunResult, + #[serde(skip_serializing_if = "Option::is_none")] + output_error: Option<&'a str>, +} + +#[derive(Serialize)] +struct ErrorRecord<'a> { + line: usize, + success: bool, + error: &'a str, +} + +#[derive(Default)] +struct Summary { + processed: usize, + successful: usize, + failed: usize, +} + +pub(crate) fn run_batch( + input: impl BufRead, + output: impl Write, + module: &Module, + scale_limits: Option, + options: &BatchOptions, +) -> Result<()> { + match scale_limits { + Some(source) => BluejaySchemaAnalyzer::with_analyzer( + source.schema, + source.schema_path, + source.query, + source.query_path, + |analyze| run_lines(input, output, module, Some(analyze), options), + )?, + None => run_lines(input, output, module, None, options), + } +} + +fn run_lines( + mut input: impl BufRead, + mut output: impl Write, + module: &Module, + analyze: Option<&ScaleFactorFn<'_>>, + options: &BatchOptions, +) -> Result<()> { + let mut summary = Summary::default(); + let mut line_bytes = Vec::new(); + let mut record = Vec::new(); + let mut line = 0; + + let outcome = loop { + line_bytes.clear(); + match input.read_until(b'\n', &mut line_bytes) { + Ok(0) => break Ok(()), + Ok(_) => {} + Err(e) => break Err(anyhow!("Couldn't read input line {}: {}", line + 1, e)), + } + line += 1; + + if line_bytes.iter().all(u8::is_ascii_whitespace) { + continue; + } + summary.processed += 1; + + let failure = match run_line(std::mem::take(&mut line_bytes), module, analyze, options) { + Ok(result) => match write_result_record(&mut record, line, &result, options.full_output) { + Err(error) => Some(format!("Line {line}: {error}")), + Ok(()) if result.success => None, + Ok(()) => Some(format!( + "The Function execution failed on line {line}. Review the logs for more information." + )), + }, + Err(error) => { + let error = format!("{error:#}"); + write_error_record(&mut record, line, &error); + Some(format!("Line {line}: {error}")) + } + }; + + if let Err(e) = output.write_all(&record) { + break Err(e.into()); + } + + match failure { + None => summary.successful += 1, + Some(reason) => { + summary.failed += 1; + if !options.continue_on_error { + break Err(anyhow!(reason)); + } + } + } + }; + + output.flush()?; + + let status = if outcome.is_ok() { + "complete" + } else { + "stopped" + }; + eprintln!( + "Batch {status}: {} inputs processed, {} successful, {} failed", + summary.processed, summary.successful, summary.failed + ); + + outcome?; + + if summary.failed > 0 { + bail!("{} of {} inputs failed", summary.failed, summary.processed); + } + + Ok(()) +} + +fn run_line( + line_bytes: Vec, + module: &Module, + analyze: Option<&ScaleFactorFn<'_>>, + options: &BatchOptions, +) -> Result { + let input = BytesContainer::new(BytesContainerType::Input, options.codec, line_bytes)?; + + let scale_factor = match (analyze, input.json_value.as_ref()) { + (Some(analyze), Some(json_value)) => analyze(json_value)?, + _ => options.default_scale_factor, + }; + + run(FunctionRunParams { + function_path: options.function_path.to_path_buf(), + input, + export: options.export, + profile_opts: None, + scale_factor, + module: module.clone(), + engine: module.engine().clone(), + }) +} + +fn write_result_record( + record: &mut Vec, + line: usize, + result: &FunctionRunResult, + full_output: bool, +) -> Result<(), String> { + record.clear(); + let output_error = result.output.encoding_error.as_deref(); + let serialized = if full_output { + serde_json::to_writer( + &mut *record, + &FullRecord { + line, + result, + output_error, + }, + ) + } else { + serde_json::to_writer( + &mut *record, + &MinimalRecord { + line, + success: result.success, + instructions: result.instructions, + memory_usage: result.memory_usage, + logs: &result.logs, + output: result.output.json_value.as_ref(), + output_error, + }, + ) + }; + + match serialized { + Ok(()) => { + record.push(b'\n'); + Ok(()) + } + Err(e) => { + let error = format!("Couldn't serialize result: {e}"); + write_error_record(record, line, &error); + Err(error) + } + } +} + +fn write_error_record(record: &mut Vec, line: usize, error: &str) { + record.clear(); + serde_json::to_writer( + &mut *record, + &ErrorRecord { + line, + success: false, + error, + }, + ) + .expect("An error record contains only a number, a bool, and a string"); + record.push(b'\n'); +} diff --git a/src/bluejay_schema_analyzer.rs b/src/bluejay_schema_analyzer.rs index 548dda96..d49f901d 100644 --- a/src/bluejay_schema_analyzer.rs +++ b/src/bluejay_schema_analyzer.rs @@ -19,6 +19,18 @@ impl BluejaySchemaAnalyzer { query_path: Option<&str>, input: &serde_json::Value, ) -> Result { + Self::with_analyzer(schema_string, schema_path, query, query_path, |analyze| { + analyze(input) + })? + } + + pub fn with_analyzer( + schema_string: &str, + schema_path: Option<&str>, + query: &str, + query_path: Option<&str>, + f: impl FnOnce(&dyn Fn(&serde_json::Value) -> Result) -> R, + ) -> Result { let document_definition = DefinitionDocument::parse(schema_string) .result .map_err(|errors| anyhow!(Error::format_errors(schema_string, schema_path, errors)))?; @@ -30,18 +42,21 @@ impl BluejaySchemaAnalyzer { .result .map_err(|errors| anyhow!(Error::format_errors(query, query_path, errors)))?; - let cache = - bluejay_validator::executable::Cache::new(&executable_document, &schema_definition); + let analyze = |input: &serde_json::Value| { + let cache = + bluejay_validator::executable::Cache::new(&executable_document, &schema_definition); + ScaleLimitsAnalyzer::analyze( + &executable_document, + &schema_definition, + None, + &Default::default(), + &cache, + input, + ) + .map_err(|e| anyhow!("Unable to analyze scale limits: {}", e.message())) + }; - ScaleLimitsAnalyzer::analyze( - &executable_document, - &schema_definition, - None, - &Default::default(), - &cache, - input, - ) - .map_err(|e| anyhow!("Unable to analyze scale limits: {}", e.message())) + Ok(f(&analyze)) } } @@ -84,6 +99,33 @@ mod tests { ); } + #[test] + fn test_with_analyzer_analyzes_many_inputs() -> Result<()> { + let schema_string = r#" + directive @scaleLimits(rate: Float!) on FIELD_DEFINITION + type Query { + cartLines: [String] @scaleLimits(rate: 0.005) + } + "#; + let query = "{ cartLines }"; + + let scale_factors = BluejaySchemaAnalyzer::with_analyzer( + schema_string, + None, + query, + None, + |analyze| -> Result> { + [100, 500, 1000] + .iter() + .map(|length| analyze(&json!({ "cartLines": vec!["line"; *length] }))) + .collect() + }, + )??; + + assert_eq!(scale_factors, vec![1.0, 2.5, 5.0]); + Ok(()) + } + #[test] fn test_analyze_schema_with_array_length_scaling() { let schema_string = r#" diff --git a/src/main.rs b/src/main.rs index cd5f41a1..8bdf8fc6 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,9 +1,11 @@ +mod batch; + use function_runner::{BytesContainer, BytesContainerType, Codec}; use wasmtime::Module; use std::{ fs::File, - io::{stdin, BufReader, Read}, + io::{stdin, stdout, BufRead, BufReader, Read}, path::PathBuf, }; @@ -32,6 +34,19 @@ struct Opts { #[clap(short, long)] input: Option, + /// Run the Function once for each line of a JSON Lines input (one JSON input per line). + /// Writes one JSON result per line to stdout and a summary to stderr. + #[clap(long, conflicts_with_all = ["profile", "profile_out", "profile_frequency"])] + batch: bool, + + /// With --batch, run all inputs even if some fail. By default, the batch stops at the first failure. + #[clap(long, requires = "batch")] + batch_continue_on_error: bool, + + /// With --batch, write the full result (including the input) for each input. + #[clap(long, requires = "batch")] + batch_full_output: bool, + /// Name of the export to invoke. #[clap(short, long, default_value = "_start")] export: String, @@ -114,7 +129,7 @@ fn read_file_to_string(file_path: &PathBuf) -> Result { fn main() -> Result<()> { let opts: Opts = Opts::parse(); - let mut input: Box = if let Some(ref input) = opts.input { + let mut input: Box = if let Some(ref input) = opts.input { Box::new(BufReader::new(File::open(input).map_err(|e| { anyhow!("Couldn't load input {:?}: {}", input, e) })?)) @@ -126,9 +141,6 @@ fn main() -> Result<()> { )); }; - let mut buffer = Vec::new(); - input.read_to_end(&mut buffer)?; - let schema_string = opts.read_schema_to_string().transpose()?; let query_string = opts.read_query_to_string().transpose()?; @@ -144,6 +156,35 @@ fn main() -> Result<()> { Codec::Json }; + if opts.batch { + let scale_limits = match (&schema_string, &query_string) { + (Some(schema), Some(query)) => Some(batch::ScaleLimitsSource { + schema, + schema_path: opts.schema_path.as_ref().and_then(|p| p.to_str()), + query, + query_path: opts.query_path.as_ref().and_then(|p| p.to_str()), + }), + _ => None, + }; + return batch::run_batch( + input, + stdout().lock(), + &module, + scale_limits, + &batch::BatchOptions { + function_path: &opts.function, + export: &opts.export, + codec, + default_scale_factor: DEFAULT_SCALE_FACTOR, + continue_on_error: opts.batch_continue_on_error, + full_output: opts.batch_full_output, + }, + ); + } + + let mut buffer = Vec::new(); + input.read_to_end(&mut buffer)?; + let input = BytesContainer::new(BytesContainerType::Input, codec, buffer)?; let scale_factor = if let (Some(schema_string), Some(query_string), Some(json_value)) = (schema_string, query_string, input.json_value.clone()) diff --git a/tests/integration_tests.rs b/tests/integration_tests.rs index 5baa864a..21aa9233 100644 --- a/tests/integration_tests.rs +++ b/tests/integration_tests.rs @@ -500,4 +500,348 @@ mod tests { Ok(()) } + + fn temp_batch_input(contents: &str) -> Result { + let file = assert_fs::NamedTempFile::new("input.jsonl")?; + file.write_str(contents)?; + + Ok(file) + } + + fn batch_records(stdout: &[u8]) -> Result> { + Ok(String::from_utf8(stdout.to_vec())? + .lines() + .map(serde_json::from_str::) + .collect::, _>>()?) + } + + #[test] + fn batch_writes_one_minimal_record_per_input() -> Result<()> { + let input_file = temp_batch_input("{\"code\":0}\n{\"code\":0}\n")?; + + let output = Command::new(cargo_bin!()) + .args([ + "--function", + "tests/fixtures/build/exit_code.wasm", + "--batch", + ]) + .arg("--input") + .arg(input_file.as_os_str()) + .output()?; + + assert!(output.status.success()); + let records = batch_records(&output.stdout)?; + assert_eq!(records.len(), 2); + assert_eq!( + records[0], + json!({ + "line": 1, + "success": true, + "instructions": records[0]["instructions"], + "memory_usage": records[0]["memory_usage"], + "logs": "", + "output": {"exit": 0}, + }) + ); + assert!(records[0]["instructions"].as_u64().unwrap() > 0); + assert_eq!(records[1]["line"], 2); + assert_eq!( + String::from_utf8(output.stderr)?, + "Batch complete: 2 inputs processed, 2 successful, 0 failed\n" + ); + + Ok(()) + } + + #[test] + fn batch_reads_stdin_and_skips_blank_lines() -> Result<()> { + let input_file = temp_batch_input("{\"code\":0}\n\n \n{\"code\":0}")?; + + let output = Command::new(cargo_bin!()) + .args([ + "--function", + "tests/fixtures/build/exit_code.wasm", + "--batch", + ]) + .stdin(Stdio::from(File::open(input_file.path())?)) + .output()?; + + assert!(output.status.success()); + let records = batch_records(&output.stdout)?; + let lines: Vec<_> = records.iter().map(|r| r["line"].clone()).collect(); + assert_eq!(lines, vec![json!(1), json!(4)]); + + Ok(()) + } + + #[test] + fn batch_stops_at_first_failure_by_default() -> Result<()> { + let input_file = temp_batch_input("{\"code\":0}\n{\"code\":1}\n{\"code\":0}\n")?; + + let output = Command::new(cargo_bin!()) + .args([ + "--function", + "tests/fixtures/build/exit_code.wasm", + "--batch", + ]) + .arg("--input") + .arg(input_file.as_os_str()) + .output()?; + + assert!(!output.status.success()); + let records = batch_records(&output.stdout)?; + assert_eq!(records.len(), 2); + assert_eq!(records[1]["line"], 2); + assert_eq!(records[1]["success"], false); + assert_eq!(records[1]["logs"], "module exited with code: 1"); + assert_eq!( + String::from_utf8(output.stderr)?, + "Batch stopped: 2 inputs processed, 1 successful, 1 failed\n\ + Error: The Function execution failed on line 2. Review the logs for more information.\n" + ); + + Ok(()) + } + + #[test] + fn batch_reports_an_invalid_function_in_the_first_record() -> Result<()> { + let input_file = temp_batch_input("{}\n{}\n")?; + + let output = Command::new(cargo_bin!()) + .args([ + "--function", + "tests/fixtures/build/invalid_import_combination.wasm", + "--batch", + ]) + .arg("--input") + .arg(input_file.as_os_str()) + .output()?; + + let error = "Invalid Function, cannot use `shopify_function_v2` and import WASI. If using Rust, change the build target to `wasm32-unknown-unknown`."; + assert!(!output.status.success()); + assert_eq!( + batch_records(&output.stdout)?, + vec![json!({"line": 1, "success": false, "error": error})] + ); + assert_eq!( + String::from_utf8(output.stderr)?, + format!( + "Batch stopped: 1 inputs processed, 0 successful, 1 failed\nError: Line 1: {error}\n" + ) + ); + + Ok(()) + } + + #[test] + fn batch_continue_on_error_runs_all_inputs_and_fails_at_the_end() -> Result<()> { + let input_file = + temp_batch_input("{\"code\":0}\n{\"code\":1}\n{\"code\":\n{\"code\":0}\n")?; + + let output = Command::new(cargo_bin!()) + .args([ + "--function", + "tests/fixtures/build/exit_code.wasm", + "--batch", + ]) + .arg("--batch-continue-on-error") + .arg("--input") + .arg(input_file.as_os_str()) + .output()?; + + assert!(!output.status.success()); + let records = batch_records(&output.stdout)?; + assert_eq!(records.len(), 4); + let success: Vec<_> = records.iter().map(|r| r["success"].clone()).collect(); + assert_eq!( + success, + vec![json!(true), json!(false), json!(false), json!(true)] + ); + assert_eq!( + records[2], + json!({ + "line": 3, + "success": false, + "error": "Invalid input JSON: EOF while parsing a value at line 2 column 0", + }) + ); + assert_eq!( + String::from_utf8(output.stderr)?, + "Batch complete: 4 inputs processed, 2 successful, 2 failed\n\ + Error: 2 of 4 inputs failed\n" + ); + + Ok(()) + } + + #[test] + fn batch_error_records_are_single_line_json() -> Result<()> { + let input_file = temp_batch_input("{\"code\":0}\n")?; + + let output = Command::new(cargo_bin!()) + .args([ + "--function", + "tests/fixtures/build/exit_code.wasm", + "--batch", + ]) + .args(["--export", "missing \"export\"\nname"]) + .arg("--input") + .arg(input_file.as_os_str()) + .output()?; + + assert!(!output.status.success()); + let records = batch_records(&output.stdout)?; + assert_eq!(records.len(), 1); + assert_eq!(records[0]["success"], false); + assert!(records[0]["error"] + .as_str() + .unwrap() + .contains("missing \"export\"\nname")); + + Ok(()) + } + + #[test] + fn batch_full_output_includes_the_full_result() -> Result<()> { + let input_file = temp_batch_input("{\"code\":0}\n")?; + + let output = Command::new(cargo_bin!()) + .args([ + "--function", + "tests/fixtures/build/exit_code.wasm", + "--batch", + ]) + .arg("--batch-full-output") + .arg("--input") + .arg(input_file.as_os_str()) + .output()?; + + assert!(output.status.success()); + let records = batch_records(&output.stdout)?; + assert_eq!(records.len(), 1); + assert_eq!(records[0]["line"], 1); + assert_eq!(records[0]["name"], "exit_code.wasm"); + assert_eq!(records[0]["input"], json!({"code": 0})); + assert_eq!(records[0]["output"], json!({"exit": 0})); + assert_eq!(records[0]["success"], true); + assert!(records[0]["size"].is_u64()); + + Ok(()) + } + + #[test] + fn batch_reports_output_that_is_not_valid_json() -> Result<()> { + let input_file = temp_batch_input("{\"code\":0}\n")?; + + let output = Command::new(cargo_bin!()) + .args([ + "--function", + "tests/fixtures/build/messagepack-invalid.wasm", + ]) + .arg("--batch") + .arg("--input") + .arg(input_file.as_os_str()) + .output()?; + + let records = batch_records(&output.stdout)?; + assert_eq!(records.len(), 1); + assert_eq!(records[0]["output"], serde_json::Value::Null); + assert!(records[0]["output_error"].is_string()); + + Ok(()) + } + + #[test] + fn batch_runs_javy_plugin_functions() -> Result<()> { + let input_file = temp_batch_input("{\"hello\":\"world\"}\n{\"hello\":\"world\"}\n")?; + + let output = Command::new(cargo_bin!()) + .args([ + "--function", + "tests/fixtures/build/js_function_javy_plugin_v3.wasm", + "--batch", + ]) + .arg("--input") + .arg(input_file.as_os_str()) + .output()?; + + assert!(output.status.success()); + let records = batch_records(&output.stdout)?; + assert_eq!(records.len(), 2); + for record in records { + assert_eq!(record["output"], json!({"hello": "world output"})); + } + + Ok(()) + } + + #[test] + fn batch_uses_schema_and_query() -> Result<()> { + let input_file = temp_batch_input( + "{\"cart\":{\"lines\":[{\"quantity\":2}]}}\n{\"cart\":{\"lines\":[]}}\n", + )?; + + let output = Command::new(cargo_bin!()) + .args(["--function", "tests/fixtures/build/noop.wasm", "--batch"]) + .args(["--schema-path", "tests/fixtures/schema/schema.graphql"]) + .args(["--query-path", "tests/fixtures/query/query.graphql"]) + .arg("--input") + .arg(input_file.as_os_str()) + .output()?; + + assert!(output.status.success()); + assert_eq!(batch_records(&output.stdout)?.len(), 2); + + Ok(()) + } + + #[test] + fn batch_fails_before_running_when_schema_is_invalid() -> Result<()> { + let schema = assert_fs::NamedTempFile::new("schema.graphql")?; + schema.write_str("type Query {")?; + let input_file = temp_batch_input("{\"cart\":{\"lines\":[]}}\n")?; + + let output = Command::new(cargo_bin!()) + .args(["--function", "tests/fixtures/build/noop.wasm", "--batch"]) + .arg("--schema-path") + .arg(schema.as_os_str()) + .args(["--query-path", "tests/fixtures/query/query.graphql"]) + .arg("--input") + .arg(input_file.as_os_str()) + .output()?; + + assert!(!output.status.success()); + assert!(output.stdout.is_empty()); + + Ok(()) + } + + #[test] + fn batch_options_require_batch() -> Result<()> { + for flag in ["--batch-continue-on-error", "--batch-full-output"] { + Command::new(cargo_bin!()) + .args(["--function", "tests/fixtures/build/exit_code.wasm", flag]) + .assert() + .failure() + .stderr(contains("--batch")); + } + + Ok(()) + } + + #[test] + fn batch_cannot_be_used_with_profiling() -> Result<()> { + Command::new(cargo_bin!()) + .args([ + "--function", + "tests/fixtures/build/exit_code.wasm", + "--batch", + "--profile", + ]) + .assert() + .failure() + .stderr(contains("cannot be used with")); + + Ok(()) + } } From af9a64af87d8d3f44cfd0a915e422e611e3dffc5 Mon Sep 17 00:00:00 2001 From: Dave Nagoda Date: Fri, 2 Oct 2026 17:43:18 -0700 Subject: [PATCH 2/5] Read single-run input before loading the Function Batch mode moved the input read below schema loading and Function compilation for every run, so a run with both a bad input and a bad Function reported the Function error instead of the input error. Read the whole input first again when the run is not a batch. Only batch mode defers the read, because it streams the input line by line. --- src/main.rs | 8 +++++--- tests/integration_tests.rs | 14 ++++++++++++++ 2 files changed, 19 insertions(+), 3 deletions(-) diff --git a/src/main.rs b/src/main.rs index 8bdf8fc6..631cdc64 100644 --- a/src/main.rs +++ b/src/main.rs @@ -141,6 +141,11 @@ fn main() -> Result<()> { )); }; + let mut buffer = Vec::new(); + if !opts.batch { + input.read_to_end(&mut buffer)?; + } + let schema_string = opts.read_schema_to_string().transpose()?; let query_string = opts.read_query_to_string().transpose()?; @@ -182,9 +187,6 @@ fn main() -> Result<()> { ); } - let mut buffer = Vec::new(); - input.read_to_end(&mut buffer)?; - let input = BytesContainer::new(BytesContainerType::Input, codec, buffer)?; let scale_factor = if let (Some(schema_string), Some(query_string), Some(json_value)) = (schema_string, query_string, input.json_value.clone()) diff --git a/tests/integration_tests.rs b/tests/integration_tests.rs index 21aa9233..c807aff8 100644 --- a/tests/integration_tests.rs +++ b/tests/integration_tests.rs @@ -144,6 +144,20 @@ mod tests { Ok(()) } + #[test] + fn input_read_error_comes_before_function_load_error() -> Result<()> { + let output = Command::new(cargo_bin!()) + .args(["--input", ".", "--function", "test/file/doesnt/exist"]) + .output()?; + + assert!(!output.status.success()); + let stderr = String::from_utf8(output.stderr)?; + assert!(stderr.starts_with("Error: "), "{stderr}"); + assert!(!stderr.contains("Couldn't load the Function"), "{stderr}"); + + Ok(()) + } + #[test] fn profile_writes_file() -> Result<()> { let (mut cmd, temp) = profile_base_cmd_in_temp_dir()?; From 51687847b5fb9656e53f80dc8038146d910be86b Mon Sep 17 00:00:00 2001 From: Dave Nagoda Date: Fri, 2 Oct 2026 17:43:27 -0700 Subject: [PATCH 3/5] Write full batch records with the JSON values as they are Full batch records serialized FunctionRunResult, whose BytesContainer flattens its JSON value into the parent. That serializer rejects arrays, strings, and booleans, changes null to {}, and writes numbers as {"$serde_json::private::Number":"123"}. Output that is not valid JSON became {} instead of null. Write the full record fields directly, with input and output taken from their JSON values, the same as minimal records. Single runs keep the existing serializer. --- README.md | 4 ++- src/batch.rs | 19 ++++++++-- tests/integration_tests.rs | 71 ++++++++++++++++++++++++++++++++------ 3 files changed, 80 insertions(+), 14 deletions(-) diff --git a/README.md b/README.md index 7878be07..77177b51 100644 --- a/README.md +++ b/README.md @@ -60,7 +60,9 @@ Batch options: - `--batch-continue-on-error`: run all inputs even if some fail. Without it, the batch stops after the first failed input. - `--batch-full-output`: write the full result for each input, the same fields - as `--json` plus `line`. + as `--json` plus `line`. `input` and `output` hold the JSON values as they + are, so `output` is `null` when the output is not valid JSON, and + `output_error` gives the reason. The exit code is `0` only if every input succeeds. `--schema-path` and `--query-path` work in batch mode; the schema and query are parsed once. diff --git a/src/batch.rs b/src/batch.rs index 649a9c5b..51cbe5d0 100644 --- a/src/batch.rs +++ b/src/batch.rs @@ -44,8 +44,14 @@ struct MinimalRecord<'a> { #[derive(Serialize)] struct FullRecord<'a> { line: usize, - #[serde(flatten)] - result: &'a FunctionRunResult, + name: &'a str, + size: u64, + memory_usage: u64, + instructions: u64, + logs: &'a str, + input: Option<&'a serde_json::Value>, + output: Option<&'a serde_json::Value>, + success: bool, #[serde(skip_serializing_if = "Option::is_none")] output_error: Option<&'a str>, } @@ -197,7 +203,14 @@ fn write_result_record( &mut *record, &FullRecord { line, - result, + name: &result.name, + size: result.size, + memory_usage: result.memory_usage, + instructions: result.instructions, + logs: &result.logs, + input: result.input.json_value.as_ref(), + output: result.output.json_value.as_ref(), + success: result.success, output_error, }, ) diff --git a/tests/integration_tests.rs b/tests/integration_tests.rs index c807aff8..6dbc939e 100644 --- a/tests/integration_tests.rs +++ b/tests/integration_tests.rs @@ -743,24 +743,75 @@ mod tests { Ok(()) } + #[test] + fn batch_keeps_json_values_of_any_type() -> Result<()> { + let values = [ + json!({"code": 0}), + json!([]), + json!([1, 2]), + json!(null), + json!(123), + json!(1.5), + json!(true), + json!("test"), + ]; + let lines = values + .iter() + .map(|value| format!("{value}\n")) + .collect::(); + let input_file = temp_batch_input(&lines)?; + + for full_output in [false, true] { + let mut cmd = Command::new(cargo_bin!()); + cmd.args(["--function", "tests/fixtures/build/noop.wasm", "--batch"]) + .arg("--input") + .arg(input_file.as_os_str()); + if full_output { + cmd.arg("--batch-full-output"); + } + let output = cmd.output()?; + + assert!(output.status.success(), "full_output: {full_output}"); + let records = batch_records(&output.stdout)?; + assert_eq!(records.len(), values.len()); + for (record, value) in records.iter().zip(&values) { + assert_eq!(record["success"], true, "{record}"); + assert_eq!(&record["output"], value, "{record}"); + if full_output { + assert_eq!(&record["input"], value, "{record}"); + } + } + } + + Ok(()) + } + #[test] fn batch_reports_output_that_is_not_valid_json() -> Result<()> { let input_file = temp_batch_input("{\"code\":0}\n")?; - let output = Command::new(cargo_bin!()) - .args([ + for full_output in [false, true] { + let mut cmd = Command::new(cargo_bin!()); + cmd.args([ "--function", "tests/fixtures/build/messagepack-invalid.wasm", + "--batch", ]) - .arg("--batch") .arg("--input") - .arg(input_file.as_os_str()) - .output()?; - - let records = batch_records(&output.stdout)?; - assert_eq!(records.len(), 1); - assert_eq!(records[0]["output"], serde_json::Value::Null); - assert!(records[0]["output_error"].is_string()); + .arg(input_file.as_os_str()); + if full_output { + cmd.arg("--batch-full-output"); + } + let output = cmd.output()?; + + let records = batch_records(&output.stdout)?; + assert_eq!(records.len(), 1); + assert_eq!(records[0]["output"], serde_json::Value::Null); + assert!(records[0]["output_error"].is_string()); + if full_output { + assert_eq!(records[0]["input"], json!({"code": 0})); + } + } Ok(()) } From ac793186c4dca65c7b20978177ccde5d56f30d95 Mon Sep 17 00:00:00 2001 From: Dave Nagoda Date: Fri, 2 Oct 2026 17:43:27 -0700 Subject: [PATCH 4/5] Reject --json with --batch --batch accepted --json and ignored it, so a caller who asked for the full result silently got minimal records. Fail with a usage error that points to --batch-full-output instead. --- README.md | 3 +++ src/main.rs | 11 ++++++++++- tests/integration_tests.rs | 16 ++++++++++++++++ 3 files changed, 29 insertions(+), 1 deletion(-) diff --git a/README.md b/README.md index 77177b51..f0572536 100644 --- a/README.md +++ b/README.md @@ -64,6 +64,9 @@ Batch options: are, so `output` is `null` when the output is not valid JSON, and `output_error` gives the reason. +`--json` cannot be used with `--batch`: batch records are already JSON. Use +`--batch-full-output` for the full result. + The exit code is `0` only if every input succeeds. `--schema-path` and `--query-path` work in batch mode; the schema and query are parsed once. Profiling is not available in batch mode. diff --git a/src/main.rs b/src/main.rs index 631cdc64..97fa0775 100644 --- a/src/main.rs +++ b/src/main.rs @@ -10,7 +10,7 @@ use std::{ }; use anyhow::{anyhow, Result}; -use clap::Parser; +use clap::{error::ErrorKind, CommandFactory, Parser}; use function_runner::{ bluejay_schema_analyzer::BluejaySchemaAnalyzer, engine::{run, FunctionRunParams, ProfileOpts}, @@ -129,6 +129,15 @@ fn read_file_to_string(file_path: &PathBuf) -> Result { fn main() -> Result<()> { let opts: Opts = Opts::parse(); + if opts.batch && opts.json { + Opts::command() + .error( + ErrorKind::ArgumentConflict, + "--json cannot be used with --batch. Batch records are already JSON; use --batch-full-output for the full result.", + ) + .exit(); + } + let mut input: Box = if let Some(ref input) = opts.input { Box::new(BufReader::new(File::open(input).map_err(|e| { anyhow!("Couldn't load input {:?}: {}", input, e) diff --git a/tests/integration_tests.rs b/tests/integration_tests.rs index 6dbc939e..aeb632f6 100644 --- a/tests/integration_tests.rs +++ b/tests/integration_tests.rs @@ -894,6 +894,22 @@ mod tests { Ok(()) } + #[test] + fn batch_cannot_be_used_with_json() -> Result<()> { + Command::new(cargo_bin!()) + .args([ + "--function", + "tests/fixtures/build/exit_code.wasm", + "--batch", + "--json", + ]) + .assert() + .failure() + .stderr(contains("--batch-full-output")); + + Ok(()) + } + #[test] fn batch_cannot_be_used_with_profiling() -> Result<()> { Command::new(cargo_bin!()) From d46b3be1afcd0003585dcf838c1d290d277d6b51 Mon Sep 17 00:00:00 2001 From: Dave Nagoda Date: Fri, 2 Oct 2026 17:43:27 -0700 Subject: [PATCH 5/5] Describe what with_analyzer checks in the README with_analyzer parses the schema and query; it does not check the query against the schema. Say so, and note that analyze can also fail when it cannot select an operation in the query. --- README.md | 15 +++++++++------ 1 file changed, 9 insertions(+), 6 deletions(-) diff --git a/README.md b/README.md index f0572536..0d5e953f 100644 --- a/README.md +++ b/README.md @@ -74,9 +74,9 @@ Profiling is not available in batch mode. ## Library usage To compute scale factors for many inputs, use -`bluejay_schema_analyzer::BluejaySchemaAnalyzer::with_analyzer`. It parses and -validates the schema and query once, then calls your closure with an `analyze` -function that returns the scale factor for one input: +`bluejay_schema_analyzer::BluejaySchemaAnalyzer::with_analyzer`. It parses the +schema and query once, then calls your closure with an `analyze` function that +returns the scale factor for one input: ```rust use function_runner::bluejay_schema_analyzer::BluejaySchemaAnalyzer; @@ -90,9 +90,12 @@ let scale_factors = BluejaySchemaAnalyzer::with_analyzer( )??; ``` -The outer `Result` holds schema and query errors. The inner `Result` holds -analysis errors for an input. The `test_with_analyzer_analyzes_many_inputs` -test in `src/bluejay_schema_analyzer.rs` runs this pattern. +The outer `Result` holds schema and query parse errors. The inner `Result` +holds analysis errors for an input, for example when `analyze` cannot select +an operation in the query. `with_analyzer` does not check the query against +the schema, so a query with an unknown field still parses. The +`test_with_analyzer_analyzes_many_inputs` test in +`src/bluejay_schema_analyzer.rs` runs this pattern. ## Development