Luigit
repositories / pi-ext

pi-ext

bugabingas pi extensions

owned by admin

extensions/ultra/spec.ts

Raw
import * as fs from "node:fs";
import * as os from "node:os";
import * as path from "node:path";
import { fileURLToPath } from "node:url";
import { StringEnum } from "@earendil-works/pi-ai";
import { CONFIG_DIR_NAME, getAgentDir } from "@earendil-works/pi-coding-agent";
import { Type } from "typebox";
import { Compile } from "typebox/compile";
import { DEFAULT_FANOUT_TOOLS, THINKING_LEVELS } from "./settings.ts";

// ---------------------------------------------------------------------------
// Schemas (TypeBox -> JSON Schema fragments)
// ---------------------------------------------------------------------------

export const DYNAMIC_EXTENSION_TOOL_NAMES = [
	"search_dynamic_extensions",
	"create_dynamic_extension",
	"copy_dynamic_extension",
	"validate_dynamic_extension",
] as const;

/** Any object accepted as an inline JSON Schema (TypeBox-compatible). */
const JsonSchemaObject = Type.Object(
	{},
	{
		additionalProperties: true,
		description: "An inline JSON Schema object for this step's JSON output.",
	},
);

/** Full `ThinkingLevel` union. */
export const ThinkingLevelSchema = StringEnum(THINKING_LEVELS, {
	description: "Overrides the tier default; max requires model support.",
});

/** A step's `schema` is either a named reference (string) or an inline object. */
const StepSchemaRef = Type.Union(
	[
		Type.String({
			description: "Key in top-level schemas.",
		}),
		JsonSchemaObject,
	],
	{
		description: "Named or inline output schema.",
	},
);

export const StepSchema = Type.Object(
	{
		summary: Type.String({
			minLength: 1,
			description: "Progress label; supports prompt interpolation.",
		}),
		prompt: Type.String({
			description:
				"Child assignment; use only supported interpolation tokens in braces.",
		}),
		tools: Type.Optional(
			Type.Array(Type.String(), {
				description:
					"Initial child tools; fanout without writeIsolation uses the ultra fanout tool policy allowlist.",
			}),
		),
		model: Type.Optional(
			Type.String({
				description:
					"Provider/model or configured tier; omitted uses session model.",
			}),
		),
		thinkingLevel: Type.Optional(ThinkingLevelSchema),
		flags: Type.Optional(
			Type.Record(
				Type.String({ pattern: "^[A-Za-z0-9][A-Za-z0-9-]*$" }),
				Type.Union([Type.Boolean(), Type.String()]),
				{
					description:
						"Pi extension flags for the child, without leading dashes; overrides inherited parent flags. Enables tools that extensions keep off by default.",
				},
			),
		),
		dynamicExtensions: Type.Optional(
			Type.Union(
				[
					Type.Array(Type.String({ minLength: 1 }), {
						description:
							"Stable dynamic-extension catalog names, loaded in declared order.",
					}),
					Type.String({
						description: "One selector resolving to catalog names.",
					}),
				],
				{
					description:
						"Purpose-built validated extensions to resolve immediately before child creation.",
				},
			),
		),
		schema: Type.Optional(StepSchemaRef),
	},
	{
		additionalProperties: false,
		description: "One child returning a JSON object.",
	},
);

/** `over` accepts a literal array or an interpolation expression string. */
const OverSchema = Type.Union(
	[
		Type.Array(Type.Unknown(), {
			description: "Literal items for one fanout agent per array element.",
		}),
		Type.String({
			description: "One selector resolving to an array.",
		}),
	],
	{
		description:
			"Fanout input: a literal array or one selector that resolves to an array.",
	},
);

const PhaseIdSchema = Type.String({
	pattern: "^[A-Za-z_][A-Za-z0-9_]*$",
	description: "Unique phase ID; args is reserved.",
});

const WhenSchema = Type.Optional(
	Type.String({
		description: "One selector; false or empty skips the phase.",
	}),
);
const ConcurrencySchema = Type.Optional(
	Type.Number({
		description: "Optional maximum concurrent agents for this phase.",
	}),
);

// Phase is a discriminated union so `over` is *structurally* required for
// `fanout` and absent-tolerant for `single` — a fanout without `over` is
// rejected by TypeBox itself, no extra semantic pass needed.
const SinglePhaseSchema = Type.Object(
	{
		id: PhaseIdSchema,
		kind: Type.Literal("single", {
			description: "One child; over is ignored.",
		}),
		when: WhenSchema,
		over: Type.Optional(OverSchema),
		step: StepSchema,
		concurrency: ConcurrencySchema,
	},
	{ description: "A phase that invokes one sub-agent." },
);

const FanoutPhaseSchema = Type.Object(
	{
		id: PhaseIdSchema,
		kind: Type.Literal("fanout", {
			description: "One child per over item; read-only without writeIsolation.",
		}),
		when: WhenSchema,
		over: OverSchema,
		step: StepSchema,
		concurrency: ConcurrencySchema,
		writeIsolation: Type.Optional(
			Type.String({
				minLength: 1,
				description: "Disjoint mutation ownership per item.",
			}),
		),
	},
	{
		description: "Parallel children; read-only unless ownership is declared.",
	},
);

export const PhaseSchema = Type.Union([SinglePhaseSchema, FanoutPhaseSchema], {
	description: "An ordered single or fanout workflow phase.",
});

export const WorkflowSpecSchema = Type.Object(
	{
		name: Type.String({
			minLength: 1,
			description: "Non-empty workflow identifier.",
		}),
		description: Type.Optional(
			Type.String({ description: "One-line workflow purpose." }),
		),
		args: Type.Optional(
			Type.Object(
				{
					hint: Type.Optional(
						Type.String({ description: "Human-facing argument usage hint." }),
					),
					params: Type.Optional(
						Type.Array(Type.String({ pattern: "^[A-Za-z_][A-Za-z0-9_]*$" }), {
							description:
								"Positional input names; last absorbs remaining command text.",
						}),
					),
				},
				{
					additionalProperties: true,
					description: "Workflow input declaration.",
				},
			),
		),
		schemas: Type.Optional(
			Type.Record(Type.String(), JsonSchemaObject, {
				description:
					"Named JSON output schemas. Every string step.schema must exactly match one key here.",
			}),
		),
		phases: Type.Array(PhaseSchema, {
			minItems: 1,
			description: "Ordered phases; references target prior phases only.",
		}),
		return: Type.Optional(
			Type.String({
				description: "One selector for the returned value.",
			}),
		),
		report: Type.Optional(
			Type.String({
				description:
					"One selector for the human report, also returned to the model.",
			}),
		),
	},
	{
		description:
			"Inline workflow. Read ultra-authoring before constructing a spec; invalid references fail before execution.",
	},
);

// ---------------------------------------------------------------------------
// Inferred types
// ---------------------------------------------------------------------------

export type ThinkingLevel = Type.Static<typeof ThinkingLevelSchema>;
export type Step = Type.Static<typeof StepSchema>;
export type Phase = Type.Static<typeof PhaseSchema>;
export type WorkflowSpec = Type.Static<typeof WorkflowSpecSchema>;

export function isDynamicExtensionBuilderStep(step: Step): boolean {
	const tools = new Set(step.tools ?? []);
	return DYNAMIC_EXTENSION_TOOL_NAMES.every((name) => tools.has(name));
}

export function workflowUsesDynamicExtensions(spec: WorkflowSpec): boolean {
	return spec.phases.some(
		(phase) =>
			isDynamicExtensionBuilderStep(phase.step) ||
			phase.step.dynamicExtensions !== undefined,
	);
}

// ---------------------------------------------------------------------------
// Compilation
// ---------------------------------------------------------------------------

/**
 * Compile a JSON Schema (a named schema from a spec's `schemas` map, or an
 * inline step schema object) into a high-performance TypeBox validator with
 * `.Check`, `.Parse`, and `.Errors`.
 */
export function compileSchema(jsonSchema: object) {
	return Compile(jsonSchema);
}

// Precompiled once: validates a whole workflow spec.
const workflowValidator = Compile(WorkflowSpecSchema);

// ---------------------------------------------------------------------------
// parseWorkflow
// ---------------------------------------------------------------------------

function formatErrors(value: unknown): string {
	return workflowValidator
		.Errors(value)
		.map((e) => `${e.instancePath || "(root)"}: ${e.message}`)
		.join("; ");
}

/**
 * Validate untrusted workflow data (a parsed object or a JSON string) against
 * the workflow schema. THROWS (never returns) on invalid input, surfacing the
 * underlying TypeBox error message. Returns the typed spec on success.
 */
export function parseWorkflow(
	json: unknown,
	fanoutToolAllowlist: readonly string[] = DEFAULT_FANOUT_TOOLS,
): WorkflowSpec {
	let value: unknown = json;
	if (typeof json === "string") {
		try {
			value = JSON.parse(json);
		} catch (err) {
			const reason = err instanceof Error ? err.message : String(err);
			throw new Error(`Invalid workflow spec: not valid JSON (${reason})`);
		}
	}

	if (!workflowValidator.Check(value)) {
		throw new Error(`Invalid workflow spec: ${formatErrors(value)}`);
	}

	// Semantic check beyond structure: a string `schema` is a name reference and
	// must resolve against the spec's `schemas` map.
	const spec = value as WorkflowSpec;
	const schemaNames = new Set(Object.keys(spec.schemas ?? {}));
	const phaseIds = new Set<string>();
	for (const phase of spec.phases) {
		if (phase.id === "args") {
			throw new Error('Invalid workflow spec: reserved phase id "args".');
		}
		if (phaseIds.has(phase.id)) {
			throw new Error(
				`Invalid workflow spec: duplicate phase id "${phase.id}".`,
			);
		}
		phaseIds.add(phase.id);

		const ref = phase.step.schema;
		if (typeof ref === "string" && !schemaNames.has(ref)) {
			const known =
				schemaNames.size > 0 ? [...schemaNames].join(", ") : "(none)";
			throw new Error(
				`Invalid workflow spec: phase "${phase.id}" references unknown schema "${ref}". ` +
					`Known schemas: ${known}.`,
			);
		}

		const tools = new Set(phase.step.tools ?? []);
		const builderCount = DYNAMIC_EXTENSION_TOOL_NAMES.filter((name) =>
			tools.has(name),
		).length;
		if (
			builderCount !== 0 &&
			builderCount !== DYNAMIC_EXTENSION_TOOL_NAMES.length
		) {
			throw new Error(
				`Invalid workflow spec: phase "${phase.id}" selects ${builderCount} dynamic-extension builder tools; select all four or none.`,
			);
		}
		if (phase.kind === "fanout" && phase.writeIsolation === undefined) {
			const allowed = new Set(fanoutToolAllowlist);
			const outsideAllowlist = [...tools].filter((tool) => !allowed.has(tool));
			if (outsideAllowlist.length > 0) {
				throw new Error(
					`Invalid workflow spec: fanout phase "${phase.id}" selects tools outside the fanout allowlist without writeIsolation: ${outsideAllowlist.join(", ")}. Allowed initial tools: ${fanoutToolAllowlist.join(", ")}. Use allowed tools, declare genuine disjoint ownership with writeIsolation, or configure ultra.fanoutToolAllowlist.`,
				);
			}
		}
	}

	return spec;
}

// ---------------------------------------------------------------------------
// Discovery
// ---------------------------------------------------------------------------

export type WorkflowSource = "project" | "global" | "bundled";

export interface DiscoveredWorkflow {
	name: string;
	source: WorkflowSource;
	path: string;
	spec: WorkflowSpec;
}

export function missingWorkflowTools(
	spec: WorkflowSpec,
	availableTools: Iterable<string>,
): string[] {
	const available = new Set(availableTools);
	return [...new Set(spec.phases.flatMap((phase) => phase.step.tools ?? []))]
		.filter((tool) => !available.has(tool))
		.sort();
}

export interface DiscoverOptions {
	/** Override the home directory for global workflow discovery in tests. */
	home?: string;
	/** Override the extension's bundled `workflows/` directory. */
	bundledDir?: string;
}

function bundledWorkflowsDir(): string {
	const here = fs.realpathSync(path.dirname(fileURLToPath(import.meta.url)));
	return path.join(here, "workflows");
}

function readWorkflowsFromDir(
	dir: string,
	source: WorkflowSource,
): DiscoveredWorkflow[] {
	let entries: string[];
	try {
		entries = fs.readdirSync(dir);
	} catch {
		// Missing directory is the normal state (e.g. no project workflows, or the
		// bundled dir not yet shipped) — silently contribute nothing.
		return [];
	}

	const out: DiscoveredWorkflow[] = [];
	for (const entry of entries) {
		if (!entry.endsWith(".json")) continue;
		const filePath = path.join(dir, entry);
		let parsed: unknown;
		try {
			parsed = JSON.parse(fs.readFileSync(filePath, "utf8"));
		} catch {
			continue; // skip unreadable / malformed files
		}
		const name = (parsed as { name?: unknown })?.name;
		if (typeof name !== "string" || name.length === 0) continue;
		out.push({ name, source, path: filePath, spec: parsed as WorkflowSpec });
	}
	return out;
}

/**
 * Discover workflows from disk, lowest-to-highest precedence:
 *   bundled (`<ext>/workflows`) < global agent workflows < project workflows.
 * A project/global name overrides a bundled name of the same id. Missing
 * directories are tolerated. Uses Node `fs`/`path` only.
 */
export function discoverWorkflows(
	cwd: string,
	opts: DiscoverOptions = {},
): DiscoveredWorkflow[] {
	const home = opts.home ?? os.homedir();
	const bundledDir = opts.bundledDir ?? bundledWorkflowsDir();

	const projectDir = path.join(cwd, CONFIG_DIR_NAME, "workflows");
	const globalDir = opts.home
		? path.join(home, CONFIG_DIR_NAME, "agent", "workflows")
		: path.join(getAgentDir(), "workflows");

	// Build the map lowest-precedence first so later sources overwrite by name.
	const byName = new Map<string, DiscoveredWorkflow>();
	for (const wf of readWorkflowsFromDir(bundledDir, "bundled"))
		byName.set(wf.name, wf);
	for (const wf of readWorkflowsFromDir(globalDir, "global"))
		byName.set(wf.name, wf);
	for (const wf of readWorkflowsFromDir(projectDir, "project"))
		byName.set(wf.name, wf);

	return [...byName.values()];
}