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
442 lines
9.9 KiB
Go
442 lines
9.9 KiB
Go
package reportbuf_test
|
|
|
|
import (
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"io/fs"
|
|
"os"
|
|
"path/filepath"
|
|
"slices"
|
|
"strconv"
|
|
"strings"
|
|
"sync"
|
|
"sync/atomic"
|
|
"testing"
|
|
"time"
|
|
|
|
"sneak.berlin/go/netwatch/internal/config"
|
|
"sneak.berlin/go/netwatch/internal/globals"
|
|
"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"
|
|
)
|
|
|
|
// TestFlushOnShutdown proves the flush-on-shutdown path: a
|
|
// report appended after start but before the periodic flush
|
|
// window must reach disk when the fx lifecycle stops. This is
|
|
// the exact case that silent data loss on restart used to
|
|
// destroy.
|
|
func TestFlushOnShutdown(t *testing.T) {
|
|
dir := t.TempDir()
|
|
t.Setenv("DATA_DIR", dir)
|
|
|
|
var buf *reportbuf.Buffer
|
|
|
|
app := fxtest.New(t,
|
|
fx.Provide(
|
|
globals.New,
|
|
logger.New,
|
|
config.New,
|
|
reportbuf.New,
|
|
),
|
|
fx.Populate(&buf),
|
|
)
|
|
|
|
app.RequireStart()
|
|
|
|
err := buf.Append(map[string]string{"probe": "shutdown"})
|
|
if err != nil {
|
|
t.Fatalf("append report: %v", err)
|
|
}
|
|
|
|
// RequireStop runs the reportbuf OnStop hook, which is the
|
|
// only code path that flushes buffered reports on shutdown.
|
|
app.RequireStop()
|
|
|
|
if !hasReportFile(t, dir) {
|
|
t.Fatal("no report file on disk after shutdown; " +
|
|
"the buffered report was lost")
|
|
}
|
|
}
|
|
|
|
// TestFailedFinalFlushFailsStop proves a final flush that cannot
|
|
// write its file makes the stop fail, which makes the process
|
|
// exit non-zero instead of dropping the buffered reports silently.
|
|
func TestFailedFinalFlushFailsStop(t *testing.T) {
|
|
dir := t.TempDir()
|
|
t.Setenv("DATA_DIR", dir)
|
|
|
|
var buf *reportbuf.Buffer
|
|
|
|
app := fxtest.New(t,
|
|
fx.Provide(
|
|
globals.New,
|
|
logger.New,
|
|
config.New,
|
|
reportbuf.New,
|
|
),
|
|
fx.Populate(&buf),
|
|
)
|
|
|
|
app.RequireStart()
|
|
|
|
err := buf.Append(map[string]string{"probe": "shutdown"})
|
|
if err != nil {
|
|
t.Fatalf("append report: %v", err)
|
|
}
|
|
|
|
// Removing the data directory leaves the final flush nowhere to
|
|
// write. A read-only directory would not do: tests run as root
|
|
// in the backend image, and root ignores the read-only bit.
|
|
err = os.RemoveAll(dir)
|
|
if err != nil {
|
|
t.Fatalf("remove data dir: %v", err)
|
|
}
|
|
|
|
err = app.Stop(t.Context())
|
|
if !errors.Is(err, fs.ErrNotExist) {
|
|
t.Fatalf("stop error = %v, want the final flush's error", err)
|
|
}
|
|
}
|
|
|
|
// startBuffer starts a Buffer through fx, as main does, with the
|
|
// DATA_DIR and DATA_DIR_MAX_BYTES the calling test has set.
|
|
func startBuffer(t *testing.T) *reportbuf.Buffer {
|
|
t.Helper()
|
|
|
|
var buf *reportbuf.Buffer
|
|
|
|
app := fxtest.New(t,
|
|
fx.Provide(
|
|
globals.New,
|
|
logger.New,
|
|
config.New,
|
|
reportbuf.New,
|
|
),
|
|
fx.Populate(&buf),
|
|
)
|
|
|
|
app.RequireStart()
|
|
t.Cleanup(app.RequireStop)
|
|
|
|
return buf
|
|
}
|
|
|
|
// lineBytes is what one report takes in the buffer: its JSON and a
|
|
// newline.
|
|
func lineBytes(t *testing.T, report any) int {
|
|
t.Helper()
|
|
|
|
line, err := json.Marshal(report)
|
|
if err != nil {
|
|
t.Fatalf("marshal report: %v", err)
|
|
}
|
|
|
|
return len(line) + 1
|
|
}
|
|
|
|
func TestAppendPastCapIsRefused(t *testing.T) {
|
|
report := map[string]string{"id": "cap"}
|
|
|
|
t.Setenv("DATA_DIR", t.TempDir())
|
|
t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(lineBytes(t, report)))
|
|
|
|
buf := startBuffer(t)
|
|
|
|
err := buf.Append(report)
|
|
if err != nil {
|
|
t.Fatalf("report that fills the cap exactly: %v", err)
|
|
}
|
|
|
|
err = buf.Append(report)
|
|
if !errors.Is(err, reportbuf.ErrFull) {
|
|
t.Fatalf("report past the cap: error = %v, want ErrFull", err)
|
|
}
|
|
}
|
|
|
|
// TestCapCountsReportFilesAlreadyInDataDir starts on a data
|
|
// directory holding a report file from an earlier run, and a file
|
|
// that is not a report, which must not count.
|
|
func TestCapCountsReportFilesAlreadyInDataDir(t *testing.T) {
|
|
const earlierBytes = 100
|
|
|
|
report := map[string]string{"id": "cap"}
|
|
dir := t.TempDir()
|
|
|
|
writeBytes(t, filepath.Join(dir, "reports-2026-01-01T00-00-00.000Z.jsonl.zst"),
|
|
earlierBytes)
|
|
writeBytes(t, filepath.Join(dir, "notes.txt"), 10*earlierBytes)
|
|
|
|
t.Setenv("DATA_DIR", dir)
|
|
t.Setenv("DATA_DIR_MAX_BYTES",
|
|
strconv.Itoa(earlierBytes+lineBytes(t, report)))
|
|
|
|
buf := startBuffer(t)
|
|
|
|
err := buf.Append(report)
|
|
if err != nil {
|
|
t.Fatalf("report that fills the cap exactly: %v", err)
|
|
}
|
|
|
|
err = buf.Append(report)
|
|
if !errors.Is(err, reportbuf.ErrFull) {
|
|
t.Fatalf("report past the cap: error = %v, want ErrFull", err)
|
|
}
|
|
}
|
|
|
|
// TestWrittenReportsCountAtFileSize checks that once reports are
|
|
// written, they count as their compressed file, not their
|
|
// uncompressed size, which frees room under the cap.
|
|
func TestWrittenReportsCountAtFileSize(t *testing.T) {
|
|
// Repetitive, so its file is far smaller than its JSON.
|
|
report := map[string]string{"id": strings.Repeat("a", 1000)}
|
|
size := lineBytes(t, report)
|
|
|
|
t.Setenv("DATA_DIR", t.TempDir())
|
|
// Room for the report twice over only if the first one counts
|
|
// at its file's size by the time the second arrives.
|
|
t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(2*size-1))
|
|
|
|
buf := startBuffer(t)
|
|
|
|
err := buf.Append(report)
|
|
if err != nil {
|
|
t.Fatalf("first report: %v", err)
|
|
}
|
|
|
|
err = buf.Flush()
|
|
if err != nil {
|
|
t.Fatalf("flush: %v", err)
|
|
}
|
|
|
|
err = buf.Append(report)
|
|
if err != nil {
|
|
t.Fatalf("second report, after the first was written: %v", err)
|
|
}
|
|
}
|
|
|
|
// TestWrittenReportsKeepCounting writes one report file after another
|
|
// under a small cap: each report must be taken while the files on disk
|
|
// leave room for it, and refused once they do not.
|
|
func TestWrittenReportsKeepCounting(t *testing.T) {
|
|
const maxBytes = 200
|
|
|
|
report := map[string]string{"id": "written"}
|
|
size := int64(lineBytes(t, report))
|
|
dir := t.TempDir()
|
|
|
|
t.Setenv("DATA_DIR", dir)
|
|
t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(maxBytes))
|
|
|
|
buf := startBuffer(t)
|
|
|
|
// Every file takes at least a byte, so they fill the cap within
|
|
// maxBytes rounds.
|
|
for range maxBytes {
|
|
used := reportFilesBytes(t, dir)
|
|
|
|
err := buf.Append(report)
|
|
if used+size > maxBytes {
|
|
if !errors.Is(err, reportbuf.ErrFull) {
|
|
t.Fatalf("with %d bytes of report files: error = %v, "+
|
|
"want ErrFull", used, err)
|
|
}
|
|
|
|
return
|
|
}
|
|
|
|
if err != nil {
|
|
t.Fatalf("with %d bytes of report files: %v", used, err)
|
|
}
|
|
|
|
err = buf.Flush()
|
|
if err != nil {
|
|
t.Fatalf("flush: %v", err)
|
|
}
|
|
}
|
|
|
|
t.Fatal("the report files never filled the cap")
|
|
}
|
|
|
|
// TestConcurrentAppendsStopAtCap appends from many goroutines at once
|
|
// with room for exactly roomFor reports: exactly that many must be
|
|
// taken, which holds only if Append checks and counts each report
|
|
// under one lock.
|
|
func TestConcurrentAppendsStopAtCap(t *testing.T) {
|
|
const (
|
|
roomFor = 5
|
|
senders = 50
|
|
)
|
|
|
|
// Large, so each Append takes long enough for the senders to
|
|
// overlap while the cap is reached.
|
|
report := map[string]string{"id": strings.Repeat("a", 1_000_000)}
|
|
|
|
t.Setenv("DATA_DIR", t.TempDir())
|
|
t.Setenv("DATA_DIR_MAX_BYTES",
|
|
strconv.Itoa(roomFor*lineBytes(t, report)))
|
|
|
|
buf := startBuffer(t)
|
|
|
|
var (
|
|
taken atomic.Int64
|
|
wg sync.WaitGroup
|
|
)
|
|
|
|
start := make(chan struct{})
|
|
|
|
for range senders {
|
|
wg.Go(func() {
|
|
<-start
|
|
|
|
err := buf.Append(report)
|
|
if err == nil {
|
|
taken.Add(1)
|
|
} else if !errors.Is(err, reportbuf.ErrFull) {
|
|
t.Errorf("append: %v", err)
|
|
}
|
|
})
|
|
}
|
|
|
|
close(start)
|
|
wg.Wait()
|
|
|
|
if got := taken.Load(); got != roomFor {
|
|
t.Fatalf("%d reports taken, want %d", got, roomFor)
|
|
}
|
|
}
|
|
|
|
// 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()
|
|
|
|
paths, err := filepath.Glob(filepath.Join(dir, "reports-*.jsonl.zst"))
|
|
if err != nil {
|
|
t.Fatalf("list report files: %v", err)
|
|
}
|
|
|
|
var total int64
|
|
|
|
for _, path := range paths {
|
|
info, statErr := os.Stat(path)
|
|
if statErr != nil {
|
|
t.Fatalf("stat %s: %v", path, statErr)
|
|
}
|
|
|
|
total += info.Size()
|
|
}
|
|
|
|
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()
|
|
|
|
err := os.WriteFile(path, make([]byte, n), 0o600)
|
|
if err != nil {
|
|
t.Fatalf("write %s: %v", path, err)
|
|
}
|
|
}
|
|
|
|
func hasReportFile(t *testing.T, dir string) bool {
|
|
t.Helper()
|
|
|
|
entries, err := os.ReadDir(dir)
|
|
if err != nil {
|
|
t.Fatalf("read data dir: %v", err)
|
|
}
|
|
|
|
for _, e := range entries {
|
|
if strings.HasSuffix(e.Name(), ".jsonl.zst") {
|
|
info, statErr := e.Info()
|
|
if statErr != nil {
|
|
t.Fatalf("stat %s: %v", e.Name(), statErr)
|
|
}
|
|
|
|
if info.Size() > 0 {
|
|
return true
|
|
}
|
|
}
|
|
}
|
|
|
|
return false
|
|
}
|