feat: integrate durable authorization audit and runtime composition
Assistant: codex Assistant-Model: gpt-6-astra Assistant-Session: 01a0e747-8f27-7242-8df8-8bc44f88c929
This commit is contained in:
parent
3e386147fd
commit
c9b6916dac
14 changed files with 767 additions and 16 deletions
|
|
@ -3,6 +3,8 @@ from __future__ import annotations
|
|||
import asyncio
|
||||
from contextlib import asynccontextmanager, suppress
|
||||
|
||||
import httpx
|
||||
|
||||
from fastapi import FastAPI, Response, status
|
||||
|
||||
from hub_core import __version__
|
||||
|
|
@ -27,7 +29,8 @@ 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
|
||||
from hub_core.security.boundary import AccessBoundary, AccessController, FactSource
|
||||
from hub_core.security.config import SecuritySettings
|
||||
|
||||
|
||||
def create_app(
|
||||
|
|
@ -37,8 +40,25 @@ def create_app(
|
|||
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,
|
||||
) -> FastAPI:
|
||||
resolved_settings = settings or RuntimeSettings.from_env()
|
||||
# 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)
|
||||
resolved_store = port_store or _create_store(resolved_settings)
|
||||
owns_store = port_store is None
|
||||
resolved_repo_projection_client = repo_projection_client
|
||||
|
|
@ -84,15 +104,19 @@ def create_app(
|
|||
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()
|
||||
try:
|
||||
yield
|
||||
finally:
|
||||
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",
|
||||
|
|
|
|||
77
hub_core/security/audit.py
Normal file
77
hub_core/security/audit.py
Normal file
|
|
@ -0,0 +1,77 @@
|
|||
"""Synchronous authorization custody using Audit Core's eight-field contract."""
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
from uuid import uuid4
|
||||
|
||||
import httpx
|
||||
|
||||
from hub_core.security.identity import AccessFailure, require_https
|
||||
|
||||
|
||||
class AuditCoreSink:
|
||||
"""No allow returns before durable remote custody acknowledges its record.
|
||||
|
||||
A lost receipt blocks execution, even if the archive already stored the
|
||||
attempt. This is an authorization-attempt journal, not a mutation outbox.
|
||||
"""
|
||||
|
||||
def __init__(self, *, base_url: str, token_file: Path, client: httpx.AsyncClient):
|
||||
require_https(base_url)
|
||||
self.base_url = base_url.rstrip("/")
|
||||
self.token_file, self.client = token_file, client
|
||||
|
||||
async def readiness(self) -> None:
|
||||
response = await self.client.get(self.base_url + "/readyz", timeout=2,
|
||||
follow_redirects=False)
|
||||
response.raise_for_status()
|
||||
receipt = response.json()
|
||||
if (receipt.get("status") != "ok" or receipt.get("durable") is not True
|
||||
or receipt.get("custody_class") not in {"archive", "operational"}):
|
||||
raise ValueError("operational audit custody is required")
|
||||
|
||||
async def append(self, record: dict) -> None:
|
||||
try:
|
||||
async with asyncio.timeout(3):
|
||||
# Probe each time, so a receiver's development fallback cannot
|
||||
# be mistaken for admitted custody through a cached readiness.
|
||||
await self.readiness()
|
||||
token = self.token_file.read_text().strip()
|
||||
if not token or not token.isascii() or any(c.isspace() for c in token):
|
||||
raise ValueError("invalid sender credential")
|
||||
correlation = record.get("correlation_id")
|
||||
outcome = record.get("outcome")
|
||||
if not isinstance(correlation, str) or not correlation or outcome not in {
|
||||
"authorized", "denied", "refused",
|
||||
}:
|
||||
raise ValueError("invalid authorization audit record")
|
||||
event = {
|
||||
"id": str(uuid4()),
|
||||
"type": "hub.access." + outcome,
|
||||
"source": "hub-core",
|
||||
"subject": "hub-access:" + correlation,
|
||||
"tenant": "tenant:platform",
|
||||
"correlation_id": correlation,
|
||||
"occurred_at": datetime.now(timezone.utc).isoformat(),
|
||||
"data": record,
|
||||
}
|
||||
raw = json.dumps(event, ensure_ascii=False, allow_nan=False).encode()
|
||||
if len(raw) > 256 * 1024:
|
||||
raise ValueError("audit envelope exceeds receiver limit")
|
||||
response = await self.client.post(
|
||||
self.base_url + "/v1/events", content=raw, timeout=2,
|
||||
follow_redirects=False,
|
||||
headers={"Authorization": f"Bearer {token}",
|
||||
"Content-Type": "application/json", "Idempotency-Key": event["id"]},
|
||||
)
|
||||
expected = {200: "duplicate", 202: "accepted"}.get(response.status_code)
|
||||
receipt = response.json()
|
||||
if (expected is None or receipt.get("status") != expected
|
||||
or not isinstance(receipt.get("reference"), str) or not receipt["reference"]):
|
||||
raise ValueError("audit custody not acknowledged")
|
||||
except Exception as exc:
|
||||
# Never return receiver bodies, credential values or private paths.
|
||||
raise AccessFailure(503, "audit_unavailable") from exc
|
||||
|
|
@ -67,6 +67,7 @@ class Decision:
|
|||
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:
|
||||
|
|
@ -146,6 +147,7 @@ class AccessController:
|
|||
"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:
|
||||
|
|
@ -196,6 +198,8 @@ class AccessBoundary:
|
|||
return
|
||||
correlation = str(uuid4())
|
||||
context = None
|
||||
action = None
|
||||
digest = None
|
||||
try:
|
||||
headers = Request(scope).headers.getlist("authorization")
|
||||
if len(headers) != 1 or not headers[0].startswith("Bearer "):
|
||||
|
|
@ -250,6 +254,11 @@ class AccessBoundary:
|
|||
"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")
|
||||
|
|
|
|||
62
hub_core/security/config.py
Normal file
62
hub_core/security/config.py
Normal file
|
|
@ -0,0 +1,62 @@
|
|||
"""Explicit runtime composition: owner facts remain an injected trust adapter."""
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
from dataclasses import dataclass
|
||||
from pathlib import Path
|
||||
|
||||
import httpx
|
||||
|
||||
from hub_core.security.audit import AuditCoreSink
|
||||
from hub_core.security.boundary import AccessController, FactSource
|
||||
from hub_core.security.identity import OIDCVerifier, require_https
|
||||
from hub_core.security.policy import FlexPolicy
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class SecuritySettings:
|
||||
issuer: str
|
||||
audience: str
|
||||
root_subject: str
|
||||
policy_url: str
|
||||
policy_caller: str
|
||||
policy_token_file: Path
|
||||
policy_keys_file: Path
|
||||
audit_url: str
|
||||
audit_token_file: Path
|
||||
|
||||
def __post_init__(self):
|
||||
for url in (self.issuer, self.policy_url, self.audit_url):
|
||||
require_https(url)
|
||||
if not self.audience or not self.root_subject:
|
||||
raise ValueError("explicit audience and immutable root subject required")
|
||||
if not self.policy_caller.startswith("system:serviceaccount:"):
|
||||
raise ValueError("explicit policy workload principal required")
|
||||
for path in (self.policy_token_file, self.policy_keys_file, self.audit_token_file):
|
||||
if not isinstance(path, Path) or not path.is_absolute():
|
||||
raise ValueError("absolute credential/trust paths required")
|
||||
if self.policy_token_file == self.audit_token_file:
|
||||
raise ValueError("policy and audit require separate credentials")
|
||||
|
||||
@classmethod
|
||||
def from_env(cls) -> SecuritySettings | None:
|
||||
fields = tuple(cls.__dataclass_fields__)
|
||||
values = {field: os.getenv("HUB_CORE_SECURITY_" + field.upper(), "") for field in fields}
|
||||
if not any(values.values()):
|
||||
return None
|
||||
missing = [field for field, value in values.items() if not value]
|
||||
if missing:
|
||||
raise ValueError("incomplete Hub security configuration: " + ", ".join(missing))
|
||||
return cls(**{field: Path(value) if field.endswith("_file") else value
|
||||
for field, value in values.items()})
|
||||
|
||||
def compose(self, *, facts: FactSource, client: httpx.AsyncClient) -> AccessController:
|
||||
return AccessController(
|
||||
identity=OIDCVerifier(issuer=self.issuer, audience=self.audience, client=client),
|
||||
facts=facts,
|
||||
policy=FlexPolicy(base_url=self.policy_url, client=client,
|
||||
caller_token_file=self.policy_token_file,
|
||||
trusted_keys_file=self.policy_keys_file, caller=self.policy_caller),
|
||||
audit=AuditCoreSink(base_url=self.audit_url, token_file=self.audit_token_file, client=client),
|
||||
root_issuer=self.issuer, root_subject=self.root_subject,
|
||||
)
|
||||
|
|
@ -92,6 +92,20 @@ def _time(value: str) -> datetime:
|
|||
def verify_decision(envelope: dict, *, request: dict, keys: dict,
|
||||
caller: str, now: datetime | None = None) -> Decision:
|
||||
verify_signature(envelope, keys)
|
||||
# The exact signed artifact is retained for independent verification. Do
|
||||
# not hide secret-shaped fields in its serialized audit representation.
|
||||
def check_fields(value):
|
||||
if isinstance(value, dict):
|
||||
for key, item in value.items():
|
||||
if any(fragment in key.lower() for fragment in (
|
||||
"password", "secret", "token", "credential", "private_key",
|
||||
)):
|
||||
raise ValueError("sensitive decision field is outside the audit profile")
|
||||
check_fields(item)
|
||||
elif isinstance(value, list):
|
||||
for item in value:
|
||||
check_fields(item)
|
||||
check_fields(envelope)
|
||||
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"]):
|
||||
|
|
@ -132,7 +146,8 @@ def verify_decision(envelope: dict, *, request: dict, keys: dict,
|
|||
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)
|
||||
return Decision(envelope["effect"] == "allow", envelope["id"], provenance["policy_version"],
|
||||
caller, go_json(envelope).decode())
|
||||
|
||||
|
||||
class FlexPolicy:
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue