Implement radar-telemt-dcs aggregation endpoint and UI integration
Publish telemt-api gateway Docker image / test (push) Successful in 9s
Publish telemt-api gateway Docker image / build-and-push (push) Successful in 1m53s

- Added new API route `/api/agg/radar-telemt-dcs` to aggregate DC status data from multiple upstreams, including metrics like coverage percentage and RTT.
- Implemented handler logic in `handlers.go` and corresponding tests in `handlers_test.go` to ensure correct data retrieval and response formatting.
- Updated the frontend to fetch and display radar DC data, enhancing the user interface with a new section for Telemt ME snapshots.
- Enhanced documentation in `AGGREGATE.md` and `README.md` to reflect the new functionality and usage details.
This commit is contained in:
Denozordec
2026-04-12 13:07:14 +07:00
parent d021e4b1d7
commit ea2ecb44d2
8 changed files with 357 additions and 9 deletions
+1 -1
View File
@@ -1,6 +1,6 @@
# telemt-api
HTTP‑шлюз на Go для [Telemt Control API](docs/API.md): один порт, **белый список IP (CIDR)**, маршруты вида `/api/{alias}/…``{base_url}/v1/…`, опционально **Mihomo**`/api/{alias}/mihomo/…` к external-controller (см. [docs/GATEWAY_RUN.md](docs/GATEWAY_RUN.md#mihomo-external-controller)), агрегация нескольких инстансов — [`/api/agg/…`](docs/AGGREGATE.md), live SSE поток — `/api/live/events`, **радар DC Telegram**`GET /api/radar/statuses` и `GET /api/radar/ping-dc` (см. [docs/GATEWAY_RUN.md](docs/GATEWAY_RUN.md#radar-dc-telegram)), метрики Prometheus на `/metrics`. **Web UI** (SvelteKit) встроен в тот же процесс/образ: статика на `/`, API на `/api/…` и `/health`, раздел **Mihomo** на `/servers/{alias}/mihomo`, **Радар DC** на `/radar`.
HTTP‑шлюз на Go для [Telemt Control API](docs/API.md): один порт, **белый список IP (CIDR)**, маршруты вида `/api/{alias}/…``{base_url}/v1/…`, опционально **Mihomo**`/api/{alias}/mihomo/…` к external-controller (см. [docs/GATEWAY_RUN.md](docs/GATEWAY_RUN.md#mihomo-external-controller)), агрегация нескольких инстансов — [`/api/agg/…`](docs/AGGREGATE.md) (в т.ч. `GET /api/agg/radar-telemt-dcs``stats/dcs` по всем нодам для радара), live SSE поток — `/api/live/events`, **радар DC Telegram**`GET /api/radar/statuses` и `GET /api/radar/ping-dc` (см. [docs/GATEWAY_RUN.md](docs/GATEWAY_RUN.md#radar-dc-telegram)), метрики Prometheus на `/metrics`. **Web UI** (SvelteKit) встроен в тот же процесс/образ: статика на `/`, API на `/api/…` и `/health`, раздел **Mihomo** на `/servers/{alias}/mihomo`, **Радар DC** на `/radar`.
## Быстрый старт (Linux)
+2
View File
@@ -4,6 +4,7 @@
- Большинство маршрутов агрегации используют **`GET /v1/stats/users`** на каждом сервере из конфигурации.
- **`GET /api/agg/fleet-status`** дополнительно вызывает на каждом upstream **`GET /v1/health`** и **`GET /v1/system/info`** (параллельно по серверам).
- **`GET /api/agg/radar-telemt-dcs`** — на каждом upstream параллельно **`GET /v1/stats/dcs`** (снимок ME / DC для панели «Радар DC»); см. [API.md](API.md) про `minimal_runtime_enabled` и поля `DcStatusData`.
**Единицы трафика в агрегатах:** поля `*_megabytes` — это **двоичные мегабайты (MiB)**, 1 MiB = 1024² октетов (как у Telemt в ответе считаются октеты, шлюз делит на MiB для удобства).
@@ -33,6 +34,7 @@
| GET | `/api/agg/users` | Объединённый список пользователей с `by_server`, суммарным `total_megabytes` и **смерженными лимитами** (см. ниже). |
| GET | `/api/agg/user/{username}` | Один пользователь в том же формате, что элементы `/api/agg/users` (без списка всех). Имя в пути: `[A-Za-z0-9_.-]+`. Ответ **`404`**, если пользователь не найден ни на одном успешном upstream. |
| GET | `/api/agg/fleet-status` | По каждому алиасу: параллельно health + system/info; в `data.servers[]` — статусы подзапросов и тела `health` / `system_info` при успехе. См. [AGGREGATE_OPENAPI.yaml](AGGREGATE_OPENAPI.yaml). |
| GET | `/api/agg/radar-telemt-dcs` | По каждому алиасу: параллельно `GET /v1/stats/dcs`; в `data.servers[]``alias`, `ok`, при успехе объект `data` (поля `middle_proxy_enabled`, `reason`, `dcs[]` с `dc`, `coverage_pct`, `rtt_ms` и т.д.). |
| GET | `/api/agg/incidents` | Нормализованный snapshot инцидентов для triage-панели: `critical/warning/info`, `affected_aliases`, рекомендуемые `actions` (runbook/deep links), счётчики по severity. |
Все методы — **GET**; действует тот же whitelist, что и для остального API шлюза.
+21
View File
@@ -113,6 +113,8 @@ func (h *Handler) dispatch(w http.ResponseWriter, r *http.Request, sub string) {
h.handleUsers(w, r)
case sub == "fleet-status":
h.handleFleetStatus(w, r)
case sub == "radar-telemt-dcs":
h.handleRadarTelemtDcs(w, r)
case sub == "incidents":
h.handleIncidents(w, r)
case strings.HasPrefix(sub, "user/"):
@@ -275,6 +277,25 @@ func (h *Handler) handleUsers(w http.ResponseWriter, r *http.Request) {
writeAggOK(w, partial, data)
}
func (h *Handler) handleRadarTelemtDcs(w http.ResponseWriter, r *http.Request) {
aliases, err := h.resolveAliases(r)
if err != nil {
writeBadRequest(w, err)
return
}
ctx, cancel := context.WithTimeout(r.Context(), 60*time.Second)
defer cancel()
data := FetchRadarTelemtDcs(ctx, h.Client, h.Parsed, aliases)
partial := false
for _, s := range data.Servers {
if !s.OK {
partial = true
break
}
}
writeAggOK(w, partial, data)
}
func (h *Handler) handleFleetStatus(w http.ResponseWriter, r *http.Request) {
aliases, err := h.resolveAliases(r)
if err != nil {
+59
View File
@@ -113,6 +113,65 @@ func TestHandlerFleetStatus(t *testing.T) {
}
}
func TestHandlerRadarTelemtDcs(t *testing.T) {
up := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path != "/v1/stats/dcs" {
http.NotFound(w, r)
return
}
_ = json.NewEncoder(w).Encode(map[string]any{
"ok": true,
"data": map[string]any{
"middle_proxy_enabled": true,
"generated_at_epoch_secs": 1,
"dcs": []map[string]any{{
"dc": 1, "rtt_ms": 42.5, "coverage_pct": 100.0,
"alive_writers": 3, "required_writers": 3, "load": 0,
}},
},
"revision": "rd",
})
}))
defer up.Close()
cfg := &config.Config{
Servers: []config.Server{
{Alias: "test", BaseURL: up.URL, PathPrefix: "/v1"},
},
}
if err := cfg.Validate(); err != nil {
t.Fatal(err)
}
parsed, err := cfg.Parse()
if err != nil {
t.Fatal(err)
}
h := NewHandler(parsed, up.Client(), nil, 0)
req := httptest.NewRequest(http.MethodGet, "/api/agg/radar-telemt-dcs?aliases=test", nil)
rec := httptest.NewRecorder()
h.ServeHTTP(rec, req)
if rec.Code != http.StatusOK {
t.Fatalf("status %d body %s", rec.Code, rec.Body.String())
}
var env struct {
OK bool `json:"ok"`
Data RadarTelemtDcsData `json:"data"`
}
if err := json.Unmarshal(rec.Body.Bytes(), &env); err != nil {
t.Fatal(err)
}
if !env.OK || len(env.Data.Servers) != 1 {
t.Fatalf("envelope: %+v", env)
}
row := env.Data.Servers[0]
if row.Alias != "test" || !row.OK || row.Data == nil || !row.Data.MiddleProxyEnabled {
t.Fatalf("row: %+v", row)
}
if len(row.Data.Dcs) != 1 || row.Data.Dcs[0].DC != 1 {
t.Fatalf("dcs: %+v", row.Data.Dcs)
}
}
func TestHandlerUserOne(t *testing.T) {
up := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path != "/v1/stats/users" {
+44
View File
@@ -0,0 +1,44 @@
package aggregate
import (
"context"
"net/http"
"sort"
"sync"
"github.com/telemt/telemt-api/internal/config"
)
const statsDcsPath = "stats/dcs"
// FetchRadarTelemtDcs calls GET /v1/stats/dcs on each Telemt in parallel.
func FetchRadarTelemtDcs(ctx context.Context, client *http.Client, parsed *config.Parsed, aliases []string) RadarTelemtDcsData {
if len(aliases) == 0 {
return RadarTelemtDcsData{}
}
rows := make([]RadarTelemtDcsServer, len(aliases))
var wg sync.WaitGroup
for i, alias := range aliases {
i, alias := i, alias
wg.Add(1)
go func() {
defer wg.Done()
data, meta := FetchTelemtGET[DcStatusPayload](ctx, client, parsed, alias, statsDcsPath)
row := RadarTelemtDcsServer{
Alias: alias,
OK: meta.OK,
HTTPStatus: meta.HTTPStatus,
LatencyMs: meta.LatencyMs,
Error: meta.Error,
Revision: meta.Revision,
}
if meta.OK {
row.Data = &data
}
rows[i] = row
}()
}
wg.Wait()
sort.Slice(rows, func(i, j int) bool { return rows[i].Alias < rows[j].Alias })
return RadarTelemtDcsData{Servers: rows}
}
+34
View File
@@ -206,3 +206,37 @@ type TopUserByUniqueIPs struct {
Username string `json:"username"`
UniqueIPs uint64 `json:"unique_ips"`
}
// DcStatusRow mirrors Telemt GET /v1/stats/dcs data.dcs[] (subset for radar UI).
type DcStatusRow struct {
DC int `json:"dc"`
RttMs *float64 `json:"rtt_ms"`
CoveragePct float64 `json:"coverage_pct"`
AliveWriters int `json:"alive_writers"`
RequiredWriters int `json:"required_writers"`
Load int `json:"load"`
}
// DcStatusPayload mirrors Telemt GET /v1/stats/dcs data object.
type DcStatusPayload struct {
MiddleProxyEnabled bool `json:"middle_proxy_enabled"`
Reason *string `json:"reason"`
GeneratedAtEpochSecs uint64 `json:"generated_at_epoch_secs"`
Dcs []DcStatusRow `json:"dcs"`
}
// RadarTelemtDcsServer is one upstream stats/dcs outcome for /api/agg/radar-telemt-dcs.
type RadarTelemtDcsServer struct {
Alias string `json:"alias"`
OK bool `json:"ok"`
HTTPStatus int `json:"http_status,omitempty"`
LatencyMs int64 `json:"latency_ms,omitempty"`
Error string `json:"error,omitempty"`
Revision string `json:"revision,omitempty"`
Data *DcStatusPayload `json:"data,omitempty"`
}
// RadarTelemtDcsData is aggregate payload for /api/agg/radar-telemt-dcs.
type RadarTelemtDcsData struct {
Servers []RadarTelemtDcsServer `json:"servers"`
}
+44
View File
@@ -180,6 +180,50 @@ export async function fetchAggFleetStatus(params?: { aliases?: string }): Promis
return body as AggEnvelope<components['schemas']['FleetStatusData']>;
}
/** Payload GET /api/agg/radar-telemt-dcs (снимок GET /v1/stats/dcs с каждой ноды). */
export type AggRadarTelemtDcsRow = {
alias: string;
ok: boolean;
http_status?: number;
latency_ms?: number;
error?: string;
revision?: string;
data?: {
middle_proxy_enabled: boolean;
reason?: string | null;
generated_at_epoch_secs?: number;
dcs?: {
dc: number;
rtt_ms?: number | null;
coverage_pct: number;
alive_writers: number;
required_writers: number;
load: number;
}[];
};
};
export type AggRadarTelemtDcsData = {
servers: AggRadarTelemtDcsRow[];
};
export async function fetchAggRadarTelemtDcs(params?: {
aliases?: string;
}): Promise<AggEnvelope<AggRadarTelemtDcsData>> {
const q = new URLSearchParams();
if (params?.aliases) q.set('aliases', params.aliases);
const url = `${gatewayBase()}/api/agg/radar-telemt-dcs${q.toString() ? `?${q}` : ''}`;
const res = await fetch(url);
const body = (await parseJson(res)) as Record<string, unknown> | null;
if (!res.ok) {
throw new ApiError(`radar-telemt-dcs HTTP ${res.status}`, res.status, body);
}
if (!body || body.ok !== true) {
throw new ApiError('radar-telemt-dcs: ok !== true', res.status, body);
}
return body as AggEnvelope<AggRadarTelemtDcsData>;
}
export async function fetchAggUniqueIps(params?: {
aliases?: string;
geo?: boolean;
+152 -8
View File
@@ -1,6 +1,12 @@
<script lang="ts">
import { onMount } from 'svelte';
import { fetchRadarPingDC, fetchRadarStatuses, type RadarPingRow } from '$lib/api/client.js';
import {
fetchAggRadarTelemtDcs,
fetchRadarPingDC,
fetchRadarStatuses,
type AggRadarTelemtDcsRow,
type RadarPingRow
} from '$lib/api/client.js';
import * as Card from '$lib/components/ui/card/index.js';
import { Button } from '$lib/components/ui/button/index.js';
import * as Table from '$lib/components/ui/table/index.js';
@@ -50,6 +56,25 @@
});
let pingFrom = $state<string | null>(null);
let telemtLoading = $state(false);
let telemtErr = $state<string | null>(null);
let telemtPartial = $state(false);
let telemtUpdated = $state<string | null>(null);
let telemtServers = $state<AggRadarTelemtDcsRow[]>([]);
const telemtDcNums = [1, 2, 3, 4, 5] as const;
function dcRowFind(row: AggRadarTelemtDcsRow, dc: number) {
return row.data?.dcs?.find((d) => d.dc === dc);
}
function telemtCoverageClass(cov: number | undefined): string {
if (cov === undefined) return 'text-muted-foreground';
if (cov >= 90) return 'text-emerald-600 font-medium';
if (cov >= 50) return 'text-amber-600 font-medium';
return 'text-red-600 font-medium';
}
function buildMatrix(data: unknown[]) {
const regions: Record<string, { ok: number; total: number }> = {
SE: { ok: 0, total: 0 },
@@ -190,12 +215,33 @@
return { text: `${ms} мс`, barPct: Math.min(ms / 3, 100), color };
}
async function loadTelemtDcs() {
telemtLoading = true;
telemtErr = null;
try {
const env = await fetchAggRadarTelemtDcs();
telemtServers = env.data.servers ?? [];
telemtPartial = !!env.partial;
telemtUpdated = new Date().toLocaleTimeString('ru-RU', { hour: '2-digit', minute: '2-digit' });
} catch (e) {
telemtErr = e instanceof Error ? e.message : String(e);
telemtServers = [];
telemtPartial = false;
} finally {
telemtLoading = false;
}
}
async function refreshRadarAndTelemt() {
await Promise.all([loadRadar(), loadTelemtDcs()]);
}
onMount(() => {
void loadRadar();
void refreshRadarAndTelemt();
});
</script>
<div class="mx-auto flex max-w-5xl flex-col gap-6">
<div class="flex w-full flex-col gap-6">
<div class="flex flex-wrap items-center justify-between gap-3">
<div class="flex items-center gap-2">
<RadarIcon class="text-primary size-8" />
@@ -209,8 +255,15 @@
</p>
</div>
</div>
<Button variant="outline" size="sm" disabled={radarLoading} onclick={() => loadRadar()}>
<RefreshCwIcon class="mr-1 size-4 {radarLoading ? 'animate-spin' : ''}" />
<Button
variant="outline"
size="sm"
disabled={radarLoading || telemtLoading}
onclick={() => void refreshRadarAndTelemt()}
>
<RefreshCwIcon
class="mr-1 size-4 {radarLoading || telemtLoading ? 'animate-spin' : ''}"
/>
Обновить
</Button>
</div>
@@ -265,6 +318,96 @@
</Card.Content>
</Card.Root>
<Card.Root>
<Card.Header>
<Card.Title class="text-base">Telemt ME — снимок по нодам</Card.Title>
<Card.Description>
<code class="text-xs">GET /api/agg/radar-telemt-dcs</code> — параллельно
<code class="text-xs">GET /v1/stats/dcs</code> на каждом upstream (как на странице «Состояние» ноды).
Требуется включённый minimal runtime API на Telemt; иначе в ячейках будет причина отключения.
{#if telemtUpdated}
<span class="text-foreground"> · обновлено {telemtUpdated}</span>
{/if}
</Card.Description>
</Card.Header>
<Card.Content>
{#if telemtErr}
<Alert variant="destructive">
<AlertTitle>Ошибка загрузки stats/dcs</AlertTitle>
<AlertDescription>{telemtErr}</AlertDescription>
</Alert>
{:else if telemtLoading && telemtServers.length === 0}
<p class="text-muted-foreground text-sm">Загрузка…</p>
{:else if telemtServers.length === 0}
<p class="text-muted-foreground text-sm">Нет серверов в конфигурации шлюза.</p>
{:else}
{#if telemtPartial}
<Alert class="mb-4">
<InfoIcon class="size-4" />
<AlertTitle class="text-sm">Частичные данные</AlertTitle>
<AlertDescription class="text-xs">
Не все ноды ответили успешно — смотрите ошибки в строках.
</AlertDescription>
</Alert>
{/if}
<div class="overflow-x-auto rounded-md border border-border">
<Table.Root>
<Table.Header>
<Table.Row>
<Table.Head class="w-28">Нода</Table.Head>
{#each telemtDcNums as dc (dc)}
<Table.Head class="min-w-[5.5rem] text-center">DC{dc}</Table.Head>
{/each}
</Table.Row>
</Table.Header>
<Table.Body>
{#each telemtServers as srv (srv.alias)}
<Table.Row>
<Table.Cell class="font-mono text-sm font-medium">{srv.alias}</Table.Cell>
{#each telemtDcNums as dc (dc)}
<Table.Cell class="align-top text-center text-sm">
{#if !srv.ok}
<span class="text-destructive text-xs leading-tight" title={srv.error ?? ''}
>ошибка</span>
{:else if !srv.data?.middle_proxy_enabled}
<span
class="text-muted-foreground text-xs leading-tight"
title={srv.data?.reason ?? ''}
>
{srv.data?.reason === 'feature_disabled'
? 'ME API off'
: (srv.data?.reason ?? 'нет данных')}
</span>
{:else}
{@const cell = dcRowFind(srv, dc)}
{#if cell}
<div class={telemtCoverageClass(cell.coverage_pct)}>
{cell.coverage_pct.toFixed(0)}%
</div>
{#if cell.rtt_ms != null && cell.rtt_ms !== undefined}
<div class="text-muted-foreground mt-0.5 text-xs tabular-nums">
{Number(cell.rtt_ms).toFixed(0)} ms
</div>
{/if}
{:else}
<span class="text-muted-foreground"></span>
{/if}
{/if}
</Table.Cell>
{/each}
</Table.Row>
{/each}
</Table.Body>
</Table.Root>
</div>
<p class="text-muted-foreground mt-3 text-xs">
Покрытие — доля alive writers к required для DC (Telemt). Это не сырой TCP-пинг с интернета, а
состояние middle proxy на процессе Telemt.
</p>
{/if}
</Card.Content>
</Card.Root>
<Card.Root>
<Card.Header>
<Card.Title class="flex items-center gap-2 text-base">
@@ -272,9 +415,10 @@
Диагностика TCP с хоста шлюза
</Card.Title>
<Card.Description>
Запрос <code class="text-xs">/api/radar/ping-dc</code> — TCP :443 до каждого DC с таймаутом 2 с
(как <code class="text-xs">ping_proxy.php</code>). Поле «from» в ответе: источник метки на
сервере.
Запрос <code class="text-xs">/api/radar/ping-dc</code> выполняется на <strong>процессе шлюза</strong>
(не на каждой ноде Telemt): TCP :443 до каждого DC, таймаут 2 с (как
<code class="text-xs">ping_proxy.php</code>). Для вида «с каждой ноды» используйте таблицу ME выше.
Поле «from» в ответе — метка источника.
{#if pingFrom}
<span class="mt-1 block text-foreground">from: {pingFrom}</span>
{/if}