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) { // 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() 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() paths, err := filepath.Glob(filepath.Join(dir, "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(paths)) for _, path := range paths { compressed, readErr := os.ReadFile(path) if readErr != nil { t.Fatalf("read %s: %v", path, readErr) } data, decErr := dec.DecodeAll(compressed, nil) if decErr != nil { t.Fatalf("decompress %s: %v", path, 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 }