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";
},
};
}