"""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 SCHEMA = "govoplan.release-run" SCHEMA_VERSION = 1 MAX_RECORD_BYTES = 2 * 1024 * 1024 MAX_PLAN_STEPS = 512 MAX_REPOSITORY_VERSIONS = 128 MAX_EVENTS = 256 MAX_REQUESTS = 128 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_retry_requested", } class ReleaseRunError(RuntimeError): """Base error for durable release run state.""" class ReleaseRunNotFound(ReleaseRunError): """Raised when a run does not exist.""" 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) -> 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()))) 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) 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": [], }, } _append_event(record, event_type="run_created") _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, )[:limit] summaries: list[dict[str, Any]] = [] for path in paths: try: record = self._read(path.stem) except (ReleaseRunCorrupt, ReleaseRunNotFound): summaries.append( { "run_id": path.stem, "status": "unavailable", "error": "Run record is unavailable or failed integrity validation.", } ) 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"], } ) 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): return _view(record) for step in record["state"]["steps"]: if step["state"] != "running": continue 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) _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, reconciled: bool = False, ) -> 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): 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")) and not reconciled ): raise ReleaseRunConflict( "An interrupted mutating step must be reconciled before retry." ) step.update( { "state": "pending", "attempt_fingerprint": None, "started_at": None, "finished_at": None, "result_code": None, } ) _remember_request(record, request_fingerprint, "retry", step_id) _append_event( record, event_type="step_retry_requested", step_id=step_id, result_code="reconciled" if reconciled else "known_failure", ) 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) step, _plan_step = _step_pair(record, step_id) if step["attempt_fingerprint"] == attempt_fingerprint and step["state"] in { "running", "succeeded", }: return _view(record) available, reason = _step_available(record, step_id) if not available: raise ReleaseRunConflict(reason) step.update( { "state": "running", "attempt_count": step["attempt_count"] + 1, "attempt_fingerprint": attempt_fingerprint, "started_at": _timestamp(), "finished_at": None, "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) if ( step["attempt_fingerprint"] == attempt_fingerprint and step["state"] == target_state and step["result_code"] == result_code ): return _view(record) 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, } ) _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: 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 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) return payload 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) self.root.mkdir(mode=0o700, parents=True, exist_ok=True) _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.") 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." ) 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 try: path.chmod(0o600) except OSError as exc: raise ReleaseRunCorrupt( "Release run record permissions could not be secured." ) from exc directory_flags = os.O_RDONLY if hasattr(os, "O_DIRECTORY"): directory_flags |= os.O_DIRECTORY directory_descriptor = os.open(self.root, directory_flags) try: os.fsync(directory_descriptor) finally: os.close(directory_descriptor) 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 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 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 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", }: 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) events = state.get("events") if not isinstance(events, list) or len(events) > MAX_EVENTS: raise ValueError previous_sequence = 0 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 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_REQUESTS: raise ValueError for request in requests: if ( not isinstance(request, dict) or set(request) != {"fingerprint", "operation", "step_id"} or not _is_sha256(request.get("fingerprint")) or request.get("operation") not in {"resume", "retry"} ): 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 if len({request["fingerprint"] for request in requests}) != len(requests): 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, "disabled_reason": disabled_reason, } ) view["state"]["steps"] = enriched_steps view["state"].pop("processed_requests", 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.", False, ) 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 _remember_request( record: dict[str, Any], fingerprint: str, operation: str, step_id: str | None ) -> None: record["state"]["processed_requests"].append( { "fingerprint": fingerprint, "operation": operation, "step_id": step_id, } ) record["state"]["processed_requests"] = record["state"]["processed_requests"][ -MAX_REQUESTS: ] def _request_seen( record: dict[str, Any], fingerprint: str, operation: str, step_id: 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 for request in matches ): return True raise ReleaseRunConflict( "The request identifier was already used for a different run-state action." ) 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 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 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") )