feat: consume Repo Manager classification publisher
Some checks failed
CI Smoke / host-smoke (push) Successful in 0s
CI Smoke / pytest-smoke (push) Failing after 1s

Assistant: codex
Assistant-Model: gpt-5.6-sol
Assistant-Session: 01a053ff-1d6f-7fe2-ac1c-a6eb40a42a0c
This commit is contained in:
tegwick 2026-09-01 00:49:30 +02:00
parent b2f430ffb1
commit f582773f5a
6 changed files with 201 additions and 3 deletions

View file

@ -1,6 +1,7 @@
from __future__ import annotations
from contextlib import asynccontextmanager
import asyncio
from contextlib import asynccontextmanager, suppress
from fastapi import FastAPI, Response, status
@ -17,6 +18,7 @@ from hub_core.runtime.repository_navigation import (
from hub_core.runtime.repository_navigation_routes import (
create_repository_navigation_router,
)
from hub_core.runtime.repo_manager_client import HTTPRepoProjectionClient
from hub_core.runtime.store import InMemoryPortStore, PortStore
from hub_core.runtime.validation import ContractValidator
from hub_core.runtime.workload_projection import (
@ -37,8 +39,20 @@ def create_app(
resolved_settings = settings or RuntimeSettings.from_env()
resolved_store = port_store or _create_store(resolved_settings)
owns_store = port_store is None
resolved_repo_projection_client = repo_projection_client
owns_repo_projection_client = False
if (
resolved_repo_projection_client is None
and resolved_settings.repo_manager_base_url is not None
):
resolved_repo_projection_client = HTTPRepoProjectionClient(
resolved_settings.repo_manager_base_url,
api_token=resolved_settings.repo_manager_api_token,
timeout_seconds=resolved_settings.repo_manager_timeout_seconds,
)
owns_repo_projection_client = True
repository_navigation = RepositoryNavigationService(
client=repo_projection_client,
client=resolved_repo_projection_client,
store=resolved_store,
)
workload_projection = WorkloadProjectionService(
@ -48,19 +62,33 @@ def create_app(
@asynccontextmanager
async def lifespan(_: FastAPI):
if repo_projection_client is not None:
refresh_task: asyncio.Task[None] | None = None
if resolved_repo_projection_client is not None:
try:
await repository_navigation.refresh()
except ProjectionRejected:
# The readiness dependency reports the rejected or absent
# projection while the API remains available for diagnosis.
pass
if resolved_settings.repo_projection_refresh_seconds > 0:
refresh_task = asyncio.create_task(
_refresh_repository_projection(
repository_navigation,
resolved_settings.repo_projection_refresh_seconds,
)
)
if workload_projection_client is not None:
try:
await workload_projection.refresh()
except WorkloadProjectionRejected:
pass
yield
if refresh_task is not None:
refresh_task.cancel()
with suppress(asyncio.CancelledError):
await refresh_task
if owns_repo_projection_client:
await resolved_repo_projection_client.aclose() # type: ignore[union-attr]
if owns_store and (closer := getattr(resolved_store, "aclose", None)):
await closer()
@ -117,6 +145,19 @@ def create_app(
return app
async def _refresh_repository_projection(
service: RepositoryNavigationService,
interval_seconds: float,
) -> None:
while True:
await asyncio.sleep(interval_seconds)
try:
await service.refresh()
except ProjectionRejected:
# The service records stale state and diagnostics; retry on schedule.
pass
def _create_store(settings: RuntimeSettings) -> PortStore:
if settings.backend == "memory":
return InMemoryPortStore()

View file

@ -28,12 +28,20 @@ class RuntimeSettings:
mcp_transport: str = "http"
database_url: str | None = None
api_token: str | None = None
repo_manager_base_url: str | None = None
repo_manager_api_token: str | None = None
repo_manager_timeout_seconds: float = 10.0
repo_projection_refresh_seconds: float = 300.0
v2_groups: frozenset[str] = frozenset()
v2_write_groups: frozenset[str] = frozenset()
legacy_write_groups: frozenset[str] = frozenset()
legacy_health: bool = False
def __post_init__(self) -> None:
if self.repo_manager_timeout_seconds <= 0:
raise ValueError("Repo Manager timeout must be positive")
if self.repo_projection_refresh_seconds < 0:
raise ValueError("repository projection refresh interval must not be negative")
overlap = self.v2_write_groups & self.legacy_write_groups
if overlap:
joined = ", ".join(sorted(overlap))
@ -61,6 +69,14 @@ class RuntimeSettings:
mcp_transport=os.getenv("HUB_CORE_MCP_TRANSPORT", "http"),
database_url=os.getenv("HUB_CORE_DATABASE_URL") or os.getenv("DATABASE_URL"),
api_token=os.getenv("HUB_CORE_API_TOKEN") or os.getenv("CORE_HUB_API_TOKEN"),
repo_manager_base_url=os.getenv("HUB_CORE_REPO_MANAGER_BASE_URL"),
repo_manager_api_token=os.getenv("HUB_CORE_REPO_MANAGER_API_TOKEN"),
repo_manager_timeout_seconds=float(
os.getenv("HUB_CORE_REPO_MANAGER_TIMEOUT_SECONDS", "10")
),
repo_projection_refresh_seconds=float(
os.getenv("HUB_CORE_REPO_PROJECTION_REFRESH_SECONDS", "300")
),
v2_groups=_env_set("HUB_CORE_V2_GROUPS"),
v2_write_groups=_env_set("HUB_CORE_V2_WRITE_GROUPS"),
legacy_write_groups=_env_set("CORE_HUB_V2_WRITE_GROUPS"),

View file

@ -0,0 +1,50 @@
from __future__ import annotations
from collections.abc import Mapping
from typing import Any
from urllib.parse import urlparse
import httpx
class HTTPRepoProjectionClient:
"""HTTP adapter for Repo Manager's read-only ``port.repo`` publisher."""
def __init__(
self,
base_url: str,
*,
api_token: str | None = None,
timeout_seconds: float = 10.0,
transport: httpx.AsyncBaseTransport | None = None,
) -> None:
parsed = urlparse(base_url)
if parsed.scheme not in {"http", "https"} or not parsed.netloc:
raise ValueError("Repo Manager base URL must be an absolute HTTP(S) origin")
if timeout_seconds <= 0:
raise ValueError("Repo Manager timeout must be positive")
headers = {"Authorization": f"Bearer {api_token}"} if api_token else None
self._client = httpx.AsyncClient(
base_url=base_url.rstrip("/"),
headers=headers,
timeout=httpx.Timeout(timeout_seconds),
follow_redirects=True,
transport=transport,
)
async def fetch_classification_page(
self, cursor: str | None
) -> Mapping[str, Any]:
params = {"cursor": cursor} if cursor is not None else None
response = await self._client.get(
"/ports/repositories/classifications",
params=params,
)
response.raise_for_status()
payload = response.json()
if not isinstance(payload, dict):
raise ValueError("Repo Manager classification page must be a JSON object")
return payload
async def aclose(self) -> None:
await self._client.aclose()