state-hub/api/services/repository_rename.py
tegwick 82ea38b180
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: add repository rename lifecycle API
Assistant: codex
Assistant-Model: gpt-5.6-sol
Assistant-Session: 01a049a4-ee9f-78e1-9d66-2cb0f9bea3e3
2026-08-29 03:17:37 +02:00

1325 lines
52 KiB
Python

"""Fail-closed, UUID-addressed repository rename lifecycle (STATE-WP-0085-T03)."""
from __future__ import annotations
import base64
import hashlib
import hmac
import json
import uuid
from copy import deepcopy
from dataclasses import dataclass
from datetime import datetime, timedelta, timezone
from typing import Any
from sqlalchemy import or_, select
from sqlalchemy.exc import IntegrityError
from sqlalchemy.ext.asyncio import AsyncSession
from api.config import settings
from api.models.agent_message import AgentMessage
from api.models.capability_catalog import CapabilityCatalog
from api.models.decision import Decision
from api.models.fabric_graph import FabricGraphNode
from api.models.interface_change import InterfaceChange
from api.models.managed_repo import ManagedRepo
from api.models.progress_event import ProgressEvent
from api.models.repository_rename import (
RepositoryForgeIdentity,
RepositoryRenameOperation,
RepositorySlug,
)
from api.models.sbom_entry import SBOMEntry
from api.models.sbom_snapshot import SBOMSnapshot
from api.models.service_catalog import ServiceFirstParty
from api.models.task import Task
from api.models.token_event import TokenEvent
from api.models.workplan import Workplan
from api.schemas.repository_rename import (
ForgeIdentityVerifyRequest,
RepositoryRenameOperationCreate,
RepositoryRenamePhaseApply,
RepositoryRenamePreflightRequest,
)
from api.services.forge_repository import (
ForgeRepositoryConflict,
ForgeRepositoryGateway,
ForgeRepositorySnapshot,
ForgeRepositoryUnreadable,
)
ACTIVE_OPERATION_PHASES = {
"draft",
"preflighted",
"forge-renamed",
"statehub-rebound",
"source-synced",
"consumers-verified",
"rollback-preflight",
}
FORWARD_PHASES = [
"preflighted",
"forge-renamed",
"statehub-rebound",
"source-synced",
"consumers-verified",
"completed",
]
class RenameLifecycleError(RuntimeError):
status_code = 409
code = "repository_rename_conflict"
def __init__(self, message: str, *, details: dict[str, Any] | None = None):
super().__init__(message)
self.details = details or {}
class RenameNotFound(RenameLifecycleError):
status_code = 404
code = "repository_rename_not_found"
class RenamePreconditionFailed(RenameLifecycleError):
status_code = 412
code = "repository_rename_precondition_failed"
class RenameServiceUnavailable(RenameLifecycleError):
status_code = 503
code = "repository_rename_service_unavailable"
def utcnow() -> datetime:
return datetime.now(timezone.utc)
def _json_default(value: Any) -> str:
if isinstance(value, datetime):
return value.astimezone(timezone.utc).isoformat()
if isinstance(value, uuid.UUID):
return str(value)
if hasattr(value, "value"):
return str(value.value)
return str(value)
def canonical_json(value: Any) -> bytes:
return json.dumps(
value,
sort_keys=True,
separators=(",", ":"),
ensure_ascii=True,
default=_json_default,
).encode("utf-8")
def checksum(value: Any) -> str:
return hashlib.sha256(canonical_json(value)).hexdigest()
def _b64encode(value: bytes) -> str:
return base64.urlsafe_b64encode(value).rstrip(b"=").decode("ascii")
def _b64decode(value: str) -> bytes:
return base64.urlsafe_b64decode(value + "=" * (-len(value) % 4))
def _sign_preflight(payload: dict[str, Any]) -> str:
secret = settings.repository_rename_preflight_secret
if not secret:
raise RenameServiceUnavailable(
"Repository rename preflight signing is not configured"
)
encoded = _b64encode(canonical_json(payload))
signature = hmac.new(secret.encode(), encoded.encode(), hashlib.sha256).digest()
return f"{encoded}.{_b64encode(signature)}"
def _verify_preflight_token(token: str) -> dict[str, Any]:
secret = settings.repository_rename_preflight_secret
if not secret:
raise RenameServiceUnavailable(
"Repository rename preflight signing is not configured"
)
try:
encoded, supplied = token.split(".", 1)
expected = hmac.new(secret.encode(), encoded.encode(), hashlib.sha256).digest()
if not hmac.compare_digest(expected, _b64decode(supplied)):
raise ValueError("signature")
payload = json.loads(_b64decode(encoded))
expires_at = datetime.fromisoformat(payload["expires_at"])
except (ValueError, KeyError, json.JSONDecodeError) as exc:
raise RenamePreconditionFailed("Invalid repository rename preflight token") from exc
if expires_at <= utcnow():
raise RenamePreconditionFailed("Repository rename preflight token has expired")
return payload
def confirmation_for(repo_id: uuid.UUID, old_slug: str, new_slug: str) -> str:
return f"rename:{repo_id}:{old_slug}:{new_slug}"
def rollback_confirmation_for(operation_id: uuid.UUID) -> str:
return f"rollback:{operation_id}"
async def _id_rows(session: AsyncSession, query) -> list[str]:
result = await session.execute(query)
return sorted(str(value) for value in result.scalars().all())
def _record_summary(ids: list[str]) -> dict[str, Any]:
return {"count": len(ids), "ids": ids, "checksum": checksum(ids)}
async def collect_continuity_baseline(
session: AsyncSession, repo_id: uuid.UUID
) -> dict[str, Any]:
workplans = await _id_rows(
session, select(Workplan.id).where(Workplan.repo_id == repo_id)
)
workplan_uuids = [uuid.UUID(value) for value in workplans]
tasks = await _id_rows(
session,
select(Task.id).where(
Task.workplan_id.in_(workplan_uuids) if workplan_uuids else False
),
)
task_uuids = [uuid.UUID(value) for value in tasks]
progress = await _id_rows(
session,
select(ProgressEvent.id).where(
or_(
ProgressEvent.workplan_id.in_(workplan_uuids)
if workplan_uuids
else False,
ProgressEvent.task_id.in_(task_uuids) if task_uuids else False,
)
),
)
decisions = await _id_rows(
session,
select(Decision.id).where(
Decision.workplan_id.in_(workplan_uuids) if workplan_uuids else False
),
)
token_events = await _id_rows(
session,
select(TokenEvent.id).where(
or_(
TokenEvent.repo_id == repo_id,
TokenEvent.workplan_id.in_(workplan_uuids)
if workplan_uuids
else False,
TokenEvent.task_id.in_(task_uuids) if task_uuids else False,
)
),
)
sbom_snapshots = await _id_rows(
session, select(SBOMSnapshot.id).where(SBOMSnapshot.repo_id == repo_id)
)
sbom_entries = await _id_rows(
session, select(SBOMEntry.id).where(SBOMEntry.repo_id == repo_id)
)
services = await _id_rows(
session,
select(ServiceFirstParty.service_id).where(ServiceFirstParty.repo_id == repo_id),
)
capabilities = await _id_rows(
session,
select(CapabilityCatalog.id).where(CapabilityCatalog.repo_id == repo_id),
)
interface_changes = await _id_rows(
session,
select(InterfaceChange.id).where(InterfaceChange.repo_id == repo_id),
)
bindings = await _id_rows(
session,
select(Workplan.id).where(
Workplan.repo_id == repo_id,
Workplan.backing_relative_path.is_not(None),
),
)
records = {
"workplans": _record_summary(workplans),
"tasks": _record_summary(tasks),
"progress_events": _record_summary(progress),
"decisions": _record_summary(decisions),
"token_events": _record_summary(token_events),
"sbom_snapshots": _record_summary(sbom_snapshots),
"sbom_entries": _record_summary(sbom_entries),
"services": _record_summary(services),
"capabilities": _record_summary(capabilities),
"interface_changes": _record_summary(interface_changes),
"workplan_bindings": _record_summary(bindings),
}
records["continuity_checksum"] = checksum(records)
return records
async def _active_work(session: AsyncSession, repo_id: uuid.UUID) -> dict[str, Any]:
workplans_result = await session.execute(
select(Workplan.id, Workplan.slug, Workplan.status)
.where(
Workplan.repo_id == repo_id,
Workplan.status.not_in(("finished", "archived")),
)
.order_by(Workplan.slug)
)
workplans = [
{"id": str(row.id), "slug": row.slug, "status": row.status}
for row in workplans_result.all()
]
workplan_ids = [uuid.UUID(row["id"]) for row in workplans]
tasks_result = await session.execute(
select(Task.id, Task.workplan_id, Task.record_id, Task.status)
.where(
Task.workplan_id.in_(workplan_ids) if workplan_ids else False,
Task.status.not_in(("done", "cancel")),
)
.order_by(Task.id)
)
tasks = [
{
"id": str(row.id),
"workplan_id": str(row.workplan_id),
"record_id": row.record_id,
"status": getattr(row.status, "value", str(row.status)),
}
for row in tasks_result.all()
]
return {"workplans": workplans, "tasks": tasks}
async def _affected_records(
session: AsyncSession, repo_id: uuid.UUID, old_slug: str
) -> dict[str, Any]:
messages = await _id_rows(
session,
select(AgentMessage.id).where(
or_(AgentMessage.from_agent == old_slug, AgentMessage.to_agent == old_slug)
),
)
interface_mentions = await _id_rows(
session,
select(InterfaceChange.id).where(
InterfaceChange.affected_repo_slugs.contains([old_slug])
),
)
fabric_nodes = await _id_rows(
session,
select(FabricGraphNode.id).where(
or_(
FabricGraphNode.repo_slug == old_slug,
FabricGraphNode.source_repo_slug == old_slug,
)
),
)
return {
"messages": _record_summary(messages),
"interface_change_mentions": _record_summary(interface_mentions),
"fabric_nodes": _record_summary(fabric_nodes),
"external_handoffs": [
item
for item, ids in (
("interface-change-consumers", interface_mentions),
("fabric-graph-projections", fabric_nodes),
)
if ids
],
}
def _new_remote_url(repo: ManagedRepo, identity: RepositoryForgeIdentity, new_slug: str) -> str:
base = str(identity.forge_instance).rstrip("/")
suffix = ".git" if (repo.remote_url or "").endswith(".git") else ""
return f"{base}/{identity.forge_owner}/{new_slug}{suffix}"
async def _load_repo_identity(
session: AsyncSession, repo_id: uuid.UUID
) -> tuple[ManagedRepo, RepositoryForgeIdentity | None]:
repo = await session.get(ManagedRepo, repo_id)
if repo is None:
raise RenameNotFound(f"Repository {repo_id} was not found")
result = await session.execute(
select(RepositoryForgeIdentity).where(
RepositoryForgeIdentity.repo_id == repo_id
)
)
return repo, result.scalar_one_or_none()
async def verify_forge_identity(
session: AsyncSession,
gateway: ForgeRepositoryGateway,
repo_id: uuid.UUID,
body: ForgeIdentityVerifyRequest,
) -> RepositoryForgeIdentity:
repo, identity = await _load_repo_identity(session, repo_id)
if identity is None:
raise RenamePreconditionFailed("Repository has no Forge identity row")
try:
snapshot = await gateway.inspect(
instance=body.forge_instance,
owner=body.forge_owner,
name=repo.slug,
)
except ForgeRepositoryUnreadable as exc:
raise RenamePreconditionFailed(str(exc)) from exc
if snapshot is None:
raise RenamePreconditionFailed("Forge repository is absent or unreadable")
if snapshot.repository_id != body.forge_repository_id:
raise RenamePreconditionFailed(
"Forge repository ID does not match the asserted immutable ID",
details={
"expected_forge_repository_id": body.forge_repository_id,
"actual_forge_repository_id": snapshot.repository_id,
},
)
if identity.verification_state == "verified":
actual = (
identity.provider,
identity.forge_instance,
identity.forge_owner,
identity.forge_repository_id,
)
supplied = (
body.provider,
body.forge_instance.rstrip("/"),
body.forge_owner,
body.forge_repository_id,
)
if actual != supplied:
raise RenamePreconditionFailed("Verified Forge identity is immutable")
return identity
identity.provider = body.provider
identity.forge_instance = body.forge_instance.rstrip("/")
identity.forge_owner = body.forge_owner
identity.forge_repository_id = body.forge_repository_id
identity.verification_state = "verified"
identity.verified_at = utcnow()
identity.verified_by = body.verified_by
identity.verification_evidence = {
"forge_snapshot": snapshot.as_dict(),
"verified_repo_slug": repo.slug,
}
await session.commit()
await session.refresh(identity)
return identity
async def build_preflight(
session: AsyncSession,
gateway: ForgeRepositoryGateway,
repo_id: uuid.UUID,
body: RepositoryRenamePreflightRequest,
*,
issue_token: bool = True,
) -> dict[str, Any]:
repo, identity = await _load_repo_identity(session, repo_id)
now = utcnow()
blockers: list[dict[str, Any]] = []
warnings: list[dict[str, Any]] = []
old_snapshot: ForgeRepositorySnapshot | None = None
target_snapshot: ForgeRepositorySnapshot | None = None
if issue_token and not settings.repository_rename_preflight_secret:
blockers.append(
{
"code": "preflight_signing_unavailable",
"message": "Repository rename preflight signing is not configured",
}
)
if repo.slug == body.new_slug:
blockers.append({"code": "same_slug", "message": "Target slug equals canonical slug"})
claimed = await session.execute(
select(RepositorySlug).where(RepositorySlug.slug == body.new_slug)
)
claimed_row = claimed.scalar_one_or_none()
old_claim = await session.execute(
select(RepositorySlug).where(RepositorySlug.slug == repo.slug)
)
old_claim_row = old_claim.scalar_one_or_none()
if (
old_claim_row is None
or old_claim_row.repo_id != repo.id
or old_claim_row.kind != "canonical"
):
blockers.append(
{
"code": "statehub_canonical_inconsistent",
"message": "Managed repository slug is not its protected canonical route",
}
)
if claimed_row is not None:
blockers.append(
{
"code": "statehub_slug_claimed",
"message": "Target slug is already protected in State Hub",
"repo_id": str(claimed_row.repo_id),
"kind": claimed_row.kind,
}
)
active_result = await session.execute(
select(RepositoryRenameOperation).where(
RepositoryRenameOperation.repo_id == repo_id,
RepositoryRenameOperation.phase.in_(ACTIVE_OPERATION_PHASES),
)
)
active_operations = list(active_result.scalars().all())
if active_operations:
blockers.append(
{
"code": "active_rename_operation",
"message": "Repository already has an active rename operation",
"operation_ids": [str(item.id) for item in active_operations],
}
)
target_operation_result = await session.execute(
select(RepositoryRenameOperation.id, RepositoryRenameOperation.repo_id).where(
RepositoryRenameOperation.new_slug == body.new_slug,
RepositoryRenameOperation.phase.in_(ACTIVE_OPERATION_PHASES),
RepositoryRenameOperation.repo_id != repo_id,
)
)
target_operations = list(target_operation_result.all())
if target_operations:
blockers.append(
{
"code": "target_held_by_active_operation",
"message": "Target slug is reserved by another active rename operation",
"operations": [
{"operation_id": str(row.id), "repo_id": str(row.repo_id)}
for row in target_operations
],
}
)
if body.queued_edge_writes:
blockers.append(
{
"code": "queued_edge_writes",
"message": "Caller reported edge writes that must be replayed or resolved first",
"count": len(body.queued_edge_writes),
}
)
if identity is None or identity.verification_state != "verified":
blockers.append(
{
"code": "forge_identity_unverified",
"message": "Repository Forge identity is not verified",
}
)
else:
try:
old_snapshot = await gateway.inspect(
instance=str(identity.forge_instance),
owner=str(identity.forge_owner),
name=repo.slug,
)
target_snapshot = await gateway.inspect(
instance=str(identity.forge_instance),
owner=str(identity.forge_owner),
name=body.new_slug,
)
except ForgeRepositoryUnreadable as exc:
blockers.append(
{"code": "forge_unreadable", "message": str(exc)}
)
if old_snapshot is None:
blockers.append(
{"code": "forge_source_absent", "message": "Canonical Forge repository is absent or unreadable"}
)
else:
if old_snapshot.repository_id != identity.forge_repository_id:
blockers.append(
{
"code": "wrong_forge_repository_id",
"message": "Forge repository ID differs from verified identity",
"expected": identity.forge_repository_id,
"actual": old_snapshot.repository_id,
}
)
if not old_snapshot.projection_readable or not old_snapshot.projection_source_present:
blockers.append(
{
"code": "projection_unreadable",
"message": "Forge work-record projection source is not readable",
}
)
if target_snapshot is not None:
blockers.append(
{
"code": "forge_target_claimed",
"message": "Target Forge repository name is already claimed",
"forge_repository_id": target_snapshot.repository_id,
}
)
baselines = await collect_continuity_baseline(session, repo_id)
active_work = await _active_work(session, repo_id)
affected = await _affected_records(session, repo_id, repo.slug)
affected.update(
{
"workplan_bindings": baselines["workplan_bindings"],
"workplans": baselines["workplans"],
"tasks": baselines["tasks"],
"progress_events": baselines["progress_events"],
"decisions": baselines["decisions"],
"token_events": baselines["token_events"],
"interface_changes": baselines["interface_changes"],
"sbom_snapshots": baselines["sbom_snapshots"],
"sbom_entries": baselines["sbom_entries"],
"services": baselines["services"],
"capability_entries": baselines["capabilities"],
}
)
if active_work["workplans"]:
warnings.append(
{
"code": "active_work_present",
"message": "Active work must be quiesced or explicitly coordinated during cutover",
"workplan_count": len(active_work["workplans"]),
"task_count": len(active_work["tasks"]),
}
)
identity_dict = (
{
"id": str(identity.id),
"verification_state": identity.verification_state,
"provider": identity.provider,
"forge_instance": identity.forge_instance,
"forge_owner": identity.forge_owner,
"forge_repository_id": identity.forge_repository_id,
}
if identity is not None
else None
)
new_remote = _new_remote_url(repo, identity, body.new_slug) if identity else None
proposed = [
{"phase": "forge-renamed", "system": "forgejo", "field": "name", "from": repo.slug, "to": body.new_slug},
{"phase": "statehub-rebound", "system": "state-hub", "field": "managed_repos.slug", "from": repo.slug, "to": body.new_slug},
{"phase": "statehub-rebound", "system": "state-hub", "field": "managed_repos.remote_url", "from": repo.remote_url, "to": new_remote},
{"phase": "statehub-rebound", "system": "state-hub", "field": "repository_slugs", "from": f"{repo.slug}:canonical", "to": f"{repo.slug}:alias,{body.new_slug}:canonical"},
]
retained = [
{"field": "managed_repos.id", "value": str(repo.id)},
{"field": "repository_forge_identities.forge_repository_id", "value": identity.forge_repository_id if identity else None},
{"field": "repository_slugs.alias", "value": repo.slug},
{"field": "managed_repos.local_path", "value": repo.local_path},
{"field": "managed_repos.host_paths", "value": repo.host_paths},
{"field": "telemetry_and_work_record_ids", "value": baselines["continuity_checksum"]},
]
stable_report = {
"schema_version": "state-hub.repository-rename-preflight.v1",
"repo_id": str(repo.id),
"old_slug": repo.slug,
"new_slug": body.new_slug,
"safe_to_apply": not blockers,
"blockers": blockers,
"warnings": warnings,
"current": {
"statehub": {
"repo_id": str(repo.id),
"slug": repo.slug,
"name": repo.name,
"remote_url": repo.remote_url,
"local_path": repo.local_path,
"host_paths": repo.host_paths,
"canonical_available": bool(
old_claim_row
and old_claim_row.repo_id == repo.id
and old_claim_row.kind == "canonical"
),
"slug_registry": (
{
"id": str(old_claim_row.id),
"kind": old_claim_row.kind,
"protected": old_claim_row.protected,
}
if old_claim_row
else None
),
},
"forge_identity": identity_dict,
"forge_available": old_snapshot is not None,
"forge": old_snapshot.as_dict() if old_snapshot else None,
},
"target": {
"slug": body.new_slug,
"statehub_available": claimed_row is None,
"forge_available": target_snapshot is None,
"forge": target_snapshot.as_dict() if target_snapshot else None,
},
"baselines": baselines,
"active_work": active_work,
"affected": affected,
"queued_edge_writes": [item.model_dump(mode="json") for item in body.queued_edge_writes],
"proposed_mutations": proposed,
"retained_history": retained,
}
report_checksum = checksum(stable_report)
expires_at = now + timedelta(seconds=settings.repository_rename_preflight_ttl_seconds)
token = None
if issue_token and not blockers:
token = _sign_preflight(
{
"version": 1,
"repo_id": str(repo.id),
"old_slug": repo.slug,
"new_slug": body.new_slug,
"forge_repository_id": identity.forge_repository_id if identity else None,
"source_commit": old_snapshot.head_commit if old_snapshot else None,
"report_checksum": report_checksum,
"issued_at": now.isoformat(),
"expires_at": expires_at.isoformat(),
}
)
return {
**stable_report,
"report_checksum": report_checksum,
"preflight_token": token,
"preflighted_at": now,
"expires_at": expires_at if token else None,
}
async def create_operation(
session: AsyncSession,
gateway: ForgeRepositoryGateway,
repo_id: uuid.UUID,
body: RepositoryRenameOperationCreate,
) -> RepositoryRenameOperation:
token = _verify_preflight_token(body.preflight_token)
if token.get("repo_id") != str(repo_id) or token.get("new_slug") != body.new_slug:
raise RenamePreconditionFailed("Preflight token does not address this rename")
preflight = await build_preflight(
session,
gateway,
repo_id,
RepositoryRenamePreflightRequest(
new_slug=body.new_slug, queued_edge_writes=body.queued_edge_writes
),
issue_token=False,
)
if not preflight["safe_to_apply"]:
raise RenamePreconditionFailed(
"Repository rename preflight is no longer safe",
details={"blockers": preflight["blockers"]},
)
if preflight["report_checksum"] != token.get("report_checksum"):
raise RenamePreconditionFailed(
"Repository rename preflight evidence is stale",
details={
"expected_report_checksum": token.get("report_checksum"),
"actual_report_checksum": preflight["report_checksum"],
},
)
old_slug = preflight["old_slug"]
if body.confirmation != confirmation_for(repo_id, old_slug, body.new_slug):
raise RenamePreconditionFailed("Explicit repository rename confirmation is incorrect")
repo, identity = await _load_repo_identity(session, repo_id)
if identity is None or identity.verification_state != "verified":
raise RenamePreconditionFailed("Repository Forge identity is not verified")
forge = preflight["current"]["forge"]
now = utcnow()
operation = RepositoryRenameOperation(
repo_id=repo.id,
forge_identity_id=identity.id,
forge_identity_state="verified",
expected_provider=str(identity.provider),
expected_forge_instance=str(identity.forge_instance),
expected_forge_owner=str(identity.forge_owner),
expected_forge_repository_id=int(identity.forge_repository_id),
expected_source_commit=forge["head_commit"],
expected_default_branch=forge["default_branch"],
old_slug=repo.slug,
new_slug=body.new_slug,
old_coordinates={
"slug": repo.slug,
"name": repo.name,
"remote_url": repo.remote_url,
"local_path": repo.local_path,
"host_paths": repo.host_paths,
},
new_coordinates={
"slug": body.new_slug,
"name": body.new_slug,
"remote_url": _new_remote_url(repo, identity, body.new_slug),
"local_path": repo.local_path,
"host_paths": repo.host_paths,
},
phase="preflighted",
actor=body.actor,
phase_changed_at=now,
preflighted_at=datetime.fromisoformat(token["issued_at"]),
preflight_expires_at=datetime.fromisoformat(token["expires_at"]),
evidence={
"preflight": {
"report_checksum": preflight["report_checksum"],
"baselines": preflight["baselines"],
"affected": preflight["affected"],
"warnings": preflight["warnings"],
},
"phases": {
"preflighted": {"at": now.isoformat(), "resumed": False}
},
},
)
session.add(operation)
try:
await session.commit()
except IntegrityError as exc:
await session.rollback()
raise RenamePreconditionFailed(
"A conflicting repository rename operation was created"
) from exc
await session.refresh(operation)
return operation
async def load_operation(
session: AsyncSession,
repo_id: uuid.UUID,
operation_id: uuid.UUID,
*,
for_update: bool = False,
) -> RepositoryRenameOperation:
query = select(RepositoryRenameOperation).where(
RepositoryRenameOperation.id == operation_id,
RepositoryRenameOperation.repo_id == repo_id,
)
if for_update:
query = query.with_for_update()
result = await session.execute(query)
operation = result.scalar_one_or_none()
if operation is None:
raise RenameNotFound(f"Repository rename operation {operation_id} was not found")
return operation
async def list_operations(
session: AsyncSession, repo_id: uuid.UUID, *, active_only: bool = False
) -> list[RepositoryRenameOperation]:
repo = await session.get(ManagedRepo, repo_id)
if repo is None:
raise RenameNotFound(f"Repository {repo_id} was not found")
query = (
select(RepositoryRenameOperation)
.where(RepositoryRenameOperation.repo_id == repo_id)
.order_by(RepositoryRenameOperation.created_at.desc())
)
if active_only:
query = query.where(RepositoryRenameOperation.phase.in_(ACTIVE_OPERATION_PHASES))
result = await session.execute(query)
return list(result.scalars().all())
def _assert_snapshot(operation: RepositoryRenameOperation, snapshot: ForgeRepositorySnapshot) -> None:
if snapshot.repository_id != operation.expected_forge_repository_id:
raise RenamePreconditionFailed(
"Forge repository ID changed",
details={
"expected": operation.expected_forge_repository_id,
"actual": snapshot.repository_id,
},
)
if snapshot.default_branch != operation.expected_default_branch:
raise RenamePreconditionFailed("Forge default branch changed")
if snapshot.head_commit != operation.expected_source_commit:
raise RenamePreconditionFailed(
"Forge default-branch head moved",
details={
"expected": operation.expected_source_commit,
"actual": snapshot.head_commit,
},
)
if not snapshot.projection_readable or not snapshot.projection_source_present:
raise RenamePreconditionFailed("Forge work-record projection is unreadable")
async def _forge_at(
gateway: ForgeRepositoryGateway,
operation: RepositoryRenameOperation,
name: str,
) -> ForgeRepositorySnapshot | None:
try:
return await gateway.inspect(
instance=operation.expected_forge_instance,
owner=operation.expected_forge_owner,
name=name,
)
except ForgeRepositoryUnreadable as exc:
raise RenamePreconditionFailed(str(exc)) from exc
async def _apply_forge_rename(
gateway: ForgeRepositoryGateway, operation: RepositoryRenameOperation
) -> dict[str, Any]:
old = await _forge_at(gateway, operation, operation.old_slug)
new = await _forge_at(gateway, operation, operation.new_slug)
if new is not None:
_assert_snapshot(operation, new)
if old is not None:
raise RenamePreconditionFailed("Both old and new Forge names are claimed")
return {"resumed": True, "forge": new.as_dict()}
if old is None:
raise RenamePreconditionFailed("Forge repository is absent at both expected names")
_assert_snapshot(operation, old)
try:
renamed = await gateway.rename(
instance=operation.expected_forge_instance,
owner=operation.expected_forge_owner,
old_name=operation.old_slug,
new_name=operation.new_slug,
)
except (ForgeRepositoryUnreadable, ForgeRepositoryConflict) as exc:
raise RenamePreconditionFailed(str(exc)) from exc
_assert_snapshot(operation, renamed)
return {"resumed": False, "forge": renamed.as_dict()}
async def _apply_statehub_rebind(
session: AsyncSession, operation: RepositoryRenameOperation
) -> dict[str, Any]:
repo = await session.get(ManagedRepo, operation.repo_id, with_for_update=True)
if repo is None:
raise RenameNotFound("Managed repository disappeared")
if repo.id != operation.repo_id:
raise RenamePreconditionFailed("Managed repository UUID changed")
old_result = await session.execute(
select(RepositorySlug)
.where(RepositorySlug.slug == operation.old_slug)
.with_for_update()
)
new_result = await session.execute(
select(RepositorySlug)
.where(RepositorySlug.slug == operation.new_slug)
.with_for_update()
)
old_slug = old_result.scalar_one_or_none()
new_slug = new_result.scalar_one_or_none()
if repo.slug == operation.new_slug:
if (
old_slug is None
or old_slug.repo_id != repo.id
or old_slug.kind != "alias"
or new_slug is None
or new_slug.repo_id != repo.id
or new_slug.kind != "canonical"
):
raise RenamePreconditionFailed("State Hub rename state is internally inconsistent")
return {"resumed": True, "repo_id": str(repo.id)}
if repo.slug != operation.old_slug:
raise RenamePreconditionFailed("Managed repository slug changed outside this operation")
if old_slug is None or old_slug.repo_id != repo.id or old_slug.kind != "canonical":
raise RenamePreconditionFailed("Old canonical slug record is missing")
if new_slug is not None:
raise RenamePreconditionFailed("Target State Hub slug became claimed")
old_slug.kind = "alias"
old_slug.protected = True
old_slug.source_operation_id = operation.id
await session.flush()
session.add(
RepositorySlug(
repo_id=repo.id,
slug=operation.new_slug,
kind="canonical",
protected=True,
source_operation_id=operation.id,
)
)
repo.slug = operation.new_slug
repo.name = operation.new_coordinates["name"]
repo.remote_url = operation.new_coordinates["remote_url"]
await session.flush()
return {"resumed": False, "repo_id": str(repo.id)}
async def verify_operation(
session: AsyncSession,
gateway: ForgeRepositoryGateway,
operation: RepositoryRenameOperation,
) -> dict[str, Any]:
checks: list[dict[str, Any]] = []
def add(name: str, ok: bool, expected: Any, actual: Any) -> None:
checks.append({"name": name, "ok": ok, "expected": expected, "actual": actual})
repo = await session.get(ManagedRepo, operation.repo_id)
add("managed_repository_exists", repo is not None, str(operation.repo_id), str(repo.id) if repo else None)
add("managed_repository_uuid", bool(repo and repo.id == operation.repo_id), str(operation.repo_id), str(repo.id) if repo else None)
expected_name = operation.old_slug if operation.phase == "preflighted" else operation.new_slug
if operation.phase == "rolled-back":
expected_name = operation.old_slug
try:
forge = await _forge_at(gateway, operation, expected_name)
except RenameLifecycleError as exc:
forge = None
add("forge_readable", False, True, str(exc))
else:
add("forge_readable", forge is not None, True, forge is not None)
if forge is not None:
add("forge_repository_id", forge.repository_id == operation.expected_forge_repository_id, operation.expected_forge_repository_id, forge.repository_id)
add("source_commit", forge.head_commit == operation.expected_source_commit, operation.expected_source_commit, forge.head_commit)
add("default_branch", forge.default_branch == operation.expected_default_branch, operation.expected_default_branch, forge.default_branch)
add("projection_readable", forge.projection_readable and forge.projection_source_present, True, forge.projection_readable and forge.projection_source_present)
expected_canonical = operation.old_slug if operation.phase in {"preflighted", "forge-renamed", "rolled-back"} else operation.new_slug
canonical_result = await session.execute(
select(RepositorySlug).where(
RepositorySlug.repo_id == operation.repo_id,
RepositorySlug.kind == "canonical",
)
)
canonical = canonical_result.scalar_one_or_none()
add("canonical_route", bool(canonical and canonical.slug == expected_canonical), expected_canonical, canonical.slug if canonical else None)
current = await collect_continuity_baseline(session, operation.repo_id)
baseline = (operation.evidence.get("preflight") or {}).get("baselines") or {}
expected_checksum = baseline.get("continuity_checksum", "")
missing: dict[str, list[str]] = {}
counts: dict[str, dict[str, int]] = {}
for record_type, baseline_record in baseline.items():
if record_type == "continuity_checksum" or not isinstance(baseline_record, dict):
continue
current_record = current.get(record_type) or {"ids": [], "count": 0}
absent = sorted(set(baseline_record.get("ids") or []) - set(current_record.get("ids") or []))
if absent:
missing[record_type] = absent
counts[record_type] = {
"baseline": int(baseline_record.get("count") or 0),
"current": int(current_record.get("count") or 0),
}
# New telemetry may legitimately arrive while a phased operation is being
# executed. Continuity means every baseline identity still exists; it does
# not freeze the repository's append-only history at the preflight count.
add("relationship_continuity", not missing, {"missing": {}}, {"missing": missing})
add(
"record_counts_non_decreasing",
all(item["current"] >= item["baseline"] for item in counts.values()),
{name: item["baseline"] for name, item in counts.items()},
{name: item["current"] for name, item in counts.items()},
)
return {
"operation_id": operation.id,
"repo_id": operation.repo_id,
"phase": operation.phase,
"ok": all(item["ok"] for item in checks),
"checks": checks,
"baseline_checksum": expected_checksum,
"current_checksum": current["continuity_checksum"],
}
def _record_phase(operation: RepositoryRenameOperation, phase: str, evidence: dict[str, Any]) -> None:
now = utcnow()
journal = deepcopy(operation.evidence or {})
phases = dict(journal.get("phases") or {})
# Verification payloads use native UUID/datetime values for API response
# typing. The operation journal is JSONB, so normalize at the boundary
# instead of leaking storage concerns through every verifier.
safe_evidence = json.loads(canonical_json(evidence))
phases[phase] = {"at": now.isoformat(), **safe_evidence}
journal["phases"] = phases
operation.evidence = journal
operation.phase = phase
operation.phase_changed_at = now
operation.error_code = None
operation.error_message = None
operation.error_details = None
operation.error_at = None
if phase == "completed":
operation.completed_at = now
if phase == "rolled-back":
operation.rolled_back_at = now
def _phase_index(phase: str) -> int:
try:
return FORWARD_PHASES.index(phase)
except ValueError:
return -1
async def apply_phase(
session: AsyncSession,
gateway: ForgeRepositoryGateway,
repo_id: uuid.UUID,
operation_id: uuid.UUID,
requested_phase: str,
body: RepositoryRenamePhaseApply,
) -> tuple[RepositoryRenameOperation, bool]:
if requested_phase not in FORWARD_PHASES[1:]:
raise RenamePreconditionFailed(f"Unsupported forward phase '{requested_phase}'")
operation = await load_operation(session, repo_id, operation_id, for_update=True)
if body.confirmation != confirmation_for(repo_id, operation.old_slug, operation.new_slug):
raise RenamePreconditionFailed("Explicit repository rename confirmation is incorrect")
current_index = _phase_index(operation.phase)
requested_index = _phase_index(requested_phase)
if current_index >= requested_index >= 0:
return operation, True
if operation.phase != body.expected_phase:
raise RenamePreconditionFailed(
"Repository rename phase compare-and-set failed",
details={"expected": body.expected_phase, "actual": operation.phase},
)
if requested_index != current_index + 1:
raise RenamePreconditionFailed(
"Repository rename phases cannot be skipped",
details={"current": operation.phase, "requested": requested_phase},
)
try:
if requested_phase == "forge-renamed":
evidence = await _apply_forge_rename(gateway, operation)
elif requested_phase == "statehub-rebound":
forge = await _forge_at(gateway, operation, operation.new_slug)
if forge is None:
raise RenamePreconditionFailed("Forge rename is not observable")
_assert_snapshot(operation, forge)
evidence = await _apply_statehub_rebind(session, operation)
elif requested_phase == "source-synced":
if not body.evidence:
raise RenamePreconditionFailed(
"Source synchronization requires operator evidence"
)
forge = await _forge_at(gateway, operation, operation.new_slug)
if forge is None:
raise RenamePreconditionFailed("Renamed Forge repository is absent")
_assert_snapshot(operation, forge)
evidence = {"forge": forge.as_dict(), "operator_evidence": body.evidence}
elif requested_phase == "consumers-verified":
if not body.checks or not all(body.checks.values()):
raise RenamePreconditionFailed(
"Consumer verification requires a non-empty set of passing checks",
details={"checks": body.checks},
)
verification = await verify_operation(session, gateway, operation)
if not verification["ok"]:
raise RenamePreconditionFailed(
"Continuity verification failed before consumer acceptance",
details=verification,
)
evidence = {"checks": body.checks, "verification": verification, "operator_evidence": body.evidence}
else: # completed
verification = await verify_operation(session, gateway, operation)
if not verification["ok"]:
raise RenamePreconditionFailed(
"Repository rename verification failed", details=verification
)
evidence = {"verification": verification}
_record_phase(operation, requested_phase, evidence)
await session.commit()
await session.refresh(operation)
return operation, False
except RenameLifecycleError as exc:
# Roll back every State Hub write attempted by the phase. A Forge
# rename may already have committed externally; the unchanged journal
# phase is precisely what lets the next request discover and resume it.
await session.rollback()
failed = await load_operation(session, repo_id, operation_id, for_update=True)
failed.error_code = exc.code
failed.error_message = str(exc)
failed.error_details = exc.details
failed.error_at = utcnow()
await session.commit()
raise
except IntegrityError as exc:
await session.rollback()
failure = RenamePreconditionFailed(
"Repository rename phase lost a database compare-and-set race"
)
failed = await load_operation(session, repo_id, operation_id, for_update=True)
failed.error_code = failure.code
failed.error_message = str(failure)
failed.error_details = failure.details
failed.error_at = utcnow()
await session.commit()
raise failure from exc
async def rollback_preflight(
session: AsyncSession,
gateway: ForgeRepositoryGateway,
repo_id: uuid.UUID,
operation_id: uuid.UUID,
*,
expected_phase: str,
confirmation: str,
) -> tuple[RepositoryRenameOperation, dict[str, Any]]:
operation = await load_operation(session, repo_id, operation_id, for_update=True)
if operation.phase == "rolled-back":
return operation, {"rollback_from_phase": "rolled-back", "safe_to_rollback": True, "blockers": [], "irreversible": []}
if confirmation != rollback_confirmation_for(operation.id):
raise RenamePreconditionFailed("Explicit rollback confirmation is incorrect")
if operation.phase != expected_phase:
raise RenamePreconditionFailed(
"Rollback preflight phase compare-and-set failed",
details={"expected": expected_phase, "actual": operation.phase},
)
if operation.phase == "rollback-preflight":
rollback_data = (operation.evidence or {}).get("rollback_preflight") or {}
return operation, rollback_data
competing_result = await session.execute(
select(RepositoryRenameOperation.id).where(
RepositoryRenameOperation.repo_id == repo_id,
RepositoryRenameOperation.id != operation.id,
RepositoryRenameOperation.phase.in_(ACTIVE_OPERATION_PHASES),
)
)
competing = [str(value) for value in competing_result.scalars().all()]
old = await _forge_at(gateway, operation, operation.old_slug)
new = await _forge_at(gateway, operation, operation.new_slug)
blockers: list[dict[str, Any]] = []
if competing:
blockers.append(
{
"code": "newer_active_rename_operation",
"operation_ids": competing,
}
)
rollback_from_phase = operation.phase
if operation.phase == "preflighted":
if old is not None and new is None:
try:
_assert_snapshot(operation, old)
except RenameLifecycleError as exc:
blockers.append({"code": exc.code, "message": str(exc), **exc.details})
elif old is None and new is not None:
# Forge committed but the phase journal did not. This is the
# principal interruption case the operation ID must recover from.
rollback_from_phase = "forge-renamed-unrecorded"
try:
_assert_snapshot(operation, new)
except RenameLifecycleError as exc:
blockers.append({"code": exc.code, "message": str(exc), **exc.details})
elif old is None:
blockers.append({"code": "forge_repository_absent_at_both_names"})
else:
blockers.append({"code": "both_forge_slugs_claimed"})
else:
if old is not None:
blockers.append({"code": "old_forge_slug_claimed", "forge_repository_id": old.repository_id})
if new is None:
blockers.append({"code": "renamed_forge_repository_absent"})
else:
try:
_assert_snapshot(operation, new)
except RenameLifecycleError as exc:
blockers.append({"code": exc.code, "message": str(exc), **exc.details})
irreversible = [
{"code": "external_caches", "message": "External caches and redirects may retain the new coordinate."},
{"code": "consumer_changes", "message": "Consumer commits made after cutover are not rewritten automatically."},
{"code": "telemetry_history", "message": "The attempted operation and phase evidence remain durable."},
]
data = {
"rollback_from_phase": rollback_from_phase,
"safe_to_rollback": not blockers,
"blockers": blockers,
"irreversible": irreversible,
}
if blockers:
raise RenamePreconditionFailed("Repository rename rollback is unsafe", details=data)
journal = deepcopy(operation.evidence or {})
journal["rollback_preflight"] = data
operation.evidence = journal
_record_phase(operation, "rollback-preflight", data)
await session.commit()
await session.refresh(operation)
return operation, data
async def _rollback_statehub(
session: AsyncSession, operation: RepositoryRenameOperation
) -> dict[str, Any]:
repo = await session.get(ManagedRepo, operation.repo_id, with_for_update=True)
if repo is None:
raise RenameNotFound("Managed repository disappeared")
if repo.slug == operation.old_slug:
return {"resumed": True, "repo_id": str(repo.id)}
if repo.slug != operation.new_slug:
raise RenamePreconditionFailed("Managed repository slug changed outside this operation")
old_result = await session.execute(
select(RepositorySlug).where(RepositorySlug.slug == operation.old_slug).with_for_update()
)
new_result = await session.execute(
select(RepositorySlug).where(RepositorySlug.slug == operation.new_slug).with_for_update()
)
old_slug = old_result.scalar_one_or_none()
new_slug = new_result.scalar_one_or_none()
if not old_slug or old_slug.repo_id != repo.id or old_slug.kind != "alias":
raise RenamePreconditionFailed("Protected old alias is unavailable for rollback")
if not new_slug or new_slug.repo_id != repo.id or new_slug.kind != "canonical":
raise RenamePreconditionFailed("New canonical slug is unavailable for rollback")
new_slug.kind = "alias"
new_slug.protected = True
new_slug.source_operation_id = operation.id
await session.flush()
old_slug.kind = "canonical"
old_slug.protected = True
old_slug.source_operation_id = operation.id
repo.slug = operation.old_slug
repo.name = operation.old_coordinates["name"]
repo.remote_url = operation.old_coordinates["remote_url"]
await session.flush()
return {"resumed": False, "repo_id": str(repo.id)}
async def apply_rollback(
session: AsyncSession,
gateway: ForgeRepositoryGateway,
repo_id: uuid.UUID,
operation_id: uuid.UUID,
*,
expected_phase: str,
confirmation: str,
) -> tuple[RepositoryRenameOperation, bool]:
operation = await load_operation(session, repo_id, operation_id, for_update=True)
if operation.phase == "rolled-back":
return operation, True
if operation.phase != expected_phase or operation.phase != "rollback-preflight":
raise RenamePreconditionFailed(
"Rollback phase compare-and-set failed",
details={"expected": expected_phase, "actual": operation.phase},
)
if confirmation != rollback_confirmation_for(operation.id):
raise RenamePreconditionFailed("Explicit rollback confirmation is incorrect")
rollback_from = ((operation.evidence or {}).get("rollback_preflight") or {}).get("rollback_from_phase")
if not rollback_from:
raise RenamePreconditionFailed("Rollback preflight evidence is missing")
old = await _forge_at(gateway, operation, operation.old_slug)
new = await _forge_at(gateway, operation, operation.new_slug)
forge_resumed = False
if rollback_from != "preflighted":
if old is not None:
_assert_snapshot(operation, old)
if new is not None:
raise RenamePreconditionFailed("Both Forge coordinates are claimed during rollback")
forge_resumed = True
else:
if new is None:
raise RenamePreconditionFailed("Forge repository is absent during rollback")
_assert_snapshot(operation, new)
try:
restored = await gateway.rename(
instance=operation.expected_forge_instance,
owner=operation.expected_forge_owner,
old_name=operation.new_slug,
new_name=operation.old_slug,
)
except (ForgeRepositoryUnreadable, ForgeRepositoryConflict) as exc:
raise RenamePreconditionFailed(str(exc)) from exc
_assert_snapshot(operation, restored)
statehub = await _rollback_statehub(session, operation)
_record_phase(
operation,
"rolled-back",
{"rollback_from_phase": rollback_from, "forge_resumed": forge_resumed, "statehub": statehub},
)
await session.commit()
await session.refresh(operation)
return operation, False