Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -258,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:
Expand Down
196 changes: 167 additions & 29 deletions crates/oj/src/start_dev.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<std::sync::atomic::AtomicBool>,
/// 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<std::collections::HashSet<PathBuf>>,
/// 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 {
Expand Down Expand Up @@ -395,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,
Expand Down Expand Up @@ -626,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<PathBuf>,
/// 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<tokio::task::JoinHandle<()>>,
/// 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<tokio::task::JoinHandle<bool>>,
/// 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,
}
Expand All @@ -644,7 +666,7 @@ fn spawn_mw_invalidate(
state: &StartState,
paths: &std::collections::HashSet<PathBuf>,
created: &std::collections::HashSet<PathBuf>,
) -> Option<tokio::task::JoinHandle<()>> {
) -> Option<tokio::task::JoinHandle<bool>> {
let port = state.plugin_serve.mw_port()?;
let changed: Vec<(String, &'static str)> = paths
.iter()
Expand All @@ -654,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<PathBuf> {
std::fs::read_to_string(cache.join("closure.json"))
.ok()
.and_then(|s| serde_json::from_str::<Vec<PathBuf>>(&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
Expand Down Expand Up @@ -774,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) => {
Expand Down Expand Up @@ -832,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();
Expand Down Expand Up @@ -870,6 +925,13 @@ async fn rebundle_worker(
};
let paths: Vec<PathBuf> = 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);
Expand All @@ -891,31 +953,51 @@ 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) {
Some((_, pinned)) => Some(pinned),
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
Expand Down Expand Up @@ -946,20 +1028,45 @@ async fn rebundle_worker(
" oj start: regen outputs changed ({}), invalidating worker",
names.join(", ")
);
oj_server::notify_plugin_mw_invalidate(port, &regen_changed).await;
let _ = oj_server::notify_plugin_mw_invalidate(port, &regen_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);
Expand All @@ -971,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" }
);
}
}
}
Expand Down Expand Up @@ -2713,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
Expand Down
18 changes: 12 additions & 6 deletions crates/oj_server/src/assets/plugin-host.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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;
}
Expand Down
19 changes: 18 additions & 1 deletion crates/oj_server/src/assets/start/bundle-client.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -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 = {
Expand Down Expand Up @@ -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))
Expand Down
Loading
Loading