Skip to main content

nmbrs_runtime/polydat_nodes/
runtime_context.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Runtime context nodes (SRD 12 §"Runtime context nodes").
5//!
6//! These nodes project stable runtime surfaces — the current
7//! phase, the current cycle ordinal, the value of a dynamic
8//! control, the active rate-limiter target, the active fiber
9//! count — into GK-readable wires. They are the read-side of
10//! the reification principle (SRD 10 §"GK as the unified access
11//! surface"): any value a workload might want to read is reached
12//! through a Polydat binding, not a side channel.
13//!
14//! Like the metric nodes (see `metrics.rs`), these are
15//! non-deterministic context projections — their output changes
16//! between cycles by definition, so the constant-folder is not
17//! allowed to collapse them. They read from globals / thread
18//! locals the runtime sets during bootstrap and on every cycle
19//! tick.
20
21use std::future::Future;
22use std::sync::atomic::{AtomicU64, Ordering};
23use std::sync::{Arc, LazyLock, Mutex, RwLock};
24
25// All nodes here are authored via `#[polydat::polydat_node]` (fully-qualified
26// `polydat::…` paths), so no `polydat::ast` imports are needed in module
27// code (the in-module tests import what they need locally).
28
29use nmbrs_metrics::component::Component;
30
31// =========================================================================
32// Global session-root handle + per-fiber task-local context
33// =========================================================================
34
35/// Global handle to the session's component root. Set by the
36/// runner during scenario bootstrap so context nodes can resolve
37/// `control(...)` reads against the live tree.
38static SESSION_ROOT: LazyLock<Mutex<Option<Arc<RwLock<Component>>>>> =
39    LazyLock::new(|| Mutex::new(None));
40
41/// Install the session root for every runtime-context node that
42/// reads tree state. Call once at scenario bootstrap; subsequent
43/// calls overwrite.
44/// Test-only serialization for the process-global session root.
45///
46/// `SESSION_ROOT` is one global; the parallel test runner is many threads.
47/// Every test that INSTALLS a root — directly, or by constructing a
48/// [`crate::session::Session`], whose constructor installs one as a side
49/// effect — must hold this guard, and so must every test that READS controls
50/// through the root. It lives here, beside the global it protects, precisely
51/// because the offender that motivated it was in another module: a
52/// session-construction test stomping the control tests' root through the
53/// constructor side effect, which a module-private test lock could never
54/// exclude.
55#[cfg(test)]
56pub(crate) fn session_root_test_guard() -> std::sync::MutexGuard<'static, ()> {
57    static TEST_LOCK: Mutex<()> = Mutex::new(());
58    TEST_LOCK.lock().unwrap_or_else(|e| e.into_inner())
59}
60
61pub fn set_session_root(root: Arc<RwLock<Component>>) {
62    *SESSION_ROOT.lock().unwrap_or_else(|e| e.into_inner()) = Some(root);
63}
64
65fn session_root() -> Option<Arc<RwLock<Component>>> {
66    SESSION_ROOT
67        .lock()
68        .unwrap_or_else(|e| e.into_inner())
69        .clone()
70}
71
72/// Public accessor for the runner-installed session root.
73///
74/// Used by integration layers that need to walk the component
75/// tree from outside this crate — the TUI's control-edit
76/// handler, the web API's control endpoints, the
77/// `dryrun=controls` renderer, etc. Returns `None` when no
78/// session is running (pre-bootstrap, after teardown, tests
79/// that never called [`set_session_root`]).
80pub fn session_root_handle() -> Option<Arc<RwLock<Component>>> {
81    session_root()
82}
83
84/// Per-fiber execution context carried across the async call
85/// chain via a [`tokio::task_local!`] binding. Thread-locals are
86/// unsafe here because tokio's work-stealing scheduler can
87/// migrate a task between worker threads at any `.await`; a
88/// task-local is bound to the task itself and survives
89/// migration.
90pub struct FiberContext {
91    /// Name of the phase this fiber is running under. Arc'd so
92    /// two fibers under the same phase share one allocation.
93    pub phase: Arc<str>,
94    /// Current cycle ordinal. `AtomicU64` rather than `Cell` so
95    /// the context is `Sync` — tokio requires task-local values
96    /// to be `Send + Sync`.
97    pub cycle: AtomicU64,
98    /// SRD-89 — the **current component's** control resolution, for this
99    /// fiber's execution: control name → live erased handle, walked up **once**
100    /// from the fiber's own phase component
101    /// ([`nmbrs_metrics::component::Component::control_snapshot`], first-match /
102    /// nearest-wins) and read lock-free thereafter. It holds live handles, not
103    /// values, so a servo retarget on the shared control is observed through it.
104    /// Because each execution's fibers carry **their own** phase component's
105    /// snapshot, a same-named control (`concurrency` / `rate`) resolves to that
106    /// execution's instance — uniformly in single-run and concurrent runs, with
107    /// no session-root walk and no cross-talk. Empty for fibers spawned without
108    /// a component (a control read then falls back to the session-root walk).
109    pub controls: ControlMap,
110}
111
112/// A lock-free, shareable snapshot of resolved control handles (see
113/// [`FiberContext::controls`]).
114pub type ControlMap =
115    Arc<std::collections::HashMap<String, Arc<dyn nmbrs_metrics::controls::ErasedControl>>>;
116
117/// Snapshot the controls visible from `component` (its own controls plus any
118/// `Subtree`-scoped ancestor control, nearest-wins) into a lock-free
119/// [`ControlMap`]. Deadlock-safe: [`Component::control_snapshot`] acquires one
120/// tier's lock at a time, never nested — so it is safe on the cadence-contended
121/// component tree. Computed once per phase, shared across the phase's fibers.
122pub fn snapshot_controls(component: &Arc<RwLock<Component>>) -> ControlMap {
123    Arc::new(nmbrs_metrics::component::Component::control_snapshot(
124        component,
125    ))
126}
127
128/// An empty [`ControlMap`] for fibers / call sites that have no component (the
129/// read then falls back to the session-root walk).
130pub fn empty_controls() -> ControlMap {
131    Arc::new(std::collections::HashMap::new())
132}
133
134tokio::task_local! {
135    /// Task-local fiber context. Set once per fiber at spawn
136    /// time via [`with_fiber_context`]; updated per cycle via
137    /// [`set_task_cycle`]. Reads outside a scope (e.g. unit
138    /// tests, non-fiber code paths) silently return defaults.
139    static FIBER_CTX: FiberContext;
140}
141
142/// Wrap a fiber's async body in a `FiberContext` scope. Every
143/// runtime-context node read performed inside the future sees
144/// the phase name and cycle counter set here.
145///
146/// Cycle starts at 0 and is updated via [`set_task_cycle`] on
147/// every iteration of the fiber's loop.
148pub async fn with_fiber_context<F>(phase: Arc<str>, controls: ControlMap, fut: F) -> F::Output
149where
150    F: Future,
151{
152    FIBER_CTX
153        .scope(
154            FiberContext {
155                phase,
156                cycle: AtomicU64::new(0),
157                controls,
158            },
159            fut,
160        )
161        .await
162}
163
164/// Resolve a control from the running fiber's **current component** snapshot
165/// (lock-free). `None` outside a fiber scope, or when the snapshot doesn't carry
166/// `name` — the caller then falls back to the session-root walk.
167fn current_phase_control(name: &str) -> Option<Arc<dyn nmbrs_metrics::controls::ErasedControl>> {
168    FIBER_CTX
169        .try_with(|ctx| ctx.controls.get(name).cloned())
170        .ok()
171        .flatten()
172}
173
174/// Update the cycle counter in the enclosing [`FiberContext`].
175/// Safe to call outside a scope — the update is a no-op if
176/// there is no active fiber context (e.g. when the node is
177/// evaluated from a unit test that didn't install one).
178pub fn set_task_cycle(cycle: u64) {
179    let _ = FIBER_CTX.try_with(|ctx| ctx.cycle.store(cycle, Ordering::Relaxed));
180}
181
182fn task_phase() -> Option<Arc<str>> {
183    FIBER_CTX.try_with(|ctx| ctx.phase.clone()).ok()
184}
185
186fn task_cycle() -> u64 {
187    FIBER_CTX
188        .try_with(|ctx| ctx.cycle.load(Ordering::Relaxed))
189        .unwrap_or(0)
190}
191
192// =========================================================================
193// control_set(name, value) — GK-driven write into a control
194// =========================================================================
195
196/// Capture the enclosing DSL binding name at node construction, for write
197/// attribution (`ControlOrigin::Polydat { binding }` — surfaces in control
198/// logs so operators attribute a change to a specific binding, not just
199/// "from GK"). The build context names the binding the node is built for;
200/// falls back to the control name when it names none (a node built outside
201/// a binding, e.g. a library test).
202fn capture_binding(ctx: &polydat::dsl::factory::BuildContext, name: &str) -> String {
203    ctx.binding()
204        .map(str::to_string)
205        .unwrap_or_else(|| name.to_string())
206}
207
208/// Polydat write node: submit an f64 write against the named control via the
209/// session root's walk-up. Returns `1` if dispatched, `0` if not — either no
210/// session root is installed (outside a running scenario / pure-kernel test)
211/// or the control already reads exactly the requested value (idempotent
212/// writes are elided; see the fixpoint check in the body).
213///
214/// Signature: `control_set(name: const str, value: f64) -> u64`. Authored
215/// via `#[polydat::polydat_node]` (SRD-80b).
216///
217/// Writes are **non-blocking** — the node spawns a tokio task that calls
218/// [`nmbrs_metrics::controls::ErasedControl::set_f64`] and does not await it
219/// (awaiting would deadlock the single-threaded runtime). The control
220/// layer's confirmed-apply still runs in the background; failures log but do
221/// not stall the issuing fiber. SRD 23 §"Mutation entry points →
222/// GK-driven feedback loops".
223///
224/// Purity is `SideChannel(Other)`: the write is an observable side effect on
225/// runtime control state, so the node is never const-folded or deduped — it
226/// dispatches on every evaluation.
227#[polydat::polydat_node(
228    category = Context,
229    purity = SideChannel(Other),
230)]
231fn control_set(
232    name: Const<&str>,
233    #[poly_const(capture_binding, from = (ctx, name))] binding: &String,
234    value: f64,
235) -> u64 {
236    // SRD-89 — a workload's `control_set(...)` must write to ITS OWN phase-tier
237    // control (`concurrency` / `rate`), not a neighbour's. Resolve through the
238    // SAME path as the read nodes ([`resolve_control`]: the running fiber's
239    // current-component snapshot, else the session-root walk-up) so the write
240    // targets exactly the control a subsequent read sees. Resolution happens
241    // HERE, inside the fiber's `FIBER_CTX` scope — the background write task
242    // below does NOT inherit `FIBER_CTX`, so it cannot re-resolve; we capture
243    // the live handle and hand it over. If nothing resolves (pre-bootstrap /
244    // pure-kernel test) report "not dispatched" (the no-root → 0 contract)
245    // without spawning.
246    let Some(erased) = resolve_control(name.0) else {
247        return 0;
248    };
249    // Idempotent-write elision (SRD-93-era servo support): a feedback
250    // binding that re-evaluates per cycle recomputes the SAME target
251    // between changes of its inputs, and re-dispatching that write
252    // buys nothing while costing a spawn + the confirmed-apply
253    // pipeline every cycle. The committed gauge projection is the
254    // fixpoint check: when the control already reads exactly the
255    // requested value, the write is elided and the node reports 0
256    // (not dispatched — same code as the no-root case). Exact
257    // equality on purpose: between input changes the recomputation is
258    // bit-identical, and any real change, however small, must
259    // dispatch. Controls without a reified gauge (`gauge_f64() ==
260    // None`) never match and keep the always-dispatch behavior. An
261    // in-flight uncommitted write can let one duplicate through —
262    // harmless, it commits the same value.
263    if erased.gauge_f64() == Some(value) {
264        return 0;
265    }
266    let name = name.0.to_string();
267    let binding = binding.clone();
268
269    // Dispatch the write on a background task — the fiber cannot await an
270    // async `set` without blocking the runtime. The handle is already resolved,
271    // so the task needs no execution context.
272    tokio::spawn(async move {
273        let origin = nmbrs_metrics::controls::ControlOrigin::Polydat { binding };
274        if let Err(e) = erased.set_f64(value, origin).await {
275            polydat::audit::warn(&format!("control_set({name}, {value}) failed: {e}"));
276        }
277    });
278
279    1
280}
281
282// =========================================================================
283// control(name) — read a dynamic control's gauge-projection
284// =========================================================================
285
286/// Read the current committed value of a dynamic control as an
287/// `f64`, projected through the control's reified gauge.
288///
289/// Signature: `control(name: String) -> f64`
290///
291/// Resolves by walking up the component tree from the session
292/// root for the first declaration that matches the given name
293/// (honoring branch scope). Returns `0.0` in every case where
294/// no numeric value is available:
295///
296/// - Name doesn't resolve to any declared control.
297/// - Control is declared but was built without
298///   [`nmbrs_metrics::controls::ControlBuilder::reify_as_gauge`],
299///   i.e. no f64 projection was registered.
300/// - Control has a projection but the current value's
301///   `to_f64` returned `None` (e.g. an enum-valued control
302///   sitting on a `None`-projected variant).
303///
304/// This is deliberately silent rather than error-raising: a
305/// workload running at cycle time cannot usefully "handle" a
306/// missing control, and the alternative (panicking or
307/// returning a sentinel like `NaN`) propagates poison
308/// downstream. For typed reads of non-numeric controls use
309/// [`ControlStr`] (`control_str`) / [`ControlBool`]
310/// (`control_bool`), which have explicit defaults for the same
311/// missing-control cases. Operators verify controls exist via
312/// `dryrun=controls` before running.
313
314/// Resolve a control's erased handle for a cycle-time read or write (SRD-89).
315///
316/// **Per-execution first, lock-free:** consults the current execution's
317/// per-phase control map (`execution_context::current_control`) — an `ArcSwap`
318/// load + `HashMap` get, touching no component lock and walking no tree. Under
319/// concurrent in-process executions (SRD-88) this is what makes a same-named
320/// control (`concurrency` / `rate`) resolve to THIS execution's own instance,
321/// so a neighbour's servo can't drive this phase (the SRD-89 §2b cross-talk).
322///
323/// **Fallback — the global `SESSION_ROOT` walk:** when there is no scoped
324/// execution (single-run / dryrun / TUI control-edit), no map yet, or a
325/// session-tier control not in the per-exec map, resolve by walking up from the
326/// installed session root via [`Component::find_control_erased_up`]. This is the
327/// pre-SRD-89 path, so single-run output is byte-identical (axiom A1). The walk
328/// holds a single, non-nested read guard on the session root (which has no
329/// parent), released before the caller projects the value.
330fn resolve_control(name: &str) -> Option<Arc<dyn nmbrs_metrics::controls::ErasedControl>> {
331    // Resolve from the running fiber's **current component** — its phase-tier
332    // control snapshot (walk-up from its own phase component, first match), read
333    // lock-free. Uniform for single-run and concurrent: each execution's fibers
334    // carry their own snapshot, so a same-named control resolves to that
335    // execution's instance. No session-root start.
336    if let Some(handle) = current_phase_control(name) {
337        return Some(handle);
338    }
339    // Fallback: no fiber scope (a control edit from OUTSIDE a running phase —
340    // the TUI `e` prompt, the web control endpoint, `dryrun=controls`). Walk up
341    // from the session root for session-tier controls.
342    let root = session_root()?;
343    let guard = root.read().ok()?;
344    guard.find_control_erased_up(name)
345}
346
347/// Read a dynamic control's current reified-gauge value as `f64`. `0.0` when
348/// the control can't be resolved (no session / per-exec context, absent) or has
349/// no reified gauge. The single read path for every control-reader node
350/// (`control` / `control_u64` / `control_bool` / `rate` / `concurrency`).
351fn control_gauge_f64(name: &str) -> f64 {
352    resolve_control(name)
353        .and_then(|c| c.gauge_f64())
354        .unwrap_or(0.0)
355}
356
357/// Read a dynamic control's current value as its human-readable string
358/// rendering (the erased-control `value_string()`), or `""` when absent.
359fn control_value_string(name: &str) -> String {
360    resolve_control(name)
361        .map(|c| c.value_string())
362        .unwrap_or_default()
363}
364
365/// Read a dynamic control's current reified-gauge value as `f64`, by
366/// walk-up from the live session root. `0.0` when there is no session root,
367/// the control is absent, or it has no reified gauge.
368///
369/// Signature: `control(name: const str) -> f64`. Intrinsically
370/// `Nondeterministic`: the control value changes over the run (operator
371/// edits, SRD-86 servo retargets, per-coordinate reruns), so the node is
372/// never const-folded, re-evaluated on every pull, and the compiler
373/// propagates that volatility to every downstream wire — `load :=
374/// control("concurrency")` (and anything reading `load`) tracks the live
375/// value without the author flagging `volatile`.
376#[polydat::polydat_node(
377    category = Context,
378    purity = Nondeterministic("reads a live dynamic control value; changes over the run"),
379)]
380fn control(name: Const<&str>) -> f64 {
381    control_gauge_f64(name.0)
382}
383
384// =========================================================================
385// control_u64(name) / control_str(name) — typed read sugar
386// =========================================================================
387
388/// Read a dynamic control's current value and cast to `u64`.
389///
390/// Signature: `control_u64(name: const str) -> u64`. Resolves the control by
391/// walk-up from the session root, reads its reified-gauge f64 projection and
392/// casts to u64 (saturating at 0 for negatives). Missing / gauge-less
393/// controls return 0. For integer-valued controls (`concurrency`,
394/// `max_retries`) read as a cycle-time parameter — no need to pipe
395/// `control("…")` through `f64_to_u64`.
396///
397/// Authored via `#[polydat::polydat_node]` (SRD-80b) — purity is the
398/// hygienic attribute below. `Nondeterministic` because the control value
399/// changes over the run (operator edits, SRD-86 servo retargets,
400/// per-coordinate reruns): the node is never const-folded, re-evaluated on
401/// every pull, and the compiler propagates that volatility to every
402/// downstream wire, so `load := control_u64("concurrency")` tracks the live
403/// value without the author flagging `volatile`.
404#[polydat::polydat_node(
405    category = Context,
406    purity = Nondeterministic("reads a live dynamic control value; changes over the run"),
407)]
408fn control_u64(name: Const<&str>) -> u64 {
409    let v = control_gauge_f64(name.0);
410    if v < 0.0 { 0 } else { v as u64 }
411}
412
413/// Read a dynamic control's current value as a boolean — `true` iff its
414/// reified-gauge value is non-zero. Missing / unreified controls → `false`.
415///
416/// Signature: `control_bool(name: const str) -> bool`. Intrinsically
417/// `Nondeterministic` (live control read) — see [`control_u64`].
418#[polydat::polydat_node(
419    category = Context,
420    purity = Nondeterministic("reads a live dynamic control value; changes over the run"),
421)]
422fn control_bool(name: Const<&str>) -> bool {
423    control_gauge_f64(name.0) != 0.0
424}
425
426/// Read a dynamic control's current value as its human-readable string
427/// (the erased-control `value_string()` rendering). Missing controls → `""`.
428///
429/// Signature: `control_str(name: const str) -> str`. Intrinsically
430/// `Nondeterministic` (live control read) — see [`control_u64`].
431#[polydat::polydat_node(
432    category = Context,
433    purity = Nondeterministic("reads a live dynamic control value; changes over the run"),
434)]
435fn control_str(name: Const<&str>) -> String {
436    control_value_string(name.0)
437}
438
439// =========================================================================
440// rate() / concurrency() — thin aliases over control(...)
441// =========================================================================
442
443/// Sugar for `control("rate")` — the current rate-limiter target (ops/sec).
444/// Intrinsically `Nondeterministic` (live control read) — see [`control`].
445#[polydat::polydat_node(
446    category = Context,
447    purity = Nondeterministic("reads the live rate control; changes over the run"),
448)]
449fn rate() -> f64 {
450    control_gauge_f64("rate")
451}
452
453/// Sugar for `control("concurrency")` — the current fiber count for the
454/// nearest phase. Intrinsically `Nondeterministic` — see [`control`].
455#[polydat::polydat_node(
456    category = Context,
457    purity = Nondeterministic("reads the live concurrency control; changes over the run"),
458)]
459fn concurrency() -> f64 {
460    control_gauge_f64("concurrency")
461}
462
463// =========================================================================
464// phase() — current phase name (thread-local)
465// =========================================================================
466
467/// Current phase name. Reads a thread-local set by the phase executor;
468/// `""` when unset (outside a cycle, or in tests that install no phase).
469///
470/// Intrinsically `Nondeterministic`: the per-fiber phase name varies across
471/// fibers and over the run, so the node is never const-folded.
472#[polydat::polydat_node(
473    category = Context,
474    purity = Nondeterministic("reads the per-fiber phase name; varies across fibers and over the run"),
475)]
476fn phase() -> String {
477    task_phase().map(|s| s.to_string()).unwrap_or_default()
478}
479
480// =========================================================================
481// phase_start_millis() / phase_elapsed_millis() — phase-scoped clock
482// =========================================================================
483
484/// Epoch millis at which the CURRENT phase started; `0` outside a phase body.
485///
486/// The phase-scoped counterpart to `session_start_millis()`. Both exist because
487/// they answer different questions and a session origin cannot substitute for a
488/// phase one in a run that sweeps many phases.
489///
490/// The origin is established once per phase by the executor's `run_phase` and
491/// carried across fiber spawns, so every op under the phase reads the same
492/// value — including concurrent sibling phases, which each see their own
493/// (the origin is a task-local, not shared execution state).
494///
495/// Intrinsically `Nondeterministic`: the origin differs per phase and per run,
496/// so the node must never be const-folded into one phase's value.
497#[polydat::polydat_node(
498    category = Context,
499    purity = Nondeterministic("reads the current phase's start time; differs per phase and per run"),
500)]
501fn phase_start_millis() -> u64 {
502    crate::execution_context::current_phase_start_ms().unwrap_or(0)
503}
504
505/// Milliseconds elapsed since the CURRENT phase started; `0` outside a phase.
506///
507/// `current_epoch_millis() - phase_start_millis()` spelled as one node, which
508/// is the form workloads actually want — a phase-relative duration available
509/// while the phase is still running, rather than only at its completion pull.
510///
511/// Intrinsically `Nondeterministic`: it grows monotonically within a phase.
512#[polydat::polydat_node(
513    category = Context,
514    purity = Nondeterministic("monotonic elapsed time within the current phase"),
515)]
516fn phase_elapsed_millis() -> u64 {
517    let Some(start) = crate::execution_context::current_phase_start_ms() else {
518        return 0;
519    };
520    let now = std::time::SystemTime::now()
521        .duration_since(std::time::UNIX_EPOCH)
522        .map(|d| d.as_millis() as u64)
523        .unwrap_or(0);
524    now.saturating_sub(start)
525}
526
527// =========================================================================
528// cycle() — current cycle ordinal (thread-local)
529// =========================================================================
530
531/// Current cycle ordinal. Reads a thread-local set by the phase executor.
532/// For bindings that already declare `cycle` as a named input this is
533/// redundant; it exists so bindings which never named `cycle` explicitly can
534/// still reach it (SRD 10's "cycle is not magic" rule — the node is context,
535/// not a privileged input).
536///
537/// Intrinsically `Nondeterministic`: the cycle ordinal changes every cycle,
538/// so the node is never const-folded.
539#[polydat::polydat_node(
540    category = Context,
541    purity = Nondeterministic("reads the per-fiber cycle ordinal; changes every cycle"),
542)]
543fn cycle() -> u64 {
544    task_cycle()
545}
546
547// =========================================================================
548// Registration
549// =========================================================================
550//
551// Every node in this module is authored via `#[polydat::polydat_node]`
552// (SRD-80b) — the control readers (`control` / `control_u64` /
553// `control_bool` / `control_str` / `rate` / `concurrency`), the fiber-context
554// readers (`phase` / `cycle`), and the `control_set` writer. Each macro
555// emits its own FuncSig + builder + `inventory::submit!`, so there is no
556// hand-written `signatures()` / `build_node()` / `register_nodes!` here.
557
558#[cfg(test)]
559mod tests {
560    // The async tests hold the `serial_test()` guard (a pure `Mutex<()>`)
561    // across `.await` to serialize global session-root installs; the
562    // awaited code never locks it, so there's no deadlock.
563    #![allow(clippy::await_holding_lock)]
564    use super::*;
565    use nmbrs_metrics::controls::{BranchScope, ControlBuilder};
566    use nmbrs_metrics::labels::Labels;
567    use polydat::ast::Value;
568    use std::collections::HashMap;
569    use std::sync::MutexGuard;
570
571    /// Serializes test access to the global `SESSION_ROOT` and the
572    /// thread-locals — the crate-wide guard hoisted next to the global
573    /// itself, so session-construction tests in OTHER modules can hold the
574    /// same lock (see `session_root_test_guard`).
575    fn serial_test() -> MutexGuard<'static, ()> {
576        super::session_root_test_guard()
577    }
578
579    /// Build a session root with a declared control, install it
580    /// as the global, and return the root handle so callers can
581    /// mutate the control further.
582    fn install_session_with_control(name: &str, initial: u32) -> Arc<RwLock<Component>> {
583        let root = Component::root(Labels::empty().with("session", "t"), HashMap::new());
584        root.read().unwrap().controls().declare(
585            ControlBuilder::new(name, initial)
586                .reify_as_gauge(|v| Some(*v as f64))
587                .branch_scope(BranchScope::Subtree)
588                .build(),
589        );
590        set_session_root(root.clone());
591        root
592    }
593
594    #[test]
595    fn control_reads_current_value() {
596        let _g = serial_test();
597        install_session_with_control("rate", 500);
598        let mut k =
599            polydat::dsl::compile_polydat_interpreter("x := control(\"rate\")").expect("compile");
600        assert_eq!(k.pull_ref("x").as_f64(), 500.0);
601    }
602
603    #[test]
604    fn control_missing_name_returns_zero() {
605        let _g = serial_test();
606        install_session_with_control("rate", 500);
607        let mut k = polydat::dsl::compile_polydat_interpreter("x := control(\"not_declared\")")
608            .expect("compile");
609        assert_eq!(k.pull_ref("x").as_f64(), 0.0);
610    }
611
612    // Live re-read after a write is covered end-to-end by
613    // `fiber_writes_control_via_control_set_and_reads_back` (the integration
614    // test, which yields for the async commit then re-pulls `control(...)`);
615    // `control_u64_is_volatile_not_const_folded` covers the const-fold
616    // property here. A unit test that pulls immediately after an async write
617    // would race the background commit, so it lives at the integration tier.
618
619    #[test]
620    fn rate_node_is_alias_of_control_rate() {
621        let _g = serial_test();
622        install_session_with_control("rate", 750);
623        let mut k = polydat::dsl::compile_polydat_interpreter("x := rate()").expect("compile");
624        assert_eq!(k.pull_ref("x").as_f64(), 750.0);
625    }
626
627    #[test]
628    fn concurrency_node_reads_concurrency_control() {
629        let _g = serial_test();
630        install_session_with_control("concurrency", 32);
631        let mut k =
632            polydat::dsl::compile_polydat_interpreter("x := concurrency()").expect("compile");
633        assert_eq!(k.pull_ref("x").as_f64(), 32.0);
634    }
635
636    #[tokio::test]
637    async fn phase_and_cycle_read_from_task_locals() {
638        let phase_arc: Arc<str> = Arc::from("rampup");
639        with_fiber_context(phase_arc.clone(), empty_controls(), async {
640            set_task_cycle(4242);
641            let mut k = polydat::dsl::compile_polydat_interpreter("p := phase()\nc := cycle()")
642                .expect("compile phase/cycle");
643            assert_eq!(k.pull_ref("p").as_str(), "rampup");
644            assert_eq!(k.pull_ref("c").as_u64(), 4242);
645        })
646        .await;
647    }
648
649    #[test]
650    fn phase_is_empty_outside_fiber_scope() {
651        // Reading outside a fiber context — e.g. from a unit test or a
652        // non-fiber call site — silently returns the empty string rather
653        // than panicking (the task_local's `try_with` Err maps to default).
654        let mut k = polydat::dsl::compile_polydat_interpreter("p := phase()").expect("compile");
655        assert_eq!(k.pull_ref("p").as_str(), "");
656    }
657
658    #[test]
659    fn cycle_is_zero_outside_fiber_scope() {
660        let mut k = polydat::dsl::compile_polydat_interpreter("c := cycle()").expect("compile");
661        assert_eq!(k.pull_ref("c").as_u64(), 0);
662    }
663
664    #[tokio::test]
665    async fn set_task_cycle_is_noop_outside_scope() {
666        // A stray call with no active scope must not panic.
667        set_task_cycle(99);
668        let mut k = polydat::dsl::compile_polydat_interpreter("c := cycle()").expect("compile");
669        assert_eq!(k.pull_ref("c").as_u64(), 0);
670    }
671
672    // ---- control_set ------------------------------------------
673
674    #[tokio::test]
675    async fn control_set_writes_through_converter_and_reaches_committed() {
676        let _g = serial_test();
677        // Install a root with a concurrency control that accepts
678        // f64 writes via an explicit from_f64 converter.
679        let root = Component::root(
680            Labels::empty().with("session", "s_cs"),
681            std::collections::HashMap::new(),
682        );
683        let c: nmbrs_metrics::controls::Control<u32> =
684            nmbrs_metrics::controls::ControlBuilder::new("concurrency", 4u32)
685                .reify_as_gauge(|v| Some(*v as f64))
686                .from_f64(|v| {
687                    if v < 0.0 || v > u32::MAX as f64 {
688                        Err(format!("out of range: {v}"))
689                    } else {
690                        Ok(v as u32)
691                    }
692                })
693                .branch_scope(nmbrs_metrics::controls::BranchScope::Subtree)
694                .build();
695        root.read().unwrap().controls().declare(c.clone());
696        set_session_root(root);
697
698        // Issue the write through the macro-authored node, built via the same
699        // factory route the compiler uses (under a binding scope).
700        let ctx = polydat::dsl::factory::BuildContext::with_binding("feedback_loop");
701        let consts = [polydat::dsl::factory::ConstArg::Str("concurrency".into())];
702        let node = polydat::dsl::factory::build_node(&ctx, "control_set", &[], &[], &consts)
703            .expect("control_set should build");
704        let mut out = [Value::None];
705        node.eval(&[Value::F64(64.0)], &mut out);
706        assert_eq!(out[0].as_u64(), 1, "write should report submitted");
707
708        // The write is async; yield a few times for the spawned
709        // task to run through validate → fanout → commit.
710        for _ in 0..10 {
711            tokio::task::yield_now().await;
712            if c.value() == 64u32 {
713                break;
714            }
715        }
716        assert_eq!(c.value(), 64u32);
717        let committed = c.get();
718        assert!(matches!(
719            committed.origin,
720            nmbrs_metrics::controls::ControlOrigin::Polydat { .. }
721        ));
722
723        // Idempotent-write elision: the control now reads 64.0, so a
724        // second write of the SAME value is a fixpoint — elided, 0
725        // (not dispatched), revision unmoved. A different value
726        // dispatches again.
727        let rev_before = c.get().rev;
728        node.eval(&[Value::F64(64.0)], &mut out);
729        assert_eq!(
730            out[0].as_u64(),
731            0,
732            "write of the committed value must be elided"
733        );
734        assert_eq!(
735            c.get().rev,
736            rev_before,
737            "an elided write must not touch the control"
738        );
739        node.eval(&[Value::F64(32.0)], &mut out);
740        assert_eq!(out[0].as_u64(), 1, "a real change dispatches");
741        for _ in 0..10 {
742            tokio::task::yield_now().await;
743            if c.value() == 32u32 {
744                break;
745            }
746        }
747        assert_eq!(c.value(), 32u32);
748    }
749
750    #[test]
751    fn control_u64_casts_gauge_to_integer() {
752        let _g = serial_test();
753        install_session_with_control("concurrency", 64);
754        // `control_u64` is macro-authored — compile + pull it end to end.
755        let mut k = polydat::dsl::compile_polydat_interpreter("x := control_u64(\"concurrency\")")
756            .expect("compile control_u64");
757        assert_eq!(k.pull_ref("x").as_u64(), 64);
758    }
759
760    #[test]
761    fn control_u64_missing_name_returns_zero() {
762        let _g = serial_test();
763        install_session_with_control("concurrency", 5);
764        let mut k = polydat::dsl::compile_polydat_interpreter("x := control_u64(\"not_there\")")
765            .expect("compile control_u64");
766        assert_eq!(k.pull_ref("x").as_u64(), 0);
767    }
768
769    #[test]
770    fn control_u64_is_volatile_not_const_folded() {
771        let _g = serial_test();
772        install_session_with_control("concurrency", 32);
773        // Intrinsic volatility (the bug this fixes): a `Nondeterministic`
774        // node is never const-folded, so the wire stays a live dynamic
775        // output re-read on every pull. A Pure reader would fold `x` to a
776        // compile-time constant — which is exactly how the old hand-written
777        // node (no purity override) cached a stale first value.
778        let k = polydat::dsl::compile_polydat_interpreter("x := control_u64(\"concurrency\")")
779            .expect("compile control_u64");
780        assert!(
781            k.get_constant("x").is_none(),
782            "control_u64 must be volatile — its wire must NOT be const-folded",
783        );
784    }
785
786    #[test]
787    fn control_bool_projects_gauge_to_boolean() {
788        let _g = serial_test();
789        install_session_with_control("enabled", 1);
790        let mut k = polydat::dsl::compile_polydat_interpreter("x := control_bool(\"enabled\")")
791            .expect("compile");
792        assert!(k.pull_ref("x").as_bool());
793    }
794
795    #[test]
796    fn control_bool_zero_is_false() {
797        let _g = serial_test();
798        install_session_with_control("enabled", 0);
799        let mut k = polydat::dsl::compile_polydat_interpreter("x := control_bool(\"enabled\")")
800            .expect("compile");
801        assert!(!k.pull_ref("x").as_bool());
802    }
803
804    #[test]
805    fn control_bool_missing_name_is_false() {
806        let _g = serial_test();
807        install_session_with_control("enabled", 1);
808        let mut k = polydat::dsl::compile_polydat_interpreter("x := control_bool(\"absent\")")
809            .expect("compile");
810        assert!(!k.pull_ref("x").as_bool());
811    }
812
813    #[test]
814    fn control_str_renders_value_string() {
815        let _g = serial_test();
816        install_session_with_control("concurrency", 42);
817        let mut k = polydat::dsl::compile_polydat_interpreter("x := control_str(\"concurrency\")")
818            .expect("compile");
819        // u32's Debug rendering is its decimal representation.
820        assert_eq!(k.pull_ref("x").as_str(), "42");
821    }
822
823    #[test]
824    fn control_str_missing_name_returns_empty() {
825        let _g = serial_test();
826        install_session_with_control("concurrency", 42);
827        let mut k = polydat::dsl::compile_polydat_interpreter("x := control_str(\"log_level\")")
828            .expect("compile");
829        assert_eq!(k.pull_ref("x").as_str(), "");
830    }
831
832    #[tokio::test]
833    async fn control_set_records_compile_time_binding_attribution() {
834        let _g = serial_test();
835        // Build a control that accepts f64 writes and install
836        // the session root.
837        let root = Component::root(
838            Labels::empty().with("session", "attr"),
839            std::collections::HashMap::new(),
840        );
841        let c: nmbrs_metrics::controls::Control<f64> =
842            nmbrs_metrics::controls::ControlBuilder::new("rate", 100.0)
843                .reify_as_gauge(|v| Some(*v))
844                .from_f64(Ok)
845                .branch_scope(nmbrs_metrics::controls::BranchScope::Subtree)
846                .build();
847        root.read().unwrap().controls().declare(c.clone());
848        set_session_root(root.clone());
849
850        // Simulate a compiler constructing a control_set factory
851        // under a binding scope named `rate_adj`. We can't call
852        // `build_node` from inside the same crate's private
853        // factory path directly, so reach into the nodes
854        // registration helper via the same build-by-name route
855        // the compiler uses.
856        let ctx = polydat::dsl::factory::BuildContext::with_binding("rate_adj");
857        let consts = [polydat::dsl::factory::ConstArg::Str("rate".into())];
858        let node = polydat::dsl::factory::build_node(&ctx, "control_set", &[], &[], &consts)
859            .expect("control_set should build");
860        let mut out = [Value::None];
861        node.eval(&[Value::F64(4242.0)], &mut out);
862
863        // Let the spawned write complete.
864        for _ in 0..40 {
865            tokio::time::sleep(std::time::Duration::from_millis(5)).await;
866            if c.value() == 4242.0 {
867                break;
868            }
869        }
870        assert_eq!(c.value(), 4242.0);
871        match c.get().origin {
872            nmbrs_metrics::controls::ControlOrigin::Polydat { ref binding } => {
873                assert_eq!(
874                    binding, "rate_adj",
875                    "attribution should be the DSL binding name, not the control name"
876                );
877            }
878            other => panic!("expected Polydat origin, got {other:?}"),
879        }
880    }
881
882    #[test]
883    fn control_set_returns_zero_without_session_root() {
884        let _g = serial_test();
885        // Explicitly clear the session root so the node can't
886        // resolve anything.
887        *SESSION_ROOT.lock().unwrap_or_else(|e| e.into_inner()) = None;
888
889        let consts = [polydat::dsl::factory::ConstArg::Str("anything".into())];
890        let node = polydat::dsl::factory::build_node(
891            &polydat::dsl::factory::BuildContext::default(),
892            "control_set",
893            &[],
894            &[],
895            &consts,
896        )
897        .expect("control_set should build");
898        let mut out = [Value::None];
899        node.eval(&[Value::F64(1.0)], &mut out);
900        assert_eq!(out[0].as_u64(), 0);
901    }
902
903    // ---- SRD-89: per-execution control isolation -------------------
904
905    /// Build a one-entry control map carrying a `concurrency` control fixed at
906    /// `val`, as `install_controls` expects (name → erased handle).
907    /// A standalone phase-tier component declaring a `concurrency` control at
908    /// `val` — the per-execution phase component a fiber resolves against.
909    fn component_with_concurrency(val: u32) -> Arc<RwLock<Component>> {
910        let comp = Component::root(Labels::empty().with("phase", "p"), HashMap::new());
911        comp.read()
912            .unwrap_or_else(|e| e.into_inner())
913            .controls()
914            .declare(
915                ControlBuilder::new("concurrency", val)
916                    .reify_as_gauge(|v| Some(*v as f64))
917                    .branch_scope(BranchScope::Subtree)
918                    .build(),
919            );
920        comp
921    }
922
923    #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
924    async fn per_execution_control_map_isolates_concurrent_reads() {
925        // SRD-89 §2b — two executions sharing one session each declare a
926        // `concurrency` control with a DISTINCT value on their own phase
927        // component. Under the pre-SRD-89 shared-`SESSION_ROOT` model both
928        // resolve to one instance (deterministic cross-talk — a servo
929        // retargeting one drives the other). The per-execution control map
930        // isolates them: each execution reads its OWN setting. This is the unit
931        // encoding of the walker's `control`/`multiservo` failure.
932        let _g = serial_test();
933        // Clear the global root so resolution can ONLY come from each fiber's
934        // own current-component snapshot (a leftover sibling root must not mask
935        // a miss).
936        *SESSION_ROOT.lock().unwrap_or_else(|e| e.into_inner()) = None;
937
938        // Two executions, each with its OWN phase component declaring
939        // `concurrency` at a distinct value.
940        let comp_a = component_with_concurrency(2);
941        let comp_b = component_with_concurrency(32);
942        let phase: Arc<str> = Arc::from("p");
943
944        // Each fiber resolves through its own component snapshot (FIBER_CTX).
945        let a_val = with_fiber_context(phase.clone(), snapshot_controls(&comp_a), async {
946            control_gauge_f64("concurrency")
947        })
948        .await;
949        let b_val = with_fiber_context(phase.clone(), snapshot_controls(&comp_b), async {
950            control_gauge_f64("concurrency")
951        })
952        .await;
953
954        assert_eq!(a_val, 2.0, "execution A must read its OWN concurrency (2)");
955        assert_eq!(
956            b_val, 32.0,
957            "execution B must read its OWN concurrency (32), not A's",
958        );
959    }
960
961    #[tokio::test(flavor = "multi_thread", worker_threads = 1)]
962    async fn control_read_falls_back_to_session_root_without_exec_context() {
963        // A1 — outside any scoped execution (single-run / dryrun / TUI), the
964        // read path is byte-identical to before SRD-89: it resolves via the
965        // global SESSION_ROOT walk. No per-exec map exists, so the fallback is
966        // the only path.
967        let _g = serial_test();
968        install_session_with_control("concurrency", 7);
969        // No execution_context::scope here — current_control() returns None.
970        assert_eq!(control_gauge_f64("concurrency"), 7.0);
971    }
972}