Доработка документации

This commit is contained in:
Maxim
2026-07-17 15:57:05 +03:00
parent 5b403d7ee4
commit 9fa3d172a8
22 changed files with 1371 additions and 481 deletions

View File

@@ -1,3 +1,5 @@
// Package adapter содержит кольцевой буфер последних принятых байт GPIO —
// источник данных для HTTP API (/api/latest, /api/history, /api/health).
package adapter
import (
@@ -6,6 +8,9 @@ import (
"time"
)
// RingBuffer — потокобезопасный кольцевой буфер байт фиксированного размера.
// Помимо самих данных ведёт статистику приёма: время последней записи,
// суммарные счётчики и текущую скорость (байт/с и бит/с).
type RingBuffer struct {
mu sync.RWMutex
data []byte
@@ -34,7 +39,7 @@ type timestampedByte struct {
timestamp time.Time
}
// =====================
// NewRingBuffer создаёт буфер на size элементов.
func NewRingBuffer(size int) *RingBuffer {
rb := &RingBuffer{
data: make([]byte, size),
@@ -48,7 +53,8 @@ func NewRingBuffer(size int) *RingBuffer {
return rb
}
// =====================
// Write добавляет байт в буфер (при переполнении затирается самый старый)
// и обновляет статистику: счётчики, время последней записи, скорость.
func (rb *RingBuffer) Write(b byte) {
rb.mu.Lock()
@@ -98,7 +104,8 @@ func (rb *RingBuffer) updateSpeedPeriodically(now time.Time) {
rb.lastCalcBytes = currentTotal
}
// Скользящее окно
// addTimestampedByte ведёт скользящее окно временных меток за последнюю секунду
// и пересчитывает по нему скорость (перетирая значение периодического механизма).
func (rb *RingBuffer) addTimestampedByte(now time.Time) {
rb.recentMu.Lock()
defer rb.recentMu.Unlock()
@@ -131,7 +138,8 @@ func (rb *RingBuffer) addTimestampedByte(now time.Time) {
}
}
// Получить скорость на основе скользящего окна (альтернативный метод)
// GetWindowSpeed возвращает скорость приёма, рассчитанную по скользящему окну
// в 1 секунду. Альтернатива значениям из Stats; текущим API не используется.
func (rb *RingBuffer) GetWindowSpeed() (bytesPerSec, bitsPerSec float64) {
rb.recentMu.Lock()
defer rb.recentMu.Unlock()
@@ -161,12 +169,12 @@ func (rb *RingBuffer) GetWindowSpeed() (bytesPerSec, bitsPerSec float64) {
return
}
// =====================
// LastWriteTime возвращает время последней записи в буфер.
func (rb *RingBuffer) LastWriteTime() time.Time {
return time.Unix(0, rb.lastWrite.Load())
}
// =====================
// GetLatest возвращает последний записанный байт; false — если буфер пуст.
func (rb *RingBuffer) GetLatest() (byte, bool) {
rb.mu.RLock()
defer rb.mu.RUnlock()
@@ -183,7 +191,7 @@ func (rb *RingBuffer) GetLatest() (byte, bool) {
return rb.data[idx], true
}
// =====================
// GetLast возвращает до n последних байт в порядке от старых к новым.
func (rb *RingBuffer) GetLast(n int) []byte {
rb.mu.RLock()
defer rb.mu.RUnlock()
@@ -207,12 +215,13 @@ func (rb *RingBuffer) GetLast(n int) []byte {
return res
}
// =====================
// IsAlive сообщает, была ли запись в буфер за последние timeout.
func (rb *RingBuffer) IsAlive(timeout time.Duration) bool {
return time.Since(time.Unix(0, rb.lastWrite.Load())) < timeout
}
// =====================
// Stats возвращает статистику буфера для /api/health. Ключи: size, filled,
// last_write, total_bytes, total_bits, bytes_per_sec, bits_per_sec.
func (rb *RingBuffer) Stats() map[string]interface{} {
rb.mu.RLock()
filled := rb.count

View File

@@ -1,3 +1,7 @@
// Package audio измеряет уровень звука с микрофона через arecord (ALSA):
// USB-устройство находится автоматически, PCM-поток читается блоками по 100 мс,
// уровень считается как RMS и публикуется по шкале 0100 (поле sound_level
// в /api/health). Падение arecord не фатально — захват перезапускается.
package audio
import (
@@ -15,6 +19,7 @@ import (
"time"
)
// Monitor — фоновый измеритель уровня звука с микрофона.
type Monitor struct {
mu sync.RWMutex
currentLevel int
@@ -22,16 +27,21 @@ type Monitor struct {
wg sync.WaitGroup
}
// NewMonitor создаёт монитор; захват запускается отдельным вызовом Start.
func NewMonitor() *Monitor {
return &Monitor{}
}
// GetLevel возвращает текущий уровень звука по шкале 0100
// (RMS/32768 × 100 × 4 с отсечкой сверху); 0 — если захват не работает.
func (m *Monitor) GetLevel() int {
m.mu.RLock()
defer m.mu.RUnlock()
return m.currentLevel
}
// Start запускает захват звука в фоне. Ошибка (микрофон недоступен) не
// критична для вызывающего кода — сервер продолжает работать без звука.
func (m *Monitor) Start() error {
ctx, cancel := context.WithCancel(context.Background())
m.cancel = cancel
@@ -156,6 +166,8 @@ func (m *Monitor) handleCrash(ctx context.Context) {
}
}
// findUSBDevice выбирает ALSA-устройство по выводу arecord -l: сначала USB-аудио,
// затем любое устройство захвата, иначе "default".
func (m *Monitor) findUSBDevice() string {
cmd := exec.Command("arecord", "-l")
output, err := cmd.Output()
@@ -201,6 +213,7 @@ func (m *Monitor) findUSBDevice() string {
return "default"
}
// Stop останавливает захват и дожидается завершения горутины чтения.
func (m *Monitor) Stop() {
if m.cancel != nil {
m.cancel()

View File

@@ -1,3 +1,7 @@
// Package logger реализует хранение данных GPIO: разбор байта, бинарные
// почасовые логи с ротацией, человекочитаемые и событийные логи, watchdog
// тишины и retention (автоочистку старых файлов). Пути хранения зависят от
// режима работы (deb-пакет или разработка) и вычисляются в paths.go.
package logger
import "time"

View File

@@ -7,6 +7,8 @@ import (
"time"
)
// DataLogger пишет сэмплы в бинарный файл с буферизацией: буфер 64 КБ
// сбрасывается раз в секунду либо при заполнении, после сброса — fsync.
type DataLogger struct {
file *os.File
mu sync.Mutex
@@ -16,11 +18,14 @@ type DataLogger struct {
closeCh chan struct{}
}
// Sample — одно измерение: метка времени в микросекундах Unix и сырой байт GPIO.
type Sample struct {
Timestamp int64 // microseconds
Value byte
}
// NewDataLogger открывает файл basePath на дозапись и запускает фоновый
// цикл сброса буфера. Именование файлов и смену часа обеспечивает RotatingLogger.
func NewDataLogger(basePath string) (*DataLogger, error) {
// путь будет формироваться через rotation
f, err := os.OpenFile(basePath, os.O_CREATE|os.O_APPEND|os.O_WRONLY, 0644)
@@ -40,6 +45,8 @@ func NewDataLogger(basePath string) (*DataLogger, error) {
return d, nil
}
// Write добавляет сэмпл в буфер в формате 9 байт:
// uint64 little-endian (микросекунды Unix) + 1 байт значения.
func (d *DataLogger) Write(s Sample) {
d.mu.Lock()
defer d.mu.Unlock()
@@ -91,6 +98,7 @@ func (d *DataLogger) flushLoop() {
}
}
// Close останавливает цикл сброса, дописывает остаток буфера и закрывает файл.
func (d *DataLogger) Close() error {
close(d.closeCh)
d.flushTick.Stop()

View File

@@ -8,12 +8,16 @@ import (
"time"
)
// EventLogger пишет системные события параллельно в два файла: events.log
// (JSON-строки для машинной обработки) и events_human.log (текст; его читает
// /api/log/events). Словарь событий — docs/data-formats.md.
type EventLogger struct {
file *os.File
humanFile *os.File
mu sync.Mutex
}
// NewEventLogger открывает оба файла событий на дозапись.
func NewEventLogger() (*EventLogger, error) {
// Получаем пути к файлам
jsonPath, err := GetEventLogPath()
@@ -45,6 +49,7 @@ func NewEventLogger() (*EventLogger, error) {
}, nil
}
// Event записывает событие name в оба файла с текущим временем и sync.
func (l *EventLogger) Event(name string) {
l.mu.Lock()
defer l.mu.Unlock()
@@ -74,6 +79,7 @@ func (l *EventLogger) Event(name string) {
l.humanFile.Sync()
}
// Close закрывает оба файла событий.
func (l *EventLogger) Close() error {
l.file.Close()
l.humanFile.Close()

View File

@@ -5,11 +5,14 @@ import (
"sync"
)
// HumanLogger пишет человекочитаемый лог GPIO (gpio_human.log) с немедленным
// sync после каждой строки. Частоту записей регулирует вызывающий код (PipeReader).
type HumanLogger struct {
file *os.File
mu sync.Mutex
}
// NewHumanLogger открывает gpio_human.log на дозапись.
func NewHumanLogger() (*HumanLogger, error) {
path, err := GetGPIOLogPath()
if err != nil {
@@ -24,6 +27,7 @@ func NewHumanLogger() (*HumanLogger, error) {
return &HumanLogger{file: f}, nil
}
// Write дописывает готовую строку в лог.
func (l *HumanLogger) Write(data string) {
l.mu.Lock()
defer l.mu.Unlock()
@@ -32,6 +36,7 @@ func (l *HumanLogger) Write(data string) {
l.file.Sync()
}
// Close закрывает файл лога.
func (l *HumanLogger) Close() error {
return l.file.Close()
}

View File

@@ -5,6 +5,9 @@ import (
"time"
)
// Monitor — watchdog потока данных: PipeReader отмечает каждую запись через
// RecordWrite, фоновый цикл сравнивает длительность тишины с порогами и пишет
// события ТИШИНА_5МИН/ТИШИНА_10МИН.
type Monitor struct {
eventLogger *EventLogger
lastWrite time.Time
@@ -14,6 +17,7 @@ type Monitor struct {
watchdogInterval time.Duration
}
// NewMonitor создаёт watchdog с порогами из config и сразу запускает фоновый цикл.
func NewMonitor(eventLogger *EventLogger, config *Config) *Monitor {
m := &Monitor{
eventLogger: eventLogger,
@@ -34,6 +38,7 @@ func (m *Monitor) SetSilenceThresholds(silence5Min, silence10Min time.Duration)
m.silence10Min = silence10Min
}
// RecordWrite отмечает факт приёма данных (сбрасывает отсчёт тишины).
func (m *Monitor) RecordWrite() {
m.mu.Lock()
m.lastWrite = time.Now()

View File

@@ -1,5 +1,7 @@
package logger
// GPIOData — разобранный байт GPIO-шины: количество обнаружений (биты 05),
// сила сигнала (биты 67) и признак превышения порога тревоги.
type GPIOData struct {
RawValue byte // сырое значение
Count byte // количество обнаружений (биты 0-5)

View File

@@ -14,6 +14,9 @@ import (
// возраст или размер говорят удалить. Защита от полного стирания директории.
const minKeepFiles = 2
// Retention — фоновая FIFO-очистка бинарных логов по двум независимым лимитам:
// возрасту (maxAgeHours) и суммарному размеру (maxSizeBytes). Любой лимит
// отключается нулём; минимум minKeepFiles свежих файлов сохраняется всегда.
type Retention struct {
maxAgeHours int
maxSizeBytes int64
@@ -21,6 +24,7 @@ type Retention struct {
stopCh chan struct{}
}
// NewRetention создаёт очистку с лимитами; запуск — Start или StartWithInterval.
func NewRetention(maxAgeHours int, maxSizeMB int, eventLogger *EventLogger) *Retention {
return &Retention{
maxAgeHours: maxAgeHours,
@@ -30,11 +34,13 @@ func NewRetention(maxAgeHours int, maxSizeMB int, eventLogger *EventLogger) *Ret
}
}
// Start запускает очистку со стандартным интервалом 15 минут.
func (r *Retention) Start() {
// Запускаем проверку каждые 15 минут
r.StartWithInterval(15 * time.Minute)
}
// StartWithInterval запускает фоновый цикл очистки с заданным интервалом;
// первая очистка выполняется через 1 минуту после запуска.
func (r *Retention) StartWithInterval(interval time.Duration) {
// Запускаем проверку с указанным интервалом
ticker := time.NewTicker(interval)
@@ -56,10 +62,13 @@ func (r *Retention) StartWithInterval(interval time.Duration) {
})
}
// Stop останавливает фоновый цикл очистки.
func (r *Retention) Stop() {
close(r.stopCh)
}
// Cleanup выполняет один проход очистки: сначала по возрасту, затем по размеру.
// Удаления фиксируются суммарным событием ОЧИСТКАОГОВ_УДАЛЕНО_N_ФАЙЛОВ_M_MB.
func (r *Retention) Cleanup() {
// Если оба лимита отключены, ничего не делаем
if r.maxAgeHours <= 0 && r.maxSizeBytes <= 0 {

View File

@@ -7,6 +7,9 @@ import (
"time"
)
// RotatingLogger — обёртка над DataLogger с почасовой ротацией: держит файл
// текущего часа (gpio-YYYY-MM-DD-HH.bin) и по таймеру переоткрывает его при
// смене часа, фиксируя событие ПЛАНОВАЯ_РОТАЦИЯ.
type RotatingLogger struct {
dataLogger *DataLogger
currentHour int
@@ -16,6 +19,8 @@ type RotatingLogger struct {
checkInterval time.Duration
}
// NewRotatingLogger открывает файл текущего часа в каталоге бинарных данных
// и запускает цикл проверки ротации (checkInterval; 0 — раз в минуту).
func NewRotatingLogger(eventLogger *EventLogger, checkInterval time.Duration) (*RotatingLogger, error) {
// Получаем директорию для бинарных данных
dataDir, err := GetDataLogsDir()
@@ -83,6 +88,7 @@ func (r *RotatingLogger) rotate() error {
return nil
}
// Write передаёт сэмпл текущему DataLogger.
func (r *RotatingLogger) Write(s Sample) {
r.mu.Lock()
logger := r.dataLogger
@@ -102,6 +108,7 @@ func (r *RotatingLogger) rotationLoop() {
}
}
// Close закрывает текущий DataLogger (с финальным сбросом буфера).
func (r *RotatingLogger) Close() error {
r.mu.Lock()
defer r.mu.Unlock()

View File

@@ -1,3 +1,5 @@
// Package pipe читает поток байт GPIO из именованного канала (FIFO) и раздаёт
// их потребителям: кольцевому буферу, бинарному и текстовым логгерам, watchdog.
package pipe
import (
@@ -10,6 +12,10 @@ import (
"gpio-monitor/internal/logger"
)
// PipeReader — единственный «производитель» данных в системе: следит за FIFO,
// переподключается при обрывах и раздаёт каждый принятый байт потребителям.
// Ведёт машину состояний пайпа (unknown/found/not_found/error/disconnected)
// и пишет события смены состояния без дублей.
type PipeReader struct {
buf *adapter.RingBuffer
dataLogger *logger.RotatingLogger
@@ -31,6 +37,8 @@ type PipeReader struct {
lastLoggedState string // для отслеживания изменений
}
// NewPipeReader создаёт читателя с настройками анти-спама из config.
// Любой из потребителей (dataLogger, humanLogger, monitor, eventLog) может быть nil.
func NewPipeReader(
buf *adapter.RingBuffer,
dataLogger *logger.RotatingLogger,
@@ -63,6 +71,9 @@ func (pr *PipeReader) logStateChange(newState string, eventName string) {
}
}
// Start запускает чтение pipePath в отдельной горутине и сразу возвращается.
// Если FIFO отсутствует — ждёт его появления (проверка каждые 2 с); при обрыве
// чтения переоткрывает файл через 1 с. Останавливается по отмене ctx.
func (pr *PipeReader) Start(ctx context.Context, pipePath string) {
go func() {
buffer := make([]byte, 4096)
@@ -184,6 +195,8 @@ func (pr *PipeReader) Start(ctx context.Context, pipePath string) {
}()
}
// handleAlert пишет алерт превышения (count > 10) в event- и human-логи,
// не чаще одного раза в alertCooldownSec для одинакового значения count.
func (pr *PipeReader) handleAlert(data logger.GPIOData, timestamp time.Time) {
// Anti-spam: не чаще 1 алерта в N секунд для одинакового количества
key := data.Count