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