/** * Foreground run controller for pi-crew. * * Owns the three helpers consumed by the team-tool / subagent-tools / * team-commands modules: * • `startForegroundRun` — spawn a foreground team run with watchdog, * abort signal, and completion-handling (status notification + widget * refresh + crew.run.* event emission). * • `abortForegroundRun` — abort a foreground team run by runId. * • `openLiveSidebar` — open the per-run live sidebar UI overlay. * * These three are conceptually one unit (all manipulate the same * foreground-team-run controllers + live sidebar state), so they share * a single file. */ import type { ExtensionAPI, ExtensionContext } from "@earendil-works/pi-coding-agent"; import { loadConfig } from "../../config/config.ts"; import { loadRunManifestById, updateRunStatus } from "../../state/stores/state-store.ts"; import { setWorkingIndicator } from "../../ui/pi-ui-compat.ts"; import { requestPowerbarUpdate } from "../../ui/powerbar-publisher.ts"; import { updateCrewWidget } from "../../ui/widget/index.ts"; import { logInternalError } from "../../utils/internal-error.ts"; import { safeAbort } from "../../utils/safe-abort.ts"; import type { RegistrationContext } from "./registration-types.ts"; /** * Install the foreground-run controller helpers into the registration * context. Re-call safe — each call rebinds `ctx.startForegroundRun`, * `ctx.abortForegroundRun`, and `ctx.openLiveSidebar` to fresh closures * (the underlying state lives on `ctx`). */ export function installForegroundRunController(pi: ExtensionAPI, ctx: RegistrationContext): void { ctx.openLiveSidebar = (extCtx, runId) => { void openLiveSidebarImpl(pi, ctx, extCtx, runId); }; ctx.startForegroundRun = (extCtx, runner, runId) => startForegroundRunImpl(pi, ctx, extCtx, runner, runId); ctx.abortForegroundRun = (runId) => { const controller = ctx.foregroundTeamRunControllers.get(runId); if (!controller) return false; // safe-abort: cancel can race a completed run whose spawned children // already exited — Node's child_process signal listener would throw // AbortError out of abort() and kill the host (uncaughtException). safeAbort(controller, "fg-run.abort-by-run-id"); return true; }; } /** * Open the live sidebar for a specific run. * * Defers the actual install to registration/ui.ts via a lazy dynamic * import — only invoked when a sidebar is actually requested, never at * module load. */ async function openLiveSidebarImpl( pi: ExtensionAPI, ctx: RegistrationContext, extensionCtx: ExtensionContext, runId: string, ): Promise { // LAZY: registration/ui pulls in transcript-viewer + heavy UI modules. const uiModule = await import("./ui.ts"); await uiModule.installLiveSidebar(extensionCtx, runId, ctx.uiState, { pi, widgetState: ctx.widgetState, getManifestCache: ctx.getManifestCache, getRunSnapshotCache: ctx.getRunSnapshotCache, isCleanedUp: () => ctx.cleanedUp, getCurrentCtx: () => ctx.currentCtx, }); } /** * Spawn a foreground team run. * * The runner executes on the next macrotask (`setImmediate`) so this * function returns synchronously to Pi. The runner receives an AbortSignal * that fires when: * • `abortForegroundRun(runId)` is called (e.g., from a tool cancel), * • `cleanupRuntime()` fires (session shutdown with reason=quit/reload). * * A foreground-watchdog is started per runId to surface hung runs to the * assistant. The .finally block handles: abort propagation, working-message * clearing, status notification, run-completed session entry, and crew.run.* * event emission. */ function startForegroundRunImpl( pi: ExtensionAPI, ctx: RegistrationContext, extensionCtx: ExtensionContext, runner: (signal?: AbortSignal) => Promise, runId?: string, ): void { const ownerGeneration = ctx.captureSessionGeneration(); // Issue #62: tool-dispatched runs receive a SPREAD COPY of the session // context (withSessionId() always returns { ...ctx }), so the object- // identity check (isContextCurrent) was false from the very first moment // and the whole completion block below was silently skipped for every // foreground run started via the team/subagent tools — stale working // message, no crew:run-completed entry, no crew.run.* metrics. Liveness // for the reporting half is therefore session-ID based; the UI clear is // unconditional (try/catch guards a disposed ctx). const ownerSessionId = extensionCtx.sessionManager?.getSessionId?.(); const controller = new AbortController(); const key = runId ?? Symbol(); ctx.foregroundTeamRunControllers.set(key, controller); if (extensionCtx.hasUI) { setWorkingIndicator(extensionCtx, { frames: ["⣾", "⣽", "⣻", "⢿", "⡿", "⣟", "⣯", "⣷"], intervalMs: 80, }); extensionCtx.ui.setWorkingMessage(runId ? `pi-crew foreground run ${runId}...` : "pi-crew foreground run..."); } if (runId) { void import("../../runtime/foreground-watchdog.ts") // LAZY: defer watchdog module until first foreground run .then(({ startForegroundWatchdog }) => { startForegroundWatchdog({ pi, cwd: extensionCtx.cwd, runId }); }) .catch((error) => { logInternalError("register.foreground-watchdog-import", error); }); } setImmediate(() => { void runner(controller.signal) .catch((error) => { const message = error instanceof Error ? error.message : String(error); if (runId) { try { const loaded = loadRunManifestById(extensionCtx.cwd, runId); if ( loaded && loaded.manifest.status !== "completed" && loaded.manifest.status !== "failed" && loaded.manifest.status !== "cancelled" && loaded.manifest.status !== "blocked" ) updateRunStatus(loaded.manifest, "failed", message); } catch (statusError) { logInternalError("register.foreground-run-failure", statusError, `runId=${runId}`); } } if (ctx.isOwnerSessionCurrent(ownerGeneration, ownerSessionId)) { extensionCtx.ui.notify(`pi-crew foreground run failed: ${message}`, "error"); } else { logInternalError("register.foreground-run-failure", error, `runId=${runId} context disposed`); } }) .finally(() => { ctx.foregroundTeamRunControllers.delete(key); if (runId) { void import("../../runtime/foreground-watchdog.ts") // LAZY: defer watchdog module until first foreground run .then(({ stopWatchdog }) => { stopWatchdog(runId); }) .catch((error) => logInternalError("register.foreground-watchdog", error, `runId=${runId}`)); } const ownerCurrent = ctx.isOwnerSessionCurrent(ownerGeneration, ownerSessionId); // Issue #62: clear the working message UNCONDITIONALLY. The `hasUI` // getter THROWS on a disposed ctx, so touch it inside try/catch — // a stale ctx just means the host already reset its UI; a live one // must never keep spinning a finished run's label. try { if (extensionCtx.hasUI) { setWorkingIndicator(extensionCtx); extensionCtx.ui.setWorkingMessage(); } } catch { /* disposed ctx — host-side UI reset already happened */ } if (ownerCurrent && runId) { const loaded = loadRunManifestById(extensionCtx.cwd, runId); const status = loaded?.manifest.status ?? "finished"; const level = status === "failed" || status === "blocked" ? "error" : status === "cancelled" ? "warning" : "info"; extensionCtx.ui.notify( `pi-crew run ${runId} ${status}. Use /team-summary ${runId} or /team-status ${runId}.`, level as "info" | "warning" | "error", ); pi.appendEntry("crew:run-completed", { runId, team: loaded?.manifest.team, workflow: loaded?.manifest.workflow, goal: loaded?.manifest.goal, status, taskCount: loaded?.tasks.length, durationMs: loaded?.manifest.createdAt ? Date.now() - new Date(loaded.manifest.createdAt).getTime() : 0, timestamp: Date.now(), }); const eventType = status === "completed" ? "crew.run.completed" : status === "failed" || status === "blocked" ? "crew.run.failed" : status === "cancelled" ? "crew.run.cancelled" : undefined; if (eventType) { pi.events?.emit?.(eventType, { runId, team: loaded?.manifest.team, workflow: loaded?.manifest.workflow, status, taskCount: loaded?.tasks.length, goal: loaded?.manifest.goal, durationMs: loaded?.manifest.createdAt ? Date.now() - new Date(loaded.manifest.createdAt).getTime() : 0, }); } // Emit per-task events so task-level metrics (crew.task.count, // crew.task.duration_ms, crew.task.tokens_total) record actual values. const terminalTaskStatuses = new Set(["completed", "failed", "needs_attention"]); for (const task of loaded?.tasks ?? []) { if (!terminalTaskStatuses.has(task.status)) continue; const taskEventType = task.status === "completed" ? "crew.task.completed" : task.status === "failed" ? "crew.task.failed" : "crew.task.needs_attention"; pi.events?.emit?.(taskEventType, { runId, team: loaded?.manifest.team, workflow: loaded?.manifest.workflow, taskId: task.id, taskCount: loaded?.tasks.length, role: task.role, durationMs: task.startedAt && task.finishedAt ? new Date(task.finishedAt).getTime() - new Date(task.startedAt).getTime() : 0, tokens: (task.usage?.input ?? 0) + (task.usage?.output ?? 0), }); } } if (ownerCurrent && ctx.currentCtx) { const config = loadConfig(ctx.currentCtx.cwd).config.ui; updateCrewWidget( ctx.currentCtx, ctx.widgetState, config, ctx.getManifestCache(ctx.currentCtx.cwd), ctx.getRunSnapshotCache(ctx.currentCtx.cwd), ); requestPowerbarUpdate( pi.events, ctx.currentCtx.cwd, config, ctx.getManifestCache(ctx.currentCtx.cwd), ctx.getRunSnapshotCache(ctx.currentCtx.cwd), ctx.currentCtx, ctx.widgetState.notificationCount ?? 0, ); } }); }); }