From 8651d8b1ffb38cbfe0a83778ee25809f06e980bc Mon Sep 17 00:00:00 2001 From: LowderPlay Date: Sat, 22 Aug 2026 02:03:09 +0500 Subject: [PATCH] feat: safe dpi traceroutes (#84) * feat: safe dpi traceroutes * chore: bump versions * feat: dynamic verdict * fix: change titles * refactor: simplify methods * fix: block message * fix: narrow cdn block * fix: title --- Cargo.lock | 4 +- frontend/src/lib/api/probe.ts | 68 ++- .../src/lib/components/ResultPanel.svelte | 20 +- .../lib/components/result/ProbeTable.svelte | 164 +++--- .../result/ResultStatusHeader.svelte | 70 ++- frontend/src/routes/check/+page.svelte | 13 +- frontend/src/routes/kb/probing/+page.svelte | 55 +- probe/Cargo.toml | 2 +- probe/README.md | 6 +- probe/openwrt/cheburprobe.config | 3 - probe/openwrt/cheburprobe.init | 4 - probe/openwrt/luci/config.js | 8 - probe/src/dpi_hop.rs | 488 ++++++++++++++++++ probe/src/main.rs | 305 +++++++---- probe/src/traceroute.rs | 43 +- reports/src/probe.rs | 31 +- website/Cargo.toml | 2 +- ...260821000000_remove_control_traceroute.sql | 3 + website/probe-hosts.toml | 10 +- website/src/api/probe.rs | 127 +++-- website/src/mqtt.rs | 19 +- 21 files changed, 1093 insertions(+), 352 deletions(-) create mode 100644 probe/src/dpi_hop.rs create mode 100644 website/migrations/20260821000000_remove_control_traceroute.sql diff --git a/Cargo.lock b/Cargo.lock index d6e1f82..069eefb 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2632,7 +2632,7 @@ dependencies = [ [[package]] name = "probe" -version = "0.2.0" +version = "0.3.0" dependencies = [ "anyhow", "clap", @@ -4706,7 +4706,7 @@ dependencies = [ [[package]] name = "website" -version = "1.2.1" +version = "1.2.2" dependencies = [ "dotenvy", "env_logger", diff --git a/frontend/src/lib/api/probe.ts b/frontend/src/lib/api/probe.ts index 7ca7afc..4f22c3a 100644 --- a/frontend/src/lib/api/probe.ts +++ b/frontend/src/lib/api/probe.ts @@ -26,22 +26,27 @@ export type DnsObservation = { | { type: "Error"; message: string }; }; +export type ProbeVerdict = + | "uncertain" + | "dns_spoofing" + | "sni_block" + | "tspu_block" + | "whitelist" + | "ok"; + +export type DisplayProbeVerdict = ProbeVerdict | "cdn_block"; +export type ResolvedProbeVerdict = Exclude; + export type ProbeResult = { job_id: string; probe_id: string; region?: string | null; provider?: string | null; asn?: string | null; - verdicts: ( - | "uncertain" - | "dns_spoofing" - | "sni_block" - | "tspu_block" - | "whitelist" - | "ok" - )[]; + verdicts: ProbeVerdict[]; host_results: ProbeHostResult[] | null; target_hop: number | null; + dpi_hop: number | null; dns: { spoofing_detected: boolean; suspicious_provider_count: number; @@ -51,6 +56,53 @@ export type ProbeResult = { } | null; }; +export function displayProbeVerdicts( + probe: ProbeResult, + isStaticCdn: boolean, +): DisplayProbeVerdict[] { + return probe.verdicts.map((verdict) => + isStaticCdn && probe.host_results?.length !== 0 && verdict === "ok" + ? "cdn_block" + : verdict, + ); +} + +const verdictPriority: DisplayProbeVerdict[] = [ + "tspu_block", + "sni_block", + "dns_spoofing", + "whitelist", + "cdn_block", + "ok", + "uncertain", +]; + +export function selectProbeVerdict( + probes: ProbeResult[], + isStaticCdn: boolean, +): ResolvedProbeVerdict | null { + if (probes.length === 0) return null; + + const votes = new Map(); + for (const probe of probes) { + for (const verdict of new Set(displayProbeVerdicts(probe, isStaticCdn))) { + votes.set(verdict, (votes.get(verdict) ?? 0) + 1); + } + } + + let winner: DisplayProbeVerdict | null = null; + let winningVotes = 0; + for (const verdict of verdictPriority) { + const count = votes.get(verdict) ?? 0; + if (count > winningVotes) { + winner = verdict; + winningVotes = count; + } + } + + return winner === "uncertain" ? null : winner; +} + export type ProbeStatus = { id: string; target: string; diff --git a/frontend/src/lib/components/ResultPanel.svelte b/frontend/src/lib/components/ResultPanel.svelte index 9406865..62ec845 100644 --- a/frontend/src/lib/components/ResultPanel.svelte +++ b/frontend/src/lib/components/ResultPanel.svelte @@ -13,6 +13,7 @@ import { ShieldCheck, } from "@lucide/svelte"; import type { CheckResult } from "$lib/api/check"; +import type { ResolvedProbeVerdict } from "$lib/api/probe"; import DetailRow from "./DetailRow.svelte"; import ResultIpList from "./result/ResultIpList.svelte"; import ResultStatusHeader from "./result/ResultStatusHeader.svelte"; @@ -21,23 +22,34 @@ import ResultTargetCard from "./result/ResultTargetCard.svelte"; type Provider = { name: string; networks: { cidr: string }[] }; type ResultTheme = "blocked" | "clean" | "whitelist"; +type ResultVerdict = ResolvedProbeVerdict | "blocked"; let { result, + probeVerdict = null, }: { result: CheckResult; + probeVerdict?: ResolvedProbeVerdict | null; } = $props(); const valueClass = "text-right text-sm font-medium text-neutral-200"; const alertValueClass = `${valueClass} text-red-500`; const successValueClass = `${valueClass} text-green-500`; -const theme = $derived( +const staticVerdict = $derived( result.whitelist && !result.domain ? "whitelist" : result.found ? "blocked" - : "clean", + : "ok", +); +const verdict = $derived(probeVerdict ?? staticVerdict); +const theme = $derived( + verdict === "whitelist" + ? "whitelist" + : verdict === "ok" + ? "clean" + : "blocked", ); const panelClass = $derived( theme === "whitelist" @@ -67,7 +79,7 @@ const providerCidrs = (provider: Provider) =>
- +
{#if result.whitelist} diff --git a/frontend/src/lib/components/result/ProbeTable.svelte b/frontend/src/lib/components/result/ProbeTable.svelte index 4560f74..2ce1270 100644 --- a/frontend/src/lib/components/result/ProbeTable.svelte +++ b/frontend/src/lib/components/result/ProbeTable.svelte @@ -8,17 +8,23 @@ import { CircleX, LoaderCircle, ShieldCheck, + TriangleAlert, } from "@lucide/svelte"; -import type { DnsObservation, ProbeResult, ProbeStatus } from "$lib/api/probe"; +import { + type DnsObservation, + displayProbeVerdicts, + type ProbeResult, + type ProbeStatus, +} from "$lib/api/probe"; let { probes, status, - isStaticBlocked, + isStaticCdn, }: { probes: ProbeResult[]; status: ProbeStatus; - isStaticBlocked: boolean; + isStaticCdn: boolean; } = $props(); let expandedRows = $state>({}); @@ -110,19 +116,6 @@ const verdictStyles = { border: "border-neutral-400/20", }, }; - -type VerdictStyle = keyof typeof verdictStyles; - -function probeVerdicts( - probe: ProbeResult, - isStaticBlocked: boolean, -): VerdictStyle[] { - return probe.verdicts.map((verdict) => - isStaticBlocked && probe.host_results?.length !== 0 && verdict === "ok" - ? "cdn_block" - : verdict, - ); -}
@@ -182,7 +175,7 @@ function probeVerdicts( {#each probes as probe (probe.probe_id)} - {@const verdicts = probeVerdicts(probe, isStaticBlocked)} + {@const verdicts = displayProbeVerdicts(probe, isStaticCdn)} {@const isExpanded = !!expandedRows[probe.probe_id]} {#each verdicts as verdict} {@const style = verdictStyles[verdict]} - {#if verdict === "whitelist"} - event.stopPropagation()} - class={`inline-flex items-center gap-1.5 px-2 py-1 rounded border ${style.bg} ${style.border} ${style.class} text-xs font-bold transition-colors hover:bg-amber-500/20 hover:border-amber-500/50 hover:text-amber-400`} - > - - {style.text} - - {:else} -
- - {style.text} -
- {/if} + event.stopPropagation()} + class={`inline-flex items-center gap-1.5 px-2 py-1 rounded border ${style.bg} ${style.border} ${style.class} text-xs font-bold transition-opacity hover:opacity-80`} + > + + {style.text} + {/each}
@@ -234,27 +218,27 @@ function probeVerdicts( {#if isExpanded} - + {/if} {#if probe.dns}
-
+

@@ -331,53 +315,57 @@ function probeVerdicts(

{/if} -
-

- CDN-проверка -

-
-
- {#each probe.host_results as host} -
-
- - Сервер {host.host_id} - ({host.host === "Blacklist" ? "в заблокированных" : "в доступных"} - диапазонах) - - - {#if host.probe_evidence.type === 'Good'} - Успешно - {:else if host.probe_evidence.type === 'ClientHello'} - Блокировка после ClientHello - {:else if host.probe_evidence.type === 'DataTimeout'} - Таймаут получения данных, получено - {host.probe_evidence.bytes} - байт - {:else if host.probe_evidence.type === 'ConnectionError'} - Ошибка подключения - {/if} - + CDN-проверка + +
+
+ {#each probe.host_results as host} +
+
+ + Сервер {host.host_id} + ({host.host === "Blacklist" ? "в заблокированных" : "в доступных"} + диапазонах) + + + {#if host.probe_evidence.type === 'Good'} + Успешно + {:else if host.probe_evidence.type === 'ClientHello'} + Блокировка после ClientHello + {:else if host.probe_evidence.type === 'DataTimeout'} + Таймаут получения данных, получено + {host.probe_evidence.bytes} + байт + {:else if host.probe_evidence.type === 'ConnectionError'} + Ошибка подключения + {/if} + +
+ {#if host.probe_evidence.type === 'Good'} + + {:else if host.probe_evidence.type === 'ClientHello'} + + {:else if host.probe_evidence.type === 'DataTimeout'} + + {:else} + + {/if}
- {#if host.probe_evidence.type === 'Good'} - - {:else if host.probe_evidence.type === 'ClientHello'} - - {:else if host.probe_evidence.type === 'DataTimeout'} - - {:else} - - {/if} -
- {/each} -
+ {/each} +
+ {/if} {/if} diff --git a/frontend/src/lib/components/result/ResultStatusHeader.svelte b/frontend/src/lib/components/result/ResultStatusHeader.svelte index e83feed..76ad5c0 100644 --- a/frontend/src/lib/components/result/ResultStatusHeader.svelte +++ b/frontend/src/lib/components/result/ResultStatusHeader.svelte @@ -1,39 +1,61 @@ diff --git a/frontend/src/routes/check/+page.svelte b/frontend/src/routes/check/+page.svelte index 94c0b73..62a4b44 100644 --- a/frontend/src/routes/check/+page.svelte +++ b/frontend/src/routes/check/+page.svelte @@ -6,6 +6,7 @@ import { CheckRequestError, fetchCheck } from "$lib/api/check"; import { type ProbeResult, type ProbeStatus, + selectProbeVerdict, startProbeSSE, } from "$lib/api/probe"; import EmptyResult from "$lib/components/EmptyResult.svelte"; @@ -104,6 +105,14 @@ $effect(() => { const error = $derived( (checkQuery.error instanceof CheckRequestError && checkQuery.error) || null, ); +const liveVerdict = $derived( + checkQuery.data && probeQuery.data + ? selectProbeVerdict( + probeQuery.data.probes, + checkQuery.data.providers.length > 0, + ) + : null, +); @@ -121,13 +130,13 @@ const error = $derived(
{:else if checkQuery.data} - + {#if shouldProbe && probeQuery.data && probeQuery.data.status.online_probes > 0} 0} /> {/if} diff --git a/frontend/src/routes/kb/probing/+page.svelte b/frontend/src/routes/kb/probing/+page.svelte index a9faca5..ad89c72 100644 --- a/frontend/src/routes/kb/probing/+page.svelte +++ b/frontend/src/routes/kb/probing/+page.svelte @@ -61,25 +61,36 @@ import KbNote from "$lib/components/kb/KbNote.svelte"; завершить TLS-обмен и получить минимальный объём данных. Для проверки IP-адреса без домена этот этап пропускается.

- - + При проверке целевого IP-адреса TCP-трассировка начинается со следующего + после DPI прыжка — dpi_hop + 1. По умолчанию проверяются три + прыжка, на каждом из которых одновременно выполняются три попытки. Первый + ответ ICMP Time Exceeded, TCP RST или успешное соединение немедленно + завершает трассировку. + + +

+ Если DPI был измерен, но ни одна попытка после него не получила ответа, + сканер выставляет вердикт «ТСПУ Блок» и показывает измеренный прыжок DPI. + Любой ответ после DPI означает, что пакеты прошли дальше, поэтому блокировка + IP-адреса на ТСПУ не подтверждается. Если прыжок, на котором установлен DPI, + для нужной версии IP не удалось определить, трассировка пропускается и этот + вердикт не выставляется. +

+ + Измерение выполняется независимо для IPv4 и IPv6. Разрыв TCP-соединения во + время калибровки делает результат этой версии IP недействительным, чтобы не + принимать поведение контрольного сервера за работу DPI. + Проверка DNS

@@ -137,12 +148,12 @@ import KbNote from "$lib/components/kb/KbNote.svelte"; получает SNI Блок. - + – TCP-трассировка, начатая сразу после прыжка с DPI, не получила ни + ICMP-ответа, ни TCP RST и не смогла подключиться к цели. Это указывает на + возможную блокировку IP-адреса на оборудовании оператора. +

  • , + pub hops: Vec, +} + +#[derive(Debug, Clone)] +pub struct DpiHopProbeHop { + pub ttl: u8, + pub router: Option, + pub outcome: DpiHopProbeHopOutcome, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum DpiHopProbeHopOutcome { + IcmpTimeExceeded, + Timeout, + TcpClosed, +} + +pub async fn detect_dpi_hop(config: DpiHopProbeConfig) -> io::Result { + tokio::task::spawn_blocking(move || detect_dpi_hop_blocking(config)) + .await + .map_err(io::Error::other)? +} + +pub fn detect_dpi_hop_blocking(config: DpiHopProbeConfig) -> io::Result { + if config.max_ttl == 0 { + return Err(io::Error::new( + io::ErrorKind::InvalidInput, + "max_ttl must be greater than zero", + )); + } + + let client_hello = make_client_hello(&config.control_sni)?; + let (domain, protocol) = match config.target { + SocketAddr::V4(_) => (Domain::IPV4, Protocol::ICMPV4), + SocketAddr::V6(_) => (Domain::IPV6, Protocol::ICMPV6), + }; + let icmp = Socket::new(domain, Type::RAW, Some(protocol))?; + icmp.set_read_timeout(Some(config.hop_timeout))?; + + let mut tcp = TcpStream::connect_timeout(&config.target, config.connect_timeout)?; + tcp.set_nodelay(true)?; + let local_addr = tcp.local_addr()?; + if !same_ip_family(local_addr, config.target) { + return Err(io::Error::other( + "DPI hop probe local and target address families differ", + )); + } + + tcp.write_all(&client_hello)?; + let mut hops = Vec::with_capacity(config.max_ttl as usize); + let mut max_icmp_time_exceeded_ttl = None; + + for ttl in 1..=config.max_ttl { + let mut payload = [0u8; PROBE_BYTES]; + rand::thread_rng().fill_bytes(&mut payload); + drain_socket(&icmp)?; + + if !send_with_ttl(&mut tcp, config.target, &payload, ttl)? { + hops.push(DpiHopProbeHop { + ttl, + router: None, + outcome: DpiHopProbeHopOutcome::TcpClosed, + }); + break; + } + + let router = + listen_for_time_exceeded(&icmp, local_addr, config.target, config.hop_timeout)?; + let outcome = if router.is_some() { + max_icmp_time_exceeded_ttl = Some(ttl); + DpiHopProbeHopOutcome::IcmpTimeExceeded + } else if peer_closed(&tcp)? { + DpiHopProbeHopOutcome::TcpClosed + } else { + DpiHopProbeHopOutcome::Timeout + }; + hops.push(DpiHopProbeHop { + ttl, + router, + outcome, + }); + if outcome == DpiHopProbeHopOutcome::TcpClosed { + break; + } + } + + Ok(DpiHopProbeResult { + target: config.target, + local_addr, + client_hello_bytes: client_hello.len(), + max_icmp_time_exceeded_ttl, + hops, + }) +} + +fn make_client_hello(sni: &str) -> io::Result> { + let server_name = ServerName::try_from(sni.to_owned()).map_err(|_| { + io::Error::new( + io::ErrorKind::InvalidInput, + "control_sni is not a valid DNS name", + ) + })?; + let config = ClientConfig::builder() + .with_root_certificates(RootCertStore::empty()) + .with_no_client_auth(); + let mut connection = ClientConnection::new(Arc::new(config), server_name) + .map_err(|error| io::Error::other(format!("create TLS client connection: {error}")))?; + let mut client_hello = Vec::new(); + connection + .write_tls(&mut client_hello) + .map_err(|error| io::Error::other(format!("write TLS ClientHello: {error}")))?; + if client_hello.is_empty() { + return Err(io::Error::other("rustls did not produce a ClientHello")); + } + Ok(client_hello) +} + +fn same_ip_family(left: SocketAddr, right: SocketAddr) -> bool { + matches!( + (left, right), + (SocketAddr::V4(_), SocketAddr::V4(_)) | (SocketAddr::V6(_), SocketAddr::V6(_)) + ) +} + +fn send_with_ttl( + tcp: &mut TcpStream, + target: SocketAddr, + payload: &[u8], + ttl: u8, +) -> io::Result { + let previous_ttl = { + let socket = SockRef::from(&*tcp); + match target { + SocketAddr::V4(_) => socket.ttl_v4()?, + SocketAddr::V6(_) => socket.unicast_hops_v6()?, + } + }; + set_ttl(tcp, target, ttl as u32)?; + let write_result = tcp.write_all(payload); + let restore_result = set_ttl(tcp, target, previous_ttl); + match write_result { + Ok(()) => { + restore_result?; + Ok(true) + } + Err(error) + if matches!( + error.kind(), + io::ErrorKind::ConnectionReset | io::ErrorKind::BrokenPipe + ) => + { + let _ = restore_result; + Ok(false) + } + Err(error) => { + let _ = restore_result; + Err(error) + } + } +} + +fn set_ttl(tcp: &TcpStream, target: SocketAddr, ttl: u32) -> io::Result<()> { + let socket = SockRef::from(tcp); + match target { + SocketAddr::V4(_) => socket.set_ttl_v4(ttl), + SocketAddr::V6(_) => socket.set_unicast_hops_v6(ttl), + } +} + +fn listen_for_time_exceeded( + icmp: &Socket, + local_addr: SocketAddr, + target: SocketAddr, + timeout: Duration, +) -> io::Result> { + let deadline = Instant::now() + timeout; + let mut buffer = [0u8; 65535]; + loop { + let remaining = deadline.saturating_duration_since(Instant::now()); + if remaining.is_zero() { + return Ok(None); + } + icmp.set_read_timeout(Some(remaining))?; + match recv_socket(icmp, &mut buffer) { + Ok((bytes, source)) => { + if let Some(router) = + match_time_exceeded(&buffer[..bytes], source, local_addr, target) + { + return Ok(Some(router)); + } + } + Err(error) + if error.kind() == io::ErrorKind::WouldBlock + || error.kind() == io::ErrorKind::TimedOut => + { + return Ok(None); + } + Err(error) => return Err(error), + } + } +} + +fn peer_closed(tcp: &TcpStream) -> io::Result { + tcp.set_nonblocking(true)?; + let mut byte = [0u8; 1]; + let result = match tcp.peek(&mut byte) { + Ok(0) => Ok(true), + Ok(_) => Ok(false), + Err(error) if error.kind() == io::ErrorKind::WouldBlock => Ok(false), + Err(error) + if matches!( + error.kind(), + io::ErrorKind::ConnectionReset | io::ErrorKind::BrokenPipe + ) => + { + Ok(true) + } + Err(error) => Err(error), + }; + let _ = tcp.set_nonblocking(false); + result +} + +fn drain_socket(socket: &Socket) -> io::Result { + let previous_timeout = socket.read_timeout()?; + socket.set_nonblocking(true)?; + let mut drained = 0; + let mut buffer = [0u8; 65535]; + loop { + match recv_socket(socket, &mut buffer) { + Ok(_) => drained += 1, + Err(error) if error.kind() == io::ErrorKind::WouldBlock => break, + Err(error) => { + let _ = socket.set_nonblocking(false); + let _ = socket.set_read_timeout(previous_timeout); + return Err(error); + } + } + } + socket.set_nonblocking(false)?; + socket.set_read_timeout(previous_timeout)?; + Ok(drained) +} + +fn recv_socket(socket: &Socket, buffer: &mut [u8]) -> io::Result<(usize, Option)> { + // `recv_from` initializes exactly the returned prefix of the buffer. + let uninitialized = unsafe { + std::slice::from_raw_parts_mut(buffer.as_mut_ptr().cast::>(), buffer.len()) + }; + let (bytes, source) = socket.recv_from(uninitialized)?; + Ok((bytes, source.as_socket())) +} + +fn match_time_exceeded( + packet: &[u8], + source: Option, + local_addr: SocketAddr, + target: SocketAddr, +) -> Option { + match (local_addr, target) { + (SocketAddr::V4(local), SocketAddr::V4(target)) => { + match_ipv4_time_exceeded(packet, local, target).map(IpAddr::V4) + } + (SocketAddr::V6(local), SocketAddr::V6(target)) => { + match_ipv6_time_exceeded(packet, source, local, target) + } + _ => None, + } +} + +fn match_ipv4_time_exceeded( + packet: &[u8], + local_addr: SocketAddrV4, + target: SocketAddrV4, +) -> Option { + let Ok(outer) = LaxSlicedPacket::from_ip(packet) else { + return None; + }; + let router = match outer.net { + Some(LaxNetSlice::Ipv4(ip)) => ip.header().source_addr(), + _ => return None, + }; + let Some(TransportSlice::Icmpv4(icmp)) = outer.transport else { + return None; + }; + if matches!( + icmp.icmp_type(), + Icmpv4Type::TimeExceeded(icmpv4::TimeExceededCode::TtlExceededInTransit) + ) && matching_quoted_tcp_tuple( + icmp.payload(), + SocketAddr::V4(local_addr), + SocketAddr::V4(target), + ) { + Some(router) + } else { + None + } +} + +fn match_ipv6_time_exceeded( + packet: &[u8], + source: Option, + local_addr: SocketAddrV6, + target: SocketAddrV6, +) -> Option { + let (icmp, outer_router) = if packet.first().map(|byte| byte >> 4) == Some(6) { + let Ok(outer) = LaxSlicedPacket::from_ip(packet) else { + return None; + }; + let router = match outer.net { + Some(LaxNetSlice::Ipv6(ip)) => IpAddr::V6(ip.header().source_addr()), + _ => return None, + }; + let Some(TransportSlice::Icmpv6(icmp)) = outer.transport else { + return None; + }; + (icmp, Some(router)) + } else { + let Ok(icmp) = Icmpv6Slice::from_slice(packet) else { + return None; + }; + (icmp, None) + }; + if !matches!( + icmp.icmp_type(), + Icmpv6Type::TimeExceeded(icmpv6::TimeExceededCode::HopLimitExceeded) + ) || !matching_quoted_tcp_tuple( + icmp.payload(), + SocketAddr::V6(local_addr), + SocketAddr::V6(target), + ) { + return None; + } + outer_router.or_else(|| match source { + Some(SocketAddr::V6(source)) => Some(IpAddr::V6(*source.ip())), + _ => None, + }) +} + +fn matching_quoted_tcp_tuple(packet: &[u8], local_addr: SocketAddr, target: SocketAddr) -> bool { + let Ok(quoted) = LaxSlicedPacket::from_ip(packet) else { + return false; + }; + let payload = match (quoted.net, local_addr, target) { + (Some(LaxNetSlice::Ipv4(ip)), SocketAddr::V4(local), SocketAddr::V4(target)) + if ip.header().source_addr() == *local.ip() + && ip.header().destination_addr() == *target.ip() + && ip.payload().ip_number == IpNumber::TCP => + { + ip.payload().payload + } + (Some(LaxNetSlice::Ipv6(ip)), SocketAddr::V6(local), SocketAddr::V6(target)) + if ip.header().source_addr() == *local.ip() + && ip.header().destination_addr() == *target.ip() + && ip.payload().ip_number == IpNumber::TCP => + { + ip.payload().payload + } + _ => return false, + }; + let Some(ports) = payload.get(..4) else { + return false; + }; + u16::from_be_bytes([ports[0], ports[1]]) == local_addr.port() + && u16::from_be_bytes([ports[2], ports[3]]) == target.port() +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn matches_icmp_time_exceeded_quote_by_flow_tuple() { + let local = SocketAddrV4::new(Ipv4Addr::new(192, 0, 2, 10), 45_000); + let target = SocketAddrV4::new(Ipv4Addr::new(203, 0, 113, 10), 443); + let router = Ipv4Addr::new(198, 51, 100, 1); + let mut packet = vec![0; 20 + 8 + 20 + 4]; + packet[0] = 0x45; + packet[2..4].copy_from_slice(&52u16.to_be_bytes()); + packet[8] = 64; + packet[9] = IpNumber::ICMP.0; + packet[12..16].copy_from_slice(&router.octets()); + packet[20] = 11; + packet[28] = 0x45; + packet[30..32].copy_from_slice(&24u16.to_be_bytes()); + packet[36] = 1; + packet[37] = IpNumber::TCP.0; + packet[40..44].copy_from_slice(&local.ip().octets()); + packet[44..48].copy_from_slice(&target.ip().octets()); + packet[48..50].copy_from_slice(&local.port().to_be_bytes()); + packet[50..52].copy_from_slice(&target.port().to_be_bytes()); + + assert_eq!( + match_time_exceeded(&packet, None, SocketAddr::V4(local), SocketAddr::V4(target)), + Some(IpAddr::V4(router)) + ); + } + + #[test] + fn rejects_icmp_time_exceeded_for_another_port() { + let local = SocketAddrV4::new(Ipv4Addr::new(192, 0, 2, 10), 45_000); + let target = SocketAddrV4::new(Ipv4Addr::new(203, 0, 113, 10), 443); + let mut packet = vec![0; 20 + 8 + 20 + 4]; + packet[0] = 0x45; + packet[2..4].copy_from_slice(&52u16.to_be_bytes()); + packet[9] = IpNumber::ICMP.0; + packet[20] = 11; + packet[28] = 0x45; + packet[30..32].copy_from_slice(&24u16.to_be_bytes()); + packet[37] = IpNumber::TCP.0; + packet[40..44].copy_from_slice(&local.ip().octets()); + packet[44..48].copy_from_slice(&target.ip().octets()); + packet[48..50].copy_from_slice(&local.port().to_be_bytes()); + packet[50..52].copy_from_slice(&444u16.to_be_bytes()); + + assert_eq!( + match_time_exceeded(&packet, None, SocketAddr::V4(local), SocketAddr::V4(target)), + None + ); + } + + #[test] + fn matches_ipv6_icmp_time_exceeded_quote_by_flow_tuple() { + let local = "[2001:db8::1]:45000".parse::().unwrap(); + let target = "[2001:db8::10]:443".parse::().unwrap(); + let router = "2001:db8::ff".parse::().unwrap(); + let mut packet = vec![0; 8 + 40 + 4]; + packet[0] = 3; + packet[8] = 0x60; + packet[14] = IpNumber::TCP.0; + packet[16..32].copy_from_slice(&local.ip().octets()); + packet[32..48].copy_from_slice(&target.ip().octets()); + packet[48..50].copy_from_slice(&local.port().to_be_bytes()); + packet[50..52].copy_from_slice(&target.port().to_be_bytes()); + + assert_eq!( + match_time_exceeded( + &packet, + Some(SocketAddr::V6(SocketAddrV6::new(router, 0, 0, 0))), + SocketAddr::V6(local), + SocketAddr::V6(target) + ), + Some(IpAddr::V6(router)) + ); + + packet[50..52].copy_from_slice(&444u16.to_be_bytes()); + assert_eq!( + match_time_exceeded( + &packet, + Some(SocketAddr::V6(SocketAddrV6::new(router, 0, 0, 0))), + SocketAddr::V6(local), + SocketAddr::V6(target) + ), + None + ); + } +} diff --git a/probe/src/main.rs b/probe/src/main.rs index 9c9f9c9..17b752e 100644 --- a/probe/src/main.rs +++ b/probe/src/main.rs @@ -1,29 +1,34 @@ mod dns; +mod dpi_hop; mod sni; 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, TcpTracerouteOutcome}; +use log::{debug, error, info, warn}; +use reports::probe::{DpiProbeConfig, ProbeConfig, ProbeResult, ProbeStatus, ProbeTask}; use rumqttc::{ AsyncClient, Event, Incoming, LastWill, MqttOptions, NetworkOptions, QoS, Transport, }; -use std::collections::HashSet; -use std::net::{IpAddr, Ipv4Addr, Ipv6Addr}; +use std::net::IpAddr; use std::sync::Arc; use std::time::{Duration, Instant}; use tokio::sync::RwLock; const CONFIG_TOPIC: &str = "probe/config/v1"; +const MQTT_MAX_PACKET_SIZE: usize = 1024 * 1024; #[derive(Clone)] struct LoadedProbeConfig { config: ProbeConfig, - control_hosts_v4: Vec, - control_hosts_v6: Vec, + dpi_hop_v4: Option, + dpi_hop_v6: Option, +} + +#[derive(Clone, Copy, Default)] +struct DpiHops { + v4: Option, + v6: Option, } #[derive(Parser, Debug, Clone)] @@ -47,14 +52,8 @@ struct Args { #[arg(long, env = "MAX_CONCURRENT_TASKS", default_value_t = 8)] max_concurrent_tasks: usize, - #[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] @@ -64,27 +63,24 @@ async fn main() -> Result<()> { if args.max_concurrent_tasks == 0 { bail!("max_concurrent_tasks must be greater than zero"); } - 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 { online: false, probe_id: &args.probe_id, version: env!("CARGO_PKG_VERSION"), + dpi_hop_v4: None, + dpi_hop_v6: None, })?; let mut options = MqttOptions::new(&args.probe_id, &args.mqtt_host, args.mqtt_port); options.set_transport(mqtt_transport(&args.mqtt_host)?); options.set_credentials("probe", &args.probe_token); options.set_keep_alive(Duration::from_secs(10)); + options.set_max_packet_size(MQTT_MAX_PACKET_SIZE, MQTT_MAX_PACKET_SIZE); options.set_last_will(LastWill::new( status_topic.clone(), offline_status, @@ -100,7 +96,7 @@ async fn main() -> Result<()> { eventloop.set_network_options(network_options); wait_for_connection(&mut eventloop).await; - publish_status(&client, &status_topic, &args, true).await?; + publish_status(&client, &status_topic, &args, true, DpiHops::default()).await?; client.subscribe(CONFIG_TOPIC, QoS::AtLeastOnce).await?; client .subscribe("probe/tasks/v1/+", QoS::AtLeastOnce) @@ -115,8 +111,15 @@ async fn main() -> Result<()> { match eventloop.poll().await { Ok(Event::Incoming(Incoming::Publish(publish))) => { if publish.topic == CONFIG_TOPIC { - if let Err(error) = update_config(&config, &publish.payload).await { - warn!("failed to update probe config: {error}"); + match update_config(&config, &publish.payload).await { + Ok(dpi_hops) => { + if let Err(error) = + publish_status(&client, &status_topic, &args, true, dpi_hops).await + { + warn!("failed to publish probe status with DPI hop: {error}"); + } + } + Err(error) => warn!("failed to update probe config: {error}"), } } else { let client = client.clone(); @@ -159,7 +162,16 @@ async fn main() -> Result<()> { Err(error) => { error!("mqtt connection error: {error}"); wait_for_connection(&mut eventloop).await; - publish_status(&client, &status_topic, &args, true).await?; + let dpi_hops = + config + .read() + .await + .as_ref() + .map_or_else(DpiHops::default, |config| DpiHops { + v4: config.dpi_hop_v4, + v6: config.dpi_hop_v6, + }); + publish_status(&client, &status_topic, &args, true, dpi_hops).await?; client.subscribe(CONFIG_TOPIC, QoS::AtLeastOnce).await?; client .subscribe("probe/tasks/v1/+", QoS::AtLeastOnce) @@ -183,7 +195,7 @@ fn mqtt_transport(mqtt_host: &str) -> Result { async fn update_config( config: &Arc>>, payload: &[u8], -) -> Result<()> { +) -> Result { let value: ProbeConfig = serde_json::from_slice(payload).context("decode probe config")?; if value.dns_samples_per_protocol == 0 { bail!("dns_samples_per_protocol must be greater than zero"); @@ -191,34 +203,103 @@ async fn update_config( if !(1..=4).contains(&value.dns_spoofing_provider_threshold) { bail!("dns_spoofing_provider_threshold must be between 1 and 4"); } - let mut control_hosts_v4 = HashSet::new(); - let mut control_hosts_v6 = HashSet::new(); - if value.traceroute_enabled { - for domain in &value.control_hosts { - match tokio::net::lookup_host((domain.as_str(), 443)).await { - Ok(addresses) => { - for address in addresses { - match address.ip() { - IpAddr::V4(address) => { - control_hosts_v4.insert(address); - } - IpAddr::V6(address) => { - control_hosts_v6.insert(address); - } - } - } - } - Err(error) => warn!("failed to resolve control host {domain}: {error}"), - } + validate_dpi_probe_config(value.dpi_probe.as_ref())?; + let dpi_hops = measure_dpi_hops(value.dpi_probe.as_ref()).await; + let loaded = LoadedProbeConfig { + config: value, + dpi_hop_v4: dpi_hops.v4, + dpi_hop_v6: dpi_hops.v6, + }; + debug!( + "measured DPI hops: IPv4={:?}, IPv6={:?}", + loaded.dpi_hop_v4, loaded.dpi_hop_v6 + ); + *config.write().await = Some(loaded); + info!("updated retained probe config"); + Ok(dpi_hops) +} + +fn validate_dpi_probe_config(config: Option<&DpiProbeConfig>) -> Result<()> { + let Some(config) = config else { + return Ok(()); + }; + if config.sni.trim().is_empty() { + bail!("dpi_probe.sni must not be empty"); + } + if config.target_v4.port() == 0 { + bail!("dpi_probe.target_v4 port must be greater than zero"); + } + if config.target_v6.port() == 0 { + bail!("dpi_probe.target_v6 port must be greater than zero"); + } + if config.connect_timeout_ms == 0 { + bail!("dpi_probe.connect_timeout_ms must be greater than zero"); + } + if config.hop_timeout_ms == 0 { + bail!("dpi_probe.hop_timeout_ms must be greater than zero"); + } + if config.max_ttl == 0 { + bail!("dpi_probe.max_ttl must be greater than zero"); + } + if config.post_dpi_hop_limit == 0 { + bail!("dpi_probe.post_dpi_hop_limit must be greater than zero"); + } + Ok(()) +} + +async fn measure_dpi_hops(config: Option<&DpiProbeConfig>) -> DpiHops { + let Some(config) = config else { + debug!("DPI hop measurement is not configured"); + return DpiHops::default(); + }; + let common = |target| dpi_hop::DpiHopProbeConfig { + target, + control_sni: config.sni.clone(), + max_ttl: config.max_ttl, + connect_timeout: Duration::from_millis(config.connect_timeout_ms), + hop_timeout: Duration::from_millis(config.hop_timeout_ms), + }; + let (v4, v6) = tokio::join!( + measure_dpi_hop(common(config.target_v4.into())), + measure_dpi_hop(common(config.target_v6.into())), + ); + DpiHops { v4, v6 } +} + +async fn measure_dpi_hop(config: dpi_hop::DpiHopProbeConfig) -> Option { + let target = config.target; + match dpi_hop::detect_dpi_hop(config).await { + Ok(result) => dpi_hop_from_result(&result), + Err(error) => { + warn!("failed to measure DPI hop for {target}: {error}"); + None } } - *config.write().await = Some(LoadedProbeConfig { - config: value, - control_hosts_v4: control_hosts_v4.into_iter().collect(), - control_hosts_v6: control_hosts_v6.into_iter().collect(), - }); - info!("updated retained probe config"); - Ok(()) +} + +fn dpi_hop_from_result(result: &dpi_hop::DpiHopProbeResult) -> Option { + debug!( + "DPI probe completed: target={}, local={}, ClientHello={} bytes", + result.target, result.local_addr, result.client_hello_bytes + ); + for hop in &result.hops { + debug!( + "DPI probe {} TTL {}: router={:?}, outcome={:?}", + result.target, hop.ttl, hop.router, hop.outcome + ); + } + if let Some(closed_hop) = result + .hops + .iter() + .find(|hop| hop.outcome == dpi_hop::DpiHopProbeHopOutcome::TcpClosed) + { + warn!( + "DPI hop measurement for {} is invalid: TCP connection closed at TTL {}", + result.target, closed_hop.ttl + ); + return None; + } + result.max_icmp_time_exceeded_ttl } async fn wait_for_connection(eventloop: &mut rumqttc::EventLoop) { @@ -244,11 +325,14 @@ async fn publish_status( topic: &str, args: &Args, online: bool, + dpi_hops: DpiHops, ) -> Result<()> { let payload = serde_json::to_vec(&ProbeStatus { online, probe_id: &args.probe_id, version: env!("CARGO_PKG_VERSION"), + dpi_hop_v4: dpi_hops.v4, + dpi_hop_v6: dpi_hops.v6, })?; client @@ -283,25 +367,15 @@ async fn handle_task( let traceroute_enabled = config .as_ref() .is_some_and(|config| config.config.traceroute_enabled); - let control_targets = config.as_ref().map_or_else(Vec::new, |config| { - if !traceroute_enabled { - return Vec::new(); - } - let mut rng = rand::thread_rng(); - match task.ip { - IpAddr::V4(_) => config - .control_hosts_v4 - .choose_multiple(&mut rng, args.traceroute_control_hosts) - .copied() - .map(IpAddr::V4) - .collect(), - IpAddr::V6(_) => config - .control_hosts_v6 - .choose_multiple(&mut rng, args.traceroute_control_hosts) - .copied() - .map(IpAddr::V6) - .collect(), - } + let dpi_hop = config.as_ref().and_then(|config| match task.ip { + IpAddr::V4(_) => config.dpi_hop_v4, + IpAddr::V6(_) => config.dpi_hop_v6, + }); + let traceroute_range = config.as_ref().and_then(|config| { + let hop_limit = config.config.dpi_probe.as_ref()?.post_dpi_hop_limit; + dpi_hop? + .checked_add(1) + .map(|start_hop| (start_hop, hop_limit)) }); let sni_check = sni::check_sni( config.as_ref().map(|config| &config.config), @@ -311,11 +385,12 @@ async fn handle_task( task.timeout_ms, ); let target_traceroute = async { - if traceroute_enabled { - traceroute::tcp_traceroute(task.ip, args.traceroute_max_hops, args.traceroute_retries) - .await - } else { - None + match (traceroute_enabled, traceroute_range) { + (true, Some((start_hop, hop_limit))) => { + traceroute::tcp_traceroute(task.ip, start_hop, hop_limit, args.traceroute_retries) + .await + } + _ => None, } }; let dns_samples_per_protocol = config @@ -333,22 +408,7 @@ async fn handle_task( dns_samples_per_protocol, dns_spoofing_provider_threshold, ); - let control_traceroute = async { - 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, dns) = - tokio::join!(sni_check, target_traceroute, control_traceroute, dns_check); + let (responses, target_traceroute, dns) = tokio::join!(sni_check, target_traceroute, dns_check); let responses = responses?; client @@ -359,10 +419,71 @@ async fn handle_task( serde_json::to_vec(&ProbeResult { responses, target_traceroute, - control_traceroute, + dpi_hop, dns, })?, ) .await .context("publish probe result") } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn decodes_separate_dpi_targets() { + let config: ProbeConfig = serde_json::from_value(serde_json::json!({ + "version": "1", + "task_timeout_ms": 15_000, + "published_at": "2026-08-21T00:00:00Z", + "hosts": [], + "dpi_probe": { + "sni": "example.com", + "target_v4": "203.0.113.10:443", + "target_v6": "[2001:db8::10]:443", + "connect_timeout_ms": 5_000, + "hop_timeout_ms": 1_000, + "max_ttl": 15 + } + })) + .unwrap(); + let dpi = config.dpi_probe.unwrap(); + + assert_eq!(dpi.target_v4.to_string(), "203.0.113.10:443"); + assert_eq!(dpi.target_v6.to_string(), "[2001:db8::10]:443"); + assert_eq!(dpi.post_dpi_hop_limit, 3); + } + + #[test] + fn status_payload_contains_separate_dpi_hops() { + let status = ProbeStatus { + online: true, + probe_id: "probe-1", + version: "1.0.0", + dpi_hop_v4: Some(4), + dpi_hop_v6: Some(6), + }; + let value = serde_json::to_value(status).unwrap(); + + assert_eq!(value["dpi_hop_v4"], 4); + assert_eq!(value["dpi_hop_v6"], 6); + } + + #[test] + fn tcp_closed_invalidates_only_its_measurement() { + let result = dpi_hop::DpiHopProbeResult { + target: "[2001:db8::10]:443".parse().unwrap(), + local_addr: "[2001:db8::1]:45000".parse().unwrap(), + client_hello_bytes: 256, + max_icmp_time_exceeded_ttl: Some(4), + hops: vec![dpi_hop::DpiHopProbeHop { + ttl: 5, + router: None, + outcome: dpi_hop::DpiHopProbeHopOutcome::TcpClosed, + }], + }; + + assert_eq!(dpi_hop_from_result(&result), None); + } +} diff --git a/probe/src/traceroute.rs b/probe/src/traceroute.rs index d84174e..3647d6c 100644 --- a/probe/src/traceroute.rs +++ b/probe/src/traceroute.rs @@ -14,10 +14,11 @@ const HOP_TIMEOUT: Duration = Duration::from_secs(1); pub async fn tcp_traceroute( target: IpAddr, - max_hops: u8, + start_hop: u8, + hop_limit: u8, retries: u8, ) -> Option { - tokio::task::spawn_blocking(move || trace_blocking(target, max_hops, retries)) + tokio::task::spawn_blocking(move || trace_blocking(target, start_hop, hop_limit, retries)) .await .map_err(|error| log::warn!("TCP traceroute task failed for {target}: {error}")) .ok()? @@ -25,15 +26,24 @@ pub async fn tcp_traceroute( .ok() } -fn trace_blocking(target: IpAddr, max_hops: u8, retries: u8) -> io::Result { +fn trace_blocking( + target: IpAddr, + start_hop: u8, + hop_limit: u8, + retries: u8, +) -> io::Result { + if start_hop == 0 || hop_limit == 0 || retries == 0 { + return Err(io::Error::new( + io::ErrorKind::InvalidInput, + "start_hop, hop_limit, and retries must be greater than zero", + )); + } let (domain, icmp_protocol) = match target { IpAddr::V4(_) => (Domain::IPV4, Protocol::ICMPV4), IpAddr::V6(_) => (Domain::IPV6, Protocol::ICMPV6), }; let receiver = Socket::new(domain, Type::RAW, Some(icmp_protocol))?; - let mut last_icmp_hop = None; - - for ttl in 1..=max_hops { + for ttl in post_dpi_hops(start_hop, hop_limit) { let destination = SockAddr::from(SocketAddr::new(target, HTTPS_PORT)); let mut tcp_attempts = Vec::with_capacity(retries as usize); let mut last_error = None; @@ -72,7 +82,12 @@ fn trace_blocking(target: IpAddr, max_hops: u8, retries: u8) -> io::Result last_icmp_hop = Some(ttl), + HopResponse::IcmpTimeExceeded => { + return Ok(TcpTracerouteResult { + target, + result: TcpTracerouteOutcome::IcmpTimeExceeded { hop: ttl }, + }); + } HopResponse::Rst => { return Ok(TcpTracerouteResult { target, @@ -91,12 +106,14 @@ fn trace_blocking(target: IpAddr, max_hops: u8, retries: u8) -> io::Result impl Iterator { + (0..hop_limit).map_while(move |offset| start_hop.checked_add(offset)) +} + fn is_connect_in_progress(error: &io::Error) -> bool { error.kind() == io::ErrorKind::WouldBlock || error.raw_os_error() == Some(rustix::io::Errno::INPROGRESS.raw_os_error()) @@ -303,4 +320,10 @@ mod tests { assert!(matching_ipv6_time_exceeded(&packet, target, 42_000)); assert!(!matching_ipv6_time_exceeded(&packet, target, 42_001)); } + + #[test] + fn post_dpi_hop_range_starts_at_requested_hop_and_is_bounded() { + assert_eq!(post_dpi_hops(6, 3).collect::>(), vec![6, 7, 8]); + assert_eq!(post_dpi_hops(254, 3).collect::>(), vec![254, 255]); + } } diff --git a/reports/src/probe.rs b/reports/src/probe.rs index 28e977b..103cd18 100644 --- a/reports/src/probe.rs +++ b/reports/src/probe.rs @@ -1,11 +1,15 @@ use serde::{Deserialize, Serialize}; -use std::net::IpAddr; +use std::net::{IpAddr, SocketAddrV4, SocketAddrV6}; #[derive(Clone, Serialize, Deserialize)] pub struct ProbeStatus<'a> { pub online: bool, pub probe_id: &'a str, pub version: &'a str, + #[serde(default)] + pub dpi_hop_v4: Option, + #[serde(default)] + pub dpi_hop_v6: Option, } #[derive(Clone, Serialize, Deserialize)] @@ -16,12 +20,28 @@ pub struct ProbeConfig { pub hosts: Vec, #[serde(default)] pub traceroute_enabled: bool, - #[serde(default)] - pub control_hosts: Vec, #[serde(default = "default_dns_samples_per_protocol")] pub dns_samples_per_protocol: u8, #[serde(default = "default_dns_spoofing_provider_threshold")] pub dns_spoofing_provider_threshold: u8, + #[serde(default)] + pub dpi_probe: Option, +} + +#[derive(Clone, Serialize, Deserialize)] +pub struct DpiProbeConfig { + pub sni: String, + pub target_v4: SocketAddrV4, + pub target_v6: SocketAddrV6, + pub connect_timeout_ms: u64, + pub hop_timeout_ms: u64, + pub max_ttl: u8, + #[serde(default = "default_post_dpi_hop_limit")] + pub post_dpi_hop_limit: u8, +} + +pub const fn default_post_dpi_hop_limit() -> u8 { + 3 } pub const fn default_dns_samples_per_protocol() -> u8 { @@ -64,7 +84,7 @@ pub struct ProbeResultEvent { pub probe_id: String, pub host_results: Vec, pub target_traceroute: Option, - pub control_traceroute: Option, + pub dpi_hop: Option, pub dns: Option, } @@ -72,7 +92,8 @@ pub struct ProbeResultEvent { pub struct ProbeResult { pub responses: Option>, pub target_traceroute: Option, - pub control_traceroute: Option, + #[serde(default)] + pub dpi_hop: Option, #[serde(default)] pub dns: Option, } diff --git a/website/Cargo.toml b/website/Cargo.toml index d3c0957..9e60d7f 100644 --- a/website/Cargo.toml +++ b/website/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "website" -version = "1.2.1" +version = "1.2.2" edition = "2024" [dependencies] diff --git a/website/migrations/20260821000000_remove_control_traceroute.sql b/website/migrations/20260821000000_remove_control_traceroute.sql new file mode 100644 index 0000000..2653d06 --- /dev/null +++ b/website/migrations/20260821000000_remove_control_traceroute.sql @@ -0,0 +1,3 @@ +ALTER TABLE probe_reports + DROP COLUMN control_hop_count, + DROP COLUMN control_trace_result; diff --git a/website/probe-hosts.toml b/website/probe-hosts.toml index 6233842..47cbb6a 100644 --- a/website/probe-hosts.toml +++ b/website/probe-hosts.toml @@ -2,7 +2,15 @@ timeout_sec = 3 min_data = 65536 dns_samples_per_protocol = 3 dns_spoofing_provider_threshold = 2 -control_hosts = ["kinopoisk.ru", "wildberries.ru", "yandex.ru", "mail.ru"] + +[dpi_probe] +sni = "rutracker.org" +target_v4 = "188.120.251.237:443" +target_v6 = "[2a01:230:4:ca4::2]:443" +connect_timeout_ms = 5000 +hop_timeout_ms = 1000 +max_ttl = 15 +post_dpi_hop_limit = 3 [[hosts]] id = "hil-hetzner" # hil-speed.hetzner.com diff --git a/website/src/api/probe.rs b/website/src/api/probe.rs index 94cbc78..da0796d 100644 --- a/website/src/api/probe.rs +++ b/website/src/api/probe.rs @@ -107,7 +107,6 @@ pub async fn probe_query( Ok(result) => { responded_probes.insert(result.probe_id.clone()); let target_traceroute = result.target_traceroute.clone(); - let control_traceroute = result.control_traceroute.clone(); let reporter_info = match fetch_probe_reporter_info(&result.probe_id, &pool).await { Ok(info) => info, Err(error) => { @@ -123,7 +122,6 @@ pub async fn probe_query( query_id, &response, target_traceroute.as_ref(), - control_traceroute.as_ref(), &pool, ).await { warn!("api: failed to save probe report for query {id}: {error}"); @@ -157,7 +155,7 @@ pub fn build_probe_response( &raw.host_results, config, raw.target_traceroute.as_ref(), - raw.control_traceroute.as_ref(), + raw.dpi_hop, raw.dns.as_ref(), ); let target_hop = @@ -197,6 +195,7 @@ pub fn build_probe_response( "verdicts": verdicts, "host_results": host_results, "target_hop": target_hop, + "dpi_hop": raw.dpi_hop, "dns": raw.dns, }) } @@ -205,7 +204,6 @@ async fn insert_probe_report( query_id: Uuid, response: &Value, target_traceroute: Option<&TcpTracerouteResult>, - control_traceroute: Option<&TcpTracerouteResult>, pool: &PgPool, ) -> Result<(), sqlx::Error> { let probe_id = response @@ -223,7 +221,6 @@ async fn insert_probe_report( }) .unwrap_or_else(|| vec!["uncertain"]); let (target_hop_count, target_trace_result) = traceroute_columns(target_traceroute); - let (control_hop_count, control_trace_result) = traceroute_columns(control_traceroute); let Some(probe_id) = probe_id else { warn!("api: ignoring probe report with non-numeric probe_id"); @@ -238,20 +235,16 @@ async fn insert_probe_report( verdicts, result, target_hop_count, - target_trace_result, - control_hop_count, - control_trace_result + target_trace_result ) - VALUES ($1, $2, $3, $4, $5, $6, $7, $8) + VALUES ($1, $2, $3, $4, $5, $6) ON CONFLICT (query_id, probe_id) DO UPDATE SET date = NOW(), verdicts = EXCLUDED.verdicts, result = EXCLUDED.result, target_hop_count = EXCLUDED.target_hop_count, - target_trace_result = EXCLUDED.target_trace_result, - control_hop_count = EXCLUDED.control_hop_count, - control_trace_result = EXCLUDED.control_trace_result + target_trace_result = EXCLUDED.target_trace_result "#, ) .bind(query_id) @@ -260,8 +253,6 @@ async fn insert_probe_report( .bind(response) .bind(target_hop_count) .bind(target_trace_result) - .bind(control_hop_count) - .bind(control_trace_result) .execute(pool) .await?; @@ -299,26 +290,36 @@ fn build_probe_verdicts( results: &[HostProbeResult], config: &ProbeConfig, target_traceroute: Option<&TcpTracerouteResult>, - control_traceroute: Option<&TcpTracerouteResult>, + dpi_hop: Option, dns: Option<&reports::probe::DnsProbeResult>, ) -> Vec<&'static str> { - let mut verdicts = Vec::new(); - let dns_spoofing = dns.is_some_and(|result| result.spoofing_detected); - if let ( - Some(TcpTracerouteResult { - result: TcpTracerouteOutcome::IcmpTimeExceeded { hop: target_hop }, - .. - }), - Some(TcpTracerouteResult { - result: TcpTracerouteOutcome::IcmpTimeExceeded { hop: control_hop }, - .. - }), - ) = (target_traceroute, control_traceroute) - && target_hop < control_hop - { - verdicts.push("tspu_block"); + if results.is_empty() && target_traceroute.is_none() && dns.is_none() { + return vec!["uncertain"]; } + let tspu_block = dpi_hop.is_some() + && target_traceroute + .is_some_and(|traceroute| matches!(&traceroute.result, TcpTracerouteOutcome::Timeout)); + let host_verdict = build_host_verdict(results, config); + let dns_spoofing = dns.is_some_and(|result| result.spoofing_detected); + + let mut verdicts = [ + tspu_block.then_some("tspu_block"), + (!matches!(host_verdict, "ok" | "uncertain")).then_some(host_verdict), + dns_spoofing.then_some("dns_spoofing"), + ] + .into_iter() + .flatten() + .collect::>(); + + if verdicts.is_empty() { + verdicts.push(host_verdict); + } + + verdicts +} + +fn build_host_verdict(results: &[HostProbeResult], config: &ProbeConfig) -> &'static str { let matched = results .iter() .filter_map(|result| { @@ -330,7 +331,7 @@ fn build_probe_verdicts( }) .collect::>(); - let host_verdict = if results.is_empty() { + if results.is_empty() { "ok" } else if matched.is_empty() { "uncertain" @@ -382,19 +383,7 @@ fn build_probe_verdicts( } else { "uncertain" } - }; - - if !matches!(host_verdict, "ok" | "uncertain") { - verdicts.push(host_verdict); } - if dns_spoofing { - verdicts.push("dns_spoofing"); - } - if verdicts.is_empty() { - verdicts.push(host_verdict); - } - - verdicts } fn publish_error_status(error: PublishError) -> Status { @@ -429,40 +418,46 @@ fn is_strict_majority(total: usize, count: usize) -> bool { mod tests { use super::*; - fn icmp_trace(hop: u8) -> TcpTracerouteResult { - TcpTracerouteResult { - target: "192.0.2.1".parse().unwrap(), - result: TcpTracerouteOutcome::IcmpTimeExceeded { hop }, - } + #[test] + fn no_checks_returns_uncertain() { + assert_eq!( + build_probe_verdicts(&[], &empty_config(), None, None, None), + vec!["uncertain"] + ); } #[test] - fn tspu_block_requires_an_earlier_target_icmp_hop() { - let target = icmp_trace(3); - let control = icmp_trace(5); + fn tspu_block_requires_a_dpi_hop_and_post_dpi_timeout() { + let timeout = TcpTracerouteResult { + target: "192.0.2.1".parse().unwrap(), + result: TcpTracerouteOutcome::Timeout, + }; assert_eq!( - build_probe_verdicts(&[], &empty_config(), Some(&target), Some(&control), None), + build_probe_verdicts(&[], &empty_config(), Some(&timeout), Some(4), None), vec!["tspu_block"] ); - - let target = icmp_trace(5); assert_eq!( - build_probe_verdicts(&[], &empty_config(), Some(&target), Some(&control), None), + build_probe_verdicts(&[], &empty_config(), Some(&timeout), None, None), vec!["ok"] ); } #[test] - fn tspu_block_requires_two_icmp_outcomes() { - let target = TcpTracerouteResult { - target: "192.0.2.1".parse().unwrap(), - result: TcpTracerouteOutcome::Connected { hop: 2 }, - }; - let control = icmp_trace(5); - assert_eq!( - build_probe_verdicts(&[], &empty_config(), Some(&target), Some(&control), None), - vec!["ok"] - ); + fn any_post_dpi_response_clears_tspu_block() { + for result in [ + TcpTracerouteOutcome::IcmpTimeExceeded { hop: 5 }, + TcpTracerouteOutcome::Connected { hop: 5 }, + TcpTracerouteOutcome::Rst { hop: 5 }, + ] { + let target = TcpTracerouteResult { + target: "192.0.2.1".parse().unwrap(), + result, + }; + assert_eq!( + build_probe_verdicts(&[], &empty_config(), Some(&target), Some(4), None), + vec!["ok"] + ); + } } #[test] @@ -501,10 +496,10 @@ mod tests { published_at: String::new(), hosts: vec![], traceroute_enabled: true, - control_hosts: vec![], dns_samples_per_protocol: reports::probe::default_dns_samples_per_protocol(), dns_spoofing_provider_threshold: reports::probe::default_dns_spoofing_provider_threshold(), + dpi_probe: None, } } } diff --git a/website/src/mqtt.rs b/website/src/mqtt.rs index b17c53b..c38ec52 100644 --- a/website/src/mqtt.rs +++ b/website/src/mqtt.rs @@ -1,6 +1,8 @@ use log::{info, warn}; use reports::probe::HostType; -use reports::probe::{Host, ProbeConfig, ProbeResult, ProbeResultEvent, ProbeStatus, ProbeTask}; +use reports::probe::{ + DpiProbeConfig, Host, ProbeConfig, ProbeResult, ProbeResultEvent, ProbeStatus, ProbeTask, +}; use rocket::serde::json::serde_json; use rumqttc::{AsyncClient, Event as MqttEvent, Incoming, MqttOptions, QoS}; use serde::Deserialize; @@ -12,6 +14,8 @@ use std::net::IpAddr; use std::sync::Arc; use std::time::Duration; +const MQTT_MAX_PACKET_SIZE: usize = 1024 * 1024; + const DEFAULT_PROBE_HOSTS: &str = include_str!("../probe-hosts.toml"); #[derive(Debug)] @@ -66,7 +70,7 @@ struct ProbeHostsFile { #[serde(default = "reports::probe::default_dns_spoofing_provider_threshold")] dns_spoofing_provider_threshold: u8, #[serde(default)] - control_hosts: Vec, + dpi_probe: Option, hosts: Vec, } @@ -93,10 +97,10 @@ impl MqttPublisher { published_at: Utc::now().to_rfc3339(), hosts: Vec::new(), traceroute_enabled: false, - control_hosts: Vec::new(), dns_samples_per_protocol: reports::probe::default_dns_samples_per_protocol(), dns_spoofing_provider_threshold: reports::probe::default_dns_spoofing_provider_threshold(), + dpi_probe: None, } })); let admin_token = match std::env::var("MQTT_ADMIN_TOKEN") { @@ -124,6 +128,7 @@ impl MqttPublisher { let mut options = MqttOptions::new(client_id, host.clone(), port); options.set_credentials("admin", admin_token); options.set_keep_alive(Duration::from_secs(10)); + options.set_max_packet_size(MQTT_MAX_PACKET_SIZE, MQTT_MAX_PACKET_SIZE); let (client, mut eventloop) = AsyncClient::new(options, 100); let event_sessions = sessions.clone(); @@ -278,17 +283,17 @@ fn load_probe_config(task_timeout_ms: u64) -> Result traceroute_enabled: std::env::var("PROBE_TRACEROUTE_ENABLED") .ok() .is_some_and(|value| value == "1" || value.eq_ignore_ascii_case("true")), - control_hosts: config.control_hosts, dns_samples_per_protocol: config.dns_samples_per_protocol, dns_spoofing_provider_threshold: config.dns_spoofing_provider_threshold, + dpi_probe: config.dpi_probe, }) } struct ParsedProbeConfig { hosts: Vec, - control_hosts: Vec, dns_samples_per_protocol: u8, dns_spoofing_provider_threshold: u8, + dpi_probe: Option, } fn parse_probe_hosts(contents: &str) -> Result { @@ -309,9 +314,9 @@ fn parse_probe_hosts(contents: &str) -> Result Ok(ParsedProbeConfig { hosts, - control_hosts: config.control_hosts, dns_samples_per_protocol: config.dns_samples_per_protocol, dns_spoofing_provider_threshold: config.dns_spoofing_provider_threshold, + dpi_probe: config.dpi_probe, }) } @@ -369,7 +374,7 @@ async fn dispatch_probe_result( probe_id: probe_id.to_string(), host_results: result.responses.unwrap_or_default(), target_traceroute: result.target_traceroute, - control_traceroute: result.control_traceroute, + dpi_hop: result.dpi_hop, dns: result.dns, }); }