fix(backend): rate-limit and cap report ingest, drop wildcard CORS (closes #20)
check / check (push) Successful in 1m1s
check / check (push) Successful in 1m1s
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) with ErrFull, answered with 507; the count starts from the files already in DATA_DIR, and reports not yet written count at their uncompressed size. CORS adds nothing unless CORS_ALLOWED_ORIGINS lists origins. A limit that is not a positive number stops the server from starting. Model: opus-5-5
This commit is contained in:
@@ -6,11 +6,13 @@ import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io/fs"
|
||||
"log/slog"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
@@ -27,8 +29,16 @@ const (
|
||||
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
|
||||
@@ -44,8 +54,13 @@ type Buffer struct {
|
||||
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
|
||||
@@ -60,9 +75,10 @@ func New(
|
||||
}
|
||||
|
||||
b := &Buffer{
|
||||
dataDir: dir,
|
||||
done: make(chan struct{}),
|
||||
log: params.Logger.Get(),
|
||||
dataDir: dir,
|
||||
done: make(chan struct{}),
|
||||
log: params.Logger.Get(),
|
||||
maxBytes: params.Config.DataDirMaxBytes,
|
||||
}
|
||||
|
||||
lc.Append(fx.Hook{
|
||||
@@ -72,6 +88,12 @@ func New(
|
||||
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
|
||||
@@ -97,15 +119,27 @@ func New(
|
||||
}
|
||||
|
||||
// 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.
|
||||
// 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')
|
||||
|
||||
@@ -178,8 +212,7 @@ func (b *Buffer) drainBuf() []byte {
|
||||
// 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 := 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.
|
||||
@@ -214,10 +247,51 @@ func (b *Buffer) writeFile(data []byte) error {
|
||||
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
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user