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}