Skip to main content

tenferro_cpu/backend/
execution_scope.rs

1//! Callback-lifetime admission for sequential high-level CPU operations.
2
3use std::cell::RefCell;
4use std::marker::PhantomData;
5use std::rc::Rc;
6use std::sync::Arc;
7
8use super::{CpuBackend, CpuRuntimeIdentity, CPU_BACKEND};
9use crate::arbiter::{fresh_execution_owner, has_active_execution, ResourcePermit};
10use crate::engine::CpuEngine;
11use crate::provider::CpuOperationEntry;
12use crate::resource_domain::CpuResourceDomain;
13use crate::CpuDomainOwnership;
14use tenferro_tensor::SessionEntryError;
15
16struct Scope {
17    identity: CpuRuntimeIdentity,
18    engine: Arc<CpuEngine>,
19    permit: Arc<ResourcePermit>,
20    operation_active: bool,
21}
22
23thread_local! {
24    static SCOPE: RefCell<Option<Scope>> = const { RefCell::new(None) };
25}
26
27/// What CPU execution the current thread is inside.
28///
29/// An owner that serializes callers with a blocking lock (such as an eager
30/// runtime's backend owner) must not wait on that lock while this thread holds
31/// a CPU execution permit: another thread may hold the owner and wait for the
32/// permit (#1946 F1).
33///
34/// # Examples
35///
36/// ```
37/// use tenferro_cpu::{current_cpu_execution, CpuBackend, CpuThreadExecution};
38///
39/// assert_eq!(current_cpu_execution(), CpuThreadExecution::Idle);
40/// let backend = CpuBackend::with_threads(1)?;
41/// assert_eq!(backend.install(current_cpu_execution)?, CpuThreadExecution::Active);
42/// let in_scope = backend.with_execution_scope(current_cpu_execution)?;
43/// assert_eq!(in_scope, CpuThreadExecution::SharedScope);
44/// # Ok::<(), tenferro_tensor::Error>(())
45/// ```
46#[derive(Clone, Copy, Debug, PartialEq, Eq)]
47pub enum CpuThreadExecution {
48    /// No CPU execution holds a permit on this thread.
49    Idle,
50    /// A shared execution scope holds a permit, and none of its operations is
51    /// running: an operation or session may still be admitted under it.
52    SharedScope,
53    /// A CPU operation or session is running; nested entry is rejected.
54    Active,
55}
56
57/// Report what CPU execution the current thread is inside, including a
58/// managed Rayon worker running a session or scope callback.
59///
60/// # Examples
61///
62/// ```
63/// use tenferro_cpu::{current_cpu_execution, CpuThreadExecution};
64///
65/// assert_eq!(current_cpu_execution(), CpuThreadExecution::Idle);
66/// ```
67pub fn current_cpu_execution() -> CpuThreadExecution {
68    let idle_scope = SCOPE.with(|slot| {
69        slot.borrow()
70            .as_ref()
71            .is_some_and(|scope| !scope.operation_active)
72    });
73    if idle_scope {
74        CpuThreadExecution::SharedScope
75    } else if has_active_execution() {
76        CpuThreadExecution::Active
77    } else {
78        CpuThreadExecution::Idle
79    }
80}
81
82struct ScopeGuard;
83
84impl Drop for ScopeGuard {
85    fn drop(&mut self) {
86        SCOPE.with(|slot| slot.borrow_mut().take());
87    }
88}
89
90// The operation loan belongs to this callback thread, never a child worker.
91pub(super) struct OperationGuard(PhantomData<Rc<()>>);
92
93impl Drop for OperationGuard {
94    fn drop(&mut self) {
95        SCOPE.with(|slot| {
96            if let Some(scope) = slot.borrow_mut().as_mut() {
97                scope.operation_active = false;
98            }
99        });
100    }
101}
102
103pub(super) enum ExecutionAdmission {
104    Standalone(ResourcePermit),
105    Shared(Arc<ResourcePermit>, OperationGuard),
106}
107
108impl ExecutionAdmission {
109    pub(super) fn permit(&self) -> &ResourcePermit {
110        match self {
111            Self::Standalone(permit) => permit,
112            Self::Shared(permit, _) => permit,
113        }
114    }
115}
116
117pub(crate) fn is_entered(domain: &CpuResourceDomain, permit: &ResourcePermit) -> bool {
118    SCOPE.with(|slot| {
119        slot.borrow().as_ref().is_some_and(|scope| {
120            scope.operation_active
121                && std::ptr::eq(scope.engine.domain(), domain)
122                && std::ptr::eq(scope.permit.as_ref(), permit)
123        })
124    })
125}
126
127impl CpuBackend {
128    /// Run sequential high-level CPU work in one entered execution scope.
129    ///
130    /// Clones of this immutable backend witness may execute ordinary tensor,
131    /// eager/AD and prepared trace operations in the callback without installing
132    /// the executor again. Construct the eager/traced runtime from such a clone.
133    /// Each operation still owns its usual exclusive buffer/cache borrow. The
134    /// scope holds the resource permit, including BLAS provider exclusion, until
135    /// return or unwind. Only Tenferro-managed CPU executors are supported.
136    ///
137    /// Enter the scope and prepare inputs before starting a steady-state timer.
138    /// This does not remove intrinsic output allocation or operation dispatch.
139    ///
140    /// # Examples
141    ///
142    /// ```
143    /// use tenferro_cpu::CpuBackend;
144    /// use tenferro_tensor::{BackendSessionHost, Tensor, TensorRead};
145    ///
146    /// let owner = CpuBackend::with_threads(1)?;
147    /// let mut operations = owner.clone();
148    /// let x = Tensor::from_vec_col_major(vec![2], vec![1.0_f64, 3.0])?;
149    /// let y = owner.with_execution_scope(|| -> Result<_, tenferro_tensor::Error> {
150    ///     let y = operations.with_backend_session(|session| {
151    ///         session.add_read(TensorRead::from_tensor(&x), TensorRead::from_tensor(&x))
152    ///     })??;
153    ///     Ok(operations.with_backend_session(|session| {
154    ///         session.add_read(TensorRead::from_tensor(&y), TensorRead::from_tensor(&x))
155    ///     })??)
156    /// })??;
157    /// assert_eq!(y.as_slice::<f64>()?, &[3.0, 9.0]);
158    /// # Ok::<(), Box<dyn std::error::Error>>(())
159    /// ```
160    ///
161    /// # Errors
162    ///
163    /// Returns [`crate::Error::RuntimeState`] if a scope or CPU execution is
164    /// already active, or [`crate::Error::Unsupported`] for an externally managed
165    /// executor. Poisoned admission state is reported as
166    /// [`crate::Error::SessionEntry`]. Executor admission errors retain their
167    /// typed source in [`crate::Error::BackendSource`]. A callback's return value,
168    /// including its own error result, is returned unchanged inside this
169    /// method's result.
170    ///
171    /// Sessions opened inside the callback with a different backend witness fail
172    /// with [`tenferro_tensor::SessionEntryError::IncompatibleContext`], and a
173    /// session opened from inside another session fails with
174    /// [`tenferro_tensor::SessionEntryError::Reentered`]; neither runs its
175    /// callback.
176    ///
177    /// # Panics
178    ///
179    /// A panic in the callback propagates after releasing the scope and permit.
180    pub fn with_execution_scope<R: Send>(
181        &self,
182        operation: impl FnOnce() -> R + Send,
183    ) -> crate::Result<R> {
184        const OP: &str = "CpuBackend::with_execution_scope";
185        if has_active_execution() || SCOPE.with(|slot| slot.borrow().is_some()) {
186            return Err(crate::Error::runtime_state(
187                OP,
188                "CPU execution is already active; open the shared scope outside active scopes and backend sessions",
189            ));
190        }
191        if self.engine.domain().ownership() != CpuDomainOwnership::Managed {
192            return Err(crate::Error::unsupported(
193                OP,
194                "shared execution scopes require a Tenferro-managed CPU domain; use ordinary operation entry for external domains",
195            ));
196        }
197        let owner = fresh_execution_owner().ok_or(SessionEntryError::Reentered {
198            backend: CPU_BACKEND,
199        })?;
200        let permit = Arc::new(self.acquire_execution_permit(owner)?);
201        let entry = CpuOperationEntry::new(self.engine.domain(), &permit)
202            .with_batch_policy(self.batch_policy);
203        entry
204            .enter(entry.preferred_engine_mode(), |_| {
205                SCOPE.with(|slot| {
206                    *slot.borrow_mut() = Some(Scope {
207                        identity: self.runtime_identity.clone(),
208                        engine: Arc::clone(&self.engine),
209                        permit: Arc::clone(&permit),
210                        operation_active: false,
211                    });
212                });
213                let _guard = ScopeGuard;
214                operation()
215            })
216            .map_err(|error| crate::Error::backend_source(OP, error))
217    }
218
219    /// Admit one CPU operation or session, before any user callback runs.
220    ///
221    /// Inside an active execution scope this reuses the scope's permit; outside
222    /// one it acquires a fresh permit, waiting in FIFO order behind other
223    /// threads that hold overlapping CPU resources.
224    pub(super) fn execution_admission(&self) -> Result<ExecutionAdmission, SessionEntryError> {
225        let shared = SCOPE.with(|slot| {
226            let mut slot = slot.borrow_mut();
227            let Some(scope) = slot.as_mut() else {
228                return Ok(None);
229            };
230            if scope.operation_active {
231                // An operation of this scope is running; nested entry falls
232                // through to the reentry check below.
233                return Ok(None);
234            }
235            if scope.identity != self.runtime_identity {
236                return Err(SessionEntryError::IncompatibleContext {
237                    backend: CPU_BACKEND,
238                    message: "operation backend does not match the active execution scope; \
239                              use a clone of the scope's backend witness"
240                        .to_owned(),
241                });
242            }
243            scope.operation_active = true;
244            Ok(Some(ExecutionAdmission::Shared(
245                Arc::clone(&scope.permit),
246                OperationGuard(PhantomData),
247            )))
248        })?;
249        if let Some(shared) = shared {
250            return Ok(shared);
251        }
252        let owner = fresh_execution_owner().ok_or(SessionEntryError::Reentered {
253            backend: CPU_BACKEND,
254        })?;
255        Ok(ExecutionAdmission::Standalone(
256            self.acquire_execution_permit(owner)?,
257        ))
258    }
259}