106 lines
2.7 KiB
Go
106 lines
2.7 KiB
Go
package logger
|
||
|
||
import (
|
||
"encoding/binary"
|
||
"os"
|
||
"sync"
|
||
"time"
|
||
)
|
||
|
||
// DataLogger пишет сэмплы в бинарный файл с буферизацией: буфер 64 КБ
|
||
// сбрасывается раз в секунду либо при заполнении, после сброса — fsync.
|
||
type DataLogger struct {
|
||
file *os.File
|
||
mu sync.Mutex
|
||
buffer []byte
|
||
bufferSize int
|
||
flushTick *time.Ticker
|
||
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)
|
||
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
|
||
}
|
||
|
||
// Write добавляет сэмпл в буфер в формате 9 байт:
|
||
// uint64 little-endian (микросекунды Unix) + 1 байт значения.
|
||
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
|
||
}
|
||
}
|
||
}
|
||
|
||
// Close останавливает цикл сброса, дописывает остаток буфера и закрывает файл.
|
||
func (d *DataLogger) Close() error {
|
||
close(d.closeCh)
|
||
d.flushTick.Stop()
|
||
return d.file.Close()
|
||
} |