feat: make repository reads alias-aware
Assistant: codex Assistant-Model: gpt-5.6-sol Assistant-Session: 01a049a4-ee9f-78e1-9d66-2cb0f9bea3e3
This commit is contained in:
parent
639b9aed08
commit
2e2ae1e5d0
19 changed files with 982 additions and 94 deletions
155
api/services/repository_aliases.py
Normal file
155
api/services/repository_aliases.py
Normal file
|
|
@ -0,0 +1,155 @@
|
|||
"""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))
|
||||
|
|
@ -16,6 +16,7 @@ from sqlalchemy.ext.asyncio import AsyncSession
|
|||
from api.models.managed_repo import ManagedRepo
|
||||
from api.models.task import Task, TaskPriority, TaskStatus
|
||||
from api.models.work_record_identifier_alias import WorkRecordIdentifierAlias
|
||||
from api.services.repository_aliases import resolve_repository_slug
|
||||
from api.models.workplan import Workplan
|
||||
|
||||
PLAN_SCHEMA = "repo-manager.identifier-migration-plan.v1"
|
||||
|
|
@ -384,11 +385,14 @@ async def repair_absent_prederivation_projection(
|
|||
|
||||
outcome = "repaired"
|
||||
async with session.begin():
|
||||
repo = await session.scalar(
|
||||
select(ManagedRepo).where(ManagedRepo.slug == repo_slug).with_for_update()
|
||||
)
|
||||
if repo is None:
|
||||
resolution = await resolve_repository_slug(session, repo_slug, required=False)
|
||||
if resolution is None:
|
||||
raise IdentifierMigrationError(f"repository projection is absent: {repo_slug}")
|
||||
repo = await session.scalar(
|
||||
select(ManagedRepo)
|
||||
.where(ManagedRepo.id == resolution.repo.id)
|
||||
.with_for_update()
|
||||
)
|
||||
if repo.topic_id != workplan_payload["topic_id"]:
|
||||
raise IdentifierMigrationError("repair topic does not match repository projection")
|
||||
|
||||
|
|
@ -462,6 +466,7 @@ async def repair_absent_prederivation_projection(
|
|||
|
||||
async def _assert_projection_preconditions(
|
||||
session: AsyncSession,
|
||||
repository_id: uuid.UUID,
|
||||
repo_slug: str,
|
||||
replacements: list[dict[str, Any]],
|
||||
*,
|
||||
|
|
@ -473,8 +478,7 @@ async def _assert_projection_preconditions(
|
|||
if mapping["kind"] == "workplan":
|
||||
source_query = text(
|
||||
"SELECT workplans.id FROM workplans "
|
||||
"JOIN managed_repos ON managed_repos.id = workplans.repo_id "
|
||||
"WHERE managed_repos.slug = :repo_slug AND workplans.id = :record_id "
|
||||
"WHERE workplans.repo_id = :repository_id AND workplans.id = :record_id "
|
||||
"FOR UPDATE"
|
||||
)
|
||||
target_query = text("SELECT id FROM workplans WHERE id = :record_id")
|
||||
|
|
@ -482,14 +486,13 @@ async def _assert_projection_preconditions(
|
|||
source_query = text(
|
||||
"SELECT tasks.id FROM tasks "
|
||||
"JOIN workplans ON workplans.id = tasks.workplan_id "
|
||||
"JOIN managed_repos ON managed_repos.id = workplans.repo_id "
|
||||
"WHERE managed_repos.slug = :repo_slug AND tasks.id = :record_id "
|
||||
"WHERE workplans.repo_id = :repository_id AND tasks.id = :record_id "
|
||||
"FOR UPDATE"
|
||||
)
|
||||
target_query = text("SELECT id FROM tasks WHERE id = :record_id")
|
||||
source = await session.execute(
|
||||
source_query,
|
||||
{"repo_slug": repo_slug, "record_id": source_id},
|
||||
{"repository_id": repository_id, "record_id": source_id},
|
||||
)
|
||||
if source.scalar_one_or_none() is None:
|
||||
raise IdentifierMigrationError(
|
||||
|
|
@ -514,8 +517,11 @@ async def apply_repository_identifier_migration(
|
|||
raise IdentifierMigrationError("migration requires a fresh database session")
|
||||
|
||||
async with session.begin():
|
||||
resolution = await resolve_repository_slug(session, repo_slug, required=False)
|
||||
if resolution is None:
|
||||
raise IdentifierMigrationError(f"repository projection is absent: {repo_slug}")
|
||||
await _assert_projection_preconditions(
|
||||
session, repo_slug, replacements, reverse=False
|
||||
session, resolution.repo.id, repo_slug, replacements, reverse=False
|
||||
)
|
||||
aliases = {
|
||||
alias.old_id: alias
|
||||
|
|
@ -596,6 +602,9 @@ async def reverse_repository_identifier_migration(
|
|||
raise IdentifierMigrationError("migration requires a fresh database session")
|
||||
|
||||
async with session.begin():
|
||||
resolution = await resolve_repository_slug(session, repo_slug, required=False)
|
||||
if resolution is None:
|
||||
raise IdentifierMigrationError(f"repository projection is absent: {repo_slug}")
|
||||
aliases = list(
|
||||
(
|
||||
await session.execute(
|
||||
|
|
@ -622,7 +631,7 @@ async def reverse_repository_identifier_migration(
|
|||
f"durable alias is not applied for {mapping['record_id']}"
|
||||
)
|
||||
await _assert_projection_preconditions(
|
||||
session, repo_slug, replacements, reverse=True
|
||||
session, resolution.repo.id, repo_slug, replacements, reverse=True
|
||||
)
|
||||
|
||||
for kind in ("task", "workplan"):
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue