Files
go-service/cmd/server/main.go

737 lines
21 KiB
Go
Raw 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.

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/logger"
"gpio-monitor/internal/pipe"
)
// ===== EMBED WEB =====
//go:embed web/dist web/css web/fonts 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)
}
// HandleLogFiles - список доступных лог-файлов
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 - чтение событий с пагинацией
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 файла с пагинацией
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",
})
}
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()
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"),
})
}
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] запрос от %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
}
}
}
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)
// 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)
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 поток с указанным URL
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
}
}
}