import type { AgentMessage } from "@earendil-works/pi-agent-core"; import type { ToolResultMessage } from "@earendil-works/pi-ai"; import type { ExtensionAPI, ExtensionCommandContext, ExtensionContext, SessionStartEvent, } from "@earendil-works/pi-coding-agent"; import type { OverlayHandle } from "@earendil-works/pi-tui"; import { DEFAULT_CONFIG, loadTodoObserverConfig, type LoadedTodoObserverConfig, type TodoObserverConfig, } from "./config.ts"; import { formatBranchContextChunks, formatTurnContext, userMessageText } from "./context.ts"; import { TodoObserver } from "./observer.ts"; import { TodoPanel } from "./todo-panel.ts"; import { TODO_STATE_ENTRY, normalizeTodos, type ObserverStatus, type TodoNode, type TodoObserverState, } from "./types.ts"; interface Observation { label: string; context: string; /** Recalculate from this state instead of the currently displayed state. */ seedTodos?: TodoNode[]; } class ObservationWorker { private queue: Observation[] = []; private draining = false; private disposed = false; private observer: TodoObserver | undefined; constructor( private readonly cwd: string, private readonly config: TodoObserverConfig, private readonly getTodos: () => TodoNode[], private readonly publish: (todos: TodoNode[], label: string) => void, private readonly setStatus: (status: ObserverStatus) => void, ) {} enqueue(observation: Observation): void { if (this.disposed || !observation.context.trim()) return; this.queue.push(observation); this.updateWorkingStatus(); void this.drain(); } dispose(): void { if (this.disposed) return; this.disposed = true; this.queue = []; this.observer?.dispose(); this.observer = undefined; } private updateWorkingStatus(): void { const waiting = this.queue.length; this.setStatus({ kind: "observing", message: waiting > 1 ? `observer updating (${waiting} waiting)…` : "observer updating…", }); } private async drain(): Promise { if (this.draining || this.disposed) return; this.draining = true; let lastObservationFailed = false; try { let observer = this.observer; if (!observer) { observer = await TodoObserver.create(this.cwd, this.config); this.observer = observer; } while (!this.disposed && this.queue.length > 0) { const observation = this.queue.shift()!; this.updateWorkingStatus(); try { const seed = observation.seedTodos ?? this.getTodos(); const nextTodos = await observer.observe(observation.label, observation.context, seed); if (!this.disposed) this.publish(nextTodos, observation.label); lastObservationFailed = false; } catch (error) { if (this.disposed) return; lastObservationFailed = true; const message = error instanceof Error ? error.message : String(error); this.setStatus({ kind: "error", message }); observer.dispose(); this.observer = undefined; // A later queued turn gets a fresh observer session seeded from current state. if (this.queue.length > 0) { observer = await TodoObserver.create(this.cwd, this.config); this.observer = observer; } } } if (!this.disposed && !lastObservationFailed) this.setStatus({ kind: "idle", message: "observer ready" }); } catch (error) { if (!this.disposed) { const message = error instanceof Error ? error.message : String(error); this.setStatus({ kind: "error", message }); this.queue = []; } } finally { this.draining = false; } } } function restoreState( ctx: ExtensionContext, maxTodos: number, ): { todos: TodoNode[]; found: boolean } { let restored: TodoNode[] = []; let found = false; for (const entry of ctx.sessionManager.getBranch()) { if (entry.type !== "custom" || entry.customType !== TODO_STATE_ENTRY) continue; const state = entry.data as Partial | undefined; if (!state || !Array.isArray(state.todos)) continue; restored = normalizeTodos(state.todos, maxTodos); found = true; } return { todos: restored, found }; } /** Recover the latest non-empty snapshot if a bad recalculation published an empty list. */ function restoreMostRecentNonEmptyState(ctx: ExtensionContext, maxTodos: number): TodoNode[] { let restored: TodoNode[] = []; for (const entry of ctx.sessionManager.getBranch()) { if (entry.type !== "custom" || entry.customType !== TODO_STATE_ENTRY) continue; const state = entry.data as Partial | undefined; if (!state || !Array.isArray(state.todos)) continue; const candidate = normalizeTodos(state.todos, maxTodos); if (candidate.length > 0) restored = candidate; } return restored; } export default function todoObserverExtension(pi: ExtensionAPI): void { let config = DEFAULT_CONFIG; let loadedConfig: LoadedTodoObserverConfig | undefined; let todos: TodoNode[] = []; let panel: TodoPanel | undefined; let overlayHandle: OverlayHandle | undefined; let worker: ObservationWorker | undefined; let currentCwd = ""; let latestUserText: string | undefined; let panelManuallyHidden = false; let runtimeActive = false; const setPanelStatus = (status: ObserverStatus) => panel?.setStatus(status); const publishTodos = (nextTodos: TodoNode[], observedTurn: string) => { if (!runtimeActive) return; todos = normalizeTodos(nextTodos, config.observer.maxTodos); panel?.setTodos(todos); const state: TodoObserverState = { version: 1, todos, updatedAt: new Date().toISOString(), observedTurn, }; try { pi.appendEntry(TODO_STATE_ENTRY, state); } catch { // The main runtime may have been replaced while the observer request was in flight. } }; const createWorker = () => { worker?.dispose(); worker = new ObservationWorker( currentCwd, config, () => todos.map((todo) => ({ ...todo })), publishTodos, setPanelStatus, ); }; const enqueueFullBranch = ( ctx: ExtensionContext, label: string, seedTodos: TodoNode[], ) => { const chunks = formatBranchContextChunks(ctx, config.context); for (let index = 0; index < chunks.length; index++) { worker?.enqueue({ label: `${label} (${index + 1}/${chunks.length})`, context: chunks[index]!, seedTodos: index === 0 ? seedTodos : undefined, }); } if (chunks.length === 0) setPanelStatus({ kind: "idle", message: "no session history to recalculate" }); }; const mountOverlay = (ctx: ExtensionContext) => { if (ctx.mode !== "tui") return; void ctx.ui .custom( (tui, theme) => { panel = new TodoPanel(tui, theme, todos); return panel; }, { overlay: true, overlayOptions: { anchor: "top-right", width: config.sidebar.width, maxHeight: config.sidebar.maxHeight, margin: { top: 0, right: 0 }, nonCapturing: true, visible: (terminalWidth) => terminalWidth >= config.sidebar.minTerminalWidth, }, onHandle: (handle) => { overlayHandle = handle; if (panelManuallyHidden) handle.setHidden(true); }, }, ) .catch((error) => { if (!runtimeActive) return; const message = error instanceof Error ? error.message : String(error); ctx.ui.notify(`Todo observer overlay failed: ${message}`, "error"); }); }; const startRuntime = (event: SessionStartEvent, ctx: ExtensionContext) => { if (ctx.mode !== "tui") return; runtimeActive = true; currentCwd = ctx.cwd; latestUserText = undefined; loadedConfig = loadTodoObserverConfig(ctx.cwd, ctx.isProjectTrusted()); config = loadedConfig.config; panelManuallyHidden = !config.sidebar.showOnStart; for (const warning of loadedConfig.warnings) ctx.ui.notify(`Todo observer config: ${warning}`, "warning"); const restored = restoreState(ctx, config.observer.maxTodos); todos = restored.todos; panel?.setTodos(todos); mountOverlay(ctx); createWorker(); // Bootstrap older sessions that predate this extension or lack an observer snapshot. if (config.observer.bootstrapOnStart && !restored.found && ctx.sessionManager.getBranch().length > 0) { enqueueFullBranch(ctx, `bootstrap after ${event.reason}`, []); } }; const hideSidebar = (ctx: ExtensionContext) => { if (ctx.mode !== "tui" || !overlayHandle) { ctx.ui.notify("Todo observer sidebar is only available in interactive mode", "warning"); return; } panelManuallyHidden = true; overlayHandle.setHidden(true); }; const showSidebar = (ctx: ExtensionContext) => { if (ctx.mode !== "tui" || !overlayHandle) { ctx.ui.notify("Todo observer sidebar is only available in interactive mode", "warning"); return; } panelManuallyHidden = false; overlayHandle.setHidden(false); }; const recalculateTodos = async (_args: string, ctx: ExtensionCommandContext) => { if (ctx.mode !== "tui") { ctx.ui.notify("Todo observer is only active in interactive mode", "warning"); return; } // Recalculation must reconcile history, not erase it. Seed a fresh observer // with the displayed list, or recover the latest non-empty session snapshot // if an earlier bad recalculation already published an empty list. const historicalSeed = todos.length > 0 ? todos.map((todo) => ({ ...todo })) : restoreMostRecentNonEmptyState(ctx, config.observer.maxTodos); createWorker(); enqueueFullBranch( ctx, "manual full-session recalculation; retain historical completed todos", historicalSeed, ); }; pi.on("session_start", async (event, ctx) => startRuntime(event, ctx)); pi.on("message_end", async (event) => { if (!runtimeActive || event.message.role !== "user") return; latestUserText = userMessageText(event.message as AgentMessage, config.context); }); pi.on("turn_end", async (event) => { if (!runtimeActive || !worker) return; const context = formatTurnContext( latestUserText, event.message as AgentMessage, event.toolResults as ToolResultMessage[], config.context, ); latestUserText = undefined; worker.enqueue({ label: `main turn ${event.turnIndex + 1}`, context }); }); pi.on("session_tree", async (_event, ctx) => { if (ctx.mode !== "tui") return; const restored = restoreState(ctx, config.observer.maxTodos); todos = restored.todos; panel?.setTodos(todos); currentCwd = ctx.cwd; latestUserText = undefined; createWorker(); enqueueFullBranch(ctx, "main session branch changed", restored.todos); }); pi.on("session_shutdown", async () => { runtimeActive = false; worker?.dispose(); worker = undefined; overlayHandle?.hide(); overlayHandle = undefined; panel = undefined; }); pi.registerCommand("todo-clear", { description: "Clear all todos while keeping the observer active", handler: async (_args, ctx) => { if (ctx.mode !== "tui") { ctx.ui.notify("Todo observer is only active in interactive mode", "warning"); return; } worker?.dispose(); publishTodos([], "manual clear"); createWorker(); setPanelStatus({ kind: "idle", message: "todos cleared" }); }, }); pi.registerCommand("todo-recalc", { description: "Recalculate todos from the complete active session branch", handler: recalculateTodos, }); // Backward-compatible alias for the original refresh command. pi.registerCommand("todo-observer-refresh", { description: "Recalculate todos from the complete active session branch", handler: recalculateTodos, }); pi.registerCommand("todo-hide", { description: "Hide the todo sidebar while continuing to track todos", handler: async (_args, ctx) => hideSidebar(ctx), }); pi.registerCommand("todo-show", { description: "Show the todo sidebar", handler: async (_args, ctx) => showSidebar(ctx), }); pi.registerCommand("todo-sidebar", { description: "Toggle the persistent todo observer sidebar", handler: async (_args, ctx) => { if (panelManuallyHidden) showSidebar(ctx); else hideSidebar(ctx); }, }); pi.registerCommand("todo-config", { description: "Show the effective todo observer configuration", handler: async (_args, ctx) => { const paths = [loadedConfig?.globalPath, loadedConfig?.projectPath].filter(Boolean).join("\n"); ctx.ui.notify( `Observer: ${config.provider}/${config.model}:${config.thinkingLevel}\nSidebar: ${config.sidebar.width} cols, visible at ${config.sidebar.minTerminalWidth}+ cols\nConfig:\n${paths || "not loaded"}`, "info", ); }, }); pi.registerShortcut("ctrl+alt+t", { description: "Toggle todo observer sidebar", handler: async (ctx) => { if (panelManuallyHidden) showSidebar(ctx); else hideSidebar(ctx); }, }); }