use rusty_time_clock::{ClockRead, SystemClock, net};
use rusty_time_core::ntp::{self, HEADER_LEN, LeapIndicator, Mode, NtpPacket, NtpTimestamp};
use rusty_time_core::server::{
ClientHandle, ClientTable, Disposition, RateLimitConfig, ResponseMode,
};
use rusty_time_nts::aead::NtsKeys;
use rusty_time_nts::cookie::{COOKIE_NONCE_LEN, KeyRing, MasterKey};
use rusty_time_nts::ef::{self, NONCE_LEN, UNIQUE_ID_LEN};
use rusty_time_nts::records::{self, record_type};
use rusty_time_nts::tls::rustls;
use rusty_time_nts::{AEAD_AES_SIV_CMAC_256, ALPN, NEXT_PROTO_NTPV4};
use crate::store::{MASTER_KEY_SLOTS, Store, StoredMasterKey};
use std::io::{Read, Write};
use std::net::{SocketAddr, TcpListener, UdpSocket};
use std::sync::{Arc, Mutex};
use std::time::Duration;
const CLIENT_TABLE_CAPACITY: usize = 16_384;
const STATUS_INTERVAL: Duration = Duration::from_secs(60);
fn client_key(peer: SocketAddr) -> std::net::IpAddr {
peer.ip()
}
const KE_COOKIE_COUNT: usize = 8;
const MAX_REPLY_COOKIES: usize = 8;
const COOKIE_FIELD_HINT: usize = 112;
pub struct ServeOptions {
pub ntp_bind: String,
pub ke_bind: String,
pub stratum: u8,
pub nts: bool,
pub cert_pem: Option<String>,
pub key_pem: Option<String>,
pub nts_name: String,
pub write_cert: Option<String>,
pub state_path: Option<String>,
pub state_passphrase: Option<String>,
pub gateway_bind: Option<String>,
pub gateway_assets: Option<String>,
pub rate_limit: RateLimitConfig,
pub control_path: String,
}
impl ServeOptions {
pub fn parse(args: &[String]) -> Result<ServeOptions, String> {
let mut opts = ServeOptions {
ntp_bind: "0.0.0.0:123".into(),
ke_bind: "0.0.0.0:4460".into(),
stratum: 1,
nts: false,
cert_pem: None,
key_pem: None,
nts_name: "localhost".into(),
write_cert: None,
state_path: None,
state_passphrase: std::env::var("RUSTY_TIME_STATE_PASSPHRASE").ok(),
gateway_bind: None,
gateway_assets: None,
rate_limit: RateLimitConfig::default(),
control_path: crate::control::default_path(),
};
let mut it = args.iter();
while let Some(flag) = it.next() {
let value = |v: Option<&String>| -> Result<String, String> {
v.cloned().ok_or(format!("{flag} needs a value"))
};
match flag.as_str() {
"--nts" => opts.nts = true,
"--bind" => opts.ntp_bind = value(it.next())?,
"--ke-bind" => opts.ke_bind = value(it.next())?,
"--nts-name" => opts.nts_name = value(it.next())?,
"--write-cert" => opts.write_cert = Some(value(it.next())?),
"--state" => opts.state_path = Some(value(it.next())?),
"--control" => opts.control_path = value(it.next())?,
"--gateway" => opts.gateway_bind = Some(value(it.next())?),
"--gateway-assets" => opts.gateway_assets = Some(value(it.next())?),
"--ratelimit-interval" => {
opts.rate_limit.interval_log2 = value(it.next())?
.parse()
.map_err(|_| "--ratelimit-interval: not a number".to_string())?;
}
"--ratelimit-burst" => {
opts.rate_limit.burst = value(it.next())?
.parse()
.map_err(|_| "--ratelimit-burst: not a number".to_string())?;
}
"--ratelimit-global" => {
opts.rate_limit.global_rate_hz = value(it.next())?
.parse()
.map_err(|_| "--ratelimit-global: not a number".to_string())?;
opts.rate_limit.global_burst = opts.rate_limit.global_rate_hz * 2.0;
}
"--no-ratelimit" => {
opts.rate_limit = RateLimitConfig {
interval_log2: -20,
burst: 1_000_000,
leak_shift: 0,
global_rate_hz: 0.0,
global_burst: 0.0,
};
}
"--stratum" => {
opts.stratum = value(it.next())?
.parse()
.map_err(|_| "--stratum: not a number".to_string())?;
}
"--cert" => {
let path = value(it.next())?;
opts.cert_pem =
Some(std::fs::read_to_string(&path).map_err(|e| format!("{path}: {e}"))?);
}
"--key" => {
let path = value(it.next())?;
opts.key_pem =
Some(std::fs::read_to_string(&path).map_err(|e| format!("{path}: {e}"))?);
}
other => return Err(format!("unknown flag '{other}'")),
}
}
if opts.stratum == 0 || opts.stratum > 15 {
return Err("--stratum must be 1..=15".into());
}
Ok(opts)
}
}
fn load_or_create_ring(opts: &ServeOptions) -> Result<(KeyRing, Option<Store>), String> {
let Some(path) = &opts.state_path else {
return Ok((fresh_key_ring()?, None));
};
let passphrase = opts.state_passphrase.as_deref().ok_or(
"--state needs a passphrase: set RUSTY_TIME_STATE_PASSPHRASE (the state file holds \
NTS master keys, which forge every cookie we ever minted)",
)?;
let mut store = Store::open(path, passphrase.as_bytes()).map_err(|e| e.to_string())?;
let stored = store.all_master_keys().map_err(|e| e.to_string())?;
let mut ring = KeyRing::new(MASTER_KEY_SLOTS);
if stored.is_empty() {
let fresh = fresh_key_ring()?;
if let Some(k) = fresh.current() {
store
.put_master_key(
0,
&StoredMasterKey {
id: k.id,
key: k.key,
},
)
.map_err(|e| e.to_string())?;
store.flush().map_err(|e| e.to_string())?;
}
println!("rtimed serve: minted a new NTS master key; state at {path}");
return Ok((fresh, Some(store)));
}
for k in &stored {
ring.rotate_in(MasterKey {
id: k.id,
key: k.key,
});
}
println!(
"rtimed serve: restored {} NTS master key(s) from {path}; cookies minted before this \
restart remain valid",
stored.len()
);
Ok((ring, Some(store)))
}
pub fn fresh_key_ring() -> Result<KeyRing, String> {
let mut ring = KeyRing::new(3);
let mut key = [0u8; 32];
rusty_time_nts::ke::fill_random(&mut key).map_err(|e| e.to_string())?;
let mut id_bytes = [0u8; 4];
rusty_time_nts::ke::fill_random(&mut id_bytes).map_err(|e| e.to_string())?;
ring.rotate_in(MasterKey {
id: u32::from_be_bytes(id_bytes),
key,
});
Ok(ring)
}
pub fn run(opts: &ServeOptions) -> i32 {
let (ring, store) = match load_or_create_ring(opts) {
Ok(v) => v,
Err(e) => {
eprintln!("rtimed serve: {e}");
return 1;
}
};
let state = Arc::new(Mutex::new(ServerState {
clients: ClientTable::new(CLIENT_TABLE_CAPACITY, opts.rate_limit),
ring,
stratum: opts.stratum,
started_unix: std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0),
}));
let _store = store;
{
let ctl_state = Arc::clone(&state);
let ctl_path = opts.control_path.clone();
std::thread::spawn(move || {
if let Err(e) = crate::control::serve(&ctl_path, ctl_state) {
eprintln!("rtimed serve: control plane unavailable: {e}");
}
});
match rusty_time_api::control_endpoint(&opts.control_path) {
rusty_time_api::ControlEndpoint::UnixPath(p) => {
println!("rtimed: control socket at {p}")
}
rusty_time_api::ControlEndpoint::Loopback(port) => println!(
"rtimed: control on 127.0.0.1:{port} (from '{}')",
opts.control_path
),
}
}
if let Some(bind) = &opts.gateway_bind {
let gw_state = Arc::clone(&state);
let bind = bind.clone();
let assets = opts.gateway_assets.clone();
println!(
"rtimed: gateway (NTP over HTTP) on http://{bind}/{}",
if assets.is_some() {
" with wasm assets"
} else {
" (status page only; pass --gateway-assets for the wasm demo)"
}
);
std::thread::spawn(move || {
if let Err(e) = crate::gateway::serve(&bind, gw_state, assets) {
eprintln!("rtimed serve: gateway unavailable: {e}");
}
});
}
if opts.nts {
let tls_config = match server_tls_config(opts) {
Ok(c) => c,
Err(e) => {
eprintln!("rtimed serve: TLS setup: {e}");
return 1;
}
};
let listener = match TcpListener::bind(&opts.ke_bind) {
Ok(l) => l,
Err(e) => {
eprintln!("rtimed serve: binding NTS-KE {}: {e}", opts.ke_bind);
return 1;
}
};
let ke_state = Arc::clone(&state);
let ntp_port = opts
.ntp_bind
.rsplit(':')
.next()
.and_then(|p| p.parse::<u16>().ok())
.unwrap_or(123);
println!("rtimed: NTS-KE listening on {}", opts.ke_bind);
std::thread::spawn(move || ke_accept_loop(listener, tls_config, ke_state, ntp_port));
}
let socket = match crate::service::activated_udp_socket() {
Some(s) => {
println!("rtimed: using the socket passed by the service manager");
s
}
None => match UdpSocket::bind(&opts.ntp_bind) {
Ok(s) => s,
Err(e) => {
eprintln!("rtimed serve: binding NTP {}: {e}", opts.ntp_bind);
return 1;
}
},
};
println!(
"rtimed: NTP listening on {} (stratum {}, rate limit 1 per {} s burst {})",
opts.ntp_bind,
opts.stratum,
2f64.powi(opts.rate_limit.interval_log2 as i32),
opts.rate_limit.burst
);
let caps = rusty_time_clock::capabilities();
crate::service::notify_ready(&format!(
"serving NTP on {} (stratum {}{})",
opts.ntp_bind,
opts.stratum,
if caps.can_discipline {
""
} else {
", clock read-only"
}
));
ntp_serve_loop(&socket, &state);
0
}
fn server_tls_config(opts: &ServeOptions) -> Result<Arc<rustls::ServerConfig>, String> {
use rusty_time_nts::tls::pki_types::{CertificateDer, PrivateKeyDer, pem::PemObject};
let (certs, key): (Vec<CertificateDer<'static>>, PrivateKeyDer<'static>) =
match (&opts.cert_pem, &opts.key_pem) {
(Some(cert_pem), Some(key_pem)) => {
let certs = CertificateDer::pem_slice_iter(cert_pem.as_bytes())
.collect::<Result<Vec<_>, _>>()
.map_err(|e| format!("parsing --cert: {e}"))?;
let key = PrivateKeyDer::from_pem_slice(key_pem.as_bytes())
.map_err(|e| format!("parsing --key: {e}"))?;
(certs, key)
}
(None, None) => {
eprintln!(
"rtimed serve: no --cert/--key given; generating a SELF-SIGNED certificate \
for '{}'. Development only โ clients must be told to trust it.",
opts.nts_name
);
let ck = oxitls_rcgen::generate_self_signed_p256(&[opts.nts_name.as_str()])
.map_err(|e| format!("generating self-signed certificate: {e}"))?;
if let Some(path) = &opts.write_cert {
std::fs::write(path, &ck.cert_pem)
.map_err(|e| format!("writing {path}: {e}"))?;
println!("rtimed serve: wrote certificate to {path}");
}
(
vec![CertificateDer::from(ck.cert_der)],
PrivateKeyDer::try_from(ck.pkcs8_der)
.map_err(|e| format!("self-signed key: {e}"))?,
)
}
_ => return Err("--cert and --key must be given together".into()),
};
let config = rusty_time_nts::tls::server_config(certs, key, &[ALPN])
.map_err(|e| format!("building server TLS config: {e}"))?;
Ok(Arc::new(config))
}
fn ke_accept_loop(
listener: TcpListener,
config: Arc<rustls::ServerConfig>,
state: Arc<Mutex<ServerState>>,
ntp_port: u16,
) {
for stream in listener.incoming() {
let Ok(stream) = stream else { continue };
let config = Arc::clone(&config);
let state = Arc::clone(&state);
std::thread::spawn(move || {
if let Err(e) = handle_ke(stream, config, state, ntp_port) {
eprintln!("rtimed serve: NTS-KE session: {e}");
}
});
}
}
fn handle_ke(
stream: std::net::TcpStream,
config: Arc<rustls::ServerConfig>,
state: Arc<Mutex<ServerState>>,
ntp_port: u16,
) -> Result<(), String> {
stream
.set_read_timeout(Some(Duration::from_secs(10)))
.map_err(|e| e.to_string())?;
stream
.set_write_timeout(Some(Duration::from_secs(10)))
.map_err(|e| e.to_string())?;
let conn = rustls::ServerConnection::new(config).map_err(|e| e.to_string())?;
let mut tls = rustls::StreamOwned::new(conn, stream);
let mut buf = Vec::new();
let mut chunk = [0u8; 2048];
loop {
match tls.read(&mut chunk) {
Ok(0) => break,
Ok(n) => {
buf.extend_from_slice(&chunk[..n]);
if buf.len() > 64 * 1024 {
return Err("client sent an oversized NTS-KE request".into());
}
if records::records(&buf)
.any(|r| matches!(r, Ok(rec) if rec.record_type == record_type::END_OF_MESSAGE))
{
break;
}
}
Err(e) if e.kind() == std::io::ErrorKind::UnexpectedEof => break,
Err(e) => return Err(e.to_string()),
}
}
let mut proto_ok = false;
let mut aead_ok = false;
for record in records::records(&buf) {
let record = record.map_err(|e| e.to_string())?;
match record.record_type {
record_type::NEXT_PROTOCOL => {
proto_ok = record
.body
.as_chunks::<2>()
.0
.iter()
.any(|c| u16::from_be_bytes(*c) == NEXT_PROTO_NTPV4);
}
record_type::AEAD_ALGORITHM => {
aead_ok = record
.body
.as_chunks::<2>()
.0
.iter()
.any(|c| u16::from_be_bytes(*c) == AEAD_AES_SIV_CMAC_256);
}
_ => {}
}
}
if !proto_ok || !aead_ok {
let mut out = Vec::new();
records::write_record(&mut out, true, record_type::ERROR, &1u16.to_be_bytes());
records::write_record(&mut out, true, record_type::END_OF_MESSAGE, &[]);
let _ = tls.write_all(&out);
let _ = tls.flush();
return Ok(());
}
let keys = export_server_keys(&tls.conn)?;
let mut out = Vec::new();
records::write_record(
&mut out,
true,
record_type::NEXT_PROTOCOL,
&NEXT_PROTO_NTPV4.to_be_bytes(),
);
records::write_record(
&mut out,
true,
record_type::AEAD_ALGORITHM,
&AEAD_AES_SIV_CMAC_256.to_be_bytes(),
);
if ntp_port != 123 {
records::write_record(
&mut out,
false,
record_type::PORT_NEGOTIATION,
&ntp_port.to_be_bytes(),
);
}
{
let guard = state.lock().map_err(|_| "server state poisoned")?;
for _ in 0..KE_COOKIE_COUNT {
let mut nonce = [0u8; COOKIE_NONCE_LEN];
rusty_time_nts::ke::fill_random(&mut nonce).map_err(|e| e.to_string())?;
let cookie = rusty_time_nts::cookie::mint(&guard.ring, &keys, &nonce)
.map_err(|e| e.to_string())?;
records::write_record(&mut out, false, record_type::NEW_COOKIE, &cookie);
}
}
records::write_record(&mut out, true, record_type::END_OF_MESSAGE, &[]);
tls.write_all(&out).map_err(|e| e.to_string())?;
tls.flush().map_err(|e| e.to_string())?;
Ok(())
}
fn export_server_keys(conn: &rustls::ServerConnection) -> Result<NtsKeys, String> {
let context_for = |direction: u8| {
[
(NEXT_PROTO_NTPV4 >> 8) as u8,
NEXT_PROTO_NTPV4 as u8,
(AEAD_AES_SIV_CMAC_256 >> 8) as u8,
AEAD_AES_SIV_CMAC_256 as u8,
direction,
]
};
let label = b"EXPORTER-network-time-security";
let c2s = conn
.export_keying_material([0u8; 32], label, Some(&context_for(0x00)))
.map_err(|e| e.to_string())?;
let s2c = conn
.export_keying_material([0u8; 32], label, Some(&context_for(0x01)))
.map_err(|e| e.to_string())?;
Ok(NtsKeys { c2s, s2c })
}
pub struct ServerState {
pub clients: ClientTable<std::net::IpAddr>,
pub ring: KeyRing,
pub stratum: u8,
pub started_unix: u64,
}
fn ntp_serve_loop(socket: &UdpSocket, state: &Arc<Mutex<ServerState>>) {
let clock = SystemClock;
let mut bufs = vec![[0u8; 1024]; net::BATCH_SIZE];
let mut received = Vec::with_capacity(net::BATCH_SIZE);
let mut scratch = net::BatchScratch::without_timestamps();
let mut send_scratch = net::BatchScratch::without_timestamps();
let batch_send = std::env::var_os("RUSTY_TIME_NO_BATCH_SEND").is_none();
let report_batches = std::env::var_os("RUSTY_TIME_BATCH_STATS").is_some();
let mut batch_calls: u64 = 0;
let mut batch_datagrams: u64 = 0;
let mut batch_max: usize = 0;
let mut next_batch_report = std::time::Instant::now() + Duration::from_secs(1);
let mut replies: Vec<(Reply, SocketAddr)> = Vec::with_capacity(net::BATCH_SIZE);
let mut next_status = std::time::Instant::now() + STATUS_INTERVAL;
loop {
if std::time::Instant::now() >= next_status {
next_status = std::time::Instant::now() + STATUS_INTERVAL;
if let Ok(guard) = state.lock() {
let stats = guard.clients.stats;
crate::service::notify_status(&format!(
"{} requests, {} answered, {} rate-limited, {} clients",
stats.requests,
stats.responses,
stats.dropped_rate_limit,
guard.clients.len()
));
}
}
match net::wait_readable(socket, Duration::from_millis(500)) {
Ok(true) => {}
Ok(false) => continue,
Err(_) => continue,
}
let count = match net::recv_batch(socket, &mut bufs, &mut scratch, &mut received) {
Ok(n) => n,
Err(_) => continue,
};
if count == 0 {
continue;
}
if report_batches {
batch_calls += 1;
batch_datagrams += count as u64;
batch_max = batch_max.max(count);
if std::time::Instant::now() >= next_batch_report && batch_calls > 0 {
next_batch_report = std::time::Instant::now() + Duration::from_secs(1);
eprintln!(
"batch: {} receives, {} datagrams, mean {:.2}, max {}",
batch_calls,
batch_datagrams,
batch_datagrams as f64 / batch_calls as f64,
batch_max
);
}
}
let recv_ts = match clock.wall_parts() {
Ok((s, n)) => unix_parts_to_ntp(s, n),
Err(_) => continue,
};
replies.clear();
if let Ok(mut guard) = state.lock() {
for i in 0..count {
let Some(item) = received.get(i).copied() else {
break;
};
let request = &bufs[i][..item.len.min(bufs[i].len())];
if let Some(reply) = build_reply_in(request, item.peer, recv_ts, &mut guard, &clock)
{
replies.push((reply, item.peer));
}
}
}
if replies.is_empty() {
continue;
}
let sent = if batch_send {
net::send_batch_by(
socket,
&replies,
|(reply, _)| reply.bytes.as_slice(),
|(_, peer)| *peer,
&mut send_scratch,
)
.unwrap_or(0)
} else {
let mut n = 0;
for (reply, peer) in &replies {
if socket.send_to(reply.bytes.as_slice(), peer).is_err() {
break;
}
n += 1;
}
n
};
if sent > 0
&& let Ok((s, n)) = clock.wall_parts()
{
let transmit = unix_parts_to_ntp(s, n);
if let Ok(mut guard) = state.lock() {
for (reply, _) in replies.iter().take(sent) {
guard.clients.note_transmit_at(reply.handle, transmit);
}
}
}
}
}
fn kiss_of_death(request: &NtpPacket, recv_ts: NtpTimestamp) -> [u8; HEADER_LEN] {
NtpPacket {
leap: LeapIndicator::NoWarning,
version: request.version,
mode: Mode::Server,
stratum: 0, poll: request.poll,
precision: -20,
root_delay: ntp::NtpShort(0),
root_dispersion: ntp::NtpShort(0),
reference_id: *b"RATE",
reference_ts: NtpTimestamp::ZERO,
origin_ts: request.transmit_ts,
receive_ts: recv_ts,
transmit_ts: recv_ts,
}
.to_bytes()
}
fn debug_xleave() -> bool {
static ON: std::sync::OnceLock<bool> = std::sync::OnceLock::new();
*ON.get_or_init(|| std::env::var("RUSTY_TIME_DEBUG_XLEAVE").is_ok())
}
fn unix_parts_to_ntp(secs: i64, nanos: u32) -> NtpTimestamp {
NtpTimestamp::from_unix(secs, nanos)
}
pub enum ReplyBytes {
Plain([u8; HEADER_LEN]),
Extended(Vec<u8>),
}
impl ReplyBytes {
pub fn as_slice(&self) -> &[u8] {
match self {
ReplyBytes::Plain(buf) => buf,
ReplyBytes::Extended(v) => v,
}
}
}
pub struct Reply {
pub bytes: ReplyBytes,
pub handle: ClientHandle,
}
impl std::ops::Deref for Reply {
type Target = [u8];
fn deref(&self) -> &[u8] {
self.bytes.as_slice()
}
}
pub fn build_reply(
request: &[u8],
peer: SocketAddr,
recv_ts: NtpTimestamp,
state: &Arc<Mutex<ServerState>>,
clock: &SystemClock,
) -> Option<Reply> {
let mut guard = state.lock().ok()?;
build_reply_in(request, peer, recv_ts, &mut guard, clock)
}
pub fn build_reply_in(
request: &[u8],
peer: SocketAddr,
recv_ts: NtpTimestamp,
guard: &mut ServerState,
clock: &SystemClock,
) -> Option<Reply> {
if request.len() < HEADER_LEN {
guard.clients.note_refused();
return None;
}
let parsed = match NtpPacket::parse(request) {
Ok(p) => p,
Err(_) => {
guard.clients.note_refused();
return None;
}
};
if parsed.mode != Mode::Client {
guard.clients.note_refused();
return None;
}
let now_mono = clock.mono_s().ok()?;
let key = client_key(peer);
let locked = {
let (disposition, handle) = guard.clients.admit_handle(&key, now_mono);
if disposition != Disposition::Respond {
(
disposition,
handle,
ResponseMode::Basic,
0u8,
recv_ts,
recv_ts,
)
} else {
let mode = guard.clients.response_mode_at(handle, parsed.origin_ts);
let stratum = guard.stratum;
let (mut receive_field, mut transmit_field) = match mode {
ResponseMode::Basic => (recv_ts, recv_ts),
ResponseMode::Interleaved { prev_transmit } => (recv_ts, prev_transmit),
};
rusty_time_core::server::mark_server_timestamps(
&mut receive_field,
&mut transmit_field,
);
guard
.clients
.note_response_at(handle, recv_ts, receive_field);
(
disposition,
handle,
mode,
stratum,
receive_field,
transmit_field,
)
}
};
let (disposition, handle, mode, stratum, receive_field, transmit_field) = locked;
match disposition {
Disposition::Respond => {}
Disposition::KissOfDeath => {
return Some(Reply {
bytes: ReplyBytes::Plain(kiss_of_death(&parsed, recv_ts)),
handle,
});
}
Disposition::Drop => return None,
}
let mut cookie: Option<&[u8]> = None;
let mut unique_id: Option<&[u8]> = None;
let mut placeholders = 0usize;
let mut auth_field: Option<ef::Field> = None;
for field in ef::fields(request) {
match field.field_type {
ef::field_type::NTS_COOKIE => cookie = Some(field.body),
ef::field_type::UNIQUE_IDENTIFIER => unique_id = Some(field.body),
ef::field_type::NTS_COOKIE_PLACEHOLDER => placeholders += 1,
ef::field_type::NTS_AUTHENTICATOR => {
auth_field = Some(field);
break;
}
_ => {}
}
}
if matches!(mode, ResponseMode::Interleaved { .. }) && debug_xleave() {
eprintln!(
"xleave: reported_tx_age={:+.6}s (should be ~1 poll interval)",
recv_ts.seconds_since(transmit_field),
);
}
let origin_field = match mode {
ResponseMode::Basic => parsed.transmit_ts,
ResponseMode::Interleaved { .. } => parsed.receive_ts,
};
let mut header = NtpPacket {
leap: LeapIndicator::NoWarning,
version: parsed.version,
mode: Mode::Server,
stratum,
poll: parsed.poll,
precision: -20,
root_delay: ntp::NtpShort(0),
root_dispersion: ntp::NtpShort::from_seconds(0.000_1),
reference_id: *b"RSTY",
reference_ts: recv_ts,
origin_ts: origin_field,
receive_ts: receive_field,
transmit_ts: transmit_field,
};
let Some(auth_field) = auth_field else {
if mode == ResponseMode::Basic {
let mut tx = clock
.wall_parts()
.ok()
.map(|(s, n)| unix_parts_to_ntp(s, n))?;
let mut rx = header.receive_ts;
rusty_time_core::server::mark_server_timestamps(&mut rx, &mut tx);
header.transmit_ts = tx;
}
return Some(Reply {
bytes: ReplyBytes::Plain(header.to_bytes()),
handle,
});
};
let (cookie, unique_id) = (cookie?, unique_id?);
let keys = rusty_time_nts::cookie::redeem(&guard.ring, cookie).ok()?;
verify_client_authenticator(request, auth_field, &keys.c2s)?;
let want = (1 + placeholders).min(MAX_REPLY_COOKIES);
let mut plaintext = Vec::with_capacity(want * COOKIE_FIELD_HINT);
let mut nonce_bytes = [0u8; COOKIE_NONCE_LEN * MAX_REPLY_COOKIES];
rusty_time_nts::ke::fill_random(&mut nonce_bytes[..want * COOKIE_NONCE_LEN]).ok()?;
let mut nonces = [[0u8; COOKIE_NONCE_LEN]; MAX_REPLY_COOKIES];
for (i, slot) in nonces.iter_mut().enumerate().take(want) {
slot.copy_from_slice(&nonce_bytes[i * COOKIE_NONCE_LEN..(i + 1) * COOKIE_NONCE_LEN]);
}
{
rusty_time_nts::cookie::mint_fields_into(
&guard.ring,
&keys,
&nonces[..want],
&mut plaintext,
)
.ok()?;
}
if mode == ResponseMode::Basic {
let mut tx = clock
.wall_parts()
.ok()
.map(|(s, n)| unix_parts_to_ntp(s, n))?;
let mut rx = header.receive_ts;
rusty_time_core::server::mark_server_timestamps(&mut rx, &mut tx);
header.transmit_ts = tx;
}
let mut reply = Vec::with_capacity(HEADER_LEN + 64 + want * COOKIE_FIELD_HINT + 64);
reply.extend_from_slice(&header.to_bytes());
let mut uid = [0u8; UNIQUE_ID_LEN];
let n = unique_id.len().min(UNIQUE_ID_LEN);
uid[..n].copy_from_slice(&unique_id[..n]);
ef::write_field(&mut reply, ef::field_type::UNIQUE_IDENTIFIER, &uid[..n]);
let mut nonce = [0u8; NONCE_LEN];
rusty_time_nts::ke::fill_random(&mut nonce).ok()?;
let ciphertext = rusty_time_nts::aead::seal(&keys.s2c, &[&reply, &nonce], &plaintext).ok()?;
ef::write_authenticator(&mut reply, &nonce, &ciphertext);
Some(Reply {
bytes: ReplyBytes::Extended(reply),
handle,
})
}
#[cfg(test)]
#[path = "server_tests.rs"]
mod tests;
fn verify_client_authenticator(request: &[u8], auth: ef::Field<'_>, c2s: &[u8; 32]) -> Option<()> {
if auth.body.len() < 4 {
return None;
}
let nonce_len = u16::from_be_bytes([auth.body[0], auth.body[1]]) as usize;
let ct_len = u16::from_be_bytes([auth.body[2], auth.body[3]]) as usize;
let ct_start = 4 + nonce_len.next_multiple_of(4);
let ct_end = ct_start.checked_add(ct_len)?;
if nonce_len == 0 || ct_end > auth.body.len() {
return None;
}
let nonce = auth.body.get(4..4 + nonce_len)?;
let ciphertext = auth.body.get(ct_start..ct_end)?;
let aad = request.get(..auth.offset)?;
rusty_time_nts::aead::open(c2s, &[aad, nonce], ciphertext).ok()?;
Some(())
}