Доработка конфигурации логов

This commit is contained in:
Maxim
2026-06-15 15:41:54 +03:00
parent 87ab5d48de
commit 9a17569ac0
6 changed files with 270 additions and 55 deletions

View File

@@ -159,9 +159,69 @@ func cors(next http.HandlerFunc) http.HandlerFunc {
} }
func main() { func main() {
// Параметры командной строки
pipePath := flag.String("pipe", "/tmp/gpio_pipe", "путь к pipe") 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() 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 (события системы) // 1. EventLogger (события системы)
@@ -185,11 +245,12 @@ func main() {
defer humanLogger.Close() defer humanLogger.Close()
eventLogger.Event("HUMAN_ЛОГЕРОТОВ") eventLogger.Event("HUMAN_ЛОГЕРОТОВ")
// 2. Monitor (watchdog) // 2. Monitor (watchdog) с настройками
monitor := logger.NewMonitor(eventLogger) monitor := logger.NewMonitor(eventLogger, config)
monitor.SetSilenceThresholds(config.Silence5Min, config.Silence10Min)
// 3. DataLogger с ротацией // 3. DataLogger с ротацией (с интервалом проверки)
dataLogger, err := logger.NewRotatingLogger(eventLogger) dataLogger, err := logger.NewRotatingLogger(eventLogger, rotationCheckInterval)
if err != nil { if err != nil {
eventLogger.Event("ОШИБКАНИЦИАЛИЗАЦИИ_DATA_ЛОГЕРА") eventLogger.Event("ОШИБКАНИЦИАЛИЗАЦИИ_DATA_ЛОГЕРА")
log.Fatal("Ошибка инициализации DataLogger:", err) log.Fatal("Ошибка инициализации DataLogger:", err)
@@ -197,16 +258,20 @@ func main() {
defer dataLogger.Close() defer dataLogger.Close()
eventLogger.Event("DATA_ЛОГЕРОТОВ") eventLogger.Event("DATA_ЛОГЕРОТОВ")
// 3.5 Retention (автоматическая очистка) // 3.5 Retention (автоматическая очистка) с настройками
retention := logger.NewRetention(48, 5000, eventLogger) if *retentionHours > 0 || *retentionMB > 0 {
retention.Start() retention := logger.NewRetention(*retentionHours, *retentionMB, eventLogger)
eventLogger.Event("RETENTION_ГОТОВ") retention.StartWithInterval(retentionInterval)
eventLogger.Event("RETENTION_ГОТОВ")
} else {
log.Println("Retention отключен (часы и MB = 0)")
}
// 4. RingBuffer // 4. RingBuffer с настройкой размера
buf := adapter.NewRingBuffer(1024 * 10) buf := adapter.NewRingBuffer(*bufferSize)
// 5. PipeReader // 5. PipeReader с настройками
pipeReader := pipe.NewPipeReader(buf, dataLogger, humanLogger, monitor, eventLogger) pipeReader := pipe.NewPipeReader(buf, dataLogger, humanLogger, monitor, eventLogger, config)
ctx, cancel := context.WithCancel(context.Background()) ctx, cancel := context.WithCancel(context.Background())
defer cancel() defer cancel()
pipeReader.Start(ctx, *pipePath) pipeReader.Start(ctx, *pipePath)
@@ -229,7 +294,8 @@ func main() {
case "/stream": case "/stream":
cors(api.HandleStream)(w, r) cors(api.HandleStream)(w, r)
case "/cam": case "/cam":
handleCamProxy(w, r) // Передаем URL камеры из конфигурации
handleCamProxyWithURL(w, r, *cameraURL)
default: default:
http.Error(w, "not found", 404) http.Error(w, "not found", 404)
} }
@@ -242,7 +308,56 @@ func main() {
} }
http.Handle("/", http.FileServer(http.FS(webFS))) http.Handle("/", http.FileServer(http.FS(webFS)))
log.Println("Server started on :8080") log.Printf("Server started on %s", *serverPort)
log.Println("Camera proxy available at /api/cam") log.Printf("Camera proxy available at /api/cam (source: %s)", *cameraURL)
log.Fatal(http.ListenAndServe(":8080", nil)) 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
}
}
}

38
internal/logger/config.go Normal file
View File

@@ -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,
}
}

View File

@@ -6,20 +6,34 @@ import (
) )
type Monitor struct { type Monitor struct {
eventLogger *EventLogger eventLogger *EventLogger
lastWrite time.Time lastWrite time.Time
mu sync.RWMutex 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{ m := &Monitor{
eventLogger: eventLogger, eventLogger: eventLogger,
lastWrite: time.Now(), lastWrite: time.Now(),
silence5Min: config.Silence5Min,
silence10Min: config.Silence10Min,
watchdogInterval: config.WatchdogIntervalMin,
} }
go m.watchdogLoop() go m.watchdogLoop()
return m 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() { func (m *Monitor) RecordWrite() {
m.mu.Lock() m.mu.Lock()
m.lastWrite = time.Now() m.lastWrite = time.Now()
@@ -27,15 +41,19 @@ func (m *Monitor) RecordWrite() {
} }
func (m *Monitor) watchdogLoop() { func (m *Monitor) watchdogLoop() {
ticker := time.NewTicker(1 * time.Minute) ticker := time.NewTicker(m.watchdogInterval)
defer ticker.Stop()
for range ticker.C { for range ticker.C {
m.mu.RLock() m.mu.RLock()
silence := time.Since(m.lastWrite) silence := time.Since(m.lastWrite)
silence5Min := m.silence5Min
silence10Min := m.silence10Min
m.mu.RUnlock() m.mu.RUnlock()
if silence > 10*time.Minute { if silence > silence10Min {
m.eventLogger.Event("ТИШИНА_10МИН") m.eventLogger.Event("ТИШИНА_10МИН")
} else if silence > 5*time.Minute { } else if silence > silence5Min {
m.eventLogger.Event("ТИШИНА_5МИН") m.eventLogger.Event("ТИШИНА_5МИН")
} }
} }

View File

@@ -10,9 +10,10 @@ import (
) )
type Retention struct { type Retention struct {
maxAgeHours int maxAgeHours int
maxSizeBytes int64 maxSizeBytes int64
eventLogger *EventLogger eventLogger *EventLogger
stopCh chan struct{}
} }
func NewRetention(maxAgeHours int, maxSizeMB int, eventLogger *EventLogger) *Retention { func NewRetention(maxAgeHours int, maxSizeMB int, eventLogger *EventLogger) *Retention {
@@ -20,15 +21,27 @@ func NewRetention(maxAgeHours int, maxSizeMB int, eventLogger *EventLogger) *Ret
maxAgeHours: maxAgeHours, maxAgeHours: maxAgeHours,
maxSizeBytes: int64(maxSizeMB) * 1024 * 1024, maxSizeBytes: int64(maxSizeMB) * 1024 * 1024,
eventLogger: eventLogger, eventLogger: eventLogger,
stopCh: make(chan struct{}),
} }
} }
func (r *Retention) Start() { func (r *Retention) Start() {
// Запускаем проверку каждые 15 минут // Запускаем проверку каждые 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() { go func() {
for range ticker.C { for {
r.Cleanup() 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() { func (r *Retention) Cleanup() {
// Если оба лимита отключены, ничего не делаем
if r.maxAgeHours <= 0 && r.maxSizeBytes <= 0 {
return
}
// Получаем директорию с данными // Получаем директорию с данными
dataDir, err := GetDataLogsDir() dataDir, err := GetDataLogsDir()
if err != nil { if err != nil {
@@ -52,11 +74,17 @@ func (r *Retention) Cleanup() {
r.eventLogger.Event("ЗАПУЩЕНА_ОЧИСТКАОГОВ") r.eventLogger.Event("ЗАПУЩЕНА_ОЧИСТКАОГОВ")
} }
// 1. Удаляем старые файлы // 1. Удаляем старые файлы (если включено)
deletedByAge := r.cleanByAge(dataDir) deletedByAge := 0
if r.maxAgeHours > 0 {
deletedByAge = r.cleanByAge(dataDir)
}
// 2. Проверяем общий размер и удаляем самые старые если превышен лимит // 2. Проверяем общий размер и удаляем самые старые если превышен лимит (если включено)
deletedBySize := r.cleanBySize(dataDir) deletedBySize := 0
if r.maxSizeBytes > 0 {
deletedBySize = r.cleanBySize(dataDir)
}
if r.eventLogger != nil && (deletedByAge > 0 || deletedBySize > 0) { if r.eventLogger != nil && (deletedByAge > 0 || deletedBySize > 0) {
r.eventLogger.Event("УДАЛЕНОАЙЛОВ") r.eventLogger.Event("УДАЛЕНОАЙЛОВ")

View File

@@ -8,23 +8,30 @@ import (
) )
type RotatingLogger struct { type RotatingLogger struct {
dataLogger *DataLogger dataLogger *DataLogger
currentHour int currentHour int
baseDir string baseDir string
mu sync.Mutex mu sync.Mutex
eventLogger *EventLogger eventLogger *EventLogger
checkInterval time.Duration
} }
func NewRotatingLogger(eventLogger *EventLogger) (*RotatingLogger, error) { func NewRotatingLogger(eventLogger *EventLogger, checkInterval time.Duration) (*RotatingLogger, error) {
// Получаем директорию для бинарных данных // Получаем директорию для бинарных данных
dataDir, err := GetDataLogsDir() dataDir, err := GetDataLogsDir()
if err != nil { if err != nil {
return nil, err return nil, err
} }
// Интервал проверки по умолчанию
if checkInterval == 0 {
checkInterval = 1 * time.Minute
}
r := &RotatingLogger{ r := &RotatingLogger{
baseDir: dataDir, baseDir: dataDir,
eventLogger: eventLogger, eventLogger: eventLogger,
checkInterval: checkInterval,
} }
if err := r.rotate(); err != nil { if err := r.rotate(); err != nil {
@@ -87,7 +94,9 @@ func (r *RotatingLogger) Write(s Sample) {
} }
func (r *RotatingLogger) rotationLoop() { func (r *RotatingLogger) rotationLoop() {
ticker := time.NewTicker(1 * time.Minute) ticker := time.NewTicker(r.checkInterval)
defer ticker.Stop()
for range ticker.C { for range ticker.C {
r.rotate() r.rotate()
} }

View File

@@ -21,6 +21,10 @@ type PipeReader struct {
lastData logger.GPIOData lastData logger.GPIOData
lastLogTime time.Time lastLogTime time.Time
lastAlertTime map[byte]time.Time lastAlertTime map[byte]time.Time
// Настройки
alertCooldownSec time.Duration
humanLogIntervalSec time.Duration
} }
func NewPipeReader( func NewPipeReader(
@@ -29,14 +33,17 @@ func NewPipeReader(
humanLogger *logger.HumanLogger, humanLogger *logger.HumanLogger,
monitor *logger.Monitor, monitor *logger.Monitor,
eventLog *logger.EventLogger, eventLog *logger.EventLogger,
config *logger.Config,
) *PipeReader { ) *PipeReader {
return &PipeReader{ return &PipeReader{
buf: buf, buf: buf,
dataLogger: dataLogger, dataLogger: dataLogger,
humanLogger: humanLogger, humanLogger: humanLogger,
monitor: monitor, monitor: monitor,
eventLog: eventLog, eventLog: eventLog,
lastAlertTime: make(map[byte]time.Time), 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 { if pr.humanLogger != nil {
shouldLog := false shouldLog := false
@@ -109,8 +116,8 @@ func (pr *PipeReader) Start(ctx context.Context, pipePath string) {
shouldLog = true shouldLog = true
} }
// Или если прошло больше 5 секунд с последнего лога // Или если прошло больше N секунд с последнего лога
if now.Sub(pr.lastLogTime) >= 5*time.Second { if now.Sub(pr.lastLogTime) >= pr.humanLogIntervalSec {
shouldLog = true 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) { func (pr *PipeReader) handleAlert(data logger.GPIOData, timestamp time.Time) {
// Anti-spam: не чаще 1 алерта в 10 секунд для одинакового количества // Anti-spam: не чаще 1 алерта в N секунд для одинакового количества
key := data.Count key := data.Count
if last, exists := pr.lastAlertTime[key]; exists { if last, exists := pr.lastAlertTime[key]; exists {
if timestamp.Sub(last) < 10*time.Second { if timestamp.Sub(last) < pr.alertCooldownSec {
return return
} }
} }