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
This commit was merged in pull request #66.
This commit is contained in:
2026-09-29 08:05:26 +02:00
parent d2f219ca19
commit 6022cc8b02
5 changed files with 116 additions and 15 deletions
+77 -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,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.
func reportFilesBytes(t *testing.T, dir string) int64 {
t.Helper()
@@ -338,6 +370,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()