// Package reportbuf accumulates telemetry reports in memory // and periodically flushes them to zstd-compressed JSONL files, // deleting the oldest files to keep them under a size cap. package reportbuf import ( "bytes" "context" "encoding/json" "errors" "fmt" "io/fs" "log/slog" "os" "path/filepath" "strings" "sync" "sync/atomic" "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 // Report files are named filePrefix + timestamp + "-" + number + // fileSuffix; see writeFile. filePrefix = "reports-" fileSuffix = ".jsonl.zst" ) // ErrFull is returned by Append when the reports not yet written // leave no room for the report under the configured maximum size, // however many report files are deleted. var ErrFull = errors.New("report files at their size cap") // Params defines the dependencies for Buffer. type Params struct { fx.In Config *config.Config Logger *logger.Logger } // reportFile is a report file that may be deleted to make room, with // the size it counts for in usedBytes. type reportFile struct { name string size int64 } // 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{} // files are the report files that may be deleted to make room, // oldest first: those in dataDir at start, then each one this // buffer writes, once it is complete. A file still being written // is not among them. filesBytes is their total size. files []reportFile filesBytes int64 log *slog.Logger maxBytes int64 mu sync.Mutex // now is the clock report files are named by: time.Now, except // in tests that need two flushes to share a timestamp. now func() time.Time // seq numbers the report files, so that two named in the same // millisecond still get different names. seq atomic.Uint64 stopOnce sync.Once // usedBytes is what Append checks against maxBytes: the size // of the report files in dataDir, plus the reports not yet // written to one at their uncompressed size. usedBytes int64 } // 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(), maxBytes: params.Config.DataDirMaxBytes, now: time.Now, } 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) } // Report files left by earlier runs count too, and are // the first to be deleted to make room. files, err := reportFiles(b.dataDir) if err != nil { return err } b.mu.Lock() b.files = files for _, f := range files { b.filesBytes += f.size } b.usedBytes = b.filesBytes // The files may be past the cap, if it was lowered since // the last run. b.deleteOldestFiles(0) b.mu.Unlock() 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 line would take usedBytes past maxBytes, the // oldest report files are deleted to make room; it stores nothing // and returns ErrFull if that cannot make room. 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) } lineBytes := int64(len(line)) + 1 // with its newline b.mu.Lock() b.deleteOldestFiles(lineBytes) if b.usedBytes+lineBytes > b.maxBytes { b.mu.Unlock() return ErrFull } b.usedBytes += lineBytes 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 } // deleteOldestFiles deletes report files, oldest first, until n more // bytes fit under maxBytes. It deletes none when the reports not yet // written leave no room for n even with every file gone, since that // would lose the files for nothing. The caller must hold b.mu. func (b *Buffer) deleteOldestFiles(n int64) { for b.usedBytes+n > b.maxBytes && len(b.files) > 0 { if b.usedBytes-b.filesBytes+n > b.maxBytes { return } f := b.files[0] b.files = b.files[1:] b.filesBytes -= f.size // A file already gone, deleted by hand, has freed its room too. err := os.Remove(filepath.Join(b.dataDir, f.name)) if err != nil && !errors.Is(err, fs.ErrNotExist) { // The file is still there, so it still counts. It is // not tried again until the next start. b.log.Error("delete report file failed", "file", f.name, "error", err) continue } b.usedBytes -= f.size b.log.Info("deleted report file to make room", "file", f.name, "bytes", f.size) } } // 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 { // The timestamp comes first, so the names sort by time; the number // after it tells apart files named in the same millisecond. ts := b.now().UTC().Format("2006-01-02T15-04-05.000Z") name := fmt.Sprintf("%s%s-%d%s", filePrefix, ts, b.seq.Add(1), fileSuffix) path := filepath.Join(b.dataDir, name) // path is built from the operator-supplied dataDir plus a // generated timestamp and number, 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) } info, err := f.Stat() if err != nil { return fmt.Errorf("stat report file: %w", err) } err = f.Close() if err != nil { return fmt.Errorf("close report file: %w", err) } // The reports counted at their uncompressed size while they // waited; now they count as the file, which from here on may be // deleted to make room. After a failed write they stay counted as // they were, which errs toward refusing reports early rather than // letting the files pass the cap. b.mu.Lock() b.usedBytes += info.Size() - int64(len(data)) b.files = append(b.files, reportFile{name: name, size: info.Size()}) b.filesBytes += info.Size() b.mu.Unlock() return nil } // reportFiles returns the report files in dir, oldest first: // os.ReadDir sorts them by name, and the names sort by time. func reportFiles(dir string) ([]reportFile, error) { entries, err := os.ReadDir(dir) if err != nil { return nil, fmt.Errorf("read data dir: %w", err) } files := make([]reportFile, 0, len(entries)) for _, entry := range entries { name := entry.Name() if !strings.HasPrefix(name, filePrefix) || !strings.HasSuffix(name, fileSuffix) { continue } info, err := entry.Info() if err != nil { return nil, fmt.Errorf("stat report file: %w", err) } files = append(files, reportFile{name: name, size: info.Size()}) } return files, nil }