diff --git a/api/services/task_record_id_backfill.py b/api/services/task_record_id_backfill.py index f37aecb..c7d43ec 100644 --- a/api/services/task_record_id_backfill.py +++ b/api/services/task_record_id_backfill.py @@ -114,3 +114,78 @@ async def backfill_task_record_ids( if not dry_run: row.record_id = record_id return report + + +async def backfill_from_forge( + session: Any, + repo_slugs: list[str], + *, + forge_base: str | None = None, + dry_run: bool = True, +) -> BackfillReport: + """Backfill from repositories cloned out of the forge. + + The local-path variant above needs a workstation checkout, which central + does not have and should not depend on: `ADR-012` decision 1 makes the forge + the projection source, and a backfill sourced from someone's laptop would + reintroduce exactly the coupling that ADR removes. + + Central can clone the forge directly, so it reads the pairing from the same + place it derives everything else. + """ + import tempfile + + from api.services.forge_projection import ( + DEFAULT_FORGE_BASE, + ForgeDeriveError, + _run_git, + ) + + base = forge_base or DEFAULT_FORGE_BASE + report = BackfillReport() + pairs: dict[str, str] = {} + + for slug in repo_slugs: + url = f"{base.rstrip('/')}/{slug}.git" + with tempfile.TemporaryDirectory(prefix=f"backfill-{slug}-") as tmp: + try: + _run_git("clone", "--depth", "1", "--quiet", url, tmp) + except (ForgeDeriveError, Exception): + # A repository that cannot be read contributes nothing. It must + # not silently reduce what the rest can identify. + continue + repo_pairs, repo_report = collect_pairs([Path(tmp)]) + report.scanned_files += repo_report.scanned_files + report.conflicts.extend(repo_report.conflicts) + for uid, rid in repo_pairs.items(): + prior = pairs.get(uid) + if prior and prior != rid: + report.conflicts.append({"uuid": uid, "first": prior, "second": rid}) + continue + pairs[uid] = rid + + report.pairs_found = len(pairs) + if not pairs: + return report + + from sqlalchemy import select + + from api.models.task import Task + + rows = list((await session.execute(select(Task))).scalars()) + by_id = {str(r.id): r for r in rows} + for uid, rid in pairs.items(): + row = by_id.get(uid) + if row is None: + report.unmatched_uuids += 1 + continue + if row.record_id == rid: + report.already_set += 1 + continue + if row.record_id and row.record_id != rid: + report.conflicts.append({"uuid": uid, "first": row.record_id, "second": rid}) + continue + report.updated += 1 + if not dry_run: + row.record_id = rid + return report