enterprise_delivery.py 43 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032
  1. #!/usr/bin/env python3
  2. """P3-WP07 local, fail-closed enterprise delivery evidence controls."""
  3. from __future__ import annotations
  4. import argparse
  5. import contextlib
  6. import fcntl
  7. import hashlib
  8. import json
  9. import os
  10. import re
  11. import stat
  12. import tempfile
  13. from datetime import UTC, datetime, timedelta
  14. from pathlib import Path
  15. from typing import Any
  16. from cryptography.hazmat.primitives import serialization
  17. from cryptography.hazmat.primitives.asymmetric.ed25519 import (
  18. Ed25519PrivateKey,
  19. Ed25519PublicKey,
  20. )
  21. ROOT = Path(__file__).resolve().parents[2]
  22. HEAD = "20260811_495"
  23. ENVIRONMENTS = ("development", "test", "staging", "production")
  24. TRUST_STORE = ROOT / "deployment/enterprise/trust_store.json"
  25. MAX_APPROVAL_WINDOW = timedelta(hours=24)
  26. CLOCK_SKEW = timedelta(minutes=5)
  27. LOCAL_TEST_PRIVATE_KEY = bytes.fromhex(
  28. "4f3b2775a4df921336c3882f071d9edce4a6e946b928345f7bf68b55d4eaa3d8"
  29. )
  30. class DeliveryError(ValueError):
  31. pass
  32. def canonical(value: Any) -> bytes:
  33. return json.dumps(
  34. value, ensure_ascii=False, sort_keys=True, separators=(",", ":")
  35. ).encode()
  36. def sha_bytes(value: bytes) -> str:
  37. return hashlib.sha256(value).hexdigest()
  38. def sha_file(path: Path) -> str:
  39. return sha_bytes(safe_bytes(path))
  40. def safe_bytes(path: Path) -> bytes:
  41. """Read one stable regular file through directory FDs without link traversal."""
  42. path = path.absolute()
  43. parts = path.parts
  44. if len(parts) < 2 or parts[0] != os.sep:
  45. raise DeliveryError(f"evidence path must be absolute: {path}")
  46. descriptor = os.open(
  47. os.sep, os.O_RDONLY | os.O_DIRECTORY | getattr(os, "O_NOFOLLOW", 0)
  48. )
  49. try:
  50. for component in parts[1:-1]:
  51. next_descriptor = os.open(
  52. component,
  53. os.O_RDONLY | os.O_DIRECTORY | getattr(os, "O_NOFOLLOW", 0),
  54. dir_fd=descriptor,
  55. )
  56. info = os.fstat(next_descriptor)
  57. if not stat.S_ISDIR(info.st_mode):
  58. os.close(next_descriptor)
  59. raise DeliveryError(f"evidence path contains a non-directory: {path}")
  60. os.close(descriptor)
  61. descriptor = next_descriptor
  62. name = parts[-1]
  63. before_name = os.stat(name, dir_fd=descriptor, follow_symlinks=False)
  64. if not stat.S_ISREG(before_name.st_mode) or before_name.st_nlink != 1:
  65. raise DeliveryError(f"evidence file must be an unlinked regular file: {path}")
  66. handle = os.open(
  67. name, os.O_RDONLY | getattr(os, "O_NOFOLLOW", 0), dir_fd=descriptor
  68. )
  69. with os.fdopen(handle, "rb") as stream:
  70. before = os.fstat(stream.fileno())
  71. if (
  72. before.st_ino != before_name.st_ino
  73. or before.st_dev != before_name.st_dev
  74. or before.st_nlink != 1
  75. ):
  76. raise DeliveryError(f"evidence file changed while reading: {path}")
  77. value = stream.read()
  78. after = os.fstat(stream.fileno())
  79. except OSError as exc:
  80. raise DeliveryError(f"unsafe evidence path (symbolic links are forbidden): {path}") from exc
  81. finally:
  82. os.close(descriptor)
  83. if (
  84. before.st_ino != after.st_ino
  85. or before.st_dev != after.st_dev
  86. or after.st_nlink != 1
  87. ):
  88. raise DeliveryError(f"evidence file changed while reading: {path}")
  89. return value
  90. def read_json(path: Path) -> dict[str, Any]:
  91. return read_json_bytes(safe_bytes(path), str(path))
  92. def read_json_bytes(value: bytes, label: str) -> dict[str, Any]:
  93. try:
  94. parsed = json.loads(value)
  95. except json.JSONDecodeError as exc:
  96. raise DeliveryError(f"invalid JSON evidence: {label}") from exc
  97. if not isinstance(parsed, dict):
  98. raise DeliveryError(f"JSON object required: {label}")
  99. return parsed
  100. def require_digest(value: str, field: str) -> str:
  101. if len(value) != 64 or set(value) - set("0123456789abcdef"):
  102. raise DeliveryError(f"{field} must be a lowercase sha256 digest")
  103. return value
  104. def safe_state_root(path: Path) -> Path:
  105. path = path.absolute()
  106. if path.is_symlink() or any(
  107. part.is_symlink() for part in [path, *path.parents] if part.exists()
  108. ):
  109. raise DeliveryError("state root and ancestors must not be symbolic links")
  110. path.mkdir(mode=0o700, parents=True, exist_ok=True)
  111. info = os.lstat(path)
  112. if not stat.S_ISDIR(info.st_mode) or info.st_mode & 0o077:
  113. raise DeliveryError("state root must be a private directory")
  114. return path
  115. def private_dir_fd(root: Path, *parts: str) -> tuple[Path, int]:
  116. """Open an owned state directory without following any controlled path part."""
  117. current = safe_state_root(root)
  118. descriptor = os.open(
  119. current, os.O_RDONLY | os.O_DIRECTORY | getattr(os, "O_NOFOLLOW", 0)
  120. )
  121. try:
  122. for part in parts:
  123. if not part or "/" in part or part in {".", ".."}:
  124. raise DeliveryError("invalid state directory component")
  125. with contextlib.suppress(FileExistsError):
  126. os.mkdir(part, mode=0o700, dir_fd=descriptor)
  127. next_descriptor = os.open(
  128. part,
  129. os.O_RDONLY
  130. | os.O_DIRECTORY
  131. | getattr(os, "O_NOFOLLOW", 0),
  132. dir_fd=descriptor,
  133. )
  134. info = os.fstat(next_descriptor)
  135. if not stat.S_ISDIR(info.st_mode) or info.st_mode & 0o077:
  136. os.close(next_descriptor)
  137. raise DeliveryError("state subdirectory must be a private directory")
  138. os.close(descriptor)
  139. descriptor = next_descriptor
  140. current = current / part
  141. except OSError as exc:
  142. os.close(descriptor)
  143. raise DeliveryError("unsafe state subdirectory") from exc
  144. except Exception:
  145. os.close(descriptor)
  146. raise
  147. return current, descriptor
  148. def state_read_json(descriptor: int, name: str) -> dict[str, Any]:
  149. if "/" in name or name in {"", ".", ".."}:
  150. raise DeliveryError("invalid state file name")
  151. try:
  152. handle = os.open(
  153. name, os.O_RDONLY | getattr(os, "O_NOFOLLOW", 0), dir_fd=descriptor
  154. )
  155. with os.fdopen(handle, "rb") as stream:
  156. before = os.fstat(stream.fileno())
  157. if not stat.S_ISREG(before.st_mode) or before.st_nlink != 1:
  158. raise DeliveryError("state file must be an unlinked regular file")
  159. value = stream.read()
  160. after = os.fstat(stream.fileno())
  161. except FileNotFoundError:
  162. raise
  163. except OSError as exc:
  164. raise DeliveryError("unsafe state file") from exc
  165. if before.st_ino != after.st_ino or after.st_nlink != 1:
  166. raise DeliveryError("state file changed while reading")
  167. try:
  168. parsed = json.loads(value)
  169. except json.JSONDecodeError as exc:
  170. raise DeliveryError("invalid state JSON") from exc
  171. if not isinstance(parsed, dict):
  172. raise DeliveryError("state JSON object required")
  173. return parsed
  174. def state_write_json(descriptor: int, name: str, value: dict[str, Any]) -> None:
  175. """Atomic, fsync-backed replacement within a previously verified directory fd."""
  176. if "/" in name or name in {"", ".", ".."}:
  177. raise DeliveryError("invalid state file name")
  178. temporary = f".tmp-{os.getpid()}-{os.urandom(8).hex()}"
  179. try:
  180. handle = os.open(
  181. temporary,
  182. os.O_WRONLY
  183. | os.O_CREAT
  184. | os.O_EXCL
  185. | getattr(os, "O_NOFOLLOW", 0),
  186. 0o600,
  187. dir_fd=descriptor,
  188. )
  189. with os.fdopen(handle, "wb") as stream:
  190. stream.write(
  191. json.dumps(value, ensure_ascii=False, indent=2, sort_keys=True).encode()
  192. + b"\n"
  193. )
  194. stream.flush()
  195. os.fsync(stream.fileno())
  196. os.replace(temporary, name, src_dir_fd=descriptor, dst_dir_fd=descriptor)
  197. os.fsync(descriptor)
  198. except OSError as exc:
  199. raise DeliveryError("unsafe state write") from exc
  200. finally:
  201. with contextlib.suppress(FileNotFoundError):
  202. os.unlink(temporary, dir_fd=descriptor)
  203. def state_file(root: Path, name: str) -> Path:
  204. """Only for non-state output paths; controlled state uses descriptor helpers."""
  205. path = root / name
  206. if path.exists() and (path.is_symlink() or os.lstat(path).st_nlink != 1):
  207. raise DeliveryError("state file must not be linked")
  208. return path
  209. class Lock:
  210. def __init__(self, root: Path) -> None:
  211. self.root = root
  212. self.stream: Any | None = None
  213. self.descriptor: int | None = None
  214. def __enter__(self) -> Lock:
  215. _, self.descriptor = private_dir_fd(self.root)
  216. descriptor = os.open(
  217. ".lock",
  218. os.O_CREAT | os.O_RDWR | getattr(os, "O_NOFOLLOW", 0),
  219. 0o600,
  220. dir_fd=self.descriptor,
  221. )
  222. info = os.fstat(descriptor)
  223. if not stat.S_ISREG(info.st_mode) or info.st_nlink != 1:
  224. os.close(descriptor)
  225. os.close(self.descriptor)
  226. raise DeliveryError("state lock must be an unlinked regular file")
  227. self.stream = os.fdopen(descriptor, "r+")
  228. fcntl.flock(self.stream.fileno(), fcntl.LOCK_EX)
  229. return self
  230. def __exit__(self, *_: Any) -> None:
  231. assert self.stream
  232. fcntl.flock(self.stream.fileno(), fcntl.LOCK_UN)
  233. self.stream.close()
  234. assert self.descriptor is not None
  235. os.close(self.descriptor)
  236. def write_json(path: Path, value: dict[str, Any]) -> None:
  237. path.parent.mkdir(mode=0o700, parents=True, exist_ok=True)
  238. descriptor, temporary = tempfile.mkstemp(prefix=".tmp-", dir=path.parent)
  239. try:
  240. with os.fdopen(descriptor, "wb") as stream:
  241. stream.write(
  242. json.dumps(value, ensure_ascii=False, indent=2, sort_keys=True).encode()
  243. + b"\n"
  244. )
  245. stream.flush()
  246. os.fsync(stream.fileno())
  247. os.chmod(temporary, 0o600)
  248. os.replace(temporary, path)
  249. directory = os.open(path.parent, os.O_RDONLY)
  250. try:
  251. os.fsync(directory)
  252. finally:
  253. os.close(directory)
  254. finally:
  255. if os.path.exists(temporary):
  256. os.unlink(temporary)
  257. def evidence_time(value: dict[str, Any], label: str) -> None:
  258. try:
  259. issued = datetime.fromisoformat(str(value["issued_at"]).replace("Z", "+00:00"))
  260. expiry = datetime.fromisoformat(str(value["expires_at"]).replace("Z", "+00:00"))
  261. except (KeyError, ValueError, TypeError) as exc:
  262. raise DeliveryError(f"{label} evidence timestamps required") from exc
  263. now = datetime.now(UTC)
  264. if issued.astimezone(UTC) > now or expiry.astimezone(UTC) <= now or issued >= expiry:
  265. raise DeliveryError(f"{label} evidence expired or not yet valid")
  266. if value.get("provider") not in {"fake_local_only", "external"}:
  267. raise DeliveryError(f"{label} evidence provider required")
  268. def approval_time(value: dict[str, Any]) -> None:
  269. issued_text, expires_text = value.get("issued_at"), value.get("expires_at")
  270. if not isinstance(issued_text, str) or not isinstance(expires_text, str):
  271. raise DeliveryError("approval evidence timestamps required")
  272. rfc3339_utc = r"^\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}(?:\.\d{1,6})?Z$"
  273. if not re.fullmatch(rfc3339_utc, issued_text) or not re.fullmatch(
  274. rfc3339_utc, expires_text
  275. ):
  276. raise DeliveryError("approval evidence timestamps must be UTC RFC3339")
  277. try:
  278. issued = datetime.fromisoformat(issued_text.replace("Z", "+00:00"))
  279. expires = datetime.fromisoformat(expires_text.replace("Z", "+00:00"))
  280. except ValueError as exc:
  281. raise DeliveryError("approval evidence timestamps required") from exc
  282. now = datetime.now(UTC)
  283. if (
  284. issued > now + CLOCK_SKEW
  285. or expires <= now
  286. or issued >= expires
  287. or expires - issued > MAX_APPROVAL_WINDOW
  288. ):
  289. raise DeliveryError("approval evidence expired, not yet valid, or window exceeds limit")
  290. def validate_sbom(sbom: dict[str, Any], artifact_sha: str) -> None:
  291. if (
  292. sbom.get("bomFormat") != "CycloneDX"
  293. or not isinstance(sbom.get("specVersion"), str)
  294. or not str(sbom.get("serialNumber", "")).startswith("urn:uuid:")
  295. or not isinstance(sbom.get("version"), int)
  296. or not isinstance(sbom.get("components"), list)
  297. ):
  298. raise DeliveryError("CycloneDX SBOM schema minimum required")
  299. properties = sbom.get("metadata", {}).get("properties", [])
  300. if not isinstance(properties, list) or not any(
  301. isinstance(item, dict)
  302. and item.get("name") == "dataops:artifact-sha256"
  303. and item.get("value") == artifact_sha
  304. for item in properties
  305. ):
  306. raise DeliveryError("CycloneDX SBOM must bind actual artifact")
  307. def approved_signer(provider: str, key_id: str, test_only: bool) -> dict[str, Any]:
  308. store = read_json(TRUST_STORE)
  309. configured = store.get("providers", {}).get(provider)
  310. if not isinstance(configured, dict) or not configured.get("enabled"):
  311. raise DeliveryError("signing provider disabled or not pre-approved")
  312. if configured.get("test_only") and not test_only:
  313. raise DeliveryError("test-only signing provider requires explicit test-only mode")
  314. signer = configured.get("approved_signers", {}).get(key_id)
  315. if not isinstance(signer, dict):
  316. raise DeliveryError("signer key id is not in controlled trust store")
  317. if signer.get("algorithm") != "Ed25519" or not signer.get("trust_root"):
  318. raise DeliveryError("signer trust root is incomplete")
  319. try:
  320. public_key = bytes.fromhex(str(signer["public_key"]))
  321. except (KeyError, TypeError, ValueError) as exc:
  322. raise DeliveryError("controlled signer public key is invalid") from exc
  323. if len(public_key) != 32 or sha_bytes(public_key) != signer.get("public_key_sha256"):
  324. raise DeliveryError("controlled signer public key digest is invalid")
  325. return signer
  326. def approved_backup_encryption(provider: str, key_reference: str, test_only: bool) -> None:
  327. configured = read_json(TRUST_STORE).get("backup_encryption_providers", {}).get(provider)
  328. if not isinstance(configured, dict) or not configured.get("enabled"):
  329. raise DeliveryError("backup encryption provider disabled or not pre-approved")
  330. if configured.get("test_only") and not test_only:
  331. raise DeliveryError("test-only backup encryption requires explicit test-only mode")
  332. if key_reference not in configured.get("approved_key_references", []):
  333. raise DeliveryError("backup key reference is not pre-approved")
  334. def verify_controlled_signature(
  335. evidence: dict[str, Any], signed: dict[str, Any], label: str, test_only: bool
  336. ) -> None:
  337. if any(field in evidence for field in ("public_key", "trust_root", "approved_signers")):
  338. raise DeliveryError(f"{label} must not carry a caller supplied trust root")
  339. provider, key_id = evidence.get("provider"), evidence.get("key_id")
  340. if not isinstance(provider, str) or not isinstance(key_id, str):
  341. raise DeliveryError(f"{label} provider and key id required")
  342. if evidence.get("algorithm") != "Ed25519":
  343. raise DeliveryError(f"{label} algorithm must be Ed25519")
  344. try:
  345. signature = bytes.fromhex(str(evidence["signature"]))
  346. if len(signature) != 64:
  347. raise ValueError("Ed25519 signature length")
  348. signer = approved_signer(provider, key_id, test_only)
  349. Ed25519PublicKey.from_public_bytes(
  350. bytes.fromhex(str(signer["public_key"]))
  351. ).verify(signature, canonical(signed))
  352. except (KeyError, TypeError, ValueError) as exc:
  353. raise DeliveryError(f"invalid {label}") from exc
  354. except Exception as exc:
  355. raise DeliveryError(f"{label} verification failed") from exc
  356. def candidate_evidence(args: argparse.Namespace) -> dict[str, Any]:
  357. # Each caller supplied evidence file is opened, inode-checked and read once.
  358. raw = {name: safe_bytes(getattr(args, name)) for name in ("artifact", "sbom", "scan", "license", "signature")}
  359. artifact_sha = sha_bytes(raw["artifact"])
  360. sbom_sha = sha_bytes(raw["sbom"])
  361. sbom = read_json_bytes(raw["sbom"], "sbom")
  362. validate_sbom(sbom, artifact_sha)
  363. signature = read_json_bytes(raw["signature"], "signature")
  364. signature_sha = sha_bytes(raw["signature"])
  365. scan = read_json_bytes(raw["scan"], "scan")
  366. license_value = read_json_bytes(raw["license"], "license")
  367. if any(field in signature for field in ("public_key", "trust_root", "approved_signers")):
  368. raise DeliveryError("signature evidence must not carry a caller supplied trust root")
  369. evidence_time(scan, "scan")
  370. evidence_time(license_value, "license")
  371. provider = signature.get("provider")
  372. key_id = signature.get("key_id")
  373. if not isinstance(provider, str) or not isinstance(key_id, str):
  374. raise DeliveryError("signature provider and key id required")
  375. if (
  376. scan.get("artifact_sha256") != artifact_sha
  377. or scan.get("sbom_sha256") != sbom_sha
  378. or scan.get("signature_sha256") != signature_sha
  379. or scan.get("provider") != provider
  380. or scan.get("key_id") != key_id
  381. or scan.get("status") != "clear"
  382. ):
  383. raise DeliveryError("scan evidence does not clear actual artifact")
  384. if (
  385. license_value.get("artifact_sha256") != artifact_sha
  386. or license_value.get("sbom_sha256") != sbom_sha
  387. or license_value.get("signature_sha256") != signature_sha
  388. or license_value.get("provider") != provider
  389. or license_value.get("key_id") != key_id
  390. or license_value.get("status") != "clear"
  391. ):
  392. raise DeliveryError("license evidence does not clear actual SBOM")
  393. signed = {"artifact_sha256": artifact_sha, "sbom_sha256": sbom_sha}
  394. if (
  395. signature.get("algorithm") != "Ed25519"
  396. or signature.get("provider") != provider
  397. or signature.get("artifact_sha256") != artifact_sha
  398. or signature.get("sbom_sha256") != sbom_sha
  399. ):
  400. raise DeliveryError("signature evidence does not bind actual artifact and SBOM")
  401. try:
  402. evidence_time(signature, "signature")
  403. signer = approved_signer(provider, key_id, bool(args.test_only))
  404. public_key = bytes.fromhex(str(signer["public_key"]))
  405. Ed25519PublicKey.from_public_bytes(public_key).verify(
  406. bytes.fromhex(str(signature["signature"])), canonical(signed)
  407. )
  408. except (KeyError, ValueError, TypeError) as exc:
  409. raise DeliveryError("invalid signature evidence") from exc
  410. except Exception as exc:
  411. raise DeliveryError("signature evidence verification failed") from exc
  412. return {
  413. "artifact_sha256": artifact_sha,
  414. "sbom_sha256": sbom_sha,
  415. "scan_sha256": sha_bytes(raw["scan"]),
  416. "license_sha256": sha_bytes(raw["license"]),
  417. "signature_sha256": signature_sha,
  418. }
  419. def command_candidate(args: argparse.Namespace) -> dict[str, Any]:
  420. root = safe_state_root(args.state_dir)
  421. if args.environment not in ENVIRONMENTS or not args.approval_ref.strip():
  422. raise DeliveryError("known environment and approval reference required")
  423. evidence = candidate_evidence(args)
  424. value = {
  425. "environment": args.environment,
  426. "version": args.version,
  427. "approval_ref": args.approval_ref,
  428. "migration_head": HEAD,
  429. **evidence,
  430. }
  431. value["bundle_digest"] = sha_bytes(
  432. canonical(
  433. {
  434. "artifact_sha256": evidence["artifact_sha256"],
  435. "sbom_sha256": evidence["sbom_sha256"],
  436. "openapi_sha256": sha_file(ROOT / "docs/architecture/OPENAPI.yaml"),
  437. "config_schema_sha256": sha_file(
  438. ROOT / "deployment/helm/dataops-platform/values.schema.json"
  439. ),
  440. "migration_head": HEAD,
  441. }
  442. )
  443. )
  444. value["candidate_digest"] = sha_bytes(canonical(value))
  445. with Lock(root):
  446. _, candidates_fd = private_dir_fd(root, "candidates")
  447. try:
  448. filename = f"{value['candidate_digest']}.json"
  449. try:
  450. existing = state_read_json(candidates_fd, filename)
  451. except FileNotFoundError:
  452. existing = None
  453. if existing is not None and existing != value:
  454. raise DeliveryError("candidate collision")
  455. if existing is None:
  456. state_write_json(candidates_fd, filename, value)
  457. finally:
  458. os.close(candidates_fd)
  459. return {
  460. "status": "candidate_created",
  461. "candidate_digest": value["candidate_digest"],
  462. }
  463. def load_candidate(root: Path, digest: str) -> dict[str, Any]:
  464. require_digest(digest, "candidate_digest")
  465. _, descriptor = private_dir_fd(root, "candidates")
  466. try:
  467. return state_read_json(descriptor, f"{digest}.json")
  468. finally:
  469. os.close(descriptor)
  470. def current_record(root: Path, environment: str) -> dict[str, Any]:
  471. _, descriptor = private_dir_fd(root, "current")
  472. try:
  473. try:
  474. return state_read_json(descriptor, f"{environment}.json")
  475. except FileNotFoundError:
  476. return {"current_digest": "none", "bundle_digest": "none"}
  477. finally:
  478. os.close(descriptor)
  479. def write_current(root: Path, environment: str, value: dict[str, Any]) -> None:
  480. _, descriptor = private_dir_fd(root, "current")
  481. try:
  482. state_write_json(descriptor, f"{environment}.json", value)
  483. finally:
  484. os.close(descriptor)
  485. def append_history(root: Path, environment: str, value: dict[str, Any]) -> None:
  486. _, descriptor = private_dir_fd(root, "history", environment)
  487. try:
  488. state_write_json(descriptor, f"{value['current_digest']}.json", value)
  489. finally:
  490. os.close(descriptor)
  491. def command_promote(args: argparse.Namespace) -> dict[str, Any]:
  492. root = safe_state_root(args.state_dir)
  493. target_index = ENVIRONMENTS.index(args.environment)
  494. candidate = load_candidate(root, args.candidate_digest)
  495. if candidate.get("environment") != args.environment:
  496. raise DeliveryError("candidate environment mismatch")
  497. with Lock(root):
  498. current = current_record(root, args.environment)
  499. if current["current_digest"] == args.candidate_digest:
  500. return {
  501. "status": "promoted_replayed",
  502. "current_digest": args.candidate_digest,
  503. }
  504. if target_index:
  505. if current["current_digest"] != "none":
  506. raise DeliveryError("current digest CAS conflict")
  507. previous = current_record(root, ENVIRONMENTS[target_index - 1])
  508. if previous.get("current_digest") != args.expected_base_digest:
  509. raise DeliveryError(
  510. "promotion must use adjacent prior environment current digest"
  511. )
  512. if previous.get("bundle_digest") != candidate.get("bundle_digest"):
  513. raise DeliveryError(
  514. "promotion must retain the exact adjacent artifact bundle"
  515. )
  516. elif current["current_digest"] != args.expected_base_digest:
  517. raise DeliveryError("current digest CAS conflict")
  518. value = {
  519. "current_digest": args.candidate_digest,
  520. "previous_digest": args.expected_base_digest
  521. if target_index
  522. else current["current_digest"],
  523. "approval_ref": candidate["approval_ref"],
  524. "environment": args.environment,
  525. "artifact_sha256": candidate["artifact_sha256"],
  526. "bundle_digest": candidate["bundle_digest"],
  527. "previous_bundle_digest": previous.get("bundle_digest")
  528. if target_index
  529. else current.get("bundle_digest"),
  530. }
  531. write_current(root, args.environment, value)
  532. append_history(root, args.environment, value)
  533. return {"status": "promoted", "current_digest": args.candidate_digest}
  534. def command_rollback(args: argparse.Namespace) -> dict[str, Any]:
  535. root = safe_state_root(args.state_dir)
  536. with Lock(root):
  537. current = current_record(root, args.environment)
  538. if not args.approval_ref.strip() or args.target_digest != current.get(
  539. "previous_digest"
  540. ):
  541. raise DeliveryError(
  542. "rollback requires approval and exact recorded previous digest"
  543. )
  544. target = load_candidate(root, args.target_digest)
  545. if target.get("bundle_digest") != current.get("previous_bundle_digest", target.get("bundle_digest")):
  546. raise DeliveryError("rollback must retain recorded artifact bundle history")
  547. value = {
  548. **current,
  549. "current_digest": args.target_digest,
  550. "previous_digest": current["current_digest"],
  551. "previous_bundle_digest": current.get("bundle_digest"),
  552. "artifact_sha256": target["artifact_sha256"],
  553. "bundle_digest": target["bundle_digest"],
  554. "rollback_approval_ref": args.approval_ref,
  555. }
  556. write_current(root, args.environment, value)
  557. append_history(root, args.environment, value)
  558. return {"status": "rollback_recorded", "current_digest": args.target_digest}
  559. def command_backup(args: argparse.Namespace) -> dict[str, Any]:
  560. if args.environment == "production" and args.encryption_provider != "external":
  561. raise DeliveryError("production backup requires an approved external encryption provider")
  562. if args.encryption_provider not in {
  563. "fake",
  564. "external",
  565. "disabled",
  566. } or not args.key_reference.startswith(("kms://", "hsm://", "secretref://")):
  567. raise DeliveryError("safe provider and key reference required")
  568. approved_backup_encryption(
  569. args.encryption_provider, args.key_reference, bool(args.test_only)
  570. )
  571. if args.encryption_provider == "external" and (
  572. args.backup_signature is None or not args.signer_key_id
  573. ):
  574. raise DeliveryError("external backup requires a controlled signer key id and signature")
  575. if args.encryption_provider == "fake" and not args.test_only:
  576. raise DeliveryError("fake backup encryption requires explicit test-only mode")
  577. directory = safe_state_root(args.backup_dir)
  578. release = {
  579. "path": args.release_artifact.name,
  580. "sha256": sha_file(args.release_artifact),
  581. }
  582. entries = {
  583. "postgres": {**release, "mode": "logical_or_snapshot"},
  584. "neo4j": {**release, "mode": "snapshot"},
  585. "minio": {**release, "mode": "snapshot"},
  586. "config": {
  587. "path": args.config_artifact.name,
  588. "sha256": sha_file(args.config_artifact),
  589. "mode": "digest_only",
  590. },
  591. "key_references": {"reference": args.key_reference, "mode": "reference_only"},
  592. }
  593. manifest = {
  594. "environment": args.environment,
  595. "components": entries,
  596. "encryption_provider": args.encryption_provider,
  597. "key_reference": args.key_reference,
  598. "migration_head": HEAD,
  599. "signature_provider": "external"
  600. if args.encryption_provider == "external"
  601. else "fake_local_only",
  602. "signer_key_id": args.signer_key_id
  603. if args.encryption_provider == "external"
  604. else "wp07-local-test",
  605. }
  606. manifest["manifest_sha256"] = sha_bytes(canonical(manifest))
  607. signature_payload = {
  608. "manifest_sha256": manifest["manifest_sha256"],
  609. "encryption_provider": manifest["encryption_provider"],
  610. "key_reference": manifest["key_reference"],
  611. }
  612. if args.encryption_provider == "external":
  613. supplied = read_json(args.backup_signature)
  614. if (
  615. supplied.get("provider") != "external"
  616. or supplied.get("key_id") != args.signer_key_id
  617. ):
  618. raise DeliveryError("external backup signature signer mismatch")
  619. verify_controlled_signature(supplied, signature_payload, "backup signature", False)
  620. manifest["integrity_signature"] = {
  621. "algorithm": "Ed25519",
  622. "provider": "external",
  623. "key_id": args.signer_key_id,
  624. "signature": supplied["signature"],
  625. }
  626. else:
  627. signer = approved_signer("fake_local_only", "wp07-local-test", True)
  628. private = Ed25519PrivateKey.from_private_bytes(LOCAL_TEST_PRIVATE_KEY)
  629. manifest["integrity_signature"] = {
  630. "algorithm": "Ed25519",
  631. "provider": "fake_local_only",
  632. "key_id": "wp07-local-test",
  633. "signature": private.sign(canonical(signature_payload)).hex(),
  634. }
  635. if signer.get("public_key") != private.public_key().public_bytes(
  636. serialization.Encoding.Raw, serialization.PublicFormat.Raw
  637. ).hex():
  638. raise DeliveryError("local test signer does not match controlled trust store")
  639. path = state_file(directory, "backup-manifest.json")
  640. write_json(path, manifest)
  641. _, backup_fd = private_dir_fd(directory, "backup")
  642. try:
  643. state_write_json(backup_fd, "latest.json", manifest)
  644. finally:
  645. os.close(backup_fd)
  646. return {"status": "planned", "manifest_path": str(path)}
  647. def command_restore(args: argparse.Namespace) -> dict[str, Any]:
  648. value = read_json(args.backup_manifest)
  649. check = {
  650. key: value[key]
  651. for key in value
  652. if key not in {"manifest_sha256", "integrity_signature"}
  653. }
  654. if value.get("manifest_sha256") != sha_bytes(canonical(check)):
  655. raise DeliveryError("backup manifest digest mismatch")
  656. try:
  657. signature = value["integrity_signature"]
  658. if not isinstance(signature, dict) or signature.get("algorithm") != "Ed25519":
  659. raise DeliveryError("backup signature algorithm invalid")
  660. if (
  661. signature.get("provider") != value.get("signature_provider")
  662. or signature.get("key_id") != value.get("signer_key_id")
  663. ):
  664. raise DeliveryError("backup signature signer identity mismatch")
  665. if (
  666. value.get("encryption_provider") == "external"
  667. and value.get("signature_provider") != "external"
  668. ) or (
  669. value.get("encryption_provider") == "fake"
  670. and value.get("signature_provider") != "fake_local_only"
  671. ):
  672. raise DeliveryError("backup encryption and signer provider mismatch")
  673. signature_payload = {
  674. "manifest_sha256": value["manifest_sha256"],
  675. "encryption_provider": value["encryption_provider"],
  676. "key_reference": value["key_reference"],
  677. }
  678. verify_controlled_signature(
  679. signature, signature_payload, "backup signature", bool(args.test_only)
  680. )
  681. except DeliveryError:
  682. raise
  683. except (KeyError, TypeError, ValueError) as exc:
  684. raise DeliveryError("backup signature invalid") from exc
  685. except Exception as exc:
  686. raise DeliveryError("backup signature verification failed") from exc
  687. if args.target == "production" or args.target != args.confirm_target:
  688. raise DeliveryError("restore requires exact new non-production target")
  689. if args.environment == "production" and value.get("encryption_provider") != "external":
  690. raise DeliveryError("production restore requires an externally encrypted backup")
  691. approved_backup_encryption(
  692. str(value.get("encryption_provider")), str(value.get("key_reference")), bool(args.test_only)
  693. )
  694. components = value.get("components", {})
  695. if set(components) != {"postgres", "neo4j", "minio", "config", "key_references"}:
  696. raise DeliveryError("backup component set incomplete")
  697. parent = args.backup_manifest.parent
  698. for name, item in components.items():
  699. if name == "key_references":
  700. if not str(item.get("reference", "")).startswith(
  701. ("kms://", "hsm://", "secretref://")
  702. ):
  703. raise DeliveryError("backup key reference invalid")
  704. continue
  705. path = parent / str(item["path"])
  706. if sha_file(path) != item["sha256"]:
  707. raise DeliveryError("backup component checksum mismatch")
  708. if value.get("migration_head") != HEAD:
  709. raise DeliveryError("backup compatibility migration mismatch")
  710. restore_root = safe_state_root(args.backup_manifest.parent)
  711. with Lock(restore_root):
  712. _, restore_fd = private_dir_fd(restore_root, "restore")
  713. try:
  714. try:
  715. state_read_json(restore_fd, f"{args.target}.json")
  716. except FileNotFoundError:
  717. pass
  718. else:
  719. raise DeliveryError("restore target has already been used")
  720. state_write_json(
  721. restore_fd,
  722. f"{args.target}.json",
  723. {
  724. "target": args.target,
  725. "manifest_sha256": value["manifest_sha256"],
  726. "status": "restore_to_new_target_planned",
  727. },
  728. )
  729. finally:
  730. os.close(restore_fd)
  731. return {"status": "restore_to_new_target_planned", "target": args.target}
  732. def command_compatibility(args: argparse.Namespace) -> dict[str, Any]:
  733. def release(artifact_path: Path, sbom_path: Path, label: str) -> dict[str, Any]:
  734. artifact = safe_bytes(artifact_path)
  735. sbom_bytes = safe_bytes(sbom_path)
  736. artifact_sha = sha_bytes(artifact)
  737. sbom = read_json_bytes(sbom_bytes, f"{label} SBOM")
  738. validate_sbom(sbom, artifact_sha)
  739. components = sbom.get("components", [])
  740. identity = [
  741. {
  742. "bom-ref": component.get("bom-ref", ""),
  743. "name": component.get("name", ""),
  744. "version": component.get("version", ""),
  745. "purl": component.get("purl", ""),
  746. }
  747. for component in components
  748. if isinstance(component, dict)
  749. ]
  750. if len(identity) != len(components):
  751. raise DeliveryError(f"{label} SBOM component identity is malformed")
  752. identity.sort(key=canonical)
  753. return {
  754. "artifact_sha256": artifact_sha,
  755. "sbom_sha256": sha_bytes(sbom_bytes),
  756. "sbom_serial": sbom.get("serialNumber"),
  757. "component_identity_sha256": sha_bytes(canonical(identity)),
  758. }
  759. releases = {
  760. "base": release(args.base_artifact, args.base_sbom, "base"),
  761. "candidate": release(args.candidate_artifact, args.candidate_sbom, "candidate"),
  762. "rollback": release(args.rollback_artifact, args.rollback_sbom, "rollback"),
  763. }
  764. def contract(path: Path | None, label: str, bound_release: dict[str, Any]) -> dict[str, Any]:
  765. if path is None:
  766. raise DeliveryError("base, candidate and rollback contract manifests are required")
  767. supplied = read_json(path)
  768. source: dict[str, Path] = {}
  769. for field in ("openapi", "config", "chart", "values"):
  770. value = supplied.get(field)
  771. if not isinstance(value, str) or not value:
  772. raise DeliveryError(f"{label} compatibility contract {field} is required")
  773. source[field] = Path(value)
  774. if (
  775. supplied.get("artifact_sha256") != bound_release["artifact_sha256"]
  776. or supplied.get("sbom_sha256") != bound_release["sbom_sha256"]
  777. or supplied.get("component_identity_sha256")
  778. != bound_release["component_identity_sha256"]
  779. ):
  780. raise DeliveryError(f"{label} contract does not bind actual artifact/SBOM identity")
  781. chain = supplied.get("migration_chain")
  782. head = supplied.get("migration_head")
  783. if (
  784. not isinstance(head, str)
  785. or not isinstance(chain, list)
  786. or not chain
  787. or not all(isinstance(item, str) and item for item in chain)
  788. or len(set(chain)) != len(chain)
  789. or chain[-1] != head
  790. ):
  791. raise DeliveryError(f"{label} migration chain and head are required")
  792. raw = {field: safe_bytes(item) for field, item in source.items()}
  793. config = read_json_bytes(raw["config"], f"{label} config schema")
  794. required = config.get("required", [])
  795. if not isinstance(required, list) or not all(isinstance(item, str) for item in required):
  796. raise DeliveryError(f"{label} config required list is invalid")
  797. value = {
  798. "openapi_sha256": sha_bytes(raw["openapi"]),
  799. "config_schema_sha256": sha_bytes(raw["config"]),
  800. "chart_sha256": sha_bytes(raw["chart"]),
  801. "values_sha256": sha_bytes(raw["values"]),
  802. "migration_head": head,
  803. "migration_chain": chain,
  804. "openapi_operations": sorted(
  805. re.findall(r"^ (/[^:\s]+):", raw["openapi"].decode(errors="strict"), re.MULTILINE)
  806. ),
  807. "config_required": sorted(required),
  808. }
  809. value["contract_sha256"] = sha_bytes(canonical(value))
  810. return value
  811. contracts = {
  812. "base": contract(args.base_contract, "base", releases["base"]),
  813. "candidate": contract(args.candidate_contract, "candidate", releases["candidate"]),
  814. "rollback": contract(args.rollback_contract, "rollback", releases["rollback"]),
  815. }
  816. base, candidate_contract, rollback_contract = (
  817. contracts["base"], contracts["candidate"], contracts["rollback"]
  818. )
  819. if not set(base["openapi_operations"]).issubset(candidate_contract["openapi_operations"]):
  820. raise DeliveryError("candidate OpenAPI removes a base operation")
  821. if not set(candidate_contract["config_required"]).issubset(base["config_required"]):
  822. raise DeliveryError("candidate config introduces a newly required setting")
  823. if candidate_contract["migration_chain"][: len(base["migration_chain"])] != base["migration_chain"]:
  824. raise DeliveryError("candidate migration chain is not an ordered extension of base")
  825. if (
  826. base["chart_sha256"] != candidate_contract["chart_sha256"]
  827. or base["values_sha256"] != candidate_contract["values_sha256"]
  828. or rollback_contract != base
  829. ):
  830. raise DeliveryError("base/candidate/rollback deployment contract is incompatible")
  831. if releases["rollback"] != releases["base"]:
  832. raise DeliveryError("rollback must restore the exact tested base artifact and SBOM")
  833. if releases["candidate"] != releases["base"]:
  834. if args.compatibility_approval is None:
  835. raise DeliveryError("changed artifact/SBOM requires a bound compatibility approval")
  836. approval = read_json(args.compatibility_approval)
  837. approval_time(approval)
  838. approval_payload = {
  839. "base": releases["base"],
  840. "candidate": releases["candidate"],
  841. "base_contract_sha256": base["contract_sha256"],
  842. "candidate_contract_sha256": candidate_contract["contract_sha256"],
  843. "approval_ref": approval.get("approval_ref"),
  844. "issued_at": approval.get("issued_at"),
  845. "expires_at": approval.get("expires_at"),
  846. }
  847. if (
  848. approval.get("base") != releases["base"]
  849. or approval.get("candidate") != releases["candidate"]
  850. or approval.get("base_contract_sha256") != base["contract_sha256"]
  851. or approval.get("candidate_contract_sha256")
  852. != candidate_contract["contract_sha256"]
  853. or not str(approval.get("approval_ref", "")).strip()
  854. ):
  855. raise DeliveryError("compatibility approval does not bind actual release pair")
  856. verify_controlled_signature(
  857. approval, approval_payload, "compatibility approval", False
  858. )
  859. value = {
  860. "base_version": args.base_version,
  861. "candidate_version": args.candidate_version,
  862. "migration_head": candidate_contract["migration_head"],
  863. "contracts": contracts,
  864. "openapi_sha256": candidate_contract["openapi_sha256"],
  865. "config_schema_sha256": candidate_contract["config_schema_sha256"],
  866. "chart_sha256": candidate_contract["chart_sha256"],
  867. "releases": releases,
  868. "comparison": {
  869. "base_to_candidate": "verified",
  870. "candidate_to_rollback": "verified_exact_base",
  871. "openapi": "base_operations_retained",
  872. "config_schema": "no_new_required_setting",
  873. "chart_values": "base_candidate_equal",
  874. "migration": "ancestry_and_rollback_checked",
  875. },
  876. "rollback": releases["rollback"],
  877. }
  878. value["matrix_sha256"] = sha_bytes(canonical(value))
  879. write_json(args.output, value)
  880. return {"status": "compatible", "matrix_sha256": value["matrix_sha256"]}
  881. def parser() -> argparse.ArgumentParser:
  882. parser = argparse.ArgumentParser()
  883. sub = parser.add_subparsers(dest="command", required=True)
  884. c = sub.add_parser("candidate")
  885. c.add_argument("--state-dir", type=Path, required=True)
  886. c.add_argument("--environment", required=True)
  887. c.add_argument("--version", required=True)
  888. c.add_argument("--artifact", type=Path, required=True)
  889. c.add_argument("--sbom", type=Path, required=True)
  890. c.add_argument("--scan", type=Path, required=True)
  891. c.add_argument("--license", type=Path, required=True)
  892. c.add_argument("--signature", type=Path, required=True)
  893. c.add_argument("--approval-ref", required=True)
  894. c.add_argument("--test-only", action="store_true")
  895. p = sub.add_parser("promote")
  896. p.add_argument("--state-dir", type=Path, required=True)
  897. p.add_argument("--environment", required=True)
  898. p.add_argument("--candidate-digest", required=True)
  899. p.add_argument("--expected-base-digest", required=True)
  900. r = sub.add_parser("rollback")
  901. r.add_argument("--state-dir", type=Path, required=True)
  902. r.add_argument("--environment", required=True)
  903. r.add_argument("--target-digest", required=True)
  904. r.add_argument("--approval-ref", required=True)
  905. b = sub.add_parser("backup-plan")
  906. b.add_argument("--backup-dir", type=Path, required=True)
  907. b.add_argument("--release-artifact", type=Path, required=True)
  908. b.add_argument("--config-artifact", type=Path, required=True)
  909. b.add_argument("--key-reference", required=True)
  910. b.add_argument("--encryption-provider", required=True)
  911. b.add_argument("--environment", default="test")
  912. b.add_argument("--signer-key-id")
  913. b.add_argument("--backup-signature", type=Path)
  914. b.add_argument("--test-only", action="store_true")
  915. rs = sub.add_parser("restore-plan")
  916. rs.add_argument("--backup-manifest", type=Path, required=True)
  917. rs.add_argument("--target", required=True)
  918. rs.add_argument("--confirm-target", required=True)
  919. rs.add_argument("--environment", required=True)
  920. rs.add_argument("--test-only", action="store_true")
  921. cm = sub.add_parser("compatibility")
  922. cm.add_argument("--output", type=Path, required=True)
  923. cm.add_argument("--base-version", required=True)
  924. cm.add_argument("--candidate-version", required=True)
  925. cm.add_argument("--base-artifact", type=Path, required=True)
  926. cm.add_argument("--base-sbom", type=Path, required=True)
  927. cm.add_argument("--candidate-artifact", type=Path, required=True)
  928. cm.add_argument("--candidate-sbom", type=Path, required=True)
  929. cm.add_argument("--rollback-artifact", type=Path, required=True)
  930. cm.add_argument("--rollback-sbom", type=Path, required=True)
  931. cm.add_argument("--compatibility-approval", type=Path)
  932. cm.add_argument("--base-contract", type=Path)
  933. cm.add_argument("--candidate-contract", type=Path)
  934. cm.add_argument("--rollback-contract", type=Path)
  935. return parser
  936. def main() -> int:
  937. args = parser().parse_args()
  938. try:
  939. result = {
  940. "candidate": command_candidate,
  941. "promote": command_promote,
  942. "rollback": command_rollback,
  943. "backup-plan": command_backup,
  944. "restore-plan": command_restore,
  945. "compatibility": command_compatibility,
  946. }[args.command](args)
  947. print(json.dumps(result, sort_keys=True))
  948. return 0
  949. except DeliveryError as exc:
  950. print(json.dumps({"status": "rejected", "error": str(exc)}, sort_keys=True))
  951. return 2
  952. if __name__ == "__main__":
  953. raise SystemExit(main())