Luigit
repositories / pi-ext

pi-ext

bugabingas pi extensions

owned by admin

extensions/ultra/dynamic-validator-worker.ts

Raw
import { once } from "node:events";
import { lstatSync, readFileSync } from "node:fs";
import { mkdtemp, rm, writeFile } from "node:fs/promises";
import { registerHooks } from "node:module";
import { tmpdir } from "node:os";
import { join, resolve } from "node:path";
import { fileURLToPath, pathToFileURL } from "node:url";
import {
	artifactFiles,
	hashArtifacts as hashArtifactTree,
} from "./dynamic-artifacts.ts";
import {
	BUILTIN_TOOL_NAMES,
	type DynamicManifest,
	ULTRA_RESERVED_TOOL_NAMES,
	VALIDATOR_WARNING_TARGET_MS,
	type ValidationAttestation,
	type ValidationDiagnostic,
	type WorkerValidationResult,
} from "./dynamic-contract.ts";

export type {
	DynamicManifest,
	ValidationAttestation,
	ValidationDiagnostic,
	WorkerValidationResult,
} from "./dynamic-contract.ts";
export { VALIDATOR_WARNING_TARGET_MS } from "./dynamic-contract.ts";

interface ValidationState {
	draftPath: string;
	diagnostics: ValidationDiagnostic[];
	manifest?: DynamicManifest;
}

type PiCodingAgent = typeof import("@earendil-works/pi-coding-agent");
type PiAiCompat = typeof import("@earendil-works/pi-ai/compat");

const PI_SPECIFIER = /^(?:@earendil-works\/|typebox(?:\/|$))/u;

/**
 * A generated package has no `node_modules` ancestor, so the parent passes the
 * running Pi installation and this worker resolves Pi modules from there.
 * Static imports would hoist above the hook, hence dynamic imports.
 */
function anchorPiModules(packageRoot: string | undefined): void {
	if (!packageRoot) return;
	const here = new URL("./", import.meta.url).href;
	const anchor = pathToFileURL(join(packageRoot, "package.json")).href;
	registerHooks({
		resolve(specifier, context, next) {
			if (PI_SPECIFIER.test(specifier) && context.parentURL?.startsWith(here))
				return next(specifier, { ...context, parentURL: anchor });
			return next(specifier, context);
		},
	});
}

let piModulesPromise: Promise<[PiCodingAgent, PiAiCompat]> | undefined;

function piModules(): Promise<[PiCodingAgent, PiAiCompat]> {
	if (!piModulesPromise) {
		anchorPiModules(process.env.ULTRA_PI_PACKAGE_ROOT);
		piModulesPromise = Promise.all([
			import("@earendil-works/pi-coding-agent"),
			import("@earendil-works/pi-ai/compat"),
		]);
	}
	return piModulesPromise;
}

function diagnostic(
	state: ValidationState,
	value: Omit<ValidationDiagnostic, "workingCopyPath">,
): void {
	state.diagnostics.push({ ...value, workingCopyPath: state.draftPath });
}

function validManifest(value: unknown): value is DynamicManifest {
	if (!value || typeof value !== "object" || Array.isArray(value)) return false;
	const manifest = value as Record<string, unknown>;
	return (
		typeof manifest.name === "string" &&
		/^[a-z0-9]+(?:-[a-z0-9]+)*$/u.test(manifest.name) &&
		typeof manifest.description === "string" &&
		manifest.description.trim().length > 0 &&
		(manifest.overrides === undefined ||
			(Array.isArray(manifest.overrides) &&
				manifest.overrides.every((item) => typeof item === "string"))) &&
		["topic", "scope"].every(
			(field) =>
				manifest[field] === undefined || typeof manifest[field] === "string",
		) &&
		[
			"tags",
			"nonGoals",
			"capabilities",
			"requiredExecutables",
			"optionalExecutables",
		].every(
			(field) =>
				manifest[field] === undefined ||
				(Array.isArray(manifest[field]) &&
					(manifest[field] as unknown[]).every(
						(item) => typeof item === "string",
					)),
		)
	);
}

function checkManifest(state: ValidationState): void {
	const path = join(state.draftPath, "manifest.json");
	try {
		const value: unknown = JSON.parse(readFileSync(path, "utf8"));
		if (!validManifest(value))
			throw new Error(
				"Expected a safe kebab-case name, nonempty description, and optional string[] overrides.",
			);
		state.manifest = value;
		for (const override of value.overrides ?? []) {
			if (
				!/^(tool|command|flag|shortcut|renderer|entry-renderer|provider):[^:]+$/u.test(
					override,
				) ||
				(override.startsWith("tool:") &&
					ULTRA_RESERVED_TOOL_NAMES.has(override.slice(5)))
			)
				diagnostic(state, {
					gate: "manifest",
					code: "manifest_override",
					status: "error",
					path,
					actual: override,
					remediation:
						"Use a qualified override; Ultra orchestration tools cannot be overridden.",
				});
		}
	} catch (error) {
		diagnostic(state, {
			gate: "manifest",
			code: "manifest_schema",
			status: "error",
			path,
			actual: String(error),
			remediation:
				"Provide a readable manifest with a safe name and description.",
		});
	}
}

function checkArtifacts(state: ValidationState): boolean {
	let safe = true;
	try {
		for (const path of artifactFiles(state.draftPath)) {
			if (lstatSync(path).isSymbolicLink()) {
				safe = false;
				diagnostic(state, {
					gate: "artifacts",
					code: "artifact_symlink",
					status: "error",
					path,
					remediation: "Replace symlinks with regular files.",
				});
			}
		}
		const path = join(state.draftPath, "index.ts");
		if (!lstatSync(path).isFile())
			throw new Error("index.ts is not a regular file");
		readFileSync(path);
	} catch (error) {
		safe = false;
		diagnostic(state, {
			gate: "artifacts",
			code: "main_entry",
			status: "error",
			path: join(state.draftPath, "index.ts"),
			actual: String(error),
			remediation:
				"Provide a readable regular index.ts and readable artifact tree.",
		});
	}
	return safe;
}

async function checkLoad(state: ValidationState): Promise<void> {
	const indexPath = join(state.draftPath, "index.ts");
	const agentDir = await mkdtemp(join(tmpdir(), "ultra-validator-agent-"));
	let session:
		| Awaited<ReturnType<PiCodingAgent["createAgentSession"]>>["session"]
		| undefined;
	try {
		const [
			{ createAgentSession, DefaultResourceLoader, SessionManager },
			{ getModels },
		] = await piModules();
		const loader = new DefaultResourceLoader({
			cwd: state.draftPath,
			agentDir,
			noExtensions: true,
			noSkills: true,
			noPromptTemplates: true,
			noThemes: true,
			noContextFiles: true,
			additionalExtensionPaths: [indexPath],
		});
		await loader.reload();
		const loaded = loader.getExtensions();
		if (loaded.errors.length > 0) {
			diagnostic(state, {
				gate: "real-loader",
				code: "extension_load",
				status: "warning",
				path: indexPath,
				actual: loaded.errors.map((error) => String(error)),
				remediation: "Repair Pi load errors before using this extension.",
			});
			return;
		}
		const generatedTools = loaded.extensions.flatMap((extension) => [
			...extension.tools.keys(),
		]);
		for (const name of generatedTools) {
			if (
				ULTRA_RESERVED_TOOL_NAMES.has(name) ||
				(BUILTIN_TOOL_NAMES.has(name) &&
					!(state.manifest?.overrides ?? []).includes(`tool:${name}`))
			)
				diagnostic(state, {
					gate: "real-loader",
					code: ULTRA_RESERVED_TOOL_NAMES.has(name)
						? "generated_tool_reserved"
						: "generated_tool_override_undeclared",
					status: "error",
					path: indexPath,
					actual: name,
					remediation:
						"Rename the tool or declare the exact allowed built-in override.",
				});
		}
		const model = getModels("openai")[0] ?? getModels("anthropic")[0];
		if (!model)
			throw new Error("Pi has no built-in model for an in-memory session.");
		({ session } = await createAgentSession({
			cwd: state.draftPath,
			agentDir,
			model,
			resourceLoader: loader,
			sessionManager: SessionManager.inMemory(state.draftPath),
			noTools: "builtin",
		}));
		const errors: unknown[] = [];
		await session.bindExtensions({
			mode: "print",
			onError: (error) => errors.push(String(error)),
		});
		if (errors.length > 0)
			diagnostic(state, {
				gate: "real-loader",
				code: "extension_session",
				status: "warning",
				path: indexPath,
				actual: errors,
				remediation: "Repair print-mode session errors before use.",
			});
		// Tools may be registered during session_start rather than at module load.
		for (const name of loaded.extensions
			.flatMap((extension) => [...extension.tools.keys()])
			.filter((name) => !generatedTools.includes(name))) {
			if (
				ULTRA_RESERVED_TOOL_NAMES.has(name) ||
				(BUILTIN_TOOL_NAMES.has(name) &&
					!(state.manifest?.overrides ?? []).includes(`tool:${name}`))
			)
				diagnostic(state, {
					gate: "real-loader",
					code: ULTRA_RESERVED_TOOL_NAMES.has(name)
						? "generated_tool_reserved"
						: "generated_tool_override_undeclared",
					status: "error",
					path: indexPath,
					actual: name,
					remediation:
						"Rename the tool or declare the exact allowed built-in override.",
				});
		}
		await session.extensionRunner.emit({
			type: "session_shutdown",
			reason: "quit",
		});
	} catch (error) {
		diagnostic(state, {
			gate: "real-loader",
			code: "real_loader_session",
			status: "warning",
			path: indexPath,
			actual:
				error instanceof Error ? (error.stack ?? error.message) : String(error),
			remediation:
				"Repair Pi loading/session behavior before consuming this extension.",
		});
	} finally {
		session?.dispose();
		await rm(agentDir, { recursive: true, force: true });
	}
}

async function checkTests(state: ValidationState): Promise<void> {
	const tests = artifactFiles(state.draftPath).filter((path) =>
		/\.test\.[cm]?[jt]s$/u.test(path),
	);
	if (tests.length === 0) return;
	try {
		const { run } = await import("node:test");
		const failures: unknown[] = [];
		let success = false;
		const stream = run({ files: tests, concurrency: false, isolation: "none" });
		stream.on("test:fail", (failure) => failures.push(failure));
		stream.on("test:summary", (summary) => {
			success = summary.success;
		});
		stream.resume();
		await once(stream, "end");
		if (!success)
			diagnostic(state, {
				gate: "tests",
				code: "tests_failed",
				status: "warning",
				output: JSON.stringify(failures).slice(0, 8_000),
				remediation: "Repair failing tests before consuming this extension.",
			});
	} catch (error) {
		diagnostic(state, {
			gate: "tests",
			code: "tests_failed",
			status: "warning",
			actual: String(error),
			remediation: "Repair test runner errors before consuming this extension.",
		});
	}
}

export function hashArtifacts(root: string): string {
	return hashArtifactTree(root);
}

export async function validateDraft(
	draftPath: string,
	onPreflight?: (result: WorkerValidationResult) => Promise<void>,
): Promise<WorkerValidationResult> {
	const started = performance.now();
	const state: ValidationState = {
		draftPath: resolve(draftPath),
		diagnostics: [],
	};
	checkManifest(state);
	const safe = checkArtifacts(state);
	let hash: string | undefined;
	try {
		hash = hashArtifacts(state.draftPath);
	} catch (error) {
		diagnostic(state, {
			gate: "hash",
			code: "artifact_hash",
			status: "error",
			actual: String(error),
			remediation: "Remove symlinks and make all artifacts readable.",
		});
	}
	if (onPreflight) await onPreflight(resultFor(state, hash, started, true));
	if (safe && hash && state.manifest) {
		try {
			await checkLoad(state);
			await checkTests(state);
		} catch (error) {
			diagnostic(state, {
				gate: "advisory",
				code: "advisory_check_failed",
				status: "warning",
				actual:
					error instanceof Error
						? (error.stack ?? error.message)
						: String(error),
				remediation:
					"Repair feedback infrastructure or run the checks manually before consuming.",
			});
		}
	}
	try {
		if (hashArtifacts(state.draftPath) !== hash) {
			hash = undefined;
			diagnostic(state, {
				gate: "hash",
				code: "artifact_changed",
				status: "error",
				remediation: "Working copy changed during validation; retry.",
			});
		}
	} catch (error) {
		hash = undefined;
		diagnostic(state, {
			gate: "hash",
			code: "artifact_hash",
			status: "error",
			actual: String(error),
			remediation: "Make artifacts readable and retry.",
		});
	}
	return resultFor(state, hash, started);
}

function resultFor(
	state: ValidationState,
	hash: string | undefined,
	started: number,
	preflight = false,
): WorkerValidationResult {
	const durationMs = performance.now() - started;
	if (
		durationMs > VALIDATOR_WARNING_TARGET_MS &&
		!state.diagnostics.some((item) => item.code === "validator_warning_target")
	)
		diagnostic(state, {
			gate: "validator-runtime",
			code: "validator_warning_target",
			status: "warning",
			limitMs: VALIDATOR_WARNING_TARGET_MS,
			elapsedMs: durationMs,
			remediation:
				"Keep feedback fast; this warning does not block publication.",
		});
	const ok =
		!!state.manifest &&
		!!hash &&
		!state.diagnostics.some((item) => item.status === "error");
	const attestation: ValidationAttestation | undefined =
		ok && hash
			? {
					sourceHash: hash,
					validatedAt: new Date().toISOString(),
					platform: process.platform,
					testedPlatforms: [process.platform],
					nodeVersion: process.version,
					warningTargetMs: VALIDATOR_WARNING_TARGET_MS,
					gates: preflight
						? ["manifest", "artifacts", "hash"]
						: ["manifest", "artifacts", "real-loader", "tests", "hash"],
				}
			: undefined;
	return {
		ok,
		diagnostics: state.diagnostics,
		...(state.manifest ? { manifest: state.manifest } : {}),
		...(hash ? { hash } : {}),
		...(attestation ? { attestation } : {}),
		durationMs,
	};
}

async function emitWorkerResult(result: WorkerValidationResult): Promise<void> {
	const resultPath = process.env.ULTRA_VALIDATOR_RESULT_PATH;
	if (resultPath) await writeFile(resultPath, JSON.stringify(result), "utf8");
	console.log(`ULTRA_VALIDATION_RESULT=${JSON.stringify(result)}`);
}

const invokedPath = process.argv[1] ? resolve(process.argv[1]) : undefined;
if (invokedPath === resolve(fileURLToPath(import.meta.url))) {
	const draftPath = process.argv[2];
	if (!draftPath) {
		console.error("usage: dynamic-validator-worker.ts <working-copy-path>");
		process.exitCode = 2;
	} else {
		validateDraft(draftPath, async (result) => {
			const resultPath = process.env.ULTRA_VALIDATOR_RESULT_PATH;
			if (resultPath)
				await writeFile(
					`${resultPath}.preflight`,
					JSON.stringify(result),
					"utf8",
				);
		}).then(
			(result) => emitWorkerResult(result),
			async (error) => {
				await emitWorkerResult({
					ok: false,
					diagnostics: [
						{
							gate: "validator",
							code: "validator_crash",
							status: "error",
							actual:
								error instanceof Error
									? (error.stack ?? error.message)
									: String(error),
							remediation: "Repair validator failure before retrying.",
							workingCopyPath: draftPath,
						},
					],
					durationMs: 0,
				});
				process.exitCode = 1;
			},
		);
	}
}