/** * pi-tool-wal — Tool Write-Ahead Log extension for Pi * * Records every tool-call intent before execution, then reconciles the actual * result afterward. Persists to SQLite at ~/.pi/agent/tool-wal/wal.db so the * log survives across sessions for audit and crash recovery. * * Isolation (single DB file, logical scopes): * project — stable project_key (git remote / explicit id / cwd fallback) * session — session_id within that project * global — entire database (explicit only) * * project_key prefers git identity so the same repo keeps one key across * machines even when absolute cwd differs. Override with TOOL_WAL_PROJECT_ID * or .pi/wal-project-id when needed. * * Default /wal views are project-scoped so one project does not see another. * * Lifecycle: * tool_execution_start -> insert pending (covers calls blocked before tool_call) * tool_call -> upsert pending with final (possibly mutated) input * tool_result -> complete (success / failed / partial / blocked) * tool_execution_end -> complete if still pending (blocked / immediate errors) * session_start -> open DB, orphan stale pending from prior crashes * session_shutdown -> orphan still-pending ops for this session, close DB * * Commands: * /wal this project summary + recent * /wal status session + project (+ global) counts * /wal pending [scope] pending ops (scope: session|project|global) * /wal recent [n] last n ops in this project * /wal session current session ops * /wal project all sessions in this project * /wal projects list known projects * /wal global [n] cross-project recent * /wal show full record * /wal orphaned orphaned ops (this project) * /wal prune [days] delete old completed rows (this project; add "global") * /wal path database path */ import type { ExtensionAPI, ExtensionContext, ToolCallEvent, ToolExecutionEndEvent, ToolExecutionStartEvent, ToolResultEvent, } from "@earendil-works/pi-coding-agent"; import { resolveProjectIdentity, ToolWalStore } from "./db.ts"; import { clearProjectIdentityCache } from "./project-id.ts"; import type { OperationStatus, QueryScope, ToolOperationRow } from "./types.ts"; import { DEFAULT_ORPHAN_AGE_MS, DEFAULT_RECENT_LIMIT } from "./types.ts"; const STATUS_LABEL: Record = { pending: "pending", success: "success", failed: "failed", partial: "partial", blocked: "blocked", orphaned: "orphaned", }; function shortId(id: string, n = 8): string { return id.length <= n ? id : id.slice(0, n); } function formatTime(ms: number | null | undefined): string { if (ms == null) return "-"; return new Date(ms).toISOString().replace("T", " ").replace(/\.\d{3}Z$/, "Z"); } function summarizeInput(inputJson: string, max = 80): string { try { const parsed = JSON.parse(inputJson) as unknown; const flat = parsed && typeof parsed === "object" && !Array.isArray(parsed) ? Object.entries(parsed as Record) .map(([k, v]) => { const val = typeof v === "string" ? v : v == null ? String(v) : JSON.stringify(v); return `${k}=${val}`; }) .join(" ") : JSON.stringify(parsed); return flat.length > max ? `${flat.slice(0, max)}…` : flat; } catch { return inputJson.length > max ? `${inputJson.slice(0, max)}…` : inputJson; } } function formatRow(row: ToolOperationRow, verbose = false): string { const dur = row.duration_ms != null ? `${row.duration_ms}ms` : row.status === "pending" ? "…" : "-"; const expected = row.expected_ok == null ? "" : row.expected_ok ? " ok" : " !expected"; const head = [ STATUS_LABEL[row.status].padEnd(8), row.tool_name.padEnd(12), shortId(row.tool_call_id).padEnd(8), shortId(row.session_id ?? "-", 8).padEnd(8), dur.padStart(7), expected, ].join(" "); if (!verbose) { return `${head} ${summarizeInput(row.input_json)}`; } const lines = [ `${head}`, ` id: ${row.id}`, ` toolCallId: ${row.tool_call_id}`, ` project: ${row.project_key ?? "-"}`, ` projectSrc: ${row.project_key_source ?? "-"}`, ` session: ${row.session_id ?? "-"}`, ` sessionFile: ${row.session_file ?? "-"}`, ` cwd: ${row.cwd ?? "-"}`, ` started: ${formatTime(row.started_at)}`, ` ended: ${formatTime(row.ended_at)}`, ` input: ${summarizeInput(row.input_json, 200)}`, ]; if (row.error_message) lines.push(` error: ${row.error_message}`); if (row.mismatch_reason) lines.push(` mismatch: ${row.mismatch_reason}`); if (row.output_text) { const out = row.output_text.length > 300 ? `${row.output_text.slice(0, 300)}…` : row.output_text; lines.push(` output: ${out.replace(/\n/g, "\\n")}`); } return lines.join("\n"); } function formatStats( label: string, stats: ReturnType, ): string { return [ `${label}:`, ` total=${stats.total} pending=${stats.pending} success=${stats.success} partial=${stats.partial}`, ` failed=${stats.failed} blocked=${stats.blocked} orphaned=${stats.orphaned}`, ` expected_ok=${stats.expectedOk} expected_not_ok=${stats.expectedNotOk}`, ].join("\n"); } interface WalMeta { sessionId: string | null; sessionFile: string | null; cwd: string; projectKey: string; projectKeySource: string; projectDetail?: string; } function sessionMeta(ctx: ExtensionContext): WalMeta { const sm = ctx.sessionManager; const cwd = ctx.cwd; const identity = resolveProjectIdentity(cwd); return { sessionId: sm.getSessionId?.() ?? null, sessionFile: sm.getSessionFile?.() ?? null, cwd: identity.cwd, projectKey: identity.projectKey, projectKeySource: identity.source, projectDetail: identity.detail, }; } function scopeFilter( scope: QueryScope, meta: WalMeta, ): { sessionId?: string; projectKey?: string } { if (scope === "session") { return meta.sessionId ? { sessionId: meta.sessionId } : { projectKey: meta.projectKey }; } if (scope === "project") { return { projectKey: meta.projectKey }; } return {}; } function parseScope(token: string | undefined, fallback: QueryScope): QueryScope { if (!token) return fallback; const t = token.toLowerCase(); if (t === "session" || t === "s") return "session"; if (t === "project" || t === "p") return "project"; if (t === "global" || t === "g" || t === "all") return "global"; return fallback; } function safeNotify( ctx: ExtensionContext, message: string, level: "info" | "warning" | "error" = "info", ): void { if (!ctx.hasUI) return; try { ctx.ui.notify(message, level); } catch { // UI may be unavailable mid-shutdown } } function updateStatusLine(ctx: ExtensionContext, store: ToolWalStore): void { if (!ctx.hasUI || !store.isOpen) return; try { const meta = sessionMeta(ctx); const pending = store.pendingCount({ sessionId: meta.sessionId ?? undefined, }); if (pending > 0) { ctx.ui.setStatus("tool-wal", `wal:${pending} pending`); } else { ctx.ui.setStatus("tool-wal", undefined); } } catch { // ignore status line failures } } export default function (pi: ExtensionAPI) { const store = new ToolWalStore(); let currentTurnIndex: number | null = null; let activeMeta: WalMeta | null = null; const ensureOpen = (): boolean => { if (store.isOpen) return true; try { store.open(); return true; } catch (err) { console.error("[tool-wal] failed to open database:", err); return false; } }; const beginFromEvent = ( ctx: ExtensionContext, args: { toolCallId: string; toolName: string; input: unknown; source: "tool_execution_start" | "tool_call"; }, ) => { const meta = sessionMeta(ctx); activeMeta = meta; store.beginOperation({ toolCallId: args.toolCallId, toolName: args.toolName, input: args.input ?? {}, projectKey: meta.projectKey, projectKeySource: meta.projectKeySource, sessionId: meta.sessionId, sessionFile: meta.sessionFile, cwd: meta.cwd, turnIndex: currentTurnIndex, source: args.source, }); }; pi.on("session_start", async (event, ctx) => { try { // Re-resolve git/explicit identity after reload or cwd change. clearProjectIdentityCache(); store.open(); const meta = sessionMeta(ctx); activeMeta = meta; // Recover abandoned pending rows from previous process crashes. // Keep very recent pending on reload to avoid racing in-flight ops. const orphaned = event.reason === "reload" ? store.markOrphaned({ olderThanMs: DEFAULT_ORPHAN_AGE_MS, excludeSessionId: meta.sessionId ?? undefined, }) : store.markOrphaned({ olderThanMs: DEFAULT_ORPHAN_AGE_MS, }); if (orphaned > 0) { safeNotify( ctx, `tool-wal: marked ${orphaned} stale pending op(s) as orphaned`, "warning", ); } updateStatusLine(ctx, store); } catch (err) { console.error("[tool-wal] session_start error:", err); safeNotify( ctx, `tool-wal: init failed: ${err instanceof Error ? err.message : String(err)}`, "error", ); } }); pi.on("session_shutdown", async (_event, ctx) => { try { if (store.isOpen) { const meta = sessionMeta(ctx); // Anything still pending for this session will not get a result. store.markOrphaned({ sessionId: meta.sessionId ?? undefined, olderThanMs: 0, }); if (ctx.hasUI) { try { ctx.ui.setStatus("tool-wal", undefined); } catch { // ignore } } store.close(); } } catch (err) { console.error("[tool-wal] session_shutdown error:", err); } finally { activeMeta = null; currentTurnIndex = null; } }); pi.on("turn_start", async (event) => { currentTurnIndex = event.turnIndex ?? null; }); // Earliest hook: covers tools blocked by an earlier tool_call handler. pi.on("tool_execution_start", async (event: ToolExecutionStartEvent, ctx) => { if (!ensureOpen()) return; try { beginFromEvent(ctx, { toolCallId: event.toolCallId, toolName: event.toolName, input: event.args ?? {}, source: "tool_execution_start", }); updateStatusLine(ctx, store); } catch (err) { console.error("[tool-wal] tool_execution_start error:", err); } }); // True write-ahead intent: final args after tool_call mutations. pi.on("tool_call", async (event: ToolCallEvent, ctx) => { if (!ensureOpen()) return; try { beginFromEvent(ctx, { toolCallId: event.toolCallId, toolName: event.toolName, input: event.input ?? {}, source: "tool_call", }); updateStatusLine(ctx, store); } catch (err) { console.error("[tool-wal] tool_call error:", err); } // Never block — this extension is audit-only. }); pi.on("tool_result", async (event: ToolResultEvent, ctx) => { if (!ensureOpen()) return; try { const meta = sessionMeta(ctx); store.completeOperation({ toolCallId: event.toolCallId, toolName: event.toolName, input: event.input, isError: event.isError, content: event.content, details: event.details, source: "tool_result", sessionId: meta.sessionId, projectKey: meta.projectKey, }); updateStatusLine(ctx, store); } catch (err) { console.error("[tool-wal] tool_result error:", err); } }); // Fallback when tool_result never fires (blocked / immediate error path). pi.on("tool_execution_end", async (event: ToolExecutionEndEvent, ctx) => { if (!ensureOpen()) return; try { const meta = sessionMeta(ctx); const existing = store.getByToolCallId(event.toolCallId, { sessionId: meta.sessionId, }); if (existing && existing.status !== "pending") return; store.completeOperation({ toolCallId: event.toolCallId, toolName: event.toolName, isError: event.isError, content: event.result?.content, details: event.result?.details, source: "tool_execution_end", sessionId: meta.sessionId, projectKey: meta.projectKey, }); updateStatusLine(ctx, store); } catch (err) { console.error("[tool-wal] tool_execution_end error:", err); } }); pi.registerCommand("wal", { description: "Tool WAL: status | pending | recent | session | project | projects | global | show | orphaned | prune | path", handler: async (args, ctx) => { if (!ensureOpen()) { ctx.ui.notify("tool-wal: database unavailable", "error"); return; } const raw = (args ?? "").trim(); const [cmd, ...rest] = raw.length === 0 ? ["summary"] : raw.split(/\s+/); const meta = sessionMeta(ctx); activeMeta = meta; const printLines = (text: string) => { console.log(`[tool-wal]\n${text}`); if (ctx.hasUI) { ctx.ui.notify( text.length > 1500 ? `${text.slice(0, 1500)}\n…(see console)` : text, "info", ); } }; const print = (text: string) => { if (ctx.hasUI) ctx.ui.notify(text, "info"); else console.log(text); }; switch (cmd.toLowerCase()) { case "summary": case "s": { const sessionStats = store.stats({ sessionId: meta.sessionId ?? undefined, }); const projectStats = store.stats({ projectKey: meta.projectKey }); const recent = store.query({ projectKey: meta.projectKey, limit: 8, }); const body = [ `project: ${meta.projectKey}`, `source: ${meta.projectKeySource}${meta.projectDetail ? ` (${meta.projectDetail})` : ""}`, `cwd: ${meta.cwd}`, `session: ${meta.sessionId ?? "(none)"}`, "", formatStats("This session", sessionStats), formatStats("This project (all sessions)", projectStats), "", "Recent (this project):", recent.length === 0 ? " (none)" : recent.map((r) => ` ${formatRow(r)}`).join("\n"), "", `db: ${store.dbPath}`, 'Tip: /wal global for cross-project; /wal projects to list.', ].join("\n"); printLines(body); return; } case "status": case "stats": { const sessionStats = store.stats({ sessionId: meta.sessionId ?? undefined, }); const projectStats = store.stats({ projectKey: meta.projectKey }); const globalStats = store.stats({}); printLines( [ `project: ${meta.projectKey} [${meta.projectKeySource}]`, formatStats("This session", sessionStats), formatStats("This project", projectStats), formatStats("Global (all projects)", globalStats), ].join("\n"), ); return; } case "pending": { const scope = parseScope(rest[0], "project"); const filter = scopeFilter(scope, meta); const rows = store.query({ ...filter, status: "pending", limit: 50, }); printLines( rows.length === 0 ? `No pending operations (${scope}).` : [ `Pending (${scope}):`, ...rows.map((r) => ` ${formatRow(r)}`), ].join("\n"), ); return; } case "recent": { // /wal recent [n] [scope] let n = DEFAULT_RECENT_LIMIT; let scope: QueryScope = "project"; if (rest[0] && /^\d+$/.test(rest[0])) { n = Math.max(1, Math.min(Number(rest[0]), 200)); scope = parseScope(rest[1], "project"); } else if (rest[0]) { scope = parseScope(rest[0], "project"); if (rest[1] && /^\d+$/.test(rest[1])) { n = Math.max(1, Math.min(Number(rest[1]), 200)); } } const filter = scopeFilter(scope, meta); const rows = store.query({ ...filter, limit: n }); printLines( rows.length === 0 ? `No operations recorded (${scope}).` : [ `Last ${rows.length} operations (${scope}):`, ...rows.map((r) => ` ${formatRow(r)}`), ].join("\n"), ); return; } case "session": { const rows = store.query({ sessionId: meta.sessionId ?? undefined, limit: 100, }); printLines( [ `Session ${meta.sessionId ?? "(unknown)"} @ ${meta.projectKey}:`, rows.length === 0 ? " (no operations)" : rows.map((r) => ` ${formatRow(r)}`).join("\n"), ].join("\n"), ); return; } case "project": { const rows = store.query({ projectKey: meta.projectKey, limit: 100, }); printLines( [ `Project ${meta.projectKey}`, `cwd: ${meta.cwd}`, rows.length === 0 ? " (no operations)" : rows.map((r) => ` ${formatRow(r)}`).join("\n"), ].join("\n"), ); return; } case "projects": { const projects = store.listProjects(30); printLines( projects.length === 0 ? "No projects recorded yet." : [ "Projects in WAL:", ...projects.map((p) => { const mark = p.project_key === meta.projectKey ? "*" : " "; const src = p.project_key_source ?? "?"; return ( ` ${mark} ${p.project_key}\n` + ` src=${src} cwd=${p.cwd ?? "?"} ops=${p.total} ` + `sessions=${p.sessions} pending=${p.pending}` ); }), "", "* = current project", ].join("\n"), ); return; } case "global": case "all": { const n = Math.max( 1, Math.min(Number(rest[0]) || DEFAULT_RECENT_LIMIT, 200), ); const rows = store.query({ limit: n }); const globalStats = store.stats({}); printLines( [ formatStats("Global", globalStats), "", `Last ${rows.length} operations (all projects):`, rows.length === 0 ? " (none)" : rows .map( (r) => ` ${shortId(r.project_key ?? "?", 24).padEnd(24)} ${formatRow(r)}`, ) .join("\n"), ].join("\n"), ); return; } case "show": { const key = rest[0]; if (!key) { print("Usage: /wal show "); return; } const byId = store.getById(key); const row = byId ?? store.getByToolCallId(key, { sessionId: meta.sessionId }) ?? store.getByToolCallId(key, { projectKey: meta.projectKey }) ?? store.getByToolCallId(key); if (!row) { const candidates = store .query({ projectKey: meta.projectKey, limit: 200 }) .filter( (r) => r.id.startsWith(key) || r.tool_call_id.startsWith(key) || (r.session_id?.startsWith(key) ?? false), ); if (candidates.length === 1) { printLines(formatRow(candidates[0], true)); return; } if (candidates.length > 1) { printLines( [ `Multiple matches for "${key}" in this project:`, ...candidates.slice(0, 10).map((r) => ` ${formatRow(r)}`), ].join("\n"), ); return; } print(`No operation found for "${key}"`); return; } printLines(formatRow(row, true)); return; } case "orphaned": { const scope = parseScope(rest[0], "project"); const filter = scopeFilter(scope, meta); const rows = store.query({ ...filter, status: "orphaned", limit: 50, }); printLines( rows.length === 0 ? `No orphaned operations (${scope}).` : [ `Orphaned (${scope}):`, ...rows.map((r) => ` ${formatRow(r)}`), ].join("\n"), ); return; } case "prune": { // /wal prune [days] [global] let days = 30; let global = false; for (const tok of rest) { if (/^\d+$/.test(tok)) days = Math.max(1, Number(tok)); else if (tok.toLowerCase() === "global" || tok.toLowerCase() === "all") { global = true; } } const removed = store.prune(days * 24 * 60 * 60 * 1000, { projectKey: global ? undefined : meta.projectKey, }); print( global ? `Pruned ${removed} completed op(s) older than ${days} day(s) (global).` : `Pruned ${removed} completed op(s) older than ${days} day(s) (project ${meta.projectKey}).`, ); return; } case "path": case "db": { printLines( [ `db: ${store.dbPath}`, `project: ${meta.projectKey}`, `source: ${meta.projectKeySource}`, `detail: ${meta.projectDetail ?? "-"}`, `cwd: ${meta.cwd}`, `session: ${meta.sessionId ?? "(none)"}`, `file: ${meta.sessionFile ?? "(none)"}`, "", "Stable id overrides (first match):", " 1. env TOOL_WAL_PROJECT_ID", " 2. .pi/wal-project-id", " 3. .pi/tool-wal.json { projectId }", " 4. git remote origin (+ path from toplevel)", " 5. cwd encoding (machine-local fallback)", ].join("\n"), ); return; } case "help": case "?": { printLines( [ "Tool write-ahead log — scopes: session < project < global", " /wal this project summary + recent", " /wal status session + project + global counts", " /wal pending [scope] pending (default: project)", " /wal recent [n] [scope] recent ops (default: project)", " /wal session current session only", " /wal project all sessions in this project", " /wal projects list projects in the DB", " /wal global [n] cross-project recent", " /wal show full record", " /wal orphaned [scope] orphaned ops", " /wal prune [days] [global] delete old completed rows", " /wal path db + current project/session ids", ].join("\n"), ); return; } default: { print(`Unknown /wal subcommand "${cmd}". Try /wal help`); } } }, }); }