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; export type Step = Type.Static; export type Phase = Type.Static; export type WorkflowSpec = Type.Static; 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(); 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[] { 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 (`/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(); 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()]; }