check / check (push) Successful in 1m57s
When a report would take the report files past DATA_DIR_MAX_BYTES, reportbuf now deletes the oldest report files until it fits, and does the same at start when files left by an earlier run are already past it. A file joins the files that may be deleted, at its place by name, only once it is completely written, so a file still being written is never deleted. A report is refused with 507 only when the reports waiting to be written fill the cap on their own, and then no file is deleted. The reports of a failed write stop counting, and the part of its file written is removed. A file whose deletion fails keeps counting; one already deleted by hand counts as freed. Model: opus-5-5
420 lines
10 KiB
Go
420 lines
10 KiB
Go
// 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"
|
|
"slices"
|
|
"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 waiting to be
|
|
// written leave no room for the report under the configured maximum
|
|
// size, however many report files are deleted.
|
|
var ErrFull = errors.New("reports waiting to be written fill the 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{}
|
|
// fileCreated is called with each report file once it is
|
|
// created, before anything is written to it: it does nothing,
|
|
// except in tests that hold the write open or make it fail.
|
|
fileCreated func(f *os.File)
|
|
// files are the report files that may be deleted to make room,
|
|
// in name order, which is oldest first: those in dataDir at
|
|
// start, and each one this buffer writes, put in at its place by
|
|
// name once it is complete, even when an older file's write
|
|
// completes after a newer one's. 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 waiting to
|
|
// be 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{}),
|
|
fileCreated: func(*os.File) {},
|
|
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 waiting
|
|
// to be 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 writes data, reports drained from the buffer, to a new
|
|
// 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)
|
|
|
|
size, err := b.createFile(filepath.Join(b.dataDir, name), data)
|
|
|
|
// The reports no longer wait to be written, so they stop counting
|
|
// at their uncompressed size. If the write failed they are lost;
|
|
// otherwise they count as the file, which from here on may be
|
|
// deleted to make room.
|
|
b.mu.Lock()
|
|
defer b.mu.Unlock()
|
|
|
|
b.usedBytes -= int64(len(data))
|
|
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
b.usedBytes += size
|
|
b.filesBytes += size
|
|
|
|
// At its place by name, not at the end: another write, of a newer
|
|
// file, may have completed while this one was being written.
|
|
i, _ := slices.BinarySearchFunc(b.files, name,
|
|
func(f reportFile, target string) int {
|
|
return strings.Compare(f.name, target)
|
|
})
|
|
b.files = slices.Insert(b.files, i, reportFile{name: name, size: size})
|
|
|
|
return nil
|
|
}
|
|
|
|
// createFile creates the file at path holding data compressed with
|
|
// zstd, and returns its size. If the write fails once the file is
|
|
// created, it removes the file, so that a failed write leaves nothing
|
|
// behind to take room.
|
|
func (b *Buffer) createFile(path string, data []byte) (int64, error) {
|
|
// 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 0, fmt.Errorf("create report file: %w", err)
|
|
}
|
|
|
|
b.fileCreated(f)
|
|
|
|
size, err := writeCompressed(f, data)
|
|
if err != nil {
|
|
_ = f.Close()
|
|
|
|
return 0, errors.Join(err, os.Remove(path))
|
|
}
|
|
|
|
err = f.Close()
|
|
if err != nil {
|
|
err = fmt.Errorf("close report file: %w", err)
|
|
|
|
return 0, errors.Join(err, os.Remove(path))
|
|
}
|
|
|
|
return size, nil
|
|
}
|
|
|
|
// writeCompressed writes data to f compressed with zstd, and returns
|
|
// the size of f.
|
|
func writeCompressed(f *os.File, data []byte) (int64, error) {
|
|
enc, err := zstd.NewWriter(f)
|
|
if err != nil {
|
|
return 0, fmt.Errorf("create zstd encoder: %w", err)
|
|
}
|
|
|
|
_, err = enc.Write(data)
|
|
if err != nil {
|
|
_ = enc.Close()
|
|
|
|
return 0, fmt.Errorf("write compressed data: %w", err)
|
|
}
|
|
|
|
err = enc.Close()
|
|
if err != nil {
|
|
return 0, fmt.Errorf("close zstd encoder: %w", err)
|
|
}
|
|
|
|
info, err := f.Stat()
|
|
if err != nil {
|
|
return 0, fmt.Errorf("stat report file: %w", err)
|
|
}
|
|
|
|
return info.Size(), 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
|
|
}
|