Files
go-service/internal/pipe/reader.go
2026-07-17 15:57:05 +03:00

228 lines
7.2 KiB
Go
Raw Permalink 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 читает поток байт GPIO из именованного канала (FIFO) и раздаёт
// их потребителям: кольцевому буферу, бинарному и текстовым логгерам, watchdog.
package pipe
import (
"context"
"fmt"
"os"
"time"
"gpio-monitor/internal/adapter"
"gpio-monitor/internal/logger"
)
// PipeReader — единственный «производитель» данных в системе: следит за FIFO,
// переподключается при обрывах и раздаёт каждый принятый байт потребителям.
// Ведёт машину состояний пайпа (unknown/found/not_found/error/disconnected)
// и пишет события смены состояния без дублей.
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
// Настройки
alertCooldownSec time.Duration
humanLogIntervalSec time.Duration
// Флаги состояния
pipeState string // "unknown", "found", "not_found"
lastLoggedState string // для отслеживания изменений
}
// NewPipeReader создаёт читателя с настройками анти-спама из config.
// Любой из потребителей (dataLogger, humanLogger, monitor, eventLog) может быть nil.
func NewPipeReader(
buf *adapter.RingBuffer,
dataLogger *logger.RotatingLogger,
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),
alertCooldownSec: config.AlertCooldownSec,
humanLogIntervalSec: config.HumanLogIntervalSec,
pipeState: "unknown",
lastLoggedState: "unknown",
}
}
func (pr *PipeReader) logStateChange(newState string, eventName string) {
if newState != pr.lastLoggedState {
if pr.eventLog != nil {
pr.eventLog.Event(eventName)
}
pr.lastLoggedState = newState
}
}
// Start запускает чтение pipePath в отдельной горутине и сразу возвращается.
// Если FIFO отсутствует — ждёт его появления (проверка каждые 2 с); при обрыве
// чтения переоткрывает файл через 1 с. Останавливается по отмене ctx.
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.pipeState != "not_found" {
pr.pipeState = "not_found"
pr.logStateChange("not_found", "PIPE_НЕ_НАЙДЕН")
}
time.Sleep(2 * time.Second)
continue
}
f, err := os.OpenFile(pipePath, os.O_RDONLY, 0)
if err != nil {
if pr.pipeState != "error" {
pr.pipeState = "error"
pr.logStateChange("error", "ОШИБКА_ОТКРЫТИЯ_PIPE")
}
time.Sleep(time.Second)
continue
}
// Pipe успешно открыт
if pr.pipeState != "found" {
pr.pipeState = "found"
pr.logStateChange("found", "PIPE_ПОДКЛЮЧЕН")
}
for {
n, err := f.Read(buffer)
if err != nil {
f.Close()
// Отключаемся только если были подключены
if pr.pipeState == "found" {
pr.pipeState = "disconnected"
pr.logStateChange("disconnected", "PIPE_ОТКЛЮЧЕН")
}
break
}
// Если были в состоянии ошибки/отключения, восстанавливаемся
if pr.pipeState != "found" {
pr.pipeState = "found"
pr.logStateChange("found", "PIPE_ПОДКЛЮЧЕН")
}
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 лог (при изменении состояния или раз в N секунд)
if pr.humanLogger != nil {
shouldLog := false
// Логируем если изменилось количество или сила
if data.Count != pr.lastData.Count || data.Strength != pr.lastData.Strength {
shouldLog = true
}
// Или если прошло больше N секунд с последнего лога
if now.Sub(pr.lastLogTime) >= pr.humanLogIntervalSec {
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)
}
}()
}
// handleAlert пишет алерт превышения (count > 10) в event- и human-логи,
// не чаще одного раза в alertCooldownSec для одинакового значения count.
func (pr *PipeReader) handleAlert(data logger.GPIOData, timestamp time.Time) {
// Anti-spam: не чаще 1 алерта в N секунд для одинакового количества
key := data.Count
if last, exists := pr.lastAlertTime[key]; exists {
if timestamp.Sub(last) < pr.alertCooldownSec {
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)
}
}