FazBrowse GitHub Viewer | Trending |
URL:
| Home
Tools: [Download Repo ZIP]   [Original HTTPS Page]

Keep automatic GC requests local to each interpreter by 1ndahous3 · Pull Request #8902 · RustPython/RustPython · GitHub

Repository navigation

Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension .rs  (3) All 1 file type selected
Viewed files
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Unified
Split
Hide whitespace
Diff view
Unified
Split
Hide whitespace
99 changes: 95 additions & 4 deletions crates/vm/src/gc_state.rs
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters. Learn more about bidirectional Unicode characters
Original file line number Diff line number Diff line change
Expand Up @@ -135,9 +135,7 @@ impl GcGeneration {
/// empties its own interpreter's objects; another interpreter's stay behind with
/// the count already zeroed, and untracking one of those must not wrap.
fn release_count(count: &AtomicUsize) {
if count.load(Ordering::Relaxed) > 0 {
count.fetch_sub(1, Ordering::Relaxed);
}
let _ = count.try_update(Ordering::Relaxed, Ordering::Relaxed, |n| n.checked_sub(1));
}

/// Whether `owner`'s collections act on `obj`.
Expand Down Expand Up @@ -552,7 +550,7 @@ impl GcState {
// (e.g. a lazily-initialized frame locals cell) that another
// thread is blocked on with no way to reach a safepoint —
// a deadlock. At a safepoint no such lock is held.
crate::signal::schedule_gc();
gc.scheduled.store(true, Ordering::Relaxed);
return false;
}
// Without threading there is no safepoint to defer to and no other
Expand All @@ -574,6 +572,8 @@ impl GcState {
force: bool,
) -> CollectResult {
if !force && !gc.is_enabled() {
#[cfg(feature = "threading")]
gc.scheduled.store(false, Ordering::Relaxed);
return CollectResult::default();
}

Expand All @@ -582,6 +582,12 @@ impl GcState {
return CollectResult::default();
};

// A busy collector must not consume another interpreter's request.
// Clear only after acquiring the lock, before callbacks can request
// a later collection.
#[cfg(feature = "threading")]
gc.scheduled.store(false, Ordering::Relaxed);

let start_time = cfg_select! {
target_arch = "wasm32" => (),
_ => std::time::Instant::now(),
Expand Down Expand Up @@ -1268,6 +1274,9 @@ pub struct GcInterpreterState {
pub generations: [GcGeneration; 3],
/// GC enabled flag
enabled: AtomicBool,
/// Automatic collection requested at this interpreter's next safepoint.
#[cfg(feature = "threading")]
scheduled: AtomicBool,
/// Debug flags
debug: AtomicU32,
/// Uncollectable objects saved by this interpreter's collections, drained
Expand All @@ -1289,6 +1298,8 @@ impl GcInterpreterState {
GcGeneration::new(0), // old[1]
],
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()),
Expand All @@ -1304,6 +1315,14 @@ impl GcInterpreterState {
self.enabled.load(Ordering::Relaxed)
}

/// Leave requests pending while a collector is busy, without repeatedly
/// entering the bytecode loop's slow path during its Python callbacks.
#[cfg(feature = "threading")]
#[inline]
pub(crate) fn collection_ready(&self) -> bool {
self.scheduled.load(Ordering::Relaxed) && !gc_state().collecting.is_locked()
}

/// Enable GC
pub fn enable(&self) {
self.enabled.store(true, Ordering::Relaxed);
Expand Down Expand Up @@ -1524,4 +1543,76 @@ mod tests {
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));
// An isolated allocation counter makes crossing the threshold
// deterministic without depending on the rest of the test process.
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);
}
}
16 changes: 2 additions & 14 deletions crates/vm/src/signal.rs
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters. Learn more about bidirectional Unicode characters
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,6 @@ bitflagset::bitflag! {
enum EvalBreakerFlag {
Signal = 0,
Qsbr = 1,
Gc = 2,
Stop = 3,
Finalizing = 4,
}
Expand Down Expand Up @@ -189,20 +188,10 @@ mod mt {
EVAL_BREAKER.insert(EvalBreakerFlag::Finalizing);
}

/// Schedule an automatic collection to run at the next bytecode safepoint.
pub(crate) fn schedule_gc() {
EVAL_BREAKER.insert(EvalBreakerFlag::Gc);
}

/// Clear the scheduled-GC bit, returning whether it had been set.
pub(crate) fn take_gc_scheduled() -> bool {
EVAL_BREAKER.remove(EvalBreakerFlag::Gc)
}

/// Drop every process-wide eval-breaker bit. Tests that assert a single
/// thread's `stop_requested` must not trip `eval_breaker_pending` have to
/// start from a clean word: cargo's Windows runner shares the process
/// across `#[test]` functions, so a sibling can leave SIGNAL/QSBR/GC/STOP.
/// across `#[test]` functions, so a sibling can leave SIGNAL/QSBR/STOP.
#[cfg(test)]
pub(crate) fn clear_eval_breaker_for_test() {
EVAL_BREAKER.clear();
Expand All @@ -211,8 +200,7 @@ mod mt {

#[cfg(feature = "threading")]
pub(crate) use mt::{
clear_qsbr_bit, clear_stop_bit, qsbr_bit_set, schedule_gc, set_finalizing_bit, set_qsbr_bit,
set_stop_bit, take_gc_scheduled,
clear_qsbr_bit, clear_stop_bit, qsbr_bit_set, set_finalizing_bit, set_qsbr_bit, set_stop_bit,
};

#[cfg(all(test, feature = "threading"))]
Expand Down
4 changes: 2 additions & 2 deletions crates/vm/src/vm/mod.rs
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters. Learn more about bidirectional Unicode characters
Original file line number Diff line number Diff line change
Expand Up @@ -3545,7 +3545,7 @@ impl VirtualMachine {
#[inline]
pub(crate) fn eval_breaker_tripped(&self) -> bool {
#[cfg(feature = "threading")]
if thread::stop_requested_for_current_thread() {
if thread::stop_requested_for_current_thread() || self.state.gc.collection_ready() {
return true;
}
#[cfg(not(target_arch = "wasm32"))]
Expand Down Expand Up @@ -3594,7 +3594,7 @@ impl VirtualMachine {
/// against a thread blocked on a lock this thread would otherwise hold.
#[cfg(feature = "threading")]
pub(crate) fn run_scheduled_gc(&self) {
if crate::signal::take_gc_scheduled() {
if self.state.gc.collection_ready() {
self.state.gc.collect(0);
}
}
Expand Down
Loading

Back | FazBrowse Home | New Git URL