diff --git a/README.md b/README.md index acbb62d..b1620ac 100644 --- a/README.md +++ b/README.md @@ -212,7 +212,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 8b17dd2..d7637ca 100644 --- a/TODO.md +++ b/TODO.md @@ -23,6 +23,14 @@ 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 not yet written leave no + room for it on their own, and then no file is deleted. A file that cannot be + deleted still counts until the next start; one already deleted by hand counts + as freed - 2026-09-29: the container sets up its own data directory (issue #75): `bin/entrypoint.sh`, still as root, creates `DATA_DIR` if missing and gives it and `/data` to the `netwatch` user with mode 750 before starting the backend diff --git a/backend/README.md b/backend/README.md index 82ca05e..17537a5 100644 --- a/backend/README.md +++ b/backend/README.md @@ -80,7 +80,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) | @@ -148,11 +148,15 @@ 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. + 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 not yet written leave + no room for it 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 @@ -170,7 +174,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/reportbuf/reportbuf.go b/backend/internal/reportbuf/reportbuf.go index eeb0184..4494001 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,8 +38,9 @@ const ( fileSuffix = ".jsonl.zst" ) -// ErrFull is returned by Append when storing the report would -// take the report files past the configured maximum size. +// ErrFull is returned by Append when the reports not yet written +// leave no room for the report under the configured maximum size, +// however many report files are deleted. var ErrFull = errors.New("report files at their size cap") // Params defines the dependencies for Buffer. @@ -49,15 +51,28 @@ 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{} + // 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 @@ -97,12 +112,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 +159,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 +174,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 +205,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 not yet +// 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() { @@ -270,25 +335,28 @@ func (b *Buffer) writeFile(data []byte) error { } // 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. + // waited; now they count as the file, which from here on may be + // deleted to make room. 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.files = append(b.files, reportFile{name: name, size: info.Size()}) + b.filesBytes += info.Size() b.mu.Unlock() return 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 +367,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 0b56b7e..c5d58e3 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,124 @@ 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") + } +} + +// TestUncountedReportFileIsNeverDeleted: only the report files counted, +// those in DATA_DIR at start and those the buffer has finished +// writing, are deleted to make room. A file still being written is not +// among them, so it is kept even when that means refusing a report. A +// write cannot be paused here, so a report file put in DATA_DIR after +// start, which is not counted either, stands in for one being written. +func TestUncountedReportFileIsNeverDeleted(t *testing.T) { + report := map[string]string{"id": "uncounted"} + dir := t.TempDir() + uncounted := reportFilePath(dir, 1) + + t.Setenv("DATA_DIR", dir) + 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) + } + + writeBytes(t, uncounted, 100) + + err = buf.Append(report) + if !errors.Is(err, reportbuf.ErrFull) { + t.Fatalf("report past the cap: error = %v, want ErrFull", err) + } + + if !exists(t, uncounted) { + t.Fatal("report file deleted, though the buffer had not counted it") + } } // TestConcurrentAppendsStopAtCap appends from many goroutines at once @@ -407,6 +602,14 @@ func readReportFiles(t *testing.T, dir string) []string { 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() @@ -416,6 +619,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()