From 75b522452c92924009410dd4cb51e927cdf96d8f Mon Sep 17 00:00:00 2001 From: Suliman Abdulrazzaq Date: Fri, 25 Sep 2026 20:10:02 +0300 Subject: [PATCH 1/2] fix: serialize writes of files shared by wheels installed in parallel Several wheels may contain the same file, e.g. the __init__.py of a pkgutil-style namespace package. Since wheels are installed in parallel, two threads could write such a file at the same time, which could result in a file with interleaved contents, i.e. a corrupted installation. Therefore, write each file under a lock that is specific to its path. --- src/poetry/installation/wheel_installer.py | 30 ++++++- tests/installation/test_wheel_installer.py | 98 ++++++++++++++++++++++ 2 files changed, 124 insertions(+), 4 deletions(-) diff --git a/src/poetry/installation/wheel_installer.py b/src/poetry/installation/wheel_installer.py index 50db97b9650..54d964791a4 100644 --- a/src/poetry/installation/wheel_installer.py +++ b/src/poetry/installation/wheel_installer.py @@ -4,10 +4,12 @@ import os import platform import sys +import threading from functools import cached_property from pathlib import Path from typing import TYPE_CHECKING +from weakref import WeakValueDictionary from installer import install from installer.destinations import SchemeDictionaryDestination @@ -31,6 +33,25 @@ from poetry.utils.env import Env +# Several wheels may contain the same file (e.g. the __init__.py of a +# pkgutil-style namespace package). Since wheels are installed in parallel, +# two threads may write such a file at the same time, which results in a file +# with interleaved contents. Therefore, each file is written under a lock +# that is specific to its path. +_file_locks: WeakValueDictionary[str, threading.Lock] = WeakValueDictionary() +_file_locks_lock = threading.Lock() + + +def _file_lock(path: str) -> threading.Lock: + key = os.path.normcase(path) + with _file_locks_lock: + lock = _file_locks.get(key) + if lock is None: + lock = _file_locks[key] = threading.Lock() + + return lock + + class WheelDestination(SchemeDictionaryDestination): @cached_property def _abspath_scheme_cache(self) -> dict[Scheme, str]: @@ -100,11 +121,12 @@ def write_to_fs( # that two threads try to create the directory. parent_folder.mkdir(parents=True, exist_ok=True) - with target_path.open("wb") as f: - hash_, size = copyfileobj_with_hashing(stream, f, self.hash_algorithm) + with _file_lock(target_path_str): + with target_path.open("wb") as f: + hash_, size = copyfileobj_with_hashing(stream, f, self.hash_algorithm) - if is_executable: - make_file_executable(target_path) + if is_executable: + make_file_executable(target_path) return RecordEntry(path, Hash(self.hash_algorithm, hash_), size) diff --git a/tests/installation/test_wheel_installer.py b/tests/installation/test_wheel_installer.py index bdd38867b86..5835f199ac9 100644 --- a/tests/installation/test_wheel_installer.py +++ b/tests/installation/test_wheel_installer.py @@ -1,10 +1,14 @@ from __future__ import annotations import re +import threading +import zipfile +from concurrent.futures import ThreadPoolExecutor from pathlib import Path from typing import TYPE_CHECKING +import installer.utils import pytest from poetry.core.constraints.version import parse_constraint @@ -15,7 +19,10 @@ if TYPE_CHECKING: + from typing import BinaryIO + from pytest import TempPathFactory + from pytest_mock import MockerFixture from tests.types import FixtureDirGetter @@ -144,3 +151,94 @@ def test_no_path_traversal_via_symlink( assert target.read_text(encoding="utf-8") == "original" else: assert not list(target_dir.iterdir()) + + +def _build_wheel(directory: Path, name: str, files: dict[str, bytes]) -> Path: + dist_info = f"{name}-0.1.dist-info" + files = { + **files, + f"{dist_info}/WHEEL": ( + b"Wheel-Version: 1.0\nRoot-Is-Purelib: true\nTag: py3-none-any\n" + ), + f"{dist_info}/METADATA": ( + f"Metadata-Version: 2.1\nName: {name}\nVersion: 0.1\n".encode() + ), + } + files[f"{dist_info}/RECORD"] = ( + "\n".join([f"{k},," for k in files] + [f"{dist_info}/RECORD,,"]) + "\n" + ).encode() + + wheel = directory / f"{name}-0.1-py3-none-any.whl" + with zipfile.ZipFile(wheel, "w") as z: + for k, v in files.items(): + z.writestr(k, v) + + return wheel + + +def test_parallel_installation_of_file_contained_in_several_wheels( + tmp_path: Path, env: MockEnv, mocker: MockerFixture +) -> None: + """ + Several wheels may contain the same file, e.g. the __init__.py + of a pkgutil-style namespace package. If such wheels are installed + in parallel, the file must not end up with interleaved contents. + """ + shared_file = "namespace/__init__.py" + content_first = b"# first\n" * 100 + content_second = b"# second\n" + wheel_first = _build_wheel(tmp_path, "first", {shared_file: content_first}) + wheel_second = _build_wheel(tmp_path, "second", {shared_file: content_second}) + target = Path(env.paths["purelib"]) / shared_file + + first_write_started = threading.Event() + second_write_finished = threading.Event() + copyfileobj_with_hashing = installer.utils.copyfileobj_with_hashing + + class PausingWriter: + """Write a part of the data and pause to give the other thread the chance + to write the same file before the rest of the data is written.""" + + def __init__(self, dest: BinaryIO) -> None: + self._dest = dest + + def write(self, data: bytes) -> int: + written = self._dest.write(data[:10]) + self._dest.flush() + first_write_started.set() + # If writes of the same file are serialized, + # the second write cannot finish and the wait times out. + second_write_finished.wait(timeout=1) + return written + self._dest.write(data[10:]) + + def copy(source: BinaryIO, dest: BinaryIO, hash_algorithm: str) -> tuple[str, int]: + if Path(dest.name) != target: + return copyfileobj_with_hashing(source, dest, hash_algorithm) + if not first_write_started.is_set(): + return copyfileobj_with_hashing( + source, + PausingWriter(dest), # type: ignore[arg-type] + hash_algorithm, + ) + result = copyfileobj_with_hashing(source, dest, hash_algorithm) + second_write_finished.set() + return result + + mocker.patch("installer.utils.copyfileobj_with_hashing", new=copy) + + wheel_installer = WheelInstaller(env) + + def install_second() -> None: + assert first_write_started.wait(timeout=10) + wheel_installer.install(wheel_second) + + with ThreadPoolExecutor(max_workers=2) as executor: + futures = [ + executor.submit(wheel_installer.install, wheel_first), + executor.submit(install_second), + ] + for future in futures: + future.result() + + # The file written last wins. + assert target.read_bytes() == content_second From 36762e99d4ab40f1a7ebcd62273e1d67a2c534e8 Mon Sep 17 00:00:00 2001 From: Suliman Abdulrazzaq Date: Thu, 1 Oct 2026 23:14:05 +0300 Subject: [PATCH 2/2] fix: take the lock for a shared file from the real path of the scheme directory --- src/poetry/installation/wheel_installer.py | 21 +++++++++- tests/installation/test_wheel_installer.py | 46 ++++++++++++++++++++++ 2 files changed, 66 insertions(+), 1 deletion(-) diff --git a/src/poetry/installation/wheel_installer.py b/src/poetry/installation/wheel_installer.py index 54d964791a4..7a80c29c2fc 100644 --- a/src/poetry/installation/wheel_installer.py +++ b/src/poetry/installation/wheel_installer.py @@ -57,6 +57,10 @@ class WheelDestination(SchemeDictionaryDestination): def _abspath_scheme_cache(self) -> dict[Scheme, str]: return {} + @cached_property + def _realpath_scheme_cache(self) -> dict[Scheme, str]: + return {} + def _abspath_scheme_dir(self, scheme: Scheme) -> str: # The scheme directory is fixed per destination, so normalize it once # and cache it instead of recomputing os.path.abspath() for every @@ -67,6 +71,16 @@ def _abspath_scheme_dir(self, scheme: Scheme) -> str: target_dir = cache[scheme] = os.path.abspath(self.scheme_dict[scheme]) return target_dir + def _realpath_scheme_dir(self, scheme: Scheme) -> str: + # Different schemes can point to the same directory through a symlink, + # e.g. platlib is lib64/... and lib64 is a symlink to lib on Fedora. + # Resolve the scheme directory once per scheme, not every file. + cache = self._realpath_scheme_cache + real_dir = cache.get(scheme) + if real_dir is None: + real_dir = cache[scheme] = os.path.realpath(self.scheme_dict[scheme]) + return real_dir + def write_to_fs( self, scheme: Scheme, @@ -121,7 +135,12 @@ def write_to_fs( # that two threads try to create the directory. parent_folder.mkdir(parents=True, exist_ok=True) - with _file_lock(target_path_str): + # The lock must be the same for the same file, whichever scheme + # (and whichever spelling of its directory) the wheel installs it to. + lock_key = ( + self._realpath_scheme_dir(scheme) + target_path_str[len(target_dir) :] + ) + with _file_lock(lock_key): with target_path.open("wb") as f: hash_, size = copyfileobj_with_hashing(stream, f, self.hash_algorithm) diff --git a/tests/installation/test_wheel_installer.py b/tests/installation/test_wheel_installer.py index 5835f199ac9..1110753e3fe 100644 --- a/tests/installation/test_wheel_installer.py +++ b/tests/installation/test_wheel_installer.py @@ -1,5 +1,7 @@ from __future__ import annotations +import io +import os import re import threading import zipfile @@ -13,6 +15,9 @@ from poetry.core.constraints.version import parse_constraint +import poetry.installation.wheel_installer as wheel_installer_module + +from poetry.installation.wheel_installer import WheelDestination from poetry.installation.wheel_installer import WheelInstaller from poetry.utils._compat import WINDOWS from poetry.utils.env import MockEnv @@ -242,3 +247,44 @@ def install_second() -> None: # The file written last wins. assert target.read_bytes() == content_second + + +def test_lock_is_shared_by_scheme_directories_behind_a_symlink( + tmp_path: Path, mocker: MockerFixture +) -> None: + """ + On Fedora a venv's platlib is under lib64, which is a symlink to lib. + A purelib wheel and a platlib wheel that contain the same file must + be written under the same lock. + """ + lib = tmp_path / "lib" + purelib = lib / "site-packages" + purelib.mkdir(parents=True) + try: + os.symlink(lib, tmp_path / "lib64", target_is_directory=True) + except (OSError, NotImplementedError): + pytest.skip("symlinks are not available") + platlib = tmp_path / "lib64" / "site-packages" + + file_lock = mocker.spy(wheel_installer_module, "_file_lock") + destination = WheelDestination( + {"purelib": str(purelib), "platlib": str(platlib)}, + interpreter="python", + script_kind="posix", + ) + destination.write_to_fs( + installer.utils.Scheme("purelib"), + "namespace/__init__.py", + io.BytesIO(b"# shared\n"), + False, + ) + destination.write_to_fs( + installer.utils.Scheme("platlib"), + "namespace/__init__.py", + io.BytesIO(b"# shared\n"), + False, + ) + + keys = [call.args[0] for call in file_lock.call_args_list] + assert len(keys) == 2 + assert keys[0] == keys[1]