Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
48 changes: 38 additions & 10 deletions engine/skills/reflect/scripts/corpus_scan.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@

Usage:
corpus_scan.py <keyword-regex> [--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.
Expand Down Expand Up @@ -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):
Expand All @@ -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
Expand Down Expand Up @@ -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(
Expand All @@ -181,14 +205,15 @@ 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/<project>/agent-transcripts/<uuid>/<uuid>.jsonl
for p in _grep_match_files(
["find", cursor_root, "-path", "*/agent-transcripts/*/*.jsonl"],
pattern,
hours,
max_file_bytes,
):
results.append(("cursor", p, "local"))
return results
Expand Down Expand Up @@ -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)
Expand All @@ -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()
Expand Down
78 changes: 78 additions & 0 deletions engine/skills/reflect/scripts/tests/test_corpus_scan.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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()
Loading