From 2e2ae1e5d0ca476b5ed520addc6c6c47a822ce3f Mon Sep 17 00:00:00 2001 From: tegwick Date: Sat, 29 Aug 2026 10:47:51 +0200 Subject: [PATCH] feat: make repository reads alias-aware Assistant: codex Assistant-Model: gpt-5.6-sol Assistant-Session: 01a049a4-ee9f-78e1-9d66-2cb0f9bea3e3 --- WORK-RECORDS.md | 2 +- api/routers/capability_requests.py | 110 ++++++- api/routers/fabric.py | 8 +- api/routers/interface_changes.py | 41 ++- api/routers/messages.py | 150 ++++++++- api/routers/repo_goals.py | 8 +- api/routers/repos.py | 120 +++++-- api/routers/sbom.py | 41 ++- api/routers/services.py | 18 +- api/routers/token_events.py | 10 + api/schemas/managed_repo.py | 10 + api/services/repository_aliases.py | 155 ++++++++++ .../work_record_identifier_migration.py | 31 +- dashboard/src/repos.md | 2 +- dashboard/src/repos/[slug].md | 10 +- scripts/consistency_check.py | 3 + tests/test_repository_alias_routing.py | 292 ++++++++++++++++++ .../test_work_record_identifier_migration.py | 52 ++++ ...85-repository-lineage-preserving-rename.md | 13 +- 19 files changed, 982 insertions(+), 94 deletions(-) create mode 100644 api/services/repository_aliases.py create mode 100644 tests/test_repository_alias_routing.py diff --git a/WORK-RECORDS.md b/WORK-RECORDS.md index 2ada1ac..86db8fd 100644 --- a/WORK-RECORDS.md +++ b/WORK-RECORDS.md @@ -324,7 +324,7 @@ | task | STATE-WP-0085-T01 | done | — | workplans/STATE-WP-0085-repository-lineage-preserving-rename.md | | task | STATE-WP-0085-T02 | done | — | workplans/STATE-WP-0085-repository-lineage-preserving-rename.md | | task | STATE-WP-0085-T03 | done | — | workplans/STATE-WP-0085-repository-lineage-preserving-rename.md | -| task | STATE-WP-0085-T04 | todo | — | workplans/STATE-WP-0085-repository-lineage-preserving-rename.md | +| task | STATE-WP-0085-T04 | done | — | workplans/STATE-WP-0085-repository-lineage-preserving-rename.md | | task | STATE-WP-0085-T05 | todo | — | workplans/STATE-WP-0085-repository-lineage-preserving-rename.md | | task | STATE-WP-0085-T06 | todo | — | workplans/STATE-WP-0085-repository-lineage-preserving-rename.md | | task | STATE-WP-0085-T07 | todo | — | workplans/STATE-WP-0085-repository-lineage-preserving-rename.md | diff --git a/api/routers/capability_requests.py b/api/routers/capability_requests.py index 3192004..271ed41 100644 --- a/api/routers/capability_requests.py +++ b/api/routers/capability_requests.py @@ -1,7 +1,7 @@ import re import uuid from datetime import datetime, timezone -from fastapi import HTTPException +from fastapi import APIRouter, Depends, HTTPException, Query, status from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession @@ -11,7 +11,6 @@ from api.models.agent_message import AgentMessage from api.models.capability_catalog import CapabilityCatalog from api.models.capability_request import CapabilityRequest from api.models.domain import Domain -from api.models.managed_repo import ManagedRepo from api.models.task import Task from api.services.suggestion_relevance import bump_matching_for_capability_request from api.schemas.capability_request import ( @@ -22,12 +21,15 @@ from api.schemas.capability_request import ( CapabilityRequestRead, CapabilityRequestReroute, CapabilityRequestStatusPatch, + CatalogCreate, + CatalogPatch, + CatalogRead, ) from hub_core.routers.capabilities import ( - create_capability_catalog_router, create_capability_request_read_router, create_capability_request_write_router, ) +from api.services.repository_aliases import resolve_repository_slug # --------------------------------------------------------------------------- @@ -394,12 +396,102 @@ def _check_transition(current: str, target: str) -> None: ) -router = create_capability_catalog_router( - get_session, - domain_model=Domain, - repo_model=ManagedRepo, - catalog_model=CapabilityCatalog, +router = APIRouter(tags=["capability-requests"]) + + +async def _catalog_domain(slug: str, session: AsyncSession) -> Domain: + domain = ( + await session.execute(select(Domain).where(Domain.slug == slug)) + ).scalar_one_or_none() + if domain is None: + raise HTTPException(status_code=404, detail=f"Domain '{slug}' not found") + return domain + + +@router.post( + "/capability-catalog/", + response_model=CatalogRead, + status_code=status.HTTP_201_CREATED, ) +async def create_catalog_entry( + body: CatalogCreate, + session: AsyncSession = Depends(get_session), +) -> CapabilityCatalog: + domain = await _catalog_domain(body.domain, session) + repo_id = None + if body.repo_slug: + repo_id = (await resolve_repository_slug(session, body.repo_slug)).repo.id + entry = CapabilityCatalog( + domain_id=domain.id, + repo_id=repo_id, + capability_type=body.capability_type, + title=body.title, + description=body.description, + keywords=body.keywords, + ) + session.add(entry) + try: + await session.commit() + except Exception: + await session.rollback() + raise HTTPException( + status_code=409, + detail=( + f"Catalog entry '{body.title}' for type '{body.capability_type}' " + f"already exists in domain '{body.domain}'" + ), + ) + await session.refresh(entry) + return entry + + +@router.get("/capability-catalog/", response_model=list[CatalogRead]) +async def list_catalog( + domain: str | None = Query(None), + capability_type: str | None = Query(None), + status_filter: str | None = Query(None, alias="status"), + repo_slug: str | None = Query(None), + session: AsyncSession = Depends(get_session), +) -> list[CapabilityCatalog]: + query = select(CapabilityCatalog).order_by(CapabilityCatalog.created_at.desc()) + if domain: + query = query.where( + CapabilityCatalog.domain_id == (await _catalog_domain(domain, session)).id + ) + if capability_type: + query = query.where(CapabilityCatalog.capability_type == capability_type) + if repo_slug: + query = query.where( + CapabilityCatalog.repo_id + == (await resolve_repository_slug(session, repo_slug)).repo.id + ) + if status_filter and status_filter != "all": + query = query.where(CapabilityCatalog.status == status_filter) + elif not status_filter: + query = query.where(CapabilityCatalog.status == "active") + return list((await session.execute(query)).scalars().all()) + + +@router.patch("/capability-catalog/{entry_id}", response_model=CatalogRead) +async def patch_catalog_entry( + entry_id: uuid.UUID, + body: CatalogPatch, + session: AsyncSession = Depends(get_session), +) -> CapabilityCatalog: + entry = await session.get(CapabilityCatalog, entry_id) + if entry is None: + raise HTTPException(status_code=404, detail=f"Catalog entry '{entry_id}' not found") + if body.repo_slug is not None: + entry.repo_id = (await resolve_repository_slug(session, body.repo_slug)).repo.id + for field in ("description", "keywords", "status"): + value = getattr(body, field) + if value is not None: + setattr(entry, field, value) + await session.commit() + await session.refresh(entry) + return entry + + router.include_router( create_capability_request_read_router( get_session, @@ -432,4 +524,4 @@ router.include_router( after_dispute=_notify_on_dispute, after_reroute=_notify_on_reroute, ) -) \ No newline at end of file +) diff --git a/api/routers/fabric.py b/api/routers/fabric.py index c389342..bb7c921 100644 --- a/api/routers/fabric.py +++ b/api/routers/fabric.py @@ -23,6 +23,7 @@ from api.services.fabric_graph import ( record_fabric_graph_error, split_graph_ingest_body, ) +from api.services.repository_aliases import resolve_repository_slug router = APIRouter(prefix="/fabric", tags=["fabric"]) @@ -157,7 +158,12 @@ async def list_graph_nodes( if domain: query = query.where(FabricGraphNode.domain_slug == domain) if repo: - query = query.where(FabricGraphNode.repo_slug == repo) + resolution = await resolve_repository_slug(session, repo, required=False) + query = query.where( + FabricGraphNode.repo_slug.in_( + resolution.slug_values if resolution else (repo,) + ) + ) if canonical_category: query = query.where(FabricGraphNode.canon_category == canonical_category) if fabric_id: diff --git a/api/routers/interface_changes.py b/api/routers/interface_changes.py index d4985a3..078557e 100644 --- a/api/routers/interface_changes.py +++ b/api/routers/interface_changes.py @@ -8,13 +8,17 @@ from sqlalchemy.ext.asyncio import AsyncSession from api.database import get_session from api.models.agent_message import AgentMessage from api.models.interface_change import InterfaceChange -from api.models.managed_repo import ManagedRepo from api.models.progress_event import ProgressEvent from api.schemas.interface_change import ( InterfaceChangeCreate, InterfaceChangePatch, InterfaceChangeRead, ) +from api.services.repository_aliases import ( + affected_slug_predicate, + canonicalize_repository_slug, + resolve_repository_slug, +) router = APIRouter(prefix="/interface-changes", tags=["interface-changes"]) @@ -32,7 +36,12 @@ async def create_interface_change( if body.change_type not in _VALID_CHANGE_TYPES: raise HTTPException(status_code=422, detail=f"change_type must be one of {sorted(_VALID_CHANGE_TYPES)}") - repo = await _repo_by_slug(body.repo_slug, session) + resolution = await resolve_repository_slug(session, body.repo_slug) + repo = resolution.repo + affected_repo_slugs = [ + await canonicalize_repository_slug(session, slug) + for slug in body.affected_repo_slugs + ] change = InterfaceChange( repo_id=repo.id, interface_type=body.interface_type, @@ -40,7 +49,7 @@ async def create_interface_change( title=body.title, description=body.description, affected_paths=body.affected_paths, - affected_repo_slugs=body.affected_repo_slugs, + affected_repo_slugs=affected_repo_slugs, planned_for=body.planned_for, author=body.author, status="draft", @@ -68,7 +77,11 @@ async def list_interface_changes( if change_type: q = q.where(InterfaceChange.change_type == change_type) if affected_repo: - q = q.where(InterfaceChange.affected_repo_slugs.contains([affected_repo])) + resolution = await resolve_repository_slug( + session, affected_repo, required=False + ) + values = resolution.slug_values if resolution else (affected_repo,) + q = q.where(affected_slug_predicate(InterfaceChange.affected_repo_slugs, values)) result = await session.execute(q) return [InterfaceChangeRead.from_orm_with_slug(c) for c in result.scalars().all()] @@ -94,7 +107,13 @@ async def patch_interface_change( status_code=409, detail=f"Cannot edit a change with status '{change.status}'. Only draft records are mutable.", ) - for field, value in body.model_dump(exclude_unset=True).items(): + payload = body.model_dump(exclude_unset=True) + if payload.get("affected_repo_slugs") is not None: + payload["affected_repo_slugs"] = [ + await canonicalize_repository_slug(session, slug) + for slug in payload["affected_repo_slugs"] + ] + for field, value in payload.items(): setattr(change, field, value) await session.commit() await session.refresh(change) @@ -119,12 +138,13 @@ async def publish_interface_change( # Send inbox notifications to agents of affected repos affected = change.affected_repo_slugs or [] for slug in affected: + target_slug = await canonicalize_repository_slug(session, slug) paths_summary = ", ".join(change.affected_paths[:5]) if change.affected_paths else "see description" if len(change.affected_paths) > 5: paths_summary += f" (+{len(change.affected_paths) - 5} more)" msg = AgentMessage( from_agent=change.repo.slug, - to_agent=slug, + to_agent=target_slug, subject=f"[{change.change_type.upper()}] {change.title}", body=( f"**Interface change published by `{change.repo.slug}`**\n\n" @@ -174,12 +194,9 @@ async def resolve_interface_change( return InterfaceChangeRead.from_orm_with_slug(change) -async def _repo_by_slug(slug: str, session: AsyncSession) -> ManagedRepo: - result = await session.execute(select(ManagedRepo).where(ManagedRepo.slug == slug)) - repo = result.scalar_one_or_none() - if repo is None: - raise HTTPException(status_code=404, detail=f"Repo '{slug}' not found") - return repo +async def _repo_by_slug(slug: str, session: AsyncSession): + resolution = await resolve_repository_slug(session, slug) + return resolution.repo async def _get_or_404(change_id: uuid.UUID, session: AsyncSession) -> InterfaceChange: diff --git a/api/routers/messages.py b/api/routers/messages.py index 6da5630..d8334cb 100644 --- a/api/routers/messages.py +++ b/api/routers/messages.py @@ -1,7 +1,153 @@ +from datetime import datetime, timezone + +from fastapi import APIRouter, Depends, HTTPException, status +from sqlalchemy import or_, select +from sqlalchemy.ext.asyncio import AsyncSession + from api.database import get_session from api.models.agent_message import AgentMessage -from hub_core.routers.messages import create_messages_router +from api.schemas.agent_message import MessageCreate, MessageRead, MessageReply +from api.services.repository_aliases import ( + canonicalize_repository_slug, + resolve_repository_slug, +) +from hub_core.message_identity import resolve_message_reference +from hub_core.models.message_identity_alias import MessageIdentityAlias + +router = APIRouter(prefix="/messages", tags=["messages"]) + + +async def _get_message(reference: str, session: AsyncSession) -> AgentMessage: + message_id = await resolve_message_reference( + session, reference, alias_model=MessageIdentityAlias + ) + if message_id is None: + raise HTTPException( + status_code=404, detail=f"Message reference {reference!r} not found" + ) + message = await session.get(AgentMessage, message_id) + if message is None: + raise HTTPException( + status_code=404, detail=f"Message reference {reference!r} not found" + ) + return message + + +@router.post("/", response_model=MessageRead, status_code=status.HTTP_201_CREATED) +async def send_message( + body: MessageCreate, + session: AsyncSession = Depends(get_session), +) -> AgentMessage: + if body.thread_id and await session.get(AgentMessage, body.thread_id) is None: + raise HTTPException(status_code=404, detail=f"Thread root {body.thread_id} not found") + payload = body.model_dump() + payload["from_agent"] = await canonicalize_repository_slug(session, body.from_agent) + payload["to_agent"] = await canonicalize_repository_slug(session, body.to_agent) + message = AgentMessage(**payload) + session.add(message) + await session.commit() + await session.refresh(message) + return message + + +@router.get("/", response_model=list[MessageRead]) +async def list_messages( + to_agent: str | None = None, + from_agent: str | None = None, + unread_only: bool = False, + limit: int = 50, + session: AsyncSession = Depends(get_session), +) -> list[AgentMessage]: + query = select(AgentMessage).where(AgentMessage.archived_at.is_(None)) + if to_agent: + resolution = await resolve_repository_slug(session, to_agent, required=False) + values = resolution.slug_values if resolution else (to_agent,) + query = query.where( + or_(AgentMessage.to_agent.in_(values), AgentMessage.to_agent == "broadcast") + ) + if from_agent: + resolution = await resolve_repository_slug(session, from_agent, required=False) + values = resolution.slug_values if resolution else (from_agent,) + query = query.where(AgentMessage.from_agent.in_(values)) + if unread_only: + query = query.where(AgentMessage.read_at.is_(None)) + result = await session.execute( + query.order_by(AgentMessage.created_at.desc()).limit(limit) + ) + return list(result.scalars().all()) + + +@router.get("/thread/{thread_id}", response_model=list[MessageRead]) +async def get_thread( + thread_id: str, + session: AsyncSession = Depends(get_session), +) -> list[AgentMessage]: + resolved = await resolve_message_reference( + session, thread_id, alias_model=MessageIdentityAlias + ) + if resolved is None: + raise HTTPException( + status_code=404, detail=f"Message reference {thread_id!r} not found" + ) + result = await session.execute( + select(AgentMessage) + .where(or_(AgentMessage.id == resolved, AgentMessage.thread_id == resolved)) + .order_by(AgentMessage.created_at) + ) + return list(result.scalars().all()) + + +@router.patch("/{message_id}/read", response_model=MessageRead) +async def mark_read( + message_id: str, + session: AsyncSession = Depends(get_session), +) -> AgentMessage: + message = await _get_message(message_id, session) + if message.read_at is None: + message.read_at = datetime.now(timezone.utc) + await session.commit() + await session.refresh(message) + return message + + +@router.patch("/{message_id}/archive", response_model=MessageRead) +async def archive_message( + message_id: str, + session: AsyncSession = Depends(get_session), +) -> AgentMessage: + message = await _get_message(message_id, session) + message.archived_at = datetime.now(timezone.utc) + if message.read_at is None: + message.read_at = message.archived_at + await session.commit() + await session.refresh(message) + return message + + +@router.post( + "/{message_id}/reply", + response_model=MessageRead, + status_code=status.HTTP_201_CREATED, +) +async def reply_to_message( + message_id: str, + body: MessageReply, + session: AsyncSession = Depends(get_session), +) -> AgentMessage: + original = await _get_message(message_id, session) + if original.read_at is None: + original.read_at = datetime.now(timezone.utc) + reply = AgentMessage( + from_agent=await canonicalize_repository_slug(session, body.from_agent), + to_agent=await canonicalize_repository_slug(session, original.from_agent), + subject=f"Re: {original.subject}", + body=body.body, + thread_id=original.thread_id or original.id, + ) + session.add(reply) + await session.commit() + await session.refresh(reply) + return reply -router = create_messages_router(get_session, message_model=AgentMessage) __all__ = ["router"] diff --git a/api/routers/repo_goals.py b/api/routers/repo_goals.py index f836b91..601684e 100644 --- a/api/routers/repo_goals.py +++ b/api/routers/repo_goals.py @@ -8,16 +8,14 @@ from api.database import get_session from api.models.managed_repo import ManagedRepo from api.models.repo_goal import RepoGoal, RepoGoalStatus from api.schemas.repo_goal import RepoGoalCreate, RepoGoalRead, RepoGoalUpdate +from api.services.repository_aliases import resolve_repository_slug router = APIRouter(prefix="/repo-goals", tags=["repo-goals"]) async def _resolve_repo(repo_slug: str, session: AsyncSession) -> ManagedRepo: - result = await session.execute(select(ManagedRepo).where(ManagedRepo.slug == repo_slug)) - repo = result.scalar_one_or_none() - if repo is None: - raise HTTPException(status_code=404, detail=f"Repo '{repo_slug}' not found") - return repo + resolution = await resolve_repository_slug(session, repo_slug) + return resolution.repo @router.get("/", response_model=list[RepoGoalRead]) diff --git a/api/routers/repos.py b/api/routers/repos.py index 0cf8ede..dec36f6 100644 --- a/api/routers/repos.py +++ b/api/routers/repos.py @@ -54,6 +54,13 @@ from api.services.sbom_nexus import SBOMNexusError from api.services.sbom_nexus import get_json as get_sbom_nexus_json from api.services.sbom_nexus import reads_from_nexus from api.services.repository_identity import stage_initial_repository_identity +from api.services.repository_aliases import ( + RepositorySlugResolution, + affected_slug_predicate, + repository_resolution_for_id, + resolve_repository_slug, + stale_external_references, +) from hub_core.routers.repos import create_repos_router router = APIRouter(prefix="/repos", tags=["repos"]) @@ -140,14 +147,14 @@ async def list_repos( if business_stake: q = q.where(ManagedRepo.business_stake.contains([business_stake])) result = await session.execute(q) - return await _project_repo_reads(list(result.scalars().all())) + return await _project_repo_reads(session, list(result.scalars().all())) @router.post("/", response_model=RepoRead, status_code=status.HTTP_201_CREATED) async def register_repo( body: RepoCreate, session: AsyncSession = Depends(get_session), -) -> ManagedRepo: +) -> RepoRead: domain_result = await session.execute(select(Domain).where(Domain.slug == body.domain_slug)) domain_obj = domain_result.scalar_one_or_none() if domain_obj is None: @@ -197,7 +204,7 @@ async def register_repo( ) from exc await session.refresh(repo) await _publish_repo_registered(repo, body, domain_obj) - return repo + return (await _project_repo_reads(session, [repo]))[0] @router.post("/onboard", response_model=RepoOnboardResult) @@ -483,7 +490,8 @@ async def get_repo_doi( Results are cached by fingerprint. Pass ?force_refresh=true to bypass the cache. """ - repo = await _get_repo_by_slug(slug, session) + resolution = await resolve_repository_slug(session, slug) + repo = resolution.repo sbom_projections = await _sbom_projection_map() domain_result = await session.execute(select(Domain).where(Domain.id == repo.domain_id)) domain_obj = domain_result.scalar_one_or_none() @@ -517,7 +525,7 @@ async def get_repo_doi( if not force_refresh and cached and cached.fingerprint == fp and cached.criteria: return DoIReport( - repo_slug=slug, + repo_slug=repo.slug, tier=cached.tier, core_pass=cached.core_pass, standard_pass=cached.standard_pass, @@ -556,11 +564,11 @@ async def get_repo_doi( async def get_repo_by_id( repo_id: uuid.UUID, session: AsyncSession = Depends(get_session), -) -> ManagedRepo: +) -> RepoRead: repo = await session.get(ManagedRepo, repo_id) if repo is None: raise HTTPException(status_code=404, detail=f"Repo '{repo_id}' not found") - return repo + return (await _project_repo_reads(session, [repo]))[0] @router.get("/scope-health", response_model=list[RepoScopeHealth]) @@ -615,9 +623,10 @@ async def update_repo_with_classification( slug: str, body: RepoUpdate, session: AsyncSession = Depends(get_session), -) -> ManagedRepo: +) -> RepoRead: """Patch repo metadata including classification spine fields.""" - repo = await _get_repo_by_slug(slug, session) + resolution = await resolve_repository_slug(session, slug) + repo = resolution.repo payload = body.model_dump(exclude_unset=True) requested_domain_slug = payload.pop("domain_slug", None) if requested_domain_slug is not None: @@ -658,7 +667,7 @@ async def update_repo_with_classification( setattr(repo, field, value) await session.commit() await session.refresh(repo) - return repo + return (await _project_repo_reads(session, [repo], requested=resolution))[0] @router.get("/{slug}", response_model=RepoRead) @@ -666,28 +675,53 @@ async def get_repo_with_sbom_projection( slug: str, session: AsyncSession = Depends(get_session), ) -> RepoRead: - repo = await _get_repo_by_slug(slug, session) - return (await _project_repo_reads([repo]))[0] + resolution = await resolve_repository_slug(session, slug) + return ( + await _project_repo_reads( + session, + [resolution.repo], + requested=resolution, + include_stale_external=True, + ) + )[0] router.include_router( _core_repo_router( include_collection_routes=False, include_lookup_routes=False, + include_slug_routes=False, ) ) +@router.post("/{slug}/paths", response_model=RepoRead) +async def register_repo_path( + slug: str, + body: RepoPathRegister, + session: AsyncSession = Depends(get_session), +) -> RepoRead: + resolution = await resolve_repository_slug(session, slug) + repo = resolution.repo + host_paths = dict(repo.host_paths or {}) + host_paths[body.host] = body.path + repo.host_paths = host_paths + await session.commit() + await session.refresh(repo) + return (await _project_repo_reads(session, [repo], requested=resolution))[0] + + @router.patch("/{slug}/archive", response_model=RepoRead) async def archive_repo( slug: str, session: AsyncSession = Depends(get_session), -) -> ManagedRepo: - repo = await _get_repo_by_slug(slug, session) +) -> RepoRead: + resolution = await resolve_repository_slug(session, slug) + repo = resolution.repo repo.status = "archived" await session.commit() await session.refresh(repo) - return repo + return (await _project_repo_reads(session, [repo], requested=resolution))[0] @router.get("/{slug}/dispatch", response_model=RepoDispatch) @@ -701,7 +735,8 @@ async def get_repo_dispatch( call it at session start to discover what work is pending without needing to read state-hub summary or scan workplan files manually. """ - repo = await _get_repo_by_slug(slug, session) + resolution = await resolve_repository_slug(session, slug) + repo = resolution.repo # Active goal goal_result = await session.execute( @@ -765,7 +800,9 @@ async def get_repo_dispatch( ic_result = await session.execute( select(InterfaceChange).where( InterfaceChange.status == "published", - InterfaceChange.affected_repo_slugs.contains([slug]), + affected_slug_predicate( + InterfaceChange.affected_repo_slugs, resolution.slug_values + ), ).order_by(InterfaceChange.published_at.desc()) ) pending_changes = [ @@ -794,7 +831,12 @@ async def get_repo_dispatch( ) return RepoDispatch( - repo_slug=slug, + repo_slug=resolution.canonical_slug, + requested_slug=resolution.requested_slug, + canonical_slug=resolution.canonical_slug, + slug_status=resolution.slug_status, + aliases=list(resolution.aliases), + stale_external_references=await stale_external_references(session, resolution), active_goal=active_goal, active_workplans=dispatch_workstreams, human_interventions=all_interventions, @@ -819,7 +861,8 @@ async def sync_repo_consistency( Returns the raw JSON output from consistency_check.py. Query param ?fix=false to run check-only without writing. """ - repo = await _get_repo_by_slug(slug, session) + resolution = await resolve_repository_slug(session, slug) + repo = resolution.repo hostname = socket.gethostname() host_paths = repo.host_paths or {} @@ -828,13 +871,13 @@ async def sync_repo_consistency( raise HTTPException( status_code=503, detail=( - f"No accessible path for repo '{slug}' on host '{hostname}'. " - f"Register with: POST /repos/{slug}/paths/" + f"No accessible path for repo '{repo.slug}' on host '{hostname}'. " + f"Register with: POST /repos/{repo.slug}/paths/" ), ) script = Path(__file__).parent.parent.parent / "scripts" / "consistency_check.py" - cmd = [sys.executable, str(script), "--repo", slug, "--json", + cmd = [sys.executable, str(script), "--repo", repo.slug, "--json", "--api-base", settings.api_base] if fix: cmd.append("--fix") @@ -853,11 +896,8 @@ async def sync_repo_consistency( async def _get_repo_by_slug(slug: str, session: AsyncSession) -> ManagedRepo: - result = await session.execute(select(ManagedRepo).where(ManagedRepo.slug == slug)) - repo = result.scalar_one_or_none() - if repo is None: - raise HTTPException(status_code=404, detail=f"Repo '{slug}' not found") - return repo + resolution = await resolve_repository_slug(session, slug) + return resolution.repo def _repo_doi_dict(repo: ManagedRepo, domain_slug: str | None) -> dict: @@ -899,11 +939,35 @@ def _projected_last_sbom_at( return str(repo.last_sbom_at) if repo.last_sbom_at else None -async def _project_repo_reads(repositories: list[ManagedRepo]) -> list[RepoRead]: +async def _project_repo_reads( + session: AsyncSession, + repositories: list[ManagedRepo], + *, + requested: RepositorySlugResolution | None = None, + include_stale_external: bool = False, +) -> list[RepoRead]: projections = await _sbom_projection_map() result: list[RepoRead] = [] for repository in repositories: + resolution = ( + requested + if requested is not None and requested.repo.id == repository.id + else await repository_resolution_for_id(session, repository) + ) read = RepoRead.model_validate(repository) + read = read.model_copy( + update={ + "requested_slug": resolution.requested_slug, + "canonical_slug": resolution.canonical_slug, + "slug_status": resolution.slug_status, + "aliases": list(resolution.aliases), + "stale_external_references": ( + await stale_external_references(session, resolution) + if include_stale_external + else [] + ), + } + ) if repository.slug in projections: read = read.model_copy( update={ diff --git a/api/routers/sbom.py b/api/routers/sbom.py index 53e7bcb..a5cb550 100644 --- a/api/routers/sbom.py +++ b/api/routers/sbom.py @@ -27,6 +27,7 @@ from api.services.sbom_nexus import ( reads_from_nexus, writes_to_nexus, ) +from api.services.repository_aliases import resolve_repository_slug router = APIRouter(prefix="/sbom", tags=["sbom"]) logger = logging.getLogger(__name__) @@ -68,9 +69,12 @@ async def ingest_sbom( session: AsyncSession = Depends(get_session), ) -> dict: """Create a new SBOM snapshot for a repo. Previous snapshots are retained.""" - repo = await _get_repo_by_slug(body.repo_slug, session) + resolution = await resolve_repository_slug(session, body.repo_slug) + repo = resolution.repo if writes_to_nexus(): - payload = await _nexus_post("/sbom/ingest/", body=body.model_dump(mode="json")) + nexus_body = body.model_dump(mode="json") + nexus_body["repo_slug"] = resolution.canonical_slug + payload = await _nexus_post("/sbom/ingest/", body=nexus_body) try: snapshot_at = datetime.fromisoformat( payload["snapshot_at"].replace("Z", "+00:00") @@ -126,7 +130,7 @@ async def ingest_sbom( await session.commit() await _meter_compat(session, request, "POST", "/sbom/ingest/") return { - "repo_slug": body.repo_slug, + "repo_slug": resolution.canonical_slug, "snapshot_id": str(snap.id), "ingested": len(body.entries), "snapshot_at": now.isoformat(), @@ -142,6 +146,10 @@ async def list_snapshots( """List SBOM snapshots, newest first. Optionally filter by repo.""" await _meter_compat(session, request, "GET", "/sbom/snapshots/") if reads_from_nexus(): + if repo_slug: + repo_slug = ( + await resolve_repository_slug(session, repo_slug) + ).canonical_slug payload = await _nexus_get( "/sbom/snapshots/", params={"repo_slug": repo_slug} if repo_slug else None, @@ -207,6 +215,10 @@ async def list_sbom_entries( """Return entries from the latest snapshot per repo (default) or filter by repo.""" await _meter_compat(session, request, "GET", "/sbom/") if reads_from_nexus(): + if repo_slug: + repo_slug = ( + await resolve_repository_slug(session, repo_slug) + ).canonical_slug params = { key: value for key, value in { @@ -299,10 +311,11 @@ async def get_repo_sbom( session: AsyncSession = Depends(get_session), ) -> SBOMRepoView: """Return the latest snapshot entries for a specific repo.""" - repo = await _get_repo_by_slug(repo_slug, session) + resolution = await resolve_repository_slug(session, repo_slug) + repo = resolution.repo await _meter_compat(session, request, "GET", "/sbom/{repo_slug}") if reads_from_nexus(): - payload = await _nexus_get(f"/sbom/{repo_slug}") + payload = await _nexus_get(f"/sbom/{resolution.canonical_slug}") payload["entries"] = [ _translate_entry(entry, repo.id) for entry in payload.get("entries", []) ] @@ -322,7 +335,7 @@ async def get_repo_sbom( ) entries = list(rows.scalars().all()) return SBOMRepoView( - repo_slug=repo_slug, + repo_slug=resolution.canonical_slug, last_sbom_at=repo.last_sbom_at, entry_count=len(entries), entries=[SBOMEntryRead.model_validate(e) for e in entries], @@ -330,11 +343,8 @@ async def get_repo_sbom( async def _get_repo_by_slug(slug: str, session: AsyncSession) -> ManagedRepo: - result = await session.execute(select(ManagedRepo).where(ManagedRepo.slug == slug)) - repo = result.scalar_one_or_none() - if repo is None: - raise HTTPException(status_code=404, detail=f"Repo '{slug}' not found") - return repo + resolution = await resolve_repository_slug(session, slug) + return resolution.repo async def _nexus_get(path: str, *, params: dict | None = None): @@ -355,10 +365,11 @@ async def _local_repo_ids(items: list[dict], session: AsyncSession) -> dict[str, slugs = {item.get("repo_slug") for item in items} if None in slugs: raise HTTPException(status_code=502, detail="SBOM Nexus response omitted repo_slug") - result = await session.execute( - select(ManagedRepo.slug, ManagedRepo.id).where(ManagedRepo.slug.in_(slugs)) - ) - repo_ids = dict(result.all()) + repo_ids: dict[str, uuid.UUID] = {} + for slug in slugs: + resolution = await resolve_repository_slug(session, slug, required=False) + if resolution is not None: + repo_ids[slug] = resolution.repo.id missing = sorted(slugs - repo_ids.keys()) if missing: raise HTTPException( diff --git a/api/routers/services.py b/api/routers/services.py index 763da99..982ab1b 100644 --- a/api/routers/services.py +++ b/api/routers/services.py @@ -13,7 +13,6 @@ from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.orm import selectinload from api.database import get_session -from api.models.managed_repo import ManagedRepo from api.models.service_catalog import ( ServiceCatalog, ServiceCloud, @@ -22,6 +21,7 @@ from api.models.service_catalog import ( ServiceThirdParty, ) from api.schemas.service import ServiceCatalogRead, ServiceUpsert +from api.services.repository_aliases import resolve_repository_slug router = APIRouter(prefix="/services", tags=["services"]) @@ -42,6 +42,7 @@ async def list_services( development_type: str | None = None, maturity_level: int | None = None, status: str | None = None, + repo_slug: str | None = None, session: AsyncSession = Depends(get_session), ) -> list[ServiceCatalog]: q = select(ServiceCatalog).options(*_WITH_EXTENSIONS) @@ -53,6 +54,11 @@ async def list_services( q = q.where(ServiceCatalog.maturity_level == maturity_level) if status: q = q.where(ServiceCatalog.status == status) + if repo_slug: + resolution = await resolve_repository_slug(session, repo_slug) + q = q.join(ServiceFirstParty).where( + ServiceFirstParty.repo_id == resolution.repo.id + ) q = q.order_by(ServiceCatalog.name.asc()) result = await session.execute(q) return list(result.scalars().all()) @@ -131,12 +137,10 @@ async def _apply_extensions(svc: ServiceCatalog, body: ServiceUpsert, session: A if body.first_party is not None: data = body.first_party.model_dump(exclude={"repo_slug"}) if body.first_party.repo_slug and not data.get("repo_id"): - repo = (await session.execute( - select(ManagedRepo).where(ManagedRepo.slug == body.first_party.repo_slug) - )).scalar_one_or_none() - if repo is None: - raise HTTPException(status_code=404, detail=f"Repo '{body.first_party.repo_slug}' not found") - data["repo_id"] = repo.id + resolution = await resolve_repository_slug( + session, body.first_party.repo_slug + ) + data["repo_id"] = resolution.repo.id await _upsert_ext(ServiceFirstParty, svc.id, data, session) diff --git a/api/routers/token_events.py b/api/routers/token_events.py index 6dcd853..802e173 100644 --- a/api/routers/token_events.py +++ b/api/routers/token_events.py @@ -13,6 +13,7 @@ from api.models.task import Task from api.models.token_event import TokenEvent from api.models.workplan import Workplan from api.services.legacy_compat import meter_legacy_query_param +from api.services.repository_aliases import resolve_repository_slug from api.schemas.token_event import ( RepoTokenSummary, TokenAggregateRow, @@ -288,6 +289,7 @@ async def get_token_summary( @router.get("/by-repo/", response_model=list[RepoTokenSummary]) async def get_tokens_by_repo( + repo_slug: str | None = None, measurement_kind: str | None = None, source_provider: str | None = None, since: datetime | None = None, @@ -304,6 +306,12 @@ async def get_tokens_by_repo( Only events that resolve to a repo are included. """ + requested_repo_id = None + if repo_slug: + requested_repo_id = ( + await resolve_repository_slug(session, repo_slug) + ).repo.id + # Fetch all events, workstreams, repos in three queries (avoids N+1) events_result = await session.execute( _filter_query( @@ -341,6 +349,8 @@ async def get_tokens_by_repo( rid = resolve_repo_id(e) if not rid or rid not in repo_map: continue + if requested_repo_id is not None and rid != requested_repo_id: + continue if rid not in groups: groups[rid] = { "repo_id": rid, diff --git a/api/schemas/managed_repo.py b/api/schemas/managed_repo.py index e160ab2..ff1f5cf 100644 --- a/api/schemas/managed_repo.py +++ b/api/schemas/managed_repo.py @@ -114,6 +114,11 @@ class RepoRead(CoreRepoRead, ClassificationFields): last_state_synced_at: datetime | None = None created_at: datetime updated_at: datetime + requested_slug: str | None = None + canonical_slug: str | None = None + slug_status: Literal["canonical", "alias"] = "canonical" + aliases: list[str] = Field(default_factory=list) + stale_external_references: list[dict[str, Any]] = Field(default_factory=list) class DispatchTask(BaseModel): @@ -158,6 +163,11 @@ class ScopeIssueDetail(BaseModel): class RepoDispatch(BaseModel): repo_slug: str + requested_slug: str + canonical_slug: str + slug_status: Literal["canonical", "alias"] + aliases: list[str] = Field(default_factory=list) + stale_external_references: list[dict[str, Any]] = Field(default_factory=list) active_goal: dict[str, Any] | None active_workplans: list[DispatchWorkplan] human_interventions: list[DispatchTask] diff --git a/api/services/repository_aliases.py b/api/services/repository_aliases.py new file mode 100644 index 0000000..8fe52cf --- /dev/null +++ b/api/services/repository_aliases.py @@ -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)) diff --git a/api/services/work_record_identifier_migration.py b/api/services/work_record_identifier_migration.py index 0b773ae..a49c6da 100644 --- a/api/services/work_record_identifier_migration.py +++ b/api/services/work_record_identifier_migration.py @@ -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"): diff --git a/dashboard/src/repos.md b/dashboard/src/repos.md index e9b11ef..fbf041d 100644 --- a/dashboard/src/repos.md +++ b/dashboard/src/repos.md @@ -268,7 +268,7 @@ if (domainBlocks.length === 0) { ${rows.map(r => html` - ${r.repo} + ${r.repo} ${r.category} ${r.capTags} ${_doiBadge(r._doiTier)} diff --git a/dashboard/src/repos/[slug].md b/dashboard/src/repos/[slug].md index 3a090d6..73a6b89 100644 --- a/dashboard/src/repos/[slug].md +++ b/dashboard/src/repos/[slug].md @@ -21,9 +21,17 @@ if (raw.error) { const name = raw.name || raw.slug || repoSlug; display(html`

Repo · ${name}

`); display(html`

← Repos  |  ← Token Cost

`); + if (raw.slug_status === "alias") { + display(html`
+ Repository alias: ${raw.requested_slug} now resolves to + ${raw.canonical_slug}. + Historical records keep the former value. +
`); + } const FIELD_ORDER = [ - "id","slug","name","domain_slug","status","description", + "id","slug","requested_slug","canonical_slug","slug_status","aliases", + "name","domain_slug","status","description", "local_path","remote_url","git_fingerprint", "sbom_source","last_sbom_at","last_state_synced_at", "created_at","updated_at", diff --git a/scripts/consistency_check.py b/scripts/consistency_check.py index 01f2143..a5c6a07 100644 --- a/scripts/consistency_check.py +++ b/scripts/consistency_check.py @@ -1200,6 +1200,9 @@ def check_repo( ) return report + # Repository aliases are accepted at the API boundary, but every new + # binding, inbox write, and generated file must use the current identity. + repo_slug = str(repo.get("canonical_slug") or repo.get("slug") or repo_slug) repo_id: str = repo["id"] repo_path: str = resolve_repo_path(repo, repo_path_override) report = ConsistencyReport(repo_slug=repo_slug, repo_path=repo_path) diff --git a/tests/test_repository_alias_routing.py b/tests/test_repository_alias_routing.py new file mode 100644 index 0000000..800c6ed --- /dev/null +++ b/tests/test_repository_alias_routing.py @@ -0,0 +1,292 @@ +from __future__ import annotations + +from datetime import datetime, timezone + +import pytest +from sqlalchemy import func, select +from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker + +from api.models.agent_message import AgentMessage +from api.models.fabric_graph import FabricGraphImport, FabricGraphNode +from api.models.interface_change import InterfaceChange +from api.models.managed_repo import ManagedRepo +from api.models.repository_rename import RepositorySlug +from tests.conftest import ( + create_test_domain, + create_test_repo, + create_test_workplan, +) + + +async def _rename_in_state_hub(test_engine, repo_id: str) -> None: + factory = async_sessionmaker(test_engine, class_=AsyncSession, expire_on_commit=False) + async with factory() as session: + repo = await session.get(ManagedRepo, repo_id) + old = ( + await session.execute( + select(RepositorySlug).where(RepositorySlug.slug == "flex-auth") + ) + ).scalar_one() + old.kind = "alias" + old.protected = True + repo.slug = "access-engine" + repo.name = "Access Engine" + session.add( + RepositorySlug( + repo_id=repo.id, + slug="access-engine", + kind="canonical", + protected=True, + ) + ) + await session.commit() + + +@pytest.mark.asyncio +async def test_alias_lookup_dispatch_bindings_and_external_staleness( + client, test_engine +): + domain = await create_test_domain(client) + repo = await create_test_repo(client, domain_slug=domain["slug"], slug="flex-auth") + workplan = await create_test_workplan( + client, + repo_id=repo["id"], + slug="FLEX-WP-0001", + title="Security stack migration", + ) + task = await client.post( + "/tasks/", + json={ + "workplan_id": workplan["id"], + "title": "Preserve identity", + "status": "todo", + "priority": "high", + }, + ) + assert task.status_code == 201, task.text + + await _rename_in_state_hub(test_engine, repo["id"]) + factory = async_sessionmaker(test_engine, class_=AsyncSession, expire_on_commit=False) + async with factory() as session: + change = InterfaceChange( + repo_id=repo["id"], + interface_type="rest_api", + change_type="breaking", + title="Historical consumer notice", + description="Recorded before cutover", + affected_paths=["/authorize"], + affected_repo_slugs=["flex-auth"], + status="published", + published_at=datetime.now(timezone.utc), + author="pytest", + ) + import_run = FabricGraphImport( + source_repo_slug="railiance-fabric", + content_hash="a" * 64, + graph_json={}, + validation_status="valid", + is_latest=True, + ) + session.add_all([change, import_run]) + await session.flush() + session.add( + FabricGraphNode( + import_id=import_run.id, + source_repo_slug="railiance-fabric", + graph_id="repo:flex-auth", + kind="repository", + name="Flex Auth", + repo_slug="flex-auth", + domain_slug="infotech", + lifecycle="active", + ) + ) + await session.commit() + + old_lookup = await client.get("/repos/flex-auth") + new_lookup = await client.get("/repos/access-engine") + assert old_lookup.status_code == new_lookup.status_code == 200 + assert old_lookup.json()["id"] == new_lookup.json()["id"] == repo["id"] + assert old_lookup.json()["requested_slug"] == "flex-auth" + assert old_lookup.json()["canonical_slug"] == "access-engine" + assert old_lookup.json()["slug_status"] == "alias" + assert old_lookup.json()["aliases"] == ["flex-auth"] + assert new_lookup.json()["slug_status"] == "canonical" + assert old_lookup.json()["stale_external_references"] == [ + { + "owner": "railiance-fabric", + "surface": "fabric_graph_nodes", + "field": "repo_slug", + "value": "flex-auth", + "count": 1, + "status": "stale", + "handoff": "owner update or re-ingest required", + } + ] + + dispatch = await client.get("/repos/access-engine/dispatch") + assert dispatch.status_code == 200, dispatch.text + assert dispatch.json()["repo_slug"] == "access-engine" + assert dispatch.json()["pending_interface_changes"][0]["title"] == change.title + + bound = await client.get(f"/workplans/{workplan['id']}") + assert bound.status_code == 200 + assert bound.json()["id"] == workplan["id"] + assert bound.json()["repo_id"] == repo["id"] + bound_task = await client.get(f"/tasks/{task.json()['id']}") + assert bound_task.status_code == 200 + assert bound_task.json()["id"] == task.json()["id"] + + +@pytest.mark.asyncio +async def test_message_history_is_immutable_and_old_slug_write_replays_once( + client, test_engine +): + domain = await create_test_domain(client) + repo = await create_test_repo(client, domain_slug=domain["slug"], slug="flex-auth") + await _rename_in_state_hub(test_engine, repo["id"]) + factory = async_sessionmaker(test_engine, class_=AsyncSession, expire_on_commit=False) + async with factory() as session: + historical = AgentMessage( + from_agent="security-review", + to_agent="flex-auth", + subject="Before cutover", + body="Historical recipient must not be rewritten", + ) + session.add(historical) + await session.commit() + historical_id = historical.id + + history = await client.get("/messages/", params={"to_agent": "access-engine"}) + assert history.status_code == 200 + assert [item["id"] for item in history.json()] == [str(historical_id)] + assert history.json()[0]["to_agent"] == "flex-auth" + + payload = { + "from_agent": "flex-auth", + "to_agent": "flex-auth", + "subject": "Queued during cutover", + "body": "Replay exactly once", + } + headers = {"Idempotency-Key": "repo-rename-old-slug-message"} + first = await client.post("/messages/", json=payload, headers=headers) + replay = await client.post("/messages/", json=payload, headers=headers) + assert first.status_code == replay.status_code == 201 + assert first.json()["id"] == replay.json()["id"] + assert first.json()["from_agent"] == "access-engine" + assert first.json()["to_agent"] == "access-engine" + + async with factory() as session: + assert await session.scalar( + select(func.count()) + .select_from(AgentMessage) + .where(AgentMessage.subject == "Queued during cutover") + ) == 1 + unchanged = await session.get(AgentMessage, historical_id) + assert unchanged.to_agent == "flex-auth" + + +@pytest.mark.asyncio +async def test_alias_aware_interface_catalog_sbom_and_telemetry_reads( + client, test_engine +): + domain = await create_test_domain(client) + repo = await create_test_repo(client, domain_slug=domain["slug"], slug="flex-auth") + await _rename_in_state_hub(test_engine, repo["id"]) + + interface = await client.post( + "/interface-changes/", + json={ + "repo_slug": "flex-auth", + "interface_type": "rest_api", + "change_type": "additive", + "title": "Canonical write", + "description": "New references use the current slug", + "affected_repo_slugs": ["flex-auth"], + }, + ) + assert interface.status_code == 201, interface.text + assert interface.json()["repo_slug"] == "access-engine" + assert interface.json()["affected_repo_slugs"] == ["access-engine"] + affected = await client.get( + "/interface-changes/", params={"affected_repo": "flex-auth"} + ) + assert [item["id"] for item in affected.json()] == [interface.json()["id"]] + + goal = await client.post( + "/repo-goals/", + json={ + "repo_id": repo["id"], + "title": "Preserve authorization continuity", + "description": "Keep Net Kingdom consumers online", + }, + ) + assert goal.status_code == 201, goal.text + goals = await client.get("/repo-goals/", params={"repo_slug": "flex-auth"}) + assert [item["id"] for item in goals.json()] == [goal.json()["id"]] + assert goals.json()[0]["repo_slug"] == "access-engine" + + capability = await client.post( + "/capability-catalog/", + json={ + "domain": "infotech", + "repo_slug": "flex-auth", + "capability_type": "authorization", + "title": "Policy decisions", + "description": "Net Kingdom authorization", + "keywords": ["authorize"], + }, + ) + assert capability.status_code == 201, capability.text + catalog = await client.get( + "/capability-catalog/", params={"repo_slug": "access-engine"} + ) + assert [item["id"] for item in catalog.json()] == [capability.json()["id"]] + assert catalog.json()[0]["repo_slug"] == "access-engine" + + service = await client.post( + "/services/catalog", + json={ + "slug": "access-engine-api", + "name": "Access Engine API", + "hosting_type": "self_hosted", + "development_type": "first_party", + "first_party": {"repo_slug": "flex-auth", "owning_domain": "infotech"}, + "self_hosted": {}, + }, + ) + assert service.status_code == 201, service.text + services = await client.get("/services/catalog", params={"repo_slug": "flex-auth"}) + assert [item["id"] for item in services.json()] == [service.json()["id"]] + + ingested = await client.post( + "/sbom/ingest/", + json={ + "repo_slug": "flex-auth", + "entries": [ + { + "package_name": "opa", + "package_version": "1.0", + "ecosystem": "go", + } + ], + }, + ) + assert ingested.status_code == 200, ingested.text + assert ingested.json()["repo_slug"] == "access-engine" + sbom = await client.get("/sbom/flex-auth") + assert sbom.status_code == 200 + assert sbom.json()["repo_slug"] == "access-engine" + + event = await client.post( + "/token-events/", + json={"repo_id": repo["id"], "tokens_in": 12, "tokens_out": 3}, + ) + assert event.status_code == 201, event.text + summary = await client.get( + "/token-events/by-repo/", params={"repo_slug": "flex-auth"} + ) + assert summary.status_code == 200 + assert len(summary.json()) == 1 + assert summary.json()[0]["repo_slug"] == "access-engine" + assert summary.json()[0]["tokens_total"] == 15 diff --git a/tests/test_work_record_identifier_migration.py b/tests/test_work_record_identifier_migration.py index 4dd2d3a..2189fbc 100644 --- a/tests/test_work_record_identifier_migration.py +++ b/tests/test_work_record_identifier_migration.py @@ -12,6 +12,7 @@ from api.models import Base from api.models.domain import Domain from api.models.managed_repo import ManagedRepo from api.models.progress_event import ProgressEvent +from api.models.repository_rename import RepositorySlug from api.models.task import Task from api.models.topic import Topic from api.models.token_event import TokenEvent @@ -276,6 +277,57 @@ async def test_repository_migration_cascades_and_reverses(test_engine): assert all(alias.reversed_at is not None for alias in aliases) +@pytest.mark.asyncio +async def test_repository_migration_accepts_protected_prior_slug(test_engine): + factory = async_sessionmaker(test_engine, class_=AsyncSession, expire_on_commit=False) + ids = await _seed_projection(factory, "flex-auth") + plan = _sealed_plan("flex-auth", ids["workplan"], ids["task"]) + + async with factory() as session: + repo = ( + await session.execute( + select(ManagedRepo).where(ManagedRepo.slug == "flex-auth") + ) + ).scalar_one() + repo.slug = "access-engine" + session.add_all( + [ + RepositorySlug( + repo_id=repo.id, + slug="flex-auth", + kind="alias", + protected=True, + ), + RepositorySlug( + repo_id=repo.id, + slug="access-engine", + kind="canonical", + protected=True, + ), + ] + ) + await session.commit() + + async with factory() as session: + result = await apply_repository_identifier_migration(session, plan, "flex-auth") + assert result.direction == "forward" + + async with factory() as session: + aliases = list((await session.execute(select(WorkRecordIdentifierAlias))).scalars()) + assert len(aliases) == 2 + assert {alias.repo_slug for alias in aliases} == {"flex-auth"} + repo = ( + await session.execute( + select(ManagedRepo).where(ManagedRepo.slug == "access-engine") + ) + ).scalar_one() + assert repo.id is not None + + async with factory() as session: + result = await reverse_repository_identifier_migration(session, plan, "flex-auth") + assert result.direction == "reverse" + + @pytest.mark.asyncio async def test_repository_migration_is_atomic_when_a_source_is_missing(test_engine): factory = async_sessionmaker(test_engine, class_=AsyncSession, expire_on_commit=False) diff --git a/workplans/STATE-WP-0085-repository-lineage-preserving-rename.md b/workplans/STATE-WP-0085-repository-lineage-preserving-rename.md index 6d01d36..06de423 100644 --- a/workplans/STATE-WP-0085-repository-lineage-preserving-rename.md +++ b/workplans/STATE-WP-0085-repository-lineage-preserving-rename.md @@ -298,7 +298,7 @@ guarding, and full migration upgrade/downgrade. The full Python suite passes ```task id: STATE-WP-0085-T04 -status: todo +status: done priority: high state_hub_task_id: "dcf7102e-ed7a-503e-bf78-db2149f72453" ``` @@ -331,6 +331,17 @@ Acceptance: - new bindings use the new slug while workplan and task UUIDs remain unchanged; - external stale references are named rather than silently rewritten. +Implemented one protected-slug resolver and applied it across repository reads, +dispatch, consistency/path binding, inbox routing, interface changes, +work-record identifier migration, goals, capability and service catalogs, SBOM, +token summaries, Fabric projection reads, and dashboard links. Repository +responses expose requested/canonical slug, canonical-or-alias status, aliases, +and named stale external projections. Historical message and interface values +remain unchanged; new writes are canonicalized. An idempotent old-slug message +replay is proven to persist once, while workplan/task UUIDs and repository FKs +remain stable. Verification: 775 Python tests and the 70-page dashboard build +pass (the build retains one pre-existing `/docs/intakes` broken-link warning). + ## Add the State Hub CLI and HelixForge orchestration contract ```task