use std::collections::{HashMap, HashSet};
use std::io::BufRead;
use std::sync::{Arc, Mutex};
use std::time::Duration;
use anyhow::Result;
use crate::costs::{EndpointUsage, UsageSnapshot};
use crate::events::TestEvent;
use crate::reporting::Reporter;
use crate::runner::{RunReport, ScenarioRunner};
use crate::scenario::{AssertDefinition, ScenarioConfig, TestGroup};
const DEFAULT_THREAD_CAP: usize = 64;
const DEFAULT_AUTO_CEILING: usize = 8;
const MIN_HEADROOM_BYTES: u64 = 256 * 1024 * 1024;
const MAX_LAUNCH_RETRIES: u32 = 3;
const BACKOFF_BASE_MS: u64 = 600;
const BACKOFF_STEP_MS: u64 = 600;
const MAX_BACKOFF_MS: u64 = 4000;
const MIN_FOOTPRINT: f64 = 128.0 * 1024.0 * 1024.0; const MAX_FOOTPRINT: f64 = 2.0 * 1024.0 * 1024.0 * 1024.0;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ParallelMode {
Manual(usize),
Auto {
min: u32,
max: u32,
},
}
#[must_use]
pub fn mode_from_cli(parallel: Option<u32>, min: u32, max: u32) -> ParallelMode {
parallel.map_or_else(
|| ParallelMode::Auto {
min: min.max(1),
max,
},
|k| ParallelMode::Manual(k.max(1) as usize),
)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct MemoryInfo {
pub total: u64,
pub available: u64,
}
#[derive(Clone)]
pub struct ScenarioFile {
pub label: String,
pub config: ScenarioConfig,
pub definitions: Vec<AssertDefinition>,
pub tests: Vec<TestGroup>,
}
impl ScenarioFile {
fn concurrency_key(&self, index: usize) -> String {
self.config
.concurrency_group
.clone()
.unwrap_or_else(|| format!("<file {index}>"))
}
}
#[derive(Clone)]
pub struct RunOptions {
pub mode: ParallelMode,
pub reporter: Arc<Reporter>,
pub memory: Option<MemoryInfo>,
}
#[derive(Debug, Default)]
pub struct ParallelRun {
pub report: RunReport,
pub per_test: Vec<(String, UsageSnapshot)>,
pub global: UsageSnapshot,
}
#[allow(clippy::cast_possible_truncation, clippy::significant_drop_tightening)]
pub fn run_scenarios(files: Vec<ScenarioFile>, opts: RunOptions) -> Result<ParallelRun> {
let reporter = opts.reporter;
let total_tests: u32 = files.iter().map(|f| f.tests.len() as u32).sum();
reporter.emit(&TestEvent::RunStarted { total_tests })?;
if files.is_empty() {
reporter.emit(&TestEvent::RunFinished {
tests_passed: 0,
tests_failed: 0,
steps_passed: 0,
steps_failed: 0,
steps_skipped: 0,
total_cost: 0.0,
total_tokens: 0,
total_input_tokens: 0,
total_output_tokens: 0,
total_cached_input_tokens: 0,
total_cache_creation_input_tokens: 0,
models: Vec::new(),
total_calls: 0,
})?;
return Ok(ParallelRun::default());
}
let n = files.len();
let mut results: Vec<Option<ParallelRun>> = Vec::new();
for _ in 0..n {
results.push(None);
}
let state = Arc::new(Mutex::new(SchedulerState {
files,
results,
started: vec![false; n],
active_groups: HashSet::new(),
}));
let gate = Arc::new(Mutex::new(make_gate(&opts.mode, opts.memory, n)));
let workers = gate.lock().unwrap().threads;
let mut threads: Vec<std::thread::JoinHandle<()>> = Vec::new();
for _ in 0..workers {
let state = Arc::clone(&state);
let gate = Arc::clone(&gate);
let reporter = Arc::clone(&reporter);
threads.push(std::thread::spawn(move || {
worker_loop(&state, &gate, &reporter);
}));
}
for thread in threads {
thread.join().unwrap();
}
let run = {
let st = state.lock().unwrap();
let mut report = RunReport::default();
let mut per_test: Vec<(String, UsageSnapshot)> = Vec::new();
let mut globals: Vec<UsageSnapshot> = Vec::new();
for r in &st.results {
if let Some(run) = r.as_ref() {
report.tests_passed += run.report.tests_passed;
report.tests_failed += run.report.tests_failed;
report.passed += run.report.passed;
report.failed += run.report.failed;
report.skipped += run.report.skipped;
report.details.extend(run.report.details.clone());
per_test.extend(run.per_test.clone());
globals.push(run.global.clone());
}
}
ParallelRun {
report,
per_test,
global: merge_globals(&globals),
}
};
reporter.emit(&TestEvent::RunFinished {
tests_passed: run.report.tests_passed,
tests_failed: run.report.tests_failed,
steps_passed: run.report.passed,
steps_failed: run.report.failed,
steps_skipped: run.report.skipped,
total_cost: run.global.total_cost,
total_tokens: run.global.total_tokens,
total_input_tokens: run.global.total_input_tokens,
total_output_tokens: run.global.total_output_tokens,
total_cached_input_tokens: run.global.total_cached_input_tokens,
total_cache_creation_input_tokens: run.global.total_cache_creation_input_tokens,
models: run.global.models.clone(),
total_calls: run.global.total_calls,
})?;
Ok(run)
}
struct Gate {
active: usize,
limit: usize,
lower: usize,
threads: usize,
ramp_ceiling: usize,
memory_total: Option<u64>,
available: Option<u64>,
footprint: f64,
footprint_count: u32,
retry_backoff: u64,
}
fn make_gate(mode: &ParallelMode, memory: Option<MemoryInfo>, files_len: usize) -> Gate {
match mode {
ParallelMode::Manual(k) => {
let k = *k;
let threads = k.max(1).min(files_len).max(1);
let limit = k.max(1).min(threads);
Gate {
active: 0,
limit,
lower: limit,
threads,
ramp_ceiling: limit,
memory_total: memory.map(|m| m.total),
available: memory.map(|m| m.available),
footprint: 0.0,
footprint_count: 0,
retry_backoff: 0,
}
}
ParallelMode::Auto { min, max } => {
let min = *min;
let max = *max;
let lower = min.max(1) as usize;
let user_max: Option<usize> = if max == 0 { None } else { Some(max as usize) };
let threads = user_max.unwrap_or(DEFAULT_THREAD_CAP).min(files_len).max(1);
let ramp_ceiling = user_max.unwrap_or(DEFAULT_AUTO_CEILING).min(threads).max(1);
Gate {
active: 0,
limit: lower.min(threads).max(1),
lower: lower.min(threads).max(1),
threads,
ramp_ceiling,
memory_total: memory.map(|m| m.total),
available: memory.map(|m| m.available),
footprint: 0.0,
footprint_count: 0,
retry_backoff: 0,
}
}
}
}
impl Gate {
fn can_launch(&self) -> bool {
self.active < self.limit && !self.memory_blocked()
}
#[allow(
clippy::unnecessary_unwrap,
clippy::cast_precision_loss,
clippy::cast_possible_truncation,
clippy::cast_sign_loss
)]
fn memory_blocked(&self) -> bool {
let Some(available) = self.available else {
return false;
};
if available < MIN_HEADROOM_BYTES {
return true;
}
if self.footprint_count > 0 && self.footprint > 0.0 {
available < self.footprint as u64
} else {
false
}
}
#[allow(clippy::missing_const_for_fn)]
fn launch_started(&mut self, before: Option<u64>) {
if before.is_some() {
self.available = before;
}
self.active += 1;
}
#[allow(
clippy::unnecessary_unwrap,
clippy::cast_precision_loss,
clippy::suboptimal_flops
)]
fn launch_finished(&mut self, before: Option<u64>, after: Option<u64>) {
if before.is_some() {
self.available = before;
}
if after.is_some() {
self.available = after;
}
if before.is_some() && after.is_some() {
let delta = before.unwrap().saturating_sub(after.unwrap());
if delta > 0 {
let d = delta as f64;
if self.footprint_count == 0 {
self.footprint = d;
} else {
self.footprint = self.footprint * 0.7 + d * 0.3;
}
self.footprint = self.footprint.clamp(MIN_FOOTPRINT, MAX_FOOTPRINT);
self.footprint_count = (self.footprint_count + 1).min(20);
}
}
}
#[allow(
clippy::unnecessary_unwrap,
clippy::cast_precision_loss,
clippy::cast_possible_truncation,
clippy::cast_sign_loss
)]
fn on_success(&mut self) {
if self.footprint_count > 0 && self.footprint > 0.0 && self.memory_total.is_some() {
let total = self.memory_total.unwrap();
let available = self.available.unwrap_or(total);
if available > 0 {
let capacity = (available as f64 / self.footprint) as usize;
self.limit = capacity.max(self.lower).min(self.threads);
return;
}
}
self.limit = (self.limit + 1).min(self.ramp_ceiling).max(self.lower);
}
fn on_launch_failure(&mut self) {
self.limit = (self.limit / 2).max(self.lower);
self.retry_backoff = (self.retry_backoff + BACKOFF_STEP_MS).min(MAX_BACKOFF_MS);
}
#[allow(clippy::missing_const_for_fn)]
fn release(&mut self, after: Option<u64>) {
if after.is_some() {
self.available = after;
}
self.active = self.active.saturating_sub(1);
}
}
enum Claim {
Take(usize),
Wait,
Done,
}
struct SchedulerState {
files: Vec<ScenarioFile>,
results: Vec<Option<ParallelRun>>,
started: Vec<bool>,
active_groups: HashSet<String>,
}
impl SchedulerState {
fn claim(&mut self) -> Claim {
for i in 0..self.files.len() {
if self.started[i] {
continue;
}
let key = self.files[i].concurrency_key(i);
if self.active_groups.contains(&key) {
continue;
}
self.started[i] = true;
self.active_groups.insert(key);
return Claim::Take(i);
}
if self.started.iter().all(|b| *b) {
Claim::Done
} else {
Claim::Wait
}
}
fn complete(&mut self, i: usize) {
let key = self.files[i].concurrency_key(i);
self.active_groups.remove(&key);
}
}
enum Decision {
Done(ParallelRun),
Retry(u64),
GiveUp(String),
}
#[allow(clippy::significant_drop_tightening)]
fn worker_loop(
state: &Arc<Mutex<SchedulerState>>,
gate: &Arc<Mutex<Gate>>,
reporter: &Arc<Reporter>,
) {
loop {
let claimed = state.lock().unwrap().claim();
match claimed {
Claim::Done => break,
Claim::Wait => std::thread::sleep(Duration::from_millis(25)),
Claim::Take(i) => {
let file = state.lock().unwrap().files[i].clone();
let run = run_file_with_retries(&file, gate, reporter);
let mut st = state.lock().unwrap();
st.results[i] = Some(run);
st.complete(i);
}
}
}
}
#[allow(clippy::significant_drop_tightening)]
fn run_file_with_retries(
file: &ScenarioFile,
gate: &Arc<Mutex<Gate>>,
reporter: &Arc<Reporter>,
) -> ParallelRun {
let mut attempt: u32 = 0;
loop {
attempt += 1;
let before = available_memory_now();
{
loop {
let mut g = gate.lock().unwrap();
if g.can_launch() {
g.launch_started(before);
break;
}
std::thread::sleep(Duration::from_millis(25));
}
}
let result = run_one_file(file, reporter);
let after = available_memory_now();
let decision = {
let mut g = gate.lock().unwrap();
g.launch_finished(before, after);
match &result {
Ok(_) => {
g.on_success();
g.release(after);
Decision::Done(result.unwrap())
}
Err(e) => {
let oom = is_retryable_oom(e);
if oom && attempt < MAX_LAUNCH_RETRIES {
let backoff = g.retry_backoff.max(BACKOFF_BASE_MS);
g.on_launch_failure();
g.release(after);
Decision::Retry(backoff)
} else {
g.on_launch_failure();
g.release(after);
Decision::GiveUp(e.clone())
}
}
}
};
match decision {
Decision::Done(run) => return run,
Decision::Retry(ms) => {
reporter.warn(format!(
"{}: browser launch failed (likely out of memory); retrying ({attempt}/{MAX_LAUNCH_RETRIES}) in {ms}ms",
file.label,
));
std::thread::sleep(Duration::from_millis(ms));
}
Decision::GiveUp(e) => {
reporter.error(format!("{}: {e}", file.label));
return synthesized_failed_report(file);
}
}
}
}
fn run_one_file(file: &ScenarioFile, reporter: &Arc<Reporter>) -> Result<ParallelRun, String> {
let runner = ScenarioRunner::with_reporter_parallel(
file.config.clone(),
file.definitions.clone(),
Arc::clone(reporter),
);
match runner.run(&file.tests) {
Ok(report) => {
let usage = runner.usage_tracker();
Ok(ParallelRun {
report,
per_test: usage.per_test_snapshots(),
global: usage.global_snapshot(),
})
}
Err(e) => Err(e.to_string()),
}
}
#[must_use]
fn is_retryable_oom(err: &str) -> bool {
const KEYWORDS: [&str; 8] = [
"memory",
"cannot allocate",
"out of memory",
"killed",
"oom",
"resource temporarily unavailable",
"failed to allocate",
"no memory",
];
let e = err.to_lowercase();
KEYWORDS.iter().any(|kw| e.contains(kw))
}
#[allow(clippy::cast_possible_truncation)]
fn synthesized_failed_report(file: &ScenarioFile) -> ParallelRun {
let n = file.tests.len() as u32;
ParallelRun {
report: RunReport {
tests_passed: 0,
tests_failed: n,
passed: 0,
failed: n,
skipped: 0,
details: Vec::new(),
},
per_test: Vec::new(),
global: UsageSnapshot::default(),
}
}
#[must_use]
fn merge_globals(snapshots: &[UsageSnapshot]) -> UsageSnapshot {
let mut endpoints: HashMap<String, EndpointUsage> = HashMap::new();
for snapshot in snapshots {
for (name, usage) in &snapshot.endpoints {
let acc = endpoints.entry(name.clone()).or_default();
acc.calls += usage.calls;
acc.input_tokens += usage.input_tokens;
acc.output_tokens += usage.output_tokens;
acc.cached_input_tokens += usage.cached_input_tokens;
acc.cache_creation_input_tokens += usage.cache_creation_input_tokens;
acc.cost += usage.cost;
acc.models.extend(usage.models.iter().cloned());
}
}
UsageSnapshot::from_endpoints(&endpoints)
}
#[must_use]
#[allow(clippy::unnecessary_unwrap)]
pub async fn probe_memory_async() -> Option<MemoryInfo> {
let total = run_sh_async("free -b | awk '/^Mem:/{print $2}'").await;
let available = run_sh_async("free -b | awk '/^Mem:/{print $7}'").await;
if total.is_some() && available.is_some() {
let t = total.unwrap().trim().parse::<u64>().ok();
let a = available.unwrap().trim().parse::<u64>().ok();
if t.is_some() && a.is_some() {
return Some(MemoryInfo {
total: t.unwrap(),
available: a.unwrap(),
});
}
}
let total_mac = run_sh_async("sysctl -n hw.memsize").await;
let page_size = run_sh_async("sysctl -n hw.pagesize").await;
let pages =
run_sh_async("vm_stat | awk '/^Pages free:/{gsub(/[^0-9]/, \"\", $3); print $3}'").await;
if total_mac.is_some() && pages.is_some() && page_size.is_some() {
let t = total_mac.unwrap().trim().parse::<u64>().ok();
let p = pages.unwrap().trim().parse::<u64>().ok();
let ps = page_size.unwrap().trim().parse::<u64>().ok();
if t.is_some() && p.is_some() && ps.is_some() {
return Some(MemoryInfo {
total: t.unwrap(),
available: p.unwrap() * ps.unwrap(),
});
}
}
None
}
#[must_use]
fn available_memory_now() -> Option<u64> {
run_sh_sync("free -b 2>/dev/null | awk '/^Mem:/{print $7}'")
.and_then(|s| s.trim().parse::<u64>().ok())
}
#[must_use]
async fn run_sh_async(script: &str) -> Option<String> {
let output = tokio::process::Command::new("sh")
.args(["-c", script])
.output()
.await
.ok()?;
if !output.status.success() {
return None;
}
Some(String::from_utf8_lossy(&output.stdout).trim().to_string())
}
#[must_use]
fn run_sh_sync(script: &str) -> Option<String> {
let mut child = std::process::Command::new("sh")
.args(["-c", script])
.stdout(std::process::Stdio::piped())
.stderr(std::process::Stdio::inherit())
.spawn()
.ok()?;
let stdout = child.stdout.as_mut()?;
let mut reader = std::io::BufReader::new(stdout);
let mut line = String::new();
let _ = reader.read_line(&mut line);
if line.is_empty() {
None
} else {
Some(line.trim().to_string())
}
}
#[cfg(test)]
mod tests {
use std::sync::Arc;
use crate::costs::{EndpointUsage, UsageSnapshot};
use crate::scenario::ScenarioConfig;
use super::{
is_retryable_oom, make_gate, merge_globals, mode_from_cli, run_sh_sync, Claim,
ParallelMode, RunOptions, ScenarioFile, SchedulerState,
};
fn file(label: &str, group: Option<&str>) -> ScenarioFile {
ScenarioFile {
label: label.to_owned(),
config: ScenarioConfig {
concurrency_group: group.map(std::borrow::ToOwned::to_owned),
..ScenarioConfig::default()
},
definitions: Vec::new(),
tests: Vec::new(),
}
}
fn state(files: Vec<ScenarioFile>) -> SchedulerState {
let n = files.len();
SchedulerState {
files,
results: Vec::new(),
started: vec![false; n],
active_groups: std::collections::HashSet::new(),
}
}
#[test]
fn test_distinct_groups_claim_in_parallel() {
let mut s = state(vec![file("a", None), file("b", Some("x")), file("c", None)]);
assert!(matches!(s.claim(), Claim::Take(0)));
assert!(matches!(s.claim(), Claim::Take(1)));
assert!(matches!(s.claim(), Claim::Take(2)));
s.complete(1);
assert!(matches!(s.claim(), Claim::Done));
}
#[test]
fn test_same_group_blocks_until_completed() {
let mut s = state(vec![
file("a", Some("g")),
file("b", Some("g")),
file("c", None),
]);
assert!(matches!(s.claim(), Claim::Take(0)));
assert!(
matches!(s.claim(), Claim::Take(2)),
"a different group still runs while 'g' is active"
);
assert!(
matches!(s.claim(), Claim::Wait),
"b is blocked by a's group"
);
s.complete(0);
assert!(
matches!(s.claim(), Claim::Take(1)),
"b runs after a finishes"
);
s.complete(1);
assert!(matches!(s.claim(), Claim::Done));
}
#[test]
fn test_no_group_means_own_group() {
let mut s = state(vec![file("a", None), file("b", None)]);
assert!(matches!(s.claim(), Claim::Take(0)));
assert!(
matches!(s.claim(), Claim::Take(1)),
"no-group files run in parallel"
);
}
#[test]
fn test_mode_from_cli_manual_wins() {
assert_eq!(mode_from_cli(Some(5), 1, 0), ParallelMode::Manual(5));
}
#[test]
fn test_mode_from_cli_auto_with_defaults() {
assert_eq!(
mode_from_cli(None, 1, 0),
ParallelMode::Auto { min: 1, max: 0 }
);
}
#[test]
fn test_gate_manual_is_fixed() {
let gate = make_gate(&ParallelMode::Manual(3), None, 10);
assert_eq!(gate.limit, 3);
assert_eq!(gate.threads, 3);
assert_eq!(gate.lower, 3);
}
#[test]
fn test_gate_auto_ramps_from_min_toward_ceiling() {
let gate = make_gate(&ParallelMode::Auto { min: 1, max: 0 }, None, 100);
assert_eq!(gate.lower, 1);
assert_eq!(gate.limit, 1);
assert_eq!(gate.ramp_ceiling, 8);
assert_eq!(gate.threads, 64);
}
#[test]
fn test_gate_auto_respects_user_max() {
let gate = make_gate(&ParallelMode::Auto { min: 1, max: 8 }, None, 100);
assert_eq!(gate.threads, 8);
assert_eq!(gate.ramp_ceiling, 8);
}
#[test]
fn test_gate_memory_guard_blocks_low_headroom() {
let gate = make_gate(
&ParallelMode::Auto { min: 1, max: 4 },
Some(crate::parallel::MemoryInfo {
total: 1_000_000_000,
available: 50_000_000, }),
4,
);
assert!(gate.memory_blocked());
assert!(!gate.can_launch());
}
#[test]
fn test_gate_memory_guard_allows_headroom() {
let gate = make_gate(
&ParallelMode::Auto { min: 1, max: 4 },
Some(crate::parallel::MemoryInfo {
total: 1_000_000_000,
available: 900_000_000,
}),
4,
);
assert!(!gate.memory_blocked());
assert!(gate.can_launch());
}
#[test]
fn test_gate_memory_guard_ignores_huge_host_total() {
let gate = make_gate(
&ParallelMode::Auto { min: 1, max: 4 },
Some(crate::parallel::MemoryInfo {
total: 500_000_000_000, available: 100_000_000, }),
4,
);
assert!(gate.memory_blocked(), "must block despite a 500 GB 'total'");
}
#[test]
fn test_gate_memory_guard_blocks_when_no_room_for_one_footprint() {
let mut gate = make_gate(
&ParallelMode::Auto { min: 1, max: 4 },
Some(crate::parallel::MemoryInfo {
total: 8_000_000_000,
available: 300_000_000, }),
4,
);
gate.footprint = 500.0 * 1024.0 * 1024.0;
gate.footprint_count = 3;
assert!(gate.memory_blocked());
}
#[test]
fn test_oom_failure_halves_limit_and_sets_backoff() {
let mut gate = make_gate(&ParallelMode::Auto { min: 1, max: 16 }, None, 100);
gate.limit = 16;
gate.on_launch_failure();
assert_eq!(gate.limit, 8);
assert!(gate.retry_backoff > 0);
}
#[test]
fn test_is_retryable_oom_matches_memory_errors() {
assert!(is_retryable_oom("failed to launch browser: out of memory"));
assert!(is_retryable_oom("cannot allocate memory for page"));
assert!(!is_retryable_oom("Chrome binary not found"));
assert!(!is_retryable_oom("invalid URL"));
}
#[test]
fn test_run_options_cloneable() {
let _ = RunOptions {
mode: ParallelMode::Manual(2),
reporter: Arc::new(crate::reporting::Reporter::default()),
memory: None,
};
}
#[test]
fn test_launch_finished_learns_footprint_from_delta() {
let mut gate = make_gate(
&ParallelMode::Auto { min: 1, max: 4 },
Some(crate::parallel::MemoryInfo {
total: 8_000_000_000,
available: 8_000_000_000,
}),
4,
);
gate.launch_finished(Some(1_000_000_000), Some(600_000_000));
assert!((gate.footprint - 400_000_000.0).abs() < 1.0);
assert_eq!(gate.footprint_count, 1);
assert_eq!(gate.available, Some(600_000_000));
}
#[test]
fn test_on_success_raises_to_memory_capacity() {
let mut gate = make_gate(
&ParallelMode::Auto { min: 1, max: 0 },
Some(crate::parallel::MemoryInfo {
total: 8_000_000_000,
available: 4_000_000_000,
}),
100,
);
assert_eq!(gate.limit, 1);
gate.footprint = 1_000_000_000.0;
gate.footprint_count = 3;
gate.on_success();
assert_eq!(gate.limit, 4);
}
#[test]
fn test_on_success_ramps_when_memory_unknown() {
let mut gate = make_gate(&ParallelMode::Auto { min: 1, max: 0 }, None, 100);
assert_eq!(gate.limit, 1);
gate.on_success();
assert_eq!(gate.limit, 2, "ramps up one at a time");
}
#[test]
fn test_can_launch_respects_active_limit() {
let mut gate = make_gate(&ParallelMode::Manual(2), None, 10);
assert!(gate.can_launch());
gate.launch_started(Some(1_000_000_000));
assert!(gate.can_launch(), "one of two slots free");
gate.launch_started(Some(1_000_000_000));
assert!(!gate.can_launch(), "both manual slots in use");
gate.release(Some(1_000_000_000));
assert!(gate.can_launch(), "slot freed after release");
}
#[test]
fn test_merge_globals_sums_endpoint_counters() {
let mut snap1 = UsageSnapshot::default();
let mut snap2 = UsageSnapshot::default();
snap1.endpoints.insert(
"a".to_owned(),
EndpointUsage {
calls: 1,
input_tokens: 100,
output_tokens: 50,
cached_input_tokens: 20,
cache_creation_input_tokens: 5,
cost: 0.01,
models: std::iter::once("m1".to_owned()).collect(),
},
);
snap2.endpoints.insert(
"a".to_owned(),
EndpointUsage {
calls: 2,
input_tokens: 200,
output_tokens: 100,
cached_input_tokens: 0,
cache_creation_input_tokens: 0,
cost: 0.02,
models: std::iter::once("m1".to_owned()).collect(),
},
);
snap2.endpoints.insert(
"b".to_owned(),
EndpointUsage {
calls: 1,
input_tokens: 10,
output_tokens: 5,
cached_input_tokens: 0,
cache_creation_input_tokens: 0,
cost: 0.001,
models: std::iter::once("m2".to_owned()).collect(),
},
);
let merged = merge_globals(&[snap1, snap2]);
assert_eq!(merged.endpoints.len(), 2);
let a = merged.endpoints.get("a").unwrap();
assert_eq!(a.calls, 3);
assert_eq!(a.input_tokens, 300);
assert_eq!(a.cached_input_tokens, 20);
assert_eq!(a.cache_creation_input_tokens, 5);
assert!((a.cost - 0.03).abs() < 0.0001);
let b = merged.endpoints.get("b").unwrap();
assert_eq!(b.calls, 1);
assert_eq!(merged.models, vec!["m1".to_owned(), "m2".to_owned()]);
}
#[test]
fn test_run_sh_sync_captures_output() {
let Some(out) = run_sh_sync("echo hello") else {
return; };
assert_eq!(out, "hello");
}
#[test]
fn test_run_scenarios_batch_emits_one_run_event_pair() {
use crate::reporting::{ColorMode, Level, Reporter};
use crate::scenario::{TestGroup, TestStep};
let id = std::process::id();
let log_path = std::env::temp_dir().join(format!("lbt-parallel-{id}.ndjson"));
let reporter = Arc::new(
Reporter::new(
Level::Error,
ColorMode::Never,
Some(&log_path),
None,
None,
false,
)
.ok()
.unwrap(),
);
let mut files: Vec<ScenarioFile> = Vec::new();
for (label, url) in [("a", "http://127.0.0.1:9/"), ("b", "http://127.0.0.1:9/")] {
files.push(ScenarioFile {
label: label.to_owned(),
config: ScenarioConfig::default(),
definitions: Vec::new(),
tests: vec![TestGroup {
name: label.to_owned(),
start_url: None,
auto_navigate: None,
base_url: None,
timeout_secs: Some(5),
browser_headless: Some(true),
viewport_width: None,
viewport_height: None,
budget: None,
endpoint: None,
steps: vec![TestStep::Navigate {
url: url.to_owned(),
wait_after_ms: None,
}],
}],
});
}
let run = crate::parallel::run_scenarios(
files,
RunOptions {
mode: ParallelMode::Manual(2),
reporter: Arc::clone(&reporter),
memory: None,
},
)
.ok()
.unwrap();
assert_eq!(
run.report.tests_passed + run.report.tests_failed,
2,
"both files contributed exactly one test"
);
reporter.finish().ok().unwrap();
let text = std::fs::read_to_string(&log_path).ok().unwrap();
let started = text
.lines()
.filter(|l| l.contains("\"type\":\"run_started\""))
.count();
let finished = text
.lines()
.filter(|l| l.contains("\"type\":\"run_finished\""))
.count();
assert_eq!(started, 1, "exactly one RunStarted for the batch");
assert_eq!(finished, 1, "exactly one RunFinished for the batch");
}
}