Skip to main content

nmbrs_runtime/
execution_context.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! SRD-88 — the task-local **`ExecutionContext`**.
5//!
6//! The load-bearing seam for running multiple executions concurrently in one
7//! process, all sharing one session. A deeply-nested fiber finds *its*
8//! execution's isolated state — `exec_id`, stop flag, and (as the
9//! de-globalization proceeds in later pushes) observer / output channel /
10//! scene tree — by reading this task-local context instead of a
11//! process-global static.
12//!
13//! **Axiom A1 — additive de-globalization.** The process-globals
14//! (`SESSION_STOP`, `GLOBAL_OBSERVER`, `CHANNEL`, …) remain the **process
15//! default**. Code that does not run inside [`scope`] (bootstrap, the CLI,
16//! single-run `nmbrs run`, tests) reads the default — behavior is identical to
17//! before this seam existed. Only concurrent executions scope a context and
18//! get isolation, so the migration is safe and incremental: each accessor that
19//! learns to consult the context is a no-op until someone scopes one.
20//!
21//! **Axiom A2 — `exec_id` is the only new global.** Allocating `exec_id` needs
22//! one process-global counter so two concurrent *fresh* executions can't both
23//! claim the same id. [`alloc_exec_id`] is that single synchronization point.
24
25use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
26use std::sync::{Arc, RwLock};
27
28use arc_swap::ArcSwapOption;
29
30use crate::scene_tree::SceneTree;
31
32/// Per-execution context (SRD-88 §3). Grows as the de-globalization proceeds;
33/// Push 1 carries the `exec_id`, the per-execution stop flag, and the
34/// per-execution observer (display/lifecycle routing).
35pub struct ExecutionContext {
36    /// This execution's id — the partition key for its metrics, log lines,
37    /// and checkpoint events in the shared session store (SRD-77).
38    pub exec_id: u64,
39    /// This execution's stop flag. A stop scoped here halts *this* execution
40    /// only; a global Ctrl-C (`session_signals::SESSION_STOP`) still halts
41    /// every execution (see [`crate::session_signals::stop_requested`]).
42    pub stop: Arc<AtomicBool>,
43    /// This execution's observer (lifecycle + log routing). `None` falls back
44    /// to the process-global `GLOBAL_OBSERVER` (axiom A1) — so a context with
45    /// no observer set behaves exactly as the single-run path.
46    pub observer: Option<Arc<dyn crate::observer::RunObserver>>,
47    /// This execution's scene tree (its own scenario structure + lifecycle
48    /// status). **Late-bound:** the pre-map walker installs it *during* the
49    /// run, so it is interior-mutable. Empty until installed; reads fall back
50    /// to the process-global `GLOBAL_TREE` (A1). The session is the common
51    /// root (SRD-88) — each execution's tree derives from the shared session
52    /// scope; this slot holds the per-execution derived structure.
53    pub scene_tree: ArcSwapOption<RwLock<SceneTree>>,
54    /// This execution's output channel (SRD-87 op-output / log / status /
55    /// raster buckets). `None` falls back to the process-global `CHANNEL`
56    /// (axiom A1). A concurrent in-process execution scopes its own (e.g. a
57    /// `CaptureChannel`) so its op stdout is captured per-execution instead of
58    /// colliding on the one process fd — what in-process example verification
59    /// needs to check `expect` regexes against each execution's output.
60    pub channel: Option<Arc<dyn crate::output_channel::OutputChannel>>,
61}
62
63impl ExecutionContext {
64    /// A fresh context with a freshly-allocated `exec_id` and its own stop
65    /// flag, no per-execution observer (falls back to the global — A1).
66    pub fn new() -> Arc<Self> {
67        Arc::new(Self {
68            exec_id: alloc_exec_id(),
69            stop: Arc::new(AtomicBool::new(false)),
70            observer: None,
71            scene_tree: ArcSwapOption::const_empty(),
72            channel: None,
73        })
74    }
75
76    /// A fresh context that routes lifecycle/log through its own `observer`
77    /// instead of the process-global. Used by concurrent executions so each
78    /// folds its own outcome independently.
79    pub fn with_observer(observer: Arc<dyn crate::observer::RunObserver>) -> Arc<Self> {
80        Arc::new(Self {
81            exec_id: alloc_exec_id(),
82            stop: Arc::new(AtomicBool::new(false)),
83            observer: Some(observer),
84            scene_tree: ArcSwapOption::const_empty(),
85            channel: None,
86        })
87    }
88
89    /// A fresh context routing both lifecycle/log (`observer`) AND the SRD-87
90    /// output buckets (`channel`) through its own per-execution sinks. Used by
91    /// in-process example verification so each concurrent execution's op stdout
92    /// is captured separately.
93    pub fn with_observer_and_channel(
94        observer: Arc<dyn crate::observer::RunObserver>,
95        channel: Arc<dyn crate::output_channel::OutputChannel>,
96    ) -> Arc<Self> {
97        Arc::new(Self {
98            exec_id: alloc_exec_id(),
99            stop: Arc::new(AtomicBool::new(false)),
100            observer: Some(observer),
101            scene_tree: ArcSwapOption::const_empty(),
102            channel: Some(channel),
103        })
104    }
105}
106
107tokio::task_local! {
108    static EXEC_CTX: Arc<ExecutionContext>;
109
110    /// SRD-100 P1c — the scene node of the phase the **current task** is
111    /// executing. Set by `run_phase` (via [`with_current_phase`]) for the
112    /// phase body's duration and carried across fiber spawns by
113    /// [`propagate`], so an *ambient* emit — `running_phase_indent`, the
114    /// errorhandler / metrics-diag log bridges, poll / validation progress —
115    /// nests under ITS OWN phase's depth instead of a global
116    /// first-running-by-DFS guess.
117    ///
118    /// Deliberately a **separate task-local**, not a field on the shared
119    /// per-execution [`ExecutionContext`]: concurrent sibling phases run in
120    /// ONE execution (so they share one `ExecutionContext` `Arc`) but each
121    /// must resolve to ITS node — a shared field would be stomped. Task-locals
122    /// are per-task, so each phase body (and its propagated fibers) sees its
123    /// own. `None` outside any phase body (e.g. the metrics scheduler thread);
124    /// readers fall back to the first-running DFS heuristic (A1).
125    static CURRENT_PHASE: crate::scene_tree::SceneNodeId;
126
127    /// Epoch millis at which the phase the **current task** is executing
128    /// started. Scoped by `run_phase` next to [`CURRENT_PHASE`] and carried
129    /// across fiber spawns by [`propagate`], so a phase-scoped clock is
130    /// available to everything running under the phase — including the ops,
131    /// which is the whole point.
132    ///
133    /// Same reason as `CURRENT_PHASE` for being a task-local rather than a
134    /// field on the shared per-execution context: concurrent sibling phases
135    /// share one `ExecutionContext` but must each resolve to THEIR start.
136    static CURRENT_PHASE_START_MS: u64;
137}
138
139/// The `exec_id` allocator — SRD-88 axiom A2, the one unavoidable global.
140/// Monotonic per process; the first allocated id is `1`, matching the legacy
141/// single-execution session's `exec_id`.
142static NEXT_EXEC_ID: AtomicU64 = AtomicU64::new(1);
143
144/// Allocate the next process-unique `exec_id`.
145pub fn alloc_exec_id() -> u64 {
146    NEXT_EXEC_ID.fetch_add(1, Ordering::Relaxed)
147}
148
149/// The current task's [`ExecutionContext`], if running inside a [`scope`].
150pub fn try_current() -> Option<Arc<ExecutionContext>> {
151    EXEC_CTX.try_with(|c| c.clone()).ok()
152}
153
154/// The current execution's `exec_id`, or `1` (the legacy single-execution
155/// default) outside any execution scope.
156pub fn current_exec_id() -> u64 {
157    EXEC_CTX.try_with(|c| c.exec_id).unwrap_or(1)
158}
159
160/// The current execution's stop flag, if scoped (`None` outside a scope, so
161/// the caller falls back to the process-global stop — A1).
162pub fn current_stop() -> Option<Arc<AtomicBool>> {
163    EXEC_CTX.try_with(|c| c.stop.clone()).ok()
164}
165
166/// The current execution's observer, if a scoped context set one. `None`
167/// outside a scope OR when the scoped context has no observer — the caller
168/// falls back to the process-global `GLOBAL_OBSERVER` (A1).
169pub fn current_observer() -> Option<Arc<dyn crate::observer::RunObserver>> {
170    EXEC_CTX.try_with(|c| c.observer.clone()).ok().flatten()
171}
172
173/// The current execution's output channel, if a scoped context set one.
174/// `None` outside a scope OR when the scoped context has no channel — the
175/// caller falls back to the process-global `CHANNEL` (A1).
176pub fn current_channel() -> Option<Arc<dyn crate::output_channel::OutputChannel>> {
177    EXEC_CTX.try_with(|c| c.channel.clone()).ok().flatten()
178}
179
180/// The current execution's scene tree, if scoped AND installed (`None`
181/// otherwise — the caller falls back to the process-global `GLOBAL_TREE`, A1).
182pub fn current_scene_tree() -> Option<Arc<RwLock<SceneTree>>> {
183    EXEC_CTX
184        .try_with(|c| c.scene_tree.load_full())
185        .ok()
186        .flatten()
187}
188
189/// Install `tree` into the current execution's context, returning `true` if a
190/// context was scoped (so the caller installs into the global only when there
191/// is no execution scope). The pre-map walker calls this so each execution's
192/// lifecycle mutations land on its own tree.
193pub fn install_scene_tree(tree: Arc<RwLock<SceneTree>>) -> bool {
194    EXEC_CTX
195        .try_with(|c| c.scene_tree.store(Some(tree)))
196        .is_ok()
197}
198
199/// Run `fut` with `ctx` as the task-local [`ExecutionContext`] for the whole
200/// future. Concurrent executions each scope their own context, so their
201/// task-local reads (stop flag, `exec_id`, …) resolve independently.
202pub async fn scope<F: std::future::Future>(ctx: Arc<ExecutionContext>, fut: F) -> F::Output {
203    EXEC_CTX.scope(ctx, fut).await
204}
205
206/// Run `fut` as the body of the phase at `scene_node_id` (SRD-100 P1c): scopes
207/// the task-local [`CURRENT_PHASE`] for the future's duration. `run_phase`
208/// wraps itself with this so its whole body — the activity loop and every
209/// fiber it [`propagate`]s — resolves to this phase's node, and any ambient
210/// emit nests under its depth.
211pub async fn with_current_phase<F: std::future::Future>(
212    scene_node_id: crate::scene_tree::SceneNodeId,
213    phase_start_ms: u64,
214    fut: F,
215) -> F::Output {
216    CURRENT_PHASE_START_MS
217        .scope(phase_start_ms, CURRENT_PHASE.scope(scene_node_id, fut))
218        .await
219}
220
221/// Epoch millis at which the phase the current task is executing started, if
222/// set. `None` outside any phase body (the metrics scheduler thread, CLI paths,
223/// unit tests that never scoped one) — a caller then has no phase to be
224/// relative to and must say so rather than invent an origin.
225pub fn current_phase_start_ms() -> Option<u64> {
226    CURRENT_PHASE_START_MS.try_with(|ms| *ms).ok()
227}
228
229/// The scene node of the phase the current task is executing, if set (inside a
230/// `run_phase` body or a fiber it propagated). `None` on tasks/threads outside
231/// any phase (e.g. the metrics scheduler) — readers fall back to the
232/// first-running DFS heuristic (A1). Distinct from [`current_scene_tree`]: that
233/// is the execution's whole tree; this is the one node the task is *in*.
234pub fn current_phase_node() -> Option<crate::scene_tree::SceneNodeId> {
235    CURRENT_PHASE.try_with(|id| *id).ok()
236}
237
238/// Wrap `fut` so it runs under the CURRENT execution context — captured **now**
239/// (at the call site, which is still inside the parent's scope) and
240/// re-established as the task-local inside a freshly-spawned task. A no-op
241/// pass-through when there is no current context (single-run / CLI / tests —
242/// axiom A1).
243///
244/// `tokio::spawn` / `JoinSet::spawn` start a NEW task that does **not** inherit
245/// the parent's task-locals, so the per-execution context (observer, scene
246/// tree, stop flag, `exec_id`) would be lost in the spawned per-cycle fibers.
247/// Wrapping each spawned future with `propagate` carries the context across the
248/// spawn boundary, so a fiber deep inside a concurrent execution still resolves
249/// to *its* execution's state.
250pub fn propagate<F>(fut: F) -> impl std::future::Future<Output = F::Output> + Send
251where
252    F: std::future::Future + Send + 'static,
253    F::Output: Send,
254{
255    // Capture BOTH task-locals now, while still inside the parent's scope.
256    let ctx = try_current();
257    let phase = current_phase_node();
258    let phase_start = current_phase_start_ms();
259    async move {
260        // Re-establish the executing-phase (inner) inside the execution
261        // context (outer), so a fiber deep inside a phase still resolves to
262        // both its execution AND its phase node (SRD-100 P1c).
263        let inner = async move {
264            let with_phase = async move {
265                match phase {
266                    Some(p) => CURRENT_PHASE.scope(p, fut).await,
267                    None => fut.await,
268                }
269            };
270            match phase_start {
271                Some(ms) => CURRENT_PHASE_START_MS.scope(ms, with_phase).await,
272                None => with_phase.await,
273            }
274        };
275        match ctx {
276            Some(c) => scope(c, inner).await,
277            None => inner.await,
278        }
279    }
280}
281
282#[cfg(test)]
283mod tests {
284    use super::*;
285
286    /// The phase origin is per-task, so concurrent sibling phases each read
287    /// THEIR own start rather than whichever ran last.
288    #[tokio::test]
289    async fn phase_start_is_scoped_and_absent_outside_a_phase() {
290        assert_eq!(
291            current_phase_start_ms(),
292            None,
293            "outside a phase there is no origin to be relative to"
294        );
295        let inside = with_current_phase(1, 111, async { current_phase_start_ms() }).await;
296        assert_eq!(inside, Some(111));
297        let sibling = with_current_phase(2, 222, async { current_phase_start_ms() }).await;
298        assert_eq!(sibling, Some(222));
299        assert_eq!(
300            current_phase_start_ms(),
301            None,
302            "the scope must not leak past the phase body"
303        );
304    }
305
306    /// Fibers are spawned tasks, which do NOT inherit task-locals — the phase
307    /// clock has to survive `propagate` or every op would read no origin.
308    #[tokio::test]
309    async fn phase_start_survives_propagate_into_a_spawned_task() {
310        let got = with_current_phase(3, 333, async {
311            tokio::spawn(propagate(async { current_phase_start_ms() }))
312                .await
313                .expect("spawned task")
314        })
315        .await;
316        assert_eq!(
317            got,
318            Some(333),
319            "a propagated fiber must resolve its own phase's origin"
320        );
321    }
322
323    #[test]
324    fn alloc_exec_id_is_monotonic_and_unique() {
325        let a = alloc_exec_id();
326        let b = alloc_exec_id();
327        assert!(b > a, "exec_id must be monotonic: {a} then {b}");
328    }
329
330    #[tokio::test]
331    async fn outside_a_scope_defaults_to_legacy_single_execution() {
332        // A1: no context scoped → legacy defaults, no isolation surface.
333        assert_eq!(current_exec_id(), 1);
334        assert!(current_stop().is_none());
335        assert!(try_current().is_none());
336    }
337
338    #[tokio::test]
339    async fn propagate_carries_context_across_spawn() {
340        // A `tokio::spawn`ed task does NOT inherit the parent's task-locals, so
341        // a bare spawn would see the default exec_id (1). `propagate` captures
342        // the current context and re-establishes it inside the spawned task.
343        let ctx = ExecutionContext::new();
344        let id = ctx.exec_id;
345        let (bare, wrapped) = scope(ctx, async move {
346            let bare = tokio::spawn(async { current_exec_id() }).await.unwrap();
347            let wrapped = tokio::spawn(propagate(async { current_exec_id() }))
348                .await
349                .unwrap();
350            (bare, wrapped)
351        })
352        .await;
353        assert_eq!(bare, 1, "a bare spawn loses the context (sees the default)");
354        assert_eq!(
355            wrapped, id,
356            "propagate carries the exec_id across the spawn"
357        );
358    }
359
360    #[tokio::test]
361    async fn current_phase_node_defaults_to_none_and_scopes() {
362        // Outside any phase body the ambient phase is unset (so
363        // `running_phase_indent` falls back to the DFS heuristic).
364        assert_eq!(current_phase_node(), None);
365        // Inside `with_current_phase` it resolves to the scoped node, and
366        // reverts after the body returns.
367        let inside = with_current_phase(7, 0, async { current_phase_node() }).await;
368        assert_eq!(inside, Some(7));
369        assert_eq!(current_phase_node(), None, "the scope reverts on exit");
370    }
371
372    #[tokio::test]
373    async fn propagate_carries_current_phase_across_spawn() {
374        // SRD-100 P1c — a bare spawn inside a phase body loses the ambient
375        // phase; `propagate` re-establishes it so a fiber deep inside the
376        // activity still indents to ITS phase's node.
377        let (bare, wrapped) = with_current_phase(42, 0, async {
378            let bare = tokio::spawn(async { current_phase_node() }).await.unwrap();
379            let wrapped = tokio::spawn(propagate(async { current_phase_node() }))
380                .await
381                .unwrap();
382            (bare, wrapped)
383        })
384        .await;
385        assert_eq!(bare, None, "a bare spawn loses the ambient phase");
386        assert_eq!(
387            wrapped,
388            Some(42),
389            "propagate carries the phase node across the spawn"
390        );
391    }
392
393    #[tokio::test]
394    async fn propagate_carries_both_exec_ctx_and_phase() {
395        // The two task-locals compose: a propagated fiber sees BOTH its
396        // execution's exec_id AND its phase node.
397        let ctx = ExecutionContext::new();
398        let id = ctx.exec_id;
399        let (eid, phase) = scope(
400            ctx,
401            with_current_phase(9, 0, async {
402                tokio::spawn(propagate(async {
403                    (current_exec_id(), current_phase_node())
404                }))
405                .await
406                .unwrap()
407            }),
408        )
409        .await;
410        assert_eq!(eid, id, "propagate carries exec_id");
411        assert_eq!(phase, Some(9), "propagate carries the phase node");
412    }
413
414    // Holds a std lock across `.await`: the awaited `scope(...)`
415    // futures run inline (task-local scope, no inter-task yield), so
416    // there is no other task that needs the lock — it's held only to
417    // keep the process-global stop flag clear for the whole assertion
418    // window, serialized against `flag_starts_unset_and_responds_to_request`.
419    #[allow(clippy::await_holding_lock)]
420    #[tokio::test]
421    async fn per_execution_stop_is_isolated() {
422        // `stop_requested()` ORs the never-reset process-global
423        // `SESSION_STOP`; serialize with the other global-flag test and
424        // clear it so a sibling test's stop can't masquerade as B's.
425        let _guard = crate::session_signals::STOP_GLOBAL_TEST_LOCK
426            .lock()
427            .unwrap_or_else(|e| e.into_inner());
428        crate::session_signals::clear_session_stop_for_test();
429
430        let a = ExecutionContext::new();
431        let b = ExecutionContext::new();
432        assert_ne!(
433            a.exec_id, b.exec_id,
434            "concurrent executions get distinct ids"
435        );
436
437        // Stop A only.
438        a.stop.store(true, Ordering::Relaxed);
439
440        // Inside A's scope the stop is observed; inside B's it is not — a stop
441        // scoped to one execution does NOT halt its concurrent sibling.
442        let a_id = scope(a.clone(), async { current_exec_id() }).await;
443        let a_stopped = scope(a.clone(), async {
444            crate::session_signals::stop_requested()
445        })
446        .await;
447        let b_stopped = scope(b.clone(), async {
448            crate::session_signals::stop_requested()
449        })
450        .await;
451
452        assert_eq!(a_id, a.exec_id, "the scoped exec_id resolves to A's");
453        assert!(a_stopped, "A observes its own stop inside A's scope");
454        assert!(
455            !b_stopped,
456            "B must NOT see A's stop — executions are isolated"
457        );
458    }
459}