mod selector;
pub use selector::{Selector, SyncParameters};
pub use source::SourceInfo;
pub mod ff;
mod ring_buffer;
use std::{net::SocketAddr, sync::Arc};
pub use ring_buffer::RingBuffer;
use crate::daemon::{
clock_parameters::ClockParameters, event, logging, receiver_stream::RoutableEvent,
selected_clock::SelectedClockSource,
};
pub mod source;
#[derive(Debug, Clone, bon::Builder)]
pub struct ClockSyncAlgorithm {
amazon_time_sync: Option<source::AmazonTimeSync>,
#[builder(default)]
ntp_sources: Vec<source::NtpSource>,
phc: Option<source::Phc>,
selected_clock: Arc<SelectedClockSource>,
selector: Selector,
#[builder(default)]
seq: u64,
}
impl ClockSyncAlgorithm {
pub fn set_amazon_time_sync(&mut self, source: source::AmazonTimeSync) {
assert!(self.amazon_time_sync.is_none());
self.amazon_time_sync = Some(source);
}
pub fn amazon_time_sync(&self) -> Option<&source::AmazonTimeSync> {
self.amazon_time_sync.as_ref()
}
pub fn add_ntp_source(&mut self, source: source::NtpSource) {
assert!(
self.ntp_sources
.iter()
.find(|s| s.socket_address() == source.socket_address())
.is_none(),
"duplicate addr"
);
self.ntp_sources.push(source);
}
pub fn remove_ntp_source(&mut self, addr: &SocketAddr) {
let removed_index = self
.ntp_sources
.iter()
.position(|s| s.socket_address() == *addr)
.unwrap();
self.ntp_sources.swap_remove(removed_index);
}
pub fn ntp_sources(&self) -> &[source::NtpSource] {
&self.ntp_sources
}
pub fn set_phc(&mut self, phc: source::Phc) {
assert!(self.phc.is_none());
self.phc = Some(phc);
}
pub fn phc(&self) -> Option<&source::Phc> {
self.phc.as_ref()
}
pub fn init_repro(&mut self) {
logging::ffevents::log_init(self.next_seq());
}
pub fn feed(&mut self, routable_event: RoutableEvent) -> Option<&SyncParameters> {
self.feed_repro(routable_event)
}
pub fn current_source_params(&self) -> Option<&SyncParameters> {
self.selector.current()
}
fn next_seq(&mut self) -> u64 {
let seq = self.seq;
self.seq += 1;
seq
}
fn feed_inner(&mut self, routable_event: RoutableEvent) -> Option<&SyncParameters> {
let alg_output = match routable_event {
RoutableEvent::AmazonTimeSync(event) => {
let amazon_time_sync = self.amazon_time_sync.as_mut()?;
Self::feed_amazon_time_sync(amazon_time_sync, event)
}
RoutableEvent::NtpSource(sender_address, event) => {
Self::feed_ntp_source(&mut self.ntp_sources, sender_address, event)
}
RoutableEvent::Phc(device_path, event) => {
let phc = self.phc.as_mut().unwrap();
Self::feed_phc(phc, device_path, event)
}
};
let (clock_parameters, source_info) = alg_output?;
let output = self.selector.update(clock_parameters, source_info.clone());
if output.is_some() {
Self::update_selected_clock(&self.selected_clock, &source_info);
}
output
}
#[expect(
clippy::needless_pass_by_value,
reason = "serializing event before logging is unergonomic"
)]
fn feed_repro(&mut self, routable_event: RoutableEvent) -> Option<&SyncParameters> {
let seq = self.next_seq();
let output = self.feed_inner(routable_event.clone());
logging::ffevents::log_feed(&routable_event, output.map(|sp| &sp.clock_parameters), seq);
output
}
fn feed_amazon_time_sync(
amazon_time_sync: &mut source::AmazonTimeSync,
event: event::Ntp,
) -> Option<(&ClockParameters, SourceInfo)> {
let stratum = event.data().stratum;
let address = SocketAddr::from(amazon_time_sync.source_address());
amazon_time_sync
.feed(event)
.map(|params| (params, SourceInfo::AmazonTimeSync(address, stratum)))
}
fn feed_ntp_source(
ntp_sources: &mut [source::NtpSource],
sender_address: SocketAddr,
event: event::Ntp,
) -> Option<(&ClockParameters, SourceInfo)> {
let stratum = event.data().stratum;
ntp_sources
.iter_mut()
.find(|source| source.socket_address() == sender_address)
.and_then(|source| source.feed(event))
.map(|params| (params, SourceInfo::NtpSource(sender_address, stratum)))
}
fn feed_phc(
phc: &mut source::Phc,
device_path: source::DevicePath,
event: event::Phc,
) -> Option<(&ClockParameters, SourceInfo)> {
phc.feed(event)
.map(|params| (params, SourceInfo::Phc(device_path)))
}
fn update_selected_clock(selected_clock: &Arc<SelectedClockSource>, source_info: &SourceInfo) {
match source_info {
SourceInfo::AmazonTimeSync(address, stratum)
| SourceInfo::NtpSource(address, stratum) => {
selected_clock.set_to_server(address.ip(), *stratum);
}
SourceInfo::Phc(_) => selected_clock.set_to_phc(),
}
}
pub fn handle_disruption(&mut self) {
let Self {
amazon_time_sync,
ntp_sources,
phc,
selected_clock,
selector,
seq,
} = self;
selected_clock.set_to_none();
if let Some(amazon_time_sync) = amazon_time_sync {
amazon_time_sync.handle_disruption();
}
for source in ntp_sources {
source.handle_disruption();
}
if let Some(phc) = phc {
phc.handle_disruption();
}
selector.handle_disruption();
tracing::info!("Handled clock disruption event.");
let current_seq = *seq;
*seq += 1;
logging::ffevents::log_disruption(current_seq);
}
}
#[cfg(test)]
mod tests {
use core::str;
use std::{net::Ipv4Addr, str::FromStr};
use rstest::rstest;
use crate::daemon::{
event::Stratum,
selected_clock::ClockSource,
time::{Duration, Instant, TscCount, tsc::Skew},
};
use super::*;
#[test]
#[tracing_test::traced_test]
fn feed_serializes_events() {
let event = event::Ntp::builder()
.counter_pre(TscCount::new(500))
.counter_post(TscCount::new(1000))
.ntp_data(event::NtpData {
server_recv_time: Instant::from_days(1),
server_send_time: Instant::from_days(2),
root_delay: Duration::from_micros(50),
root_dispersion: Duration::from_millis(17),
stratum: Stratum::TWO,
})
.build()
.unwrap();
let event = RoutableEvent::AmazonTimeSync(event);
let mut csa = ClockSyncAlgorithm::builder()
.amazon_time_sync(source::AmazonTimeSync::new(Skew::from_ppm(15.0)))
.ntp_sources(vec![])
.selected_clock(Arc::new(SelectedClockSource::default()))
.selector(Selector::new(Skew::from_ppm(15.0)))
.build();
let clock_parameters = csa.feed(event.clone());
assert!(clock_parameters.is_none());
use crate::daemon::logging::ffevents::types::{ClockParametersLog, RoutableEventLog};
let serialized_event = serde_json::to_string(&RoutableEventLog::from(&event)).unwrap();
let serialized_output = serde_json::to_string(
&clock_parameters.map(|sp| ClockParametersLog::from(&sp.clock_parameters)),
)
.unwrap();
let serialized_event = serialized_event.replace("\"", r#"\""#);
let serialized_output = serialized_output.replace("\"", r#"\""#);
assert!(logs_contain(&serialized_event));
assert!(logs_contain(&serialized_output));
}
#[rstest]
#[case(SourceInfo::AmazonTimeSync("169.254.169.123:123".parse().unwrap(), Stratum::TWO), ClockSource::Server(Ipv4Addr::from_str("169.254.169.123").unwrap().into()), Stratum::TWO)]
#[case(SourceInfo::NtpSource("169.254.169.101:123".parse().unwrap(), Stratum::ONE), ClockSource::Server(Ipv4Addr::from_str("169.254.169.101").unwrap().into()), Stratum::ONE)]
#[case(SourceInfo::NtpSource("[2001:db8::1:1234]:123".parse().unwrap(), Stratum::TWO), ClockSource::Server(Ipv4Addr::from_str("199.132.19.175").unwrap().into()), Stratum::TWO)]
#[case(SourceInfo::Phc("/dev/ptp0".into()), ClockSource::Phc, Stratum::Unspecified)]
fn update_selected_clock(
#[case] source_info: SourceInfo,
#[case] expected_clock_source: ClockSource,
#[case] expected_stratum: Stratum,
) {
let selected_clock_source = Arc::new(SelectedClockSource::default());
ClockSyncAlgorithm::update_selected_clock(&selected_clock_source, &source_info);
let (clock_source, stratum) = selected_clock_source.get();
assert_eq!(clock_source, expected_clock_source);
assert_eq!(stratum, expected_stratum);
}
fn make_ntp_event(pre: i64, post: i64) -> event::Ntp {
event::Ntp::builder()
.counter_pre(TscCount::new(pre))
.counter_post(TscCount::new(post))
.ntp_data(event::NtpData {
server_recv_time: Instant::from_days(1),
server_send_time: Instant::from_days(2),
root_delay: Duration::from_micros(50),
root_dispersion: Duration::from_millis(17),
stratum: Stratum::TWO,
})
.build()
.unwrap()
}
fn make_csa() -> ClockSyncAlgorithm {
ClockSyncAlgorithm::builder()
.selected_clock(Arc::new(SelectedClockSource::default()))
.selector(Selector::new(Skew::from_ppm(15.0)))
.build()
}
#[test]
fn add_pool_source_and_feed() {
let mut csa = make_csa();
let pool_domain = "pool.ntp.org".to_string();
let addr: SocketAddr = "192.0.2.1:123".parse().unwrap();
csa.add_ntp_source(source::NtpSource::new_with_pool(
addr,
pool_domain,
Skew::from_ppm(15.0),
));
let event = RoutableEvent::NtpSource(addr, make_ntp_event(500, 1000));
let result = csa.feed(event);
assert!(result.is_none());
}
#[test]
fn remove_pool_source_removes_correct_source() {
let mut csa = make_csa();
let pool_domain = "pool.ntp.org".to_string();
let addr: SocketAddr = "192.0.2.1:123".parse().unwrap();
csa.add_ntp_source(source::NtpSource::new_with_pool(
addr,
pool_domain,
Skew::from_ppm(15.0),
));
assert_eq!(csa.ntp_sources().len(), 1);
csa.remove_ntp_source(&addr);
assert_eq!(csa.ntp_sources().len(), 0);
}
}