diff --git a/TODO.md b/TODO.md index adbc424..a217ce0 100644 --- a/TODO.md +++ b/TODO.md @@ -23,6 +23,11 @@ latest run passes. # Completed Steps +- 2026-09-29: report file names can no longer collide (issue #61): each is + `reports--.jsonl.zst`, where the number counts the files + written since the server started, 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 - 2026-09-29: bounded the report endpoint (issue #20): `POST /api/v1/reports` still needs no credentials, but each client address, as resolved through `TRUSTED_PROXIES`, may send `REPORTS_PER_MINUTE` (default 60) reports a diff --git a/backend/README.md b/backend/README.md index 095af11..85dfebc 100644 --- a/backend/README.md +++ b/backend/README.md @@ -102,9 +102,12 @@ which `netwatch` owns. ### Report storage -Reports are written as `reports-.jsonl.zst` files in `DATA_DIR`. -Each file contains one JSON object per line, compressed with zstd. Files are -created with `O_EXCL` to prevent overwrites. +Reports are written as `reports--.jsonl.zst` files in +`DATA_DIR`. The timestamp is in UTC to the millisecond, so the names sort by +time; the number counts the files the server has written since it started, so +two files written in the same millisecond still get different names. Each file +contains one JSON object per line, compressed with zstd. Files are created with +`O_EXCL` to prevent overwrites. ### Report limits diff --git a/backend/internal/reportbuf/reportbuf.go b/backend/internal/reportbuf/reportbuf.go index 318b77f..5e19f85 100644 --- a/backend/internal/reportbuf/reportbuf.go +++ b/backend/internal/reportbuf/reportbuf.go @@ -14,6 +14,7 @@ import ( "path/filepath" "strings" "sync" + "sync/atomic" "time" "sneak.berlin/go/netwatch/internal/config" @@ -30,7 +31,8 @@ const ( dirPerms fs.FileMode = 0o750 filePerms fs.FileMode = 0o640 - // Report files are named filePrefix + timestamp + fileSuffix. + // Report files are named filePrefix + timestamp + "-" + number + + // fileSuffix; see writeFile. filePrefix = "reports-" fileSuffix = ".jsonl.zst" ) @@ -56,6 +58,9 @@ type Buffer struct { log *slog.Logger maxBytes int64 mu sync.Mutex + // 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 not yet @@ -211,11 +216,14 @@ func (b *Buffer) drainBuf() []byte { // writeFile creates a 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 := time.Now().UTC().Format("2006-01-02T15-04-05.000Z") - path := filepath.Join(b.dataDir, filePrefix+ts+fileSuffix) + 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 - // 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 path, os.O_WRONLY|os.O_CREATE|os.O_EXCL, diff --git a/backend/internal/reportbuf/reportbuf_test.go b/backend/internal/reportbuf/reportbuf_test.go index 90d3a94..5926201 100644 --- a/backend/internal/reportbuf/reportbuf_test.go +++ b/backend/internal/reportbuf/reportbuf_test.go @@ -3,9 +3,11 @@ package reportbuf_test import ( "encoding/json" "errors" + "fmt" "io/fs" "os" "path/filepath" + "slices" "strconv" "strings" "sync" @@ -18,6 +20,7 @@ import ( "sneak.berlin/go/netwatch/internal/logger" "sneak.berlin/go/netwatch/internal/reportbuf" + "github.com/klauspost/compress/zstd" "go.uber.org/fx" "go.uber.org/fx/fxtest" ) @@ -210,10 +213,6 @@ func TestWrittenReportsCountAtFileSize(t *testing.T) { 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) if err != nil { 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) } - // 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() if err != nil { t.Fatalf("flush: %v", err) @@ -315,6 +310,60 @@ 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) { + // A pair of flushes may straddle a millisecond, which proves + // nothing, so pairs are flushed until one falls within one. + const maxPairs = 1000 + + dir := t.TempDir() + t.Setenv("DATA_DIR", dir) + + buf := startBuffer(t) + flushes := 0 + + for pair := 1; ; pair++ { + start := time.Now().Truncate(time.Millisecond) + + for range 2 { + flushes++ + + err := buf.Append(map[string]int{"id": flushes}) + if err != nil { + t.Fatalf("append report %d: %v", flushes, err) + } + + err = buf.Flush() + if err != nil { + t.Fatalf("flush %d: %v", flushes, err) + } + } + + if time.Now().Truncate(time.Millisecond).Equal(start) { + break + } + + if pair == maxPairs { + t.Fatalf("no pair of flushes fell within one millisecond "+ + "in %d tries", maxPairs) + } + } + + 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. func reportFilesBytes(t *testing.T, dir string) int64 { t.Helper() @@ -338,6 +387,43 @@ func reportFilesBytes(t *testing.T, dir string) int64 { 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) { t.Helper()