diff --git a/packages/coding-agent/test/interactive-mode-todo-clear.test.ts b/packages/coding-agent/test/interactive-mode-todo-clear.test.ts index 14a569138..407b7d2fd 100644 --- a/packages/coding-agent/test/interactive-mode-todo-clear.test.ts +++ b/packages/coding-agent/test/interactive-mode-todo-clear.test.ts @@ -165,6 +165,7 @@ describe("InteractiveMode todo HUD persistence", () => { mode.setTodos(session.getTodoPhases()); await mode.init(); + vi.useFakeTimers(); eventBus.emit(TASK_SUBAGENT_LIFECYCLE_CHANNEL, { id: "ReviewFixer", index: 0, @@ -173,6 +174,7 @@ describe("InteractiveMode todo HUD persistence", () => { status: "completed", detached: true, }); + vi.advanceTimersByTime(100); const task = session.getTodoPhases()[0]?.tasks[0]; expect(task?.status).toBe("completed"); diff --git a/python/robomp/.env.example b/python/robomp/.env.example index 164c2715b..9087d7ef7 100644 --- a/python/robomp/.env.example +++ b/python/robomp/.env.example @@ -114,6 +114,11 @@ ROBOMP_MAX_CONCURRENCY=8 ROBOMP_TASK_TIMEOUT_SECONDS=2400 ROBOMP_TASK_TIMEOUT_HARD_GRACE_SECONDS=60 ROBOMP_REQUEST_TIMEOUT_SECONDS=120 +# Strip node_modules + the per-workspace bun cache after every task run (and +# once at boot). Deps are reinstalled at each run start, so this only costs a +# re-download per run — leaving it on keeps disk usage bounded by concurrency +# instead of by the number of open issues. +ROBOMP_RECLAIM_WORKSPACE_CACHES=true # ============================================================================= diff --git a/python/robomp/src/config.py b/python/robomp/src/config.py index 2679fd5d3..a2ff7d545 100644 --- a/python/robomp/src/config.py +++ b/python/robomp/src/config.py @@ -154,6 +154,15 @@ class Settings(BaseSettings): natives_cache_max_bytes: int = Field(4 * 1024**3, alias="ROBOMP_NATIVES_CACHE_MAX_BYTES") natives_cache_gc_interval_seconds: float = Field(3600.0, alias="ROBOMP_NATIVES_CACHE_GC_INTERVAL_SECONDS") + # Post-run workspace cache reclamation. Every task run reinstalls + # node_modules (`ensure_workspace_dependencies`), so between runs the + # checkout's node_modules and the workspace-private bun install cache are + # dead weight — multiple GB per issue that would otherwise persist until + # the issue closes and exhaust the disk. When enabled, the worker strips + # them after every event and WorkerPool.start() sweeps all workspaces once + # at boot. Costs a dependency re-download on the next run for that issue. + reclaim_workspace_caches: bool = Field(True, alias="ROBOMP_RECLAIM_WORKSPACE_CACHES") + @field_validator("bot_login", mode="after") @classmethod def _require_bot_login(cls, value: str) -> str: diff --git a/python/robomp/src/queue.py b/python/robomp/src/queue.py index 1174631c5..d64e07d7f 100644 --- a/python/robomp/src/queue.py +++ b/python/robomp/src/queue.py @@ -91,6 +91,12 @@ class WorkerPool: async def start(self) -> None: await self._reap_all_slots() + if self.settings.reclaim_workspace_caches: + # Crash leftovers: strip dependency caches from every workspace + # before the dispatcher can touch any of them again. + swept = await asyncio.to_thread(self.sandbox.reclaim_all_caches) + if swept: + log.info("workspace cache sweep", extra={"workspaces": swept}) recovered = self.db.reset_stuck_running() if recovered: log.info("recovered stuck events", extra={"count": recovered}) @@ -329,9 +335,32 @@ class WorkerPool: _reap_slot(slot_uid) finally: self._slot_pool.release(slot_uid) + await self._reclaim_event_caches(row) await self._release(row) clear_current_event(token) + async def _reclaim_event_caches(self, row: EventRow) -> None: + """Drop the workspace's dependency caches now that its event is over. + + Runs before `_release` so the issue key is still in `_inflight`: + nothing can re-enter `ensure_workspace` for this issue mid-reclaim. + Skipped during shutdown (the row resumes right after restart and the + startup sweep covers it). Best-effort: a reclaim failure never fails + the event. + """ + if not self.settings.reclaim_workspace_caches or self._shutting_down: + return + repo, sep, number = (row.issue_key or "").rpartition("#") + if not sep or not repo or not number.isdigit(): + return + try: + reclaimed = await asyncio.to_thread(self.sandbox.reclaim_workspace_caches, repo=repo, number=int(number)) + except OSError as exc: + log.warning("workspace cache reclaim failed", extra={"key": row.issue_key, "err": str(exc)}) + return + if reclaimed: + log.info("workspace caches reclaimed", extra={"key": row.issue_key}) + async def _dispatch_and_mark(self, row: EventRow, *, slot_uid: int | None = None) -> None: await self._dispatch(row, slot_uid=slot_uid) if row.delivery_id in self._cancelled: diff --git a/python/robomp/src/sandbox.py b/python/robomp/src/sandbox.py index bf7a8b565..7aae0eef6 100644 --- a/python/robomp/src/sandbox.py +++ b/python/robomp/src/sandbox.py @@ -48,6 +48,7 @@ import signal import stat import subprocess import threading +from collections.abc import Iterable from dataclasses import dataclass from pathlib import Path from typing import Any, Protocol @@ -668,6 +669,83 @@ def _chown_workspace(ws_root: Path, slot_uid: int | None) -> None: ) +# ---------- workspace cache reclamation ---------- + + +_TRASH_PREFIX = ".trash-" +# ws/repo/node_modules is depth 2 from ws_root's repo dir; nested workspace +# installs (packages/x/node_modules, python/x/web/node_modules) sit at 3-4. +_NODE_MODULES_SCAN_DEPTH = 4 + + +def _find_node_modules(repo_dir: Path, *, max_depth: int = _NODE_MODULES_SCAN_DEPTH) -> list[Path]: + """Locate `node_modules` dirs in a checkout without descending into them. + + Depth-limited (bun's hoisted linker keeps everything at the root; nested + workspace installs sit a couple of levels down) and prunes `.git` plus the + matches themselves, so the walk stays cheap on a large tree. + """ + found: list[Path] = [] + if not repo_dir.is_dir(): + return found + base_depth = len(repo_dir.parts) + for current, dirnames, _files in os.walk(repo_dir): + current_path = Path(current) + if "node_modules" in dirnames: + found.append(current_path / "node_modules") + if len(current_path.parts) - base_depth + 1 >= max_depth: + dirnames[:] = [] + else: + dirnames[:] = [d for d in dirnames if d not in (".git", "node_modules")] + return found + + +def _stage_workspace_trash(ws_root: Path) -> tuple[Path, ...]: + """Rename reclaimable cache dirs into `.trash-*` staging dirs. + + Rename is atomic and near-free on the same filesystem, so a caller can + hold a lock across this and defer the slow ``rmtree`` to after the lock + is dropped. Targets: every ``node_modules`` in the checkout, the + workspace-private XDG cache (bun install cache lives there) and the + tmpdir — all re-created by ``ensure_workspace`` + + ``ensure_workspace_dependencies`` on the next run. Session transcripts, + context, artifacts and the git worktree are never touched. + + Returns every staged trash dir, including leftovers from a previous + interrupted reclaim. + """ + if not ws_root.is_dir(): + return () + staged = [p for p in ws_root.iterdir() if p.name.startswith(_TRASH_PREFIX)] + candidates = [ + ws_root / ".omp-xdg" / "cache", + ws_root / ".omp-tmp", + *_find_node_modules(ws_root / "repo"), + ] + trash_root: Path | None = None + for index, victim in enumerate(candidates): + try: + st = victim.lstat() + except FileNotFoundError: + continue + if not (stat.S_ISDIR(st.st_mode) or stat.S_ISLNK(st.st_mode)): + continue + if trash_root is None: + trash_root = ws_root / f"{_TRASH_PREFIX}{secrets.token_hex(4)}" + trash_root.mkdir(mode=0o700) + staged.append(trash_root) + try: + victim.rename(trash_root / f"{index}-{victim.name}") + except OSError as exc: + log.warning("cache reclaim rename failed", extra={"path": str(victim), "err": str(exc)}) + return tuple(staged) + + +def _purge_trash(staged: Iterable[Path]) -> None: + for path in staged: + shutil.rmtree(path, ignore_errors=True) + + # ---------- SandboxManager ---------- @@ -1023,6 +1101,46 @@ class SandboxManager: if ws_root.exists(): shutil.rmtree(ws_root, ignore_errors=True) + def reclaim_workspace_caches(self, *, repo: str, number: int) -> bool: + """Strip re-creatable dependency caches from an idle workspace. + + Every task run reinstalls ``node_modules`` (see + ``host_tools.ensure_workspace_dependencies``), so between runs the + checkout's ``node_modules`` and the workspace-private bun install + cache are dead weight — multiple GB per issue that would otherwise + persist until the issue closes, which is exactly how the host runs + out of disk. ``--continue`` resumes are unaffected: session + transcripts, context, artifacts and the worktree survive. + + The rename pass runs under the per-repo lock (serialized against + ``ensure_workspace``); the slow deletes happen after the lock is + dropped. Returns True when anything was reclaimed. + """ + with self._repo_lock(repo): + staged = _stage_workspace_trash(self.workspace_root(repo, number)) + _purge_trash(staged) + return bool(staged) + + def reclaim_all_caches(self) -> int: + """Sweep dependency caches from every workspace under ``root``. + + Crash-leftover recovery: called from ``WorkerPool.start()`` before + the dispatch loop comes online, so no task can be touching a + workspace and the per-repo locks are deliberately skipped. Returns + the number of workspaces that had something to reclaim. + """ + if not self.root.is_dir(): + return 0 + count = 0 + for entry in sorted(self.root.iterdir()): + if entry.name == "_pool" or entry.name.startswith(".") or not entry.is_dir(): + continue + staged = _stage_workspace_trash(entry) + if staged: + _purge_trash(staged) + count += 1 + return count + __all__ = [ "GitCommandError",