// Package reportbuf accumulates telemetry reports in memory // and periodically flushes them to zstd-compressed JSONL files. package reportbuf import ( "bytes" "context" "encoding/json" "fmt" "io/fs" "log/slog" "os" "path/filepath" "sync" "time" "sneak.berlin/go/netwatch/internal/config" "sneak.berlin/go/netwatch/internal/logger" "github.com/klauspost/compress/zstd" "go.uber.org/fx" ) const ( flushSizeThreshold = 10 << 20 // 10 MiB flushInterval = 1 * time.Minute defaultDataDir = "./data/reports" dirPerms fs.FileMode = 0o750 filePerms fs.FileMode = 0o640 ) // Params defines the dependencies for Buffer. type Params struct { fx.In Config *config.Config Logger *logger.Logger } // Buffer accumulates JSON lines in memory and flushes them // to zstd-compressed files on disk. type Buffer struct { buf bytes.Buffer dataDir string done chan struct{} log *slog.Logger mu sync.Mutex stopOnce sync.Once } // New creates a Buffer and registers lifecycle hooks to // manage the data directory and flush goroutine. func New( lc fx.Lifecycle, params Params, ) (*Buffer, error) { dir := params.Config.DataDir if dir == "" { dir = defaultDataDir } b := &Buffer{ dataDir: dir, done: make(chan struct{}), log: params.Logger.Get(), } lc.Append(fx.Hook{ OnStart: func(_ context.Context) error { err := os.MkdirAll(b.dataDir, dirPerms) if err != nil { return fmt.Errorf("create data dir: %w", err) } go b.flushLoop() return nil }, OnStop: func(_ context.Context) error { // stopOnce makes OnStop idempotent: a second // invocation must not close an already-closed channel // (which would panic) or flush again. var err error b.stopOnce.Do(func() { close(b.done) err = b.flushLocked() }) // A failed final flush fails the stop, so the process // exits non-zero. return err }, }) return b, nil } // Append marshals v as a single JSON line and appends it to // the buffer. If the buffer reaches the size threshold, it is // drained and written to disk asynchronously. func (b *Buffer) Append(v any) error { line, err := json.Marshal(v) if err != nil { return fmt.Errorf("marshal report: %w", err) } b.mu.Lock() b.buf.Write(line) b.buf.WriteByte('\n') if b.buf.Len() >= flushSizeThreshold { data := b.drainBuf() b.mu.Unlock() go func() { writeErr := b.writeFile(data) if writeErr != nil { b.log.Error("flush reports failed", "error", writeErr) } }() return nil } b.mu.Unlock() return nil } // flushLoop runs a ticker that periodically flushes buffered // data to disk until the done channel is closed. func (b *Buffer) flushLoop() { ticker := time.NewTicker(flushInterval) defer ticker.Stop() for { select { case <-ticker.C: err := b.flushLocked() if err != nil { b.log.Error("flush reports failed", "error", err) } case <-b.done: return } } } // flushLocked acquires the lock, drains the buffer, and // writes the data to a compressed file. func (b *Buffer) flushLocked() error { b.mu.Lock() if b.buf.Len() == 0 { b.mu.Unlock() return nil } data := b.drainBuf() b.mu.Unlock() return b.writeFile(data) } // drainBuf copies the buffer contents and resets it. // The caller must hold b.mu. func (b *Buffer) drainBuf() []byte { data := make([]byte, b.buf.Len()) copy(data, b.buf.Bytes()) b.buf.Reset() return data } // writeFile creates a timestamped zstd-compressed JSONL file // in the data directory. func (b *Buffer) writeFile(data []byte) error { ts := time.Now().UTC().Format("2006-01-02T15-04-05.000Z") name := fmt.Sprintf("reports-%s.jsonl.zst", ts) path := filepath.Join(b.dataDir, name) // path is built from the operator-supplied dataDir plus a // generated timestamp, so it carries no external input. f, err := os.OpenFile( //nolint:gosec // see comment above path, os.O_WRONLY|os.O_CREATE|os.O_EXCL, filePerms, ) if err != nil { return fmt.Errorf("create report file: %w", err) } // Closes the file on the early returns below. The success // path closes it explicitly to check the error; closing it // a second time here is harmless. defer func() { _ = f.Close() }() enc, err := zstd.NewWriter(f) if err != nil { return fmt.Errorf("create zstd encoder: %w", err) } _, err = enc.Write(data) if err != nil { _ = enc.Close() return fmt.Errorf("write compressed data: %w", err) } err = enc.Close() if err != nil { return fmt.Errorf("close zstd encoder: %w", err) } err = f.Close() if err != nil { return fmt.Errorf("close report file: %w", err) } return nil }