mirror of
https://github.com/RooVetGit/Roo-Code.git
synced 2026-10-11 03:38:15 +00:00
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+).
This commit is contained in:
parent
4c173af275
commit
b5e3a3d898
1 changed files with 210 additions and 168 deletions
|
|
@ -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<TaskMetadata | null> {
|
||||
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<boolean> {
|
||||
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<boolean> {
|
||||
// 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 <base>/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<boolean> => {
|
||||
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<boolean> => {
|
||||
// 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<number>((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<boolean> => {
|
||||
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?.(
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue