Skip to content

Audit and enforce safe blocking synchronization for Python threads #64

Description

@jakub-graczyk

Problem and proposed change

Rust code can run on a thread currently attached to Python. If that thread waits for a Rust lock while the lock holder must call Python before releasing it, neither thread can continue and the process can freeze.

Review every existing call to a blocking method listed in clippy.toml. Each call must:

  • use a PyO3 attachment-aware method;
  • run inside Python::detach;
  • have a caller that guarantees the thread is detached; or
  • use a documented lock whose guarded work never calls Python or destroys Python-owned values.

Exemption: the SharedRetrySession mutex. DowngradingConsistencyRetrySession logs while the lock is held, so a Python log handler can run under it. A deadlock needs a handler that calls back into the same retry session, which we consider unrealistic, so it was decided not to fix it (#200).

Keep the Clippy policy enabled so new unreviewed calls and obsolete exemptions fail CI.

Context

PR #194 implements most initial fixes and enforcement:

#194

Configured blocking methods:

# Blocking on a lock while attached to the interpreter can deadlock: the holder
# may have lost the GIL inside Python and needs it back to release the lock.
# Attached code waits through pyo3's `*_py_attached` methods. A plain call is
# only fine on a detached thread, inside `py.detach`, or when every holder of
# that lock releases it without needing the GIL (pure Rust under the lock: no
# Python call, no Python allocation, no `Py` drop). Each such call carries an
# `#[expect]` saying which.
disallowed-methods = [
{ path = "std::sync::Mutex::lock", reason = "use lock_py_attached, or #[expect] saying why this thread is detached or why holders release the lock without needing the GIL" },
{ path = "std::sync::RwLock::read", reason = "use read_py_attached, or #[expect] saying why this thread is detached or why holders release the lock without needing the GIL" },
{ path = "std::sync::RwLock::write", reason = "use write_py_attached, or #[expect] saying why this thread is detached or why holders release the lock without needing the GIL" },
{ path = "std::sync::Condvar::wait", reason = "wait inside py.detach and #[expect] this lint" },
{ path = "std::sync::Condvar::wait_while", reason = "wait inside py.detach and #[expect] this lint" },
{ path = "std::sync::Condvar::wait_timeout", reason = "wait inside py.detach and #[expect] this lint" },
{ path = "std::sync::Condvar::wait_timeout_while", reason = "wait inside py.detach and #[expect] this lint" },
{ path = "std::sync::OnceLock::get_or_init", reason = "use get_or_init_py_attached" },
{ path = "std::sync::OnceLock::get_or_try_init", reason = "use PyOnceLock::get_or_try_init(py, ..), or call inside py.detach and #[expect] this lint" },
{ path = "std::sync::OnceLock::set", reason = "call inside py.detach and #[expect] this lint; PyOnceLock::set can also block while attached" },
{ path = "pyo3::sync::PyOnceLock::set", reason = "blocks while attached if another thread is initializing; use get_or_init(py, ..)" },
{ path = "std::sync::OnceLock::wait", reason = "wait inside py.detach and #[expect] this lint" },
{ path = "std::sync::Once::call_once", reason = "use call_once_py_attached" },
{ path = "std::sync::Once::call_once_force", reason = "use call_once_force_py_attached" },
{ path = "std::sync::Once::wait", reason = "wait inside py.detach and #[expect] this lint" },
{ path = "std::sync::Once::wait_force", reason = "wait inside py.detach and #[expect] this lint" },
{ path = "tokio::sync::Mutex::blocking_lock", reason = "never block a Python thread on a tokio lock" },
{ path = "tokio::sync::Mutex::blocking_lock_owned", reason = "never block a Python thread on a tokio lock" },
{ path = "tokio::sync::RwLock::blocking_read", reason = "never block a Python thread on a tokio lock" },
{ path = "tokio::sync::RwLock::blocking_write", reason = "never block a Python thread on a tokio lock" },
{ path = "tokio::sync::oneshot::Receiver::blocking_recv", reason = "block inside py.detach and #[expect] this lint" },
{ path = "tokio::sync::mpsc::Receiver::blocking_recv", reason = "block inside py.detach and #[expect] this lint" },
{ path = "tokio::sync::mpsc::Receiver::blocking_recv_many", reason = "block inside py.detach and #[expect] this lint" },
{ path = "tokio::sync::mpsc::Sender::blocking_send", reason = "block inside py.detach and #[expect] this lint" },
{ path = "tokio::sync::mpsc::UnboundedReceiver::blocking_recv", reason = "block inside py.detach and #[expect] this lint" },
{ path = "tokio::sync::mpsc::UnboundedReceiver::blocking_recv_many", reason = "block inside py.detach and #[expect] this lint" },
{ path = "tokio::sync::broadcast::Receiver::blocking_recv", reason = "block inside py.detach and #[expect] this lint" },
{ path = "tokio::runtime::Runtime::block_on", reason = "block inside py.detach and #[expect] this lint" },
{ path = "tokio::runtime::Handle::block_on", reason = "block inside py.detach and #[expect] this lint" },
{ path = "tokio::runtime::Runtime::shutdown_timeout", reason = "shut down inside py.detach and #[expect] this lint" },
]

CI treats Clippy warnings as errors:

SCYLLA_URI: "127.0.0.2:9042"
SCYLLA_URI2: "127.0.0.3:9042"
SCYLLA_URI3: "127.0.0.4:9042"
RUSTFLAGS: "-Dwarnings --cfg scylla_unstable"

.PHONY: clippy
clippy:
RUSTFLAGS="-Dwarnings --cfg scylla_unstable" cargo clippy --all-targets --features python-extension,openssl-010,rust-tests

PyO3 is pinned to 0.29, so suggested replacement methods must exist in that version:

[dependencies]
pyo3 = { version = "=0.29.0", features = ["bigdecimal", "num-bigint", "uuid", "chrono", "time", "py-clone"] }
tokio = { version = "1.42.0", features = ["full"] }

Current deadlock regression harness:

import subprocess
import sys
import pytest
from helpers.exit_scenarios import CONTACT_POINT
CHILD_TIMEOUT = 30.0
def _run_child(script: str, *args: str) -> subprocess.CompletedProcess[str]:
"""Run `script` in a fresh interpreter; a child that never exits is a deadlock."""
# Unbuffered, so a hung child's output is not lost in its buffers.
command = [sys.executable, "-u", "-c", script, *args]
with subprocess.Popen(command, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True) as child:
try:
stdout, stderr = child.communicate(timeout=CHILD_TIMEOUT)
except subprocess.TimeoutExpired:
# `subprocess.run` drops a timed-out child's output on POSIX; collect it here.
child.kill()
stdout, stderr = child.communicate()
pytest.fail(
f"the child never exited within {CHILD_TIMEOUT}s: it deadlocked\nstdout: {stdout!r}\nstderr: {stderr!r}"
)
return subprocess.CompletedProcess(command, child.returncode, stdout, stderr)
_AFTER_SUBINTERPRETER = """
try:
import _interpreters as interpreters
except ImportError:
import _xxsubinterpreters as interpreters
import socket
from scylla.session import SessionBuilder
interpreters.destroy(interpreters.create())
listener = socket.socket()
listener.bind(("127.0.0.1", 0))
listener.listen()
future = SessionBuilder().contact_points([listener.getsockname()]).connection_timeout(3600).connect()
try:
future.result(timeout=0)
except TimeoutError:
print("timed out")
"""
def test_attachment_check_survives_a_subinterpreter() -> None:
"""Waiting on a pending future after a subinterpreter was created.
Creating a subinterpreter makes `PyGILState_Check` return 1 on every thread, so an
attachment check built on it takes detached threads for attached ones: this panicked
with "still attached inside py.detach", or in release builds let a detached tokio
worker lock as if attached. It runs in a child, as the subinterpreter's effect
is process-wide and permanent.
"""
completed = _run_child(_AFTER_SUBINTERPRETER)
assert "panicked" not in completed.stderr, completed.stderr
assert completed.returncode == 0, f"the child exited with {completed.returncode}\n{completed.stderr}"
assert "timed out" in completed.stdout, completed.stderr
_CALLBACK_WHILE_PARKING = """
import asyncio
import socket
import threading
from scylla.session import SessionBuilder
parking_entered = threading.Event()
callback_entered = threading.Event()
class HookedLoop(asyncio.SelectorEventLoop):
armed = False
callback_seen = None
def create_future(self):
future = super().create_future()
if self.armed and self.callback_seen is None:
parking_entered.set()
self.callback_seen = callback_entered.wait(5)
return future
listener = socket.socket()
listener.bind(("127.0.0.1", 0))
listener.listen()
def fail_the_connection_once_parked():
conn, _ = listener.accept()
assert parking_entered.wait(5)
conn.close()
listener.close()
threading.Thread(target=fail_the_connection_once_parked, daemon=True).start()
async def main():
# Well under CHILD_TIMEOUT, so a helper that never closes the connection is not taken for a deadlock.
future = SessionBuilder().contact_points([listener.getsockname()]).connection_timeout(10).connect()
future.on_error(lambda _error: callback_entered.set())
loop.armed = True
try:
await future
except Exception:
pass
assert loop.callback_seen
loop = HookedLoop()
loop.run_until_complete(main())
print("done")
"""
def test_callback_completion_while_the_event_loop_parks() -> None:
"""Firing a callback and waking the awaiter while the event loop is parking on it.
Parking holds the waker's lock across `loop.create_future()`, which here waits with
the GIL released until the callback has fired. The completion thread then wakes the
waker: if it waited for that lock still holding the GIL, the loop could never return
from `create_future` to release it, and the whole process would hang.
"""
completed = _run_child(_CALLBACK_WHILE_PARKING)
assert completed.returncode == 0, f"the child exited with {completed.returncode}\n{completed.stderr}"
assert "done" in completed.stdout, completed.stderr
_CONCURRENT_NEXT = """
import sys
import threading
import time
from scylla.results import RowFactory
from scylla.session import SessionBuilder
factory_entered = threading.Event()
second_started = threading.Event()
errors = []
threading.excepthook = lambda args: errors.append(args.exc_value)
class BlockingFactory(RowFactory):
def build(self, columns):
factory_entered.set()
assert second_started.wait(5)
# Give the second thread time to reach the iterator's lock.
time.sleep(0.2)
return {column.column_name: column.value for column in columns}
session = SessionBuilder().contact_points([(sys.argv[1], int(sys.argv[2]))]).connect().result(timeout=10)
result = session.execute("SELECT release_version FROM system.local", factory=BlockingFactory()).result(timeout=10)
iterator = result.iter_current_page()
rows = []
def first():
rows.append(next(iterator))
def second():
second_started.set()
try:
next(iterator)
except StopIteration:
pass
t1 = threading.Thread(target=first)
t1.start()
assert factory_entered.wait(5)
t2 = threading.Thread(target=second)
t2.start()
t1.join(5)
t2.join(5)
if errors:
raise errors[0]
assert not t1.is_alive()
assert not t2.is_alive()
assert len(rows) == 1
print("done")
"""
# Best-effort: Python cannot tell when a thread is blocked on a Rust mutex, so the
# contention is near-certain thanks to the sleep, not guaranteed. Good enough for us.
@pytest.mark.requires_db
def test_concurrent_next_while_row_factory_waits() -> None:
"""Two threads calling `next` on one row iterator while the row factory waits.
Thread A holds the iterator's lock while its factory waits with the GIL released.
Thread B enters `next`, takes the GIL and waits for the lock. If B kept the GIL while
waiting, A could never return from the factory to release the lock, and the whole
process would hang. It runs in a child, so that hang is caught by a timeout.
"""
host, port = CONTACT_POINT
completed = _run_child(_CONCURRENT_NEXT, host, str(port))
assert completed.returncode == 0, f"the child exited with {completed.returncode}\n{completed.stderr}"
assert "done" in completed.stdout, completed.stderr

Split from #197:

#197

Acceptance criteria

  • Every existing use of a method listed in clippy.toml is reviewed and classified.
  • New uses fail Clippy without a narrowly scoped #[expect].
  • Each exemption explains the caller context and lock invariant.
  • Stale #[expect] attributes fail CI.
  • Clippy messages recommend APIs available in pinned dependencies.
  • Regression tests cover iterator contention, callback completion while parking, and Batcher contention.
  • Potential hangs run in subprocesses with external watchdogs.
  • Ordering is forced deterministically where practical.
  • GIL-enabled and free-threaded test configurations pass.
  • No public Python API changes.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    bugSomething isn't working

    Type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions