From 5f6dc1697c33af577b8ed209068917657363ed59 Mon Sep 17 00:00:00 2001 From: Lowder Date: Mon, 17 Aug 2026 01:42:29 +0500 Subject: [PATCH] fix: control sampling --- .../lib/components/result/ProbeTable.svelte | 7 +- probe/README.md | 4 +- probe/src/main.rs | 46 +++++-- probe/src/traceroute.rs | 120 ++++++++++-------- 4 files changed, 113 insertions(+), 64 deletions(-) diff --git a/frontend/src/lib/components/result/ProbeTable.svelte b/frontend/src/lib/components/result/ProbeTable.svelte index 0144b77..14eced3 100644 --- a/frontend/src/lib/components/result/ProbeTable.svelte +++ b/frontend/src/lib/components/result/ProbeTable.svelte @@ -180,7 +180,9 @@ const verdictStyles = { {:else}
- Блокировка на ТСПУ не обнаружена + Блокировка на ТСПУ не обнаружена после + {probe.target_hop} + прыжков
{/if}
@@ -200,7 +202,8 @@ const verdictStyles = { {:else if host.probe_evidence.type === 'ClientHello'} Блокировка после ClientHello {:else if host.probe_evidence.type === 'DataTimeout'} - Таймаут получения данных, получено{host.probe_evidence.bytes} + Таймаут получения данных, получено + {host.probe_evidence.bytes} байт {:else if host.probe_evidence.type === 'ConnectionError'} Ошибка подключения diff --git a/probe/README.md b/probe/README.md index 745416d..e1b7579 100644 --- a/probe/README.md +++ b/probe/README.md @@ -123,9 +123,11 @@ docker run --rm \ | `--probe-token`, `PROBE_TOKEN` | Секретный токен сканера. | обязательно | | `--max-concurrent-tasks`, `MAX_CONCURRENT_TASKS` | Максимальное количество одновременных заданий. | `8` | | `--traceroute-max-hops`, `TRACEROUTE_MAX_HOPS` | Максимальный TTL для TCP traceroute. | `5` | +| `--traceroute-retries`, `TRACEROUTE_RETRIES` | Количество одновременных TCP-попыток на каждом TTL. | `3` | +| `--traceroute-control-hosts`, `TRACEROUTE_CONTROL_HOSTS` | Максимальное количество случайных контрольных IP для одновременной трассировки. | `3` | | `RUST_LOG` | Уровень логирования. | `info` | -`MAX_CONCURRENT_TASKS` и `TRACEROUTE_MAX_HOPS` должны быть больше нуля. Для получения ICMP-ответов traceroute процессу требуется capability `CAP_NET_RAW`; systemd unit и Docker-образ настраивают её автоматически. +`MAX_CONCURRENT_TASKS`, `TRACEROUTE_MAX_HOPS`, `TRACEROUTE_RETRIES` и `TRACEROUTE_CONTROL_HOSTS` должны быть больше нуля. Для получения ICMP-ответов traceroute процессу требуется capability `CAP_NET_RAW`; systemd unit и Docker-образ настраивают её автоматически. ## Как работает проверка diff --git a/probe/src/main.rs b/probe/src/main.rs index 4afee48..348104a 100644 --- a/probe/src/main.rs +++ b/probe/src/main.rs @@ -3,9 +3,10 @@ mod traceroute; use anyhow::{Context, Result, bail}; use clap::Parser; +use futures::future::join_all; use log::{error, info, warn}; use rand::seq::SliceRandom; -use reports::probe::{ProbeConfig, ProbeResult, ProbeStatus, ProbeTask}; +use reports::probe::{ProbeConfig, ProbeResult, ProbeStatus, ProbeTask, TcpTracerouteOutcome}; use rumqttc::{ AsyncClient, Event, Incoming, LastWill, MqttOptions, NetworkOptions, QoS, Transport, }; @@ -47,6 +48,12 @@ struct Args { #[arg(long, env = "TRACEROUTE_MAX_HOPS", default_value_t = 5)] traceroute_max_hops: u8, + + #[arg(long, env = "TRACEROUTE_RETRIES", default_value_t = 3)] + traceroute_retries: u8, + + #[arg(long, env = "TRACEROUTE_CONTROL_HOSTS", default_value_t = 3)] + traceroute_control_hosts: usize, } #[tokio::main] @@ -59,6 +66,12 @@ async fn main() -> Result<()> { if args.traceroute_max_hops == 0 { bail!("traceroute_max_hops must be greater than zero"); } + if args.traceroute_retries == 0 { + bail!("traceroute_retries must be greater than zero"); + } + if args.traceroute_control_hosts == 0 { + bail!("traceroute_control_hosts must be greater than zero"); + } let status_topic = format!("probe/status/v1/{}", args.probe_id); let offline_status = serde_json::to_vec(&ProbeStatus { @@ -258,19 +271,21 @@ async fn handle_task( let result_topic = format!("probe/results/v1/{job_id}/{}", args.probe_id); let config = config.read().await.clone(); - let control_target = config.as_ref().and_then(|config| { + let control_targets = config.as_ref().map_or_else(Vec::new, |config| { let mut rng = rand::thread_rng(); match task.ip { IpAddr::V4(_) => config .control_hosts_v4 - .choose(&mut rng) + .choose_multiple(&mut rng, args.traceroute_control_hosts) .copied() - .map(IpAddr::V4), + .map(IpAddr::V4) + .collect(), IpAddr::V6(_) => config .control_hosts_v6 - .choose(&mut rng) + .choose_multiple(&mut rng, args.traceroute_control_hosts) .copied() - .map(IpAddr::V6), + .map(IpAddr::V6) + .collect(), } }); let sni_check = sni::check_sni( @@ -280,12 +295,21 @@ async fn handle_task( job_id, task.timeout_ms, ); - let target_traceroute = traceroute::tcp_traceroute(task.ip, args.traceroute_max_hops); + let target_traceroute = + traceroute::tcp_traceroute(task.ip, args.traceroute_max_hops, args.traceroute_retries); let control_traceroute = async { - match control_target { - Some(target) => traceroute::tcp_traceroute(target, args.traceroute_max_hops).await, - None => None, - } + join_all(control_targets.into_iter().map(|target| { + traceroute::tcp_traceroute(target, args.traceroute_max_hops, args.traceroute_retries) + })) + .await + .into_iter() + .flatten() + .min_by_key(|trace| match trace.result { + TcpTracerouteOutcome::Rst { hop } + | TcpTracerouteOutcome::Connected { hop } + | TcpTracerouteOutcome::IcmpTimeExceeded { hop } => hop, + TcpTracerouteOutcome::Timeout => u8::MAX, + }) }; let (responses, target_traceroute, control_traceroute) = tokio::join!(sni_check, target_traceroute, control_traceroute); diff --git a/probe/src/traceroute.rs b/probe/src/traceroute.rs index 53cac39..d84174e 100644 --- a/probe/src/traceroute.rs +++ b/probe/src/traceroute.rs @@ -12,8 +12,12 @@ use std::time::{Duration, Instant}; const HTTPS_PORT: u16 = 443; const HOP_TIMEOUT: Duration = Duration::from_secs(1); -pub async fn tcp_traceroute(target: IpAddr, max_hops: u8) -> Option { - tokio::task::spawn_blocking(move || trace_blocking(target, max_hops)) +pub async fn tcp_traceroute( + target: IpAddr, + max_hops: u8, + retries: u8, +) -> Option { + tokio::task::spawn_blocking(move || trace_blocking(target, max_hops, retries)) .await .map_err(|error| log::warn!("TCP traceroute task failed for {target}: {error}")) .ok()? @@ -21,7 +25,7 @@ pub async fn tcp_traceroute(target: IpAddr, max_hops: u8) -> Option io::Result { +fn trace_blocking(target: IpAddr, max_hops: u8, retries: u8) -> io::Result { let (domain, icmp_protocol) = match target { IpAddr::V4(_) => (Domain::IPV4, Protocol::ICMPV4), IpAddr::V6(_) => (Domain::IPV6, Protocol::ICMPV6), @@ -30,35 +34,44 @@ fn trace_blocking(target: IpAddr, max_hops: u8) -> io::Result tcp.set_ttl_v4(ttl as u32)?, - IpAddr::V6(_) => tcp.set_unicast_hops_v6(ttl as u32)?, - } - let unspecified = match target { - IpAddr::V4(_) => SocketAddr::new(IpAddr::V4(Ipv4Addr::UNSPECIFIED), 0), - IpAddr::V6(_) => SocketAddr::new(IpAddr::V6(Ipv6Addr::UNSPECIFIED), 0), - }; - tcp.bind(&SockAddr::from(unspecified))?; let destination = SockAddr::from(SocketAddr::new(target, HTTPS_PORT)); - if let Err(error) = tcp.connect(&destination) { - if error.kind() == io::ErrorKind::ConnectionRefused { - return Ok(TcpTracerouteResult { - target, - result: TcpTracerouteOutcome::Rst { hop: ttl }, - }); + let mut tcp_attempts = Vec::with_capacity(retries as usize); + let mut last_error = None; + for _ in 0..retries { + let tcp = Socket::new(domain, Type::STREAM, Some(Protocol::TCP))?; + tcp.set_nonblocking(true)?; + match target { + IpAddr::V4(_) => tcp.set_ttl_v4(ttl as u32)?, + IpAddr::V6(_) => tcp.set_unicast_hops_v6(ttl as u32)?, } - if !is_connect_in_progress(&error) { - return Err(error); + let unspecified = match target { + IpAddr::V4(_) => SocketAddr::new(IpAddr::V4(Ipv4Addr::UNSPECIFIED), 0), + IpAddr::V6(_) => SocketAddr::new(IpAddr::V6(Ipv6Addr::UNSPECIFIED), 0), + }; + tcp.bind(&SockAddr::from(unspecified))?; + if let Err(error) = tcp.connect(&destination) { + if error.kind() == io::ErrorKind::ConnectionRefused { + return Ok(TcpTracerouteResult { + target, + result: TcpTracerouteOutcome::Rst { hop: ttl }, + }); + } + if !is_connect_in_progress(&error) { + last_error = Some(error); + continue; + } } + let source_port = tcp + .local_addr()? + .as_socket() + .map(|addr| addr.port()) + .unwrap_or(0); + tcp_attempts.push((tcp, source_port)); } - let source_port = tcp - .local_addr()? - .as_socket() - .map(|addr| addr.port()) - .unwrap_or(0); - match wait_for_hop_response(&receiver, &tcp, target, source_port)? { + if tcp_attempts.is_empty() { + return Err(last_error.unwrap_or_else(|| io::Error::other("no traceroute attempts"))); + } + match wait_for_hop_response(&receiver, &tcp_attempts, target)? { HopResponse::IcmpTimeExceeded => last_icmp_hop = Some(ttl), HopResponse::Rst => { return Ok(TcpTracerouteResult { @@ -98,45 +111,44 @@ enum HopResponse { fn wait_for_hop_response( receiver: &Socket, - tcp: &Socket, + tcp_attempts: &[(Socket, u16)], target: IpAddr, - source_port: u16, ) -> io::Result { const ICMP_KEY: usize = 1; - const TCP_KEY: usize = 2; let poller = Poller::new()?; - // SAFETY: both sockets remain alive and are removed from the poller before this function exits. + // SAFETY: all sockets remain alive until polling and cleanup finish. unsafe { poller.add(receiver, Event::readable(ICMP_KEY))?; - if let Err(error) = poller.add(tcp, Event::writable(TCP_KEY)) { - poller.delete(receiver)?; - return Err(error); + for (index, (tcp, _)) in tcp_attempts.iter().enumerate() { + if let Err(error) = poller.add(tcp, Event::writable(index + 2)) { + let _ = poller.delete(receiver); + for (added, _) in tcp_attempts.iter().take(index) { + let _ = poller.delete(added); + } + return Err(error); + } } } - let result = wait_on_poller(&poller, receiver, tcp, target, source_port); - let tcp_delete = poller.delete(tcp); - let receiver_delete = poller.delete(receiver); - let cleanup = tcp_delete.and(receiver_delete); - match result { - Ok(response) => cleanup.map(|()| response), - Err(error) => Err(error), + let result = wait_on_poller(&poller, receiver, tcp_attempts, target); + for (tcp, _) in tcp_attempts { + let _ = poller.delete(tcp); } + let _ = poller.delete(receiver); + result } fn wait_on_poller( poller: &Poller, receiver: &Socket, - tcp: &Socket, + tcp_attempts: &[(Socket, u16)], target: IpAddr, - source_port: u16, ) -> io::Result { const ICMP_KEY: usize = 1; - const TCP_KEY: usize = 2; let deadline = Instant::now() + HOP_TIMEOUT; - let mut watch_tcp = true; + let mut watching_tcp = vec![true; tcp_attempts.len()]; let mut buffer = [0u8; 2048]; let mut events = Events::new(); @@ -153,20 +165,28 @@ fn wait_on_poller( Err(error) => return Err(error), } - if watch_tcp && events.iter().any(|event| event.key == TCP_KEY) { - match tcp.take_error()? { + for event in events.iter().filter(|event| event.key >= 2) { + let index = event.key - 2; + if !watching_tcp.get(index).copied().unwrap_or(false) { + continue; + } + match tcp_attempts[index].0.take_error()? { Some(error) if error.kind() == io::ErrorKind::ConnectionRefused => { return Ok(HopResponse::Rst); } None => return Ok(HopResponse::Connected), - Some(_) => watch_tcp = false, + Some(_) => watching_tcp[index] = false, } } if events.iter().any(|event| event.key == ICMP_KEY) { let mut raw = receiver; match raw.read(&mut buffer) { - Ok(bytes) if is_matching_time_exceeded(&buffer[..bytes], target, source_port) => { + Ok(bytes) + if tcp_attempts.iter().any(|(_, source_port)| { + is_matching_time_exceeded(&buffer[..bytes], target, *source_port) + }) => + { return Ok(HopResponse::IcmpTimeExceeded); } Ok(_) => poller.modify(receiver, Event::readable(ICMP_KEY))?,