fix(backend): give each report file a name of its own (closes #61)
check / check (push) Successful in 32s

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
This commit is contained in:
2026-09-29 03:25:11 +00:00
parent ced1956b06
commit 0987817aae
4 changed files with 116 additions and 14 deletions
+5
View File
@@ -23,6 +23,11 @@ latest run passes.
# Completed Steps
- 2026-09-29: report file names can no longer collide (issue #61): each is
`reports-<timestamp>-<number>.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: nginx listens on `PORT` (issue #26), 8080 when unset or empty: the
nginx image renders `nginx.conf` as a template at container start, filling in
`PORT` and no other variable. `bin/entrypoint.sh` refuses to start when `PORT`
+6 -3
View File
@@ -102,9 +102,12 @@ which `netwatch` owns.
### Report storage
Reports are written as `reports-<timestamp>.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-<timestamp>-<number>.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
+11 -3
View File
@@ -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,
+94 -8
View File
@@ -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()