Luigit
repositories / pi-ext

pi-ext

bugabingas pi extensions

owned by admin

extensions/ultra/implementation.ts

Raw
// ultra runtime: loaded only when the human invokes or completes /ultra.

import { existsSync } from "node:fs";
import { resolve } from "node:path";
import { fileURLToPath } from "node:url";
import {
	createAgentSession,
	DefaultResourceLoader,
	type ExtensionAPI,
	type ExtensionCommandContext,
	type ExtensionContext,
	getAgentDir,
	parseArgs,
	SessionManager,
} from "@earendil-works/pi-coding-agent";
import { Type } from "typebox";
import { createActionSummarizer } from "./action-summarizer.ts";
import { renderUltraCall } from "./call.ts";
import { formatUsage } from "./display.ts";
import {
	BUILTIN_TOOL_NAMES,
	loadSubagentFlagCatalog,
	makeResolveModel,
	makeSessionFactory,
	makeSubagentSessionManager,
	type SessionFactoryDeps,
	subagentFlagPrompt,
} from "./drivers.ts";
import type { DynamicExtensionResolver } from "./dynamic-types.ts";
import { runWorkflow, type WorkflowResult } from "./engine.ts";
import {
	createSessionWorkflowJournal,
	type SessionWorkflowJournal,
} from "./journal.ts";
import { initialReducerState } from "./progress.ts";
import {
	type BaseStepOptions,
	makeStepRunner,
	type RunnerContext,
	resolveModelTier,
	type StepDrivers,
	type StepRunner,
} from "./runner.ts";
import {
	effectiveModelTiers,
	loadUltraSettings,
	type UltraSettings,
} from "./settings.ts";
import {
	type DiscoveredWorkflow,
	DYNAMIC_EXTENSION_TOOL_NAMES,
	discoverWorkflows,
	missingWorkflowTools,
	type WorkflowSpec,
	WorkflowSpecSchema,
	workflowUsesDynamicExtensions,
} from "./spec.ts";
import { span } from "./src/debug.ts";
import { renderUltraResult, UltraOverlayComponent } from "./ui.ts";
import {
	completeUltraCommand,
	formatWorkflowToolResult,
	makeToolProgressSink,
	parseCommandLine,
	parseUltraCommand,
	resolveWorkflowInput,
} from "./wire.ts";

interface SubagentExtensionSelection {
	all: boolean;
	paths: string[];
	pathKeys: Set<string>;
}

function pathKey(path: string): string {
	const absolute = resolve(path).replaceAll("\\", "/");
	return process.platform === "win32" ? absolute.toLowerCase() : absolute;
}

function notifyWorkflowCompletion(
	ctx: ExtensionContext,
	result: WorkflowResult,
): void {
	if (!ctx.hasUI) return;
	const dropped = result.phases.reduce(
		(count, phase) => count + phase.dropped,
		0,
	);
	const status = result.aborted
		? "aborted"
		: dropped > 0
			? "incomplete"
			: "complete";
	const usage = formatUsage(result.tokenUsage).summary;
	ctx.ui.notify(
		`ultra: "${result.workflow}" ${status} · ${result.phases.length} phases · ${usage}.`,
		status === "complete" ? "info" : "warning",
	);
}

export function resolveSubagentExtensions(
	setting: UltraSettings["subagentExtensions"],
): SubagentExtensionSelection {
	if (setting === "all") return { all: true, paths: [], pathKeys: new Set() };
	if (setting === "none") return { all: false, paths: [], pathKeys: new Set() };

	const paths = setting.map((name) => {
		if (!/^[a-z0-9][a-z0-9-]*$/u.test(name))
			throw new Error(`ultra: invalid sub-agent extension name: ${name}`);
		// Sibling packages sit next to this one: `index.js` when generated, `index.ts` in source.
		const path = ["index.js", "index.ts"]
			.map((entry) =>
				fileURLToPath(new URL(`../${name}/${entry}`, import.meta.url)),
			)
			.find((candidate) => existsSync(candidate));
		if (!path) throw new Error(`ultra: unknown sub-agent extension: ${name}`);
		return path;
	});
	return { all: false, paths, pathKeys: new Set(paths.map(pathKey)) };
}

export function createUltraRuntime(pi: ExtensionAPI) {
	const settings = (ctx: ExtensionContext): UltraSettings => {
		const cfg = loadUltraSettings(pi, ctx);
		return {
			...cfg,
			modelTiers: effectiveModelTiers(cfg, ctx.model?.provider),
		};
	};

	// Workflow files are tiny and human-authored during live sessions. Rescan at
	// each use so a newly saved workflow immediately works without `/reload`.
	const discover = (cwd: string): DiscoveredWorkflow[] =>
		discoverWorkflows(cwd);

	const activateWorkflowTool = () => {
		const active = pi.getActiveTools();
		if (!active.includes("run_workflow")) {
			pi.setActiveTools([...active, "run_workflow"]);
		}
	};

	const deactivateWorkflowTool = () => {
		const active = pi.getActiveTools();
		if (active.includes("run_workflow")) {
			pi.setActiveTools(active.filter((name) => name !== "run_workflow"));
		}
	};

	const availableSubagentTools = (
		extensions: SubagentExtensionSelection,
	): string[] =>
		pi
			.getAllTools()
			.filter(
				(tool) =>
					tool.name !== "run_workflow" &&
					(extensions.all ||
						BUILTIN_TOOL_NAMES.has(tool.name) ||
						tool.sourceInfo.source === "builtin" ||
						extensions.pathKeys.has(pathKey(tool.sourceInfo.path))),
			)
			.map((tool) => tool.name);

	const workflowError = (
		spec: WorkflowSpec,
		extensions: SubagentExtensionSelection,
		cfg: UltraSettings,
	): string | undefined => {
		try {
			for (const phase of spec.phases)
				resolveModelTier(phase.step, cfg.modelTiers);
		} catch (error) {
			return error instanceof Error ? error.message : String(error);
		}
		const available = [
			...availableSubagentTools(extensions),
			...DYNAMIC_EXTENSION_TOOL_NAMES,
		];
		const missing = missingWorkflowTools(spec, available);
		return missing.length > 0
			? `ultra: workflow "${spec.name}" requires unavailable tools: ${missing.join(", ")}`
			: undefined;
	};

	const createJournal = (
		ctx: ExtensionContext,
		spec: WorkflowSpec,
		cfg: UltraSettings,
		args?: Record<string, unknown>,
	) =>
		createSessionWorkflowJournal({
			cwd: ctx.cwd,
			spec,
			args,
			modelRouting: spec.phases.map(({ step }) => {
				const resolved = resolveModelTier(step, cfg.modelTiers);
				return {
					model:
						resolved.model ??
						(ctx.model ? `${ctx.model.provider}/${ctx.model.id}` : undefined),
					thinkingLevel: resolved.thinkingLevel,
				};
			}),
			entries: ctx.sessionManager.getBranch(),
			appendEntry: (customType, data) => pi.appendEntry(customType, data),
		});

	// Assemble a real `StepRunner` from the injected SDK drivers + this ctx.
	const buildStepRunner = (
		ctx: ExtensionContext,
		spec: WorkflowSpec,
		cfg: UltraSettings,
		journal: SessionWorkflowJournal,
		extensions: SubagentExtensionSelection,
		dynamicExtensions?: DynamicExtensionResolver,
	): StepRunner => {
		// The injected-dep types are deliberately broader than the SDK's specific
		// signatures (so the factories unit-test with simple fakes); the real SDK
		// values are adapted to them here at the boundary. Runtime correctness of
		// these calls is the Task-10 live gate.
		const drivers: StepDrivers = {
			sessionFactory: makeSessionFactory({
				createAgentSession:
					createAgentSession as unknown as SessionFactoryDeps["createAgentSession"],
				DefaultResourceLoader:
					DefaultResourceLoader as unknown as SessionFactoryDeps["DefaultResourceLoader"],
				getAgentDir,
				cwd: ctx.cwd,
				loadExtensions: extensions.all,
				extensionPaths: extensions.paths,
				loadSkills: cfg.subagentSkills === "all",
				// Children inherit the parent's CLI extension flags, parsed exactly as Pi parses them.
				inheritedFlags: parseArgs(process.argv.slice(2)).unknownFlags,
			}),
			resolveModel: makeResolveModel({
				modelRegistry: ctx.modelRegistry,
				defaultModel: ctx.model,
			}),
		};
		const manageSubagentSession = makeSubagentSessionManager({
			SessionManager,
			cwd: ctx.cwd,
			parentSession: ctx.sessionManager.getSessionFile(),
			sessionDir: ctx.sessionManager.getSessionDir(),
			retention: cfg.subagentSessionRetention,
		});
		const ctxLike: RunnerContext = {
			newSessionManager(sessionKey, agentId) {
				const retained = journal.getSession(sessionKey);
				const child = manageSubagentSession({
					name: `ultra ${spec.name} ${journal.id.slice(0, 8)} ${agentId}`,
					retained,
				});
				if (
					child.ref &&
					(!retained ||
						child.ref.id !== retained.id ||
						child.ref.path !== retained.path)
				)
					journal.setSession(sessionKey, child.ref);
				return child;
			},
		};
		const baseOpts: BaseStepOptions = {
			maxRetries: cfg.maxRetries,
			schemas: spec.schemas,
			modelTiers: cfg.modelTiers,
			actionSummarizer: createActionSummarizer(ctx, cfg.actionSummarizerModel),
			dynamicExtensionResolver: dynamicExtensions,
		};
		return makeStepRunner(drivers, ctxLike, baseOpts);
	};

	const dynamicResolver = async (
		ctx: ExtensionContext,
		spec: WorkflowSpec,
		journal: SessionWorkflowJournal,
		extensions: SubagentExtensionSelection,
	): Promise<DynamicExtensionResolver | undefined> => {
		if (!workflowUsesDynamicExtensions(spec)) return undefined;
		const { createDynamicExtensionCoordinator } = await import(
			"./dynamic-catalog.ts"
		);
		return createDynamicExtensionCoordinator({
			agentDir: getAgentDir(),
			parentSessionId: ctx.sessionManager.getSessionId(),
			workflowRunId: journal.id,
			availableTools: availableSubagentTools(extensions),
		});
	};

	const prepareRun = async (
		ctx: ExtensionContext,
		spec: WorkflowSpec,
		cfg: UltraSettings,
		args: Record<string, unknown> | undefined,
		extensions: SubagentExtensionSelection,
	): Promise<{
		journal: SessionWorkflowJournal;
		dynamicExtensions?: DynamicExtensionResolver;
		stepRunner: StepRunner;
	}> => {
		const journal = createJournal(ctx, spec, cfg, args);
		const dynamicExtensions = await dynamicResolver(
			ctx,
			spec,
			journal,
			extensions,
		);
		return {
			journal,
			dynamicExtensions,
			stepRunner: buildStepRunner(
				ctx,
				spec,
				cfg,
				journal,
				extensions,
				dynamicExtensions,
			),
		};
	};

	// ── `run_workflow` tool (model-callable, mid-turn, read-only onUpdate surface) ──
	pi.registerTool({
		name: "run_workflow",
		label: "run_workflow",
		description:
			"run saved or inline multi-agent workflows. An inline spec takes precedence when both spec and name are supplied; args provide workflow inputs. " +
			"return selected output, failures, usage, and fullOutputPath for complete evidence; results above 50KB or 2000 lines return a bounded file reference.",
		promptSnippet: "run saved or inline multi-agent workflows.",
		promptGuidelines: [
			"Use run_workflow when decomposition into independent or sequential child-agent steps reduces context or verification risk.",
			"Use dynamic extensions when a workflow needs a reusable, task-specific capability; read the ultra-authoring skill for how to build and load them.",
		],
		parameters: Type.Object({
			name: Type.Optional(
				Type.String({
					description: "Saved workflow name (project/global/bundled).",
				}),
			),
			spec: Type.Optional(WorkflowSpecSchema),
			args: Type.Optional(
				Type.Record(Type.String(), Type.Unknown(), {
					description: "Workflow inputs, available as {args.NAME}.",
				}),
			),
		}),
		renderCall: renderUltraCall,
		renderResult: renderUltraResult,
		async execute(_toolCallId, params, signal, onUpdate, ctx) {
			const finish = span?.("workflow.execute");
			try {
				const cfg = settings(ctx);
				if (!cfg.enabled) {
					finish?.("finish", { outcome: "disabled" });
					return {
						content: [
							{
								type: "text",
								text: "ultra is disabled (ultra.enabled = false).",
							},
						],
						details: initialReducerState(),
					};
				}
				const { spec, args } = resolveWorkflowInput(
					params,
					discover(ctx.cwd),
					undefined,
					cfg.fanoutToolAllowlist,
				);
				const extensions = resolveSubagentExtensions(cfg.subagentExtensions);
				const error = workflowError(spec, extensions, cfg);
				if (error) throw new Error(error);
				const prepared = await prepareRun(ctx, spec, cfg, args, extensions);
				// onProgress folds into the reducer + pushes through `onUpdate`; headless
				// (no `onUpdate`) degrades to a no-op sink that still tracks final state.
				const progress = makeToolProgressSink(onUpdate ?? (() => {}), {
					workflowName: spec.name,
					workflowDescription: spec.description,
				});
				const result = await runWorkflow(spec, args, {
					stepRunner: prepared.stepRunner,
					concurrency: cfg.concurrency,
					journal: prepared.journal,
					dynamicExtensions: prepared.dynamicExtensions,
					signal, // the SDK execute() signal: stops the pool + fans session.abort()
					onProgress: progress.sink,
				});
				const output = await formatWorkflowToolResult(result);
				notifyWorkflowCompletion(ctx, result);
				finish?.("finish", {
					outcome: result.aborted
						? "aborted"
						: result.phases.some((phase) => phase.dropped > 0)
							? "incomplete"
							: "completed",
				});
				return {
					content: [{ type: "text", text: output.text }],
					details: {
						...progress.finish({
							...result,
							fullOutputPath: output.fullOutputPath,
						}),
						fullOutputPath: output.fullOutputPath,
					},
					terminate: false, // the agent keeps reasoning over the results
				};
			} catch (error) {
				finish?.("error", { outcome: "failed" });
				throw error;
			}
		},
	});

	deactivateWorkflowTool();

	const sendExecPrompt = (request: string) => {
		activateWorkflowTool();
		pi.sendUserMessage(`# ultra exec\n\n# user request\n\n${request}`);
	};

	// One extension load per cwd and selection; children load the same set per step.
	const flagCatalogs = new Map<
		string,
		ReturnType<typeof loadSubagentFlagCatalog>
	>();
	const subagentFlagCatalog = (ctx: ExtensionContext) => {
		const extensions = resolveSubagentExtensions(
			settings(ctx).subagentExtensions,
		);
		const key = JSON.stringify([ctx.cwd, extensions.all, extensions.paths]);
		let catalog = flagCatalogs.get(key);
		if (!catalog) {
			catalog = loadSubagentFlagCatalog({
				DefaultResourceLoader:
					DefaultResourceLoader as unknown as SessionFactoryDeps["DefaultResourceLoader"],
				getAgentDir,
				cwd: ctx.cwd,
				loadExtensions: extensions.all,
				extensionPaths: extensions.paths,
			});
			flagCatalogs.set(key, catalog);
			// A failed load is retried on the next turn instead of being cached.
			catalog.catch(() => flagCatalogs.delete(key));
		}
		return catalog;
	};

	// ── `/ultra [exec | run]` command ──
	return {
		activateWorkflowTool,
		/** Undefined when no sub-agent extension owns both tools and flags, or loading failed. */
		subagentFlagSection: async (ctx: ExtensionContext) => {
			const catalog = await subagentFlagCatalog(ctx).catch(() => []);
			return catalog.length ? subagentFlagPrompt(catalog) : undefined;
		},
		getArgumentCompletions: (prefix: string) =>
			completeUltraCommand(discover(process.cwd()), prefix),
		handleCommand: async (argString: string, ctx: ExtensionCommandContext) => {
			const cfg = settings(ctx);
			if (!cfg.enabled) {
				if (ctx.hasUI)
					ctx.ui.notify(
						"ultra is disabled (ultra.enabled = false).",
						"warning",
					);
				return "disabled";
			}
			// The command surface is the focus-grabbing overlay — interactive only.
			// Headless / non-UI callers use the `run_workflow` tool instead.
			if (!ctx.hasUI) {
				return "skipped";
			}
			const command = parseUltraCommand(argString);
			if (command.subcommand === "activate") {
				activateWorkflowTool();
				ctx.ui.notify("ultra: run_workflow tool enabled.", "info");
				return "completed";
			}
			if (command.subcommand === "exec") {
				sendExecPrompt(command.instructions);
				return "completed";
			}
			if (command.subcommand === "invalid") {
				ctx.ui.notify(
					"Usage: /ultra | /ultra exec <instructions> | /ultra run <workflow> [args]",
					"warning",
				);
				return "invalid";
			}

			const discovered = discover(ctx.cwd);
			const parsed = parseCommandLine(command.invocation);
			let resolved: { spec: WorkflowSpec; args?: Record<string, unknown> };
			try {
				resolved = resolveWorkflowInput(
					{ name: parsed.name, args: parsed.args },
					discovered,
					parsed.rest,
					cfg.fanoutToolAllowlist,
				);
			} catch (err) {
				ctx.ui.notify(
					err instanceof Error ? err.message : String(err),
					"error",
				);
				return "invalid";
			}
			const { spec, args: runArgs } = resolved;
			let extensions: SubagentExtensionSelection;
			try {
				extensions = resolveSubagentExtensions(cfg.subagentExtensions);
			} catch (error) {
				ctx.ui.notify(
					error instanceof Error ? error.message : String(error),
					"error",
				);
				return "failed";
			}
			const error = workflowError(spec, extensions, cfg);
			if (error) {
				ctx.ui.notify(error, "error");
				return "invalid";
			}
			const prepared = await prepareRun(ctx, spec, cfg, runArgs, extensions);

			// The command owns its AbortController (idle ⇒ ctx.signal is undefined).
			// Global Esc / per-agent abort fire it; the engine fans session.abort().
			const controller = new AbortController();
			const outcome = await ctx.ui.custom<
				| { status: "completed"; result: WorkflowResult }
				| { status: "aborted"; result?: WorkflowResult }
				| { status: "failed"; message: string }
				| null
			>(
				(tui, theme, keybindings, done) => {
					const overlay = new UltraOverlayComponent({
						host: tui,
						theme,
						onRunAbort: () => controller.abort(),
						workflowName: spec.name,
						workflowDescription: spec.description,
						toolUi: tui,
						toolsExpanded: () => ctx.ui.getToolsExpanded(),
						hideToolPreviews: cfg.actionSummarizerModel !== null,
						isToolsExpandKey: (data) =>
							keybindings.matches(data, "app.tools.expand"),
						onToolsExpand: () =>
							ctx.ui.setToolsExpanded(!ctx.ui.getToolsExpanded()),
						cwd: ctx.cwd,
					});
					// P3: return the Component SYNCHRONOUSLY; run the engine DETACHED, or
					// the overlay never renders. `done` resolves the custom() promise.
					void (async () => {
						try {
							const r = await runWorkflow(spec, runArgs, {
								stepRunner: prepared.stepRunner,
								concurrency: cfg.concurrency,
								journal: prepared.journal,
								dynamicExtensions: prepared.dynamicExtensions,
								signal: controller.signal,
								onProgress: (e) => overlay.handleProgress(e),
							});
							done(
								controller.signal.aborted
									? { status: "aborted", result: r }
									: { status: "completed", result: r },
							);
						} catch (error) {
							done(
								controller.signal.aborted
									? { status: "aborted" }
									: {
											status: "failed",
											message:
												error instanceof Error ? error.message : String(error),
										},
							);
						}
					})();
					return overlay;
				},
				{
					overlay: true,
					overlayOptions: {
						anchor: "bottom-left",
						width: "100%",
						maxHeight: "100%",
					},
				},
			);

			if (!outcome) {
				ctx.ui.notify(`ultra: "${spec.name}" aborted.`, "warning");
				return "aborted";
			}
			if (outcome.status === "failed") {
				ctx.ui.notify(outcome.message, "error");
				return "failed";
			}
			const result = outcome.result;
			if (!result) {
				ctx.ui.notify(`ultra: "${spec.name}" aborted.`, "warning");
				return "aborted";
			}
			const output = await formatWorkflowToolResult(result);
			pi.sendMessage({
				customType: "ultra-result",
				content: output.text,
				display: true,
				details: { ...result, fullOutputPath: output.fullOutputPath },
			});
			notifyWorkflowCompletion(ctx, result);
			return result.aborted
				? "aborted"
				: result.phases.some((phase) => phase.dropped > 0)
					? "incomplete"
					: "completed";
		},
	};
}