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}