fix hook OOM, lift worker cap, fix antigravity .agents path and frontmatter
#791: per-repo fcntl flock in _rebuild_code prevents concurrent hook rebuilds from exhausting memory; changed_paths wired through so only modified files are re-extracted; stale nodes evicted on deletion; SIGALRM watchdog with GRAPHIFY_REBUILD_TIMEOUT; Darwin-aware RLIMIT_DATA memory cap #792: remove hard 8-worker cap (GRAPHIFY_MAX_WORKERS env var); add --max-workers, --token-budget, --max-concurrency, --api-timeout CLI flags to graphify extract; fix ollama API key gate for loopback URLs; explicit timeout on OpenAI client (GRAPHIFY_API_TIMEOUT, default 600s); per-chunk progress prints during extraction #453 + #785: rename .agent -> .agents throughout antigravity install/uninstall; add trigger:always_on YAML frontmatter to _ANTIGRAVITY_RULES so Antigravity recognises the rules file Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Sonnet 4.6
parent
8e6483511d
commit
dc50979a68
+128
-21
@@ -142,7 +142,7 @@ _PLATFORM_CONFIG: dict[str, dict] = {
|
||||
},
|
||||
"antigravity": {
|
||||
"skill_file": "skill.md",
|
||||
"skill_dst": Path.home() / ".agent" / "skills" / "graphify" / "SKILL.md",
|
||||
"skill_dst": Path(".agents") / "skills" / "graphify" / "SKILL.md",
|
||||
"claude_md": False,
|
||||
},
|
||||
"windows": {
|
||||
@@ -474,10 +474,15 @@ def vscode_uninstall(project_dir: Path | None = None) -> None:
|
||||
print(f" {instructions} -> deleted (was empty after removal)")
|
||||
|
||||
|
||||
_ANTIGRAVITY_RULES_PATH = Path(".agent") / "rules" / "graphify.md"
|
||||
_ANTIGRAVITY_WORKFLOW_PATH = Path(".agent") / "workflows" / "graphify.md"
|
||||
_ANTIGRAVITY_RULES_PATH = Path(".agents") / "rules" / "graphify.md"
|
||||
_ANTIGRAVITY_WORKFLOW_PATH = Path(".agents") / "workflows" / "graphify.md"
|
||||
|
||||
_ANTIGRAVITY_RULES = """\
|
||||
---
|
||||
trigger: always_on
|
||||
description: Always consult the graphify knowledge graph at graphify-out/ before answering codebase or architecture questions.
|
||||
---
|
||||
|
||||
## graphify
|
||||
|
||||
This project has a graphify knowledge graph at graphify-out/.
|
||||
@@ -486,7 +491,7 @@ Rules:
|
||||
- Before answering architecture or codebase questions, read graphify-out/GRAPH_REPORT.md for god nodes and community structure
|
||||
- If graphify-out/wiki/index.md exists, navigate it instead of reading raw files
|
||||
- If the graphify MCP server is active, utilize tools like `query_graph`, `get_node`, and `shortest_path` for precise architecture navigation instead of falling back to `grep`
|
||||
- If the MCP server is not active, the CLI equivalents are `graphify query "<question>"`, `graphify path "<A>" "<B>"`, and `graphify explain "<concept>"` — prefer these over grep for cross-module questions
|
||||
- If the MCP server is not active, the CLI equivalents are `graphify query "<question>"`, `graphify path "<A>" "<B>"`, and `graphify explain "<concept>"` - prefer these over grep for cross-module questions
|
||||
- After modifying code files in this session, run `graphify update .` to keep the graph current (AST-only, no API cost)
|
||||
"""
|
||||
|
||||
@@ -498,7 +503,7 @@ description: Turn any folder of files into a navigable knowledge graph
|
||||
|
||||
# Workflow: graphify
|
||||
|
||||
Follow the graphify skill installed at ~/.agent/skills/graphify/SKILL.md to run the full pipeline.
|
||||
Follow the graphify skill installed at ~/.agents/skills/graphify/SKILL.md to run the full pipeline.
|
||||
|
||||
If no path argument is given, use `.` (current directory).
|
||||
"""
|
||||
@@ -573,7 +578,7 @@ def _antigravity_install(project_dir: Path) -> None:
|
||||
install(platform="antigravity")
|
||||
|
||||
# 1.5. Inject YAML frontmatter for native Antigravity tool discovery
|
||||
skill_dst = Path.home() / _PLATFORM_CONFIG["antigravity"]["skill_dst"]
|
||||
skill_dst = _PLATFORM_CONFIG["antigravity"]["skill_dst"]
|
||||
if skill_dst.exists():
|
||||
content = skill_dst.read_text(encoding="utf-8")
|
||||
if not content.startswith("---\n"):
|
||||
@@ -636,7 +641,7 @@ def _antigravity_uninstall(project_dir: Path) -> None:
|
||||
print(f"graphify workflow removed from {wf_path.resolve()}")
|
||||
|
||||
# Remove skill file
|
||||
skill_dst = Path.home() / _PLATFORM_CONFIG["antigravity"]["skill_dst"]
|
||||
skill_dst = _PLATFORM_CONFIG["antigravity"]["skill_dst"]
|
||||
if skill_dst.exists():
|
||||
skill_dst.unlink()
|
||||
print(f"graphify skill removed from {skill_dst}")
|
||||
@@ -1165,6 +1170,10 @@ def main() -> None:
|
||||
print(" extract <path> headless full extraction (AST + semantic LLM) for CI/scripts")
|
||||
print(" --backend B gemini|kimi|claude|openai|ollama (default: whichever API key is set)")
|
||||
print(" --model M override backend default model")
|
||||
print(" --max-workers N AST extraction subprocess count (default: cpu_count)")
|
||||
print(" --token-budget N per-chunk token cap for semantic extraction (default: 60000)")
|
||||
print(" --max-concurrency N parallel semantic chunks in flight (default: 4; set 1 for local LLMs)")
|
||||
print(" --api-timeout S per-request timeout in seconds for the LLM client (default: 600)")
|
||||
print(" --out DIR output dir (default: <path>); writes <DIR>/graphify-out/")
|
||||
print(" --google-workspace export .gdoc/.gsheet/.gslides shortcuts via gws before extraction")
|
||||
print(" --no-cluster skip clustering, write raw extraction only")
|
||||
@@ -1203,8 +1212,8 @@ def main() -> None:
|
||||
print(" trae uninstall remove graphify section from AGENTS.md")
|
||||
print(" trae-cn install write graphify section to AGENTS.md (Trae CN)")
|
||||
print(" trae-cn uninstall remove graphify section from AGENTS.md")
|
||||
print(" antigravity install write .agent/rules + .agent/workflows + skill (Google Antigravity)")
|
||||
print(" antigravity uninstall remove .agent/rules, .agent/workflows, and skill")
|
||||
print(" antigravity install write .agents/rules + .agents/workflows + skill (Google Antigravity)")
|
||||
print(" antigravity uninstall remove .agents/rules, .agents/workflows, and skill")
|
||||
print(" hermes install write skill to ~/.hermes/skills/graphify/ (Hermes)")
|
||||
print(" hermes uninstall remove skill from ~/.hermes/skills/graphify/")
|
||||
print(" kiro install write skill to .kiro/skills/graphify/ + steering file (Kiro IDE/CLI)")
|
||||
@@ -1706,7 +1715,10 @@ def main() -> None:
|
||||
sys.exit(1)
|
||||
from graphify.watch import _rebuild_code
|
||||
print(f"Re-extracting code files in {watch_path} (no LLM needed)...")
|
||||
ok = _rebuild_code(watch_path, force=force)
|
||||
# Interactive CLI: block on the per-repo lock rather than skip, so the
|
||||
# user sees their explicit `graphify update` complete instead of
|
||||
# exiting silently when a hook-driven rebuild happens to be running.
|
||||
ok = _rebuild_code(watch_path, force=force, block_on_lock=True)
|
||||
if ok:
|
||||
print("Code graph updated. For doc/paper/image changes run /graphify --update in your AI assistant.")
|
||||
if not (
|
||||
@@ -2132,8 +2144,10 @@ def main() -> None:
|
||||
# has an API key set.
|
||||
if len(sys.argv) < 3:
|
||||
print(
|
||||
"Usage: graphify extract <path> [--backend gemini|kimi|claude|openai] "
|
||||
"[--out DIR] [--google-workspace] [--no-cluster]",
|
||||
"Usage: graphify extract <path> [--backend gemini|kimi|claude|openai|ollama] "
|
||||
"[--model M] [--out DIR] [--google-workspace] [--no-cluster] "
|
||||
"[--max-workers N] [--token-budget N] [--max-concurrency N] "
|
||||
"[--api-timeout S]",
|
||||
file=sys.stderr,
|
||||
)
|
||||
sys.exit(1)
|
||||
@@ -2151,6 +2165,34 @@ def main() -> None:
|
||||
google_workspace = False
|
||||
global_merge = False
|
||||
global_repo_tag: str | None = None
|
||||
# Performance/tuning knobs (issue #792). None means "use library default".
|
||||
cli_max_workers: int | None = None
|
||||
cli_token_budget: int | None = None
|
||||
cli_max_concurrency: int | None = None
|
||||
cli_api_timeout: float | None = None
|
||||
|
||||
def _parse_int(name: str, raw: str) -> int:
|
||||
try:
|
||||
v = int(raw)
|
||||
except ValueError:
|
||||
print(f"error: {name} must be a positive integer (got {raw!r})", file=sys.stderr)
|
||||
sys.exit(2)
|
||||
if v <= 0:
|
||||
print(f"error: {name} must be > 0 (got {v})", file=sys.stderr)
|
||||
sys.exit(2)
|
||||
return v
|
||||
|
||||
def _parse_float(name: str, raw: str) -> float:
|
||||
try:
|
||||
v = float(raw)
|
||||
except ValueError:
|
||||
print(f"error: {name} must be a positive number (got {raw!r})", file=sys.stderr)
|
||||
sys.exit(2)
|
||||
if v <= 0:
|
||||
print(f"error: {name} must be > 0 (got {v})", file=sys.stderr)
|
||||
sys.exit(2)
|
||||
return v
|
||||
|
||||
args = sys.argv[3:]
|
||||
i = 0
|
||||
while i < len(args):
|
||||
@@ -2177,9 +2219,32 @@ def main() -> None:
|
||||
global_merge = True; i += 1
|
||||
elif a == "--as" and i + 1 < len(args):
|
||||
global_repo_tag = args[i + 1]; i += 2
|
||||
elif a == "--max-workers" and i + 1 < len(args):
|
||||
cli_max_workers = _parse_int("--max-workers", args[i + 1]); i += 2
|
||||
elif a.startswith("--max-workers="):
|
||||
cli_max_workers = _parse_int("--max-workers", a.split("=", 1)[1]); i += 1
|
||||
elif a == "--token-budget" and i + 1 < len(args):
|
||||
cli_token_budget = _parse_int("--token-budget", args[i + 1]); i += 2
|
||||
elif a.startswith("--token-budget="):
|
||||
cli_token_budget = _parse_int("--token-budget", a.split("=", 1)[1]); i += 1
|
||||
elif a == "--max-concurrency" and i + 1 < len(args):
|
||||
cli_max_concurrency = _parse_int("--max-concurrency", args[i + 1]); i += 2
|
||||
elif a.startswith("--max-concurrency="):
|
||||
cli_max_concurrency = _parse_int("--max-concurrency", a.split("=", 1)[1]); i += 1
|
||||
elif a == "--api-timeout" and i + 1 < len(args):
|
||||
cli_api_timeout = _parse_float("--api-timeout", args[i + 1]); i += 2
|
||||
elif a.startswith("--api-timeout="):
|
||||
cli_api_timeout = _parse_float("--api-timeout", a.split("=", 1)[1]); i += 1
|
||||
else:
|
||||
i += 1
|
||||
|
||||
# CLI flag wins over env var. Setting GRAPHIFY_API_TIMEOUT here so
|
||||
# _call_openai_compat picks it up without needing a new kwarg path.
|
||||
if cli_api_timeout is not None:
|
||||
os.environ["GRAPHIFY_API_TIMEOUT"] = str(cli_api_timeout)
|
||||
if cli_max_workers is not None:
|
||||
os.environ["GRAPHIFY_MAX_WORKERS"] = str(cli_max_workers)
|
||||
|
||||
# Backend resolution. If user did not pass --backend, sniff env.
|
||||
# If backend was explicitly requested, validate its key is present
|
||||
# and surface a clear error early — don't let extract_corpus_parallel
|
||||
@@ -2210,11 +2275,31 @@ def main() -> None:
|
||||
)
|
||||
sys.exit(1)
|
||||
if not _get_backend_api_key(backend):
|
||||
print(
|
||||
f"error: backend '{backend}' requires {_format_backend_env_keys(backend)} to be set.",
|
||||
file=sys.stderr,
|
||||
)
|
||||
sys.exit(1)
|
||||
# Ollama on a loopback URL ignores auth entirely; don't block
|
||||
# the run just because OLLAMA_API_KEY is unset (issue #792).
|
||||
# extract_files_direct already prints a warning and substitutes
|
||||
# a placeholder key in that case.
|
||||
allow_no_key = False
|
||||
if backend == "ollama":
|
||||
from urllib.parse import urlparse
|
||||
ollama_url = os.environ.get(
|
||||
"OLLAMA_BASE_URL",
|
||||
_BACKENDS["ollama"].get("base_url", ""),
|
||||
)
|
||||
try:
|
||||
host = (urlparse(ollama_url).hostname or "").lower()
|
||||
except Exception:
|
||||
host = ""
|
||||
allow_no_key = (
|
||||
host in ("localhost", "127.0.0.1", "::1")
|
||||
or host.startswith("127.")
|
||||
)
|
||||
if not allow_no_key:
|
||||
print(
|
||||
f"error: backend '{backend}' requires {_format_backend_env_keys(backend)} to be set.",
|
||||
file=sys.stderr,
|
||||
)
|
||||
sys.exit(1)
|
||||
|
||||
# Resolve output dir. The user-facing contract is "<out>/graphify-out/"
|
||||
# so a fresh checkout writes graphify-out/ at the project root, matching
|
||||
@@ -2279,9 +2364,12 @@ def main() -> None:
|
||||
ast_result: dict = {"nodes": [], "edges": [], "input_tokens": 0, "output_tokens": 0}
|
||||
if code_files:
|
||||
from graphify.extract import extract as _ast_extract
|
||||
ast_kwargs: dict = {"cache_root": target}
|
||||
if cli_max_workers is not None:
|
||||
ast_kwargs["max_workers"] = cli_max_workers
|
||||
print(f"[graphify extract] AST extraction on {len(code_files)} code files...")
|
||||
try:
|
||||
ast_result = _ast_extract(code_files, cache_root=target)
|
||||
ast_result = _ast_extract(code_files, **ast_kwargs)
|
||||
except Exception as exc:
|
||||
print(f"[graphify extract] AST extraction failed: {exc}", file=sys.stderr)
|
||||
ast_result = {"nodes": [], "edges": [], "input_tokens": 0, "output_tokens": 0}
|
||||
@@ -2312,12 +2400,31 @@ def main() -> None:
|
||||
|
||||
if uncached_paths:
|
||||
print(f"[graphify extract] semantic extraction on {len(uncached_paths)} files via {backend}...")
|
||||
corpus_kwargs: dict = {
|
||||
"backend": backend,
|
||||
"model": model,
|
||||
"root": target,
|
||||
}
|
||||
if cli_token_budget is not None:
|
||||
corpus_kwargs["token_budget"] = cli_token_budget
|
||||
if cli_max_concurrency is not None:
|
||||
corpus_kwargs["max_concurrency"] = cli_max_concurrency
|
||||
|
||||
# Minimal progress callback so the CLI is no longer silent
|
||||
# during long local-inference runs (issue #792 addendum).
|
||||
_total_chunks = {"n": 0}
|
||||
def _progress(idx: int, total: int, _result: dict) -> None:
|
||||
_total_chunks["n"] = total
|
||||
print(
|
||||
f"[graphify extract] chunk {idx + 1}/{total} done",
|
||||
flush=True,
|
||||
)
|
||||
corpus_kwargs["on_chunk_done"] = _progress
|
||||
|
||||
try:
|
||||
fresh = _extract_corpus_parallel(
|
||||
[Path(p) for p in uncached_paths],
|
||||
backend=backend,
|
||||
model=model,
|
||||
root=target,
|
||||
**corpus_kwargs,
|
||||
)
|
||||
except ImportError as exc:
|
||||
print(f"error: {exc}", file=sys.stderr)
|
||||
|
||||
+18
-2
@@ -5555,7 +5555,22 @@ def _extract_parallel(
|
||||
import concurrent.futures
|
||||
|
||||
if max_workers is None:
|
||||
max_workers = min(os.cpu_count() or 4, len(uncached_work), 8)
|
||||
# Honour GRAPHIFY_MAX_WORKERS env override; otherwise scale to the
|
||||
# full CPU. The historical `, 8)` cap was a safety bound for laptops
|
||||
# in 2023 — on a 32-thread workstation it costs a 4x slowdown
|
||||
# (issue #792). Capping at len(uncached_work) keeps small jobs
|
||||
# from spawning useless idle workers.
|
||||
env_raw = os.environ.get("GRAPHIFY_MAX_WORKERS", "").strip()
|
||||
env_cap = None
|
||||
if env_raw:
|
||||
try:
|
||||
v = int(env_raw)
|
||||
if v > 0:
|
||||
env_cap = v
|
||||
except ValueError:
|
||||
pass
|
||||
cpu_cap = env_cap if env_cap is not None else (os.cpu_count() or 4)
|
||||
max_workers = min(cpu_cap, len(uncached_work))
|
||||
|
||||
root_str = str(effective_root)
|
||||
work_items = [(idx, str(path), root_str) for idx, path in uncached_work]
|
||||
@@ -5656,7 +5671,8 @@ def extract(
|
||||
subdirectory so the cache stays at ./graphify-out/cache/.
|
||||
parallel: if True and there are >= _PARALLEL_THRESHOLD uncached files,
|
||||
use ProcessPoolExecutor for multi-core extraction.
|
||||
max_workers: max subprocess count. Defaults to min(cpu_count, 8).
|
||||
max_workers: max subprocess count. Defaults to cpu_count (or the
|
||||
value of GRAPHIFY_MAX_WORKERS if set), bounded by len(uncached_work).
|
||||
"""
|
||||
_check_tree_sitter_version()
|
||||
_raise_recursion_limit()
|
||||
|
||||
+25
-7
@@ -69,7 +69,7 @@ _GRAPHIFY_LOG="${HOME}/.cache/graphify-rebuild.log"
|
||||
mkdir -p "$(dirname "$_GRAPHIFY_LOG")"
|
||||
echo "[graphify hook] launching background rebuild (log: $_GRAPHIFY_LOG)"
|
||||
nohup $GRAPHIFY_PYTHON -c "
|
||||
import os, sys
|
||||
import os, signal, sys
|
||||
from pathlib import Path
|
||||
|
||||
changed_raw = os.environ.get('GRAPHIFY_CHANGED', '')
|
||||
@@ -81,10 +81,17 @@ if not changed:
|
||||
print(f'[graphify hook] {len(changed)} file(s) changed - rebuilding graph...')
|
||||
|
||||
try:
|
||||
import os as _os
|
||||
from graphify.watch import _rebuild_code
|
||||
_force = _os.environ.get('GRAPHIFY_FORCE', '').lower() in ('1', 'true', 'yes')
|
||||
_rebuild_code(Path('.'), force=_force)
|
||||
from graphify.watch import _rebuild_code, _apply_resource_limits
|
||||
_apply_resource_limits()
|
||||
_timeout = int(os.environ.get('GRAPHIFY_REBUILD_TIMEOUT', '600'))
|
||||
if _timeout > 0 and hasattr(signal, 'SIGALRM'):
|
||||
signal.signal(signal.SIGALRM, lambda *_: (_ for _ in ()).throw(TimeoutError(f'graphify rebuild exceeded {_timeout}s')))
|
||||
signal.alarm(_timeout)
|
||||
_force = os.environ.get('GRAPHIFY_FORCE', '').lower() in ('1', 'true', 'yes')
|
||||
_rebuild_code(Path('.'), changed_paths=changed, force=_force)
|
||||
except TimeoutError as exc:
|
||||
print(f'[graphify hook] {exc}')
|
||||
sys.exit(1)
|
||||
except Exception as exc:
|
||||
print(f'[graphify hook] Rebuild failed: {exc}')
|
||||
sys.exit(1)
|
||||
@@ -125,12 +132,23 @@ _GRAPHIFY_LOG="${HOME}/.cache/graphify-rebuild.log"
|
||||
mkdir -p "$(dirname "$_GRAPHIFY_LOG")"
|
||||
echo "[graphify] Branch switched - launching background rebuild (log: $_GRAPHIFY_LOG)"
|
||||
nohup $GRAPHIFY_PYTHON -c "
|
||||
from graphify.watch import _rebuild_code
|
||||
from graphify.watch import _rebuild_code, _apply_resource_limits
|
||||
from pathlib import Path
|
||||
import os, sys
|
||||
import os, signal, sys
|
||||
try:
|
||||
_apply_resource_limits()
|
||||
_timeout = int(os.environ.get('GRAPHIFY_REBUILD_TIMEOUT', '600'))
|
||||
if _timeout > 0 and hasattr(signal, 'SIGALRM'):
|
||||
signal.signal(signal.SIGALRM, lambda *_: (_ for _ in ()).throw(TimeoutError(f'graphify rebuild exceeded {_timeout}s')))
|
||||
signal.alarm(_timeout)
|
||||
_force = os.environ.get('GRAPHIFY_FORCE', '').lower() in ('1', 'true', 'yes')
|
||||
# post-checkout: branch switch can touch arbitrary files; full rebuild path
|
||||
# (no changed_paths) is correct here. The flock inside _rebuild_code still
|
||||
# prevents pile-ups when commit + checkout fire back-to-back.
|
||||
_rebuild_code(Path('.'), force=_force)
|
||||
except TimeoutError as exc:
|
||||
print(f'[graphify] {exc}')
|
||||
sys.exit(1)
|
||||
except Exception as exc:
|
||||
print(f'[graphify] Rebuild failed: {exc}')
|
||||
sys.exit(1)
|
||||
|
||||
+15
-1
@@ -229,7 +229,21 @@ def _call_openai_compat(
|
||||
f"Run: pip install {pkg_hint}"
|
||||
) from exc
|
||||
|
||||
client = OpenAI(api_key=api_key, base_url=base_url)
|
||||
# Local backends (ollama, llama.cpp, vLLM) routinely take >60s for a
|
||||
# single chunk on a large model — far longer than the openai SDK's
|
||||
# default. Honour GRAPHIFY_API_TIMEOUT (seconds) for explicit override;
|
||||
# default to 600s, which is long enough for a 31B model on a 16k chunk
|
||||
# but still bounds runaway connections (issue #792 addendum).
|
||||
timeout_raw = os.environ.get("GRAPHIFY_API_TIMEOUT", "").strip()
|
||||
timeout_s: float = 600.0
|
||||
if timeout_raw:
|
||||
try:
|
||||
v = float(timeout_raw)
|
||||
if v > 0:
|
||||
timeout_s = v
|
||||
except ValueError:
|
||||
pass
|
||||
client = OpenAI(api_key=api_key, base_url=base_url, timeout=timeout_s)
|
||||
kwargs: dict = {
|
||||
"model": model,
|
||||
"messages": [
|
||||
|
||||
+152
-7
@@ -1,5 +1,6 @@
|
||||
# monitor a folder and auto-trigger --update when files change
|
||||
from __future__ import annotations
|
||||
import contextlib
|
||||
import json
|
||||
import os
|
||||
import sys
|
||||
@@ -9,6 +10,74 @@ from pathlib import Path
|
||||
_GRAPHIFY_OUT = os.environ.get("GRAPHIFY_OUT", "graphify-out")
|
||||
|
||||
|
||||
@contextlib.contextmanager
|
||||
def _rebuild_lock(out_dir: Path, *, blocking: bool = False):
|
||||
"""Per-repo advisory lock around a rebuild.
|
||||
|
||||
Yields True if acquired, False if another rebuild is already running and
|
||||
``blocking`` is False. Uses fcntl.flock so the lock is released
|
||||
automatically if the process is killed (no stale-lock cleanup needed).
|
||||
|
||||
Falls back to a no-op yield(True) on platforms without fcntl (Windows).
|
||||
"""
|
||||
try:
|
||||
import fcntl
|
||||
except ImportError:
|
||||
yield True
|
||||
return
|
||||
|
||||
out_dir.mkdir(parents=True, exist_ok=True)
|
||||
lock_path = out_dir / ".rebuild.lock"
|
||||
fh = open(lock_path, "a", encoding="utf-8")
|
||||
try:
|
||||
flags = fcntl.LOCK_EX if blocking else (fcntl.LOCK_EX | fcntl.LOCK_NB)
|
||||
try:
|
||||
fcntl.flock(fh.fileno(), flags)
|
||||
except BlockingIOError:
|
||||
yield False
|
||||
return
|
||||
try:
|
||||
fh.write(str(os.getpid()))
|
||||
fh.flush()
|
||||
except OSError:
|
||||
pass
|
||||
yield True
|
||||
finally:
|
||||
try:
|
||||
fcntl.flock(fh.fileno(), fcntl.LOCK_UN)
|
||||
except OSError:
|
||||
pass
|
||||
fh.close()
|
||||
|
||||
|
||||
def _apply_resource_limits() -> None:
|
||||
"""Best-effort nice + memory cap. Called from inline hook scripts.
|
||||
|
||||
GRAPHIFY_REBUILD_MEMORY_LIMIT_MB caps RSS-ish memory. Uses RLIMIT_DATA on
|
||||
macOS (RLIMIT_AS is unreliable under Apple's libmalloc) and RLIMIT_AS on
|
||||
Linux. Silently skips if the platform doesn't support it.
|
||||
"""
|
||||
try:
|
||||
os.nice(10)
|
||||
except (OSError, AttributeError):
|
||||
pass
|
||||
mb = os.environ.get("GRAPHIFY_REBUILD_MEMORY_LIMIT_MB", "").strip()
|
||||
if not mb:
|
||||
return
|
||||
try:
|
||||
limit = int(mb) * 1024 * 1024
|
||||
except ValueError:
|
||||
return
|
||||
try:
|
||||
import resource
|
||||
which = resource.RLIMIT_DATA if sys.platform == "darwin" else resource.RLIMIT_AS
|
||||
soft, hard = resource.getrlimit(which)
|
||||
new_hard = hard if hard != resource.RLIM_INFINITY and hard < limit else limit
|
||||
resource.setrlimit(which, (limit, new_hard))
|
||||
except (ImportError, ValueError, OSError):
|
||||
pass
|
||||
|
||||
|
||||
def _git_head() -> str | None:
|
||||
"""Return current git HEAD commit hash, or None outside a repo."""
|
||||
import subprocess as _sp
|
||||
@@ -46,20 +115,55 @@ def _relativize_source_files(payload: dict, root: Path) -> None:
|
||||
continue
|
||||
|
||||
|
||||
def _rebuild_code(watch_path: Path, *, follow_symlinks: bool = False, force: bool = False) -> bool:
|
||||
def _rebuild_code(
|
||||
watch_path: Path,
|
||||
*,
|
||||
changed_paths: list[Path] | None = None,
|
||||
follow_symlinks: bool = False,
|
||||
force: bool = False,
|
||||
acquire_lock: bool = True,
|
||||
block_on_lock: bool = False,
|
||||
) -> bool:
|
||||
"""Re-run AST extraction + build + cluster + report for code files. No LLM needed.
|
||||
|
||||
When ``force`` is True the node-count safety check in ``to_json`` is bypassed
|
||||
so the rebuilt graph overwrites graph.json even if it has fewer nodes.
|
||||
Use this after refactors that legitimately delete code.
|
||||
|
||||
Returns True on success, False on error.
|
||||
When ``changed_paths`` is provided, only those files are re-extracted; nodes
|
||||
for unchanged files are preserved from the existing graph. Deleted paths
|
||||
in ``changed_paths`` (paths that no longer exist on disk) are dropped from
|
||||
the preserved set. When ``changed_paths`` is None the full code corpus is
|
||||
re-extracted (used by the watcher and post-checkout hook).
|
||||
|
||||
``acquire_lock`` (default True) takes a non-blocking per-repo flock around
|
||||
the rebuild so concurrent post-commit hooks across multiple repos do not
|
||||
pile up. Returns False with a log line if the lock is held. Pass
|
||||
``block_on_lock=True`` to wait instead of skip (used by the interactive
|
||||
``graphify update`` CLI).
|
||||
|
||||
Returns True on success, False on error or skipped-due-to-lock.
|
||||
"""
|
||||
out = watch_path / _GRAPHIFY_OUT
|
||||
if acquire_lock:
|
||||
with _rebuild_lock(out, blocking=block_on_lock) as got:
|
||||
if not got:
|
||||
print("[graphify watch] Rebuild already in progress for "
|
||||
f"{watch_path.resolve()} - skipping.")
|
||||
return False
|
||||
return _rebuild_code(
|
||||
watch_path,
|
||||
changed_paths=changed_paths,
|
||||
follow_symlinks=follow_symlinks,
|
||||
force=force,
|
||||
acquire_lock=False,
|
||||
)
|
||||
|
||||
watch_root = watch_path.resolve()
|
||||
project_root = Path.cwd().resolve() if not watch_path.is_absolute() else watch_root
|
||||
report_root = _report_root_label(watch_path)
|
||||
try:
|
||||
from graphify.extract import extract
|
||||
from graphify.extract import extract, _get_extractor
|
||||
from graphify.detect import detect
|
||||
from graphify.build import build_from_json
|
||||
from graphify.cluster import cluster, score_all
|
||||
@@ -71,7 +175,6 @@ def _rebuild_code(watch_path: Path, *, follow_symlinks: bool = False, force: boo
|
||||
code_files = [Path(f) for f in detected['files']['code']]
|
||||
|
||||
# Include document files that have AST extractors (e.g. .md, .mdx, .qmd)
|
||||
from graphify.extract import _get_extractor
|
||||
for doc_file in detected['files'].get('document', []):
|
||||
p = Path(doc_file)
|
||||
if _get_extractor(p) is not None:
|
||||
@@ -81,21 +184,63 @@ def _rebuild_code(watch_path: Path, *, follow_symlinks: bool = False, force: boo
|
||||
print("[graphify watch] No code files found - nothing to rebuild.")
|
||||
return False
|
||||
|
||||
# Incremental path: when the caller passed an explicit change list,
|
||||
# extract only changed-and-still-existing files. Deleted paths are
|
||||
# tracked separately so their stale nodes can be evicted below.
|
||||
deleted_paths: set[str] = set()
|
||||
if changed_paths is not None:
|
||||
code_set = {p.resolve() for p in code_files}
|
||||
wanted: list[Path] = []
|
||||
for raw in changed_paths:
|
||||
cand = (watch_root / raw).resolve() if not raw.is_absolute() else raw.resolve()
|
||||
if cand.exists() and cand in code_set:
|
||||
wanted.append(cand)
|
||||
else:
|
||||
# File was deleted, renamed away, or filtered out by detect
|
||||
# (e.g. .gitignore, vendored). Either way, evict any
|
||||
# preserved nodes that still claim this source path.
|
||||
try:
|
||||
deleted_paths.add(str(cand.relative_to(project_root)))
|
||||
except ValueError:
|
||||
deleted_paths.add(str(cand))
|
||||
if not wanted and not deleted_paths:
|
||||
print("[graphify watch] No tracked code files in change set - skipping rebuild.")
|
||||
return True
|
||||
extract_targets = wanted
|
||||
else:
|
||||
extract_targets = code_files
|
||||
|
||||
commit = _git_head()
|
||||
result = extract(code_files, cache_root=watch_root)
|
||||
result = extract(extract_targets, cache_root=watch_root) if extract_targets else {
|
||||
"nodes": [], "edges": [], "hyperedges": [],
|
||||
"input_tokens": 0, "output_tokens": 0,
|
||||
}
|
||||
|
||||
# Preserve semantic nodes/edges from a previous full run.
|
||||
# AST-only rebuild replaces nodes for changed files; everything else is kept.
|
||||
# Filter by node ID membership in the new AST output, not by file_type —
|
||||
# INFERRED/AMBIGUOUS nodes extracted from code files also carry file_type="code"
|
||||
# and would be wrongly dropped by a file_type-based filter.
|
||||
out = watch_path / _GRAPHIFY_OUT
|
||||
# When the caller supplied changed_paths, also evict preserved nodes whose
|
||||
# source_file matches a path that was changed (re-extracted) or deleted —
|
||||
# otherwise the old nodes for those files would survive forever.
|
||||
existing_graph = out / "graph.json"
|
||||
if existing_graph.exists():
|
||||
try:
|
||||
existing = json.loads(existing_graph.read_text(encoding="utf-8"))
|
||||
new_ast_ids = {n["id"] for n in result["nodes"]}
|
||||
preserved_nodes = [n for n in existing.get("nodes", []) if n["id"] not in new_ast_ids]
|
||||
evict_sources: set[str] = set(deleted_paths)
|
||||
if changed_paths is not None:
|
||||
for p in extract_targets:
|
||||
try:
|
||||
evict_sources.add(str(p.relative_to(project_root)))
|
||||
except ValueError:
|
||||
evict_sources.add(str(p))
|
||||
preserved_nodes = [
|
||||
n for n in existing.get("nodes", [])
|
||||
if n["id"] not in new_ast_ids
|
||||
and (not evict_sources or n.get("source_file") not in evict_sources)
|
||||
]
|
||||
all_ids = new_ast_ids | {n["id"] for n in preserved_nodes}
|
||||
preserved_edges = [
|
||||
e for e in existing.get("links", existing.get("edges", []))
|
||||
|
||||
Reference in New Issue
Block a user