diff --git a/cmd/server/main.go b/cmd/server/main.go index a411aa7..2af661c 100644 --- a/cmd/server/main.go +++ b/cmd/server/main.go @@ -159,9 +159,69 @@ func cors(next http.HandlerFunc) http.HandlerFunc { } 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("Retention: %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 (события системы) @@ -185,11 +245,12 @@ func main() { defer humanLogger.Close() eventLogger.Event("HUMAN_ЛОГЕР_ГОТОВ") - // 2. Monitor (watchdog) - monitor := logger.NewMonitor(eventLogger) + // 2. Monitor (watchdog) с настройками + monitor := logger.NewMonitor(eventLogger, config) + monitor.SetSilenceThresholds(config.Silence5Min, config.Silence10Min) - // 3. DataLogger с ротацией - dataLogger, err := logger.NewRotatingLogger(eventLogger) + // 3. DataLogger с ротацией (с интервалом проверки) + dataLogger, err := logger.NewRotatingLogger(eventLogger, rotationCheckInterval) if err != nil { eventLogger.Event("ОШИБКА_ИНИЦИАЛИЗАЦИИ_DATA_ЛОГЕРА") log.Fatal("Ошибка инициализации DataLogger:", err) @@ -197,16 +258,20 @@ func main() { defer dataLogger.Close() eventLogger.Event("DATA_ЛОГЕР_ГОТОВ") - // 3.5 Retention (автоматическая очистка) - retention := logger.NewRetention(48, 5000, eventLogger) - retention.Start() - eventLogger.Event("RETENTION_ГОТОВ") + // 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(1024 * 10) + // 4. RingBuffer с настройкой размера + buf := adapter.NewRingBuffer(*bufferSize) - // 5. PipeReader - pipeReader := pipe.NewPipeReader(buf, dataLogger, humanLogger, monitor, eventLogger) + // 5. PipeReader с настройками + pipeReader := pipe.NewPipeReader(buf, dataLogger, humanLogger, monitor, eventLogger, config) ctx, cancel := context.WithCancel(context.Background()) defer cancel() pipeReader.Start(ctx, *pipePath) @@ -229,7 +294,8 @@ func main() { case "/stream": cors(api.HandleStream)(w, r) case "/cam": - handleCamProxy(w, r) + // Передаем URL камеры из конфигурации + handleCamProxyWithURL(w, r, *cameraURL) default: http.Error(w, "not found", 404) } @@ -242,7 +308,56 @@ func main() { } http.Handle("/", http.FileServer(http.FS(webFS))) - log.Println("Server started on :8080") - log.Println("Camera proxy available at /api/cam") - log.Fatal(http.ListenAndServe(":8080", nil)) + log.Printf("Server started on %s", *serverPort) + log.Printf("Camera proxy available at /api/cam (source: %s)", *cameraURL) + log.Fatal(http.ListenAndServe(*serverPort, nil)) } + +// handleCamProxyWithURL проксирует MJPEG поток с указанным URL +func handleCamProxyWithURL(w http.ResponseWriter, r *http.Request, cameraURL string) { + log.Printf("[cam] proxy request from %s to %s", r.RemoteAddr, cameraURL) + + resp, err := http.Get(cameraURL) + if err != nil { + log.Printf("[cam] camera unavailable: %v", err) + http.Error(w, "camera unavailable", 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, "stream unsupported", 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] client disconnected") + return + } + flusher.Flush() + } + if err != nil { + if err != io.EOF { + log.Printf("[cam] stream ended: %v", err) + } + return + } + } +} \ No newline at end of file diff --git a/internal/logger/config.go b/internal/logger/config.go new file mode 100644 index 0000000..6dfc497 --- /dev/null +++ b/internal/logger/config.go @@ -0,0 +1,38 @@ +package logger + +import "time" + +// Config содержит настройки для всех компонентов логирования +type Config struct { + // Retention настройки + RetentionHours int + RetentionMB int + RetentionIntervalMin time.Duration + + // Мониторинг тишины + Silence5Min time.Duration + Silence10Min time.Duration + WatchdogIntervalMin time.Duration + + // Анти-спам и интервалы + AlertCooldownSec time.Duration + HumanLogIntervalSec time.Duration + + // Буфер + BufferSize int +} + +// DefaultConfig возвращает конфигурацию по умолчанию +func DefaultConfig() *Config { + return &Config{ + RetentionHours: 48, + RetentionMB: 5000, + RetentionIntervalMin: 15 * time.Minute, + Silence5Min: 5 * time.Minute, + Silence10Min: 10 * time.Minute, + WatchdogIntervalMin: 1 * time.Minute, + AlertCooldownSec: 10 * time.Second, + HumanLogIntervalSec: 5 * time.Second, + BufferSize: 10240, + } +} \ No newline at end of file diff --git a/internal/logger/monitor.go b/internal/logger/monitor.go index ad30392..3e96f4b 100644 --- a/internal/logger/monitor.go +++ b/internal/logger/monitor.go @@ -6,20 +6,34 @@ import ( ) type Monitor struct { - eventLogger *EventLogger - lastWrite time.Time - mu sync.RWMutex + eventLogger *EventLogger + lastWrite time.Time + mu sync.RWMutex + silence5Min time.Duration + silence10Min time.Duration + watchdogInterval time.Duration } -func NewMonitor(eventLogger *EventLogger) *Monitor { +func NewMonitor(eventLogger *EventLogger, config *Config) *Monitor { m := &Monitor{ - eventLogger: eventLogger, - lastWrite: time.Now(), + eventLogger: eventLogger, + lastWrite: time.Now(), + silence5Min: config.Silence5Min, + silence10Min: config.Silence10Min, + watchdogInterval: config.WatchdogIntervalMin, } go m.watchdogLoop() return m } +// SetSilenceThresholds позволяет изменить пороги тишины после создания +func (m *Monitor) SetSilenceThresholds(silence5Min, silence10Min time.Duration) { + m.mu.Lock() + defer m.mu.Unlock() + m.silence5Min = silence5Min + m.silence10Min = silence10Min +} + func (m *Monitor) RecordWrite() { m.mu.Lock() m.lastWrite = time.Now() @@ -27,15 +41,19 @@ func (m *Monitor) RecordWrite() { } func (m *Monitor) watchdogLoop() { - ticker := time.NewTicker(1 * time.Minute) + ticker := time.NewTicker(m.watchdogInterval) + defer ticker.Stop() + for range ticker.C { m.mu.RLock() silence := time.Since(m.lastWrite) + silence5Min := m.silence5Min + silence10Min := m.silence10Min m.mu.RUnlock() - if silence > 10*time.Minute { + if silence > silence10Min { m.eventLogger.Event("ТИШИНА_10МИН") - } else if silence > 5*time.Minute { + } else if silence > silence5Min { m.eventLogger.Event("ТИШИНА_5МИН") } } diff --git a/internal/logger/retention.go b/internal/logger/retention.go index feda26b..badd7d3 100644 --- a/internal/logger/retention.go +++ b/internal/logger/retention.go @@ -10,9 +10,10 @@ import ( ) type Retention struct { - maxAgeHours int - maxSizeBytes int64 - eventLogger *EventLogger + maxAgeHours int + maxSizeBytes int64 + eventLogger *EventLogger + stopCh chan struct{} } func NewRetention(maxAgeHours int, maxSizeMB int, eventLogger *EventLogger) *Retention { @@ -20,15 +21,27 @@ func NewRetention(maxAgeHours int, maxSizeMB int, eventLogger *EventLogger) *Ret maxAgeHours: maxAgeHours, maxSizeBytes: int64(maxSizeMB) * 1024 * 1024, eventLogger: eventLogger, + stopCh: make(chan struct{}), } } func (r *Retention) Start() { // Запускаем проверку каждые 15 минут - ticker := time.NewTicker(15 * time.Minute) + r.StartWithInterval(15 * time.Minute) +} + +func (r *Retention) StartWithInterval(interval time.Duration) { + // Запускаем проверку с указанным интервалом + ticker := time.NewTicker(interval) go func() { - for range ticker.C { - r.Cleanup() + for { + select { + case <-ticker.C: + r.Cleanup() + case <-r.stopCh: + ticker.Stop() + return + } } }() @@ -38,7 +51,16 @@ func (r *Retention) Start() { }) } +func (r *Retention) Stop() { + close(r.stopCh) +} + func (r *Retention) Cleanup() { + // Если оба лимита отключены, ничего не делаем + if r.maxAgeHours <= 0 && r.maxSizeBytes <= 0 { + return + } + // Получаем директорию с данными dataDir, err := GetDataLogsDir() if err != nil { @@ -52,11 +74,17 @@ func (r *Retention) Cleanup() { r.eventLogger.Event("ЗАПУЩЕНА_ОЧИСТКА_ЛОГОВ") } - // 1. Удаляем старые файлы - deletedByAge := r.cleanByAge(dataDir) + // 1. Удаляем старые файлы (если включено) + deletedByAge := 0 + if r.maxAgeHours > 0 { + deletedByAge = r.cleanByAge(dataDir) + } - // 2. Проверяем общий размер и удаляем самые старые если превышен лимит - deletedBySize := r.cleanBySize(dataDir) + // 2. Проверяем общий размер и удаляем самые старые если превышен лимит (если включено) + deletedBySize := 0 + if r.maxSizeBytes > 0 { + deletedBySize = r.cleanBySize(dataDir) + } if r.eventLogger != nil && (deletedByAge > 0 || deletedBySize > 0) { r.eventLogger.Event("УДАЛЕНО_ФАЙЛОВ") diff --git a/internal/logger/rotation.go b/internal/logger/rotation.go index 93a193c..cbddd70 100644 --- a/internal/logger/rotation.go +++ b/internal/logger/rotation.go @@ -8,23 +8,30 @@ import ( ) type RotatingLogger struct { - dataLogger *DataLogger - currentHour int - baseDir string - mu sync.Mutex - eventLogger *EventLogger + dataLogger *DataLogger + currentHour int + baseDir string + mu sync.Mutex + eventLogger *EventLogger + checkInterval time.Duration } -func NewRotatingLogger(eventLogger *EventLogger) (*RotatingLogger, error) { +func NewRotatingLogger(eventLogger *EventLogger, checkInterval time.Duration) (*RotatingLogger, error) { // Получаем директорию для бинарных данных dataDir, err := GetDataLogsDir() if err != nil { return nil, err } + // Интервал проверки по умолчанию + if checkInterval == 0 { + checkInterval = 1 * time.Minute + } + r := &RotatingLogger{ - baseDir: dataDir, - eventLogger: eventLogger, + baseDir: dataDir, + eventLogger: eventLogger, + checkInterval: checkInterval, } if err := r.rotate(); err != nil { @@ -87,7 +94,9 @@ func (r *RotatingLogger) Write(s Sample) { } func (r *RotatingLogger) rotationLoop() { - ticker := time.NewTicker(1 * time.Minute) + ticker := time.NewTicker(r.checkInterval) + defer ticker.Stop() + for range ticker.C { r.rotate() } diff --git a/internal/pipe/reader.go b/internal/pipe/reader.go index 3803aa9..f9c9738 100644 --- a/internal/pipe/reader.go +++ b/internal/pipe/reader.go @@ -21,6 +21,10 @@ type PipeReader struct { lastData logger.GPIOData lastLogTime time.Time lastAlertTime map[byte]time.Time + + // Настройки + alertCooldownSec time.Duration + humanLogIntervalSec time.Duration } func NewPipeReader( @@ -29,14 +33,17 @@ func NewPipeReader( humanLogger *logger.HumanLogger, monitor *logger.Monitor, eventLog *logger.EventLogger, + config *logger.Config, ) *PipeReader { return &PipeReader{ - buf: buf, - dataLogger: dataLogger, - humanLogger: humanLogger, - monitor: monitor, - eventLog: eventLog, - lastAlertTime: make(map[byte]time.Time), + buf: buf, + dataLogger: dataLogger, + humanLogger: humanLogger, + monitor: monitor, + eventLog: eventLog, + lastAlertTime: make(map[byte]time.Time), + alertCooldownSec: config.AlertCooldownSec, + humanLogIntervalSec: config.HumanLogIntervalSec, } } @@ -100,7 +107,7 @@ func (pr *PipeReader) Start(ctx context.Context, pipePath string) { }) } - // 3. Human-readable лог (при изменении состояния или раз в 5 секунд) + // 3. Human-readable лог (при изменении состояния или раз в N секунд) if pr.humanLogger != nil { shouldLog := false @@ -109,8 +116,8 @@ func (pr *PipeReader) Start(ctx context.Context, pipePath string) { shouldLog = true } - // Или если прошло больше 5 секунд с последнего лога - if now.Sub(pr.lastLogTime) >= 5*time.Second { + // Или если прошло больше N секунд с последнего лога + if now.Sub(pr.lastLogTime) >= pr.humanLogIntervalSec { shouldLog = true } @@ -149,10 +156,10 @@ func (pr *PipeReader) Start(ctx context.Context, pipePath string) { } func (pr *PipeReader) handleAlert(data logger.GPIOData, timestamp time.Time) { - // Anti-spam: не чаще 1 алерта в 10 секунд для одинакового количества + // Anti-spam: не чаще 1 алерта в N секунд для одинакового количества key := data.Count if last, exists := pr.lastAlertTime[key]; exists { - if timestamp.Sub(last) < 10*time.Second { + if timestamp.Sub(last) < pr.alertCooldownSec { return } }