diff --git a/ddtrace/internal/ipc.py b/ddtrace/internal/ipc.py index 94de9af0ed7..8afe5f05139 100644 --- a/ddtrace/internal/ipc.py +++ b/ddtrace/internal/ipc.py @@ -4,6 +4,7 @@ from pathlib import Path import secrets import tempfile +import types import typing from ddtrace.internal import forksafe @@ -16,20 +17,40 @@ MAX_FILE_SIZE = 8192 +# The shared file is always opened in binary mode, whichever locking backend is +# in use. +SharedFile = typing.IO[bytes] + + +class _DummyFile(io.BytesIO): + """Stand-in yielded when the shared file is unavailable. + + Callers still get a usable file object, so they run to completion instead of + having to handle an exception, but nothing written to it is shared with any + other process. put_unlocked reports writes to it as failures so that callers + which retry, or which remember what they have already sent, are not told the + data was stored. + """ + class BaseLock: - def __init__(self, file: typing.IO[typing.Any]): + def __init__(self, file: SharedFile) -> None: self.file = file - def acquire(self): ... + def acquire(self) -> None: ... - def release(self): ... + def release(self) -> None: ... - def __enter__(self): + def __enter__(self) -> "BaseLock": self.acquire() return self - def __exit__(self, exc_type, exc_value, exc_tb): + def __exit__( + self, + exc_type: typing.Optional[type[BaseException]], + exc_value: typing.Optional[BaseException], + exc_tb: typing.Optional[types.TracebackType], + ) -> None: self.release() @@ -41,14 +62,14 @@ def __exit__(self, exc_type, exc_value, exc_tb): class BaseUnixLock(BaseLock): __acquire_mode__: typing.Optional[int] = None - def acquire(self): + def acquire(self) -> None: if self.__acquire_mode__ is None: msg = f"Cannot use lock of type {type(self)} directly" raise ValueError(msg) fcntl.lockf(self.file, self.__acquire_mode__) - def release(self): + def release(self) -> None: fcntl.lockf(self.file, fcntl.LOCK_UN) class ReadLock(BaseUnixLock): @@ -61,6 +82,11 @@ class WriteLock(BaseUnixLock): except ModuleNotFoundError: # Availability: Windows + # + # The defs in this branch are deliberately left unannotated: their bodies + # call msvcrt/_winapi APIs that only exist in the Windows stubs, and mypy + # runs against Linux/macOS, where checking them would report every call as + # a missing attribute. import msvcrt class BaseWinLock(BaseLock): @@ -105,7 +131,13 @@ def __init__(self, name: typing.Optional[str] = None) -> None: str(TMPDIR / (name or secrets.token_hex(8))) if TMPDIR is not None else None ) if self.filename is not None: - Path(self.filename).touch(exist_ok=True) + try: + Path(self.filename).touch(exist_ok=True) + except OSError: + # The temp dir is not writable, or the file exists but belongs to + # another user (e.g. it was created before dropping privileges). + log.debug("Cannot create shared file %s; disabling it", self.filename, exc_info=True) + self.filename = None # Thread-level lock to serialize access within the same process. # POSIX advisory file locks (fcntl.lockf) do NOT block threads within # the same process, so concurrent put/snatchall from different threads @@ -113,7 +145,9 @@ def __init__(self, name: typing.Optional[str] = None) -> None: # forksafe.Lock() resets after fork so children don't inherit a locked mutex. self._file_thread_lock = forksafe.Lock() - def put_unlocked(self, f: typing.BinaryIO, data: str) -> bool: + def put_unlocked(self, f: SharedFile, data: str) -> bool: + if isinstance(f, _DummyFile): + return False f.seek(0, os.SEEK_END) dt = (data + "\x00").encode() if f.tell() + len(dt) <= MAX_FILE_SIZE: @@ -132,7 +166,7 @@ def put(self, data: str) -> bool: except Exception: # nosec return False - def peekall_unlocked(self, f: typing.BinaryIO) -> list[str]: + def peekall_unlocked(self, f: SharedFile) -> list[str]: f.seek(0) return data.decode().split("\x00") if (data := f.read().strip(b"\x00")) else [] @@ -161,7 +195,7 @@ def snatchall(self) -> list[str]: except Exception: # nosec return [] - def clear_unlocked(self, f: typing.BinaryIO) -> None: + def clear_unlocked(self, f: SharedFile) -> None: f.seek(0) f.truncate() @@ -176,31 +210,49 @@ def clear(self) -> None: except Exception: # nosec pass + def _try_open(self, filename: str, mode: str) -> typing.Optional[SharedFile]: + # Open the shared file, returning None if it cannot be opened. A + # PermissionError is not transient (the file belongs to another user, + # e.g. it was created before the process dropped privileges), so the + # shared file is disabled for good rather than failing on every access. + try: + return open_file(filename, mode) + except PermissionError: + log.debug("Cannot open shared file %s; disabling it", filename, exc_info=True) + self.filename = None + except OSError: + log.debug("Cannot open shared file %s", filename, exc_info=True) + return None + @contextmanager - def lock_shared(self): + def lock_shared(self) -> typing.Iterator[SharedFile]: """Context manager to acquire a shared/read lock on the file.""" if self.filename is None: - # No writable temp dir (e.g. readOnlyRootFilesystem). Yield a dummy - # file-like so the context manager always yields and callers get [] from peek. - yield io.BytesIO(b"") + # No writable temp dir (e.g. readOnlyRootFilesystem). + yield _DummyFile() return with self._file_thread_lock: - with open_file(self.filename, "rb") as f, ReadLock(f): + if (f := self._try_open(self.filename, "rb")) is None: + yield _DummyFile() + return + with f, ReadLock(f): yield f @contextmanager - def lock_exclusive(self): + def lock_exclusive(self) -> typing.Iterator[SharedFile]: """Context manager to acquire an exclusive/write lock on the file.""" if self.filename is None: - # No writable temp dir (e.g. readOnlyRootFilesystem). Yield a dummy - # file-like so the context manager always yields and callers run without failing. - yield io.BytesIO(b"") + # No writable temp dir (e.g. readOnlyRootFilesystem). + yield _DummyFile() return # Acquire the thread-level lock first to prevent same-process threads # from bypassing the POSIX file lock (fcntl.lockf only blocks across # processes, not within the same process). with self._file_thread_lock: - with open_file(self.filename, "r+b") as f, WriteLock(f): + if (f := self._try_open(self.filename, "r+b")) is None: + yield _DummyFile() + return + with f, WriteLock(f): yield f # Flush before releasing the lock. Here we first release the lock, # then close the file. If a read happens in between these two diff --git a/tests/internal/test_ipc.py b/tests/internal/test_ipc.py new file mode 100644 index 00000000000..9681b401af2 --- /dev/null +++ b/tests/internal/test_ipc.py @@ -0,0 +1,135 @@ +import os +from pathlib import Path + +import pytest + +from ddtrace.internal import ipc +from ddtrace.internal.ipc import SharedStringFile + + +def test_shared_string_file_roundtrip(tmp_path, monkeypatch): + monkeypatch.setattr(ipc, "TMPDIR", tmp_path) + + ssf = SharedStringFile("roundtrip") + assert ssf.put("hello") + assert ssf.put("world") + assert ssf.peekall() == ["hello", "world"] + assert ssf.snatchall() == ["hello", "world"] + assert ssf.peekall() == [] + + +def test_shared_string_file_uncreatable(tmp_path, monkeypatch): + # A shared file we cannot even create degrades to a no-op instead of raising. + # A regular file standing in for the temp dir gives ENOTDIR, which -- unlike + # a read-only directory -- root cannot bypass either. + blocker = tmp_path / "not-a-dir" + blocker.touch() + monkeypatch.setattr(ipc, "TMPDIR", blocker) + + ssf = SharedStringFile("nope") + + assert ssf.filename is None + assert not ssf.put("hello") + assert ssf.peekall() == [] + assert ssf.snatchall() == [] + with ssf.lock_exclusive() as f: + assert ssf.peekall_unlocked(f) == [] + + +@pytest.mark.parametrize("mode", ["r+b", "rb"]) +def test_shared_string_file_permission_error_disables(tmp_path, monkeypatch, mode): + # A file we cannot open (e.g. owned by another user after dropping + # privileges) must not make callers fail, and must not be retried. + monkeypatch.setattr(ipc, "TMPDIR", tmp_path) + ssf = SharedStringFile("denied") + assert ssf.put("hello") + + calls = [] + + def denied(path, _mode): + calls.append(path) + raise PermissionError(13, "Permission denied", path) + + monkeypatch.setattr(ipc, "open_file", denied) + + lock = ssf.lock_exclusive if mode == "r+b" else ssf.lock_shared + with lock() as f: + assert ssf.peekall_unlocked(f) == [] + + assert ssf.filename is None + assert len(calls) == 1 + + # No further attempts to open the file are made. + assert ssf.peekall() == [] + assert not ssf.put("world") + with lock() as f: + assert ssf.peekall_unlocked(f) == [] + assert len(calls) == 1 + + +def test_shared_string_file_transient_error_not_disabling(tmp_path, monkeypatch): + # A transient failure degrades the current operation only. + monkeypatch.setattr(ipc, "TMPDIR", tmp_path) + ssf = SharedStringFile("transient") + assert ssf.put("hello") + + filename = ssf.filename + real_open_file = ipc.open_file + + def missing(path, mode): + raise FileNotFoundError(2, "No such file or directory", path) + + monkeypatch.setattr(ipc, "open_file", missing) + assert ssf.peekall() == [] + assert ssf.filename == filename + + monkeypatch.setattr(ipc, "open_file", real_open_file) + assert ssf.peekall() == ["hello"] + + +@pytest.mark.skipif( + not hasattr(os, "getuid") or os.getuid() == 0, + reason="needs POSIX permissions enforced against a non-root user", +) +def test_shared_string_file_existing_file_not_owned(tmp_path, monkeypatch): + # The real-world scenario: the shared file already exists and the current + # process has no permission on it. + monkeypatch.setattr(ipc, "TMPDIR", tmp_path) + path = tmp_path / "unowned" + path.touch() + os.chmod(path, 0o000) + + ssf = SharedStringFile("unowned") + assert ssf.filename is None or Path(ssf.filename) == path + + # Whatever failed, no exception escapes to the caller. + assert ssf.peekall() == [] + assert not ssf.put("hello") + with ssf.lock_exclusive() as f: + assert ssf.peekall_unlocked(f) == [] + assert ssf.filename is None + + +# errno 13 makes OSError construct a PermissionError, so the transient case uses +# EIO to stay a plain OSError and exercise the other branch of _try_open. +@pytest.mark.parametrize("errno_", [13, 5], ids=["permission-denied", "transient"]) +def test_shared_string_file_fallback_write_reports_failure(tmp_path, monkeypatch, errno_): + # A write that lands in the fallback must report failure. Callers such as + # Config._add_extra_service memoize on a True return and never retry, so + # claiming success would silently drop the value for good. + monkeypatch.setattr(ipc, "TMPDIR", tmp_path) + ssf = SharedStringFile("fallback") + + def failing(path, mode): + raise OSError(errno_, "nope", path) + + monkeypatch.setattr(ipc, "open_file", failing) + + assert not ssf.put("hello") + with ssf.lock_exclusive() as f: + assert not ssf.put_unlocked(f, "hello") + + # And nothing was actually stored: the real file is still empty. + monkeypatch.undo() + monkeypatch.setattr(ipc, "TMPDIR", tmp_path) + assert SharedStringFile("fallback").peekall() == []