Skip to main content

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}