From 114c31777f0a0b798e71e5a333411eca3a602087 Mon Sep 17 00:00:00 2001 From: sneak Date: Sat, 3 Oct 2026 13:21:42 +0000 Subject: [PATCH] Delete the oldest report files to stay under the size cap (closes #54) When a report would take the report files past DATA_DIR_MAX_BYTES, reportbuf now deletes the oldest report files until it fits, and does the same at start when files left by an earlier run are already past it. A file joins the files that may be deleted only once it is completely written, so a file still being written is never deleted. A report is refused with 507 only when the reports waiting to be written fill the cap on their own, and then no file is deleted. The reports of a failed write stop counting, and the part of its file written is removed. A file whose deletion fails keeps counting; one already deleted by hand counts as freed. Model: opus-5-5 --- README.md | 2 +- TODO.md | 9 + backend/README.md | 28 +- backend/internal/handlers/report.go | 7 +- backend/internal/handlers/report_test.go | 6 +- backend/internal/reportbuf/export_test.go | 12 +- backend/internal/reportbuf/reportbuf.go | 212 ++++++++--- backend/internal/reportbuf/reportbuf_test.go | 379 +++++++++++++++++-- 8 files changed, 546 insertions(+), 109 deletions(-) diff --git a/README.md b/README.md index e604338..d8fbffa 100644 --- a/README.md +++ b/README.md @@ -220,7 +220,7 @@ What the [upaas](https://git.eeqj.de/sneak/upaas) app for netwatch needs: - `REPORTS_PER_MINUTE`, default `60`: reports each client address may send a minute - `DATA_DIR_MAX_BYTES`, default `1073741824` (1 GiB): the most room the - report files may take + report files may take; the oldest are deleted to stay under it - `CORS_ALLOWED_ORIGINS`, default empty: other origins whose pages may call the API - `DEBUG`, default `false`: debug logging diff --git a/TODO.md b/TODO.md index 536f37f..3d6e60a 100644 --- a/TODO.md +++ b/TODO.md @@ -23,6 +23,15 @@ latest run passes. # Completed Steps +- 2026-10-03: `DATA_DIR_MAX_BYTES` is now how much of the report files is kept + (issue #54): when a report would take them past it, the oldest report files + are deleted to make room, each deletion logged, and at start files already + past it are deleted the same way. A file still being written is never deleted. + A report is refused with 507 only when the reports waiting to be written fill + the cap on their own, and then no file is deleted. The reports of a failed + write stop counting, and the part of its file written is removed. A file that + cannot be deleted still counts until the next start; one already deleted by + hand counts as freed - 2026-10-03: the Go tests run with the race detector and coverage (issue #88): `backend/script/test` runs `go test -timeout 30s -race -cover ./...` and, if that fails, runs it again with `-v` and fails. Go's `-timeout` bounds the diff --git a/backend/README.md b/backend/README.md index 0161cc4..42494ac 100644 --- a/backend/README.md +++ b/backend/README.md @@ -83,7 +83,7 @@ Internal packages in `internal/` follow standard Go project layout: | `BIND_ADDRESS` | empty | IP address to listen on; empty listens on every interface | | `PORT` | `8080` | HTTP listen port | | `DATA_DIR` | `./data/reports` | Directory for compressed reports | -| `DATA_DIR_MAX_BYTES` | `1073741824` (1 GiB) | Largest total size of the report files in `DATA_DIR`; see [Report limits](#report-limits) | +| `DATA_DIR_MAX_BYTES` | `1073741824` (1 GiB) | Most bytes of report files kept in `DATA_DIR`, oldest deleted first; see [Report limits](#report-limits) | | `DEBUG` | `false` | Enable debug logging | | `TRUSTED_PROXIES` | loopback + RFC1918 | Comma-separated CIDRs whose `X-Forwarded-For` / `X-Real-IP` headers are trusted for client IP resolution | | `REPORTS_PER_MINUTE` | `60` | Reports each client address may send a minute; see [Report limits](#report-limits) | @@ -128,10 +128,11 @@ Reports are written as `reports--.jsonl.zst` files in `DATA_DIR`. The timestamp is in UTC to the millisecond, so the names sort by time. The number starts at 1 when the server starts and goes up by one for each file the server starts to write, so two files written in the same millisecond -still get different names. A failed write uses up its number, leaving a gap in -the numbers if the file could not be created and otherwise a file under that -number that may be incomplete. Each file contains one JSON object per line, -compressed with zstd. Files are created with `O_EXCL` to prevent overwrites. +still get different names. A failed write uses up its number and leaves a gap in +the numbers: its file, if it was created, is removed. The file stays, counted +toward `DATA_DIR_MAX_BYTES` from the next start, only if removing it fails too. +Each file contains one JSON object per line, compressed with zstd. Files are +created with `O_EXCL` to prevent overwrites. ### Report limits @@ -151,11 +152,17 @@ credentials, so it is bounded instead. Both refusals below answer with the same `X-RateLimit-Reset` headers. - **Size cap.** The report files in `DATA_DIR` may total at most `DATA_DIR_MAX_BYTES`, counting the files already there at start. Reports - waiting in memory count at their uncompressed size until they are written, so - a report that would take the total past the cap is refused with 507, and - nothing of it is stored. Deleting report files frees room only at the next - start, when the files are counted again. The default of 1 GiB is small enough - for any host; set it to the space you can give `DATA_DIR`. + waiting in memory count at their uncompressed size until they are written; + those lost to a failed write stop counting, and the part of its file written + is removed. When a report would take the total past the cap, the oldest report + files are deleted to make room, and each deletion is logged with the file's + name and size; a file still being written is never deleted. A report is + refused with 507, and nothing of it is stored, only when the reports waiting + to be written fill the cap on their own, and then no file is deleted. At + start, report files past the cap, as after lowering it, are deleted the same + way. So the cap is how much of the newest reports is kept: the default of 1 + GiB is small enough for any host; set it to the space you can give + `DATA_DIR`. ### CORS @@ -173,7 +180,6 @@ starting, with an error naming `CORS_ALLOWED_ORIGINS`. - Add integration test that POSTs a report and verifies the compressed output - Add report decompression/query endpoint - Add metrics (Prometheus) for buffer size, flush count, report count -- Add retention policy to prune old report files ## License diff --git a/backend/internal/handlers/report.go b/backend/internal/handlers/report.go index e7dd4e4..8aa5ad9 100644 --- a/backend/internal/handlers/report.go +++ b/backend/internal/handlers/report.go @@ -86,11 +86,12 @@ func (s *Handlers) decodeErrorStatus(err error) int { } // appendErrorStatus logs a failure to store a report and returns -// the status to send: 507 when the report files are at their size -// cap, otherwise 500. +// the status to send: 507 when the reports waiting to be written fill +// the size cap, otherwise 500. func (s *Handlers) appendErrorStatus(err error) int { if errors.Is(err, reportbuf.ErrFull) { - s.log.Warn("report refused: report files at their size cap") + s.log.Warn("report refused: " + + "reports waiting to be written fill the size cap") return http.StatusInsufficientStorage } diff --git a/backend/internal/handlers/report_test.go b/backend/internal/handlers/report_test.go index c906b52..bd6345c 100644 --- a/backend/internal/handlers/report_test.go +++ b/backend/internal/handlers/report_test.go @@ -125,9 +125,9 @@ func TestHandleReportStorageFailureIsNon2xx(t *testing.T) { } } -// TestHandleReportFullIs507 checks the answer when the report files -// are at their size cap: 507 and the usual error body, which tells -// the client nothing more. +// TestHandleReportFullIs507 checks the answer when the reports waiting +// to be written fill the size cap: 507 and the usual error body, which +// tells the client nothing more. func TestHandleReportFullIs507(t *testing.T) { t.Parallel() diff --git a/backend/internal/reportbuf/export_test.go b/backend/internal/reportbuf/export_test.go index 845feff..4408b9a 100644 --- a/backend/internal/reportbuf/export_test.go +++ b/backend/internal/reportbuf/export_test.go @@ -1,6 +1,9 @@ package reportbuf -import "time" +import ( + "os" + "time" +) // FlushSizeThreshold exposes the buffer size at which Append starts // writing a report file to the external tests. @@ -17,3 +20,10 @@ func (b *Buffer) Flush() error { func (b *Buffer) StopClock(at time.Time) { b.now = func() time.Time { return at } } + +// OnFileCreated makes the buffer call fn with each report file it +// writes from now on, once the file is created and before anything is +// written to it. +func (b *Buffer) OnFileCreated(fn func(f *os.File)) { + b.fileCreated = fn +} diff --git a/backend/internal/reportbuf/reportbuf.go b/backend/internal/reportbuf/reportbuf.go index eeb0184..4449a1f 100644 --- a/backend/internal/reportbuf/reportbuf.go +++ b/backend/internal/reportbuf/reportbuf.go @@ -1,5 +1,6 @@ // Package reportbuf accumulates telemetry reports in memory -// and periodically flushes them to zstd-compressed JSONL files. +// and periodically flushes them to zstd-compressed JSONL files, +// deleting the oldest files to keep them under a size cap. package reportbuf import ( @@ -37,9 +38,10 @@ const ( fileSuffix = ".jsonl.zst" ) -// ErrFull is returned by Append when storing the report would -// take the report files past the configured maximum size. -var ErrFull = errors.New("report files at their size cap") +// ErrFull is returned by Append when the reports waiting to be +// written leave no room for the report under the configured maximum +// size, however many report files are deleted. +var ErrFull = errors.New("reports waiting to be written fill the size cap") // Params defines the dependencies for Buffer. type Params struct { @@ -49,15 +51,32 @@ type Params struct { Logger *logger.Logger } +// reportFile is a report file that may be deleted to make room, with +// the size it counts for in usedBytes. +type reportFile struct { + name string + size int64 +} + // Buffer accumulates JSON lines in memory and flushes them // to zstd-compressed files on disk. type Buffer struct { - buf bytes.Buffer - dataDir string - done chan struct{} - log *slog.Logger - maxBytes int64 - mu sync.Mutex + buf bytes.Buffer + dataDir string + done chan struct{} + // fileCreated is called with each report file once it is + // created, before anything is written to it: it does nothing, + // except in tests that hold the write open or make it fail. + fileCreated func(f *os.File) + // files are the report files that may be deleted to make room, + // oldest first: those in dataDir at start, then each one this + // buffer writes, once it is complete. A file still being written + // is not among them. filesBytes is their total size. + files []reportFile + filesBytes int64 + 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 @@ -66,8 +85,8 @@ type Buffer struct { 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 - // written to one at their uncompressed size. + // of the report files in dataDir, plus the reports waiting to + // be written to one at their uncompressed size. usedBytes int64 } @@ -83,11 +102,12 @@ func New( } b := &Buffer{ - dataDir: dir, - done: make(chan struct{}), - log: params.Logger.Get(), - maxBytes: params.Config.DataDirMaxBytes, - now: time.Now, + dataDir: dir, + done: make(chan struct{}), + fileCreated: func(*os.File) {}, + log: params.Logger.Get(), + maxBytes: params.Config.DataDirMaxBytes, + now: time.Now, } lc.Append(fx.Hook{ @@ -97,12 +117,28 @@ func New( return fmt.Errorf("create data dir: %w", err) } - // Report files left by earlier runs count too. - b.usedBytes, err = reportFilesSize(b.dataDir) + // Report files left by earlier runs count too, and are + // the first to be deleted to make room. + files, err := reportFiles(b.dataDir) if err != nil { return err } + b.mu.Lock() + + b.files = files + for _, f := range files { + b.filesBytes += f.size + } + + b.usedBytes = b.filesBytes + + // The files may be past the cap, if it was lowered since + // the last run. + b.deleteOldestFiles(0) + + b.mu.Unlock() + go b.flushLoop() return nil @@ -128,9 +164,10 @@ func New( } // Append marshals v as a single JSON line and appends it to -// the buffer. It stores nothing and returns ErrFull if the line -// would take usedBytes past maxBytes. If the buffer reaches the -// size threshold, it is drained and written to disk +// the buffer. If the line would take usedBytes past maxBytes, the +// oldest report files are deleted to make room; it stores nothing +// and returns ErrFull if that cannot make room. If the buffer +// reaches the size threshold, it is drained and written to disk // asynchronously. func (b *Buffer) Append(v any) error { line, err := json.Marshal(v) @@ -142,6 +179,8 @@ func (b *Buffer) Append(v any) error { b.mu.Lock() + b.deleteOldestFiles(lineBytes) + if b.usedBytes+lineBytes > b.maxBytes { b.mu.Unlock() @@ -171,6 +210,37 @@ func (b *Buffer) Append(v any) error { return nil } +// deleteOldestFiles deletes report files, oldest first, until n more +// bytes fit under maxBytes. It deletes none when the reports waiting +// to be written leave no room for n even with every file gone, since +// that would lose the files for nothing. The caller must hold b.mu. +func (b *Buffer) deleteOldestFiles(n int64) { + for b.usedBytes+n > b.maxBytes && len(b.files) > 0 { + if b.usedBytes-b.filesBytes+n > b.maxBytes { + return + } + + f := b.files[0] + b.files = b.files[1:] + b.filesBytes -= f.size + + // A file already gone, deleted by hand, has freed its room too. + err := os.Remove(filepath.Join(b.dataDir, f.name)) + if err != nil && !errors.Is(err, fs.ErrNotExist) { + // The file is still there, so it still counts. It is + // not tried again until the next start. + b.log.Error("delete report file failed", + "file", f.name, "error", err) + + continue + } + + b.usedBytes -= f.size + b.log.Info("deleted report file to make room", + "file", f.name, "bytes", f.size) + } +} + // flushLoop runs a ticker that periodically flushes buffered // data to disk until the done channel is closed. func (b *Buffer) flushLoop() { @@ -217,15 +287,41 @@ func (b *Buffer) drainBuf() []byte { return data } -// writeFile creates a timestamped zstd-compressed JSONL file -// in the data directory. +// writeFile writes data, reports drained from the buffer, to a new +// timestamped zstd-compressed JSONL file in the data directory. func (b *Buffer) writeFile(data []byte) error { // 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) + size, err := b.createFile(filepath.Join(b.dataDir, name), data) + + // The reports no longer wait to be written, so they stop counting + // at their uncompressed size. If the write failed they are lost; + // otherwise they count as the file, which from here on may be + // deleted to make room. + b.mu.Lock() + defer b.mu.Unlock() + + b.usedBytes -= int64(len(data)) + + if err != nil { + return err + } + + b.usedBytes += size + b.files = append(b.files, reportFile{name: name, size: size}) + b.filesBytes += size + + return nil +} + +// createFile creates the file at path holding data compressed with +// zstd, and returns its size. If the write fails once the file is +// created, it removes the file, so that a failed write leaves nothing +// behind to take room. +func (b *Buffer) createFile(path string, data []byte) (int64, error) { // path is built from the operator-supplied dataDir plus a // generated timestamp and number, so it carries no external input. f, err := os.OpenFile( //nolint:gosec // see comment above @@ -234,61 +330,65 @@ func (b *Buffer) writeFile(data []byte) error { filePerms, ) if err != nil { - return fmt.Errorf("create report file: %w", err) + return 0, fmt.Errorf("create report file: %w", err) } - // Closes the file on the early returns below. The success - // path closes it explicitly to check the error; closing it - // a second time here is harmless. - defer func() { _ = f.Close() }() + b.fileCreated(f) + size, err := writeCompressed(f, data) + if err != nil { + _ = f.Close() + + return 0, errors.Join(err, os.Remove(path)) + } + + err = f.Close() + if err != nil { + err = fmt.Errorf("close report file: %w", err) + + return 0, errors.Join(err, os.Remove(path)) + } + + return size, nil +} + +// writeCompressed writes data to f compressed with zstd, and returns +// the size of f. +func writeCompressed(f *os.File, data []byte) (int64, error) { enc, err := zstd.NewWriter(f) if err != nil { - return fmt.Errorf("create zstd encoder: %w", err) + return 0, fmt.Errorf("create zstd encoder: %w", err) } _, err = enc.Write(data) if err != nil { _ = enc.Close() - return fmt.Errorf("write compressed data: %w", err) + return 0, fmt.Errorf("write compressed data: %w", err) } err = enc.Close() if err != nil { - return fmt.Errorf("close zstd encoder: %w", err) + return 0, fmt.Errorf("close zstd encoder: %w", err) } info, err := f.Stat() if err != nil { - return fmt.Errorf("stat report file: %w", err) + return 0, fmt.Errorf("stat report file: %w", err) } - err = f.Close() - if err != nil { - return fmt.Errorf("close report file: %w", err) - } - - // The reports counted at their uncompressed size while they - // waited; now they count as the file. After a failed write they - // stay counted as they were, which errs toward refusing reports - // early rather than letting the files pass the cap. - b.mu.Lock() - b.usedBytes += info.Size() - int64(len(data)) - b.mu.Unlock() - - return nil + return info.Size(), nil } -// reportFilesSize returns the total size of the report files in -// dir. -func reportFilesSize(dir string) (int64, error) { +// reportFiles returns the report files in dir, oldest first: +// os.ReadDir sorts them by name, and the names sort by time. +func reportFiles(dir string) ([]reportFile, error) { entries, err := os.ReadDir(dir) if err != nil { - return 0, fmt.Errorf("read data dir: %w", err) + return nil, fmt.Errorf("read data dir: %w", err) } - var total int64 + files := make([]reportFile, 0, len(entries)) for _, entry := range entries { name := entry.Name() @@ -299,11 +399,11 @@ func reportFilesSize(dir string) (int64, error) { info, err := entry.Info() if err != nil { - return 0, fmt.Errorf("stat report file: %w", err) + return nil, fmt.Errorf("stat report file: %w", err) } - total += info.Size() + files = append(files, reportFile{name: name, size: info.Size()}) } - return total, nil + return files, nil } diff --git a/backend/internal/reportbuf/reportbuf_test.go b/backend/internal/reportbuf/reportbuf_test.go index f9b020e..d9722f9 100644 --- a/backend/internal/reportbuf/reportbuf_test.go +++ b/backend/internal/reportbuf/reportbuf_test.go @@ -139,11 +139,19 @@ func lineBytes(t *testing.T, report any) int { 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) - t.Setenv("DATA_DIR", t.TempDir()) - t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(lineBytes(t, report))) + writeBytes(t, earlier, 1) + + t.Setenv("DATA_DIR", dir) + t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(1+lineBytes(t, report))) buf := startBuffer(t) @@ -156,24 +164,38 @@ func TestAppendPastCapIsRefused(t *testing.T) { 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") + } } -// 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 +// 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": "cap"} + 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, filepath.Join(dir, "reports-2026-01-01T00-00-00.000Z.jsonl.zst"), - earlierBytes) - writeBytes(t, filepath.Join(dir, "notes.txt"), 10*earlierBytes) + 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(earlierBytes+lineBytes(t, report))) + strconv.Itoa(3*fileBytes+lineBytes(t, report))) buf := startBuffer(t) @@ -182,9 +204,55 @@ func TestCapCountsReportFilesAlreadyInDataDir(t *testing.T) { 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 !errors.Is(err, reportbuf.ErrFull) { - t.Fatalf("report past the cap: error = %v, want ErrFull", err) + 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) + } } } @@ -195,8 +263,9 @@ 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", 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)) @@ -217,13 +286,22 @@ func TestWrittenReportsCountAtFileSize(t *testing.T) { 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") + } } -// 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 +// 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)) @@ -234,23 +312,24 @@ func TestWrittenReportsKeepCounting(t *testing.T) { 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) + for range reports { + before := 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", before, err) } - if err != nil { - t.Fatalf("with %d bytes of report files: %v", used, 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() @@ -258,8 +337,217 @@ func TestWrittenReportsKeepCounting(t *testing.T) { t.Fatalf("flush: %v", err) } } +} - t.Fatal("the report files never filled the cap") +// 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 the write of one report file +// open while a second write completes, then sends a report that needs +// room. Deleting either file would make it, and the one being written is +// the older, but only the complete one may be deleted. Once the first +// write is complete, its file is deleted when room is needed. +func TestFileBeingWrittenIsNeverDeleted(t *testing.T) { + report := map[string]string{"id": "writing"} + + t.Setenv("DATA_DIR", t.TempDir()) + // Room for two reports waiting to be written, but not for two + // beside a report file. + t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(2*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("first report: %v", err) + } + + flushed := make(chan error) + + go func() { flushed <- buf.Flush() }() + + writing := <-created + + // Only the first write is held; the second goes through, and + // so do the writes after it, the final one at stop included. + var complete string + + buf.OnFileCreated(func(f *os.File) { complete = f.Name() }) + + // Errorf, not Fatalf, until the first write is released, so that a + // failure here does not leave it held. + err = buf.Append(report) + if err != nil { + t.Errorf("second report: %v", err) + } + + err = buf.Flush() + if err != nil { + t.Errorf("flush of the second report: %v", err) + } + + err = buf.Append(report) + if err != nil { + t.Errorf("report that needs room: %v", err) + } + + if !exists(t, writing) { + t.Error("report file deleted while it was being written") + } + + if exists(t, complete) { + t.Error("complete report file kept, though room was needed") + } + + close(release) + + err = <-flushed + if err != nil { + t.Fatalf("flush of the first report: %v", err) + } + + err = buf.Append(report) + if err != nil { + t.Fatalf("report after the first 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 @@ -492,6 +780,14 @@ func readReportFiles(dir string) ([]string, error) { return contents, nil } +// 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() @@ -501,6 +797,21 @@ func writeBytes(t *testing.T, path string, n int) { } } +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()