Luigit
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".',
	);
}