feat: traceroutes packets capture

This commit is contained in:
Lowder
2026-10-08 21:08:25 +05:00
parent c023d3dfb5
commit 9db493bb0a
9 changed files with 794 additions and 32 deletions
Generated
+1 -1
View File
@@ -2488,7 +2488,7 @@ dependencies = [
[[package]]
name = "probe"
version = "0.7.1"
version = "0.7.2"
dependencies = [
"anyhow",
"clap",
+138 -4
View File
@@ -39,8 +39,38 @@ type Hop = {
address: string | null;
reverse_names: string[];
outcome: string;
tcp_diagnostics?: TcpDiagnostics;
};
type TcpPacket = {
observed_ms: number;
source: string;
destination: string;
ttl: number | null;
ip_id: number | null;
sequence: number;
acknowledgment: number;
window: number;
flags: string[];
timestamp: number | null;
timestamp_echo: number | null;
options_hex: string;
payload_bytes: number;
};
type TcpDiagnostics = {
packets: TcpPacket[];
capture_error: string | null;
connect_error: string | null;
send_error: string | null;
truncated: boolean;
client_hello_sent_ms: number | null;
junk_sent_ms: number | null;
};
type DpiHop = {
ttl: number;
src: string | null;
outcome: string;
tcp_diagnostics?: TcpDiagnostics;
};
type DpiHop = { ttl: number; src: string | null; outcome: string };
type CommandType =
| "resubscribe_tasks"
| "traceroute"
@@ -393,6 +423,90 @@ const outcomeLabel: Record<string, string> = {
};
</script>
{#snippet tcpDiagnostics(diagnostics: TcpDiagnostics | undefined)}
{#if diagnostics}
<details
class="mt-2 rounded-lg border border-neutral-700 p-3 text-xs"
open={diagnostics.packets.some((packet) => packet.flags.includes("RST") || packet.flags.includes("FIN"))}
>
<summary class="cursor-pointer text-neutral-300">
TCP: {diagnostics.packets.length} входящих пакетов{diagnostics.truncated ? " (список сокращён)" : ""}
</summary>
{#if diagnostics.capture_error}
<p class="mt-2 text-amber-300">
Ошибка захвата: {diagnostics.capture_error}
</p>
{/if}
{#if diagnostics.connect_error}
<p class="mt-2 text-amber-300">
Ошибка соединения: {diagnostics.connect_error}
</p>
{/if}
{#if diagnostics.send_error}
<p class="mt-2 text-amber-300">
Ошибка отправки: {diagnostics.send_error}
</p>
{/if}
<p class="mt-2 text-neutral-400">
Время приблизительное, от начала захвата.
{#if diagnostics.client_hello_sent_ms !== null}
ClientHello: {diagnostics.client_hello_sent_ms} мс; данные:
{diagnostics.junk_sent_ms ?? "—"}
мс.
{/if}
Δ TTL сравнивает пакет с SYN-ACK того же соединения; это не расстояние
до DPI.
</p>
<div class="mt-3 overflow-x-auto">
<table class="w-full whitespace-nowrap text-left font-mono">
<thead class="text-neutral-500">
<tr>
<th class="pr-4">мс</th>
<th class="pr-4">Источник → цель</th>
<th class="pr-4">Флаги</th>
<th class="pr-4">TTL / Δ</th>
<th class="pr-4">SEQ / ACK</th>
<th class="pr-4">Окно / IP ID</th>
<th class="pr-4">TS / echo</th>
<th>Байты / опции</th>
</tr>
</thead>
<tbody>
{#each diagnostics.packets as packet}
{@const baseline = diagnostics.packets.find((p) => p.destination === packet.destination && p.flags.includes("SYN") && p.flags.includes("ACK"))}
<tr
class="border-t border-neutral-800"
class:text-amber-300={packet.flags.includes("RST") || packet.flags.includes("FIN")}
>
<td class="py-2 pr-4">{packet.observed_ms}</td>
<td class="pr-4">{packet.source} → {packet.destination}</td>
<td class="pr-4">{packet.flags.join(" ") || "—"}</td>
<td class="pr-4">
{packet.ttl ?? "—"}
/
{packet.ttl !== null && baseline?.ttl != null ? packet.ttl - baseline.ttl : "—"}
</td>
<td class="pr-4">{packet.sequence}/ {packet.acknowledgment}</td>
<td class="pr-4">{packet.window} / {packet.ip_id ?? "—"}</td>
<td class="pr-4">
{packet.timestamp ?? "—"}
/ {packet.timestamp_echo ?? "—"}
</td>
<td>{packet.payload_bytes} / {packet.options_hex || "—"}</td>
</tr>
{/each}
</tbody>
</table>
</div>
{#if diagnostics.packets.some((p) => p.ttl === null)}
<p class="mt-2 text-neutral-500">
Для части пакетов IP-заголовок недоступен; TTL не измерен.
</p>
{/if}
</details>
{/if}
{/snippet}
{#snippet dpiHops(hops: DpiHop[])}
<ol class="space-y-1">
{#each hops as hop}
@@ -404,6 +518,11 @@ const outcomeLabel: Record<string, string> = {
<span class="text-xs text-neutral-500"
>{outcomeLabel[hop.outcome] ?? hop.outcome}</span
>
{#if hop.tcp_diagnostics}
<div class="col-span-3 min-w-0">
{@render tcpDiagnostics(hop.tcp_diagnostics)}
</div>
{/if}
</li>
{/each}
</ol>
@@ -782,9 +901,10 @@ PROBE_TOKEN={created.token}</pre>
></label
>
<p class="mb-4 text-sm text-neutral-400">
TCP-соединение к хосту на порту 443, ClientHello с указанным SNI,
затем пакеты с возрастающим TTL. Хост разрешается сканером;
используется первый IP адрес.
Для каждого TTL создаётся новое TCP-соединение к хосту на порту 443.
ClientHello с указанным SNI отправляется с обычным TTL; через 100 мс
отправляются случайные данные с проверяемым TTL. Хост разрешается
сканером; используется первый IP адрес.
</p>
{/if}
<div class="grid gap-3 sm:grid-cols-[1fr_2fr_100px_auto] sm:items-end">
@@ -823,6 +943,15 @@ PROBE_TOKEN={created.token}</pre>
{#if busy.endsWith(":command")}
<p class="mt-5 text-sm text-cyan-300">Ожидание ответа сканера…</p>
{/if}
{#if result && result.type !== "error"}
<button
type="button"
class="btn-secondary mt-4"
onclick={() => result && void copy(JSON.stringify(result, null, 2))}
>
Копировать JSON результата
</button>
{/if}
{#if result?.type === "error"}
<p class="mt-5 text-sm text-red-300">{result.message}</p>
{:else if result?.type === "resubscribe_tasks"}
@@ -875,6 +1004,11 @@ PROBE_TOKEN={created.token}</pre>
<span class="text-xs text-neutral-500"
>{outcomeLabel[hop.outcome] ?? hop.outcome}</span
>
{#if hop.tcp_diagnostics}
<div class="col-span-3 min-w-0">
{@render tcpDiagnostics(hop.tcp_diagnostics)}
</div>
{/if}
</li>
{/each}
</ol>
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "probe"
version = "0.7.1"
version = "0.7.2"
edition = "2024"
license-file = "../LICENSE"
description = "Dynamic network probe daemon for Cheburcheck"
+169 -15
View File
@@ -1,3 +1,4 @@
use crate::packet_capture::{TcpCapture, connect_with_port};
use etherparse::{
Icmpv4Type, Icmpv6Slice, Icmpv6Type, IpNumber, LaxNetSlice, LaxSlicedPacket, TransportSlice,
icmpv4, icmpv6,
@@ -24,6 +25,7 @@ pub struct DpiHopProbeConfig {
pub hop_timeout: Duration,
pub accept_rst_fin_as_timeout: bool,
pub accept_ack_as_timeout: bool,
pub collect_tcp_metadata: bool,
}
#[derive(Debug, Clone)]
@@ -72,7 +74,38 @@ pub fn detect_dpi_hop_blocking(config: DpiHopProbeConfig) -> io::Result<DpiHopPr
// subsequently changed socket TTL. Use an independent connection for
// every TTL so each hop corresponds to an actual packet and cannot
// affect later hops.
let mut tcp = TcpStream::connect_timeout(&config.target, config.connect_timeout)?;
let capture = config
.collect_tcp_metadata
.then(|| TcpCapture::start(config.target));
let (port, connection) = if capture.is_some() {
connect_with_port(config.target, config.connect_timeout)
} else {
(
None,
TcpStream::connect_timeout(&config.target, config.connect_timeout),
)
};
let mut tcp = match connection {
Ok(tcp) => tcp,
Err(error) => {
let Some(capture) = capture else {
return Err(error);
};
let mut diagnostics = capture.finish(&port.into_iter().collect::<Vec<_>>());
diagnostics.connect_error = Some(error.to_string());
hops.push(DpiHopProbeHop {
ttl,
router: None,
outcome: if error.kind() == io::ErrorKind::ConnectionRefused {
DpiHopProbeHopOutcome::TcpClosed
} else {
DpiHopProbeHopOutcome::Timeout
},
tcp_diagnostics: Some(diagnostics),
});
break;
}
};
tcp.set_nodelay(true)?;
tcp.set_write_timeout(Some(config.connect_timeout))?;
let local_addr = tcp.local_addr()?;
@@ -97,24 +130,50 @@ pub fn detect_dpi_hop_blocking(config: DpiHopProbeConfig) -> io::Result<DpiHopPr
result_local_addr = Some(local_addr);
client_hello_bytes = Some(client_hello.len());
}
tcp.write_all(&client_hello)?;
// A blocking DPI may silently discard the ClientHello, so waiting for
// a TCP acknowledgement would deadlock the measurement. Give the DPI
// a short interval to classify this new flow before sending the
// TTL-limited payload instead. This is needed for every attempt because
// each TTL deliberately uses a fresh TCP connection.
// Classify the flow with a normal-TTL ClientHello. Limiting this first
// packet triggers low-TTL filtering even for control SNI on some paths.
let hello_sent_ms = capture.as_ref().map(TcpCapture::elapsed_ms);
if let Err(error) = tcp.write_all(&client_hello) {
if capture.is_none() {
return Err(error);
}
// Continue observing after a fast reset to catch competing server
// packets as well. Automatic measurements retain their old behavior.
std::thread::sleep(config.hop_timeout);
let mut diagnostics = finish_capture(capture, local_addr, hello_sent_ms, None);
if let Some(diagnostics) = &mut diagnostics {
diagnostics.send_error = Some(error.to_string());
}
hops.push(DpiHopProbeHop {
ttl,
router: None,
outcome: if matches!(
error.kind(),
io::ErrorKind::ConnectionReset | io::ErrorKind::BrokenPipe
) {
DpiHopProbeHopOutcome::TcpClosed
} else {
DpiHopProbeHopOutcome::Timeout
},
tcp_diagnostics: diagnostics,
});
break;
}
std::thread::sleep(DPI_CLASSIFICATION_DELAY);
drain_socket(&icmp)?;
let mut payload = [0u8; PROBE_BYTES];
rand::thread_rng().fill_bytes(&mut payload);
if !send_with_ttl(&mut tcp, config.target, &payload, ttl)? {
let junk_sent_ms = capture.as_ref().map(TcpCapture::elapsed_ms);
let payload_sent = send_with_ttl(&mut tcp, config.target, &payload, ttl)?;
if !payload_sent {
if capture.is_some() {
std::thread::sleep(config.hop_timeout);
}
hops.push(DpiHopProbeHop {
ttl,
router: None,
outcome: classify_hop(None, true, false, &config),
tcp_diagnostics: finish_capture(capture, local_addr, hello_sent_ms, junk_sent_ms),
});
if config.accept_rst_fin_as_timeout {
continue;
@@ -134,6 +193,7 @@ pub fn detect_dpi_hop_blocking(config: DpiHopProbeConfig) -> io::Result<DpiHopPr
ttl,
router,
outcome,
tcp_diagnostics: finish_capture(capture, local_addr, hello_sent_ms, junk_sent_ms),
});
if outcome.invalidates_measurement() {
break;
@@ -142,15 +202,42 @@ pub fn detect_dpi_hop_blocking(config: DpiHopProbeConfig) -> io::Result<DpiHopPr
Ok(DpiHopProbeResult {
target: config.target,
local_addr: result_local_addr
.ok_or_else(|| io::Error::other("DPI hop probe made no TTL attempts"))?,
client_hello_bytes: client_hello_bytes
.ok_or_else(|| io::Error::other("DPI hop probe produced no ClientHello"))?,
local_addr: result_local_addr.unwrap_or_else(|| {
SocketAddr::new(
if config.target.is_ipv4() {
IpAddr::V4(Ipv4Addr::UNSPECIFIED)
} else {
IpAddr::V6(std::net::Ipv6Addr::UNSPECIFIED)
},
0,
)
}),
client_hello_bytes: client_hello_bytes.unwrap_or(0),
max_icmp_time_exceeded_ttl,
hops,
})
}
fn finish_capture(
capture: Option<TcpCapture>,
local: SocketAddr,
hello_ms: Option<u64>,
junk_ms: Option<u64>,
) -> Option<reports::probe::TcpDiagnostics> {
capture.map(|capture| {
let mut diagnostics = capture.finish(&[local.port()]);
for packet in &mut diagnostics.packets {
// Raw IPv6 sockets may omit the destination IP header.
if packet.destination.ip().is_unspecified() {
packet.destination.set_ip(local.ip());
}
}
diagnostics.client_hello_sent_ms = hello_ms;
diagnostics.junk_sent_ms = junk_ms;
diagnostics
})
}
fn make_client_hello(sni: &str) -> io::Result<Vec<u8>> {
let server_name = ServerName::try_from(sni.to_owned()).map_err(|_| {
io::Error::new(
@@ -470,6 +557,72 @@ fn matching_quoted_tcp_tuple(packet: &[u8], local_addr: SocketAddr, target: Sock
mod tests {
use super::*;
#[test]
#[ignore = "requires CAP_NET_RAW and local socket access"]
fn live_clienthello_keeps_default_ttl_and_junk_uses_tested_ttl() {
use std::io::Read;
let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
let target = listener.local_addr().unwrap();
let hello_len = make_client_hello("example.com").unwrap().len();
let server = std::thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(2)))
.unwrap();
let mut hello = vec![0; hello_len];
stream.read_exact(&mut hello).unwrap();
assert_eq!(hello[0], 22); // TLS handshake record
let mut junk = [0; PROBE_BYTES];
stream.read_exact(&mut junk).unwrap();
stream.write_all(b"response").unwrap();
std::thread::sleep(Duration::from_millis(250));
});
let receiver = Socket::new(Domain::IPV4, Type::RAW, Some(Protocol::TCP)).unwrap();
receiver
.set_read_timeout(Some(Duration::from_secs(2)))
.unwrap();
let observer = std::thread::spawn(move || {
let mut packets = Vec::new();
let mut buffer = [0; 65535];
while packets.len() < 2 {
let (bytes, _) = recv_socket(&receiver, &mut buffer).unwrap();
let packet = LaxSlicedPacket::from_ip(&buffer[..bytes]).unwrap();
if let (Some(LaxNetSlice::Ipv4(ip)), Some(TransportSlice::Tcp(tcp))) =
(packet.net, packet.transport)
{
if tcp.destination_port() == target.port() && !tcp.payload().is_empty() {
packets.push((ip.header().ttl(), tcp.payload().len()));
}
}
}
packets
});
let default_ttl = Socket::new(Domain::IPV4, Type::STREAM, Some(Protocol::TCP))
.unwrap()
.ttl_v4()
.unwrap() as u8;
let trace = detect_dpi_hop_blocking(DpiHopProbeConfig {
target,
max_ttl: 1,
hop_timeout: Duration::from_millis(100),
collect_tcp_metadata: true,
..probe_config()
})
.unwrap();
let packets = observer.join().unwrap();
server.join().unwrap();
assert_eq!(packets, [(default_ttl, hello_len), (1, PROBE_BYTES)]);
let metadata = trace.hops[0].tcp_diagnostics.as_ref().unwrap();
assert_eq!(metadata.capture_error, None);
assert!(
metadata
.packets
.iter()
.any(|p| p.flags.contains(&"SYN".into()))
);
assert!(metadata.junk_sent_ms.unwrap() >= metadata.client_hello_sent_ms.unwrap() + 100);
}
fn probe_config() -> DpiHopProbeConfig {
DpiHopProbeConfig {
target: "192.0.2.1:443".parse().unwrap(),
@@ -479,6 +632,7 @@ mod tests {
hop_timeout: Duration::from_secs(1),
accept_rst_fin_as_timeout: false,
accept_ack_as_timeout: false,
collect_tcp_metadata: false,
}
}
+48 -11
View File
@@ -1,5 +1,6 @@
mod dns;
mod dpi_hop;
mod packet_capture;
mod sni;
mod traceroute;
mod update;
@@ -461,7 +462,7 @@ async fn update_config(
bail!("dns_spoofing_provider_threshold must be between 1 and 4");
}
validate_dpi_probe_config(value.dpi_probe.as_ref())?;
let dpi_hops = measure_dpi_hops(value.dpi_probe.as_ref()).await;
let dpi_hops = measure_dpi_hops(value.dpi_probe.as_ref(), false).await;
let loaded = LoadedProbeConfig {
config: value,
dpi_hop_v4: dpi_hops.v4,
@@ -494,7 +495,15 @@ async fn run_dpi_command(
.dpi_probe
.as_ref()
.context("DPI hop measurement is not configured")?;
let hops = measure_dpi_hops(Some(dpi)).await;
let hops = measure_dpi_hops(Some(dpi), true).await;
let mut status_hops = hops.clone();
for hop in status_hops
.hops_v4
.iter_mut()
.chain(&mut status_hops.hops_v6)
{
hop.tcp_diagnostics = None;
}
let mut guard = config.write().await;
let loaded = guard.as_mut().context("probe config is not loaded")?;
if loaded.config.version != snapshot.config.version
@@ -506,8 +515,8 @@ async fn run_dpi_command(
}
loaded.dpi_hop_v4 = hops.v4;
loaded.dpi_hop_v6 = hops.v6;
loaded.dpi_hops_v4 = hops.hops_v4.clone();
loaded.dpi_hops_v6 = hops.hops_v6.clone();
loaded.dpi_hops_v4 = status_hops.hops_v4.clone();
loaded.dpi_hops_v6 = status_hops.hops_v6.clone();
Ok((
ProbeCommandResult::RemeasureDpiHop {
dpi_hop_v4: hops.v4,
@@ -515,7 +524,7 @@ async fn run_dpi_command(
dpi_hops_v4: hops.hops_v4.clone(),
dpi_hops_v6: hops.hops_v6.clone(),
},
Some(hops),
Some(status_hops),
))
}
ProbeCommand::SniTraceroute {
@@ -553,6 +562,7 @@ async fn run_dpi_command(
accept_ack_as_timeout: dpi_config
.as_ref()
.is_some_and(|dpi| dpi.accept_ack_as_timeout),
collect_tcp_metadata: true,
})
.await?;
Ok((
@@ -598,7 +608,7 @@ fn validate_dpi_probe_config(config: Option<&DpiProbeConfig>) -> Result<()> {
Ok(())
}
async fn measure_dpi_hops(config: Option<&DpiProbeConfig>) -> DpiHops {
async fn measure_dpi_hops(config: Option<&DpiProbeConfig>, collect_tcp_metadata: bool) -> DpiHops {
let Some(config) = config else {
debug!("DPI hop measurement is not configured");
return DpiHops::default();
@@ -611,6 +621,7 @@ async fn measure_dpi_hops(config: Option<&DpiProbeConfig>) -> DpiHops {
hop_timeout: Duration::from_millis(config.hop_timeout_ms),
accept_rst_fin_as_timeout: config.accept_rst_fin_as_timeout,
accept_ack_as_timeout: config.accept_ack_as_timeout,
collect_tcp_metadata,
};
let (v4, v6) = tokio::join!(
measure_dpi_hop(common(config.target_v4.into())),
@@ -648,11 +659,12 @@ fn dpi_hop_from_result(result: &dpi_hop::DpiHopProbeResult) -> Option<u8> {
result.target, hop.ttl, hop.router, hop.outcome
);
}
if let Some(invalid_hop) = result
.hops
.iter()
.find(|hop| hop.outcome.invalidates_measurement())
{
if let Some(invalid_hop) = result.hops.iter().find(|hop| {
hop.outcome.invalidates_measurement()
|| hop.tcp_diagnostics.as_ref().is_some_and(|diagnostics| {
diagnostics.connect_error.is_some() || diagnostics.send_error.is_some()
})
}) {
warn!(
"DPI hop measurement for {} is invalid: {:?} at TTL {}",
result.target, invalid_hop.outcome, invalid_hop.ttl
@@ -927,17 +939,20 @@ mod tests {
ttl: 1,
router: Some("192.0.2.1".parse().unwrap()),
outcome: reports::probe::DpiProbeHopOutcome::IcmpTimeExceeded,
tcp_diagnostics: None,
}],
dpi_hops_v6: vec![
reports::probe::DpiProbeHop {
ttl: 1,
router: Some("2001:db8::1".parse().unwrap()),
outcome: reports::probe::DpiProbeHopOutcome::IcmpTimeExceeded,
tcp_diagnostics: None,
},
reports::probe::DpiProbeHop {
ttl: 2,
router: None,
outcome: reports::probe::DpiProbeHopOutcome::TcpClosed,
tcp_diagnostics: None,
},
],
};
@@ -976,6 +991,7 @@ mod tests {
ttl: 5,
router: None,
outcome,
tcp_diagnostics: None,
}],
};
@@ -994,9 +1010,30 @@ mod tests {
ttl: 1,
router: None,
outcome: dpi_hop::DpiHopProbeHopOutcome::Timeout,
tcp_diagnostics: None,
}],
};
assert_eq!(dpi_hop_from_result(&result), Some(0));
}
#[test]
fn manual_connect_error_is_not_reported_as_dpi_hop_zero() {
let result = dpi_hop::DpiHopProbeResult {
target: "192.0.2.1:443".parse().unwrap(),
local_addr: "0.0.0.0:0".parse().unwrap(),
client_hello_bytes: 0,
max_icmp_time_exceeded_ttl: None,
hops: vec![reports::probe::DpiProbeHop {
ttl: 1,
router: None,
outcome: reports::probe::DpiProbeHopOutcome::Timeout,
tcp_diagnostics: Some(reports::probe::TcpDiagnostics {
connect_error: Some("connect timed out".into()),
..Default::default()
}),
}],
};
assert_eq!(dpi_hop_from_result(&result), None);
}
}
+395
View File
@@ -0,0 +1,395 @@
use etherparse::{LaxNetSlice, LaxSlicedPacket, TcpOptionElement, TcpSlice, TransportSlice};
use reports::probe::{TcpDiagnostics, TcpPacketMetadata};
use socket2::{Domain, Protocol, Socket, Type};
use std::io;
use std::mem::MaybeUninit;
use std::net::{IpAddr, SocketAddr};
use std::sync::{
Arc,
atomic::{AtomicBool, Ordering},
};
use std::thread::JoinHandle;
use std::time::{Duration, Instant};
const MAX_PACKETS: usize = 32;
/// A separate reader observes TCP headers without consuming the stream's data.
/// It runs before connect so a SYN-ACK or an immediate reset is not missed.
pub struct TcpCapture {
started: Instant,
stop: Arc<AtomicBool>,
reader: Option<JoinHandle<TcpDiagnostics>>,
error: Option<String>,
}
impl TcpCapture {
pub fn start(target: SocketAddr) -> Self {
let started = Instant::now();
let stop = Arc::new(AtomicBool::new(false));
let socket = (|| {
let domain = if target.is_ipv4() {
Domain::IPV4
} else {
Domain::IPV6
};
let socket = Socket::new(domain, Type::RAW, Some(Protocol::TCP))?;
socket.set_read_timeout(Some(Duration::from_millis(20)))?;
Ok::<_, io::Error>(socket)
})();
match socket {
Ok(socket) => {
let reader_stop = stop.clone();
let reader = std::thread::spawn(move || {
let mut result = TcpDiagnostics::default();
let mut buffer = vec![MaybeUninit::<u8>::uninit(); 65535];
let mut shutdown_reads = 0;
loop {
if reader_stop.load(Ordering::Relaxed) {
shutdown_reads += 1;
if shutdown_reads > 256 {
result.truncated = true;
break;
}
if let Err(e) = socket.set_nonblocking(true) {
result.capture_error = Some(e.to_string());
break;
}
}
match socket.recv_from(&mut buffer) {
Ok((len, source)) => {
// SAFETY: recv_from initialized the returned prefix.
let bytes = unsafe {
std::slice::from_raw_parts(buffer.as_ptr().cast::<u8>(), len)
};
if let Some(packet) = parse_packet(
bytes,
source.as_socket(),
target,
started.elapsed().as_millis() as u64,
) {
// Filter to the exact local port(s) at finish. Bound the
// temporary queue too, in case other traces run concurrently.
if result.packets.len() < 256 {
result.packets.push(packet);
} else {
result.truncated = true;
}
}
}
Err(e)
if matches!(
e.kind(),
io::ErrorKind::WouldBlock | io::ErrorKind::TimedOut
) =>
{
if reader_stop.load(Ordering::Relaxed) {
break;
}
}
Err(e) if e.kind() == io::ErrorKind::Interrupted => {}
Err(e) => {
result.capture_error = Some(e.to_string());
break;
}
}
}
result
});
Self {
started,
stop,
reader: Some(reader),
error: None,
}
}
Err(e) => Self {
started,
stop,
reader: None,
error: Some(e.to_string()),
},
}
}
pub fn elapsed_ms(&self) -> u64 {
self.started.elapsed().as_millis() as u64
}
pub fn finish(mut self, ports: &[u16]) -> TcpDiagnostics {
self.stop.store(true, Ordering::Relaxed);
let mut result = match self.reader.take() {
Some(reader) => reader.join().unwrap_or_else(|_| TcpDiagnostics {
capture_error: Some("TCP capture reader panicked".into()),
..Default::default()
}),
None => TcpDiagnostics {
capture_error: self.error.take(),
..Default::default()
},
};
result
.packets
.retain(|p| ports.contains(&p.destination.port()));
if result.packets.len() > MAX_PACKETS {
result.packets.truncate(MAX_PACKETS);
result.truncated = true;
}
result
}
}
/// Binding before connecting preserves the local port even when connect fails.
pub fn connect_with_port(
target: SocketAddr,
timeout: Duration,
) -> (Option<u16>, io::Result<std::net::TcpStream>) {
let mut port = None;
let result = (|| {
let domain = if target.is_ipv4() {
Domain::IPV4
} else {
Domain::IPV6
};
let socket = Socket::new(domain, Type::STREAM, Some(Protocol::TCP))?;
let unspecified = if target.is_ipv4() {
IpAddr::V4(std::net::Ipv4Addr::UNSPECIFIED)
} else {
IpAddr::V6(std::net::Ipv6Addr::UNSPECIFIED)
};
socket.bind(&SocketAddr::new(unspecified, 0).into())?;
port = socket.local_addr()?.as_socket().map(|a| a.port());
socket.connect_timeout(&target.into(), timeout)?;
Ok(socket.into())
})();
(port, result)
}
impl Drop for TcpCapture {
fn drop(&mut self) {
self.stop.store(true, Ordering::Relaxed);
if let Some(reader) = self.reader.take() {
let _ = reader.join();
}
}
}
fn parse_packet(
bytes: &[u8],
source: Option<SocketAddr>,
target: SocketAddr,
observed_ms: u64,
) -> Option<TcpPacketMetadata> {
// Linux raw IPv6 sockets deliver the TCP segment without an IPv6 header.
// recv_from supplies its source; the destination IP is unavailable there.
let (src, dst, ttl, ip_id, tcp) =
if target.is_ipv6() && source.is_some_and(|s| s.ip() == target.ip()) {
match LaxSlicedPacket::from_ip(bytes) {
Ok(packet) if matches!(packet.net, Some(LaxNetSlice::Ipv6(_))) => ip_tcp(packet)?,
_ => (
target.ip(),
IpAddr::V6(std::net::Ipv6Addr::UNSPECIFIED),
None,
None,
TcpSlice::from_slice(bytes).ok()?,
),
}
} else {
ip_tcp(LaxSlicedPacket::from_ip(bytes).ok()?)?
};
if src != target.ip() || tcp.source_port() != target.port() {
return None;
}
let mut flags = Vec::new();
for (set, name) in [
(tcp.syn(), "SYN"),
(tcp.ack(), "ACK"),
(tcp.rst(), "RST"),
(tcp.fin(), "FIN"),
(tcp.psh(), "PSH"),
(tcp.urg(), "URG"),
(tcp.ece(), "ECE"),
(tcp.cwr(), "CWR"),
] {
if set {
flags.push(name.to_owned());
}
}
let mut timestamp = None;
let mut timestamp_echo = None;
for option in tcp.options_iterator().flatten() {
if let TcpOptionElement::Timestamp(value, echo) = option {
timestamp = Some(value);
timestamp_echo = Some(echo);
}
}
Some(TcpPacketMetadata {
observed_ms,
source: SocketAddr::new(src, tcp.source_port()),
destination: SocketAddr::new(dst, tcp.destination_port()),
ttl,
ip_id,
sequence: tcp.sequence_number(),
acknowledgment: tcp.acknowledgment_number(),
window: tcp.window_size(),
flags,
timestamp,
timestamp_echo,
options_hex: tcp.options().iter().map(|b| format!("{b:02x}")).collect(),
payload_bytes: tcp.payload().len(),
})
}
#[cfg(test)]
mod tests {
use super::*;
use etherparse::PacketBuilder;
#[test]
fn captures_reset_fingerprint_without_payload() {
let mut bytes = Vec::new();
PacketBuilder::ipv4([192, 0, 2, 1], [192, 0, 2, 2], 117)
.tcp(443, 45000, 1234, 2048)
.rst()
.ack(5678)
.options(&[TcpOptionElement::Timestamp(100, 200)])
.unwrap()
.write(&mut bytes, b"private payload")
.unwrap();
bytes[4..6].copy_from_slice(&42u16.to_be_bytes());
let packet = parse_packet(&bytes, None, "192.0.2.1:443".parse().unwrap(), 12).unwrap();
assert_eq!(packet.flags, ["ACK", "RST"]);
assert_eq!(packet.ttl, Some(117));
assert_eq!(packet.ip_id, Some(42));
assert_eq!(
(packet.sequence, packet.acknowledgment, packet.window),
(1234, 5678, 2048)
);
assert_eq!(
(packet.timestamp, packet.timestamp_echo),
(Some(100), Some(200))
);
assert_eq!(packet.payload_bytes, 15);
let json = serde_json::to_string(&packet).unwrap();
assert!(!json.contains("private payload"));
assert!(parse_packet(&bytes, None, "192.0.2.3:443".parse().unwrap(), 0).is_none());
assert!(parse_packet(&bytes, None, "192.0.2.1:444".parse().unwrap(), 0).is_none());
assert!(parse_packet(&bytes[..22], None, "192.0.2.1:443".parse().unwrap(), 0).is_none());
}
#[test]
fn captures_ipv6_fin_with_or_without_ip_header() {
let target: SocketAddr = "[2001:db8::1]:443".parse().unwrap();
let dst: std::net::Ipv6Addr = "2001:db8::2".parse().unwrap();
let IpAddr::V6(src) = target.ip() else {
unreachable!()
};
let mut bytes = Vec::new();
PacketBuilder::ipv6(src.octets(), dst.octets(), 55)
.tcp(443, 45000, 9, 123)
.fin()
.ack(10)
.write(&mut bytes, &[])
.unwrap();
let full = parse_packet(&bytes, Some(target), target, 20).unwrap();
assert_eq!(full.ttl, Some(55));
assert_eq!(full.destination.ip(), IpAddr::V6(dst));
assert_eq!(full.flags, ["ACK", "FIN"]);
let bare = parse_packet(&bytes[40..], Some(target), target, 21).unwrap();
assert_eq!(bare.ttl, None);
assert_eq!(bare.sequence, full.sequence);
assert!(bare.destination.ip().is_unspecified());
}
#[test]
fn legacy_hops_remain_compatible() {
let json = r#"{"ttl":1,"src":null,"outcome":"tcp_closed"}"#;
let hop: reports::probe::DpiProbeHop = serde_json::from_str(json).unwrap();
assert!(hop.tcp_diagnostics.is_none());
assert_eq!(serde_json::to_string(&hop).unwrap(), json);
let diagnostics = TcpDiagnostics {
packets: vec![],
capture_error: Some("permission denied".into()),
..Default::default()
};
let hop = reports::probe::DpiProbeHop {
tcp_diagnostics: Some(diagnostics),
..hop
};
let restored: reports::probe::DpiProbeHop =
serde_json::from_str(&serde_json::to_string(&hop).unwrap()).unwrap();
assert_eq!(restored.tcp_diagnostics, hop.tcp_diagnostics);
}
#[test]
#[ignore = "requires CAP_NET_RAW and local socket access"]
fn live_capture_records_syn_ack_and_reset() {
use std::io::{Read, Write};
let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
let target = listener.local_addr().unwrap();
let server = std::thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
stream.write_all(b"x").unwrap();
std::thread::sleep(Duration::from_millis(50));
socket2::SockRef::from(&stream)
.set_linger(Some(Duration::ZERO))
.unwrap();
});
let capture = TcpCapture::start(target);
let (port, connection) = connect_with_port(target, Duration::from_secs(1));
let mut stream = connection.unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(1)))
.unwrap();
let mut byte = [0];
stream.read_exact(&mut byte).unwrap();
assert_eq!(
stream.read(&mut byte).unwrap_err().kind(),
io::ErrorKind::ConnectionReset
);
server.join().unwrap();
let result = capture.finish(&[port.unwrap()]);
assert_eq!(result.capture_error, None);
assert!(
result
.packets
.iter()
.any(|p| p.flags.contains(&"SYN".into()) && p.flags.contains(&"ACK".into()))
);
assert!(
result
.packets
.iter()
.any(|p| p.flags.contains(&"RST".into()))
);
assert!(
result
.packets
.iter()
.all(|p| p.ttl.is_some() && p.source == target)
);
}
}
fn ip_tcp(
packet: LaxSlicedPacket<'_>,
) -> Option<(IpAddr, IpAddr, Option<u8>, Option<u16>, TcpSlice<'_>)> {
let (src, dst, ttl, id) = match packet.net? {
LaxNetSlice::Ipv4(ip) => (
IpAddr::V4(ip.header().source_addr()),
IpAddr::V4(ip.header().destination_addr()),
Some(ip.header().ttl()),
Some(ip.header().identification()),
),
LaxNetSlice::Ipv6(ip) => (
IpAddr::V6(ip.header().source_addr()),
IpAddr::V6(ip.header().destination_addr()),
Some(ip.header().hop_limit()),
None,
),
_ => return None,
};
let Some(TransportSlice::Tcp(tcp)) = packet.transport else {
return None;
};
Some((src, dst, ttl, id, tcp))
}
+5
View File
@@ -1,3 +1,4 @@
use crate::packet_capture::TcpCapture;
use etherparse::{
Icmpv4Type, Icmpv6Slice, Icmpv6Type, IpNumber, LaxNetSlice, LaxSlicedPacket, TransportSlice,
icmpv4, icmpv6,
@@ -74,6 +75,7 @@ fn trace_all_blocking(
let receiver = Socket::new(domain, Type::RAW, Some(protocol))?;
let mut hops = Vec::new();
for ttl in 1..=max_hops {
let capture = TcpCapture::start(SocketAddr::new(target, HTTPS_PORT));
let attempts = connect_attempts(target, ttl, retries)?;
let response = if attempts.is_empty() {
HopResponse::Timeout
@@ -95,6 +97,9 @@ fn trace_all_blocking(
address,
reverse_names: Vec::new(),
outcome,
tcp_diagnostics: Some(
capture.finish(&attempts.iter().map(|(_, port)| *port).collect::<Vec<_>>()),
),
});
if complete {
break;
+36
View File
@@ -26,6 +26,40 @@ pub struct DpiProbeHop {
#[serde(rename = "src")]
pub router: Option<IpAddr>,
pub outcome: DpiProbeHopOutcome,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub tcp_diagnostics: Option<TcpDiagnostics>,
}
/// Incoming TCP headers observed during one manual hop attempt. No payload is retained.
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Default)]
#[serde(default)]
pub struct TcpDiagnostics {
pub packets: Vec<TcpPacketMetadata>,
pub capture_error: Option<String>,
pub truncated: bool,
pub connect_error: Option<String>,
pub send_error: Option<String>,
pub client_hello_sent_ms: Option<u64>,
pub junk_sent_ms: Option<u64>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct TcpPacketMetadata {
/// Approximate observation time relative to the start of this hop's capture.
pub observed_ms: u64,
pub source: std::net::SocketAddr,
pub destination: std::net::SocketAddr,
/// Absent when the platform's raw IPv6 socket omits the IP header.
pub ttl: Option<u8>,
pub ip_id: Option<u16>,
pub sequence: u32,
pub acknowledgment: u32,
pub window: u16,
pub flags: Vec<String>,
pub timestamp: Option<u32>,
pub timestamp_echo: Option<u32>,
pub options_hex: String,
pub payload_bytes: usize,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
@@ -243,6 +277,8 @@ pub struct ManualTracerouteHop {
pub address: Option<IpAddr>,
pub reverse_names: Vec<String>,
pub outcome: ManualTracerouteOutcome,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub tcp_diagnostics: Option<TcpDiagnostics>,
}
#[derive(Clone, Copy, Serialize, Deserialize)]
+1
View File
@@ -143,6 +143,7 @@ mod tests {
ttl: 1,
router: Some("192.0.2.1".parse().unwrap()),
outcome: reports::probe::DpiProbeHopOutcome::IcmpTimeExceeded,
tcp_diagnostics: None,
}];
let probes = vec![
ProbeMetadata {