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