very fast at protocol indexer with flexible filtering, xrpc queries, cursor-backed event stream, and more, built on fjall
rust fjall at-protocol atproto indexer
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988//! socket-coupled egress capabilities for externally influenced endpoints.//! [`PublicEndpoint`] carries syntax proof only. The public HTTP and websocket//! resolvers carry the address proof to the connector: every fresh connection//! resolves the complete answer set, rejects it if any answer is not public, and//! only then gives the addresses to the socket implementation.
//! the IPv4/IPv6 classifier below is adapted from Svix's `is_allowed` and helper//! logic, pinned to commit `805721db7165a00275c92460fae4363f2fb9b720` in//! `svix-webhooks/server/svix-server/src/core/webhook_http_client.rs`. Svix is//! MIT licensed; its exact notice is retained in//! `licenses/svix-webhooks-MIT.txt`. The range policy and its mapped-IPv6//! canonicalization are intentionally retained rather than replaced with an//! ad-hoc decoder.
use std::future::Future;use std::net::{IpAddr, Ipv4Addr, Ipv6Addr, SocketAddr};use std::pin::Pin;use std::sync::{Arc, OnceLock};use std::time::Duration;
use futures::StreamExt;use ipnet::{Ipv4Net, Ipv6Net};use reqwest::dns::{Addrs, Name, Resolve, Resolving};use reqwest::redirect::Policy;use smol_str::SmolStr;use thiserror::Error;use url::{Host, Url};
pub(crate) const DNS_LOOKUP_TIMEOUT: Duration = Duration::from_secs(5);const PUBLIC_REQUEST_TIMEOUT: Duration = Duration::from_secs(30);/// a peer that drops our SYN answers nothing, so a blackholed connect would/// otherwise spend the whole request budget before the next candidate source/// gets a turn. applies to every public client, and never raised above the/// caller's own request timeout.const PUBLIC_CONNECT_TIMEOUT: Duration = Duration::from_secs(5);/// handle resolution runs on the read path of every mini doc, and one of its/// steps fetches a host the account itself names, so it gets its own budget.pub(crate) const HANDLE_REQUEST_TIMEOUT: Duration = Duration::from_secs(3);pub(crate) const HANDLE_CONNECT_TIMEOUT: Duration = Duration::from_secs(2);
/// A hostname or public IP literal accepted by the public egress syntax policy.////// This value contains no port, path, query, fragment, or userinfo. DNS names/// are lower-case IDNA ASCII and have no trailing dot. A literal is accepted/// only when its address is public according to [`is_private_ip`].#[derive(Debug, Clone, PartialEq, Eq, Hash)]pub(crate) struct PublicHost(SmolStr);
#[derive(Debug, Clone, PartialEq, Eq, Error)]pub(crate) enum HostError { #[error("PDS hostname is empty")] Empty, #[error("PDS hostname must not include a port, path, query, fragment, or userinfo")] NotHostname, #[error("PDS hostname is localhost or a local-only domain")] LocalName, #[error("PDS hostname is a non-public IP address")] NonPublicIp, #[error("PDS hostname is not a valid DNS name")] InvalidDnsName, #[error("PDS source URL must use wss with no port, path, query, fragment, or userinfo")] NonCanonicalUrl,}
impl PublicHost { /// Parse a host authority without resolving it. pub(crate) fn parse(input: &str) -> Result<Self, HostError> { let host = input.strip_suffix('.').unwrap_or(input); if host.is_empty() { return Err(HostError::Empty); } if host.contains('@') || (host.starts_with('[') && !host.ends_with(']')) { return Err(HostError::NotHostname); }
let literal_ip = host .strip_prefix('[') .and_then(|host| host.strip_suffix(']')) .unwrap_or(host); if let Ok(ip) = literal_ip.parse::<IpAddr>() { return (!is_private_ip(ip)) .then(|| Self(ip.to_string().into())) .ok_or(HostError::NonPublicIp); } if host.contains(':') { return Err(HostError::NotHostname); }
// Url owns IDNA and authority parsing. Constructing a host-only URL // also rejects malformed percent escapes and authority syntax. let parsed = Url::parse(&format!("wss://{host}/")).map_err(|_| HostError::InvalidDnsName)?; if !parsed.username().is_empty() || parsed.password().is_some() || parsed.port().is_some() || parsed.path() != "/" || parsed.query().is_some() || parsed.fragment().is_some() { return Err(HostError::NotHostname); } let host = parsed.host_str().ok_or(HostError::InvalidDnsName)?; if let Ok(ip) = host.parse::<IpAddr>() { return (!is_private_ip(ip)) .then(|| Self(ip.to_string().into())) .ok_or(HostError::NonPublicIp); } if matches!(host, "localhost" | "local" | "internal") || host.ends_with(".localhost") || host.ends_with(".local") || host.ends_with(".internal") { return Err(HostError::LocalName); } if !is_dns_name(host) { return Err(HostError::InvalidDnsName); }
Ok(Self(host.into())) }
pub(crate) fn as_str(&self) -> &str { &self.0 }
pub(crate) fn firehose_url(&self) -> Url { let mut url = Url::parse("wss://localhost/").expect("constant PDS URL must parse"); let host = self .0 .parse::<Ipv6Addr>() .map(|_| format!("[{}]", self.0)) .unwrap_or_else(|_| self.0.to_string()); url.set_host(Some(&host)) .expect("public host must be valid in a URL"); url }
/// Construct the canonical HTTPS base URL from this already-validated host. /// URL serialization owns IPv6 bracket handling; callers do not format an /// authority string by hand. #[cfg(feature = "indexer")] pub(crate) fn https_endpoint(&self) -> PublicEndpoint { let mut url = Url::parse("https://localhost/").expect("constant PDS URL must parse"); let host = self .0 .parse::<Ipv6Addr>() .map(|_| format!("[{}]", self.0)) .unwrap_or_else(|_| self.0.to_string()); url.set_host(Some(&host)) .expect("public host must be valid in a URL"); PublicEndpoint { url } }}
/// A URL whose authority passed the public syntax policy.#[derive(Debug, Clone, PartialEq, Eq)]pub(crate) struct PublicEndpoint { url: Url,}
#[derive(Debug, Clone, PartialEq, Eq, Error)]pub(crate) enum EndpointError { #[error("public endpoint scheme is not supported: {0}")] UnsupportedScheme(SmolStr), #[error("public endpoint must not include userinfo or a fragment")] UnsafeAuthority, #[error("public endpoint has no host")] MissingHost, #[error(transparent)] Host(#[from] HostError),}
impl PublicEndpoint { /// Parse a direct HTTP endpoint where a caller's protocol explicitly /// permits cleartext HTTP. HTTPS remains accepted. pub(crate) fn parse_http(url: &Url) -> Result<Self, EndpointError> { Self::parse_schemes(url, |scheme| matches!(scheme, "http" | "https")) }
/// Parse the public websocket endpoint. Public firehose connections are /// TLS-only; the trusted raw path remains available to operator config. pub(crate) fn parse_websocket(url: &Url) -> Result<Self, EndpointError> { Self::parse_schemes(url, |scheme| scheme == "wss") }
fn parse_schemes( url: &Url, supported: impl FnOnce(&str) -> bool, ) -> Result<Self, EndpointError> { if !supported(url.scheme()) { return Err(EndpointError::UnsupportedScheme(url.scheme().into())); } if !url.username().is_empty() || url.password().is_some() || url.fragment().is_some() { return Err(EndpointError::UnsafeAuthority); } // use Url::host rather than host_str so bracketed IPv6 and mapped IPv6 // literals are represented as addresses before the shared classifier runs. let host = match url.host().ok_or(EndpointError::MissingHost)? { Host::Domain(host) => PublicHost::parse(host)?, Host::Ipv4(ip) => PublicHost::parse(&ip.to_string())?, Host::Ipv6(ip) => PublicHost::parse(&ip.to_string())?, }; let mut normalized = url.clone(); let host_for_url = host .as_str() .parse::<Ipv6Addr>() .map(|_| format!("[{}]", host.as_str())) .unwrap_or_else(|_| host.as_str().to_owned()); normalized .set_host(Some(&host_for_url)) .map_err(|_| EndpointError::Host(HostError::InvalidDnsName))?; Ok(Self { url: normalized }) }
pub(crate) fn url(&self) -> &Url { &self.url }
#[cfg(feature = "indexer")] pub(crate) fn with_path(&self, path: &str) -> Self { let mut endpoint = self.clone(); endpoint.url.set_path(path); endpoint }
#[cfg(feature = "indexer")] pub(crate) fn append_query_pair(&mut self, key: &str, value: &str) { self.url.query_pairs_mut().append_pair(key, value); }}
/// The system resolver used by both public HTTP and websocket connectors.////// Every lookup validates the complete answer set before returning any address./// The injected lookup makes the policy independently testable without live DNS.#[derive(Clone)]pub(crate) struct PublicResolver { lookup: Arc<dyn Lookup>,}
trait Lookup: Send + Sync { fn lookup( &self, host: String, port: u16, ) -> Pin<Box<dyn Future<Output = Result<Vec<SocketAddr>, String>> + Send>>;}
struct SystemLookup;
impl Lookup for SystemLookup { fn lookup( &self, host: String, port: u16, ) -> Pin<Box<dyn Future<Output = Result<Vec<SocketAddr>, String>> + Send>> { Box::pin(async move { tokio::net::lookup_host((host, port)) .await .map(|addresses| addresses.collect()) .map_err(|error| error.to_string()) }) }}
impl PublicResolver { pub(crate) fn system() -> Self { Self { lookup: Arc::new(SystemLookup), } }
async fn resolve_addresses( &self, host: &str, port: u16, ) -> Result<Vec<SocketAddr>, ResolveError> { let addresses = tokio::time::timeout( DNS_LOOKUP_TIMEOUT, self.lookup.lookup(host.to_owned(), port), ) .await .map_err(|_| ResolveError::Timeout)? .map_err(ResolveError::Lookup)?; validate_addresses(addresses) }}
#[derive(Debug, Error)]enum ResolveError { #[error("DNS lookup timed out")] Timeout, #[error("DNS lookup failed: {0}")] Lookup(String), #[error("DNS lookup returned no addresses")] Empty, #[error("DNS answer contains non-public address: {0}")] NonPublic(IpAddr),}
fn validate_addresses(addresses: Vec<SocketAddr>) -> Result<Vec<SocketAddr>, ResolveError> { if addresses.is_empty() { return Err(ResolveError::Empty); } if let Some(address) = addresses.iter().find(|address| is_private_ip(address.ip())) { return Err(ResolveError::NonPublic(address.ip())); } Ok(addresses)}
impl Resolve for PublicResolver { fn resolve(&self, name: Name) -> Resolving { let resolver = self.clone(); let host = name.as_str().to_owned(); Box::pin(async move { let addresses = resolver.resolve_addresses(&host, 0).await?; Ok(Box::new(addresses.into_iter()) as Addrs) }) }}
impl tokio_websockets::resolver::Resolver for PublicResolver { async fn resolve(&self, host: &str, port: u16) -> Result<SocketAddr, tokio_websockets::Error> { self.resolve_addresses(host, port) .await .map_err(|_| tokio_websockets::Error::CannotResolveHost)? .into_iter() .next() .ok_or(tokio_websockets::Error::CannotResolveHost) }}
const MAX_PUBLIC_HTTP_BODY_BYTES: usize = 2 * 1024 * 1024;
/// The only public HTTP request capability. Its raw reqwest client is private;/// callers must provide a [`PublicEndpoint`] and cannot send an unchecked URL.#[derive(Clone)]pub(crate) struct PublicHttpClient { client: Arc<reqwest::Client>,}
#[derive(Debug, Error)]pub enum PublicHttpError { #[error(transparent)] UnsafeEndpoint(#[from] EndpointError), #[error(transparent)] Reqwest(#[from] reqwest::Error), #[error("public HTTP response body exceeds {limit} bytes")] BodyTooLarge { limit: usize },}
impl PublicHttpClient { #[cfg(any(feature = "indexer", test))] pub(crate) fn get(&self, endpoint: &PublicEndpoint) -> reqwest::RequestBuilder { self.client.get(endpoint.url().clone()) }
/// Execute a streaming request after checking its final URL. The wrapped /// client owns the socket-coupled resolver, so callers cannot substitute a /// raw transport after obtaining endpoint syntax proof. #[cfg(feature = "indexer")] pub(crate) async fn execute( &self, request: reqwest::Request, ) -> Result<reqwest::Response, PublicHttpError> { PublicEndpoint::parse_http(request.url()).map_err(PublicHttpError::UnsafeEndpoint)?; self.client.execute(request).await.map_err(Into::into) }
/// Send a Jacquard request only after proving its complete public endpoint /// syntax. The response is consumed through a bounded decoded stream. pub(crate) async fn send_http_bounded( &self, request: http::Request<Vec<u8>>, max_body_bytes: usize, ) -> Result<http::Response<Vec<u8>>, PublicHttpError> { send_public_http(&self.client, request, max_body_bytes).await }
pub(crate) async fn send_http( &self, request: http::Request<Vec<u8>>, ) -> Result<http::Response<Vec<u8>>, PublicHttpError> { self.send_http_bounded(request, MAX_PUBLIC_HTTP_BODY_BYTES) .await }
#[cfg(test)] pub(crate) fn test(client: reqwest::Client) -> Self { Self { client: Arc::new(client), } }}
impl jacquard_common::http_client::HttpClient for PublicHttpClient { type Error = PublicHttpError;
async fn send_http( &self, request: http::Request<Vec<u8>>, ) -> Result<http::Response<Vec<u8>>, Self::Error> { PublicHttpClient::send_http(self, request).await }}
/// Public transport for legacy XRPC adapters. Unlike reqwest's DNS hook alone,/// this validates the request URI first because literal IPs bypass DNS. The/// response is consumed through a bounded decoded stream.async fn send_public_http( client: &reqwest::Client, request: http::Request<Vec<u8>>, max_body_bytes: usize,) -> Result<http::Response<Vec<u8>>, PublicHttpError> { let url = Url::parse(&request.uri().to_string()).map_err(|_| { PublicHttpError::UnsafeEndpoint(EndpointError::Host(HostError::InvalidDnsName)) })?; PublicEndpoint::parse_http(&url).map_err(PublicHttpError::UnsafeEndpoint)?; send_reqwest(client, request, max_body_bytes).await}
async fn send_reqwest( client: &reqwest::Client, request: http::Request<Vec<u8>>, max_body_bytes: usize,) -> Result<http::Response<Vec<u8>>, PublicHttpError> { let (parts, body) = request.into_parts(); let mut request = client .request(parts.method, parts.uri.to_string()) .body(body); for (name, value) in &parts.headers { request = request.header(name.as_str(), value.as_bytes()); } let response = request.send().await?; let status = response.status(); let headers = response.headers().clone(); let mut stream = response.bytes_stream(); let mut body = bytes::BytesMut::new(); while let Some(chunk) = stream.next().await { let chunk = chunk?; if body.len() + chunk.len() > max_body_bytes { return Err(PublicHttpError::BodyTooLarge { limit: max_body_bytes, }); } body.extend_from_slice(&chunk); } let mut output = http::Response::builder().status(status); for (name, value) in &headers { output = output.header(name, value); } Ok(output .body(body.freeze().to_vec()) .expect("response status and headers are valid"))}
/// Build the public transport used by legacy XRPC adapters that cannot carry/// [`PublicEndpoint`] through their trait surface. Callers must validate the/// endpoint before handing requests to this client; direct CAR callers use the/// typed `ThrottledHttpClient::get` wrapper.fn public_http_transport(timeout: Duration, connect_timeout: Duration) -> reqwest::Client { let connect_timeout = connect_timeout.min(timeout); reqwest::Client::builder() .user_agent(concat!( env!("CARGO_PKG_NAME"), "/", env!("CARGO_PKG_VERSION") )) .dns_resolver(Arc::new(PublicResolver::system())) .no_proxy() .redirect(Policy::none()) .timeout(timeout) .connect_timeout(connect_timeout) .gzip(true) .brotli(true) .zstd(true) .build() .expect("public HTTP client configuration should be valid")}
pub(crate) fn public_http_with_timeout(timeout: Duration) -> PublicHttpClient { PublicHttpClient { client: Arc::new(public_http_transport(timeout, PUBLIC_CONNECT_TIMEOUT)), }}
/// the same public egress validation as [`public_http`], on the handle budget.pub(crate) fn handle_http() -> PublicHttpClient { static CLIENT: OnceLock<PublicHttpClient> = OnceLock::new(); CLIENT .get_or_init(|| PublicHttpClient { client: Arc::new(public_http_transport( HANDLE_REQUEST_TIMEOUT, HANDLE_CONNECT_TIMEOUT, )), }) .clone()}
pub(crate) fn public_http() -> PublicHttpClient { static CLIENT: OnceLock<PublicHttpClient> = OnceLock::new(); CLIENT .get_or_init(|| public_http_with_timeout(PUBLIC_REQUEST_TIMEOUT)) .clone()}
pub(crate) fn trusted_http() -> &'static reqwest::Client { static CLIENT: OnceLock<reqwest::Client> = OnceLock::new(); CLIENT.get_or_init(|| { reqwest::Client::builder() .user_agent(concat!( env!("CARGO_PKG_NAME"), "/", env!("CARGO_PKG_VERSION") )) .gzip(true) .brotli(true) .zstd(true) // Trust is scoped to the configured origin, not arbitrary redirect // targets selected by that service. .redirect(Policy::none()) .build() .expect("trusted HTTP client configuration should be valid") })}
/// Svix's conservative public-address classifier. Keep the canonicalization/// first: an IPv4-mapped IPv6 literal is an IPv4 address for policy purposes.fn is_allowed(addr: IpAddr) -> bool { match addr.to_canonical() { IpAddr::V4(addr) => { !(addr.is_private() || addr.is_loopback() || addr.is_link_local() || addr.is_broadcast() || addr.is_multicast() || addr.is_documentation() || is_shared(addr) || is_reserved(addr) || is_benchmarking(addr) || starts_with_zero(addr)) } IpAddr::V6(addr) => { !(addr.is_multicast() || addr.is_loopback() || addr.is_unspecified() || is_unicast_link_local(addr) || is_site_local(addr) || is_unique_local(addr) || is_documentation_v6(addr) || is_discard_v6(addr) || is_iana_non_global_v6(addr) || is_6to4_or_nat64_v6(addr) || is_srv6(addr) || is_transition_or_tunnel_v6(addr)) } }}
/// True for private, reserved, or special-use addresses that must never be/// reached through the public egress capability.pub(crate) fn is_private_ip(ip: IpAddr) -> bool { !is_allowed(ip)}
#[inline(always)]fn is_shared(addr: Ipv4Addr) -> bool { static CGNAT: Ipv4Net = Ipv4Net::new_assert(Ipv4Addr::new(100, 64, 0, 0), 10); CGNAT.contains(&addr)}
#[inline(always)]fn is_reserved(addr: Ipv4Addr) -> bool { static IETF_RESERVED_1: Ipv4Net = Ipv4Net::new_assert(Ipv4Addr::new(192, 0, 0, 0), 24); static IETF_RESERVED_2: Ipv4Net = Ipv4Net::new_assert(Ipv4Addr::new(192, 88, 99, 0), 24); static CLASS_E: Ipv4Net = Ipv4Net::new_assert(Ipv4Addr::new(240, 0, 0, 0), 4);
// revisit allowing class E addresses if the deployment environment ever // makes that protocol space routable. IETF_RESERVED_1.contains(&addr) || IETF_RESERVED_2.contains(&addr) || CLASS_E.contains(&addr)}
#[inline(always)]fn is_benchmarking(addr: Ipv4Addr) -> bool { static BENCHMARKING: Ipv4Net = Ipv4Net::new_assert(Ipv4Addr::new(198, 18, 0, 0), 15); BENCHMARKING.contains(&addr)}
#[inline(always)]fn starts_with_zero(addr: Ipv4Addr) -> bool { addr.octets()[0] == 0}
#[inline(always)]fn is_unicast_link_local(addr: Ipv6Addr) -> bool { // technically only fe80::/64 are link-local, but fe80::/10 are all reserved static V6_LL: Ipv6Net = Ipv6Net::new_assert(Ipv6Addr::new(0xfe80, 0, 0, 0, 0, 0, 0, 0), 10); V6_LL.contains(&addr)}
#[inline(always)]fn is_site_local(addr: Ipv6Addr) -> bool { // Deprecated site-local space remains unroutable and is retained as a // separate rule so future link-local edits cannot accidentally reopen it. static SITE_LOCAL: Ipv6Net = Ipv6Net::new_assert(Ipv6Addr::new(0xfec0, 0, 0, 0, 0, 0, 0, 0), 10); SITE_LOCAL.contains(&addr)}
#[inline(always)]fn is_unique_local(addr: Ipv6Addr) -> bool { static ULA: Ipv6Net = Ipv6Net::new_assert(Ipv6Addr::new(0xfc00, 0, 0, 0, 0, 0, 0, 0), 7); ULA.contains(&addr)}
#[inline(always)]fn is_documentation_v6(addr: Ipv6Addr) -> bool { static V6_DOCUMENTATION_1: Ipv6Net = Ipv6Net::new_assert(Ipv6Addr::new(0x2001, 0xdb8, 0, 0, 0, 0, 0, 0), 32); static V6_DOCUMENTATION_2: Ipv6Net = Ipv6Net::new_assert(Ipv6Addr::new(0x3fff, 0, 0, 0, 0, 0, 0, 0), 20); V6_DOCUMENTATION_1.contains(&addr) || V6_DOCUMENTATION_2.contains(&addr)}
#[inline(always)]fn is_discard_v6(addr: Ipv6Addr) -> bool { static V6_DISCARD: Ipv6Net = Ipv6Net::new_assert(Ipv6Addr::new(0x100, 0, 0, 0, 0, 0, 0, 0), 64); V6_DISCARD.contains(&addr)}
#[inline(always)]fn is_iana_non_global_v6(addr: Ipv6Addr) -> bool { // IANA special-purpose allocations: benchmarking and ORCHID/ORCHIDv2. static BENCHMARKING: Ipv6Net = Ipv6Net::new_assert(Ipv6Addr::new(0x2001, 0x2, 0, 0, 0, 0, 0, 0), 48); static ORCHID: Ipv6Net = Ipv6Net::new_assert(Ipv6Addr::new(0x2001, 0x10, 0, 0, 0, 0, 0, 0), 28); static ORCHID_V2: Ipv6Net = Ipv6Net::new_assert(Ipv6Addr::new(0x2001, 0x20, 0, 0, 0, 0, 0, 0), 28); BENCHMARKING.contains(&addr) || ORCHID.contains(&addr) || ORCHID_V2.contains(&addr)}
#[inline(always)]fn is_6to4_or_nat64_v6(addr: Ipv6Addr) -> bool { static IP6TO4: Ipv6Net = Ipv6Net::new_assert(Ipv6Addr::new(0x2002, 0, 0, 0, 0, 0, 0, 0), 16); static NAT64: Ipv6Net = Ipv6Net::new_assert(Ipv6Addr::new(0x64, 0xff9b, 0, 0, 0, 0, 0, 0), 96); static LOCAL_NAT64: Ipv6Net = Ipv6Net::new_assert(Ipv6Addr::new(0x64, 0xff9b, 0x1, 0, 0, 0, 0, 0), 48); IP6TO4.contains(&addr) || NAT64.contains(&addr) || LOCAL_NAT64.contains(&addr)}
#[inline(always)]fn is_srv6(addr: Ipv6Addr) -> bool { static SRV6: Ipv6Net = Ipv6Net::new_assert(Ipv6Addr::new(0x5f00, 0, 0, 0, 0, 0, 0, 0), 16); SRV6.contains(&addr)}
/// Hydrant defense-in-depth extensions beyond the copied Svix baseline./// Deprecated compatibility/translation forms and Teredo can preserve or hide/// IPv4 routing semantics, so deny each whole prefix rather than decoding it.#[inline(always)]fn is_transition_or_tunnel_v6(addr: Ipv6Addr) -> bool { static IPV4_COMPATIBLE: Ipv6Net = Ipv6Net::new_assert(Ipv6Addr::new(0, 0, 0, 0, 0, 0, 0, 0), 96); static IPV4_TRANSLATED: Ipv6Net = Ipv6Net::new_assert(Ipv6Addr::new(0, 0, 0, 0, 0xffff, 0, 0, 0), 96); static TEREDO: Ipv6Net = Ipv6Net::new_assert(Ipv6Addr::new(0x2001, 0, 0, 0, 0, 0, 0, 0), 32); IPV4_COMPATIBLE.contains(&addr) || IPV4_TRANSLATED.contains(&addr) || TEREDO.contains(&addr)}
fn is_dns_name(host: &str) -> bool { host.len() <= 253 && host.split('.').all(|label| { !label.is_empty() && label.len() <= 63 && !label.starts_with('-') && !label.ends_with('-') && label .bytes() .all(|byte| byte.is_ascii_alphanumeric() || byte == b'-') })}
#[cfg(test)]mod tests { use super::*; use axum::{Router, body::Body, http::StatusCode, response::IntoResponse, routing::get}; use std::sync::Mutex; use tokio_websockets::{ClientBuilder, ServerBuilder};
#[test] fn the_handle_budget_is_tighter_than_the_public_one() { assert!(HANDLE_REQUEST_TIMEOUT < PUBLIC_REQUEST_TIMEOUT); assert!(HANDLE_CONNECT_TIMEOUT < HANDLE_REQUEST_TIMEOUT); assert!(PUBLIC_CONNECT_TIMEOUT < PUBLIC_REQUEST_TIMEOUT); }
struct ScriptedLookup(Mutex<Vec<Vec<SocketAddr>>>);
impl Lookup for ScriptedLookup { fn lookup( &self, _host: String, _port: u16, ) -> Pin<Box<dyn Future<Output = Result<Vec<SocketAddr>, String>> + Send>> { let result = self.0.lock().unwrap().pop().unwrap_or_default(); Box::pin(std::future::ready(Ok(result))) } }
fn resolver(answer: Vec<SocketAddr>) -> PublicResolver { PublicResolver { lookup: Arc::new(ScriptedLookup(Mutex::new(vec![answer]))), } }
fn addr(ip: &str) -> SocketAddr { SocketAddr::new(ip.parse().unwrap(), 443) }
#[tokio::test] async fn mixed_answers_are_rejected_as_a_set() { let resolver = resolver(vec![addr("93.184.216.34"), addr("127.0.0.1")]); assert!(matches!( resolver.resolve_addresses("pds.example", 443).await, Err(ResolveError::NonPublic(_)) )); }
#[tokio::test] async fn websocket_public_resolver_rejects_private_before_tcp() { let resolver = resolver(vec![addr("127.0.0.1")]); let uri = "ws://public.example/".parse().unwrap(); let result = ClientBuilder::from_uri(uri) .resolver(resolver) .connect() .await; assert!(matches!( result, Err(tokio_websockets::Error::CannotResolveHost) )); }
#[derive(Clone)] struct FixedSocketResolver(SocketAddr);
impl tokio_websockets::resolver::Resolver for FixedSocketResolver { async fn resolve( &self, _host: &str, _port: u16, ) -> Result<SocketAddr, tokio_websockets::Error> { Ok(self.0) } }
#[tokio::test] async fn websocket_handshake_retains_original_uri_host() { let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); let address = listener.local_addr().unwrap(); let server = tokio::spawn(async move { let (tcp, _) = listener.accept().await.unwrap(); let (request, _ws) = ServerBuilder::new().accept(tcp).await.unwrap(); // tokio-websockets exposes the request-target URI as a path, while // the authority is carried by Host. This is the value TLS/SNI and // HTTP routing must preserve when the resolver supplies a socket. assert_eq!( request.headers().get(http::header::HOST).unwrap(), "public.example" ); });
let uri = "ws://public.example/handshake".parse().unwrap(); let (_ws, _response) = ClientBuilder::from_uri(uri) .resolver(FixedSocketResolver(address)) .connect() .await .unwrap(); server.await.unwrap(); }
#[tokio::test] async fn empty_answers_are_rejected() { assert!(matches!( resolver(Vec::new()) .resolve_addresses("pds.example", 443) .await, Err(ResolveError::Empty) )); }
#[tokio::test] async fn mixed_ipv4_and_ipv6_answers_are_rejected_as_a_set() { let resolver = resolver(vec![addr("93.184.216.34"), addr("::ffff:127.0.0.1")]); assert!(matches!( resolver.resolve_addresses("pds.example", 443).await, Err(ResolveError::NonPublic(_)) )); }
#[test] fn svix_is_allowed_regressions() { // RFC 1918 for address in [ "10.254.0.0", "192.168.10.65", "172.16.10.65", // zero-prefix (including unspecified) "0.1.2.3", "0.0.0.0", // loopback and link-local "127.0.0.1", "127.254.254.254", "169.254.45.1", // broadcast and multicast "255.255.255.255", "224.1.2.3", "233.252.0.16", // protocol-reserved, docs, CGNAT, benchmarking, class E "192.0.0.0", "192.0.0.255", "192.0.2.255", "198.51.100.65", "203.0.113.6", "100.100.0.0", "192.88.99.3", "198.18.0.0", "240.240.240.240", ] { assert!(!is_allowed(address.parse().unwrap()), "accepted {address}"); } assert!(is_allowed("1.1.1.1".parse().unwrap())); assert!(is_allowed("8.8.8.8".parse().unwrap()));
for address in [ "::0", "::1", "2001:db8::1", "2001:2::dead:beef", "2001:10::dead:beef", "2001:20::dead:beef", "3fff::dead:beef", "fe80::dead:beef", "febf:ffff::dead:beef", "fec0::dead:beef", "feff:ffff::dead:beef", "ff00::dead:beef", "fdee:3a45:9057::1", "fd7a:115c:a1e0::1", // all 6to4 and both reserved NAT64 prefixes are denied. "2002:c000:0204::dead:beef", "2002:c0a8:0001::dead:beef", "64:ff9b::808:808", "64:ff9b::c0a8:0", "64:ff9b:1::dead:beef", // Hydrant defense-in-depth: deprecated compatibility/translation // and Teredo prefixes are denied without embedded-address decoding. "::8.8.8.8", "::ffff:0:8.8.8.8", "2001:0:ffff:ffff:ffff:ffff:ffff:ffff", "5f00::dead:beef", "100::1:2:3:4", ] { assert!(!is_allowed(address.parse().unwrap()), "accepted {address}"); } assert!(is_allowed("2001:4860:4860::8888".parse().unwrap())); assert!(is_allowed("::ffff:8.8.8.8".parse().unwrap())); assert!(!is_allowed("::ffff:127.0.0.1".parse().unwrap())); }
#[test] fn public_endpoint_normalizes_host_and_rejects_unsafe_authority() { let endpoint = PublicEndpoint::parse_http(&Url::parse("https://BÜCHER.Example./xrpc?x=1").unwrap()) .unwrap(); assert_eq!(endpoint.url().host_str(), Some("xn--bcher-kva.example")); assert_eq!( endpoint.url().as_str(), "https://xn--bcher-kva.example/xrpc?x=1" ); assert!(PublicEndpoint::parse_http(&Url::parse("https://user@example/").unwrap()).is_err()); assert!( PublicEndpoint::parse_http(&Url::parse("https://example/#fragment").unwrap()).is_err() ); assert!(PublicEndpoint::parse_http(&Url::parse("https://127.0.0.1/").unwrap()).is_err()); assert!( PublicEndpoint::parse_http(&Url::parse("https://[::ffff:127.0.0.1]/").unwrap()) .is_err() ); assert!( PublicEndpoint::parse_http(&Url::parse("https://[::ffff:8.8.8.8]/").unwrap()).is_ok() ); assert!(PublicEndpoint::parse_http(&Url::parse("ftp://example/").unwrap()).is_err()); assert!(PublicEndpoint::parse_websocket(&Url::parse("ws://example/").unwrap()).is_err()); assert!(PublicEndpoint::parse_websocket(&Url::parse("wss://example/").unwrap()).is_ok()); assert!(PublicHost::parse("local").is_err()); assert!(PublicHost::parse("example:443").is_err()); for disguised_loopback in ["2130706433", "017700000001", "0x7f000001"] { assert!( PublicHost::parse(disguised_loopback).is_err(), "accepted disguised loopback {disguised_loopback}" ); } }
#[tokio::test] async fn public_client_no_proxy_disables_configured_proxy() { let app = Router::new().route("/", get(|| async { "direct" })); let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); let address = listener.local_addr().unwrap(); let server = tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); let client = reqwest::Client::builder() // If no_proxy is removed, this dead local proxy makes the request fail. .proxy(reqwest::Proxy::all("http://127.0.0.1:9").unwrap()) .no_proxy() .redirect(Policy::none()) .resolve("public.example", address) .build() .unwrap(); let endpoint = PublicEndpoint::parse_http(&Url::parse("http://public.example/").unwrap()).unwrap(); let response = PublicHttpClient::test(client) .get(&endpoint) .send() .await .unwrap(); assert_eq!(response.text().await.unwrap(), "direct"); server.abort(); }
#[tokio::test] async fn public_http_does_not_follow_redirects_to_private_literals() { let app = Router::new().route( "/", get(|| async { ( StatusCode::FOUND, [(http::header::LOCATION, "http://127.0.0.1:9/private")], ) .into_response() }), ); let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); let address = listener.local_addr().unwrap(); let server = tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); let client = reqwest::Client::builder() .proxy(reqwest::Proxy::all("http://127.0.0.1:9").unwrap()) .no_proxy() .redirect(Policy::none()) .resolve("public.example", address) .build() .unwrap(); let endpoint = PublicEndpoint::parse_http(&Url::parse("http://public.example/").unwrap()).unwrap(); let response = PublicHttpClient::test(client) .get(&endpoint) .send() .await .unwrap(); assert_eq!(response.status(), StatusCode::FOUND); server.abort(); }
#[tokio::test] async fn public_jacquard_response_body_is_capped_before_materializing() { let app = Router::new().route( "/", get(|| async { Body::from(vec![b'x'; MAX_PUBLIC_HTTP_BODY_BYTES + 1]).into_response() }), ); let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); let address = listener.local_addr().unwrap(); let server = tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); let client = reqwest::Client::builder() .no_proxy() .resolve("public.example", address) .build() .unwrap(); let request = http::Request::builder() .uri("http://public.example/") .body(Vec::new()) .unwrap(); let result = PublicHttpClient::test(client).send_http(request).await; assert!(matches!(result, Err(PublicHttpError::BodyTooLarge { .. }))); server.abort(); }}