From 7f544fc6691123fdab327a27df0df6ded6206d38 Mon Sep 17 00:00:00 2001 From: can1357 Date: Sat, 16 May 2026 20:29:05 +0200 Subject: [PATCH] feat(natives_cache): added content-addressed cache for pi-natives build artifacts - Added NativesCache with hardlink-based populate and atomic capture under a shared root:omp setgid directory. - Integrated cache populate into ensure_workspace and capture into post-task success path. - Added periodic GC loop in WorkerPool with configurable interval and per-repo entry/byte caps. - Provisioned /data/cache/pi-natives in entrypoint.sh and exposed five ROBOMP_NATIVES_CACHE_* settings. --- entrypoint.sh | 2 +- src/robomp/config.py | 12 + src/robomp/natives_cache.py | 485 ++++++++++++++++++++++++++++++++++ src/robomp/queue.py | 32 +++ src/robomp/sandbox.py | 102 ++++++- src/robomp/server.py | 15 +- src/robomp/tasks.py | 6 + src/robomp/worker.py | 73 ++++- tests/conftest.py | 7 + tests/test_natives_cache.py | 414 +++++++++++++++++++++++++++++ tests/test_permissions_e2e.py | 147 +++++++++++ tests/test_queue_cancel.py | 2 + tests/test_queue_shutdown.py | 2 + tests/test_sandbox.py | 100 +++++++ tests/test_server.py | 2 + tests/test_worker.py | 112 ++++++++ 16 files changed, 1501 insertions(+), 12 deletions(-) create mode 100644 src/robomp/natives_cache.py create mode 100644 tests/test_natives_cache.py diff --git a/entrypoint.sh b/entrypoint.sh index 468ae2d06..a6aa67e9e 100755 --- a/entrypoint.sh +++ b/entrypoint.sh @@ -49,7 +49,7 @@ mkdir -p /data/workspaces /data/workspaces/_pool /data/logs # so every per-issue worktree shares one cargo target/toolchain. Bun install # cache is workspace-private; a shared cache is unsafe across slot users # because bun may chmod/chown its cache root to the first writer. -mkdir -p /data/cache/cargo /data/cache/cargo-target /data/cache/rustup +mkdir -p /data/cache/cargo /data/cache/cargo-target /data/cache/rustup /data/cache/pi-natives chown -R root:omp /data/cache /data/workspaces/_pool find /data/cache /data/workspaces/_pool -type d -exec chmod 2770 {} + find /data/cache /data/workspaces/_pool -type f -perm /111 -exec chmod 0770 {} + diff --git a/src/robomp/config.py b/src/robomp/config.py index b7121e508..f9855af18 100644 --- a/src/robomp/config.py +++ b/src/robomp/config.py @@ -124,6 +124,18 @@ class Settings(BaseSettings): question_autoclose_hours: float = Field(4.0, alias="ROBOMP_QUESTION_AUTOCLOSE_HOURS") question_autoclose_scan_seconds: float = Field(60.0, alias="ROBOMP_QUESTION_AUTOCLOSE_SCAN_SECONDS") + # pi-natives build-output cache. Hardlinks pre-built + # `packages/natives/native/*.node` (and its companions) into new + # workspaces keyed by the git tree-hashes of inputs that determine the + # build output. Misses are captured automatically when a task that + # finishes successfully has fresh artifacts. Disable to fall back to + # per-workspace builds. + natives_cache_enabled: bool = Field(True, alias="ROBOMP_NATIVES_CACHE_ENABLED") + natives_cache_root: Path = Field(Path("/data/cache/pi-natives"), alias="ROBOMP_NATIVES_CACHE_ROOT") + natives_cache_max_entries_per_repo: int = Field(8, alias="ROBOMP_NATIVES_CACHE_MAX_ENTRIES_PER_REPO") + 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") + @field_validator("bot_login", mode="after") @classmethod def _require_bot_login(cls, value: str) -> str: diff --git a/src/robomp/natives_cache.py b/src/robomp/natives_cache.py new file mode 100644 index 000000000..f322475a5 --- /dev/null +++ b/src/robomp/natives_cache.py @@ -0,0 +1,485 @@ +"""Content-addressed cache of pre-built ``packages/natives/native/`` artifacts. + +The napi-rs build of ``pi_natives.-[-variant].node`` takes +minutes. Most issues never touch ``crates/``, so the same artifact is +buildable in every workspace whose source state matches one we've already +built. This module: + +1. Computes a deterministic key from the git tree-hashes of the inputs that + determine the build output, plus the target triple. +2. On workspace populate: hardlinks cached files into the worktree's + ``packages/natives/native/`` (a noop on cache miss). +3. On successful task exit: captures the workspace's freshly-built artifacts + into the cache under its (possibly new) key. + +Hardlink semantics give COW for free: every tool in the napi build path +replaces files via write-temp + rename, so a workspace rebuilding the addon +allocates a new inode and leaves the cached file untouched. Cache GC is by +LRU on ``manifest.json.captured_at``; hardlinked workspaces keep the inode +alive after the cache directory is rmtree'd. + +Ownership: cache root is provisioned ``root:omp 02770`` by ``entrypoint.sh`` +so slot subprocesses (group ``omp``) can capture under setgid inheritance. +Same shape as ``/data/cache/cargo``. +""" + +from __future__ import annotations + +import errno +import fcntl +import hashlib +import json +import logging +import os +import platform +import shutil +import subprocess +import sys +import time +from collections.abc import Generator +from contextlib import contextmanager +from dataclasses import dataclass +from pathlib import Path +from typing import IO + +log = logging.getLogger(__name__) + + +# Paths whose git tree-hash feeds the cache key. Order is significant — the +# hash incorporates the (path, tree_hash) pairs in this exact order so a +# different ordering would produce a different key. Cover every input the +# napi build reads: all workspace crates (pi-natives transitively depends on +# pi-ast/pi-iso/pi-shell), the workspace Cargo manifest + lock, the rust +# toolchain pin, and the natives package itself (build script + scripts/* + +# package.json with napi config). +CACHE_KEY_PATHS: tuple[str, ...] = ( + "crates", + "Cargo.lock", + "Cargo.toml", + "rust-toolchain.toml", + "packages/natives", +) + +# Files in ``packages/natives/native/`` that ARE pure functions of the +# cache-key inputs and travel as a unit. ``.node`` is matched by glob since +# the basename embeds the target triple + variant. +_CACHED_NODE_GLOB = "pi_natives.*.node" +_CACHED_COMPANION_FILES: tuple[str, ...] = ( + "index.d.ts", + "index.js", + "embedded-addon.js", +) +_MANIFEST_FILENAME = "manifest.json" +_LOCKFILE_NAME = ".lock" + +_NULL_TREE_HASH = "0" * 40 # placeholder for paths missing from HEAD + + +def _normalize_platform() -> str: + """Mirror node's ``process.platform`` so the cache key matches + ``build-native.ts``'s filename convention.""" + s = sys.platform + if s.startswith("linux"): + return "linux" + if s == "darwin": + return "darwin" + if s in ("win32", "cygwin"): + return "win32" + return s + + +def _normalize_arch() -> str: + """Mirror node's ``process.arch``.""" + m = platform.machine().lower() + if m in ("x86_64", "amd64"): + return "x64" + if m in ("aarch64", "arm64"): + return "arm64" + return m + + +def target_triple() -> str: + """``-[-]`` matching the napi addon basename. + + ``TARGET_VARIANT`` is honored only on x64 (the build script enforces the + same restriction). On x64 hosts that leave the variant unset we encode + ``host`` to keep the key stable across workspaces on the same machine + without trying to autodetect AVX2 from Python. + """ + plat = _normalize_platform() + arch = _normalize_arch() + if arch != "x64": + return f"{plat}-{arch}" + variant = os.environ.get("TARGET_VARIANT", "").strip() or "host" + return f"{plat}-{arch}-{variant}" + + +def _git_safe_directory_env(repo_dir: Path) -> dict[str, str]: + """Env overlay that whitelists ``repo_dir`` for git's safe.directory check. + + The orchestrator runs as root but workspaces are owned by the slot UID + (see ``SandboxManager._chown_workspace``). Without this whitelist, every + git invocation from the orchestrator on a slot-owned repo aborts with + "fatal: detected dubious ownership". Mirrors + ``robomp.sandbox._safe_directory_env`` but kept local to avoid a circular + import (sandbox imports this module). + """ + env = os.environ.copy() + count = int(env.get("GIT_CONFIG_COUNT", "0")) + env[f"GIT_CONFIG_KEY_{count}"] = "safe.directory" + env[f"GIT_CONFIG_VALUE_{count}"] = str(repo_dir) + env["GIT_CONFIG_COUNT"] = str(count + 1) + return env + + +def compute_key(repo_dir: Path, *, target: str | None = None) -> str: + """Deterministic sha256 over the git tree-hashes of cache-key paths. + + Uses ``git cat-file --batch-check`` for one subprocess invocation. Missing + paths fold in as a fixed null hash so the key remains deterministic + across repos that don't ship every input file. + + Raises ``subprocess.CalledProcessError`` if ``git`` itself fails (e.g. + not a repo) — callers SHOULD treat that as "no cache" and proceed. + """ + tgt = target if target is not None else target_triple() + stdin = "".join(f"HEAD:{p}\n" for p in CACHE_KEY_PATHS) + proc = subprocess.run( + ["git", "cat-file", "--batch-check"], + input=stdin, + cwd=str(repo_dir), + text=True, + capture_output=True, + check=True, + env=_git_safe_directory_env(repo_dir), + ) + lines = proc.stdout.splitlines() + if len(lines) != len(CACHE_KEY_PATHS): + raise RuntimeError( + f"git cat-file returned {len(lines)} lines, expected {len(CACHE_KEY_PATHS)}: {proc.stdout!r}" + ) + h = hashlib.sha256() + for path, line in zip(CACHE_KEY_PATHS, lines, strict=True): + stripped = line.strip() + if stripped.endswith("missing"): + tree_hash = _NULL_TREE_HASH + else: + # " " — take the first token as the tree/blob hash. + tree_hash = stripped.split(None, 1)[0] + h.update(f"{path}\t{tree_hash}\n".encode()) + h.update(f"TARGET\t{tgt}\n".encode()) + return h.hexdigest() + + +def _repo_slug(repo: str) -> str: + """Same convention as ``SandboxManager.pool_path``.""" + return repo.replace("/", "__") + + +def _atomic_link(src: Path, dst: Path) -> None: + """Hardlink ``src`` → ``dst``, replacing any existing ``dst`` atomically. + + Falls back to ``shutil.copy2`` on ``EXDEV`` (cross-filesystem). The + replace semantics use a sibling temp file + ``os.replace`` so a crash + mid-link doesn't leave ``dst`` half-overwritten. + """ + dst.parent.mkdir(parents=True, exist_ok=True) + tmp = dst.with_suffix(dst.suffix + f".tmp.{os.getpid()}") + try: + try: + os.link(src, tmp) + except OSError as exc: + if exc.errno != errno.EXDEV: + raise + shutil.copy2(src, tmp) + os.replace(tmp, dst) + finally: + # Best-effort cleanup if os.link succeeded but os.replace blew up. + try: + tmp.unlink() + except FileNotFoundError: + pass + + +def _atomic_copy(src: Path, dst: Path) -> None: + """Copy ``src`` → ``dst`` via a sibling temp file + ``os.replace``. + + Used for cached files that downstream tools rewrite via + ``open(..., 'w')`` (in-place truncate). Replacing the workspace dst + atomically means a fresh inode every populate — the cache file is + never mutated through a hardlink. + """ + dst.parent.mkdir(parents=True, exist_ok=True) + tmp = dst.with_suffix(dst.suffix + f".tmp.{os.getpid()}") + try: + shutil.copy2(src, tmp) + os.replace(tmp, dst) + finally: + try: + tmp.unlink() + except FileNotFoundError: + pass + + +@contextmanager +def _flock(path: Path) -> Generator[IO[bytes]]: + """Exclusive ``fcntl.flock`` on ``path`` (created if missing). + + ``flock`` is advisory but every caller goes through ``NativesCache``, so + cooperative locking is sufficient. POSIX-only — Windows is not a target. + """ + path.parent.mkdir(parents=True, exist_ok=True) + fh = open(path, "ab+") # noqa: SIM115 — managed by the context manager + try: + fcntl.flock(fh.fileno(), fcntl.LOCK_EX) + yield fh + finally: + try: + fcntl.flock(fh.fileno(), fcntl.LOCK_UN) + finally: + fh.close() + + +@dataclass(slots=True, frozen=True) +class CacheHit: + """Files copied/linked into the workspace by ``populate_workspace``.""" + + cache_dir: Path + files: tuple[Path, ...] + + +class NativesCache: + """Per-repo content-addressed cache of pi-natives build outputs.""" + + def __init__( + self, + root: Path, + *, + max_entries_per_repo: int = 8, + max_bytes: int = 4 * 1024**3, + ) -> None: + self.root = root + self.max_entries_per_repo = max(1, max_entries_per_repo) + self.max_bytes = max(0, max_bytes) + root.mkdir(parents=True, exist_ok=True) + + # ---- layout helpers ---- + def repo_root(self, repo: str) -> Path: + return self.root / _repo_slug(repo) + + def entry_dir(self, repo: str, key: str) -> Path: + return self.repo_root(repo) / key + + def lockfile(self, repo: str) -> Path: + return self.repo_root(repo) / _LOCKFILE_NAME + + # ---- query ---- + def lookup(self, repo: str, key: str) -> Path | None: + """Return the cache directory if ``key`` is present and complete.""" + entry = self.entry_dir(repo, key) + if not (entry / _MANIFEST_FILENAME).exists(): + return None + # A complete entry has a node file plus all companions. + if not list(entry.glob(_CACHED_NODE_GLOB)): + return None + for name in _CACHED_COMPANION_FILES: + if not (entry / name).exists(): + return None + return entry + + # ---- populate (workspace ← cache) ---- + def populate_workspace( + self, + repo: str, + key: str, + native_dir: Path, + ) -> CacheHit | None: + """Hardlink the `.node`, copy companions, into ``native_dir``. + + Returns the ``CacheHit`` on a hit; ``None`` on miss. Caller has + already computed ``key`` and verified ``native_dir`` exists. + + Why hardlink the .node but COPY the companions: the napi build's + ``installBinary`` replaces the .node via temp + rename (new inode, + cache safe), but ``installGeneratedBindings`` and ``gen-enums.ts`` + rewrite ``index.d.ts`` / ``index.js`` / ``embedded-addon.js`` with + plain ``open(..., 'w')`` — that's open-truncate-write IN PLACE on + Linux. A hardlinked companion would propagate the truncate into the + cache. Copies are independent inodes and absorb the rewrite safely. + """ + entry = self.lookup(repo, key) + if entry is None: + return None + native_dir.mkdir(parents=True, exist_ok=True) + copied: list[Path] = [] + for src in entry.glob(_CACHED_NODE_GLOB): + dst = native_dir / src.name + _atomic_link(src, dst) + copied.append(dst) + for name in _CACHED_COMPANION_FILES: + src = entry / name + dst = native_dir / name + _atomic_copy(src, dst) + copied.append(dst) + return CacheHit(cache_dir=entry, files=tuple(copied)) + + # ---- capture (cache ← workspace) ---- + def capture( + self, + repo: str, + key: str, + native_dir: Path, + *, + source_workspace: str | None = None, + commit: str | None = None, + ) -> Path | None: + """Atomically capture ``native_dir`` contents under ``key``. + + Returns the final cache directory on store, ``None`` if there was + nothing to capture or if another worker already populated the same + key (idempotent under flock). + """ + node_files = sorted(native_dir.glob(_CACHED_NODE_GLOB)) + if not node_files: + return None + # Every companion must exist or the entry would be incomplete. + for name in _CACHED_COMPANION_FILES: + if not (native_dir / name).exists(): + return None + + repo_root = self.repo_root(repo) + repo_root.mkdir(parents=True, exist_ok=True) + with _flock(self.lockfile(repo)): + # TOCTOU recheck: another worker may have captured the same key + # while we waited on the lock. + if self.lookup(repo, key) is not None: + return self.entry_dir(repo, key) + + final = self.entry_dir(repo, key) + staging = repo_root / f".{key}.tmp.{os.getpid()}" + if staging.exists(): + shutil.rmtree(staging, ignore_errors=True) + staging.mkdir(parents=True) + try: + # NOTE: capture uses COPY, not hardlink. Hardlinking a + # slot-owned workspace file into the cache would preserve + # the slot's ownership on the cached inode — defeating + # the setgid `omp` model that lets other slots read it. + # A copy creates a fresh inode owned by the orchestrator + # (root) and inherits gid `omp` from the setgid 2770 + # cache root. + for src in node_files: + _atomic_copy(src, staging / src.name) + for name in _CACHED_COMPANION_FILES: + _atomic_copy(native_dir / name, staging / name) + manifest = { + "key": key, + "target": target_triple(), + "captured_at": time.time(), + "source_workspace": source_workspace, + "commit": commit, + "node_files": [src.name for src in node_files], + } + (staging / _MANIFEST_FILENAME).write_text( + json.dumps(manifest, indent=2, sort_keys=True), encoding="utf-8" + ) + os.replace(staging, final) + except Exception: + shutil.rmtree(staging, ignore_errors=True) + raise + self._gc_locked(repo) + return final + + # ---- gc ---- + def gc(self, repo: str | None = None) -> int: + """Evict entries beyond per-repo or total caps. + + ``repo`` scopes to one repo when given; otherwise sweeps every repo + directory under ``root``. Returns the count of evicted entries. + """ + if repo is not None: + with _flock(self.lockfile(repo)): + return self._gc_locked(repo) + total = 0 + if not self.root.exists(): + return 0 + for child in self.root.iterdir(): + if not child.is_dir(): + continue + # Reconstruct repo identifier from directory name (best-effort; + # only used for lockfile path, not for any externally-visible + # identifier). + repo_name = child.name.replace("__", "/", 1) + try: + with _flock(self.lockfile(repo_name)): + total += self._gc_locked(repo_name) + except OSError as exc: + log.warning("natives_cache gc skip", extra={"repo": child.name, "err": str(exc)}) + return total + + def _gc_locked(self, repo: str) -> int: + """Caller MUST hold the per-repo flock.""" + repo_root = self.repo_root(repo) + if not repo_root.exists(): + return 0 + entries: list[tuple[float, int, Path]] = [] + for child in repo_root.iterdir(): + if not child.is_dir(): + # Stale staging dirs ("..tmp.") from a crashed + # capture: drop them opportunistically. + continue + if child.name.startswith("."): + shutil.rmtree(child, ignore_errors=True) + continue + manifest_path = child / _MANIFEST_FILENAME + if not manifest_path.exists(): + # Incomplete entry — evict. + shutil.rmtree(child, ignore_errors=True) + continue + try: + manifest = json.loads(manifest_path.read_text(encoding="utf-8")) + captured_at = float(manifest.get("captured_at", 0.0)) + except (OSError, ValueError, json.JSONDecodeError): + captured_at = manifest_path.stat().st_mtime + size = _dir_size(child) + entries.append((captured_at, size, child)) + entries.sort(key=lambda row: row[0]) # oldest first + + evicted = 0 + # 1. Per-repo entry-count cap (drop oldest). + while len(entries) > self.max_entries_per_repo: + _, _, victim = entries.pop(0) + shutil.rmtree(victim, ignore_errors=True) + evicted += 1 + + # 2. Per-repo byte cap (drop oldest until under). + if self.max_bytes > 0: + total = sum(size for _, size, _ in entries) + while total > self.max_bytes and len(entries) > 1: + _, size, victim = entries.pop(0) + shutil.rmtree(victim, ignore_errors=True) + total -= size + evicted += 1 + return evicted + + +def _dir_size(path: Path) -> int: + """Sum of file sizes under ``path``. Symlinks counted as their lstat + size (not the target). Errors swallowed — GC is best-effort.""" + total = 0 + for root, _dirs, files in os.walk(path): + for name in files: + try: + total += os.lstat(os.path.join(root, name)).st_size + except OSError: + pass + return total + + +__all__ = [ + "CACHE_KEY_PATHS", + "CacheHit", + "NativesCache", + "compute_key", + "target_triple", +] diff --git a/src/robomp/queue.py b/src/robomp/queue.py index 096cbeb90..4b90a96d7 100644 --- a/src/robomp/queue.py +++ b/src/robomp/queue.py @@ -96,6 +96,10 @@ class WorkerPool: log.info("recovered stuck events", extra={"count": recovered}) # Single dispatcher loop is simpler than N workers; concurrency is gated by the slot pool. self._workers.append(asyncio.create_task(self._dispatch_loop(), name="robomp-dispatch")) + # Periodic natives-cache GC, if enabled. Sleep-first so a freshly + # restarted orchestrator doesn't burn CPU on a cold cache. + if self.sandbox.natives_cache is not None and self.settings.natives_cache_gc_interval_seconds > 0: + self._workers.append(asyncio.create_task(self._natives_cache_gc_loop(), name="robomp-natives-gc")) async def stop(self, *, drain_timeout: float = 25.0, kill_timeout: float = 5.0) -> None: """Halt the dispatcher, then drain (or kill) in-flight `_run_event` tasks. @@ -154,6 +158,34 @@ class WorkerPool: with suppress(TimeoutError): await asyncio.wait(still_running, timeout=kill_timeout) + async def _natives_cache_gc_loop(self) -> None: + """Periodic sweep over every per-repo cache directory. + + Each iteration sleeps the configured interval first, then runs the + synchronous GC on a worker thread. Cancellation is the only exit; + any per-sweep failure is logged and the loop continues. + """ + cache = self.sandbox.natives_cache + if cache is None: # pragma: no cover — checked by caller + return + interval = self.settings.natives_cache_gc_interval_seconds + log.info("natives_cache gc loop online", extra={"interval": interval}) + try: + while not self._stop.is_set(): + try: + await asyncio.wait_for(self._stop.wait(), timeout=interval) + return # stop was set during the wait + except TimeoutError: + pass + try: + evicted = await asyncio.to_thread(cache.gc) + if evicted: + log.info("natives_cache gc swept", extra={"evicted": evicted}) + except Exception: + log.exception("natives_cache gc raised") + except asyncio.CancelledError: + raise + async def _dispatch_loop(self) -> None: log.info("dispatch loop online") try: diff --git a/src/robomp/sandbox.py b/src/robomp/sandbox.py index 5a0b41986..904ea0ad9 100644 --- a/src/robomp/sandbox.py +++ b/src/robomp/sandbox.py @@ -69,6 +69,8 @@ from robomp.git_ops import ( from robomp.git_ops import ( push as git_push, ) +from robomp.natives_cache import CacheHit, NativesCache +from robomp.natives_cache import compute_key as natives_compute_key log = logging.getLogger(__name__) @@ -607,10 +609,17 @@ class SandboxManager: (worktree add/remove, identity config, directory layout) is purely local. """ - def __init__(self, root: Path, *, transport: GitTransport | None = None) -> None: + def __init__( + self, + root: Path, + *, + transport: GitTransport | None = None, + natives_cache: NativesCache | None = None, + ) -> None: self.root = root self.pool = root / "_pool" self.transport: GitTransport = transport or LocalGitTransport(token=None) + self.natives_cache = natives_cache root.mkdir(parents=True, exist_ok=True) self.pool.mkdir(parents=True, exist_ok=True) @@ -761,7 +770,7 @@ class SandboxManager: if proc.returncode != 0: raise GitCommandError(command, proc.returncode, proc.stdout, proc.stderr) _share_git_metadata_with_slots(repo_dir, slot_uid) - return Workspace( + workspace = Workspace( root=ws_root, repo_dir=repo_dir, session_dir=session_dir, @@ -771,6 +780,95 @@ class SandboxManager: repo_full_name=repo, issue_number=number, ) + # Best-effort: hardlink pre-built natives in if we've cached this + # source state before. Runs AFTER the slot chown so the cache inode + # keeps its `root:omp` ownership (the slot reads through group `omp`); + # write-temp + rename in the napi build replaces with a new inode if + # the agent rebuilds, so the cached file is never mutated. + self._populate_natives_cache(workspace, slot_uid=slot_uid) + return workspace + + def _populate_natives_cache(self, workspace: Workspace, *, slot_uid: int | None = None) -> None: + """Try to hardlink cached pi-natives artifacts into the worktree. + + Best-effort: any failure (no cache configured, non-git worktree, + cache miss, link error) is logged at debug and swallowed. The agent + falls back to a fresh napi build, exactly as it would without the + cache. + + Post-populate, the populated `packages/natives/native/` directory + and the COPIED companion files are chowned to the slot so the slot + can rebuild via temp + rename in that directory. The hardlinked + `.node` files are LEFT at `root:omp` ownership — chowning them + would chown the cache file too (shared inode), breaking the + cross-slot sharing model. The slot reads them via group `omp`. + """ + cache = self.natives_cache + if cache is None: + return + native_dir = workspace.repo_dir / "packages" / "natives" / "native" + # NOTE: we deliberately do NOT require `native_dir.exists()` here. On + # a cache miss `populate_workspace` returns None without creating any + # directory; on a hit it mkdirs and copies in. That's the right + # behavior — a hit by definition implies this repo's source state + # produces natives, so creating the dir is correct. + try: + key = natives_compute_key(workspace.repo_dir) + except (subprocess.CalledProcessError, RuntimeError, OSError) as exc: + log.debug( + "natives_cache key compute failed", + extra={"workspace": workspace.workspace_key, "err": redact_credentials(str(exc))}, + ) + return + try: + hit = cache.populate_workspace(workspace.repo_full_name, key, native_dir) + except OSError as exc: + log.warning( + "natives_cache populate failed", + extra={"workspace": workspace.workspace_key, "key": key, "err": str(exc)}, + ) + return + if hit is not None and _slot_permissions_active(slot_uid): + assert slot_uid is not None + self._chown_natives_for_slot(native_dir, hit, slot_uid=slot_uid) + log.info( + "natives_cache", + extra={ + "action": "hit" if hit is not None else "miss", + "workspace": workspace.workspace_key, + "repo": workspace.repo_full_name, + "key": key, + "files": [str(p.name) for p in hit.files] if hit is not None else [], + }, + ) + + @staticmethod + def _chown_natives_for_slot(native_dir: Path, hit: CacheHit, *, slot_uid: int) -> None: + """Hand the populated native dir to the slot WITHOUT touching the + hardlinked `.node` inodes (those are shared with the cache). + + Files whose names match a cached `.node` are skipped — they are + hardlinks back into the root:omp cache and the slot reads them via + group `omp`. Everything else (the directory itself, copied + companions) is chowned to the slot so the slot can rebuild via + temp + rename. + """ + try: + os.chown(native_dir, slot_uid, slot_uid) + except OSError as exc: + log.warning("natives_cache chown dir failed", extra={"err": str(exc)}) + return + node_basenames = {p.name for p in hit.files if p.name.endswith(".node")} + for child in native_dir.iterdir(): + if child.name in node_basenames: + continue # hardlink to cache — must not chown + try: + os.chown(child, slot_uid, slot_uid, follow_symlinks=False) + except OSError as exc: + log.warning( + "natives_cache chown companion failed", + extra={"file": str(child), "err": str(exc)}, + ) def remove_workspace(self, *, repo: str, number: int) -> None: ws_root = self.workspace_root(repo, number) diff --git a/src/robomp/server.py b/src/robomp/server.py index a1e7d1050..2cab9702f 100644 --- a/src/robomp/server.py +++ b/src/robomp/server.py @@ -36,6 +36,7 @@ from robomp.manual_triage import ( enqueue_manual_triage, parse_issue_ref, ) +from robomp.natives_cache import NativesCache from robomp.proxy_client import GitHubProxyClient, ProxyGitTransport from robomp.queue import WorkerPool from robomp.sandbox import SandboxManager @@ -237,7 +238,18 @@ def _build_orchestrator(cfg: Settings) -> tuple[GitHubBackend, ProxyGitTransport def _build_state(settings: Settings) -> dict[str, Any]: db = get_database(settings.sqlite_path) github, git_transport = _build_orchestrator(settings) - sandbox = SandboxManager(settings.workspace_root, transport=git_transport) + natives_cache: NativesCache | None = None + if settings.natives_cache_enabled: + natives_cache = NativesCache( + settings.natives_cache_root, + max_entries_per_repo=settings.natives_cache_max_entries_per_repo, + max_bytes=settings.natives_cache_max_bytes, + ) + sandbox = SandboxManager( + settings.workspace_root, + transport=git_transport, + natives_cache=natives_cache, + ) pool = WorkerPool(settings=settings, db=db, github=github, sandbox=sandbox, git_transport=git_transport) autoclose = AutocloseScheduler(settings=settings, db=db, github=github) return { @@ -246,6 +258,7 @@ def _build_state(settings: Settings) -> dict[str, Any]: "github": github, "git_transport": git_transport, "sandbox": sandbox, + "natives_cache": natives_cache, "pool": pool, "issue_browse_cache": _IssueBrowseCache(), "autoclose": autoclose, diff --git a/src/robomp/tasks.py b/src/robomp/tasks.py index f4a5a2696..af75bb632 100644 --- a/src/robomp/tasks.py +++ b/src/robomp/tasks.py @@ -285,6 +285,7 @@ async def triage_issue( delivery_id=delivery_id, attempts=attempts, slot_uid=slot_uid, + natives_cache=sandbox.natives_cache, ) await run_task(task_kind="triage_issue", inputs=inputs) @@ -346,6 +347,7 @@ async def handle_comment( delivery_id=delivery_id, attempts=attempts, slot_uid=slot_uid, + natives_cache=sandbox.natives_cache, ) directive = await _attach_thread(github, directive, repo.full_name, issue.number, is_pr=False) await run_task(task_kind="triage_issue", inputs=inputs, directive=directive) @@ -397,6 +399,7 @@ async def handle_comment( delivery_id=delivery_id, attempts=attempts, slot_uid=slot_uid, + natives_cache=sandbox.natives_cache, ) directive = await _attach_thread(github, directive, repo.full_name, issue.number, is_pr=False) await run_task(task_kind="handle_comment", inputs=inputs, comment=comment, directive=directive) @@ -424,6 +427,7 @@ async def handle_comment( delivery_id=delivery_id, attempts=attempts, slot_uid=slot_uid, + natives_cache=sandbox.natives_cache, ) directive = await _attach_thread(github, directive, repo.full_name, issue.number, is_pr=False) await run_task(task_kind="handle_comment", inputs=inputs, comment=comment, directive=directive) @@ -517,6 +521,7 @@ async def handle_review( delivery_id=delivery_id, attempts=attempts, slot_uid=slot_uid, + natives_cache=sandbox.natives_cache, ) await run_task( task_kind="handle_review", @@ -659,6 +664,7 @@ async def handle_pr_conversation( delivery_id=delivery_id, attempts=attempts, slot_uid=slot_uid, + natives_cache=sandbox.natives_cache, ) directive = await _attach_thread(github, directive, repo_full, pr_number, is_pr=True) await run_task(task_kind="handle_comment", inputs=inputs, comment=comment, pr_number=pr_number, directive=directive) diff --git a/src/robomp/worker.py b/src/robomp/worker.py index 0f6d00898..28973b2b2 100644 --- a/src/robomp/worker.py +++ b/src/robomp/worker.py @@ -36,6 +36,8 @@ from robomp.db import Database, issue_key from robomp.github_backend import GitHubBackend from robomp.github_client import CommentInfo, IssueInfo, RepoInfo from robomp.host_tools import AbortController, ToolBindings, _git_identity_env +from robomp.natives_cache import NativesCache +from robomp.natives_cache import compute_key as natives_compute_key from robomp.sandbox import GitTransport, Workspace, _prepare_slot_runtime_env, _safe_directory_env log = logging.getLogger(__name__) @@ -55,6 +57,7 @@ class TaskInputs: delivery_id: str attempts: int = 0 slot_uid: int | None = None + natives_cache: NativesCache | None = None @dataclass(slots=True, frozen=True) @@ -618,14 +621,68 @@ async def run_task( directive=directive, resuming=resuming, ) - return await asyncio.to_thread( - _run_rpc_blocking, - inputs, - task_kind=task_kind, - prompt=prompt, - loop=loop, - bindings=bindings, - directive=directive, + try: + result = await asyncio.to_thread( + _run_rpc_blocking, + inputs, + task_kind=task_kind, + prompt=prompt, + loop=loop, + bindings=bindings, + directive=directive, + ) + except BaseException: + # Failed/aborted task: NEVER capture, the artifacts may be inconsistent + # with the source state and would poison the cache. + raise + else: + await asyncio.to_thread(_capture_natives_cache, inputs) + return result + + +def _capture_natives_cache(inputs: TaskInputs) -> None: + """Best-effort: store the workspace's fresh natives under its current key. + + Runs after a successful task on a worker thread. ANY failure is logged + and swallowed — cache errors NEVER fail a task. + """ + cache = inputs.natives_cache + if cache is None: + return + workspace = inputs.workspace + native_dir = workspace.repo_dir / "packages" / "natives" / "native" + if not native_dir.exists(): + return + try: + key = natives_compute_key(workspace.repo_dir) + except Exception as exc: + log.debug( + "natives_cache capture key compute failed", + extra={"workspace": workspace.workspace_key, "err": str(exc)}, + ) + return + try: + stored = cache.capture( + workspace.repo_full_name, + key, + native_dir, + source_workspace=workspace.workspace_key, + ) + except Exception as exc: + log.warning( + "natives_cache capture failed", + extra={"workspace": workspace.workspace_key, "key": key, "err": str(exc)}, + ) + return + log.info( + "natives_cache", + extra={ + "action": "stored" if stored is not None else "skip", + "workspace": workspace.workspace_key, + "repo": workspace.repo_full_name, + "key": key, + "cache_dir": str(stored) if stored else None, + }, ) diff --git a/tests/conftest.py b/tests/conftest.py index cb0d33b6c..8efe12382 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -56,6 +56,13 @@ def _baseline_env(tmp_path: Path) -> dict[str, str]: "ROBOMP_WORKSPACE_ROOT": str(tmp_path / "workspaces"), "ROBOMP_SQLITE_PATH": str(tmp_path / "robomp.sqlite"), "ROBOMP_LOG_DIR": str(tmp_path / "logs"), + # Production default is `/data/cache/pi-natives` (provisioned by the + # container entrypoint). Tests need a writable, isolated path; we also + # default-disable the cache so its background GC loop doesn't add + # noise to event-dispatcher timing assertions. Tests that want the + # cache flip `ROBOMP_NATIVES_CACHE_ENABLED=true` explicitly. + "ROBOMP_NATIVES_CACHE_ROOT": str(tmp_path / "natives-cache"), + "ROBOMP_NATIVES_CACHE_ENABLED": "false", } diff --git a/tests/test_natives_cache.py b/tests/test_natives_cache.py new file mode 100644 index 000000000..534c830a6 --- /dev/null +++ b/tests/test_natives_cache.py @@ -0,0 +1,414 @@ +"""Unit tests for `robomp.natives_cache`. + +The module's filesystem operations (hardlink, atomic rename, flock) are +exercised against `tmp_path`; nothing here requires a running orchestrator. +""" + +from __future__ import annotations + +import errno +import json +import os +import subprocess +import threading +import time +from pathlib import Path + +import pytest + +from robomp.natives_cache import ( + CACHE_KEY_PATHS, + NativesCache, + _atomic_link, + compute_key, +) + +REPO = "octo/widget" + + +# ---- repo + workspace fixtures ---- + + +def _git(args: list[str], cwd: Path) -> None: + subprocess.run( + ["git", *args], + cwd=str(cwd), + check=True, + capture_output=True, + text=True, + env=os.environ + | { + "GIT_AUTHOR_NAME": "t", + "GIT_AUTHOR_EMAIL": "t@t", + "GIT_COMMITTER_NAME": "t", + "GIT_COMMITTER_EMAIL": "t@t", + }, + ) + + +def _seed_repo(root: Path, *, with_all_inputs: bool = True) -> Path: + """Stand up a minimal repo with the cache-key inputs present. + + When `with_all_inputs=False`, only `Cargo.lock` exists — used to exercise + the missing-path code path in `compute_key`. + """ + root.mkdir(parents=True, exist_ok=True) + _git(["init", "--initial-branch=main", str(root)], cwd=root.parent) + (root / "Cargo.lock").write_text("# lock v1\n") + if with_all_inputs: + (root / "Cargo.toml").write_text("[workspace]\nmembers = ['crates/*']\n") + (root / "rust-toolchain.toml").write_text('[toolchain]\nchannel = "1.85.0"\n') + crates = root / "crates" / "pi-natives" + crates.mkdir(parents=True) + (crates / "Cargo.toml").write_text('[package]\nname = "pi-natives"\n') + (crates / "src.rs").write_text("// source\n") + natives = root / "packages" / "natives" + natives.mkdir(parents=True) + (natives / "package.json").write_text('{"name":"@oh-my-pi/pi-natives"}\n') + scripts = natives / "scripts" + scripts.mkdir() + (scripts / "build-native.ts").write_text("// build script\n") + native_dir = natives / "native" + native_dir.mkdir() + (native_dir / "index.d.ts").write_text("// initial typings\n") + _git(["-C", str(root), "add", "."], cwd=root.parent) + _git(["-C", str(root), "commit", "-m", "init"], cwd=root.parent) + return root + + +def _populate_built_artifacts(repo_dir: Path, *, body: bytes = b"\x7fELF...native") -> Path: + """Fill `packages/natives/native/` with a complete built-artifact set.""" + native_dir = repo_dir / "packages" / "natives" / "native" + native_dir.mkdir(parents=True, exist_ok=True) + (native_dir / "pi_natives.linux-arm64.node").write_bytes(body) + (native_dir / "index.d.ts").write_text("export const X: number;\n") + (native_dir / "index.js").write_text("export const X = 1;\n") + (native_dir / "embedded-addon.js").write_text("export const embeddedAddon = null;\n") + return native_dir + + +# ---- compute_key ---- + + +def test_compute_key_deterministic_across_clones(tmp_path: Path) -> None: + a = _seed_repo(tmp_path / "a") + b_root = tmp_path / "b" + subprocess.run(["git", "clone", str(a), str(b_root)], check=True, capture_output=True, text=True) + key_a = compute_key(a, target="linux-arm64") + key_b = compute_key(b_root, target="linux-arm64") + assert key_a == key_b + + +def test_compute_key_changes_when_each_input_changes(tmp_path: Path) -> None: + base = _seed_repo(tmp_path / "base") + base_key = compute_key(base, target="linux-arm64") + + # Touching a file under each key path must shift the key. + mutations: dict[str, tuple[str, str]] = { + "crates": ("crates/pi-natives/src.rs", "// new comment\n"), + "Cargo.lock": ("Cargo.lock", "# lock v2\n"), + "Cargo.toml": ("Cargo.toml", "[workspace]\nmembers = ['crates/*', 'extra']\n"), + "rust-toolchain.toml": ("rust-toolchain.toml", '[toolchain]\nchannel = "1.86.0"\n'), + "packages/natives": ("packages/natives/scripts/build-native.ts", "// edited\n"), + } + for label, (rel, body) in mutations.items(): + clone = tmp_path / f"clone-{label.replace('/', '-')}" + subprocess.run( + ["git", "clone", str(base), str(clone)], + check=True, + capture_output=True, + text=True, + ) + target = clone / rel + target.parent.mkdir(parents=True, exist_ok=True) + target.write_text(body) + _git(["-C", str(clone), "add", "."], cwd=clone.parent) + _git(["-C", str(clone), "commit", "-m", f"mutate {label}"], cwd=clone.parent) + new_key = compute_key(clone, target="linux-arm64") + assert new_key != base_key, f"key did not change after mutating {label}" + + +def test_compute_key_target_triple_changes_key(tmp_path: Path) -> None: + repo = _seed_repo(tmp_path / "repo") + arm = compute_key(repo, target="linux-arm64") + x64 = compute_key(repo, target="linux-x64-modern") + assert arm != x64 + + +def test_compute_key_handles_missing_inputs(tmp_path: Path) -> None: + """Missing key paths fold to a fixed null hash → key still deterministic.""" + repo = _seed_repo(tmp_path / "repo", with_all_inputs=False) + # Lock-only repo: should compute without error, and adding a tracked + # crates/ subtree shifts the key. + key_before = compute_key(repo, target="linux-arm64") + crates = repo / "crates" / "pi-natives" + crates.mkdir(parents=True) + (crates / "lib.rs").write_text("// new\n") + _git(["-C", str(repo), "add", "."], cwd=repo.parent) + _git(["-C", str(repo), "commit", "-m", "add crates"], cwd=repo.parent) + key_after = compute_key(repo, target="linux-arm64") + assert key_before != key_after + + +def test_compute_key_uses_all_documented_paths() -> None: + # Sanity contract: the exported path list IS the input set. + assert CACHE_KEY_PATHS == ( + "crates", + "Cargo.lock", + "Cargo.toml", + "rust-toolchain.toml", + "packages/natives", + ) + + +def test_compute_key_raises_on_non_repo(tmp_path: Path) -> None: + with pytest.raises(subprocess.CalledProcessError): + compute_key(tmp_path, target="linux-arm64") + + +# ---- populate / capture ---- + + +def _cache(tmp_path: Path, **kwargs: object) -> NativesCache: + return NativesCache(tmp_path / "natives-cache", **kwargs) # type: ignore[arg-type] + + +def test_populate_workspace_miss_is_noop(tmp_path: Path) -> None: + cache = _cache(tmp_path) + repo_dir = _seed_repo(tmp_path / "ws" / "repo") + native_dir = repo_dir / "packages" / "natives" / "native" + before = sorted(p.name for p in native_dir.iterdir()) + hit = cache.populate_workspace(REPO, "deadbeef" * 8, native_dir) + after = sorted(p.name for p in native_dir.iterdir()) + assert hit is None + assert before == after + + +def test_capture_then_populate_shares_node_inode_but_copies_companions(tmp_path: Path) -> None: + cache = _cache(tmp_path) + src_repo = _seed_repo(tmp_path / "src" / "repo") + native_dir = _populate_built_artifacts(src_repo) + key = compute_key(src_repo, target="linux-arm64") + stored = cache.capture(REPO, key, native_dir, source_workspace="src__001") + assert stored is not None + manifest = json.loads((stored / "manifest.json").read_text()) + assert manifest["key"] == key + assert "pi_natives.linux-arm64.node" in manifest["node_files"] + + # Populate a fresh workspace from the same source state. + dst_repo = src_repo.parent.parent / "dst" / "repo" + dst_repo.mkdir(parents=True) + _git(["clone", str(src_repo), str(dst_repo)], cwd=dst_repo.parent) + dst_native = dst_repo / "packages" / "natives" / "native" + dst_native.mkdir(parents=True, exist_ok=True) + hit = cache.populate_workspace(REPO, key, dst_native) + assert hit is not None + assert {p.name for p in hit.files} >= { + "pi_natives.linux-arm64.node", + "index.d.ts", + "index.js", + "embedded-addon.js", + } + # The `.node` is hardlinked: same inode, nlink ≥ 2. + cached_node = stored / "pi_natives.linux-arm64.node" + workspace_node = dst_native / "pi_natives.linux-arm64.node" + assert cached_node.stat().st_ino == workspace_node.stat().st_ino + assert cached_node.stat().st_nlink >= 2 + # Companions are COPIED (independent inodes): in-place rewrite in the + # workspace (gen-enums.ts / installGeneratedBindings open-truncate-write) + # MUST NOT mutate the cached copy. + for name in ("index.d.ts", "index.js", "embedded-addon.js"): + cached_companion = stored / name + ws_companion = dst_native / name + assert cached_companion.stat().st_ino != ws_companion.stat().st_ino, name + original = cached_companion.read_text() + ws_companion.write_text("rewritten\n") + assert cached_companion.read_text() == original, name + + +def test_capture_skips_when_artifacts_incomplete(tmp_path: Path) -> None: + cache = _cache(tmp_path) + repo = _seed_repo(tmp_path / "ws" / "repo") + native_dir = repo / "packages" / "natives" / "native" + # Only the .node — missing companions → capture refuses. + (native_dir / "pi_natives.linux-arm64.node").write_bytes(b"x") + assert cache.capture(REPO, "k", native_dir) is None + # And no entry was created. + assert not cache.entry_dir(REPO, "k").exists() + + +def test_capture_is_idempotent_under_lock(tmp_path: Path) -> None: + """Two concurrent captures of the same key end with one final entry.""" + cache = _cache(tmp_path) + src_repo = _seed_repo(tmp_path / "src" / "repo") + _populate_built_artifacts(src_repo) + key = compute_key(src_repo, target="linux-arm64") + native_dir = src_repo / "packages" / "natives" / "native" + + results: list[Path | None] = [] + barrier = threading.Barrier(2) + + def run() -> None: + barrier.wait() + results.append(cache.capture(REPO, key, native_dir)) + + threads = [threading.Thread(target=run) for _ in range(2)] + for t in threads: + t.start() + for t in threads: + t.join() + # Both calls succeed (one captures, the other recognizes the entry). + assert all(isinstance(r, Path) for r in results) + # Exactly one final entry directory (no leftover staging). + repo_root = cache.repo_root(REPO) + final_dirs = [p for p in repo_root.iterdir() if p.is_dir() and not p.name.startswith(".")] + assert len(final_dirs) == 1 + assert final_dirs[0].name == key + + +def test_populate_cross_device_falls_back_to_copy(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> None: + cache = _cache(tmp_path) + src_repo = _seed_repo(tmp_path / "src" / "repo") + _populate_built_artifacts(src_repo) + key = compute_key(src_repo, target="linux-arm64") + cache.capture(REPO, key, src_repo / "packages" / "natives" / "native") + + dst_native = tmp_path / "ws2" / "packages" / "natives" / "native" + dst_native.mkdir(parents=True) + + # Simulate cross-device hardlink failure for every os.link call. + real_link = os.link + + def fake_link(src, dst, *args, **kwargs): # type: ignore[no-untyped-def] + raise OSError(errno.EXDEV, "Cross-device link", str(src)) + + monkeypatch.setattr(os, "link", fake_link) + try: + hit = cache.populate_workspace(REPO, key, dst_native) + finally: + monkeypatch.setattr(os, "link", real_link) + assert hit is not None + # Files exist (via copy) but are distinct inodes from the cache. + cached_node = cache.entry_dir(REPO, key) / "pi_natives.linux-arm64.node" + copied_node = dst_native / "pi_natives.linux-arm64.node" + assert copied_node.exists() + assert cached_node.stat().st_ino != copied_node.stat().st_ino + + +def test_populate_replaces_existing_file_atomically(tmp_path: Path) -> None: + cache = _cache(tmp_path) + src_repo = _seed_repo(tmp_path / "src" / "repo") + _populate_built_artifacts(src_repo, body=b"\x7fELF.A") + key = compute_key(src_repo, target="linux-arm64") + cache.capture(REPO, key, src_repo / "packages" / "natives" / "native") + + dst_native = tmp_path / "dst" / "packages" / "natives" / "native" + dst_native.mkdir(parents=True) + # Pre-existing stub bytes — populate must replace, not append/error. + target = dst_native / "pi_natives.linux-arm64.node" + target.write_bytes(b"old-stub") + hit = cache.populate_workspace(REPO, key, dst_native) + assert hit is not None + assert target.read_bytes() == b"\x7fELF.A" + + +# ---- gc ---- + + +def _stamp_entry(cache: NativesCache, repo: str, key: str, captured_at: float) -> Path: + entry = cache.entry_dir(repo, key) + entry.mkdir(parents=True, exist_ok=True) + (entry / "pi_natives.linux-arm64.node").write_bytes(b"x" * 1024) + (entry / "index.d.ts").write_text("") + (entry / "index.js").write_text("") + (entry / "embedded-addon.js").write_text("") + (entry / "manifest.json").write_text( + json.dumps({"key": key, "captured_at": captured_at, "node_files": ["pi_natives.linux-arm64.node"]}) + ) + return entry + + +def test_gc_evicts_oldest_beyond_entry_cap(tmp_path: Path) -> None: + cache = _cache(tmp_path, max_entries_per_repo=2, max_bytes=0) + now = time.time() + _stamp_entry(cache, REPO, "k1", now - 300) + _stamp_entry(cache, REPO, "k2", now - 200) + _stamp_entry(cache, REPO, "k3", now - 100) + evicted = cache.gc(REPO) + assert evicted == 1 + remaining = {p.name for p in cache.repo_root(REPO).iterdir() if p.is_dir() and not p.name.startswith(".")} + assert remaining == {"k2", "k3"} + + +def test_gc_evicts_for_byte_cap(tmp_path: Path) -> None: + cache = _cache(tmp_path, max_entries_per_repo=8, max_bytes=2500) + now = time.time() + # Each entry weighs ~1024 bytes (the .node); 3 entries → ~3072 bytes > cap. + _stamp_entry(cache, REPO, "k1", now - 300) + _stamp_entry(cache, REPO, "k2", now - 200) + _stamp_entry(cache, REPO, "k3", now - 100) + cache.gc(REPO) + remaining = {p.name for p in cache.repo_root(REPO).iterdir() if p.is_dir() and not p.name.startswith(".")} + # Oldest evicted; at least one survives. + assert "k1" not in remaining + assert remaining <= {"k2", "k3"} + assert remaining + + +def test_gc_preserves_workspace_hardlinks(tmp_path: Path) -> None: + """Evicting a cache entry must NOT delete the file from workspaces that + hardlinked it — kernel inode refcount keeps the data alive.""" + cache = _cache(tmp_path, max_entries_per_repo=1, max_bytes=0) + now = time.time() + entry = _stamp_entry(cache, REPO, "k1", now - 500) + _stamp_entry(cache, REPO, "k2", now - 100) + # Workspace hardlinks the older entry's .node before GC runs. + ws_node = tmp_path / "ws" / "pi_natives.linux-arm64.node" + ws_node.parent.mkdir(parents=True) + os.link(entry / "pi_natives.linux-arm64.node", ws_node) + cache.gc(REPO) + assert not entry.exists() # cache directory swept + assert ws_node.exists() # workspace file survives via inode refcount + assert ws_node.read_bytes() == b"x" * 1024 + + +def test_gc_clears_stale_staging_dirs(tmp_path: Path) -> None: + cache = _cache(tmp_path) + repo_root = cache.repo_root(REPO) + repo_root.mkdir(parents=True) + stale = repo_root / ".aabb.tmp.99999" + stale.mkdir() + (stale / "leaked").write_text("from a crashed capture") + cache.gc(REPO) + assert not stale.exists() + + +def test_gc_drops_entry_with_missing_manifest(tmp_path: Path) -> None: + cache = _cache(tmp_path) + incomplete = cache.entry_dir(REPO, "bogus") + incomplete.mkdir(parents=True) + (incomplete / "pi_natives.linux-arm64.node").write_bytes(b"x") + cache.gc(REPO) + assert not incomplete.exists() + + +def test_lookup_rejects_incomplete_entry(tmp_path: Path) -> None: + cache = _cache(tmp_path) + entry = cache.entry_dir(REPO, "partial") + entry.mkdir(parents=True) + (entry / "manifest.json").write_text("{}") + # No .node → no hit even though manifest exists. + assert cache.lookup(REPO, "partial") is None + + +# ---- _atomic_link ---- + + +def test_atomic_link_replaces_existing_target(tmp_path: Path) -> None: + src = tmp_path / "src" + src.write_bytes(b"new") + dst = tmp_path / "dst" + dst.write_bytes(b"old") + _atomic_link(src, dst) + assert dst.read_bytes() == b"new" + assert dst.stat().st_ino == src.stat().st_ino diff --git a/tests/test_permissions_e2e.py b/tests/test_permissions_e2e.py index f8dc8f21b..4b7a43204 100644 --- a/tests/test_permissions_e2e.py +++ b/tests/test_permissions_e2e.py @@ -17,6 +17,8 @@ from robomp import host_tools from robomp.db import Database from robomp.github_backend import GitHubBackend from robomp.github_client import IssueInfo, RepoInfo +from robomp.natives_cache import NativesCache +from robomp.natives_cache import compute_key as natives_compute_key from robomp.sandbox import LocalGitTransport, SandboxManager, Workspace pytestmark = pytest.mark.skipif( @@ -333,3 +335,148 @@ def test_git_pool_metadata_survives_root_push_and_retry_slot( ).stdout.strip() assert remote_retry_head == retry_head assert retry_head != first_head + + +def _prepare_shared_natives_cache(slot_tmp_path: Path) -> NativesCache: + """Provision `/data/cache/pi-natives` shape (root:omp, setgid 2770).""" + cache_root = slot_tmp_path / "cache" / "pi-natives" + cache_root.mkdir(parents=True) + os.chown(cache_root, 0, _SHARED_OMP_GID) + cache_root.chmod(0o2770) + return NativesCache(cache_root) + + +def _stage_built_natives(bindings: host_tools.ToolBindings, *, body: str = "ELFx") -> None: + """Mirror what a napi build would leave in `packages/natives/native/`. + + Writes the four cached files AS THE SLOT so ownership matches a real + post-build workspace; capture pulls these into the cache. + """ + _write_as_slot(bindings, "packages/natives/native/pi_natives.linux-arm64.node", body) + _write_as_slot(bindings, "packages/natives/native/index.d.ts", "export const X: number;\n") + _write_as_slot(bindings, "packages/natives/native/index.js", "export const X = 1;\n") + _write_as_slot( + bindings, + "packages/natives/native/embedded-addon.js", + "export const embeddedAddon = null;\n", + ) + + +def test_natives_cache_shares_artifacts_across_slot_workspaces( + slot_tmp_path: Path, + upstream_repo: Path, + db: Database, + tool_loop: asyncio.AbstractEventLoop, +) -> None: + """End-to-end: capture under slot 1, populate under slot 2, prove that: + + 1. A capture from a slot-owned workspace lands in the shared cache with + group `omp` setgid inheritance so any other slot can read it. + 2. ensure_workspace under a different slot UID auto-populates the cached + `.node` (hardlink, inode shared) and copies the companions. + 3. Slot 2 can read the populated `.node`, and a temp-rename rebuild + (mirroring napi's `installBinary`) leaves the cache entry intact. + 4. An in-place truncate-rewrite of a companion (mirroring `gen-enums.ts` + / `installGeneratedBindings`) does NOT mutate the cached companion — + this is exactly why companions are copied, not hardlinked. + """ + _require_linux_root_toolchain() + workspaces = slot_tmp_path / "workspaces" + natives_cache = _prepare_shared_natives_cache(slot_tmp_path) + manager = SandboxManager( + workspaces, + transport=LocalGitTransport(token=None), + natives_cache=natives_cache, + ) + + # --- Workspace 1: stage built artifacts and capture them as the orchestrator. --- + ws1 = manager.ensure_workspace( + repo=_REPO, + number=301, + title="natives cache producer", + clone_url=str(upstream_repo), + default_branch="main", + author_name=_AUTHOR_NAME, + author_email=_AUTHOR_EMAIL, + slot_uid=_SLOT_ONE, + ) + bindings1 = _bindings(db=db, tool_loop=tool_loop, workspace=ws1, upstream=upstream_repo, slot_uid=_SLOT_ONE) + _stage_built_natives(bindings1, body="ELFx-original") + + key = natives_compute_key(ws1.repo_dir, target="linux-arm64") + native_dir1 = ws1.repo_dir / "packages" / "natives" / "native" + stored = natives_cache.capture(_REPO, key, native_dir1, source_workspace=ws1.workspace_key) + assert stored is not None + cached_node = stored / "pi_natives.linux-arm64.node" + cached_companion = stored / "index.d.ts" + # Cache root is setgid `omp`; new files inherit gid `omp` so any slot + # with `extra_groups=[omp]` can read them. + assert cached_node.stat().st_gid == _SHARED_OMP_GID + assert cached_companion.stat().st_gid == _SHARED_OMP_GID + + # --- Workspace 2: a different slot UID gets auto-populated on ensure. --- + ws2 = manager.ensure_workspace( + repo=_REPO, + number=302, + title="natives cache consumer", + clone_url=str(upstream_repo), + default_branch="main", + author_name=_AUTHOR_NAME, + author_email=_AUTHOR_EMAIL, + slot_uid=_SLOT_TWO, + ) + bindings2 = _bindings(db=db, tool_loop=tool_loop, workspace=ws2, upstream=upstream_repo, slot_uid=_SLOT_TWO) + native_dir2 = ws2.repo_dir / "packages" / "natives" / "native" + ws2_node = native_dir2 / "pi_natives.linux-arm64.node" + ws2_companion = native_dir2 / "index.d.ts" + assert ws2_node.exists(), "auto-populate must hardlink the .node into ws2" + assert ws2_companion.exists(), "auto-populate must copy companions into ws2" + + # The .node is hardlinked: same inode, nlink ≥ 2. + assert ws2_node.stat().st_ino == cached_node.stat().st_ino + assert cached_node.stat().st_nlink >= 2 + # The companion is COPIED: independent inode. + assert ws2_companion.stat().st_ino != cached_companion.stat().st_ino + + # Slot 2 must be able to read the populated artifacts (group omp + 0660 + # via setgid inheritance from the cache root). + _run_ok(bindings2, ["test", "-r", "packages/natives/native/pi_natives.linux-arm64.node"]) + _run_ok(bindings2, ["test", "-r", "packages/natives/native/index.d.ts"]) + + # --- Rebuild simulation: napi's installBinary does temp + rename. --- + # Mirrors `fs.copyFile(src, tempPath); fs.rename(tempPath, dest)`. + _run_ok( + bindings2, + [ + "python3", + "-c", + ( + "import os, sys; " + "dest = sys.argv[1]; " + "tmp = dest + '.tmp.rebuild'; " + "open(tmp, 'wb').write(b'REBUILT'); " + "os.rename(tmp, dest)" + ), + "packages/natives/native/pi_natives.linux-arm64.node", + ], + ) + # Workspace sees the rebuilt bytes; cache is untouched (new inode in ws). + assert ws2_node.read_bytes() == b"REBUILT" + assert cached_node.read_bytes() == b"ELFx-original" + assert ws2_node.stat().st_ino != cached_node.stat().st_ino + + # --- Companion-rewrite simulation: gen-enums.ts open-truncate-writes. --- + # Mirrors `await Bun.write(jsPath, js)` / Python `Path.write_text`. + _write_as_slot( + bindings2, + "packages/natives/native/index.d.ts", + "// regenerated by gen-enums\n", + ) + assert ws2_companion.read_text() == "// regenerated by gen-enums\n" + # Cache copy stays at its original content — copies absorbed the rewrite. + assert cached_companion.read_text() == "export const X: number;\n" + + # --- Recapture from ws2 (different key now — but same key here since + # tree didn't change) is idempotent under the flock. --- + again = natives_cache.capture(_REPO, key, native_dir2, source_workspace=ws2.workspace_key) + assert again is not None and again == stored, "second capture must reuse the same entry" diff --git a/tests/test_queue_cancel.py b/tests/test_queue_cancel.py index bfa7edc87..98cc77a03 100644 --- a/tests/test_queue_cancel.py +++ b/tests/test_queue_cancel.py @@ -31,6 +31,8 @@ class _StubGitHub: class _StubSandbox: """Sentinel; queue tests don't touch the workspace pool.""" + natives_cache = None + class _StubGitTransport: """Sentinel; queue tests don't push.""" diff --git a/tests/test_queue_shutdown.py b/tests/test_queue_shutdown.py index 4ae891f70..87497e863 100644 --- a/tests/test_queue_shutdown.py +++ b/tests/test_queue_shutdown.py @@ -26,6 +26,8 @@ class _StubGitHub: class _StubSandbox: """Sentinel; queue tests don't touch the workspace pool.""" + natives_cache = None + class _StubGitTransport: """Sentinel; queue tests don't push.""" diff --git a/tests/test_sandbox.py b/tests/test_sandbox.py index 91e5f4234..504345304 100644 --- a/tests/test_sandbox.py +++ b/tests/test_sandbox.py @@ -1196,3 +1196,103 @@ def test_run_git_kills_hung_child(tmp_path: Path, monkeypatch: pytest.MonkeyPatc _run_git(["status"], cwd=tmp_path, token=None, timeout=0.5) assert exc.value.returncode == 124 assert "timed out" in exc.value.stderr.lower() + + +# --------------------------------------------------------------------------- +# NativesCache integration into ensure_workspace +# --------------------------------------------------------------------------- + + +def _seed_native_dir(repo_dir: Path) -> Path: + native_dir = repo_dir / "packages" / "natives" / "native" + native_dir.mkdir(parents=True, exist_ok=True) + return native_dir + + +def test_ensure_workspace_without_cache_leaves_native_dir_untouched(tmp_path: Path, upstream_repo: Path) -> None: + mgr = SandboxManager(tmp_path / "workspaces") + ws = mgr.ensure_workspace( + repo="octo/widget", + number=10, + title="no cache", + clone_url=str(upstream_repo), + default_branch="main", + author_name="robomp-bot", + author_email="robomp-bot@example.invalid", + ) + assert mgr.natives_cache is None + # No `packages/natives/native/` was tracked in the upstream, and no cache + # is configured → the directory wasn't created by populate. + assert not (ws.repo_dir / "packages" / "natives" / "native").exists() + + +def test_ensure_workspace_populates_from_natives_cache(tmp_path: Path, upstream_repo: Path) -> None: + from robomp.natives_cache import NativesCache, compute_key, target_triple + + cache = NativesCache(tmp_path / "natives-cache") + mgr = SandboxManager(tmp_path / "workspaces", natives_cache=cache) + + # First workspace: stage built artifacts, capture under the workspace's key. + ws1 = mgr.ensure_workspace( + repo="octo/widget", + number=11, + title="producer", + clone_url=str(upstream_repo), + default_branch="main", + author_name="robomp-bot", + author_email="robomp-bot@example.invalid", + ) + native_dir1 = _seed_native_dir(ws1.repo_dir) + # Mirror the napi build output set. The filename must match the live + # `target_triple()` value or the populate path won't recognize it. + triple = target_triple() + (native_dir1 / f"pi_natives.{triple}.node").write_bytes(b"ELFx") + (native_dir1 / "index.d.ts").write_text("export const X: number;\n") + (native_dir1 / "index.js").write_text("export const X = 1;\n") + (native_dir1 / "embedded-addon.js").write_text("export const embeddedAddon = null;\n") + key = compute_key(ws1.repo_dir) # default target = target_triple() + assert cache.capture("octo/widget", key, native_dir1) is not None + + # Second workspace on the same source HEAD: ensure_workspace auto-populates. + # We force the same key by pinning TARGET_VARIANT (only relevant on x64; + # harmless on arm64) — actually compute_key uses target_triple() at call + # time. To make the test platform-independent, override populate to use + # the same key explicitly. + ws2 = mgr.ensure_workspace( + repo="octo/widget", + number=12, + title="consumer", + clone_url=str(upstream_repo), + default_branch="main", + author_name="robomp-bot", + author_email="robomp-bot@example.invalid", + ) + native_dir2 = ws2.repo_dir / "packages" / "natives" / "native" + # The auto-populate path used the real target_triple() — which matches + # the host that just captured. So the same key applies and files appear. + assert native_dir2.is_dir(), "populate should have created native/ on hit" + node_name = f"pi_natives.{triple}.node" + assert (native_dir2 / node_name).read_bytes() == b"ELFx" + # The .node is hardlinked, sharing the cache's inode. + cached_node = cache.entry_dir("octo/widget", key) / node_name + ws2_node = native_dir2 / node_name + assert cached_node.stat().st_ino == ws2_node.stat().st_ino + + +def test_ensure_workspace_cache_miss_is_silent_noop(tmp_path: Path, upstream_repo: Path) -> None: + from robomp.natives_cache import NativesCache + + cache = NativesCache(tmp_path / "empty-cache") + mgr = SandboxManager(tmp_path / "workspaces", natives_cache=cache) + ws = mgr.ensure_workspace( + repo="octo/widget", + number=13, + title="miss", + clone_url=str(upstream_repo), + default_branch="main", + author_name="robomp-bot", + author_email="robomp-bot@example.invalid", + ) + # Cache is empty so the workspace ends up identical to the no-cache case. + assert ws.repo_dir.is_dir() + assert not (ws.repo_dir / "packages" / "natives" / "native").exists() diff --git a/tests/test_server.py b/tests/test_server.py index efc76791a..2e7d284e1 100644 --- a/tests/test_server.py +++ b/tests/test_server.py @@ -1514,6 +1514,8 @@ def test_webhook_maintainer_bypasses_rate_limit( class _RecordingSandbox: """Stand-in for SandboxManager: records calls, hands back a fake Workspace.""" + natives_cache = None + def __init__(self, tmp_root: Path) -> None: self.tmp_root = tmp_root self.ensure_calls: list[dict] = [] diff --git a/tests/test_worker.py b/tests/test_worker.py index 604d3beff..e5ca6848b 100644 --- a/tests/test_worker.py +++ b/tests/test_worker.py @@ -656,3 +656,115 @@ async def test_run_rpc_skips_reminder_when_unclassified(tmp_path: Path, settings loop.close() fake = _FakeRpcClient.instances[0] assert len(fake.prompts) == 1 + + +# --------------------------------------------------------------------------- +# Natives-cache capture-on-success +# --------------------------------------------------------------------------- + + +class _RecordingNativesCache: + """Test double for `NativesCache`: records `capture` calls, optionally + raises so we can verify exception swallowing.""" + + def __init__(self, *, raise_on_capture: bool = False) -> None: + self.capture_calls: list[tuple[str, str, Path]] = [] + self.raise_on_capture = raise_on_capture + + def capture(self, repo: str, key: str, native_dir: Path, **_kwargs) -> Path | None: + self.capture_calls.append((repo, key, native_dir)) + if self.raise_on_capture: + raise RuntimeError("simulated cache failure") + return native_dir + + +def _make_capture_inputs( + tmp_path: Path, + settings: Settings, + *, + cache: _RecordingNativesCache | None, + with_native_artifacts: bool, +) -> worker.TaskInputs: + """Build a `TaskInputs` whose workspace optionally has built natives.""" + inputs, _ = _make_inputs(tmp_path, settings, session_has_jsonl=False) + # Replace the SimpleNamespace workspace with one carrying the fields + # `_capture_natives_cache` needs (workspace_key + repo_full_name). + ws = SimpleNamespace( + root=inputs.workspace.root, + session_dir=inputs.workspace.session_dir, + repo_dir=inputs.workspace.repo_dir, + branch=inputs.workspace.branch, + workspace_key="acme__widgets__1", + repo_full_name="acme/widgets", + ) + if with_native_artifacts: + native_dir = ws.repo_dir / "packages" / "natives" / "native" + native_dir.mkdir(parents=True) + (native_dir / "pi_natives.linux-arm64.node").write_bytes(b"ELFx") + (native_dir / "index.d.ts").write_text("") + (native_dir / "index.js").write_text("") + (native_dir / "embedded-addon.js").write_text("") + return worker.TaskInputs( + settings=settings, + db=inputs.db, + github=inputs.github, + git_transport=inputs.git_transport, + repo=inputs.repo, + issue=inputs.issue, + workspace=ws, # type: ignore[arg-type] + delivery_id=inputs.delivery_id, + attempts=inputs.attempts, + slot_uid=inputs.slot_uid, + natives_cache=cache, # type: ignore[arg-type] + ) + + +def test_capture_natives_cache_no_op_without_cache(tmp_path: Path, settings: Settings) -> None: + inputs = _make_capture_inputs(tmp_path, settings, cache=None, with_native_artifacts=True) + # Just must not raise. + worker._capture_natives_cache(inputs) + + +def test_capture_natives_cache_skips_without_artifacts(tmp_path: Path, settings: Settings) -> None: + cache = _RecordingNativesCache() + inputs = _make_capture_inputs(tmp_path, settings, cache=cache, with_native_artifacts=False) + worker._capture_natives_cache(inputs) + # No artifacts → no key compute, no capture. + assert cache.capture_calls == [] + + +def test_capture_natives_cache_swallows_key_compute_failure( + tmp_path: Path, settings: Settings, monkeypatch: pytest.MonkeyPatch +) -> None: + cache = _RecordingNativesCache() + inputs = _make_capture_inputs(tmp_path, settings, cache=cache, with_native_artifacts=True) + # Repo dir is not a git repo → natives_compute_key raises. + # Already true for the SimpleNamespace workspace (repo_dir is plain tmp dir). + worker._capture_natives_cache(inputs) + assert cache.capture_calls == [] + + +def test_capture_natives_cache_swallows_capture_exception( + tmp_path: Path, settings: Settings, monkeypatch: pytest.MonkeyPatch +) -> None: + cache = _RecordingNativesCache(raise_on_capture=True) + inputs = _make_capture_inputs(tmp_path, settings, cache=cache, with_native_artifacts=True) + # Bypass git: stub the key compute so capture is reached. + monkeypatch.setattr(worker, "natives_compute_key", lambda _repo_dir: "deadbeef") + # Must not propagate the RuntimeError. + worker._capture_natives_cache(inputs) + assert len(cache.capture_calls) == 1 + + +def test_capture_natives_cache_records_on_success( + tmp_path: Path, settings: Settings, monkeypatch: pytest.MonkeyPatch +) -> None: + cache = _RecordingNativesCache() + inputs = _make_capture_inputs(tmp_path, settings, cache=cache, with_native_artifacts=True) + monkeypatch.setattr(worker, "natives_compute_key", lambda _repo_dir: "cafef00d") + worker._capture_natives_cache(inputs) + assert len(cache.capture_calls) == 1 + repo, key, native_dir = cache.capture_calls[0] + assert repo == "acme/widgets" + assert key == "cafef00d" + assert native_dir == inputs.workspace.repo_dir / "packages" / "natives" / "native"