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

642 lines
19 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/*.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
}
// Вычисляем смещение (нумерация страниц с 1)
// Страница 1 → самые ранние записи (начало файла)
// Страница N → самые поздние записи (конец файла)
start := (page - 1) * pageSize
end := start + pageSize
if end > total {
end = total
}
// Получаем строки для текущей страницы
var resultLines []string
if start < total {
resultLines = lines[start:end]
} else {
resultLines = []string{}
}
// Парсим строки для JSON
type Event struct {
Time string `json:"time"`
Event string `json:"event"`
}
var events []Event
for _, line := range resultLines {
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,
})
}
// 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
}
// Защита от 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"`
Count byte `json:"count"`
Strength byte `json:"strength"`
StrengthName string `json:"strength_name"`
Amplitude byte `json:"amplitude"` // Старшие 2 бита (0-3)
Signal byte `json:"signal"` // Младшие 6 бит (0-63)
}
var samples []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]
data := logger.ParseGPIO(value)
samples = append(samples, Sample{
Timestamp: ts,
Time: time.UnixMicro(ts).Format("02.01.2006 15:04:05"),
Value: value,
Count: data.Count,
Strength: data.Strength,
StrengthName: logger.GetStrengthName(data.Strength),
Amplitude: (value >> 6) & 0x3, // Старшие 2 бита (0-3) - АМПЛИТУДА
Signal: value & 0x3F, // Младшие 6 бит (0-63) - СИГНАЛ
})
}
// Ограничиваем количество точек для производительности
if len(samples) > 10000 {
step := len(samples) / 10000
var filtered []Sample
for i := 0; i < len(samples); i += step {
filtered = append(filtered, samples[i])
}
samples = filtered
}
writeJSON(w, map[string]any{
"filename": filename,
"total": len(samples),
"samples": samples,
})
}
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
}
}
}