1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
use crate::arp::NetworkManager;
use crate::error::Result;
use crate::oui::lookup_vendor;
use crate::ping::PingScanner;
use futures::stream::{self, StreamExt};
use ipnet::Ipv4Net;
use std::collections::HashMap;
use std::net::IpAddr;
use std::sync::{
atomic::{AtomicUsize, Ordering},
Arc,
};
#[cfg(feature = "mdns")]
mod mdns_fusion;
mod types;
pub use types::{DiscoverPhase, DiscoverProgress, FoundBy, Host, NameSource};
pub struct DiscoverEngine {
ping_scanner: PingScanner,
concurrency: usize,
ping_timeout_ms: u64,
dns_timeout_ms: u64,
}
impl DiscoverEngine {
pub fn new(concurrency: usize) -> Self {
Self::new_with_timeouts(
concurrency,
crate::DEFAULT_PING_TIMEOUT_MS,
crate::DEFAULT_DNS_TIMEOUT_MS,
)
}
pub fn new_with_timeouts(
concurrency: usize,
ping_timeout_ms: u64,
dns_timeout_ms: u64,
) -> Self {
// Clamp both ends. `.max(1)` alone left the upper bound to
// `Semaphore::new`, which asserts `permits <= usize::MAX >> 3` --
// so a caller passing `usize::MAX` got a panic rather than an
// error, from a constructor that returns `Self` and cannot report
// one. Absurd input, but this is public API and a process abort is
// the wrong failure mode for it.
let concurrency = concurrency.clamp(1, crate::MAX_CONCURRENCY);
Self {
ping_scanner: PingScanner::new(concurrency),
concurrency,
ping_timeout_ms: ping_timeout_ms.max(1),
dns_timeout_ms: dns_timeout_ms.max(1),
}
}
pub async fn scan_subnet(&self, subnet: Ipv4Net, resolve: bool) -> Result<Vec<Host>> {
self.scan_subnet_with_progress(subnet, resolve, None).await
}
/// Ping-sweep a subnet, then resolve names for the hosts that answered.
///
/// Enforces the /16 cap itself rather than trusting the caller. This is
/// public API re-exported at the crate root, and it used to collect every
/// address of whatever `Ipv4Net` it was given straight into a `Vec` --
/// `0.0.0.0/0` is 4,294,967,294 entries, roughly 73 GB, allocated before
/// a single packet is sent. `Ops` checked, the engine did not.
pub async fn scan_subnet_with_progress(
&self,
subnet: Ipv4Net,
resolve: bool,
progress: Option<Arc<dyn Fn(DiscoverProgress) + Send + Sync>>,
) -> Result<Vec<Host>> {
crate::ops::validation::ensure_subnet_limit(&subnet, &subnet.to_string())?;
let ips: Vec<IpAddr> = subnet.hosts().map(IpAddr::V4).collect();
// 0) Start the mDNS browse now, so it overlaps everything below.
// See discover/mdns_fusion.rs for why this exists and what it costs.
#[cfg(feature = "mdns")]
let mdns_browse = mdns_fusion::spawn_browse();
// 1) Ping first (fast), to avoid reverse-DNS work on dead hosts.
let total = ips.len();
let completed = Arc::new(AtomicUsize::new(0));
let found = Arc::new(AtomicUsize::new(0));
let ping_results = stream::iter(ips)
.map(|ip| {
let scanner = self.ping_scanner.clone();
let timeout_ms = self.ping_timeout_ms;
let completed = completed.clone();
let found = found.clone();
let progress = progress.clone();
async move {
let res = scanner.ping(ip, timeout_ms).await;
let done = completed.fetch_add(1, Ordering::SeqCst) + 1;
// Capture `found` *after* both atomic updates settle.
// Using a single load for both branches avoids the
// previous inconsistency where an unreachable IP could
// be reported with a `found` count that was updated by
// a concurrent task between this task's two loads.
if res.alive {
found.fetch_add(1, Ordering::SeqCst);
}
let found_snapshot = found.load(Ordering::SeqCst);
if let Some(cb) = &progress {
if res.alive || done == total || done.is_multiple_of(10) {
cb(DiscoverProgress {
phase: DiscoverPhase::Ping,
completed: done,
total,
found: found_snapshot,
ip,
});
}
}
res
}
})
.buffer_unordered(self.concurrency)
.collect::<Vec<crate::ping::PingResult>>()
.await;
let alive: Vec<crate::ping::PingResult> =
ping_results.into_iter().filter(|r| r.alive).collect();
// 2) Load ARP/neighbor table once and reuse it.
//
// On a blocking thread: this shells out to `arp` on Windows and
// macOS, and running it inline parked a runtime worker on the child
// process for every discover and sweep.
let arp_map: HashMap<IpAddr, crate::arp::ArpEntry> =
tokio::task::spawn_blocking(NetworkManager::get_arp_table)
.await
.unwrap_or_else(|_| Ok(Vec::new()))
.unwrap_or_default()
.into_iter()
.map(|e| (e.ip, e))
.collect();
// 3) Optionally reverse-DNS alive hosts using a single resolver.
let hostname_map: HashMap<IpAddr, Option<String>> = if resolve {
let ips = alive.iter().map(|r| r.ip).collect::<Vec<_>>();
let concurrency = self.concurrency.min(32);
let total = ips.len();
let completed = Arc::new(AtomicUsize::new(0));
let resolved = Arc::new(AtomicUsize::new(0));
stream::iter(ips)
.map(|ip| {
let dns_timeout_ms = self.dns_timeout_ms;
let completed = completed.clone();
let resolved = resolved.clone();
let progress = progress.clone();
async move {
let name =
crate::dns::reverse_lookup_best_effort_timeout(ip, dns_timeout_ms)
.await;
let done = completed.fetch_add(1, Ordering::SeqCst) + 1;
if name.is_some() {
resolved.fetch_add(1, Ordering::SeqCst);
}
let resolved_count = resolved.load(Ordering::SeqCst);
if let Some(cb) = &progress {
if name.is_some() || done == total || done.is_multiple_of(5) {
cb(DiscoverProgress {
phase: DiscoverPhase::Resolve,
completed: done,
total,
found: resolved_count,
ip,
});
}
}
(ip, name)
}
})
.buffer_unordered(concurrency)
.collect::<Vec<(IpAddr, Option<String>)>>()
.await
.into_iter()
.collect()
} else {
HashMap::new()
};
// 4) Fuse in mDNS names for hosts the reverse lookup left unnamed.
// Fill-blanks-only, so nothing that resolves today changes. Placed
// after `hostname_map` so ARP-only neighbours are named too -- a
// device that never answered a probe is the one least likely to have
// a PTR record.
#[cfg(feature = "mdns")]
let mdns_names = mdns_fusion::collect_names(mdns_browse).await;
#[cfg(not(feature = "mdns"))]
let mdns_names: HashMap<IpAddr, String> = HashMap::new();
let build = |ip: IpAddr, rtt_ms: Option<u64>, found_by: FoundBy| {
let mac_entry = arp_map.get(&ip);
let mac_str = mac_entry.map(|e| e.mac.to_string());
let vendor = mac_entry
.and_then(|e| e.vendor.clone())
.or_else(|| mac_str.as_deref().and_then(lookup_vendor));
// One path for both builds: without the mdns feature the map is
// simply empty, so the fallback never fires.
let (hostname, hostname_source) = match hostname_map.get(&ip).cloned().unwrap_or(None) {
Some(name) => (Some(name), Some(NameSource::Reverse)),
None => match mdns_names.get(&ip) {
Some(name) => (Some(name.clone()), Some(NameSource::Mdns)),
None => (None, None),
},
};
Host {
ip,
hostname,
mac: mac_str,
vendor,
rtt_ms,
found_by,
hostname_source,
}
};
let mut hosts: Vec<Host> = alive
.iter()
.map(|r| build(r.ip, r.rtt_ms, FoundBy::Probe))
.collect();
// Anything the OS has an ARP entry for is on this link, whether or
// not it answered us. Plenty of devices do not: consumer IoT
// routinely ignores ICMP, and Windows drops echo requests by
// default. Discarding them meant discovery reported a fraction of
// the network -- measured at 13 of 25 known devices on one ordinary
// LAN -- while the table needed to find them was already loaded, and
// used only to decorate the hosts that had replied.
let answered: std::collections::HashSet<IpAddr> = alive.iter().map(|r| r.ip).collect();
let mut neighbors: Vec<IpAddr> = arp_map
.keys()
.copied()
.filter(|ip| !answered.contains(ip))
.filter(|ip| match ip {
// Only within the range that was asked for. The neighbour
// table spans every interface, so it holds addresses from
// other subnets entirely.
IpAddr::V4(v4) => is_host_address(&subnet, *v4),
IpAddr::V6(_) => false,
})
.filter(|ip| arp_map.get(ip).is_none_or(|e| e.mac.bytes()[0] & 1 == 0))
.collect();
neighbors.sort_unstable();
hosts.extend(
neighbors
.into_iter()
.map(|ip| build(ip, None, FoundBy::Neighbor)),
);
Ok(hosts)
}
}
/// An address in `subnet` that a device can hold: not the network address
/// and not the broadcast address, except on /31 and /32 where every address
/// is usable. The Windows neighbour table lists x.x.x.255
/// (ff:ff:ff:ff:ff:ff), and discovery reported it as a host.
fn is_host_address(subnet: &Ipv4Net, ip: std::net::Ipv4Addr) -> bool {
subnet.contains(&ip)
&& (subnet.prefix_len() >= 31 || (ip != subnet.network() && ip != subnet.broadcast()))
}
#[cfg(test)]
mod host_address_tests {
use super::is_host_address;
#[test]
fn network_and_broadcast_are_not_hosts() {
let net = "192.168.1.0/24".parse().unwrap();
assert!(!is_host_address(&net, "192.168.1.0".parse().unwrap()));
assert!(!is_host_address(&net, "192.168.1.255".parse().unwrap()));
assert!(is_host_address(&net, "192.168.1.1".parse().unwrap()));
assert!(is_host_address(&net, "192.168.1.254".parse().unwrap()));
assert!(!is_host_address(&net, "192.168.2.1".parse().unwrap()));
}
#[test]
fn every_address_counts_on_a_31_or_32() {
let p2p = "10.0.0.0/31".parse().unwrap();
assert!(is_host_address(&p2p, "10.0.0.0".parse().unwrap()));
assert!(is_host_address(&p2p, "10.0.0.1".parse().unwrap()));
let single = "10.0.0.5/32".parse().unwrap();
assert!(is_host_address(&single, "10.0.0.5".parse().unwrap()));
}
}