diff --git a/engine/skills/reflect/scripts/corpus_scan.py b/engine/skills/reflect/scripts/corpus_scan.py index aa966ef..d44ac4d 100644 --- a/engine/skills/reflect/scripts/corpus_scan.py +++ b/engine/skills/reflect/scripts/corpus_scan.py @@ -14,7 +14,7 @@ Usage: corpus_scan.py [--hours N] [--include-remote HOST,...|all] - [--pull-dir DIR] [--out FILE] + [--pull-dir DIR] [--out FILE] [--max-file-bytes N] corpus_scan.py "e2e|playwright|ci-regression" --hours 24 Local-only: lists matching files, no SSH, nothing pulled. @@ -73,6 +73,8 @@ INVOKER_CONFIG = os.path.expanduser("~/.invoker/config.json") SSH_KEY_DEFAULT = os.path.expanduser("~/.ssh/id_ed25519") +DEFAULT_MAX_FILE_BYTES = 64 * 1024 * 1024 +GREP_TIMEOUT_SECONDS = 15 def _grep_count(pattern, path): @@ -96,20 +98,41 @@ def _mtime_minutes(hours): return max(1, round(hours * 60)) -def _grep_match_files(root_glob_cmd, pattern, hours): +def _grep_match_files(root_glob_cmd, pattern, hours, max_file_bytes=DEFAULT_MAX_FILE_BYTES): """Local discovery: files under a find(1) expression, modified within - `hours`, whose content matches `pattern` (grep -l, not read-into-Python).""" + `hours`, whose content matches `pattern` (grep -l, not read-into-Python). + + One transcript must never take the whole scan down with it. A file + larger than `max_file_bytes` (None disables the cap) is skipped before + grep runs, and a grep that exceeds GREP_TIMEOUT_SECONDS is skipped + after; both print a stderr line naming the path so the skip is visible + in the progress log, never silent. A 146 MB Codex rollout once raised + TimeoutExpired straight out of this loop and killed the run.""" find_cmd = root_glob_cmd + ["-mmin", f"-{_mtime_minutes(hours)}"] try: found = subprocess.run(find_cmd, capture_output=True, text=True, timeout=30).stdout.splitlines() - except Exception: + except Exception as exc: + print(f"skip: find failed, returning no files for {shlex.join(find_cmd)}: {exc!r}", file=sys.stderr) return [] matched = [] for p in found: p = p.strip() if not p: continue - r = subprocess.run(["grep", "-l", "-i", "-E", pattern, p], capture_output=True, text=True, timeout=15) + if max_file_bytes is not None: + try: + size = os.path.getsize(p) + except OSError: + size = None + if size is not None and size > max_file_bytes: + print(f"skip: {p} is {size} bytes, above --max-file-bytes {max_file_bytes}", file=sys.stderr) + continue + try: + r = subprocess.run(["grep", "-l", "-i", "-E", pattern, p], capture_output=True, text=True, + timeout=GREP_TIMEOUT_SECONDS) + except subprocess.TimeoutExpired: + print(f"skip: grep timed out after {GREP_TIMEOUT_SECONDS}s on {p}", file=sys.stderr) + continue if r.returncode == 0: matched.append(p) return matched @@ -157,20 +180,21 @@ def split_sidechain(paths): return kept, dict(skipped) -def discover_local(pattern, hours, include_sidechain=False): +def discover_local(pattern, hours, include_sidechain=False, max_file_bytes=DEFAULT_MAX_FILE_BYTES): """Returns [(kind, path, 'local')] for Claude Code + Codex + Cursor sessions on this machine modified in the last `hours` hours whose content matches `pattern`. Cursor has no token fields — audit still records thrash/signal counts, never cost. Claude subagent (sidechain) transcripts are dropped and counted on stderr as - subagent_sessions_skipped unless include_sidechain=True.""" + subagent_sessions_skipped unless include_sidechain=True. Files above + `max_file_bytes` are skipped with a stderr line (see _grep_match_files).""" home = os.path.expanduser("~") claude_root = os.path.join(home, ".claude", "projects") codex_root = os.path.join(home, ".codex", "sessions") cursor_root = os.path.join(home, ".cursor", "projects") results = [] if os.path.isdir(claude_root): - claude_paths = _grep_match_files(["find", claude_root, "-iname", "*.jsonl"], pattern, hours) + claude_paths = _grep_match_files(["find", claude_root, "-iname", "*.jsonl"], pattern, hours, max_file_bytes) if not include_sidechain: claude_paths, skipped = split_sidechain(claude_paths) print( @@ -181,7 +205,7 @@ def discover_local(pattern, hours, include_sidechain=False): for p in claude_paths: results.append(("claude", p, "local")) if os.path.isdir(codex_root): - for p in _grep_match_files(["find", codex_root, "-iname", "rollout-*.jsonl"], pattern, hours): + for p in _grep_match_files(["find", codex_root, "-iname", "rollout-*.jsonl"], pattern, hours, max_file_bytes): results.append(("codex", p, "local")) if os.path.isdir(cursor_root): # ~/.cursor/projects//agent-transcripts//.jsonl @@ -189,6 +213,7 @@ def discover_local(pattern, hours, include_sidechain=False): ["find", cursor_root, "-path", "*/agent-transcripts/*/*.jsonl"], pattern, hours, + max_file_bytes, ): results.append(("cursor", p, "local")) return results @@ -339,6 +364,9 @@ def main(): ap.add_argument("--out", default="corpus_scan_results.json") ap.add_argument("--extra-signal", action="append", default=[], help="name=pattern, repeatable, added on top of DEFAULT_SIGNALS") + ap.add_argument("--max-file-bytes", type=int, default=DEFAULT_MAX_FILE_BYTES, + help="skip (with a stderr line) any local transcript larger than this many bytes; " + f"default {DEFAULT_MAX_FILE_BYTES} (64 MB)") args = ap.parse_args() signals = dict(DEFAULT_SIGNALS) @@ -348,7 +376,7 @@ def main(): signals[n] = p t0 = time.time() - files = [(k, p, h) for k, p, h in discover_local(args.pattern, args.hours)] + files = [(k, p, h) for k, p, h in discover_local(args.pattern, args.hours, max_file_bytes=args.max_file_bytes)] print(f"local: {len(files)} matching file(s) in the last {args.hours}h", file=sys.stderr) targets = load_remote_targets() diff --git a/engine/skills/reflect/scripts/tests/test_corpus_scan.py b/engine/skills/reflect/scripts/tests/test_corpus_scan.py index 575f210..4618bc6 100644 --- a/engine/skills/reflect/scripts/tests/test_corpus_scan.py +++ b/engine/skills/reflect/scripts/tests/test_corpus_scan.py @@ -12,6 +12,8 @@ that's the whole reason this script exists over token_audit.py/ top_sessions.py. """ +import contextlib +import io import os import sys import unittest @@ -135,5 +137,81 @@ def test_ignores_entries_without_timestamp(self): self.assertEqual(len(buckets), 0) +class TestGrepMatchFilesSkips(unittest.TestCase): + """_grep_match_files must survive one slow or oversized transcript: skip + it with an explicit stderr line naming the path, and still return every + other match. A 146 MB Codex rollout once raised TimeoutExpired out of + the per-file grep and killed the whole corpus scan.""" + + def setUp(self): + import tempfile + + self.tmp = tempfile.TemporaryDirectory() + self.fast = os.path.join(self.tmp.name, "fast.jsonl") + self.slow = os.path.join(self.tmp.name, "slow.jsonl") + self.other = os.path.join(self.tmp.name, "other.jsonl") + for path, body in ((self.fast, "hit\n"), (self.slow, "hit\n"), (self.other, "hit\n")): + with open(path, "w", encoding="utf-8") as handle: + handle.write(body) + self.real_run = corpus_scan.subprocess.run + + def tearDown(self): + corpus_scan.subprocess.run = self.real_run + self.tmp.cleanup() + + def _fake_run_factory(self, timeout_on): + import subprocess as sp + + listed = "\n".join([self.fast, self.slow, self.other]) + "\n" + + def fake_run(cmd, **kwargs): + if cmd[0] == "find": + return sp.CompletedProcess(cmd, 0, stdout=listed, stderr="") + if cmd[-1] == timeout_on: + raise sp.TimeoutExpired(cmd, kwargs.get("timeout")) + return sp.CompletedProcess(cmd, 0, stdout=cmd[-1] + "\n", stderr="") + + return fake_run + + def _run(self, **kwargs): + err = io.StringIO() + with contextlib.redirect_stderr(err): + matched = corpus_scan._grep_match_files(["find", self.tmp.name], "hit", 24, **kwargs) + return matched, err.getvalue() + + def test_timeout_on_one_file_skips_it_and_keeps_the_rest(self): + corpus_scan.subprocess.run = self._fake_run_factory(timeout_on=self.slow) + matched, err = self._run() + self.assertEqual(matched, [self.fast, self.other]) + self.assertIn(self.slow, err) + self.assertIn("timed out", err) + self.assertIn(str(corpus_scan.GREP_TIMEOUT_SECONDS), err) + + def test_file_above_max_bytes_is_skipped_with_stderr_line(self): + corpus_scan.subprocess.run = self._fake_run_factory(timeout_on=None) + with open(self.slow, "w", encoding="utf-8") as handle: + handle.write("x" * 100) + matched, err = self._run(max_file_bytes=50) + self.assertEqual(matched, [self.fast, self.other]) + self.assertIn(self.slow, err) + self.assertIn("max-file-bytes", err) + self.assertIn("50", err) + + def test_default_cap_is_64_mb(self): + self.assertEqual(corpus_scan.DEFAULT_MAX_FILE_BYTES, 64 * 1024 * 1024) + + def test_find_failure_logs_to_stderr_and_returns_empty(self): + import subprocess as sp + + def boom(cmd, **kwargs): + raise sp.TimeoutExpired(cmd, kwargs.get("timeout")) + + corpus_scan.subprocess.run = boom + matched, err = self._run() + self.assertEqual(matched, []) + self.assertIn("find", err) + self.assertIn(self.tmp.name, err) + + if __name__ == "__main__": unittest.main()