Files
go-service/internal/pipe/reader.go

179 lines
4.4 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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
// Для дедупликации логов
lastData logger.GPIOData
lastLogTime time.Time
lastAlertTime 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,
lastAlertTime: 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:
}
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++ {
rawByte := buffer[i]
now := time.Now()
// Парсим данные
data := logger.ParseGPIO(rawByte)
// 1. В RAM для UI (сырое значение)
pr.buf.Write(rawByte)
// 2. На диск для архива (бинарный)
if pr.dataLogger != nil {
pr.dataLogger.Write(logger.Sample{
Timestamp: now.UnixMicro(),
Value: rawByte,
})
}
// 3. Human-readable лог (при изменении состояния или раз в 5 секунд)
if pr.humanLogger != nil {
shouldLog := false
// Логируем если изменилось количество или сила
if data.Count != pr.lastData.Count || data.Strength != pr.lastData.Strength {
shouldLog = true
}
// Или если прошло больше 5 секунд с последнего лога
if now.Sub(pr.lastLogTime) >= 5*time.Second {
shouldLog = true
}
if shouldLog {
strengthName := logger.GetStrengthName(data.Strength)
humanLine := fmt.Sprintf(
"[%s] Обнаружение: %d объектов | Сила: %s (уровень %d) | Сырое: 0x%02X (%d)\n",
now.Format("2006-01-02 15:04:05.000"),
data.Count,
strengthName,
data.Strength,
rawByte,
rawByte,
)
pr.humanLogger.Write(humanLine)
pr.lastData = data
pr.lastLogTime = now
}
}
// 4. Алерт если количество обнаружений > 10
if data.IsHighCount {
pr.handleAlert(data, now)
}
// 5. Обновляем watchdog
if pr.monitor != nil {
pr.monitor.RecordWrite()
}
}
}
time.Sleep(1 * time.Second)
}
}()
}
func (pr *PipeReader) handleAlert(data logger.GPIOData, timestamp time.Time) {
// Anti-spam: не чаще 1 алерта в 10 секунд для одинакового количества
key := data.Count
if last, exists := pr.lastAlertTime[key]; exists {
if timestamp.Sub(last) < 10*time.Second {
return
}
}
pr.lastAlertTime[key] = timestamp
strengthName := logger.GetStrengthName(data.Strength)
// 1. В event log (JSON)
if pr.eventLog != nil {
pr.eventLog.Event(fmt.Sprintf("ОБНАРУЖЕНО_%d_ОБЪЕКТОВ_СИЛА_%d", data.Count, data.Strength))
}
// 2. В human-readable лог с предупреждением
if pr.humanLogger != nil {
alertLine := fmt.Sprintf(
"[%s] ВНИМАНИЕ: Обнаружено превышение! %d объектов | Сила: %s (уровень %d)\n",
timestamp.Format("2006-01-02 15:04:05.000"),
data.Count,
strengthName,
data.Strength,
)
pr.humanLogger.Write(alertLine)
}
}