fix: control sampling

This commit is contained in:
Lowder
2026-08-17 01:42:29 +05:00
parent 2992f92bfa
commit 5f6dc1697c
4 changed files with 113 additions and 64 deletions
@@ -180,7 +180,9 @@ const verdictStyles = {
</div>
{:else}
<div class="text-md text-neutral-200 mb-2">
Блокировка на ТСПУ <b>не обнаружена</b>
Блокировка на ТСПУ <b>не обнаружена</b> после
{probe.target_hop}
прыжков
</div>
{/if}
<div class="grid grid-cols-1 md:grid-cols-2 gap-4">
@@ -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'}
Ошибка подключения
+3 -1
View File
@@ -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-образ настраивают её автоматически.
## Как работает проверка
+35 -11
View File
@@ -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);
+70 -50
View File
@@ -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<TcpTracerouteResult> {
tokio::task::spawn_blocking(move || trace_blocking(target, max_hops))
pub async fn tcp_traceroute(
target: IpAddr,
max_hops: u8,
retries: u8,
) -> Option<TcpTracerouteResult> {
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<TcpTracerout
.ok()
}
fn trace_blocking(target: IpAddr, max_hops: u8) -> io::Result<TcpTracerouteResult> {
fn trace_blocking(target: IpAddr, max_hops: u8, retries: u8) -> io::Result<TcpTracerouteResult> {
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<TcpTracerouteResul
let mut last_icmp_hop = None;
for ttl in 1..=max_hops {
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)?,
}
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<HopResponse> {
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<HopResponse> {
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))?,