/** * SQLite persistence for the tool write-ahead log. * * Uses Node's built-in node:sqlite (DatabaseSync). No extra npm dependency. * * Storage layout: * Single DB file: ~/.pi/agent/tool-wal/wal.db * Logical isolation columns: * project_key — stable id (git remote / explicit / cwd fallback) * project_key_source — explicit | git | cwd * session_id — one conversation within that project * cwd — local path at write time (machine-specific) * * Scope hierarchy (narrow → wide): session → project → global. */ import { createHash, randomUUID } from "node:crypto"; import { mkdirSync } from "node:fs"; import { homedir } from "node:os"; import { dirname, join } from "node:path"; import { DatabaseSync } from "node:sqlite"; import { projectKeyFromCwd, resolveProjectIdentity, } from "./project-id.ts"; import type { BeginOperationInput, CompleteOperationInput, OperationStats, OperationStatus, ProjectSummary, QueryOptions, ToolOperationRow, } from "./types.ts"; import { MAX_STORED_CHARS } from "./types.ts"; export { projectKeyFromCwd, resolveProjectIdentity } from "./project-id.ts"; const SCHEMA_VERSION = 3; /** Tables only — indexes are applied after migrate() so legacy DBs can gain columns first. */ const TABLES_SQL = ` CREATE TABLE IF NOT EXISTS meta ( key TEXT PRIMARY KEY, value TEXT NOT NULL ); CREATE TABLE IF NOT EXISTS tool_operations ( id TEXT PRIMARY KEY, tool_call_id TEXT NOT NULL, project_key TEXT, project_key_source TEXT, session_id TEXT, session_file TEXT, cwd TEXT, tool_name TEXT NOT NULL, input_json TEXT NOT NULL, input_hash TEXT, status TEXT NOT NULL, started_at INTEGER NOT NULL, ended_at INTEGER, duration_ms INTEGER, is_error INTEGER, output_text TEXT, output_json TEXT, details_json TEXT, error_message TEXT, expected_ok INTEGER, mismatch_reason TEXT, turn_index INTEGER, created_at INTEGER NOT NULL, updated_at INTEGER NOT NULL ); `; const INDEXES_SQL = ` CREATE UNIQUE INDEX IF NOT EXISTS idx_tool_ops_session_call ON tool_operations(session_id, tool_call_id); CREATE INDEX IF NOT EXISTS idx_tool_ops_project ON tool_operations(project_key); CREATE INDEX IF NOT EXISTS idx_tool_ops_project_started ON tool_operations(project_key, started_at DESC); CREATE INDEX IF NOT EXISTS idx_tool_ops_status ON tool_operations(status); CREATE INDEX IF NOT EXISTS idx_tool_ops_started ON tool_operations(started_at DESC); CREATE INDEX IF NOT EXISTS idx_tool_ops_tool_name ON tool_operations(tool_name); CREATE INDEX IF NOT EXISTS idx_tool_ops_tool_call_id ON tool_operations(tool_call_id); `; export function defaultDbPath(): string { return join(homedir(), ".pi", "agent", "tool-wal", "wal.db"); } function stableStringify(value: unknown): string { const seen = new WeakSet(); const normalize = (v: unknown): unknown => { if (v === null || typeof v !== "object") return v; if (seen.has(v as object)) return "[Circular]"; seen.add(v as object); if (Array.isArray(v)) return v.map(normalize); const obj = v as Record; const out: Record = {}; for (const key of Object.keys(obj).sort()) { out[key] = normalize(obj[key]); } return out; }; try { return JSON.stringify(normalize(value)); } catch { return JSON.stringify(String(value)); } } export function hashInput(input: unknown): string { return createHash("sha256").update(stableStringify(input)).digest("hex"); } function truncate(text: string, max = MAX_STORED_CHARS): string { if (text.length <= max) return text; return `${text.slice(0, max)}\n…[truncated ${text.length - max} chars]`; } function safeJson(value: unknown, max = MAX_STORED_CHARS): string { return truncate(stableStringify(value), max); } function extractText(content: unknown): string { if (typeof content === "string") return content; if (!Array.isArray(content)) { if (content == null) return ""; return typeof content === "object" ? stableStringify(content) : String(content); } const parts: string[] = []; for (const block of content) { if (!block || typeof block !== "object") continue; const typed = block as { type?: string; text?: string }; if (typed.type === "text" && typeof typed.text === "string") { parts.push(typed.text); } } return parts.join("\n"); } function extractErrorMessage(content: unknown, details: unknown, isError: boolean): string | null { if (!isError) return null; const text = extractText(content).trim(); if (text) return truncate(text, 4_096); if (details && typeof details === "object") { const d = details as Record; for (const key of ["error", "message", "reason"]) { if (typeof d[key] === "string" && d[key]) return truncate(d[key] as string, 4_096); } } return "Tool reported an error"; } function looksBlocked(errorMessage: string | null, details: unknown): boolean { const haystacks: string[] = []; if (errorMessage) haystacks.push(errorMessage); if (details && typeof details === "object") { haystacks.push(stableStringify(details)); } const joined = haystacks.join("\n").toLowerCase(); return ( joined.includes("blocked") || joined.includes("tool execution was blocked") || joined.includes("permission denied by extension") ); } function looksPartial(details: unknown, content: unknown): boolean { if (details && typeof details === "object") { const d = details as Record; if (d.truncated === true || d.truncated === "true") return true; if (d.cancelled === true) return true; } const text = extractText(content); if (/truncated/i.test(text) && /output was truncated|truncated to/i.test(text)) { return true; } return false; } function evaluateExpected(args: { status: OperationStatus; isError: boolean; inputHashAtStart: string | null; inputAtComplete: unknown | undefined; details: unknown; content: unknown; }): { expectedOk: boolean; mismatchReason: string | null } { if (args.status === "blocked") { return { expectedOk: false, mismatchReason: "Tool call was blocked before execution" }; } if (args.status === "failed") { return { expectedOk: false, mismatchReason: "Tool execution reported an error" }; } if (args.status === "orphaned") { return { expectedOk: false, mismatchReason: "Operation never received a completion event" }; } if (args.inputAtComplete !== undefined && args.inputHashAtStart) { const completeHash = hashInput(args.inputAtComplete); if (completeHash !== args.inputHashAtStart) { return { expectedOk: false, mismatchReason: "Input at completion differs from write-ahead intent", }; } } if (args.status === "partial") { return { expectedOk: true, mismatchReason: "Completed with partial/truncated output", }; } return { expectedOk: true, mismatchReason: null }; } function emptyStats(): OperationStats { return { total: 0, pending: 0, success: 0, failed: 0, partial: 0, blocked: 0, orphaned: 0, expectedOk: 0, expectedNotOk: 0, }; } function rowToStats(base: unknown): OperationStats { const row = (base ?? {}) as Record; const n = (v: number | null | undefined) => Number(v ?? 0); return { total: n(row.total), pending: n(row.pending), success: n(row.success), failed: n(row.failed), partial: n(row.partial), blocked: n(row.blocked), orphaned: n(row.orphaned), expectedOk: n(row.expectedOk), expectedNotOk: n(row.expectedNotOk), }; } const STATS_SELECT = ` COUNT(*) AS total, SUM(CASE WHEN status = 'pending' THEN 1 ELSE 0 END) AS pending, SUM(CASE WHEN status = 'success' THEN 1 ELSE 0 END) AS success, SUM(CASE WHEN status = 'failed' THEN 1 ELSE 0 END) AS failed, SUM(CASE WHEN status = 'partial' THEN 1 ELSE 0 END) AS partial, SUM(CASE WHEN status = 'blocked' THEN 1 ELSE 0 END) AS blocked, SUM(CASE WHEN status = 'orphaned' THEN 1 ELSE 0 END) AS orphaned, SUM(CASE WHEN expected_ok = 1 THEN 1 ELSE 0 END) AS expectedOk, SUM(CASE WHEN expected_ok = 0 THEN 1 ELSE 0 END) AS expectedNotOk `; export class ToolWalStore { readonly dbPath: string; private db: DatabaseSync | null = null; constructor(dbPath: string = defaultDbPath()) { this.dbPath = dbPath; } get isOpen(): boolean { return this.db !== null; } open(): void { if (this.db) return; mkdirSync(dirname(this.dbPath), { recursive: true }); const db = new DatabaseSync(this.dbPath); // Reliability pragmas: durable WAL log for the log itself. db.exec("PRAGMA journal_mode = WAL;"); db.exec("PRAGMA synchronous = NORMAL;"); db.exec("PRAGMA busy_timeout = 5000;"); db.exec("PRAGMA foreign_keys = ON;"); db.exec(TABLES_SQL); this.migrate(db); db.exec(INDEXES_SQL); this.db = db; } private tableColumns(db: DatabaseSync): Set { const cols = db.prepare("PRAGMA table_info(tool_operations)").all() as Array<{ name: string; }>; return new Set(cols.map((c) => c.name)); } private setSchemaVersion(db: DatabaseSync, version: number): void { db.prepare( `INSERT INTO meta(key, value) VALUES ('schema_version', ?) ON CONFLICT(key) DO UPDATE SET value = excluded.value`, ).run(String(version)); } private migrate(db: DatabaseSync): void { const versionRow = db .prepare("SELECT value FROM meta WHERE key = 'schema_version'") .get() as { value?: string } | undefined; let version = versionRow?.value ? Number(versionRow.value) : 0; if (!Number.isFinite(version)) version = 0; // Fresh DB. if (!versionRow) { const names = this.tableColumns(db); if (!names.has("project_key")) { db.exec("ALTER TABLE tool_operations ADD COLUMN project_key TEXT"); } if (!names.has("project_key_source")) { db.exec("ALTER TABLE tool_operations ADD COLUMN project_key_source TEXT"); } db.prepare( "INSERT INTO meta(key, value) VALUES ('schema_version', ?), ('created_at', ?)", ).run(String(SCHEMA_VERSION), String(Date.now())); return; } if (version < 2) { const names = this.tableColumns(db); if (!names.has("project_key")) { db.exec("ALTER TABLE tool_operations ADD COLUMN project_key TEXT"); } // Backfill machine-local keys from cwd (best effort for legacy rows). const missing = db .prepare( `SELECT id, cwd FROM tool_operations WHERE project_key IS NULL AND cwd IS NOT NULL AND cwd != ''`, ) .all() as Array<{ id: string; cwd: string }>; const update = db.prepare( "UPDATE tool_operations SET project_key = ?, updated_at = ? WHERE id = ?", ); const now = Date.now(); for (const row of missing) { try { update.run(projectKeyFromCwd(row.cwd), now, row.id); } catch { // skip unresolvable cwd } } version = 2; this.setSchemaVersion(db, version); } if (version < 3) { const names = this.tableColumns(db); if (!names.has("project_key_source")) { db.exec("ALTER TABLE tool_operations ADD COLUMN project_key_source TEXT"); } // Mark legacy cwd-encoded keys; leave others null until rewritten by new writes. db.prepare( `UPDATE tool_operations SET project_key_source = 'cwd' WHERE project_key_source IS NULL AND project_key IS NOT NULL AND project_key LIKE '--%--'`, ).run(); version = 3; this.setSchemaVersion(db, version); } } close(): void { if (!this.db) return; try { this.db.close(); } catch { // ignore double-close } this.db = null; } private requireDb(): DatabaseSync { if (!this.db) { throw new Error("Tool WAL database is not open"); } return this.db; } /** * Record intent before (or at the start of) tool execution. * Idempotent per (session_id, tool_call_id). */ beginOperation(input: BeginOperationInput): ToolOperationRow { const db = this.requireDb(); const now = Date.now(); const sessionId = input.sessionId ?? null; const cwd = input.cwd ?? null; const resolved = input.projectKey != null ? { projectKey: input.projectKey, source: input.projectKeySource ?? null, } : cwd ? (() => { const id = resolveProjectIdentity(cwd); return { projectKey: id.projectKey, source: id.source }; })() : { projectKey: null, source: null }; const projectKey = resolved.projectKey; const projectKeySource = resolved.source; const inputJson = safeJson(input.input); const inputHash = hashInput(input.input); const existing = db .prepare( `SELECT * FROM tool_operations WHERE tool_call_id = ? AND ((session_id IS NULL AND ? IS NULL) OR session_id = ?) LIMIT 1`, ) .get(input.toolCallId, sessionId, sessionId) as ToolOperationRow | undefined; if (existing) { // Refresh intent from tool_call (final mutated args). Do not downgrade completed rows. if (existing.status === "pending" && input.source === "tool_call") { db.prepare( `UPDATE tool_operations SET tool_name = ?, input_json = ?, input_hash = ?, project_key = COALESCE(?, project_key), project_key_source = COALESCE(?, project_key_source), cwd = COALESCE(?, cwd), session_file = COALESCE(?, session_file), turn_index = COALESCE(?, turn_index), updated_at = ? WHERE id = ?`, ).run( input.toolName, inputJson, inputHash, projectKey, projectKeySource, cwd, input.sessionFile ?? null, input.turnIndex ?? null, now, existing.id, ); return this.getById(existing.id)!; } return existing; } const id = randomUUID(); db.prepare( `INSERT INTO tool_operations ( id, tool_call_id, project_key, project_key_source, session_id, session_file, cwd, tool_name, input_json, input_hash, status, started_at, ended_at, duration_ms, is_error, output_text, output_json, details_json, error_message, expected_ok, mismatch_reason, turn_index, created_at, updated_at ) VALUES ( ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'pending', ?, NULL, NULL, NULL, NULL, NULL, NULL, NULL, NULL, NULL, ?, ?, ? )`, ).run( id, input.toolCallId, projectKey, projectKeySource, sessionId, input.sessionFile ?? null, cwd, input.toolName, inputJson, inputHash, now, input.turnIndex ?? null, now, now, ); return this.getById(id)!; } /** * Finalize a pending operation after execution (or block / immediate error). */ completeOperation(input: CompleteOperationInput): ToolOperationRow | null { const db = this.requireDb(); const now = Date.now(); let row: ToolOperationRow | undefined; if (input.sessionId) { row = db .prepare( `SELECT * FROM tool_operations WHERE tool_call_id = ? AND session_id = ? ORDER BY started_at DESC LIMIT 1`, ) .get(input.toolCallId, input.sessionId) as ToolOperationRow | undefined; } if (!row && input.projectKey) { row = db .prepare( `SELECT * FROM tool_operations WHERE tool_call_id = ? AND project_key = ? ORDER BY started_at DESC LIMIT 1`, ) .get(input.toolCallId, input.projectKey) as ToolOperationRow | undefined; } if (!row) { row = db .prepare( `SELECT * FROM tool_operations WHERE tool_call_id = ? ORDER BY started_at DESC LIMIT 1`, ) .get(input.toolCallId) as ToolOperationRow | undefined; } if (!row) { // Completion without a prior begin — still record for audit. row = this.beginOperation({ toolCallId: input.toolCallId, toolName: input.toolName, input: input.input ?? {}, sessionId: input.sessionId, projectKey: input.projectKey, source: "tool_call", }); } // tool_result is authoritative; ignore a later tool_execution_end. if (row.status !== "pending" && input.source === "tool_execution_end") { return row; } // Do not overwrite a terminal state with another completion except tool_result upgrading. if (row.status !== "pending" && input.source !== "tool_result") { return row; } const errorMessage = input.errorMessage ?? extractErrorMessage(input.content, input.details, input.isError); let status: OperationStatus = input.status ?? (input.isError ? looksBlocked(errorMessage, input.details) ? "blocked" : "failed" : looksPartial(input.details, input.content) ? "partial" : "success"); // Non-zero bash exit is a failure even if isError is false in some paths. if ( status === "success" && input.details && typeof input.details === "object" && typeof (input.details as { exitCode?: unknown }).exitCode === "number" && (input.details as { exitCode: number }).exitCode !== 0 ) { status = "failed"; } const { expectedOk, mismatchReason } = evaluateExpected({ status, isError: input.isError || status === "failed" || status === "blocked", inputHashAtStart: row.input_hash, inputAtComplete: input.input, details: input.details, content: input.content, }); const duration = Math.max(0, now - row.started_at); const outputText = truncate(extractText(input.content)); const outputJson = input.content === undefined ? null : safeJson(input.content); const detailsJson = input.details === undefined ? null : safeJson(input.details); db.prepare( `UPDATE tool_operations SET status = ?, ended_at = ?, duration_ms = ?, is_error = ?, output_text = ?, output_json = ?, details_json = ?, error_message = ?, expected_ok = ?, mismatch_reason = ?, updated_at = ? WHERE id = ?`, ).run( status, now, duration, input.isError || status === "failed" || status === "blocked" ? 1 : 0, outputText || null, outputJson, detailsJson, errorMessage, expectedOk ? 1 : 0, mismatchReason, now, row.id, ); return this.getById(row.id); } /** * Mark stale pending rows as orphaned (crash / interrupted session recovery). */ markOrphaned(options: { olderThanMs?: number; sessionId?: string; projectKey?: string; excludeSessionId?: string; } = {}): number { const db = this.requireDb(); const now = Date.now(); const olderThan = options.olderThanMs ?? 0; const cutoff = now - olderThan; let sql = ` UPDATE tool_operations SET status = 'orphaned', ended_at = ?, duration_ms = CASE WHEN started_at IS NOT NULL THEN MAX(0, ? - started_at) ELSE NULL END, is_error = 1, expected_ok = 0, mismatch_reason = COALESCE(mismatch_reason, 'Left pending without completion'), error_message = COALESCE(error_message, 'Orphaned pending operation'), updated_at = ? WHERE status = 'pending' AND started_at <= ? `; const params: Array = [now, now, now, cutoff]; if (options.sessionId) { sql += " AND session_id = ?"; params.push(options.sessionId); } if (options.projectKey) { sql += " AND project_key = ?"; params.push(options.projectKey); } if (options.excludeSessionId) { sql += " AND (session_id IS NULL OR session_id != ?)"; params.push(options.excludeSessionId); } const result = db.prepare(sql).run(...params) as { changes: number }; return result.changes ?? 0; } getById(id: string): ToolOperationRow | null { const db = this.requireDb(); const row = db.prepare("SELECT * FROM tool_operations WHERE id = ?").get(id) as | ToolOperationRow | undefined; return row ?? null; } getByToolCallId( toolCallId: string, options: { sessionId?: string | null; projectKey?: string | null } = {}, ): ToolOperationRow | null { const db = this.requireDb(); if (options.sessionId != null) { const row = db .prepare( `SELECT * FROM tool_operations WHERE tool_call_id = ? AND session_id = ? ORDER BY started_at DESC LIMIT 1`, ) .get(toolCallId, options.sessionId) as ToolOperationRow | undefined; return row ?? null; } if (options.projectKey != null) { const row = db .prepare( `SELECT * FROM tool_operations WHERE tool_call_id = ? AND project_key = ? ORDER BY started_at DESC LIMIT 1`, ) .get(toolCallId, options.projectKey) as ToolOperationRow | undefined; return row ?? null; } const row = db .prepare( `SELECT * FROM tool_operations WHERE tool_call_id = ? ORDER BY started_at DESC LIMIT 1`, ) .get(toolCallId) as ToolOperationRow | undefined; return row ?? null; } query(options: QueryOptions = {}): ToolOperationRow[] { const db = this.requireDb(); const where: string[] = []; const params: Array = []; if (options.sessionId) { where.push("session_id = ?"); params.push(options.sessionId); } else if (options.projectKey) { where.push("project_key = ?"); params.push(options.projectKey); } if (options.toolName) { where.push("tool_name = ?"); params.push(options.toolName); } if (options.sinceMs) { where.push("started_at >= ?"); params.push(options.sinceMs); } if (options.status) { const statuses = Array.isArray(options.status) ? options.status : [options.status]; where.push(`status IN (${statuses.map(() => "?").join(", ")})`); params.push(...statuses); } const limit = Math.max(1, Math.min(options.limit ?? 50, 500)); const offset = Math.max(0, options.offset ?? 0); const sql = ` SELECT * FROM tool_operations ${where.length ? `WHERE ${where.join(" AND ")}` : ""} ORDER BY started_at DESC LIMIT ? OFFSET ? `; params.push(limit, offset); return db.prepare(sql).all(...params) as ToolOperationRow[]; } /** * Stats for a scope. Prefer sessionId > projectKey > global. */ stats(options: { sessionId?: string; projectKey?: string } | string = {}): OperationStats { const db = this.requireDb(); // Back-compat: stats(sessionId?: string) const opts = typeof options === "string" ? { sessionId: options } : (options ?? {}); if (opts.sessionId) { return rowToStats( db .prepare( `SELECT ${STATS_SELECT} FROM tool_operations WHERE session_id = ?`, ) .get(opts.sessionId), ); } if (opts.projectKey) { return rowToStats( db .prepare( `SELECT ${STATS_SELECT} FROM tool_operations WHERE project_key = ?`, ) .get(opts.projectKey), ); } return rowToStats(db.prepare(`SELECT ${STATS_SELECT} FROM tool_operations`).get()); } listProjects(limit = 50): ProjectSummary[] { const db = this.requireDb(); const rows = db .prepare( `SELECT project_key, MAX(project_key_source) AS project_key_source, MAX(cwd) AS cwd, COUNT(*) AS total, SUM(CASE WHEN status = 'pending' THEN 1 ELSE 0 END) AS pending, COUNT(DISTINCT session_id) AS sessions, MAX(started_at) AS last_started_at FROM tool_operations WHERE project_key IS NOT NULL GROUP BY project_key ORDER BY last_started_at DESC LIMIT ?`, ) .all(Math.max(1, Math.min(limit, 200))) as Array<{ project_key: string; project_key_source: string | null; cwd: string | null; total: number; pending: number; sessions: number; last_started_at: number | null; }>; return rows.map((r) => ({ project_key: r.project_key, project_key_source: r.project_key_source, cwd: r.cwd, total: Number(r.total ?? 0), pending: Number(r.pending ?? 0), sessions: Number(r.sessions ?? 0), last_started_at: r.last_started_at == null ? null : Number(r.last_started_at), })); } /** * Delete completed rows older than the given age. Never deletes pending. * Scope with projectKey when provided. */ prune(olderThanMs: number, options: { projectKey?: string } = {}): number { const db = this.requireDb(); const cutoff = Date.now() - olderThanMs; if (options.projectKey) { const result = db .prepare( `DELETE FROM tool_operations WHERE status != 'pending' AND project_key = ? AND COALESCE(ended_at, started_at) < ?`, ) .run(options.projectKey, cutoff) as { changes: number }; return result.changes ?? 0; } const result = db .prepare( `DELETE FROM tool_operations WHERE status != 'pending' AND COALESCE(ended_at, started_at) < ?`, ) .run(cutoff) as { changes: number }; return result.changes ?? 0; } /** Pending count for the optional status line (defaults to session scope). */ pendingCount(options: { sessionId?: string; projectKey?: string } | string = {}): number { const db = this.requireDb(); const opts = typeof options === "string" ? { sessionId: options } : (options ?? {}); if (opts.sessionId) { const row = db .prepare( `SELECT COUNT(*) AS c FROM tool_operations WHERE status = 'pending' AND session_id = ?`, ) .get(opts.sessionId) as { c: number }; return Number(row?.c ?? 0); } if (opts.projectKey) { const row = db .prepare( `SELECT COUNT(*) AS c FROM tool_operations WHERE status = 'pending' AND project_key = ?`, ) .get(opts.projectKey) as { c: number }; return Number(row?.c ?? 0); } const row = db .prepare(`SELECT COUNT(*) AS c FROM tool_operations WHERE status = 'pending'`) .get() as { c: number }; return Number(row?.c ?? 0); } } export { emptyStats };