diff --git a/src/poetry/installation/wheel_installer.py b/src/poetry/installation/wheel_installer.py index 50db97b9650..7a80c29c2fc 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,11 +33,34 @@ 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]: 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 @@ -46,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, @@ -100,11 +135,17 @@ 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) - - if is_executable: - make_file_executable(target_path) + # 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) + + 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..1110753e3fe 100644 --- a/tests/installation/test_wheel_installer.py +++ b/tests/installation/test_wheel_installer.py @@ -1,21 +1,33 @@ from __future__ import annotations +import io +import os 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 +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 if TYPE_CHECKING: + from typing import BinaryIO + from pytest import TempPathFactory + from pytest_mock import MockerFixture from tests.types import FixtureDirGetter @@ -144,3 +156,135 @@ 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 + + +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]