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.
This commit is contained in:
can1357
2026-05-16 20:29:05 +02:00
parent dc2eb47297
commit 7f544fc669
16 changed files with 1501 additions and 12 deletions
+1 -1
View File
@@ -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 {} +
+12
View File
@@ -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:
+485
View File
@@ -0,0 +1,485 @@
"""Content-addressed cache of pre-built ``packages/natives/native/`` artifacts.
The napi-rs build of ``pi_natives.<platform>-<arch>[-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:
"""``<platform>-<arch>[-<variant>]`` 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:
# "<hash> <type> <size>" — 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 (".<key>.tmp.<pid>") 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",
]
+32
View File
@@ -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:
+100 -2
View File
@@ -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)
+14 -1
View File
@@ -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,
+6
View File
@@ -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)
+65 -8
View File
@@ -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,
},
)
+7
View File
@@ -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",
}
+414
View File
@@ -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
+147
View File
@@ -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"
+2
View File
@@ -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."""
+2
View File
@@ -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."""
+100
View File
@@ -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()
+2
View File
@@ -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] = []
+112
View File
@@ -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"