Implements the four residual conformance checks left open by the T04
minimal vertical:
- C2: GET /ports/registry/registrations/{hub_slug} resolves missing (404),
ambiguous (shared reuse_surface_id across hub_slugs), and stale
(deprecated/retired descriptor) registrations; a new .../audit route
exposes queryable registration history from the existing in-memory
history and the PostgreSQL runtime_audit_ledger.
- C7: harness proof that disabled compatibility groups deny access
(404) with no fixture credentials involved, matching the existing
fail-closed compat router behavior.
- C9: harness proof plus a dedicated test that /readyz degrades only on
an unavailable configured dependency while unrelated disabled
projections stay non-blocking.
- C10: ContractValidator now negotiates contract_version_min/max against
the runtime's contract version and rejects incompatible or inverted
ranges with an explicit 422 instead of silently accepting them.
HUB-WP-0009 is now finished. HUB-WP-0006 is marked blocked: its only open
task (T06) has no remaining hub-core code path and waits on an external
Forgejo identity/production deployment gate. HUB-WP-0011 is marked
blocked: T02/T03 already waited on external credential/deployment
review, and T01 needs a source/destination ownership and retention
decision against live message data before it can be implemented safely.
Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
Assistant: claude-code
Assistant-Model: sonnet
Assistant-Process: 310936@bnt-lap001
Assistant-Session: 00cd9abe-09a0-416b-88e0-f907b9101629
330 lines
12 KiB
Python
330 lines
12 KiB
Python
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import hashlib
|
|
import json
|
|
from collections.abc import Mapping
|
|
from copy import deepcopy
|
|
from datetime import datetime, timezone
|
|
from typing import Any, Protocol
|
|
from uuid import UUID, uuid4
|
|
|
|
from hub_core.contracts import CONTRACT_VERSION
|
|
from hub_core.runtime.models import (
|
|
EventCommand,
|
|
MessageCommand,
|
|
PortAccepted,
|
|
PortCollection,
|
|
PortRecord,
|
|
Provenance,
|
|
RegistryRegistration,
|
|
)
|
|
from hub_core.runtime.repository_navigation import NavigationProjection
|
|
from hub_core.runtime.workload_projection import WorkloadProjection
|
|
|
|
|
|
class PortStore(Protocol):
|
|
"""Persistence boundary for the initial hub-core runtime ports."""
|
|
|
|
backend_name: str
|
|
|
|
async def readiness_checks(self) -> dict[str, str]: ...
|
|
|
|
async def register_extension(
|
|
self,
|
|
registration: RegistryRegistration,
|
|
correlation_id: UUID,
|
|
) -> PortAccepted: ...
|
|
|
|
async def send_message(self, command: MessageCommand) -> PortAccepted: ...
|
|
|
|
async def list_messages(self, address: str, conversation_id: UUID | None) -> PortCollection: ...
|
|
|
|
async def append_progress(self, command: EventCommand) -> PortAccepted: ...
|
|
|
|
async def append_interaction(self, command: EventCommand) -> PortAccepted: ...
|
|
|
|
async def query_projection(self, projection_id: str) -> PortRecord | None: ...
|
|
|
|
async def get_repository_navigation(self) -> NavigationProjection | None: ...
|
|
|
|
async def replace_repository_navigation(
|
|
self, projection: NavigationProjection
|
|
) -> None: ...
|
|
|
|
async def mark_repository_navigation_stale(
|
|
self, *, checked_at: datetime, diagnostic: Mapping[str, Any]
|
|
) -> None: ...
|
|
|
|
async def get_workload_projection(self) -> WorkloadProjection | None: ...
|
|
|
|
async def replace_workload_projection(self, projection: WorkloadProjection) -> None: ...
|
|
|
|
async def mark_workload_projection_stale(
|
|
self, *, checked_at: datetime, diagnostic: Mapping[str, Any]
|
|
) -> None: ...
|
|
|
|
async def resolve_registration(self, hub_slug: str) -> PortRecord | None: ...
|
|
|
|
async def list_registration_audit(self, hub_slug: str) -> PortCollection: ...
|
|
|
|
|
|
class InMemoryPortStore:
|
|
"""Deterministic ephemeral backend for local runtime and conformance tests.
|
|
|
|
Production readiness rejects this backend unless explicitly allowed. The
|
|
store deliberately keeps progress and interaction event families separate.
|
|
"""
|
|
|
|
backend_name = "memory"
|
|
|
|
def __init__(self) -> None:
|
|
self._lock = asyncio.Lock()
|
|
self._registrations: dict[str, dict[str, Any]] = {}
|
|
self._registration_audit: dict[str, list[dict[str, Any]]] = {}
|
|
self._messages: list[dict[str, Any]] = []
|
|
self._progress_events: list[dict[str, Any]] = []
|
|
self._interaction_events: list[dict[str, Any]] = []
|
|
self._repository_navigation: NavigationProjection | None = None
|
|
self._workload_projection: WorkloadProjection | None = None
|
|
|
|
async def readiness_checks(self) -> dict[str, str]:
|
|
return {"database": "not_applicable"}
|
|
|
|
async def register_extension(
|
|
self,
|
|
registration: RegistryRegistration,
|
|
correlation_id: UUID,
|
|
) -> PortAccepted:
|
|
hub_slug = str(registration.descriptor["hub_slug"])
|
|
value = registration.model_dump(mode="json")
|
|
async with self._lock:
|
|
duplicate = self._registrations.get(hub_slug) == value
|
|
self._registrations[hub_slug] = deepcopy(value)
|
|
self._registration_audit.setdefault(hub_slug, []).append(
|
|
{
|
|
"id": str(uuid4()),
|
|
"action": "registry.duplicate" if duplicate else "registry.accepted",
|
|
"hub_slug": hub_slug,
|
|
"correlation_id": str(correlation_id),
|
|
"recorded_at": _now().isoformat(),
|
|
}
|
|
)
|
|
return PortAccepted(
|
|
id=hub_slug,
|
|
status="duplicate" if duplicate else "accepted",
|
|
correlation_id=correlation_id,
|
|
)
|
|
|
|
async def resolve_registration(self, hub_slug: str) -> PortRecord | None:
|
|
async with self._lock:
|
|
value = self._registrations.get(hub_slug)
|
|
if value is None:
|
|
return None
|
|
registrations = deepcopy(self._registrations)
|
|
return _resolve_registration_record("hub-core-memory", hub_slug, value, registrations)
|
|
|
|
async def list_registration_audit(self, hub_slug: str) -> PortCollection:
|
|
async with self._lock:
|
|
entries = deepcopy(self._registration_audit.get(hub_slug, []))
|
|
return PortCollection(
|
|
items=[self._record("registration_audit", entry) for entry in entries]
|
|
)
|
|
|
|
async def send_message(self, command: MessageCommand) -> PortAccepted:
|
|
message_id = uuid4()
|
|
value = {
|
|
"id": str(message_id),
|
|
"created_at": _now().isoformat(),
|
|
**command.model_dump(mode="json"),
|
|
}
|
|
async with self._lock:
|
|
self._messages.append(value)
|
|
return PortAccepted(
|
|
id=str(message_id),
|
|
status="accepted",
|
|
correlation_id=command.correlation_id,
|
|
)
|
|
|
|
async def list_messages(self, address: str, conversation_id: UUID | None) -> PortCollection:
|
|
async with self._lock:
|
|
values = [
|
|
deepcopy(message)
|
|
for message in self._messages
|
|
if address in message["to_addresses"]
|
|
and (
|
|
conversation_id is None
|
|
or message.get("conversation_id") == str(conversation_id)
|
|
)
|
|
]
|
|
return PortCollection(items=[self._record("message", value) for value in values])
|
|
|
|
async def append_progress(self, command: EventCommand) -> PortAccepted:
|
|
return await self._append_event(command, self._progress_events, "progress")
|
|
|
|
async def append_interaction(self, command: EventCommand) -> PortAccepted:
|
|
return await self._append_event(command, self._interaction_events, "interaction")
|
|
|
|
async def query_projection(self, projection_id: str) -> PortRecord | None:
|
|
async with self._lock:
|
|
sources: Mapping[str, Any] = {
|
|
"hub_registry": list(self._registrations.values()),
|
|
"messages": self._messages,
|
|
"progress_events": self._progress_events,
|
|
"interaction_events": self._interaction_events,
|
|
}
|
|
if projection_id not in sources:
|
|
return None
|
|
items = deepcopy(sources[projection_id])
|
|
return self._record(
|
|
projection_id,
|
|
{
|
|
"projection_id": projection_id,
|
|
"items": items,
|
|
"rebuild_from": _rebuild_sources(projection_id),
|
|
},
|
|
)
|
|
|
|
async def get_repository_navigation(self) -> NavigationProjection | None:
|
|
async with self._lock:
|
|
return deepcopy(self._repository_navigation)
|
|
|
|
async def replace_repository_navigation(
|
|
self, projection: NavigationProjection
|
|
) -> None:
|
|
async with self._lock:
|
|
self._repository_navigation = deepcopy(projection)
|
|
|
|
async def mark_repository_navigation_stale(
|
|
self, *, checked_at: datetime, diagnostic: Mapping[str, Any]
|
|
) -> None:
|
|
async with self._lock:
|
|
current = self._repository_navigation
|
|
if current is None:
|
|
return
|
|
self._repository_navigation = NavigationProjection(
|
|
projection_status="stale",
|
|
source_snapshot=deepcopy(current.source_snapshot),
|
|
source_checked_at=checked_at,
|
|
rebuilt_at=current.rebuilt_at,
|
|
content_hash=current.content_hash,
|
|
repositories=deepcopy(current.repositories),
|
|
facets=deepcopy(current.facets),
|
|
diagnostics=(deepcopy(dict(diagnostic)),),
|
|
)
|
|
|
|
async def get_workload_projection(self) -> WorkloadProjection | None:
|
|
async with self._lock:
|
|
return deepcopy(self._workload_projection)
|
|
|
|
async def replace_workload_projection(self, projection: WorkloadProjection) -> None:
|
|
async with self._lock:
|
|
self._workload_projection = deepcopy(projection)
|
|
|
|
async def mark_workload_projection_stale(
|
|
self, *, checked_at: datetime, diagnostic: Mapping[str, Any]
|
|
) -> None:
|
|
async with self._lock:
|
|
current = self._workload_projection
|
|
if current is None:
|
|
return
|
|
self._workload_projection = WorkloadProjection(
|
|
projection_status="stale",
|
|
source_snapshot=deepcopy(current.source_snapshot),
|
|
source_checked_at=checked_at,
|
|
rebuilt_at=current.rebuilt_at,
|
|
content_hash=current.content_hash,
|
|
workloads=deepcopy(current.workloads),
|
|
diagnostics=(deepcopy(dict(diagnostic)),),
|
|
)
|
|
|
|
async def _append_event(
|
|
self,
|
|
command: EventCommand,
|
|
target: list[dict[str, Any]],
|
|
family: str,
|
|
) -> PortAccepted:
|
|
event_id = uuid4()
|
|
value = {
|
|
"id": str(event_id),
|
|
"family": family,
|
|
"recorded_at": _now().isoformat(),
|
|
**command.model_dump(mode="json"),
|
|
}
|
|
async with self._lock:
|
|
target.append(value)
|
|
return PortAccepted(
|
|
id=str(event_id),
|
|
status="accepted",
|
|
correlation_id=command.correlation_id,
|
|
)
|
|
|
|
def _record(self, kind: str, value: dict[str, Any]) -> PortRecord:
|
|
encoded = json.dumps(value, sort_keys=True, separators=(",", ":")).encode()
|
|
record_id = str(value.get("id") or kind)
|
|
return PortRecord(
|
|
id=record_id,
|
|
data=deepcopy(value),
|
|
provenance=Provenance(
|
|
source_system="hub-core-memory",
|
|
source_ref=f"memory://{kind}/{record_id}",
|
|
schema_version=CONTRACT_VERSION,
|
|
content_hash=hashlib.sha256(encoded).hexdigest(),
|
|
indexed_at=_now(),
|
|
),
|
|
)
|
|
|
|
|
|
def _now() -> datetime:
|
|
return datetime.now(timezone.utc)
|
|
|
|
|
|
def _resolve_registration_record(
|
|
source_system: str,
|
|
hub_slug: str,
|
|
value: Mapping[str, Any],
|
|
registrations: Mapping[str, Mapping[str, Any]],
|
|
) -> PortRecord:
|
|
descriptor = value["descriptor"]
|
|
reuse_surface_id = descriptor.get("reuse_surface_id")
|
|
ambiguous_with = sorted(
|
|
other_slug
|
|
for other_slug, other_value in registrations.items()
|
|
if other_slug != hub_slug
|
|
and other_value["descriptor"].get("reuse_surface_id") == reuse_surface_id
|
|
)
|
|
stale = descriptor.get("status") in {"deprecated", "retired"}
|
|
if ambiguous_with:
|
|
resolution = "ambiguous"
|
|
elif stale:
|
|
resolution = "stale"
|
|
else:
|
|
resolution = "ok"
|
|
data = {
|
|
"hub_slug": hub_slug,
|
|
"resolution": resolution,
|
|
"ambiguous_with": ambiguous_with,
|
|
"descriptor": deepcopy(dict(descriptor)),
|
|
"manifest": deepcopy(dict(value["manifest"])),
|
|
}
|
|
encoded = json.dumps(data, sort_keys=True, separators=(",", ":")).encode()
|
|
return PortRecord(
|
|
id=hub_slug,
|
|
data=data,
|
|
provenance=Provenance(
|
|
source_system=source_system,
|
|
source_ref=f"{source_system.replace('hub-core-', '')}://registration/{hub_slug}",
|
|
schema_version=CONTRACT_VERSION,
|
|
content_hash=hashlib.sha256(encoded).hexdigest(),
|
|
indexed_at=_now(),
|
|
),
|
|
)
|
|
|
|
|
|
def _rebuild_sources(projection_id: str) -> list[str]:
|
|
return {
|
|
"hub_registry": ["hub_descriptors", "hub_manifests"],
|
|
"messages": ["messages"],
|
|
"progress_events": ["progress_events"],
|
|
"interaction_events": ["interaction_events"],
|
|
}[projection_id]
|