Harden production authorization and service auth
Assistant: codex Assistant-Model: gpt-5.6-sol Assistant-Session: 01a0217e-8c4c-7383-be6b-f50a6e485306
This commit is contained in:
parent
f579f3761c
commit
70371649af
20 changed files with 1268 additions and 54 deletions
|
|
@ -12,7 +12,7 @@ from typing import Any
|
|||
|
||||
_LABEL = re.compile(r"^[a-z][a-z0-9-]{0,79}$")
|
||||
_DECISION_REF = re.compile(r"^[A-Z][A-Z0-9-]{2,80}$")
|
||||
_DELIVERY_RESULTS = {"delivered", "failed", "skipped-no-topic"}
|
||||
_DELIVERY_RESULTS = {"delivered", "failed", "queued", "skipped-no-topic"}
|
||||
_VERIFY_RESULT = re.compile(r"^(positive|negative):(pass|fail)$")
|
||||
_ERROR_RESULT = re.compile(
|
||||
r"^failed-(Catalog|Decision|PolicyGuard|Backend|Provisioning|Verification|Delivery)Error$"
|
||||
|
|
|
|||
314
src/secrets_engine/authorization.py
Normal file
314
src/secrets_engine/authorization.py
Normal file
|
|
@ -0,0 +1,314 @@
|
|||
"""Fail-closed consumer validation for flex-auth action authorizations.
|
||||
|
||||
The canonical contract is flex-auth revision c473f19. State Hub does not yet
|
||||
provide the durable authoritative endpoint, so this module validates supplied
|
||||
objects but does not resolve or enable production actions by itself.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
import json
|
||||
import uuid
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime, timezone
|
||||
from typing import Any
|
||||
|
||||
from secrets_engine.catalog import CatalogEntry
|
||||
from secrets_engine.errors import DecisionError
|
||||
|
||||
SCHEMA_VERSION = "0.1"
|
||||
AUTHORITY = "state-hub"
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class ValidatedActionAuthorization:
|
||||
authorization_id: str
|
||||
decision_id: str
|
||||
action: str
|
||||
subject_id: str
|
||||
expires_at: str
|
||||
|
||||
|
||||
def build_action_request(
|
||||
entry: CatalogEntry,
|
||||
action: str,
|
||||
*,
|
||||
subject_id: str,
|
||||
subject_type: str,
|
||||
purpose: str,
|
||||
fields: list[str] | tuple[str, ...] = (),
|
||||
policy_targets: list[str] | tuple[str, ...] = (),
|
||||
auth_targets: list[str] | tuple[str, ...] = (),
|
||||
request_id: str = "",
|
||||
) -> dict[str, Any]:
|
||||
"""Build the exact normalized secrets-engine profile for flex-auth."""
|
||||
if not action or not subject_id or not subject_type or not purpose:
|
||||
raise DecisionError(
|
||||
"action request requires action, subject id/type, and purpose"
|
||||
)
|
||||
request: dict[str, Any] = {}
|
||||
if request_id:
|
||||
request["id"] = request_id
|
||||
request.update(
|
||||
{
|
||||
"subject": {"id": subject_id, "type": subject_type},
|
||||
"action": action,
|
||||
"resource": {
|
||||
"id": f"catalog:{entry.id}",
|
||||
"type": "secret-catalog-lane",
|
||||
"system": "secrets-engine",
|
||||
"attributes": {
|
||||
"stage": entry.stage,
|
||||
"fields": sorted(set(fields)),
|
||||
"policy_targets": sorted(set(policy_targets)),
|
||||
"auth_targets": sorted(set(auth_targets)),
|
||||
},
|
||||
},
|
||||
"context": {"purpose": purpose},
|
||||
}
|
||||
)
|
||||
return request
|
||||
|
||||
|
||||
def _required_dict(container: dict[str, Any], name: str) -> dict[str, Any]:
|
||||
value = container.get(name)
|
||||
if not isinstance(value, dict):
|
||||
raise DecisionError(f"action authorization requires object '{name}'")
|
||||
return value
|
||||
|
||||
|
||||
def _required_text(container: dict[str, Any], name: str) -> str:
|
||||
value = container.get(name)
|
||||
if not isinstance(value, str) or not value:
|
||||
raise DecisionError(f"action authorization requires non-empty '{name}'")
|
||||
return value
|
||||
|
||||
|
||||
def _parse_time(value: object, name: str) -> datetime:
|
||||
if not isinstance(value, str):
|
||||
raise DecisionError(f"action authorization requires timestamp '{name}'")
|
||||
try:
|
||||
parsed = datetime.fromisoformat(value.replace("Z", "+00:00"))
|
||||
except ValueError as e:
|
||||
raise DecisionError(f"action authorization has invalid timestamp '{name}'") from e
|
||||
if parsed.tzinfo is None:
|
||||
raise DecisionError(f"action authorization timestamp '{name}' needs timezone")
|
||||
return parsed.astimezone(timezone.utc)
|
||||
|
||||
|
||||
def _sorted_map(value: object) -> dict[str, Any]:
|
||||
if not isinstance(value, dict):
|
||||
return {}
|
||||
return {key: _canonical_map_value(value[key]) for key in sorted(value)}
|
||||
|
||||
|
||||
def _canonical_map_value(value: Any) -> Any:
|
||||
if isinstance(value, dict):
|
||||
return _sorted_map(value)
|
||||
if isinstance(value, list):
|
||||
return [_canonical_map_value(item) for item in value]
|
||||
return value
|
||||
|
||||
|
||||
def _subject_ref(value: object) -> dict[str, Any]:
|
||||
if not isinstance(value, dict):
|
||||
raise DecisionError("action authorization subject must be an object")
|
||||
subject: dict[str, Any] = {"id": _required_text(value, "id")}
|
||||
for name in ("type", "tenant"):
|
||||
if value.get(name):
|
||||
subject[name] = _required_text(value, name)
|
||||
if value.get("attributes") is not None:
|
||||
subject["attributes"] = _sorted_map(value.get("attributes"))
|
||||
return subject
|
||||
|
||||
|
||||
def _resource_ref(value: object) -> dict[str, Any]:
|
||||
if not isinstance(value, dict):
|
||||
raise DecisionError("action authorization resource must be an object")
|
||||
resource: dict[str, Any] = {"id": _required_text(value, "id")}
|
||||
for name in ("type", "system", "tenant"):
|
||||
if value.get(name):
|
||||
resource[name] = _required_text(value, name)
|
||||
if value.get("attributes") is not None:
|
||||
resource["attributes"] = _sorted_map(value.get("attributes"))
|
||||
return resource
|
||||
|
||||
|
||||
def canonical_check_request(request: object) -> dict[str, Any]:
|
||||
"""Match Go encoding/json field order used by flex-auth request digests."""
|
||||
if not isinstance(request, dict):
|
||||
raise DecisionError("action authorization request must be an object")
|
||||
canonical: dict[str, Any] = {}
|
||||
if request.get("id"):
|
||||
canonical["id"] = _required_text(request, "id")
|
||||
if request.get("tenant"):
|
||||
canonical["tenant"] = _required_text(request, "tenant")
|
||||
canonical["subject"] = _subject_ref(request.get("subject"))
|
||||
canonical["action"] = _required_text(request, "action")
|
||||
canonical["resource"] = _resource_ref(request.get("resource"))
|
||||
if request.get("context") is not None:
|
||||
canonical["context"] = _sorted_map(request.get("context"))
|
||||
if request.get("caring_context") is not None:
|
||||
canonical["caring_context"] = _canonical_map_value(
|
||||
request.get("caring_context")
|
||||
)
|
||||
if request.get("policy_version"):
|
||||
canonical["policy_version"] = _required_text(request, "policy_version")
|
||||
return canonical
|
||||
|
||||
|
||||
def request_digest(request: object) -> str:
|
||||
canonical = canonical_check_request(request)
|
||||
encoded = json.dumps(
|
||||
canonical, ensure_ascii=False, separators=(",", ":")
|
||||
).encode("utf-8")
|
||||
return "sha256:" + hashlib.sha256(encoded).hexdigest()
|
||||
|
||||
|
||||
def _require_exact_target_sets(request: dict[str, Any]) -> None:
|
||||
resource = _required_dict(request, "resource")
|
||||
attributes = resource.get("attributes", {})
|
||||
if not isinstance(attributes, dict):
|
||||
raise DecisionError("action authorization resource attributes must be an object")
|
||||
for name in ("fields", "policy_targets", "auth_targets"):
|
||||
values = attributes.get(name, [])
|
||||
if not isinstance(values, list) or not all(
|
||||
isinstance(item, str) and item for item in values
|
||||
):
|
||||
raise DecisionError(f"action authorization target set '{name}' is invalid")
|
||||
if values != sorted(set(values)):
|
||||
raise DecisionError(
|
||||
f"action authorization target set '{name}' must be sorted and unique"
|
||||
)
|
||||
|
||||
|
||||
def validate_action_authorization(
|
||||
envelope: object,
|
||||
expected_request: object,
|
||||
*,
|
||||
accepted_policy_packages: set[str],
|
||||
accepted_policy_versions: set[str],
|
||||
minimum_approval_count: int = 1,
|
||||
now: datetime | None = None,
|
||||
) -> ValidatedActionAuthorization:
|
||||
"""Validate exact request binding and dual control; never parse prose."""
|
||||
if minimum_approval_count < 1:
|
||||
raise DecisionError("minimum approval count must be positive")
|
||||
if not accepted_policy_packages or not accepted_policy_versions:
|
||||
raise DecisionError("accepted flex-auth policy package/version is required")
|
||||
if not isinstance(envelope, dict):
|
||||
raise DecisionError("action authorization must be an object")
|
||||
if envelope.get("schema_version") != SCHEMA_VERSION:
|
||||
raise DecisionError("unsupported action authorization schema version")
|
||||
authorization_id = _required_text(envelope, "id")
|
||||
try:
|
||||
parsed_authorization_id = uuid.UUID(authorization_id)
|
||||
except ValueError as e:
|
||||
raise DecisionError("action authorization id must be a canonical UUID") from e
|
||||
if str(parsed_authorization_id) != authorization_id:
|
||||
raise DecisionError("action authorization id must be a canonical UUID")
|
||||
if envelope.get("status") != "approved":
|
||||
raise DecisionError("action authorization status is not approved")
|
||||
if envelope.get("superseded_by"):
|
||||
raise DecisionError("action authorization is superseded")
|
||||
provenance = _required_dict(envelope, "provenance")
|
||||
if provenance.get("authority") != AUTHORITY:
|
||||
raise DecisionError("action authorization authority is not State Hub")
|
||||
|
||||
request = canonical_check_request(envelope.get("request"))
|
||||
expected = canonical_check_request(expected_request)
|
||||
_require_exact_target_sets(request)
|
||||
_require_exact_target_sets(expected)
|
||||
if request != expected:
|
||||
raise DecisionError("action authorization request does not exactly match action")
|
||||
|
||||
validity = _required_dict(envelope, "validity")
|
||||
expires = _parse_time(validity.get("expires_at"), "expires_at")
|
||||
not_before = (
|
||||
_parse_time(validity.get("not_before"), "not_before")
|
||||
if validity.get("not_before") is not None
|
||||
else None
|
||||
)
|
||||
current = (now or datetime.now(timezone.utc)).astimezone(timezone.utc)
|
||||
if not_before is not None and current < not_before:
|
||||
raise DecisionError("action authorization window has not started")
|
||||
if current >= expires:
|
||||
raise DecisionError("action authorization has expired")
|
||||
|
||||
approvals = _required_dict(envelope, "approvals")
|
||||
required_count = approvals.get("required_count")
|
||||
entries = approvals.get("entries")
|
||||
if not isinstance(required_count, int) or required_count < 1:
|
||||
raise DecisionError("action authorization approval count is invalid")
|
||||
if required_count < minimum_approval_count:
|
||||
raise DecisionError("action authorization approval threshold is insufficient")
|
||||
if not isinstance(entries, list):
|
||||
raise DecisionError("action authorization approval entries are invalid")
|
||||
approvers: set[str] = set()
|
||||
for entry in entries:
|
||||
if not isinstance(entry, dict):
|
||||
raise DecisionError("action authorization approval entry is invalid")
|
||||
subject_id = _required_text(entry, "subject_id")
|
||||
approved_at = _parse_time(entry.get("approved_at"), "approved_at")
|
||||
if approved_at > current or approved_at >= expires:
|
||||
raise DecisionError("action authorization approval time is outside window")
|
||||
if not_before is not None and approved_at < not_before:
|
||||
raise DecisionError("action authorization approval time is outside window")
|
||||
if subject_id in approvers:
|
||||
raise DecisionError("action authorization contains duplicate approver")
|
||||
approvers.add(subject_id)
|
||||
if len(approvers) < required_count:
|
||||
raise DecisionError("action authorization has insufficient distinct approvals")
|
||||
|
||||
decision = _required_dict(envelope, "decision")
|
||||
if decision.get("effect") != "allow":
|
||||
raise DecisionError("flex-auth decision effect is not allow")
|
||||
decision_id = _required_text(decision, "id")
|
||||
if request.get("id") and decision.get("request_id") != request["id"]:
|
||||
raise DecisionError("flex-auth decision request id does not match request")
|
||||
binding = _required_dict(decision, "binding")
|
||||
bound_request: dict[str, Any] = {}
|
||||
if binding.get("tenant"):
|
||||
bound_request["tenant"] = binding["tenant"]
|
||||
bound_request.update(
|
||||
{
|
||||
"subject": binding.get("subject"),
|
||||
"action": binding.get("action"),
|
||||
"resource": binding.get("resource"),
|
||||
"context": binding.get("context", {}),
|
||||
}
|
||||
)
|
||||
expected_bound: dict[str, Any] = {}
|
||||
if request.get("tenant"):
|
||||
expected_bound["tenant"] = request["tenant"]
|
||||
expected_bound.update(
|
||||
{
|
||||
"subject": request["subject"],
|
||||
"action": request["action"],
|
||||
"resource": request["resource"],
|
||||
"context": request.get("context", {}),
|
||||
}
|
||||
)
|
||||
if canonical_check_request(bound_request) != canonical_check_request(
|
||||
expected_bound
|
||||
):
|
||||
raise DecisionError("flex-auth decision binding does not match request")
|
||||
if binding.get("request_digest") != request_digest(request):
|
||||
raise DecisionError("flex-auth request digest does not match request")
|
||||
if _subject_ref(decision.get("subject")) != request["subject"]:
|
||||
raise DecisionError("flex-auth decision subject does not match request")
|
||||
if _resource_ref(decision.get("resource")) != request["resource"]:
|
||||
raise DecisionError("flex-auth decision resource does not match request")
|
||||
decision_provenance = _required_dict(decision, "provenance")
|
||||
if decision_provenance.get("policy_package") not in accepted_policy_packages:
|
||||
raise DecisionError("flex-auth policy package is not accepted")
|
||||
if decision_provenance.get("policy_version") not in accepted_policy_versions:
|
||||
raise DecisionError("flex-auth policy version is not accepted")
|
||||
|
||||
return ValidatedActionAuthorization(
|
||||
authorization_id=authorization_id,
|
||||
decision_id=decision_id,
|
||||
action=request["action"],
|
||||
subject_id=request["subject"]["id"],
|
||||
expires_at=expires.isoformat(),
|
||||
)
|
||||
|
|
@ -21,8 +21,10 @@ Every privileged action is decision-gated and writes non-secret evidence.
|
|||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import os
|
||||
import sys
|
||||
from pathlib import Path
|
||||
from urllib.parse import urlparse
|
||||
|
||||
from secrets_engine import __version__
|
||||
from secrets_engine.apply import apply_plan
|
||||
|
|
@ -83,8 +85,28 @@ def _privileged_evidence(
|
|||
)
|
||||
|
||||
|
||||
def _require_lane_approval(cfg: Config, entry):
|
||||
"""Resolve and enforce the lane approval for a privileged live action."""
|
||||
def _unsafe_local_demo_enabled(cfg: Config) -> bool:
|
||||
"""Return true only for an explicit, offline, loopback-only demo."""
|
||||
host = (urlparse(cfg.bao_addr).hostname or "").lower()
|
||||
return (
|
||||
os.environ.get("SECRETS_ENGINE_UNSAFE_DEMO") == "1"
|
||||
and not cfg.hub_url
|
||||
and host in {"127.0.0.1", "localhost", "::1"}
|
||||
)
|
||||
|
||||
|
||||
def _require_lane_approval(cfg: Config, entry, action: str = ""):
|
||||
"""Resolve approval for a live action, failing production closed.
|
||||
|
||||
The durable State Hub action-authorization endpoint is not available yet.
|
||||
Production therefore cannot rely on a coarse lane decision. The one narrow
|
||||
exception is an explicit offline demo against a loopback OpenBao instance.
|
||||
"""
|
||||
if entry.stage == "prod" and not _unsafe_local_demo_enabled(cfg):
|
||||
raise DecisionError(
|
||||
f"production action '{action or 'unknown'}' requires a durable "
|
||||
"State Hub action authorization; live production remains disabled"
|
||||
)
|
||||
if not entry.approval_required():
|
||||
return None
|
||||
decision = resolve_decision(
|
||||
|
|
@ -207,7 +229,7 @@ def cmd_apply(cfg: Config, args) -> int:
|
|||
return 0
|
||||
|
||||
with _privileged_evidence(cfg, entry, "apply") as evidence:
|
||||
decision = _require_lane_approval(cfg, entry)
|
||||
decision = _require_lane_approval(cfg, entry, "apply")
|
||||
evidence.mark_approved(decision)
|
||||
plan = build_plan(
|
||||
entry, args.stage, decision_id=decision.id if decision else ""
|
||||
|
|
@ -236,7 +258,7 @@ def cmd_provision(cfg: Config, args) -> int:
|
|||
raise ProvisioningError(
|
||||
f"lane '{entry.id}' is stage '{entry.stage}', not '{args.stage}'"
|
||||
)
|
||||
decision = _require_lane_approval(cfg, entry)
|
||||
decision = _require_lane_approval(cfg, entry, "provision")
|
||||
evidence.mark_approved(decision)
|
||||
client = OpenBaoClient.resolve(
|
||||
cfg.bao_addr, bootstrap_token_file=args.bootstrap_token_file
|
||||
|
|
@ -269,7 +291,7 @@ def cmd_verify(cfg: Config, args) -> int:
|
|||
"negative_requested": negative,
|
||||
},
|
||||
) as evidence:
|
||||
decision = _require_lane_approval(cfg, entry)
|
||||
decision = _require_lane_approval(cfg, entry, "verify")
|
||||
evidence.mark_approved(decision)
|
||||
client = OpenBaoClient.resolve(
|
||||
cfg.bao_addr, bootstrap_token_file=args.bootstrap_token_file
|
||||
|
|
@ -345,7 +367,7 @@ def cmd_handoff(cfg: Config, args) -> int:
|
|||
raise ProvisioningError(
|
||||
f"lane '{entry.id}' is {entry.kind}; handoff needs auth-capability"
|
||||
)
|
||||
decision = _require_lane_approval(cfg, entry)
|
||||
decision = _require_lane_approval(cfg, entry, "handoff")
|
||||
evidence.mark_approved(decision)
|
||||
client = OpenBaoClient.resolve(
|
||||
cfg.bao_addr, bootstrap_token_file=args.bootstrap_token_file
|
||||
|
|
@ -397,7 +419,7 @@ def cmd_exec(cfg: Config, args) -> int:
|
|||
},
|
||||
) as evidence:
|
||||
# require approval + readiness before running.
|
||||
decision = _require_lane_approval(cfg, entry)
|
||||
decision = _require_lane_approval(cfg, entry, "exec")
|
||||
evidence.mark_approved(decision)
|
||||
if not args.command:
|
||||
from secrets_engine.errors import DeliveryError
|
||||
|
|
@ -476,7 +498,7 @@ def cmd_revoke(cfg: Config, args) -> int:
|
|||
with _privileged_evidence(
|
||||
cfg, entry, "revoke", detail={"operation": plan.operation}
|
||||
) as evidence:
|
||||
decision = _require_lane_approval(cfg, entry)
|
||||
decision = _require_lane_approval(cfg, entry, "deactivate")
|
||||
evidence.mark_approved(decision)
|
||||
client = OpenBaoClient.resolve(
|
||||
cfg.bao_addr, bootstrap_token_file=args.bootstrap_token_file
|
||||
|
|
@ -524,7 +546,7 @@ def cmd_lifecycle(cfg: Config, args) -> int:
|
|||
"live destroy is disabled until an exact-action destruction "
|
||||
"approval contract is available; use --dry-run to inspect targets"
|
||||
)
|
||||
decision = _require_lane_approval(cfg, entry)
|
||||
decision = _require_lane_approval(cfg, entry, args.operation)
|
||||
evidence.mark_approved(decision)
|
||||
client = OpenBaoClient.resolve(
|
||||
cfg.bao_addr, bootstrap_token_file=args.bootstrap_token_file
|
||||
|
|
|
|||
|
|
@ -47,6 +47,7 @@ class EvidenceWriter:
|
|||
topic_id: str = ""
|
||||
workstream_id: str = ""
|
||||
author: str = "secrets-engine"
|
||||
repo_slug: str = "secrets-engine"
|
||||
actor: str = field(default_factory=lambda: os.environ.get("USER", "unknown"))
|
||||
|
||||
def __post_init__(self) -> None:
|
||||
|
|
@ -91,8 +92,13 @@ class EvidenceWriter:
|
|||
}
|
||||
self._append_local(record)
|
||||
if hub_requested:
|
||||
delivery_result = self._post_hub(
|
||||
action, result, catalog_id, stage, decision_id
|
||||
delivery = self._post_hub(
|
||||
action,
|
||||
result,
|
||||
catalog_id,
|
||||
stage,
|
||||
decision_id,
|
||||
record_id=record_id,
|
||||
)
|
||||
# Append-only companion evidence makes an unavailable State Hub
|
||||
# visible without rewriting or delaying the primary local record.
|
||||
|
|
@ -102,23 +108,32 @@ class EvidenceWriter:
|
|||
"related_record_id": record_id,
|
||||
"ts": datetime.now(timezone.utc).isoformat(),
|
||||
"action": "evidence-delivery",
|
||||
"result": delivery_result,
|
||||
"result": delivery.status,
|
||||
"actor": self.actor,
|
||||
"catalog_id": catalog_id,
|
||||
"stage": stage,
|
||||
"decision_id": decision_id,
|
||||
"detail": {},
|
||||
"detail": {"outbox_id": delivery.outbox_id}
|
||||
if delivery.outbox_id
|
||||
else {},
|
||||
"hub_delivery_requested": False,
|
||||
}
|
||||
)
|
||||
return record
|
||||
|
||||
def _post_hub(
|
||||
self, action: str, result: str, catalog_id: str, stage: str, decision_id: str
|
||||
) -> str:
|
||||
self,
|
||||
action: str,
|
||||
result: str,
|
||||
catalog_id: str,
|
||||
stage: str,
|
||||
decision_id: str,
|
||||
*,
|
||||
record_id: str,
|
||||
) -> "HubDelivery":
|
||||
"""Best-effort progress note; return a non-secret delivery outcome."""
|
||||
if not self.topic_id:
|
||||
return "skipped-no-topic"
|
||||
return HubDelivery("skipped-no-topic")
|
||||
summary = f"secrets-engine {action}: {result}"
|
||||
if catalog_id:
|
||||
summary += f" [{catalog_id}{'/' + stage if stage else ''}]"
|
||||
|
|
@ -133,17 +148,47 @@ class EvidenceWriter:
|
|||
if decision_id:
|
||||
payload["detail"] = {"decision_id": decision_id, "catalog_id": catalog_id}
|
||||
try:
|
||||
idempotency_key = f"secrets-engine:{record_id}"
|
||||
req = urllib.request.Request(
|
||||
self.hub_url.rstrip("/") + "/progress/",
|
||||
data=json.dumps(payload).encode(),
|
||||
headers={"Content-Type": "application/json"},
|
||||
headers={
|
||||
"Content-Type": "application/json",
|
||||
"Idempotency-Key": idempotency_key,
|
||||
"X-StateHub-Source-Agent": self.author,
|
||||
"X-StateHub-Repo-Slug": self.repo_slug,
|
||||
},
|
||||
method="POST",
|
||||
)
|
||||
urllib.request.urlopen(req, timeout=3).read()
|
||||
return "delivered"
|
||||
response = urllib.request.urlopen(req, timeout=3)
|
||||
body = response.read()
|
||||
status = getattr(response, "status", 200)
|
||||
if status == 202:
|
||||
try:
|
||||
receipt = json.loads(body or b"{}")
|
||||
except (json.JSONDecodeError, TypeError):
|
||||
receipt = {}
|
||||
if isinstance(receipt, dict) and receipt.get("queued") is True:
|
||||
outbox_id = receipt.get("outbox_id", "")
|
||||
if isinstance(outbox_id, str):
|
||||
try:
|
||||
if str(uuid.UUID(outbox_id)) == outbox_id:
|
||||
return HubDelivery("queued", outbox_id=outbox_id)
|
||||
except ValueError:
|
||||
pass
|
||||
return HubDelivery("failed")
|
||||
return HubDelivery("delivered")
|
||||
except (urllib.error.URLError, OSError, ValueError):
|
||||
# Hub being offline must never block secret work or leak anything.
|
||||
return "failed"
|
||||
return HubDelivery("failed")
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class HubDelivery:
|
||||
"""Non-secret outcome returned by State Hub or its edge relay."""
|
||||
|
||||
status: str
|
||||
outbox_id: str = ""
|
||||
|
||||
|
||||
@dataclass
|
||||
|
|
|
|||
217
src/secrets_engine/service_auth.py
Normal file
217
src/secrets_engine/service_auth.py
Normal file
|
|
@ -0,0 +1,217 @@
|
|||
"""Explicit KeyCape service-JWT provider for future OpenBao JWT login.
|
||||
|
||||
This module implements the accepted KeyCape consumer contract without wiring it
|
||||
into OpenBao yet. The platform owner still needs to publish the exact OpenBao
|
||||
JWT auth mount and role. Keeping provider selection separate prevents an auth
|
||||
failure from falling back to bootstrap, operator, or AppRole credentials.
|
||||
|
||||
JWT parsing here is a claim preflight, not signature verification. OpenBao must
|
||||
verify the RS256 signature against the configured issuer before issuing a token.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import base64
|
||||
import binascii
|
||||
import json
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
from typing import Any, Callable
|
||||
from urllib.error import HTTPError, URLError
|
||||
from urllib.parse import urlencode
|
||||
from urllib.request import Request, urlopen
|
||||
|
||||
from secrets_engine.errors import BackendError
|
||||
from secrets_engine.openbao import read_strict_token_file
|
||||
|
||||
CLIENT_ID = "secrets-engine-openbao"
|
||||
SUBJECT = "service:secrets-engine"
|
||||
PRINCIPAL_TYPE = "service"
|
||||
TENANT = "tenant:coulomb"
|
||||
ROLE = "secrets-engine"
|
||||
SCOPE = "openbao:login"
|
||||
MAX_TOKEN_SECONDS = 15 * 60
|
||||
RENEW_WINDOW_SECONDS = 3 * 60
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class KeyCapeServiceAuthConfig:
|
||||
token_url: str
|
||||
issuer: str
|
||||
client_secret_file: Path
|
||||
client_id: str = CLIENT_ID
|
||||
subject: str = SUBJECT
|
||||
audience: str = CLIENT_ID
|
||||
tenant: str = TENANT
|
||||
required_role: str = ROLE
|
||||
scope: str = SCOPE
|
||||
timeout_seconds: float = 10.0
|
||||
|
||||
def __post_init__(self) -> None:
|
||||
if not self.token_url.startswith("https://"):
|
||||
raise BackendError("KeyCape token URL must use HTTPS")
|
||||
if not self.issuer.startswith("https://"):
|
||||
raise BackendError("KeyCape issuer must use HTTPS")
|
||||
if self.client_id != CLIENT_ID or self.audience != CLIENT_ID:
|
||||
raise BackendError("KeyCape service client/audience must match accepted contract")
|
||||
if (
|
||||
self.subject != SUBJECT
|
||||
or self.tenant != TENANT
|
||||
or self.required_role != ROLE
|
||||
or self.scope != SCOPE
|
||||
):
|
||||
raise BackendError("KeyCape service identity claims must match accepted contract")
|
||||
if self.timeout_seconds <= 0:
|
||||
raise BackendError("KeyCape timeout must be positive")
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class ServiceJWT:
|
||||
token: str = field(repr=False)
|
||||
issued_at: int
|
||||
expires_at: int
|
||||
claims: dict[str, Any] = field(repr=False)
|
||||
|
||||
def needs_renewal(self, now: datetime | None = None) -> bool:
|
||||
current = int((now or datetime.now(timezone.utc)).timestamp())
|
||||
return self.expires_at - current <= RENEW_WINDOW_SECONDS
|
||||
|
||||
|
||||
def _decode_segment(value: str, label: str) -> dict[str, Any]:
|
||||
try:
|
||||
padded = value + "=" * (-len(value) % 4)
|
||||
decoded = base64.urlsafe_b64decode(padded.encode("ascii"))
|
||||
parsed = json.loads(decoded)
|
||||
except (UnicodeEncodeError, binascii.Error, json.JSONDecodeError) as e:
|
||||
raise BackendError(f"KeyCape JWT has invalid {label}") from e
|
||||
if not isinstance(parsed, dict):
|
||||
raise BackendError(f"KeyCape JWT {label} must be an object")
|
||||
return parsed
|
||||
|
||||
|
||||
def preflight_service_jwt(
|
||||
token: str,
|
||||
config: KeyCapeServiceAuthConfig,
|
||||
*,
|
||||
now: datetime | None = None,
|
||||
) -> ServiceJWT:
|
||||
"""Validate non-cryptographic JWT shape/claims before OpenBao login."""
|
||||
parts = token.split(".")
|
||||
if len(parts) != 3 or not all(parts):
|
||||
raise BackendError("KeyCape access token is not a compact JWT")
|
||||
header = _decode_segment(parts[0], "header")
|
||||
claims = _decode_segment(parts[1], "payload")
|
||||
if header.get("alg") != "RS256":
|
||||
raise BackendError("KeyCape JWT algorithm is not RS256")
|
||||
|
||||
exact_claims = {
|
||||
"iss": config.issuer,
|
||||
"sub": config.subject,
|
||||
"aud": config.audience,
|
||||
"principal_type": PRINCIPAL_TYPE,
|
||||
"tenant": config.tenant,
|
||||
"groups": [],
|
||||
}
|
||||
for name, expected in exact_claims.items():
|
||||
if claims.get(name) != expected:
|
||||
raise BackendError(f"KeyCape JWT claim '{name}' does not match contract")
|
||||
if claims.get("roles") != [config.required_role]:
|
||||
raise BackendError("KeyCape JWT roles do not match contract")
|
||||
if claims.get("scope") != config.scope:
|
||||
raise BackendError("KeyCape JWT scope does not match contract")
|
||||
|
||||
assurance = claims.get("assurance")
|
||||
if not isinstance(assurance, dict) or (
|
||||
assurance.get("aal") != "AAL1"
|
||||
or assurance.get("method") != "client_secret"
|
||||
or assurance.get("mfa") is not False
|
||||
or assurance.get("source") != "key-cape"
|
||||
):
|
||||
raise BackendError("KeyCape JWT assurance does not match contract")
|
||||
|
||||
issued_at = claims.get("iat")
|
||||
expires_at = claims.get("exp")
|
||||
if (
|
||||
not isinstance(issued_at, int)
|
||||
or isinstance(issued_at, bool)
|
||||
or not isinstance(expires_at, int)
|
||||
or isinstance(expires_at, bool)
|
||||
):
|
||||
raise BackendError("KeyCape JWT iat/exp must be integer timestamps")
|
||||
current = int((now or datetime.now(timezone.utc)).timestamp())
|
||||
if issued_at > current + 60:
|
||||
raise BackendError("KeyCape JWT issue time is in the future")
|
||||
if expires_at <= current:
|
||||
raise BackendError("KeyCape JWT has expired")
|
||||
if expires_at <= issued_at or expires_at - issued_at > MAX_TOKEN_SECONDS:
|
||||
raise BackendError("KeyCape JWT lifetime exceeds accepted 15-minute bound")
|
||||
return ServiceJWT(
|
||||
token=token,
|
||||
issued_at=issued_at,
|
||||
expires_at=expires_at,
|
||||
claims=dict(claims),
|
||||
)
|
||||
|
||||
|
||||
Transport = Callable[..., Any]
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class KeyCapeServiceAuthProvider:
|
||||
config: KeyCapeServiceAuthConfig
|
||||
transport: Transport = field(default=urlopen, repr=False, compare=False)
|
||||
|
||||
def exchange(self, *, now: datetime | None = None) -> ServiceJWT:
|
||||
"""Perform one client-credentials exchange; never retry or fall back."""
|
||||
secret = read_strict_token_file(
|
||||
self.config.client_secret_file,
|
||||
purpose="KeyCape client secret",
|
||||
)
|
||||
basic = base64.b64encode(
|
||||
f"{self.config.client_id}:{secret}".encode("utf-8")
|
||||
).decode("ascii")
|
||||
body = urlencode(
|
||||
{"grant_type": "client_credentials", "scope": self.config.scope}
|
||||
).encode("ascii")
|
||||
request = Request(
|
||||
self.config.token_url,
|
||||
data=body,
|
||||
method="POST",
|
||||
headers={
|
||||
"Authorization": f"Basic {basic}",
|
||||
"Content-Type": "application/x-www-form-urlencoded",
|
||||
"Accept": "application/json",
|
||||
},
|
||||
)
|
||||
try:
|
||||
response = self.transport(request, timeout=self.config.timeout_seconds)
|
||||
with response:
|
||||
status = getattr(response, "status", 200)
|
||||
raw = response.read()
|
||||
except HTTPError as e:
|
||||
raise BackendError(f"KeyCape token exchange failed with HTTP {e.code}") from e
|
||||
except (URLError, TimeoutError, OSError) as e:
|
||||
raise BackendError("KeyCape token exchange failed") from e
|
||||
if status != 200:
|
||||
raise BackendError(f"KeyCape token exchange failed with HTTP {status}")
|
||||
try:
|
||||
payload = json.loads(raw)
|
||||
except (UnicodeDecodeError, json.JSONDecodeError) as e:
|
||||
raise BackendError("KeyCape token response is not valid JSON") from e
|
||||
if not isinstance(payload, dict):
|
||||
raise BackendError("KeyCape token response must be an object")
|
||||
if payload.get("id_token") or payload.get("refresh_token"):
|
||||
raise BackendError("KeyCape service exchange returned a forbidden extra token")
|
||||
token = payload.get("access_token")
|
||||
if not isinstance(token, str) or not token:
|
||||
raise BackendError("KeyCape token response has no access token")
|
||||
if str(payload.get("token_type", "")).lower() != "bearer":
|
||||
raise BackendError("KeyCape token response type is not Bearer")
|
||||
expires_in = payload.get("expires_in")
|
||||
if (
|
||||
not isinstance(expires_in, int)
|
||||
or isinstance(expires_in, bool)
|
||||
or not 0 < expires_in <= MAX_TOKEN_SECONDS
|
||||
):
|
||||
raise BackendError("KeyCape token response lifetime is outside contract")
|
||||
return preflight_service_jwt(token, self.config, now=now)
|
||||
Loading…
Add table
Add a link
Reference in a new issue