// 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) } }