// Package reportbuf accumulates telemetry reports in memory // and periodically flushes them to zstd-compressed JSONL files. package reportbuf import ( "bytes" "context" "encoding/json" "errors" "fmt" "io/fs" "log/slog" "os" "path/filepath" "strings" "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 // Report files are named filePrefix + timestamp + fileSuffix. filePrefix = "reports-" fileSuffix = ".jsonl.zst" ) // ErrFull is returned by Append when storing the report would // take the report files past the configured maximum size. 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 } // 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 maxBytes int64 mu sync.Mutex 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, } 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. b.usedBytes, err = reportFilesSize(b.dataDir) if err != nil { return 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. It stores nothing and returns ErrFull if the line // would take usedBytes past maxBytes. 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() 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 } // 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") path := filepath.Join(b.dataDir, filePrefix+ts+fileSuffix) // 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) } 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. 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.mu.Unlock() return nil } // reportFilesSize returns the total size of the report files in // dir. func reportFilesSize(dir string) (int64, error) { entries, err := os.ReadDir(dir) if err != nil { return 0, fmt.Errorf("read data dir: %w", err) } var total int64 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 0, fmt.Errorf("stat report file: %w", err) } total += info.Size() } return total, nil }