feat: admin interface

This commit is contained in:
Lowder
2026-09-30 12:56:26 +05:00
parent 0bbd2be8ca
commit 816674c17a
13 changed files with 1459 additions and 45 deletions
Generated
+8 -6
View File
@@ -1223,7 +1223,7 @@ dependencies = [
"idna",
"ipnet",
"jni",
"rand 0.10.2",
"rand 0.10.3",
"rustls 0.23.44",
"thiserror 2.0.20",
"tinyvec",
@@ -1246,7 +1246,7 @@ dependencies = [
"jni",
"once_cell",
"prefix-trie 0.8.4",
"rand 0.10.2",
"rand 0.10.3",
"ring",
"thiserror 2.0.20",
"tinyvec",
@@ -1271,7 +1271,7 @@ dependencies = [
"ndk-context",
"once_cell",
"parking_lot",
"rand 0.10.2",
"rand 0.10.3",
"resolv-conf",
"rustls 0.23.44",
"smallvec",
@@ -2603,7 +2603,7 @@ dependencies = [
"bytes",
"getrandom 0.4.3",
"lru-slab",
"rand 0.10.2",
"rand 0.10.3",
"rand_pcg",
"ring",
"rustc-hash",
@@ -2674,9 +2674,9 @@ dependencies = [
[[package]]
name = "rand"
version = "0.10.2"
version = "0.10.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c7f5fa3a058cd35567ef9bfa5e75732bee0f9e4c55fa90477bef2dfcdbc4be80"
checksum = "65c9fb96cbc91e3478eaae79a69fcd3f1ae4ad052e471fe6732fff548984b4af"
dependencies = [
"chacha20",
"getrandom 0.4.3",
@@ -4517,6 +4517,7 @@ dependencies = [
"governor",
"log",
"querying",
"rand 0.10.3",
"reports",
"reqwest",
"rocket",
@@ -4528,6 +4529,7 @@ dependencies = [
"serde",
"sqlx",
"toml",
"uuid",
]
[[package]]
+15
View File
@@ -72,6 +72,21 @@ HTTP_PORT=80 docker compose up --build
`DATABASE_INTERVAL_SECONDS`, а интервал повтора после ошибки скачивания — через
`DATABASE_RETRY_INTERVAL_SECONDS` (по умолчанию 300 секунд).
## Управление сканерами
Задайте `MQTT_ADMIN_PASSWORD` в `.env` и откройте `/admin/probes`. Пароль защищает
административный API и после ввода сохраняется в localStorage браузера. Кнопка
«Выйти» удаляет его из localStorage. Для доставки команд также требуется
`MQTT_ADMIN_TOKEN` — учётные данные веб-сервера для MQTT.
На странице можно создавать сканеры, менять их метаданные и флаги, отправлять
команды переподписки и traceroute, а также запускать проверку обновлений для
всех подключённых сканеров или одного сканера. `PROBE_ID` и `PROBE_TOKEN` нового сканера
показываются сразу после создания; сохраните токен, так как повторно он не
отображается. Конфигурация `website/probe-hosts.toml` монтируется в контейнер
веб-сервера: кнопка перезагрузки читает файл с диска и публикует новую
конфигурацию в MQTT.
---
## Вклад
+3
View File
@@ -32,6 +32,8 @@ services:
DATABASE_CACHE_DIR: /var/cache/cheburcheck/databases
DATABASE_RETRY_INTERVAL_SECONDS: "${DATABASE_RETRY_INTERVAL_SECONDS:-300}"
MQTT_ADMIN_TOKEN: "${MQTT_ADMIN_TOKEN}"
MQTT_ADMIN_PASSWORD: "${MQTT_ADMIN_PASSWORD}"
PROBE_CONFIG_PATH: /app/probe-hosts.toml
NODE_STATS_KEY: "${NODE_STATS_KEY}"
MQTT_HOST: rmqtt
MQTT_PORT: 11883
@@ -46,6 +48,7 @@ services:
GITHUB_TOKEN: "${GITHUB_TOKEN:-}"
volumes:
- database-cache:/var/cache/cheburcheck/databases
- ./website/probe-hosts.toml:/app/probe-hosts.toml:ro
expose:
- "8000"
+24
View File
@@ -0,0 +1,24 @@
export async function adminRequest<T>(
password: string,
path: string,
method = "GET",
body?: unknown,
): Promise<T> {
const response = await fetch("/api/v1/admin" + path, {
method,
cache: "no-store",
headers: {
Authorization: "Bearer " + password,
Accept: "application/json",
...(body ? { "Content-Type": "application/json" } : {}),
},
...(body ? { body: JSON.stringify(body) } : {}),
});
if (!response.ok)
throw new Error(
response.status === 401
? "Неверный пароль администратора"
: "Ошибка " + response.status + ": " + response.statusText,
);
return response.json() as Promise<T>;
}
@@ -0,0 +1,853 @@
<script lang="ts">
import {
Copy,
Download,
Plus,
Radio,
RadioTower,
RefreshCw,
Route,
Settings2,
Trash2,
} from "@lucide/svelte";
import {
createMutation,
createQuery,
useQueryClient,
} from "@tanstack/svelte-query";
import { onMount } from "svelte";
import { adminRequest } from "$lib/api/admin";
type Probe = {
id: number;
name: string;
region: string | null;
asn: string | null;
provider: string | null;
hidden: boolean;
disable_traceroutes: boolean;
cdn_unblocked: boolean;
last_connected_at: string | null;
online: boolean;
version: string | null;
bundle_type: string | null;
dpi_hop_v4: number | null;
dpi_hop_v6: number | null;
};
type Hop = {
ttl: number;
address: string | null;
reverse_names: string[];
outcome: string;
};
type CommandResult =
| { type: "traceroute"; target: string; hops: Hop[] }
| { type: "resubscribe_tasks"; requested: boolean }
| { type: "error"; message: string };
type Form = Pick<
Probe,
| "name"
| "region"
| "asn"
| "provider"
| "hidden"
| "disable_traceroutes"
| "cdn_unblocked"
>;
const emptyForm = (): Form => ({
name: "",
region: "",
asn: "",
provider: "",
hidden: false,
disable_traceroutes: false,
cdn_unblocked: false,
});
const passwordKey = "cheburcheck:admin-password";
let password = $state("");
let enteredPassword = $state("");
let signingIn = $state(false);
let onlineOnly = $state(false);
let busy = $state("");
let error = $state("");
let notice = $state("");
let form = $state<Form>(emptyForm());
let editing = $state<number | null>(null);
let created = $state<{ id: number; token: string } | null>(null);
let selected = $state<number | null>(null);
let target = $state("");
let maxHops = $state(30);
let result = $state<CommandResult | null>(null);
const queryClient = useQueryClient();
const probeKey = ["admin", "probes"] as const;
const probesQuery = createQuery(() => ({
queryKey: probeKey,
queryFn: () => api<Probe[]>("/probes"),
enabled: password.length > 0,
staleTime: 15_000,
refetchInterval: 30_000,
}));
const probes = $derived(probesQuery.data ?? []);
const sortedProbes = $derived(
[...probes].sort(
(a, b) => Number(b.online) - Number(a.online) || a.id - b.id,
),
);
const visibleProbes = $derived(
onlineOnly ? sortedProbes.filter((probe) => probe.online) : sortedProbes,
);
const loading = $derived(probesQuery.isFetching);
const createProbe = createMutation(() => ({
mutationFn: (input: Form) =>
api<{ id: number; token: string }>("/probes", "POST", input),
onSuccess: () => queryClient.invalidateQueries({ queryKey: probeKey }),
}));
const updateProbe = createMutation(() => ({
mutationFn: ({ id, input }: { id: number; input: Form }) =>
api<Probe>("/probes/" + id, "PUT", input),
onSuccess: () => queryClient.invalidateQueries({ queryKey: probeKey }),
}));
const removeProbe = createMutation(() => ({
mutationFn: (id: number) =>
api<{ removed: boolean }>("/probes/" + id, "DELETE"),
onSuccess: () => queryClient.invalidateQueries({ queryKey: probeKey }),
}));
const reloadProbeConfig = createMutation(() => ({
mutationFn: () =>
api<{ hosts: unknown[]; published_at: string }>(
"/probe-config/reload",
"POST",
),
}));
const requestUpdateCheck = createMutation(() => ({
mutationFn: (id: number | null) =>
api<{ requested: boolean }>(
id === null ? "/probes/update-check" : "/probes/" + id + "/update-check",
"POST",
),
}));
const sendProbeCommand = createMutation(() => ({
mutationFn: ({
id,
type,
target,
maxHops,
}: {
id: number;
type: "resubscribe_tasks" | "traceroute";
target: string;
maxHops: number;
}) =>
api<CommandResult>(
"/probes/" + id + "/commands",
"POST",
type === "traceroute" ? { type, target, max_hops: maxHops } : { type },
),
}));
$effect(() => {
if (probesQuery.isError)
error = String(
probesQuery.error instanceof Error
? probesQuery.error.message
: probesQuery.error,
);
});
function api<T>(
path: string,
method = "GET",
body?: unknown,
auth = password,
): Promise<T> {
return adminRequest<T>(auth, path, method, body);
}
async function signIn() {
const candidate = enteredPassword.trim();
error = "";
signingIn = true;
try {
queryClient.removeQueries({ queryKey: ["admin", "login"] });
const rows = await queryClient.fetchQuery({
queryKey: ["admin", "login"],
queryFn: () => api<Probe[]>("/probes", "GET", undefined, candidate),
});
queryClient.setQueryData(probeKey, rows);
password = candidate;
localStorage.setItem(passwordKey, candidate);
} catch (e) {
error = String(e instanceof Error ? e.message : e);
} finally {
queryClient.removeQueries({ queryKey: ["admin", "login"] });
signingIn = false;
}
}
function signOut() {
localStorage.removeItem(passwordKey);
password = "";
enteredPassword = "";
queryClient.removeQueries({ queryKey: probeKey });
created = null;
result = null;
error = "";
}
onMount(() => {
const saved = localStorage.getItem(passwordKey);
if (saved) {
password = saved;
enteredPassword = saved;
}
});
async function save() {
if (!form.name.trim()) {
error = "Укажите имя сканера";
return;
}
busy = "save";
error = "";
notice = "";
try {
if (editing === null) {
created = await createProbe.mutateAsync(form);
notice = `Сканер #${created.id} создан. Сохраните данные для подключения.`;
} else {
await updateProbe.mutateAsync({ id: editing, input: form });
notice = `Сканер #${editing} обновлён`;
}
form = emptyForm();
editing = null;
} catch (e) {
error = String(e instanceof Error ? e.message : e);
} finally {
busy = "";
}
}
function edit(probe: Probe) {
editing = probe.id;
form = {
name: probe.name,
region: probe.region,
asn: probe.asn,
provider: probe.provider,
hidden: probe.hidden,
disable_traceroutes: probe.disable_traceroutes,
cdn_unblocked: probe.cdn_unblocked,
};
document.getElementById("probe-form")?.scrollIntoView({ behavior: "smooth" });
}
async function toggle(
probe: Probe,
key: "hidden" | "disable_traceroutes" | "cdn_unblocked",
) {
busy = `${probe.id}:${key}`;
error = "";
try {
await updateProbe.mutateAsync({
id: probe.id,
input: {
name: probe.name,
region: probe.region,
asn: probe.asn,
provider: probe.provider,
hidden: probe.hidden,
disable_traceroutes: probe.disable_traceroutes,
cdn_unblocked: probe.cdn_unblocked,
[key]: !probe[key],
},
});
} catch (e) {
error = String(e instanceof Error ? e.message : e);
} finally {
busy = "";
}
}
async function removeConfirmed(probe: Probe) {
if (
!window.confirm(
"Удалить сканер #" +
probe.id +
" (" +
probe.name +
")? Связанные результаты сканера будут удалены. Обычные отчёты сохранятся без привязки к сканеру.",
)
)
return;
busy = probe.id + ":remove";
error = "";
notice = "";
try {
await removeProbe.mutateAsync(probe.id);
if (editing === probe.id) {
editing = null;
form = emptyForm();
}
if (selected === probe.id) {
selected = null;
result = null;
}
if (created?.id === probe.id) created = null;
notice = "Сканер #" + probe.id + " удалён";
} catch (e) {
error = String(e instanceof Error ? e.message : e);
} finally {
busy = "";
}
}
async function reloadConfig() {
busy = "reload";
error = "";
notice = "";
try {
const config = await reloadProbeConfig.mutateAsync();
notice = `Конфигурация отправлена: ${config.hosts.length} хостов, ${new Date(config.published_at).toLocaleString("ru-RU")}`;
} catch (e) {
error = String(e instanceof Error ? e.message : e);
} finally {
busy = "";
}
}
async function updateCheck(id: number | null) {
busy = id === null ? "update:all" : id + ":update";
error = "";
notice = "";
try {
await requestUpdateCheck.mutateAsync(id);
notice =
id === null
? "Запрос проверки обновлений отправлен всем подключённым сканерам"
: "Запрос проверки обновлений отправлен сканеру #" + id;
} catch (e) {
error = String(e instanceof Error ? e.message : e);
} finally {
busy = "";
}
}
async function command(id: number, type: "resubscribe_tasks" | "traceroute") {
selected = id;
result = null;
busy = `${id}:command`;
error = "";
try {
result = await sendProbeCommand.mutateAsync({
id,
type,
target: target.trim(),
maxHops,
});
} catch (e) {
error = String(e instanceof Error ? e.message : e);
} finally {
busy = "";
}
}
async function copy(value: string) {
await navigator.clipboard.writeText(value);
notice = "Скопировано";
}
function asnTarget(asn: string): string {
const value = asn.trim().toUpperCase();
return value.startsWith("AS") ? value : "AS" + value;
}
const outcomeLabel: Record<string, string> = {
icmp_time_exceeded: "Промежуточный узел",
rst: "TCP RST",
connected: "Подключено",
timeout: "Нет ответа",
};
</script>
<svelte:head
><title>Сканеры · Cheburcheck</title>
<meta name="robots" content="noindex, nofollow"></svelte:head
>
<div class="space-y-7 text-neutral-100">
<div class="flex flex-wrap items-end justify-between gap-4">
<h1 class="mt-2 text-3xl font-bold tracking-tight">Сканеры</h1>
{#if password}
<button type="button" class="btn" onclick={signOut}>Выйти</button>
{/if}
</div>
{#if !password || error === "Неверный пароль администратора"}
<form
class="panel max-w-md space-y-4"
onsubmit={(event) => { event.preventDefault(); void signIn(); }}
>
<label class="block text-sm font-medium" for="admin-password"
>Пароль администратора</label
>
<input
id="admin-password"
class="input"
type="password"
autocomplete="current-password"
bind:value={enteredPassword}
required
>
{#if error}
<p role="alert" class="text-sm text-red-300">{error}</p>
{/if}
<button type="submit" class="btn-primary" disabled={signingIn}>
{signingIn ? "Проверка…" : "Открыть"}
</button>
</form>
{:else}
<div class="flex flex-wrap gap-3">
<button
type="button"
class="btn flex items-center gap-2"
onclick={() => void probesQuery.refetch()}
disabled={loading}
>
<RefreshCw size={16} />
Обновить список
</button><button
type="button"
class="btn flex items-center gap-2"
onclick={() => void reloadConfig()}
disabled={busy !== ""}
>
<RadioTower size={16} />
Перезагрузить probe-hosts.toml
</button>
<button
type="button"
class="btn flex items-center gap-2"
onclick={() => void updateCheck(null)}
disabled={busy !== ""}
>
<Download size={16} />
Проверить обновления всех
</button>
</div>
{#if notice}
<div
role="status"
class="rounded-lg border border-emerald-800 bg-emerald-950/40 px-4 py-3 text-sm text-emerald-300"
>
{notice}
</div>
{/if}
{#if error}
<div
role="alert"
class="rounded-lg border border-red-800 bg-red-950/40 px-4 py-3 text-sm text-red-300"
>
{error}
</div>
{/if}
<section class="panel overflow-x-auto">
<div class="mb-4 flex flex-wrap items-center justify-between gap-3">
<h2 class="text-lg font-semibold">Текущие сканеры</h2>
<div class="flex flex-wrap items-center gap-4">
<label class="check text-xs">
<input type="checkbox" bind:checked={onlineOnly}>
Только онлайн
</label>
<span class="text-xs text-neutral-500"
>{probes.length}
всего · {probes.filter((p) => p.online).length} онлайн</span
>
</div>
</div>
{#if visibleProbes.length === 0}
<p class="py-8 text-center text-sm text-neutral-500">
{loading ? "Загрузка…" : onlineOnly && probes.length > 0 ? "Нет сканеров онлайн" : "Сканеров пока нет"}
</p>
{:else}
<table class="w-full min-w-[780px] text-left text-sm">
<thead
class="border-b border-neutral-800 text-xs uppercase tracking-wider text-neutral-500"
>
<tr>
<th class="py-3 pr-3">Сканер</th>
<th class="px-3">Состояние</th>
<th class="px-3">Метаданные</th>
<th class="px-3">Флаги</th>
<th class="pl-3">Действия</th>
</tr>
</thead>
<tbody>
{#each visibleProbes as probe (probe.id)}
<tr
class="border-b border-neutral-800/70 align-top last:border-0"
>
<td class="py-4 pr-3">
<div class="font-semibold">#{probe.id}</div>
<div class="mt-1 text-xs text-neutral-500">
<p>{probe.name}</p>
{probe.last_connected_at ? `${new Date(probe.last_connected_at).toLocaleString("ru-RU")}` : "Не подключался"}
</div>
</td>
<td class="px-3 py-4">
<span
class:!text-emerald-300={probe.online}
class="inline-flex items-center gap-1 text-neutral-500"
><Radio size={14} />
{probe.online ? "Онлайн" : "Офлайн"}</span
>
<div class="mt-1 text-xs text-neutral-500">
{probe.version ?? "Версия неизвестна"}
{probe.bundle_type ? `· ${probe.bundle_type}` : ""}
</div>
{#if probe.dpi_hop_v4 || probe.dpi_hop_v6}
<div class="mt-1 text-xs text-neutral-500">
v4: <b>{probe.dpi_hop_v4 ?? "—"}</b>; v6:
<b>{probe.dpi_hop_v6 ?? "—"}</b>
</div>
{/if}
</td>
<td class="px-3 py-4 text-xs text-neutral-400">
{probe.region || "Регион не указан"}<br>
{probe.provider || "Провайдер не указан"}
{#if probe.asn}
·
<a
class="text-cyan-300 underline hover:text-cyan-200"
href={"/check?target=" + encodeURIComponent(asnTarget(probe.asn))}
>{probe.asn}</a
>
{/if}
</td>
<td class="px-3 py-4">
<div class="flex flex-col items-start gap-1">
<button
type="button"
class="flag"
class:active={probe.hidden}
disabled={busy !== ""}
onclick={() => void toggle(probe, "hidden")}
>
{probe.hidden ? "Скрыт" : "Публичный"}
</button><button
type="button"
class="flag"
class:active={probe.disable_traceroutes}
disabled={busy !== ""}
onclick={() => void toggle(probe, "disable_traceroutes")}
>
{probe.disable_traceroutes ? "Трассировка выкл." : "Трассировка вкл."}
</button><button
type="button"
class="flag"
class:active={probe.cdn_unblocked}
disabled={busy !== ""}
onclick={() => void toggle(probe, "cdn_unblocked")}
>
{probe.cdn_unblocked ? "CDN доступен" : "CDN блокируется"}
</button>
</div>
</td>
<td class="pl-3 py-4">
<div class="flex flex-col items-start gap-2">
<button
type="button"
class="link"
onclick={() => edit(probe)}
>
<Settings2 size={14} />
Изменить
</button><button
type="button"
class="link"
disabled={busy !== ""}
onclick={() => void command(probe.id, "resubscribe_tasks")}
>
<RefreshCw size={14} />
Переподписать
</button><button
type="button"
class="link"
disabled={busy !== ""}
onclick={() => void updateCheck(probe.id)}
>
<Download size={14} />
Проверить обновление
</button><button
type="button"
class="link"
onclick={() => { selected = probe.id; result = null; document.getElementById("commands")?.scrollIntoView({ behavior: "smooth" }); }}
>
<Route size={14} />
Traceroute
</button>
<button
type="button"
class="link remove-link"
disabled={busy !== ""}
onclick={() => void removeConfirmed(probe)}
>
<Trash2 size={14} />
Удалить
</button>
</div>
</td>
</tr>
{/each}
</tbody>
</table>
{/if}
</section>
{#if created}
<section class="panel border-emerald-800">
<h2 class="font-semibold text-emerald-300">Данные нового сканера</h2>
<p class="mt-1 text-xs text-neutral-400">
Токен показан только сейчас. Скопируйте его перед закрытием страницы.
</p>
<div class="mt-4 flex items-start gap-3">
<pre
class="min-w-0 flex-1 overflow-x-auto rounded-lg bg-black/60 p-4 text-sm text-cyan-200"
>PROBE_ID={created.id}
PROBE_TOKEN={created.token}</pre>
<button
type="button"
class="btn"
aria-label="Скопировать переменные"
onclick={() => void copy(`PROBE_ID=${created?.id}\nPROBE_TOKEN=${created?.token}`)}
>
<Copy size={17} />
</button>
</div>
</section>
{/if}
<section id="probe-form" class="panel">
<h2 class="mb-5 text-lg font-semibold">
{editing === null ? "Добавить сканер" : `Изменить сканер #${editing}`}
</h2>
<form
class="space-y-4"
onsubmit={(event) => { event.preventDefault(); void save(); }}
>
<div class="grid gap-4 sm:grid-cols-2">
<label class="field"
>Имя<input
class="input"
bind:value={form.name}
maxlength="255"
required
></label
><label class="field"
>Регион<input
class="input"
bind:value={form.region}
maxlength="255"
></label
><label class="field"
>Провайдер<input
class="input"
bind:value={form.provider}
maxlength="255"
></label
><label class="field"
>ASN<input
class="input"
bind:value={form.asn}
maxlength="32"
></label
>
</div>
<div class="flex flex-wrap gap-5 text-sm">
<label class="check"
><input type="checkbox" bind:checked={form.hidden}>
Скрыт</label
><label class="check"
><input type="checkbox" bind:checked={form.disable_traceroutes}>
Отключить трассировки</label
><label class="check"
><input type="checkbox" bind:checked={form.cdn_unblocked}>
CDN доступен</label
>
</div>
<div class="flex gap-3">
<button
type="submit"
class="btn-primary flex items-center gap-2"
disabled={busy !== ""}
>
<Plus size={16} />{editing === null ? "Создать" : "Сохранить"}
</button>
{#if editing !== null}
<button
type="button"
class="btn"
onclick={() => { editing = null; form = emptyForm(); }}
>
Отмена
</button>
{/if}
</div>
</form>
</section>
<section id="commands" class="panel">
<h2 class="mb-2 text-lg font-semibold">Команды и трассировка</h2>
<p class="mb-5 text-sm text-neutral-400">
Команда выполняется выбранным сканером. Ответ может занять до минуты.
</p>
<div class="grid gap-3 sm:grid-cols-[1fr_2fr_100px_auto] sm:items-end">
<label class="field"
>Сканер<select class="input" bind:value={selected}>
<option value={null}>Выберите</option>
{#each probes as probe}
<option value={probe.id}>#{probe.id} {probe.name}</option>
{/each}
</select></label
><label class="field"
>IP адрес цели<input
class="input"
type="text"
placeholder="1.1.1.1"
bind:value={target}
></label
><label class="field"
>Макс. hops<input
class="input"
type="number"
min="1"
max="64"
bind:value={maxHops}
></label
><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")}
>
<Route size={16} />
Запустить
</button>
</div>
{#if busy.endsWith(":command")}
<p class="mt-5 text-sm text-cyan-300">Ожидание ответа сканера…</p>
{/if}
{#if result?.type === "error"}
<p class="mt-5 text-sm text-red-300">{result.message}</p>
{:else if result?.type === "resubscribe_tasks"}
<p class="mt-5 text-sm text-emerald-300">
Запрос на переподписку получен сканером.
</p>
{:else if result?.type === "traceroute"}
<div class="mt-6">
<h3 class="mb-4 font-semibold">Маршрут до {result.target}</h3>
<ol class="space-y-1">
{#each result.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"
class:!border-emerald-800={hop.outcome === "connected" || hop.outcome === "rst"}
>
<span class="font-mono text-neutral-500">{hop.ttl}</span>
<div>
<span class="font-mono" class:text-neutral-500={!hop.address}
>{hop.address ?? "* * *"}</span
>
{#if hop.reverse_names.length}
<span class="ml-3 text-xs text-neutral-400"
>{hop.reverse_names.join(", ")}</span
>
{/if}
</div>
<span class="text-xs text-neutral-500"
>{outcomeLabel[hop.outcome] ?? hop.outcome}</span
>
</li>
{/each}
</ol>
{#if result.hops.length === 0}
<p class="text-sm text-neutral-500">Ответов нет.</p>
{/if}
</div>
{/if}
</section>
{/if}
</div>
<style>
.panel {
border: 1px solid #303030;
border-radius: 14px;
background: #171717;
padding: 1.5rem;
}
.input {
display: block;
width: 100%;
margin-top: 0.4rem;
border: 1px solid #404040;
border-radius: 8px;
background: #0a0a0a;
padding: 0.65rem 0.75rem;
color: #eee;
outline: none;
}
.input:focus {
border-color: #22d3ee;
}
.field {
display: block;
font-size: 0.8rem;
color: #aaa;
}
.check {
display: flex;
align-items: center;
gap: 0.5rem;
color: #d4d4d4;
}
.check input {
accent-color: #22d3ee;
}
.btn,
.btn-primary {
border-radius: 8px;
padding: 0.65rem 0.9rem;
font-size: 0.85rem;
cursor: pointer;
}
.btn {
border: 1px solid #404040;
background: #262626;
color: #e5e5e5;
}
.btn-primary {
border: 1px solid #0891b2;
background: #0891b2;
color: white;
font-weight: 600;
}
button:disabled {
opacity: 0.45;
cursor: not-allowed;
}
.link {
display: inline-flex;
align-items: center;
gap: 0.4rem;
color: #67e8f9;
font-size: 0.8rem;
cursor: pointer;
}
.link:hover {
text-decoration: underline;
}
.link.remove-link {
color: #f87171;
}
.flag {
border: 1px solid #404040;
border-radius: 5px;
padding: 0.15rem 0.4rem;
color: #a3a3a3;
font-size: 0.7rem;
cursor: pointer;
}
.flag.active {
border-color: #0e7490;
color: #67e8f9;
}
</style>
+1 -1
View File
@@ -32,7 +32,7 @@ server {
location /api/v1/ {
proxy_buffering off;
proxy_cache off;
proxy_read_timeout 60s;
proxy_read_timeout 75s;
proxy_pass http://website_backend;
}
+2
View File
@@ -20,4 +20,6 @@ dotenvy = { version = "0.15.7" }
governor = { version = "0.6", features = ["dashmap"] }
rumqttc = "0.24"
toml = "0.8"
uuid = { version = "1", features = ["v4"] }
reqwest = { workspace = true }
rand = "0.10.3"
@@ -0,0 +1,2 @@
-- Keep historical reports when their reporter is removed.
ALTER TABLE reports ALTER COLUMN reporter DROP NOT NULL;
+340
View File
@@ -0,0 +1,340 @@
use crate::mqtt::{MqttPublisher, ProbeStatusSnapshot};
use rand::distr::{Alphanumeric, SampleString};
use reports::probe::{ProbeCommand, ProbeCommandResult};
use rocket::http::Status;
use rocket::request::{FromRequest, Outcome};
use rocket::serde::json::Json;
use rocket::{Request, State};
use serde::{Deserialize, Serialize};
use sqlx::PgPool;
use sqlx::types::chrono::{DateTime, Utc};
use std::net::IpAddr;
pub struct Admin;
#[rocket::async_trait]
impl<'r> FromRequest<'r> for Admin {
type Error = ();
async fn from_request(request: &'r Request<'_>) -> Outcome<Self, Self::Error> {
let configured = std::env::var("MQTT_ADMIN_PASSWORD")
.ok()
.filter(|s| !s.is_empty());
let supplied = request
.headers()
.get_one("Authorization")
.and_then(|s| s.strip_prefix("Bearer "));
match (configured, supplied) {
(Some(expected), Some(actual))
if constant_time_eq(expected.as_bytes(), actual.as_bytes()) =>
{
Outcome::Success(Admin)
}
_ => Outcome::Error((Status::Unauthorized, ())),
}
}
}
fn constant_time_eq(a: &[u8], b: &[u8]) -> bool {
let mut difference = a.len() ^ b.len();
for (left, right) in a.iter().zip(b.iter()) {
difference |= (left ^ right) as usize;
}
difference == 0
}
#[derive(Serialize, sqlx::FromRow)]
pub struct ProbeRow {
id: i32,
name: String,
region: Option<String>,
asn: Option<String>,
provider: Option<String>,
hidden: bool,
disable_traceroutes: bool,
cdn_unblocked: bool,
last_connected_at: Option<DateTime<Utc>>,
online: bool,
version: Option<String>,
bundle_type: Option<String>,
dpi_hop_v4: Option<i16>,
dpi_hop_v6: Option<i16>,
}
async fn rows(pool: &PgPool, mqtt: &MqttPublisher) -> Result<Vec<ProbeRow>, Status> {
let mut rows: Vec<ProbeRow> = sqlx::query_as(
"SELECT id, name, region, asn, provider, hidden, disable_traceroutes, cdn_unblocked, last_connected_at,
false AS online, NULL::text AS version, NULL::text AS bundle_type, NULL::smallint AS dpi_hop_v4, NULL::smallint AS dpi_hop_v6
FROM reporters ORDER BY id"
).fetch_all(pool).await.map_err(|error| {
log::error!("failed to load admin probe list: {error}");
Status::InternalServerError
})?;
let statuses = mqtt.probe_statuses().await;
for row in &mut rows {
if let Some(ProbeStatusSnapshot {
online,
version,
bundle_type,
dpi_hop_v4,
dpi_hop_v6,
}) = statuses.get(&row.id.to_string())
{
row.online = *online;
row.version = Some(version.clone());
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);
}
}
Ok(rows)
}
#[get("/probes")]
pub async fn list(
_admin: Admin,
pool: &State<PgPool>,
mqtt: &State<MqttPublisher>,
) -> Result<Json<Vec<ProbeRow>>, Status> {
Ok(Json(rows(pool, mqtt).await?))
}
#[derive(Deserialize)]
pub struct ProbeInput {
name: String,
region: Option<String>,
asn: Option<String>,
provider: Option<String>,
#[serde(default)]
hidden: bool,
#[serde(default)]
disable_traceroutes: bool,
#[serde(default)]
cdn_unblocked: bool,
}
fn validate(input: &ProbeInput) -> Result<(), Status> {
if input.name.trim().is_empty()
|| input.name.len() > 255
|| input.region.as_ref().is_some_and(|s| s.len() > 255)
|| input.asn.as_ref().is_some_and(|s| s.len() > 32)
|| input.provider.as_ref().is_some_and(|s| s.len() > 255)
{
return Err(Status::BadRequest);
}
Ok(())
}
#[derive(Serialize)]
pub struct CreatedProbe {
id: i32,
token: String,
}
#[post("/probes", format = "json", data = "<input>")]
pub async fn create(
_admin: Admin,
input: Json<ProbeInput>,
pool: &State<PgPool>,
) -> Result<Json<CreatedProbe>, Status> {
validate(&input)?;
let token = Alphanumeric.sample_string(&mut rand::rng(), 16);
let id: i32 = sqlx::query_scalar("INSERT INTO reporters (name, token, region, asn, provider, hidden, disable_traceroutes, cdn_unblocked)
VALUES ($1,$2,$3,$4,$5,$6,$7,$8) RETURNING id")
.bind(input.name.trim()).bind(&token).bind(&input.region).bind(&input.asn).bind(&input.provider)
.bind(input.hidden).bind(input.disable_traceroutes).bind(input.cdn_unblocked)
.fetch_one(&**pool).await.map_err(|_| Status::InternalServerError)?;
Ok(Json(CreatedProbe { id, token }))
}
#[put("/probes/<id>", format = "json", data = "<input>")]
pub async fn update(
_admin: Admin,
id: i32,
input: Json<ProbeInput>,
pool: &State<PgPool>,
mqtt: &State<MqttPublisher>,
) -> Result<Json<ProbeRow>, Status> {
validate(&input)?;
let updated: Option<i32> = sqlx::query_scalar("UPDATE reporters SET name=$2, region=$3, asn=$4, provider=$5, hidden=$6, disable_traceroutes=$7, cdn_unblocked=$8 WHERE id=$1 RETURNING id")
.bind(id).bind(input.name.trim()).bind(&input.region).bind(&input.asn).bind(&input.provider)
.bind(input.hidden).bind(input.disable_traceroutes).bind(input.cdn_unblocked)
.fetch_optional(&**pool).await.map_err(|_| Status::InternalServerError)?;
if updated.is_none() {
return Err(Status::NotFound);
}
rows(pool, mqtt)
.await?
.into_iter()
.find(|row| row.id == id)
.map(Json)
.ok_or(Status::NotFound)
}
#[derive(Serialize)]
pub struct RemovedProbe {
removed: bool,
}
#[delete("/probes/<id>")]
pub async fn remove(
_admin: Admin,
id: i32,
pool: &State<PgPool>,
) -> Result<Json<RemovedProbe>, Status> {
let mut tx = pool.begin().await.map_err(|error| {
log::error!("failed to begin probe removal {id}: {error}");
Status::InternalServerError
})?;
let exists: Option<i32> =
sqlx::query_scalar("SELECT id FROM reporters WHERE id = $1 FOR UPDATE")
.bind(id)
.fetch_optional(&mut *tx)
.await
.map_err(|error| {
log::error!("failed to lock probe {id} for removal: {error}");
Status::InternalServerError
})?;
if exists.is_none() {
return Err(Status::NotFound);
}
sqlx::query("UPDATE reports SET reporter = NULL WHERE reporter = $1")
.bind(id)
.execute(&mut *tx)
.await
.map_err(|error| {
log::error!("failed to detach reports from probe {id}: {error}");
Status::InternalServerError
})?;
sqlx::query("DELETE FROM reporters WHERE id = $1")
.bind(id)
.execute(&mut *tx)
.await
.map_err(|error| {
log::error!("failed to delete probe {id}: {error}");
Status::InternalServerError
})?;
tx.commit().await.map_err(|error| {
log::error!("failed to commit probe removal {id}: {error}");
Status::InternalServerError
})?;
Ok(Json(RemovedProbe { removed: true }))
}
#[post("/probe-config/reload")]
pub async fn reload(
_admin: Admin,
mqtt: &State<MqttPublisher>,
) -> Result<Json<reports::probe::ProbeConfig>, Status> {
mqtt.reload_probe_config()
.await
.map(Json)
.map_err(|_| Status::InternalServerError)
}
#[derive(Serialize)]
pub struct UpdateCheckResponse {
requested: bool,
}
#[post("/probes/update-check")]
pub async fn update_all_probes(
_admin: Admin,
mqtt: &State<MqttPublisher>,
) -> Result<Json<UpdateCheckResponse>, Status> {
mqtt.request_probe_update(None)
.await
.map_err(|_| Status::ServiceUnavailable)?;
Ok(Json(UpdateCheckResponse { requested: true }))
}
#[post("/probes/<id>/update-check")]
pub async fn update_one_probe(
_admin: Admin,
id: i32,
pool: &State<PgPool>,
mqtt: &State<MqttPublisher>,
) -> Result<Json<UpdateCheckResponse>, Status> {
let exists: bool = sqlx::query_scalar("SELECT EXISTS(SELECT 1 FROM reporters WHERE id=$1)")
.bind(id)
.fetch_one(&**pool)
.await
.map_err(|_| Status::InternalServerError)?;
if !exists {
return Err(Status::NotFound);
}
mqtt.request_probe_update(Some(&id.to_string()))
.await
.map_err(|_| Status::ServiceUnavailable)?;
Ok(Json(UpdateCheckResponse { requested: true }))
}
#[derive(Deserialize)]
#[serde(tag = "type", rename_all = "snake_case")]
pub enum CommandInput {
ResubscribeTasks,
Traceroute { target: IpAddr, max_hops: u8 },
}
#[post("/probes/<id>/commands", format = "json", data = "<input>")]
pub async fn command(
_admin: Admin,
id: i32,
input: Json<CommandInput>,
pool: &State<PgPool>,
mqtt: &State<MqttPublisher>,
) -> Result<Json<ProbeCommandResult>, Status> {
let exists: bool = sqlx::query_scalar("SELECT EXISTS(SELECT 1 FROM reporters WHERE id=$1)")
.bind(id)
.fetch_one(&**pool)
.await
.map_err(|_| Status::InternalServerError)?;
if !exists {
return Err(Status::NotFound);
}
let command = match input.into_inner() {
CommandInput::ResubscribeTasks => ProbeCommand::ResubscribeTasks,
CommandInput::Traceroute { target, max_hops } if (1..=64).contains(&max_hops) => {
ProbeCommand::Traceroute { target, max_hops }
}
CommandInput::Traceroute { .. } => return Err(Status::BadRequest),
};
mqtt.send_command(&id.to_string(), command)
.await
.map(Json)
.map_err(|error| match error {
crate::mqtt::PublishError::CommandTimeout => Status::GatewayTimeout,
_ => Status::ServiceUnavailable,
})
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn validates_probe_fields() {
let mut input = ProbeInput {
name: "test".into(),
region: None,
asn: None,
provider: None,
hidden: false,
disable_traceroutes: false,
cdn_unblocked: false,
};
assert!(validate(&input).is_ok());
input.name = " ".into();
assert!(validate(&input).is_err());
input.name = "test".into();
input.asn = Some("x".repeat(33));
assert!(validate(&input).is_err());
}
#[test]
fn compares_passwords_without_accepting_prefixes() {
assert!(constant_time_eq(b"secret", b"secret"));
assert!(!constant_time_eq(b"secret", b"secre"));
assert!(!constant_time_eq(b"secret", b"secret2"));
}
}
+2 -1
View File
@@ -140,7 +140,7 @@ pub async fn probe_query(
}
};
let is_ip_target = domain.is_none();
let probe_config = mqtt.probe_config();
let probe_config = mqtt.probe_config().await;
if domain.is_none() && !probe_config.traceroute_enabled {
return Err(Status::BadRequest);
}
@@ -648,6 +648,7 @@ fn publish_error_status(error: PublishError) -> Status {
| PublishError::Serialize(_)
| PublishError::Subscribe(_)
| PublishError::Publish(_) => Status::InternalServerError,
PublishError::CommandTimeout => Status::GatewayTimeout,
}
}
+14
View File
@@ -1,5 +1,6 @@
#[macro_use]
extern crate rocket;
mod admin;
mod agency;
mod api;
mod database_refresh;
@@ -121,6 +122,19 @@ async fn rocket() -> _ {
],
)
.mount("/agency", routes![agency::upload_report])
.mount(
"/api/v1/admin",
routes![
admin::list,
admin::create,
admin::update,
admin::remove,
admin::reload,
admin::update_all_probes,
admin::update_one_probe,
admin::command
],
)
.mount("/mqtt", routes![mqtt_auth::auth, mqtt_auth::acl])
.mount(
"/api/v1/probe-updates",
+181 -37
View File
@@ -1,7 +1,8 @@
use log::{info, warn};
use reports::probe::HostType;
use reports::probe::{
DpiProbeConfig, Host, ProbeConfig, ProbeResult, ProbeResultEvent, ProbeStatus, ProbeTask,
DpiProbeConfig, Host, ProbeCommand, ProbeCommandResult, ProbeConfig, ProbeResult,
ProbeResultEvent, ProbeStatus, ProbeTask,
};
use rocket::serde::json::serde_json;
use rumqttc::{AsyncClient, Event as MqttEvent, Incoming, MqttOptions, QoS};
@@ -16,8 +17,6 @@ use std::time::Duration;
const MQTT_MAX_PACKET_SIZE: usize = 1024 * 1024;
const DEFAULT_PROBE_HOSTS: &str = include_str!("../probe-hosts.toml");
#[derive(Debug)]
pub enum PublishError {
NotConfigured,
@@ -26,6 +25,7 @@ pub enum PublishError {
Serialize(serde_json::Error),
Subscribe(rumqttc::ClientError),
Publish(rumqttc::ClientError),
CommandTimeout,
}
impl fmt::Display for PublishError {
@@ -45,6 +45,7 @@ impl fmt::Display for PublishError {
write!(formatter, "failed to subscribe to results: {error}")
}
PublishError::Publish(error) => write!(formatter, "failed to publish task: {error}"),
PublishError::CommandTimeout => write!(formatter, "probe command timed out"),
}
}
}
@@ -54,7 +55,18 @@ pub struct MqttPublisher {
client: Option<AsyncClient>,
sessions: Arc<rocket::tokio::sync::RwLock<HashMap<String, ProbeResultSender>>>,
probe_statuses: ProbeStatuses,
probe_config: Arc<ProbeConfig>,
probe_config: Arc<rocket::tokio::sync::RwLock<ProbeConfig>>,
command_sessions: Arc<
rocket::tokio::sync::Mutex<
HashMap<
String,
(
String,
rocket::tokio::sync::oneshot::Sender<ProbeCommandResult>,
),
>,
>,
>,
task_timeout_ms: u64,
}
@@ -99,20 +111,23 @@ impl MqttPublisher {
let sessions = Arc::new(rocket::tokio::sync::RwLock::new(HashMap::new()));
let probe_statuses = Arc::new(rocket::tokio::sync::RwLock::new(HashMap::new()));
let task_timeout_ms = task_timeout_ms_from_env();
let probe_config = Arc::new(load_probe_config(task_timeout_ms).unwrap_or_else(|error| {
warn!("failed to load probe config: {error}");
ProbeConfig {
version: env!("CARGO_PKG_VERSION").to_string(),
task_timeout_ms,
published_at: Utc::now().to_rfc3339(),
hosts: Vec::new(),
traceroute_enabled: false,
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 probe_config = Arc::new(rocket::tokio::sync::RwLock::new(
load_probe_config(task_timeout_ms).unwrap_or_else(|error| {
warn!("failed to load probe config: {error}");
ProbeConfig {
version: env!("CARGO_PKG_VERSION").to_string(),
task_timeout_ms,
published_at: Utc::now().to_rfc3339(),
hosts: Vec::new(),
traceroute_enabled: false,
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 command_sessions = Arc::new(rocket::tokio::sync::Mutex::new(HashMap::new()));
let admin_token = match std::env::var("MQTT_ADMIN_TOKEN") {
Ok(token) if !token.is_empty() => token,
_ => {
@@ -122,6 +137,7 @@ impl MqttPublisher {
sessions,
probe_statuses,
probe_config,
command_sessions,
task_timeout_ms: task_timeout_ms_from_env(),
};
}
@@ -145,12 +161,26 @@ impl MqttPublisher {
let event_probe_statuses = probe_statuses.clone();
let config_client = client.clone();
let event_probe_config = probe_config.clone();
let event_command_sessions = command_sessions.clone();
rocket::tokio::spawn(async move {
loop {
match eventloop.poll().await {
Ok(MqttEvent::Incoming(Incoming::ConnAck(_))) => {
if let Err(error) = config_client
.subscribe("probe/status/v1/+", QoS::AtLeastOnce)
.await
{
warn!("failed to subscribe to probe status updates: {error}");
}
if let Err(error) = config_client
.subscribe("probe/command-results/v1/+/+", QoS::AtLeastOnce)
.await
{
warn!("failed to subscribe to command results: {error}");
}
if let Err(error) =
publish_probe_config(&config_client, event_probe_config.as_ref()).await
publish_probe_config(&config_client, &*event_probe_config.read().await)
.await
{
warn!("failed to publish retained probe config: {error}");
}
@@ -164,6 +194,23 @@ impl MqttPublisher {
&publish.payload,
)
.await;
if let Some((probe_id, command_id)) =
parse_command_result_topic(&publish.topic)
{
if let Ok(result) =
serde_json::from_slice::<ProbeCommandResult>(&publish.payload)
{
let mut sessions = event_command_sessions.lock().await;
if sessions
.get(command_id)
.is_some_and(|(expected, _)| expected == probe_id)
{
if let Some((_, sender)) = sessions.remove(command_id) {
let _ = sender.send(result);
}
}
}
}
}
Ok(_) => {}
Err(error) => {
@@ -174,22 +221,13 @@ impl MqttPublisher {
}
});
let status_client = client.clone();
rocket::tokio::spawn(async move {
if let Err(error) = status_client
.subscribe("probe/status/v1/+", QoS::AtLeastOnce)
.await
{
warn!("failed to subscribe to probe status updates: {error}");
}
});
info!("mqtt publisher configured for {host}:{port}");
Self {
client: Some(client),
sessions,
probe_statuses,
probe_config,
command_sessions,
task_timeout_ms,
}
}
@@ -215,8 +253,58 @@ impl MqttPublisher {
self.probe_statuses.read().await.clone()
}
pub fn probe_config(&self) -> Arc<ProbeConfig> {
self.probe_config.clone()
pub async fn probe_config(&self) -> ProbeConfig {
self.probe_config.read().await.clone()
}
pub async fn reload_probe_config(&self) -> Result<ProbeConfig, PublishError> {
let config = load_probe_config(self.task_timeout_ms)?;
let client = self.client.as_ref().ok_or(PublishError::NotConfigured)?;
publish_probe_config(client, &config).await?;
*self.probe_config.write().await = config.clone();
Ok(config)
}
pub async fn request_probe_update(&self, probe_id: Option<&str>) -> Result<(), PublishError> {
let client = self.client.as_ref().ok_or(PublishError::NotConfigured)?;
let topic = probe_update_topic(probe_id);
client
.publish(topic, QoS::AtLeastOnce, false, b"check".as_slice())
.await
.map_err(PublishError::Publish)
}
pub async fn send_command(
&self,
probe_id: &str,
command: ProbeCommand,
) -> Result<ProbeCommandResult, PublishError> {
let client = self.client.as_ref().ok_or(PublishError::NotConfigured)?;
let command_id = uuid::Uuid::new_v4().to_string();
let (sender, receiver) = rocket::tokio::sync::oneshot::channel();
self.command_sessions
.lock()
.await
.insert(command_id.clone(), (probe_id.to_owned(), sender));
let payload = serde_json::to_vec(&command).map_err(PublishError::Serialize)?;
let published = client
.publish(
format!("probe/commands/v1/{probe_id}/{command_id}"),
QoS::AtLeastOnce,
false,
payload,
)
.await;
if let Err(error) = published {
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;
self.command_sessions.lock().await.remove(&command_id);
result
.ok()
.and_then(Result::ok)
.ok_or(PublishError::CommandTimeout)
}
pub async fn subscribe_probe_results(
@@ -296,12 +384,17 @@ async fn publish_probe_config(
}
fn load_probe_config(task_timeout_ms: u64) -> Result<ProbeConfig, PublishError> {
let config = if let Some(path) = std::env::var_os("PROBE_CONFIG_PATH") {
let contents = std::fs::read_to_string(path).map_err(PublishError::Config)?;
parse_probe_hosts(&contents)?
} else {
parse_probe_hosts(DEFAULT_PROBE_HOSTS)?
};
let path = std::env::var_os("PROBE_CONFIG_PATH")
.map(std::path::PathBuf::from)
.unwrap_or_else(|| {
if std::path::Path::new("website/probe-hosts.toml").exists() {
"website/probe-hosts.toml".into()
} else {
"probe-hosts.toml".into()
}
});
let contents = std::fs::read_to_string(path).map_err(PublishError::Config)?;
let config = parse_probe_hosts(&contents)?;
Ok(ProbeConfig {
version: env!("CARGO_PKG_VERSION").to_string(),
task_timeout_ms,
@@ -415,6 +508,13 @@ async fn dispatch_probe_result(
}
}
fn probe_update_topic(probe_id: Option<&str>) -> String {
match probe_id {
Some(id) => format!("probe/update/v1/{id}"),
None => "probe/update/v1".to_string(),
}
}
fn parse_probe_status_topic(topic: &str) -> Option<&str> {
let mut parts = topic.split('/');
match (
@@ -429,6 +529,28 @@ fn parse_probe_status_topic(topic: &str) -> Option<&str> {
}
}
fn parse_command_result_topic(topic: &str) -> Option<(&str, &str)> {
let mut parts = topic.split('/');
match (
parts.next(),
parts.next(),
parts.next(),
parts.next(),
parts.next(),
parts.next(),
) {
(
Some("probe"),
Some("command-results"),
Some("v1"),
Some(probe_id),
Some(command_id),
None,
) if !probe_id.is_empty() && !command_id.is_empty() => Some((probe_id, command_id)),
_ => None,
}
}
fn parse_probe_result_topic(topic: &str) -> Option<(&str, &str)> {
let mut parts = topic.split('/');
match (
@@ -457,6 +579,28 @@ fn task_timeout_ms_from_env() -> u64 {
mod tests {
use super::*;
#[test]
fn update_check_uses_global_or_targeted_topic() {
assert_eq!(probe_update_topic(None), "probe/update/v1");
assert_eq!(probe_update_topic(Some("42")), "probe/update/v1/42");
}
#[test]
fn command_result_topics_require_exact_probe_and_command() {
assert_eq!(
parse_command_result_topic("probe/command-results/v1/42/abc"),
Some(("42", "abc"))
);
assert_eq!(
parse_command_result_topic("probe/command-results/v1/42/abc/extra"),
None
);
assert_eq!(
parse_command_result_topic("probe/command-results/v1/42/"),
None
);
}
#[rocket::async_test]
async fn status_snapshot_keeps_offline_node_metadata() {
let statuses = Arc::new(rocket::tokio::sync::RwLock::new(HashMap::new()));
+14
View File
@@ -109,6 +109,20 @@ pub async fn acl(
if request.username != "probe" || request.clientid.is_empty() {
return Json(MqttAuthResponse::deny());
}
let Ok(reporter_id) = request.clientid.parse::<i32>() else {
return Json(MqttAuthResponse::deny());
};
let active =
sqlx::query_scalar::<_, bool>("SELECT EXISTS(SELECT 1 FROM reporters WHERE id = $1)")
.bind(reporter_id)
.fetch_optional(&**pool)
.await
.ok()
.flatten()
.unwrap_or(false);
if !active {
return Json(MqttAuthResponse::deny());
}
match request.access {
1 if can_probe_subscribe(request.clientid, request.topic, pool).await => {