package main import ( "context" "embed" "encoding/json" "flag" "fmt" "io" "io/fs" "log" "net/http" "strings" "time" "gpio-monitor/internal/adapter" "gpio-monitor/internal/logger" "gpio-monitor/internal/pipe" ) // ===== EMBED WEB ===== // //go:embed web/dist web/css web/*.html web/*.png var webFiles embed.FS type API struct { buf *adapter.RingBuffer startTime time.Time } func writeJSON(w http.ResponseWriter, v any) { w.Header().Set("Content-Type", "application/json") json.NewEncoder(w).Encode(v) } func (a *API) HandleHealth(w http.ResponseWriter, r *http.Request) { latest, ok := a.buf.GetLatest() idle := time.Since(a.buf.LastWriteTime()) status := "На связи" if idle > 15*time.Second { status = "Нет связи" } writeJSON(w, map[string]any{ "status": status, "uptime_sec": time.Since(a.startTime).Seconds(), "last_data_ms": idle.Milliseconds(), "pipe_alive": a.buf.IsAlive(3 * time.Second), "has_data": ok, "latest": latest, "stats": a.buf.Stats(), }) } func (a *API) HandleLatest(w http.ResponseWriter, r *http.Request) { latest, ok := a.buf.GetLatest() if !ok { writeJSON(w, map[string]string{"error": "no data"}) return } writeJSON(w, map[string]any{ "latest": latest, "history": a.buf.GetLast(10), }) } func (a *API) HandleHistory(w http.ResponseWriter, r *http.Request) { raw := a.buf.GetLast(100) out := make([]int, len(raw)) for i, v := range raw { out[i] = int(v) } writeJSON(w, map[string]any{ "bytes": out, }) } func (a *API) HandleStream(w http.ResponseWriter, r *http.Request) { // Проверяем доступность MJPEG потока в go2rtc resp, err := http.Get("http://localhost:1984/api/streams?src=cam_mjpeg") camAvailable := err == nil && resp.StatusCode == 200 if resp != nil { resp.Body.Close() } writeJSON(w, map[string]interface{}{ "cam": "/api/cam", "available": camAvailable, "source": fmt.Sprintf("http://%s:1984/api/stream.mjpeg?src=cam_mjpeg", strings.Split(r.Host, ":")[0]), }) } // Прокси для MJPEG потока — просто ретранслирует готовый поток из go2rtc func handleCamProxy(w http.ResponseWriter, r *http.Request) { log.Printf("[cam] proxy request from %s", r.RemoteAddr) resp, err := http.Get("http://127.0.0.1:1984/api/stream.mjpeg?src=cam_mjpeg") if err != nil { log.Printf("[cam] go2rtc unavailable: %v", err) http.Error(w, "camera unavailable", 503) return } defer resp.Body.Close() // Прокидываем заголовки for k, vv := range resp.Header { for _, v := range vv { w.Header().Add(k, v) } } w.Header().Set("Access-Control-Allow-Origin", "*") w.WriteHeader(resp.StatusCode) flusher, ok := w.(http.Flusher) if !ok { http.Error(w, "stream unsupported", 500) return } buf := make([]byte, 32*1024) for { n, err := resp.Body.Read(buf) if n > 0 { _, err = w.Write(buf[:n]) if err != nil { log.Printf("[cam] client disconnected") return } flusher.Flush() } if err != nil { if err != io.EOF { log.Printf("[cam] stream ended: %v", err) } return } } } func cors(next http.HandlerFunc) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { w.Header().Set("Access-Control-Allow-Origin", "*") w.Header().Set("Access-Control-Allow-Methods", "GET, OPTIONS") w.Header().Set("Access-Control-Allow-Headers", "Content-Type") if r.Method == "OPTIONS" { w.WriteHeader(200) return } next(w, r) } } func main() { // Параметры командной строки pipePath := flag.String("pipe", "/tmp/gpio_pipe", "путь к pipe") // Параметры retention (хранения) retentionHours := flag.Int("retention-hours", 48, "количество часов хранения логов (0 = отключено)") retentionMB := flag.Int("retention-mb", 5000, "максимальный размер логов в MB (0 = отключено)") retentionIntervalMin := flag.Int("retention-interval", 15, "интервал проверки retention в минутах") // Параметры ротации файлов rotationCheckIntervalMin := flag.Int("rotation-check-interval", 1, "интервал проверки ротации в минутах") // Параметры буфера bufferSize := flag.Int("buffer-size", 10240, "размер кольцевого буфера (количество элементов)") // Параметры мониторинга silence5minAlert := flag.Int("silence-5min", 5, "время тишины для алерта 5 минут (в минутах)") silence10minAlert := flag.Int("silence-10min", 10, "время тишины для алерта 10 минут (в минутах)") watchdogIntervalMin := flag.Int("watchdog-interval", 1, "интервал проверки watchdog в минутах") // Параметры анти-спама alertCooldownSec := flag.Int("alert-cooldown", 10, "задержка между одинаковыми алертами в секундах") humanLogIntervalSec := flag.Int("human-log-interval", 5, "интервал записи human-readable логов в секундах") // Параметры сервера serverPort := flag.String("port", ":8080", "порт для веб-сервера") cameraURL := flag.String("camera-url", "http://127.0.0.1:1984/api/stream.mjpeg?src=cam_mjpeg", "URL камеры для прокси") flag.Parse() // Конвертируем в time.Duration retentionInterval := time.Duration(*retentionIntervalMin) * time.Minute rotationCheckInterval := time.Duration(*rotationCheckIntervalMin) * time.Minute silence5min := time.Duration(*silence5minAlert) * time.Minute silence10min := time.Duration(*silence10minAlert) * time.Minute watchdogInterval := time.Duration(*watchdogIntervalMin) * time.Minute alertCooldown := time.Duration(*alertCooldownSec) * time.Second humanLogInterval := time.Duration(*humanLogIntervalSec) * time.Second // Сохраняем параметры для передачи в логгеры config := &logger.Config{ RetentionHours: *retentionHours, RetentionMB: *retentionMB, RetentionIntervalMin: retentionInterval, Silence5Min: silence5min, Silence10Min: silence10min, WatchdogIntervalMin: watchdogInterval, AlertCooldownSec: alertCooldown, HumanLogIntervalSec: humanLogInterval, BufferSize: *bufferSize, } // Логируем конфигурацию log.Printf("=== КОНФИГУРАЦИЯ ===") log.Printf("Retention: %d часов, %d MB (проверка каждые %d мин)", *retentionHours, *retentionMB, *retentionIntervalMin) log.Printf("Ротация: проверка каждые %d мин", *rotationCheckIntervalMin) log.Printf("Буфер: %d элементов", *bufferSize) log.Printf("Алерты тишины: %d мин и %d мин (проверка каждые %d мин)", *silence5minAlert, *silence10minAlert, *watchdogIntervalMin) log.Printf("Анти-спам: %d сек, Human-лог: %d сек", *alertCooldownSec, *humanLogIntervalSec) log.Printf("Сервер: %s, Камера: %s", *serverPort, *cameraURL) log.Printf("====================") // ===== ИНИЦИАЛИЗАЦИЯ ЛОГГЕРОВ ===== // 1. EventLogger (события системы) eventLogger, err := logger.NewEventLogger() if err != nil { log.Fatal("Ошибка инициализации EventLogger:", err) } defer eventLogger.Close() eventLogger.Event("СЕРВЕР_ЗАПУЩЕН") // Выводим пути для информации if logsDir, err := logger.GetLogsDir(); err == nil { log.Printf("Логи сохраняются в: %s", logsDir) } // 1.5 Human-readable логгер humanLogger, err := logger.NewHumanLogger() if err != nil { log.Fatal("Ошибка инициализации HumanLogger:", err) } defer humanLogger.Close() eventLogger.Event("HUMAN_ЛОГЕР_ГОТОВ") // 2. Monitor (watchdog) с настройками monitor := logger.NewMonitor(eventLogger, config) monitor.SetSilenceThresholds(config.Silence5Min, config.Silence10Min) // 3. DataLogger с ротацией (с интервалом проверки) dataLogger, err := logger.NewRotatingLogger(eventLogger, rotationCheckInterval) if err != nil { eventLogger.Event("ОШИБКА_ИНИЦИАЛИЗАЦИИ_DATA_ЛОГЕРА") log.Fatal("Ошибка инициализации DataLogger:", err) } defer dataLogger.Close() eventLogger.Event("DATA_ЛОГЕР_ГОТОВ") // 3.5 Retention (автоматическая очистка) с настройками if *retentionHours > 0 || *retentionMB > 0 { retention := logger.NewRetention(*retentionHours, *retentionMB, eventLogger) retention.StartWithInterval(retentionInterval) eventLogger.Event("RETENTION_ГОТОВ") } else { log.Println("Retention отключен (часы и MB = 0)") } // 4. RingBuffer с настройкой размера buf := adapter.NewRingBuffer(*bufferSize) // 5. PipeReader с настройками pipeReader := pipe.NewPipeReader(buf, dataLogger, humanLogger, monitor, eventLogger, config) ctx, cancel := context.WithCancel(context.Background()) defer cancel() pipeReader.Start(ctx, *pipePath) // ===== API И WEB СЕРВЕР ===== api := &API{ buf: buf, startTime: time.Now(), } // ===== API ===== http.Handle("/api/", http.StripPrefix("/api", http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { switch r.URL.Path { case "/latest": cors(api.HandleLatest)(w, r) case "/history": cors(api.HandleHistory)(w, r) case "/health": cors(api.HandleHealth)(w, r) case "/stream": cors(api.HandleStream)(w, r) case "/cam": // Передаем URL камеры из конфигурации handleCamProxyWithURL(w, r, *cameraURL) default: http.Error(w, "not found", 404) } }))) // ===== WEB (EMBEDDED) ===== webFS, err := fs.Sub(webFiles, "web") if err != nil { log.Fatal(err) } http.Handle("/", http.FileServer(http.FS(webFS))) log.Printf("Server started on %s", *serverPort) log.Printf("Camera proxy available at /api/cam (source: %s)", *cameraURL) log.Fatal(http.ListenAndServe(*serverPort, nil)) } // handleCamProxyWithURL проксирует MJPEG поток с указанным URL func handleCamProxyWithURL(w http.ResponseWriter, r *http.Request, cameraURL string) { log.Printf("[cam] proxy request from %s to %s", r.RemoteAddr, cameraURL) resp, err := http.Get(cameraURL) if err != nil { log.Printf("[cam] camera unavailable: %v", err) http.Error(w, "camera unavailable", 503) return } defer resp.Body.Close() // Прокидываем заголовки for k, vv := range resp.Header { for _, v := range vv { w.Header().Add(k, v) } } w.Header().Set("Access-Control-Allow-Origin", "*") w.WriteHeader(resp.StatusCode) flusher, ok := w.(http.Flusher) if !ok { http.Error(w, "stream unsupported", 500) return } buf := make([]byte, 32*1024) for { n, err := resp.Body.Read(buf) if n > 0 { _, err = w.Write(buf[:n]) if err != nil { log.Printf("[cam] client disconnected") return } flusher.Flush() } if err != nil { if err != io.EOF { log.Printf("[cam] stream ended: %v", err) } return } } }