check / check (push) Successful in 42s
POST /api/v1/reports stays unauthenticated but is bounded. Each client address, as the trusted-proxy logic resolves it, may send REPORTS_PER_MINUTE reports a minute (default 60, counted by go-chi/httprate over a sliding minute); past that it gets 429 with Retry-After. reportbuf refuses a report that would take the report files past DATA_DIR_MAX_BYTES (default 1 GiB), counting the files already in DATA_DIR and unwritten reports at their uncompressed size; the handler answers 507. CORS adds nothing unless CORS_ALLOWED_ORIGINS lists origins. A limit that is not a positive number, or an origin that is not a plain scheme://host[:port], stops the server from starting. Model: opus-5-5
298 lines
6.5 KiB
Go
298 lines
6.5 KiB
Go
// 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
|
|
}
|