use crate::common::linked_list::LinkedList;
use crate::common::lock::{PyMutex, PyRwLock};
use crate::object::{GC_NO_OWNER, GC_PERMANENT, GC_REACHABLE, GC_UNTRACKED, GcLink, GcOwner};
use crate::{AsObject, PyObject, PyObjectRef};
use core::ptr::NonNull;
use core::sync::atomic::{AtomicBool, AtomicU16, AtomicU32, AtomicUsize, Ordering};
fn elapsed_secs(
#[cfg(target_arch = "wasm32")] _start: (),
#[cfg(not(target_arch = "wasm32"))] start: std::time::Instant,
) -> f64 {
cfg_select! {
target_arch = "wasm32" => 0.0,
_ => start.elapsed().as_secs_f64(),
}
}
bitflags::bitflags! {
#[derive(Copy, Clone, Debug, Default, PartialEq, Eq)]
pub struct GcDebugFlags: u32 {
const STATS = 1 << 0;
const COLLECTABLE = 1 << 1;
const UNCOLLECTABLE = 1 << 2;
const SAVEALL = 1 << 5;
const LEAK = Self::COLLECTABLE.bits() | Self::UNCOLLECTABLE.bits() | Self::SAVEALL.bits();
}
}
#[derive(Clone, Copy, Debug, Default)]
pub struct CollectResult {
pub collected: usize,
pub uncollectable: usize,
pub candidates: usize,
pub duration: f64,
}
#[derive(Clone, Copy, Debug, Default)]
pub struct GcStats {
pub collections: usize,
pub collected: usize,
pub uncollectable: usize,
pub candidates: usize,
pub duration: f64,
}
pub struct GcGeneration {
threshold: AtomicU32,
stats: PyMutex<GcStats>,
}
impl GcGeneration {
#[must_use]
pub const fn new(threshold: u32) -> Self {
Self {
threshold: AtomicU32::new(threshold),
stats: PyMutex::new(GcStats {
collections: 0,
collected: 0,
uncollectable: 0,
candidates: 0,
duration: 0.0,
}),
}
}
pub fn threshold(&self) -> u32 {
self.threshold.load(Ordering::Relaxed)
}
pub fn set_threshold(&self, value: u32) {
self.threshold.store(value, Ordering::Relaxed);
}
pub fn stats(&self) -> GcStats {
let guard = self.stats.lock();
GcStats {
collections: guard.collections,
collected: guard.collected,
uncollectable: guard.uncollectable,
candidates: guard.candidates,
duration: guard.duration,
}
}
pub fn update_stats(
&self,
collected: usize,
uncollectable: usize,
candidates: usize,
duration: f64,
) {
let mut guard = self.stats.lock();
guard.collections += 1;
guard.collected += collected;
guard.uncollectable += uncollectable;
guard.candidates += candidates;
guard.duration += duration;
}
#[cfg(all(unix, feature = "threading"))]
unsafe fn reinit_stats_after_fork(&self) {
unsafe { crate::common::lock::reinit_mutex_after_fork(&self.stats) };
}
}
fn release_count(count: &AtomicUsize) {
let _ = count.try_update(Ordering::Relaxed, Ordering::Relaxed, |n| n.checked_sub(1));
}
fn is_owned_by(obj: &PyObject, owner: GcOwner) -> bool {
let obj_owner = obj.gc_owner();
obj_owner == owner || obj_owner == GC_NO_OWNER
}
#[derive(Clone, Copy, PartialEq, Eq, Hash)]
struct GcPtr(NonNull<PyObject>);
#[derive(Default)]
struct GcPtrHasher(u64);
impl core::hash::Hasher for GcPtrHasher {
fn finish(&self) -> u64 {
self.0
}
fn write_usize(&mut self, value: usize) {
let mut z = (value as u64).wrapping_add(0x9E37_79B9_7F4A_7C15);
z = (z ^ (z >> 30)).wrapping_mul(0xBF58_476D_1CE4_E5B9);
z = (z ^ (z >> 27)).wrapping_mul(0x94D0_49BB_1331_11EB);
self.0 = z ^ (z >> 31);
}
fn write(&mut self, bytes: &[u8]) {
for &byte in bytes {
self.0 = (self.0 ^ u64::from(byte)).wrapping_mul(0x0100_0000_01B3);
}
}
}
type GcBuildHasher = core::hash::BuildHasherDefault<GcPtrHasher>;
type GcSet<T> = std::collections::HashSet<T, GcBuildHasher>;
type GcMap<K, V> = std::collections::HashMap<K, V, GcBuildHasher>;
#[cfg(feature = "threading")]
struct CollectStopTheWorld {
stopped: Vec<crate::common::rc::PyRc<crate::vm::PyGlobalState>>,
admission: Option<parking_lot::MutexGuard<'static, ()>>,
restarted: bool,
}
#[cfg(feature = "threading")]
impl CollectStopTheWorld {
fn new() -> Self {
if !crate::vm::thread::current_vm_is_set() {
return Self {
stopped: Vec::new(),
admission: None,
restarted: true,
};
}
let mut guard = Self {
stopped: Vec::new(),
admission: Some(crate::vm::runtime::lock_admission_for_stop()),
restarted: false,
};
for state in crate::vm::runtime::live_interpreter_states() {
state.stop_the_world.stop_the_world(&state);
guard.stopped.push(state);
}
guard
}
fn restart(&mut self) {
if self.restarted {
return;
}
self.restarted = true;
for state in self.stopped.iter().rev() {
state.stop_the_world.start_the_world(state);
}
self.admission = None;
}
#[cfg(all(unix, debug_assertions))]
fn is_stopped(&self) -> bool {
!self.stopped.is_empty()
}
}
#[cfg(feature = "threading")]
impl Drop for CollectStopTheWorld {
fn drop(&mut self) {
self.restart();
}
}
pub struct GcState {
generation_lists: [PyRwLock<LinkedList<GcLink>>; 3],
permanent_list: PyRwLock<LinkedList<GcLink>>,
counts: [AtomicUsize; 3],
permanent_count: AtomicUsize,
collecting: PyMutex<()>,
next_owner: AtomicU16,
retired: PyMutex<Vec<GcOwner>>,
}
#[cfg(feature = "threading")]
unsafe impl Send for GcState {}
#[cfg(feature = "threading")]
unsafe impl Sync for GcState {}
impl Default for GcState {
fn default() -> Self {
Self::new()
}
}
impl GcState {
#[must_use]
pub const fn new() -> Self {
Self {
generation_lists: [
PyRwLock::new(LinkedList::new()),
PyRwLock::new(LinkedList::new()),
PyRwLock::new(LinkedList::new()),
],
permanent_list: PyRwLock::new(LinkedList::new()),
counts: [
AtomicUsize::new(0),
AtomicUsize::new(0),
AtomicUsize::new(0),
],
permanent_count: AtomicUsize::new(0),
collecting: PyMutex::new(()),
next_owner: AtomicU16::new(GC_NO_OWNER + 1),
retired: PyMutex::new(Vec::new()),
}
}
fn alloc_owner(&self) -> GcOwner {
self.next_owner
.try_update(Ordering::Relaxed, Ordering::Relaxed, |next| {
next.checked_add(1)
})
.unwrap_or(GC_NO_OWNER)
}
fn retire_owner(&self, owner: GcOwner) {
if owner == GC_NO_OWNER {
return;
}
self.retired.lock().push(owner);
}
pub fn get_count(&self) -> (usize, usize, usize) {
(
self.counts[0].load(Ordering::Relaxed),
self.counts[1].load(Ordering::Relaxed),
self.counts[2].load(Ordering::Relaxed),
)
}
pub unsafe fn track_object(&self, obj: NonNull<PyObject>, owner: GcOwner) {
let obj_ref = unsafe { obj.as_ref() };
obj_ref.set_gc_tracked();
obj_ref.set_gc_generation(0);
obj_ref.set_gc_owner(owner);
self.generation_lists[0].write().push_front(obj);
self.counts[0].fetch_add(1, Ordering::Relaxed);
}
unsafe fn track_object_fresh(&self, obj: NonNull<PyObject>, owner: GcOwner) {
let obj_ref = unsafe { obj.as_ref() };
obj_ref.init_gc_tracked_bit();
obj_ref.set_gc_generation(0);
obj_ref.set_gc_owner(owner);
self.generation_lists[0].write().push_front(obj);
self.counts[0].fetch_add(1, Ordering::Relaxed);
}
unsafe fn track_pair_fresh(&self, a: NonNull<PyObject>, b: NonNull<PyObject>, owner: GcOwner) {
debug_assert_ne!(a, b);
for obj in [a, b] {
let obj_ref = unsafe { obj.as_ref() };
obj_ref.init_gc_tracked_bit();
obj_ref.set_gc_generation(0);
obj_ref.set_gc_owner(owner);
}
{
let mut list = self.generation_lists[0].write();
list.push_front(a);
list.push_front(b);
}
self.counts[0].fetch_add(2, Ordering::Relaxed);
}
pub unsafe fn untrack_object(&self, obj: NonNull<PyObject>) {
let obj_ref = unsafe { obj.as_ref() };
loop {
let obj_gen = obj_ref.gc_generation();
let (list_lock, count) = if obj_gen <= 2 {
(
&self.generation_lists[obj_gen as usize] as &PyRwLock<LinkedList<GcLink>>,
&self.counts[obj_gen as usize],
)
} else if obj_gen == GC_PERMANENT {
(&self.permanent_list, &self.permanent_count)
} else {
return; };
let mut list = list_lock.write();
if obj_ref.gc_generation() != obj_gen {
drop(list);
continue; }
if unsafe { list.remove(obj) }.is_some() {
release_count(count);
obj_ref.clear_gc_tracked();
obj_ref.set_gc_generation(GC_UNTRACKED);
} else {
eprintln!(
"GC WARNING: untrack_object failed to remove obj={obj:p} from gen={obj_gen}, \
tracked={}, gc_gen={}",
obj_ref.is_gc_tracked(),
obj_ref.gc_generation()
);
}
return;
}
}
pub fn get_objects(&self, generation: Option<i32>, owner: GcOwner) -> Vec<PyObjectRef> {
fn collect_from_list(
list: &LinkedList<GcLink>,
owner: GcOwner,
) -> impl Iterator<Item = PyObjectRef> + '_ {
list.iter()
.filter(move |obj| is_owned_by(obj, owner))
.filter_map(|obj| obj.try_to_owned())
}
match generation {
None => {
let mut result = Vec::new();
for gen_list in &self.generation_lists {
result.extend(collect_from_list(&gen_list.read(), owner));
}
result.extend(collect_from_list(&self.permanent_list.read(), owner));
result
}
Some(g) if (0..=2).contains(&g) => {
let guard = self.generation_lists[g as usize].read();
collect_from_list(&guard, owner).collect()
}
_ => Vec::new(),
}
}
fn maybe_collect(&self, gc: &GcInterpreterState) -> bool {
if !gc.is_enabled() {
return false;
}
let count0 = self.counts[0].load(Ordering::Relaxed) as u32;
let threshold0 = gc.generations[0].threshold();
if threshold0 > 0 && count0 >= threshold0 {
#[cfg(feature = "threading")]
{
gc.scheduled.store(true, Ordering::Relaxed);
return false;
}
#[cfg(not(feature = "threading"))]
{
self.collect_inner(gc, 0, false);
return true;
}
}
false
}
fn collect_inner(
&self,
gc: &GcInterpreterState,
generation: usize,
force: bool,
) -> CollectResult {
if !force && !gc.is_enabled() {
#[cfg(feature = "threading")]
gc.scheduled.store(false, Ordering::Relaxed);
return CollectResult::default();
}
let Some(_guard) = self.collecting.try_lock() else {
return CollectResult::default();
};
#[cfg(feature = "threading")]
gc.scheduled.store(false, Ordering::Relaxed);
let start_time = cfg_select! {
target_arch = "wasm32" => (),
_ => std::time::Instant::now(),
};
core::sync::atomic::fence(Ordering::SeqCst);
let generation = generation.min(2);
let debug = gc.get_debug();
crate::builtins::type_::type_cache_clear();
#[cfg(feature = "threading")]
crate::object::qsbr::QSBR.process();
#[cfg(feature = "threading")]
let mut stw = CollectStopTheWorld::new();
let gen_locks: Vec<_> = (0..=generation)
.map(|i| self.generation_lists[i].read())
.collect();
let owner = gc.owner;
let retired = {
let mut retired = self.retired.lock().clone();
retired.sort_unstable();
retired
};
let mut candidate_ptrs: Vec<GcPtr> = Vec::new();
for gen_list in &gen_locks {
for obj in gen_list.iter() {
if retired.binary_search(&obj.gc_owner()).is_ok() {
obj.set_gc_owner(GC_NO_OWNER);
}
let strong_count = obj.strong_count();
if strong_count > 0 && is_owned_by(obj, owner) && !obj.is_gc_collecting() {
obj.start_gc_refs(strong_count);
candidate_ptrs.push(GcPtr(NonNull::from(obj)));
}
}
}
if generation == 2 && !retired.is_empty() {
for obj in self.permanent_list.read().iter() {
if retired.binary_search(&obj.gc_owner()).is_ok() {
obj.set_gc_owner(GC_NO_OWNER);
}
}
self.retired
.lock()
.retain(|tag| retired.binary_search(tag).is_err());
}
if candidate_ptrs.is_empty() {
let reset_end = if generation >= 2 { 2 } else { generation + 1 };
for count in self.counts.iter().take(reset_end) {
count.store(0, Ordering::Relaxed);
}
let duration = elapsed_secs(start_time);
gc.generations[generation].update_stats(0, 0, 0, duration);
return CollectResult {
collected: 0,
uncollectable: 0,
candidates: 0,
duration,
};
}
let candidates = candidate_ptrs.len();
if debug.contains(GcDebugFlags::STATS) {
eprintln!("gc: collecting {candidates} objects from generations 0..={generation}");
}
let mut referent_ptrs: Vec<NonNull<PyObject>> = Vec::new();
let mut referent_ranges: GcMap<GcPtr, (usize, usize)> = GcMap::default();
for &ptr in &candidate_ptrs {
let obj = unsafe { ptr.0.as_ref() };
if obj.strong_count() == 0 {
continue;
}
let start = referent_ptrs.len();
unsafe { obj.gc_extend_referent_ptrs(&mut referent_ptrs) };
let end = referent_ptrs.len();
for &child_ptr in &referent_ptrs[start..end] {
let child = unsafe { child_ptr.as_ref() };
if child.is_gc_collecting() {
child.subtract_gc_ref();
}
}
referent_ranges.insert(ptr, (start, end));
}
let mut worklist: Vec<GcPtr> = Vec::new();
for &ptr in &candidate_ptrs {
let obj = unsafe { ptr.0.as_ref() };
if obj.gc_refs() > 0 {
obj.mark_gc_reachable();
worklist.push(ptr);
}
}
while let Some(ptr) = worklist.pop() {
let obj = unsafe { ptr.0.as_ref() };
if obj.is_gc_tracked() {
let computed;
let children: &[NonNull<PyObject>] = match referent_ranges.get(&ptr) {
Some(&(start, end)) => &referent_ptrs[start..end],
None => {
computed = unsafe { obj.gc_get_referent_ptrs() };
&computed
}
};
for &child_ptr in children {
let child = unsafe { child_ptr.as_ref() };
if child.is_gc_collecting() && child.mark_gc_reachable() {
worklist.push(GcPtr(child_ptr));
}
}
}
}
let mut reachable: Vec<GcPtr> = Vec::new();
let mut unreachable: Vec<GcPtr> = Vec::new();
for &ptr in &candidate_ptrs {
let obj = unsafe { ptr.0.as_ref() };
if obj.gc_refs() == GC_REACHABLE {
reachable.push(ptr);
} else {
unreachable.push(ptr);
}
obj.end_gc_refs();
}
#[cfg(all(unix, feature = "threading", debug_assertions))]
if stw.is_stopped() {
let unreachable_set: GcSet<GcPtr> = unreachable.iter().copied().collect();
let mut cur = crate::vm::thread::get_current_frame();
while !cur.is_null() {
let iframe = unsafe { &*cur };
if let Some(fo) = iframe.frame_obj() {
let obj = fo.as_object();
let ptr = GcPtr(NonNull::from(obj));
debug_assert!(
!unreachable_set.contains(&ptr),
"running frame {obj:p} classified unreachable during GC"
);
}
cur = iframe.previous();
}
}
if debug.contains(GcDebugFlags::STATS) {
eprintln!(
"gc: {} reachable, {} unreachable",
reachable.len(),
unreachable.len()
);
}
let survivor_refs: Vec<PyObjectRef> = reachable
.iter()
.filter_map(|ptr| {
let obj = unsafe { ptr.0.as_ref() };
obj.try_to_owned()
})
.collect();
let unreachable_refs: Vec<crate::PyObjectRef> = unreachable
.iter()
.filter_map(|ptr| {
let obj = unsafe { ptr.0.as_ref() };
obj.try_to_owned()
})
.collect();
#[cfg(feature = "threading")]
stw.restart();
if unreachable.is_empty() {
drop(gen_locks);
self.promote_survivors(generation, &survivor_refs);
let reset_end = if generation >= 2 { 2 } else { generation + 1 };
for count in self.counts.iter().take(reset_end) {
count.store(0, Ordering::Relaxed);
}
let duration = elapsed_secs(start_time);
gc.generations[generation].update_stats(0, 0, candidates, duration);
return CollectResult {
collected: 0,
uncollectable: 0,
candidates,
duration,
};
}
drop(gen_locks);
if unreachable_refs.is_empty() {
self.promote_survivors(generation, &survivor_refs);
let reset_end = if generation >= 2 { 2 } else { generation + 1 };
for count in self.counts.iter().take(reset_end) {
count.store(0, Ordering::Relaxed);
}
let duration = elapsed_secs(start_time);
gc.generations[generation].update_stats(0, 0, candidates, duration);
return CollectResult {
collected: 0,
uncollectable: 0,
candidates,
duration,
};
}
let initial_counts: GcMap<GcPtr, usize> = unreachable_refs
.iter()
.map(|obj| {
let ptr = GcPtr(core::ptr::NonNull::from(obj.as_ref()));
(ptr, obj.strong_count())
})
.collect();
let mut all_callbacks: Vec<(crate::PyRef<crate::object::PyWeak>, crate::PyObjectRef)> =
Vec::new();
for obj_ref in &unreachable_refs {
let callbacks = obj_ref.gc_clear_weakrefs_collect_callbacks();
all_callbacks.extend(callbacks);
}
for (wr, cb) in all_callbacks {
if let Some(Err(e)) = crate::vm::thread::with_vm(&cb, |vm| cb.call((wr.clone(),), vm)) {
crate::vm::thread::with_vm(&cb, |vm| {
vm.run_unraisable(e.clone(), Some("weakref callback".to_owned()), cb.clone());
});
}
}
for obj_ref in &unreachable_refs {
obj_ref.try_call_finalizer();
}
let mut resurrected_set: GcSet<GcPtr> = GcSet::default();
let unreachable_set: GcSet<GcPtr> = unreachable.iter().copied().collect();
for obj in &unreachable_refs {
let ptr = GcPtr(core::ptr::NonNull::from(obj.as_ref()));
let initial = initial_counts.get(&ptr).copied().unwrap_or(1);
if obj.strong_count() > initial {
resurrected_set.insert(ptr);
}
}
let mut worklist: Vec<GcPtr> = resurrected_set.iter().copied().collect();
while let Some(ptr) = worklist.pop() {
let obj = unsafe { ptr.0.as_ref() };
let referent_ptrs = unsafe { obj.gc_get_referent_ptrs() };
for child_ptr in referent_ptrs {
let child_gc_ptr = GcPtr(child_ptr);
if unreachable_set.contains(&child_gc_ptr) && resurrected_set.insert(child_gc_ptr) {
worklist.push(child_gc_ptr);
}
}
}
let (resurrected, truly_dead): (Vec<_>, Vec<_>) =
unreachable_refs.into_iter().partition(|obj| {
let ptr = GcPtr(core::ptr::NonNull::from(obj.as_ref()));
resurrected_set.contains(&ptr)
});
if debug.contains(GcDebugFlags::STATS) {
eprintln!(
"gc: {} resurrected, {} truly dead",
resurrected.len(),
truly_dead.len()
);
}
let collected = {
let dead_ptrs: GcSet<usize> = truly_dead
.iter()
.map(|obj| obj.as_ref() as *const PyObject as usize)
.collect();
let instance_dict_count = truly_dead
.iter()
.filter(|obj| {
if let Some(dict_ref) = obj.dict() {
dead_ptrs.contains(&(dict_ref.as_object() as *const PyObject as usize))
} else {
false
}
})
.count();
truly_dead.len() - instance_dict_count
};
self.promote_survivors(generation, &survivor_refs);
drop(survivor_refs);
self.promote_survivors(generation, &resurrected);
drop(resurrected);
if debug.contains(GcDebugFlags::COLLECTABLE) {
for obj in &truly_dead {
eprintln!(
"gc: collectable <{} {:p}>",
obj.class().name(),
obj.as_ref()
);
}
}
if debug.contains(GcDebugFlags::SAVEALL) {
self.promote_survivors(generation, &truly_dead);
let mut garbage_guard = gc.garbage.lock();
for obj_ref in &truly_dead {
garbage_guard.push(obj_ref.clone());
}
}
if !truly_dead.is_empty() {
let save_all = debug.contains(GcDebugFlags::SAVEALL);
let mut late_resurrected: GcSet<GcPtr> = GcSet::default();
if !save_all {
let mut expected_counts: GcMap<GcPtr, usize> = GcMap::default();
for obj_ref in &truly_dead {
let obj = obj_ref.as_ref();
if obj.is_gc_tracked() {
unsafe { self.untrack_object(NonNull::from(obj)) };
}
expected_counts.insert(GcPtr(NonNull::from(obj)), 1);
}
let mut referents: GcMap<GcPtr, Vec<NonNull<PyObject>>> = GcMap::default();
for obj_ref in &truly_dead {
let referent_ptrs = unsafe { obj_ref.gc_get_referent_ptrs() };
for child_ptr in &referent_ptrs {
if let Some(n) = expected_counts.get_mut(&GcPtr(*child_ptr)) {
*n += 1;
}
}
referents.insert(GcPtr(NonNull::from(obj_ref.as_ref())), referent_ptrs);
}
let mut worklist: Vec<GcPtr> = Vec::new();
for obj_ref in &truly_dead {
let ptr = GcPtr(NonNull::from(obj_ref.as_ref()));
if obj_ref.strong_count() > expected_counts[&ptr]
&& late_resurrected.insert(ptr)
{
worklist.push(ptr);
}
}
while let Some(ptr) = worklist.pop() {
let Some(referent_ptrs) = referents.get(&ptr) else {
continue;
};
for child_ptr in referent_ptrs {
let child = GcPtr(*child_ptr);
if expected_counts.contains_key(&child) && late_resurrected.insert(child) {
worklist.push(child);
}
}
}
#[expect(
clippy::iter_over_hash_type,
reason = "Iteration order doesn't matter here"
)]
for &ptr in &late_resurrected {
let owner = unsafe { ptr.0.as_ref() }.gc_owner();
unsafe { self.track_object(ptr.0, owner) };
}
}
rustpython_common::refcount::with_deferred_drops(|| {
if !save_all {
for obj_ref in &truly_dead {
let obj = obj_ref.as_ref();
if late_resurrected.contains(&GcPtr(NonNull::from(obj))) {
continue;
}
if obj.gc_has_clear() {
let edges = unsafe { obj.gc_clear() };
drop(edges);
}
}
}
drop(truly_dead);
});
}
let reset_end = if generation >= 2 { 2 } else { generation + 1 };
for count in self.counts.iter().take(reset_end) {
count.store(0, Ordering::Relaxed);
}
let duration = elapsed_secs(start_time);
gc.generations[generation].update_stats(collected, 0, candidates, duration);
CollectResult {
collected,
uncollectable: 0,
candidates,
duration,
}
}
fn promote_survivors(&self, from_gen: usize, survivors: &[PyObjectRef]) {
let next_gen = (from_gen + 1).min(2);
for batch in survivors.chunks(256) {
for src_gen in 0..next_gen {
let mut src = self.generation_lists[src_gen].write();
let mut dst = self.generation_lists[next_gen].write();
let mut promoted = 0;
for obj_ref in batch {
let obj = obj_ref.as_ref();
if obj.gc_generation() as usize != src_gen || !obj.is_gc_tracked() {
continue;
}
let ptr = NonNull::from(obj);
if unsafe { src.remove(ptr) }.is_some() {
dst.push_front(ptr);
obj.set_gc_generation(next_gen as u8);
promoted += 1;
}
}
if promoted != 0 {
let _ = self.counts[src_gen].try_update(
Ordering::Relaxed,
Ordering::Relaxed,
|count| Some(count.saturating_sub(promoted)),
);
self.counts[next_gen].fetch_add(promoted, Ordering::Relaxed);
}
}
}
}
pub fn get_freeze_count(&self) -> usize {
self.permanent_count.load(Ordering::Relaxed)
}
fn freeze(&self, owner: GcOwner) {
let mut count = 0usize;
for (gen_idx, gen_list) in self.generation_lists.iter().enumerate() {
let mut list = gen_list.write();
let mut perm = self.permanent_list.write();
let moving: Vec<_> = list
.iter()
.filter(|obj| is_owned_by(obj, owner))
.map(NonNull::from)
.collect();
for ptr in moving {
if unsafe { list.remove(ptr) }.is_none() {
continue;
}
perm.push_front(ptr);
unsafe { ptr.as_ref().set_gc_generation(GC_PERMANENT) };
count += 1;
release_count(&self.counts[gen_idx]);
}
}
self.permanent_count.fetch_add(count, Ordering::Relaxed);
}
fn unfreeze(&self, owner: GcOwner) {
let mut count = 0usize;
{
let mut gen2 = self.generation_lists[2].write();
let mut perm_list = self.permanent_list.write();
let moving: Vec<_> = perm_list
.iter()
.filter(|obj| is_owned_by(obj, owner))
.map(NonNull::from)
.collect();
for ptr in moving {
if unsafe { perm_list.remove(ptr) }.is_none() {
continue;
}
gen2.push_front(ptr);
unsafe { ptr.as_ref().set_gc_generation(2) };
count += 1;
}
let _ = self.permanent_count.try_update(
Ordering::Relaxed,
Ordering::Relaxed,
|permanent| Some(permanent.saturating_sub(count)),
);
}
self.counts[2].fetch_add(count, Ordering::Relaxed);
}
#[cfg(all(unix, feature = "threading"))]
pub unsafe fn reinit_after_fork(&self) {
use crate::common::lock::{reinit_mutex_after_fork, reinit_rwlock_after_fork};
unsafe {
reinit_mutex_after_fork(&self.collecting);
reinit_mutex_after_fork(&self.retired);
for rw in &self.generation_lists {
reinit_rwlock_after_fork(rw);
}
reinit_rwlock_after_fork(&self.permanent_list);
}
}
}
pub struct GcInterpreterState {
owner: GcOwner,
pub generations: [GcGeneration; 3],
enabled: AtomicBool,
#[cfg(feature = "threading")]
scheduled: AtomicBool,
debug: AtomicU32,
pub garbage: PyMutex<Vec<PyObjectRef>>,
pub py_garbage: crate::builtins::PyListRef,
pub py_callbacks: crate::builtins::PyListRef,
}
impl GcInterpreterState {
pub fn new(ctx: &crate::vm::Context) -> Self {
Self {
owner: gc_state().alloc_owner(),
generations: [
GcGeneration::new(2000), GcGeneration::new(10), GcGeneration::new(0), ],
enabled: AtomicBool::new(true),
#[cfg(feature = "threading")]
scheduled: AtomicBool::new(false),
debug: AtomicU32::new(0),
garbage: PyMutex::new(Vec::new()),
py_garbage: ctx.new_list(Vec::new()),
py_callbacks: ctx.new_list(Vec::new()),
}
}
pub fn is_enabled(&self) -> bool {
self.enabled.load(Ordering::Relaxed)
}
#[cfg(feature = "threading")]
#[inline]
pub(crate) fn collection_ready(&self) -> bool {
self.scheduled.load(Ordering::Relaxed) && !gc_state().collecting.is_locked()
}
pub fn enable(&self) {
self.enabled.store(true, Ordering::Relaxed);
}
pub fn disable(&self) {
self.enabled.store(false, Ordering::Relaxed);
}
pub fn get_debug(&self) -> GcDebugFlags {
GcDebugFlags::from_bits_truncate(self.debug.load(Ordering::SeqCst))
}
pub fn set_debug(&self, flags: GcDebugFlags) {
self.debug.store(flags.bits(), Ordering::SeqCst);
}
pub fn get_threshold(&self) -> (u32, u32, u32) {
(
self.generations[0].threshold(),
self.generations[1].threshold(),
self.generations[2].threshold(),
)
}
pub fn set_threshold(&self, t0: u32, t1: Option<u32>, t2: Option<u32>) {
self.generations[0].set_threshold(t0);
if let Some(t1) = t1 {
self.generations[1].set_threshold(t1);
}
if let Some(t2) = t2 {
self.generations[2].set_threshold(t2);
}
}
pub fn get_stats(&self) -> [GcStats; 3] {
[
self.generations[0].stats(),
self.generations[1].stats(),
self.generations[2].stats(),
]
}
pub fn collect(&self, generation: usize) -> CollectResult {
gc_state().collect_inner(self, generation, false)
}
pub fn collect_force(&self, generation: usize) -> CollectResult {
gc_state().collect_inner(self, generation, true)
}
pub fn get_objects(&self, generation: Option<i32>) -> Vec<PyObjectRef> {
gc_state().get_objects(generation, self.owner)
}
pub fn freeze(&self) {
gc_state().freeze(self.owner);
}
pub fn unfreeze(&self) {
gc_state().unfreeze(self.owner);
}
#[cfg(all(unix, feature = "threading"))]
pub unsafe fn reinit_after_fork(&self) {
unsafe {
crate::common::lock::reinit_mutex_after_fork(&self.garbage);
for generation in &self.generations {
generation.reinit_stats_after_fork();
}
}
}
}
impl Drop for GcInterpreterState {
fn drop(&mut self) {
gc_state().retire_owner(self.owner);
}
}
#[must_use]
pub fn current_owner() -> GcOwner {
crate::vm::thread::current_gc_state().map_or(GC_NO_OWNER, |gc| unsafe { gc.as_ref() }.owner)
}
pub(crate) unsafe fn track_new_object(obj: NonNull<PyObject>) {
let state = gc_state();
let Some(gc) = crate::vm::thread::current_gc_state() else {
unsafe { state.track_object_fresh(obj, GC_NO_OWNER) };
return;
};
let gc = unsafe { gc.as_ref() };
unsafe { state.track_object_fresh(obj, gc.owner) };
state.maybe_collect(gc);
}
pub(crate) unsafe fn track_new_pair(obj: NonNull<PyObject>, frame: NonNull<PyObject>) {
let state = gc_state();
let Some(gc) = crate::vm::thread::current_gc_state() else {
unsafe { state.track_pair_fresh(obj, frame, GC_NO_OWNER) };
return;
};
let gc = unsafe { gc.as_ref() };
unsafe { state.track_pair_fresh(obj, frame, gc.owner) };
state.maybe_collect(gc);
}
pub fn gc_state() -> &'static GcState {
rustpython_common::static_cell! {
static GC_STATE: GcState;
}
GC_STATE.get_or_init(GcState::new)
}
#[cfg(test)]
mod tests {
use super::*;
fn interpreter_state() -> GcInterpreterState {
GcInterpreterState::new(crate::vm::Context::genesis())
}
#[test]
fn gc_state_default() {
let state = interpreter_state();
assert!(state.is_enabled());
assert_eq!(state.get_debug(), GcDebugFlags::empty());
assert_eq!(state.get_threshold(), (2000, 10, 0));
}
#[test]
fn gc_enable_disable() {
let state = interpreter_state();
assert!(state.is_enabled());
state.disable();
assert!(!state.is_enabled());
state.enable();
assert!(state.is_enabled());
}
#[test]
fn gc_threshold() {
let state = interpreter_state();
state.set_threshold(100, Some(20), Some(30));
assert_eq!(state.get_threshold(), (100, 20, 30));
}
#[test]
fn gc_debug_flags() {
let state = interpreter_state();
state.set_debug(GcDebugFlags::STATS | GcDebugFlags::COLLECTABLE);
assert_eq!(
state.get_debug(),
GcDebugFlags::STATS | GcDebugFlags::COLLECTABLE
);
}
#[test]
fn gc_owner_tags_are_distinct_while_live() {
let first = interpreter_state();
let second = interpreter_state();
assert_ne!(first.owner, second.owner);
assert_ne!(first.owner, GC_NO_OWNER);
assert_ne!(second.owner, GC_NO_OWNER);
}
#[cfg(feature = "threading")]
#[test]
fn automatic_gc_request_stays_with_allocating_interpreter() {
let first = crate::Interpreter::without_stdlib(Default::default());
let second = crate::Interpreter::without_stdlib(Default::default());
first.enter(|vm| vm.state.gc.scheduled.store(false, Ordering::Relaxed));
second.enter(|vm| vm.state.gc.scheduled.store(false, Ordering::Relaxed));
let allocations = GcState::new();
allocations.counts[0].store(1, Ordering::Relaxed);
first.enter(|vm| {
vm.state.gc.set_threshold(1, None, None);
assert!(!allocations.maybe_collect(&vm.state.gc));
});
second.enter(|vm| {
assert!(!vm.state.gc.scheduled.load(Ordering::Relaxed));
vm.run_scheduled_gc();
});
first.enter(|vm| {
assert!(vm.state.gc.scheduled.swap(false, Ordering::Relaxed));
});
}
#[cfg(feature = "threading")]
#[test]
fn automatic_gc_request_survives_busy_collector() {
let state = interpreter_state();
let _guard = gc_state().collecting.lock();
state.scheduled.store(true, Ordering::Relaxed);
state.collect(0);
assert!(state.scheduled.load(Ordering::Relaxed));
assert!(!state.collection_ready());
state.disable();
state.collect(0);
assert!(!state.scheduled.load(Ordering::Relaxed));
}
#[test]
fn release_count_does_not_wrap_during_reset() {
use std::sync::Barrier;
let count = AtomicUsize::new(1);
let start = Barrier::new(3);
let finish = Barrier::new(3);
let mut underflows = 0;
std::thread::scope(|scope| {
for reset in [false, true] {
let (count, start, finish) = (&count, &start, &finish);
scope.spawn(move || {
for _ in 0..10_000 {
start.wait();
if reset {
count.store(0, Ordering::Relaxed);
} else {
release_count(count);
}
finish.wait();
}
});
}
for _ in 0..10_000 {
count.store(1, Ordering::Relaxed);
start.wait();
finish.wait();
underflows += usize::from(count.load(Ordering::Relaxed) > 1);
}
});
assert_eq!(underflows, 0);
}
#[test]
fn survivor_promotion_handles_mixed_and_stale_generations() {
struct Heap {
state: GcState,
objects: Vec<PyObjectRef>,
}
impl Drop for Heap {
fn drop(&mut self) {
for obj in &self.objects {
unsafe { self.state.untrack_object(NonNull::from(obj.as_ref())) };
}
}
}
let ctx = crate::vm::Context::genesis();
let _collector = gc_state().collecting.lock();
let mut heap = Heap {
state: GcState::new(),
objects: Vec::new(),
};
for _ in 0..600 {
let obj: PyObjectRef = ctx.new_list(Vec::new()).into();
let ptr = NonNull::from(obj.as_ref());
unsafe {
gc_state().untrack_object(ptr);
heap.state.track_object(ptr, 1);
}
heap.objects.push(obj);
}
let Heap { state, objects } = &heap;
state.promote_survivors(0, &objects[..200]);
state.promote_survivors(1, &objects[..100]);
assert_eq!(state.get_count(), (400, 100, 100));
objects[250].set_gc_owner(2);
state.freeze(2);
unsafe { state.untrack_object(NonNull::from(objects[500].as_ref())) };
state.counts[0].store(0, Ordering::Relaxed);
state.promote_survivors(2, objects);
assert_eq!(state.get_count(), (0, 0, 598));
assert_eq!(state.get_freeze_count(), 1);
for (index, obj) in objects.iter().enumerate() {
let expected = match index {
250 => GC_PERMANENT,
500 => GC_UNTRACKED,
_ => 2,
};
assert_eq!(obj.gc_generation(), expected);
}
assert_eq!(state.generation_lists[0].read().iter().count(), 0);
assert_eq!(state.generation_lists[1].read().iter().count(), 0);
assert_eq!(state.generation_lists[2].read().iter().count(), 598);
let promoted: GcSet<_> = state.generation_lists[2]
.read()
.iter()
.map(NonNull::from)
.collect();
for obj in objects.iter().filter(|obj| obj.gc_generation() == 2) {
assert!(promoted.contains(&NonNull::from(obj.as_ref())));
}
}
}