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    /// Combined global usage across all files.
171    pub global: UsageSnapshot,
172}
173
174/// Runs every file concurrently, respecting the concurrency mode and
175/// per-file groups, each on its own isolated browser. Emits one
176/// `RunStarted`/`RunFinished` pair for the batch and merges cost/usage.
177///
178/// # Errors
179///
180/// Returns an error only for internal scheduling failures; a file whose
181/// browser cannot be launched (after retries) is recorded as a failed
182/// report rather than aborting the batch.
183#[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        // Shadow with per-iteration clones so each closure captures its own
226        // Arc (the loop variable itself must not move into the closure).
227        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    // Merge per-file reports + usage while holding the scheduler lock, then
239    // release it (end of the block) before emitting the batch RunFinished.
240    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
283// ── Adaptive gate ─────────────────────────────────────────────────────
284
285/// The adaptive controller shared by every worker. All fields are guarded
286/// by the `Mutex` the workers lock around it.
287struct Gate {
288    /// Number of browsers currently running.
289    active: usize,
290    /// Current concurrency limit; moves in `[lower, threads]` as the
291    /// controller learns the machine's memory behaviour.
292    limit: usize,
293    /// Hard floor for the limit (from `--parallel-min`).
294    lower: usize,
295    /// Number of worker threads; also the ceiling the limit may reach.
296    threads: usize,
297    /// Ceiling for the limit when memory is unknown (ramp target).
298    ramp_ceiling: usize,
299    /// Total physical memory, if probed.
300    memory_total: Option<u64>,
301    /// Last known available memory, if probed.
302    available: Option<u64>,
303    /// Learned average per-browser footprint, bytes (0 = unknown).
304    footprint: f64,
305    /// Number of footprint samples folded into the average.
306    footprint_count: u32,
307    /// Current backoff (ms) applied before retrying a failed launch.
308    retry_backoff: u64,
309}
310
311/// Computes the initial gate state from the mode and (optional) memory.
312fn 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    /// Whether a new browser may be launched right now: below the
356    /// concurrency limit and with enough memory headroom.
357    fn can_launch(&self) -> bool {
358        self.active < self.limit && !self.memory_blocked()
359    }
360
361    /// True when a new browser should not be launched yet: either less than an
362    /// absolute [`MIN_HEADROOM_BYTES`] is free, or there is not room for one
363    /// more browser of the learned footprint. Uses absolute free bytes (never a
364    /// fraction of `total`), so it stays correct inside VMs/containers where
365    /// `total` may report the host's much larger RAM.
366    #[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    /// Reserves a launch slot.
387    #[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    /// Folds the memory delta around one run into the learned footprint and
396    /// refreshes the last-known available memory.
397    #[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    /// After a successful run, try to raise the limit toward the
425    /// memory-derived capacity, or ramp up one at a time when memory is
426    /// unknown.
427    #[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    /// After an out-of-memory-style launch failure, back off: halve the
447    /// limit (never below the floor) and increase the retry backoff.
448    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    /// Frees a launch slot and refreshes known available memory.
454    #[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
463// ── Scheduler ─────────────────────────────────────────────────────────
464
465/// Outcome of [`SchedulerState::claim`].
466enum Claim {
467    /// File `usize` is assigned to this worker.
468    Take(usize),
469    /// Work remains but every unstarted file is blocked by a concurrency
470    /// group — retry shortly.
471    Wait,
472    /// Every file has been started; no more work.
473    Done,
474}
475
476/// Shared scheduler state, guarded by the mutex every worker locks.
477struct SchedulerState {
478    files: Vec<ScenarioFile>,
479    results: Vec<Option<ParallelRun>>,
480    started: Vec<bool>,
481    active_groups: HashSet<String>,
482}
483
484impl SchedulerState {
485    /// Reserves the first unstarted file whose concurrency group is not
486    /// currently running, or reports that the worker should wait / that all
487    /// work is done.
488    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    /// Releases file `i`'s concurrency group so blocked files can start.
509    fn complete(&mut self, i: usize) {
510        let key = self.files[i].concurrency_key(i);
511        self.active_groups.remove(&key);
512    }
513}
514
515/// Per-file run decision, computed while holding the gate.
516enum Decision {
517    /// Success; return the run.
518    Done(ParallelRun),
519    /// Memory-style failure to retry after `u64` ms.
520    Retry(u64),
521    /// Non-retryable failure; record as failed.
522    GiveUp(String),
523}
524
525/// Worker loop: repeatedly claim a runnable file, run it on its own
526/// isolated browser (gated + retried), and record the result.
527#[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/// Runs one file, waiting for a launch slot, and retrying memory-style
550/// launch failures with backoff.
551#[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        // Wait for a launch slot (concurrency limit + memory headroom). The
563        // gate guard drops at the end of the inner block, before we sleep.
564        {
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
620/// Runs one scenario file on a fresh isolated browser. Returns the run on
621/// success, or the launch error on failure (so the caller can classify and
622/// retry).
623fn 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/// Whether a browser-launch error looks like memory exhaustion and is worth
643/// retrying after backing off.
644#[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/// A report that fails every test in `file` (used when the browser cannot be
661/// launched after retries).
662#[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/// Merges a batch of per-file global usage snapshots into one, summing
680/// per-endpoint counters so the cost report is accurate across files.
681#[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// ── Memory probing ────────────────────────────────────────────────────
700
701/// Best-effort probe of total/available physical memory (Linux then macOS).
702/// Returns `None` when unavailable; the scheduler then relies on trial and
703/// error.
704#[must_use]
705#[allow(clippy::unnecessary_unwrap)]
706pub async fn probe_memory_async() -> Option<MemoryInfo> {
707    // Linux: `free -b` prints "Mem:  <total> <used> <free> <shared> <buff> <cache> <available>".
708    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    // macOS: total bytes from sysctl; free pages × page size (from sysctl,
722    // or vm_stat). `Pages free` is printed with a trailing period, so strip
723    // non-digits.
724    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/// Synchronous best-effort probe of currently available memory, used by
743/// workers to learn the per-browser footprint and refresh the headroom
744/// guard. Returns `None` when unavailable (e.g. non-Linux workers).
745#[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/// Runs a `sh -c` script asynchronously and returns trimmed stdout.
752#[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/// Runs a `sh -c` script synchronously and returns the first trimmed line.
766#[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        // Unknown memory: ceiling falls back to DEFAULT_AUTO_CEILING (8).
891        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, // below MIN_HEADROOM_BYTES
909            }),
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        // VM/container: total reports the host's huge RAM, but available is
933        // small. The guard must block based on absolute available bytes, not
934        // a fraction of total.
935        let gate = make_gate(
936            &ParallelMode::Auto { min: 1, max: 4 },
937            Some(crate::parallel::MemoryInfo {
938                total: 500_000_000_000, // 500 GB host
939                available: 100_000_000, // 100 MB left in the cgroup
940            }),
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, // above the absolute floor...
953            }),
954            4,
955        );
956        // ...but one learned browser needs ~500 MB, so it must still block.
957        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        // Before = 1 GiB free, after = 600 MiB free → one browser ≈ 400 MiB.
999        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        // Learn a 1 GiB footprint, then capacity = 4 GiB / 1 GiB = 4.
1017        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; // sh unavailable on this platform — nothing to assert
1100        };
1101        assert_eq!(out, "hello");
1102    }
1103
1104    /// Runs the real worker pool against two failing scenario files and
1105    /// asserts the batch machinery: both files' tests are accounted for, and
1106    /// exactly one `RunStarted` + one `RunFinished` are emitted (each file
1107    /// suppresses its own, so the batch owns the pair). Uses an unreachable
1108    /// URL so files fail fast with or without Chrome installed.
1109    #[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        // Both files' single test must be accounted for in the merged report.
1166        assert_eq!(
1167            run.report.tests_passed + run.report.tests_failed,
1168            2,
1169            "both files contributed exactly one test"
1170        );
1171
1172        // Exactly one batch-level RunStarted and one RunFinished.
1173        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}