1#[cfg(feature = "threading")]
2use super::StopTheWorldState;
3use super::{
4 Context, PyConfig, PyGlobalState, VirtualMachine,
5 runtime::{self, InterpreterWhence},
6 setting::Settings,
7 thread,
8};
9use crate::{
10 PyResult, builtins, common::rc::PyRc, frozen::FrozenModule, getpath, py_freeze, stdlib::atexit,
11 vm::PyBaseExceptionRef,
12};
13use alloc::collections::BTreeMap;
14use core::sync::atomic::Ordering;
15
16type InitFunc = Box<dyn FnOnce(&mut VirtualMachine)>;
17
18const EXITCODE_FLUSH_FAILURE: u32 = 120;
21
22pub struct InterpreterBuilder {
37 settings: Settings,
38 pub ctx: PyRc<Context>,
39 module_defs: Vec<&'static builtins::PyModuleDef>,
40 frozen_modules: Vec<(&'static str, FrozenModule)>,
41 init_hooks: Vec<InitFunc>,
42}
43
44struct InitializeVmOpts<'a> {
46 settings: Settings,
47 ctx: PyRc<Context>,
48 module_defs: Vec<&'static builtins::PyModuleDef>,
49 frozen_modules: Vec<(&'static str, FrozenModule)>,
50 init_hooks: Vec<InitFunc>,
51 is_main: bool,
52 whence: InterpreterWhence,
53 parent_state: Option<&'a PyGlobalState>,
55 interp_config: runtime::InterpreterConfig,
56}
57
58fn initialize_vm<F>(opts: InitializeVmOpts<'_>, init: F) -> (VirtualMachine, PyRc<PyGlobalState>)
60where
61 F: FnOnce(&mut VirtualMachine),
62{
63 let InitializeVmOpts {
64 settings,
65 ctx,
66 module_defs,
67 frozen_modules,
68 init_hooks,
69 is_main,
70 whence,
71 parent_state,
72 interp_config,
73 } = opts;
74 use crate::codecs::CodecsRegistry;
75 use crate::common::lock::PyMutex;
76 use crate::warn::WarningsState;
77 use core::sync::atomic::{AtomicBool, AtomicI64, AtomicU64};
78 use crossbeam_utils::atomic::AtomicCell;
79
80 #[cfg(feature = "threading")]
82 thread::install_blocking_wait_hook();
83
84 let (config, all_module_defs, frozen, int_max_str_digits) = if let Some(parent) = parent_state {
85 let int_max_str_digits = AtomicCell::new(parent.int_max_str_digits.load());
87 (
88 parent.config.clone(),
89 parent.module_defs.clone(),
90 parent.frozen.clone(),
91 int_max_str_digits,
92 )
93 } else {
94 super::init_hash_secret(settings.hash_seed);
96 let paths = getpath::init_path_config(&settings);
97 let config = PyConfig::new(settings, paths);
98
99 let mut all_module_defs: BTreeMap<&'static str, &'static builtins::PyModuleDef> =
101 crate::stdlib::builtin_module_defs(&ctx)
102 .into_iter()
103 .chain(module_defs)
104 .map(|def| (def.name.as_str(), def))
105 .collect();
106
107 if let Some(&sysconfigdata_def) = all_module_defs.get("_sysconfigdata") {
109 use std::sync::OnceLock;
110 static SYSCONFIGDATA_NAME: OnceLock<&'static str> = OnceLock::new();
111 let leaked_name = *SYSCONFIGDATA_NAME.get_or_init(|| {
112 let name = crate::stdlib::sys::sysconfigdata_name();
113 Box::leak(name.into_boxed_str())
114 });
115 all_module_defs.insert(leaked_name, sysconfigdata_def);
116 }
117
118 let int_max_str_digits = AtomicCell::new(match config.settings.int_max_str_digits {
119 -1 => 4300,
120 other => other,
121 } as usize);
122
123 let mut frozen: std::collections::HashMap<
124 &'static str,
125 FrozenModule,
126 rapidhash::quality::RandomState,
127 > = core_frozen_inits().collect();
128 frozen.extend(frozen_modules);
129
130 (config, all_module_defs, frozen, int_max_str_digits)
131 };
132
133 let codec_registry = CodecsRegistry::new(&ctx);
135 let warnings = WarningsState::init_state(&ctx);
136
137 let interpreter_id = runtime::alloc_interpreter_id();
138 let runtime_root_id = parent_state.map_or(interpreter_id, |parent| parent.runtime_root_id);
139
140 #[cfg(feature = "threading")]
144 let main_thread_ident = AtomicCell::new(parent_state.map_or(0, |p| p.main_thread_ident.load()));
145
146 let feature_flags = interp_config.feature_flags();
147 let own_gil = interp_config.own_gil();
148
149 let global_state = PyRc::new(PyGlobalState {
151 gc: crate::gc_state::GcInterpreterState::new(&ctx),
152 interpreter_id,
153 runtime_root_id,
154 whence,
155 is_main,
156 config,
157 module_defs: all_module_defs,
158 frozen,
159 stacksize: AtomicCell::new(0),
160 thread_count: AtomicCell::new(0),
161 atexit_funcs: PyMutex::default(),
162 audit_hooks: PyMutex::default(),
163 codec_registry,
164 struct_format_cache: crate::buffer::FormatSpecCache::default(),
165 finalizing: AtomicBool::new(false),
166 #[cfg(feature = "threading")]
167 finalizing_thread_ident: AtomicCell::new(0),
168 warnings,
169 override_frozen_modules: AtomicCell::new(0),
170 before_forkers: PyMutex::default(),
171 after_forkers_child: PyMutex::default(),
172 after_forkers_parent: PyMutex::default(),
173 int_max_str_digits,
174 switch_interval: AtomicCell::new(0.005),
175 global_trace_func: PyMutex::default(),
176 global_profile_func: PyMutex::default(),
177 type_mutex: PyMutex::default(),
178 #[cfg(feature = "threading")]
179 main_thread_ident,
180 #[cfg(feature = "threading")]
181 thread_frames: parking_lot::Mutex::new(std::collections::HashMap::new()),
182 #[cfg(feature = "threading")]
183 thread_handles: parking_lot::Mutex::new(Vec::new()),
184 #[cfg(feature = "threading")]
185 shutdown_handles: parking_lot::Mutex::new(Vec::new()),
186 monitoring: PyMutex::default(),
187 monitoring_events: AtomicCell::new(0),
188 instrumentation_version: AtomicU64::new(0),
189 #[cfg(feature = "threading")]
190 stop_the_world: StopTheWorldState::new(),
191 feature_flags,
192 own_gil,
193 running_main: AtomicBool::new(false),
194 ready: AtomicBool::new(false),
195 id_refcount: AtomicI64::new(0),
196 require_idref: AtomicBool::new(false),
197 });
198
199 let mut vm = VirtualMachine::new(ctx, global_state);
202
203 for hook in init_hooks {
205 hook(&mut vm);
206 }
207
208 init(&mut vm);
210
211 runtime::register_interpreter(&vm.state);
217
218 let vm_guard = thread::VmBootstrapGuard::new(&vm);
222 vm.initialize();
223 vm.state.ready.store(true, Ordering::Release);
224 drop(vm_guard);
225
226 let global_state = vm.state.clone();
228 (vm, global_state)
229}
230
231fn create_subinterpreter_from_parent(
233 parent: &VirtualMachine,
234 config: runtime::InterpreterConfig,
235) -> Result<Interpreter, &'static str> {
236 config.check()?;
237 #[cfg(feature = "threading")]
243 let _restore_parent = {
244 let saved = thread::current_vm_is_set().then(thread::save_current_thread);
245 scopeguard::guard(saved, |saved| {
246 if let Some(saved) = saved {
247 thread::restore_current_thread(saved);
248 }
249 })
250 };
251
252 let (vm, global_state) = initialize_vm(
253 InitializeVmOpts {
254 settings: Settings::default(),
256 ctx: parent.ctx.clone(),
257 module_defs: Vec::new(),
258 frozen_modules: Vec::new(),
259 init_hooks: Vec::new(),
260 is_main: false,
261 whence: InterpreterWhence::Stdlib,
262 parent_state: Some(&parent.state),
263 interp_config: config,
264 },
265 |_| {},
266 );
267 let interp = Interpreter { global_state, vm };
268 interp.enter(|vm| {
270 let _ = vm.ensure_main_module();
271 });
272 Ok(interp)
273}
274
275impl InterpreterBuilder {
276 #[must_use]
278 pub fn new() -> Self {
279 Self {
280 settings: Settings::default(),
281 ctx: Context::genesis().clone(),
282 module_defs: Vec::new(),
283 frozen_modules: Vec::new(),
284 init_hooks: Vec::new(),
285 }
286 }
287
288 #[must_use]
292 pub fn settings(mut self, settings: Settings) -> Self {
293 self.settings = settings;
294 self
295 }
296
297 #[must_use]
310 pub fn add_native_module(self, def: &'static builtins::PyModuleDef) -> Self {
311 self.add_native_modules(&[def])
312 }
313
314 #[must_use]
327 pub fn add_native_modules(mut self, defs: &[&'static builtins::PyModuleDef]) -> Self {
328 self.module_defs.extend_from_slice(defs);
329 self
330 }
331
332 #[must_use]
349 pub fn init_hook<F>(mut self, init: F) -> Self
350 where
351 F: FnOnce(&mut VirtualMachine) + 'static,
352 {
353 self.init_hooks.push(Box::new(init));
354 self
355 }
356
357 #[must_use]
371 pub fn add_frozen_modules<I>(mut self, frozen: I) -> Self
372 where
373 I: IntoIterator<Item = (&'static str, FrozenModule)>,
374 {
375 self.frozen_modules.extend(frozen);
376 self
377 }
378
379 #[must_use]
383 pub fn build(self) -> Interpreter {
384 let (vm, global_state) = initialize_vm(
385 InitializeVmOpts {
386 settings: self.settings,
387 ctx: self.ctx,
388 module_defs: self.module_defs,
389 frozen_modules: self.frozen_modules,
390 init_hooks: self.init_hooks,
391 is_main: true,
392 whence: InterpreterWhence::Runtime,
393 parent_state: None,
394 interp_config: runtime::InterpreterConfig::MAIN,
395 },
396 |_| {}, );
398 Interpreter { global_state, vm }
399 }
400
401 #[must_use]
403 pub fn interpreter(self) -> Interpreter {
404 self.build()
405 }
406}
407
408impl Default for InterpreterBuilder {
409 fn default() -> Self {
410 Self::new()
411 }
412}
413
414pub struct Interpreter {
439 pub global_state: PyRc<PyGlobalState>,
440 vm: VirtualMachine,
441}
442
443impl Interpreter {
444 #[must_use]
455 pub fn builder(settings: Settings) -> InterpreterBuilder {
456 InterpreterBuilder::new().settings(settings)
457 }
458
459 #[must_use]
464 pub fn without_stdlib(settings: Settings) -> Self {
465 Self::with_init(settings, |_| {})
466 }
467
468 pub fn with_init<F>(settings: Settings, init: F) -> Self
472 where
473 F: FnOnce(&mut VirtualMachine),
474 {
475 let (vm, global_state) = initialize_vm(
476 InitializeVmOpts {
477 settings,
478 ctx: Context::genesis().clone(),
479 module_defs: Vec::new(),
480 frozen_modules: Vec::new(),
481 init_hooks: Vec::new(),
482 is_main: true,
483 whence: InterpreterWhence::Runtime,
484 parent_state: None,
485 interp_config: runtime::InterpreterConfig::MAIN,
486 },
487 init,
488 );
489 Self { global_state, vm }
490 }
491
492 #[inline]
494 #[must_use]
495 pub fn id(&self) -> i64 {
496 self.global_state.interpreter_id
497 }
498
499 #[inline]
501 #[must_use]
502 pub fn whence(&self) -> InterpreterWhence {
503 self.global_state.whence
504 }
505
506 #[inline]
511 #[must_use]
512 pub fn is_main(&self) -> bool {
513 self.global_state.is_main
514 }
515
516 #[inline]
521 #[must_use]
522 pub fn is_process_main(&self) -> bool {
523 runtime::main_interpreter_id() == Some(self.id())
524 }
525
526 #[cfg(feature = "threading")]
532 #[must_use]
533 pub fn create_owned_subinterpreter(&self) -> i64 {
534 self.create_owned_subinterpreter_with_config(runtime::InterpreterConfig::ISOLATED)
535 .expect("the isolated config is always valid")
536 }
537
538 #[cfg(feature = "threading")]
540 pub fn create_owned_subinterpreter_with_config(
541 &self,
542 config: runtime::InterpreterConfig,
543 ) -> Result<i64, &'static str> {
544 Ok(runtime::store_owned_interpreter(
545 self.create_subinterpreter_with_config(config)?,
546 ))
547 }
548
549 #[must_use]
559 pub fn create_subinterpreter(&self) -> Self {
560 self.create_subinterpreter_with_config(runtime::InterpreterConfig::ISOLATED)
561 .expect("the isolated config is always valid")
562 }
563
564 pub fn create_subinterpreter_from_vm(
566 parent: &VirtualMachine,
567 config: runtime::InterpreterConfig,
568 ) -> Result<Self, &'static str> {
569 create_subinterpreter_from_parent(parent, config)
570 }
571
572 pub fn create_subinterpreter_with_config(
574 &self,
575 config: runtime::InterpreterConfig,
576 ) -> Result<Self, &'static str> {
577 create_subinterpreter_from_parent(&self.vm, config)
578 }
579
580 #[cfg(feature = "threading")]
582 pub fn new_thread(&self) -> thread::ThreadedVirtualMachine {
583 self.vm.new_thread()
584 }
585
586 pub fn enter<F, R>(&self, f: F) -> R
596 where
597 F: FnOnce(&VirtualMachine) -> R,
598 {
599 thread::enter_vm(&self.vm, || f(&self.vm))
600 }
601
602 pub fn enter_and_expect<F, R>(&self, f: F, msg: &str) -> R
609 where
610 F: FnOnce(&VirtualMachine) -> PyResult<R>,
611 {
612 self.enter(|vm| {
613 let result = f(vm);
614 vm.expect_pyresult(result, msg)
615 })
616 }
617
618 pub fn run<F>(self, f: F) -> u32
627 where
628 F: FnOnce(&VirtualMachine) -> PyResult<()>,
629 {
630 let res = self.enter(|vm| f(vm));
631 self.finalize(res.err())
632 }
633
634 pub fn finalize(self, exc: Option<PyBaseExceptionRef>) -> u32 {
649 self.enter(|vm| {
650 let mut flush_status = vm.flush_std();
651
652 let exit_code = if let Some(exc) = exc {
654 vm.handle_exit_exception(exc)
655 } else {
656 0
657 };
658
659 #[cfg(feature = "threading")]
663 if let Ok(threading) = vm.import("threading", 0)
664 && let Ok(shutdown) = threading.get_attr("_shutdown", vm)
665 && let Err(e) = shutdown.call((), vm)
666 {
667 vm.run_unraisable(
668 e,
669 Some("Exception ignored in threading shutdown".to_owned()),
670 threading,
671 );
672 }
673
674 atexit::_run_exitfuncs(vm);
677
678 #[cfg(feature = "threading")]
682 finalize_subinterpreters(vm);
683
684 #[cfg(feature = "threading")]
689 {
690 vm.state.stop_the_world.stop_the_world(&vm.state);
691 vm.state
692 .finalizing_thread_ident
693 .store(crate::stdlib::_thread::get_ident());
694 }
695 vm.state.finalizing.store(true, Ordering::Release);
696 #[cfg(feature = "threading")]
697 {
698 thread::set_other_threads_shutting_down(&vm.state);
699 vm.state.stop_the_world.start_the_world(&vm.state);
700 crate::signal::set_finalizing_bit();
701 }
702
703 vm.state.gc.collect_force(2);
705
706 vm.finalize_modules();
709
710 let interpreter_id = vm.state.interpreter_id;
715 crate::stdlib::_interpchannels::clear_interpreter(interpreter_id);
716 crate::stdlib::_interpqueues::clear_interpreter(interpreter_id);
717
718 if vm.flush_std() < 0 && flush_status == 0 {
719 flush_status = -1;
720 }
721
722 let exit_code = if exit_code == 0 && flush_status < 0 {
724 EXITCODE_FLUSH_FAILURE
725 } else {
726 exit_code
727 };
728
729 #[cfg(feature = "threading")]
732 crate::object::qsbr::QSBR.process();
733
734 exit_code
735 })
736 }
737}
738
739#[cfg(feature = "threading")]
742fn finalize_subinterpreters(vm: &VirtualMachine) {
743 if !vm.state.is_main {
744 return;
745 }
746 let root_id = vm.state.runtime_root_id;
747 if runtime::owned_interpreter_ids_for(root_id).is_empty() {
749 return;
750 }
751 let message = vm
753 .ctx
754 .new_str("remaining subinterpreters; close them with Interpreter.close()");
755 let _ = crate::warn::warn(
756 message.into(),
757 Some(vm.ctx.exceptions.runtime_warning.to_owned()),
758 0,
759 None,
760 vm,
761 );
762 while let Some(id) = runtime::owned_interpreter_ids_for(root_id)
765 .into_iter()
766 .next()
767 {
768 let _ = runtime::destroy_owned_interpreter(id);
769 }
770}
771
772fn core_frozen_inits() -> impl Iterator<Item = (&'static str, FrozenModule)> {
773 let iter = core::iter::empty();
774 macro_rules! ext_modules {
775 ($iter:ident, $($t:tt)*) => {
776 let $iter = $iter.chain(py_freeze!($($t)*));
777 };
778 }
779
780 ext_modules!(
784 iter,
785 dir = "../../Lib/python_builtins",
786 crate_name = "rustpython_compiler_core"
787 );
788
789 ext_modules!(
795 iter,
796 dir = "../../Lib/core_modules",
797 crate_name = "rustpython_compiler_core"
798 );
799
800 let mut entries: Vec<_> = iter.collect();
802
803 if let Some(hello_code) = entries
805 .iter()
806 .find(|(n, _)| *n == "__hello__")
807 .map(|(_, m)| m.code)
808 {
809 entries.push((
810 "__hello_alias__",
811 FrozenModule {
812 code: hello_code,
813 package: false,
814 },
815 ));
816 entries.push((
817 "__phello_alias__",
818 FrozenModule {
819 code: hello_code,
820 package: true,
821 },
822 ));
823 entries.push((
824 "__phello_alias__.spam",
825 FrozenModule {
826 code: hello_code,
827 package: false,
828 },
829 ));
830 entries.push((
831 "__hello_only__",
832 FrozenModule {
833 code: hello_code,
834 package: false,
835 },
836 ));
837 }
838 if let Some(code) = entries
839 .iter()
840 .find(|(n, _)| *n == "__phello__")
841 .map(|(_, m)| m.code)
842 {
843 entries.push((
844 "__phello__.__init__",
845 FrozenModule {
846 code,
847 package: false,
848 },
849 ));
850 }
851 if let Some(code) = entries
852 .iter()
853 .find(|(n, _)| *n == "__phello__.ham")
854 .map(|(_, m)| m.code)
855 {
856 entries.push((
857 "__phello__.ham.__init__",
858 FrozenModule {
859 code,
860 package: false,
861 },
862 ));
863 }
864 entries.into_iter()
865}
866
867#[cfg(test)]
868mod tests {
869 use super::*;
870 use crate::{
871 AsObject, PyObjectRef,
872 builtins::{PyStr, int},
873 vm::{MAIN_INTERPRETER_ID, runtime},
874 };
875 use malachite_bigint::ToBigInt;
876
877 #[test]
878 fn add_py_integers() {
879 Interpreter::without_stdlib(Default::default()).enter(|vm| {
880 let a: PyObjectRef = vm.ctx.new_int(33_i32).into();
881 let b: PyObjectRef = vm.ctx.new_int(12_i32).into();
882 let res = vm._add(&a, &b).unwrap();
883 let value = int::get_value(&res);
884 assert_eq!(*value, 45_i32.to_bigint().unwrap());
885 })
886 }
887
888 #[test]
889 fn multiply_str() {
890 Interpreter::without_stdlib(Default::default()).enter(|vm| {
891 let a = vm.new_pyobj(crate::common::ascii!("Hello "));
892 let b = vm.new_pyobj(4_i32);
893 let res = vm._mul(&a, &b).unwrap();
894 let value = res.downcast_ref::<PyStr>().unwrap();
895 assert_eq!(value.as_wtf8(), "Hello Hello Hello Hello ")
896 })
897 }
898
899 #[test]
901 fn main_interpreter_identity() {
902 let main = Interpreter::without_stdlib(Default::default());
903 assert!(main.is_main());
904 assert_eq!(main.whence(), InterpreterWhence::Runtime);
905 assert!(
906 runtime::list_interpreters()
907 .iter()
908 .any(|info| info.id == main.id() && info.whence == InterpreterWhence::Runtime)
909 );
910 assert!(main.id() >= MAIN_INTERPRETER_ID);
913 }
914
915 #[test]
917 fn create_subinterpreter_registers_distinct_ids() {
918 let main = Interpreter::without_stdlib(Default::default());
919 let sub1 = main.create_subinterpreter();
920 let sub2 = main.create_subinterpreter();
921
922 assert!(main.is_main());
923 assert!(!sub1.is_main());
924 assert!(!sub2.is_main());
925 assert_eq!(sub1.whence(), InterpreterWhence::Stdlib);
926 assert_eq!(sub2.whence(), InterpreterWhence::Stdlib);
927 assert_ne!(main.id(), sub1.id());
928 assert_ne!(main.id(), sub2.id());
929 assert_ne!(sub1.id(), sub2.id());
930
931 let ids: Vec<i64> = runtime::list_interpreters()
932 .into_iter()
933 .map(|i| i.id)
934 .collect();
935 assert!(ids.contains(&main.id()));
936 assert!(ids.contains(&sub1.id()));
937 assert!(ids.contains(&sub2.id()));
938 }
939
940 fn wait_until_unregistered(id: i64) {
947 use core::time::Duration;
948 use std::time::Instant;
949
950 let deadline = Instant::now() + Duration::from_secs(30);
951 while runtime::lookup_interpreter(id).is_some() {
952 assert!(
953 Instant::now() < deadline,
954 "interpreter {id} still registered long after its last reference"
955 );
956 std::thread::yield_now();
957 }
958 }
959
960 #[cfg(feature = "threading")]
966 #[test]
967 fn registering_waits_for_an_in_flight_stop() {
968 use core::time::Duration;
969 use std::sync::mpsc;
970
971 let admission = runtime::lock_admission_for_stop();
973
974 let (tx, rx) = mpsc::channel();
975 let worker = std::thread::spawn(move || {
976 let interp = Interpreter::without_stdlib(Default::default());
977 tx.send(interp.id()).expect("receiver is alive");
978 interp
979 });
980
981 assert!(
982 matches!(
983 rx.recv_timeout(Duration::from_millis(200)),
984 Err(mpsc::RecvTimeoutError::Timeout)
985 ),
986 "an interpreter registered while a stop-the-world was in flight"
987 );
988
989 drop(admission);
990 let id = rx
991 .recv_timeout(Duration::from_secs(30))
992 .expect("registration proceeds once the world restarts");
993 assert!(runtime::lookup_interpreter(id).is_some());
994 drop(worker.join().expect("worker did not panic"));
995 wait_until_unregistered(id);
996 }
997
998 #[test]
1000 fn drop_subinterpreter_unregisters() {
1001 let main = Interpreter::without_stdlib(Default::default());
1002 let sub_id = {
1003 let sub = main.create_subinterpreter();
1004 let id = sub.id();
1005 assert!(runtime::lookup_interpreter(id).is_some());
1006 id
1007 };
1008 wait_until_unregistered(sub_id);
1009 assert!(runtime::lookup_interpreter(main.id()).is_some());
1010 }
1011
1012 #[test]
1014 fn subinterpreters_isolate_modules() {
1015 let main = Interpreter::without_stdlib(Default::default());
1016 let sub = main.create_subinterpreter();
1017
1018 let (main_sys_ptr, main_builtins_ptr, main_ctx_ptr, main_state_ptr) = main.enter(|vm| {
1019 (
1020 vm.sys_module.as_object() as *const _,
1021 vm.builtins.as_object() as *const _,
1022 PyRc::as_ptr(&vm.ctx),
1023 PyRc::as_ptr(&vm.state),
1024 )
1025 });
1026 let (sub_sys_ptr, sub_builtins_ptr, sub_ctx_ptr, sub_state_ptr) = sub.enter(|vm| {
1027 (
1028 vm.sys_module.as_object() as *const _,
1029 vm.builtins.as_object() as *const _,
1030 PyRc::as_ptr(&vm.ctx),
1031 PyRc::as_ptr(&vm.state),
1032 )
1033 });
1034
1035 assert_ne!(main_sys_ptr, sub_sys_ptr);
1036 assert_ne!(main_builtins_ptr, sub_builtins_ptr);
1037 assert_ne!(main_state_ptr, sub_state_ptr);
1039 assert_eq!(main_ctx_ptr, sub_ctx_ptr);
1041 }
1042
1043 #[test]
1045 fn subinterpreters_behaviorally_isolate_builtins_and_sys_modules() {
1046 const PROBE: &str = "__rustpython_subinterpreter_isolation_probe__";
1047
1048 let main = Interpreter::without_stdlib(Default::default());
1049 let sub = main.create_subinterpreter();
1050
1051 main.enter(|vm| {
1052 vm.builtins
1053 .set_attr(PROBE, vm.ctx.new_int(11_i32), vm)
1054 .unwrap();
1055 vm.sys_module
1056 .get_attr("modules", vm)
1057 .unwrap()
1058 .set_item(PROBE, vm.ctx.new_int(12_i32).into(), vm)
1059 .unwrap();
1060 });
1061
1062 sub.enter(|vm| {
1063 assert!(vm.builtins.get_attr(PROBE, vm).is_err());
1064 let modules = vm.sys_module.get_attr("modules", vm).unwrap();
1065 assert!(modules.get_item(PROBE, vm).is_err());
1066
1067 vm.builtins
1068 .set_attr(PROBE, vm.ctx.new_int(21_i32), vm)
1069 .unwrap();
1070 modules
1071 .set_item(PROBE, vm.ctx.new_int(22_i32).into(), vm)
1072 .unwrap();
1073 });
1074
1075 main.enter(|vm| {
1076 let builtin_probe = vm.builtins.get_attr(PROBE, vm).unwrap();
1077 assert_eq!(*int::get_value(&builtin_probe), 11_i32.to_bigint().unwrap());
1078
1079 let module_probe = vm
1080 .sys_module
1081 .get_attr("modules", vm)
1082 .unwrap()
1083 .get_item(PROBE, vm)
1084 .unwrap();
1085 assert_eq!(*int::get_value(&module_probe), 12_i32.to_bigint().unwrap());
1086 });
1087 }
1088
1089 #[test]
1092 fn create_subinterpreter_while_parent_entered() {
1093 let main = Interpreter::without_stdlib(Default::default());
1094 main.enter(|vm| {
1095 let before = vm.state.interpreter_id;
1096 let sub = main.create_subinterpreter();
1097 assert_ne!(sub.id(), before);
1098 assert_eq!(vm.state.interpreter_id, before);
1100 let n: PyObjectRef = vm.ctx.new_int(7_i32).into();
1102 assert_eq!(int::get_value(&n), &7_i32.to_bigint().unwrap());
1103 drop(sub);
1105 });
1106 }
1107
1108 #[test]
1110 fn sequential_enter_main_and_sub() {
1111 let main = Interpreter::without_stdlib(Default::default());
1112 let sub = main.create_subinterpreter();
1113
1114 main.enter(|vm| {
1115 assert!(vm.state.is_main_interpreter());
1116 let a: PyObjectRef = vm.ctx.new_int(1_i32).into();
1117 let b: PyObjectRef = vm.ctx.new_int(2_i32).into();
1118 let res = vm._add(&a, &b).unwrap();
1119 assert_eq!(*int::get_value(&res), 3_i32.to_bigint().unwrap());
1120 });
1121 sub.enter(|vm| {
1122 assert!(!vm.state.is_main_interpreter());
1123 let a: PyObjectRef = vm.ctx.new_int(10_i32).into();
1124 let b: PyObjectRef = vm.ctx.new_int(5_i32).into();
1125 let res = vm._mul(&a, &b).unwrap();
1126 assert_eq!(*int::get_value(&res), 50_i32.to_bigint().unwrap());
1127 });
1128 main.enter(|vm| {
1130 assert!(vm.state.is_main_interpreter());
1131 });
1132 }
1133
1134 #[cfg(feature = "threading")]
1136 #[test]
1137 fn concurrent_main_and_subinterpreter_threads() {
1138 use alloc::sync::Arc;
1139 use core::sync::atomic::{AtomicUsize, Ordering};
1140
1141 let main = Interpreter::without_stdlib(Default::default());
1142 let sub = main.create_subinterpreter();
1143 let counter = Arc::new(AtomicUsize::new(0));
1144
1145 let c1 = Arc::clone(&counter);
1146 let h_main = main.enter(|vm| {
1147 let thread_vm = vm.new_thread();
1148 let c = Arc::clone(&c1);
1149 std::thread::spawn(move || {
1150 thread_vm.run(|vm| {
1151 for _ in 0..100 {
1152 let a: PyObjectRef = vm.ctx.new_int(1_i32).into();
1153 let b: PyObjectRef = vm.ctx.new_int(1_i32).into();
1154 let _ = vm._add(&a, &b).unwrap();
1155 c.fetch_add(1, Ordering::Relaxed);
1156 }
1157 assert!(vm.state.is_main_interpreter());
1158 });
1159 })
1160 });
1161
1162 let c2 = Arc::clone(&counter);
1163 let h_sub = sub.enter(|vm| {
1164 let thread_vm = vm.new_thread();
1165 let c = Arc::clone(&c2);
1166 std::thread::spawn(move || {
1167 thread_vm.run(|vm| {
1168 for _ in 0..100 {
1169 let a: PyObjectRef = vm.ctx.new_int(2_i32).into();
1170 let b: PyObjectRef = vm.ctx.new_int(3_i32).into();
1171 let _ = vm._mul(&a, &b).unwrap();
1172 c.fetch_add(1, Ordering::Relaxed);
1173 }
1174 assert!(!vm.state.is_main_interpreter());
1175 });
1176 })
1177 });
1178
1179 h_main.join().expect("main worker panicked");
1180 h_sub.join().expect("sub worker panicked");
1181 assert_eq!(counter.load(Ordering::Relaxed), 200);
1182 }
1183
1184 #[cfg(feature = "threading")]
1186 #[test]
1187 fn main_and_subinterpreter_run_sections_overlap() {
1188 use alloc::sync::Arc;
1189 use core::time::Duration;
1190 use std::{
1191 sync::{Condvar, Mutex},
1192 time::Instant,
1193 };
1194
1195 #[derive(Default)]
1196 struct OverlapState {
1197 entered: usize,
1198 release: bool,
1199 }
1200
1201 let main = Interpreter::without_stdlib(Default::default());
1202 let sub = main.create_subinterpreter();
1203 let state = Arc::new((Mutex::new(OverlapState::default()), Condvar::new()));
1204
1205 let spawn_worker = |interpreter: &Interpreter| {
1206 let state = Arc::clone(&state);
1207 interpreter.enter(|vm| {
1208 let thread_vm = vm.new_thread();
1209 std::thread::spawn(move || {
1210 thread_vm.run(|vm| {
1211 let a: PyObjectRef = vm.ctx.new_int(20_i32).into();
1212 let b: PyObjectRef = vm.ctx.new_int(22_i32).into();
1213 assert_eq!(
1214 *int::get_value(&vm._add(&a, &b).unwrap()),
1215 42_i32.to_bigint().unwrap()
1216 );
1217
1218 let (lock, ready) = &*state;
1219 {
1220 let mut state = lock.lock().unwrap();
1221 state.entered += 1;
1222 ready.notify_all();
1223 }
1224 loop {
1229 vm.check_signals().unwrap();
1230 let state = lock.lock().unwrap();
1231 if state.release {
1232 break;
1233 }
1234 let _ = ready.wait_timeout(state, Duration::from_millis(1)).unwrap();
1235 }
1236 });
1237 })
1238 })
1239 };
1240
1241 let main_worker = spawn_worker(&main);
1242 let sub_worker = spawn_worker(&sub);
1243
1244 let (lock, ready) = &*state;
1245 let deadline = Instant::now() + Duration::from_secs(30);
1246 let mut state_guard = lock.lock().unwrap();
1247 while state_guard.entered < 2 {
1248 let now = Instant::now();
1249 if now >= deadline {
1250 break;
1251 }
1252 let (next, _) = ready.wait_timeout(state_guard, deadline - now).unwrap();
1253 state_guard = next;
1254 }
1255 let overlapped = state_guard.entered == 2;
1256 state_guard.release = true;
1257 ready.notify_all();
1258 drop(state_guard);
1259
1260 main_worker.join().expect("main worker panicked");
1261 sub_worker.join().expect("subinterpreter worker panicked");
1262 assert!(
1263 overlapped,
1264 "main and subinterpreter run sections were serialized"
1265 );
1266 }
1267
1268 #[cfg(feature = "threading")]
1270 #[test]
1271 fn busy_main_interpreter_does_not_block_subinterpreter() {
1272 use alloc::sync::Arc;
1273 use core::{
1274 sync::atomic::{AtomicBool, Ordering},
1275 time::Duration,
1276 };
1277 use std::time::Instant;
1278
1279 let main = Interpreter::without_stdlib(Default::default());
1280 let sub = main.create_subinterpreter();
1281 let main_started = Arc::new(AtomicBool::new(false));
1282 let sub_finished = Arc::new(AtomicBool::new(false));
1283
1284 let main_started_worker = Arc::clone(&main_started);
1285 let sub_finished_worker = Arc::clone(&sub_finished);
1286 let main_worker = main.enter(|vm| {
1287 let thread_vm = vm.new_thread();
1288 std::thread::spawn(move || {
1289 thread_vm.run(|vm| {
1290 main_started_worker.store(true, Ordering::Release);
1291 let deadline = Instant::now() + Duration::from_secs(30);
1292 let mut operations = 0;
1293 while !sub_finished_worker.load(Ordering::Acquire) && Instant::now() < deadline
1294 {
1295 let a: PyObjectRef = vm.ctx.new_int(20_i32).into();
1296 let b: PyObjectRef = vm.ctx.new_int(22_i32).into();
1297 let result = vm._add(&a, &b).unwrap();
1298 assert_eq!(*int::get_value(&result), 42_i32.to_bigint().unwrap());
1299 operations += 1;
1300 vm.check_signals().unwrap();
1305 std::thread::yield_now();
1306 }
1307 (sub_finished_worker.load(Ordering::Acquire), operations)
1308 })
1309 })
1310 });
1311
1312 let main_started_worker = Arc::clone(&main_started);
1313 let sub_finished_worker = Arc::clone(&sub_finished);
1314 let sub_worker = sub.enter(|vm| {
1315 let thread_vm = vm.new_thread();
1316 std::thread::spawn(move || {
1317 while !main_started_worker.load(Ordering::Acquire) {
1318 std::thread::yield_now();
1319 }
1320 thread_vm.run(|vm| {
1321 let a: PyObjectRef = vm.ctx.new_int(6_i32).into();
1322 let b: PyObjectRef = vm.ctx.new_int(7_i32).into();
1323 let result = vm._mul(&a, &b).unwrap();
1324 assert_eq!(*int::get_value(&result), 42_i32.to_bigint().unwrap());
1325 sub_finished_worker.store(true, Ordering::Release);
1326 });
1327 })
1328 });
1329
1330 let (sub_progressed_while_main_was_busy, main_operations) =
1331 main_worker.join().expect("main worker panicked");
1332 sub_worker.join().expect("subinterpreter worker panicked");
1333
1334 assert!(main_operations > 0);
1335 assert!(
1336 sub_progressed_while_main_was_busy,
1337 "subinterpreter made no progress until the busy main interpreter exited"
1338 );
1339 }
1340
1341 #[cfg(feature = "threading")]
1343 #[test]
1344 fn subinterpreter_new_thread_shares_sub_state() {
1345 let main = Interpreter::without_stdlib(Default::default());
1346 let sub = main.create_subinterpreter();
1347 let sub_id = sub.id();
1348
1349 let handle = sub.enter(|vm| {
1350 let thread_vm = vm.new_thread();
1351 std::thread::spawn(move || {
1352 thread_vm.run(|vm| {
1353 assert_eq!(vm.state.interpreter_id, sub_id);
1354 assert!(!vm.state.is_main_interpreter());
1355 });
1356 })
1357 });
1358 handle.join().expect("thread panicked");
1359 }
1360
1361 #[cfg(feature = "rustpython-compiler")]
1363 #[test]
1364 fn subinterpreter_runs_python_code() {
1365 use crate::compiler::Mode;
1366
1367 let main = Interpreter::without_stdlib(Default::default());
1368 let sub = main.create_subinterpreter();
1369
1370 sub.enter(|vm| {
1371 let scope = vm.new_scope_with_builtins();
1372 let source = "x = 40 + 2\n";
1373 let code = vm
1374 .compile(source, Mode::Exec, "<sub>")
1375 .map_err(|err| err.into_pyexception(vm, Some(source)))
1376 .unwrap();
1377 vm.run_code_obj(code, scope.clone()).unwrap();
1378 let x = scope.globals.get_item("x", vm).unwrap();
1379 assert_eq!(*int::get_value(&x), 42_i32.to_bigint().unwrap());
1380 });
1381 }
1382
1383 fn run(vm: &VirtualMachine, scope: &crate::scope::Scope, source: &str) {
1386 let code = vm
1387 .compile(source, crate::compiler::Mode::Exec, "<test>")
1388 .map_err(|err| err.into_pyexception(vm, Some(source)))
1389 .unwrap();
1390 vm.run_code_obj(code, scope.clone()).unwrap();
1391 }
1392
1393 #[test]
1394 fn subinterpreter_subclasses_are_scoped_to_their_interpreter() {
1395 use crate::scope::Scope;
1396
1397 fn lists_subclass(vm: &VirtualMachine, scope: &Scope, name: &str) -> bool {
1398 run(
1399 vm,
1400 scope,
1401 &format!("found = any(c.__name__ == {name:?} for c in int.__subclasses__())\n"),
1402 );
1403 let found = scope.globals.get_item("found", vm).unwrap();
1404 found.try_to_bool(vm).unwrap()
1405 }
1406
1407 let main = Interpreter::without_stdlib(Default::default());
1408 let sub = main.create_subinterpreter();
1409
1410 let main_scope = main.enter(|vm| {
1413 let scope = vm.new_scope_with_builtins();
1414 run(vm, &scope, "class MainOnly(int): pass\n");
1415 scope
1416 });
1417 let sub_scope = sub.enter(|vm| {
1418 let scope = vm.new_scope_with_builtins();
1419 run(vm, &scope, "class SubOnly(int): pass\n");
1420 scope
1421 });
1422
1423 main.enter(|vm| {
1424 assert!(lists_subclass(vm, &main_scope, "MainOnly"));
1425 assert!(!lists_subclass(vm, &main_scope, "SubOnly"));
1426 assert!(lists_subclass(vm, &main_scope, "bool"));
1429 });
1430 sub.enter(|vm| {
1431 assert!(lists_subclass(vm, &sub_scope, "SubOnly"));
1432 assert!(!lists_subclass(vm, &sub_scope, "MainOnly"));
1433 assert!(lists_subclass(vm, &sub_scope, "bool"));
1434 });
1435
1436 main.enter(|_| drop(main_scope));
1437 sub.enter(|_| drop(sub_scope));
1438 }
1439
1440 #[test]
1442 fn collections_only_reach_the_collecting_interpreter() {
1443 use core::time::Duration;
1444 use std::time::Instant;
1445
1446 const CYCLE: &str = "class Node:\n pass\n\
1447 a = Node()\n\
1448 b = Node()\n\
1449 a.other = b\n\
1450 b.other = a\n\
1451 del a\n\
1452 del b\n";
1453
1454 fn live_nodes(vm: &VirtualMachine) -> usize {
1455 vm.state
1456 .gc
1457 .get_objects(None)
1458 .iter()
1459 .filter(|obj| &*obj.class().name() == "Node")
1460 .count()
1461 }
1462
1463 let main = Interpreter::without_stdlib(Default::default());
1464 let sub = main.create_subinterpreter();
1465
1466 let sub_scope = sub.enter(|vm| {
1467 let scope = vm.new_scope_with_builtins();
1468 run(vm, &scope, CYCLE);
1469 assert_eq!(live_nodes(vm), 2);
1470 scope
1471 });
1472
1473 let deadline = Instant::now() + Duration::from_secs(30);
1481 while !main.enter(|vm| vm.state.gc.collect_force(2).candidates > 0) {
1482 assert!(
1483 Instant::now() < deadline,
1484 "no collection ran in the parent interpreter"
1485 );
1486 std::thread::sleep(Duration::from_millis(5));
1487 }
1488 sub.enter(|vm| assert_eq!(live_nodes(vm), 2));
1489
1490 sub.enter(|_| drop(sub_scope));
1491 }
1492
1493 #[test]
1495 fn get_objects_only_reports_the_calling_interpreter() {
1496 fn tracks_class(vm: &VirtualMachine, name: &str) -> bool {
1497 vm.state
1498 .gc
1499 .get_objects(None)
1500 .iter()
1501 .any(|obj| &*obj.class().name() == name)
1502 }
1503
1504 let main = Interpreter::without_stdlib(Default::default());
1505 let sub = main.create_subinterpreter();
1506
1507 let main_scope = main.enter(|vm| {
1508 let scope = vm.new_scope_with_builtins();
1509 run(vm, &scope, "class MainNode:\n pass\nkeep = MainNode()\n");
1510 scope
1511 });
1512 let sub_scope = sub.enter(|vm| {
1513 let scope = vm.new_scope_with_builtins();
1514 run(vm, &scope, "class SubNode:\n pass\nkeep = SubNode()\n");
1515 scope
1516 });
1517
1518 main.enter(|vm| {
1519 assert!(tracks_class(vm, "MainNode"));
1520 assert!(!tracks_class(vm, "SubNode"));
1521 });
1522 sub.enter(|vm| {
1523 assert!(tracks_class(vm, "SubNode"));
1524 assert!(!tracks_class(vm, "MainNode"));
1525 });
1526
1527 main.enter(|_| drop(main_scope));
1528 sub.enter(|_| drop(sub_scope));
1529 }
1530
1531 #[cfg(feature = "threading")]
1533 #[test]
1534 fn runtime_owned_interpreter_lifecycle() {
1535 let main = Interpreter::without_stdlib(Default::default());
1536 let sub = main.create_subinterpreter();
1537 let id = sub.id();
1538
1539 assert_eq!(runtime::store_owned_interpreter(sub), id);
1540 assert!(runtime::is_owned_interpreter(id));
1541 assert!(runtime::lookup_interpreter(id).is_some());
1542 assert!(runtime::owned_interpreter_count() >= 1);
1545
1546 let reclaimed = runtime::take_owned_interpreter(id).expect("owned by runtime");
1549 assert_eq!(reclaimed.id(), id);
1550 assert!(!runtime::is_owned_interpreter(id));
1551 assert!(runtime::lookup_interpreter(id).is_some());
1552 assert!(runtime::take_owned_interpreter(id).is_none());
1553
1554 drop(reclaimed);
1556 wait_until_unregistered(id);
1557 }
1558
1559 #[cfg(feature = "threading")]
1561 #[test]
1562 fn create_owned_subinterpreter_returns_id() {
1563 let main = Interpreter::without_stdlib(Default::default());
1564 let id = main.create_owned_subinterpreter();
1565 assert!(runtime::is_owned_interpreter(id));
1566 assert_ne!(id, main.id());
1567
1568 let sub = runtime::take_owned_interpreter(id).expect("owned by runtime");
1569 assert_eq!(sub.id(), id);
1570 assert!(!sub.is_main());
1571 }
1572
1573 #[cfg(feature = "threading")]
1576 #[test]
1577 fn owned_subinterpreter_cleanup_is_scoped_to_its_runtime() {
1578 let main1 = Interpreter::without_stdlib(Default::default());
1579 let main2 = Interpreter::without_stdlib(Default::default());
1580 let sub1 = main1.create_owned_subinterpreter();
1581 let sub2 = main2.create_owned_subinterpreter();
1582
1583 main1.enter(finalize_subinterpreters);
1584
1585 assert!(!runtime::is_owned_interpreter(sub1));
1586 assert!(runtime::is_owned_interpreter(sub2));
1587
1588 let sub2 = runtime::take_owned_interpreter(sub2).expect("owned by second runtime");
1589 let _ = sub2.finalize(None);
1590 }
1591
1592 #[cfg(all(feature = "threading", feature = "rustpython-compiler"))]
1597 #[test]
1598 fn gc_collect_is_safe_while_another_interpreter_runs() {
1599 use crate::compiler::Mode;
1600 use alloc::sync::Arc;
1601 use core::{
1602 sync::atomic::{AtomicBool, Ordering},
1603 time::Duration,
1604 };
1605 use std::time::Instant;
1606
1607 const CHURN: &str = "\
1610for _ in range(40):
1611 a = {}
1612 b = {'peer': a}
1613 a['peer'] = b
1614";
1615
1616 let main = Interpreter::without_stdlib(Default::default());
1617 let sub = main.create_subinterpreter();
1618 let stop = Arc::new(AtomicBool::new(false));
1619
1620 let run_source = |vm: &VirtualMachine, source: &str| {
1621 let scope = vm.new_scope_with_builtins();
1622 let code = vm
1623 .compile(source, Mode::Exec, "<churn>")
1624 .map_err(|err| err.into_pyexception(vm, Some(source)))
1625 .unwrap();
1626 vm.run_code_obj(code, scope).unwrap();
1627 };
1628
1629 let stop_worker = Arc::clone(&stop);
1631 let churner = sub.enter(|vm| {
1632 let thread_vm = vm.new_thread();
1633 std::thread::spawn(move || {
1634 thread_vm.run(|vm| {
1635 while !stop_worker.load(Ordering::Acquire) {
1636 run_source(vm, CHURN);
1637 }
1638 });
1639 })
1640 });
1641
1642 main.enter(|vm| {
1644 run_source(vm, CHURN);
1645 let deadline = Instant::now() + Duration::from_secs(2);
1646 let mut collections = 0;
1647 while Instant::now() < deadline && collections < 20 {
1648 vm.state.gc.collect_force(2);
1649 collections += 1;
1650 }
1651 assert!(collections > 0);
1652 });
1653
1654 stop.store(true, Ordering::Release);
1655 churner.join().expect("churn worker panicked");
1656 }
1657
1658 #[cfg(all(feature = "threading", feature = "rustpython-compiler"))]
1663 #[test]
1664 fn stop_the_world_parks_threads_of_another_interpreter() {
1665 use crate::compiler::Mode;
1666 use alloc::sync::Arc;
1667 use core::{
1668 sync::atomic::{AtomicBool, AtomicU64, Ordering},
1669 time::Duration,
1670 };
1671
1672 let main = Interpreter::without_stdlib(Default::default());
1673 let sub = main.create_subinterpreter();
1674 let sub_state = sub.enter(|vm| vm.state.clone());
1675
1676 let progress = Arc::new(AtomicU64::new(0));
1677 let stop = Arc::new(AtomicBool::new(false));
1678
1679 let progress_worker = Arc::clone(&progress);
1682 let stop_worker = Arc::clone(&stop);
1683 let worker = sub.enter(|vm| {
1684 let thread_vm = vm.new_thread();
1685 std::thread::spawn(move || {
1686 thread_vm.run(|vm| {
1687 let source = "x = 1 + 1\n";
1688 let code = vm
1689 .compile(source, Mode::Exec, "<spin>")
1690 .map_err(|err| err.into_pyexception(vm, Some(source)))
1691 .unwrap();
1692 while !stop_worker.load(Ordering::Acquire) {
1693 let scope = vm.new_scope_with_builtins();
1694 vm.run_code_obj(code.clone(), scope).unwrap();
1695 progress_worker.fetch_add(1, Ordering::Release);
1696 }
1697 });
1698 })
1699 });
1700
1701 while progress.load(Ordering::Acquire) == 0 {
1703 std::thread::yield_now();
1704 }
1705
1706 main.enter(|_vm| {
1707 sub_state.stop_the_world.stop_the_world(&sub_state);
1710
1711 let parked_at = progress.load(Ordering::Acquire);
1712 std::thread::sleep(Duration::from_millis(50));
1713 assert_eq!(
1714 progress.load(Ordering::Acquire),
1715 parked_at,
1716 "subinterpreter thread kept running while its world was stopped"
1717 );
1718
1719 sub_state.stop_the_world.start_the_world(&sub_state);
1720 });
1721
1722 let resumed_from = progress.load(Ordering::Acquire);
1724 while progress.load(Ordering::Acquire) == resumed_from {
1725 std::thread::yield_now();
1726 }
1727
1728 stop.store(true, Ordering::Release);
1729 worker.join().expect("worker panicked");
1730 }
1731
1732 #[cfg(all(feature = "threading", feature = "rustpython-compiler"))]
1735 #[test]
1736 fn stop_the_world_parks_when_global_stop_bit_is_cleared() {
1737 use crate::compiler::Mode;
1738 use alloc::sync::Arc;
1739 use core::{
1740 sync::atomic::{AtomicBool, AtomicU64, Ordering},
1741 time::Duration,
1742 };
1743
1744 let main = Interpreter::without_stdlib(Default::default());
1745 let sub = main.create_subinterpreter();
1746 let sub_state = sub.enter(|vm| vm.state.clone());
1747
1748 let progress = Arc::new(AtomicU64::new(0));
1749 let stop = Arc::new(AtomicBool::new(false));
1750 let flip = Arc::new(AtomicBool::new(true));
1751
1752 let progress_worker = Arc::clone(&progress);
1753 let stop_worker = Arc::clone(&stop);
1754 let worker = sub.enter(|vm| {
1755 let thread_vm = vm.new_thread();
1756 std::thread::spawn(move || {
1757 thread_vm.run(|vm| {
1758 let source = "x = 1 + 1\n";
1759 let code = vm
1760 .compile(source, Mode::Exec, "<spin>")
1761 .map_err(|err| err.into_pyexception(vm, Some(source)))
1762 .unwrap();
1763 while !stop_worker.load(Ordering::Acquire) {
1764 let scope = vm.new_scope_with_builtins();
1765 vm.run_code_obj(code.clone(), scope).unwrap();
1766 progress_worker.fetch_add(1, Ordering::Release);
1767 }
1768 });
1769 })
1770 });
1771
1772 let deadline = std::time::Instant::now() + Duration::from_secs(10);
1773 while progress.load(Ordering::Acquire) == 0 {
1774 assert!(
1775 std::time::Instant::now() < deadline,
1776 "worker never started making progress"
1777 );
1778 std::thread::yield_now();
1779 }
1780
1781 let flip_flag = Arc::clone(&flip);
1782 let flicker_cycles = Arc::new(AtomicU64::new(0));
1783 let flicker_cycles_thread = Arc::clone(&flicker_cycles);
1784 let flicker = std::thread::spawn(move || {
1785 while flip_flag.load(Ordering::Acquire) {
1786 crate::signal::set_stop_bit();
1787 crate::signal::clear_stop_bit();
1788 flicker_cycles_thread.fetch_add(1, Ordering::Release);
1789 }
1790 });
1791
1792 while flicker_cycles.load(Ordering::Acquire) == 0 {
1793 assert!(
1794 std::time::Instant::now() < deadline,
1795 "STOP_BIT flicker never ran"
1796 );
1797 std::thread::yield_now();
1798 }
1799
1800 let (tx, rx) = std::sync::mpsc::channel();
1803 let stopper = std::thread::spawn(move || {
1804 sub_state.stop_the_world.stop_the_world(&sub_state);
1805 let stopped = tx.send(());
1806 (sub_state, stopped)
1807 });
1808
1809 let stopped = rx.recv_timeout(Duration::from_secs(10));
1810
1811 let parked_at = progress.load(Ordering::Acquire);
1812 std::thread::sleep(Duration::from_millis(50));
1813 let parked_after = progress.load(Ordering::Acquire);
1814
1815 flip.store(false, Ordering::Release);
1818 if stopped.is_err() {
1819 stop.store(true, Ordering::Release);
1820 }
1821 flicker.join().expect("stop-bit flicker panicked");
1822 let (stop_state, sent) = stopper.join().expect("stopper panicked");
1823 if stopped.is_ok() {
1824 stop_state.stop_the_world.start_the_world(&stop_state);
1825 }
1826 sent.expect("send");
1827 assert!(
1828 stopped.is_ok(),
1829 "stop-the-world hung while STOP_BIT was flickering"
1830 );
1831 assert_eq!(
1832 parked_after, parked_at,
1833 "subinterpreter thread kept running after STOP_BIT was cleared"
1834 );
1835
1836 let resumed_from = progress.load(Ordering::Acquire);
1837 let resume_deadline = std::time::Instant::now() + Duration::from_secs(10);
1838 while progress.load(Ordering::Acquire) == resumed_from {
1839 assert!(
1840 std::time::Instant::now() < resume_deadline,
1841 "worker never resumed after start_the_world"
1842 );
1843 std::thread::yield_now();
1844 }
1845
1846 stop.store(true, Ordering::Release);
1847 worker.join().expect("worker panicked");
1848 }
1849
1850 #[cfg(feature = "threading")]
1854 #[test]
1855 fn eval_breaker_tripped_when_stop_requested() {
1856 let interp = Interpreter::without_stdlib(Default::default());
1857 interp.enter(|vm| {
1858 crate::signal::clear_eval_breaker_for_test();
1859 assert!(
1860 crate::vm::thread::set_stop_requested_for_current_thread(true),
1861 "current thread has no stop_requested flag"
1862 );
1863 let tripped = vm.eval_breaker_tripped();
1864 crate::vm::thread::set_stop_requested_for_current_thread(false);
1865 assert!(tripped);
1866 });
1867 }
1868
1869 #[cfg(all(feature = "threading", feature = "rustpython-compiler"))]
1875 #[test]
1876 fn nested_enter_of_subinterpreter_is_stoppable() {
1877 use crate::compiler::Mode;
1878 use alloc::sync::Arc;
1879 use core::{
1880 sync::atomic::{AtomicBool, AtomicU64, Ordering},
1881 time::Duration,
1882 };
1883
1884 let main = Interpreter::without_stdlib(Default::default());
1885 let sub = main.create_subinterpreter();
1886 let sub_state = sub.enter(|vm| vm.state.clone());
1887
1888 let progress = Arc::new(AtomicU64::new(0));
1889 let stop = Arc::new(AtomicBool::new(false));
1890
1891 let progress_worker = Arc::clone(&progress);
1893 let stop_worker = Arc::clone(&stop);
1894 let main_vm = main.enter(|vm| vm.new_thread());
1895 let sub_vm = sub.enter(|vm| vm.new_thread());
1896 let worker = std::thread::spawn(move || {
1897 main_vm.run(|_main| {
1898 sub_vm.run(|vm| {
1899 let source = "x = 1 + 1\n";
1900 let code = vm
1901 .compile(source, Mode::Exec, "<nested>")
1902 .map_err(|err| err.into_pyexception(vm, Some(source)))
1903 .unwrap();
1904 while !stop_worker.load(Ordering::Acquire) {
1905 let scope = vm.new_scope_with_builtins();
1906 vm.run_code_obj(code.clone(), scope).unwrap();
1907 progress_worker.fetch_add(1, Ordering::Release);
1908 }
1909 });
1910 });
1911 });
1912
1913 while progress.load(Ordering::Acquire) == 0 {
1914 std::thread::yield_now();
1915 }
1916
1917 sub_state.stop_the_world.stop_the_world(&sub_state);
1918 let parked_at = progress.load(Ordering::Acquire);
1919 std::thread::sleep(Duration::from_millis(50));
1920 assert_eq!(
1921 progress.load(Ordering::Acquire),
1922 parked_at,
1923 "nested subinterpreter thread kept running while the sub's world was stopped"
1924 );
1925 sub_state.stop_the_world.start_the_world(&sub_state);
1926
1927 let resumed_from = progress.load(Ordering::Acquire);
1928 while progress.load(Ordering::Acquire) == resumed_from {
1929 std::thread::yield_now();
1930 }
1931
1932 stop.store(true, Ordering::Release);
1933 worker.join().expect("nested worker panicked");
1934 }
1935
1936 #[cfg(feature = "threading")]
1943 #[test]
1944 fn a_thread_blocked_on_a_lock_does_not_stall_stop_the_world() {
1945 use super::super::thread::ThreadState;
1946 use crate::common::lock::PyDetachingRwLock;
1947 use alloc::sync::Arc;
1948 use core::{
1949 sync::atomic::{AtomicU64, Ordering},
1950 time::Duration,
1951 };
1952
1953 let interp = Interpreter::without_stdlib(Default::default());
1954 let state = interp.enter(|vm| vm.state.clone());
1955
1956 let lock: Arc<PyDetachingRwLock<()>> = Arc::new(PyDetachingRwLock::new(()));
1957 let worker_ident = Arc::new(AtomicU64::new(0));
1960
1961 let held = lock.write();
1963
1964 let worker_lock = Arc::clone(&lock);
1965 let published_ident = Arc::clone(&worker_ident);
1966 let worker = interp.enter(|vm| {
1967 let thread_vm = vm.new_thread();
1968 std::thread::spawn(move || {
1969 thread_vm.run(|_vm| {
1970 published_ident.store(crate::stdlib::_thread::get_ident(), Ordering::Release);
1971 let _read = worker_lock.read();
1972 });
1973 })
1974 });
1975
1976 let deadline = std::time::Instant::now() + Duration::from_secs(10);
1985 let blocked_detached = |ident| {
1986 state.thread_frames.lock().get(&ident).is_some_and(|slot| {
1987 slot.state.load(Ordering::Acquire) == ThreadState::Detached as i32
1988 })
1989 };
1990 loop {
1991 match worker_ident.load(Ordering::Acquire) {
1992 ident if ident != 0 && blocked_detached(ident) => break,
1993 _ => assert!(
1994 std::time::Instant::now() < deadline,
1995 "the worker never detached for the contended acquire"
1996 ),
1997 }
1998 std::thread::yield_now();
1999 }
2000
2001 let (tx, rx) = std::sync::mpsc::channel();
2004 let stop_state = state;
2005 let stopper = std::thread::spawn(move || {
2006 stop_state.stop_the_world.stop_the_world(&stop_state);
2007 let stopped = tx.send(());
2008 stop_state.stop_the_world.start_the_world(&stop_state);
2009 stopped
2010 });
2011
2012 let stopped = rx.recv_timeout(Duration::from_secs(10));
2013
2014 drop(held);
2017 assert!(
2018 stopped.is_ok(),
2019 "stop-the-world did not complete while a thread was blocked on a lock"
2020 );
2021 stopper.join().expect("stopper panicked").expect("send");
2022 worker.join().expect("worker panicked");
2023 }
2024
2025 #[cfg(feature = "threading")]
2034 #[test]
2035 fn a_callback_inside_a_detached_call_waits_for_the_world() {
2036 use alloc::sync::Arc;
2037 use core::{
2038 sync::atomic::{AtomicBool, Ordering},
2039 time::Duration,
2040 };
2041
2042 let interp = Interpreter::without_stdlib(Default::default());
2043 let state = interp.enter(|vm| vm.state.clone());
2044
2045 let detached = Arc::new(AtomicBool::new(false));
2046 let ran = Arc::new(AtomicBool::new(false));
2047 let go = Arc::new(AtomicBool::new(false));
2048
2049 let worker_detached = Arc::clone(&detached);
2050 let worker_ran = Arc::clone(&ran);
2051 let worker_go = Arc::clone(&go);
2052 let worker = interp.enter(|vm| {
2053 let thread_vm = vm.new_thread();
2054 std::thread::spawn(move || {
2055 thread_vm.run(|vm| {
2056 vm.allow_threads(|| {
2057 worker_detached.store(true, Ordering::Release);
2058 while !worker_go.load(Ordering::Acquire) {
2062 std::thread::yield_now();
2063 }
2064 vm.attach_for_callback(|| worker_ran.store(true, Ordering::Release));
2065 });
2066 });
2067 })
2068 });
2069
2070 while !detached.load(Ordering::Acquire) {
2071 std::thread::yield_now();
2072 }
2073
2074 let (stopped_tx, stopped_rx) = std::sync::mpsc::channel();
2077 let (release_tx, release_rx) = std::sync::mpsc::channel();
2078 let stop_state = state;
2079 let stopper = std::thread::spawn(move || {
2080 stop_state.stop_the_world.stop_the_world(&stop_state);
2081 stopped_tx.send(()).expect("send");
2082 release_rx.recv().expect("recv");
2083 stop_state.stop_the_world.start_the_world(&stop_state);
2084 });
2085
2086 stopped_rx
2087 .recv_timeout(Duration::from_secs(10))
2088 .expect("stop-the-world did not complete");
2089
2090 go.store(true, Ordering::Release);
2094 std::thread::sleep(Duration::from_millis(200));
2095 let ran_while_stopped = ran.load(Ordering::Acquire);
2096
2097 release_tx.send(()).expect("send");
2100 stopper.join().expect("stopper panicked");
2101 worker.join().expect("worker panicked");
2102
2103 assert!(
2104 !ran_while_stopped,
2105 "a callback ran Python while the world was stopped"
2106 );
2107 assert!(
2108 ran.load(Ordering::Acquire),
2109 "the callback never ran once the world started again"
2110 );
2111 }
2112
2113 #[test]
2115 fn process_main_id_recorded_and_stable() {
2116 let main = Interpreter::without_stdlib(Default::default());
2119 let recorded = runtime::main_interpreter_id().expect("a process main exists");
2120
2121 let _sub = main.create_subinterpreter();
2123 let _main2 = Interpreter::without_stdlib(Default::default());
2124 assert_eq!(runtime::main_interpreter_id(), Some(recorded));
2125 }
2126
2127 #[test]
2130 fn second_interpreter_reuses_process_hash_secret() {
2131 let settings = Settings {
2132 hash_seed: Some(7),
2133 ..Settings::default()
2134 };
2135 let first = Interpreter::without_stdlib(settings);
2136 first.enter(|vm| {
2137 vm.ctx
2138 .intern_str("zz_hash_seed_regression")
2139 .to_object()
2140 .hash(vm)
2141 .unwrap();
2142 });
2143
2144 let settings = Settings {
2145 hash_seed: Some(8),
2146 ..Settings::default()
2147 };
2148 let second = Interpreter::without_stdlib(settings);
2149 second.enter(|vm| {
2150 let interned = vm.ctx.intern_str("zz_hash_seed_regression").to_object();
2151 let dict = vm.ctx.new_dict();
2152 dict.set_item(&*interned, vm.ctx.new_int(1).into(), vm)
2153 .unwrap();
2154
2155 let scope = vm.new_scope_with_builtins();
2156 scope.globals.set_item("d", dict.into(), vm).unwrap();
2157 run(
2158 vm,
2159 &scope,
2160 "k = 'zz_hash_seed_' + 'regression'\nresult = k in d\n",
2161 );
2162 let result = scope.globals.get_item("result", vm).unwrap();
2163 assert!(result.try_to_bool(vm).unwrap());
2164 });
2165 }
2166}