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