Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
51 changes: 46 additions & 5 deletions src/poetry/installation/wheel_installer.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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()
Comment thread
sourcery-ai[bot] marked this conversation as resolved.
_file_locks_lock = threading.Lock()


def _file_lock(path: str) -> threading.Lock:
key = os.path.normcase(path)
Comment thread
radoering marked this conversation as resolved.
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
Expand All @@ -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,
Expand Down Expand Up @@ -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)

Expand Down
144 changes: 144 additions & 0 deletions tests/installation/test_wheel_installer.py
Original file line number Diff line number Diff line change
@@ -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

Expand Down Expand Up @@ -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]
Loading