from __future__ import annotations import asyncio from contextlib import asynccontextmanager, suppress import httpx from fastapi import FastAPI, Response, status from hub_core import __version__ from hub_core.runtime.config import RuntimeSettings from hub_core.runtime.compat import SQLCompatibilityStore, create_compatibility_router from hub_core.runtime.models import HealthResponse, ReadinessResponse from hub_core.runtime.ports import create_ports_router from hub_core.runtime.repository_navigation import ( ProjectionRejected, RepoProjectionClient, RepositoryNavigationService, ) 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 ( WorkloadProjectionClient, WorkloadProjectionRejected, WorkloadProjectionService, ) from hub_core.runtime.workload_projection_routes import create_workload_projection_router from hub_core.security.boundary import AccessBoundary, AccessController, FactSource, ACCESS_DEPENDENCIES from hub_core.security.config import SecuritySettings from hub_core.security.browser import BrowserSettings, BrowserSessions, create_browser_router def create_app( *, settings: RuntimeSettings | None = None, port_store: PortStore | None = None, repo_projection_client: RepoProjectionClient | None = None, workload_projection_client: WorkloadProjectionClient | None = None, access_controller: AccessController | None = None, access_facts: FactSource | None = None, security_settings: SecuritySettings | None = None, browser_settings: BrowserSettings | None = None, ) -> FastAPI: resolved_settings = settings or RuntimeSettings.from_env() if access_controller is not None and not resolved_settings.enforce_access: raise ValueError("an access controller requires enforcement mode") # Importing this module also constructs the standalone app. Environment # configuration is activated only by an explicit owner-facts composition; # without that adapter the default app stays closed, not import-broken. if security_settings is None and access_facts is not None: security_settings = SecuritySettings.from_env() security_client = None if access_controller is not None and (access_facts is not None or security_settings is not None): raise ValueError("choose an access controller or owner-facts composition") if security_settings is not None or access_facts is not None: if not resolved_settings.enforce_access: raise ValueError("security composition requires enforcement mode") if security_settings is None or access_facts is None: raise ValueError("security composition requires configuration and authoritative owner facts") security_client = httpx.AsyncClient(trust_env=False) access_controller = security_settings.compose(facts=access_facts, client=security_client) if browser_settings is not None and (not resolved_settings.enforce_access or access_controller is None): raise ValueError("browser sessions require an enforced access controller") browser = BrowserSessions(settings=browser_settings, controller=access_controller) if browser_settings else None resolved_store = port_store or _create_store(resolved_settings) owns_store = port_store is None outcome_delivery = (access_controller is not None and callable(getattr(access_controller.audit, "append_outcome", None)) and callable(getattr(resolved_store, "deliver_outcomes", 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=resolved_repo_projection_client, store=resolved_store, ) workload_projection = WorkloadProjectionService( client=workload_projection_client, store=resolved_store, ) @asynccontextmanager async def lifespan(_: FastAPI): 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 outcome_task = asyncio.create_task(_deliver_outcomes(resolved_store, access_controller.audit)) if outcome_delivery else None try: yield finally: if outcome_task is not None: outcome_task.cancel() with suppress(asyncio.CancelledError): await outcome_task if browser is not None: browser.clear() if refresh_task is not None: refresh_task.cancel() with suppress(asyncio.CancelledError): await refresh_task if security_client is not None: await security_client.aclose() 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() app = FastAPI( title="Hub Core Runtime", version=__version__, description="HelixForge hub framework and named-port runtime.", lifespan=lifespan, ) app.state.settings = resolved_settings app.state.port_store = resolved_store app.state.compat_store = ( SQLCompatibilityStore(resolved_store) if resolved_store.backend_name == "postgresql" else None ) app.state.contract_validator = ContractValidator() app.state.repository_navigation = repository_navigation app.state.workload_projection = workload_projection app.state.access_controller = access_controller app.state.browser_sessions = browser if browser is not None: app.include_router(create_browser_router()) if resolved_settings.enforce_access: app.add_middleware(AccessBoundary, host=app, controller=access_controller, browser=browser) @app.get("/healthz", response_model=HealthResponse, tags=["system"]) async def healthz() -> HealthResponse: return HealthResponse( service="core-hub" if resolved_settings.legacy_health else "hub-core", version=__version__, ) @app.get("/readyz", response_model=ReadinessResponse, tags=["system"]) async def readyz(response: Response) -> ReadinessResponse: dependency_checks = { **await resolved_store.readiness_checks(), **await repository_navigation.readiness_checks(), **await workload_projection.readiness_checks(), } if resolved_settings.enforce_access: dependency_checks.update(access_controller.readiness_checks() if access_controller else { **{"access_" + name: "unavailable" for name in ACCESS_DEPENDENCIES}, "access_profile": "unavailable", }) if callable(getattr(resolved_store, "deliver_outcomes", None)): dependency_checks["outcome_delivery"] = (await resolved_store.outcome_readiness()) if outcome_delivery else "unavailable" ready = resolved_settings.is_ready(resolved_store.backend_name) and all( value in {"ok", "not_applicable"} for value in dependency_checks.values() ) if not ready: response.status_code = status.HTTP_503_SERVICE_UNAVAILABLE return ReadinessResponse( status="ok" if ready else "degraded", checks={ **resolved_settings.readiness_checks(resolved_store.backend_name), **dependency_checks, }, ) # Exact projection routes precede the generic /projections/{projection_id} # route so Starlette dispatch cannot shadow them. if resolved_settings.statehub_inbox_reads: from hub_core.runtime.inbox_projection import create_inbox_projection_router app.include_router(create_inbox_projection_router()) app.include_router(create_workload_projection_router()) app.include_router(create_ports_router()) app.include_router(create_repository_navigation_router()) app.include_router(create_compatibility_router()) return app async def _deliver_outcomes(store, sink) -> None: while True: try: await store.deliver_outcomes(sink) except Exception: # DB outages leave rows durable; readiness exposes missing/stale data. pass await asyncio.sleep(1) 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() if settings.backend == "postgresql": if not settings.database_url: raise RuntimeError("HUB_CORE_DATABASE_URL is required for PostgreSQL backend") from hub_core.runtime.postgres_store import PostgresPortStore return PostgresPortStore.from_url(settings.database_url) raise RuntimeError(f"Unsupported HUB_CORE_BACKEND '{settings.backend}'") app = create_app()