fix(webserver): registry-owned stat poll replaces fs.watchFile baseline race
The dev bundle watch missed rebuilds that landed while the registry was
constructing: fs.watchFile captures its comparison baseline with an
ASYNCHRONOUS first stat, so a write racing that window is absorbed into
the baseline and never reported. Standalone repro missed 24/400 same-tick
rewrites; the CI flake in web-plugins.spec.ts ('watch mode: a bundle
content change re-hashes the row...') was exactly this — the spec writes
immediately after createHostWebPluginRegistry returns.
The watch now polls from one registry-owned setInterval against a stat
baseline the scan itself captures synchronously, stat-before-read: a
write landing between stat and read leaves the hash newer than the
baseline (next tick re-hashes to the same rev, no spurious notify); a
write landing after the read leaves the baseline older (next tick
detects and notifies). No blind window. The poll iterates the live
table, so rescans retarget the watch for free and dispose clears one
timer. Stress: real-registry same-tick rewrite 0/600 missed (was 1/300);
spec watch tests 0/50.
New regression test pins the same-tick-as-construction write. Its
rewrite deliberately differs in size from the seed: a same-millisecond
same-size rewrite is invisible to any mtime+size poll (coarse fs
timestamps) — a stat-polling limit, not this regression. Loading-model
Agent Note updated in both languages (pair re-recorded).
This commit is contained in:
@@ -22,8 +22,7 @@
|
||||
*/
|
||||
|
||||
import { createHash } from 'node:crypto'
|
||||
import { readFileSync, unwatchFile, watchFile } from 'node:fs'
|
||||
import type { Stats } from 'node:fs'
|
||||
import { readFileSync, statSync } from 'node:fs'
|
||||
import { dirname, join } from 'node:path'
|
||||
import type { Context } from 'cordis'
|
||||
|
||||
@@ -106,10 +105,13 @@ export interface WebPluginRegistryDeps {
|
||||
/** Sink for rescan failures (the initial scan throws instead — misconfiguration fails loud at load). */
|
||||
onError: (err: Error) => void
|
||||
/**
|
||||
* Dev-mode bundle watching: stat-poll every scanned row's client bundle
|
||||
* (fs.watchFile — polling by design: network mounts deliver no inotify
|
||||
* events) and re-hash + notify onRebuilt subscribers on change. Absent =
|
||||
* no watching (prod composition).
|
||||
* Dev-mode bundle watching: one registry-owned interval stat-polls every
|
||||
* scanned row's client bundle (polling by design: network mounts deliver no
|
||||
* inotify events) and re-hashes + notifies onRebuilt subscribers on change.
|
||||
* Each row's stat baseline is captured synchronously before its content is
|
||||
* hashed, so a rebuild landing while the registry constructs is still
|
||||
* detected on the first tick (fs.watchFile's asynchronous baseline lost
|
||||
* that window). Absent = no watching (prod composition).
|
||||
*/
|
||||
watch?: {
|
||||
/** Stat-poll interval in milliseconds; default 500 (the build-side watcher's polling default). */
|
||||
@@ -128,6 +130,15 @@ interface DshClientDeclaration {
|
||||
interface WebPluginRecord {
|
||||
entry: WebBootEntry
|
||||
clientPath: string
|
||||
/**
|
||||
* Bundle stat captured immediately BEFORE the content read that produced
|
||||
* `entry.rev` — the watch baseline. The stat→read order makes a write
|
||||
* racing the scan converge instead of being absorbed: landing between stat
|
||||
* and read leaves the hash newer than the baseline (next tick re-hashes to
|
||||
* the same rev, no spurious notify); landing after the read leaves the
|
||||
* baseline older (next tick detects, re-hashes, notifies).
|
||||
*/
|
||||
stat: { mtimeMs: number; size: number }
|
||||
}
|
||||
|
||||
/** Narrow an unknown parsed JSON value to the dshClient declaration, throwing on malformed fields. */
|
||||
@@ -211,57 +222,62 @@ export function createHostWebPluginRegistry(deps: WebPluginRegistryDeps): HostWe
|
||||
const rebuilt = (id: string): string | undefined => {
|
||||
const record = table.get(id)
|
||||
if (record === undefined) return undefined
|
||||
// stat BEFORE read, like scan(): a write racing this pair converges (see
|
||||
// WebPluginRecord.stat) instead of desynchronizing baseline and rev.
|
||||
const stat = statSync(record.clientPath)
|
||||
const rev = shortHash(readFileSync(record.clientPath))
|
||||
record.stat = { mtimeMs: stat.mtimeMs, size: stat.size }
|
||||
record.entry = graphRow(id, rev, record.entry.inject, record.entry.immediately === true)
|
||||
graph = composeGraph(table)
|
||||
return rev
|
||||
}
|
||||
|
||||
// Dev bundle watch: one fs.watchFile stat poll per table row. A torn read
|
||||
// of a half-written bundle self-heals — the ongoing write keeps changing
|
||||
// the stats, so the next poll tick re-hashes the completed file.
|
||||
const watched = new Map<string, { path: string; listener: (curr: Stats, prev: Stats) => void }>()
|
||||
const syncWatches = (): void => {
|
||||
if (watchInterval === undefined) return
|
||||
for (const [id, watch] of watched) {
|
||||
if (table.get(id)?.clientPath === watch.path) continue
|
||||
unwatchFile(watch.path, watch.listener)
|
||||
watched.delete(id)
|
||||
}
|
||||
// Dev bundle watch: one registry-owned setInterval stat-polls every table
|
||||
// row against the record's own baseline. fs.watchFile is unusable here: it
|
||||
// captures its comparison baseline with an ASYNCHRONOUS first stat, so a
|
||||
// rebuild landing between scan()'s content read and that stat is absorbed
|
||||
// into the baseline and never reported — and the missed window is exactly
|
||||
// registry construction, when a dev build is most likely to be finishing.
|
||||
// The record baseline has no such window: scan()/rebuilt() stat before they
|
||||
// read, so any write the hash missed is newer than the baseline and lands
|
||||
// on the next tick. A torn read of a half-written bundle self-heals the
|
||||
// same way — the ongoing write keeps changing the stats.
|
||||
const pollTick = (): void => {
|
||||
for (const [id, record] of table) {
|
||||
if (watched.has(id)) continue
|
||||
const listener = (curr: Stats, prev: Stats): void => {
|
||||
// fs.watchFile fires on any stat delta (atime included); only content
|
||||
// signals count. An all-zero curr means the file vanished mid-rebuild
|
||||
// — the completing write fires the next tick, so skipping is safe.
|
||||
if (curr.mtimeMs === prev.mtimeMs && curr.size === prev.size) return
|
||||
if (curr.mtimeMs === 0) return
|
||||
const before = table.get(id)?.entry.rev
|
||||
let rev: string | undefined
|
||||
let stat: { mtimeMs: number; size: number }
|
||||
try {
|
||||
stat = statSync(record.clientPath)
|
||||
} catch (error) {
|
||||
const code = (error as NodeJS.ErrnoException).code
|
||||
if (code === 'ENOENT') continue // mid-rename window; the completed write lands on a later tick
|
||||
deps.onError(error instanceof Error ? error : new Error(String(error)))
|
||||
continue
|
||||
}
|
||||
if (stat.mtimeMs === record.stat.mtimeMs && stat.size === record.stat.size) continue
|
||||
const before = record.entry.rev
|
||||
let rev: string | undefined
|
||||
try {
|
||||
rev = rebuilt(id)
|
||||
} catch (error) {
|
||||
const code = (error as NodeJS.ErrnoException).code
|
||||
if (code === 'ENOENT') continue // vanished between stat and read; same self-heal
|
||||
deps.onError(error instanceof Error ? error : new Error(String(error)))
|
||||
continue
|
||||
}
|
||||
if (rev === undefined || rev === before) continue
|
||||
for (const notify of rebuildListeners) {
|
||||
// A throwing subscriber must not skip later subscribers or escape
|
||||
// into the timer callback (that would kill the process).
|
||||
try {
|
||||
rev = rebuilt(id)
|
||||
notify(id, rev)
|
||||
} catch (error) {
|
||||
const code = (error as NodeJS.ErrnoException).code
|
||||
if (code === 'ENOENT') return // mid-rename window; the completed write fires the next poll tick
|
||||
deps.onError(error instanceof Error ? error : new Error(String(error)))
|
||||
return
|
||||
}
|
||||
if (rev === undefined || rev === before) return
|
||||
for (const notify of rebuildListeners) {
|
||||
// A throwing subscriber must not escape the fs.watchFile callback
|
||||
// (that would skip later subscribers and can kill the process).
|
||||
try {
|
||||
notify(id, rev)
|
||||
} catch (error) {
|
||||
deps.onError(error instanceof Error ? error : new Error(String(error)))
|
||||
}
|
||||
}
|
||||
}
|
||||
watchFile(record.clientPath, { interval: watchInterval, persistent: false }, listener)
|
||||
watched.set(id, { path: record.clientPath, listener })
|
||||
}
|
||||
}
|
||||
syncWatches()
|
||||
const pollTimer = watchInterval === undefined ? undefined : setInterval(pollTick, watchInterval)
|
||||
pollTimer?.unref()
|
||||
|
||||
let pending = false
|
||||
const unsubscribe = deps.ctx.on('internal/plugin', () => {
|
||||
@@ -270,9 +286,10 @@ export function createHostWebPluginRegistry(deps: WebPluginRegistryDeps): HostWe
|
||||
queueMicrotask(() => {
|
||||
pending = false
|
||||
try {
|
||||
// The poll iterates `table` directly, so the swap also retargets the
|
||||
// watch: fresh records carry fresh stat baselines from scan().
|
||||
table = scan(deps)
|
||||
graph = composeGraph(table)
|
||||
syncWatches()
|
||||
} catch (error) {
|
||||
// Keep serving the previous graph: a mid-flight rescan failure must not
|
||||
// take down the boot manifest for plugins that were fine.
|
||||
@@ -291,8 +308,7 @@ export function createHostWebPluginRegistry(deps: WebPluginRegistryDeps): HostWe
|
||||
},
|
||||
dispose: () => {
|
||||
unsubscribe()
|
||||
for (const { path, listener } of watched.values()) unwatchFile(path, listener)
|
||||
watched.clear()
|
||||
if (pollTimer !== undefined) clearInterval(pollTimer)
|
||||
rebuildListeners.clear()
|
||||
},
|
||||
}
|
||||
@@ -314,8 +330,13 @@ function scan(deps: WebPluginRegistryDeps): Map<string, WebPluginRecord> {
|
||||
throw new Error(`web-plugins: ${name} declares dshClient but exports no "./client" bundle`)
|
||||
}
|
||||
const clientPath = join(dirname(pkgPath), clientRel)
|
||||
const stat = statSync(clientPath)
|
||||
const rev = shortHash(readFileSync(clientPath))
|
||||
table.set(name, { entry: graphRow(name, rev, decl.inject, decl.immediately === true), clientPath })
|
||||
table.set(name, {
|
||||
entry: graphRow(name, rev, decl.inject, decl.immediately === true),
|
||||
clientPath,
|
||||
stat: { mtimeMs: stat.mtimeMs, size: stat.size },
|
||||
})
|
||||
}
|
||||
return table
|
||||
}
|
||||
|
||||
@@ -147,6 +147,27 @@ describe('createHostWebPluginRegistry', () => {
|
||||
expect(rebuilds).toHaveLength(1)
|
||||
})
|
||||
|
||||
it('watch mode: a write landing during registry construction is still detected (regression: fs.watchFile baseline absorption)', async () => {
|
||||
// The old fs.watchFile watch captured its comparison baseline with an
|
||||
// ASYNCHRONOUS first stat; a rewrite in the same tick as construction was
|
||||
// absorbed into that baseline and never reported (the CI flake). The
|
||||
// record-baseline poll stats synchronously before hashing, so this exact
|
||||
// timing must now always notify.
|
||||
const { deps, root } = makeDeps([{ name: 'watched', pkg: webDecl() }])
|
||||
deps.watch = { intervalMs: 20 }
|
||||
const registry = createHostWebPluginRegistry(deps)
|
||||
const rebuilds: string[] = []
|
||||
registry.onRebuilt(id => rebuilds.push(id))
|
||||
// Same tick as construction — inside the old watch's blind window. The
|
||||
// rewrite deliberately differs in SIZE from the seed: a same-millisecond
|
||||
// same-size rewrite is invisible to any mtime+size poll by construction
|
||||
// (coarse filesystem timestamps), which is a stat-polling limit, not the
|
||||
// regression under test.
|
||||
writeFileSync(join(root, 'watched', 'lib', 'client.js'), '// same-tick rewritten contents')
|
||||
await vi.waitFor(() => { expect(rebuilds).toEqual(['watched']) }, { timeout: 5000 })
|
||||
registry.dispose()
|
||||
})
|
||||
|
||||
it('rejects a non-positive or non-integer watch interval at build time', () => {
|
||||
for (const intervalMs of [0, -5, 1.5]) {
|
||||
const { deps } = makeDeps([{ name: 'p', pkg: webDecl() }])
|
||||
|
||||
Reference in New Issue
Block a user