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
94 changes: 73 additions & 21 deletions ddtrace/internal/ipc.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
from pathlib import Path
import secrets
import tempfile
import types
import typing

from ddtrace.internal import forksafe
Expand All @@ -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()


Expand All @@ -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):
Expand All @@ -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):
Expand Down Expand Up @@ -105,15 +131,23 @@ 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
# would race on the Python write buffer vs OS flush window.
# 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:
Expand All @@ -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 []

Expand Down Expand Up @@ -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()

Expand All @@ -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
Comment thread
P403n1x87 marked this conversation as resolved.
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
Expand Down
135 changes: 135 additions & 0 deletions tests/internal/test_ipc.py
Original file line number Diff line number Diff line change
@@ -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() == []
Loading