From 45dbc30301f9495c65a395ffe3b218ef20a30a08 Mon Sep 17 00:00:00 2001 From: Abdessamad Derraz <3028866+Abdess@users.noreply.github.com> Date: Sat, 8 Aug 2026 04:24:03 +0200 Subject: [PATCH] feat: lock pack artifacts during pipeline runs --- scripts/common.py | 46 +++++++++++++ scripts/generate_pack.py | 38 ++++++++--- scripts/pipeline.py | 32 +++++---- tests/test_artifact_lock.py | 130 ++++++++++++++++++++++++++++++++++++ 4 files changed, 223 insertions(+), 23 deletions(-) create mode 100644 tests/test_artifact_lock.py diff --git a/scripts/common.py b/scripts/common.py index d6117080..c6a287cb 100644 --- a/scripts/common.py +++ b/scripts/common.py @@ -6,6 +6,7 @@ and file resolution - eliminates DRY violations across scripts. from __future__ import annotations +import contextlib import hashlib import json import os @@ -258,6 +259,15 @@ def load_data_dir_registry(platforms_dir: str = "platforms") -> dict: return data.get("data_directories", {}) +def load_platform_registry(platforms_dir: str = "platforms") -> dict: + """Return the `platforms:` mapping from _registry.yml, empty if absent.""" + registry_path = os.path.join(platforms_dir, "_registry.yml") + if not os.path.exists(registry_path): + return {} + with open(registry_path) as f: + return (yaml.safe_load(f) or {}).get("platforms", {}) + + def list_registered_platforms( platforms_dir: str = "platforms", include_archived: bool = False, @@ -1424,3 +1434,39 @@ def write_provenance_snapshot( "entries": sorted(entries, key=lambda e: (e["dat"], e["name"])), } return write_if_changed(path, json.dumps(snapshot, indent=2) + "\n") + + +class ArtifactLockBusy(RuntimeError): + """Raised when another process already holds the artifact directory.""" + + +@contextlib.contextmanager +def artifact_lock(directory: str, exclusive: bool = True): + """Serialize access to a shared artifact directory across processes. + + Two pipeline runs building the same dist/ leave readers looking at + half-written ZIPs, which surfaces as BadZipFile far from its cause. + Writers take the lock exclusively, readers share it. On platforms + without flock the lock is a no-op. + """ + try: + import fcntl + except ImportError: + yield + return + + os.makedirs(directory, exist_ok=True) + lock_path = os.path.join(directory, ".lock") + mode = fcntl.LOCK_EX if exclusive else fcntl.LOCK_SH + with open(lock_path, "w") as handle: + try: + fcntl.flock(handle, mode | fcntl.LOCK_NB) + except OSError as exc: + raise ArtifactLockBusy( + f"{directory} is in use by another run " + f"(lock: {lock_path}). Wait for it to finish." + ) from exc + try: + yield + finally: + fcntl.flock(handle, fcntl.LOCK_UN) diff --git a/scripts/generate_pack.py b/scripts/generate_pack.py index 9a94a1db..d429fc33 100644 --- a/scripts/generate_pack.py +++ b/scripts/generate_pack.py @@ -13,6 +13,7 @@ and 3-tier storage (embedded/external/user_provided). from __future__ import annotations import argparse +import contextlib import hashlib import json import os @@ -26,7 +27,9 @@ from pathlib import Path sys.path.insert(0, os.path.dirname(__file__)) from common import ( + ArtifactLockBusy, MANUFACTURER_PREFIXES, + artifact_lock, build_target_cores_cache, build_zip_contents_index, check_inside_zip, @@ -2486,6 +2489,17 @@ def _run_manifest_mode( print(f" ERROR: {e}") +@contextlib.contextmanager +def _pack_output_lock(output_dir: str, exclusive: bool = True): + """Hold the output directory for the duration of a pack run.""" + try: + with artifact_lock(output_dir, exclusive=exclusive): + yield + except ArtifactLockBusy as exc: + print(f"ERROR: {exc}") + sys.exit(1) + + def _run_verify_packs(args): """Extract each pack and verify file paths + hashes.""" import shutil @@ -2870,7 +2884,8 @@ def main(): # Quick-exit modes: --verify-packs alone = verify existing packs only # Combined with --all-variants, generation runs first then verify if args.verify_packs and not args.all_variants: - _run_verify_packs(args) + with _pack_output_lock(args.output_dir, exclusive=False): + _run_verify_packs(args) return if args.manifest_targets: generate_target_manifests( @@ -3028,16 +3043,17 @@ def main(): args, groups, db, zip_contents, emu_profiles, target_cores_cache ) else: - _run_platform_packs( - args, - groups, - db, - zip_contents, - data_registry, - emu_profiles, - target_cores_cache, - system_filter, - ) + with _pack_output_lock(args.output_dir): + _run_platform_packs( + args, + groups, + db, + zip_contents, + data_registry, + emu_profiles, + target_cores_cache, + system_filter, + ) # Manifest generation (JSON inventory for install.py) diff --git a/scripts/pipeline.py b/scripts/pipeline.py index 9825993a..9e46969d 100644 --- a/scripts/pipeline.py +++ b/scripts/pipeline.py @@ -29,6 +29,9 @@ import sys import time from pathlib import Path +sys.path.insert(0, str(Path(__file__).parent)) +from common import ArtifactLockBusy, artifact_lock + def run(cmd: list[str], label: str) -> tuple[bool, str]: """Run a command. Returns (success, captured_output).""" @@ -326,18 +329,23 @@ def main(): # picked up by pack verification and shipped in releases. out_dir = Path(args.output_dir) if out_dir.is_dir(): - stale = [ - p for p in out_dir.iterdir() - if p.is_file() and ( - p.suffix == ".zip" - or ".zip." in p.name - or p.name == "SHA256SUMS.txt" - ) - ] - for p in stale: - p.unlink() - if stale: - print(f"Purged {len(stale)} stale pack file(s) from {out_dir}/") + try: + with artifact_lock(str(out_dir)): + stale = [ + p for p in out_dir.iterdir() + if p.is_file() and ( + p.suffix == ".zip" + or ".zip." in p.name + or p.name == "SHA256SUMS.txt" + ) + ] + for p in stale: + p.unlink() + if stale: + print(f"Purged {len(stale)} stale pack file(s) from {out_dir}/") + except ArtifactLockBusy as exc: + print(f"ERROR: {exc}") + sys.exit(1) pack_cmd = [ sys.executable, diff --git a/tests/test_artifact_lock.py b/tests/test_artifact_lock.py new file mode 100644 index 00000000..032336f5 --- /dev/null +++ b/tests/test_artifact_lock.py @@ -0,0 +1,130 @@ +#!/usr/bin/env python3 +"""Concurrency guard on the pack output directory. + +Two runs building the same dist/ leave readers looking at half-written +ZIPs, which surfaces as BadZipFile far from its cause. +""" + +from __future__ import annotations + +import os +import subprocess +import sys +import tempfile +import unittest +from pathlib import Path + +REPO_ROOT = Path(__file__).resolve().parent.parent +sys.path.insert(0, str(REPO_ROOT / "scripts")) + +from common import ArtifactLockBusy, artifact_lock # noqa: E402 + + +class ArtifactLockTest(unittest.TestCase): + def setUp(self): + self.dir = tempfile.mkdtemp() + + def test_writer_excludes_writer(self): + with artifact_lock(self.dir): + with self.assertRaises(ArtifactLockBusy): + with artifact_lock(self.dir): + pass + + def test_writer_excludes_reader(self): + with artifact_lock(self.dir): + with self.assertRaises(ArtifactLockBusy): + with artifact_lock(self.dir, exclusive=False): + pass + + def test_readers_share(self): + with artifact_lock(self.dir, exclusive=False): + with artifact_lock(self.dir, exclusive=False): + pass + + def test_lock_is_released_on_exit(self): + with artifact_lock(self.dir): + pass + with artifact_lock(self.dir): + pass + + def test_lock_is_released_on_error(self): + with self.assertRaises(ValueError): + with artifact_lock(self.dir): + raise ValueError("boom") + with artifact_lock(self.dir): + pass + + def test_separate_directories_do_not_contend(self): + other = tempfile.mkdtemp() + with artifact_lock(self.dir): + with artifact_lock(other): + pass + + def test_lock_survives_a_missing_directory(self): + missing = os.path.join(self.dir, "not-created-yet") + with artifact_lock(missing): + self.assertTrue(os.path.isdir(missing)) + + +class PackLockCliTest(unittest.TestCase): + """generate_pack and pipeline must refuse a directory held elsewhere.""" + + def setUp(self): + self.dir = tempfile.mkdtemp() + + def _run(self, argv: list[str]) -> subprocess.CompletedProcess: + return subprocess.run( + [sys.executable, *argv], + capture_output=True, + text=True, + cwd=str(REPO_ROOT), + timeout=300, + ) + + def test_generate_pack_refuses_locked_output(self): + with artifact_lock(self.dir): + proc = self._run( + [ + "scripts/generate_pack.py", + "--platform", + "misterfpga", + "--output-dir", + self.dir, + "--offline", + ] + ) + self.assertEqual(proc.returncode, 1, proc.stdout + proc.stderr) + self.assertIn("is in use by another run", proc.stdout) + + def test_verify_packs_refuses_a_writer(self): + with artifact_lock(self.dir): + proc = self._run( + [ + "scripts/generate_pack.py", + "--platform", + "misterfpga", + "--verify-packs", + "--output-dir", + self.dir, + ] + ) + self.assertEqual(proc.returncode, 1, proc.stdout + proc.stderr) + self.assertIn("is in use by another run", proc.stdout) + + def test_pipeline_refuses_locked_output(self): + with artifact_lock(self.dir): + proc = self._run( + [ + "scripts/pipeline.py", + "--offline", + "--skip-docs", + "--output-dir", + self.dir, + ] + ) + self.assertEqual(proc.returncode, 1, proc.stdout + proc.stderr) + self.assertIn("is in use by another run", proc.stdout) + + +if __name__ == "__main__": + unittest.main()