railiance-platform/scripts/custody-projection.py

949 lines
34 KiB
Python
Raw Normal View History

#!/usr/bin/env python3
"""Data-driven, attended ephemeral credential custody lifecycle.
The command never prints or persists bearer values. Live projection requires a
current Whitehat broker-readiness receipt in State Hub and an exact attended
confirmation. Cleanup is scoped only by the validated projection contract and
receipt.
"""
from __future__ import annotations
import argparse
import json
import os
import secrets
import stat
import subprocess
import sys
import tempfile
import time
from datetime import UTC, datetime
from pathlib import Path
from typing import Any, Callable
SCRIPT_DIR = Path(__file__).resolve().parent
if str(SCRIPT_DIR) not in sys.path:
sys.path.insert(0, str(SCRIPT_DIR))
from custody_contract import ( # noqa: E402
CLEANUP_INTERFACE,
PROJECTION_INTERFACE,
WORKPLAN_ID,
ContractError,
canonical_json,
contract_digest,
current_broker_receipt,
digest,
load_json,
parse_time,
resource_names,
validate_cleanup_receipt,
validate_projection_contract,
validate_projection_receipt,
)
from remote_exec import RemoteExecutionError, run_remote # noqa: E402
DEFAULT_API_BASE = os.environ.get("STATE_HUB_URL", "http://127.0.0.1:8000")
class ProcedureError(RuntimeError):
pass
def utc_now() -> datetime:
return datetime.now(UTC)
def rfc3339(value: datetime) -> str:
return value.astimezone(UTC).replace(microsecond=0).isoformat().replace("+00:00", "Z")
def safe_run(
command: list[str],
*,
label: str,
env: dict[str, str] | None = None,
input_text: str | None = None,
allow_missing: bool = False,
runner: Callable[..., subprocess.CompletedProcess[str]] = subprocess.run,
) -> subprocess.CompletedProcess[str]:
completed = runner(
command,
text=True,
input=input_text,
capture_output=True,
env=env,
check=False,
)
if completed.returncode and not allow_missing:
raise ProcedureError(f"{label} failed (exit {completed.returncode})")
return completed
def secure_unlink(path: Path) -> None:
if not path.exists():
return
size = path.stat().st_size
with path.open("r+b", buffering=0) as handle:
handle.write(b"\0" * size)
handle.flush()
os.fsync(handle.fileno())
path.unlink()
def write_receipt(path: Path, value: dict[str, Any]) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
descriptor, temporary_name = tempfile.mkstemp(prefix=f".{path.name}.", dir=path.parent)
temporary = Path(temporary_name)
try:
os.fchmod(descriptor, 0o600)
with os.fdopen(descriptor, "w", encoding="utf-8") as handle:
handle.write(json.dumps(value, indent=2, sort_keys=True) + "\n")
handle.flush()
os.fsync(handle.fileno())
os.replace(temporary, path)
os.chmod(path, 0o600)
finally:
if temporary.exists():
temporary.unlink()
class Operator:
def __init__(self, contract: dict[str, Any], token_file: Path) -> None:
mode = stat.S_IMODE(token_file.stat().st_mode) if token_file.exists() else None
if mode != 0o600:
raise ProcedureError(f"OpenBao token file must exist with mode 0600: {token_file}")
token = token_file.read_text(encoding="utf-8").splitlines()[0].strip()
if not token:
raise ProcedureError("OpenBao token file is empty")
self.contract = contract
self.remote = contract["authority"]["remote"]
self.bao_env = dict(
os.environ,
BAO_ADDR=os.environ.get("BAO_ADDR", "https://bao.coulomb.social"),
BAO_TOKEN=token,
)
def bao(
self,
args: list[str],
*,
label: str,
input_text: str | None = None,
allow_missing: bool = False,
) -> subprocess.CompletedProcess[str]:
return safe_run(
["bao", *args],
label=label,
env=self.bao_env,
input_text=input_text,
allow_missing=allow_missing,
)
def kubectl(
self,
args: list[str],
*,
label: str,
input_text: str | None = None,
allow_missing: bool = False,
) -> subprocess.CompletedProcess[str]:
try:
return run_remote(
self.remote,
["kubectl", *args],
label=label,
input_text=input_text,
allow_missing=allow_missing,
)
except RemoteExecutionError as exc:
raise ProcedureError(str(exc)) from exc
def bao_exists(self, args: list[str], *, label: str) -> bool:
result = self.bao(args, label=label, allow_missing=True)
if result.returncode == 0:
return True
error = result.stderr.lower()
if "404" in error or "no value found" in error or "no policy named" in error:
return False
raise ProcedureError(f"{label} failed without a confirmed not-found response")
def kubectl_exists(self, args: list[str], *, label: str) -> bool:
result = self.kubectl(args, label=label, allow_missing=True)
return result.returncode == 0 and bool(result.stdout.strip())
def registry(self) -> list[dict[str, Any]]:
authority = self.contract["authority"]
result = self.bao(
["kv", "get", f"-field={authority['registry_field']}", authority["registry_path"]],
label="read sender registry",
)
try:
value = json.loads(result.stdout)
except json.JSONDecodeError as exc:
raise ProcedureError("sender registry is not valid JSON") from exc
if not isinstance(value, list):
raise ProcedureError("sender registry must be a list")
return value
def write_registry(self, registry: list[dict[str, Any]]) -> None:
authority = self.contract["authority"]
with tempfile.TemporaryDirectory(prefix="railiance-custody-registry-") as directory:
root = Path(directory)
os.chmod(root, 0o700)
document = root / "registry.json"
document.write_text(canonical_json(registry), encoding="utf-8")
os.chmod(document, 0o600)
try:
self.bao(
[
"kv",
"patch",
authority["registry_path"],
f"{authority['registry_field']}=@{document}",
],
label="write sender registry",
)
finally:
secure_unlink(document)
def identity_names(contract: dict[str, Any]) -> set[str]:
return {item["sender_name"] for item in contract["identities"]}
def token_paths(contract: dict[str, Any]) -> dict[str, str]:
authority = contract["authority"]
prefix = authority["kv_prefix"].strip("/")
return {
item["handle"]: f"{authority['kv_mount']}/{prefix}/{item['handle']}"
for item in contract["identities"]
}
def add_temporary_identities(
registry: list[dict[str, Any]], contract: dict[str, Any], values: dict[str, str]
) -> list[dict[str, Any]]:
names = identity_names(contract)
present = {str(item.get("name")) for item in registry if isinstance(item, dict)}
if present & names:
raise ProcedureError(f"temporary identities already exist: {sorted(present & names)}")
updated = list(registry)
expires = contract["window"]["expires_at"]
for identity in contract["identities"]:
updated.append(
{
"name": identity["sender_name"],
"tokens": [values[identity["handle"]]],
"sources": ["whitehat-security"],
"tenants": [identity["tenant"]],
"may_write": True,
"may_read": True,
"expires_at": expires,
}
)
return updated
def remove_temporary_identities(
registry: list[dict[str, Any]], contract: dict[str, Any]
) -> tuple[list[dict[str, Any]], list[str]]:
names = identity_names(contract)
removed = sorted(
str(item.get("name"))
for item in registry
if isinstance(item, dict) and item.get("name") in names
)
updated = [
item
for item in registry
if not isinstance(item, dict) or item.get("name") not in names
]
return updated, removed
def build_policy(contract: dict[str, Any]) -> str:
paths = token_paths(contract)
blocks: list[str] = []
for path in paths.values():
mount, relative = path.split("/", 1)
blocks.extend(
[
f'path "{mount}/data/{relative}" {{',
' capabilities = ["read"]',
"}",
"",
f'path "{mount}/metadata/{relative}" {{',
' capabilities = ["read"]',
"}",
"",
]
)
return "\n".join(blocks).rstrip() + "\n"
def build_manifest(contract: dict[str, Any], lease_id: str) -> str:
names = resource_names(contract)
authority = contract["authority"]
runner = contract["runner"]
engagement = contract["engagement_id"]
annotations = {
"custody.railiance/engagement": engagement,
"custody.railiance/lease-id": lease_id,
"custody.railiance/contract-sha256": contract_digest(contract),
}
document = {
"apiVersion": "v1",
"kind": "List",
"items": [
{
"apiVersion": "external-secrets.io/v1",
"kind": "ClusterSecretStore",
"metadata": {"name": names["store"], "annotations": annotations},
"spec": {
"provider": {
"vault": {
"server": "http://openbao.openbao.svc:8200",
"path": authority["kv_mount"],
"version": "v2",
"auth": {
"kubernetes": {
"mountPath": "kubernetes",
"role": names["role"],
"serviceAccountRef": {
"name": authority["eso_service_account"],
"namespace": authority["eso_namespace"],
},
}
},
}
},
"conditions": [{"namespaces": [runner["namespace"]]}],
},
},
{
"apiVersion": "external-secrets.io/v1",
"kind": "ExternalSecret",
"metadata": {
"name": names["external_secret"],
"namespace": runner["namespace"],
"annotations": annotations,
},
"spec": {
"refreshInterval": "1m",
"secretStoreRef": {"kind": "ClusterSecretStore", "name": names["store"]},
"target": {
"name": runner["secret_name"],
"creationPolicy": "Owner",
"deletionPolicy": "Delete",
},
"data": [
{
"secretKey": identity["handle"],
"remoteRef": {
"key": f"{authority['kv_prefix'].strip('/')}/{identity['handle']}",
"property": "token",
},
}
for identity in sorted(contract["identities"], key=lambda item: item["handle"])
],
},
},
],
}
return canonical_json(document)
def time_state(contract: dict[str, Any], now: datetime) -> str:
starts = parse_time(contract["window"]["starts_at"], "window.starts_at")
cutoff = parse_time(contract["window"]["projection_cutoff"], "window.projection_cutoff")
if now < starts:
return "before-window"
if now > cutoff:
return "projection-cutoff-passed"
return "projection-window-open"
def target_ready(operator: Operator, contract: dict[str, Any]) -> str:
target = contract["target"]
image = operator.kubectl(
[
"-n",
target["namespace"],
"get",
"deploy",
target["deployment"],
"-o",
f"jsonpath={{.spec.template.spec.containers[?(@.name==\"{target['container']}\")].image}}",
],
label="read target image",
).stdout.strip()
if not image.endswith("@" + target["image_digest"]):
raise ProcedureError("target is not on the contract-approved image digest")
ready = operator.kubectl(
[
"-n",
target["namespace"],
"get",
"deploy",
target["deployment"],
"-o",
"jsonpath={.status.readyReplicas}/{.status.replicas}",
],
label="read target readiness",
).stdout.strip()
if ready != "1/1":
raise ProcedureError(f"target is not 1/1 Ready (observed {ready or 'unknown'})")
return ready
def resource_presence(operator: Operator, contract: dict[str, Any]) -> dict[str, Any]:
names = resource_names(contract)
runner = contract["runner"]
paths = token_paths(contract)
identities = sorted(
identity_names(contract)
& {
str(item.get("name"))
for item in operator.registry()
if isinstance(item, dict)
}
)
exact_paths = sorted(
handle
for handle, path in paths.items()
if operator.bao_exists(
["kv", "metadata", "get", "-format=json", path],
label=f"check exact {handle} KV metadata",
)
)
resources: list[str] = []
checks = (
("store", ["get", "clustersecretstore", names["store"], "-o", "name"]),
(
"external_secret",
["-n", runner["namespace"], "get", "externalsecret", names["external_secret"], "-o", "name"],
),
(
"mounted_secret",
["-n", runner["namespace"], "get", "secret", runner["secret_name"], "-o", "name"],
),
)
for label, args in checks:
if operator.kubectl_exists(args, label=f"check exact {label}"):
resources.append(label)
return {
"temporary_identities_present": identities,
"exact_token_paths_present": exact_paths,
"projection_resources_present": resources,
}
def live_preflight(
operator: Operator,
contract: dict[str, Any],
*,
api_base: str = DEFAULT_API_BASE,
now: datetime | None = None,
) -> dict[str, Any]:
current = (now or utc_now()).astimezone(UTC)
operator.bao(["status", "-format=json"], label="read OpenBao status")
target = target_ready(operator, contract)
runner = contract["runner"]
authority = contract["authority"]
operator.kubectl(["get", "ns", runner["namespace"], "-o", "name"], label="verify runner namespace")
operator.kubectl(
["-n", authority["eso_namespace"], "get", "sa", authority["eso_service_account"], "-o", "name"],
label="verify ESO service account",
)
presence = resource_presence(operator, contract)
broker_ready = True
broker_reason = None
try:
current_broker_receipt(contract, api_base, now=current)
except ContractError as exc:
broker_ready = False
broker_reason = str(exc)
return {
"interface": "railiance.custody-preflight",
"version": 1,
"engagement_id": contract["engagement_id"],
"projection_contract_digest": contract_digest(contract),
"time_state": time_state(contract, current),
"broker_ready": broker_ready,
"broker_gate": "pass" if broker_ready else "blocked",
"broker_reason": broker_reason,
"target_image_matches": True,
"target_ready": target,
"openbao_initialized": True,
"openbao_sealed": False,
**presence,
"secret_values_observed": False,
}
def write_exact_values(operator: Operator, contract: dict[str, Any], values: dict[str, str]) -> None:
with tempfile.TemporaryDirectory(prefix="railiance-custody-values-") as directory:
root = Path(directory)
os.chmod(root, 0o700)
files: list[Path] = []
try:
for handle, path in token_paths(contract).items():
document = root / handle
document.write_text(values[handle], encoding="utf-8")
os.chmod(document, 0o600)
files.append(document)
operator.bao(
[
"kv",
"put",
path,
f"token=@{document}",
f"engagement_id={contract['engagement_id']}",
f"expires_at={contract['window']['expires_at']}",
],
label=f"write exact {handle} KV path",
)
finally:
for document in files:
secure_unlink(document)
def configure_projection(operator: Operator, contract: dict[str, Any], lease_id: str) -> None:
names = resource_names(contract)
authority = contract["authority"]
runner = contract["runner"]
operator.bao(
["policy", "write", names["policy"], "-"],
label="write exact ESO policy",
input_text=build_policy(contract),
)
operator.bao(
[
"write",
f"auth/kubernetes/role/{names['role']}",
f"bound_service_account_names={authority['eso_service_account']}",
f"bound_service_account_namespaces={authority['eso_namespace']}",
f"policies={names['policy']}",
"ttl=5m",
"max_ttl=15m",
],
label="write exact ESO Kubernetes role",
)
operator.kubectl(
["apply", "-f", "-"],
label="apply exact projection manifest",
input_text=build_manifest(contract, lease_id),
)
operator.kubectl(
["wait", "--for=condition=Ready", f"clustersecretstore/{names['store']}", "--timeout=90s"],
label="wait for exact store",
)
operator.kubectl(
[
"-n",
runner["namespace"],
"wait",
"--for=condition=Ready",
f"externalsecret/{names['external_secret']}",
"--timeout=90s",
],
label="wait for exact ExternalSecret",
)
keys = operator.kubectl(
[
"-n",
runner["namespace"],
"get",
"secret",
runner["secret_name"],
"-o",
'go-template={{range $k, $_ := .data}}{{$k}}{{"\\n"}}{{end}}',
],
label="verify projected key names",
).stdout.splitlines()
expected = sorted(identity["handle"] for identity in contract["identities"])
if sorted(keys) != expected:
raise ProcedureError(f"projected Secret has unexpected key names: {sorted(keys)}")
def force_sender_sync_and_restart(operator: Operator, contract: dict[str, Any]) -> None:
target = contract["target"]
namespace = target["namespace"]
external_secret = target["sender_external_secret"]
before = operator.kubectl(
["-n", namespace, "get", "secret", external_secret, "-o", "jsonpath={.metadata.resourceVersion}"],
label="read sender Secret resource version",
).stdout.strip()
operator.kubectl(
[
"-n",
namespace,
"annotate",
"externalsecret",
external_secret,
f"force-sync={int(time.time())}",
"--overwrite",
],
label="force sender registry sync",
)
deadline = time.monotonic() + 120
while time.monotonic() < deadline:
current = operator.kubectl(
["-n", namespace, "get", "secret", external_secret, "-o", "jsonpath={.metadata.resourceVersion}"],
label="poll sender Secret resource version",
).stdout.strip()
if current and current != before:
break
time.sleep(2)
else:
raise ProcedureError("sender Secret did not refresh within 120 seconds")
operator.kubectl(
["-n", namespace, "rollout", "restart", f"deploy/{target['deployment']}"],
label="restart sender registry reader",
)
operator.kubectl(
["-n", namespace, "rollout", "status", f"deploy/{target['deployment']}", "--timeout=120s"],
label="wait for target rollout",
)
def delete_projection(operator: Operator, contract: dict[str, Any]) -> None:
names = resource_names(contract)
runner = contract["runner"]
operator.kubectl(
["-n", runner["namespace"], "delete", "externalsecret", names["external_secret"], "--ignore-not-found"],
label="delete exact ExternalSecret",
)
operator.kubectl(
["-n", runner["namespace"], "delete", "secret", runner["secret_name"], "--ignore-not-found"],
label="delete exact mounted Secret",
)
operator.kubectl(
["delete", "clustersecretstore", names["store"], "--ignore-not-found"],
label="delete exact store",
)
for handle, path in token_paths(contract).items():
if operator.bao_exists(
["kv", "metadata", "get", "-format=json", path],
label=f"check exact {handle} metadata before delete",
):
operator.bao(["kv", "metadata", "delete", path], label=f"delete exact {handle} KV path")
if operator.bao_exists(["read", f"auth/kubernetes/role/{names['role']}"] , label="check exact role"):
operator.bao(["delete", f"auth/kubernetes/role/{names['role']}"] , label="delete exact role")
if operator.bao_exists(["policy", "read", names["policy"]], label="check exact policy"):
operator.bao(["policy", "delete", names["policy"]], label="delete exact policy")
def verify_cleanup_absent(operator: Operator, contract: dict[str, Any]) -> None:
operator.bao(["status", "-format=json"], label="verify OpenBao after cleanup")
operator.kubectl(["get", "ns", contract["runner"]["namespace"], "-o", "name"], label="verify Kubernetes after cleanup")
presence = resource_presence(operator, contract)
if any(presence.values()):
raise ProcedureError(f"engagement resources remain after cleanup: {presence}")
def cleanup_scope(operator: Operator, contract: dict[str, Any]) -> dict[str, Any]:
registry = operator.registry()
updated, removed = remove_temporary_identities(registry, contract)
if removed:
operator.write_registry(updated)
force_sender_sync_and_restart(operator, contract)
delete_projection(operator, contract)
verify_cleanup_absent(operator, contract)
target_ready(operator, contract)
return {
"removed_identities": removed,
"removed_resources": sorted(
[
*token_paths(contract).values(),
*resource_names(contract).values(),
]
),
}
def transactional(
steps: list[Callable[[], None]], rollback: Callable[[], Any]
) -> None:
try:
for step in steps:
step()
except Exception as original:
try:
rollback()
except Exception as cleanup_error:
raise ProcedureError(
"projection failed and exact cleanup could not be proven"
) from cleanup_error
raise original
def resource_uids(
operator: Operator, contract: dict[str, Any], *, require_all: bool = True
) -> dict[str, str]:
names = resource_names(contract)
runner = contract["runner"]
requests = {
"store": [
"get", "clustersecretstore", names["store"], "--ignore-not-found",
"-o", "jsonpath={.metadata.uid}",
],
"external_secret": [
"-n", runner["namespace"], "get", "externalsecret", names["external_secret"],
"--ignore-not-found",
"-o", "jsonpath={.metadata.uid}",
],
"mounted_secret": [
"-n", runner["namespace"], "get", "secret", runner["secret_name"],
"--ignore-not-found",
"-o", "jsonpath={.metadata.uid}",
],
}
result = {
key: operator.kubectl(args, label=f"read exact {key} UID").stdout.strip()
for key, args in requests.items()
}
if require_all and any(not value for value in result.values()):
raise ProcedureError("projection resource UID evidence is incomplete")
return {key: value for key, value in result.items() if value}
def verify_receipt_resource_scope(
operator: Operator, contract: dict[str, Any], receipt: dict[str, Any]
) -> dict[str, str]:
present = resource_uids(operator, contract, require_all=False)
expected = receipt["resources"]["uids"]
mismatched = {
key: value
for key, value in present.items()
if expected.get(key) != value
}
if mismatched:
raise ProcedureError(
"live projection resource UID differs from the cleanup receipt; refusing deletion"
)
return present
def project(
operator: Operator,
contract: dict[str, Any],
*,
api_base: str = DEFAULT_API_BASE,
now: datetime | None = None,
lease_id: str | None = None,
) -> dict[str, Any]:
current = (now or utc_now()).astimezone(UTC)
if time_state(contract, current) != "projection-window-open":
raise ProcedureError("projection is outside the approved projection window")
# This is intentionally before token generation and every mutation.
broker_receipt = current_broker_receipt(contract, api_base, now=current)
preflight = live_preflight(operator, contract, api_base=api_base, now=current)
if not preflight["broker_ready"]:
raise ProcedureError("consumer broker is not ready")
if any(
preflight[key]
for key in (
"temporary_identities_present",
"exact_token_paths_present",
"projection_resources_present",
)
):
raise ProcedureError("cleanup is required before a new projection")
active_lease = lease_id or f"custody:{secrets.token_hex(16)}"
values = {identity["handle"]: secrets.token_urlsafe(48) for identity in contract["identities"]}
if len(set(values.values())) != len(values):
raise ProcedureError("credential generator returned duplicate values")
try:
transactional(
[
lambda: write_exact_values(operator, contract, values),
lambda: operator.write_registry(
add_temporary_identities(operator.registry(), contract, values)
),
lambda: configure_projection(operator, contract, active_lease),
lambda: force_sender_sync_and_restart(operator, contract),
],
lambda: cleanup_scope(operator, contract),
)
finally:
values.clear()
projected_at = rfc3339(utc_now())
base = {
"interface": PROJECTION_INTERFACE,
"version": 1,
"workplan_id": WORKPLAN_ID,
"state": "projected",
"lease_id": active_lease,
"engagement_id": contract["engagement_id"],
"target": {
"id": contract["target"]["id"],
"revision": contract["target"]["revision"],
"image_digest": contract["target"]["image_digest"],
},
"projection_contract_digest": contract_digest(contract),
"broker_receipt_digest": digest(broker_receipt),
"projected_at": projected_at,
"expires_at": contract["window"]["expires_at"],
"identities": sorted(
(
{
"handle": item["handle"],
"role": item["role"],
"sender_name": item["sender_name"],
"mount_path": item["mount_path"],
}
for item in contract["identities"]
),
key=lambda item: item["handle"],
),
"resources": {
"names": resource_names(contract),
"uids": resource_uids(operator, contract),
},
"cleanup_authority": "railiance-platform",
"secret_values_observed": False,
}
receipt = {**base, "receipt_id": f"sha256:{digest(base)}"}
validate_projection_receipt(receipt, contract)
return receipt
def projection_status(
operator: Operator, contract: dict[str, Any], receipt: dict[str, Any]
) -> dict[str, Any]:
validate_projection_receipt(receipt, contract)
presence = resource_presence(operator, contract)
counts = [
len(presence["temporary_identities_present"]),
len(presence["exact_token_paths_present"]),
len(presence["projection_resources_present"]),
]
state = "active" if counts == [2, 2, 3] else "absent" if counts == [0, 0, 0] else "partial"
uid_match = False
present_uids = resource_uids(operator, contract, require_all=False)
uid_match = bool(present_uids) and all(
receipt["resources"]["uids"].get(key) == value
for key, value in present_uids.items()
)
if present_uids and not uid_match:
state = "mismatched"
return {
"interface": "railiance.custody-projection-status",
"version": 1,
"engagement_id": contract["engagement_id"],
"lease_id": receipt["lease_id"],
"state": state,
"resource_uids_match": uid_match,
**presence,
"target_ready": target_ready(operator, contract),
"secret_values_observed": False,
}
def cleanup(
operator: Operator, contract: dict[str, Any], projection: dict[str, Any]
) -> dict[str, Any]:
validate_projection_receipt(projection, contract)
verify_receipt_resource_scope(operator, contract, projection)
result = cleanup_scope(operator, contract)
receipt = {
"interface": CLEANUP_INTERFACE,
"version": 1,
"workplan_id": WORKPLAN_ID,
"state": "cleaned",
"lease_id": projection["lease_id"],
"engagement_id": contract["engagement_id"],
"projection_receipt_id": projection["receipt_id"],
"cleaned_at": rfc3339(utc_now()),
"removed_resources": result["removed_resources"],
"target_ready": True,
"secret_values_observed": False,
}
validate_cleanup_receipt(receipt, projection, contract)
return receipt
def assert_confirm(value: str | None, expected: str) -> None:
if value != expected:
raise ProcedureError(f"live command requires --confirm {expected}")
def assert_expired_cleanup(contract: dict[str, Any], now: datetime) -> None:
expires = parse_time(contract["window"]["expires_at"], "window.expires_at")
if now.astimezone(UTC) <= expires:
raise ProcedureError("cleanup-expired refuses before receipt expiry")
def main() -> int:
parser = argparse.ArgumentParser()
parser.add_argument(
"command",
choices=["validate", "render", "preflight", "project", "status", "cleanup", "cleanup-expired"],
)
parser.add_argument("--contract", type=Path, required=True)
parser.add_argument("--receipt", type=Path)
parser.add_argument("--receipt-out", type=Path)
parser.add_argument("--token-file", type=Path, default=Path.home() / ".local/openbao/platform-admin.token")
parser.add_argument("--state-hub", default=DEFAULT_API_BASE)
parser.add_argument("--confirm")
args = parser.parse_args()
try:
contract = validate_projection_contract(load_json(args.contract))
if args.command == "validate":
result = {
"valid": True,
"projection_contract_digest": contract_digest(contract),
"resource_names": resource_names(contract),
"secret_values_observed": False,
}
elif args.command == "render":
result = {
"policy": build_policy(contract),
"manifest": json.loads(build_manifest(contract, "custody:" + "0" * 32)),
"secret_values_observed": False,
}
else:
projection = None
if args.command in {"status", "cleanup", "cleanup-expired"}:
if not args.receipt:
raise ProcedureError(f"{args.command} requires --receipt")
projection = validate_projection_receipt(load_json(args.receipt), contract)
if args.command == "cleanup-expired":
assert_expired_cleanup(contract, utc_now())
if args.command == "project":
# Reject an unready consumer before reading platform-admin
# custody material. project() repeats this immutable gate.
current_broker_receipt(contract, args.state_hub)
operator = Operator(contract, args.token_file)
if args.command == "preflight":
result = live_preflight(operator, contract, api_base=args.state_hub)
elif args.command == "project":
assert_confirm(args.confirm, f"{contract['engagement_id']}:attended")
result = project(operator, contract, api_base=args.state_hub)
else:
assert projection is not None
if args.command == "status":
result = projection_status(operator, contract, projection)
else:
if args.command == "cleanup-expired":
assert_confirm(args.confirm, f"{contract['engagement_id']}:expired")
else:
assert_confirm(args.confirm, f"{contract['engagement_id']}:cleanup")
result = cleanup(operator, contract, projection)
if args.receipt_out:
write_receipt(args.receipt_out, result)
print(json.dumps(result, indent=2, sort_keys=True))
return 0
except (ContractError, OSError, ProcedureError, RemoteExecutionError) as exc:
print(f"custody projection failed: {exc}", file=sys.stderr)
return 1
if __name__ == "__main__":
raise SystemExit(main())