Delete the oldest report files to stay under the size cap (closes #54)
check / check (push) Successful in 1m57s

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, at its place by name,
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
This commit was merged in pull request #90.
This commit is contained in:
2026-10-03 17:51:08 +02:00
parent 9e4d3fdb54
commit e4df415676
8 changed files with 638 additions and 109 deletions
+11 -1
View File
@@ -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
}
+166 -56
View File
@@ -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 (
@@ -12,6 +13,7 @@ import (
"log/slog"
"os"
"path/filepath"
"slices"
"strings"
"sync"
"sync/atomic"
@@ -37,9 +39,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 +52,34 @@ 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,
// in name order, which is oldest first: those in dataDir at
// start, and each one this buffer writes, put in at its place by
// name once it is complete, even when an older file's write
// completes after a newer one's. 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 +88,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 +105,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 +120,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 +167,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 +182,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 +213,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 +290,48 @@ 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.filesBytes += size
// At its place by name, not at the end: another write, of a newer
// file, may have completed while this one was being written.
i, _ := slices.BinarySearchFunc(b.files, name,
func(f reportFile, target string) int {
return strings.Compare(f.name, target)
})
b.files = slices.Insert(b.files, i, reportFile{name: name, size: 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 +340,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 +409,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
}
+427 -34
View File
@@ -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()
@@ -524,3 +835,85 @@ func hasReportFile(t *testing.T, dir string) bool {
return false
}
// TestFilesDeletedOldestFirstWhenWritesOverlap holds the write of an
// older report file open until a newer one's write completes, then
// releases it. When room is needed, the older file is deleted first,
// though its write was the last to complete.
func TestFilesDeletedOldestFirstWhenWritesOverlap(t *testing.T) {
const maxBytes = 1000
report := map[string]string{"id": "overlap"}
t.Setenv("DATA_DIR", t.TempDir())
t.Setenv("DATA_DIR_MAX_BYTES", strconv.Itoa(maxBytes))
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("older report: %v", err)
}
flushed := make(chan error)
go func() { flushed <- buf.Flush() }()
older := <-created
// Only the older write is held; the newer one goes through, and
// so do the writes after it, the final one at stop included.
var newer string
buf.OnFileCreated(func(f *os.File) { newer = f.Name() })
// Errorf, not Fatalf, until the older write is released, so that a
// failure here does not leave it held.
err = buf.Append(report)
if err != nil {
t.Errorf("newer report: %v", err)
}
err = buf.Flush()
if err != nil {
t.Errorf("flush of the newer report: %v", err)
}
close(release)
err = <-flushed
if err != nil {
t.Fatalf("flush of the older report: %v", err)
}
info, err := os.Stat(newer)
if err != nil {
t.Fatalf("stat %s: %v", newer, err)
}
// A report that fits beside the newer file alone, so deleting the
// older one makes exactly the room it needs.
pad := maxBytes - int(info.Size()) - lineBytes(t, map[string]string{"id": ""})
err = buf.Append(map[string]string{"id": strings.Repeat("a", pad)})
if err != nil {
t.Fatalf("report that needs room: %v", err)
}
if exists(t, older) {
t.Fatal("older report file kept when room was needed")
}
if !exists(t, newer) {
t.Fatal("newer report file deleted before the older one")
}
}