1use std::collections::HashMap;
32use std::fmt;
33use std::future::Future;
34use std::pin::Pin;
35use std::sync::atomic::{AtomicU64, Ordering};
36use std::sync::{Arc, Mutex, OnceLock};
37use std::time::{Duration, Instant};
38
39use futures::future::join_all;
40use tokio::sync::Mutex as AsyncMutex;
41
42use crate::instruments::gauge::ValueGauge;
43use crate::labels::Labels;
44use crate::snapshot::MetricSet;
45
46#[derive(Debug, Default)]
61pub struct RevAllocator {
62 counter: AtomicU64,
63}
64
65impl RevAllocator {
66 pub fn new() -> Self {
67 Self::default()
68 }
69
70 pub fn next(&self) -> u64 {
71 self.counter.fetch_add(1, Ordering::AcqRel).wrapping_add(1)
72 }
73
74 pub fn global() -> &'static Arc<RevAllocator> {
79 static GLOBAL: OnceLock<Arc<RevAllocator>> = OnceLock::new();
80 GLOBAL.get_or_init(|| Arc::new(RevAllocator::new()))
81 }
82}
83
84#[derive(Clone, Debug, PartialEq, Eq)]
92pub enum ControlOrigin {
93 Launch,
95 Cli,
97 Tui,
99 Polydat { binding: String },
101 Governor { source: String },
104 Api { source: String },
106 Test,
108}
109
110#[derive(Clone, Debug)]
112pub struct Versioned<T: Clone> {
113 pub value: T,
114 pub rev: u64,
118 pub updated_at: Instant,
120 pub origin: ControlOrigin,
121}
122
123#[derive(Clone, Debug, PartialEq, Eq)]
129pub enum SetError {
130 ValidationFailed(String),
133 ApplyFailed(Vec<ApplyFailure>),
136 FinalViolation { scope: String },
140}
141
142#[derive(Clone, Copy, Debug, PartialEq, Eq)]
148pub enum BranchScope {
149 Local,
154 Subtree,
160}
161
162#[derive(Clone, Debug, PartialEq, Eq)]
163pub struct ApplyFailure {
164 pub applier_index: usize,
169 pub message: String,
172}
173
174impl fmt::Display for SetError {
175 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
176 match self {
177 Self::ValidationFailed(msg) => write!(f, "validation failed: {msg}"),
178 Self::ApplyFailed(failures) => {
179 write!(f, "{} applier(s) failed:", failures.len())?;
180 for fail in failures {
181 write!(f, " [#{}: {}]", fail.applier_index, fail.message)?;
182 }
183 Ok(())
184 }
185 Self::FinalViolation { scope } => {
186 write!(
187 f,
188 "control is declared final at scope '{scope}'; runtime writes are rejected"
189 )
190 }
191 }
192 }
193}
194
195impl std::error::Error for SetError {}
196
197pub trait ControlApplier<T>: Send + Sync + 'static
210where
211 T: Clone + Send + Sync + 'static,
212{
213 fn apply(&self, value: T) -> Pin<Box<dyn Future<Output = Result<(), String>> + Send + '_>>;
214}
215
216pub struct SyncApplier<T, F>
221where
222 T: Clone + Send + Sync + 'static,
223 F: Fn(T) -> Result<(), String> + Send + Sync + 'static,
224{
225 f: F,
226 _marker: std::marker::PhantomData<fn(T)>,
227}
228
229impl<T, F> SyncApplier<T, F>
230where
231 T: Clone + Send + Sync + 'static,
232 F: Fn(T) -> Result<(), String> + Send + Sync + 'static,
233{
234 pub fn new(f: F) -> Self {
235 Self {
236 f,
237 _marker: std::marker::PhantomData,
238 }
239 }
240}
241
242impl<T, F> ControlApplier<T> for SyncApplier<T, F>
243where
244 T: Clone + Send + Sync + 'static,
245 F: Fn(T) -> Result<(), String> + Send + Sync + 'static,
246{
247 fn apply(&self, value: T) -> Pin<Box<dyn Future<Output = Result<(), String>> + Send + '_>> {
248 let out = (self.f)(value);
249 Box::pin(async move { out })
250 }
251}
252
253type Validator<T> = Box<dyn Fn(&T) -> Result<(), String> + Send + Sync>;
258type ToF64<T> = Box<dyn Fn(&T) -> Option<f64> + Send + Sync>;
259type FromF64<T> = Box<dyn Fn(f64) -> Result<T, String> + Send + Sync>;
260
261pub struct Control<T: Clone + Send + Sync + 'static> {
262 inner: Arc<ControlInner<T>>,
263}
264
265impl<T: Clone + Send + Sync + 'static> Clone for Control<T> {
266 fn clone(&self) -> Self {
267 Self {
268 inner: self.inner.clone(),
269 }
270 }
271}
272
273struct ControlInner<T: Clone + Send + Sync + 'static> {
274 name: String,
275 write_lock: AsyncMutex<()>,
279 committed: std::sync::RwLock<Versioned<T>>,
283 appliers: Mutex<Vec<Arc<dyn ControlApplier<T>>>>,
284 validator: OnceLock<Validator<T>>,
285 apply_timeout: Duration,
286 rev_allocator: Arc<RevAllocator>,
287 gauge: Option<GaugeReification<T>>,
293 final_at_scope: Option<String>,
299 branch_scope: BranchScope,
302 from_f64: Option<FromF64<T>>,
307}
308
309struct GaugeReification<T> {
310 gauge: Arc<ValueGauge>,
311 to_f64: ToF64<T>,
312}
313
314impl<T: Clone + Send + Sync + 'static> Control<T> {
315 pub fn name(&self) -> &str {
318 &self.inner.name
319 }
320
321 pub fn get(&self) -> Versioned<T> {
323 self.inner
324 .committed
325 .read()
326 .unwrap_or_else(|e| e.into_inner())
327 .clone()
328 }
329
330 pub fn value(&self) -> T {
333 self.get().value
334 }
335
336 pub fn register_applier<A>(&self, applier: A) -> usize
342 where
343 A: ControlApplier<T>,
344 {
345 let arc: Arc<dyn ControlApplier<T>> = Arc::new(applier);
346 let mut g = self
347 .inner
348 .appliers
349 .lock()
350 .unwrap_or_else(|e| e.into_inner());
351 g.push(arc);
352 g.len() - 1
353 }
354
355 pub fn applier_count(&self) -> usize {
357 self.inner
358 .appliers
359 .lock()
360 .unwrap_or_else(|e| e.into_inner())
361 .len()
362 }
363
364 pub async fn set(&self, value: T, origin: ControlOrigin) -> Result<u64, SetError> {
377 let _guard = self.inner.write_lock.lock().await;
378
379 if let Some(ref scope) = self.inner.final_at_scope
384 && origin != ControlOrigin::Launch
385 {
386 return Err(SetError::FinalViolation {
387 scope: scope.clone(),
388 });
389 }
390
391 if let Some(validator) = self.inner.validator.get() {
393 validator(&value).map_err(SetError::ValidationFailed)?;
394 }
395
396 let appliers: Vec<Arc<dyn ControlApplier<T>>> = {
400 let g = self
401 .inner
402 .appliers
403 .lock()
404 .unwrap_or_else(|e| e.into_inner());
405 g.clone()
406 };
407
408 let timeout = self.inner.apply_timeout;
409 let futures = appliers.iter().enumerate().map(|(idx, applier)| {
410 let v = value.clone();
411 let applier = applier.clone();
412 async move {
413 let fut = applier.apply(v);
414 match tokio::time::timeout(timeout, fut).await {
415 Ok(Ok(())) => Ok(idx),
416 Ok(Err(msg)) => Err(ApplyFailure {
417 applier_index: idx,
418 message: msg,
419 }),
420 Err(_) => Err(ApplyFailure {
421 applier_index: idx,
422 message: format!("apply timed out after {:?}", timeout),
423 }),
424 }
425 }
426 });
427
428 let results = join_all(futures).await;
429 let failures: Vec<ApplyFailure> = results.into_iter().filter_map(|r| r.err()).collect();
430
431 if !failures.is_empty() {
432 return Err(SetError::ApplyFailed(failures));
433 }
434
435 let rev = self.inner.rev_allocator.next();
437 let versioned = Versioned {
438 value: value.clone(),
439 rev,
440 updated_at: Instant::now(),
441 origin,
442 };
443 *self
444 .inner
445 .committed
446 .write()
447 .unwrap_or_else(|e| e.into_inner()) = versioned;
448 self.publish_gauge(&value);
453
454 Ok(rev)
455 }
456
457 pub fn gauge_f64(&self) -> Option<f64> {
463 let g = self.inner.gauge.as_ref()?;
464 let v = self.get().value;
465 (g.to_f64)(&v)
466 }
467
468 pub fn reified_gauge(&self) -> Option<Arc<ValueGauge>> {
475 self.inner.gauge.as_ref().map(|g| g.gauge.clone())
476 }
477
478 pub fn is_final(&self) -> bool {
482 self.inner.final_at_scope.is_some()
483 }
484
485 pub fn final_scope(&self) -> Option<&str> {
488 self.inner.final_at_scope.as_deref()
489 }
490
491 pub fn branch_scope(&self) -> BranchScope {
493 self.inner.branch_scope
494 }
495
496 fn publish_gauge(&self, value: &T) {
497 if let Some(ref g) = self.inner.gauge
498 && let Some(f) = (g.to_f64)(value)
499 {
500 g.gauge.set(f);
501 }
502 }
503}
504
505pub struct ControlBuilder<T: Clone + Send + Sync + 'static> {
513 name: String,
514 initial: T,
515 rev_allocator: Arc<RevAllocator>,
516 apply_timeout: Duration,
517 validator: Option<Validator<T>>,
518 reify: Option<ToF64<T>>,
519 final_at_scope: Option<String>,
520 branch_scope: BranchScope,
521 from_f64: Option<FromF64<T>>,
522}
523
524impl<T: Clone + Send + Sync + 'static> ControlBuilder<T> {
525 pub fn new(name: &str, initial: T) -> Self {
526 Self {
527 name: name.to_string(),
528 initial,
529 rev_allocator: RevAllocator::global().clone(),
530 apply_timeout: Duration::from_secs(5),
535 validator: None,
536 reify: None,
537 final_at_scope: None,
538 branch_scope: BranchScope::Local,
539 from_f64: None,
540 }
541 }
542
543 pub fn rev_allocator(mut self, allocator: Arc<RevAllocator>) -> Self {
544 self.rev_allocator = allocator;
545 self
546 }
547
548 pub fn apply_timeout(mut self, d: Duration) -> Self {
549 self.apply_timeout = d;
550 self
551 }
552
553 pub fn validator<F>(mut self, f: F) -> Self
557 where
558 F: Fn(&T) -> Result<(), String> + Send + Sync + 'static,
559 {
560 self.validator = Some(Box::new(f));
561 self
562 }
563
564 pub fn reify_as_gauge<F>(mut self, to_f64: F) -> Self
576 where
577 F: Fn(&T) -> Option<f64> + Send + Sync + 'static,
578 {
579 self.reify = Some(Box::new(to_f64));
580 self
581 }
582
583 pub fn final_at_scope(mut self, scope_name: impl Into<String>) -> Self {
595 self.final_at_scope = Some(scope_name.into());
596 self
597 }
598
599 pub fn branch_scope(mut self, scope: BranchScope) -> Self {
608 self.branch_scope = scope;
609 self
610 }
611
612 pub fn from_f64<F>(mut self, f: F) -> Self
625 where
626 F: Fn(f64) -> Result<T, String> + Send + Sync + 'static,
627 {
628 self.from_f64 = Some(Box::new(f));
629 self
630 }
631
632 pub fn build(self) -> Control<T> {
633 let versioned = Versioned {
634 value: self.initial.clone(),
635 rev: 0,
636 updated_at: Instant::now(),
637 origin: ControlOrigin::Launch,
638 };
639 let validator_slot: OnceLock<Validator<T>> = OnceLock::new();
640 if let Some(v) = self.validator {
641 let _ = validator_slot.set(v);
642 }
643 let gauge = self.reify.map(|to_f64| {
647 let gauge = Arc::new(ValueGauge::new(Labels::of("control", &self.name)));
648 if let Some(f) = to_f64(&self.initial) {
649 gauge.set(f);
650 }
651 GaugeReification { gauge, to_f64 }
652 });
653 Control {
654 inner: Arc::new(ControlInner {
655 name: self.name,
656 write_lock: AsyncMutex::new(()),
657 committed: std::sync::RwLock::new(versioned),
658 appliers: Mutex::new(Vec::new()),
659 validator: validator_slot,
660 apply_timeout: self.apply_timeout,
661 rev_allocator: self.rev_allocator,
662 gauge,
663 final_at_scope: self.final_at_scope,
664 branch_scope: self.branch_scope,
665 from_f64: self.from_f64,
666 }),
667 }
668 }
669}
670
671pub trait ErasedControl: Send + Sync {
680 fn name(&self) -> &str;
681 fn rev(&self) -> u64;
682 fn origin(&self) -> ControlOrigin;
683 fn applier_count(&self) -> usize;
684 fn value_string(&self) -> String;
687 fn value_type_name(&self) -> &'static str;
690 fn gauge_f64(&self) -> Option<f64>;
695 fn has_reified_gauge(&self) -> bool;
697 fn is_final(&self) -> bool;
700 fn final_scope(&self) -> Option<String>;
702 fn branch_scope(&self) -> BranchScope;
704 fn accepts_f64_writes(&self) -> bool;
707 fn set_f64(
713 &self,
714 value: f64,
715 origin: ControlOrigin,
716 ) -> Pin<Box<dyn Future<Output = Result<u64, SetError>> + Send>>;
717}
718
719impl<T> ErasedControl for Control<T>
720where
721 T: Clone + Send + Sync + fmt::Debug + 'static,
722{
723 fn name(&self) -> &str {
724 Control::name(self)
725 }
726
727 fn rev(&self) -> u64 {
728 self.get().rev
729 }
730
731 fn origin(&self) -> ControlOrigin {
732 self.get().origin
733 }
734
735 fn applier_count(&self) -> usize {
736 Control::applier_count(self)
737 }
738
739 fn value_string(&self) -> String {
740 format!("{:?}", self.get().value)
741 }
742
743 fn value_type_name(&self) -> &'static str {
744 std::any::type_name::<T>()
745 }
746
747 fn gauge_f64(&self) -> Option<f64> {
748 Control::gauge_f64(self)
749 }
750
751 fn has_reified_gauge(&self) -> bool {
752 self.inner.gauge.is_some()
753 }
754
755 fn is_final(&self) -> bool {
756 Control::is_final(self)
757 }
758
759 fn final_scope(&self) -> Option<String> {
760 Control::final_scope(self).map(|s| s.to_string())
761 }
762
763 fn branch_scope(&self) -> BranchScope {
764 Control::branch_scope(self)
765 }
766
767 fn accepts_f64_writes(&self) -> bool {
768 self.inner.from_f64.is_some()
769 }
770
771 fn set_f64(
772 &self,
773 value: f64,
774 origin: ControlOrigin,
775 ) -> Pin<Box<dyn Future<Output = Result<u64, SetError>> + Send>> {
776 let inner = self.inner.clone();
777 let self_clone = Control { inner };
778 Box::pin(async move {
779 let converter = match self_clone.inner.from_f64.as_ref() {
780 Some(c) => c,
781 None => {
782 return Err(SetError::ValidationFailed(format!(
783 "control '{}' has no f64 setter registered — \
784 declare one via ControlBuilder::from_f64",
785 self_clone.inner.name,
786 )));
787 }
788 };
789 let typed = match converter(value) {
790 Ok(t) => t,
791 Err(msg) => return Err(SetError::ValidationFailed(msg)),
792 };
793 self_clone.set(typed, origin).await
794 })
795 }
796}
797
798#[derive(Default)]
807pub struct ControlRegistry {
808 entries: std::sync::RwLock<HashMap<String, Arc<dyn std::any::Any + Send + Sync>>>,
809 erased: std::sync::RwLock<HashMap<String, Arc<dyn ErasedControl>>>,
810}
811
812impl ControlRegistry {
813 pub fn new() -> Self {
814 Self::default()
815 }
816
817 pub fn declare<T>(&self, control: Control<T>)
821 where
822 T: Clone + Send + Sync + fmt::Debug + 'static,
823 {
824 let name = control.name().to_string();
825 let erased: Arc<dyn ErasedControl> = Arc::new(control.clone());
826 let typed: Arc<dyn std::any::Any + Send + Sync> = Arc::new(control);
827
828 let mut entries = self.entries.write().unwrap_or_else(|e| e.into_inner());
829 let mut erased_map = self.erased.write().unwrap_or_else(|e| e.into_inner());
830 assert!(
831 !entries.contains_key(&name),
832 "ControlRegistry: duplicate control declaration for '{name}'",
833 );
834 entries.insert(name.clone(), typed);
835 erased_map.insert(name, erased);
836 }
837
838 pub fn get<T>(&self, name: &str) -> Option<Control<T>>
842 where
843 T: Clone + Send + Sync + 'static,
844 {
845 let entries = self.entries.read().unwrap_or_else(|e| e.into_inner());
846 let any = entries.get(name)?.clone();
847 drop(entries);
848 any.downcast::<Control<T>>().ok().map(|arc| (*arc).clone())
849 }
850
851 pub fn get_erased(&self, name: &str) -> Option<Arc<dyn ErasedControl>> {
854 self.erased
855 .read()
856 .unwrap_or_else(|e| e.into_inner())
857 .get(name)
858 .cloned()
859 }
860
861 pub fn list(&self) -> Vec<Arc<dyn ErasedControl>> {
864 self.erased
865 .read()
866 .unwrap_or_else(|e| e.into_inner())
867 .values()
868 .cloned()
869 .collect()
870 }
871
872 pub fn len(&self) -> usize {
873 self.entries.read().unwrap_or_else(|e| e.into_inner()).len()
874 }
875
876 pub fn is_empty(&self) -> bool {
877 self.len() == 0
878 }
879
880 pub fn snapshot_gauges(&self, base_labels: &Labels, captured_at: Instant) -> MetricSet {
903 let mut set = MetricSet::at(captured_at, Duration::ZERO);
904 let erased = self.erased.read().unwrap_or_else(|e| e.into_inner());
905 for ctl in erased.values() {
906 if let Some(value) = ctl.gauge_f64() {
908 let family_name = format!("control_{}", ctl.name());
909 let labels = base_labels.with("control", ctl.name());
910 set.insert_gauge(&family_name, labels, value, captured_at);
911 }
912 let info_family = format!("control_info_{}", ctl.name());
917 let info_labels = base_labels
918 .with("control", ctl.name())
919 .with("value", ctl.value_string());
920 set.insert_gauge(&info_family, info_labels, 1.0, captured_at);
921 }
922 set
923 }
924}
925
926#[cfg(test)]
931mod tests {
932 use super::*;
933 use std::sync::atomic::{AtomicU32, AtomicUsize};
934
935 #[test]
938 fn rev_allocator_is_monotonic() {
939 let a = RevAllocator::new();
940 let r1 = a.next();
941 let r2 = a.next();
942 let r3 = a.next();
943 assert!(r1 < r2 && r2 < r3);
944 assert!(r1 >= 1);
947 }
948
949 fn build_u32(name: &str, initial: u32) -> Control<u32> {
952 ControlBuilder::new(name, initial)
953 .apply_timeout(Duration::from_secs(1))
954 .build()
955 }
956
957 #[tokio::test]
958 async fn initial_value_is_seeded() {
959 let c = build_u32("concurrency", 16);
960 let got = c.get();
961 assert_eq!(got.value, 16);
962 assert_eq!(got.rev, 0);
963 assert_eq!(got.origin, ControlOrigin::Launch);
964 }
965
966 #[tokio::test]
967 async fn set_with_no_appliers_commits_and_advances_rev() {
968 let c = build_u32("concurrency", 4);
969 let rev = c.set(32, ControlOrigin::Test).await.unwrap();
970 assert!(rev >= 1);
971 let v = c.get();
972 assert_eq!(v.value, 32);
973 assert_eq!(v.rev, rev);
974 assert_eq!(v.origin, ControlOrigin::Test);
975 }
976
977 #[tokio::test]
978 async fn two_successive_sets_produce_strictly_increasing_revs() {
979 let c = build_u32("c", 1);
980 let r1 = c.set(2, ControlOrigin::Test).await.unwrap();
981 let r2 = c.set(3, ControlOrigin::Test).await.unwrap();
982 assert!(r1 < r2);
983 }
984
985 #[tokio::test]
988 async fn validator_rejects_bad_values() {
989 let c: Control<u32> = ControlBuilder::new("concurrency", 4)
990 .validator(|v| {
991 if *v == 0 {
992 Err("must be > 0".into())
993 } else if *v > 10_000 {
994 Err("too large".into())
995 } else {
996 Ok(())
997 }
998 })
999 .build();
1000
1001 match c.set(0, ControlOrigin::Test).await {
1002 Err(SetError::ValidationFailed(msg)) => assert!(msg.contains("must be > 0")),
1003 other => panic!("expected ValidationFailed, got {other:?}"),
1004 }
1005 assert_eq!(c.value(), 4);
1007
1008 c.set(32, ControlOrigin::Test).await.unwrap();
1010 assert_eq!(c.value(), 32);
1011 }
1012
1013 #[tokio::test]
1016 async fn single_applier_sees_new_value() {
1017 let seen = Arc::new(AtomicU32::new(0));
1018 let c = build_u32("c", 0);
1019 let seen_clone = seen.clone();
1020 c.register_applier(SyncApplier::new(move |v: u32| {
1021 seen_clone.store(v, Ordering::SeqCst);
1022 Ok(())
1023 }));
1024 c.set(42, ControlOrigin::Test).await.unwrap();
1025 assert_eq!(seen.load(Ordering::SeqCst), 42);
1026 }
1027
1028 #[tokio::test]
1029 async fn multiple_appliers_all_see_same_value_on_success() {
1030 let c = build_u32("c", 0);
1031 let counts: Vec<Arc<AtomicU32>> = (0..5).map(|_| Arc::new(AtomicU32::new(0))).collect();
1032 for counter in &counts {
1033 let c2 = counter.clone();
1034 c.register_applier(SyncApplier::new(move |v: u32| {
1035 c2.store(v, Ordering::SeqCst);
1036 Ok(())
1037 }));
1038 }
1039 c.set(99, ControlOrigin::Test).await.unwrap();
1040 for counter in &counts {
1041 assert_eq!(counter.load(Ordering::SeqCst), 99);
1042 }
1043 }
1044
1045 #[tokio::test]
1048 async fn any_applier_error_fails_the_set_and_reports_index() {
1049 let c = build_u32("c", 0);
1050 c.register_applier(SyncApplier::new(|_| Ok(())));
1052 c.register_applier(SyncApplier::new(|_| Err("subsystem X offline".into())));
1053 c.register_applier(SyncApplier::new(|_| Ok(())));
1054
1055 match c.set(1, ControlOrigin::Test).await {
1056 Err(SetError::ApplyFailed(failures)) => {
1057 assert_eq!(failures.len(), 1);
1058 assert_eq!(failures[0].applier_index, 1);
1059 assert!(failures[0].message.contains("subsystem X offline"));
1060 }
1061 other => panic!("expected ApplyFailed, got {other:?}"),
1062 }
1063 assert_eq!(c.value(), 0);
1065 assert_eq!(c.get().rev, 0);
1066 }
1067
1068 #[tokio::test]
1069 async fn multiple_applier_failures_all_reported() {
1070 let c = build_u32("c", 0);
1071 c.register_applier(SyncApplier::new(|_| Err("first".into())));
1072 c.register_applier(SyncApplier::new(|_| Ok(())));
1073 c.register_applier(SyncApplier::new(|_| Err("third".into())));
1074
1075 match c.set(1, ControlOrigin::Test).await {
1076 Err(SetError::ApplyFailed(failures)) => {
1077 assert_eq!(failures.len(), 2);
1078 let indices: Vec<usize> = failures.iter().map(|f| f.applier_index).collect();
1079 assert!(indices.contains(&0));
1080 assert!(indices.contains(&2));
1081 }
1082 other => panic!("expected ApplyFailed with 2 failures, got {other:?}"),
1083 }
1084 assert_eq!(c.value(), 0);
1085 }
1086
1087 #[tokio::test]
1088 async fn applier_timeout_fails_the_set() {
1089 let c: Control<u32> = ControlBuilder::new("c", 0)
1090 .apply_timeout(Duration::from_millis(50))
1091 .build();
1092
1093 struct SlowApplier;
1094 impl ControlApplier<u32> for SlowApplier {
1095 fn apply(
1096 &self,
1097 _v: u32,
1098 ) -> Pin<Box<dyn Future<Output = Result<(), String>> + Send + '_>> {
1099 Box::pin(async {
1100 tokio::time::sleep(Duration::from_secs(10)).await;
1101 Ok(())
1102 })
1103 }
1104 }
1105 c.register_applier(SlowApplier);
1106
1107 match c.set(1, ControlOrigin::Test).await {
1108 Err(SetError::ApplyFailed(failures)) => {
1109 assert_eq!(failures.len(), 1);
1110 assert!(
1111 failures[0].message.contains("timed out"),
1112 "message = {}",
1113 failures[0].message,
1114 );
1115 }
1116 other => panic!("expected timeout as ApplyFailed, got {other:?}"),
1117 }
1118 assert_eq!(c.value(), 0);
1119 }
1120
1121 #[tokio::test]
1124 async fn concurrent_writers_serialize_through_write_lock() {
1125 let c = Arc::new(build_u32("c", 0));
1130 let apply_count = Arc::new(AtomicUsize::new(0));
1131 let ac = apply_count.clone();
1132 c.register_applier(SyncApplier::new(move |_: u32| {
1133 ac.fetch_add(1, Ordering::SeqCst);
1134 std::thread::sleep(Duration::from_millis(5));
1137 Ok(())
1138 }));
1139
1140 let c1 = c.clone();
1141 let c2 = c.clone();
1142 let h1 = tokio::spawn(async move { c1.set(10, ControlOrigin::Test).await });
1143 let h2 = tokio::spawn(async move { c2.set(20, ControlOrigin::Test).await });
1144
1145 let r1 = h1.await.unwrap().unwrap();
1146 let r2 = h2.await.unwrap().unwrap();
1147 assert_ne!(r1, r2, "revs must be distinct across concurrent writes");
1148 assert_eq!(apply_count.load(Ordering::SeqCst), 2);
1149 let final_v = c.value();
1151 assert!(final_v == 10 || final_v == 20);
1152 }
1153
1154 #[test]
1157 fn registry_declare_and_typed_lookup() {
1158 let reg = ControlRegistry::new();
1159 let c = build_u32("concurrency", 16);
1160 reg.declare(c);
1161 let looked_up: Control<u32> = reg.get("concurrency").unwrap();
1162 assert_eq!(looked_up.value(), 16);
1163 assert_eq!(reg.len(), 1);
1164 }
1165
1166 #[test]
1167 fn registry_get_with_wrong_type_returns_none() {
1168 let reg = ControlRegistry::new();
1169 reg.declare(build_u32("concurrency", 8));
1170 let wrong: Option<Control<u64>> = reg.get("concurrency");
1172 assert!(wrong.is_none());
1173 }
1174
1175 #[test]
1176 fn registry_missing_control_returns_none() {
1177 let reg = ControlRegistry::new();
1178 assert!(reg.get::<u32>("nope").is_none());
1179 assert!(reg.get_erased("nope").is_none());
1180 }
1181
1182 #[test]
1183 #[should_panic(expected = "duplicate control declaration")]
1184 fn registry_duplicate_declaration_panics() {
1185 let reg = ControlRegistry::new();
1186 reg.declare(build_u32("c", 1));
1187 reg.declare(build_u32("c", 2));
1188 }
1189
1190 #[tokio::test]
1191 async fn erased_view_surfaces_name_rev_and_value_string() {
1192 let reg = ControlRegistry::new();
1193 let c = build_u32("concurrency", 8);
1194 reg.declare(c.clone());
1195 c.set(16, ControlOrigin::Test).await.unwrap();
1196
1197 let erased = reg.get_erased("concurrency").unwrap();
1198 assert_eq!(erased.name(), "concurrency");
1199 assert!(erased.rev() >= 1);
1200 assert_eq!(erased.value_string(), "16");
1201 assert_eq!(erased.origin(), ControlOrigin::Test);
1202 assert_eq!(erased.applier_count(), 0);
1203 }
1204
1205 #[test]
1206 fn registry_list_returns_all_controls() {
1207 let reg = ControlRegistry::new();
1208 reg.declare(build_u32("a", 1));
1209 reg.declare(build_u32("b", 2));
1210 let listed = reg.list();
1211 assert_eq!(listed.len(), 2);
1212 let mut names: Vec<&str> = listed.iter().map(|c| c.name()).collect();
1213 names.sort();
1214 assert_eq!(names, vec!["a", "b"]);
1215 }
1216
1217 #[tokio::test]
1222 async fn reified_gauge_seeds_with_initial_value() {
1223 let c: Control<u32> = ControlBuilder::new("concurrency", 32u32)
1224 .reify_as_gauge(|v| Some(*v as f64))
1225 .build();
1226 assert_eq!(c.gauge_f64(), Some(32.0));
1230 }
1231
1232 #[tokio::test]
1233 async fn reified_gauge_updates_on_commit() {
1234 let c: Control<u32> = ControlBuilder::new("concurrency", 8u32)
1235 .reify_as_gauge(|v| Some(*v as f64))
1236 .build();
1237 c.set(64, ControlOrigin::Test).await.unwrap();
1238 assert_eq!(c.gauge_f64(), Some(64.0));
1239 }
1240
1241 #[tokio::test]
1242 async fn reified_gauge_unchanged_on_apply_failure() {
1243 let c: Control<u32> = ControlBuilder::new("concurrency", 8u32)
1244 .reify_as_gauge(|v| Some(*v as f64))
1245 .build();
1246 c.register_applier(SyncApplier::new(|_: u32| Err("no".into())));
1247 let _ = c.set(999, ControlOrigin::Test).await;
1248 assert_eq!(c.gauge_f64(), Some(8.0));
1250 }
1251
1252 #[tokio::test]
1253 async fn unreified_control_has_no_gauge() {
1254 let c = build_u32("c", 1);
1255 assert!(c.gauge_f64().is_none());
1256 assert!(c.reified_gauge().is_none());
1257 }
1258
1259 #[tokio::test]
1260 async fn to_f64_can_suppress_samples() {
1261 #[derive(Clone, Debug, PartialEq)]
1265 enum Mode {
1266 Off,
1267 On(u32),
1268 }
1269 let c: Control<Mode> = ControlBuilder::new("explain", Mode::Off)
1270 .reify_as_gauge(|m| match m {
1271 Mode::Off => None,
1272 Mode::On(n) => Some(*n as f64),
1273 })
1274 .build();
1275 assert!(c.gauge_f64().is_none());
1276 c.set(Mode::On(7), ControlOrigin::Test).await.unwrap();
1277 assert_eq!(c.gauge_f64(), Some(7.0));
1278 c.set(Mode::Off, ControlOrigin::Test).await.unwrap();
1279 assert!(c.gauge_f64().is_none());
1280 }
1281
1282 #[tokio::test]
1283 async fn registry_snapshot_emits_numeric_gauge_for_reified_controls() {
1284 let reg = ControlRegistry::new();
1285 reg.declare(
1286 ControlBuilder::new("concurrency", 4u32)
1287 .reify_as_gauge(|v| Some(*v as f64))
1288 .build(),
1289 );
1290 reg.declare(ControlBuilder::new("non_reified", 99u32).build());
1291 let base = crate::labels::Labels::of("phase", "rampup");
1292 let now = std::time::Instant::now();
1293 let snap = reg.snapshot_gauges(&base, now);
1294
1295 let family = snap
1297 .family("control_concurrency")
1298 .expect("reified control should produce a numeric gauge family");
1299 let metric = family.metrics().next().unwrap();
1300 assert_eq!(metric.labels().get("phase"), Some("rampup"));
1301 assert_eq!(metric.labels().get("control"), Some("concurrency"));
1302
1303 assert!(snap.family("control_non_reified").is_none());
1306 }
1307
1308 #[tokio::test]
1309 async fn registry_snapshot_emits_info_family_for_every_control() {
1310 let reg = ControlRegistry::new();
1317 reg.declare(
1318 ControlBuilder::new("concurrency", 4u32)
1319 .reify_as_gauge(|v| Some(*v as f64))
1320 .build(),
1321 );
1322 reg.declare(ControlBuilder::new("enabled", true).build());
1323 reg.declare(ControlBuilder::new("errors_policy", "retry".to_string()).build());
1324
1325 let base = crate::labels::Labels::of("phase", "bulk");
1326 let now = std::time::Instant::now();
1327 let snap = reg.snapshot_gauges(&base, now);
1328
1329 assert!(snap.family("control_concurrency").is_some());
1331 let info = snap
1332 .family("control_info_concurrency")
1333 .expect("info family should accompany the numeric gauge");
1334 let m = info.metrics().next().unwrap();
1335 assert_eq!(m.labels().get("value"), Some("4"));
1336 assert_eq!(m.labels().get("control"), Some("concurrency"));
1337
1338 assert!(snap.family("control_enabled").is_none());
1340 let info = snap
1341 .family("control_info_enabled")
1342 .expect("bool control should emit info family");
1343 let m = info.metrics().next().unwrap();
1344 assert_eq!(m.labels().get("value"), Some("true"));
1345
1346 let info = snap
1351 .family("control_info_errors_policy")
1352 .expect("string control should emit info family");
1353 let m = info.metrics().next().unwrap();
1354 assert_eq!(m.labels().get("value"), Some("\"retry\""));
1355 }
1356
1357 #[tokio::test]
1358 async fn works_with_non_copy_value_type() {
1359 #[derive(Clone, Debug, PartialEq, Eq)]
1362 struct Targets {
1363 hosts: Vec<String>,
1364 }
1365
1366 let c: Control<Targets> = ControlBuilder::new(
1367 "targets",
1368 Targets {
1369 hosts: vec!["a".into()],
1370 },
1371 )
1372 .build();
1373
1374 let seen: Arc<Mutex<Option<Targets>>> = Arc::new(Mutex::new(None));
1375 let seen_c = seen.clone();
1376 c.register_applier(SyncApplier::new(move |t: Targets| {
1377 *seen_c.lock().unwrap() = Some(t);
1378 Ok(())
1379 }));
1380
1381 let new_val = Targets {
1382 hosts: vec!["a".into(), "b".into()],
1383 };
1384 c.set(new_val.clone(), ControlOrigin::Test).await.unwrap();
1385 assert_eq!(c.value(), new_val);
1386 assert_eq!(*seen.lock().unwrap(), Some(new_val));
1387 }
1388
1389 #[tokio::test]
1392 async fn final_control_rejects_non_launch_writes() {
1393 let c: Control<u32> = ControlBuilder::new("concurrency", 1u32)
1394 .final_at_scope("ddl_phase")
1395 .build();
1396 assert!(c.is_final());
1397 assert_eq!(c.final_scope(), Some("ddl_phase"));
1398
1399 for origin in [ControlOrigin::Test, ControlOrigin::Tui, ControlOrigin::Cli] {
1401 match c.set(42, origin.clone()).await {
1402 Err(SetError::FinalViolation { scope }) => {
1403 assert_eq!(scope, "ddl_phase");
1404 }
1405 other => panic!("expected FinalViolation for {origin:?}, got {other:?}"),
1406 }
1407 }
1408 assert_eq!(c.value(), 1u32);
1410 assert_eq!(c.get().rev, 0);
1411 }
1412
1413 #[tokio::test]
1414 async fn final_control_accepts_launch_seed() {
1415 let c: Control<u32> = ControlBuilder::new("concurrency", 4u32)
1420 .final_at_scope("ddl_phase")
1421 .build();
1422 let rev = c.set(1u32, ControlOrigin::Launch).await.unwrap();
1423 assert!(rev >= 1);
1424 assert_eq!(c.value(), 1u32);
1425 }
1426
1427 #[tokio::test]
1428 async fn final_rejection_runs_before_validation() {
1429 let c: Control<u32> = ControlBuilder::new("concurrency", 1u32)
1434 .final_at_scope("ddl_phase")
1435 .validator(|v| {
1436 if *v == 99 {
1437 Err("no 99s".into())
1438 } else {
1439 Ok(())
1440 }
1441 })
1442 .build();
1443 match c.set(99, ControlOrigin::Test).await {
1444 Err(SetError::FinalViolation { scope }) => {
1445 assert_eq!(scope, "ddl_phase");
1446 }
1447 other => panic!("expected FinalViolation, got {other:?}"),
1448 }
1449 }
1450
1451 #[tokio::test]
1452 async fn nonfinal_control_is_default() {
1453 let c = build_u32("concurrency", 1);
1454 assert!(!c.is_final());
1455 assert!(c.final_scope().is_none());
1456 }
1457
1458 #[test]
1461 fn default_branch_scope_is_local() {
1462 let c = build_u32("c", 1);
1463 assert_eq!(c.branch_scope(), BranchScope::Local);
1464 }
1465
1466 #[test]
1467 fn branch_scope_subtree_is_preserved_in_erased() {
1468 let c: Control<u32> = ControlBuilder::new("hdr_sigdigs", 3u32)
1469 .branch_scope(BranchScope::Subtree)
1470 .build();
1471 assert_eq!(c.branch_scope(), BranchScope::Subtree);
1472
1473 let reg = ControlRegistry::new();
1474 reg.declare(c);
1475 let erased = reg.get_erased("hdr_sigdigs").unwrap();
1476 assert_eq!(erased.branch_scope(), BranchScope::Subtree);
1477 }
1478
1479 #[tokio::test]
1482 async fn set_f64_rejects_without_converter() {
1483 let c: Control<u32> = ControlBuilder::new("concurrency", 4u32).build();
1484 let reg = ControlRegistry::new();
1485 reg.declare(c);
1486 let erased = reg.get_erased("concurrency").unwrap();
1487 assert!(!erased.accepts_f64_writes());
1488 match erased.set_f64(8.0, ControlOrigin::Test).await {
1489 Err(SetError::ValidationFailed(msg)) => {
1490 assert!(msg.contains("no f64 setter"), "got: {msg}");
1491 }
1492 other => panic!("expected ValidationFailed, got {other:?}"),
1493 }
1494 }
1495
1496 #[tokio::test]
1497 async fn set_f64_writes_through_converter_to_typed_control() {
1498 let c: Control<u32> = ControlBuilder::new("concurrency", 4u32)
1499 .from_f64(|v| {
1500 if v < 0.0 || v > u32::MAX as f64 {
1501 Err(format!("concurrency out of range: {v}"))
1502 } else {
1503 Ok(v as u32)
1504 }
1505 })
1506 .build();
1507 let reg = ControlRegistry::new();
1508 reg.declare(c.clone());
1509 let erased = reg.get_erased("concurrency").unwrap();
1510 assert!(erased.accepts_f64_writes());
1511
1512 let rev = erased.set_f64(64.0, ControlOrigin::Test).await.unwrap();
1513 assert!(rev >= 1);
1514 assert_eq!(c.value(), 64u32);
1515 }
1516
1517 #[tokio::test]
1518 async fn set_f64_converter_error_is_surfaced() {
1519 let c: Control<u32> = ControlBuilder::new("concurrency", 4u32)
1520 .from_f64(|v| {
1521 if v < 0.0 {
1522 Err(format!("negative: {v}"))
1523 } else {
1524 Ok(v as u32)
1525 }
1526 })
1527 .build();
1528 let reg = ControlRegistry::new();
1529 reg.declare(c.clone());
1530 let erased = reg.get_erased("concurrency").unwrap();
1531 match erased.set_f64(-1.0, ControlOrigin::Test).await {
1532 Err(SetError::ValidationFailed(msg)) => {
1533 assert!(msg.contains("negative"), "got: {msg}");
1534 }
1535 other => panic!("expected ValidationFailed, got {other:?}"),
1536 }
1537 assert_eq!(c.value(), 4u32);
1539 }
1540
1541 #[tokio::test]
1542 async fn set_f64_respects_final_scope() {
1543 let c: Control<u32> = ControlBuilder::new("concurrency", 4u32)
1544 .from_f64(|v| Ok(v as u32))
1545 .final_at_scope("ddl_phase")
1546 .build();
1547 let reg = ControlRegistry::new();
1548 reg.declare(c.clone());
1549 let erased = reg.get_erased("concurrency").unwrap();
1550 match erased.set_f64(8.0, ControlOrigin::Test).await {
1551 Err(SetError::FinalViolation { scope }) => {
1552 assert_eq!(scope, "ddl_phase");
1553 }
1554 other => panic!("expected FinalViolation, got {other:?}"),
1555 }
1556 }
1557}