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