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 | |
| 4 | use oxedyne_fe2o3_core::prelude::*; |
| 5 | use crate::{ |
| 6 | time::{CalClock, CalClockZone, CalClockDuration}, |
| 7 | index::{time_basis::{TimeIndexInterval, TimeIndex}, TimeLong}, |
| 8 | }; |
| 9 | use 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. |
| 17 | pub 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)] |
| 29 | pub enum NtpAlgorithm { |
| 30 | Rfc1305, |
| 31 | Prodigitum, // placeholder |
| 32 | } |
| 33 | |
| 34 | #[derive(Debug, Clone)] |
| 35 | struct NtpServerConnection { |
| 36 | address: IpAddr, |
| 37 | socket_addr: SocketAddr, |
| 38 | timeout: Duration, |
| 39 | } |
| 40 | |
| 41 | #[derive(Debug, Clone)] |
| 42 | pub 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 | |
| 52 | impl 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)] |
| 83 | pub 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)] |
| 113 | pub 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)] |
| 125 | pub struct NtpOffsets { |
| 126 | values: Vec<f64>, |
| 127 | mean: f64, |
| 128 | std_dev: f64, |
| 129 | } |
| 130 | |
| 131 | impl 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 | |
| 183 | impl 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 | } |