Skip to content

Commit 9bf8afa

Browse files
authored
Fix stale commit extraction resume failures (#4)
Reconcile completed commit plans before staging so stale aggregate failures do not force expensive re-extraction.
2 parents d88e722 + 98eab7f commit 9bf8afa

2 files changed

Lines changed: 167 additions & 9 deletions

File tree

scripts/streaming_conveyor.py

Lines changed: 82 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -606,8 +606,11 @@ def manifest_complete_commit_ranges(
606606
"""
607607
if range_size <= 0:
608608
raise ValueError(f"range_size must be positive, got {range_size}")
609-
if any(k == f"{repo}::commits" or k.startswith(f"{repo}::r")
610-
for k in manifest.failed):
609+
# A repo-level failure records an earlier extraction attempt, not a unit of
610+
# range coverage. Once the authoritative plan and every range prove exact
611+
# coverage it is stale. A failed range remains authoritative and must keep
612+
# completion fail-closed until that same range key is marked done.
613+
if any(k.startswith(f"{repo}::r") for k in manifest.failed):
611614
return None
612615

613616
plan = manifest.done.get(commit_plan_key(repo))
@@ -660,6 +663,52 @@ def manifest_complete_commit_ranges(
660663
return tuple(ranges)
661664

662665

666+
def mark_commit_stream_complete(
667+
repo: str,
668+
manifest: Manifest,
669+
manifest_lock: threading.Lock | None,
670+
complete_ranges: Sequence[tuple[int, int]],
671+
) -> dict:
672+
"""Persist aggregate completion proven by the plan and exact range coverage.
673+
674+
``Manifest.mark_done`` also removes an earlier aggregate failure. Keeping
675+
this as a derived summary avoids treating ``repo::commits`` as a second,
676+
independent source of completion truth.
677+
"""
678+
if not complete_ranges:
679+
raise ValueError(f"cannot mark {repo}::commits complete without ranges")
680+
plan = manifest.done.get(commit_plan_key(repo))
681+
if not isinstance(plan, dict) or plan.get("source") != "commit_plan":
682+
raise RuntimeError(
683+
f"cannot mark {repo}::commits complete without authoritative commit_plan"
684+
)
685+
try:
686+
n_records = int(plan["n_records"])
687+
except (KeyError, TypeError, ValueError) as exc:
688+
raise RuntimeError(
689+
f"cannot mark {repo}::commits complete: invalid commit_plan n_records"
690+
) from exc
691+
if complete_ranges[0][0] != 0 or complete_ranges[-1][1] != n_records:
692+
raise RuntimeError(
693+
f"cannot mark {repo}::commits complete: ranges do not cover "
694+
f"[0, {n_records})"
695+
)
696+
info = {
697+
"source": "commits",
698+
"repo": repo,
699+
"complete": True,
700+
"completion_proof": "commit_plan_exact_range_coverage",
701+
"n_records": n_records,
702+
"range_count": len(complete_ranges),
703+
}
704+
if manifest_lock is None:
705+
manifest.mark_done(f"{repo}::commits", info)
706+
else:
707+
with manifest_lock:
708+
manifest.mark_done(f"{repo}::commits", info)
709+
return info
710+
711+
663712
def manifest_done_commit_intervals(repo: str, manifest: Manifest) -> tuple[tuple[int, int], ...]:
664713
"""Return validated done commit intervals for ``repo``.
665714
@@ -1295,6 +1344,7 @@ def should_stage_repo_from_manifest(
12951344
manifest: Manifest,
12961345
range_size: int,
12971346
only_repos: set[str] | None,
1347+
manifest_lock: threading.Lock | None = None,
12981348
) -> bool:
12991349
"""Return False when the manifest proves this repo needs no extraction.
13001350
@@ -1310,10 +1360,20 @@ def should_stage_repo_from_manifest(
13101360
return False
13111361

13121362
code_needed = streams in {"both", "code"} and not manifest.is_done(code_key(repo))
1313-
commits_needed = (
1314-
streams in {"both", "commits"}
1315-
and manifest_complete_commit_ranges(repo, manifest, range_size) is None
1316-
)
1363+
commits_needed = False
1364+
if streams in {"both", "commits"}:
1365+
complete_ranges = manifest_complete_commit_ranges(repo, manifest, range_size)
1366+
commits_needed = complete_ranges is None
1367+
if complete_ranges is not None:
1368+
# This callback can skip extraction entirely, so reconcile the
1369+
# derived aggregate sentinel here instead of waiting for
1370+
# run_commits_half, which will never be called for this repo.
1371+
mark_commit_stream_complete(
1372+
repo,
1373+
manifest,
1374+
manifest_lock,
1375+
complete_ranges,
1376+
)
13171377
return code_needed or commits_needed
13181378

13191379

@@ -1872,6 +1932,12 @@ def run_commits_half(
18721932
if complete_ranges is not None:
18731933
first = complete_ranges[0][0]
18741934
last = complete_ranges[-1][1]
1935+
mark_commit_stream_complete(
1936+
repo,
1937+
manifest,
1938+
manifest_lock,
1939+
complete_ranges,
1940+
)
18751941
_log(
18761942
f"SKIP (done) {repo}::commits: manifest covers "
18771943
f"{len(complete_ranges)} range(s) [{first}:{last}); "
@@ -2246,7 +2312,15 @@ def handle_commit_range(
22462312
# True iff EVERY range for this repo is now marked done in the manifest
22472313
# (covers resume-skipped + newly-done; excludes cancelled/failed). Drives
22482314
# temp + extract-cache retention in process_one_repo.
2249-
all_ranges_done = manifest_covers_commit_span(repo, manifest, n_records)
2315+
complete_ranges = manifest_complete_commit_ranges(repo, manifest, range_size)
2316+
all_ranges_done = complete_ranges is not None
2317+
if complete_ranges is not None:
2318+
mark_commit_stream_complete(
2319+
repo,
2320+
manifest,
2321+
manifest_lock,
2322+
complete_ranges,
2323+
)
22502324
return done, failed, all_ranges_done
22512325

22522326

@@ -2987,6 +3061,7 @@ def should_process(repo: str) -> bool:
29873061
manifest=manifest,
29883062
range_size=args.range_size,
29893063
only_repos=only_repos,
3064+
manifest_lock=manifest_lock,
29903065
)
29913066
if should_stage:
29923067
ensure_min_free_disk(

tests/test_streaming_conveyor_progress.py

Lines changed: 85 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -164,6 +164,52 @@ def test_manifest_complete_commit_ranges_uses_authoritative_record_count(tmp_pat
164164
) == ((0, 500), (500, 823))
165165

166166

167+
def test_manifest_complete_commit_ranges_ignores_stale_aggregate_failure(tmp_path):
168+
import streaming_conveyor
169+
170+
manifest = streaming_conveyor.Manifest(
171+
path=tmp_path / "_done.json",
172+
done={
173+
"repo::commit_plan": {"source": "commit_plan", "n_records": 823},
174+
"repo::r0": {"range": [0, 500], "source": "commits"},
175+
"repo::r500": {"range": [500, 823], "source": "commits"},
176+
},
177+
failed={
178+
"repo::commits": {
179+
"stage": "extract_git_history",
180+
"detail": "stale failure from an earlier attempt",
181+
}
182+
},
183+
)
184+
185+
assert streaming_conveyor.manifest_complete_commit_ranges(
186+
"repo", manifest, 500
187+
) == ((0, 500), (500, 823))
188+
189+
190+
def test_manifest_complete_commit_ranges_rejects_failed_range(tmp_path):
191+
import streaming_conveyor
192+
193+
manifest = streaming_conveyor.Manifest(
194+
path=tmp_path / "_done.json",
195+
done={
196+
"repo::commit_plan": {"source": "commit_plan", "n_records": 823},
197+
"repo::r0": {"range": [0, 500], "source": "commits"},
198+
"repo::r500": {"range": [500, 823], "source": "commits"},
199+
},
200+
failed={
201+
"repo::r500": {
202+
"stage": "pack",
203+
"detail": "range is not authoritatively complete",
204+
}
205+
},
206+
)
207+
208+
assert streaming_conveyor.manifest_complete_commit_ranges(
209+
"repo", manifest, 500
210+
) is None
211+
212+
167213
def test_manifest_complete_commit_ranges_refuses_exact_multiple(tmp_path):
168214
import streaming_conveyor
169215

@@ -740,7 +786,12 @@ def test_run_commits_half_skips_extract_when_manifest_proves_complete(tmp_path):
740786
"repo::r0": {"range": [0, 500], "source": "commits"},
741787
"repo::r500": {"range": [500, 612], "source": "commits"},
742788
},
743-
failed={},
789+
failed={
790+
"repo::commits": {
791+
"stage": "extract_git_history",
792+
"detail": "stale failure from an earlier attempt",
793+
}
794+
},
744795
)
745796

746797
with ThreadPoolExecutor(max_workers=1) as pool:
@@ -764,6 +815,15 @@ def test_run_commits_half_skips_extract_when_manifest_proves_complete(tmp_path):
764815
)
765816

766817
assert (done, failed, all_done) == (0, 0, True)
818+
assert manifest.done["repo::commits"] == {
819+
"source": "commits",
820+
"repo": "repo",
821+
"complete": True,
822+
"completion_proof": "commit_plan_exact_range_coverage",
823+
"n_records": 612,
824+
"range_count": 2,
825+
}
826+
assert "repo::commits" not in manifest.failed
767827

768828

769829
def test_process_one_repo_cleans_partial_intermediates_by_default(tmp_path, monkeypatch):
@@ -1060,6 +1120,15 @@ def runner(
10601120
range_keys = sorted(key for key in manifest.done if key.startswith("repo::r"))
10611121
assert range_keys == ["repo::r0", "repo::r1", "repo::r2"]
10621122
assert manifest.done["repo::commit_plan"]["n_records"] == 3
1123+
assert manifest.done["repo::commits"] == {
1124+
"source": "commits",
1125+
"repo": "repo",
1126+
"complete": True,
1127+
"completion_proof": "commit_plan_exact_range_coverage",
1128+
"n_records": 3,
1129+
"range_count": 3,
1130+
}
1131+
assert "repo::commits" not in manifest.failed
10631132
assert len(observed_kwargs) == 3
10641133
reader = DedupStore(str(db), near=False, commit_every=1)
10651134
try:
@@ -1302,7 +1371,12 @@ def test_should_stage_repo_skips_manifest_proven_complete_both_streams(tmp_path)
13021371
"repo::r0": {"range": [0, 500], "source": "commits"},
13031372
"repo::r500": {"range": [500, 612], "source": "commits"},
13041373
},
1305-
failed={},
1374+
failed={
1375+
"repo::commits": {
1376+
"stage": "extract_git_history",
1377+
"detail": "stale failure from an earlier attempt",
1378+
}
1379+
},
13061380
)
13071381

13081382
assert not streaming_conveyor.should_stage_repo_from_manifest(
@@ -1313,6 +1387,15 @@ def test_should_stage_repo_skips_manifest_proven_complete_both_streams(tmp_path)
13131387
range_size=500,
13141388
only_repos=None,
13151389
)
1390+
assert manifest.done["repo::commits"] == {
1391+
"source": "commits",
1392+
"repo": "repo",
1393+
"complete": True,
1394+
"completion_proof": "commit_plan_exact_range_coverage",
1395+
"n_records": 612,
1396+
"range_count": 2,
1397+
}
1398+
assert "repo::commits" not in manifest.failed
13161399
assert streaming_conveyor.should_stage_repo_from_manifest(
13171400
"repo",
13181401
streams="both",

0 commit comments

Comments
 (0)