1use std::path::PathBuf;
4use std::sync::Arc;
5use std::time::{Duration, Instant};
6
7use futures::stream::{FuturesUnordered, StreamExt};
8use hickory_resolver::TokioResolver;
9
10use crate::error::Result;
11use crate::limit::{Pacer, PacingLimits};
12use crate::lookup::outcome::{Finding, Reason, Source, Status};
13use crate::lookup::referral;
14use crate::lookup::registry::{Freshness, ServiceMap};
15use crate::lookup::resolve::{self, DnsVerdict};
16use crate::lookup::whois::{self, Server, Servers};
17use crate::lookup::{rdap, registration};
18use crate::tld::Suffix;
19use crate::user_agent;
20
21#[derive(Debug, Clone)]
22pub struct Settings {
23 pub pacing: PacingLimits,
24 pub timeout: Duration,
25 pub cache_path: PathBuf,
26 pub refresh: bool,
27 pub source_policy: SourcePolicy,
28 pub registry_servers: Option<PathBuf>,
29 pub text_servers: Option<PathBuf>,
30 pub replace_servers: bool,
31 pub allow_referrals: bool,
32}
33
34#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
35pub enum SourcePolicy {
36 #[default]
37 Auto,
38 Registry,
39 Text,
40 Dns,
42}
43
44const RESOLVER_HOST: &str = "dns.resolver.local";
45
46#[derive(Debug)]
47pub struct Engine {
48 client: reqwest::Client,
49 resolver: TokioResolver,
50 services: ServiceMap,
51 servers: Servers,
52 referral_cache: tokio::sync::Mutex<std::collections::HashMap<String, Option<String>>>,
54 pacer: Arc<Pacer>,
55 settings: Settings,
56}
57
58impl Engine {
59 pub async fn build(settings: Settings) -> Result<Self> {
60 let resolver = resolve::build(settings.timeout)?;
61
62 let client = reqwest::Client::builder()
63 .user_agent(user_agent())
64 .timeout(settings.timeout)
65 .connect_timeout(settings.timeout)
66 .dns_resolver(resolve::HttpResolver::new(&resolver))
67 .redirect(reqwest::redirect::Policy::none())
69 .https_only(true)
70 .build()
71 .map_err(|source| crate::Error::NetworkUnreachable {
72 source: Box::new(source),
73 })?;
74
75 let skip_published = !matches!(
76 settings.source_policy,
77 SourcePolicy::Auto | SourcePolicy::Registry
78 ) || (settings.replace_servers && settings.registry_servers.is_some());
79 let (mut services, _freshness) = if skip_published {
80 (ServiceMap::default(), Freshness::Cached)
81 } else {
82 ServiceMap::load(&client, &settings.cache_path, settings.refresh).await?
83 };
84 if let Some(path) = &settings.registry_servers {
85 services.merge(ServiceMap::from_file(path)?);
86 }
87
88 let mut servers = if settings.replace_servers && settings.text_servers.is_some() {
89 Servers::default()
90 } else {
91 Servers::bundled()?
92 };
93 if let Some(path) = &settings.text_servers {
94 servers.merge(Servers::from_file(path)?);
95 }
96
97 Ok(Self {
98 pacer: Arc::new(Pacer::new(settings.pacing.clone())),
99 servers,
100 referral_cache: tokio::sync::Mutex::new(std::collections::HashMap::new()),
101 client,
102 resolver,
103 services,
104 settings,
105 })
106 }
107
108 pub async fn check_name(&self, name: &str, suffix: &Suffix) -> Finding {
109 let started = Instant::now();
110 let domain = format!("{name}.{suffix}");
111 if domain.len() > 253 {
113 return Finding::unknown(name, suffix, Reason::NotRegistrable, started.elapsed());
114 }
115 let mut reason = Reason::NoService;
117
118 if matches!(
119 self.settings.source_policy,
120 SourcePolicy::Auto | SourcePolicy::Registry
121 ) && let Some(services) = self.services.for_suffix(suffix.as_str())
122 {
123 let (answer, responder) = rdap::query(
124 &self.client,
125 &self.pacer,
126 services,
127 &domain,
128 self.settings.timeout,
129 )
130 .await;
131
132 match answer {
133 rdap::Verdict::Available => {
134 return Finding {
135 domain,
136 name: name.to_owned(),
137 suffix: suffix.clone(),
138 status: Status::Available,
139 source: Some(Source::Registry),
140 elapsed: started.elapsed(),
141 responder,
142 registration: None,
143 };
144 }
145 rdap::Verdict::Taken(body) => {
146 return Finding {
147 domain,
148 name: name.to_owned(),
149 suffix: suffix.clone(),
150 status: Status::Taken,
151 source: Some(Source::Registry),
152 elapsed: started.elapsed(),
153 responder,
154 registration: Some(registration::parse(&body)),
155 };
156 }
157 rdap::Verdict::Unknown(refused) => {
158 if self.settings.source_policy == SourcePolicy::Registry {
159 return Finding::unknown(name, suffix, refused, started.elapsed());
160 }
161 reason = refused;
162 }
163 }
164 } else if self.settings.source_policy == SourcePolicy::Registry {
165 return Finding::unknown(name, suffix, Reason::NoService, started.elapsed());
166 }
167
168 let text_allowed = matches!(
170 self.settings.source_policy,
171 SourcePolicy::Auto | SourcePolicy::Text
172 );
173 if text_allowed && let Some(server) = self.servers.for_suffix(suffix.as_str()) {
174 match whois::query(
175 &self.resolver,
176 &self.pacer,
177 server,
178 &domain,
179 self.settings.timeout,
180 )
181 .await
182 {
183 whois::Verdict::Available => {
184 return Finding {
185 domain,
186 name: name.to_owned(),
187 suffix: suffix.clone(),
188 status: Status::Available,
189 source: Some(Source::Text),
190 elapsed: started.elapsed(),
191 responder: Some(server.host.clone()),
192 registration: None,
193 };
194 }
195 whois::Verdict::Taken { .. } => {
196 return Finding {
197 domain,
198 name: name.to_owned(),
199 suffix: suffix.clone(),
200 status: Status::Taken,
201 source: Some(Source::Text),
202 elapsed: started.elapsed(),
203 responder: Some(server.host.clone()),
204 registration: None,
205 };
206 }
207 whois::Verdict::Unknown(refused) => {
208 if self.settings.source_policy == SourcePolicy::Text {
209 return Finding::unknown(name, suffix, refused, started.elapsed());
210 }
211 if !matches!(refused, Reason::NoService) {
212 reason = refused;
213 }
214 }
215 }
216 }
217
218 if text_allowed
220 && self.settings.allow_referrals
221 && let Some(found) = self.referred_server(suffix).await
222 {
223 let server = Server {
224 host: found,
225 available_phrase: String::new(),
226 };
227 match whois::query(
228 &self.resolver,
229 &self.pacer,
230 &server,
231 &domain,
232 self.settings.timeout,
233 )
234 .await
235 {
236 whois::Verdict::Available => {
237 return Finding {
238 domain,
239 name: name.to_owned(),
240 suffix: suffix.clone(),
241 status: Status::Available,
242 source: Some(Source::Text),
243 elapsed: started.elapsed(),
244 responder: Some(server.host),
245 registration: None,
246 };
247 }
248 whois::Verdict::Taken { .. } => {
249 return Finding {
250 domain,
251 name: name.to_owned(),
252 suffix: suffix.clone(),
253 status: Status::Taken,
254 source: Some(Source::Text),
255 elapsed: started.elapsed(),
256 responder: Some(server.host),
257 registration: None,
258 };
259 }
260 whois::Verdict::Unknown(refused) => {
261 if !matches!(refused, Reason::NoService) {
262 reason = refused;
263 }
264 }
265 }
266 }
267
268 if self.settings.source_policy == SourcePolicy::Text {
269 return Finding::unknown(name, suffix, reason, started.elapsed());
270 }
271
272 let dns_lease = self.pacer.acquire(RESOLVER_HOST).await.ok();
274 let dns_verdict = resolve::query(&self.resolver, &domain).await;
275 drop(dns_lease);
276
277 match dns_verdict {
279 DnsVerdict::InUse => Finding {
280 domain,
281 name: name.to_owned(),
282 suffix: suffix.clone(),
283 status: Status::Taken,
284 source: Some(Source::Dns),
285 elapsed: started.elapsed(),
286 responder: None,
287 registration: None,
288 },
289 DnsVerdict::Absent => Finding::unknown(name, suffix, reason, started.elapsed()),
290 DnsVerdict::NoAnswer => {
291 if matches!(reason, Reason::NoService) {
292 reason = Reason::Unreachable;
293 }
294 Finding::unknown(name, suffix, reason, started.elapsed())
295 }
296 }
297 }
298
299 pub async fn check_domain(
301 &self,
302 catalog: &crate::tld::Catalog,
303 domain: &str,
304 ) -> Option<Finding> {
305 let (name, suffix) = catalog.split_domain(domain)?;
306 Some(self.check_name(&name, &suffix).await)
307 }
308
309 pub async fn sweep(&self, names: &[String], suffixes: &[Suffix]) -> Vec<Finding> {
311 let window = self
312 .settings
313 .pacing
314 .total_concurrency
315 .max(1)
316 .saturating_mul(2);
317 let mut pending = names
318 .iter()
319 .flat_map(|name| suffixes.iter().map(move |suffix| (name, suffix)));
320
321 let mut work = FuturesUnordered::new();
322 let mut findings = Vec::new();
323
324 for (name, suffix) in pending.by_ref().take(window) {
325 work.push(self.check_name(name, suffix));
326 }
327 while let Some(finding) = work.next().await {
328 findings.push(finding);
329 if let Some((name, suffix)) = pending.next() {
330 work.push(self.check_name(name, suffix));
331 }
332 }
333
334 findings.sort_by(|a, b| a.domain.cmp(&b.domain));
335 findings
336 }
337
338 async fn referred_server(&self, suffix: &Suffix) -> Option<String> {
340 let key = suffix.as_str().to_owned();
341 let mut cache = self.referral_cache.lock().await;
342 if let Some(found) = cache.get(&key) {
343 return found.clone();
344 }
345 let delegation = referral::query(
346 &self.resolver,
347 &self.pacer,
348 suffix.as_str(),
349 self.settings.timeout,
350 )
351 .await;
352 let text_host = delegation.text_host.clone();
353 cache.insert(key, text_host.clone());
354 text_host
355 }
356
357 pub async fn dns_records(&self, domain: &str) -> crate::lookup::DnsRecords {
358 resolve::dns_records(&self.resolver, domain).await
359 }
360
361 pub async fn paused(&self) -> Vec<crate::limit::PausedHost> {
362 self.pacer.paused_hosts().await
363 }
364}
365
366#[cfg(test)]
367mod tests {
368 use tempfile::tempdir;
369
370 use super::*;
371 use crate::lookup::outcome::Tally;
372 use crate::tld::Catalog;
373
374 const UNSERVED: &str = "zzzz-no-such-extension";
375 const ALSO_UNSERVED: &str = "yyyy-no-such-extension";
376
377 fn suffix(value: &str) -> Suffix {
378 Suffix::parse(value).expect("the test suffix parses")
379 }
380
381 fn grounded_settings(source_policy: SourcePolicy, dir: &std::path::Path) -> Settings {
382 Settings {
383 pacing: PacingLimits::default(),
384 timeout: Duration::from_secs(1),
385 cache_path: dir.join("servers.json"),
386 refresh: false,
387 source_policy,
388 registry_servers: None,
389 text_servers: None,
390 replace_servers: false,
391 allow_referrals: false,
392 }
393 }
394
395 async fn text_only_engine(dir: &std::path::Path) -> Engine {
396 Engine::build(grounded_settings(SourcePolicy::Text, dir))
397 .await
398 .expect("an engine that never leaves the machine")
399 }
400
401 async fn registry_only_engine(dir: &std::path::Path) -> Engine {
402 let list = dir.join("services.json");
403 std::fs::write(
404 &list,
405 r#"{"services":[[["com"],["https://rdap.example.test/com/"]]]}"#,
406 )
407 .expect("the service list is written");
408
409 Engine::build(Settings {
410 registry_servers: Some(list),
411 replace_servers: true,
412 ..grounded_settings(SourcePolicy::Registry, dir)
413 })
414 .await
415 .expect("an engine that never leaves the machine")
416 }
417
418 #[tokio::test]
419 async fn asking_the_registry_only_reports_unknown_when_the_extension_has_no_service() {
420 let dir = tempdir().expect("temp dir");
421 let engine = registry_only_engine(dir.path()).await;
422
423 let finding = engine.check_name("example", &suffix(UNSERVED)).await;
424
425 assert_eq!(finding.status, Status::Unknown(Reason::NoService));
426 assert!(!finding.is_available());
427 assert_eq!(finding.source, None);
428 assert_eq!(finding.responder, None);
429 assert_eq!(finding.domain, format!("example.{UNSERVED}"));
430 assert_eq!(finding.name, "example");
431 }
432
433 #[tokio::test]
434 async fn the_text_protocol_alone_reports_unknown_when_no_server_answers_for_the_extension() {
435 let dir = tempdir().expect("temp dir");
436 let engine = text_only_engine(dir.path()).await;
437
438 let finding = engine.check_name("example", &suffix(UNSERVED)).await;
439
440 assert_eq!(finding.status, Status::Unknown(Reason::NoService));
441 assert!(!finding.is_available());
442 assert_eq!(finding.source, None);
443 }
444
445 #[tokio::test]
446 async fn a_sweep_nothing_can_answer_reports_every_row_unknown_and_none_free() {
447 let dir = tempdir().expect("temp dir");
448 let engine = text_only_engine(dir.path()).await;
449
450 let findings = engine
451 .sweep(
452 &["beta".to_owned(), "alpha".to_owned()],
453 &[suffix(UNSERVED), suffix(ALSO_UNSERVED)],
454 )
455 .await;
456
457 assert_eq!(findings.len(), 4);
458 assert!(findings.iter().all(|finding| finding.status.is_unknown()));
459
460 let tally = Tally::of(&findings);
461 assert_eq!(tally.available, 0);
462 assert_eq!(tally.unknown, 4);
463 }
464
465 #[tokio::test]
466 async fn a_sweep_hands_back_one_row_per_pair_in_domain_order() {
467 let dir = tempdir().expect("temp dir");
468 let engine = text_only_engine(dir.path()).await;
469
470 let findings = engine
471 .sweep(
472 &["beta".to_owned(), "alpha".to_owned()],
473 &[suffix(UNSERVED), suffix(ALSO_UNSERVED)],
474 )
475 .await;
476
477 let domains: Vec<&str> = findings
478 .iter()
479 .map(|finding| finding.domain.as_str())
480 .collect();
481 let mut expected = domains.clone();
482 expected.sort_unstable();
483 assert_eq!(domains, expected);
484 assert_eq!(
485 domains.first(),
486 Some(&format!("alpha.{ALSO_UNSERVED}").as_str())
487 );
488 }
489
490 #[tokio::test]
491 async fn an_empty_sweep_asks_nothing_and_returns_nothing() {
492 let dir = tempdir().expect("temp dir");
493 let engine = text_only_engine(dir.path()).await;
494
495 assert!(engine.sweep(&[], &[suffix(UNSERVED)]).await.is_empty());
496 assert!(engine.sweep(&["alpha".to_owned()], &[]).await.is_empty());
497 }
498
499 #[tokio::test]
500 async fn a_name_typed_in_full_is_checked_exactly_as_given() {
501 let dir = tempdir().expect("temp dir");
502 let engine = text_only_engine(dir.path()).await;
503 let catalog = Catalog::bundled().expect("the bundled catalog parses");
504
505 let finding = engine
506 .check_domain(&catalog, &format!("example.{UNSERVED}"))
507 .await
508 .expect("a domain with an extension splits");
509
510 assert_eq!(finding.domain, format!("example.{UNSERVED}"));
511 assert_eq!(finding.suffix.as_str(), UNSERVED);
512 assert!(finding.status.is_unknown());
513 }
514
515 #[tokio::test]
516 async fn a_bare_name_with_no_extension_is_not_checked_as_a_domain() {
517 let dir = tempdir().expect("temp dir");
518 let engine = text_only_engine(dir.path()).await;
519 let catalog = Catalog::bundled().expect("the bundled catalog parses");
520
521 assert!(engine.check_domain(&catalog, "example").await.is_none());
522 }
523
524 #[tokio::test]
525 async fn a_fresh_engine_holds_no_registry_back() {
526 let dir = tempdir().expect("temp dir");
527 let engine = text_only_engine(dir.path()).await;
528
529 assert!(engine.paused().await.is_empty());
530 }
531
532 #[test]
533 fn the_default_source_uses_the_registry_first() {
534 assert_eq!(SourcePolicy::default(), SourcePolicy::Auto);
535 }
536
537 #[test]
538 fn settings_carry_everything_a_run_needs() {
539 let settings = Settings {
540 pacing: PacingLimits::default(),
541 timeout: Duration::from_secs(10),
542 cache_path: PathBuf::from("/tmp/reserve/servers.json"),
543 refresh: false,
544 source_policy: SourcePolicy::Auto,
545 registry_servers: None,
546 text_servers: None,
547 replace_servers: false,
548 allow_referrals: true,
549 };
550 assert_eq!(settings.timeout, Duration::from_secs(10));
551 assert!(!settings.refresh);
552 }
553}