mirror of
https://github.com/LowderPlay/cheburcheck.git
synced 2026-10-04 12:48:06 +03:00
feat: dpi traceroute commands and debug (#102)
This commit is contained in:
Generated
+2
-2
@@ -2488,7 +2488,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "probe"
|
||||
version = "0.6.8"
|
||||
version = "0.7.0"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"clap",
|
||||
@@ -4509,7 +4509,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "website"
|
||||
version = "1.4.3"
|
||||
version = "1.4.4"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"dotenvy",
|
||||
|
||||
@@ -40,9 +40,30 @@ type Hop = {
|
||||
reverse_names: string[];
|
||||
outcome: string;
|
||||
};
|
||||
type DpiHop = { ttl: number; src: string | null; outcome: string };
|
||||
type CommandType =
|
||||
| "resubscribe_tasks"
|
||||
| "traceroute"
|
||||
| "remeasure_dpi_hop"
|
||||
| "sni_traceroute";
|
||||
type CommandResult =
|
||||
| { type: "traceroute"; target: string; hops: Hop[] }
|
||||
| { type: "resubscribe_tasks"; requested: boolean }
|
||||
| {
|
||||
type: "remeasure_dpi_hop";
|
||||
dpi_hop_v4: number | null;
|
||||
dpi_hop_v6: number | null;
|
||||
dpi_hops_v4: DpiHop[];
|
||||
dpi_hops_v6: DpiHop[];
|
||||
}
|
||||
| {
|
||||
type: "sni_traceroute";
|
||||
host: string;
|
||||
sni: string;
|
||||
target: string;
|
||||
dpi_hop: number | null;
|
||||
hops: DpiHop[];
|
||||
}
|
||||
| { type: "error"; message: string };
|
||||
type Form = Pick<
|
||||
Probe,
|
||||
@@ -77,6 +98,8 @@ let created = $state<{ id: number; token: string } | null>(null);
|
||||
let selected = $state<number | null>(null);
|
||||
let target = $state("");
|
||||
let maxHops = $state(30);
|
||||
let traceType = $state<"traceroute" | "sni_traceroute">("traceroute");
|
||||
let sni = $state("");
|
||||
let result = $state<CommandResult | null>(null);
|
||||
|
||||
const queryClient = useQueryClient();
|
||||
@@ -133,16 +156,22 @@ const sendProbeCommand = createMutation(() => ({
|
||||
type,
|
||||
target,
|
||||
maxHops,
|
||||
sni,
|
||||
}: {
|
||||
id: number;
|
||||
type: "resubscribe_tasks" | "traceroute";
|
||||
type: CommandType;
|
||||
target: string;
|
||||
maxHops: number;
|
||||
sni: string;
|
||||
}) =>
|
||||
api<CommandResult>(
|
||||
"/probes/" + id + "/commands",
|
||||
"POST",
|
||||
type === "traceroute" ? { type, target, max_hops: maxHops } : { type },
|
||||
type === "traceroute"
|
||||
? { type, target, max_hops: maxHops }
|
||||
: type === "sni_traceroute"
|
||||
? { type, host: target, sni, max_hops: maxHops }
|
||||
: { type },
|
||||
),
|
||||
}));
|
||||
$effect(() => {
|
||||
@@ -324,7 +353,7 @@ async function updateCheck(id: number | null) {
|
||||
busy = "";
|
||||
}
|
||||
}
|
||||
async function command(id: number, type: "resubscribe_tasks" | "traceroute") {
|
||||
async function command(id: number, type: CommandType) {
|
||||
selected = id;
|
||||
result = null;
|
||||
busy = `${id}:command`;
|
||||
@@ -335,7 +364,11 @@ async function command(id: number, type: "resubscribe_tasks" | "traceroute") {
|
||||
type,
|
||||
target: target.trim(),
|
||||
maxHops,
|
||||
sni: sni.trim(),
|
||||
});
|
||||
if (result.type === "remeasure_dpi_hop") {
|
||||
await queryClient.invalidateQueries({ queryKey: probeKey });
|
||||
}
|
||||
} catch (e) {
|
||||
error = String(e instanceof Error ? e.message : e);
|
||||
} finally {
|
||||
@@ -354,10 +387,33 @@ const outcomeLabel: Record<string, string> = {
|
||||
icmp_time_exceeded: "Промежуточный узел",
|
||||
rst: "TCP RST",
|
||||
connected: "Подключено",
|
||||
tcp_closed: "TCP закрыт",
|
||||
tcp_acknowledged: "TCP подтверждён",
|
||||
timeout: "Нет ответа",
|
||||
};
|
||||
</script>
|
||||
|
||||
{#snippet dpiHops(hops: DpiHop[])}
|
||||
<ol class="space-y-1">
|
||||
{#each hops as hop}
|
||||
<li
|
||||
class="grid grid-cols-[2.5rem_1fr_auto] gap-3 rounded-lg border border-neutral-800 bg-black/20 px-3 py-2 text-sm"
|
||||
>
|
||||
<span class="font-mono text-neutral-500">{hop.ttl}</span>
|
||||
<span class="font-mono">{hop.src ?? "* * *"}</span>
|
||||
<span class="text-xs text-neutral-500"
|
||||
>{outcomeLabel[hop.outcome] ?? hop.outcome}</span
|
||||
>
|
||||
</li>
|
||||
{/each}
|
||||
</ol>
|
||||
{#if hops.length === 0}
|
||||
<p class="text-sm text-neutral-500">
|
||||
Измерение не удалось: ответов нет. Подробности в логах сканера.
|
||||
</p>
|
||||
{/if}
|
||||
{/snippet}
|
||||
|
||||
<svelte:head
|
||||
><title>Сканеры · Cheburcheck</title>
|
||||
<meta name="robots" content="noindex, nofollow"></svelte:head
|
||||
@@ -494,7 +550,7 @@ const outcomeLabel: Record<string, string> = {
|
||||
{probe.version ?? "Версия неизвестна"}
|
||||
{probe.bundle_type ? `· ${probe.bundle_type}` : ""}
|
||||
</div>
|
||||
{#if probe.dpi_hop_v4 || probe.dpi_hop_v6}
|
||||
{#if probe.dpi_hop_v4 !== null || probe.dpi_hop_v6 !== null}
|
||||
<div class="mt-1 text-xs text-neutral-500">
|
||||
v4: <b>{probe.dpi_hop_v4 ?? "—"}</b>; v6:
|
||||
<b>{probe.dpi_hop_v6 ?? "—"}</b>
|
||||
@@ -559,6 +615,14 @@ const outcomeLabel: Record<string, string> = {
|
||||
>
|
||||
<RefreshCw size={14} />
|
||||
Переподписать
|
||||
</button><button
|
||||
type="button"
|
||||
class="link"
|
||||
disabled={busy !== "" || !probe.online}
|
||||
onclick={() => void command(probe.id, "remeasure_dpi_hop")}
|
||||
>
|
||||
<RefreshCw size={14} />
|
||||
Перемерить DPI hop
|
||||
</button><button
|
||||
type="button"
|
||||
class="link"
|
||||
@@ -688,8 +752,41 @@ PROBE_TOKEN={created.token}</pre>
|
||||
<section id="commands" class="panel">
|
||||
<h2 class="mb-2 text-lg font-semibold">Команды и трассировка</h2>
|
||||
<p class="mb-5 text-sm text-neutral-400">
|
||||
Команда выполняется выбранным сканером. Ответ может занять до минуты.
|
||||
Команда выполняется выбранным сканером. DPI и SNI измерения могут занять
|
||||
несколько минут.
|
||||
</p>
|
||||
<div class="mb-4 flex flex-wrap items-end gap-3">
|
||||
<label class="field"
|
||||
>Режим трассировки
|
||||
<select class="input" bind:value={traceType}>
|
||||
<option value="traceroute">TCP traceroute</option>
|
||||
<option value="sni_traceroute">SNI traceroute (DPI)</option>
|
||||
</select>
|
||||
</label>
|
||||
<button
|
||||
type="button"
|
||||
class="btn"
|
||||
disabled={selected === null || busy !== ""}
|
||||
onclick={() => selected !== null && void command(selected, "remeasure_dpi_hop")}
|
||||
>
|
||||
Перемерить DPI hop
|
||||
</button>
|
||||
</div>
|
||||
{#if traceType === "sni_traceroute"}
|
||||
<label class="field mb-4"
|
||||
>SNI<input
|
||||
class="input"
|
||||
type="text"
|
||||
placeholder="rutracker.org"
|
||||
bind:value={sni}
|
||||
></label
|
||||
>
|
||||
<p class="mb-4 text-sm text-neutral-400">
|
||||
TCP-соединение к хосту на порту 443, ClientHello с указанным SNI,
|
||||
затем пакеты с возрастающим TTL. Хост разрешается сканером;
|
||||
используется первый IP адрес.
|
||||
</p>
|
||||
{/if}
|
||||
<div class="grid gap-3 sm:grid-cols-[1fr_2fr_100px_auto] sm:items-end">
|
||||
<label class="field"
|
||||
>Сканер<select class="input" bind:value={selected}>
|
||||
@@ -699,7 +796,7 @@ PROBE_TOKEN={created.token}</pre>
|
||||
{/each}
|
||||
</select></label
|
||||
><label class="field"
|
||||
>IP адрес цели<input
|
||||
>{traceType === "sni_traceroute" ? "Хост или IP адрес цели" : "IP адрес цели"}<input
|
||||
class="input"
|
||||
type="text"
|
||||
placeholder="1.1.1.1"
|
||||
@@ -716,8 +813,8 @@ PROBE_TOKEN={created.token}</pre>
|
||||
><button
|
||||
type="button"
|
||||
class="btn-primary flex items-center justify-center gap-2"
|
||||
disabled={selected === null || !target.trim() || busy !== ""}
|
||||
onclick={() => selected !== null && void command(selected, "traceroute")}
|
||||
disabled={selected === null || !target.trim() || (traceType === "sni_traceroute" && !sni.trim()) || !Number.isInteger(maxHops) || maxHops < 1 || maxHops > 64 || busy !== ""}
|
||||
onclick={() => selected !== null && void command(selected, traceType)}
|
||||
>
|
||||
<Route size={16} />
|
||||
Запустить
|
||||
@@ -732,6 +829,29 @@ PROBE_TOKEN={created.token}</pre>
|
||||
<p class="mt-5 text-sm text-emerald-300">
|
||||
Запрос на переподписку получен сканером.
|
||||
</p>
|
||||
{:else if result?.type === "remeasure_dpi_hop"}
|
||||
<p class="mt-5 text-sm text-emerald-300">
|
||||
DPI hop перемерен: v4 {result.dpi_hop_v4 ?? "—"}; v6
|
||||
{result.dpi_hop_v6 ?? "—"}.
|
||||
</p>
|
||||
<div class="mt-4 grid gap-4 sm:grid-cols-2">
|
||||
{#each [{ label: "IPv4", hops: result.dpi_hops_v4 }, { label: "IPv6", hops: result.dpi_hops_v6 }] as trace}
|
||||
<div>
|
||||
<h3 class="mb-2 font-semibold">{trace.label}</h3>
|
||||
{@render dpiHops(trace.hops)}
|
||||
</div>
|
||||
{/each}
|
||||
</div>
|
||||
{:else if result?.type === "sni_traceroute"}
|
||||
<div class="mt-6">
|
||||
<h3 class="mb-2 font-semibold">
|
||||
SNI маршрут до {result.host} ({result.target})
|
||||
</h3>
|
||||
<p class="mb-4 text-sm text-neutral-400">
|
||||
SNI: {result.sni}; DPI hop: {result.dpi_hop ?? "—"}
|
||||
</p>
|
||||
{@render dpiHops(result.hops)}
|
||||
</div>
|
||||
{:else if result?.type === "traceroute"}
|
||||
<div class="mt-6">
|
||||
<h3 class="mb-4 font-semibold">Маршрут до {result.target}</h3>
|
||||
|
||||
+1
-1
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "probe"
|
||||
version = "0.6.8"
|
||||
version = "0.7.0"
|
||||
edition = "2024"
|
||||
license-file = "../LICENSE"
|
||||
description = "Dynamic network probe daemon for Cheburcheck"
|
||||
|
||||
+78
-26
@@ -22,6 +22,8 @@ pub struct DpiHopProbeConfig {
|
||||
pub max_ttl: u8,
|
||||
pub connect_timeout: Duration,
|
||||
pub hop_timeout: Duration,
|
||||
pub accept_rst_fin_as_timeout: bool,
|
||||
pub accept_ack_as_timeout: bool,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
@@ -33,26 +35,9 @@ pub struct DpiHopProbeResult {
|
||||
pub hops: Vec<DpiHopProbeHop>,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct DpiHopProbeHop {
|
||||
pub ttl: u8,
|
||||
pub router: Option<IpAddr>,
|
||||
pub outcome: DpiHopProbeHopOutcome,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub enum DpiHopProbeHopOutcome {
|
||||
IcmpTimeExceeded,
|
||||
Timeout,
|
||||
TcpClosed,
|
||||
TcpAcknowledged,
|
||||
}
|
||||
|
||||
impl DpiHopProbeHopOutcome {
|
||||
pub const fn invalidates_measurement(self) -> bool {
|
||||
matches!(self, Self::TcpClosed | Self::TcpAcknowledged)
|
||||
}
|
||||
}
|
||||
pub use reports::probe::{
|
||||
DpiProbeHop as DpiHopProbeHop, DpiProbeHopOutcome as DpiHopProbeHopOutcome,
|
||||
};
|
||||
|
||||
pub async fn detect_dpi_hop(config: DpiHopProbeConfig) -> io::Result<DpiHopProbeResult> {
|
||||
tokio::task::spawn_blocking(move || detect_dpi_hop_blocking(config))
|
||||
@@ -89,6 +74,7 @@ pub fn detect_dpi_hop_blocking(config: DpiHopProbeConfig) -> io::Result<DpiHopPr
|
||||
// affect later hops.
|
||||
let mut tcp = TcpStream::connect_timeout(&config.target, config.connect_timeout)?;
|
||||
tcp.set_nodelay(true)?;
|
||||
tcp.set_write_timeout(Some(config.connect_timeout))?;
|
||||
let local_addr = tcp.local_addr()?;
|
||||
if !same_ip_family(local_addr, config.target) {
|
||||
return Err(io::Error::other(
|
||||
@@ -128,14 +114,19 @@ pub fn detect_dpi_hop_blocking(config: DpiHopProbeConfig) -> io::Result<DpiHopPr
|
||||
hops.push(DpiHopProbeHop {
|
||||
ttl,
|
||||
router: None,
|
||||
outcome: DpiHopProbeHopOutcome::TcpClosed,
|
||||
outcome: classify_hop(None, true, false, &config),
|
||||
});
|
||||
if config.accept_rst_fin_as_timeout {
|
||||
continue;
|
||||
}
|
||||
break;
|
||||
}
|
||||
|
||||
let router =
|
||||
listen_for_time_exceeded(&icmp, local_addr, config.target, config.hop_timeout)?;
|
||||
let outcome = classify_hop(router, peer_closed(&tcp)?, tcp_payload_acknowledged(&tcp)?);
|
||||
let closed = peer_closed(&tcp)?;
|
||||
let acknowledged = !closed && tcp_payload_acknowledged(&tcp)?;
|
||||
let outcome = classify_hop(router, closed, acknowledged, &config);
|
||||
if outcome == DpiHopProbeHopOutcome::IcmpTimeExceeded {
|
||||
max_icmp_time_exceeded_ttl = Some(ttl);
|
||||
}
|
||||
@@ -278,13 +269,22 @@ fn classify_hop(
|
||||
router: Option<IpAddr>,
|
||||
peer_closed: bool,
|
||||
payload_acknowledged: bool,
|
||||
config: &DpiHopProbeConfig,
|
||||
) -> DpiHopProbeHopOutcome {
|
||||
if router.is_some() {
|
||||
DpiHopProbeHopOutcome::IcmpTimeExceeded
|
||||
} else if peer_closed {
|
||||
DpiHopProbeHopOutcome::TcpClosed
|
||||
if config.accept_rst_fin_as_timeout {
|
||||
DpiHopProbeHopOutcome::Timeout
|
||||
} else {
|
||||
DpiHopProbeHopOutcome::TcpClosed
|
||||
}
|
||||
} else if payload_acknowledged {
|
||||
DpiHopProbeHopOutcome::TcpAcknowledged
|
||||
if config.accept_ack_as_timeout {
|
||||
DpiHopProbeHopOutcome::Timeout
|
||||
} else {
|
||||
DpiHopProbeHopOutcome::TcpAcknowledged
|
||||
}
|
||||
} else {
|
||||
DpiHopProbeHopOutcome::Timeout
|
||||
}
|
||||
@@ -470,10 +470,22 @@ fn matching_quoted_tcp_tuple(packet: &[u8], local_addr: SocketAddr, target: Sock
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
fn probe_config() -> DpiHopProbeConfig {
|
||||
DpiHopProbeConfig {
|
||||
target: "192.0.2.1:443".parse().unwrap(),
|
||||
control_sni: "example.com".into(),
|
||||
max_ttl: 15,
|
||||
connect_timeout: Duration::from_secs(5),
|
||||
hop_timeout: Duration::from_secs(1),
|
||||
accept_rst_fin_as_timeout: false,
|
||||
accept_ack_as_timeout: false,
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn acknowledged_payload_marks_direct_tcp_delivery() {
|
||||
assert_eq!(
|
||||
classify_hop(None, false, true),
|
||||
classify_hop(None, false, true, &probe_config()),
|
||||
DpiHopProbeHopOutcome::TcpAcknowledged
|
||||
);
|
||||
assert!(DpiHopProbeHopOutcome::TcpAcknowledged.invalidates_measurement());
|
||||
@@ -482,11 +494,51 @@ mod tests {
|
||||
#[test]
|
||||
fn icmp_response_takes_precedence_over_tcp_state() {
|
||||
assert_eq!(
|
||||
classify_hop(Some(IpAddr::V4(Ipv4Addr::LOCALHOST)), true, true),
|
||||
classify_hop(
|
||||
Some(IpAddr::V4(Ipv4Addr::LOCALHOST)),
|
||||
true,
|
||||
true,
|
||||
&probe_config()
|
||||
),
|
||||
DpiHopProbeHopOutcome::IcmpTimeExceeded
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn tcp_timeout_options_are_independent_and_preserve_icmp_precedence() {
|
||||
for accept_rst_fin_as_timeout in [false, true] {
|
||||
for accept_ack_as_timeout in [false, true] {
|
||||
let config = DpiHopProbeConfig {
|
||||
accept_rst_fin_as_timeout,
|
||||
accept_ack_as_timeout,
|
||||
..probe_config()
|
||||
};
|
||||
let closed = if accept_rst_fin_as_timeout {
|
||||
DpiHopProbeHopOutcome::Timeout
|
||||
} else {
|
||||
DpiHopProbeHopOutcome::TcpClosed
|
||||
};
|
||||
let acknowledged = if accept_ack_as_timeout {
|
||||
DpiHopProbeHopOutcome::Timeout
|
||||
} else {
|
||||
DpiHopProbeHopOutcome::TcpAcknowledged
|
||||
};
|
||||
assert_eq!(classify_hop(None, true, false, &config), closed);
|
||||
assert_eq!(classify_hop(None, true, true, &config), closed);
|
||||
assert_eq!(classify_hop(None, false, true, &config), acknowledged);
|
||||
assert_eq!(
|
||||
classify_hop(None, false, false, &config),
|
||||
DpiHopProbeHopOutcome::Timeout
|
||||
);
|
||||
assert_eq!(
|
||||
classify_hop(Some(IpAddr::V4(Ipv4Addr::LOCALHOST)), true, true, &config),
|
||||
DpiHopProbeHopOutcome::IcmpTimeExceeded
|
||||
);
|
||||
assert!(!DpiHopProbeHopOutcome::Timeout.invalidates_measurement());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn matches_icmp_time_exceeded_quote_by_flow_tuple() {
|
||||
let local = SocketAddrV4::new(Ipv4Addr::new(192, 0, 2, 10), 45_000);
|
||||
|
||||
+217
-5
@@ -31,12 +31,16 @@ struct LoadedProbeConfig {
|
||||
config: ProbeConfig,
|
||||
dpi_hop_v4: Option<u8>,
|
||||
dpi_hop_v6: Option<u8>,
|
||||
dpi_hops_v4: Vec<reports::probe::DpiProbeHop>,
|
||||
dpi_hops_v6: Vec<reports::probe::DpiProbeHop>,
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Default)]
|
||||
#[derive(Clone, Default)]
|
||||
struct DpiHops {
|
||||
v4: Option<u8>,
|
||||
v6: Option<u8>,
|
||||
hops_v4: Vec<reports::probe::DpiProbeHop>,
|
||||
hops_v6: Vec<reports::probe::DpiProbeHop>,
|
||||
}
|
||||
|
||||
#[derive(Parser, Debug, Clone)]
|
||||
@@ -136,6 +140,8 @@ async fn main() -> Result<()> {
|
||||
bundle_type: Some(args.bundle_type),
|
||||
dpi_hop_v4: None,
|
||||
dpi_hop_v6: None,
|
||||
dpi_hops_v4: Vec::new(),
|
||||
dpi_hops_v6: Vec::new(),
|
||||
})?;
|
||||
|
||||
let mut options = MqttOptions::new(&args.probe_id, &args.mqtt_host, args.mqtt_port);
|
||||
@@ -203,6 +209,9 @@ async fn main() -> Result<()> {
|
||||
let probe_id = args.probe_id.clone();
|
||||
let payload = publish.payload.to_vec();
|
||||
let retries = args.traceroute_retries;
|
||||
let command_args = args.clone();
|
||||
let command_status_topic = status_topic.clone();
|
||||
let command_config = config.clone();
|
||||
let semaphore = task_semaphore.clone();
|
||||
tokio::spawn(async move {
|
||||
let result = match serde_json::from_slice::<ProbeCommand>(&payload) {
|
||||
@@ -233,6 +242,39 @@ async fn main() -> Result<()> {
|
||||
},
|
||||
}
|
||||
}
|
||||
Ok(
|
||||
command @ (ProbeCommand::RemeasureDpiHop
|
||||
| ProbeCommand::SniTraceroute { .. }),
|
||||
) => match semaphore.acquire_owned().await {
|
||||
Ok(_permit) => {
|
||||
match run_dpi_command(command, &command_config).await {
|
||||
Ok((result, status)) => {
|
||||
if let Some(hops) = status {
|
||||
if let Err(error) = publish_status(
|
||||
&client,
|
||||
&command_status_topic,
|
||||
&command_args,
|
||||
true,
|
||||
hops,
|
||||
)
|
||||
.await
|
||||
{
|
||||
warn!(
|
||||
"failed to publish remeasured DPI status: {error}"
|
||||
);
|
||||
}
|
||||
}
|
||||
result
|
||||
}
|
||||
Err(error) => ProbeCommandResult::Error {
|
||||
message: error.to_string(),
|
||||
},
|
||||
}
|
||||
}
|
||||
Err(error) => ProbeCommandResult::Error {
|
||||
message: error.to_string(),
|
||||
},
|
||||
},
|
||||
Err(error) => ProbeCommandResult::Error {
|
||||
message: format!("invalid command: {error}"),
|
||||
},
|
||||
@@ -300,6 +342,8 @@ async fn main() -> Result<()> {
|
||||
.map_or_else(DpiHops::default, |config| DpiHops {
|
||||
v4: config.dpi_hop_v4,
|
||||
v6: config.dpi_hop_v6,
|
||||
hops_v4: config.dpi_hops_v4.clone(),
|
||||
hops_v6: config.dpi_hops_v6.clone(),
|
||||
});
|
||||
publish_status(&client, &status_topic, &args, true, dpi_hops).await?;
|
||||
client.subscribe(CONFIG_TOPIC, QoS::AtLeastOnce).await?;
|
||||
@@ -422,6 +466,8 @@ async fn update_config(
|
||||
config: value,
|
||||
dpi_hop_v4: dpi_hops.v4,
|
||||
dpi_hop_v6: dpi_hops.v6,
|
||||
dpi_hops_v4: dpi_hops.hops_v4.clone(),
|
||||
dpi_hops_v6: dpi_hops.hops_v6.clone(),
|
||||
};
|
||||
debug!(
|
||||
"measured DPI hops: IPv4={:?}, IPv6={:?}",
|
||||
@@ -432,6 +478,98 @@ async fn update_config(
|
||||
Ok(dpi_hops)
|
||||
}
|
||||
|
||||
async fn run_dpi_command(
|
||||
command: ProbeCommand,
|
||||
config: &Arc<RwLock<Option<LoadedProbeConfig>>>,
|
||||
) -> Result<(ProbeCommandResult, Option<DpiHops>)> {
|
||||
match command {
|
||||
ProbeCommand::RemeasureDpiHop => {
|
||||
let snapshot = config
|
||||
.read()
|
||||
.await
|
||||
.clone()
|
||||
.context("probe config is not loaded")?;
|
||||
let dpi = snapshot
|
||||
.config
|
||||
.dpi_probe
|
||||
.as_ref()
|
||||
.context("DPI hop measurement is not configured")?;
|
||||
let hops = measure_dpi_hops(Some(dpi)).await;
|
||||
let mut guard = config.write().await;
|
||||
let loaded = guard.as_mut().context("probe config is not loaded")?;
|
||||
if loaded.config.version != snapshot.config.version
|
||||
|| loaded.config.published_at != snapshot.config.published_at
|
||||
|| serde_json::to_vec(&loaded.config.dpi_probe)?
|
||||
!= serde_json::to_vec(&snapshot.config.dpi_probe)?
|
||||
{
|
||||
bail!("probe config changed during measurement; retry the 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();
|
||||
Ok((
|
||||
ProbeCommandResult::RemeasureDpiHop {
|
||||
dpi_hop_v4: hops.v4,
|
||||
dpi_hop_v6: hops.v6,
|
||||
dpi_hops_v4: hops.hops_v4.clone(),
|
||||
dpi_hops_v6: hops.hops_v6.clone(),
|
||||
},
|
||||
Some(hops),
|
||||
))
|
||||
}
|
||||
ProbeCommand::SniTraceroute {
|
||||
host,
|
||||
sni,
|
||||
max_hops,
|
||||
} => {
|
||||
if !(1..=64).contains(&max_hops) {
|
||||
bail!("max_hops must be between 1 and 64");
|
||||
}
|
||||
reports::probe::validate_sni_traceroute(&host, &sni).map_err(anyhow::Error::msg)?;
|
||||
// Resolve on the probe so the route reflects its own network/DNS view.
|
||||
let target = tokio::time::timeout(
|
||||
Duration::from_secs(10),
|
||||
tokio::net::lookup_host((host.as_str(), 443)),
|
||||
)
|
||||
.await
|
||||
.context("host resolution timed out")??
|
||||
.next()
|
||||
.context("host resolved to no addresses")?;
|
||||
let dpi_config = config
|
||||
.read()
|
||||
.await
|
||||
.as_ref()
|
||||
.and_then(|loaded| loaded.config.dpi_probe.clone());
|
||||
let trace = dpi_hop::detect_dpi_hop(dpi_hop::DpiHopProbeConfig {
|
||||
target,
|
||||
control_sni: sni.clone(),
|
||||
max_ttl: max_hops,
|
||||
connect_timeout: Duration::from_secs(5),
|
||||
hop_timeout: Duration::from_secs(1),
|
||||
accept_rst_fin_as_timeout: dpi_config
|
||||
.as_ref()
|
||||
.is_some_and(|dpi| dpi.accept_rst_fin_as_timeout),
|
||||
accept_ack_as_timeout: dpi_config
|
||||
.as_ref()
|
||||
.is_some_and(|dpi| dpi.accept_ack_as_timeout),
|
||||
})
|
||||
.await?;
|
||||
Ok((
|
||||
ProbeCommandResult::SniTraceroute {
|
||||
host,
|
||||
sni,
|
||||
target,
|
||||
dpi_hop: dpi_hop_from_result(&trace),
|
||||
hops: trace.hops,
|
||||
},
|
||||
None,
|
||||
))
|
||||
}
|
||||
_ => bail!("unsupported DPI command"),
|
||||
}
|
||||
}
|
||||
|
||||
fn validate_dpi_probe_config(config: Option<&DpiProbeConfig>) -> Result<()> {
|
||||
let Some(config) = config else {
|
||||
return Ok(());
|
||||
@@ -471,21 +609,30 @@ async fn measure_dpi_hops(config: Option<&DpiProbeConfig>) -> DpiHops {
|
||||
max_ttl: config.max_ttl,
|
||||
connect_timeout: Duration::from_millis(config.connect_timeout_ms),
|
||||
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,
|
||||
};
|
||||
let (v4, v6) = tokio::join!(
|
||||
measure_dpi_hop(common(config.target_v4.into())),
|
||||
measure_dpi_hop(common(config.target_v6.into())),
|
||||
);
|
||||
DpiHops { v4, v6 }
|
||||
DpiHops {
|
||||
v4: v4.0,
|
||||
v6: v6.0,
|
||||
hops_v4: v4.1,
|
||||
hops_v6: v6.1,
|
||||
}
|
||||
}
|
||||
|
||||
async fn measure_dpi_hop(config: dpi_hop::DpiHopProbeConfig) -> Option<u8> {
|
||||
async fn measure_dpi_hop(
|
||||
config: dpi_hop::DpiHopProbeConfig,
|
||||
) -> (Option<u8>, Vec<reports::probe::DpiProbeHop>) {
|
||||
let target = config.target;
|
||||
match dpi_hop::detect_dpi_hop(config).await {
|
||||
Ok(result) => dpi_hop_from_result(&result),
|
||||
Ok(result) => (dpi_hop_from_result(&result), result.hops),
|
||||
Err(error) => {
|
||||
warn!("failed to measure DPI hop for {target}: {error}");
|
||||
None
|
||||
(None, Vec::new())
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -547,6 +694,8 @@ async fn publish_status(
|
||||
bundle_type: Some(args.bundle_type),
|
||||
dpi_hop_v4: dpi_hops.v4,
|
||||
dpi_hop_v6: dpi_hops.v6,
|
||||
dpi_hops_v4: dpi_hops.hops_v4,
|
||||
dpi_hops_v6: dpi_hops.hops_v6,
|
||||
})?;
|
||||
|
||||
client
|
||||
@@ -695,6 +844,31 @@ mod tests {
|
||||
assert_eq!(command_id("probe/commands/v1/42/trace-1/extra", "42"), None);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn remeasurement_requires_a_loaded_dpi_configuration() {
|
||||
let config = Arc::new(RwLock::new(None));
|
||||
let error = run_dpi_command(ProbeCommand::RemeasureDpiHop, &config)
|
||||
.await
|
||||
.err()
|
||||
.unwrap();
|
||||
assert!(error.to_string().contains("not loaded"));
|
||||
assert!(config.read().await.is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn dpi_commands_decode() {
|
||||
let command: ProbeCommand =
|
||||
serde_json::from_str(r#"{"type":"remeasure_dpi_hop"}"#).unwrap();
|
||||
assert!(matches!(command, ProbeCommand::RemeasureDpiHop));
|
||||
let command: ProbeCommand = serde_json::from_str(
|
||||
r#"{"type":"sni_traceroute","host":"example.com","sni":"blocked.example"}"#,
|
||||
)
|
||||
.unwrap();
|
||||
assert!(
|
||||
matches!(command, ProbeCommand::SniTraceroute { host, sni, max_hops: 30 } if host == "example.com" && sni == "blocked.example")
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn manual_commands_decode() {
|
||||
let command: ProbeCommand =
|
||||
@@ -730,6 +904,14 @@ mod tests {
|
||||
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);
|
||||
assert!(!dpi.accept_rst_fin_as_timeout);
|
||||
assert!(!dpi.accept_ack_as_timeout);
|
||||
let mut value = serde_json::to_value(dpi).unwrap();
|
||||
value["accept_rst_fin_as_timeout"] = true.into();
|
||||
value["accept_ack_as_timeout"] = true.into();
|
||||
let dpi: DpiProbeConfig = serde_json::from_value(value).unwrap();
|
||||
assert!(dpi.accept_rst_fin_as_timeout);
|
||||
assert!(dpi.accept_ack_as_timeout);
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -741,12 +923,42 @@ mod tests {
|
||||
bundle_type: Some("debian"),
|
||||
dpi_hop_v4: Some(4),
|
||||
dpi_hop_v6: Some(6),
|
||||
dpi_hops_v4: vec![reports::probe::DpiProbeHop {
|
||||
ttl: 1,
|
||||
router: Some("192.0.2.1".parse().unwrap()),
|
||||
outcome: reports::probe::DpiProbeHopOutcome::IcmpTimeExceeded,
|
||||
}],
|
||||
dpi_hops_v6: vec![
|
||||
reports::probe::DpiProbeHop {
|
||||
ttl: 1,
|
||||
router: Some("2001:db8::1".parse().unwrap()),
|
||||
outcome: reports::probe::DpiProbeHopOutcome::IcmpTimeExceeded,
|
||||
},
|
||||
reports::probe::DpiProbeHop {
|
||||
ttl: 2,
|
||||
router: None,
|
||||
outcome: reports::probe::DpiProbeHopOutcome::TcpClosed,
|
||||
},
|
||||
],
|
||||
};
|
||||
let value = serde_json::to_value(status).unwrap();
|
||||
|
||||
assert_eq!(value["dpi_hop_v4"], 4);
|
||||
assert_eq!(value["dpi_hop_v6"], 6);
|
||||
assert_eq!(value["bundle_type"], "debian");
|
||||
assert_eq!(
|
||||
value["dpi_hops_v4"][0],
|
||||
serde_json::json!({
|
||||
"ttl": 1, "src": "192.0.2.1", "outcome": "icmp_time_exceeded"
|
||||
})
|
||||
);
|
||||
assert_eq!(value["dpi_hops_v6"][0]["src"], "2001:db8::1");
|
||||
assert_eq!(
|
||||
value["dpi_hops_v6"][1],
|
||||
serde_json::json!({
|
||||
"ttl": 2, "src": null, "outcome": "tcp_closed"
|
||||
})
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
||||
@@ -12,6 +12,35 @@ pub struct ProbeStatus<'a> {
|
||||
pub dpi_hop_v4: Option<u8>,
|
||||
#[serde(default)]
|
||||
pub dpi_hop_v6: Option<u8>,
|
||||
#[serde(default)]
|
||||
pub dpi_hops_v4: Vec<DpiProbeHop>,
|
||||
#[serde(default)]
|
||||
pub dpi_hops_v6: Vec<DpiProbeHop>,
|
||||
}
|
||||
|
||||
/// One attempted DPI-probe TTL, including unsuccessful measurements.
|
||||
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
|
||||
pub struct DpiProbeHop {
|
||||
pub ttl: u8,
|
||||
/// Observed ICMP response source; absent when no packet source was captured.
|
||||
#[serde(rename = "src")]
|
||||
pub router: Option<IpAddr>,
|
||||
pub outcome: DpiProbeHopOutcome,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
|
||||
#[serde(rename_all = "snake_case")]
|
||||
pub enum DpiProbeHopOutcome {
|
||||
IcmpTimeExceeded,
|
||||
Timeout,
|
||||
TcpClosed,
|
||||
TcpAcknowledged,
|
||||
}
|
||||
|
||||
impl DpiProbeHopOutcome {
|
||||
pub const fn invalidates_measurement(self) -> bool {
|
||||
matches!(self, Self::TcpClosed | Self::TcpAcknowledged)
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Serialize, Deserialize)]
|
||||
@@ -40,6 +69,12 @@ pub struct DpiProbeConfig {
|
||||
pub max_ttl: u8,
|
||||
#[serde(default = "default_post_dpi_hop_limit")]
|
||||
pub post_dpi_hop_limit: u8,
|
||||
/// Treat TCP RST/FIN as a timeout and continue measuring later TTLs.
|
||||
#[serde(default)]
|
||||
pub accept_rst_fin_as_timeout: bool,
|
||||
/// Treat an acknowledged payload as a timeout and continue measuring later TTLs.
|
||||
#[serde(default)]
|
||||
pub accept_ack_as_timeout: bool,
|
||||
}
|
||||
|
||||
pub const fn default_post_dpi_hop_limit() -> u8 {
|
||||
@@ -161,6 +196,13 @@ pub enum ProbeCommand {
|
||||
max_hops: u8,
|
||||
},
|
||||
ResubscribeTasks,
|
||||
RemeasureDpiHop,
|
||||
SniTraceroute {
|
||||
host: String,
|
||||
sni: String,
|
||||
#[serde(default = "default_manual_max_hops")]
|
||||
max_hops: u8,
|
||||
},
|
||||
}
|
||||
|
||||
pub const fn default_manual_max_hops() -> u8 {
|
||||
@@ -177,6 +219,19 @@ pub enum ProbeCommandResult {
|
||||
ResubscribeTasks {
|
||||
requested: bool,
|
||||
},
|
||||
RemeasureDpiHop {
|
||||
dpi_hop_v4: Option<u8>,
|
||||
dpi_hop_v6: Option<u8>,
|
||||
dpi_hops_v4: Vec<DpiProbeHop>,
|
||||
dpi_hops_v6: Vec<DpiProbeHop>,
|
||||
},
|
||||
SniTraceroute {
|
||||
host: String,
|
||||
sni: String,
|
||||
target: std::net::SocketAddr,
|
||||
dpi_hop: Option<u8>,
|
||||
hops: Vec<DpiProbeHop>,
|
||||
},
|
||||
Error {
|
||||
message: String,
|
||||
},
|
||||
@@ -222,3 +277,62 @@ pub enum ProbeEvidence {
|
||||
DataTimeout { bytes: u32 },
|
||||
Good,
|
||||
}
|
||||
|
||||
/// Manual SNI traces accept a bare hostname or IP and a DNS SNI, on port 443.
|
||||
pub fn validate_sni_traceroute(host: &str, sni: &str) -> Result<(), &'static str> {
|
||||
fn dns_name(value: &str) -> bool {
|
||||
!value.is_empty()
|
||||
&& value.len() <= 253
|
||||
&& value
|
||||
.strip_suffix('.')
|
||||
.unwrap_or(value)
|
||||
.split('.')
|
||||
.all(|label| {
|
||||
!label.is_empty()
|
||||
&& label.len() <= 63
|
||||
&& !label.starts_with('-')
|
||||
&& !label.ends_with('-')
|
||||
&& label
|
||||
.bytes()
|
||||
.all(|b| b.is_ascii_alphanumeric() || b == b'-')
|
||||
})
|
||||
}
|
||||
if host.parse::<IpAddr>().is_err() && !dns_name(host) {
|
||||
return Err("host must be a bare hostname or IP address");
|
||||
}
|
||||
if sni.parse::<IpAddr>().is_ok() || !dns_name(sni) {
|
||||
return Err("sni must be a DNS name");
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod command_tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn validates_separate_host_and_sni() {
|
||||
for host in ["example.com", "192.0.2.1", "2001:db8::1"] {
|
||||
assert!(validate_sni_traceroute(host, "blocked.example").is_ok());
|
||||
}
|
||||
for host in [
|
||||
"",
|
||||
"https://example.com",
|
||||
"example.com:443",
|
||||
"bad host",
|
||||
"-invalid.example",
|
||||
] {
|
||||
assert!(validate_sni_traceroute(host, "blocked.example").is_err());
|
||||
}
|
||||
for sni in [
|
||||
"",
|
||||
"192.0.2.1",
|
||||
"bad/name",
|
||||
"bad..name",
|
||||
"bad_.name",
|
||||
"bad...",
|
||||
] {
|
||||
assert!(validate_sni_traceroute("example.com", sni).is_err());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+1
-1
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "website"
|
||||
version = "1.4.3"
|
||||
version = "1.4.4"
|
||||
edition = "2024"
|
||||
|
||||
[dependencies]
|
||||
|
||||
@@ -11,6 +11,9 @@ connect_timeout_ms = 5000
|
||||
hop_timeout_ms = 1000
|
||||
max_ttl = 15
|
||||
post_dpi_hop_limit = 3
|
||||
# Continue DPI probing by treating TCP RST/FIN or payload ACK as timeouts.
|
||||
accept_rst_fin_as_timeout = false
|
||||
accept_ack_as_timeout = false
|
||||
|
||||
[[hosts]]
|
||||
id = "hil-hetzner" # hil-speed.hetzner.com
|
||||
|
||||
+37
-1
@@ -59,6 +59,10 @@ pub struct ProbeRow {
|
||||
bundle_type: Option<String>,
|
||||
dpi_hop_v4: Option<i16>,
|
||||
dpi_hop_v6: Option<i16>,
|
||||
#[sqlx(skip)]
|
||||
dpi_hops_v4: Vec<reports::probe::DpiProbeHop>,
|
||||
#[sqlx(skip)]
|
||||
dpi_hops_v6: Vec<reports::probe::DpiProbeHop>,
|
||||
}
|
||||
|
||||
async fn rows(pool: &PgPool, mqtt: &MqttPublisher) -> Result<Vec<ProbeRow>, Status> {
|
||||
@@ -78,6 +82,8 @@ async fn rows(pool: &PgPool, mqtt: &MqttPublisher) -> Result<Vec<ProbeRow>, Stat
|
||||
bundle_type,
|
||||
dpi_hop_v4,
|
||||
dpi_hop_v6,
|
||||
dpi_hops_v4,
|
||||
dpi_hops_v6,
|
||||
}) = statuses.get(&row.id.to_string())
|
||||
{
|
||||
row.online = *online;
|
||||
@@ -85,6 +91,8 @@ async fn rows(pool: &PgPool, mqtt: &MqttPublisher) -> Result<Vec<ProbeRow>, Stat
|
||||
row.bundle_type = bundle_type.clone();
|
||||
row.dpi_hop_v4 = dpi_hop_v4.map(i16::from);
|
||||
row.dpi_hop_v6 = dpi_hop_v6.map(i16::from);
|
||||
row.dpi_hops_v4 = dpi_hops_v4.clone();
|
||||
row.dpi_hops_v6 = dpi_hops_v6.clone();
|
||||
}
|
||||
}
|
||||
Ok(rows)
|
||||
@@ -273,7 +281,16 @@ pub async fn update_one_probe(
|
||||
#[serde(tag = "type", rename_all = "snake_case")]
|
||||
pub enum CommandInput {
|
||||
ResubscribeTasks,
|
||||
Traceroute { target: IpAddr, max_hops: u8 },
|
||||
RemeasureDpiHop,
|
||||
SniTraceroute {
|
||||
host: String,
|
||||
sni: String,
|
||||
max_hops: u8,
|
||||
},
|
||||
Traceroute {
|
||||
target: IpAddr,
|
||||
max_hops: u8,
|
||||
},
|
||||
}
|
||||
|
||||
#[post("/probes/<id>/commands", format = "json", data = "<input>")]
|
||||
@@ -294,6 +311,25 @@ pub async fn command(
|
||||
}
|
||||
let command = match input.into_inner() {
|
||||
CommandInput::ResubscribeTasks => ProbeCommand::ResubscribeTasks,
|
||||
CommandInput::RemeasureDpiHop => ProbeCommand::RemeasureDpiHop,
|
||||
CommandInput::SniTraceroute {
|
||||
host,
|
||||
sni,
|
||||
max_hops,
|
||||
} => {
|
||||
let host = host.trim().to_owned();
|
||||
let sni = sni.trim().to_owned();
|
||||
if !(1..=64).contains(&max_hops)
|
||||
|| reports::probe::validate_sni_traceroute(&host, &sni).is_err()
|
||||
{
|
||||
return Err(Status::BadRequest);
|
||||
}
|
||||
ProbeCommand::SniTraceroute {
|
||||
host,
|
||||
sni,
|
||||
max_hops,
|
||||
}
|
||||
}
|
||||
CommandInput::Traceroute { target, max_hops } if (1..=64).contains(&max_hops) => {
|
||||
ProbeCommand::Traceroute { target, max_hops }
|
||||
}
|
||||
|
||||
@@ -59,6 +59,8 @@ struct NodeStatus {
|
||||
bundle_type: Option<String>,
|
||||
dpi_hop_v4: Option<u8>,
|
||||
dpi_hop_v6: Option<u8>,
|
||||
dpi_hops_v4: Vec<reports::probe::DpiProbeHop>,
|
||||
dpi_hops_v6: Vec<reports::probe::DpiProbeHop>,
|
||||
}
|
||||
|
||||
#[get("/nodes")]
|
||||
@@ -111,6 +113,12 @@ fn build_node_statuses(
|
||||
bundle_type: status.and_then(|status| status.bundle_type.clone()),
|
||||
dpi_hop_v4: status.and_then(|status| status.dpi_hop_v4),
|
||||
dpi_hop_v6: status.and_then(|status| status.dpi_hop_v6),
|
||||
dpi_hops_v4: status
|
||||
.map(|status| status.dpi_hops_v4.clone())
|
||||
.unwrap_or_default(),
|
||||
dpi_hops_v6: status
|
||||
.map(|status| status.dpi_hops_v6.clone())
|
||||
.unwrap_or_default(),
|
||||
}
|
||||
})
|
||||
.collect()
|
||||
@@ -131,6 +139,11 @@ mod tests {
|
||||
|
||||
#[test]
|
||||
fn builds_online_and_offline_json_nodes() {
|
||||
let hops = vec![reports::probe::DpiProbeHop {
|
||||
ttl: 1,
|
||||
router: Some("192.0.2.1".parse().unwrap()),
|
||||
outcome: reports::probe::DpiProbeHopOutcome::IcmpTimeExceeded,
|
||||
}];
|
||||
let probes = vec![
|
||||
ProbeMetadata {
|
||||
id: 1,
|
||||
@@ -157,6 +170,8 @@ mod tests {
|
||||
bundle_type: Some("debian".to_string()),
|
||||
dpi_hop_v4: Some(5),
|
||||
dpi_hop_v6: None,
|
||||
dpi_hops_v4: hops.clone(),
|
||||
dpi_hops_v6: Vec::new(),
|
||||
},
|
||||
)]);
|
||||
|
||||
@@ -177,6 +192,8 @@ mod tests {
|
||||
bundle_type: Some("debian".to_string()),
|
||||
dpi_hop_v4: Some(5),
|
||||
dpi_hop_v6: None,
|
||||
dpi_hops_v4: hops.clone(),
|
||||
dpi_hops_v6: Vec::new(),
|
||||
},
|
||||
NodeStatus {
|
||||
probe_id: 2,
|
||||
@@ -190,6 +207,8 @@ mod tests {
|
||||
bundle_type: None,
|
||||
dpi_hop_v4: None,
|
||||
dpi_hop_v6: None,
|
||||
dpi_hops_v4: Vec::new(),
|
||||
dpi_hops_v6: Vec::new(),
|
||||
},
|
||||
]
|
||||
);
|
||||
|
||||
+80
-1
@@ -81,6 +81,8 @@ pub struct ProbeStatusSnapshot {
|
||||
pub bundle_type: Option<String>,
|
||||
pub dpi_hop_v4: Option<u8>,
|
||||
pub dpi_hop_v6: Option<u8>,
|
||||
pub dpi_hops_v4: Vec<reports::probe::DpiProbeHop>,
|
||||
pub dpi_hops_v6: Vec<reports::probe::DpiProbeHop>,
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
@@ -280,6 +282,23 @@ impl MqttPublisher {
|
||||
command: ProbeCommand,
|
||||
) -> Result<ProbeCommandResult, PublishError> {
|
||||
let client = self.client.as_ref().ok_or(PublishError::NotConfigured)?;
|
||||
let timeout = match &command {
|
||||
ProbeCommand::SniTraceroute { max_hops, .. } => {
|
||||
Duration::from_secs(70 + u64::from(*max_hops) * 7)
|
||||
}
|
||||
ProbeCommand::RemeasureDpiHop => {
|
||||
let config = self.probe_config.read().await;
|
||||
Duration::from_millis(config.dpi_probe.as_ref().map_or(60_000, |dpi| {
|
||||
60_000
|
||||
+ u64::from(dpi.max_ttl)
|
||||
* (dpi
|
||||
.connect_timeout_ms
|
||||
.saturating_add(dpi.hop_timeout_ms)
|
||||
.saturating_add(100))
|
||||
}))
|
||||
}
|
||||
_ => Duration::from_secs(60),
|
||||
};
|
||||
let command_id = uuid::Uuid::new_v4().to_string();
|
||||
let (sender, receiver) = rocket::tokio::sync::oneshot::channel();
|
||||
self.command_sessions
|
||||
@@ -299,7 +318,7 @@ impl MqttPublisher {
|
||||
self.command_sessions.lock().await.remove(&command_id);
|
||||
return Err(PublishError::Publish(error));
|
||||
}
|
||||
let result = rocket::tokio::time::timeout(Duration::from_secs(60), receiver).await;
|
||||
let result = rocket::tokio::time::timeout(timeout, receiver).await;
|
||||
self.command_sessions.lock().await.remove(&command_id);
|
||||
result
|
||||
.ok()
|
||||
@@ -474,6 +493,8 @@ async fn dispatch_probe_status(probe_statuses: &ProbeStatuses, topic: &str, payl
|
||||
bundle_type: status.bundle_type.map(str::to_string),
|
||||
dpi_hop_v4: status.dpi_hop_v4,
|
||||
dpi_hop_v6: status.dpi_hop_v6,
|
||||
dpi_hops_v4: status.dpi_hops_v4,
|
||||
dpi_hops_v6: status.dpi_hops_v6,
|
||||
},
|
||||
);
|
||||
}
|
||||
@@ -579,6 +600,33 @@ fn task_timeout_ms_from_env() -> u64 {
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn global_dpi_timeout_options_survive_config_parsing() {
|
||||
let contents = r#"
|
||||
timeout_sec = 3
|
||||
min_data = 65536
|
||||
hosts = []
|
||||
[dpi_probe]
|
||||
sni = "example.com"
|
||||
target_v4 = "192.0.2.1:443"
|
||||
target_v6 = "[2001:db8::1]:443"
|
||||
connect_timeout_ms = 5000
|
||||
hop_timeout_ms = 1000
|
||||
max_ttl = 15
|
||||
"#;
|
||||
let config = parse_probe_hosts(contents).unwrap();
|
||||
let dpi = config.dpi_probe.unwrap();
|
||||
assert!(!dpi.accept_rst_fin_as_timeout);
|
||||
assert!(!dpi.accept_ack_as_timeout);
|
||||
|
||||
let enabled =
|
||||
format!("{contents}\naccept_rst_fin_as_timeout = true\naccept_ack_as_timeout = true\n");
|
||||
let config = parse_probe_hosts(&enabled).unwrap();
|
||||
let payload = serde_json::to_value(config.dpi_probe.unwrap()).unwrap();
|
||||
assert_eq!(payload["accept_rst_fin_as_timeout"], true);
|
||||
assert_eq!(payload["accept_ack_as_timeout"], true);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn update_check_uses_global_or_targeted_topic() {
|
||||
assert_eq!(probe_update_topic(None), "probe/update/v1");
|
||||
@@ -620,10 +668,39 @@ mod tests {
|
||||
bundle_type: Some("openwrt".to_string()),
|
||||
dpi_hop_v4: Some(4),
|
||||
dpi_hop_v6: Some(6),
|
||||
dpi_hops_v4: Vec::new(),
|
||||
dpi_hops_v6: Vec::new(),
|
||||
})
|
||||
);
|
||||
}
|
||||
|
||||
#[rocket::async_test]
|
||||
async fn status_snapshot_keeps_invalid_dpi_measurement_hops() {
|
||||
let statuses = Arc::new(rocket::tokio::sync::RwLock::new(HashMap::new()));
|
||||
dispatch_probe_status(
|
||||
&statuses,
|
||||
"probe/status/v1/42",
|
||||
br#"{"online":true,"probe_id":"42","version":"1.2.3","dpi_hop_v4":null,"dpi_hops_v4":[{"ttl":1,"src":"192.0.2.1","outcome":"icmp_time_exceeded"},{"ttl":2,"src":null,"outcome":"tcp_closed"}],"dpi_hops_v6":[{"ttl":1,"src":"2001:db8::1","outcome":"icmp_time_exceeded"}]}"#,
|
||||
).await;
|
||||
let statuses = statuses.read().await;
|
||||
let status = statuses.get("42").unwrap();
|
||||
assert_eq!(status.dpi_hop_v4, None);
|
||||
assert_eq!(status.dpi_hops_v4.len(), 2);
|
||||
assert_eq!(
|
||||
status.dpi_hops_v4[0].router,
|
||||
Some("192.0.2.1".parse().unwrap())
|
||||
);
|
||||
assert_eq!(
|
||||
status.dpi_hops_v4[1].outcome,
|
||||
reports::probe::DpiProbeHopOutcome::TcpClosed
|
||||
);
|
||||
assert_eq!(status.dpi_hops_v4[1].router, None);
|
||||
assert_eq!(
|
||||
status.dpi_hops_v6[0].router,
|
||||
Some("2001:db8::1".parse().unwrap())
|
||||
);
|
||||
}
|
||||
|
||||
#[rocket::async_test]
|
||||
async fn empty_retained_status_removes_snapshot() {
|
||||
let statuses = Arc::new(rocket::tokio::sync::RwLock::new(HashMap::from([(
|
||||
@@ -634,6 +711,8 @@ mod tests {
|
||||
bundle_type: None,
|
||||
dpi_hop_v4: None,
|
||||
dpi_hop_v6: None,
|
||||
dpi_hops_v4: Vec::new(),
|
||||
dpi_hops_v6: Vec::new(),
|
||||
},
|
||||
)])));
|
||||
|
||||
|
||||
Reference in New Issue
Block a user