Assistant: codex Assistant-Model: gpt-6-astra Assistant-Session: 01a0e747-8f27-7242-8df8-8bc44f88c929
315 lines
14 KiB
Python
315 lines
14 KiB
Python
"""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 = ""
|
|
signed_envelope: str | None = None
|
|
|
|
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,
|
|
"signed_decision": decision.signed_envelope,
|
|
"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, body_timeout: float = 10, browser=None):
|
|
if not 0 < body_timeout <= 10:
|
|
raise ValueError("body timeout must be positive and at most ten seconds")
|
|
self.app, self.host, self.controller = app, host, controller
|
|
self.body_timeout = body_timeout
|
|
self.browser = browser
|
|
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
|
|
if self.browser is not None:
|
|
from hub_core.security.browser import BROWSER_ROUTES
|
|
if (scope["method"], scope["path"]) in BROWSER_ROUTES:
|
|
response = await self.browser.handle(Request(scope, receive))
|
|
await response(scope, receive, send)
|
|
return
|
|
correlation = str(uuid4())
|
|
context = None
|
|
cookie_auth = False
|
|
action = None
|
|
digest = None
|
|
try:
|
|
request = Request(scope, receive)
|
|
headers = request.headers.getlist("authorization")
|
|
if self.browser is not None and not headers:
|
|
token = self.browser.access_token(request)
|
|
cookie_auth = True
|
|
else:
|
|
if self.browser is not None:
|
|
from hub_core.security.browser import SESSION_COOKIE
|
|
if SESSION_COOKIE in request.cookies:
|
|
raise AccessFailure(400, "mixed_browser_credentials")
|
|
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
|
|
try:
|
|
async with asyncio.timeout(self.body_timeout):
|
|
async for chunk in request.stream():
|
|
size += len(chunk)
|
|
if size > 1024 * 1024:
|
|
raise AccessFailure(413, "request_too_large")
|
|
chunks.append(chunk)
|
|
except TimeoutError as exc:
|
|
raise AccessFailure(408, "request_body_timeout") from exc
|
|
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")
|
|
if cookie_auth:
|
|
self.browser.access_token(request) # Recheck expiry/logout after owner awaits.
|
|
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,
|
|
"principal_type": actor.principal_type if actor else None,
|
|
"action": action,
|
|
"resource_digest": hashlib.sha256(scope["path"].encode()).hexdigest(),
|
|
"request_digest": digest,
|
|
"target_tenant": "tenant:platform",
|
|
})
|
|
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)
|