fix(backend): give each report file a name of its own (closes #61)
check / check (push) Successful in 1m7s
check / check (push) Successful in 1m7s
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 buffer reads the time through a clock the new test stops, so its two flushes share one timestamp on every run; the 1 ms pauses earlier tests used to dodge the collision are gone. Model: opus-5-5
This commit is contained in:
@@ -1,7 +1,15 @@
|
||||
package reportbuf
|
||||
|
||||
import "time"
|
||||
|
||||
// Flush writes the buffered reports to a file now, as the periodic
|
||||
// flush does, so tests need not wait a minute for it.
|
||||
func (b *Buffer) Flush() error {
|
||||
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 }
|
||||
}
|
||||
|
||||
@@ -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,12 @@ type Buffer struct {
|
||||
log *slog.Logger
|
||||
maxBytes int64
|
||||
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
|
||||
// usedBytes is what Append checks against maxBytes: the size
|
||||
// of the report files in dataDir, plus the reports not yet
|
||||
@@ -79,6 +87,7 @@ func New(
|
||||
done: make(chan struct{}),
|
||||
log: params.Logger.Get(),
|
||||
maxBytes: params.Config.DataDirMaxBytes,
|
||||
now: time.Now,
|
||||
}
|
||||
|
||||
lc.Append(fx.Hook{
|
||||
@@ -211,11 +220,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 {
|
||||
ts := time.Now().UTC().Format("2006-01-02T15-04-05.000Z")
|
||||
path := filepath.Join(b.dataDir, filePrefix+ts+fileSuffix)
|
||||
// The timestamp comes first, so the names sort by time; the number
|
||||
// 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
|
||||
// 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,
|
||||
|
||||
@@ -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()
|
||||
|
||||
|
||||
Reference in New Issue
Block a user