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
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    /// @docgen Cached so one sweep asks IANA at most once per extension.
53    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            // @docgen A registry that answers with a redirect could otherwise steer the request to an internal address.
68            .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        // @docgen The 253-octet limit applies to the whole name, so two valid halves can still join into an invalid one.
112        if domain.len() > 253 {
113            return Finding::unknown(name, suffix, Reason::NotRegistrable, started.elapsed());
114        }
115        // @docgen DNS silence must never overwrite a reason a registry actually gave, so the most informative one is kept.
116        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        // @docgen For many country registries the text protocol is the only thing that answers, and it can prove a name free.
169        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        // @docgen A bundled host can go stale, so IANA is asked who serves the extension today before giving up.
219        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        // @docgen The resolver needs a permit like every other source, or the whole fan-out lands on it at once.
273        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        // @docgen A DNS miss proves nothing: a registered name that was never delegated looks identical to a free one.
278        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    /// @docgen A name the user spelled out in full is checked as given, whatever extensions the run selected.
300    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    /// @docgen The pacer bounds requests in flight, not futures allocated, so the product is streamed rather than collected.
310    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    /// @docgen The lock is held across the query so the first caller asks IANA and the rest wait for its answer.
339    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}