feat: lock pack artifacts during pipeline runs

This commit is contained in:
Abdessamad Derraz committed 2026-08-08 04:24:03 +02:00
1 parent f008568108
commit 45dbc30301
4 files changed
+223 -23

No files matched your search

+46
View File
@@ -6,6 +6,7 @@ and file resolution - eliminates DRY violations across scripts.
from __future__ import annotations from __future__ import annotations
import contextlib
import hashlib import hashlib
import json import json
import os import os
@@ -258,6 +259,15 @@ def load_data_dir_registry(platforms_dir: str = "platforms") -> dict:
return data.get("data_directories", {}) 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( def list_registered_platforms(
platforms_dir: str = "platforms", platforms_dir: str = "platforms",
include_archived: bool = False, include_archived: bool = False,
@@ -1424,3 +1434,39 @@ def write_provenance_snapshot(
"entries": sorted(entries, key=lambda e: (e["dat"], e["name"])), "entries": sorted(entries, key=lambda e: (e["dat"], e["name"])),
} }
return write_if_changed(path, json.dumps(snapshot, indent=2) + "\n") 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)
+27 -11
View File
@@ -13,6 +13,7 @@ and 3-tier storage (embedded/external/user_provided).
from __future__ import annotations from __future__ import annotations
import argparse import argparse
import contextlib
import hashlib import hashlib
import json import json
import os import os
@@ -26,7 +27,9 @@ from pathlib import Path
sys.path.insert(0, os.path.dirname(__file__)) sys.path.insert(0, os.path.dirname(__file__))
from common import ( from common import (
ArtifactLockBusy,
MANUFACTURER_PREFIXES, MANUFACTURER_PREFIXES,
artifact_lock,
build_target_cores_cache, build_target_cores_cache,
build_zip_contents_index, build_zip_contents_index,
check_inside_zip, check_inside_zip,
@@ -2486,6 +2489,17 @@ def _run_manifest_mode(
print(f" ERROR: {e}") 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): def _run_verify_packs(args):
"""Extract each pack and verify file paths + hashes.""" """Extract each pack and verify file paths + hashes."""
import shutil import shutil
@@ -2870,7 +2884,8 @@ def main():
# Quick-exit modes: --verify-packs alone = verify existing packs only # Quick-exit modes: --verify-packs alone = verify existing packs only
# Combined with --all-variants, generation runs first then verify # Combined with --all-variants, generation runs first then verify
if args.verify_packs and not args.all_variants: 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 return
if args.manifest_targets: if args.manifest_targets:
generate_target_manifests( generate_target_manifests(
@@ -3028,16 +3043,17 @@ def main():
args, groups, db, zip_contents, emu_profiles, target_cores_cache args, groups, db, zip_contents, emu_profiles, target_cores_cache
) )
else: else:
_run_platform_packs( with _pack_output_lock(args.output_dir):
args, _run_platform_packs(
groups, args,
db, groups,
zip_contents, db,
data_registry, zip_contents,
emu_profiles, data_registry,
target_cores_cache, emu_profiles,
system_filter, target_cores_cache,
) system_filter,
)
# Manifest generation (JSON inventory for install.py) # Manifest generation (JSON inventory for install.py)
+20 -12
View File
@@ -29,6 +29,9 @@ import sys
import time import time
from pathlib import Path 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]: def run(cmd: list[str], label: str) -> tuple[bool, str]:
"""Run a command. Returns (success, captured_output).""" """Run a command. Returns (success, captured_output)."""
@@ -326,18 +329,23 @@ def main():
# picked up by pack verification and shipped in releases. # picked up by pack verification and shipped in releases.
out_dir = Path(args.output_dir) out_dir = Path(args.output_dir)
if out_dir.is_dir(): if out_dir.is_dir():
stale = [ try:
p for p in out_dir.iterdir() with artifact_lock(str(out_dir)):
if p.is_file() and ( stale = [
p.suffix == ".zip" p for p in out_dir.iterdir()
or ".zip." in p.name if p.is_file() and (
or p.name == "SHA256SUMS.txt" p.suffix == ".zip"
) or ".zip." in p.name
] or p.name == "SHA256SUMS.txt"
for p in stale: )
p.unlink() ]
if stale: for p in stale:
print(f"Purged {len(stale)} stale pack file(s) from {out_dir}/") 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 = [ pack_cmd = [
sys.executable, sys.executable,
+130
View File
@@ -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()