diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 4338d658..a3e50f41 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -227,6 +227,8 @@ jobs: run: node e2e/scan-starvation.mjs - name: thousands of plugin fs.watch handles register and deliver under a write storm run: node e2e/fs-watch-many-watchers.mjs + - name: top-level files and dirs created after boot are watched (new shared/ tree shape) + run: node e2e/watch-new-root-dir.mjs - name: mid-boot SIGTERM kills plugin children; next boot reaps a SIGKILLed server's orphans run: node e2e/orphan-reap.mjs - name: worker build forms (?worker, inline, url, new URL) and worker-only chunks skip document helpers @@ -256,6 +258,8 @@ jobs: run: node e2e/ssr-concurrent.mjs - name: start dev holds the page reload behind the editor HMR gate until the flush run: node e2e/start-hmr-gate.mjs + - name: start dev skips rebundle and reload for changes outside every served graph + run: node e2e/start-rebundle-gate.mjs # i need to fix this or automate build: diff --git a/crates/oj/src/start_dev.rs b/crates/oj/src/start_dev.rs index 83692d68..7e81493a 100644 --- a/crates/oj/src/start_dev.rs +++ b/crates/oj/src/start_dev.rs @@ -67,6 +67,19 @@ struct StartState { /// stale fallback. In an Arc so the handler closure captures only the flag, /// not the whole state (PluginServe is state-held: that would cycle). runner_dirty: Arc, + /// The last client bundle's complete input set (closure.json: every module + /// id the bundler touched, plus watch files and config deps). A batch that + /// misses it entirely cannot change the bundle, so the rebundle is skipped + /// (rollup watch semantics; Vite's module graph plays this role). Empty + /// means unknown, which always rebundles. + client_closure: std::sync::RwLock>, + /// False after a failed rebundle: the closure may be a torso (the build + /// stopped early), so every next batch rebundles until one succeeds. + bundle_ok: std::sync::atomic::AtomicBool, + /// The dev server's SSR bridge, here for its module-graph lookup: files + /// the generic pipeline serves (a dev worker script fetched by URL) are + /// in no other graph the reload decision consults. + ssr: oj_server::SsrBridge, } impl StartState { @@ -160,8 +173,8 @@ pub async fn start_dev( // plugin children (workerd) spawn well before the server listens. let host_slot = std::sync::Arc::new(std::sync::OnceLock::new()); tokio::spawn(oj_server::close_plugins_on_shutdown(host_slot.clone())); - let built_task = tokio::spawn( - oj_server::DevServer { + let built_task = tokio::spawn({ + let dev_server = oj_server::DevServer { root: root.clone(), port, host, @@ -170,9 +183,16 @@ pub async fn start_dev( no_cache: false, lazy: false, mode: Some(mode.clone()), + }; + async move { + let built = dev_server.build_app().await?; + // Subscribed the moment the server (and its watcher) exists, not + // after the client bundle joins: an edit landing during the rest + // of boot queues for the Start watcher instead of being lost. + let watch_rx = built.watch_feed.subscribe(); + anyhow::Ok((built, watch_rx)) } - .build_app(), - ); + }); // The two codegen steps run SERIALIZED on purpose, not as a lost // concurrency opportunity: the route-tree generator writes @@ -209,10 +229,7 @@ pub async fn start_dev( let (reload_tx, _) = broadcast::channel::<()>(16); let (bundle_res, built_res) = tokio::join!(bundle, built_task); let pinned = bundle_res??; - let built = built_res??; - // Subscribed before the rest of boot, so an edit landing while the engine - // comes up queues for the watcher thread instead of being lost. - let watch_rx = built.watch_feed.subscribe(); + let (built, watch_rx) = built_res??; oj_server::boot_phase("bundle+build joined"); // The in-process Start runner: an embedded engine whose module host runs // the dev server's SSR pipeline (StartHost); it needs the built app's @@ -391,6 +408,10 @@ pub async fn start_dev( }) .collect(), ), + client_closure: std::sync::RwLock::new(load_client_closure(&cache)), + // The boot bundle succeeded (start_dev aborts otherwise). + bundle_ok: std::sync::atomic::AtomicBool::new(true), + ssr: built.ssr.clone(), }); // The activation handler: when the plugin middleware comes up after boot, @@ -622,10 +643,15 @@ struct PendingRebundle { /// Every changed path of the merged batches (relevant and not), for the /// regen decisions; the HMR gate recorded them at the watcher event. paths: std::collections::HashSet, - /// The settled worker-invalidate sends already in flight for the merged - /// batches: this run's browser reload waits for them, so a reload never - /// lands while the worker environments still hold stale modules. - invalidates: Vec>, + /// The worker-invalidate sends (early and settled) in flight for the + /// merged batches: this run's browser reload waits for them, so a reload + /// never lands while the worker environments still hold stale modules, + /// and their replies say whether any worker graph actually served a + /// changed module. + invalidates: Vec>, + /// A relevant file was created: resolution can shift onto a new file even + /// when no bundled input changed, so creates always rebundle. + any_created: bool, /// The newest merged batch number: its done-frame clears the editor pill. batch: u64, } @@ -640,7 +666,7 @@ fn spawn_mw_invalidate( state: &StartState, paths: &std::collections::HashSet, created: &std::collections::HashSet, -) -> Option> { +) -> Option> { let port = state.plugin_serve.mw_port()?; let changed: Vec<(String, &'static str)> = paths .iter() @@ -650,9 +676,37 @@ fn spawn_mw_invalidate( if changed.is_empty() { return None; } - Some(rt.spawn(async move { - oj_server::notify_plugin_mw_invalidate(port, &changed).await; - })) + Some(rt.spawn(async move { oj_server::notify_plugin_mw_invalidate(port, &changed).await })) +} + +/// The last successful bundle's input set, written by bundle-client.mjs. +/// Empty on any read or parse failure: unknown always rebundles. +fn load_client_closure(cache: &Path) -> std::collections::HashSet { + std::fs::read_to_string(cache.join("closure.json")) + .ok() + .and_then(|s| serde_json::from_str::>(&s).ok()) + .map(|v| v.into_iter().collect()) + .unwrap_or_default() +} + +/// Whether a batch can change the client bundle: a path in the closure (the +/// watcher spelling, or its real path, since the bundler records resolved +/// paths), a create (resolution can shift onto a new file), a failed previous +/// bundle, or no closure to consult. +fn client_bundle_affected(state: &StartState, paths: &[PathBuf], any_created: bool) -> bool { + if any_created || !state.bundle_ok.load(std::sync::atomic::Ordering::SeqCst) { + return true; + } + let closure = state + .client_closure + .read() + .unwrap_or_else(std::sync::PoisonError::into_inner); + if closure.is_empty() { + return true; + } + paths.iter().any(|p| { + closure.contains(p) || std::fs::canonicalize(p).is_ok_and(|real| closure.contains(&real)) + }) } /// The files the regen steps write, which the watcher deliberately never @@ -770,8 +824,11 @@ fn spawn_start_watcher( // again below; the host dedups by content identity, so the repeat // send only costs work when a write landed inside the settle // window and actually changed content. The OS watcher's own - // delivery latency is the remaining, unclosable window. - let _ = spawn_mw_invalidate(&rt, &state, &paths, &created); + // delivery latency is the remaining, unclosable window. The + // handle is kept: the host's reply says whether a worker graph + // matched, and the dedup means the settled send alone can read + // as a miss for a change this early send already carried. + let early_invalidate = spawn_mw_invalidate(&rt, &state, &paths, &created); loop { match rx.recv_timeout(std::time::Duration::from_millis(50)) { Ok(ev) => { @@ -828,7 +885,9 @@ fn spawn_start_watcher( .unwrap_or_else(std::sync::PoisonError::into_inner); let run = slot.get_or_insert_with(PendingRebundle::default); run.paths.extend(paths.iter().cloned()); + run.invalidates.extend(early_invalidate); run.invalidates.extend(invalidate); + run.any_created |= created.iter().any(|p| watch_relevant(p)); run.batch = batch; } wake.notify_one(); @@ -866,6 +925,13 @@ async fn rebundle_worker( }; let paths: Vec = run.paths.into_iter().collect(); let regen_files = regen_output_files(&root, &cache); + // The gate: a batch whose paths miss the last bundle's input closure + // cannot change the client bundle (rollup watch semantics; Vite only + // acts on files its module graphs know). Regen checks still run below + // either way: a server fn can live in a file the client never + // imports, since gen-resolver scans all of src. + let needs_bundle = client_bundle_affected(&state, &paths, run.any_created); + let no_invalidates_sent = run.invalidates.is_empty(); let client = { let (r, c, m) = (root.clone(), cache.clone(), state.mode.clone()); let env = Arc::clone(&state.script_env); @@ -887,7 +953,10 @@ async fn rebundle_worker( if routes_changed || server_fn_changed { let _ = generate_server_fn_resolver(&r, &c, &env); } - let pinned = if bundle_client_entry(&r, &c, &env).is_err() { + // A route-set change rewrites the generated tree, a bundled + // input, even when the triggering paths missed the closure. + let bundling = needs_bundle || routes_changed; + let pinned = if !bundling || bundle_client_entry(&r, &c, &env).is_err() { None } else { match start_bundle_store(&r, &m).persist(&c) { @@ -895,23 +964,40 @@ async fn rebundle_worker( None => oj_cache::start_bundle::PinnedBundle::from_build_dir(&c), } }; - (routes_now, pinned) + (routes_now, server_fn_changed, bundling, pinned) }) }; let side = async { // A reload onto stale worker modules would re-render old content: - // this run's settled invalidates complete before the signal. + // this run's settled invalidates complete before the signal, and + // their replies say whether any worker graph served a change (a + // failed send reads as yes). + let mut hit = false; for handle in run.invalidates { - let _ = handle.await; + hit |= handle.await.unwrap_or(true); } + hit }; - let (client, _) = tokio::join!(client, side); - if let Ok((routes_now, pinned)) = client { + let (client, worker_hit) = tokio::join!(client, side); + let (mut bundled, mut server_fn_changed) = (false, false); + if let Ok((routes_now, fn_changed, bundling, pinned)) = client { prev_routes = routes_now; - if let Some(pinned) = pinned { - let pinned = Arc::new(pinned); - state.gzip.retain(|h| pinned.has_hash(h)); - *state.bundle.write().unwrap() = pinned; + server_fn_changed = fn_changed; + bundled = bundling; + if bundling { + state + .bundle_ok + .store(pinned.is_some(), std::sync::atomic::Ordering::SeqCst); + if let Some(pinned) = pinned { + let pinned = Arc::new(pinned); + state.gzip.retain(|h| pinned.has_hash(h)); + *state.bundle.write().unwrap() = pinned; + *state + .client_closure + .write() + .unwrap_or_else(std::sync::PoisonError::into_inner) = + load_client_closure(&cache); + } } } // Regenerated outputs are edits the watcher never forwards: push the @@ -942,20 +1028,45 @@ async fn rebundle_worker( " oj start: regen outputs changed ({}), invalidating worker", names.join(", ") ); - oj_server::notify_plugin_mw_invalidate(port, ®en_changed).await; + let _ = oj_server::notify_plugin_mw_invalidate(port, ®en_changed).await; } } // The non-lazy runner reloads strictly after the regen completed (its // re-import must see the new generated files, as the inline loop - // guaranteed) and before the browser reload. - if !state.lazy_runner() { - state.engine.reload().await; - } + // guaranteed) and before the browser reload. A batch that bundled or + // regenerated reloads it whole, as before; one that did neither only + // revalidates, which drops exactly the stale records (and answers + // whether the engine served any of the changed files at all). + let engine_hit = if !state.lazy_runner() { + if bundled || server_fn_changed || !regen_changed.is_empty() { + state.engine.reload().await; + true + } else { + state.engine.revalidate().await + } + } else if no_invalidates_sent { + // Lazy with the middleware down: the fallback runner is serving, + // so its graph is the only reload signal left. + state.engine.revalidate().await + } else { + false + }; + // Reload the browser only when something it can observe changed: the + // bundle, a regen output, a worker-served module, an engine-served + // one, or a module the generic pipeline served (a dev worker script). + // A change outside every graph (a README, a config of some other + // tool) reloads nothing, as in Vite. + let reload = bundled + || server_fn_changed + || !regen_changed.is_empty() + || worker_hit + || engine_hit + || paths.iter().any(|p| state.ssr.graph_knows_file(p)); // The hold was taken at the watcher event; by now the editor's flush // may have consumed it, in which case this run's reload goes out with // the fresh bundle instead of waiting for a flush that already came. let held = state.gate.as_ref().is_some_and(|g| g.reload_is_held()); - if !held { + if reload && !held { let _ = state.reload_tx.send(()); } let modules = client_module_count(&cache); @@ -967,10 +1078,18 @@ async fn rebundle_worker( Some(0), true, )); - if held { - println!(" oj start: rebuilt, reload held for the editor's flush"); + if !reload { + println!(" oj start: change outside every served graph, no rebundle, no reload"); + } else if held { + println!( + " oj start: {}, reload held for the editor's flush", + if bundled { "rebuilt" } else { "updated" } + ); } else { - println!(" oj start: rebuilt, reloading"); + println!( + " oj start: {}, reloading", + if bundled { "rebuilt" } else { "updated" } + ); } } } @@ -2709,6 +2828,29 @@ mod tests { ); } + // The rebundle gate's input set: closure.json parses to the file set, and + // anything unreadable reads as empty, which the gate treats as unknown + // (always rebundle) — a torn or missing closure must never suppress one. + #[test] + fn client_closure_parses_and_defaults_to_empty() { + let dir = tmp("closure-load"); + assert!( + load_client_closure(&dir).is_empty(), + "missing file is empty" + ); + std::fs::write(dir.join("closure.json"), "{ not json").unwrap(); + assert!(load_client_closure(&dir).is_empty(), "torn file is empty"); + std::fs::write( + dir.join("closure.json"), + r#"["/app/src/a.tsx","/app/styles/app.css"]"#, + ) + .unwrap(); + let set = load_client_closure(&dir); + assert!(set.contains(Path::new("/app/src/a.tsx"))); + assert!(set.contains(Path::new("/app/styles/app.css"))); + assert_eq!(set.len(), 2); + } + // The rebundle worker tracks the regen outputs' content against what the // worker environments last saw and pushes only real moves: rewritten // content and a newly created file count, an untouched file (or one diff --git a/crates/oj_server/src/assets/plugin-host.mjs b/crates/oj_server/src/assets/plugin-host.mjs index 2cd8da51..eb6216ae 100644 --- a/crates/oj_server/src/assets/plugin-host.mjs +++ b/crates/oj_server/src/assets/plugin-host.mjs @@ -2229,6 +2229,7 @@ async function invalidateEnvironments(environments, watcher, changes) { unmatchedChangeLogged.add(file); process.stderr.write(`${OJ} plugin host: change to ${file} matched no module in any runner environment\n`); } + return matched.size; } // Vite's prepareError: the payload hot.send({type:"error"}) carries. @@ -3219,13 +3220,18 @@ async function ensureConfigureServerMiddleware() { // Answer only once the invalidation is done, so the Rust side's POST // completing means a next request cannot be served from stale modules. // Serialized: the early (pre-settle) and settled sends for one edit - // batch must not interleave their module-graph walks. - invalidateQueue = invalidateQueue + // batch must not interleave their module-graph walks. The reply says + // how many changes matched a runner-backed graph (null: the walk + // threw, the caller must assume a hit): the Rust side reloads the + // browser only when something actually served changed. + const run = invalidateQueue .then(() => invalidateEnvironments(server.environments, fileWatcher, changes)) - .catch(() => {}); - await invalidateQueue; - res.statusCode = 204; - res.end(); + .catch(() => null); + invalidateQueue = run; + const matchedCount = await run; + res.statusCode = 200; + res.setHeader("content-type", "application/json"); + res.end(JSON.stringify({ matched: matchedCount === null ? null : matchedCount || 0 })); }); return; } diff --git a/crates/oj_server/src/assets/start/bundle-client.mjs b/crates/oj_server/src/assets/start/bundle-client.mjs index 5b537ac8..a1df4bad 100644 --- a/crates/oj_server/src/assets/start/bundle-client.mjs +++ b/crates/oj_server/src/assets/start/bundle-client.mjs @@ -117,6 +117,15 @@ async function main() { visitedIds.add(id); return null; }, + // The complete input set, from the bundler itself: getModuleIds() covers + // every module the build touched, including the kinds the transform hook + // and the output walk miss (css imports among them). closure.json must + // not under-report: the dev server skips rebundles for paths outside it. + buildEnd() { + try { + for (const id of this.getModuleIds()) visitedIds.add(id); + } catch {} + }, }; const serverFnClient = { @@ -210,7 +219,15 @@ async function main() { for (const id of Object.keys(o.modules ?? {})) closure.add(id); } const closureFiles = [...closure] - .map((id) => String(id).split("?")[0]) + // An asset module's id is virtual (\0oj-css:/abs/x.css and friends), but + // its FILE is an input: raw/inline embed its bytes, css lands in + // css-urls.json. The closure must carry the file, or an edit to it never + // rebundles. + .map((id) => { + const s = String(id); + return s.startsWith("\0oj-") ? s.slice(s.indexOf(":") + 1) : s; + }) + .map((id) => id.split("?")[0]) .filter((p) => !p.startsWith("\0")) .map((p) => (isAbsolute(p) ? p : resolve(APP, p))) .filter((p) => !p.startsWith(HERE) && !/(^|\/)routeTree\.gen\.[jt]sx?$/.test(p)) diff --git a/crates/oj_server/src/plugin_mw.rs b/crates/oj_server/src/plugin_mw.rs index 3fe2f070..f972a176 100644 --- a/crates/oj_server/src/plugin_mw.rs +++ b/crates/oj_server/src/plugin_mw.rs @@ -49,20 +49,34 @@ pub(crate) fn stream_reqwest_response(resp: reqwest::Response) -> Response { } // Tell the plugin middleware server that files changed so it can invalidate -// module graphs and send HMR; type is "update" | "create" | "delete". Fire-and-forget. -pub async fn notify_plugin_mw_invalidate(port: u16, changes: &[(String, &'static str)]) { +// module graphs and send HMR; type is "update" | "create" | "delete". +// Returns whether a runner-backed graph matched one of the changes; a host +// that cannot tell (an error reply, an unreachable host, a null count) reads +// as true, so the caller's reload decision only ever errs toward reloading. +pub async fn notify_plugin_mw_invalidate(port: u16, changes: &[(String, &'static str)]) -> bool { let client = plugin_mw_client(); let changes: Vec = changes .iter() .map(|(path, kind)| serde_json::json!({ "path": path, "type": kind })) .collect(); let body = serde_json::json!({ "changes": changes }).to_string(); - let _ = client + let resp = client .post(format!("http://127.0.0.1:{port}/__oj_invalidate")) .header(header::CONTENT_TYPE, "application/json") .body(body) .send() .await; + let Ok(resp) = resp else { return true }; + let Ok(body) = resp.bytes().await else { + return true; + }; + let Ok(parsed) = serde_json::from_slice::(&body) else { + return true; + }; + match parsed.get("matched") { + Some(serde_json::Value::Number(n)) => n.as_u64().map(|n| n > 0).unwrap_or(true), + _ => true, + } } pub(crate) fn plugin_mw_client() -> &'static reqwest::Client { diff --git a/crates/oj_server/src/ssr.rs b/crates/oj_server/src/ssr.rs index 8a932a5f..b7baa131 100644 --- a/crates/oj_server/src/ssr.rs +++ b/crates/oj_server/src/ssr.rs @@ -40,6 +40,14 @@ impl SsrBridge { self.state.engine_registry.clone() } + /// Whether the generic pipeline ever served this file as a module (its + /// module graph; a Start worker script fetched by URL lands here while + /// staying outside both the client bundle and the Start engine's graph). + pub fn graph_knows_file(&self, file: &Path) -> bool { + let url = crate::rewrite::url_of(&self.state.root, file); + self.state.graph.lock().unwrap().contains(Path::new(&url)) + } + pub async fn resolve(&self, importer: &str, spec: &str) -> Result { ssr_resolve_inner(&self.state, importer, spec).await } diff --git a/crates/oj_server/src/tests.rs b/crates/oj_server/src/tests.rs index 1a0289d3..faae7679 100644 --- a/crates/oj_server/src/tests.rs +++ b/crates/oj_server/src/tests.rs @@ -1118,8 +1118,9 @@ fn watch_ignored_globs_match_relative_and_absolute_paths() { } #[test] -fn watch_feed_drops_ignored_paths_and_dropped_subscribers() { +fn watch_feed_drops_excluded_paths_and_dropped_subscribers() { let root = Path::new("/app"); + let cache = Path::new("/app/relocated-cache"); let ignored = watch_ignored_patterns(root, &["**/generated/**".to_string()]); let feed = WatchFeed::default(); let ev = |paths: &[&str]| { @@ -1132,17 +1133,35 @@ fn watch_feed_drops_ignored_paths_and_dropped_subscribers() { drop(dropped); feed.publish( - &ev(&["/app/shared/a.ts", "/app/src/generated/b.ts"]), + &ev(&[ + "/app/shared/a.ts", + "/app/src/generated/b.ts", + "/app/relocated-cache/v1/start/chunk.js", + ]), &ignored, root, + cache, + ); + feed.publish(&ev(&["/app/src/generated/c.ts"]), &ignored, root, cache); + feed.publish( + &ev(&["/app/relocated-cache/v1/mod.json"]), + &ignored, + root, + cache, + ); + + feed.publish( + &ev(&["/app/vite.config.ts.timestamp-1791586920568-4c4081376eef.mjs"]), + &ignored, + root, + cache, ); - feed.publish(&ev(&["/app/src/generated/c.ts"]), &ignored, root); - let got = kept.try_recv().expect("the unignored path is delivered"); + let got = kept.try_recv().expect("the unexcluded path is delivered"); assert_eq!(got.paths, vec![PathBuf::from("/app/shared/a.ts")]); assert!( kept.try_recv().is_err(), - "an all-ignored event is not delivered" + "all-ignored, cache-dir and config-temp events are not delivered" ); assert_eq!( feed.0.lock().unwrap().len(), diff --git a/crates/oj_server/src/watch.rs b/crates/oj_server/src/watch.rs index c90098bf..59e289b1 100644 --- a/crates/oj_server/src/watch.rs +++ b/crates/oj_server/src/watch.rs @@ -414,21 +414,28 @@ pub(crate) enum WatchMsg { pub struct WatchFeed(pub(crate) Arc>>>); impl WatchFeed { - /// Every event from now on, minus `server.watch.ignored` paths; dropping - /// the receiver unsubscribes. + /// Every event from now on, minus excluded paths (`server.watch.ignored` + /// and the cache dir); dropping the receiver unsubscribes. pub fn subscribe(&self) -> std::sync::mpsc::Receiver { let (tx, rx) = std::sync::mpsc::channel(); self.0.lock().unwrap().push(tx); rx } - pub(crate) fn publish(&self, ev: ¬ify::Event, ignored: &[glob::Pattern], root: &Path) { + pub(crate) fn publish( + &self, + ev: ¬ify::Event, + ignored: &[glob::Pattern], + root: &Path, + cache_base: &Path, + ) { let mut subscribers = self.0.lock().unwrap(); if subscribers.is_empty() { return; } let mut ev = ev.clone(); - ev.paths.retain(|p| !is_watch_ignored(ignored, root, p)); + ev.paths + .retain(|p| !watch_excluded(ignored, root, cache_base, p)); if ev.paths.is_empty() { return; } @@ -436,6 +443,39 @@ impl WatchFeed { } } +/// What every consumer of watcher events skips: `server.watch.ignored` plus +/// the resolved cache dir, which Vite puts in chokidar's `ignored` -- the +/// name-based skip in `is_unwatched_dir` misses an `OJ_CACHE_DIR` inside the +/// root, whose writes would echo every compile back into the watcher -- plus +/// Vite's own config-loader temp. +pub(crate) fn watch_excluded( + ignored: &[glob::Pattern], + root: &Path, + cache_base: &Path, + path: &Path, +) -> bool { + path.starts_with(cache_base) + || path + .file_name() + .and_then(|n| n.to_str()) + .is_some_and(is_config_timestamp_temp) + || is_watch_ignored(ignored, root, path) +} + +/// Vite's `loadConfigFromBundledFile` temp (`.timestamp--.mjs`). +/// Vite prefers `node_modules/.vite-temp/` for it, but skips that dir under +/// Deno -- oj's runtime -- and writes beside the config instead, so every +/// config load drops a short-lived file at the root the watch would see; +/// under Start, whose rebundle loads the config, that was a rebuild loop. +fn is_config_timestamp_temp(name: &str) -> bool { + name.ends_with(".mjs") + && name.find(".timestamp-").is_some_and(|i| { + name.as_bytes() + .get(i + ".timestamp-".len()) + .is_some_and(u8::is_ascii_digit) + }) +} + /// Vite's ensureWatchedFile: a served file OUTSIDE the root is not covered by /// the root watch, so its directory goes to the watcher thread (the directory, /// not the file: editors save by rename-replace, which strands an inode watch). @@ -471,21 +511,36 @@ fn is_unwatched_dir(name: &std::ffi::OsStr) -> bool { ) } -/// Watch each top-level root entry except the unwatched dirs; fall back to a -/// recursive root watch only if nothing else could be watched. -fn watch_root(watcher: &mut notify::RecommendedWatcher, root: &Path) -> notify::Result<()> { +/// Watch each top-level root entry except the unwatched dirs, plus the root +/// itself non-recursively: new top-level entries are born under it, and the +/// per-entry watches only cover what existed at boot (chokidar watches the +/// root, so Vite sees new children). Falls back to a recursive root watch +/// only if nothing else could be watched. +fn watch_root( + watcher: &mut notify::RecommendedWatcher, + root: &Path, + cache_base: &Path, +) -> notify::Result<()> { use notify::{RecursiveMode, Watcher}; - let mut watched_any = false; + let root_watched = watcher.watch(root, RecursiveMode::NonRecursive).is_ok(); + let mut watched_any = root_watched; if let Ok(entries) = std::fs::read_dir(root) { for entry in entries.flatten() { - if is_unwatched_dir(&entry.file_name()) { + if is_unwatched_dir(&entry.file_name()) || entry.path() == cache_base { continue; } let path = entry.path(); let mode = if path.is_dir() { RecursiveMode::Recursive - } else { + } else if !root_watched || entry.file_type().is_ok_and(|t| t.is_symlink()) { + // A plain top-level file is already covered by the root watch; + // a second watch would deliver each edit twice (ContentChanges + // only drops metadata noise). A symlinked entry still needs its + // own watch: it follows to the target inode, whose edits never + // touch the root directory entry. RecursiveMode::NonRecursive + } else { + continue; }; if watcher.watch(&path, mode).is_ok() { watched_any = true; @@ -499,6 +554,61 @@ fn watch_root(watcher: &mut notify::RecommendedWatcher, root: &Path) -> notify:: } } +/// A directory created at the top level after boot, caught by the root's +/// non-recursive watch: watch it recursively and report the files it already +/// holds as creates (chokidar emits `add` per existing file when a new dir +/// appears), since writes landing between the mkdir and this watch were +/// never seen. Returns those files. A dir RENAMED into the root (a staged +/// tree moved into place) arrives as `Modify(Name)`, not `Create` -- chokidar +/// raises `addDir` for both -- so renames are adopted too; rename-From paths +/// and plain file renames fall through the `is_dir` check, and re-watching an +/// already-watched dir is idempotent (notify keys watches by path). +fn adopt_created_root_dirs( + watcher: &mut notify::RecommendedWatcher, + root: &Path, + cache_base: &Path, + ev: ¬ify::Event, +) -> Vec { + use notify::{RecursiveMode, Watcher}; + if !matches!( + ev.kind, + notify::EventKind::Create(_) + | notify::EventKind::Modify(notify::event::ModifyKind::Name(_)) + ) { + return Vec::new(); + } + let mut found = Vec::new(); + for path in &ev.paths { + if path.parent() != Some(root) + || path == cache_base + || path.file_name().is_none_or(is_unwatched_dir) + || !path.is_dir() + { + continue; + } + if watcher.watch(path, RecursiveMode::Recursive).is_ok() { + scan_files(path, &mut found); + } + } + found +} + +fn scan_files(dir: &Path, out: &mut Vec) { + let Ok(entries) = std::fs::read_dir(dir) else { + return; + }; + for entry in entries.flatten() { + let path = entry.path(); + if path.is_dir() { + if !is_unwatched_dir(&entry.file_name()) { + scan_files(&path, out); + } + } else { + out.push(path); + } + } +} + /// One debounced set of changes. #[derive(Default)] struct Batch { @@ -567,11 +677,21 @@ pub(crate) fn spawn_watcher(state: Arc, rx: std::sync::mpsc::Receiv let tx = state.watch_tx.clone(); let feed = state.watch_feed.clone(); - let (ignored, root) = (state.watch_ignored.clone(), state.root.clone()); + // Canonicalized so a relative or symlinked `OJ_CACHE_DIR` still + // matches the absolute paths the watcher reports. + let cache_base = { + let base = oj_cache::cache_base(&state.root); + std::fs::canonicalize(&base).unwrap_or(base) + }; + let (ignored, root, cache) = ( + state.watch_ignored.clone(), + state.root.clone(), + cache_base.clone(), + ); let mut watcher = match notify::recommended_watcher(move |ev: notify::Result| { if let Ok(ev) = &ev { - feed.publish(ev, &ignored, &root); + feed.publish(ev, &ignored, &root, &cache); } let _ = tx.send(WatchMsg::Fs(ev)); }) { @@ -582,7 +702,7 @@ pub(crate) fn spawn_watcher(state: Arc, rx: std::sync::mpsc::Receiv } }; let mut served_dirs: std::collections::HashSet = Default::default(); - if let Err(err) = watch_root(&mut watcher, &state.root) { + if let Err(err) = watch_root(&mut watcher, &state.root, &cache_base) { eprintln!("oj: cannot watch {}: {err}", state.root.display()); return; } @@ -603,13 +723,27 @@ pub(crate) fn spawn_watcher(state: Arc, rx: std::sync::mpsc::Receiv Err(_) => break, }; let mut batch = Batch::default(); - batch.add(&mut changes, &first); + ingest_event( + &state, + &mut watcher, + &cache_base, + &mut changes, + &mut batch, + &first, + ); if batch.paths.is_empty() { continue; } loop { match rx.recv_timeout(Duration::from_millis(debounce_ms)) { - Ok(WatchMsg::Fs(Ok(ev))) => batch.add(&mut changes, &ev), + Ok(WatchMsg::Fs(Ok(ev))) => ingest_event( + &state, + &mut watcher, + &cache_base, + &mut changes, + &mut batch, + &ev, + ), Ok(WatchMsg::Fs(Err(_))) => {} Ok(WatchMsg::Dir(dir)) => watch_served_dir(&mut watcher, &mut served_dirs, dir), Err(RecvTimeoutError::Timeout) => break, @@ -619,7 +753,7 @@ pub(crate) fn spawn_watcher(state: Arc, rx: std::sync::mpsc::Receiv let Batch { paths, mut created } = batch; let paths: Vec = paths .into_iter() - .filter(|p| !is_watch_ignored(&state.watch_ignored, &state.root, p)) + .filter(|p| !watch_excluded(&state.watch_ignored, &state.root, &cache_base, p)) .collect(); if paths.is_empty() { continue; @@ -630,3 +764,27 @@ pub(crate) fn spawn_watcher(state: Arc, rx: std::sync::mpsc::Receiv } }); } + +/// One raw watcher event into the current batch: a created top-level dir is +/// adopted first, and its existing files enter this same batch (and the feed, +/// which only hears the notify callback) as a synthesized create. +fn ingest_event( + state: &ServerState, + watcher: &mut notify::RecommendedWatcher, + cache_base: &Path, + changes: &mut ContentChanges, + batch: &mut Batch, + ev: ¬ify::Event, +) { + let adopted = adopt_created_root_dirs(watcher, &state.root, cache_base, ev); + if !adopted.is_empty() { + let mut synth = + notify::Event::new(notify::EventKind::Create(notify::event::CreateKind::Any)); + synth.paths = adopted; + state + .watch_feed + .publish(&synth, &state.watch_ignored, &state.root, cache_base); + batch.add(changes, &synth); + } + batch.add(changes, ev); +} diff --git a/e2e/start-rebundle-gate.mjs b/e2e/start-rebundle-gate.mjs new file mode 100644 index 00000000..b5643cd9 --- /dev/null +++ b/e2e/start-rebundle-gate.mjs @@ -0,0 +1,106 @@ +// SPDX-License-Identifier: MIT +// Copyright (c) 2026 Raphael Amorim + +// The Start rebundle gate: a change whose paths miss every served graph (the +// client bundle's input closure, the engine's and the workers' module +// graphs) must not rebundle the client or reload the browser -- Vite does +// nothing for a file its graphs never served. A change to a bundled module +// must still rebundle and serve fresh. Run with a built target/debug/oj (or +// OJ_BIN). +import { spawn, execSync } from "node:child_process"; +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; +import { fileURLToPath } from "node:url"; +import { settles, waitUp } from "./util.mjs"; + +const here = path.dirname(fileURLToPath(import.meta.url)); +const repo = path.join(here, ".."); +const fixture = path.join(here, "fixtures", "start-app"); +const oj = process.env.OJ_BIN ?? path.join(repo, "target", "debug", "oj"); +const PORT = 6861; +const sleep = (ms) => new Promise((r) => setTimeout(r, ms)); + +const installed = + fs.existsSync(path.join(fixture, "node_modules", "@tanstack", "react-start")) && + fs.existsSync(path.join(fixture, "node_modules", "rolldown")); +if (!installed) { + console.log("SKIP start rebundle gate: fixture deps not installed"); + console.log(" enable with: (cd e2e/fixtures/start-app && npm install)"); + process.exit(0); +} + +if (!process.env.OJ_BIN) execSync("cargo build -p oj", { cwd: repo, stdio: "inherit" }); + +const tmp = fs.mkdtempSync(path.join(os.tmpdir(), "oj-start-gate-")); +const app = path.join(tmp, "app"); +fs.mkdirSync(app); +for (const f of ["src", "public", "styles", "packages", "tsconfig.json", "package.json", "vite.config.ts"]) { + fs.cpSync(path.join(fixture, f), path.join(app, f), { recursive: true }); +} +fs.symlinkSync(path.join(fixture, "node_modules"), path.join(app, "node_modules"), "dir"); +// Present before boot, so its edit below is an update of a known file, not a +// create (creates rebundle by design: resolution can shift onto new files). +fs.writeFileSync(path.join(app, "README.md"), "# fixture\n"); + +let log = ""; +const srv = spawn(oj, ["dev", app, "--port", String(PORT)], { stdio: ["ignore", "pipe", "pipe"] }); +srv.stdout.on("data", (c) => (log += c)); +srv.stderr.on("data", (c) => (log += c)); + +const get = async (route) => { + const res = await fetch(`http://localhost:${PORT}${route}`); + return { status: res.status, body: await res.text() }; +}; + +let failed = false; +try { + await waitUp(`http://localhost:${PORT}/`, { proc: srv }); + const about = await get("/about"); + if (about.status !== 200 || !about.body.includes("about-page-marker")) { + throw new Error(`/about did not render (${about.status})`); + } + + // A change outside every served graph: no rebundle, no reload. + const bundledBefore = log.split("client bundled").length; + fs.appendFileSync(path.join(app, "README.md"), "\nan edit outside every graph\n"); + await settles(async () => log.includes("change outside every served graph"), { timeoutMs: 10000 }); + await sleep(1500); + const bundledAfter = log.split("client bundled").length; + if (bundledAfter !== bundledBefore) { + throw new Error(`README edit rebundled the client (${bundledAfter - bundledBefore} run(s)); log tail:\n${log.slice(-3000)}`); + } + console.log("gate: README edit skipped the rebundle and the reload"); + + // A bundled module still rebundles and reloads. (The served document's + // freshness is the engine's own path and not asserted here: the base + // behavior predates the gate and is unchanged by it.) + const aboutFile = path.join(app, "src", "routes", "about.tsx"); + fs.writeFileSync(aboutFile, fs.readFileSync(aboutFile, "utf8").replace("about-page-marker", "about-page-edited")); + await settles(async () => log.split("client bundled").length > bundledAfter, { timeoutMs: 20000 }); + if (log.split("client bundled").length === bundledAfter) { + throw new Error(`route edit did not rebundle; log tail:\n${log.slice(-3000)}`); + } + await settles(async () => log.includes("rebuilt, reloading"), { timeoutMs: 10000 }); + if (!log.includes("rebuilt, reloading")) { + throw new Error(`route edit did not reload; log tail:\n${log.slice(-3000)}`); + } + console.log("gate: bundled-module edit rebundled and reloaded"); + console.log("START-REBUNDLE-GATE E2E PASSED"); +} catch (err) { + failed = true; + console.error("START-REBUNDLE-GATE E2E FAILED:", err.message); +} finally { + srv.kill("SIGKILL"); + await sleep(300); + for (let i = 0; ; i++) { + try { + fs.rmSync(tmp, { recursive: true, force: true }); + break; + } catch (e) { + if (i >= 20) break; + Atomics.wait(new Int32Array(new SharedArrayBuffer(4)), 0, 0, 100); + } + } +} +process.exit(failed ? 1 : 0); diff --git a/e2e/watch-new-root-dir.mjs b/e2e/watch-new-root-dir.mjs new file mode 100644 index 00000000..d570c77c --- /dev/null +++ b/e2e/watch-new-root-dir.mjs @@ -0,0 +1,102 @@ +// SPDX-License-Identifier: MIT +// Copyright (c) 2026 Raphael Amorim + +// Top-level entries created AFTER boot must be watched: the per-entry root +// watches only cover what read_dir saw at startup, so a new root-level file +// or a freshly mkdir'd top-level directory (an agent adding a shared/ tree to +// a running server) served fine but never produced watcher events -- edits +// there kept the old module until restart. Vite's chokidar watches the root +// itself and picks up new children; this drives both shapes end to end. +import { spawn, execSync } from "node:child_process"; +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; +import assert from "node:assert/strict"; +import { createRequire } from "node:module"; +import { fileURLToPath } from "node:url"; +import { waitUp } from "./util.mjs"; + +const here = path.dirname(fileURLToPath(import.meta.url)); +const repo = path.join(here, ".."); +const oj = process.env.OJ_BIN ?? path.join(repo, "target", "debug", "oj"); +const { chromium } = createRequire(path.join(here, "x.js"))("playwright"); +const sleep = (ms) => new Promise((r) => setTimeout(r, ms)); +const PORT = 5499; + +if (!process.env.OJ_BIN) execSync("cargo build -p oj", { cwd: repo, stdio: "inherit" }); + +const app = fs.mkdtempSync(path.join(os.tmpdir(), "oj-watch-newdir-")); +fs.mkdirSync(path.join(app, "src"), { recursive: true }); +fs.writeFileSync(path.join(app, "package.json"), JSON.stringify({ name: "watch-newdir", version: "1.0.0" })); +fs.writeFileSync(path.join(app, "src", "main.js"), `document.title = "v1"; window.__READY = true;\n`); +fs.writeFileSync( + path.join(app, "index.html"), + `t`, +); + +let failed = false; +const srv = spawn(oj, ["dev", app, "--port", String(PORT)], { stdio: "ignore" }); +let browser; +try { + await waitUp(`http://localhost:${PORT}/`, { proc: srv }); + browser = await chromium.launch(); + const page = await browser.newPage(); + const errors = []; + page.on("pageerror", (e) => errors.push(String(e))); + await page.goto(`http://localhost:${PORT}/`, { timeout: 30000 }); + await page.waitForFunction(() => window.__READY === true, { timeout: 10000 }); + + // A root-level file created after boot: the import lands via the (always + // watched) src/main.js edit; the edit AFTER that only fires if the root + // watch reports its new children. + fs.writeFileSync(path.join(app, "banner.js"), `export const banner = "B1";\n`); + fs.writeFileSync( + path.join(app, "src", "main.js"), + `import { banner } from "../banner.js";\ndocument.title = banner; window.__READY = true;\n`, + ); + await page.waitForFunction(() => document.title === "B1", { timeout: 10000 }); + fs.writeFileSync(path.join(app, "banner.js"), `export const banner = "B2";\n`); + await page.waitForFunction(() => document.title === "B2", { timeout: 10000 }); + console.log("new root-level file: edit after creation reloads"); + + // A top-level directory created after boot, then an edit inside it. + fs.mkdirSync(path.join(app, "lib")); + fs.writeFileSync(path.join(app, "lib", "dep.js"), `export const label = "L1";\n`); + fs.writeFileSync( + path.join(app, "src", "main.js"), + `import { banner } from "../banner.js";\nimport { label } from "../lib/dep.js";\ndocument.title = banner + "-" + label; window.__READY = true;\n`, + ); + await page.waitForFunction(() => document.title === "B2-L1", { timeout: 10000 }); + fs.writeFileSync(path.join(app, "lib", "dep.js"), `export const label = "L2";\n`); + await page.waitForFunction(() => document.title === "B2-L2", { timeout: 10000 }); + console.log("new top-level dir: edit after mkdir reloads"); + + // A staged tree RENAMED into the root (same filesystem, so a true rename): + // notify reports it as Modify(Name), not Create, and it must be adopted the + // same way. + const staging = fs.mkdtempSync(path.join(os.tmpdir(), "oj-watch-stage-")); + fs.mkdirSync(path.join(staging, "pkg")); + fs.writeFileSync(path.join(staging, "pkg", "mod.js"), `export const tag = "P1";\n`); + fs.renameSync(path.join(staging, "pkg"), path.join(app, "pkg")); + fs.writeFileSync( + path.join(app, "src", "main.js"), + `import { banner } from "../banner.js";\nimport { label } from "../lib/dep.js";\nimport { tag } from "../pkg/mod.js";\ndocument.title = banner + "-" + label + "-" + tag; window.__READY = true;\n`, + ); + await page.waitForFunction(() => document.title === "B2-L2-P1", { timeout: 10000 }); + fs.writeFileSync(path.join(app, "pkg", "mod.js"), `export const tag = "P2";\n`); + await page.waitForFunction(() => document.title === "B2-L2-P2", { timeout: 10000 }); + fs.rmSync(staging, { recursive: true, force: true }); + console.log("renamed-in top-level dir: edit after rename reloads"); + + assert.equal(errors.length, 0, `page errors: ${errors.join("|")}`); + console.log("WATCH-NEW-ROOT-DIR E2E PASSED"); +} catch (err) { + failed = true; + console.error("WATCH-NEW-ROOT-DIR E2E FAILED:", err.message); +} finally { + if (browser) await browser.close(); + srv.kill("SIGKILL"); + await sleep(300); + fs.rmSync(app, { recursive: true, force: true }); +} +process.exit(failed ? 1 : 0);