1use 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
70pub type PostBuildFuture =
74 std::pin::Pin<Box<dyn std::future::Future<Output = anyhow::Result<()>> + Send + 'static>>;
75
76pub 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#[derive(Debug, Clone)]
107pub enum BuildDriver {
108 InProcess,
115 Shell(String),
119}
120
121impl BuildDriver {
122 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
139pub const LEGACY_SHELL_BUILD: &str = "bun run build";
142
143#[derive(Debug, Clone)]
145pub struct WatchOptions {
146 pub watch_dir: PathBuf,
148 pub build: BuildDriver,
150 pub build_out_dir: PathBuf,
153 pub state_dir: PathBuf,
155 pub debounce: Duration,
157 pub initial_build: bool,
159 pub build_env: Vec<(String, String)>,
164}
165
166impl WatchOptions {
167 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
218pub 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 pub fn with_post_build(mut self, f: PostBuildFn) -> Self {
250 self.post_build = Some(f);
251 self
252 }
253
254 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 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 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 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 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 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 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 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 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 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
470pub 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
486pub 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 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 #[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 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 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 assert!(calls[0].ends_with("gen-0"), "hook receives gen dir, got {:?}", calls[0]);
728 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 let served = watcher.rebuild().await.unwrap();
750 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 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 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 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 tokio::time::sleep(Duration::from_millis(100)).await;
815 write(&workload.path().join("src/index.txt"), "<h1>B</h1>");
816
817 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}