use std::time::{Duration, SystemTime};
use bytes::Bytes;
use microsandbox_protocol::codec;
use microsandbox_protocol::core::ClockSync;
use microsandbox_protocol::message::{Message, MessageType};
use microsandbox_types::GuestClockPolicy;
use tokio::task::JoinHandle;
use crate::relay::{ControlWrite, ControlWriter};
use crate::{RuntimeError, RuntimeResult};
const CLOCK_SYNC_POLL_INTERVAL: Duration = Duration::from_secs(2);
const CLOCK_SYNC_INTERVAL: Duration = Duration::from_secs(60);
const CLOCK_SYNC_WAKE_THRESHOLD: Duration = Duration::from_secs(6);
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum RestoreActivationMode {
IdentityAndClock,
IdentityOnly,
}
impl RestoreActivationMode {
pub(crate) fn for_policy(policy: GuestClockPolicy) -> Self {
match policy {
GuestClockPolicy::Sync => Self::IdentityAndClock,
GuestClockPolicy::Off => Self::IdentityOnly,
}
}
pub(crate) fn description(self) -> &'static str {
match self {
Self::IdentityAndClock => "identity-and-clock activation",
Self::IdentityOnly => "VM Generation ID activation",
}
}
pub(crate) fn install(
self,
vm: &msb_krun::VmControl,
id: msb_krun::VmGenerationId,
) -> Option<msb_krun::VmGenerationRequest> {
match self {
Self::IdentityAndClock => vm.install_vm_generation_and_clock(id),
Self::IdentityOnly => vm.install_vm_generation_id(id),
}
}
}
pub(crate) fn spawn_clock_sync_task(
agent_tx: ControlWriter,
policy: GuestClockPolicy,
already_synchronized: bool,
) -> Option<JoinHandle<()>> {
policy
.is_sync()
.then(|| tokio::spawn(clock_sync_task(agent_tx, already_synchronized)))
}
async fn clock_sync_task(agent_tx: ControlWriter, already_synchronized: bool) {
let mut last_wall = SystemTime::now();
let mut last_sync = if already_synchronized {
last_wall
} else {
match send_clock_sync(&agent_tx).await {
Ok(sent_at) => sent_at,
Err(err) => {
tracing::debug!(error = %err, "agent relay: initial clock sync failed");
return;
}
}
};
let mut interval = tokio::time::interval(CLOCK_SYNC_POLL_INTERVAL);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
loop {
interval.tick().await;
let now = SystemTime::now();
let wall_gap = now
.duration_since(last_wall)
.unwrap_or(CLOCK_SYNC_WAKE_THRESHOLD);
let since_sync = now.duration_since(last_sync).unwrap_or(CLOCK_SYNC_INTERVAL);
if wall_gap >= CLOCK_SYNC_WAKE_THRESHOLD || since_sync >= CLOCK_SYNC_INTERVAL {
match send_clock_sync(&agent_tx).await {
Ok(sent_at) => last_sync = sent_at,
Err(err) => {
tracing::debug!(error = %err, "agent relay: clock sync task exiting");
break;
}
}
}
last_wall = now;
}
}
async fn send_clock_sync(agent_tx: &ControlWriter) -> RuntimeResult<SystemTime> {
let now = SystemTime::now();
agent_tx
.send(ControlWrite::clock_sync()?)
.await
.map_err(|_| RuntimeError::Custom("agent relay ring writer channel closed".into()))?;
Ok(now)
}
pub(crate) fn current_clock_sync_frame() -> RuntimeResult<Bytes> {
let elapsed = SystemTime::now()
.duration_since(SystemTime::UNIX_EPOCH)
.map_err(|e| RuntimeError::Custom(format!("clock sync before Unix epoch: {e}")))?;
let unix_time_nanos = u64::try_from(elapsed.as_nanos()).map_err(|_| {
RuntimeError::Custom("clock sync timestamp does not fit in u64 nanoseconds".into())
})?;
encode_clock_sync_frame(unix_time_nanos)
}
pub(crate) fn encode_clock_sync_frame(unix_time_nanos: u64) -> RuntimeResult<Bytes> {
let sync = ClockSync { unix_time_nanos };
let msg = Message::with_payload(MessageType::ClockSync, 0, &sync)
.map_err(|e| RuntimeError::Custom(format!("encode clock sync: {e}")))?;
let mut buf = Vec::new();
codec::encode_to_buf(&msg, &mut buf)
.map_err(|e| RuntimeError::Custom(format!("encode clock sync frame: {e}")))?;
Ok(Bytes::from(buf))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn restore_activation_steps_the_clock_only_when_synchronizing() {
assert_eq!(
RestoreActivationMode::for_policy(GuestClockPolicy::Sync),
RestoreActivationMode::IdentityAndClock
);
assert_eq!(
RestoreActivationMode::for_policy(GuestClockPolicy::Off),
RestoreActivationMode::IdentityOnly
);
}
}