state-hub/api/services/task_record_id_backfill.py
tegwick cdd5cef373
All checks were successful
CI Smoke / host-smoke (push) Successful in 0s
CI Smoke / container-smoke (push) Successful in 1s
Build and Publish Multi-Context Image / build-and-push (push) Successful in 26s
fix(backfill): source task identities from the forge, not a workstation
The first implementation took local filesystem paths. Central has no workstation
checkout and must not depend on one: ADR-012 decision 1 makes the forge the
projection source, and a backfill reading someone's laptop would reintroduce the
exact coupling that ADR removes.

Surfaced concretely — central's postgres is not reachable from the workstation
(only the API tunnel, which is HTTP), so the local-path variant cannot reach the
database it needs to update, while the pod can clone the forge and already holds
the connection.

A repository that cannot be cloned contributes nothing rather than reducing what
the rest can identify.

Refs STATE-WP-0083-T06

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>

Assistant: claude-code
Assistant-Model: opus
Assistant-Process: 2583210@bnt-lap001
Assistant-Session: f2bff2d5-e9b2-4338-92ca-10282a927006
2026-08-26 02:11:23 +02:00

191 lines
6.5 KiB
Python

"""Backfill canonical record ids onto existing task rows (STATE-WP-0083-T06).
Only the repository files hold the mapping. A file task declares both its
canonical id and the projection UUID it was registered under:
```task
id: CUST-WP-0067-T01
state_hub_task_id: "f3608db4-..."
```
so the pairing can be read directly rather than guessed from titles. Anything a
file does not claim is left alone: a task row whose canonical id cannot be
established keeps `record_id` null, and the reset continues to refuse to act on
it. An unknown identity must stay unknown rather than be inferred.
"""
from __future__ import annotations
import re
from dataclasses import dataclass, field
from pathlib import Path
from typing import Any
_TASK_BLOCK_RE = re.compile(r"```task\s*\n(.*?)\n```", re.DOTALL)
_ID_RE = re.compile(r"^id:\s*(\S+)", re.MULTILINE)
_UUID_RE = re.compile(r'state_hub_task_id:\s*"?([0-9a-f-]{36})"?')
@dataclass
class BackfillReport:
scanned_files: int = 0
pairs_found: int = 0
updated: int = 0
already_set: int = 0
conflicts: list[dict[str, str]] = field(default_factory=list)
unmatched_uuids: int = 0
def to_dict(self) -> dict[str, Any]:
return {
"schema": "state-hub.task-record-id-backfill.v1",
"scanned_files": self.scanned_files,
"pairs_found": self.pairs_found,
"updated": self.updated,
"already_set": self.already_set,
"unmatched_uuids": self.unmatched_uuids,
"conflicts": self.conflicts,
}
def collect_pairs(roots: list[Path]) -> tuple[dict[str, str], BackfillReport]:
"""Map projection UUID -> canonical record id, from workplan files."""
report = BackfillReport()
pairs: dict[str, str] = {}
for root in roots:
wp_dir = root / "workplans"
if not wp_dir.is_dir():
continue
for path in sorted(wp_dir.rglob("*.md")):
if path.name.startswith("."):
continue
try:
text = path.read_text(encoding="utf-8")
except (OSError, UnicodeDecodeError):
continue
report.scanned_files += 1
for block in _TASK_BLOCK_RE.finditer(text):
body = block.group(1)
rid = _ID_RE.search(body)
uid = _UUID_RE.search(body)
if not rid or not uid:
continue
record_id, task_uuid = rid.group(1).strip(), uid.group(1)
prior = pairs.get(task_uuid)
if prior and prior != record_id:
# One UUID claimed by two canonical ids: a duplicate
# registration. Recording it and skipping is the only safe
# option — picking one would fabricate an identity.
report.conflicts.append(
{"uuid": task_uuid, "first": prior, "second": record_id}
)
continue
pairs[task_uuid] = record_id
report.pairs_found = len(pairs)
return pairs, report
async def backfill_task_record_ids(
session: Any, roots: list[Path], *, dry_run: bool = True
) -> BackfillReport:
from sqlalchemy import select
from api.models.task import Task
pairs, report = collect_pairs(roots)
if not pairs:
return report
rows = list((await session.execute(select(Task))).scalars())
by_id = {str(r.id): r for r in rows}
for task_uuid, record_id in pairs.items():
row = by_id.get(task_uuid)
if row is None:
report.unmatched_uuids += 1
continue
if row.record_id == record_id:
report.already_set += 1
continue
if row.record_id and row.record_id != record_id:
report.conflicts.append(
{"uuid": task_uuid, "first": row.record_id, "second": record_id}
)
continue
report.updated += 1
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