use degenbot_core::diag;
use degenbot_core::{op_info, op_warn};
use std::collections::VecDeque;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, OnceLock};
use degenbot_config::FleetConfig;
use parking_lot::{Mutex, RwLock};
use crate::role::CordonClass;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum FleetPosture {
Nominal,
Cordoned,
}
#[derive(Debug, Clone, Copy, PartialEq)]
pub enum EnterReason {
EventBurst {
events: u64,
},
DutySpike {
duty_percent: f64,
},
LaneDeath,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PostureCause {
LaneDeath,
}
#[derive(Debug, Clone, Copy, PartialEq)]
pub enum PostureChange {
Held,
Entered(EnterReason),
Exited,
}
#[derive(Debug, Clone, Copy, PartialEq)]
pub struct PosturePolicy {
pub enter_events: usize,
pub enter_window_ms: u64,
pub duty_percent: f64,
pub duty_window_ms: u64,
pub exit_clean_ms: u64,
pub sim_intake_floor_override: Option<usize>,
}
impl PosturePolicy {
#[must_use]
pub const fn doc_defaults() -> Self {
Self {
enter_events: 2,
enter_window_ms: 1_000,
duty_percent: 2.0,
duty_window_ms: 5_000,
exit_clean_ms: 10_000,
sim_intake_floor_override: None,
}
}
#[must_use]
pub fn from_config(cfg: &FleetConfig) -> Self {
Self {
enter_events: cfg.cordon_enter_events,
enter_window_ms: cfg.cordon_enter_window_ms,
duty_percent: cfg.cordon_duty_percent,
duty_window_ms: cfg.cordon_duty_window_ms,
exit_clean_ms: cfg.cordon_exit_clean_ms,
sim_intake_floor_override: cfg.cordon_sim_intake_floor,
}
}
#[must_use]
pub fn sim_intake_cap(&self, slot_cap: usize) -> usize {
self.sim_intake_floor_override
.unwrap_or(slot_cap / 2)
.min(slot_cap)
.max(1)
}
#[must_use]
pub fn patched_with(self, patch: PosturePolicyPatch) -> Self {
Self {
enter_events: patch.enter_events.unwrap_or(self.enter_events),
enter_window_ms: patch.enter_window_ms.unwrap_or(self.enter_window_ms),
duty_percent: patch.duty_percent.unwrap_or(self.duty_percent),
duty_window_ms: patch.duty_window_ms.unwrap_or(self.duty_window_ms),
exit_clean_ms: patch.exit_clean_ms.unwrap_or(self.exit_clean_ms),
sim_intake_floor_override: match patch.sim_intake_floor_override {
None => self.sim_intake_floor_override,
Some(floor) => floor,
},
}
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq)]
pub struct PosturePolicyPatch {
pub enter_events: Option<usize>,
pub enter_window_ms: Option<u64>,
pub duty_percent: Option<f64>,
pub duty_window_ms: Option<u64>,
pub exit_clean_ms: Option<u64>,
pub sim_intake_floor_override: Option<Option<usize>>,
}
impl PosturePolicyPatch {
#[must_use]
pub fn is_empty(&self) -> bool {
self.enter_events.is_none()
&& self.enter_window_ms.is_none()
&& self.duty_percent.is_none()
&& self.duty_window_ms.is_none()
&& self.exit_clean_ms.is_none()
&& self.sim_intake_floor_override.is_none()
}
pub fn validate(&self) -> Result<(), PostureRetuneError> {
if self.is_empty() {
return Err(PostureRetuneError::EmptyPatch);
}
if self.enter_events.is_some_and(|v| v < 1) {
return Err(PostureRetuneError::EnterEvents(
self.enter_events.unwrap_or_default(),
));
}
if self.enter_window_ms.is_some_and(|v| v == 0) {
return Err(PostureRetuneError::EnterWindow(0));
}
if let Some(duty) = self.duty_percent {
if !(duty > 0.0 && duty <= 100.0) {
return Err(PostureRetuneError::DutyPercent(duty));
}
}
if self.duty_window_ms.is_some_and(|v| v == 0) {
return Err(PostureRetuneError::DutyWindow(0));
}
if self.exit_clean_ms.is_some_and(|v| v == 0) {
return Err(PostureRetuneError::ExitClean(0));
}
if let Some(Some(floor)) = self.sim_intake_floor_override {
if floor < 1 {
return Err(PostureRetuneError::SimIntakeFloor(floor));
}
}
Ok(())
}
}
#[derive(Debug, Clone, Copy, PartialEq, thiserror::Error)]
pub enum PostureRetuneError {
#[error("at least one cordon threshold key is required (empty patch)")]
EmptyPatch,
#[error("cordon_enter_events must be >= 1, got {0}")]
EnterEvents(usize),
#[error("cordon_enter_window_ms must be > 0 ms, got {0} ms")]
EnterWindow(u64),
#[error("cordon_duty_percent must be in (0.0, 100.0], got {0}")]
DutyPercent(f64),
#[error("cordon_duty_window_ms must be > 0 ms, got {0} ms")]
DutyWindow(u64),
#[error("cordon_exit_clean_ms must be > 0 ms, got {0} ms")]
ExitClean(u64),
#[error("cordon_sim_intake_floor must be >= 1, got {0}")]
SimIntakeFloor(usize),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct ThrottleSample {
pub events: u64,
pub throttled_usec: u64,
pub elapsed_usec: u64,
}
impl ThrottleSample {
#[must_use]
pub const fn is_clean(&self) -> bool {
self.events == 0 && self.throttled_usec == 0
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct PostureCounters {
pub entered: u64,
pub exited: u64,
pub intake_suppressed: u64,
pub lane_deaths: u64,
}
#[derive(Debug, Clone, Copy)]
struct Sample {
now_ms: u64,
sample: ThrottleSample,
}
#[derive(Debug)]
pub struct PostureStateMachine {
state: FleetPosture,
policy: PosturePolicy,
samples: VecDeque<Sample>,
last_unclean_ms: Option<u64>,
counters: PostureCounters,
lane_death_hold: bool,
}
impl PostureStateMachine {
#[must_use]
pub fn new(policy: PosturePolicy) -> Self {
Self {
state: FleetPosture::Nominal,
policy,
samples: VecDeque::new(),
last_unclean_ms: None,
counters: PostureCounters::default(),
lane_death_hold: false,
}
}
#[must_use]
pub const fn state(&self) -> FleetPosture {
self.state
}
#[must_use]
pub const fn counters(&self) -> &PostureCounters {
&self.counters
}
#[must_use]
pub const fn lane_death_held(&self) -> bool {
self.lane_death_hold
}
#[must_use]
pub const fn policy(&self) -> &PosturePolicy {
&self.policy
}
pub fn set_policy(&mut self, policy: PosturePolicy) {
self.policy = policy;
}
pub fn observe(&mut self, now_ms: u64, sample: ThrottleSample) -> PostureChange {
self.samples.push_back(Sample { now_ms, sample });
self.prune(now_ms);
if !sample.is_clean() {
self.last_unclean_ms = Some(now_ms);
}
match self.state {
FleetPosture::Nominal => self.maybe_enter(now_ms),
FleetPosture::Cordoned => self.maybe_exit(now_ms),
}
}
pub fn observe_cause(&mut self, cause: PostureCause) -> PostureChange {
match cause {
PostureCause::LaneDeath => {
self.counters.lane_deaths += 1;
if self.lane_death_hold {
return PostureChange::Held;
}
self.lane_death_hold = true;
if self.state == FleetPosture::Cordoned {
op_warn!(domain = pump, lane_deaths = self.counters.lane_deaths,
"lane-death HOLD upgrades an existing cordon โ sticky, clean-window exit disabled"
);
return PostureChange::Held;
}
self.enter(EnterReason::LaneDeath)
}
}
}
pub const fn note_intake_suppressed(&mut self) {
self.counters.intake_suppressed += 1;
}
#[must_use]
pub const fn admits_lease(&self, class: CordonClass) -> bool {
!matches!(
(self.state, class),
(FleetPosture::Cordoned, CordonClass::Deferrable)
)
}
#[must_use]
pub fn sim_intake_cap(&self, slot_cap: usize) -> usize {
match self.state {
FleetPosture::Nominal => slot_cap,
FleetPosture::Cordoned => self.policy.sim_intake_cap(slot_cap),
}
}
fn prune(&mut self, now_ms: u64) {
let window = self.policy.duty_window_ms.max(self.policy.enter_window_ms);
while let Some(front) = self.samples.front() {
if now_ms.saturating_sub(front.now_ms) > window {
self.samples.pop_front();
} else {
break;
}
}
}
fn maybe_enter(&mut self, now_ms: u64) -> PostureChange {
let burst: u64 = self
.samples
.iter()
.filter(|s| now_ms.saturating_sub(s.now_ms) <= self.policy.enter_window_ms)
.map(|s| s.sample.events)
.sum();
if burst >= u64::try_from(self.policy.enter_events.max(1)).unwrap_or(1) {
return self.enter(EnterReason::EventBurst { events: burst });
}
let (events, throttled_usec, elapsed_usec) = self.duty_window_totals();
let _ = events;
if elapsed_usec > 0 {
#[expect(
clippy::cast_precision_loss,
reason = "duty percent is an f64 metric by definition (ยตs/ยตs ratio)"
)]
let duty_percent = throttled_usec as f64 / elapsed_usec as f64 * 100.0;
if duty_percent > self.policy.duty_percent {
return self.enter(EnterReason::DutySpike { duty_percent });
}
}
PostureChange::Held
}
fn duty_window_totals(&self) -> (u64, u64, u64) {
self.samples
.iter()
.fold((0, 0, 0), |(ev, th, el), Sample { sample, .. }| {
(
ev + sample.events,
th + sample.throttled_usec,
el + sample.elapsed_usec,
)
})
}
fn enter(&mut self, reason: EnterReason) -> PostureChange {
self.state = FleetPosture::Cordoned;
self.counters.entered += 1;
op_warn!(domain = pump, reason = ?reason,
entered = self.counters.entered,
"cordon ENTER โ deferrable intake held, sim intake floored; in-flight units complete"
);
PostureChange::Entered(reason)
}
fn maybe_exit(&mut self, now_ms: u64) -> PostureChange {
if self.lane_death_hold {
return PostureChange::Held;
}
let dirty_recently = self
.last_unclean_ms
.is_some_and(|last| now_ms.saturating_sub(last) < self.policy.exit_clean_ms);
if dirty_recently {
return PostureChange::Held;
}
self.state = FleetPosture::Nominal;
self.counters.exited += 1;
op_info!(
domain = pump,
exited = self.counters.exited,
clean_ms = self.policy.exit_clean_ms,
"cordon EXIT after clean-window hysteresis"
);
PostureChange::Exited
}
}
#[derive(Debug)]
pub struct PostureWatch {
shared: Arc<FeedShared>,
seen: AtomicU64,
}
impl PostureWatch {
#[must_use]
pub fn current(&self) -> FleetPosture {
self.shared.state.read().posture
}
#[must_use]
pub fn has_changed(&self) -> bool {
self.shared.state.read().seq != self.seen.load(Ordering::Relaxed)
}
#[must_use]
pub fn take_if_changed(&self) -> Option<FleetPosture> {
let state = self.shared.state.read();
if state.seq == self.seen.load(Ordering::Relaxed) {
return None;
}
self.seen.store(state.seq, Ordering::Relaxed);
Some(state.posture)
}
}
#[derive(Debug)]
struct FeedShared {
state: RwLock<FeedState>,
}
#[derive(Debug, Clone, Copy)]
struct FeedState {
posture: FleetPosture,
seq: u64,
}
#[derive(Debug)]
pub struct PostureOwner {
machine: Mutex<PostureStateMachine>,
broadcast: Arc<FeedShared>,
}
impl PostureOwner {
#[must_use]
pub fn new(policy: PosturePolicy) -> Self {
Self {
machine: Mutex::new(PostureStateMachine::new(policy)),
broadcast: Arc::new(FeedShared {
state: RwLock::new(FeedState {
posture: FleetPosture::Nominal,
seq: 0,
}),
}),
}
}
pub fn observe_throttle(&self, now_ms: u64, sample: ThrottleSample) -> PostureChange {
let (change, posture) = {
let mut machine = self.machine.lock();
let change = machine.observe(now_ms, sample);
(change, machine.state())
};
if !matches!(change, PostureChange::Held) {
self.publish(posture);
}
change
}
pub fn observe_cause(&self, cause: PostureCause) -> PostureChange {
let (change, posture) = {
let mut machine = self.machine.lock();
let change = machine.observe_cause(cause);
(change, machine.state())
};
if !matches!(change, PostureChange::Held) {
self.publish(posture);
}
change
}
#[must_use]
pub fn current(&self) -> FleetPosture {
self.machine.lock().state()
}
#[must_use]
pub fn policy(&self) -> PosturePolicy {
*self.machine.lock().policy()
}
#[must_use]
pub fn subscribe(&self) -> PostureWatch {
let seen = self.broadcast.state.read().seq;
PostureWatch {
shared: Arc::clone(&self.broadcast),
seen: AtomicU64::new(seen),
}
}
pub fn retune(&self, new_policy: PosturePolicy) {
let posture = {
let mut machine = self.machine.lock();
machine.set_policy(new_policy);
machine.state()
};
self.publish(posture);
}
#[must_use]
pub fn admits_lease(&self, class: CordonClass) -> bool {
self.machine.lock().admits_lease(class)
}
pub fn note_intake_suppressed(&self) {
self.machine.lock().note_intake_suppressed();
}
#[must_use]
pub fn sim_intake_cap(&self, slot_cap: usize) -> usize {
self.machine.lock().sim_intake_cap(slot_cap)
}
#[must_use]
pub fn counters(&self) -> PostureCounters {
*self.machine.lock().counters()
}
#[must_use]
pub fn lane_death_held(&self) -> bool {
self.machine.lock().lane_death_held()
}
fn publish(&self, posture: FleetPosture) {
let mut state = self.broadcast.state.write();
if state.posture != posture {
state.posture = posture;
state.seq = state.seq.wrapping_add(1);
}
}
}
static PROCESS_OWNER: OnceLock<PostureOwner> = OnceLock::new();
#[must_use]
pub fn install_process_owner(policy: PosturePolicy) -> &'static PostureOwner {
if PROCESS_OWNER.set(PostureOwner::new(policy)).is_err() {
diag!(
domain = pump,
"process owner already installed โ first-wins, keeping the existing owner"
);
}
process()
}
#[must_use]
pub fn process() -> &'static PostureOwner {
PROCESS_OWNER.get_or_init(|| PostureOwner::new(PosturePolicy::doc_defaults()))
}
#[cfg(test)]
mod tests {
use super::*;
fn policy() -> PosturePolicy {
PosturePolicy {
enter_events: 2,
enter_window_ms: 1_000,
duty_percent: 2.0,
duty_window_ms: 5_000,
exit_clean_ms: 10_000,
sim_intake_floor_override: None,
}
}
fn sm() -> PostureStateMachine {
PostureStateMachine::new(policy())
}
fn sample(events: u64, throttled_usec: u64, elapsed_usec: u64) -> ThrottleSample {
ThrottleSample {
events,
throttled_usec,
elapsed_usec,
}
}
#[test]
fn nominal_is_the_boot_state() {
assert_eq!(sm().state(), FleetPosture::Nominal);
}
#[test]
fn a_lane_death_cordons_immediately() {
let mut m = sm();
assert_eq!(
m.observe_cause(PostureCause::LaneDeath),
PostureChange::Entered(EnterReason::LaneDeath)
);
assert_eq!(m.state(), FleetPosture::Cordoned);
assert_eq!(m.counters().lane_deaths, 1);
assert_eq!(m.counters().entered, 1);
}
#[test]
fn the_lane_death_cordon_is_sticky_across_clean_windows() {
let mut m = sm();
m.observe_cause(PostureCause::LaneDeath);
for t in (0..30_u64).map(|i| 10_000 + i * 1_000) {
assert_eq!(
m.observe(t, sample(0, 0, 1_000_000)),
PostureChange::Held,
"a clean window must never lift a lane-death cordon"
);
}
assert_eq!(m.state(), FleetPosture::Cordoned);
assert_eq!(m.counters().exited, 0, "the sticky cordon never exits");
}
#[test]
fn a_lane_death_upgrades_a_throttle_cordon_to_sticky() {
let mut m = sm();
m.observe(100, sample(2, 0, 1_000_000));
assert_eq!(m.state(), FleetPosture::Cordoned);
assert_eq!(m.counters().entered, 1);
assert_eq!(
m.observe_cause(PostureCause::LaneDeath),
PostureChange::Held,
"no state transition โ the cordon was already up"
);
assert_eq!(m.counters().lane_deaths, 1);
assert_eq!(m.counters().entered, 1, "no second enter counted");
for t in (0..30_u64).map(|i| 10_000 + i * 1_000) {
m.observe(t, sample(0, 0, 1_000_000));
}
assert_eq!(m.state(), FleetPosture::Cordoned);
assert_eq!(m.counters().exited, 0);
}
#[test]
fn repeated_lane_deaths_count_but_do_not_re_enter() {
let mut m = sm();
m.observe_cause(PostureCause::LaneDeath);
for _ in 0..3 {
assert_eq!(
m.observe_cause(PostureCause::LaneDeath),
PostureChange::Held
);
}
assert_eq!(m.counters().lane_deaths, 4);
assert_eq!(m.counters().entered, 1);
}
#[test]
fn lane_death_cordon_effects_match_the_throttle_vocabulary() {
let mut m = sm();
m.observe_cause(PostureCause::LaneDeath);
assert!(!m.admits_lease(CordonClass::Deferrable));
assert!(m.admits_lease(CordonClass::Never));
assert!(m.admits_lease(CordonClass::SimPool));
assert_eq!(m.sim_intake_cap(4), 2);
}
#[test]
fn the_owner_publishes_the_lane_death_transition() {
let owner = PostureOwner::new(policy());
let watch = owner.subscribe();
assert_eq!(watch.current(), FleetPosture::Nominal);
owner.observe_cause(PostureCause::LaneDeath);
assert_eq!(watch.current(), FleetPosture::Cordoned);
assert_eq!(owner.current(), FleetPosture::Cordoned);
}
#[test]
fn enters_on_event_burst_within_the_window() {
let mut m = sm();
assert_eq!(m.observe(100, sample(1, 0, 1_000_000)), PostureChange::Held);
let change = m.observe(600, sample(1, 0, 500_000));
assert_eq!(
change,
PostureChange::Entered(EnterReason::EventBurst { events: 2 })
);
assert_eq!(m.state(), FleetPosture::Cordoned);
assert_eq!(m.counters().entered, 1);
assert!(matches!(
m.observe(700, sample(1, 0, 100_000)),
PostureChange::Held
));
assert_eq!(m.counters().entered, 1);
}
#[test]
fn burst_outside_the_enter_window_does_not_cordon() {
let mut m = sm();
m.observe(0, sample(1, 0, 1_000_000));
m.observe(3_000, sample(1, 0, 2_000_000));
m.observe(3_100, sample(0, 0, 100_000));
assert_eq!(m.state(), FleetPosture::Nominal);
}
#[test]
fn enters_on_duty_spike_over_the_duty_window() {
let mut m = sm();
let change = m.observe(1_000, sample(0, 500_000, 1_000_000));
assert!(matches!(
change,
PostureChange::Entered(EnterReason::DutySpike { .. })
));
assert_eq!(m.state(), FleetPosture::Cordoned);
assert_eq!(m.counters().entered, 1);
}
#[test]
fn a_rising_duty_only_crosses_once_the_trailing_total_exceeds_the_threshold() {
let mut m = sm();
for t in 1..5_u64 {
assert_eq!(
m.observe(t * 1_000, sample(0, 500, 1_000_000)),
PostureChange::Held
);
}
let change = m.observe(5_000, sample(0, 5_000_000, 1_000_000));
assert!(matches!(
change,
PostureChange::Entered(EnterReason::DutySpike { .. })
));
}
#[test]
fn sub_threshold_duty_never_cordons() {
let mut m = sm();
for t in 0..10_u64 {
m.observe((t + 1) * 1_000, sample(0, 10_000, 1_000_000));
}
assert_eq!(m.state(), FleetPosture::Nominal);
}
#[test]
fn exits_only_after_the_full_clean_hysteresis() {
let mut m = sm();
m.observe(100, sample(2, 0, 100_000)); assert_eq!(m.state(), FleetPosture::Cordoned);
let mut t = 200;
while t < 10_000 {
assert_eq!(m.observe(t, sample(0, 0, 1_000_000)), PostureChange::Held);
t += 1_000;
}
assert_eq!(m.state(), FleetPosture::Cordoned);
let change = m.observe(10_200, sample(0, 0, 100_000));
assert_eq!(change, PostureChange::Exited);
assert_eq!(m.state(), FleetPosture::Nominal);
assert_eq!(m.counters().exited, 1);
}
#[test]
fn a_dirty_window_restarts_the_clean_clock() {
let mut m = sm();
m.observe(0, sample(2, 0, 100_000));
let mut t = 1_000;
while t < 9_000 {
m.observe(t, sample(0, 0, 1_000_000));
t += 1_000;
}
m.observe(9_000, sample(1, 0, 100_000));
assert_eq!(m.state(), FleetPosture::Cordoned);
let mut t = 10_000;
while t < 18_000 {
m.observe(t, sample(0, 0, 1_000_000));
t += 1_000;
}
assert_eq!(m.state(), FleetPosture::Cordoned, "only 9 s clean");
m.observe(19_100, sample(0, 0, 100_000));
assert_eq!(
m.state(),
FleetPosture::Nominal,
"10 s clean since the reset"
);
}
#[test]
fn cordon_effects_match_the_sign_off_table() {
let mut m = sm();
assert_eq!(m.sim_intake_cap(4), 4);
assert!(m.admits_lease(CordonClass::Never));
assert!(m.admits_lease(CordonClass::SimPool));
assert!(m.admits_lease(CordonClass::Deferrable));
m.observe(0, sample(2, 0, 100_000));
assert_eq!(m.state(), FleetPosture::Cordoned);
assert!(!m.admits_lease(CordonClass::Deferrable));
assert!(m.admits_lease(CordonClass::SimPool));
assert!(m.admits_lease(CordonClass::Never));
assert_eq!(m.sim_intake_cap(4), 2, "floor = half the slot cap");
assert_eq!(m.sim_intake_cap(1), 1, "floored at one");
}
#[test]
fn override_intake_floor_wins_and_caps_at_the_slot_cap() {
let p = PosturePolicy {
sim_intake_floor_override: Some(7),
..policy()
};
assert_eq!(p.sim_intake_cap(4), 4, "nothing above the slot cap");
let p = PosturePolicy {
sim_intake_floor_override: Some(0),
..policy()
};
assert_eq!(p.sim_intake_cap(4), 1, "floored at one");
}
#[test]
fn suppressed_intake_is_counted_for_the_tuning_loop() {
let mut m = sm();
m.note_intake_suppressed();
m.note_intake_suppressed();
assert_eq!(m.counters().intake_suppressed, 2);
}
#[test]
fn typed_config_projects_all_the_q5_amendment_thresholds() {
let cfg = degenbot_config::BotConfig::default();
let p = PosturePolicy::from_config(&cfg.fleet);
assert_eq!(p.enter_events, 2);
assert_eq!(p.enter_window_ms, 1_000);
assert!((p.duty_percent - 2.0).abs() < 1e-9, "duty default is 2%");
assert_eq!(p.duty_window_ms, 5_000);
assert_eq!(p.exit_clean_ms, 10_000);
assert_eq!(p.sim_intake_floor_override, None);
}
#[test]
fn owner_transitions_publish_only_on_change() {
let owner = PostureOwner::new(policy());
let watch = owner.subscribe();
owner.observe_throttle(0, sample(0, 0, 1_000));
assert!(!watch.has_changed());
assert_eq!(watch.current(), FleetPosture::Nominal);
owner.observe_throttle(100, sample(1, 0, 1_000));
assert!(!watch.has_changed());
owner.observe_throttle(600, sample(1, 0, 500_000));
assert!(watch.has_changed());
assert_eq!(watch.take_if_changed(), Some(FleetPosture::Cordoned));
assert_eq!(watch.take_if_changed(), None, "one edge per transition");
owner.observe_throttle(700, sample(1, 0, 100_000));
assert!(!watch.has_changed());
assert_eq!(watch.current(), FleetPosture::Cordoned);
}
#[test]
fn owner_current_reflects_the_machine_including_exit_hysteresis() {
let owner = PostureOwner::new(policy());
assert_eq!(owner.current(), FleetPosture::Nominal);
owner.observe_throttle(0, sample(3, 0, 1_000));
assert_eq!(owner.current(), FleetPosture::Cordoned);
assert_eq!(owner.counters().entered, 1);
let mut now = 1_000;
loop {
owner.observe_throttle(now, sample(0, 0, 1_000));
if owner.current() == FleetPosture::Nominal {
break;
}
now += 1_000;
assert!(now <= 60_000, "the cordon never lifted");
}
assert_eq!(owner.counters().exited, 1);
}
#[test]
fn first_wins_process_install_keeps_the_existing_owner() {
let first = install_process_owner(policy());
let second = install_process_owner(PosturePolicy {
enter_events: 99,
..policy()
});
assert!(
std::ptr::eq(first, second),
"first-wins: a later install returns the existing owner"
);
assert!(std::ptr::eq(first, process()));
assert_eq!(
second.policy().enter_events,
first.policy().enter_events,
"the losing install's policy never landed"
);
}
#[test]
fn retune_swaps_thresholds_and_the_next_observe_rederives() {
let owner = PostureOwner::new(policy());
let watch = owner.subscribe();
owner.observe_throttle(0, sample(1, 0, 1_000_000));
assert_eq!(owner.current(), FleetPosture::Nominal);
owner.retune(PosturePolicy {
enter_events: 1,
..policy()
});
assert_eq!(owner.policy().enter_events, 1, "the retune swapped");
assert!(!watch.has_changed(), "a no-effect retune is not an edge");
owner.observe_throttle(2_000, sample(1, 0, 100_000));
assert_eq!(owner.current(), FleetPosture::Cordoned);
assert_eq!(watch.take_if_changed(), Some(FleetPosture::Cordoned));
}
#[test]
fn two_owners_are_fully_independent_hermetic_isolation() {
let a = PostureOwner::new(policy());
let b = PostureOwner::new(policy());
let watch_a = a.subscribe();
let watch_b = b.subscribe();
a.observe_throttle(0, sample(3, 0, 1_000));
assert_eq!(a.current(), FleetPosture::Cordoned);
assert_eq!(b.current(), FleetPosture::Nominal);
assert_eq!(watch_a.take_if_changed(), Some(FleetPosture::Cordoned));
assert_eq!(
watch_b.take_if_changed(),
None,
"no posture leaks across owners"
);
b.observe_throttle(1_000, sample(0, 0, 1_000));
assert_eq!(watch_b.take_if_changed(), None);
assert_eq!(watch_a.take_if_changed(), None);
}
#[test]
fn owner_read_throughs_match_the_machine_semantics() {
let owner = PostureOwner::new(PosturePolicy {
sim_intake_floor_override: Some(1),
..policy()
});
assert!(owner.admits_lease(CordonClass::Deferrable));
assert_eq!(owner.sim_intake_cap(8), 8, "nominal intake is the cap");
owner.observe_throttle(0, sample(3, 0, 1_000));
assert!(!owner.admits_lease(CordonClass::Deferrable));
assert!(owner.admits_lease(CordonClass::Never));
assert!(owner.admits_lease(CordonClass::SimPool));
assert_eq!(owner.sim_intake_cap(8), 1, "the cordon floor override");
assert_eq!(owner.counters().entered, 1);
}
#[test]
fn an_empty_patch_is_rejected() {
assert!(PosturePolicyPatch::default().validate().is_err());
assert!(PosturePolicyPatch::default().is_empty());
assert_eq!(
PosturePolicyPatch::default().validate(),
Err(PostureRetuneError::EmptyPatch),
"the empty-patch refusal is its own typed error"
);
}
#[test]
fn every_threshold_rule_is_a_typed_rejection() {
assert_eq!(
PosturePolicyPatch {
enter_events: Some(0),
..PosturePolicyPatch::default()
}
.validate(),
Err(PostureRetuneError::EnterEvents(0)),
"enter_events >= 1"
);
assert_eq!(
PosturePolicyPatch {
enter_window_ms: Some(0),
..PosturePolicyPatch::default()
}
.validate(),
Err(PostureRetuneError::EnterWindow(0)),
"windows must be > 0 ms"
);
assert_eq!(
PosturePolicyPatch {
duty_window_ms: Some(0),
..PosturePolicyPatch::default()
}
.validate(),
Err(PostureRetuneError::DutyWindow(0))
);
assert_eq!(
PosturePolicyPatch {
exit_clean_ms: Some(0),
..PosturePolicyPatch::default()
}
.validate(),
Err(PostureRetuneError::ExitClean(0))
);
for duty in [0.0, -1.0, 100.5, f64::NAN, f64::INFINITY] {
let rejected = PosturePolicyPatch {
duty_percent: Some(duty),
..PosturePolicyPatch::default()
}
.validate();
assert!(
matches!(rejected, Err(PostureRetuneError::DutyPercent(_))),
"duty {duty} must be refused as DutyPercent, got {rejected:?}"
);
}
assert_eq!(
PosturePolicyPatch {
sim_intake_floor_override: Some(Some(0)),
..PosturePolicyPatch::default()
}
.validate(),
Err(PostureRetuneError::SimIntakeFloor(0)),
"an explicit floor must be >= 1"
);
}
#[test]
fn boundary_values_are_admitted() {
let validated = PosturePolicyPatch {
enter_events: Some(1),
enter_window_ms: Some(1),
duty_percent: Some(100.0),
duty_window_ms: Some(1),
exit_clean_ms: Some(1),
sim_intake_floor_override: Some(Some(1)),
}
.validate();
assert_eq!(
validated,
Ok(()),
"inclusive bounds are legal (1 event, 1 ms windows, 100.0% duty, floor 1)"
);
let cleared = PosturePolicyPatch {
sim_intake_floor_override: Some(None),
..PosturePolicyPatch::default()
}
.validate();
assert_eq!(
cleared,
Ok(()),
"clearing the floor is a legal one-key patch"
);
}
#[test]
fn patched_with_touches_only_supplied_keys() {
let base = policy();
let patched = base.patched_with(PosturePolicyPatch {
enter_events: Some(7),
..PosturePolicyPatch::default()
});
assert_eq!(patched.enter_events, 7, "the supplied key landed");
assert_eq!(patched.enter_window_ms, base.enter_window_ms);
assert!(
(patched.duty_percent - base.duty_percent).abs() < f64::EPSILON,
"an absent key keeps the current duty percent"
);
assert_eq!(patched.duty_window_ms, base.duty_window_ms);
assert_eq!(patched.exit_clean_ms, base.exit_clean_ms);
assert_eq!(
patched.sim_intake_floor_override, base.sim_intake_floor_override,
"an absent key keeps the current value"
);
}
#[test]
fn patched_with_distinguishes_floor_set_clear_and_absent() {
let base = PosturePolicy {
sim_intake_floor_override: Some(3),
..policy()
};
assert_eq!(
base.patched_with(PosturePolicyPatch::default())
.sim_intake_floor_override,
Some(3)
);
assert_eq!(
base.patched_with(PosturePolicyPatch {
sim_intake_floor_override: Some(Some(1)),
..PosturePolicyPatch::default()
})
.sim_intake_floor_override,
Some(1)
);
assert_eq!(
base.patched_with(PosturePolicyPatch {
sim_intake_floor_override: Some(None),
..PosturePolicyPatch::default()
})
.sim_intake_floor_override,
None
);
}
#[test]
fn a_validated_patch_retunes_the_owner_end_to_end() {
let owner = PostureOwner::new(policy());
let patch = PosturePolicyPatch {
duty_percent: Some(5.0),
..PosturePolicyPatch::default()
};
assert_eq!(patch.validate(), Ok(()), "the channel validated the patch");
let effective = owner.policy().patched_with(patch);
owner.retune(effective);
assert!(
(owner.policy().duty_percent - 5.0).abs() < f64::EPSILON,
"the retuned duty percent is live"
);
assert_eq!(owner.policy().enter_events, policy().enter_events);
}
}