extern crate alloc;
use alloc::boxed::Box;
use core::cell::Cell;
use super::super::stamp::{Kairos, ToU16};
use super::super::time::TimeSource;
use super::error::{InvalidStationId, SkewExceeded};
use super::fold;
use super::run::KairosRun;
use super::state::ClockState;
use super::{ClockConfig, ClockStats};
pub struct LocalClock<TS: TimeSource = Box<dyn TimeSource>> {
time_source: TS,
station_id: u32,
state: Cell<u128>,
max_forward_skew: Cell<u64>,
max_backward_skew: Cell<u64>,
max_skew_warning_threshold: u64,
}
impl<TS: TimeSource> LocalClock<TS> {
pub fn new(time_source: TS, station_id: u32) -> Result<Self, InvalidStationId> {
Self::with_warning_threshold(
time_source,
station_id,
ClockConfig::default().max_skew_warning_threshold,
)
}
pub fn with_warning_threshold(
time_source: TS,
station_id: u32,
max_skew_warning_threshold: u64,
) -> Result<Self, InvalidStationId> {
if station_id == 0 {
return Err(InvalidStationId);
}
Ok(Self {
time_source,
station_id,
state: Cell::new(ClockState::UNINIT.0),
max_forward_skew: Cell::new(0),
max_backward_skew: Cell::new(0),
max_skew_warning_threshold,
})
}
#[must_use]
pub const fn station_id(&self) -> u32 {
self.station_id
}
pub fn now<T: ToU16>(&self, kairotic: T) -> Kairos {
let (physical, logical) = self.advance(None);
Kairos::new(physical, logical, self.station_id, kairotic)
}
pub fn after<T: ToU16>(&self, previous: Kairos, kairotic: T) -> Kairos {
let (physical, logical) = self.advance(Some(previous));
Kairos::new(physical, logical, self.station_id, kairotic)
}
pub fn now_run<T: ToU16>(&self, len: core::num::NonZeroU32, kairotic: T) -> KairosRun {
let advance = fold::advance_run(
ClockState(self.state.get()).last(),
self.time_source.now(),
len.get(),
);
if advance.forward_skew > self.max_forward_skew.get() {
self.max_forward_skew.set(advance.forward_skew);
}
if advance.backward_skew > self.max_backward_skew.get() {
self.max_backward_skew.set(advance.backward_skew);
}
self.state
.set(ClockState::minted(advance.physical, advance.logical).0);
KairosRun::new(
Kairos::new(
advance.first_physical,
advance.first_logical,
self.station_id,
kairotic,
),
len.get(),
)
}
pub fn observe(&self, remote: Kairos) {
let _ = self.advance(Some(remote));
}
pub fn try_observe(&self, remote: Kairos, max_forward_skew: u64) -> Result<(), SkewExceeded> {
let current_physical = self.time_source.now();
let observed_forward_skew = remote.physical().saturating_sub(current_physical);
if observed_forward_skew > max_forward_skew {
return Err(SkewExceeded {
observed_forward_skew,
max_forward_skew,
});
}
let _ = self.advance_from(current_physical, Some(remote));
Ok(())
}
#[must_use]
pub fn stats(&self) -> ClockStats {
let attempts = u32::from(ClockState(self.state.get()).last().is_some());
ClockStats {
max_forward_skew_ns: self.max_forward_skew.get(),
max_backward_skew_ns: self.max_backward_skew.get(),
skew_warning_threshold_ns: self.max_skew_warning_threshold,
peak_attempts: attempts,
}
}
fn advance(&self, previous: Option<Kairos>) -> (u64, u16) {
self.advance_from(self.time_source.now(), previous)
}
fn advance_from(&self, current_physical: u64, previous: Option<Kairos>) -> (u64, u16) {
let remote = previous.map(|prev| (prev.physical(), prev.logical()));
let advance = fold::advance(
ClockState(self.state.get()).last(),
current_physical,
remote,
);
if advance.forward_skew > self.max_forward_skew.get() {
self.max_forward_skew.set(advance.forward_skew);
}
if advance.backward_skew > self.max_backward_skew.get() {
self.max_backward_skew.set(advance.backward_skew);
}
self.state
.set(ClockState::minted(advance.physical, advance.logical).0);
(advance.physical, advance.logical)
}
}
pub type DynLocalClock = LocalClock<Box<dyn TimeSource>>;