import { spawn } from "node:child_process"; import { createHash, randomUUID } from "node:crypto"; import { existsSync, readFileSync } from "node:fs"; import { access, copyFile, mkdir, open, readdir, readFile, rename, rm, stat, writeFile, } from "node:fs/promises"; import { homedir } from "node:os"; import { dirname, extname, join, posix, relative, resolve, sep, win32, } from "node:path"; import { fileURLToPath } from "node:url"; import type { ToolDefinition } from "@earendil-works/pi-coding-agent"; import { Type } from "typebox"; import { artifactFiles, hashArtifacts } from "./dynamic-artifacts.ts"; import { type DynamicManifest, VALIDATOR_WATCHDOG_MS, type ValidationAttestation, type ValidationDiagnostic, type WorkerValidationResult, } from "./dynamic-contract.ts"; import type { DynamicExtensionResolver, ResolvedDynamicExtension, UsedDynamicExtension, } from "./dynamic-types.ts"; import { DYNAMIC_EXTENSION_TOOL_NAMES } from "./spec.ts"; const MAX_TOOL_OUTPUT = 12_000; const PI_PACKAGE_NAME = "@earendil-works/pi-coding-agent"; /** * The validator worker ships next to this module with the same file extension: * `.ts` in the source checkout, `.js` inside a generated package. */ export function defaultWorkerPath(): string { const here = fileURLToPath(import.meta.url); return join(dirname(here), `dynamic-validator-worker${extname(here)}`); } let piPackageRootCache: string | undefined | null = null; /** * Locate the running Pi installation so a standalone worker process can import * Pi modules without a repository `node_modules` ancestor. */ export function piPackageRoot( env: NodeJS.ProcessEnv = process.env, mainScript: string | undefined = process.argv[1], ): string | undefined { if (env.ULTRA_PI_PACKAGE_ROOT) return env.ULTRA_PI_PACKAGE_ROOT; if (piPackageRootCache !== null) return piPackageRootCache; let dir = mainScript ? dirname(resolve(mainScript)) : undefined; while (dir) { try { const pkg = JSON.parse( readFileSync(join(dir, "package.json"), "utf8"), ) as { name?: unknown }; if (pkg.name === PI_PACKAGE_NAME) break; } catch { // Not a package boundary. } const parent = dirname(dir); dir = parent === dir ? undefined : parent; } piPackageRootCache = dir; return dir; } const NAME_PATTERN = /^[a-z0-9]+(?:-[a-z0-9]+)*$/u; export interface DynamicCatalogPaths { catalog: string; cache: string; runs: string; locks: string; } export function dynamicCatalogPaths(options: { agentDir: string; platform?: NodeJS.Platform; env?: NodeJS.ProcessEnv; home?: string; }): DynamicCatalogPaths { const platform = options.platform ?? process.platform; const env = options.env ?? process.env; const home = options.home ?? homedir(); const platformPath = platform === "win32" ? win32 : posix; let platformCache: string; if (platform === "win32") platformCache = env.LOCALAPPDATA ?? platformPath.join(home, "AppData", "Local"); else if (platform === "darwin") platformCache = platformPath.join(home, "Library", "Caches"); else platformCache = env.XDG_CACHE_HOME ?? platformPath.join(home, ".cache"); const cache = platformPath.join( platformCache, "pi", "ultra", "dynamic-extensions", ); return { catalog: platformPath.join( options.agentDir, "ultra", "dynamic-extensions", "catalog", ), cache, runs: platformPath.join(cache, "runs"), locks: platformPath.join(cache, "locks"), }; } interface CurrentRevision { name: string; revision: string; description: string; topic: string; publishedAt: string; } interface SearchResult { name: string; currentHash: string; integrityStatus: "intact" | "mismatch"; description: string; topic: string; tags: string[]; testedPlatforms: string[]; sourcePath: string; } interface BuilderState { query: string; candidates: string[]; ownedName?: string; } interface DraftMetadata { name: string; baseHash: string | null; createdAt: string; } export interface DynamicCoordinatorOptions { agentDir: string; parentSessionId: string; workflowRunId: string; availableTools: string[]; paths?: DynamicCatalogPaths; workerPath?: string; } export function createDynamicExtensionCoordinator( options: DynamicCoordinatorOptions, ): DynamicExtensionResolver { return new DynamicCoordinator(options); } class DynamicCoordinator implements DynamicExtensionResolver { readonly paths: DynamicCatalogPaths; readonly runRoot: string; readonly workerPath: string; private initialized?: Promise; constructor(private readonly options: DynamicCoordinatorOptions) { this.paths = options.paths ?? dynamicCatalogPaths(options); this.runRoot = join( this.paths.runs, safeSegment(options.parentSessionId), safeSegment(options.workflowRunId), ); this.workerPath = options.workerPath ?? defaultWorkerPath(); } private init(): Promise { if (!this.initialized) { this.initialized = (async () => { await Promise.all([ mkdir(this.paths.catalog, { recursive: true }), mkdir(this.paths.runs, { recursive: true }), mkdir(this.paths.locks, { recursive: true }), ]); })(); } return this.initialized; } async resolve( names: string[], args: { phase: string; agentId: string; signal?: AbortSignal }, ): Promise { await this.init(); const unique = new Set(); const resolved: ResolvedDynamicExtension[] = []; for (const name of names) { assertName(name); if (unique.has(name)) throw new Error( `ultra: dynamic extension "${name}" appears more than once in phase "${args.phase}".`, ); unique.add(name); resolved.push( await withNameLock(this.paths, name, args.signal, async () => { const current = await readCurrent(this.paths, name); if (!current) throw new Error( `ultra: dynamic extension "${name}" has no canonical revision. Run a builder with all four dynamic-extension tools.`, ); const sourcePath = join(this.paths.catalog, name, current.revision); const verified = verifyCurrent(sourcePath, current); if (!verified.ok) throw new Error( `ultra: dynamic extension "${name}" is ${verified.status} at ${sourcePath}; call validate_dynamic_extension on that exact path.`, ); const snapshot = join( this.runRoot, safeSegment(args.agentId), "snapshots", name, current.revision, ); if (!existsSync(snapshot)) { await mkdir(dirname(snapshot), { recursive: true }); await copyArtifacts(sourcePath, snapshot); const copiedHash = hashArtifacts(snapshot); if (copiedHash !== current.revision) { await rm(snapshot, { recursive: true, force: true }); throw new Error( `ultra: snapshot hash mismatch for "${name}": expected ${current.revision}, got ${copiedHash}.`, ); } } const manifest = await readManifest(sourcePath); return { name, description: manifest.description, revision: current.revision, selector: `${name}@${current.revision}`, canonicalPath: sourcePath, sourcePath: snapshot, loadPath: join(snapshot, "index.ts"), overrides: [...(manifest.overrides ?? [])], }; }), ); } return resolved; } async builder( agentId: string, signal?: AbortSignal, ): Promise<{ tools: ToolDefinition[]; prompt: string; activeTools: string[]; }> { await this.init(); const workspace = join(this.runRoot, safeSegment(agentId)); await mkdir(workspace, { recursive: true }); const tools = makeBuilderTools({ paths: this.paths, workspace, workerPath: this.workerPath, defaultSignal: signal, }); const prompt = await readFile( new URL("./prompts/dynamic-extension-builder.md", import.meta.url), "utf8", ); return { tools, prompt: prompt.trim(), activeTools: [ ...new Set([ ...this.options.availableTools.filter( (name) => name !== "run_workflow", ), ...DYNAMIC_EXTENSION_TOOL_NAMES, ]), ], }; } loaded( extensions: ResolvedDynamicExtension[], _sessionId: string, ): UsedDynamicExtension[] { return extensions.map((extension) => ({ name: extension.name, description: extension.description, selector: extension.selector, sourcePath: extension.sourcePath, debugLogPaths: [], })); } async complete(): Promise { for (const agent of await readdir(this.runRoot, { withFileTypes: true, }).catch(() => [])) { if (!agent.isDirectory()) continue; const agentRoot = join(this.runRoot, agent.name); await rm(join(agentRoot, "snapshots"), { recursive: true, force: true }); if ((await readdir(agentRoot).catch(() => [])).length === 0) await rm(agentRoot, { recursive: true, force: true }); } if ((await readdir(this.runRoot).catch(() => [])).length === 0) await rm(this.runRoot, { recursive: true, force: true }); } } function safeSegment(value: string): string { return value.replace(/[^A-Za-z0-9._-]/gu, "_"); } function assertName(name: string): void { if (!NAME_PATTERN.test(name)) throw new Error( `ultra: invalid dynamic extension name "${name}"; use a safe kebab-case path segment.`, ); } function lockRoot(paths: DynamicCatalogPaths): string { const catalogPath = resolve(paths.catalog).replaceAll("\\", "/"); const key = process.platform === "win32" ? catalogPath.toLowerCase() : catalogPath; const hash = createHash("sha256").update(key).digest("hex"); return join(paths.locks, hash); } async function withNameLock( paths: DynamicCatalogPaths, name: string, signal: AbortSignal | undefined, work: () => Promise, ): Promise { const lock = join(lockRoot(paths), `${name}.lock`); await mkdir(dirname(lock), { recursive: true }); const started = Date.now(); for (;;) { if (signal?.aborted) throw abortError(); try { await mkdir(lock); await writeFile( join(lock, "owner.json"), JSON.stringify({ pid: process.pid, acquiredAt: new Date().toISOString(), }), ); break; } catch (error) { if ((error as NodeJS.ErrnoException).code !== "EEXIST") throw error; const value = await stat(lock).catch(() => undefined); if ( value && Date.now() - value.mtimeMs > VALIDATOR_WATCHDOG_MS * 2 && !(await lockOwnerActive(lock)) ) { await rm(lock, { recursive: true, force: true }); continue; } if (Date.now() - started > VALIDATOR_WATCHDOG_MS) throw new Error( `ultra: timed out waiting for publication lock for ${name}.`, ); await abortableDelay(40, signal); } } try { return await work(); } finally { await rm(lock, { recursive: true, force: true }); } } async function lockOwnerActive(lock: string): Promise { try { const owner = JSON.parse( await readFile(join(lock, "owner.json"), "utf8"), ) as { pid?: unknown }; if (typeof owner.pid !== "number") return false; try { process.kill(owner.pid, 0); return true; } catch (error) { return (error as NodeJS.ErrnoException).code === "EPERM"; } } catch { return false; } } function abortError(): Error { const error = new Error("aborted"); error.name = "AbortError"; return error; } function abortableDelay(ms: number, signal?: AbortSignal): Promise { if (signal?.aborted) return Promise.reject(abortError()); return new Promise((resolveDelay, reject) => { const timer = setTimeout(resolveDelay, ms); const aborted = () => { clearTimeout(timer); reject(abortError()); }; signal?.addEventListener("abort", aborted, { once: true }); if (signal) setTimeout(() => signal.removeEventListener("abort", aborted), ms + 1); }); } async function readCurrent( paths: DynamicCatalogPaths, name: string, ): Promise { try { const current = JSON.parse( await readFile(join(paths.catalog, name, "current.json"), "utf8"), ) as CurrentRevision; return current.name === name && /^[a-f0-9]{64}$/u.test(current.revision) ? current : undefined; } catch { return undefined; } } async function readManifest(path: string): Promise { return JSON.parse(await readFile(join(path, "manifest.json"), "utf8")); } function verifyCurrent( sourcePath: string, current: CurrentRevision, ): { ok: boolean; status: "intact" | "mismatch" } { try { const ok = hashArtifacts(sourcePath) === current.revision; return { ok, status: ok ? "intact" : "mismatch" }; } catch { return { ok: false, status: "mismatch" }; } } async function copyArtifacts( source: string, destination: string, ): Promise { for (const file of artifactFiles(source)) { const sourceFile = resolve(file); const relativePath = relative(source, sourceFile); if (relativePath === ".ultra-validation.json") continue; const target = join(destination, relativePath); await mkdir(dirname(target), { recursive: true }); await copyFile(sourceFile, target); } } async function atomicWriteJson(path: string, value: unknown): Promise { await mkdir(dirname(path), { recursive: true }); const temporary = `${path}.${process.pid}.${randomUUID()}.tmp`; await writeFile(temporary, `${JSON.stringify(value, null, 2)}\n`, "utf8"); const handle = await open(temporary, "r+"); try { await handle.sync(); } finally { await handle.close(); } await rename(temporary, path); } interface BuilderToolOptions { paths: DynamicCatalogPaths; workspace: string; workerPath: string; defaultSignal?: AbortSignal; } function makeBuilderTools(options: BuilderToolOptions): ToolDefinition[] { const statePath = join(options.workspace, "builder-state.json"); const readState = async (): Promise => { try { return JSON.parse(await readFile(statePath, "utf8")); } catch { return undefined; } }; const assertOwnership = async ( name: string, ): Promise => { const state = await readState(); if (state?.ownedName && state.ownedName !== name) throw new Error( `this builder owns dynamic extension "${state.ownedName}" and cannot modify "${name}".`, ); return state; }; const result = (details: unknown) => ({ content: [ { type: "text" as const, text: bounded(JSON.stringify(details, null, 2)), }, ], details, }); return [ { name: "search_dynamic_extensions", label: "Search dynamic extensions", description: "Search the validated shared dynamic-extension catalog before creating or improving one. Returns bounded ranked matches with names, revisions, descriptions, topics, tags, and paths.", parameters: Type.Object( { query: Type.String({ minLength: 1, description: "Capability, topic, or keywords to match against published catalog entries.", }), }, { additionalProperties: false }, ), async execute(_id, params) { const input = params as { query: string }; const matches = await searchCatalog(options.paths, input.query); const previous = await readState(); await atomicWriteJson(statePath, { query: input.query, candidates: matches.map((match) => match.name), ...(previous?.ownedName ? { ownedName: previous.ownedName } : {}), } satisfies BuilderState); return result({ matches }); }, }, { name: "create_dynamic_extension", label: "Create dynamic extension", description: "Scaffold a new extension in a unique editable working copy. Existing copies are preserved; catalog ownership and safe names still apply.", parameters: Type.Object( { name: Type.String({ minLength: 1, description: "Stable kebab-case dynamic-extension name.", }), description: Type.String({ minLength: 1, description: "Human-readable capability description.", }), topic: Type.Optional( Type.String({ description: "Optional search topic." }), ), scope: Type.Optional( Type.String({ description: "Optional scope notes." }), ), nonGoals: Type.Optional(Type.Array(Type.String())), tags: Type.Optional(Type.Array(Type.String())), consideredCandidates: Type.Optional( Type.Array( Type.Object( { name: Type.String({ minLength: 1, description: "Rejected catalog candidate name.", }), rationale: Type.String({ minLength: 1, description: "Concrete reason the candidate cannot satisfy the requested capability.", }), }, { additionalProperties: false }, ), { description: "Optional catalog alternatives considered.", }, ), ), noMatchRationale: Type.Optional( Type.String({ minLength: 1, description: "Optional reason for creating a new extension.", }), ), }, { additionalProperties: false }, ), async execute(_id, params) { const input = params as CreateParams; assertName(input.name); const searched = (await assertOwnership(input.name)) ?? { query: input.name, candidates: [], }; if (await readCurrent(options.paths, input.name)) throw new Error( `dynamic extension "${input.name}" already exists; copy and improve it.`, ); const draft = await newWorkingCopy(options.workspace, input.name); await scaffoldDraft(draft, input); await atomicWriteJson(join(dirname(draft), "draft.json"), { name: input.name, baseHash: null, createdAt: new Date().toISOString(), } satisfies DraftMetadata); await atomicWriteJson(statePath, { ...searched, ownedName: input.name, } satisfies BuilderState); return result({ name: input.name, workingCopyPath: draft }); }, }, { name: "copy_dynamic_extension", label: "Copy dynamic extension", description: "Copy a canonical extension into a unique editable working copy without replacing earlier copies. Returns its stable name, base revision hash, and workingCopyPath.", parameters: Type.Object( { name: Type.String({ minLength: 1, description: "Stable name of the published catalog extension to copy.", }), }, { additionalProperties: false }, ), async execute(_id, params, signal) { const input = params as { name: string }; assertName(input.name); const state = await assertOwnership(input.name); const effectiveSignal = signal ?? options.defaultSignal; const checkout = await withNameLock( options.paths, input.name, effectiveSignal, async () => { const current = await readCurrent(options.paths, input.name); if (!current) throw new Error(`unknown dynamic extension: ${input.name}`); const source = join( options.paths.catalog, input.name, current.revision, ); const verified = verifyCurrent(source, current); if (!verified.ok) throw new Error( `canonical source is ${verified.status} at ${source}; validate it before copying.`, ); const draft = await newWorkingCopy(options.workspace, input.name); await copyArtifacts(source, draft); await atomicWriteJson(join(dirname(draft), "draft.json"), { name: input.name, baseHash: current.revision, createdAt: new Date().toISOString(), } satisfies DraftMetadata); return { current, draft }; }, ); await atomicWriteJson(statePath, { query: state?.query ?? input.name, candidates: state?.candidates ?? [input.name], ownedName: input.name, } satisfies BuilderState); return result({ name: input.name, baseHash: checkout.current.revision, workingCopyPath: checkout.draft, }); }, }, { name: "validate_dynamic_extension", label: "Validate dynamic extension", description: "Check structural safety and bounded Pi load/tests. Publish eligible working copies with advisory warnings. Returns bounded publication handoff; full diagnostics remain in details.", parameters: Type.Object( { path: Type.String({ minLength: 1, description: "Absolute or current-working-directory-relative path to the editable working-copy source.", }), }, { additionalProperties: false }, ), async execute(_id, params, signal) { const input = params as { path: string }; const manifest = await readManifest(resolve(input.path)).catch( () => undefined, ); if (manifest) await assertOwnership(manifest.name); if (manifest) { const state = await readState(); await atomicWriteJson(statePath, { query: state?.query ?? manifest.topic ?? manifest.name, candidates: state?.candidates ?? [manifest.name], ownedName: manifest.name, } satisfies BuilderState); } const validation = await validateAndPublish({ paths: options.paths, draftPath: resolve(input.path), workerPath: options.workerPath, signal: signal ?? options.defaultSignal, }); return { content: [ { type: "text" as const, text: bounded( JSON.stringify( { status: validation.status, name: validation.name ?? manifest?.name, workingCopyPath: resolve(input.path), revision: validation.revision, diagnostics: ( validation.diagnostics as ValidationDiagnostic[] ) .slice(0, 8) .map(({ code, status, remediation }) => ({ code, status, nextAction: remediation, })), diagnosticCount: ( validation.diagnostics as ValidationDiagnostic[] ).length, }, null, 2, ), ), }, ], details: validation, }; }, }, ] as ToolDefinition[]; } function bounded(value: string): string { return value.length <= MAX_TOOL_OUTPUT ? value : `${value.slice(0, MAX_TOOL_OUTPUT)}\n[truncated]`; } function normalizeTokens(value: string): string[] { return value .replace(/([a-z0-9])([A-Z])/gu, "$1 $2") .toLowerCase() .split(/[^a-z0-9]+/gu) .filter(Boolean) .map((token) => token.length > 4 && token.endsWith("ies") ? `${token.slice(0, -3)}y` : token.length > 3 && token.endsWith("s") ? token.slice(0, -1) : token, ); } function readmeIndexText(readme: string): string { const wanted = new Set(["purpose", "scope", "usage"]); const lines: string[] = []; let include = false; let inCode = false; for (const line of readme.split(/\r?\n/u)) { if (line.startsWith("```")) { inCode = !inCode; continue; } if (inCode) continue; const heading = /^##\s+(.+)$/u.exec(line); if (heading) { include = wanted.has(heading[1].trim().toLowerCase()); continue; } if (include) lines.push(line); } return lines.join(" "); } async function searchCatalog( paths: DynamicCatalogPaths, query: string, ): Promise { await mkdir(paths.catalog, { recursive: true }); const queryPhrase = query.toLowerCase().trim(); const queryTokens = normalizeTokens(query); const ranked: Array<{ score: number; result: SearchResult }> = []; for (const entry of await readdir(paths.catalog, { withFileTypes: true })) { if (!entry.isDirectory()) continue; const current = await readCurrent(paths, entry.name); if (!current) continue; const sourcePath = join(paths.catalog, entry.name, current.revision); let manifest: DynamicManifest; try { manifest = await readManifest(sourcePath); } catch { continue; } const readme = await readFile(join(sourcePath, "README.md"), "utf8").catch( () => "", ); const verified = verifyCurrent(sourcePath, current); const attestation = await readAttestation(sourcePath); const nameTopic = `${manifest.name} ${manifest.topic ?? ""} ${manifest.description}`; const tags = (manifest.tags ?? []).join(" "); const scope = manifest.scope ?? ""; const prose = readmeIndexText(readme); const fields = [nameTopic, tags, scope, prose]; let score = fields.some((field) => field.toLowerCase().includes(queryPhrase), ) ? 10_000 : 0; const weights = [100, 40, 15, 4]; for (const [index, field] of fields.entries()) { const tokens = new Set(normalizeTokens(field)); for (const token of queryTokens) if (tokens.has(token)) score += weights[index]; } if (score === 0) continue; ranked.push({ score, result: { name: manifest.name, currentHash: current.revision, integrityStatus: verified.status, description: manifest.description, topic: manifest.topic ?? "", tags: manifest.tags ?? [], testedPlatforms: attestation?.testedPlatforms ?? (attestation ? [attestation.platform] : []), sourcePath, }, }); } return ranked .sort( (left, right) => right.score - left.score || left.result.name.localeCompare(right.result.name), ) .slice(0, 5) .map(({ result }) => result); } async function readAttestation( sourcePath: string, ): Promise { try { return JSON.parse( await readFile(join(sourcePath, ".ultra-validation.json"), "utf8"), ); } catch { return undefined; } } interface CreateParams { name: string; description: string; topic?: string; scope?: string; nonGoals?: string[]; tags?: string[]; consideredCandidates?: Array<{ name: string; rationale: string }>; noMatchRationale?: string; } async function newWorkingCopy( workspace: string, name: string, ): Promise { const parent = join(workspace, "drafts", name); await mkdir(parent, { recursive: true }); const root = join(parent, randomUUID()); await mkdir(root); return join(root, "source"); } async function scaffoldDraft( draft: string, params: CreateParams, ): Promise { await mkdir(draft, { recursive: true }); const toolName = params.name.replaceAll("-", "_"); const manifest: DynamicManifest = { name: params.name, description: params.description, ...(params.topic ? { topic: params.topic } : {}), ...(params.scope ? { scope: params.scope } : {}), ...(params.nonGoals ? { nonGoals: params.nonGoals } : {}), ...(params.tags ? { tags: [...new Set(params.tags)] } : {}), }; const alternatives = (params.consideredCandidates ?? []) .map((candidate) => `- **${candidate.name}**: ${candidate.rationale}`) .join("\n") || params.noMatchRationale || ""; const files: Record = { "manifest.json": `${JSON.stringify(manifest, null, 2)}\n`, "package.json": `${JSON.stringify( { name: `@ultra-dynamic/${params.name}`, private: true, type: "module", main: "index.ts", }, null, 2, )}\n`, "README.md": `# ${params.name}\n\n${params.description}\n\n${params.scope ?? ""}${alternatives ? `\n\n${alternatives}` : ""}\n\nLoad via \`dynamicExtensions\`; review validator diagnostics before consuming.\nHuman promotion required for a repository extension.\n`, "core.ts": `export function executeCapability(value: string): { value: string } {\n\treturn { value };\n}\n`, "index.ts": `import type { ExtensionAPI } from "@earendil-works/pi-coding-agent";\nimport { executeCapability } from "./core.ts";\n\nexport default function dynamicExtension(pi: ExtensionAPI) {\n\tpi.registerTool({\n\t\tname: "${toolName}",\n\t\tlabel: "${params.name}",\n\t\tdescription: ${JSON.stringify(params.description)},\n\t\tparameters: {\n\t\t\ttype: "object",\n\t\t\tproperties: { value: { type: "string" } },\n\t\t\trequired: ["value"],\n\t\t\tadditionalProperties: false,\n\t\t},\n\t\tasync execute(_id, params) {\n\t\t\tconst result = executeCapability(params.value);\n\t\t\treturn { content: [{ type: "text", text: JSON.stringify(result) }], details: result };\n\t\t},\n\t});\n\tpi.on("session_start", () => {\n\t\tconst active = pi.getActiveTools();\n\t\tif (!active.includes("${toolName}")) pi.setActiveTools([...active, "${toolName}"]);\n\t});\n}\n`, }; await Promise.all( Object.entries(files).map(([name, content]) => writeFile(join(draft, name), content, "utf8"), ), ); } async function readDraftMetadata( draftPath: string, ): Promise { try { return JSON.parse( await readFile(join(dirname(draftPath), "draft.json"), "utf8"), ); } catch { return undefined; } } type ValidatorResult = WorkerValidationResult & { aborted?: boolean }; function parseValidatorResult(text: string): ValidatorResult { const value: unknown = JSON.parse(text); if ( !value || typeof value !== "object" || Array.isArray(value) || typeof (value as { ok?: unknown }).ok !== "boolean" || !Array.isArray((value as { diagnostics?: unknown }).diagnostics) ) throw new Error("validator result has an invalid shape"); return value as ValidatorResult; } async function readValidatorResult( path: string, ): Promise { try { return parseValidatorResult(await readFile(path, "utf8")); } catch { return undefined; } } function validatorProtocolFailure( draftPath: string, output: string, actual?: string, ): ValidatorResult { return failedValidation(draftPath, { gate: "validator", code: "validator_protocol", status: "error", ...(actual ? { actual } : {}), output: bounded(output), remediation: "Inspect validator output and repair the worker protocol failure.", workingCopyPath: draftPath, }); } async function runValidator(args: { workerPath: string; draftPath: string; signal?: AbortSignal; }): Promise { if (args.signal?.aborted) return abortedValidation( args.draftPath, "validation aborted before launch", ); return new Promise((resolveResult) => { const resultPath = join( dirname(args.draftPath), `.validation-result-${randomUUID()}.json`, ); const preflightPath = `${resultPath}.preflight`; const child = spawn(process.execPath, [args.workerPath, args.draftPath], { cwd: args.draftPath, env: { ...process.env, ...(piPackageRoot() ? { ULTRA_PI_PACKAGE_ROOT: piPackageRoot() } : {}), CI: "true", ULTRA_VALIDATOR_TEMP: dirname(args.draftPath), ULTRA_VALIDATOR_RESULT_PATH: resultPath, }, detached: process.platform !== "win32", stdio: ["ignore", "pipe", "pipe"], }); let output = ""; const append = (chunk: Buffer): void => { output = `${output}${chunk.toString("utf8")}`.slice(-64_000); }; child.stdout.on("data", append); child.stderr.on("data", append); let finished = false; let stopping = false; const finish = (value: ValidatorResult): void => { if (finished) return; finished = true; clearTimeout(timer); args.signal?.removeEventListener("abort", onAbort); void rm(resultPath, { force: true }).catch(() => {}); void rm(preflightPath, { force: true }).catch(() => {}); resolveResult(value); }; const stop = async (reason: "aborted" | "timeout"): Promise => { if (finished || stopping) return; stopping = true; await killProcessTree(child.pid); if (reason === "aborted") { finish(abortedValidation(args.draftPath, "validation aborted")); return; } const preflight = await readValidatorResult(preflightPath); const warning: ValidationDiagnostic = { gate: "validator", code: "validator_timeout", status: "warning", limitMs: VALIDATOR_WATCHDOG_MS, output: bounded(output), remediation: "Repair hanging Pi load, tests, or benchmark before consuming this extension.", workingCopyPath: args.draftPath, }; finish( preflight?.ok && preflight.hash && preflight.manifest && preflight.attestation ? { ...preflight, diagnostics: [...preflight.diagnostics, warning] } : failedValidation(args.draftPath, { ...warning, status: "error" }), ); }; const timer = setTimeout(() => void stop("timeout"), VALIDATOR_WATCHDOG_MS); const onAbort = () => void stop("aborted"); args.signal?.addEventListener("abort", onAbort, { once: true }); child.on("error", (error) => finish( failedValidation(args.draftPath, { gate: "validator", code: "validator_spawn", status: "error", actual: error.message, remediation: "Ensure the current Node executable can launch the bundled validator.", workingCopyPath: args.draftPath, }), ), ); child.on("close", () => { if (stopping) return; void (async () => { const fileResult = await readValidatorResult(resultPath); if (fileResult) { finish(fileResult); return; } const marker = "ULTRA_VALIDATION_RESULT="; const index = output.lastIndexOf(marker); if (index < 0) { finish(validatorProtocolFailure(args.draftPath, output)); return; } try { const line = output.slice(index + marker.length).split(/\r?\n/u)[0]; finish(parseValidatorResult(line)); } catch (error) { finish( validatorProtocolFailure( args.draftPath, output, error instanceof Error ? error.message : String(error), ), ); } })(); }); }); } async function killProcessTree(pid: number | undefined): Promise { if (!pid) return; if (process.platform === "win32") { await new Promise((resolveKill) => { const killer = spawn("taskkill", ["/pid", String(pid), "/T", "/F"], { stdio: "ignore", windowsHide: true, }); killer.once("error", () => resolveKill()); killer.once("close", () => resolveKill()); }); return; } try { process.kill(-pid, "SIGKILL"); } catch { try { process.kill(pid, "SIGKILL"); } catch {} } } function failedValidation( _draftPath: string, diagnostic: ValidationDiagnostic, ): WorkerValidationResult { return { ok: false, diagnostics: [diagnostic], durationMs: 0 }; } function abortedValidation( draftPath: string, message: string, ): WorkerValidationResult & { aborted: true } { return { ok: false, aborted: true, diagnostics: [ { gate: "validator", code: "validator_aborted", status: "error", actual: message, remediation: "Resume the retained working copy and validate again.", workingCopyPath: draftPath, }, ], durationMs: 0, }; } function publicationFailure( validation: WorkerValidationResult, draftPath: string, code: string, expected: unknown, actual: unknown, remediation: string, ): Record { return { status: "failed", diagnostics: [ ...validation.diagnostics, { gate: "publication", code, status: "error", expected, actual, path: join(draftPath, "manifest.json"), remediation, workingCopyPath: draftPath, }, ], workingCopyPath: draftPath, }; } async function validateAndPublish(args: { paths: DynamicCatalogPaths; draftPath: string; workerPath: string; signal?: AbortSignal; }): Promise> { await access(args.draftPath); const validation = await runValidator(args); if ( !validation.ok || !validation.manifest || !validation.hash || !validation.attestation ) return { status: validation.aborted ? "aborted" : "failed", diagnostics: validation.diagnostics, workingCopyPath: args.draftPath, }; const manifest = validation.manifest; const revision = validation.hash; let attestation = validation.attestation; if ( !NAME_PATTERN.test(manifest.name) || attestation.sourceHash !== revision || !(await readManifest(args.draftPath) .then((value) => value.name === manifest.name) .catch(() => false)) || (() => { try { return hashArtifacts(args.draftPath) === revision; } catch { return false; } })() === false ) return publicationFailure( validation, args.draftPath, "validation_identity_or_hash", revision, manifest.name, "Repair manifest identity, artifact symlinks, or changed files; validate again.", ); const metadata = await readDraftMetadata(args.draftPath); const catalogRelative = relative(args.paths.catalog, args.draftPath); const catalogParts = catalogRelative.split(sep); const canonicalName = !catalogRelative.startsWith("..") && catalogParts.length === 2 ? catalogParts[0] : undefined; const identityName = metadata?.name ?? canonicalName; if (identityName && identityName !== manifest.name) return publicationFailure( validation, args.draftPath, "extension_identity_changed", identityName, manifest.name, "Restore the working copy's original name; create a new working copy for a new identity.", ); return withNameLock(args.paths, manifest.name, args.signal, async () => { const current = await readCurrent(args.paths, manifest.name); const currentSource = current ? join(args.paths.catalog, manifest.name, current.revision) : undefined; const manualCurrent = currentSource && resolve(currentSource) === resolve(args.draftPath); const baseHash = manualCurrent ? current?.revision : metadata?.baseHash; if ((current?.revision ?? null) !== (baseHash ?? null)) { const diagnostic: ValidationDiagnostic = { gate: "publication", code: "stale_base", status: "error", expected: current?.revision ?? null, actual: baseHash ?? null, path: currentSource, remediation: "Copy the current published extension into another working copy; reconcile changes and validate again.", workingCopyPath: args.draftPath, }; return { status: "stale", name: manifest.name, baseHash: baseHash ?? null, currentHash: current?.revision ?? null, currentPath: currentSource, diagnostics: [...validation.diagnostics, diagnostic], workingCopyPath: args.draftPath, }; } if (currentSource && current) { const oldAttestation = await readAttestation(currentSource); attestation = { ...attestation, testedPlatforms: [ ...new Set([ ...(oldAttestation?.testedPlatforms ?? (oldAttestation ? [oldAttestation.platform] : [])), ...attestation.testedPlatforms, ]), ], }; } const nameRoot = join(args.paths.catalog, manifest.name); const target = join(nameRoot, revision); if (!existsSync(target)) { const staging = `${target}.tmp-${randomUUID()}`; await copyArtifacts(args.draftPath, staging); if (hashArtifacts(staging) !== revision) { await rm(staging, { recursive: true, force: true }); return publicationFailure( validation, args.draftPath, "artifact_hash_changed", revision, null, "Working copy changed during publication; validate again.", ); } await writeFile( join(staging, ".ultra-validation.json"), `${JSON.stringify(attestation, null, 2)}\n`, "utf8", ); await mkdir(nameRoot, { recursive: true }); await rename(staging, target); } else { if (hashArtifacts(target) !== revision) return publicationFailure( validation, args.draftPath, "catalog_hash_mismatch", revision, null, "Canonical revision is corrupted; do not overwrite it.", ); await atomicWriteJson( join(target, ".ultra-validation.json"), attestation, ); } await atomicWriteJson(join(nameRoot, "current.json"), { name: manifest.name, revision, description: manifest.description, topic: manifest.topic ?? "", publishedAt: new Date().toISOString(), } satisfies CurrentRevision); if (metadata && !manualCurrent) { try { await atomicWriteJson(join(dirname(args.draftPath), "draft.json"), { ...metadata, baseHash: revision, } satisfies DraftMetadata); } catch (error) { return { status: "validated", name: manifest.name, revision, sourcePath: target, workingCopyPath: args.draftPath, diagnostics: [ ...validation.diagnostics, { gate: "publication", code: "working_copy_bookkeeping", status: "warning", actual: error instanceof Error ? error.message : String(error), remediation: "Publication succeeded; copy the current revision before another edit because this working copy may retain its old base.", workingCopyPath: args.draftPath, }, ], }; } } return { status: "validated", name: manifest.name, revision, description: manifest.description, sourcePath: target, workingCopyPath: args.draftPath, diagnostics: validation.diagnostics, }; }); }