// gpio-monitor-server — HTTP-сервер мониторинга GPIO-шины Raspberry Pi. // Читает поток байт из FIFO pipe, хранит их в кольцевом буфере и бинарных // почасовых логах, отдаёт JSON API и встроенный веб-дашборд на :8080. // Конфигурация — только флагами командной строки (см. main или docs/operations.md). package main import ( "context" "embed" "encoding/binary" "encoding/json" "flag" "fmt" "io" "io/fs" "log" "net/http" "os" "path/filepath" "sort" "strconv" "strings" "time" "gpio-monitor/internal/adapter" "gpio-monitor/internal/audio" "gpio-monitor/internal/logger" "gpio-monitor/internal/pipe" ) // ===== EMBED WEB ===== //go:embed web/dist web/css web/fonts web/*.html web/*.png var webFiles embed.FS // API — состояние HTTP-хендлеров: кольцевой буфер, время старта сервера // и монитор звука. Спецификация ответов — docs/api.md. type API struct { buf *adapter.RingBuffer startTime time.Time audioMonitor *audio.Monitor } func writeJSON(w http.ResponseWriter, v any) { w.Header().Set("Content-Type", "application/json") json.NewEncoder(w).Encode(v) } // HandleLogFiles отдаёт список бинарных лог-файлов gpio-*.bin, // отсортированный по времени изменения (новые первыми). func (a *API) HandleLogFiles(w http.ResponseWriter, r *http.Request) { dataDir, err := logger.GetDataLogsDir() if err != nil { http.Error(w, err.Error(), 500) return } files, err := filepath.Glob(filepath.Join(dataDir, "gpio-*.bin")) if err != nil { http.Error(w, err.Error(), 500) return } type LogFileInfo struct { Name string `json:"name"` Path string `json:"path"` Size int64 `json:"size"` ModTime time.Time `json:"mod_time"` Time string `json:"time"` Date string `json:"date"` Hour int `json:"hour"` IsActive bool `json:"is_active"` } var result []LogFileInfo for _, file := range files { info, err := os.Stat(file) if err != nil { continue } base := filepath.Base(file) parts := strings.Split(strings.TrimSuffix(base, ".bin"), "-") var fileTime time.Time if len(parts) == 5 { year, _ := strconv.Atoi(parts[1]) month, _ := strconv.Atoi(parts[2]) day, _ := strconv.Atoi(parts[3]) hour, _ := strconv.Atoi(parts[4]) fileTime = time.Date(year, time.Month(month), day, hour, 0, 0, 0, time.Local) } else { fileTime = info.ModTime() } result = append(result, LogFileInfo{ Name: base, Path: file, Size: info.Size(), ModTime: info.ModTime(), Time: fileTime.Format("02.01.2006 15:00"), Date: fileTime.Format("2006-01-02"), Hour: fileTime.Hour(), IsActive: false, }) } // Сортируем по времени sort.Slice(result, func(i, j int) bool { return result[i].ModTime.After(result[j].ModTime) }) writeJSON(w, result) } // HandleLogEvents отдаёт события из events_human.log с пагинацией // (page, page_size; страница 1 — самые новые). func (a *API) HandleLogEvents(w http.ResponseWriter, r *http.Request) { // Параметры пагинации page := 1 if p := r.URL.Query().Get("page"); p != "" { if parsed, err := strconv.Atoi(p); err == nil && parsed > 0 { page = parsed } } pageSize := 100 if ps := r.URL.Query().Get("page_size"); ps != "" { if parsed, err := strconv.Atoi(ps); err == nil && parsed > 0 { pageSize = parsed } } eventsPath, err := logger.GetEventHumanLogPath() if err != nil { http.Error(w, err.Error(), 500) return } content, err := os.ReadFile(eventsPath) if err != nil { http.Error(w, err.Error(), 500) return } // Разбиваем на строки lines := strings.Split(string(content), "\n") // Удаляем пустые строки var nonEmptyLines []string for _, line := range lines { if line != "" { nonEmptyLines = append(nonEmptyLines, line) } } lines = nonEmptyLines total := len(lines) totalPages := (total + pageSize - 1) / pageSize if totalPages == 0 { totalPages = 1 } // Корректируем номер страницы if page > totalPages { page = totalPages } if page < 1 { page = 1 } start := total - page*pageSize if start < 0 { start = 0 } end := start + pageSize if end > total { end = total } // Получаем строки для текущей страницы var pageLines []string if start < total && start >= 0 { pageLines = lines[start:end] } else { pageLines = []string{} } var reversedLines []string for i := len(pageLines) - 1; i >= 0; i-- { reversedLines = append(reversedLines, pageLines[i]) } // Парсим строки для JSON type Event struct { Time string `json:"time"` Event string `json:"event"` } var events []Event for _, line := range reversedLines { if line == "" { continue } parts := strings.SplitN(line, "] EVENT: ", 2) if len(parts) == 2 { events = append(events, Event{ Time: strings.TrimPrefix(parts[0], "["), Event: parts[1], }) } } writeJSON(w, map[string]any{ "events": events, "total": total, "page": page, "page_size": pageSize, "total_pages": totalPages, "has_previous": page > 1, "has_next": page < totalPages, "returned": len(events), "start_index": start + 1, "end_index": end, "order": "newest_first", }) } // HandleLogData отдаёт сэмплы одного .bin-файла (параметр file) с пагинацией // (page, page_size; страница 1 — самые новые). Формат записи — 9 байт, // см. docs/data-formats.md. Файл читается в память целиком. func (a *API) HandleLogData(w http.ResponseWriter, r *http.Request) { filename := r.URL.Query().Get("file") if filename == "" { http.Error(w, "отсутствует параметр file", 400) return } // Параметры пагинации page := 1 if p := r.URL.Query().Get("page"); p != "" { if parsed, err := strconv.Atoi(p); err == nil && parsed > 0 { page = parsed } } pageSize := 100 if ps := r.URL.Query().Get("page_size"); ps != "" { if parsed, err := strconv.Atoi(ps); err == nil && parsed > 0 { pageSize = parsed } } // Защита от path traversal filename = filepath.Base(filename) dataDir, err := logger.GetDataLogsDir() if err != nil { http.Error(w, err.Error(), 500) return } filePath := filepath.Join(dataDir, filename) if _, err := os.Stat(filePath); os.IsNotExist(err) { http.Error(w, "файл не найден", 404) return } // Читаем бинарный файл f, err := os.Open(filePath) if err != nil { http.Error(w, err.Error(), 500) return } defer f.Close() type Sample struct { Timestamp int64 `json:"ts"` Time string `json:"time"` Value byte `json:"value"` Amplitude byte `json:"amplitude"` // Старшие 2 бита (0-3) Signal byte `json:"signal"` // Младшие 6 бит (0-63) StrengthName string `json:"strength_name"` } // Сначала читаем все сэмплы var allSamples []Sample buf := make([]byte, 9) // 8 байт timestamp + 1 байт value for { n, err := f.Read(buf) if err == io.EOF { break } if err != nil { http.Error(w, err.Error(), 500) return } if n < 9 { continue } ts := int64(binary.LittleEndian.Uint64(buf[0:8])) value := buf[8] amplitude := (value >> 6) & 0x3 signal := value & 0x3F allSamples = append(allSamples, Sample{ Timestamp: ts, Time: time.UnixMicro(ts).Format("02.01.2006 15:04:05"), Value: value, Amplitude: amplitude, Signal: signal, StrengthName: logger.GetStrengthName(amplitude), }) } total := len(allSamples) if total == 0 { writeJSON(w, map[string]any{ "filename": filename, "total": 0, "page": 1, "page_size": pageSize, "total_pages": 1, "has_previous": false, "has_next": false, "start_index": 0, "end_index": 0, "samples": []Sample{}, "stats": map[string]any{ "total_points": 0, "amplitude_over_0": 0, "signal_over_10": 0, }, "order": "newest_first", }) return } totalPages := (total + pageSize - 1) / pageSize if totalPages == 0 { totalPages = 1 } // Корректируем номер страницы if page > totalPages { page = totalPages } if page < 1 { page = 1 } start := total - page*pageSize if start < 0 { start = 0 } end := start + pageSize if end > total { end = total } // Получаем сэмплы для текущей страницы var pageSamples []Sample if start < total && start >= 0 { pageSamples = allSamples[start:end] } else { pageSamples = []Sample{} } // Разворачиваем, чтобы новые были сверху var reversedSamples []Sample for i := len(pageSamples) - 1; i >= 0; i-- { reversedSamples = append(reversedSamples, pageSamples[i]) } // Считаем статистику amplitudeOver0 := 0 signalOver10 := 0 for _, s := range allSamples { if s.Amplitude > 0 { amplitudeOver0++ } if s.Signal > 10 { signalOver10++ } } writeJSON(w, map[string]any{ "filename": filename, "total": total, "page": page, "page_size": pageSize, "total_pages": totalPages, "has_previous": page > 1, "has_next": page < totalPages, "start_index": start + 1, "end_index": end, "samples": reversedSamples, "stats": map[string]any{ "total_points": total, "amplitude_over_0": amplitudeOver0, "signal_over_10": signalOver10, }, "order": "newest_first", }) } // HandleHealth отдаёт статус сервера: связь с pipe (пороги 15 с и 3 с), // uptime, статистику буфера и уровень звука. 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 = "Нет связи" } now := time.Now() soundLevel := 0 if a.audioMonitor != nil { soundLevel = a.audioMonitor.GetLevel() } 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(), "server_time": now.Format("02.01.2006 15:04"), "sound_level": soundLevel, }) } // HandleLatest отдаёт последний байт и до 10 предыдущих (от старых к новым). 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), }) } // HandleHistory отдаёт до 300 последних байт для графика (от старых к новым). func (a *API) HandleHistory(w http.ResponseWriter, r *http.Request) { raw := a.buf.GetLast(300) out := make([]int, len(raw)) for i, v := range raw { out[i] = int(v) } writeJSON(w, map[string]any{ "bytes": out, }) } // HandleStream отдаёт статус видеопотока камеры. Внимание: доступность // проверяется по захардкоженному localhost:1984 без учёта флага -camera-url // (известное поведение, docs/tech-debt.md). 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]), }) } // handleCamProxy — НЕ ИСПОЛЬЗУЕТСЯ: роутер вызывает handleCamProxyWithURL, // этот вариант с захардкоженным URL остался как мёртвый код (docs/tech-debt.md). func handleCamProxy(w http.ResponseWriter, r *http.Request) { log.Printf("[cam] запрос от %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 недоступен: %v", err) http.Error(w, "камера недоступна", 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, "поток не поддерживается", 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] клиент отключился") return } flusher.Flush() } if err != nil { if err != io.EOF { log.Printf("[cam] поток завершен: %v", err) } return } } } // cors разрешает кросс-доменные GET-запросы с любых источников // (рассчитано на доверенную локальную сеть). 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("Хранение: %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) // 4.5 Монитор уровня звука с микрофона audioMonitor := audio.NewMonitor() if err := audioMonitor.Start(); err != nil { log.Printf("Аудио-монитор не запущен (микрофон недоступен?): %v", err) } defer audioMonitor.Stop() // 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(), audioMonitor: audioMonitor, } // ===== 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) case "/log/files": cors(api.HandleLogFiles)(w, r) case "/log/data": cors(api.HandleLogData)(w, r) case "/log/events": cors(api.HandleLogEvents)(w, r) 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("Сервер запущен на %s", *serverPort) log.Printf("Прокси камеры доступен по адресу /api/cam (источник: %s)", *cameraURL) log.Fatal(http.ListenAndServe(*serverPort, nil)) } // handleCamProxyWithURL проксирует MJPEG-поток с cameraURL (флаг -camera-url), // ретранслируя заголовки и сбрасывая буфер после каждого чанка 32 КБ. func handleCamProxyWithURL(w http.ResponseWriter, r *http.Request, cameraURL string) { log.Printf("[cam] запрос от %s к %s", r.RemoteAddr, cameraURL) resp, err := http.Get(cameraURL) if err != nil { log.Printf("[cam] камера недоступна: %v", err) http.Error(w, "камера недоступна", 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, "поток не поддерживается", 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] клиент отключился") return } flusher.Flush() } if err != nil { if err != io.EOF { log.Printf("[cam] поток завершен: %v", err) } return } } }