refactor(web): rewrite webserver as a plain route-registration plugin

HttpServerService provides ctx.httpServer: register(route) -> disposer
(duplicate patterns throw), tapIndex transforms in registration order, and
the bound port; matching is exact > longest prefix > static dist fallback
(403/405/SPA semantics preserved). The server listens on activation, answers
per-request failures with 400 + a log line instead of exiting the process,
and knows no harness concepts — the boot graph, bundle routes, SSE channel,
and /api prefix all moved to their owning plugins.
This commit is contained in:
imccyu
2026-07-25 01:18:29 +08:00
parent 3cad7f6957
commit f499872cec
10 changed files with 194 additions and 1449 deletions
+161 -209
View File
@@ -1,232 +1,184 @@
/**
* @deepseek-ai/dsh-host-webserver — the web-shape HTTP carrier: node:http server
* routing /api/* to an injected fetch-shaped handler (node:http ↔ WHATWG
* bridge with SSE streamed out chunk by chunk) and everything else to static
* file serving. Web (browser) shape only — Electron loads dist over file://
* and carries fetch over an IPC bridge, not this server. This package never
* prints: the URL line belongs to the shell.
* @deepseek-ai/dsh-host-webserver — plain HTTP route-registration plugin: a
* node:http server plus the `httpServer` service (named-route registry + index
* transform taps + static dist fallback). Knows no harness concepts — every
* feature surface (API bridge, plugin bundles, SSE) is a route some other
* plugin registers. Web (browser) shape only — Electron loads dist over
* file:// and carries fetch over an IPC bridge, not this server. This package
* never prints: the URL line belongs to the shell.
*/
import { createServer } from 'node:http'
import type { IncomingMessage, ServerResponse } from 'node:http'
import type { IncomingMessage, ServerResponse, Server } from 'node:http'
import { readFile } from 'node:fs/promises'
import type { AddressInfo } from 'node:net'
import { dirname } from 'node:path'
import { Context, Service } from 'cordis'
import z from 'schemastery'
import { serveStatic } from './static.ts'
import { createPluginEventChannel } from './plugin-events.ts'
import type { HostWebPluginRegistry, WebBootGraph } from './web-plugins.ts'
export { createHostWebPluginRegistry } from './web-plugins.ts'
export type {
HostWebPluginRegistry, LoaderEntryView, LoaderView, WebBootEntry, WebBootGraph, WebPluginRegistryDeps,
} from './web-plugins.ts'
export type { PluginEventChannel, PluginEventFrame } from './plugin-events.ts'
/** Options for startWebServer. */
export interface WebServerOptions {
/** Address or hostname to listen on. */
host: string
/** Port to listen on; zero requests an OS-assigned port. */
port: number
/**
* Absolute path of index.html inside the static root — the caller resolves
* it (dist location is workspace knowledge of the shell, not this package's).
*/
distIndex: string
/** Fetch-shaped API carrier; /api/*-prefixed requests are bridged to it. */
apiHandler: { fetch: typeof fetch }
/**
* Web plugin table. When present, every index.html response carries the
* `window.__DSH_BOOT__` entry graph script, `/plugins/<id>/client.js` serves
* each fetch entry's client bundle, and `GET /plugins/events` streams graph/
* rebuilt frames (SSE) — rebuilt frames ride the registry's own bundle-watch
* notifications (`onRebuilt`). Absent = all three surfaces off (carrier-only
* use).
*/
webPlugins?: Pick<HostWebPluginRegistry, 'graph' | 'clientPath' | 'onRebuilt'>
declare module 'cordis' {
interface Context {
httpServer: HttpServerService
}
}
/** Listening web server handle. */
export interface RunningWebServer {
/** The listening port, including the OS-assigned value when options.port is zero. */
/** Route match kind: 'exact' matches the pathname verbatim; 'prefix' p matches p and p/<anything>. */
export type WebRouteKind = 'exact' | 'prefix'
/** One named route registration. */
export interface WebRoute {
kind: WebRouteKind
/** Absolute pathname, no trailing slash. */
path: string
/** Owns the full response lifecycle (may hold the response open, e.g. SSE). */
handler: (req: IncomingMessage, res: ServerResponse) => void | Promise<void>
}
/** Gateway config: listen address plus the static dist anchor (injected by the composing app, never self-resolved). */
export interface Config {
/** Listen host; the two supported values are loopback and all-interfaces. */
host: '127.0.0.1' | '0.0.0.0'
/** Listen port; zero requests an OS-assigned port. */
port: number
/**
* Shutdown: close + closeAllConnections (SSE connections never end on their
* own; without the force-close, close() would hang). Idempotent.
*/
close(): Promise<void>
/** Absolute path of index.html inside the static root (dist location is workspace knowledge of the app). */
distIndex: string
}
/**
* Start the web-shape HTTP server on the caller-selected host and port.
* Routing: /api/* → apiHandler bridge; non-GET/HEAD → 405; everything else →
* static with the step1-locked semantics (403 traversal, SPA fallback 200).
* A listen failure (EADDRINUSE…) rejects — the shell decides how to exit; a
* server error after listen goes to onError. A request whose handling throws
* (malformed %-escapes, a client dropping mid-body) is answered 400 — or the
* socket destroyed when headers are already out — and reported to onError;
* it never becomes an unhandled rejection.
* @param options - port, static root anchor, and the API carrier.
* @param onError - sink for post-listen server errors and per-request handling failures.
* @returns the running server handle once listening.
* The web-shape HTTP carrier service. Activation listens immediately (route
* registration order carries no request-facing semantics: named routes are
* composed to be disjoint, and the static dist fallback answers anything not
* yet claimed during the boot window). A listen failure throws out of init —
* a FAILED fiber the boot's fail-loud sweep reports.
*/
export function startWebServer(options: WebServerOptions, onError: (err: Error) => void): Promise<RunningWebServer> {
const { host, port, distIndex, apiHandler, webPlugins } = options
const distRoot = dirname(distIndex)
const renderIndex = webPlugins === undefined ? undefined : async (): Promise<string> => {
const html = await readFile(distIndex, 'utf8')
return injectBootManifest(html, webPlugins.graph())
}
const pluginEvents = webPlugins === undefined ? undefined : createPluginEventChannel()
// Rebuilt frames come from the registry's own bundle watch (dev mode); a
// prod registry without watching simply never notifies.
const unsubscribeRebuilt = webPlugins !== undefined && pluginEvents !== undefined
? webPlugins.onRebuilt((id, rev) => { pluginEvents.broadcast({ type: 'rebuilt', id, rev }) })
: undefined
export class HttpServerService extends Service {
static Config: z<Config> = z.object({
host: z.union([z.const('127.0.0.1'), z.const('0.0.0.0')]).required(),
port: z.natural().max(65535).required(),
distIndex: z.string().required(),
})
const handle = async (req: IncomingMessage, res: ServerResponse): Promise<void> => {
/* v8 ignore next -- `?? '/'` arm: node:http always sets url on server
requests; the field is only optional on the client-side IncomingMessage type */
const rawPath = new URL(req.url ?? '/', 'http://x').pathname
if (rawPath.startsWith('/api/')) {
await bridge(req, res, apiHandler)
return
}
if (req.method !== 'GET' && req.method !== 'HEAD') {
res.writeHead(405)
res.end()
return
}
if (webPlugins !== undefined && pluginEvents !== undefined && rawPath === '/plugins/events') {
pluginEvents.connect(res, webPlugins.graph())
return
}
if (webPlugins !== undefined && rawPath.startsWith('/plugins/') && rawPath.endsWith('/client.js')) {
await servePluginBundle(decodeURIComponent(rawPath), res, webPlugins)
return
}
await serveStatic(decodeURIComponent(rawPath), res, distRoot, distIndex, renderIndex)
private readonly exact = new Map<string, WebRoute>()
private readonly prefixes = new Map<string, WebRoute>()
private readonly indexTaps: ((html: string) => string)[] = []
private readonly distRoot: string
private readonly distIndex: string
private server!: Server
private listenedPort!: number
constructor(ctx: Context, private config: Config) {
super(ctx, 'httpServer')
this.distIndex = config.distIndex
this.distRoot = dirname(config.distIndex)
}
// Last-resort guard: handle() rejecting would otherwise be an unhandled
// rejection, and one malformed request (a bad %-escape hitting
// decodeURIComponent, a client dropping mid-body) would kill the whole
// process. Nothing after this catch can throw again on the same response.
const server = createServer((req, res) => {
handle(req, res).catch((err: unknown) => {
onError(err instanceof Error ? err : new Error(String(err)))
if (res.headersSent) {
res.destroy()
/** The listening port (the OS-assigned value when config.port is 0). */
get port(): number {
return this.listenedPort
}
/**
* Register a named route. Duplicate (kind, path) throws — route patterns are
* a composition-level contract, so a collision is a misconfiguration.
* @param route - kind, path, and the owning handler.
* @returns the disposer removing the route.
*/
register(route: WebRoute): () => void {
const table = route.kind === 'exact' ? this.exact : this.prefixes
if (table.has(route.path)) {
throw new Error(`webserver: duplicate ${route.kind} route "${route.path}"`)
}
table.set(route.path, route)
return () => { table.delete(route.path) }
}
/**
* Register an index.html transform, applied to every index response in
* registration order.
* @param transform - pure html-to-html function.
* @returns the disposer removing the transform.
*/
tapIndex(transform: (html: string) => string): () => void {
this.indexTaps.push(transform)
return () => {
const at = this.indexTaps.indexOf(transform)
if (at !== -1) this.indexTaps.splice(at, 1)
}
}
/** Listen; resolves once the socket is bound (rejection = FAILED fiber). */
async [Service.init](): Promise<void> {
const handle = async (req: IncomingMessage, res: ServerResponse): Promise<void> => {
/* v8 ignore next -- `?? '/'` arm: node:http always sets url on server
requests; the field is only optional on the client-side IncomingMessage type */
const rawPath = new URL(req.url ?? '/', 'http://x').pathname
const route = this.match(rawPath)
if (route !== undefined) {
await route.handler(req, res)
return
}
res.writeHead(400)
res.end()
})
})
let closing: Promise<void> | undefined
const close = (): Promise<void> => (closing ??= new Promise((resolveClose) => {
unsubscribeRebuilt?.()
server.close(() => { resolveClose() })
server.closeAllConnections()
}))
return new Promise((resolveListen, rejectListen) => {
server.once('error', rejectListen)
server.listen(port, host, () => {
server.off('error', rejectListen)
server.on('error', onError)
resolveListen({ port: (server.address() as AddressInfo).port, close })
})
})
}
/**
* Inject the boot entry graph into index.html: `window.__DSH_BOOT__` as the
* first script in <head> (before the shell bundle reads it). `<` is escaped in
* the JSON so plugin-controlled strings cannot break out of the script element.
* @param html - the index.html source.
* @param graph - the composed entry graph from the registry.
* @returns the html with the graph script injected.
*/
export function injectBootManifest(html: string, graph: WebBootGraph): string {
const json = JSON.stringify(graph).replaceAll('<', '\\u003c')
const script = `<script>window.__DSH_BOOT__ = ${json}</script>`
const head = html.indexOf('<head>')
if (head !== -1) return `${html.slice(0, head + 6)}${script}${html.slice(head + 6)}`
// Headless fixture pages may lack <head>; prepending keeps the read-before-shell ordering.
return `${script}${html}`
}
/**
* Serve one plugin client bundle from the registry table (unknown id = 404;
* the id may contain a scope slash). The `?rev=` query is a cache-busting
* parameter only — serving ignores it; `no-cache` makes the browser revalidate
* so a stale rev never sticks.
*/
async function servePluginBundle(
pathname: string, res: ServerResponse, webPlugins: Pick<HostWebPluginRegistry, 'clientPath'>,
): Promise<void> {
const id = pathname.slice('/plugins/'.length, -'/client.js'.length)
const path = webPlugins.clientPath(id)
if (path === undefined) {
res.writeHead(404)
res.end()
return
}
try {
const body = await readFile(path)
res.writeHead(200, { 'content-type': 'text/javascript; charset=utf-8', 'cache-control': 'no-cache' })
res.end(body)
} catch {
// Registered but unreadable (bundle not built yet): loud 404 beats a silent SPA-fallback HTML page.
res.writeHead(404)
res.end()
}
}
/** Bridge one node:http request to the WHATWG fetch handler (client close aborts; SSE bodies stream out chunk by chunk). */
async function bridge(req: IncomingMessage, res: ServerResponse, apiHandler: { fetch: typeof fetch }): Promise<void> {
const abort = new AbortController()
// Client-disconnect detection MUST hang off the response, not the request:
// since Node 16, IncomingMessage 'close' fires as soon as the request body is
// fully consumed (immediately for a bodyless GET), which would abort every SSE
// stream right after open. ServerResponse 'close' fires on connection teardown;
// writableEnded distinguishes a normal end() from the client going away.
res.on('close', () => {
if (!res.writableEnded) abort.abort()
})
const chunks: Buffer[] = []
for await (const chunk of req) chunks.push(chunk as Buffer)
/* v8 ignore next 3 -- `??` arms: node:http always sets url/method on server
requests; the fields are only optional on the client-side IncomingMessage type */
const request = new Request(new URL(req.url ?? '/', 'http://dsh.internal'), {
method: req.method ?? 'GET',
headers: Object.fromEntries(Object.entries(req.headers).filter(([, v]) => typeof v === 'string') as [string, string][]),
...chunks.length > 0 ? { body: Buffer.concat(chunks) } : {},
signal: abort.signal,
})
const response = await apiHandler.fetch(request)
res.writeHead(response.status, Object.fromEntries(response.headers.entries()))
if (response.body === null) {
res.end()
return
}
for await (const chunk of response.body) {
// Backpressure: a false return means the socket buffer is full — wait for drain
// instead of buffering unboundedly (slow/suspended SSE consumers). 'close' also
// resolves so a mid-wait disconnect can't park this loop forever; the close
// handler above aborts the handler stream, which then ends the iteration.
if (!res.write(chunk)) {
await new Promise<void>((resolve) => {
const done = (): void => {
res.off('drain', done)
res.off('close', done)
resolve()
}
res.once('drain', done)
res.once('close', done)
})
// Static fallback keeps the pre-plugin semantics: non-GET/HEAD is 405,
// traversal 403, miss falls back to index.html 200 (SPA routing).
if (req.method !== 'GET' && req.method !== 'HEAD') {
res.writeHead(405)
res.end()
return
}
await serveStatic(decodeURIComponent(rawPath), res, this.distRoot, this.distIndex, () => this.renderIndex())
}
// Last-resort guard: handle() rejecting would otherwise be an unhandled
// rejection killing the process on one malformed request (bad %-escape,
// client dropping mid-body). Per-request failures log and answer 400 —
// never a process exit.
this.server = createServer((req, res) => {
handle(req, res).catch((err: unknown) => {
this.ctx.logger.warn(err instanceof Error ? err : new Error(String(err)))
if (res.headersSent) {
res.destroy()
return
}
res.writeHead(400)
res.end()
})
})
await new Promise<void>((resolve, reject) => {
this.server.once('error', reject)
this.server.listen(this.config.port, this.config.host, () => {
this.server.off('error', reject)
this.server.on('error', (err) => { this.ctx.logger.error(err) })
this.listenedPort = (this.server.address() as AddressInfo).port
resolve()
})
})
// close + closeAllConnections: held-open responses (SSE) never end on
// their own; without the force-close, close() would hang teardown.
this.ctx.effect(() => () => new Promise<void>((resolve) => {
this.server.close(() => { resolve() })
this.server.closeAllConnections()
}), 'httpServer.listen')
}
/** Longest-prefix-wins over the prefix table after an exact-table miss. */
private match(pathname: string): WebRoute | undefined {
const exact = this.exact.get(pathname)
if (exact !== undefined) return exact
let best: WebRoute | undefined
for (const [prefix, route] of this.prefixes) {
if (pathname !== prefix && !pathname.startsWith(`${prefix}/`)) continue
if (best === undefined || prefix.length > best.path.length) best = route
}
return best
}
/** Index body: dist index.html through the registered taps in order. */
private async renderIndex(): Promise<string> {
let html = await readFile(this.distIndex, 'utf8')
for (const transform of this.indexTaps) html = transform(html)
return html
}
res.end()
}
export default HttpServerService
+20 -18
View File
@@ -15,28 +15,30 @@ export const name = 'host-webserver-invariant'
export const inject = ['invariants']
/**
* Owned relation: the web plugin registry's boot entry graph must stay
* self-consistent — every row must resolve a clientPath under the same id
* (the /plugins/<id>/client.js URL it advertises would otherwise 404 on a
* browser that just received the graph). Checked synchronously on every
* rescan trigger (cordis 'internal/plugin'): graph() and clientPath() read
* the same table object, so the relation is self-consistent at any instant —
* no need to wait out the registry's own debounced rescan. The registry
* arrives through the context key the assembly publishes it under.
* Owned relation: route registrations and their disposers must stay
* symmetric — after the owning fiber of a registered route unloads, the
* route table must no longer answer for its path (a stale route would keep
* serving a disposed plugin's handler). Checked on every fiber teardown
* (cordis 'internal/plugin'): the service's own registry state is compared
* against the set of live fibers' registrations indirectly, by probing that
* dispose really removed the entry — the register() disposer contract.
*/
const install: InvariantInstaller = (ctx, fail) => {
ctx.on('internal/plugin', () => {
const registry = ctx.get('webPlugins') as
| {
graph(): { entries: { id: string; url: string }[] }
clientPath(id: string): string | undefined
}
const server = ctx.get('httpServer') as
| { register(route: { kind: 'exact'; path: string; handler: () => void }): () => void }
| undefined
if (registry === undefined) return // carrier-only deployments never publish the registry
for (const row of registry.graph().entries) {
if (registry.clientPath(row.id) === undefined) {
fail(`web plugin graph row "${row.id}" advertises ${row.url} but resolves no client bundle path — the served __DSH_BOOT__ would 404 on fetch`)
}
if (server === undefined) return // no webserver row in this composition
// Register/dispose probe on a reserved path: if dispose leaves the route
// behind, a second register throws the duplicate error — the asymmetry.
// Each register(probe)() is one register+dispose cycle, so the probe never
// leaves residue; a leftover from the first cycle makes the second throw.
const probe = { kind: 'exact' as const, path: '/__dsh_invariant_probe__', handler: () => {} }
try {
server.register(probe)()
server.register(probe)()
} catch {
fail('httpServer.register() disposer left the route registered — route table and fiber lifecycles diverged')
}
}, { global: true })
}
@@ -1,56 +0,0 @@
/**
* `/plugins/events` SSE channel: the system-side push surface for the client
* entry graph (connect → current graph frame; dev rebuild → rebuilt frame).
* Presentation-only wire — frames never enter the session log (distinct from
* the /api/* session SSE, which is api-contract territory). Connections are
* plain node:http responses held in a set; the server's closeAllConnections
* tears them down on shutdown.
*/
import type { ServerResponse } from 'node:http'
import type { WebBootGraph } from './web-plugins.ts'
/** One `/plugins/events` frame: the full graph on connect, or one rebuilt bundle notice. */
export type PluginEventFrame =
| { type: 'graph'; graph: WebBootGraph }
| { type: 'rebuilt'; id: string; rev: string }
/** Broadcast surface owned by the webserver routing layer. */
export interface PluginEventChannel {
/** Adopt one incoming SSE request: writes the SSE preamble and the current-graph frame, then keeps the response open. */
connect(res: ServerResponse, graph: WebBootGraph): void
/** Push one frame to every open connection. */
broadcast(frame: PluginEventFrame): void
}
/** Serialize one frame as an SSE data line. */
function sseData(frame: PluginEventFrame): string {
return `data: ${JSON.stringify(frame)}\n\n`
}
/**
* Create the channel (one per running server).
* @returns the connect/broadcast surface.
*/
export function createPluginEventChannel(): PluginEventChannel {
const connections = new Set<ServerResponse>()
return {
connect(res, graph) {
res.writeHead(200, {
'content-type': 'text/event-stream',
'cache-control': 'no-cache',
'connection': 'keep-alive',
})
// Comment line on open so clients/proxies see a live channel even when
// no rebuild ever happens; EventSource frame parsing skips it naturally.
res.write(': connected\n\n')
res.write(sseData({ type: 'graph', graph }))
connections.add(res)
res.on('close', () => { connections.delete(res) })
},
broadcast(frame) {
const line = sseData(frame)
for (const res of connections) res.write(line)
},
}
}
-360
View File
@@ -1,360 +0,0 @@
/**
* HostWebPluginRegistry: composes the client entry graph served as
* `window.__DSH_BOOT__` ({rev, entries}). Every row is discovered among the
* host Loader's loaded entries by its package.json `dshClient` declaration
* (all client plugin packages arrive by fetch — one uniform bundle shape),
* resolving each one's client bundle path from `exports["./client"]` and
* hashing the bundle content into a `rev` (cache busting + HMR diff anchor).
* `inject` edges and the `immediately` prefetch mark come from the manifest
* (dshClient — the package owns its dependency edges and its boot tier); the
* composition layer contributes only the roster. The webserver consumes the
* table to emit the boot graph and to serve `GET /plugins/<id>/client.js`;
* in dev mode the registry additionally stat-polls each scanned bundle file
* and re-hashes + notifies `onRebuilt` subscribers on change (the rebuild
* signal is the registry's own observation — no builder protocol exists).
*
* The vendored loader emits no "entry loaded" event (only `loader/entry-init`,
* which fires at Entry construction before import/apply), so the registry
* scans `loader.entries()` and rescans on cordis `internal/plugin` (fiber
* create/dispose), microtask-debounced. Plugin-set changes take effect on
* restart per the config-source ruling; the subscription only keeps the table
* fresh within a process lifetime.
*/
import { createHash } from 'node:crypto'
import { readFileSync, statSync, type Stats } from 'node:fs'
import { dirname, join } from 'node:path'
import type { Context } from 'cordis'
/** One composed client entry (`window.__DSH_BOOT__.entries` row). */
export interface WebBootEntry {
/** Entry name == package name. */
id: string
/** Bundle URL served by this webserver (`/plugins/<id>/client.js?rev=<rev>`). */
url: string
/** Bundle content hash (sha1, shortened). */
rev: string
/** Package-name dependency edges from the manifest (dshClient.inject), informational (preflight/HMR display). */
inject?: string[]
/** Boot phase-one prefetch tier: the shell fetches these bundles in parallel before creating entries. */
immediately?: boolean
}
/** The composed entry graph: injected into index.html and pushed on /plugins/events connect. */
export interface WebBootGraph {
/** Consistency anchor over all rows: changes whenever any entry row changes. */
rev: string
/** All composed entries (order carries no semantics; governance ordering is the client Loader's job). */
entries: WebBootEntry[]
}
/** The web plugin table consumed by the boot injection, the bundle endpoint, and the rebuild channel. */
export interface HostWebPluginRegistry {
/** Current composed entry graph (stable object between changes). */
graph(): WebBootGraph
/**
* Absolute path of an entry's client bundle.
* @param id - entry id (package name).
* @returns the path, or undefined for an unknown id.
*/
clientPath(id: string): string | undefined
/**
* Re-hash one entry's bundle: updates the row's rev/url and the graph rev.
* The dev bundle watch calls this on every observed file change.
* @param id - entry id (package name).
* @returns the new bundle rev, or undefined for an unknown id.
*/
rebuilt(id: string): string | undefined
/**
* Subscribe to bundle rebuilds observed by the dev watch (only fires when
* the re-hash produced a different rev — an unchanged bundle is silent).
* @param listener - receives the entry id and its new bundle rev.
* @returns the unsubscriber.
*/
onRebuilt(listener: (id: string, rev: string) => void): () => void
/** Remove the loader subscription, all bundle watches, and all rebuild listeners. */
dispose(): void
}
/** Structural view of a loader entry (webserver keeps zero workspace dependencies; cordis stays a type-only peer). */
export interface LoaderEntryView {
options: { name: string }
/** Present once the entry's plugin fiber exists (import succeeded and apply ran/started). */
fiber?: unknown
/** True when the entry or an owning group is disabled. */
disabled: boolean
}
/** Structural view of the host Loader (entry enumeration is all the registry needs). */
export interface LoaderView {
entries(): Iterable<LoaderEntryView>
}
/** Dependencies injected by the assembly layer. */
export interface WebPluginRegistryDeps {
/** Host root context; used only to subscribe `internal/plugin` for rescans. */
ctx: Context
/** The host Loader owning the plugin entries. */
loader: LoaderView
/**
* Resolve a package specifier to its package.json absolute path (assembly
* passes `createRequire(...).resolve(`${name}/package.json`)`); injected so
* the registry makes no module-resolution assumptions of its own.
*/
resolvePkgJson: (name: string) => string
/** 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
* with an explicit stat baseline (polling by design: network mounts deliver
* no inotify events) and re-hash + notify onRebuilt subscribers on change.
* Absent = no watching (prod composition).
*/
watch?: {
/** Stat-poll interval in milliseconds; default 500 (the build-side watcher's polling default). */
intervalMs?: number
}
}
/** package.json `dshClient` declaration shape (file boundary — validated field by field). */
interface DshClientDeclaration {
inject?: string[]
platform: string
/** Boot phase-one prefetch mark; absent means lazy (fetched on demand). */
immediately?: boolean
}
interface WebPluginRecord {
entry: WebBootEntry
clientPath: string
}
interface WatchedBundle {
path: string
mtimeMs: number
size: number
dirty: boolean
}
/** Narrow an unknown parsed JSON value to the dshClient declaration, throwing on malformed fields. */
function parseDshClient(name: string, value: unknown): DshClientDeclaration | undefined {
if (value === undefined) return undefined
if (typeof value !== 'object' || value === null) {
throw new Error(`web-plugins: ${name} has a non-object dshClient declaration`)
}
const decl = value as Record<string, unknown>
if (typeof decl.platform !== 'string') {
throw new Error(`web-plugins: ${name} dshClient.platform must be a string`)
}
if (decl.inject !== undefined && (!Array.isArray(decl.inject) || decl.inject.some(i => typeof i !== 'string'))) {
throw new Error(`web-plugins: ${name} dshClient.inject must be a string array`)
}
if (decl.immediately !== undefined && typeof decl.immediately !== 'boolean') {
throw new Error(`web-plugins: ${name} dshClient.immediately must be a boolean`)
}
return {
platform: decl.platform,
...(decl.inject !== undefined ? { inject: decl.inject as string[] } : {}),
...(decl.immediately !== undefined ? { immediately: decl.immediately } : {}),
}
}
/** Resolve `exports["./client"]` to a relative path, accepting the string and one-level conditional forms. */
function clientExportOf(name: string, exportsField: unknown): string | undefined {
if (typeof exportsField !== 'object' || exportsField === null) return undefined
const client = (exportsField as Record<string, unknown>)['./client']
if (client === undefined) return undefined
if (typeof client === 'string') return client
if (typeof client === 'object' && client !== null) {
const fallback = (client as Record<string, unknown>).default
if (typeof fallback === 'string') return fallback
}
throw new Error(`web-plugins: ${name} exports["./client"] has an unsupported shape`)
}
/** sha1 content hash shortened to 12 hex chars (bundle rev / graph rev). */
function shortHash(input: string | Buffer): string {
return createHash('sha1').update(input).digest('hex').slice(0, 12)
}
/** Graph row for one bundle rev (url carries the rev as its cache-busting query). */
function graphRow(id: string, rev: string, inject: string[] | undefined, immediately: boolean): WebBootEntry {
return {
id,
url: `/plugins/${id}/client.js?rev=${rev}`,
rev,
...(inject !== undefined ? { inject } : {}),
...(immediately ? { immediately: true } : {}),
}
}
/** Compose the graph value from the current table. */
function composeGraph(table: Map<string, WebPluginRecord>): WebBootGraph {
const entries = [...table.values()].map(record => record.entry)
return { rev: shortHash(JSON.stringify(entries)), entries }
}
/**
* Build the web plugin registry: scan once synchronously (a malformed
* declaration, an unbuilt bundle, or an invalid watch interval throws here —
* load-time fail loud), then rescan on `internal/plugin`, microtask-debounced
* (failures go to `deps.onError`). With `deps.watch`, every scanned bundle
* file is stat-polled and a content change re-hashes the row and notifies
* `onRebuilt` subscribers.
* @param deps - loader view, resolution hook, error sink, and optional dev watch (see {@link WebPluginRegistryDeps}).
* @returns the registry handle.
*/
export function createHostWebPluginRegistry(deps: WebPluginRegistryDeps): HostWebPluginRegistry {
const watchInterval = deps.watch === undefined ? undefined : deps.watch.intervalMs ?? 500
if (watchInterval !== undefined && (!Number.isInteger(watchInterval) || watchInterval <= 0)) {
throw new Error(`web-plugins: watch.intervalMs must be a positive integer (got ${String(deps.watch?.intervalMs)})`)
}
const stageWatches = (
candidateTable: Map<string, WebPluginRecord>,
currentWatches: Map<string, WatchedBundle>,
): Map<string, WatchedBundle> => {
const candidateWatches = new Map<string, WatchedBundle>()
if (watchInterval === undefined) return candidateWatches
for (const [id, record] of candidateTable) {
const current = currentWatches.get(id)
if (current?.path === record.clientPath) {
candidateWatches.set(id, { ...current })
continue
}
const baseline = statSync(record.clientPath)
candidateWatches.set(id, {
path: record.clientPath,
mtimeMs: baseline.mtimeMs,
size: baseline.size,
dirty: false,
})
}
return candidateWatches
}
let table = scan(deps)
let graph = composeGraph(table)
let watched = stageWatches(table, new Map())
const rebuildListeners = new Set<(id: string, rev: string) => void>()
const rebuilt = (id: string): string | undefined => {
const record = table.get(id)
if (record === undefined) return undefined
const rev = shortHash(readFileSync(record.clientPath))
record.entry = graphRow(id, rev, record.entry.inject, record.entry.immediately === true)
graph = composeGraph(table)
return rev
}
// Dev bundle watch: capture every row's baseline synchronously before the
// registry is returned, then poll those baselines. fs.watchFile establishes
// its first baseline asynchronously, so an immediate rebuild can otherwise
// become the baseline and disappear without an observed delta.
const pollWatches = (): void => {
for (const [id, watch] of watched) {
let current: Stats
try {
current = statSync(watch.path)
} catch (error) {
const code = (error as NodeJS.ErrnoException).code
if (code === 'ENOENT') {
watch.dirty = true
continue
}
deps.onError(error instanceof Error ? error : new Error(String(error)))
continue
}
if (!watch.dirty && current.mtimeMs === watch.mtimeMs && current.size === watch.size) continue
const before = table.get(id)?.entry.rev
let rev: string | undefined
try {
rev = rebuilt(id)
} catch (error) {
const code = (error as NodeJS.ErrnoException).code
if (code === 'ENOENT') {
watch.dirty = true
continue
}
watch.mtimeMs = current.mtimeMs
watch.size = current.size
deps.onError(error instanceof Error ? error : new Error(String(error)))
continue
}
watch.mtimeMs = current.mtimeMs
watch.size = current.size
watch.dirty = false
if (rev === undefined || rev === before) continue
for (const notify of rebuildListeners) {
// A throwing subscriber must not skip later subscribers or escape the
// polling callback into the process event loop.
try {
notify(id, rev)
} catch (error) {
deps.onError(error instanceof Error ? error : new Error(String(error)))
}
}
}
}
const watchTimer = watchInterval === undefined ? undefined : setInterval(pollWatches, watchInterval)
watchTimer?.unref()
let pending = false
const unsubscribe = deps.ctx.on('internal/plugin', () => {
if (pending) return
pending = true
queueMicrotask(() => {
pending = false
try {
const candidateTable = scan(deps)
const candidateGraph = composeGraph(candidateTable)
const candidateWatches = stageWatches(candidateTable, watched)
table = candidateTable
graph = candidateGraph
watched = candidateWatches
} catch (error) {
// Keep serving the previous graph: a mid-flight rescan failure must not
// take down the boot manifest for plugins that were fine.
deps.onError(error instanceof Error ? error : new Error(String(error)))
}
})
})
return {
graph: () => graph,
clientPath: id => table.get(id)?.clientPath,
rebuilt,
onRebuilt: (listener) => {
rebuildListeners.add(listener)
return () => { rebuildListeners.delete(listener) }
},
dispose: () => {
unsubscribe()
if (watchTimer !== undefined) clearInterval(watchTimer)
watched.clear()
rebuildListeners.clear()
},
}
}
/** One full table build from the loader's current entries (bundle content is hashed here — an unreadable bundle throws). */
function scan(deps: WebPluginRegistryDeps): Map<string, WebPluginRecord> {
const table = new Map<string, WebPluginRecord>()
for (const entry of deps.loader.entries()) {
if (entry.fiber === undefined || entry.disabled) continue
const name = entry.options.name
if (table.has(name)) continue
const pkgPath = deps.resolvePkgJson(name)
const pkg = JSON.parse(readFileSync(pkgPath, 'utf8')) as Record<string, unknown>
const decl = parseDshClient(name, pkg.dshClient)
if (decl === undefined || decl.platform !== 'web') continue
const clientRel = clientExportOf(name, pkg.exports)
if (clientRel === undefined) {
throw new Error(`web-plugins: ${name} declares dshClient but exports no "./client" bundle`)
}
const clientPath = join(dirname(pkgPath), clientRel)
const rev = shortHash(readFileSync(clientPath))
table.set(name, { entry: graphRow(name, rev, decl.inject, decl.immediately === true), clientPath })
}
return table
}