diff --git a/cmd/server/main.go b/cmd/server/main.go index 14efac7..ca10c3a 100644 --- a/cmd/server/main.go +++ b/cmd/server/main.go @@ -10,11 +10,12 @@ import ( "io/fs" "log" "net/http" - "os" "strings" "time" "gpio-monitor/internal/adapter" + "gpio-monitor/internal/logger" + "gpio-monitor/internal/pipe" ) // ===== EMBED WEB ===== @@ -125,7 +126,6 @@ func handleCamProxy(w http.ResponseWriter, r *http.Request) { for { n, err := resp.Body.Read(buf) - if n > 0 { _, err = w.Write(buf[:n]) if err != nil { @@ -134,7 +134,6 @@ func handleCamProxy(w http.ResponseWriter, r *http.Request) { } flusher.Flush() } - if err != nil { if err != io.EOF { log.Printf("[cam] stream ended: %v", err) @@ -144,43 +143,6 @@ func handleCamProxy(w http.ResponseWriter, r *http.Request) { } } -// ===== PIPE READER ===== -func startPipeReader(ctx context.Context, buf *adapter.RingBuffer, path string) { - go func() { - buffer := make([]byte, 64) - - for { - select { - case <-ctx.Done(): - return - default: - } - - f, err := os.OpenFile(path, os.O_RDONLY, 0) - if err != nil { - time.Sleep(time.Second) - continue - } - - log.Println("pipe connected") - - for { - n, err := f.Read(buffer) - if err != nil { - f.Close() - break - } - - for i := 0; i < n; i++ { - buf.Write(buffer[i]) - } - } - - log.Println("pipe disconnected, retry...") - } - }() -} - func cors(next http.HandlerFunc) http.HandlerFunc { return func(w http.ResponseWriter, r *http.Request) { w.Header().Set("Access-Control-Allow-Origin", "*") @@ -197,15 +159,49 @@ func cors(next http.HandlerFunc) http.HandlerFunc { } func main() { - pipePath := flag.String("pipe", "/tmp/gpio_pipe", "pipe path") + pipePath := flag.String("pipe", "/tmp/gpio_pipe", "путь к pipe") flag.Parse() + // ===== ИНИЦИАЛИЗАЦИЯ ЛОГГЕРОВ ===== + + // 1. EventLogger (события системы) + eventLogger, err := logger.NewEventLogger("logs/events.log") + if err != nil { + log.Fatal("Ошибка инициализации EventLogger:", err) + } + defer eventLogger.Close() + eventLogger.Event("СЕРВЕР_ЗАПУЩЕН") + + // 1.5 Human-readable логгер + humanLogger, err := logger.NewHumanLogger("logs/gpio_human.log") + if err != nil { + log.Fatal("Ошибка инициализации HumanLogger:", err) + } + defer humanLogger.Close() + eventLogger.Event("HUMAN_ЛОГЕР_ГОТОВ") + + // 2. Monitor (watchdog) + monitor := logger.NewMonitor(eventLogger) + + // 3. DataLogger с ротацией + dataLogger, err := logger.NewRotatingLogger("logs/data", eventLogger) + if err != nil { + eventLogger.Event("ОШИБКА_ИНИЦИАЛИЗАЦИИ_DATA_ЛОГЕРА") + log.Fatal("Ошибка инициализации DataLogger:", err) + } + defer dataLogger.Close() + eventLogger.Event("DATA_ЛОГЕР_ГОТОВ") + + // 4. RingBuffer buf := adapter.NewRingBuffer(1024 * 10) + + // 5. PipeReader + pipeReader := pipe.NewPipeReader(buf, dataLogger, humanLogger, monitor, eventLogger) ctx, cancel := context.WithCancel(context.Background()) defer cancel() + pipeReader.Start(ctx, *pipePath) - startPipeReader(ctx, buf, *pipePath) - + // ===== API И WEB СЕРВЕР ===== api := &API{ buf: buf, startTime: time.Now(), @@ -234,7 +230,6 @@ func main() { if err != nil { log.Fatal(err) } - http.Handle("/", http.FileServer(http.FS(webFS))) log.Println("Server started on :8080") diff --git a/internal/logger/data_logger.go b/internal/logger/data_logger.go new file mode 100644 index 0000000..3d9cc92 --- /dev/null +++ b/internal/logger/data_logger.go @@ -0,0 +1,98 @@ +package logger + +import ( + "encoding/binary" + "os" + "sync" + "time" +) + +type DataLogger struct { + file *os.File + mu sync.Mutex + buffer []byte + bufferSize int + flushTick *time.Ticker + closeCh chan struct{} +} + +type Sample struct { + Timestamp int64 // microseconds + Value byte +} + +func NewDataLogger(basePath string) (*DataLogger, error) { + // путь будет формироваться через rotation + f, err := os.OpenFile(basePath, os.O_CREATE|os.O_APPEND|os.O_WRONLY, 0644) + if err != nil { + return nil, err + } + + d := &DataLogger{ + file: f, + buffer: make([]byte, 0, 64*1024), // 64KB буфер + bufferSize: 0, + flushTick: time.NewTicker(1 * time.Second), + closeCh: make(chan struct{}), + } + + go d.flushLoop() + return d, nil +} + +func (d *DataLogger) Write(s Sample) { + d.mu.Lock() + defer d.mu.Unlock() + + // Формат: [timestamp uint64][value byte] + tsBuf := make([]byte, 8) + binary.LittleEndian.PutUint64(tsBuf, uint64(s.Timestamp)) + + d.buffer = append(d.buffer, tsBuf...) + d.buffer = append(d.buffer, s.Value) + d.bufferSize += 9 + + // Если буфер переполнен - сбрасываем немедленно + if d.bufferSize >= 64*1024 { + d.flush() + } +} + +func (d *DataLogger) flush() { + if d.bufferSize == 0 { + return + } + + _, err := d.file.Write(d.buffer[:d.bufferSize]) + if err != nil { + // тут должен быть event + return + } + + d.file.Sync() // важно для durability + + d.buffer = d.buffer[:0] + d.bufferSize = 0 +} + +func (d *DataLogger) flushLoop() { + for { + select { + case <-d.flushTick.C: + d.mu.Lock() + d.flush() + d.mu.Unlock() + case <-d.closeCh: + d.mu.Lock() + d.flush() + d.mu.Unlock() + return + } + } +} + +func (d *DataLogger) Close() error { + close(d.closeCh) + d.flushTick.Stop() + return d.file.Close() +} \ No newline at end of file diff --git a/internal/logger/event_logger.go b/internal/logger/event_logger.go new file mode 100644 index 0000000..5e9dec5 --- /dev/null +++ b/internal/logger/event_logger.go @@ -0,0 +1,78 @@ +package logger + +import ( + "encoding/json" + "fmt" + "os" + "path/filepath" + "sync" + "time" +) + +type EventLogger struct { + file *os.File + humanFile *os.File + mu sync.Mutex +} + +func NewEventLogger(path string) (*EventLogger, error) { + // Создаём директорию для логов + dir := filepath.Dir(path) + if err := os.MkdirAll(dir, 0755); err != nil { + return nil, err + } + + // JSON файл для событий + f, err := os.OpenFile(path, os.O_CREATE|os.O_APPEND|os.O_WRONLY, 0644) + if err != nil { + return nil, err + } + + // Human-readable файл для событий + humanPath := filepath.Join(dir, "events_human.log") + humanF, err := os.OpenFile(humanPath, os.O_CREATE|os.O_APPEND|os.O_WRONLY, 0644) + if err != nil { + f.Close() + return nil, err + } + + return &EventLogger{ + file: f, + humanFile: humanF, + }, nil +} + +func (l *EventLogger) Event(name string) { + l.mu.Lock() + defer l.mu.Unlock() + + now := time.Now() + + // JSON формат (для машин) + ev := struct { + Timestamp int64 `json:"ts"` + Time string `json:"time"` + Event string `json:"event"` + }{ + Timestamp: now.Unix(), + Time: now.Format("2006-01-02 15:04:05"), + Event: name, + } + + data, _ := json.Marshal(ev) + data = append(data, '\n') + l.file.Write(data) + + // Human-readable формат + humanLine := fmt.Sprintf("[%s] EVENT: %s\n", now.Format("2006-01-02 15:04:05.000"), name) + l.humanFile.WriteString(humanLine) + + l.file.Sync() + l.humanFile.Sync() +} + +func (l *EventLogger) Close() error { + l.file.Close() + l.humanFile.Close() + return nil +} diff --git a/internal/logger/human_logger.go b/internal/logger/human_logger.go new file mode 100644 index 0000000..4b4a97a --- /dev/null +++ b/internal/logger/human_logger.go @@ -0,0 +1,36 @@ +package logger + +import ( + "os" + "sync" +) + +type HumanLogger struct { + file *os.File + mu sync.Mutex +} + +func NewHumanLogger(path string) (*HumanLogger, error) { + if err := os.MkdirAll("logs", 0755); err != nil { + return nil, err + } + + f, err := os.OpenFile(path, os.O_CREATE|os.O_APPEND|os.O_WRONLY, 0644) + if err != nil { + return nil, err + } + + return &HumanLogger{file: f}, nil +} + +func (l *HumanLogger) Write(data string) { + l.mu.Lock() + defer l.mu.Unlock() + + l.file.WriteString(data) + l.file.Sync() +} + +func (l *HumanLogger) Close() error { + return l.file.Close() +} diff --git a/internal/logger/monitor.go b/internal/logger/monitor.go new file mode 100644 index 0000000..ad30392 --- /dev/null +++ b/internal/logger/monitor.go @@ -0,0 +1,42 @@ +package logger + +import ( + "sync" + "time" +) + +type Monitor struct { + eventLogger *EventLogger + lastWrite time.Time + mu sync.RWMutex +} + +func NewMonitor(eventLogger *EventLogger) *Monitor { + m := &Monitor{ + eventLogger: eventLogger, + lastWrite: time.Now(), + } + go m.watchdogLoop() + return m +} + +func (m *Monitor) RecordWrite() { + m.mu.Lock() + m.lastWrite = time.Now() + m.mu.Unlock() +} + +func (m *Monitor) watchdogLoop() { + ticker := time.NewTicker(1 * time.Minute) + for range ticker.C { + m.mu.RLock() + silence := time.Since(m.lastWrite) + m.mu.RUnlock() + + if silence > 10*time.Minute { + m.eventLogger.Event("ТИШИНА_10МИН") + } else if silence > 5*time.Minute { + m.eventLogger.Event("ТИШИНА_5МИН") + } + } +} \ No newline at end of file diff --git a/internal/logger/rotation.go b/internal/logger/rotation.go new file mode 100644 index 0000000..ed459a8 --- /dev/null +++ b/internal/logger/rotation.go @@ -0,0 +1,98 @@ +package logger + +import ( + "fmt" + "os" + "path/filepath" + "sync" + "time" +) + +type RotatingLogger struct { + dataLogger *DataLogger + currentHour int + baseDir string + mu sync.Mutex + eventLogger *EventLogger +} + +func NewRotatingLogger(baseDir string, eventLogger *EventLogger) (*RotatingLogger, error) { + if err := os.MkdirAll(baseDir, 0755); err != nil { + return nil, err + } + + r := &RotatingLogger{ + baseDir: baseDir, + eventLogger: eventLogger, + } + + if err := r.rotate(); err != nil { + return nil, err + } + + go r.rotationLoop() + return r, nil +} + +func (r *RotatingLogger) getFilename(hour int) string { + now := time.Now() + return filepath.Join(r.baseDir, fmt.Sprintf("gpio-%04d-%02d-%02d-%02d.bin", + now.Year(), now.Month(), now.Day(), hour)) +} + +func (r *RotatingLogger) rotate() error { + r.mu.Lock() + defer r.mu.Unlock() + + now := time.Now() + newHour := now.Hour() + + // Если уже правильный час и логгер существует - ок + if r.dataLogger != nil && r.currentHour == newHour { + return nil + } + + // Закрываем старый + if r.dataLogger != nil { + r.dataLogger.Close() + } + + // Открываем новый + filename := r.getFilename(newHour) + dataLogger, err := NewDataLogger(filename) + if err != nil { + r.eventLogger.Event("ROTATION_FAILED") + return err + } + + r.dataLogger = dataLogger + r.currentHour = newHour + r.eventLogger.Event("ROTATION_COMPLETE") + return nil +} + +func (r *RotatingLogger) Write(s Sample) { + r.mu.Lock() + logger := r.dataLogger + r.mu.Unlock() + + if logger != nil { + logger.Write(s) + } +} + +func (r *RotatingLogger) rotationLoop() { + ticker := time.NewTicker(1 * time.Minute) + for range ticker.C { + r.rotate() + } +} + +func (r *RotatingLogger) Close() error { + r.mu.Lock() + defer r.mu.Unlock() + if r.dataLogger != nil { + return r.dataLogger.Close() + } + return nil +} \ No newline at end of file diff --git a/internal/pipe/reader.go b/internal/pipe/reader.go new file mode 100644 index 0000000..9c28132 --- /dev/null +++ b/internal/pipe/reader.go @@ -0,0 +1,146 @@ +package pipe + +import ( + "context" + "fmt" + "os" + "time" + + "gpio-monitor/internal/adapter" + "gpio-monitor/internal/logger" +) + +type PipeReader struct { + buf *adapter.RingBuffer + dataLogger *logger.RotatingLogger + humanLogger *logger.HumanLogger + monitor *logger.Monitor + eventLog *logger.EventLogger + lastAlert map[byte]time.Time // для предотвращения спама алертов +} + +func NewPipeReader( + buf *adapter.RingBuffer, + dataLogger *logger.RotatingLogger, + humanLogger *logger.HumanLogger, + monitor *logger.Monitor, + eventLog *logger.EventLogger, +) *PipeReader { + return &PipeReader{ + buf: buf, + dataLogger: dataLogger, + humanLogger: humanLogger, + monitor: monitor, + eventLog: eventLog, + lastAlert: make(map[byte]time.Time), + } +} + +func (pr *PipeReader) Start(ctx context.Context, pipePath string) { + go func() { + buffer := make([]byte, 4096) + + for { + select { + case <-ctx.Done(): + return + default: + } + + // Проверяем существует ли pipe + if _, err := os.Stat(pipePath); os.IsNotExist(err) { + if pr.eventLog != nil { + pr.eventLog.Event("PIPE_НЕ_НАЙДЕН") + } + time.Sleep(2 * time.Second) + continue + } + + f, err := os.OpenFile(pipePath, os.O_RDONLY, 0) + if err != nil { + if pr.eventLog != nil { + pr.eventLog.Event("ОШИБКА_ОТКРЫТИЯ_PIPE") + } + time.Sleep(time.Second) + continue + } + + if pr.eventLog != nil { + pr.eventLog.Event("PIPE_ПОДКЛЮЧЕН") + } + + for { + n, err := f.Read(buffer) + if err != nil { + f.Close() + if pr.eventLog != nil { + pr.eventLog.Event("PIPE_ОТКЛЮЧЕН") + } + break + } + + // Обрабатываем каждый байт + for i := 0; i < n; i++ { + b := buffer[i] + now := time.Now() + + // 1. В RAM для UI + pr.buf.Write(b) + + // 2. На диск для архива (бинарный) + if pr.dataLogger != nil { + pr.dataLogger.Write(logger.Sample{ + Timestamp: now.UnixMicro(), + Value: b, + }) + } + + // 3. Human-readable лог (только когда есть данные) + if pr.humanLogger != nil { + humanLine := fmt.Sprintf("[%s] Значение GPIO: %d (0x%02X)\n", + now.Format("2006-01-02 15:04:05.000"), + b, b) + pr.humanLogger.Write(humanLine) + } + + // 4. Алерт если значение > 10 + if b > 10 { + pr.handleAlert(b, now) + } + + // 5. Обновляем watchdog + if pr.monitor != nil { + pr.monitor.RecordWrite() + } + } + } + + // Пауза перед переподключением + time.Sleep(1 * time.Second) + } + }() +} + +func (pr *PipeReader) handleAlert(value byte, timestamp time.Time) { + // Anti-spam: не чаще 1 алерта в секунду для одного значения + if last, exists := pr.lastAlert[value]; exists { + if timestamp.Sub(last) < 1*time.Second { + return + } + } + pr.lastAlert[value] = timestamp + + // 1. В event log (JSON) + if pr.eventLog != nil { + pr.eventLog.Event(fmt.Sprintf("ВЫСОКОЕ_ЗНАЧЕНИЕ_%d", value)) + } + + // 2. В human-readable лог с предупреждением (только на русском) + if pr.humanLogger != nil { + alertLine := fmt.Sprintf("[%s] ВНИМАНИЕ: Обнаружено высокое значение! GPIO = %d (>10)\n", + timestamp.Format("2006-01-02 15:04:05.000"), value) + pr.humanLogger.Write(alertLine) + } + + // 3. В консоль больше НЕ пишем (всё уже в логах) +} \ No newline at end of file