From c88c5da0632b0718941a3ac293a9624fc36092ab Mon Sep 17 00:00:00 2001 From: clawbot <35+clawbot@noreply.example.org> Date: Tue, 29 Sep 2026 02:28:52 +0000 Subject: [PATCH] fix(backend): give each report file a name of its own (closes #61) 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 counts the files written since the server started, so names still sort by time and never repeat. The new test flushes pairs until one falls within one millisecond; the 1 ms pauses earlier tests used to dodge the collision are gone. Model: opus-5-5 --- TODO.md | 5 + backend/README.md | 9 +- backend/internal/reportbuf/reportbuf.go | 14 ++- backend/internal/reportbuf/reportbuf_test.go | 102 +++++++++++++++++-- 4 files changed, 116 insertions(+), 14 deletions(-) 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()