use crate::discipline::{ChangeVerdict, ClockCommand, Discipline, DisciplineConfig, Plan};
use crate::filter::{Sample, SampleRegister};
pub const REGISTER_CAPACITY: usize = 64;
#[derive(Clone, Copy, Debug)]
pub struct Estimate {
pub offset_s: f64,
pub freq_ppm: Option<f64>,
pub sd_s: f64,
pub samples: usize,
}
pub struct MultiController {
registers: Vec<SampleRegister>,
discipline: Discipline,
freq_cmd_ppm: f64,
drain_ppm: f64,
drain_remaining_s: f64,
last_plan_mono_s: Option<f64>,
drain_booked_until: Option<f64>,
unapplied: Option<Unapplied>,
}
#[derive(Clone, Copy, Debug)]
struct Unapplied {
mono_s: f64,
dfreq_ppm: f64,
step_s: f64,
freq_cmd_before: f64,
drain_ppm_before: f64,
drain_remaining_before: f64,
}
pub struct ControllerStep {
pub plan: Plan,
pub applied_ppm: f64,
pub estimate_offset_s: f64,
pub estimate_freq_ppm: Option<f64>,
pub estimate_sd_s: f64,
pub samples_used: usize,
pub verdict: ChangeVerdict,
}
fn new_register(config: &DisciplineConfig) -> SampleRegister {
let mut r = SampleRegister::new(REGISTER_CAPACITY);
r.set_weight_floor_ratio(config.weight_floor_ratio);
r.set_offset_weight_floor_ratio(config.offset_weight_floor_ratio);
r.set_offset_age_halflife_s(config.offset_age_halflife_s);
r.set_offset_weight_dispersion_k(config.offset_weight_dispersion_k);
r.set_slope_density_weighting(config.slope_density_weighting);
r.set_adaptive_window(config.adaptive_window);
r
}
impl MultiController {
pub fn new(config: DisciplineConfig, sources: usize) -> Self {
MultiController {
registers: (0..sources.max(1)).map(|_| new_register(&config)).collect(),
discipline: Discipline::new(config),
freq_cmd_ppm: 0.0,
drain_ppm: 0.0,
drain_remaining_s: 0.0,
last_plan_mono_s: None,
drain_booked_until: None,
unapplied: None,
}
}
pub fn revert_last_plan(&mut self) -> bool {
let Some(u) = self.unapplied.take() else {
return false;
};
for r in &mut self.registers {
r.slew_samples(u.mono_s, -u.dfreq_ppm, -u.step_s);
}
self.freq_cmd_ppm = u.freq_cmd_before;
self.drain_ppm = u.drain_ppm_before;
self.drain_remaining_s = u.drain_remaining_before;
true
}
pub fn confirm_last_plan(&mut self) {
self.unapplied = None;
}
pub fn freq_ppm(&self) -> f64 {
self.freq_cmd_ppm
}
pub fn drain_ppm(&self) -> f64 {
self.drain_ppm
}
pub fn applied_ppm(&self) -> f64 {
self.freq_cmd_ppm + self.drain_ppm
}
pub fn poll_log2(&self) -> i8 {
self.discipline.poll_log2()
}
pub fn poll_interval_s(&self) -> f64 {
(2.0f64).powi(self.discipline.poll_log2() as i32)
}
pub fn samples(&self) -> usize {
self.registers[0].len()
}
pub fn samples_from(&self, index: usize) -> usize {
self.registers[index].len()
}
pub fn sources(&self) -> usize {
self.registers.len()
}
pub fn preload_frequency(&mut self, freq_ppm: f64) {
self.freq_cmd_ppm = freq_ppm;
}
pub fn retry_interval_s(&self) -> f64 {
self.discipline.retry_interval_s()
}
pub fn drain_completes_at(&self) -> Option<f64> {
let last = self.last_plan_mono_s?;
if self.drain_ppm == 0.0 || self.drain_remaining_s <= 0.0 {
return None;
}
Some(last + self.drain_remaining_s / (self.drain_ppm.abs() * 1e-6))
}
pub fn poll_drain(&mut self, mono_now_s: f64) -> Option<ClockCommand> {
let completes_at = self.drain_completes_at()?;
if mono_now_s < completes_at {
return None;
}
let last = self.drain_booked_until?;
let consumed = self.drain_ppm * 1e-6 * (mono_now_s - last).max(0.0);
if consumed != 0.0 {
for r in &mut self.registers {
r.slew_samples(mono_now_s, 0.0, consumed);
}
}
self.drain_ppm = 0.0;
self.drain_remaining_s = 0.0;
self.last_plan_mono_s = Some(mono_now_s);
self.drain_booked_until = Some(mono_now_s);
Some(ClockCommand::Slew {
freq_ppm: self.freq_cmd_ppm,
drain_offset: 0.0,
drain_rate_ppm: 0.0,
})
}
fn settle_drain(&mut self, mono_now_s: f64) {
let drained = match self.drain_booked_until {
Some(last) if self.drain_ppm != 0.0 => {
self.drain_ppm * 1e-6 * (mono_now_s - last).max(0.0)
}
_ => 0.0,
};
if drained != 0.0 {
for r in &mut self.registers {
r.slew_samples(mono_now_s, 0.0, drained);
}
self.drain_remaining_s = (self.drain_remaining_s - drained.abs()).max(0.0);
}
if self.drain_booked_until.is_some() {
self.drain_booked_until = Some(mono_now_s);
}
}
pub fn observe(&mut self, index: usize, mono_now_s: f64, sample: Sample) -> Estimate {
self.settle_drain(mono_now_s);
self.registers[index].push(sample);
self.estimate(index, mono_now_s)
}
pub fn estimate(&mut self, index: usize, mono_now_s: f64) -> Estimate {
let samples = self.registers[index].len();
match self.registers[index].regress(mono_now_s) {
Some(e) => Estimate {
offset_s: e.offset,
freq_ppm: e.freq_ppm,
sd_s: e.offset_sd.max(1e-7),
samples,
},
None => match self.registers[index].best() {
Some(best) => Estimate {
offset_s: best.offset,
freq_ppm: None,
sd_s: (best.delay / 2.0).max(1e-7),
samples,
},
None => Estimate {
offset_s: 0.0,
freq_ppm: None,
sd_s: 1e-3,
samples,
},
},
}
}
pub fn steer(&mut self, est: Estimate, mono_now_s: f64, leap_pending: bool) -> ControllerStep {
let (offset, freq, sd) = (est.offset_s, est.freq_ppm, est.sd_s);
let plan = self
.discipline
.on_estimate_with_leap(offset, freq, sd, leap_pending);
let freq_cmd_new = self.discipline.freq_ppm();
let freq_cmd_before = self.freq_cmd_ppm;
let drain_ppm_before = self.drain_ppm;
let drain_remaining_before = self.drain_remaining_s;
let dfreq_ppm = freq_cmd_new - self.freq_cmd_ppm;
let step_s = match plan.command {
ClockCommand::Step { add_seconds } => add_seconds,
ClockCommand::Slew { .. } => 0.0,
};
match plan.command {
ClockCommand::Step { .. } => {
self.drain_ppm = 0.0;
self.drain_remaining_s = 0.0;
}
ClockCommand::Slew {
drain_offset,
drain_rate_ppm,
..
} => {
self.drain_ppm = drain_rate_ppm.copysign(drain_offset);
self.drain_remaining_s = drain_offset.abs();
}
}
for r in &mut self.registers {
r.slew_samples(mono_now_s, dfreq_ppm, step_s);
}
self.freq_cmd_ppm = freq_cmd_new;
self.last_plan_mono_s = Some(mono_now_s);
self.drain_booked_until = Some(mono_now_s);
self.unapplied = Some(Unapplied {
mono_s: mono_now_s,
dfreq_ppm,
step_s,
freq_cmd_before,
drain_ppm_before,
drain_remaining_before,
});
ControllerStep {
plan,
applied_ppm: self.freq_cmd_ppm + self.drain_ppm,
estimate_offset_s: offset,
estimate_freq_ppm: freq,
samples_used: est.samples,
estimate_sd_s: sd,
verdict: plan.verdict,
}
}
}
pub struct SyncController {
inner: MultiController,
}
impl SyncController {
pub fn new(config: DisciplineConfig) -> Self {
SyncController {
inner: MultiController::new(config, 1),
}
}
pub fn on_sample(&mut self, mono_now_s: f64, sample: Sample) -> ControllerStep {
self.on_sample_with_leap(mono_now_s, sample, false)
}
pub fn on_sample_with_leap(
&mut self,
mono_now_s: f64,
sample: Sample,
leap_pending: bool,
) -> ControllerStep {
let est = self.inner.observe(0, mono_now_s, sample);
self.inner.steer(est, mono_now_s, leap_pending)
}
pub fn revert_last_plan(&mut self) -> bool {
self.inner.revert_last_plan()
}
pub fn confirm_last_plan(&mut self) {
self.inner.confirm_last_plan()
}
pub fn freq_ppm(&self) -> f64 {
self.inner.freq_ppm()
}
pub fn drain_ppm(&self) -> f64 {
self.inner.drain_ppm()
}
pub fn applied_ppm(&self) -> f64 {
self.inner.applied_ppm()
}
pub fn poll_log2(&self) -> i8 {
self.inner.poll_log2()
}
pub fn samples(&self) -> usize {
self.inner.samples()
}
pub fn preload_frequency(&mut self, freq_ppm: f64) {
self.inner.preload_frequency(freq_ppm)
}
pub fn retry_interval_s(&self) -> f64 {
self.inner.retry_interval_s()
}
pub fn drain_completes_at(&self) -> Option<f64> {
self.inner.drain_completes_at()
}
pub fn poll_drain(&mut self, mono_now_s: f64) -> Option<ClockCommand> {
self.inner.poll_drain(mono_now_s)
}
}
#[cfg(test)]
mod refusal_tests {
use super::*;
use crate::filter::Sample;
fn feed(c: &mut SyncController, n: usize, base: f64) {
for i in 0..n {
let t = 16.0 * (i as f64 + 1.0);
c.on_sample(
t,
Sample {
t,
offset: base - 20e-6 * t,
delay: 200e-6,
dispersion: 1e-6,
},
);
}
}
#[test]
fn a_refused_command_leaves_no_trace() {
let cfg = DisciplineConfig::default();
let mut applied = SyncController::new(cfg);
let mut refused = SyncController::new(cfg);
feed(&mut applied, 12, 0.010);
feed(&mut refused, 12, 0.010);
let t = 16.0 * 13.0;
let sample = Sample {
t,
offset: 0.010 - 20e-6 * t,
delay: 200e-6,
dispersion: 1e-6,
};
let before_freq = refused.freq_ppm();
let before_drain = refused.drain_ppm();
applied.on_sample(t, sample);
applied.confirm_last_plan();
refused.on_sample(t, sample);
assert!(refused.revert_last_plan(), "there was a plan to revert");
assert_eq!(
refused.freq_ppm(),
before_freq,
"the frequency command survived a refusal"
);
assert_eq!(
refused.drain_ppm(),
before_drain,
"the drain survived a refusal"
);
let t2 = 16.0 * 14.0;
let next = Sample {
t: t2,
offset: 0.010 - 20e-6 * t2,
delay: 200e-6,
dispersion: 1e-6,
};
let a = applied.on_sample(t2, next);
let r = refused.on_sample(t2, next);
assert_ne!(
a.applied_ppm, r.applied_ppm,
"a reverted controller behaved identically to one that applied its command, so the revert did not actually restore the books"
);
}
#[test]
fn reverting_nothing_is_a_no_op() {
let mut c = SyncController::new(DisciplineConfig::default());
assert!(!c.revert_last_plan(), "nothing has been planned yet");
feed(&mut c, 6, 0.001);
let freq = c.freq_ppm();
assert!(c.revert_last_plan());
assert!(!c.revert_last_plan(), "a second revert must do nothing");
assert_ne!(freq, f64::NAN);
}
}
#[cfg(test)]
mod tests {
use super::*;
fn closed_loop(base_freq_ppm: f64, initial_err_s: f64, polls: usize) -> f64 {
let mut controller = SyncController::new(DisciplineConfig {
iburst: false,
makestep_threshold: None,
..DisciplineConfig::default()
});
let mut err = initial_err_s;
let mut t = 0.0f64;
for _ in 0..polls {
let mono = t + err;
let step = controller.on_sample(
mono,
Sample {
t: mono,
offset: -err,
delay: 0.0002,
dispersion: 0.0,
},
);
let dt = step.plan.next_poll_s;
err += (base_freq_ppm + step.applied_ppm) * 1e-6 * dt;
t += dt;
}
err
}
#[test]
fn the_controller_converges_on_a_drifting_clock() {
let residual = closed_loop(40.0, 0.030, 80);
assert!(
residual.abs() < 1e-4,
"did not converge: {residual} s remaining"
);
}
#[test]
fn it_converges_from_either_direction() {
for (drift, start) in [
(40.0, 0.030),
(-40.0, -0.030),
(100.0, -0.050),
(-15.0, 0.010),
] {
let residual = closed_loop(drift, start, 120);
assert!(
residual.abs() < 1e-3,
"drift {drift} ppm from {start} s left {residual} s"
);
}
}
#[test]
fn frequency_and_drain_stay_separate() {
let mut controller = SyncController::new(DisciplineConfig {
iburst: false,
makestep_threshold: None,
..DisciplineConfig::default()
});
let step = controller.on_sample(
0.0,
Sample {
t: 0.0,
offset: 0.001,
delay: 0.0002,
dispersion: 0.0,
},
);
assert!(
controller.drain_ppm() != 0.0,
"a non-zero offset should start a drain"
);
assert_eq!(
step.applied_ppm,
controller.freq_ppm() + controller.drain_ppm(),
"the applied total must be exactly the two parts"
);
}
#[test]
fn a_failed_exchange_retries_at_the_burst_spacing_not_the_poll_interval() {
let controller = SyncController::new(DisciplineConfig {
min_poll: 4, iburst: true,
..DisciplineConfig::default()
});
let retry = controller.retry_interval_s();
assert!(
retry <= 4.0,
"a cold-start retry waited {retry} s; the burst spacing is the point"
);
}
#[test]
fn once_the_burst_is_spent_retries_use_the_poll_interval() {
let mut controller = SyncController::new(DisciplineConfig {
min_poll: 4,
iburst: true,
makestep_threshold: None,
..DisciplineConfig::default()
});
for i in 0..8 {
controller.on_sample(
i as f64 * 2.0,
Sample {
t: i as f64 * 2.0,
offset: 1e-6,
delay: 0.0002,
dispersion: 0.0,
},
);
}
assert!(
controller.retry_interval_s() >= 16.0,
"after the burst, retries must back off to the poll interval"
);
}
#[test]
fn a_drain_ends_when_its_budget_is_spent() {
let mut controller = SyncController::new(DisciplineConfig {
iburst: false,
makestep_threshold: None,
..DisciplineConfig::default()
});
controller.on_sample(
0.0,
Sample {
t: 0.0,
offset: 0.010,
delay: 0.0002,
dispersion: 0.0,
},
);
let ends = controller
.drain_completes_at()
.expect("a drain should be running");
assert!(ends > 0.0, "drain has no completion time");
assert!(controller.poll_drain(ends - 1e-6).is_none());
assert!(controller.applied_ppm() != 0.0);
assert!(controller.poll_drain(ends).is_some());
assert!(controller.poll_drain(ends + 1.0).is_none());
assert_eq!(
controller.drain_ppm(),
0.0,
"a spent drain must stop slewing the clock"
);
}
#[test]
fn a_late_retirement_books_what_the_clock_actually_received() {
let mut controller = SyncController::new(DisciplineConfig {
iburst: false,
makestep_threshold: None,
..DisciplineConfig::default()
});
controller.on_sample(
0.0,
Sample {
t: 0.0,
offset: 0.010,
delay: 0.0002,
dispersion: 0.0,
},
);
let rate = controller.drain_ppm();
let ends = controller.drain_completes_at().expect("a drain");
let late = 0.5;
controller.poll_drain(ends + late);
let step = controller.on_sample(
ends + late,
Sample {
t: ends + late,
offset: 0.0,
delay: 0.0002,
dispersion: 0.0,
},
);
let overrun = rate.abs() * 1e-6 * late;
assert!(
overrun > 1e-6,
"test is vacuous unless the overrun is meaningful"
);
assert!(
step.estimate_offset_s.abs() < 0.010,
"late retirement lost correction the clock had already received: estimate {} s",
step.estimate_offset_s
);
}
#[test]
fn a_preloaded_frequency_is_the_starting_point() {
let mut controller = SyncController::new(DisciplineConfig::default());
controller.preload_frequency(-12.5);
assert!((controller.freq_ppm() + 12.5).abs() < 1e-12);
}
}