Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
74 changes: 74 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,80 @@ 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`. `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.

`--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.

## Library usage

To compute scale factors for many inputs, use
`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;

let scale_factors = BluejaySchemaAnalyzer::with_analyzer(
&schema,
Some("schema.graphql"),
&query,
Some("input.graphql"),
|analyze| inputs.iter().map(|input| analyze(input)).collect::<anyhow::Result<Vec<f64>>>(),
)??;
```

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

Building requires a rust toolchain of `1.66.0` to `1.67.0`. `cargo install --path . --locked` will build
Expand Down
257 changes: 257 additions & 0 deletions src/batch.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,257 @@
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<f64> + '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,
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>,
}

#[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<ScaleLimitsSource>,
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<u8>,
module: &Module,
analyze: Option<&ScaleFactorFn<'_>>,
options: &BatchOptions,
) -> Result<FunctionRunResult> {
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<u8>,
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,
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,
},
)
} 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<u8>, 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');
}
Loading
Loading