#[cfg(all(not(unix), feature = "threading"))]
use super::FramePtr;
#[cfg(feature = "threading")]
use crate::PyObjectRef;
#[cfg(feature = "threading")]
use crate::builtins::PyBaseExceptionRef;
#[cfg(feature = "threading")]
use alloc::sync::Arc;
#[cfg(feature = "threading")]
use rustpython_common::lock::PyMutex;
use crate::frame::InterpreterFrame;
#[cfg(feature = "threading")]
use crate::vm::PyGlobalState;
use crate::{AsObject, PyObject, VirtualMachine};
#[cfg(all(unix, feature = "threading"))]
use crate::{Py, frame::FrameObject};
#[cfg(all(unix, feature = "threading"))]
use core::sync::atomic::AtomicPtr;
use core::{
cell::{Cell, RefCell},
ptr::NonNull,
sync::atomic::{AtomicUsize, Ordering},
};
use itertools::Itertools;
#[cfg(feature = "threading")]
use std::collections::HashMap;
use std::thread_local;
#[cfg(feature = "threading")]
#[repr(i32)]
#[derive(Copy, Clone, Debug, PartialEq, Eq)]
pub enum ThreadState {
Detached = 0,
Attached = 1,
Suspended = 2,
ShuttingDown = 3,
}
#[cfg(feature = "threading")]
impl ThreadState {
#[must_use]
pub const fn from_i32(v: i32) -> Option<Self> {
match v {
0 => Some(Self::Detached),
1 => Some(Self::Attached),
2 => Some(Self::Suspended),
3 => Some(Self::ShuttingDown),
_ => None,
}
}
}
#[cfg(feature = "threading")]
pub struct ThreadSlot {
#[cfg(unix)]
pub top_frame: AtomicPtr<Py<FrameObject>>,
pub top_iframe: AtomicUsize,
#[cfg(not(unix))]
pub frames: parking_lot::Mutex<Vec<FramePtr>>,
pub exception: crate::PyAtomicRef<Option<crate::exceptions::types::PyBaseException>>,
pub trace_func: PyMutex<PyObjectRef>,
pub profile_func: PyMutex<PyObjectRef>,
pub state: core::sync::atomic::AtomicI32,
pub stop_requested: core::sync::atomic::AtomicBool,
pub thread: std::thread::Thread,
pub(crate) qsbr: Arc<crate::object::qsbr::QsbrSlot>,
}
#[cfg(feature = "threading")]
pub type CurrentFrameSlot = Arc<ThreadSlot>;
struct FrameSlotCache {
current_frame: AtomicUsize,
#[cfg(all(unix, feature = "threading"))]
top_frame: Cell<*const AtomicPtr<Py<FrameObject>>>,
#[cfg(feature = "threading")]
top_iframe: Cell<*const AtomicUsize>,
}
thread_local! {
pub(super) static VM_STACK: RefCell<Vec<NonNull<VirtualMachine>>> = Vec::with_capacity(1).into();
#[cfg(feature = "threading")]
static GILSTATE_VM: RefCell<Option<Box<ThreadedVirtualMachine>>> = const { RefCell::new(None) };
pub(crate) static COROUTINE_ORIGIN_TRACKING_DEPTH: Cell<u32> = const { Cell::new(0) };
#[cfg(feature = "threading")]
static INTERP_THREAD_SLOTS: RefCell<HashMap<i64, CurrentFrameSlot>> =
RefCell::new(HashMap::new());
#[cfg(feature = "threading")]
static CURRENT_THREAD_SLOT: RefCell<Option<CurrentFrameSlot>> = const { RefCell::new(None) };
pub(crate) static FRAME_SLOT_CACHE: FrameSlotCache = const {
FrameSlotCache {
current_frame: AtomicUsize::new(0),
#[cfg(all(unix, feature = "threading"))]
top_frame: Cell::new(core::ptr::null()),
#[cfg(feature = "threading")]
top_iframe: Cell::new(core::ptr::null()),
}
};
#[cfg(feature = "threading")]
static CURRENT_STOP_REQUESTED: Cell<*const core::sync::atomic::AtomicBool> =
const { Cell::new(core::ptr::null()) };
}
#[must_use]
pub fn current_vm_is_set() -> bool {
VM_STACK.with(|vms| !vms.borrow().is_empty())
}
pub fn with_current_vm<R>(f: impl FnOnce(&VirtualMachine) -> R) -> R {
VM_STACK.with(|vms| {
let vm = vms
.borrow()
.last()
.copied()
.expect("call with_current_vm() but no current VM is attached");
f(unsafe { vm.as_ref() })
})
}
fn set_current_vm<R>(vm: &VirtualMachine, f: impl FnOnce() -> R) -> R {
#[cfg(feature = "threading")]
let switched = begin_interpreter_section(vm);
VM_STACK.with(|vms| {
vms.borrow_mut().push(vm.into());
scopeguard::defer! {
vms.borrow_mut().pop();
#[cfg(feature = "threading")]
end_interpreter_section(switched);
}
f()
})
}
pub(crate) fn current_gc_state() -> Option<NonNull<crate::gc_state::GcInterpreterState>> {
VM_STACK
.try_with(|vms| {
let vm = vms.try_borrow().ok()?.last().copied()?;
Some(NonNull::from(&unsafe { vm.as_ref() }.state.gc))
})
.ok()
.flatten()
}
pub fn try_with_current_vm<R>(f: impl FnOnce(&VirtualMachine) -> R) -> Option<R> {
VM_STACK.with(|vms| {
let vm = vms.borrow().last().copied()?;
Some(f(unsafe { vm.as_ref() }))
})
}
pub fn enter_vm<R>(vm: &VirtualMachine, f: impl FnOnce() -> R) -> R {
set_current_vm(vm, f)
}
#[must_use]
pub(crate) struct VmBootstrapGuard {
#[cfg(feature = "threading")]
switched: bool,
}
impl VmBootstrapGuard {
pub(crate) fn new(vm: &VirtualMachine) -> Self {
#[cfg(feature = "threading")]
let switched = begin_interpreter_section(vm);
VM_STACK.with(|vms| vms.borrow_mut().push(vm.into()));
Self {
#[cfg(feature = "threading")]
switched,
}
}
}
impl Drop for VmBootstrapGuard {
fn drop(&mut self) {
VM_STACK.with(|vms| {
vms.borrow_mut().pop();
});
#[cfg(feature = "threading")]
end_interpreter_section(self.switched);
}
}
#[cfg(feature = "threading")]
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum CurrentVmAttachState {
AlreadyAttached,
Attached,
}
#[cfg(feature = "threading")]
pub struct SavedThreadState {
vm_stack: Vec<NonNull<VirtualMachine>>,
gilstate_vm: Option<Box<ThreadedVirtualMachine>>,
}
#[cfg(feature = "threading")]
#[must_use = "the saved thread state must be restored"]
pub fn save_current_thread() -> SavedThreadState {
let vm_stack = VM_STACK.with(|vms| core::mem::take(&mut *vms.borrow_mut()));
assert!(
!vm_stack.is_empty(),
"save_current_thread() called without an attached VM"
);
let gilstate_vm = GILSTATE_VM.with(|gilstate_vm| gilstate_vm.borrow_mut().take());
detach_thread();
SavedThreadState {
vm_stack,
gilstate_vm,
}
}
#[cfg(feature = "threading")]
pub fn restore_current_thread(state: SavedThreadState) {
assert!(
!current_vm_is_set(),
"restore_current_thread() called with an attached VM"
);
let SavedThreadState {
vm_stack,
gilstate_vm,
} = state;
let vm = vm_stack
.last()
.copied()
.expect("saved thread state has no VM");
GILSTATE_VM.with(|current| {
let mut current = current.borrow_mut();
assert!(
current.is_none(),
"restore_current_thread() called with a GILState VM"
);
*current = gilstate_vm;
});
let vm = unsafe { vm.as_ref() };
init_thread_slot_if_needed(vm);
attach_thread(vm);
VM_STACK.with(|vms| *vms.borrow_mut() = vm_stack);
}
#[cfg(feature = "threading")]
pub fn attach_current_thread(
make_vm: impl FnOnce() -> ThreadedVirtualMachine,
) -> CurrentVmAttachState {
if current_vm_is_set() {
return CurrentVmAttachState::AlreadyAttached;
}
GILSTATE_VM.with(|gilstate_vm| {
let mut gilstate_vm = gilstate_vm.borrow_mut();
let threaded_vm = gilstate_vm.get_or_insert_with(|| Box::new(make_vm()));
let vm = &threaded_vm.vm;
vm.c_stack_soft_limit
.set(VirtualMachine::calculate_c_stack_soft_limit());
init_thread_slot_if_needed(vm);
attach_thread(vm);
VM_STACK.with(|vms| {
debug_assert!(vms.borrow().is_empty());
vms.borrow_mut().push(vm.into());
});
});
CurrentVmAttachState::Attached
}
#[cfg(feature = "threading")]
pub fn release_current_thread(state: CurrentVmAttachState) {
if state == CurrentVmAttachState::AlreadyAttached {
return;
}
let gilstate_vm = GILSTATE_VM.with(|gilstate_vm| gilstate_vm.borrow_mut().take());
drop(gilstate_vm);
VM_STACK.with(|vms| {
vms.borrow_mut()
.pop()
.expect("release_current_thread() called without an attached VM");
});
detach_thread();
}
#[cfg(feature = "threading")]
fn init_thread_slot_if_needed(vm: &VirtualMachine) {
let slot = ensure_thread_slot(vm);
set_current_thread_slot(slot);
}
#[cfg(feature = "threading")]
fn ensure_thread_slot(vm: &VirtualMachine) -> CurrentFrameSlot {
let interp_id = vm.state.interpreter_id;
INTERP_THREAD_SLOTS.with(|slots| {
let mut slots = slots.borrow_mut();
if let Some(existing) = slots.get(&interp_id) {
return existing.clone();
}
let thread_id = crate::stdlib::_thread::get_ident();
let mut registry = vm.state.thread_frames.lock();
let new_slot = Arc::new(ThreadSlot {
#[cfg(unix)]
top_frame: AtomicPtr::new(core::ptr::null_mut()),
top_iframe: AtomicUsize::new(0),
#[cfg(not(unix))]
frames: parking_lot::Mutex::new(Vec::new()),
exception: crate::PyAtomicRef::from(None::<PyBaseExceptionRef>),
trace_func: PyMutex::new(
vm.state
.global_trace_func
.lock()
.clone()
.unwrap_or_else(|| vm.ctx.none()),
),
profile_func: PyMutex::new(
vm.state
.global_profile_func
.lock()
.clone()
.unwrap_or_else(|| vm.ctx.none()),
),
state: core::sync::atomic::AtomicI32::new(
if vm.state.stop_the_world.requested.load(Ordering::Acquire) {
ThreadState::Suspended as i32
} else {
ThreadState::Detached as i32
},
),
stop_requested: core::sync::atomic::AtomicBool::new(false),
thread: std::thread::current(),
qsbr: crate::object::qsbr::QSBR.register(),
});
registry.insert(thread_id, new_slot.clone());
drop(registry);
slots.insert(interp_id, new_slot.clone());
new_slot
})
}
#[cfg(feature = "threading")]
#[must_use]
pub fn current_thread_slot() -> Option<CurrentFrameSlot> {
CURRENT_THREAD_SLOT.with(|slot| slot.borrow().clone())
}
#[cfg(feature = "threading")]
fn set_current_thread_slot(slot: CurrentFrameSlot) {
FRAME_SLOT_CACHE.with(|cache| {
#[cfg(unix)]
cache.top_frame.set(&slot.top_frame);
cache.top_iframe.set(&slot.top_iframe);
});
CURRENT_STOP_REQUESTED.with(|c| c.set(&slot.stop_requested));
CURRENT_THREAD_SLOT.with(|current| {
*current.borrow_mut() = Some(slot);
});
}
#[cfg(feature = "threading")]
fn current_slot_is_attached() -> bool {
CURRENT_THREAD_SLOT.with(|slot| {
slot.borrow()
.as_ref()
.is_some_and(|s| s.state.load(Ordering::Acquire) == ThreadState::Attached as i32)
})
}
#[cfg(feature = "threading")]
fn begin_interpreter_section(vm: &VirtualMachine) -> bool {
let target = ensure_thread_slot(vm);
let already_current = CURRENT_THREAD_SLOT.with(|slot| {
slot.borrow()
.as_ref()
.is_some_and(|s| Arc::ptr_eq(s, &target))
});
if already_current && current_slot_is_attached() {
return false;
}
if !already_current && current_slot_is_attached() {
detach_thread();
}
set_current_thread_slot(target);
attach_thread(vm);
true
}
#[cfg(feature = "threading")]
fn end_interpreter_section(switched: bool) {
if !switched {
return;
}
if current_slot_is_attached() {
detach_thread();
}
if let Some(vm_ptr) = VM_STACK.with(|vms| vms.borrow().last().copied()) {
let vm = unsafe { vm_ptr.as_ref() };
set_current_thread_slot(ensure_thread_slot(vm));
attach_thread(vm);
}
}
#[cfg(feature = "threading")]
fn wait_while_suspended(slot: &ThreadSlot) -> u64 {
let mut wait_yields = 0u64;
while slot.state.load(Ordering::Acquire) == ThreadState::Suspended as i32 {
wait_yields = wait_yields.saturating_add(1);
std::thread::park();
}
wait_yields
}
#[cfg(feature = "threading")]
fn hang_thread() -> ! {
loop {
std::thread::park();
}
}
#[cfg(feature = "threading")]
pub fn hang_current_thread(state: &PyGlobalState) -> ! {
CURRENT_THREAD_SLOT.with(|slot| {
if let Some(s) = slot.borrow().as_ref() {
let prev = s
.state
.swap(ThreadState::ShuttingDown as i32, Ordering::AcqRel);
if prev == ThreadState::Attached as i32 {
crate::object::qsbr::QSBR.offline(&s.qsbr);
}
s.stop_requested.store(false, Ordering::Release);
}
});
if state.stop_the_world.requested.load(Ordering::Acquire) {
state.stop_the_world.notify_thread_gone();
}
hang_thread();
}
#[cfg(feature = "threading")]
pub fn set_other_threads_shutting_down(state: &PyGlobalState) {
let current = crate::stdlib::_thread::get_ident();
let registry = state.thread_frames.lock();
#[expect(
clippy::iter_over_hash_type,
reason = "Iteration order doesn't matter here"
)]
for (&id, slot) in registry.iter() {
if id == current {
continue;
}
slot.stop_requested.store(false, Ordering::Release);
slot.state
.store(ThreadState::ShuttingDown as i32, Ordering::Release);
slot.thread.unpark();
}
}
#[cfg(feature = "threading")]
fn attach_thread(vm: &VirtualMachine) {
CURRENT_THREAD_SLOT.with(|slot| {
if let Some(s) = slot.borrow().as_ref() {
super::stw_trace(format_args!("attach begin"));
loop {
match s.state.compare_exchange(
ThreadState::Detached as i32,
ThreadState::Attached as i32,
Ordering::AcqRel,
Ordering::Relaxed,
) {
Ok(_) => {
crate::object::qsbr::QSBR.online(&s.qsbr);
super::stw_trace(format_args!("attach DETACHED->ATTACHED"));
break;
}
Err(state) => match ThreadState::from_i32(state) {
Some(ThreadState::Suspended) => {
super::stw_trace(format_args!("attach wait-suspended"));
let wait_yields = wait_while_suspended(s);
vm.state.stop_the_world.add_attach_wait_yields(wait_yields);
}
Some(ThreadState::ShuttingDown) => {
super::stw_trace(format_args!("attach hang shutting-down"));
hang_thread();
}
_ => {
debug_assert!(false, "unexpected thread state in attach: {state}");
break;
}
},
}
}
}
});
suspend_if_needed(&vm.state);
}
#[cfg(feature = "threading")]
fn detach_thread() {
CURRENT_THREAD_SLOT.with(|slot| {
if let Some(s) = slot.borrow().as_ref() {
match s.state.compare_exchange(
ThreadState::Attached as i32,
ThreadState::Detached as i32,
Ordering::AcqRel,
Ordering::Acquire,
) {
Ok(_) => {
crate::object::qsbr::QSBR.offline(&s.qsbr);
}
Err(state) => {
debug_assert!(
matches!(ThreadState::from_i32(state), Some(ThreadState::Detached)),
"unexpected thread state in detach: {state}"
);
return;
}
}
super::stw_trace(format_args!("detach ATTACHED->DETACHED"));
}
});
}
#[cfg(feature = "threading")]
pub fn allow_threads<R>(vm: &VirtualMachine, f: impl FnOnce() -> R) -> R {
let should_transition = CURRENT_THREAD_SLOT.with(|slot| {
slot.borrow()
.as_ref()
.is_some_and(|s| s.state.load(Ordering::Acquire) == ThreadState::Attached as i32)
});
if !should_transition {
return f();
}
detach_thread();
let reattach_guard = scopeguard::guard(vm, attach_thread);
let result = f();
drop(reattach_guard);
result
}
#[cfg(not(feature = "threading"))]
pub fn allow_threads<R>(_vm: &VirtualMachine, f: impl FnOnce() -> R) -> R {
f()
}
#[cfg(feature = "threading")]
pub fn attach_for_callback<R>(vm: &VirtualMachine, f: impl FnOnce() -> R) -> R {
let should_transition = CURRENT_THREAD_SLOT.with(|slot| {
slot.borrow()
.as_ref()
.is_some_and(|s| s.state.load(Ordering::Acquire) != ThreadState::Attached as i32)
});
if !should_transition {
return f();
}
attach_thread(vm);
let redetach_guard = scopeguard::guard((), |()| detach_thread());
let result = f();
drop(redetach_guard);
result
}
#[cfg(not(feature = "threading"))]
pub fn attach_for_callback<R>(_vm: &VirtualMachine, f: impl FnOnce() -> R) -> R {
f()
}
#[cfg(feature = "threading")]
fn wait_detached_from_interpreter(wait: &dyn Fn()) {
let current = VM_STACK
.try_with(|vms| vms.try_borrow().ok()?.last().copied())
.ok()
.flatten();
match current {
Some(vm) => allow_threads(unsafe { vm.as_ref() }, wait),
None => wait(),
}
}
#[cfg(feature = "threading")]
pub(crate) fn install_blocking_wait_hook() {
rustpython_common::lock::set_blocking_wait_hook(wait_detached_from_interpreter);
}
#[cfg(feature = "threading")]
pub fn suspend_if_needed(state: &PyGlobalState) {
let should_suspend = CURRENT_THREAD_SLOT.with(|slot| {
slot.borrow()
.as_ref()
.is_some_and(|s| s.stop_requested.load(Ordering::Relaxed))
});
if should_suspend {
do_suspend(state);
}
}
#[cfg(feature = "threading")]
#[cold]
fn do_suspend(state: &PyGlobalState) {
let stw = &state.stop_the_world;
CURRENT_THREAD_SLOT.with(|slot| {
let borrowed = slot.borrow();
let Some(s) = borrowed.as_ref() else {
return;
};
let park = {
let _registry = state.thread_frames.lock();
if stw.requested.load(Ordering::Acquire) {
Some(s.state.compare_exchange(
ThreadState::Attached as i32,
ThreadState::Suspended as i32,
Ordering::AcqRel,
Ordering::Acquire,
))
} else {
s.stop_requested.store(false, Ordering::Release);
None
}
};
match park {
None => {
super::stw_trace(format_args!("suspend skip not-requested"));
return;
}
Some(Ok(_)) => {
s.stop_requested.store(false, Ordering::Release);
}
Some(Err(state)) => match ThreadState::from_i32(state) {
Some(ThreadState::Detached) => {
super::stw_trace(format_args!("suspend skip DETACHED"));
return;
}
Some(ThreadState::Suspended) => {
s.stop_requested.store(false, Ordering::Release);
super::stw_trace(format_args!("suspend skip already-suspended"));
return;
}
Some(ThreadState::ShuttingDown) => {
s.stop_requested.store(false, Ordering::Release);
super::stw_trace(format_args!("suspend hang shutting-down"));
hang_thread();
}
_ => {
debug_assert!(false, "unexpected thread state in suspend: {state}");
return;
}
},
}
super::stw_trace(format_args!("suspend ATTACHED->SUSPENDED"));
stw.notify_suspended();
super::stw_trace(format_args!("suspend notified-requester"));
let wait_yields = wait_while_suspended(s);
stw.add_suspend_wait_yields(wait_yields);
loop {
match s.state.compare_exchange(
ThreadState::Detached as i32,
ThreadState::Attached as i32,
Ordering::AcqRel,
Ordering::Acquire,
) {
Ok(_) => break,
Err(state) => match ThreadState::from_i32(state) {
Some(ThreadState::Suspended) => {
let extra_wait = wait_while_suspended(s);
stw.add_suspend_wait_yields(extra_wait);
}
Some(ThreadState::Attached) => break,
Some(ThreadState::ShuttingDown) => {
super::stw_trace(format_args!("suspend resume hang shutting-down"));
hang_thread();
}
_ => {
debug_assert!(false, "unexpected post-suspend state: {state}");
break;
}
},
}
}
s.stop_requested.store(false, Ordering::Release);
super::stw_trace(format_args!("suspend resume -> ATTACHED"));
});
}
#[cfg(feature = "threading")]
#[inline]
#[must_use]
pub fn stop_requested_for_current_thread() -> bool {
CURRENT_STOP_REQUESTED.with(|cached| {
let flag = cached.get();
!flag.is_null() && unsafe { &*flag }.load(Ordering::Relaxed)
})
}
#[cfg(all(test, feature = "threading"))]
pub(crate) fn set_stop_requested_for_current_thread(value: bool) -> bool {
CURRENT_STOP_REQUESTED.with(|cached| {
let flag = cached.get();
if flag.is_null() {
return false;
}
unsafe { &*flag }.store(value, Ordering::Release);
true
})
}
#[cfg(feature = "threading")]
pub(crate) fn qsbr_break_requested() -> bool {
CURRENT_THREAD_SLOT.with(|slot| {
slot.borrow()
.as_ref()
.is_some_and(|s| s.qsbr.requested.load(Ordering::Relaxed))
})
}
#[cfg(feature = "threading")]
pub(crate) fn qsbr_checkpoint() {
use crate::object::qsbr::QSBR;
CURRENT_THREAD_SLOT.with(|slot| {
if let Some(s) = slot.borrow().as_ref() {
s.qsbr.requested.store(false, Ordering::Relaxed);
QSBR.quiescent_state(&s.qsbr);
}
});
QSBR.process();
}
#[cfg(all(feature = "threading", debug_assertions))]
pub(crate) fn debug_assert_current_thread_attached() {
CURRENT_THREAD_SLOT.with(|slot| {
if let Some(s) = slot.borrow().as_ref() {
debug_assert_eq!(
s.state.load(Ordering::Relaxed),
ThreadState::Attached as i32,
"type cache read while thread not ATTACHED"
);
}
});
}
#[cfg(all(not(unix), feature = "threading"))]
pub fn push_thread_frame(fp: FramePtr) {
CURRENT_THREAD_SLOT.with(|slot| {
if let Some(s) = slot.borrow().as_ref() {
s.frames.lock().push(fp);
} else {
debug_assert!(
false,
"push_thread_frame called without initialized thread slot"
);
}
});
}
#[cfg(all(not(unix), feature = "threading"))]
pub fn pop_thread_frame() {
CURRENT_THREAD_SLOT.with(|slot| {
if let Some(s) = slot.borrow().as_ref() {
s.frames.lock().pop();
} else {
debug_assert!(
false,
"pop_thread_frame called without initialized thread slot"
);
}
});
}
#[must_use]
#[allow(clippy::not_unsafe_ptr_arg_deref)]
pub fn set_current_frame(frame: *const InterpreterFrame) -> *const InterpreterFrame {
FRAME_SLOT_CACHE.with(|cache| {
#[cfg(feature = "threading")]
{
let slot = cache.top_iframe.get();
if !slot.is_null() {
unsafe { &*slot }.store(frame as usize, Ordering::Relaxed);
}
#[cfg(unix)]
{
let slot = cache.top_frame.get();
if !slot.is_null() {
let fo_ptr = if frame.is_null() {
core::ptr::null_mut()
} else {
let frame_obj = unsafe { (*frame).frame_obj() };
frame_obj.map_or(core::ptr::null_mut(), |py| {
py as *const Py<FrameObject> as *mut Py<FrameObject>
})
};
unsafe { &*slot }.store(fo_ptr, Ordering::Relaxed);
}
}
}
cache.current_frame.swap(frame as usize, Ordering::Relaxed)
}) as *const InterpreterFrame
}
#[inline(always)]
#[must_use]
pub fn set_current_frame_nosave(frame: *const InterpreterFrame) -> *const InterpreterFrame {
FRAME_SLOT_CACHE.with(|cache| cache.current_frame.swap(frame as usize, Ordering::Relaxed))
as *const InterpreterFrame
}
#[must_use]
pub fn get_current_frame() -> *const InterpreterFrame {
FRAME_SLOT_CACHE.with(|cache| cache.current_frame.load(Ordering::Relaxed))
as *const InterpreterFrame
}
#[cfg(feature = "threading")]
pub fn update_thread_exception(exc: Option<PyBaseExceptionRef>) {
CURRENT_THREAD_SLOT.with(|slot| {
if let Some(s) = slot.borrow().as_ref() {
let _old = unsafe { s.exception.swap(exc) };
}
});
}
#[cfg(feature = "threading")]
pub fn get_all_current_exceptions(vm: &VirtualMachine) -> Vec<(u64, Option<PyBaseExceptionRef>)> {
let registry = vm.state.thread_frames.lock();
registry
.iter()
.map(|(id, slot)| (*id, slot.exception.load_owned()))
.collect()
}
#[cfg(feature = "threading")]
pub fn cleanup_current_thread_frames(vm: &VirtualMachine) {
let thread_id = crate::stdlib::_thread::get_ident();
let interp_id = vm.state.interpreter_id;
let slot_for_interp = INTERP_THREAD_SLOTS.with(|slots| slots.borrow_mut().remove(&interp_id));
let current_slot = CURRENT_THREAD_SLOT.with(|slot| slot.borrow().as_ref().cloned());
let slot_to_clean = slot_for_interp.or(current_slot);
if let Some(slot) = &slot_to_clean {
let _ = slot.state.compare_exchange(
ThreadState::Attached as i32,
ThreadState::Detached as i32,
Ordering::AcqRel,
Ordering::Acquire,
);
}
let _removed = if let Some(slot) = &slot_to_clean {
let mut registry = vm.state.thread_frames.lock();
match registry.get(&thread_id) {
Some(registered) if Arc::ptr_eq(registered, slot) => registry.remove(&thread_id),
_ => None,
}
} else {
None
};
if let Some(slot) = &_removed
&& vm.state.stop_the_world.requested.load(Ordering::Acquire)
&& thread_id != vm.state.stop_the_world.requester_ident()
&& slot.state.load(Ordering::Relaxed) != ThreadState::Suspended as i32
{
vm.state.stop_the_world.notify_thread_gone();
}
CURRENT_THREAD_SLOT.with(|s| {
let clear = match (s.borrow().as_ref(), slot_to_clean.as_ref()) {
(Some(cur), Some(cleaned)) => Arc::ptr_eq(cur, cleaned),
(Some(_), None) => false,
(None, _) => false,
};
if clear {
*s.borrow_mut() = None;
#[cfg(feature = "threading")]
FRAME_SLOT_CACHE.with(|cache| {
#[cfg(unix)]
cache.top_frame.set(core::ptr::null());
cache.top_iframe.set(core::ptr::null());
});
#[cfg(feature = "threading")]
CURRENT_STOP_REQUESTED.with(|c| c.set(core::ptr::null()));
}
});
}
#[cfg(feature = "threading")]
pub fn reinit_frame_slot_after_fork(vm: &VirtualMachine) {
let current_ident = crate::stdlib::_thread::get_ident();
#[cfg(not(unix))]
let current_frames: Vec<FramePtr> = {
let mut current_frames = Vec::new();
let mut cur = get_current_frame();
while !cur.is_null() {
let iframe = unsafe { &*cur };
if let Some(fo) = iframe.frame_obj() {
current_frames.push(FramePtr(unsafe {
NonNull::new_unchecked(fo as *const _ as *mut _)
}));
}
cur = iframe.previous.load(Ordering::Relaxed) as *const InterpreterFrame;
}
current_frames.reverse();
current_frames
};
#[cfg(unix)]
let top_fo_ptr = {
let top_iframe = get_current_frame();
if top_iframe.is_null() {
core::ptr::null_mut()
} else {
match unsafe { (*top_iframe).frame_obj() } {
Some(fo) => fo as *const Py<FrameObject> as *mut Py<FrameObject>,
None => core::ptr::null_mut(),
}
}
};
let top_iframe_ptr = get_current_frame() as usize;
let new_slot = Arc::new(ThreadSlot {
#[cfg(unix)]
top_frame: AtomicPtr::new(top_fo_ptr),
top_iframe: AtomicUsize::new(top_iframe_ptr),
#[cfg(not(unix))]
frames: parking_lot::Mutex::new(current_frames),
exception: crate::PyAtomicRef::from(vm.topmost_exception()),
trace_func: PyMutex::new(vm.trace_func.borrow().clone()),
profile_func: PyMutex::new(vm.profile_func.borrow().clone()),
state: core::sync::atomic::AtomicI32::new(ThreadState::Attached as i32),
stop_requested: core::sync::atomic::AtomicBool::new(false),
thread: std::thread::current(),
qsbr: crate::object::qsbr::QSBR.register(),
});
FRAME_SLOT_CACHE.with(|cache| {
#[cfg(unix)]
cache.top_frame.set(&new_slot.top_frame);
cache.top_iframe.set(&new_slot.top_iframe);
});
#[cfg(feature = "threading")]
CURRENT_STOP_REQUESTED.with(|c| c.set(&new_slot.stop_requested));
let mut registry = vm.state.thread_frames.lock();
registry.clear();
registry.insert(current_ident, new_slot.clone());
drop(registry);
CURRENT_THREAD_SLOT.with(|s| {
*s.borrow_mut() = Some(new_slot.clone());
});
INTERP_THREAD_SLOTS.with(|slots| {
slots.borrow_mut().insert(vm.state.interpreter_id, new_slot);
});
}
#[cfg(feature = "threading")]
pub fn purge_other_interpreter_slots_after_fork(keep_id: i64) {
INTERP_THREAD_SLOTS.with(|slots| {
slots.borrow_mut().retain(|&id, _| id == keep_id);
});
}
#[cfg(feature = "threading")]
fn top_slot_is_attached() -> bool {
current_slot_is_attached()
}
#[cfg(not(feature = "threading"))]
fn top_slot_is_attached() -> bool {
true
}
enum WithVmTarget {
AlreadyCurrent(NonNull<VirtualMachine>),
NeedsSwitch(NonNull<VirtualMachine>),
}
pub fn with_vm<F, R>(obj: &PyObject, f: F) -> Option<R>
where
F: Fn(&VirtualMachine) -> R,
{
let vm_owns_obj = |interp: NonNull<VirtualMachine>| {
let vm = unsafe { interp.as_ref() };
obj.fast_isinstance(vm.ctx.types.object_type)
};
let target = VM_STACK.with(|vms| {
let vms = vms.borrow();
if let Some(top) = vms.last().copied()
&& vm_owns_obj(top)
&& top_slot_is_attached()
{
return Some(WithVmTarget::AlreadyCurrent(top));
}
let interp = match vms.iter().copied().exactly_one() {
Ok(x) => {
debug_assert!(vm_owns_obj(x));
x
}
Err(mut others) => others.find(|x| vm_owns_obj(*x))?,
};
Some(WithVmTarget::NeedsSwitch(interp))
})?;
match target {
WithVmTarget::AlreadyCurrent(interp) => {
let vm = unsafe { interp.as_ref() };
Some(f(vm))
}
WithVmTarget::NeedsSwitch(interp) => {
let vm = unsafe { interp.as_ref() };
Some(set_current_vm(vm, || f(vm)))
}
}
}
#[must_use = "ThreadedVirtualMachine does nothing unless you move it to another thread and call .run()"]
#[cfg(feature = "threading")]
pub struct ThreadedVirtualMachine {
pub(super) vm: VirtualMachine,
}
#[cfg(feature = "threading")]
impl ThreadedVirtualMachine {
pub fn make_spawn_func<F, R>(self, f: F) -> impl FnOnce() -> R
where
F: FnOnce(&VirtualMachine) -> R,
{
move || self.run(f)
}
pub fn run<F, R>(&self, f: F) -> R
where
F: FnOnce(&VirtualMachine) -> R,
{
let vm = &self.vm;
vm.c_stack_soft_limit
.set(VirtualMachine::calculate_c_stack_soft_limit());
enter_vm(vm, || f(vm))
}
}
impl VirtualMachine {
#[cfg(feature = "threading")]
pub fn start_thread<F, R>(&self, f: F) -> std::thread::JoinHandle<R>
where
F: Send + 'static + FnOnce(&Self) -> R,
R: Send + 'static,
{
let func = self.new_thread().make_spawn_func(f);
std::thread::spawn(func)
}
#[cfg(feature = "threading")]
pub fn new_thread(&self) -> ThreadedVirtualMachine {
let global_trace = self.state.global_trace_func.lock().clone();
let global_profile = self.state.global_profile_func.lock().clone();
let use_tracing = global_trace.is_some() || global_profile.is_some();
let vm = Self {
builtins: self.builtins.clone(),
sys_module: self.sys_module.clone(),
ctx: self.ctx.clone(),
datastack: core::cell::UnsafeCell::new(crate::datastack::DataStack::new()),
wasm_id: self.wasm_id.clone(),
exceptions: RefCell::default(),
import_func: self.import_func.clone(),
importlib: self.importlib.clone(),
profile_func: RefCell::new(global_profile.unwrap_or_else(|| self.ctx.none())),
trace_func: RefCell::new(global_trace.unwrap_or_else(|| self.ctx.none())),
use_tracing: Cell::new(use_tracing),
what_event: Cell::new(None),
tracing_depth: Cell::new(0),
recursion_limit: self.recursion_limit.clone(),
signal_handlers: core::cell::OnceCell::new(),
signal_rx: None,
repr_guards: RefCell::default(),
state: self.state.clone(),
initialized: self.initialized,
recursion_depth: Cell::new(0),
#[cfg(any(miri, target_env = "musl"))]
native_recursion_depth: Cell::new(0),
c_stack_soft_limit: Cell::new(Self::calculate_c_stack_soft_limit()),
async_gen_firstiter: RefCell::new(None),
async_gen_finalizer: RefCell::new(None),
asyncio_running_loop: RefCell::new(None),
asyncio_running_task: RefCell::new(None),
context_stack: RefCell::default(),
callable_cache: self.callable_cache.clone(),
pending_tailcall_frame: Cell::new(None),
pending_tailcall_owner: core::cell::UnsafeCell::new(None),
pending_gen_resume: core::cell::UnsafeCell::new(None),
trampoline_stack: core::cell::UnsafeCell::new(Vec::new()),
};
ThreadedVirtualMachine { vm }
}
}