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