From ac383c58ce5d54f32a0a293c3ec1ca627935fef1 Mon Sep 17 00:00:00 2001 From: Leilei Zhang Date: Wed, 30 Sep 2026 18:49:16 +0800 Subject: [PATCH 1/2] Guard reentrant Python apartment teardown Reject final apartment close during synchronous native callbacks and preserve anonymously dropped contexts for explicit same-thread recovery. Cover generated PropertySet MapChanged callbacks in isolated STA/MTA subprocesses. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> --- bindings/py/README.md | 26 ++ bindings/py/dynwinrt.pyi | 2 + bindings/py/src/async_runtime.rs | 3 +- bindings/py/src/implementation.rs | 5 +- bindings/py/src/runtime.rs | 258 ++++++++++++++++- tests/e2e/runners/py_apartment_reentrancy.py | 279 +++++++++++++++++++ tests/e2e/runners/py_runner.py | 37 +++ tests/e2e/typecheck/python_generated_api.py | 2 + 8 files changed, 597 insertions(+), 15 deletions(-) create mode 100644 tests/e2e/runners/py_apartment_reentrancy.py diff --git a/bindings/py/README.md b/bindings/py/README.md index 01e4f74d..008241d2 100644 --- a/bindings/py/README.md +++ b/bindings/py/README.md @@ -507,6 +507,32 @@ model are supported. Requesting a conflicting model raises `OSError` with each successful call, including `S_FALSE`, must be paired with one `ro_uninitialize()` call on the same thread. +Do not close the final managed apartment from a synchronous native-to-Python +callback (for example, a `PropertySet.map_changed` handler). `close()`, +`__exit__()`, and `ro_uninitialize()` raise `RuntimeError` **before** +`RoUninitialize` in that situation. A named `RoApartment` remains active; wait +for both the callback and its outer native call to return, release any retained +native callback values, then retry `apartment.close()` on the owner thread. +Non-final nested initializations may still be balanced inside the callback. +`ro_uninitialize()` cannot consume an active `RoApartment` initialization +without a matching `ro_initialize()` call. + +If a final `RoApartment` is dropped during a callback (including an unnamed +`with RoApartment():` whose `__exit__` failed), its initialization stays in a +same-thread pending lease rather than being uninitialized inside the native +stack. After the native call returns, recover the lease explicitly: + +```python +apartment = RoApartment.recover_pending() # on the original OS thread +# Release any remaining native PropertySet / callback references first. +apartment.close() +``` + +`recover_pending()` raises if called inside a callback or if no lease is +pending on that thread; calling it from another thread cannot consume the +owner's lease. It does not close or release anything automatically. Retain a +named apartment when possible so `.close()` can simply be retried. + WinRT is never initialized implicitly. A call on a thread without an apartment raises `OSError` with `CO_E_NOTINITIALIZED` in `error.winerror`; its message explains how to open one. diff --git a/bindings/py/dynwinrt.pyi b/bindings/py/dynwinrt.pyi index 973f5ac5..159da564 100644 --- a/bindings/py/dynwinrt.pyi +++ b/bindings/py/dynwinrt.pyi @@ -84,6 +84,8 @@ class RoApartment: self, exc_type: object, exc_value: object, traceback: object ) -> Literal[False]: ... def close(self) -> None: ... + @staticmethod + def recover_pending() -> RoApartment: ... def __repr__(self) -> str: ... @final diff --git a/bindings/py/src/async_runtime.rs b/bindings/py/src/async_runtime.rs index f8d2ce24..db82d24b 100644 --- a/bindings/py/src/async_runtime.rs +++ b/bindings/py/src/async_runtime.rs @@ -8,7 +8,7 @@ use std::sync::{Arc, Mutex, MutexGuard}; use crate::errors::{ map_dynwinrt_error, map_dynwinrt_error_with_context, map_windows_error_with_context, }; -use crate::runtime::DynWinRTValue; +use crate::runtime::{DynWinRTValue, NativeCallbackGuard}; use pyo3::exceptions::{PyRuntimeError, PyTypeError}; use pyo3::prelude::*; use pyo3::types::PyList; @@ -870,6 +870,7 @@ impl DynWinRTAsyncWithProgress { let weak_dispatcher = Arc::downgrade(&dispatcher); let progress_callback: dynwinrt::ProgressCallback = Box::new(move |value| { + let _callback_guard = NativeCallbackGuard::enter(); Python::attach(|py| { let Some(dispatcher) = weak_dispatcher.upgrade() else { return; diff --git a/bindings/py/src/implementation.rs b/bindings/py/src/implementation.rs index 1ca498c0..43476c5c 100644 --- a/bindings/py/src/implementation.rs +++ b/bindings/py/src/implementation.rs @@ -19,8 +19,8 @@ use windows::core::{Error, HRESULT}; use crate::errors::map_windows_error; use crate::runtime::{ - DynWinRTMethodSig, DynWinRTType, DynWinRTValue, PYWINRT_E_UNRAISABLE_PYTHON_EXCEPTION, WinGUID, - native_outputs, wrap_python_callback_context, + DynWinRTMethodSig, DynWinRTType, DynWinRTValue, NativeCallbackGuard, + PYWINRT_E_UNRAISABLE_PYTHON_EXCEPTION, WinGUID, native_outputs, wrap_python_callback_context, }; const RO_E_CLOSED: HRESULT = HRESULT(0x80000013_u32 as i32); @@ -192,6 +192,7 @@ impl CallbackCell { if self.interpreter.stopping.load(Ordering::Acquire) { return Err(closed_error()); } + let _callback_guard = NativeCallbackGuard::enter(); Python::try_attach(|py| { if self.interpreter.stopping.load(Ordering::Acquire) { return Err(closed_error()); diff --git a/bindings/py/src/runtime.rs b/bindings/py/src/runtime.rs index e03532bf..165d44e9 100644 --- a/bindings/py/src/runtime.rs +++ b/bindings/py/src/runtime.rs @@ -1,6 +1,9 @@ // Copyright (c) Microsoft Corporation. // Licensed under the MIT License. +use std::cell::{Cell, RefCell}; +use std::marker::PhantomData; +use std::rc::Rc; use std::sync::{Arc, Mutex}; use dynwinrt; @@ -117,6 +120,72 @@ fn set_typed_field( // Runtime initialization // ====================================================================== +#[derive(Default)] +struct ManagedApartments { + contexts: usize, + manual: usize, + pending: Vec, +} + +impl ManagedApartments { + fn count(&self) -> usize { + self.contexts + self.manual + self.pending.len() + } +} + +thread_local! { + static MANAGED_APARTMENTS: RefCell = RefCell::new(ManagedApartments::default()); + static NATIVE_CALLBACK_DEPTH: Cell = const { Cell::new(0) }; +} + +pub(crate) struct NativeCallbackGuard { + counted: bool, + _owner_thread: PhantomData>, +} + +impl NativeCallbackGuard { + pub(crate) fn enter() -> Self { + let counted = NATIVE_CALLBACK_DEPTH.with(|depth| { + let Some(next) = depth.get().checked_add(1) else { + return false; + }; + depth.set(next); + true + }); + Self { + counted, + _owner_thread: PhantomData, + } + } +} + +impl Drop for NativeCallbackGuard { + fn drop(&mut self) { + if self.counted { + NATIVE_CALLBACK_DEPTH.with(|depth| depth.set(depth.get() - 1)); + } + } +} + +fn in_native_callback() -> bool { + NATIVE_CALLBACK_DEPTH.with(|depth| depth.get() != 0) +} + +fn would_uninitialize_final_apartment() -> bool { + in_native_callback() && MANAGED_APARTMENTS.with(|state| state.borrow().count() <= 1) +} + +fn check_reentrant_uninitialize(operation: &str) -> PyResult<()> { + if would_uninitialize_final_apartment() { + return Err(PyRuntimeError::new_err(format!( + "{operation} cannot risk uninitializing the final COM apartment during a synchronous \ + native callback on this thread; retry after the callback and native invocation \ + return. Recover a dropped context with RoApartment.recover_pending() on this thread." + ))); + } + Ok(()) +} + #[pyclass] pub struct WinAppSDKContext(pub(crate) dynwinrt::WinAppSdkContext); @@ -132,6 +201,7 @@ impl WinAppSDKContext { pub struct RoApartment { apartment_type: i32, active: bool, + close_rejected: bool, } /// `apartment_type` used when Python omits it: the multithreaded apartment. @@ -162,21 +232,55 @@ impl RoApartment { )); } unsafe { RoInitialize(ro_init_type(self.apartment_type)) }.map_err(map_windows_error)?; + MANAGED_APARTMENTS.with(|state| state.borrow_mut().contexts += 1); self.active = true; Ok(()) } - fn uninitialize(&mut self) { - if self.active { - unsafe { windows::Win32::System::WinRT::RoUninitialize() }; - self.active = false; + fn uninitialize(&mut self, operation: &str) -> PyResult<()> { + if !self.active { + return Ok(()); + } + if let Err(error) = check_reentrant_uninitialize(operation) { + self.close_rejected = true; + return Err(error); } + MANAGED_APARTMENTS.with(|state| state.borrow_mut().contexts -= 1); + self.active = false; + self.close_rejected = false; + unsafe { windows::Win32::System::WinRT::RoUninitialize() }; + Ok(()) } } impl Drop for RoApartment { fn drop(&mut self) { - self.uninitialize(); + if !self.active { + return; + } + let pending = would_uninitialize_final_apartment(); + MANAGED_APARTMENTS.with(|state| { + let mut state = state.borrow_mut(); + state.contexts -= 1; + if pending { + state.pending.push(self.apartment_type); + } + }); + self.active = false; + if pending { + if !self.close_rejected { + Python::try_attach(|py| { + PyRuntimeError::new_err( + "RoApartment was dropped during a synchronous native callback; its \ + COM apartment remains active. Call RoApartment.recover_pending().close() \ + on this thread after the callback and native invocation return.", + ) + .write_unraisable(py, None); + }); + } + } else { + unsafe { windows::Win32::System::WinRT::RoUninitialize() }; + } } } @@ -188,6 +292,7 @@ impl RoApartment { Self { apartment_type: apartment_type.unwrap_or(DEFAULT_APARTMENT_TYPE), active: false, + close_rejected: false, } } @@ -201,13 +306,35 @@ impl RoApartment { _exc_type: &Bound<'_, PyAny>, _exc_value: &Bound<'_, PyAny>, _traceback: &Bound<'_, PyAny>, - ) -> bool { - self.uninitialize(); - false + ) -> PyResult { + self.uninitialize("RoApartment.__exit__()")?; + Ok(false) } - fn close(&mut self) { - self.uninitialize(); + fn close(&mut self) -> PyResult<()> { + self.uninitialize("RoApartment.close()") + } + + #[staticmethod] + fn recover_pending() -> PyResult { + if in_native_callback() { + return Err(PyRuntimeError::new_err( + "recover_pending() must be called on the owner thread after the synchronous \ + native callback and invocation return", + )); + } + MANAGED_APARTMENTS.with(|state| { + let mut state = state.borrow_mut(); + let apartment_type = state.pending.pop().ok_or_else(|| { + PyRuntimeError::new_err("no dropped RoApartment is pending on this thread") + })?; + state.contexts += 1; + Ok(Self { + apartment_type, + active: true, + close_rejected: false, + }) + }) } fn __repr__(&self) -> String { @@ -228,13 +355,30 @@ pub fn init_winappsdk(major: u32, minor: u32) -> PyResult { #[pyfunction] pub fn ro_initialize(apartment_type: Option) -> PyResult<()> { let init_type = ro_init_type(apartment_type.unwrap_or(DEFAULT_APARTMENT_TYPE)); - unsafe { RoInitialize(init_type) }.map_err(map_windows_error) + unsafe { RoInitialize(init_type) }.map_err(map_windows_error)?; + MANAGED_APARTMENTS.with(|state| state.borrow_mut().manual += 1); + Ok(()) } #[pyfunction] -pub fn ro_uninitialize() { +pub fn ro_uninitialize() -> PyResult<()> { use windows::Win32::System::WinRT::RoUninitialize; + check_reentrant_uninitialize("ro_uninitialize()")?; + MANAGED_APARTMENTS.with(|state| { + let mut state = state.borrow_mut(); + if state.manual == 0 && (state.contexts != 0 || !state.pending.is_empty()) { + return Err(PyRuntimeError::new_err( + "ro_uninitialize() has no matching ro_initialize() on this thread; close the \ + RoApartment or recover a dropped context with RoApartment.recover_pending()", + )); + } + if state.manual != 0 { + state.manual -= 1; + } + Ok(()) + })?; unsafe { RoUninitialize() }; + Ok(()) } // ====================================================================== @@ -348,6 +492,7 @@ pub fn register_xaml_runtime_class( 0x8001010Eu32 as i32, ))); } + let _callback_guard = NativeCallbackGuard::enter(); Python::attach(|py| { let result = (|| -> PyResult { let invocation_context = context.call_method0(py, "copy")?; @@ -724,6 +869,7 @@ impl DynWinRTOverrideInterface { if std::thread::current().id() != thread_id { return windows::core::HRESULT(0x8001010Eu32 as i32); } + let _callback_guard = NativeCallbackGuard::enter(); Python::attach(|py| { let result = (|| -> PyResult<()> { let invocation_context = context.call_method0(py, "copy")?; @@ -752,6 +898,7 @@ impl DynWinRTOverrideInterface { if std::thread::current().id() != thread_id { return windows::core::HRESULT(0x8001010Eu32 as i32); } + let _callback_guard = NativeCallbackGuard::enter(); Python::attach(|py| { let result = (|| -> PyResult<(f32, f32)> { let invocation_context = context.call_method0(py, "copy")?; @@ -1461,6 +1608,7 @@ impl DynWinRTValue { let callback = wrap_python_callback_context(py, callback)?; let progress_cb: dynwinrt::ProgressCallback = Box::new(move |val: dynwinrt::WinRTValue| { + let _callback_guard = NativeCallbackGuard::enter(); Python::attach(|py| { let result = (|| -> PyResult<()> { let py_val = Py::new(py, DynWinRTValue::new(val))?; @@ -2359,6 +2507,7 @@ fn create_python_delegate( ) -> PyResult { let delegate_callback: dynwinrt::delegate::DelegateCallback = Box::new(move |args: &[dynwinrt::WinRTValue]| { + let _callback_guard = NativeCallbackGuard::enter(); Python::attach(|py| { let result = (|| -> PyResult<()> { let py_args = args @@ -2480,6 +2629,7 @@ impl DynWinRtElementFactory { let get_callbacks = callbacks.clone(); let get_callback: dynwinrt::ElementFactoryGetCallback = Box::new(move |args| { + let _callback_guard = NativeCallbackGuard::enter(); Python::attach(|py| { let (callback, error_target) = { let callbacks = get_callbacks.lock().map_err(|_| E_FAIL)?; @@ -2507,6 +2657,7 @@ impl DynWinRtElementFactory { let recycle_callbacks = callbacks.clone(); let recycle_callback: dynwinrt::ElementFactoryRecycleCallback = Box::new(move |args| { + let _callback_guard = NativeCallbackGuard::enter(); Python::attach(|py| { let (callback, error_target) = { let callbacks = match recycle_callbacks.lock() { @@ -2606,6 +2757,89 @@ mod tests { use std::ffi::c_void; use std::sync::atomic::{AtomicU32, Ordering}; + #[test] + fn reentrant_apartment_guard_preserves_nested_and_manual_initializations() { + Python::initialize(); + std::thread::spawn(|| { + Python::attach(|py| { + let mut apartment = RoApartment::new(Some(RO_INIT_SINGLETHREADED.0)); + apartment.initialize().unwrap(); + let outer_callback = NativeCallbackGuard::enter(); + let inner_callback = NativeCallbackGuard::enter(); + let error = apartment.uninitialize("RoApartment.close()").unwrap_err(); + assert!(error.is_instance_of::(py)); + assert!(error.to_string().contains("retry after the callback")); + assert!(apartment.active); + + let mut nested = RoApartment::new(Some(RO_INIT_SINGLETHREADED.0)); + nested.initialize().unwrap(); + nested.uninitialize("RoApartment.close()").unwrap(); + ro_initialize(Some(RO_INIT_SINGLETHREADED.0)).unwrap(); + ro_uninitialize().unwrap(); + assert!(ro_uninitialize().is_err()); + assert!(apartment.active); + drop(inner_callback); + assert!(apartment.uninitialize("RoApartment.close()").is_err()); + drop(outer_callback); + + assert!(ro_uninitialize().is_err()); + apartment.uninitialize("RoApartment.close()").unwrap(); + assert!(!apartment.active); + apartment.uninitialize("RoApartment.close()").unwrap(); + let mut other_model = RoApartment::new(Some(RO_INIT_MULTITHREADED.0)); + other_model.initialize().unwrap(); + other_model.uninitialize("RoApartment.close()").unwrap(); + }); + }) + .join() + .unwrap(); + } + + #[test] + fn dropped_final_apartment_requires_owner_thread_recovery() { + Python::initialize(); + std::thread::spawn(|| { + Python::attach(|py| { + let mut apartment = RoApartment::new(Some(RO_INIT_SINGLETHREADED.0)); + apartment.initialize().unwrap(); + let callback = NativeCallbackGuard::enter(); + apartment + .uninitialize("RoApartment.__exit__()") + .unwrap_err(); + drop(apartment); + assert!(RoApartment::recover_pending().is_err()); + drop(callback); + + let wrong_thread_error = py.detach(|| { + std::thread::spawn(|| { + Python::attach(|py| { + RoApartment::recover_pending() + .err() + .expect("wrong-thread recovery should fail") + .value(py) + .str() + .unwrap() + .to_str() + .unwrap() + .to_owned() + }) + }) + .join() + .unwrap() + }); + assert!(wrong_thread_error.contains("on this thread")); + + let mut recovered = RoApartment::recover_pending().unwrap(); + assert!(recovered.active); + assert!(RoApartment::recover_pending().is_err()); + recovered.uninitialize("RoApartment.close()").unwrap(); + assert!(!recovered.active); + }); + }) + .join() + .unwrap(); + } + #[derive(Default)] struct QueryCounts { queries: AtomicU32, diff --git a/tests/e2e/runners/py_apartment_reentrancy.py b/tests/e2e/runners/py_apartment_reentrancy.py new file mode 100644 index 00000000..d55dbb81 --- /dev/null +++ b/tests/e2e/runners/py_apartment_reentrancy.py @@ -0,0 +1,279 @@ +# Copyright (c) Microsoft Corporation. +# Licensed under the MIT License. + +"""Isolated, real PropertySet MapChanged/apartment lifetime regressions.""" + +import argparse +import os +import sys +import threading + + +def blocked_close(action): + try: + action() + except RuntimeError as error: + message = str(error) + assert 'synchronous native callback' in message, message + assert 'retry after the callback' in message, message + else: + raise AssertionError('final apartment uninitialization succeeded in a native callback') + + +def run_case(generated, mode, scenario): + sys.path.insert(0, os.path.dirname(os.path.abspath(generated))) + import dynwinrt as dw + from python_bindings import PropertySet + + def switch_models(): + with dw.RoApartment(1 - mode) as opposite: + assert 'active=true' in repr(opposite) + + def dispose(mapping, token): + mapping.off_map_changed(token) + dw.release_projected(mapping) + + if scenario in ('close', 'exit'): + apartment = dw.RoApartment(mode) + apartment.__enter__() + events = [] + retained = [] + + def on_changed(sender, args): + events.append(args.key) + if not retained: + retained.extend((sender, args)) + assert sender[args.key] is None + if scenario == 'close': + blocked_close(apartment.close) + else: + blocked_close(lambda: apartment.__exit__(None, None, None)) + assert 'active=true' in repr(apartment) + assert args.key in sender + + mapping = PropertySet() + token = mapping.on_map_changed(on_changed) + mapping.insert('first', None) + mapping.insert('second', None) + assert events == ['first', 'second'], events + sender, args = retained + assert args.key == 'first' and sender['first'] is None + dispose(mapping, token) + dw.release_projected(args) + dw.release_projected(sender) + if scenario == 'close': + apartment.close() + apartment.close() + else: + assert apartment.__exit__(None, None, None) is False + assert 'active=false' in repr(apartment) + switch_models() + + elif scenario == 'nested': + outer = dw.RoApartment(mode) + inner = dw.RoApartment(mode) + outer.__enter__() + inner.__enter__() + events = [] + + def on_changed(sender, args): + events.append(args.key) + inner.close() + assert 'active=false' in repr(inner) + blocked_close(outer.close) + assert 'active=true' in repr(outer) + assert sender[args.key] is None + + mapping = PropertySet() + token = mapping.on_map_changed(on_changed) + mapping.insert('nested', None) + assert events == ['nested'], events + dispose(mapping, token) + outer.close() + switch_models() + + elif scenario == 'manual': + dw.ro_initialize(mode) + dw.ro_initialize(mode) + events = [] + + def on_changed(sender, args): + events.append(args.key) + dw.ro_uninitialize() + blocked_close(dw.ro_uninitialize) + assert sender[args.key] is None + + mapping = PropertySet() + token = mapping.on_map_changed(on_changed) + mapping.insert('manual', None) + assert events == ['manual'], events + dispose(mapping, token) + dw.ro_uninitialize() + switch_models() + + elif scenario == 'pending': + outer = dw.RoApartment(mode) + outer.__enter__() + events = [] + + def on_changed(sender, args): + events.append(args.key) + + def unnamed_context(): + with dw.RoApartment(mode): + outer.close() + + blocked_close(unnamed_context) + assert 'active=false' in repr(outer) + try: + dw.RoApartment.recover_pending() + except RuntimeError as error: + assert 'after' in str(error), error + else: + raise AssertionError('pending apartment recovered inside its native callback') + assert sender[args.key] is None + + mapping = PropertySet() + token = mapping.on_map_changed(on_changed) + mapping.insert('pending', None) + assert events == ['pending'], events + + wrong_thread = [] + + def recover_on_wrong_thread(): + try: + dw.RoApartment.recover_pending() + except RuntimeError as error: + wrong_thread.append(str(error)) + + thread = threading.Thread(target=recover_on_wrong_thread) + thread.start() + thread.join(timeout=5) + assert not thread.is_alive() and len(wrong_thread) == 1, wrong_thread + assert 'on this thread' in wrong_thread[0], wrong_thread + + recovered = dw.RoApartment.recover_pending() + assert 'active=true' in repr(recovered) + try: + dw.RoApartment.recover_pending() + except RuntimeError as error: + assert 'no dropped RoApartment' in str(error), error + else: + raise AssertionError('recovered the same pending apartment twice') + dispose(mapping, token) + recovered.close() + switch_models() + + elif scenario == 'drop': + outer = dw.RoApartment(mode) + outer.__enter__() + inner = dw.RoApartment(mode) + inner.__enter__() + unraisable = [] + original_hook = sys.unraisablehook + sys.unraisablehook = unraisable.append + + def on_changed(sender, args): + outer.close() + inner_holder.pop() + assert sender[args.key] is None + + inner_holder = [inner] + del inner + mapping = PropertySet() + token = mapping.on_map_changed(on_changed) + try: + mapping.insert('dropped', None) + finally: + sys.unraisablehook = original_hook + assert len(unraisable) == 1, unraisable + assert 'recover_pending()' in str(unraisable[0].exc_value) + unraisable.clear() + recovered = dw.RoApartment.recover_pending() + dispose(mapping, token) + recovered.close() + switch_models() + + elif scenario == 'pending_error': + outer = dw.RoApartment(mode) + outer.__enter__() + unraisable = [] + original_hook = sys.unraisablehook + sys.unraisablehook = unraisable.append + + def on_changed(_sender, _args): + with dw.RoApartment(mode): + outer.close() + + mapping = PropertySet() + token = mapping.on_map_changed(on_changed) + try: + try: + mapping.insert('pending-error', None) + except OSError: + pass + finally: + sys.unraisablehook = original_hook + assert len(unraisable) == 1, unraisable + assert isinstance(unraisable[0].exc_value, RuntimeError) + assert 'synchronous native callback' in str(unraisable[0].exc_value) + unraisable.clear() + assert 'active=false' in repr(outer) + recovered = dw.RoApartment.recover_pending() + dispose(mapping, token) + recovered.close() + switch_models() + + elif scenario == 'exception': + apartment = dw.RoApartment(mode) + apartment.__enter__() + unraisable = [] + original_hook = sys.unraisablehook + sys.unraisablehook = unraisable.append + + def on_changed(_sender, _args): + blocked_close(apartment.close) + raise ValueError('handler failed after blocked close') + + mapping = PropertySet() + token = mapping.on_map_changed(on_changed) + try: + try: + mapping.insert('failed', None) + except OSError: + pass + finally: + sys.unraisablehook = original_hook + assert len(unraisable) == 1, unraisable + assert isinstance(unraisable[0].exc_value, ValueError) + assert 'handler failed after blocked close' in str(unraisable[0].exc_value) + mapping.off_map_changed(token) + unraisable.clear() + events = [] + token = mapping.on_map_changed(lambda _sender, args: events.append(args.key)) + mapping.insert('after', None) + assert events == ['after'], events + dispose(mapping, token) + apartment.close() + switch_models() + + else: + raise AssertionError(f'unknown apartment scenario: {scenario}') + + print(f'PASS {scenario} apartment={mode}', flush=True) + + +if __name__ == '__main__': + parser = argparse.ArgumentParser() + parser.add_argument('--generated', required=True) + parser.add_argument('--mode', type=int, choices=[0, 1], required=True) + parser.add_argument( + '--scenario', + choices=[ + 'close', 'exit', 'nested', 'manual', 'pending', 'drop', + 'pending_error', 'exception', + ], + required=True, + ) + args = parser.parse_args() + run_case(args.generated, args.mode, args.scenario) diff --git a/tests/e2e/runners/py_runner.py b/tests/e2e/runners/py_runner.py index 7b807cbc..3d24fd40 100644 --- a/tests/e2e/runners/py_runner.py +++ b/tests/e2e/runners/py_runner.py @@ -19,6 +19,7 @@ import re import sys import os +import subprocess import threading @@ -1414,6 +1415,42 @@ def observe_null(sender, args): dw.release_projected(null_events[0][1]) dw.release_projected(null_events[0][0]) del obj['null-element'] + if cls.__name__ == 'PropertySet': + runner = os.path.join( + os.path.dirname(__file__), 'py_apartment_reentrancy.py' + ) + for mode in (dw.RO_INIT_SINGLETHREADED, dw.RO_INIT_MULTITHREADED): + for scenario in ( + 'close', 'exit', 'nested', 'manual', 'pending', + 'drop', 'pending_error', 'exception', + ): + try: + child = subprocess.run( + [ + sys.executable, '-I', runner, + '--generated', generated_dir, + '--mode', str(mode), + '--scenario', scenario, + ], + capture_output=True, text=True, + timeout=30, check=False, + ) + except subprocess.TimeoutExpired as error: + cr['error'] = ( + f'{scenario} apartment={mode} timed out in ' + f'a native callback: {error}' + ) + return cr + if ( + child.returncode + or f'PASS {scenario} apartment={mode}' not in child.stdout + ): + cr['error'] = ( + f'{scenario} apartment={mode} failed with ' + f'0x{child.returncode & 0xffffffff:08X}: ' + f'{child.stdout} {child.stderr}' + ) + return cr cr['pass'] = True elif kind == 'work_item_callback_passthrough': diff --git a/tests/e2e/typecheck/python_generated_api.py b/tests/e2e/typecheck/python_generated_api.py index d7427187..d2b61024 100644 --- a/tests/e2e/typecheck/python_generated_api.py +++ b/tests/e2e/typecheck/python_generated_api.py @@ -17,6 +17,7 @@ DynWinRTMethodSig, DynWinRTType, DynWinRTValue, + RoApartment, WinGUID, ) from dynwinrt.values import ( @@ -92,6 +93,7 @@ def check_structural_async_compatibility() -> None: def check_runtime_stubs() -> None: + assert_type(RoApartment.recover_pending(), RoApartment) iid: WinGUID = WinGUID.parse("00000000-0000-0000-c000-000000000046") interface: DynWinRTType = DynWinRTType.register_interface("IUnknown", iid) signature: DynWinRTMethodSig = DynWinRTMethodSig().add_out( From 33bac688524e234599a9b6ebe882b9271c76abb1 Mon Sep 17 00:00:00 2001 From: Leilei Zhang Date: Wed, 30 Sep 2026 19:31:11 +0800 Subject: [PATCH 2/2] Keep Python apartment leases on their creating OS thread Reject foreign-thread apartment operations before native teardown. Retain foreign drops in an owner-recoverable queue without attaching Python, and fail closed during owner TLS teardown. Cover both architecture callback paths and exceptional queue states. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com> --- bindings/py/README.md | 19 +- bindings/py/src/runtime.rs | 541 +++++++++++++++++-- tests/e2e/runners/py_apartment_reentrancy.py | 105 +++- tests/e2e/runners/py_runner.py | 22 + 4 files changed, 613 insertions(+), 74 deletions(-) diff --git a/bindings/py/README.md b/bindings/py/README.md index 008241d2..88aea086 100644 --- a/bindings/py/README.md +++ b/bindings/py/README.md @@ -505,7 +505,9 @@ with RoApartment(RO_INIT_SINGLETHREADED): model are supported. Requesting a conflicting model raises `OSError` with `RPC_E_CHANGED_MODE`. The low-level `ro_initialize()` API remains available, but each successful call, including `S_FALSE`, must be paired with one -`ro_uninitialize()` call on the same thread. +`ro_uninitialize()` call on the same thread. `ro_uninitialize()` rejects calls +without a matching `dynwinrt.ro_initialize()` on that thread, including calls +intended to balance an initialization made by another library. Do not close the final managed apartment from a synchronous native-to-Python callback (for example, a `PropertySet.map_changed` handler). `close()`, @@ -514,8 +516,7 @@ callback (for example, a `PropertySet.map_changed` handler). `close()`, for both the callback and its outer native call to return, release any retained native callback values, then retry `apartment.close()` on the owner thread. Non-final nested initializations may still be balanced inside the callback. -`ro_uninitialize()` cannot consume an active `RoApartment` initialization -without a matching `ro_initialize()` call. +`ro_uninitialize()` cannot consume a `RoApartment` initialization. If a final `RoApartment` is dropped during a callback (including an unnamed `with RoApartment():` whose `__exit__` failed), its initialization stays in a @@ -531,7 +532,17 @@ apartment.close() `recover_pending()` raises if called inside a callback or if no lease is pending on that thread; calling it from another thread cannot consume the owner's lease. It does not close or release anything automatically. Retain a -named apartment when possible so `.close()` can simply be retried. +named apartment when possible so `.close()` can simply be retried. An implicit +drop without an earlier failed close emits a native stderr diagnostic. + +An apartment is bound to the OS thread on which `RoApartment()` was created. +Calling `__enter__()`, `close()`, or `__exit__()` on another thread raises +`RuntimeError` without changing its state; retry on the creating thread. +If its last Python reference is instead dropped on a different thread, native +uninitialization is **not** attempted there: a native stderr diagnostic identifies +the pending lease, which the creating thread can explicitly retrieve using +`RoApartment.recover_pending()` and close after native references are released. +Do not depend on Python garbage collection to close an apartment. WinRT is never initialized implicitly. A call on a thread without an apartment raises `OSError` with `CO_E_NOTINITIALIZED` in `error.winerror`; its message diff --git a/bindings/py/src/runtime.rs b/bindings/py/src/runtime.rs index 165d44e9..ad2e0b87 100644 --- a/bindings/py/src/runtime.rs +++ b/bindings/py/src/runtime.rs @@ -4,7 +4,8 @@ use std::cell::{Cell, RefCell}; use std::marker::PhantomData; use std::rc::Rc; -use std::sync::{Arc, Mutex}; +use std::sync::{Arc, Mutex, MutexGuard}; +use std::thread::{self, ThreadId}; use dynwinrt; use pyo3::exceptions::{PyIndexError, PyOverflowError, PyRuntimeError, PyTypeError}; @@ -133,9 +134,61 @@ impl ManagedApartments { } } +struct ForeignDropQueue { + owner_alive: bool, + pending: Vec, +} + +struct OwnerDropQueue(Arc>); + +impl OwnerDropQueue { + fn new() -> Self { + Self(Arc::new(Mutex::new(ForeignDropQueue { + owner_alive: true, + pending: Vec::new(), + }))) + } +} + +impl Drop for OwnerDropQueue { + fn drop(&mut self) { + let (pending, poisoned) = { + let (mut queue, poisoned) = lock_foreign_drops(&self.0); + queue.owner_alive = false; + (queue.pending.len(), poisoned) + }; + if poisoned { + native_apartment_diagnostic("RoApartment owner-thread pending state was poisoned"); + } + if pending != 0 { + native_apartment_diagnostic(&format!( + "{} dropped RoApartment initialization(s) were not recovered before their \ + owner OS thread exited", + pending + )); + } + } +} + +fn lock_foreign_drops(queue: &Mutex) -> (MutexGuard<'_, ForeignDropQueue>, bool) { + match queue.lock() { + Ok(guard) => (guard, false), + Err(poisoned) => { + queue.clear_poison(); + (poisoned.into_inner(), true) + } + } +} + +fn native_apartment_diagnostic(message: &str) { + use std::io::Write; + let _ = writeln!(std::io::stderr(), "{message}"); +} + thread_local! { static MANAGED_APARTMENTS: RefCell = RefCell::new(ManagedApartments::default()); static NATIVE_CALLBACK_DEPTH: Cell = const { Cell::new(0) }; + static FOREIGN_DROPS: OwnerDropQueue = OwnerDropQueue::new(); } pub(crate) struct NativeCallbackGuard { @@ -167,16 +220,35 @@ impl Drop for NativeCallbackGuard { } } -fn in_native_callback() -> bool { - NATIVE_CALLBACK_DEPTH.with(|depth| depth.get() != 0) -} - -fn would_uninitialize_final_apartment() -> bool { - in_native_callback() && MANAGED_APARTMENTS.with(|state| state.borrow().count() <= 1) -} - fn check_reentrant_uninitialize(operation: &str) -> PyResult<()> { - if would_uninitialize_final_apartment() { + let callback_active = NATIVE_CALLBACK_DEPTH + .try_with(|depth| depth.get() != 0) + .map_err(|_| { + PyRuntimeError::new_err(format!( + "{operation}: owner-thread callback state is unavailable during teardown" + )) + })?; + let final_apartment = if callback_active { + MANAGED_APARTMENTS + .try_with(|state| { + state + .try_borrow() + .map(|state| state.count() <= 1) + .map_err(|_| { + PyRuntimeError::new_err( + "owner-thread apartment state is already in use; retry after the callback", + ) + }) + }) + .map_err(|_| { + PyRuntimeError::new_err( + "owner-thread apartment state is unavailable during teardown", + ) + })?? + } else { + false + }; + if final_apartment { return Err(PyRuntimeError::new_err(format!( "{operation} cannot risk uninitializing the final COM apartment during a synchronous \ native callback on this thread; retry after the callback and native invocation \ @@ -197,11 +269,13 @@ impl WinAppSDKContext { } } -#[pyclass(unsendable)] +#[pyclass] pub struct RoApartment { apartment_type: i32, active: bool, close_rejected: bool, + owner_thread: ThreadId, + foreign_drops: Arc>, } /// `apartment_type` used when Python omits it: the multithreaded apartment. @@ -225,7 +299,30 @@ fn ro_init_type(apartment_type: i32) -> RO_INIT_TYPE { } impl RoApartment { + fn enqueue_owner_recovery(&self) -> (bool, bool) { + let (mut queue, poisoned) = lock_foreign_drops(&self.foreign_drops); + if queue.owner_alive { + // The owner's contexts count still owns this native initialization. + queue.pending.push(self.apartment_type); + (true, poisoned) + } else { + (false, poisoned) + } + } + + fn check_owner_thread(&self, operation: &str) -> PyResult<()> { + if thread::current().id() != self.owner_thread { + return Err(PyRuntimeError::new_err(format!( + "{operation} must run on the OS thread where RoApartment was created; \ + its apartment is unchanged. Retry on that owner thread, or recover a \ + guard dropped on another thread with RoApartment.recover_pending()." + ))); + } + Ok(()) + } + fn initialize(&mut self) -> PyResult<()> { + self.check_owner_thread("RoApartment.__enter__()")?; if self.active { return Err(PyRuntimeError::new_err( "the COM apartment context is already active", @@ -238,6 +335,7 @@ impl RoApartment { } fn uninitialize(&mut self, operation: &str) -> PyResult<()> { + self.check_owner_thread(operation)?; if !self.active { return Ok(()); } @@ -245,7 +343,26 @@ impl RoApartment { self.close_rejected = true; return Err(error); } - MANAGED_APARTMENTS.with(|state| state.borrow_mut().contexts -= 1); + MANAGED_APARTMENTS + .try_with(|state| { + let mut state = state.try_borrow_mut().map_err(|_| { + PyRuntimeError::new_err( + "owner-thread apartment state is already in use; retry close()", + ) + })?; + if state.contexts == 0 { + return Err(PyRuntimeError::new_err( + "RoApartment has no matching initialization on its owner thread", + )); + } + state.contexts -= 1; + Ok(()) + }) + .map_err(|_| { + PyRuntimeError::new_err( + "owner-thread apartment state is unavailable during teardown", + ) + })??; self.active = false; self.close_rejected = false; unsafe { windows::Win32::System::WinRT::RoUninitialize() }; @@ -258,28 +375,78 @@ impl Drop for RoApartment { if !self.active { return; } - let pending = would_uninitialize_final_apartment(); - MANAGED_APARTMENTS.with(|state| { - let mut state = state.borrow_mut(); + if thread::current().id() != self.owner_thread { + let (queued, poisoned) = self.enqueue_owner_recovery(); + self.active = false; + if poisoned { + native_apartment_diagnostic( + "RoApartment owner-thread pending state was poisoned; its native \ + initialization remains recoverable on the owner thread.", + ); + } + let message = if queued { + "RoApartment was dropped on a different OS thread; its COM apartment remains \ + active on the creating thread. Call RoApartment.recover_pending().close() \ + there after native calls and retained references have been released." + } else { + "RoApartment was dropped after its owner OS thread exited; its COM apartment \ + cannot be recovered on a different thread." + }; + native_apartment_diagnostic(message); + return; + } + let disposition = MANAGED_APARTMENTS.try_with(|state| { + let Ok(mut state) = state.try_borrow_mut() else { + return None; + }; + if state.contexts == 0 { + return None; + } + let callback_active = NATIVE_CALLBACK_DEPTH.try_with(|depth| depth.get() != 0); + let pending = match callback_active { + Ok(active) => active && state.count() <= 1, + Err(_) => true, + }; state.contexts -= 1; if pending { state.pending.push(self.apartment_type); } + Some((pending, callback_active.is_err())) }); self.active = false; - if pending { - if !self.close_rejected { - Python::try_attach(|py| { - PyRuntimeError::new_err( - "RoApartment was dropped during a synchronous native callback; its \ - COM apartment remains active. Call RoApartment.recover_pending().close() \ - on this thread after the callback and native invocation return.", - ) - .write_unraisable(py, None); - }); + match disposition { + Ok(Some((true, true))) => native_apartment_diagnostic( + "RoApartment callback state was unavailable during owner-thread teardown; \ + its COM initialization remains pending without native uninitialization.", + ), + Ok(Some((true, false))) if !self.close_rejected => native_apartment_diagnostic( + "RoApartment was dropped during a synchronous native callback; its \ + COM apartment remains active. Call RoApartment.recover_pending().close() \ + on this thread after the callback and native invocation return.", + ), + Ok(Some((true, false))) => {} + Ok(Some((false, false))) => { + unsafe { windows::Win32::System::WinRT::RoUninitialize() }; + } + _ => { + let (queued, poisoned) = self.enqueue_owner_recovery(); + if poisoned { + native_apartment_diagnostic( + "RoApartment owner-thread pending state was poisoned during teardown", + ); + } + if queued { + native_apartment_diagnostic( + "RoApartment owner-thread apartment state was unavailable during Drop; \ + native initialization remains pending for recover_pending()", + ); + } else { + native_apartment_diagnostic( + "RoApartment owner-thread state and recovery queue were unavailable \ + during teardown; native uninitialization was not attempted", + ); + } } - } else { - unsafe { windows::Win32::System::WinRT::RoUninitialize() }; } } } @@ -288,12 +455,21 @@ impl Drop for RoApartment { impl RoApartment { #[new] #[pyo3(signature = (apartment_type=None))] - fn new(apartment_type: Option) -> Self { - Self { + fn new(apartment_type: Option) -> PyResult { + let foreign_drops = FOREIGN_DROPS + .try_with(|queue| queue.0.clone()) + .map_err(|_| { + PyRuntimeError::new_err( + "RoApartment cannot be created after its owner-thread state was destroyed", + ) + })?; + Ok(Self { apartment_type: apartment_type.unwrap_or(DEFAULT_APARTMENT_TYPE), active: false, close_rejected: false, - } + owner_thread: thread::current().id(), + foreign_drops, + }) } fn __enter__(mut slf: PyRefMut<'_, Self>) -> PyResult> { @@ -317,31 +493,75 @@ impl RoApartment { #[staticmethod] fn recover_pending() -> PyResult { - if in_native_callback() { + let callback_active = NATIVE_CALLBACK_DEPTH + .try_with(|depth| depth.get() != 0) + .map_err(|_| { + PyRuntimeError::new_err( + "owner-thread callback state is unavailable during teardown", + ) + })?; + if callback_active { return Err(PyRuntimeError::new_err( "recover_pending() must be called on the owner thread after the synchronous \ native callback and invocation return", )); } - MANAGED_APARTMENTS.with(|state| { - let mut state = state.borrow_mut(); - let apartment_type = state.pending.pop().ok_or_else(|| { - PyRuntimeError::new_err("no dropped RoApartment is pending on this thread") + let foreign_drops = FOREIGN_DROPS + .try_with(|queue| queue.0.clone()) + .map_err(|_| { + PyRuntimeError::new_err( + "owner-thread apartment recovery queue is unavailable during teardown", + ) })?; - state.contexts += 1; - Ok(Self { - apartment_type, - active: true, - close_rejected: false, + let local = MANAGED_APARTMENTS + .try_with(|state| -> PyResult> { + let mut state = state.try_borrow_mut().map_err(|_| { + PyRuntimeError::new_err("owner-thread apartment state is already in use") + })?; + let apartment_type = state.pending.pop(); + if apartment_type.is_some() { + state.contexts += 1; + } + Ok(apartment_type) }) + .map_err(|_| { + PyRuntimeError::new_err( + "owner-thread apartment state is unavailable during teardown", + ) + })??; + let apartment_type = match local { + Some(apartment_type) => apartment_type, + None => { + let (pending, poisoned) = { + let (mut queue, poisoned) = lock_foreign_drops(&foreign_drops); + (queue.pending.pop(), poisoned) + }; + if poisoned { + native_apartment_diagnostic( + "RoApartment owner-thread pending state was poisoned; recovering its \ + native initialization before any COM uninitialization.", + ); + } + pending.ok_or_else(|| { + PyRuntimeError::new_err("no dropped RoApartment is pending on this thread") + })? + } + }; + Ok(Self { + apartment_type, + active: true, + close_rejected: false, + owner_thread: thread::current().id(), + foreign_drops, }) } - fn __repr__(&self) -> String { - format!( + fn __repr__(&self) -> PyResult { + self.check_owner_thread("RoApartment.__repr__()")?; + Ok(format!( "RoApartment(apartment_type={}, active={})", self.apartment_type, self.active - ) + )) } } @@ -364,19 +584,23 @@ pub fn ro_initialize(apartment_type: Option) -> PyResult<()> { pub fn ro_uninitialize() -> PyResult<()> { use windows::Win32::System::WinRT::RoUninitialize; check_reentrant_uninitialize("ro_uninitialize()")?; - MANAGED_APARTMENTS.with(|state| { - let mut state = state.borrow_mut(); - if state.manual == 0 && (state.contexts != 0 || !state.pending.is_empty()) { - return Err(PyRuntimeError::new_err( - "ro_uninitialize() has no matching ro_initialize() on this thread; close the \ - RoApartment or recover a dropped context with RoApartment.recover_pending()", - )); - } - if state.manual != 0 { + MANAGED_APARTMENTS + .try_with(|state| { + let mut state = state.try_borrow_mut().map_err(|_| { + PyRuntimeError::new_err("owner-thread apartment state is already in use") + })?; + if state.manual == 0 { + return Err(PyRuntimeError::new_err( + "ro_uninitialize() has no matching ro_initialize() on this thread; close the \ + RoApartment or recover a dropped context with RoApartment.recover_pending()", + )); + } state.manual -= 1; - } - Ok(()) - })?; + Ok(()) + }) + .map_err(|_| { + PyRuntimeError::new_err("owner-thread apartment state is unavailable during teardown") + })??; unsafe { RoUninitialize() }; Ok(()) } @@ -2756,13 +2980,36 @@ mod tests { use pyo3::types::PyDict; use std::ffi::c_void; use std::sync::atomic::{AtomicU32, Ordering}; + use std::time::Duration; + + struct LateDropProbe { + apartment: Option, + state_was_destroyed: Arc, + creation_was_rejected: Arc, + } + + impl Drop for LateDropProbe { + fn drop(&mut self) { + self.state_was_destroyed.store( + MANAGED_APARTMENTS.try_with(|_| ()).is_err(), + Ordering::SeqCst, + ); + self.creation_was_rejected + .store(RoApartment::new(None).is_err(), Ordering::SeqCst); + drop(self.apartment.take()); + } + } + + thread_local! { + static LATE_APARTMENT_DROP: RefCell> = const { RefCell::new(None) }; + } #[test] fn reentrant_apartment_guard_preserves_nested_and_manual_initializations() { Python::initialize(); std::thread::spawn(|| { Python::attach(|py| { - let mut apartment = RoApartment::new(Some(RO_INIT_SINGLETHREADED.0)); + let mut apartment = RoApartment::new(Some(RO_INIT_SINGLETHREADED.0)).unwrap(); apartment.initialize().unwrap(); let outer_callback = NativeCallbackGuard::enter(); let inner_callback = NativeCallbackGuard::enter(); @@ -2771,7 +3018,7 @@ mod tests { assert!(error.to_string().contains("retry after the callback")); assert!(apartment.active); - let mut nested = RoApartment::new(Some(RO_INIT_SINGLETHREADED.0)); + let mut nested = RoApartment::new(Some(RO_INIT_SINGLETHREADED.0)).unwrap(); nested.initialize().unwrap(); nested.uninitialize("RoApartment.close()").unwrap(); ro_initialize(Some(RO_INIT_SINGLETHREADED.0)).unwrap(); @@ -2786,7 +3033,7 @@ mod tests { apartment.uninitialize("RoApartment.close()").unwrap(); assert!(!apartment.active); apartment.uninitialize("RoApartment.close()").unwrap(); - let mut other_model = RoApartment::new(Some(RO_INIT_MULTITHREADED.0)); + let mut other_model = RoApartment::new(Some(RO_INIT_MULTITHREADED.0)).unwrap(); other_model.initialize().unwrap(); other_model.uninitialize("RoApartment.close()").unwrap(); }); @@ -2800,7 +3047,7 @@ mod tests { Python::initialize(); std::thread::spawn(|| { Python::attach(|py| { - let mut apartment = RoApartment::new(Some(RO_INIT_SINGLETHREADED.0)); + let mut apartment = RoApartment::new(Some(RO_INIT_SINGLETHREADED.0)).unwrap(); apartment.initialize().unwrap(); let callback = NativeCallbackGuard::enter(); apartment @@ -2840,6 +3087,182 @@ mod tests { .unwrap(); } + #[test] + fn foreign_close_and_drop_preserve_owner_thread_retry() { + let _serial = crate::errors::UNRAISABLE_HOOK_TEST_LOCK.lock().unwrap(); + Python::initialize(); + std::thread::spawn(|| { + let mut apartment = RoApartment::new(Some(RO_INIT_SINGLETHREADED.0)).unwrap(); + apartment.initialize().unwrap(); + let apartment = std::thread::spawn(move || { + assert!(apartment.uninitialize("RoApartment.close()").is_err()); + assert!(apartment.active); + apartment + }) + .join() + .unwrap(); + assert!(apartment.active); + + let callback = NativeCallbackGuard::enter(); + std::thread::spawn(move || drop(apartment)).join().unwrap(); + assert!(RoApartment::recover_pending().is_err()); + drop(callback); + + let mut recovered = RoApartment::recover_pending().unwrap(); + assert!(recovered.active); + recovered.uninitialize("RoApartment.close()").unwrap(); + let mut other_model = RoApartment::new(Some(RO_INIT_MULTITHREADED.0)).unwrap(); + other_model.initialize().unwrap(); + other_model.uninitialize("RoApartment.close()").unwrap(); + }) + .join() + .unwrap(); + } + + #[test] + fn poisoned_foreign_drop_queue_keeps_lease_recoverable() { + let _serial = crate::errors::UNRAISABLE_HOOK_TEST_LOCK.lock().unwrap(); + Python::initialize(); + std::thread::spawn(|| { + let mut apartment = RoApartment::new(Some(RO_INIT_SINGLETHREADED.0)).unwrap(); + apartment.initialize().unwrap(); + let queue = apartment.foreign_drops.clone(); + let poisoned = std::panic::catch_unwind(|| { + let _guard = queue.lock().unwrap(); + panic!("poison the owner-thread recovery queue"); + }); + assert!(poisoned.is_err()); + + std::thread::spawn(move || drop(apartment)).join().unwrap(); + assert!(!queue.is_poisoned()); + let poisoned_again = std::panic::catch_unwind(|| { + let _guard = queue.lock().unwrap(); + panic!("poison the owner-thread queue before recovery"); + }); + assert!(poisoned_again.is_err()); + let mut recovered = RoApartment::recover_pending().unwrap(); + assert!(!queue.is_poisoned()); + recovered.uninitialize("RoApartment.close()").unwrap(); + assert!(RoApartment::recover_pending().is_err()); + }) + .join() + .unwrap(); + } + + #[test] + fn closed_owner_queue_rejects_foreign_drop_without_native_cleanup() { + let _serial = crate::errors::UNRAISABLE_HOOK_TEST_LOCK.lock().unwrap(); + Python::initialize(); + let (sender, receiver) = std::sync::mpsc::channel(); + let owner = std::thread::spawn(move || { + let mut apartment = RoApartment::new(Some(RO_INIT_SINGLETHREADED.0)).unwrap(); + apartment.initialize().unwrap(); + sender.send(apartment).unwrap(); + }); + let apartment = receiver.recv().unwrap(); + owner.join().unwrap(); + + let queue = apartment.foreign_drops.clone(); + assert!(!lock_foreign_drops(&queue).0.owner_alive); + drop(apartment); + let (queue, _) = lock_foreign_drops(&queue); + assert!(queue.pending.is_empty()); + } + + #[test] + fn foreign_drop_does_not_wait_for_owners_python_gil() { + Python::initialize(); + std::thread::spawn(|| { + Python::attach(|_py| { + let mut apartment = RoApartment::new(Some(RO_INIT_SINGLETHREADED.0)).unwrap(); + apartment.initialize().unwrap(); + let callback = NativeCallbackGuard::enter(); + let (sent, received) = std::sync::mpsc::channel(); + let worker = std::thread::spawn(move || { + drop(apartment); + sent.send(()).unwrap(); + }); + assert!( + received.recv_timeout(Duration::from_secs(3)).is_ok(), + "foreign Drop waited for the owner to release Python's GIL" + ); + worker.join().unwrap(); + assert!(RoApartment::recover_pending().is_err()); + drop(callback); + let mut recovered = RoApartment::recover_pending().unwrap(); + recovered.uninitialize("RoApartment.close()").unwrap(); + }); + }) + .join() + .unwrap(); + } + + #[test] + fn owner_queue_closes_with_unrecovered_foreign_drop() { + let (sent, received) = std::sync::mpsc::channel(); + let (continue_owner, owner_waits) = std::sync::mpsc::channel::<()>(); + let owner = std::thread::spawn(move || { + let mut apartment = RoApartment::new(Some(RO_INIT_SINGLETHREADED.0)).unwrap(); + apartment.initialize().unwrap(); + sent.send(apartment).unwrap(); + owner_waits.recv().unwrap(); + }); + let apartment = received.recv().unwrap(); + let queue = apartment.foreign_drops.clone(); + drop(apartment); + assert_eq!(lock_foreign_drops(&queue).0.pending.len(), 1); + continue_owner.send(()).unwrap(); + owner.join().unwrap(); + let (queue, _) = lock_foreign_drops(&queue); + assert!(!queue.owner_alive); + assert_eq!(queue.pending.len(), 1); + } + + #[test] + fn owner_drop_after_apartment_tls_teardown_fails_closed() { + let state_was_destroyed = Arc::new(std::sync::atomic::AtomicBool::new(false)); + let creation_was_rejected = Arc::new(std::sync::atomic::AtomicBool::new(false)); + let observed = state_was_destroyed.clone(); + let rejected = creation_was_rejected.clone(); + std::thread::spawn(move || { + LATE_APARTMENT_DROP.with(|_| ()); + let mut apartment = RoApartment::new(Some(RO_INIT_SINGLETHREADED.0)).unwrap(); + apartment.initialize().unwrap(); + LATE_APARTMENT_DROP.with(|slot| { + *slot.borrow_mut() = Some(LateDropProbe { + apartment: Some(apartment), + state_was_destroyed: observed, + creation_was_rejected: rejected, + }); + }); + }) + .join() + .unwrap(); + assert!(state_was_destroyed.load(Ordering::SeqCst)); + assert!(creation_was_rejected.load(Ordering::SeqCst)); + } + + #[test] + fn manual_uninitialize_requires_matching_owner_thread_initialization() { + Python::initialize(); + std::thread::spawn(|| { + assert!(ro_uninitialize().is_err()); + ro_initialize(Some(RO_INIT_SINGLETHREADED.0)).unwrap(); + assert!( + std::thread::spawn(|| ro_uninitialize().is_err()) + .join() + .unwrap() + ); + ro_uninitialize().unwrap(); + + let mut apartment = RoApartment::new(Some(RO_INIT_MULTITHREADED.0)).unwrap(); + apartment.initialize().unwrap(); + apartment.uninitialize("RoApartment.close()").unwrap(); + }) + .join() + .unwrap(); + } + #[derive(Default)] struct QueryCounts { queries: AtomicU32, diff --git a/tests/e2e/runners/py_apartment_reentrancy.py b/tests/e2e/runners/py_apartment_reentrancy.py index d55dbb81..bb8f72f6 100644 --- a/tests/e2e/runners/py_apartment_reentrancy.py +++ b/tests/e2e/runners/py_apartment_reentrancy.py @@ -33,6 +33,22 @@ def dispose(mapping, token): mapping.off_map_changed(token) dw.release_projected(mapping) + def foreign_error(action): + errors = [] + + def on_foreign_thread(): + try: + action() + except BaseException as error: + errors.append(error) + + thread = threading.Thread(target=on_foreign_thread) + thread.start() + thread.join(timeout=5) + assert not thread.is_alive() and len(errors) == 1, errors + assert type(errors[0]) is RuntimeError, errors[0] + return str(errors[0]) + if scenario in ('close', 'exit'): apartment = dw.RoApartment(mode) apartment.__enter__() @@ -169,9 +185,6 @@ def recover_on_wrong_thread(): outer.__enter__() inner = dw.RoApartment(mode) inner.__enter__() - unraisable = [] - original_hook = sys.unraisablehook - sys.unraisablehook = unraisable.append def on_changed(sender, args): outer.close() @@ -182,13 +195,7 @@ def on_changed(sender, args): del inner mapping = PropertySet() token = mapping.on_map_changed(on_changed) - try: - mapping.insert('dropped', None) - finally: - sys.unraisablehook = original_hook - assert len(unraisable) == 1, unraisable - assert 'recover_pending()' in str(unraisable[0].exc_value) - unraisable.clear() + mapping.insert('dropped', None) recovered = dw.RoApartment.recover_pending() dispose(mapping, token) recovered.close() @@ -224,6 +231,81 @@ def on_changed(_sender, _args): recovered.close() switch_models() + elif scenario == 'foreign_enter': + apartment = dw.RoApartment(mode) + message = foreign_error(apartment.__enter__) + assert 'OS thread where RoApartment was created' in message, message + assert 'active=false' in repr(apartment) + apartment.__enter__() + mapping = PropertySet() + dw.release_projected(mapping) + apartment.close() + switch_models() + + elif scenario in ('foreign_close', 'foreign_exit', 'foreign_manual'): + apartment = dw.RoApartment(mode) + apartment.__enter__() + events = [] + mapping = PropertySet() + token = mapping.on_map_changed(lambda _sender, args: events.append(args.key)) + if scenario == 'foreign_close': + message = foreign_error(apartment.close) + elif scenario == 'foreign_exit': + message = foreign_error(lambda: apartment.__exit__(None, None, None)) + else: + message = foreign_error(dw.ro_uninitialize) + try: + dw.ro_uninitialize() + except RuntimeError as error: + assert 'no matching ro_initialize()' in str(error), error + else: + raise AssertionError('manual uninitialize consumed a RoApartment context') + if scenario == 'foreign_manual': + assert 'no matching ro_initialize()' in message, message + else: + assert 'OS thread where RoApartment was created' in message, message + assert 'active=true' in repr(apartment) + mapping.insert('owner-still-active', None) + assert events == ['owner-still-active'], events + dispose(mapping, token) + apartment.close() + switch_models() + + elif scenario == 'foreign_drop': + apartment = dw.RoApartment(mode) + apartment.__enter__() + events = [] + + def on_changed(_sender, args): + events.append(args.key) + try: + dw.RoApartment.recover_pending() + except RuntimeError as error: + assert 'native callback' in str(error), error + else: + raise AssertionError('foreign-dropped apartment recovered in a native callback') + + mapping = PropertySet() + token = mapping.on_map_changed(on_changed) + holder = [apartment] + del apartment + + def drop_on_foreign_thread(): + value = holder.pop() + del value + + thread = threading.Thread(target=drop_on_foreign_thread) + thread.start() + thread.join(timeout=5) + assert not thread.is_alive() and not holder, holder + mapping.insert('owner-still-active', None) + assert events == ['owner-still-active'], events + dispose(mapping, token) + recovered = dw.RoApartment.recover_pending() + assert 'active=true' in repr(recovered) + recovered.close() + switch_models() + elif scenario == 'exception': apartment = dw.RoApartment(mode) apartment.__enter__() @@ -271,7 +353,8 @@ def on_changed(_sender, _args): '--scenario', choices=[ 'close', 'exit', 'nested', 'manual', 'pending', 'drop', - 'pending_error', 'exception', + 'pending_error', 'exception', 'foreign_enter', 'foreign_close', + 'foreign_exit', 'foreign_manual', 'foreign_drop', ], required=True, ) diff --git a/tests/e2e/runners/py_runner.py b/tests/e2e/runners/py_runner.py index 3d24fd40..dc124662 100644 --- a/tests/e2e/runners/py_runner.py +++ b/tests/e2e/runners/py_runner.py @@ -1423,6 +1423,8 @@ def observe_null(sender, args): for scenario in ( 'close', 'exit', 'nested', 'manual', 'pending', 'drop', 'pending_error', 'exception', + 'foreign_enter', 'foreign_close', 'foreign_exit', + 'foreign_manual', 'foreign_drop', ): try: child = subprocess.run( @@ -1451,6 +1453,26 @@ def observe_null(sender, args): f'{child.stdout} {child.stderr}' ) return cr + if ( + scenario == 'foreign_drop' + and 'RoApartment was dropped on a different OS thread' + not in child.stderr + ): + cr['error'] = ( + f'{scenario} apartment={mode} did not report the ' + f'foreign drop: {child.stderr!r}' + ) + return cr + if ( + scenario == 'drop' + and 'RoApartment was dropped during a synchronous native callback' + not in child.stderr + ): + cr['error'] = ( + f'{scenario} apartment={mode} did not report the ' + f'pending drop: {child.stderr!r}' + ) + return cr cr['pass'] = True elif kind == 'work_item_callback_passthrough':