use std::path::PathBuf;
use std::sync::Arc;
use std::time::{Duration, Instant};
use futures::stream::{FuturesUnordered, StreamExt};
use hickory_resolver::TokioResolver;
use crate::error::Result;
use crate::limit::{Pacer, PacingLimits};
use crate::lookup::outcome::{Finding, Reason, Source, Status};
use crate::lookup::referral;
use crate::lookup::registry::{Freshness, ServiceMap};
use crate::lookup::resolve::{self, DnsVerdict};
use crate::lookup::whois::{self, Server, Servers};
use crate::lookup::{rdap, registration};
use crate::tld::Suffix;
use crate::user_agent;
#[derive(Debug, Clone)]
pub struct Settings {
pub pacing: PacingLimits,
pub timeout: Duration,
pub cache_path: PathBuf,
pub refresh: bool,
pub source_policy: SourcePolicy,
pub registry_servers: Option<PathBuf>,
pub text_servers: Option<PathBuf>,
pub replace_servers: bool,
pub allow_referrals: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum SourcePolicy {
#[default]
Auto,
Registry,
Text,
Dns,
}
use crate::limit::RESOLVER_HOST;
#[derive(Debug)]
pub struct Engine {
client: reqwest::Client,
resolver: TokioResolver,
services: ServiceMap,
servers: Servers,
referral_cache: tokio::sync::Mutex<
std::collections::HashMap<String, Arc<tokio::sync::OnceCell<Option<String>>>>,
>,
pacer: Arc<Pacer>,
settings: Settings,
}
impl Engine {
pub async fn build(settings: Settings) -> Result<Self> {
let resolver = resolve::build(settings.timeout)?;
let client = reqwest::Client::builder()
.user_agent(user_agent())
.timeout(settings.timeout)
.connect_timeout(settings.timeout)
.dns_resolver(resolve::HttpResolver::new(&resolver))
.redirect(reqwest::redirect::Policy::none())
.https_only(true)
.build()
.map_err(|source| crate::Error::NetworkUnreachable {
source: Box::new(source),
})?;
let skip_published = !matches!(
settings.source_policy,
SourcePolicy::Auto | SourcePolicy::Registry
) || (settings.replace_servers && settings.registry_servers.is_some());
let (mut services, _freshness) = if skip_published {
(ServiceMap::default(), Freshness::Cached)
} else {
ServiceMap::load(&client, &settings.cache_path, settings.refresh).await?
};
if let Some(path) = &settings.registry_servers {
services.merge(ServiceMap::from_file(path)?);
}
let text_possible = matches!(
settings.source_policy,
SourcePolicy::Auto | SourcePolicy::Text
);
let skip_bundled_servers =
!text_possible || (settings.replace_servers && settings.text_servers.is_some());
let mut servers = if skip_bundled_servers {
Servers::default()
} else {
Servers::bundled()?
};
if let Some(path) = &settings.text_servers {
servers.merge(Servers::from_file(path)?);
}
Ok(Self {
pacer: Arc::new(Pacer::new(settings.pacing.clone())),
servers,
referral_cache: tokio::sync::Mutex::new(std::collections::HashMap::new()),
client,
resolver,
services,
settings,
})
}
pub async fn check_name(&self, name: &str, suffix: &Suffix) -> Finding {
let started = Instant::now();
let domain = format!("{name}.{suffix}");
if domain.len() > 253 {
return Finding::unknown(name, suffix, Reason::NotRegistrable, started.elapsed());
}
let mut reason = Reason::NoService;
if matches!(
self.settings.source_policy,
SourcePolicy::Auto | SourcePolicy::Registry
) && let Some(services) = self.services.for_suffix(suffix.as_str())
{
let (answer, responder) = rdap::query(
&self.client,
&self.pacer,
services,
&domain,
self.settings.timeout,
)
.await;
match answer {
rdap::Verdict::Available => {
return Finding {
domain,
name: name.to_owned(),
suffix: suffix.clone(),
status: Status::Available,
source: Some(Source::Registry),
elapsed: started.elapsed(),
responder,
registration: None,
};
}
rdap::Verdict::Taken(body) => {
return Finding {
domain,
name: name.to_owned(),
suffix: suffix.clone(),
status: Status::Taken,
source: Some(Source::Registry),
elapsed: started.elapsed(),
responder,
registration: Some(registration::parse(&body)),
};
}
rdap::Verdict::Unknown(refused) => {
if self.settings.source_policy == SourcePolicy::Registry {
return Finding::unknown(name, suffix, refused, started.elapsed());
}
reason = refused;
}
}
} else if self.settings.source_policy == SourcePolicy::Registry {
return Finding::unknown(name, suffix, Reason::NoService, started.elapsed());
}
let text_allowed = matches!(
self.settings.source_policy,
SourcePolicy::Auto | SourcePolicy::Text
);
if text_allowed && let Some(server) = self.servers.for_suffix(suffix.as_str()) {
match whois::query(
&self.resolver,
&self.pacer,
server,
&domain,
self.settings.timeout,
)
.await
{
whois::Verdict::Available => {
return Finding {
domain,
name: name.to_owned(),
suffix: suffix.clone(),
status: Status::Available,
source: Some(Source::Text),
elapsed: started.elapsed(),
responder: Some(crate::lookup::outcome::scrub(&server.host)),
registration: None,
};
}
whois::Verdict::Taken { .. } => {
return Finding {
domain,
name: name.to_owned(),
suffix: suffix.clone(),
status: Status::Taken,
source: Some(Source::Text),
elapsed: started.elapsed(),
responder: Some(crate::lookup::outcome::scrub(&server.host)),
registration: None,
};
}
whois::Verdict::Unknown(refused) => {
if self.settings.source_policy == SourcePolicy::Text {
return Finding::unknown(name, suffix, refused, started.elapsed());
}
if !matches!(refused, Reason::NoService) {
reason = refused;
}
}
}
}
if text_allowed
&& self.settings.allow_referrals
&& let Some(found) = self
.referred_server(suffix)
.await
.filter(|host| crate::lookup::registry::is_public_host(host))
{
let server = Server {
host: found,
available_phrase: String::new(),
};
match whois::query(
&self.resolver,
&self.pacer,
&server,
&domain,
self.settings.timeout,
)
.await
{
whois::Verdict::Available => {
return Finding {
domain,
name: name.to_owned(),
suffix: suffix.clone(),
status: Status::Available,
source: Some(Source::Text),
elapsed: started.elapsed(),
responder: Some(crate::lookup::outcome::scrub(&server.host)),
registration: None,
};
}
whois::Verdict::Taken { .. } => {
return Finding {
domain,
name: name.to_owned(),
suffix: suffix.clone(),
status: Status::Taken,
source: Some(Source::Text),
elapsed: started.elapsed(),
responder: Some(crate::lookup::outcome::scrub(&server.host)),
registration: None,
};
}
whois::Verdict::Unknown(refused) => {
if !matches!(refused, Reason::NoService) {
reason = refused;
}
}
}
}
if self.settings.source_policy == SourcePolicy::Text {
return Finding::unknown(name, suffix, reason, started.elapsed());
}
let dns_lease = self.pacer.acquire(RESOLVER_HOST).await.ok();
let dns_verdict = resolve::query(&self.resolver, &domain).await;
drop(dns_lease);
match dns_verdict {
DnsVerdict::InUse => Finding {
domain,
name: name.to_owned(),
suffix: suffix.clone(),
status: Status::Taken,
source: Some(Source::Dns),
elapsed: started.elapsed(),
responder: None,
registration: None,
},
DnsVerdict::Absent => Finding::unknown(name, suffix, reason, started.elapsed()),
DnsVerdict::NoAnswer => {
if matches!(reason, Reason::NoService) {
reason = Reason::Unreachable;
}
Finding::unknown(name, suffix, reason, started.elapsed())
}
}
}
pub async fn check_domain(
&self,
catalog: &crate::tld::Catalog,
domain: &str,
) -> Option<Finding> {
let (name, suffix) = catalog.split_domain(domain)?;
Some(self.check_name(&name, &suffix).await)
}
pub async fn sweep(
&self,
names: &[String],
suffixes: &[Suffix],
mut answered: impl FnMut(&Finding),
) -> Vec<Finding> {
let window = self
.settings
.pacing
.total_concurrency
.max(1)
.saturating_mul(2);
let mut pending = names
.iter()
.flat_map(|name| suffixes.iter().map(move |suffix| (name, suffix)));
let mut work = FuturesUnordered::new();
let mut findings = Vec::new();
for (name, suffix) in pending.by_ref().take(window) {
work.push(self.check_name(name, suffix));
}
while let Some(finding) = work.next().await {
answered(&finding);
findings.push(finding);
if let Some((name, suffix)) = pending.next() {
work.push(self.check_name(name, suffix));
}
}
findings.sort_by(|a, b| a.domain.cmp(&b.domain));
findings
}
async fn referred_server(&self, suffix: &Suffix) -> Option<String> {
let cell = {
let mut cache = self.referral_cache.lock().await;
Arc::clone(
cache
.entry(suffix.as_str().to_owned())
.or_insert_with(|| Arc::new(tokio::sync::OnceCell::new())),
)
};
cell.get_or_init(|| async {
referral::query(
&self.resolver,
&self.pacer,
suffix.as_str(),
self.settings.timeout,
)
.await
.text_host
})
.await
.clone()
}
pub async fn dns_records(&self, domain: &str) -> crate::lookup::DnsRecords {
resolve::dns_records(&self.resolver, domain).await
}
pub async fn paused(&self) -> Vec<crate::limit::PausedHost> {
self.pacer.paused_hosts().await
}
}
#[cfg(test)]
mod tests {
use tempfile::tempdir;
use super::*;
use crate::lookup::outcome::Tally;
use crate::tld::Catalog;
const UNSERVED: &str = "zzzz-no-such-extension";
const ALSO_UNSERVED: &str = "yyyy-no-such-extension";
fn suffix(value: &str) -> Suffix {
Suffix::parse(value).expect("the test suffix parses")
}
fn grounded_settings(source_policy: SourcePolicy, dir: &std::path::Path) -> Settings {
Settings {
pacing: PacingLimits::default(),
timeout: Duration::from_secs(1),
cache_path: dir.join("servers.json"),
refresh: false,
source_policy,
registry_servers: None,
text_servers: None,
replace_servers: false,
allow_referrals: false,
}
}
async fn text_only_engine(dir: &std::path::Path) -> Engine {
Engine::build(grounded_settings(SourcePolicy::Text, dir))
.await
.expect("an engine that never leaves the machine")
}
async fn registry_only_engine(dir: &std::path::Path) -> Engine {
let list = dir.join("services.json");
std::fs::write(
&list,
r#"{"services":[[["com"],["https://rdap.example.test/com/"]]]}"#,
)
.expect("the service list is written");
Engine::build(Settings {
registry_servers: Some(list),
replace_servers: true,
..grounded_settings(SourcePolicy::Registry, dir)
})
.await
.expect("an engine that never leaves the machine")
}
#[tokio::test]
async fn asking_the_registry_only_reports_unknown_when_the_extension_has_no_service() {
let dir = tempdir().expect("temp dir");
let engine = registry_only_engine(dir.path()).await;
let finding = engine.check_name("example", &suffix(UNSERVED)).await;
assert_eq!(finding.status, Status::Unknown(Reason::NoService));
assert!(!finding.is_available());
assert_eq!(finding.source, None);
assert_eq!(finding.responder, None);
assert_eq!(finding.domain, format!("example.{UNSERVED}"));
assert_eq!(finding.name, "example");
}
#[tokio::test]
async fn the_text_protocol_alone_reports_unknown_when_no_server_answers_for_the_extension() {
let dir = tempdir().expect("temp dir");
let engine = text_only_engine(dir.path()).await;
let finding = engine.check_name("example", &suffix(UNSERVED)).await;
assert_eq!(finding.status, Status::Unknown(Reason::NoService));
assert!(!finding.is_available());
assert_eq!(finding.source, None);
}
#[tokio::test]
async fn a_sweep_nothing_can_answer_reports_every_row_unknown_and_none_free() {
let dir = tempdir().expect("temp dir");
let engine = text_only_engine(dir.path()).await;
let mut answered = 0_usize;
let findings = engine
.sweep(
&["beta".to_owned(), "alpha".to_owned()],
&[suffix(UNSERVED), suffix(ALSO_UNSERVED)],
|_| answered += 1,
)
.await;
assert_eq!(
answered,
findings.len(),
"every answer is handed to the caller as it lands"
);
assert_eq!(findings.len(), 4);
assert!(findings.iter().all(|finding| finding.status.is_unknown()));
let tally = Tally::of(&findings);
assert_eq!(tally.available, 0);
assert_eq!(tally.unknown, 4);
}
#[tokio::test]
async fn a_sweep_hands_back_one_row_per_pair_in_domain_order() {
let dir = tempdir().expect("temp dir");
let engine = text_only_engine(dir.path()).await;
let findings = engine
.sweep(
&["beta".to_owned(), "alpha".to_owned()],
&[suffix(UNSERVED), suffix(ALSO_UNSERVED)],
|_| {},
)
.await;
let domains: Vec<&str> = findings
.iter()
.map(|finding| finding.domain.as_str())
.collect();
let mut expected = domains.clone();
expected.sort_unstable();
assert_eq!(domains, expected);
assert_eq!(
domains.first(),
Some(&format!("alpha.{ALSO_UNSERVED}").as_str())
);
}
#[tokio::test]
async fn an_empty_sweep_asks_nothing_and_returns_nothing() {
let dir = tempdir().expect("temp dir");
let engine = text_only_engine(dir.path()).await;
let mut answered = 0_usize;
assert!(
engine
.sweep(&[], &[suffix(UNSERVED)], |_| answered += 1)
.await
.is_empty()
);
assert!(
engine
.sweep(&["alpha".to_owned()], &[], |_| answered += 1)
.await
.is_empty()
);
assert_eq!(answered, 0, "nothing to check means nothing is reported");
}
#[tokio::test]
async fn a_name_typed_in_full_is_checked_exactly_as_given() {
let dir = tempdir().expect("temp dir");
let engine = text_only_engine(dir.path()).await;
let catalog = Catalog::bundled().expect("the bundled catalog parses");
let finding = engine
.check_domain(&catalog, &format!("example.{UNSERVED}"))
.await
.expect("a domain with an extension splits");
assert_eq!(finding.domain, format!("example.{UNSERVED}"));
assert_eq!(finding.suffix.as_str(), UNSERVED);
assert!(finding.status.is_unknown());
}
#[tokio::test]
async fn a_bare_name_with_no_extension_is_not_checked_as_a_domain() {
let dir = tempdir().expect("temp dir");
let engine = text_only_engine(dir.path()).await;
let catalog = Catalog::bundled().expect("the bundled catalog parses");
assert!(engine.check_domain(&catalog, "example").await.is_none());
}
#[tokio::test]
async fn a_fresh_engine_holds_no_registry_back() {
let dir = tempdir().expect("temp dir");
let engine = text_only_engine(dir.path()).await;
assert!(engine.paused().await.is_empty());
}
#[test]
fn the_default_source_uses_the_registry_first() {
assert_eq!(SourcePolicy::default(), SourcePolicy::Auto);
}
#[test]
fn settings_carry_everything_a_run_needs() {
let settings = Settings {
pacing: PacingLimits::default(),
timeout: Duration::from_secs(10),
cache_path: PathBuf::from("/tmp/reserve/servers.json"),
refresh: false,
source_policy: SourcePolicy::Auto,
registry_servers: None,
text_servers: None,
replace_servers: false,
allow_referrals: true,
};
assert_eq!(settings.timeout, Duration::from_secs(10));
assert!(!settings.refresh);
}
}