From daea781f85f0b48e50274a9a14fee6ee08caabe4 Mon Sep 17 00:00:00 2001 From: "Jesper L. Nielsen" Date: Thu, 1 Oct 2026 10:58:43 +0200 Subject: [PATCH] fix: keep the daemon's timers on a monotonic clock MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Every deadline the daemon keeps — record TTLs and refresh points, probe and announce schedules, retransmissions, delayed responses, resolver timeouts and the periodic IP check — was a millisecond value derived from `SystemTime`. A step of the system clock therefore moved all of them: a backwards step parked them until the wall clock had caught up, so a daemon that had registered its services just before the step never announced them and never noticed a new address; a forwards step fired them all at once and made every cached record look expired. Seen on an embedded device whose RTC driver loads a few seconds after the responder starts. The kernel re-seats the wall clock from the RTC at that point, and with an RTC holding a date in the past the responder went silent on every interface for the rest of the uptime, while `register()` had already reported success. Use `std::time::Instant` for points in time and `Duration` for the deltas between them. `current_time_millis()` is removed; the fields and parameters that carried its value are `Instant` now, and the `0` sentinels that meant "no time" are `Option`. Differences use the standard saturating operations, so the two places that could underflow when a record was already past its expiry are covered as well. TTLs on the wire are still remaining seconds computed from the record's creation, so packets are unchanged. Nothing in the crate needs wall-clock time. Co-Authored-By: Claude Fable 5.1 --- src/dns_cache.rs | 51 +++++++------ src/dns_parser.rs | 134 ++++++++++++++++++---------------- src/lib.rs | 10 --- src/service_daemon.rs | 166 ++++++++++++++++++++++-------------------- src/service_info.rs | 82 ++++++++++++--------- 5 files changed, 235 insertions(+), 208 deletions(-) diff --git a/src/dns_cache.rs b/src/dns_cache.rs index 395ada05..5dabe9ea 100644 --- a/src/dns_cache.rs +++ b/src/dns_cache.rs @@ -5,7 +5,6 @@ #[cfg(feature = "logging")] use crate::log::{debug, trace}; use crate::{ - current_time_millis, dns_parser::{DnsAddress, DnsPointer, DnsRecordBox, DnsSrv, InterfaceId, RRType}, service_info::{split_sub_domain, MyIntf}, ScopedIp, @@ -13,6 +12,7 @@ use crate::{ use std::{ collections::{HashMap, HashSet}, ops::BitOr, + time::{Duration, Instant}, }; /// Bitflags-style type for filtering by IP version. @@ -184,7 +184,7 @@ impl DnsCache { pub(crate) fn service_verify_queries( &mut self, instance: &str, - expire_at: Option, + expire_at: Option, ) -> Vec<(String, RRType)> { let Some(srv_vec) = self.srv.get_mut(instance) else { return Vec::new(); @@ -227,7 +227,7 @@ impl DnsCache { &mut self, intf: &MyIntf, incoming: DnsRecordBox, - timers: &mut Vec, + timers: &mut Vec, is_for_us: bool, ) -> Option<(&DnsRecordIntf, bool)> { let entry_name = incoming.get_name().to_string(); @@ -336,7 +336,10 @@ impl DnsCache { /// Iterates all ADDR records and remove ones that expired. /// Returns the expired ones in a map of names and addresses. - pub(crate) fn evict_expired_addr(&mut self, now: u64) -> HashMap> { + pub(crate) fn evict_expired_addr( + &mut self, + now: Instant, + ) -> HashMap> { let mut removed = HashMap::new(); self.addr.retain(|_, records| { @@ -364,7 +367,10 @@ impl DnsCache { /// returns the set of expired instance names for each ty_domain. /// /// An instance in the returned set indicates its PTR and/or SRV record has expired. - pub(crate) fn evict_expired_services(&mut self, now: u64) -> HashMap> { + pub(crate) fn evict_expired_services( + &mut self, + now: Instant, + ) -> HashMap> { let mut expired_instances = HashMap::new(); // Check all ty_domain in the cache by following all PTR records, regardless @@ -477,8 +483,8 @@ impl DnsCache { /// Checks refresh due for PTR records of `ty_domain`. /// Returns all updated refresh time. - pub(crate) fn refresh_due_ptr(&mut self, ty_domain: &str) -> HashSet { - let now = current_time_millis(); + pub(crate) fn refresh_due_ptr(&mut self, ty_domain: &str) -> HashSet { + let now = Instant::now(); // Check all PTR records for this ty_domain. self.ptr @@ -496,8 +502,8 @@ impl DnsCache { pub(crate) fn refresh_due_srv_txt( &mut self, ty_domain: &str, - ) -> (HashMap>, HashSet) { - let now = current_time_millis(); + ) -> (HashMap>, HashSet) { + let now = Instant::now(); let instances: Vec<_> = self .ptr @@ -518,7 +524,7 @@ impl DnsCache { let mut new_timers = HashSet::new(); for instance in instances { // Check SRV records. - let refresh_timers: HashSet = self + let refresh_timers: HashSet = self .srv .get_mut(instance) .into_iter() @@ -535,7 +541,7 @@ impl DnsCache { } // Check TXT records. - let refresh_timers: HashSet = self + let refresh_timers: HashSet = self .txt .get_mut(instance) .into_iter() @@ -557,8 +563,11 @@ impl DnsCache { /// Returns the set of `host`, where refreshing the A / AAAA records is due /// for a `ty_domain`. - pub(crate) fn refresh_due_hosts(&mut self, ty_domain: &str) -> (HashSet, HashSet) { - let now = current_time_millis(); + pub(crate) fn refresh_due_hosts( + &mut self, + ty_domain: &str, + ) -> (HashSet, HashSet) { + let now = Instant::now(); let instances: Vec<_> = self .ptr @@ -597,7 +606,7 @@ impl DnsCache { let mut refresh_due = HashSet::new(); let mut new_timers = HashSet::new(); for hostname in hostnames_browsed { - let refresh_timers: HashSet = self + let refresh_timers: HashSet = self .addr .get_mut(&hostname.to_lowercase()) .into_iter() @@ -620,7 +629,7 @@ impl DnsCache { &mut self, hostname: &str, ) -> HashSet<(String, ScopedIp)> { - let now = current_time_millis(); + let now = Instant::now(); self.addr .get_mut(hostname) @@ -654,7 +663,7 @@ impl DnsCache { &'a self, name: &str, qtype: RRType, - now: u64, + now: Instant, ) -> Vec<&'a DnsRecordIntf> { let records_opt = match qtype { RRType::PTR => self.get_ptr(name), @@ -829,9 +838,9 @@ impl DnsCache { fn apply_cache_flush( incoming: &DnsRecordBox, existing_records: &mut [DnsRecordIntf], - timers: &mut Vec, + timers: &mut Vec, ) { - let now = current_time_millis(); + let now = Instant::now(); let class = incoming.get_class(); let rtype = incoming.get_type(); @@ -847,8 +856,8 @@ fn apply_cache_flush( if class == r.record.get_class() && rtype == r.record.get_type() - && now > r.record.get_created() + 1000 - && r.record.get_expire() > now + 1000 + && now > r.record.get_created() + Duration::from_millis(1000) + && r.record.get_expire() > now + Duration::from_millis(1000) { should_flush = true; @@ -864,7 +873,7 @@ fn apply_cache_flush( if should_flush { trace!("FLUSH one record: {:?}", &r.record); - let new_expire = now + 1000; + let new_expire = now + Duration::from_millis(1000); r.record.set_expire(new_expire); // Add a timer so the run loop will handle this expire. diff --git a/src/dns_parser.rs b/src/dns_parser.rs index ca79e3f3..631f1ea8 100644 --- a/src/dns_parser.rs +++ b/src/dns_parser.rs @@ -7,7 +7,6 @@ #[cfg(feature = "logging")] use crate::log::{debug, trace}; -use crate::current_time_millis; use crate::error::{e_fmt, Error, Result}; use crate::service_info::{decode_txt, is_unicast_link_local, DnsRegistry, MyIntf, ServiceInfo}; @@ -25,6 +24,7 @@ use std::{ hash::Hash, net::{IpAddr, Ipv4Addr, Ipv6Addr}, str, + time::{Duration, Instant}, }; /// Represents a network interface identifier defined by the OS. @@ -455,13 +455,15 @@ impl DnsEntryExt for DnsQuestion { #[derive(Debug, Clone)] pub struct DnsRecord { pub(crate) entry: DnsEntry, - ttl: u32, // in seconds, 0 means this record should not be cached - created: u64, // UNIX time in millis - expires: u64, // expires at this UNIX time in millis + ttl: u32, // in seconds, 0 means this record should not be cached + /// When this record was created (received or registered). + created: Instant, + /// When this record expires. + expires: Instant, /// Support re-query an instance before its PTR record expires. /// See https://datatracker.ietf.org/doc/html/rfc6762#section-5.2 - refresh: u64, // UNIX time in millis + refresh: Instant, /// If conflict resolution decides to change the name, this is the new one. new_name: Option, @@ -469,7 +471,7 @@ pub struct DnsRecord { impl DnsRecord { fn new(name: &str, ty: RRType, class: u16, ttl: u32) -> Self { - let created = current_time_millis(); + let created = Instant::now(); // From RFC 6762 section 5.2: // "... The querier should plan to issue a query at 80% of the record @@ -492,31 +494,31 @@ impl DnsRecord { self.ttl } - pub const fn get_expire_time(&self) -> u64 { + pub const fn get_expire_time(&self) -> Instant { self.expires } - pub const fn get_refresh_time(&self) -> u64 { + pub const fn get_refresh_time(&self) -> Instant { self.refresh } - pub const fn is_expired(&self, now: u64) -> bool { + pub fn is_expired(&self, now: Instant) -> bool { now >= self.expires } /// Returns whether record expires in 1 second. /// /// This is useful because mDNS sets TTL to 1 (not 0) for expiring records. - pub const fn expires_soon(&self, now: u64) -> bool { - now + 1000 >= self.expires + pub fn expires_soon(&self, now: Instant) -> bool { + now + Duration::from_millis(1000) >= self.expires } - pub const fn refresh_due(&self, now: u64) -> bool { + pub fn refresh_due(&self, now: Instant) -> bool { now >= self.refresh } /// Returns whether `now` (in millis) has passed half of TTL. - pub fn halflife_passed(&self, now: u64) -> bool { + pub fn halflife_passed(&self, now: Instant) -> bool { let halflife = get_expiration_time(self.created, self.ttl, 50); now > halflife } @@ -532,7 +534,7 @@ impl DnsRecord { } /// Returns if this record is due for refresh. If yes, `refresh` time is updated. - pub fn refresh_maybe(&mut self, now: u64) -> bool { + pub fn refresh_maybe(&mut self, now: Instant) -> bool { if self.is_expired(now) || !self.refresh_due(now) { return false; } @@ -563,18 +565,19 @@ impl DnsRecord { } /// Returns the remaining TTL in seconds - fn get_remaining_ttl(&self, now: u64) -> u32 { - let remaining_millis = get_expiration_time(self.created, self.ttl, 100) - now; - cmp::max(0, remaining_millis / 1000) as u32 + fn get_remaining_ttl(&self, now: Instant) -> u32 { + get_expiration_time(self.created, self.ttl, 100) + .saturating_duration_since(now) + .as_secs() as u32 } /// Return the absolute time for this record being created - pub const fn get_created(&self) -> u64 { + pub const fn get_created(&self) -> Instant { self.created } - /// Set the absolute expiration time in millis - fn set_expire(&mut self, expire_at: u64) { + /// Set the expiration time + fn set_expire(&mut self, expire_at: Instant) { self.expires = expire_at; } @@ -592,11 +595,9 @@ impl DnsRecord { } /// Modify TTL to reflect the remaining life time from `now`. - pub fn update_ttl(&mut self, now: u64) { - if now > self.created { - let elapsed = now - self.created; - self.ttl -= (elapsed / 1000) as u32; - } + pub fn update_ttl(&mut self, now: Instant) { + let elapsed = now.saturating_duration_since(self.created); + self.ttl = self.ttl.saturating_sub(elapsed.as_secs() as u32); } pub fn set_new_name(&mut self, new_name: String) { @@ -696,33 +697,33 @@ pub trait DnsRecordExt: fmt::Debug { self.get_record_mut().reset_ttl(other.get_record()); } - fn get_created(&self) -> u64 { + fn get_created(&self) -> Instant { self.get_record().get_created() } - fn get_expire(&self) -> u64 { + fn get_expire(&self) -> Instant { self.get_record().get_expire_time() } - fn set_expire(&mut self, expire_at: u64) { + fn set_expire(&mut self, expire_at: Instant) { self.get_record_mut().set_expire(expire_at); } /// Set expire as `expire_at` if it is sooner than the current `expire`. - fn set_expire_sooner(&mut self, expire_at: u64) { + fn set_expire_sooner(&mut self, expire_at: Instant) { if expire_at < self.get_expire() { self.get_record_mut().set_expire(expire_at); } } /// Returns true if the record expires in 1 second from `now`. - fn expires_soon(&self, now: u64) -> bool { + fn expires_soon(&self, now: Instant) -> bool { self.get_record().expires_soon(now) } /// Given `now`, if the record is due to refresh, this method updates the refresh time /// and returns the new refresh time. Otherwise, returns None. - fn updated_refresh_time(&mut self, now: u64) -> Option { + fn updated_refresh_time(&mut self, now: Instant) -> Option { if self.get_record_mut().refresh_maybe(now) { Some(self.get_record().get_refresh_time()) } else { @@ -1438,7 +1439,10 @@ impl DnsOutPacket { /// Writes a record (answer, authoritative answer, additional). /// /// In error cases nothing is written to the packet. - fn write_record(&mut self, record_ext: &dyn DnsRecordExt, now: u64) -> WriteResult { + /// `now` is `None` for records whose full TTL is written (authorities, + /// additionals); for answers it is the time the answer was added, so the + /// remaining TTL is written instead. + fn write_record(&mut self, record_ext: &dyn DnsRecordExt, now: Option) -> WriteResult { let start_size = self.size(); let record = record_ext.get_record(); @@ -1451,10 +1455,9 @@ impl DnsOutPacket { self.write_short(record.entry.class); } - if now == 0 { - self.write_u32(record.ttl); - } else { - self.write_u32(record.get_remaining_ttl(now)); + match now { + None => self.write_u32(record.ttl), + Some(now) => self.write_u32(record.get_remaining_ttl(now)), } // Placeholder for record size @@ -1797,7 +1800,8 @@ pub struct DnsOutgoing { id: u16, multicast: bool, questions: Vec, - answers: Vec<(DnsRecordBox, u64)>, + /// Answers with the time they were added (`None`: write the full TTL). + answers: Vec<(DnsRecordBox, Option)>, authorities: Vec, additionals: Vec, known_answer_count: i64, // for internal maintenance only @@ -1822,7 +1826,7 @@ impl DnsOutgoing { } /// For testing purposes only. - pub(crate) fn _answers(&self) -> &[(DnsRecordBox, u64)] { + pub(crate) fn _answers(&self) -> &[(DnsRecordBox, Option)] { &self.answers } @@ -1906,7 +1910,7 @@ impl DnsOutgoing { /// A workaround as Rust doesn't allow us to pass DnsRecordBox in as `impl DnsRecordExt` pub fn add_answer_box(&mut self, answer_box: DnsRecordBox) { - self.answers.push((answer_box, 0)); + self.answers.push((answer_box, None)); } pub fn add_authority(&mut self, record: DnsRecordBox) { @@ -1943,18 +1947,20 @@ impl DnsOutgoing { return false; } - self.add_answer_at_time(answer, 0) + self.add_answer_at_time(answer, None) } /// Returns true if `answer` is added to the outgoing msg. /// Returns false if the answer is expired `now` hence not added. - /// If `now` is 0, do not check if the answer expires. + /// If `now` is `None`, do not check if the answer expires, and write its + /// full TTL on the wire. pub fn add_answer_at_time( &mut self, answer: impl DnsRecordExt + Send + 'static, - now: u64, + now: Option, ) -> bool { - if now == 0 || !answer.get_record().is_expired(now) { + let expired = now.is_some_and(|now| answer.get_record().is_expired(now)); + if !expired { trace!("add_answer push: {:?}", &answer); self.answers.push((answer.boxed(), now)); return true; @@ -2128,13 +2134,13 @@ impl DnsOutgoing { for auth in self.authorities.iter() { builder.add(Section::Authority, |packet| { - packet.write_record(auth.as_ref(), 0) + packet.write_record(auth.as_ref(), None) }); } for addi in self.additionals.iter() { builder.add(Section::Additional, |packet| { - packet.write_record(addi.as_ref(), 0) + packet.write_record(addi.as_ref(), None) }); } @@ -2852,12 +2858,12 @@ const fn u32_from_be_slice(s: &[u8]) -> u32 { u32::from_be_bytes(u8_array) } -/// Returns the UNIX time in millis at which this record will have expired -/// by a certain percentage. -const fn get_expiration_time(created: u64, ttl: u32, percent: u32) -> u64 { - // 'created' is in millis, 'ttl' is in seconds, hence: +/// Returns the time at which this record will have expired by a certain +/// percentage of its TTL. +fn get_expiration_time(created: Instant, ttl: u32, percent: u32) -> Instant { + // 'ttl' is in seconds, hence: // ttl * 1000 * (percent / 100) => ttl * percent * 10 - created + (ttl as u64 * percent as u64 * 10) + created + Duration::from_millis(ttl as u64 * percent as u64 * 10) } #[cfg(test)] @@ -3016,7 +3022,7 @@ mod tests { 0xaaaa5555, "test-service".to_string(), ), - 0, + None, ); let packets = out.to_packets(MAX_PKT_DEFAULT, IPV6); assert_eq!(packets.len(), 1); @@ -3039,7 +3045,7 @@ mod tests { 0xaaaa5555, "test-service.local".to_string(), ), - 0, + None, ); out.add_answer_at_time( DnsPointer::new( @@ -3049,7 +3055,7 @@ mod tests { 0xffffffff, "test-service.local".to_string(), ), - 0, + None, ); let packets = out.to_packets(MAX_PKT_DEFAULT, IPV6); assert_eq!(packets.len(), 1); @@ -3112,7 +3118,7 @@ mod tests { 0, format!("{long_label}._test._tcp.local."), ), - 0, + None, ); out.add_answer_at_time( DnsPointer::new( @@ -3122,7 +3128,7 @@ mod tests { 0, "ok._test._tcp.local.".to_string(), ), - 0, + None, ); let packets = out.to_packets(MAX_PKT_DEFAULT, IPV6); @@ -3480,7 +3486,7 @@ mod tests { let mut out = DnsOutgoing::new(FLAGS_QR_RESPONSE); for i in 0..ANSWER_COUNT { - out.add_answer_at_time(ptr_answer(i), 0); + out.add_answer_at_time(ptr_answer(i), None); } let packets = out.to_packets(MAX_PKT_DEFAULT, IPV6); @@ -3549,12 +3555,12 @@ mod tests { #[test] fn test_dns_outgoing_oversized_record_sent_alone() { let mut out = DnsOutgoing::new(FLAGS_QR_RESPONSE); - out.add_answer_at_time(ptr_answer(0), 0); + out.add_answer_at_time(ptr_answer(0), None); out.add_answer_at_time( DnsTxt::new("big._spill._tcp.local.", CLASS_IN, 4500, vec![b'x'; 2000]), - 0, + None, ); - out.add_answer_at_time(ptr_answer(1), 0); + out.add_answer_at_time(ptr_answer(1), None); let packets = out.to_packets(MAX_PKT_DEFAULT, IPV6); assert_eq!(packets.len(), 3, "the big record needs a packet to itself"); @@ -3580,7 +3586,7 @@ mod tests { #[test] fn test_dns_outgoing_record_over_absolute_ceiling_dropped() { let mut out = DnsOutgoing::new(FLAGS_QR_RESPONSE); - out.add_answer_at_time(ptr_answer(0), 0); + out.add_answer_at_time(ptr_answer(0), None); out.add_answer_at_time( DnsTxt::new( "huge._spill._tcp.local.", @@ -3588,9 +3594,9 @@ mod tests { 4500, vec![b'x'; MAX_PKT_ABSOLUTE_IPV6], ), - 0, + None, ); - out.add_answer_at_time(ptr_answer(1), 0); + out.add_answer_at_time(ptr_answer(1), None); let packets = out.to_packets(MAX_PKT_DEFAULT, IPV6); for packet in &packets { @@ -3611,7 +3617,7 @@ mod tests { fn test_dns_outgoing_all_sections_spill() { let mut out = DnsOutgoing::new(FLAGS_QR_RESPONSE); for i in 0..40 { - out.add_answer_at_time(ptr_answer(i), 0); + out.add_answer_at_time(ptr_answer(i), None); } for i in 40..80 { out.add_authority(Box::new(ptr_answer(i))); @@ -3656,7 +3662,7 @@ mod tests { "negative.local.".to_string(), bitmap.clone(), ), - 0, + None, ); let packets = out.to_packets(MAX_PKT_DEFAULT, IPV6); assert_eq!(packets.len(), 1); diff --git a/src/lib.rs b/src/lib.rs index fb98feb0..bd5243b4 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -197,13 +197,3 @@ pub use flume::Receiver; /// Errors returned by the receiving methods of `Receiver`. Re-export from `flume` crate. pub use flume::{RecvError, RecvTimeoutError, TryRecvError}; - -use std::time::SystemTime; - -/// Returns the current time in milliseconds since the UNIX epoch. -pub(crate) fn current_time_millis() -> u64 { - SystemTime::now() - .duration_since(SystemTime::UNIX_EPOCH) - .expect("failed to get current UNIX time") - .as_millis() as u64 -} diff --git a/src/service_daemon.rs b/src/service_daemon.rs index 5e50bcf7..cf9c91b2 100644 --- a/src/service_daemon.rs +++ b/src/service_daemon.rs @@ -31,7 +31,6 @@ #[cfg(feature = "logging")] use crate::log::{debug, error, trace}; use crate::{ - current_time_millis, dns_cache::{DnsCache, IpType}, dns_parser::{ ip_address_rr_type, max_pkt_absolute, DnsAddress, DnsEntryExt, DnsIncoming, DnsNSec, @@ -57,7 +56,7 @@ use std::{ fmt, io, net::{IpAddr, Ipv4Addr, Ipv6Addr, SocketAddr, SocketAddrV4, SocketAddrV6, UdpSocket}, str, thread, - time::Duration, + time::{Duration, Instant}, vec, }; @@ -951,17 +950,16 @@ fn new_socket(addr: SocketAddr, non_block: bool) -> Result { Ok(fd) } -/// Specify a UNIX timestamp in millis to run `command` for the next time. +/// Specify when to run `command` for the next time. struct ReRun { - /// UNIX timestamp in millis. - next_time: u64, + next_time: Instant, command: Command, } /// A query response deferred per RFC 6762 §6 (shared response). struct DelayedResponse { - /// UNIX timestamp in millis at which to send `out`. - next_time: u64, + /// When to send `out`. + next_time: Instant, out: DnsOutgoing, if_index: u32, is_ipv4: bool, @@ -1152,7 +1150,7 @@ struct Zeroconf { /// /// The timestamps are set at the future timestamp when the command should timeout. /// `hostname` is case-insensitive and stored in lowercase. - hostname_resolvers: HashMap, Option)>, // + hostname_resolvers: HashMap, Option)>, // /// All repeating transmissions. retransmissions: Vec, @@ -1189,7 +1187,7 @@ struct Zeroconf { /// /// When the run loop goes through a single iteration, it will /// set its timeout to the earliest timer in this list. - timers: BinaryHeap>, + timers: BinaryHeap>, status: DaemonStatus, @@ -1535,27 +1533,29 @@ impl Zeroconf { } // Setup timer for IP checks. + // `None` when IP checks are disabled. let mut next_ip_check = if self.ip_check_interval > 0 { - current_time_millis() + self.ip_check_interval + Some(Instant::now() + Duration::from_millis(self.ip_check_interval)) } else { - 0 + None }; - if next_ip_check > 0 { - self.add_timer(next_ip_check); + if let Some(t) = next_ip_check { + self.add_timer(t); } // Start the run loop. let mut events = mio::Events::with_capacity(1024); loop { - let now = current_time_millis(); + let now = Instant::now(); let earliest_timer = self.peek_earliest_timer(); let timeout = earliest_timer.map(|timer| { // If `timer` already passed, set `timeout` to be 1ms. - let millis = if timer > now { timer - now } else { 1 }; - Duration::from_millis(millis) + timer + .saturating_duration_since(now) + .max(Duration::from_millis(1)) }); // Process incoming packets, command events and optional timeout. @@ -1565,7 +1565,7 @@ impl Zeroconf { Err(e) => debug!("failed to select from sockets: {}", e), } - let now = current_time_millis(); + let now = Instant::now(); // Remove the timers if already passed. self.pop_timers_till(now); @@ -1642,7 +1642,7 @@ impl Zeroconf { self.increase_counter(Counter::CacheRefreshAddr, query_count); // check and evict expired records in our cache - let now = current_time_millis(); + let now = Instant::now(); // Notify service listeners about the expired records. let expired_services = self.cache.evict_expired_services(now); @@ -1671,9 +1671,10 @@ impl Zeroconf { self.probing_handler(); // check IP changes if next_ip_check is reached. - if now >= next_ip_check && next_ip_check > 0 { - next_ip_check = now + self.ip_check_interval; - self.add_timer(next_ip_check); + if next_ip_check.is_some_and(|t| now >= t) { + let t = now + Duration::from_millis(self.ip_check_interval); + next_ip_check = Some(t); + self.add_timer(t); self.check_ip_changes(); } @@ -1817,20 +1818,20 @@ impl Zeroconf { } } - fn add_timer(&mut self, next_time: u64) { + fn add_timer(&mut self, next_time: Instant) { self.timers.push(Reverse(next_time)); } - fn peek_earliest_timer(&self) -> Option { + fn peek_earliest_timer(&self) -> Option { self.timers.peek().map(|Reverse(v)| *v) } - fn _pop_earliest_timer(&mut self) -> Option { + fn _pop_earliest_timer(&mut self) -> Option { self.timers.pop().map(|Reverse(v)| v) } /// Pop all timers that are already passed till `now`. - fn pop_timers_till(&mut self, now: u64) { + fn pop_timers_till(&mut self, now: Instant) { while let Some(Reverse(v)) = self.timers.peek() { if *v > now { break; @@ -2347,9 +2348,10 @@ impl Zeroconf { // RFC 6762 section 8.3. // ..The Multicast DNS responder MUST send at least two unsolicited // responses, one second apart. - let next_time = current_time_millis() - + ANNOUNCE_SECOND_DELAY_MILLIS - + fastrand::u64(0..ANNOUNCE_SECOND_JITTER_MILLIS); + let next_time = Instant::now() + + Duration::from_millis( + ANNOUNCE_SECOND_DELAY_MILLIS + fastrand::u64(0..ANNOUNCE_SECOND_JITTER_MILLIS), + ); for if_index in outgoing_intfs { self.add_retransmission( next_time, @@ -2362,7 +2364,7 @@ impl Zeroconf { /// Send probings or finish them if expired. Notify waiting services. fn probing_handler(&mut self) { - let now = current_time_millis(); + let now = Instant::now(); let mut invalid_intf_addrs = HashSet::new(); for (if_index, intf) in self.my_intfs.iter() { @@ -2440,8 +2442,10 @@ impl Zeroconf { if announced_v4 || announced_v6 { let next_time = now - + ANNOUNCE_SECOND_DELAY_MILLIS - + fastrand::u64(0..ANNOUNCE_SECOND_JITTER_MILLIS); + + Duration::from_millis( + ANNOUNCE_SECOND_DELAY_MILLIS + + fastrand::u64(0..ANNOUNCE_SECOND_JITTER_MILLIS), + ); let command = Command::RegisterResend(info.get_fullname().to_string(), *if_index); self.retransmissions.push(ReRun { next_time, command }); @@ -2496,14 +2500,14 @@ impl Zeroconf { 0, fullname.to_string(), ), - 0, + None, ); if let Some(sub) = info.get_subtype() { trace!("Adding subdomain {}", sub); out.add_answer_at_time( DnsPointer::new(sub, RRType::PTR, CLASS_IN, 0, fullname.to_string()), - 0, + None, ); } @@ -2517,7 +2521,7 @@ impl Zeroconf { info.get_port(), hostname.to_string(), ), - 0, + None, ); out.add_answer_at_time( DnsTxt::new( @@ -2526,7 +2530,7 @@ impl Zeroconf { 0, info.generate_txt(), ), - 0, + None, ); let if_addrs = if is_ipv4 { @@ -2549,7 +2553,7 @@ impl Zeroconf { address, intf.into(), ), - 0, + None, ); } @@ -2574,7 +2578,7 @@ impl Zeroconf { listener: Sender, timeout: Option, ) { - let real_timeout = timeout.map(|t| current_time_millis() + t); + let real_timeout = timeout.map(|t| Instant::now() + Duration::from_millis(t)); self.hostname_resolvers .insert(hostname.to_lowercase(), (listener, real_timeout)); if let Some(t) = real_timeout { @@ -2618,7 +2622,7 @@ impl Zeroconf { /// Sends out a list of `questions` (i.e. DNS questions) via multicast. fn send_query_vec(&self, questions: &[(&str, RRType)]) { let mut out = DnsOutgoing::new(FLAGS_QR_QUERY); - let now = current_time_millis(); + let now = Instant::now(); for (name, qtype) in questions { out.add_question(name, *qtype); @@ -2812,7 +2816,7 @@ impl Zeroconf { /// pending instance keeps being queried for as long as the browse is /// active. fn query_unresolved_instances(&mut self, ty_domain: &str) { - let now = current_time_millis(); + let now = Instant::now(); let mut instances = Vec::new(); if let Some(records) = self.cache.get_ptr(ty_domain) { for record in records.iter().filter(|r| !r.record.expires_soon(now)) { @@ -2835,7 +2839,7 @@ impl Zeroconf { &mut self, ty_domain: &str, sender: &Sender, - now: u64, + now: Instant, ) { let mut resolved: HashSet = HashSet::new(); let mut unresolved: HashSet = HashSet::new(); @@ -2917,7 +2921,7 @@ impl Zeroconf { fn add_pending_resolve(&mut self, instance: String) { if !self.pending_resolves.contains(&instance) { - let next_time = current_time_millis() + RESOLVE_RETRY_BASE_MILLIS; + let next_time = Instant::now() + Duration::from_millis(RESOLVE_RETRY_BASE_MILLIS); self.add_retransmission(next_time, Command::Resolve(instance.clone(), 1)); self.pending_resolves.insert(instance); } @@ -2929,7 +2933,7 @@ impl Zeroconf { ty_domain: &str, fullname: &str, ) -> Result { - let now = current_time_millis(); + let now = Instant::now(); let mut resolved_service = ResolvedService { ty_domain: ty_domain.to_string(), sub_ty_domain: None, @@ -3063,7 +3067,7 @@ impl Zeroconf { /// Deal with incoming response packets. All answers /// are held in the cache, and listeners are notified. fn handle_response(&mut self, mut msg: DnsIncoming, if_index: u32) { - let now = current_time_millis(); + let now = Instant::now(); // remove records that are expired. let mut record_predicate = |record: &DnsRecordBox| { @@ -3309,7 +3313,7 @@ impl Zeroconf { // } // Probing again with the new names. - let create_time = current_time_millis() + fastrand::u64(0..250); + let create_time = Instant::now() + Duration::from_millis(fastrand::u64(0..250)); let waiting_services = probe.waiting_services.clone(); @@ -3366,7 +3370,7 @@ impl Zeroconf { let mut unresolved: HashSet = HashSet::new(); let mut removed_instances = HashMap::new(); - let now = current_time_millis(); + let now = Instant::now(); for (ty_domain, records) in self.cache.all_ptr().iter() { if !self.service_queriers.contains_key(ty_domain) { @@ -3552,7 +3556,7 @@ impl Zeroconf { self.increase_counter(Counter::KnownAnswerSuppression, out.known_answer_count()); let delay = fastrand::u64(SHARED_RESPONSE_DELAY_MIN_MILLIS..SHARED_RESPONSE_DELAY_MAX_MILLIS); - let next_time = current_time_millis() + delay; + let next_time = Instant::now() + Duration::from_millis(delay); self.delayed_responses.push(DelayedResponse { next_time, out, @@ -3630,7 +3634,7 @@ impl Zeroconf { // records in its Authority Section, and we MUST defend our // records immediately so the prober detects the conflict. if let Some(dns_registry) = self.dns_registry_map.get_mut(&if_index) { - dns_registry.apply_multicast_rate_limit(out, current_time_millis(), is_ipv4); + dns_registry.apply_multicast_rate_limit(out, Instant::now(), is_ipv4); } } @@ -3682,7 +3686,7 @@ impl Zeroconf { }; if let Some(dns_registry) = self.dns_registry_map.get_mut(&if_index) { - dns_registry.apply_multicast_rate_limit(&mut out, current_time_millis(), is_ipv4); + dns_registry.apply_multicast_rate_limit(&mut out, Instant::now(), is_ipv4); } if out.answers_count() == 0 { return; @@ -3735,7 +3739,7 @@ impl Zeroconf { } } - fn add_retransmission(&mut self, next_time: u64, command: Command) { + fn add_retransmission(&mut self, next_time: Instant, command: Command) { self.retransmissions.push(ReRun { next_time, command }); self.add_timer(next_time); } @@ -3922,7 +3926,7 @@ impl Zeroconf { return; } - let now = current_time_millis(); + let now = Instant::now(); if !repeating { // Binds a `listener` to querying mDNS domain type `ty`. // @@ -3946,7 +3950,10 @@ impl Zeroconf { // RFC 6762 §5.2: delay the first query by a random jitter. let jitter = fastrand::u64(INITIAL_QUERY_DELAY_MIN_MILLIS..INITIAL_QUERY_DELAY_MAX_MILLIS); - self.add_retransmission(now + jitter, Command::Browse(ty, 1, cache_only, listener)); + self.add_retransmission( + now + Duration::from_millis(jitter), + Command::Browse(ty, 1, cache_only, listener), + ); return; } @@ -3956,7 +3963,7 @@ impl Zeroconf { self.increase_counter(Counter::Browse, 1); - let next_time = now + (next_delay * 1000) as u64; + let next_time = now + Duration::from_millis((next_delay * 1000) as u64); let max_delay = 60 * 60; let delay = cmp::min(next_delay * 2, max_delay); self.add_retransmission(next_time, Command::Browse(ty, delay, cache_only, listener)); @@ -3981,7 +3988,7 @@ impl Zeroconf { ); return; } - let now = current_time_millis(); + let now = Instant::now(); if !repeating { self.add_hostname_resolver(hostname.to_owned(), listener.clone(), timeout); // if we already have the records in our cache, just send them @@ -3991,7 +3998,7 @@ impl Zeroconf { let jitter = fastrand::u64(INITIAL_QUERY_DELAY_MIN_MILLIS..INITIAL_QUERY_DELAY_MAX_MILLIS); self.add_retransmission( - now + jitter, + now + Duration::from_millis(jitter), Command::ResolveHostname(hostname, 1, listener, None), ); return; @@ -4000,7 +4007,7 @@ impl Zeroconf { self.send_query_vec(&[(&hostname, RRType::A), (&hostname, RRType::AAAA)]); self.increase_counter(Counter::ResolveHostname, 1); - let next_time = now + u64::from(next_delay) * 1000; + let next_time = now + Duration::from_millis(u64::from(next_delay) * 1000); let max_delay = 60 * 60; let delay = cmp::min(next_delay * 2, max_delay); @@ -4027,7 +4034,7 @@ impl Zeroconf { // // Back off exponentially let next_delay = RESOLVE_RETRY_BASE_MILLIS << try_count; - let next_time = current_time_millis() + next_delay; + let next_time = Instant::now() + Duration::from_millis(next_delay); self.add_retransmission(next_time, Command::Resolve(instance, try_count + 1)); } else { // This fast-path retry chain is ending. @@ -4054,7 +4061,7 @@ impl Zeroconf { let packet = self.unregister_service(&info, intf, &sock.pktinfo); // repeat for one time just in case some peers miss the message if !repeating && !packet.is_empty() { - let next_time = current_time_millis() + 120; + let next_time = Instant::now() + Duration::from_millis(120); self.retransmissions.push(ReRun { next_time, command: Command::UnregisterResend(packet, *if_index, true), @@ -4067,7 +4074,7 @@ impl Zeroconf { if let Some(sock) = self.ipv6_sock.as_ref() { let packet = self.unregister_service(&info, intf, &sock.pktinfo); if !repeating && !packet.is_empty() { - let next_time = current_time_millis() + 120; + let next_time = Instant::now() + Duration::from_millis(120); self.retransmissions.push(ReRun { next_time, command: Command::UnregisterResend(packet, *if_index, false), @@ -4236,11 +4243,11 @@ impl Zeroconf { though its TTL may indicate that it is not yet due to expire, that record SHOULD be promptly flushed from the cache. */ - let now = current_time_millis(); + let now = Instant::now(); let expire_at = if repeating { None } else { - Some(now + timeout.as_millis() as u64) + Some(now + Duration::from_millis(timeout.as_millis() as u64)) }; // send query for the resource records. @@ -4257,7 +4264,10 @@ impl Zeroconf { self.add_timer(new_expire); // ensure a check for the new expire time. // schedule a resend 1 second later - self.add_retransmission(now + 1000, Command::Verify(instance, timeout)); + self.add_retransmission( + now + Duration::from_millis(1000), + Command::Verify(instance, timeout), + ); } } } @@ -4741,7 +4751,7 @@ fn call_service_listener( } fn call_hostname_resolution_listener( - listeners_map: &HashMap, Option)>, + listeners_map: &HashMap, Option)>, hostname: &str, event: HostnameResolutionEvent, ) { @@ -5043,7 +5053,7 @@ fn prepare_announce( let mut probing_count = 0; let mut out = DnsOutgoing::new(FLAGS_QR_RESPONSE | FLAGS_AA); - let create_time = current_time_millis() + fastrand::u64(0..250); + let create_time = Instant::now() + Duration::from_millis(fastrand::u64(0..250)); out.add_answer_at_time( DnsPointer::new( @@ -5053,7 +5063,7 @@ fn prepare_announce( info.get_other_ttl(), service_fullname.to_string(), ), - 0, + None, ); if let Some(sub) = info.get_subtype() { @@ -5066,7 +5076,7 @@ fn prepare_announce( info.get_other_ttl(), service_fullname.to_string(), ), - 0, + None, ); } @@ -5090,7 +5100,7 @@ fn prepare_announce( if !info.requires_probe() || dns_registry.is_probing_done(&srv, info.get_fullname(), create_time) { - out.add_answer_at_time(srv, 0); + out.add_answer_at_time(srv, None); } else { probing_count += 1; } @@ -5111,7 +5121,7 @@ fn prepare_announce( if !info.requires_probe() || dns_registry.is_probing_done(&txt, info.get_fullname(), create_time) { - out.add_answer_at_time(txt, 0); + out.add_answer_at_time(txt, None); } else { probing_count += 1; } @@ -5136,7 +5146,7 @@ fn prepare_announce( if !info.requires_probe() || dns_registry.is_probing_done(&dns_addr, info.get_fullname(), create_time) { - out.add_answer_at_time(dns_addr, 0); + out.add_answer_at_time(dns_addr, None); } else { probing_count += 1; } @@ -5162,7 +5172,7 @@ fn announce_service_on_intf( if let Some(mut out) = prepare_announce(info, intf, dns_registry, is_ipv4) { // RFC 6762 §6: a record MUST NOT be multicast on an interface more than // once per second. Announcements are unsolicited multicast responses. - dns_registry.apply_multicast_rate_limit(&mut out, current_time_millis(), is_ipv4); + dns_registry.apply_multicast_rate_limit(&mut out, Instant::now(), is_ipv4); if out.answers_count() > 0 { let _ = send_dns_outgoing(&out, intf, sock, port, None, None)?; } @@ -5240,8 +5250,8 @@ fn hostname_change(original: &str) -> String { /// that are finished. fn check_probing( dns_registry: &mut DnsRegistry, - timers: &mut BinaryHeap>, - now: u64, + timers: &mut BinaryHeap>, + now: Instant, ) -> (DnsOutgoing, Vec) { let mut expired_probes = Vec::new(); let mut out = DnsOutgoing::new(FLAGS_QR_QUERY); @@ -5498,8 +5508,8 @@ mod tests { .probing .values_mut() { - probe.start_time = crate::current_time_millis() - 1000; - probe.next_send = 0; + probe.start_time = Instant::now() - Duration::from_millis(1000); + probe.next_send = probe.start_time; } daemon.probing_handler(); assert_eq!( @@ -6912,11 +6922,11 @@ mod tests { let mut out = DnsOutgoing::new(FLAGS_QR_RESPONSE | FLAGS_AA); out.add_answer_at_time( DnsPointer::new(&ty_domain, RRType::PTR, CLASS_IN, ttl, instance.clone()), - 0, + None, ); out.add_answer_at_time( DnsSrv::new(&instance, CLASS_IN, ttl, 0, 0, port, host.clone()), - 0, + None, ); out.to_data_on_wire(MAX_PKT_DEFAULT, true) }; @@ -6933,7 +6943,7 @@ mod tests { IpAddr::V4(intf_ip), if_id.clone(), ), - 0, + None, ); out.to_data_on_wire(MAX_PKT_DEFAULT, true) }; @@ -7090,8 +7100,8 @@ mod tests { .probing .values_mut() { - probe.start_time = crate::current_time_millis() - 1000; - probe.next_send = 0; + probe.start_time = Instant::now() - Duration::from_millis(1000); + probe.next_send = probe.start_time; } daemon.probing_handler(); assert_eq!( diff --git a/src/service_info.rs b/src/service_info.rs index 4a65c719..4e5d8cc4 100644 --- a/src/service_info.rs +++ b/src/service_info.rs @@ -14,6 +14,7 @@ use std::{ fmt, net::{IpAddr, Ipv4Addr}, str::FromStr, + time::{Duration, Instant}, }; #[cfg(feature = "serde")] @@ -1021,14 +1022,14 @@ pub(crate) struct Probe { pub(crate) waiting_services: HashSet, /// The time (T) to send the first query . - pub(crate) start_time: u64, + pub(crate) start_time: Instant, /// The time to send the next (including the first) query. - pub(crate) next_send: u64, + pub(crate) next_send: Instant, } impl Probe { - pub(crate) fn new(start_time: u64) -> Self { + pub(crate) fn new(start_time: Instant) -> Self { // RFC 6762: https://datatracker.ietf.org/doc/html/rfc6762#section-8.1: // // "250 ms after the first query, the host should send a second; then, @@ -1071,7 +1072,7 @@ impl Probe { /// Compares with `incoming` records. Postpone probe and retry if we yield. pub(crate) fn tiebreaking(&mut self, msg: &DnsIncoming, probe_name: &str) { - let now = crate::current_time_millis(); + let now = Instant::now(); // Only do tiebreaking if probe already started. // This check also helps avoid redo tiebreaking if start time @@ -1116,8 +1117,8 @@ impl Probe { match cmp_result { cmp::Ordering::Less => { debug!("tiebreaking '{probe_name}': LOST, will wait for one second",); - self.start_time = now + 1000; // wait and restart. - self.next_send = now + 1000; + self.start_time = now + Duration::from_millis(1000); // wait and restart. + self.next_send = now + Duration::from_millis(1000); } ordering => { debug!("tiebreaking '{probe_name}': {:?}", ordering); @@ -1125,15 +1126,15 @@ impl Probe { } } - pub(crate) fn update_next_send(&mut self, now: u64) { - self.next_send = now + 250; + pub(crate) fn update_next_send(&mut self, now: Instant) { + self.next_send = now + Duration::from_millis(250); } /// Returns whether this probe is finished. - pub(crate) fn expired(&self, now: u64) -> bool { + pub(crate) fn expired(&self, now: Instant) -> bool { // The 2nd query is T + 250ms, the 3rd query is T + 500ms, // The expire time is T + 750ms - now >= self.start_time + 750 + now >= self.start_time + Duration::from_millis(750) } } @@ -1156,7 +1157,7 @@ pub(crate) struct DnsRegistry { pub(crate) active: HashMap>, /// timers of the newly added probes. - pub(crate) new_timers: Vec, + pub(crate) new_timers: Vec, /// Mapping from original names to new names. pub(crate) name_changes: HashMap, @@ -1170,10 +1171,10 @@ pub(crate) struct DnsRegistry { /// carries both address families, but they are distinct multicast groups /// (`224.0.0.251` and `ff02::fb`) reaching potentially different listeners, /// so sending a record on one group must not throttle it on the other. - pub(crate) last_multicast_v4: HashMap, + pub(crate) last_multicast_v4: HashMap, /// Same as [`Self::last_multicast_v4`] but for this interface's IPv6 group. - pub(crate) last_multicast_v6: HashMap, + pub(crate) last_multicast_v6: HashMap, } impl DnsRegistry { @@ -1205,7 +1206,7 @@ impl DnsRegistry { pub(crate) fn apply_multicast_rate_limit( &mut self, out: &mut DnsOutgoing, - now: u64, + now: Instant, is_ipv4: bool, ) { let last_multicast = if is_ipv4 { @@ -1216,7 +1217,10 @@ impl DnsRegistry { // Prune stale entries so the map stays bounded across name changes; // any record older than the one-second window is irrelevant now. - last_multicast.retain(|_, last| now.saturating_sub(*last) < MULTICAST_RATE_LIMIT_MILLIS); + last_multicast.retain(|_, last| { + now.saturating_duration_since(*last) + < Duration::from_millis(MULTICAST_RATE_LIMIT_MILLIS) + }); out.retain_answers(|record| keep_after_rate_limit(last_multicast, record, now)); @@ -1238,7 +1242,7 @@ impl DnsRegistry { &mut self, answer: &T, service_name: &str, - start_time: u64, + start_time: Instant, ) -> bool where T: DnsRecordExt + Send + 'static, @@ -1296,7 +1300,7 @@ impl DnsRegistry { &mut self, original: &str, new_name: &str, - probe_time: u64, + probe_time: Instant, ) -> bool { let mut found_records = Vec::new(); let mut new_timer_added = false; @@ -1369,13 +1373,18 @@ pub(crate) const MULTICAST_RATE_LIMIT_MILLIS: u64 = 1000; /// Returns whether `record` may still be multicast under the RFC 6762 section 6 /// rate limit, updating `last_multicast` to `now` when it is kept. fn keep_after_rate_limit( - last_multicast: &mut HashMap, + last_multicast: &mut HashMap, record: &DnsRecordBox, - now: u64, + now: Instant, ) -> bool { let key = rate_limit_key(record); match last_multicast.get(&key) { - Some(last) if now.saturating_sub(*last) < MULTICAST_RATE_LIMIT_MILLIS => false, + Some(last) + if now.saturating_duration_since(*last) + < Duration::from_millis(MULTICAST_RATE_LIMIT_MILLIS) => + { + false + } _ => { last_multicast.insert(key, now); true @@ -1512,11 +1521,14 @@ impl ResolvedService { #[cfg(test)] mod tests { - use super::{decode_txt, encode_txt, u8_slice_to_hex, DnsRegistry, ServiceInfo, TxtProperty}; + use super::{ + decode_txt, encode_txt, u8_slice_to_hex, DnsRegistry, Instant, ServiceInfo, TxtProperty, + }; use crate::dns_parser::{DnsOutgoing, DnsPointer, RRType, CLASS_IN, FLAGS_QR_RESPONSE}; use crate::{IfKind, IfPredicate}; use if_addrs::{IfAddr, IfOperStatus, Ifv4Addr, Ifv6Addr, Interface}; use std::net::{Ipv4Addr, Ipv6Addr}; + use std::time::Duration; /// RFC 6762 section 6: the same record must not be multicast on an /// interface more than once per second, but is allowed again after a @@ -1535,12 +1547,12 @@ mod tests { 4500, "inst._test._tcp.local.".to_string(), ), - 0, + None, ); out }; - let now = 1_000_000; + let now = Instant::now(); // First multicast at `now`: the record passes through. let mut out = build_out(); @@ -1549,12 +1561,12 @@ mod tests { // Again 500ms later: the record is throttled (dropped). let mut out = build_out(); - registry.apply_multicast_rate_limit(&mut out, now + 500, true); + registry.apply_multicast_rate_limit(&mut out, now + Duration::from_millis(500), true); assert_eq!(out.answers_count(), 0); // Exactly 1 second after the first send: allowed again. let mut out = build_out(); - registry.apply_multicast_rate_limit(&mut out, now + 1000, true); + registry.apply_multicast_rate_limit(&mut out, now + Duration::from_millis(1000), true); assert_eq!(out.answers_count(), 1); } @@ -1577,12 +1589,12 @@ mod tests { 4500, "inst._test._tcp.local.".to_string(), ), - 0, + None, ); out }; - let now = 1_000_000; + let now = Instant::now(); // Multicast the record on IPv4: passes through. let mut out = build_out(); @@ -1598,12 +1610,12 @@ mod tests { // A second IPv4 send within the window is still throttled, confirming // the IPv6 send did not reset (or get charged to) the IPv4 bucket. let mut out = build_out(); - registry.apply_multicast_rate_limit(&mut out, now + 500, true); + registry.apply_multicast_rate_limit(&mut out, now + Duration::from_millis(500), true); assert_eq!(out.answers_count(), 0); // Likewise a second IPv6 send within the window is throttled. let mut out = build_out(); - registry.apply_multicast_rate_limit(&mut out, now + 500, false); + registry.apply_multicast_rate_limit(&mut out, now + Duration::from_millis(500), false); assert_eq!(out.answers_count(), 0); } @@ -1634,28 +1646,28 @@ mod tests { ) }; - let now = 1_000_000; + let now = Instant::now(); // Send the PTR answer once so it is throttled going forward. let mut out = DnsOutgoing::new(FLAGS_QR_RESPONSE); - out.add_answer_at_time(ptr_answer(), 0); + out.add_answer_at_time(ptr_answer(), None); registry.apply_multicast_rate_limit(&mut out, now, true); assert_eq!(out.answers_count(), 1); // 100ms later: PTR answer is throttled, and `extra` rides along as an // additional. With no answer surviving, nothing is sent. let mut out = DnsOutgoing::new(FLAGS_QR_RESPONSE); - out.add_answer_at_time(ptr_answer(), 0); + out.add_answer_at_time(ptr_answer(), None); out.add_additional_answer(extra()); - registry.apply_multicast_rate_limit(&mut out, now + 100, true); + registry.apply_multicast_rate_limit(&mut out, now + Duration::from_millis(100), true); assert_eq!(out.answers_count(), 0); // 200ms later: `extra` is now requested as a real answer. It must pass, // because it was never actually multicast above (only carried as an // unsent additional), so the 1-second limit does not apply to it. let mut out = DnsOutgoing::new(FLAGS_QR_RESPONSE); - out.add_answer_at_time(extra(), 0); - registry.apply_multicast_rate_limit(&mut out, now + 200, true); + out.add_answer_at_time(extra(), None); + registry.apply_multicast_rate_limit(&mut out, now + Duration::from_millis(200), true); assert_eq!(out.answers_count(), 1); }