diff --git a/Cargo.lock b/Cargo.lock index ad48ddf..5b77e1f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2488,7 +2488,7 @@ dependencies = [ [[package]] name = "probe" -version = "0.7.1" +version = "0.7.2" dependencies = [ "anyhow", "clap", diff --git a/frontend/src/routes/admin/probes/+page.svelte b/frontend/src/routes/admin/probes/+page.svelte index b195f20..df641cf 100644 --- a/frontend/src/routes/admin/probes/+page.svelte +++ b/frontend/src/routes/admin/probes/+page.svelte @@ -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 = { }; +{#snippet tcpDiagnostics(diagnostics: TcpDiagnostics | undefined)} + {#if diagnostics} +
packet.flags.includes("RST") || packet.flags.includes("FIN"))} + > + + TCP: {diagnostics.packets.length} входящих пакетов{diagnostics.truncated ? " (список сокращён)" : ""} + + {#if diagnostics.capture_error} +

+ Ошибка захвата: {diagnostics.capture_error} +

+ {/if} + {#if diagnostics.connect_error} +

+ Ошибка соединения: {diagnostics.connect_error} +

+ {/if} + {#if diagnostics.send_error} +

+ Ошибка отправки: {diagnostics.send_error} +

+ {/if} +

+ Время приблизительное, от начала захвата. + {#if diagnostics.client_hello_sent_ms !== null} + ClientHello: {diagnostics.client_hello_sent_ms} мс; данные: + {diagnostics.junk_sent_ms ?? "—"} + мс. + {/if} + Δ TTL сравнивает пакет с SYN-ACK того же соединения; это не расстояние + до DPI. +

+
+ + + + + + + + + + + + + + + {#each diagnostics.packets as packet} + {@const baseline = diagnostics.packets.find((p) => p.destination === packet.destination && p.flags.includes("SYN") && p.flags.includes("ACK"))} + + + + + + + + + + + {/each} + +
мсИсточник → цельФлагиTTL / ΔSEQ / ACKОкно / IP IDTS / echoБайты / опции
{packet.observed_ms}{packet.source} → {packet.destination}{packet.flags.join(" ") || "—"} + {packet.ttl ?? "—"} + / + {packet.ttl !== null && baseline?.ttl != null ? packet.ttl - baseline.ttl : "—"} + {packet.sequence}/ {packet.acknowledgment}{packet.window} / {packet.ip_id ?? "—"} + {packet.timestamp ?? "—"} + / {packet.timestamp_echo ?? "—"} + {packet.payload_bytes} / {packet.options_hex || "—"}
+
+ {#if diagnostics.packets.some((p) => p.ttl === null)} +

+ Для части пакетов IP-заголовок недоступен; TTL не измерен. +

+ {/if} +
+ {/if} +{/snippet} + {#snippet dpiHops(hops: DpiHop[])}
    {#each hops as hop} @@ -404,6 +518,11 @@ const outcomeLabel: Record = { {outcomeLabel[hop.outcome] ?? hop.outcome} + {#if hop.tcp_diagnostics} +
    + {@render tcpDiagnostics(hop.tcp_diagnostics)} +
    + {/if} {/each}
@@ -782,9 +901,10 @@ PROBE_TOKEN={created.token} >

- TCP-соединение к хосту на порту 443, ClientHello с указанным SNI, - затем пакеты с возрастающим TTL. Хост разрешается сканером; - используется первый IP адрес. + Для каждого TTL создаётся новое TCP-соединение к хосту на порту 443. + ClientHello с указанным SNI отправляется с обычным TTL; через 100 мс + отправляются случайные данные с проверяемым TTL. Хост разрешается + сканером; используется первый IP адрес.

{/if}
@@ -823,6 +943,15 @@ PROBE_TOKEN={created.token} {#if busy.endsWith(":command")}

Ожидание ответа сканера…

{/if} + {#if result && result.type !== "error"} + + {/if} {#if result?.type === "error"}

{result.message}

{:else if result?.type === "resubscribe_tasks"} @@ -875,6 +1004,11 @@ PROBE_TOKEN={created.token} {outcomeLabel[hop.outcome] ?? hop.outcome} + {#if hop.tcp_diagnostics} +
+ {@render tcpDiagnostics(hop.tcp_diagnostics)} +
+ {/if} {/each} diff --git a/probe/Cargo.toml b/probe/Cargo.toml index cd21f3b..ed57d5e 100644 --- a/probe/Cargo.toml +++ b/probe/Cargo.toml @@ -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" diff --git a/probe/src/dpi_hop.rs b/probe/src/dpi_hop.rs index efa53ef..ad404d4 100644 --- a/probe/src/dpi_hop.rs +++ b/probe/src/dpi_hop.rs @@ -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 tcp, + Err(error) => { + let Some(capture) = capture else { + return Err(error); + }; + let mut diagnostics = capture.finish(&port.into_iter().collect::>()); + 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 io::Result io::Result, + local: SocketAddr, + hello_ms: Option, + junk_ms: Option, +) -> Option { + 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> { 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, } } diff --git a/probe/src/main.rs b/probe/src/main.rs index 6126c9d..9b2c1c9 100644 --- a/probe/src/main.rs +++ b/probe/src/main.rs @@ -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 { 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); + } } diff --git a/probe/src/packet_capture.rs b/probe/src/packet_capture.rs new file mode 100644 index 0000000..d47f68e --- /dev/null +++ b/probe/src/packet_capture.rs @@ -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, + reader: Option>, + error: Option, +} + +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::::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::(), 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, io::Result) { + 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, + target: SocketAddr, + observed_ms: u64, +) -> Option { + // 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, Option, 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)) +} diff --git a/probe/src/traceroute.rs b/probe/src/traceroute.rs index 7c6f359..0e2ac14 100644 --- a/probe/src/traceroute.rs +++ b/probe/src/traceroute.rs @@ -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::>()), + ), }); if complete { break; diff --git a/reports/src/probe.rs b/reports/src/probe.rs index 1e34fa8..71f0652 100644 --- a/reports/src/probe.rs +++ b/reports/src/probe.rs @@ -26,6 +26,40 @@ pub struct DpiProbeHop { #[serde(rename = "src")] pub router: Option, pub outcome: DpiProbeHopOutcome, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub tcp_diagnostics: Option, +} + +/// 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, + pub capture_error: Option, + pub truncated: bool, + pub connect_error: Option, + pub send_error: Option, + pub client_hello_sent_ms: Option, + pub junk_sent_ms: Option, +} + +#[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, + pub ip_id: Option, + pub sequence: u32, + pub acknowledgment: u32, + pub window: u16, + pub flags: Vec, + pub timestamp: Option, + pub timestamp_echo: Option, + 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, pub reverse_names: Vec, pub outcome: ManualTracerouteOutcome, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub tcp_diagnostics: Option, } #[derive(Clone, Copy, Serialize, Deserialize)] diff --git a/website/src/api/node_stats.rs b/website/src/api/node_stats.rs index cab4b71..015fc0f 100644 --- a/website/src/api/node_stats.rs +++ b/website/src/api/node_stats.rs @@ -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 {