1#[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#[derive(Debug, Clone, Default, PartialEq, Eq)]
56pub struct SpecSource {
57 pub path: Option<PathBuf>,
59 pub toml: Option<String>,
61 pub tables: Vec<DesignTableFile>,
63}
64
65impl SpecSource {
66 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 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#[derive(Debug, Clone, PartialEq, Eq)]
115pub struct Provenance {
116 engine: RecordedBuild,
117 host: RecordedBuild,
118 arguments: Vec<String>,
119}
120
121impl Provenance {
122 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 pub fn engine(&self) -> &RecordedBuild {
136 &self.engine
137 }
138
139 pub fn host(&self) -> &RecordedBuild {
141 &self.host
142 }
143
144 pub fn arguments(&self) -> &[String] {
146 &self.arguments
147 }
148
149 #[cfg(test)]
151 pub(crate) fn with_engine(self, engine: RecordedBuild) -> Self {
152 Self { engine, ..self }
153 }
154}
155
156#[derive(Debug, Clone)]
161#[non_exhaustive]
162pub struct SweepOptions {
163 pub concurrency: Concurrency,
165 pub memory_budget: Option<u64>,
168 pub gpu_memory_budget: Option<u64>,
171 pub control: SweepControl,
173 pub shard: Shard,
175 pub resume: bool,
179 pub retry_failed: bool,
183 pub(crate) active_runs: Option<ActiveRuns>,
185 pub spec_source: SpecSource,
187 pub provenance: Provenance,
189}
190
191impl SweepOptions {
192 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 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#[derive(Debug, Clone, PartialEq)]
221pub struct SweepOutline {
222 pub model: String,
224 pub backend: Backend,
226 pub configs: Option<u64>,
228 pub replicates: u64,
230 pub runs: u64,
232 pub blocks: Vec<PlannedBlock>,
234 pub shard: Shard,
236 pub skipped: u64,
238 pub pending: u64,
240 pub layout: ExecutionLayout,
242 pub projected_bytes: u64,
244 pub series_rows: u64,
246 pub stat_columns: Vec<String>,
248 pub reducer_columns: Vec<String>,
250 pub dry_run: bool,
252 pub search: Option<SearchOutline>,
254}
255
256#[derive(Debug, Clone, PartialEq, Eq)]
258pub enum SweepWarning {
259 Plan(PlanWarning),
261 BuildChanged {
269 role: BuildRole,
271 recorded: Box<RecordedBuild>,
273 current: Box<RecordedBuild>,
275 between_shards: bool,
277 },
278 MissingRuns {
280 count: u64,
282 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
355pub enum SweepEnd {
356 Planned,
358 Complete,
360 Aborted,
362 DeviceLost,
364}
365
366#[derive(Debug, Clone, PartialEq)]
368pub struct SweepReport {
369 pub outline: SweepOutline,
371 pub end: SweepEnd,
373 pub counts: ResultCounts,
375 pub elapsed: Duration,
377 pub output_dir: Option<PathBuf>,
379}
380
381#[derive(Debug, Clone, PartialEq)]
383pub struct SweepRecord {
384 pub report: SweepReport,
386 pub manifest: Manifest,
388 pub files: Option<SweepFiles>,
390 pub search: Option<SearchReport>,
392}
393
394#[derive(Debug)]
398#[non_exhaustive]
399pub enum ExploreError {
400 Plan(PlanError),
402 Capacity(CapacityError),
404 Probe(ProbeError),
406 Measure(MeasureError),
408 Resume(ResumeError),
410 Output(OutputError),
412 Execution(ExecutionError),
414 Search(SearchPlanError),
416 #[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#[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#[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 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#[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#[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#[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#[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#[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#[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
765pub(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 pub(crate) folder: Option<&'a Path>,
776 pub(crate) dry_run: bool,
778}
779
780pub(crate) struct SweepPreparation {
782 plan: Arc<Plan>,
783 probe: ProbeReport,
784 measure: Arc<MeasurePlan>,
785 pending: Vec<PlannedRun>,
787 resumed: Option<ResumeScan>,
789 #[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 started: Instant,
799 started_unix_ms: u64,
800}
801
802impl SweepPreparation {
803 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 (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 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 #[cfg(any(target_arch = "wasm32", test))]
921 pub(crate) fn pending(&self) -> &[PlannedRun] {
922 &self.pending
923 }
924
925 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 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 #[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 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 #[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 #[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 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 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
1106pub(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
1142fn 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
1158pub(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
1179pub(crate) struct ManifestParts<'a> {
1181 pub(crate) mode: ManifestMode,
1182 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 pub(crate) recorded: Option<&'a Manifest>,
1190 pub(crate) search: Option<ManifestSearch>,
1191}
1192
1193pub(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 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
1286pub(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#[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
1321pub(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 #[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 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 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 #[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 #[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 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 #[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}