Skip to main content

llm_browser_testkit/
parallel.rs

1//! Concurrent execution of multiple scenario files.
2//!
3//! A single `llm-browser-testkit run a.toml b.toml c.toml` invocation runs
4//! each file on its own **isolated browser** (separate Chrome process, so
5//! cookies, localStorage, and other session state can never leak between
6//! files). Each file's own tests always run sequentially on that file's
7//! browser.
8//!
9//! # Concurrency control
10//!
11//! - `[config] concurrency_group = "name"` makes every file that declares
12//!   the same group mutually exclusive: they never execute at the same
13//!   time. A file without a group gets its own implicit group, so distinct
14//!   files run in parallel by default.
15//! - Concurrency is bounded by a [`ParallelMode`]: an exact manual count
16//!   (`--parallel N`) or an **auto** mode that adapts to the machine.
17//!
18//! # Auto scaling
19//!
20//! In auto mode the orchestrator never assumes a fixed per-browser memory
21//! footprint. It instead **learns** it by trial and error:
22//!
23//! - it probes available system memory once and keeps a running estimate of
24//!   how many bytes one browser actually holds (`available_before` −
25//!   `available_after` around each run), clamped to sane bounds;
26//! - after each successful file it raises the concurrency limit toward
27//!   `capacity = available / learned_footprint`, and otherwise ramps up one
28//!   at a time;
29//! - it **guards** launches: a worker will not start a new browser while
30//!   less than an absolute [`MIN_HEADROOM_BYTES`] is free, or there is no
31//!   room for one more browser of the learned footprint. It never uses a
32//!   fraction of `total` memory, so it stays correct in VMs/containers
33//!   where `total` reports the host's much larger RAM;
34//! - when a launch fails with an out-of-memory-style error it **halves the
35//!   limit** and **retries the file** with backoff, up to
36//!   [`MAX_LAUNCH_RETRIES`].
37//!
38//! Bounds come from `--parallel-min` / `--parallel-max` (default `1` /
39//! unlimited). When system memory cannot be measured the auto ceiling falls
40//! back to a conservative [`DEFAULT_AUTO_CEILING`].
41
42use 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
55/// Worker-thread cap used when `--parallel-max` is unlimited, so we never
56/// spawn an unbounded number of threads. The adaptive limit still throttles
57/// the number of *running browsers* below this, and the memory guard /
58/// learned footprint stop it before the machine is saturated.
59const DEFAULT_THREAD_CAP: usize = 64;
60
61/// Fallback auto ceiling when system memory cannot be measured (roughly "a
62/// handful of browsers", so the default is genuinely parallel), overridable
63/// with `--parallel-max`.
64const DEFAULT_AUTO_CEILING: usize = 8;
65
66/// Absolute memory floor (bytes) kept free before new browsers are gated.
67/// Deliberately NOT a fraction of `total` memory: on VMs/containers `total`
68/// often reports the host's RAM (hundreds of GB), so a percentage guard
69/// would never fire — or a small cgroup limit would be blocked forever
70/// against the huge host figure. An absolute floor + "room for one learned
71/// browser" is what actually prevents Chrome from `OOM`-ing the process.
72const MIN_HEADROOM_BYTES: u64 = 256 * 1024 * 1024; // 256 MiB
73
74/// How many times a memory-style browser-launch failure is retried before
75/// the file is reported as failed.
76const MAX_LAUNCH_RETRIES: u32 = 3;
77
78/// Base and per-retry step for the backoff before a retried launch.
79const BACKOFF_BASE_MS: u64 = 600;
80const BACKOFF_STEP_MS: u64 = 600;
81const MAX_BACKOFF_MS: u64 = 4000;
82
83/// Bounds for the learned per-browser footprint estimate (bytes), rejecting
84/// noise from pages that balloon or transient dips.
85const MIN_FOOTPRINT: f64 = 128.0 * 1024.0 * 1024.0; // 128 MiB
86const MAX_FOOTPRINT: f64 = 2.0 * 1024.0 * 1024.0 * 1024.0; // 2 GiB
87
88/// How concurrency is chosen for a batch of files.
89#[derive(Debug, Clone, Copy, PartialEq, Eq)]
90pub enum ParallelMode {
91    /// Run at exactly this many files concurrently (`--parallel N`).
92    Manual(usize),
93    /// Adapt to available memory and learned browser footprint, never
94    /// exceeding `max` (0 = unlimited) or dropping below `min`.
95    Auto {
96        /// Hard floor for concurrency (from `--parallel-min`).
97        min: u32,
98        /// Hard ceiling for concurrency; `0` = unlimited (from
99        /// `--parallel-max`).
100        max: u32,
101    },
102}
103
104/// Builds the [`ParallelMode`] from CLI flags. `--parallel N` selects manual
105/// mode; otherwise auto mode with the given min/max bounds.
106#[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/// Best-effort snapshot of total / available physical memory.
118#[derive(Debug, Clone, Copy, PartialEq, Eq)]
119pub struct MemoryInfo {
120    /// Total physical memory, bytes.
121    pub total: u64,
122    /// Currently available (free) memory, bytes.
123    pub available: u64,
124}
125
126/// Runs one scenario file, ready for the scheduler. CLI overrides applied
127/// and viewport matrix expanded by the caller.
128#[derive(Clone)]
129pub struct ScenarioFile {
130    /// Human-readable label (typically the file path) used in logs.
131    pub label: String,
132    /// Effective scenario configuration for this file.
133    pub config: ScenarioConfig,
134    /// Reusable assertion definitions from this file.
135    pub definitions: Vec<AssertDefinition>,
136    /// The (viewport-expanded) tests to run.
137    pub tests: Vec<TestGroup>,
138}
139
140impl ScenarioFile {
141    /// The concurrency key for this file: its declared `concurrency_group`,
142    /// or a unique per-file key so distinct files run in parallel by
143    /// default.
144    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/// Options controlling a parallel batch run.
153#[derive(Clone)]
154pub struct RunOptions {
155    /// How concurrency is chosen (manual or auto-scaling).
156    pub mode: ParallelMode,
157    /// Shared event sink (thread-safe: every sink is mutex-guarded).
158    pub reporter: Arc<Reporter>,
159    /// Best-effort memory snapshot; `None` falls back to trial-and-error.
160    pub memory: Option<MemoryInfo>,
161}
162
163/// Result of a parallel batch run: the merged report plus combined usage.
164#[derive(Debug, Default)]
165pub struct ParallelRun {
166    /// Merged test/step results across all files.
167    pub report: RunReport,
168    /// Per-test usage snapshots across all files.
169    pub per_test: Vec<(String, UsageSnapshot)>,
170    /// Per-file usage snapshots (file label → usage), in batch order.
171    pub per_file: Vec<(String, UsageSnapshot)>,
172    /// Combined global usage across all files.
173    pub global: UsageSnapshot,
174}
175
176/// Runs every file concurrently, respecting the concurrency mode and
177/// per-file groups, each on its own isolated browser. Emits one
178/// `RunStarted`/`RunFinished` pair for the batch and merges cost/usage.
179///
180/// # Errors
181///
182/// Returns an error only for internal scheduling failures; a file whose
183/// browser cannot be launched (after retries) is recorded as a failed
184/// report rather than aborting the batch.
185#[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        // Shadow with per-iteration clones so each closure captures its own
228        // Arc (the loop variable itself must not move into the closure).
229        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    // Merge per-file reports + usage while holding the scheduler lock, then
241    // release it (end of the block) before emitting the batch RunFinished.
242    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
290// ── Adaptive gate ─────────────────────────────────────────────────────
291
292/// The adaptive controller shared by every worker. All fields are guarded
293/// by the `Mutex` the workers lock around it.
294struct Gate {
295    /// Number of browsers currently running.
296    active: usize,
297    /// Current concurrency limit; moves in `[lower, threads]` as the
298    /// controller learns the machine's memory behaviour.
299    limit: usize,
300    /// Hard floor for the limit (from `--parallel-min`).
301    lower: usize,
302    /// Number of worker threads; also the ceiling the limit may reach.
303    threads: usize,
304    /// Ceiling for the limit when memory is unknown (ramp target).
305    ramp_ceiling: usize,
306    /// Total physical memory, if probed.
307    memory_total: Option<u64>,
308    /// Last known available memory, if probed.
309    available: Option<u64>,
310    /// Learned average per-browser footprint, bytes (0 = unknown).
311    footprint: f64,
312    /// Number of footprint samples folded into the average.
313    footprint_count: u32,
314    /// Current backoff (ms) applied before retrying a failed launch.
315    retry_backoff: u64,
316}
317
318/// Computes the initial gate state from the mode and (optional) memory.
319fn 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    /// Whether a new browser may be launched right now: below the
363    /// concurrency limit and with enough memory headroom.
364    fn can_launch(&self) -> bool {
365        self.active < self.limit && !self.memory_blocked()
366    }
367
368    /// True when a new browser should not be launched yet: either less than an
369    /// absolute [`MIN_HEADROOM_BYTES`] is free, or there is not room for one
370    /// more browser of the learned footprint. Uses absolute free bytes (never a
371    /// fraction of `total`), so it stays correct inside VMs/containers where
372    /// `total` may report the host's much larger RAM.
373    #[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    /// Reserves a launch slot.
394    #[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    /// Folds the memory delta around one run into the learned footprint and
403    /// refreshes the last-known available memory.
404    #[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    /// After a successful run, try to raise the limit toward the
432    /// memory-derived capacity, or ramp up one at a time when memory is
433    /// unknown.
434    #[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    /// After an out-of-memory-style launch failure, back off: halve the
454    /// limit (never below the floor) and increase the retry backoff.
455    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    /// Frees a launch slot and refreshes known available memory.
461    #[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
470// ── Scheduler ─────────────────────────────────────────────────────────
471
472/// Outcome of [`SchedulerState::claim`].
473enum Claim {
474    /// File `usize` is assigned to this worker.
475    Take(usize),
476    /// Work remains but every unstarted file is blocked by a concurrency
477    /// group — retry shortly.
478    Wait,
479    /// Every file has been started; no more work.
480    Done,
481}
482
483/// Shared scheduler state, guarded by the mutex every worker locks.
484struct SchedulerState {
485    files: Vec<ScenarioFile>,
486    results: Vec<Option<ParallelRun>>,
487    started: Vec<bool>,
488    active_groups: HashSet<String>,
489}
490
491impl SchedulerState {
492    /// Reserves the first unstarted file whose concurrency group is not
493    /// currently running, or reports that the worker should wait / that all
494    /// work is done.
495    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    /// Releases file `i`'s concurrency group so blocked files can start.
516    fn complete(&mut self, i: usize) {
517        let key = self.files[i].concurrency_key(i);
518        self.active_groups.remove(&key);
519    }
520}
521
522/// Per-file run decision, computed while holding the gate.
523enum Decision {
524    /// Success; return the run.
525    Done(ParallelRun),
526    /// Memory-style failure to retry after `u64` ms.
527    Retry(u64),
528    /// Non-retryable failure; record as failed.
529    GiveUp(String),
530}
531
532/// Worker loop: repeatedly claim a runnable file, run it on its own
533/// isolated browser (gated + retried), and record the result.
534#[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/// Runs one file, waiting for a launch slot, and retrying memory-style
557/// launch failures with backoff.
558#[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        // Wait for a launch slot (concurrency limit + memory headroom). The
570        // gate guard drops at the end of the inner block, before we sleep.
571        {
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
627/// Runs one scenario file on a fresh isolated browser. Returns the run on
628/// success, or the launch error on failure (so the caller can classify and
629/// retry).
630fn 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/// Whether a browser-launch error looks like memory exhaustion and is worth
651/// retrying after backing off.
652#[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/// A report that fails every test in `file` (used when the browser cannot be
669/// launched after retries).
670#[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/// Merges a batch of per-file global usage snapshots into one, summing
689/// per-endpoint counters so the cost report is accurate across files.
690#[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// ── Memory probing ────────────────────────────────────────────────────
709
710/// Best-effort probe of total/available physical memory (Linux then macOS).
711/// Returns `None` when unavailable; the scheduler then relies on trial and
712/// error.
713#[must_use]
714#[allow(clippy::unnecessary_unwrap)]
715pub async fn probe_memory_async() -> Option<MemoryInfo> {
716    // Linux: `free -b` prints "Mem:  <total> <used> <free> <shared> <buff> <cache> <available>".
717    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    // macOS: total bytes from sysctl; free pages × page size (from sysctl,
731    // or vm_stat). `Pages free` is printed with a trailing period, so strip
732    // non-digits.
733    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/// Synchronous best-effort probe of currently available memory, used by
752/// workers to learn the per-browser footprint and refresh the headroom
753/// guard. Returns `None` when unavailable (e.g. non-Linux workers).
754#[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/// Runs a `sh -c` script asynchronously and returns trimmed stdout.
761#[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/// Runs a `sh -c` script synchronously and returns the first trimmed line.
775#[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        // Unknown memory: ceiling falls back to DEFAULT_AUTO_CEILING (8).
900        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, // below MIN_HEADROOM_BYTES
918            }),
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        // VM/container: total reports the host's huge RAM, but available is
942        // small. The guard must block based on absolute available bytes, not
943        // a fraction of total.
944        let gate = make_gate(
945            &ParallelMode::Auto { min: 1, max: 4 },
946            Some(crate::parallel::MemoryInfo {
947                total: 500_000_000_000, // 500 GB host
948                available: 100_000_000, // 100 MB left in the cgroup
949            }),
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, // above the absolute floor...
962            }),
963            4,
964        );
965        // ...but one learned browser needs ~500 MB, so it must still block.
966        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        // Before = 1 GiB free, after = 600 MiB free → one browser ≈ 400 MiB.
1008        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        // Learn a 1 GiB footprint, then capacity = 4 GiB / 1 GiB = 4.
1026        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; // sh unavailable on this platform — nothing to assert
1109        };
1110        assert_eq!(out, "hello");
1111    }
1112
1113    /// Runs the real worker pool against two failing scenario files and
1114    /// asserts the batch machinery: both files' tests are accounted for, and
1115    /// exactly one `RunStarted` + one `RunFinished` are emitted (each file
1116    /// suppresses its own, so the batch owns the pair). Uses an unreachable
1117    /// URL so files fail fast with or without Chrome installed.
1118    #[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        // Both files' single test must be accounted for in the merged report.
1175        assert_eq!(
1176            run.report.tests_passed + run.report.tests_failed,
1177            2,
1178            "both files contributed exactly one test"
1179        );
1180
1181        // Exactly one batch-level RunStarted and one RunFinished.
1182        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}