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