Skip to main content

henad_explore/
sweep.rs

1//! Sweeps that plan a spec against a model, run every planned run and write the results to a directory or to memory.
2
3#[cfg(not(target_arch = "wasm32"))]
4use std::borrow::Cow;
5use std::fmt;
6#[cfg(not(target_arch = "wasm32"))]
7use std::io::{self, Write};
8use std::path::{Path, PathBuf};
9use std::sync::Arc;
10use std::time::Duration;
11
12use web_time::Instant;
13
14use henad_compute::entry::ModelEntry;
15use henad_compute::gpu::GpuContext;
16use henad_core::explore::measure::{MeasureError, MeasurePlan};
17use henad_core::explore::outcome::PlannedRun;
18#[cfg(not(target_arch = "wasm32"))]
19use henad_core::explore::outcome::RunOutcome;
20use henad_core::explore::plan::{Plan, PlanError, PlanWarning, PlannedBlock, Shard};
21use henad_core::explore::search::SearchReport;
22use henad_core::explore::spec::SweepSpec;
23use henad_core::metadata::Backend;
24use henad_core::provenance::BuildInfo;
25
26use crate::exec::{
27    ActiveRuns, BatchEnd, Concurrency, ExecutionBudget, ExecutionError, ExecutionLayout, SweepControl, choose_layout,
28    gpu_memory_budget,
29};
30#[cfg(not(target_arch = "wasm32"))]
31use crate::exec::{Executor, RunRequest, RunSink};
32#[cfg(not(target_arch = "wasm32"))]
33use crate::handle::SweepOutput;
34#[cfg(not(target_arch = "wasm32"))]
35use crate::output::OutputWriter;
36use crate::output::manifest::{
37    BuildRole, FORMAT, FORMAT_VERSION, Manifest, ManifestBlock, ManifestColumns, ManifestDesignTable,
38    ManifestExecution, ManifestMode, ManifestModel, ManifestPlan, ManifestRuntime, ManifestSearch, ManifestSeeds,
39    ManifestSession, ManifestSpecSource, ManifestStatus, ManifestTimestamps, RecordedBuild, ResultCounts, now_unix_ms,
40    rfc3339,
41};
42use crate::output::memory::SweepFiles;
43use crate::output::resume::{ResumeError, ResumeScan};
44use crate::output::runs_csv::column_names;
45use crate::output::{OutputDir, OutputError, runs_csv, series_csv};
46use crate::probe::{CapacityError, ProbeError, ProbeReport, TimedProbe, check_capacity};
47#[cfg(not(target_arch = "wasm32"))]
48use crate::progress::ProgressMeter;
49use crate::progress::{Progress, ProgressEvent};
50use crate::schema::{backend_name, schema_json};
51use crate::search_run::{SearchOutline, SearchPlanError};
52use crate::spec_file::{DesignTableFile, ExecutionTable, SpecFile};
53
54/// Spec file a sweep was read from.
55#[derive(Debug, Clone, Default, PartialEq, Eq)]
56pub struct SpecSource {
57    /// Path of the file as given, `None` for a spec built from flags.
58    pub path: Option<PathBuf>,
59    /// Text of the file as written.
60    pub toml: Option<String>,
61    /// Design tables the file reads.
62    pub tables: Vec<DesignTableFile>,
63}
64
65impl SpecSource {
66    /// Returns the source of `file`, read from `path` as the text `toml`.
67    pub(crate) fn loaded(path: &Path, toml: String, file: &SpecFile) -> Self {
68        Self {
69            path: Some(path.to_owned()),
70            toml: Some(toml),
71            tables: file.tables(),
72        }
73    }
74}
75
76impl From<&SpecSource> for ManifestSpecSource {
77    fn from(source: &SpecSource) -> Self {
78        Self {
79            path: source.path.as_ref().map(|path| path.display().to_string()),
80            toml: source.toml.clone(),
81            tables: source
82                .tables
83                .iter()
84                .map(|table| ManifestDesignTable {
85                    path: table.path.display().to_string(),
86                    fnv1a64: hex(table.fnv1a64),
87                })
88                .collect(),
89        }
90    }
91}
92
93impl From<&ManifestSpecSource> for SpecSource {
94    /// Reads the source a manifest records. A table hash that is not 16 hexadecimal digits leaves its table out.
95    fn from(source: &ManifestSpecSource) -> Self {
96        Self {
97            path: source.path.as_ref().map(PathBuf::from),
98            toml: source.toml.clone(),
99            tables: source
100                .tables
101                .iter()
102                .filter_map(|table| {
103                    Some(DesignTableFile {
104                        path: PathBuf::from(&table.path),
105                        fnv1a64: u64::from_str_radix(&table.fnv1a64, 16).ok()?,
106                    })
107                })
108                .collect(),
109        }
110    }
111}
112
113/// Builds of Henad and of the host binary, and the host's command line, for the manifest.
114#[derive(Debug, Clone, PartialEq, Eq)]
115pub struct Provenance {
116    engine: RecordedBuild,
117    host: RecordedBuild,
118    arguments: Vec<String>,
119}
120
121impl Provenance {
122    /// Returns the provenance of a sweep that the host binary `host` runs with the command line `arguments`.
123    ///
124    /// The engine's build is [`ENGINE_BUILD`](crate::ENGINE_BUILD), as [`RecordedBuild::engine`] records it. A host
125    /// passes `henad_core::build_info!()` from its own crate.
126    pub fn new(host: BuildInfo, arguments: Vec<String>) -> Self {
127        Self {
128            engine: RecordedBuild::engine(),
129            host: RecordedBuild::from(&host),
130            arguments,
131        }
132    }
133
134    /// Henad's build, as [`RecordedBuild::engine`] records it.
135    pub fn engine(&self) -> &RecordedBuild {
136        &self.engine
137    }
138
139    /// Build of the host binary.
140    pub fn host(&self) -> &RecordedBuild {
141        &self.host
142    }
143
144    /// Command line of the host binary.
145    pub fn arguments(&self) -> &[String] {
146        &self.arguments
147    }
148
149    /// Returns the provenance with `engine` instead of Henad's own build.
150    #[cfg(test)]
151    pub(crate) fn with_engine(self, engine: RecordedBuild) -> Self {
152        Self { engine, ..self }
153    }
154}
155
156/// Settings of a sweep that never change its results.
157///
158/// [`Self::new`] returns the defaults, and a caller sets the other fields by assignment. A spec file's `[execution]`
159/// table goes in through [`Self::apply_execution`], before the caller's own settings.
160#[derive(Debug, Clone)]
161#[non_exhaustive]
162pub struct SweepOptions {
163    /// Number of runs stepped at once, as CPU lanes or GPU tracks.
164    pub concurrency: Concurrency,
165    /// Host memory budget in bytes for all live runs together, `None` for no limit. The lanes are sized from the probed
166    /// run, and a probed run larger than the budget leaves one lane.
167    pub memory_budget: Option<u64>,
168    /// Device memory budget in bytes for all live GPU runs together, `None` for the device's largest buffer. A run
169    /// larger than the budget runs alone.
170    pub gpu_memory_budget: Option<u64>,
171    /// Switch that pauses or aborts the sweep from another thread.
172    pub control: SweepControl,
173    /// Share of the plan's runs that this sweep executes.
174    pub shard: Shard,
175    /// Whether to add to a directory that holds runs of the same plan, running only the runs it lacks.
176    ///
177    /// A directory that holds no results starts afresh, and so does a sweep held in memory.
178    pub resume: bool,
179    /// Whether a resume reruns the runs that ended on a fault. A run that timed out always runs again.
180    ///
181    /// Note that a sweep held in memory resumes nothing, so no run runs again.
182    pub retry_failed: bool,
183    /// Table that lists each run in progress, `None` when nothing watches the runs.
184    pub(crate) active_runs: Option<ActiveRuns>,
185    /// Spec file the sweep was read from, for the manifest.
186    pub spec_source: SpecSource,
187    /// Build of the host and its command line, for the manifest.
188    pub provenance: Provenance,
189}
190
191impl SweepOptions {
192    /// Returns the options at their defaults, recording `provenance` in every manifest the sweep writes.
193    pub fn new(provenance: Provenance) -> Self {
194        Self {
195            concurrency: Concurrency::Auto,
196            memory_budget: None,
197            gpu_memory_budget: None,
198            control: SweepControl::new(),
199            shard: Shard::WHOLE,
200            resume: false,
201            retry_failed: false,
202            active_runs: None,
203            spec_source: SpecSource::default(),
204            provenance,
205        }
206    }
207
208    /// Copies the concurrency and memory settings of a spec's `[execution]` table into the options.
209    ///
210    /// Note that this overwrites all three. A caller with its own settings applies the table first and those settings
211    /// after it.
212    pub fn apply_execution(&mut self, execution: &ExecutionTable) {
213        self.concurrency = execution.concurrent;
214        self.memory_budget = execution.memory;
215        self.gpu_memory_budget = execution.gpu_memory;
216    }
217}
218
219/// Size and layout of a planned sweep or search.
220#[derive(Debug, Clone, PartialEq)]
221pub struct SweepOutline {
222    /// Id of the model.
223    pub model: String,
224    /// Backend the model runs on.
225    pub backend: Backend,
226    /// Number of configs in a sweep's plan, `None` for a search, whose budget is [`SearchOutline::max_evaluations`].
227    pub configs: Option<u64>,
228    /// Number of runs of each config, or of each evaluation of a search.
229    pub replicates: u64,
230    /// Number of runs in the whole plan, or in a search's whole budget.
231    pub runs: u64,
232    /// Blocks of the plan, each with the seed its design drew from.
233    pub blocks: Vec<PlannedBlock>,
234    /// Share of the plan's runs that this sweep executes.
235    pub shard: Shard,
236    /// Number of runs in the shard that the output directory already holds and a resume keeps.
237    pub skipped: u64,
238    /// Number of runs this sweep executes.
239    pub pending: u64,
240    /// Lanes or tracks the runs are spread over.
241    pub layout: ExecutionLayout,
242    /// Projected size in bytes of all live runs together.
243    pub projected_bytes: u64,
244    /// Number of rows that the pending runs add to `series.csv` once each reaches its final tick.
245    pub series_rows: u64,
246    /// Names of the stat columns of `series.csv`, before CSV escaping.
247    pub stat_columns: Vec<String>,
248    /// Names of the reducer columns of `runs.csv`, before CSV escaping.
249    pub reducer_columns: Vec<String>,
250    /// Whether the sweep is planned and probed alone, with nothing run.
251    pub dry_run: bool,
252    /// Budget and space of a search, `None` for a sweep.
253    pub search: Option<SearchOutline>,
254}
255
256/// Warning about a sweep that still runs, but likely not as intended.
257#[derive(Debug, Clone, PartialEq, Eq)]
258pub enum SweepWarning {
259    /// A warning of the plan.
260    Plan(PlanWarning),
261    /// A directory whose sessions so far ran the build `recorded` for `role`, and a build `current` that differs.
262    ///
263    /// A resume and a search resume compare the build that runs, `current`, against every build recorded by the
264    /// sessions that wrote runs. A merge compares the builds of the lowest shard that records builds with those of the
265    /// other shards. It sets `recorded` to a build of that shard and `current` to a build of another shard that matches
266    /// none of that shard's builds, and sets `between_shards`. Two builds that
267    /// [are treated as the same build](RecordedBuild::reads_as) are reported once, as builds that cannot be told apart.
268    BuildChanged {
269        /// Role of both builds, engine or model.
270        role: BuildRole,
271        /// Build that a session of the directory records, or a build of the lowest shard that records builds.
272        recorded: Box<RecordedBuild>,
273        /// Build that runs, or a build of another shard.
274        current: Box<RecordedBuild>,
275        /// Whether `current` is another shard's build, met by a merge, instead of the build that runs.
276        between_shards: bool,
277    },
278    /// Runs of the plan that no merged directory holds.
279    MissingRuns {
280        /// Number of missing runs.
281        count: u64,
282        /// Lowest missing run ids in ascending order, at most [`MAX_LISTED_RUNS`](crate::merge::MAX_LISTED_RUNS).
283        first: Vec<u64>,
284    },
285}
286
287impl fmt::Display for SweepWarning {
288    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
289        match self {
290            Self::Plan(warning) => warning.fmt(f),
291            Self::BuildChanged {
292                role,
293                recorded,
294                current,
295                between_shards,
296            } => {
297                let (recorded_text, current_text) = (recorded.describe(), current.describe());
298                let role_name = role.as_str();
299                if recorded.reads_as(current) {
300                    let subject = if *between_shards {
301                        format!("the shards ran the same {role_name} build")
302                    } else {
303                        format!("the directory's runs so far came from this {role_name} build")
304                    };
305                    write!(
306                        f,
307                        "cannot tell whether {subject}, {current_text}. Neither build records a commit or a source \
308                         hash"
309                    )?;
310                } else if *between_shards {
311                    write!(
312                        f,
313                        "the shards ran different {role_name} builds: {recorded_text} in one, and {current_text} in \
314                         another"
315                    )?;
316                } else {
317                    write!(
318                        f,
319                        "the directory's runs so far came from the {role_name} build {recorded_text}, and this build \
320                         is {current_text}"
321                    )?;
322                }
323                if recorded_text == current_text
324                    && let Some(difference) = recorded.crate_difference(current)
325                {
326                    write!(f, ". They differ in {difference}")?;
327                }
328                if *role == BuildRole::Model
329                    && let Some(advice) = unidentified_model_advice(recorded, current)
330                {
331                    write!(f, ". {advice}")?;
332                }
333                Ok(())
334            }
335            Self::MissingRuns { count, first } => {
336                let first: Vec<String> = first.iter().map(ToString::to_string).collect();
337                let (runs, verb, pronoun) = if *count == 1 {
338                    ("run", "is", "it")
339                } else {
340                    ("runs", "are", "them")
341                };
342                write!(
343                    f,
344                    "{count} {runs} of the plan {verb} missing, starting with {}. Resume the merged directory to run \
345                     {pronoun}",
346                    first.join(", ")
347                )
348            }
349        }
350    }
351}
352
353/// End of a sweep that did not fail.
354#[derive(Debug, Clone, Copy, PartialEq, Eq)]
355pub enum SweepEnd {
356    /// A dry run, planned and probed with nothing run.
357    Planned,
358    /// Every run of the shard is written.
359    Complete,
360    /// The control aborted the sweep. A resume runs the rest.
361    Aborted,
362    /// The GPU device was lost. A resume runs the rest.
363    DeviceLost,
364}
365
366/// Result of a sweep that did not fail.
367#[derive(Debug, Clone, PartialEq)]
368pub struct SweepReport {
369    /// Outline of the sweep, as [`ProgressEvent::Planned`] carries it.
370    pub outline: SweepOutline,
371    /// End of the sweep.
372    pub end: SweepEnd,
373    /// Rows in `runs.csv` once the sweep ends, by status, the rows a resume kept included. Zero for a dry run.
374    pub counts: ResultCounts,
375    /// Time from the start of planning to the end of the sweep.
376    pub elapsed: Duration,
377    /// Directory the results were written to, `None` for a dry run or a sweep held in memory.
378    pub output_dir: Option<PathBuf>,
379}
380
381/// Report, manifest and files of a sweep that ran to its end or was aborted.
382#[derive(Debug, Clone, PartialEq)]
383pub struct SweepRecord {
384    /// Report of the sweep, as [`ProgressEvent::Ended`] carries it.
385    pub report: SweepReport,
386    /// Manifest with the sweep's final status, as `manifest.json` holds it.
387    pub manifest: Manifest,
388    /// Files of a sweep held in memory, `None` for a sweep written to a directory.
389    pub files: Option<SweepFiles>,
390    /// Standing of a search at its end, `None` for a sweep.
391    pub search: Option<SearchReport>,
392}
393
394/// A sweep that cannot run to its end.
395///
396/// The variants differ between targets, and a match outside this crate ends in a wildcard arm.
397#[derive(Debug)]
398#[non_exhaustive]
399pub enum ExploreError {
400    /// A spec that the model rejects.
401    Plan(PlanError),
402    /// Configs the device cannot host.
403    Capacity(CapacityError),
404    /// A probe build that failed.
405    Probe(ProbeError),
406    /// Reducers that do not bind to the columns of the probe build.
407    Measure(MeasureError),
408    /// A directory the sweep cannot resume into.
409    Resume(ResumeError),
410    /// Results that cannot be written.
411    Output(OutputError),
412    /// A batch of runs that cannot run.
413    Execution(ExecutionError),
414    /// A search spec that the model or the options reject.
415    Search(SearchPlanError),
416    /// A GPU model with no device, on a machine where no device can be acquired.
417    #[cfg(not(target_arch = "wasm32"))]
418    Device(crate::device::DeviceError),
419}
420
421impl fmt::Display for ExploreError {
422    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
423        match self {
424            Self::Plan(_) => f.write_str("cannot plan the sweep"),
425            Self::Capacity(_) => f.write_str("the sweep does not fit this GPU"),
426            Self::Probe(_) => f.write_str("cannot probe the model"),
427            Self::Measure(_) => f.write_str("cannot bind the reducers to the model's stat columns"),
428            Self::Resume(_) => f.write_str("cannot resume the sweep"),
429            Self::Output(_) => f.write_str("cannot write the results"),
430            Self::Execution(_) => f.write_str("cannot run the sweep"),
431            Self::Search(_) => f.write_str("cannot plan the search"),
432            #[cfg(not(target_arch = "wasm32"))]
433            Self::Device(_) => f.write_str("cannot acquire a GPU device for the sweep"),
434        }
435    }
436}
437
438impl std::error::Error for ExploreError {
439    fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
440        match self {
441            Self::Plan(error) => Some(error),
442            Self::Capacity(error) => Some(error),
443            Self::Probe(error) => Some(error),
444            Self::Measure(error) => Some(error),
445            Self::Resume(error) => Some(error),
446            Self::Output(error) => Some(error),
447            Self::Execution(error) => Some(error),
448            Self::Search(error) => Some(error),
449            #[cfg(not(target_arch = "wasm32"))]
450            Self::Device(error) => Some(error),
451        }
452    }
453}
454
455impl From<PlanError> for ExploreError {
456    fn from(error: PlanError) -> Self {
457        Self::Plan(error)
458    }
459}
460
461impl From<CapacityError> for ExploreError {
462    fn from(error: CapacityError) -> Self {
463        Self::Capacity(error)
464    }
465}
466
467impl From<ProbeError> for ExploreError {
468    fn from(error: ProbeError) -> Self {
469        Self::Probe(error)
470    }
471}
472
473impl From<MeasureError> for ExploreError {
474    fn from(error: MeasureError) -> Self {
475        Self::Measure(error)
476    }
477}
478
479impl From<ResumeError> for ExploreError {
480    fn from(error: ResumeError) -> Self {
481        Self::Resume(error)
482    }
483}
484
485impl From<OutputError> for ExploreError {
486    fn from(error: OutputError) -> Self {
487        Self::Output(error)
488    }
489}
490
491impl From<ExecutionError> for ExploreError {
492    fn from(error: ExecutionError) -> Self {
493        Self::Execution(error)
494    }
495}
496
497impl From<SearchPlanError> for ExploreError {
498    fn from(error: SearchPlanError) -> Self {
499        Self::Search(error)
500    }
501}
502
503/// Runs the sweep, or the search that a spec's `[search]` table describes, and blocks until it ends. Native only.
504///
505/// A sweep plans the spec, checks every config against the device, probes the first config without a fault and the
506/// last config, and chooses a layout from the larger probe. It then writes the manifest with status `running`,
507/// streams each run to `runs.csv` and `series.csv` in plan order, rebuilds `summary.csv` from `runs.csv`, and
508/// replaces the manifest with its final status. Runs reach the output in plan order whatever order they finish in,
509/// and a search's runs in the order that it requests them.
510///
511/// `gpu` is the device a GPU model runs on, shared with the host, its
512/// [`FaultSink`](henad_compute::fault::FaultSink) included. Note that the sink holds one fault, and whichever side
513/// reads it first takes it. A fault that the sweep takes first ends every live run, whichever side raised it, and a
514/// fault that the host takes first leaves the runs going. A host that renders on its device passes `None`. When `gpu`
515/// is `None`, a device for a GPU model is acquired once the spec is planned, sized to its [`ModelEntry::gpu_needs`].
516///
517/// The manifest records the adapter of the device the sweep runs on when its context carries
518/// [`RuntimeInfo`](henad_compute::runtime_info::RuntimeInfo). A context acquired here carries it, and a host's own
519/// context carries it once built with [`GpuContext::with_runtime_info`]. A bare [`SweepSpec`] carries no execution
520/// settings. [`LoadedSpec`](crate::spec_file::LoadedSpec) keeps a spec file's `[execution]` table for
521/// [`SweepOptions::apply_execution`].
522///
523/// With `options.resume`, a directory holding runs of the same plan keeps the runs that its [`ResumeScan`]
524/// keeps, and the sweep runs the rest. Both tables then list every run in order of its id, as a sweep run in one go
525/// would. A sweep held in memory starts afresh whatever `options.resume` and `options.retry_failed` say.
526/// A run that faults is recorded with its status, and the sweep carries on. A sweep that fails once its manifest is
527/// written marks the manifest `failed` when it can.
528///
529/// # Errors
530///
531/// Returns [`ExploreError`] when the spec cannot be planned, no device can be acquired, a config does not fit the
532/// device, the probe build fails, the output directory holds results and is not resumed, another sweep, search or
533/// merge is writing to it, the directory cannot be resumed, the results cannot be written, or a batch cannot run.
534#[cfg(not(target_arch = "wasm32"))]
535pub fn run_spec(
536    model: &ModelEntry,
537    gpu: Option<&GpuContext>,
538    spec: &SweepSpec,
539    output: SweepOutput,
540    options: &SweepOptions,
541    progress: &mut dyn Progress,
542) -> Result<SweepRecord, ExploreError> {
543    use crate::search_run::{run_search_in_memory, run_search_into_directory};
544
545    let (planned, device) = plan_then_acquire(model, gpu, spec, crate::device::acquire_headless)?;
546    let gpu = device.as_deref();
547    let runtime = ManifestRuntime::new(gpu.and_then(GpuContext::runtime_info));
548    let dir = match output {
549        SweepOutput::Memory => None,
550        SweepOutput::Directory(dir) => Some(dir),
551    };
552    let folder = dir.as_deref();
553    let inputs = SweepInputs {
554        entry: model,
555        gpu,
556        runtime: &runtime,
557        spec,
558        source: &options.spec_source,
559        provenance: &options.provenance,
560        options,
561        folder,
562        dry_run: false,
563    };
564    match (folder, planned) {
565        (None, SpecPlan::Sweep(plan)) => run_in_memory(&inputs, Some(plan), progress),
566        (Some(dir), SpecPlan::Sweep(plan)) => run_into_directory(&inputs, Some(plan), dir, progress),
567        (None, SpecPlan::Search(plan)) => run_search_in_memory(&inputs, Some(plan), progress),
568        (Some(dir), SpecPlan::Search(plan)) => run_search_into_directory(&inputs, Some(plan), dir, progress),
569    }
570}
571
572/// Plan of a spec, for a sweep or a search.
573#[cfg(not(target_arch = "wasm32"))]
574enum SpecPlan {
575    Sweep(Arc<Plan>),
576    Search(Arc<crate::search_run::SearchPlan>),
577}
578
579#[cfg(not(target_arch = "wasm32"))]
580impl SpecPlan {
581    /// Plans `spec` against `entry`, as a search when it has a `[search]` table.
582    fn new(entry: &ModelEntry, spec: &SweepSpec) -> Result<Self, ExploreError> {
583        let schema = entry.schema();
584        Ok(if spec.search.is_some() {
585            Self::Search(Arc::new(crate::search_run::SearchPlan::new(spec, &schema)?))
586        } else {
587            Self::Sweep(Arc::new(spec.plan(&schema)?))
588        })
589    }
590}
591
592/// Plans `spec` and probes its configs as `--dry-run` does, writing nothing. Native only.
593///
594/// `gpu` is a device the host shares with the probe builds, its [`FaultSink`](henad_compute::fault::FaultSink)
595/// included, as [`run_spec`] describes. When `gpu` is `None`, a device for a GPU model is acquired for its
596/// [`ModelEntry::gpu_needs`] once the spec is planned, since a probe builds the model. With `folder` and without
597/// `options.resume`, rejects a folder that holds results, as `--out` does. With both, reads the folder as a resume
598/// would, counts the runs it would skip and run, and compares the recorded builds.
599///
600/// # Errors
601///
602/// Returns [`ExploreError`] when the spec cannot be planned, no device can be acquired, a config does not fit the
603/// device, the probe build fails, `folder` holds results and is not resumed, or `folder` cannot be resumed.
604#[cfg(not(target_arch = "wasm32"))]
605pub fn plan_spec(
606    model: &ModelEntry,
607    gpu: Option<&GpuContext>,
608    spec: &SweepSpec,
609    folder: Option<&Path>,
610    options: &SweepOptions,
611    progress: &mut dyn Progress,
612) -> Result<SweepReport, ExploreError> {
613    use crate::search_run::SearchPreparation;
614
615    let (planned, device) = plan_then_acquire(model, gpu, spec, crate::device::acquire_headless)?;
616    let gpu = device.as_deref();
617    let runtime = ManifestRuntime::new(gpu.and_then(GpuContext::runtime_info));
618    let inputs = SweepInputs {
619        entry: model,
620        gpu,
621        runtime: &runtime,
622        spec,
623        source: &options.spec_source,
624        provenance: &options.provenance,
625        options,
626        folder,
627        dry_run: true,
628    };
629    let report = match planned {
630        SpecPlan::Search(plan) => {
631            let preparation = SearchPreparation::new(&inputs, Some(plan), None)?;
632            preparation.announce(&inputs, progress);
633            preparation.report(SweepEnd::Planned, ResultCounts::default(), None)
634        }
635        SpecPlan::Sweep(plan) => {
636            let preparation = SweepPreparation::new(&inputs, Some(plan), None)?;
637            preparation.announce(&inputs, progress);
638            preparation.report(SweepEnd::Planned, ResultCounts::default(), None)
639        }
640    };
641    progress.report(&ProgressEvent::Ended(&report));
642    Ok(report)
643}
644
645/// Plans `spec` against `model`, then returns the plan and the device a sweep of it runs on, acquired through
646/// `acquire` for a GPU model when `gpu` is `None`.
647///
648/// A spec that the model rejects reports its own error before any device is acquired, on a machine without an adapter
649/// as well.
650///
651/// # Errors
652///
653/// Returns [`ExploreError`] when the spec cannot be planned, and [`ExploreError::Device`] when a GPU model needs a
654/// device and `acquire` fails.
655#[cfg(not(target_arch = "wasm32"))]
656fn plan_then_acquire<'a>(
657    model: &ModelEntry,
658    gpu: Option<&'a GpuContext>,
659    spec: &SweepSpec,
660    acquire: impl FnOnce(henad_compute::gpu::GpuNeeds) -> Result<GpuContext, crate::device::DeviceError>,
661) -> Result<(SpecPlan, Option<Cow<'a, GpuContext>>), ExploreError> {
662    let planned = SpecPlan::new(model, spec)?;
663    let device = device_or_acquired(model, gpu, acquire)?;
664    Ok((planned, device))
665}
666
667/// Returns the device a sweep of `entry` runs on: `gpu` for a GPU model, or a device acquired for its needs when
668/// `gpu` is `None`.
669///
670/// Returns `None` for a CPU model, whatever `gpu` is. A CPU sweep uses no device, and its manifest records no
671/// adapter.
672///
673/// # Errors
674///
675/// Returns [`ExploreError::Device`] when a GPU model needs a device and no device can be acquired.
676#[cfg(not(target_arch = "wasm32"))]
677pub(crate) fn sweep_device<'a>(
678    entry: &ModelEntry,
679    gpu: Option<&'a GpuContext>,
680) -> Result<Option<Cow<'a, GpuContext>>, ExploreError> {
681    device_or_acquired(entry, gpu, crate::device::acquire_headless)
682}
683
684/// Returns the device a sweep of `entry` runs on, as [`sweep_device`] does with `acquire` instead of
685/// [`acquire_headless`](crate::device::acquire_headless).
686///
687/// # Errors
688///
689/// Returns [`ExploreError::Device`] when a GPU model needs a device and `acquire` fails.
690#[cfg(not(target_arch = "wasm32"))]
691fn device_or_acquired<'a>(
692    entry: &ModelEntry,
693    gpu: Option<&'a GpuContext>,
694    acquire: impl FnOnce(henad_compute::gpu::GpuNeeds) -> Result<GpuContext, crate::device::DeviceError>,
695) -> Result<Option<Cow<'a, GpuContext>>, ExploreError> {
696    match (entry.gpu_needs(), gpu) {
697        (None, _) => Ok(None),
698        (Some(_), Some(ctx)) => Ok(Some(Cow::Borrowed(ctx))),
699        (Some(needs), None) => {
700            let ctx = acquire(needs).map_err(ExploreError::Device)?;
701            Ok(Some(Cow::Owned(ctx)))
702        }
703    }
704}
705
706/// Runs the plan of `inputs.spec` into the directory `output_dir`, as [`run_spec`] does, and returns its record.
707///
708/// `plan`, when given, is the plan of `inputs.spec`, and the spec is planned here otherwise.
709///
710/// # Errors
711///
712/// Returns the errors of [`run_spec`].
713#[cfg(not(target_arch = "wasm32"))]
714pub(crate) fn run_into_directory(
715    inputs: &SweepInputs<'_>,
716    plan: Option<Arc<Plan>>,
717    output_dir: &Path,
718    progress: &mut dyn Progress,
719) -> Result<SweepRecord, ExploreError> {
720    let preparation = SweepPreparation::new(inputs, plan, None)?;
721    preparation.announce(inputs, progress);
722    let record = preparation.write_directory(inputs, output_dir, progress)?;
723    progress.report(&ProgressEvent::Ended(&record.report));
724    Ok(record)
725}
726
727/// Runs the plan of `inputs.spec`, holding its four files in memory, and returns its record.
728///
729/// `plan`, when given, is the plan of `inputs.spec`, and the spec is planned here otherwise. The files hold the bytes
730/// a directory would. The caller passes `inputs` with no `folder`.
731///
732/// # Errors
733///
734/// Returns the errors of [`run_spec`] that do not come from a directory.
735#[cfg(not(target_arch = "wasm32"))]
736pub(crate) fn run_in_memory(
737    inputs: &SweepInputs<'_>,
738    plan: Option<Arc<Plan>>,
739    progress: &mut dyn Progress,
740) -> Result<SweepRecord, ExploreError> {
741    use crate::output::memory::memory_writer;
742
743    let preparation = SweepPreparation::new(inputs, plan, None)?;
744    preparation.announce(inputs, progress);
745    let mut manifest = preparation.manifest(inputs)?;
746    let writer = memory_writer(
747        &preparation.plan,
748        inputs.entry.param_descriptors(),
749        &preparation.measure,
750    )?;
751    let (end, writer) = preparation.run_pending(inputs, writer, progress)?;
752    let counts = writer.counts();
753    let end = finish_manifest(&mut manifest, end, counts);
754    let files = SweepFiles::assemble(writer, &manifest)?;
755    let record = SweepRecord {
756        report: preparation.report(end, counts, None),
757        manifest,
758        files: Some(files),
759        search: None,
760    };
761    progress.report(&ProgressEvent::Ended(&record.report));
762    Ok(record)
763}
764
765/// Model, device, spec and records of one sweep, and the host's settings.
766pub(crate) struct SweepInputs<'a> {
767    pub(crate) entry: &'a ModelEntry,
768    pub(crate) gpu: Option<&'a GpuContext>,
769    pub(crate) runtime: &'a ManifestRuntime,
770    pub(crate) spec: &'a SweepSpec,
771    pub(crate) source: &'a SpecSource,
772    pub(crate) provenance: &'a Provenance,
773    pub(crate) options: &'a SweepOptions,
774    /// Directory the results go in or a dry run reads, `None` for a sweep held in memory.
775    pub(crate) folder: Option<&'a Path>,
776    /// Whether the sweep is planned and probed and stops there, writing nothing.
777    pub(crate) dry_run: bool,
778}
779
780/// A sweep planned, checked and probed, with its layout chosen.
781pub(crate) struct SweepPreparation {
782    plan: Arc<Plan>,
783    probe: ProbeReport,
784    measure: Arc<MeasurePlan>,
785    /// Runs of the shard that this sweep executes, in plan order.
786    pending: Vec<PlannedRun>,
787    /// Directory the sweep resumes, `None` for a sweep that starts afresh.
788    resumed: Option<ResumeScan>,
789    /// Directory a resume holds locked from before its scan until its last write, `None` for a dry run or a sweep
790    /// that starts afresh.
791    #[cfg_attr(
792        target_arch = "wasm32",
793        expect(dead_code, reason = "a sweep in a browser writes no directory")
794    )]
795    locked: Option<OutputDir>,
796    outline: SweepOutline,
797    /// Clock reading when planning started.
798    started: Instant,
799    started_unix_ms: u64,
800}
801
802impl SweepPreparation {
803    /// Checks and probes `plan`, the plan of `inputs.spec`, and chooses its layout. With no `plan`, the spec is
804    /// planned first. `probe`, when given, replaces the report that [`ProbeReport::for_plan`] returns for the plan,
805    /// and its clock readings replace the ones taken on entry.
806    pub(crate) fn new(
807        inputs: &SweepInputs<'_>,
808        plan: Option<Arc<Plan>>,
809        probe: Option<TimedProbe>,
810    ) -> Result<Self, ExploreError> {
811        let (entry, gpu, options) = (inputs.entry, inputs.gpu, inputs.options);
812        debug_assert!(
813            inputs.spec.search.is_none(),
814            "a spec with a [search] table runs as a search"
815        );
816        let (started, started_unix_ms, probe) = match probe {
817            Some(timed) => (timed.started, timed.started_unix_ms, Some(timed.report)),
818            None => (Instant::now(), now_unix_ms(), None),
819        };
820        let plan = match plan {
821            Some(plan) => plan,
822            None => Arc::new(inputs.spec.plan(&entry.schema())?),
823        };
824        let resume_dir = inputs
825            .folder
826            .filter(|output_dir| options.resume && OutputDir::holds_results(output_dir));
827        let locked = match (resume_dir, inputs.folder) {
828            (None, Some(output_dir)) => {
829                OutputDir::check_free(output_dir)?;
830                None
831            }
832            // The directory is locked before the scan and held until the last write. Otherwise another writer could
833            // change the tables between the two.
834            (Some(output_dir), _) if !inputs.dry_run => Some(OutputDir::open(output_dir)?),
835            _ => None,
836        };
837        if let Some(ctx) = gpu {
838            check_capacity(entry, &plan, &ctx.device.limits())?;
839        }
840        let probe = match probe {
841            Some(probe) => probe,
842            None => ProbeReport::for_plan(entry, gpu, &plan)?,
843        };
844        let measure = MeasurePlan::new(plan.run_settings(), plan.measure_settings(), probe.columns.clone())?;
845        let resumed = resume_dir
846            .map(|output_dir| {
847                let columns = column_names(entry.param_descriptors(), plan.actions(), measure.reducers().names());
848                ResumeScan::read(
849                    output_dir,
850                    &plan,
851                    options.shard,
852                    options.retry_failed,
853                    &runs_csv::header_line(&columns),
854                    &series_csv::header_line(measure.columns()),
855                )
856            })
857            .transpose()?;
858        let pending: Vec<PlannedRun> = plan
859            .runs_in_shard(options.shard)
860            .filter(|run| {
861                resumed
862                    .as_ref()
863                    .is_none_or(|scan| !scan.finished().contains(&run.run_id))
864            })
865            .collect();
866
867        // The layout is sized from the first config that builds without a fault or the last config, whichever holds
868        // more.
869        let last_probe = ProbeReport::for_last_config(entry, gpu, &plan, &probe);
870        let sizing_probe = match &last_probe {
871            Some(last_probe) if last_probe.footprint() > probe.footprint() => last_probe,
872            _ => &probe,
873        };
874        let pending_count = pending.len() as u64;
875        let (layout, projected_bytes) = sized_layout(entry, gpu, options, sizing_probe, pending_count)?;
876        let outline = SweepOutline {
877            model: entry.id().to_owned(),
878            backend: entry.metadata().backend,
879            configs: Some(plan.configs().len() as u64),
880            replicates: plan.replicates(),
881            runs: plan.run_count(),
882            blocks: plan.blocks().to_vec(),
883            shard: options.shard,
884            skipped: resumed.as_ref().map_or(0, |scan| scan.finished().len() as u64),
885            pending: pending_count,
886            layout,
887            projected_bytes,
888            series_rows: pending_count.saturating_mul(measure.series_row_count()),
889            stat_columns: (0..measure.columns().len())
890                .map(|column| measure.columns().name(column).to_owned())
891                .collect(),
892            reducer_columns: measure.reducers().names().to_vec(),
893            dry_run: inputs.dry_run,
894            search: None,
895        };
896        Ok(Self {
897            plan,
898            probe,
899            measure: Arc::new(measure),
900            pending,
901            resumed,
902            locked,
903            outline,
904            started,
905            started_unix_ms,
906        })
907    }
908
909    #[cfg(any(target_arch = "wasm32", test))]
910    pub(crate) fn plan(&self) -> &Arc<Plan> {
911        &self.plan
912    }
913
914    #[cfg(any(target_arch = "wasm32", test))]
915    pub(crate) fn measure(&self) -> &Arc<MeasurePlan> {
916        &self.measure
917    }
918
919    /// Runs this sweep executes, in plan order.
920    #[cfg(any(target_arch = "wasm32", test))]
921    pub(crate) fn pending(&self) -> &[PlannedRun] {
922        &self.pending
923    }
924
925    /// Reports the outline, then each warning of the plan and of a resume under a build that differs from the engine
926    /// or model build of `inputs`.
927    pub(crate) fn announce(&self, inputs: &SweepInputs<'_>, progress: &mut dyn Progress) {
928        progress.report(&ProgressEvent::Planned(&self.outline));
929        for warning in self.warnings(inputs.provenance, inputs.entry) {
930            progress.report(&ProgressEvent::Warned(&warning));
931        }
932    }
933
934    /// Returns the warnings of the plan, and of a resume under a build that differs from the engine build of
935    /// `provenance` or the model build of `entry`.
936    fn warnings(&self, provenance: &Provenance, entry: &ModelEntry) -> Vec<SweepWarning> {
937        let mut warnings: Vec<SweepWarning> = self.plan.warnings().iter().cloned().map(SweepWarning::Plan).collect();
938        if let Some(scan) = &self.resumed {
939            warnings.extend(build_warnings(&scan.recorded, provenance, entry));
940        }
941        warnings
942    }
943
944    /// Writes the manifest with status `running`, runs every pending run into `output_dir` and replaces the manifest
945    /// with the sweep's final status.
946    #[cfg(not(target_arch = "wasm32"))]
947    fn write_directory(
948        &self,
949        inputs: &SweepInputs<'_>,
950        output_dir: &Path,
951        progress: &mut dyn Progress,
952    ) -> Result<SweepRecord, ExploreError> {
953        let created;
954        let dir = if let Some(dir) = &self.locked {
955            dir
956        } else {
957            created = OutputDir::create(output_dir)?;
958            &created
959        };
960        let mut manifest = self.manifest(inputs)?;
961        dir.write_manifest(&manifest)?;
962        let (end, counts) = match self.run_into(dir, inputs, progress) {
963            Ok(finished) => finished,
964            Err(error) => {
965                manifest.fail(now_unix_ms());
966                // The sweep's own error is the one to report. A manifest that cannot be written says `running`.
967                drop(dir.write_manifest(&manifest));
968                return Err(error);
969            }
970        };
971        let end = finish_manifest(&mut manifest, end, counts);
972        dir.write_manifest(&manifest)?;
973        Ok(SweepRecord {
974            report: self.report(end, counts, Some(output_dir.to_owned())),
975            manifest,
976            files: None,
977            search: None,
978        })
979    }
980
981    /// Repairs a resumed directory, runs every pending run into `dir`, puts the tables in run order and rebuilds the
982    /// summary.
983    ///
984    /// Returns the end of the batch and the counts of the rows `runs.csv` holds.
985    #[cfg(not(target_arch = "wasm32"))]
986    fn run_into(
987        &self,
988        dir: &OutputDir,
989        inputs: &SweepInputs<'_>,
990        progress: &mut dyn Progress,
991    ) -> Result<(BatchEnd, ResultCounts), ExploreError> {
992        if let Some(scan) = &self.resumed {
993            scan.repair(dir)?;
994        }
995        let params = inputs.entry.param_descriptors();
996        let writer = if self.resumed.is_some() {
997            dir.append_writer(&self.plan, params)?
998        } else {
999            dir.open_writer(&self.plan, params, &self.measure)?
1000        };
1001        let (end, writer) = self.run_pending(inputs, writer, progress)?;
1002        let written = writer.counts();
1003        writer.finish().map_err(|source| OutputError::Write {
1004            path: dir.path().to_owned(),
1005            source,
1006        })?;
1007        let kept = self.resumed.as_ref().map(ResumeScan::counts).unwrap_or_default();
1008        let last_kept = self.resumed.as_ref().and_then(|scan| scan.finished().last().copied());
1009        if self.pending.first().is_some_and(|run| Some(run.run_id) < last_kept) {
1010            dir.order_tables()?;
1011        }
1012        dir.write_summary()?;
1013        Ok((end, kept + written))
1014    }
1015
1016    /// Runs every pending run and commits each to `writer` in plan order, reporting each to `progress`.
1017    ///
1018    /// Returns the end of the batch and the writer.
1019    #[cfg(not(target_arch = "wasm32"))]
1020    fn run_pending<W: Write>(
1021        &self,
1022        inputs: &SweepInputs<'_>,
1023        writer: OutputWriter<W>,
1024        progress: &mut dyn Progress,
1025    ) -> Result<(BatchEnd, OutputWriter<W>), ExploreError> {
1026        let options = inputs.options;
1027        let executor = Executor::new(
1028            inputs.entry,
1029            inputs.gpu,
1030            Arc::clone(&self.measure),
1031            self.outline.layout,
1032            options.control.clone(),
1033        )?
1034        .with_timeout(self.plan.run_settings().timeout)
1035        .with_active_runs(options.active_runs.clone())
1036        .with_gpu_memory_budget(options.gpu_memory_budget);
1037        let requests: Vec<RunRequest<'_>> = self
1038            .pending
1039            .iter()
1040            .map(|&run| RunRequest::planned(&self.plan, run))
1041            .collect();
1042        let mut sink = SweepSink {
1043            writer,
1044            progress,
1045            meter: ProgressMeter::new(self.outline.pending),
1046        };
1047        let end = executor.run_batch(&requests, &mut sink)?;
1048        Ok((end, sink.writer))
1049    }
1050
1051    /// Returns the manifest of the sweep while it runs.
1052    ///
1053    /// A resume keeps the start time, sessions and merged shards the directory recorded.
1054    ///
1055    /// # Errors
1056    ///
1057    /// Returns [`OutputError::Manifest`] when the spec cannot be written as JSON.
1058    pub(crate) fn manifest(&self, inputs: &SweepInputs<'_>) -> Result<Manifest, OutputError> {
1059        running_manifest(
1060            inputs,
1061            &ManifestParts {
1062                mode: ManifestMode::Sweep,
1063                plan: &self.plan,
1064                probe: &self.probe,
1065                outline: &self.outline,
1066                manifest_plan: self.manifest_plan(),
1067                started_unix_ms: self.started_unix_ms,
1068                recorded: self.resumed.as_ref().map(|scan| &scan.recorded),
1069                search: None,
1070            },
1071        )
1072    }
1073
1074    /// Returns the record of the plan's hashes, counts and blocks.
1075    fn manifest_plan(&self) -> ManifestPlan {
1076        ManifestPlan {
1077            plan_hash: hex(self.plan.plan_hash()),
1078            results_fingerprint: hex(self.plan.results_fingerprint()),
1079            configs: self.outline.configs,
1080            replicates: self.outline.replicates,
1081            runs: self.outline.runs,
1082            blocks: self
1083                .plan
1084                .blocks()
1085                .iter()
1086                .map(|block| ManifestBlock {
1087                    design: block.design.as_str().to_owned(),
1088                    configs: block.configs.end - block.configs.start,
1089                    design_seed: block.design_seed,
1090                })
1091                .collect(),
1092        }
1093    }
1094
1095    pub(crate) fn report(&self, end: SweepEnd, counts: ResultCounts, output_dir: Option<PathBuf>) -> SweepReport {
1096        SweepReport {
1097            outline: self.outline.clone(),
1098            end,
1099            counts,
1100            elapsed: self.started.elapsed(),
1101            output_dir,
1102        }
1103    }
1104}
1105
1106/// Returns the layout for `runs` runs of `entry` sized from `sizing_probe`, and the bytes the live runs are projected
1107/// to hold.
1108///
1109/// Note that a probe sizes its buffers to rayon's global pool, and a run in a lane to the lane's pool. When the
1110/// layout splits the workers into lanes, the probe is built again on a pool as wide as a lane and the lanes are
1111/// budgeted from that build.
1112///
1113/// # Errors
1114///
1115/// Returns [`ExploreError::Probe`] when the probe cannot be built again on a lane's pool.
1116pub(crate) fn sized_layout(
1117    entry: &ModelEntry,
1118    gpu: Option<&GpuContext>,
1119    options: &SweepOptions,
1120    sizing_probe: &ProbeReport,
1121    runs: u64,
1122) -> Result<(ExecutionLayout, u64), ExploreError> {
1123    let backend = entry.metadata().backend;
1124    let detected = ExecutionBudget::detect();
1125    let unbudgeted = choose_layout(options.concurrency, &detected, backend, sizing_probe, runs);
1126    let rebuilt = if unbudgeted.cpu_lanes > 1 && unbudgeted.threads_per_lane != detected.workers {
1127        Some(sizing_probe.rebuilt_on(entry, unbudgeted.threads_per_lane)?)
1128    } else {
1129        None
1130    };
1131    let resources = ExecutionBudget {
1132        memory_budget: options.memory_budget,
1133        gpu_memory_budget: gpu.map(|ctx| gpu_memory_budget(options.gpu_memory_budget, ctx)),
1134        ..detected
1135    };
1136    let lane_probe = rebuilt.as_ref().unwrap_or(sizing_probe);
1137    let layout = choose_layout(options.concurrency, &resources, backend, lane_probe, runs);
1138    let run_probe = if layout.cpu_lanes > 1 { lane_probe } else { sizing_probe };
1139    Ok((layout, layout.projected_bytes(run_probe)))
1140}
1141
1142/// Returns the advice for a model build that records neither a commit nor a source hash, checking the build that
1143/// runs first, or `None` when `recorded` and `current` both record a commit or a source hash.
1144///
1145/// An entry that was never inserted into a [`ModelSet`](henad_compute::entry::ModelSet) records no build at all, and
1146/// a build script alone cannot provide a build.
1147fn unidentified_model_advice(recorded: &RecordedBuild, current: &RecordedBuild) -> Option<&'static str> {
1148    let unidentified = [current, recorded].into_iter().find(|build| !build.is_identified())?;
1149    Some(if unidentified.package.is_empty() {
1150        "The model's entry records no build. Insert it into a `ModelSet` built with `henad::build_info!()` in a crate \
1151         whose build script calls `henad_build::stamp_commit()`"
1152    } else {
1153        "The model's build records no commit or source hash. Call `henad_build::stamp_commit()` in the build script \
1154         of the crate that registers it"
1155    })
1156}
1157
1158/// Returns a [`SweepWarning::BuildChanged`] for each build that [`Manifest::recorded_builds`] lists for `recorded` and
1159/// differs from the engine build of `provenance`, or from the build that registered `entry`.
1160pub(crate) fn build_warnings(recorded: &Manifest, provenance: &Provenance, entry: &ModelEntry) -> Vec<SweepWarning> {
1161    let model = RecordedBuild::from(entry.source());
1162    [(BuildRole::Engine, provenance.engine()), (BuildRole::Model, &model)]
1163        .into_iter()
1164        .flat_map(|(role, current)| {
1165            recorded
1166                .recorded_builds(role)
1167                .into_iter()
1168                .filter(|build| !build.same_build(current))
1169                .map(move |build| SweepWarning::BuildChanged {
1170                    role,
1171                    recorded: Box::new(build),
1172                    current: Box::new(current.clone()),
1173                    between_shards: false,
1174                })
1175        })
1176        .collect()
1177}
1178
1179/// Parts of a manifest that differ between a sweep and a search.
1180pub(crate) struct ManifestParts<'a> {
1181    pub(crate) mode: ManifestMode,
1182    /// Plan the runs come from. For a search, the plan of its fixed values alone.
1183    pub(crate) plan: &'a Plan,
1184    pub(crate) probe: &'a ProbeReport,
1185    pub(crate) outline: &'a SweepOutline,
1186    pub(crate) manifest_plan: ManifestPlan,
1187    pub(crate) started_unix_ms: u64,
1188    /// Manifest of the directory a resume adds to, `None` for a fresh start.
1189    pub(crate) recorded: Option<&'a Manifest>,
1190    pub(crate) search: Option<ManifestSearch>,
1191}
1192
1193/// Returns the manifest, with status `running`, of the sweep or search that `inputs` and `parts` describe.
1194///
1195/// A resume keeps the start time, sessions and merged shards the directory recorded.
1196///
1197/// # Errors
1198///
1199/// Returns [`OutputError::Manifest`] when the spec cannot be written as JSON.
1200pub(crate) fn running_manifest(inputs: &SweepInputs<'_>, parts: &ManifestParts<'_>) -> Result<Manifest, OutputError> {
1201    let (entry, provenance, options) = (inputs.entry, inputs.provenance, inputs.options);
1202    let backend = backend_name(entry.metadata().backend).to_owned();
1203    let mut spec_file = SpecFile::from(inputs.spec);
1204    spec_file.execution = ExecutionTable {
1205        concurrent: options.concurrency,
1206        memory: options.memory_budget,
1207        gpu_memory: options.gpu_memory_budget,
1208    };
1209    let outline = parts.outline;
1210    let layout = outline.layout;
1211    let seeds = parts.plan.seed_settings();
1212    let spec = serde_json::to_value(&spec_file).map_err(OutputError::Manifest)?;
1213    let session = ManifestSession {
1214        started: rfc3339(parts.started_unix_ms),
1215        commit: provenance.engine().commit.clone(),
1216        skipped: outline.skipped,
1217        ran: 0,
1218        engine: Some(provenance.engine().clone()),
1219        host: Some(provenance.host().clone()),
1220        model_source: Some(RecordedBuild::from(entry.source())),
1221    };
1222    let (started_unix_ms, sessions, merged_shards) = match parts.recorded {
1223        Some(recorded) => {
1224            let mut recorded = recorded.clone();
1225            recorded.record_session_engines();
1226            let status = recorded.status;
1227            let mut sessions = recorded.sessions;
1228            // A session whose process ended before `Manifest::finish` never had its runs counted.
1229            if matches!(status, ManifestStatus::Running | ManifestStatus::Failed)
1230                && let Some(last) = sessions.last_mut()
1231            {
1232                last.ran = outline.skipped.saturating_sub(last.skipped);
1233            }
1234            sessions.push(session);
1235            (recorded.timestamps.started_unix_ms, sessions, recorded.merged_shards)
1236        }
1237        None => (parts.started_unix_ms, vec![session], None),
1238    };
1239    Ok(Manifest {
1240        format: FORMAT.to_owned(),
1241        format_version: FORMAT_VERSION,
1242        mode: parts.mode,
1243        status: ManifestStatus::Running,
1244        engine: provenance.engine().clone(),
1245        model: ManifestModel {
1246            id: entry.id().to_owned(),
1247            name: entry.name().to_owned(),
1248            backend: backend.clone(),
1249            schema_hash: hex(parts.plan.schema_hash()),
1250            schema: schema_json(entry, Some(parts.probe)),
1251            replays_exactly: entry.metadata().replays_exactly,
1252        },
1253        spec,
1254        spec_source: inputs.source.into(),
1255        argv: provenance.arguments().to_vec(),
1256        plan: parts.manifest_plan.clone(),
1257        seeds: ManifestSeeds {
1258            root: seeds.root,
1259            scheme: seeds.scheme.as_str().to_owned(),
1260            formula: seeds.scheme.formula().to_owned(),
1261        },
1262        columns: ManifestColumns {
1263            stats: outline.stat_columns.clone(),
1264            reducers: outline.reducer_columns.clone(),
1265        },
1266        shard: options.shard.into(),
1267        execution: ManifestExecution {
1268            backend,
1269            concurrency: options.concurrency.to_string(),
1270            cpu_lanes: layout.cpu_lanes,
1271            threads_per_lane: layout.threads_per_lane,
1272            gpu_tracks: layout.gpu_tracks,
1273            projected_bytes: outline.projected_bytes,
1274            memory_budget: options.memory_budget,
1275            gpu_memory_budget: options.gpu_memory_budget,
1276        },
1277        runtime: inputs.runtime.clone(),
1278        timestamps: ManifestTimestamps::started_at(started_unix_ms),
1279        sessions,
1280        results: None,
1281        merged_shards,
1282        search: parts.search.clone(),
1283    })
1284}
1285
1286/// Marks `manifest` as ended after a batch that ended as `end` with `counts` rows written, and returns the end of the
1287/// sweep.
1288pub(crate) fn finish_manifest(manifest: &mut Manifest, end: BatchEnd, counts: ResultCounts) -> SweepEnd {
1289    let (status, end) = match end {
1290        BatchEnd::Complete => (ManifestStatus::Complete, SweepEnd::Complete),
1291        BatchEnd::Aborted => (ManifestStatus::Aborted, SweepEnd::Aborted),
1292        BatchEnd::DeviceLost => (ManifestStatus::Incomplete, SweepEnd::DeviceLost),
1293    };
1294    manifest.finish(status, counts, now_unix_ms());
1295    end
1296}
1297
1298/// Sink that writes each run to the output and reports the sweep's progress.
1299#[cfg(not(target_arch = "wasm32"))]
1300struct SweepSink<'s, W: Write> {
1301    writer: OutputWriter<W>,
1302    progress: &'s mut dyn Progress,
1303    meter: ProgressMeter,
1304}
1305
1306#[cfg(not(target_arch = "wasm32"))]
1307impl<W: Write> RunSink for SweepSink<'_, W> {
1308    fn commit(&mut self, outcome: RunOutcome) -> io::Result<()> {
1309        self.writer.write_run(&outcome)?;
1310        self.progress.report(&ProgressEvent::RunCommitted(&outcome));
1311        Ok(())
1312    }
1313
1314    fn finished(&mut self, outcome: &RunOutcome) {
1315        if let Some(update) = self.meter.record_finished_run(outcome.status) {
1316            self.progress.report(&ProgressEvent::Progressed(update));
1317        }
1318    }
1319}
1320
1321/// Returns `hash` as 16 hexadecimal digits.
1322pub(crate) fn hex(hash: u64) -> String {
1323    format!("{hash:016x}")
1324}
1325
1326#[cfg(test)]
1327mod tests {
1328    use std::num::NonZeroUsize;
1329    use std::path::PathBuf;
1330
1331    use henad_core::explore::design::DesignKind;
1332    use henad_core::explore::factor::{FactorSpec, LevelSpec};
1333    use henad_core::explore::spec::{BlockSpec, SweepSpec};
1334
1335    use std::path::Path;
1336
1337    use henad_compute::entry::{ModelEntry, register_grid_model};
1338    use henad_core::provenance::BuildInfo;
1339    use henad_models::game_of_life::GameOfLifeModel;
1340
1341    use super::{
1342        ExploreError, SpecSource, SweepEnd, SweepOptions, SweepOutline, SweepReport, SweepWarning, plan_spec,
1343        plan_then_acquire, run_spec,
1344    };
1345    use crate::device::DeviceError;
1346    use crate::exec::Concurrency;
1347    use crate::handle::SweepOutput;
1348    use crate::output::manifest::{BuildRole, Manifest, ManifestStatus, RecordedBuild};
1349    use crate::output::{MANIFEST_FILE, OutputError, RUNS_FILE, SERIES_FILE, SUMMARY_FILE};
1350    use crate::probe::ProbeReport;
1351    use crate::progress::{NoProgress, Progress, ProgressEvent};
1352    use crate::spec_file::SpecFile;
1353    use crate::tests::support::{ScratchDir, entry, provenance};
1354
1355    /// Progress that keeps one line per event, leaving out the updates that depend on the clock.
1356    #[derive(Debug, Default)]
1357    struct Recorded(Vec<String>);
1358
1359    impl Progress for Recorded {
1360        fn report(&mut self, event: &ProgressEvent<'_>) {
1361            let line = match event {
1362                ProgressEvent::Planned(outline) => format!("planned {}", outline.runs),
1363                ProgressEvent::Warned(warning) => format!("warned {warning}"),
1364                ProgressEvent::RunCommitted(outcome) => format!("run {}", outcome.run.run_id),
1365                ProgressEvent::Progressed(_) | ProgressEvent::SearchBatchTold(_) => return,
1366                ProgressEvent::Ended(report) => format!("ended {:?}", report.end),
1367            };
1368            self.0.push(line);
1369        }
1370    }
1371
1372    /// Returns 2 configs of 2 replicates of SIR on a 12 by 12 grid, 6 steps each.
1373    fn small_sweep() -> SweepSpec {
1374        let mut spec = SweepSpec::new("sir");
1375        spec.fixed = vec![
1376            ("grid_width".to_owned(), "12".to_owned()),
1377            ("grid_height".to_owned(), "12".to_owned()),
1378        ];
1379        spec.run.steps = 6;
1380        spec.run.replicates = 2;
1381        spec.seeds.root = 3;
1382        spec.blocks = vec![BlockSpec {
1383            design: DesignKind::Factorial,
1384            factors: vec![FactorSpec::param(
1385                "infection_rate",
1386                LevelSpec::Values(vec!["0.2".to_owned(), "0.4".to_owned()]),
1387            )],
1388            design_seed: None,
1389        }];
1390        spec
1391    }
1392
1393    fn options() -> SweepOptions {
1394        SweepOptions::new(provenance())
1395    }
1396
1397    /// Runs `spec` over `entry` into `output_dir` with `options`, reporting to `progress`.
1398    fn run_into(
1399        entry: &ModelEntry,
1400        spec: &SweepSpec,
1401        output_dir: &Path,
1402        options: &SweepOptions,
1403        progress: &mut dyn Progress,
1404    ) -> Result<SweepReport, ExploreError> {
1405        run_spec(
1406            entry,
1407            None,
1408            spec,
1409            SweepOutput::Directory(output_dir.to_owned()),
1410            options,
1411            progress,
1412        )
1413        .map(|record| record.report)
1414    }
1415
1416    #[test]
1417    fn a_dry_run_plans_and_probes_and_writes_nothing() {
1418        let sir = entry("sir", None);
1419        let scratch = ScratchDir::new("dry-run");
1420        let mut progress = Recorded::default();
1421        let report = plan_spec(
1422            &sir,
1423            None,
1424            &small_sweep(),
1425            Some(scratch.path()),
1426            &options(),
1427            &mut progress,
1428        )
1429        .expect("the dry run plans");
1430        assert_eq!(report.end, SweepEnd::Planned);
1431        assert_eq!(
1432            (report.outline.configs, report.outline.runs, report.outline.series_rows),
1433            (Some(2), 4, 4 * 7)
1434        );
1435        assert!(report.outline.dry_run && report.outline.projected_bytes > 0);
1436        assert_eq!((report.counts.rows, report.output_dir), (0, None));
1437        assert!(!scratch.path().exists(), "a dry run writes nothing");
1438        assert_eq!(progress.0, ["planned 4", "ended Planned"]);
1439    }
1440
1441    /// Progress that keeps the outline a sweep announces.
1442    #[derive(Default)]
1443    struct Outlined(Option<SweepOutline>);
1444
1445    impl Progress for Outlined {
1446        fn report(&mut self, event: &ProgressEvent<'_>) {
1447            if let ProgressEvent::Planned(outline) = event {
1448                self.0 = Some((*outline).clone());
1449            }
1450        }
1451    }
1452
1453    /// Checks that a dry run plans what the corresponding sweep or search announces, a dry run's flag aside.
1454    #[test]
1455    fn plan_spec_matches_a_dry_run() {
1456        let sir = entry("sir", None);
1457        let scratch = ScratchDir::new("plan-spec");
1458        let mut search = small_sweep();
1459        search.blocks.clear();
1460        search.search = Some(henad_core::explore::search::SearchSpec {
1461            algorithm: henad_core::explore::search::SearchAlgorithm::Random,
1462            max_evaluations: 4,
1463            batch_size: 2,
1464            objective: Some(henad_core::explore::search::Objective {
1465                column: "Infected:max".to_owned(),
1466                goal: henad_core::explore::search::Goal::Minimize,
1467                aggregate: henad_core::explore::search::Aggregate::Mean,
1468            }),
1469            space: vec![FactorSpec::param(
1470                "infection_rate",
1471                LevelSpec::Range {
1472                    min: 0.1,
1473                    max: 0.5,
1474                    step: None,
1475                },
1476            )],
1477        });
1478        for (name, spec) in [("sweep", small_sweep()), ("search", search)] {
1479            let planned = plan_spec(&sir, None, &spec, None, &options(), &mut NoProgress).expect("the spec plans");
1480            let mut announced = Outlined::default();
1481            let report =
1482                run_into(&sir, &spec, &scratch.path().join(name), &options(), &mut announced).expect("the spec runs");
1483            assert_eq!(report.end, SweepEnd::Complete, "{name}");
1484            let announced = announced.0.expect("a sweep announces its outline");
1485            assert!(planned.outline.dry_run && !announced.dry_run, "{name}");
1486            assert_eq!(
1487                SweepOutline {
1488                    dry_run: false,
1489                    ..planned.outline
1490                },
1491                announced,
1492                "{name}"
1493            );
1494        }
1495    }
1496
1497    #[test]
1498    fn the_projected_bytes_are_measured_on_a_pool_as_wide_as_a_lane() {
1499        let ants = entry("ants", None);
1500        let mut spec = SweepSpec::new("ants");
1501        spec.fixed = vec![
1502            ("num_agents".to_owned(), "300".to_owned()),
1503            ("world_width".to_owned(), "64".to_owned()),
1504            ("world_height".to_owned(), "64".to_owned()),
1505        ];
1506        spec.run.replicates = 2;
1507        let two_lanes = SweepOptions {
1508            concurrency: Concurrency::Fixed(NonZeroUsize::new(2).expect("2 is above 0")),
1509            ..options()
1510        };
1511        let report = plan_spec(&ants, None, &spec, None, &two_lanes, &mut NoProgress).expect("the dry run plans");
1512        let layout = report.outline.layout;
1513        assert_eq!(layout.cpu_lanes, 2);
1514
1515        let plan = spec.plan(&ants.schema()).expect("a valid spec");
1516        let probe = ProbeReport::for_plan(&ants, None, &plan).expect("ants builds");
1517        let lane = probe
1518            .rebuilt_on(&ants, layout.threads_per_lane)
1519            .expect("ants builds on a lane's pool");
1520        assert_eq!(report.outline.projected_bytes, 2 * lane.heap_bytes);
1521    }
1522
1523    #[test]
1524    fn the_projection_takes_the_larger_of_the_first_and_last_configs() {
1525        let life = entry("game_of_life", None);
1526        for widths in [["16", "512"], ["512", "16"]] {
1527            let mut spec = SweepSpec::new("game_of_life");
1528            spec.fixed = vec![("grid_height".to_owned(), "256".to_owned())];
1529            spec.blocks = vec![BlockSpec {
1530                design: DesignKind::Factorial,
1531                factors: vec![FactorSpec::param(
1532                    "grid_width",
1533                    LevelSpec::Values(widths.map(str::to_owned).to_vec()),
1534                )],
1535                design_seed: None,
1536            }];
1537            let one_lane = SweepOptions {
1538                concurrency: Concurrency::Fixed(NonZeroUsize::MIN),
1539                ..options()
1540            };
1541            let report = plan_spec(&life, None, &spec, None, &one_lane, &mut NoProgress).expect("the dry run plans");
1542            let plan = spec.plan(&life.schema()).expect("a valid spec");
1543            let wide = usize::from(widths[0] == "16");
1544            let run = plan.run(wide as u64).expect("each config has a run");
1545            let params = &plan.config(run.config_id).expect("the run's config").params;
1546            let largest = ProbeReport::build(&life, None, params, Some(run.seed)).expect("the model builds");
1547            assert_eq!(report.outline.projected_bytes, largest.heap_bytes, "widths {widths:?}");
1548        }
1549    }
1550
1551    #[test]
1552    fn a_sweep_records_its_plan_and_ends_complete() {
1553        let sir = entry("sir", None);
1554        let spec = small_sweep();
1555        let scratch = ScratchDir::new("manifest");
1556        let source = SpecSource {
1557            path: Some(PathBuf::from("specs/small.toml")),
1558            toml: Some("model = \"sir\"\n".to_owned()),
1559            tables: Vec::new(),
1560        };
1561        let mut progress = Recorded::default();
1562        let sweep_options = SweepOptions {
1563            spec_source: source.clone(),
1564            ..options()
1565        };
1566        let report = run_into(&sir, &spec, scratch.path(), &sweep_options, &mut progress).expect("the sweep runs");
1567        assert_eq!(report.end, SweepEnd::Complete);
1568        assert_eq!(report.output_dir.as_deref(), Some(scratch.path()));
1569        assert_eq!(
1570            progress.0,
1571            ["planned 4", "run 0", "run 1", "run 2", "run 3", "ended Complete"]
1572        );
1573        for file in [RUNS_FILE, SERIES_FILE, SUMMARY_FILE] {
1574            assert!(scratch.path().join(file).is_file(), "{file}");
1575        }
1576
1577        let text = std::fs::read_to_string(scratch.path().join(MANIFEST_FILE)).expect("the manifest is written");
1578        let manifest: Manifest = serde_json::from_str(&text).expect("the manifest reads back");
1579        assert_eq!(manifest.status, ManifestStatus::Complete);
1580        assert_eq!(manifest.results, Some(report.counts));
1581        assert_eq!((manifest.plan.configs, manifest.plan.runs), (Some(2), 4));
1582        assert_eq!(manifest.plan.blocks[0].design, "factorial");
1583        assert_eq!(manifest.plan.plan_hash.len(), 16);
1584        assert_eq!(manifest.seeds.root, 3);
1585        assert_eq!(manifest.columns.stats, ["Susceptible", "Infected", "Recovered"]);
1586        assert_eq!(manifest.columns.reducers.len(), 12);
1587        assert_eq!(manifest.spec_source.toml, source.toml);
1588        assert_eq!(manifest.spec_source.path.as_deref(), Some("specs/small.toml"));
1589        assert_eq!(manifest.sessions[0].ran, 4);
1590        assert!(manifest.timestamps.finished_unix_ms >= Some(manifest.timestamps.started_unix_ms));
1591        assert!(manifest.model.schema["stat_columns"].is_array());
1592        let written: SpecFile = serde_json::from_value(manifest.spec).expect("the spec reads back");
1593        let mut sorted = spec.clone();
1594        sorted.fixed.sort();
1595        assert_eq!(
1596            written.into_spec().expect("a valid spec"),
1597            sorted,
1598            "fixed values come back by id"
1599        );
1600
1601        let again = run_into(&sir, &spec, scratch.path(), &sweep_options, &mut NoProgress)
1602            .expect_err("the directory holds results");
1603        assert!(
1604            matches!(again, ExploreError::Output(OutputError::HoldsResults { .. })),
1605            "{again:?}"
1606        );
1607    }
1608
1609    #[test]
1610    fn an_aborted_sweep_says_so_and_keeps_its_headers() {
1611        let sir = entry("sir", None);
1612        let scratch = ScratchDir::new("aborted");
1613        let sweep_options = options();
1614        sweep_options.control.abort();
1615        let report = run_into(&sir, &small_sweep(), scratch.path(), &sweep_options, &mut NoProgress)
1616            .expect("an aborted sweep is not an error");
1617        assert_eq!((report.end, report.counts.rows), (SweepEnd::Aborted, 0));
1618        let text = std::fs::read_to_string(scratch.path().join(MANIFEST_FILE)).expect("the manifest is written");
1619        let manifest: Manifest = serde_json::from_str(&text).expect("the manifest reads back");
1620        assert_eq!(manifest.status, ManifestStatus::Aborted);
1621        for file in [RUNS_FILE, SUMMARY_FILE] {
1622            let table = std::fs::read_to_string(scratch.path().join(file)).expect("the table is written");
1623            assert_eq!(table.lines().count(), 1, "{file} holds its header alone");
1624        }
1625    }
1626
1627    /// Returns the model build warning for `recorded` against `current`.
1628    fn model_warning(recorded: &RecordedBuild, current: &RecordedBuild, between_shards: bool) -> String {
1629        SweepWarning::BuildChanged {
1630            role: BuildRole::Model,
1631            recorded: Box::new(recorded.clone()),
1632            current: Box::new(current.clone()),
1633            between_shards,
1634        }
1635        .to_string()
1636    }
1637
1638    #[test]
1639    fn a_warning_between_unidentified_builds_says_they_cannot_be_told_apart() {
1640        let bare = RecordedBuild::from(register_grid_model::<GameOfLifeModel>().source());
1641        let text = model_warning(&bare, &bare, false);
1642        assert!(
1643            text.starts_with("cannot tell whether the directory's runs so far came from this model build"),
1644            "{text}"
1645        );
1646        assert!(
1647            text.ends_with(
1648                "Insert it into a `ModelSet` built with `henad::build_info!()` in a crate whose build script calls \
1649                 `henad_build::stamp_commit()`"
1650            ),
1651            "an entry from no set records no build: {text}"
1652        );
1653        let text = model_warning(&bare, &bare, true);
1654        assert!(
1655            text.starts_with("cannot tell whether the shards ran the same model build"),
1656            "{text}"
1657        );
1658
1659        let unstamped = RecordedBuild::from(&BuildInfo::__from_env(
1660            "my-model",
1661            "0.1.0",
1662            Some(""),
1663            None,
1664            Some(""),
1665            None,
1666            false,
1667        ));
1668        let stamped = RecordedBuild::from(entry("game_of_life", None).source());
1669        let text = model_warning(&stamped, &unstamped, false);
1670        assert!(
1671            text.starts_with("the directory's runs so far came from the model build"),
1672            "{text}"
1673        );
1674        assert!(
1675            text.ends_with("Call `henad_build::stamp_commit()` in the build script of the crate that registers it"),
1676            "a build script can stamp a set's build: {text}"
1677        );
1678    }
1679
1680    /// Checks that a GPU spec that the model rejects fails as a plan, before any device is acquired.
1681    #[test]
1682    fn a_gpu_spec_is_planned_before_its_device_is_acquired() {
1683        let gpu_sir = henad_models::example_models()
1684            .get("gpu_sir")
1685            .cloned()
1686            .expect("the model is registered");
1687        let mut spec = SweepSpec::new("gpu_sir");
1688        spec.fixed = vec![("infection_rate".to_owned(), "2".to_owned())];
1689        let mut acquisitions = 0;
1690        let mut refuse = |_| {
1691            acquisitions += 1;
1692            Err(DeviceError::BelowBaseline {
1693                adapter: "test".to_owned(),
1694                limit: "max_buffer_size",
1695            })
1696        };
1697        let refused = plan_then_acquire(&gpu_sir, None, &spec, &mut refuse).err();
1698        assert!(matches!(refused, Some(ExploreError::Plan(_))), "{refused:?}");
1699        let planned = plan_spec(&gpu_sir, None, &spec, None, &options(), &mut NoProgress);
1700        assert!(matches!(planned, Err(ExploreError::Plan(_))), "{planned:?}");
1701        let ran = run_spec(&gpu_sir, None, &spec, SweepOutput::Memory, &options(), &mut NoProgress);
1702        assert!(matches!(ran, Err(ExploreError::Plan(_))), "{ran:?}");
1703
1704        spec.fixed.clear();
1705        let unacquired = plan_then_acquire(&gpu_sir, None, &spec, &mut refuse).err();
1706        assert!(matches!(unacquired, Some(ExploreError::Device(_))), "{unacquired:?}");
1707        assert_eq!(acquisitions, 1, "only the spec the model accepts asks for a device");
1708    }
1709
1710    #[test]
1711    fn missing_runs_are_counted_in_words() {
1712        let missing = |count| SweepWarning::MissingRuns { count, first: vec![3] }.to_string();
1713        assert_eq!(
1714            missing(1),
1715            "1 run of the plan is missing, starting with 3. Resume the merged directory to run it"
1716        );
1717        assert_eq!(
1718            missing(2),
1719            "2 runs of the plan are missing, starting with 3. Resume the merged directory to run them"
1720        );
1721    }
1722}