feat: safe dpi traceroutes (#84)

* feat: safe dpi traceroutes

* chore: bump versions

* feat: dynamic verdict

* fix: change titles

* refactor: simplify methods

* fix: block message

* fix: narrow cdn block

* fix: title
This commit is contained in:
LowderPlay
2026-08-22 02:03:09 +05:00
committed by GitHub
parent 4d05f18c4c
commit 8651d8b1ff
21 changed files with 1093 additions and 352 deletions
Generated
+2 -2
View File
@@ -2632,7 +2632,7 @@ dependencies = [
[[package]]
name = "probe"
version = "0.2.0"
version = "0.3.0"
dependencies = [
"anyhow",
"clap",
@@ -4706,7 +4706,7 @@ dependencies = [
[[package]]
name = "website"
version = "1.2.1"
version = "1.2.2"
dependencies = [
"dotenvy",
"env_logger",
+60 -8
View File
@@ -26,22 +26,27 @@ export type DnsObservation = {
| { type: "Error"; message: string };
};
export type ProbeVerdict =
| "uncertain"
| "dns_spoofing"
| "sni_block"
| "tspu_block"
| "whitelist"
| "ok";
export type DisplayProbeVerdict = ProbeVerdict | "cdn_block";
export type ResolvedProbeVerdict = Exclude<DisplayProbeVerdict, "uncertain">;
export type ProbeResult = {
job_id: string;
probe_id: string;
region?: string | null;
provider?: string | null;
asn?: string | null;
verdicts: (
| "uncertain"
| "dns_spoofing"
| "sni_block"
| "tspu_block"
| "whitelist"
| "ok"
)[];
verdicts: ProbeVerdict[];
host_results: ProbeHostResult[] | null;
target_hop: number | null;
dpi_hop: number | null;
dns: {
spoofing_detected: boolean;
suspicious_provider_count: number;
@@ -51,6 +56,53 @@ export type ProbeResult = {
} | null;
};
export function displayProbeVerdicts(
probe: ProbeResult,
isStaticCdn: boolean,
): DisplayProbeVerdict[] {
return probe.verdicts.map((verdict) =>
isStaticCdn && probe.host_results?.length !== 0 && verdict === "ok"
? "cdn_block"
: verdict,
);
}
const verdictPriority: DisplayProbeVerdict[] = [
"tspu_block",
"sni_block",
"dns_spoofing",
"whitelist",
"cdn_block",
"ok",
"uncertain",
];
export function selectProbeVerdict(
probes: ProbeResult[],
isStaticCdn: boolean,
): ResolvedProbeVerdict | null {
if (probes.length === 0) return null;
const votes = new Map<DisplayProbeVerdict, number>();
for (const probe of probes) {
for (const verdict of new Set(displayProbeVerdicts(probe, isStaticCdn))) {
votes.set(verdict, (votes.get(verdict) ?? 0) + 1);
}
}
let winner: DisplayProbeVerdict | null = null;
let winningVotes = 0;
for (const verdict of verdictPriority) {
const count = votes.get(verdict) ?? 0;
if (count > winningVotes) {
winner = verdict;
winningVotes = count;
}
}
return winner === "uncertain" ? null : winner;
}
export type ProbeStatus = {
id: string;
target: string;
+16 -4
View File
@@ -13,6 +13,7 @@ import {
ShieldCheck,
} from "@lucide/svelte";
import type { CheckResult } from "$lib/api/check";
import type { ResolvedProbeVerdict } from "$lib/api/probe";
import DetailRow from "./DetailRow.svelte";
import ResultIpList from "./result/ResultIpList.svelte";
import ResultStatusHeader from "./result/ResultStatusHeader.svelte";
@@ -21,23 +22,34 @@ import ResultTargetCard from "./result/ResultTargetCard.svelte";
type Provider = { name: string; networks: { cidr: string }[] };
type ResultTheme = "blocked" | "clean" | "whitelist";
type ResultVerdict = ResolvedProbeVerdict | "blocked";
let {
result,
probeVerdict = null,
}: {
result: CheckResult;
probeVerdict?: ResolvedProbeVerdict | null;
} = $props();
const valueClass = "text-right text-sm font-medium text-neutral-200";
const alertValueClass = `${valueClass} text-red-500`;
const successValueClass = `${valueClass} text-green-500`;
const theme = $derived<ResultTheme>(
const staticVerdict = $derived<ResultVerdict>(
result.whitelist && !result.domain
? "whitelist"
: result.found
? "blocked"
: "clean",
: "ok",
);
const verdict = $derived<ResultVerdict>(probeVerdict ?? staticVerdict);
const theme = $derived<ResultTheme>(
verdict === "whitelist"
? "whitelist"
: verdict === "ok"
? "clean"
: "blocked",
);
const panelClass = $derived(
theme === "whitelist"
@@ -67,7 +79,7 @@ const providerCidrs = (provider: Provider) =>
<div class="mt-8 space-y-6">
<div class="grid grid-cols-1 md:grid-cols-2 gap-4">
<div class={`border p-4 rounded-lg flex items-center ${panelClass}`}>
<ResultStatusHeader {theme} blocked={result.blocked} />
<ResultStatusHeader {verdict} />
</div>
<ResultTargetCard
targetType={result.targetType}
@@ -171,7 +183,7 @@ const providerCidrs = (provider: Provider) =>
{#if result.whitelist}
<DetailRow
label="Белый список (?)"
label="Исключение для CDN блокировки (?)"
href="/kb/whitelist"
icon={ShieldCheck}
>
@@ -8,17 +8,23 @@ import {
CircleX,
LoaderCircle,
ShieldCheck,
TriangleAlert,
} from "@lucide/svelte";
import type { DnsObservation, ProbeResult, ProbeStatus } from "$lib/api/probe";
import {
type DnsObservation,
displayProbeVerdicts,
type ProbeResult,
type ProbeStatus,
} from "$lib/api/probe";
let {
probes,
status,
isStaticBlocked,
isStaticCdn,
}: {
probes: ProbeResult[];
status: ProbeStatus;
isStaticBlocked: boolean;
isStaticCdn: boolean;
} = $props();
let expandedRows = $state<Record<string, boolean>>({});
@@ -110,19 +116,6 @@ const verdictStyles = {
border: "border-neutral-400/20",
},
};
type VerdictStyle = keyof typeof verdictStyles;
function probeVerdicts(
probe: ProbeResult,
isStaticBlocked: boolean,
): VerdictStyle[] {
return probe.verdicts.map((verdict) =>
isStaticBlocked && probe.host_results?.length !== 0 && verdict === "ok"
? "cdn_block"
: verdict,
);
}
</script>
<div class="mt-8 space-y-4">
@@ -182,7 +175,7 @@ function probeVerdicts(
</thead>
<tbody>
{#each probes as probe (probe.probe_id)}
{@const verdicts = probeVerdicts(probe, isStaticBlocked)}
{@const verdicts = displayProbeVerdicts(probe, isStaticCdn)}
{@const isExpanded = !!expandedRows[probe.probe_id]}
<tr
class="border-b border-neutral-800/50 hover:bg-neutral-800/20 transition-colors cursor-pointer select-none"
@@ -203,23 +196,14 @@ function probeVerdicts(
<div class="flex flex-wrap gap-1.5">
{#each verdicts as verdict}
{@const style = verdictStyles[verdict]}
{#if verdict === "whitelist"}
<a
href="/kb/whitelist"
onclick={(event) => event.stopPropagation()}
class={`inline-flex items-center gap-1.5 px-2 py-1 rounded border ${style.bg} ${style.border} ${style.class} text-xs font-bold transition-colors hover:bg-amber-500/20 hover:border-amber-500/50 hover:text-amber-400`}
>
<style.icon size={14} />
{style.text}
</a>
{:else}
<div
class={`inline-flex items-center gap-1.5 px-2 py-1 rounded border ${style.bg} ${style.border} ${style.class} text-xs font-bold`}
>
<style.icon size={14} />
{style.text}
</div>
{/if}
<a
href={verdict === "whitelist" ? "/kb/whitelist" : "/kb/probing"}
onclick={(event) => event.stopPropagation()}
class={`inline-flex items-center gap-1.5 px-2 py-1 rounded border ${style.bg} ${style.border} ${style.class} text-xs font-bold transition-opacity hover:opacity-80`}
>
<style.icon size={14} />
{style.text}
</a>
{/each}
</div>
</td>
@@ -234,27 +218,27 @@ function probeVerdicts(
{#if isExpanded}
<tr class="bg-neutral-900/30">
<td colspan="4" class="p-4 border-b border-neutral-800/50">
<!-- {#if probe.verdicts.includes("tspu_block")}
{#if probe.verdicts.includes("tspu_block") && probe.dpi_hop !== null}
<div
class="mb-3 flex items-center gap-2 rounded-md border border-red-500/50 bg-red-500/15 px-3 py-2 text-red-200"
>
<TriangleAlert size={18} class="shrink-0 text-red-400" />
<span class="font-bold">
Блокировка ТСПУ обнаружена после
{probe.target_hop}
Блокировка на ТСПУ обнаружена после
{probe.dpi_hop}
прыжка
</span>
</div>
{:else}
<div class="text-md text-neutral-200 mb-2">
Блокировка на ТСПУ <b>не обнаружена</b> после
{probe.target_hop}
прыжков
{:else if probe.target_hop !== null}
<div class="text-md text-neutral-200 mb-3">
Блокировка IP на ТСПУ <b>не обнаружена</b>
</div>
{/if} -->
{/if}
{#if probe.dns}
<div class="mb-4">
<div class="mb-2 flex items-center justify-between gap-3">
<div
class="border-t border-neutral-800 py-2 flex items-center justify-between gap-3"
>
<h4
class="text-xs font-bold uppercase tracking-wide text-neutral-300"
>
@@ -331,53 +315,57 @@ function probeVerdicts(
</p>
</div>
{/if}
<div class="mb-2 border-t border-neutral-800 pt-4">
<h4
class="text-xs font-bold uppercase tracking-wide text-neutral-300"
{#if probe.host_results?.length}
<div
class="border-t border-neutral-800 py-2 flex items-center justify-between gap-3"
>
CDN-проверка
</h4>
</div>
<div class="grid grid-cols-1 md:grid-cols-2 gap-4">
{#each probe.host_results as host}
<div
class="flex items-center justify-between p-2 rounded bg-neutral-800/30 border border-neutral-700/30"
<h4
class="text-xs font-bold uppercase tracking-wide text-neutral-300"
>
<div class="flex flex-col">
<span class="text-xs font-bold text-neutral-400">
Сервер {host.host_id}
({host.host === "Blacklist" ? "в заблокированных" : "в доступных"}
диапазонах)
</span>
<span class="text-xs text-neutral-200">
{#if host.probe_evidence.type === 'Good'}
Успешно
{:else if host.probe_evidence.type === 'ClientHello'}
Блокировка после ClientHello
{:else if host.probe_evidence.type === 'DataTimeout'}
Таймаут получения данных, получено
{host.probe_evidence.bytes}
байт
{:else if host.probe_evidence.type === 'ConnectionError'}
Ошибка подключения
{/if}
</span>
CDN-проверка
</h4>
</div>
<div class="grid grid-cols-1 md:grid-cols-2 gap-4">
{#each probe.host_results as host}
<div
class="flex items-center justify-between p-2 rounded bg-neutral-800/30 border border-neutral-700/30"
>
<div class="flex flex-col">
<span class="text-xs font-bold text-neutral-400">
Сервер {host.host_id}
({host.host === "Blacklist" ? "в заблокированных" : "в доступных"}
диапазонах)
</span>
<span class="text-xs text-neutral-200">
{#if host.probe_evidence.type === 'Good'}
Успешно
{:else if host.probe_evidence.type === 'ClientHello'}
Блокировка после ClientHello
{:else if host.probe_evidence.type === 'DataTimeout'}
Таймаут получения данных, получено
{host.probe_evidence.bytes}
байт
{:else if host.probe_evidence.type === 'ConnectionError'}
Ошибка подключения
{/if}
</span>
</div>
{#if host.probe_evidence.type === 'Good'}
<CircleCheck size={14} class="text-green-500" />
{:else if host.probe_evidence.type === 'ClientHello'}
<CircleX size={14} class="text-red-500" />
{:else if host.probe_evidence.type === 'DataTimeout'}
<CircleX size={14} class="text-orange-500" />
{:else}
<CircleQuestionMark
size={14}
class="text-neutral-500"
/>
{/if}
</div>
{#if host.probe_evidence.type === 'Good'}
<CircleCheck size={14} class="text-green-500" />
{:else if host.probe_evidence.type === 'ClientHello'}
<CircleX size={14} class="text-red-500" />
{:else if host.probe_evidence.type === 'DataTimeout'}
<CircleX size={14} class="text-orange-500" />
{:else}
<CircleQuestionMark
size={14}
class="text-neutral-500"
/>
{/if}
</div>
{/each}
</div>
{/each}
</div>
{/if}
</td>
</tr>
{/if}
@@ -1,39 +1,61 @@
<script lang="ts">
import { ShieldAlert, ShieldCheck, ShieldX } from "@lucide/svelte";
import type { ResolvedProbeVerdict } from "$lib/api/probe";
type ResultTheme = "blocked" | "clean" | "whitelist";
type ResultVerdict = ResolvedProbeVerdict | "blocked";
let {
theme,
blocked,
verdict,
}: {
theme: ResultTheme;
blocked: boolean;
verdict: ResultVerdict;
} = $props();
const title = $derived(
theme === "whitelist"
? "Белый список"
: theme === "blocked"
? "Заблокирован"
: "Доступен",
);
const subtitle = $derived(
theme === "whitelist"
? "Ресурс находится в белом списке"
: theme === "blocked"
? "Ресурс был найден в списках блокировок"
: "Ограничений не обнаружено",
);
const labels: Record<ResultVerdict, { title: string; subtitle: string }> = {
blocked: {
title: "Заблокирован",
subtitle: "Ресурс был найден в списках блокировок",
},
tspu_block: {
title: "Заблокирован",
subtitle: "Сканеры обнаружили блокировку TCP на уровне ТСПУ",
},
sni_block: {
title: "Заблокирован",
subtitle: "Сканеры обнаружили блокировку по имени домена в SNI",
},
dns_spoofing: {
title: "Заблокирован",
subtitle: "Сканеры обнаружили подмену ответов DNS для данного домена",
},
whitelist: {
title: "Исключение для CDN",
subtitle:
"Домен снимает ограничение 16-20 КБ при подключении к заблокированным CDN",
},
cdn_block: {
title: "Заблокирован",
subtitle: "Сканеры обнаружили блокировку CDN (16-20 КБ)",
},
ok: {
title: "Не ограничен",
subtitle: "Ограничений не обнаружено",
},
};
const title = $derived(labels[verdict].title);
const subtitle = $derived(labels[verdict].subtitle);
const accentClass = $derived(
theme === "whitelist"
verdict === "whitelist"
? "text-[#f0b100]"
: theme === "blocked"
? "text-red-500"
: "text-green-500",
: verdict === "ok"
? "text-green-500"
: "text-red-500",
);
const StatusIcon = $derived(
theme === "whitelist" ? ShieldAlert : blocked ? ShieldX : ShieldCheck,
verdict === "whitelist"
? ShieldAlert
: verdict === "ok"
? ShieldCheck
: ShieldX,
);
</script>
+11 -2
View File
@@ -6,6 +6,7 @@ import { CheckRequestError, fetchCheck } from "$lib/api/check";
import {
type ProbeResult,
type ProbeStatus,
selectProbeVerdict,
startProbeSSE,
} from "$lib/api/probe";
import EmptyResult from "$lib/components/EmptyResult.svelte";
@@ -104,6 +105,14 @@ $effect(() => {
const error = $derived(
(checkQuery.error instanceof CheckRequestError && checkQuery.error) || null,
);
const liveVerdict = $derived(
checkQuery.data && probeQuery.data
? selectProbeVerdict(
probeQuery.data.probes,
checkQuery.data.providers.length > 0,
)
: null,
);
</script>
<SearchForm />
@@ -121,13 +130,13 @@ const error = $derived(
<ErrorMessage status={error.status} reason={error.message} />
</div>
{:else if checkQuery.data}
<ResultPanel result={checkQuery.data} />
<ResultPanel result={checkQuery.data} probeVerdict={liveVerdict} />
{#if shouldProbe && probeQuery.data && probeQuery.data.status.online_probes > 0}
<ProbeTable
probes={probeQuery.data.probes}
status={probeQuery.data.status}
isStaticBlocked={checkQuery.data.blocked}
isStaticCdn={checkQuery.data.providers.length > 0}
/>
{/if}
+33 -22
View File
@@ -61,25 +61,36 @@ import KbNote from "$lib/components/kb/KbNote.svelte";
завершить TLS-обмен и получить минимальный объём данных. Для проверки
IP-адреса без домена этот этап пропускается.
</p>
<!-- <li>
Сканер открывает TCP-соединения к разрешённому IP-адресу на порту 443,
последовательно увеличивая TTL от 1 до 5. На каждом TTL по умолчанию
одновременно отправляются три попытки: это уменьшает влияние потерь
пакетов и ограничения частоты ICMP-ответов. Ответ ICMP Time Exceeded
переводит проверку к следующему TTL, а TCP RST или успешное соединение
завершают трассировку.
<ul class="list-disc">
<li>
При получении конфигурации сканер отдельно для IPv4 и IPv6 измеряет
положение DPI в сети. Он подключается к контрольному серверу на порту 443,
отправляет TLS ClientHello с контрольным SNI, а затем посылает TCP-пакеты
с последовательно увеличивающимся TTL. Ответы ICMP Time Exceeded
показывают, какие прыжки (hops) пакет успел пройти. Результат сохраняется
отдельно для каждой версии IP.
</li>
<li>
Из списка известных доступных хостов случайно выбираются не более трёх
IP-адресов той же версии протокола, что и цель. Все контрольные
трассировки идут параллельно, а для сравнения используется наименьший
измеренный контрольный TTL.
</li> -->
<!-- <p>
Если трассировки цели и контроля завершились ответом ICMP Time Exceeded, но
ответ для цели пришёл на меньшем TTL, сканер выставляет вердикт «ТСПУ Блок».
В результатах показывается номер последнего достигнутого перехода для цели.
</p> -->
При проверке целевого IP-адреса TCP-трассировка начинается со следующего
после DPI прыжка — <code>dpi_hop + 1</code>. По умолчанию проверяются три
прыжка, на каждом из которых одновременно выполняются три попытки. Первый
ответ ICMP Time Exceeded, TCP RST или успешное соединение немедленно
завершает трассировку.
</li>
</ul>
<p>
Если DPI был измерен, но ни одна попытка после него не получила ответа,
сканер выставляет вердикт «ТСПУ Блок» и показывает измеренный прыжок DPI.
Любой ответ после DPI означает, что пакеты прошли дальше, поэтому блокировка
IP-адреса на ТСПУ не подтверждается. Если прыжок, на котором установлен DPI,
для нужной версии IP не удалось определить, трассировка пропускается и этот
вердикт не выставляется.
</p>
<KbNote variant="info">
Измерение выполняется независимо для IPv4 и IPv6. Разрыв TCP-соединения во
время калибровки делает результат этой версии IP недействительным, чтобы не
принимать поведение контрольного сервера за работу DPI.
</KbNote>
<KbHeading id="проверка-dns" title="Проверка DNS"> Проверка DNS </KbHeading>
<p>
@@ -137,12 +148,12 @@ import KbNote from "$lib/components/kb/KbNote.svelte";
получает
<b>SNI Блок</b>.
</li>
<!-- <li>
<li>
<b>ТСПУ Блок</b>
– TCP-трассировка до цели остановилась раньше контрольной трассировки.
Обычно означает, что соединение разрывается на оборудовании оператора
(ТСПУ) до того как соединение дошло до цели.
</li> -->
– TCP-трассировка, начатая сразу после прыжка с DPI, не получила ни
ICMP-ответа, ни TCP RST и не смогла подключиться к цели. Это указывает на
возможную блокировку IP-адреса на оборудовании оператора.
</li>
<li>
<!-- biome-ignore format: link punctuation -->
<a
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "probe"
version = "0.2.0"
version = "0.3.0"
edition = "2024"
license-file = "../LICENSE"
description = "Dynamic network probe daemon for Cheburcheck"
+2 -4
View File
@@ -166,12 +166,10 @@ docker run --rm \
| `--probe-id`, `PROBE_ID` | ID сканера. | обязательно |
| `--probe-token`, `PROBE_TOKEN` | Секретный токен сканера. | обязательно |
| `--max-concurrent-tasks`, `MAX_CONCURRENT_TASKS` | Максимальное количество одновременных заданий. | `8` |
| `--traceroute-max-hops`, `TRACEROUTE_MAX_HOPS` | Максимальный TTL для TCP traceroute. | `5` |
| `--traceroute-retries`, `TRACEROUTE_RETRIES` | Количество одновременных TCP-попыток на каждом TTL. | `3` |
| `--traceroute-control-hosts`, `TRACEROUTE_CONTROL_HOSTS` | Максимальное количество случайных контрольных IP для одновременной трассировки. | `3` |
| `RUST_LOG` | Уровень логирования. | `info` |
`MAX_CONCURRENT_TASKS`, `TRACEROUTE_MAX_HOPS`, `TRACEROUTE_RETRIES` и `TRACEROUTE_CONTROL_HOSTS` должны быть больше нуля. Для получения ICMP-ответов traceroute процессу требуется capability `CAP_NET_RAW`; systemd unit и Docker-образ настраивают её автоматически.
`MAX_CONCURRENT_TASKS` и `TRACEROUTE_RETRIES` должны быть больше нуля. Для получения ICMP-ответов traceroute процессу требуется capability `CAP_NET_RAW`; systemd unit и Docker-образ настраивают её автоматически.
## Как работает проверка
@@ -180,7 +178,7 @@ docker run --rm \
1. публикует retained-статус `online` в MQTT;
2. подписывается на конфигурацию динамического сканирования;
3. получает задания на проверку доменов и IP-адресов;
4. параллельно запускает SNI-проверки (для доменов), TCP traceroute до цели и контрольный TCP traceroute;
4. параллельно запускает SNI-проверки (для доменов) и TCP traceroute до цели, начиная со следующего после DPI узла;
5. отправляет результат обратно в Cheburcheck.
Для каждого тестового хоста сканер открывает TCP-соединение, начинает TLS-handshake с проверяемым доменом в SNI, затем отправляет простой HTTP GET-запрос.
-3
View File
@@ -6,8 +6,5 @@ config cheburprobe 'main'
option mqtt_port '443'
option connection_timeout '30'
option max_concurrent_tasks '8'
option traceroute_max_hops '5'
option traceroute_retries '3'
option traceroute_control_hosts '3'
option log_level 'info'
-4
View File
@@ -15,9 +15,7 @@ start_service() {
config_get mqtt_port main mqtt_port '443'
config_get connection_timeout main connection_timeout '30'
config_get max_concurrent_tasks main max_concurrent_tasks '8'
config_get traceroute_max_hops main traceroute_max_hops '5'
config_get traceroute_retries main traceroute_retries '3'
config_get traceroute_control_hosts main traceroute_control_hosts '3'
config_get log_level main log_level 'info'
procd_open_instance
@@ -30,9 +28,7 @@ start_service() {
MQTT_PORT="$mqtt_port" \
MQTT_CONNECTION_TIMEOUT_SECS="$connection_timeout" \
MAX_CONCURRENT_TASKS="$max_concurrent_tasks" \
TRACEROUTE_MAX_HOPS="$traceroute_max_hops" \
TRACEROUTE_RETRIES="$traceroute_retries" \
TRACEROUTE_CONTROL_HOSTS="$traceroute_control_hosts" \
RUST_LOG="$log_level"
procd_set_param stdout 1
procd_set_param stderr 1
-8
View File
@@ -44,18 +44,10 @@ return view.extend({
o.datatype = 'and(uinteger,min(1))';
o.default = '8';
o = s.option(form.Value, 'traceroute_max_hops', _('Traceroute maximum hops'));
o.datatype = 'and(uinteger,min(1))';
o.default = '5';
o = s.option(form.Value, 'traceroute_retries', _('Traceroute retries'));
o.datatype = 'and(uinteger,min(1))';
o.default = '3';
o = s.option(form.Value, 'traceroute_control_hosts', _('Traceroute control hosts'));
o.datatype = 'and(uinteger,min(1))';
o.default = '3';
o = s.option(form.ListValue, 'log_level', _('Log level'));
o.value('error', _('Error'));
o.value('warn', _('Warning'));
+488
View File
@@ -0,0 +1,488 @@
use etherparse::{
Icmpv4Type, Icmpv6Slice, Icmpv6Type, IpNumber, LaxNetSlice, LaxSlicedPacket, TransportSlice,
icmpv4, icmpv6,
};
use rand::RngCore;
use rustls::pki_types::ServerName;
use rustls::{ClientConfig, ClientConnection, RootCertStore};
use socket2::{Domain, Protocol, SockRef, Socket, Type};
use std::io::{self, Write};
use std::mem::MaybeUninit;
use std::net::{IpAddr, Ipv4Addr, SocketAddr, SocketAddrV4, SocketAddrV6, TcpStream};
use std::sync::Arc;
use std::time::{Duration, Instant};
const PROBE_BYTES: usize = 256;
#[derive(Debug, Clone)]
pub struct DpiHopProbeConfig {
pub target: SocketAddr,
pub control_sni: String,
pub max_ttl: u8,
pub connect_timeout: Duration,
pub hop_timeout: Duration,
}
#[derive(Debug, Clone)]
pub struct DpiHopProbeResult {
pub target: SocketAddr,
pub local_addr: SocketAddr,
pub client_hello_bytes: usize,
pub max_icmp_time_exceeded_ttl: Option<u8>,
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,
}
pub async fn detect_dpi_hop(config: DpiHopProbeConfig) -> io::Result<DpiHopProbeResult> {
tokio::task::spawn_blocking(move || detect_dpi_hop_blocking(config))
.await
.map_err(io::Error::other)?
}
pub fn detect_dpi_hop_blocking(config: DpiHopProbeConfig) -> io::Result<DpiHopProbeResult> {
if config.max_ttl == 0 {
return Err(io::Error::new(
io::ErrorKind::InvalidInput,
"max_ttl must be greater than zero",
));
}
let client_hello = make_client_hello(&config.control_sni)?;
let (domain, protocol) = match config.target {
SocketAddr::V4(_) => (Domain::IPV4, Protocol::ICMPV4),
SocketAddr::V6(_) => (Domain::IPV6, Protocol::ICMPV6),
};
let icmp = Socket::new(domain, Type::RAW, Some(protocol))?;
icmp.set_read_timeout(Some(config.hop_timeout))?;
let mut tcp = TcpStream::connect_timeout(&config.target, config.connect_timeout)?;
tcp.set_nodelay(true)?;
let local_addr = tcp.local_addr()?;
if !same_ip_family(local_addr, config.target) {
return Err(io::Error::other(
"DPI hop probe local and target address families differ",
));
}
tcp.write_all(&client_hello)?;
let mut hops = Vec::with_capacity(config.max_ttl as usize);
let mut max_icmp_time_exceeded_ttl = None;
for ttl in 1..=config.max_ttl {
let mut payload = [0u8; PROBE_BYTES];
rand::thread_rng().fill_bytes(&mut payload);
drain_socket(&icmp)?;
if !send_with_ttl(&mut tcp, config.target, &payload, ttl)? {
hops.push(DpiHopProbeHop {
ttl,
router: None,
outcome: DpiHopProbeHopOutcome::TcpClosed,
});
break;
}
let router =
listen_for_time_exceeded(&icmp, local_addr, config.target, config.hop_timeout)?;
let outcome = if router.is_some() {
max_icmp_time_exceeded_ttl = Some(ttl);
DpiHopProbeHopOutcome::IcmpTimeExceeded
} else if peer_closed(&tcp)? {
DpiHopProbeHopOutcome::TcpClosed
} else {
DpiHopProbeHopOutcome::Timeout
};
hops.push(DpiHopProbeHop {
ttl,
router,
outcome,
});
if outcome == DpiHopProbeHopOutcome::TcpClosed {
break;
}
}
Ok(DpiHopProbeResult {
target: config.target,
local_addr,
client_hello_bytes: client_hello.len(),
max_icmp_time_exceeded_ttl,
hops,
})
}
fn make_client_hello(sni: &str) -> io::Result<Vec<u8>> {
let server_name = ServerName::try_from(sni.to_owned()).map_err(|_| {
io::Error::new(
io::ErrorKind::InvalidInput,
"control_sni is not a valid DNS name",
)
})?;
let config = ClientConfig::builder()
.with_root_certificates(RootCertStore::empty())
.with_no_client_auth();
let mut connection = ClientConnection::new(Arc::new(config), server_name)
.map_err(|error| io::Error::other(format!("create TLS client connection: {error}")))?;
let mut client_hello = Vec::new();
connection
.write_tls(&mut client_hello)
.map_err(|error| io::Error::other(format!("write TLS ClientHello: {error}")))?;
if client_hello.is_empty() {
return Err(io::Error::other("rustls did not produce a ClientHello"));
}
Ok(client_hello)
}
fn same_ip_family(left: SocketAddr, right: SocketAddr) -> bool {
matches!(
(left, right),
(SocketAddr::V4(_), SocketAddr::V4(_)) | (SocketAddr::V6(_), SocketAddr::V6(_))
)
}
fn send_with_ttl(
tcp: &mut TcpStream,
target: SocketAddr,
payload: &[u8],
ttl: u8,
) -> io::Result<bool> {
let previous_ttl = {
let socket = SockRef::from(&*tcp);
match target {
SocketAddr::V4(_) => socket.ttl_v4()?,
SocketAddr::V6(_) => socket.unicast_hops_v6()?,
}
};
set_ttl(tcp, target, ttl as u32)?;
let write_result = tcp.write_all(payload);
let restore_result = set_ttl(tcp, target, previous_ttl);
match write_result {
Ok(()) => {
restore_result?;
Ok(true)
}
Err(error)
if matches!(
error.kind(),
io::ErrorKind::ConnectionReset | io::ErrorKind::BrokenPipe
) =>
{
let _ = restore_result;
Ok(false)
}
Err(error) => {
let _ = restore_result;
Err(error)
}
}
}
fn set_ttl(tcp: &TcpStream, target: SocketAddr, ttl: u32) -> io::Result<()> {
let socket = SockRef::from(tcp);
match target {
SocketAddr::V4(_) => socket.set_ttl_v4(ttl),
SocketAddr::V6(_) => socket.set_unicast_hops_v6(ttl),
}
}
fn listen_for_time_exceeded(
icmp: &Socket,
local_addr: SocketAddr,
target: SocketAddr,
timeout: Duration,
) -> io::Result<Option<IpAddr>> {
let deadline = Instant::now() + timeout;
let mut buffer = [0u8; 65535];
loop {
let remaining = deadline.saturating_duration_since(Instant::now());
if remaining.is_zero() {
return Ok(None);
}
icmp.set_read_timeout(Some(remaining))?;
match recv_socket(icmp, &mut buffer) {
Ok((bytes, source)) => {
if let Some(router) =
match_time_exceeded(&buffer[..bytes], source, local_addr, target)
{
return Ok(Some(router));
}
}
Err(error)
if error.kind() == io::ErrorKind::WouldBlock
|| error.kind() == io::ErrorKind::TimedOut =>
{
return Ok(None);
}
Err(error) => return Err(error),
}
}
}
fn peer_closed(tcp: &TcpStream) -> io::Result<bool> {
tcp.set_nonblocking(true)?;
let mut byte = [0u8; 1];
let result = match tcp.peek(&mut byte) {
Ok(0) => Ok(true),
Ok(_) => Ok(false),
Err(error) if error.kind() == io::ErrorKind::WouldBlock => Ok(false),
Err(error)
if matches!(
error.kind(),
io::ErrorKind::ConnectionReset | io::ErrorKind::BrokenPipe
) =>
{
Ok(true)
}
Err(error) => Err(error),
};
let _ = tcp.set_nonblocking(false);
result
}
fn drain_socket(socket: &Socket) -> io::Result<usize> {
let previous_timeout = socket.read_timeout()?;
socket.set_nonblocking(true)?;
let mut drained = 0;
let mut buffer = [0u8; 65535];
loop {
match recv_socket(socket, &mut buffer) {
Ok(_) => drained += 1,
Err(error) if error.kind() == io::ErrorKind::WouldBlock => break,
Err(error) => {
let _ = socket.set_nonblocking(false);
let _ = socket.set_read_timeout(previous_timeout);
return Err(error);
}
}
}
socket.set_nonblocking(false)?;
socket.set_read_timeout(previous_timeout)?;
Ok(drained)
}
fn recv_socket(socket: &Socket, buffer: &mut [u8]) -> io::Result<(usize, Option<SocketAddr>)> {
// `recv_from` initializes exactly the returned prefix of the buffer.
let uninitialized = unsafe {
std::slice::from_raw_parts_mut(buffer.as_mut_ptr().cast::<MaybeUninit<u8>>(), buffer.len())
};
let (bytes, source) = socket.recv_from(uninitialized)?;
Ok((bytes, source.as_socket()))
}
fn match_time_exceeded(
packet: &[u8],
source: Option<SocketAddr>,
local_addr: SocketAddr,
target: SocketAddr,
) -> Option<IpAddr> {
match (local_addr, target) {
(SocketAddr::V4(local), SocketAddr::V4(target)) => {
match_ipv4_time_exceeded(packet, local, target).map(IpAddr::V4)
}
(SocketAddr::V6(local), SocketAddr::V6(target)) => {
match_ipv6_time_exceeded(packet, source, local, target)
}
_ => None,
}
}
fn match_ipv4_time_exceeded(
packet: &[u8],
local_addr: SocketAddrV4,
target: SocketAddrV4,
) -> Option<Ipv4Addr> {
let Ok(outer) = LaxSlicedPacket::from_ip(packet) else {
return None;
};
let router = match outer.net {
Some(LaxNetSlice::Ipv4(ip)) => ip.header().source_addr(),
_ => return None,
};
let Some(TransportSlice::Icmpv4(icmp)) = outer.transport else {
return None;
};
if matches!(
icmp.icmp_type(),
Icmpv4Type::TimeExceeded(icmpv4::TimeExceededCode::TtlExceededInTransit)
) && matching_quoted_tcp_tuple(
icmp.payload(),
SocketAddr::V4(local_addr),
SocketAddr::V4(target),
) {
Some(router)
} else {
None
}
}
fn match_ipv6_time_exceeded(
packet: &[u8],
source: Option<SocketAddr>,
local_addr: SocketAddrV6,
target: SocketAddrV6,
) -> Option<IpAddr> {
let (icmp, outer_router) = if packet.first().map(|byte| byte >> 4) == Some(6) {
let Ok(outer) = LaxSlicedPacket::from_ip(packet) else {
return None;
};
let router = match outer.net {
Some(LaxNetSlice::Ipv6(ip)) => IpAddr::V6(ip.header().source_addr()),
_ => return None,
};
let Some(TransportSlice::Icmpv6(icmp)) = outer.transport else {
return None;
};
(icmp, Some(router))
} else {
let Ok(icmp) = Icmpv6Slice::from_slice(packet) else {
return None;
};
(icmp, None)
};
if !matches!(
icmp.icmp_type(),
Icmpv6Type::TimeExceeded(icmpv6::TimeExceededCode::HopLimitExceeded)
) || !matching_quoted_tcp_tuple(
icmp.payload(),
SocketAddr::V6(local_addr),
SocketAddr::V6(target),
) {
return None;
}
outer_router.or_else(|| match source {
Some(SocketAddr::V6(source)) => Some(IpAddr::V6(*source.ip())),
_ => None,
})
}
fn matching_quoted_tcp_tuple(packet: &[u8], local_addr: SocketAddr, target: SocketAddr) -> bool {
let Ok(quoted) = LaxSlicedPacket::from_ip(packet) else {
return false;
};
let payload = match (quoted.net, local_addr, target) {
(Some(LaxNetSlice::Ipv4(ip)), SocketAddr::V4(local), SocketAddr::V4(target))
if ip.header().source_addr() == *local.ip()
&& ip.header().destination_addr() == *target.ip()
&& ip.payload().ip_number == IpNumber::TCP =>
{
ip.payload().payload
}
(Some(LaxNetSlice::Ipv6(ip)), SocketAddr::V6(local), SocketAddr::V6(target))
if ip.header().source_addr() == *local.ip()
&& ip.header().destination_addr() == *target.ip()
&& ip.payload().ip_number == IpNumber::TCP =>
{
ip.payload().payload
}
_ => return false,
};
let Some(ports) = payload.get(..4) else {
return false;
};
u16::from_be_bytes([ports[0], ports[1]]) == local_addr.port()
&& u16::from_be_bytes([ports[2], ports[3]]) == target.port()
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn matches_icmp_time_exceeded_quote_by_flow_tuple() {
let local = SocketAddrV4::new(Ipv4Addr::new(192, 0, 2, 10), 45_000);
let target = SocketAddrV4::new(Ipv4Addr::new(203, 0, 113, 10), 443);
let router = Ipv4Addr::new(198, 51, 100, 1);
let mut packet = vec![0; 20 + 8 + 20 + 4];
packet[0] = 0x45;
packet[2..4].copy_from_slice(&52u16.to_be_bytes());
packet[8] = 64;
packet[9] = IpNumber::ICMP.0;
packet[12..16].copy_from_slice(&router.octets());
packet[20] = 11;
packet[28] = 0x45;
packet[30..32].copy_from_slice(&24u16.to_be_bytes());
packet[36] = 1;
packet[37] = IpNumber::TCP.0;
packet[40..44].copy_from_slice(&local.ip().octets());
packet[44..48].copy_from_slice(&target.ip().octets());
packet[48..50].copy_from_slice(&local.port().to_be_bytes());
packet[50..52].copy_from_slice(&target.port().to_be_bytes());
assert_eq!(
match_time_exceeded(&packet, None, SocketAddr::V4(local), SocketAddr::V4(target)),
Some(IpAddr::V4(router))
);
}
#[test]
fn rejects_icmp_time_exceeded_for_another_port() {
let local = SocketAddrV4::new(Ipv4Addr::new(192, 0, 2, 10), 45_000);
let target = SocketAddrV4::new(Ipv4Addr::new(203, 0, 113, 10), 443);
let mut packet = vec![0; 20 + 8 + 20 + 4];
packet[0] = 0x45;
packet[2..4].copy_from_slice(&52u16.to_be_bytes());
packet[9] = IpNumber::ICMP.0;
packet[20] = 11;
packet[28] = 0x45;
packet[30..32].copy_from_slice(&24u16.to_be_bytes());
packet[37] = IpNumber::TCP.0;
packet[40..44].copy_from_slice(&local.ip().octets());
packet[44..48].copy_from_slice(&target.ip().octets());
packet[48..50].copy_from_slice(&local.port().to_be_bytes());
packet[50..52].copy_from_slice(&444u16.to_be_bytes());
assert_eq!(
match_time_exceeded(&packet, None, SocketAddr::V4(local), SocketAddr::V4(target)),
None
);
}
#[test]
fn matches_ipv6_icmp_time_exceeded_quote_by_flow_tuple() {
let local = "[2001:db8::1]:45000".parse::<SocketAddrV6>().unwrap();
let target = "[2001:db8::10]:443".parse::<SocketAddrV6>().unwrap();
let router = "2001:db8::ff".parse::<std::net::Ipv6Addr>().unwrap();
let mut packet = vec![0; 8 + 40 + 4];
packet[0] = 3;
packet[8] = 0x60;
packet[14] = IpNumber::TCP.0;
packet[16..32].copy_from_slice(&local.ip().octets());
packet[32..48].copy_from_slice(&target.ip().octets());
packet[48..50].copy_from_slice(&local.port().to_be_bytes());
packet[50..52].copy_from_slice(&target.port().to_be_bytes());
assert_eq!(
match_time_exceeded(
&packet,
Some(SocketAddr::V6(SocketAddrV6::new(router, 0, 0, 0))),
SocketAddr::V6(local),
SocketAddr::V6(target)
),
Some(IpAddr::V6(router))
);
packet[50..52].copy_from_slice(&444u16.to_be_bytes());
assert_eq!(
match_time_exceeded(
&packet,
Some(SocketAddr::V6(SocketAddrV6::new(router, 0, 0, 0))),
SocketAddr::V6(local),
SocketAddr::V6(target)
),
None
);
}
}
+213 -92
View File
@@ -1,29 +1,34 @@
mod dns;
mod dpi_hop;
mod sni;
mod traceroute;
use anyhow::{Context, Result, bail};
use clap::Parser;
use futures::future::join_all;
use log::{error, info, warn};
use rand::seq::SliceRandom;
use reports::probe::{ProbeConfig, ProbeResult, ProbeStatus, ProbeTask, TcpTracerouteOutcome};
use log::{debug, error, info, warn};
use reports::probe::{DpiProbeConfig, ProbeConfig, ProbeResult, ProbeStatus, ProbeTask};
use rumqttc::{
AsyncClient, Event, Incoming, LastWill, MqttOptions, NetworkOptions, QoS, Transport,
};
use std::collections::HashSet;
use std::net::{IpAddr, Ipv4Addr, Ipv6Addr};
use std::net::IpAddr;
use std::sync::Arc;
use std::time::{Duration, Instant};
use tokio::sync::RwLock;
const CONFIG_TOPIC: &str = "probe/config/v1";
const MQTT_MAX_PACKET_SIZE: usize = 1024 * 1024;
#[derive(Clone)]
struct LoadedProbeConfig {
config: ProbeConfig,
control_hosts_v4: Vec<Ipv4Addr>,
control_hosts_v6: Vec<Ipv6Addr>,
dpi_hop_v4: Option<u8>,
dpi_hop_v6: Option<u8>,
}
#[derive(Clone, Copy, Default)]
struct DpiHops {
v4: Option<u8>,
v6: Option<u8>,
}
#[derive(Parser, Debug, Clone)]
@@ -47,14 +52,8 @@ struct Args {
#[arg(long, env = "MAX_CONCURRENT_TASKS", default_value_t = 8)]
max_concurrent_tasks: usize,
#[arg(long, env = "TRACEROUTE_MAX_HOPS", default_value_t = 5)]
traceroute_max_hops: u8,
#[arg(long, env = "TRACEROUTE_RETRIES", default_value_t = 3)]
traceroute_retries: u8,
#[arg(long, env = "TRACEROUTE_CONTROL_HOSTS", default_value_t = 3)]
traceroute_control_hosts: usize,
}
#[tokio::main]
@@ -64,27 +63,24 @@ async fn main() -> Result<()> {
if args.max_concurrent_tasks == 0 {
bail!("max_concurrent_tasks must be greater than zero");
}
if args.traceroute_max_hops == 0 {
bail!("traceroute_max_hops must be greater than zero");
}
if args.traceroute_retries == 0 {
bail!("traceroute_retries must be greater than zero");
}
if args.traceroute_control_hosts == 0 {
bail!("traceroute_control_hosts must be greater than zero");
}
let status_topic = format!("probe/status/v1/{}", args.probe_id);
let offline_status = serde_json::to_vec(&ProbeStatus {
online: false,
probe_id: &args.probe_id,
version: env!("CARGO_PKG_VERSION"),
dpi_hop_v4: None,
dpi_hop_v6: None,
})?;
let mut options = MqttOptions::new(&args.probe_id, &args.mqtt_host, args.mqtt_port);
options.set_transport(mqtt_transport(&args.mqtt_host)?);
options.set_credentials("probe", &args.probe_token);
options.set_keep_alive(Duration::from_secs(10));
options.set_max_packet_size(MQTT_MAX_PACKET_SIZE, MQTT_MAX_PACKET_SIZE);
options.set_last_will(LastWill::new(
status_topic.clone(),
offline_status,
@@ -100,7 +96,7 @@ async fn main() -> Result<()> {
eventloop.set_network_options(network_options);
wait_for_connection(&mut eventloop).await;
publish_status(&client, &status_topic, &args, true).await?;
publish_status(&client, &status_topic, &args, true, DpiHops::default()).await?;
client.subscribe(CONFIG_TOPIC, QoS::AtLeastOnce).await?;
client
.subscribe("probe/tasks/v1/+", QoS::AtLeastOnce)
@@ -115,8 +111,15 @@ async fn main() -> Result<()> {
match eventloop.poll().await {
Ok(Event::Incoming(Incoming::Publish(publish))) => {
if publish.topic == CONFIG_TOPIC {
if let Err(error) = update_config(&config, &publish.payload).await {
warn!("failed to update probe config: {error}");
match update_config(&config, &publish.payload).await {
Ok(dpi_hops) => {
if let Err(error) =
publish_status(&client, &status_topic, &args, true, dpi_hops).await
{
warn!("failed to publish probe status with DPI hop: {error}");
}
}
Err(error) => warn!("failed to update probe config: {error}"),
}
} else {
let client = client.clone();
@@ -159,7 +162,16 @@ async fn main() -> Result<()> {
Err(error) => {
error!("mqtt connection error: {error}");
wait_for_connection(&mut eventloop).await;
publish_status(&client, &status_topic, &args, true).await?;
let dpi_hops =
config
.read()
.await
.as_ref()
.map_or_else(DpiHops::default, |config| DpiHops {
v4: config.dpi_hop_v4,
v6: config.dpi_hop_v6,
});
publish_status(&client, &status_topic, &args, true, dpi_hops).await?;
client.subscribe(CONFIG_TOPIC, QoS::AtLeastOnce).await?;
client
.subscribe("probe/tasks/v1/+", QoS::AtLeastOnce)
@@ -183,7 +195,7 @@ fn mqtt_transport(mqtt_host: &str) -> Result<Transport> {
async fn update_config(
config: &Arc<RwLock<Option<LoadedProbeConfig>>>,
payload: &[u8],
) -> Result<()> {
) -> Result<DpiHops> {
let value: ProbeConfig = serde_json::from_slice(payload).context("decode probe config")?;
if value.dns_samples_per_protocol == 0 {
bail!("dns_samples_per_protocol must be greater than zero");
@@ -191,34 +203,103 @@ async fn update_config(
if !(1..=4).contains(&value.dns_spoofing_provider_threshold) {
bail!("dns_spoofing_provider_threshold must be between 1 and 4");
}
let mut control_hosts_v4 = HashSet::new();
let mut control_hosts_v6 = HashSet::new();
if value.traceroute_enabled {
for domain in &value.control_hosts {
match tokio::net::lookup_host((domain.as_str(), 443)).await {
Ok(addresses) => {
for address in addresses {
match address.ip() {
IpAddr::V4(address) => {
control_hosts_v4.insert(address);
}
IpAddr::V6(address) => {
control_hosts_v6.insert(address);
}
}
}
}
Err(error) => warn!("failed to resolve control host {domain}: {error}"),
}
validate_dpi_probe_config(value.dpi_probe.as_ref())?;
let dpi_hops = measure_dpi_hops(value.dpi_probe.as_ref()).await;
let loaded = LoadedProbeConfig {
config: value,
dpi_hop_v4: dpi_hops.v4,
dpi_hop_v6: dpi_hops.v6,
};
debug!(
"measured DPI hops: IPv4={:?}, IPv6={:?}",
loaded.dpi_hop_v4, loaded.dpi_hop_v6
);
*config.write().await = Some(loaded);
info!("updated retained probe config");
Ok(dpi_hops)
}
fn validate_dpi_probe_config(config: Option<&DpiProbeConfig>) -> Result<()> {
let Some(config) = config else {
return Ok(());
};
if config.sni.trim().is_empty() {
bail!("dpi_probe.sni must not be empty");
}
if config.target_v4.port() == 0 {
bail!("dpi_probe.target_v4 port must be greater than zero");
}
if config.target_v6.port() == 0 {
bail!("dpi_probe.target_v6 port must be greater than zero");
}
if config.connect_timeout_ms == 0 {
bail!("dpi_probe.connect_timeout_ms must be greater than zero");
}
if config.hop_timeout_ms == 0 {
bail!("dpi_probe.hop_timeout_ms must be greater than zero");
}
if config.max_ttl == 0 {
bail!("dpi_probe.max_ttl must be greater than zero");
}
if config.post_dpi_hop_limit == 0 {
bail!("dpi_probe.post_dpi_hop_limit must be greater than zero");
}
Ok(())
}
async fn measure_dpi_hops(config: Option<&DpiProbeConfig>) -> DpiHops {
let Some(config) = config else {
debug!("DPI hop measurement is not configured");
return DpiHops::default();
};
let common = |target| dpi_hop::DpiHopProbeConfig {
target,
control_sni: config.sni.clone(),
max_ttl: config.max_ttl,
connect_timeout: Duration::from_millis(config.connect_timeout_ms),
hop_timeout: Duration::from_millis(config.hop_timeout_ms),
};
let (v4, v6) = tokio::join!(
measure_dpi_hop(common(config.target_v4.into())),
measure_dpi_hop(common(config.target_v6.into())),
);
DpiHops { v4, v6 }
}
async fn measure_dpi_hop(config: dpi_hop::DpiHopProbeConfig) -> Option<u8> {
let target = config.target;
match dpi_hop::detect_dpi_hop(config).await {
Ok(result) => dpi_hop_from_result(&result),
Err(error) => {
warn!("failed to measure DPI hop for {target}: {error}");
None
}
}
*config.write().await = Some(LoadedProbeConfig {
config: value,
control_hosts_v4: control_hosts_v4.into_iter().collect(),
control_hosts_v6: control_hosts_v6.into_iter().collect(),
});
info!("updated retained probe config");
Ok(())
}
fn dpi_hop_from_result(result: &dpi_hop::DpiHopProbeResult) -> Option<u8> {
debug!(
"DPI probe completed: target={}, local={}, ClientHello={} bytes",
result.target, result.local_addr, result.client_hello_bytes
);
for hop in &result.hops {
debug!(
"DPI probe {} TTL {}: router={:?}, outcome={:?}",
result.target, hop.ttl, hop.router, hop.outcome
);
}
if let Some(closed_hop) = result
.hops
.iter()
.find(|hop| hop.outcome == dpi_hop::DpiHopProbeHopOutcome::TcpClosed)
{
warn!(
"DPI hop measurement for {} is invalid: TCP connection closed at TTL {}",
result.target, closed_hop.ttl
);
return None;
}
result.max_icmp_time_exceeded_ttl
}
async fn wait_for_connection(eventloop: &mut rumqttc::EventLoop) {
@@ -244,11 +325,14 @@ async fn publish_status(
topic: &str,
args: &Args,
online: bool,
dpi_hops: DpiHops,
) -> Result<()> {
let payload = serde_json::to_vec(&ProbeStatus {
online,
probe_id: &args.probe_id,
version: env!("CARGO_PKG_VERSION"),
dpi_hop_v4: dpi_hops.v4,
dpi_hop_v6: dpi_hops.v6,
})?;
client
@@ -283,25 +367,15 @@ async fn handle_task(
let traceroute_enabled = config
.as_ref()
.is_some_and(|config| config.config.traceroute_enabled);
let control_targets = config.as_ref().map_or_else(Vec::new, |config| {
if !traceroute_enabled {
return Vec::new();
}
let mut rng = rand::thread_rng();
match task.ip {
IpAddr::V4(_) => config
.control_hosts_v4
.choose_multiple(&mut rng, args.traceroute_control_hosts)
.copied()
.map(IpAddr::V4)
.collect(),
IpAddr::V6(_) => config
.control_hosts_v6
.choose_multiple(&mut rng, args.traceroute_control_hosts)
.copied()
.map(IpAddr::V6)
.collect(),
}
let dpi_hop = config.as_ref().and_then(|config| match task.ip {
IpAddr::V4(_) => config.dpi_hop_v4,
IpAddr::V6(_) => config.dpi_hop_v6,
});
let traceroute_range = config.as_ref().and_then(|config| {
let hop_limit = config.config.dpi_probe.as_ref()?.post_dpi_hop_limit;
dpi_hop?
.checked_add(1)
.map(|start_hop| (start_hop, hop_limit))
});
let sni_check = sni::check_sni(
config.as_ref().map(|config| &config.config),
@@ -311,11 +385,12 @@ async fn handle_task(
task.timeout_ms,
);
let target_traceroute = async {
if traceroute_enabled {
traceroute::tcp_traceroute(task.ip, args.traceroute_max_hops, args.traceroute_retries)
.await
} else {
None
match (traceroute_enabled, traceroute_range) {
(true, Some((start_hop, hop_limit))) => {
traceroute::tcp_traceroute(task.ip, start_hop, hop_limit, args.traceroute_retries)
.await
}
_ => None,
}
};
let dns_samples_per_protocol = config
@@ -333,22 +408,7 @@ async fn handle_task(
dns_samples_per_protocol,
dns_spoofing_provider_threshold,
);
let control_traceroute = async {
join_all(control_targets.into_iter().map(|target| {
traceroute::tcp_traceroute(target, args.traceroute_max_hops, args.traceroute_retries)
}))
.await
.into_iter()
.flatten()
.min_by_key(|trace| match trace.result {
TcpTracerouteOutcome::Rst { hop }
| TcpTracerouteOutcome::Connected { hop }
| TcpTracerouteOutcome::IcmpTimeExceeded { hop } => hop,
TcpTracerouteOutcome::Timeout => u8::MAX,
})
};
let (responses, target_traceroute, control_traceroute, dns) =
tokio::join!(sni_check, target_traceroute, control_traceroute, dns_check);
let (responses, target_traceroute, dns) = tokio::join!(sni_check, target_traceroute, dns_check);
let responses = responses?;
client
@@ -359,10 +419,71 @@ async fn handle_task(
serde_json::to_vec(&ProbeResult {
responses,
target_traceroute,
control_traceroute,
dpi_hop,
dns,
})?,
)
.await
.context("publish probe result")
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn decodes_separate_dpi_targets() {
let config: ProbeConfig = serde_json::from_value(serde_json::json!({
"version": "1",
"task_timeout_ms": 15_000,
"published_at": "2026-08-21T00:00:00Z",
"hosts": [],
"dpi_probe": {
"sni": "example.com",
"target_v4": "203.0.113.10:443",
"target_v6": "[2001:db8::10]:443",
"connect_timeout_ms": 5_000,
"hop_timeout_ms": 1_000,
"max_ttl": 15
}
}))
.unwrap();
let dpi = config.dpi_probe.unwrap();
assert_eq!(dpi.target_v4.to_string(), "203.0.113.10:443");
assert_eq!(dpi.target_v6.to_string(), "[2001:db8::10]:443");
assert_eq!(dpi.post_dpi_hop_limit, 3);
}
#[test]
fn status_payload_contains_separate_dpi_hops() {
let status = ProbeStatus {
online: true,
probe_id: "probe-1",
version: "1.0.0",
dpi_hop_v4: Some(4),
dpi_hop_v6: Some(6),
};
let value = serde_json::to_value(status).unwrap();
assert_eq!(value["dpi_hop_v4"], 4);
assert_eq!(value["dpi_hop_v6"], 6);
}
#[test]
fn tcp_closed_invalidates_only_its_measurement() {
let result = dpi_hop::DpiHopProbeResult {
target: "[2001:db8::10]:443".parse().unwrap(),
local_addr: "[2001:db8::1]:45000".parse().unwrap(),
client_hello_bytes: 256,
max_icmp_time_exceeded_ttl: Some(4),
hops: vec![dpi_hop::DpiHopProbeHop {
ttl: 5,
router: None,
outcome: dpi_hop::DpiHopProbeHopOutcome::TcpClosed,
}],
};
assert_eq!(dpi_hop_from_result(&result), None);
}
}
+33 -10
View File
@@ -14,10 +14,11 @@ const HOP_TIMEOUT: Duration = Duration::from_secs(1);
pub async fn tcp_traceroute(
target: IpAddr,
max_hops: u8,
start_hop: u8,
hop_limit: u8,
retries: u8,
) -> Option<TcpTracerouteResult> {
tokio::task::spawn_blocking(move || trace_blocking(target, max_hops, retries))
tokio::task::spawn_blocking(move || trace_blocking(target, start_hop, hop_limit, retries))
.await
.map_err(|error| log::warn!("TCP traceroute task failed for {target}: {error}"))
.ok()?
@@ -25,15 +26,24 @@ pub async fn tcp_traceroute(
.ok()
}
fn trace_blocking(target: IpAddr, max_hops: u8, retries: u8) -> io::Result<TcpTracerouteResult> {
fn trace_blocking(
target: IpAddr,
start_hop: u8,
hop_limit: u8,
retries: u8,
) -> io::Result<TcpTracerouteResult> {
if start_hop == 0 || hop_limit == 0 || retries == 0 {
return Err(io::Error::new(
io::ErrorKind::InvalidInput,
"start_hop, hop_limit, and retries must be greater than zero",
));
}
let (domain, icmp_protocol) = match target {
IpAddr::V4(_) => (Domain::IPV4, Protocol::ICMPV4),
IpAddr::V6(_) => (Domain::IPV6, Protocol::ICMPV6),
};
let receiver = Socket::new(domain, Type::RAW, Some(icmp_protocol))?;
let mut last_icmp_hop = None;
for ttl in 1..=max_hops {
for ttl in post_dpi_hops(start_hop, hop_limit) {
let destination = SockAddr::from(SocketAddr::new(target, HTTPS_PORT));
let mut tcp_attempts = Vec::with_capacity(retries as usize);
let mut last_error = None;
@@ -72,7 +82,12 @@ fn trace_blocking(target: IpAddr, max_hops: u8, retries: u8) -> io::Result<TcpTr
return Err(last_error.unwrap_or_else(|| io::Error::other("no traceroute attempts")));
}
match wait_for_hop_response(&receiver, &tcp_attempts, target)? {
HopResponse::IcmpTimeExceeded => last_icmp_hop = Some(ttl),
HopResponse::IcmpTimeExceeded => {
return Ok(TcpTracerouteResult {
target,
result: TcpTracerouteOutcome::IcmpTimeExceeded { hop: ttl },
});
}
HopResponse::Rst => {
return Ok(TcpTracerouteResult {
target,
@@ -91,12 +106,14 @@ fn trace_blocking(target: IpAddr, max_hops: u8, retries: u8) -> io::Result<TcpTr
Ok(TcpTracerouteResult {
target,
result: last_icmp_hop
.map(|hop| TcpTracerouteOutcome::IcmpTimeExceeded { hop })
.unwrap_or(TcpTracerouteOutcome::Timeout),
result: TcpTracerouteOutcome::Timeout,
})
}
fn post_dpi_hops(start_hop: u8, hop_limit: u8) -> impl Iterator<Item = u8> {
(0..hop_limit).map_while(move |offset| start_hop.checked_add(offset))
}
fn is_connect_in_progress(error: &io::Error) -> bool {
error.kind() == io::ErrorKind::WouldBlock
|| error.raw_os_error() == Some(rustix::io::Errno::INPROGRESS.raw_os_error())
@@ -303,4 +320,10 @@ mod tests {
assert!(matching_ipv6_time_exceeded(&packet, target, 42_000));
assert!(!matching_ipv6_time_exceeded(&packet, target, 42_001));
}
#[test]
fn post_dpi_hop_range_starts_at_requested_hop_and_is_bounded() {
assert_eq!(post_dpi_hops(6, 3).collect::<Vec<_>>(), vec![6, 7, 8]);
assert_eq!(post_dpi_hops(254, 3).collect::<Vec<_>>(), vec![254, 255]);
}
}
+26 -5
View File
@@ -1,11 +1,15 @@
use serde::{Deserialize, Serialize};
use std::net::IpAddr;
use std::net::{IpAddr, SocketAddrV4, SocketAddrV6};
#[derive(Clone, Serialize, Deserialize)]
pub struct ProbeStatus<'a> {
pub online: bool,
pub probe_id: &'a str,
pub version: &'a str,
#[serde(default)]
pub dpi_hop_v4: Option<u8>,
#[serde(default)]
pub dpi_hop_v6: Option<u8>,
}
#[derive(Clone, Serialize, Deserialize)]
@@ -16,12 +20,28 @@ pub struct ProbeConfig {
pub hosts: Vec<Host>,
#[serde(default)]
pub traceroute_enabled: bool,
#[serde(default)]
pub control_hosts: Vec<String>,
#[serde(default = "default_dns_samples_per_protocol")]
pub dns_samples_per_protocol: u8,
#[serde(default = "default_dns_spoofing_provider_threshold")]
pub dns_spoofing_provider_threshold: u8,
#[serde(default)]
pub dpi_probe: Option<DpiProbeConfig>,
}
#[derive(Clone, Serialize, Deserialize)]
pub struct DpiProbeConfig {
pub sni: String,
pub target_v4: SocketAddrV4,
pub target_v6: SocketAddrV6,
pub connect_timeout_ms: u64,
pub hop_timeout_ms: u64,
pub max_ttl: u8,
#[serde(default = "default_post_dpi_hop_limit")]
pub post_dpi_hop_limit: u8,
}
pub const fn default_post_dpi_hop_limit() -> u8 {
3
}
pub const fn default_dns_samples_per_protocol() -> u8 {
@@ -64,7 +84,7 @@ pub struct ProbeResultEvent {
pub probe_id: String,
pub host_results: Vec<HostProbeResult>,
pub target_traceroute: Option<TcpTracerouteResult>,
pub control_traceroute: Option<TcpTracerouteResult>,
pub dpi_hop: Option<u8>,
pub dns: Option<DnsProbeResult>,
}
@@ -72,7 +92,8 @@ pub struct ProbeResultEvent {
pub struct ProbeResult {
pub responses: Option<Vec<HostProbeResult>>,
pub target_traceroute: Option<TcpTracerouteResult>,
pub control_traceroute: Option<TcpTracerouteResult>,
#[serde(default)]
pub dpi_hop: Option<u8>,
#[serde(default)]
pub dns: Option<DnsProbeResult>,
}
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "website"
version = "1.2.1"
version = "1.2.2"
edition = "2024"
[dependencies]
@@ -0,0 +1,3 @@
ALTER TABLE probe_reports
DROP COLUMN control_hop_count,
DROP COLUMN control_trace_result;
+9 -1
View File
@@ -2,7 +2,15 @@ timeout_sec = 3
min_data = 65536
dns_samples_per_protocol = 3
dns_spoofing_provider_threshold = 2
control_hosts = ["kinopoisk.ru", "wildberries.ru", "yandex.ru", "mail.ru"]
[dpi_probe]
sni = "rutracker.org"
target_v4 = "188.120.251.237:443"
target_v6 = "[2a01:230:4:ca4::2]:443"
connect_timeout_ms = 5000
hop_timeout_ms = 1000
max_ttl = 15
post_dpi_hop_limit = 3
[[hosts]]
id = "hil-hetzner" # hil-speed.hetzner.com
+61 -66
View File
@@ -107,7 +107,6 @@ pub async fn probe_query(
Ok(result) => {
responded_probes.insert(result.probe_id.clone());
let target_traceroute = result.target_traceroute.clone();
let control_traceroute = result.control_traceroute.clone();
let reporter_info = match fetch_probe_reporter_info(&result.probe_id, &pool).await {
Ok(info) => info,
Err(error) => {
@@ -123,7 +122,6 @@ pub async fn probe_query(
query_id,
&response,
target_traceroute.as_ref(),
control_traceroute.as_ref(),
&pool,
).await {
warn!("api: failed to save probe report for query {id}: {error}");
@@ -157,7 +155,7 @@ pub fn build_probe_response(
&raw.host_results,
config,
raw.target_traceroute.as_ref(),
raw.control_traceroute.as_ref(),
raw.dpi_hop,
raw.dns.as_ref(),
);
let target_hop =
@@ -197,6 +195,7 @@ pub fn build_probe_response(
"verdicts": verdicts,
"host_results": host_results,
"target_hop": target_hop,
"dpi_hop": raw.dpi_hop,
"dns": raw.dns,
})
}
@@ -205,7 +204,6 @@ async fn insert_probe_report(
query_id: Uuid,
response: &Value,
target_traceroute: Option<&TcpTracerouteResult>,
control_traceroute: Option<&TcpTracerouteResult>,
pool: &PgPool,
) -> Result<(), sqlx::Error> {
let probe_id = response
@@ -223,7 +221,6 @@ async fn insert_probe_report(
})
.unwrap_or_else(|| vec!["uncertain"]);
let (target_hop_count, target_trace_result) = traceroute_columns(target_traceroute);
let (control_hop_count, control_trace_result) = traceroute_columns(control_traceroute);
let Some(probe_id) = probe_id else {
warn!("api: ignoring probe report with non-numeric probe_id");
@@ -238,20 +235,16 @@ async fn insert_probe_report(
verdicts,
result,
target_hop_count,
target_trace_result,
control_hop_count,
control_trace_result
target_trace_result
)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
VALUES ($1, $2, $3, $4, $5, $6)
ON CONFLICT (query_id, probe_id)
DO UPDATE SET
date = NOW(),
verdicts = EXCLUDED.verdicts,
result = EXCLUDED.result,
target_hop_count = EXCLUDED.target_hop_count,
target_trace_result = EXCLUDED.target_trace_result,
control_hop_count = EXCLUDED.control_hop_count,
control_trace_result = EXCLUDED.control_trace_result
target_trace_result = EXCLUDED.target_trace_result
"#,
)
.bind(query_id)
@@ -260,8 +253,6 @@ async fn insert_probe_report(
.bind(response)
.bind(target_hop_count)
.bind(target_trace_result)
.bind(control_hop_count)
.bind(control_trace_result)
.execute(pool)
.await?;
@@ -299,26 +290,36 @@ fn build_probe_verdicts(
results: &[HostProbeResult],
config: &ProbeConfig,
target_traceroute: Option<&TcpTracerouteResult>,
control_traceroute: Option<&TcpTracerouteResult>,
dpi_hop: Option<u8>,
dns: Option<&reports::probe::DnsProbeResult>,
) -> Vec<&'static str> {
let mut verdicts = Vec::new();
let dns_spoofing = dns.is_some_and(|result| result.spoofing_detected);
if let (
Some(TcpTracerouteResult {
result: TcpTracerouteOutcome::IcmpTimeExceeded { hop: target_hop },
..
}),
Some(TcpTracerouteResult {
result: TcpTracerouteOutcome::IcmpTimeExceeded { hop: control_hop },
..
}),
) = (target_traceroute, control_traceroute)
&& target_hop < control_hop
{
verdicts.push("tspu_block");
if results.is_empty() && target_traceroute.is_none() && dns.is_none() {
return vec!["uncertain"];
}
let tspu_block = dpi_hop.is_some()
&& target_traceroute
.is_some_and(|traceroute| matches!(&traceroute.result, TcpTracerouteOutcome::Timeout));
let host_verdict = build_host_verdict(results, config);
let dns_spoofing = dns.is_some_and(|result| result.spoofing_detected);
let mut verdicts = [
tspu_block.then_some("tspu_block"),
(!matches!(host_verdict, "ok" | "uncertain")).then_some(host_verdict),
dns_spoofing.then_some("dns_spoofing"),
]
.into_iter()
.flatten()
.collect::<Vec<_>>();
if verdicts.is_empty() {
verdicts.push(host_verdict);
}
verdicts
}
fn build_host_verdict(results: &[HostProbeResult], config: &ProbeConfig) -> &'static str {
let matched = results
.iter()
.filter_map(|result| {
@@ -330,7 +331,7 @@ fn build_probe_verdicts(
})
.collect::<Vec<_>>();
let host_verdict = if results.is_empty() {
if results.is_empty() {
"ok"
} else if matched.is_empty() {
"uncertain"
@@ -382,19 +383,7 @@ fn build_probe_verdicts(
} else {
"uncertain"
}
};
if !matches!(host_verdict, "ok" | "uncertain") {
verdicts.push(host_verdict);
}
if dns_spoofing {
verdicts.push("dns_spoofing");
}
if verdicts.is_empty() {
verdicts.push(host_verdict);
}
verdicts
}
fn publish_error_status(error: PublishError) -> Status {
@@ -429,40 +418,46 @@ fn is_strict_majority(total: usize, count: usize) -> bool {
mod tests {
use super::*;
fn icmp_trace(hop: u8) -> TcpTracerouteResult {
TcpTracerouteResult {
target: "192.0.2.1".parse().unwrap(),
result: TcpTracerouteOutcome::IcmpTimeExceeded { hop },
}
#[test]
fn no_checks_returns_uncertain() {
assert_eq!(
build_probe_verdicts(&[], &empty_config(), None, None, None),
vec!["uncertain"]
);
}
#[test]
fn tspu_block_requires_an_earlier_target_icmp_hop() {
let target = icmp_trace(3);
let control = icmp_trace(5);
fn tspu_block_requires_a_dpi_hop_and_post_dpi_timeout() {
let timeout = TcpTracerouteResult {
target: "192.0.2.1".parse().unwrap(),
result: TcpTracerouteOutcome::Timeout,
};
assert_eq!(
build_probe_verdicts(&[], &empty_config(), Some(&target), Some(&control), None),
build_probe_verdicts(&[], &empty_config(), Some(&timeout), Some(4), None),
vec!["tspu_block"]
);
let target = icmp_trace(5);
assert_eq!(
build_probe_verdicts(&[], &empty_config(), Some(&target), Some(&control), None),
build_probe_verdicts(&[], &empty_config(), Some(&timeout), None, None),
vec!["ok"]
);
}
#[test]
fn tspu_block_requires_two_icmp_outcomes() {
let target = TcpTracerouteResult {
target: "192.0.2.1".parse().unwrap(),
result: TcpTracerouteOutcome::Connected { hop: 2 },
};
let control = icmp_trace(5);
assert_eq!(
build_probe_verdicts(&[], &empty_config(), Some(&target), Some(&control), None),
vec!["ok"]
);
fn any_post_dpi_response_clears_tspu_block() {
for result in [
TcpTracerouteOutcome::IcmpTimeExceeded { hop: 5 },
TcpTracerouteOutcome::Connected { hop: 5 },
TcpTracerouteOutcome::Rst { hop: 5 },
] {
let target = TcpTracerouteResult {
target: "192.0.2.1".parse().unwrap(),
result,
};
assert_eq!(
build_probe_verdicts(&[], &empty_config(), Some(&target), Some(4), None),
vec!["ok"]
);
}
}
#[test]
@@ -501,10 +496,10 @@ mod tests {
published_at: String::new(),
hosts: vec![],
traceroute_enabled: true,
control_hosts: vec![],
dns_samples_per_protocol: reports::probe::default_dns_samples_per_protocol(),
dns_spoofing_provider_threshold:
reports::probe::default_dns_spoofing_provider_threshold(),
dpi_probe: None,
}
}
}
+12 -7
View File
@@ -1,6 +1,8 @@
use log::{info, warn};
use reports::probe::HostType;
use reports::probe::{Host, ProbeConfig, ProbeResult, ProbeResultEvent, ProbeStatus, ProbeTask};
use reports::probe::{
DpiProbeConfig, Host, ProbeConfig, ProbeResult, ProbeResultEvent, ProbeStatus, ProbeTask,
};
use rocket::serde::json::serde_json;
use rumqttc::{AsyncClient, Event as MqttEvent, Incoming, MqttOptions, QoS};
use serde::Deserialize;
@@ -12,6 +14,8 @@ use std::net::IpAddr;
use std::sync::Arc;
use std::time::Duration;
const MQTT_MAX_PACKET_SIZE: usize = 1024 * 1024;
const DEFAULT_PROBE_HOSTS: &str = include_str!("../probe-hosts.toml");
#[derive(Debug)]
@@ -66,7 +70,7 @@ struct ProbeHostsFile {
#[serde(default = "reports::probe::default_dns_spoofing_provider_threshold")]
dns_spoofing_provider_threshold: u8,
#[serde(default)]
control_hosts: Vec<String>,
dpi_probe: Option<DpiProbeConfig>,
hosts: Vec<ProbeHostEntry>,
}
@@ -93,10 +97,10 @@ impl MqttPublisher {
published_at: Utc::now().to_rfc3339(),
hosts: Vec::new(),
traceroute_enabled: false,
control_hosts: Vec::new(),
dns_samples_per_protocol: reports::probe::default_dns_samples_per_protocol(),
dns_spoofing_provider_threshold:
reports::probe::default_dns_spoofing_provider_threshold(),
dpi_probe: None,
}
}));
let admin_token = match std::env::var("MQTT_ADMIN_TOKEN") {
@@ -124,6 +128,7 @@ impl MqttPublisher {
let mut options = MqttOptions::new(client_id, host.clone(), port);
options.set_credentials("admin", admin_token);
options.set_keep_alive(Duration::from_secs(10));
options.set_max_packet_size(MQTT_MAX_PACKET_SIZE, MQTT_MAX_PACKET_SIZE);
let (client, mut eventloop) = AsyncClient::new(options, 100);
let event_sessions = sessions.clone();
@@ -278,17 +283,17 @@ fn load_probe_config(task_timeout_ms: u64) -> Result<ProbeConfig, PublishError>
traceroute_enabled: std::env::var("PROBE_TRACEROUTE_ENABLED")
.ok()
.is_some_and(|value| value == "1" || value.eq_ignore_ascii_case("true")),
control_hosts: config.control_hosts,
dns_samples_per_protocol: config.dns_samples_per_protocol,
dns_spoofing_provider_threshold: config.dns_spoofing_provider_threshold,
dpi_probe: config.dpi_probe,
})
}
struct ParsedProbeConfig {
hosts: Vec<Host>,
control_hosts: Vec<String>,
dns_samples_per_protocol: u8,
dns_spoofing_provider_threshold: u8,
dpi_probe: Option<DpiProbeConfig>,
}
fn parse_probe_hosts(contents: &str) -> Result<ParsedProbeConfig, PublishError> {
@@ -309,9 +314,9 @@ fn parse_probe_hosts(contents: &str) -> Result<ParsedProbeConfig, PublishError>
Ok(ParsedProbeConfig {
hosts,
control_hosts: config.control_hosts,
dns_samples_per_protocol: config.dns_samples_per_protocol,
dns_spoofing_provider_threshold: config.dns_spoofing_provider_threshold,
dpi_probe: config.dpi_probe,
})
}
@@ -369,7 +374,7 @@ async fn dispatch_probe_result(
probe_id: probe_id.to_string(),
host_results: result.responses.unwrap_or_default(),
target_traceroute: result.target_traceroute,
control_traceroute: result.control_traceroute,
dpi_hop: result.dpi_hop,
dns: result.dns,
});
}