repositories / pi-ext
pi-ext
bugabingas pi extensions
owned by admin
extensions/ultra/spec.ts
Rawimport * 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()];
}