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
25 changes: 13 additions & 12 deletions src/lemoncrow/infra/code_intel/zoekt/indexer.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,8 @@
from time import time
from typing import Any

from .server import _git_dirs, _read_git_head

_TEXT_SUFFIXES = {
".py",
".ts",
Expand Down Expand Up @@ -91,34 +93,33 @@ def index_age_seconds(self) -> int:
snapshot = self.ensure_snapshot()
return int(max(0, time() - snapshot.indexed_at))

def _snapshot_cache_path(self) -> Path:
return self.repo_root / ".git" / "lemoncrow" / "zoekt_snapshot.json"
def _snapshot_cache_path(self) -> Path | None:
# The checkout's own git admin directory: a linked worktree's `.git` is a
# file, and its admin directory goes away with `git worktree remove`.
dirs = _git_dirs(self.repo_root)
return None if dirs is None else dirs[0] / "lemoncrow" / "zoekt_snapshot.json"

def _current_head(self) -> str | None:
head_file = self.repo_root / ".git" / "HEAD"
try:
ref = head_file.read_text(encoding="utf-8").strip()
if ref.startswith("ref: "):
ref_path = self.repo_root / ".git" / ref[5:]
return ref_path.read_text(encoding="utf-8").strip()
return ref
except OSError:
return None
return _read_git_head(self.repo_root)

def _load_snapshot_from_disk(self) -> ZoektIndexSnapshot | None:
cache_path = self._snapshot_cache_path()
if cache_path is None:
return None
try:
data = json.loads(cache_path.read_text(encoding="utf-8"))
stored_head = data.get("head")
current_head = self._current_head()
if current_head is None or stored_head != current_head:
return None
return ZoektIndexSnapshot.from_dict(data["snapshot"])
except (FileNotFoundError, json.JSONDecodeError, KeyError, TypeError, ValueError):
except (OSError, KeyError, TypeError, ValueError):
return None

def _save_snapshot_to_disk(self, snapshot: ZoektIndexSnapshot) -> None:
cache_path = self._snapshot_cache_path()
if cache_path is None:
return
current_head = self._current_head()
payload = json.dumps({"head": current_head, "snapshot": snapshot.to_dict()})
tmp = cache_path.with_suffix(".tmp")
Expand Down
104 changes: 88 additions & 16 deletions src/lemoncrow/pro/capabilities/code_context/engine.py
Original file line number Diff line number Diff line change
Expand Up @@ -2673,6 +2673,14 @@ def _resolve_index_max_workers() -> int:
_AUTOSYNC_GIT_HEAD_TIMEOUT_S = 3.0


@dataclass(frozen=True)
class _IndexedBaseline:
"""The tree signature and git HEAD an index run started from; None where unrecorded."""

signature: str | None = None
head: str | None = None


def _resolve_autosync_index_max_workers() -> int:
"""Worker count for background autosync indexing.

Expand Down Expand Up @@ -4357,6 +4365,9 @@ def _index_repo_unsafe(
) -> IndexStats:
"""Unlocked inner — callers must hold ``self._autosync_lock``."""
self.db_path.parent.mkdir(parents=True, exist_ok=True)
# Read before the scan, so a change landing mid-scan differs from the recorded
# baseline and the next autosync check reindexes it rather than taking it as done.
baseline = _IndexedBaseline(self._source_tree_signature(), self._autosync_git_head())
all_files = [
path
for path in iter_source_files(
Expand Down Expand Up @@ -4608,6 +4619,7 @@ def _index_repo_unsafe(
with self._connect() as conn:
self._init_schema(conn)
self._stamp_indexer_semantics_version(conn)
self._record_indexed_baseline(conn, baseline)

# Compute+persist the centrality map for the new index_version NOW so the
# O(edges) power iteration is charged to indexing, never to the first
Expand Down Expand Up @@ -4687,6 +4699,44 @@ def _stamp_indexer_semantics_version(self, conn: sqlite3.Connection) -> None:
(str(_CODE_INDEXER_SEMANTICS_VERSION),),
)

def _indexed_baseline_key(self) -> str:
# Per repo: a shared db_path holds several repos' indexes.
return f"autosync_baseline:{self.repo_id}"

def _record_indexed_baseline(self, conn: sqlite3.Connection, baseline: _IndexedBaseline) -> None:
conn.execute(
"""
INSERT INTO engine_state(key, value) VALUES (?, ?)
ON CONFLICT(key) DO UPDATE SET value = excluded.value
""",
(self._indexed_baseline_key(), json.dumps({"signature": baseline.signature, "head": baseline.head})),
)

def _indexed_baseline(self) -> _IndexedBaseline:
"""What the last completed index run of this repo started from.

The autosync loop measures change against this, not against its own first
look at the tree: an engine is replaced whenever the index version moves, and
a replacement that took the tree as it found it would treat every change made
since the last index run as already indexed.
"""
try:
with self._connect() as conn:
self._init_schema(conn)
row = conn.execute(
"SELECT value FROM engine_state WHERE key = ?", (self._indexed_baseline_key(),)
).fetchone()
data = json.loads(str(row["value"])) if row is not None else {}
except (sqlite3.Error, ValueError):
return _IndexedBaseline()
if not isinstance(data, dict):
return _IndexedBaseline()
signature, head = data.get("signature"), data.get("head")
return _IndexedBaseline(
signature if isinstance(signature, str) and signature else None,
head if isinstance(head, str) and head else None,
)

def _delete_files_index(self, conn: sqlite3.Connection, rels: list[str]) -> None:
"""Remove every indexed row for the files *rels*, in one pass per table.

Expand Down Expand Up @@ -14923,13 +14973,12 @@ def _maybe_autosync_reindex_locked(self, *, known_change: str | None = None) ->
"""
if known_change is not None:
self._autosync_state = "syncing"
if not self._run_index_subprocess():
if not self._reindex_and_adopt_baseline():
# Reindex failed; leave the signature/pending state stale so the
# next poll retries instead of recording a failed sync as done.
self._autosync_state = "idle"
self._record_autosync_event(event="reindex", reason=known_change, reindexed=False)
return False
self._autosync_signature = self._source_tree_signature()
self._autosync_last_sync_ms = int(time.time() * 1000)
self._autosync_pending_events = 0
self._autosync_state = "idle"
Expand All @@ -14938,13 +14987,20 @@ def _maybe_autosync_reindex_locked(self, *, known_change: str | None = None) ->
return True

current_signature = self._source_tree_signature()
reason = "source_signature_changed"
if self._autosync_signature is None:
self._autosync_signature = current_signature
self._autosync_last_sync_ms = int(time.time() * 1000)
self._autosync_state = "idle"
self._record_autosync_event(event="bootstrap", reason="seed_signature", reindexed=False)
return False
if current_signature == self._autosync_signature:
recorded = self._indexed_baseline().signature
if recorded is None:
# Indexed before baselines were recorded: nothing says what it covers.
reason = "no_indexed_baseline"
else:
self._autosync_signature = recorded
if current_signature == recorded:
self._autosync_last_sync_ms = int(time.time() * 1000)
self._autosync_state = "idle"
self._record_autosync_event(event="bootstrap", reason="indexed_baseline", reindexed=False)
return False
elif current_signature == self._autosync_signature:
self._autosync_state = "idle"
self._autosync_pending_events = 0
self._record_autosync_event(event="full_check", reason="unchanged", reindexed=False)
Expand All @@ -14957,18 +15013,34 @@ def _maybe_autosync_reindex_locked(self, *, known_change: str | None = None) ->
self._record_autosync_event(event="change_detected", reason="within_debounce_window", reindexed=False)
return False
self._autosync_state = "syncing"
if not self._run_index_subprocess():
if not self._reindex_and_adopt_baseline():
# Reindex failed; leave the signature/pending state stale so the next
# poll retries instead of recording a failed sync as complete.
self._autosync_state = "idle"
self._record_autosync_event(event="reindex", reason="source_signature_changed", reindexed=False)
self._record_autosync_event(event="reindex", reason=reason, reindexed=False)
return False
self._autosync_signature = self._source_tree_signature()
self._autosync_last_sync_ms = int(time.time() * 1000)
self._autosync_pending_events = 0
self._autosync_state = "idle"
self._autosync_reindex_count += 1
self._record_autosync_event(event="reindex", reason="source_signature_changed", reindexed=True)
self._record_autosync_event(event="reindex", reason=reason, reindexed=True)
return True

def _reindex_and_adopt_baseline(self) -> bool:
"""Run the index subprocess; on success, measure from the baseline it recorded.

A run can succeed without recording one -- a seeded worktree index it declines
to rebuild still carries the main checkout's baseline. Adopting that would
re-detect the same change and reindex again on every tick, so a run that left
the baseline untouched falls back to the tree and HEAD as they are now.
"""
before = self._indexed_baseline()
if not self._run_index_subprocess():
return False
after = self._indexed_baseline()
recorded = after if after != before else _IndexedBaseline()
self._autosync_signature = recorded.signature or self._source_tree_signature()
self._autosync_head = recorded.head or self._autosync_git_head() or self._autosync_head
return True

def _maybe_refresh_zoekt_index(self) -> None:
Expand Down Expand Up @@ -15205,15 +15277,15 @@ def _autosync_tick(self, now_ms: int) -> None:

def _autosync_poll_locked(self, now_ms: int) -> None:
head = self._autosync_git_head()
if head is not None and self._autosync_head is not None and head != self._autosync_head:
if head is not None and self._autosync_head is None:
# Like the tree signature: start from the HEAD the index was built at.
self._autosync_head = self._indexed_baseline().head or head
if head is not None and head != self._autosync_head:
# On failure keep the old HEAD, so the next tick retries.
if self._maybe_autosync_reindex_locked(known_change="head_moved"):
self._autosync_head = head
# The reindex reseeded the tree signature: that is a full check.
self._autosync_last_full_check_ms = now_ms
return
if head is not None:
self._autosync_head = head
last = self._autosync_last_full_check_ms
if last is None or now_ms - last >= max(self._autosync_poll_ms, _AUTOSYNC_MIN_POLL_MS):
self._autosync_last_full_check_ms = now_ms
Expand Down
112 changes: 110 additions & 2 deletions tests/core/test_code_context.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@
CallGraphNode,
traverse_call_graph,
)
from lemoncrow.pro.capabilities.code_context.engine import _CODE_INDEXER_SEMANTICS_VERSION
from lemoncrow.pro.capabilities.code_context.engine import _CODE_INDEXER_SEMANTICS_VERSION, _IndexedBaseline
from lemoncrow.pro.capabilities.code_context.models import SymbolRecord, TextMatch
from lemoncrow.pro.capabilities.code_context.output_policy import TRUNCATION_MARKER
from lemoncrow.pro.code_intel.cross_lang.runner import CrossLangRunner
Expand Down Expand Up @@ -2166,10 +2166,20 @@ def signature() -> str:
self.tree_walks += 1
return real_signature()

def record_baseline() -> None:
# What a completed index run leaves behind; the walk is the subprocess's, not the loop's.
with engine._connect() as conn:
engine._init_schema(conn)
engine._record_indexed_baseline(conn, _IndexedBaseline(real_signature(), engine._autosync_git_head()))

def run_index_subprocess(*, force: bool = False) -> bool:
self.reindexes.append(self.tree_walks)
return reindex_results.pop(0) if reindex_results else True
ok = reindex_results.pop(0) if reindex_results else True
if ok:
record_baseline()
return ok

record_baseline()
monkeypatch.setattr(engine, "_source_tree_signature", signature)
monkeypatch.setattr(engine, "_run_index_subprocess", run_index_subprocess)
monkeypatch.setattr(engine, "index_ready", lambda: True)
Expand Down Expand Up @@ -2324,6 +2334,104 @@ def is_alive(self) -> bool:
assert probe.reindexes == []


def _indexed_git_fixture(tmp_path: Path) -> CodeContextEngine:
repo = tmp_path / "repo"
_init_git_fixture_repo(repo)
_write_fixture_repo(repo)
_commit_all(repo, "initial")
engine = CodeContextEngine(repo, db_path=tmp_path / "code.sqlite", autosync_enabled=False)
engine.index_repo()
return engine


def _replacement_engine(engine: CodeContextEngine) -> CodeContextEngine:
"""What the engine cache builds once the index version moves: a new engine on the same index."""
return CodeContextEngine(engine.repo_root, db_path=engine.db_path, autosync_enabled=False)


def _finds(engine: CodeContextEngine, name: str) -> bool:
return any(s.symbol_name == name for s in engine.search_symbols(name, mode="lexical", limit=5, auto_index=False))


def test_a_replacement_engine_indexes_a_commit_made_since_the_last_index_run(tmp_path: Path) -> None:
engine = _indexed_git_fixture(tmp_path)
(engine.repo_root / "src" / "orders.py").write_text("class CommittedBeforeItLooked:\n pass\n", encoding="utf-8")
_commit_all(engine.repo_root, "move HEAD")

replacement = _replacement_engine(engine)
replacement._autosync_tick(0)

assert _finds(replacement, "CommittedBeforeItLooked")
assert replacement._autosync_history[-1]["reason"] == "head_moved"


def test_a_replacement_engine_indexes_a_working_tree_change_made_since_the_last_index_run(tmp_path: Path) -> None:
engine = _indexed_git_fixture(tmp_path)
(engine.repo_root / "src" / "orders.py").write_text("class EditedBeforeItLooked:\n pass\n", encoding="utf-8")

replacement = _replacement_engine(engine)
replacement._autosync_tick(0)

assert _finds(replacement, "EditedBeforeItLooked")
assert replacement._autosync_history[-1]["reason"] == "source_signature_changed"


def test_a_change_made_during_an_index_run_is_reindexed_by_the_next_check(
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
) -> None:
repo = tmp_path / "repo"
_init_git_fixture_repo(repo)
_write_fixture_repo(repo)
_commit_all(repo, "initial")
engine = CodeContextEngine(repo, db_path=tmp_path / "code.sqlite", autosync_enabled=False)
real_extract = engine._parallel_extract

def extract_while_the_tree_changes(*args: object, **kwargs: object) -> object:
results = real_extract(*args, **kwargs) # type: ignore[arg-type]
(repo / "src" / "orders.py").write_text("class WrittenMidScan:\n pass\n", encoding="utf-8")
return results

monkeypatch.setattr(engine, "_parallel_extract", extract_while_the_tree_changes)
engine.index_repo()
monkeypatch.undo()
assert not _finds(engine, "WrittenMidScan")

replacement = _replacement_engine(engine)
replacement._autosync_tick(0)

assert _finds(replacement, "WrittenMidScan")


def test_a_reindex_that_records_no_baseline_is_not_rerun_on_every_tick(
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
) -> None:
engine, _ = _autosync_probe_engine(tmp_path, monkeypatch, git=True)
# A seeded worktree index carries the main checkout's baseline, and an index run that
# declines to rebuild it (stale semantics, main not yet rebuilt) succeeds but records nothing.
with engine._connect() as conn:
engine._record_indexed_baseline(conn, _IndexedBaseline("main-signature", "0" * 40))
runs: list[int] = []
monkeypatch.setattr(engine, "_run_index_subprocess", lambda *, force=False: runs.append(1) or True)

for minute in range(3):
engine._autosync_tick(minute * _AUTOSYNC_MINUTE_MS)

assert len(runs) == 1


def test_an_index_without_a_recorded_baseline_is_reindexed_on_the_first_check(
tmp_path: Path, monkeypatch: pytest.MonkeyPatch
) -> None:
engine, probe = _autosync_probe_engine(tmp_path, monkeypatch)
with engine._connect() as conn:
conn.execute("DELETE FROM engine_state WHERE key = ?", (engine._indexed_baseline_key(),))

engine._autosync_tick(0)

assert len(probe.reindexes) == 1
assert engine._autosync_history[-1]["reason"] == "no_indexed_baseline"


def test_incremental_index_noop_does_not_bump_version(tmp_path: Path) -> None:
_write_fixture_repo(tmp_path)
engine = CodeContextEngine(tmp_path, db_path=tmp_path / "code.sqlite")
Expand Down
Loading
Loading