import type { Component, TUI } from "@earendil-works/pi-tui"; import { formatWorkflowProgress, formatWorkflowProgressStatus, formatWorkflowRunSummary, } from "./workflow-display.js"; import type { WorkflowAgentChild, WorkflowLogEntry, WorkflowMeta, WorkflowProgressEvent, WorkflowProgressState, WorkflowRunResult, } from "./workflow-types.js"; export type WorkflowProgressStatus = | "running" | "completed" | "stopped" | "error"; export type WorkflowAgentChildSource = | Iterable | ReadonlyMap; export interface WorkflowProgressEnvelope extends WorkflowProgressState { result: unknown; duration: number; status: WorkflowProgressStatus; error?: string; } export interface WorkflowProgressCompletion { envelope: WorkflowProgressEnvelope; summary: string; } export interface WorkflowProgressErrorSummary { envelope: WorkflowProgressEnvelope; summary: string; } export interface WorkflowProgressTrackerOptions { startedAt?: number; initialMeta?: WorkflowMeta; maxLogEntries?: number; now?: () => number; } export interface WorkflowProgressUiApi { setStatus?: (key: string, text: string | undefined) => void; setWidget?: ( key: string, content: string[] | undefined, options?: WorkflowProgressWidgetOptions ) => void; setThemedWidget?: WorkflowProgressThemedWidgetSetter; } export interface WorkflowProgressTheme { fg(color: string, text: string): string; } export interface WorkflowProgressWidgetOptions { placement?: "aboveEditor" | "belowEditor"; } export type WorkflowProgressWidgetFactory = ( tui: TUI, theme: WorkflowProgressTheme ) => Component & { dispose?(): void }; export type WorkflowProgressThemedWidgetSetter = ( key: string, content: WorkflowProgressWidgetFactory | undefined, options?: WorkflowProgressWidgetOptions ) => void; export interface WorkflowProgressUiContext { hasUI?: boolean; ui?: WorkflowProgressUiApi; } export interface WorkflowProgressUiUpdate { content: Array<{ type: "text"; text: string }>; details: WorkflowProgressEnvelope; } export interface WorkflowProgressUiOptions { onUpdate?: (result: WorkflowProgressUiUpdate) => void; statusKey?: string; widgetKey?: string; widgetPlacement?: "aboveEditor" | "belowEditor"; } export interface WorkflowProgressUiController { update(envelope: WorkflowProgressEnvelope): void; clear(): void; } export interface WorkflowProgressErrorOptions { stopped?: boolean; summary?: string; message?: string; } const DEFAULT_INITIAL_WORKFLOW_META: WorkflowMeta = { name: "pending", description: "Pending workflow parse", }; const WORKFLOW_PROGRESS_LOG_HISTORY_LIMIT = 20; const WORKFLOW_PROGRESS_SPINNER = [ "✳", "✴", "✵", "✶", "✷", "✸", "✹", "✺", "✻", "✼", "✽", ]; const WORKFLOW_PROGRESS_SPINNER_INTERVAL_MS = 100; const WORKFLOW_PROGRESS_CALL_LINE_PATTERN = /^(\s+#\d+ )(○|●|✓|✗|-)( .*)$/; const WORKFLOW_PROGRESS_PHASE_LINE_PATTERN = /^(\s+)(▶|✓| )( .*)$/; const WORKFLOW_PROGRESS_PHASE_SUMMARY_PATTERN = / \d+\/\d+(?: · .*)?$/; export class WorkflowProgressTracker { private readonly childAgents = new Map(); private readonly maxLogEntries: number; private readonly now: () => number; private readonly startedAt: number; private latestProgress: WorkflowProgressState; constructor(options: WorkflowProgressTrackerOptions = {}) { this.now = options.now ?? Date.now; this.startedAt = options.startedAt ?? this.now(); this.maxLogEntries = options.maxLogEntries ?? WORKFLOW_PROGRESS_LOG_HISTORY_LIMIT; this.latestProgress = { meta: options.initialMeta ?? DEFAULT_INITIAL_WORKFLOW_META, phases: [], logs: [], agents: [], agentCalls: [], }; } updateFromProgressEvent( event: WorkflowProgressEvent ): WorkflowProgressEnvelope { this.latestProgress = mergeWorkflowProgress( this.latestProgress, event, this.childAgents.values(), this.maxLogEntries ); return this.getEnvelope(); } updateChildAgent(child: WorkflowAgentChild): WorkflowProgressEnvelope { this.childAgents.set(child.id, child); this.latestProgress = { ...this.latestProgress, agents: mergeWorkflowAgents( this.latestProgress.agents, this.childAgents.values() ), }; return this.getEnvelope(); } getEnvelope( status: WorkflowProgressStatus = "running", fields: Partial = {} ): WorkflowProgressEnvelope { return { meta: fields.meta ?? this.latestProgress.meta, phases: fields.phases ?? this.latestProgress.phases, logs: fields.logs ?? this.latestProgress.logs, agents: fields.agents ?? this.latestProgress.agents, agentCalls: fields.agentCalls ?? this.latestProgress.agentCalls, result: fields.result ?? null, duration: fields.duration ?? this.now() - this.startedAt, status, currentPhase: fields.currentPhase ?? this.latestProgress.currentPhase, currentTask: fields.currentTask ?? this.latestProgress.currentTask, error: fields.error, }; } getProgressText(status: WorkflowProgressStatus = "running"): string { return formatWorkflowProgress(this.getEnvelope(status)); } complete(result: WorkflowRunResult): WorkflowProgressCompletion { const envelope = this.getEnvelope("completed", { meta: result.meta, phases: result.phases, logs: result.logs, agents: mergeWorkflowAgents( result.agentChildren, this.childAgents.values() ), agentCalls: result.agentCalls, result: result.value, }); this.latestProgress = envelope; return { envelope, summary: formatWorkflowRunSummary(result), }; } error( error: unknown, options: WorkflowProgressErrorOptions = {} ): WorkflowProgressErrorSummary { const message = options.message ?? (error instanceof Error ? error.message : String(error)); const envelope = this.getEnvelope(options.stopped ? "stopped" : "error", { error: message, }); this.latestProgress = envelope; return { envelope, summary: options.summary ?? (options.stopped ? "Workflow cancelled." : "Workflow failed."), }; } } export function createWorkflowProgressTracker( options?: WorkflowProgressTrackerOptions ): WorkflowProgressTracker { return new WorkflowProgressTracker(options); } export function createWorkflowProgressUi( ctx: WorkflowProgressUiContext, options: WorkflowProgressUiOptions = {} ): WorkflowProgressUiController | undefined { if (ctx.hasUI === false) { return undefined; } const { onUpdate, statusKey = "workflow", widgetKey = "workflow-progress", widgetPlacement = "aboveEditor", } = options; const ui = ctx.ui; const setThemedWidgetCapability = ui?.setThemedWidget ?? (ui?.setWidget as WorkflowProgressThemedWidgetSetter | undefined); let lastWidgetText: string | undefined; let lastPartialText: string | undefined; let lastStatusText: string | undefined; let latestEnvelope: WorkflowProgressEnvelope | undefined; let latestEnvelopeUpdatedAt = 0; let tui: TUI | undefined; let widgetFrame = 0; let widgetInterval: ReturnType | undefined; let widgetRegistered = false; let widgetRegisteredWithThemedApi = false; let themedWidgetUnavailable = false; let active = false; const stopSpinnerTimer = () => { if (widgetInterval) { clearInterval(widgetInterval); widgetInterval = undefined; } }; const ensureSpinnerTimer = () => { if (widgetInterval) { return; } widgetInterval = setInterval(() => { widgetFrame += 1; tui?.requestRender?.(); }, WORKFLOW_PROGRESS_SPINNER_INTERVAL_MS); if (typeof widgetInterval === "object") { widgetInterval.unref?.(); } }; const clear = () => { stopSpinnerTimer(); if (!active) { return; } if (widgetRegisteredWithThemedApi) { ui?.setThemedWidget?.(widgetKey, undefined); } else { ui?.setWidget?.(widgetKey, undefined); } ui?.setStatus?.(statusKey, undefined); lastWidgetText = undefined; lastStatusText = undefined; latestEnvelope = undefined; latestEnvelopeUpdatedAt = 0; tui = undefined; widgetFrame = 0; widgetRegistered = false; widgetRegisteredWithThemedApi = false; active = false; }; const setPlainWidget = (widgetText: string) => { if (widgetText === lastWidgetText) { return; } ui?.setWidget?.(widgetKey, widgetText.split("\n"), { placement: widgetPlacement, }); lastWidgetText = widgetText; widgetRegisteredWithThemedApi = false; active = true; }; const setThemedWidget = ( envelope: WorkflowProgressEnvelope, widgetText: string ) => { latestEnvelope = envelope; latestEnvelopeUpdatedAt = Date.now(); if (themedWidgetUnavailable || !setThemedWidgetCapability) { setPlainWidget(widgetText); return; } if (widgetRegistered) { const widgetTextChanged = widgetText !== lastWidgetText; lastWidgetText = widgetText; ensureSpinnerTimer(); if (widgetTextChanged) { tui?.requestRender?.(); } active = true; return; } try { setThemedWidgetCapability( widgetKey, (nextTui, theme) => { tui = nextTui; return { render: (_width: number) => latestEnvelope ? renderThemedWorkflowProgress( latestEnvelope, theme, widgetFrame, latestEnvelopeUpdatedAt ) : [], invalidate: () => { widgetRegistered = false; tui = undefined; }, }; }, { placement: widgetPlacement } ); lastWidgetText = widgetText; widgetRegistered = true; widgetRegisteredWithThemedApi = Boolean(ui?.setThemedWidget); active = true; ensureSpinnerTimer(); } catch { themedWidgetUnavailable = true; widgetRegistered = false; tui = undefined; stopSpinnerTimer(); setPlainWidget(widgetText); } }; return { update(envelope) { if (isTerminalWorkflowProgressStatus(envelope.status)) { clear(); return; } const partialText = formatWorkflowProgressStatus(envelope); if (partialText !== lastPartialText && onUpdate) { onUpdate({ content: [{ type: "text", text: partialText }], details: envelope, }); lastPartialText = partialText; } const widgetText = formatWorkflowProgress(envelope); setThemedWidget(envelope, widgetText); if (partialText !== lastStatusText) { ui?.setStatus?.(statusKey, partialText); lastStatusText = partialText; active = true; } }, clear, }; } function renderThemedWorkflowProgress( envelope: WorkflowProgressEnvelope, theme: WorkflowProgressTheme, frame: number, envelopeUpdatedAt: number ): string[] { const spinner = WORKFLOW_PROGRESS_SPINNER[frame % WORKFLOW_PROGRESS_SPINNER.length]; const now = Date.now(); const workflowDuration = envelope.duration + Math.max(0, now - envelopeUpdatedAt); const liveEnvelope = { ...envelope, duration: workflowDuration }; const phaseOccurrences = new Map(); return formatWorkflowProgress(liveEnvelope) .split("\n") .map((line) => renderThemedWorkflowProgressLine( line, envelope, theme, spinner, now, workflowDuration, phaseOccurrences ) ); } function renderThemedWorkflowProgressLine( line: string, envelope: WorkflowProgressEnvelope, theme: WorkflowProgressTheme, spinner: string, now: number, workflowDuration: number, phaseOccurrences: Map ): string { if (line === "") { return line; } if (line.startsWith("◆ Workflow:")) { return colorWorkflowErrorCounts( theme, theme.fg("accent", `${spinner}${line.slice(1)}`) ); } if (line.startsWith(" log:") || line.startsWith(" active:")) { return theme.fg("dim", line); } if (line.trimStart().startsWith("… ")) { return theme.fg("dim", line); } if (line.startsWith(" error:")) { return theme.fg("error", line); } if (line.startsWith(" ")) { return theme.fg("dim", line); } const callMatch = WORKFLOW_PROGRESS_CALL_LINE_PATTERN.exec(line); if (callMatch) { const [, prefix, icon, rest] = callMatch; const call = getWorkflowProgressCall(envelope, prefix); const duration = getScopedWorkflowProgressDuration(call, now); switch (icon) { case "●": return `${theme.fg("dim", prefix)}${theme.fg( "accent", `${spinner} ${rest.trimStart()}…` )} ${theme.fg( "dim", formatWorkflowProgressElapsed(duration ?? workflowDuration) )}`; case "✓": return `${theme.fg("dim", prefix)}${theme.fg( "success", icon )}${theme.fg("dim", rest)}${formatCompletedWorkflowProgressElapsed( duration, theme )}`; case "✗": return `${theme.fg("dim", prefix)}${theme.fg("error", `${icon}${rest}`)}${formatCompletedWorkflowProgressElapsed( duration, theme )}`; case "-": return `${theme.fg("dim", line)}${formatCompletedWorkflowProgressElapsed( duration, theme )}`; default: return theme.fg("dim", line); } } const phaseMatch = WORKFLOW_PROGRESS_PHASE_LINE_PATTERN.exec(line); if (phaseMatch) { const [, indent, marker, rest] = phaseMatch; const phaseName = formatActivePhaseLabel(rest); const occurrence = phaseOccurrences.get(phaseName) ?? 0; phaseOccurrences.set(phaseName, occurrence + 1); const phase = envelope.phases .filter((entry) => entry.name === phaseName) .sort((left, right) => left.index - right.index)[occurrence]; const duration = getScopedWorkflowProgressDuration(phase, now); if (phase?.completedAt !== undefined) { return `${indent}${theme.fg("success", "✓")}${theme.fg( "dim", rest )}${formatCompletedWorkflowProgressElapsed(duration, theme)}`; } if (marker === "▶") { return `${indent}${theme.fg( "accent", `${spinner} ${phaseName}…` )} ${theme.fg( "dim", formatWorkflowProgressElapsed(duration ?? workflowDuration) )}`; } if (marker === "✓") { return `${indent}${theme.fg("success", marker)}${theme.fg( "dim", rest )}${formatCompletedWorkflowProgressElapsed(duration, theme)}`; } return theme.fg("dim", line); } return colorWorkflowErrorCounts(theme, line); } function colorWorkflowErrorCounts( theme: WorkflowProgressTheme, line: string ): string { return line.replace(/\d+ errors?/g, (match) => theme.fg("error", match)); } function formatActivePhaseLabel(line: string): string { return line.trimStart().replace(WORKFLOW_PROGRESS_PHASE_SUMMARY_PATTERN, ""); } function getWorkflowProgressCall( envelope: WorkflowProgressEnvelope, prefix: string ) { const index = Number.parseInt(prefix.trimStart().slice(1), 10) - 1; return envelope.agentCalls?.find((call) => call.index === index); } function getScopedWorkflowProgressDuration( scope: { startedAt?: number; completedAt?: number } | undefined, now: number ): number | undefined { if ( typeof scope?.startedAt !== "number" || !Number.isFinite(scope.startedAt) ) { return undefined; } const endedAt = typeof scope.completedAt === "number" && Number.isFinite(scope.completedAt) ? scope.completedAt : now; return Math.max(0, endedAt - scope.startedAt); } function formatCompletedWorkflowProgressElapsed( duration: number | undefined, theme: WorkflowProgressTheme ): string { return duration === undefined ? "" : ` ${theme.fg("dim", formatWorkflowProgressElapsed(duration))}`; } function formatWorkflowProgressElapsed(duration: number | undefined): string { if ( typeof duration !== "number" || !Number.isFinite(duration) || duration < 0 ) { return "0.0s"; } return `${(duration / 1000).toFixed(1)}s`; } function isTerminalWorkflowProgressStatus( status: WorkflowProgressStatus ): boolean { return status === "completed" || status === "stopped" || status === "error"; } export function mergeWorkflowProgress( previous: WorkflowProgressState, event: WorkflowProgressEvent, childAgents: WorkflowAgentChildSource = [], maxLogEntries = WORKFLOW_PROGRESS_LOG_HISTORY_LIMIT ): WorkflowProgressState { return { meta: event.meta, currentPhase: event.currentPhase, currentTask: event.currentTask, phases: event.phases, logs: mergeWorkflowLogs(previous.logs, event.logs, maxLogEntries), agents: mergeWorkflowAgents(event.agents, childAgents), agentCalls: event.agentCalls ?? previous.agentCalls, }; } export function mergeWorkflowLogs( previous: WorkflowLogEntry[], next: WorkflowLogEntry[], limit = WORKFLOW_PROGRESS_LOG_HISTORY_LIMIT ): WorkflowLogEntry[] { if (next.length === 0) { return previous; } const logsByIndex = new Map(); for (const log of previous) { logsByIndex.set(log.index, log); } for (const log of next) { logsByIndex.set(log.index, log); } return [...logsByIndex.values()] .sort((left, right) => left.index - right.index) .slice(-limit); } export function mergeWorkflowAgents( agents: WorkflowAgentChild[], childAgents: WorkflowAgentChildSource ): WorkflowAgentChild[] { const merged = new Map(); for (const agent of agents) { merged.set(agent.id, agent); } for (const agent of getWorkflowAgentChildValues(childAgents)) { merged.set(agent.id, agent); } return [...merged.values()]; } function getWorkflowAgentChildValues( childAgents: WorkflowAgentChildSource ): Iterable { if (isWorkflowAgentChildMap(childAgents)) { return childAgents.values(); } return childAgents; } function isWorkflowAgentChildMap( childAgents: WorkflowAgentChildSource ): childAgents is ReadonlyMap { return typeof (childAgents as { get?: unknown }).get === "function"; } export function workflowProgressResult( envelope: WorkflowProgressEnvelope, header: string ) { const safeEnvelope = toJsonSafeWorkflowEnvelope(envelope); return { content: [ { type: "text" as const, text: `${header}\n\n${JSON.stringify(safeEnvelope, null, 2)}`, }, ], details: safeEnvelope, }; } export function toJsonSafeWorkflowEnvelope( envelope: WorkflowProgressEnvelope ): WorkflowProgressEnvelope { return toJsonSafeValue(envelope) as WorkflowProgressEnvelope; } function toJsonSafeValue( value: unknown, seen = new WeakMap(), path = "$" ): unknown { if (value === null) { return null; } switch (typeof value) { case "string": case "boolean": return value; case "number": return Number.isFinite(value) ? value : null; case "bigint": return unsupportedJsonValue("bigint", value.toString()); case "function": return unsupportedJsonValue("function"); case "symbol": return unsupportedJsonValue("symbol", value.description); case "undefined": return undefined; case "object": break; } const objectValue = value as object; const existingPath = seen.get(objectValue); if (existingPath) { return unsupportedJsonValue("circular", existingPath); } seen.set(objectValue, path); if (Array.isArray(value)) { const result = value.map((item, index) => toJsonSafeValue(item, seen, `${path}[${index}]`) ); seen.delete(objectValue); return result; } const result: Record = {}; for (const key of Object.keys(value)) { const descriptor = Object.getOwnPropertyDescriptor(value, key); if (!descriptor) { continue; } if (!("value" in descriptor)) { result[key] = unsupportedJsonValue("accessor"); continue; } const safeValue = toJsonSafeValue(descriptor.value, seen, `${path}.${key}`); if (safeValue !== undefined) { result[key] = safeValue; } } seen.delete(objectValue); return result; } function unsupportedJsonValue(type: string, value?: string) { return { unsupported: true, type, ...(value === undefined ? {} : { value }), }; }