From 7220fe5b8b5727aaaa5d5ad3a378f20f131d9a1a Mon Sep 17 00:00:00 2001 From: Naeel Date: Mon, 6 Apr 2026 19:14:57 +0300 Subject: [PATCH] =?UTF-8?q?feat:=20IoT=20Admin=20Stats=20page=20=E2=80=94?= =?UTF-8?q?=20/iot-admin=20(v0.1.70)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - GET /iot-admin — HTML страница администратора (go:embed) - GET /iot-admin/stats — JSON API с данными (Bearer ADMIN_STATS_TOKEN) - Источники: Kafka consumer lag, K8s pod statuses, PostgreSQL per-tenant stats - Авторизация: ADMIN_STATS_TOKEN env var - Auto-refresh каждые 30 секунд - Nubes brand style --- deployments/k8s/operator.yaml | 11 +- internal/api/admin_embed.go | 27 + internal/api/handler/handler.go | 5 +- .../api/handler/iot_admin_stats_handler.go | 214 +++++++ internal/api/router.go | 5 + internal/api/ui/iot-admin.html | 547 ++++++++++++++++++ internal/storage/iotpg/iot_telemetry_store.go | 93 +++ main.go | 13 +- 8 files changed, 906 insertions(+), 9 deletions(-) create mode 100644 internal/api/admin_embed.go create mode 100644 internal/api/handler/iot_admin_stats_handler.go create mode 100644 internal/api/ui/iot-admin.html diff --git a/deployments/k8s/operator.yaml b/deployments/k8s/operator.yaml index 0156f45..ee664d3 100644 --- a/deployments/k8s/operator.yaml +++ b/deployments/k8s/operator.yaml @@ -1,4 +1,4 @@ -# Изменено: 2026-04-05 (добавлен IOT_PG_DSN, версия v0.1.59) +# Изменено: 2026-04-06 (добавлены KAFKA_BROKERS, ADMIN_STATS_TOKEN, версия v0.1.70) # Деплой sless оператора в кластер. # Состав: # - ConfigMap: не-секретные env vars (S3_ENDPOINT, REGISTRY_HOST и т.д.) @@ -33,6 +33,8 @@ data: # EXTERNAL_URL — если задан, URL функции = EXTERNAL_URL/fn/{namespace}/{name} # Позволяет обойтись без wildcard DNS *.fn.kube5s.ru EXTERNAL_URL: "https://sless.kube5s.ru" + # KAFKA_BROKERS — адрес Kafka для чтения consumer lag на странице администратора + KAFKA_BROKERS: "kafka.sless.svc.cluster.local:9092" --- # Secret создаётся отдельно через kubectl (не коммитить секреты в git!) # Описание ключей: @@ -75,7 +77,7 @@ spec: - name: operator # При обновлении версии оператора — менять тег здесь (не latest!) # v0.1.59 — добавлено сохранение телеметрии в IoT Postgres (per-tenant DB) - image: pearlharbor.registryk8s.services.ngcloud.ru/naeel/sless-operator:v0.1.59 + image: pearlharbor.registryk8s.services.ngcloud.ru/naeel/sless-operator:v0.1.70 # Always — чтобы всегда тянуть по точному тегу (не кешировать старый) imagePullPolicy: Always ports: @@ -94,6 +96,11 @@ spec: - secretRef: name: iot-postgres-secret optional: true + env: + # ADMIN_STATS_TOKEN — токен доступа к /iot-admin/stats (страница администратора). + # Менять на уникальный: kubectl set env deploy/sless-operator ADMIN_STATS_TOKEN= -n sless + - name: ADMIN_STATS_TOKEN + value: "iot-admin-sless-2026" readinessProbe: httpGet: path: /healthz diff --git a/internal/api/admin_embed.go b/internal/api/admin_embed.go new file mode 100644 index 0000000..892f3ec --- /dev/null +++ b/internal/api/admin_embed.go @@ -0,0 +1,27 @@ +// Создано: 2026-04-06 +// admin_embed.go — встраивает HTML страницы администратора IoT в бинарник через go:embed. +// +// Страница /iot-admin доступна без JWT — данные не содержит. +// Все данные загружаются через /iot-admin/stats (защищён ADMIN_STATS_TOKEN). +// Почему go:embed: единый деплой, нет отдельных pod-ов, нет nginx drift. + +package api + +import ( + _ "embed" + "net/http" +) + +// iotAdminHTML — бинарное содержимое страницы администратора IoT, встроенное при сборке. +// +//go:embed ui/iot-admin.html +var iotAdminHTML []byte + +// ServeIoTAdmin обрабатывает GET /iot-admin — отдаёт HTML страницу администратора. +// Auth не нужен для HTML — сама страница ничего не содержит, только UI оболочка. +func ServeIoTAdmin(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "text/html; charset=utf-8") + w.Header().Set("Cache-Control", "no-cache, must-revalidate") + w.WriteHeader(http.StatusOK) + _, _ = w.Write(iotAdminHTML) +} diff --git a/internal/api/handler/handler.go b/internal/api/handler/handler.go index dde0f62..3c7041a 100644 --- a/internal/api/handler/handler.go +++ b/internal/api/handler/handler.go @@ -46,7 +46,10 @@ type Handler struct { PG *postgres.Store // IoTPG — хранилище IoT телеметрии (per-tenant Postgres). nil если IOT_PG_DSN не задан. IoTPG *iotpg.IoTPostgresStore - Log *slog.Logger + // KafkaBrokers — адреса Kafka брокеров (KAFKA_BROKERS env var). + // Используется страницей администратора для чтения consumer lag. + KafkaBrokers string + Log *slog.Logger } // writeJSON отправляет JSON-ответ с указанным статусом. diff --git a/internal/api/handler/iot_admin_stats_handler.go b/internal/api/handler/iot_admin_stats_handler.go new file mode 100644 index 0000000..15072ad --- /dev/null +++ b/internal/api/handler/iot_admin_stats_handler.go @@ -0,0 +1,214 @@ +// Создано: 2026-04-06 +// iot_admin_stats_handler.go — handler для страницы администратора IoT. +// +// Endpoints: +// GET /iot-admin/stats — JSON с агрегированной статистикой (защищён ADMIN_STATS_TOKEN) +// +// Источники данных: +// - PostgreSQL (IoTPG): counts per tenant, last 1h/24h, latest rows +// - Kafka: consumer lag (latest offset - committed offset для group iot-pg-consumer) +// - K8s: статус подов iot-mqtt-bridge и iot-kafka-consumer +// +// Авторизация: Bearer из env ADMIN_STATS_TOKEN. +// Если ADMIN_STATS_TOKEN не задан — endpoint возвращает 503. + +package handler + +import ( + "context" + "fmt" + "net/http" + "os" + "strings" + "time" + + kafka "github.com/segmentio/kafka-go" + corev1 "k8s.io/api/core/v1" + "sigs.k8s.io/controller-runtime/pkg/client" +) + +// iotAdminPodStatus — краткая информация о k8s pod для страницы администратора. +type iotAdminPodStatus struct { + Name string `json:"name"` + Phase string `json:"phase"` + Ready bool `json:"ready"` + Restarts int32 `json:"restarts"` + Age string `json:"age"` +} + +// iotAdminKafkaStats — информация о Kafka топике и consumer lag. +type iotAdminKafkaStats struct { + LatestOffset int64 `json:"latest_offset"` + CommittedOffset int64 `json:"committed_offset"` + ConsumerLag int64 `json:"consumer_lag"` + Error string `json:"error,omitempty"` +} + +// AdminStats обрабатывает GET /iot-admin/stats. +// Проверяет Bearer-токен из ADMIN_STATS_TOKEN, затем собирает и возвращает статистику. +func (h *Handler) AdminStats(w http.ResponseWriter, r *http.Request) { + adminToken := os.Getenv("ADMIN_STATS_TOKEN") + if adminToken == "" { + writeJSON(w, http.StatusServiceUnavailable, errResp("admin stats not configured: ADMIN_STATS_TOKEN not set")) + return + } + + authHeader := r.Header.Get("Authorization") + if !strings.HasPrefix(authHeader, "Bearer ") || strings.TrimPrefix(authHeader, "Bearer ") != adminToken { + writeJSON(w, http.StatusUnauthorized, errResp("unauthorized")) + return + } + + ctx, cancel := context.WithTimeout(r.Context(), 15*time.Second) + defer cancel() + + result := map[string]any{ + "collected_at": time.Now().UTC(), + } + + // PostgreSQL: статистика по всем tenant + if h.IoTPG != nil { + pgStats, err := h.IoTPG.GetAdminStats(ctx) + if err != nil { + result["postgres"] = map[string]any{"reachable": false, "error": err.Error()} + } else { + result["postgres"] = pgStats + } + } else { + result["postgres"] = map[string]any{"reachable": false, "error": "IoTPG not configured"} + } + + // Kafka: consumer lag для топика iot.telemetry / группы iot-pg-consumer + result["kafka"] = h.collectIotKafkaLag(ctx) + + // K8s: статус подов bridge и consumer + result["pods"] = h.collectIotPodStatuses(ctx) + + writeJSON(w, http.StatusOK, result) +} + +// collectIotKafkaLag получает latest offset топика и committed offset consumer group, +// вычисляет lag = latest - committed. +// Topic: "iot.telemetry", Consumer Group: "iot-pg-consumer". +func (h *Handler) collectIotKafkaLag(ctx context.Context) iotAdminKafkaStats { + if h.KafkaBrokers == "" { + return iotAdminKafkaStats{Error: "KAFKA_BROKERS not configured"} + } + + brokers := strings.Split(h.KafkaBrokers, ",") + brokerAddr := kafka.TCP(brokers...) + + kc := &kafka.Client{ + Addr: brokerAddr, + Timeout: 5 * time.Second, + } + + const topic = "iot.telemetry" + const group = "iot-pg-consumer" + + // Получаем latest offset (конец лога — сколько всего сообщений прошло) + offsetsResp, err := kc.ListOffsets(ctx, &kafka.ListOffsetsRequest{ + Addr: brokerAddr, + Topics: map[string][]kafka.OffsetRequest{ + topic: {kafka.LastOffsetOf(0)}, + }, + }) + if err != nil { + return iotAdminKafkaStats{Error: fmt.Sprintf("list offsets: %v", err)} + } + + var latestOffset int64 + if partitions, ok := offsetsResp.Topics[topic]; ok && len(partitions) > 0 { + if partitions[0].Error == nil { + latestOffset = partitions[0].LastOffset + } + } + + // Получаем committed offset consumer group (что consumer уже обработал) + fetchResp, err := kc.OffsetFetch(ctx, &kafka.OffsetFetchRequest{ + Addr: brokerAddr, + GroupID: group, + Topics: map[string][]int{topic: {0}}, + }) + if err != nil { + return iotAdminKafkaStats{ + LatestOffset: latestOffset, + Error: fmt.Sprintf("offset fetch: %v", err), + } + } + + var committedOffset int64 + if partitions, ok := fetchResp.Topics[topic]; ok && len(partitions) > 0 { + if partitions[0].Error == nil { + committedOffset = partitions[0].CommittedOffset + } + } + + lag := latestOffset - committedOffset + if lag < 0 { + lag = 0 + } + + return iotAdminKafkaStats{ + LatestOffset: latestOffset, + CommittedOffset: committedOffset, + ConsumerLag: lag, + } +} + +// collectIotPodStatuses собирает статус k8s pods для bridge и consumer по label app={name}. +func (h *Handler) collectIotPodStatuses(ctx context.Context) map[string]any { + result := map[string]any{} + + for _, appLabel := range []string{"iot-mqtt-bridge", "iot-kafka-consumer"} { + podList := &corev1.PodList{} + if err := h.K8s.List(ctx, podList, + client.InNamespace("sless"), + client.MatchingLabels{"app": appLabel}, + ); err != nil { + result[appLabel] = map[string]any{"error": err.Error()} + continue + } + if len(podList.Items) == 0 { + result[appLabel] = map[string]any{"status": "not found"} + continue + } + + pod := podList.Items[0] + var restarts int32 + for _, cs := range pod.Status.ContainerStatuses { + restarts += cs.RestartCount + } + ready := false + for _, cond := range pod.Status.Conditions { + if cond.Type == corev1.PodReady && cond.Status == corev1.ConditionTrue { + ready = true + } + } + + result[appLabel] = iotAdminPodStatus{ + Name: pod.Name, + Phase: string(pod.Status.Phase), + Ready: ready, + Restarts: restarts, + Age: iotFormatAge(pod.CreationTimestamp.Time), + } + } + + return result +} + +// iotFormatAge возвращает человекочитаемый возраст (s/m/h/d) pod-а. +func iotFormatAge(created time.Time) string { + d := time.Since(created) + switch { + case d < time.Minute: + return fmt.Sprintf("%ds", int(d.Seconds())) + case d < time.Hour: + return fmt.Sprintf("%dm", int(d.Minutes())) + case d < 24*time.Hour: + return fmt.Sprintf("%dh", int(d.Hours())) + default: + return fmt.Sprintf("%dd", int(d.Hours()/24)) + } +} diff --git a/internal/api/router.go b/internal/api/router.go index 7472b82..1723257 100644 --- a/internal/api/router.go +++ b/internal/api/router.go @@ -44,6 +44,11 @@ func NewRouter(h *handler.Handler, log *slog.Logger) http.Handler { // IoT Консоль — статический HTML, публично доступен r.HandleFunc("/console", ServeIoTConsole).Methods(http.MethodGet) + // IoT Admin — страница администратора (HTML без auth + JSON API с ADMIN_STATS_TOKEN) + // Не для конечных пользователей: показывает Kafka lag, pod statuses, PG stats per tenant. + r.HandleFunc("/iot-admin", ServeIoTAdmin).Methods(http.MethodGet) + r.HandleFunc("/iot-admin/stats", h.AdminStats).Methods(http.MethodGet) + // Публичный прокси для вызова HTTP-триггеров — без auth токена // Все HTTP методы разрешены (GET/POST/PUT/... — решает сама функция) r.PathPrefix("/fn/{namespace}/{name}").HandlerFunc(h.InvokeFunction) diff --git a/internal/api/ui/iot-admin.html b/internal/api/ui/iot-admin.html new file mode 100644 index 0000000..2195a13 --- /dev/null +++ b/internal/api/ui/iot-admin.html @@ -0,0 +1,547 @@ + + + + + + + Nubes IoT Admin + + + + + + + + + +
+ +
+

Доступ для администратора

+

Введите ADMIN_STATS_TOKEN для просмотра статистики IoT pipeline.

+ + + +
+ + + +
+ + + + diff --git a/internal/storage/iotpg/iot_telemetry_store.go b/internal/storage/iotpg/iot_telemetry_store.go index be00cba..798b9e3 100644 --- a/internal/storage/iotpg/iot_telemetry_store.go +++ b/internal/storage/iotpg/iot_telemetry_store.go @@ -225,6 +225,99 @@ func (s *IoTPostgresStore) QueryTelemetry(ctx context.Context, namespace, device return result, rows.Err() } +// TenantPGStats — статистика телеметрии одного tenant за разные периоды. +type TenantPGStats struct { + Namespace string `json:"namespace"` + DBName string `json:"db_name"` + Total int64 `json:"total"` + Last1h int64 `json:"last_1h"` + Last24h int64 `json:"last_24h"` + Latest []TelemetryRow `json:"latest"` + Error string `json:"error,omitempty"` +} + +// PostgresAdminStats — агрегированная статистика по всем tenant для страницы администратора. +type PostgresAdminStats struct { + Tenants []TenantPGStats `json:"tenants"` + TotalAll int64 `json:"total_all"` + Reachable bool `json:"reachable"` +} + +// GetAdminStats собирает статистику по всем tenant из management DB. +// Используется только страницей администратора — не для tenant API. +func (s *IoTPostgresStore) GetAdminStats(ctx context.Context) (*PostgresAdminStats, error) { + // Список всех тенантов из management DB + nsRows, err := s.adminDB.QueryContext(ctx, `SELECT namespace FROM tenant_credentials ORDER BY namespace`) + if err != nil { + return nil, fmt.Errorf("iotpg: list tenants: %w", err) + } + defer nsRows.Close() + + var namespaces []string + for nsRows.Next() { + var ns string + if err := nsRows.Scan(&ns); err != nil { + return nil, err + } + namespaces = append(namespaces, ns) + } + if err := nsRows.Err(); err != nil { + return nil, err + } + + result := &PostgresAdminStats{ + Reachable: true, + Tenants: make([]TenantPGStats, 0, len(namespaces)), + } + + for _, ns := range namespaces { + stats := TenantPGStats{ + Namespace: ns, + DBName: tenantDBName(ns), + } + + tenantDB, err := s.getTenantDB(ctx, ns) + if err != nil { + stats.Error = err.Error() + result.Tenants = append(result.Tenants, stats) + continue + } + + // Counts: total, last 1h, last 24h — одним запросом + err = tenantDB.QueryRowContext(ctx, ` +SELECT + COUNT(*), + COUNT(*) FILTER (WHERE ts > NOW() - INTERVAL '1 hour'), + COUNT(*) FILTER (WHERE ts > NOW() - INTERVAL '24 hours') +FROM iot_telemetry`).Scan(&stats.Total, &stats.Last1h, &stats.Last24h) + if err != nil { + stats.Error = err.Error() + result.Tenants = append(result.Tenants, stats) + continue + } + result.TotalAll += stats.Total + + // Последние 5 сообщений для предпросмотра + latestRows, err := tenantDB.QueryContext(ctx, + `SELECT id, device_id, ts, payload FROM iot_telemetry ORDER BY ts DESC LIMIT 5`) + if err == nil { + defer latestRows.Close() + for latestRows.Next() { + var r TelemetryRow + var rawPayload []byte + if err := latestRows.Scan(&r.ID, &r.DeviceID, &r.Ts, &rawPayload); err == nil { + r.Payload = json.RawMessage(rawPayload) + stats.Latest = append(stats.Latest, r) + } + } + } + + result.Tenants = append(result.Tenants, stats) + } + + return result, nil +} + // isDBNotExistErr проверяет что ошибка — «database does not exist» (PostgreSQL code 3D000). // Используется в QueryTelemetry: если DB нет — просто нет данных, не ошибка системы. func isDBNotExistErr(err error) bool { diff --git a/main.go b/main.go index cf4c5d1..d02ce75 100644 --- a/main.go +++ b/main.go @@ -227,12 +227,13 @@ func main() { // REST API сервер — запускается параллельно с operator manager apiHandler := slessapi.NewRouter(&handler.Handler{ - K8s: mgr.GetClient(), - Scheme: mgr.GetScheme(), - S3: s3Client, - PG: pg, - IoTPG: iotPGStore, - Log: log, + K8s: mgr.GetClient(), + Scheme: mgr.GetScheme(), + S3: s3Client, + PG: pg, + IoTPG: iotPGStore, + KafkaBrokers: os.Getenv("KAFKA_BROKERS"), + Log: log, }, log) go func() {