Files
go-service/internal/logger/data_logger.go

98 lines
1.8 KiB
Go

package logger
import (
"encoding/binary"
"os"
"sync"
"time"
)
type DataLogger struct {
file *os.File
mu sync.Mutex
buffer []byte
bufferSize int
flushTick *time.Ticker
closeCh chan struct{}
}
type Sample struct {
Timestamp int64 // microseconds
Value byte
}
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
}
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
}
}
}
func (d *DataLogger) Close() error {
close(d.closeCh)
d.flushTick.Stop()
return d.file.Close()
}