rein-aharness/scripts/metered-native-owner.py

348 lines
21 KiB
Python
Raw Normal View History

"""Single attended, non-retrying Railiance proof under the six frozen approvals.
Non-secret bundle is staged separately. Credentials arrive only on SSH stdin,
live in owner-only tmpfs for this process, and never enter the bound child except
through Secrets Engine's approved provider/worker delivery.
"""
from __future__ import annotations
import argparse
import contextlib
import copy
import hashlib
import io
import json
import os
from pathlib import Path
import sqlite3
import stat
import subprocess
import sys
import tempfile
import time
from datetime import datetime, timezone
from dataclasses import replace
PACKET_SHA = "22e640428e3a7ae8fc2363cf193db98381cd94e4be7ab009bc6703658d6fc544"
PROVIDER = "glas-claude-agent-dev-anthropic"
WORKER = "activity-core-metered-worker-token"
IDS = {
PROVIDER: dict(apply="0c05bd0a-f81f-451d-840c-5565628e2edc", verify="47f118a3-a86c-43ad-969d-42e09a0f45bb", exec="7b32443a-a817-400c-a130-01b9ef04c8ee"),
WORKER: dict(apply="2ce76d7f-01c3-4b63-8476-8d1230679769", verify="c2af4bcf-15b4-4305-aac1-acc7e4dcb099", exec="dc22666a-d5d7-46bc-9490-ebe9fe09dc38"),
}
LIMITS = dict(token_ttl=300, token_max_ttl=900, secret_id_ttl=300, secret_id_num_uses=1, token_num_uses=8)
OWNER = Path("/home/tegwick/hfact/owner-metered")
TARGET = Path("/home/tegwick/hfact/targets/hfact-glas-proof")
BASE_HEAD = "679d23e0517707ba63e25399b5262f0d2f37315d"
OLD_FILES = {"owner.json": "e0d3fb84649fdca302eccd9415f3cc2beaee84bf07bd83205cdffbe93e5573bc", "spend-policy.json": "f31c585916de1ea8bd8e48c72421803dfde1015e6201704b9320f77d2c545d9c"}
DEFINITION = "5bae5505-77f1-5ba3-8dfa-2e10ed15531e"
RECEIPT_NAME = "execution.json"
def require(value, code):
if not value:
raise ValueError(code)
def private(path, directory=False):
s = path.lstat()
require(s.st_uid == os.getuid() and stat.S_IMODE(s.st_mode) == (0o700 if directory else 0o600) and (stat.S_ISDIR(s.st_mode) if directory else stat.S_ISREG(s.st_mode)), "private_path_required")
def write_new(path, value):
fd = os.open(path, os.O_CREAT | os.O_EXCL | os.O_WRONLY | os.O_NOFOLLOW, 0o600)
with os.fdopen(fd, "w") as stream:
stream.write(value)
stream.flush()
os.fsync(stream.fileno())
def run(args):
p = subprocess.run(args, capture_output=True, text=True, timeout=30)
require(p.returncode == 0, "metadata_command_failed")
return p.stdout.strip()
def save(bundle, receipt):
path = bundle / RECEIPT_NAME
temp = path.with_suffix(".tmp")
temp.write_text(json.dumps(receipt, indent=2) + "\n")
temp.chmod(0o600)
temp.replace(path)
def queue_rows():
# Read-only owner metadata via the existing authenticated cluster transport.
code = '''import asyncio,json
from activity_core.db import make_engine
from sqlalchemy import text
async def main():
e=make_engine()
async with e.connect() as c:
rows=await c.execute(text("SELECT id,state,attempt,claim_owner,target_repo,harness_profile_ref,repository_grant,triggering_event_id FROM ops_runs WHERE activity_definition_id='5bae5505-77f1-5ba3-8dfa-2e10ed15531e' OR labels @> '[\\\"hfact-metered\\\"]'::jsonb"))
print(json.dumps([dict(x) for x in rows.mappings()],default=str))
await e.dispose()
asyncio.run(main())'''
return json.loads(run(["/usr/local/bin/kubectl", "-n", "activity-core", "exec", "deploy/actcore-api", "--", "python", "-c", code]))
def install_host(bundle):
"""Run with the immutable admitted runtime; preserve every old ledger byte."""
from rein_aharness.spend_admission import SpendLedger, SpendPolicy
from rein_aharness.request_admission import RequestLedger
import yaml
private(OWNER, True)
require(not (bundle / "installation.json").exists(), "prior_installation_requires_reconciliation")
require(not queue_rows(), "metered_queue_not_empty")
require(run(["git", "-C", str(TARGET), "rev-parse", "HEAD"]) == BASE_HEAD and not run(["git", "-C", str(TARGET), "status", "--porcelain"]), "target_not_pristine")
for filename, expected in OLD_FILES.items():
private(OWNER / filename)
require(hashlib.sha256((OWNER / filename).read_bytes()).hexdigest() == expected, "old_host_pin_drift")
candidate = yaml.safe_load((bundle / "catalog" / (PROVIDER + ".yaml")).read_text())["delivery_config"]["exec_owner"]
for filename in OLD_FILES:
require(hashlib.sha256((bundle / filename).read_bytes()).hexdigest() == candidate["files"][str(OWNER / filename)]["sha256"], "candidate_pin_drift")
old = OWNER / "spend.sqlite3"
private(old)
db = sqlite3.connect(old)
try:
db.execute("BEGIN EXCLUSIVE")
counts = {table: db.execute("SELECT count(*) FROM " + table).fetchone()[0] for table in ("reservations", "request_routes", "request_reservations")}
require(not any(counts.values()), "prior_liabilities_require_reconciliation")
policy = SpendPolicy.load(bundle / "spend-policy.json")
ledger = SpendLedger(Path(candidate["environment"]["AGENT_HARNESS_SPEND_LEDGER"]), policy)
require(not ledger.path.exists(), "new_ledger_already_exists")
for filename in OLD_FILES:
write_new(OWNER / (filename + ".pre-renewal-20260927"), (OWNER / filename).read_text())
ledger.initialize()
RequestLedger(ledger).initialize()
for filename in OLD_FILES:
staged = OWNER / (filename + ".renewal-staged")
write_new(staged, (bundle / filename).read_text())
staged.replace(OWNER / filename)
# Changing both pins retires the old catalog's dispatch. There is no
# active metered worker; old ledger and files remain preserved.
check = subprocess.run([*candidate["command"], "--check"], env=candidate["environment"], cwd=candidate["cwd"], capture_output=True, text=True, timeout=30)
require(check.returncode == 0, "installed_owner_check_failed")
receipt = dict(status="installed", observed_at=datetime.now(timezone.utc).isoformat(), old_counts=counts, old_ledger_preserved=True, new_ledger=str(ledger.path), old_dispatch_pins_retired=True, owner_check=json.loads(check.stdout), dispatches=0)
write_new(bundle / "installation.json", json.dumps(receipt, indent=2) + "\n")
finally:
db.rollback()
db.close()
def prepare_configs(bundle, private_dir, trust=None):
import yaml
from secrets_engine.config import Config
from secrets_engine.catalog import get_entry
from secrets_engine.approval_consume import _expected_request
packet = (bundle / "native-pdp-inputs.json").read_bytes()
require(hashlib.sha256(packet).hexdigest() == PACKET_SHA, "frozen_packet_drift")
data = json.loads(packet)
rows = data if isinstance(data, list) else data["requests"]
expected = {(r["catalog"], r["action"]): r["request"] for r in rows}
require(set(expected) == {(lane, action) for lane in IDS for action in IDS[lane]}, "six_exact_requests_required")
configs, entries = {}, {}
for action in ("apply", "verify", "exec"):
folder = private_dir / action
folder.mkdir(mode=0o700)
for lane in IDS:
doc = yaml.safe_load((bundle / "catalog" / (lane + ".yaml")).read_text())
doc = copy.deepcopy(doc)
doc["approval"]["authorization_id"] = IDS[lane][action]
write_new(folder / (lane + ".yaml"), yaml.safe_dump(doc, sort_keys=False))
cfg = replace(Config.load(), catalog_dir=folder, evidence_dir=bundle / "native-evidence", hub_url="", bao_addr="http://127.0.0.1:28200", approval_url="http://127.0.0.1:28281", approval_token_file=None, approval_client_secret_file=private_dir / "client-secret", keycape_token_url="https://kc.coulomb.social/token", keycape_issuer="https://kc.coulomb.social", keycape_client_secret_file=None, openbao_jwt_login_file=None, authorization_subject_id="secrets-engine", authorization_subject_type="service", authorization_policy_package="secrets-engine.catalog-lane.lifecycle", authorization_policy_version="v2", authorization_min_approvals=1, pdp_url="http://127.0.0.1:28282", pdp_token_file=private_dir / "pdp-caller", clock_trust_file=trust)
configs[action] = cfg
for lane in IDS:
entry = get_entry(folder, lane)
actual = _expected_request(cfg, entry, action, fields=() if action == "apply" else tuple(entry.fields), policy_targets=(entry.policy_name,), auth_targets=(entry.role_name,))
require(actual == expected[lane, action], "frozen_action_request_drift")
entries[lane, action] = entry
return configs, entries
def clock_admission(bundle, directory):
from railiance_clock.admission import admit
from urllib.request import urlopen
custody = json.loads((bundle / "clock-public.json").read_text())
require(hashlib.sha256(custody["public_key_pem"].encode()).hexdigest() == "bd583446b5ed61d086806b2a0c5aaf33a875b751e45599e75335d1f415be609a", "clock_key_drift")
with urlopen("http://127.0.0.1:8787/healthz", timeout=5) as response:
epoch = json.load(response)["epoch"]
require(epoch == "b1164ccb-a4c2-4cc8-adf8-1d5597de697b", "clock_epoch_requires_readmission")
return admit(directory=directory, authority_id="railiance01", environment="prod", kid=custody["kid"], epoch=epoch, policy_id="railiance01-online-v1", public_key_pem=custody["public_key_pem"], endpoint="http://127.0.0.1:8787/v1/time-samples", limits=dict(max_age_ns=5000000000, max_rtt_ns=2000000000, max_width_ns=3000000000, max_server_error_ns=500000000, timer_ppm=1000, timer_resolution_ns=1000, suspend_tolerance_ns=5000000, trust_session_ns=900000000000), admission_ref="CCR-2026-0028; REINAH-WP-0003; attended proof 2026-09-27", transport="admitted-loopback")
def verify_cleanup(client, entry, other):
from urllib.request import Request, urlopen
from urllib.error import HTTPError
with client.approle_session(entry.role_name) as session:
token = session.client.token
for path in (f"platform/data/{other.path}", "platform/metadata/workloads", "platform/data/workloads/secrets-engine/approval-client"):
require(session.client.token_capabilities(path, token=token) == ["deny"], "sibling_or_metadata_authority")
require(session.revocation_succeeded, "session_revocation_failed")
try:
with urlopen(Request(client.addr + "/v1/auth/token/lookup-self", headers={"X-Vault-Token": token}), timeout=15):
raise ValueError("revoked_token_still_usable")
except HTTPError as error:
require(error.code == 403, "revocation_not_definitive")
finally:
del token
return dict(sibling_denied=True, metadata_denied=True, session_revoked=True, revoked_lookup_status=403)
def validate_queued(rows, trigger):
require(len(rows) == 1 and rows[0]["state"] == "open" and rows[0]["attempt"] == 0 and rows[0]["claim_owner"] is None and rows[0]["target_repo"] == "hfact-glas-proof" and rows[0]["harness_profile_ref"] == "harness.agent-dev-local@1.1.1" and rows[0]["repository_grant"] == {"version": "1", "allowed_paths": ["PROOF.md"], "commit_count": {"min": 1, "max": 1}, "publish": False} and rows[0]["triggering_event_id"] == trigger, "natural_queue_binding_refused")
def execute(bundle, resume=False):
from secrets_engine.openbao import OpenBaoClient
from secrets_engine.approval_consume import authorize_action
from secrets_engine.application_time import read_window
from secrets_engine.cli import build_parser
from secrets_engine.exec_owner import validate_delivery_target
from secrets_engine.plan import build_plan
global RECEIPT_NAME
prior = None
if resume:
prior = json.loads((bundle / "execution.json").read_text())
require(prior.get("phase") == "one_trigger_started" and prior.get("failure_code") == "natural_queue_binding_refused" and prior.get("trigger", {}).get("trigger_key") == "manual-dab6c4c7-db03-4887-a17b-9d99b752482e", "unexpected_prior_failure")
require([(r["catalog"], r["action"], r["approval_id"], r["exit_code"]) for r in prior["actions"]] == [(lane, action, IDS[lane][action], 0) for action in ("apply", "verify") for lane in IDS], "prior_native_actions_not_verified")
RECEIPT_NAME = "execution-resume.json"
require(not (bundle / RECEIPT_NAME).exists(), "prior_attempt_requires_reconciliation")
receipt = dict(status="failed", phase="preflight", observed_at=datetime.now(timezone.utc).isoformat(), actions=[], automatic_retry=False)
save(bundle, receipt)
payload = {}
directory = None
try:
require(run(["git", "-C", str(TARGET), "rev-parse", "HEAD"]) == BASE_HEAD and not run(["git", "-C", str(TARGET), "status", "--porcelain"]), "target_not_pristine")
require(json.loads((bundle / "installation.json").read_text())["status"] == "installed", "installation_receipt_required")
if resume:
validate_queued(queue_rows(), prior["trigger"]["trigger_key"])
with sqlite3.connect(f"file:{OWNER}/spend-tool-proof-renewal-20260927.sqlite3?mode=ro", uri=True) as db:
require(all(db.execute("SELECT count(*) FROM " + table).fetchone()[0] == 0 for table in ("reservations", "request_routes", "request_reservations")), "prior_paid_dispatch_requires_reconciliation")
receipt.update(prior_receipt="execution.json", completed_actions_not_replayed=["provider/apply", "worker/apply", "provider/verify", "worker/verify", "queue/trigger"], trigger=prior["trigger"])
else:
require(not queue_rows(), "metered_queue_not_empty")
runtime = Path("/run/user") / str(os.getuid())
private(runtime, True)
require(run(["findmnt", "-n", "-o", "FSTYPE", "-T", str(runtime)]) == "tmpfs", "tmpfs_required")
with tempfile.TemporaryDirectory(prefix="metered-native-", dir=runtime) as name:
directory = Path(name)
private(directory, True)
payload = json.loads(sys.stdin.buffer.read(131073))
require(set(payload) == {"client_secret", "negative_token", "backend_token", "pdp_token"}, "credential_envelope_invalid")
for key, filename in (("client_secret", "client-secret"), ("negative_token", "negative-token"), ("pdp_token", "pdp-caller")):
require(isinstance(payload[key], str) and 0 < len(payload[key]) < 32768, "credential_envelope_invalid")
write_new(directory / filename, payload.pop(key))
os.environ["BAO_TOKEN"] = payload.pop("backend_token")
os.environ.pop("VAULT_TOKEN", None)
os.environ["BAO_ADDR"] = "http://127.0.0.1:28200"
trust = clock_admission(bundle, directory)
configs, entries = prepare_configs(bundle, directory, trust)
read_window(configs["exec"])
command = entries[PROVIDER, "exec"].delivery_config["exec_owner"]["command"]
validate_delivery_target(entries[PROVIDER, "exec"], "ANTHROPIC_API_KEY", command, "exec-env")
client = OpenBaoClient.resolve(configs["exec"].bao_addr)
identity = json.loads(run([client.bao_bin, "token", "lookup", "-format=json"]))["data"]
policies = set(identity["policies"]) | set(identity.get("identity_policies", []))
require("platform-admin" in policies and "root" not in policies and identity.get("entity_id") and 0 < identity["ttl"] <= 3600, "attended_operator_required")
# Refuse drift/existing state before any native consume. An already
# applied lane needs reconciliation, never a replay of this procedure.
for lane in IDS:
entry = entries[lane, "apply"]
if resume:
require(client.read_policy(entry.policy_name).strip() == build_plan(entry, "prod").policy_hcl.strip(), "applied_policy_drift")
role = json.loads(run([client.bao_bin, "read", "-format=json", "auth/approle/role/" + entry.role_name]))["data"]
require(all(role.get(k) == v for k, v in LIMITS.items()) and role.get("token_policies") == [entry.policy_name], "role_limits_drift")
else:
require(client.read_policy(entry.policy_name) is None and not client.approle_exists(entry.role_name), "existing_lane_requires_reconciliation")
for lane in IDS:
for action in (("exec",) if resume else IDS[lane]):
entry = entries[lane, action]
require(authorize_action(configs[action], entry, action, fields=() if action == "apply" else tuple(entry.fields), policy_targets=(entry.policy_name,), auth_targets=(entry.role_name,)) is not None, "native_authorization_missing")
receipt["phase"] = "remaining_exec_checks_passed" if resume else "all_six_native_checks_passed"
save(bundle, receipt)
for action in (() if resume else ("apply", "verify")):
for lane in IDS:
receipt["phase"] = lane + "_" + action + "_started"
save(bundle, receipt)
argv = ["apply", lane, "--stage", "prod", "--auth", "env"] if action == "apply" else ["verify", lane, "--auth", "env", "--negative-token-file", str(directory / "negative-token")]
args = build_parser().parse_args(argv)
with contextlib.redirect_stdout(io.StringIO()), contextlib.redirect_stderr(io.StringIO()):
code = args.func(configs[action], args)
require(code == 0, "native_action_failed")
row = dict(catalog=lane, action=action, approval_id=IDS[lane][action], exit_code=code)
entry = entries[lane, action]
if action == "apply":
role = json.loads(run([client.bao_bin, "read", "-format=json", "auth/approle/role/" + entry.role_name]))["data"]
require(all(role.get(k) == v for k, v in LIMITS.items()) and role.get("token_policies") == [entry.policy_name], "role_limits_drift")
require(client.read_policy(entry.policy_name).strip() == build_plan(entry, "prod").policy_hcl.strip(), "applied_policy_drift")
row["limits"] = LIMITS
else:
row.update(verify_cleanup(client, entry, entries[WORKER if lane == PROVIDER else PROVIDER, action]))
receipt["actions"].append(row)
save(bundle, receipt)
# The owner triggers exactly once only after both lanes verify.
# Persist intent first; any uncertainty blocks re-trigger/re-exec.
if not resume:
receipt["phase"] = "one_trigger_started"
save(bundle, receipt)
from urllib.request import Request, urlopen
with urlopen(Request("http://127.0.0.1:8010/activity-definitions/" + DEFINITION + "/trigger", method="POST"), timeout=30) as response:
require(response.status == 202, "trigger_refused")
receipt["trigger"] = json.load(response)
for _ in range(20):
rows = queue_rows()
if rows:
break
time.sleep(1)
validate_queued(rows, receipt["trigger"]["trigger_key"])
receipt["queued_run"] = rows[0]
receipt["phase"] = "one_exec_started"
save(bundle, receipt)
args = build_parser().parse_args(["exec", "--catalog", PROVIDER, "--auth", "env", "--mode", "exec-env", "--", *command])
args.command = args.command[1:]
captured = io.StringIO()
with contextlib.redirect_stdout(captured), contextlib.redirect_stderr(io.StringIO()):
code = args.func(configs["exec"], args)
# Only the closed owner result schema may leave this envelope.
results = []
for line in captured.getvalue().splitlines():
try:
row = json.loads(line)
except ValueError:
continue
if isinstance(row, dict) and set(row) <= {"ok", "claimed", "empty", "run_id", "ops_state", "code"} and "ok" in row:
results.append(row)
receipt.update(exec_exit_code=code, owner_results=results, phase="one_exec_returned")
receipt["actions"].extend(dict(catalog=lane, action="exec", approval_id=IDS[lane]["exec"], exit_code=code) for lane in IDS)
receipt["status"] = "passed" if code == 0 and len(results) == 1 and results[0].get("claimed") and results[0].get("ok") else "attempt_failed"
receipt["private_runtime_removed"] = True
except Exception as error:
receipt["failure_type"] = type(error).__name__
if type(error) is ValueError and str(error).replace("_", "").isalnum():
receipt["failure_code"] = str(error)
finally:
payload.clear()
os.environ.pop("BAO_TOKEN", None)
if directory is not None:
receipt["private_runtime_removed"] = not directory.exists()
save(bundle, receipt)
return 0 if receipt["status"] == "passed" else 1
if __name__ == "__main__":
parser = argparse.ArgumentParser()
parser.add_argument("mode", choices=("validate", "install", "execute", "resume"))
parser.add_argument("bundle", type=Path)
args = parser.parse_args()
if args.mode == "validate":
with tempfile.TemporaryDirectory() as directory:
prepare_configs(args.bundle, Path(directory))
print("Six frozen request bindings verified")
elif args.mode == "install":
install_host(args.bundle)
else:
raise SystemExit(execute(args.bundle, resume=args.mode == "resume"))