rustpython_vm/vm/thread.rs
1#[cfg(all(not(unix), feature = "threading"))]
2use super::FramePtr;
3#[cfg(feature = "threading")]
4use crate::PyObjectRef;
5#[cfg(feature = "threading")]
6use crate::builtins::PyBaseExceptionRef;
7#[cfg(feature = "threading")]
8use alloc::sync::Arc;
9#[cfg(feature = "threading")]
10use rustpython_common::lock::PyMutex;
11
12use crate::frame::InterpreterFrame;
13#[cfg(feature = "threading")]
14use crate::vm::PyGlobalState;
15use crate::{AsObject, PyObject, VirtualMachine};
16#[cfg(all(unix, feature = "threading"))]
17use crate::{Py, frame::FrameObject};
18#[cfg(all(unix, feature = "threading"))]
19use core::sync::atomic::AtomicPtr;
20use core::{
21 cell::{Cell, RefCell},
22 ptr::NonNull,
23 sync::atomic::{AtomicUsize, Ordering},
24};
25use itertools::Itertools;
26#[cfg(feature = "threading")]
27use std::collections::HashMap;
28use std::thread_local;
29
30/// Thread states for stop-the-world support (`_Py_THREAD_*`).
31///
32/// DETACHED: not executing Python bytecode (in native code, or idle)
33/// ATTACHED: actively executing Python bytecode
34/// SUSPENDED: parked by a stop-the-world request
35/// SHUTTING_DOWN: interpreter is finalizing; the OS thread must hang
36/// (`_PyThreadState_HangThread`) and must not look done to `_ThreadHandle`.
37#[cfg(feature = "threading")]
38#[repr(i32)]
39#[derive(Copy, Clone, Debug, PartialEq, Eq)]
40pub enum ThreadState {
41 Detached = 0,
42 Attached = 1,
43 Suspended = 2,
44 ShuttingDown = 3,
45}
46
47#[cfg(feature = "threading")]
48impl ThreadState {
49 #[must_use]
50 pub const fn from_i32(v: i32) -> Option<Self> {
51 match v {
52 0 => Some(Self::Detached),
53 1 => Some(Self::Attached),
54 2 => Some(Self::Suspended),
55 3 => Some(Self::ShuttingDown),
56 _ => None,
57 }
58 }
59}
60
61/// Per-thread shared state for sys._current_frames() and sys._current_exceptions().
62/// The exception field uses atomic operations for lock-free cross-thread reads.
63#[cfg(feature = "threading")]
64pub struct ThreadSlot {
65 /// Top of the owning thread's Python call stack, published for
66 /// cross-thread readers (`sys._current_frames`, cross-thread `f_back`).
67 /// The rest of the stack is reachable via each frame's `previous` pointer.
68 /// Written lock-free on the hot push/pop path with relaxed ordering; every
69 /// cross-thread read runs under stop-the-world, which parks the owning
70 /// thread at a safepoint and supplies the happens-before edge, so the
71 /// pointer and the frames it reaches are quiescent and alive at read time.
72 #[cfg(unix)]
73 pub top_frame: AtomicPtr<Py<FrameObject>>,
74 /// Raw InterpreterFrame pointer, published alongside top_frame so
75 /// cross-thread readers (sys._current_frames) can materialize
76 /// stack-allocated frames that have no FrameObject.
77 pub top_iframe: AtomicUsize,
78 /// Raw frame pointers, valid while the owning thread's call stack is active.
79 /// Readers must hold the Mutex and convert to FrameObjectRef inside the lock.
80 /// Stands in for `top_frame` where that field is not built, so a reader
81 /// that finds no `top_iframe` still has the frames to answer from.
82 #[cfg(not(unix))]
83 pub frames: parking_lot::Mutex<Vec<FramePtr>>,
84 pub exception: crate::PyAtomicRef<Option<crate::exceptions::types::PyBaseException>>,
85 /// `tstate->c_traceobj`. Cross-thread source of truth for sys.gettrace /
86 /// sys._settraceallthreads.
87 pub trace_func: PyMutex<PyObjectRef>,
88 /// `tstate->c_profileobj`. Cross-thread source of truth for sys.getprofile /
89 /// sys._setprofileallthreads.
90 pub profile_func: PyMutex<PyObjectRef>,
91 /// Thread state for stop-the-world: DETACHED / ATTACHED / SUSPENDED / SHUTTING_DOWN
92 pub state: core::sync::atomic::AtomicI32,
93 /// Per-thread stop request bit (eval breaker equivalent).
94 pub stop_requested: core::sync::atomic::AtomicBool,
95 /// Handle for waking this thread from park in stop-the-world paths.
96 pub thread: std::thread::Thread,
97 /// QSBR state for deferred memory reclamation.
98 pub(crate) qsbr: Arc<crate::object::qsbr::QsbrSlot>,
99}
100
101#[cfg(feature = "threading")]
102pub type CurrentFrameSlot = Arc<ThreadSlot>;
103
104/// Coalesced per-thread frame-publishing state, touched on every
105/// `enter_iframe`/`exit_iframe`. Bundling `current_frame` together with
106/// the cached `top_frame`/`top_iframe` slot pointers means
107/// `set_current_frame` needs a single `thread_local!.with()` call
108/// (one `_tlv_get_addr` on macOS) instead of two or three separate
109/// ones — each `.with()` on a distinct `thread_local!` is its own TLS
110/// lookup even though the bodies are just a cached-pointer store.
111struct FrameSlotCache {
112 /// Current top frame for signal-safe traceback walking.
113 /// Stores a `*const InterpreterFrame` as `usize`.
114 /// Read by faulthandler's signal handler to dump tracebacks without
115 /// accessing RefCell or locks. Uses AtomicUsize for async-signal-safety.
116 current_frame: AtomicUsize,
117 /// Cached pointer to this thread's `ThreadSlot::top_frame`, so the hot
118 /// push/pop path can publish the top frame with a single relaxed store and
119 /// no `CURRENT_THREAD_SLOT` RefCell borrow. Null until the slot is
120 /// initialized; the `Arc<ThreadSlot>` in `CURRENT_THREAD_SLOT` keeps the
121 /// pointee alive until `cleanup_current_thread_frames` clears this.
122 #[cfg(all(unix, feature = "threading"))]
123 top_frame: Cell<*const AtomicPtr<Py<FrameObject>>>,
124 /// Cached pointer to this thread's `ThreadSlot::top_iframe` for the hot
125 /// light-frame push/pop path. The slot's Arc keeps the pointee alive.
126 #[cfg(feature = "threading")]
127 top_iframe: Cell<*const AtomicUsize>,
128}
129
130thread_local! {
131 pub(super) static VM_STACK: RefCell<Vec<NonNull<VirtualMachine>>> = Vec::with_capacity(1).into();
132
133 /// Thread state created through the GILState-style C API.
134 ///
135 /// This is separate from the current VM stack: it only means "attached now",
136 /// while this owns the per-thread VM that may be detached and re-attached.
137 /// Despite the historical CPython "GILState" name, this does not model a
138 /// GIL; it stores the VM used by that compatibility API.
139 ///
140 /// The Box keeps the VM address stable while VM_STACK holds a raw pointer to it.
141 /// This matters when release_current_thread() moves the owner out of TLS and
142 /// drops it while the VM is still current, so object destructors can still find
143 /// their VM.
144 #[cfg(feature = "threading")]
145 static GILSTATE_VM: RefCell<Option<Box<ThreadedVirtualMachine>>> = const { RefCell::new(None) };
146
147 pub(crate) static COROUTINE_ORIGIN_TRACKING_DEPTH: Cell<u32> = const { Cell::new(0) };
148
149 /// Per-interpreter thread slots for this OS thread (PEP 734 multi-interpreter).
150 ///
151 /// CPython keeps a `PyThreadState` per (thread, interpreter) pair. RustPython
152 /// mirrors that: each interpreter's `PyGlobalState.thread_frames` gets its own
153 /// [`ThreadSlot`] for this OS thread. `CURRENT_THREAD_SLOT` always points at
154 /// the slot for the currently entered interpreter.
155 #[cfg(feature = "threading")]
156 static INTERP_THREAD_SLOTS: RefCell<HashMap<i64, CurrentFrameSlot>> =
157 RefCell::new(HashMap::new());
158
159 /// Current thread's slot for the currently entered interpreter.
160 #[cfg(feature = "threading")]
161 static CURRENT_THREAD_SLOT: RefCell<Option<CurrentFrameSlot>> = const { RefCell::new(None) };
162
163 pub(crate) static FRAME_SLOT_CACHE: FrameSlotCache = const {
164 FrameSlotCache {
165 current_frame: AtomicUsize::new(0),
166 #[cfg(all(unix, feature = "threading"))]
167 top_frame: Cell::new(core::ptr::null()),
168 #[cfg(feature = "threading")]
169 top_iframe: Cell::new(core::ptr::null()),
170 }
171 };
172
173 /// Cached pointer to this thread's `ThreadSlot::stop_requested`, for the
174 /// safepoint the dispatch loop takes once per instruction. Reading it
175 /// through `CURRENT_THREAD_SLOT` costs a `RefCell` borrow — two stores to
176 /// thread-local memory — where this costs one relaxed load. The slot's Arc
177 /// keeps the pointee alive, as with the frame pointers above.
178 #[cfg(feature = "threading")]
179 static CURRENT_STOP_REQUESTED: Cell<*const core::sync::atomic::AtomicBool> =
180 const { Cell::new(core::ptr::null()) };
181
182}
183
184#[must_use]
185pub fn current_vm_is_set() -> bool {
186 VM_STACK.with(|vms| !vms.borrow().is_empty())
187}
188
189pub fn with_current_vm<R>(f: impl FnOnce(&VirtualMachine) -> R) -> R {
190 VM_STACK.with(|vms| {
191 let vm = vms
192 .borrow()
193 .last()
194 .copied()
195 .expect("call with_current_vm() but no current VM is attached");
196 // SAFETY: entries in VM_STACK either borrow a VM for the dynamic
197 // scope of a set_current_vm()/enter_vm() call or point at GILSTATE_VM.
198 f(unsafe { vm.as_ref() })
199 })
200}
201
202fn set_current_vm<R>(vm: &VirtualMachine, f: impl FnOnce() -> R) -> R {
203 // Attach to this VM's interpreter, detaching the enclosing one if this is a
204 // switch between interpreters on the same OS thread.
205 #[cfg(feature = "threading")]
206 let switched = begin_interpreter_section(vm);
207
208 VM_STACK.with(|vms| {
209 vms.borrow_mut().push(vm.into());
210 scopeguard::defer! {
211 vms.borrow_mut().pop();
212 #[cfg(feature = "threading")]
213 end_interpreter_section(switched);
214 }
215 f()
216 })
217}
218
219/// Pointer to the GC state of the interpreter running on this thread.
220///
221/// The pointee belongs to the `PyGlobalState` of the VM on top of `VM_STACK`,
222/// which is borrowed for the whole `set_current_vm` scope — so the pointer stays
223/// valid as long as the caller remains inside that scope.
224pub(crate) fn current_gc_state() -> Option<NonNull<crate::gc_state::GcInterpreterState>> {
225 // Reached from every tracked allocation, including ones a thread-local
226 // destructor makes while the VM stack is being torn down, so neither a
227 // destroyed key nor an outstanding borrow may panic here.
228 VM_STACK
229 .try_with(|vms| {
230 let vm = vms.try_borrow().ok()?.last().copied()?;
231 // SAFETY: entries in VM_STACK either borrow a VM for the dynamic
232 // scope of a set_current_vm()/enter_vm() call or point at GILSTATE_VM.
233 Some(NonNull::from(&unsafe { vm.as_ref() }.state.gc))
234 })
235 .ok()
236 .flatten()
237}
238
239pub fn try_with_current_vm<R>(f: impl FnOnce(&VirtualMachine) -> R) -> Option<R> {
240 VM_STACK.with(|vms| {
241 let vm = vms.borrow().last().copied()?;
242 // SAFETY: entries in VM_STACK either borrow a VM for the dynamic
243 // scope of a set_current_vm()/enter_vm() call or point at GILSTATE_VM.
244 Some(f(unsafe { vm.as_ref() }))
245 })
246}
247
248pub fn enter_vm<R>(vm: &VirtualMachine, f: impl FnOnce() -> R) -> R {
249 // Attach/detach is handled by `set_current_vm`, which pairs it with the
250 // VM_STACK push so that switching interpreters mid-stack stays consistent.
251 set_current_vm(vm, f)
252}
253
254/// RAII counterpart to `enter_vm`, for code that runs Python bytecode across
255/// several statements interspersed with `&mut VirtualMachine` calls
256/// (`VirtualMachine::initialize`), where a single closure-based `enter_vm`
257/// scope cannot be expressed because the borrow checker won't let a closure
258/// hold `&mut VirtualMachine` at the same time `enter_vm` reborrows it as
259/// `&VirtualMachine`. Construction only needs a transient `&VirtualMachine`
260/// borrow, so it can be dropped before subsequent `&mut` use.
261///
262/// Without this, code that runs Python bytecode before any `enter_vm` scope
263/// exists would leave the thread not ATTACHED, making lock-free type cache
264/// reads unsound.
265#[must_use]
266pub(crate) struct VmBootstrapGuard {
267 #[cfg(feature = "threading")]
268 switched: bool,
269}
270
271impl VmBootstrapGuard {
272 pub(crate) fn new(vm: &VirtualMachine) -> Self {
273 #[cfg(feature = "threading")]
274 let switched = begin_interpreter_section(vm);
275
276 VM_STACK.with(|vms| vms.borrow_mut().push(vm.into()));
277
278 Self {
279 #[cfg(feature = "threading")]
280 switched,
281 }
282 }
283}
284
285impl Drop for VmBootstrapGuard {
286 fn drop(&mut self) {
287 VM_STACK.with(|vms| {
288 vms.borrow_mut().pop();
289 });
290
291 #[cfg(feature = "threading")]
292 end_interpreter_section(self.switched);
293 }
294}
295
296#[cfg(feature = "threading")]
297#[derive(Clone, Copy, Debug, Eq, PartialEq)]
298pub enum CurrentVmAttachState {
299 AlreadyAttached,
300 Attached,
301}
302
303/// State preserved while the current native thread is detached from its VM.
304#[cfg(feature = "threading")]
305pub struct SavedThreadState {
306 vm_stack: Vec<NonNull<VirtualMachine>>,
307 gilstate_vm: Option<Box<ThreadedVirtualMachine>>,
308}
309
310/// Detach the current native thread and preserve its VM context for restoration.
311#[cfg(feature = "threading")]
312#[must_use = "the saved thread state must be restored"]
313pub fn save_current_thread() -> SavedThreadState {
314 let vm_stack = VM_STACK.with(|vms| core::mem::take(&mut *vms.borrow_mut()));
315 assert!(
316 !vm_stack.is_empty(),
317 "save_current_thread() called without an attached VM"
318 );
319 let gilstate_vm = GILSTATE_VM.with(|gilstate_vm| gilstate_vm.borrow_mut().take());
320 detach_thread();
321 SavedThreadState {
322 vm_stack,
323 gilstate_vm,
324 }
325}
326
327/// Restore a VM context previously returned by [`save_current_thread`].
328#[cfg(feature = "threading")]
329pub fn restore_current_thread(state: SavedThreadState) {
330 assert!(
331 !current_vm_is_set(),
332 "restore_current_thread() called with an attached VM"
333 );
334 let SavedThreadState {
335 vm_stack,
336 gilstate_vm,
337 } = state;
338 let vm = vm_stack
339 .last()
340 .copied()
341 .expect("saved thread state has no VM");
342
343 GILSTATE_VM.with(|current| {
344 let mut current = current.borrow_mut();
345 assert!(
346 current.is_none(),
347 "restore_current_thread() called with a GILState VM"
348 );
349 *current = gilstate_vm;
350 });
351
352 // SAFETY: borrowed VMs remain alive for the dynamic save/restore scope,
353 // while an owned GILState VM was restored above before this dereference.
354 let vm = unsafe { vm.as_ref() };
355 // Point CURRENT_THREAD_SLOT at the restored interpreter before attach.
356 // After subinterpreter bootstrap, CURRENT may still refer to the temporary
357 // subinterpreter slot (DETACHED); attaching that would leave the parent
358 // slot detached and later confuse outermost detach.
359 init_thread_slot_if_needed(vm);
360 attach_thread(vm);
361 VM_STACK.with(|vms| *vms.borrow_mut() = vm_stack);
362}
363
364/// Attach the current native thread to a RustPython VM until
365/// `release_current_thread()` is called.
366#[cfg(feature = "threading")]
367pub fn attach_current_thread(
368 make_vm: impl FnOnce() -> ThreadedVirtualMachine,
369) -> CurrentVmAttachState {
370 if current_vm_is_set() {
371 return CurrentVmAttachState::AlreadyAttached;
372 }
373
374 GILSTATE_VM.with(|gilstate_vm| {
375 let mut gilstate_vm = gilstate_vm.borrow_mut();
376 let threaded_vm = gilstate_vm.get_or_insert_with(|| Box::new(make_vm()));
377 let vm = &threaded_vm.vm;
378
379 vm.c_stack_soft_limit
380 .set(VirtualMachine::calculate_c_stack_soft_limit());
381
382 init_thread_slot_if_needed(vm);
383
384 attach_thread(vm);
385
386 VM_STACK.with(|vms| {
387 debug_assert!(vms.borrow().is_empty());
388 vms.borrow_mut().push(vm.into());
389 });
390 });
391
392 CurrentVmAttachState::Attached
393}
394
395#[cfg(feature = "threading")]
396pub fn release_current_thread(state: CurrentVmAttachState) {
397 if state == CurrentVmAttachState::AlreadyAttached {
398 return;
399 }
400
401 let gilstate_vm = GILSTATE_VM.with(|gilstate_vm| gilstate_vm.borrow_mut().take());
402 drop(gilstate_vm);
403
404 VM_STACK.with(|vms| {
405 vms.borrow_mut()
406 .pop()
407 .expect("release_current_thread() called without an attached VM");
408 });
409
410 detach_thread();
411}
412
413/// Ensure this OS thread has a [`ThreadSlot`] registered with `vm`'s interpreter
414/// and make it the current slot.
415///
416/// Called automatically by `enter_vm()` / `VmBootstrapGuard` whenever a VM
417/// becomes current. Switching between interpreters on the same OS thread swaps
418/// `CURRENT_THREAD_SLOT` to that interpreter's slot (creating one if needed).
419#[cfg(feature = "threading")]
420fn init_thread_slot_if_needed(vm: &VirtualMachine) {
421 let slot = ensure_thread_slot(vm);
422 set_current_thread_slot(slot);
423}
424
425/// Look up (creating if needed) this thread's [`ThreadSlot`] for `vm`'s
426/// interpreter, without making it the current slot.
427#[cfg(feature = "threading")]
428fn ensure_thread_slot(vm: &VirtualMachine) -> CurrentFrameSlot {
429 let interp_id = vm.state.interpreter_id;
430 INTERP_THREAD_SLOTS.with(|slots| {
431 let mut slots = slots.borrow_mut();
432 if let Some(existing) = slots.get(&interp_id) {
433 return existing.clone();
434 }
435
436 let thread_id = crate::stdlib::_thread::get_ident();
437 let mut registry = vm.state.thread_frames.lock();
438 let new_slot = Arc::new(ThreadSlot {
439 #[cfg(unix)]
440 top_frame: AtomicPtr::new(core::ptr::null_mut()),
441 top_iframe: AtomicUsize::new(0),
442 #[cfg(not(unix))]
443 frames: parking_lot::Mutex::new(Vec::new()),
444 exception: crate::PyAtomicRef::from(None::<PyBaseExceptionRef>),
445 trace_func: PyMutex::new(
446 vm.state
447 .global_trace_func
448 .lock()
449 .clone()
450 .unwrap_or_else(|| vm.ctx.none()),
451 ),
452 profile_func: PyMutex::new(
453 vm.state
454 .global_profile_func
455 .lock()
456 .clone()
457 .unwrap_or_else(|| vm.ctx.none()),
458 ),
459 state: core::sync::atomic::AtomicI32::new(
460 if vm.state.stop_the_world.requested.load(Ordering::Acquire) {
461 // Match init_threadstate(): new thread-state starts
462 // suspended while stop-the-world is active.
463 ThreadState::Suspended as i32
464 } else {
465 ThreadState::Detached as i32
466 },
467 ),
468 stop_requested: core::sync::atomic::AtomicBool::new(false),
469 thread: std::thread::current(),
470 qsbr: crate::object::qsbr::QSBR.register(),
471 });
472 registry.insert(thread_id, new_slot.clone());
473 drop(registry);
474 slots.insert(interp_id, new_slot.clone());
475 new_slot
476 })
477}
478
479/// The current thread's `ThreadSlot` for the entered interpreter, if any.
480#[cfg(feature = "threading")]
481#[must_use]
482pub fn current_thread_slot() -> Option<CurrentFrameSlot> {
483 CURRENT_THREAD_SLOT.with(|slot| slot.borrow().clone())
484}
485
486/// Make `slot` the current thread slot (and the cached top-frame pointer).
487#[cfg(feature = "threading")]
488fn set_current_thread_slot(slot: CurrentFrameSlot) {
489 FRAME_SLOT_CACHE.with(|cache| {
490 #[cfg(unix)]
491 cache.top_frame.set(&slot.top_frame);
492 cache.top_iframe.set(&slot.top_iframe);
493 });
494 CURRENT_STOP_REQUESTED.with(|c| c.set(&slot.stop_requested));
495 CURRENT_THREAD_SLOT.with(|current| {
496 *current.borrow_mut() = Some(slot);
497 });
498}
499
500/// Whether the current thread slot is ATTACHED.
501#[cfg(feature = "threading")]
502fn current_slot_is_attached() -> bool {
503 CURRENT_THREAD_SLOT.with(|slot| {
504 slot.borrow()
505 .as_ref()
506 .is_some_and(|s| s.state.load(Ordering::Acquire) == ThreadState::Attached as i32)
507 })
508}
509
510/// Attach this thread to `vm`'s interpreter for the duration of a section,
511/// detaching whichever interpreter it was attached to (≈ `_PyThreadState_Swap`).
512///
513/// A thread must never be ATTACHED to two interpreters at once: stop-the-world
514/// treats an ATTACHED slot as "running this interpreter's bytecode" and a
515/// DETACHED slot as parkable without cooperation, so running interpreter B's
516/// code while B's slot is DETACHED would let a collector conclude B is stopped
517/// while this thread keeps mutating the (process-global) object graph.
518///
519/// Returns whether the attachment changed, i.e. whether the matching
520/// [`end_interpreter_section`] must undo it.
521#[cfg(feature = "threading")]
522fn begin_interpreter_section(vm: &VirtualMachine) -> bool {
523 let target = ensure_thread_slot(vm);
524 let already_current = CURRENT_THREAD_SLOT.with(|slot| {
525 slot.borrow()
526 .as_ref()
527 .is_some_and(|s| Arc::ptr_eq(s, &target))
528 });
529 if already_current && current_slot_is_attached() {
530 // Nested section in the same interpreter: already attached.
531 return false;
532 }
533 if !already_current && current_slot_is_attached() {
534 detach_thread();
535 }
536 set_current_thread_slot(target);
537 attach_thread(vm);
538 true
539}
540
541/// Undo [`begin_interpreter_section`]: detach this interpreter and re-attach the
542/// enclosing one, if any. Call after the VM has been popped from `VM_STACK`.
543#[cfg(feature = "threading")]
544fn end_interpreter_section(switched: bool) {
545 if !switched {
546 return;
547 }
548 if current_slot_is_attached() {
549 detach_thread();
550 }
551 // The enclosing section, if any, is the VM now on top of the stack.
552 if let Some(vm_ptr) = VM_STACK.with(|vms| vms.borrow().last().copied()) {
553 // SAFETY: entries on VM_STACK are valid for their enter/set_current_vm scope.
554 let vm = unsafe { vm_ptr.as_ref() };
555 set_current_thread_slot(ensure_thread_slot(vm));
556 attach_thread(vm);
557 }
558}
559
560/// Transition DETACHED → ATTACHED. Blocks if the thread was SUSPENDED by
561/// a stop-the-world request (like `_PyThreadState_Attach` + `tstate_wait_attach`).
562#[cfg(feature = "threading")]
563fn wait_while_suspended(slot: &ThreadSlot) -> u64 {
564 let mut wait_yields = 0u64;
565 while slot.state.load(Ordering::Acquire) == ThreadState::Suspended as i32 {
566 wait_yields = wait_yields.saturating_add(1);
567 std::thread::park();
568 }
569 wait_yields
570}
571
572/// `PyThread_hang_thread`: park this OS thread forever.
573#[cfg(feature = "threading")]
574fn hang_thread() -> ! {
575 loop {
576 std::thread::park();
577 }
578}
579
580/// `_PyThreadState_HangThread`: this thread may no longer run Python.
581///
582/// Mark the slot shutting-down (so later stop-the-world requests do not wait
583/// for it) and never return. The matching `_ThreadHandle` stays not-done, so
584/// `Thread.is_alive()` remains true for a daemon forced off during finalize.
585#[cfg(feature = "threading")]
586pub fn hang_current_thread(state: &PyGlobalState) -> ! {
587 CURRENT_THREAD_SLOT.with(|slot| {
588 if let Some(s) = slot.borrow().as_ref() {
589 let prev = s
590 .state
591 .swap(ThreadState::ShuttingDown as i32, Ordering::AcqRel);
592 if prev == ThreadState::Attached as i32 {
593 crate::object::qsbr::QSBR.offline(&s.qsbr);
594 }
595 s.stop_requested.store(false, Ordering::Release);
596 }
597 });
598 if state.stop_the_world.requested.load(Ordering::Acquire) {
599 state.stop_the_world.notify_thread_gone();
600 }
601 hang_thread();
602}
603
604/// `_PyThreadState_SetShuttingDown` on every non-current thread.
605///
606/// Call while the world is stopped. Wake parked threads so they observe
607/// `SHUTTING_DOWN` and hang on the next attach, instead of resuming Python.
608#[cfg(feature = "threading")]
609pub fn set_other_threads_shutting_down(state: &PyGlobalState) {
610 let current = crate::stdlib::_thread::get_ident();
611 let registry = state.thread_frames.lock();
612
613 #[expect(
614 clippy::iter_over_hash_type,
615 reason = "Iteration order doesn't matter here"
616 )]
617 for (&id, slot) in registry.iter() {
618 if id == current {
619 continue;
620 }
621 slot.stop_requested.store(false, Ordering::Release);
622 slot.state
623 .store(ThreadState::ShuttingDown as i32, Ordering::Release);
624 slot.thread.unpark();
625 }
626}
627
628#[cfg(feature = "threading")]
629fn attach_thread(vm: &VirtualMachine) {
630 CURRENT_THREAD_SLOT.with(|slot| {
631 if let Some(s) = slot.borrow().as_ref() {
632 super::stw_trace(format_args!("attach begin"));
633 loop {
634 match s.state.compare_exchange(
635 ThreadState::Detached as i32,
636 ThreadState::Attached as i32,
637 Ordering::AcqRel,
638 Ordering::Relaxed,
639 ) {
640 Ok(_) => {
641 crate::object::qsbr::QSBR.online(&s.qsbr);
642 super::stw_trace(format_args!("attach DETACHED->ATTACHED"));
643 break;
644 }
645 Err(state) => match ThreadState::from_i32(state) {
646 Some(ThreadState::Suspended) => {
647 // Parked by stop-the-world — wait until released to DETACHED
648 super::stw_trace(format_args!("attach wait-suspended"));
649 let wait_yields = wait_while_suspended(s);
650 vm.state.stop_the_world.add_attach_wait_yields(wait_yields);
651 // Retry CAS
652 }
653 Some(ThreadState::ShuttingDown) => {
654 super::stw_trace(format_args!("attach hang shutting-down"));
655 hang_thread();
656 }
657 _ => {
658 debug_assert!(false, "unexpected thread state in attach: {state}");
659 break;
660 }
661 },
662 }
663 }
664 }
665 });
666 // A stop-the-world may have been requested while this thread was detached.
667 // Honoring it here (rather than only at the next bytecode safepoint) keeps
668 // a thread doing rapid allow_threads calls from re-attaching and running
669 // past the requester forever, which would stall stop-the-world. Done
670 // outside the CURRENT_THREAD_SLOT borrow above because suspend re-borrows
671 // it. Safe against a concurrent start_the_world: suspend_if_needed decides
672 // whether to park under the registry lock, so it never parks after the
673 // request has been withdrawn.
674 suspend_if_needed(&vm.state);
675}
676
677/// Transition ATTACHED → DETACHED (like `_PyThreadState_Detach`).
678#[cfg(feature = "threading")]
679fn detach_thread() {
680 CURRENT_THREAD_SLOT.with(|slot| {
681 if let Some(s) = slot.borrow().as_ref() {
682 match s.state.compare_exchange(
683 ThreadState::Attached as i32,
684 ThreadState::Detached as i32,
685 Ordering::AcqRel,
686 Ordering::Acquire,
687 ) {
688 Ok(_) => {
689 crate::object::qsbr::QSBR.offline(&s.qsbr);
690 }
691 Err(state) => {
692 debug_assert!(
693 matches!(ThreadState::from_i32(state), Some(ThreadState::Detached)),
694 "unexpected thread state in detach: {state}"
695 );
696 return;
697 }
698 }
699 super::stw_trace(format_args!("detach ATTACHED->DETACHED"));
700 }
701 });
702}
703
704/// Temporarily transition the current thread ATTACHED → DETACHED while
705/// running `f`, then re-attach afterwards. This allows `stop_the_world`
706/// to park this thread during blocking operations.
707///
708/// `Py_BEGIN_ALLOW_THREADS` / `Py_END_ALLOW_THREADS` equivalent.
709#[cfg(feature = "threading")]
710pub fn allow_threads<R>(vm: &VirtualMachine, f: impl FnOnce() -> R) -> R {
711 // Preserve save/restore semantics:
712 // only detach if this call observed ATTACHED at entry, and always restore
713 // on unwind.
714 let should_transition = CURRENT_THREAD_SLOT.with(|slot| {
715 slot.borrow()
716 .as_ref()
717 .is_some_and(|s| s.state.load(Ordering::Acquire) == ThreadState::Attached as i32)
718 });
719 if !should_transition {
720 return f();
721 }
722
723 detach_thread();
724 let reattach_guard = scopeguard::guard(vm, attach_thread);
725 let result = f();
726 drop(reattach_guard);
727 result
728}
729
730/// No-op on non-threading builds.
731#[cfg(not(feature = "threading"))]
732pub fn allow_threads<R>(_vm: &VirtualMachine, f: impl FnOnce() -> R) -> R {
733 f()
734}
735
736/// Run `f` with this thread attached, then return it to where it was.
737///
738/// The inverse of [`allow_threads`], for a callback that has to run Python from
739/// inside a call the thread detached for — a handshake callback reaching a
740/// Python `sni_callback`, say. Running that detached would execute Python on a
741/// thread a stop-the-world requester counts as parked. `PyGILState_Ensure` and
742/// `PyGILState_Release` bracket such a callback for the same reason.
743///
744/// A thread already attached, or one with no interpreter to attach to, just
745/// runs `f`. A thread a stop-the-world has already moved to SUSPENDED parks
746/// here until the world starts again, because [`attach_thread`] treats that
747/// state as the wait it is; that is the point of routing through it rather than
748/// testing for DETACHED alone.
749#[cfg(feature = "threading")]
750pub fn attach_for_callback<R>(vm: &VirtualMachine, f: impl FnOnce() -> R) -> R {
751 let should_transition = CURRENT_THREAD_SLOT.with(|slot| {
752 slot.borrow()
753 .as_ref()
754 .is_some_and(|s| s.state.load(Ordering::Acquire) != ThreadState::Attached as i32)
755 });
756 if !should_transition {
757 return f();
758 }
759
760 attach_thread(vm);
761 // Detach again even if `f` unwinds, so the `allow_threads` this is nested
762 // inside still finds the state it left behind.
763 let redetach_guard = scopeguard::guard((), |()| detach_thread());
764 let result = f();
765 drop(redetach_guard);
766 result
767}
768
769/// No-op on non-threading builds.
770#[cfg(not(feature = "threading"))]
771pub fn attach_for_callback<R>(_vm: &VirtualMachine, f: impl FnOnce() -> R) -> R {
772 f()
773}
774
775/// Wait for a lock the way a blocking call waits: detached, so a
776/// stop-the-world requester never has to wait for this thread to reach a
777/// safepoint it cannot reach while blocked.
778///
779/// Threads with no interpreter to leave — a native thread, or one whose
780/// locals are already being destroyed — simply block.
781///
782/// Detaching cannot park the one thread that can start the world again:
783/// [`park_detached_threads`](super::StopTheWorldState) skips the requester's
784/// slot outright, by thread id, and [`suspend_if_needed`] keys off a stop bit
785/// never set for it. That exemption is wider than the one `_PyEval_StopTheWorld`
786/// gives, where only an ATTACHED requester is skipped and a DETACHED one is
787/// suspended like any other thread — so this rests on a local invariant rather
788/// than on the reference behavior.
789#[cfg(feature = "threading")]
790fn wait_detached_from_interpreter(wait: &dyn Fn()) {
791 // Read the VM out before waiting: attaching afterwards reaches for the
792 // same thread locals, which must not still be borrowed here.
793 let current = VM_STACK
794 .try_with(|vms| vms.try_borrow().ok()?.last().copied())
795 .ok()
796 .flatten();
797 match current {
798 // SAFETY: entries in VM_STACK either borrow a VM for the dynamic
799 // scope of a set_current_vm()/enter_vm() call or point at GILSTATE_VM.
800 Some(vm) => allow_threads(unsafe { vm.as_ref() }, wait),
801 None => wait(),
802 }
803}
804
805/// Teach the lock types how to detach this thread. Idempotent, so every
806/// interpreter can call it while initializing.
807#[cfg(feature = "threading")]
808pub(crate) fn install_blocking_wait_hook() {
809 rustpython_common::lock::set_blocking_wait_hook(wait_detached_from_interpreter);
810}
811
812/// Called from check_signals when stop-the-world is requested.
813/// Transitions ATTACHED → SUSPENDED and waits until released
814/// (like `_PyThreadState_Suspend` + `_PyThreadState_Attach`).
815#[cfg(feature = "threading")]
816pub fn suspend_if_needed(state: &PyGlobalState) {
817 let should_suspend = CURRENT_THREAD_SLOT.with(|slot| {
818 slot.borrow()
819 .as_ref()
820 .is_some_and(|s| s.stop_requested.load(Ordering::Relaxed))
821 });
822 if should_suspend {
823 do_suspend(state);
824 }
825}
826
827#[cfg(feature = "threading")]
828#[cold]
829fn do_suspend(state: &PyGlobalState) {
830 let stw = &state.stop_the_world;
831 CURRENT_THREAD_SLOT.with(|slot| {
832 let borrowed = slot.borrow();
833 let Some(s) = borrowed.as_ref() else {
834 return;
835 };
836
837 // Decide whether to park while holding the thread registry. Both edges
838 // of `requested` are written under that lock: `init_thread_countdown`
839 // sets it, and `start_the_world` clears it and then releases every
840 // SUSPENDED thread without letting go. Publishing SUSPENDED here is
841 // therefore either seen by that release pass or never reached, which
842 // leaves the requester the only writer that takes a thread out of
843 // SUSPENDED. A completion check that observed this thread parked cannot
844 // then be invalidated by the thread resuming on its own.
845 let park = {
846 let _registry = state.thread_frames.lock();
847 if stw.requested.load(Ordering::Acquire) {
848 Some(s.state.compare_exchange(
849 ThreadState::Attached as i32,
850 ThreadState::Suspended as i32,
851 Ordering::AcqRel,
852 Ordering::Acquire,
853 ))
854 } else {
855 // The stop already ended; this thread's request bit is stale.
856 s.stop_requested.store(false, Ordering::Release);
857 None
858 }
859 };
860
861 match park {
862 None => {
863 super::stw_trace(format_args!("suspend skip not-requested"));
864 return;
865 }
866 Some(Ok(_)) => {
867 // Consumed this thread's stop request bit.
868 s.stop_requested.store(false, Ordering::Release);
869 }
870 Some(Err(state)) => match ThreadState::from_i32(state) {
871 Some(ThreadState::Detached) => {
872 // Leaving VM; caller will re-check on next entry.
873 super::stw_trace(format_args!("suspend skip DETACHED"));
874 return;
875 }
876 Some(ThreadState::Suspended) => {
877 // Already parked by another path.
878 s.stop_requested.store(false, Ordering::Release);
879 super::stw_trace(format_args!("suspend skip already-suspended"));
880 return;
881 }
882 Some(ThreadState::ShuttingDown) => {
883 s.stop_requested.store(false, Ordering::Release);
884 super::stw_trace(format_args!("suspend hang shutting-down"));
885 hang_thread();
886 }
887 _ => {
888 debug_assert!(false, "unexpected thread state in suspend: {state}");
889 return;
890 }
891 },
892 }
893 super::stw_trace(format_args!("suspend ATTACHED->SUSPENDED"));
894
895 // Notify the stop-the-world requester that we've parked. The registry
896 // is released first: the requester's wait loop takes the notify mutex
897 // and then the registry, so taking them the other way round here would
898 // invert the order.
899 stw.notify_suspended();
900 super::stw_trace(format_args!("suspend notified-requester"));
901
902 // Wait until start_the_world sets us back to DETACHED
903 let wait_yields = wait_while_suspended(s);
904 stw.add_suspend_wait_yields(wait_yields);
905
906 // Re-attach (DETACHED → ATTACHED), tstate_wait_attach CAS loop.
907 loop {
908 match s.state.compare_exchange(
909 ThreadState::Detached as i32,
910 ThreadState::Attached as i32,
911 Ordering::AcqRel,
912 Ordering::Acquire,
913 ) {
914 Ok(_) => break,
915 Err(state) => match ThreadState::from_i32(state) {
916 Some(ThreadState::Suspended) => {
917 let extra_wait = wait_while_suspended(s);
918 stw.add_suspend_wait_yields(extra_wait);
919 }
920 Some(ThreadState::Attached) => break,
921 Some(ThreadState::ShuttingDown) => {
922 super::stw_trace(format_args!("suspend resume hang shutting-down"));
923 hang_thread();
924 }
925 _ => {
926 debug_assert!(false, "unexpected post-suspend state: {state}");
927 break;
928 }
929 },
930 }
931 }
932 s.stop_requested.store(false, Ordering::Release);
933 super::stw_trace(format_args!("suspend resume -> ATTACHED"));
934 });
935}
936
937#[cfg(feature = "threading")]
938#[inline]
939#[must_use]
940pub fn stop_requested_for_current_thread() -> bool {
941 CURRENT_STOP_REQUESTED.with(|cached| {
942 let flag = cached.get();
943 // SAFETY: the pointer is non-null only while `CURRENT_THREAD_SLOT`
944 // holds the `Arc<ThreadSlot>` that owns the flag; both are cleared
945 // together in `cleanup_current_thread_frames`.
946 !flag.is_null() && unsafe { &*flag }.load(Ordering::Relaxed)
947 })
948}
949
950#[cfg(all(test, feature = "threading"))]
951pub(crate) fn set_stop_requested_for_current_thread(value: bool) -> bool {
952 CURRENT_STOP_REQUESTED.with(|cached| {
953 let flag = cached.get();
954 if flag.is_null() {
955 return false;
956 }
957 // SAFETY: same lifetime as `stop_requested_for_current_thread`.
958 unsafe { &*flag }.store(value, Ordering::Release);
959 true
960 })
961}
962
963/// Whether the QSBR subsystem asked this thread to pass a checkpoint.
964/// A missed or racing read of this flag is harmless: the pending
965/// retirement is still processed at the next checkpoint or by the GC
966/// backstop.
967#[cfg(feature = "threading")]
968pub(crate) fn qsbr_break_requested() -> bool {
969 CURRENT_THREAD_SLOT.with(|slot| {
970 slot.borrow()
971 .as_ref()
972 .is_some_and(|s| s.qsbr.requested.load(Ordering::Relaxed))
973 })
974}
975
976/// Pass a QSBR checkpoint: the calling thread holds no borrowed cache
977/// pointers here (instruction boundary), so mark it quiescent and try to
978/// free retired allocations.
979#[cfg(feature = "threading")]
980pub(crate) fn qsbr_checkpoint() {
981 use crate::object::qsbr::QSBR;
982 CURRENT_THREAD_SLOT.with(|slot| {
983 if let Some(s) = slot.borrow().as_ref() {
984 s.qsbr.requested.store(false, Ordering::Relaxed);
985 QSBR.quiescent_state(&s.qsbr);
986 }
987 });
988 QSBR.process();
989}
990
991/// Debug check: lock-free type-cache reads are only sound on threads that
992/// are registered with QSBR and currently ATTACHED.
993#[cfg(all(feature = "threading", debug_assertions))]
994pub(crate) fn debug_assert_current_thread_attached() {
995 CURRENT_THREAD_SLOT.with(|slot| {
996 if let Some(s) = slot.borrow().as_ref() {
997 debug_assert_eq!(
998 s.state.load(Ordering::Relaxed),
999 ThreadState::Attached as i32,
1000 "type cache read while thread not ATTACHED"
1001 );
1002 }
1003 });
1004}
1005
1006/// Push a frame pointer onto the current thread's shared frame stack.
1007/// The pointed-to frame must remain alive until the matching pop.
1008///
1009/// Only used on non-unix threading builds; unix builds publish the top frame
1010/// through `set_current_frame` writing `ThreadSlot::top_frame`.
1011#[cfg(all(not(unix), feature = "threading"))]
1012pub fn push_thread_frame(fp: FramePtr) {
1013 CURRENT_THREAD_SLOT.with(|slot| {
1014 if let Some(s) = slot.borrow().as_ref() {
1015 s.frames.lock().push(fp);
1016 } else {
1017 debug_assert!(
1018 false,
1019 "push_thread_frame called without initialized thread slot"
1020 );
1021 }
1022 });
1023}
1024
1025/// Pop a frame from the current thread's shared frame stack.
1026/// Called when a frame is exited.
1027#[cfg(all(not(unix), feature = "threading"))]
1028pub fn pop_thread_frame() {
1029 CURRENT_THREAD_SLOT.with(|slot| {
1030 if let Some(s) = slot.borrow().as_ref() {
1031 s.frames.lock().pop();
1032 } else {
1033 debug_assert!(
1034 false,
1035 "pop_thread_frame called without initialized thread slot"
1036 );
1037 }
1038 });
1039}
1040
1041/// Set the current thread's top InterpreterFrame pointer.
1042/// Returns the previous pointer so it can be restored on pop.
1043#[must_use]
1044#[allow(clippy::not_unsafe_ptr_arg_deref)]
1045pub fn set_current_frame(frame: *const InterpreterFrame) -> *const InterpreterFrame {
1046 FRAME_SLOT_CACHE.with(|cache| {
1047 // Publish the top frame for cross-thread readers (faulthandler,
1048 // sys._current_frames).
1049 #[cfg(feature = "threading")]
1050 {
1051 let slot = cache.top_iframe.get();
1052 if !slot.is_null() {
1053 unsafe { &*slot }.store(frame as usize, Ordering::Relaxed);
1054 }
1055 #[cfg(unix)]
1056 {
1057 let slot = cache.top_frame.get();
1058 if !slot.is_null() {
1059 let fo_ptr = if frame.is_null() {
1060 core::ptr::null_mut()
1061 } else {
1062 let frame_obj = unsafe { (*frame).frame_obj() };
1063 frame_obj.map_or(core::ptr::null_mut(), |py| {
1064 py as *const Py<FrameObject> as *mut Py<FrameObject>
1065 })
1066 };
1067 unsafe { &*slot }.store(fo_ptr, Ordering::Relaxed);
1068 }
1069 }
1070 }
1071 cache.current_frame.swap(frame as usize, Ordering::Relaxed)
1072 }) as *const InterpreterFrame
1073}
1074
1075/// Lightweight version that only writes to TLS `current_frame`, returning
1076/// the previous value. Does not update cross-thread top_frame (that's
1077/// updated by `set_current_frame` for FrameObject-based calls).
1078#[inline(always)]
1079#[must_use]
1080pub fn set_current_frame_nosave(frame: *const InterpreterFrame) -> *const InterpreterFrame {
1081 FRAME_SLOT_CACHE.with(|cache| cache.current_frame.swap(frame as usize, Ordering::Relaxed))
1082 as *const InterpreterFrame
1083}
1084
1085/// Get the current thread's top InterpreterFrame pointer.
1086/// Used by faulthandler's signal handler to start traceback walking.
1087#[must_use]
1088pub fn get_current_frame() -> *const InterpreterFrame {
1089 FRAME_SLOT_CACHE.with(|cache| cache.current_frame.load(Ordering::Relaxed))
1090 as *const InterpreterFrame
1091}
1092
1093/// Update the current thread's exception slot atomically (no locks).
1094/// Called from push_exception/pop_exception/set_exception.
1095#[cfg(feature = "threading")]
1096pub fn update_thread_exception(exc: Option<PyBaseExceptionRef>) {
1097 CURRENT_THREAD_SLOT.with(|slot| {
1098 if let Some(s) = slot.borrow().as_ref() {
1099 // SAFETY: Called only from the owning thread. The old ref is dropped
1100 // here on the owning thread, which is safe.
1101 let _old = unsafe { s.exception.swap(exc) };
1102 }
1103 });
1104}
1105
1106/// Collect all threads' current exceptions for sys._current_exceptions().
1107/// Acquires the global registry lock briefly, then reads each slot's exception atomically.
1108#[cfg(feature = "threading")]
1109pub fn get_all_current_exceptions(vm: &VirtualMachine) -> Vec<(u64, Option<PyBaseExceptionRef>)> {
1110 let registry = vm.state.thread_frames.lock();
1111 registry
1112 .iter()
1113 .map(|(id, slot)| (*id, slot.exception.load_owned()))
1114 .collect()
1115}
1116
1117/// Cleanup thread slot for the current thread in `vm`'s interpreter.
1118/// Called at thread exit (or when leaving an interpreter permanently).
1119#[cfg(feature = "threading")]
1120pub fn cleanup_current_thread_frames(vm: &VirtualMachine) {
1121 let thread_id = crate::stdlib::_thread::get_ident();
1122 let interp_id = vm.state.interpreter_id;
1123
1124 // Prefer the slot registered for this interpreter; fall back to CURRENT.
1125 let slot_for_interp = INTERP_THREAD_SLOTS.with(|slots| slots.borrow_mut().remove(&interp_id));
1126 let current_slot = CURRENT_THREAD_SLOT.with(|slot| slot.borrow().as_ref().cloned());
1127 let slot_to_clean = slot_for_interp.or(current_slot);
1128
1129 // A dying thread should not remain logically ATTACHED while its
1130 // thread-state slot is being removed.
1131 if let Some(slot) = &slot_to_clean {
1132 let _ = slot.state.compare_exchange(
1133 ThreadState::Attached as i32,
1134 ThreadState::Detached as i32,
1135 Ordering::AcqRel,
1136 Ordering::Acquire,
1137 );
1138 }
1139
1140 // Guard against OS thread-id reuse races: only remove the registry entry
1141 // if it still points at this thread's own slot.
1142 let _removed = if let Some(slot) = &slot_to_clean {
1143 let mut registry = vm.state.thread_frames.lock();
1144 match registry.get(&thread_id) {
1145 Some(registered) if Arc::ptr_eq(registered, slot) => registry.remove(&thread_id),
1146 _ => None,
1147 }
1148 } else {
1149 None
1150 };
1151
1152 if let Some(slot) = &_removed
1153 && vm.state.stop_the_world.requested.load(Ordering::Acquire)
1154 && thread_id != vm.state.stop_the_world.requester_ident()
1155 && slot.state.load(Ordering::Relaxed) != ThreadState::Suspended as i32
1156 {
1157 // A non-requester thread disappeared while stop-the-world is pending.
1158 // Unblock requester countdown progress.
1159 vm.state.stop_the_world.notify_thread_gone();
1160 }
1161
1162 // If CURRENT pointed at the cleaned slot, clear it (and top-frame cache).
1163 CURRENT_THREAD_SLOT.with(|s| {
1164 let clear = match (s.borrow().as_ref(), slot_to_clean.as_ref()) {
1165 (Some(cur), Some(cleaned)) => Arc::ptr_eq(cur, cleaned),
1166 (Some(_), None) => false,
1167 (None, _) => false,
1168 };
1169 if clear {
1170 *s.borrow_mut() = None;
1171 #[cfg(feature = "threading")]
1172 FRAME_SLOT_CACHE.with(|cache| {
1173 #[cfg(unix)]
1174 cache.top_frame.set(core::ptr::null());
1175 cache.top_iframe.set(core::ptr::null());
1176 });
1177 #[cfg(feature = "threading")]
1178 CURRENT_STOP_REQUESTED.with(|c| c.set(core::ptr::null()));
1179 }
1180 });
1181}
1182
1183/// Reinitialize thread slot after fork. Called in child process.
1184/// Creates a fresh slot and registers it for the current thread,
1185/// preserving the current thread's frames from the signal-safe frame chain.
1186///
1187/// Precondition: `reinit_locks_after_fork()` has already reset all
1188/// VmState locks to unlocked.
1189#[cfg(feature = "threading")]
1190pub fn reinit_frame_slot_after_fork(vm: &VirtualMachine) {
1191 let current_ident = crate::stdlib::_thread::get_ident();
1192 // On non-unix, rebuild the shared frame stack (bottom-to-top) from the
1193 // current thread's frame chain, which walks top-to-bottom via `previous`.
1194 #[cfg(not(unix))]
1195 let current_frames: Vec<FramePtr> = {
1196 let mut current_frames = Vec::new();
1197 let mut cur = get_current_frame();
1198 while !cur.is_null() {
1199 // SAFETY: the forking thread's chain frames are alive.
1200 let iframe = unsafe { &*cur };
1201 if let Some(fo) = iframe.frame_obj() {
1202 current_frames.push(FramePtr(unsafe {
1203 NonNull::new_unchecked(fo as *const _ as *mut _)
1204 }));
1205 }
1206 cur = iframe.previous.load(Ordering::Relaxed) as *const InterpreterFrame;
1207 }
1208 current_frames.reverse();
1209 current_frames
1210 };
1211 #[cfg(unix)]
1212 let top_fo_ptr = {
1213 let top_iframe = get_current_frame();
1214 if top_iframe.is_null() {
1215 core::ptr::null_mut()
1216 } else {
1217 match unsafe { (*top_iframe).frame_obj() } {
1218 Some(fo) => fo as *const Py<FrameObject> as *mut Py<FrameObject>,
1219 None => core::ptr::null_mut(),
1220 }
1221 }
1222 };
1223 let top_iframe_ptr = get_current_frame() as usize;
1224 let new_slot = Arc::new(ThreadSlot {
1225 // The surviving child thread keeps executing its current frame chain.
1226 // Only publish heavy frames for signal safety.
1227 #[cfg(unix)]
1228 top_frame: AtomicPtr::new(top_fo_ptr),
1229 top_iframe: AtomicUsize::new(top_iframe_ptr),
1230 #[cfg(not(unix))]
1231 frames: parking_lot::Mutex::new(current_frames),
1232 exception: crate::PyAtomicRef::from(vm.topmost_exception()),
1233 trace_func: PyMutex::new(vm.trace_func.borrow().clone()),
1234 profile_func: PyMutex::new(vm.profile_func.borrow().clone()),
1235 state: core::sync::atomic::AtomicI32::new(ThreadState::Attached as i32),
1236 stop_requested: core::sync::atomic::AtomicBool::new(false),
1237 thread: std::thread::current(),
1238 qsbr: crate::object::qsbr::QSBR.register(),
1239 });
1240 FRAME_SLOT_CACHE.with(|cache| {
1241 #[cfg(unix)]
1242 cache.top_frame.set(&new_slot.top_frame);
1243 cache.top_iframe.set(&new_slot.top_iframe);
1244 });
1245 #[cfg(feature = "threading")]
1246 CURRENT_STOP_REQUESTED.with(|c| c.set(&new_slot.stop_requested));
1247
1248 // Lock is safe: reinit_locks_after_fork() already reset it to unlocked.
1249 let mut registry = vm.state.thread_frames.lock();
1250 registry.clear();
1251 registry.insert(current_ident, new_slot.clone());
1252 drop(registry);
1253
1254 CURRENT_THREAD_SLOT.with(|s| {
1255 *s.borrow_mut() = Some(new_slot.clone());
1256 });
1257 INTERP_THREAD_SLOTS.with(|slots| {
1258 slots.borrow_mut().insert(vm.state.interpreter_id, new_slot);
1259 });
1260}
1261
1262/// Drop this thread's cached slots for every interpreter except `keep_id`.
1263///
1264/// After `fork()` only the calling thread survives, and the other
1265/// interpreters' registries are cleared; a cached slot would otherwise stay
1266/// current for an interpreter that no longer lists it, hiding the thread from
1267/// that interpreter's stop-the-world. The next enter builds a fresh slot.
1268#[cfg(feature = "threading")]
1269pub fn purge_other_interpreter_slots_after_fork(keep_id: i64) {
1270 INTERP_THREAD_SLOTS.with(|slots| {
1271 slots.borrow_mut().retain(|&id, _| id == keep_id);
1272 });
1273}
1274
1275/// Whether the interpreter on top of `VM_STACK` is currently ATTACHED on this
1276/// thread. Without the `threading` feature there is no attach state to check.
1277#[cfg(feature = "threading")]
1278fn top_slot_is_attached() -> bool {
1279 current_slot_is_attached()
1280}
1281#[cfg(not(feature = "threading"))]
1282fn top_slot_is_attached() -> bool {
1283 true
1284}
1285
1286/// Which VM `with_vm` found for `obj`, and whether this thread is already
1287/// attached to it (see the fast path below).
1288enum WithVmTarget {
1289 /// `interp` is on top of `VM_STACK`, so this thread is already ATTACHED
1290 /// to it (see [`begin_interpreter_section`]'s invariant): no section
1291 /// switch is needed.
1292 AlreadyCurrent(NonNull<VirtualMachine>),
1293 /// `interp` owns `obj` but is not the top of `VM_STACK` (a nested,
1294 /// currently-detached interpreter), so a real attach/detach section is
1295 /// required.
1296 NeedsSwitch(NonNull<VirtualMachine>),
1297}
1298
1299pub fn with_vm<F, R>(obj: &PyObject, f: F) -> Option<R>
1300where
1301 F: Fn(&VirtualMachine) -> R,
1302{
1303 let vm_owns_obj = |interp: NonNull<VirtualMachine>| {
1304 // SAFETY: all references in VM_STACK should be valid
1305 let vm = unsafe { interp.as_ref() };
1306 obj.fast_isinstance(vm.ctx.types.object_type)
1307 };
1308 // `with_vm` runs on every teardown of an object with a `__del__` slot or a
1309 // weakref callback (drop_slow_inner / try_call_finalizer / gc_state), which
1310 // for `__del__` objects and weakrefs collected during a GC pass means it
1311 // runs on essentially every such drop. The overwhelming majority of those
1312 // drops happen from within Python bytecode executing on this very thread,
1313 // i.e. `obj`'s owning interpreter is already the one on top of `VM_STACK`.
1314 // `begin_interpreter_section` (via `set_current_vm`) only ever needs to run
1315 // for the rare case where the object belongs to a *different* interpreter
1316 // than the one currently attached (a nested subinterpreter scenario) or no
1317 // interpreter is attached at all (object dropped on a thread outside any
1318 // VM, e.g. during shutdown or from a plain Rust thread) — in the fast case
1319 // we can skip it entirely and call `f` directly.
1320 let target = VM_STACK.with(|vms| {
1321 let vms = vms.borrow();
1322 // Fast path: at most one interpreter is ever ATTACHED per OS thread,
1323 // and it is always the one on top of `VM_STACK` — every push onto
1324 // VM_STACK (set_current_vm, VmBootstrapGuard) is paired with an attach
1325 // *before* the push, and every pop is paired with a detach (or a
1326 // re-attach of the newly-exposed top) in `end_interpreter_section`.
1327 // So if `obj`'s owning interpreter is the current top, this thread is
1328 // already attached to it and there is nothing for
1329 // `begin_interpreter_section` to do: no INTERP_THREAD_SLOTS lookup, no
1330 // Arc clone, no atomic state transition.
1331 // The top may still be DETACHED while the thread sits inside an
1332 // `allow_threads` section (a blocking call that dropped an object
1333 // with `__del__`); running `f` there would execute Python on a thread
1334 // a stop-the-world requester counts as parked, so that case takes the
1335 // full attach path below.
1336 if let Some(top) = vms.last().copied()
1337 && vm_owns_obj(top)
1338 && top_slot_is_attached()
1339 {
1340 return Some(WithVmTarget::AlreadyCurrent(top));
1341 }
1342 let interp = match vms.iter().copied().exactly_one() {
1343 Ok(x) => {
1344 debug_assert!(vm_owns_obj(x));
1345 x
1346 }
1347 Err(mut others) => others.find(|x| vm_owns_obj(*x))?,
1348 };
1349 Some(WithVmTarget::NeedsSwitch(interp))
1350 })?;
1351 match target {
1352 WithVmTarget::AlreadyCurrent(interp) => {
1353 // SAFETY: `interp` is (or was, at the point it was read above) the
1354 // top of VM_STACK for this thread, so it is valid for at least the
1355 // dynamic scope of the enclosing set_current_vm()/enter_vm() call,
1356 // which contains this whole function call.
1357 let vm = unsafe { interp.as_ref() };
1358 Some(f(vm))
1359 }
1360 WithVmTarget::NeedsSwitch(interp) => {
1361 // SAFETY: all references in VM_STACK should be valid, and should not be changed or moved
1362 // at least until this function returns and the stack unwinds to an enter_vm() call
1363 let vm = unsafe { interp.as_ref() };
1364 Some(set_current_vm(vm, || f(vm)))
1365 }
1366 }
1367}
1368
1369#[must_use = "ThreadedVirtualMachine does nothing unless you move it to another thread and call .run()"]
1370#[cfg(feature = "threading")]
1371pub struct ThreadedVirtualMachine {
1372 pub(super) vm: VirtualMachine,
1373}
1374
1375#[cfg(feature = "threading")]
1376impl ThreadedVirtualMachine {
1377 /// Create a `FnOnce()` that can easily be passed to a function like [`std::thread::Builder::spawn`]
1378 ///
1379 /// # Note
1380 ///
1381 /// If you return a `PyObjectRef` (or a type that contains one) from `F`, and don't `join()`
1382 /// on the thread this `FnOnce` runs in, there is a possibility that that thread will panic
1383 /// as `PyObjectRef`'s `Drop` implementation tries to run the `__del__` destructor of a
1384 /// Python object but finds that it's not in the context of any vm.
1385 pub fn make_spawn_func<F, R>(self, f: F) -> impl FnOnce() -> R
1386 where
1387 F: FnOnce(&VirtualMachine) -> R,
1388 {
1389 move || self.run(f)
1390 }
1391
1392 /// Run a function in this thread context
1393 ///
1394 /// # Note
1395 ///
1396 /// If you return a `PyObjectRef` (or a type that contains one) from `F`, and don't return the object
1397 /// to the parent thread and then `join()` on the `JoinHandle` (or similar), there is a possibility that
1398 /// the current thread will panic as `PyObjectRef`'s `Drop` implementation tries to run the `__del__`
1399 /// destructor of a python object but finds that it's not in the context of any vm.
1400 pub fn run<F, R>(&self, f: F) -> R
1401 where
1402 F: FnOnce(&VirtualMachine) -> R,
1403 {
1404 let vm = &self.vm;
1405 // Each spawned thread has its own native stack bounds. Recompute the
1406 // soft limit here instead of inheriting the parent thread's value.
1407 vm.c_stack_soft_limit
1408 .set(VirtualMachine::calculate_c_stack_soft_limit());
1409 enter_vm(vm, || f(vm))
1410 }
1411}
1412
1413impl VirtualMachine {
1414 /// Start a new thread with access to the same interpreter.
1415 ///
1416 /// # Note
1417 ///
1418 /// If you return a `PyObjectRef` (or a type that contains one) from `F`, and don't `join()`
1419 /// on the thread, there is a possibility that that thread will panic as `PyObjectRef`'s `Drop`
1420 /// implementation tries to run the `__del__` destructor of a python object but finds that it's
1421 /// not in the context of any vm.
1422 #[cfg(feature = "threading")]
1423 pub fn start_thread<F, R>(&self, f: F) -> std::thread::JoinHandle<R>
1424 where
1425 F: Send + 'static + FnOnce(&Self) -> R,
1426 R: Send + 'static,
1427 {
1428 let func = self.new_thread().make_spawn_func(f);
1429 std::thread::spawn(func)
1430 }
1431
1432 /// Create a new VM thread that can be passed to a function like [`std::thread::spawn`]
1433 /// to use the same interpreter on a different thread. Note that if you just want to
1434 /// use this with `thread::spawn`, you can use
1435 /// [`vm.start_thread()`](`VirtualMachine::start_thread`) as a convenience.
1436 ///
1437 /// # Usage
1438 ///
1439 /// ```
1440 /// # rustpython_vm::Interpreter::without_stdlib(Default::default()).enter(|vm| {
1441 /// use std::thread::Builder;
1442 /// let handle = Builder::new()
1443 /// .name("my thread :)".into())
1444 /// .spawn(vm.new_thread().make_spawn_func(|vm| vm.ctx.none()))
1445 /// .expect("couldn't spawn thread");
1446 /// let returned_obj = handle.join().expect("thread panicked");
1447 /// assert!(vm.is_none(&returned_obj));
1448 /// # })
1449 /// ```
1450 ///
1451 /// Note: this function is safe, but running the returned ThreadedVirtualMachine in the same
1452 /// thread context (i.e. with the same thread-local storage) doesn't have any
1453 /// specific guaranteed behavior.
1454 #[cfg(feature = "threading")]
1455 pub fn new_thread(&self) -> ThreadedVirtualMachine {
1456 let global_trace = self.state.global_trace_func.lock().clone();
1457 let global_profile = self.state.global_profile_func.lock().clone();
1458 let use_tracing = global_trace.is_some() || global_profile.is_some();
1459
1460 let vm = Self {
1461 builtins: self.builtins.clone(),
1462 sys_module: self.sys_module.clone(),
1463 ctx: self.ctx.clone(),
1464 datastack: core::cell::UnsafeCell::new(crate::datastack::DataStack::new()),
1465 wasm_id: self.wasm_id.clone(),
1466 exceptions: RefCell::default(),
1467 import_func: self.import_func.clone(),
1468 importlib: self.importlib.clone(),
1469 profile_func: RefCell::new(global_profile.unwrap_or_else(|| self.ctx.none())),
1470 trace_func: RefCell::new(global_trace.unwrap_or_else(|| self.ctx.none())),
1471 use_tracing: Cell::new(use_tracing),
1472 what_event: Cell::new(None),
1473 tracing_depth: Cell::new(0),
1474 recursion_limit: self.recursion_limit.clone(),
1475 signal_handlers: core::cell::OnceCell::new(),
1476 signal_rx: None,
1477 repr_guards: RefCell::default(),
1478 state: self.state.clone(),
1479 initialized: self.initialized,
1480 recursion_depth: Cell::new(0),
1481 #[cfg(any(miri, target_env = "musl"))]
1482 native_recursion_depth: Cell::new(0),
1483 c_stack_soft_limit: Cell::new(Self::calculate_c_stack_soft_limit()),
1484 async_gen_firstiter: RefCell::new(None),
1485 async_gen_finalizer: RefCell::new(None),
1486 asyncio_running_loop: RefCell::new(None),
1487 asyncio_running_task: RefCell::new(None),
1488 context_stack: RefCell::default(),
1489 callable_cache: self.callable_cache.clone(),
1490 pending_tailcall_frame: Cell::new(None),
1491 pending_tailcall_owner: core::cell::UnsafeCell::new(None),
1492 pending_gen_resume: core::cell::UnsafeCell::new(None),
1493 trampoline_stack: core::cell::UnsafeCell::new(Vec::new()),
1494 };
1495 ThreadedVirtualMachine { vm }
1496 }
1497}