Доработка. Журналирование

This commit is contained in:
Maxim
2026-06-09 11:00:20 +03:00
parent ac0e3944ab
commit 354156e557
7 changed files with 537 additions and 44 deletions

View File

@@ -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")

View File

@@ -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()
}

View File

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

View File

@@ -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()
}

View File

@@ -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МИН")
}
}
}

View File

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

146
internal/pipe/reader.go Normal file
View File

@@ -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. В консоль больше НЕ пишем (всё уже в логах)
}