2 Commits
Author SHA1 Message Date
clawbot a911353023 nginx: trust X-Forwarded-For only from TRUSTED_PROXIES (closes #64)
check / check (push) Successful in 52s
nginx trusted X-Forwarded-For from every RFC1918 address, so a client
reaching it from one could write a new address on each request and
get a fresh rate-limit allowance. The container's TRUSTED_PROXIES now
names the reverse proxies nginx trusts, none by default.
bin/entrypoint.sh makes each entry a CIDR, checks it with the new
"netwatch-server check-cidr", which runs the server's own
TRUSTED_PROXIES parsing, and writes one set_real_ip_from line per
entry into /etc/nginx/trusted-proxies.conf, which nginx.conf includes.
The backend is started with TRUSTED_PROXIES=127.0.0.1/32, since nginx
is its only client. The viewport test mounts an empty file there.

Model: opus-5-5
2026-09-29 06:22:56 +00:00
clawbot 6022cc8b02 fix(backend): give each report file a name of its own (closes #61)
check / check (push) Successful in 13s
Report files were named by a millisecond timestamp and created with
O_EXCL, so two flushes in the same millisecond, such as a flush for
size and the final flush at shutdown, got the same name and the second
failed, losing its reports. Each name now carries a number after the
timestamp that goes up by one for each file the server starts to
write, so names still sort by time and never repeat within a run. A
failed write uses up its number, leaving a gap if the file could not
be created and otherwise a file under that number that may be
incomplete.

Model: opus-5-5
2026-09-29 08:05:26 +02:00
5 changed files with 116 additions and 15 deletions
+7
View File
@@ -32,6 +32,13 @@ latest run passes.
an entry that is not an IP address or CIDR, as `netwatch-server check-cidr` an entry that is not an IP address or CIDR, as `netwatch-server check-cidr`
finds; it starts the backend with `TRUSTED_PROXIES=127.0.0.1/32`, since nginx finds; it starts the backend with `TRUSTED_PROXIES=127.0.0.1/32`, since nginx
is its only client is its only client
- 2026-09-29: report file names can no longer collide (issue #61): each is
`reports-<timestamp>-<number>.jsonl.zst`, where the number goes up by one for
each file the server starts to write, so two flushes in the same millisecond,
such as a flush for size and the final flush at shutdown, each get a file of
their own instead of the second one failing. A failed write uses up its
number, leaving a gap if the file could not be created and otherwise a file
under that number that may be incomplete.
- 2026-09-29: ready to run under upaas (issue #59): the image has a - 2026-09-29: ready to run under upaas (issue #59): the image has a
`HEALTHCHECK` that requests `/.well-known/healthcheck` through nginx on the `HEALTHCHECK` that requests `/.well-known/healthcheck` through nginx on the
port from `PORT`. The backend no longer reads a bad `PORT` as 0 or a bad port from `PORT`. The backend no longer reads a bad `PORT` as 0 or a bad
+8 -3
View File
@@ -118,9 +118,14 @@ it as this server parses its own `TRUSTED_PROXIES`.
### Report storage ### Report storage
Reports are written as `reports-<timestamp>.jsonl.zst` files in `DATA_DIR`. Reports are written as `reports-<timestamp>-<number>.jsonl.zst` files in
Each file contains one JSON object per line, compressed with zstd. Files are `DATA_DIR`. The timestamp is in UTC to the millisecond, so the names sort by
created with `O_EXCL` to prevent overwrites. time. The number starts at 1 when the server starts and goes up by one for each
file the server starts to write, so two files written in the same millisecond
still get different names. A failed write uses up its number, leaving a gap in
the numbers if the file could not be created and otherwise a file under that
number that may be incomplete. Each file contains one JSON object per line,
compressed with zstd. Files are created with `O_EXCL` to prevent overwrites.
### Report limits ### Report limits
@@ -1,7 +1,15 @@
package reportbuf package reportbuf
import "time"
// Flush writes the buffered reports to a file now, as the periodic // Flush writes the buffered reports to a file now, as the periodic
// flush does, so tests need not wait a minute for it. // flush does, so tests need not wait a minute for it.
func (b *Buffer) Flush() error { func (b *Buffer) Flush() error {
return b.flushLocked() return b.flushLocked()
} }
// StopClock makes every report file the buffer writes from now on
// carry the timestamp at, as if all were written in one millisecond.
func (b *Buffer) StopClock(at time.Time) {
b.now = func() time.Time { return at }
}
+16 -4
View File
@@ -14,6 +14,7 @@ import (
"path/filepath" "path/filepath"
"strings" "strings"
"sync" "sync"
"sync/atomic"
"time" "time"
"sneak.berlin/go/netwatch/internal/config" "sneak.berlin/go/netwatch/internal/config"
@@ -30,7 +31,8 @@ const (
dirPerms fs.FileMode = 0o750 dirPerms fs.FileMode = 0o750
filePerms fs.FileMode = 0o640 filePerms fs.FileMode = 0o640
// Report files are named filePrefix + timestamp + fileSuffix. // Report files are named filePrefix + timestamp + "-" + number +
// fileSuffix; see writeFile.
filePrefix = "reports-" filePrefix = "reports-"
fileSuffix = ".jsonl.zst" fileSuffix = ".jsonl.zst"
) )
@@ -56,6 +58,12 @@ type Buffer struct {
log *slog.Logger log *slog.Logger
maxBytes int64 maxBytes int64
mu sync.Mutex 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 stopOnce sync.Once
// usedBytes is what Append checks against maxBytes: the size // usedBytes is what Append checks against maxBytes: the size
// of the report files in dataDir, plus the reports not yet // of the report files in dataDir, plus the reports not yet
@@ -79,6 +87,7 @@ func New(
done: make(chan struct{}), done: make(chan struct{}),
log: params.Logger.Get(), log: params.Logger.Get(),
maxBytes: params.Config.DataDirMaxBytes, maxBytes: params.Config.DataDirMaxBytes,
now: time.Now,
} }
lc.Append(fx.Hook{ lc.Append(fx.Hook{
@@ -211,11 +220,14 @@ func (b *Buffer) drainBuf() []byte {
// writeFile creates a timestamped zstd-compressed JSONL file // writeFile creates a timestamped zstd-compressed JSONL file
// in the data directory. // in the data directory.
func (b *Buffer) writeFile(data []byte) error { func (b *Buffer) writeFile(data []byte) error {
ts := time.Now().UTC().Format("2006-01-02T15-04-05.000Z") // The timestamp comes first, so the names sort by time; the number
path := filepath.Join(b.dataDir, filePrefix+ts+fileSuffix) // 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 // path is built from the operator-supplied dataDir plus a
// generated timestamp, so it carries no external input. // generated timestamp and number, so it carries no external input.
f, err := os.OpenFile( //nolint:gosec // see comment above f, err := os.OpenFile( //nolint:gosec // see comment above
path, path,
os.O_WRONLY|os.O_CREATE|os.O_EXCL, os.O_WRONLY|os.O_CREATE|os.O_EXCL,
+77 -8
View File
@@ -3,9 +3,11 @@ package reportbuf_test
import ( import (
"encoding/json" "encoding/json"
"errors" "errors"
"fmt"
"io/fs" "io/fs"
"os" "os"
"path/filepath" "path/filepath"
"slices"
"strconv" "strconv"
"strings" "strings"
"sync" "sync"
@@ -18,6 +20,7 @@ import (
"sneak.berlin/go/netwatch/internal/logger" "sneak.berlin/go/netwatch/internal/logger"
"sneak.berlin/go/netwatch/internal/reportbuf" "sneak.berlin/go/netwatch/internal/reportbuf"
"github.com/klauspost/compress/zstd"
"go.uber.org/fx" "go.uber.org/fx"
"go.uber.org/fx/fxtest" "go.uber.org/fx/fxtest"
) )
@@ -210,10 +213,6 @@ func TestWrittenReportsCountAtFileSize(t *testing.T) {
t.Fatalf("flush: %v", err) t.Fatalf("flush: %v", err)
} }
// The second report is written at shutdown, and must not land in
// the first file's millisecond (see TestWrittenReportsKeepCounting).
time.Sleep(time.Millisecond)
err = buf.Append(report) err = buf.Append(report)
if err != nil { if err != nil {
t.Fatalf("second report, after the first was written: %v", err) t.Fatalf("second report, after the first was written: %v", err)
@@ -254,10 +253,6 @@ func TestWrittenReportsKeepCounting(t *testing.T) {
t.Fatalf("with %d bytes of report files: %v", used, err) t.Fatalf("with %d bytes of report files: %v", used, err)
} }
// Report files are named to the millisecond; two in the same
// one collide (https://git.eeqj.de/sneak/netwatch/issues/61).
time.Sleep(time.Millisecond)
err = buf.Flush() err = buf.Flush()
if err != nil { if err != nil {
t.Fatalf("flush: %v", err) t.Fatalf("flush: %v", err)
@@ -315,6 +310,43 @@ func TestConcurrentAppendsStopAtCap(t *testing.T) {
} }
} }
// TestTwoFlushesInOneMillisecond flushes twice within one millisecond,
// as a flush for size and the final flush at shutdown can: each flush
// must write a file of its own, and the files must hold every report.
func TestTwoFlushesInOneMillisecond(t *testing.T) {
const flushes = 2
dir := t.TempDir()
t.Setenv("DATA_DIR", dir)
buf := startBuffer(t)
buf.StopClock(time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC))
for id := 1; id <= flushes; id++ {
err := buf.Append(map[string]int{"id": id})
if err != nil {
t.Fatalf("append report %d: %v", id, err)
}
err = buf.Flush()
if err != nil {
t.Fatalf("flush %d: %v", id, err)
}
}
files := readReportFiles(t, dir)
if len(files) != flushes {
t.Fatalf("%d report files after %d flushes", len(files), flushes)
}
for id := 1; id <= flushes; id++ {
want := fmt.Sprintf(`{"id":%d}`+"\n", id)
if !slices.Contains(files, want) {
t.Fatalf("no report file holds report %d alone", id)
}
}
}
// reportFilesBytes returns the total size of the report files in dir. // reportFilesBytes returns the total size of the report files in dir.
func reportFilesBytes(t *testing.T, dir string) int64 { func reportFilesBytes(t *testing.T, dir string) int64 {
t.Helper() t.Helper()
@@ -338,6 +370,43 @@ func reportFilesBytes(t *testing.T, dir string) int64 {
return total return total
} }
// readReportFiles returns the decompressed contents of each report
// file in dir.
func readReportFiles(t *testing.T, dir string) []string {
t.Helper()
files := os.DirFS(dir)
names, err := fs.Glob(files, "reports-*.jsonl.zst")
if err != nil {
t.Fatalf("list report files: %v", err)
}
dec, err := zstd.NewReader(nil)
if err != nil {
t.Fatalf("create zstd decoder: %v", err)
}
defer dec.Close()
contents := make([]string, 0, len(names))
for _, name := range names {
compressed, readErr := fs.ReadFile(files, name)
if readErr != nil {
t.Fatalf("read %s: %v", name, readErr)
}
data, decErr := dec.DecodeAll(compressed, nil)
if decErr != nil {
t.Fatalf("decompress %s: %v", name, decErr)
}
contents = append(contents, string(data))
}
return contents
}
func writeBytes(t *testing.T, path string, n int) { func writeBytes(t *testing.T, path string, n int) {
t.Helper() t.Helper()