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 } // TestAppendPastCapIsRefused fills the cap with a report not yet // written. The next report is refused, and the report file already in // DATA_DIR is kept: it is smaller than a report, so deleting it could // not make room. func TestAppendPastCapIsRefused(t *testing.T) { report := map[string]string{"id": "cap"} dir := t.TempDir() earlier := reportFilePath(dir, 1) writeBytes(t, earlier, 1) t.Setenv("DATA_DIR", dir) t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(1+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) } if !exists(t, earlier) { t.Fatal("report file deleted, though that could not make room") } } // TestOldestReportFileDeletedFirst starts on a data directory holding // report files from an earlier run, and a file that is not a report, // which neither counts nor is ever deleted. Nothing is deleted while // there is room; then only the oldest report file is. func TestOldestReportFileDeletedFirst(t *testing.T) { const fileBytes = 100 report := map[string]string{"id": "oldest"} dir := t.TempDir() oldest := reportFilePath(dir, 1) kept := []string{ reportFilePath(dir, 2), reportFilePath(dir, 3), filepath.Join(dir, "notes.txt"), } writeBytes(t, oldest, fileBytes) for _, path := range kept { writeBytes(t, path, fileBytes) } t.Setenv("DATA_DIR", dir) // Room for the three report files and one report. t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(3*fileBytes+lineBytes(t, report))) buf := startBuffer(t) err := buf.Append(report) if err != nil { t.Fatalf("report that fills the cap exactly: %v", err) } if !exists(t, oldest) { t.Fatal("oldest report file deleted while there was room") } err = buf.Append(report) if err != nil { t.Fatalf("report past the cap: %v", err) } if exists(t, oldest) { t.Fatal("oldest report file kept when room was needed") } for _, path := range kept { if !exists(t, path) { t.Fatalf("%s deleted; only the oldest report file should be", path) } } } // TestStartDeletesFilesPastCap starts on report files past the cap, // as after the cap is lowered: the oldest are deleted until the rest // fit. func TestStartDeletesFilesPastCap(t *testing.T) { const fileBytes = 100 dir := t.TempDir() oldest := reportFilePath(dir, 1) kept := []string{reportFilePath(dir, 2), reportFilePath(dir, 3)} writeBytes(t, oldest, fileBytes) for _, path := range kept { writeBytes(t, path, fileBytes) } t.Setenv("DATA_DIR", dir) t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(len(kept)*fileBytes)) startBuffer(t) if exists(t, oldest) { t.Fatal("oldest report file kept, though the files were past the cap") } for _, path := range kept { if !exists(t, path) { t.Fatalf("%s deleted, though the rest fit without it", path) } } } // 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) dir := t.TempDir() t.Setenv("DATA_DIR", dir) // 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) } if !hasReportFile(t, dir) { t.Fatal("first report's file deleted, though the second fit beside it") } } // TestWrittenFilesDeletedToMakeRoom writes one report file after // another under a small cap. Every report must be taken; files are // deleted only when the report would not fit beside them, and the // files kept leave room for it. func TestWrittenFilesDeletedToMakeRoom(t *testing.T) { const ( maxBytes = 200 // One file each, which take far more than maxBytes together. reports = 50 ) 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) for range reports { before := reportFilesBytes(t, dir) err := buf.Append(report) if err != nil { t.Fatalf("with %d bytes of report files: %v", before, err) } after := reportFilesBytes(t, dir) if before+size <= maxBytes && after != before { t.Fatalf("files deleted, though the report fit beside "+ "their %d bytes", before) } if after+size > maxBytes { t.Fatalf("%d bytes of report files kept, leaving no room "+ "for the report", after) } err = buf.Flush() if err != nil { t.Fatalf("flush: %v", err) } } } // TestFailedDeletionStillCounts makes deleting the oldest report file // fail. It is still there, so it still takes room, and the next oldest // is deleted in its place. func TestFailedDeletionStillCounts(t *testing.T) { const fileBytes = 100 report := map[string]string{"id": "stuck"} dir := t.TempDir() stuck := reportFilePath(dir, 1) next := reportFilePath(dir, 2) newest := reportFilePath(dir, 3) for _, path := range []string{stuck, next, newest} { writeBytes(t, path, fileBytes) } t.Setenv("DATA_DIR", dir) t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(3*fileBytes)) buf := startBuffer(t) // A directory that is not empty cannot be deleted, even by root, // which the tests run as in the backend image. err := os.Remove(stuck) if err != nil { t.Fatalf("remove %s: %v", stuck, err) } err = os.Mkdir(stuck, 0o750) if err != nil { t.Fatalf("make directory %s: %v", stuck, err) } writeBytes(t, filepath.Join(stuck, "file"), 1) err = buf.Append(report) if err != nil { t.Fatalf("report past the cap: %v", err) } if exists(t, next) { t.Fatal("next oldest report file kept: the failed deletion " + "counted as making room") } if !exists(t, newest) { t.Fatal("newest report file deleted, though deleting one made room") } } // TestFileDeletedByHandFreesRoom deletes the oldest report file by // hand after start. When room is needed, its room counts as freed, so // no other file is deleted. func TestFileDeletedByHandFreesRoom(t *testing.T) { const fileBytes = 100 report := map[string]string{"id": "by-hand"} dir := t.TempDir() gone := reportFilePath(dir, 1) kept := reportFilePath(dir, 2) writeBytes(t, gone, fileBytes) writeBytes(t, kept, fileBytes) t.Setenv("DATA_DIR", dir) t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(2*fileBytes)) buf := startBuffer(t) err := os.Remove(gone) if err != nil { t.Fatalf("remove %s: %v", gone, err) } err = buf.Append(report) if err != nil { t.Fatalf("report past the cap: %v", err) } if !exists(t, kept) { t.Fatal("report file deleted, though the one deleted by hand " + "had made room") } } // TestFileBeingWrittenIsNeverDeleted holds a write open, its file // created but not complete, while a report needs room. The file may // not be deleted to make it, so the report is refused. Once the write // is complete, the file is deleted when room is needed. func TestFileBeingWrittenIsNeverDeleted(t *testing.T) { report := map[string]string{"id": "writing"} t.Setenv("DATA_DIR", t.TempDir()) t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(lineBytes(t, report))) buf := startBuffer(t) created := make(chan string) release := make(chan struct{}) buf.OnFileCreated(func(f *os.File) { created <- f.Name() <-release }) err := buf.Append(report) if err != nil { t.Fatalf("report that fills the cap exactly: %v", err) } flushed := make(chan error) go func() { flushed <- buf.Flush() }() writing := <-created // Errorf, not Fatalf, until the write is released, so that a // failure here does not leave it held. err = buf.Append(report) if !errors.Is(err, reportbuf.ErrFull) { t.Errorf("report while the first was being written: "+ "error = %v, want ErrFull", err) } if !exists(t, writing) { t.Error("report file deleted while it was being written") } close(release) err = <-flushed if err != nil { t.Fatalf("flush: %v", err) } // Writes from here on, the final one at stop included, go // through unheld. buf.OnFileCreated(func(*os.File) {}) err = buf.Append(report) if err != nil { t.Fatalf("report after the write was complete: %v", err) } if exists(t, writing) { t.Fatal("complete report file kept when room was needed") } } // TestFailedWriteStopsCounting makes a write fail once its file is // created. Its reports are lost, so they stop counting, and the part of // the file written is removed, so it takes no room. func TestFailedWriteStopsCounting(t *testing.T) { report := map[string]string{"id": "failed"} t.Setenv("DATA_DIR", t.TempDir()) t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(lineBytes(t, report))) buf := startBuffer(t) var failed string // Closing the file under the write makes the write fail. buf.OnFileCreated(func(f *os.File) { failed = f.Name() _ = f.Close() }) err := buf.Append(report) if err != nil { t.Fatalf("report that fills the cap exactly: %v", err) } err = buf.Flush() if err == nil { t.Fatal("flush succeeded, though its file was closed under it") } if exists(t, failed) { t.Fatal("file of the failed write kept") } // Writes from here on, the final one at stop included, succeed. buf.OnFileCreated(func(*os.File) {}) err = buf.Append(report) if err != nil { t.Fatalf("report after the failed write: %v", err) } } // 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 } // reportFilePath returns the path in dir of a report file named as // written on the given day of January 2026, so that a lower day sorts // as older. func reportFilePath(dir string, day int) string { return filepath.Join(dir, fmt.Sprintf("reports-2026-01-%02dT00-00-00.000Z-1.jsonl.zst", day)) } 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 exists(t *testing.T, path string) bool { t.Helper() _, err := os.Stat(path) if errors.Is(err, fs.ErrNotExist) { return false } if err != nil { t.Fatalf("stat %s: %v", path, err) } return true } 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 }