repositories / pi-ext
pi-ext
bugabingas pi extensions
owned by admin
extensions/ultra/wire.ts
Raw// ultra — wiring helpers (Task 9, offline-testable half).
//
// The pieces of the extension surface that carry real logic but no SDK
// coupling: the tool-surface progress sink (folds engine events with the shared
// reducer and pushes them through `onUpdate`), the `getArgumentCompletions`
// filter, and name-vs-inline spec resolution. `implementation.ts` wires these
// to the real SDK after `/ultra` is first used.
import { mkdtemp, writeFile } from "node:fs/promises";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { truncateHead } from "@earendil-works/pi-coding-agent";
import type { WorkflowResult } from "./engine.ts";
import {
initialReducerState,
type ProgressMeta,
type ReducerState,
reduceProgress,
type UltraProgressEvent,
} from "./progress.ts";
import {
type DiscoveredWorkflow,
parseWorkflow,
type WorkflowSpec,
} from "./spec.ts";
export { formatUltraResultMarkdown } from "./result.ts";
/** One bounded model-facing envelope for both surfaces, with retrievable evidence. */
export async function formatWorkflowToolResult(
result: WorkflowResult,
): Promise<{ text: string; fullOutputPath: string }> {
const directory = await mkdtemp(join(tmpdir(), "ultra-output-"));
const fullOutputPath = join(directory, "result.json");
await writeFile(fullOutputPath, JSON.stringify(result, null, 2), {
mode: 0o600,
});
const dropped = result.phases.reduce(
(count, phase) => count + phase.dropped,
0,
);
// This describes execution, not the quality or semantic success of the artifact.
const executionStatus = result.aborted
? "aborted"
: dropped > 0
? "incomplete"
: "finished";
const { phaseResults: _phaseResults, phaseFailures, ...output } = result;
const text = JSON.stringify(
{
...output,
executionStatus,
fullOutputPath,
phaseFailures: Object.fromEntries(
Object.entries(phaseFailures).map(([phase, failures]) => [
phase,
failures.map(
({ item: _item, transcriptTail: _tail, ...failure }) => failure,
),
]),
),
},
null,
2,
);
if (!truncateHead(text).truncated) return { text, fullOutputPath };
return {
text: JSON.stringify({
truncated: true,
executionStatus,
dropped,
message:
"Workflow output exceeded 50KB or 2000 lines. Read the full output selectively from the file.",
fullOutputPath,
}),
fullOutputPath,
};
}
// ---------------------------------------------------------------------------
// Tool-surface progress sink
// ---------------------------------------------------------------------------
/** The minimal `AgentToolResult` shape `onUpdate` consumes; `renderResult`
* reads `details` (validated P2). */
export interface ToolUpdate {
content: { type: "text"; text: string }[];
details: ReducerState;
}
/** One-line fallback summary (the live board is drawn from `details` by
* `renderResult`; this `content` is just the textual fallback). */
function summarize(state: ReducerState): string {
const phase = state.phase ? `${state.phase}: ` : "";
return [
`${phase}${state.done}/${state.total} done`,
state.running ? `${state.running} active` : "",
state.queued ? `${state.queued} queued` : "",
state.later ? `${state.later} later` : "",
state.dropped ? `${state.dropped} dropped` : "",
]
.filter(Boolean)
.join(", ");
}
/**
* Build the `run_workflow` tool's progress sink: each engine `onProgress` event is
* folded with the shared reducer and pushed through `onUpdate` as the full
* `AgentToolResult` shape (so `renderResult` can read fresh `details` every
* tick). `getState` exposes the accumulated board for the final tool return.
*/
export function makeToolProgressSink(
onUpdate: (update: ToolUpdate) => void,
meta: ProgressMeta = {},
): {
sink: (event: UltraProgressEvent) => void;
finish: (result: WorkflowResult) => ReducerState;
getState: () => ReducerState;
} {
let state = initialReducerState(meta);
return {
sink: (event) => {
// Strip `controls` before folding: the tool surface ignores them, and
// Pi STRUCTURED-CLONES the `details` it receives (both per-tick via
// `onUpdate` and the final tool result). `controls` are functions, so
// leaving them in throws "The object can not be cloned.",
// which propagates back into every step that emitted an event and
// fails the whole run. Only the command overlay (never serialized)
// keeps controls.
const safe =
"agentId" in event
? { ...event, controls: undefined, prompt: undefined }
: event;
state = reduceProgress(state, safe);
onUpdate({
content: [{ type: "text", text: summarize(state) }],
details: state,
});
},
finish: (result) => {
state = { ...state, workflowResult: result };
return state;
},
getState: () => state,
};
}
// ---------------------------------------------------------------------------
// Argument completions
// ---------------------------------------------------------------------------
/** Minimal SDK `AutocompleteItem` shape used by command completion. */
export interface AutocompleteItem {
value: string;
label: string;
description?: string;
}
/** Filter discovered workflow names by a (case-insensitive) prefix; `null` when
* nothing matches (the SDK convention for "no completions"). */
export function completeWorkflowNames(
discovered: readonly DiscoveredWorkflow[],
prefix: string,
): AutocompleteItem[] | null {
const p = prefix.toLowerCase();
const items = discovered
.filter((w) => w.name.toLowerCase().startsWith(p))
.map((w) => ({ value: w.name, label: w.name, description: w.source }));
return items.length > 0 ? items : null;
}
const ULTRA_SUBCOMMANDS: AutocompleteItem[] = [
{
value: "exec",
label: "exec",
description: "execute instructions through an authored workflow",
},
{
value: "run",
label: "run",
description: "run a saved workflow",
},
];
export function completeUltraCommand(
discovered: readonly DiscoveredWorkflow[],
prefix: string,
): AutocompleteItem[] | null {
const input = prefix.trimStart();
if (!/\s/.test(input)) {
const p = input.toLowerCase();
const items = ULTRA_SUBCOMMANDS.filter((item) => item.value.startsWith(p));
return items.length > 0 ? items : null;
}
const match = input.match(/^run\s+(.*)$/i);
if (!match) return null;
const workflows = completeWorkflowNames(discovered, match[1]);
return (
workflows?.map((item) => ({
...item,
value: `run ${item.value}`,
})) ?? null
);
}
export type UltraCommand =
| { subcommand: "activate" }
| { subcommand: "exec"; instructions: string }
| { subcommand: "run"; invocation: string }
| { subcommand: "invalid" };
export function parseUltraCommand(raw: string): UltraCommand {
const input = raw.trim();
if (!input) return { subcommand: "activate" };
const space = input.search(/\s/);
const subcommand = space === -1 ? input : input.slice(0, space);
const rest = space === -1 ? "" : input.slice(space + 1).trim();
if (subcommand === "exec" && rest)
return { subcommand: "exec", instructions: rest };
if (subcommand === "run" && rest)
return { subcommand: "run", invocation: rest };
return { subcommand: "invalid" };
}
// ---------------------------------------------------------------------------
// Spec resolution
// ---------------------------------------------------------------------------
/** The `run_workflow` tool / `/ultra run` command input. */
export interface WorkflowInput {
name?: string;
spec?: unknown;
args?: Record<string, unknown>;
}
/**
* Parse a raw `/ultra run <name> [args]` invocation: the first whitespace token is
* the workflow name; the remainder is either a trailing JSON object (used as the
* run args directly) or bare text returned verbatim as `rest` for positional
* mapping against the spec's declared `args.params` (see `mapPositionalArgs`).
* The name always resolves regardless of the remainder's shape.
*/
export function parseCommandLine(raw: string): {
name?: string;
args?: Record<string, unknown>;
rest?: string;
} {
const trimmed = raw.trim();
if (!trimmed) return {};
const space = trimmed.search(/\s/);
if (space === -1) return { name: trimmed };
const name = trimmed.slice(0, space);
const rest = trimmed.slice(space + 1).trim();
if (!rest) return { name };
try {
const parsed = JSON.parse(rest);
if (
parsed !== null &&
typeof parsed === "object" &&
!Array.isArray(parsed)
) {
return { name, args: parsed as Record<string, unknown>, rest };
}
} catch {
// non-JSON trailing text → kept as `rest` for positional mapping
}
return { name, rest };
}
/**
* Map a bare command-line remainder to named args using a workflow's declared
* positional `params`: whitespace-split tokens fill `params` in order, and the
* last declared param soaks any remaining tokens so a trailing free-text arg
* (e.g. a feature description) survives intact. Returns `undefined` when nothing
* maps (empty remainder or no declared params). Only `/ultra run` uses
* this — `run_workflow` passes a structured `args` object directly.
*/
export function mapPositionalArgs(
rest: string,
params: readonly string[],
): Record<string, unknown> | undefined {
const tokens = rest.trim().split(/\s+/).filter(Boolean);
if (tokens.length === 0 || params.length === 0) return undefined;
const out: Record<string, unknown> = {};
params.forEach((name, i) => {
if (i >= tokens.length) return;
// The last declared param soaks the remaining tokens so a trailing
// free-text arg (e.g. a feature description) survives intact; earlier
// params take exactly one token each.
out[name] = i === params.length - 1 ? tokens.slice(i).join(" ") : tokens[i];
});
return Object.keys(out).length > 0 ? out : undefined;
}
/** Resolved run input: a validated spec + its run args. */
export interface ResolvedWorkflow {
spec: WorkflowSpec;
args?: Record<string, unknown>;
}
/**
* Resolve a run request to a validated spec. An inline `spec` (ephemeral mode)
* wins and is parsed/validated; otherwise a `name` (persisted mode) is looked up
* among discovered workflows. Throws on an unknown name, a malformed inline
* spec, or neither being supplied.
*/
export function resolveWorkflowInput(
input: WorkflowInput,
discovered: readonly DiscoveredWorkflow[],
rest?: string,
fanoutToolAllowlist?: readonly string[],
): ResolvedWorkflow {
if (input.spec !== undefined) {
return {
spec: parseWorkflow(input.spec, fanoutToolAllowlist),
args: input.args,
};
}
if (input.name) {
const found = discovered.find((w) => w.name === input.name);
if (!found) {
const known = discovered.map((w) => w.name).join(", ") || "(none)";
throw new Error(
`ultra: unknown workflow "${input.name}". Available: ${known}.`,
);
}
// Validate the discovered spec at resolve time (not at discovery, so a
// malformed file still lists in the picker) — `/ultra run <name>` then fails
// fast with parseWorkflow's TypeBox path error instead of crashing the
// engine mid-run on a raw, unvalidated object.
const spec = parseWorkflow(found.spec, fanoutToolAllowlist);
let args = input.args;
// Positional fallback: a bare `/ultra run <name> <tok...>` line (no JSON args
// object) maps to the spec's declared `args.params` in order.
if (args === undefined && rest && spec.args?.params?.length) {
args = mapPositionalArgs(rest, spec.args.params);
}
return { spec, args };
}
throw new Error(
'ultra: provide either a workflow "name" or an inline "spec".',
);
}