Files
go-service/cmd/server/main.go
2026-07-17 15:57:05 +03:00

774 lines
24 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

// 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
}
}
}