diff --git a/app-tests/git-leak/README.md b/app-tests/git-leak/README.md index a569975db..a2bab848c 100644 --- a/app-tests/git-leak/README.md +++ b/app-tests/git-leak/README.md @@ -51,6 +51,33 @@ Gate-coverage matrix (what each flagship test actually does): | `test_boot_loads_all_scopes` | **baseline → gate (PR4)** | PASSES with the loose default target; set `BOOT_TARGET_SECONDS` low (plan: 120 @ 50) on PR4 to gate the parallel-boot fix | | `test_repeat_sync_rss_stays_bounded` | **RSS guard** | PASSES; an RSS-budget guard against per-sync allocation leaks (the cache *count* can't grow for any impl, so there is no count assertion — see below) | | `test_server_recovers_after_postgres_bounce` | **guard (PER-15065 + gap publishes)** | PASSES on this branch (which has #915); guards the in-place broadcaster reconnect and that a scope PUT *during* the outage is buffered/replayed, not dropped | +| `test_delete_recreate_storm` | **guard (lock re-mint)** | PASSES — rapid delete/re-create of the same source serializes on the repo lock and ends with clean caches; guards 89e090be | +| `test_randomized_churn_holds_invariants` | **guard (seeded churn)** | PASSES — seeded random put/refresh/settled-delete churn holds invariants at every settle point (replay a failure with `CHURN_SEED=`); repoints and delete-vs-inflight-sync races are deliberately excluded, both lifted when PR3's fleet purge lands | +| `test_delete_during_hung_fetch_no_crash` | **guard (use-after-free)** | PASSES — deleting a scope whose clone is hung never crashes a worker; guards the use-after-free class 89e090be fixed | +| `test_delete_during_hung_fetch_returns_bounded` | **gate (PR3)** | FAILS without the PR3 fetch timeout — the purge waits on the repo lock and a hung clone holds that lock indefinitely, so the DELETE never returns in bounded time | +| `test_repoint_during_inflight_fetch_drains_old_source` | **gate (PR3, update path)** | Green half passes today — repointing while the old source's clone is hung still serves the new source; red half FAILS without PR3's update-path purge — the old source's cache entries never drain | +| `test_multiworker_churn_drains_every_worker` | **gate (PR3, broadcast)** | FAILS without PR3's broadcast purge — cache purges are process-local, so a worker whose caches were populated by something other than the DELETE it served (e.g. the leader's watcher syncs) leaks permanently; the HIGH finding from the PR2 review, as a gate | +| `test_warm_boot_reuses_clones` | **guard (S2)** | PASSES — a restart with intact clones must serve without re-cloning | +| `test_corrupt_clone_recovers_without_clone_loop` | **guard (S3/T7)** | PASSES — emptying a clone's object store in place while the server holds a warm cached handle is detected as invalid and recovers through the invalid-repo branch with exactly one re-clone, no serve-500s wedge and no re-clone loop; verifies the gutted-object-store detection fix | +| `test_orphan_clone_dir_is_reclaimed` | **gate (orphan sweep, unowned)** | FAILS — a clone dir with no live scope is never reclaimed; no orphan sweep exists yet (PR3+, currently unowned) | +| `test_redis_wiped_boot_reclaims_clones` | **gate (orphan sweep, unowned)** | FAILS — after a scope-store wipe, on-disk clones referencing nothing are never reclaimed; same missing-sweep class as the orphan-dir gate | +| `test_boot_with_unreachable_remotes_still_serves_healthy` | **gate (PR3) — watch this flip** | FAILS without the PR3 fetch timeout — unreachable remotes present at boot hang the preload/first-sync clones and starve the executor, so a healthy scope can't serve; boot-time cousin of the offline gate | +| `test_shard_reconfig_still_serves_but_orphans_old_clones` | **half-gate (S5, orphan sweep)** | Green half PASSES — serving survives a `SCOPES_REPO_CLONES_SHARDS` reconfig (re-clone under new ids); red half FAILS — the old-shard dirs are orphaned until the orphan sweep lands | +| `test_force_push_rewrite_recovers` | **characterization** | PASSES — a force-pushed (rewritten) head is picked up on refresh, pinning today's behavior (pygit2's forced default fetch refspec plus `set_target` moving the local ref) | +| `test_deleted_branch_keeps_serving_last_head` | **characterization** | PASSES — deleting the tracked branch upstream doesn't crash anything; fetch doesn't prune, so OPAL silently keeps serving the last known head (documented, not necessarily desirable, behavior) | + +## Invariants +`invariants.py` defines quiescence invariants I1-I6, checked at every `opal`-fixture +teardown (I1-I4 and I6 via `check_invariants`; I5 lives directly in `conftest.py`, +since it needs the worker pids captured at fixture setup): I1 no orphan clone dir on +disk, I2 no cached repo handle with no dir on disk, I3 no `repos_last_fetched` key +with no live scope, I4 no `repo_locks` key with no live scope, I5 the worker pid set +is unchanged across the test (proving no crash/respawn), I6 the scopes API still +answers 200. Two markers let a test opt out deliberately instead of by accident: +`@pytest.mark.invariant_exempt("I1", ...)` for red-gate tests that knowingly leave a +violation behind (named per test, not a blanket skip), and +`@pytest.mark.allow_worker_restart` for tests that restart `opal_server` on purpose, +which would otherwise trip I5. Notes on the guards: - `test_repeat_sync_rss_stays_bounded` — clone paths are keyed by the repo URL, diff --git a/app-tests/git-leak/conftest.py b/app-tests/git-leak/conftest.py index 547c22404..5e5e4ae61 100644 --- a/app-tests/git-leak/conftest.py +++ b/app-tests/git-leak/conftest.py @@ -10,6 +10,7 @@ list_seeded_repos, worker_pids, ) +from invariants import check_invariants def pytest_addoption(parser): @@ -67,7 +68,7 @@ def stack(request, repo_count): @pytest.fixture() -def opal(stack) -> OpalServerClient: +def opal(stack, request) -> OpalServerClient: # The compose stack is session-scoped (one server for the whole run), but # scopes must not leak between tests: clone paths are keyed by repo URL, so # a scope left behind by one test shares a cache entry with any later test @@ -84,8 +85,18 @@ def opal(stack) -> OpalServerClient: assert ( len(worker_pids()) == 1 ), f"expected a single-worker stack, found workers {sorted(worker_pids())}" + pids_before = worker_pids() yield stack stack.delete_all_scopes() + exempt_marker = request.node.get_closest_marker("invariant_exempt") + exempt = set(exempt_marker.args) if exempt_marker else set() + if request.node.get_closest_marker("allow_worker_restart") is None: + assert worker_pids() == pids_before, ( + f"I5: worker pid-set changed during the test " + f"({sorted(pids_before)} -> {sorted(worker_pids())}) — a crash/" + f"respawn resets the caches and can fake a clean drain" + ) + check_invariants(stack, exempt=exempt) @pytest.fixture() diff --git a/app-tests/git-leak/docker-compose.yml b/app-tests/git-leak/docker-compose.yml index 66b5d7a59..3dcc2bc42 100644 --- a/app-tests/git-leak/docker-compose.yml +++ b/app-tests/git-leak/docker-compose.yml @@ -114,6 +114,10 @@ services: # (references/debug-pubsub.md §3-4) — a single worker can't tell #915's # in-place broadcaster reconnect from a plain worker respawn. UVICORN_NUM_WORKERS: "${OPAL_TEST_WORKERS:-1}" + # Shard-count knob for the reshard boot test (test_boot_states.py): + # changing SCOPES_REPO_CLONES_SHARDS between boots moves every source_id, + # orphaning all previous clone dirs — a real ops event worth gating. + OPAL_SCOPES_REPO_CLONES_SHARDS: "${OPAL_TEST_SHARDS:-1}" OPAL_SCOPES: "1" OPAL_REDIS_URL: "redis://redis:6379" OPAL_BROADCAST_URI: "postgres://opal:opal@postgres:5432/opal" diff --git a/app-tests/git-leak/helpers.py b/app-tests/git-leak/helpers.py index 64a32abbe..31972b69f 100644 --- a/app-tests/git-leak/helpers.py +++ b/app-tests/git-leak/helpers.py @@ -2,7 +2,7 @@ import subprocess import time from pathlib import Path -from typing import Dict, List +from typing import Any, Dict, List import requests @@ -50,7 +50,7 @@ def wait_healthy(self, timeout: int = 180) -> None: time.sleep(2) raise RuntimeError(f"opal-server not healthy in {timeout}s (last: {last})") - def stats(self, samples: int = 3, interval: float = 0.1) -> Dict[str, int]: + def stats(self, samples: int = 3, interval: float = 0.1) -> Dict[str, Any]: """Read the git-fetcher cache stats, merged across a few reads. The stack runs a single uvicorn worker (see docker-compose.yml), so the @@ -59,15 +59,29 @@ def stats(self, samples: int = 3, interval: float = 0.1) -> Dict[str, int]: the ``max`` per key only smooths over a read that races an in-flight mutation; it is not relied on to paper over multi-worker nondeterminism (which the single-worker setup removes outright). + + Merge semantics: numeric counts take the max across samples (peak + smoothing); ``pid`` and the ``*_keys`` lists are last-wins. The + count fields and their paired ``*_keys`` lists are therefore only + guaranteed mutually consistent at ``samples=1`` — which is what + every consistency-sensitive consumer (the invariant checker, + per-pid sampling) uses. Do not assert + ``len(stats["repos_keys"]) == stats["repos"]`` on a multi-sample + merge. """ - merged: Dict[str, int] = {} + merged: Dict[str, Any] = {} for i in range(max(1, samples)): resp = requests.get( f"{self.base_url}/internal/git-fetcher-cache-stats", timeout=10 ) resp.raise_for_status() for key, value in resp.json().items(): - merged[key] = max(merged.get(key, 0), value) + if key != "pid" and isinstance(value, (int, float)): + merged[key] = max(merged.get(key, 0), value) + else: + # pid and the *_keys lists: last-wins (single-worker stack + # makes every read hit the same worker anyway) + merged[key] = value if i < samples - 1: time.sleep(interval) return merged @@ -234,7 +248,12 @@ def create_repo(self, name: str) -> None: return resp = requests.post( f"{self.base_url}/api/v1/user/repos", - json={"name": name, "private": False, "auto_init": True}, + json={ + "name": name, + "private": False, + "auto_init": True, + "default_branch": "main", + }, auth=self._auth, timeout=10, ) @@ -250,6 +269,39 @@ def delete_repo(self, name: str) -> None: resp.raise_for_status() +class RepoMutator: + """Host-side git mutations against a bed Gitea repo (force-push, branch + ops) — the remote-transition tests' hands. + + Uses GitPython over the published host port with admin basic-auth. + """ + + def __init__(self, name: str, workdir: Path): + import git as gitpython + + self._url = ( + f"http://{GITEA_USER}:{GITEA_PASSWORD}@localhost:13000/" + f"{GITEA_USER}/{name}.git" + ) + self._clone = gitpython.Repo.clone_from(self._url, str(workdir / name)) + with self._clone.config_writer() as cw: + cw.set_value("user", "name", "bed-mutator") + cw.set_value("user", "email", "bed@test.local") + + def force_push_rewrite(self, branch: str = "main") -> None: + self._clone.git.checkout(branch) + self._clone.git.commit("--amend", "--allow-empty", "-m", "rewritten history") + self._clone.git.push("--force", "origin", branch) + + def push_new_branch(self, branch: str) -> None: + self._clone.git.checkout("-b", branch) + self._clone.git.commit("--allow-empty", "-m", f"seed {branch}") + self._clone.git.push("origin", branch) + + def delete_remote_branch(self, branch: str) -> None: + self._clone.git.push("origin", f":{branch}") + + def gitea_repo_url(name: str) -> str: # url reachable from inside the opal_server container return f"{GITEA_INTERNAL_URL}/{GITEA_USER}/{name}.git" @@ -396,6 +448,47 @@ def bounce_postgres(down_seconds: int = 5, during=None) -> None: compose("up", "-d", "--wait", "--no-recreate", "postgres") +def wait_until(predicate, timeout: float, interval: float = 2.0) -> bool: + """Poll ``predicate()`` until truthy or ``timeout`` elapses. + + Swallows transient exceptions from the predicate (a stats read + racing a restart is not a verdict) — only the final state decides. + """ + deadline = time.time() + timeout + while time.time() < deadline: + try: + if predicate(): + return True + except Exception: + pass + time.sleep(interval) + try: + return bool(predicate()) + except Exception: + return False + + +def stats_by_pid(opal, min_pids: int = 2, attempts: int = 200, interval: float = 0.05): + """Sample the stats endpoint repeatedly; keep the LATEST snapshot per pid. + + Requests land on arbitrary workers, so repeated single-sample reads + eventually observe each worker. Returns {pid: latest_stats} once min_pids + distinct pids are seen (plus one grace sample for freshness), or after + `attempts` samples; the caller decides whether the count is sufficient. + """ + seen = {} + grace_sweep_done = False + for _ in range(attempts): + snap = opal.stats(samples=1) + seen[snap["pid"]] = snap + if len(seen) >= min_pids: + if grace_sweep_done: + break + grace_sweep_done = True + time.sleep(interval) + return seen + + def list_seeded_repos(count: int) -> List[str]: return [f"policy-repo-{i:04d}" for i in range(count)] diff --git a/app-tests/git-leak/invariants.py b/app-tests/git-leak/invariants.py new file mode 100644 index 000000000..f8163d241 --- /dev/null +++ b/app-tests/git-leak/invariants.py @@ -0,0 +1,89 @@ +"""Quiescence invariants I1-I6 for the git-leak bed. + +Asserted at every test's teardown (see the ``opal`` fixture) after +``delete_all_scopes()``, and callable mid-test. Red-gate tests that +knowingly leave violations exempt specific IDs via +``@pytest.mark.invariant_exempt("I1", ...)``. + +Deviation from the design spec: I1's "none missing" direction is NOT +asserted — scopes clone lazily, so a just-created scope legitimately has +no dir yet. Only the orphan direction (dir with no live scope) signals a +leak. +""" +import hashlib + +import requests +from helpers import compose + +BASE_DIR_IN_CONTAINER = "/opal/git_sources" + + +def source_id(url: str, branch: str = "main", shards: int = 1) -> str: + """Host-side mirror of GitPolicyFetcher.source_id (sha256(url) + shard).""" + base = hashlib.sha256(url.encode("utf-8")).hexdigest() + index = hashlib.sha256(branch.encode("utf-8")).digest()[0] % shards + return f"{base}-{index}" + + +def clone_dirs(service: str = "opal_server") -> set: + out = compose( + "exec", + "-T", + service, + "sh", + "-c", + f"ls -1 {BASE_DIR_IN_CONTAINER} 2>/dev/null || true", + ).stdout + return {line.strip() for line in out.splitlines() if line.strip()} + + +def live_source_ids(opal, shards: int = 1) -> set: + resp = requests.get(f"{opal.base_url}/scopes", timeout=30) + resp.raise_for_status() + return { + source_id(s["policy"]["url"], s["policy"].get("branch", "main"), shards) + for s in resp.json() + } + + +def check_invariants(opal, exempt=frozenset(), shards: int = 1) -> None: + exempt = set(exempt) + failures = [] + live = live_source_ids(opal, shards) + disk = clone_dirs() + stats = opal.stats(samples=1) + + if "I1" not in exempt: + orphans = disk - live + if orphans: + failures.append( + f"I1: {len(orphans)} orphan clone dir(s) on disk, e.g. {sorted(orphans)[:3]}" + ) + if "I2" not in exempt: + cached_dirs = {k.rsplit("/", 1)[-1] for k in stats.get("repos_keys", [])} + stray = cached_dirs - disk + if stray: + failures.append( + f"I2: cached repo handle(s) with no dir on disk: {sorted(stray)[:3]}" + ) + if "I3" not in exempt: + stray = set(stats.get("repos_last_fetched_keys", [])) - live + if stray: + failures.append( + f"I3: repos_last_fetched key(s) with no live scope: {sorted(stray)[:3]}" + ) + if "I4" not in exempt: + stray = set(stats.get("repo_locks_keys", [])) - live + if stray: + failures.append( + f"I4: repo_locks key(s) with no live scope: {sorted(stray)[:3]}" + ) + if "I6" not in exempt: + if requests.get(f"{opal.base_url}/scopes", timeout=10).status_code != 200: + failures.append("I6: scopes API not answering 200") + + assert not failures, ( + "invariant violations:\n " + + "\n ".join(failures) + + f"\nstats={stats}\ndisk(sample)={sorted(disk)[:10]}" + ) diff --git a/app-tests/git-leak/pytest.ini b/app-tests/git-leak/pytest.ini index 0d501534f..df0c03a10 100644 --- a/app-tests/git-leak/pytest.ini +++ b/app-tests/git-leak/pytest.ini @@ -9,3 +9,7 @@ # marker makes that guarantee explicit — it pins the rootdir to this directory # when run in place, independent of the root config — while the root run (from # the repo root) still never descends here. + +markers = + invariant_exempt(ids): exempt named invariants (I1..I6) from the opal fixture's teardown check for this (red-gate) test + allow_worker_restart: this test restarts opal_server on purpose; skip the I5 pid-stability teardown check diff --git a/app-tests/git-leak/test_boot_states.py b/app-tests/git-leak/test_boot_states.py new file mode 100644 index 000000000..99e8548e3 --- /dev/null +++ b/app-tests/git-leak/test_boot_states.py @@ -0,0 +1,247 @@ +"""Start-state tests: what the server boots INTO (S2-S7 in the design spec). + +The cold-empty start (S1) is covered by test_boot.py. +""" +import time + +import pytest +from helpers import ( + HEALTHY_PROBE_REPO, + compose, + gitea_repo_url, + list_seeded_repos, + make_repo_unreachable, + wait_until, +) +from invariants import clone_dirs + + +def _clone_log_count() -> int: + return compose("logs", "--no-log-prefix", "opal_server").stdout.count( + "Cloning repo" + ) + + +@pytest.mark.timeout(900) +@pytest.mark.allow_worker_restart +def test_warm_boot_reuses_clones(opal, repo_count): + """A restart with intact clones must serve without re-cloning (S2).""" + n = min(repo_count, 10) + for i, repo in enumerate(list_seeded_repos(n)): + opal.put_scope(f"warm-{i}", gitea_repo_url(repo)) + for i in range(n): + assert wait_until( + lambda i=i: opal.get_scope_policy(f"warm-{i}").status_code == 200, + timeout=600, + ), f"warm-{i} never served before the restart" + + clones_before = _clone_log_count() + compose("restart", "opal_server") + opal.wait_healthy() + + for i in range(n): + assert wait_until( + lambda i=i: opal.get_scope_policy(f"warm-{i}").status_code == 200, + timeout=300, + ), f"warm-{i} not served after warm restart" + assert ( + _clone_log_count() == clones_before + ), "warm boot re-cloned instead of reusing the on-disk clones" + + +@pytest.mark.timeout(900) +def test_corrupt_clone_recovers_without_clone_loop(opal): + """S3/T7: emptying a clone's object store in place (objects/ kept as an + empty dir, so pygit2.discover_repository still finds the repo) while the + server holds a warm cached handle must be detected as invalid (gutted + object store: refs intact, head object unreadable) and recover through + the invalid-repo branch with exactly one re-clone — no serve-500s- + forever wedge and no re-clone loop.""" + repo = list_seeded_repos(3)[2] + opal.put_scope("dirty", gitea_repo_url(repo)) + assert wait_until( + lambda: opal.get_scope_policy("dirty").status_code == 200, timeout=300 + ) + + invalid_before = compose("logs", "--no-log-prefix", "opal_server").stdout.count( + "Deleting invalid repo" + ) + + # empty the object store's CONTENTS in place, keeping the objects/ dir + # itself so discover_repository still resolves the repo; the server keeps + # its cached pygit2 handle (warm cache) — the gutted-object-store case + # (deleting the objects/ node instead would route recovery through the + # repo-not-found -> _clone() branch and never touch the cached handle) + compose( + "exec", + "-T", + "opal_server", + "sh", + "-c", + 'for d in /opal/git_sources/*/; do rm -rf "$d/.git/objects"/*; done', + ) + opal.refresh_all() + + assert wait_until( + lambda: opal.get_scope_policy("dirty").status_code == 200, timeout=300 + ), "scope never recovered after clone corruption" + + invalid_after = compose("logs", "--no-log-prefix", "opal_server").stdout.count( + "Deleting invalid repo" + ) + assert invalid_after > invalid_before, ( + "recovery did not take the invalid-repo (warm cached handle) branch — " + "the corruption failed to exercise the Bug A path" + ) + + clones_1 = _clone_log_count() + opal.refresh_all() + time.sleep(10) + clones_2 = _clone_log_count() + assert clones_2 == clones_1, ( + f"re-clone loop: clone count kept growing after recovery " + f"({clones_1} -> {clones_2})" + ) + + +@pytest.mark.timeout(900) +@pytest.mark.invariant_exempt("I1") +def test_orphan_clone_dir_is_reclaimed(opal): + """RED until an orphan sweep exists (PR3+, currently unowned): a clone dir + with no live scope must eventually be removed.""" + fake_sid = "f" * 64 + "-0" + compose( + "exec", + "-T", + "opal_server", + "sh", + "-c", + f"mkdir -p /opal/git_sources/{fake_sid} && touch /opal/git_sources/{fake_sid}/junk", + ) + try: + opal.refresh_all() + assert wait_until( + lambda: fake_sid not in clone_dirs(), timeout=60 + ), "orphan clone dir never reclaimed (needs an orphan sweep)" + finally: + # red gate leaves state on purpose; clean it so later tests' I1 holds + compose( + "exec", + "-T", + "opal_server", + "sh", + "-c", + f"rm -rf /opal/git_sources/{fake_sid}", + ) + + +@pytest.mark.timeout(900) +@pytest.mark.allow_worker_restart +@pytest.mark.invariant_exempt("I1") +def test_redis_wiped_boot_reclaims_clones(opal): + """RED until the orphan sweep (same class as the orphan-dir gate): after a + scope-store wipe, on-disk clones reference nothing and must be + reclaimed.""" + opal.put_scope("wipe-0", gitea_repo_url(list_seeded_repos(1)[0])) + assert wait_until( + lambda: opal.get_scope_policy("wipe-0").status_code == 200, timeout=300 + ), "wipe-0 never served before the wipe" + try: + compose("stop", "opal_server") + compose("exec", "-T", "redis", "redis-cli", "FLUSHALL") + compose("start", "opal_server") + opal.wait_healthy() + assert wait_until( + lambda: clone_dirs() == set(), timeout=60 + ), f"clones of the wiped scope store never reclaimed: {sorted(clone_dirs())[:5]}" + finally: + opal.hard_reset() + # hard_reset never touches git_sources/, and post-FLUSHALL no scope + # record exists to route a DELETE's rmtree at the leftover clone — + # remove it explicitly so later tests' I1 stays meaningful. The scope + # store is empty here, so a blanket clean is safe. + compose( + "exec", + "-T", + "opal_server", + "sh", + "-c", + "rm -rf /opal/git_sources/*", + ) + + +@pytest.mark.timeout(1200) +@pytest.mark.allow_worker_restart +@pytest.mark.invariant_exempt("I1", "I3", "I4") +def test_boot_with_unreachable_remotes_still_serves_healthy(opal): + """RED until PR3 (fetch timeout) — WATCH THIS FLIP when PR3 merges. + + Boot-time cousin of the offline gate: unreachable remotes present at boot + hang the preload/first-sync clones and starve the executor, so a healthy + scope can't serve. + """ + for i in range(10): + opal.put_scope(f"down-{i}", make_repo_unreachable(f"down-{i}-repo")) + opal.put_scope("boot-healthy", gitea_repo_url(HEALTHY_PROBE_REPO)) + try: + compose("restart", "opal_server") + opal.wait_healthy(timeout=300) + assert wait_until( + lambda: opal.get_scope_policy("boot-healthy").status_code == 200, + timeout=300, + ), "healthy scope starved at boot by unreachable remotes (PR3 gate)" + finally: + opal.hard_reset() + # hard_reset flushes Redis but never touches git_sources/ — the + # blackhole scopes' partial clone dirs would poison later tests' I1 + # checks. The scope store is empty post-reset, so a blanket clean is + # safe. (Same rationale as test_redis_wiped_boot_reclaims_clones.) + compose( + "exec", + "-T", + "opal_server", + "sh", + "-c", + "rm -rf /opal/git_sources/*", + ) + + +@pytest.mark.timeout(1200) +@pytest.mark.allow_worker_restart +@pytest.mark.invariant_exempt("I1") +def test_shard_reconfig_still_serves_but_orphans_old_clones(opal, tmp_path): + """S5: SCOPES_REPO_CLONES_SHARDS reconfig moves every source_id. + GREEN half: serving must survive the reshard (re-clone under new ids). + RED half (until the orphan sweep): the old-shard dirs are orphaned.""" + import os + + from invariants import live_source_ids + + opal.put_scope("shard-0", gitea_repo_url(list_seeded_repos(1)[0])) + assert wait_until( + lambda: opal.get_scope_policy("shard-0").status_code == 200, timeout=300 + ) + # preserve the shards=1 clones across the recreate (recreate wipes the fs) + compose("cp", "opal_server:/opal/git_sources", str(tmp_path / "saved")) + + os.environ["OPAL_TEST_SHARDS"] = "4" + try: + compose("up", "-d", "--no-deps", "--force-recreate", "opal_server") + opal.wait_healthy() + # restore the old-shard dirs next to whatever the new boot creates + compose("cp", str(tmp_path / "saved") + "/.", "opal_server:/opal/git_sources") + opal.refresh_all() + assert wait_until( + lambda: opal.get_scope_policy("shard-0").status_code == 200, + timeout=300, + ), "scope stopped serving after the shard reconfig (green half broken!)" + assert wait_until( + lambda: clone_dirs() <= live_source_ids(opal, shards=4), timeout=60 + ), ( + "old-shard clone dirs orphaned after reshard (red half — orphan " + f"sweep gate): {sorted(clone_dirs())[:5]}" + ) + finally: + os.environ["OPAL_TEST_SHARDS"] = "1" + compose("up", "-d", "--no-deps", "--force-recreate", "opal_server") + opal.wait_healthy() diff --git a/app-tests/git-leak/test_leak.py b/app-tests/git-leak/test_leak.py index 971823cdd..61137392c 100644 --- a/app-tests/git-leak/test_leak.py +++ b/app-tests/git-leak/test_leak.py @@ -178,6 +178,7 @@ def test_shared_repo_survives_sibling_scope_delete(opal, repo_count): @pytest.mark.timeout(900) +@pytest.mark.invariant_exempt("I1", "I3", "I4") def test_scope_repoint_releases_old_repo_cache(opal, repo_count): """Re-pointing a scope to a new repo URL, then deleting it, must drain ALL cache entries — including the orphaned old URL's. diff --git a/app-tests/git-leak/test_remote_transitions.py b/app-tests/git-leak/test_remote_transitions.py new file mode 100644 index 000000000..9c1ae80ed --- /dev/null +++ b/app-tests/git-leak/test_remote_transitions.py @@ -0,0 +1,60 @@ +"""Remote-side transitions (T16): the tenant rewrites or deletes what we track. + +Both are CHARACTERIZATION tests — they pin today's behavior. If one fails, +the behavior changed: decide deliberately, don't just update the assert. +""" +import pytest +from helpers import GiteaAdmin, RepoMutator, gitea_repo_url, wait_until + + +@pytest.mark.timeout(900) +def test_force_push_rewrite_recovers(opal, gitea_admin, tmp_path): + """A force-pushed (rewritten) head must be picked up on refresh: pygit2's + default fetch refspec is forced, and set_target just moves the local + ref.""" + gitea_admin.create_repo("mutation-force-push") + try: + opal.put_scope("fp", gitea_repo_url("mutation-force-push")) + assert wait_until( + lambda: opal.get_scope_policy("fp").status_code == 200, timeout=300 + ) + old_bundle = opal.get_scope_policy("fp").json() + + RepoMutator("mutation-force-push", tmp_path).force_push_rewrite() + opal.refresh_all() + + assert wait_until( + lambda: opal.get_scope_policy("fp").status_code == 200 + and opal.get_scope_policy("fp").json().get("hash") + != old_bundle.get("hash"), + timeout=300, + ), "rewritten head never served after refresh" + finally: + opal.delete_scope("fp") + gitea_admin.delete_repo("mutation-force-push") + + +@pytest.mark.timeout(900) +def test_deleted_branch_keeps_serving_last_head(opal, gitea_admin, tmp_path): + """Deleting the tracked branch upstream must not crash anything; fetch + doesn't prune, so OPAL silently keeps serving the last known head — + documented (not necessarily desirable) behavior.""" + gitea_admin.create_repo("mutation-branch-del") + mut = RepoMutator("mutation-branch-del", tmp_path) + mut.push_new_branch("extra") + try: + opal.put_scope("bd", gitea_repo_url("mutation-branch-del"), branch="extra") + assert wait_until( + lambda: opal.get_scope_policy("bd").status_code == 200, timeout=300 + ) + + mut.delete_remote_branch("extra") + opal.refresh_all() + + # settle, then pin: still healthy, still serving the last head + assert wait_until( + lambda: opal.get_scope_policy("bd").status_code == 200, timeout=120 + ), "scope stopped serving after upstream branch deletion" + finally: + opal.delete_scope("bd") + gitea_admin.delete_repo("mutation-branch-del") diff --git a/app-tests/git-leak/test_resilience.py b/app-tests/git-leak/test_resilience.py index d2df056f0..401fa4aed 100644 --- a/app-tests/git-leak/test_resilience.py +++ b/app-tests/git-leak/test_resilience.py @@ -41,6 +41,8 @@ def _wait_served(opal, scope_id: str, timeout: int): @pytest.mark.timeout(420) +@pytest.mark.allow_worker_restart +@pytest.mark.invariant_exempt("I1") def test_offline_repo_does_not_block_healthy_scopes(opal, repo_count): """Unreachable repos must not stop a healthy scope from serving. diff --git a/app-tests/git-leak/test_transitions.py b/app-tests/git-leak/test_transitions.py new file mode 100644 index 000000000..71ba7be61 --- /dev/null +++ b/app-tests/git-leak/test_transitions.py @@ -0,0 +1,222 @@ +"""Transition tests: interleavings the end-state gates never drive. + +See docs/superpowers/specs/2026-07-13-scope-fetcher-lifecycle-tests-design.md. +""" +import time + +import pytest +from helpers import ( + compose, + gitea_repo_url, + list_seeded_repos, + make_repo_unreachable, + wait_until, + worker_pids, +) +from invariants import clone_dirs, live_source_ids, source_id + + +@pytest.mark.timeout(900) +def test_delete_recreate_storm(opal, repo_count): + """Rapid delete/re-create of the same source must serialize on the repo + lock (lock re-mint path) and end with clean caches. + + Guards 89e090be. + """ + url = gitea_repo_url(list_seeded_repos(1)[0]) + for i in range(5): + opal.put_scope("storm", url) + assert wait_until( + lambda: opal.get_scope_policy("storm").status_code == 200, timeout=300 + ), f"round {i}: recreated scope never served" + opal.delete_scope("storm") + assert wait_until( + lambda: opal.stats(samples=1)["repo_locks"] == 0, timeout=60 + ), f"caches did not drain after the storm: {opal.stats()}" + + +@pytest.mark.timeout(1200) +def test_randomized_churn_holds_invariants(opal, repo_count): + """Seeded random put/refresh churn with settled deletes; invariants must + hold at every settle point. Replay a failure with CHURN_SEED=. + + Two deliberate constraints, both lifted when PR3's fleet purge lands + (no silent caps): + - 'repoint' ops are EXCLUDED: a repoint orphans the old source's cache + entries by design today (the red repoint gate covers it). A `put` on a + live scope therefore reuses that scope's existing repo. + - Deletes run only at round END, after every live scope has settled: a + DELETE racing an in-flight sync loses its purge (the sync re-populates + the caches for the dead scope — same PR3 class; proven deterministically + by seed 309006536 during this test's development; deterministic red-gate + coverage lands in test_repoint_during_inflight_fetch_drains_old_source). + A bounded residual window remains (a served scope's re-sync can still be + in flight); in practice the settle polling latency dwarfs a tiny repo's + sync time. + """ + import os + import random + + seed = int(os.environ.get("CHURN_SEED", "0")) or random.randrange(1, 2**31) + print(f"\nCHURN_SEED={seed}") + rng = random.Random(seed) + repos = list_seeded_repos(min(repo_count, 6)) + live = {} + + for round_no in range(4): + # burst: puts + refreshes only (deletes deferred to round end) + for _ in range(10): + op = rng.choice(["put", "refresh"]) + sid_ = f"rand-{rng.randrange(3)}" + if op == "put": + repo = live.get(sid_) or rng.choice(repos) + opal.put_scope(sid_, gitea_repo_url(repo)) + live[sid_] = repo + else: + opal.refresh_all() + # settle every live scope before any delete + for sid_ in list(live): + assert wait_until( + lambda s=sid_: opal.get_scope_policy(s).status_code == 200, + timeout=300, + ), f"round {round_no}: live scope {sid_} never settled (seed {seed})" + # settled deletes: each live scope has a coin-flip chance to go + for sid_ in list(live): + if rng.random() < 0.5: + opal.delete_scope(sid_) + live.pop(sid_) + opal.delete_scope(f"ghost-{round_no}") # delete-missing stays a 204 no-op + drained = wait_until( + lambda: opal.stats(samples=1)["repo_locks"] + <= len({r for r in live.values()}), + timeout=120, + ) + assert drained, ( + f"round {round_no}: locks exceed live sources (seed {seed}): " + f"{opal.stats()}" + ) + + +@pytest.mark.timeout(900) +@pytest.mark.allow_worker_restart +@pytest.mark.invariant_exempt("I1", "I3", "I4") +def test_delete_during_hung_fetch_no_crash(opal): + """Deleting a scope whose clone is hung must never crash a worker (the use- + after-free class 89e090be fixed). + + GREEN since that fix. + """ + import requests as _requests + + pids = worker_pids() + opal.put_scope("hung-nc", make_repo_unreachable("hung-nc-repo")) + time.sleep(5) # let the leader's clone start and block on the blackhole + try: + try: + opal.delete_scope("hung-nc") + except _requests.RequestException: + pass # a slow/blocked DELETE is the OTHER test's concern + assert ( + worker_pids() == pids + ), "worker crashed/respawned during delete-vs-hung-fetch" + finally: + opal.hard_reset() + + +@pytest.mark.timeout(900) +@pytest.mark.allow_worker_restart +@pytest.mark.invariant_exempt("I1", "I3", "I4") +def test_delete_during_hung_fetch_returns_bounded(opal): + """RED until PR3 (fetch timeout): the purge waits on the repo lock, and a + hung clone holds that lock indefinitely, so the DELETE hangs with it.""" + import requests as _requests + + opal.put_scope("hung-b", make_repo_unreachable("hung-b-repo")) + time.sleep(5) + try: + start = time.time() + resp = _requests.delete(f"{opal.base_url}/scopes/hung-b", timeout=90) + assert ( + resp.status_code in (200, 204) and time.time() - start < 90 + ), "DELETE of a hung-fetch scope did not return in bounded time" + finally: + opal.hard_reset() + + +@pytest.mark.timeout(900) +@pytest.mark.allow_worker_restart +@pytest.mark.invariant_exempt("I1", "I3", "I4") +def test_repoint_during_inflight_fetch_drains_old_source(opal, repo_count): + """RED until PR3 (update-path purge). + + Repointing a scope while its old source's clone is hung must still + serve the new source (green half) and eventually drop the old + source's cache entries (red half). + """ + repo_b = list_seeded_repos(2)[1] + old_url = make_repo_unreachable("repoint-hang-repo") + opal.put_scope("rp", old_url) + time.sleep(5) # old source's clone is now in flight, holding its lock + try: + opal.put_scope("rp", gitea_repo_url(repo_b)) # repoint while hung + assert wait_until( + lambda: opal.get_scope_policy("rp").status_code == 200, timeout=300 + ), "repointed scope never served its new source" + + old_sid = source_id(old_url) + + def _old_entries_gone(): + s = opal.stats(samples=1) + return ( + old_sid not in s["repo_locks_keys"] + and old_sid not in s["repos_last_fetched_keys"] + ) + + assert wait_until(_old_entries_gone, timeout=60), ( + f"old source {old_sid[:12]}… cache entries leaked after repoint " + f"(PR3 update-path purge gate): {opal.stats()}" + ) + finally: + opal.hard_reset() + + +@pytest.mark.timeout(1200) +def test_multiworker_churn_drains_every_worker(opal_multiworker, repo_count): + """RED until PR3 (broadcast purge): cache purges are process-local, so any + worker whose caches were populated by something other than the DELETE it + serves leaks permanently. + + Who populates what: the LEADER accumulates handles/locks via its watcher's + syncs (scopes/task.py); ANY worker additionally caches a pygit2 handle when + it serves a policy bundle (make_bundle -> _get_current_branch_head -> + _get_repo). The purge runs only on whichever worker happens to serve the + DELETE — every accumulation on a different worker outlives the scope. In + this test only the leader accumulates (nothing GETs bundles), so the + leader's retained entries are the observable leak. The HIGH finding from + the PR2 review, as a gate. + """ + from helpers import stats_by_pid + + opal = opal_multiworker + n = min(repo_count, 10) + for i, repo in enumerate(list_seeded_repos(n)): + opal.put_scope(f"mw-{i}", gitea_repo_url(repo)) + assert wait_until( + lambda: any(s["repos"] >= 1 for s in stats_by_pid(opal, attempts=40).values()), + timeout=600, + ), "no worker ever populated its repo cache" + + for i in range(n): + opal.delete_scope(f"mw-{i}") + + def _every_worker_drained(): + snaps = stats_by_pid(opal, min_pids=2, attempts=60) + return len(snaps) >= 2 and all( + s["repo_locks"] == 0 and s["repos"] == 0 and s["repos_last_fetched"] == 0 + for s in snaps.values() + ) + + assert wait_until(_every_worker_drained, timeout=120), ( + "a worker kept caches the churn's DELETEs never reached (PR3 gate): " + f"{ {p: {k: v for k, v in s.items() if isinstance(v, int)} for p, s in stats_by_pid(opal).items()} }" + ) diff --git a/packages/opal-server/opal_server/debug_stats.py b/packages/opal-server/opal_server/debug_stats.py index 5439ce950..249548d55 100644 --- a/packages/opal-server/opal_server/debug_stats.py +++ b/packages/opal-server/opal_server/debug_stats.py @@ -3,6 +3,7 @@ Used only by the off-by-default /internal stats endpoint so tests can observe the cache growth that the memory-leak fix (PR2) eliminates. """ +import os from pathlib import Path from typing import Dict, List, Optional @@ -21,13 +22,28 @@ def _read_rss_kb() -> int: return 0 -def git_fetcher_cache_stats() -> Dict[str, int]: - """Sizes of the three process-global GitPolicyFetcher caches + RSS.""" +def git_fetcher_cache_stats() -> Dict: + """Sizes + keys of the three process-global GitPolicyFetcher caches, RSS, + and the worker pid (per-process caches: the pid identifies WHICH worker + answered, so multi-worker bed tests can assert per-worker drain).""" + # Snapshot each cache once (dict.copy() is a single C-level operation, + # atomic under the GIL): this handler runs on a Starlette worker thread + # while the caches are mutated on the event-loop/executor threads, so + # iterating the live dicts can raise "dictionary changed size during + # iteration", and reading len() and keys() separately can return a + # self-contradictory count/keys pair. + repo_locks = GitPolicyFetcher.repo_locks.copy() + repos = GitPolicyFetcher.repos.copy() + repos_last_fetched = GitPolicyFetcher.repos_last_fetched.copy() return { - "repo_locks": len(GitPolicyFetcher.repo_locks), - "repos": len(GitPolicyFetcher.repos), - "repos_last_fetched": len(GitPolicyFetcher.repos_last_fetched), + "pid": os.getpid(), + "repo_locks": len(repo_locks), + "repos": len(repos), + "repos_last_fetched": len(repos_last_fetched), "rss_kb": _read_rss_kb(), + "repo_locks_keys": sorted(repo_locks.keys()), + "repos_keys": sorted(repos.keys()), + "repos_last_fetched_keys": sorted(repos_last_fetched.keys()), } @@ -57,5 +73,5 @@ def register_internal_stats_route( include_in_schema=False, dependencies=dependencies or [], ) - def _git_fetcher_cache_stats() -> Dict[str, int]: + def _git_fetcher_cache_stats() -> Dict: return git_fetcher_cache_stats() diff --git a/packages/opal-server/opal_server/git_fetcher.py b/packages/opal-server/opal_server/git_fetcher.py index 923bc423e..ad4a391eb 100644 --- a/packages/opal-server/opal_server/git_fetcher.py +++ b/packages/opal-server/opal_server/git_fetcher.py @@ -4,6 +4,7 @@ import hashlib import shutil import time +from contextlib import asynccontextmanager from pathlib import Path from typing import Optional, cast @@ -140,13 +141,24 @@ def __init__( f"Initializing git fetcher: scope_id={scope_id}, url={redact_url(source.url)}, branch={self._source.branch}, source_id={self._source_id}" ) - async def _get_repo_lock(self): - # Previous file based implementation worked across multiple processes/threads, but wasn't fair (next acquiree is random) - # This implementation works only within the same process/thread, but is fair (next acquiree is the earliest to enter the lock) - lock = GitPolicyFetcher.repo_locks[ - self._source_id - ] = GitPolicyFetcher.repo_locks.get(self._source_id, asyncio.Lock()) - return lock + @staticmethod + @asynccontextmanager + async def lock_source(source_id: str): + """Serialize all mutation of a source's clone dir and cached handles. + + Locks are minted on demand into ``repo_locks`` (asyncio.Lock: process- + local but fair, unlike the previous file-based lock). A scope delete + pops the dict entry while holding the lock, so after acquiring we must + re-check that ``repo_locks`` still maps ``source_id`` to the lock we + acquired — a waiter woken after a delete would otherwise proceed under + the stale lock, unserialized against holders of the freshly-minted one. + """ + while True: + lock = GitPolicyFetcher.repo_locks.setdefault(source_id, asyncio.Lock()) + async with lock: + if GitPolicyFetcher.repo_locks.get(source_id) is lock: + yield + return async def _was_fetched_after(self, t: datetime.datetime): last_fetched = GitPolicyFetcher.repos_last_fetched.get(self._source_id, None) @@ -168,15 +180,17 @@ async def fetch_and_notify_on_changes( - if the hinted commit hash is provided and is already found in the local clone we use this hint to avoid an necessary fetch. """ - repo_lock = await self._get_repo_lock() - async with repo_lock: + async with GitPolicyFetcher.lock_source(self._source_id): with tracer.trace( "git_policy_fetcher.fetch_and_notify_on_changes", resource=self._scope_id, ): if self._discover_repository(self._repo_path): logger.debug("Repo found at {path}", path=self._repo_path) - repo = self._get_valid_repo() + # The probe opens/parses a fresh Repository handle from + # disk — off the event loop so a slow disk can't stall + # every other request being served on this worker. + repo = await run_sync(self._get_valid_repo) if repo is not None: should_fetch = await self._should_fetch( repo, @@ -188,13 +202,21 @@ async def fetch_and_notify_on_changes( logger.debug( f"Fetching remote (force_fetch={force_fetch}): {self._remote} ({redact_url(self._source.url)})" ) - GitPolicyFetcher.repos_last_fetched[ - self._source_id - ] = datetime.datetime.now() + # Record the START time but write it only on + # success: a failed fetch must not look "fresh" + # to _was_fetched_after(), or it suppresses the + # forced refresh a webhook just asked for. The + # start time (not completion) is what req_time + # comparisons need: a fetch that STARTED after + # the request already satisfies it. + fetch_started = datetime.datetime.now() await run_sync( repo.remotes[self._remote].fetch, callbacks=self._auth_callbacks, ) + GitPolicyFetcher.repos_last_fetched[ + self._source_id + ] = fetch_started logger.debug( f"Fetch completed: {redact_url(self._source.url)}" ) @@ -203,11 +225,23 @@ async def fetch_and_notify_on_changes( await self._notify_on_changes(repo) return else: - # repo dir exists but invalid -> we must delete the directory + # repo dir exists but invalid -> drop the cached handle + # FIRST (it is the thing judging the dir invalid; kept, + # it would re-invalidate the fresh clone on every sync + # -> infinite re-clone loop), then delete the directory. logger.warning( "Deleting invalid repo: {path}", path=self._repo_path ) - shutil.rmtree(self._repo_path) + GitPolicyFetcher.forget_repo(str(self._repo_path)) + try: + await run_sync(shutil.rmtree, str(self._repo_path)) + except FileNotFoundError: + pass # already gone — the intended end state + except OSError as e: + logger.warning( + f"Failed to remove clone dir " + f"{self._repo_path}: {e!r}" + ) else: logger.info("Repo not found at {path}", path=self._repo_path) @@ -219,11 +253,25 @@ def _discover_repository(self, path: Path) -> bool: return discover_repository(str(path)) and git_path.exists() async def _clone(self): + if self._repo_path.exists(): + # A failed/interrupted clone leaves a partial dir; + # clone_repository refuses a non-empty destination, which would + # wedge every retry for this source. + try: + await run_sync(shutil.rmtree, str(self._repo_path)) + except FileNotFoundError: + pass # already gone — the intended end state + except OSError as e: + logger.warning(f"Failed to remove clone dir {self._repo_path}: {e!r}") logger.info( "Cloning repo at '{url}' to '{path}'", url=redact_url(self._source.url), path=self._repo_path, ) + # Same start-time rule as the fetch path above: the clone's + # negotiation reflects remote state at clone START, so that is the + # timestamp req_time comparisons need. + clone_started = datetime.datetime.now() try: repo: Repository = await run_sync( clone_repository, @@ -235,6 +283,12 @@ async def _clone(self): logger.exception(f"Could not clone repo at {redact_url(self._source.url)}") else: logger.info(f"Clone completed: {redact_url(self._source.url)}") + # Cache the fresh handle so the next sync's _get_repo() reuses it + # instead of reopening (or hitting a stale predecessor). + GitPolicyFetcher.repos[str(self._repo_path)] = repo + # A reclone just downloaded current remote state — record it so + # _was_fetched_after() doesn't force a redundant fetch next cycle. + GitPolicyFetcher.repos_last_fetched[self._source_id] = clone_started await self._notify_on_changes(repo) def _get_repo(self) -> Repository: @@ -247,7 +301,37 @@ def _get_valid_repo(self) -> Optional[Repository]: try: repo = self._get_repo() RepoInterface.verify_found_repo_matches_remote(repo, self._source.url) - return repo + # A clone can be discoverable yet unusable: refs and config + # intact but the object store gutted (crash mid-gc, disk + # corruption). A fetch then negotiates "up to date" against the + # intact refs and downloads nothing, so without this check the + # scope serves 500s forever with no self-heal. Validate that the + # tracked branch's head object is actually readable FROM DISK: + # the check must use a short-lived fresh handle, because the + # cached warm handle keeps deleted pack files readable through + # its open mmaps (unlink does not invalidate them) and would + # report the object as present. Partial corruption deeper in + # the tree is NOT caught here (that would need fsck-grade + # checks). + probe = Repository(str(self._repo_path)) + try: + try: + ref = probe.lookup_reference( + f"refs/remotes/{self._remote}/{self._source.branch}" + ) + except KeyError: + # Branch not fetched yet — the fetch path handles that. + return repo + if probe.get(ref.target) is None: + logger.warning( + "Repo at {path} has refs but an unreadable object " + "store (missing head object) — treating as invalid", + path=self._repo_path, + ) + return None + return repo + finally: + probe.free() except pygit2.GitError: logger.warning("Invalid repo at: {path}", path=self._repo_path) return None @@ -322,10 +406,22 @@ async def _notify_on_changes(self, repo: Repository): local_branch.set_target(new_revision) def _get_current_branch_head(self) -> str: - repo = self._get_repo() - head_commit_hash = RepoInterface.get_commit_hash( - repo, self._source.branch, self._remote - ) + # Opened fresh per call instead of using the shared cached handle: + # this runs on executor threads (run_sync(make_bundle) in the policy- + # bundle route) and outside lock_source, where the cached handle can + # be free()'d concurrently by a scope delete or invalid-repo recovery. + # asyncio locks don't exclude executor threads — sharing the handle + # here is a use-after-free. Same fresh-probe pattern as + # _get_valid_repo's disk-truth check. + repo = Repository(str(self._repo_path)) + try: + head_commit_hash = RepoInterface.get_commit_hash( + repo, self._source.branch, self._remote + ) + finally: + free = getattr(repo, "free", None) + if callable(free): + free() if not head_commit_hash: logger.error("Could not find current branch head") raise ValueError("Could not find current branch head") @@ -369,6 +465,30 @@ def base_dir(base_dir: Path) -> Path: def repo_clone_path(base_dir: Path, source: GitPolicyScopeSource) -> Path: return GitPolicyFetcher.base_dir(base_dir) / GitPolicyFetcher.source_id(source) + @staticmethod + def forget_repo(path: str) -> None: + """Drop the cached repository for a clone path and release its handles. + + The cached ``pygit2.Repository`` keeps OS file descriptors and mmapped + pack indexes open; without this, a deleted scope's repo pins memory and + inodes for the lifetime of the process even after the clone is removed. + ``Repository.free()`` is called only when available (the pinned pygit2 + always has it; the guard defends against test doubles and future API + changes); otherwise the dropped reference is reclaimed by GC. + """ + repo = GitPolicyFetcher.repos.pop(path, None) + if repo is None: + return + free = getattr(repo, "free", None) + if callable(free): + try: + free() + except Exception as e: + logger.warning( + f"pygit2 Repository.free() failed for {path}: {e!r}; " + "relying on GC to release the handles" + ) + class GitCallback(RemoteCallbacks): def __init__(self, source: GitPolicyScopeSource): diff --git a/packages/opal-server/opal_server/policy/watcher/task.py b/packages/opal-server/opal_server/policy/watcher/task.py index d420d183c..a2b3177f6 100644 --- a/packages/opal-server/opal_server/policy/watcher/task.py +++ b/packages/opal-server/opal_server/policy/watcher/task.py @@ -5,6 +5,7 @@ from fastapi_websocket_pubsub import Topic from fastapi_websocket_pubsub.pub_sub_server import PubSubEndpoint +from opal_common.http_utils import redact_url_in_text from opal_common.logger import logger from opal_common.sources.base_policy_source import BasePolicySource from opal_server.config import opal_server_config @@ -28,11 +29,19 @@ async def __aexit__(self, exc_type, exc, tb): async def _on_webhook(self, topic: Topic, data: Any): logger.info(f"Webhook listener triggered ({len(self._webhook_tasks)})") - for task in self._webhook_tasks: - if task.done(): - # Clean references to finished tasks - self._webhook_tasks.remove(task) - + # Rebuild rather than remove-while-iterating: list.remove() inside a + # `for t in self._webhook_tasks` loop skips the element after each removal, + # so finished tasks accumulate. Retrieve exceptions before dropping the + # references — otherwise a failed trigger() is only reported by asyncio's + # generic "Task exception was never retrieved" at GC time. + for t in self._webhook_tasks: + if t.done() and not t.cancelled() and t.exception() is not None: + # git exceptions can embed a credentialed remote URL verbatim; + # scrub here as well as in the global log patcher, which only + # runs once configure_logs() has installed it. + exc_text = redact_url_in_text(repr(t.exception())) + logger.error(f"Webhook trigger task failed: {exc_text}") + self._webhook_tasks = [t for t in self._webhook_tasks if not t.done()] self._webhook_tasks.append(asyncio.create_task(self.trigger(topic, data))) async def _listen_to_webhook_notifications(self): @@ -73,7 +82,7 @@ async def stop(self): for task in self._tasks + self._webhook_tasks: if not task.done(): task.cancel() - await asyncio.gather(*self._tasks, return_exceptions=True) + await asyncio.gather(*self._tasks, *self._webhook_tasks, return_exceptions=True) async def trigger(self, topic: Topic, data: Any): """Triggers the policy watcher from outside to check for changes (git diff --git a/packages/opal-server/opal_server/scopes/api.py b/packages/opal-server/opal_server/scopes/api.py index d54b1074c..1e3f5385c 100644 --- a/packages/opal-server/opal_server/scopes/api.py +++ b/packages/opal-server/opal_server/scopes/api.py @@ -14,7 +14,7 @@ ) from fastapi.responses import RedirectResponse from fastapi_websocket_pubsub import PubSubEndpoint -from git import InvalidGitRepositoryError +from git import InvalidGitRepositoryError, NoSuchPathError from opal_common.async_utils import run_sync from opal_common.authentication.authz import ( require_peer_type, @@ -44,6 +44,7 @@ from opal_server.data.data_update_publisher import DataUpdatePublisher from opal_server.git_fetcher import GitPolicyFetcher from opal_server.scopes.scope_repository import ScopeNotFoundError, ScopeRepository +from opal_server.scopes.service import ScopesService def verify_private_key(private_key: str, key_format: EncryptionKeyFormat) -> bool: @@ -80,6 +81,7 @@ def init_scope_router( scopes: ScopeRepository, authenticator: JWTAuthenticator, pubsub_endpoint: PubSubEndpoint, + scopes_service: ScopesService, ): router = APIRouter() @@ -171,8 +173,16 @@ async def delete_scope( logger.error(f"Unauthorized to delete scope: {repr(ex)}") raise - # TODO: This should also asynchronously clean the repo from the disk (if it's not used by other scopes) - await scopes.delete(scope_id) + try: + # Deletes the scope and also cleans the repo clone from disk and the + # GitPolicyFetcher in-memory caches (unless another scope shares them). + # The cache purge is process-local best-effort: it runs on whichever + # worker serves this DELETE; a fleet-wide purge (leader included) is + # tracked for PR3 of the leak series. + await scopes_service.delete_scope(scope_id) + except ScopeNotFoundError: + # Deleting a missing scope was always a silent no-op (204); keep it. + pass return Response(status_code=status.HTTP_204_NO_CONTENT) @@ -270,7 +280,14 @@ async def get_scope_policy( try: return await run_sync(fetcher.make_bundle, base_hash) - except (InvalidGitRepositoryError, pygit2.GitError, ValueError): + except ( + InvalidGitRepositoryError, + # A concurrent delete/recovery can rmtree the clone dir before + # Repo() opens it — fall back to the default scope, not a 500. + NoSuchPathError, + pygit2.GitError, + ValueError, + ): logger.warning( "Requested scope {scope_id} has invalid repo, returning default scope", scope_id=scope_id, @@ -295,6 +312,7 @@ async def _generate_default_scope_bundle(scope_id: str) -> PolicyBundle: except ( ScopeNotFoundError, InvalidGitRepositoryError, + NoSuchPathError, pygit2.GitError, ValueError, ): diff --git a/packages/opal-server/opal_server/scopes/service.py b/packages/opal-server/opal_server/scopes/service.py index 637e9f6c5..7455a3fe2 100644 --- a/packages/opal-server/opal_server/scopes/service.py +++ b/packages/opal-server/opal_server/scopes/service.py @@ -7,6 +7,7 @@ import git from ddtrace import tracer from fastapi_websocket_pubsub import PubSubEndpoint +from opal_common.async_utils import run_sync from opal_common.git_utils.commit_viewer import VersionedFile from opal_common.http_utils import redact_url from opal_common.logger import logger @@ -18,7 +19,11 @@ create_policy_update, create_update_all_directories_in_repo, ) -from opal_server.scopes.scope_repository import Scope, ScopeRepository +from opal_server.scopes.scope_repository import ( + Scope, + ScopeNotFoundError, + ScopeRepository, +) def is_rego_source_file( @@ -163,26 +168,129 @@ async def delete_scope(self, scope_id: str): with tracer.trace("scopes_service.delete_scope", resource=scope_id): logger.info(f"Delete scope: {scope_id}") scope = await self._scopes.get(scope_id) - url = scope.policy.url - scopes = await self._scopes.all() - remove_repo_clone = True + if not isinstance(scope.policy, GitPolicyScopeSource): + # Mirrors sync_scope: only git sources have a clone dir and + # fetcher caches to clean up. + logger.warning( + f"Scope {scope_id} has a non-git policy source, " + "deleting the scope record only" + ) + await self._scopes.delete(scope_id) + return - for scope in scopes: - if scope.scope_id != scope_id and scope.policy.url == url: - logger.info( - f"found another scope with same remote url ({scope.scope_id}), skipping clone deletion" + deleted_source = scope.policy + deleted_source_id = GitPolicyFetcher.source_id(deleted_source) + scope_dir = GitPolicyFetcher.repo_clone_path(self._base_dir, deleted_source) + + # Clone dir, the `repos` handle cache, and `repos_last_fetched` are + # all keyed by source_id (= the clone path). A sibling only shares + # storage when it resolves to the same source_id; same url with a + # different branch can shard to a different source_id (and a + # different clone dir) when SCOPES_REPO_CLONES_SHARDS > 1, so gate on + # source_id, not url — otherwise the deleted scope's clone + pygit2 + # handle leak. + # + # Serialize record-delete + sibling-check + purge on the source + # lock: two concurrent sibling deletes could otherwise both read + # the other as still-live (the check and the record delete span + # separate store round-trips) and BOTH skip the purge, orphaning + # the clone and cache entries permanently. Deleting our record + # before re-checking makes the last deleter see no sharer. + async with GitPolicyFetcher.lock_source(deleted_source_id): + try: + await self._scopes.delete(scope_id) + finally: + # The purge must stay reachable even when the record + # delete raises an ambiguous outcome (committed server- + # side, error surfaced to the client): the client retry + # then hits ScopeNotFoundError -> 204 no-op, and a purge + # gated on a clean delete would be permanently orphaned. + # If the delete genuinely failed the record survives, + # over-purging self-heals (the scope re-clones on its + # next sync), and the error still propagates as a 500. + await self._purge_source_cache_if_unshared( + deleted_source_id, scope_dir, scope_id ) - remove_repo_clone = False - break - - if remove_repo_clone: - scope_dir = GitPolicyFetcher.repo_clone_path( - self._base_dir, cast(GitPolicyScopeSource, scope.policy) - ) - shutil.rmtree(scope_dir, ignore_errors=True) - await self._scopes.delete(scope_id) + async def _purge_source_cache_if_unshared( + self, deleted_source_id: str, scope_dir: Path, scope_id: str + ): + """Remove the clone dir and the GitPolicyFetcher cache entries keyed by + ``deleted_source_id``, unless a surviving scope still shares them. + + Must run under ``lock_source(deleted_source_id)``, after scope_id's + record was deleted (so the last of two concurrent sibling deleters + sees no sharer and purges). + """ + sharing_scope_id = await self._find_scope_sharing_source( + deleted_source_id, scope_id + ) + if sharing_scope_id is not None: + logger.info( + f"Scope {sharing_scope_id} shares the same clone " + "(source id), skipping clone deletion" + ) + return + # NOTE (PR3): delete must ultimately route through the leader + # like put/refresh — only the leader should mutate the shared + # clone tree. Today this purge is process-local best-effort + # (the leader's caches leak until PR3's broadcast purge), and + # a non-leader DELETE rmtree's the shared tree unserialized + # against the leader's in-flight fetches (cross-process; + # bounded and self-healing, but an invariant break). + # The same class exists IN-process: a sync that loaded the + # scope before this delete and acquires the fresh lock after + # the purge re-clones and re-populates the caches for the dead + # scope (found by the bed's randomized churn driver; + # deterministic seed recorded there). PR3's purge routing must + # include a scope-liveness check before clone. + try: + await run_sync(shutil.rmtree, str(scope_dir)) + except FileNotFoundError: + pass # never cloned (or already gone) — nothing to clean + except OSError as e: + logger.warning( + f"Failed to remove clone dir {scope_dir} of deleted " + f"scope {scope_id}: {e!r}" + ) + GitPolicyFetcher.forget_repo(str(scope_dir)) + GitPolicyFetcher.repos_last_fetched.pop(deleted_source_id, None) + # Popped while the lock is held: lock_source waiters re-check + # the dict entry after acquiring and retry on the fresh lock. + GitPolicyFetcher.repo_locks.pop(deleted_source_id, None) + + async def _find_scope_sharing_source( + self, deleted_source_id: str, scope_id: str + ) -> Optional[str]: + """Return the id of a surviving scope that shares deleted_source_id, or + None (= safe to purge). + + Must not be able to skip the purge by raising: scope_id's record + is already deleted, so a client retry is a 204 no-op + (ScopeNotFoundError) and the purge becomes permanently + unreachable. all() does a full scan + Scope.parse_raw — a + transient store error or one malformed record throws. Over- + purging self-heals (a surviving sibling re-clones on its next + sync); under-purging is a permanent leak. + """ + try: + return next( + ( + s.scope_id + for s in await self._scopes.all() + if s.scope_id != scope_id + and isinstance(s.policy, GitPolicyScopeSource) + and GitPolicyFetcher.source_id(s.policy) == deleted_source_id + ), + None, + ) + except Exception as e: + logger.warning( + f"sibling check failed after deleting scope " + f"{scope_id}; purging defensively: {e!r}" + ) + return None async def sync_scopes(self, only_poll_updates=False, notify_on_changes=True): with tracer.trace("scopes_service.sync_scopes"): @@ -205,24 +313,37 @@ async def sync_scopes(self, only_poll_updates=False, notify_on_changes=True): skipped_scopes.append(scope) continue - try: - await self.sync_scope( - scope=scope, - force_fetch=True, - notify_on_changes=notify_on_changes, - ) - except Exception as e: - logger.exception(f"sync_scope failed for {scope.scope_id}") - + await self._sync_snapshotted_scope( + scope, force_fetch=True, notify_on_changes=notify_on_changes + ) fetched_source_ids.add(src_id) for scope in skipped_scopes: # No need to refetch the same repo, just check for changes - try: - await self.sync_scope( - scope=scope, - force_fetch=False, - notify_on_changes=notify_on_changes, - ) - except Exception as e: - logger.exception(f"sync_scope failed for {scope.scope_id}") + await self._sync_snapshotted_scope( + scope, force_fetch=False, notify_on_changes=notify_on_changes + ) + + async def _sync_snapshotted_scope( + self, scope: Scope, force_fetch: bool, notify_on_changes: bool + ): + """Sync one scope taken from a (possibly stale) all() snapshot, + swallowing per-scope errors so the sweep continues. + + Passes scope_id (not the snapshot object) so sync_scope re-gets + fresh state right before use — a delete that landed after the + snapshot surfaces as ScopeNotFoundError instead of re-cloning a + dead scope's repo and re-populating the fetcher caches for it. + """ + try: + await self.sync_scope( + scope_id=scope.scope_id, + force_fetch=force_fetch, + notify_on_changes=notify_on_changes, + ) + except ScopeNotFoundError: + logger.info( + f"scope {scope.scope_id} was deleted while sync was queued, skipping" + ) + except Exception: + logger.exception(f"sync_scope failed for {scope.scope_id}") diff --git a/packages/opal-server/opal_server/server.py b/packages/opal-server/opal_server/server.py index 0946d3d2c..7c9e4552f 100644 --- a/packages/opal-server/opal_server/server.py +++ b/packages/opal-server/opal_server/server.py @@ -4,6 +4,7 @@ import sys import traceback from functools import partial +from pathlib import Path from typing import List, Optional from fastapi import Depends, FastAPI, Request @@ -40,6 +41,7 @@ from opal_server.scopes.api import init_scope_router from opal_server.scopes.loader import load_scopes from opal_server.scopes.scope_repository import ScopeRepository +from opal_server.scopes.service import ScopesService from opal_server.security.api import init_security_router from opal_server.security.jwks import JwksStaticEndpoint from opal_server.statistics import OpalStatistics, init_statistics_router @@ -270,8 +272,15 @@ def _configure_api_routes(self, app: FastAPI): ) if opal_server_config.SCOPES: + scopes_service = ScopesService( + base_dir=Path(opal_server_config.BASE_DIR), + scopes=self._scopes, + pubsub_endpoint=self.pubsub.endpoint, + ) app.include_router( - init_scope_router(self._scopes, authenticator, self.pubsub.endpoint), + init_scope_router( + self._scopes, authenticator, self.pubsub.endpoint, scopes_service + ), tags=["Scopes"], prefix="/scopes", ) diff --git a/packages/opal-server/opal_server/tests/debug_stats_endpoint_test.py b/packages/opal-server/opal_server/tests/debug_stats_endpoint_test.py index 87e392910..00cc3247a 100644 --- a/packages/opal-server/opal_server/tests/debug_stats_endpoint_test.py +++ b/packages/opal-server/opal_server/tests/debug_stats_endpoint_test.py @@ -23,7 +23,16 @@ def test_endpoint_present_when_enabled(): resp = client.get("/internal/git-fetcher-cache-stats") assert resp.status_code == 200 body = resp.json() - assert set(body) == {"repo_locks", "repos", "repos_last_fetched", "rss_kb"} + assert set(body) == { + "pid", + "repo_locks", + "repos", + "repos_last_fetched", + "rss_kb", + "repo_locks_keys", + "repos_keys", + "repos_last_fetched_keys", + } def test_endpoint_applies_passed_dependencies(): @@ -77,7 +86,16 @@ def test_server_wiring_mounts_endpoint_when_flag_enabled(): client = TestClient(_build_server().app) resp = client.get("/internal/git-fetcher-cache-stats") assert resp.status_code == 200 - assert set(resp.json()) == {"repo_locks", "repos", "repos_last_fetched", "rss_kb"} + assert set(resp.json()) == { + "pid", + "repo_locks", + "repos", + "repos_last_fetched", + "rss_kb", + "repo_locks_keys", + "repos_keys", + "repos_last_fetched_keys", + } def test_server_wiring_omits_endpoint_when_flag_disabled(): diff --git a/packages/opal-server/opal_server/tests/debug_stats_test.py b/packages/opal-server/opal_server/tests/debug_stats_test.py index c67fb3530..b2efc5987 100644 --- a/packages/opal-server/opal_server/tests/debug_stats_test.py +++ b/packages/opal-server/opal_server/tests/debug_stats_test.py @@ -1,3 +1,4 @@ +import os import sys from opal_server.config import opal_server_config @@ -26,3 +27,50 @@ def test_stats_report_dict_sizes(monkeypatch): def test_internal_stats_flag_defaults_off(): assert opal_server_config.DEBUG_INTERNAL_STATS is False + + +def test_stats_include_pid_and_cache_keys(monkeypatch): + monkeypatch.setattr(GitPolicyFetcher, "repos", {"/clones/x": object()}) + monkeypatch.setattr(GitPolicyFetcher, "repos_last_fetched", {"sid-1": "ts"}) + monkeypatch.setattr(GitPolicyFetcher, "repo_locks", {"sid-1": object()}) + + stats = git_fetcher_cache_stats() + + assert stats["pid"] == os.getpid() + assert stats["repos_keys"] == ["/clones/x"] + assert stats["repos_last_fetched_keys"] == ["sid-1"] + assert stats["repo_locks_keys"] == ["sid-1"] + + +class _ChurningDict(dict): + """Simulates a dict being resized by another thread mid-iteration: any + direct keys() iteration blows up; only an atomic .copy() snapshot is safe.""" + + def keys(self): + raise RuntimeError("dictionary changed size during iteration (simulated)") + + +def test_stats_survive_concurrent_cache_mutation(monkeypatch): + monkeypatch.setattr(GitPolicyFetcher, "repos", _ChurningDict({"/c/x": object()})) + monkeypatch.setattr( + GitPolicyFetcher, "repos_last_fetched", _ChurningDict({"sid-1": "ts"}) + ) + monkeypatch.setattr(GitPolicyFetcher, "repo_locks", _ChurningDict({"sid-1": 1})) + + stats = git_fetcher_cache_stats() # must not iterate the live dicts + + assert stats["repos_keys"] == ["/c/x"] + assert stats["repo_locks_keys"] == ["sid-1"] + assert stats["repos_last_fetched_keys"] == ["sid-1"] + + +def test_stats_counts_derive_from_the_same_snapshot_as_keys(monkeypatch): + monkeypatch.setattr(GitPolicyFetcher, "repos", {"/c/a": object(), "/c/b": object()}) + monkeypatch.setattr(GitPolicyFetcher, "repos_last_fetched", {"s1": "t", "s2": "t"}) + monkeypatch.setattr(GitPolicyFetcher, "repo_locks", {"s1": object()}) + + stats = git_fetcher_cache_stats() + + assert stats["repos"] == len(stats["repos_keys"]) == 2 + assert stats["repos_last_fetched"] == len(stats["repos_last_fetched_keys"]) == 2 + assert stats["repo_locks"] == len(stats["repo_locks_keys"]) == 1 diff --git a/packages/opal-server/opal_server/tests/delete_scope_cache_purge_test.py b/packages/opal-server/opal_server/tests/delete_scope_cache_purge_test.py new file mode 100644 index 000000000..d8164e53b --- /dev/null +++ b/packages/opal-server/opal_server/tests/delete_scope_cache_purge_test.py @@ -0,0 +1,413 @@ +import asyncio + +import pytest +from opal_common.schemas.policy_source import GitPolicyScopeSource, NoAuthData +from opal_common.schemas.scopes import Scope +from opal_server.git_fetcher import GitPolicyFetcher +from opal_server.scopes.scope_repository import ScopeNotFoundError +from opal_server.scopes.service import ScopesService + + +class FakeScopeRepository: + def __init__(self, scopes): + self._scopes = {s.scope_id: s for s in scopes} + + async def get(self, scope_id): + await asyncio.sleep(0) + if scope_id not in self._scopes: + raise ScopeNotFoundError(scope_id) + return self._scopes[scope_id] + + async def all(self): + await asyncio.sleep(0) + return list(self._scopes.values()) + + async def delete(self, scope_id): + await asyncio.sleep(0) + self._scopes.pop(scope_id, None) + + +def _scope(scope_id, url, branch="main"): + return Scope( + scope_id=scope_id, + policy=GitPolicyScopeSource( + source_type="git", + url=url, + branch=branch, + auth=NoAuthData(auth_type="none"), + ), + data={"entries": []}, + ) + + +@pytest.fixture(autouse=True) +def clear_caches(): + GitPolicyFetcher.repos.clear() + GitPolicyFetcher.repos_last_fetched.clear() + GitPolicyFetcher.repo_locks.clear() + yield + GitPolicyFetcher.repos.clear() + GitPolicyFetcher.repos_last_fetched.clear() + GitPolicyFetcher.repo_locks.clear() + + +@pytest.mark.asyncio +async def test_delete_unique_scope_purges_caches(tmp_path, monkeypatch): + scope = _scope("only", "https://git/repo-a.git") + repo = FakeScopeRepository([scope]) + svc = ScopesService(base_dir=tmp_path, scopes=repo, pubsub_endpoint=None) + + src = scope.policy + sid = GitPolicyFetcher.source_id(src) + clone_path = str(GitPolicyFetcher.repo_clone_path(tmp_path, src)) + GitPolicyFetcher.repos[clone_path] = object() + GitPolicyFetcher.repos_last_fetched[sid] = "ts" + GitPolicyFetcher.repo_locks[sid] = asyncio.Lock() + + monkeypatch.setattr( + "opal_server.scopes.service.shutil.rmtree", lambda *a, **k: None + ) + + await svc.delete_scope("only") + + assert clone_path not in GitPolicyFetcher.repos + assert sid not in GitPolicyFetcher.repos_last_fetched + assert sid not in GitPolicyFetcher.repo_locks + + +@pytest.mark.asyncio +async def test_delete_keeps_caches_when_sibling_shares_source(tmp_path, monkeypatch): + a = _scope("a", "https://git/shared.git") + b = _scope("b", "https://git/shared.git") # same url+branch -> same source_id + repo = FakeScopeRepository([a, b]) + svc = ScopesService(base_dir=tmp_path, scopes=repo, pubsub_endpoint=None) + + sid = GitPolicyFetcher.source_id(a.policy) + clone_path = str(GitPolicyFetcher.repo_clone_path(tmp_path, a.policy)) + GitPolicyFetcher.repos[clone_path] = object() + GitPolicyFetcher.repos_last_fetched[sid] = "ts" + GitPolicyFetcher.repo_locks[sid] = asyncio.Lock() + + rmtree_calls = [] + monkeypatch.setattr( + "opal_server.scopes.service.shutil.rmtree", + lambda p, **k: rmtree_calls.append(p), + ) + + await svc.delete_scope("a") + + assert rmtree_calls == [] # sibling shares the source id; clone must survive + assert clone_path in GitPolicyFetcher.repos + assert sid in GitPolicyFetcher.repos_last_fetched + assert sid in GitPolicyFetcher.repo_locks + + +@pytest.mark.asyncio +async def test_concurrent_sibling_deletes_still_purge(tmp_path, monkeypatch): + """Two concurrent DELETEs of source-sharing scopes must not BOTH skip the + purge (TOCTOU on the sibling check) — whichever finishes last must + purge.""" + a = _scope("a", "https://git/shared.git") + b = _scope("b", "https://git/shared.git") + repo = FakeScopeRepository([a, b]) + svc = ScopesService(base_dir=tmp_path, scopes=repo, pubsub_endpoint=None) + sid = GitPolicyFetcher.source_id(a.policy) + clone_path = str(GitPolicyFetcher.repo_clone_path(tmp_path, a.policy)) + GitPolicyFetcher.repos[clone_path] = object() + GitPolicyFetcher.repos_last_fetched[sid] = "ts" + rmtree_calls = [] + monkeypatch.setattr( + "opal_server.scopes.service.shutil.rmtree", + lambda p, **k: rmtree_calls.append(str(p)), + ) + + await asyncio.gather(svc.delete_scope("a"), svc.delete_scope("b")) + + assert rmtree_calls == [ + clone_path + ], "concurrent sibling deletes both skipped the purge (TOCTOU)" + assert clone_path not in GitPolicyFetcher.repos + assert sid not in GitPolicyFetcher.repos_last_fetched + + +class _AllRaisesAfterDeleteRepository(FakeScopeRepository): + """All() blows up only on the post-delete sibling re-check (e.g. a + transient store error or one malformed record failing parse_raw).""" + + def __init__(self, scopes): + super().__init__(scopes) + self.deleted_once = False + + async def delete(self, scope_id): + await super().delete(scope_id) + self.deleted_once = True + + async def all(self): + if self.deleted_once: + raise RuntimeError("store scan failed (malformed record)") + return await super().all() + + +@pytest.mark.asyncio +async def test_sibling_check_failure_still_purges(tmp_path, monkeypatch): + """If the post-delete sibling re-check raises, the purge must still run: + + the record is already deleted, so the client's retry is a 204 no-op + (ScopeNotFoundError) and a skipped purge is a permanent leak. Over- + purging is safe — a surviving sibling re-clones on its next sync. + """ + scope = _scope("only", "https://git/repo-a.git") + repo = _AllRaisesAfterDeleteRepository([scope]) + svc = ScopesService(base_dir=tmp_path, scopes=repo, pubsub_endpoint=None) + + sid = GitPolicyFetcher.source_id(scope.policy) + clone_path = str(GitPolicyFetcher.repo_clone_path(tmp_path, scope.policy)) + GitPolicyFetcher.repos[clone_path] = object() + GitPolicyFetcher.repos_last_fetched[sid] = "ts" + + rmtree_calls = [] + monkeypatch.setattr( + "opal_server.scopes.service.shutil.rmtree", + lambda p, **k: rmtree_calls.append(str(p)), + ) + + from opal_common.logger import logger as opal_logger + + records = [] + sink_id = opal_logger.add(lambda m: records.append(str(m)), level="WARNING") + try: + # Must not raise: the re-check failure is downgraded to a warning. + await svc.delete_scope("only") + finally: + opal_logger.remove(sink_id) + + assert rmtree_calls == [clone_path], ( + "sibling re-check failure skipped the purge — permanent leak " + "(retry is a 204 no-op, the purge is unreachable)" + ) + assert clone_path not in GitPolicyFetcher.repos + assert sid not in GitPolicyFetcher.repos_last_fetched + assert any( + "sibling check failed" in r and "purging defensively" in r for r in records + ), f"defensive purge not logged: {records}" + + +@pytest.mark.asyncio +async def test_delete_purges_when_sibling_shares_url_but_not_source( + tmp_path, monkeypatch +): + """Same url, different branch, sharded clones (SCOPES_REPO_CLONES_SHARDS>1) + resolve to different source_ids -> different clone dirs. + + Deleting one must still purge its own clone + caches; the url- + sharing sibling lives elsewhere. + """ + # shards=4: branch "main" -> index 1, "dev" -> index 3 (distinct source_ids). + monkeypatch.setattr( + "opal_server.git_fetcher.opal_server_config.SCOPES_REPO_CLONES_SHARDS", 4 + ) + a = _scope("a", "https://git/shared.git", branch="main") + b = _scope("b", "https://git/shared.git", branch="dev") # same url, diff source_id + assert GitPolicyFetcher.source_id(a.policy) != GitPolicyFetcher.source_id(b.policy) + + repo = FakeScopeRepository([a, b]) + svc = ScopesService(base_dir=tmp_path, scopes=repo, pubsub_endpoint=None) + + sid_a = GitPolicyFetcher.source_id(a.policy) + clone_path_a = str(GitPolicyFetcher.repo_clone_path(tmp_path, a.policy)) + GitPolicyFetcher.repos[clone_path_a] = object() + GitPolicyFetcher.repos_last_fetched[sid_a] = "ts" + + rmtree_calls = [] + monkeypatch.setattr( + "opal_server.scopes.service.shutil.rmtree", + lambda p, **k: rmtree_calls.append(str(p)), + ) + + await svc.delete_scope("a") + + assert rmtree_calls == [clone_path_a] # its own clone dir removed + assert clone_path_a not in GitPolicyFetcher.repos + assert sid_a not in GitPolicyFetcher.repos_last_fetched + + +@pytest.mark.asyncio +async def test_delete_serializes_against_inflight_repo_lock(tmp_path, monkeypatch): + """The purge must wait for the repo lock held by an in-flight fetch — + otherwise it rmtree's the clone and free()s the pygit2 handle out from + under the fetch thread.""" + scope = _scope("only", "https://git/repo-a.git") + repo = FakeScopeRepository([scope]) + svc = ScopesService(base_dir=tmp_path, scopes=repo, pubsub_endpoint=None) + + sid = GitPolicyFetcher.source_id(scope.policy) + clone_path = str(GitPolicyFetcher.repo_clone_path(tmp_path, scope.policy)) + GitPolicyFetcher.repos[clone_path] = object() + + rmtree_calls = [] + monkeypatch.setattr( + "opal_server.scopes.service.shutil.rmtree", + lambda p, **k: rmtree_calls.append(str(p)), + ) + + # Simulate an in-flight fetch holding the repo lock. + lock = GitPolicyFetcher.repo_locks.setdefault(sid, asyncio.Lock()) + await lock.acquire() + try: + delete_task = asyncio.create_task(svc.delete_scope("only")) + for _ in range(10): # give the delete every chance to (wrongly) proceed + await asyncio.sleep(0) + assert not delete_task.done(), "delete_scope did not wait for the repo lock" + assert rmtree_calls == [], "purge ran while the fetch held the repo lock" + finally: + lock.release() + + await asyncio.wait_for(delete_task, timeout=5) + assert rmtree_calls == [clone_path] + assert clone_path not in GitPolicyFetcher.repos + assert sid not in GitPolicyFetcher.repo_locks + + +@pytest.mark.asyncio +async def test_lock_source_waiter_retries_after_delete_pops_entry(): + """A waiter queued on the old lock must not proceed under it once a delete + popped the entry — it retries and serializes on the freshly-minted lock.""" + sid = "some-source-id" + events = [] + + async def deleter(): + async with GitPolicyFetcher.lock_source(sid): + events.append("deleter-in") + await asyncio.sleep(0.01) # let the waiter queue on this lock + GitPolicyFetcher.repo_locks.pop(sid, None) + events.append("deleter-out") + + async def waiter(): + async with GitPolicyFetcher.lock_source(sid): + events.append("waiter-in") + # We must hold the *current* dict entry, not the popped one. + assert GitPolicyFetcher.repo_locks.get(sid) is not None + + await asyncio.gather(deleter(), waiter()) + assert events == ["deleter-in", "deleter-out", "waiter-in"] + + +class _AmbiguousDeleteRepository(FakeScopeRepository): + """Delete commits server-side but the client sees an error (classic + dropped-connection/timeout ambiguous store outcome).""" + + async def delete(self, scope_id): + await super().delete(scope_id) + raise ConnectionError("connection dropped after the delete committed") + + +@pytest.mark.asyncio +async def test_purge_still_runs_when_record_delete_raises_ambiguously( + tmp_path, monkeypatch +): + """If the record delete commits but raises to the caller, the purge must + still run: the record is gone, so a client retry is a 204 no-op + (ScopeNotFoundError) and a purge gated on a clean delete is permanently + orphaned. The error must still propagate (the client sees a 500 and can + retry).""" + scope = _scope("only", "https://git/repo-a.git") + repo = _AmbiguousDeleteRepository([scope]) + svc = ScopesService(base_dir=tmp_path, scopes=repo, pubsub_endpoint=None) + + sid = GitPolicyFetcher.source_id(scope.policy) + clone_path = str(GitPolicyFetcher.repo_clone_path(tmp_path, scope.policy)) + GitPolicyFetcher.repos[clone_path] = object() + GitPolicyFetcher.repos_last_fetched[sid] = "ts" + + rmtree_calls = [] + monkeypatch.setattr( + "opal_server.scopes.service.shutil.rmtree", + lambda p, **k: rmtree_calls.append(str(p)), + ) + + with pytest.raises(ConnectionError): + await svc.delete_scope("only") + + assert rmtree_calls == [clone_path], ( + "record-delete failure skipped the purge — the retry is a 204 no-op, " + "so the clone and cache entries leak permanently" + ) + assert clone_path not in GitPolicyFetcher.repos + assert sid not in GitPolicyFetcher.repos_last_fetched + assert sid not in GitPolicyFetcher.repo_locks + + +class _DeletedAfterSnapshotRepository(FakeScopeRepository): + """All() still returns the scope (a stale sync_scopes snapshot); get() + reports it deleted — models a DELETE landing between the two.""" + + async def get(self, scope_id): + await asyncio.sleep(0) + raise ScopeNotFoundError(scope_id) + + +@pytest.mark.asyncio +async def test_sync_scopes_skips_scope_deleted_after_snapshot(tmp_path, monkeypatch): + """A scope deleted between sync_scopes' all() snapshot and its queued + sync_scope call must not be fetched — re-cloning it would re-populate + repos/repos_last_fetched/repo_locks for a dead scope (leaked until + restart).""" + scope = _scope("dead", "https://git/repo-a.git") + repo = _DeletedAfterSnapshotRepository([scope]) + svc = ScopesService(base_dir=tmp_path, scopes=repo, pubsub_endpoint=None) + + fetch_calls = [] + + async def fake_fetch(self, *args, **kwargs): + fetch_calls.append(self._scope_id) + + monkeypatch.setattr(GitPolicyFetcher, "fetch_and_notify_on_changes", fake_fetch) + + await svc.sync_scopes() # must not raise: the delete race is expected + + assert fetch_calls == [], "sync fetched a scope whose delete already landed" + assert not GitPolicyFetcher.repos + assert not GitPolicyFetcher.repos_last_fetched + assert not GitPolicyFetcher.repo_locks + + +@pytest.mark.asyncio +async def test_recreate_after_delete_serializes_and_sees_clean_caches( + tmp_path, monkeypatch +): + """A re-create's first sync queued during a delete must run only after the + purge completed, on the freshly-minted lock, and see empty caches.""" + scope = _scope("only", "https://git/repo-a.git") + repo = FakeScopeRepository([scope]) + svc = ScopesService(base_dir=tmp_path, scopes=repo, pubsub_endpoint=None) + sid = GitPolicyFetcher.source_id(scope.policy) + clone_path = str(GitPolicyFetcher.repo_clone_path(tmp_path, scope.policy)) + GitPolicyFetcher.repos[clone_path] = object() + monkeypatch.setattr( + "opal_server.scopes.service.shutil.rmtree", lambda *a, **k: None + ) + + order = [] + gate = GitPolicyFetcher.repo_locks.setdefault(sid, asyncio.Lock()) + await gate.acquire() # simulate the in-flight fetch the delete must wait on + try: + + async def deleter(): + await svc.delete_scope("only") + order.append("delete-done") + + async def recreator(): # stands in for the re-created scope's first sync + async with GitPolicyFetcher.lock_source(sid): + order.append(("recreate-in", clone_path in GitPolicyFetcher.repos)) + + d = asyncio.create_task(deleter()) + for _ in range(5): + await asyncio.sleep(0) # deleter queues on the held lock first + r = asyncio.create_task(recreator()) + for _ in range(5): + await asyncio.sleep(0) # recreator queues behind it + finally: + gate.release() + + await asyncio.wait_for(asyncio.gather(d, r), timeout=5) + assert order == ["delete-done", ("recreate-in", False)] diff --git a/packages/opal-server/opal_server/tests/delete_scope_route_test.py b/packages/opal-server/opal_server/tests/delete_scope_route_test.py new file mode 100644 index 000000000..70ff5b7f5 --- /dev/null +++ b/packages/opal-server/opal_server/tests/delete_scope_route_test.py @@ -0,0 +1,97 @@ +import pytest +from fastapi import FastAPI +from fastapi.testclient import TestClient +from opal_common.schemas.policy_source import GitPolicyScopeSource, NoAuthData +from opal_common.schemas.scopes import Scope +from opal_server.git_fetcher import GitPolicyFetcher +from opal_server.scopes.api import init_scope_router +from opal_server.scopes.scope_repository import ScopeNotFoundError +from opal_server.scopes.service import ScopesService + + +class FakeScopeRepository: + def __init__(self, scopes): + self._scopes = {s.scope_id: s for s in scopes} + + async def get(self, scope_id): + if scope_id not in self._scopes: + raise ScopeNotFoundError(scope_id) + return self._scopes[scope_id] + + async def all(self): + return list(self._scopes.values()) + + async def delete(self, scope_id): + self._scopes.pop(scope_id, None) + + +class FakeAuthenticator: + """Mimics a JWTAuthenticator whose verifier is disabled (no public key).""" + + enabled = False + + def __call__(self): + return {} + + +def _scope(scope_id, url, branch="main"): + return Scope( + scope_id=scope_id, + policy=GitPolicyScopeSource( + source_type="git", + url=url, + branch=branch, + auth=NoAuthData(auth_type="none"), + ), + data={"entries": []}, + ) + + +def _client(repo, base_dir): + service = ScopesService(base_dir=base_dir, scopes=repo, pubsub_endpoint=None) + app = FastAPI() + app.include_router( + init_scope_router(repo, FakeAuthenticator(), None, service), + prefix="/scopes", + ) + return TestClient(app) + + +@pytest.fixture(autouse=True) +def clear_caches(): + GitPolicyFetcher.repos.clear() + GitPolicyFetcher.repos_last_fetched.clear() + GitPolicyFetcher.repo_locks.clear() + yield + GitPolicyFetcher.repos.clear() + GitPolicyFetcher.repos_last_fetched.clear() + GitPolicyFetcher.repo_locks.clear() + + +def test_delete_route_purges_fetcher_caches(tmp_path, monkeypatch): + """DELETE /scopes/{id} must flow through ScopesService.delete_scope so the + GitPolicyFetcher caches drain (the git-leak churn gate).""" + scope = _scope("only", "https://git/repo-a.git") + repo = FakeScopeRepository([scope]) + + sid = GitPolicyFetcher.source_id(scope.policy) + clone_path = str(GitPolicyFetcher.repo_clone_path(tmp_path, scope.policy)) + GitPolicyFetcher.repos[clone_path] = object() + GitPolicyFetcher.repos_last_fetched[sid] = "ts" + + monkeypatch.setattr( + "opal_server.scopes.service.shutil.rmtree", lambda *a, **k: None + ) + + resp = _client(repo, tmp_path).delete("/scopes/only") + + assert resp.status_code == 204 + assert clone_path not in GitPolicyFetcher.repos + assert sid not in GitPolicyFetcher.repos_last_fetched + + +def test_delete_route_missing_scope_stays_204(tmp_path): + """Deleting a nonexistent scope was a silent no-op (204) before the purge + wiring and must remain one.""" + resp = _client(FakeScopeRepository([]), tmp_path).delete("/scopes/ghost") + assert resp.status_code == 204 diff --git a/packages/opal-server/opal_server/tests/forget_repo_test.py b/packages/opal-server/opal_server/tests/forget_repo_test.py new file mode 100644 index 000000000..826c2fdc6 --- /dev/null +++ b/packages/opal-server/opal_server/tests/forget_repo_test.py @@ -0,0 +1,42 @@ +from opal_server.git_fetcher import GitPolicyFetcher + + +class _FakeRepo: + def __init__(self): + self.freed = False + + def free(self): + self.freed = True + + +def test_forget_repo_pops_and_frees(monkeypatch): + fake = _FakeRepo() + monkeypatch.setattr(GitPolicyFetcher, "repos", {"/clones/x": fake}) + + GitPolicyFetcher.forget_repo("/clones/x") + + assert "/clones/x" not in GitPolicyFetcher.repos + assert fake.freed is True + + +def test_forget_repo_unknown_path_is_noop(monkeypatch): + monkeypatch.setattr(GitPolicyFetcher, "repos", {}) + GitPolicyFetcher.forget_repo("/clones/missing") # must not raise + assert GitPolicyFetcher.repos == {} + + +def test_forget_repo_without_free_method(monkeypatch): + monkeypatch.setattr(GitPolicyFetcher, "repos", {"/clones/y": object()}) + GitPolicyFetcher.forget_repo("/clones/y") # object() has no .free(); must not raise + assert "/clones/y" not in GitPolicyFetcher.repos + + +class _ExplodingRepo: + def free(self): + raise RuntimeError("free failed") + + +def test_forget_repo_free_raises_entry_still_popped(monkeypatch): + monkeypatch.setattr(GitPolicyFetcher, "repos", {"/clones/z": _ExplodingRepo()}) + GitPolicyFetcher.forget_repo("/clones/z") # must swallow (and log) the error + assert "/clones/z" not in GitPolicyFetcher.repos diff --git a/packages/opal-server/opal_server/tests/invalid_repo_recovery_test.py b/packages/opal-server/opal_server/tests/invalid_repo_recovery_test.py new file mode 100644 index 000000000..2e4c5b346 --- /dev/null +++ b/packages/opal-server/opal_server/tests/invalid_repo_recovery_test.py @@ -0,0 +1,310 @@ +from pathlib import Path + +import pygit2 +import pytest +from opal_common.schemas.policy_source import GitPolicyScopeSource, NoAuthData +from opal_server.git_fetcher import GitPolicyFetcher + + +def _make_fetcher(base_dir: Path, scope_id: str, url: str) -> GitPolicyFetcher: + source = GitPolicyScopeSource( + source_type="git", + url=url, + branch="main", + auth=NoAuthData(auth_type="none"), + ) + return GitPolicyFetcher(base_dir=base_dir, scope_id=scope_id, source=source) + + +@pytest.fixture(autouse=True) +def _reset_class_state(): + GitPolicyFetcher.repos.clear() + GitPolicyFetcher.repos_last_fetched.clear() + GitPolicyFetcher.repo_locks.clear() + yield + GitPolicyFetcher.repos.clear() + GitPolicyFetcher.repos_last_fetched.clear() + GitPolicyFetcher.repo_locks.clear() + + +class _BrokenRepo: + """Simulates a cached pygit2 handle whose backing clone dir went bad.""" + + freed = False + + @property + def remotes(self): + raise pygit2.GitError("stale handle: backing files gone") + + def free(self): + self.freed = True + + +@pytest.mark.asyncio +async def test_recovery_forgets_stale_cached_handle(monkeypatch, tmp_path): + """Bug A: without forget_repo in the recovery branch, the broken cached + handle survives and re-invalidates every fresh clone -> infinite loop.""" + fetcher = _make_fetcher(tmp_path, "s", "https://example.com/r.git") + path = str(fetcher._repo_path) + broken = _BrokenRepo() + GitPolicyFetcher.repos[path] = broken + monkeypatch.setattr(fetcher, "_discover_repository", lambda p: True) + monkeypatch.setattr("opal_server.git_fetcher.shutil.rmtree", lambda p, **k: None) + clone_calls = [] + + async def fake_clone(): + clone_calls.append(True) + + monkeypatch.setattr(fetcher, "_clone", fake_clone) + + await fetcher.fetch_and_notify_on_changes() + + assert clone_calls == [True] + assert GitPolicyFetcher.repos.get(path) is not broken, ( + "recovery left the stale handle cached — next sync re-invalidates " + "the fresh clone (infinite re-clone loop)" + ) + assert broken.freed is True, "recovery evicted the handle without free()ing it" + + +@pytest.mark.asyncio +async def test_clone_caches_the_fresh_handle(monkeypatch, tmp_path): + fetcher = _make_fetcher(tmp_path, "s", "https://example.com/r.git") + fresh = object() + monkeypatch.setattr( + "opal_server.git_fetcher.clone_repository", lambda *a, **k: fresh + ) + + async def _no_notify(repo): + return None + + monkeypatch.setattr(fetcher, "_notify_on_changes", _no_notify) + + await fetcher._clone() + + assert GitPolicyFetcher.repos.get(str(fetcher._repo_path)) is fresh, ( + "fresh clone's handle not cached — _get_repo would reopen (or worse, " + "return a stale entry) on the next sync" + ) + + +@pytest.mark.asyncio +async def test_clone_clears_partial_dir_before_cloning(monkeypatch, tmp_path): + """U9: a failed clone leaves a partial dir; clone_repository refuses a + non-empty destination, wedging every retry.""" + fetcher = _make_fetcher(tmp_path, "s", "https://example.com/r.git") + fetcher._repo_path.mkdir(parents=True) + (fetcher._repo_path / "leftover.pack").write_text("partial clone debris") + seen = {} + + def fake_clone(url, path, callbacks=None): + p = Path(path) + seen["nonempty"] = p.exists() and any(p.iterdir()) + return object() + + monkeypatch.setattr("opal_server.git_fetcher.clone_repository", fake_clone) + + async def _no_notify(repo): + return None + + monkeypatch.setattr(fetcher, "_notify_on_changes", _no_notify) + + await fetcher._clone() + + assert seen["nonempty"] is False, ( + "partial dir not cleared before clone — a real clone_repository " + "raises on a non-empty destination, wedging the scope forever" + ) + + +class _RemoteStub: + def __init__(self, url): + self.url = url + self.name = "origin" + + +class _Ref: + target = "deadbeef" * 5 + + +class _WarmHandle: + """Cached handle over a gutted clone: its open mmaps keep the deleted + pack files readable (unlink does not invalidate them), so through THIS + handle the repo looks perfectly healthy.""" + + def __init__(self, url): + self.remotes = [_RemoteStub(url)] + self.freed = False + + def lookup_reference(self, name): + return _Ref() + + def get(self, oid): + return object() # mmap still serves the object + + def free(self): + self.freed = True + + +class _GuttedProbe: + """What a fresh on-disk handle sees: refs intact, head object missing.""" + + def __init__(self, path): + self.freed = False + + def lookup_reference(self, name): + return _Ref() + + def get(self, oid): + return None # object store gutted on disk + + def free(self): + self.freed = True + + +@pytest.mark.asyncio +async def test_gutted_object_store_triggers_recovery(monkeypatch, tmp_path): + """Refs intact + objects missing on disk must be treated as invalid even + though the warm cached handle still reads everything via its mmaps: + fetch would negotiate "up to date" and the scope would serve 500s + forever otherwise.""" + url = "https://example.com/r.git" + fetcher = _make_fetcher(tmp_path, "s", url) + path = str(fetcher._repo_path) + warm = _WarmHandle(url) + GitPolicyFetcher.repos[path] = warm + probes = [] + + def fake_repository(p): + probe = _GuttedProbe(p) + probes.append(probe) + return probe + + monkeypatch.setattr("opal_server.git_fetcher.Repository", fake_repository) + monkeypatch.setattr(fetcher, "_discover_repository", lambda p: True) + monkeypatch.setattr("opal_server.git_fetcher.shutil.rmtree", lambda p, **k: None) + clone_calls = [] + + async def fake_clone(): + clone_calls.append(True) + + monkeypatch.setattr(fetcher, "_clone", fake_clone) + + await fetcher.fetch_and_notify_on_changes() + + assert clone_calls == [True], "gutted clone was not routed to recovery" + assert path not in GitPolicyFetcher.repos + assert warm.freed is True, "recovery evicted the warm handle without free()" + assert probes and all(p.freed for p in probes), "disk probe handle leaked" + + +class _PoisonRepo: + """Any attribute access means the shared cached handle was touched.""" + + def __getattr__(self, name): + raise AssertionError( + "shared cached handle was read outside lock_source (UAF hazard)" + ) + + +def test_branch_head_does_not_touch_shared_handle(monkeypatch, tmp_path): + fetcher = _make_fetcher(tmp_path, "s", "https://example.com/r.git") + path = str(fetcher._repo_path) + GitPolicyFetcher.repos[path] = _PoisonRepo() + fresh = object() + monkeypatch.setattr("opal_server.git_fetcher.Repository", lambda p: fresh) + monkeypatch.setattr( + "opal_server.git_fetcher.RepoInterface.get_commit_hash", + lambda repo, branch, remote: "abc123" if repo is fresh else None, + ) + + assert fetcher._get_current_branch_head() == "abc123" + + +class _HealthyCachedRepo: + """The warm cached handle _get_valid_repo returns when the disk probe + confirms the repo is healthy.""" + + def __init__(self, url): + self.remotes = [_RemoteStub(url)] + self.freed = False + + def free(self): + self.freed = True + + +class _BranchNotFetchedProbe: + """Fresh on-disk probe when the tracked branch has never been fetched (e.g. + right after a scope's remote branch config changed).""" + + def __init__(self, path): + self.freed = False + + def lookup_reference(self, name): + raise KeyError(name) + + def free(self): + self.freed = True + + +def test_get_valid_repo_tolerates_branch_not_yet_fetched(monkeypatch, tmp_path): + """The fetch path is responsible for missing branches -- the probe must not + treat that as corruption and must still return the cached repo.""" + url = "https://example.com/r.git" + fetcher = _make_fetcher(tmp_path, "s", url) + path = str(fetcher._repo_path) + cached = _HealthyCachedRepo(url) + GitPolicyFetcher.repos[path] = cached + probes = [] + + def fake_repository(p): + probe = _BranchNotFetchedProbe(p) + probes.append(probe) + return probe + + monkeypatch.setattr("opal_server.git_fetcher.Repository", fake_repository) + + result = fetcher._get_valid_repo() + + assert result is cached + assert probes and all(p.freed for p in probes), "disk probe handle leaked" + + +class _HealthyProbe: + """Fresh on-disk probe when the head object is actually readable from disk + (the healthy case).""" + + def __init__(self, path): + self.freed = False + + def lookup_reference(self, name): + return _Ref() + + def get(self, oid): + return object() # readable from disk + + def free(self): + self.freed = True + + +def test_get_valid_repo_returns_cached_when_probe_confirms_healthy( + monkeypatch, tmp_path +): + url = "https://example.com/r.git" + fetcher = _make_fetcher(tmp_path, "s", url) + path = str(fetcher._repo_path) + cached = _HealthyCachedRepo(url) + GitPolicyFetcher.repos[path] = cached + probes = [] + + def fake_repository(p): + probe = _HealthyProbe(p) + probes.append(probe) + return probe + + monkeypatch.setattr("opal_server.git_fetcher.Repository", fake_repository) + + result = fetcher._get_valid_repo() + + assert result is cached + assert probes and all(p.freed for p in probes), "disk probe handle leaked" diff --git a/packages/opal-server/opal_server/tests/scope_policy_fallback_test.py b/packages/opal-server/opal_server/tests/scope_policy_fallback_test.py new file mode 100644 index 000000000..5034bea65 --- /dev/null +++ b/packages/opal-server/opal_server/tests/scope_policy_fallback_test.py @@ -0,0 +1,116 @@ +"""GET /scopes/{scope_id}/policy fallback when the clone dir vanishes. + +A concurrent delete (or invalid-repo recovery) can rmtree a live scope's +clone between the scope-record read and make_bundle opening the repo. +The route must fall back to the default scope's bundle, not 500. (Known +limitation, tracked for PR3: a live scope is briefly served the default +bundle instead of a retryable error.) +""" + +import pytest +from fastapi import FastAPI +from fastapi.testclient import TestClient +from git import NoSuchPathError +from opal_common.schemas.policy import PolicyBundle +from opal_common.schemas.policy_source import GitPolicyScopeSource, NoAuthData +from opal_common.schemas.scopes import Scope +from opal_server.git_fetcher import GitPolicyFetcher +from opal_server.scopes.api import init_scope_router +from opal_server.scopes.scope_repository import ScopeNotFoundError +from opal_server.scopes.service import ScopesService + + +class FakeScopeRepository: + def __init__(self, scopes): + self._scopes = {s.scope_id: s for s in scopes} + + async def get(self, scope_id): + if scope_id not in self._scopes: + raise ScopeNotFoundError(scope_id) + return self._scopes[scope_id] + + async def all(self): + return list(self._scopes.values()) + + async def delete(self, scope_id): + self._scopes.pop(scope_id, None) + + +class FakeAuthenticator: + """Mimics a JWTAuthenticator whose verifier is disabled (no public key).""" + + enabled = False + + def __call__(self): + return {} + + +def _scope(scope_id, url, branch="main"): + return Scope( + scope_id=scope_id, + policy=GitPolicyScopeSource( + source_type="git", + url=url, + branch=branch, + auth=NoAuthData(auth_type="none"), + ), + data={"entries": []}, + ) + + +def _client(repo, base_dir): + service = ScopesService(base_dir=base_dir, scopes=repo, pubsub_endpoint=None) + app = FastAPI() + app.include_router( + init_scope_router(repo, FakeAuthenticator(), None, service), + prefix="/scopes", + ) + return TestClient(app) + + +def _default_bundle(): + return PolicyBundle( + manifest=[], hash="default-head", data_modules=[], policy_modules=[] + ) + + +def test_live_scope_clone_vanish_falls_back_to_default_bundle(tmp_path, monkeypatch): + """A live scope whose clone dir vanished mid-request (NoSuchPathError from + make_bundle) must be served the default scope's bundle, not a 500.""" + live = _scope("live", "https://git/live.git") + default = _scope("default", "https://git/default.git") + repo = FakeScopeRepository([live, default]) + + def fake_make_bundle(self, base_hash): + if self._scope_id == "live": + raise NoSuchPathError(str(tmp_path / "gone")) + return _default_bundle() + + monkeypatch.setattr(GitPolicyFetcher, "make_bundle", fake_make_bundle) + monkeypatch.setattr( + "opal_server.scopes.api.opal_server_config.BASE_DIR", str(tmp_path) + ) + + resp = _client(repo, tmp_path).get("/scopes/live/policy") + + assert resp.status_code == 200 + assert resp.json()["hash"] == "default-head" + + +def test_missing_scope_still_falls_back_to_default_bundle(tmp_path, monkeypatch): + """The pre-existing scope-not-found fallback must keep working alongside + the clone-vanish branch.""" + default = _scope("default", "https://git/default.git") + repo = FakeScopeRepository([default]) + + monkeypatch.setattr( + GitPolicyFetcher, "make_bundle", lambda self, base_hash: _default_bundle() + ) + monkeypatch.setattr( + "opal_server.scopes.api.opal_server_config.BASE_DIR", str(tmp_path) + ) + + resp = _client(repo, tmp_path).get("/scopes/ghost/policy") + + assert resp.status_code == 200 + assert resp.json()["hash"] == "default-head" diff --git a/packages/opal-server/opal_server/tests/test_git_fetcher_repos_last_fetched.py b/packages/opal-server/opal_server/tests/test_git_fetcher_repos_last_fetched.py index fd12e7f35..e27e08ece 100644 --- a/packages/opal-server/opal_server/tests/test_git_fetcher_repos_last_fetched.py +++ b/packages/opal-server/opal_server/tests/test_git_fetcher_repos_last_fetched.py @@ -1,6 +1,7 @@ import datetime from pathlib import Path +import pygit2 import pytest from opal_common.schemas.policy_source import GitPolicyScopeSource, NoAuthData from opal_server.git_fetcher import GitPolicyFetcher @@ -88,3 +89,81 @@ async def test_force_fetch_not_downgraded_by_sibling_source(monkeypatch, tmp_pat ) is True ) + + +class _FailingRemote: + def fetch(self, *args, **kwargs): + raise pygit2.GitError("network down") + + +class _OkRemote: + def fetch(self, *args, **kwargs): + return None + + +class _FakeRepo: + def __init__(self, remote): + self.remotes = {"origin": remote} + + +@pytest.mark.asyncio +async def test_failed_fetch_does_not_stamp_last_fetched(monkeypatch, tmp_path): + """Bug B: stamping before the fetch means a FAILED fetch still records + "fetched now", which suppresses the next webhook-requested forced refresh + via _was_fetched_after().""" + fetcher = _make_fetcher(tmp_path, "scope_x", "https://example.com/repo-x.git") + monkeypatch.setattr(fetcher, "_discover_repository", lambda path: True) + monkeypatch.setattr(fetcher, "_get_valid_repo", lambda: _FakeRepo(_FailingRemote())) + + with pytest.raises(pygit2.GitError): + await fetcher.fetch_and_notify_on_changes(force_fetch=True) + + assert fetcher._source_id not in GitPolicyFetcher.repos_last_fetched + + +@pytest.mark.asyncio +async def test_successful_fetch_stamps_last_fetched(monkeypatch, tmp_path): + fetcher = _make_fetcher(tmp_path, "scope_x", "https://example.com/repo-x.git") + monkeypatch.setattr(fetcher, "_discover_repository", lambda path: True) + monkeypatch.setattr(fetcher, "_get_valid_repo", lambda: _FakeRepo(_OkRemote())) + + async def _no_notify(repo): + return None + + monkeypatch.setattr(fetcher, "_notify_on_changes", _no_notify) + await fetcher.fetch_and_notify_on_changes(force_fetch=True) + + assert fetcher._source_id in GitPolicyFetcher.repos_last_fetched + + +@pytest.mark.asyncio +async def test_stamp_records_fetch_start_time_not_completion(monkeypatch, tmp_path): + """The stamp must be the fetch START time: a fetch that STARTED after a + refresh's req_time already satisfies it; completion time would wrongly + suppress refreshes requested mid-fetch.""" + import time as _time + + class _SlowRemote: + def __init__(self): + self.entered_at = None + + def fetch(self, *args, **kwargs): + self.entered_at = datetime.datetime.now() + _time.sleep(0.05) + + remote = _SlowRemote() + fetcher = _make_fetcher(tmp_path, "scope_x", "https://example.com/repo-x.git") + monkeypatch.setattr(fetcher, "_discover_repository", lambda path: True) + monkeypatch.setattr(fetcher, "_get_valid_repo", lambda: _FakeRepo(remote)) + + async def _no_notify(repo): + return None + + monkeypatch.setattr(fetcher, "_notify_on_changes", _no_notify) + await fetcher.fetch_and_notify_on_changes(force_fetch=True) + + stamp = GitPolicyFetcher.repos_last_fetched[fetcher._source_id] + assert stamp <= remote.entered_at, ( + "stamp is later than the fetch's entry time — completion time was " + "stored instead of start time" + ) diff --git a/packages/opal-server/opal_server/tests/webhook_task_cleanup_test.py b/packages/opal-server/opal_server/tests/webhook_task_cleanup_test.py new file mode 100644 index 000000000..12d5284db --- /dev/null +++ b/packages/opal-server/opal_server/tests/webhook_task_cleanup_test.py @@ -0,0 +1,120 @@ +import asyncio + +import pytest +from opal_server.policy.watcher.task import BasePolicyWatcherTask + + +class _Watcher(BasePolicyWatcherTask): + async def trigger(self, topic, data): + return None # fast no-op so the created task finishes quickly + + +@pytest.mark.asyncio +async def test_done_tasks_are_all_removed(): + w = _Watcher(pubsub_endpoint=None) + + async def _done(): + return None + + # three already-finished tasks pre-loaded into the list + finished = [asyncio.create_task(_done()) for _ in range(3)] + await asyncio.gather(*finished) + w._webhook_tasks = list(finished) + + await w._on_webhook("webhook", None) + + # all 3 done ones removed... + remaining_done = [t for t in w._webhook_tasks if t in finished] + assert remaining_done == [], f"stale done tasks leaked: {remaining_done}" + # ...and exactly the one freshly scheduled trigger task survives. + survivors = [t for t in w._webhook_tasks if t not in finished] + assert len(w._webhook_tasks) == 1 + assert len(survivors) == 1, f"new trigger task not scheduled: {w._webhook_tasks}" + + # Await the trigger task directly (not sleep(0)) so the test doesn't depend + # on scheduler tick ordering. + await asyncio.gather(*survivors) + + +class _FailingWatcher(BasePolicyWatcherTask): + async def trigger(self, topic, data): + raise RuntimeError("trigger blew up") + + +@pytest.mark.asyncio +async def test_failed_trigger_exception_is_retrieved_and_logged(): + """Sweeping a failed trigger task must retrieve and log its exception, not + silently drop the reference (asyncio's 'exception was never retrieved').""" + from opal_common.logger import logger as opal_logger + + w = _FailingWatcher(pubsub_endpoint=None) + + await w._on_webhook("webhook", None) # schedules a trigger that raises + failed = list(w._webhook_tasks) + await asyncio.gather(*failed, return_exceptions=True) + + records = [] + sink_id = opal_logger.add(lambda m: records.append(str(m)), level="ERROR") + try: + await w._on_webhook("webhook", None) # sweep must log the failure + finally: + opal_logger.remove(sink_id) + + assert failed[0] not in w._webhook_tasks, "failed task not swept" + assert any( + "Webhook trigger task failed" in r and "trigger blew up" in r for r in records + ), f"failure not logged: {records}" + + await asyncio.gather(*w._webhook_tasks, return_exceptions=True) + + +class _CredentialLeakingWatcher(BasePolicyWatcherTask): + async def trigger(self, topic, data): + raise RuntimeError("fetch https://user:secret@host/repo failed") + + +@pytest.mark.asyncio +async def test_failed_trigger_log_redacts_credentialed_url(): + """A git exception can embed a credentialed remote URL verbatim; the + sweep's failure log must not leak it.""" + from opal_common.logger import logger as opal_logger + + w = _CredentialLeakingWatcher(pubsub_endpoint=None) + + await w._on_webhook("webhook", None) # schedules a trigger that raises + failed = list(w._webhook_tasks) + await asyncio.gather(*failed, return_exceptions=True) + + records = [] + sink_id = opal_logger.add(lambda m: records.append(str(m)), level="ERROR") + try: + await w._on_webhook("webhook", None) # sweep must log the failure + finally: + opal_logger.remove(sink_id) + + assert any("Webhook trigger task failed" in r for r in records), records + assert any("://***@" in r for r in records), records + assert not any("secret" in r for r in records), records + + await asyncio.gather(*w._webhook_tasks, return_exceptions=True) + + +class _HangingWatcher(BasePolicyWatcherTask): + def __init__(self, *args, **kwargs): + super().__init__(*args, **kwargs) + self.started = asyncio.Event() + + async def trigger(self, topic, data): + self.started.set() + await asyncio.Event().wait() # hangs until cancelled + + +@pytest.mark.asyncio +async def test_stop_cancels_and_gathers_inflight_trigger(): + w = _HangingWatcher(pubsub_endpoint=None) + await w._on_webhook("webhook", None) + await asyncio.wait_for(w.started.wait(), timeout=5) + + await w.stop() # must cancel the hung trigger AND await it (gather) + + assert all(t.cancelled() for t in w._webhook_tasks)