Skip to main content

degenbot_workers/
posture.rs

1//! The `Nominal ⇄ Cordoned` posture FSM (design doc §6) — the first
2//! consumer of `degenbot.cgroup.throttled` (`cpu_budget::cgroup_throttle_delta`).
3//!
4//! Thresholds are typed, runtime-tunable config keys (the sign-off
5//! amendment 2026-09-09: enter triggers, exit window, and cordon effects
6//! calibrated from soak data via the operator channel — never share
7//! arithmetic, which is [`crate::budget`]'s authority).
8//!
9//! Entering/exiting is LOUD: a transition fires a structured log line plus
10//! the posture counters (never a silent degrade). Time is passed in as
11//! monotonic milliseconds so the FSM is deterministic under test; callers
12//! feed `Instant::now()` deltas from their throttle poller.
13
14use 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/// Process-level fleet posture. NOT a slot state (design doc §3.2): it
26/// gates lease transitions, it never sheds a running unit.
27#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
28pub enum FleetPosture {
29    /// Normal operation: every dispatchable role leases freely.
30    Nominal,
31    /// Under cgroup throttling: no new leases for cordon-deferrable roles,
32    /// sim intake floored; in-flight units COMPLETE (T7/T8), pins are never
33    /// shed, the merge pin and ambient I/O are never cordoned.
34    Cordoned,
35}
36
37/// Why the posture entered cordon (exported as the transition's cause).
38#[derive(Debug, Clone, Copy, PartialEq)]
39pub enum EnterReason {
40    /// ≥ `enter_events` throttle events within the rolling enter window.
41    EventBurst {
42        /// The events observed in the window.
43        events: u64,
44    },
45    /// Throttled-time duty exceeded `duty_percent` over the duty window.
46    DutySpike {
47        /// Measured duty percent.
48        duty_percent: f64,
49    },
50    /// A lane died mid-flight (FF-T4, Z6XTDX — the DECIDED option (a)
51    /// input). Entered with its OWN exit discipline: see
52    /// [`PostureCause::LaneDeath`].
53    LaneDeath,
54}
55
56/// A typed non-throttle posture cause (FF-T4, Z6XTDX — the DECIDED
57/// option (a): the input is typed AT the posture owner, not a sample
58/// it has to infer from; the failure taxonomy stays CLOSED per
59/// ADR-040's per-bucket reactions).
60#[derive(Debug, Clone, Copy, PartialEq, Eq)]
61pub enum PostureCause {
62    /// A lane died mid-flight: enter the cordon IMMEDIATELY (no sample
63    /// hysteresis — the fleet is degraded NOW) with its OWN exit
64    /// discipline: the clean window NEVER lifts a lane-death cordon
65    /// (the lane is still dead) — the cordon is STICKY until a fresh
66    /// process. In-flight paths get terminal receipts; the process
67    /// stays alive (the cordoned posture holds deferrable intake and
68    /// floors sim intake — the §6 cordon effects).
69    LaneDeath,
70}
71
72/// The posture change a [`PostureStateMachine::observe`] tick produced.
73#[derive(Debug, Clone, Copy, PartialEq)]
74pub enum PostureChange {
75    /// No transition (the FSM does not re-enter while already cordoned).
76    Held,
77    /// Nominal → Cordoned (loud: log at warn + counter).
78    Entered(EnterReason),
79    /// Cordoned → Nominal after the clean-window hysteresis (log at info).
80    Exited,
81}
82
83/// Typed thresholds (Q5 amendment) — [`PosturePolicy::from_config`] is the
84/// boot source; the operator channel re-derives from the same schema keys.
85#[derive(Debug, Clone, Copy, PartialEq)]
86pub struct PosturePolicy {
87    /// Enter trigger (a): >= this many throttle events in the enter window.
88    pub enter_events: usize,
89    /// Enter window (a): the rolling burst window.
90    pub enter_window_ms: u64,
91    /// Enter trigger (b): throttled-time duty percent over the duty window
92    /// (`2.0` = >2%).
93    pub duty_percent: f64,
94    /// Enter trigger (b) window.
95    pub duty_window_ms: u64,
96    /// Exit: this much continuous clean time required after the last dirty
97    /// sample (hysteresis prevents flapping; §6: 10 s).
98    pub exit_clean_ms: u64,
99    /// Cordon effect (b): sim intake cap while cordoned; `None` = half the
100    /// slot cap (in-flight sims are never cancelled).
101    pub sim_intake_floor_override: Option<usize>,
102}
103
104impl PosturePolicy {
105    /// The design-doc §6 defaults (2 events / 1 s, >2% / 5 s, 10 s clean),
106    /// used when a host boots without a typed config (hermetic runs).
107    #[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    /// The typed-config projection — every threshold is a schema key.
120    #[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    /// The cordon sim-intake cap: override or half the slot cap, floored
133    /// at 1, never above the cap (§6 effect (b)).
134    #[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    /// Apply a validated [`PosturePolicyPatch`] to `self`, producing the
143    /// effective policy (the JCI2FW Part B re-tune channel's only write
144    /// path: current policy + supplied fields). Pure — the caller feeds the
145    /// result to [`PostureOwner::retune`]. Call [`PosturePolicyPatch::validate`]
146    /// FIRST; this projection never checks semantics.
147    #[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                // Key absent: keep the current override.
157                None => self.sim_intake_floor_override,
158                // Key present: set it — `Some(v)` = explicit floor,
159                // `None` = cleared (back to half the slot cap).
160                Some(floor) => floor,
161            },
162        }
163    }
164}
165
166/// A partial re-tune request over the six typed thresholds (JCI2FW Part B,
167/// the operator channel's wire shape): every field is `None` = "key not
168/// supplied — keep the current value". `sim_intake_floor_override` is
169/// doubly-`Option`: the OUTER `None` is key-absent, and the inner
170/// `Some(None)` is the operator supplying the key's `None` value (clear the
171/// override, back to half the slot cap — the typed key itself is
172/// `opt usize`).
173///
174/// The semantic rules (windows > 0, duty percent in range, floors >= 1,
175/// non-empty patch) are encoded ONCE, in [`Self::validate`] — callers
176/// REJECT, never clamp silently.
177#[derive(Debug, Clone, Copy, Default, PartialEq)]
178pub struct PosturePolicyPatch {
179    /// `cordon_enter_events`: the burst trigger count.
180    pub enter_events: Option<usize>,
181    /// `cordon_enter_window_ms`: the burst window.
182    pub enter_window_ms: Option<u64>,
183    /// `cordon_duty_percent`: the duty trigger percent.
184    pub duty_percent: Option<f64>,
185    /// `cordon_duty_window_ms`: the duty window.
186    pub duty_window_ms: Option<u64>,
187    /// `cordon_exit_clean_ms`: the exit hysteresis.
188    pub exit_clean_ms: Option<u64>,
189    /// `cordon_sim_intake_floor`: outer `None` = key absent; inner
190    /// `Some(None)` = clear the override; `Some(Some(n))` = explicit floor.
191    pub sim_intake_floor_override: Option<Option<usize>>,
192}
193
194impl PosturePolicyPatch {
195    /// Whether NO key was supplied (an empty patch — rejected by
196    /// [`Self::validate`]; a re-tune that changes nothing must never look
197    /// like a successful one).
198    #[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    /// The ONE encoding of the re-tune's semantic rules. Every violation is
209    /// a typed [`PostureRetuneError`] naming the offending key and value —
210    /// the channel rejects, it never clamps.
211    ///
212    /// # Errors
213    ///
214    /// [`PostureRetuneError::EmptyPatch`] when no key was supplied, or the
215    /// per-key range error for the first offending value.
216    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            // The sane inclusive range: a measured duty percent is a
230            // (0.0, 100.0] quantity — 0 would cordon on ANY throttled
231            // microsecond and >100 or non-finite cannot be a duty.
232            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/// Why the re-tune channel refused a [`PosturePolicyPatch`] (JCI2FW Part
252/// B). The rules live once, in [`PosturePolicyPatch::validate`]; the wire
253/// layer maps these to its typed channel error verbatim.
254#[derive(Debug, Clone, Copy, PartialEq, thiserror::Error)]
255pub enum PostureRetuneError {
256    /// No threshold key was supplied — an empty patch is refused.
257    #[error("at least one cordon threshold key is required (empty patch)")]
258    EmptyPatch,
259    /// `cordon_enter_events` below the floor of 1.
260    #[error("cordon_enter_events must be >= 1, got {0}")]
261    EnterEvents(usize),
262    /// `cordon_enter_window_ms` must be a positive window.
263    #[error("cordon_enter_window_ms must be > 0 ms, got {0} ms")]
264    EnterWindow(u64),
265    /// `cordon_duty_percent` outside the sane (0.0, 100.0] range.
266    #[error("cordon_duty_percent must be in (0.0, 100.0], got {0}")]
267    DutyPercent(f64),
268    /// `cordon_duty_window_ms` must be a positive window.
269    #[error("cordon_duty_window_ms must be > 0 ms, got {0} ms")]
270    DutyWindow(u64),
271    /// `cordon_exit_clean_ms` must be a positive window.
272    #[error("cordon_exit_clean_ms must be > 0 ms, got {0} ms")]
273    ExitClean(u64),
274    /// `cordon_sim_intake_floor` below the floor of 1.
275    #[error("cordon_sim_intake_floor must be >= 1, got {0}")]
276    SimIntakeFloor(usize),
277}
278
279/// One throttle-poll delta: `cgroup_throttle_delta()`'s counters plus the
280/// elapsed wall time of the poll interval.
281#[derive(Debug, Clone, Copy, PartialEq, Eq)]
282pub struct ThrottleSample {
283    /// `nr_throttled` delta since the last sample.
284    pub events: u64,
285    /// `throttled_usec` delta since the last sample.
286    pub throttled_usec: u64,
287    /// Elapsed wall time (µs) of the poll interval backing the deltas.
288    pub elapsed_usec: u64,
289}
290
291impl ThrottleSample {
292    /// A clean sample: no throttle events, no throttled time.
293    #[must_use]
294    pub const fn is_clean(&self) -> bool {
295        self.events == 0 && self.throttled_usec == 0
296    }
297}
298
299/// Loud-transition counters (exported with the posture metrics so
300/// thresholds are tuned against measurements, §6).
301#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
302pub struct PostureCounters {
303    /// Cordons entered.
304    pub entered: u64,
305    /// Cordons exited.
306    pub exited: u64,
307    /// Lease grants denied while cordoned (deferrable intake held + sim
308    /// intake suppression above the floor).
309    pub intake_suppressed: u64,
310    /// Lane-death causes observed (FF-T4 — the sticky cordons, incl.
311    /// upgrades of an existing throttle cordon).
312    pub lane_deaths: u64,
313}
314
315#[derive(Debug, Clone, Copy)]
316struct Sample {
317    now_ms: u64,
318    sample: ThrottleSample,
319}
320
321/// The posture state machine. Deterministic: time is monotonic ms fed by
322/// the caller.
323#[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    /// A lane-death cordon is STICKY (FF-T4): once set, the clean-window
331    /// exit refuses to lift the cordon — the lane is still dead. Only
332    /// a fresh process clears it (the operator restart path); there is
333    /// deliberately no in-process clear API (a sticky cordon that
334    /// silently un-sticks is exactly the silent-narrow class §10
335    /// forbids).
336    lane_death_hold: bool,
337}
338
339impl PostureStateMachine {
340    /// A fresh machine in [`FleetPosture::Nominal`].
341    #[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    /// Current posture.
354    #[must_use]
355    pub const fn state(&self) -> FleetPosture {
356        self.state
357    }
358
359    /// Loud-transition counters (metrics export surface).
360    #[must_use]
361    pub const fn counters(&self) -> &PostureCounters {
362        &self.counters
363    }
364
365    /// The sticky lane-death hold (FF-T4): once set, no clean window can
366    /// ever lift it — the `FleetHost` Faulted transition (TB4QGX T6) keys on
367    /// THIS typed latch, never on elapsed cordon time.
368    #[must_use]
369    pub const fn lane_death_held(&self) -> bool {
370        self.lane_death_hold
371    }
372
373    /// The active policy (operator-channel re-tune reads it through here).
374    #[must_use]
375    pub const fn policy(&self) -> &PosturePolicy {
376        &self.policy
377    }
378
379    /// Swap the policy (the operator channel's re-tune). Pure: the state
380    /// and the trailing sample window are KEPT — a retune never fabricates
381    /// samples, so the next [`Self::observe`] re-derives the posture under
382    /// the new thresholds. Semantic validation of the new policy is the
383    /// re-tune caller's job (the JCI2FW Part B channel).
384    pub fn set_policy(&mut self, policy: PosturePolicy) {
385        self.policy = policy;
386    }
387
388    /// Feed one throttle-poll delta at `now_ms`. Returns whether the
389    /// posture transitioned (and why) — the caller surfaces that loudly.
390    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    /// Feed one typed non-throttle cause (FF-T4, Z6XTDX — the DECIDED
404    /// option (a) input). A [`PostureCause::LaneDeath`] enters the
405    /// cordon from ANY state with its OWN exit discipline: the cordon
406    /// is sticky (the clean window never lifts it — see
407    /// [`Self::maybe_exit`]), in-flight paths get terminal receipts at
408    /// the detection site, and the process stays alive under the §6
409    /// cordon effects. Idempotent while the hold is already set (the
410    /// counter still counts every cause).
411    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                    // An existing (throttle) cordon UPGRADES to the
421                    // sticky hold: no state transition, but the cause
422                    // is new and loud.
423                    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    /// A lease grant was denied because of the posture (counter only — the
434    /// dispatcher supplies the typed error).
435    pub const fn note_intake_suppressed(&mut self) {
436        self.counters.intake_suppressed += 1;
437    }
438
439    /// Whether the posture admits new lease intake for `class` right now
440    /// (§6: `Never` and `SimPool` lease freely — sim is only intake-FLOORED;
441    /// `Deferrable` is held while cordoned).
442    #[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    /// The sim intake cap in the current posture (§6 effect (b)).
451    #[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        // Trigger (a): event burst inside the enter window.
472        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        // Trigger (b): duty percent over the trailing duty window.
482        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        // The lane-death hold is STICKY (FF-T4): the clean window never
521        // lifts it — the lane is still dead. Throttle samples keep
522        // feeding the machine (harmless bookkeeping); only a fresh
523        // process clears the hold.
524        if self.lane_death_hold {
525            return PostureChange::Held;
526        }
527        // Exit: `exit_clean_ms` of clean time since the last dirty sample
528        // (hysteresis; §6: 10 s of clean windows).
529        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// ---- the ONE process-level fleet posture owner (JCI2FW Part A) ------------
548
549/// Watch-style subscription to the fleet posture feed — the workers-crate
550/// equivalent of a `tokio::sync::watch` receiver (the crate carries no
551/// tokio; this is the `parking_lot` pattern its other feeds use). One
552/// producer (the [`PostureOwner`]), many independent consumers: every
553/// `FleetHost` subscribes at boot for the T7 shed trigger, tests subscribe
554/// to observe transitions, and the Part B operator channel will drive the
555/// owner directly. A consumer sees the latest posture and REAL transitions
556/// only — a `Held` tick or a no-effect retune never raises an edge.
557#[derive(Debug)]
558pub struct PostureWatch {
559    shared: Arc<FeedShared>,
560    seen: AtomicU64,
561}
562
563impl PostureWatch {
564    /// The latest posture (reads through — never consumes the edge).
565    #[must_use]
566    pub fn current(&self) -> FleetPosture {
567        self.shared.state.read().posture
568    }
569
570    /// Did a transition occur since the last drain? (Does not consume.)
571    #[must_use]
572    pub fn has_changed(&self) -> bool {
573        self.shared.state.read().seq != self.seen.load(Ordering::Relaxed)
574    }
575
576    /// Drain: the posture if a transition occurred since the last drain,
577    /// else `None` (consuming the edge — exactly-once per transition).
578    #[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/// The broadcast cell behind the feed: the posture plus a monotonic
590/// transition sequence (the edge counter the watches diff against).
591#[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/// The ONE owner of a fleet posture state machine — a shared, thread-safe
603/// shell around the pure [`PostureStateMachine`]. Every consumer (every
604/// `FleetHost`, the block pump's throttle feed, tests) consults THE SAME
605/// instance, so there is exactly one posture per process (per hermetic
606/// test scope) and no host-local mirrors to drift apart.
607///
608/// Sync: the machine sits behind a `parking_lot::Mutex` — every consult is
609/// a short read-through critical section, never held across an await
610/// (the crate has none); the transition feed is the [`PostureWatch`]
611/// broadcast above.
612#[derive(Debug)]
613pub struct PostureOwner {
614    machine: Mutex<PostureStateMachine>,
615    broadcast: Arc<FeedShared>,
616}
617
618impl PostureOwner {
619    /// A fresh owner in [`FleetPosture::Nominal`]. Hermetic tests build
620    /// their own owner and inject it via `FleetBoot::owner` — NEVER the
621    /// process global ([`process`]/[`install_process_owner`]); posture
622    /// leaking across tests is a failure class (7KAPBB).
623    #[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    /// Feed one throttle-poll delta. Publishes to the feed ONLY on a real
637    /// transition (`Held` ticks are silent — a subscriber never sees a
638    /// spurious edge).
639    ///
640    /// # Feeder-site contract (TB4QGX T3, ADR-044)
641    /// A BOT-side caller of this method MUST also wake the fleet hosts on a
642    /// non-`Held` change (the degenbot-bot host waker,
643    /// `arb_engine::fleet_wake::wake_hosts`), which emits ONE untrusted,
644    /// seq-stamped `PostureEdge` per host. This crate cannot know about host
645    /// channels (layering), so the wake is the caller's obligation; the
646    /// host's `BackstopTick` bounds the damage if a feeder forgets, and the
647    /// hint never carries a posture value — hosts re-read the live owner.
648    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    /// Feed one typed non-throttle cause (FF-T4, Z6XTDX). Publishes to
661    /// the feed on a real transition like [`Self::observe_throttle`]
662    /// (an idempotent hold-upgrade returns `Held` and stays silent —
663    /// the detection site owns the loud lane-death log).
664    ///
665    /// # Feeder-site contract (TB4QGX T3, ADR-044)
666    /// A BOT-side caller of this method MUST also wake the fleet hosts on a
667    /// non-`Held` change; see [`Self::observe_throttle`]. On the lane-death
668    /// (`Faulted`) arm the wake still fires, and the host drains its held
669    /// receipts terminally instead of parking them.
670    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    /// The current posture.
683    #[must_use]
684    pub fn current(&self) -> FleetPosture {
685        self.machine.lock().state()
686    }
687
688    /// The active policy (read-through; the Part B operator channel reads
689    /// and re-tunes through here).
690    #[must_use]
691    pub fn policy(&self) -> PosturePolicy {
692        *self.machine.lock().policy()
693    }
694
695    /// Subscribe a watch: the receiver starts at the CURRENT posture with
696    /// no pending edge (it observes only transitions from here on).
697    #[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    /// Swap the policy (the Part B operator channel's entry point). The
707    /// swap is atomic under the machine lock and keeps the state + the
708    /// trailing sample window; the feed re-publishes the current posture
709    /// so it mirrors the machine post-swap (a no-op unless the posture
710    /// itself changed — the feed carries only real transitions, and the
711    /// next `observe_throttle` re-derives the posture under the new
712    /// thresholds). Semantic validation of the new policy is the caller's
713    /// job.
714    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    /// Whether the posture admits new lease intake for `class` right now
724    /// (read-through — the dispatcher's enqueue/T-table gates and the seat
725    /// hosts' admission all consult this).
726    #[must_use]
727    pub fn admits_lease(&self, class: CordonClass) -> bool {
728        self.machine.lock().admits_lease(class)
729    }
730
731    /// Count a lease grant denied because of the posture (read-through —
732    /// the tuning loop's suppression metric).
733    pub fn note_intake_suppressed(&self) {
734        self.machine.lock().note_intake_suppressed();
735    }
736
737    /// The sim intake cap in the current posture (read-through).
738    #[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    /// Loud-transition counters snapshot (the tuning loop's metrics).
744    #[must_use]
745    pub fn counters(&self) -> PostureCounters {
746        *self.machine.lock().counters()
747    }
748
749    /// The sticky lane-death hold (read-through). A host keys its Faulted
750    /// transition on this TYPED latch, never on elapsed cordon time (a long
751    /// recoverable EventBurst/Duty cordon must not fault).
752    #[must_use]
753    pub fn lane_death_held(&self) -> bool {
754        self.machine.lock().lane_death_held()
755    }
756
757    /// Mirror the machine's state into the feed — compare-then-publish, so
758    /// the sequence (and every subscriber's edge) moves ONLY on a real
759    /// posture change.
760    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
769/// The process-level posture owner (the KAHU5W holder pattern): ONE
770/// instance per process, installed by the FIRST fleet boot (first-wins —
771/// later installs log at debug and return the existing owner).
772static PROCESS_OWNER: OnceLock<PostureOwner> = OnceLock::new();
773
774/// Install the process-level owner with `policy` (first-wins: the first
775/// fleet boot wins; later calls log at debug and return the existing
776/// owner). Production installs happen at `FleetHost::boot` — BEFORE the
777/// first throttle feed reaches [`process`]. Hermetic tests never call
778/// this: they inject fresh owners via `FleetBoot::owner`.
779#[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/// The process-level owner, or a doc-default owner when nothing was
791/// installed yet (the holder's default stance: tests and standalone
792/// constructions observe schema defaults). `FleetBoot` without an injected
793/// owner installs its policy here first-wins at boot time, which is why a
794/// production feed never lands on the doc-default stance.
795#[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    /// FF-T4 — the DECIDED option (a): a typed `PostureCause`
833    /// enters the cordon IMMEDIATELY (no sample hysteresis).
834    #[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    /// The lane-death cordon is STICKY: the clean window never lifts it
847    /// (the lane is still dead) — its own exit discipline, deliberately
848    /// different from the throttle hysteresis.
849    #[test]
850    fn the_lane_death_cordon_is_sticky_across_clean_windows() {
851        let mut m = sm();
852        m.observe_cause(PostureCause::LaneDeath);
853        // A long run of clean samples past the full exit window: the
854        // throttle hysteresis would EXIT here — the lane-death hold
855        // refuses (never a silent un-cordon).
856        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    /// A lane death UPGRADES an existing throttle cordon: no state
868    /// transition (already Cordoned), but the hold turns sticky and the
869    /// cause is counted — loud, never silent.
870    #[test]
871    fn a_lane_death_upgrades_a_throttle_cordon_to_sticky() {
872        let mut m = sm();
873        // Enter via the throttle hysteresis first.
874        m.observe(100, sample(2, 0, 1_000_000));
875        assert_eq!(m.state(), FleetPosture::Cordoned);
876        assert_eq!(m.counters().entered, 1);
877        // The lane death upgrades the hold.
878        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        // And the upgrade is sticky: clean windows no longer lift it.
886        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    /// Idempotent while the hold is already set (every cause counts).
894    #[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    /// The cordon effects hold under a lane-death cause exactly as under
909    /// a throttle cause (the §6 vocabulary: deferrable intake held,
910    /// sim intake floored, in-flight units complete).
911    #[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        // The sim intake floor (a Cordoned cap of 4 slots -> 2, the
919        // same shape a throttle cordon produces).
920        assert_eq!(m.sim_intake_cap(4), 2);
921    }
922
923    /// The owner publishes the lane-death transition to the feed (a
924    /// subscriber sees the Cordoned edge exactly once).
925    #[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        // One lone event inside the window: below the threshold.
939        assert_eq!(m.observe(100, sample(1, 0, 1_000_000)), PostureChange::Held);
940        // The second event inside the same trailing 1 s window -> enter,
941        // loudly typed.
942        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        // Already cordoned: further dirty samples do not re-enter.
950        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        // The 1 s window expired; two lone events > 1 s apart never burst.
962        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        // A 1 s poll interval whose throttled time was 50% — measured over
971        // the trailing 5 s window that is 50% duty, far above the 2%
972        // threshold: the duty trigger fires on the first observation.
973        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        // Sub-threshold ticks (0.05% each) hold the posture...
986        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        // ...until one more dirty tick crosses the window total above 2%.
993        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        // 1% duty (design doc: steady state 0.18%) for ten windows.
1004        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)); // burst enter
1014        assert_eq!(m.state(), FleetPosture::Cordoned);
1015        // Clean time below the hysteresis keeps the cordon.
1016        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        // Past 10 s clean: exit.
1023        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        // A lone dirty sample resets the clean-window clock.
1039        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        // Nominal: full sim cap; deferrable admitted.
1059        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        // Cordon: deferrable held, sim floored at half the cap, Never free.
1067        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        // Defaults mirror design doc §6: 2 events / 1 s, >2% / 5 s, 10 s.
1101        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    // ---- PostureOwner (JCI2FW Part A) -------------------------------------
1110
1111    #[test]
1112    fn owner_transitions_publish_only_on_change() {
1113        let owner = PostureOwner::new(policy());
1114        let watch = owner.subscribe();
1115        // A clean sample: no transition, no publication.
1116        owner.observe_throttle(0, sample(0, 0, 1_000));
1117        assert!(!watch.has_changed());
1118        assert_eq!(watch.current(), FleetPosture::Nominal);
1119        // One lone event below the threshold: still no edge.
1120        owner.observe_throttle(100, sample(1, 0, 1_000));
1121        assert!(!watch.has_changed());
1122        // The bursting event crosses: exactly one edge, and the drain
1123        // consumes it (exactly-once per transition).
1124        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        // Already cordoned: dirty ticks are Held — never re-published.
1129        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        // The full clean hysteresis lifts the cordon (10 s of clean ticks).
1142        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        // One lone event: sub-threshold under the boot policy.
1178        owner.observe_throttle(0, sample(1, 0, 1_000_000));
1179        assert_eq!(owner.current(), FleetPosture::Nominal);
1180        // Retune to a 1-event trigger: no posture change, no edge — but
1181        // the next lone (clean) sample cordons under the new thresholds.
1182        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 cordons; b never hears about it.
1200        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's own feed stays silent on a clean tick, and the feeds stay
1210        // independent in both directions.
1211        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    // ---- the operator re-tune patch (JCI2FW Part B) -----------------------
1233
1234    #[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        // The sane inclusive duty range (0.0, 100.0]: 0, negatives, >100,
1282        // and non-finite values are all refused, never clamped.
1283        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        // The inclusive edges pass: 1 event, 1 ms windows, the 100.0% duty
1308        // ceiling, a floor of exactly 1, and a CLEARED floor.
1309        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        // Key absent: the current override is kept.
1363        assert_eq!(
1364            base.patched_with(PosturePolicyPatch::default())
1365                .sim_intake_floor_override,
1366            Some(3)
1367        );
1368        // Key present with a value: the override is set.
1369        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        // Key present with the key's None value: the override is CLEARED
1378        // (back to half the slot cap at read time).
1379        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}