state-hub/api/services/repository_aliases.py
tegwick 2e2ae1e5d0
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 25s
feat: make repository reads alias-aware
Assistant: codex
Assistant-Model: gpt-5.6-sol
Assistant-Session: 01a049a4-ee9f-78e1-9d66-2cb0f9bea3e3
2026-08-29 10:47:51 +02:00

155 lines
5.1 KiB
Python

"""Canonical repository identity resolution across current and prior slugs.
The slug registry is the lookup boundary. Historical records deliberately keep
the slug they recorded; callers use ``slug_values`` when they need an identity-
wide read and ``canonical_slug`` when they create a new reference.
"""
from __future__ import annotations
import uuid
from dataclasses import dataclass
from fastapi import HTTPException
from sqlalchemy import func, or_, select
from sqlalchemy.ext.asyncio import AsyncSession
from api.models.fabric_graph import FabricGraphEdge, FabricGraphImport, FabricGraphNode
from api.models.managed_repo import ManagedRepo
from api.models.repository_rename import RepositorySlug
@dataclass(frozen=True)
class RepositorySlugResolution:
repo: ManagedRepo
requested_slug: str
canonical_slug: str
slug_status: str
aliases: tuple[str, ...]
source_operation_id: uuid.UUID | None = None
@property
def slug_values(self) -> tuple[str, ...]:
return (self.canonical_slug, *self.aliases)
async def resolve_repository_slug(
session: AsyncSession,
slug: str,
*,
required: bool = True,
) -> RepositorySlugResolution | None:
"""Resolve a current or protected prior slug to one repository UUID."""
record = (
await session.execute(
select(RepositorySlug).where(RepositorySlug.slug == slug)
)
).scalar_one_or_none()
# Compatibility for databases upgraded before the identity backfill has run.
if record is None:
repo = (
await session.execute(select(ManagedRepo).where(ManagedRepo.slug == slug))
).scalar_one_or_none()
if repo is None:
if required:
raise HTTPException(status_code=404, detail=f"Repo '{slug}' not found")
return None
return RepositorySlugResolution(
repo=repo,
requested_slug=slug,
canonical_slug=repo.slug,
slug_status="canonical",
aliases=(),
)
repo = await session.get(ManagedRepo, record.repo_id)
if repo is None:
raise HTTPException(
status_code=409,
detail=f"Repository slug registry entry '{slug}' has no repository",
)
records = list(
(
await session.execute(
select(RepositorySlug)
.where(RepositorySlug.repo_id == repo.id)
.order_by(RepositorySlug.kind.desc(), RepositorySlug.slug)
)
).scalars()
)
canonicals = [item for item in records if item.kind == "canonical"]
if len(canonicals) != 1 or canonicals[0].slug != repo.slug:
raise HTTPException(
status_code=409,
detail=f"Repository '{repo.id}' has inconsistent canonical slug state",
)
return RepositorySlugResolution(
repo=repo,
requested_slug=slug,
canonical_slug=repo.slug,
slug_status=record.kind,
aliases=tuple(item.slug for item in records if item.kind == "alias"),
source_operation_id=record.source_operation_id,
)
async def canonicalize_repository_slug(session: AsyncSession, value: str) -> str:
"""Canonicalize a value only when it is a registered repository identity."""
resolution = await resolve_repository_slug(session, value, required=False)
return resolution.canonical_slug if resolution is not None else value
async def repository_resolution_for_id(
session: AsyncSession, repo: ManagedRepo
) -> RepositorySlugResolution:
return await resolve_repository_slug(session, repo.slug) # type: ignore[return-value]
async def stale_external_references(
session: AsyncSession,
resolution: RepositorySlugResolution,
) -> list[dict[str, object]]:
"""Name external projections that still contain an historical slug.
These are handoffs to the projection owner, never implicit rewrite targets.
"""
if not resolution.aliases:
return []
aliases = list(resolution.aliases)
checks = (
(FabricGraphImport, "source_repo_slug"),
(FabricGraphNode, "source_repo_slug"),
(FabricGraphNode, "repo_slug"),
(FabricGraphEdge, "source_repo_slug"),
)
stale: list[dict[str, object]] = []
for model, field_name in checks:
field = getattr(model, field_name)
rows = (
await session.execute(
select(field, func.count()).where(field.in_(aliases)).group_by(field)
)
).all()
for value, count in rows:
stale.append(
{
"owner": "railiance-fabric",
"surface": model.__tablename__,
"field": field_name,
"value": value,
"count": count,
"status": "stale",
"handoff": "owner update or re-ingest required",
}
)
return stale
def affected_slug_predicate(field: object, slugs: tuple[str, ...]):
"""Match a JSONB slug-list against any name in one repository lineage."""
return or_(*(field.contains([slug]) for slug in slugs))