From 31e91294872d381a44df38231b2e906a5076b72c Mon Sep 17 00:00:00 2001 From: pcvantol Date: Fri, 18 Sep 2026 00:03:38 +0200 Subject: [PATCH 1/3] Add bounded installed Forge update controller --- README.md | 10 + .../FORGE_INSTALLED_UPDATE_RUNBOOK.md | 103 ++ scripts/update_installed_forge.py | 992 ++++++++++++++++++ tests/test_installed_forge_update.py | 283 +++++ 4 files changed, 1388 insertions(+) create mode 100644 docs/operations/FORGE_INSTALLED_UPDATE_RUNBOOK.md create mode 100644 scripts/update_installed_forge.py create mode 100644 tests/test_installed_forge_update.py 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..2254dd8 --- /dev/null +++ b/docs/operations/FORGE_INSTALLED_UPDATE_RUNBOOK.md @@ -0,0 +1,103 @@ +# 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. Validate the terminal release receipt and wheel without importing the + wheel. +2. Create an isolated versioned runtime slot outside the source checkout with + the explicit Python interpreter and install only the local wheel using + `--no-index --no-deps`. +3. Read the candidate's version, distribution, module, prefix, and interpreter + from an isolated process. +4. Acquire the installation-operation lock and the canonical + `forge-runtime-mutation.lock`; reject an active runtime process, dispatcher, + Mission, scheduler submission, 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. Atomically point the stable resolver at a maintenance fence, apply the same + Forge-owned migration to the live selected data root, and repeat the full + preservation check. +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 live migration, failure restores the retained 2.7.21 command route. + A hard interruption may leave the explicit maintenance fence; resuming the + same operation reconciles it from durable evidence. +- From the first schema-38 readback onward, the old binary is never selected. + The resolver remains fenced until the 2.7.22 candidate is verified, or it + already resolves to that candidate after an atomic activation interruption. +- A completed receipt is idempotently returned only when its request digest + matches exactly. + +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, +and interruption before migration, after migration, and during activation. + +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..d3e970a --- /dev/null +++ b/scripts/update_installed_forge.py @@ -0,0 +1,992 @@ +#!/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 +from contextlib import ExitStack, contextmanager +from dataclasses import asdict, dataclass +from datetime import datetime, timezone +from email.parser import BytesParser +from hashlib import sha256 +import json +import os +from pathlib import Path +import shutil +import sqlite3 +import subprocess +import sys +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 +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", +}) +ACTIVE_MISSION_STATES = frozenset({ + "READY", "ACTIVE", "WAITING_FOR_EXECUTION", "WAITING_FOR_EVIDENCE", + "READY_TO_CONTINUE", "INTEGRATION_RUNNING", +}) +ACTIVE_SUBMISSION_STATES = frozenset({ + "CREATED", "SUBMITTED", "ACCEPTED", "EXECUTING", "RECEIPT_AVAILABLE", +}) +TERMINAL_PERMIT_STATES = frozenset({"CONSUMED", "CANCELLED", "EXPIRED", "FAILED", "REVOKED"}) + + +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 file_digest(path: Path) -> str: + if path.is_symlink() or not path.is_file(): + raise InstalledForgeUpdateError(f"required regular file is unavailable: {path}") + digest = sha256() + with path.open("rb") as handle: + for chunk in iter(lambda: handle.read(1024 * 1024), b""): + digest.update(chunk) + 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}") + if create: + path.mkdir(parents=True, exist_ok=True) + if path.is_symlink() or 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) + temporary = path.with_name(f".{path.name}.tmp-{os.getpid()}") + if temporary.exists() or temporary.is_symlink(): + temporary.unlink() + descriptor = os.open(temporary, os.O_WRONLY | os.O_CREAT | os.O_EXCL, 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]: + if path.is_symlink() or not path.is_file(): + raise InstalledForgeUpdateError(f"durable JSON evidence is unavailable or unsafe: {path}") + try: + value = json.loads(path.read_text(encoding="utf-8")) + except (OSError, 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-{os.getpid()}") + if temporary.exists() or temporary.is_symlink(): + temporary.unlink() + temporary.symlink_to(target) + os.replace(temporary, path) + + +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", + } + 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": str(pathlib.Path(sys.executable).resolve()), + "prefix": str(pathlib.Path(sys.prefix).resolve()), +}, sort_keys=True)) +""" + result = _run((str(interpreter), "-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(not value or any(character not in "abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789._-" for character in value) + 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 validate_qualified_artifact(request: UpdateRequest) -> dict[str, Any]: + wheel = Path(request.wheel) + if file_digest(wheel) != request.wheel_sha256: + raise InstalledForgeUpdateError("wheel digest does not match the selected qualified artifact") + 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") + try: + with zipfile.ZipFile(wheel) as archive: + names = archive.namelist() + if any(name.startswith("/") or ".." in Path(name).parts for name in names): + raise InstalledForgeUpdateError("wheel contains an unsafe member path") + metadata_names = [name for name in names if name.endswith(".dist-info/METADATA")] + if len(metadata_names) != 1: + raise InstalledForgeUpdateError("wheel metadata is missing or ambiguous") + metadata = BytesParser().parsebytes(archive.read(metadata_names[0])) + except zipfile.BadZipFile as error: + raise InstalledForgeUpdateError("wheel is not a valid ZIP artifact") from error + if metadata.get("Name") != "forge-autonomy" or metadata.get("Version") != request.version: + raise InstalledForgeUpdateError("wheel package metadata does not match Forge 2.7.22") + + receipt_path = Path(request.qualification_receipt) + if file_digest(receipt_path) != request.qualification_receipt_sha256: + raise InstalledForgeUpdateError("qualification receipt digest changed") + receipt = _read_json(receipt_path) + qualification = receipt.get("qualification") + artifacts = receipt.get("artifacts") + if ( + 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 not isinstance(qualification, dict) + or qualification.get("exact_main_sha") != request.product_source + or qualification.get("qualification") != "forge-production-distribution" + or not isinstance(artifacts, dict) + or artifacts.get("wheel") != request.wheel_sha256 + ): + raise InstalledForgeUpdateError("release-complete qualification does not bind the exact wheel and source") + return { + "wheel": str(wheel), "wheel_sha256": request.wheel_sha256, + "receipt": str(receipt_path), "receipt_sha256": request.qualification_receipt_sha256, + "release_operation_id": receipt.get("operation_id"), + } + + +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) -> dict[str, Any]: + if path.is_symlink() or not path.is_file(): + raise InstalledForgeUpdateError(f"runtime database is unavailable or unsafe: {path}") + try: + connection = sqlite3.connect(path.resolve().as_uri() + "?mode=ro", uri=True) + connection.row_factory = sqlite3.Row + 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_%'" + )) + 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 "connection" in locals(): + connection.close() + snapshot = { + "database": str(path), "integrity_check": integrity, + "foreign_key_check": foreign_keys, "user_version": user_version, + "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 [] + active_missions = [row for row in missions if isinstance(row, Mapping) and row.get("status") in ACTIVE_MISSION_STATES] + if active_missions: + raise InstalledForgeUpdateError("Forge has active or automatically resumable Mission state") + submissions = writer.get("submissions") or [] + if any(isinstance(row, Mapping) and row.get("state") in ACTIVE_SUBMISSION_STATES for row in submissions): + raise InstalledForgeUpdateError("Forge has an active 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 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") + temporary = destination.with_name(f".{destination.name}.tmp-{os.getpid()}") + if temporary.exists(): + temporary.unlink() + 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 + + +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.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]: + updated = {**state, **evidence, "phase": phase, "updated_at": _now()} + history = list(state.get("history", [])) + history.append({"phase": 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 or parent 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 = validate_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: + raise InstalledForgeUpdateError("immutable candidate slot belongs to a different request") + elif self.slot.exists() and not self.slot.is_symlink(): + identity = installed_identity(self.slot / "bin" / "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) + ): + raise InstalledForgeUpdateError("unreceipted candidate slot cannot be safely adopted") + _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, "staged_at": _now(), + "recovered_after_atomic_slot_move": True, + }) + else: + stage = self.runtime_root / "slots" / f".stage-{self.request.operation_id}" + _safe_directory(stage.parent, create=True) + if stage.is_symlink(): + raise InstalledForgeUpdateError("candidate staging path is unsafe") + _run((self.request.base_python, "-m", "venv", str(stage)), cwd=self.runtime_root) + candidate_python = stage / "bin" / "python" + _run((str(candidate_python), "-m", "pip", "install", "--no-index", "--no-deps", self.request.wheel), + cwd=self.runtime_root) + 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(stage.resolve()) + os.sep) + or Path(str(identity.get("prefix"))).resolve() != stage.resolve() + ): + raise InstalledForgeUpdateError("candidate slot identity is inconsistent") + if self.slot.exists() or self.slot.is_symlink(): + raise InstalledForgeUpdateError("candidate slot appeared concurrently") + os.replace(stage, self.slot) + 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) + os.sep): + raise InstalledForgeUpdateError("staged candidate readback changed") + if not self.slot_receipt.exists(): + _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, "staged_at": _now(), + }) + return self._advance(state, "STAGED", artifact_qualification=qualification, candidate=identity, + candidate_slot=str(self.slot)) + + 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(): + if resolver.is_symlink() or not resolver.is_file() or file_digest(resolver) != 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") + shutil.copy2(resolver, self.legacy_entrypoint) + os.chmod(self.legacy_entrypoint, resolver.stat().st_mode & 0o777) + 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) + 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": + return state + root = self.operation_root / "qualification-copy" + _safe_directory(root, create=True) + database = root / "forge.db" + temporary = root / f".forge.db.tmp-{os.getpid()}" + shutil.copy2(self.backup_path, temporary) + os.replace(temporary, database) + for suffix in ("-wal", "-shm"): + sidecar = root / ("forge.db" + suffix) + if sidecar.exists() and not sidecar.is_symlink(): + sidecar.unlink() + instance = _safe_directory(root / "instance", create=True) + marker = instance / "runtime-instance.json" + marker.write_text(self.request.runtime_id + "\n", encoding="utf-8") + os.chmod(marker, 0o600) + 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_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), + "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 _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") + output = _candidate_migrate(self.slot / "bin" / "forge", self.data_root, cwd=self.runtime_root) + after = database_snapshot(self.database) + preservation = verify_preservation(before, after, self.request) + state = self._advance( + state, "MIGRATED", live_migration={ + **preservation, "migrated_at": _now(), "candidate_output": output, + "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) + 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 identity.get("sys_executable") != str((self.slot / "bin" / "python").resolve()) + or not str(identity.get("module", "")).startswith(str(self.slot) + os.sep) + or status.get("product_version") != self.request.version + or status.get("data_root") != self.request.data_root + 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"], + "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.""" + current = database_snapshot(self.database) + if current.get("user_version") == SCHEMA_BEFORE and self.legacy_entrypoint.exists(): + self._restore_legacy_before_migration(state, error) + return + candidate = (self.slot / "bin" / "forge").resolve() + if not self.current.is_symlink() or _resolved_link(self.current) != candidate: + _replace_symlink(self.current, os.path.relpath(self.fenced_resolver, self.runtime_root)) + self._advance( + state, state.get("phase", "RECOVERY_PENDING"), + safety_disposition="SCHEMA_38_OLD_BINARY_FENCED", last_error=str(error), + ) + + def run(self) -> dict[str, Any]: + _safe_directory(self.data_root) + _safe_directory(self.runtime_root) + _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) + if receipt.get("request_digest") != self.request.digest: + raise InstalledForgeUpdateError("completed update receipt conflicts with this request") + return receipt + + state = self._stage(state) + self._interrupt("stage") + update_lock = self.runtime_root / "locks" / "installation-update.lock" + runtime_lock = self.data_root / "forge-runtime-mutation.lock" + with ExitStack() as locks: + locks.enter_context(exclusive_lock(update_lock)) + locks.enter_context(exclusive_lock(runtime_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") + 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..4723c66 --- /dev/null +++ b/tests/test_installed_forge_update.py @@ -0,0 +1,283 @@ +"""Qualification for the bounded product-owned installed Forge update.""" + +from __future__ import annotations + +import importlib.util +from hashlib import sha256 +import json +from pathlib import Path +import sqlite3 +import sys +import tempfile +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) + 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" + with zipfile.ZipFile(self.wheel, "w") as archive: + archive.writestr("forge_autonomy-2.7.22.dist-info/METADATA", metadata) + self.wheel_digest = update.file_digest(self.wheel) + 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": "release-1", + "artifacts": {"wheel": self.wheel_digest, "sdist": "sha256:" + "d" * 64}, + "qualification": { + "exact_main_sha": "a" * 40, + "qualification": "forge-production-distribution", + }, + }, 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 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, "does not bind"): + 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_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._fence(state) + with patch.object(update, "database_snapshot", return_value={"user_version": 37}): + 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(), candidate.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) + + +if __name__ == "__main__": + unittest.main() From e36c393ba1dcaf8f33f4cbae6a32234ce2a30ab0 Mon Sep 17 00:00:00 2001 From: pcvantol Date: Fri, 18 Sep 2026 00:34:54 +0200 Subject: [PATCH 2/3] Harden installed Forge update recovery --- .../FORGE_INSTALLED_UPDATE_RUNBOOK.md | 50 +- scripts/update_installed_forge.py | 888 ++++++++++++++---- tests/test_installed_forge_update.py | 261 ++++- 3 files changed, 982 insertions(+), 217 deletions(-) diff --git a/docs/operations/FORGE_INSTALLED_UPDATE_RUNBOOK.md b/docs/operations/FORGE_INSTALLED_UPDATE_RUNBOOK.md index 2254dd8..44878e0 100644 --- a/docs/operations/FORGE_INSTALLED_UPDATE_RUNBOOK.md +++ b/docs/operations/FORGE_INSTALLED_UPDATE_RUNBOOK.md @@ -35,17 +35,21 @@ interpreter, version, and bytes. ## Safety sequence -1. Validate the terminal release receipt and wheel without importing the - wheel. +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 and install only the local wheel using - `--no-index --no-deps`. + 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 and the canonical - `forge-runtime-mutation.lock`; reject an active runtime process, dispatcher, - Mission, scheduler submission, provider-generation permit, planning queue, - reset, or conflicting operation. +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. @@ -56,9 +60,11 @@ interpreter, version, and bytes. 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. Atomically point the stable resolver at a maintenance fence, apply the same - Forge-owned migration to the live selected data root, and repeat the full - preservation check. +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. @@ -77,14 +83,18 @@ authorized read-only Forge-to-EP check. Re-run the same exact operation ID and arguments. A conflicting request is rejected. -- Before live migration, failure restores the retained 2.7.21 command route. +- 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 onward, the old binary is never selected. - The resolver remains fenced until the 2.7.22 candidate is verified, or it - already resolves to that candidate after an atomic activation interruption. -- A completed receipt is idempotently returned only when its request digest - matches exactly. +- 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 @@ -96,7 +106,11 @@ separate authority and compatibility proof. 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, -and interruption before migration, after migration, and during activation. +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 diff --git a/scripts/update_installed_forge.py b/scripts/update_installed_forge.py index d3e970a..9d7df71 100644 --- a/scripts/update_installed_forge.py +++ b/scripts/update_installed_forge.py @@ -15,18 +15,26 @@ from __future__ import annotations import argparse -from contextlib import ExitStack, contextmanager +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 @@ -39,6 +47,12 @@ 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", @@ -50,14 +64,11 @@ "schema_version", "migration_version", "last_migration", "forge_version", "database_version", "last_access_at", "integrity_status", }) -ACTIVE_MISSION_STATES = frozenset({ - "READY", "ACTIVE", "WAITING_FOR_EXECUTION", "WAITING_FOR_EVIDENCE", - "READY_TO_CONTINUE", "INTEGRATION_RUNNING", +SAFE_MISSION_STATES = frozenset({ + "BLOCKED", "FAILED", "COMPLETED", "ARCHIVED", "INTEGRATION_BLOCKED", "INTEGRATION_COMPLETE", }) -ACTIVE_SUBMISSION_STATES = frozenset({ - "CREATED", "SUBMITTED", "ACCEPTED", "EXECUTING", "RECEIPT_AVAILABLE", -}) -TERMINAL_PERMIT_STATES = frozenset({"CONSUMED", "CANCELLED", "EXPIRED", "FAILED", "REVOKED"}) +SAFE_SUBMISSION_STATES = frozenset({"RECONCILED", "BLOCKED", "FAILED", "SUPERSEDED"}) +TERMINAL_PERMIT_STATES = frozenset({"INVALIDATED", "INVALIDATED_BY_OPERATIONAL_RESET"}) class InstalledForgeUpdateError(RuntimeError): @@ -76,32 +87,79 @@ 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: - if path.is_symlink() or not path.is_file(): - raise InstalledForgeUpdateError(f"required regular file is unavailable: {path}") + _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() - with path.open("rb") as handle: - for chunk in iter(lambda: handle.read(1024 * 1024), b""): + 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) - if path.is_symlink() or not path.is_dir(): + _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) - temporary = path.with_name(f".{path.name}.tmp-{os.getpid()}") - if temporary.exists() or temporary.is_symlink(): - temporary.unlink() - descriptor = os.open(temporary, os.O_WRONLY | os.O_CREAT | os.O_EXCL, 0o600) + 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)) @@ -115,11 +173,9 @@ def _atomic_json(path: Path, value: object) -> None: def _read_json(path: Path) -> dict[str, Any]: - if path.is_symlink() or not path.is_file(): - raise InstalledForgeUpdateError(f"durable JSON evidence is unavailable or unsafe: {path}") try: - value = json.loads(path.read_text(encoding="utf-8")) - except (OSError, json.JSONDecodeError) as error: + 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}") @@ -128,11 +184,13 @@ def _read_json(path: Path) -> dict[str, Any]: def _replace_symlink(path: Path, target: str) -> None: _safe_directory(path.parent, create=True) - temporary = path.with_name(f".{path.name}.link-{os.getpid()}") - if temporary.exists() or temporary.is_symlink(): - temporary.unlink() + temporary = path.with_name(f".{path.name}.link-{secrets.token_hex(16)}") temporary.symlink_to(target) - os.replace(temporary, path) + try: + os.replace(temporary, path) + finally: + if temporary.is_symlink(): + temporary.unlink() def _resolved_link(path: Path) -> Path: @@ -147,6 +205,7 @@ def _environment() -> dict[str, str]: "LANG": "C.UTF-8", "LC_ALL": "C.UTF-8", "PYTHONNOUSERSITE": "1", + "PYTHONDONTWRITEBYTECODE": "1", } if os.environ.get("HOME"): environment["HOME"] = os.environ["HOME"] @@ -177,11 +236,12 @@ def installed_identity(interpreter: Path, *, cwd: Path) -> dict[str, Any]: "version": canonical_version(), "distribution_version": importlib.metadata.version("forge-autonomy"), "module": str(pathlib.Path(forge.__file__).resolve()), - "sys_executable": str(pathlib.Path(sys.executable).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), "-I", "-c", program), cwd=cwd) + result = _run((str(interpreter), "-B", "-I", "-c", program), cwd=cwd) try: identity = json.loads(result.stdout) except json.JSONDecodeError as error: @@ -215,8 +275,11 @@ class UpdateRequest: def validate(self) -> None: identifiers = (self.operation_id, self.runtime_id, self.installation_id) - if any(not value or any(character not in "abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789._-" for character in value) - for value in identifiers): + 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") @@ -241,51 +304,174 @@ def digest(self) -> str: return _digest_bytes(_json_bytes(asdict(self))) -def validate_qualified_artifact(request: UpdateRequest) -> dict[str, Any]: +def _validated_wheel(request: UpdateRequest) -> tuple[bytes, dict[str, str]]: wheel = Path(request.wheel) - if file_digest(wheel) != request.wheel_sha256: - raise InstalledForgeUpdateError("wheel digest does not match the selected qualified artifact") 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(wheel) as archive: - names = archive.namelist() - if any(name.startswith("/") or ".." in Path(name).parts for name in names): - raise InstalledForgeUpdateError("wheel contains an unsafe member path") - metadata_names = [name for name in names if name.endswith(".dist-info/METADATA")] - if len(metadata_names) != 1: + 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_names[0])) - except zipfile.BadZipFile as error: - raise InstalledForgeUpdateError("wheel is not a valid ZIP artifact") from error - if metadata.get("Name") != "forge-autonomy" or metadata.get("Version") != request.version: - raise InstalledForgeUpdateError("wheel package metadata does not match Forge 2.7.22") - + 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) - if file_digest(receipt_path) != request.qualification_receipt_sha256: + receipt_bytes = _read_regular_bytes(receipt_path) + if _digest_bytes(receipt_bytes) != request.qualification_receipt_sha256: raise InstalledForgeUpdateError("qualification receipt digest changed") - receipt = _read_json(receipt_path) + 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 ( - 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 not isinstance(qualification, dict) + 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 not isinstance(artifacts, dict) - or artifacts.get("wheel") != request.wheel_sha256 + 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 qualification does not bind the exact wheel and source") - return { - "wheel": str(wheel), "wheel_sha256": request.wheel_sha256, + 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: @@ -304,13 +490,15 @@ def _table_digest(connection: sqlite3.Connection, table: str) -> tuple[int, str] return len(rows), _digest_bytes(_json_bytes(rows)) -def database_snapshot(path: Path) -> dict[str, Any]: +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 = sqlite3.connect(path.resolve().as_uri() + "?mode=ro", uri=True) + connection = existing_connection or sqlite3.connect(path.resolve().as_uri() + "?mode=ro", uri=True) connection.row_factory = sqlite3.Row - connection.execute("PRAGMA query_only=ON") + 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]) @@ -355,7 +543,7 @@ def database_snapshot(path: Path) -> dict[str, Any]: except sqlite3.Error as error: raise InstalledForgeUpdateError("runtime database readback failed") from error finally: - if "connection" in locals(): + if owns_connection and "connection" in locals(): connection.close() snapshot = { "database": str(path), "integrity_check": integrity, @@ -402,12 +590,11 @@ def assert_quiescent(snapshot: Mapping[str, Any]) -> None: ): raise InstalledForgeUpdateError("Forge dispatcher is not durably idle") missions = writer.get("missions") or [] - active_missions = [row for row in missions if isinstance(row, Mapping) and row.get("status") in ACTIVE_MISSION_STATES] - if active_missions: - raise InstalledForgeUpdateError("Forge has active or automatically resumable Mission state") + 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(isinstance(row, Mapping) and row.get("state") in ACTIVE_SUBMISSION_STATES for row in submissions): - raise InstalledForgeUpdateError("Forge has an active scheduler submission") + 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") @@ -493,9 +680,9 @@ 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") - temporary = destination.with_name(f".{destination.name}.tmp-{os.getpid()}") - if temporary.exists(): - temporary.unlink() + 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) @@ -541,6 +728,87 @@ def _candidate_migrate(executable: Path, data_root: Path, *, cwd: Path) -> dict[ 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, @@ -561,6 +829,7 @@ def __init__( 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" @@ -586,9 +855,14 @@ def _state(self) -> dict[str, Any]: return state def _advance(self, state: dict[str, Any], phase: str, **evidence: object) -> dict[str, Any]: - updated = {**state, **evidence, "phase": phase, "updated_at": _now()} + 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", [])) - history.append({"phase": phase, "at": updated["updated_at"]}) + 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 @@ -609,7 +883,7 @@ def _processes(self) -> Sequence[str]: pid, parent = int(fields[0]), int(fields[1]) except ValueError: continue - if pid in own or parent in own: + if pid in own: continue processes.append(fields[2]) return processes @@ -624,58 +898,79 @@ def _assert_no_runtime_process(self) -> None: raise InstalledForgeUpdateError("a selected Forge runtime process is still active") def _stage(self, state: dict[str, Any]) -> dict[str, Any]: - qualification = validate_qualified_artifact(self.request) + 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: - raise InstalledForgeUpdateError("immutable candidate slot belongs to a different request") - elif self.slot.exists() and not self.slot.is_symlink(): - identity = installed_identity(self.slot / "bin" / "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) + receipt.get("request_digest") != self.request.digest + or receipt.get("wheel_manifest_digest") != qualification["wheel_manifest_digest"] ): - raise InstalledForgeUpdateError("unreceipted candidate slot cannot be safely adopted") - _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, "staged_at": _now(), - "recovered_after_atomic_slot_move": True, - }) + 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: - stage = self.runtime_root / "slots" / f".stage-{self.request.operation_id}" - _safe_directory(stage.parent, create=True) - if stage.is_symlink(): + # 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") - _run((self.request.base_python, "-m", "venv", str(stage)), cwd=self.runtime_root) - candidate_python = stage / "bin" / "python" - _run((str(candidate_python), "-m", "pip", "install", "--no-index", "--no-deps", self.request.wheel), - cwd=self.runtime_root) + 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(stage.resolve()) + os.sep) - or Path(str(identity.get("prefix"))).resolve() != stage.resolve() + 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") - if self.slot.exists() or self.slot.is_symlink(): - raise InstalledForgeUpdateError("candidate slot appeared concurrently") - os.replace(stage, self.slot) - 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) + os.sep): - raise InstalledForgeUpdateError("staged candidate readback changed") - if not self.slot_receipt.exists(): _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, "staged_at": _now(), + "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)) + 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) @@ -683,13 +978,23 @@ def _adopt_resolver(self, state: dict[str, Any]) -> dict[str, Any]: _safe_directory(self.runtime_root / "bin", create=True) _safe_directory(self.runtime_root / "fenced", create=True) if not self.legacy_entrypoint.exists(): - if resolver.is_symlink() or not resolver.is_file() or file_digest(resolver) != self.request.resolver_sha256: + 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") - shutil.copy2(resolver, self.legacy_entrypoint) - os.chmod(self.legacy_entrypoint, resolver.stat().st_mode & 0o777) + 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 = ( @@ -731,6 +1036,27 @@ def _backup(self, state: dict[str, Any], before: Mapping[str, Any]) -> dict[str, 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() @@ -741,29 +1067,50 @@ def _backup(self, state: dict[str, Any], before: Mapping[str, Any]) -> dict[str, 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" - temporary = root / f".forge.db.tmp-{os.getpid()}" - shutil.copy2(self.backup_path, temporary) + 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.write_text(self.request.runtime_id + "\n", encoding="utf-8") - os.chmod(marker, 0o600) + 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_migrate(self.slot / "bin" / "forge", root, cwd=self.runtime_root) + 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"], }) @@ -782,17 +1129,84 @@ def _restore_legacy_before_migration(self, state: dict[str, Any], error: Excepti 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") - output = _candidate_migrate(self.slot / "bin" / "forge", self.data_root, cwd=self.runtime_root) - after = database_snapshot(self.database) + after = self._install_qualified_database(before) preservation = verify_preservation(before, after, self.request) state = self._advance( state, "MIGRATED", live_migration={ - **preservation, "migrated_at": _now(), "candidate_output": output, + **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", @@ -801,6 +1215,7 @@ def _migrate_live(self, state: dict[str, Any], before: Mapping[str, Any]) -> tup 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={ @@ -832,10 +1247,10 @@ def _activate(self, state: dict[str, Any], after: Mapping[str, Any]) -> dict[str if ( identity.get("version") != self.request.version or version != self.request.version - or identity.get("sys_executable") != str((self.slot / "bin" / "python").resolve()) - or not str(identity.get("module", "")).startswith(str(self.slot) + os.sep) + 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 status.get("data_root") != self.request.data_root + 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) ): @@ -853,92 +1268,181 @@ def _activate(self, state: dict[str, Any], after: Mapping[str, Any]) -> dict[str def _secure_failure(self, state: dict[str, Any], error: Exception) -> None: """Leave a pre-migration legacy route or a schema-38-safe candidate/fence.""" - current = database_snapshot(self.database) - if current.get("user_version") == SCHEMA_BEFORE and self.legacy_entrypoint.exists(): + 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 - candidate = (self.slot / "bin" / "forge").resolve() - if not self.current.is_symlink() or _resolved_link(self.current) != candidate: - _replace_symlink(self.current, os.path.relpath(self.fenced_resolver, self.runtime_root)) + _replace_symlink(self.current, os.path.relpath(self.fenced_resolver, self.runtime_root)) self._advance( state, state.get("phase", "RECOVERY_PENDING"), - safety_disposition="SCHEMA_38_OLD_BINARY_FENCED", last_error=str(error), + 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") + + 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) - _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) - if receipt.get("request_digest") != self.request.digest: - raise InstalledForgeUpdateError("completed update receipt conflicts with this request") - return receipt - - state = self._stage(state) - self._interrupt("stage") update_lock = self.runtime_root / "locks" / "installation-update.lock" runtime_lock = self.data_root / "forge-runtime-mutation.lock" - with ExitStack() as locks: - locks.enter_context(exclusive_lock(update_lock)) - locks.enter_context(exclusive_lock(runtime_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") + 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 - except Exception as error: + + 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: - self._secure_failure(state, error) - except Exception: - pass - raise + 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: diff --git a/tests/test_installed_forge_update.py b/tests/test_installed_forge_update.py index 4723c66..ca1e4d8 100644 --- a/tests/test_installed_forge_update.py +++ b/tests/test_installed_forge_update.py @@ -2,13 +2,16 @@ 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 @@ -33,7 +36,7 @@ class InstalledForgeUpdateTests(unittest.TestCase): def setUp(self) -> None: self.temporary = tempfile.TemporaryDirectory() self.addCleanup(self.temporary.cleanup) - self.root = Path(self.temporary.name) + 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() @@ -43,17 +46,63 @@ def setUp(self) -> None: 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: - archive.writestr("forge_autonomy-2.7.22.dist-info/METADATA", metadata) + 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": "release-1", - "artifacts": {"wheel": self.wheel_digest, "sdist": "sha256:" + "d" * 64}, + "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( @@ -150,11 +199,24 @@ def _controller(self) -> object: 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, "does not bind"): + with self.assertRaisesRegex(update.InstalledForgeUpdateError, "noncanonical"): update.validate_qualified_artifact(changed) self.wheel.write_bytes(b"changed") with self.assertRaisesRegex(update.InstalledForgeUpdateError, "wheel digest"): @@ -210,6 +272,81 @@ def test_concurrent_update_lock_is_rejected(self) -> None: 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() @@ -227,8 +364,11 @@ def test_crash_before_migration_restores_the_legacy_route(self) -> None: 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}): + 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) @@ -259,7 +399,7 @@ def test_crash_during_atomic_activation_leaves_the_candidate_active(self) -> Non ) with patch.object(update, "database_snapshot", return_value={"user_version": 38}): controller._secure_failure(state, RuntimeError("interrupted")) - self.assertEqual(self.resolver.resolve(), candidate.resolve()) + self.assertEqual(self.resolver.resolve(), controller.fenced_resolver.resolve()) def test_preservation_rejects_domain_loss_and_unexpected_schema_objects(self) -> None: self._installed_schema37() @@ -278,6 +418,113 @@ def test_preservation_rejects_domain_loss_and_unexpected_schema_objects(self) -> 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_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) + if __name__ == "__main__": unittest.main() From 21ad57bca56734a8ab33fe24e14e855436e24e26 Mon Sep 17 00:00:00 2001 From: pcvantol Date: Fri, 18 Sep 2026 00:40:39 +0200 Subject: [PATCH 3/3] Bind completed update to schema fingerprint --- scripts/update_installed_forge.py | 20 ++++++++++++++++++++ tests/test_installed_forge_update.py | 21 +++++++++++++++++++++ 2 files changed, 41 insertions(+) diff --git a/scripts/update_installed_forge.py b/scripts/update_installed_forge.py index 9d7df71..03398c1 100644 --- a/scripts/update_installed_forge.py +++ b/scripts/update_installed_forge.py @@ -505,6 +505,10 @@ def database_snapshot(path: Path, *, existing_connection: sqlite3.Connection | N 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 @@ -548,6 +552,7 @@ def database_snapshot(path: Path, *, existing_connection: sqlite3.Connection | N 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": { @@ -614,6 +619,19 @@ def assert_quiescent(snapshot: Mapping[str, Any]) -> None: 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") @@ -1261,6 +1279,7 @@ def _activate(self, state: dict[str, Any], after: Mapping[str, Any]) -> dict[str 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", @@ -1361,6 +1380,7 @@ def _verify_complete(self, state: Mapping[str, Any], receipt: Mapping[str, Any]) 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) diff --git a/tests/test_installed_forge_update.py b/tests/test_installed_forge_update.py index ca1e4d8..6eaaa70 100644 --- a/tests/test_installed_forge_update.py +++ b/tests/test_installed_forge_update.py @@ -453,6 +453,21 @@ def test_quiescence_uses_explicit_safe_state_allowlists(self) -> None: 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() @@ -524,6 +539,12 @@ def test_exact_published_wheel_end_to_end_when_requested(self) -> None: 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__":