1use std::collections::{HashMap, HashSet};
43use std::io::BufRead;
44use std::sync::{Arc, Mutex};
45use std::time::Duration;
46
47use anyhow::Result;
48
49use crate::costs::{EndpointUsage, UsageSnapshot};
50use crate::events::TestEvent;
51use crate::reporting::Reporter;
52use crate::runner::{RunReport, ScenarioRunner};
53use crate::scenario::{AssertDefinition, ScenarioConfig, TestGroup};
54
55const DEFAULT_THREAD_CAP: usize = 64;
60
61const DEFAULT_AUTO_CEILING: usize = 8;
65
66const MIN_HEADROOM_BYTES: u64 = 256 * 1024 * 1024; const MAX_LAUNCH_RETRIES: u32 = 3;
77
78const BACKOFF_BASE_MS: u64 = 600;
80const BACKOFF_STEP_MS: u64 = 600;
81const MAX_BACKOFF_MS: u64 = 4000;
82
83const 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)]
90pub enum ParallelMode {
91 Manual(usize),
93 Auto {
96 min: u32,
98 max: u32,
101 },
102}
103
104#[must_use]
107pub fn mode_from_cli(parallel: Option<u32>, min: u32, max: u32) -> ParallelMode {
108 parallel.map_or_else(
109 || ParallelMode::Auto {
110 min: min.max(1),
111 max,
112 },
113 |k| ParallelMode::Manual(k.max(1) as usize),
114 )
115}
116
117#[derive(Debug, Clone, Copy, PartialEq, Eq)]
119pub struct MemoryInfo {
120 pub total: u64,
122 pub available: u64,
124}
125
126#[derive(Clone)]
129pub struct ScenarioFile {
130 pub label: String,
132 pub config: ScenarioConfig,
134 pub definitions: Vec<AssertDefinition>,
136 pub tests: Vec<TestGroup>,
138}
139
140impl ScenarioFile {
141 fn concurrency_key(&self, index: usize) -> String {
145 self.config
146 .concurrency_group
147 .clone()
148 .unwrap_or_else(|| format!("<file {index}>"))
149 }
150}
151
152#[derive(Clone)]
154pub struct RunOptions {
155 pub mode: ParallelMode,
157 pub reporter: Arc<Reporter>,
159 pub memory: Option<MemoryInfo>,
161}
162
163#[derive(Debug, Default)]
165pub struct ParallelRun {
166 pub report: RunReport,
168 pub per_test: Vec<(String, UsageSnapshot)>,
170 pub global: UsageSnapshot,
172}
173
174#[allow(clippy::cast_possible_truncation, clippy::significant_drop_tightening)]
184pub fn run_scenarios(files: Vec<ScenarioFile>, opts: RunOptions) -> Result<ParallelRun> {
185 let reporter = opts.reporter;
186 let total_tests: u32 = files.iter().map(|f| f.tests.len() as u32).sum();
187
188 reporter.emit(&TestEvent::RunStarted { total_tests })?;
189
190 if files.is_empty() {
191 reporter.emit(&TestEvent::RunFinished {
192 tests_passed: 0,
193 tests_failed: 0,
194 steps_passed: 0,
195 steps_failed: 0,
196 steps_skipped: 0,
197 total_cost: 0.0,
198 total_tokens: 0,
199 total_input_tokens: 0,
200 total_output_tokens: 0,
201 total_cached_input_tokens: 0,
202 models: Vec::new(),
203 total_calls: 0,
204 })?;
205 return Ok(ParallelRun::default());
206 }
207
208 let n = files.len();
209 let mut results: Vec<Option<ParallelRun>> = Vec::new();
210 for _ in 0..n {
211 results.push(None);
212 }
213 let state = Arc::new(Mutex::new(SchedulerState {
214 files,
215 results,
216 started: vec![false; n],
217 active_groups: HashSet::new(),
218 }));
219 let gate = Arc::new(Mutex::new(make_gate(&opts.mode, opts.memory, n)));
220 let workers = gate.lock().unwrap().threads;
221
222 let mut threads: Vec<std::thread::JoinHandle<()>> = Vec::new();
223 for _ in 0..workers {
224 let state = Arc::clone(&state);
227 let gate = Arc::clone(&gate);
228 let reporter = Arc::clone(&reporter);
229 threads.push(std::thread::spawn(move || {
230 worker_loop(&state, &gate, &reporter);
231 }));
232 }
233 for thread in threads {
234 thread.join().unwrap();
235 }
236
237 let run = {
240 let st = state.lock().unwrap();
241 let mut report = RunReport::default();
242 let mut per_test: Vec<(String, UsageSnapshot)> = Vec::new();
243 let mut globals: Vec<UsageSnapshot> = Vec::new();
244 for r in &st.results {
245 if let Some(run) = r.as_ref() {
246 report.tests_passed += run.report.tests_passed;
247 report.tests_failed += run.report.tests_failed;
248 report.passed += run.report.passed;
249 report.failed += run.report.failed;
250 report.skipped += run.report.skipped;
251 report.details.extend(run.report.details.clone());
252 per_test.extend(run.per_test.clone());
253 globals.push(run.global.clone());
254 }
255 }
256 ParallelRun {
257 report,
258 per_test,
259 global: merge_globals(&globals),
260 }
261 };
262
263 reporter.emit(&TestEvent::RunFinished {
264 tests_passed: run.report.tests_passed,
265 tests_failed: run.report.tests_failed,
266 steps_passed: run.report.passed,
267 steps_failed: run.report.failed,
268 steps_skipped: run.report.skipped,
269 total_cost: run.global.total_cost,
270 total_tokens: run.global.total_tokens,
271 total_input_tokens: run.global.total_input_tokens,
272 total_output_tokens: run.global.total_output_tokens,
273 total_cached_input_tokens: run.global.total_cached_input_tokens,
274 models: run.global.models.clone(),
275 total_calls: run.global.total_calls,
276 })?;
277
278 Ok(run)
279}
280
281struct Gate {
286 active: usize,
288 limit: usize,
291 lower: usize,
293 threads: usize,
295 ramp_ceiling: usize,
297 memory_total: Option<u64>,
299 available: Option<u64>,
301 footprint: f64,
303 footprint_count: u32,
305 retry_backoff: u64,
307}
308
309fn make_gate(mode: &ParallelMode, memory: Option<MemoryInfo>, files_len: usize) -> Gate {
311 match mode {
312 ParallelMode::Manual(k) => {
313 let k = *k;
314 let threads = k.max(1).min(files_len).max(1);
315 let limit = k.max(1).min(threads);
316 Gate {
317 active: 0,
318 limit,
319 lower: limit,
320 threads,
321 ramp_ceiling: limit,
322 memory_total: memory.map(|m| m.total),
323 available: memory.map(|m| m.available),
324 footprint: 0.0,
325 footprint_count: 0,
326 retry_backoff: 0,
327 }
328 }
329 ParallelMode::Auto { min, max } => {
330 let min = *min;
331 let max = *max;
332 let lower = min.max(1) as usize;
333 let user_max: Option<usize> = if max == 0 { None } else { Some(max as usize) };
334 let threads = user_max.unwrap_or(DEFAULT_THREAD_CAP).min(files_len).max(1);
335 let ramp_ceiling = user_max.unwrap_or(DEFAULT_AUTO_CEILING).min(threads).max(1);
336 Gate {
337 active: 0,
338 limit: lower.min(threads).max(1),
339 lower: lower.min(threads).max(1),
340 threads,
341 ramp_ceiling,
342 memory_total: memory.map(|m| m.total),
343 available: memory.map(|m| m.available),
344 footprint: 0.0,
345 footprint_count: 0,
346 retry_backoff: 0,
347 }
348 }
349 }
350}
351
352impl Gate {
353 fn can_launch(&self) -> bool {
356 self.active < self.limit && !self.memory_blocked()
357 }
358
359 #[allow(
365 clippy::unnecessary_unwrap,
366 clippy::cast_precision_loss,
367 clippy::cast_possible_truncation,
368 clippy::cast_sign_loss
369 )]
370 fn memory_blocked(&self) -> bool {
371 let Some(available) = self.available else {
372 return false;
373 };
374 if available < MIN_HEADROOM_BYTES {
375 return true;
376 }
377 if self.footprint_count > 0 && self.footprint > 0.0 {
378 available < self.footprint as u64
379 } else {
380 false
381 }
382 }
383
384 #[allow(clippy::missing_const_for_fn)]
386 fn launch_started(&mut self, before: Option<u64>) {
387 if before.is_some() {
388 self.available = before;
389 }
390 self.active += 1;
391 }
392
393 #[allow(
396 clippy::unnecessary_unwrap,
397 clippy::cast_precision_loss,
398 clippy::suboptimal_flops
399 )]
400 fn launch_finished(&mut self, before: Option<u64>, after: Option<u64>) {
401 if before.is_some() {
402 self.available = before;
403 }
404 if after.is_some() {
405 self.available = after;
406 }
407 if before.is_some() && after.is_some() {
408 let delta = before.unwrap().saturating_sub(after.unwrap());
409 if delta > 0 {
410 let d = delta as f64;
411 if self.footprint_count == 0 {
412 self.footprint = d;
413 } else {
414 self.footprint = self.footprint * 0.7 + d * 0.3;
415 }
416 self.footprint = self.footprint.clamp(MIN_FOOTPRINT, MAX_FOOTPRINT);
417 self.footprint_count = (self.footprint_count + 1).min(20);
418 }
419 }
420 }
421
422 #[allow(
426 clippy::unnecessary_unwrap,
427 clippy::cast_precision_loss,
428 clippy::cast_possible_truncation,
429 clippy::cast_sign_loss
430 )]
431 fn on_success(&mut self) {
432 if self.footprint_count > 0 && self.footprint > 0.0 && self.memory_total.is_some() {
433 let total = self.memory_total.unwrap();
434 let available = self.available.unwrap_or(total);
435 if available > 0 {
436 let capacity = (available as f64 / self.footprint) as usize;
437 self.limit = capacity.max(self.lower).min(self.threads);
438 return;
439 }
440 }
441 self.limit = (self.limit + 1).min(self.ramp_ceiling).max(self.lower);
442 }
443
444 fn on_launch_failure(&mut self) {
447 self.limit = (self.limit / 2).max(self.lower);
448 self.retry_backoff = (self.retry_backoff + BACKOFF_STEP_MS).min(MAX_BACKOFF_MS);
449 }
450
451 #[allow(clippy::missing_const_for_fn)]
453 fn release(&mut self, after: Option<u64>) {
454 if after.is_some() {
455 self.available = after;
456 }
457 self.active = self.active.saturating_sub(1);
458 }
459}
460
461enum Claim {
465 Take(usize),
467 Wait,
470 Done,
472}
473
474struct SchedulerState {
476 files: Vec<ScenarioFile>,
477 results: Vec<Option<ParallelRun>>,
478 started: Vec<bool>,
479 active_groups: HashSet<String>,
480}
481
482impl SchedulerState {
483 fn claim(&mut self) -> Claim {
487 for i in 0..self.files.len() {
488 if self.started[i] {
489 continue;
490 }
491 let key = self.files[i].concurrency_key(i);
492 if self.active_groups.contains(&key) {
493 continue;
494 }
495 self.started[i] = true;
496 self.active_groups.insert(key);
497 return Claim::Take(i);
498 }
499 if self.started.iter().all(|b| *b) {
500 Claim::Done
501 } else {
502 Claim::Wait
503 }
504 }
505
506 fn complete(&mut self, i: usize) {
508 let key = self.files[i].concurrency_key(i);
509 self.active_groups.remove(&key);
510 }
511}
512
513enum Decision {
515 Done(ParallelRun),
517 Retry(u64),
519 GiveUp(String),
521}
522
523#[allow(clippy::significant_drop_tightening)]
526fn worker_loop(
527 state: &Arc<Mutex<SchedulerState>>,
528 gate: &Arc<Mutex<Gate>>,
529 reporter: &Arc<Reporter>,
530) {
531 loop {
532 let claimed = state.lock().unwrap().claim();
533 match claimed {
534 Claim::Done => break,
535 Claim::Wait => std::thread::sleep(Duration::from_millis(25)),
536 Claim::Take(i) => {
537 let file = state.lock().unwrap().files[i].clone();
538 let run = run_file_with_retries(&file, gate, reporter);
539 let mut st = state.lock().unwrap();
540 st.results[i] = Some(run);
541 st.complete(i);
542 }
543 }
544 }
545}
546
547#[allow(clippy::significant_drop_tightening)]
550fn run_file_with_retries(
551 file: &ScenarioFile,
552 gate: &Arc<Mutex<Gate>>,
553 reporter: &Arc<Reporter>,
554) -> ParallelRun {
555 let mut attempt: u32 = 0;
556 loop {
557 attempt += 1;
558 let before = available_memory_now();
559
560 {
563 loop {
564 let mut g = gate.lock().unwrap();
565 if g.can_launch() {
566 g.launch_started(before);
567 break;
568 }
569 std::thread::sleep(Duration::from_millis(25));
570 }
571 }
572
573 let result = run_one_file(file, reporter);
574 let after = available_memory_now();
575
576 let decision = {
577 let mut g = gate.lock().unwrap();
578 g.launch_finished(before, after);
579 match &result {
580 Ok(_) => {
581 g.on_success();
582 g.release(after);
583 Decision::Done(result.unwrap())
584 }
585 Err(e) => {
586 let oom = is_retryable_oom(e);
587 if oom && attempt < MAX_LAUNCH_RETRIES {
588 let backoff = g.retry_backoff.max(BACKOFF_BASE_MS);
589 g.on_launch_failure();
590 g.release(after);
591 Decision::Retry(backoff)
592 } else {
593 g.on_launch_failure();
594 g.release(after);
595 Decision::GiveUp(e.clone())
596 }
597 }
598 }
599 };
600
601 match decision {
602 Decision::Done(run) => return run,
603 Decision::Retry(ms) => {
604 reporter.warn(format!(
605 "{}: browser launch failed (likely out of memory); retrying ({attempt}/{MAX_LAUNCH_RETRIES}) in {ms}ms",
606 file.label,
607 ));
608 std::thread::sleep(Duration::from_millis(ms));
609 }
610 Decision::GiveUp(e) => {
611 reporter.error(format!("{}: {e}", file.label));
612 return synthesized_failed_report(file);
613 }
614 }
615 }
616}
617
618fn run_one_file(file: &ScenarioFile, reporter: &Arc<Reporter>) -> Result<ParallelRun, String> {
622 let runner = ScenarioRunner::with_reporter_parallel(
623 file.config.clone(),
624 file.definitions.clone(),
625 Arc::clone(reporter),
626 );
627 match runner.run(&file.tests) {
628 Ok(report) => {
629 let usage = runner.usage_tracker();
630 Ok(ParallelRun {
631 report,
632 per_test: usage.per_test_snapshots(),
633 global: usage.global_snapshot(),
634 })
635 }
636 Err(e) => Err(e.to_string()),
637 }
638}
639
640#[must_use]
643fn is_retryable_oom(err: &str) -> bool {
644 const KEYWORDS: [&str; 8] = [
645 "memory",
646 "cannot allocate",
647 "out of memory",
648 "killed",
649 "oom",
650 "resource temporarily unavailable",
651 "failed to allocate",
652 "no memory",
653 ];
654 let e = err.to_lowercase();
655 KEYWORDS.iter().any(|kw| e.contains(kw))
656}
657
658#[allow(clippy::cast_possible_truncation)]
661fn synthesized_failed_report(file: &ScenarioFile) -> ParallelRun {
662 let n = file.tests.len() as u32;
663 ParallelRun {
664 report: RunReport {
665 tests_passed: 0,
666 tests_failed: n,
667 passed: 0,
668 failed: n,
669 skipped: 0,
670 details: Vec::new(),
671 },
672 per_test: Vec::new(),
673 global: UsageSnapshot::default(),
674 }
675}
676
677#[must_use]
680fn merge_globals(snapshots: &[UsageSnapshot]) -> UsageSnapshot {
681 let mut endpoints: HashMap<String, EndpointUsage> = HashMap::new();
682 for snapshot in snapshots {
683 for (name, usage) in &snapshot.endpoints {
684 let acc = endpoints.entry(name.clone()).or_default();
685 acc.calls += usage.calls;
686 acc.input_tokens += usage.input_tokens;
687 acc.output_tokens += usage.output_tokens;
688 acc.cached_input_tokens += usage.cached_input_tokens;
689 acc.cost += usage.cost;
690 acc.models.extend(usage.models.iter().cloned());
691 }
692 }
693 UsageSnapshot::from_endpoints(&endpoints)
694}
695
696#[must_use]
702#[allow(clippy::unnecessary_unwrap)]
703pub async fn probe_memory_async() -> Option<MemoryInfo> {
704 let total = run_sh_async("free -b | awk '/^Mem:/{print $2}'").await;
706 let available = run_sh_async("free -b | awk '/^Mem:/{print $7}'").await;
707 if total.is_some() && available.is_some() {
708 let t = total.unwrap().trim().parse::<u64>().ok();
709 let a = available.unwrap().trim().parse::<u64>().ok();
710 if t.is_some() && a.is_some() {
711 return Some(MemoryInfo {
712 total: t.unwrap(),
713 available: a.unwrap(),
714 });
715 }
716 }
717
718 let total_mac = run_sh_async("sysctl -n hw.memsize").await;
722 let page_size = run_sh_async("sysctl -n hw.pagesize").await;
723 let pages =
724 run_sh_async("vm_stat | awk '/^Pages free:/{gsub(/[^0-9]/, \"\", $3); print $3}'").await;
725 if total_mac.is_some() && pages.is_some() && page_size.is_some() {
726 let t = total_mac.unwrap().trim().parse::<u64>().ok();
727 let p = pages.unwrap().trim().parse::<u64>().ok();
728 let ps = page_size.unwrap().trim().parse::<u64>().ok();
729 if t.is_some() && p.is_some() && ps.is_some() {
730 return Some(MemoryInfo {
731 total: t.unwrap(),
732 available: p.unwrap() * ps.unwrap(),
733 });
734 }
735 }
736 None
737}
738
739#[must_use]
743fn available_memory_now() -> Option<u64> {
744 run_sh_sync("free -b 2>/dev/null | awk '/^Mem:/{print $7}'")
745 .and_then(|s| s.trim().parse::<u64>().ok())
746}
747
748#[must_use]
750async fn run_sh_async(script: &str) -> Option<String> {
751 let output = tokio::process::Command::new("sh")
752 .args(["-c", script])
753 .output()
754 .await
755 .ok()?;
756 if !output.status.success() {
757 return None;
758 }
759 Some(String::from_utf8_lossy(&output.stdout).trim().to_string())
760}
761
762#[must_use]
764fn run_sh_sync(script: &str) -> Option<String> {
765 let mut child = std::process::Command::new("sh")
766 .args(["-c", script])
767 .stdout(std::process::Stdio::piped())
768 .stderr(std::process::Stdio::inherit())
769 .spawn()
770 .ok()?;
771 let stdout = child.stdout.as_mut()?;
772 let mut reader = std::io::BufReader::new(stdout);
773 let mut line = String::new();
774 let _ = reader.read_line(&mut line);
775 if line.is_empty() {
776 None
777 } else {
778 Some(line.trim().to_string())
779 }
780}
781
782#[cfg(test)]
783mod tests {
784 use std::sync::Arc;
785
786 use crate::costs::{EndpointUsage, UsageSnapshot};
787 use crate::scenario::ScenarioConfig;
788
789 use super::{
790 is_retryable_oom, make_gate, merge_globals, mode_from_cli, run_sh_sync, Claim,
791 ParallelMode, RunOptions, ScenarioFile, SchedulerState,
792 };
793
794 fn file(label: &str, group: Option<&str>) -> ScenarioFile {
795 ScenarioFile {
796 label: label.to_owned(),
797 config: ScenarioConfig {
798 concurrency_group: group.map(std::borrow::ToOwned::to_owned),
799 ..ScenarioConfig::default()
800 },
801 definitions: Vec::new(),
802 tests: Vec::new(),
803 }
804 }
805
806 fn state(files: Vec<ScenarioFile>) -> SchedulerState {
807 let n = files.len();
808 SchedulerState {
809 files,
810 results: Vec::new(),
811 started: vec![false; n],
812 active_groups: std::collections::HashSet::new(),
813 }
814 }
815
816 #[test]
817 fn test_distinct_groups_claim_in_parallel() {
818 let mut s = state(vec![file("a", None), file("b", Some("x")), file("c", None)]);
819 assert!(matches!(s.claim(), Claim::Take(0)));
820 assert!(matches!(s.claim(), Claim::Take(1)));
821 assert!(matches!(s.claim(), Claim::Take(2)));
822 s.complete(1);
823 assert!(matches!(s.claim(), Claim::Done));
824 }
825
826 #[test]
827 fn test_same_group_blocks_until_completed() {
828 let mut s = state(vec![
829 file("a", Some("g")),
830 file("b", Some("g")),
831 file("c", None),
832 ]);
833 assert!(matches!(s.claim(), Claim::Take(0)));
834 assert!(
835 matches!(s.claim(), Claim::Take(2)),
836 "a different group still runs while 'g' is active"
837 );
838 assert!(
839 matches!(s.claim(), Claim::Wait),
840 "b is blocked by a's group"
841 );
842 s.complete(0);
843 assert!(
844 matches!(s.claim(), Claim::Take(1)),
845 "b runs after a finishes"
846 );
847 s.complete(1);
848 assert!(matches!(s.claim(), Claim::Done));
849 }
850
851 #[test]
852 fn test_no_group_means_own_group() {
853 let mut s = state(vec![file("a", None), file("b", None)]);
854 assert!(matches!(s.claim(), Claim::Take(0)));
855 assert!(
856 matches!(s.claim(), Claim::Take(1)),
857 "no-group files run in parallel"
858 );
859 }
860
861 #[test]
862 fn test_mode_from_cli_manual_wins() {
863 assert_eq!(mode_from_cli(Some(5), 1, 0), ParallelMode::Manual(5));
864 }
865
866 #[test]
867 fn test_mode_from_cli_auto_with_defaults() {
868 assert_eq!(
869 mode_from_cli(None, 1, 0),
870 ParallelMode::Auto { min: 1, max: 0 }
871 );
872 }
873
874 #[test]
875 fn test_gate_manual_is_fixed() {
876 let gate = make_gate(&ParallelMode::Manual(3), None, 10);
877 assert_eq!(gate.limit, 3);
878 assert_eq!(gate.threads, 3);
879 assert_eq!(gate.lower, 3);
880 }
881
882 #[test]
883 fn test_gate_auto_ramps_from_min_toward_ceiling() {
884 let gate = make_gate(&ParallelMode::Auto { min: 1, max: 0 }, None, 100);
885 assert_eq!(gate.lower, 1);
886 assert_eq!(gate.limit, 1);
887 assert_eq!(gate.ramp_ceiling, 8);
889 assert_eq!(gate.threads, 64);
890 }
891
892 #[test]
893 fn test_gate_auto_respects_user_max() {
894 let gate = make_gate(&ParallelMode::Auto { min: 1, max: 8 }, None, 100);
895 assert_eq!(gate.threads, 8);
896 assert_eq!(gate.ramp_ceiling, 8);
897 }
898
899 #[test]
900 fn test_gate_memory_guard_blocks_low_headroom() {
901 let gate = make_gate(
902 &ParallelMode::Auto { min: 1, max: 4 },
903 Some(crate::parallel::MemoryInfo {
904 total: 1_000_000_000,
905 available: 50_000_000, }),
907 4,
908 );
909 assert!(gate.memory_blocked());
910 assert!(!gate.can_launch());
911 }
912
913 #[test]
914 fn test_gate_memory_guard_allows_headroom() {
915 let gate = make_gate(
916 &ParallelMode::Auto { min: 1, max: 4 },
917 Some(crate::parallel::MemoryInfo {
918 total: 1_000_000_000,
919 available: 900_000_000,
920 }),
921 4,
922 );
923 assert!(!gate.memory_blocked());
924 assert!(gate.can_launch());
925 }
926
927 #[test]
928 fn test_gate_memory_guard_ignores_huge_host_total() {
929 let gate = make_gate(
933 &ParallelMode::Auto { min: 1, max: 4 },
934 Some(crate::parallel::MemoryInfo {
935 total: 500_000_000_000, available: 100_000_000, }),
938 4,
939 );
940 assert!(gate.memory_blocked(), "must block despite a 500 GB 'total'");
941 }
942
943 #[test]
944 fn test_gate_memory_guard_blocks_when_no_room_for_one_footprint() {
945 let mut gate = make_gate(
946 &ParallelMode::Auto { min: 1, max: 4 },
947 Some(crate::parallel::MemoryInfo {
948 total: 8_000_000_000,
949 available: 300_000_000, }),
951 4,
952 );
953 gate.footprint = 500.0 * 1024.0 * 1024.0;
955 gate.footprint_count = 3;
956 assert!(gate.memory_blocked());
957 }
958
959 #[test]
960 fn test_oom_failure_halves_limit_and_sets_backoff() {
961 let mut gate = make_gate(&ParallelMode::Auto { min: 1, max: 16 }, None, 100);
962 gate.limit = 16;
963 gate.on_launch_failure();
964 assert_eq!(gate.limit, 8);
965 assert!(gate.retry_backoff > 0);
966 }
967
968 #[test]
969 fn test_is_retryable_oom_matches_memory_errors() {
970 assert!(is_retryable_oom("failed to launch browser: out of memory"));
971 assert!(is_retryable_oom("cannot allocate memory for page"));
972 assert!(!is_retryable_oom("Chrome binary not found"));
973 assert!(!is_retryable_oom("invalid URL"));
974 }
975
976 #[test]
977 fn test_run_options_cloneable() {
978 let _ = RunOptions {
979 mode: ParallelMode::Manual(2),
980 reporter: Arc::new(crate::reporting::Reporter::default()),
981 memory: None,
982 };
983 }
984
985 #[test]
986 fn test_launch_finished_learns_footprint_from_delta() {
987 let mut gate = make_gate(
988 &ParallelMode::Auto { min: 1, max: 4 },
989 Some(crate::parallel::MemoryInfo {
990 total: 8_000_000_000,
991 available: 8_000_000_000,
992 }),
993 4,
994 );
995 gate.launch_finished(Some(1_000_000_000), Some(600_000_000));
997 assert!((gate.footprint - 400_000_000.0).abs() < 1.0);
998 assert_eq!(gate.footprint_count, 1);
999 assert_eq!(gate.available, Some(600_000_000));
1000 }
1001
1002 #[test]
1003 fn test_on_success_raises_to_memory_capacity() {
1004 let mut gate = make_gate(
1005 &ParallelMode::Auto { min: 1, max: 0 },
1006 Some(crate::parallel::MemoryInfo {
1007 total: 8_000_000_000,
1008 available: 4_000_000_000,
1009 }),
1010 100,
1011 );
1012 assert_eq!(gate.limit, 1);
1013 gate.footprint = 1_000_000_000.0;
1015 gate.footprint_count = 3;
1016 gate.on_success();
1017 assert_eq!(gate.limit, 4);
1018 }
1019
1020 #[test]
1021 fn test_on_success_ramps_when_memory_unknown() {
1022 let mut gate = make_gate(&ParallelMode::Auto { min: 1, max: 0 }, None, 100);
1023 assert_eq!(gate.limit, 1);
1024 gate.on_success();
1025 assert_eq!(gate.limit, 2, "ramps up one at a time");
1026 }
1027
1028 #[test]
1029 fn test_can_launch_respects_active_limit() {
1030 let mut gate = make_gate(&ParallelMode::Manual(2), None, 10);
1031 assert!(gate.can_launch());
1032 gate.launch_started(Some(1_000_000_000));
1033 assert!(gate.can_launch(), "one of two slots free");
1034 gate.launch_started(Some(1_000_000_000));
1035 assert!(!gate.can_launch(), "both manual slots in use");
1036 gate.release(Some(1_000_000_000));
1037 assert!(gate.can_launch(), "slot freed after release");
1038 }
1039
1040 #[test]
1041 fn test_merge_globals_sums_endpoint_counters() {
1042 let mut snap1 = UsageSnapshot::default();
1043 let mut snap2 = UsageSnapshot::default();
1044 snap1.endpoints.insert(
1045 "a".to_owned(),
1046 EndpointUsage {
1047 calls: 1,
1048 input_tokens: 100,
1049 output_tokens: 50,
1050 cached_input_tokens: 20,
1051 cost: 0.01,
1052 models: std::iter::once("m1".to_owned()).collect(),
1053 },
1054 );
1055 snap2.endpoints.insert(
1056 "a".to_owned(),
1057 EndpointUsage {
1058 calls: 2,
1059 input_tokens: 200,
1060 output_tokens: 100,
1061 cached_input_tokens: 0,
1062 cost: 0.02,
1063 models: std::iter::once("m1".to_owned()).collect(),
1064 },
1065 );
1066 snap2.endpoints.insert(
1067 "b".to_owned(),
1068 EndpointUsage {
1069 calls: 1,
1070 input_tokens: 10,
1071 output_tokens: 5,
1072 cached_input_tokens: 0,
1073 cost: 0.001,
1074 models: std::iter::once("m2".to_owned()).collect(),
1075 },
1076 );
1077 let merged = merge_globals(&[snap1, snap2]);
1078 assert_eq!(merged.endpoints.len(), 2);
1079 let a = merged.endpoints.get("a").unwrap();
1080 assert_eq!(a.calls, 3);
1081 assert_eq!(a.input_tokens, 300);
1082 assert_eq!(a.cached_input_tokens, 20);
1083 assert!((a.cost - 0.03).abs() < 0.0001);
1084 let b = merged.endpoints.get("b").unwrap();
1085 assert_eq!(b.calls, 1);
1086 assert_eq!(merged.models, vec!["m1".to_owned(), "m2".to_owned()]);
1087 }
1088
1089 #[test]
1090 fn test_run_sh_sync_captures_output() {
1091 let Some(out) = run_sh_sync("echo hello") else {
1092 return; };
1094 assert_eq!(out, "hello");
1095 }
1096
1097 #[test]
1103 fn test_run_scenarios_batch_emits_one_run_event_pair() {
1104 use crate::reporting::{ColorMode, Level, Reporter};
1105 use crate::scenario::{TestGroup, TestStep};
1106
1107 let id = std::process::id();
1108 let log_path = std::env::temp_dir().join(format!("lbt-parallel-{id}.ndjson"));
1109 let reporter = Arc::new(
1110 Reporter::new(
1111 Level::Error,
1112 ColorMode::Never,
1113 Some(&log_path),
1114 None,
1115 None,
1116 false,
1117 )
1118 .ok()
1119 .unwrap(),
1120 );
1121
1122 let mut files: Vec<ScenarioFile> = Vec::new();
1123 for (label, url) in [("a", "http://127.0.0.1:9/"), ("b", "http://127.0.0.1:9/")] {
1124 files.push(ScenarioFile {
1125 label: label.to_owned(),
1126 config: ScenarioConfig::default(),
1127 definitions: Vec::new(),
1128 tests: vec![TestGroup {
1129 name: label.to_owned(),
1130 start_url: None,
1131 auto_navigate: None,
1132 base_url: None,
1133 timeout_secs: Some(5),
1134 browser_headless: Some(true),
1135 viewport_width: None,
1136 viewport_height: None,
1137 budget: None,
1138 endpoint: None,
1139 steps: vec![TestStep::Navigate {
1140 url: url.to_owned(),
1141 wait_after_ms: None,
1142 }],
1143 }],
1144 });
1145 }
1146
1147 let run = crate::parallel::run_scenarios(
1148 files,
1149 RunOptions {
1150 mode: ParallelMode::Manual(2),
1151 reporter: Arc::clone(&reporter),
1152 memory: None,
1153 },
1154 )
1155 .ok()
1156 .unwrap();
1157
1158 assert_eq!(
1160 run.report.tests_passed + run.report.tests_failed,
1161 2,
1162 "both files contributed exactly one test"
1163 );
1164
1165 reporter.finish().ok().unwrap();
1167 let text = std::fs::read_to_string(&log_path).ok().unwrap();
1168 let started = text
1169 .lines()
1170 .filter(|l| l.contains("\"type\":\"run_started\""))
1171 .count();
1172 let finished = text
1173 .lines()
1174 .filter(|l| l.contains("\"type\":\"run_finished\""))
1175 .count();
1176 assert_eq!(started, 1, "exactly one RunStarted for the batch");
1177 assert_eq!(finished, 1, "exactly one RunFinished for the batch");
1178 }
1179}