Skip to main content

reserve_core/lookup/
engine.rs

1//! Running a sweep. The registry is asked first because only it can prove a name free.
2
3use 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    /// @docgen DNS can prove a name is taken and never that one is free.
41    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    /// @docgen One cell per extension, so a sweep asks IANA at most once for each without the map lock spanning the network call.
53    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            // @docgen A registry that answers with a redirect could otherwise steer the request to an internal address.
70            .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        // @docgen Under a registry-only or DNS-only run nothing ever reads this table, so parsing 884 entries is startup cost for nothing.
91        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        // @docgen The 253-octet limit applies to the whole name, so two valid halves can still join into an invalid one.
121        if domain.len() > 253 {
122            return Finding::unknown(name, suffix, Reason::NotRegistrable, started.elapsed());
123        }
124        // @docgen DNS silence must never overwrite a reason a registry actually gave, so the most informative one is kept.
125        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        // @docgen For many country registries the text protocol is the only thing that answers, and it can prove a name free.
178        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        // @docgen A bundled host can go stale, so IANA is asked who serves the extension today before giving up.
228        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        // @docgen The resolver needs a permit like every other source, or the whole fan-out lands on it at once.
285        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        // @docgen A DNS miss proves nothing: a registered name that was never delegated looks identical to a free one.
290        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    /// @docgen A name the user spelled out in full is checked as given, whatever extensions the run selected.
312    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    /// @docgen The pacer bounds requests in flight, not futures allocated, so the product is streamed rather than collected.
322    /// @docgen The caller is handed each answer as it lands so a waiting person can be shown progress; the core still owns no stream.
323    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    /// @docgen The map lock is dropped before the query, so extensions resolve in parallel while callers sharing one still wait for a single ask.
358    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}