#![allow(clippy::await_holding_lock)]
use std::collections::HashMap;
use std::sync::Arc;
use std::time::Duration;
use nmbrs_metrics::component::Component;
use nmbrs_metrics::controls::{BranchScope, ControlBuilder, ControlOrigin};
use nmbrs_metrics::labels::Labels;
use nmbrs_runtime::polydat_nodes::runtime_context::{
empty_controls, set_session_root, set_task_cycle, snapshot_controls, with_fiber_context,
};
use polydat::ast::{PolydatNode, Value};
use polydat::library::param_helpers::{InRange, IsPositive, Required, ThisOr};
use std::sync::Mutex;
static TEST_LOCK: Mutex<()> = Mutex::new(());
fn build_session_with_concurrency(initial: u32) -> Arc<std::sync::RwLock<Component>> {
let root = Component::root(Labels::empty().with("session", "integ"), HashMap::new());
root.read().unwrap().controls().declare(
ControlBuilder::new("concurrency", initial)
.reify_as_gauge(|v| Some(*v as f64))
.branch_scope(BranchScope::Subtree)
.from_f64(|v| {
if !(0.0..=10_000.0).contains(&v) {
Err(format!("concurrency out of range: {v}"))
} else {
Ok(v as u32)
}
})
.build(),
);
set_session_root(root.clone());
root
}
#[tokio::test]
async fn fiber_reads_phase_and_cycle_from_task_context() {
let _g = TEST_LOCK.lock().unwrap();
let phase: Arc<str> = Arc::from("rampup");
with_fiber_context(phase.clone(), empty_controls(), async {
for cycle in [0u64, 1, 17, 999] {
set_task_cycle(cycle);
let mut k = polydat::dsl::compile_polydat("p := phase()\nc := cycle()")
.expect("compile phase/cycle");
assert_eq!(k.pull("p").as_str(), "rampup");
assert_eq!(k.pull("c").as_u64(), cycle);
}
})
.await;
}
#[tokio::test]
async fn param_helpers_pass_happy_values() {
let _g = TEST_LOCK.lock().unwrap();
let n = Required::new("cycles".to_string());
let mut out = [Value::None];
n.eval(&[Value::U64(10_000)], &mut out);
assert_eq!(out[0].as_u64(), 10_000);
let n = IsPositive::new("rate".to_string());
let mut out = [Value::None];
n.eval(&[Value::U64(1)], &mut out);
assert_eq!(out[0].as_u64(), 1);
let n = InRange::new(1, 100);
let mut out = [Value::None];
n.eval(&[Value::U64(50)], &mut out);
assert_eq!(out[0].as_u64(), 50);
let n = ThisOr::new();
let mut out = [Value::None];
n.eval(&[Value::U64(7), Value::U64(99)], &mut out);
assert_eq!(out[0].as_u64(), 7);
n.eval(&[Value::None, Value::U64(99)], &mut out);
assert_eq!(out[0].as_u64(), 99);
}
#[tokio::test]
async fn fiber_reads_control_through_context() {
let _g = TEST_LOCK.lock().unwrap();
let root = build_session_with_concurrency(8);
let phase: Arc<str> = Arc::from("rampup");
with_fiber_context(phase, snapshot_controls(&root), async {
set_task_cycle(0);
let mut k = polydat::dsl::compile_polydat(
"c := control(\"concurrency\")\n\
u := control_u64(\"concurrency\")\n\
s := control_str(\"concurrency\")",
)
.expect("compile control readers");
assert_eq!(k.pull("c").as_f64(), 8.0);
assert_eq!(k.pull("u").as_u64(), 8);
assert_eq!(k.pull("s").as_str(), "8");
})
.await;
let _ = root;
}
#[tokio::test]
async fn fiber_writes_control_via_control_set_and_reads_back() {
let _g = TEST_LOCK.lock().unwrap();
let root = build_session_with_concurrency(8);
let phase: Arc<str> = Arc::from("rampup");
with_fiber_context(phase, snapshot_controls(&root), async {
let ctx = polydat::dsl::factory::BuildContext::with_binding("integration_feedback_loop");
let consts = [polydat::dsl::factory::ConstArg::Str("concurrency".into())];
let writer = polydat::dsl::factory::build_node(&ctx, "control_set", &[], &[], &consts)
.expect("build control_set");
let mut write_out = [Value::None];
writer.eval(&[Value::F64(42.0)], &mut write_out);
assert_eq!(write_out[0].as_u64(), 1, "write should report submitted");
let mut k =
polydat::dsl::compile_polydat("r := control(\"concurrency\")").expect("compile");
let mut observed = 0.0;
for _ in 0..40 {
tokio::time::sleep(Duration::from_millis(5)).await;
observed = k.pull("r").as_f64();
if observed == 42.0 {
break;
}
}
assert_eq!(
observed, 42.0,
"control_set's write should commit and be visible via read"
);
})
.await;
let control: nmbrs_metrics::controls::Control<u32> =
root.read().unwrap().controls().get("concurrency").unwrap();
let versioned = control.get();
assert_eq!(versioned.value, 42);
assert!(
matches!(versioned.origin, ControlOrigin::Polydat { ref binding } if binding == "integration_feedback_loop"),
"expected Polydat origin tagged with the feedback_loop binding, got {:?}",
versioned.origin,
);
}
#[tokio::test]
async fn control_set_out_of_range_leaves_value_unchanged() {
let _g = TEST_LOCK.lock().unwrap();
let root = build_session_with_concurrency(16);
let control: nmbrs_metrics::controls::Control<u32> =
root.read().unwrap().controls().get("concurrency").unwrap();
let phase: Arc<str> = Arc::from("rampup");
with_fiber_context(phase, snapshot_controls(&root), async {
let consts = [polydat::dsl::factory::ConstArg::Str("concurrency".into())];
let writer = polydat::dsl::factory::build_node(
&polydat::dsl::factory::BuildContext::default(),
"control_set",
&[],
&[],
&consts,
)
.expect("build control_set");
let mut write_out = [Value::None];
writer.eval(&[Value::F64(99_999.0)], &mut write_out);
assert_eq!(write_out[0].as_u64(), 1);
tokio::time::sleep(Duration::from_millis(30)).await;
})
.await;
assert_eq!(control.value(), 16);
assert_eq!(control.get().rev, 0);
}
#[tokio::test]
async fn branch_scoped_control_resolves_from_descendant_fiber() {
let _g = TEST_LOCK.lock().unwrap();
let root = Component::root(Labels::empty().with("session", "integ_bs"), HashMap::new());
root.read().unwrap().controls().declare(
ControlBuilder::new("hdr_sigdigs", 4u32)
.reify_as_gauge(|v| Some(*v as f64))
.branch_scope(BranchScope::Subtree)
.build(),
);
set_session_root(root.clone());
let phase: Arc<str> = Arc::from("any_phase");
with_fiber_context(phase, snapshot_controls(&root), async {
let mut k =
polydat::dsl::compile_polydat("r := control(\"hdr_sigdigs\")").expect("compile");
assert_eq!(k.pull("r").as_f64(), 4.0);
})
.await;
}