Skip to main content

mesofact_dev/
watcher.rs

1//! File-watch + auto-rebuild loop for `mesofact-static` workloads.
2//!
3//! Watches `<workload>/src/` via [`notify`], debounces edits, rebuilds
4//! ([`BuildDriver`] — in-process by default, `sh -c` when `workload.toml`
5//! declares a `[build] command`), then snapshots
6//! `<workload>/dist/` into `<workload>/.mesofact-dev/gen-<N>/` and flips a
7//! [`DistPointer`](super::DistPointer) so the running server starts serving
8//! the new artifact on the next request. Build stdout/stderr inherits the
9//! parent's, so output shows up in the operator's terminal or the Run-tab
10//! log surface.
11//!
12//! Atomicity story: each generation is a separate directory. `dist/` is
13//! copied into `.mesofact-dev/gen-<N>-staging/` and then atomic-renamed to
14//! `.mesofact-dev/gen-<N>/` on the same filesystem; the pointer flip is a
15//! single `RwLock` write; in-flight reads keep using the `PathBuf` they
16//! already cloned, so a request that started reading
17//! `gen-<N-1>/html/index.html` doesn't observe a torn write. `dist/` itself
18//! is left intact so concurrent publishers (mesofact-publisher to R2,
19//! qed-run to pond/MinIO) can read it without racing the watcher. GC keeps
20//! the last two generations.
21//!
22//! @yah:ticket(R255-S5, "Decide tier-1 object store: s3s-fs vs serve-off-disk vs self-mock")
23//! @yah:assignee(agent:claude)
24//! @yah:at(2026-05-25T20:08:09Z)
25//! @yah:kind(spike)
26//! @yah:status(review)
27//! @yah:parent(R255)
28//! @yah:next("decide on the single fact: does the almanac/runtime fetch via S3 API calls or plain CDN GET? if S3, tier-1 needs s3s-fs regardless")
29//! @yah:next("if s3s-fs: run the existing publish_to_local_sim unchanged against it so tier-1 becomes a strict subset of tier-2 (only diff: in-process s3s-fs vs containerized MinIO+Caddy+yubaba)")
30//! @yah:next("reject self-mock: reimplementing SigV4 + bucket-policy + error shapes drifts from MinIO and defeats the point")
31//! @yah:gotcha("license-check s3s before adopting — must be MIT/BSD/Apache-2.0/ISC (believed Apache-2.0)")
32//! @yah:assumes("tier-1 read path is plain HTTP-GET (Caddy proxies the public bucket), so serve-off-disk is faithful for reads; only the build->PUT->read publish contract is unexercised at tier 1")
33//! @yah:handoff("Spike closed: serve-off-disk wins for tier-1 (dev). Read path is plain HTTP GET via axum ServeDir — browsers never call S3 APIs. The build→PUT→serve publish contract is exercised at tier-2 (sim+MinIO, R256 T1–T5 in review). No s3s-fs needed at tier-1; self-mock rejected as before. The @yah:assumes fact was correct. Only a tier-2 concern: R256-F8 tracks the publish-to-MinIO watcher sink for hot-ish sim reload, which uses the real MinIO container (not s3s-fs).")
34//! @yah:verify("No code change needed — tier-1 is already serve-off-disk. Verify the assumes by checking mesofact-dev routes: nothing in app/yah/web/src/ or mesofact-runtime makes S3 API calls client-side.")
35//!
36//! @yah:ticket(R256-F8, "Parameterize watcher sink: DistPointer (serve off disk) vs publish-to-MinIO")
37//! @yah:assignee(agent:claude)
38//! @yah:at(2026-05-25T20:30:17Z)
39//! @yah:status(review)
40//! @yah:parent(R256)
41//! @yah:next("introduce a sink abstraction so the same watch->rebuild loop can target either the in-process DistPointer (dev tier) or publish_to_local_sim (sim tier -> MinIO container)")
42//! @yah:next("this is what gives the containerized sim a hot-ish reload (edit -> rebuild -> republish -> Caddy serves new artifact) WITHOUT putting mesofact-dev inside a container")
43//! @yah:next("keep the host-side watcher as the only dev-mode component; the sim containers stay pure mesofact-core + Caddy + MinIO")
44//! @yah:assumes("today the watcher's only sink is the in-process DistPointer (build_and_swap -> pointer.set at watcher.rs:253); there is no publish sink")
45//! @arch:see(.yah/docs/working/mesofact-dev-camp-embedding.md)
46//! @yah:handoff("PostBuildFn + PostBuildFuture type aliases added to watcher.rs. Watcher gains post_build: Option<PostBuildFn> field and with_post_build(self, f) builder. build_and_swap calls the hook after the pointer flip — failure logs a warning but does not fail the rebuild (in-process dev server stays healthy). No new dep on cloud: the hook is a generic async closure; camp.rs (which already depends on both mesofact-dev and cloud) will wire publish_to_local_sim into the closure. Two new tests: post_build_hook_receives_gen_dir_on_success + post_build_hook_failure_does_not_fail_rebuild. All 20 mesofact-dev tests pass; cargo check cloud+yah+desktop clean.")
47//! @yah:verify("cargo test -p mesofact-dev --locked  # 20 passed")
48//! @yah:verify("cargo check -p cloud -p yah -p desktop --locked")
49//!
50//! @arch:see(app/yah/cli/src/camp.rs)
51//!
52
53use std::{
54    path::{Path, PathBuf},
55    process::Stdio,
56    sync::atomic::{AtomicBool, Ordering},
57    sync::Arc,
58    time::Duration,
59};
60
61use anyhow::{Context, Result};
62use notify::{Event, EventKind, RecursiveMode, Watcher as _};
63use serde::Deserialize;
64use tokio::process::Command;
65use tokio::sync::mpsc;
66use tracing::{error, info, warn};
67
68use crate::DistPointer;
69
70// ── Sink hook ─────────────────────────────────────────────────────────────────
71
72/// Boxed future returned by a [`PostBuildFn`].
73pub type PostBuildFuture =
74    std::pin::Pin<Box<dyn std::future::Future<Output = anyhow::Result<()>> + Send + 'static>>;
75
76/// Optional post-build hook. Called with the snapshot directory (the `gen-N/`
77/// directory, whose `html/` subdirectory is what the DistPointer points at)
78/// after every successful build, immediately after the pointer flip.
79///
80/// Failure is logged as a warning but does not roll back the pointer flip —
81/// the in-process server continues to serve the new snapshot; only the
82/// secondary sink (e.g. MinIO publish for the sim tier) is stale.
83///
84/// Wired by camp for the sim tier: the closure calls
85/// `cloud::publish_to_local_sim` so the Caddy+MinIO stack serves the updated
86/// artifact without requiring a full container restart. The dev tier leaves
87/// this `None` — the DistPointer alone is sufficient.
88pub type PostBuildFn =
89    Box<dyn Fn(std::path::PathBuf) -> PostBuildFuture + Send + Sync + 'static>;
90
91const DEBOUNCE_MS: u64 = 200;
92const GENS_TO_KEEP: usize = 2;
93const STATE_DIR_NAME: &str = ".mesofact-dev";
94
95/// How a rebuild is produced.
96///
97/// The dev loop needs a bundler; the prod binary must never link one (W225
98/// §2/§3). That constraint says where `mesofact-build` may be linked — it
99/// does not say the dev tier has to reach it through a *third executable*,
100/// and until R759-T4 it did: the default was `sh -c "bun run build"`, so the
101/// documented "two binaries and no Node" loop actually needed bun on PATH,
102/// and the release ships only `mesofact` and `mesofact-dev`
103/// (`.yah/qed/release-build.toml`). [`BuildDriver::InProcess`] closes that:
104/// dev is the tier that is *allowed* to be fat, so it links the pipeline and
105/// calls it directly.
106#[derive(Debug, Clone)]
107pub enum BuildDriver {
108    /// Run `mesofact-build`'s pipeline in this process. Needs no package
109    /// manager and no Node — the pipeline materializes `node_modules` from
110    /// the project's lockfile itself.
111    ///
112    /// Only constructible with the `build` feature; [`WatchOptions`] falls
113    /// back to [`BuildDriver::Shell`] without it.
114    InProcess,
115    /// `sh -c <command>`, run from the workload root. What `workload.toml`'s
116    /// `[build] command` selects, and what a project with its own bundler
117    /// step wants.
118    Shell(String),
119}
120
121impl BuildDriver {
122    /// What `workload.toml` gets when it declares no `[build] command`.
123    ///
124    /// Without the `build` feature there is no pipeline linked in, so this
125    /// keeps the historical shell default rather than silently doing
126    /// nothing.
127    pub fn default_for_workload() -> Self {
128        #[cfg(feature = "build")]
129        {
130            Self::InProcess
131        }
132        #[cfg(not(feature = "build"))]
133        {
134            Self::Shell(LEGACY_SHELL_BUILD.to_string())
135        }
136    }
137}
138
139/// The pre-R759-T4 default. Retained as the `--no-default-features` fallback
140/// and named so the one place it still applies is greppable.
141pub const LEGACY_SHELL_BUILD: &str = "bun run build";
142
143/// Knobs for the watch loop.
144#[derive(Debug, Clone)]
145pub struct WatchOptions {
146    /// Directory to watch recursively for source edits.
147    pub watch_dir: PathBuf,
148    /// How a rebuild is produced.
149    pub build: BuildDriver,
150    /// Directory the build writes into (the parent of `html/`). Relative
151    /// paths resolve against the workload root.
152    pub build_out_dir: PathBuf,
153    /// Where generation snapshots live. Defaults to `<workload>/.mesofact-dev`.
154    pub state_dir: PathBuf,
155    /// Coalesce events landing within this window into one rebuild.
156    pub debounce: Duration,
157    /// Run an initial build at startup even if `dist/html/` already exists.
158    pub initial_build: bool,
159    /// Extra env vars injected into the build subprocess (R490-F7, R584-T1).
160    /// Carries the dev S3 coords (`R2_ENDPOINT`, `R2_BUCKET`, camp-injected
161    /// creds) so a workload's build-time `r2` reads resolve against the
162    /// camp's dev-tier S3 driver instead of real R2.
163    pub build_env: Vec<(String, String)>,
164}
165
166impl WatchOptions {
167    /// Best-effort defaults: watch `<workload>/src`, build in-process, output
168    /// to `<workload>/dist`, snapshot under `<workload>/.mesofact-dev`. Reads
169    /// `<workload>/workload.toml` if present to override `build.command` /
170    /// `build.out_dir` (so the dev server agrees with the production
171    /// reconciler).
172    ///
173    /// A declared `[build] command` still wins and still runs through `sh
174    /// -c`; every workload in this repo declares one, so the in-process
175    /// default only changes what a workload that declares *nothing* gets —
176    /// which is the scaffolded project `mesofact new` emits.
177    pub fn defaults_for_workload(workload: &Path) -> Self {
178        let mut opts = Self {
179            watch_dir: workload.join("src"),
180            build: BuildDriver::default_for_workload(),
181            build_out_dir: workload.join("dist"),
182            state_dir: workload.join(STATE_DIR_NAME),
183            debounce: Duration::from_millis(DEBOUNCE_MS),
184            initial_build: true,
185            build_env: Vec::new(),
186        };
187        if let Ok(text) = std::fs::read_to_string(workload.join("workload.toml")) {
188            if let Ok(parsed) = toml::from_str::<WorkloadTomlPartial>(&text) {
189                if let Some(build) = parsed.build {
190                    if let Some(cmd) = build.command {
191                        opts.build = BuildDriver::Shell(cmd);
192                    }
193                    if let Some(out) = build.out_dir {
194                        opts.build_out_dir = if out.is_absolute() {
195                            out
196                        } else {
197                            workload.join(out)
198                        };
199                    }
200                }
201            }
202        }
203        opts
204    }
205}
206
207#[derive(Deserialize, Default)]
208struct WorkloadTomlPartial {
209    build: Option<BuildPartial>,
210}
211
212#[derive(Deserialize, Default)]
213struct BuildPartial {
214    command: Option<String>,
215    out_dir: Option<PathBuf>,
216}
217
218/// File-watch + rebuild loop. Construct with [`Watcher::new`] and drive with
219/// [`Watcher::run`]; pair with a [`crate::Server`] sharing the same
220/// [`DistPointer`].
221///
222/// For the sim tier, set a post-build publish hook with [`Watcher::with_post_build`]
223/// so each successful rebuild also pushes the artifact to MinIO. The pointer
224/// flip and the MinIO publish are independent: a failing publish logs a warning
225/// but does not prevent the pointer from advancing.
226pub struct Watcher {
227    workload: PathBuf,
228    pointer: DistPointer,
229    options: WatchOptions,
230    post_build: Option<PostBuildFn>,
231}
232
233impl Watcher {
234    pub fn new(workload: impl Into<PathBuf>, pointer: DistPointer, options: WatchOptions) -> Self {
235        Self {
236            workload: workload.into(),
237            pointer,
238            options,
239            post_build: None,
240        }
241    }
242
243    /// Attach a post-build publish hook (sim tier).
244    ///
245    /// After every successful build the hook receives the snapshot directory
246    /// (`gen-N/`, whose `html/` is what the pointer already points at). Camp
247    /// uses this to call `cloud::publish_to_local_sim` so Caddy+MinIO stays
248    /// in sync without a container restart.
249    pub fn with_post_build(mut self, f: PostBuildFn) -> Self {
250        self.post_build = Some(f);
251        self
252    }
253
254    /// Drives the watch loop until the notify channel closes or the inner
255    /// task is cancelled. Logs failures rather than exiting — a failed
256    /// rebuild leaves the previous snapshot in place.
257    pub async fn run(self) -> Result<()> {
258        tokio::fs::create_dir_all(&self.options.state_dir)
259            .await
260            .with_context(|| {
261                format!("creating state dir {}", self.options.state_dir.display())
262            })?;
263
264        // Pre-seed the gen counter from existing snapshot dirs so a restart
265        // doesn't try to rename into a name that already exists.
266        let mut next_gen = next_gen_number(&self.options.state_dir).await?;
267
268        let (event_tx, mut event_rx) = mpsc::unbounded_channel::<Event>();
269
270        // Keep the watcher alive for the lifetime of this task.
271        let mut watcher = notify::recommended_watcher(move |res: notify::Result<Event>| match res {
272            Ok(ev) => {
273                let _ = event_tx.send(ev);
274            }
275            Err(e) => warn!(error = %e, "notify watcher error"),
276        })?;
277
278        if !self.options.watch_dir.exists() {
279            warn!(
280                dir = %self.options.watch_dir.display(),
281                "watch dir missing — auto-rebuild disabled until it appears",
282            );
283        } else {
284            watcher
285                .watch(&self.options.watch_dir, RecursiveMode::Recursive)
286                .with_context(|| {
287                    format!("watching {}", self.options.watch_dir.display())
288                })?;
289            info!(dir = %self.options.watch_dir.display(), "watching for changes");
290        }
291
292        if self.options.initial_build {
293            info!("running initial build");
294            if let Err(e) = self.build_and_swap(&mut next_gen).await {
295                error!(error = format!("{e:#}"), "initial build failed; serving existing snapshot if any");
296            }
297        }
298
299        // Debounce loop.
300        loop {
301            let Some(first) = event_rx.recv().await else {
302                info!("watcher channel closed");
303                break;
304            };
305            let mut actionable = is_actionable(&first);
306
307            let deadline = tokio::time::Instant::now() + self.options.debounce;
308            loop {
309                let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
310                if remaining.is_zero() {
311                    break;
312                }
313                match tokio::time::timeout(remaining, event_rx.recv()).await {
314                    Ok(Some(ev)) => {
315                        actionable |= is_actionable(&ev);
316                    }
317                    Ok(None) => return Ok(()),
318                    Err(_) => break,
319                }
320            }
321
322            if !actionable {
323                continue;
324            }
325
326            info!("source change detected; rebuilding");
327            if let Err(e) = self.build_and_swap(&mut next_gen).await {
328                error!(error = format!("{e:#}"), "rebuild failed; keeping previous snapshot");
329            }
330        }
331
332        Ok(())
333    }
334
335    /// One-shot rebuild — run the build command, snapshot, swap. Exposed
336    /// for the reconciler to drive builds explicitly (R255-T3).
337    pub async fn rebuild(&self) -> Result<PathBuf> {
338        let mut next_gen = next_gen_number(&self.options.state_dir).await?;
339        self.build_and_swap(&mut next_gen).await
340    }
341
342    /// Produce `dist/` however [`BuildDriver`] says to.
343    async fn run_build(&self) -> Result<()> {
344        match &self.options.build {
345            BuildDriver::Shell(command) => {
346                let status = Command::new("sh")
347                    .arg("-c")
348                    .arg(command)
349                    .current_dir(&self.workload)
350                    .envs(self.options.build_env.iter().cloned())
351                    .stdout(Stdio::inherit())
352                    .stderr(Stdio::inherit())
353                    .status()
354                    .await
355                    .with_context(|| format!("spawning build: {command}"))?;
356                if !status.success() {
357                    anyhow::bail!("build exited with {}", status);
358                }
359                Ok(())
360            }
361            BuildDriver::InProcess => self.run_build_in_process().await,
362        }
363    }
364
365    #[cfg(feature = "build")]
366    async fn run_build_in_process(&self) -> Result<()> {
367        // `build_env` exists for the shell child (R490-F7 hands the dev S3
368        // coords to a build subprocess). In-process there is no child to
369        // inherit them, so set them on this process for the duration of the
370        // build — the pipeline's own R2 reads look them up through `env`
371        // exactly as the child did. Rebuilds are serialized by the watch
372        // loop, so nothing else in this process races the assignment.
373        for (k, v) in &self.options.build_env {
374            std::env::set_var(k, v);
375        }
376        mesofact::build::pipeline::build(mesofact::build::pipeline::BuildOptions {
377            project_root: self.workload.clone(),
378            out_dir: Some(self.options.build_out_dir.clone()),
379            build_id: None,
380            // `Auto`: materialize `node_modules` from the project's lockfile
381            // only when it is missing. This is the step that makes the loop
382            // work with no package manager and no Node on PATH.
383            install: mesofact::build::pipeline::InstallMode::Auto,
384        })
385        .await
386        .context("in-process build")?;
387        Ok(())
388    }
389
390    #[cfg(not(feature = "build"))]
391    async fn run_build_in_process(&self) -> Result<()> {
392        anyhow::bail!(
393            "this mesofact-dev was built without the `build` feature, so it cannot build in-process; declare `[build] command` in workload.toml"
394        )
395    }
396
397    async fn build_and_swap(&self, next_gen: &mut u64) -> Result<PathBuf> {
398        self.run_build().await?;
399
400        let html_src = self.options.build_out_dir.join("html");
401        if !html_src.is_dir() {
402            anyhow::bail!(
403                "build did not produce {} (expected html/ under {})",
404                html_src.display(),
405                self.options.build_out_dir.display(),
406            );
407        }
408
409        tokio::fs::create_dir_all(&self.options.state_dir)
410            .await
411            .with_context(|| format!("creating state dir {}", self.options.state_dir.display()))?;
412
413        let n = *next_gen;
414        *next_gen = n + 1;
415        let gen_dir = self.options.state_dir.join(format!("gen-{n}"));
416        if gen_dir.exists() {
417            tokio::fs::remove_dir_all(&gen_dir).await.ok();
418        }
419        // Copy build_out_dir into a staging dir, then atomic-rename staging
420        // → gen-N. Leaves build_out_dir intact for concurrent publishers
421        // (mesofact-publisher to R2, qed-run to pond/MinIO) that read the
422        // same dist/ tree. Atomicity of the gen-N flip is preserved by the
423        // staging → gen-N rename, which is a single-fs rename of a directory
424        // that no one else knows about yet.
425        let staging = self.options.state_dir.join(format!("gen-{n}-staging"));
426        if staging.exists() {
427            tokio::fs::remove_dir_all(&staging).await.ok();
428        }
429        copy_dir_recursive(&self.options.build_out_dir, &staging)
430            .await
431            .with_context(|| {
432                format!(
433                    "copying {} -> {}",
434                    self.options.build_out_dir.display(),
435                    staging.display(),
436                )
437            })?;
438        tokio::fs::rename(&staging, &gen_dir)
439            .await
440            .with_context(|| {
441                format!(
442                    "renaming {} -> {}",
443                    staging.display(),
444                    gen_dir.display(),
445                )
446            })?;
447
448        let served = gen_dir.join("html");
449        self.pointer.set(served.clone());
450        info!(gen = n, served = %served.display(), "snapshot ready; pointer updated");
451
452        // Optional post-build hook: publish to sim-tier object store (MinIO).
453        // Failure is non-fatal — the pointer is already flipped and the
454        // in-process server serves the new snapshot. A stale sim tier logs a
455        // warning so operators notice without breaking the dev loop.
456        if let Some(hook) = &self.post_build {
457            if let Err(e) = hook(gen_dir.clone()).await {
458                warn!(error = format!("{e:#}"), gen = n, "post-build hook failed");
459            }
460        }
461
462        if let Err(e) = gc_generations(&self.options.state_dir, GENS_TO_KEEP).await {
463            warn!(error = %e, "gc of old generations failed");
464        }
465
466        Ok(served)
467    }
468}
469
470/// Spawn a watch loop on a background tokio task. Returns a handle that
471/// stops the loop when dropped (the watcher's internal channel closes when
472/// the task exits, which is fine for a process-lifetime dev server).
473pub fn spawn(watcher: Watcher) -> WatcherHandle {
474    let running = Arc::new(AtomicBool::new(true));
475    let flag = Arc::clone(&running);
476    let join = tokio::spawn(async move {
477        let result = watcher.run().await;
478        flag.store(false, Ordering::SeqCst);
479        if let Err(e) = result {
480            error!(error = format!("{e:#}"), "watcher exited with error");
481        }
482    });
483    WatcherHandle { running, _join: join }
484}
485
486/// Handle that keeps a spawned watcher alive. Dropping it does not cancel
487/// the task (tokio detaches on drop); use [`WatcherHandle::is_running`] to
488/// observe lifecycle.
489pub struct WatcherHandle {
490    running: Arc<AtomicBool>,
491    _join: tokio::task::JoinHandle<()>,
492}
493
494impl WatcherHandle {
495    pub fn is_running(&self) -> bool {
496        self.running.load(Ordering::SeqCst)
497    }
498}
499
500fn is_actionable(ev: &Event) -> bool {
501    matches!(
502        ev.kind,
503        EventKind::Create(_) | EventKind::Modify(_) | EventKind::Remove(_)
504    )
505}
506
507async fn copy_dir_recursive(src: &Path, dst: &Path) -> Result<()> {
508    tokio::fs::create_dir_all(dst)
509        .await
510        .with_context(|| format!("creating {}", dst.display()))?;
511    let mut stack: Vec<(PathBuf, PathBuf)> = vec![(src.to_path_buf(), dst.to_path_buf())];
512    while let Some((s, d)) = stack.pop() {
513        let mut rd = tokio::fs::read_dir(&s)
514            .await
515            .with_context(|| format!("reading {}", s.display()))?;
516        while let Some(ent) = rd.next_entry().await? {
517            let ft = ent.file_type().await?;
518            let from = ent.path();
519            let to = d.join(ent.file_name());
520            if ft.is_dir() {
521                tokio::fs::create_dir_all(&to)
522                    .await
523                    .with_context(|| format!("creating {}", to.display()))?;
524                stack.push((from, to));
525            } else {
526                tokio::fs::copy(&from, &to)
527                    .await
528                    .with_context(|| format!("copying {} -> {}", from.display(), to.display()))?;
529            }
530        }
531    }
532    Ok(())
533}
534
535async fn next_gen_number(state_dir: &Path) -> Result<u64> {
536    let mut max_seen: Option<u64> = None;
537    let mut rd = match tokio::fs::read_dir(state_dir).await {
538        Ok(rd) => rd,
539        Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(0),
540        Err(e) => return Err(e).context("reading state dir for gen numbering"),
541    };
542    while let Some(ent) = rd.next_entry().await? {
543        let name = ent.file_name();
544        let s = name.to_string_lossy();
545        if let Some(rest) = s.strip_prefix("gen-") {
546            if let Ok(n) = rest.parse::<u64>() {
547                max_seen = Some(max_seen.map_or(n, |m| m.max(n)));
548            }
549        }
550    }
551    Ok(max_seen.map(|m| m + 1).unwrap_or(0))
552}
553
554async fn gc_generations(state_dir: &Path, keep: usize) -> Result<()> {
555    let mut entries: Vec<(u64, PathBuf)> = Vec::new();
556    let mut rd = tokio::fs::read_dir(state_dir).await?;
557    while let Some(ent) = rd.next_entry().await? {
558        let s = ent.file_name();
559        let s = s.to_string_lossy();
560        if let Some(rest) = s.strip_prefix("gen-") {
561            if let Ok(n) = rest.parse::<u64>() {
562                entries.push((n, ent.path()));
563            }
564        }
565    }
566    entries.sort_by_key(|e| std::cmp::Reverse(e.0));
567    for (_, path) in entries.into_iter().skip(keep) {
568        if let Err(e) = tokio::fs::remove_dir_all(&path).await {
569            warn!(error = %e, dir = %path.display(), "failed to gc snapshot");
570        }
571    }
572    Ok(())
573}
574
575#[cfg(test)]
576mod tests {
577    use super::*;
578    use std::time::Duration;
579    use tempfile::tempdir;
580
581    fn write(path: &Path, body: &str) {
582        if let Some(parent) = path.parent() {
583            std::fs::create_dir_all(parent).unwrap();
584        }
585        std::fs::write(path, body).unwrap();
586    }
587
588    /// Build script: `mkdir -p dist/html && cat src/index.txt > dist/html/index.html`.
589    /// Idempotent, no external deps; lets us simulate a real bun build.
590    const FAKE_BUILD: &str = "mkdir -p dist/html && cp src/index.txt dist/html/index.html";
591
592    #[tokio::test]
593    async fn defaults_for_workload_reads_workload_toml() {
594        let dir = tempdir().unwrap();
595        write(
596            &dir.path().join("workload.toml"),
597            r#"
598kind = "mesofact-static"
599[build]
600command = "echo built"
601out_dir = "outdir"
602"#,
603        );
604        let opts = WatchOptions::defaults_for_workload(dir.path());
605        assert!(
606            matches!(&opts.build, BuildDriver::Shell(c) if c == "echo built"),
607            "a declared [build] command must still win over the in-process default: {:?}",
608            opts.build
609        );
610        assert_eq!(opts.build_out_dir, dir.path().join("outdir"));
611        assert_eq!(opts.watch_dir, dir.path().join("src"));
612        assert_eq!(opts.state_dir, dir.path().join(".mesofact-dev"));
613    }
614
615    #[tokio::test]
616    async fn defaults_for_workload_without_toml_uses_hardcoded_defaults() {
617        let dir = tempdir().unwrap();
618        let opts = WatchOptions::defaults_for_workload(dir.path());
619        assert_eq!(opts.build_out_dir, dir.path().join("dist"));
620
621        // R759-T4: a workload that declares no build command builds
622        // in-process. This is what makes `mesofact new` + `mesofact-dev .`
623        // work with only the two shipped binaries — the old `bun run build`
624        // default needed bun on PATH and a `mesofact-build` the release does
625        // not ship.
626        #[cfg(feature = "build")]
627        assert!(
628            matches!(opts.build, BuildDriver::InProcess),
629            "{:?}",
630            opts.build
631        );
632        #[cfg(not(feature = "build"))]
633        assert!(
634            matches!(&opts.build, BuildDriver::Shell(c) if c == LEGACY_SHELL_BUILD),
635            "{:?}",
636            opts.build
637        );
638    }
639
640    #[tokio::test]
641    async fn next_gen_number_seeds_from_existing_snapshots() {
642        let dir = tempdir().unwrap();
643        for n in [0u64, 3, 7] {
644            std::fs::create_dir_all(dir.path().join(format!("gen-{n}"))).unwrap();
645        }
646        std::fs::create_dir_all(dir.path().join("not-a-gen")).unwrap();
647        let n = next_gen_number(dir.path()).await.unwrap();
648        assert_eq!(n, 8);
649    }
650
651    #[tokio::test]
652    async fn next_gen_number_returns_zero_for_missing_dir() {
653        let dir = tempdir().unwrap();
654        let missing = dir.path().join("nope");
655        assert_eq!(next_gen_number(&missing).await.unwrap(), 0);
656    }
657
658    #[tokio::test]
659    async fn gc_keeps_last_n_generations() {
660        let dir = tempdir().unwrap();
661        for n in 0..5u64 {
662            std::fs::create_dir_all(dir.path().join(format!("gen-{n}"))).unwrap();
663        }
664        gc_generations(dir.path(), 2).await.unwrap();
665        assert!(dir.path().join("gen-4").is_dir());
666        assert!(dir.path().join("gen-3").is_dir());
667        assert!(!dir.path().join("gen-2").exists());
668        assert!(!dir.path().join("gen-0").exists());
669    }
670
671    #[tokio::test]
672    async fn rebuild_runs_build_command_and_swaps_pointer() {
673        let workload = tempdir().unwrap();
674        write(&workload.path().join("src/index.txt"), "<h1>A</h1>");
675
676        let pointer = DistPointer::new(workload.path().join("dist").join("html"));
677        let mut opts = WatchOptions::defaults_for_workload(workload.path());
678        opts.build = BuildDriver::Shell(FAKE_BUILD.to_string());
679        let watcher = Watcher::new(workload.path(), pointer.clone(), opts);
680
681        let served = watcher.rebuild().await.unwrap();
682        assert!(served.starts_with(workload.path().join(".mesofact-dev")));
683        assert_eq!(pointer.current(), served);
684        let body = std::fs::read_to_string(served.join("index.html")).unwrap();
685        assert!(body.contains("<h1>A</h1>"));
686
687        // dist/ must remain intact for concurrent publishers (R492-B1).
688        let dist_html = workload.path().join("dist").join("html").join("index.html");
689        let dist_body = std::fs::read_to_string(&dist_html)
690            .expect("dist/html/index.html should still exist after rebuild");
691        assert_eq!(dist_body, body, "dist/ and gen-N/ must hold the same bytes");
692
693        // Staging dir must be cleaned up by the rename.
694        let staging = workload.path().join(".mesofact-dev").join("gen-0-staging");
695        assert!(!staging.exists(), "staging dir should be gone after rename");
696    }
697
698    #[tokio::test]
699    async fn post_build_hook_receives_gen_dir_on_success() {
700        use std::sync::{Arc, Mutex};
701
702        let workload = tempdir().unwrap();
703        write(&workload.path().join("src/index.txt"), "<h1>hook</h1>");
704
705        let pointer = DistPointer::new(workload.path().join("dist").join("html"));
706        let mut opts = WatchOptions::defaults_for_workload(workload.path());
707        opts.build = BuildDriver::Shell(FAKE_BUILD.to_string());
708
709        let received: Arc<Mutex<Vec<PathBuf>>> = Arc::new(Mutex::new(vec![]));
710        let recv2 = Arc::clone(&received);
711        let hook: PostBuildFn = Box::new(move |gen_dir: PathBuf| {
712            let inner = Arc::clone(&recv2);
713            Box::pin(async move {
714                inner.lock().unwrap().push(gen_dir);
715                Ok(())
716            })
717        });
718
719        let watcher = Watcher::new(workload.path(), pointer.clone(), opts)
720            .with_post_build(hook);
721
722        watcher.rebuild().await.unwrap();
723
724        let calls = received.lock().unwrap();
725        assert_eq!(calls.len(), 1, "hook should be called once per rebuild");
726        // The hook receives the gen-N directory (parent of html/).
727        assert!(calls[0].ends_with("gen-0"), "hook receives gen dir, got {:?}", calls[0]);
728        // pointer points at gen-0/html.
729        assert_eq!(pointer.current(), calls[0].join("html"));
730    }
731
732    #[tokio::test]
733    async fn post_build_hook_failure_does_not_fail_rebuild() {
734        let workload = tempdir().unwrap();
735        write(&workload.path().join("src/index.txt"), "<h1>hook-err</h1>");
736
737        let pointer = DistPointer::new(workload.path().join("dist").join("html"));
738        let mut opts = WatchOptions::defaults_for_workload(workload.path());
739        opts.build = BuildDriver::Shell(FAKE_BUILD.to_string());
740
741        let hook: PostBuildFn = Box::new(move |_gen_dir: PathBuf| {
742            Box::pin(async move { anyhow::bail!("publish failed (test)") })
743        });
744
745        let watcher = Watcher::new(workload.path(), pointer.clone(), opts)
746            .with_post_build(hook);
747
748        // rebuild() should succeed even though the hook fails.
749        let served = watcher.rebuild().await.unwrap();
750        // Pointer still flipped.
751        assert_eq!(pointer.current(), served);
752    }
753
754    #[tokio::test]
755    async fn rebuild_fails_when_html_dir_missing() {
756        let workload = tempdir().unwrap();
757        write(&workload.path().join("src/index.txt"), "x");
758
759        let pointer = DistPointer::new(workload.path().join("dist").join("html"));
760        let mut opts = WatchOptions::defaults_for_workload(workload.path());
761        // Build "succeeds" but writes the wrong place (dist/, not dist/html/).
762        opts.build = BuildDriver::Shell("mkdir -p dist && echo x > dist/marker".to_string());
763        let watcher = Watcher::new(workload.path(), pointer.clone(), opts);
764
765        let err = watcher.rebuild().await.unwrap_err();
766        assert!(err.to_string().contains("html/"));
767        // Pointer unchanged because rebuild bailed before set().
768        assert_eq!(pointer.current(), workload.path().join("dist").join("html"));
769    }
770
771    #[tokio::test]
772    async fn rebuild_fails_when_build_exits_nonzero() {
773        let workload = tempdir().unwrap();
774        write(&workload.path().join("src/index.txt"), "x");
775
776        let pointer = DistPointer::new(workload.path().join("dist").join("html"));
777        let mut opts = WatchOptions::defaults_for_workload(workload.path());
778        opts.build = BuildDriver::Shell("exit 7".to_string());
779        let watcher = Watcher::new(workload.path(), pointer.clone(), opts);
780
781        let err = watcher.rebuild().await.unwrap_err();
782        assert!(err.to_string().contains("exited"));
783    }
784
785    #[tokio::test]
786    async fn watcher_run_rebuilds_on_file_change() {
787        let workload = tempdir().unwrap();
788        write(&workload.path().join("src/index.txt"), "<h1>A</h1>");
789
790        let pointer = DistPointer::new(workload.path().join("dist").join("html"));
791        let mut opts = WatchOptions::defaults_for_workload(workload.path());
792        opts.build = BuildDriver::Shell(FAKE_BUILD.to_string());
793        opts.debounce = Duration::from_millis(50);
794        let watcher = Watcher::new(workload.path(), pointer.clone(), opts);
795
796        let task = tokio::spawn(async move { watcher.run().await });
797
798        // Wait for the initial build to flip the pointer off the original
799        // dist/html path.
800        let initial_path = workload.path().join("dist").join("html");
801        for _ in 0..40 {
802            if pointer.current() != initial_path {
803                break;
804            }
805            tokio::time::sleep(Duration::from_millis(50)).await;
806        }
807        let first = pointer.current();
808        assert_ne!(first, initial_path, "initial build did not swap pointer");
809        let body = std::fs::read_to_string(first.join("index.html")).unwrap();
810        assert!(body.contains("<h1>A</h1>"));
811
812        // Edit the source. Add a tiny sleep to let notify register the
813        // watch before the modification.
814        tokio::time::sleep(Duration::from_millis(100)).await;
815        write(&workload.path().join("src/index.txt"), "<h1>B</h1>");
816
817        // Wait for pointer to advance again.
818        for _ in 0..60 {
819            if pointer.current() != first {
820                break;
821            }
822            tokio::time::sleep(Duration::from_millis(100)).await;
823        }
824        let second = pointer.current();
825        assert_ne!(second, first, "file edit did not trigger rebuild");
826        let body = std::fs::read_to_string(second.join("index.html")).unwrap();
827        assert!(body.contains("<h1>B</h1>"));
828
829        task.abort();
830    }
831}