hub-core/hub_core/runtime/store.py
tegwick e89d621f18
Some checks are pending
CI Smoke / host-smoke (push) Waiting to run
CI Smoke / pytest-smoke (push) Waiting to run
Close HUB-WP-0009 conformance gaps (C2, C7, C9, C10); mark blocked workplans
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
2026-09-27 23:59:46 +02:00

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]