From b5e3a3d898a8498304651c32e52cdfd8ef075fc6 Mon Sep 17 00:00:00 2001 From: Hannes Rudolph Date: Mon, 26 Jan 2026 15:32:36 -0700 Subject: [PATCH] perf: optimize task history retention purge with parallel processing - Add parallel metadata reads with 50 concurrent operations - Add parallel deletions with 10 concurrent operations - Batch read all metadata first, then filter, then delete - Reduce retry attempts from 3 to 1 (remove aggressive retries with sleeps) - Extract helper functions for cleaner code (readTaskMetadata, pathExists, removeDir) This significantly improves purge performance for users with many tasks (~9000+). --- src/utils/task-history-retention.ts | 378 +++++++++++++++------------- 1 file changed, 210 insertions(+), 168 deletions(-) diff --git a/src/utils/task-history-retention.ts b/src/utils/task-history-retention.ts index d59b979bd0..6d17d5a6a2 100644 --- a/src/utils/task-history-retention.ts +++ b/src/utils/task-history-retention.ts @@ -2,6 +2,7 @@ import * as vscode from "vscode" import * as path from "path" import * as fs from "fs/promises" import type { Dirent } from "fs" +import pLimit from "p-limit" import { TASK_HISTORY_RETENTION_OPTIONS, type TaskHistoryRetentionSetting } from "@roo-code/types" @@ -16,8 +17,119 @@ export type PurgeResult = { cutoff: number | null } +/** Concurrency limit for parallel metadata reads */ +const METADATA_READ_CONCURRENCY = 50 + +/** Concurrency limit for parallel task deletions */ +const DELETION_CONCURRENCY = 10 + +/** + * Task metadata read result for batch processing + */ +interface TaskMetadata { + taskId: string + taskDir: string + ts: number | null + isOrphan: boolean + mtime: number | null +} + +/** + * Read metadata for a single task directory. + * Returns null if the task directory should be skipped. + */ +async function readTaskMetadata(taskId: string, tasksDir: string): Promise { + const taskDir = path.join(tasksDir, taskId) + const metadataPath = path.join(taskDir, GlobalFileNames.taskMetadata) + + let ts: number | null = null + let isOrphan = false + let mtime: number | null = null + + // Try to read metadata file + try { + const raw = await fs.readFile(metadataPath, "utf8") + const meta: unknown = JSON.parse(raw) + const maybeTs = Number( + typeof meta === "object" && meta !== null && "ts" in meta ? (meta as { ts: unknown }).ts : undefined, + ) + if (Number.isFinite(maybeTs)) { + ts = maybeTs + } + } catch { + // Missing or invalid metadata + } + + // Check for orphan directories (checkpoint-only) - only if no valid timestamp + if (ts === null) { + try { + const childEntries = await fs.readdir(taskDir, { withFileTypes: true }) + const visibleNames = childEntries.map((e) => e.name).filter((n) => !n.startsWith(".")) + const hasCheckpointsDir = childEntries.some((e) => e.isDirectory() && e.name === "checkpoints") + const nonCheckpointVisible = visibleNames.filter((n) => n !== "checkpoints") + const hasMetadataFile = visibleNames.includes(GlobalFileNames.taskMetadata) + if (hasCheckpointsDir && nonCheckpointVisible.length === 0 && !hasMetadataFile) { + isOrphan = true + } + } catch { + // Ignore errors + } + + // Get mtime as fallback for tasks without valid ts + if (!isOrphan) { + try { + const stat = await fs.stat(taskDir) + mtime = stat.mtime.getTime() + } catch { + // Can't stat - skip this task + return null + } + } + } + + return { taskId, taskDir, ts, isOrphan, mtime } +} + +/** + * Check if path exists + */ +async function pathExists(p: string): Promise { + try { + await fs.access(p) + return true + } catch { + return false + } +} + +/** + * Simplified directory removal - one attempt with fallback + * Removed aggressive retries and sleeps for performance + */ +async function removeDir(dir: string): Promise { + // First attempt: standard recursive remove + try { + await fs.rm(dir, { recursive: true, force: true }) + } catch { + // ignore + } + + if (!(await pathExists(dir))) return true + + // Fallback: try removing checkpoints first (common stubborn directory) + try { + await fs.rm(path.join(dir, "checkpoints"), { recursive: true, force: true }) + await fs.rm(dir, { recursive: true, force: true }) + } catch { + // ignore + } + + return !(await pathExists(dir)) +} + /** * Purge old task directories under /tasks based on task_metadata.json ts value. + * Optimized for performance with parallel metadata reads and parallel deletions. * Executes best-effort deletes; errors are logged and skipped. * * @param retention Retention setting: "never" | "90" | "60" | "30" | "7" | "3" or number of days @@ -44,7 +156,9 @@ export async function purgeOldTasks( const logv = (msg: string) => { if (verbose) log?.(msg) } - logv(`[Retention] Starting purge with retention=${retention} (${days} day(s))${dryRun ? " (dry run)" : ""}`) + logv( + `[Retention] Starting optimized purge with retention=${retention} (${days} day(s))${dryRun ? " (dry run)" : ""}`, + ) let basePath: string @@ -71,190 +185,118 @@ export async function purgeOldTasks( } const taskDirs = entries.filter((d) => d.isDirectory()) + const totalTasks = taskDirs.length - logv(`[Retention] Found ${taskDirs.length} task director${taskDirs.length === 1 ? "y" : "ies"} under ${tasksDir}`) + logv(`[Retention] Found ${totalTasks} task director${totalTasks === 1 ? "y" : "ies"} under ${tasksDir}`) - // Small helpers - const pathExists = async (p: string): Promise => { - try { - await fs.access(p) - return true - } catch { - return false - } + if (totalTasks === 0) { + return { purgedCount: 0, cutoff } } - const sleep = (ms: number) => new Promise((r) => setTimeout(r, ms)) + // Phase 1: Batch read all metadata in parallel + logv(`[Retention] Phase 1: Reading metadata for ${totalTasks} tasks (concurrency: ${METADATA_READ_CONCURRENCY})`) + const metadataLimit = pLimit(METADATA_READ_CONCURRENCY) - // Aggressive recursive remove with retries; also directly clears checkpoints if needed - const removeDirAggressive = async (dir: string): Promise => { - // Try up to 3 passes with short delays - for (let attempt = 1; attempt <= 3; attempt++) { - try { - await fs.rm(dir, { recursive: true, force: true }) - } catch { - // ignore and try more targeted cleanup below - } + const metadataResults = await Promise.all( + taskDirs.map((d) => metadataLimit(() => readTaskMetadata(d.name, tasksDir))), + ) - // Verify - if (!(await pathExists(dir))) return true + // Phase 2: Filter tasks that need deletion + const tasksToDelete: Array<{ metadata: TaskMetadata; reason: string }> = [] - // Targeted cleanup for stubborn checkpoint-only directories - try { - await fs.rm(path.join(dir, "checkpoints"), { recursive: true, force: true }) - } catch { - // ignore - } - - // Remove children one by one in case some FS impls struggle with rm -r - try { - const entries = await fs.readdir(dir, { withFileTypes: true }) - for (const entry of entries) { - const entryPath = path.join(dir, entry.name) - try { - if (entry.isDirectory()) { - await fs.rm(entryPath, { recursive: true, force: true }) - } else { - await fs.unlink(entryPath) - } - } catch { - // ignore individual failures; we'll retry the parent - } - } - } catch { - // ignore - } - - // Final attempt this pass - try { - await fs.rm(dir, { recursive: true, force: true }) - } catch { - // ignore - } - - if (!(await pathExists(dir))) return true - - // Backoff a bit before next attempt - await sleep(50 * attempt) - } - - return !(await pathExists(dir)) - } - - const results: number[] = [] - - for (const d of taskDirs) { - const taskDir = path.join(tasksDir, d.name) - const metadataPath = path.join(taskDir, GlobalFileNames.taskMetadata) - - let ts: number | null = null - - // First try to get a timestamp from task_metadata.json (if present) - try { - const raw = await fs.readFile(metadataPath, "utf8") - const meta: unknown = JSON.parse(raw) - const maybeTs = Number( - typeof meta === "object" && meta !== null && "ts" in meta ? (meta as { ts: unknown }).ts : undefined, - ) - if (Number.isFinite(maybeTs)) { - ts = maybeTs - } - } catch { - // Missing or invalid metadata; we'll fall back to directory mtime. - } + for (const metadata of metadataResults) { + if (!metadata) continue let shouldDelete = false let reason = "" - // Check for checkpoint-only orphan directories (delete regardless of age) - try { - const childEntries = await fs.readdir(taskDir, { withFileTypes: true }) - const visibleNames = childEntries.map((e) => e.name).filter((n) => !n.startsWith(".")) - const hasCheckpointsDir = childEntries.some((e) => e.isDirectory() && e.name === "checkpoints") - const nonCheckpointVisible = visibleNames.filter((n) => n !== "checkpoints") - const hasMetadataFile = visibleNames.includes(GlobalFileNames.taskMetadata) - if (hasCheckpointsDir && nonCheckpointVisible.length === 0 && !hasMetadataFile) { - shouldDelete = true - reason = "orphan checkpoints_only" - } - } catch { - // Ignore errors while scanning children; proceed with normal logic - } - - if (!shouldDelete && ts !== null && ts < cutoff) { - // Normal case: metadata has a valid ts older than cutoff + // Check orphan directories (delete regardless of age) + if (metadata.isOrphan) { shouldDelete = true - reason = `ts=${ts}` - } else if (!shouldDelete) { - // Orphan/legacy case: no valid ts; fall back to directory mtime - try { - const stat = await fs.stat(taskDir) - const mtimeMs = stat.mtime.getTime() - if (mtimeMs < cutoff) { - shouldDelete = true - reason = `no valid ts, mtime=${stat.mtime.toISOString()}` - } - } catch { - // If we can't stat the directory, skip it. - } + reason = "orphan checkpoints_only" + } + // Check by timestamp + else if (metadata.ts !== null && metadata.ts < cutoff) { + shouldDelete = true + reason = `ts=${metadata.ts}` + } + // Check by mtime fallback + else if (metadata.ts === null && metadata.mtime !== null && metadata.mtime < cutoff) { + shouldDelete = true + reason = `no valid ts, mtime=${new Date(metadata.mtime).toISOString()}` } - if (!shouldDelete) { - results.push(0) - continue + if (shouldDelete) { + tasksToDelete.push({ metadata, reason }) } - - if (dryRun) { - logv(`[Retention][DRY RUN] Would delete task ${d.name} (${reason}) @ ${taskDir}`) - results.push(1) - continue - } - - // Attempt deletion using provider callback (for full cleanup) or direct rm - let deletionError: unknown | null = null - let deleted = false - try { - if (deleteTaskById) { - logv(`[Retention] Deleting task ${d.name} via provider @ ${taskDir} (${reason})`) - await deleteTaskById(d.name, taskDir) - // Provider callback handles full cleanup; check if directory is gone - deleted = !(await pathExists(taskDir)) - } else { - logv(`[Retention] Deleting task ${d.name} via fs.rm @ ${taskDir} (${reason})`) - await fs.rm(taskDir, { recursive: true, force: true }) - deleted = !(await pathExists(taskDir)) - } - } catch (e) { - deletionError = e - } - - // If directory still exists after initial attempt, try aggressive cleanup with retries - if (!deleted) { - deleted = await removeDirAggressive(taskDir) - } - - if (!deleted) { - // Did not actually remove; report the most relevant error - if (deletionError) { - log?.( - `[Retention] Failed to delete task ${d.name} @ ${taskDir}: ${ - deletionError instanceof Error ? deletionError.message : String(deletionError) - } (directory still present)`, - ) - } else { - log?.( - `[Retention] Failed to delete task ${d.name} @ ${taskDir}: directory still present after cleanup attempts`, - ) - } - results.push(0) - continue - } - - logv(`[Retention] Deleted task ${d.name} (${reason}) @ ${taskDir}`) - results.push(1) } - const purged = results.reduce((sum, n) => sum + n, 0) + logv(`[Retention] Phase 2: ${tasksToDelete.length} of ${totalTasks} tasks marked for deletion`) + + if (tasksToDelete.length === 0) { + log?.(`[Retention] No tasks met purge criteria${dryRun ? " (dry run)" : ""}`) + return { purgedCount: 0, cutoff } + } + + // Phase 3: Delete tasks in parallel + if (dryRun) { + for (const { metadata, reason } of tasksToDelete) { + logv(`[Retention][DRY RUN] Would delete task ${metadata.taskId} (${reason}) @ ${metadata.taskDir}`) + } + log?.( + `[Retention] Would purge ${tasksToDelete.length} task(s) (dry run); cutoff=${new Date(cutoff).toISOString()}`, + ) + return { purgedCount: tasksToDelete.length, cutoff } + } + + logv(`[Retention] Phase 3: Deleting ${tasksToDelete.length} tasks (concurrency: ${DELETION_CONCURRENCY})`) + const deleteLimit = pLimit(DELETION_CONCURRENCY) + + const deleteResults = await Promise.all( + tasksToDelete.map(({ metadata, reason }) => + deleteLimit(async (): Promise => { + let deleted = false + + try { + if (deleteTaskById) { + logv( + `[Retention] Deleting task ${metadata.taskId} via provider @ ${metadata.taskDir} (${reason})`, + ) + await deleteTaskById(metadata.taskId, metadata.taskDir) + deleted = !(await pathExists(metadata.taskDir)) + } else { + logv(`[Retention] Deleting task ${metadata.taskId} via fs.rm @ ${metadata.taskDir} (${reason})`) + await fs.rm(metadata.taskDir, { recursive: true, force: true }) + deleted = !(await pathExists(metadata.taskDir)) + } + } catch (e) { + // Primary deletion failed, try fallback + logv( + `[Retention] Primary deletion failed for ${metadata.taskId}: ${ + e instanceof Error ? e.message : String(e) + }`, + ) + } + + // Fallback: simplified removal + if (!deleted) { + deleted = await removeDir(metadata.taskDir) + } + + if (!deleted) { + log?.( + `[Retention] Failed to delete task ${metadata.taskId} @ ${metadata.taskDir}: directory still present`, + ) + } else { + logv(`[Retention] Deleted task ${metadata.taskId} (${reason}) @ ${metadata.taskDir}`) + } + + return deleted + }), + ), + ) + + const purged = deleteResults.filter(Boolean).length if (purged > 0) { log?.(