Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_datime/src/time/ntp_advanced.rs

17.0 KiB, 101 runs

created by r1870400018:8570, which is this file's identity for as long as the history lasts, whatever it is later renamed to

download · who wrote it · its history

1//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
2//! Anthropic Claude
3
4use oxedyne_fe2o3_core::prelude::*;
5use crate::{
6 time::{CalClock, CalClockZone, CalClockDuration},
7 index::{time_basis::{TimeIndexInterval, TimeIndex}, TimeLong},
8};
9use std::{
10 collections::{HashMap, BTreeMap},
11 net::{IpAddr, SocketAddr, UdpSocket},
12 time::{Duration, Instant},
13};
14
15/// Polls several servers and applies the RFC 1305 intersection, selection and
16/// combine algorithms, so one bad server does not carry the answer.
17pub struct AdvancedNtpClient {
18 servers: Vec<NtpServerConnection>,
19 server_stats: HashMap<IpAddr, NtpServerStats>,
20 offsets: HashMap<IpAddr, NtpOffsets>,
21 algorithm: NtpAlgorithm,
22 min_servers: usize,
23 min_true_chimers: usize,
24 timeout: Duration,
25 best_server: Option<IpAddr>,
26}
27
28#[derive(Debug, Clone, PartialEq)]
29pub enum NtpAlgorithm {
30 Rfc1305,
31 Prodigitum, // placeholder
32}
33
34#[derive(Debug, Clone)]
35struct NtpServerConnection {
36 address: IpAddr,
37 socket_addr: SocketAddr,
38 timeout: Duration,
39}
40
41#[derive(Debug, Clone)]
42pub struct NtpData {
43 address: IpAddr,
44 offset: CalClockDuration,
45 delay: CalClockDuration,
46 root_distance: i64,
47 correctness_interval: TimeIndexInterval<TimeLong>,
48 stratum: u8,
49 precision: i8,
50}
51
52impl NtpData {
53 pub fn address(&self) -> IpAddr {
54 self.address
55 }
56
57 pub fn offset(&self) -> &CalClockDuration {
58 &self.offset
59 }
60
61 pub fn delay(&self) -> &CalClockDuration {
62 &self.delay
63 }
64
65 pub fn root_distance(&self) -> i64 {
66 self.root_distance
67 }
68
69 pub fn correctness_interval(&self) -> &TimeIndexInterval<TimeLong> {
70 &self.correctness_interval
71 }
72
73 pub fn stratum(&self) -> u8 {
74 self.stratum
75 }
76
77 pub fn precision(&self) -> i8 {
78 self.precision
79 }
80}
81
82#[derive(Debug, Clone)]
83pub struct NtpServerStats {
84 #[allow(dead_code)]
85 offset_history: Vec<f64>,
86 #[allow(dead_code)]
87 delay_history: Vec<f64>,
88 #[allow(dead_code)]
89 mean_offset: f64,
90 #[allow(dead_code)]
91 std_dev_offset: f64,
92 #[allow(dead_code)]
93 mean_delay: f64,
94 #[allow(dead_code)]
95 std_dev_delay: f64,
96 #[allow(dead_code)]
97 sample_count: usize,
98 #[allow(dead_code)]
99 root_distance: f64,
100 #[allow(dead_code)]
101 stratum: u8,
102 #[allow(dead_code)]
103 precision: i8,
104 #[allow(dead_code)]
105 last_poll_time: Option<Instant>,
106 #[allow(dead_code)]
107 is_reachable: bool,
108 #[allow(dead_code)]
109 falseticker_count: usize,
110}
111
112#[derive(Debug, Clone)]
113pub struct AdvancedNtpResult {
114 pub network_time: CalClock,
115 pub offset: CalClockDuration, // network_time - local_time
116 pub delay: CalClockDuration,
117 pub root_distance: f64, // accuracy estimate
118 pub servers_used: usize,
119 pub servers_total: usize,
120 pub best_server: Option<IpAddr>,
121 pub jitter: f64,
122}
123
124#[derive(Debug, Clone)]
125pub struct NtpOffsets {
126 values: Vec<f64>,
127 mean: f64,
128 std_dev: f64,
129}
130
131impl NtpOffsets {
132 pub fn new() -> Self {
133 Self {
134 values: Vec::new(),
135 mean: 0.0,
136 std_dev: 0.0,
137 }
138 }
139
140 pub fn add(&mut self, value: f64) {
141 self.values.push(value);
142
143 // Keep only last 8 samples
144 if self.values.len() > 8 {
145 self.values.remove(0);
146 }
147
148 self.recalculate_stats();
149 }
150
151 fn recalculate_stats(&mut self) {
152 if self.values.is_empty() {
153 return;
154 }
155
156 self.mean = self.values.iter().sum::<f64>() / self.values.len() as f64;
157
158 if self.values.len() > 1 {
159 let variance: f64 = self.values
160 .iter()
161 .map(|x| (x - self.mean).powi(2))
162 .sum::<f64>() / (self.values.len() - 1) as f64;
163 self.std_dev = variance.sqrt();
164 } else {
165 self.std_dev = 0.0;
166 }
167 }
168
169 pub fn std_dev(&self) -> f64 {
170 self.std_dev
171 }
172
173 pub fn mean(&self) -> f64 {
174 self.mean
175 }
176
177 pub fn count(&self) -> usize {
178 self.values.len()
179 }
180}
181
182
183impl AdvancedNtpClient {
184 pub fn new(servers: Vec<IpAddr>) -> Outcome<Self> {
185 if servers.len() < 4 {
186 return Err(err!(
187 "Advanced NTP requires at least 4 servers, got {}",
188 servers.len();
189 Invalid, Input
190 ));
191 }
192
193 let server_connections = servers
194 .into_iter()
195 .map(|addr| NtpServerConnection {
196 address: addr,
197 socket_addr: SocketAddr::new(addr, 123),
198 timeout: Duration::from_secs(2),
199 })
200 .collect();
201
202 Ok(Self {
203 servers: server_connections,
204 server_stats: HashMap::new(),
205 offsets: HashMap::new(),
206 algorithm: NtpAlgorithm::Rfc1305,
207 min_servers: 4,
208 min_true_chimers: 2,
209 timeout: Duration::from_secs(2),
210 best_server: None,
211 })
212 }
213
214 pub fn with_algorithm(mut self, algorithm: NtpAlgorithm) -> Self {
215 self.algorithm = algorithm;
216 self
217 }
218
219 pub fn with_min_servers(mut self, min_servers: usize) -> Self {
220 self.min_servers = min_servers;
221 self
222 }
223
224 pub fn with_timeout(mut self, timeout: Duration) -> Self {
225 self.timeout = timeout;
226 for server in &mut self.servers {
227 server.timeout = timeout;
228 }
229 self
230 }
231
232 pub fn synchronize(&mut self, zone: CalClockZone) -> Outcome<AdvancedNtpResult> {
233 // Step 1: Poll all servers
234 let data_list = res!(self.poll_all_servers());
235
236 if data_list.is_empty() {
237 return Err(err!("Could not get any NTP data from any of the given servers"; Network));
238 }
239
240 // Step 2: Apply RFC 1305 algorithms
241 match self.algorithm {
242 NtpAlgorithm::Rfc1305 => self.rfc1305_algorithm(data_list, zone),
243 NtpAlgorithm::Prodigitum => self.prodigitum_algorithm(data_list, zone),
244 }
245 }
246
247 fn poll_all_servers(&mut self) -> Outcome<Vec<NtpData>> {
248 let mut data_list = Vec::new();
249 let servers = self.servers.clone(); // Clone to avoid borrowing issues
250
251 for server in &servers {
252 match self.poll_single_server(server) {
253 Ok(data) => {
254 // Update offset history
255 let offset_ns = data.offset.nanoseconds();
256 let offset_series = self.offsets.entry(data.address).or_insert_with(NtpOffsets::new);
257 offset_series.add(offset_ns as f64);
258
259 data_list.push(data);
260 },
261 Err(_e) => {
262 // Mark server as unreachable and continue
263 self.mark_server_unreachable(&server.address);
264 }
265 }
266 }
267
268 Ok(data_list)
269 }
270
271 fn poll_single_server(&self, server: &NtpServerConnection) -> Outcome<NtpData> {
272 // Create UDP socket
273 let socket = res!(UdpSocket::bind("0.0.0.0:0").map_err(|e| {
274 err!("Failed to create UDP socket: {}", e; Network)
275 }));
276
277 res!(socket.set_read_timeout(Some(server.timeout)).map_err(|e| {
278 err!("Failed to set socket timeout: {}", e; Network)
279 }));
280
281 // Create NTP request packet
282 let request_packet = self.create_ntp_request();
283 let send_time = Instant::now();
284
285 // Send request
286 res!(socket.send_to(&request_packet, server.socket_addr).map_err(|e| {
287 err!("Failed to send NTP request to {}: {}", server.address, e; Network)
288 }));
289
290 // Receive response
291 let mut buffer = [0u8; 48];
292 let recv_time = Instant::now();
293
294 let (bytes_received, _) = res!(socket.recv_from(&mut buffer).map_err(|e| {
295 err!("Failed to receive NTP response from {}: {}", server.address, e; Network)
296 }));
297
298 if bytes_received != 48 {
299 return Err(err!(
300 "Invalid NTP response size from {}: {} bytes",
301 server.address,
302 bytes_received;
303 Invalid, Network
304 ));
305 }
306
307 // Parse response and calculate timing
308 self.parse_ntp_response(&buffer, send_time, recv_time, server.address)
309 }
310
311 fn create_ntp_request(&self) -> [u8; 48] {
312 let mut packet = [0u8; 48];
313
314 // Set LI (0), VN (4), Mode (3) - Client mode
315 packet[0] = 0x1b; // 00 100 011
316
317 // Set poll interval (6 = 64 seconds)
318 packet[2] = 6;
319
320 // Set precision (-20 = microsecond precision)
321 packet[3] = 0xec; // -20 in two's complement
322
323 // Transmit timestamp (current time)
324 let now = std::time::SystemTime::now()
325 .duration_since(std::time::UNIX_EPOCH)
326 .unwrap_or_default();
327
328 let ntp_timestamp = (now.as_secs() + 2208988800) as u64; // Convert to NTP epoch
329 let fraction = ((now.subsec_nanos() as u64) << 32) / 1_000_000_000;
330
331 // Write transmit timestamp (bytes 40-47)
332 packet[40..44].copy_from_slice(&(ntp_timestamp as u32).to_be_bytes());
333 packet[44..48].copy_from_slice(&(fraction as u32).to_be_bytes());
334
335 packet
336 }
337
338 fn parse_ntp_response(
339 &self,
340 packet: &[u8; 48],
341 send_time: Instant,
342 recv_time: Instant,
343 server_addr: IpAddr,
344 ) -> Outcome<NtpData> {
345 // Extract NTP packet fields
346 let stratum = packet[1];
347 let precision = packet[3] as i8;
348
349 // Extract timestamps
350 let root_delay = u32::from_be_bytes([packet[4], packet[5], packet[6], packet[7]]) as f64 / 65536.0;
351 let root_dispersion = u32::from_be_bytes([packet[8], packet[9], packet[10], packet[11]]) as f64 / 65536.0;
352
353 // Extract server timestamps
354 let _reference_ts = self.extract_ntp_timestamp(&packet[16..24]);
355 let _originate_ts = self.extract_ntp_timestamp(&packet[24..32]);
356 let receive_ts = self.extract_ntp_timestamp(&packet[32..40]);
357 let transmit_ts = self.extract_ntp_timestamp(&packet[40..48]);
358
359 // Calculate timing
360 let rtt = recv_time.duration_since(send_time).as_secs_f64();
361 let delay = rtt;
362
363 // Calculate offset using NTP algorithm
364 let t1 = send_time.elapsed().as_secs_f64(); // Local send time
365 let t2 = receive_ts; // Server receive time
366 let t3 = transmit_ts; // Server transmit time
367 let t4 = recv_time.elapsed().as_secs_f64(); // Local receive time
368
369 let offset_seconds = ((t2 - t1) + (t3 - t4)) / 2.0;
370
371 // Convert to CalClockDuration
372 let offset = CalClockDuration::from_seconds(offset_seconds as i64);
373 let ntp_delay = CalClockDuration::from_seconds(delay as i64);
374
375 // Calculate root distance in nanoseconds
376 let root_distance_ns = ((root_delay + root_dispersion + delay) * 1_000_000_000.0) as i64;
377
378 // Calculate correctness interval in nanoseconds
379 let offset_ns = (offset_seconds * 1_000_000_000.0) as i64;
380 let half_delay_ns = (delay * 1_000_000_000.0 / 2.0) as i64;
381
382 let interval_start = TimeIndex::new(TimeLong::new(offset_ns - half_delay_ns));
383 let interval_finish = TimeIndex::new(TimeLong::new(offset_ns + half_delay_ns));
384 let correctness_interval = res!(TimeIndexInterval::new(interval_start, interval_finish));
385
386 Ok(NtpData {
387 address: server_addr,
388 offset,
389 delay: ntp_delay,
390 root_distance: root_distance_ns,
391 correctness_interval,
392 stratum,
393 precision,
394 })
395 }
396
397 fn extract_ntp_timestamp(&self, bytes: &[u8]) -> f64 {
398 let seconds = u32::from_be_bytes([bytes[0], bytes[1], bytes[2], bytes[3]]) as u64;
399 let fraction = u32::from_be_bytes([bytes[4], bytes[5], bytes[6], bytes[7]]) as u64;
400
401 // Convert from NTP epoch (1900) to Unix epoch (1970)
402 let unix_seconds = seconds.saturating_sub(2208988800);
403 let fractional_seconds = fraction as f64 / (1u64 << 32) as f64;
404
405 unix_seconds as f64 + fractional_seconds
406 }
407
408 #[allow(dead_code)]
409 fn calculate_jitter(&self, server_addr: IpAddr, current_delay: f64) -> f64 {
410 if let Some(stats) = self.server_stats.get(&server_addr) {
411 if stats.delay_history.len() > 1 {
412 return stats.std_dev_delay;
413 }
414 }
415 current_delay * 0.1 // Default estimate
416 }
417
418 fn rfc1305_algorithm(
419 &mut self,
420 data_list: Vec<NtpData>,
421 zone: CalClockZone,
422 ) -> Outcome<AdvancedNtpResult> {
423 // Step 1: Intersection Algorithm
424 let intersection_interval = res!(self.get_intersection_interval(&data_list));
425 let true_chimers = res!(self.get_true_chimers(&data_list, &intersection_interval));
426
427 if true_chimers.len() < self.min_true_chimers {
428 return Err(err!(
429 "Only {} true chimers found, minimum {} required",
430 true_chimers.len(),
431 self.min_true_chimers;
432 Invalid, Network
433 ));
434 }
435
436 // Step 2: Selection Algorithm
437 let survivors = res!(self.get_survivors(&true_chimers));
438
439 if survivors.is_empty() {
440 return Err(err!("No servers survived selection algorithm"; Invalid, Network));
441 }
442
443 // Step 3: Combine Algorithm - sort by root distance
444 let mut sorted_survivors = survivors;
445 sorted_survivors.sort_by(|a, b| a.root_distance.cmp(&b.root_distance));
446
447 let best_data = &sorted_survivors[0];
448
449 // Convert to CalClock
450 let network_time = res!(self.data_to_calclock(best_data, zone));
451
452 // Calculate results
453 let local_time = res!(CalClock::now(network_time.zone().clone()));
454 let offset = res!(network_time.duration_until(&local_time));
455
456 self.best_server = Some(best_data.address);
457
458 Ok(AdvancedNtpResult {
459 network_time,
460 offset,
461 delay: best_data.delay.clone(),
462 root_distance: best_data.root_distance as f64,
463 servers_used: 1,
464 servers_total: self.servers.len(),
465 best_server: Some(best_data.address),
466 jitter: 0.0, // Calculated from offset history
467 })
468 }
469
470
471 fn mark_server_unreachable(&mut self, address: &IpAddr) {
472 if let Some(stats) = self.server_stats.get_mut(address) {
473 stats.is_reachable = false;
474 stats.falseticker_count += 1;
475 }
476 }
477
478 pub fn server_statistics(&self) -> &HashMap<IpAddr, NtpServerStats> {
479 &self.server_stats
480 }
481
482 pub fn best_server(&self) -> Option<IpAddr> {
483 self.best_server
484 }
485
486 fn get_intersection_interval(&self, data_list: &[NtpData]) -> Outcome<TimeIndexInterval<TimeLong>> {
487
488 let m = self.servers.len();
489 let mut tuples = BTreeMap::new();
490
491 // Create ordered tuple map from intervals
492 for data in data_list {
493 let offset = data.offset.nanoseconds();
494 let interval = data.correctness_interval();
495
496 // Add interval start (-1), offset (0), and interval finish (+1)
497 tuples.insert(interval.start().time().value(), -1);
498 tuples.insert(offset, 0);
499 tuples.insert(interval.finish().time().value(), 1);
500 }
501
502 let mut f = 0; // Number of falsetickers
503
504 loop {
505 #[allow(unused_assignments)]
506 let mut lower = 0i64; // Will be overwritten before use
507 #[allow(unused_assignments)]
508 let mut upper = 0i64; // Will be overwritten before use
509 let mut endcount = 0;
510 let mut midcount = 0;
511
512 // Go forwards
513 for (key, value) in &tuples {
514 endcount -= value;
515
516 if endcount >= (m - f) as i32 {
517 lower = *key;
518
519 // Now go backwards
520 endcount = 0;
521 for (key2, value2) in tuples.iter().rev() {
522 endcount += value2;
523
524 if endcount >= (m - f) as i32 {
525 upper = *key2;
526 if lower <= upper && midcount <= f as i32 {
527 let start = TimeIndex::new(TimeLong::new(lower));
528 let finish = TimeIndex::new(TimeLong::new(upper));
529 return TimeIndexInterval::new(start, finish);
530 }
531 }
532
533 if *value2 == 0 {
534 midcount += 1;
535 }
536 }
537 break;
538 }
539
540 if *value == 0 {
541 midcount += 1;
542 }
543 }
544
545 f += 1;
546
547 if f >= m / 2 {
548 return Err(err!("Number of falsetickers {} is at least half the number {} of NTP servers polled. Formal correctness cannot be achieved.", f, m; Invalid, Network));
549 }
550 }
551 }
552
553 fn get_true_chimers(&self, data_list: &[NtpData], intersection_interval: &TimeIndexInterval<TimeLong>) -> Outcome<Vec<NtpData>> {
554 let mut result = Vec::new();
555
556 for data in data_list {
557 let interval = data.correctness_interval();
558
559 // Check if interval overlaps with intersection interval
560 if !interval.overlaps(intersection_interval) {
561 continue;
562 }
563
564 result.push(data.clone());
565 }
566
567 Ok(result)
568 }
569
570 fn get_survivors(&self, true_chimers: &[NtpData]) -> Outcome<Vec<NtpData>> {
571 let mut result = true_chimers.to_vec();
572
573 if result.len() <= self.min_true_chimers {
574 return Ok(result);
575 }
576
577 let mut have_pruned = true;
578 while have_pruned && result.len() > self.min_true_chimers {
579 have_pruned = false;
580 let mut peer_jitter = BTreeMap::new();
581 let mut select_jitter = BTreeMap::new();
582
583 for data1 in &result {
584 let theta1 = data1.offset.nanoseconds();
585 let lambda1 = data1.root_distance;
586
587 // Calculate peer jitter (standard deviation from offset history)
588 let offset_series = self.offsets.get(&data1.address);
589 let peer_jitter_value = if let Some(series) = offset_series {
590 series.std_dev() as i64
591 } else {
592 0
593 };
594 peer_jitter.insert(peer_jitter_value, data1.clone());
595
596 // Calculate select jitter (RMS of differences)
597 let mut d = Vec::new();
598 for data2 in true_chimers {
599 let theta2 = data2.offset.nanoseconds();
600 d.push(((theta1 - theta2).abs() * lambda1) as i64);
601 }
602
603 let select_jitter_value = self.rms(&d) as i64;
604 select_jitter.insert(select_jitter_value, data1.clone());
605 }
606
607 // Get maximum select jitter and minimum peer jitter
608 if let (Some((phi_s_max, worst_select)), Some((phi_r_min, _))) =
609 (select_jitter.iter().rev().next(), peer_jitter.iter().next()) {
610
611 if phi_s_max > phi_r_min {
612 // Remove the worst server
613 result.retain(|data| data.address != worst_select.address);
614 have_pruned = true;
615 }
616 }
617 }
618
619 Ok(result)
620 }
621
622 fn rms(&self, series: &[i64]) -> f64 {
623 if series.is_empty() {
624 return 0.0;
625 }
626
627 let mut sum_x2 = 0i64;
628 for &x in series {
629 sum_x2 += x * x; // Fixed: was overwriting instead of accumulating
630 }
631
632 ((sum_x2 as f64) / (series.len() as f64)).sqrt()
633 }
634
635 fn data_to_calclock(&self, data: &NtpData, zone: CalClockZone) -> Outcome<CalClock> {
636 let current_time = ok!(std::time::SystemTime::now()
637 .duration_since(std::time::UNIX_EPOCH)
638 .map_err(|e| err!("System time error: {}", e; System)));
639
640 let offset_seconds = data.offset.total_seconds() as f64;
641 let adjusted_time = current_time.as_secs_f64() + offset_seconds;
642 CalClock::from_unix_timestamp_seconds(adjusted_time as i64, zone)
643 }
644
645 /// Placeholder for Prodigitum algorithm.
646 fn prodigitum_algorithm(
647 &mut self,
648 data_list: Vec<NtpData>,
649 zone: CalClockZone,
650 ) -> Outcome<AdvancedNtpResult> {
651 // For now, fall back to RFC 1305
652 self.rfc1305_algorithm(data_list, zone)
653 }
654}