pub mod adjust_data;
pub mod clock_adjust;
pub mod clock_state_writer;
use std::path::Path;
use std::sync::{Arc, Mutex};
use tokio::sync::watch;
use tokio_util::sync::CancellationToken;
use tracing::{debug, info};
use crate::daemon::MAX_DISPERSION_GROWTH_PPB;
use crate::daemon::async_ring_buffer::Receiver;
use crate::daemon::clock_parameters::ClockParameters;
use crate::daemon::clock_state::adjust_data::ClockAdjustData;
#[cfg(not(feature = "test-side-by-side"))]
use crate::daemon::clock_state::clock_adjust::KAPIClockAdjuster;
#[cfg(feature = "test-side-by-side")]
use crate::daemon::clock_state::clock_adjust::NoopClockAdjuster;
use crate::daemon::clock_state::clock_adjust::{ClockAdjust, ClockAdjuster, ClockCorrection};
use crate::daemon::clock_state::clock_state_writer::ClockStateWriter;
use crate::daemon::clock_state::clock_state_writer::{ClockStateWrite, SafeShmWriter};
use crate::daemon::clock_sync_algorithm::SyncParameters;
use crate::daemon::io::ClockDisruptionEvent;
use crate::daemon::io::tsc::ReadTscImpl;
use crate::daemon::io::vmclock::{State as VMClockState, VMClockParams};
use crate::daemon::logging::synchronization::{
log_clock_status_change, log_sync_snapshot, log_system_clock_step, log_system_clock_summary,
};
use crate::daemon::time::clocks::{ClockBound, MonotonicCoarse};
use crate::daemon::time::{Clock, Duration};
use crate::shm::{
CLOCKBOUND_SHM_DEFAULT_PATH_V0, CLOCKBOUND_SHM_DEFAULT_PATH_V1, ClockErrorBoundLayoutVersion,
ClockStatus, ShmWriter,
};
const FREE_RUNNING_GRACE_PERIOD: Duration = Duration::from_secs(60);
const CLOCK_ADJUST_SUMMARY_PERIOD_SECS: u64 = 600;
pub(crate) struct ClockState {
state_writer: Box<dyn ClockStateWrite>,
clock_adjuster: Box<dyn ClockAdjust>,
sync_parameters: Option<SyncParameters>,
clock_status: ClockStatus,
clock_disruption_receiver: watch::Receiver<ClockDisruptionEvent>,
clock_params_receiver: Receiver<SyncParameters>,
cancellation_token: CancellationToken,
vmclock_state: Option<Arc<Mutex<VMClockState>>>,
adjust_data: ClockAdjustData,
}
impl ClockState {
pub fn new(
clock_state_writer: Box<dyn ClockStateWrite>,
clock_adjuster: Box<dyn ClockAdjust>,
clock_params_receiver: Receiver<SyncParameters>,
clock_disruption_receiver: watch::Receiver<ClockDisruptionEvent>,
cancellation_token: CancellationToken,
clock_status: ClockStatus,
vmclock_state: Option<Arc<Mutex<VMClockState>>>,
) -> Self {
Self {
state_writer: clock_state_writer,
clock_adjuster,
clock_params_receiver,
clock_disruption_receiver,
sync_parameters: None,
clock_status,
cancellation_token,
vmclock_state,
adjust_data: ClockAdjustData::default(),
}
}
pub fn construct(
clock_params_receiver: Receiver<SyncParameters>,
clock_disruption_receiver: watch::Receiver<ClockDisruptionEvent>,
cancellation_token: CancellationToken,
vmclock_params: Option<VMClockParams>,
clock_status: ClockStatus,
) -> Self {
let clock_disruption_support_enabled = vmclock_params.is_some();
let disruption_marker = vmclock_params
.as_ref()
.map_or(0, |params| params.disruption_marker);
let vmclock_state = vmclock_params.map(|params| params.shared_state);
let shm_writer_0 = ShmWriter::new(
Path::new(CLOCKBOUND_SHM_DEFAULT_PATH_V0),
ClockErrorBoundLayoutVersion::V2,
ClockErrorBoundLayoutVersion::V2,
)
.unwrap();
let safe_shm_writer_0 = SafeShmWriter::new(shm_writer_0);
let shm_writer_1 = ShmWriter::new(
Path::new(CLOCKBOUND_SHM_DEFAULT_PATH_V1),
ClockErrorBoundLayoutVersion::V3,
ClockErrorBoundLayoutVersion::V3,
)
.unwrap();
let safe_shm_writer_1 = SafeShmWriter::new(shm_writer_1);
let clock_state_writer: ClockStateWriter<SafeShmWriter> = ClockStateWriter::builder()
.clock_disruption_support_enabled(clock_disruption_support_enabled)
.shm_writer_0(safe_shm_writer_0)
.shm_writer_1(safe_shm_writer_1)
.max_drift_ppb(MAX_DISPERSION_GROWTH_PPB)
.disruption_marker(disruption_marker)
.build();
#[cfg(not(feature = "test-side-by-side"))]
let clock_adjuster: ClockAdjuster<KAPIClockAdjuster> =
ClockAdjuster::new(KAPIClockAdjuster);
#[cfg(feature = "test-side-by-side")]
let clock_adjuster: ClockAdjuster<NoopClockAdjuster> =
ClockAdjuster::new(NoopClockAdjuster);
Self::new(
Box::new(clock_state_writer),
Box::new(clock_adjuster),
clock_params_receiver,
clock_disruption_receiver,
cancellation_token,
clock_status,
vmclock_state,
)
}
fn is_vmclock_in_failed_state(&self) -> bool {
self.vmclock_state
.as_ref()
.is_some_and(|state| *state.lock().unwrap() == VMClockState::Failed)
}
pub async fn run(&mut self) {
let mut snapshot_interval = tokio::time::interval(tokio::time::Duration::from_mins(1));
snapshot_interval.reset_after(tokio::time::Duration::from_secs(8));
let mut clock_adjust_summary_interval = tokio::time::interval(
tokio::time::Duration::from_secs(CLOCK_ADJUST_SUMMARY_PERIOD_SECS),
);
clock_adjust_summary_interval.reset();
debug!("Starting run for ClockState.");
self.state_writer.initialize_ceb_v2_shm();
loop {
tokio::select! {
biased; Ok(()) = self.clock_disruption_receiver.changed() => {
self.handle_disruption();
}
source_params = self.clock_params_receiver.recv() => {
self.handle_sync_parameters(source_params.unwrap());
},
_ = snapshot_interval.tick() => {
log_sync_snapshot(self.sync_parameters.as_ref(), self.clock_status);
},
_ = clock_adjust_summary_interval.tick() => {
self.flush_clock_adjust_summary();
},
() = self.cancellation_token.cancelled() => {
debug!("Received shutdown signal. Exiting.");
break;
},
}
}
debug!("ClockState runner exiting.");
}
fn determine_clock_error_bound_v2_status(&self, parameters: &ClockParameters) -> ClockStatus {
let mut clock_status = self.clock_adjuster.get_clock_realtime_status();
let time_since_parameters_updated = MonotonicCoarse.get_time() - parameters.as_of_monotonic;
if clock_status == ClockStatus::Synchronized
&& time_since_parameters_updated >= FREE_RUNNING_GRACE_PERIOD
{
clock_status = ClockStatus::FreeRunning;
}
clock_status
}
fn determine_clock_error_bound_v3_status(clock_params_age: Duration) -> ClockStatus {
if clock_params_age < FREE_RUNNING_GRACE_PERIOD {
ClockStatus::Synchronized
} else {
ClockStatus::FreeRunning
}
}
fn set_clock_status(status: &mut ClockStatus, new: ClockStatus) {
let previous = *status;
*status = new;
if new != previous {
info!("Clock status changed from {previous:?} to {new:?}.");
log_clock_status_change(previous, new);
}
}
pub fn handle_sync_parameters(&mut self, sync_parameters: SyncParameters) {
self.sync_parameters = Some(sync_parameters);
let vmclock_failed = self.is_vmclock_in_failed_state();
if let Some(sp) = &self.sync_parameters {
let parameters = &sp.clock_parameters;
match self.clock_adjuster.handle_clock_parameters(parameters) {
ClockCorrection::Step { step_ns } => log_system_clock_step(step_ns),
ClockCorrection::Smooth { correction_ns } => {
self.adjust_data.record(correction_ns);
}
}
let clockbound_clock = ClockBound::new(parameters.clone(), ReadTscImpl);
let clock_params_age_v3 = clockbound_clock.get_time() - parameters.time;
let mut clock_status_v3 =
ClockState::determine_clock_error_bound_v3_status(clock_params_age_v3);
if vmclock_failed {
clock_status_v3 = ClockStatus::Unknown;
}
Self::set_clock_status(&mut self.clock_status, clock_status_v3);
self.state_writer
.handle_clock_parameters_shm1(parameters, clock_status_v3);
let mut clock_status_v2 = self.determine_clock_error_bound_v2_status(parameters);
if vmclock_failed {
clock_status_v2 = ClockStatus::Unknown;
}
self.state_writer
.handle_clock_parameters_shm0(parameters, clock_status_v2);
}
}
fn flush_clock_adjust_summary(&mut self) {
log_system_clock_summary(CLOCK_ADJUST_SUMMARY_PERIOD_SECS, &self.adjust_data);
self.adjust_data = ClockAdjustData::default();
}
pub fn handle_disruption(&mut self) {
let Self {
clock_adjuster,
state_writer: clock_state_writer,
clock_params_receiver,
clock_disruption_receiver,
sync_parameters,
clock_status,
cancellation_token: _,
vmclock_state: _,
adjust_data: _,
} = self;
let val = clock_disruption_receiver.borrow_and_update().clone();
if let Some(disruption_marker) = val.disruption_marker {
if let Some(sp) = sync_parameters {
clock_state_writer.handle_disruption(&sp.clock_parameters, disruption_marker);
}
*sync_parameters = None;
clock_params_receiver.handle_disruption();
clock_adjuster.handle_disruption(disruption_marker);
Self::set_clock_status(clock_status, ClockStatus::Disrupted);
tracing::debug!("Handled clock disruption event.");
}
}
}
#[cfg(test)]
mod tests {
use mockall::predicate::eq;
use rstest::rstest;
use crate::{
daemon::{
async_ring_buffer,
clock_state::{clock_adjust::MockClockAdjust, clock_state_writer::MockClockStateWrite},
clock_sync_algorithm::SourceInfo,
time::{Duration, Instant, TscCount, tsc::Period},
},
shm::ClockStatus,
};
use super::*;
fn get_sample_clock_parameters() -> ClockParameters {
ClockParameters {
tsc_count: TscCount::new(0),
time: Instant::new(0),
clock_error_bound: Duration::new(0),
period: Period::from_seconds(0.0),
period_max_error: Period::from_seconds(0.0),
as_of_monotonic: Instant::new(0),
}
}
#[tokio::test]
async fn handle_disruption() {
let disruption_marker = 123;
let cancellation_token = CancellationToken::new();
let mut mock_clock_adjuster: MockClockAdjust = MockClockAdjust::new();
let clock_parameters = ClockParameters {
tsc_count: TscCount::new(1000),
time: Instant::from_nanos(1000),
clock_error_bound: Duration::from_nanos(1000),
period: Period::from_seconds(1e-9),
period_max_error: Period::from_seconds(1e-11),
as_of_monotonic: Instant::from_nanos(1000),
};
let expected_clock_parameters = clock_parameters.clone();
mock_clock_adjuster
.expect_handle_disruption()
.once()
.with(eq(disruption_marker))
.return_const(());
let mut mock_clock_state_writer: MockClockStateWrite = MockClockStateWrite::new();
mock_clock_state_writer
.expect_handle_disruption()
.once()
.withf(move |param: &ClockParameters, marker: &u64| {
*param == expected_clock_parameters && *marker == 123
})
.return_const(());
let (_tx, rx) = async_ring_buffer::create(1);
let (clock_disruption_sender, clock_disruption_receiver) =
watch::channel::<ClockDisruptionEvent>(ClockDisruptionEvent::default());
let mut clock_state = ClockState::new(
Box::new(mock_clock_state_writer),
Box::new(mock_clock_adjuster),
rx,
clock_disruption_receiver,
cancellation_token,
ClockStatus::Unknown,
None,
);
clock_state.sync_parameters = Some(SyncParameters {
clock_parameters,
source_info: SourceInfo::Phc("/dev/ptp0".into()),
selected_at: Instant::new(0),
selected_at_clock_error_bound: Duration::new(0),
});
clock_disruption_sender
.send(ClockDisruptionEvent {
disruption_marker: Some(disruption_marker),
})
.unwrap();
clock_state.handle_disruption();
}
#[tokio::test(start_paused = true)]
async fn handle_sync_parameters_with_parameters() {
let mut sequence = mockall::Sequence::new();
let expected_clock_error_bound_v2_status = ClockStatus::Synchronized;
let expected_clock_error_bound_v3_status = ClockStatus::Synchronized;
let mut expected_clock_params = get_sample_clock_parameters();
expected_clock_params.as_of_monotonic = MonotonicCoarse.get_time();
let cancellation_token = CancellationToken::new();
let mut mock_clock_adjuster: MockClockAdjust = MockClockAdjust::new();
let expected_clock_params_clone = expected_clock_params.clone();
mock_clock_adjuster
.expect_handle_clock_parameters()
.once()
.withf(move |actual_clock_params| *actual_clock_params == expected_clock_params_clone)
.in_sequence(&mut sequence)
.return_const(ClockCorrection::Smooth { correction_ns: 0 });
let mut mock_clock_state_writer: MockClockStateWrite = MockClockStateWrite::new();
let expected_clock_params_clone = expected_clock_params.clone();
mock_clock_state_writer
.expect_handle_clock_parameters_shm1()
.once()
.withf(move |actual_clock_params, actual_clock_status| {
*actual_clock_params == expected_clock_params_clone
&& *actual_clock_status == expected_clock_error_bound_v3_status
})
.in_sequence(&mut sequence)
.return_const(());
mock_clock_adjuster
.expect_get_clock_realtime_status()
.once()
.in_sequence(&mut sequence)
.return_const(ClockStatus::Synchronized);
let expected_clock_params_clone = expected_clock_params.clone();
mock_clock_state_writer
.expect_handle_clock_parameters_shm0()
.once()
.withf(move |actual_clock_params, actual_clock_status| {
*actual_clock_params == expected_clock_params_clone
&& *actual_clock_status == expected_clock_error_bound_v2_status
})
.in_sequence(&mut sequence)
.return_const(());
let (_tx, rx) = async_ring_buffer::create(1);
let (_, clock_disruption_receiver) =
watch::channel::<ClockDisruptionEvent>(ClockDisruptionEvent::default());
let mut clock_state = ClockState::new(
Box::new(mock_clock_state_writer),
Box::new(mock_clock_adjuster),
rx,
clock_disruption_receiver,
cancellation_token,
ClockStatus::Unknown,
None,
);
let sync_parameters = SyncParameters {
clock_parameters: expected_clock_params,
source_info: SourceInfo::Phc("/dev/ptp0".into()),
selected_at: Instant::new(0),
selected_at_clock_error_bound: Duration::new(0),
};
assert_eq!(clock_state.sync_parameters, None);
clock_state.handle_sync_parameters(sync_parameters.clone());
assert_eq!(clock_state.sync_parameters, Some(sync_parameters));
}
#[tokio::test(start_paused = true)]
async fn handle_sync_parameters_vmclock_failed_forces_unknown() {
let mut expected_clock_params = get_sample_clock_parameters();
expected_clock_params.as_of_monotonic = MonotonicCoarse.get_time();
let cancellation_token = CancellationToken::new();
let mut mock_clock_adjuster: MockClockAdjust = MockClockAdjust::new();
mock_clock_adjuster
.expect_handle_clock_parameters()
.once()
.return_const(ClockCorrection::Smooth { correction_ns: 0 });
mock_clock_adjuster
.expect_get_clock_realtime_status()
.once()
.return_const(ClockStatus::Synchronized);
let mut mock_clock_state_writer: MockClockStateWrite = MockClockStateWrite::new();
mock_clock_state_writer
.expect_handle_clock_parameters_shm1()
.once()
.withf(|_clock_params, clock_status| *clock_status == ClockStatus::Unknown)
.return_const(());
mock_clock_state_writer
.expect_handle_clock_parameters_shm0()
.once()
.withf(|_clock_params, clock_status| *clock_status == ClockStatus::Unknown)
.return_const(());
let (_tx, rx) = async_ring_buffer::create(1);
let (_, clock_disruption_receiver) =
watch::channel::<ClockDisruptionEvent>(ClockDisruptionEvent::default());
let vmclock_state = Some(Arc::new(Mutex::new(VMClockState::Failed)));
let mut clock_state = ClockState::new(
Box::new(mock_clock_state_writer),
Box::new(mock_clock_adjuster),
rx,
clock_disruption_receiver,
cancellation_token,
ClockStatus::Unknown,
vmclock_state,
);
clock_state.handle_sync_parameters(SyncParameters {
clock_parameters: expected_clock_params,
source_info: SourceInfo::Phc("/dev/ptp0".into()),
selected_at: Instant::new(0),
selected_at_clock_error_bound: Duration::new(0),
});
}
#[rstest]
#[case::vmclock_none(None)]
#[case::vmclock_running(Some(Arc::new(Mutex::new(VMClockState::Running))))]
#[tokio::test(start_paused = true)]
async fn handle_sync_parameters_nominal(
#[case] vmclock_state: Option<Arc<Mutex<VMClockState>>>,
) {
let mut expected_clock_params = get_sample_clock_parameters();
expected_clock_params.as_of_monotonic = MonotonicCoarse.get_time();
let cancellation_token = CancellationToken::new();
let mut mock_clock_adjuster: MockClockAdjust = MockClockAdjust::new();
mock_clock_adjuster
.expect_handle_clock_parameters()
.once()
.return_const(ClockCorrection::Smooth { correction_ns: 0 });
mock_clock_adjuster
.expect_get_clock_realtime_status()
.once()
.return_const(ClockStatus::Synchronized);
let mut mock_clock_state_writer: MockClockStateWrite = MockClockStateWrite::new();
mock_clock_state_writer
.expect_handle_clock_parameters_shm1()
.once()
.withf(|_clock_params, clock_status| *clock_status == ClockStatus::Synchronized)
.return_const(());
mock_clock_state_writer
.expect_handle_clock_parameters_shm0()
.once()
.withf(|_clock_params, clock_status| *clock_status == ClockStatus::Synchronized)
.return_const(());
let (_tx, rx) = async_ring_buffer::create(1);
let (_, clock_disruption_receiver) =
watch::channel::<ClockDisruptionEvent>(ClockDisruptionEvent::default());
let mut clock_state = ClockState::new(
Box::new(mock_clock_state_writer),
Box::new(mock_clock_adjuster),
rx,
clock_disruption_receiver,
cancellation_token,
ClockStatus::Unknown,
vmclock_state,
);
clock_state.handle_sync_parameters(SyncParameters {
clock_parameters: expected_clock_params,
source_info: SourceInfo::Phc("/dev/ptp0".into()),
selected_at: Instant::new(0),
selected_at_clock_error_bound: Duration::new(0),
});
}
#[tokio::test]
async fn vmclock_failed_reflects_shared_state() {
let cancellation_token = CancellationToken::new();
let (_, clock_disruption_receiver) =
watch::channel::<ClockDisruptionEvent>(ClockDisruptionEvent::default());
let build = |vmclock_state: Option<Arc<Mutex<VMClockState>>>| {
let (_, rx) = async_ring_buffer::create(1);
ClockState::new(
Box::new(MockClockStateWrite::new()),
Box::new(MockClockAdjust::new()),
rx,
clock_disruption_receiver.clone(),
cancellation_token.clone(),
ClockStatus::Unknown,
vmclock_state,
)
};
assert!(!build(None).is_vmclock_in_failed_state());
assert!(
!build(Some(Arc::new(Mutex::new(VMClockState::Running)))).is_vmclock_in_failed_state()
);
assert!(
build(Some(Arc::new(Mutex::new(VMClockState::Failed)))).is_vmclock_in_failed_state()
);
}
#[rstest]
#[case::synchronized_stays_synchronized_params_0sec_old(
Duration::from_secs(0),
ClockStatus::Synchronized,
ClockStatus::Synchronized
)]
#[case::synchronized_stays_synchronized_params_30sec_old(
Duration::from_secs(30),
ClockStatus::Synchronized,
ClockStatus::Synchronized
)]
#[case::synchronized_goes_freerunning_params_60sec_old(
Duration::from_secs(60),
ClockStatus::Synchronized,
ClockStatus::FreeRunning
)]
#[case::synchronized_goes_freerunning_params_90sec_old(
Duration::from_secs(90),
ClockStatus::Synchronized,
ClockStatus::FreeRunning
)]
#[case::unknown_stays_unknown_params_0sec_old(
Duration::from_secs(0),
ClockStatus::Unknown,
ClockStatus::Unknown
)]
#[case::unknown_stays_unknown_params_90sec_old(
Duration::from_secs(90),
ClockStatus::Unknown,
ClockStatus::Unknown
)]
#[case::disrupted_stays_disrupted_params_0sec_old(
Duration::from_secs(0),
ClockStatus::Disrupted,
ClockStatus::Disrupted
)]
#[case::disrupted_stays_disrupted_params_90sec_old(
Duration::from_secs(90),
ClockStatus::Disrupted,
ClockStatus::Disrupted
)]
#[tokio::test]
async fn determine_clock_error_bound_v2_status(
#[case] clock_params_age: Duration,
#[case] clock_adjust_status: ClockStatus,
#[case] expected_clock_error_bound_v2_status: ClockStatus,
) {
let cancellation_token = CancellationToken::new();
let mut mock_clock_adjuster: MockClockAdjust = MockClockAdjust::new();
let (_, clock_disruption_receiver) =
watch::channel::<ClockDisruptionEvent>(ClockDisruptionEvent::default());
mock_clock_adjuster
.expect_get_clock_realtime_status()
.once()
.return_const(clock_adjust_status);
let clock_state = ClockState::new(
Box::new(MockClockStateWrite::new()),
Box::new(mock_clock_adjuster),
async_ring_buffer::create(1).1,
clock_disruption_receiver,
cancellation_token,
ClockStatus::Unknown,
None,
);
let mut clock_parameters = get_sample_clock_parameters();
clock_parameters.as_of_monotonic = MonotonicCoarse.get_time() - clock_params_age;
let res = clock_state.determine_clock_error_bound_v2_status(&clock_parameters);
assert_eq!(res, expected_clock_error_bound_v2_status);
}
#[rstest]
#[case::synchronized_params_0sec_old(Duration::from_secs(0), ClockStatus::Synchronized)]
#[case::synchronized_params_30sec_old(Duration::from_secs(30), ClockStatus::Synchronized)]
#[case::freerunning_params_60sec_old(Duration::from_secs(60), ClockStatus::FreeRunning)]
#[case::freerunning_params_90sec_old(Duration::from_secs(90), ClockStatus::FreeRunning)]
fn determine_clock_error_bound_v3_status(
#[case] clock_params_age: Duration,
#[case] expected_clock_error_bound_v3_status: ClockStatus,
) {
assert_eq!(
ClockState::determine_clock_error_bound_v3_status(clock_params_age),
expected_clock_error_bound_v3_status
);
}
#[test]
fn set_clock_status_logs_only_on_transition() {
use crate::daemon::logging::StructuredLog;
use crate::daemon::logging::test_support::with_layer;
use serde_json::Value;
let (output, _tmp) = with_layer(StructuredLog::Synchronization, "0.1.0-test", || {
let mut status = ClockStatus::Unknown;
ClockState::set_clock_status(&mut status, ClockStatus::Synchronized);
ClockState::set_clock_status(&mut status, ClockStatus::Synchronized);
assert_eq!(status, ClockStatus::Synchronized);
});
let lines: Vec<&str> = output.lines().collect();
assert_eq!(
lines.len(),
1,
"only the actual transition should be logged"
);
let entry: Value = serde_json::from_str(lines[0]).unwrap();
assert_eq!(entry["event"], "clock_status_change");
assert_eq!(entry["clock_status_change"]["from"], "UNKNOWN");
assert_eq!(entry["clock_status_change"]["to"], "SYNCHRONIZED");
}
}