use std::time::{Duration, Instant};
use super::incarnation::Incarnation;
#[derive(Clone, Copy, PartialEq, Eq, Debug, serde::Serialize, serde::Deserialize)]
pub enum AttestedStatus {
Ready,
NotReady,
ProviderUnknown,
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub enum Continuity {
Unestablished,
Established,
Expired,
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub enum ProjectedReadiness {
Ready,
NotReady,
Unknown,
}
pub const fn project(status: AttestedStatus, continuity: Continuity) -> ProjectedReadiness {
match (status, continuity) {
(_, Continuity::Expired) => ProjectedReadiness::Unknown,
(AttestedStatus::Ready, Continuity::Established) => ProjectedReadiness::Ready,
(AttestedStatus::Ready, Continuity::Unestablished) => ProjectedReadiness::Unknown,
(AttestedStatus::NotReady, _) => ProjectedReadiness::NotReady,
(AttestedStatus::ProviderUnknown, _) => ProjectedReadiness::Unknown,
}
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub enum DisruptReason {
PathFailed,
IncarnationSuperseded,
GenerationChanged,
ScopeValidationFailed,
}
#[derive(Clone, Copy, Debug)]
pub struct ReadinessObservation {
pub attested_status: AttestedStatus,
pub estimated_start: Option<Duration>,
pub source_incarnation: Incarnation,
pub capability_generation: u64,
pub last_seq: u64,
pub promised_cadence: Duration,
pub continuity: Continuity,
pub locally_observed_at: Instant,
}
impl ReadinessObservation {
pub const fn projected(&self) -> ProjectedReadiness {
project(self.attested_status, self.continuity)
}
}
#[derive(Clone, Copy, Debug)]
pub struct DeliveredBeat {
pub attested_status: AttestedStatus,
pub estimated_start: Option<Duration>,
pub source_incarnation: Incarnation,
pub capability_generation: u64,
pub seq: u64,
pub promised_cadence: Duration,
pub continuity_bearing: bool,
}
#[derive(Debug)]
pub struct ObservationCell {
observation: Option<ReadinessObservation>,
continuity: Continuity,
deadline: Instant,
own_interval: Duration,
factor: u32,
last_disrupt: Option<DisruptReason>,
}
impl ObservationCell {
pub fn register(now: Instant, own_interval: Duration, factor: u32) -> Self {
Self {
observation: None,
continuity: Continuity::Unestablished,
deadline: now + own_interval.saturating_mul(factor),
own_interval,
factor,
last_disrupt: None,
}
}
fn window(&self, promised_cadence: Duration) -> Duration {
promised_cadence
.max(self.own_interval)
.saturating_mul(self.factor)
}
pub fn update_interval(&mut self, own_interval: Duration) {
if own_interval == self.own_interval {
return;
}
if self.continuity == Continuity::Expired {
self.own_interval = own_interval;
return;
}
let old_window = self.deadline_window(self.own_interval);
self.own_interval = own_interval;
let new_window = self.deadline_window(own_interval);
if new_window >= old_window {
self.deadline += new_window - old_window;
} else if let Some(deadline) = self.deadline.checked_sub(old_window - new_window) {
self.deadline = deadline;
}
}
fn deadline_window(&self, own_interval: Duration) -> Duration {
match self.continuity {
Continuity::Established => {
let promised = self
.observation
.as_ref()
.map(|obs| obs.promised_cadence)
.unwrap_or(Duration::ZERO);
promised.max(own_interval).saturating_mul(self.factor)
}
Continuity::Unestablished | Continuity::Expired => {
own_interval.saturating_mul(self.factor)
}
}
}
pub fn on_admitted_beat(&mut self, now: Instant, beat: DeliveredBeat) {
let crossed_incarnation = self
.observation
.is_some_and(|obs| obs.source_incarnation != beat.source_incarnation);
let crossed_generation = self
.observation
.is_some_and(|obs| obs.capability_generation != beat.capability_generation);
if beat.continuity_bearing {
self.continuity = Continuity::Established;
self.deadline = now + self.window(beat.promised_cadence);
self.last_disrupt = None;
} else if crossed_incarnation && self.continuity == Continuity::Established {
self.continuity = Continuity::Expired;
self.last_disrupt = Some(DisruptReason::IncarnationSuperseded);
} else if crossed_generation {
self.continuity = Continuity::Unestablished;
self.deadline = now + self.own_interval.saturating_mul(self.factor);
self.last_disrupt = Some(DisruptReason::GenerationChanged);
}
self.observation = Some(ReadinessObservation {
attested_status: beat.attested_status,
estimated_start: beat.estimated_start,
source_incarnation: beat.source_incarnation,
capability_generation: beat.capability_generation,
last_seq: beat.seq,
promised_cadence: beat.promised_cadence,
continuity: self.continuity,
locally_observed_at: now,
});
}
pub fn expire_if_due(&mut self, now: Instant) {
if self.continuity != Continuity::Expired && now >= self.deadline {
self.continuity = Continuity::Expired;
if let Some(obs) = &mut self.observation {
obs.continuity = Continuity::Expired;
}
}
}
pub fn disrupt(&mut self, reason: DisruptReason) {
self.continuity = Continuity::Expired;
self.last_disrupt = Some(reason);
if let Some(obs) = &mut self.observation {
obs.continuity = Continuity::Expired;
}
}
pub fn projected(&self) -> ProjectedReadiness {
match &self.observation {
None => ProjectedReadiness::Unknown,
Some(obs) => project(obs.attested_status, self.continuity),
}
}
pub const fn continuity(&self) -> Continuity {
self.continuity
}
pub fn observation(&self) -> Option<&ReadinessObservation> {
self.observation.as_ref()
}
pub const fn last_disrupt(&self) -> Option<DisruptReason> {
self.last_disrupt
}
}
#[cfg(test)]
mod tests {
use super::*;
const K: u32 = 3;
const D: Duration = Duration::from_millis(100);
fn beat(seq: u64, status: AttestedStatus, bearing: bool) -> DeliveredBeat {
DeliveredBeat {
attested_status: status,
estimated_start: None,
source_incarnation: Incarnation::new(1),
capability_generation: 4,
seq,
promised_cadence: Duration::from_millis(100),
continuity_bearing: bearing,
}
}
#[test]
fn interval_update_reanchors_deadline_without_resetting_continuity() {
let t0 = Instant::now();
let mut cell = ObservationCell::register(t0, Duration::from_millis(200), K);
let t1 = t0 + Duration::from_millis(50);
cell.on_admitted_beat(t1, beat(1, AttestedStatus::Ready, true));
assert_eq!(cell.continuity(), Continuity::Established);
cell.update_interval(Duration::from_millis(100));
assert_eq!(cell.continuity(), Continuity::Established);
cell.expire_if_due(t1 + Duration::from_millis(299));
assert_eq!(cell.continuity(), Continuity::Established);
cell.expire_if_due(t1 + Duration::from_millis(301));
assert_eq!(
cell.continuity(),
Continuity::Expired,
"the tightened window expires at the re-anchored deadline",
);
let mut cell = ObservationCell::register(t0, Duration::from_millis(200), K);
cell.on_admitted_beat(t1, beat(1, AttestedStatus::Ready, true));
cell.update_interval(Duration::from_millis(400));
cell.expire_if_due(t1 + Duration::from_millis(1199));
assert_eq!(
cell.continuity(),
Continuity::Established,
"the loosened window holds past the old deadline",
);
cell.expire_if_due(t1 + Duration::from_millis(1201));
assert_eq!(cell.continuity(), Continuity::Expired);
}
#[test]
fn interval_update_reanchors_an_unestablished_deadline_by_own_d_not_promised() {
let t0 = Instant::now();
let mut cell = ObservationCell::register(t0, D, K);
let mut warm = beat(1, AttestedStatus::NotReady, false);
warm.promised_cadence = Duration::from_secs(1);
cell.on_admitted_beat(t0, warm);
assert_eq!(cell.continuity(), Continuity::Unestablished);
cell.update_interval(Duration::from_millis(200));
cell.expire_if_due(t0 + Duration::from_millis(599));
assert_eq!(
cell.continuity(),
Continuity::Unestablished,
"the loosened establishment deadline has not fired yet",
);
cell.expire_if_due(t0 + Duration::from_millis(600));
assert_eq!(
cell.continuity(),
Continuity::Expired,
"the establishment deadline re-anchored to own D × factor",
);
}
#[test]
fn generation_change_starts_a_fresh_observation() {
let t0 = Instant::now();
let mut cell = ObservationCell::register(t0, D, K);
cell.on_admitted_beat(t0, beat(10, AttestedStatus::Ready, true));
assert_eq!(cell.projected(), ProjectedReadiness::Ready);
let mut regen = beat(11, AttestedStatus::Ready, false);
regen.capability_generation = 5;
cell.on_admitted_beat(t0 + D, regen);
assert_eq!(cell.continuity(), Continuity::Unestablished);
assert_eq!(cell.projected(), ProjectedReadiness::Unknown);
assert_eq!(cell.last_disrupt(), Some(DisruptReason::GenerationChanged));
cell.expire_if_due(t0 + D + D * K);
assert_eq!(cell.continuity(), Continuity::Expired);
let mut live = beat(12, AttestedStatus::Ready, true);
live.capability_generation = 5;
cell.on_admitted_beat(t0 + D * 5, live);
assert_eq!(cell.projected(), ProjectedReadiness::Ready);
let mut pessimist = ObservationCell::register(t0, D, K);
pessimist.on_admitted_beat(t0, beat(1, AttestedStatus::NotReady, true));
let mut regen_nr = beat(2, AttestedStatus::NotReady, false);
regen_nr.capability_generation = 5;
pessimist.on_admitted_beat(t0 + D, regen_nr);
assert_eq!(pessimist.projected(), ProjectedReadiness::NotReady);
}
#[test]
fn live_beat_under_a_new_generation_establishes_immediately() {
let t0 = Instant::now();
let mut cell = ObservationCell::register(t0, D, K);
cell.on_admitted_beat(t0, beat(10, AttestedStatus::Ready, true));
assert_eq!(cell.projected(), ProjectedReadiness::Ready);
let mut live_regen = beat(11, AttestedStatus::Ready, true);
live_regen.capability_generation = 5;
cell.on_admitted_beat(t0 + D, live_regen);
assert_eq!(
cell.continuity(),
Continuity::Established,
"a live beat under the new generation IS the establishment",
);
assert_eq!(cell.projected(), ProjectedReadiness::Ready);
assert_eq!(cell.last_disrupt(), None);
assert_eq!(
cell.observation().unwrap().capability_generation,
5,
"the observation tracks the new generation",
);
cell.expire_if_due(t0 + D + Duration::from_millis(299));
assert_eq!(cell.projected(), ProjectedReadiness::Ready);
cell.expire_if_due(t0 + D + Duration::from_millis(300));
assert_eq!(cell.projected(), ProjectedReadiness::Unknown);
}
#[test]
fn projection_table_is_pinned_exactly() {
use AttestedStatus::*;
use Continuity::*;
use ProjectedReadiness as P;
let table = [
(Ready, Unestablished, P::Unknown), (Ready, Established, P::Ready),
(Ready, Expired, P::Unknown),
(NotReady, Unestablished, P::NotReady), (NotReady, Established, P::NotReady),
(NotReady, Expired, P::Unknown),
(ProviderUnknown, Unestablished, P::Unknown),
(ProviderUnknown, Established, P::Unknown),
(ProviderUnknown, Expired, P::Unknown),
];
for (status, continuity, expected) in table {
assert_eq!(
project(status, continuity),
expected,
"project({status:?}, {continuity:?})",
);
}
}
#[test]
fn registration_starts_unestablished_and_expires_at_the_deadline() {
let t0 = Instant::now();
let mut cell = ObservationCell::register(t0, D, K);
assert_eq!(cell.continuity(), Continuity::Unestablished);
assert_eq!(cell.projected(), ProjectedReadiness::Unknown);
cell.expire_if_due(t0 + D * K - Duration::from_millis(1));
assert_eq!(cell.continuity(), Continuity::Unestablished);
cell.expire_if_due(t0 + D * K);
assert_eq!(cell.continuity(), Continuity::Expired);
assert_eq!(cell.projected(), ProjectedReadiness::Unknown);
}
#[test]
fn warm_start_ready_projects_unknown_but_notready_projects_immediately() {
let t0 = Instant::now();
let mut ready_cell = ObservationCell::register(t0, D, K);
ready_cell.on_admitted_beat(t0, beat(100, AttestedStatus::Ready, false));
assert_eq!(ready_cell.continuity(), Continuity::Unestablished);
assert_eq!(ready_cell.projected(), ProjectedReadiness::Unknown);
let mut notready_cell = ObservationCell::register(t0, D, K);
notready_cell.on_admitted_beat(t0, beat(100, AttestedStatus::NotReady, false));
assert_eq!(notready_cell.projected(), ProjectedReadiness::NotReady);
}
#[test]
fn continuity_bearing_beat_establishes_and_ready_projects() {
let t0 = Instant::now();
let mut cell = ObservationCell::register(t0, D, K);
cell.on_admitted_beat(t0, beat(100, AttestedStatus::Ready, false));
cell.on_admitted_beat(t0 + D, beat(101, AttestedStatus::Ready, true));
assert_eq!(cell.continuity(), Continuity::Established);
assert_eq!(cell.projected(), ProjectedReadiness::Ready);
}
#[test]
fn cached_newer_beats_never_extend_the_establishment_deadline() {
let t0 = Instant::now();
let mut cell = ObservationCell::register(t0, D, K);
cell.on_admitted_beat(t0, beat(100, AttestedStatus::NotReady, false));
assert_eq!(cell.projected(), ProjectedReadiness::NotReady);
cell.on_admitted_beat(t0 + D, beat(101, AttestedStatus::NotReady, false));
cell.on_admitted_beat(t0 + D * 2, beat(102, AttestedStatus::NotReady, false));
cell.expire_if_due(t0 + D * K);
assert_eq!(cell.continuity(), Continuity::Expired);
assert_eq!(cell.projected(), ProjectedReadiness::Unknown);
}
#[test]
fn established_expires_on_silence_and_a_live_beat_recovers() {
let t0 = Instant::now();
let mut cell = ObservationCell::register(t0, D, K);
cell.on_admitted_beat(t0, beat(1, AttestedStatus::Ready, true));
cell.expire_if_due(t0 + Duration::from_millis(299));
assert_eq!(cell.projected(), ProjectedReadiness::Ready);
cell.expire_if_due(t0 + Duration::from_millis(300));
assert_eq!(cell.continuity(), Continuity::Expired);
assert_eq!(cell.projected(), ProjectedReadiness::Unknown);
let t1 = t0 + Duration::from_millis(400);
cell.on_admitted_beat(t1, beat(2, AttestedStatus::Ready, false));
assert_eq!(cell.projected(), ProjectedReadiness::Unknown);
cell.on_admitted_beat(t1 + D, beat(3, AttestedStatus::Ready, true));
assert_eq!(cell.projected(), ProjectedReadiness::Ready);
}
#[test]
fn window_is_k_times_max_of_promised_cadence_and_own_interval() {
let t0 = Instant::now();
let own_d = Duration::from_millis(500);
let mut cell = ObservationCell::register(t0, own_d, K);
cell.on_admitted_beat(t0, beat(1, AttestedStatus::Ready, true));
cell.expire_if_due(t0 + Duration::from_millis(1499));
assert_eq!(cell.projected(), ProjectedReadiness::Ready);
cell.expire_if_due(t0 + Duration::from_millis(1500));
assert_eq!(cell.projected(), ProjectedReadiness::Unknown);
let mut slow = ObservationCell::register(t0, D, K);
let mut b = beat(1, AttestedStatus::Ready, true);
b.promised_cadence = Duration::from_secs(1);
slow.on_admitted_beat(t0, b);
slow.expire_if_due(t0 + Duration::from_millis(2999));
assert_eq!(slow.projected(), ProjectedReadiness::Ready);
slow.expire_if_due(t0 + Duration::from_secs(3));
assert_eq!(slow.projected(), ProjectedReadiness::Unknown);
}
#[test]
fn continuity_never_carries_across_incarnations() {
let t0 = Instant::now();
let mut cell = ObservationCell::register(t0, D, K);
cell.on_admitted_beat(t0, beat(50, AttestedStatus::Ready, true));
assert_eq!(cell.projected(), ProjectedReadiness::Ready);
let mut restarted = beat(1, AttestedStatus::Ready, false);
restarted.source_incarnation = Incarnation::new(2);
cell.on_admitted_beat(t0 + D, restarted);
assert_eq!(cell.continuity(), Continuity::Expired);
assert_eq!(cell.projected(), ProjectedReadiness::Unknown);
assert_eq!(
cell.last_disrupt(),
Some(DisruptReason::IncarnationSuperseded),
);
let mut live = beat(2, AttestedStatus::Ready, true);
live.source_incarnation = Incarnation::new(2);
cell.on_admitted_beat(t0 + D * 2, live);
assert_eq!(cell.projected(), ProjectedReadiness::Ready);
}
#[test]
fn incarnation_crossing_expires_only_established_optimism() {
let t0 = Instant::now();
let mut cell = ObservationCell::register(t0, D, K);
cell.on_admitted_beat(t0, beat(1, AttestedStatus::NotReady, false));
assert_eq!(cell.continuity(), Continuity::Unestablished);
assert_eq!(cell.projected(), ProjectedReadiness::NotReady);
let mut restarted = beat(2, AttestedStatus::NotReady, false);
restarted.source_incarnation = Incarnation::new(2);
cell.on_admitted_beat(t0 + D, restarted);
assert_eq!(
cell.continuity(),
Continuity::Unestablished,
"a new-incarnation warm-start does not force-expire an unestablished cell",
);
assert_eq!(
cell.projected(),
ProjectedReadiness::NotReady,
"cached pessimism still projects — pessimism is safe across the boundary",
);
assert_eq!(
cell.last_disrupt(),
None,
"no established optimism was revoked",
);
cell.expire_if_due(t0 + D * K);
assert_eq!(cell.continuity(), Continuity::Expired);
assert_eq!(cell.projected(), ProjectedReadiness::Unknown);
let mut established = ObservationCell::register(t0, D, K);
established.on_admitted_beat(t0, beat(1, AttestedStatus::Ready, true));
assert_eq!(established.projected(), ProjectedReadiness::Ready);
let mut warm_new_inc = beat(2, AttestedStatus::Ready, false);
warm_new_inc.source_incarnation = Incarnation::new(2);
established.on_admitted_beat(t0 + D, warm_new_inc);
assert_eq!(established.continuity(), Continuity::Expired);
assert_eq!(
established.last_disrupt(),
Some(DisruptReason::IncarnationSuperseded),
);
}
#[test]
fn disrupt_expires_immediately_with_reason() {
let t0 = Instant::now();
let mut cell = ObservationCell::register(t0, D, K);
cell.on_admitted_beat(t0, beat(1, AttestedStatus::Ready, true));
cell.disrupt(DisruptReason::PathFailed);
assert_eq!(cell.continuity(), Continuity::Expired);
assert_eq!(cell.projected(), ProjectedReadiness::Unknown);
assert_eq!(cell.last_disrupt(), Some(DisruptReason::PathFailed));
let obs = cell.observation().unwrap();
assert_eq!(obs.continuity, Continuity::Expired);
}
}