diff --git a/README.md b/README.md index d908f43..7ed884b 100644 --- a/README.md +++ b/README.md @@ -106,6 +106,16 @@ The exact instance-ID comparison is consistency evidence, not a newly invented cryptographic peer identity; preflight reports the latter as `NOT_ASSERTED`. `CONFIGURED` is therefore never presented as `LIVE_READY`. +## Bounded installed maintenance + +The product-owned external controller for the selected Forge 2.7.21 to 2.7.22 +installation transition is documented in the +[installed update runbook](docs/operations/FORGE_INSTALLED_UPDATE_RUNBOOK.md). +It stages an exact qualified wheel, backs up and migrates the selected runtime, +and atomically activates a versioned slot. It is not packaged as a Forge +Runtime command and does not replace Forge Platform's normal installer/update +composition boundary. + ## Operational-history reset Forge schema 38 provides a bounded, product-owned maintenance service under diff --git a/docs/operations/FORGE_INSTALLED_UPDATE_RUNBOOK.md b/docs/operations/FORGE_INSTALLED_UPDATE_RUNBOOK.md new file mode 100644 index 0000000..44878e0 --- /dev/null +++ b/docs/operations/FORGE_INSTALLED_UPDATE_RUNBOOK.md @@ -0,0 +1,117 @@ +# Forge installed update controller + +Status: bounded product-owned maintenance provisioner for the selected +Forge 2.7.21 to 2.7.22 transition. + +This controller closes one concrete product provisioning gap. It is not the +universal Forge Platform installer, an installer UI, a new service supervisor, +or a Forge Runtime self-updater. Forge Platform remains the owner of normal +cross-product installation and update composition. Forge retains the concrete +runtime selection, schema migration, recovery, verification, and cleanup +semantics that the platform must eventually invoke through a qualified adapter. + +The external entry point is +[`scripts/update_installed_forge.py`](../../scripts/update_installed_forge.py). +It is deliberately excluded from the `forge-autonomy` wheel and therefore +cannot be called by the installed runtime as a second self-installer. + +## Supported operation + +The controller accepts only one explicitly bound existing installation and one +exact `forge-autonomy` 2.7.22 wheel. Every invocation binds: + +- operation, runtime, installation, and peer-configuration identities; +- data root, runtime root, current command resolver, legacy interpreter, and + candidate Python interpreter; +- product version and original product source; +- wheel and terminal release-qualification receipt, including their SHA-256 + digests; +- protected installation-controller source and file digest. + +It rejects a changed target, artifact, receipt, resolver, peer binding, writer +state, or concurrent maintenance owner. Unknown historical installer +provenance remains unknown; adoption records only the observed entry point, +interpreter, version, and bytes. + +## Safety sequence + +1. Read the wheel once through a no-follow descriptor; validate its digest, + canonical RECORD, purelib tag, package metadata, member allowlist, and the + exact terminal release/reconciliation receipt without importing it. +2. Create an isolated versioned runtime slot outside the source checkout with + the explicit Python interpreter. Extract only the already validated bytes + into a pip-free virtual environment, then verify every installed file and + the generated command wrapper. An interrupted unreceipted slot is + quarantined and rebuilt, never trusted in place. +3. Read the candidate's version, distribution, module, prefix, and interpreter + from an isolated process. +4. Acquire the installation-operation lock, canonical + `forge-runtime-mutation.lock`, and bootstrap `locks/runtime.lock`; reject an + active runtime process, dispatcher, non-terminal Mission or scheduler + submission, active provider-generation permit, planning queue, reset, or + conflicting operation. +5. Adopt the selected legacy command entry point behind a stable product-owned + resolver. Before migration it still resolves to the byte-equal retained + legacy entry point. +6. Take a SQLite backup through the backup API while the writer fence is held. + Verify `integrity_check`, foreign keys, digest, and the complete logical + pre-migration snapshot. +7. Migrate an isolated copy of that backup through the staged candidate's + normal `forge ... server init` path. Require schema 38, the exact new reset + table set, an idle reset control row, and byte-logical preservation of every + pre-existing domain table and protected metadata/binding. +8. Point the stable resolver at a maintenance fence. Take an exclusive SQLite + writer boundary, prove the live database is still byte-logically identical + to the backed-up snapshot, and atomically install the already Forge-migrated + database copy. The replacement remains read-only until activation and final + receipt persistence; full preservation and sidecar checks run again. +9. Atomically select the candidate slot and read back the exact installed CLI, + module, interpreter, version, runtime identity, data root, schema, and peer + binding. No service or historical Mission is started. +10. Persist one protected operation receipt below + `artifacts/installation/` and retain the verified database + backup below `backups/installation/`. + +The controller never calls operational-reset prepare/apply/resume/finish, +Mission intake, Action derivation, a planner/provider, or an EP submission. +Credential material is neither read nor archived. A later installed CLI +preflight resolves the unchanged credential reference only for its separately +authorized read-only Forge-to-EP check. + +## Interruption and resume + +Re-run the same exact operation ID and arguments. A conflicting request is +rejected. + +- Before atomic database replacement, failure restores the retained 2.7.21 + command route only when the complete live snapshot still equals `before`. + A hard interruption may leave the explicit maintenance fence; resuming the + same operation reconciles it from durable evidence. +- From the first schema-38 readback or any ambiguous partial state onward, the + old binary is never selected. Any caught activation/readback failure selects + the maintenance fence; replay reconciles the exact protected candidate. +- A completed receipt is idempotently returned only after revalidating its + bytes, backup, release receipt, controller, wheel manifest, slot, resolver, + installed identity, target binding, schema, and database integrity. Later + legitimate runtime history does not invalidate the historical update + receipt. + +The backup is recovery evidence, not permission for an automatic database +rollback. Restoring it after later security, budget, or external effects needs +separate authority and compatibility proof. + +## Qualification boundary + +`tests/test_installed_forge_update.py` covers exact release binding, target and +writer rejection, the real Forge schema-37 to schema-38 migrator, preservation +of Missions, allocations, reviews, execution receipts, governance grants, +configuration and identity, concurrent-operation exclusion, resolver adoption, +canonical receipt shape, path safety, exact slot contents, exclusive atomic +database replacement, late-writer rejection, and interruption before +migration, after migration, and during activation. An opt-in test runs the +entire route and replay against the exact published wheel and terminal release +receipt. + +Production use additionally requires protected merge/check evidence for the +exact controller source, a terminal release-complete receipt, and live +post-activation installed CLI, reset-preview, and authenticated EP readback. diff --git a/scripts/update_installed_forge.py b/scripts/update_installed_forge.py new file mode 100644 index 0000000..03398c1 --- /dev/null +++ b/scripts/update_installed_forge.py @@ -0,0 +1,1516 @@ +#!/usr/bin/env python3 +"""Bounded product-owned update controller for one installed Forge runtime. + +This is an external maintenance provisioner, not a Forge Runtime command. It +stages one exact qualified wheel in an immutable slot, takes the canonical +runtime mutation lock, creates and verifies a SQLite backup, qualifies Forge's +own migration on an isolated copy, fences the legacy command resolver, invokes +the candidate's owning migrator, and atomically activates the candidate. + +The durable operation is intentionally narrow: callers must supply every +installation, artifact, source, identity, interpreter, and resolver binding. +It never starts a service, Mission, planner, provider, reset, or EP submission. +""" + +from __future__ import annotations + +import argparse +import base64 +from contextlib import contextmanager +import csv +from dataclasses import asdict, dataclass +from datetime import datetime, timezone +from email.parser import BytesParser +from hashlib import sha256 +from io import BytesIO, StringIO +import json +import os +from pathlib import Path +import re +import secrets +import shlex +import shutil +import sqlite3 +import stat +import subprocess +import sys +import tempfile +from typing import Any, Callable, Iterator, Mapping, Sequence +import zipfile + +try: + import fcntl +except ImportError: # pragma: no cover - supported installation target is POSIX. + fcntl = None # type: ignore[assignment] + + +CONTRACT_VERSION = "forge-installed-update/v1" +SCHEMA_BEFORE = 37 +SCHEMA_AFTER = 38 +PHASE_ORDER = { + phase: index for index, phase in enumerate(( + "PREPARED", "STAGED", "ADOPTED", "BACKED_UP", "MIGRATION_QUALIFIED", + "FENCED", "MIGRATED", "ACTIVATING", "ACTIVATED", "COMPLETE", + )) +} +NEW_SCHEMA_38_TABLES = frozenset({ + "operational_reset_state", + "operational_reset_operations", + "operational_reset_audit", + "operational_reset_tombstones", + "operational_reset_artifact_steps", +}) +VOLATILE_METADATA_KEYS = frozenset({ + "schema_version", "migration_version", "last_migration", "forge_version", + "database_version", "last_access_at", "integrity_status", +}) +SAFE_MISSION_STATES = frozenset({ + "BLOCKED", "FAILED", "COMPLETED", "ARCHIVED", "INTEGRATION_BLOCKED", "INTEGRATION_COMPLETE", +}) +SAFE_SUBMISSION_STATES = frozenset({"RECONCILED", "BLOCKED", "FAILED", "SUPERSEDED"}) +TERMINAL_PERMIT_STATES = frozenset({"INVALIDATED", "INVALIDATED_BY_OPERATIONAL_RESET"}) + + +class InstalledForgeUpdateError(RuntimeError): + """The selected installation cannot be updated without weakening a gate.""" + + +def _now() -> str: + return datetime.now(timezone.utc).isoformat().replace("+00:00", "Z") + + +def _json_bytes(value: object) -> bytes: + return (json.dumps(value, sort_keys=True, separators=(",", ":"), ensure_ascii=False) + "\n").encode() + + +def _digest_bytes(value: bytes) -> str: + return "sha256:" + sha256(value).hexdigest() + + +def _assert_no_symlink_components(path: Path) -> None: + if not path.is_absolute(): + raise InstalledForgeUpdateError(f"path must be absolute: {path}") + current = Path(path.anchor) + for part in path.parts[1:]: + current /= part + try: + mode = current.lstat().st_mode + except FileNotFoundError: + break + if stat.S_ISLNK(mode): + raise InstalledForgeUpdateError(f"path contains a symbolic-link component: {current}") + + +def _read_regular_bytes(path: Path) -> bytes: + _assert_no_symlink_components(path) + flags = os.O_RDONLY | getattr(os, "O_NOFOLLOW", 0) + try: + descriptor = os.open(path, flags) + except OSError as error: + raise InstalledForgeUpdateError(f"required regular file is unavailable: {path}") from error + try: + if not stat.S_ISREG(os.fstat(descriptor).st_mode): + raise InstalledForgeUpdateError(f"required regular file is unavailable: {path}") + chunks: list[bytes] = [] + while True: + chunk = os.read(descriptor, 1024 * 1024) + if not chunk: + break + chunks.append(chunk) + return b"".join(chunks) + finally: + os.close(descriptor) + + +def file_digest(path: Path) -> str: + _assert_no_symlink_components(path) + flags = os.O_RDONLY | getattr(os, "O_NOFOLLOW", 0) + try: + descriptor = os.open(path, flags) + except OSError as error: + raise InstalledForgeUpdateError(f"required regular file is unavailable: {path}") from error + digest = sha256() + try: + if not stat.S_ISREG(os.fstat(descriptor).st_mode): + raise InstalledForgeUpdateError(f"required regular file is unavailable: {path}") + while True: + chunk = os.read(descriptor, 1024 * 1024) + if not chunk: + break + digest.update(chunk) + finally: + os.close(descriptor) + return "sha256:" + digest.hexdigest() + + +def _safe_directory(path: Path, *, create: bool = False) -> Path: + if not path.is_absolute(): + raise InstalledForgeUpdateError(f"path must be absolute: {path}") + _assert_no_symlink_components(path) + if create: + path.mkdir(parents=True, exist_ok=True) + _assert_no_symlink_components(path) + if not path.is_dir(): + raise InstalledForgeUpdateError(f"directory is unavailable or unsafe: {path}") + return path + + +def _atomic_json(path: Path, value: object) -> None: + _safe_directory(path.parent, create=True) + descriptor, name = tempfile.mkstemp(prefix=f".{path.name}.tmp-", dir=path.parent) + temporary = Path(name) + os.chmod(temporary, 0o600) + try: + with os.fdopen(descriptor, "wb") as handle: + handle.write(_json_bytes(value)) + handle.flush() + os.fsync(handle.fileno()) + os.replace(temporary, path) + os.chmod(path, 0o600) + finally: + if temporary.exists(): + temporary.unlink() + + +def _read_json(path: Path) -> dict[str, Any]: + try: + value = json.loads(_read_regular_bytes(path)) + except (OSError, UnicodeDecodeError, json.JSONDecodeError) as error: + raise InstalledForgeUpdateError(f"durable JSON evidence is unreadable: {path}") from error + if not isinstance(value, dict): + raise InstalledForgeUpdateError(f"durable JSON evidence is not an object: {path}") + return value + + +def _replace_symlink(path: Path, target: str) -> None: + _safe_directory(path.parent, create=True) + temporary = path.with_name(f".{path.name}.link-{secrets.token_hex(16)}") + temporary.symlink_to(target) + try: + os.replace(temporary, path) + finally: + if temporary.is_symlink(): + temporary.unlink() + + +def _resolved_link(path: Path) -> Path: + if not path.is_symlink(): + raise InstalledForgeUpdateError(f"managed resolver is not a symlink: {path}") + return path.resolve(strict=True) + + +def _environment() -> dict[str, str]: + environment = { + "PATH": "/usr/bin:/bin:/usr/sbin:/sbin", + "LANG": "C.UTF-8", + "LC_ALL": "C.UTF-8", + "PYTHONNOUSERSITE": "1", + "PYTHONDONTWRITEBYTECODE": "1", + } + if os.environ.get("HOME"): + environment["HOME"] = os.environ["HOME"] + return environment + + +def _run(arguments: Sequence[str], *, cwd: Path, timeout: int = 300) -> subprocess.CompletedProcess[str]: + try: + return subprocess.run( + tuple(arguments), cwd=cwd, env=_environment(), text=True, + stdout=subprocess.PIPE, stderr=subprocess.PIPE, timeout=timeout, check=True, + ) + except (OSError, subprocess.CalledProcessError, subprocess.TimeoutExpired) as error: + stdout = getattr(error, "stdout", "") or "" + stderr = getattr(error, "stderr", "") or "" + diagnostic = (stdout + "\n" + stderr).strip()[-2000:] + raise InstalledForgeUpdateError( + f"bounded subprocess failed: {arguments[0]}: {diagnostic or type(error).__name__}" + ) from error + + +def installed_identity(interpreter: Path, *, cwd: Path) -> dict[str, Any]: + program = """ +import importlib.metadata, json, pathlib, sys +import forge +from forge._version import canonical_version +print(json.dumps({ + "version": canonical_version(), + "distribution_version": importlib.metadata.version("forge-autonomy"), + "module": str(pathlib.Path(forge.__file__).resolve()), + "sys_executable": sys.executable, + "resolved_sys_executable": str(pathlib.Path(sys.executable).resolve()), + "prefix": str(pathlib.Path(sys.prefix).resolve()), +}, sort_keys=True)) +""" + result = _run((str(interpreter), "-B", "-I", "-c", program), cwd=cwd) + try: + identity = json.loads(result.stdout) + except json.JSONDecodeError as error: + raise InstalledForgeUpdateError("installed identity readback is malformed") from error + if not isinstance(identity, dict): + raise InstalledForgeUpdateError("installed identity readback is malformed") + return identity + + +@dataclass(frozen=True) +class UpdateRequest: + operation_id: str + version: str + product_source: str + wheel: str + wheel_sha256: str + qualification_receipt: str + qualification_receipt_sha256: str + controller_source: str + controller_sha256: str + data_root: str + runtime_root: str + runtime_id: str + installation_id: str + peer_configuration_digest: str + resolver: str + resolver_sha256: str + existing_interpreter: str + existing_version: str + base_python: str + + def validate(self) -> None: + identifiers = (self.operation_id, self.runtime_id, self.installation_id) + if any( + re.fullmatch(r"[A-Za-z0-9][A-Za-z0-9._-]{0,127}", value) is None + or value in {".", ".."} + for value in identifiers + ): + raise InstalledForgeUpdateError("operation and installation identities must be filesystem-safe") + if self.version != "2.7.22" or self.existing_version != "2.7.21": + raise InstalledForgeUpdateError("this bounded controller supports only the selected 2.7.21 to 2.7.22 update") + for label, digest in ( + ("wheel", self.wheel_sha256), ("qualification receipt", self.qualification_receipt_sha256), + ("controller", self.controller_sha256), ("resolver", self.resolver_sha256), + ("peer configuration", self.peer_configuration_digest), + ): + if not digest.startswith("sha256:") or len(digest) != 71: + raise InstalledForgeUpdateError(f"{label} digest is invalid") + if len(self.product_source) != 40 or len(self.controller_source) != 40: + raise InstalledForgeUpdateError("source revisions must be exact 40-character revisions") + paths = ( + self.wheel, self.qualification_receipt, self.data_root, self.runtime_root, + self.resolver, self.existing_interpreter, self.base_python, + ) + if any(not Path(value).is_absolute() for value in paths): + raise InstalledForgeUpdateError("all installation paths must be absolute") + + @property + def digest(self) -> str: + return _digest_bytes(_json_bytes(asdict(self))) + + +def _validated_wheel(request: UpdateRequest) -> tuple[bytes, dict[str, str]]: + wheel = Path(request.wheel) + expected_name = f"forge_autonomy-{request.version}-py3-none-any.whl" + if wheel.name != expected_name: + raise InstalledForgeUpdateError("wheel filename does not match the selected product and version") + wheel_bytes = _read_regular_bytes(wheel) + if _digest_bytes(wheel_bytes) != request.wheel_sha256: + raise InstalledForgeUpdateError("wheel digest does not match the selected qualified artifact") + dist_info = f"forge_autonomy-{request.version}.dist-info" + try: + with zipfile.ZipFile(BytesIO(wheel_bytes)) as archive: + infos = archive.infolist() + names = [info.filename for info in infos] + if len(names) != len(set(names)): + raise InstalledForgeUpdateError("wheel contains duplicate member paths") + for info in infos: + parts = Path(info.filename).parts + mode = info.external_attr >> 16 + if ( + info.filename.startswith("/") or "\\" in info.filename or ".." in parts + or not parts or parts[0] not in {"forge", dist_info} + or stat.S_ISLNK(mode) + ): + raise InstalledForgeUpdateError("wheel contains an unsafe or unexpected member path") + metadata_name = f"{dist_info}/METADATA" + wheel_metadata_name = f"{dist_info}/WHEEL" + record_name = f"{dist_info}/RECORD" + if any(name not in names for name in (metadata_name, wheel_metadata_name, record_name)): + raise InstalledForgeUpdateError("wheel metadata is missing or ambiguous") + metadata = BytesParser().parsebytes(archive.read(metadata_name)) + wheel_metadata = BytesParser().parsebytes(archive.read(wheel_metadata_name)) + if ( + metadata.get("Name") != "forge-autonomy" + or metadata.get("Version") != request.version + or wheel_metadata.get("Root-Is-Purelib") != "true" + or "py3-none-any" not in wheel_metadata.get_all("Tag", []) + ): + raise InstalledForgeUpdateError("wheel package metadata does not match the supported Forge artifact") + rows = list(csv.reader(StringIO(archive.read(record_name).decode("utf-8")))) + if any(len(row) != 3 for row in rows): + raise InstalledForgeUpdateError("wheel RECORD is malformed") + records = {row[0]: (row[1], row[2]) for row in rows} + files = [info for info in infos if not info.is_dir()] + if len(records) != len(rows) or set(records) != {info.filename for info in files}: + raise InstalledForgeUpdateError("wheel RECORD does not exactly enumerate the artifact") + manifest: dict[str, str] = {} + for info in files: + payload = archive.read(info.filename) + manifest[info.filename] = _digest_bytes(payload) + recorded_hash, recorded_size = records[info.filename] + if info.filename == record_name: + if recorded_hash or recorded_size: + raise InstalledForgeUpdateError("wheel RECORD self-entry is not canonical") + continue + encoded = base64.urlsafe_b64encode(sha256(payload).digest()).rstrip(b"=").decode("ascii") + if recorded_hash != f"sha256={encoded}" or recorded_size != str(len(payload)): + raise InstalledForgeUpdateError("wheel RECORD does not bind an artifact member") + except (UnicodeDecodeError, zipfile.BadZipFile) as error: + raise InstalledForgeUpdateError("wheel is not a valid canonical ZIP artifact") from error + return wheel_bytes, manifest + + +def _qualified_artifact(request: UpdateRequest) -> tuple[dict[str, Any], bytes, dict[str, str]]: + wheel_bytes, manifest = _validated_wheel(request) + expected_name = f"forge_autonomy-{request.version}-py3-none-any.whl" + sdist_name = f"forge_autonomy-{request.version}.tar.gz" + receipt_path = Path(request.qualification_receipt) + receipt_bytes = _read_regular_bytes(receipt_path) + if _digest_bytes(receipt_bytes) != request.qualification_receipt_sha256: + raise InstalledForgeUpdateError("qualification receipt digest changed") + try: + receipt = json.loads(receipt_bytes) + except (UnicodeDecodeError, json.JSONDecodeError) as error: + raise InstalledForgeUpdateError("qualification receipt is malformed") from error + if not isinstance(receipt, dict): + raise InstalledForgeUpdateError("qualification receipt is malformed") + qualification = receipt.get("qualification") + artifacts = receipt.get("artifacts") + publication = receipt.get("publication_receipt") + cleanup = receipt.get("cleanup") + release = publication.get("github_release") if isinstance(publication, dict) else None + cleanup_release = cleanup.get("github_release") if isinstance(cleanup, dict) else None + expected_top = { + "product", "component", "version", "source_revision", "operation_id", "policy_revision", + "artifacts", "state", "qualification", "publication_receipt", "cleanup", + } + expected_publication = { + "github_release", "observed_artifact_digests", "original_release_run_conclusion", + "original_release_run_id", "product_source_revision", "readback", "reconciliation_contract", + "reconciliation_run_id", "registry", "release_controller_source", + } + expected_cleanup = { + "github_release", "operation_local_cleanup", "original_release_run_id", + "reconciliation_contract", "reconciliation_run_id", "release_controller_source", "result", + } + expected_public_release = {"api_url", "database_id", "node_id", "tag", "target_commitish"} + expected_cleanup_release = { + "api_url", "database_id", "draft", "node_id", "tag", "tag_commit", "target_commitish", + } + sdist_digest = artifacts.get("sdist") if isinstance(artifacts, dict) else None + release_controller = publication.get("release_controller_source") if isinstance(publication, dict) else None + original_run = publication.get("original_release_run_id") if isinstance(publication, dict) else None + reconciliation_run = publication.get("reconciliation_run_id") if isinstance(publication, dict) else None + database_id = release.get("database_id") if isinstance(release, dict) else None + expected_api = f"https://api.github.com/repos/pcvantol/forge/releases/{database_id}" + exact_artifacts = {"wheel": request.wheel_sha256, "sdist": sdist_digest} + exact_qualified = {f"dist/{expected_name}": request.wheel_sha256, f"dist/{sdist_name}": sdist_digest} + exact_observed = {expected_name: request.wheel_sha256, sdist_name: sdist_digest} + if ( + set(receipt) != expected_top + or receipt.get("state") != "RELEASE_COMPLETE" + or receipt.get("product") != "forge" or receipt.get("component") != "forge-autonomy" + or receipt.get("version") != request.version or receipt.get("source_revision") != request.product_source + or receipt.get("operation_id") != f"forge-release-{request.version}-{request.product_source}" + or receipt.get("policy_revision") != "forge-bootstrap-release-cadence-v2" + or not isinstance(artifacts, dict) or set(artifacts) != {"wheel", "sdist"} + or artifacts != exact_artifacts or not isinstance(sdist_digest, str) + or re.fullmatch(r"sha256:[0-9a-f]{64}", sdist_digest) is None + or not isinstance(qualification, dict) or set(qualification) != { + "artifact_digests", "exact_main_sha", "qualification" + } + or qualification.get("exact_main_sha") != request.product_source + or qualification.get("qualification") != "forge-production-distribution" + or qualification.get("artifact_digests") != exact_qualified + or not isinstance(publication, dict) or set(publication) != expected_publication + or publication.get("product_source_revision") != request.product_source + or publication.get("readback") != "PASS" or publication.get("registry") != "pypi" + or publication.get("original_release_run_conclusion") != "failure" + or publication.get("reconciliation_contract") != "forge-existing-release-reconciliation/v1" + or publication.get("observed_artifact_digests") != exact_observed + or not isinstance(original_run, str) or not original_run.isdigit() + or not isinstance(reconciliation_run, str) or not reconciliation_run.isdigit() + or not isinstance(release_controller, str) or re.fullmatch(r"[0-9a-f]{40}", release_controller) is None + or not isinstance(release, dict) or set(release) != expected_public_release + or release.get("api_url") != expected_api or not isinstance(database_id, int) or database_id <= 0 + or not isinstance(release.get("node_id"), str) or not release.get("node_id") + or release.get("tag") != f"forge-v{request.version}" + or release.get("target_commitish") != request.product_source + or not isinstance(cleanup, dict) or set(cleanup) != expected_cleanup + or cleanup.get("result") != "COMPLETE" or cleanup.get("operation_local_cleanup") != "COMPLETE" + or cleanup.get("original_release_run_id") != original_run + or cleanup.get("reconciliation_run_id") != reconciliation_run + or cleanup.get("release_controller_source") != release_controller + or cleanup.get("reconciliation_contract") != "forge-existing-release-reconciliation/v1" + or not isinstance(cleanup_release, dict) or set(cleanup_release) != expected_cleanup_release + or cleanup_release.get("api_url") != release.get("api_url") + or cleanup_release.get("database_id") != database_id + or cleanup_release.get("node_id") != release.get("node_id") + or cleanup_release.get("draft") is not False + or cleanup_release.get("tag") != release.get("tag") + or cleanup_release.get("target_commitish") != request.product_source + or cleanup_release.get("tag_commit") != request.product_source + ): + raise InstalledForgeUpdateError("release-complete publication, policy, or cleanup lineage is noncanonical") + evidence = { + "wheel": str(Path(request.wheel)), "wheel_sha256": request.wheel_sha256, + "wheel_manifest_digest": _digest_bytes(_json_bytes(manifest)), + "receipt": str(receipt_path), "receipt_sha256": request.qualification_receipt_sha256, + "release_operation_id": receipt.get("operation_id"), + "release_controller_source": release_controller, + "original_release_run_id": original_run, "reconciliation_run_id": reconciliation_run, + } + return evidence, wheel_bytes, manifest + + +def validate_qualified_artifact(request: UpdateRequest) -> dict[str, Any]: + evidence, _, _ = _qualified_artifact(request) + return evidence + + +def _sqlite_value(value: object) -> object: + if isinstance(value, bytes): + return {"bytes_hex": value.hex()} + return value + + +def _table_digest(connection: sqlite3.Connection, table: str) -> tuple[int, str]: + quoted = '"' + table.replace('"', '""') + '"' + rows = [ + [_sqlite_value(value) for value in row] + for row in connection.execute(f"SELECT * FROM {quoted}").fetchall() + ] + rows.sort(key=lambda row: json.dumps(row, sort_keys=True, separators=(",", ":"), ensure_ascii=False)) + return len(rows), _digest_bytes(_json_bytes(rows)) + + +def database_snapshot(path: Path, *, existing_connection: sqlite3.Connection | None = None) -> dict[str, Any]: + if path.is_symlink() or not path.is_file(): + raise InstalledForgeUpdateError(f"runtime database is unavailable or unsafe: {path}") + owns_connection = existing_connection is None + try: + connection = existing_connection or sqlite3.connect(path.resolve().as_uri() + "?mode=ro", uri=True) + connection.row_factory = sqlite3.Row + if owns_connection: + connection.execute("PRAGMA query_only=ON") + integrity = connection.execute("PRAGMA integrity_check").fetchone()[0] + foreign_keys = [tuple(row) for row in connection.execute("PRAGMA foreign_key_check")] + user_version = int(connection.execute("PRAGMA user_version").fetchone()[0]) + tables = sorted(row[0] for row in connection.execute( + "SELECT name FROM sqlite_master WHERE type='table' AND name NOT LIKE 'sqlite_%'" + )) + schema_objects = [dict(row) for row in connection.execute( + "SELECT type,name,tbl_name,sql FROM sqlite_master " + "WHERE name NOT LIKE 'sqlite_%' ORDER BY type,name" + )] + table_metrics = { + table: {"count": count, "digest": digest} + for table in tables + for count, digest in (_table_digest(connection, table),) + } + metadata = dict(connection.execute("SELECT key,value FROM runtime_metadata")) + protected_metadata = { + key: value for key, value in metadata.items() if key not in VOLATILE_METADATA_KEYS + } + peer_row = connection.execute( + "SELECT binding_id,configuration_revision,configuration_digest,document " + "FROM execution_host_peer_configuration WHERE singleton=1" + ).fetchone() if "execution_host_peer_configuration" in tables else None + peer = None if peer_row is None else { + "binding_id": peer_row[0], "configuration_revision": peer_row[1], + "configuration_digest": peer_row[2], "document_digest": _digest_bytes(str(peer_row[3]).encode()), + } + dispatcher = [dict(row) for row in connection.execute( + "SELECT status,active_mission_id FROM dispatcher_state" + )] if "dispatcher_state" in tables else [] + missions = [dict(row) for row in connection.execute( + "SELECT mission_id,status FROM mission_state ORDER BY mission_id" + )] if "mission_state" in tables else [] + submissions = [dict(row) for row in connection.execute( + "SELECT submission_id,state FROM scheduler_submissions ORDER BY submission_id" + )] if "scheduler_submissions" in tables else [] + permits = [dict(row) for row in connection.execute( + "SELECT permit_id,state FROM planning_provider_generation_permits ORDER BY permit_id" + )] if "planning_provider_generation_permits" in tables else [] + planning = [dict(row) for row in connection.execute( + "SELECT current_queue,pending_engineering_actions,blocked_engineering_actions FROM planning_state" + )] if "planning_state" in tables else [] + reset = [dict(row) for row in connection.execute( + "SELECT dataset_generation,active_operation_id,state FROM operational_reset_state" + )] if "operational_reset_state" in tables else [] + except sqlite3.Error as error: + raise InstalledForgeUpdateError("runtime database readback failed") from error + finally: + if owns_connection and "connection" in locals(): + connection.close() + snapshot = { + "database": str(path), "integrity_check": integrity, + "foreign_key_check": foreign_keys, "user_version": user_version, + "schema_digest": _digest_bytes(_json_bytes(schema_objects)), + "metadata": metadata, "protected_metadata_digest": _digest_bytes(_json_bytes(protected_metadata)), + "tables": table_metrics, "peer": peer, + "writer_state": { + "dispatcher": dispatcher, "missions": missions, "submissions": submissions, + "generation_permits": permits, "planning": planning, "operational_reset": reset, + }, + } + logical = {key: value for key, value in snapshot.items() if key != "database"} + snapshot["content_digest"] = _digest_bytes(_json_bytes(logical)) + snapshot["snapshot_digest"] = _digest_bytes(_json_bytes(snapshot)) + return snapshot + + +def assert_selected_installation(request: UpdateRequest, snapshot: Mapping[str, Any]) -> None: + metadata = snapshot.get("metadata") + peer = snapshot.get("peer") + if not isinstance(metadata, Mapping): + raise InstalledForgeUpdateError("runtime metadata readback is missing") + if metadata.get("runtime_id") != request.runtime_id: + raise InstalledForgeUpdateError("selected data root belongs to a different runtime") + if metadata.get("installation_id") != request.installation_id: + raise InstalledForgeUpdateError("selected data root belongs to a different installation") + if snapshot.get("user_version") not in {SCHEMA_BEFORE, SCHEMA_AFTER}: + raise InstalledForgeUpdateError("selected runtime schema is outside the bounded update path") + if not isinstance(peer, Mapping) or peer.get("configuration_digest") != request.peer_configuration_digest: + raise InstalledForgeUpdateError("selected peer configuration changed") + marker = Path(request.data_root) / "instance" / "runtime-instance.json" + if marker.is_symlink() or not marker.is_file() or marker.read_text(encoding="utf-8").strip() != request.runtime_id: + raise InstalledForgeUpdateError("runtime instance marker does not bind the selected runtime") + + +def assert_quiescent(snapshot: Mapping[str, Any]) -> None: + writer = snapshot.get("writer_state") + if not isinstance(writer, Mapping): + raise InstalledForgeUpdateError("writer-state readback is missing") + dispatcher = writer.get("dispatcher") + if not isinstance(dispatcher, list) or len(dispatcher) != 1 or any( + row.get("status") != "IDLE" or row.get("active_mission_id") is not None + for row in dispatcher if isinstance(row, Mapping) + ): + raise InstalledForgeUpdateError("Forge dispatcher is not durably idle") + missions = writer.get("missions") or [] + if any(not isinstance(row, Mapping) or row.get("status") not in SAFE_MISSION_STATES for row in missions): + raise InstalledForgeUpdateError("Forge has non-terminal or non-paused Mission state") + submissions = writer.get("submissions") or [] + if any(not isinstance(row, Mapping) or row.get("state") not in SAFE_SUBMISSION_STATES for row in submissions): + raise InstalledForgeUpdateError("Forge has a non-terminal scheduler submission") + permits = writer.get("generation_permits") or [] + if any(isinstance(row, Mapping) and row.get("state") not in TERMINAL_PERMIT_STATES for row in permits): + raise InstalledForgeUpdateError("Forge has an active provider-generation permit") + for row in writer.get("planning") or []: + if not isinstance(row, Mapping): + raise InstalledForgeUpdateError("Forge planning state is malformed") + for key in ("current_queue", "pending_engineering_actions", "blocked_engineering_actions"): + try: + value = json.loads(str(row.get(key))) + except json.JSONDecodeError as error: + raise InstalledForgeUpdateError("Forge planning queue is unreadable") from error + if value: + raise InstalledForgeUpdateError("Forge planning queue is not empty") + reset = writer.get("operational_reset") or [] + if any(not isinstance(row, Mapping) or row.get("active_operation_id") is not None or row.get("state") != "IDLE" + for row in reset): + raise InstalledForgeUpdateError("Forge operational reset maintenance is active") + + +def assert_completed_schema(snapshot: Mapping[str, Any], installed_readback: Mapping[str, Any]) -> None: + tables = snapshot.get("tables") + expected_digest = installed_readback.get("database_schema_digest") + if ( + snapshot.get("user_version") != SCHEMA_AFTER + or not isinstance(tables, Mapping) + or not NEW_SCHEMA_38_TABLES.issubset(tables) + or not isinstance(expected_digest, str) + or snapshot.get("schema_digest") != expected_digest + ): + raise InstalledForgeUpdateError("completed runtime schema changed from the activated schema 38") + + +def verify_preservation(before: Mapping[str, Any], after: Mapping[str, Any], request: UpdateRequest) -> dict[str, Any]: + if after.get("integrity_check") != "ok" or after.get("foreign_key_check") != []: + raise InstalledForgeUpdateError("migrated runtime failed SQLite integrity validation") + if after.get("user_version") != SCHEMA_AFTER: + raise InstalledForgeUpdateError("Forge owning migration did not reach schema 38") + before_metadata, after_metadata = before.get("metadata"), after.get("metadata") + if not isinstance(before_metadata, Mapping) or not isinstance(after_metadata, Mapping): + raise InstalledForgeUpdateError("migration metadata readback is incomplete") + for key, expected in ( + ("runtime_id", request.runtime_id), ("installation_id", request.installation_id), + ): + if before_metadata.get(key) != expected or after_metadata.get(key) != expected: + raise InstalledForgeUpdateError(f"migration changed selected {key}") + if after.get("protected_metadata_digest") != before.get("protected_metadata_digest"): + raise InstalledForgeUpdateError("migration changed protected runtime metadata") + if after.get("peer") != before.get("peer"): + raise InstalledForgeUpdateError("migration changed the configured Forge-to-EP peer binding") + before_tables, after_tables = before.get("tables"), after.get("tables") + if not isinstance(before_tables, Mapping) or not isinstance(after_tables, Mapping): + raise InstalledForgeUpdateError("migration table readback is incomplete") + for table, metric in before_tables.items(): + if table == "runtime_metadata": + continue + if after_tables.get(table) != metric: + raise InstalledForgeUpdateError(f"migration changed historical table contents: {table}") + new_tables = set(after_tables) - set(before_tables) + if new_tables != set(NEW_SCHEMA_38_TABLES): + raise InstalledForgeUpdateError("migration produced an unexpected schema-38 table set") + expected_counts = {table: 0 for table in NEW_SCHEMA_38_TABLES} + expected_counts["operational_reset_state"] = 1 + if any(after_tables[table]["count"] != count for table, count in expected_counts.items()): + raise InstalledForgeUpdateError("migration initialized unexpected operational-reset data") + reset = after.get("writer_state", {}).get("operational_reset", []) + if reset != [{"dataset_generation": 0, "active_operation_id": None, "state": "IDLE"}]: + raise InstalledForgeUpdateError("schema-38 reset state is not an idle, fresh control record") + return { + "status": "PASS", "from_schema": SCHEMA_BEFORE, "to_schema": SCHEMA_AFTER, + "preserved_table_count": len(before_tables) - 1, + "added_tables": sorted(new_tables), + "protected_metadata_digest": after.get("protected_metadata_digest"), + "peer_configuration_digest": request.peer_configuration_digest, + } + + +@contextmanager +def exclusive_lock(path: Path) -> Iterator[None]: + if fcntl is None: + raise InstalledForgeUpdateError("POSIX installation locking is unavailable") + _safe_directory(path.parent, create=True) + if path.is_symlink(): + raise InstalledForgeUpdateError(f"installation lock path is unsafe: {path}") + with path.open("a+", encoding="utf-8") as handle: + try: + fcntl.flock(handle.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB) + except BlockingIOError as error: + raise InstalledForgeUpdateError(f"concurrent maintenance owns lock: {path}") from error + try: + yield + finally: + fcntl.flock(handle.fileno(), fcntl.LOCK_UN) + + +def _copy_sqlite_backup(source: Path, destination: Path) -> dict[str, Any]: + _safe_directory(destination.parent, create=True) + if destination.exists() or destination.is_symlink(): + raise InstalledForgeUpdateError("installation backup already exists without matching durable evidence") + descriptor, name = tempfile.mkstemp(prefix=f".{destination.name}.tmp-", dir=destination.parent) + os.close(descriptor) + temporary = Path(name) + try: + source_connection = sqlite3.connect(source.resolve().as_uri() + "?mode=ro", uri=True) + destination_connection = sqlite3.connect(temporary) + source_connection.backup(destination_connection) + destination_connection.commit() + if destination_connection.execute("PRAGMA integrity_check").fetchone()[0] != "ok": + raise InstalledForgeUpdateError("SQLite backup integrity check failed") + if destination_connection.execute("PRAGMA foreign_key_check").fetchone() is not None: + raise InstalledForgeUpdateError("SQLite backup foreign-key check failed") + destination_connection.close() + source_connection.close() + os.chmod(temporary, 0o600) + os.replace(temporary, destination) + except sqlite3.Error as error: + raise InstalledForgeUpdateError("consistent SQLite backup failed") from error + finally: + for connection_name in ("destination_connection", "source_connection"): + connection = locals().get(connection_name) + if connection is not None: + try: + connection.close() + except sqlite3.Error: + pass + if temporary.exists(): + temporary.unlink() + verification = database_snapshot(destination) + return { + "path": str(destination), "sha256": file_digest(destination), + "size": destination.stat().st_size, "integrity_check": verification["integrity_check"], + "foreign_key_check": verification["foreign_key_check"], + "snapshot_digest": verification["snapshot_digest"], + } + + +def _candidate_migrate(executable: Path, data_root: Path, *, cwd: Path) -> dict[str, Any]: + result = _run((str(executable), "--data-root", str(data_root), "server", "init"), cwd=cwd) + try: + output = json.loads(result.stdout) + except json.JSONDecodeError as error: + raise InstalledForgeUpdateError("candidate migration readback is malformed") from error + if output.get("storage_schema") != str(SCHEMA_AFTER) or output.get("initialized") is not True: + raise InstalledForgeUpdateError("candidate migration did not return schema-38 installed readback") + return output + + +def _candidate_site_packages(slot: Path) -> Path: + candidates = [path for path in (slot / "lib").glob("python*/site-packages") if path.is_dir()] + if len(candidates) != 1: + raise InstalledForgeUpdateError("candidate virtual environment has ambiguous site-packages") + return _safe_directory(candidates[0]) + + +def _entrypoint_bytes(slot: Path) -> bytes: + return ( + "#!/bin/sh\n" + f"exec {shlex.quote(str(slot / 'bin' / 'python'))} -B -m forge \"$@\"\n" + ).encode("utf-8") + + +def _install_validated_wheel(slot: Path, wheel_bytes: bytes, manifest: Mapping[str, str]) -> None: + site_packages = _candidate_site_packages(slot) + if any(site_packages.iterdir()): + raise InstalledForgeUpdateError("fresh candidate site-packages is not empty") + with zipfile.ZipFile(BytesIO(wheel_bytes)) as archive: + for info in archive.infolist(): + target = site_packages.joinpath(*Path(info.filename).parts) + if info.is_dir(): + _safe_directory(target, create=True) + continue + _safe_directory(target.parent, create=True) + descriptor = os.open( + target, + os.O_WRONLY | os.O_CREAT | os.O_EXCL | getattr(os, "O_NOFOLLOW", 0), + 0o600, + ) + with os.fdopen(descriptor, "wb") as handle: + handle.write(archive.read(info.filename)) + handle.flush() + os.fsync(handle.fileno()) + entrypoint = slot / "bin" / "forge" + descriptor = os.open( + entrypoint, + os.O_WRONLY | os.O_CREAT | os.O_EXCL | getattr(os, "O_NOFOLLOW", 0), + 0o700, + ) + with os.fdopen(descriptor, "wb") as handle: + handle.write(_entrypoint_bytes(slot)) + handle.flush() + os.fsync(handle.fileno()) + _verify_candidate_files(slot, manifest) + + +def _verify_candidate_files(slot: Path, manifest: Mapping[str, str]) -> dict[str, Any]: + site_packages = _candidate_site_packages(slot) + for cache in list(site_packages.rglob("__pycache__")): + _assert_no_symlink_components(cache) + if not cache.is_dir(): + raise InstalledForgeUpdateError("candidate bytecode cache path is unsafe") + for child in cache.rglob("*"): + _assert_no_symlink_components(child) + shutil.rmtree(cache) + actual: dict[str, str] = {} + for path in site_packages.rglob("*"): + _assert_no_symlink_components(path) + if path.is_file(): + relative = path.relative_to(site_packages).as_posix() + actual[relative] = file_digest(path) + if actual != dict(manifest): + expected = dict(manifest) + missing = sorted(set(expected) - set(actual)) + unexpected = sorted(set(actual) - set(expected)) + changed = sorted(name for name in set(actual) & set(expected) if actual[name] != expected[name]) + raise InstalledForgeUpdateError( + "candidate installed files do not match the exact wheel manifest: " + f"missing={missing[:5]}, unexpected={unexpected[:5]}, changed={changed[:5]}" + ) + entrypoint = slot / "bin" / "forge" + if _read_regular_bytes(entrypoint) != _entrypoint_bytes(slot): + raise InstalledForgeUpdateError("candidate command entry point changed") + return { + "wheel_manifest_digest": _digest_bytes(_json_bytes(dict(manifest))), + "installed_file_count": len(actual), + "entrypoint_sha256": file_digest(entrypoint), + } + + +class InstalledForgeUpdateController: + def __init__( + self, + request: UpdateRequest, + *, + process_reader: Callable[[], Sequence[str]] | None = None, + interrupt_after: str | None = None, + ) -> None: + request.validate() + self.request = request + self.data_root = Path(request.data_root) + self.runtime_root = Path(request.runtime_root) + self.database = self.data_root / "forge.db" + self.operation_root = self.data_root / "artifacts" / "installation" / request.operation_id + self.state_path = self.operation_root / "operation.json" + self.receipt_path = self.operation_root / "receipt.json" + self.backup_root = self.data_root / "backups" / "installation" / request.operation_id + self.backup_path = self.backup_root / "forge-schema37.sqlite3" + self.slot = self.runtime_root / "slots" / f"{request.version}-{request.wheel_sha256.removeprefix('sha256:')[:12]}" + self.slot_receipt = self.slot / "forge-installation-slot.json" + self.slot_claim = self.slot.parent / f".{self.slot.name}.{request.operation_id}.owner.json" + self.current = self.runtime_root / "current" + self.stable_resolver = self.runtime_root / "bin" / "forge" + self.fenced_resolver = self.runtime_root / "fenced" / "forge" + self.legacy_entrypoint = self.runtime_root / "legacy" / ( + request.resolver_sha256.removeprefix("sha256:")[:16] + "-forge" + ) + self.process_reader = process_reader or self._processes + self.interrupt_after = interrupt_after + + def _state(self) -> dict[str, Any]: + if self.state_path.exists(): + state = _read_json(self.state_path) + if state.get("contract_version") != CONTRACT_VERSION or state.get("request") != asdict(self.request): + raise InstalledForgeUpdateError("durable update operation conflicts with the requested target") + return state + state = { + "contract_version": CONTRACT_VERSION, "operation_id": self.request.operation_id, + "request": asdict(self.request), "request_digest": self.request.digest, + "phase": "PREPARED", "created_at": _now(), "updated_at": _now(), + "history": [{"phase": "PREPARED", "at": _now()}], + } + _atomic_json(self.state_path, state) + return state + + def _advance(self, state: dict[str, Any], phase: str, **evidence: object) -> dict[str, Any]: + current_phase = str(state.get("phase", "PREPARED")) + effective_phase = phase + if PHASE_ORDER.get(current_phase, -1) > PHASE_ORDER.get(phase, -1): + effective_phase = current_phase + updated = {**state, **evidence, "phase": effective_phase, "updated_at": _now()} + history = list(state.get("history", [])) + if effective_phase != current_phase or phase == current_phase: + history.append({"phase": effective_phase, "at": updated["updated_at"]}) + updated["history"] = history + _atomic_json(self.state_path, updated) + return updated + + def _interrupt(self, point: str) -> None: + if self.interrupt_after == point: + raise InstalledForgeUpdateError(f"simulated interruption after {point}") + + def _processes(self) -> Sequence[str]: + result = _run(("/bin/ps", "-axo", "pid=,ppid=,command="), cwd=self.runtime_root, timeout=30) + processes: list[str] = [] + own = {os.getpid(), os.getppid()} + for line in result.stdout.splitlines(): + fields = line.strip().split(None, 2) + if len(fields) != 3: + continue + try: + pid, parent = int(fields[0]), int(fields[1]) + except ValueError: + continue + if pid in own: + continue + processes.append(fields[2]) + return processes + + def _assert_no_runtime_process(self) -> None: + needles = ( + self.request.data_root, self.request.resolver, self.request.existing_interpreter, + str(self.runtime_root / "venv"), str(self.slot), + ) + matches = [command for command in self.process_reader() if any(needle in command for needle in needles)] + if matches: + raise InstalledForgeUpdateError("a selected Forge runtime process is still active") + + def _stage(self, state: dict[str, Any]) -> dict[str, Any]: + qualification, wheel_bytes, manifest = _qualified_artifact(self.request) + if file_digest(Path(__file__)) != self.request.controller_sha256: + raise InstalledForgeUpdateError("installation controller bytes do not match the protected candidate") + if self.slot_receipt.exists(): + receipt = _read_json(self.slot_receipt) + if ( + receipt.get("request_digest") != self.request.digest + or receipt.get("wheel_manifest_digest") != qualification["wheel_manifest_digest"] + ): + raise InstalledForgeUpdateError("immutable candidate slot belongs to a different request") + file_evidence = _verify_candidate_files(self.slot, manifest) + if receipt.get("installed_files") != file_evidence: + raise InstalledForgeUpdateError("immutable candidate slot evidence changed") + else: + # Venv entry points embed their creation path. Claim the final + # path first; an interrupted, unreceipted slot is quarantined and + # rebuilt from the pinned bytes rather than trusted or relocated. + _safe_directory(self.slot.parent, create=True) + if self.slot.is_symlink(): + raise InstalledForgeUpdateError("candidate staging path is unsafe") + claim = { + "contract_version": CONTRACT_VERSION, + "request_digest": self.request.digest, + "operation_id": self.request.operation_id, + } + claim_preexisting = self.slot_claim.exists() + if claim_preexisting: + observed_claim = _read_json(self.slot_claim) + if any(observed_claim.get(key) != value for key, value in claim.items()): + raise InstalledForgeUpdateError("candidate slot claim belongs to a different request") + else: + _atomic_json(self.slot_claim, {**claim, "created_at": _now()}) + staging_owner = self.slot / "forge-installation-staging.json" + if self.slot.exists(): + if not self.slot.is_dir(): + raise InstalledForgeUpdateError("unreceipted candidate slot is unsafe") + if staging_owner.is_file(): + owner = _read_json(staging_owner) + if any(owner.get(key) != value for key, value in claim.items()): + raise InstalledForgeUpdateError("unreceipted candidate slot belongs to a different request") + elif not claim_preexisting or any(self.slot.iterdir()): + raise InstalledForgeUpdateError("unreceipted candidate slot has no matching operation owner") + abandoned = self.slot.parent / f".{self.slot.name}.abandoned-{secrets.token_hex(16)}" + os.replace(self.slot, abandoned) + self.slot.mkdir(mode=0o700) + _atomic_json(staging_owner, {**claim, "created_at": _now()}) + _run((self.request.base_python, "-m", "venv", "--without-pip", str(self.slot)), cwd=self.runtime_root) + candidate_python = self.slot / "bin" / "python" + _install_validated_wheel(self.slot, wheel_bytes, manifest) + file_evidence = _verify_candidate_files(self.slot, manifest) + identity = installed_identity(candidate_python, cwd=self.runtime_root) + if ( + identity.get("version") != self.request.version + or identity.get("distribution_version") != self.request.version + or not str(identity.get("module", "")).startswith(str(self.slot.resolve()) + os.sep) + or Path(str(identity.get("prefix"))).resolve() != self.slot.resolve() + ): + raise InstalledForgeUpdateError("candidate slot identity is inconsistent") + _atomic_json(self.slot_receipt, { + "contract_version": CONTRACT_VERSION, "request_digest": self.request.digest, + "wheel_sha256": self.request.wheel_sha256, "product_source": self.request.product_source, + "version": self.request.version, "identity": identity, + "wheel_manifest_digest": qualification["wheel_manifest_digest"], + "installed_files": file_evidence, "staged_at": _now(), + }) + if self.slot_claim.is_file(): + self.slot_claim.unlink() + file_evidence = _verify_candidate_files(self.slot, manifest) + identity = installed_identity(self.slot / "bin" / "python", cwd=self.runtime_root) + if identity.get("version") != self.request.version or not str(identity.get("module", "")).startswith(str(self.slot.resolve()) + os.sep): + raise InstalledForgeUpdateError("staged candidate readback changed") + return self._advance(state, "STAGED", artifact_qualification=qualification, candidate=identity, + candidate_slot=str(self.slot), installed_files=file_evidence) + + def _adopt_resolver(self, state: dict[str, Any]) -> dict[str, Any]: + resolver = Path(self.request.resolver) + _safe_directory(self.runtime_root / "legacy", create=True) + _safe_directory(self.runtime_root / "bin", create=True) + _safe_directory(self.runtime_root / "fenced", create=True) + if not self.legacy_entrypoint.exists(): + legacy_bytes = _read_regular_bytes(resolver) + if _digest_bytes(legacy_bytes) != self.request.resolver_sha256: + raise InstalledForgeUpdateError("legacy command resolver changed before adoption") + identity = installed_identity(Path(self.request.existing_interpreter), cwd=self.runtime_root) + if identity.get("version") != self.request.existing_version: + raise InstalledForgeUpdateError("legacy interpreter no longer provides the selected Forge version") + descriptor = os.open( + self.legacy_entrypoint, + os.O_WRONLY | os.O_CREAT | os.O_EXCL | getattr(os, "O_NOFOLLOW", 0), + 0o700, + ) + with os.fdopen(descriptor, "wb") as handle: + handle.write(legacy_bytes) + handle.flush() + os.fsync(handle.fileno()) + if file_digest(self.legacy_entrypoint) != self.request.resolver_sha256: + raise InstalledForgeUpdateError("retained legacy entrypoint does not match the pinned resolver bytes") + elif file_digest(self.legacy_entrypoint) != self.request.resolver_sha256: + raise InstalledForgeUpdateError("retained legacy entrypoint changed") + fence = ( + "#!/bin/sh\n" + f"echo 'Forge installation maintenance is active: {self.request.operation_id}' >&2\n" + "exit 75\n" + ) + if not self.fenced_resolver.exists(): + descriptor = os.open(self.fenced_resolver, os.O_WRONLY | os.O_CREAT | os.O_EXCL, 0o755) + with os.fdopen(descriptor, "w", encoding="utf-8") as handle: + handle.write(fence) + elif self.fenced_resolver.read_text(encoding="utf-8") != fence: + raise InstalledForgeUpdateError("maintenance fence launcher changed") + if not self.current.exists() and not self.current.is_symlink(): + _replace_symlink(self.current, os.path.relpath(self.legacy_entrypoint, self.runtime_root)) + if not self.stable_resolver.exists() and not self.stable_resolver.is_symlink(): + _replace_symlink(self.stable_resolver, "../current") + elif _resolved_link(self.stable_resolver) != _resolved_link(self.current): + # The stable link follows current; compare its textual contract, not + # a transient current target, when it already exists. + if os.readlink(self.stable_resolver) != "../current": + raise InstalledForgeUpdateError("stable Forge resolver changed") + expected_external = self.stable_resolver.resolve(strict=True) + if resolver.is_symlink(): + if resolver.resolve(strict=True) != expected_external: + raise InstalledForgeUpdateError("external Forge resolver was retargeted") + else: + if file_digest(resolver) != self.request.resolver_sha256: + raise InstalledForgeUpdateError("external Forge resolver changed") + _replace_symlink(resolver, str(self.stable_resolver)) + return self._advance(state, "ADOPTED", resolver={ + "external": str(resolver), "stable": str(self.stable_resolver), + "legacy": str(self.legacy_entrypoint), "legacy_sha256": self.request.resolver_sha256, + }) + + def _backup(self, state: dict[str, Any], before: Mapping[str, Any]) -> dict[str, Any]: + existing = state.get("backup") + if isinstance(existing, Mapping): + if existing.get("path") != str(self.backup_path) or file_digest(self.backup_path) != existing.get("sha256"): + raise InstalledForgeUpdateError("durable installation backup changed") + backup = dict(existing) + recovered = database_snapshot(self.backup_path) + if recovered.get("content_digest") != before.get("content_digest"): + raise InstalledForgeUpdateError("durable installation backup no longer matches the pre-migration snapshot") + elif self.backup_path.exists() and not self.backup_path.is_symlink(): + recovered = database_snapshot(self.backup_path) + if ( + recovered.get("user_version") != SCHEMA_BEFORE + or recovered.get("integrity_check") != "ok" + or recovered.get("foreign_key_check") != [] + or recovered.get("content_digest") != before.get("content_digest") + ): + raise InstalledForgeUpdateError("unreceipted installation backup cannot be safely adopted") + backup = { + "path": str(self.backup_path), "sha256": file_digest(self.backup_path), + "size": self.backup_path.stat().st_size, + "integrity_check": recovered["integrity_check"], + "foreign_key_check": recovered["foreign_key_check"], + "snapshot_digest": recovered["snapshot_digest"], + "created_at": _now(), "recovered_after_atomic_backup_write": True, + "source_wal_present": None, "source_shm_present": None, + } + else: + backup = _copy_sqlite_backup(self.database, self.backup_path) + backup["created_at"] = _now() + backup["source_wal_present"] = self.database.with_name(self.database.name + "-wal").exists() + backup["source_shm_present"] = self.database.with_name(self.database.name + "-shm").exists() + return self._advance(state, "BACKED_UP", before=before, backup=backup) + + def _qualify_copy(self, state: dict[str, Any], before: Mapping[str, Any]) -> dict[str, Any]: + existing = state.get("migration_qualification") + if isinstance(existing, Mapping) and existing.get("status") == "PASS": + database = self.operation_root / "qualification-copy" / "forge.db" + qualified = database_snapshot(database) + verify_preservation(before, qualified, self.request) + if existing.get("after_snapshot_digest") != qualified.get("snapshot_digest"): + raise InstalledForgeUpdateError("durable migration qualification copy changed") + return state + root = self.operation_root / "qualification-copy" + _safe_directory(root, create=True) + database = root / "forge.db" + descriptor, name = tempfile.mkstemp(prefix=".forge.db.tmp-", dir=root) + os.close(descriptor) + temporary = Path(name) + shutil.copyfile(self.backup_path, temporary) + os.chmod(temporary, 0o600) + os.replace(temporary, database) + for suffix in ("-wal", "-shm"): + sidecar = root / ("forge.db" + suffix) + if sidecar.is_symlink(): + raise InstalledForgeUpdateError("qualification copy contains an unsafe SQLite sidecar") + if sidecar.exists() and not sidecar.is_symlink(): + sidecar.unlink() + instance = _safe_directory(root / "instance", create=True) + marker = instance / "runtime-instance.json" + marker_payload = (self.request.runtime_id + "\n").encode("utf-8") + if marker.exists() or marker.is_symlink(): + if _read_regular_bytes(marker) != marker_payload: + raise InstalledForgeUpdateError("qualification-copy runtime marker changed") + else: + descriptor = os.open( + marker, os.O_WRONLY | os.O_CREAT | os.O_EXCL | getattr(os, "O_NOFOLLOW", 0), 0o600, + ) + with os.fdopen(descriptor, "wb") as handle: + handle.write(marker_payload) + handle.flush() + os.fsync(handle.fileno()) + copy_before = database_snapshot(database) + if copy_before.get("content_digest") != before.get("content_digest"): + raise InstalledForgeUpdateError("isolated qualification copy does not match the consistent backup") + candidate_output = _candidate_migrate(self.slot / "bin" / "forge", root, cwd=self.runtime_root) + copy_after = database_snapshot(database) + qualification = verify_preservation(copy_before, copy_after, self.request) + qualification.update({ + "qualified_at": _now(), "copy_root": str(root), + "candidate_output": candidate_output, + "before_snapshot_digest": copy_before["snapshot_digest"], + "after_snapshot_digest": copy_after["snapshot_digest"], + }) + return self._advance(state, "MIGRATION_QUALIFIED", migration_qualification=qualification) + + def _fence(self, state: dict[str, Any]) -> dict[str, Any]: + _replace_symlink(self.current, os.path.relpath(self.fenced_resolver, self.runtime_root)) + if _resolved_link(Path(self.request.resolver)) != self.fenced_resolver.resolve(): + raise InstalledForgeUpdateError("external resolver did not enter the maintenance fence") + return self._advance(state, "FENCED", safety_disposition="LEGACY_COMMAND_FENCED") + + def _restore_legacy_before_migration(self, state: dict[str, Any], error: Exception) -> None: + _replace_symlink(self.current, os.path.relpath(self.legacy_entrypoint, self.runtime_root)) + self._advance( + state, state.get("phase", "RECOVERY_PENDING"), + safety_disposition="LEGACY_RESTORED_BEFORE_MIGRATION", last_error=str(error), + ) + + def _install_qualified_database(self, before: Mapping[str, Any]) -> dict[str, Any]: + source = self.operation_root / "qualification-copy" / "forge.db" + qualified = database_snapshot(source) + verify_preservation(before, qualified, self.request) + descriptor, name = tempfile.mkstemp(prefix=".forge.db.install-", dir=self.data_root) + os.close(descriptor) + temporary = Path(name) + live_connection: sqlite3.Connection | None = None + source_connection: sqlite3.Connection | None = None + destination_connection: sqlite3.Connection | None = None + try: + live_connection = sqlite3.connect(self.database) + live_connection.row_factory = sqlite3.Row + live_connection.execute("PRAGMA busy_timeout=0") + journal_mode = live_connection.execute("PRAGMA journal_mode=DELETE").fetchone()[0] + if str(journal_mode).lower() != "delete": + raise InstalledForgeUpdateError("live runtime journal could not enter crash-safe swap mode") + live_connection.execute("BEGIN EXCLUSIVE") + locked_live = database_snapshot(self.database, existing_connection=live_connection) + if ( + locked_live.get("user_version") != SCHEMA_BEFORE + or locked_live.get("content_digest") != before.get("content_digest") + ): + raise InstalledForgeUpdateError("live runtime changed after the qualified backup") + source_connection = sqlite3.connect(source.resolve().as_uri() + "?mode=ro", uri=True) + destination_connection = sqlite3.connect(temporary) + source_connection.backup(destination_connection) + destination_connection.commit() + destination_mode = destination_connection.execute("PRAGMA journal_mode=DELETE").fetchone()[0] + if str(destination_mode).lower() != "delete": + raise InstalledForgeUpdateError("migrated database copy has an unsafe journal mode") + destination_connection.close() + destination_connection = None + source_connection.close() + source_connection = None + installed_copy = database_snapshot(temporary) + verify_preservation(before, installed_copy, self.request) + os.chmod(temporary, 0o400) + for suffix in ("-wal", "-shm"): + sidecar = self.database.with_name(self.database.name + suffix) + if sidecar.is_symlink(): + raise InstalledForgeUpdateError("live database has an unsafe SQLite sidecar") + if sidecar.exists(): + sidecar.unlink() + self._interrupt("database_swap_prepared") + os.replace(temporary, self.database) + self._interrupt("database_swap") + directory = os.open(self.data_root, os.O_RDONLY) + try: + os.fsync(directory) + finally: + os.close(directory) + except sqlite3.Error as error: + raise InstalledForgeUpdateError("exclusive atomic installation of the migrated database failed") from error + finally: + for connection in (destination_connection, source_connection, live_connection): + if connection is not None: + try: + connection.close() + except sqlite3.Error: + pass + if temporary.exists() and not temporary.is_symlink(): + temporary.unlink() + installed = database_snapshot(self.database) + verify_preservation(before, installed, self.request) + return installed + + def _migrate_live(self, state: dict[str, Any], before: Mapping[str, Any]) -> tuple[dict[str, Any], dict[str, Any]]: + current = database_snapshot(self.database) + if current.get("user_version") == SCHEMA_BEFORE: + state = self._fence(state) + self._interrupt("fence") + after = self._install_qualified_database(before) + preservation = verify_preservation(before, after, self.request) + state = self._advance( + state, "MIGRATED", live_migration={ + **preservation, "migrated_at": _now(), + "application_mode": "ATOMIC_PRODUCT_MIGRATED_COPY", + "before_snapshot_digest": before["snapshot_digest"], + "after_snapshot_digest": after["snapshot_digest"], + }, safety_disposition="CANDIDATE_REQUIRED_SCHEMA_38", + ) + self._interrupt("migration") + return state, after + if current.get("user_version") == SCHEMA_AFTER: + preservation = verify_preservation(before, current, self.request) + os.chmod(self.database, 0o400) + if state.get("phase") not in {"MIGRATED", "ACTIVATING", "ACTIVATED", "COMPLETE"}: + state = self._advance( + state, "MIGRATED", live_migration={ + **preservation, "reconciled_at": _now(), + "before_snapshot_digest": before["snapshot_digest"], + "after_snapshot_digest": current["snapshot_digest"], + }, safety_disposition="CANDIDATE_REQUIRED_SCHEMA_38", + ) + return state, current + raise InstalledForgeUpdateError("live runtime schema changed outside the bounded operation") + + def _activate(self, state: dict[str, Any], after: Mapping[str, Any]) -> dict[str, Any]: + state = self._advance(state, "ACTIVATING", safety_disposition="CANDIDATE_ACTIVATION_IN_PROGRESS") + candidate = self.slot / "bin" / "forge" + _replace_symlink(self.current, os.path.relpath(candidate, self.runtime_root)) + self._interrupt("activation") + if _resolved_link(Path(self.request.resolver)) != candidate.resolve(): + raise InstalledForgeUpdateError("external command resolver did not activate the candidate") + identity = installed_identity(self.slot / "bin" / "python", cwd=self.runtime_root) + version = _run((self.request.resolver, "--version"), cwd=self.runtime_root).stdout.strip() + status_result = _run( + (self.request.resolver, "--data-root", self.request.data_root, "server", "status"), + cwd=self.runtime_root, + ) + try: + status = json.loads(status_result.stdout) + except json.JSONDecodeError as error: + raise InstalledForgeUpdateError("activated CLI status readback is malformed") from error + if ( + identity.get("version") != self.request.version + or version != self.request.version + or Path(str(identity.get("sys_executable", ""))).resolve() != (self.slot / "bin" / "python").resolve() + or not str(identity.get("module", "")).startswith(str(self.slot.resolve()) + os.sep) + or status.get("product_version") != self.request.version + or Path(str(status.get("data_root", ""))).resolve() != self.data_root.resolve() + or status.get("instance_id") != self.request.runtime_id + or status.get("storage_schema") != str(SCHEMA_AFTER) + ): + raise InstalledForgeUpdateError("activated Forge CLI readback does not match the selected installation") + final_snapshot = database_snapshot(self.database) + preservation = verify_preservation(state["before"], final_snapshot, self.request) + return self._advance( + state, "ACTIVATED", installed_readback={ + "identity": identity, "cli_version": version, "status": status, + "database_snapshot_digest": final_snapshot["snapshot_digest"], + "database_schema_digest": final_snapshot["schema_digest"], + "preservation": preservation, "resolver": self.request.resolver, + "resolved_executable": str(candidate.resolve()), + }, safety_disposition="CANDIDATE_ACTIVE", + ) + + def _secure_failure(self, state: dict[str, Any], error: Exception) -> None: + """Leave a pre-migration legacy route or a schema-38-safe candidate/fence.""" + try: + state = self._state() + except Exception: + pass + try: + current = database_snapshot(self.database) + except Exception: + current = {} + before = state.get("before") + if ( + current.get("user_version") == SCHEMA_BEFORE + and isinstance(before, Mapping) + and current.get("content_digest") == before.get("content_digest") + and self.legacy_entrypoint.exists() + ): + self._restore_legacy_before_migration(state, error) + return + _replace_symlink(self.current, os.path.relpath(self.fenced_resolver, self.runtime_root)) + self._advance( + state, state.get("phase", "RECOVERY_PENDING"), + safety_disposition="UNVERIFIED_OR_MIGRATED_RUNTIME_FENCED", last_error=str(error), + ) + + def _verify_complete(self, state: Mapping[str, Any], receipt: Mapping[str, Any]) -> None: + if state.get("phase") != "COMPLETE" or state.get("request_digest") != self.request.digest: + raise InstalledForgeUpdateError("completed operation state conflicts with this request") + if state.get("receipt_sha256") != file_digest(self.receipt_path): + raise InstalledForgeUpdateError("completed operation receipt changed") + qualification, _, manifest = _qualified_artifact(self.request) + if file_digest(Path(__file__)) != self.request.controller_sha256: + raise InstalledForgeUpdateError("completed operation controller bytes changed") + if receipt.get("request_digest") != self.request.digest or receipt.get("state") != "COMPLETE": + raise InstalledForgeUpdateError("completed update receipt conflicts with this request") + backup = receipt.get("backup") + if ( + not isinstance(backup, Mapping) + or backup.get("path") != str(self.backup_path) + or backup.get("sha256") != file_digest(self.backup_path) + ): + raise InstalledForgeUpdateError("completed operation backup changed") + slot_receipt = _read_json(self.slot_receipt) + file_evidence = _verify_candidate_files(self.slot, manifest) + if ( + slot_receipt.get("request_digest") != self.request.digest + or slot_receipt.get("wheel_manifest_digest") != qualification["wheel_manifest_digest"] + or slot_receipt.get("installed_files") != file_evidence + ): + raise InstalledForgeUpdateError("completed operation candidate slot changed") + candidate = self.slot / "bin" / "forge" + if _resolved_link(Path(self.request.resolver)) != candidate.resolve(): + raise InstalledForgeUpdateError("completed operation resolver no longer selects the candidate") + identity = installed_identity(self.slot / "bin" / "python", cwd=self.runtime_root) + version = _run((self.request.resolver, "--version"), cwd=self.runtime_root).stdout.strip() + status_result = _run( + (self.request.resolver, "--data-root", self.request.data_root, "server", "status"), + cwd=self.runtime_root, + ) + try: + status = json.loads(status_result.stdout) + except json.JSONDecodeError as error: + raise InstalledForgeUpdateError("completed installed CLI readback is malformed") from error + if ( + identity.get("version") != self.request.version + or identity.get("distribution_version") != self.request.version + or not str(identity.get("module", "")).startswith(str(self.slot.resolve()) + os.sep) + or version != self.request.version + or status.get("product_version") != self.request.version + or status.get("instance_id") != self.request.runtime_id + or status.get("storage_schema") != str(SCHEMA_AFTER) + or Path(str(status.get("data_root", ""))).resolve() != self.data_root.resolve() + ): + raise InstalledForgeUpdateError("completed installed CLI identity changed") + final_snapshot = database_snapshot(self.database) + assert_selected_installation(self.request, final_snapshot) + if final_snapshot.get("integrity_check") != "ok" or final_snapshot.get("foreign_key_check") != []: + raise InstalledForgeUpdateError("completed runtime database integrity changed") + before = state.get("before") + if not isinstance(before, Mapping): + raise InstalledForgeUpdateError("completed operation lacks its pre-migration snapshot") + backup_snapshot = database_snapshot(self.backup_path) + if backup_snapshot.get("content_digest") != before.get("content_digest"): + raise InstalledForgeUpdateError("completed operation backup no longer matches its source snapshot") + for key in ("migration_qualification", "live_migration"): + evidence = receipt.get(key) + if not isinstance(evidence, Mapping) or evidence.get("status") != "PASS": + raise InstalledForgeUpdateError("completed operation lacks successful migration evidence") + installed_readback = receipt.get("installed_readback") + if ( + not isinstance(installed_readback, Mapping) + or not isinstance(installed_readback.get("preservation"), Mapping) + or installed_readback["preservation"].get("status") != "PASS" + ): + raise InstalledForgeUpdateError("completed operation lacks successful installation readback") + assert_completed_schema(final_snapshot, installed_readback) + + def _restore_database_writable(self) -> None: + _assert_no_symlink_components(self.database) + if not self.database.is_file(): + raise InstalledForgeUpdateError("installed database is unavailable after activation") + os.chmod(self.database, 0o600) + + def run(self) -> dict[str, Any]: + _safe_directory(self.data_root) + _safe_directory(self.runtime_root) + update_lock = self.runtime_root / "locks" / "installation-update.lock" + runtime_lock = self.data_root / "forge-runtime-mutation.lock" + bootstrap_lock = self.data_root / "locks" / "runtime.lock" + with exclusive_lock(update_lock): + _safe_directory(self.operation_root, create=True) + os.chmod(self.operation_root, 0o700) + state = self._state() + if state.get("phase") == "COMPLETE": + receipt = _read_json(self.receipt_path) + self._verify_complete(state, receipt) + self._restore_database_writable() + return receipt + + state = self._stage(state) + self._interrupt("stage") + with exclusive_lock(runtime_lock), exclusive_lock(bootstrap_lock): + self._assert_no_runtime_process() + live = database_snapshot(self.database) + assert_selected_installation(self.request, live) + assert_quiescent(live) + try: + state = self._adopt_resolver(state) + self._interrupt("adoption") + before = state.get("before") + if not isinstance(before, Mapping): + if live.get("user_version") != SCHEMA_BEFORE: + raise InstalledForgeUpdateError("schema 38 lacks this operation's pre-migration snapshot") + before = live + assert_selected_installation(self.request, before) + if before.get("user_version") != SCHEMA_BEFORE: + raise InstalledForgeUpdateError("durable pre-migration snapshot is not schema 37") + state = self._backup(state, before) + self._interrupt("backup") + state = self._qualify_copy(state, before) + self._interrupt("qualification") + state, after = self._migrate_live(state, before) + state = self._activate(state, after) + receipt = { + "contract_version": CONTRACT_VERSION, + "operation_id": self.request.operation_id, + "request_digest": self.request.digest, + "state": "COMPLETE", + "product": "forge", + "version": self.request.version, + "product_source": self.request.product_source, + "wheel_sha256": self.request.wheel_sha256, + "controller_source": self.request.controller_source, + "controller_sha256": self.request.controller_sha256, + "runtime_id": self.request.runtime_id, + "installation_id": self.request.installation_id, + "data_root": self.request.data_root, + "backup": state["backup"], + "migration_qualification": state["migration_qualification"], + "live_migration": state["live_migration"], + "installed_readback": state["installed_readback"], + "credential_disposition": "PRESERVED_UNCHANGED", + "service_disposition": "NOT_STARTED", + "mission_disposition": "NOT_STARTED_OR_RESUMED", + "reset_disposition": "NOT_EXECUTED", + "completed_at": _now(), + } + _atomic_json(self.receipt_path, receipt) + self._advance(state, "COMPLETE", receipt_sha256=file_digest(self.receipt_path), + safety_disposition="CANDIDATE_ACTIVE") + self._restore_database_writable() + return receipt + except Exception as error: + try: + self._secure_failure(state, error) + except Exception: + pass + raise + + +def _request_from_args(args: argparse.Namespace) -> UpdateRequest: + return UpdateRequest( + operation_id=args.operation_id, version=args.version, product_source=args.product_source, + wheel=str(args.wheel), wheel_sha256=args.wheel_sha256, + qualification_receipt=str(args.qualification_receipt), + qualification_receipt_sha256=args.qualification_receipt_sha256, + controller_source=args.controller_source, controller_sha256=args.controller_sha256, + data_root=str(args.data_root), runtime_root=str(args.runtime_root), + runtime_id=args.runtime_id, installation_id=args.installation_id, + peer_configuration_digest=args.peer_configuration_digest, + resolver=str(args.resolver), resolver_sha256=args.resolver_sha256, + existing_interpreter=str(args.existing_interpreter), existing_version=args.existing_version, + base_python=str(args.base_python), + ) + + +def main(argv: Sequence[str] | None = None) -> int: + parser = argparse.ArgumentParser(description="Safely update one exact installed Forge runtime") + parser.add_argument("--operation-id", required=True) + parser.add_argument("--version", required=True) + parser.add_argument("--product-source", required=True) + parser.add_argument("--wheel", required=True, type=Path) + parser.add_argument("--wheel-sha256", required=True) + parser.add_argument("--qualification-receipt", required=True, type=Path) + parser.add_argument("--qualification-receipt-sha256", required=True) + parser.add_argument("--controller-source", required=True) + parser.add_argument("--controller-sha256", required=True) + parser.add_argument("--data-root", required=True, type=Path) + parser.add_argument("--runtime-root", required=True, type=Path) + parser.add_argument("--runtime-id", required=True) + parser.add_argument("--installation-id", required=True) + parser.add_argument("--peer-configuration-digest", required=True) + parser.add_argument("--resolver", required=True, type=Path) + parser.add_argument("--resolver-sha256", required=True) + parser.add_argument("--existing-interpreter", required=True, type=Path) + parser.add_argument("--existing-version", required=True) + parser.add_argument("--base-python", required=True, type=Path) + args = parser.parse_args(argv) + try: + receipt = InstalledForgeUpdateController(_request_from_args(args)).run() + except (InstalledForgeUpdateError, OSError, sqlite3.Error, ValueError) as error: + print(json.dumps({"status": "ERROR", "error": str(error)}, sort_keys=True)) + return 1 + print(json.dumps(receipt, sort_keys=True)) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/tests/test_installed_forge_update.py b/tests/test_installed_forge_update.py new file mode 100644 index 0000000..6eaaa70 --- /dev/null +++ b/tests/test_installed_forge_update.py @@ -0,0 +1,551 @@ +"""Qualification for the bounded product-owned installed Forge update.""" + +from __future__ import annotations + +import base64 +import importlib.util +from hashlib import sha256 +import json +import os +from pathlib import Path +import sqlite3 +import sys +import tempfile +from types import SimpleNamespace +import unittest +from unittest.mock import patch +import zipfile + +from forge.runtime import RuntimeBootstrap +from forge.runtime.operational_reset import MAINTENANCE_TABLES + + +SCRIPT = Path(__file__).parents[1] / "scripts" / "update_installed_forge.py" +SPEC = importlib.util.spec_from_file_location("forge_installed_update", SCRIPT) +assert SPEC and SPEC.loader +update = importlib.util.module_from_spec(SPEC) +sys.modules[SPEC.name] = update +SPEC.loader.exec_module(update) + + +class InstalledForgeUpdateTests(unittest.TestCase): + runtime_id = "forge-runtime-735b0321-c4bf-41cd-81d3-9ee00249254b" + installation_id = "99ede979-e8b8-48ca-9174-3257778c680f" + peer_digest = "sha256:" + "c" * 64 + + def setUp(self) -> None: + self.temporary = tempfile.TemporaryDirectory() + self.addCleanup(self.temporary.cleanup) + self.root = Path(self.temporary.name).resolve() + self.data_root = self.root / "Forge Server" + self.runtime_root = self.root / "Forge Server Runtime" + self.runtime_root.mkdir() + self.resolver = self.root / "bin" / "forge" + self.resolver.parent.mkdir() + self.resolver.write_text("#!/bin/sh\nexit 0\n", encoding="utf-8") + self.resolver.chmod(0o755) + self.wheel = self.root / "forge_autonomy-2.7.22-py3-none-any.whl" + metadata = b"Metadata-Version: 2.4\nName: forge-autonomy\nVersion: 2.7.22\n\n" + wheel_metadata = b"Wheel-Version: 1.0\nRoot-Is-Purelib: true\nTag: py3-none-any\n\n" + members = { + "forge/__init__.py": b"__version__ = '2.7.22'\n", + "forge_autonomy-2.7.22.dist-info/METADATA": metadata, + "forge_autonomy-2.7.22.dist-info/WHEEL": wheel_metadata, + } + record_name = "forge_autonomy-2.7.22.dist-info/RECORD" + record = "".join( + f"{name},sha256={base64.urlsafe_b64encode(sha256(payload).digest()).rstrip(b'=').decode()},{len(payload)}\n" + for name, payload in members.items() + ) + f"{record_name},,\n" + with zipfile.ZipFile(self.wheel, "w") as archive: + for name, payload in members.items(): + archive.writestr(name, payload) + archive.writestr(record_name, record) + self.wheel_digest = update.file_digest(self.wheel) + self.sdist_digest = "sha256:" + "d" * 64 + self.receipt = self.root / "release-complete.json" + self.receipt.write_text(json.dumps({ + "state": "RELEASE_COMPLETE", "product": "forge", "component": "forge-autonomy", + "version": "2.7.22", "source_revision": "a" * 40, + "operation_id": "forge-release-2.7.22-" + "a" * 40, + "policy_revision": "forge-bootstrap-release-cadence-v2", + "artifacts": {"wheel": self.wheel_digest, "sdist": self.sdist_digest}, + "qualification": { + "exact_main_sha": "a" * 40, + "qualification": "forge-production-distribution", + "artifact_digests": { + "dist/forge_autonomy-2.7.22-py3-none-any.whl": self.wheel_digest, + "dist/forge_autonomy-2.7.22.tar.gz": self.sdist_digest, + }, + }, + "publication_receipt": { + "product_source_revision": "a" * 40, "readback": "PASS", "registry": "pypi", + "original_release_run_conclusion": "failure", "original_release_run_id": "100", + "reconciliation_contract": "forge-existing-release-reconciliation/v1", + "reconciliation_run_id": "101", "release_controller_source": "b" * 40, + "observed_artifact_digests": { + "forge_autonomy-2.7.22-py3-none-any.whl": self.wheel_digest, + "forge_autonomy-2.7.22.tar.gz": self.sdist_digest, + }, + "github_release": { + "api_url": "https://api.github.com/repos/pcvantol/forge/releases/123", + "database_id": 123, "node_id": "release-node", "tag": "forge-v2.7.22", + "target_commitish": "a" * 40, + }, + }, + "cleanup": { + "result": "COMPLETE", "operation_local_cleanup": "COMPLETE", + "original_release_run_id": "100", "reconciliation_run_id": "101", + "reconciliation_contract": "forge-existing-release-reconciliation/v1", + "release_controller_source": "b" * 40, + "github_release": { + "api_url": "https://api.github.com/repos/pcvantol/forge/releases/123", + "database_id": 123, "node_id": "release-node", "draft": False, + "target_commitish": "a" * 40, "tag_commit": "a" * 40, "tag": "forge-v2.7.22", + }, + }, + }, sort_keys=True), encoding="utf-8") + self.request = update.UpdateRequest( + operation_id="forge-update-2722-test-001", version="2.7.22", + product_source="a" * 40, wheel=str(self.wheel), wheel_sha256=self.wheel_digest, + qualification_receipt=str(self.receipt), + qualification_receipt_sha256=update.file_digest(self.receipt), + controller_source="b" * 40, controller_sha256=update.file_digest(SCRIPT), + data_root=str(self.data_root), runtime_root=str(self.runtime_root), + runtime_id=self.runtime_id, installation_id=self.installation_id, + peer_configuration_digest=self.peer_digest, + resolver=str(self.resolver), resolver_sha256=update.file_digest(self.resolver), + existing_interpreter=sys.executable, existing_version="2.7.21", + base_python=sys.executable, + ) + + def _installed_schema37(self) -> None: + database = RuntimeBootstrap(data_root=self.data_root, forge_version="2.7.22").open() + connection = database._connection + self.runtime_id = database.runtime_identity.runtime_id + self.request = update.UpdateRequest(**{ + **self.request.__dict__, "runtime_id": self.runtime_id, + }) + connection.execute( + "INSERT INTO runtime_metadata(key,value) VALUES ('installation_id',?)", + (self.installation_id,), + ) + peer_document = json.dumps({"credential_reference": "keychain://forge.ep/consumer"}, sort_keys=True) + connection.execute( + "INSERT INTO execution_host_peer_configuration VALUES (1,'forge-ep-primary',5,?,?)", + (self.peer_digest, peer_document), + ) + connection.execute( + "INSERT INTO mission_state VALUES (?,?,?,?,?,?,?,?,?)", + ("MISSION-0001", "COMPLETED", "COMPLETED", None, None, "{}", "{}", "{}", + json.dumps({ + "mission_id": "MISSION-0001", "status": "COMPLETED", + "budgets": {"attempts": 1, "tokens": 1000}, + })), + ) + connection.execute( + "INSERT INTO mission_id_allocations VALUES (?,?,?)", + ("MISSION-0001", "2026-09-17T00:00:00Z", "test-allocation"), + ) + connection.execute( + "INSERT INTO architecture_reviews VALUES (?,?,?,?,?,?,?,?,?)", + ("review-1", "MISSION-0001", "repository://forge", "[]", "low", "low", "high", + "2026-09-17T00:00:00Z", json.dumps({"id": "review-1"})), + ) + connection.execute( + "INSERT INTO execution_receipts VALUES (?,?,?,?,?,?,?,?)", + ("receipt-1", "MISSION-0001", "engineering-platform", "run-1", "report-1", + "correlation-1", "2026-09-17T00:00:00Z", "complete"), + ) + database._insert_governance_grant( + "grant-1", self.installation_id, "operator-1", "MISSION_INTAKE", + "test-provenance", "sha256:grant-1", "2026-09-17T00:00:00Z", + ) + database._insert_governance_authority( + self.installation_id, "operator-1", "MISSION_INTAKE", "2026-09-17T00:00:00Z", + ) + connection.execute( + "INSERT OR REPLACE INTO dispatcher_state VALUES (1,'IDLE',NULL,'[]',?)", + (json.dumps({"status": "IDLE", "active_mission_id": None, "mission_sequence": []}),), + ) + connection.commit() + database.close() + (self.data_root / "instance" / "runtime-instance.json").write_text( + self.runtime_id + "\n", encoding="utf-8" + ) + + connection = sqlite3.connect(self.data_root / "forge.db") + for (trigger,) in connection.execute( + "SELECT name FROM sqlite_master WHERE type='trigger' AND name LIKE 'operational_reset_%'" + ).fetchall(): + connection.execute(f'DROP TRIGGER "{trigger}"') + for table in MAINTENANCE_TABLES: + for (trigger,) in connection.execute( + "SELECT name FROM sqlite_master WHERE type='trigger' AND tbl_name=?", (table,) + ).fetchall(): + connection.execute(f'DROP TRIGGER "{trigger}"') + for table in reversed(MAINTENANCE_TABLES): + connection.execute(f'DROP TABLE "{table}"') + connection.execute( + "UPDATE runtime_metadata SET value='37' " + "WHERE key IN ('schema_version','migration_version','last_migration','database_version')" + ) + connection.execute("PRAGMA user_version=37") + connection.commit() + connection.close() + + def _controller(self) -> object: + return update.InstalledForgeUpdateController( + self.request, process_reader=lambda: (), + ) + + def _qualified_schema38_copy(self, controller: object, before: dict[str, object]) -> None: + copy_root = controller.operation_root / "qualification-copy" + (copy_root / "instance").mkdir(parents=True) + (copy_root / "instance" / "runtime-instance.json").write_text( + self.runtime_id + "\n", encoding="utf-8" + ) + update._copy_sqlite_backup(self.data_root / "forge.db", copy_root / "forge.db") + migrated = RuntimeBootstrap(data_root=copy_root, forge_version="2.7.22").open() + migrated.close() + update.verify_preservation( + before, update.database_snapshot(copy_root / "forge.db"), self.request, + ) + + def test_exact_release_complete_artifact_is_accepted_and_mismatch_rejected(self) -> None: + evidence = update.validate_qualified_artifact(self.request) + self.assertEqual(evidence["wheel_sha256"], self.wheel_digest) + changed = update.UpdateRequest(**{**self.request.__dict__, "product_source": "f" * 40}) + with self.assertRaisesRegex(update.InstalledForgeUpdateError, "noncanonical"): + update.validate_qualified_artifact(changed) + self.wheel.write_bytes(b"changed") + with self.assertRaisesRegex(update.InstalledForgeUpdateError, "wheel digest"): + update.validate_qualified_artifact(self.request) + + def test_real_schema37_to_38_migration_preserves_history_and_bindings(self) -> None: + self._installed_schema37() + before = update.database_snapshot(self.data_root / "forge.db") + update.assert_selected_installation(self.request, before) + update.assert_quiescent(before) + backup = self.root / "backup" / "forge.sqlite3" + backup_evidence = update._copy_sqlite_backup(self.data_root / "forge.db", backup) + self.assertEqual(backup_evidence["integrity_check"], "ok") + + copy_root = self.root / "qualification-copy" + (copy_root / "instance").mkdir(parents=True) + (copy_root / "instance" / "runtime-instance.json").write_text( + self.runtime_id + "\n", encoding="utf-8" + ) + copied = copy_root / "forge.db" + copied.write_bytes(backup.read_bytes()) + migrated = RuntimeBootstrap(data_root=copy_root, forge_version="2.7.22").open() + migrated.close() + after = update.database_snapshot(copied) + preservation = update.verify_preservation(before, after, self.request) + self.assertEqual(preservation["status"], "PASS") + self.assertEqual(after["tables"]["mission_state"], before["tables"]["mission_state"]) + self.assertEqual(after["tables"]["mission_id_allocations"], before["tables"]["mission_id_allocations"]) + self.assertEqual(after["tables"]["architecture_reviews"], before["tables"]["architecture_reviews"]) + self.assertEqual(after["tables"]["execution_receipts"], before["tables"]["execution_receipts"]) + self.assertEqual(after["tables"]["governance_capability_grants"], before["tables"]["governance_capability_grants"]) + self.assertEqual(after["peer"], before["peer"]) + + def test_changed_target_and_active_writer_are_blocked(self) -> None: + self._installed_schema37() + snapshot = update.database_snapshot(self.data_root / "forge.db") + wrong = update.UpdateRequest(**{**self.request.__dict__, "runtime_id": "forge-runtime-wrong"}) + with self.assertRaisesRegex(update.InstalledForgeUpdateError, "different runtime"): + update.assert_selected_installation(wrong, snapshot) + connection = sqlite3.connect(self.data_root / "forge.db") + connection.execute( + "UPDATE dispatcher_state SET status='ACTIVE',active_mission_id='MISSION-0001' WHERE singleton=1" + ) + connection.commit() + connection.close() + with self.assertRaisesRegex(update.InstalledForgeUpdateError, "dispatcher"): + update.assert_quiescent(update.database_snapshot(self.data_root / "forge.db")) + + def test_concurrent_update_lock_is_rejected(self) -> None: + lock = self.runtime_root / "locks" / "installation-update.lock" + with update.exclusive_lock(lock): + with self.assertRaisesRegex(update.InstalledForgeUpdateError, "concurrent maintenance"): + with update.exclusive_lock(lock): + pass + + def test_global_update_lock_precedes_operation_state_and_staging(self) -> None: + self._installed_schema37() + controller = self._controller() + lock = self.runtime_root / "locks" / "installation-update.lock" + with update.exclusive_lock(lock): + with self.assertRaisesRegex(update.InstalledForgeUpdateError, "concurrent maintenance"): + controller.run() + self.assertFalse(controller.state_path.exists()) + self.assertFalse(controller.slot.exists()) + + def test_candidate_venv_is_created_at_its_final_non_relocated_slot(self) -> None: + controller = self._controller() + state = controller._state() + + def identity(_interpreter: Path, *, cwd: Path) -> dict[str, str]: + del cwd + root = controller.slot.resolve() + return { + "version": "2.7.22", "distribution_version": "2.7.22", + "module": str(root / "lib" / "python" / "site-packages" / "forge" / "__init__.py"), + "prefix": str(root), "sys_executable": str(root / "bin" / "python"), + } + + with ( + patch.object(update, "installed_identity", side_effect=identity), + patch.object(update, "_install_validated_wheel"), + patch.object(update, "_verify_candidate_files", return_value={ + "wheel_manifest_digest": update._digest_bytes(update._json_bytes({})), + "installed_file_count": 0, "entrypoint_sha256": "sha256:" + "0" * 64, + }), + patch.object(update, "_run", return_value=SimpleNamespace(stdout="", stderr="")), + ): + staged = controller._stage(state) + self.assertEqual(staged["phase"], "STAGED") + self.assertTrue(controller.slot_receipt.is_file()) + self.assertTrue((controller.slot / "forge-installation-staging.json").is_file()) + self.assertFalse((controller.slot.parent / f".stage-{self.request.operation_id}").exists()) + + def test_process_scan_does_not_hide_a_sibling_with_the_same_parent(self) -> None: + controller = self._controller() + output = ( + f"{os.getpid()} {os.getppid()} self\n" + f"99999 {os.getppid()} selected-forge-runtime\n" + ) + with patch.object(update, "_run", return_value=SimpleNamespace(stdout=output, stderr="")): + self.assertEqual(controller._processes(), ["selected-forge-runtime"]) + + def test_atomic_backup_is_adopted_after_state_write_interruption(self) -> None: + self._installed_schema37() + controller = self._controller() + state = controller._state() + before = update.database_snapshot(self.data_root / "forge.db") + update._copy_sqlite_backup(self.data_root / "forge.db", controller.backup_path) + recovered = controller._backup(state, before) + self.assertTrue(recovered["backup"]["recovered_after_atomic_backup_write"]) + self.assertEqual( + update.database_snapshot(controller.backup_path)["content_digest"], + before["content_digest"], + ) + + def test_complete_fast_path_revalidates_durable_receipt_bytes(self) -> None: + controller = self._controller() + state = controller._state() + update._atomic_json(controller.receipt_path, { + "request_digest": self.request.digest, "state": "COMPLETE", + }) + controller._advance( + state, "COMPLETE", receipt_sha256=update.file_digest(controller.receipt_path), + ) + update._atomic_json(controller.receipt_path, { + "request_digest": self.request.digest, "state": "COMPLETE", "tampered": True, + }) + with self.assertRaisesRegex(update.InstalledForgeUpdateError, "receipt changed"): + controller.run() + + def test_resolver_adoption_fences_without_editing_the_legacy_environment(self) -> None: + controller = self._controller() + state = controller._state() + with patch.object(update, "installed_identity", return_value={"version": "2.7.21"}): + state = controller._adopt_resolver(state) + self.assertTrue(self.resolver.is_symlink()) + self.assertEqual(self.resolver.resolve(), controller.legacy_entrypoint.resolve()) + legacy_bytes = controller.legacy_entrypoint.read_bytes() + state = controller._fence(state) + self.assertEqual(self.resolver.resolve(), controller.fenced_resolver.resolve()) + self.assertEqual(controller.legacy_entrypoint.read_bytes(), legacy_bytes) + + def test_crash_before_migration_restores_the_legacy_route(self) -> None: + controller = self._controller() + state = controller._state() + with patch.object(update, "installed_identity", return_value={"version": "2.7.21"}): + state = controller._adopt_resolver(state) + state = controller._advance(state, "BACKED_UP", before={"content_digest": "sha256:before"}) + state = controller._fence(state) + with patch.object(update, "database_snapshot", return_value={ + "user_version": 37, "content_digest": "sha256:before", + }): + controller._secure_failure(state, RuntimeError("interrupted")) + self.assertEqual(self.resolver.resolve(), controller.legacy_entrypoint.resolve()) + durable = update._read_json(controller.state_path) + self.assertEqual(durable["safety_disposition"], "LEGACY_RESTORED_BEFORE_MIGRATION") + + def test_crash_after_migration_never_reactivates_the_old_binary(self) -> None: + controller = self._controller() + state = controller._state() + with patch.object(update, "installed_identity", return_value={"version": "2.7.21"}): + state = controller._adopt_resolver(state) + controller._fence(state) + with patch.object(update, "database_snapshot", return_value={"user_version": 38}): + controller._secure_failure(state, RuntimeError("interrupted")) + self.assertEqual(self.resolver.resolve(), controller.fenced_resolver.resolve()) + self.assertNotEqual(self.resolver.resolve(), controller.legacy_entrypoint.resolve()) + + def test_crash_during_atomic_activation_leaves_the_candidate_active(self) -> None: + controller = self._controller() + state = controller._state() + with patch.object(update, "installed_identity", return_value={"version": "2.7.21"}): + state = controller._adopt_resolver(state) + candidate = controller.slot / "bin" / "forge" + candidate.parent.mkdir(parents=True) + candidate.write_text("#!/bin/sh\nexit 0\n", encoding="utf-8") + candidate.chmod(0o755) + update._replace_symlink( + controller.current, __import__("os").path.relpath(candidate, controller.runtime_root) + ) + with patch.object(update, "database_snapshot", return_value={"user_version": 38}): + controller._secure_failure(state, RuntimeError("interrupted")) + self.assertEqual(self.resolver.resolve(), controller.fenced_resolver.resolve()) + + def test_preservation_rejects_domain_loss_and_unexpected_schema_objects(self) -> None: + self._installed_schema37() + before = update.database_snapshot(self.data_root / "forge.db") + after = json.loads(json.dumps(before)) + after["user_version"] = 38 + after["tables"]["mission_state"]["count"] = 0 + for table in update.NEW_SCHEMA_38_TABLES: + after["tables"][table] = { + "count": 1 if table == "operational_reset_state" else 0, + "digest": "sha256:" + "0" * 64, + } + after["writer_state"]["operational_reset"] = [ + {"dataset_generation": 0, "active_operation_id": None, "state": "IDLE"} + ] + with self.assertRaisesRegex(update.InstalledForgeUpdateError, "historical table"): + update.verify_preservation(before, after, self.request) + + def test_noncanonical_receipt_and_unsafe_operation_ids_are_rejected(self) -> None: + receipt = json.loads(self.receipt.read_text(encoding="utf-8")) + receipt["publication_receipt"]["unexpected"] = "not-allowed" + self.receipt.write_text(json.dumps(receipt, sort_keys=True), encoding="utf-8") + changed = update.UpdateRequest(**{ + **self.request.__dict__, + "qualification_receipt_sha256": update.file_digest(self.receipt), + }) + with self.assertRaisesRegex(update.InstalledForgeUpdateError, "noncanonical"): + update.validate_qualified_artifact(changed) + for operation_id in (".", "..", "unsafe/name"): + with self.assertRaises(update.InstalledForgeUpdateError): + update.UpdateRequest(**{**self.request.__dict__, "operation_id": operation_id}).validate() + + def test_quiescence_uses_explicit_safe_state_allowlists(self) -> None: + self._installed_schema37() + connection = sqlite3.connect(self.data_root / "forge.db") + connection.execute( + "INSERT INTO planning_provider_generation_permits VALUES (?,?,?,?,?,?,?,?)", + ("permit-1", "provider", 1, "sha256:" + "1" * 64, "sha256:" + "2" * 64, + "INVALIDATED", "now", "now"), + ) + connection.commit() + connection.close() + update.assert_quiescent(update.database_snapshot(self.data_root / "forge.db")) + connection = sqlite3.connect(self.data_root / "forge.db") + connection.execute( + "UPDATE planning_provider_generation_permits SET state='PENDING' WHERE permit_id='permit-1'" + ) + connection.execute("UPDATE mission_state SET status='CREATED' WHERE mission_id='MISSION-0001'") + connection.commit() + connection.close() + with self.assertRaises(update.InstalledForgeUpdateError): + update.assert_quiescent(update.database_snapshot(self.data_root / "forge.db")) + + def test_completed_replay_requires_exact_activated_schema38_fingerprint(self) -> None: + snapshot = { + "user_version": 38, + "tables": {name: {} for name in update.NEW_SCHEMA_38_TABLES}, + "schema_digest": "sha256:" + "1" * 64, + } + readback = {"database_schema_digest": snapshot["schema_digest"]} + update.assert_completed_schema(snapshot, readback) + with self.assertRaisesRegex(update.InstalledForgeUpdateError, "schema changed"): + update.assert_completed_schema({**snapshot, "user_version": 37}, readback) + with self.assertRaisesRegex(update.InstalledForgeUpdateError, "schema changed"): + update.assert_completed_schema( + {**snapshot, "schema_digest": "sha256:" + "2" * 64}, readback, + ) + + def test_atomic_product_migrated_copy_rejects_late_live_mutation(self) -> None: + self._installed_schema37() + controller = self._controller() + before = update.database_snapshot(self.data_root / "forge.db") + self._qualified_schema38_copy(controller, before) + connection = sqlite3.connect(self.data_root / "forge.db") + connection.execute("INSERT INTO mission_id_allocations VALUES (?,?,?)", ("late", "now", "late")) + connection.commit() + connection.close() + with self.assertRaisesRegex(update.InstalledForgeUpdateError, "changed after"): + controller._install_qualified_database(before) + + def test_database_swap_is_crash_safe_before_atomic_replace(self) -> None: + self._installed_schema37() + controller = update.InstalledForgeUpdateController( + self.request, process_reader=lambda: (), interrupt_after="database_swap_prepared", + ) + before = update.database_snapshot(self.data_root / "forge.db") + self._qualified_schema38_copy(controller, before) + with self.assertRaisesRegex(update.InstalledForgeUpdateError, "database_swap_prepared"): + controller._install_qualified_database(before) + self.assertEqual(update.database_snapshot(self.data_root / "forge.db")["user_version"], 37) + self.assertFalse((self.data_root / "forge.db-wal").exists()) + self.assertFalse((self.data_root / "forge.db-shm").exists()) + + def test_database_swap_is_crash_safe_after_atomic_replace(self) -> None: + self._installed_schema37() + controller = update.InstalledForgeUpdateController( + self.request, process_reader=lambda: (), interrupt_after="database_swap", + ) + before = update.database_snapshot(self.data_root / "forge.db") + self._qualified_schema38_copy(controller, before) + with self.assertRaisesRegex(update.InstalledForgeUpdateError, "database_swap"): + controller._install_qualified_database(before) + after = update.database_snapshot(self.data_root / "forge.db") + self.assertEqual(after["user_version"], 38) + update.verify_preservation(before, after, self.request) + self.assertEqual((self.data_root / "forge.db").stat().st_mode & 0o777, 0o400) + self.assertFalse((self.data_root / "forge.db-wal").exists()) + self.assertFalse((self.data_root / "forge.db-shm").exists()) + + def test_exact_published_wheel_end_to_end_when_requested(self) -> None: + wheel_value = os.environ.get("FORGE_EXACT_WHEEL") + receipt_value = os.environ.get("FORGE_EXACT_RELEASE_RECEIPT") + if not wheel_value or not receipt_value: + self.skipTest("exact published wheel paths were not supplied") + self._installed_schema37() + wheel = Path(wheel_value).resolve() + receipt = Path(receipt_value).resolve() + release = json.loads(receipt.read_text(encoding="utf-8")) + legacy_interpreter = Path(os.environ.get("FORGE_LEGACY_INTERPRETER", sys.executable)) + try: + legacy_identity = update.installed_identity(legacy_interpreter, cwd=self.root) + except update.InstalledForgeUpdateError as error: + self.skipTest(f"legacy Forge interpreter was not supplied: {error}") + request = update.UpdateRequest(**{ + **self.request.__dict__, + "product_source": release["source_revision"], + "wheel": str(wheel), "wheel_sha256": update.file_digest(wheel), + "qualification_receipt": str(receipt), + "qualification_receipt_sha256": update.file_digest(receipt), + "controller_sha256": update.file_digest(SCRIPT), + "existing_interpreter": str(legacy_interpreter), + "existing_version": legacy_identity["version"], + }) + controller = update.InstalledForgeUpdateController(request, process_reader=lambda: ()) + completed = controller.run() + self.assertEqual(completed["state"], "COMPLETE") + self.assertEqual(update.database_snapshot(self.data_root / "forge.db")["user_version"], 38) + self.assertEqual(controller.run(), completed) + self.assertEqual((self.data_root / "forge.db").stat().st_mode & 0o777, 0o600) + connection = sqlite3.connect(self.data_root / "forge.db") + connection.execute("PRAGMA user_version=37") + connection.commit() + connection.close() + with self.assertRaisesRegex(update.InstalledForgeUpdateError, "schema changed"): + controller.run() + + +if __name__ == "__main__": + unittest.main()