Luigit
repositories / pi-ext

pi-ext

bugabingas pi extensions

owned by admin

extensions/ultra/dynamic-catalog.ts

Raw
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<void>;

	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<void> {
		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<ResolvedDynamicExtension[]> {
		await this.init();
		const unique = new Set<string>();
		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<void> {
		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<T>(
	paths: DynamicCatalogPaths,
	name: string,
	signal: AbortSignal | undefined,
	work: () => Promise<T>,
): Promise<T> {
	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<boolean> {
	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<void> {
	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<CurrentRevision | undefined> {
	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<DynamicManifest> {
	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<void> {
	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<void> {
	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<BuilderState | undefined> => {
		try {
			return JSON.parse(await readFile(statePath, "utf8"));
		} catch {
			return undefined;
		}
	};
	const assertOwnership = async (
		name: string,
	): Promise<BuilderState | undefined> => {
		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<SearchResult[]> {
	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<ValidationAttestation | undefined> {
	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<string> {
	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<void> {
	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<string, string> = {
		"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<DraftMetadata | undefined> {
	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<ValidatorResult | undefined> {
	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<ValidatorResult> {
	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<void> => {
			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<void> {
	if (!pid) return;
	if (process.platform === "win32") {
		await new Promise<void>((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<string, unknown> {
	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<Record<string, unknown>> {
	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,
		};
	});
}