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))?,