Компоненты: - iot-operator: controller-manager (IoTDevice CRD) + REST API (порт 9090) - mqtt-bridge: MQTT (EMQX) → Kafka bridge - kafka-consumer: Kafka → Postgres pipeline Модуль: gitea.services.ngcloud.ru/Nail/IoT Все 3 бинарника собираются, import paths адаптированы.
215 lines
6.5 KiB
Go
215 lines
6.5 KiB
Go
// Создано: 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))
|
||
}
|
||
}
|