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            models: Vec::new(),
203            total_calls: 0,
204        })?;
205        return Ok(ParallelRun::default());
206    }
207
208    let n = files.len();
209    let mut results: Vec<Option<ParallelRun>> = Vec::new();
210    for _ in 0..n {
211        results.push(None);
212    }
213    let state = Arc::new(Mutex::new(SchedulerState {
214        files,
215        results,
216        started: vec![false; n],
217        active_groups: HashSet::new(),
218    }));
219    let gate = Arc::new(Mutex::new(make_gate(&opts.mode, opts.memory, n)));
220    let workers = gate.lock().unwrap().threads;
221
222    let mut threads: Vec<std::thread::JoinHandle<()>> = Vec::new();
223    for _ in 0..workers {
224        // Shadow with per-iteration clones so each closure captures its own
225        // Arc (the loop variable itself must not move into the closure).
226        let state = Arc::clone(&state);
227        let gate = Arc::clone(&gate);
228        let reporter = Arc::clone(&reporter);
229        threads.push(std::thread::spawn(move || {
230            worker_loop(&state, &gate, &reporter);
231        }));
232    }
233    for thread in threads {
234        thread.join().unwrap();
235    }
236
237    // Merge per-file reports + usage while holding the scheduler lock, then
238    // release it (end of the block) before emitting the batch RunFinished.
239    let run = {
240        let st = state.lock().unwrap();
241        let mut report = RunReport::default();
242        let mut per_test: Vec<(String, UsageSnapshot)> = Vec::new();
243        let mut globals: Vec<UsageSnapshot> = Vec::new();
244        for r in &st.results {
245            if let Some(run) = r.as_ref() {
246                report.tests_passed += run.report.tests_passed;
247                report.tests_failed += run.report.tests_failed;
248                report.passed += run.report.passed;
249                report.failed += run.report.failed;
250                report.skipped += run.report.skipped;
251                report.details.extend(run.report.details.clone());
252                per_test.extend(run.per_test.clone());
253                globals.push(run.global.clone());
254            }
255        }
256        ParallelRun {
257            report,
258            per_test,
259            global: merge_globals(&globals),
260        }
261    };
262
263    reporter.emit(&TestEvent::RunFinished {
264        tests_passed: run.report.tests_passed,
265        tests_failed: run.report.tests_failed,
266        steps_passed: run.report.passed,
267        steps_failed: run.report.failed,
268        steps_skipped: run.report.skipped,
269        total_cost: run.global.total_cost,
270        total_tokens: run.global.total_tokens,
271        total_input_tokens: run.global.total_input_tokens,
272        total_output_tokens: run.global.total_output_tokens,
273        total_cached_input_tokens: run.global.total_cached_input_tokens,
274        models: run.global.models.clone(),
275        total_calls: run.global.total_calls,
276    })?;
277
278    Ok(run)
279}
280
281// ── Adaptive gate ─────────────────────────────────────────────────────
282
283/// The adaptive controller shared by every worker. All fields are guarded
284/// by the `Mutex` the workers lock around it.
285struct Gate {
286    /// Number of browsers currently running.
287    active: usize,
288    /// Current concurrency limit; moves in `[lower, threads]` as the
289    /// controller learns the machine's memory behaviour.
290    limit: usize,
291    /// Hard floor for the limit (from `--parallel-min`).
292    lower: usize,
293    /// Number of worker threads; also the ceiling the limit may reach.
294    threads: usize,
295    /// Ceiling for the limit when memory is unknown (ramp target).
296    ramp_ceiling: usize,
297    /// Total physical memory, if probed.
298    memory_total: Option<u64>,
299    /// Last known available memory, if probed.
300    available: Option<u64>,
301    /// Learned average per-browser footprint, bytes (0 = unknown).
302    footprint: f64,
303    /// Number of footprint samples folded into the average.
304    footprint_count: u32,
305    /// Current backoff (ms) applied before retrying a failed launch.
306    retry_backoff: u64,
307}
308
309/// Computes the initial gate state from the mode and (optional) memory.
310fn make_gate(mode: &ParallelMode, memory: Option<MemoryInfo>, files_len: usize) -> Gate {
311    match mode {
312        ParallelMode::Manual(k) => {
313            let k = *k;
314            let threads = k.max(1).min(files_len).max(1);
315            let limit = k.max(1).min(threads);
316            Gate {
317                active: 0,
318                limit,
319                lower: limit,
320                threads,
321                ramp_ceiling: limit,
322                memory_total: memory.map(|m| m.total),
323                available: memory.map(|m| m.available),
324                footprint: 0.0,
325                footprint_count: 0,
326                retry_backoff: 0,
327            }
328        }
329        ParallelMode::Auto { min, max } => {
330            let min = *min;
331            let max = *max;
332            let lower = min.max(1) as usize;
333            let user_max: Option<usize> = if max == 0 { None } else { Some(max as usize) };
334            let threads = user_max.unwrap_or(DEFAULT_THREAD_CAP).min(files_len).max(1);
335            let ramp_ceiling = user_max.unwrap_or(DEFAULT_AUTO_CEILING).min(threads).max(1);
336            Gate {
337                active: 0,
338                limit: lower.min(threads).max(1),
339                lower: lower.min(threads).max(1),
340                threads,
341                ramp_ceiling,
342                memory_total: memory.map(|m| m.total),
343                available: memory.map(|m| m.available),
344                footprint: 0.0,
345                footprint_count: 0,
346                retry_backoff: 0,
347            }
348        }
349    }
350}
351
352impl Gate {
353    /// Whether a new browser may be launched right now: below the
354    /// concurrency limit and with enough memory headroom.
355    fn can_launch(&self) -> bool {
356        self.active < self.limit && !self.memory_blocked()
357    }
358
359    /// True when a new browser should not be launched yet: either less than an
360    /// absolute [`MIN_HEADROOM_BYTES`] is free, or there is not room for one
361    /// more browser of the learned footprint. Uses absolute free bytes (never a
362    /// fraction of `total`), so it stays correct inside VMs/containers where
363    /// `total` may report the host's much larger RAM.
364    #[allow(
365        clippy::unnecessary_unwrap,
366        clippy::cast_precision_loss,
367        clippy::cast_possible_truncation,
368        clippy::cast_sign_loss
369    )]
370    fn memory_blocked(&self) -> bool {
371        let Some(available) = self.available else {
372            return false;
373        };
374        if available < MIN_HEADROOM_BYTES {
375            return true;
376        }
377        if self.footprint_count > 0 && self.footprint > 0.0 {
378            available < self.footprint as u64
379        } else {
380            false
381        }
382    }
383
384    /// Reserves a launch slot.
385    #[allow(clippy::missing_const_for_fn)]
386    fn launch_started(&mut self, before: Option<u64>) {
387        if before.is_some() {
388            self.available = before;
389        }
390        self.active += 1;
391    }
392
393    /// Folds the memory delta around one run into the learned footprint and
394    /// refreshes the last-known available memory.
395    #[allow(
396        clippy::unnecessary_unwrap,
397        clippy::cast_precision_loss,
398        clippy::suboptimal_flops
399    )]
400    fn launch_finished(&mut self, before: Option<u64>, after: Option<u64>) {
401        if before.is_some() {
402            self.available = before;
403        }
404        if after.is_some() {
405            self.available = after;
406        }
407        if before.is_some() && after.is_some() {
408            let delta = before.unwrap().saturating_sub(after.unwrap());
409            if delta > 0 {
410                let d = delta as f64;
411                if self.footprint_count == 0 {
412                    self.footprint = d;
413                } else {
414                    self.footprint = self.footprint * 0.7 + d * 0.3;
415                }
416                self.footprint = self.footprint.clamp(MIN_FOOTPRINT, MAX_FOOTPRINT);
417                self.footprint_count = (self.footprint_count + 1).min(20);
418            }
419        }
420    }
421
422    /// After a successful run, try to raise the limit toward the
423    /// memory-derived capacity, or ramp up one at a time when memory is
424    /// unknown.
425    #[allow(
426        clippy::unnecessary_unwrap,
427        clippy::cast_precision_loss,
428        clippy::cast_possible_truncation,
429        clippy::cast_sign_loss
430    )]
431    fn on_success(&mut self) {
432        if self.footprint_count > 0 && self.footprint > 0.0 && self.memory_total.is_some() {
433            let total = self.memory_total.unwrap();
434            let available = self.available.unwrap_or(total);
435            if available > 0 {
436                let capacity = (available as f64 / self.footprint) as usize;
437                self.limit = capacity.max(self.lower).min(self.threads);
438                return;
439            }
440        }
441        self.limit = (self.limit + 1).min(self.ramp_ceiling).max(self.lower);
442    }
443
444    /// After an out-of-memory-style launch failure, back off: halve the
445    /// limit (never below the floor) and increase the retry backoff.
446    fn on_launch_failure(&mut self) {
447        self.limit = (self.limit / 2).max(self.lower);
448        self.retry_backoff = (self.retry_backoff + BACKOFF_STEP_MS).min(MAX_BACKOFF_MS);
449    }
450
451    /// Frees a launch slot and refreshes known available memory.
452    #[allow(clippy::missing_const_for_fn)]
453    fn release(&mut self, after: Option<u64>) {
454        if after.is_some() {
455            self.available = after;
456        }
457        self.active = self.active.saturating_sub(1);
458    }
459}
460
461// ── Scheduler ─────────────────────────────────────────────────────────
462
463/// Outcome of [`SchedulerState::claim`].
464enum Claim {
465    /// File `usize` is assigned to this worker.
466    Take(usize),
467    /// Work remains but every unstarted file is blocked by a concurrency
468    /// group — retry shortly.
469    Wait,
470    /// Every file has been started; no more work.
471    Done,
472}
473
474/// Shared scheduler state, guarded by the mutex every worker locks.
475struct SchedulerState {
476    files: Vec<ScenarioFile>,
477    results: Vec<Option<ParallelRun>>,
478    started: Vec<bool>,
479    active_groups: HashSet<String>,
480}
481
482impl SchedulerState {
483    /// Reserves the first unstarted file whose concurrency group is not
484    /// currently running, or reports that the worker should wait / that all
485    /// work is done.
486    fn claim(&mut self) -> Claim {
487        for i in 0..self.files.len() {
488            if self.started[i] {
489                continue;
490            }
491            let key = self.files[i].concurrency_key(i);
492            if self.active_groups.contains(&key) {
493                continue;
494            }
495            self.started[i] = true;
496            self.active_groups.insert(key);
497            return Claim::Take(i);
498        }
499        if self.started.iter().all(|b| *b) {
500            Claim::Done
501        } else {
502            Claim::Wait
503        }
504    }
505
506    /// Releases file `i`'s concurrency group so blocked files can start.
507    fn complete(&mut self, i: usize) {
508        let key = self.files[i].concurrency_key(i);
509        self.active_groups.remove(&key);
510    }
511}
512
513/// Per-file run decision, computed while holding the gate.
514enum Decision {
515    /// Success; return the run.
516    Done(ParallelRun),
517    /// Memory-style failure to retry after `u64` ms.
518    Retry(u64),
519    /// Non-retryable failure; record as failed.
520    GiveUp(String),
521}
522
523/// Worker loop: repeatedly claim a runnable file, run it on its own
524/// isolated browser (gated + retried), and record the result.
525#[allow(clippy::significant_drop_tightening)]
526fn worker_loop(
527    state: &Arc<Mutex<SchedulerState>>,
528    gate: &Arc<Mutex<Gate>>,
529    reporter: &Arc<Reporter>,
530) {
531    loop {
532        let claimed = state.lock().unwrap().claim();
533        match claimed {
534            Claim::Done => break,
535            Claim::Wait => std::thread::sleep(Duration::from_millis(25)),
536            Claim::Take(i) => {
537                let file = state.lock().unwrap().files[i].clone();
538                let run = run_file_with_retries(&file, gate, reporter);
539                let mut st = state.lock().unwrap();
540                st.results[i] = Some(run);
541                st.complete(i);
542            }
543        }
544    }
545}
546
547/// Runs one file, waiting for a launch slot, and retrying memory-style
548/// launch failures with backoff.
549#[allow(clippy::significant_drop_tightening)]
550fn run_file_with_retries(
551    file: &ScenarioFile,
552    gate: &Arc<Mutex<Gate>>,
553    reporter: &Arc<Reporter>,
554) -> ParallelRun {
555    let mut attempt: u32 = 0;
556    loop {
557        attempt += 1;
558        let before = available_memory_now();
559
560        // Wait for a launch slot (concurrency limit + memory headroom). The
561        // gate guard drops at the end of the inner block, before we sleep.
562        {
563            loop {
564                let mut g = gate.lock().unwrap();
565                if g.can_launch() {
566                    g.launch_started(before);
567                    break;
568                }
569                std::thread::sleep(Duration::from_millis(25));
570            }
571        }
572
573        let result = run_one_file(file, reporter);
574        let after = available_memory_now();
575
576        let decision = {
577            let mut g = gate.lock().unwrap();
578            g.launch_finished(before, after);
579            match &result {
580                Ok(_) => {
581                    g.on_success();
582                    g.release(after);
583                    Decision::Done(result.unwrap())
584                }
585                Err(e) => {
586                    let oom = is_retryable_oom(e);
587                    if oom && attempt < MAX_LAUNCH_RETRIES {
588                        let backoff = g.retry_backoff.max(BACKOFF_BASE_MS);
589                        g.on_launch_failure();
590                        g.release(after);
591                        Decision::Retry(backoff)
592                    } else {
593                        g.on_launch_failure();
594                        g.release(after);
595                        Decision::GiveUp(e.clone())
596                    }
597                }
598            }
599        };
600
601        match decision {
602            Decision::Done(run) => return run,
603            Decision::Retry(ms) => {
604                reporter.warn(format!(
605                    "{}: browser launch failed (likely out of memory); retrying ({attempt}/{MAX_LAUNCH_RETRIES}) in {ms}ms",
606                    file.label,
607                ));
608                std::thread::sleep(Duration::from_millis(ms));
609            }
610            Decision::GiveUp(e) => {
611                reporter.error(format!("{}: {e}", file.label));
612                return synthesized_failed_report(file);
613            }
614        }
615    }
616}
617
618/// Runs one scenario file on a fresh isolated browser. Returns the run on
619/// success, or the launch error on failure (so the caller can classify and
620/// retry).
621fn run_one_file(file: &ScenarioFile, reporter: &Arc<Reporter>) -> Result<ParallelRun, String> {
622    let runner = ScenarioRunner::with_reporter_parallel(
623        file.config.clone(),
624        file.definitions.clone(),
625        Arc::clone(reporter),
626    );
627    match runner.run(&file.tests) {
628        Ok(report) => {
629            let usage = runner.usage_tracker();
630            Ok(ParallelRun {
631                report,
632                per_test: usage.per_test_snapshots(),
633                global: usage.global_snapshot(),
634            })
635        }
636        Err(e) => Err(e.to_string()),
637    }
638}
639
640/// Whether a browser-launch error looks like memory exhaustion and is worth
641/// retrying after backing off.
642#[must_use]
643fn is_retryable_oom(err: &str) -> bool {
644    const KEYWORDS: [&str; 8] = [
645        "memory",
646        "cannot allocate",
647        "out of memory",
648        "killed",
649        "oom",
650        "resource temporarily unavailable",
651        "failed to allocate",
652        "no memory",
653    ];
654    let e = err.to_lowercase();
655    KEYWORDS.iter().any(|kw| e.contains(kw))
656}
657
658/// A report that fails every test in `file` (used when the browser cannot be
659/// launched after retries).
660#[allow(clippy::cast_possible_truncation)]
661fn synthesized_failed_report(file: &ScenarioFile) -> ParallelRun {
662    let n = file.tests.len() as u32;
663    ParallelRun {
664        report: RunReport {
665            tests_passed: 0,
666            tests_failed: n,
667            passed: 0,
668            failed: n,
669            skipped: 0,
670            details: Vec::new(),
671        },
672        per_test: Vec::new(),
673        global: UsageSnapshot::default(),
674    }
675}
676
677/// Merges a batch of per-file global usage snapshots into one, summing
678/// per-endpoint counters so the cost report is accurate across files.
679#[must_use]
680fn merge_globals(snapshots: &[UsageSnapshot]) -> UsageSnapshot {
681    let mut endpoints: HashMap<String, EndpointUsage> = HashMap::new();
682    for snapshot in snapshots {
683        for (name, usage) in &snapshot.endpoints {
684            let acc = endpoints.entry(name.clone()).or_default();
685            acc.calls += usage.calls;
686            acc.input_tokens += usage.input_tokens;
687            acc.output_tokens += usage.output_tokens;
688            acc.cached_input_tokens += usage.cached_input_tokens;
689            acc.cost += usage.cost;
690            acc.models.extend(usage.models.iter().cloned());
691        }
692    }
693    UsageSnapshot::from_endpoints(&endpoints)
694}
695
696// ── Memory probing ────────────────────────────────────────────────────
697
698/// Best-effort probe of total/available physical memory (Linux then macOS).
699/// Returns `None` when unavailable; the scheduler then relies on trial and
700/// error.
701#[must_use]
702#[allow(clippy::unnecessary_unwrap)]
703pub async fn probe_memory_async() -> Option<MemoryInfo> {
704    // Linux: `free -b` prints "Mem:  <total> <used> <free> <shared> <buff> <cache> <available>".
705    let total = run_sh_async("free -b | awk '/^Mem:/{print $2}'").await;
706    let available = run_sh_async("free -b | awk '/^Mem:/{print $7}'").await;
707    if total.is_some() && available.is_some() {
708        let t = total.unwrap().trim().parse::<u64>().ok();
709        let a = available.unwrap().trim().parse::<u64>().ok();
710        if t.is_some() && a.is_some() {
711            return Some(MemoryInfo {
712                total: t.unwrap(),
713                available: a.unwrap(),
714            });
715        }
716    }
717
718    // macOS: total bytes from sysctl; free pages × page size (from sysctl,
719    // or vm_stat). `Pages free` is printed with a trailing period, so strip
720    // non-digits.
721    let total_mac = run_sh_async("sysctl -n hw.memsize").await;
722    let page_size = run_sh_async("sysctl -n hw.pagesize").await;
723    let pages =
724        run_sh_async("vm_stat | awk '/^Pages free:/{gsub(/[^0-9]/, \"\", $3); print $3}'").await;
725    if total_mac.is_some() && pages.is_some() && page_size.is_some() {
726        let t = total_mac.unwrap().trim().parse::<u64>().ok();
727        let p = pages.unwrap().trim().parse::<u64>().ok();
728        let ps = page_size.unwrap().trim().parse::<u64>().ok();
729        if t.is_some() && p.is_some() && ps.is_some() {
730            return Some(MemoryInfo {
731                total: t.unwrap(),
732                available: p.unwrap() * ps.unwrap(),
733            });
734        }
735    }
736    None
737}
738
739/// Synchronous best-effort probe of currently available memory, used by
740/// workers to learn the per-browser footprint and refresh the headroom
741/// guard. Returns `None` when unavailable (e.g. non-Linux workers).
742#[must_use]
743fn available_memory_now() -> Option<u64> {
744    run_sh_sync("free -b 2>/dev/null | awk '/^Mem:/{print $7}'")
745        .and_then(|s| s.trim().parse::<u64>().ok())
746}
747
748/// Runs a `sh -c` script asynchronously and returns trimmed stdout.
749#[must_use]
750async fn run_sh_async(script: &str) -> Option<String> {
751    let output = tokio::process::Command::new("sh")
752        .args(["-c", script])
753        .output()
754        .await
755        .ok()?;
756    if !output.status.success() {
757        return None;
758    }
759    Some(String::from_utf8_lossy(&output.stdout).trim().to_string())
760}
761
762/// Runs a `sh -c` script synchronously and returns the first trimmed line.
763#[must_use]
764fn run_sh_sync(script: &str) -> Option<String> {
765    let mut child = std::process::Command::new("sh")
766        .args(["-c", script])
767        .stdout(std::process::Stdio::piped())
768        .stderr(std::process::Stdio::inherit())
769        .spawn()
770        .ok()?;
771    let stdout = child.stdout.as_mut()?;
772    let mut reader = std::io::BufReader::new(stdout);
773    let mut line = String::new();
774    let _ = reader.read_line(&mut line);
775    if line.is_empty() {
776        None
777    } else {
778        Some(line.trim().to_string())
779    }
780}
781
782#[cfg(test)]
783mod tests {
784    use std::sync::Arc;
785
786    use crate::costs::{EndpointUsage, UsageSnapshot};
787    use crate::scenario::ScenarioConfig;
788
789    use super::{
790        is_retryable_oom, make_gate, merge_globals, mode_from_cli, run_sh_sync, Claim,
791        ParallelMode, RunOptions, ScenarioFile, SchedulerState,
792    };
793
794    fn file(label: &str, group: Option<&str>) -> ScenarioFile {
795        ScenarioFile {
796            label: label.to_owned(),
797            config: ScenarioConfig {
798                concurrency_group: group.map(std::borrow::ToOwned::to_owned),
799                ..ScenarioConfig::default()
800            },
801            definitions: Vec::new(),
802            tests: Vec::new(),
803        }
804    }
805
806    fn state(files: Vec<ScenarioFile>) -> SchedulerState {
807        let n = files.len();
808        SchedulerState {
809            files,
810            results: Vec::new(),
811            started: vec![false; n],
812            active_groups: std::collections::HashSet::new(),
813        }
814    }
815
816    #[test]
817    fn test_distinct_groups_claim_in_parallel() {
818        let mut s = state(vec![file("a", None), file("b", Some("x")), file("c", None)]);
819        assert!(matches!(s.claim(), Claim::Take(0)));
820        assert!(matches!(s.claim(), Claim::Take(1)));
821        assert!(matches!(s.claim(), Claim::Take(2)));
822        s.complete(1);
823        assert!(matches!(s.claim(), Claim::Done));
824    }
825
826    #[test]
827    fn test_same_group_blocks_until_completed() {
828        let mut s = state(vec![
829            file("a", Some("g")),
830            file("b", Some("g")),
831            file("c", None),
832        ]);
833        assert!(matches!(s.claim(), Claim::Take(0)));
834        assert!(
835            matches!(s.claim(), Claim::Take(2)),
836            "a different group still runs while 'g' is active"
837        );
838        assert!(
839            matches!(s.claim(), Claim::Wait),
840            "b is blocked by a's group"
841        );
842        s.complete(0);
843        assert!(
844            matches!(s.claim(), Claim::Take(1)),
845            "b runs after a finishes"
846        );
847        s.complete(1);
848        assert!(matches!(s.claim(), Claim::Done));
849    }
850
851    #[test]
852    fn test_no_group_means_own_group() {
853        let mut s = state(vec![file("a", None), file("b", None)]);
854        assert!(matches!(s.claim(), Claim::Take(0)));
855        assert!(
856            matches!(s.claim(), Claim::Take(1)),
857            "no-group files run in parallel"
858        );
859    }
860
861    #[test]
862    fn test_mode_from_cli_manual_wins() {
863        assert_eq!(mode_from_cli(Some(5), 1, 0), ParallelMode::Manual(5));
864    }
865
866    #[test]
867    fn test_mode_from_cli_auto_with_defaults() {
868        assert_eq!(
869            mode_from_cli(None, 1, 0),
870            ParallelMode::Auto { min: 1, max: 0 }
871        );
872    }
873
874    #[test]
875    fn test_gate_manual_is_fixed() {
876        let gate = make_gate(&ParallelMode::Manual(3), None, 10);
877        assert_eq!(gate.limit, 3);
878        assert_eq!(gate.threads, 3);
879        assert_eq!(gate.lower, 3);
880    }
881
882    #[test]
883    fn test_gate_auto_ramps_from_min_toward_ceiling() {
884        let gate = make_gate(&ParallelMode::Auto { min: 1, max: 0 }, None, 100);
885        assert_eq!(gate.lower, 1);
886        assert_eq!(gate.limit, 1);
887        // Unknown memory: ceiling falls back to DEFAULT_AUTO_CEILING (8).
888        assert_eq!(gate.ramp_ceiling, 8);
889        assert_eq!(gate.threads, 64);
890    }
891
892    #[test]
893    fn test_gate_auto_respects_user_max() {
894        let gate = make_gate(&ParallelMode::Auto { min: 1, max: 8 }, None, 100);
895        assert_eq!(gate.threads, 8);
896        assert_eq!(gate.ramp_ceiling, 8);
897    }
898
899    #[test]
900    fn test_gate_memory_guard_blocks_low_headroom() {
901        let gate = make_gate(
902            &ParallelMode::Auto { min: 1, max: 4 },
903            Some(crate::parallel::MemoryInfo {
904                total: 1_000_000_000,
905                available: 50_000_000, // below MIN_HEADROOM_BYTES
906            }),
907            4,
908        );
909        assert!(gate.memory_blocked());
910        assert!(!gate.can_launch());
911    }
912
913    #[test]
914    fn test_gate_memory_guard_allows_headroom() {
915        let gate = make_gate(
916            &ParallelMode::Auto { min: 1, max: 4 },
917            Some(crate::parallel::MemoryInfo {
918                total: 1_000_000_000,
919                available: 900_000_000,
920            }),
921            4,
922        );
923        assert!(!gate.memory_blocked());
924        assert!(gate.can_launch());
925    }
926
927    #[test]
928    fn test_gate_memory_guard_ignores_huge_host_total() {
929        // VM/container: total reports the host's huge RAM, but available is
930        // small. The guard must block based on absolute available bytes, not
931        // a fraction of total.
932        let gate = make_gate(
933            &ParallelMode::Auto { min: 1, max: 4 },
934            Some(crate::parallel::MemoryInfo {
935                total: 500_000_000_000, // 500 GB host
936                available: 100_000_000, // 100 MB left in the cgroup
937            }),
938            4,
939        );
940        assert!(gate.memory_blocked(), "must block despite a 500 GB 'total'");
941    }
942
943    #[test]
944    fn test_gate_memory_guard_blocks_when_no_room_for_one_footprint() {
945        let mut gate = make_gate(
946            &ParallelMode::Auto { min: 1, max: 4 },
947            Some(crate::parallel::MemoryInfo {
948                total: 8_000_000_000,
949                available: 300_000_000, // above the absolute floor...
950            }),
951            4,
952        );
953        // ...but one learned browser needs ~500 MB, so it must still block.
954        gate.footprint = 500.0 * 1024.0 * 1024.0;
955        gate.footprint_count = 3;
956        assert!(gate.memory_blocked());
957    }
958
959    #[test]
960    fn test_oom_failure_halves_limit_and_sets_backoff() {
961        let mut gate = make_gate(&ParallelMode::Auto { min: 1, max: 16 }, None, 100);
962        gate.limit = 16;
963        gate.on_launch_failure();
964        assert_eq!(gate.limit, 8);
965        assert!(gate.retry_backoff > 0);
966    }
967
968    #[test]
969    fn test_is_retryable_oom_matches_memory_errors() {
970        assert!(is_retryable_oom("failed to launch browser: out of memory"));
971        assert!(is_retryable_oom("cannot allocate memory for page"));
972        assert!(!is_retryable_oom("Chrome binary not found"));
973        assert!(!is_retryable_oom("invalid URL"));
974    }
975
976    #[test]
977    fn test_run_options_cloneable() {
978        let _ = RunOptions {
979            mode: ParallelMode::Manual(2),
980            reporter: Arc::new(crate::reporting::Reporter::default()),
981            memory: None,
982        };
983    }
984
985    #[test]
986    fn test_launch_finished_learns_footprint_from_delta() {
987        let mut gate = make_gate(
988            &ParallelMode::Auto { min: 1, max: 4 },
989            Some(crate::parallel::MemoryInfo {
990                total: 8_000_000_000,
991                available: 8_000_000_000,
992            }),
993            4,
994        );
995        // Before = 1 GiB free, after = 600 MiB free → one browser ≈ 400 MiB.
996        gate.launch_finished(Some(1_000_000_000), Some(600_000_000));
997        assert!((gate.footprint - 400_000_000.0).abs() < 1.0);
998        assert_eq!(gate.footprint_count, 1);
999        assert_eq!(gate.available, Some(600_000_000));
1000    }
1001
1002    #[test]
1003    fn test_on_success_raises_to_memory_capacity() {
1004        let mut gate = make_gate(
1005            &ParallelMode::Auto { min: 1, max: 0 },
1006            Some(crate::parallel::MemoryInfo {
1007                total: 8_000_000_000,
1008                available: 4_000_000_000,
1009            }),
1010            100,
1011        );
1012        assert_eq!(gate.limit, 1);
1013        // Learn a 1 GiB footprint, then capacity = 4 GiB / 1 GiB = 4.
1014        gate.footprint = 1_000_000_000.0;
1015        gate.footprint_count = 3;
1016        gate.on_success();
1017        assert_eq!(gate.limit, 4);
1018    }
1019
1020    #[test]
1021    fn test_on_success_ramps_when_memory_unknown() {
1022        let mut gate = make_gate(&ParallelMode::Auto { min: 1, max: 0 }, None, 100);
1023        assert_eq!(gate.limit, 1);
1024        gate.on_success();
1025        assert_eq!(gate.limit, 2, "ramps up one at a time");
1026    }
1027
1028    #[test]
1029    fn test_can_launch_respects_active_limit() {
1030        let mut gate = make_gate(&ParallelMode::Manual(2), None, 10);
1031        assert!(gate.can_launch());
1032        gate.launch_started(Some(1_000_000_000));
1033        assert!(gate.can_launch(), "one of two slots free");
1034        gate.launch_started(Some(1_000_000_000));
1035        assert!(!gate.can_launch(), "both manual slots in use");
1036        gate.release(Some(1_000_000_000));
1037        assert!(gate.can_launch(), "slot freed after release");
1038    }
1039
1040    #[test]
1041    fn test_merge_globals_sums_endpoint_counters() {
1042        let mut snap1 = UsageSnapshot::default();
1043        let mut snap2 = UsageSnapshot::default();
1044        snap1.endpoints.insert(
1045            "a".to_owned(),
1046            EndpointUsage {
1047                calls: 1,
1048                input_tokens: 100,
1049                output_tokens: 50,
1050                cached_input_tokens: 20,
1051                cost: 0.01,
1052                models: std::iter::once("m1".to_owned()).collect(),
1053            },
1054        );
1055        snap2.endpoints.insert(
1056            "a".to_owned(),
1057            EndpointUsage {
1058                calls: 2,
1059                input_tokens: 200,
1060                output_tokens: 100,
1061                cached_input_tokens: 0,
1062                cost: 0.02,
1063                models: std::iter::once("m1".to_owned()).collect(),
1064            },
1065        );
1066        snap2.endpoints.insert(
1067            "b".to_owned(),
1068            EndpointUsage {
1069                calls: 1,
1070                input_tokens: 10,
1071                output_tokens: 5,
1072                cached_input_tokens: 0,
1073                cost: 0.001,
1074                models: std::iter::once("m2".to_owned()).collect(),
1075            },
1076        );
1077        let merged = merge_globals(&[snap1, snap2]);
1078        assert_eq!(merged.endpoints.len(), 2);
1079        let a = merged.endpoints.get("a").unwrap();
1080        assert_eq!(a.calls, 3);
1081        assert_eq!(a.input_tokens, 300);
1082        assert_eq!(a.cached_input_tokens, 20);
1083        assert!((a.cost - 0.03).abs() < 0.0001);
1084        let b = merged.endpoints.get("b").unwrap();
1085        assert_eq!(b.calls, 1);
1086        assert_eq!(merged.models, vec!["m1".to_owned(), "m2".to_owned()]);
1087    }
1088
1089    #[test]
1090    fn test_run_sh_sync_captures_output() {
1091        let Some(out) = run_sh_sync("echo hello") else {
1092            return; // sh unavailable on this platform — nothing to assert
1093        };
1094        assert_eq!(out, "hello");
1095    }
1096
1097    /// Runs the real worker pool against two failing scenario files and
1098    /// asserts the batch machinery: both files' tests are accounted for, and
1099    /// exactly one `RunStarted` + one `RunFinished` are emitted (each file
1100    /// suppresses its own, so the batch owns the pair). Uses an unreachable
1101    /// URL so files fail fast with or without Chrome installed.
1102    #[test]
1103    fn test_run_scenarios_batch_emits_one_run_event_pair() {
1104        use crate::reporting::{ColorMode, Level, Reporter};
1105        use crate::scenario::{TestGroup, TestStep};
1106
1107        let id = std::process::id();
1108        let log_path = std::env::temp_dir().join(format!("lbt-parallel-{id}.ndjson"));
1109        let reporter = Arc::new(
1110            Reporter::new(
1111                Level::Error,
1112                ColorMode::Never,
1113                Some(&log_path),
1114                None,
1115                None,
1116                false,
1117            )
1118            .ok()
1119            .unwrap(),
1120        );
1121
1122        let mut files: Vec<ScenarioFile> = Vec::new();
1123        for (label, url) in [("a", "http://127.0.0.1:9/"), ("b", "http://127.0.0.1:9/")] {
1124            files.push(ScenarioFile {
1125                label: label.to_owned(),
1126                config: ScenarioConfig::default(),
1127                definitions: Vec::new(),
1128                tests: vec![TestGroup {
1129                    name: label.to_owned(),
1130                    start_url: None,
1131                    auto_navigate: None,
1132                    base_url: None,
1133                    timeout_secs: Some(5),
1134                    browser_headless: Some(true),
1135                    viewport_width: None,
1136                    viewport_height: None,
1137                    budget: None,
1138                    endpoint: None,
1139                    steps: vec![TestStep::Navigate {
1140                        url: url.to_owned(),
1141                        wait_after_ms: None,
1142                    }],
1143                }],
1144            });
1145        }
1146
1147        let run = crate::parallel::run_scenarios(
1148            files,
1149            RunOptions {
1150                mode: ParallelMode::Manual(2),
1151                reporter: Arc::clone(&reporter),
1152                memory: None,
1153            },
1154        )
1155        .ok()
1156        .unwrap();
1157
1158        // Both files' single test must be accounted for in the merged report.
1159        assert_eq!(
1160            run.report.tests_passed + run.report.tests_failed,
1161            2,
1162            "both files contributed exactly one test"
1163        );
1164
1165        // Exactly one batch-level RunStarted and one RunFinished.
1166        reporter.finish().ok().unwrap();
1167        let text = std::fs::read_to_string(&log_path).ok().unwrap();
1168        let started = text
1169            .lines()
1170            .filter(|l| l.contains("\"type\":\"run_started\""))
1171            .count();
1172        let finished = text
1173            .lines()
1174            .filter(|l| l.contains("\"type\":\"run_finished\""))
1175            .count();
1176        assert_eq!(started, 1, "exactly one RunStarted for the batch");
1177        assert_eq!(finished, 1, "exactly one RunFinished for the batch");
1178    }
1179}