Skip to main content

nmbrs_workload/
verify.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Workload verification — run a workload and check its output against rules
5//! embedded in the workload file. Shared by the `nmbrs check` subcommand and the
6//! example-walker test, so "how CI checks the examples" and "how a user checks
7//! their own workload" are the same code.
8//!
9//! A verification target is resolved the same way `nmbrs run` resolves
10//! `workload=…`: a directory (walk every workload under it), an existing
11//! `.yaml`/`.yml` file (run it by path), or a **bundled catalog name** such as
12//! `examples/cursors/all_cursor/enumerate` (run it by name, read its rules from the embedded
13//! source). So whatever tab-completion offers for `nmbrs check <TAB>` — local
14//! files *and* catalog names — checks the same way it runs.
15//!
16//! Two **equivalent** rule surfaces (a file may use either, or both — their
17//! cases combine):
18//!
19//! 1. **`#@` comment directives** — trailing YAML comments, inert to the
20//!    runtime:
21//!    ```text
22//!    #@ run scenario=enumerate
23//!    #@ expect 50 completed, 0 failed
24//!    #@ case overload
25//!    #@   run concurrency=32 rate=100000
26//!    #@   expect-fail error_rate_exceeded
27//!    #@ requires backend (needs a live service)
28//!    #@ session cwd            (sessions under the sandbox cwd — stick_session)
29//!    #@ again phases=probe     (a second invocation, session state preserved)
30//!    ```
31//! 2. **A `verify:` YAML block** — also inert (the runtime ignores unknown
32//!    top-level keys). Three equivalent shapes:
33//!    ```yaml
34//!    # single case (a directive map)
35//!    verify: { run: scenario=enumerate, expect: "50 completed, 0 failed" }
36//!    # a list of cases
37//!    verify:
38//!      - { case: a, run: scenario=a, expect: "..." }
39//!      - { case: b, expect-fail: "..." }
40//!    # a name-keyed map (key = case name)
41//!    verify:
42//!      a: { run: scenario=a, expect: "..." }
43//!      b: { expect: ["x", "y"] }
44//!    ```
45//!
46//! `expect` / `expect-fail` accept a single regex or a list of regexes. Each
47//! must match the run's combined stdout+stderr; `expect-fail` additionally
48//! requires a non-zero exit.
49
50use regex::Regex;
51use std::path::{Path, PathBuf};
52use std::process::Command;
53use std::time::{Duration, Instant};
54
55/// Default per-case run timeout, in seconds.
56pub const DEFAULT_TIMEOUT_SECS: u64 = 90;
57
58/// Keywords recognized inside a `verify:` directive map (vs. case names).
59const DIRECTIVE_KEYS: &[&str] = &[
60    "run",
61    "expect",
62    "expect-fail",
63    "expect_fail",
64    "requires",
65    "timeout",
66    "case",
67];
68
69/// One verification case: an invocation plus the regexes its output must match.
70pub struct VerifyCase {
71    pub name: String,
72    pub run_args: Vec<String>,
73    /// SRD-108 examples round-trip — extra invocations of the SAME
74    /// workload in the same sandbox after the first run, one arg
75    /// list per `#@ again` line, session state PRESERVED between
76    /// invocations. Every non-final invocation must succeed; the
77    /// final one feeds the expect / expect-fail rules, and the
78    /// `expect` regexes match the CONCATENATED output of all
79    /// invocations.
80    pub again: Vec<Vec<String>>,
81    /// `#@ session cwd` — omit the harness's `--session-path`
82    /// injection so sessions land under the (throwaway) sandbox
83    /// cwd's `sessions/` root. Required for behaviors keyed to
84    /// `sessions/latest` (SRD-106 `stick_session`); the sandbox
85    /// `sessions/` dir is wiped at case start so cases stay
86    /// independent.
87    pub session_cwd: bool,
88    pub expects: Vec<Regex>,
89    pub expect_fails: Vec<Regex>,
90    pub timeout: u64,
91}
92
93impl VerifyCase {
94    fn new(name: impl Into<String>) -> Self {
95        VerifyCase {
96            name: name.into(),
97            run_args: Vec::new(),
98            again: Vec::new(),
99            session_cwd: false,
100            expects: Vec::new(),
101            expect_fails: Vec::new(),
102            timeout: DEFAULT_TIMEOUT_SECS,
103        }
104    }
105}
106
107/// The parsed verification plan for one workload file.
108pub struct VerifyPlan {
109    /// `Some(reason)` ⇒ skip this file (e.g. needs external infra).
110    pub requires: Option<String>,
111    pub cases: Vec<VerifyCase>,
112}
113
114impl VerifyPlan {
115    /// Parse both rule surfaces from a workload file's text and combine them.
116    pub fn parse(src: &str) -> Result<VerifyPlan, String> {
117        let mut plan = parse_directives(src)?;
118        merge_verify_block(src, &mut plan)?;
119        Ok(plan)
120    }
121
122    /// True when the file declares no verification rules at all.
123    pub fn is_empty(&self) -> bool {
124        self.requires.is_none() && self.cases.is_empty()
125    }
126}
127
128fn compile(v: &str) -> Result<Regex, String> {
129    Regex::new(v).map_err(|e| format!("bad regex {v:?}: {e}"))
130}
131
132// ── `#@` comment directives ───────────────────────────────────────────────
133
134fn parse_directives(src: &str) -> Result<VerifyPlan, String> {
135    let mut requires: Option<String> = None;
136    let mut cases: Vec<VerifyCase> = Vec::new();
137    let mut cur: Option<VerifyCase> = None;
138
139    for raw in src.lines() {
140        let line = raw.trim_start();
141        let Some(rest) = line.strip_prefix("#@") else {
142            continue;
143        };
144        let rest = rest.trim();
145        let (kw, val) = match rest.split_once(char::is_whitespace) {
146            Some((k, v)) => (k.trim_end_matches(':'), v.trim()),
147            None => (rest.trim_end_matches(':'), ""),
148        };
149        macro_rules! case {
150            () => {
151                cur.get_or_insert_with(|| VerifyCase::new("default"))
152            };
153        }
154        match kw {
155            "requires" => requires = Some(val.to_string()),
156            "case" => {
157                if let Some(c) = cur.take() {
158                    cases.push(c);
159                }
160                cur = Some(VerifyCase::new(val));
161            }
162            "run" => case!().run_args = val.split_whitespace().map(String::from).collect(),
163            "again" => case!()
164                .again
165                .push(val.split_whitespace().map(String::from).collect()),
166            "session" => match val {
167                "cwd" => case!().session_cwd = true,
168                other => return Err(format!("unknown `#@ session {other}` (only `cwd`)")),
169            },
170            "expect" => case!().expects.push(compile(val)?),
171            "expect-fail" | "expect_fail" => case!().expect_fails.push(compile(val)?),
172            "timeout" => {
173                case!().timeout = val.parse().map_err(|_| format!("bad timeout {val:?}"))?
174            }
175            other => return Err(format!("unknown `#@ {other}` directive")),
176        }
177    }
178    if let Some(c) = cur.take() {
179        cases.push(c);
180    }
181    Ok(VerifyPlan { requires, cases })
182}
183
184// ── `verify:` YAML block ──────────────────────────────────────────────────
185
186fn merge_verify_block(src: &str, plan: &mut VerifyPlan) -> Result<(), String> {
187    // The file may not be a YAML mapping (or may be a `#@`-only file); either
188    // way, no `verify:` block to merge.
189    let Ok(doc) = serde_yaml::from_str::<serde_yaml::Value>(src) else {
190        return Ok(());
191    };
192    let Some(block) = doc.get("verify") else {
193        return Ok(());
194    };
195    match block {
196        serde_yaml::Value::Sequence(items) => {
197            for (i, item) in items.iter().enumerate() {
198                plan.cases.push(case_from_map(item, None, i)?);
199            }
200        }
201        serde_yaml::Value::Mapping(m) => {
202            let all_directives = m
203                .keys()
204                .all(|k| k.as_str().is_some_and(|s| DIRECTIVE_KEYS.contains(&s)));
205            if all_directives {
206                // A single directive map — either a file-level `requires` skip
207                // or one unnamed case.
208                if let Some(req) = m.get("requires").and_then(|v| v.as_str()) {
209                    plan.requires = Some(req.to_string());
210                } else {
211                    plan.cases.push(case_from_map(block, None, 0)?);
212                }
213            } else {
214                // A name-keyed map: each key is a case name.
215                for (k, v) in m {
216                    let name = k.as_str().ok_or("verify: case name must be a string")?;
217                    plan.cases.push(case_from_map(v, Some(name), 0)?);
218                }
219            }
220        }
221        _ => return Err("verify: must be a map or a list of cases".to_string()),
222    }
223    Ok(())
224}
225
226/// Build a case from a `verify:` entry map. `name_override` (the key in a
227/// name-keyed map) wins; else the entry's `case:` field; else `case-<idx>`.
228fn case_from_map(
229    v: &serde_yaml::Value,
230    name_override: Option<&str>,
231    idx: usize,
232) -> Result<VerifyCase, String> {
233    let m = v.as_mapping().ok_or("verify: each case must be a map")?;
234    let get = |k: &str| m.get(serde_yaml::Value::from(k));
235    let name = name_override
236        .map(String::from)
237        .or_else(|| get("case").and_then(|x| x.as_str()).map(String::from))
238        .unwrap_or_else(|| format!("case-{idx}"));
239    let mut case = VerifyCase::new(name);
240    if let Some(run) = get("run").and_then(|x| x.as_str()) {
241        case.run_args = run.split_whitespace().map(String::from).collect();
242    }
243    for s in strings_of(get("expect")) {
244        case.expects.push(compile(&s)?);
245    }
246    for s in strings_of(get("expect-fail").or_else(|| get("expect_fail"))) {
247        case.expect_fails.push(compile(&s)?);
248    }
249    if let Some(t) = get("timeout").and_then(|x| x.as_u64()) {
250        case.timeout = t;
251    }
252    Ok(case)
253}
254
255/// A scalar string or a sequence of strings → `Vec<String>`.
256fn strings_of(v: Option<&serde_yaml::Value>) -> Vec<String> {
257    match v {
258        Some(serde_yaml::Value::String(s)) => vec![s.clone()],
259        Some(serde_yaml::Value::Sequence(items)) => items
260            .iter()
261            .filter_map(|x| x.as_str().map(String::from))
262            .collect(),
263        _ => Vec::new(),
264    }
265}
266
267// ── Run + check ───────────────────────────────────────────────────────────
268
269/// What one verification case produced.
270pub enum Outcome {
271    Pass,
272    Skip(String),
273    Fail(String),
274}
275
276/// Run one case via `<binary> run workload=… <args> --session-path …` (under
277/// an in-process deadline) from `sandbox`, capture combined output, and check
278/// the rules.
279/// `workload_ref` is whatever goes after `workload=` — an absolute file path
280/// or a bundled catalog name; the subprocess resolves it exactly as a normal
281/// `nmbrs run` would.
282pub fn run_case(
283    binary: &Path,
284    workload_ref: &str,
285    sandbox: &Path,
286    label: &str,
287    case: &VerifyCase,
288) -> Result<(), String> {
289    let safe_label = label.replace(['/', ' ', ':'], "_");
290    let session = sandbox.join(format!("session-{safe_label}"));
291    // Case independence for cwd-session cases: each gets a
292    // PRIVATE working directory, so its `sessions/latest` can
293    // neither re-attach to a prior case's session nor race a
294    // concurrent case's — the check walker drives cases on
295    // worker threads, and a shared cwd made multi-invocation
296    // re-attach cases fail only in large sweeps.
297    let case_cwd;
298    let workdir: &Path = if case.session_cwd {
299        case_cwd = sandbox.join(format!("cwd-{safe_label}"));
300        let _ = std::fs::remove_dir_all(&case_cwd);
301        std::fs::create_dir_all(&case_cwd).map_err(|e| format!("create case cwd: {e}"))?;
302        &case_cwd
303    } else {
304        let _ = std::fs::remove_dir_all(&session);
305        sandbox
306    };
307
308    // The deadline is enforced here rather than by wrapping in the
309    // coreutils `timeout` program: that binary isn't a given on
310    // macOS, and on Windows `timeout.exe` is the cmd.exe delay
311    // command — it can't run a child at all.
312    let invoke = |args: &[String]| -> Result<(String, bool, bool), String> {
313        use std::io::Read as _;
314        use std::process::Stdio;
315        let mut cmd = Command::new(binary);
316        cmd.arg("run")
317            .arg(format!("workload={workload_ref}"))
318            .args(args);
319        if !case.session_cwd {
320            cmd.arg("--session-path").arg(&session);
321        }
322        let mut child = cmd
323            .current_dir(workdir)
324            .stdin(Stdio::null())
325            .stdout(Stdio::piped())
326            .stderr(Stdio::piped())
327            .spawn()
328            .map_err(|e| format!("spawn failed: {e}"))?;
329        // Drain both pipes on threads so a chatty child can't fill
330        // one while we watch the deadline, deadlocking both sides.
331        let mut out_pipe = child.stdout.take().expect("stdout piped");
332        let mut err_pipe = child.stderr.take().expect("stderr piped");
333        let out_h = std::thread::spawn(move || {
334            let mut buf = Vec::new();
335            let _ = out_pipe.read_to_end(&mut buf);
336            buf
337        });
338        let err_h = std::thread::spawn(move || {
339            let mut buf = Vec::new();
340            let _ = err_pipe.read_to_end(&mut buf);
341            buf
342        });
343        let deadline = std::time::Instant::now() + std::time::Duration::from_secs(case.timeout);
344        let (status, timed_out) = loop {
345            if let Some(st) = child.try_wait().map_err(|e| format!("wait failed: {e}"))? {
346                break (st, false);
347            }
348            if std::time::Instant::now() >= deadline {
349                let _ = child.kill();
350                let st = child.wait().map_err(|e| format!("wait failed: {e}"))?;
351                break (st, true);
352            }
353            std::thread::sleep(std::time::Duration::from_millis(25));
354        };
355        let stdout = out_h.join().unwrap_or_default();
356        let stderr = err_h.join().unwrap_or_default();
357        let combined = format!(
358            "{}{}",
359            String::from_utf8_lossy(&stdout),
360            String::from_utf8_lossy(&stderr)
361        );
362        Ok((combined, status.success(), timed_out))
363    };
364
365    let (mut combined, mut succeeded, mut timed_out) = invoke(&case.run_args)?;
366    for (i, extra) in case.again.iter().enumerate() {
367        if !succeeded || timed_out {
368            return Err(format!(
369                "invocation {} of {} failed before the `again` steps                  completed:
370{combined}",
371                i + 1,
372                case.again.len() + 1
373            ));
374        }
375        let (c, s, to) = invoke(extra)?;
376        combined.push_str(&c);
377        succeeded = s;
378        timed_out = to;
379    }
380    check_case_output(case, &combined, succeeded, timed_out)
381}
382
383/// Check one case's RUN RESULT against its expectations — the run-mechanism-
384/// agnostic half of [`run_case`]. `combined` is the merged stdout+stderr the
385/// run produced; `succeeded` is its exit success; `timed_out` is the timeout
386/// signal. Used both by the subprocess [`run_case`] and by `nmbrs`'s in-process
387/// `run_executions`-backed verification (so the SAME `expect` / `expect-fail`
388/// rules apply whether examples run as subprocesses or as concurrent in-process
389/// executions sharing one session).
390pub fn check_case_output(
391    case: &VerifyCase,
392    combined: &str,
393    succeeded: bool,
394    timed_out: bool,
395) -> Result<(), String> {
396    if timed_out {
397        return Err(format!("timed out after {}s", case.timeout));
398    }
399    if case.expect_fails.is_empty() {
400        if !succeeded {
401            let err = combined
402                .lines()
403                .find(|l| l.contains("error:") || l.contains("panic"))
404                .unwrap_or("(no error line)");
405            return Err(format!("run failed (expected success): {err}"));
406        }
407    } else {
408        if succeeded {
409            return Err("expected a failure (`expect-fail`) but the run succeeded".to_string());
410        }
411        for re in &case.expect_fails {
412            if !re.is_match(combined) {
413                return Err(format!(
414                    "expect-fail /{re}/ did not match the failure output"
415                ));
416            }
417        }
418    }
419    for re in &case.expects {
420        if !re.is_match(combined) {
421            return Err(format!("expect /{re}/ did not match the output"));
422        }
423    }
424    Ok(())
425}
426
427/// Whether a workload is REQUIRED to declare verification rules.
428///
429/// True for anything under an `examples/` directory: those files are the
430/// documented, CI-gated surface, and one arriving without rules is a
431/// regression in the documentation itself. Everywhere else rules are
432/// optional — an unruled workload is skipped, not failed.
433///
434/// Path-component matching, not substring: a workload at
435/// `/home/me/examples-scratch/w.yaml` is not under `examples/`, while
436/// `/repo/examples/optimizer/w.yaml` is.
437pub fn requires_verification_rules(run_ref: &str) -> bool {
438    std::path::Path::new(run_ref)
439        .components()
440        .any(|c| c.as_os_str() == "examples")
441}
442
443/// Verify a workload from its rule text + run reference: parse the rules, run
444/// every case, return one `(label, Outcome)` per case (or a single Skip / Fail
445/// for the whole workload). `run_ref` is what to pass as `workload=` — an
446/// absolute file path or a catalog name.
447pub fn verify_source(
448    binary: &Path,
449    label_root: &str,
450    run_ref: &str,
451    rule_text: &str,
452    sandbox: &Path,
453) -> Vec<(String, Outcome)> {
454    let plan = match VerifyPlan::parse(rule_text) {
455        Ok(p) => p,
456        Err(e) => return vec![(label_root.to_string(), Outcome::Fail(e))],
457    };
458    if let Some(reason) = plan.requires {
459        return vec![(label_root.to_string(), Outcome::Skip(reason))];
460    }
461    if plan.cases.is_empty() {
462        // A workload with no rules is not a failure in general — it is a
463        // workload nobody asked to be checked. `nmbrs check <dir>` walks every
464        // YAML it finds, and most of those (adapter workloads, operational
465        // ones like an incremental compaction sweep) exist to be RUN against real
466        // infrastructure, not to self-verify; failing them made the walk's
467        // result meaningless and trained people to ignore it.
468        //
469        // Under `examples/` the opposite holds: those files are the
470        // documentation CI gates, and one landing without rules is exactly
471        // the regression this check exists to catch. So there, silence is a
472        // failure.
473        return vec![(
474            label_root.to_string(),
475            if requires_verification_rules(run_ref) {
476                Outcome::Fail(
477                    "no verification rules — every workload under `examples/` must \
478                     declare them (add `#@ expect …` comments or a `verify:` block)"
479                        .into(),
480                )
481            } else {
482                Outcome::Skip("no verification rules — nothing to check".into())
483            },
484        )];
485    }
486    plan.cases
487        .iter()
488        .map(|c| {
489            let label = format!("{label_root}::{}", c.name);
490            let outcome = match run_case(binary, run_ref, sandbox, &label, c) {
491                Ok(()) => Outcome::Pass,
492                Err(e) => Outcome::Fail(format!("{label}: {e}")),
493            };
494            (label, outcome)
495        })
496        .collect()
497}
498
499/// Verify one workload file: read it, then run it by its absolute path. The
500/// file is read relative to *this* process's cwd, but the run reference is made
501/// absolute because each case is launched from a sandbox cwd.
502pub fn verify_file(
503    binary: &Path,
504    label_root: &str,
505    workload: &Path,
506    sandbox: &Path,
507) -> Vec<(String, Outcome)> {
508    let src = match std::fs::read_to_string(workload) {
509        Ok(s) => s,
510        Err(e) => {
511            return vec![(
512                label_root.to_string(),
513                Outcome::Fail(format!("read error: {e}")),
514            )];
515        }
516    };
517    let abs = workload
518        .canonicalize()
519        .unwrap_or_else(|_| workload.to_path_buf());
520    verify_source(binary, label_root, &abs.to_string_lossy(), &src, sandbox)
521}
522
523/// Where a verification target's rule text and run reference come from. Mirrors
524/// `nmbrs run`'s `workload=…` resolution.
525pub enum WorkloadSource {
526    /// A workload file on disk: run by (absolute) path, rules from the file.
527    File(PathBuf),
528    /// A bundled catalog workload: run by name, rules from the embedded source.
529    Catalog { name: String, source: String },
530}
531
532/// Resolve a single workload reference the way `nmbrs run` does: an existing
533/// `.yaml`/`.yml` file path first, then a bundled catalog name. `None` if it is
534/// neither. Directories are not a single workload — callers handle those
535/// separately (via [`verify_path`]).
536pub fn resolve_ref(reference: &str) -> Option<WorkloadSource> {
537    let p = Path::new(reference);
538    if p.is_file() {
539        // Absolute so the sandbox-cwd subprocess can still find it.
540        return Some(WorkloadSource::File(
541            p.canonicalize().unwrap_or_else(|_| p.to_path_buf()),
542        ));
543    }
544    crate::catalog::lookup(reference).map(|w| WorkloadSource::Catalog {
545        name: w.name.to_string(),
546        source: w.source.to_string(),
547    })
548}
549
550/// A workload reference's **declared top-level `params:`** (string scalars),
551/// following its `extends:` chain — resolved the way `nmbrs run` resolves
552/// `workload=…`. Returns `None` if the reference resolves to nothing or has
553/// no `params:` block.
554///
555/// The canonical "what params did the workload declare" accessor. The
556/// runner folds these under the CLI params (CLI wins) to form the run's
557/// effective params, so a setting works identically whether declared in the
558/// workload or passed on the command line. Also used to recognize a
559/// **console-owning adapter declared in the workload** (e.g.
560/// `params: { adapter: plotter }`) so the dashboard yields to the adapter on
561/// a TTY (SRD-41/87); the returned params carry the adapter's display-shaping
562/// keys (e.g. stdout's `filename`) so the preference is decided correctly.
563pub fn declared_params(reference: &str) -> Option<std::collections::HashMap<String, String>> {
564    let merged = match resolve_ref(reference)? {
565        // Display-shaping probe only: resolution warnings are
566        // dropped HERE because the authoritative load that follows
567        // this probe surfaces the identical warnings itself.
568        WorkloadSource::File(path) => crate::extends::load_and_merge(&path).ok()?.0,
569        WorkloadSource::Catalog { name, .. } => {
570            crate::extends::load_and_merge_bundled(crate::catalog::lookup(&name)?)
571                .ok()?
572                .0
573        }
574    };
575    let doc: serde_yaml::Value = serde_yaml::from_str(&merged).ok()?;
576    let params = doc.get("params")?.as_mapping()?;
577    let mut out = std::collections::HashMap::new();
578    for (k, v) in params {
579        let Some(key) = k.as_str() else { continue };
580        // Only scalar params shape the display decision; skip nested
581        // structures. Numbers / bools render to their lexical form.
582        let val = match v {
583            serde_yaml::Value::String(s) => s.clone(),
584            serde_yaml::Value::Number(n) => n.to_string(),
585            serde_yaml::Value::Bool(b) => b.to_string(),
586            _ => continue,
587        };
588        out.insert(key.to_string(), val);
589    }
590    Some(out)
591}
592
593/// Aggregate verification result across one or more files.
594#[derive(Default)]
595pub struct VerifySummary {
596    pub passed: usize,
597    pub skipped: Vec<String>,
598    pub failures: Vec<String>,
599    /// One entry per workload checked, in completion order — the
600    /// raw material for an end-of-run "slowest workloads" report.
601    pub timings: Vec<WorkloadTiming>,
602}
603
604/// A single workload's aggregate outcome, for live progress and the
605/// timing report. Coarser than [`Outcome`] (which is per *case*): a
606/// workload with any failing case is [`CheckStatus::Fail`]; one with no
607/// failures and at least one skip (and no pass) is [`CheckStatus::Skip`];
608/// otherwise [`CheckStatus::Pass`].
609#[derive(Debug, Clone, Copy, PartialEq, Eq)]
610pub enum CheckStatus {
611    Pass,
612    Skip,
613    Fail,
614}
615
616/// Wall-clock spent checking one workload, with its aggregate status.
617#[derive(Debug, Clone)]
618pub struct WorkloadTiming {
619    pub label: String,
620    pub elapsed: Duration,
621    pub status: CheckStatus,
622}
623
624/// Live verification progress, emitted as workloads start and finish so a
625/// caller (e.g. `nmbrs check`) can render an active/pending/done/errors
626/// status line. Invoked from worker threads — handlers must be `Sync`.
627#[derive(Debug, Clone)]
628pub enum CheckProgress {
629    /// Discovery finished — `total` workloads will be checked.
630    Begin { total: usize },
631    /// A workload began running.
632    Started { label: String },
633    /// A workload finished, with its wall-clock and aggregate status.
634    Finished {
635        label: String,
636        elapsed: Duration,
637        status: CheckStatus,
638    },
639}
640
641/// A progress handler. `&`-shared across verification worker threads.
642pub type ProgressFn<'a> = dyn Fn(CheckProgress) + Sync + 'a;
643
644/// No-op progress handler, for callers that only want the summary.
645pub fn no_progress(_: CheckProgress) {}
646
647/// Fold a file's per-case outcomes into one [`CheckStatus`].
648fn aggregate_status(outcomes: &[(String, Outcome)]) -> CheckStatus {
649    let mut saw_pass = false;
650    let mut saw_skip = false;
651    for (_, o) in outcomes {
652        match o {
653            Outcome::Fail(_) => return CheckStatus::Fail,
654            Outcome::Pass => saw_pass = true,
655            Outcome::Skip(_) => saw_skip = true,
656        }
657    }
658    if saw_skip && !saw_pass {
659        CheckStatus::Skip
660    } else {
661        CheckStatus::Pass
662    }
663}
664
665/// Verify a workload file or every `*.yaml` under a directory (recursively).
666/// Files run concurrently (cases within a file are sequential). `progress`
667/// is invoked from worker threads as each workload starts and finishes —
668/// pass [`no_progress`] for a quiet run.
669pub fn verify_path(
670    binary: &Path,
671    path: &Path,
672    sandbox: &Path,
673    progress: &ProgressFn,
674) -> VerifySummary {
675    let _ = std::fs::create_dir_all(sandbox);
676    let mut files: Vec<PathBuf> = Vec::new();
677    if path.is_dir() {
678        collect_yaml(path, &mut files);
679        files.sort();
680    } else {
681        files.push(path.to_path_buf());
682    }
683    progress(CheckProgress::Begin { total: files.len() });
684
685    let acc: std::sync::Mutex<VerifySummary> = std::sync::Mutex::new(VerifySummary::default());
686    // Work-stealing over `files`: a shared atomic cursor each worker pulls from
687    // when it frees up. Static chunking serialized the slow demos (settle /
688    // servo workloads that dwell on real-time metric windows) behind one chunk
689    // while other chunks idled; a shared queue runs them concurrently, so the
690    // wall-clock floor is the single slowest file, not a chunk-sum. Workers
691    // spend almost all their time blocked on the child `nmbrs` process, so we
692    // oversubscribe past core count (2× cores, capped) to overlap the waits.
693    let next = std::sync::atomic::AtomicUsize::new(0);
694    let workers = files.len().min(
695        std::thread::available_parallelism()
696            .map(|n| n.get())
697            .unwrap_or(4)
698            .saturating_mul(2)
699            .clamp(1, 16),
700    );
701    std::thread::scope(|s| {
702        for _ in 0..workers {
703            s.spawn(|| {
704                loop {
705                    let i = next.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
706                    let Some(f) = files.get(i) else { break };
707                    let label = f
708                        .file_name()
709                        .and_then(|n| n.to_str())
710                        .unwrap_or("?")
711                        .to_string();
712                    // Run the file (a sequence of cases) outside the lock, then
713                    // fold its outcomes under a single lock acquisition.
714                    progress(CheckProgress::Started {
715                        label: label.clone(),
716                    });
717                    let start = Instant::now();
718                    let outcomes = verify_file(binary, &label, f, sandbox);
719                    let elapsed = start.elapsed();
720                    let status = aggregate_status(&outcomes);
721                    let mut g = acc.lock().unwrap();
722                    for (lbl, outcome) in outcomes {
723                        match outcome {
724                            Outcome::Pass => g.passed += 1,
725                            Outcome::Skip(r) => g.skipped.push(format!("{lbl}: {r}")),
726                            Outcome::Fail(m) => g.failures.push(m),
727                        }
728                    }
729                    g.timings.push(WorkloadTiming {
730                        label: label.clone(),
731                        elapsed,
732                        status,
733                    });
734                    drop(g);
735                    progress(CheckProgress::Finished {
736                        label,
737                        elapsed,
738                        status,
739                    });
740                }
741            });
742        }
743    });
744    let mut sum = acc.into_inner().unwrap();
745    sum.skipped.sort();
746    sum.failures.sort();
747    sum
748}
749
750/// Verify a target named the way `nmbrs run` names workloads: a directory (walk
751/// every workload under it), an existing workload file, or a bundled catalog
752/// name (`examples/cursors/all_cursor/enumerate`, …). This is the `nmbrs check` entry point — so
753/// anything the binary can `run` by name, it can `check` by the same name.
754pub fn verify_target(
755    binary: &Path,
756    target: &str,
757    sandbox: &Path,
758    progress: &ProgressFn,
759) -> VerifySummary {
760    let p = Path::new(target);
761    if p.is_dir() {
762        return verify_path(binary, p, sandbox, progress);
763    }
764    let _ = std::fs::create_dir_all(sandbox);
765    // A single named target is one workload; emit the same Begin/Started/
766    // Finished lifecycle a directory walk does so progress rendering is
767    // uniform, and record its timing for the report.
768    progress(CheckProgress::Begin { total: 1 });
769    let label = match &resolve_ref(target) {
770        Some(WorkloadSource::File(path)) => path
771            .file_name()
772            .and_then(|n| n.to_str())
773            .unwrap_or(target)
774            .to_string(),
775        _ => target.to_string(),
776    };
777    progress(CheckProgress::Started {
778        label: label.clone(),
779    });
780    let start = Instant::now();
781    let cases: Vec<(String, Outcome)> = match resolve_ref(target) {
782        Some(WorkloadSource::File(path)) => verify_file(binary, &label, &path, sandbox),
783        Some(WorkloadSource::Catalog { name, source }) => {
784            verify_source(binary, &name, &name, &source, sandbox)
785        }
786        None => vec![(
787            target.to_string(),
788            Outcome::Fail(format!(
789                "no such workload '{target}': not a local file, not a directory, and \
790                 no bundled workload by that name (try `nmbrs describe workloads --all`).{}",
791                crate::suggest::did_you_mean(&crate::suggest::suggest_workloads(target))
792            )),
793        )],
794    };
795    let elapsed = start.elapsed();
796    let status = aggregate_status(&cases);
797    let mut sum = VerifySummary::default();
798    for (lbl, outcome) in cases {
799        match outcome {
800            Outcome::Pass => sum.passed += 1,
801            Outcome::Skip(r) => sum.skipped.push(format!("{lbl}: {r}")),
802            Outcome::Fail(m) => sum.failures.push(m),
803        }
804    }
805    sum.timings.push(WorkloadTiming {
806        label: label.clone(),
807        elapsed,
808        status,
809    });
810    progress(CheckProgress::Finished {
811        label,
812        elapsed,
813        status,
814    });
815    sum
816}
817
818/// Every `*.yaml` / `*.yml` workload file under `dir` (recursive),
819/// sorted. The discovery half of [`verify_path`], exposed so the
820/// in-process example walker (`nmbrs_runtime::verify_in_process`) finds
821/// the same files the subprocess walker does.
822pub fn collect_workload_files(dir: &Path) -> Vec<PathBuf> {
823    let mut files = Vec::new();
824    collect_yaml(dir, &mut files);
825    files.sort();
826    files
827}
828
829fn collect_yaml(dir: &Path, out: &mut Vec<PathBuf>) {
830    let Ok(entries) = std::fs::read_dir(dir) else {
831        return;
832    };
833    for e in entries.flatten() {
834        let p = e.path();
835        if p.is_dir() {
836            collect_yaml(&p, out);
837        } else if p.extension().is_some_and(|x| x == "yaml" || x == "yml") {
838            out.push(p);
839        }
840    }
841}
842
843#[cfg(test)]
844mod rules_required_tests {
845    use super::requires_verification_rules;
846
847    /// Under `examples/`, rules are the point — a file without them is a
848    /// documentation regression, so the check must fail rather than skip.
849    #[test]
850    fn examples_require_rules() {
851        assert!(requires_verification_rules(
852            "/repo/examples/optimizer/control.yaml"
853        ));
854        assert!(requires_verification_rules("examples/w.yaml"));
855        assert!(requires_verification_rules(
856            "/repo/nmbrs/examples/modules/module_test.yaml"
857        ));
858    }
859
860    /// Everywhere else rules are optional: adapter and operational workloads
861    /// exist to be RUN against real infrastructure, not to self-verify.
862    #[test]
863    fn other_locations_do_not_require_rules() {
864        assert!(!requires_verification_rules(
865            "/repo/nmbrs/workloads/cql/incremental_sweep.yaml"
866        ));
867        assert!(!requires_verification_rules("/tmp/scratch.yaml"));
868        assert!(!requires_verification_rules("some_catalog_name"));
869    }
870
871    /// Component matching, not substring — a sibling directory whose name
872    /// merely starts with "examples" is not the examples tree.
873    #[test]
874    fn matches_path_components_not_substrings() {
875        assert!(!requires_verification_rules(
876            "/home/me/examples-scratch/w.yaml"
877        ));
878        assert!(!requires_verification_rules("/home/me/myexamples/w.yaml"));
879        assert!(requires_verification_rules(
880            "/home/me/examples/scratch/w.yaml"
881        ));
882    }
883}
884
885#[cfg(test)]
886mod tests {
887    use super::*;
888
889    #[test]
890    fn comment_and_yaml_forms_are_equivalent() {
891        let comment = "ops: { a: { raw: x } }\n#@ run cycles=3\n#@ expect 0 failed\n";
892        let single = "ops: { a: { raw: x } }\nverify: { run: cycles=3, expect: \"0 failed\" }\n";
893        let listed =
894            "ops: { a: { raw: x } }\nverify:\n  - { run: cycles=3, expect: \"0 failed\" }\n";
895        let named =
896            "ops: { a: { raw: x } }\nverify:\n  smoke: { run: cycles=3, expect: \"0 failed\" }\n";
897        for src in [comment, single, listed, named] {
898            let p = VerifyPlan::parse(src).expect("parse");
899            assert_eq!(p.cases.len(), 1, "one case for: {src}");
900            assert_eq!(p.cases[0].run_args, vec!["cycles=3"], "run for: {src}");
901            assert_eq!(p.cases[0].expects.len(), 1, "expect for: {src}");
902        }
903    }
904
905    #[test]
906    fn name_keyed_map_yields_named_cases() {
907        let src = "verify:\n  alpha: { run: scenario=a, expect: \"x\" }\n  beta: { expect: [\"y\", \"z\"] }\n";
908        let p = VerifyPlan::parse(src).unwrap();
909        let names: Vec<&str> = p.cases.iter().map(|c| c.name.as_str()).collect();
910        assert!(
911            names.contains(&"alpha") && names.contains(&"beta"),
912            "names: {names:?}"
913        );
914        let beta = p.cases.iter().find(|c| c.name == "beta").unwrap();
915        assert_eq!(beta.expects.len(), 2);
916    }
917
918    #[test]
919    fn requires_block_skips_the_file() {
920        let p = VerifyPlan::parse("verify: { requires: needs a backend }\n").unwrap();
921        assert_eq!(p.requires.as_deref(), Some("needs a backend"));
922        assert!(p.cases.is_empty());
923    }
924
925    #[test]
926    fn comment_and_block_cases_combine() {
927        let src = "#@ case fromcomment\n#@   expect a\nverify:\n  fromblock: { expect: b }\n";
928        let p = VerifyPlan::parse(src).unwrap();
929        let names: Vec<&str> = p.cases.iter().map(|c| c.name.as_str()).collect();
930        assert!(
931            names.contains(&"fromcomment") && names.contains(&"fromblock"),
932            "{names:?}"
933        );
934    }
935
936    #[test]
937    fn resolve_ref_finds_files_and_rejects_unknown_names() {
938        // An on-disk file resolves to an absolute `File` source.
939        let dir = std::env::temp_dir().join(format!("nmbrs-verify-resolve-{}", std::process::id()));
940        let _ = std::fs::create_dir_all(&dir);
941        let file = dir.join("w.yaml");
942        std::fs::write(&file, "ops: { a: { raw: x } }\n").unwrap();
943        match resolve_ref(file.to_str().unwrap()) {
944            Some(WorkloadSource::File(p)) => assert!(p.is_absolute(), "absolute: {p:?}"),
945            other => panic!(
946                "expected File, got {}",
947                matches!(other, Some(WorkloadSource::Catalog { .. })) as i32
948            ),
949        }
950        // A name that is neither a file nor (in this test process) a bundled
951        // workload resolves to nothing — the CLI reports it as not found.
952        assert!(resolve_ref("definitely/not/a/workload").is_none());
953        let _ = std::fs::remove_dir_all(&dir);
954    }
955}