From 26ba422a97c77bd4d788d9877c51d33109e7cb08 Mon Sep 17 00:00:00 2001 From: mHost Developer Date: Mon, 28 Sep 2026 23:54:55 +0800 Subject: [PATCH] =?UTF-8?q?fix(dns):=20server=20=E5=B1=82=E6=9F=A5?= =?UTF-8?q?=E8=AF=A2=E5=B9=B6=E5=8F=91=E5=8C=96=EF=BC=8C=E6=B6=88=E9=99=A4?= =?UTF-8?q?=E9=98=9F=E5=A4=B4=E9=98=BB=E5=A1=9E=20(issue=20#223)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit UDP 主循环此前把 handle_dns_request(...).await 内联在 recv_from 分支里, 一条撞上慢上游的查询(config.timeout_ms 默认 3000ms)没返回之前 recv_from 根本不会被 poll,所有后续查询全部排队 —— 最坏吞吐 ≈ 1/上游 RTT。 用户感知为「部分网站突然打不开、过几秒恢复」。 改成与 proxy.rs(#76)一致的分发模型: - 主循环只做 recv_from + 分发,每个查询交给独立 task; - 请求字节用 bytes::Bytes 拷贝出栈上 4 KiB 缓冲(对齐 proxy 的 P-R5 / P-R4), 消除「原 buffer 被慢上游 await 阻塞住」的问题; - Semaphore 限流(MAX_CONCURRENT_CLIENT_QUERIES = 1024,与 proxy 对齐), try_acquire_owned 发生在 spawn **之前**,且刻意放在拷贝 / 读 resolver 之前 —— 命中上限时那些 per-query 开销全是白做的; - JoinSet 管理在途 task,shutdown 走 abort_all + drain:stop() 立即返回, 不被最长 3000ms 的上游超时拖住(见 AGENTS.md 的 exit cleanup 契约), 同时避免 detached task 对着已 drop 的 socket 刷无意义 warn; - resolver_slot 的读提到 spawn 之前,从「每 task 一次读锁」降为「每包一次」, hot-swap 逐查询生效的语义不变。 server 层与 proxy 不同,不需要 per-query 临时 socket:socket 只用于无状态的 send_to(没有 connect() + recv 配对),共享并发发送是安全的。 顺带去掉 build_answer_response 的 answer 双重 clone(参数改收 &Record, clone 只发生在 add_answer 内部那一次),cache-miss 的 answer 查询每次省一次 Record 分配。 测试: - test_dns_server_no_head_of_line_blocking —— 已在修复前的串行实现上验证过 会失败(fast 查询耗时 1.204s,阈值 500ms),修复后为毫秒级; - test_dns_server_concurrency_capped —— 洪水下断言 available_permits() 归零; - test_dns_server_stop_aborts_inflight —— 断言 stop() 在 500ms 内返回。 flood 测试里 client socket 用完即丢,不持有 2048 个 fd。本测试 fd 峰值 ~1026(主要来自 hickory 为每个在途查询开的上游 socket),与已有的 proxy 并发测试同量级(后者峰值 ~3074,长期 CI 绿)。 cargo fmt / clippy --workspace --all-targets -D warnings / 全量测试(606)全绿。 --- src-tauri/crates/mhost-dns/src/server.rs | 449 +++++++++++++++++++++-- 1 file changed, 415 insertions(+), 34 deletions(-) diff --git a/src-tauri/crates/mhost-dns/src/server.rs b/src-tauri/crates/mhost-dns/src/server.rs index 0aac6bc..9ee5080 100644 --- a/src-tauri/crates/mhost-dns/src/server.rs +++ b/src-tauri/crates/mhost-dns/src/server.rs @@ -41,6 +41,14 @@ pub(crate) const LOCAL_RULE_TTL: u32 = 300; /// UDP 缓冲区大小(EDNS(0) 协商后的最大响应长度)。 const UDP_BUF_SIZE: usize = 4096; +/// **fix(issue #223)**:server 层每个客户端查询允许的并发上限。 +/// +/// 主循环改成「recv_from + 分发」之后(issue #223),每收一个包就 spawn +/// 一个 task。为了防止恶意/异常客户端(DoS)打满资源,用 semaphore 限流; +/// 数值与 proxy 层的同名常量对齐。超过上限的 query 立即丢弃(UDP 协议允许 +/// 丢包,调用方会重试)。 +pub const MAX_CONCURRENT_CLIENT_QUERIES: usize = 1024; + /// DNS 响应缓存条目:记录列表 + 过期时间。 type CacheEntry = (Vec, Instant); @@ -95,6 +103,13 @@ pub struct DnsServer { /// 返回 0.0.0.0 / NXDOMAIN。共享 `Arc` 让 spawn 任务零锁读 + 命令 /// 层通过 `reload_ad_block_rules` 在外面 hot-reload。 ad_block_engine: Arc, + /// **fix(issue #223)**:查询并发上限信号量。主循环对每个收到的 query + /// 都 spawn task,用它防止 DoS 场景打爆 runtime。 + /// + /// 生命周期与 server 绑定(`new()` 里创建一次),`start()/stop()` 不重建 —— + /// 与 proxy 层的 `concurrency` 字段语义保持一致,这样测试可以在 server + /// 运行期间通过 `available_permits()` 观测限流是否真的生效。 + concurrency: Arc, } impl DnsServer { @@ -122,9 +137,18 @@ impl DnsServer { refresh_shutdown: Arc::new(tokio::sync::Notify::new()), cache: Arc::new(PlMutex::new(LruCache::new(cache_size))), ad_block_engine: Arc::new(AdBlockEngine::new()), + concurrency: Arc::new(tokio::sync::Semaphore::new(MAX_CONCURRENT_CLIENT_QUERIES)), }) } + /// 当前可用的并发 permit 数(用于测试 + 监控)。 + /// + /// **fix(issue #223)**:值为 0 说明所有查询槽位都被在途查询占满, + /// 此刻到达的新查询会被直接丢弃。对应 proxy 层的同名方法。 + pub fn available_permits(&self) -> usize { + self.concurrency.available_permits() + } + /// 启动 DNS 服务(异步,在后台运行)。 /// TODO: 当前仅监听 UDP,TCP 支持计划在后续迭代中添加。 /// @@ -181,8 +205,34 @@ impl DnsServer { let cache = self.cache.clone(); let ad_block_engine = self.ad_block_engine.clone(); + // **fix(issue #223)**:socket 包一层 Arc,好让 spawn 出去的 task + // 能直接 send_to 回客户端。 + // + // 注意这里和 proxy.rs 不同,**不需要** per-query 的临时 socket: + // proxy 需要临时 socket 是因为它 `connect()` 后要在同一个 socket 上 + // recv 回包,多个并发 query 会互相把回包读串。server 层是 + // 「收包 → 处理 → send_to 原地址」的无状态发送,没有 recv 配对, + // 所以共享同一个 socket 并发 send_to 是安全的。 + let socket = Arc::new(socket); + let concurrency = Arc::clone(&self.concurrency); + let handle = tokio::spawn(async move { - let mut buf = vec![0u8; UDP_BUF_SIZE]; + // **fix(issue #223)**:主循环只负责 recv_from + 分发。 + // + // 之前 `handle_dns_request(...).await` 是**内联**在这个分支里的, + // 一条慢上游查询(config.timeout_ms 默认 3000ms)没返回之前, + // `recv_from` 根本不会被 poll —— 所有后续查询全部排队等它, + // 最坏吞吐 ≈ 1/上游 RTT,系统 DNS 解析近乎停摆。 + // 表现就是 issue 描述的「部分网站突然打不开、过几秒恢复」。 + // + // 栈上 4 KiB 缓冲(对齐 proxy 的 P-R5):不再跨 await 持有借用, + // 下一轮循环可以安全复用同一块内存,零 alloc。 + let mut buf = [0u8; UDP_BUF_SIZE]; + + // **fix(issue #223)**:JoinSet 管理在途 task,让 shutdown 能 + // 确定性地 abort 它们。裸 tokio::spawn 出来的 task 是 detached 的, + // loop 退出后会继续对着已 drop 的 socket send_to 刷无意义 warn。 + let mut inflight = tokio::task::JoinSet::new(); loop { tokio::select! { @@ -190,45 +240,91 @@ impl DnsServer { let (len, src) = result .map_err(|e| DnsError::Server(format!("recv failed: {}", e)))?; - let request_data = &buf[..len]; + // 限流必须发生在 spawn **之前**:超限时根本不 spawn task, + // 避免洪水场景下还没限住就先被 OOM 击垮(同 proxy 的模型)。 + // + // 也刻意放在拷贝 / 读 resolver **之前**:命中上限时这些 + // per-query 开销全是白做的。 + let permit = match concurrency.clone().try_acquire_owned() { + Ok(p) => p, + Err(_) => { + tracing::warn!( + "DNS server concurrency cap ({}) reached, dropping query from {}", + MAX_CONCURRENT_CLIENT_QUERIES, + src + ); + continue; + } + }; + + // **fix(issue #223)**:spawn 之前必须把请求拷出来 —— + // 下一轮循环会复用 buf。bytes::Bytes 是 Arc-backed + // slice,task 持有它零额外成本(同 proxy 的 P-R4)。 + // 这一处 memcpy 换掉的是「原 buffer 被慢上游 await + // 阻塞住」的问题。 + let request = bytes::Bytes::copy_from_slice(&buf[..len]); + // 每次查询重新读 slot —— 让 refresh 任务的 hot-swap 真正生效。 + // 提到 spawn 之前执行:每收一个包读一次读锁(只是原子 + // ref bump + guard drop),比「每个 task 各读一次」少 N-1 次。 let resolver = resolver_slot.read().clone(); - let response_data = match handle_dns_request( - request_data, - &rule_engine, - &ad_block_engine, - &resolver, - &cache, - ).await { - Some(data) => data, - None => { - // 构造 FormErr 响应,保留原始 request id - if len < 2 { - continue; - } - let id = u16::from_be_bytes([buf[0], buf[1]]); - let mut header = Header::new(); - header.set_id(id); - header.set_message_type(MessageType::Response); - header.set_response_code(ResponseCode::FormErr); - let mut response = Message::new(); - response.set_header(header); - match response.to_bytes() { - Ok(data) => data, - Err(e) => { - tracing::warn!("Failed to encode FormErr response: {}", e); - continue; + + let sock = Arc::clone(&socket); + let rules = Arc::clone(&rule_engine); + let ad_block = Arc::clone(&ad_block_engine); + let query_cache = Arc::clone(&cache); + + inflight.spawn(async move { + // permit 的生命周期 = 查询的生命周期,task 结束自动归还。 + let _permit: tokio::sync::OwnedSemaphorePermit = permit; + + let response_data = match handle_dns_request( + &request, + &rules, + &ad_block, + &resolver, + &query_cache, + ).await { + Some(data) => data, + None => { + // 构造 FormErr 响应,保留原始 request id + if request.len() < 2 { + return; + } + let id = u16::from_be_bytes([request[0], request[1]]); + let mut header = Header::new(); + header.set_id(id); + header.set_message_type(MessageType::Response); + header.set_response_code(ResponseCode::FormErr); + let mut response = Message::new(); + response.set_header(header); + match response.to_bytes() { + Ok(data) => data, + Err(e) => { + tracing::warn!("Failed to encode FormErr response: {}", e); + return; + } } } + }; + + if let Err(e) = sock.send_to(&response_data, src).await { + tracing::warn!("Failed to send DNS response to {}: {}", src, e); } - }; + }); - if let Err(e) = socket.send_to(&response_data, src).await { - tracing::warn!("Failed to send DNS response to {}: {}", src, e); - } + // 非阻塞回收已完成的 task,防止 JoinSet 随服务时长无限增长 + // (已完成 task 的 output 还占着内存)。 + while inflight.try_join_next().is_some() {} } _ = shutdown_rx.recv() => { tracing::info!("DNS server received shutdown signal"); + // **fix(issue #223)**:abort 在途 task 而**不**等它们跑完 —— + // 等下去最坏要 3000ms(上游超时),会把 stop() / 应用退出 / + // 关闭 DNS 模式全部拖住(见 AGENTS.md 的 exit cleanup 契约)。 + // UDP 语义允许丢包,客户端会重试。 + inflight.abort_all(); + while inflight.join_next().await.is_some() {} break; } } @@ -553,8 +649,12 @@ async fn handle_dns_request( // **fix (P-R2, issue #90)**: `*record.clone()` 之前是 `clone Box + deref` // (两次 alloc:Box 本身 + Box 内的 Record)。`record.as_ref().clone()` // 只 alloc 一次(Record 本身),更高效。 + // **fix(issue #223)**:之前是 `answer.clone()` 再传进 + // `build_answer_response(... answer: Record)`,miss 路径每查询多一次 + // Record 分配。改成传 `&Record` 后,函数内部 `add_answer` 需要 owned + // 时的 clone 是唯一一次,省掉调用方提前备的那一份。 let answer = record.as_ref().clone(); - let response_bytes = build_answer_response(&request, &query, answer.clone()); + let response_bytes = build_answer_response(&request, &query, &answer); // TTL=0 是合法 DNS 值("不要缓存"),跳过 put 避免 // 浪费 LRU slot 且下次 peek 立刻过期被 pop。 @@ -592,7 +692,7 @@ fn build_cached_response(request: &Message, query: &Query, records: &[Record]) - } /// 构造标准 Answer 响应(带一条 Answer 记录)。 -fn build_answer_response(request: &Message, query: &Query, answer: Record) -> Option> { +fn build_answer_response(request: &Message, query: &Query, answer: &Record) -> Option> { let mut header = Header::response_from_request(request.header()); header.set_authoritative(false); header.set_recursion_available(true); @@ -600,7 +700,7 @@ fn build_answer_response(request: &Message, query: &Query, answer: Record) -> Op response.set_header(header); response.set_id(request.id()); response.add_query(query.clone()); - response.add_answer(answer); + response.add_answer(answer.clone()); match response.to_bytes() { Ok(bytes) => Some(bytes), Err(e) => { @@ -1018,6 +1118,111 @@ mod tests { buf[..len].to_vec() } + /// 返回一个当前空闲的 UDP 端口(bind 后立刻 drop)。 + /// + /// issue #223 新增的测试不能用 1053/1054/… 这类递增硬编码端口 —— + /// 同一个 test binary 里的测试是并发跑的,撞端口会表现为莫名其妙的 + /// bind 失败。 + async fn free_udp_port() -> u16 { + UdpSocket::bind("127.0.0.1:0") + .await + .unwrap() + .local_addr() + .unwrap() + .port() + } + + /// 构造一个 A 查询的 wire 格式字节。 + fn build_a_query(name: &str) -> Vec { + let query = Query::query(Name::from_utf8(name).unwrap(), RecordType::A); + let mut request = Message::new(); + request.set_id(0x1234); + request.set_recursion_desired(true); + request.add_query(query); + request.to_bytes().unwrap() + } + + /// Spawn 一个 mock upstream:`slow_name` 的查询睡 `delay` 才应答, + /// 其他名字立即应答。 + /// + /// 返回 `(upstream 地址, 收到 slow 查询的通知通道)`。后者让调用方能 + /// **确定性地**等到「slow 查询已经打到上游」的那一刻,而不是靠 sleep + /// 猜 —— 否则 server 还没收到包时就发 fast 查询,测试会假通过。 + /// + /// 上游本身也是 per-query spawn(和 DnsServer 一样), + /// 否则这个夹具自己就会串行化,把被测对象的行为掩盖掉。 + async fn spawn_name_dependent_upstream( + slow_name: &'static str, + delay: Duration, + ) -> (SocketAddr, tokio::sync::mpsc::UnboundedReceiver<()>) { + let socket = Arc::new(UdpSocket::bind("127.0.0.1:0").await.unwrap()); + let addr = socket.local_addr().unwrap(); + let (slow_tx, slow_rx) = tokio::sync::mpsc::unbounded_channel(); + + tokio::spawn(async move { + let mut buf = vec![0u8; 512]; + loop { + let Ok((len, src)) = socket.recv_from(&mut buf).await else { + break; + }; + let payload = buf[..len].to_vec(); + let sock = Arc::clone(&socket); + let tx = slow_tx.clone(); + tokio::spawn(async move { + let Ok(request) = Message::from_bytes(&payload) else { + return; + }; + let name = request + .queries() + .first() + .map(|q| q.name().to_utf8()) + .unwrap_or_default(); + if name.trim_end_matches('.') == slow_name { + let _ = tx.send(()); + tokio::time::sleep(delay).await; + } + + let mut response = Message::new(); + response.set_id(request.id()); + response.set_message_type(MessageType::Response); + response.set_op_code(OpCode::Query); + response.set_recursion_desired(true); + response.set_recursion_available(true); + response.set_response_code(ResponseCode::NoError); + for q in request.queries().iter() { + response.add_query(q.clone()); + if q.query_type() == RecordType::A { + let record = Record::from_rdata( + q.name().clone(), + 60, + RData::A(A(Ipv4Addr::new(1, 2, 3, 4))), + ); + response.add_answer(record); + } + } + if let Ok(bytes) = response.to_bytes() { + let _ = sock.send_to(&bytes, src).await; + } + }); + } + }); + (addr, slow_rx) + } + + /// Spawn 一个「黑洞」upstream:只 recv_from,永远不应答。 + /// + /// 让每个查询都占住并发 permit 直到 config.timeout_ms 超时 —— + /// 用来观察信号量上限、以及 stop() 的 abort 行为。 + async fn spawn_blackhole_upstream() -> SocketAddr { + let socket = UdpSocket::bind("127.0.0.1:0").await.unwrap(); + let addr = socket.local_addr().unwrap(); + tokio::spawn(async move { + let mut buf = vec![0u8; 512]; + while socket.recv_from(&mut buf).await.is_ok() {} + }); + addr + } + #[tokio::test] async fn test_dns_server_start_stop() { let config = DnsConfig { @@ -1041,6 +1246,182 @@ mod tests { assert!(!server.is_running()); } + /// 回归测试(issue #223):一条慢上游查询不能阻塞其他查询。 + /// + /// 修复前主循环把 `handle_dns_request(...).await` 内联在 recv_from 分支里, + /// 一条撞上慢上游的查询会把 recv_from 卡住不让 poll,后续查询全部排队 —— + /// 最坏吞吐 ≈ 1/上游 RTT。 + /// + /// 夹具:同一个 mock upstream,`slow.example.com` 睡 1200ms 才应答, + /// 其他名字立即应答。先把 slow 查询发出去并**确认它已经打到上游**, + /// 然后计时发一个 fast 查询,断言它在 500ms 内返回。 + /// 串行实现下这里会是 ~1200ms,并发实现下是毫秒级。 + #[tokio::test] + async fn test_dns_server_no_head_of_line_blocking() { + let (upstream, mut slow_seen) = + spawn_name_dependent_upstream("slow.example.com", Duration::from_millis(1200)).await; + + let port = free_udp_port().await; + // timeout_ms 必须 > 1200ms,否则 slow 查询会在测试跑完前自己超时返回, + // 让断言失去意义。 + let config = DnsConfig { + port, + upstream: vec![upstream.to_string()], + timeout_ms: 3000, + ..Default::default() + }; + let server = Arc::new(DnsServer::new(config).unwrap()); + let server_clone = server.clone(); + let handle = tokio::spawn(async move { + server_clone.start().await.unwrap(); + }); + wait_for_server_running(&server, 1000).await; + let server_addr = SocketAddr::from(([127, 0, 0, 1], port)); + + // 先把 slow 查询发出去(不 await 响应),让 server 进入「正在处理 + // 慢上游」的状态。持有 client socket 直到测试结束。 + let slow_client = UdpSocket::bind("0.0.0.0:0").await.unwrap(); + slow_client + .send_to(&build_a_query("slow.example.com"), server_addr) + .await + .unwrap(); + tokio::time::timeout(Duration::from_secs(2), slow_seen.recv()) + .await + .expect("slow query should have reached upstream") + .expect("slow-notify channel should stay open"); + + // 现在 server 里有一个耗时 1200ms 的在途查询。 + let started = Instant::now(); + let response = send_a_query(server_addr, "fast.example.com").await; + let elapsed = started.elapsed(); + + let msg = Message::from_bytes(&response).unwrap(); + assert_eq!(msg.response_code(), ResponseCode::NoError); + assert!( + !msg.answers().is_empty(), + "fast query should resolve via upstream" + ); + assert!( + elapsed < Duration::from_millis(500), + "a slow upstream query must not head-of-line block other queries: \ + fast query took {:?}, expected < 500ms (serial impl. would be ~1200ms)", + elapsed + ); + + server.stop().await.unwrap(); + handle.await.unwrap(); + drop(slow_client); + } + + /// 回归测试(issue #223):并发上限信号量生效。 + /// + /// 黑洞 upstream 让每个查询都占住 permit 直到上游超时。发 2 * MAX 个 + /// 查询后,available_permits() 必须是 0 —— 说明限流真的在卡住 spawn + /// 数量,而不是无上限地长出 2048 个 task。 + #[tokio::test] + async fn test_dns_server_concurrency_capped() { + let upstream = spawn_blackhole_upstream().await; + let port = free_udp_port().await; + let config = DnsConfig { + port, + upstream: vec![upstream.to_string()], + timeout_ms: 3000, + ..Default::default() + }; + let server = Arc::new(DnsServer::new(config).unwrap()); + let server_clone = server.clone(); + let handle = tokio::spawn(async move { + server_clone.start().await.unwrap(); + }); + wait_for_server_running(&server, 1000).await; + let server_addr = SocketAddr::from(([127, 0, 0, 1], port)); + + // 分批发送:一次性 burst 2048 个包会被内核 UDP 接收缓冲吃掉一部分 + // (macOS 默认 rcvspace 只有几 KB,一个 DNS query 约 40 字节), + // 导致实际到达的包不足 1024 个、断言 available_permits()==0 变成 + // 偶发失败。分批 + 小 sleep 让 server 的 recv 循环持续排空。 + // + let payload = build_a_query("flood.example.com"); + let total = MAX_CONCURRENT_CLIENT_QUERIES * 2; + for _ in 0..total / 64 { + for _ in 0..64 { + // socket 用完即丢,不存进 Vec 持有 2048 个:send_to 返回时 + // 数据报已经交给内核,关闭发送方不会把它撤回来,没必要占着 fd。 + // (本测试 fd 峰值 ~1026,主要来自 hickory 为每个在途查询开的 + // 上游 socket —— 和已有的 proxy 并发测试同量级,后者峰值 + // ~3074 且一直是 CI 绿的。) + let sock = UdpSocket::bind("0.0.0.0:0").await.unwrap(); + sock.send_to(&payload, server_addr).await.unwrap(); + } + tokio::time::sleep(Duration::from_millis(5)).await; + } + + // 给首批 task 时间占满 permit + tokio::time::sleep(Duration::from_millis(800)).await; + + assert_eq!( + server.available_permits(), + 0, + "concurrency cap must be saturated by the flood (started from {})", + MAX_CONCURRENT_CLIENT_QUERIES + ); + assert!(server.is_running(), "server must survive the flood"); + + server.stop().await.unwrap(); + handle.await.unwrap(); + } + + /// 回归测试(issue #223):stop() 不被在途查询拖住。 + /// + /// 黑洞 upstream + 在途查询时,shutdown 分支走 `abort_all()` 而不是 + /// 等待 task 跑完 —— 否则 stop() 最坏要等 3000ms(上游超时), + /// 会把「关闭 DNS 模式」和「应用退出」全部卡住(见 AGENTS.md 的 + /// exit / DNS cleanup 契约)。 + #[tokio::test] + async fn test_dns_server_stop_aborts_inflight() { + let upstream = spawn_blackhole_upstream().await; + let port = free_udp_port().await; + let config = DnsConfig { + port, + upstream: vec![upstream.to_string()], + timeout_ms: 3000, + ..Default::default() + }; + let server = Arc::new(DnsServer::new(config).unwrap()); + let server_clone = server.clone(); + let handle = tokio::spawn(async move { + server_clone.start().await.unwrap(); + }); + wait_for_server_running(&server, 1000).await; + let server_addr = SocketAddr::from(([127, 0, 0, 1], port)); + + for name in ["a.example.com", "b.example.com", "c.example.com"] { + let sock = UdpSocket::bind("0.0.0.0:0").await.unwrap(); + sock.send_to(&build_a_query(name), server_addr) + .await + .unwrap(); + } // socket 用完即丢,见 concurrency 测试里的说明 + // 让查询真正进入在途状态(permit 被占用) + tokio::time::sleep(Duration::from_millis(200)).await; + assert!( + server.available_permits() < MAX_CONCURRENT_CLIENT_QUERIES, + "queries should be in flight before stop()" + ); + + let started = Instant::now(); + server.stop().await.unwrap(); + let elapsed = started.elapsed(); + + assert!( + elapsed < Duration::from_millis(500), + "stop() must abort in-flight queries instead of waiting out the \ + upstream timeout, took {:?} (upstream timeout is 3000ms)", + elapsed + ); + assert!(!server.is_running()); + handle.await.unwrap(); + } + #[tokio::test] async fn test_dns_server_resolve_custom_domain() { let profile = make_profile(