From 66cff86b172fb2da758117b0f531705f40d74a24 Mon Sep 17 00:00:00 2001 From: eareimu Date: Wed, 15 Jul 2026 16:10:20 +0800 Subject: [PATCH 1/7] feat: add typed endpoint lookup request --- src/h3/lookup.rs | 1 + src/http.rs | 1 + src/mdns.rs | 1 + src/resolvers.rs | 15 ++++-- src/resolvers/deferred.rs | 18 +++++-- src/resolvers/endpoint_candidates.rs | 80 ++++++++++++++++++++++++++-- 6 files changed, 106 insertions(+), 10 deletions(-) diff --git a/src/h3/lookup.rs b/src/h3/lookup.rs index 8dfd84e..70b61f2 100644 --- a/src/h3/lookup.rs +++ b/src/h3/lookup.rs @@ -281,6 +281,7 @@ where fn lookup_endpoint_candidates<'a>( &'a self, name: &'a str, + _lookup: crate::resolvers::endpoint_candidates::EndpointLookup, ) -> crate::resolvers::endpoint_candidates::EndpointCandidateFuture<'a> { Box::pin(async move { let Some((domain, _sequence)) = diff --git a/src/http.rs b/src/http.rs index 0b13ec3..19f7851 100644 --- a/src/http.rs +++ b/src/http.rs @@ -397,6 +397,7 @@ impl crate::resolvers::endpoint_candidates::ResolveEndpointCandidates for HttpRe fn lookup_endpoint_candidates<'a>( &'a self, name: &'a str, + _lookup: crate::resolvers::endpoint_candidates::EndpointLookup, ) -> crate::resolvers::endpoint_candidates::EndpointCandidateFuture<'a> { let lookup = async move { let Some((domain, _sequence)) = diff --git a/src/mdns.rs b/src/mdns.rs index 4642031..d75c134 100644 --- a/src/mdns.rs +++ b/src/mdns.rs @@ -361,6 +361,7 @@ impl crate::resolvers::endpoint_candidates::ResolveEndpointCandidates for MdnsRe fn lookup_endpoint_candidates<'a>( &'a self, name: &'a str, + _lookup: crate::resolvers::endpoint_candidates::EndpointLookup, ) -> crate::resolvers::endpoint_candidates::EndpointCandidateFuture<'a> { Box::pin(async move { let Some((domain, _sequence)) = diff --git a/src/resolvers.rs b/src/resolvers.rs index 7aca24d..a9318c8 100644 --- a/src/resolvers.rs +++ b/src/resolvers.rs @@ -356,6 +356,7 @@ impl Resolvers { pub async fn lookup_endpoint_candidates( &self, name: &str, + lookup: crate::resolvers::endpoint_candidates::EndpointLookup, ) -> Result { let mut errors = vec![]; let mut groups = Vec::new(); @@ -369,7 +370,10 @@ impl Resolvers { continue; }; - match candidate_resolver.lookup_endpoint_candidates(name).await { + match candidate_resolver + .lookup_endpoint_candidates(name, lookup) + .await + { Ok(candidates) => groups.extend(candidates.groups), Err(error) => errors.push((entry.resolver.to_string(), error)), } @@ -413,9 +417,10 @@ impl crate::resolvers::endpoint_candidates::ResolveEndpointCandidates for Resolv fn lookup_endpoint_candidates<'a>( &'a self, name: &'a str, + lookup: crate::resolvers::endpoint_candidates::EndpointLookup, ) -> crate::resolvers::endpoint_candidates::EndpointCandidateFuture<'a> { async move { - Resolvers::lookup_endpoint_candidates(self, name) + Resolvers::lookup_endpoint_candidates(self, name, lookup) .await .map_err(io::Error::other) } @@ -589,6 +594,7 @@ mod tests { fn lookup_endpoint_candidates<'a>( &'a self, _name: &'a str, + _lookup: crate::resolvers::endpoint_candidates::EndpointLookup, ) -> crate::resolvers::endpoint_candidates::EndpointCandidateFuture<'a> { use dhttp_identity::certificate::{ CertificateChainKey, CertificateChainKind, CertificateSequence, @@ -629,7 +635,10 @@ mod tests { })); let candidates = resolvers - .lookup_endpoint_candidates("demo.dhttp.net") + .lookup_endpoint_candidates( + "demo.dhttp.net", + crate::resolvers::endpoint_candidates::EndpointLookup::default(), + ) .await .expect("candidate lookup succeeds"); diff --git a/src/resolvers/deferred.rs b/src/resolvers/deferred.rs index bd70502..4b218c5 100644 --- a/src/resolvers/deferred.rs +++ b/src/resolvers/deferred.rs @@ -5,7 +5,9 @@ use futures::{FutureExt, future::BoxFuture}; use snafu::{ResultExt, Snafu}; use tokio::sync::{Notify, OnceCell}; -use crate::resolvers::endpoint_candidates::{EndpointCandidateFuture, ResolveEndpointCandidates}; +use crate::resolvers::endpoint_candidates::{ + EndpointCandidateFuture, EndpointLookup, ResolveEndpointCandidates, +}; #[derive(Debug, Snafu)] #[snafu(module, visibility(pub))] @@ -124,8 +126,18 @@ impl ResolveEndpointCandidates for DeferredResolver where R: ResolveEndpointCandidates + 'static, { - fn lookup_endpoint_candidates<'a>(&'a self, name: &'a str) -> EndpointCandidateFuture<'a> { - async move { self.wait().await.lookup_endpoint_candidates(name).await }.boxed() + fn lookup_endpoint_candidates<'a>( + &'a self, + name: &'a str, + lookup: EndpointLookup, + ) -> EndpointCandidateFuture<'a> { + async move { + self.wait() + .await + .lookup_endpoint_candidates(name, lookup) + .await + } + .boxed() } } diff --git a/src/resolvers/endpoint_candidates.rs b/src/resolvers/endpoint_candidates.rs index cb40df5..4ba0d50 100644 --- a/src/resolvers/endpoint_candidates.rs +++ b/src/resolvers/endpoint_candidates.rs @@ -1,6 +1,6 @@ -use std::io; +use std::{io, num::NonZeroUsize}; -use dhttp_identity::certificate::CertificateChainKey; +use dhttp_identity::certificate::{CertificateChainKey, CertificateSequence}; use dquic::{ qbase::net::addr::EndpointAddr as DquicEndpointAddr, qresolve::{Resolve, Source}, @@ -21,10 +21,61 @@ pub struct EndpointCandidates { pub groups: Vec, } +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub enum SequenceQuery { + #[default] + Default, + Exact(CertificateSequence), + Limit(NonZeroUsize), + All, +} + +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +pub struct EndpointLookup { + pub sequences: SequenceQuery, + pub record_limit: Option, +} + +impl EndpointLookup { + #[must_use] + pub fn exact(sequence: CertificateSequence) -> Self { + Self { + sequences: SequenceQuery::Exact(sequence), + record_limit: None, + } + } + + #[must_use] + pub fn limit(count: NonZeroUsize) -> Self { + Self { + sequences: SequenceQuery::Limit(count), + record_limit: None, + } + } + + #[must_use] + pub fn all() -> Self { + Self { + sequences: SequenceQuery::All, + record_limit: None, + } + } + + #[must_use] + pub fn with_record_limit(mut self, count: NonZeroUsize) -> Self { + self.record_limit = Some(count); + self + } +} + pub type EndpointCandidateFuture<'a> = BoxFuture<'a, io::Result>; pub trait ResolveEndpointCandidates: Resolve { - fn lookup_endpoint_candidates<'a>(&'a self, name: &'a str) -> EndpointCandidateFuture<'a>; + fn lookup_endpoint_candidates<'a>( + &'a self, + name: &'a str, + lookup: EndpointLookup, + ) -> EndpointCandidateFuture<'a>; } pub type ArcEndpointCandidateResolver = @@ -95,12 +146,33 @@ fn effective_chain_key( #[cfg(test)] mod tests { - use std::net::SocketAddrV4; + use std::{net::SocketAddrV4, num::NonZeroUsize}; use dhttp_identity::certificate::{CertificateChainKind, CertificateSequence}; use super::*; + #[test] + fn endpoint_lookup_constructors_encode_valid_states() { + let one = NonZeroUsize::new(1).unwrap(); + let exact = CertificateSequence::from(2u8); + + assert_eq!(EndpointLookup::default().sequences, SequenceQuery::Default); + assert_eq!( + EndpointLookup::exact(exact).sequences, + SequenceQuery::Exact(exact) + ); + assert_eq!( + EndpointLookup::limit(one).sequences, + SequenceQuery::Limit(one) + ); + assert_eq!(EndpointLookup::all().sequences, SequenceQuery::All); + assert_eq!( + EndpointLookup::all().with_record_limit(one).record_limit, + Some(one) + ); + } + fn direct(addr: &str, main: bool, sequence: u32) -> DnsEndpointAddr { let socket: SocketAddrV4 = addr.parse().expect("socket addr"); let mut endpoint = DnsEndpointAddr::direct_v4(socket); From 10bd33e9c8376bafadd9c636e5f61c86a5fd3ec1 Mon Sep 17 00:00:00 2001 From: eareimu Date: Wed, 15 Jul 2026 16:13:27 +0800 Subject: [PATCH 2/7] feat: query ordered primary sequences --- src/h3/lookup.rs | 76 +++++++++++++++++++++------- src/http.rs | 73 +++++++++++++++++++++----- src/resolvers/endpoint_candidates.rs | 44 +++++++++++++++- src/resolvers/endpoint_group.rs | 22 ++++---- 4 files changed, 170 insertions(+), 45 deletions(-) diff --git a/src/h3/lookup.rs b/src/h3/lookup.rs index 70b61f2..553d8cc 100644 --- a/src/h3/lookup.rs +++ b/src/h3/lookup.rs @@ -21,16 +21,16 @@ use crate::{ const LOOKUP_API_PATH: &str = "/api/v2/lookup"; -fn lookup_url(base_url: &url::Url, name: &str, sequence: Option) -> url::Url { +fn lookup_url( + base_url: &url::Url, + name: &str, + lookup: crate::resolvers::endpoint_candidates::EndpointLookup, +) -> url::Url { let mut url = base_url .join(LOOKUP_API_PATH) .expect("h3 dns lookup api path must be valid"); url.query_pairs_mut().append_pair("host", name); - if let Some(sequence) = sequence { - let sequence_text = sequence.get().to_string(); - url.query_pairs_mut() - .append_pair("sequence", &sequence_text); - } + crate::resolvers::endpoint_candidates::append_endpoint_lookup_query(&mut url, lookup); url } @@ -244,7 +244,10 @@ where return Ok(stream.boxed()); } - let url = lookup_url(&self.base_url, domain, sequence); + let endpoint_lookup = sequence + .map(crate::resolvers::endpoint_candidates::EndpointLookup::exact) + .unwrap_or_default(); + let url = lookup_url(&self.base_url, domain, endpoint_lookup); let uri: http::Uri = url.as_str().parse().expect("URL should be valid URI"); tracing::trace!("sending lookup request to {}", self.base_url); @@ -281,16 +284,19 @@ where fn lookup_endpoint_candidates<'a>( &'a self, name: &'a str, - _lookup: crate::resolvers::endpoint_candidates::EndpointLookup, + lookup: crate::resolvers::endpoint_candidates::EndpointLookup, ) -> crate::resolvers::endpoint_candidates::EndpointCandidateFuture<'a> { Box::pin(async move { - let Some((domain, _sequence)) = + let Some((domain, sequence)) = crate::resolvers::endpoint_lookup_name_and_sequence(name) else { return Err(io::Error::other("no DNS record found")); }; - let url = lookup_url(&self.base_url, domain, None); + let lookup = sequence + .map(crate::resolvers::endpoint_candidates::EndpointLookup::exact) + .unwrap_or(lookup); + let url = lookup_url(&self.base_url, domain, lookup); let uri: http::Uri = url.as_str().parse().expect("URL should be valid URI"); let response = self .lookup_response_with_retry(uri) @@ -319,13 +325,16 @@ where #[cfg(test)] mod tests { - use std::{collections::HashMap, net::SocketAddrV4}; + use std::{collections::HashMap, net::SocketAddrV4, num::NonZeroUsize}; use super::*; - use crate::core::{ - MdnsPacket, - parser::record::endpoint::EndpointAddr as DnsEndpointAddr, - wire::{MultiResponse, ResponseRecord}, + use crate::{ + core::{ + MdnsPacket, + parser::record::endpoint::EndpointAddr as DnsEndpointAddr, + wire::{MultiResponse, ResponseRecord}, + }, + resolvers::endpoint_candidates::EndpointLookup, }; fn direct(addr: &str, main: bool, sequence: u32) -> DnsEndpointAddr { @@ -369,7 +378,7 @@ mod tests { #[test] fn h3_lookup_url_targets_v2_api_from_origin_base() { let base_url = url::Url::parse("https://dns.example.test:4433").expect("url"); - let url = lookup_url(&base_url, "demo.dhttp.net", None); + let url = lookup_url(&base_url, "demo.dhttp.net", EndpointLookup::default()); assert_eq!( url.as_str(), @@ -380,7 +389,7 @@ mod tests { #[test] fn h3_lookup_url_does_not_duplicate_v2_base_path() { let base_url = url::Url::parse("https://dns.example.test:4433/api/v2/").expect("url"); - let url = lookup_url(&base_url, "demo.dhttp.net", None); + let url = lookup_url(&base_url, "demo.dhttp.net", EndpointLookup::default()); assert_eq!( url.as_str(), @@ -394,7 +403,7 @@ mod tests { let url = lookup_url( &base_url, "demo.dhttp.net", - Some(CertificateSequence::from(3u8)), + EndpointLookup::exact(CertificateSequence::from(3u8)), ); assert_eq!( @@ -403,6 +412,37 @@ mod tests { ); } + #[test] + fn h3_lookup_url_encodes_endpoint_lookup_matrix() { + let base_url = url::Url::parse("https://dns.example.test:4433").expect("url"); + let one = NonZeroUsize::new(1).unwrap(); + let cases = [ + ( + EndpointLookup::default(), + "https://dns.example.test:4433/api/v2/lookup?host=demo.dhttp.net", + ), + ( + EndpointLookup::exact(CertificateSequence::from(2u8)), + "https://dns.example.test:4433/api/v2/lookup?host=demo.dhttp.net&sequence=2", + ), + ( + EndpointLookup::limit(one), + "https://dns.example.test:4433/api/v2/lookup?host=demo.dhttp.net&sequence_limit=1", + ), + ( + EndpointLookup::all().with_record_limit(one), + "https://dns.example.test:4433/api/v2/lookup?host=demo.dhttp.net&sequence_limit=all&record_limit=1", + ), + ]; + + for (lookup, expected) in cases { + assert_eq!( + lookup_url(&base_url, "demo.dhttp.net", lookup).as_str(), + expected + ); + } + } + #[test] fn lookup_records_selects_first_server_ordered_group() { let response = response_for( diff --git a/src/http.rs b/src/http.rs index 19f7851..be7a701 100644 --- a/src/http.rs +++ b/src/http.rs @@ -1,7 +1,6 @@ use std::{fmt::Display, io, sync::Arc}; use dashmap::DashMap; -use dhttp_identity::certificate::CertificateSequence; use dquic::{ qbase::net::addr::EndpointAddr, qresolve::{Publish, PublishFuture, Resolve, ResolveFuture, Source}, @@ -32,13 +31,13 @@ pub struct HttpResolver { cached_records: DashMap, } -fn lookup_url(base_url: &Url, name: &str, sequence: Option) -> Url { +fn lookup_url( + base_url: &Url, + name: &str, + lookup: crate::resolvers::endpoint_candidates::EndpointLookup, +) -> Url { let mut url = api_url(base_url, LOOKUP_API_PATH, name); - if let Some(sequence) = sequence { - let sequence_text = sequence.get().to_string(); - url.query_pairs_mut() - .append_pair("sequence", &sequence_text); - } + crate::resolvers::endpoint_candidates::append_endpoint_lookup_query(&mut url, lookup); url } @@ -350,7 +349,13 @@ impl Resolve for HttpResolver { let response = self .http_client - .get(lookup_url(&self.base_url, domain, sequence)) + .get(lookup_url( + &self.base_url, + domain, + sequence + .map(crate::resolvers::endpoint_candidates::EndpointLookup::exact) + .unwrap_or_default(), + )) .send() .await? .error_for_status()? @@ -397,10 +402,10 @@ impl crate::resolvers::endpoint_candidates::ResolveEndpointCandidates for HttpRe fn lookup_endpoint_candidates<'a>( &'a self, name: &'a str, - _lookup: crate::resolvers::endpoint_candidates::EndpointLookup, + lookup: crate::resolvers::endpoint_candidates::EndpointLookup, ) -> crate::resolvers::endpoint_candidates::EndpointCandidateFuture<'a> { let lookup = async move { - let Some((domain, _sequence)) = + let Some((domain, sequence)) = crate::resolvers::endpoint_lookup_name_and_sequence(name) else { return Err(Error::NoRecordFound); @@ -409,7 +414,13 @@ impl crate::resolvers::endpoint_candidates::ResolveEndpointCandidates for HttpRe let source = Source::Http { server }; let response = self .http_client - .get(lookup_url(&self.base_url, domain, None)) + .get(lookup_url( + &self.base_url, + domain, + sequence + .map(crate::resolvers::endpoint_candidates::EndpointLookup::exact) + .unwrap_or(lookup), + )) .send() .await? .error_for_status()? @@ -436,7 +447,12 @@ impl crate::resolvers::endpoint_candidates::ResolveEndpointCandidates for HttpRe #[cfg(test)] mod tests { + use std::num::NonZeroUsize; + + use dhttp_identity::certificate::CertificateSequence; + use super::*; + use crate::resolvers::endpoint_candidates::EndpointLookup; fn direct( addr: &str, @@ -498,7 +514,7 @@ mod tests { #[test] fn http_lookup_url_does_not_duplicate_v2_base_path() { let base_url = Url::parse("https://dns.example.test/api/v2/").expect("url"); - let url = lookup_url(&base_url, "demo.dhttp.net", None); + let url = lookup_url(&base_url, "demo.dhttp.net", EndpointLookup::default()); assert_eq!( url.as_str(), @@ -512,7 +528,7 @@ mod tests { let url = lookup_url( &base_url, "demo.dhttp.net", - Some(CertificateSequence::from(7u8)), + EndpointLookup::exact(CertificateSequence::from(7u8)), ); assert_eq!( @@ -520,4 +536,35 @@ mod tests { "https://dns.example.test/api/v2/lookup?host=demo.dhttp.net&sequence=7" ); } + + #[test] + fn http_lookup_url_encodes_endpoint_lookup_matrix() { + let base_url = Url::parse("https://dns.example.test").expect("url"); + let one = NonZeroUsize::new(1).unwrap(); + let cases = [ + ( + EndpointLookup::default(), + "https://dns.example.test/api/v2/lookup?host=demo.dhttp.net", + ), + ( + EndpointLookup::exact(CertificateSequence::from(2u8)), + "https://dns.example.test/api/v2/lookup?host=demo.dhttp.net&sequence=2", + ), + ( + EndpointLookup::limit(one), + "https://dns.example.test/api/v2/lookup?host=demo.dhttp.net&sequence_limit=1", + ), + ( + EndpointLookup::all().with_record_limit(one), + "https://dns.example.test/api/v2/lookup?host=demo.dhttp.net&sequence_limit=all&record_limit=1", + ), + ]; + + for (lookup, expected) in cases { + assert_eq!( + lookup_url(&base_url, "demo.dhttp.net", lookup).as_str(), + expected + ); + } + } } diff --git a/src/resolvers/endpoint_candidates.rs b/src/resolvers/endpoint_candidates.rs index 4ba0d50..9a47e69 100644 --- a/src/resolvers/endpoint_candidates.rs +++ b/src/resolvers/endpoint_candidates.rs @@ -1,6 +1,6 @@ use std::{io, num::NonZeroUsize}; -use dhttp_identity::certificate::{CertificateChainKey, CertificateSequence}; +use dhttp_identity::certificate::{CertificateChainKey, CertificateChainKind, CertificateSequence}; use dquic::{ qbase::net::addr::EndpointAddr as DquicEndpointAddr, qresolve::{Resolve, Source}, @@ -68,6 +68,25 @@ impl EndpointLookup { } } +pub(crate) fn append_endpoint_lookup_query(url: &mut url::Url, lookup: EndpointLookup) { + let mut pairs = url.query_pairs_mut(); + match lookup.sequences { + SequenceQuery::Default => {} + SequenceQuery::Exact(sequence) => { + pairs.append_pair("sequence", &sequence.get().to_string()); + } + SequenceQuery::Limit(limit) => { + pairs.append_pair("sequence_limit", &limit.get().to_string()); + } + SequenceQuery::All => { + pairs.append_pair("sequence_limit", "all"); + } + } + if let Some(limit) = lookup.record_limit { + pairs.append_pair("record_limit", &limit.get().to_string()); + } +} + pub type EndpointCandidateFuture<'a> = BoxFuture<'a, io::Result>; pub trait ResolveEndpointCandidates: Resolve { @@ -115,6 +134,9 @@ pub(crate) fn grouped_endpoint_candidates( } in records { let chain_key = effective_chain_key(&record, fallback_chain_key); + if chain_key.kind() != CertificateChainKind::Primary { + continue; + } let Ok(endpoint) = DquicEndpointAddr::try_from(record) else { continue; }; @@ -208,6 +230,26 @@ mod tests { assert_eq!(groups[1].1.len(), 1); } + #[test] + fn grouping_ignores_secondary_records() { + let groups = grouped_endpoint_candidates([ + TaggedEndpointCandidate { + tag: "secondary", + record: direct("192.0.2.20:4433", false, 1), + fallback_chain_key: None, + }, + TaggedEndpointCandidate { + tag: "primary", + record: direct("192.0.2.10:4433", true, 2), + fallback_chain_key: None, + }, + ]); + + assert_eq!(groups.len(), 1); + assert_eq!(groups[0].0.to_string(), "primary:2"); + assert_eq!(groups[0].1[0].0, "primary"); + } + #[test] fn grouping_uses_fallback_chain_key_for_unmarked_endpoint() { let endpoint = DnsEndpointAddr::direct_v4("192.0.2.60:4433".parse().unwrap()); diff --git a/src/resolvers/endpoint_group.rs b/src/resolvers/endpoint_group.rs index 9902525..c5478cc 100644 --- a/src/resolvers/endpoint_group.rs +++ b/src/resolvers/endpoint_group.rs @@ -100,37 +100,33 @@ mod tests { } #[test] - fn selected_endpoint_addrs_uses_first_chain_key_group() { + fn selected_endpoint_addrs_uses_first_primary_chain_key_group() { let secondary = direct("192.0.2.20:4433", false, 0); let primary_a = direct("192.0.2.10:4433", true, 2); let primary_b = direct("192.0.2.11:4433", true, 2); let selected = super::selected_endpoint_addrs([secondary, primary_a, primary_b]); - assert_eq!(selected.len(), 1); + assert_eq!(selected.len(), 2); assert_eq!( selected[0], - dquic::qbase::net::addr::EndpointAddr::direct("192.0.2.20:4433".parse().unwrap()) + dquic::qbase::net::addr::EndpointAddr::direct("192.0.2.10:4433".parse().unwrap()) + ); + assert_eq!( + selected[1], + dquic::qbase::net::addr::EndpointAddr::direct("192.0.2.11:4433".parse().unwrap()) ); } #[test] - fn selected_endpoint_addrs_uses_one_secondary_chain_key_group_when_no_primary_exists() { + fn selected_endpoint_addrs_returns_empty_when_no_primary_exists() { let secondary_a = direct("192.0.2.20:4433", false, 5); let secondary_b = direct("192.0.2.21:4433", false, 5); let other_secondary = direct("192.0.2.30:4433", false, 6); let selected = super::selected_endpoint_addrs([secondary_a, secondary_b, other_secondary]); - assert_eq!(selected.len(), 2); - assert_eq!( - selected[0], - dquic::qbase::net::addr::EndpointAddr::direct("192.0.2.20:4433".parse().unwrap()) - ); - assert_eq!( - selected[1], - dquic::qbase::net::addr::EndpointAddr::direct("192.0.2.21:4433".parse().unwrap()) - ); + assert!(selected.is_empty()); } #[test] From adad7c4faf4b34a6055037a26ff032307f5d2a61 Mon Sep 17 00:00:00 2001 From: eareimu Date: Wed, 15 Jul 2026 16:19:41 +0800 Subject: [PATCH 3/7] feat: merge ordered primary candidates --- src/h3/lookup.rs | 27 +++--- src/http.rs | 40 ++++----- src/mdns.rs | 91 +++++++++++++++++++- src/resolvers.rs | 122 ++++++++++++++++++++++++++- src/resolvers/endpoint_candidates.rs | 65 ++++++++++++++ 5 files changed, 309 insertions(+), 36 deletions(-) diff --git a/src/h3/lookup.rs b/src/h3/lookup.rs index 553d8cc..b49a1cc 100644 --- a/src/h3/lookup.rs +++ b/src/h3/lookup.rs @@ -305,18 +305,21 @@ where let source = Source::H3 { server: Arc::from(self.base_url.origin().ascii_serialization()), }; - let groups = LookupRecords::decode_candidate_groups(domain, response.as_ref()) - .map_err(io::Error::other)? - .into_iter() - .map(|(chain, endpoints)| EndpointCandidateGroup { - chain, - endpoints: endpoints - .into_iter() - .map(|((), endpoint)| endpoint) - .collect(), - sources: vec![source.clone()], - }) - .collect(); + let groups = crate::resolvers::endpoint_candidates::select_group_pairs( + LookupRecords::decode_candidate_groups(domain, response.as_ref()) + .map_err(io::Error::other)?, + lookup.sequences, + ) + .into_iter() + .map(|(chain, endpoints)| EndpointCandidateGroup { + chain, + endpoints: endpoints + .into_iter() + .map(|((), endpoint)| endpoint) + .collect(), + sources: vec![source.clone()], + }) + .collect(); Ok(EndpointCandidates { groups }) }) diff --git a/src/http.rs b/src/http.rs index be7a701..3c6ec2b 100644 --- a/src/http.rs +++ b/src/http.rs @@ -412,33 +412,33 @@ impl crate::resolvers::endpoint_candidates::ResolveEndpointCandidates for HttpRe }; let server = Arc::from(self.base_url.host_str().unwrap_or("")); let source = Source::Http { server }; + let lookup = sequence + .map(crate::resolvers::endpoint_candidates::EndpointLookup::exact) + .unwrap_or(lookup); let response = self .http_client - .get(lookup_url( - &self.base_url, - domain, - sequence - .map(crate::resolvers::endpoint_candidates::EndpointLookup::exact) - .unwrap_or(lookup), - )) + .get(lookup_url(&self.base_url, domain, lookup)) .send() .await? .error_for_status()? .bytes() .await?; - let groups = decode_candidate_groups(domain, response.as_ref())? - .into_iter() - .map(|(chain, endpoints)| { - crate::resolvers::endpoint_candidates::EndpointCandidateGroup { - chain, - endpoints: endpoints - .into_iter() - .map(|((), endpoint)| endpoint) - .collect(), - sources: vec![source.clone()], - } - }) - .collect(); + let groups = crate::resolvers::endpoint_candidates::select_group_pairs( + decode_candidate_groups(domain, response.as_ref())?, + lookup.sequences, + ) + .into_iter() + .map(|(chain, endpoints)| { + crate::resolvers::endpoint_candidates::EndpointCandidateGroup { + chain, + endpoints: endpoints + .into_iter() + .map(|((), endpoint)| endpoint) + .collect(), + sources: vec![source.clone()], + } + }) + .collect(); Ok(crate::resolvers::endpoint_candidates::EndpointCandidates { groups }) }; Box::pin(lookup.map_err(io::Error::other)) diff --git a/src/mdns.rs b/src/mdns.rs index d75c134..94b1a57 100644 --- a/src/mdns.rs +++ b/src/mdns.rs @@ -356,19 +356,40 @@ impl MdnsResolvers { } } +#[cfg(feature = "dquic-network")] +fn select_candidate_groups( + groups: Vec, + query: crate::resolvers::endpoint_candidates::SequenceQuery, +) -> Vec { + use crate::resolvers::endpoint_candidates::SequenceQuery; + + match query { + SequenceQuery::Default => groups.into_iter().take(3).collect(), + SequenceQuery::Exact(sequence) => groups + .into_iter() + .filter(|group| group.chain.sequence() == sequence) + .collect(), + SequenceQuery::Limit(limit) => groups.into_iter().take(limit.get()).collect(), + SequenceQuery::All => groups, + } +} + #[cfg(feature = "dquic-network")] impl crate::resolvers::endpoint_candidates::ResolveEndpointCandidates for MdnsResolvers { fn lookup_endpoint_candidates<'a>( &'a self, name: &'a str, - _lookup: crate::resolvers::endpoint_candidates::EndpointLookup, + lookup: crate::resolvers::endpoint_candidates::EndpointLookup, ) -> crate::resolvers::endpoint_candidates::EndpointCandidateFuture<'a> { Box::pin(async move { - let Some((domain, _sequence)) = + let Some((domain, sequence)) = crate::resolvers::endpoint_lookup_name_and_sequence(name) else { return Err(io::Error::other("no DNS record found")); }; + let lookup = sequence + .map(crate::resolvers::endpoint_candidates::EndpointLookup::exact) + .unwrap_or(lookup); let mut lookup_futures = FuturesUnordered::new(); let mut has_resolver = false; @@ -425,6 +446,7 @@ impl crate::resolvers::endpoint_candidates::ResolveEndpointCandidates for MdnsRe } }) .collect(); + let groups = select_candidate_groups(groups, lookup.sequences); Ok(crate::resolvers::endpoint_candidates::EndpointCandidates { groups }) }) @@ -461,3 +483,68 @@ impl Resolve for MdnsResolvers { self.query_with_sequence(domain, sequence).boxed() } } + +#[cfg(all(test, feature = "dquic-network"))] +mod tests { + use std::num::NonZeroUsize; + + use dhttp_identity::certificate::{ + CertificateChainKey, CertificateChainKind, CertificateSequence, + }; + + use super::*; + use crate::resolvers::endpoint_candidates::{ + EndpointCandidateGroup, EndpointLookup, SequenceQuery, + }; + + fn group(sequence: u8) -> EndpointCandidateGroup { + EndpointCandidateGroup { + chain: CertificateChainKey::new( + CertificateSequence::from(sequence), + CertificateChainKind::Primary, + ), + endpoints: Vec::new(), + sources: Vec::new(), + } + } + + #[test] + fn mdns_candidate_selection_applies_sequence_query_locally() { + let groups = || vec![group(2), group(1), group(3), group(4)]; + let sequences = |groups: Vec| { + groups + .into_iter() + .map(|group| group.chain.sequence().get()) + .collect::>() + }; + + assert_eq!( + sequences(select_candidate_groups(groups(), SequenceQuery::Default)), + vec![2, 1, 3] + ); + assert_eq!( + sequences(select_candidate_groups( + groups(), + SequenceQuery::Limit(NonZeroUsize::new(2).unwrap()), + )), + vec![2, 1] + ); + assert_eq!( + sequences(select_candidate_groups( + groups(), + SequenceQuery::Exact(CertificateSequence::from(3u8)), + )), + vec![3] + ); + assert_eq!( + sequences(select_candidate_groups(groups(), SequenceQuery::All)), + vec![2, 1, 3, 4] + ); + + let lookup = EndpointLookup::all().with_record_limit(NonZeroUsize::new(1).unwrap()); + assert_eq!( + sequences(select_candidate_groups(groups(), lookup.sequences)), + vec![2, 1, 3, 4] + ); + } +} diff --git a/src/resolvers.rs b/src/resolvers.rs index a9318c8..06ab5fb 100644 --- a/src/resolvers.rs +++ b/src/resolvers.rs @@ -359,7 +359,8 @@ impl Resolvers { lookup: crate::resolvers::endpoint_candidates::EndpointLookup, ) -> Result { let mut errors = vec![]; - let mut groups = Vec::new(); + let mut groups = + Vec::::new(); for entry in self.resolvers.clone() { let Some(candidate_resolver) = entry.endpoint_candidates else { @@ -374,7 +375,27 @@ impl Resolvers { .lookup_endpoint_candidates(name, lookup) .await { - Ok(candidates) => groups.extend(candidates.groups), + Ok(candidates) => { + for mut group in candidates.groups { + if let Some(existing) = groups + .iter_mut() + .find(|existing| existing.chain == group.chain) + { + for endpoint in group.endpoints.drain(..) { + if !existing.endpoints.contains(&endpoint) { + existing.endpoints.push(endpoint); + } + } + for source in group.sources.drain(..) { + if !existing.sources.contains(&source) { + existing.sources.push(source); + } + } + } else { + groups.push(group); + } + } + } Err(error) => errors.push((entry.resolver.to_string(), error)), } } @@ -647,6 +668,103 @@ mod tests { assert_eq!(candidates.groups[1].chain.to_string(), "primary:0"); } + #[cfg(feature = "resolvers")] + #[derive(Debug)] + struct CandidateSetResolver { + label: &'static str, + groups: Vec<(u8, &'static str, dquic::qresolve::Source)>, + } + + #[cfg(feature = "resolvers")] + impl fmt::Display for CandidateSetResolver { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.write_str(self.label) + } + } + + #[cfg(feature = "resolvers")] + impl dquic::qresolve::Resolve for CandidateSetResolver { + fn lookup<'l>(&'l self, _name: &'l str) -> dquic::qresolve::ResolveFuture<'l> { + use futures::{FutureExt, StreamExt, stream}; + async { Ok(stream::empty().boxed()) }.boxed() + } + } + + #[cfg(feature = "resolvers")] + impl crate::resolvers::endpoint_candidates::ResolveEndpointCandidates for CandidateSetResolver { + fn lookup_endpoint_candidates<'a>( + &'a self, + _name: &'a str, + _lookup: crate::resolvers::endpoint_candidates::EndpointLookup, + ) -> crate::resolvers::endpoint_candidates::EndpointCandidateFuture<'a> { + use dhttp_identity::certificate::{ + CertificateChainKey, CertificateChainKind, CertificateSequence, + }; + use futures::FutureExt; + + let groups = self + .groups + .iter() + .map(|(sequence, endpoint, source)| { + crate::resolvers::endpoint_candidates::EndpointCandidateGroup { + chain: CertificateChainKey::new( + CertificateSequence::from(*sequence), + CertificateChainKind::Primary, + ), + endpoints: vec![dquic::qbase::net::addr::EndpointAddr::direct( + endpoint.parse().unwrap(), + )], + sources: vec![source.clone()], + } + }) + .collect(); + async move { Ok(crate::resolvers::endpoint_candidates::EndpointCandidates { groups }) } + .boxed() + } + } + + #[cfg(feature = "resolvers")] + #[tokio::test] + async fn aggregate_endpoint_candidates_merge_duplicate_sequences_stably() { + use dquic::qresolve::Source; + + let resolvers = Resolvers::new() + .with_candidate_resolver(Arc::new(CandidateSetResolver { + label: "a", + groups: vec![ + (2, "192.0.2.20:4433", Source::System), + (1, "192.0.2.10:4433", Source::System), + ], + })) + .with_candidate_resolver(Arc::new(CandidateSetResolver { + label: "b", + groups: vec![ + (2, "192.0.2.21:4433", Source::Dht), + (3, "192.0.2.30:4433", Source::Dht), + ], + })); + + let candidates = resolvers + .lookup_endpoint_candidates( + "demo.dhttp.net", + crate::resolvers::endpoint_candidates::EndpointLookup::all(), + ) + .await + .expect("candidate lookup succeeds"); + + let sequences = candidates + .groups + .iter() + .map(|group| group.chain.sequence().get()) + .collect::>(); + assert_eq!(sequences, vec![2, 1, 3]); + assert_eq!(candidates.groups[0].endpoints.len(), 2); + assert_eq!( + candidates.groups[0].sources, + vec![Source::System, Source::Dht] + ); + } + #[cfg(feature = "resolvers")] #[test] fn dns_scheme_round_trips_supported_schemes_and_rejects_dht() { diff --git a/src/resolvers/endpoint_candidates.rs b/src/resolvers/endpoint_candidates.rs index 9a47e69..57062ac 100644 --- a/src/resolvers/endpoint_candidates.rs +++ b/src/resolvers/endpoint_candidates.rs @@ -87,6 +87,21 @@ pub(crate) fn append_endpoint_lookup_query(url: &mut url::Url, lookup: EndpointL } } +pub(crate) fn select_group_pairs( + groups: Vec<(CertificateChainKey, T)>, + query: SequenceQuery, +) -> Vec<(CertificateChainKey, T)> { + match query { + SequenceQuery::Default => groups.into_iter().take(3).collect(), + SequenceQuery::Exact(sequence) => groups + .into_iter() + .filter(|(chain, _)| chain.sequence() == sequence) + .collect(), + SequenceQuery::Limit(limit) => groups.into_iter().take(limit.get()).collect(), + SequenceQuery::All => groups, + } +} + pub type EndpointCandidateFuture<'a> = BoxFuture<'a, io::Result>; pub trait ResolveEndpointCandidates: Resolve { @@ -250,6 +265,56 @@ mod tests { assert_eq!(groups[0].1[0].0, "primary"); } + #[test] + fn sequence_query_selects_ordered_group_pairs() { + let pairs = || { + vec![ + ( + CertificateChainKey::new( + CertificateSequence::from(2u8), + CertificateChainKind::Primary, + ), + "two", + ), + ( + CertificateChainKey::new( + CertificateSequence::from(1u8), + CertificateChainKind::Primary, + ), + "one", + ), + ( + CertificateChainKey::new( + CertificateSequence::from(3u8), + CertificateChainKind::Primary, + ), + "three", + ), + ] + }; + + assert_eq!( + select_group_pairs( + pairs(), + SequenceQuery::Exact(CertificateSequence::from(1u8)), + ), + vec![( + CertificateChainKey::new( + CertificateSequence::from(1u8), + CertificateChainKind::Primary, + ), + "one", + )] + ); + assert_eq!( + select_group_pairs(pairs(), SequenceQuery::Limit(NonZeroUsize::new(2).unwrap()),) + .into_iter() + .map(|(_, value)| value) + .collect::>(), + vec!["two", "one"] + ); + } + #[test] fn grouping_uses_fallback_chain_key_for_unmarked_endpoint() { let endpoint = DnsEndpointAddr::direct_v4("192.0.2.60:4433".parse().unwrap()); From e471b8ed406a20aa69a518c3fd06445ceb242aa4 Mon Sep 17 00:00:00 2001 From: eareimu Date: Wed, 15 Jul 2026 16:26:34 +0800 Subject: [PATCH 4/7] feat: publish primary sequence metadata --- src/mdns/service.rs | 6 ++ src/publishers/packet.rs | 100 +++++++++++++++++++- src/publishers/publisher.rs | 178 +++++++++++++++++++++++++++++++----- 3 files changed, 259 insertions(+), 25 deletions(-) diff --git a/src/mdns/service.rs b/src/mdns/service.rs index 67c2cc4..d905443 100644 --- a/src/mdns/service.rs +++ b/src/mdns/service.rs @@ -226,6 +226,12 @@ impl Mdns { guard.insert(local_name, eps); } + #[cfg(test)] + pub(crate) fn published_endpoints(&self, host_name: &str) -> Option> { + let local_name = Self::local_name(self.service_name.clone(), host_name.to_owned()); + self.hosts.lock().unwrap().get(&local_name).cloned() + } + #[inline] pub(crate) fn protocol(&self) -> Arc { self.inner diff --git a/src/publishers/packet.rs b/src/publishers/packet.rs index f863f9c..2abccc4 100644 --- a/src/publishers/packet.rs +++ b/src/publishers/packet.rs @@ -20,13 +20,15 @@ pub enum EncodeAuthorityDnsPacketError { }, #[snafu(display("failed to encode endpoint address"))] EncodeEndpoint, + #[snafu(display("secondary authority cannot publish dns endpoints"))] + SecondaryAuthority, } -pub(crate) fn dns_packet_for_authority( +pub(crate) fn dns_endpoints_for_authority( authority: &dyn LocalAuthority, name: &str, endpoints: &mut dyn Iterator, -) -> Result, EncodeAuthorityDnsPacketError> { +) -> Result, EncodeAuthorityDnsPacketError> { ensure!( authority.name() == name, encode_authority_dns_packet_error::AuthorityNameMismatchSnafu @@ -42,11 +44,26 @@ pub(crate) fn dns_packet_for_authority( let Ok(mut endpoint) = DnsEndpointAddr::try_from(endpoint) else { return encode_authority_dns_packet_error::EncodeEndpointSnafu.fail(); }; - endpoint.set_main(chain.kind() == CertificateChainKind::Primary); + endpoint.set_main(true); endpoint.set_sequence(chain.sequence()); encoded.push(endpoint); } + ensure!( + encoded.is_empty() || chain.kind() == CertificateChainKind::Primary, + encode_authority_dns_packet_error::SecondaryAuthoritySnafu + ); + + Ok(encoded) +} + +pub(crate) fn dns_packet_for_authority( + authority: &dyn LocalAuthority, + name: &str, + endpoints: &mut dyn Iterator, +) -> Result, EncodeAuthorityDnsPacketError> { + let encoded = dns_endpoints_for_authority(authority, name, endpoints)?; + let mut hosts = HashMap::new(); hosts.insert(name.to_owned(), encoded); Ok(MdnsPacket::answer(0, &hosts).to_bytes()) @@ -61,7 +78,7 @@ mod tests { use futures::future::BoxFuture; use rustls::pki_types::CertificateDer; - use super::dns_packet_for_authority; + use super::{dns_endpoints_for_authority, dns_packet_for_authority}; use crate::core::parser::{ packet::be_packet, record::{RData, Type}, @@ -99,6 +116,28 @@ mod tests { } } + fn authority_with_chain( + name: &'static str, + sequence: u8, + kind: dhttp_identity::certificate::CertificateChainKind, + ) -> TestAuthority { + let mut certificate = include_bytes!("../../tests/fixtures/valid.der").to_vec(); + let marker = b"0:0:0123456789abcdef"; + let offset = certificate + .windows(marker.len()) + .position(|window| window == marker) + .expect("fixture contains dhttp subject key identifier"); + certificate[offset] = b'0' + sequence; + certificate[offset + 2] = match kind { + dhttp_identity::certificate::CertificateChainKind::Primary => b'0', + dhttp_identity::certificate::CertificateChainKind::Secondary => b'1', + }; + TestAuthority { + name, + cert_chain: vec![CertificateDer::from(certificate)], + } + } + fn endpoint() -> DquicEndpointAddr { DquicEndpointAddr::direct(SocketAddr::V4(SocketAddrV4::new( Ipv4Addr::new(203, 0, 113, 10), @@ -157,4 +196,57 @@ mod tests { assert!(remain.is_empty()); assert!(parsed.answers.is_empty()); } + + #[test] + fn dns_endpoints_for_primary_authority_stamp_sequence_metadata() { + let authority = authority_with_chain( + "client.example.com.dhttp.net", + 7, + dhttp_identity::certificate::CertificateChainKind::Primary, + ); + let mut endpoints = std::iter::once(endpoint()); + + let encoded = + dns_endpoints_for_authority(&authority, "client.example.com.dhttp.net", &mut endpoints) + .expect("primary endpoints encode"); + + assert_eq!(encoded.len(), 1); + assert!(encoded[0].is_main()); + assert_eq!(encoded[0].normalized_sequence().get(), 7); + } + + #[test] + fn dns_endpoints_for_secondary_authority_reject_non_empty_publication() { + let authority = authority_with_chain( + "client.example.com.dhttp.net", + 7, + dhttp_identity::certificate::CertificateChainKind::Secondary, + ); + let mut endpoints = std::iter::once(endpoint()); + + let error = + dns_endpoints_for_authority(&authority, "client.example.com.dhttp.net", &mut endpoints) + .expect_err("secondary publication is rejected"); + + assert_eq!( + error.to_string(), + "secondary authority cannot publish dns endpoints" + ); + } + + #[test] + fn dns_endpoints_for_secondary_authority_allows_empty_clear() { + let authority = authority_with_chain( + "client.example.com.dhttp.net", + 7, + dhttp_identity::certificate::CertificateChainKind::Secondary, + ); + let mut endpoints = std::iter::empty(); + + let encoded = + dns_endpoints_for_authority(&authority, "client.example.com.dhttp.net", &mut endpoints) + .expect("empty secondary clear is allowed"); + + assert!(encoded.is_empty()); + } } diff --git a/src/publishers/publisher.rs b/src/publishers/publisher.rs index e2217a2..f15eb5b 100644 --- a/src/publishers/publisher.rs +++ b/src/publishers/publisher.rs @@ -2,7 +2,7 @@ use std::{fmt, io, sync::Arc}; use dhttp_identity::name::Name; use dquic::qresolve::Publish; -use snafu::{IntoError, ResultExt, Snafu}; +use snafu::{OptionExt, ResultExt, Snafu}; use super::{AddressView, PublishScope}; @@ -17,6 +17,17 @@ pub enum PublisherError { #[cfg(all(feature = "mdns", feature = "dquic-network"))] #[snafu(display("all mdns publishers failed"))] Mdns { source: MdnsPublishersError }, + #[cfg(all(feature = "mdns", feature = "dquic-network"))] + #[snafu(display("failed to get mdns publisher local authority"))] + MdnsLocalAuthority { source: h3x::quic::ConnectionError }, + #[cfg(all(feature = "mdns", feature = "dquic-network"))] + #[snafu(display("anonymous endpoint cannot publish mdns records"))] + MdnsAnonymousEndpoint, + #[cfg(all(feature = "mdns", feature = "dquic-network"))] + #[snafu(display("failed to encode mdns dns records"))] + MdnsEncode { + source: crate::publishers::packet::EncodeAuthorityDnsPacketError, + }, } #[derive(Clone)] @@ -31,7 +42,10 @@ enum PublisherKind { publisher: Arc, }, #[cfg(all(feature = "mdns", feature = "dquic-network"))] - Mdns(Arc), + Mdns { + resolvers: Arc, + authority: Arc, + }, } #[cfg(all(feature = "mdns", feature = "dquic-network"))] @@ -67,7 +81,7 @@ impl fmt::Debug for Publisher { .field("publisher", publisher) .finish(), #[cfg(all(feature = "mdns", feature = "dquic-network"))] - PublisherKind::Mdns(resolvers) => f + PublisherKind::Mdns { resolvers, .. } => f .debug_struct("Publisher") .field("mdns", resolvers) .finish(), @@ -80,7 +94,7 @@ impl fmt::Display for Publisher { match &self.inner { PublisherKind::Custom { publisher, .. } => fmt::Display::fmt(publisher, f), #[cfg(all(feature = "mdns", feature = "dquic-network"))] - PublisherKind::Mdns(resolvers) => fmt::Display::fmt(resolvers, f), + PublisherKind::Mdns { resolvers, .. } => fmt::Display::fmt(resolvers, f), } } } @@ -107,9 +121,15 @@ impl Publisher { } #[cfg(all(feature = "mdns", feature = "dquic-network"))] - pub fn mdns(resolvers: Arc) -> Self { + pub fn mdns(resolvers: Arc, authority: Arc) -> Self + where + A: h3x::quic::DynWithLocalAuthority + 'static, + { Self { - inner: PublisherKind::Mdns(resolvers), + inner: PublisherKind::Mdns { + resolvers, + authority, + }, } } @@ -122,7 +142,10 @@ impl Publisher { publish_selected(publisher.as_ref(), scope, name, view).await } #[cfg(all(feature = "mdns", feature = "dquic-network"))] - PublisherKind::Mdns(resolvers) => publish_mdns(resolvers, name, view).await, + PublisherKind::Mdns { + resolvers, + authority, + } => publish_mdns(resolvers, authority.as_ref(), name, view).await, } } } @@ -153,43 +176,53 @@ where #[cfg(all(feature = "mdns", feature = "dquic-network"))] async fn publish_mdns( resolvers: &crate::mdns::MdnsResolvers, + authority_provider: &dyn h3x::quic::DynWithLocalAuthority, name: &Name<'_>, view: &V, ) -> Result<(), PublisherError> where V: AddressView + Sync, { + let authority = authority_provider + .local_authority() + .await + .context(publisher_error::MdnsLocalAuthoritySnafu)? + .context(publisher_error::MdnsAnonymousEndpointSnafu)?; + let mut no_endpoints = std::iter::empty(); + crate::publishers::packet::dns_endpoints_for_authority( + authority.as_ref(), + name.as_str(), + &mut no_endpoints, + ) + .context(publisher_error::MdnsEncodeSnafu)?; let bound_resolvers = resolvers.bound_resolvers(); if bound_resolvers.is_empty() { tracing::debug!(name = %name, "no mdns publishers currently bound"); return Ok(()); } - let mut errors = Vec::new(); - let mut succeeded = false; for bound in bound_resolvers { let scope = PublishScope::LocalLink { device: bound.device.clone().into(), family: bound.family, }; - match publish_selected(&bound.resolver, &scope, name, view).await { - Ok(()) => succeeded = true, - Err(PublisherError::Publish { source, .. }) => { - errors.push((bound.resolver.to_string(), source)); - } - Err(error) => return Err(error), - } + let mut endpoints = view.endpoints(scope.selector()); + let endpoints = crate::publishers::packet::dns_endpoints_for_authority( + authority.as_ref(), + name.as_str(), + &mut endpoints, + ) + .context(publisher_error::MdnsEncodeSnafu)?; + bound.resolver.insert_host(name.to_string(), endpoints); } - if succeeded { - Ok(()) - } else { - Err(publisher_error::MdnsSnafu.into_error(MdnsPublishersError { errors })) - } + Ok(()) } #[cfg(test)] mod tests { + #[cfg(all(feature = "mdns", feature = "dquic-network"))] + use std::sync::atomic::{AtomicUsize, Ordering}; use std::{ fmt, io, net::{Ipv4Addr, SocketAddr, SocketAddrV4}, @@ -264,6 +297,109 @@ mod tests { EndpointAddr::direct(SocketAddr::V4(SocketAddrV4::new(Ipv4Addr::from(ip), port))) } + #[cfg(all(feature = "mdns", feature = "dquic-network"))] + #[derive(Debug)] + struct RecordingAuthorityProvider { + calls: AtomicUsize, + authority: Arc, + } + + #[cfg(all(feature = "mdns", feature = "dquic-network"))] + impl h3x::quic::DynWithLocalAuthority for RecordingAuthorityProvider { + fn local_authority( + &self, + ) -> futures::future::BoxFuture< + '_, + Result< + Option>, + h3x::quic::ConnectionError, + >, + > { + self.calls.fetch_add(1, Ordering::SeqCst); + futures::future::ready(Ok(Some(self.authority.clone()))).boxed() + } + } + + #[cfg(all(feature = "mdns", feature = "dquic-network"))] + #[derive(Debug)] + struct MdnsTestAuthority { + cert_chain: Vec>, + } + + #[cfg(all(feature = "mdns", feature = "dquic-network"))] + impl dhttp_identity::identity::LocalAuthority for MdnsTestAuthority { + fn name(&self) -> &str { + "alice.dhttp.net" + } + + fn cert_chain(&self) -> &[rustls::pki_types::CertificateDer<'static>] { + &self.cert_chain + } + + fn sign( + &self, + _data: &[u8], + ) -> futures::future::BoxFuture<'_, Result, dhttp_identity::identity::SignError>> + { + futures::future::ready(Ok(Vec::new())).boxed() + } + } + + #[cfg(all(feature = "mdns", feature = "dquic-network"))] + #[tokio::test] + async fn mdns_publish_uses_publisher_owned_authority() { + use std::str::FromStr; + + let mut certificate = include_bytes!("../../tests/fixtures/valid.der").to_vec(); + let marker = b"0:0:0123456789abcdef"; + let offset = certificate + .windows(marker.len()) + .position(|window| window == marker) + .expect("fixture contains dhttp subject key identifier"); + certificate[offset] = b'7'; + let authority = Arc::new(MdnsTestAuthority { + cert_chain: vec![rustls::pki_types::CertificateDer::from(certificate)], + }); + let provider = Arc::new(RecordingAuthorityProvider { + calls: AtomicUsize::new(0), + authority, + }); + let pattern = h3x::dquic::binds::BindPattern::from_str("iface://v4.lo:0") + .expect("valid loopback pattern"); + let resolvers = Arc::new( + crate::mdns::MdnsResolvers::bind( + h3x::dquic::Network::builder().build(), + Arc::new(vec![pattern]), + "_test._udp.local", + ) + .await, + ); + let publisher = Publisher::mdns(resolvers.clone(), provider.clone()); + let view = crate::publishers::PublishAddresses::new().local_link( + "lo", + Family::V4, + [endpoint([127, 0, 0, 1], 4433)], + ); + let name = Name::try_from("alice.dhttp.net").expect("valid name"); + + publisher + .publish(&name, &view) + .await + .expect("empty mdns publication succeeds"); + + assert_eq!(provider.calls.load(Ordering::SeqCst), 1); + let bound = resolvers.bound_resolvers(); + assert!(!bound.is_empty(), "loopback mDNS resolver must be bound"); + let records = bound[0] + .resolver + .published_endpoints("alice.dhttp.net") + .expect("published host exists"); + assert_eq!(records.len(), 1); + assert!(!records[0].is_signed()); + assert!(records[0].is_main()); + assert_eq!(records[0].normalized_sequence().get(), 7); + } + #[tokio::test] async fn custom_publisher_selects_wide_area_addresses() { let wide = endpoint([203, 0, 113, 10], 4433); From d3b9fb87f4c61fbef45233665fc825193f33019c Mon Sep 17 00:00:00 2001 From: eareimu Date: Wed, 15 Jul 2026 16:51:51 +0800 Subject: [PATCH 5/7] fix: disambiguate modern endpoint records --- src/core/parser/record/endpoint.rs | 35 ++++++++++++++++++++++++++++++ 1 file changed, 35 insertions(+) diff --git a/src/core/parser/record/endpoint.rs b/src/core/parser/record/endpoint.rs index 320352a..2b5ab06 100644 --- a/src/core/parser/record/endpoint.rs +++ b/src/core/parser/record/endpoint.rs @@ -558,6 +558,14 @@ pub(crate) fn be_endpoint_addr_compat( ]; if legacy_lengths.contains(&(rdlen as usize)) { + // Modern records have variable length, so a valid modern encoding can + // legitimately be 12, 18, or 36 bytes as well. Prefer a complete + // modern parse and only fall back to the legacy layout when it fails. + if let Ok((remaining, endpoint)) = be_endpoint_addr(input) + && remaining.is_empty() + { + return Ok((remaining, endpoint)); + } return be_legacy_endpoint_addr_by_length(input, rdlen); } @@ -1025,6 +1033,33 @@ mod tests { } } + #[test] + fn compat_parser_does_not_misclassify_modern_lengths_as_legacy() { + let mut direct = EndpointAddr::direct_v4("203.0.113.10:4433".parse().unwrap()); + direct.set_main(true); + direct.set_sequence(CertificateSequence::try_from(10u32).unwrap()); + direct.set_load(Some(1.0)); + + let mut nat = EndpointAddr::nat_v4( + "198.51.100.10:4433".parse().unwrap(), + "192.0.2.10:4433".parse().unwrap(), + ); + nat.set_main(true); + nat.set_sequence(CertificateSequence::from(1u8)); + nat.set_load(Some(2.0)); + + for endpoint in [direct, nat] { + let mut buf = BytesMut::new(); + buf.put_endpoint_addr(&endpoint); + assert!([12, 18].contains(&buf.len())); + + let (remaining, decoded) = + be_endpoint_addr_compat(&buf, u16::try_from(buf.len()).unwrap()).unwrap(); + assert!(remaining.is_empty()); + assert_eq!(decoded, endpoint); + } + } + #[test] fn endpoint_signature_roundtrip_and_verify() { #[derive(Debug)] From 8a583c18d0d4d67565d6dfb9f034ad6ebc33adf6 Mon Sep 17 00:00:00 2001 From: eareimu Date: Wed, 15 Jul 2026 20:22:33 +0800 Subject: [PATCH 6/7] fix: gate resolver URL helpers by feature --- src/resolvers/endpoint_candidates.rs | 2 ++ 1 file changed, 2 insertions(+) diff --git a/src/resolvers/endpoint_candidates.rs b/src/resolvers/endpoint_candidates.rs index 57062ac..101770f 100644 --- a/src/resolvers/endpoint_candidates.rs +++ b/src/resolvers/endpoint_candidates.rs @@ -68,6 +68,7 @@ impl EndpointLookup { } } +#[cfg(any(feature = "h3", feature = "http"))] pub(crate) fn append_endpoint_lookup_query(url: &mut url::Url, lookup: EndpointLookup) { let mut pairs = url.query_pairs_mut(); match lookup.sequences { @@ -87,6 +88,7 @@ pub(crate) fn append_endpoint_lookup_query(url: &mut url::Url, lookup: EndpointL } } +#[cfg(any(feature = "h3", feature = "http", test))] pub(crate) fn select_group_pairs( groups: Vec<(CertificateChainKey, T)>, query: SequenceQuery, From 0a493fa7652e9cc80ba06bf544268c2439513c4a Mon Sep 17 00:00:00 2001 From: eareimu Date: Wed, 15 Jul 2026 20:22:34 +0800 Subject: [PATCH 7/7] chore: prepare dyns v0.7.0-beta.1 --- Cargo.toml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/Cargo.toml b/Cargo.toml index fe538e0..fd4e777 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,7 +1,7 @@ [package] name = "dyns" description = "DNS discovery and resolver support for DHTTP applications" -version = "0.6.0-beta.3" +version = "0.7.0-beta.1" edition = "2024" license = "Apache-2.0" repository = "https://github.com/genmeta/ddns"