repositories / pi-ext
pi-ext
bugabingas pi extensions
owned by admin
extensions/ultra/dynamic-validator-worker.ts
Rawimport { 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;
},
);
}
}