Luigit
repositories / pi-ext

pi-ext

bugabingas pi extensions

owned by admin

extensions/ultra/progress.ts

Raw
// ultra — live progress: pure event types, session-event mapping, and a shared
// board/transcript reducer (Task 7a). NO SDK *session* is ever constructed
// here; only the `AgentSessionEvent` *type* is imported. Both `runner.ts`
// (emits these events + binds `controls`) and `ui.ts` (folds them with the
// reducer and renders) import from this module, so the names are the contract.
//
// Field paths read from `AgentSessionEvent` are the ones validated against a
// real Pi + MiniMax run (prototypes P1/P4): streaming text arrives on
// `event.assistantMessageEvent` as `{ type: "text_delta", delta }` during
// `message_update` (NOT `event.message.content`, which is empty mid-stream); a
// terminating tool's structured payload is on `tool_execution_end.result.details`.

import type { AgentSessionEvent } from "@earendil-works/pi-coding-agent";
import { formatStructuredText, summarizeValue } from "./display.ts";
import type { WorkflowResult } from "./engine.ts";
import type { StepFailure, TokenUsage } from "./runner.ts";
import type { ThinkingLevel } from "./spec.ts";

// ---------------------------------------------------------------------------
// Event types (the `deps.onProgress` contract — design line 242)
// ---------------------------------------------------------------------------

/** Lifecycle stage of a single mapped agent event. */
export type AgentProgressKind = "prompt" | "start" | "action" | "delta" | "end";

/**
 * Live control surface for one running sub-agent, bound by the session-owning
 * runner. Only the command overlay acts on these; the tool surface ignores
 * `controls`. The methods may
 * be sync or async (`session.steer`/`session.abort` return promises that the
 * runner must `await`/`.catch()` to avoid floating rejections).
 */
export interface AgentControls {
	steer(text: string): Promise<void> | void;
	abort(): Promise<void> | void;
}

/**
 * A single progress event forwarded through `deps.onProgress`. Produced by
 * `mapSessionEvent` (for tool-using steps' session events) or constructed
 * directly by the runner/engine (coarse start/end, dropped items). `controls`
 * is attached by the runner, never by `mapSessionEvent`.
 */
export interface AgentProgressEvent {
	agentId: string;
	phase: string;
	kind: AgentProgressKind;
	/** Tool name for `action` events (tool start/end). */
	toolName?: string;
	/** Tool call id for native tool rendering. */
	toolCallId?: string;
	/** Raw tool args for native tool rendering. */
	args?: unknown;
	/** Raw tool result for native tool rendering. */
	result?: unknown;
	isError?: boolean;
	isPartial?: boolean;
	/** Short resolved task summary shown on progress rows. */
	summary?: string;
	/** Exact resolved sub-agent prompt, shown as the first overlay transcript message. */
	prompt?: string;
	/** LLM-authored board phrase for an action; takes precedence over raw args. */
	actionSummary?: string;
	/** Streamed delta text, a tool action/result preview, or a coarse status line. */
	text?: string;
	/** Terminal/error marker: `end` carries the run outcome; `action` may carry `"error"`. */
	status?: string;
	/** Requested model tier before resolution. */
	tier?: string;
	/** Resolved provider/model label. */
	model?: string;
	thinkingLevel?: ThinkingLevel;
	startedAt?: number;
	usage?: TokenUsage;
	failure?: StepFailure;
	item?: unknown;
	/** Per-agent live controls (runner-bound; absent on the tool surface). */
	controls?: AgentControls;
}

export type PhasePlanStatus =
	| "resolved"
	| "agents-pending"
	| "condition-pending"
	| "done"
	| "skipped";

export interface PlannedPhase {
	phase: string;
	status: PhasePlanStatus;
	agentIds: readonly string[];
}

export type UltraProgressEvent =
	| AgentProgressEvent
	| { kind: "plan"; phases: readonly PlannedPhase[]; startedAt: number }
	| ({ kind: "phase" } & PlannedPhase);

// ---------------------------------------------------------------------------
// mapSessionEvent — pure AgentSessionEvent -> AgentProgressEvent
// ---------------------------------------------------------------------------

const MAX_PREVIEW = 120;

function cap(text: string): string {
	return text.length > MAX_PREVIEW ? `${text.slice(0, MAX_PREVIEW)}…` : text;
}

function isRecord(value: unknown): value is Record<string, unknown> {
	return typeof value === "object" && value !== null && !Array.isArray(value);
}

/** Short preview of arbitrary JSON without dumping nested bodies. */
function previewValue(value: unknown): string {
	return summarizeValue(value);
}

/** Short preview of a tool call's args for the board/transcript. */
function previewArgs(args: unknown): string {
	if (args === undefined || args === null) return "";
	return previewValue(args);
}

function firstLine(value: unknown): string | undefined {
	return typeof value === "string" ? value.split("\n")[0]?.trim() : undefined;
}

function pathArg(args: Record<string, unknown>): string | undefined {
	return firstLine(args.path) ?? firstLine(args.file);
}

function searchArgs(args: Record<string, unknown>): string {
	const pattern = firstLine(args.pattern) ?? firstLine(args.query);
	const path = pathArg(args);
	if (pattern && path) return `${pattern} in ${path}`;
	return pattern ?? path ?? previewArgs(args);
}

function actionIntent(toolName: string, args: unknown): string {
	if (!isRecord(args)) return `Using ${toolName}`;
	const path = pathArg(args);
	switch (toolName) {
		case "bash": {
			const command = firstLine(args.command);
			return command ? `Running ${cap(command)}` : "Running command";
		}
		case "read":
			return path ? `Reading ${path}` : "Reading input";
		case "write":
			return path ? `Writing ${path}` : "Writing file";
		case "edit":
			return path ? `Editing ${path}` : "Editing file";
		case "ls":
			return path ? `Listing ${path}` : "Listing files";
		case "grep":
			return `Searching ${searchArgs(args)}`;
		case "find":
			return `Finding ${searchArgs(args)}`;
		case "structured_output":
			return structuredOutputSummary(args);
		default: {
			const detail = previewArgs(args);
			return detail ? `Using ${toolName} · ${detail}` : `Using ${toolName}`;
		}
	}
}

function plural(n: number, word: string): string {
	return `${n} ${word}${n === 1 ? "" : "s"}`;
}

function structuredOutputSummary(value: Record<string, unknown>): string {
	if (typeof value.status === "string") {
		const summary = firstLine(value.title) ?? firstLine(value.summary);
		return cap(summary ? `${value.status} · ${summary}` : value.status);
	}
	if (Array.isArray(value.findings))
		return plural(value.findings.length, "finding");
	if (typeof value.real === "boolean") {
		const title = firstLine(value.title);
		return title ? `real=${value.real} · ${cap(title)}` : `real=${value.real}`;
	}
	const keys = Object.keys(value);
	return keys.length > 0 ? `returned ${keys.join(", ")}` : "returned output";
}

/**
 * Short preview of a tool result. Prefers `result.details` — the validated
 * structured-output payload path (P1) — falling back to `result.content`, then
 * the whole result.
 */
function contentText(content: unknown): string | undefined {
	if (typeof content === "string")
		return firstLine(formatStructuredText(content, summarizeValue));
	if (!Array.isArray(content)) return undefined;
	for (const part of content) {
		if (!isRecord(part) || part.type !== "text") continue;
		const text =
			typeof part.text === "string"
				? firstLine(formatStructuredText(part.text, summarizeValue))
				: undefined;
		if (text) return text;
	}
	return undefined;
}

function previewResult(result: unknown): string {
	if (result === undefined) return "";
	if (result === null) return "null";
	if (isRecord(result)) {
		if (result.details !== undefined && result.details !== null)
			return previewValue(result.details);
		const text = contentText(result.content);
		if (text) return cap(text);
		if (result.content !== undefined) return previewValue(result.content);
	}
	return previewValue(result);
}

/**
 * Map one Pi `AgentSessionEvent` to an `UltraProgressEvent`, stamping it with
 * the owning agent id + phase. Returns `null` for events outside the live-board
 * vocabulary (turn boundaries, message start/end, compaction, queue, etc.).
 * Pure: reads only the validated field paths; never constructs a session and
 * never attaches `controls`.
 */
export function mapSessionEvent(
	event: AgentSessionEvent,
	agentId: string,
	phase: string,
): AgentProgressEvent | null {
	const base = { agentId, phase } as const;

	switch (event.type) {
		case "agent_start":
			return { ...base, kind: "start" };

		case "tool_execution_start":
			return {
				...base,
				kind: "action",
				toolName: event.toolName,
				toolCallId: event.toolCallId,
				args: event.args,
				text: previewArgs(event.args),
			};

		case "tool_execution_update":
			return {
				...base,
				kind: "action",
				toolName: event.toolName,
				toolCallId: event.toolCallId,
				args: event.args,
				result: event.partialResult,
				isPartial: true,
				text: previewResult(event.partialResult),
			};

		case "tool_execution_end":
			return {
				...base,
				kind: "action",
				toolName: event.toolName,
				toolCallId: event.toolCallId,
				result: event.result,
				isError: event.isError,
				isPartial: false,
				text: previewResult(event.result),
				status: event.isError ? "error" : undefined,
			};

		case "message_update": {
			const inner = event.assistantMessageEvent;
			if (inner.type === "text_delta" || inner.type === "thinking_delta") {
				return { ...base, kind: "delta", text: inner.delta };
			}
			return null;
		}

		case "agent_end":
			return { ...base, kind: "end", status: "done" };

		default:
			return null;
	}
}

// ---------------------------------------------------------------------------
// Reducer — folds an UltraProgressEvent stream into board + transcript state
// ---------------------------------------------------------------------------

/** Lifecycle classification of an agent board row. */
export type AgentStatus = "queued" | "running" | "done" | "dropped";

/** One live board row: an agent's status + its latest action. */
export interface AgentRow {
	agentId: string;
	phase: string;
	status: AgentStatus;
	/** Latest coarse action ("read a.ts", "grep …") — never the streamed body. */
	action: string;
	tier?: string;
	model?: string;
	thinkingLevel?: ThinkingLevel;
	startedAt?: number;
	endedAt?: number;
	usage?: TokenUsage;
	failure?: StepFailure;
	item?: unknown;
	summary?: string;
	prompt?: string;
	/** This successful step was restored from the session journal. */
	cached?: boolean;
	/** Live controls for command-overlay rows; stripped from the tool surface. */
	controls?: AgentControls;
}

/** Folded planned and live view consumed by both `ui.ts` sinks. */
export interface ProgressMeta {
	workflowName?: string;
	workflowDescription?: string;
}

export type TranscriptItem =
	| { kind: "user"; text: string }
	| { kind: "text"; text: string }
	| {
			kind: "tool";
			toolName: string;
			toolCallId?: string;
			args?: unknown;
			result?: unknown;
			isError?: boolean;
			isPartial?: boolean;
			fallback: string;
	  };

export interface ReducerState extends ProgressMeta {
	rows: readonly AgentRow[];
	/** agentId -> accumulated transcript (text/thinking deltas + tool action lines). */
	transcripts: Readonly<Record<string, string>>;
	transcriptItems: Readonly<Record<string, readonly TranscriptItem[]>>;
	/** Planned phases in specification order. */
	phases: readonly PlannedPhase[];
	/** Current engine phase; initial planning does not select one. */
	phase: string | undefined;
	startedAt: number | undefined;
	endedAt: number | undefined;
	/** Known agents. */
	total: number;
	queued: number;
	running: number;
	/** Agents that ended successfully. */
	done: number;
	/** Agents that ended dropped/failed. */
	dropped: number;
	/** Phases whose agent count or run decision still depends on earlier results. */
	later: number;
	/** Final aggregate, attached only to the completed workflow tool result. */
	workflowResult?: WorkflowResult;
}

/** A fresh, empty reducer state. */
export function initialReducerState(meta: ProgressMeta = {}): ReducerState {
	return {
		workflowName: meta.workflowName,
		workflowDescription: meta.workflowDescription,
		rows: [],
		transcripts: {},
		transcriptItems: {},
		phases: [],
		phase: undefined,
		startedAt: undefined,
		endedAt: undefined,
		total: 0,
		queued: 0,
		running: 0,
		done: 0,
		dropped: 0,
		later: 0,
		workflowResult: undefined,
	};
}

// Statuses on an `end` event that mean the agent did not produce a usable
// result. Anything else (incl. "done"/undefined on an `end`) counts as done.
const DROP_STATUSES = new Set(["dropped", "error", "aborted", "failed"]);

function classifyEnd(status: string | undefined): AgentStatus {
	return status !== undefined && DROP_STATUSES.has(status) ? "dropped" : "done";
}

/** The coarse board label for an `action` event. */
function actionLabel(event: AgentProgressEvent): string {
	if (event.actionSummary) return event.actionSummary;
	if (event.result !== undefined) {
		const result = event.text ?? previewResult(event.result);
		if (event.isError) return result ? `Failed · ${result}` : "Failed";
		return result ? `Result · ${result}` : "Finished";
	}
	if (event.toolName === "structured_output") {
		const detail = isRecord(event.args)
			? structuredOutputSummary(event.args)
			: event.text;
		return detail ? `structured output · ${detail}` : "structured output";
	}
	if (event.toolName && event.args !== undefined)
		return actionIntent(event.toolName, event.args);
	if (event.toolName && event.text) return `${event.toolName} ${event.text}`;
	if (event.toolName) return `Using ${event.toolName}`;
	return event.text ?? "";
}

/** Text contributed to the per-agent transcript by this event (may be empty). */
function transcriptText(event: AgentProgressEvent): string {
	switch (event.kind) {
		case "delta":
			return event.text ?? "";
		case "action": {
			const label = actionLabel(event);
			return label ? `\n\n> ${label}\n\n` : "";
		}
		default:
			return event.text ?? "";
	}
}

function transcriptItem(event: AgentProgressEvent): TranscriptItem | undefined {
	if (event.kind === "prompt" && event.prompt)
		return { kind: "user", text: event.prompt };
	if (event.kind === "delta") return { kind: "text", text: event.text ?? "" };
	if (event.kind !== "action" || !event.toolName) return undefined;
	const fallback = actionLabel(event);
	return {
		kind: "tool",
		toolName: event.toolName,
		toolCallId: event.toolCallId,
		args: event.args,
		result: event.result,
		isError: event.isError,
		isPartial: event.isPartial,
		fallback,
	};
}

function nextRow(
	prev: AgentRow | undefined,
	event: AgentProgressEvent,
): AgentRow {
	const base: AgentRow = prev ?? {
		agentId: event.agentId,
		phase: event.phase,
		status: "running",
		action: "",
	};
	let status = base.status;
	let action = base.action;
	let cached = base.cached;
	switch (event.kind) {
		case "prompt":
			break;
		case "start":
			status = "running";
			if (event.text === "cached") {
				cached = true;
				action = "resumed from checkpoint";
			}
			break;
		case "action":
			if (event.result === undefined || event.isError)
				action = actionLabel(event) || action;
			break;
		case "delta":
			// Token deltas feed the transcript, not the board's latest-action line.
			break;
		case "end":
			status = classifyEnd(event.status);
			break;
	}
	return {
		agentId: event.agentId,
		phase: event.phase,
		status,
		action,
		tier: event.tier ?? base.tier,
		model: event.model ?? base.model,
		thinkingLevel: event.thinkingLevel ?? base.thinkingLevel,
		startedAt: event.startedAt ?? base.startedAt,
		endedAt: event.kind === "end" ? (base.endedAt ?? Date.now()) : base.endedAt,
		usage: event.usage ?? base.usage,
		failure: event.failure ?? base.failure,
		item: event.item ?? base.item,
		summary: event.summary ?? base.summary,
		prompt: event.prompt ?? base.prompt,
		cached,
		controls: event.controls ?? base.controls,
	};
}

/** A queued row seeded before its agent session exists. */
function queuedRow(agentId: string, phase: string): AgentRow {
	return { agentId, phase, status: "queued", action: "" };
}

function withCounts(
	state: Omit<
		ReducerState,
		"total" | "queued" | "running" | "done" | "dropped" | "later"
	>,
): ReducerState {
	let queued = 0;
	let running = 0;
	let done = 0;
	let dropped = 0;
	for (const row of state.rows) {
		if (row.status === "queued") queued++;
		else if (row.status === "running") running++;
		else if (row.status === "done") done++;
		else dropped++;
	}
	const later = state.phases.filter(
		(phase) =>
			phase.status === "agents-pending" || phase.status === "condition-pending",
	).length;
	const settled =
		state.phases.length > 0 &&
		state.phases.every(
			(phase) => phase.status === "done" || phase.status === "skipped",
		);
	return {
		...state,
		endedAt: settled ? (state.endedAt ?? Date.now()) : state.endedAt,
		total: state.rows.length,
		queued,
		running,
		done,
		dropped,
		later,
	};
}

function applyPlan(
	state: ReducerState,
	event: Extract<UltraProgressEvent, { kind: "plan" }>,
): ReducerState {
	const planned = event.phases.map((phase) => ({
		...phase,
		agentIds: [...phase.agentIds],
	}));
	return withCounts({
		...state,
		phases: planned,
		phase: undefined,
		startedAt: event.startedAt,
		endedAt: undefined,
		rows: planned.flatMap((phase) =>
			phase.agentIds.map((agentId) => queuedRow(agentId, phase.phase)),
		),
	});
}

function applyPhase(
	state: ReducerState,
	event: Extract<UltraProgressEvent, { kind: "phase" }>,
): ReducerState {
	const phase = {
		phase: event.phase,
		status: event.status,
		agentIds: [...event.agentIds],
	};
	const phases = state.phases.some((entry) => entry.phase === event.phase)
		? state.phases.map((entry) => (entry.phase === event.phase ? phase : entry))
		: [...state.phases, phase];
	const keep = new Set(event.agentIds);
	const rows = state.rows.filter(
		(row) => row.phase !== event.phase || keep.has(row.agentId),
	);
	for (const agentId of event.agentIds) {
		if (!rows.some((row) => row.agentId === agentId))
			rows.push(queuedRow(agentId, event.phase));
	}
	return withCounts({
		...state,
		phases,
		phase: event.phase,
		rows,
	});
}

/**
 * Pure reducer: fold planning, phase, and agent events into one ordered board.
 * Agent rows are upserted by id; counts are always derived from current state.
 */
export function reduceProgress(
	state: ReducerState,
	event: UltraProgressEvent,
): ReducerState {
	if (event.kind === "plan") return applyPlan(state, event);
	if (event.kind === "phase") return applyPhase(state, event);

	const idx = state.rows.findIndex((row) => row.agentId === event.agentId);
	const prev = idx >= 0 ? state.rows[idx] : undefined;
	const row = nextRow(prev, event);
	const rows = idx < 0 ? [...state.rows, row] : state.rows.slice();
	if (idx >= 0) rows[idx] = row;

	const added = transcriptText(event);
	const transcripts = added
		? {
				...state.transcripts,
				[event.agentId]: (state.transcripts[event.agentId] ?? "") + added,
			}
		: state.transcripts;

	const item = transcriptItem(event);
	let transcriptItems = state.transcriptItems;
	if (item) {
		const current = state.transcriptItems[event.agentId] ?? [];
		const toolIndex =
			item.kind === "tool" && item.toolCallId
				? current.findIndex(
						(entry) =>
							entry.kind === "tool" && entry.toolCallId === item.toolCallId,
					)
				: -1;
		const next =
			toolIndex < 0
				? [...current, item]
				: current.map((entry, index) =>
						index === toolIndex && entry.kind === "tool" && item.kind === "tool"
							? {
									...entry,
									args: item.args ?? entry.args,
									result: item.result ?? entry.result,
									isError: item.isError ?? entry.isError,
									isPartial: item.isPartial ?? entry.isPartial,
									fallback: item.fallback || entry.fallback,
								}
							: entry,
					);
		transcriptItems = {
			...state.transcriptItems,
			[event.agentId]: next,
		};
	}

	const phases =
		state.phase === event.phase ||
		state.phases.some((phase) => phase.phase === event.phase)
			? state.phases
			: [
					...state.phases,
					{
						phase: event.phase,
						status: "resolved" as const,
						agentIds: [event.agentId],
					},
				];

	// Agent events can change only one row status. Updating its counters directly
	// avoids rescanning every fanout row for every streamed token while preserving
	// immutable snapshots required by tool-result structured cloning.
	let { total, queued, running, done, dropped } = state;
	if (!prev || prev.status !== row.status) {
		const adjust = (status: AgentStatus, amount: number): void => {
			if (status === "queued") queued += amount;
			else if (status === "running") running += amount;
			else if (status === "done") done += amount;
			else dropped += amount;
		};
		if (!prev) total++;
		else adjust(prev.status, -1);
		adjust(row.status, 1);
	}

	return {
		workflowName: state.workflowName,
		workflowDescription: state.workflowDescription,
		rows,
		transcripts,
		transcriptItems,
		phases,
		phase: event.phase,
		startedAt: state.startedAt,
		endedAt: state.endedAt,
		total,
		queued,
		running,
		done,
		dropped,
		later: state.later,
		workflowResult: state.workflowResult,
	};
}