Skip to content

Commit 06007c3

Browse files
committed
Stream source cache populate progress
1 parent 4dffb6f commit 06007c3

4 files changed

Lines changed: 29 additions & 9 deletions

File tree

scripts/streaming_conveyor.py

Lines changed: 11 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -1224,20 +1224,23 @@ def populate_code_source_cache(
12241224
) -> dict:
12251225
"""Populate the code source cache without running index/tokenize stages."""
12261226
started = time.monotonic()
1227+
1228+
def emit_repo_ready(repo: str, repo_dir: Path, repo_count: int) -> None:
1229+
progress.emit(
1230+
"source_cache_repo_ready",
1231+
repo=repo,
1232+
repo_dir=str(repo_dir),
1233+
source_cache_dir=str(source_cache_dir),
1234+
repo_count=repo_count,
1235+
)
1236+
12271237
report = sr.populate_source_cache(
12281238
work_root,
12291239
should_process,
12301240
source_cache_dir,
12311241
max_repos=max_repos,
1242+
on_repo_ready=emit_repo_ready,
12321243
)
1233-
for idx, item in enumerate(report["repos"], start=1):
1234-
progress.emit(
1235-
"source_cache_repo_ready",
1236-
repo=item["repo"],
1237-
repo_dir=item["path"],
1238-
source_cache_dir=str(source_cache_dir),
1239-
repo_count=idx,
1240-
)
12411244
report = {
12421245
**report,
12431246
"elapsed_s": round(time.monotonic() - started, 6),

scripts/streaming_reindex.py

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -665,6 +665,7 @@ def populate_source_cache(
665665
source_cache_dir: Path,
666666
*,
667667
max_repos: int | None = None,
668+
on_repo_ready: Callable[[str, Path, int], None] | None = None,
668669
) -> dict:
669670
"""Materialize source-cache repos without running the tokenization pipeline.
670671
@@ -684,6 +685,8 @@ def populate_source_cache(
684685
try:
685686
for repo, repo_dir in gen:
686687
repos.append({"repo": repo, "path": str(repo_dir)})
688+
if on_repo_ready is not None:
689+
on_repo_ready(repo, repo_dir, len(repos))
687690
if max_repos is not None and len(repos) >= max_repos:
688691
break
689692
finally:

tests/test_streaming_conveyor_progress.py

Lines changed: 12 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -66,15 +66,25 @@ def test_populate_code_source_cache_emits_ready_and_summary(tmp_path, monkeypatc
6666
cache = tmp_path / "source_cache"
6767
observed = {}
6868

69-
def fake_populate(work_root, should_process, source_cache_dir, *, max_repos):
69+
def fake_populate(
70+
work_root,
71+
should_process,
72+
source_cache_dir,
73+
*,
74+
max_repos,
75+
on_repo_ready,
76+
):
7077
observed["work_root"] = work_root
7178
observed["source_cache_dir"] = source_cache_dir
7279
observed["max_repos"] = max_repos
80+
observed["has_callback"] = on_repo_ready is not None
7381
repos = []
7482
for repo in ("repo-a", "repo.bare", "repo-b"):
7583
if should_process(repo):
7684
repos.append({"repo": repo, "path": str(source_cache_dir / repo)})
7785
repos = repos[:max_repos]
86+
for idx, item in enumerate(repos, start=1):
87+
on_repo_ready(item["repo"], Path(item["path"]), idx)
7888
return {
7989
"source_cache_dir": str(source_cache_dir),
8090
"repos": repos,
@@ -95,6 +105,7 @@ def fake_populate(work_root, should_process, source_cache_dir, *, max_repos):
95105
"work_root": tmp_path / "work",
96106
"source_cache_dir": cache,
97107
"max_repos": 1,
108+
"has_callback": True,
98109
}
99110
assert report["repo_count"] == 1
100111
rows = [

tests/test_streaming_reindex_run_checked.py

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -170,14 +170,17 @@ def fake_stream(work_root, should_process, *, source_cache_dir, source_cache_onl
170170

171171
monkeypatch.setattr(streaming_reindex, "stream_repo_subtrees", fake_stream)
172172

173+
ready = []
173174
report = streaming_reindex.populate_source_cache(
174175
tmp_path / "work",
175176
streaming_reindex.is_code_worktree_repo,
176177
cache,
177178
max_repos=1,
179+
on_repo_ready=lambda repo, path, count: ready.append((repo, path, count)),
178180
)
179181

180182
assert calls == [(tmp_path / "work", cache, False)]
183+
assert ready == [("repo-a", cache / "repo-a", 1)]
181184
assert report == {
182185
"source_cache_dir": str(cache),
183186
"repos": [{"repo": "repo-a", "path": str(cache / "repo-a")}],

0 commit comments

Comments
 (0)