"""Durable state for guided, explicitly executed release runs. The run store records an immutable selective-plan snapshot and a small mutable state machine. It does not execute release commands. Existing release executors remain responsible for their narrow confirmations and will be wired to this state in a later slice. """ from __future__ import annotations from copy import deepcopy from datetime import UTC, datetime import hashlib import json import os from pathlib import Path import re import stat import tempfile import threading from typing import Any import uuid from filelock import FileLock, Timeout as FileLockTimeout SCHEMA = "govoplan.release-run" SCHEMA_VERSION = 2 MAX_RECORD_BYTES = 2 * 1024 * 1024 MAX_PLAN_STEPS = 512 MAX_REPOSITORY_VERSIONS = 128 MAX_EVENTS = 256 MAX_COMMANDS = 2_048 MAX_LIST_LIMIT = 100 RUN_ID_PATTERN = re.compile(r"^rr-[0-9]{8}T[0-9]{6}Z-[0-9a-f]{12}$") STEP_ID_PATTERN = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._:@-]{0,159}$") REQUEST_ID_PATTERN = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._:@/-]{7,127}$") CHANNEL_PATTERN = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._-]{0,63}$") REPOSITORY_PATTERN = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._-]{0,127}$") WORKSPACE_FINGERPRINT_PATTERN = re.compile(r"^[0-9a-f]{64}$") VERSION_PATTERN = re.compile( r"^[0-9]+\.[0-9]+\.[0-9]+(?:[-+][0-9A-Za-z][0-9A-Za-z.-]*)?$" ) STEP_STATES = {"pending", "running", "succeeded", "failed", "interrupted"} PLAN_STATUSES = {"ready", "attention", "blocked"} PLANNED_STEP_STATUSES = {"planned", "needs-executor", "needs-version-metadata"} EVENT_TYPES = { "run_created", "run_resumed", "step_started", "step_succeeded", "step_failed", "step_interrupted", "step_reconciled", "step_retry_requested", } ATTEMPT_OUTCOMES = {"running", "succeeded", "failed", "interrupted"} RECONCILIATION_OUTCOMES = {"effect_absent", "effect_succeeded", "unresolved"} RECONCILIATION_CONFIRMATION = "RECONCILE" class ReleaseRunError(RuntimeError): """Base error for durable release run state.""" class ReleaseRunNotFound(ReleaseRunError): """Raised when a run does not exist.""" class ReleaseRunWorkspaceMismatch(ReleaseRunNotFound): """Raised when a valid record belongs to another workspace.""" class ReleaseRunConflict(ReleaseRunError): """Raised when a state transition is unsafe or no longer applicable.""" class ReleaseRunCorrupt(ReleaseRunError): """Raised when a persisted record fails its integrity or schema checks.""" class ReleaseRunStore: """Atomic, bounded local JSON store for release-run state.""" def __init__(self, root: Path, *, expected_workspace_fingerprint: str) -> None: # Keep the configured path spelling so a symlink root (or symlinked # parent) remains detectable instead of being erased by resolve(). self.root = Path(os.path.abspath(os.fspath(root.expanduser()))) if not WORKSPACE_FINGERPRINT_PATTERN.fullmatch( expected_workspace_fingerprint ): raise ReleaseRunConflict("Release workspace fingerprint is invalid.") self.expected_workspace_fingerprint = expected_workspace_fingerprint self._thread_lock = threading.RLock() self._lock_path = self.root / ".release-runs.lock" def create( self, *, input_snapshot: dict[str, Any], plan_snapshot: dict[str, Any], ) -> dict[str, Any]: immutable_input = _validated_input(input_snapshot) if ( immutable_input["workspace_fingerprint"] != self.expected_workspace_fingerprint ): raise ReleaseRunConflict( "Release run input belongs to a different workspace." ) immutable_plan = _validated_plan(plan_snapshot) now = _timestamp() run_id = ( f"rr-{now.replace('-', '').replace(':', '')[:15]}Z-{uuid.uuid4().hex[:12]}" ) immutable = { "input": immutable_input, "plan": immutable_plan, } immutable["digest"] = _immutable_digest(immutable) steps = [ { "id": item["id"], "state": "pending", "attempt_count": 0, "attempt_fingerprint": None, "started_at": None, "finished_at": None, "result_code": None, } for item in immutable_plan["dry_run_steps"] ] record: dict[str, Any] = { "schema": SCHEMA, "schema_version": SCHEMA_VERSION, "run_id": run_id, "created_at": now, "updated_at": now, "immutable": immutable, "state": { "status": "planned", "steps": steps, "events": [], "next_event_sequence": 1, "processed_requests": [], "processed_attempts": [], }, } _append_event(record, event_type="run_created") record["updated_at"] = record["state"]["events"][-1]["at"] _refresh_status(record) record["record_digest"] = _record_digest(record) with self._locked(): self._ensure_root() self._write_new(record) return _view(record) def get(self, run_id: str) -> dict[str, Any]: with self._locked(): record = self._read(run_id) return _view(record) def list(self, *, limit: int = 20) -> list[dict[str, Any]]: limit = min(max(int(limit), 1), MAX_LIST_LIMIT) with self._locked(): self._ensure_root() paths = sorted( ( path for path in self.root.iterdir() if path.name.endswith(".json") and RUN_ID_PATTERN.fullmatch(path.stem) ), key=lambda item: item.name, reverse=True, ) summaries: list[dict[str, Any]] = [] for path in paths: try: record = self._read(path.stem) except ReleaseRunWorkspaceMismatch: continue except (ReleaseRunCorrupt, ReleaseRunNotFound): summaries.append( { "run_id": path.stem, "status": "unavailable", "error": "Run record is unavailable or failed integrity validation.", } ) if len(summaries) >= limit: break continue view = _view(record) immutable_input = view["immutable"]["input"] summaries.append( { "run_id": view["run_id"], "created_at": view["created_at"], "updated_at": view["updated_at"], "status": view["state"]["status"], "channel": immutable_input["channel"], "repo_versions": immutable_input["repo_versions"], "recommended_next": view["recommended_next"], } ) if len(summaries) >= limit: break return summaries def resume(self, run_id: str, *, request_id: str) -> dict[str, Any]: request_fingerprint = _request_fingerprint(request_id) with self._locked(): record = self._read(run_id) if _request_seen(record, request_fingerprint, "resume", None, None): return _view(record) _ensure_command_capacity(record) for step in record["state"]["steps"]: if step["state"] != "running": continue attempt = _attempt_for_fingerprint( record, step["attempt_fingerprint"], step_id=step["id"] ) attempt["outcome"] = "interrupted" attempt["result_code"] = "process_interrupted" step["state"] = "interrupted" step["finished_at"] = _timestamp() step["result_code"] = "process_interrupted" _append_event( record, event_type="step_interrupted", step_id=step["id"], result_code="process_interrupted", ) _remember_request(record, request_fingerprint, "resume", None, None) _append_event(record, event_type="run_resumed") self._persist_update(record) return _view(record) def retry_step( self, run_id: str, step_id: str, *, request_id: str, ) -> dict[str, Any]: _validate_step_id(step_id) request_fingerprint = _request_fingerprint(request_id) with self._locked(): record = self._read(run_id) if _request_seen(record, request_fingerprint, "retry", step_id, None): return _view(record) step, plan_step = _step_pair(record, step_id) if step["state"] not in {"failed", "interrupted"}: raise ReleaseRunConflict( "Only a failed or interrupted release step can be prepared for retry." ) if step["state"] == "interrupted" and bool(plan_step.get("mutating")): raise ReleaseRunConflict( "Record the observed effect of an interrupted mutating step before retry." ) _ensure_command_capacity(record) retry_reason = ( "known_failure" if step["state"] == "failed" else "read_only_interruption" ) step.update( { "state": "pending", "attempt_fingerprint": None, "started_at": None, "finished_at": None, "result_code": None, } ) _remember_request(record, request_fingerprint, "retry", step_id, None) _append_event( record, event_type="step_retry_requested", step_id=step_id, result_code=retry_reason, ) self._persist_update(record) return _view(record) def reconcile_step( self, run_id: str, step_id: str, *, request_id: str, outcome: str, confirmation: str, ) -> dict[str, Any]: """Record an operator-observed outcome for an interrupted mutation.""" _validate_step_id(step_id) if outcome not in RECONCILIATION_OUTCOMES: raise ReleaseRunConflict("Release reconciliation outcome is invalid.") if confirmation != RECONCILIATION_CONFIRMATION: raise ReleaseRunConflict( f"Type {RECONCILIATION_CONFIRMATION} to record reconciliation." ) request_fingerprint = _request_fingerprint(request_id) with self._locked(): record = self._read(run_id) if _request_seen( record, request_fingerprint, "reconcile", step_id, outcome ): return _view(record) step, plan_step = _step_pair(record, step_id) if step["state"] != "interrupted" or not bool( plan_step.get("mutating") ): raise ReleaseRunConflict( "Only an interrupted mutating step can be reconciled." ) _ensure_command_capacity(record) if outcome == "effect_absent": step.update( { "state": "pending", "attempt_fingerprint": None, "started_at": None, "finished_at": None, "result_code": None, } ) elif outcome == "effect_succeeded": step.update( { "state": "succeeded", "finished_at": _timestamp(), "result_code": "reconciled_effect_succeeded", } ) _remember_request( record, request_fingerprint, "reconcile", step_id, outcome ) _append_event( record, event_type="step_reconciled", step_id=step_id, result_code=outcome, ) self._persist_update(record) return _view(record) def start_step( self, run_id: str, step_id: str, *, attempt_id: str, ) -> dict[str, Any]: """Claim one step for a future confirmed executor. This is deliberately not an HTTP endpoint in the foundation slice. """ _validate_step_id(step_id) attempt_fingerprint = _request_fingerprint(attempt_id) with self._locked(): record = self._read(run_id) existing_attempt = _find_attempt(record, attempt_fingerprint) if existing_attempt is not None: if existing_attempt["step_id"] != step_id: raise ReleaseRunConflict( "The attempt identifier was already used for another release step." ) return _view(record) step, _plan_step = _step_pair(record, step_id) available, reason = _step_available(record, step_id) if not available: raise ReleaseRunConflict(reason) _ensure_attempt_capacity(record) step.update( { "state": "running", "attempt_count": step["attempt_count"] + 1, "attempt_fingerprint": attempt_fingerprint, "started_at": _timestamp(), "finished_at": None, "result_code": None, } ) record["state"]["processed_attempts"].append( { "fingerprint": attempt_fingerprint, "step_id": step_id, "outcome": "running", "result_code": None, } ) _append_event(record, event_type="step_started", step_id=step_id) self._persist_update(record) return _view(record) def finish_step( self, run_id: str, step_id: str, *, attempt_id: str, succeeded: bool, result_code: str, ) -> dict[str, Any]: """Record a future executor outcome without persisting its output.""" _validate_step_id(step_id) attempt_fingerprint = _request_fingerprint(attempt_id) result_code = _safe_code(result_code) target_state = "succeeded" if succeeded else "failed" with self._locked(): record = self._read(run_id) step, _plan_step = _step_pair(record, step_id) attempt = _find_attempt(record, attempt_fingerprint) if attempt is None or attempt["step_id"] != step_id: raise ReleaseRunConflict( "The release attempt identifier was not claimed for this step." ) if attempt["outcome"] != "running": if ( attempt["outcome"] == target_state and attempt["result_code"] == result_code ): return _view(record) raise ReleaseRunConflict( "The release attempt already has a different recorded outcome." ) if ( step["state"] != "running" or step["attempt_fingerprint"] != attempt_fingerprint ): raise ReleaseRunConflict( "The release step is not running with this attempt identifier." ) step.update( { "state": target_state, "finished_at": _timestamp(), "result_code": result_code, } ) attempt["outcome"] = target_state attempt["result_code"] = result_code _append_event( record, event_type="step_succeeded" if succeeded else "step_failed", step_id=step_id, result_code=result_code, ) self._persist_update(record) return _view(record) def _persist_update(self, record: dict[str, Any]) -> None: self._assert_workspace(record) record["updated_at"] = _timestamp() _refresh_status(record) record["record_digest"] = _record_digest(record) _validate_record(record) self._atomic_write(self._path(record["run_id"]), record) def _write_new(self, record: dict[str, Any]) -> None: _validate_record(record) path = self._path(record["run_id"]) if path.exists() or path.is_symlink(): raise ReleaseRunConflict("The generated release run already exists.") self._atomic_write(path, record) def _read(self, run_id: str) -> dict[str, Any]: path = self._path(run_id) flags = os.O_RDONLY if hasattr(os, "O_NOFOLLOW"): flags |= os.O_NOFOLLOW try: descriptor = os.open(path, flags) except FileNotFoundError as exc: raise ReleaseRunNotFound("Release run was not found.") from exc except OSError as exc: raise ReleaseRunCorrupt( "Release run record could not be opened safely." ) from exc try: metadata = os.fstat(descriptor) if not stat.S_ISREG(metadata.st_mode): raise ReleaseRunCorrupt("Release run record is not a regular file.") if hasattr(os, "geteuid") and metadata.st_uid != os.geteuid(): raise ReleaseRunCorrupt( "Release run record is not owned by the console operator." ) if metadata.st_mode & 0o077: raise ReleaseRunCorrupt( "Release run record permissions are broader than 0600." ) if metadata.st_size > MAX_RECORD_BYTES: raise ReleaseRunCorrupt("Release run record exceeds its size limit.") with os.fdopen(descriptor, "rb") as handle: descriptor = -1 encoded = handle.read(MAX_RECORD_BYTES + 1) if len(encoded) > MAX_RECORD_BYTES: raise ReleaseRunCorrupt("Release run record exceeds its size limit.") payload = json.loads(encoded.decode("utf-8")) except (OSError, UnicodeDecodeError, json.JSONDecodeError) as exc: raise ReleaseRunCorrupt("Release run record is unreadable.") from exc finally: if descriptor >= 0: os.close(descriptor) _validate_record(payload, expected_run_id=run_id) self._assert_workspace(payload) return payload def _assert_workspace(self, record: dict[str, Any]) -> None: if ( record["immutable"]["input"]["workspace_fingerprint"] != self.expected_workspace_fingerprint ): raise ReleaseRunWorkspaceMismatch("Release run was not found.") def _path(self, run_id: str) -> Path: if not RUN_ID_PATTERN.fullmatch(run_id): raise ReleaseRunNotFound("Release run was not found.") return self.root / f"{run_id}.json" def _ensure_root(self) -> None: _reject_symlink_components(self.root) _create_private_directory_tree(self.root) _reject_symlink_components(self.root) if self.root.is_symlink() or not self.root.is_dir(): raise ReleaseRunCorrupt("Release run storage root is not a directory.") root_metadata = self.root.stat() if hasattr(os, "geteuid") and root_metadata.st_uid != os.geteuid(): raise ReleaseRunCorrupt( "Release run storage is not owned by the console operator." ) try: self.root.chmod(0o700) except OSError as exc: raise ReleaseRunCorrupt( "Release run storage permissions could not be secured." ) from exc if self.root.stat().st_mode & 0o077: raise ReleaseRunCorrupt( "Release run storage permissions are broader than 0700." ) _fsync_directory(self.root) def _atomic_write(self, path: Path, record: dict[str, Any]) -> None: encoded = _canonical_json(record) + b"\n" if len(encoded) > MAX_RECORD_BYTES: raise ReleaseRunConflict("Release run record exceeds its size limit.") descriptor = -1 temporary_path: str | None = None try: descriptor, temporary_path = tempfile.mkstemp( dir=self.root, prefix=".release-run-", suffix=".tmp" ) os.fchmod(descriptor, 0o600) with os.fdopen(descriptor, "wb") as handle: descriptor = -1 handle.write(encoded) handle.flush() os.fsync(handle.fileno()) os.replace(temporary_path, path) temporary_path = None written_metadata = path.lstat() if ( not stat.S_ISREG(written_metadata.st_mode) or written_metadata.st_mode & 0o077 or ( hasattr(os, "geteuid") and written_metadata.st_uid != os.geteuid() ) ): raise ReleaseRunCorrupt( "Release run record permissions could not be secured." ) _fsync_directory(self.root) finally: if descriptor >= 0: os.close(descriptor) if temporary_path is not None: try: os.unlink(temporary_path) except FileNotFoundError: pass def _locked(self): # type: ignore[no-untyped-def] self._ensure_root() return _CombinedLock(self._thread_lock, FileLock(str(self._lock_path))) class _CombinedLock: def __init__(self, thread_lock: threading.RLock, file_lock: FileLock) -> None: self.thread_lock = thread_lock self.file_lock = file_lock def __enter__(self) -> None: self.thread_lock.acquire() try: self.file_lock.acquire(timeout=10) except FileLockTimeout as exc: self.thread_lock.release() raise ReleaseRunConflict( "Release run storage is busy; retry the operation." ) from exc except BaseException: self.thread_lock.release() raise def __exit__(self, *exc_info: object) -> None: try: self.file_lock.release() finally: self.thread_lock.release() def _validated_input(value: dict[str, Any]) -> dict[str, Any]: if not isinstance(value, dict): raise ReleaseRunConflict("Release run input must be an object.") allowed = { "channel", "repo_versions", "online", "remote_tags", "public_catalog", "include_migrations", "workspace_fingerprint", } if set(value) != allowed: raise ReleaseRunConflict("Release run input fields are invalid.") channel = value.get("channel") repo_versions = value.get("repo_versions") if not isinstance(channel, str) or not CHANNEL_PATTERN.fullmatch(channel): raise ReleaseRunConflict("Release channel is invalid.") if ( not isinstance(repo_versions, dict) or not repo_versions or len(repo_versions) > MAX_REPOSITORY_VERSIONS ): raise ReleaseRunConflict("At least one repository version is required.") normalized_versions: dict[str, str] = {} for repo, version in sorted(repo_versions.items()): if not isinstance(repo, str) or not REPOSITORY_PATTERN.fullmatch(repo): raise ReleaseRunConflict("Release repository name is invalid.") if not isinstance(version, str) or not VERSION_PATTERN.fullmatch(version): raise ReleaseRunConflict("Release target version is invalid.") normalized_versions[repo] = version result: dict[str, Any] = { "channel": channel, "repo_versions": normalized_versions, } for field in ("online", "remote_tags", "public_catalog", "include_migrations"): if not isinstance(value.get(field), bool): raise ReleaseRunConflict(f"Release run input {field} must be boolean.") result[field] = value[field] workspace_fingerprint = value.get("workspace_fingerprint") if not isinstance( workspace_fingerprint, str ) or not WORKSPACE_FINGERPRINT_PATTERN.fullmatch(workspace_fingerprint): raise ReleaseRunConflict("Release workspace fingerprint is invalid.") result["workspace_fingerprint"] = workspace_fingerprint return result def _validated_plan(value: dict[str, Any]) -> dict[str, Any]: if not isinstance(value, dict): raise ReleaseRunConflict("Release plan snapshot must be an object.") copied = deepcopy(value) if set(copied) != { "generated_at", "target_channel", "status", "units", "compatibility", "dry_run_steps", "notes", "gate_findings", "recommended_action", "source_preflight_ready", }: raise ReleaseRunConflict("Release plan snapshot fields are invalid.") if not _is_timestamp(copied.get("generated_at")): raise ReleaseRunConflict("Release plan timestamp is invalid.") if copied.get("status") not in PLAN_STATUSES: raise ReleaseRunConflict("Release plan status is invalid.") if not isinstance( copied.get("target_channel"), str ) or not CHANNEL_PATTERN.fullmatch(copied["target_channel"]): raise ReleaseRunConflict("Release plan channel is invalid.") if not isinstance(copied.get("units"), list): raise ReleaseRunConflict("Release plan units are invalid.") if not isinstance(copied.get("compatibility"), list) or not isinstance( copied.get("gate_findings"), list ): raise ReleaseRunConflict("Release plan findings are invalid.") if not isinstance(copied.get("notes"), list) or not all( isinstance(note, str) for note in copied["notes"] ): raise ReleaseRunConflict("Release plan notes are invalid.") if not isinstance(copied.get("source_preflight_ready"), bool): raise ReleaseRunConflict("Release plan preflight state is invalid.") recommendation = copied.get("recommended_action") if recommendation is not None and not isinstance(recommendation, dict): raise ReleaseRunConflict("Release plan recommendation is invalid.") steps = copied.get("dry_run_steps") if ( not isinstance(steps, list) or not steps or len(steps) > MAX_PLAN_STEPS ): raise ReleaseRunConflict("Release plan step collection is invalid.") seen: set[str] = set() for step in steps: if not isinstance(step, dict) or set(step) != { "id", "title", "detail", "command", "cwd", "mutating", "repo", "status", }: raise ReleaseRunConflict("Release plan step is invalid.") step_id = step.get("id") _validate_step_id(step_id) if step_id in seen: raise ReleaseRunConflict("Release plan step identifiers must be unique.") seen.add(step_id) if not isinstance(step.get("mutating", False), bool): raise ReleaseRunConflict("Release plan step mutation flag is invalid.") if step.get("status") not in PLANNED_STEP_STATUSES: raise ReleaseRunConflict("Release plan step status is invalid.") for field, maximum in (("title", 500), ("detail", 2_000)): if not isinstance(step.get(field), str) or len(step[field]) > maximum: raise ReleaseRunConflict("Release plan step text is invalid.") for field in ("command", "cwd", "repo"): if step.get(field) is not None and ( not isinstance(step[field], str) or len(step[field]) > 8_192 ): raise ReleaseRunConflict("Release plan step metadata is invalid.") encoded = _canonical_json(copied) if len(encoded) > MAX_RECORD_BYTES: raise ReleaseRunConflict("Release plan snapshot exceeds its size limit.") return copied def _validate_record(value: Any, *, expected_run_id: str | None = None) -> None: try: if not isinstance(value, dict): raise ValueError if set(value) != { "schema", "schema_version", "run_id", "created_at", "updated_at", "immutable", "state", "record_digest", }: raise ValueError if ( value.get("schema") != SCHEMA or type(value.get("schema_version")) is not int or value["schema_version"] != SCHEMA_VERSION ): raise ValueError run_id = value.get("run_id") if not isinstance(run_id, str) or not RUN_ID_PATTERN.fullmatch(run_id): raise ValueError if expected_run_id is not None and run_id != expected_run_id: raise ValueError for key in ("created_at", "updated_at"): if not _is_timestamp(value.get(key)): raise ValueError created_at = value["created_at"] updated_at = value["updated_at"] if updated_at < created_at: raise ValueError immutable = value.get("immutable") if not isinstance(immutable, dict) or set(immutable) != { "input", "plan", "digest", }: raise ValueError validated_input = _validated_input(immutable["input"]) validated_plan = _validated_plan(immutable["plan"]) if immutable["input"] != validated_input or immutable["plan"] != validated_plan: raise ValueError if immutable.get("digest") != _immutable_digest(immutable): raise ValueError if validated_plan["target_channel"] != validated_input["channel"]: raise ValueError planned_units = validated_plan["units"] planned_versions = { unit.get("repo"): unit.get("target_version") for unit in planned_units if isinstance(unit, dict) and isinstance(unit.get("repo"), str) and isinstance(unit.get("target_version"), str) } if planned_versions != validated_input["repo_versions"] or len( planned_versions ) != len(planned_units): raise ValueError state = value.get("state") if not isinstance(state, dict) or set(state) != { "status", "steps", "events", "next_event_sequence", "processed_requests", "processed_attempts", }: raise ValueError steps = state.get("steps") plan_steps = validated_plan["dry_run_steps"] plan_step_ids = {plan_step["id"] for plan_step in plan_steps} if not isinstance(steps, list) or len(steps) != len(plan_steps): raise ValueError for step, plan_step in zip(steps, plan_steps, strict=True): if ( not isinstance(step, dict) or set(step) != { "id", "state", "attempt_count", "attempt_fingerprint", "started_at", "finished_at", "result_code", } or step.get("id") != plan_step["id"] ): raise ValueError if step.get("state") not in STEP_STATES: raise ValueError if type(step.get("attempt_count")) is not int or step["attempt_count"] < 0: raise ValueError fingerprint = step.get("attempt_fingerprint") if fingerprint is not None and not _is_sha256(fingerprint): raise ValueError for timestamp in (step.get("started_at"), step.get("finished_at")): if timestamp is not None and ( not isinstance(timestamp, str) or len(timestamp) > 40 ): raise ValueError result_code = step.get("result_code") if result_code is not None: _safe_code(result_code) _validate_step_state_fields(step) for timestamp in (step.get("started_at"), step.get("finished_at")): if timestamp is not None and not ( created_at <= timestamp <= updated_at ): raise ValueError first_incomplete_seen = False for step in steps: if first_incomplete_seen and step["state"] != "pending": raise ValueError if step["state"] != "succeeded": first_incomplete_seen = True events = state.get("events") if not isinstance(events, list) or len(events) > MAX_EVENTS: raise ValueError previous_sequence = 0 previous_event_at: str | None = None for event in events: if not isinstance(event, dict) or event.get("type") not in EVENT_TYPES: raise ValueError expected_event_keys = {"sequence", "at", "type"} if event.get("step_id") is not None: expected_event_keys.add("step_id") if event.get("result_code") is not None: expected_event_keys.add("result_code") if set(event) != expected_event_keys: raise ValueError sequence = event.get("sequence") if type(sequence) is not int or sequence <= previous_sequence: raise ValueError previous_sequence = sequence if not _is_timestamp(event.get("at")): raise ValueError if not (created_at <= event["at"] <= updated_at): raise ValueError if previous_event_at is not None and event["at"] < previous_event_at: raise ValueError previous_event_at = event["at"] if event.get("step_id") is not None: _validate_step_id(event["step_id"]) if event["step_id"] not in plan_step_ids: raise ValueError if event.get("result_code") is not None: _safe_code(event["result_code"]) _validate_event_fields(event) next_sequence = state.get("next_event_sequence") if type(next_sequence) is not int or next_sequence <= previous_sequence: raise ValueError requests = state.get("processed_requests") if not isinstance(requests, list) or len(requests) > MAX_COMMANDS: raise ValueError for request in requests: if ( not isinstance(request, dict) or set(request) != {"fingerprint", "operation", "step_id", "parameter"} or not _is_sha256(request.get("fingerprint")) or request.get("operation") not in {"resume", "retry", "reconcile"} ): raise ValueError if request.get("step_id") is not None: _validate_step_id(request["step_id"]) if request["step_id"] not in plan_step_ids: raise ValueError if (request["operation"] == "resume") != (request["step_id"] is None): raise ValueError parameter = request["parameter"] if request["operation"] == "reconcile": if parameter not in RECONCILIATION_OUTCOMES: raise ValueError elif parameter is not None: raise ValueError if len({request["fingerprint"] for request in requests}) != len(requests): raise ValueError attempts = state.get("processed_attempts") if not isinstance(attempts, list) or len(attempts) > MAX_COMMANDS: raise ValueError for attempt in attempts: if ( not isinstance(attempt, dict) or set(attempt) != {"fingerprint", "step_id", "outcome", "result_code"} or not _is_sha256(attempt.get("fingerprint")) or attempt.get("outcome") not in ATTEMPT_OUTCOMES ): raise ValueError _validate_step_id(attempt.get("step_id")) if attempt["step_id"] not in plan_step_ids: raise ValueError attempt_result = attempt.get("result_code") if attempt["outcome"] == "running": if attempt_result is not None: raise ValueError else: if attempt_result is None: raise ValueError _safe_code(attempt_result) if ( attempt["outcome"] == "interrupted" and attempt_result != "process_interrupted" ): raise ValueError if len({attempt["fingerprint"] for attempt in attempts}) != len(attempts): raise ValueError attempts_by_fingerprint = { attempt["fingerprint"]: attempt for attempt in attempts } attempts_per_step = { step_id: sum( attempt["step_id"] == step_id for attempt in attempts ) for step_id in plan_step_ids } for step in steps: if step["attempt_count"] != attempts_per_step[step["id"]]: raise ValueError fingerprint = step["attempt_fingerprint"] if fingerprint is None: continue attempt = attempts_by_fingerprint.get(fingerprint) if attempt is None or attempt["step_id"] != step["id"]: raise ValueError expected_outcome = step["state"] if step["result_code"] == "reconciled_effect_succeeded": if step["state"] != "succeeded" or attempt["outcome"] != "interrupted": raise ValueError elif attempt["outcome"] != expected_outcome: raise ValueError if state.get("status") != _status(value): raise ValueError if value.get("record_digest") != _record_digest(value): raise ValueError except ( KeyError, StopIteration, TypeError, ValueError, ReleaseRunConflict, ) as exc: raise ReleaseRunCorrupt( "Release run record failed schema or integrity validation." ) from exc def _view(record: dict[str, Any]) -> dict[str, Any]: view = deepcopy(record) plan_steps = view["immutable"]["plan"]["dry_run_steps"] mutable_steps = view["state"]["steps"] enriched_steps: list[dict[str, Any]] = [] for index, (step, plan_step) in enumerate( zip(mutable_steps, plan_steps, strict=True) ): available, disabled_reason = _step_available(record, step["id"]) retry_available = step["state"] == "failed" or ( step["state"] == "interrupted" and not bool(plan_step.get("mutating")) ) enriched_steps.append( { **plan_step, **{ key: value for key, value in step.items() if key != "attempt_fingerprint" }, "order": index + 1, "available": available, "retry_available": retry_available, "reconciliation_required": step["state"] == "interrupted" and bool(plan_step.get("mutating")), "disabled_reason": disabled_reason, } ) view["state"]["steps"] = enriched_steps view["state"].pop("processed_requests", None) view["state"].pop("processed_attempts", None) view["recommended_next"] = _recommended_next(record) return view def _step_available(record: dict[str, Any], step_id: str) -> tuple[bool, str]: step, plan_step = _step_pair(record, step_id) if record["immutable"]["plan"].get("status") == "blocked": return False, "The immutable plan contains release-gate blockers." if step["state"] != "pending": descriptions = { "running": "This step is already running.", "succeeded": "This step is complete.", "failed": "Request an explicit retry before running this step again.", "interrupted": "Reconcile the interrupted attempt before retry.", } return False, descriptions.get(step["state"], "This step is unavailable.") planned_status = plan_step.get("status", "planned") if planned_status != "planned": return ( False, f"The plan marks this step {planned_status}; no executor is available.", ) steps = record["state"]["steps"] index = next(index for index, item in enumerate(steps) if item["id"] == step_id) incomplete = next( (item for item in steps[:index] if item["state"] != "succeeded"), None ) if incomplete is not None: return False, f"Complete {incomplete['id']} first." return True, "" def _recommended_next(record: dict[str, Any]) -> dict[str, Any]: plan = record["immutable"]["plan"] if plan.get("status") == "blocked": recommendation = plan.get("recommended_action") if isinstance(recommendation, dict): return { "id": recommendation.get("id") or "resolve_plan_blocker", "step_id": None, "title": recommendation.get("title") or "Resolve release gate", "detail": recommendation.get("detail") or "The frozen plan is blocked.", "remediation": recommendation.get("remediation") or "Resolve the blocker and create a new run from a fresh plan.", "available": False, } return { "id": "resolve_plan_blocker", "step_id": None, "title": "Resolve release gate", "detail": "The frozen plan is blocked.", "remediation": "Resolve the blocker and create a new run from a fresh plan.", "available": False, } for step, plan_step in zip( record["state"]["steps"], plan["dry_run_steps"], strict=True ): if step["state"] == "running": return _step_recommendation( "wait_for_step", step, plan_step, "Wait for the active executor or resume the run after an interruption.", False, ) if step["state"] == "failed": return _step_recommendation( "retry_step", step, plan_step, "Review the bounded result code, then explicitly prepare this known failed step for retry.", True, ) if step["state"] == "interrupted": if bool(plan_step.get("mutating")): return _step_recommendation( "reconcile_step", step, plan_step, "Reconcile the local and remote effect before explicitly preparing a retry.", True, ) return _step_recommendation( "retry_step", step, plan_step, "The read-only attempt was interrupted and can be explicitly prepared for retry.", True, ) if step["state"] == "pending": available, reason = _step_available(record, step["id"]) return _step_recommendation( "execute_step" if available else "unavailable_step", step, plan_step, "Use the existing preview and narrowly confirmed executor controls." if available else reason, available, ) return { "id": "run_complete", "step_id": None, "title": "Release run complete", "detail": "Every recorded step succeeded.", "remediation": "Review publication and installation evidence before closing the release.", "available": False, } def _step_recommendation( recommendation_id: str, step: dict[str, Any], plan_step: dict[str, Any], remediation: str, available: bool, ) -> dict[str, Any]: return { "id": recommendation_id, "step_id": step["id"], "title": plan_step.get("title") or step["id"], "detail": plan_step.get("detail") or "Release plan step.", "remediation": remediation, "available": available, } def _step_pair( record: dict[str, Any], step_id: str ) -> tuple[dict[str, Any], dict[str, Any]]: for step, plan_step in zip( record["state"]["steps"], record["immutable"]["plan"]["dry_run_steps"], strict=True, ): if step["id"] == step_id: return step, plan_step raise ReleaseRunNotFound("Release run step was not found.") def _append_event( record: dict[str, Any], *, event_type: str, step_id: str | None = None, result_code: str | None = None, ) -> None: if event_type not in EVENT_TYPES: raise ReleaseRunConflict("Release event type is invalid.") event: dict[str, Any] = { "sequence": record["state"]["next_event_sequence"], "at": _timestamp(), "type": event_type, } if step_id is not None: _validate_step_id(step_id) event["step_id"] = step_id if result_code is not None: event["result_code"] = _safe_code(result_code) record["state"]["next_event_sequence"] += 1 record["state"]["events"].append(event) record["state"]["events"] = record["state"]["events"][-MAX_EVENTS:] def _ensure_command_capacity(record: dict[str, Any]) -> None: if len(record["state"]["processed_requests"]) >= MAX_COMMANDS: raise ReleaseRunConflict( "Release run command history reached its safety limit; create a new run." ) def _ensure_attempt_capacity(record: dict[str, Any]) -> None: if len(record["state"]["processed_attempts"]) >= MAX_COMMANDS: raise ReleaseRunConflict( "Release run attempt history reached its safety limit; create a new run." ) def _remember_request( record: dict[str, Any], fingerprint: str, operation: str, step_id: str | None, parameter: str | None, ) -> None: _ensure_command_capacity(record) record["state"]["processed_requests"].append( { "fingerprint": fingerprint, "operation": operation, "step_id": step_id, "parameter": parameter, } ) def _request_seen( record: dict[str, Any], fingerprint: str, operation: str, step_id: str | None, parameter: str | None, ) -> bool: matches = [ request for request in record["state"]["processed_requests"] if request["fingerprint"] == fingerprint ] if not matches: return False if any( request["operation"] == operation and request["step_id"] == step_id and request["parameter"] == parameter for request in matches ): return True raise ReleaseRunConflict( "The request identifier was already used for a different run-state action." ) def _find_attempt( record: dict[str, Any], fingerprint: str ) -> dict[str, Any] | None: return next( ( attempt for attempt in record["state"]["processed_attempts"] if attempt["fingerprint"] == fingerprint ), None, ) def _attempt_for_fingerprint( record: dict[str, Any], fingerprint: str | None, *, step_id: str ) -> dict[str, Any]: if fingerprint is None: raise ReleaseRunCorrupt( "Running release step has no persisted attempt identifier." ) attempt = _find_attempt(record, fingerprint) if attempt is None or attempt["step_id"] != step_id: raise ReleaseRunCorrupt( "Running release step has no matching persisted attempt." ) return attempt def _refresh_status(record: dict[str, Any]) -> None: record["state"]["status"] = _status(record) def _status(record: dict[str, Any]) -> str: if record["immutable"]["plan"].get("status") == "blocked": return "blocked" states = [step["state"] for step in record["state"]["steps"]] if states and all(state == "succeeded" for state in states): return "completed" if "running" in states: return "running" if any(state in {"failed", "interrupted"} for state in states): return "attention" if states: first_pending = next( step for step, mutable in zip( record["immutable"]["plan"]["dry_run_steps"], record["state"]["steps"], strict=True, ) if mutable["state"] == "pending" ) if first_pending.get("status", "planned") != "planned": return "blocked" return "planned" def _immutable_digest(immutable: dict[str, Any]) -> str: return hashlib.sha256( _canonical_json( {"input": immutable.get("input"), "plan": immutable.get("plan")} ) ).hexdigest() def _record_digest(record: dict[str, Any]) -> str: return hashlib.sha256( _canonical_json( {key: value for key, value in record.items() if key != "record_digest"} ) ).hexdigest() def _request_fingerprint(request_id: str) -> str: if not isinstance(request_id, str) or not REQUEST_ID_PATTERN.fullmatch(request_id): raise ReleaseRunConflict("Run-state request identifier is invalid.") return hashlib.sha256(request_id.encode("utf-8")).hexdigest() def _validate_step_id(step_id: Any) -> None: if not isinstance(step_id, str) or not STEP_ID_PATTERN.fullmatch(step_id): raise ReleaseRunConflict("Release plan step identifier is invalid.") def _safe_code(value: str) -> str: if ( not isinstance(value, str) or len(value) > 64 or not re.fullmatch(r"[a-z][a-z0-9_]*", value) ): raise ReleaseRunConflict("Release result code is invalid.") return value def _is_sha256(value: Any) -> bool: return isinstance(value, str) and bool(re.fullmatch(r"[0-9a-f]{64}", value)) def _is_timestamp(value: Any) -> bool: if not isinstance(value, str) or not re.fullmatch( r"[0-9]{4}-[0-9]{2}-[0-9]{2}T[0-9]{2}:[0-9]{2}:[0-9]{2}Z", value ): return False try: datetime.strptime(value, "%Y-%m-%dT%H:%M:%SZ") except ValueError: return False return True def _validate_step_state_fields(step: dict[str, Any]) -> None: state = step["state"] attempt_count = step["attempt_count"] fingerprint = step["attempt_fingerprint"] started_at = step["started_at"] finished_at = step["finished_at"] result_code = step["result_code"] if started_at is not None and not _is_timestamp(started_at): raise ValueError if finished_at is not None and not _is_timestamp(finished_at): raise ValueError if started_at is not None and finished_at is not None and finished_at < started_at: raise ValueError if state == "pending": if any( value is not None for value in (fingerprint, started_at, finished_at, result_code) ): raise ValueError return if attempt_count < 1 or fingerprint is None or started_at is None: raise ValueError if state == "running": if finished_at is not None or result_code is not None: raise ValueError return if finished_at is None or result_code is None: raise ValueError if state == "interrupted" and result_code != "process_interrupted": raise ValueError def _validate_event_fields(event: dict[str, Any]) -> None: event_type = event["type"] has_step = "step_id" in event has_result = "result_code" in event if event_type in {"run_created", "run_resumed"}: if has_step or has_result: raise ValueError return if not has_step: raise ValueError if event_type == "step_started": if has_result: raise ValueError return if not has_result: raise ValueError if event_type == "step_reconciled" and event["result_code"] not in ( RECONCILIATION_OUTCOMES ): raise ValueError if ( event_type == "step_interrupted" and event["result_code"] != "process_interrupted" ): raise ValueError def _create_private_directory_tree(path: Path) -> None: missing: list[Path] = [] current = path while not current.exists(): missing.append(current) if current == current.parent: break current = current.parent _validate_authority_components(current) for directory in reversed(missing): try: os.mkdir(directory, mode=0o700) except FileExistsError: pass except OSError as exc: raise ReleaseRunCorrupt( "Release run storage directory could not be created safely." ) from exc _reject_symlink_components(directory) try: metadata = directory.lstat() if not stat.S_ISDIR(metadata.st_mode): raise ReleaseRunCorrupt( "Release run storage path is not a directory." ) if hasattr(os, "geteuid") and metadata.st_uid != os.geteuid(): raise ReleaseRunCorrupt( "Release run storage is not owned by the console operator." ) directory.chmod(0o700) except OSError as exc: raise ReleaseRunCorrupt( "Release run storage permissions could not be secured." ) from exc if directory.stat().st_mode & 0o077: raise ReleaseRunCorrupt( "Release run storage permissions are broader than 0700." ) _fsync_directory(directory) _fsync_directory(directory.parent) _validate_authority_components(path.parent) def _validate_authority_components(path: Path) -> None: home = Path(os.path.abspath(os.fspath(Path.home()))) if path == home or path.is_relative_to(home): current = home parts = path.relative_to(home).parts candidates = [current] else: current = Path(path.anchor) parts = path.parts[1:] candidates = [] for part in parts: current /= part candidates.append(current) for current in candidates: try: metadata = current.lstat() except OSError as exc: raise ReleaseRunCorrupt( "Release run storage ancestry could not be validated." ) from exc if stat.S_ISLNK(metadata.st_mode) or not stat.S_ISDIR(metadata.st_mode): raise ReleaseRunCorrupt( "Release run storage ancestry is not a safe directory path." ) if hasattr(os, "geteuid") and metadata.st_uid not in {0, os.geteuid()}: raise ReleaseRunCorrupt( "Release run storage ancestry has an untrusted owner." ) writable_by_others = bool(metadata.st_mode & 0o022) sticky = bool(metadata.st_mode & stat.S_ISVTX) if writable_by_others and not sticky: raise ReleaseRunCorrupt( "Release run storage ancestry is writable by another principal." ) def _fsync_directory(path: Path) -> None: flags = os.O_RDONLY if hasattr(os, "O_DIRECTORY"): flags |= os.O_DIRECTORY try: descriptor = os.open(path, flags) try: os.fsync(descriptor) finally: os.close(descriptor) except OSError as exc: raise ReleaseRunCorrupt( "Release run storage metadata could not be persisted." ) from exc def _reject_symlink_components(path: Path) -> None: current = Path(path.anchor) for part in path.parts[1:]: current /= part if current.is_symlink(): raise ReleaseRunCorrupt( "Release run storage must not contain symbolic-link components." ) def _canonical_json(value: Any) -> bytes: try: return json.dumps( value, sort_keys=True, separators=(",", ":"), ensure_ascii=False, allow_nan=False, ).encode("utf-8") except (TypeError, ValueError) as exc: raise ReleaseRunConflict("Release run contains non-JSON data.") from exc def _timestamp() -> str: return ( datetime.now(tz=UTC).replace(microsecond=0).isoformat().replace("+00:00", "Z") )