774 lines
24 KiB
Go
774 lines
24 KiB
Go
// 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
|
||
}
|
||
}
|
||
} |