From ea2ecb44d2dbe462a9427b48e86c64d6ce56ba56 Mon Sep 17 00:00:00 2001 From: Denozordec Date: Sun, 12 Apr 2026 13:07:14 +0700 Subject: [PATCH] Implement radar-telemt-dcs aggregation endpoint and UI integration - 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. --- README.md | 2 +- docs/AGGREGATE.md | 2 + internal/aggregate/handlers.go | 21 ++++ internal/aggregate/handlers_test.go | 59 ++++++++++ internal/aggregate/radar_telemt.go | 44 ++++++++ internal/aggregate/types.go | 34 ++++++ web/src/lib/api/client.ts | 44 ++++++++ web/src/routes/radar/+page.svelte | 160 ++++++++++++++++++++++++++-- 8 files changed, 357 insertions(+), 9 deletions(-) create mode 100644 internal/aggregate/radar_telemt.go diff --git a/README.md b/README.md index 4f8c8f2..4aea746 100644 --- a/README.md +++ b/README.md @@ -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) diff --git a/docs/AGGREGATE.md b/docs/AGGREGATE.md index 61437b7..8202181 100644 --- a/docs/AGGREGATE.md +++ b/docs/AGGREGATE.md @@ -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 шлюза. diff --git a/internal/aggregate/handlers.go b/internal/aggregate/handlers.go index 531960f..b18d0a4 100644 --- a/internal/aggregate/handlers.go +++ b/internal/aggregate/handlers.go @@ -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 { diff --git a/internal/aggregate/handlers_test.go b/internal/aggregate/handlers_test.go index fcab56c..3bcc962 100644 --- a/internal/aggregate/handlers_test.go +++ b/internal/aggregate/handlers_test.go @@ -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" { diff --git a/internal/aggregate/radar_telemt.go b/internal/aggregate/radar_telemt.go new file mode 100644 index 0000000..5260041 --- /dev/null +++ b/internal/aggregate/radar_telemt.go @@ -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} +} diff --git a/internal/aggregate/types.go b/internal/aggregate/types.go index 1be7b23..8fe1f19 100644 --- a/internal/aggregate/types.go +++ b/internal/aggregate/types.go @@ -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"` +} diff --git a/web/src/lib/api/client.ts b/web/src/lib/api/client.ts index 699f049..65380eb 100644 --- a/web/src/lib/api/client.ts +++ b/web/src/lib/api/client.ts @@ -180,6 +180,50 @@ export async function fetchAggFleetStatus(params?: { aliases?: string }): Promis return body as AggEnvelope; } +/** 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> { + 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 | 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; +} + export async function fetchAggUniqueIps(params?: { aliases?: string; geo?: boolean; diff --git a/web/src/routes/radar/+page.svelte b/web/src/routes/radar/+page.svelte index b5755bc..e0b4e72 100644 --- a/web/src/routes/radar/+page.svelte +++ b/web/src/routes/radar/+page.svelte @@ -1,6 +1,12 @@ -
+
@@ -209,8 +255,15 @@

-
@@ -265,6 +318,96 @@ + + + Telemt ME — снимок по нодам + + GET /api/agg/radar-telemt-dcs — параллельно + GET /v1/stats/dcs на каждом upstream (как на странице «Состояние» ноды). + Требуется включённый minimal runtime API на Telemt; иначе в ячейках будет причина отключения. + {#if telemtUpdated} + · обновлено {telemtUpdated} + {/if} + + + + {#if telemtErr} + + Ошибка загрузки stats/dcs + {telemtErr} + + {:else if telemtLoading && telemtServers.length === 0} +

Загрузка…

+ {:else if telemtServers.length === 0} +

Нет серверов в конфигурации шлюза.

+ {:else} + {#if telemtPartial} + + + Частичные данные + + Не все ноды ответили успешно — смотрите ошибки в строках. + + + {/if} +
+ + + + Нода + {#each telemtDcNums as dc (dc)} + DC{dc} + {/each} + + + + {#each telemtServers as srv (srv.alias)} + + {srv.alias} + {#each telemtDcNums as dc (dc)} + + {#if !srv.ok} + ошибка + {:else if !srv.data?.middle_proxy_enabled} + + {srv.data?.reason === 'feature_disabled' + ? 'ME API off' + : (srv.data?.reason ?? 'нет данных')} + + {:else} + {@const cell = dcRowFind(srv, dc)} + {#if cell} +
+ {cell.coverage_pct.toFixed(0)}% +
+ {#if cell.rtt_ms != null && cell.rtt_ms !== undefined} +
+ {Number(cell.rtt_ms).toFixed(0)} ms +
+ {/if} + {:else} + + {/if} + {/if} +
+ {/each} +
+ {/each} +
+
+
+

+ Покрытие — доля alive writers к required для DC (Telemt). Это не сырой TCP-пинг с интернета, а + состояние middle proxy на процессе Telemt. +

+ {/if} +
+
+ @@ -272,9 +415,10 @@ Диагностика TCP с хоста шлюза - Запрос /api/radar/ping-dc — TCP :443 до каждого DC с таймаутом 2 с - (как ping_proxy.php). Поле «from» в ответе: источник метки на - сервере. + Запрос /api/radar/ping-dc выполняется на процессе шлюза + (не на каждой ноде Telemt): TCP :443 до каждого DC, таймаут 2 с (как + ping_proxy.php). Для вида «с каждой ноды» используйте таблицу ME выше. + Поле «from» в ответе — метка источника. {#if pingFrom} from: {pingFrom} {/if}