diff --git a/packages/coding-agent/src/cli/gc-cli.ts b/packages/coding-agent/src/cli/gc-cli.ts index 7b2e0f3ce..5e3e2aed9 100644 --- a/packages/coding-agent/src/cli/gc-cli.ts +++ b/packages/coding-agent/src/cli/gc-cli.ts @@ -2,6 +2,7 @@ import { Database } from "bun:sqlite"; import * as fs from "node:fs/promises"; import * as path from "node:path"; import { gunzipSync, gzipSync } from "node:zlib"; +import { withStatsSyncLock } from "@oh-my-pi/omp-stats/aggregator"; import { getAgentDir, getBlobsDir, @@ -429,18 +430,6 @@ function sessionArtifactsPath(sessionPath: string): string { return sessionPath.slice(0, -SESSION_SUFFIX.length); } -function sessionIdFromSessionPath(sessionPath: string): string | undefined { - const basename = path.basename(sessionPath); - if (basename.endsWith(COMPRESSED_SESSION_SUFFIX)) { - const id = basename.slice(0, -COMPRESSED_SESSION_SUFFIX.length); - return id || undefined; - } - if (basename.endsWith(SESSION_SUFFIX)) { - const id = basename.slice(0, -SESSION_SUFFIX.length); - return id || undefined; - } - return undefined; -} interface SessionLineageHeader { id: string; parentSession?: string; @@ -469,10 +458,6 @@ function sessionLineageHeaderFromText(text: string): SessionLineageHeader | unde return undefined; } -async function archivedSessionIdFromFile(file: string): Promise { - return sessionLineageHeaderFromText(await readTextIfPresent(file))?.id ?? sessionIdFromSessionPath(file); -} - async function gzipSessionFile(source: string, destination: string): Promise { await fs.mkdir(path.dirname(destination), { recursive: true }); const tempPath = `${destination}.${process.pid}.${Date.now()}.tmp`; @@ -578,7 +563,7 @@ function deleteHistoryRowsForSessions(dbPath: string, sessionIds: string[]): { d async function collectArchivedSessionIds(archiveRoot: string): Promise { const ids = new Set(); for (const file of await collectCompressedJsonlFiles(archiveRoot)) { - const id = await archivedSessionIdFromFile(file); + const id = sessionLineageHeaderFromText(await readTextIfPresent(file))?.id; if (id) ids.add(id); } return [...ids].sort(); @@ -611,6 +596,13 @@ async function cleanupHistoryRowsForArchivedSessions( const STATS_SESSION_TABLES = ["messages", "user_messages", "tool_calls", "file_offsets"] as const; const STATS_ENTRY_TABLES = ["messages", "user_messages", "tool_calls"] as const; +type StatsEntryTable = (typeof STATS_ENTRY_TABLES)[number]; + +const STATS_IDENTITY_COLUMNS: Record = { + messages: ["entry_id", "timestamp"], + user_messages: ["entry_id", "timestamp"], + tool_calls: ["entry_id", "timestamp", "tool_call_id"], +}; interface ArchivedStatsSession { path: string; @@ -622,9 +614,15 @@ interface StatsLineageNode extends ArchivedStatsSession { statsPaths: Set; } +interface StatsEntryIdentity { + entryId: string; + timestamp: number; + toolCallId: string; +} + interface StatsTransferTarget { path: string; - entryIds: string[]; + identities: Record; } interface StatsCleanupPlan { @@ -692,7 +690,7 @@ function resolveLogicalSessionMatch( const stem = matches[0] ? path.basename(matches[0].logicalRoot, SESSION_SUFFIX) : ""; const idMatches = matches.filter(match => stem === match.node.id || stem.endsWith(`_${match.node.id}`)); if (idMatches.length === 1) return idMatches[0]; - return matches.length === 1 ? matches[0] : undefined; + return undefined; } /** @@ -787,31 +785,84 @@ function buildStatsCleanupPlans( }); } -async function collectSessionMessageEntryIds(sessionPath: string): Promise { +async function collectSessionStatsIdentities( + sessionPath: string, +): Promise> { + const identities: Record = { + messages: [], + user_messages: [], + tool_calls: [], + }; const decoder = new TextDecoder(); - const entryIds: string[] = []; for await (const line of readLines(Bun.file(sessionPath).stream())) { if (line.length === 0) continue; try { - const entry = JSON.parse(decoder.decode(line)) as { type?: unknown; id?: unknown }; - if (entry.type === "message" && typeof entry.id === "string" && entry.id.length > 0) entryIds.push(entry.id); + const record: unknown = JSON.parse(decoder.decode(line)); + if ( + !record || + typeof record !== "object" || + !("type" in record) || + record.type !== "message" || + !("id" in record) || + typeof record.id !== "string" || + record.id.length === 0 || + !("message" in record) || + !record.message || + typeof record.message !== "object" + ) { + continue; + } + const message = record.message; + if (!("role" in message)) continue; + const parsedEntryTimestamp = + "timestamp" in record && typeof record.timestamp === "string" ? Date.parse(record.timestamp) : Number.NaN; + if (message.role === "user") { + identities.user_messages.push({ + entryId: record.id, + timestamp: Number.isFinite(parsedEntryTimestamp) ? parsedEntryTimestamp : 0, + toolCallId: "", + }); + continue; + } + if (message.role !== "assistant") continue; + const timestamp = + "timestamp" in message && typeof message.timestamp === "number" && Number.isFinite(message.timestamp) + ? message.timestamp + : Number.isFinite(parsedEntryTimestamp) + ? parsedEntryTimestamp + : 0; + identities.messages.push({ entryId: record.id, timestamp, toolCallId: "" }); + if (!("content" in message) || !Array.isArray(message.content)) continue; + for (const block of message.content) { + if ( + block && + typeof block === "object" && + "type" in block && + block.type === "toolCall" && + "id" in block && + typeof block.id === "string" && + block.id.length > 0 + ) { + identities.tool_calls.push({ entryId: record.id, timestamp, toolCallId: block.id }); + } + } } catch { // Stats parsing is also lenient: a malformed line cannot own a retained row. } } - return entryIds; + return identities; } async function populateStatsTransferTargets(plans: StatsCleanupPlan[]): Promise { - const entryIdsBySession = new Map>(); + const identitiesBySession = new Map>>(); for (const plan of plans) { for (const retained of plan.retainedSessions) { - let entryIds = entryIdsBySession.get(retained.path); - if (!entryIds) { - entryIds = collectSessionMessageEntryIds(retained.path); - entryIdsBySession.set(retained.path, entryIds); + let identities = identitiesBySession.get(retained.path); + if (!identities) { + identities = collectSessionStatsIdentities(retained.path); + identitiesBySession.set(retained.path, identities); } - plan.transfers.push({ path: retained.path, entryIds: await entryIds }); + plan.transfers.push({ path: retained.path, identities: await identities }); } } } @@ -830,50 +881,83 @@ function reconcileStatsRowsForSessions(dbPath: string, plans: StatsCleanupPlan[] table => tableExists(db, table) && tableHasColumn(db, table, "session_file"), ); const entryTables = STATS_ENTRY_TABLES.filter( - table => sessionTables.includes(table) && tableHasColumn(db, table, "entry_id"), + table => + sessionTables.includes(table) && + STATS_IDENTITY_COLUMNS[table].every(column => tableHasColumn(db, table, column)), ); - const missingEntryIdentity = STATS_ENTRY_TABLES.filter( - table => sessionTables.includes(table) && !tableHasColumn(db, table, "entry_id"), - ); - if ( - missingEntryIdentity.length > 0 && - plans.some(plan => plan.transfers.some(target => target.entryIds.length > 0)) - ) { - throw new Error(`stats tables missing entry_id: ${missingEntryIdentity.join(", ")}`); - } - db.run("CREATE TEMP TABLE gc_retained_entries (entry_id TEXT PRIMARY KEY, target_session_file TEXT NOT NULL)"); + + db.run(` + CREATE TEMP TABLE gc_retained_entries ( + table_name TEXT NOT NULL, + entry_id TEXT NOT NULL, + timestamp INTEGER NOT NULL, + tool_call_id TEXT NOT NULL, + target_session_file TEXT NOT NULL, + PRIMARY KEY (table_name, entry_id, timestamp, tool_call_id) + ) + `); const clearRetainedEntries = db.prepare("DELETE FROM gc_retained_entries"); - const insertRetainedEntry = db.prepare( - "INSERT OR IGNORE INTO gc_retained_entries (entry_id, target_session_file) VALUES (?, ?)", - ); - const transferStatements = entryTables.map(table => - db.prepare(` + const insertRetainedEntry = db.prepare(` + INSERT OR IGNORE INTO gc_retained_entries ( + table_name, entry_id, timestamp, tool_call_id, target_session_file + ) VALUES (?, ?, ?, ?, ?) + `); + const transferStatements = entryTables.map(table => { + const toolCallMatch = + table === "tool_calls" ? `retained.tool_call_id = ${table}.tool_call_id` : "retained.tool_call_id = ''"; + const identityMatch = ` + retained.table_name = '${table}' + AND retained.entry_id = ${table}.entry_id + AND retained.timestamp = ${table}.timestamp + AND ${toolCallMatch} + `; + return db.prepare(` UPDATE OR IGNORE ${table} SET session_file = ( - SELECT target_session_file - FROM gc_retained_entries - WHERE gc_retained_entries.entry_id = ${table}.entry_id + SELECT retained.target_session_file + FROM gc_retained_entries AS retained + WHERE ${identityMatch} ) WHERE session_file = ? - AND entry_id IN (SELECT entry_id FROM gc_retained_entries) - `), - ); - const deleteStatements = sessionTables.map(table => - db.prepare(`DELETE FROM ${table} WHERE session_file = ? OR instr(session_file, ?) = 1`), - ); + AND EXISTS ( + SELECT 1 + FROM gc_retained_entries AS retained + WHERE ${identityMatch} + ) + `); + }); + const deletionStatements = sessionTables.map(table => ({ + table, + statement: db.prepare(`DELETE FROM ${table} WHERE session_file = ? OR instr(session_file, ?) = 1`), + })); let deleted = 0; const tx = db.transaction((cleanupPlans: StatsCleanupPlan[]) => { for (const plan of cleanupPlans) { clearRetainedEntries.run(); for (const target of plan.transfers) { - for (const entryId of target.entryIds) insertRetainedEntry.run(entryId, target.path); + for (const table of entryTables) { + for (const identity of target.identities[table]) { + insertRetainedEntry.run( + table, + identity.entryId, + identity.timestamp, + identity.toolCallId, + target.path, + ); + } + } } for (const sessionPath of new Set(plan.sessionPaths)) { for (const statement of transferStatements) statement.run(sessionPath); } for (const sessionPath of new Set(plan.sessionPaths)) { const nestedPrefix = `${sessionArtifactsPath(sessionPath)}${path.sep}`; - for (const statement of deleteStatements) { + for (const { table, statement } of deletionStatements) { + const requiresTransfer = + plan.retainedSessions.length > 0 && + table !== "file_offsets" && + !entryTables.some(entryTable => entryTable === table); + if (requiresTransfer) continue; const result = statement.run(sessionPath, nestedPrefix) as SqliteRunResult; deleted += sqliteNumber(result.changes); } @@ -890,18 +974,24 @@ function reconcileStatsRowsForSessions(dbPath: string, plans: StatsCleanupPlan[] async function collectArchivedStatsSessions( archiveRoot: string, sessionsRoot: string, + onError: (file: string, error: unknown) => void, ): Promise { const sessions: ArchivedStatsSession[] = []; for (const file of await collectCompressedJsonlFiles(archiveRoot)) { const relative = path.relative(archiveRoot, file); if (!relative || relative.startsWith("..") || path.isAbsolute(relative)) continue; const sourcePath = path.join(sessionsRoot, relative.slice(0, -".gz".length)); - const header = sessionLineageHeaderFromText(await readTextIfPresent(file)); - sessions.push({ - path: sourcePath, - id: header?.id ?? sessionIdFromSessionPath(sourcePath) ?? sourcePath, - parentSession: header?.parentSession, - }); + try { + const header = sessionLineageHeaderFromText(await readTextIfPresent(file)); + if (!header) throw new Error("archive is missing a valid session header"); + sessions.push({ + path: sourcePath, + id: header.id, + parentSession: header.parentSession, + }); + } catch (error) { + onError(file, error); + } } return sessions; } @@ -929,9 +1019,19 @@ async function cleanupStatsRowsForArchivedSessions( let retainedSessions: SessionInfo[]; try { - for (const session of await collectArchivedStatsSessions(archiveRoot, getSessionsDir(options.agentDir))) { + for (const session of await collectArchivedStatsSessions( + archiveRoot, + getSessionsDir(options.agentDir), + (file, error) => { + result.errors.push(`stats cleanup scan ${file}: ${errorMessage(error)}`); + }, + )) { archivedByPath.set(path.resolve(session.path), session); } + } catch (error) { + result.errors.push(`stats cleanup scan: ${errorMessage(error)}`); + } + try { retainedSessions = await listActiveSessions(getSessionsDir(options.agentDir)); } catch (error) { result.errors.push(`stats cleanup scan: ${errorMessage(error)}`); @@ -939,10 +1039,12 @@ async function cleanupStatsRowsForArchivedSessions( } try { - const storedSessionPaths = collectStoredStatsSessionPaths(dbPath); - const plans = buildStatsCleanupPlans([...archivedByPath.values()], retainedSessions, storedSessionPaths); - await populateStatsTransferTargets(plans); - result.statsRowsDeleted = reconcileStatsRowsForSessions(dbPath, plans); + await withStatsSyncLock(dbPath, async () => { + const storedSessionPaths = collectStoredStatsSessionPaths(dbPath); + const plans = buildStatsCleanupPlans([...archivedByPath.values()], retainedSessions, storedSessionPaths); + await populateStatsTransferTargets(plans); + result.statsRowsDeleted = reconcileStatsRowsForSessions(dbPath, plans); + }); } catch (error) { result.errors.push(`stats cleanup: ${errorMessage(error)}`); } diff --git a/packages/coding-agent/test/gc-cli.test.ts b/packages/coding-agent/test/gc-cli.test.ts index 969fcabb3..b23dacc91 100644 --- a/packages/coding-agent/test/gc-cli.test.ts +++ b/packages/coding-agent/test/gc-cli.test.ts @@ -3,8 +3,9 @@ import { afterEach, beforeEach, describe, expect, spyOn, test } from "bun:test"; import * as fs from "node:fs/promises"; import * as os from "node:os"; import * as path from "node:path"; -import { gunzipSync } from "node:zlib"; -import { runGcCommand } from "@oh-my-pi/pi-coding-agent/cli/gc-cli"; +import { gunzipSync, gzipSync } from "node:zlib"; +import { withStatsSyncLock } from "@oh-my-pi/omp-stats/aggregator"; +import { type GcResult, runGcCommand } from "@oh-my-pi/pi-coding-agent/cli/gc-cli"; import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings"; import { getAgentDir, @@ -632,6 +633,8 @@ describe("runGcCommand cold-session archive", () => { const child = path.join(sessionsDir, "20260726_child-session.jsonl"); const sibling = path.join(sessionsDir, "20260725_sibling-session.jsonl"); const timestamp = "2026-06-26T12:00:00.000Z"; + const timestampMs = Date.parse(timestamp); + const collisionTimestamp = "2026-06-27T12:00:00.000Z"; const sharedUser = { type: "message", id: "shared-user", @@ -644,7 +647,12 @@ describe("runGcCommand cold-session archive", () => { id: "shared-assistant", parentId: "shared-user", timestamp, - message: { role: "assistant", content: [] }, + message: { + role: "assistant", + model: "test-model", + provider: "test-provider", + content: [{ type: "toolCall", id: "shared-tool", name: "read" }], + }, }; const parentOnlyUser = { type: "message", @@ -660,6 +668,41 @@ describe("runGcCommand cold-session archive", () => { timestamp, message: { role: "assistant", content: [] }, }; + const archivedCollisionUser = { + type: "message", + id: "collision-user", + parentId: "parent-only-assistant", + timestamp, + message: { role: "user", content: "archived collision" }, + }; + const archivedCollisionAssistant = { + type: "message", + id: "collision-assistant", + parentId: "collision-user", + timestamp, + message: { + role: "assistant", + model: "test-model", + provider: "test-provider", + content: [{ type: "toolCall", id: "collision-tool", name: "read" }], + }, + }; + const retainedCollisionUser = { + ...archivedCollisionUser, + timestamp: collisionTimestamp, + message: { role: "user", content: "different retained entry" }, + }; + const retainedCollisionAssistant = { + ...archivedCollisionAssistant, + timestamp: collisionTimestamp, + }; + const terminalAssistant = { + type: "message", + id: "terminal-assistant", + parentId: null, + timestamp, + message: { role: "assistant", content: [] }, + }; await Bun.write( parent, [ @@ -674,6 +717,9 @@ describe("runGcCommand cold-session archive", () => { sharedAssistant, parentOnlyUser, parentOnlyAssistant, + archivedCollisionUser, + archivedCollisionAssistant, + terminalAssistant, "", ] .map(entry => (typeof entry === "string" ? entry : JSON.stringify(entry))) @@ -693,6 +739,9 @@ describe("runGcCommand cold-session archive", () => { }), JSON.stringify(sharedUser), JSON.stringify(sharedAssistant), + JSON.stringify(retainedCollisionUser), + JSON.stringify(retainedCollisionAssistant), + JSON.stringify(terminalAssistant), "", ].join("\n"), ); @@ -709,6 +758,7 @@ describe("runGcCommand cold-session archive", () => { }), JSON.stringify(sharedUser), JSON.stringify(sharedAssistant), + JSON.stringify(terminalAssistant), "", ].join("\n"), ); @@ -719,13 +769,13 @@ describe("runGcCommand cold-session archive", () => { const statsDbPath = path.join(root, "stats.db"); const db = new Database(statsDbPath); db.run( - "CREATE TABLE messages (session_file TEXT NOT NULL, entry_id TEXT NOT NULL, UNIQUE(session_file, entry_id))", + "CREATE TABLE messages (session_file TEXT NOT NULL, entry_id TEXT NOT NULL, timestamp INTEGER NOT NULL, UNIQUE(session_file, entry_id))", ); db.run( - "CREATE TABLE user_messages (session_file TEXT NOT NULL, entry_id TEXT NOT NULL, UNIQUE(session_file, entry_id))", + "CREATE TABLE user_messages (session_file TEXT NOT NULL, entry_id TEXT NOT NULL, timestamp INTEGER NOT NULL, UNIQUE(session_file, entry_id))", ); db.run( - "CREATE TABLE tool_calls (session_file TEXT NOT NULL, entry_id TEXT NOT NULL, tool_call_id TEXT NOT NULL, UNIQUE(session_file, tool_call_id))", + "CREATE TABLE tool_calls (session_file TEXT NOT NULL, entry_id TEXT NOT NULL, timestamp INTEGER NOT NULL, tool_call_id TEXT NOT NULL, UNIQUE(session_file, tool_call_id))", ); db.run( "CREATE TABLE file_offsets (session_file TEXT PRIMARY KEY, offset INTEGER NOT NULL, last_modified INTEGER NOT NULL)", @@ -734,15 +784,17 @@ describe("runGcCommand cold-session archive", () => { ["messages", "shared-assistant", "parent-only-assistant"], ["user_messages", "shared-user", "parent-only-user"], ] as const) { - const insert = db.prepare(`INSERT INTO ${table} (session_file, entry_id) VALUES (?, ?)`); - insert.run(parent, sharedId); - insert.run(parent, parentOnlyId); + const insert = db.prepare(`INSERT INTO ${table} (session_file, entry_id, timestamp) VALUES (?, ?, ?)`); + insert.run(parent, sharedId, timestampMs); + insert.run(parent, parentOnlyId, timestampMs); + insert.run(parent, table === "messages" ? "collision-assistant" : "collision-user", timestampMs); } const insertToolCall = db.prepare( - "INSERT INTO tool_calls (session_file, entry_id, tool_call_id) VALUES (?, ?, ?)", + "INSERT INTO tool_calls (session_file, entry_id, timestamp, tool_call_id) VALUES (?, ?, ?, ?)", ); - insertToolCall.run(parent, "shared-assistant", "shared-tool"); - insertToolCall.run(parent, "parent-only-assistant", "parent-only-tool"); + insertToolCall.run(parent, "shared-assistant", timestampMs, "shared-tool"); + insertToolCall.run(parent, "parent-only-assistant", timestampMs, "parent-only-tool"); + insertToolCall.run(parent, "collision-assistant", timestampMs, "collision-tool"); db.prepare("INSERT INTO file_offsets (session_file, offset, last_modified) VALUES (?, ?, ?)").run(parent, 444, 1); const insertOffset = db.prepare( "INSERT INTO file_offsets (session_file, offset, last_modified) VALUES (?, ?, ?)", @@ -772,7 +824,7 @@ describe("runGcCommand cold-session archive", () => { check.close(); expect(result.archive?.archived).toBe(1); - expect(result.archive?.statsRowsDeleted).toBe(4); + expect(result.archive?.statsRowsDeleted).toBe(7); expect(result.archive?.errors).toEqual([]); expect(messages).toEqual([{ session_file: child, entry_id: "shared-assistant" }]); expect(userMessages).toEqual([{ session_file: child, entry_id: "shared-user" }]); @@ -813,6 +865,91 @@ describe("runGcCommand cold-session archive", () => { ]); }); + test("scopes incompatible retained-entry deletion decisions to each cleanup plan", async () => { + const sessionsDir = path.join(getSessionsDir(root), "project"); + await fs.mkdir(sessionsDir, { recursive: true }); + const parent = path.join(sessionsDir, "20260626_partial-parent.jsonl"); + const child = path.join(sessionsDir, "20260726_partial-child.jsonl"); + const timestamp = "2026-06-26T12:00:00.000Z"; + const sharedAssistant = { + type: "message", + id: "shared-assistant", + parentId: null, + timestamp, + message: { role: "assistant", content: [] }, + }; + await Bun.write( + parent, + [ + JSON.stringify({ type: "session", version: 3, id: "partial-parent", timestamp, cwd: "/tmp" }), + JSON.stringify(sharedAssistant), + "", + ].join("\n"), + ); + await agePath(parent, 90); + await Bun.write( + child, + [ + JSON.stringify({ + type: "session", + version: 3, + id: "partial-child", + timestamp, + cwd: "/tmp", + parentSession: parent, + }), + JSON.stringify(sharedAssistant), + "", + ].join("\n"), + ); + const unrelated = await writeSession(root, "project", "partial-unrelated", "complete", { + ageDays: 90, + filename: "20260625_partial-unrelated", + }); + + const statsDbPath = path.join(root, "stats.db"); + const db = new Database(statsDbPath); + db.run( + "CREATE TABLE messages (session_file TEXT NOT NULL, entry_id TEXT NOT NULL, timestamp INTEGER NOT NULL, UNIQUE(session_file, entry_id))", + ); + db.run("CREATE TABLE user_messages (session_file TEXT NOT NULL)"); + db.run( + "CREATE TABLE file_offsets (session_file TEXT PRIMARY KEY, offset INTEGER NOT NULL, last_modified INTEGER NOT NULL)", + ); + db.prepare("INSERT INTO messages (session_file, entry_id, timestamp) VALUES (?, ?, ?)").run( + parent, + "shared-assistant", + Date.parse(timestamp), + ); + db.prepare("INSERT INTO user_messages (session_file) VALUES (?)").run(parent); + db.prepare("INSERT INTO user_messages (session_file) VALUES (?)").run(unrelated); + db.prepare("INSERT INTO file_offsets (session_file, offset, last_modified) VALUES (?, ?, ?)").run(parent, 10, 1); + db.close(); + + const result = await runGcCommand({ + flags: { + agentDir: root, + archive: true, + coldArchiveAfterDays: 30, + retainNewestGlobal: 0, + retainNewestPerCwd: 0, + apply: true, + }, + }); + const check = new Database(statsDbPath); + const messages = check.prepare("SELECT session_file, entry_id FROM messages").all(); + const legacyRows = check.prepare("SELECT session_file FROM user_messages").all(); + const offsets = check.prepare("SELECT session_file FROM file_offsets").all(); + check.close(); + + expect(result.archive?.archived).toBe(2); + expect(result.archive?.statsRowsDeleted).toBe(2); + expect(result.archive?.errors).toEqual([]); + expect(messages).toEqual([{ session_file: child, entry_id: "shared-assistant" }]); + expect(legacyRows).toEqual([{ session_file: parent }]); + expect(offsets).toEqual([]); + }); + test("prunes stats still owned by a session's paths from before it moved", async () => { const original = await writeSession(root, "before-move", "moved-session", "complete", { filename: "20260626_moved-session", @@ -863,6 +1000,163 @@ describe("runGcCommand cold-session archive", () => { expect(remaining).toEqual(Object.fromEntries(tables.map(table => [table, []]))); }); + test("does not claim an unrelated historical path by basename alone", async () => { + const session = await writeSession(root, "project", "actual-session-id", "complete", { + ageDays: 90, + filename: "custom-name", + }); + const unrelated = path.join(root, "unrelated", path.basename(session)); + const statsDbPath = path.join(root, "stats.db"); + const db = new Database(statsDbPath); + db.run("CREATE TABLE messages (session_file TEXT NOT NULL)"); + const insert = db.prepare("INSERT INTO messages (session_file) VALUES (?)"); + insert.run(session); + insert.run(unrelated); + db.close(); + + const result = await runGcCommand({ + flags: { + agentDir: root, + archive: true, + coldArchiveAfterDays: 30, + retainNewestGlobal: 0, + retainNewestPerCwd: 0, + apply: true, + }, + }); + const check = new Database(statsDbPath); + const rows = check.prepare("SELECT session_file FROM messages").all(); + check.close(); + + expect(result.archive?.archived).toBe(1); + expect(result.archive?.statsRowsDeleted).toBe(1); + expect(result.archive?.errors).toEqual([]); + expect(rows).toEqual([{ session_file: unrelated }]); + }); + + test("continues stats cleanup past a corrupt historical archive", async () => { + const session = await writeSession(root, "project", "archive-me", "complete", { ageDays: 90 }); + const corruptArchive = path.join(root, "archive", "sessions", "older", "corrupt.jsonl.gz"); + await fs.mkdir(path.dirname(corruptArchive), { recursive: true }); + await Bun.write(corruptArchive, "not gzip"); + const statsDbPath = path.join(root, "stats.db"); + const db = new Database(statsDbPath); + db.run("CREATE TABLE messages (session_file TEXT NOT NULL)"); + db.prepare("INSERT INTO messages (session_file) VALUES (?)").run(session); + db.close(); + + const result = await runGcCommand({ + flags: { + agentDir: root, + archive: true, + coldArchiveAfterDays: 30, + retainNewestGlobal: 0, + retainNewestPerCwd: 0, + apply: true, + }, + }); + const check = new Database(statsDbPath); + const rows = check.prepare("SELECT session_file FROM messages").all(); + check.close(); + + expect(result.archive?.archived).toBe(1); + expect(result.archive?.statsRowsDeleted).toBe(1); + expect(result.archive?.errors.some(error => error.startsWith(`stats cleanup scan ${corruptArchive}: `))).toBe( + true, + ); + expect(rows).toEqual([]); + }); + + test("skips decompressible historical archives without a valid session header", async () => { + const archive = path.join(root, "archive", "sessions", "older", "headerless-session.jsonl.gz"); + await fs.mkdir(path.dirname(archive), { recursive: true }); + await Bun.write( + archive, + gzipSync(`${JSON.stringify({ type: "message", message: { role: "assistant", content: [] } })}\n`), + ); + const historicalStatsPath = path.join(root, "unrelated", "headerless-session.jsonl"); + const statsDbPath = path.join(root, "stats.db"); + const stats = new Database(statsDbPath); + stats.run("CREATE TABLE messages (session_file TEXT NOT NULL)"); + stats.prepare("INSERT INTO messages (session_file) VALUES (?)").run(historicalStatsPath); + stats.close(); + const historyDbPath = getHistoryDbPath(root); + const history = new Database(historyDbPath); + history.run("CREATE TABLE history (id INTEGER PRIMARY KEY, prompt TEXT NOT NULL, session_id TEXT)"); + history.run("INSERT INTO history (prompt, session_id) VALUES ('keep me', 'headerless-session')"); + history.close(); + + const result = await runGcCommand({ + flags: { + agentDir: root, + archive: true, + coldArchiveAfterDays: 30, + retainNewestGlobal: 0, + retainNewestPerCwd: 0, + apply: true, + }, + }); + const statsCheck = new Database(statsDbPath); + const statsRows = statsCheck.prepare("SELECT session_file FROM messages").all(); + statsCheck.close(); + const historyCheck = new Database(historyDbPath); + const historyRows = historyCheck.prepare("SELECT session_id FROM history").all(); + historyCheck.close(); + + expect(result.archive?.statsRowsDeleted).toBe(0); + expect(result.archive?.historyRowsDeleted).toBe(0); + expect(result.archive?.errors).toContain( + `stats cleanup scan ${archive}: archive is missing a valid session header`, + ); + expect(statsRows).toEqual([{ session_file: historicalStatsPath }]); + expect(historyRows).toEqual([{ session_id: "headerless-session" }]); + }); + + test("waits for the shared stats lock before archive reconciliation", async () => { + const session = await writeSession(root, "project", "archive-me", "complete", { ageDays: 90 }); + const statsDbPath = path.join(root, "stats.db"); + const db = new Database(statsDbPath); + db.run("CREATE TABLE messages (session_file TEXT NOT NULL)"); + db.prepare("INSERT INTO messages (session_file) VALUES (?)").run(session); + db.close(); + const sessionMoved = (async () => { + for await (const event of fs.watch(path.dirname(session))) { + if (event.filename === path.basename(session)) return true; + } + return false; + })(); + + let gcPromise: Promise | undefined; + let archivedWhileLocked = false; + let rowsWhileLocked: unknown[] = []; + await withStatsSyncLock(statsDbPath, async () => { + gcPromise = runGcCommand({ + flags: { + agentDir: root, + archive: true, + coldArchiveAfterDays: 30, + retainNewestGlobal: 0, + retainNewestPerCwd: 0, + apply: true, + }, + }); + archivedWhileLocked = await sessionMoved; + const lockedCheck = new Database(statsDbPath); + rowsWhileLocked = lockedCheck.prepare("SELECT session_file FROM messages").all(); + lockedCheck.close(); + }); + if (!gcPromise) throw new Error("GC did not start"); + const result = await gcPromise; + const check = new Database(statsDbPath); + const rows = check.prepare("SELECT session_file FROM messages").all(); + check.close(); + + expect(archivedWhileLocked).toBe(true); + expect(rowsWhileLocked).toEqual([{ session_file: session }]); + expect(result.archive?.statsRowsDeleted).toBe(1); + expect(rows).toEqual([]); + }); + test("reports stats cleanup failures and retries rows for already archived sessions", async () => { const session = await writeSession(root, "project", "archive-me", "complete", { ageDays: 90, diff --git a/packages/stats/src/aggregator.ts b/packages/stats/src/aggregator.ts index 2c6628db6..c252e8467 100644 --- a/packages/stats/src/aggregator.ts +++ b/packages/stats/src/aggregator.ts @@ -1,5 +1,6 @@ import * as fs from "node:fs"; -import { workerHostEntry } from "@oh-my-pi/pi-utils"; +import * as path from "node:path"; +import { getStatsDbPath, workerHostEntry } from "@oh-my-pi/pi-utils"; import { getRecentErrors as dbGetRecentErrors, getRecentRequests as dbGetRecentRequests, @@ -48,6 +49,209 @@ import type { } from "./types"; import { computeUsageWindowStats, fetchUsageSnapshots } from "./usage-windows"; +const STATS_SYNC_LOCK_RETRY_MS = 25; +const STATS_SYNC_LOCK_STALE_MS = 60 * 60 * 1000; + +interface StatsLockSnapshot { + dev: number; + ino: number; + size: number; + mtimeMs: number; + text: string; +} + +function errorCode(error: unknown): string | undefined { + if (!error || typeof error !== "object" || !("code" in error)) return undefined; + return typeof error.code === "string" ? error.code : undefined; +} + +function sameStatsLock(left: StatsLockSnapshot, right: StatsLockSnapshot): boolean { + return ( + left.dev === right.dev && + left.ino === right.ino && + left.size === right.size && + left.mtimeMs === right.mtimeMs && + left.text === right.text + ); +} + +async function readStatsLockSnapshot(lockPath: string): Promise { + try { + const before = await fs.promises.stat(lockPath); + const text = await fs.promises.readFile(lockPath, "utf8"); + const after = await fs.promises.stat(lockPath); + const first = { + dev: before.dev, + ino: before.ino, + size: before.size, + mtimeMs: before.mtimeMs, + text, + }; + const second = { + dev: after.dev, + ino: after.ino, + size: after.size, + mtimeMs: after.mtimeMs, + text, + }; + return sameStatsLock(first, second) ? second : null; + } catch (error) { + if (errorCode(error) === "ENOENT") return null; + throw error; + } +} + +function statsLockOwnerIsRunning(snapshot: StatsLockSnapshot): boolean { + const pid = Number.parseInt(snapshot.text.split(/\r?\n/, 1)[0] ?? "", 10); + if (!Number.isSafeInteger(pid) || pid <= 0) return Date.now() - snapshot.mtimeMs < STATS_SYNC_LOCK_STALE_MS; + try { + process.kill(pid, 0); + return true; + } catch (error) { + return errorCode(error) !== "ESRCH"; + } +} + +function createStatsLockToken(): string { + return `${process.pid}\n${Date.now()}\n${Math.random()}\n`; +} + +async function removeAbandonedStatsLockFile(lockPath: string): Promise { + const snapshot = await readStatsLockSnapshot(lockPath); + if (!snapshot || statsLockOwnerIsRunning(snapshot)) return false; + + const current = await readStatsLockSnapshot(lockPath); + if (!current || !sameStatsLock(snapshot, current) || statsLockOwnerIsRunning(current)) return false; + try { + await fs.promises.unlink(lockPath); + return true; + } catch (error) { + if (errorCode(error) === "ENOENT") return false; + throw error; + } +} + +async function removeOwnedStatsLockFile(lockPath: string, owned: StatsLockSnapshot): Promise { + const current = await readStatsLockSnapshot(lockPath); + if (!current || !sameStatsLock(owned, current)) return; + try { + await fs.promises.unlink(lockPath); + } catch (error) { + if (errorCode(error) !== "ENOENT") throw error; + } +} + +async function removeAbandonedStatsLock(lockPath: string): Promise { + const snapshot = await readStatsLockSnapshot(lockPath); + if (!snapshot || statsLockOwnerIsRunning(snapshot)) return false; + + const breakerPath = `${lockPath}.break`; + let breaker: fs.promises.FileHandle; + try { + breaker = await fs.promises.open(breakerPath, "wx"); + } catch (error) { + if (errorCode(error) === "EEXIST") { + await removeAbandonedStatsLockFile(breakerPath); + return false; + } + throw error; + } + + const breakerToken = createStatsLockToken(); + let breakerSnapshot: StatsLockSnapshot | null = null; + let breakerTokenWritten = false; + let removed = false; + let cleanupError: unknown; + try { + await breaker.writeFile(breakerToken); + breakerTokenWritten = true; + breakerSnapshot = await readStatsLockSnapshot(breakerPath); + if (!breakerSnapshot || breakerSnapshot.text !== breakerToken) { + throw new Error(`Stats lock breaker changed while acquiring ${breakerPath}`); + } + + const current = await readStatsLockSnapshot(lockPath); + if (current && sameStatsLock(snapshot, current) && !statsLockOwnerIsRunning(current)) { + try { + await fs.promises.unlink(lockPath); + removed = true; + } catch (error) { + if (errorCode(error) !== "ENOENT") throw error; + } + } + } finally { + try { + await breaker.close(); + } catch (error) { + cleanupError = error; + } + try { + if (!breakerSnapshot && breakerTokenWritten) { + const current = await readStatsLockSnapshot(breakerPath); + if (current?.text === breakerToken) breakerSnapshot = current; + } + if (breakerSnapshot) await removeOwnedStatsLockFile(breakerPath, breakerSnapshot); + } catch (error) { + cleanupError ??= error; + } + } + if (cleanupError) throw cleanupError; + return removed; +} + +/** + * Serialize stats ingestion and archive reconciliation across processes. + * The lock covers file discovery, parsing, and the final SQLite write so a + * parse result for a session moved by GC can never commit after cleanup. + */ +export async function withStatsSyncLock(dbPath: string, fn: () => Promise): Promise { + const lockPath = `${dbPath}.sync.lock`; + await fs.promises.mkdir(path.dirname(dbPath), { recursive: true }); + let handle: fs.promises.FileHandle | undefined; + while (!handle) { + try { + handle = await fs.promises.open(lockPath, "wx"); + } catch (error) { + if (errorCode(error) !== "EEXIST") throw error; + await removeAbandonedStatsLock(lockPath); + const { promise, resolve } = Promise.withResolvers(); + setTimeout(resolve, STATS_SYNC_LOCK_RETRY_MS); + await promise; + } + } + + const token = createStatsLockToken(); + let result: T | undefined; + let operationError: unknown; + let operationFailed = false; + let tokenWritten = false; + try { + await handle.writeFile(token); + tokenWritten = true; + result = await fn(); + } catch (error) { + operationFailed = true; + operationError = error; + } + + let cleanupError: unknown; + try { + await handle.close(); + } catch (error) { + cleanupError = error; + } + try { + const current = await fs.promises.readFile(lockPath, "utf8"); + if (!tokenWritten || current === token) await fs.promises.unlink(lockPath); + } catch (error) { + if (errorCode(error) !== "ENOENT") cleanupError ??= error; + } + + if (operationFailed) throw operationError; + if (cleanupError) throw cleanupError; + return result as T; +} + /** * Apply a freshly parsed result to the database. Runs entirely on the * main thread so the single SQLite handle owns every write. @@ -218,6 +422,10 @@ export async function smokeTestSyncWorker({ timeoutMs = 5_000 }: { timeoutMs?: n * bar walks at a steady rate). */ export async function syncAllSessions(opts?: SyncOptions): Promise<{ processed: number; files: number }> { + return withStatsSyncLock(getStatsDbPath(), () => syncAllSessionsLocked(opts)); +} + +async function syncAllSessionsLocked(opts?: SyncOptions): Promise<{ processed: number; files: number }> { await initDb(); const files = await listAllSessionFiles(); diff --git a/packages/stats/test/sync-serial.test.ts b/packages/stats/test/sync-serial.test.ts index f759587e9..f954cda3c 100644 --- a/packages/stats/test/sync-serial.test.ts +++ b/packages/stats/test/sync-serial.test.ts @@ -1,7 +1,7 @@ import { afterEach, describe, expect, it, vi } from "bun:test"; import * as fs from "node:fs/promises"; import * as path from "node:path"; -import { syncAllSessions } from "@oh-my-pi/omp-stats/aggregator"; +import { syncAllSessions, withStatsSyncLock } from "@oh-my-pi/omp-stats/aggregator"; import { getOverallStats } from "@oh-my-pi/omp-stats/db"; import { getSessionsDir } from "@oh-my-pi/pi-utils"; import { installStatsTestIsolation } from "./helpers/temp-agent"; @@ -93,4 +93,78 @@ describe("stats sync serial mode", () => { await expect(syncAllSessions({ workers: 2 })).rejects.toBe(workerProbe); expect(workerSpy).toHaveBeenCalled(); }); + + it("reclaims a dead owner's abandoned stale breaker", async () => { + const dbPath = path.join(getSessionsDir(), "stats-lock.db"); + const lockPath = `${dbPath}.sync.lock`; + const breakerPath = `${lockPath}.break`; + const deadPid = 424_242; + await fs.mkdir(path.dirname(dbPath), { recursive: true }); + await Bun.write(lockPath, `${deadPid}\n1\nprimary\n`); + await Bun.write(breakerPath, `${deadPid}\n1\nbreaker\n`); + const staleTime = new Date(Date.now() - 2 * 60 * 60 * 1000); + await fs.utimes(lockPath, staleTime, staleTime); + await fs.utimes(breakerPath, staleTime, staleTime); + vi.spyOn(process, "kill").mockImplementation(pid => { + if (pid === deadPid) throw Object.assign(new Error("dead owner"), { code: "ESRCH" }); + return true; + }); + + vi.spyOn(globalThis, "setTimeout").mockImplementation(callback => { + queueMicrotask(callback as () => void); + return 0; + }); + const result = await withStatsSyncLock(dbPath, async () => "acquired"); + + expect(result).toBe("acquired"); + expect(await Bun.file(lockPath).exists()).toBe(false); + expect(await Bun.file(breakerPath).exists()).toBe(false); + }); + + it("never reclaims a live owner-stamped breaker even when its mtime is stale", async () => { + const dbPath = path.join(getSessionsDir(), "stats-live-breaker.db"); + const lockPath = `${dbPath}.sync.lock`; + const breakerPath = `${lockPath}.break`; + const deadPid = 424_243; + const liveBreakerToken = `${process.pid}\n1\nlive-breaker\n`; + await fs.mkdir(path.dirname(dbPath), { recursive: true }); + await Bun.write(lockPath, `${deadPid}\n1\nprimary\n`); + await Bun.write(breakerPath, liveBreakerToken); + const staleTime = new Date(Date.now() - 2 * 60 * 60 * 1000); + await fs.utimes(lockPath, staleTime, staleTime); + await fs.utimes(breakerPath, staleTime, staleTime); + vi.spyOn(process, "kill").mockImplementation(pid => { + if (pid === deadPid) throw Object.assign(new Error("dead owner"), { code: "ESRCH" }); + return true; + }); + const retryScheduled = Promise.withResolvers(); + let resumeRetry: (() => void) | undefined; + let breakerReleased = false; + vi.spyOn(globalThis, "setTimeout").mockImplementation(callback => { + const run = callback as () => void; + if (breakerReleased) { + queueMicrotask(run); + } else { + resumeRetry = run; + retryScheduled.resolve(); + } + return 0; + }); + + let acquired = false; + const pending = withStatsSyncLock(dbPath, async () => { + acquired = true; + }); + try { + await retryScheduled.promise; + expect(acquired).toBe(false); + expect(await Bun.file(breakerPath).text()).toBe(liveBreakerToken); + } finally { + breakerReleased = true; + await fs.rm(breakerPath, { force: true }); + resumeRetry?.(); + await pending; + } + expect(acquired).toBe(true); + }); });