1use degenbot_core::diag;
15use degenbot_core::{op_info, op_warn};
16use std::collections::VecDeque;
17use std::sync::atomic::{AtomicU64, Ordering};
18use std::sync::{Arc, OnceLock};
19
20use degenbot_config::FleetConfig;
21use parking_lot::{Mutex, RwLock};
22
23use crate::role::CordonClass;
24
25#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
28pub enum FleetPosture {
29 Nominal,
31 Cordoned,
35}
36
37#[derive(Debug, Clone, Copy, PartialEq)]
39pub enum EnterReason {
40 EventBurst {
42 events: u64,
44 },
45 DutySpike {
47 duty_percent: f64,
49 },
50 LaneDeath,
54}
55
56#[derive(Debug, Clone, Copy, PartialEq, Eq)]
61pub enum PostureCause {
62 LaneDeath,
70}
71
72#[derive(Debug, Clone, Copy, PartialEq)]
74pub enum PostureChange {
75 Held,
77 Entered(EnterReason),
79 Exited,
81}
82
83#[derive(Debug, Clone, Copy, PartialEq)]
86pub struct PosturePolicy {
87 pub enter_events: usize,
89 pub enter_window_ms: u64,
91 pub duty_percent: f64,
94 pub duty_window_ms: u64,
96 pub exit_clean_ms: u64,
99 pub sim_intake_floor_override: Option<usize>,
102}
103
104impl PosturePolicy {
105 #[must_use]
108 pub const fn doc_defaults() -> Self {
109 Self {
110 enter_events: 2,
111 enter_window_ms: 1_000,
112 duty_percent: 2.0,
113 duty_window_ms: 5_000,
114 exit_clean_ms: 10_000,
115 sim_intake_floor_override: None,
116 }
117 }
118
119 #[must_use]
121 pub fn from_config(cfg: &FleetConfig) -> Self {
122 Self {
123 enter_events: cfg.cordon_enter_events,
124 enter_window_ms: cfg.cordon_enter_window_ms,
125 duty_percent: cfg.cordon_duty_percent,
126 duty_window_ms: cfg.cordon_duty_window_ms,
127 exit_clean_ms: cfg.cordon_exit_clean_ms,
128 sim_intake_floor_override: cfg.cordon_sim_intake_floor,
129 }
130 }
131
132 #[must_use]
135 pub fn sim_intake_cap(&self, slot_cap: usize) -> usize {
136 self.sim_intake_floor_override
137 .unwrap_or(slot_cap / 2)
138 .min(slot_cap)
139 .max(1)
140 }
141
142 #[must_use]
148 pub fn patched_with(self, patch: PosturePolicyPatch) -> Self {
149 Self {
150 enter_events: patch.enter_events.unwrap_or(self.enter_events),
151 enter_window_ms: patch.enter_window_ms.unwrap_or(self.enter_window_ms),
152 duty_percent: patch.duty_percent.unwrap_or(self.duty_percent),
153 duty_window_ms: patch.duty_window_ms.unwrap_or(self.duty_window_ms),
154 exit_clean_ms: patch.exit_clean_ms.unwrap_or(self.exit_clean_ms),
155 sim_intake_floor_override: match patch.sim_intake_floor_override {
156 None => self.sim_intake_floor_override,
158 Some(floor) => floor,
161 },
162 }
163 }
164}
165
166#[derive(Debug, Clone, Copy, Default, PartialEq)]
178pub struct PosturePolicyPatch {
179 pub enter_events: Option<usize>,
181 pub enter_window_ms: Option<u64>,
183 pub duty_percent: Option<f64>,
185 pub duty_window_ms: Option<u64>,
187 pub exit_clean_ms: Option<u64>,
189 pub sim_intake_floor_override: Option<Option<usize>>,
192}
193
194impl PosturePolicyPatch {
195 #[must_use]
199 pub fn is_empty(&self) -> bool {
200 self.enter_events.is_none()
201 && self.enter_window_ms.is_none()
202 && self.duty_percent.is_none()
203 && self.duty_window_ms.is_none()
204 && self.exit_clean_ms.is_none()
205 && self.sim_intake_floor_override.is_none()
206 }
207
208 pub fn validate(&self) -> Result<(), PostureRetuneError> {
217 if self.is_empty() {
218 return Err(PostureRetuneError::EmptyPatch);
219 }
220 if self.enter_events.is_some_and(|v| v < 1) {
221 return Err(PostureRetuneError::EnterEvents(
222 self.enter_events.unwrap_or_default(),
223 ));
224 }
225 if self.enter_window_ms.is_some_and(|v| v == 0) {
226 return Err(PostureRetuneError::EnterWindow(0));
227 }
228 if let Some(duty) = self.duty_percent {
229 if !(duty > 0.0 && duty <= 100.0) {
233 return Err(PostureRetuneError::DutyPercent(duty));
234 }
235 }
236 if self.duty_window_ms.is_some_and(|v| v == 0) {
237 return Err(PostureRetuneError::DutyWindow(0));
238 }
239 if self.exit_clean_ms.is_some_and(|v| v == 0) {
240 return Err(PostureRetuneError::ExitClean(0));
241 }
242 if let Some(Some(floor)) = self.sim_intake_floor_override {
243 if floor < 1 {
244 return Err(PostureRetuneError::SimIntakeFloor(floor));
245 }
246 }
247 Ok(())
248 }
249}
250
251#[derive(Debug, Clone, Copy, PartialEq, thiserror::Error)]
255pub enum PostureRetuneError {
256 #[error("at least one cordon threshold key is required (empty patch)")]
258 EmptyPatch,
259 #[error("cordon_enter_events must be >= 1, got {0}")]
261 EnterEvents(usize),
262 #[error("cordon_enter_window_ms must be > 0 ms, got {0} ms")]
264 EnterWindow(u64),
265 #[error("cordon_duty_percent must be in (0.0, 100.0], got {0}")]
267 DutyPercent(f64),
268 #[error("cordon_duty_window_ms must be > 0 ms, got {0} ms")]
270 DutyWindow(u64),
271 #[error("cordon_exit_clean_ms must be > 0 ms, got {0} ms")]
273 ExitClean(u64),
274 #[error("cordon_sim_intake_floor must be >= 1, got {0}")]
276 SimIntakeFloor(usize),
277}
278
279#[derive(Debug, Clone, Copy, PartialEq, Eq)]
282pub struct ThrottleSample {
283 pub events: u64,
285 pub throttled_usec: u64,
287 pub elapsed_usec: u64,
289}
290
291impl ThrottleSample {
292 #[must_use]
294 pub const fn is_clean(&self) -> bool {
295 self.events == 0 && self.throttled_usec == 0
296 }
297}
298
299#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
302pub struct PostureCounters {
303 pub entered: u64,
305 pub exited: u64,
307 pub intake_suppressed: u64,
310 pub lane_deaths: u64,
313}
314
315#[derive(Debug, Clone, Copy)]
316struct Sample {
317 now_ms: u64,
318 sample: ThrottleSample,
319}
320
321#[derive(Debug)]
324pub struct PostureStateMachine {
325 state: FleetPosture,
326 policy: PosturePolicy,
327 samples: VecDeque<Sample>,
328 last_unclean_ms: Option<u64>,
329 counters: PostureCounters,
330 lane_death_hold: bool,
337}
338
339impl PostureStateMachine {
340 #[must_use]
342 pub fn new(policy: PosturePolicy) -> Self {
343 Self {
344 state: FleetPosture::Nominal,
345 policy,
346 samples: VecDeque::new(),
347 last_unclean_ms: None,
348 counters: PostureCounters::default(),
349 lane_death_hold: false,
350 }
351 }
352
353 #[must_use]
355 pub const fn state(&self) -> FleetPosture {
356 self.state
357 }
358
359 #[must_use]
361 pub const fn counters(&self) -> &PostureCounters {
362 &self.counters
363 }
364
365 #[must_use]
369 pub const fn lane_death_held(&self) -> bool {
370 self.lane_death_hold
371 }
372
373 #[must_use]
375 pub const fn policy(&self) -> &PosturePolicy {
376 &self.policy
377 }
378
379 pub fn set_policy(&mut self, policy: PosturePolicy) {
385 self.policy = policy;
386 }
387
388 pub fn observe(&mut self, now_ms: u64, sample: ThrottleSample) -> PostureChange {
391 self.samples.push_back(Sample { now_ms, sample });
392 self.prune(now_ms);
393 if !sample.is_clean() {
394 self.last_unclean_ms = Some(now_ms);
395 }
396
397 match self.state {
398 FleetPosture::Nominal => self.maybe_enter(now_ms),
399 FleetPosture::Cordoned => self.maybe_exit(now_ms),
400 }
401 }
402
403 pub fn observe_cause(&mut self, cause: PostureCause) -> PostureChange {
412 match cause {
413 PostureCause::LaneDeath => {
414 self.counters.lane_deaths += 1;
415 if self.lane_death_hold {
416 return PostureChange::Held;
417 }
418 self.lane_death_hold = true;
419 if self.state == FleetPosture::Cordoned {
420 op_warn!(domain = pump, lane_deaths = self.counters.lane_deaths,
424 "lane-death HOLD upgrades an existing cordon — sticky, clean-window exit disabled"
425 );
426 return PostureChange::Held;
427 }
428 self.enter(EnterReason::LaneDeath)
429 }
430 }
431 }
432
433 pub const fn note_intake_suppressed(&mut self) {
436 self.counters.intake_suppressed += 1;
437 }
438
439 #[must_use]
443 pub const fn admits_lease(&self, class: CordonClass) -> bool {
444 !matches!(
445 (self.state, class),
446 (FleetPosture::Cordoned, CordonClass::Deferrable)
447 )
448 }
449
450 #[must_use]
452 pub fn sim_intake_cap(&self, slot_cap: usize) -> usize {
453 match self.state {
454 FleetPosture::Nominal => slot_cap,
455 FleetPosture::Cordoned => self.policy.sim_intake_cap(slot_cap),
456 }
457 }
458
459 fn prune(&mut self, now_ms: u64) {
460 let window = self.policy.duty_window_ms.max(self.policy.enter_window_ms);
461 while let Some(front) = self.samples.front() {
462 if now_ms.saturating_sub(front.now_ms) > window {
463 self.samples.pop_front();
464 } else {
465 break;
466 }
467 }
468 }
469
470 fn maybe_enter(&mut self, now_ms: u64) -> PostureChange {
471 let burst: u64 = self
473 .samples
474 .iter()
475 .filter(|s| now_ms.saturating_sub(s.now_ms) <= self.policy.enter_window_ms)
476 .map(|s| s.sample.events)
477 .sum();
478 if burst >= u64::try_from(self.policy.enter_events.max(1)).unwrap_or(1) {
479 return self.enter(EnterReason::EventBurst { events: burst });
480 }
481 let (events, throttled_usec, elapsed_usec) = self.duty_window_totals();
483 let _ = events;
484 if elapsed_usec > 0 {
485 #[expect(
486 clippy::cast_precision_loss,
487 reason = "duty percent is an f64 metric by definition (µs/µs ratio)"
488 )]
489 let duty_percent = throttled_usec as f64 / elapsed_usec as f64 * 100.0;
490 if duty_percent > self.policy.duty_percent {
491 return self.enter(EnterReason::DutySpike { duty_percent });
492 }
493 }
494 PostureChange::Held
495 }
496
497 fn duty_window_totals(&self) -> (u64, u64, u64) {
498 self.samples
499 .iter()
500 .fold((0, 0, 0), |(ev, th, el), Sample { sample, .. }| {
501 (
502 ev + sample.events,
503 th + sample.throttled_usec,
504 el + sample.elapsed_usec,
505 )
506 })
507 }
508
509 fn enter(&mut self, reason: EnterReason) -> PostureChange {
510 self.state = FleetPosture::Cordoned;
511 self.counters.entered += 1;
512 op_warn!(domain = pump, reason = ?reason,
513 entered = self.counters.entered,
514 "cordon ENTER — deferrable intake held, sim intake floored; in-flight units complete"
515 );
516 PostureChange::Entered(reason)
517 }
518
519 fn maybe_exit(&mut self, now_ms: u64) -> PostureChange {
520 if self.lane_death_hold {
525 return PostureChange::Held;
526 }
527 let dirty_recently = self
530 .last_unclean_ms
531 .is_some_and(|last| now_ms.saturating_sub(last) < self.policy.exit_clean_ms);
532 if dirty_recently {
533 return PostureChange::Held;
534 }
535 self.state = FleetPosture::Nominal;
536 self.counters.exited += 1;
537 op_info!(
538 domain = pump,
539 exited = self.counters.exited,
540 clean_ms = self.policy.exit_clean_ms,
541 "cordon EXIT after clean-window hysteresis"
542 );
543 PostureChange::Exited
544 }
545}
546
547#[derive(Debug)]
558pub struct PostureWatch {
559 shared: Arc<FeedShared>,
560 seen: AtomicU64,
561}
562
563impl PostureWatch {
564 #[must_use]
566 pub fn current(&self) -> FleetPosture {
567 self.shared.state.read().posture
568 }
569
570 #[must_use]
572 pub fn has_changed(&self) -> bool {
573 self.shared.state.read().seq != self.seen.load(Ordering::Relaxed)
574 }
575
576 #[must_use]
579 pub fn take_if_changed(&self) -> Option<FleetPosture> {
580 let state = self.shared.state.read();
581 if state.seq == self.seen.load(Ordering::Relaxed) {
582 return None;
583 }
584 self.seen.store(state.seq, Ordering::Relaxed);
585 Some(state.posture)
586 }
587}
588
589#[derive(Debug)]
592struct FeedShared {
593 state: RwLock<FeedState>,
594}
595
596#[derive(Debug, Clone, Copy)]
597struct FeedState {
598 posture: FleetPosture,
599 seq: u64,
600}
601
602#[derive(Debug)]
613pub struct PostureOwner {
614 machine: Mutex<PostureStateMachine>,
615 broadcast: Arc<FeedShared>,
616}
617
618impl PostureOwner {
619 #[must_use]
624 pub fn new(policy: PosturePolicy) -> Self {
625 Self {
626 machine: Mutex::new(PostureStateMachine::new(policy)),
627 broadcast: Arc::new(FeedShared {
628 state: RwLock::new(FeedState {
629 posture: FleetPosture::Nominal,
630 seq: 0,
631 }),
632 }),
633 }
634 }
635
636 pub fn observe_throttle(&self, now_ms: u64, sample: ThrottleSample) -> PostureChange {
649 let (change, posture) = {
650 let mut machine = self.machine.lock();
651 let change = machine.observe(now_ms, sample);
652 (change, machine.state())
653 };
654 if !matches!(change, PostureChange::Held) {
655 self.publish(posture);
656 }
657 change
658 }
659
660 pub fn observe_cause(&self, cause: PostureCause) -> PostureChange {
671 let (change, posture) = {
672 let mut machine = self.machine.lock();
673 let change = machine.observe_cause(cause);
674 (change, machine.state())
675 };
676 if !matches!(change, PostureChange::Held) {
677 self.publish(posture);
678 }
679 change
680 }
681
682 #[must_use]
684 pub fn current(&self) -> FleetPosture {
685 self.machine.lock().state()
686 }
687
688 #[must_use]
691 pub fn policy(&self) -> PosturePolicy {
692 *self.machine.lock().policy()
693 }
694
695 #[must_use]
698 pub fn subscribe(&self) -> PostureWatch {
699 let seen = self.broadcast.state.read().seq;
700 PostureWatch {
701 shared: Arc::clone(&self.broadcast),
702 seen: AtomicU64::new(seen),
703 }
704 }
705
706 pub fn retune(&self, new_policy: PosturePolicy) {
715 let posture = {
716 let mut machine = self.machine.lock();
717 machine.set_policy(new_policy);
718 machine.state()
719 };
720 self.publish(posture);
721 }
722
723 #[must_use]
727 pub fn admits_lease(&self, class: CordonClass) -> bool {
728 self.machine.lock().admits_lease(class)
729 }
730
731 pub fn note_intake_suppressed(&self) {
734 self.machine.lock().note_intake_suppressed();
735 }
736
737 #[must_use]
739 pub fn sim_intake_cap(&self, slot_cap: usize) -> usize {
740 self.machine.lock().sim_intake_cap(slot_cap)
741 }
742
743 #[must_use]
745 pub fn counters(&self) -> PostureCounters {
746 *self.machine.lock().counters()
747 }
748
749 #[must_use]
753 pub fn lane_death_held(&self) -> bool {
754 self.machine.lock().lane_death_held()
755 }
756
757 fn publish(&self, posture: FleetPosture) {
761 let mut state = self.broadcast.state.write();
762 if state.posture != posture {
763 state.posture = posture;
764 state.seq = state.seq.wrapping_add(1);
765 }
766 }
767}
768
769static PROCESS_OWNER: OnceLock<PostureOwner> = OnceLock::new();
773
774#[must_use]
780pub fn install_process_owner(policy: PosturePolicy) -> &'static PostureOwner {
781 if PROCESS_OWNER.set(PostureOwner::new(policy)).is_err() {
782 diag!(
783 domain = pump,
784 "process owner already installed — first-wins, keeping the existing owner"
785 );
786 }
787 process()
788}
789
790#[must_use]
796pub fn process() -> &'static PostureOwner {
797 PROCESS_OWNER.get_or_init(|| PostureOwner::new(PosturePolicy::doc_defaults()))
798}
799
800#[cfg(test)]
801mod tests {
802 use super::*;
803
804 fn policy() -> PosturePolicy {
805 PosturePolicy {
806 enter_events: 2,
807 enter_window_ms: 1_000,
808 duty_percent: 2.0,
809 duty_window_ms: 5_000,
810 exit_clean_ms: 10_000,
811 sim_intake_floor_override: None,
812 }
813 }
814
815 fn sm() -> PostureStateMachine {
816 PostureStateMachine::new(policy())
817 }
818
819 fn sample(events: u64, throttled_usec: u64, elapsed_usec: u64) -> ThrottleSample {
820 ThrottleSample {
821 events,
822 throttled_usec,
823 elapsed_usec,
824 }
825 }
826
827 #[test]
828 fn nominal_is_the_boot_state() {
829 assert_eq!(sm().state(), FleetPosture::Nominal);
830 }
831
832 #[test]
835 fn a_lane_death_cordons_immediately() {
836 let mut m = sm();
837 assert_eq!(
838 m.observe_cause(PostureCause::LaneDeath),
839 PostureChange::Entered(EnterReason::LaneDeath)
840 );
841 assert_eq!(m.state(), FleetPosture::Cordoned);
842 assert_eq!(m.counters().lane_deaths, 1);
843 assert_eq!(m.counters().entered, 1);
844 }
845
846 #[test]
850 fn the_lane_death_cordon_is_sticky_across_clean_windows() {
851 let mut m = sm();
852 m.observe_cause(PostureCause::LaneDeath);
853 for t in (0..30_u64).map(|i| 10_000 + i * 1_000) {
857 assert_eq!(
858 m.observe(t, sample(0, 0, 1_000_000)),
859 PostureChange::Held,
860 "a clean window must never lift a lane-death cordon"
861 );
862 }
863 assert_eq!(m.state(), FleetPosture::Cordoned);
864 assert_eq!(m.counters().exited, 0, "the sticky cordon never exits");
865 }
866
867 #[test]
871 fn a_lane_death_upgrades_a_throttle_cordon_to_sticky() {
872 let mut m = sm();
873 m.observe(100, sample(2, 0, 1_000_000));
875 assert_eq!(m.state(), FleetPosture::Cordoned);
876 assert_eq!(m.counters().entered, 1);
877 assert_eq!(
879 m.observe_cause(PostureCause::LaneDeath),
880 PostureChange::Held,
881 "no state transition — the cordon was already up"
882 );
883 assert_eq!(m.counters().lane_deaths, 1);
884 assert_eq!(m.counters().entered, 1, "no second enter counted");
885 for t in (0..30_u64).map(|i| 10_000 + i * 1_000) {
887 m.observe(t, sample(0, 0, 1_000_000));
888 }
889 assert_eq!(m.state(), FleetPosture::Cordoned);
890 assert_eq!(m.counters().exited, 0);
891 }
892
893 #[test]
895 fn repeated_lane_deaths_count_but_do_not_re_enter() {
896 let mut m = sm();
897 m.observe_cause(PostureCause::LaneDeath);
898 for _ in 0..3 {
899 assert_eq!(
900 m.observe_cause(PostureCause::LaneDeath),
901 PostureChange::Held
902 );
903 }
904 assert_eq!(m.counters().lane_deaths, 4);
905 assert_eq!(m.counters().entered, 1);
906 }
907
908 #[test]
912 fn lane_death_cordon_effects_match_the_throttle_vocabulary() {
913 let mut m = sm();
914 m.observe_cause(PostureCause::LaneDeath);
915 assert!(!m.admits_lease(CordonClass::Deferrable));
916 assert!(m.admits_lease(CordonClass::Never));
917 assert!(m.admits_lease(CordonClass::SimPool));
918 assert_eq!(m.sim_intake_cap(4), 2);
921 }
922
923 #[test]
926 fn the_owner_publishes_the_lane_death_transition() {
927 let owner = PostureOwner::new(policy());
928 let watch = owner.subscribe();
929 assert_eq!(watch.current(), FleetPosture::Nominal);
930 owner.observe_cause(PostureCause::LaneDeath);
931 assert_eq!(watch.current(), FleetPosture::Cordoned);
932 assert_eq!(owner.current(), FleetPosture::Cordoned);
933 }
934
935 #[test]
936 fn enters_on_event_burst_within_the_window() {
937 let mut m = sm();
938 assert_eq!(m.observe(100, sample(1, 0, 1_000_000)), PostureChange::Held);
940 let change = m.observe(600, sample(1, 0, 500_000));
943 assert_eq!(
944 change,
945 PostureChange::Entered(EnterReason::EventBurst { events: 2 })
946 );
947 assert_eq!(m.state(), FleetPosture::Cordoned);
948 assert_eq!(m.counters().entered, 1);
949 assert!(matches!(
951 m.observe(700, sample(1, 0, 100_000)),
952 PostureChange::Held
953 ));
954 assert_eq!(m.counters().entered, 1);
955 }
956
957 #[test]
958 fn burst_outside_the_enter_window_does_not_cordon() {
959 let mut m = sm();
960 m.observe(0, sample(1, 0, 1_000_000));
961 m.observe(3_000, sample(1, 0, 2_000_000));
963 m.observe(3_100, sample(0, 0, 100_000));
964 assert_eq!(m.state(), FleetPosture::Nominal);
965 }
966
967 #[test]
968 fn enters_on_duty_spike_over_the_duty_window() {
969 let mut m = sm();
970 let change = m.observe(1_000, sample(0, 500_000, 1_000_000));
974 assert!(matches!(
975 change,
976 PostureChange::Entered(EnterReason::DutySpike { .. })
977 ));
978 assert_eq!(m.state(), FleetPosture::Cordoned);
979 assert_eq!(m.counters().entered, 1);
980 }
981
982 #[test]
983 fn a_rising_duty_only_crosses_once_the_trailing_total_exceeds_the_threshold() {
984 let mut m = sm();
985 for t in 1..5_u64 {
987 assert_eq!(
988 m.observe(t * 1_000, sample(0, 500, 1_000_000)),
989 PostureChange::Held
990 );
991 }
992 let change = m.observe(5_000, sample(0, 5_000_000, 1_000_000));
994 assert!(matches!(
995 change,
996 PostureChange::Entered(EnterReason::DutySpike { .. })
997 ));
998 }
999
1000 #[test]
1001 fn sub_threshold_duty_never_cordons() {
1002 let mut m = sm();
1003 for t in 0..10_u64 {
1005 m.observe((t + 1) * 1_000, sample(0, 10_000, 1_000_000));
1006 }
1007 assert_eq!(m.state(), FleetPosture::Nominal);
1008 }
1009
1010 #[test]
1011 fn exits_only_after_the_full_clean_hysteresis() {
1012 let mut m = sm();
1013 m.observe(100, sample(2, 0, 100_000)); assert_eq!(m.state(), FleetPosture::Cordoned);
1015 let mut t = 200;
1017 while t < 10_000 {
1018 assert_eq!(m.observe(t, sample(0, 0, 1_000_000)), PostureChange::Held);
1019 t += 1_000;
1020 }
1021 assert_eq!(m.state(), FleetPosture::Cordoned);
1022 let change = m.observe(10_200, sample(0, 0, 100_000));
1024 assert_eq!(change, PostureChange::Exited);
1025 assert_eq!(m.state(), FleetPosture::Nominal);
1026 assert_eq!(m.counters().exited, 1);
1027 }
1028
1029 #[test]
1030 fn a_dirty_window_restarts_the_clean_clock() {
1031 let mut m = sm();
1032 m.observe(0, sample(2, 0, 100_000));
1033 let mut t = 1_000;
1034 while t < 9_000 {
1035 m.observe(t, sample(0, 0, 1_000_000));
1036 t += 1_000;
1037 }
1038 m.observe(9_000, sample(1, 0, 100_000));
1040 assert_eq!(m.state(), FleetPosture::Cordoned);
1041 let mut t = 10_000;
1042 while t < 18_000 {
1043 m.observe(t, sample(0, 0, 1_000_000));
1044 t += 1_000;
1045 }
1046 assert_eq!(m.state(), FleetPosture::Cordoned, "only 9 s clean");
1047 m.observe(19_100, sample(0, 0, 100_000));
1048 assert_eq!(
1049 m.state(),
1050 FleetPosture::Nominal,
1051 "10 s clean since the reset"
1052 );
1053 }
1054
1055 #[test]
1056 fn cordon_effects_match_the_sign_off_table() {
1057 let mut m = sm();
1058 assert_eq!(m.sim_intake_cap(4), 4);
1060 assert!(m.admits_lease(CordonClass::Never));
1061 assert!(m.admits_lease(CordonClass::SimPool));
1062 assert!(m.admits_lease(CordonClass::Deferrable));
1063
1064 m.observe(0, sample(2, 0, 100_000));
1065 assert_eq!(m.state(), FleetPosture::Cordoned);
1066 assert!(!m.admits_lease(CordonClass::Deferrable));
1068 assert!(m.admits_lease(CordonClass::SimPool));
1069 assert!(m.admits_lease(CordonClass::Never));
1070 assert_eq!(m.sim_intake_cap(4), 2, "floor = half the slot cap");
1071 assert_eq!(m.sim_intake_cap(1), 1, "floored at one");
1072 }
1073
1074 #[test]
1075 fn override_intake_floor_wins_and_caps_at_the_slot_cap() {
1076 let p = PosturePolicy {
1077 sim_intake_floor_override: Some(7),
1078 ..policy()
1079 };
1080 assert_eq!(p.sim_intake_cap(4), 4, "nothing above the slot cap");
1081 let p = PosturePolicy {
1082 sim_intake_floor_override: Some(0),
1083 ..policy()
1084 };
1085 assert_eq!(p.sim_intake_cap(4), 1, "floored at one");
1086 }
1087
1088 #[test]
1089 fn suppressed_intake_is_counted_for_the_tuning_loop() {
1090 let mut m = sm();
1091 m.note_intake_suppressed();
1092 m.note_intake_suppressed();
1093 assert_eq!(m.counters().intake_suppressed, 2);
1094 }
1095
1096 #[test]
1097 fn typed_config_projects_all_the_q5_amendment_thresholds() {
1098 let cfg = degenbot_config::BotConfig::default();
1099 let p = PosturePolicy::from_config(&cfg.fleet);
1100 assert_eq!(p.enter_events, 2);
1102 assert_eq!(p.enter_window_ms, 1_000);
1103 assert!((p.duty_percent - 2.0).abs() < 1e-9, "duty default is 2%");
1104 assert_eq!(p.duty_window_ms, 5_000);
1105 assert_eq!(p.exit_clean_ms, 10_000);
1106 assert_eq!(p.sim_intake_floor_override, None);
1107 }
1108
1109 #[test]
1112 fn owner_transitions_publish_only_on_change() {
1113 let owner = PostureOwner::new(policy());
1114 let watch = owner.subscribe();
1115 owner.observe_throttle(0, sample(0, 0, 1_000));
1117 assert!(!watch.has_changed());
1118 assert_eq!(watch.current(), FleetPosture::Nominal);
1119 owner.observe_throttle(100, sample(1, 0, 1_000));
1121 assert!(!watch.has_changed());
1122 owner.observe_throttle(600, sample(1, 0, 500_000));
1125 assert!(watch.has_changed());
1126 assert_eq!(watch.take_if_changed(), Some(FleetPosture::Cordoned));
1127 assert_eq!(watch.take_if_changed(), None, "one edge per transition");
1128 owner.observe_throttle(700, sample(1, 0, 100_000));
1130 assert!(!watch.has_changed());
1131 assert_eq!(watch.current(), FleetPosture::Cordoned);
1132 }
1133
1134 #[test]
1135 fn owner_current_reflects_the_machine_including_exit_hysteresis() {
1136 let owner = PostureOwner::new(policy());
1137 assert_eq!(owner.current(), FleetPosture::Nominal);
1138 owner.observe_throttle(0, sample(3, 0, 1_000));
1139 assert_eq!(owner.current(), FleetPosture::Cordoned);
1140 assert_eq!(owner.counters().entered, 1);
1141 let mut now = 1_000;
1143 loop {
1144 owner.observe_throttle(now, sample(0, 0, 1_000));
1145 if owner.current() == FleetPosture::Nominal {
1146 break;
1147 }
1148 now += 1_000;
1149 assert!(now <= 60_000, "the cordon never lifted");
1150 }
1151 assert_eq!(owner.counters().exited, 1);
1152 }
1153
1154 #[test]
1155 fn first_wins_process_install_keeps_the_existing_owner() {
1156 let first = install_process_owner(policy());
1157 let second = install_process_owner(PosturePolicy {
1158 enter_events: 99,
1159 ..policy()
1160 });
1161 assert!(
1162 std::ptr::eq(first, second),
1163 "first-wins: a later install returns the existing owner"
1164 );
1165 assert!(std::ptr::eq(first, process()));
1166 assert_eq!(
1167 second.policy().enter_events,
1168 first.policy().enter_events,
1169 "the losing install's policy never landed"
1170 );
1171 }
1172
1173 #[test]
1174 fn retune_swaps_thresholds_and_the_next_observe_rederives() {
1175 let owner = PostureOwner::new(policy());
1176 let watch = owner.subscribe();
1177 owner.observe_throttle(0, sample(1, 0, 1_000_000));
1179 assert_eq!(owner.current(), FleetPosture::Nominal);
1180 owner.retune(PosturePolicy {
1183 enter_events: 1,
1184 ..policy()
1185 });
1186 assert_eq!(owner.policy().enter_events, 1, "the retune swapped");
1187 assert!(!watch.has_changed(), "a no-effect retune is not an edge");
1188 owner.observe_throttle(2_000, sample(1, 0, 100_000));
1189 assert_eq!(owner.current(), FleetPosture::Cordoned);
1190 assert_eq!(watch.take_if_changed(), Some(FleetPosture::Cordoned));
1191 }
1192
1193 #[test]
1194 fn two_owners_are_fully_independent_hermetic_isolation() {
1195 let a = PostureOwner::new(policy());
1196 let b = PostureOwner::new(policy());
1197 let watch_a = a.subscribe();
1198 let watch_b = b.subscribe();
1199 a.observe_throttle(0, sample(3, 0, 1_000));
1201 assert_eq!(a.current(), FleetPosture::Cordoned);
1202 assert_eq!(b.current(), FleetPosture::Nominal);
1203 assert_eq!(watch_a.take_if_changed(), Some(FleetPosture::Cordoned));
1204 assert_eq!(
1205 watch_b.take_if_changed(),
1206 None,
1207 "no posture leaks across owners"
1208 );
1209 b.observe_throttle(1_000, sample(0, 0, 1_000));
1212 assert_eq!(watch_b.take_if_changed(), None);
1213 assert_eq!(watch_a.take_if_changed(), None);
1214 }
1215
1216 #[test]
1217 fn owner_read_throughs_match_the_machine_semantics() {
1218 let owner = PostureOwner::new(PosturePolicy {
1219 sim_intake_floor_override: Some(1),
1220 ..policy()
1221 });
1222 assert!(owner.admits_lease(CordonClass::Deferrable));
1223 assert_eq!(owner.sim_intake_cap(8), 8, "nominal intake is the cap");
1224 owner.observe_throttle(0, sample(3, 0, 1_000));
1225 assert!(!owner.admits_lease(CordonClass::Deferrable));
1226 assert!(owner.admits_lease(CordonClass::Never));
1227 assert!(owner.admits_lease(CordonClass::SimPool));
1228 assert_eq!(owner.sim_intake_cap(8), 1, "the cordon floor override");
1229 assert_eq!(owner.counters().entered, 1);
1230 }
1231
1232 #[test]
1235 fn an_empty_patch_is_rejected() {
1236 assert!(PosturePolicyPatch::default().validate().is_err());
1237 assert!(PosturePolicyPatch::default().is_empty());
1238 assert_eq!(
1239 PosturePolicyPatch::default().validate(),
1240 Err(PostureRetuneError::EmptyPatch),
1241 "the empty-patch refusal is its own typed error"
1242 );
1243 }
1244
1245 #[test]
1246 fn every_threshold_rule_is_a_typed_rejection() {
1247 assert_eq!(
1248 PosturePolicyPatch {
1249 enter_events: Some(0),
1250 ..PosturePolicyPatch::default()
1251 }
1252 .validate(),
1253 Err(PostureRetuneError::EnterEvents(0)),
1254 "enter_events >= 1"
1255 );
1256 assert_eq!(
1257 PosturePolicyPatch {
1258 enter_window_ms: Some(0),
1259 ..PosturePolicyPatch::default()
1260 }
1261 .validate(),
1262 Err(PostureRetuneError::EnterWindow(0)),
1263 "windows must be > 0 ms"
1264 );
1265 assert_eq!(
1266 PosturePolicyPatch {
1267 duty_window_ms: Some(0),
1268 ..PosturePolicyPatch::default()
1269 }
1270 .validate(),
1271 Err(PostureRetuneError::DutyWindow(0))
1272 );
1273 assert_eq!(
1274 PosturePolicyPatch {
1275 exit_clean_ms: Some(0),
1276 ..PosturePolicyPatch::default()
1277 }
1278 .validate(),
1279 Err(PostureRetuneError::ExitClean(0))
1280 );
1281 for duty in [0.0, -1.0, 100.5, f64::NAN, f64::INFINITY] {
1284 let rejected = PosturePolicyPatch {
1285 duty_percent: Some(duty),
1286 ..PosturePolicyPatch::default()
1287 }
1288 .validate();
1289 assert!(
1290 matches!(rejected, Err(PostureRetuneError::DutyPercent(_))),
1291 "duty {duty} must be refused as DutyPercent, got {rejected:?}"
1292 );
1293 }
1294 assert_eq!(
1295 PosturePolicyPatch {
1296 sim_intake_floor_override: Some(Some(0)),
1297 ..PosturePolicyPatch::default()
1298 }
1299 .validate(),
1300 Err(PostureRetuneError::SimIntakeFloor(0)),
1301 "an explicit floor must be >= 1"
1302 );
1303 }
1304
1305 #[test]
1306 fn boundary_values_are_admitted() {
1307 let validated = PosturePolicyPatch {
1310 enter_events: Some(1),
1311 enter_window_ms: Some(1),
1312 duty_percent: Some(100.0),
1313 duty_window_ms: Some(1),
1314 exit_clean_ms: Some(1),
1315 sim_intake_floor_override: Some(Some(1)),
1316 }
1317 .validate();
1318 assert_eq!(
1319 validated,
1320 Ok(()),
1321 "inclusive bounds are legal (1 event, 1 ms windows, 100.0% duty, floor 1)"
1322 );
1323 let cleared = PosturePolicyPatch {
1324 sim_intake_floor_override: Some(None),
1325 ..PosturePolicyPatch::default()
1326 }
1327 .validate();
1328 assert_eq!(
1329 cleared,
1330 Ok(()),
1331 "clearing the floor is a legal one-key patch"
1332 );
1333 }
1334
1335 #[test]
1336 fn patched_with_touches_only_supplied_keys() {
1337 let base = policy();
1338 let patched = base.patched_with(PosturePolicyPatch {
1339 enter_events: Some(7),
1340 ..PosturePolicyPatch::default()
1341 });
1342 assert_eq!(patched.enter_events, 7, "the supplied key landed");
1343 assert_eq!(patched.enter_window_ms, base.enter_window_ms);
1344 assert!(
1345 (patched.duty_percent - base.duty_percent).abs() < f64::EPSILON,
1346 "an absent key keeps the current duty percent"
1347 );
1348 assert_eq!(patched.duty_window_ms, base.duty_window_ms);
1349 assert_eq!(patched.exit_clean_ms, base.exit_clean_ms);
1350 assert_eq!(
1351 patched.sim_intake_floor_override, base.sim_intake_floor_override,
1352 "an absent key keeps the current value"
1353 );
1354 }
1355
1356 #[test]
1357 fn patched_with_distinguishes_floor_set_clear_and_absent() {
1358 let base = PosturePolicy {
1359 sim_intake_floor_override: Some(3),
1360 ..policy()
1361 };
1362 assert_eq!(
1364 base.patched_with(PosturePolicyPatch::default())
1365 .sim_intake_floor_override,
1366 Some(3)
1367 );
1368 assert_eq!(
1370 base.patched_with(PosturePolicyPatch {
1371 sim_intake_floor_override: Some(Some(1)),
1372 ..PosturePolicyPatch::default()
1373 })
1374 .sim_intake_floor_override,
1375 Some(1)
1376 );
1377 assert_eq!(
1380 base.patched_with(PosturePolicyPatch {
1381 sim_intake_floor_override: Some(None),
1382 ..PosturePolicyPatch::default()
1383 })
1384 .sim_intake_floor_override,
1385 None
1386 );
1387 }
1388
1389 #[test]
1390 fn a_validated_patch_retunes_the_owner_end_to_end() {
1391 let owner = PostureOwner::new(policy());
1392 let patch = PosturePolicyPatch {
1393 duty_percent: Some(5.0),
1394 ..PosturePolicyPatch::default()
1395 };
1396 assert_eq!(patch.validate(), Ok(()), "the channel validated the patch");
1397 let effective = owner.policy().patched_with(patch);
1398 owner.retune(effective);
1399 assert!(
1400 (owner.policy().duty_percent - 5.0).abs() < f64::EPSILON,
1401 "the retuned duty percent is live"
1402 );
1403 assert_eq!(owner.policy().enter_events, policy().enter_events);
1404 }
1405}