fn job_test_binary() -> std::path::PathBuf {
if let Some(path) = std::env::var_os("CARGO_BIN_EXE_camel") {
return std::path::PathBuf::from(path);
}
let target = std::env::var_os("CARGO_TARGET_DIR")
.map(std::path::PathBuf::from)
.unwrap_or_else(|| {
std::path::PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../../target")
});
let dev = target.join("debug").join("camel");
if dev.exists() {
return dev;
}
let release = target.join("release").join("camel");
assert!(
release.exists(),
"camel binary not built: run `cargo build -p camel-cli` (probed {} and {})",
dev.display(),
release.display()
);
release
}
fn run_camel_job(dir: &std::path::Path, args: &[&str]) -> (i32, String, String) {
run_camel_job_env(dir, args, &[])
}
fn write_job_fixture_config(dir: &std::path::Path) {
std::fs::write(
dir.join("Camel.toml"),
r#"[default]
routes = ["routes/*.yaml"]
log_level = "off"
watch = false
"#,
)
.expect("write Camel.toml");
}
fn write_tap_route(dir: &std::path::Path) {
std::fs::create_dir(dir.join("routes")).expect("mkdir routes");
std::fs::write(
dir.join("routes/job-route.yaml"),
r#"routes:
- id: "job-tap"
from: "direct:tap"
"#,
)
.expect("write route");
}
fn read_report(path: &std::path::Path) -> serde_json::Value {
let text = std::fs::read_to_string(path).expect("report file exists");
serde_json::from_str(text.trim()).expect("report is JSON")
}
#[test]
fn declared_args_validate_before_boot() {
let dir = tempfile::tempdir().expect("tempdir");
write_job_fixture_config(dir.path());
std::fs::write(
dir.path().join("job.job.yaml"),
r#"args:
name:
required: true
tier:
default: gold
execute:
mode: one-shot
timeout: 60s
send:
to: direct:tap
body: "ping"
routeFiles:
- routes/missing.yaml
"#,
)
.expect("write job doc");
let (code, _stdout, stderr) = run_camel_job(dir.path(), &["job.job.yaml", "--arg", "other=x"]);
assert_eq!(code, 2, "unknown --arg must exit 2; stderr:\n{stderr}");
assert!(
stderr.contains("unknown argument"),
"diagnostic must name the class; stderr:\n{stderr}"
);
assert!(
stderr.contains("other"),
"diagnostic must name the offending argument; stderr:\n{stderr}"
);
assert!(
!stderr.contains("camel-cli job failed"),
"validation must fail before boot; stderr:\n{stderr}"
);
let (code, _stdout, stderr) = run_camel_job(dir.path(), &["job.job.yaml"]);
assert_eq!(
code, 2,
"missing required arg must exit 2; stderr:\n{stderr}"
);
assert!(
stderr.contains("missing required argument"),
"diagnostic must name the class; stderr:\n{stderr}"
);
assert!(
stderr.contains("`name`"),
"diagnostic must name the missing argument; stderr:\n{stderr}"
);
let (code, _stdout, stderr) = run_camel_job(dir.path(), &["job.job.yaml", "--arg", "name=x"]);
assert_eq!(code, 2, "missing route file exits 2; stderr:\n{stderr}");
assert!(
!stderr.contains("unknown argument") && !stderr.contains("missing required argument"),
"valid args must pass validation; stderr:\n{stderr}"
);
}
#[test]
fn declared_defaults_and_explicit_values() {
let dir = tempfile::tempdir().expect("tempdir");
write_job_fixture_config(dir.path());
write_tap_route(dir.path());
std::fs::write(
dir.path().join("job.job.yaml"),
r#"args:
tier:
default: gold
execute:
mode: one-shot
timeout: 60s
capture-reply: true
send:
to: direct:tap
body: "value=${arg:tier}"
routeFiles:
- routes/job-route.yaml
"#,
)
.expect("write job doc");
let report = dir.path().join("report.json");
let (code, _stdout, stderr) = run_camel_job(
dir.path(),
&["job.job.yaml", "--report", report.to_str().expect("utf8")],
);
assert_eq!(code, 0, "default run must complete; stderr:\n{stderr}");
let json = read_report(&report);
assert_eq!(json["outcome"], "Completed", "report: {json}");
assert_eq!(json["reply"]["body"], "value=gold", "report: {json}");
let (code, _stdout, stderr) = run_camel_job(
dir.path(),
&[
"job.job.yaml",
"--report",
report.to_str().expect("utf8"),
"--arg",
"tier=silver",
],
);
assert_eq!(code, 0, "explicit run must complete; stderr:\n{stderr}");
let json = read_report(&report);
assert_eq!(json["outcome"], "Completed", "report: {json}");
assert_eq!(json["reply"]["body"], "value=silver", "report: {json}");
}
#[test]
fn legacy_args_remain_headers_with_deprecation() {
let dir = tempfile::tempdir().expect("tempdir");
write_job_fixture_config(dir.path());
write_tap_route(dir.path());
std::fs::write(
dir.path().join("job.job.yaml"),
r#"execute:
mode: one-shot
timeout: 60s
capture-reply: true
send:
to: direct:tap
body: "ping"
headers:
name: Doc
routeFiles:
- routes/job-route.yaml
"#,
)
.expect("write job doc");
let report = dir.path().join("report.json");
let (code, _stdout, stderr) = run_camel_job(
dir.path(),
&[
"job.job.yaml",
"--report",
report.to_str().expect("utf8"),
"--arg",
"name=First",
"--arg",
"name=Last",
],
);
assert_eq!(code, 0, "legacy run must complete; stderr:\n{stderr}");
let json = read_report(&report);
assert_eq!(json["outcome"], "Completed", "report: {json}");
assert_eq!(json["reply"]["headers"]["name"], "Last", "report: {json}");
assert!(
stderr.contains("--arg header injection"),
"stderr must carry the deprecation note; stderr:\n{stderr}"
);
assert!(
stderr.contains("deprecated"),
"stderr must mark the legacy behavior deprecated; stderr:\n{stderr}"
);
}
#[test]
fn declared_arg_is_not_implicit_header() {
let dir = tempfile::tempdir().expect("tempdir");
write_job_fixture_config(dir.path());
write_tap_route(dir.path());
std::fs::write(
dir.path().join("job.job.yaml"),
r#"args:
name:
required: true
execute:
mode: one-shot
timeout: 60s
capture-reply: true
send:
to: direct:tap
body: "${arg:name}"
headers:
X-Marker: m
routeFiles:
- routes/job-route.yaml
"#,
)
.expect("write job doc");
let report = dir.path().join("report.json");
let (code, _stdout, stderr) = run_camel_job(
dir.path(),
&[
"job.job.yaml",
"--report",
report.to_str().expect("utf8"),
"--arg",
"name=John",
],
);
assert_eq!(code, 0, "declared run must complete; stderr:\n{stderr}");
let json = read_report(&report);
assert_eq!(json["outcome"], "Completed", "report: {json}");
assert_eq!(json["reply"]["body"], "John", "report: {json}");
assert!(
json["reply"]["headers"].get("name").is_none(),
"declared args must not become implicit headers: {json}"
);
assert_eq!(
json["reply"]["headers"]["X-Marker"], "m",
"marker header must be present so the negative assert is not vacuous: {json}"
);
assert!(
!stderr.contains("deprecated"),
"declared documents are not the legacy path; stderr:\n{stderr}"
);
}
#[test]
fn declared_args_interpolate_all_field_positions() {
let dir = tempfile::tempdir().expect("tempdir");
write_job_fixture_config(dir.path());
write_tap_route(dir.path());
std::fs::write(
dir.path().join("job.job.yaml"),
r#"args:
target:
default: direct:tap
text:
default: hello
header:
default: gold
wait:
default: 60s
execute:
mode: one-shot
timeout: "${arg:wait}"
capture-reply: true
send:
to: "${arg:target}"
body: "${arg:text}"
headers:
X-Tier: "${arg:header}"
routeFiles:
- routes/job-route.yaml
"#,
)
.expect("write job doc");
let report = dir.path().join("report.json");
let (code, _stdout, stderr) = run_camel_job(
dir.path(),
&["job.job.yaml", "--report", report.to_str().expect("utf8")],
);
assert_eq!(code, 0, "all-field run must complete; stderr:\n{stderr}");
let json = read_report(&report);
assert_eq!(json["outcome"], "Completed", "report: {json}");
assert_eq!(json["reply"]["body"], "hello", "report: {json}");
assert_eq!(json["reply"]["headers"]["X-Tier"], "gold", "report: {json}");
}
fn write_typed_canonical_job(dir: &std::path::Path) {
std::fs::create_dir(dir.join("routes")).expect("mkdir routes");
std::fs::write(
dir.join("routes/job-route.yaml"),
r#"routes:
- id: "job-tap-in"
from: "direct:in"
steps:
- transform: {simple: "in-${body}"}
- id: "job-tap-out"
from: "direct:out"
steps:
- transform: {simple: "out-${body}"}
"#,
)
.expect("write route");
std::fs::write(
dir.join("job.job.yaml"),
r#"args:
target:
type: "enum[direct:in,direct:out]"
count:
type: int
default: "7"
verbose:
type: bool
wait:
type: int
default: "30"
execute:
mode: one-shot
timeout: "${arg:wait}s"
capture-reply: true
send:
to: "${arg:target}"
body: "n=${arg:count} v=${arg:verbose}"
headers:
tier: "${arg:target}"
routeFiles:
- routes/job-route.yaml
"#,
)
.expect("write job doc");
}
#[test]
fn typed_args_interpolate_canonical_forms_all_fields() {
let dir = tempfile::tempdir().expect("tempdir");
write_job_fixture_config(dir.path());
write_typed_canonical_job(dir.path());
let report = dir.path().join("report.json");
let (code, _stdout, stderr) = run_camel_job(
dir.path(),
&[
"job.job.yaml",
"--report",
report.to_str().expect("utf8"),
"--arg",
"verbose=false",
"--arg",
"target=direct:out",
"--arg",
"count=007",
],
);
assert_eq!(
code, 0,
"canonical-forms run must complete (timeout 30s accepted); stderr:\n{stderr}"
);
let json = read_report(&report);
assert_eq!(json["outcome"], "Completed", "report: {json}");
assert_eq!(json["reply"]["body"], "out-n=7 v=false", "report: {json}");
assert_eq!(
json["reply"]["headers"]["tier"], "direct:out",
"report: {json}"
);
}
#[test]
fn typed_coercion_failure_exits_2_before_boot() {
let dir = tempfile::tempdir().expect("tempdir");
write_job_fixture_config(dir.path());
write_typed_canonical_job(dir.path());
let report = dir.path().join("report.json");
let (code, _stdout, stderr) = run_camel_job(
dir.path(),
&[
"job.job.yaml",
"--report",
report.to_str().expect("utf8"),
"--arg",
"count=abc",
],
);
assert_eq!(code, 2, "coercion failure must exit 2; stderr:\n{stderr}");
assert!(
stderr.contains("invalid value `abc` for argument `count`"),
"diagnostic must name the raw value and the argument; stderr:\n{stderr}"
);
assert!(
stderr.contains("expected type `int`"),
"diagnostic must name the expected type; stderr:\n{stderr}"
);
assert!(
!stderr.contains("camel-cli job failed"),
"coercion must fail before boot; stderr:\n{stderr}"
);
assert!(
!report.exists(),
"coercion failure must write no report; stderr:\n{stderr}"
);
}
#[test]
fn untyped_document_behavior_unchanged() {
let dir = tempfile::tempdir().expect("tempdir");
write_job_fixture_config(dir.path());
write_tap_route(dir.path());
std::fs::write(
dir.path().join("job.job.yaml"),
r#"args:
tier:
default: gold
execute:
mode: one-shot
timeout: 60s
capture-reply: true
send:
to: direct:tap
body: "value=${arg:tier}"
routeFiles:
- routes/job-route.yaml
"#,
)
.expect("write job doc");
let report = dir.path().join("report.json");
let (code, _stdout, stderr) = run_camel_job(
dir.path(),
&[
"job.job.yaml",
"--report",
report.to_str().expect("utf8"),
"--arg",
"tier=007",
],
);
assert_eq!(code, 0, "untyped run must complete; stderr:\n{stderr}");
let json = read_report(&report);
assert_eq!(json["outcome"], "Completed", "report: {json}");
assert_eq!(
json["reply"]["body"], "value=007",
"untyped values must pass through verbatim (A2 behavior): {json}"
);
}
fn read_file_eventually(dir: &std::path::Path, name: &str) -> String {
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(2);
loop {
if let Ok(text) = std::fs::read_to_string(dir.join(name)) {
return text;
}
if std::time::Instant::now() >= deadline {
panic!("{name} missing under {} after 2 s", dir.display());
}
std::thread::sleep(std::time::Duration::from_millis(25));
}
}
#[test]
fn batch_typed_arg_coerces_and_drains() {
let dir = tempfile::tempdir().expect("tempdir");
write_job_fixture_config(dir.path());
std::fs::create_dir(dir.path().join("routes")).expect("mkdir routes");
let routes = format!(
r#"routes:
- id: "fan"
from: "direct:fan"
steps:
- to: "seda:w1"
- id: "w1"
from: "seda:w1"
steps:
- transform: {{simple: "id-${{header.batch_id}}"}}
- to: "file:{base}?fileName=tagged.txt"
"#,
base = dir.path().display()
);
std::fs::write(dir.path().join("routes/job-route.yaml"), routes).expect("write route");
std::fs::write(
dir.path().join("job.job.yaml"),
r#"args:
batch_id:
type: int
execute:
mode: batch
timeout: 60s
send:
to: direct:fan
body: "m"
headers:
batch_id: "${arg:batch_id}"
routeFiles:
- routes/job-route.yaml
"#,
)
.expect("write job doc");
let (code, stdout, stderr) =
run_camel_job(dir.path(), &["job.job.yaml", "--arg", "batch_id=007"]);
assert_eq!(
code, 0,
"batch run must complete and drain;\nstdout:\n{stdout}\nstderr:\n{stderr}"
);
let report: serde_json::Value =
serde_json::from_str(stdout.trim()).expect("stdout is the JSON report; got:\n{stdout}");
assert_eq!(report["mode"], "batch", "report: {report}");
assert_eq!(report["outcome"], "Completed", "report: {report}");
let tagged = read_file_eventually(dir.path(), "tagged.txt");
assert!(
tagged.contains("id-7"),
"tagged.txt must carry the COERCED canonical value (7, not 007); got: {tagged}"
);
}
fn run_camel_job_env(
dir: &std::path::Path,
args: &[&str],
env: &[(std::ffi::OsString, std::ffi::OsString)],
) -> (i32, String, String) {
let mut full: Vec<&str> = vec!["job"];
full.extend(args.iter().copied());
let output = std::process::Command::new(job_test_binary())
.args(full)
.envs(env.iter().cloned())
.current_dir(dir)
.output()
.expect("spawn camel binary");
(
output.status.code().unwrap_or(-1),
String::from_utf8_lossy(&output.stdout).into_owned(),
String::from_utf8_lossy(&output.stderr).into_owned(),
)
}
fn write_jobs_root_document(dir: &std::path::Path, name: &str, document: &str) {
let jobs = dir.join("jobs");
std::fs::create_dir_all(&jobs).expect("mkdir jobs");
std::fs::write(jobs.join(format!("{name}.job.yaml")), document).expect("write job doc");
}
fn jobs_root_help_document(description: &str, args: &str, send_to: &str) -> String {
format!(
r#"description: {description}
{args}execute:
mode: one-shot
timeout: 60s
send:
to: {send_to}
body: "ping"
routes: |
routes:
- id: "job-tap"
from: "direct:tap"
"#
)
}
#[test]
fn help_with_name_renders_declared_interface() {
let dir = tempfile::tempdir().expect("tempdir");
write_job_fixture_config(dir.path());
let document = jobs_root_help_document(
"Ingest the daily feed",
r#"args:
target:
required: true
description: Where to send the feed
retries:
default: "3"
description: How many attempts to make
"#,
"direct:tap",
);
write_jobs_root_document(dir.path(), "daily-sync", &document);
let (code, stdout, stderr) = run_camel_job(dir.path(), &["daily-sync", "--help"]);
assert_eq!(code, 0, "--help must exit 0; stderr:\n{stderr}");
assert!(
stdout.starts_with("daily-sync"),
"stdout must start with the stem; stdout:\n{stdout}"
);
assert!(
stdout.contains("Ingest the daily feed"),
"description must render; stdout:\n{stdout}"
);
assert!(
stdout.contains("Mode: one-shot"),
"mode row must render; stdout:\n{stdout}"
);
assert!(
stdout.contains("Sends to:"),
"send target row must render; stdout:\n{stdout}"
);
assert!(
stdout.contains(" retries string optional default=3 How many attempts to make"),
"optional argument row must render; stdout:\n{stdout}"
);
assert!(
stdout.contains(" target string required Where to send the feed"),
"required argument row must render; stdout:\n{stdout}"
);
assert!(
!stdout.contains("Usage:"),
"the interface is not clap help; stdout:\n{stdout}"
);
}
#[test]
fn help_with_name_reports_required_args_without_pairs() {
let dir = tempfile::tempdir().expect("tempdir");
write_job_fixture_config(dir.path());
let document = jobs_root_help_document(
"Ingest the daily feed",
r#"args:
target:
required: true
"#,
r#""${arg:target}""#,
);
write_jobs_root_document(dir.path(), "daily-sync", &document);
let (code, stdout, stderr) = run_camel_job(dir.path(), &["daily-sync", "--help"]);
assert_eq!(
code, 0,
"--help must exit 0 without pairs; stderr:\n{stderr}"
);
assert!(
stdout.contains("Sends to: ${arg:target}"),
"raw ${{arg:}} token must survive verbatim; stdout:\n{stdout}"
);
}
#[test]
fn help_no_args_block_prints_no_arguments() {
let dir = tempfile::tempdir().expect("tempdir");
write_job_fixture_config(dir.path());
let document = jobs_root_help_document("Do the thing quietly", "", "direct:tap");
write_jobs_root_document(dir.path(), "quiet-sync", &document);
let (code, stdout, stderr) = run_camel_job(dir.path(), &["quiet-sync", "--help"]);
assert_eq!(code, 0, "--help must exit 0; stderr:\n{stderr}");
assert!(
stdout.contains("Arguments:\n (no arguments)"),
"absent args block must render the placeholder row; stdout:\n{stdout}"
);
}
#[test]
fn help_writes_no_report_and_boots_nothing() {
let dir = tempfile::tempdir().expect("tempdir");
write_job_fixture_config(dir.path());
write_tap_route(dir.path());
write_jobs_root_document(
dir.path(),
"daily-sync",
r#"description: Ingest the daily feed
execute:
mode: one-shot
timeout: 60s
send:
to: direct:tap
body: "ping"
routeFiles:
- ../routes/job-route.yaml
"#,
);
let report = dir.path().join("out.json");
let (code, stdout, stderr) = run_camel_job_env(
dir.path(),
&["daily-sync", "--help", "--report", "out.json"],
&[(
std::ffi::OsString::from("CAMEL_JOB_SIGNAL_MARKER"),
std::ffi::OsString::from("1"),
)],
);
assert_eq!(code, 0, "--help must exit 0; stderr:\n{stderr}");
assert!(
!report.exists(),
"help must not write the report file; stdout:\n{stdout}"
);
assert!(
stdout.starts_with("daily-sync"),
"stdout is the declared interface; stdout:\n{stdout}"
);
assert!(
!stderr.contains("signal streams armed"),
"help installs no signal streams; stderr:\n{stderr}"
);
}
#[test]
fn help_unknown_name_fails_loud() {
let dir = tempfile::tempdir().expect("tempdir");
write_job_fixture_config(dir.path());
let other = jobs_root_help_document("Unrelated job", "", "direct:tap");
write_jobs_root_document(dir.path(), "other", &other);
let (code, _stdout, stderr) = run_camel_job(dir.path(), &["ghost", "--help"]);
assert_eq!(code, 2, "unknown name must exit 2; stderr:\n{stderr}");
assert!(
stderr.contains("ghost"),
"diagnostic must name the job; stderr:\n{stderr}"
);
assert!(
stderr.contains("no job `ghost` in any configured root"),
"bare-name resolution diagnostic must carry; stderr:\n{stderr}"
);
assert!(
!stderr.contains("Usage:"),
"failure is not clap help; stderr:\n{stderr}"
);
}
#[test]
fn help_malformed_document_fails_loud() {
let dir = tempfile::tempdir().expect("tempdir");
write_job_fixture_config(dir.path());
write_jobs_root_document(
dir.path(),
"broken",
r#"description: Broken on purpose
totallyUnknownField: yes
execute:
mode: one-shot
timeout: 60s
send:
to: direct:tap
routes: |
routes:
- id: "job-tap"
from: "direct:tap"
"#,
);
let (code, _stdout, stderr) = run_camel_job(dir.path(), &["broken", "--help"]);
assert_eq!(code, 2, "malformed document must exit 2; stderr:\n{stderr}");
assert!(
stderr.contains("unknown field in job document"),
"parse diagnostic must carry; stderr:\n{stderr}"
);
assert!(
!stderr.contains("Usage:"),
"failure is not clap help; stderr:\n{stderr}"
);
}
#[test]
fn help_without_name_prints_usage() {
let dir = tempfile::tempdir().expect("tempdir");
write_job_fixture_config(dir.path());
let document = jobs_root_help_document("Ingest the daily feed", "", "direct:tap");
write_jobs_root_document(dir.path(), "daily-sync", &document);
let (code, stdout, _stderr) = run_camel_job(dir.path(), &["--help"]);
assert_eq!(code, 0, "usage help must exit 0");
assert!(
stdout.contains("Usage: camel job"),
"usage line must carry; stdout:\n{stdout}"
);
let (code, stdout, _stderr) = run_camel_job(dir.path(), &[]);
assert_eq!(code, 0, "bare listing must stay exit 0");
assert!(
stdout.contains("Jobs in jobs/:"),
"bare invocation still lists; stdout:\n{stdout}"
);
}
#[test]
fn help_short_flag_behaves_like_long() {
let dir = tempfile::tempdir().expect("tempdir");
write_job_fixture_config(dir.path());
let document = jobs_root_help_document(
"Ingest the daily feed",
r#"args:
target:
required: true
description: Where to send the feed
retries:
default: "3"
description: How many attempts to make
"#,
"direct:tap",
);
write_jobs_root_document(dir.path(), "daily-sync", &document);
let (long_code, long_stdout, _long_stderr) =
run_camel_job(dir.path(), &["daily-sync", "--help"]);
let (short_code, short_stdout, short_stderr) = run_camel_job(dir.path(), &["daily-sync", "-h"]);
assert_eq!(short_code, 0, "-h must exit 0; stderr:\n{short_stderr}");
assert_eq!(long_code, 0, "--help must exit 0");
assert_eq!(
short_stdout, long_stdout,
"-h and --help must render identically"
);
}
mod exit_code_tests {
use crate::commands::job::{JobReport, exit_code_for};
#[test]
fn job_exit_code_interrupted() {
let report = JobReport {
document: "doc".to_string(),
mode: "one-shot".to_string(),
outcome: "Interrupted",
terminated_early: false,
duration_ms: 1,
reply: None,
error: Some("interrupted by signal (SIGINT/SIGTERM)".to_string()),
shutdown_error: None,
};
assert_eq!(exit_code_for(report.outcome), 2);
}
}
mod report_tests {
use crate::commands::job::{
JobReport, MIN_SHUTDOWN_BUDGET, exit_code_for, record_shutdown_failure,
};
#[test]
fn job_report_interrupted_serializes() {
let report = JobReport {
document: "doc".to_string(),
mode: "one-shot".to_string(),
outcome: "Interrupted",
terminated_early: false,
duration_ms: 1,
reply: None,
error: Some("interrupted by signal (SIGINT/SIGTERM)".to_string()),
shutdown_error: Some("shutdown failure: x".to_string()),
};
let json = serde_json::to_value(&report).expect("report must serialize");
assert_eq!(
json["outcome"],
serde_json::json!("Interrupted"),
"outcome must stay Interrupted: {json}"
);
assert!(
json["shutdown_error"].is_string(),
"shutdown detail must serialize: {json}"
);
}
#[test]
fn job_interrupted_shutdown_failure_preserves_verdict() {
let mut report = JobReport {
document: "doc".to_string(),
mode: "one-shot".to_string(),
outcome: "Interrupted",
terminated_early: false,
duration_ms: 1,
reply: None,
error: Some("interrupted by signal (SIGINT/SIGTERM)".to_string()),
shutdown_error: None,
};
record_shutdown_failure(
&mut report,
"shutdown failure: x".to_string(),
MIN_SHUTDOWN_BUDGET,
);
assert_eq!(report.outcome, "Interrupted");
assert_eq!(
report.shutdown_error.as_deref(),
Some("shutdown failure: x"),
"non-zero-budget teardown detail must be recorded"
);
assert_eq!(exit_code_for(report.outcome), 2);
}
#[test]
fn shutdown_error_serializes_alongside_error() {
let report = JobReport {
document: "doc".to_string(),
mode: "one-shot".to_string(),
outcome: "Failed",
terminated_early: false,
duration_ms: 1,
reply: None,
error: Some("pipeline failed".to_string()),
shutdown_error: Some("shutdown failure: x".to_string()),
};
let json = serde_json::to_string(&report).expect("report must serialize");
assert!(
json.contains("pipeline failed"),
"verdict error must serialize: {json}"
);
assert!(
json.contains("shutdown failure: x"),
"shutdown detail must serialize: {json}"
);
}
#[test]
fn shutdown_error_omitted_when_absent() {
let report = JobReport {
document: "doc".to_string(),
mode: "one-shot".to_string(),
outcome: "Failed",
terminated_early: false,
duration_ms: 1,
reply: None,
error: Some("pipeline failed".to_string()),
shutdown_error: None,
};
let json = serde_json::to_string(&report).expect("report must serialize");
assert!(
!json.contains("shutdown_error"),
"absent shutdown_error must be omitted: {json}"
);
}
}
mod shutdown_budget_tests {
use std::time::{Duration, Instant};
use crate::commands::job::document::JobMode;
use crate::commands::job::{MIN_SHUTDOWN_BUDGET, shutdown_budget};
#[test]
fn shutdown_budget_batch_is_remaining() {
let deadline = Instant::now() + Duration::from_secs(3);
let budget = shutdown_budget(JobMode::Batch, deadline);
assert!(
budget <= Duration::from_secs(3) && budget > Duration::from_secs(2),
"expected ~3s remaining, got {budget:?}"
);
}
#[test]
fn shutdown_budget_batch_zero_when_past() {
let deadline = Instant::now() - Duration::from_secs(1);
assert_eq!(shutdown_budget(JobMode::Batch, deadline), Duration::ZERO);
}
#[test]
fn shutdown_budget_one_shot_floored() {
let deadline = Instant::now() - Duration::from_secs(1);
assert_eq!(
shutdown_budget(JobMode::OneShot, deadline),
MIN_SHUTDOWN_BUDGET
);
}
#[test]
fn shutdown_budget_one_shot_is_remaining_when_large() {
let deadline = Instant::now() + Duration::from_secs(10);
let budget = shutdown_budget(JobMode::OneShot, deadline);
assert!(
budget <= Duration::from_secs(10) && budget > MIN_SHUTDOWN_BUDGET,
"expected ~10s remaining, got {budget:?}"
);
}
#[test]
fn job_interrupted_shutdown_budget_by_mode() {
let spent = Instant::now() - Duration::from_secs(1);
assert!(
shutdown_budget(JobMode::OneShot, spent) >= MIN_SHUTDOWN_BUDGET,
"interrupted one-shot teardown keeps the floor"
);
assert_eq!(
shutdown_budget(JobMode::Batch, spent),
Duration::ZERO,
"interrupted batch teardown has no floor"
);
let soon = Instant::now() + Duration::from_secs(3);
assert_eq!(shutdown_budget(JobMode::OneShot, soon), MIN_SHUTDOWN_BUDGET);
let batch = shutdown_budget(JobMode::Batch, soon);
assert!(
batch <= Duration::from_secs(3) && batch > Duration::from_secs(2),
"interrupted batch teardown gets the remaining deadline, got {batch:?}"
);
}
}
mod store_plan_tests {
use crate::commands::job::filter_store_source_plan;
use crate::compile::store::{StoreDocument, StoreEntryKind, VirtualDocumentStore};
fn store(plan: &[&str]) -> VirtualDocumentStore {
let job_text =
"execute:\n mode: one-shot\n timeout: 30s\n send:\n to: direct:start\n";
VirtualDocumentStore::build(
"job.job.yaml",
&[
StoreDocument {
path: "job.job.yaml".to_string(),
kind: StoreEntryKind::Job,
bytes: job_text.as_bytes().to_vec(),
},
StoreDocument {
path: "other.job.yaml".to_string(),
kind: StoreEntryKind::Job,
bytes: job_text.as_bytes().to_vec(),
},
StoreDocument {
path: "routes/a.yaml".to_string(),
kind: StoreEntryKind::Route,
bytes: "routes:\n - id: a\n from: timer:a\n"
.as_bytes()
.to_vec(),
},
],
&[],
&plan.iter().map(|p| (*p).to_string()).collect::<Vec<_>>(),
)
.expect("valid store builds")
}
#[test]
fn store_plan_entry_dropped_routes_kept() {
let mut store = store(&["job.job.yaml", "routes/a.yaml"]);
assert_eq!(filter_store_source_plan(&mut store), None);
assert_eq!(store.index.source_plan.references, vec!["routes/a.yaml"]);
}
#[test]
fn store_plan_extra_job_reference_named() {
let mut store = store(&["job.job.yaml", "routes/a.yaml", "other.job.yaml"]);
assert_eq!(
filter_store_source_plan(&mut store),
Some("other.job.yaml".to_string())
);
}
}