Skip to main content

rustpython_common/lock/
detaching.rs

1//! A reader-writer lock that lets a thread leave its interpreter before it
2//! blocks.
3//!
4//! [`RawDetachingRwLock`] carries the reasoning: what goes wrong when a thread
5//! waits for a lock while attached, why the waiter rather than the holder is
6//! the one that has to give way, and the rule that comes with fixing it.
7//!
8//! The wait itself is handed to a hook, because this crate cannot depend on the
9//! vm and so cannot detach a thread by itself. Whoever can installs it through
10//! [`set_blocking_wait_hook`]; until then, and on any thread that is not
11//! running an interpreter, a blocked acquire just blocks. Only the contended
12//! path reaches any of this — an acquire that takes the lock on its first try
13//! is the same atomic exchange it was.
14
15use super::RawRwLock;
16#[cfg(feature = "threading")]
17use core::cell::Cell;
18use lock_api::{
19    RawRwLock as RawRwLockTrait, RawRwLockDowngrade, RawRwLockUpgrade as RawRwLockUpgradeTrait,
20    RawRwLockUpgradeDowngrade,
21};
22#[cfg(feature = "threading")]
23use std::sync::OnceLock;
24
25/// Runs `wait` with the calling thread detached from its interpreter.
26#[cfg(feature = "threading")]
27pub type BlockingWaitHook = fn(wait: &dyn Fn());
28
29#[cfg(feature = "threading")]
30static BLOCKING_WAIT: OnceLock<BlockingWaitHook> = OnceLock::new();
31
32/// Install the hook that detaches a thread around a blocked lock acquire.
33///
34/// Later calls are ignored, so every interpreter in a process can call this
35/// during its own initialization.
36#[cfg(feature = "threading")]
37pub fn set_blocking_wait_hook(hook: BlockingWaitHook) {
38    let _ = BLOCKING_WAIT.set(hook);
39}
40
41#[cfg(feature = "threading")]
42std::thread_local! {
43    /// Set while this thread is inside the hook, so that a lock taken by the
44    /// hook itself — or by anything detaching and re-attaching runs — waits
45    /// plainly instead of recursing back into it.
46    static IN_HOOK: Cell<bool> = const { Cell::new(false) };
47}
48
49/// Clears [`IN_HOOK`] even if the hook unwinds.
50#[cfg(feature = "threading")]
51struct HookGuard;
52
53#[cfg(feature = "threading")]
54impl Drop for HookGuard {
55    fn drop(&mut self) {
56        let _ = IN_HOOK.try_with(|in_hook| in_hook.set(false));
57    }
58}
59
60#[cfg(all(feature = "threading", debug_assertions))]
61std::thread_local! {
62    /// Set on the one thread still running while the world is stopped.
63    static WORLD_STOPPED: Cell<bool> = const { Cell::new(false) };
64}
65
66/// Record whether this thread is the one running inside a stopped world.
67///
68/// The rule for opting a lock into detaching is that nothing reachable from a
69/// stop-the-world section takes it — a section that did could block on a lock
70/// only that same section can release. Not implementing `Traverse` states the
71/// rule to a collection; this states it to every other section, which is
72/// otherwise unchecked. Debug builds only; release builds track nothing and
73/// pay nothing.
74#[cfg(feature = "threading")]
75#[inline]
76pub fn set_world_stopped(stopped: bool) {
77    #[cfg(debug_assertions)]
78    let _ = WORLD_STOPPED.try_with(|flag| flag.set(stopped));
79    #[cfg(not(debug_assertions))]
80    let _ = stopped;
81}
82
83/// Panics if a stop-the-world section is taking one of these locks.
84#[cfg(all(feature = "threading", debug_assertions))]
85#[track_caller]
86fn assert_not_stopping_the_world() {
87    // `try_with` fails only once thread locals are being destroyed, which is
88    // not a point at which this thread is driving a stop.
89    let stopped = WORLD_STOPPED.try_with(Cell::get).unwrap_or(false);
90    assert!(
91        !stopped,
92        "a stop-the-world section took a detaching lock, which a parked thread \
93         may be holding and only this section can release"
94    );
95}
96
97#[cfg(not(all(feature = "threading", debug_assertions)))]
98#[inline(always)]
99fn assert_not_stopping_the_world() {}
100
101/// Block on `wait`, detached from this thread's interpreter if there is one.
102///
103/// Nothing spins on the way here. The lock underneath already spins before it
104/// parks, and skips that spin once a waiter has parked — the same condition
105/// `_PyMutex_LockTimed` spins under. A spin layered on top cannot read that
106/// condition, and would go on retrying a `try_lock` that reports failure for as
107/// long as a writer holds the writer bit, which it takes before it waits for
108/// readers to drain: a yield per retry for the whole of exactly the wait this
109/// exists to survive.
110#[cfg(feature = "threading")]
111#[cold]
112#[inline(never)]
113fn wait_detached(wait: impl Fn()) {
114    let Some(hook) = BLOCKING_WAIT.get() else {
115        wait();
116        return;
117    };
118    // `try_with` fails once the thread's locals are being destroyed, which is
119    // also a point at which there is no interpreter left to detach from.
120    let entered = IN_HOOK
121        .try_with(|in_hook| !in_hook.replace(true))
122        .unwrap_or(false);
123    if !entered {
124        wait();
125        return;
126    }
127    let _guard = HookGuard;
128    hook(&wait);
129}
130
131/// Without threads there is no interpreter to leave and nothing to stop.
132#[cfg(not(feature = "threading"))]
133#[inline]
134fn wait_detached(wait: impl Fn()) {
135    wait();
136}
137
138/// A reader-writer lock whose blocking acquires detach first, and which is the
139/// raw lock it wraps in every other respect.
140///
141/// Use through [`PyDetachingRwLock`](super::PyDetachingRwLock).
142///
143/// # Why this exists
144///
145/// Stopping the world means waiting until every other thread sits at
146/// SUSPENDED, and there are two ways a thread gets there:
147///
148/// - A DETACHED thread is not running interpreter code, so the requester moves
149///   it to SUSPENDED itself. The thread never finds out.
150/// - An ATTACHED thread can only suspend itself, at a safepoint — the check
151///   `check_signals` makes between bytecodes.
152///
153/// A thread blocked acquiring a lock runs no bytecode, so it reaches no
154/// safepoint. While ATTACHED it is a thread the world cannot stop for as long
155/// as it waits, and the requester waits without a bound.
156///
157/// On its own that is a pause. It becomes a deadlock as soon as the lock being
158/// waited for is held by a thread the same stop has already parked:
159///
160/// ```text
161/// A  holds the lock, blocks inside allow_threads  -> DETACHED
162/// B  requests a stop, and parks A                 -> A is SUSPENDED, holding the lock
163/// C  wants the same lock, and waits for it        -> ATTACHED, blocked
164///
165///   B waits for C to suspend           C reaches no safepoint
166///   C waits for A to release           A is parked
167///   A waits for B to start the world   B is still waiting for C
168/// ```
169///
170/// No thread in that cycle can break it, because none of them is running. It is
171/// not hypothetical: an `SSLSocket.read` against a peer that completed a
172/// handshake and then went quiet froze whole processes this way, the main
173/// thread included, so not even a Python-level timeout could fire.
174///
175/// The holder cannot be the one to give way. A lock is held across a blocking
176/// call precisely because that is what the call needs. So the waiter gives way
177/// instead: it leaves its interpreter for the duration of the wait, which is
178/// what a blocking call does anyway, and a waiter that has left is a waiter the
179/// requester can park. C detaches before it blocks, the stop completes, B
180/// finishes, A resumes and releases, and C takes the lock and attaches again.
181///
182/// # Only for locks a stop-the-world section never takes
183///
184/// The wait acquires the lock while detached, so the thread comes back holding
185/// it, and re-attaching is a point at which a stop-the-world in flight will
186/// park the thread. It is therefore parked *holding the lock*. Everything that
187/// stops the world must be able to finish without that lock: if a collection
188/// were to take it, the collection would block on a thread only the collection
189/// can release, and neither would move again.
190///
191/// So this is opt-in per lock, and the rule for opting in is that nothing
192/// reachable from a stop-the-world section takes the same lock. An object whose
193/// payload holds no references — nothing for the collector to traverse into —
194/// satisfies that; most do not.
195///
196/// Not implementing the vm's `Traverse` for this lock enforces part of that: a
197/// payload holding one cannot derive `Traverse`, so it cannot become something
198/// a collection walks into. Only that part. A collection is not the only thing
199/// that stops the world — dumping tracebacks, enumerating thread frames and
200/// forking all do — and nothing checks what those reach. For them the rule is
201/// still a convention.
202#[repr(transparent)]
203pub struct RawDetachingRwLock(RawRwLock);
204
205// SAFETY: every method forwards to the wrapped raw lock, which upholds the
206// contract; the blocking acquires only add a wait that ends with the same lock
207// acquired.
208unsafe impl RawRwLockTrait for RawDetachingRwLock {
209    #[allow(
210        clippy::declare_interior_mutable_const,
211        reason = "raw lock initializer, as in the type it wraps"
212    )]
213    const INIT: Self = Self(<RawRwLock as RawRwLockTrait>::INIT);
214
215    type GuardMarker = <RawRwLock as RawRwLockTrait>::GuardMarker;
216
217    #[inline]
218    fn lock_shared(&self) {
219        assert_not_stopping_the_world();
220        if !self.0.try_lock_shared() {
221            wait_detached(|| self.0.lock_shared());
222        }
223    }
224
225    #[inline]
226    fn try_lock_shared(&self) -> bool {
227        self.0.try_lock_shared()
228    }
229
230    #[inline]
231    unsafe fn unlock_shared(&self) {
232        unsafe { self.0.unlock_shared() }
233    }
234
235    #[inline]
236    fn lock_exclusive(&self) {
237        assert_not_stopping_the_world();
238        if !self.0.try_lock_exclusive() {
239            wait_detached(|| self.0.lock_exclusive());
240        }
241    }
242
243    #[inline]
244    fn try_lock_exclusive(&self) -> bool {
245        self.0.try_lock_exclusive()
246    }
247
248    #[inline]
249    unsafe fn unlock_exclusive(&self) {
250        unsafe { self.0.unlock_exclusive() }
251    }
252
253    #[inline]
254    fn is_locked(&self) -> bool {
255        self.0.is_locked()
256    }
257
258    #[inline]
259    fn is_locked_exclusive(&self) -> bool {
260        self.0.is_locked_exclusive()
261    }
262}
263
264// SAFETY: forwards to the wrapped raw lock.
265unsafe impl RawRwLockDowngrade for RawDetachingRwLock {
266    #[inline]
267    unsafe fn downgrade(&self) {
268        unsafe { self.0.downgrade() }
269    }
270}
271
272// SAFETY: forwards to the wrapped raw lock; `lock_upgradable` only adds a wait
273// that ends with the same lock acquired.
274//
275// `lock_upgradable` detaches for the same reason `lock_shared` does: it starts
276// from holding nothing, so the wait cannot park a thread that holds the lock.
277// `upgrade` does not, because it runs with the upgradable lock already held.
278unsafe impl RawRwLockUpgradeTrait for RawDetachingRwLock {
279    #[inline]
280    fn lock_upgradable(&self) {
281        assert_not_stopping_the_world();
282        if !self.0.try_lock_upgradable() {
283            wait_detached(|| self.0.lock_upgradable());
284        }
285    }
286
287    #[inline]
288    fn try_lock_upgradable(&self) -> bool {
289        self.0.try_lock_upgradable()
290    }
291
292    #[inline]
293    unsafe fn unlock_upgradable(&self) {
294        unsafe { self.0.unlock_upgradable() }
295    }
296
297    #[inline]
298    unsafe fn upgrade(&self) {
299        // SAFETY: the caller holds the upgradable lock, as `upgrade` requires.
300        unsafe { self.0.upgrade() }
301    }
302
303    #[inline]
304    unsafe fn try_upgrade(&self) -> bool {
305        unsafe { self.0.try_upgrade() }
306    }
307}
308
309// SAFETY: forwards to the wrapped raw lock.
310unsafe impl RawRwLockUpgradeDowngrade for RawDetachingRwLock {
311    #[inline]
312    unsafe fn downgrade_upgradable(&self) {
313        unsafe { self.0.downgrade_upgradable() }
314    }
315
316    #[inline]
317    unsafe fn downgrade_to_upgradable(&self) {
318        unsafe { self.0.downgrade_to_upgradable() }
319    }
320}
321
322// `RawRwLockRecursive` is deliberately not implemented, so that `read_recursive`
323// does not exist on these locks. It is the one blocking acquire that cannot
324// detach: a recursive read may be the re-entrant take of a lock this thread
325// already holds, and detaching there parks a thread *holding* the lock, which is
326// the deadlock this type exists to avoid. Leaving it implemented but attached
327// would instead leave an acquire that stalls stop-the-world, so neither form of
328// it belongs here.
329
330#[cfg(test)]
331mod tests {
332    #[cfg(all(feature = "threading", debug_assertions))]
333    use super::set_world_stopped;
334    #[cfg(all(feature = "threading", debug_assertions))]
335    use crate::lock::PyDetachingRwLock;
336
337    /// The opt-in rule holds for every stop-the-world section, not only the
338    /// collector that not implementing `Traverse` speaks to.
339    #[cfg(all(feature = "threading", debug_assertions))]
340    #[test]
341    fn taking_one_while_stopping_the_world_is_caught() {
342        let lock = PyDetachingRwLock::new(());
343
344        // Ordinary use, for contrast.
345        drop(lock.write());
346
347        // The panic below is the expected result, so it prints where an
348        // unexpected one would. Silencing it would mean replacing the panic hook,
349        // which is process-wide and would swallow the output of whatever else the
350        // test binary is running at the same time.
351        set_world_stopped(true);
352        let taken = std::panic::catch_unwind(core::panic::AssertUnwindSafe(|| {
353            let _guard = lock.read();
354        }));
355        set_world_stopped(false);
356
357        assert!(
358            taken.is_err(),
359            "a stop-the-world section took a detaching lock and nothing complained"
360        );
361
362        // The flag is per-thread and back to clear, so the lock still works.
363        drop(lock.write());
364    }
365}