Skip to main content

henad_explore/output/
mod.rs

1//! Output directories of a sweep, and the writer that streams each finished run into a directory.
2//!
3//! A directory holds four files. `runs.csv` has one row per run, `series.csv` the sampled stat rows of every run,
4//! `summary.csv` statistics over the replicates of each config, and `manifest.json` the settings, provenance and
5//! progress of the sweep. Both CSV tables list runs in order of their ids once a sweep ends. A search adds the
6//! tables [`search_tables`] writes.
7
8pub mod details;
9pub mod manifest;
10pub mod memory;
11pub mod read;
12pub mod resume;
13pub mod runs_csv;
14pub mod search_tables;
15pub mod series_csv;
16pub mod summary_csv;
17
18use std::fmt;
19use std::fs::{File, OpenOptions, TryLockError};
20use std::io::{self, BufReader, BufWriter, Write};
21use std::path::{Path, PathBuf};
22use std::sync::Arc;
23
24use henad_core::explore::measure::MeasurePlan;
25use henad_core::explore::outcome::RunOutcome;
26use henad_core::explore::plan::{Config, Plan};
27use henad_core::params::ParamDescriptor;
28
29use crate::exec::RunSink;
30use crate::output::manifest::{Manifest, ResultCounts};
31use crate::output::read::{ReadError, RunsCsv, SeriesScan, SeriesSegment, merge_series};
32use crate::output::runs_csv::RunsWriter;
33use crate::output::series_csv::SeriesWriter;
34use crate::output::summary_csv::{SummaryError, write_summary};
35
36/// Name of the table with one row per run.
37pub const RUNS_FILE: &str = "runs.csv";
38/// Name of the table of sampled stat rows.
39pub const SERIES_FILE: &str = "series.csv";
40/// Name of the table of statistics over the replicates of each config.
41pub const SUMMARY_FILE: &str = "summary.csv";
42/// Name of the file holding the settings, provenance and progress of a sweep.
43pub const MANIFEST_FILE: &str = "manifest.json";
44/// Name of a search's table with one row per evaluation.
45pub const EVALUATIONS_FILE: &str = "evaluations.csv";
46/// Name of a search's table with one row per batch.
47pub const BATCHES_FILE: &str = "batches.csv";
48/// Name of the table ranking every candidate of a search that scores an objective.
49pub const BEST_FILE: &str = "best.csv";
50/// Name of the table of cells that a Pattern Space Exploration fills.
51pub const ARCHIVE_FILE: &str = "archive.csv";
52/// Name of a genetic algorithm's table with one row per finished generation.
53pub const GENERATIONS_FILE: &str = "generations.csv";
54
55/// Files a sweep or search writes, any one of which marks a directory as holding results.
56const RESULT_FILES: [&str; 9] = [
57    RUNS_FILE,
58    SERIES_FILE,
59    SUMMARY_FILE,
60    MANIFEST_FILE,
61    EVALUATIONS_FILE,
62    BATCHES_FILE,
63    BEST_FILE,
64    ARCHIVE_FILE,
65    GENERATIONS_FILE,
66];
67
68/// File a writer locks while it writes to an output directory.
69pub const LOCK_FILE: &str = ".lock";
70
71/// Suffix of the file that holds a table's full replacement until it is renamed into place.
72const STAGED_SUFFIX: &str = ".staged";
73
74/// File whose presence says both staged tables are complete.
75const STAGED_MARKER: &str = "tables.staged";
76
77/// Returns the path of the staged table that replaces `file` in `dir`.
78fn staged_path(dir: &Path, file: &str) -> PathBuf {
79    dir.join(format!("{file}{STAGED_SUFFIX}"))
80}
81
82/// Returns the paths of `runs.csv` and `series.csv` in the directory at `path`, as they stand.
83///
84/// A staged table is used instead of its original while a replacement waits to be finished.
85pub(crate) fn table_paths(path: &Path) -> (PathBuf, PathBuf) {
86    let complete = path.join(STAGED_MARKER).exists();
87    let current = |file: &str| {
88        let staged = staged_path(path, file);
89        if complete && staged.exists() {
90            staged
91        } else {
92            path.join(file)
93        }
94    };
95    (current(RUNS_FILE), current(SERIES_FILE))
96}
97
98/// A directory that holds, or is about to hold, the results of one sweep, locked against other writers.
99///
100/// The lock is the operating system's advisory lock on the file [`LOCK_FILE`], released when the value is dropped. The
101/// file stays in the directory, and the next writer locks it again. The lock is released when the process that holds it
102/// ends, even when the process is killed.
103///
104/// A writer creates afresh each file it replaces. It appends to `runs.csv` and `series.csv`, and cuts them, only after
105/// [`Self::open`] has checked that neither table is a symbolic link.
106#[derive(Debug)]
107pub struct OutputDir {
108    path: PathBuf,
109    #[expect(dead_code, reason = "held for its drop, which releases the lock")]
110    lock: DirLock,
111}
112
113impl OutputDir {
114    /// Checks that the directory at `path` holds no results, without creating it.
115    ///
116    /// A symbolic link under the name of a file counts as that file, even a link that points nowhere.
117    ///
118    /// # Errors
119    ///
120    /// Returns [`OutputError::HoldsResults`] when the directory holds a file that a sweep or a search writes,
121    /// including a staged table or the marker of staged tables.
122    pub fn check_free(path: &Path) -> Result<(), OutputError> {
123        let holds_results = |file: String| {
124            Err(OutputError::HoldsResults {
125                dir: path.to_owned(),
126                file,
127            })
128        };
129        if let Some(&file) = RESULT_FILES
130            .iter()
131            .find(|&&file| std::fs::symlink_metadata(path.join(file)).is_ok())
132        {
133            return holds_results(file.to_owned());
134        }
135        // `table_paths` can read a staged table instead of its original, so a staged table counts as results.
136        let staged = std::fs::read_dir(path)
137            .into_iter()
138            .flatten()
139            .flatten()
140            .map(|entry| entry.file_name().to_string_lossy().into_owned())
141            .find(|name| name.ends_with(STAGED_SUFFIX));
142        match staged {
143            Some(file) => holds_results(file),
144            None => Ok(()),
145        }
146    }
147
148    /// Creates the directory at `path` and its parents, or uses an existing directory that holds no results, and
149    /// locks it.
150    ///
151    /// # Errors
152    ///
153    /// Returns [`OutputError::HoldsResults`] when the directory holds a file that a sweep or a search writes,
154    /// [`OutputError::Locked`] when another writer holds its lock, and [`OutputError::Write`] when it cannot be
155    /// created.
156    pub fn create(path: &Path) -> Result<Self, OutputError> {
157        Self::check_free(path)?;
158        std::fs::create_dir_all(path).map_err(write_error_at(path))?;
159        let dir = Self {
160            path: path.to_owned(),
161            lock: DirLock::acquire(path)?,
162        };
163        // Checked again under the lock. Another writer could have finished between the first check and the lock.
164        Self::check_free(path)?;
165        Ok(dir)
166    }
167
168    /// Returns whether the directory at `path` holds a file that a sweep or a search writes.
169    pub fn holds_results(path: &Path) -> bool {
170        Self::check_free(path).is_err()
171    }
172
173    /// Locks the existing directory at `path`, finishing a table replacement that a process left behind.
174    ///
175    /// # Errors
176    ///
177    /// Returns [`OutputError::Locked`] when another writer holds the directory's lock, [`OutputError::Link`] when
178    /// `runs.csv` or `series.csv` is a symbolic link, and [`OutputError::Write`] when the replacement cannot be
179    /// finished or discarded.
180    pub fn open(path: &Path) -> Result<Self, OutputError> {
181        let dir = Self {
182            path: path.to_owned(),
183            lock: DirLock::acquire(path)?,
184        };
185        dir.finish_staged()?;
186        // A resume appends to both tables and cuts them, and both operations would follow a link.
187        for file in [RUNS_FILE, SERIES_FILE] {
188            let table = path.join(file);
189            if std::fs::symlink_metadata(&table).is_ok_and(|metadata| metadata.file_type().is_symlink()) {
190                return Err(OutputError::Link { path: table });
191            }
192        }
193        Ok(dir)
194    }
195
196    /// Checks that no writer holds the lock of the directory at `path`, without taking it.
197    ///
198    /// # Errors
199    ///
200    /// Returns [`OutputError::Locked`] when another writer holds the lock.
201    pub fn check_unlocked(path: &Path) -> Result<(), OutputError> {
202        // A directory or a lock file that is missing is unlocked. The check never creates the directory or the file.
203        let Ok(file) = File::open(path.join(LOCK_FILE)) else {
204            return Ok(());
205        };
206        match file.try_lock_shared() {
207            Err(TryLockError::WouldBlock) => Err(OutputError::Locked { dir: path.to_owned() }),
208            Ok(()) | Err(TryLockError::Error(_)) => Ok(()),
209        }
210    }
211
212    /// Path of the directory.
213    pub fn path(&self) -> &Path {
214        &self.path
215    }
216
217    /// Creates `runs.csv` and `series.csv`, and returns a writer for the runs of `plan`.
218    ///
219    /// `params` are the model's parameters, and `measure` fixes the stat and reducer columns.
220    ///
221    /// # Errors
222    ///
223    /// Returns [`OutputError::Write`] when a file exists already or cannot be created, or its header cannot be written.
224    pub fn open_writer(
225        &self,
226        plan: &Arc<Plan>,
227        params: &[ParamDescriptor],
228        measure: &MeasurePlan,
229    ) -> Result<OutputWriter<BufWriter<File>>, OutputError> {
230        let runs_path = self.path.join(RUNS_FILE);
231        let runs = File::create_new(&runs_path).map_err(write_error_at(&runs_path))?;
232        let runs = RunsWriter::new(BufWriter::new(runs), params, plan.actions(), measure.reducers().names())
233            .map_err(write_error_at(&runs_path))?;
234        let series_path = self.path.join(SERIES_FILE);
235        let series = File::create_new(&series_path).map_err(write_error_at(&series_path))?;
236        let series =
237            SeriesWriter::new(BufWriter::new(series), measure.columns()).map_err(write_error_at(&series_path))?;
238        Ok(OutputWriter::new(Arc::clone(plan), runs, series))
239    }
240
241    /// Opens `runs.csv` and `series.csv`, whose headers are already written, and returns a writer that adds the runs
242    /// of `plan` to them.
243    ///
244    /// # Errors
245    ///
246    /// Returns [`OutputError::Write`] when a file cannot be opened.
247    pub fn append_writer(
248        &self,
249        plan: &Arc<Plan>,
250        params: &[ParamDescriptor],
251    ) -> Result<OutputWriter<BufWriter<File>>, OutputError> {
252        let append = |file: &str| {
253            let path = self.path.join(file);
254            OpenOptions::new()
255                .append(true)
256                .open(&path)
257                .map(BufWriter::new)
258                .map_err(write_error_at(&path))
259        };
260        let runs = RunsWriter::appending(append(RUNS_FILE)?, params);
261        let series = SeriesWriter::appending(append(SERIES_FILE)?);
262        Ok(OutputWriter::new(Arc::clone(plan), runs, series))
263    }
264
265    /// Replaces `runs.csv` and `series.csv` with the text that `write_runs` and `write_series` write.
266    ///
267    /// Each table is written in full beside the original with a `.staged` suffix. A marker file then records that both
268    /// tables are complete, both tables are renamed into place, and the marker is removed. A process that ends between
269    /// the marker and its removal leaves the replacement for [`Self::open`] to finish.
270    ///
271    /// # Errors
272    ///
273    /// Returns [`OutputError::Write`] when a write, a rename or a removal fails. The originals stay in place until
274    /// both new tables are complete.
275    pub fn replace_tables(
276        &self,
277        write_runs: impl FnOnce(&mut dyn Write) -> io::Result<()>,
278        write_series: impl FnOnce(&mut dyn Write) -> io::Result<()>,
279    ) -> Result<(), OutputError> {
280        self.write_staged(RUNS_FILE, write_runs)?;
281        self.write_staged(SERIES_FILE, write_series)?;
282        let marker = self.path.join(STAGED_MARKER);
283        create_fresh(&marker)
284            .and_then(|file| file.sync_all())
285            .map_err(write_error_at(&marker))?;
286        self.finish_staged()
287    }
288
289    /// Calls `write` to write the staged table that replaces `file`, and flushes the table to the disk.
290    fn write_staged(
291        &self,
292        file: &str,
293        write: impl FnOnce(&mut dyn Write) -> io::Result<()>,
294    ) -> Result<(), OutputError> {
295        let path = staged_path(&self.path, file);
296        let staged = create_fresh(&path).map_err(write_error_at(&path))?;
297        let mut writer = BufWriter::new(staged);
298        write(&mut writer).map_err(write_error_at(&path))?;
299        let staged = writer
300            .into_inner()
301            .map_err(|error| write_error_at(&path)(error.into_error()))?;
302        staged.sync_all().map_err(write_error_at(&path))
303    }
304
305    /// Moves both staged tables into place when the marker says they are complete, and removes staged tables left
306    /// incomplete otherwise.
307    fn finish_staged(&self) -> Result<(), OutputError> {
308        let marker = self.path.join(STAGED_MARKER);
309        let complete = marker.exists();
310        for file in [RUNS_FILE, SERIES_FILE] {
311            let staged = staged_path(&self.path, file);
312            // A link that points nowhere is a staged table to remove as well.
313            if std::fs::symlink_metadata(&staged).is_err() {
314                continue;
315            }
316            if complete {
317                let installed = self.path.join(file);
318                std::fs::rename(&staged, &installed).map_err(write_error_at(&installed))?;
319            } else {
320                std::fs::remove_file(&staged).map_err(write_error_at(&staged))?;
321            }
322        }
323        if complete {
324            std::fs::remove_file(&marker).map_err(write_error_at(&marker))?;
325        }
326        Ok(())
327    }
328
329    /// Rewrites `runs.csv` and `series.csv` with their runs in order of their ids, leaving tables already in order
330    /// alone.
331    ///
332    /// # Errors
333    ///
334    /// Returns [`OutputError::Table`] when a table cannot be read back, and [`OutputError::Write`] when the
335    /// replacement fails.
336    pub fn order_tables(&self) -> Result<(), OutputError> {
337        let runs_path = self.path.join(RUNS_FILE);
338        let series_path = self.path.join(SERIES_FILE);
339        let runs = RunsCsv::read(&runs_path).map_err(OutputError::Table)?;
340        let series = SeriesScan::read(&series_path, |_| true).map_err(OutputError::Table)?;
341        let in_order = runs.records.is_sorted_by_key(|record| record.run_id) && series.segments.len() <= 1;
342        if in_order {
343            return Ok(());
344        }
345        let (Some(runs_header), Some(series_header)) = (runs.header.as_ref(), series.header.as_ref()) else {
346            return Ok(());
347        };
348        let mut records: Vec<_> = runs.records.iter().collect();
349        records.sort_by_key(|record| record.run_id);
350        let segments: Vec<SeriesSegment> = series
351            .segments
352            .iter()
353            .map(|range| SeriesSegment {
354                path: series_path.clone(),
355                range: range.clone(),
356                input_index: 0,
357            })
358            .collect();
359        let runs_header_line = runs_csv::header_line(runs_header);
360        self.replace_tables(
361            |dest| {
362                dest.write_all(runs_header_line.as_bytes())?;
363                records
364                    .iter()
365                    .try_for_each(|record| dest.write_all(record.text.as_bytes()))
366            },
367            |dest| merge_series(dest, series_header, &segments, |_, run_id| Some(run_id)),
368        )
369    }
370
371    /// Creates the table `file` afresh, removing any existing file of that name, and returns a buffered writer over it.
372    ///
373    /// # Errors
374    ///
375    /// Returns [`OutputError::Write`] when the file cannot be removed or created.
376    pub fn create_table(&self, file: &str) -> Result<BufWriter<File>, OutputError> {
377        let path = self.path.join(file);
378        create_fresh(&path).map(BufWriter::new).map_err(write_error_at(&path))
379    }
380
381    /// Writes `manifest` to `manifest.json`. An earlier manifest is replaced by a single rename, so a reader sees
382    /// either the old or the new manifest whole.
383    ///
384    /// # Errors
385    ///
386    /// Returns [`OutputError::Manifest`] when the manifest cannot be serialized, and [`OutputError::Write`] when a
387    /// write or the rename fails.
388    pub fn write_manifest(&self, manifest: &Manifest) -> Result<(), OutputError> {
389        let text = manifest_text(manifest)?;
390        let partial = self.path.join(format!("{MANIFEST_FILE}.partial"));
391        // The file closes at the end of the block, before the rename.
392        {
393            let mut file = create_fresh(&partial).map_err(write_error_at(&partial))?;
394            file.write_all(text.as_bytes()).map_err(write_error_at(&partial))?;
395            file.sync_all().map_err(write_error_at(&partial))?;
396        }
397        let manifest_path = self.path.join(MANIFEST_FILE);
398        std::fs::rename(&partial, &manifest_path).map_err(write_error_at(&manifest_path))
399    }
400
401    /// Rebuilds `summary.csv` from `runs.csv`.
402    ///
403    /// # Errors
404    ///
405    /// Returns [`OutputError::Read`] when `runs.csv` cannot be read, [`OutputError::Write`] when `summary.csv`
406    /// cannot be written, and [`OutputError::Summary`] when `runs.csv` cannot be summarized.
407    pub fn write_summary(&self) -> Result<(), OutputError> {
408        let runs_path = self.path.join(RUNS_FILE);
409        let read_error = |source| OutputError::Read {
410            path: runs_path.clone(),
411            source,
412        };
413        let runs = File::open(&runs_path).map_err(read_error)?;
414        let summary_path = self.path.join(SUMMARY_FILE);
415        let summary = create_fresh(&summary_path).map_err(write_error_at(&summary_path))?;
416        match write_summary(BufReader::new(runs), BufWriter::new(summary)) {
417            Ok(_) => Ok(()),
418            Err(SummaryError::Read(source)) => Err(read_error(source)),
419            Err(SummaryError::Io(source)) => Err(write_error_at(&summary_path)(source)),
420            Err(error) => Err(OutputError::Summary(error)),
421        }
422    }
423}
424
425/// Returns `manifest` as the text of `manifest.json`, pretty-printed JSON ending in a line feed.
426///
427/// # Errors
428///
429/// Returns [`OutputError::Manifest`] when the manifest cannot be serialized.
430pub(crate) fn manifest_text(manifest: &Manifest) -> Result<String, OutputError> {
431    let mut text = serde_json::to_string_pretty(manifest).map_err(OutputError::Manifest)?;
432    text.push('\n');
433    Ok(text)
434}
435
436/// Returns a function that turns an I/O error on `path` into an [`OutputError::Write`].
437fn write_error_at(path: &Path) -> impl FnOnce(io::Error) -> OutputError + use<> {
438    let path = path.to_owned();
439    move |source| OutputError::Write { path, source }
440}
441
442/// Creates the file at `path` afresh. A file or a symbolic link already there is removed first, never followed.
443///
444/// # Errors
445///
446/// Returns the error of the removal or the creation.
447fn create_fresh(path: &Path) -> io::Result<File> {
448    match std::fs::remove_file(path) {
449        Err(error) if error.kind() != io::ErrorKind::NotFound => return Err(error),
450        _ => {}
451    }
452    File::create_new(path)
453}
454
455/// Advisory lock a writer holds on [`LOCK_FILE`] of an output directory, released when dropped.
456///
457/// Note that the file stays in the directory. If it were removed, a writer that had it open could lock the removed file
458/// while another writer locks a new lock file.
459#[derive(Debug)]
460struct DirLock {
461    /// File the lock is held on.
462    file: File,
463}
464
465impl DirLock {
466    /// Locks the directory at `dir`, creating its lock file when missing.
467    ///
468    /// Note that a filesystem without file locks, such as some network filesystems, leaves the directory unguarded.
469    ///
470    /// # Errors
471    ///
472    /// Returns [`OutputError::Locked`] when another writer holds the lock, and [`OutputError::Write`] when the lock
473    /// file cannot be opened.
474    fn acquire(dir: &Path) -> Result<Self, OutputError> {
475        let path = dir.join(LOCK_FILE);
476        if std::fs::symlink_metadata(&path).is_ok_and(|metadata| metadata.file_type().is_symlink()) {
477            std::fs::remove_file(&path).map_err(write_error_at(&path))?;
478        }
479        let file = OpenOptions::new()
480            .read(true)
481            .write(true)
482            .create(true)
483            .truncate(false)
484            .open(&path)
485            .map_err(write_error_at(&path))?;
486        match file.try_lock() {
487            Err(TryLockError::WouldBlock) => Err(OutputError::Locked { dir: dir.to_owned() }),
488            Ok(()) | Err(TryLockError::Error(_)) => Ok(Self { file }),
489        }
490    }
491}
492
493impl Drop for DirLock {
494    /// Releases the lock before the file closes.
495    fn drop(&mut self) {
496        drop(self.file.unlock());
497    }
498}
499
500/// Sink that streams each committed run to `runs.csv` and `series.csv`, and counts the rows by status.
501///
502/// A run's series is written and flushed before its row in `runs.csv`, so every run in `runs.csv` has its whole
503/// series written.
504#[derive(Debug)]
505pub struct OutputWriter<W: Write> {
506    plan: Arc<Plan>,
507    runs: RunsWriter<W>,
508    series: SeriesWriter<W>,
509    counts: ResultCounts,
510}
511
512impl<W: Write> OutputWriter<W> {
513    /// Returns a writer for the runs of `plan` over writers whose headers are written.
514    pub fn new(plan: Arc<Plan>, runs: RunsWriter<W>, series: SeriesWriter<W>) -> Self {
515        Self {
516            plan,
517            runs,
518            series,
519            counts: ResultCounts::default(),
520        }
521    }
522
523    /// Writes the series of `outcome`, then its row, and flushes after each write.
524    ///
525    /// # Errors
526    ///
527    /// Returns the error of a write or a flush.
528    ///
529    /// # Panics
530    ///
531    /// Panics when the run's config is not in the plan.
532    pub fn write_run(&mut self, outcome: &RunOutcome) -> io::Result<()> {
533        let Self {
534            plan,
535            runs,
536            series,
537            counts,
538        } = self;
539        let config = plan
540            .config(outcome.run.config_id)
541            .expect("a committed run's config is in its plan");
542        write_both(runs, series, counts, outcome, config)
543    }
544
545    /// Writes the series of `outcome`, a run of `config`, then its row, and flushes after each write.
546    ///
547    /// A search writes its runs this way. Its configs are in no plan.
548    ///
549    /// # Errors
550    ///
551    /// Returns the error of a write or a flush.
552    pub fn write_config_run(&mut self, outcome: &RunOutcome, config: &Config) -> io::Result<()> {
553        write_both(&mut self.runs, &mut self.series, &mut self.counts, outcome, config)
554    }
555
556    /// Counts of the rows written so far.
557    pub fn counts(&self) -> ResultCounts {
558        self.counts
559    }
560
561    /// Flushes both writers, and returns them as the writers of `runs.csv` and `series.csv`.
562    ///
563    /// # Errors
564    ///
565    /// Returns the error of a flush.
566    pub fn finish(self) -> io::Result<(W, W)> {
567        let series = self.series.into_inner()?;
568        let runs = self.runs.into_inner()?;
569        Ok((runs, series))
570    }
571}
572
573/// Writes the series of `outcome`, a run of `config`, to `series`, then its row to `runs`, flushing after each write,
574/// and counts the row in `counts`.
575fn write_both<W: Write>(
576    runs: &mut RunsWriter<W>,
577    series: &mut SeriesWriter<W>,
578    counts: &mut ResultCounts,
579    outcome: &RunOutcome,
580    config: &Config,
581) -> io::Result<()> {
582    series.write_run(outcome.run.run_id, &outcome.series)?;
583    series.flush()?;
584    runs.write_run(outcome, config)?;
585    runs.flush()?;
586    counts.count(outcome.status);
587    Ok(())
588}
589
590impl<W: Write> RunSink for OutputWriter<W> {
591    fn commit(&mut self, outcome: RunOutcome) -> io::Result<()> {
592        self.write_run(&outcome)
593    }
594}
595
596/// Results that cannot be written.
597#[derive(Debug)]
598pub enum OutputError {
599    /// Directory `dir` holds `file` from an earlier sweep or search.
600    HoldsResults {
601        /// Path of the output directory.
602        dir: PathBuf,
603        /// Name of the file found, such as `runs.csv` or a staged table.
604        file: String,
605    },
606    /// Another sweep, search or merge holds the lock of directory `dir` and writes to it.
607    Locked {
608        /// Path of the output directory.
609        dir: PathBuf,
610    },
611    /// The table at `path` is a symbolic link. A writer never follows a symbolic link.
612    Link {
613        /// Path of `runs.csv` or `series.csv` in the output directory.
614        path: PathBuf,
615    },
616    /// Reading `path` failed.
617    Read {
618        /// Path of `runs.csv`, or its file name alone for runs held in memory.
619        path: PathBuf,
620        /// Error from the reader, of kind `InvalidData` for text that is not UTF-8.
621        source: io::Error,
622    },
623    /// Creating or writing `path` failed.
624    Write {
625        /// Path of the file or directory, or the file name alone for results held in memory.
626        path: PathBuf,
627        /// Error from the operating system, or from the writer of results held in memory.
628        source: io::Error,
629    },
630    /// A table that cannot be read back.
631    Table(ReadError),
632    /// `runs.csv` cannot be summarized.
633    Summary(SummaryError),
634    /// Serializing the manifest failed.
635    Manifest(serde_json::Error),
636}
637
638impl fmt::Display for OutputError {
639    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
640        match self {
641            Self::HoldsResults { dir, file } => {
642                write!(f, "'{}' already holds results ({file})", dir.display())
643            }
644            Self::Locked { dir } => write!(
645                f,
646                "another sweep, search or merge is writing to '{}'. Wait for it to end, or stop it",
647                dir.display()
648            ),
649            Self::Link { path } => write!(
650                f,
651                "'{}' is a symbolic link. A sweep writes to no table through a link",
652                path.display()
653            ),
654            Self::Read { path, .. } => write!(f, "cannot read '{}'", path.display()),
655            Self::Write { path, .. } => write!(f, "cannot create or write '{}'", path.display()),
656            Self::Table(_) => f.write_str("cannot read a table back"),
657            Self::Summary(_) => f.write_str("cannot summarize the runs"),
658            Self::Manifest(_) => f.write_str("cannot serialize the manifest"),
659        }
660    }
661}
662
663impl std::error::Error for OutputError {
664    fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
665        match self {
666            Self::HoldsResults { .. } | Self::Locked { .. } | Self::Link { .. } => None,
667            Self::Read { source, .. } | Self::Write { source, .. } => Some(source),
668            Self::Table(error) => Some(error),
669            Self::Summary(error) => Some(error),
670            Self::Manifest(error) => Some(error),
671        }
672    }
673}
674
675#[cfg(test)]
676mod tests {
677    use std::fs;
678
679    use super::{
680        LOCK_FILE, MANIFEST_FILE, OutputDir, OutputError, RUNS_FILE, SERIES_FILE, STAGED_MARKER, staged_path,
681        table_paths,
682    };
683    use crate::tests::support::ScratchDir;
684
685    #[test]
686    fn a_replacement_cut_short_is_finished_once_both_tables_are_staged() {
687        let scratch = ScratchDir::new("staged-tables");
688        drop(OutputDir::create(scratch.path()).expect("a scratch directory"));
689        let path = scratch.path();
690        for file in [RUNS_FILE, SERIES_FILE] {
691            fs::write(path.join(file), "old\n").expect("a table writes");
692            fs::write(staged_path(path, file), "new\n").expect("a staged table writes");
693        }
694        assert_eq!(table_paths(path).0, path.join(RUNS_FILE), "no marker, no replacement");
695        drop(OutputDir::open(path).expect("the directory opens"));
696        for file in [RUNS_FILE, SERIES_FILE] {
697            assert_eq!(fs::read_to_string(path.join(file)).expect("a table reads"), "old\n");
698            assert!(
699                !staged_path(path, file).exists(),
700                "a staged table without the marker is dropped"
701            );
702        }
703
704        // The process ended after moving runs.csv into place and before series.csv.
705        fs::write(staged_path(path, SERIES_FILE), "new\n").expect("a staged table writes");
706        fs::write(path.join(RUNS_FILE), "new\n").expect("a table writes");
707        fs::write(path.join(STAGED_MARKER), "").expect("the marker writes");
708        assert_eq!(
709            table_paths(path),
710            (path.join(RUNS_FILE), staged_path(path, SERIES_FILE))
711        );
712        drop(OutputDir::open(path).expect("the directory opens"));
713        for file in [RUNS_FILE, SERIES_FILE] {
714            assert_eq!(fs::read_to_string(path.join(file)).expect("a table reads"), "new\n");
715        }
716        assert!(!path.join(STAGED_MARKER).exists());
717    }
718
719    #[test]
720    fn a_second_writer_is_refused_while_the_first_holds_the_lock() {
721        let scratch = ScratchDir::new("locked-directory");
722        let path = scratch.path();
723        let first = OutputDir::create(path).expect("a scratch directory");
724        assert!(path.join(LOCK_FILE).exists(), "the writer holds the lock file");
725        for second in [OutputDir::open(path), OutputDir::create(path)] {
726            assert!(matches!(second, Err(OutputError::Locked { .. })), "{second:?}");
727        }
728        assert!(matches!(
729            OutputDir::check_unlocked(path),
730            Err(OutputError::Locked { .. })
731        ));
732        drop(first);
733        assert!(path.join(LOCK_FILE).exists(), "the writer leaves its lock file");
734        assert!(OutputDir::check_unlocked(path).is_ok());
735        let next = OutputDir::open(path).expect("the next writer locks the file again");
736        assert!(matches!(OutputDir::open(path), Err(OutputError::Locked { .. })));
737        drop(next);
738
739        // A killed process leaves its lock file behind, and the lock is released when the process ends.
740        fs::write(path.join(LOCK_FILE), "").expect("a stale lock file writes");
741        assert!(OutputDir::check_unlocked(path).is_ok());
742        drop(OutputDir::open(path).expect("a stale lock file is taken over"));
743    }
744
745    #[cfg(unix)]
746    #[test]
747    fn a_writer_follows_no_link_in_its_directory() {
748        use std::os::unix::fs::symlink;
749
750        use henad_core::explore::spec::SweepSpec;
751
752        use super::SUMMARY_FILE;
753        use crate::exec::Concurrency;
754        use crate::tests::support::{entry, sweep};
755
756        let scratch = ScratchDir::new("planted-links");
757        fs::create_dir_all(scratch.path()).expect("a scratch directory");
758        let victim = scratch.path().join("victim.txt");
759        fs::write(&victim, "kept\n").expect("a scratch file");
760        let path = scratch.path().join("out");
761        fs::create_dir(&path).expect("a scratch directory");
762
763        symlink(&victim, path.join(format!("{MANIFEST_FILE}.partial"))).expect("a link");
764        symlink(&victim, path.join(LOCK_FILE)).expect("a link");
765        let mut spec = SweepSpec::new("game_of_life");
766        spec.fixed = vec![
767            ("grid_width".to_owned(), "8".to_owned()),
768            ("grid_height".to_owned(), "8".to_owned()),
769        ];
770        spec.run.steps = 2;
771        sweep(&entry("game_of_life", None), None, &spec, &path, Concurrency::Auto);
772        assert_eq!(fs::read_to_string(&victim).expect("the victim reads"), "kept\n");
773        let manifest = fs::symlink_metadata(path.join(MANIFEST_FILE)).expect("the manifest is written");
774        assert!(manifest.file_type().is_file());
775
776        let dangling = scratch.path().join("dangling");
777        fs::create_dir(&dangling).expect("a scratch directory");
778        symlink(scratch.path().join("nowhere.csv"), dangling.join(SUMMARY_FILE)).expect("a link");
779        let refused = OutputDir::create(&dangling);
780        assert!(
781            matches!(&refused, Err(OutputError::HoldsResults { file, .. }) if file == SUMMARY_FILE),
782            "{refused:?}"
783        );
784
785        let linked = scratch.path().join("linked");
786        fs::create_dir(&linked).expect("a scratch directory");
787        symlink(&victim, linked.join(RUNS_FILE)).expect("a link");
788        let refused = OutputDir::open(&linked);
789        assert!(matches!(refused, Err(OutputError::Link { .. })), "{refused:?}");
790        assert_eq!(fs::read_to_string(&victim).expect("the victim reads"), "kept\n");
791    }
792
793    #[test]
794    fn a_directory_holding_results_is_refused() {
795        let scratch = ScratchDir::new("holding-results");
796        let nested = scratch.path().join("a").join("b");
797        let dir = OutputDir::create(&nested).expect("a missing directory is created");
798        assert!(dir.path().is_dir());
799        drop(dir);
800        assert!(OutputDir::create(&nested).is_ok(), "an empty directory is taken");
801
802        fs::write(nested.join("notes.txt"), "kept").expect("a scratch file");
803        assert!(
804            OutputDir::create(&nested).is_ok(),
805            "a file a sweep does not write is left alone"
806        );
807
808        fs::write(nested.join(MANIFEST_FILE), "{}").expect("a scratch file");
809        let error = OutputDir::create(&nested).expect_err("a manifest marks results");
810        assert!(
811            matches!(&error, OutputError::HoldsResults { file, .. } if file == MANIFEST_FILE),
812            "{error:?}"
813        );
814        fs::remove_file(nested.join(MANIFEST_FILE)).expect("the manifest is removed");
815        fs::write(nested.join(RUNS_FILE), "").expect("a scratch file");
816        assert!(OutputDir::check_free(&nested).is_err(), "so does an empty runs.csv");
817        fs::remove_file(nested.join(RUNS_FILE)).expect("the table is removed");
818
819        for staged in [staged_path(&nested, SERIES_FILE), nested.join(STAGED_MARKER)] {
820            fs::write(&staged, "").expect("a scratch file");
821            let error = OutputDir::create(&nested).expect_err("a staged file marks results");
822            let name = staged.file_name().map(|name| name.to_string_lossy().into_owned());
823            assert!(
824                matches!(&error, OutputError::HoldsResults { file, .. } if Some(file) == name.as_ref()),
825                "{error:?}"
826            );
827            fs::remove_file(&staged).expect("the staged file is removed");
828        }
829        assert!(OutputDir::check_free(&nested).is_ok());
830    }
831}