feat: add fail-closed Hub access profile foundation
Some checks failed
CI Smoke / host-smoke (push) Successful in 0s
CI Smoke / pytest-smoke (push) Failing after 3s

Assistant: codex
Assistant-Model: gpt-6-astra
Assistant-Session: 01a0e747-8f27-7242-8df8-8bc44f88c929
This commit is contained in:
tegwick 2026-09-28 11:44:50 +02:00
parent df39fd5f43
commit 3e386147fd
35 changed files with 2009 additions and 195 deletions

View file

@ -2,6 +2,7 @@ from __future__ import annotations
import json
from typing import Any
from collections.abc import Callable
import httpx
from fastmcp import FastMCP
@ -57,8 +58,14 @@ class HubCoreMCPServer:
api_base: str,
instructions: str | None = None,
register_tools: bool = True,
token_provider: Callable[[], str] | None = None,
require_credentials: bool = False,
trailing_slash: bool = True,
) -> None:
self.api_base = api_base.rstrip("/")
self.token_provider = token_provider
self.require_credentials = require_credentials
self.trailing_slash = trailing_slash
self.mcp = FastMCP(
name=name,
instructions=instructions or "Generic FOS hub MCP server.",
@ -503,40 +510,51 @@ class HubCoreMCPServer:
try:
with self._client() as client:
response = client.get(
normalize_trailing_slash(path),
normalize_trailing_slash(path, trailing=self.trailing_slash),
params=self._clean(params or {}),
)
response.raise_for_status()
return response.json()
except httpx.HTTPStatusError as exc:
return {"error": f"API {exc.response.status_code}: {exc.response.text[:300]}"}
except Exception as exc:
return {"error": f"Request failed: {exc}"}
return {"error": f"API {exc.response.status_code}"}
except Exception:
return {"error": "Request failed"}
def _post(self, path: str, body: dict[str, Any]) -> Any:
try:
with self._client() as client:
response = client.post(normalize_trailing_slash(path), json=self._clean(body))
response = client.post(normalize_trailing_slash(path, trailing=self.trailing_slash), json=self._clean(body))
response.raise_for_status()
return response.json()
except httpx.HTTPStatusError as exc:
return {"error": f"API {exc.response.status_code}: {exc.response.text[:300]}"}
except Exception as exc:
return {"error": f"Request failed: {exc}"}
return {"error": f"API {exc.response.status_code}"}
except Exception:
return {"error": "Request failed"}
def _patch(self, path: str, body: dict[str, Any]) -> Any:
try:
with self._client() as client:
response = client.patch(normalize_trailing_slash(path), json=self._clean(body))
response = client.patch(normalize_trailing_slash(path, trailing=self.trailing_slash), json=self._clean(body))
response.raise_for_status()
return response.json()
except httpx.HTTPStatusError as exc:
return {"error": f"API {exc.response.status_code}: {exc.response.text[:300]}"}
except Exception as exc:
return {"error": f"Request failed: {exc}"}
return {"error": f"API {exc.response.status_code}"}
except Exception:
return {"error": "Request failed"}
def _client(self) -> httpx.Client:
return httpx.Client(base_url=self.api_base, timeout=30.0, follow_redirects=True)
# The host resolves a Hub-audience credential from the current invocation.
# Never retain it on the MCP server or fall back to a shared root token.
headers = {}
if self.token_provider is not None:
token = self.token_provider()
if not token or any(c.isspace() for c in token):
raise ValueError("current invocation has no Hub credential")
headers["Authorization"] = f"Bearer {token}"
elif self.require_credentials:
raise ValueError("MCP host must provide a current Hub credential")
return httpx.Client(base_url=self.api_base, timeout=30.0,
headers=headers, follow_redirects=not bool(headers))
@staticmethod
def _clean(data: dict[str, Any]) -> dict[str, Any]:

View file

@ -27,6 +27,7 @@ from hub_core.runtime.workload_projection import (
WorkloadProjectionService,
)
from hub_core.runtime.workload_projection_routes import create_workload_projection_router
from hub_core.security.boundary import AccessBoundary, AccessController
def create_app(
@ -35,6 +36,7 @@ def create_app(
port_store: PortStore | None = None,
repo_projection_client: RepoProjectionClient | None = None,
workload_projection_client: WorkloadProjectionClient | None = None,
access_controller: AccessController | None = None,
) -> FastAPI:
resolved_settings = settings or RuntimeSettings.from_env()
resolved_store = port_store or _create_store(resolved_settings)
@ -108,6 +110,9 @@ def create_app(
app.state.contract_validator = ContractValidator()
app.state.repository_navigation = repository_navigation
app.state.workload_projection = workload_projection
app.state.access_controller = access_controller
if resolved_settings.enforce_access:
app.add_middleware(AccessBoundary, host=app, controller=access_controller)
@app.get("/healthz", response_model=HealthResponse, tags=["system"])
async def healthz() -> HealthResponse:
@ -123,6 +128,8 @@ def create_app(
**await repository_navigation.readiness_checks(),
**await workload_projection.readiness_checks(),
}
if resolved_settings.enforce_access:
dependency_checks["access_profile"] = "ok" if access_controller else "unavailable"
ready = resolved_settings.is_ready(resolved_store.backend_name) and all(
value in {"ok", "not_applicable"} for value in dependency_checks.values()
)

View file

@ -110,7 +110,9 @@ def _run_api(host: str, port: int) -> None:
def _run_mcp(host: str, port: int, transport: str, api_base: str) -> None:
server = HubCoreMCPServer(name="hub-core", api_base=api_base)
server = HubCoreMCPServer(name="hub-core", api_base=api_base,
require_credentials=RuntimeSettings.from_env().enforce_access,
trailing_slash=False)
server.mcp.run(transport=transport, host=host, port=port)

View file

@ -548,6 +548,10 @@ async def _protected(
_enabled(request, group)
if write and group not in request.app.state.settings.v2_write_groups:
raise HTTPException(status_code=503, detail="compatibility group is read-only")
if request.app.state.settings.enforce_access:
if getattr(request.state, "hub_access", None) is None:
raise HTTPException(status_code=503, detail="access boundary unavailable")
return
if not authorization or not authorization.startswith("Bearer "):
raise _unauthorized("Missing bearer token")
token = authorization.removeprefix("Bearer ").strip()

View file

@ -38,9 +38,21 @@ class RuntimeSettings:
legacy_health: bool = False
statehub_inbox_reads: bool = False
statehub_inbox_agent: str = "state-hub"
access_mode: str = "auto"
@property
def enforce_access(self) -> bool:
return self.access_mode == "enforce" or (
self.access_mode == "auto" and self.environment not in {"development", "test"}
)
def __post_init__(self) -> None:
if self.statehub_inbox_reads and (self.backend != "postgresql" or not self.api_token):
if self.access_mode not in {"auto", "enforce", "development"}:
raise ValueError("unsupported access mode")
if self.access_mode == "development" and self.environment not in {"development", "test"}:
raise ValueError("development access is forbidden outside development/test")
if self.statehub_inbox_reads and (self.backend != "postgresql" or
(not self.enforce_access and not self.api_token)):
raise ValueError("State Hub inbox reads require PostgreSQL and operator token")
if self.repo_manager_timeout_seconds <= 0:
raise ValueError("Repo Manager timeout must be positive")
@ -87,6 +99,7 @@ class RuntimeSettings:
legacy_health=_env_bool("HUB_CORE_LEGACY_HEALTH", False),
statehub_inbox_reads=_env_bool("HUB_CORE_STATEHUB_INBOX_READS", False),
statehub_inbox_agent=os.getenv("HUB_CORE_STATEHUB_INBOX_AGENT", "state-hub"),
access_mode=os.getenv("HUB_CORE_ACCESS_MODE", "auto"),
)
def readiness_checks(self, store_backend: str) -> dict[str, str]:
@ -99,7 +112,8 @@ class RuntimeSettings:
"operator",
}
authorization_ready = (
not protected_groups
self.enforce_access
or not protected_groups
or bool(self.api_token)
or store_backend == "postgresql"
)

View file

@ -99,7 +99,10 @@ def create_inbox_projection_router() -> APIRouter:
settings = request.app.state.settings
token = settings.api_token
supplied = (authorization or "").removeprefix("Bearer ")
if not token or not (authorization or "").startswith("Bearer ") or not hmac.compare_digest(supplied, token):
if settings.enforce_access:
if getattr(request.state, "hub_access", None) is None:
raise HTTPException(503, "access boundary unavailable")
elif not token or not (authorization or "").startswith("Bearer ") or not hmac.compare_digest(supplied, token):
raise HTTPException(401, "inbox pilot requires operator bearer authentication",
headers={"WWW-Authenticate": "Bearer"})
if to_agent != settings.statehub_inbox_agent:

View file

@ -25,6 +25,21 @@ 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")
@ -112,6 +127,7 @@ def create_ports_router() -> APIRouter:
)
async def append_progress(
body: EventCommand,
request: Request,
store: PortStore = Depends(get_port_store),
validator: ContractValidator = Depends(get_contract_validator),
) -> PortAccepted:
@ -119,7 +135,7 @@ def create_ports_router() -> APIRouter:
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(body)
return await store.append_progress(_attribute_event(body, request))
@router.post(
"/events/interaction",
@ -130,6 +146,7 @@ def create_ports_router() -> APIRouter:
)
async def append_interaction(
body: EventCommand,
request: Request,
store: PortStore = Depends(get_port_store),
validator: ContractValidator = Depends(get_contract_validator),
) -> PortAccepted:
@ -137,7 +154,7 @@ def create_ports_router() -> APIRouter:
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(body)
return await store.append_interaction(_attribute_event(body, request))
@router.get(
"/projections/{projection_id}",

View file

@ -0,0 +1 @@
"""Hub access profile v1: deny by default, with owner-supplied trust adapters."""

View file

@ -0,0 +1,280 @@
"""Reusable HTTP enforcement for the private platform-root milestone.
Owner adapters establish identity, live account/tenant facts, signed policy
decisions and durable audit. An absent adapter never grants access.
"""
from __future__ import annotations
import asyncio
import hashlib
import json
import math
import time
from dataclasses import dataclass
from importlib.resources import files
from typing import Protocol
from uuid import uuid4
from starlette.requests import Request
from starlette.responses import JSONResponse
from starlette.routing import Match
from hub_core.security.identity import AccessFailure, Actor
PROFILE = "hub-core.access/1.0.0"
@dataclass(frozen=True)
class LiveFacts:
issuer: str
subject: str
actor_tenant: str
target_tenant: str
account_active: bool
actor_tenant_active: bool
target_tenant_active: bool
root_entitled: bool
checked_at: float
evidence_id: str
producer_addresses: frozenset[str] = frozenset()
def __post_init__(self):
for value in (self.account_active, self.actor_tenant_active,
self.target_tenant_active, self.root_entitled):
if type(value) is not bool:
raise ValueError("authoritative state must be boolean")
if type(self.checked_at) not in {int, float} or not math.isfinite(self.checked_at):
raise ValueError("finite fact observation time required")
for value in (self.issuer, self.subject, self.actor_tenant,
self.target_tenant, self.evidence_id, *self.producer_addresses):
if not isinstance(value, str) or not value:
raise ValueError("nonempty authoritative references required")
@dataclass(frozen=True)
class Authorization:
actor: Actor
action: str
resource: str
facts: LiveFacts
correlation_id: str
request_digest: str
@dataclass(frozen=True)
class Decision:
allowed: bool
decision_id: str
policy_version: str
caller: str = ""
def __post_init__(self):
if type(self.allowed) is not bool or not self.decision_id or not self.policy_version:
raise ValueError("explicit boolean decision and provenance required")
class Identity(Protocol):
async def authenticate(self, token: str) -> Actor: ...
class FactSource(Protocol):
async def resolve(self, actor: Actor, resource: str) -> LiveFacts:
"""Query authoritative account/tenant state, without cached root grants."""
...
class Policy(Protocol):
async def evaluate(self, request: Authorization) -> Decision:
"""Verify signed origin, binding, lifetime, caller and obligations."""
...
class Audit(Protocol):
async def append(self, record: dict) -> None:
"""Return only after durable acceptance; raise on delivery failure."""
...
class AccessController:
def __init__(self, *, identity: Identity, facts: FactSource, policy: Policy,
audit: Audit, root_issuer: str, root_subject: str):
if not root_issuer or not root_subject:
raise ValueError("immutable root identity is required")
self.identity, self.facts, self.policy, self.audit = identity, facts, policy, audit
self.root_identity = (root_issuer, root_subject)
async def authorize(self, token: str, action: str, resource: str,
correlation_id: str, request_digest: str) -> Authorization:
actor = await self.identity.authenticate(token)
try:
return await self._authorize_actor(actor, action, resource, correlation_id, request_digest)
except Exception as exc:
failure = exc if isinstance(exc, AccessFailure) else AccessFailure(503, "access_unavailable")
failure.actor = actor
raise failure
async def _authorize_actor(self, actor: Actor, action: str, resource: str,
correlation_id: str, request_digest: str) -> Authorization:
if actor.expires_at <= time.time():
raise AccessFailure(401, "expired_access_token")
if actor.principal_type == "human" and (
(actor.issuer, actor.subject) != self.root_identity
or actor.tenant != "tenant:platform" or actor.assurance not in {"aal2", "aal3"}
):
raise AccessFailure(403, "root_required")
facts = await self.facts.resolve(actor, resource)
if (facts.issuer, facts.subject, facts.actor_tenant) != (
actor.issuer, actor.subject, actor.tenant
) or not 0 <= time.time() - facts.checked_at <= 5 or not facts.evidence_id:
raise AccessFailure(503, "untrusted_or_stale_facts")
# v1 explicitly classifies Hub records as platform-owned. Other tenants
# require the Phase 2 resource resolver/storage contract, not a header.
if facts.target_tenant != "tenant:platform":
raise AccessFailure(403, "unsupported_target_tenant")
if not (facts.account_active and facts.actor_tenant_active and facts.target_tenant_active):
raise AccessFailure(403, "inactive_identity_or_tenant")
if actor.principal_type == "human" and not facts.root_entitled:
raise AccessFailure(403, "root_entitlement_required")
context = Authorization(actor, action, resource, facts, correlation_id, request_digest)
decision = await self.policy.evaluate(context)
await self.audit.append({
"profile": PROFILE, "correlation_id": correlation_id,
"issuer": actor.issuer, "subject": actor.subject,
"principal_type": actor.principal_type, "actor_tenant": actor.tenant,
"target_tenant": facts.target_tenant, "action": action,
"resource_digest": hashlib.sha256(resource.encode()).hexdigest(),
"request_digest": request_digest, "facts_evidence": facts.evidence_id,
"decision_id": decision.decision_id, "policy_version": decision.policy_version,
"policy_caller": decision.caller,
"outcome": "authorized" if decision.allowed else "denied",
})
if not decision.allowed:
raise AccessFailure(403, "policy_denied")
if actor.expires_at <= time.time():
raise AccessFailure(401, "expired_access_token")
if time.time() - facts.checked_at > 5:
raise AccessFailure(503, "facts_expired_during_authorization")
return context
def route_key(route, method: str) -> str:
endpoint = route.endpoint
return f"{method}:{route.path}:{endpoint.__module__}.{endpoint.__name__}"
def iter_routes(router):
for route in router.routes:
if hasattr(route, "original_router"):
yield from iter_routes(route.original_router)
else:
yield route
class AccessBoundary:
"""ASGI boundary also covers docs, redirects, unknown routes and WebSockets.
Install on an embedded host with its own explicit catalog to protect SDK
routers. Routes added without catalog admission remain denied.
"""
def __init__(self, app, *, host, controller: AccessController | None,
catalog: dict[str, str] | None = None):
self.app, self.host, self.controller = app, host, controller
self.catalog = catalog if catalog is not None else json.loads(
files("hub_core.security").joinpath("routes.json").read_text()
)["routes"]
async def __call__(self, scope, receive, send):
if scope["type"] == "websocket":
await send({"type": "websocket.close", "code": 1008})
return
if scope["type"] != "http":
await self.app(scope, receive, send)
return
# No prefix or trailing-slash exception. Detailed readiness is protected.
if scope["method"] == "GET" and scope["path"] == "/healthz":
await JSONResponse({"status": "ok"})(scope, receive, send)
return
correlation = str(uuid4())
context = None
try:
headers = Request(scope).headers.getlist("authorization")
if len(headers) != 1 or not headers[0].startswith("Bearer "):
raise AccessFailure(401, "bearer_required")
token = headers[0][7:]
if not token or len(token) > 16384 or any(c.isspace() for c in token):
raise AccessFailure(401, "invalid_access_token")
if self.controller is None:
raise AccessFailure(503, "access_dependencies_unavailable")
route = next((r for r in iter_routes(self.host)
if r.matches(scope)[0] == Match.FULL), None)
key = route_key(route, scope["method"]) if route and hasattr(route, "endpoint") else None
action = self.catalog.get(key)
if not action:
raise AccessFailure(403, "surface_not_admitted")
# Bind policy to the exact request without exposing content to PDP/audit.
request = Request(scope, receive)
chunks, size = [], 0
async for chunk in request.stream():
size += len(chunk)
if size > 1024 * 1024:
raise AccessFailure(413, "request_too_large")
chunks.append(chunk)
body = b"".join(chunks)
digest = hashlib.sha256(b"\0".join([
scope["method"].encode(), scope["path"].encode(),
scope.get("query_string", b""), body,
])).hexdigest()
async with asyncio.timeout(10):
context = await self.controller.authorize(
token, action, scope["path"], correlation, digest,
)
if body:
try:
payload = json.loads(body)
except (ValueError, UnicodeError):
payload = None # Handler owns content validation.
if isinstance(payload, dict):
for field in ("from_address", "from_agent", "author"):
if field in payload and payload[field] not in context.facts.producer_addresses:
raise AccessFailure(403, "producer_identity_mismatch")
scope.setdefault("state", {})["hub_access"] = context
except Exception as exc:
failure = exc if isinstance(exc, AccessFailure) else AccessFailure(503, "access_unavailable")
actor = context.actor if context else getattr(failure, "actor", None)
if self.controller is not None:
try:
async with asyncio.timeout(3):
await self.controller.audit.append({
"profile": PROFILE, "correlation_id": correlation,
"outcome": "refused", "reason": failure.code,
"subject": actor.subject if actor else None,
"issuer": actor.issuer if actor else None,
"actor_tenant": actor.tenant if actor else None,
})
except Exception:
failure = AccessFailure(503, "audit_unavailable")
headers = {"Cache-Control": "no-store", "X-Correlation-ID": correlation}
if failure.status == 401:
headers["WWW-Authenticate"] = "Bearer"
await JSONResponse({"detail": failure.code}, status_code=failure.status,
headers=headers)(scope, receive, send)
return
delivered = False
async def replay():
nonlocal delivered
if not delivered:
delivered = True
return {"type": "http.request", "body": body, "more_body": False}
return await receive()
async def protected_send(message):
if message["type"] == "http.response.start":
message["headers"] = [(k, v) for k, v in message.get("headers", [])
if k.lower() not in {b"cache-control", b"x-correlation-id"}]
message["headers"].extend([(b"cache-control", b"no-store"),
(b"x-correlation-id", correlation.encode())])
await send(message)
await self.app(scope, replay, protected_send)

View file

@ -0,0 +1,168 @@
"""IAM v0.3 access-token verification. No username-based authority."""
from __future__ import annotations
import asyncio
import ipaddress
import time
from dataclasses import dataclass
from urllib.parse import urlsplit
import httpx
import jwt
class AccessFailure(Exception):
def __init__(self, status: int, code: str):
self.status, self.code = status, code
super().__init__(code)
@dataclass(frozen=True)
class Actor:
issuer: str
subject: str
tenant: str
principal_type: str
assurance: str
authenticated_at: int
expires_at: int
def require_https(url: str) -> None:
parsed = urlsplit(url)
hostname = (parsed.hostname or "").rstrip(".").lower()
try:
local = ipaddress.ip_address(hostname).is_loopback
except ValueError:
local = hostname == "localhost" or hostname.endswith(".localhost")
if (parsed.scheme != "https" or not parsed.hostname or parsed.username
or parsed.password or parsed.fragment or parsed.query
or local):
raise ValueError("a non-local HTTPS trust endpoint is required")
class OIDCVerifier:
"""Discover keys at an explicitly trusted issuer; bounded, rotation-aware cache.
Supports RFC 9068 at+jwt tokens or the admitted KeyCape Bearer payload type.
ID tokens without either access-token marker are rejected.
"""
def __init__(self, *, issuer: str, audience: str, client: httpx.AsyncClient,
key_ttl: int = 60, max_token_age: int = 300):
require_https(issuer)
if not audience or not 1 <= key_ttl <= 300 or not 1 <= max_token_age <= 300:
raise ValueError("audience and bounded key/token lifetimes are required")
self.issuer, self.audience, self.client = issuer, audience, client
self.key_ttl, self.max_token_age = key_ttl, max_token_age
self._keys: dict = {}
self._loaded = 0.0
self._lock = asyncio.Lock()
async def _refresh(self) -> None:
try:
response = await self.client.get(
self.issuer.rstrip("/") + "/.well-known/openid-configuration",
timeout=3, follow_redirects=False,
)
response.raise_for_status()
discovery = response.json()
if discovery["issuer"] != self.issuer:
raise ValueError("issuer mismatch")
require_https(discovery["jwks_uri"])
response = await self.client.get(discovery["jwks_uri"], timeout=3,
follow_redirects=False)
response.raise_for_status()
keys = {}
for value in response.json()["keys"]:
if value.get("kty") != "RSA" or value.get("use", "sig") != "sig":
continue
if value.get("alg", "RS256") != "RS256":
continue
if "verify" not in value.get("key_ops", ["verify"]):
continue
kid = value["kid"]
if not isinstance(kid, str) or not kid or kid in keys:
raise ValueError("invalid key IDs")
key = jwt.PyJWK.from_dict(value, algorithm="RS256").key
if key.key_size < 2048:
raise ValueError("weak issuer key")
keys[kid] = key
if not keys:
raise ValueError("no signing keys")
self._keys, self._loaded = keys, time.monotonic()
except (httpx.HTTPError, ValueError, KeyError, TypeError, jwt.PyJWTError) as exc:
raise AccessFailure(503, "identity_unavailable") from exc
async def authenticate(self, token: str) -> Actor:
try:
header = jwt.get_unverified_header(token)
if header.get("alg") != "RS256" or not isinstance(header.get("kid"), str):
raise ValueError("unsupported token")
async with self._lock:
# At most one unknown-key refresh per second, to bound random-kid traffic.
age = time.monotonic() - self._loaded
if age >= self.key_ttl or (header["kid"] not in self._keys and age >= 1):
await self._refresh()
key = self._keys.get(header["kid"])
if key is None:
raise ValueError("unknown key")
claims = jwt.decode(token, key, algorithms=["RS256"], issuer=self.issuer,
audience=self.audience, leeway=0,
options={"require": ["iss", "sub", "aud", "exp", "iat",
"tenant", "principal_type", "groups",
"roles", "assurance"]})
if header.get("typ") != "at+jwt" and claims.get("typ") != "Bearer":
raise ValueError("not an access token")
if claims.get("typ", "Bearer") != "Bearer" or claims.get("environment") in {
"local", "development", "test",
}:
raise ValueError("unsupported token profile")
for name in ("sub", "tenant"):
if not isinstance(claims[name], str) or not claims[name]:
raise ValueError("invalid identity")
for name in ("groups", "roles"):
if not isinstance(claims[name], list) or any(
not isinstance(item, str) for item in claims[name]
):
raise ValueError("invalid IAM array")
scope = claims.get("scope", claims.get("scp"))
if not isinstance(scope, (str, list)) or (
isinstance(scope, list) and any(not isinstance(item, str) for item in scope)
):
raise ValueError("invalid scope")
assurance = claims["assurance"]
if (not isinstance(assurance, dict)
or assurance.get("level") not in {"aal1", "aal2", "aal3"}
or type(assurance.get("mfa")) is not bool
or not isinstance(assurance.get("source"), str)
or not assurance["source"]
or not isinstance(assurance.get("methods"), list)
or not all(isinstance(x, str) for x in assurance["methods"])):
raise ValueError("invalid assurance")
for value in (claims["iat"], claims["exp"], assurance.get("at")):
if type(value) is not int:
raise ValueError("integer timestamps required")
if "nbf" in claims and type(claims["nbf"]) is not int:
raise ValueError("integer not-before required")
if assurance["level"] in {"aal2", "aal3"} and not assurance["mfa"]:
raise ValueError("missing MFA evidence")
now = time.time()
if (not claims["iat"] <= now < claims["exp"]
or claims["exp"] - claims["iat"] > self.max_token_age
or not 0 <= now - assurance["at"] <= self.max_token_age):
raise ValueError("stale identity or assurance")
principal = claims["principal_type"]
if principal not in {"human", "service", "agent"}:
raise ValueError("invalid principal type")
if principal == "agent":
agent = claims.get("agent", {})
if not agent.get("id") or agent.get("mode") != "autonomous":
# Delegation needs a separately admitted actor/workload contract.
raise ValueError("unsupported delegation")
return Actor(self.issuer, claims["sub"], claims["tenant"], principal,
assurance["level"], assurance["at"], claims["exp"])
except AccessFailure:
raise
except (jwt.PyJWTError, ValueError, KeyError, TypeError, AttributeError) as exc:
raise AccessFailure(401, "invalid_access_token") from exc

177
hub_core/security/policy.py Normal file
View file

@ -0,0 +1,177 @@
"""Authenticated flex-auth client; verification precedes every allow/deny."""
from __future__ import annotations
import base64
import hashlib
import json
from datetime import datetime, timezone
from pathlib import Path
import httpx
from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PublicKey
from hub_core.security.boundary import Authorization, Decision
from hub_core.security.identity import AccessFailure, require_https
def _object(pairs):
result = {}
for key, value in pairs:
if key in result:
raise ValueError("duplicate JSON key")
result[key] = value
return result
def parse_json(raw: bytes | str):
def reject(_):
raise ValueError("non-integer numbers are outside this profile")
return json.loads(raw, object_pairs_hook=_object, parse_float=reject, parse_constant=reject)
def go_json(value) -> bytes:
"""Preserve Go struct/wire order and escape HTML as encoding/json does.
This deliberately does not sort struct fields. Maps sent in CheckRequest
are sorted separately. Unsupported floating point input fails closed.
"""
encoded = json.dumps(value, ensure_ascii=False, separators=(",", ":"), allow_nan=False)
for char, escaped in (("<", "\\u003c"), (">", "\\u003e"), ("&", "\\u0026"),
("\u2028", "\\u2028"), ("\u2029", "\\u2029")):
encoded = encoded.replace(char, escaped)
return encoded.encode()
def _sorted_maps(value):
if isinstance(value, dict):
return {k: _sorted_maps(value[k]) for k in sorted(value)}
if isinstance(value, list):
return [_sorted_maps(x) for x in value]
return value
def submitted_digest(request: dict) -> str:
# requestDigestMaterial, SubjectRef and ResourceRef are Go structs, whose
# declaration order (unlike maps) participates in the current wire contract.
material = {}
if request.get("tenant"):
material["tenant"] = request["tenant"]
for field, order in (("subject", ("id", "type", "tenant", "attributes")),
("resource", ("id", "type", "system", "tenant", "attributes"))):
if field == "resource":
material["action"] = request["action"]
material[field] = {k: _sorted_maps(request[field][k]) for k in order
if request[field].get(k)}
if request.get("context"):
material["context"] = _sorted_maps(request["context"])
return "sha256:" + hashlib.sha256(go_json(material)).hexdigest()
def verify_signature(envelope: dict, keys: dict) -> None:
signature = envelope["signature"]
if signature["mode"] != "signed" or signature["alg"] != "ed25519":
raise ValueError("signed Ed25519 decision required")
candidates = [key for key in keys["keys"] if key["kid"] == signature["kid"]]
if len(candidates) != 1 or candidates[0]["alg"] != "ed25519":
raise ValueError("untrusted signing key")
def decode(value):
return base64.b64decode(value + "=" * (-len(value) % 4), altchars=b"-_", validate=True)
key = Ed25519PublicKey.from_public_bytes(decode(candidates[0]["public_key"]))
key.verify(decode(signature["value"]), go_json({
k: v for k, v in envelope.items() if k != "signature"
}))
def _time(value: str) -> datetime:
result = datetime.fromisoformat(value.replace("Z", "+00:00"))
if result.tzinfo is None:
raise ValueError("timezone required")
return result
def verify_decision(envelope: dict, *, request: dict, keys: dict,
caller: str, now: datetime | None = None) -> Decision:
verify_signature(envelope, keys)
now = now or datetime.now(timezone.utc)
if (envelope["contract_version"] != "flex-auth.decision-record.v1"
or envelope["request_id"] != request["id"] or not envelope["id"]):
raise ValueError("invalid decision contract or correlation")
binding = envelope["binding"]
if binding["submitted_request_digest"] != submitted_digest(request):
raise ValueError("submitted request mismatch")
if binding["action"] != request["action"] or binding.get("tenant") != request["tenant"]:
raise ValueError("action or tenant mismatch")
for field, fields in (("subject", ("id", "type", "tenant")),
("resource", ("id", "type", "system", "tenant"))):
for key in fields:
if binding[field].get(key) != request[field].get(key):
raise ValueError("evaluated identity or resource mismatch")
if envelope[field] != binding[field]:
raise ValueError("inconsistent binding")
if binding.get("context", {}) != request.get("context", {}):
raise ValueError("context mismatch")
provenance = envelope["provenance"]
caller_record = provenance["caller"]
if (caller_record["mode"] != "enforce" or caller_record["principal"] != caller
or caller_record["audience"] != "flex-auth"
or _time(caller_record["not_after"]) <= now):
raise ValueError("untrusted workload caller")
if not provenance["policy_version"] or not provenance["policy_package_digest"]:
raise ValueError("missing policy provenance")
age = (now - _time(provenance["decision_time"])).total_seconds()
if not 0 <= age <= 30:
raise ValueError("stale decision")
if envelope.get("obligations"):
# No obligation is silently treated as satisfied. Owner-specific
# approval/redaction/audit handlers require a later profile revision.
raise ValueError("unsupported decision obligations")
if envelope["effect"] not in {"allow", "deny"}:
raise ValueError("unsupported effect")
if envelope["effect"] == "allow":
lifetime = envelope["lifetime"]
if (lifetime["kind"] != "ttl" or not
_time(lifetime["not_before"]) <= now < _time(lifetime["expires_at"])):
raise ValueError("invalid decision lifetime")
return Decision(envelope["effect"] == "allow", envelope["id"], provenance["policy_version"], caller)
class FlexPolicy:
def __init__(self, *, base_url: str, client: httpx.AsyncClient,
caller_token_file: Path, trusted_keys_file: Path, caller: str):
require_https(base_url)
if not caller.startswith("system:serviceaccount:"):
raise ValueError("explicit admitted workload caller required")
self.base_url, self.client = base_url.rstrip("/"), client
self.caller_token_file, self.trusted_keys_file = caller_token_file, trusted_keys_file
self.caller = caller
async def evaluate(self, request: Authorization) -> Decision:
actor, facts = request.actor, request.facts
check = {
"id": request.correlation_id, "tenant": facts.target_tenant,
"subject": {"id": actor.subject, "type": actor.principal_type,
"tenant": actor.tenant,
"attributes": {"issuer": actor.issuer, "assurance": actor.assurance}},
"action": request.action,
"resource": {"id": request.resource, "type": "hub-route", "system": "hub-core",
"tenant": facts.target_tenant},
"context": {"http_request_digest": request.request_digest,
"facts_evidence": facts.evidence_id,
"root_entitled": facts.root_entitled},
}
try:
# Reread projected credentials and owner-delivered public trust at
# every check. No remote key response can bootstrap its own trust.
token = self.caller_token_file.read_text().strip()
if not token or any(c.isspace() for c in token):
raise ValueError("invalid caller credential")
keys = parse_json(self.trusted_keys_file.read_bytes())
response = await self.client.post(self.base_url + "/v1/check", json=check,
headers={"Authorization": f"Bearer {token}"},
timeout=3, follow_redirects=False)
response.raise_for_status()
return verify_decision(parse_json(response.content), request=check,
keys=keys, caller=self.caller)
except Exception as exc:
# Neither response body nor credentials appear in the public error.
raise AccessFailure(503, "policy_unavailable_or_untrusted") from exc

View file

@ -0,0 +1,93 @@
{
"profile": "hub-core.access/1.0.0",
"status": "candidate-owner-review-required",
"routes": {
"GET:/annotation-categories:hub_core.runtime.compat.annotation_categories": "hub.hub_core.runtime.compat.annotation_categories.get",
"GET:/annotations:hub_core.runtime.compat.empty_collection": "hub.hub_core.runtime.compat.empty_collection.get",
"GET:/api-consumers:hub_core.runtime.compat.list_consumers": "hub.hub_core.runtime.compat.list_consumers.get",
"GET:/api/v2/annotation-categories:hub_core.runtime.compat.annotation_categories": "hub.hub_core.runtime.compat.annotation_categories.get",
"GET:/api/v2/annotations:hub_core.runtime.compat.empty_collection": "hub.hub_core.runtime.compat.empty_collection.get",
"GET:/api/v2/api-consumers:hub_core.runtime.compat.list_consumers": "hub.hub_core.runtime.compat.list_consumers.get",
"GET:/api/v2/decision-records:hub_core.runtime.compat.empty_collection": "hub.hub_core.runtime.compat.empty_collection.get",
"GET:/api/v2/deployment-records:hub_core.runtime.compat.empty_collection": "hub.hub_core.runtime.compat.empty_collection.get",
"GET:/api/v2/docs:hub_core.runtime.compat.docs": "hub.hub_core.runtime.compat.docs.get",
"GET:/api/v2/event-types:hub_core.runtime.compat.event_types": "hub.hub_core.runtime.compat.event_types.get",
"GET:/api/v2/hub-capability-manifests:hub_core.runtime.compat.list_manifests": "hub.hub_core.runtime.compat.list_manifests.get",
"GET:/api/v2/hub-registry:hub_core.runtime.compat.hub_registry": "hub.hub_core.runtime.compat.hub_registry.get",
"GET:/api/v2/hubs:hub_core.runtime.compat.list_hubs": "hub.hub_core.runtime.compat.list_hubs.get",
"GET:/api/v2/interaction-events:hub_core.runtime.compat.list_interactions": "hub.hub_core.runtime.compat.list_interactions.get",
"GET:/api/v2/openapi.json:hub_core.runtime.compat.openapi_json": "hub.hub_core.runtime.compat.openapi_json.get",
"GET:/api/v2/openapi.yaml:hub_core.runtime.compat.openapi_yaml": "hub.hub_core.runtime.compat.openapi_yaml.get",
"GET:/api/v2/outcome-signals:hub_core.runtime.compat.empty_collection": "hub.hub_core.runtime.compat.empty_collection.get",
"GET:/api/v2/policy-scopes:hub_core.runtime.compat.policy_scopes": "hub.hub_core.runtime.compat.policy_scopes.get",
"GET:/api/v2/requirement-candidates:hub_core.runtime.compat.empty_collection": "hub.hub_core.runtime.compat.empty_collection.get",
"GET:/api/v2/widget-types:hub_core.runtime.compat.widget_types": "hub.hub_core.runtime.compat.widget_types.get",
"GET:/api/v2/widgets:hub_core.runtime.compat.list_widgets": "hub.hub_core.runtime.compat.list_widgets.get",
"GET:/console:hub_core.runtime.compat.console": "hub.hub_core.runtime.compat.console.get",
"GET:/decision-records:hub_core.runtime.compat.empty_collection": "hub.hub_core.runtime.compat.empty_collection.get",
"GET:/deployment-records:hub_core.runtime.compat.empty_collection": "hub.hub_core.runtime.compat.empty_collection.get",
"GET:/docs/oauth2-redirect:fastapi.applications.swagger_ui_redirect": "hub.fastapi.applications.swagger_ui_redirect.get",
"GET:/docs:fastapi.applications.swagger_ui_html": "hub.fastapi.applications.swagger_ui_html.get",
"GET:/docs:hub_core.runtime.compat.docs": "hub.hub_core.runtime.compat.docs.get",
"GET:/event-types:hub_core.runtime.compat.event_types": "hub.hub_core.runtime.compat.event_types.get",
"GET:/hub-capability-manifests:hub_core.runtime.compat.list_manifests": "hub.hub_core.runtime.compat.list_manifests.get",
"GET:/hub-registry:hub_core.runtime.compat.hub_registry": "hub.hub_core.runtime.compat.hub_registry.get",
"GET:/hubs:hub_core.runtime.compat.list_hubs": "hub.hub_core.runtime.compat.list_hubs.get",
"GET:/interaction-events:hub_core.runtime.compat.list_interactions": "hub.hub_core.runtime.compat.list_interactions.get",
"GET:/openapi.json:fastapi.applications.openapi": "hub.fastapi.applications.openapi.get",
"GET:/openapi.json:hub_core.runtime.compat.openapi_json": "hub.hub_core.runtime.compat.openapi_json.get",
"GET:/openapi.yaml:hub_core.runtime.compat.openapi_yaml": "hub.hub_core.runtime.compat.openapi_yaml.get",
"GET:/outcome-signals:hub_core.runtime.compat.empty_collection": "hub.hub_core.runtime.compat.empty_collection.get",
"GET:/policy-scopes:hub_core.runtime.compat.policy_scopes": "hub.hub_core.runtime.compat.policy_scopes.get",
"GET:/ports/messaging/messages:hub_core.runtime.ports.list_messages": "hub.hub_core.runtime.ports.list_messages.get",
"GET:/ports/projections/repository-navigation/facets/{facet_kind}/{facet_value}:hub_core.runtime.repository_navigation_routes.query_facet": "hub.hub_core.runtime.repository_navigation_routes.query_facet.get",
"GET:/ports/projections/repository-navigation/repositories:hub_core.runtime.repository_navigation_routes.query_repositories": "hub.hub_core.runtime.repository_navigation_routes.query_repositories.get",
"GET:/ports/projections/statehub-inbox:hub_core.runtime.inbox_projection.inbox": "hub.hub_core.runtime.inbox_projection.inbox.get",
"GET:/ports/projections/workloads/resolve:hub_core.runtime.workload_projection_routes.resolve_workload": "hub.hub_core.runtime.workload_projection_routes.resolve_workload.get",
"GET:/ports/projections/workloads:hub_core.runtime.workload_projection_routes.query_workloads": "hub.hub_core.runtime.workload_projection_routes.query_workloads.get",
"GET:/ports/projections/{projection_id}:hub_core.runtime.ports.query_projection": "hub.hub_core.runtime.ports.query_projection.get",
"GET:/ports/registry/registrations/{hub_slug}/audit:hub_core.runtime.ports.registration_audit": "hub.hub_core.runtime.ports.registration_audit.get",
"GET:/ports/registry/registrations/{hub_slug}:hub_core.runtime.ports.resolve_registration": "hub.hub_core.runtime.ports.resolve_registration.get",
"GET:/readyz:hub_core.runtime.app.readyz": "hub.hub_core.runtime.app.readyz.get",
"GET:/redoc:fastapi.applications.redoc_html": "hub.fastapi.applications.redoc_html.get",
"GET:/requirement-candidates:hub_core.runtime.compat.empty_collection": "hub.hub_core.runtime.compat.empty_collection.get",
"GET:/widget-types:hub_core.runtime.compat.widget_types": "hub.hub_core.runtime.compat.widget_types.get",
"GET:/widgets:hub_core.runtime.compat.list_widgets": "hub.hub_core.runtime.compat.list_widgets.get",
"HEAD:/docs/oauth2-redirect:fastapi.applications.swagger_ui_redirect": "hub.fastapi.applications.swagger_ui_redirect.head",
"HEAD:/docs:fastapi.applications.swagger_ui_html": "hub.fastapi.applications.swagger_ui_html.head",
"HEAD:/openapi.json:fastapi.applications.openapi": "hub.fastapi.applications.openapi.head",
"HEAD:/redoc:fastapi.applications.redoc_html": "hub.fastapi.applications.redoc_html.head",
"PATCH:/api/v2/hub-capability-manifests/{manifest_id}:hub_core.runtime.compat.patch_manifest": "hub.hub_core.runtime.compat.patch_manifest.patch",
"PATCH:/hub-capability-manifests/{manifest_id}:hub_core.runtime.compat.patch_manifest": "hub.hub_core.runtime.compat.patch_manifest.patch",
"POST:/annotations:hub_core.runtime.compat.accept_deferred": "hub.hub_core.runtime.compat.accept_deferred.post",
"POST:/api-consumers/{consumer_id}/api-keys:hub_core.runtime.compat.create_key": "hub.hub_core.runtime.compat.create_key.post",
"POST:/api-consumers:hub_core.runtime.compat.create_consumer": "hub.hub_core.runtime.compat.create_consumer.post",
"POST:/api/v2/annotations:hub_core.runtime.compat.accept_deferred": "hub.hub_core.runtime.compat.accept_deferred.post",
"POST:/api/v2/api-consumers/{consumer_id}/api-keys:hub_core.runtime.compat.create_key": "hub.hub_core.runtime.compat.create_key.post",
"POST:/api/v2/api-consumers:hub_core.runtime.compat.create_consumer": "hub.hub_core.runtime.compat.create_consumer.post",
"POST:/api/v2/decision-records:hub_core.runtime.compat.accept_deferred": "hub.hub_core.runtime.compat.accept_deferred.post",
"POST:/api/v2/deployment-records:hub_core.runtime.compat.accept_deferred": "hub.hub_core.runtime.compat.accept_deferred.post",
"POST:/api/v2/hub-capability-manifests/{manifest_id}/activate:hub_core.runtime.compat.activate_manifest": "hub.hub_core.runtime.compat.activate_manifest.post",
"POST:/api/v2/hub-capability-manifests:hub_core.runtime.compat.create_manifest": "hub.hub_core.runtime.compat.create_manifest.post",
"POST:/api/v2/hubs:hub_core.runtime.compat.create_hub": "hub.hub_core.runtime.compat.create_hub.post",
"POST:/api/v2/interaction-events:hub_core.runtime.compat.create_interaction": "hub.hub_core.runtime.compat.create_interaction.post",
"POST:/api/v2/outcome-signals:hub_core.runtime.compat.accept_deferred": "hub.hub_core.runtime.compat.accept_deferred.post",
"POST:/api/v2/requirement-candidates:hub_core.runtime.compat.accept_deferred": "hub.hub_core.runtime.compat.accept_deferred.post",
"POST:/api/v2/token:hub_core.runtime.compat.token": "hub.hub_core.runtime.compat.token.post",
"POST:/api/v2/widgets:hub_core.runtime.compat.create_widget": "hub.hub_core.runtime.compat.create_widget.post",
"POST:/decision-records:hub_core.runtime.compat.accept_deferred": "hub.hub_core.runtime.compat.accept_deferred.post",
"POST:/deployment-records:hub_core.runtime.compat.accept_deferred": "hub.hub_core.runtime.compat.accept_deferred.post",
"POST:/hub-capability-manifests/{manifest_id}/activate:hub_core.runtime.compat.activate_manifest": "hub.hub_core.runtime.compat.activate_manifest.post",
"POST:/hub-capability-manifests:hub_core.runtime.compat.create_manifest": "hub.hub_core.runtime.compat.create_manifest.post",
"POST:/hubs:hub_core.runtime.compat.create_hub": "hub.hub_core.runtime.compat.create_hub.post",
"POST:/interaction-events:hub_core.runtime.compat.create_interaction": "hub.hub_core.runtime.compat.create_interaction.post",
"POST:/outcome-signals:hub_core.runtime.compat.accept_deferred": "hub.hub_core.runtime.compat.accept_deferred.post",
"POST:/ports/events/interaction:hub_core.runtime.ports.append_interaction": "hub.hub_core.runtime.ports.append_interaction.post",
"POST:/ports/events/progress:hub_core.runtime.ports.append_progress": "hub.hub_core.runtime.ports.append_progress.post",
"POST:/ports/messaging/messages:hub_core.runtime.ports.send_message": "hub.hub_core.runtime.ports.send_message.post",
"POST:/ports/registry/registrations:hub_core.runtime.ports.register_extension": "hub.hub_core.runtime.ports.register_extension.post",
"POST:/requirement-candidates:hub_core.runtime.compat.accept_deferred": "hub.hub_core.runtime.compat.accept_deferred.post",
"POST:/token:hub_core.runtime.compat.token": "hub.hub_core.runtime.compat.token.post",
"POST:/widgets:hub_core.runtime.compat.create_widget": "hub.hub_core.runtime.compat.create_widget.post"
}
}