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}