// Создано: 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)) } }