hub-core/hub_core/runtime/ports.py
tegwick 3e386147fd
Some checks failed
CI Smoke / host-smoke (push) Successful in 0s
CI Smoke / pytest-smoke (push) Failing after 3s
feat: add fail-closed Hub access profile foundation
Assistant: codex
Assistant-Model: gpt-6-astra
Assistant-Session: 01a0e747-8f27-7242-8df8-8bc44f88c929
2026-09-28 11:44:50 +02:00

174 lines
6.1 KiB
Python

from __future__ import annotations
from uuid import UUID
from fastapi import APIRouter, Depends, Header, HTTPException, Request, status
from jsonschema import ValidationError
from hub_core.runtime.models import (
EventCommand,
MessageCommand,
PortAccepted,
PortCollection,
PortRecord,
RegistryRegistration,
)
from hub_core.runtime.store import PortStore
from hub_core.runtime.validation import ContractValidator
def get_port_store(request: Request) -> PortStore:
return request.app.state.port_store
def get_contract_validator(request: Request) -> ContractValidator:
return request.app.state.contract_validator
def _attribute_event(body: EventCommand, request: Request) -> EventCommand:
context = getattr(request.state, "hub_access", None)
if context is None:
return body
# Reserved server provenance overrides any payload assertion. Domain
# subject_refs remain business data and are never authentication evidence.
return body.model_copy(update={"payload": {**body.payload, "_hub_access": {
"issuer": context.actor.issuer, "subject": context.actor.subject,
"principal_type": context.actor.principal_type,
"actor_tenant": context.actor.tenant,
"target_tenant": context.facts.target_tenant,
"correlation_id": context.correlation_id,
}}})
def create_ports_router() -> APIRouter:
router = APIRouter(prefix="/ports")
@router.post(
"/registry/registrations",
response_model=PortAccepted,
status_code=status.HTTP_202_ACCEPTED,
tags=["registry"],
openapi_extra={"x-port-id": "port.registry", "x-direction": "in"},
)
async def register_extension(
body: RegistryRegistration,
x_correlation_id: UUID = Header(alias="X-Correlation-ID"),
store: PortStore = Depends(get_port_store),
validator: ContractValidator = Depends(get_contract_validator),
) -> PortAccepted:
try:
validator.validate_registration(body)
except (ValidationError, ValueError) as exc:
raise HTTPException(status_code=422, detail=str(exc)) from exc
return await store.register_extension(body, x_correlation_id)
@router.get(
"/registry/registrations/{hub_slug}",
response_model=PortRecord,
tags=["registry"],
openapi_extra={"x-port-id": "port.registry", "x-direction": "out"},
)
async def resolve_registration(
hub_slug: str,
store: PortStore = Depends(get_port_store),
) -> PortRecord:
resolved = await store.resolve_registration(hub_slug)
if resolved is None:
raise HTTPException(status_code=404, detail=f"Registration '{hub_slug}' not found")
return resolved
@router.get(
"/registry/registrations/{hub_slug}/audit",
response_model=PortCollection,
tags=["registry"],
openapi_extra={"x-port-id": "port.registry", "x-direction": "out"},
)
async def registration_audit(
hub_slug: str,
store: PortStore = Depends(get_port_store),
) -> PortCollection:
audit = await store.list_registration_audit(hub_slug)
if not audit.items:
raise HTTPException(status_code=404, detail=f"No audit history for '{hub_slug}'")
return audit
@router.get(
"/messaging/messages",
response_model=PortCollection,
tags=["messaging"],
openapi_extra={"x-port-id": "port.messaging", "x-direction": "out"},
)
async def list_messages(
address: str,
conversation_id: UUID | None = None,
store: PortStore = Depends(get_port_store),
) -> PortCollection:
return await store.list_messages(address, conversation_id)
@router.post(
"/messaging/messages",
response_model=PortAccepted,
status_code=status.HTTP_202_ACCEPTED,
tags=["messaging"],
openapi_extra={"x-port-id": "port.messaging", "x-direction": "in"},
)
async def send_message(
body: MessageCommand,
store: PortStore = Depends(get_port_store),
) -> PortAccepted:
return await store.send_message(body)
@router.post(
"/events/progress",
response_model=PortAccepted,
status_code=status.HTTP_202_ACCEPTED,
tags=["events"],
openapi_extra={"x-port-id": "port.events.progress", "x-direction": "in"},
)
async def append_progress(
body: EventCommand,
request: Request,
store: PortStore = Depends(get_port_store),
validator: ContractValidator = Depends(get_contract_validator),
) -> PortAccepted:
try:
validator.validate_event_family(body.event_type, "progress")
except ValueError as exc:
raise HTTPException(status_code=422, detail=str(exc)) from exc
return await store.append_progress(_attribute_event(body, request))
@router.post(
"/events/interaction",
response_model=PortAccepted,
status_code=status.HTTP_202_ACCEPTED,
tags=["events"],
openapi_extra={"x-port-id": "port.events.interaction", "x-direction": "in"},
)
async def append_interaction(
body: EventCommand,
request: Request,
store: PortStore = Depends(get_port_store),
validator: ContractValidator = Depends(get_contract_validator),
) -> PortAccepted:
try:
validator.validate_event_family(body.event_type, "interaction")
except ValueError as exc:
raise HTTPException(status_code=422, detail=str(exc)) from exc
return await store.append_interaction(_attribute_event(body, request))
@router.get(
"/projections/{projection_id}",
response_model=PortRecord,
tags=["projections"],
openapi_extra={"x-port-id": "port.projection.query", "x-direction": "out"},
)
async def query_projection(
projection_id: str,
store: PortStore = Depends(get_port_store),
) -> PortRecord:
projection = await store.query_projection(projection_id)
if projection is None:
raise HTTPException(status_code=404, detail=f"Projection '{projection_id}' not found")
return projection
return router